Skip to content

Commit a0648a9

Browse files
committed
feat(tencent): support cos bucket-dump
1 parent af3675d commit a0648a9

6 files changed

Lines changed: 454 additions & 0 deletions

File tree

pkg/providers/tencent/cos/client.go

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import (
1515
)
1616

1717
const defaultServiceEndpoint = "http://service.cos.myqcloud.com"
18+
const defaultBucketEndpointFormat = "https://%s.cos.%s.myqcloud.com"
1819

1920
type Option func(*Client)
2021

@@ -119,6 +120,72 @@ func (c *Client) ListBuckets(ctx context.Context) (*ListBucketsResponse, error)
119120
return &out, nil
120121
}
121122

123+
func (c *Client) ListObjects(ctx context.Context, bucket, region, marker string, maxKeys int) (ListObjectsResponse, error) {
124+
if ctx == nil {
125+
ctx = context.Background()
126+
}
127+
if err := c.credential.Validate(); err != nil {
128+
return ListObjectsResponse{}, err
129+
}
130+
bucket = strings.TrimSpace(bucket)
131+
if bucket == "" {
132+
return ListObjectsResponse{}, fmt.Errorf("tencent cos client: empty bucket")
133+
}
134+
region = strings.TrimSpace(region)
135+
if region == "" || region == "all" {
136+
return ListObjectsResponse{}, fmt.Errorf("tencent cos client: empty region")
137+
}
138+
if maxKeys <= 0 {
139+
maxKeys = 1000
140+
}
141+
142+
u, err := c.bucketURL(bucket, region)
143+
if err != nil {
144+
return ListObjectsResponse{}, err
145+
}
146+
query := u.Query()
147+
query.Set("max-keys", fmt.Sprintf("%d", maxKeys))
148+
if marker = strings.TrimSpace(marker); marker != "" {
149+
query.Set("marker", marker)
150+
}
151+
u.RawQuery = query.Encode()
152+
153+
httpResp, err := c.retryPolicy.Do(ctx, true, func() (*http.Response, error) {
154+
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil)
155+
if err != nil {
156+
return nil, err
157+
}
158+
if err := Sign(req, c.credential, c.now().UTC()); err != nil {
159+
return nil, err
160+
}
161+
return c.httpClient.Do(req)
162+
})
163+
if err != nil {
164+
return ListObjectsResponse{}, err
165+
}
166+
if httpResp == nil {
167+
return ListObjectsResponse{}, fmt.Errorf("tencent cos client: empty response")
168+
}
169+
defer closeResponse(httpResp)
170+
171+
body, err := io.ReadAll(httpResp.Body)
172+
if err != nil {
173+
return ListObjectsResponse{}, fmt.Errorf("read tencent cos response: %w", err)
174+
}
175+
if err := decodeError(httpResp, body); err != nil {
176+
return ListObjectsResponse{}, err
177+
}
178+
179+
var out ListObjectsResponse
180+
if len(body) == 0 {
181+
return out, nil
182+
}
183+
if err := xml.Unmarshal(body, &out); err != nil {
184+
return ListObjectsResponse{}, fmt.Errorf("decode tencent cos response: %w", err)
185+
}
186+
return out, nil
187+
}
188+
122189
func (c *Client) serviceURL() (*url.URL, error) {
123190
rawURL := strings.TrimSpace(c.serviceEndpoint)
124191
if rawURL == "" {
@@ -137,6 +204,21 @@ func (c *Client) serviceURL() (*url.URL, error) {
137204
return u, nil
138205
}
139206

207+
func (c *Client) bucketURL(bucket, region string) (*url.URL, error) {
208+
rawURL := fmt.Sprintf(defaultBucketEndpointFormat, bucket, region)
209+
u, err := url.Parse(rawURL)
210+
if err != nil {
211+
return nil, fmt.Errorf("tencent cos client: invalid bucket endpoint %q: %w", rawURL, err)
212+
}
213+
if u.Scheme == "" || u.Host == "" {
214+
return nil, fmt.Errorf("tencent cos client: invalid bucket endpoint %q", rawURL)
215+
}
216+
if strings.TrimSpace(u.Path) == "" {
217+
u.Path = "/"
218+
}
219+
return u, nil
220+
}
221+
140222
func closeResponse(resp *http.Response) {
141223
if resp == nil || resp.Body == nil {
142224
return
Lines changed: 105 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,105 @@
1+
package cos
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"strings"
7+
8+
"github.com/404tk/cloudtoolkit/utils"
9+
"github.com/404tk/cloudtoolkit/utils/logger"
10+
"github.com/404tk/cloudtoolkit/utils/processbar"
11+
)
12+
13+
func (d *Driver) ListObjects(ctx context.Context, buckets map[string]string) {
14+
for bucket, region := range buckets {
15+
resp, err := d.listObjectsPage(ctx, bucket, region, "", 100)
16+
if err != nil {
17+
logger.Error(fmt.Sprintf("List Objects in %s failed: %s", bucket, err))
18+
continue
19+
}
20+
21+
if len(resp.Objects) == 0 {
22+
logger.Error(fmt.Sprintf("No Objects found in %s.", bucket))
23+
continue
24+
}
25+
logger.Warning(fmt.Sprintf("%d objects found in %s.", len(resp.Objects), bucket))
26+
27+
fmt.Printf("\n%-70s\t%-10s\n", "Key", "Size")
28+
fmt.Printf("%-70s\t%-10s\n", "---", "----")
29+
for _, object := range resp.Objects {
30+
fmt.Printf("%-70s\t%-10s\n", object.Key, utils.ParseBytes(object.Size))
31+
}
32+
fmt.Println()
33+
34+
select {
35+
case <-ctx.Done():
36+
return
37+
default:
38+
}
39+
}
40+
}
41+
42+
func (d *Driver) TotalObjects(ctx context.Context, buckets map[string]string) {
43+
tracker := processbar.NewCountTracker()
44+
defer tracker.Finish()
45+
46+
for bucket, region := range buckets {
47+
count, err := d.countBucketObjects(ctx, bucket, region, tracker)
48+
if err != nil {
49+
logger.Error(fmt.Sprintf("List Objects in %s failed: %s", bucket, err))
50+
return
51+
}
52+
fmt.Printf("\r")
53+
logger.Warning(fmt.Sprintf("%s has %d objects.", bucket, count))
54+
}
55+
}
56+
57+
func (d *Driver) listObjectsPage(ctx context.Context, bucket, region, marker string, maxKeys int) (ListObjectsResponse, error) {
58+
client := d.client()
59+
if client == nil {
60+
return ListObjectsResponse{}, fmt.Errorf("tencent cos: nil client")
61+
}
62+
return client.ListObjects(ctx, bucket, normalizeRegion(region), marker, maxKeys)
63+
}
64+
65+
func (d *Driver) countBucketObjects(ctx context.Context, bucket, region string, tracker *processbar.CountTracker) (int, error) {
66+
count := 0
67+
marker := ""
68+
for {
69+
resp, err := d.listObjectsPage(ctx, bucket, region, marker, 1000)
70+
if err != nil {
71+
return 0, err
72+
}
73+
count += len(resp.Objects)
74+
if tracker != nil {
75+
tracker.Update(bucket, count)
76+
}
77+
if !resp.IsTruncated {
78+
return count, nil
79+
}
80+
81+
next := nextMarker(resp)
82+
if next == "" {
83+
return 0, fmt.Errorf("tencent cos: truncated response for bucket %s missing continuation marker", bucket)
84+
}
85+
marker = next
86+
}
87+
}
88+
89+
func nextMarker(resp ListObjectsResponse) string {
90+
if marker := strings.TrimSpace(resp.NextMarker); marker != "" {
91+
return marker
92+
}
93+
if len(resp.Objects) == 0 {
94+
return ""
95+
}
96+
return strings.TrimSpace(resp.Objects[len(resp.Objects)-1].Key)
97+
}
98+
99+
func normalizeRegion(region string) string {
100+
region = strings.TrimSpace(region)
101+
if region == "" || region == "all" {
102+
return ""
103+
}
104+
return region
105+
}
Lines changed: 133 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,133 @@
1+
package cos
2+
3+
import (
4+
"context"
5+
"io"
6+
"net/http"
7+
"strings"
8+
"testing"
9+
"time"
10+
11+
"github.com/404tk/cloudtoolkit/pkg/providers/tencent/api"
12+
"github.com/404tk/cloudtoolkit/pkg/providers/tencent/auth"
13+
)
14+
15+
func TestClientListObjectsUsesBucketScopedEndpoint(t *testing.T) {
16+
client := NewClient(
17+
auth.New("AKIDEXAMPLE", "SECRETKEYEXAMPLE", "TOKENEXAMPLE"),
18+
WithHTTPClient(&http.Client{
19+
Transport: roundTripFunc(func(r *http.Request) (*http.Response, error) {
20+
if r.Method != http.MethodGet {
21+
t.Fatalf("unexpected method: %s", r.Method)
22+
}
23+
if r.URL.Host != "examplebucket-1250000000.cos.ap-guangzhou.myqcloud.com" {
24+
t.Fatalf("unexpected host: %s", r.URL.Host)
25+
}
26+
if r.URL.Path != "/" {
27+
t.Fatalf("unexpected path: %s", r.URL.Path)
28+
}
29+
if got := r.URL.Query().Get("max-keys"); got != "100" {
30+
t.Fatalf("unexpected max-keys: %s", got)
31+
}
32+
if got := r.URL.Query().Get("marker"); got != "" {
33+
t.Fatalf("unexpected marker: %s", got)
34+
}
35+
if r.Header.Get("x-cos-security-token") != "TOKENEXAMPLE" {
36+
t.Fatalf("unexpected token header: %q", r.Header.Get("x-cos-security-token"))
37+
}
38+
authHeader := r.Header.Get("Authorization")
39+
if authHeader == "" {
40+
t.Fatal("missing authorization header")
41+
}
42+
if !strings.Contains(authHeader, "q-url-param-list=max-keys") {
43+
t.Fatalf("unexpected authorization header: %s", authHeader)
44+
}
45+
return &http.Response{
46+
StatusCode: http.StatusOK,
47+
Header: make(http.Header),
48+
Body: io.NopCloser(strings.NewReader(`
49+
<ListBucketResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/">
50+
<Name>examplebucket-1250000000</Name>
51+
<MaxKeys>100</MaxKeys>
52+
<IsTruncated>false</IsTruncated>
53+
<Contents>
54+
<Key>alpha.txt</Key>
55+
<Size>12</Size>
56+
</Contents>
57+
</ListBucketResult>`)),
58+
Request: r,
59+
}, nil
60+
}),
61+
}),
62+
WithRetryPolicy(api.RetryPolicy{MaxAttempts: 1}),
63+
WithClock(func() time.Time { return time.Date(2026, 4, 19, 12, 0, 0, 0, time.UTC) }),
64+
)
65+
66+
resp, err := client.ListObjects(context.Background(), "examplebucket-1250000000", "ap-guangzhou", "", 100)
67+
if err != nil {
68+
t.Fatalf("ListObjects() error = %v", err)
69+
}
70+
if resp.Name != "examplebucket-1250000000" || len(resp.Objects) != 1 || resp.Objects[0].Key != "alpha.txt" || resp.Objects[0].Size != 12 {
71+
t.Fatalf("unexpected response: %+v", resp)
72+
}
73+
}
74+
75+
func TestDriverCountBucketObjectsPaginatesMarker(t *testing.T) {
76+
client := NewClient(
77+
auth.New("AKIDEXAMPLE", "SECRETKEYEXAMPLE", ""),
78+
WithHTTPClient(&http.Client{
79+
Transport: roundTripFunc(func(r *http.Request) (*http.Response, error) {
80+
if r.URL.Host != "examplebucket-1250000000.cos.ap-guangzhou.myqcloud.com" {
81+
t.Fatalf("unexpected host: %s", r.URL.Host)
82+
}
83+
switch r.URL.RawQuery {
84+
case "max-keys=1000":
85+
return &http.Response{
86+
StatusCode: http.StatusOK,
87+
Header: make(http.Header),
88+
Body: io.NopCloser(strings.NewReader(`
89+
<ListBucketResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/">
90+
<Name>examplebucket-1250000000</Name>
91+
<MaxKeys>1000</MaxKeys>
92+
<IsTruncated>true</IsTruncated>
93+
<Contents><Key>a.txt</Key><Size>1</Size></Contents>
94+
<Contents><Key>b.txt</Key><Size>2</Size></Contents>
95+
</ListBucketResult>`)),
96+
Request: r,
97+
}, nil
98+
case "marker=b.txt&max-keys=1000", "max-keys=1000&marker=b.txt":
99+
return &http.Response{
100+
StatusCode: http.StatusOK,
101+
Header: make(http.Header),
102+
Body: io.NopCloser(strings.NewReader(`
103+
<ListBucketResult xmlns="http://s3.amazonaws.com/doc/2006-03-01/">
104+
<Name>examplebucket-1250000000</Name>
105+
<MaxKeys>1000</MaxKeys>
106+
<IsTruncated>false</IsTruncated>
107+
<Contents><Key>c.txt</Key><Size>3</Size></Contents>
108+
</ListBucketResult>`)),
109+
Request: r,
110+
}, nil
111+
default:
112+
t.Fatalf("unexpected query: %s", r.URL.RawQuery)
113+
return nil, nil
114+
}
115+
}),
116+
}),
117+
WithRetryPolicy(api.RetryPolicy{MaxAttempts: 1}),
118+
WithClock(func() time.Time { return time.Date(2026, 4, 19, 12, 0, 0, 0, time.UTC) }),
119+
)
120+
121+
driver := &Driver{
122+
Credential: auth.New("AKIDEXAMPLE", "SECRETKEYEXAMPLE", ""),
123+
Client: client,
124+
}
125+
126+
count, err := driver.countBucketObjects(context.Background(), "examplebucket-1250000000", "ap-guangzhou", nil)
127+
if err != nil {
128+
t.Fatalf("countBucketObjects() error = %v", err)
129+
}
130+
if count != 3 {
131+
t.Fatalf("unexpected object count: %d", count)
132+
}
133+
}

pkg/providers/tencent/cos/types.go

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,22 @@ type COSBucket struct {
1313
CreationDate string `xml:"CreationDate"`
1414
}
1515

16+
type ListObjectsResponse struct {
17+
XMLName xml.Name `xml:"ListBucketResult"`
18+
Name string `xml:"Name"`
19+
Prefix string `xml:"Prefix"`
20+
Marker string `xml:"Marker"`
21+
NextMarker string `xml:"NextMarker"`
22+
MaxKeys int `xml:"MaxKeys"`
23+
IsTruncated bool `xml:"IsTruncated"`
24+
Objects []COSObject `xml:"Contents"`
25+
}
26+
27+
type COSObject struct {
28+
Key string `xml:"Key"`
29+
Size int64 `xml:"Size"`
30+
}
31+
1632
type errorResponse struct {
1733
XMLName xml.Name `xml:"Error"`
1834
Code string `xml:"Code"`

0 commit comments

Comments
 (0)