Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions drivers/halalcloud_open/driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ type HalalCloudOpen struct {
sdkClient *sdkClient.Client
sdkUserFileService *sdkUserFile.UserFileService
sdkUserService *sdkUser.UserService
offlineTaskService offlineTaskService
uploadThread int
}

Expand Down
12 changes: 11 additions & 1 deletion drivers/halalcloud_open/driver_init.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,12 @@ package halalcloudopen

import (
"context"
"net/http"
"time"

"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/halalcloud/golang-sdk-lite/halalcloud/apiclient"
sdkOffline "github.com/halalcloud/golang-sdk-lite/halalcloud/services/offline"
sdkUser "github.com/halalcloud/golang-sdk-lite/halalcloud/services/user"
sdkUserFile "github.com/halalcloud/golang-sdk-lite/halalcloud/services/userfile"
)
Expand Down Expand Up @@ -36,10 +38,18 @@ func (d *HalalCloudOpen) Init(ctx context.Context) error {
host = "openapi.2dland.cn"
}

client := apiclient.NewClient(nil, host, d.Addition.ClientID, d.Addition.ClientSecret, d.halalCommon, apiclient.WithTimeout(time.Second*time.Duration(timeout)))
// All SDK services share the quota associated with these credentials.
httpClient := &http.Client{
Transport: &rateLimitedTransport{
base: http.DefaultTransport,
limiter: halalCloudAPILimiter(host, d.Addition.ClientID),
},
}
client := apiclient.NewClient(httpClient, host, d.Addition.ClientID, d.Addition.ClientSecret, d.halalCommon, apiclient.WithTimeout(time.Second*time.Duration(timeout)))
d.sdkClient = client
d.sdkUserFileService = sdkUserFile.NewUserFileService(client)
d.sdkUserService = sdkUser.NewUserService(client)
d.offlineTaskService = sdkOffline.NewOfflineTaskService(client)
userInfo, err := d.sdkUserService.Get(ctx, &sdkUser.User{})
if err != nil {
return err
Expand Down
118 changes: 118 additions & 0 deletions drivers/halalcloud_open/offline.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
package halalcloudopen

import (
"context"
"errors"
"fmt"
"strings"

"github.com/OpenListTeam/OpenList/v4/internal/model"
sdkModel "github.com/halalcloud/golang-sdk-lite/halalcloud/model"
sdkOffline "github.com/halalcloud/golang-sdk-lite/halalcloud/services/offline"
)

// A larger page reduces requests during status polling.
const offlineTaskListPageSize int64 = 200

type offlineTaskService interface {
Add(ctx context.Context, req *sdkOffline.UserTask) (*sdkOffline.UserTask, error)
List(ctx context.Context, req *sdkOffline.OfflineTaskListRequest) (*sdkOffline.OfflineTaskListResponse, error)
Delete(ctx context.Context, req *sdkOffline.OfflineTaskDeleteRequest) (*sdkOffline.OfflineTaskDeleteResponse, error)
}

// OfflineDownload creates a URL-based task that writes directly into parentDir.
func (d *HalalCloudOpen) OfflineDownload(ctx context.Context, fileURL string, parentDir model.Obj) (*sdkOffline.UserTask, error) {
if d.offlineTaskService == nil {
return nil, errors.New("HalalCloudOpen offline task service is not initialized")
}

fileURL = strings.TrimSpace(fileURL)
if fileURL == "" {
return nil, errors.New("HalalCloudOpen offline download URL is empty")
}

// Let HalalCloud dispatch the URL to its supported task type.
task, err := d.offlineTaskService.Add(ctx, &sdkOffline.UserTask{
Url: fileURL,
SavePath: parentDir.GetPath(),
})
if err != nil {
return nil, fmt.Errorf("failed to create HalalCloudOpen offline task: %w", err)
}
if task == nil || strings.TrimSpace(task.Identity) == "" {
return nil, errors.New("failed to create HalalCloudOpen offline task: empty task identity")
}
// Identity is the user-task handle shared by Add, List, and Delete.
return task, nil
}

// OfflineList returns every user task, following the opaque pagination token
// until the API reports that traversal is complete.
func (d *HalalCloudOpen) OfflineList(ctx context.Context) ([]*sdkOffline.UserTask, error) {
if d.offlineTaskService == nil {
return nil, errors.New("HalalCloudOpen offline task service is not initialized")
}

tasks := make([]*sdkOffline.UserTask, 0)
token := ""
seenTokens := make(map[string]struct{})
for {
resp, err := d.offlineTaskService.List(ctx, &sdkOffline.OfflineTaskListRequest{
ListInfo: &sdkModel.ScanListRequest{
Limit: offlineTaskListPageSize,
Token: token,
},
})
if err != nil {
return nil, fmt.Errorf("failed to list HalalCloudOpen offline tasks: %w", err)
}
if resp == nil {
return nil, errors.New("failed to list HalalCloudOpen offline tasks: empty response")
}
tasks = append(tasks, resp.Tasks...)

if resp.ListInfo == nil || resp.ListInfo.Token == "" {
break
}
// ListInfo.Token is opaque. Guard against a malformed response repeating
// a token so one status poll cannot loop forever.
nextToken := resp.ListInfo.Token
if nextToken == token {
return nil, errors.New("failed to list HalalCloudOpen offline tasks: pagination token did not advance")
}
if _, ok := seenTokens[nextToken]; ok {
return nil, errors.New("failed to list HalalCloudOpen offline tasks: pagination token repeated")
}
seenTokens[nextToken] = struct{}{}
token = nextToken
}
return tasks, nil
}

// DeleteOfflineTasks removes task records and optionally their downloaded files.
func (d *HalalCloudOpen) DeleteOfflineTasks(ctx context.Context, taskIDs []string, deleteFiles bool) error {
if d.offlineTaskService == nil {
return errors.New("HalalCloudOpen offline task service is not initialized")
}

identities := make([]string, 0, len(taskIDs))
for _, taskID := range taskIDs {
if taskID = strings.TrimSpace(taskID); taskID != "" {
identities = append(identities, taskID)
}
}
if len(identities) == 0 {
return nil
}

_, err := d.offlineTaskService.Delete(ctx, &sdkOffline.OfflineTaskDeleteRequest{
Identity: identities,
DeleteFiles: deleteFiles,
})
if err != nil {
return fmt.Errorf("failed to delete HalalCloudOpen offline tasks: %w", err)
}
return nil
}

var _ offlineTaskService = (*sdkOffline.OfflineTaskService)(nil)
188 changes: 188 additions & 0 deletions drivers/halalcloud_open/offline_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,188 @@
package halalcloudopen

import (
"context"
"errors"
"strings"
"testing"

"github.com/OpenListTeam/OpenList/v4/internal/model"
sdkModel "github.com/halalcloud/golang-sdk-lite/halalcloud/model"
sdkOffline "github.com/halalcloud/golang-sdk-lite/halalcloud/services/offline"
)

type fakeOfflineTaskService struct {
add func(context.Context, *sdkOffline.UserTask) (*sdkOffline.UserTask, error)
list func(context.Context, *sdkOffline.OfflineTaskListRequest) (*sdkOffline.OfflineTaskListResponse, error)
delete func(context.Context, *sdkOffline.OfflineTaskDeleteRequest) (*sdkOffline.OfflineTaskDeleteResponse, error)
}

func (f *fakeOfflineTaskService) Add(ctx context.Context, req *sdkOffline.UserTask) (*sdkOffline.UserTask, error) {
if f.add == nil {
return nil, errors.New("unexpected Add call")
}
return f.add(ctx, req)
}

func (f *fakeOfflineTaskService) List(ctx context.Context, req *sdkOffline.OfflineTaskListRequest) (*sdkOffline.OfflineTaskListResponse, error) {
if f.list == nil {
return nil, errors.New("unexpected List call")
}
return f.list(ctx, req)
}

func (f *fakeOfflineTaskService) Delete(ctx context.Context, req *sdkOffline.OfflineTaskDeleteRequest) (*sdkOffline.OfflineTaskDeleteResponse, error) {
if f.delete == nil {
return nil, errors.New("unexpected Delete call")
}
return f.delete(ctx, req)
}

func TestOfflineDownload(t *testing.T) {
var got *sdkOffline.UserTask
driver := &HalalCloudOpen{
offlineTaskService: &fakeOfflineTaskService{
add: func(_ context.Context, req *sdkOffline.UserTask) (*sdkOffline.UserTask, error) {
got = req
return &sdkOffline.UserTask{Identity: "task-id"}, nil
},
},
}
parent := &model.Object{Path: "/downloads", IsFolder: true}

task, err := driver.OfflineDownload(context.Background(), " https://example.com/file ", parent)
if err != nil {
t.Fatalf("OfflineDownload() error = %v", err)
}
if task.Identity != "task-id" {
t.Fatalf("OfflineDownload() identity = %q, want %q", task.Identity, "task-id")
}
if got == nil {
t.Fatal("OfflineDownload() did not call the SDK service")
}
if got.Url != "https://example.com/file" {
t.Errorf("Add request URL = %q, want trimmed task URL", got.Url)
}
if got.SavePath != "/downloads" {
t.Errorf("Add request SavePath = %q, want %q", got.SavePath, "/downloads")
}
}

func TestOfflineDownloadRejectsEmptyURL(t *testing.T) {
called := false
driver := &HalalCloudOpen{
offlineTaskService: &fakeOfflineTaskService{
add: func(_ context.Context, _ *sdkOffline.UserTask) (*sdkOffline.UserTask, error) {
called = true
return nil, nil
},
},
}
parent := &model.Object{Path: "/downloads", IsFolder: true}

_, err := driver.OfflineDownload(context.Background(), " ", parent)
if err == nil || !strings.Contains(err.Error(), "URL is empty") {
t.Fatalf("OfflineDownload() error = %v, want empty URL error", err)
}
if called {
t.Fatal("OfflineDownload() called the SDK service for an empty URL")
}
}

func TestOfflineDownloadRequiresTaskIdentity(t *testing.T) {
driver := &HalalCloudOpen{
offlineTaskService: &fakeOfflineTaskService{
add: func(_ context.Context, _ *sdkOffline.UserTask) (*sdkOffline.UserTask, error) {
return &sdkOffline.UserTask{}, nil
},
},
}

_, err := driver.OfflineDownload(context.Background(), "https://example.com/file", &model.Object{Path: "/"})
if err == nil || !strings.Contains(err.Error(), "empty task identity") {
t.Fatalf("OfflineDownload() error = %v, want empty identity error", err)
}
}

func TestOfflineListPaginates(t *testing.T) {
var requests []*sdkOffline.OfflineTaskListRequest
driver := &HalalCloudOpen{
offlineTaskService: &fakeOfflineTaskService{
list: func(_ context.Context, req *sdkOffline.OfflineTaskListRequest) (*sdkOffline.OfflineTaskListResponse, error) {
requests = append(requests, req)
if req.ListInfo.Token == "" {
return &sdkOffline.OfflineTaskListResponse{
Tasks: []*sdkOffline.UserTask{{Identity: "first"}},
ListInfo: &sdkModel.ScanListRequest{Token: "next"},
}, nil
}
return &sdkOffline.OfflineTaskListResponse{
Tasks: []*sdkOffline.UserTask{{Identity: "second"}},
ListInfo: &sdkModel.ScanListRequest{},
}, nil
},
},
}

tasks, err := driver.OfflineList(context.Background())
if err != nil {
t.Fatalf("OfflineList() error = %v", err)
}
if len(tasks) != 2 || tasks[0].Identity != "first" || tasks[1].Identity != "second" {
t.Fatalf("OfflineList() tasks = %#v, want first and second tasks", tasks)
}
if len(requests) != 2 {
t.Fatalf("OfflineList() request count = %d, want 2", len(requests))
}
if requests[0].ListInfo.Limit != offlineTaskListPageSize || requests[0].ListInfo.Token != "" {
t.Errorf("first List request = %#v, want initial page", requests[0].ListInfo)
}
if requests[1].ListInfo.Limit != offlineTaskListPageSize || requests[1].ListInfo.Token != "next" {
t.Errorf("second List request = %#v, want next page", requests[1].ListInfo)
}
}

func TestOfflineListRejectsRepeatedToken(t *testing.T) {
driver := &HalalCloudOpen{
offlineTaskService: &fakeOfflineTaskService{
list: func(_ context.Context, req *sdkOffline.OfflineTaskListRequest) (*sdkOffline.OfflineTaskListResponse, error) {
if req.ListInfo.Token == "" {
return &sdkOffline.OfflineTaskListResponse{ListInfo: &sdkModel.ScanListRequest{Token: "same"}}, nil
}
return &sdkOffline.OfflineTaskListResponse{ListInfo: &sdkModel.ScanListRequest{Token: "same"}}, nil
},
},
}

_, err := driver.OfflineList(context.Background())
if err == nil || !strings.Contains(err.Error(), "pagination token did not advance") {
t.Fatalf("OfflineList() error = %v, want repeated token error", err)
}
}

func TestDeleteOfflineTasks(t *testing.T) {
var got *sdkOffline.OfflineTaskDeleteRequest
driver := &HalalCloudOpen{
offlineTaskService: &fakeOfflineTaskService{
delete: func(_ context.Context, req *sdkOffline.OfflineTaskDeleteRequest) (*sdkOffline.OfflineTaskDeleteResponse, error) {
got = req
return &sdkOffline.OfflineTaskDeleteResponse{Count: 2}, nil
},
},
}

if err := driver.DeleteOfflineTasks(context.Background(), []string{" first ", "", "second"}, false); err != nil {
t.Fatalf("DeleteOfflineTasks() error = %v", err)
}
if got == nil {
t.Fatal("DeleteOfflineTasks() did not call the SDK service")
}
if len(got.Identity) != 2 || got.Identity[0] != "first" || got.Identity[1] != "second" {
t.Errorf("Delete request identities = %#v, want trimmed non-empty identities", got.Identity)
}
if got.DeleteFiles {
t.Error("Delete request DeleteFiles = true, want false")
}
}

var _ offlineTaskService = (*fakeOfflineTaskService)(nil)
37 changes: 37 additions & 0 deletions drivers/halalcloud_open/rate_limit.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package halalcloudopen

import (
"net/http"
"strings"
"sync"
"time"

"golang.org/x/time/rate"
)

const halalCloudAPIRequestInterval = time.Second

// API quotas are shared by mounts that use the same credentials.
var halalCloudAPILimiters sync.Map

func halalCloudAPILimiter(host, clientID string) *rate.Limiter {
key := strings.ToLower(strings.TrimSpace(host)) + "\x00" + strings.TrimSpace(clientID)
limiter, _ := halalCloudAPILimiters.LoadOrStore(
key,
rate.NewLimiter(rate.Every(halalCloudAPIRequestInterval), 1),
)
return limiter.(*rate.Limiter)
}

// rateLimitedTransport applies the credential quota to every SDK service.
type rateLimitedTransport struct {
base http.RoundTripper
limiter *rate.Limiter
}

func (t *rateLimitedTransport) RoundTrip(req *http.Request) (*http.Response, error) {
if err := t.limiter.Wait(req.Context()); err != nil {
return nil, err
}
return t.base.RoundTrip(req)
}
Loading