diff --git a/pkg/storage/fs/kiteworks/kiteworks.go b/pkg/storage/fs/kiteworks/kiteworks.go index d4592595e14..805077b3fde 100644 --- a/pkg/storage/fs/kiteworks/kiteworks.go +++ b/pkg/storage/fs/kiteworks/kiteworks.go @@ -1,14 +1,21 @@ package kiteworks import ( + "bytes" "context" + "encoding/json" "errors" + "fmt" "io" + "math" "net/http" "net/url" + "path" "strings" + "sync" provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" + types "github.com/cs3org/go-cs3apis/cs3/types/v1beta1" "github.com/mitchellh/mapstructure" "github.com/rs/zerolog" @@ -39,6 +46,8 @@ type Driver struct { apiToken string storageID string log zerolog.Logger + // TODO: lock tokens are ephemeral; lost on restart and not shared across nodes. + locks sync.Map // nodeID → *provider.Lock } // New returns a read-only Kiteworks storage driver. @@ -83,22 +92,30 @@ func (d *Driver) toResourceInfo(fi *kwlib.FileInfo, spaceID, spaceRootPath strin if !strings.HasPrefix(relPath, "/") { relPath = "/" + relPath } + spaceRoot := &provider.ResourceId{ + StorageId: d.storageID, + SpaceId: spaceID, + OpaqueId: spaceID, + } ri := &provider.ResourceInfo{ Id: &provider.ResourceId{ StorageId: d.storageID, SpaceId: spaceID, OpaqueId: fi.ID, }, - Name: fi.Name, - Path: relPath, - Etag: fi.ETag(), - Mtime: utils.TimeToTS(fi.MTime()), - PermissionSet: &provider.ResourcePermissions{ - Stat: true, - GetPath: true, - InitiateFileDownload: fi.IsFile(), - ListContainer: fi.IsDir(), - }, + Space: &provider.StorageSpace{Root: spaceRoot}, + Name: fi.Name, + Path: relPath, + Etag: fi.ETag(), + Mtime: utils.TimeToTS(fi.MTime()), + PermissionSet: permissionSet(fi), + } + if fi.ParentID != nil && *fi.ParentID != "" && *fi.ParentID != "0" { + ri.ParentId = &provider.ResourceId{ + StorageId: d.storageID, + SpaceId: spaceID, + OpaqueId: *fi.ParentID, + } } if fi.IsDir() { ri.Type = provider.ResourceType_RESOURCE_TYPE_CONTAINER @@ -112,11 +129,17 @@ func (d *Driver) toResourceInfo(fi *kwlib.FileInfo, spaceID, spaceRootPath strin return ri } -// Capabilities declares kiteworks read-only: every write method rejects with -// NotSupported. The zero value is the declaration, so any capability added later -// stays false here without an edit. +// Capabilities omits Trash, Sharing and ArbitraryMetadata: those methods still +// reject with NotSupported. func (d *Driver) Capabilities(_ context.Context) storage.Capabilities { - return storage.Capabilities{} + return storage.Capabilities{ + Upload: true, + CreateContainer: true, + Delete: true, + Move: true, + Versioning: true, + Locking: true, + } } // --- Read methods --- @@ -130,9 +153,20 @@ func (d *Driver) ListStorageSpaces(ctx context.Context, _ []*provider.ListStorag return nil, err } + u, hasUser := ctxpkg.ContextGetUser(ctx) + spaces := make([]*provider.StorageSpace, 0, len(dirs.Data)) for i := range dirs.Data { fi := &dirs.Data[i] + opaque := utils.AppendPlainToOpaque(nil, "spaceAlias", "project/"+fi.Name) + if hasUser && u.GetId().GetOpaqueId() != "" { + grants := map[string]*provider.ResourcePermissions{ + u.Id.OpaqueId: spaceRole(fi), + } + if b, err := json.Marshal(grants); err == nil { + opaque.Map["grants"] = &types.OpaqueEntry{Decoder: "json", Value: b} + } + } spaces = append(spaces, &provider.StorageSpace{ Id: &provider.StorageSpaceId{OpaqueId: fi.ID}, Name: fi.Name, @@ -144,7 +178,7 @@ func (d *Driver) ListStorageSpaces(ctx context.Context, _ []*provider.ListStorag }, RootInfo: d.toResourceInfo(fi, fi.ID, fi.Path), Mtime: utils.TimeToTS(fi.MTime()), - Opaque: utils.AppendPlainToOpaque(nil, "spaceAlias", "project/"+fi.Name), + Opaque: opaque, }) } return spaces, nil @@ -158,7 +192,6 @@ func (d *Driver) resolveRef(ctx context.Context, ref *provider.Reference) (nodeI if nodeID == "" { nodeID = spaceID } - relPath := strings.Trim(strings.TrimPrefix(ref.GetPath(), "./"), "/.") if relPath == "" { return nodeID, spaceID, nil @@ -199,6 +232,9 @@ func (d *Driver) spaceRootPath(ctx context.Context, spaceID string) (string, err } func (d *Driver) GetMD(ctx context.Context, ref *provider.Reference, _, _ []string) (*provider.ResourceInfo, error) { + if vr, ok := parseVersionRef(ref); ok { + return d.getVersionMD(ctx, ref, vr) + } nodeID, spaceID, err := d.resolveRef(ctx, ref) if err != nil { return nil, err @@ -251,7 +287,7 @@ func (d *Driver) ListFolder(ctx context.Context, ref *provider.Reference, _, _ [ for i := range items { ri := d.toResourceInfo(&items[i], spaceID, rootPath) // ocdav does path.Join(requestPath, info.Path) for ListFolder results, - // so Path must be just the filename — not the space-root-relative path. + // so Path must be just the filename, not the space-root-relative path. ri.Path = ri.Name infos = append(infos, ri) } @@ -259,6 +295,9 @@ func (d *Driver) ListFolder(ctx context.Context, ref *provider.Reference, _, _ [ } func (d *Driver) Download(ctx context.Context, ref *provider.Reference, openReaderFunc func(*provider.ResourceInfo) bool) (*provider.ResourceInfo, io.ReadCloser, error) { + if vr, ok := parseVersionRef(ref); ok { + return d.downloadVersion(ctx, ref, vr, openReaderFunc) + } nodeID, spaceID, err := d.resolveRef(ctx, ref) if err != nil { return nil, nil, err @@ -300,31 +339,63 @@ func (d *Driver) ListGrants(_ context.Context, _ *provider.Reference) ([]*provid return []*provider.Grant{}, nil } -func (d *Driver) GetQuota(ctx context.Context, _ *provider.Reference) (uint64, uint64, uint64, error) { - q, err := d.client(ctx).GetQuotaInfo() +// TODO: this maps only one of Kiteworks' two quota models. +// +// 1. System quota: a per-folder limit on a top-level folder. Matches the oCIS per-space +// quota. +// 2. Total folder quota: one pool shared by every folder the account owns. No oCIS +// equivalent. +// +// The choice is per-folder and defaults to total folder quota. So out of the box every +// space reports the same total and used (the account pool), and an unlimited pool makes +// every space report unrestricted. +// +// Both come back under the same field names (storage_quota, storage_used, +// storage_available), so this endpoint cannot tell them apart. Read useFolderQuota from +// GET /rest/folders/{id} and treat the pooled case as "no real space quota". +// +// Enabling system quota: +// - Admin > Users > Profiles > Standard > Collaboration > Enable Folders That Use System Quota +// - Per folder, at creation time: Folder quota > Use system storage quota +// +// Total folder quota is set on the same profile page. +func (d *Driver) GetQuota(ctx context.Context, ref *provider.Reference) (uint64, uint64, uint64, error) { + nodeID, _, err := d.resolveRef(ctx, ref) + if err != nil || nodeID == "" { + return 0, 0, math.MaxUint64, nil + } + q, err := d.client(ctx).GetFolderQuota(nodeID) if err != nil { - // Non-fatal for read-only driver; return zero quota rather than failing. - return 0, 0, 0, nil + // The endpoint requires file_add on the folder, so viewers get a 403. + // Quota is informational, so report unlimited rather than fail the caller. + return 0, 0, math.MaxUint64, nil } - total := uint64(q.FolderQuotaAllowed) - used := uint64(q.FolderQuotaUsed) - var remaining uint64 - if total > used { - remaining = total - used + total := uint64(max(q.StorageQuota, 0)) + used := uint64(max(q.StorageUsed, 0)) + // total=0 means no quota is applied to the folder. Signal unlimited remaining + // so the graph service doesn't compute 0/0 = NaN and report "exceeded". + if total == 0 { + return 0, used, math.MaxUint64, nil } - return total, used, remaining, nil -} - -func (d *Driver) GetLock(_ context.Context, _ *provider.Reference) (*provider.Lock, error) { - return nil, nil + return total, used, uint64(max(q.StorageAvailable, 0)), nil } -func (d *Driver) ListRevisions(_ context.Context, _ *provider.Reference) ([]*provider.FileVersion, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") -} - -func (d *Driver) DownloadRevision(_ context.Context, _ *provider.Reference, _ string, _ func(*provider.ResourceInfo) bool) (*provider.ResourceInfo, io.ReadCloser, error) { - return nil, nil, errtypes.NotSupported("kiteworks: read-only driver") +func (d *Driver) GetLock(ctx context.Context, ref *provider.Reference) (*provider.Lock, error) { + nodeID, _, err := d.resolveRef(ctx, ref) + if err != nil { + return nil, err + } + fi, err := d.client(ctx).GetFileByID(nodeID) + if err != nil { + return nil, err + } + if !fi.Locked { + return nil, nil + } + if v, ok := d.locks.Load(nodeID); ok { + return v.(*provider.Lock), nil + } + return &provider.Lock{LockId: "kw-" + nodeID, Type: provider.LockType_LOCK_TYPE_EXCL}, nil } func (d *Driver) ListRecycle(_ context.Context, _ *provider.Reference, _, _ string) ([]*provider.RecycleItem, error) { @@ -337,20 +408,186 @@ func (d *Driver) CreateReference(_ context.Context, _ string, _ *url.URL) error return errtypes.NotSupported("kiteworks: read-only driver") } -func (d *Driver) CreateDir(_ context.Context, _ *provider.Reference) (*storage.CreateDirResult, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") +func (d *Driver) CreateDir(ctx context.Context, ref *provider.Reference) (*storage.CreateDirResult, error) { + parentRef := &provider.Reference{ + ResourceId: ref.GetResourceId(), + Path: path.Dir(ref.GetPath()), + } + parentID, spaceID, err := d.resolveRef(ctx, parentRef) + if err != nil { + return nil, err + } + + name := path.Base(ref.GetPath()) + if name == "" || name == "." { + return nil, errtypes.BadRequest("kiteworks: CreateDir requires a folder name") + } + + folderID, err := d.client(ctx).CreateFolder(parentID, kwlib.CreateDirRequest{Name: name}) + if err != nil { + return nil, err + } + return &storage.CreateDirResult{ + SpaceID: spaceID, + ResourceID: &provider.ResourceId{ + StorageId: d.storageID, + SpaceId: spaceID, + OpaqueId: folderID, + }, + }, nil } -func (d *Driver) TouchFile(_ context.Context, _ *provider.Reference, _ bool, _ string) (*storage.TouchFileResult, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") +func (d *Driver) TouchFile(ctx context.Context, ref *provider.Reference, _ bool, _ string) (*storage.TouchFileResult, error) { + parentRef := &provider.Reference{ + ResourceId: ref.GetResourceId(), + Path: path.Dir(ref.GetPath()), + } + parentID, spaceID, err := d.resolveRef(ctx, parentRef) + if err != nil { + return nil, err + } + + name := path.Base(ref.GetPath()) + if name == "" || name == "." { + return nil, errtypes.BadRequest("kiteworks: TouchFile requires a filename") + } + + c := d.client(ctx) + upload, err := c.InitializeUpload(parentID, name, 0, 1) + if err != nil { + return nil, err + } + fi, err := c.UploadChunk(ctx, upload.URI, name, bytes.NewReader(nil), 0, 0, true) + if err != nil { + return nil, err + } + return &storage.TouchFileResult{ + SpaceID: spaceID, + ResourceID: &provider.ResourceId{ + StorageId: d.storageID, + SpaceId: spaceID, + OpaqueId: fi.ID, + }, + }, nil } -func (d *Driver) Delete(_ context.Context, _ *provider.Reference) (*storage.DeleteResult, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") +func (d *Driver) Delete(ctx context.Context, ref *provider.Reference) (*storage.DeleteResult, error) { + nodeID, spaceID, err := d.resolveRef(ctx, ref) + if err != nil { + return nil, err + } + + c := d.client(ctx) + err = c.DeleteFolder(nodeID) + if err != nil { + var ce *kwlib.ClientError + if !errors.As(err, &ce) || ce.StatusCode != http.StatusNotFound { + return nil, err + } + if err = c.DeleteFile(nodeID); err != nil { + return nil, err + } + } + d.locks.Delete(nodeID) + return &storage.DeleteResult{ + ResourceId: &provider.ResourceId{ + StorageId: d.storageID, + SpaceId: spaceID, + OpaqueId: nodeID, + }, + }, nil } -func (d *Driver) Move(_ context.Context, _, _ *provider.Reference) (*storage.MoveResult, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") +// rawNodeInfo fetches the raw KW FileInfo for a node, trying folder first. +// Used by write methods that need ParentID or type without a full toResourceInfo conversion. +func (d *Driver) rawNodeInfo(ctx context.Context, nodeID string) (*kwlib.FileInfo, error) { + c := d.client(ctx) + fi, err := c.GetFolderByID(nodeID) + if err == nil { + return fi, nil + } + var ce *kwlib.ClientError + if !errors.As(err, &ce) || (ce.StatusCode != http.StatusNotFound && ce.StatusCode != http.StatusForbidden) { + return nil, err + } + return c.GetFileByID(nodeID) +} + +func moveNode(c *kwlib.APIClient, fi *kwlib.FileInfo, parentID string) error { + if fi.IsDir() { + return c.MoveFolder(fi.ID, parentID) + } + _, err := c.Move(fi, &kwlib.FileInfo{ID: parentID}, false) + return err +} + +func renameNode(c *kwlib.APIClient, fi *kwlib.FileInfo, name string) error { + if fi.IsDir() { + _, err := c.RenameFolder(fi, name) + return err + } + _, err := c.RenameFile(fi, name, false) + return err +} + +func (d *Driver) Move(ctx context.Context, src, dst *provider.Reference) (*storage.MoveResult, error) { + dstName := path.Base(dst.GetPath()) + if dstName == "" || dstName == "." || dstName == "/" { + return nil, errtypes.BadRequest("kiteworks: Move requires a destination name") + } + + srcNodeID, _, err := d.resolveRef(ctx, src) + if err != nil { + return nil, err + } + + dstParentID, _, err := d.resolveRef(ctx, &provider.Reference{ + ResourceId: dst.GetResourceId(), + Path: path.Dir(dst.GetPath()), + }) + if err != nil { + return nil, err + } + if dstParentID == srcNodeID { + return nil, errtypes.BadRequest("kiteworks: cannot move a resource into itself") + } + + srcFI, err := d.rawNodeInfo(ctx, srcNodeID) + if err != nil { + return nil, err + } + + srcParentID := "" + if srcFI.ParentID != nil { + srcParentID = *srcFI.ParentID + } + + // KW has no combined move+rename. Rename first: it is the undoable step, and + // the new name cannot collide in a parent the node has not left yet. + c := d.client(ctx) + renamed := dstName != srcFI.Name + if renamed { + if err := renameNode(c, srcFI, dstName); err != nil { + return nil, err + } + } + + if srcParentID != dstParentID { + if err := moveNode(c, srcFI, dstParentID); err != nil { + if renamed { + if rbErr := renameNode(c, srcFI, srcFI.Name); rbErr != nil { + d.log.Error().Err(rbErr).Str("nodeID", srcNodeID).Msg("could not restore the original name after a failed move") + return nil, fmt.Errorf("kiteworks: move failed and %q is left renamed to %q: %w", srcFI.Name, dstName, err) + } + } + return nil, err + } + } + + return &storage.MoveResult{ + OldReference: src, + NewReference: dst, + }, nil } func (d *Driver) InitiateUpload(_ context.Context, _ *provider.Reference, _ int64, _ map[string]string) (map[string]string, error) { @@ -361,26 +598,110 @@ func (d *Driver) Upload(_ context.Context, _ storage.UploadRequest, _ storage.Up return nil, errtypes.NotSupported("kiteworks: read-only driver") } +// MarkProcessing is a no-op: KW controls its own file metadata so there is no +// reliable way to mark a node as "in-flight" without abusing the checkout-lock +// API (which would block concurrent uploads and leave orphaned locks on crash). +// Trade-off: a file touched by TouchFile is downloadable as a zero-byte stub +// until CommitUpload replaces it. func (d *Driver) MarkProcessing(_ context.Context, _ *provider.Reference, _ bool, _ string) error { - return errtypes.NotSupported("kiteworks: read-only driver") + return nil } -func (d *Driver) CommitUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.UploadSource) error { - return errtypes.NotSupported("kiteworks: read-only driver") +// TODO: large uploads still fail for the client in the sync path. +// +// There are two legs, and their speeds are unrelated: +// +// client -> oCIS TUS PATCHes, staged to a local bin file +// oCIS -> KW this method, 8 MiB chunks over a KW upload session +// +// async: no problem. CommitUpload runs after the client is already done, so it +// is bound to no request. +// +// sync: problem. The whole oCIS -> KW transfer happens inside the client's last +// chunk request. That transfer can take minutes; the request is not meant to +// live that long. +// +// Ideas: +// 1. Always run the KW driver with async enabled. asyncfileuploads is not +// wired into revaconfig, so the commit currently always runs inline. +// 2. Make sync partly async: answer the client before the file is fully in KW. +// Sensible, but needs verifying against real clients. +// 3. Let the last chunk request run arbitrarily long for big files. Doubtful, +// timeouts exist for a reason. +// 4. Commit chunk by chunk instead of at the end. Far from the current +// architecture, and TUS resume does not map onto KW's sequential chunks. +// 5. Upload client -> KW directly and only register the file here. Drops +// antivirus and checksums. +func (d *Driver) CommitUpload(ctx context.Context, ref *provider.Reference, _ string, source storage.UploadSource) error { + if source.Body == nil { + return errtypes.BadRequest("kiteworks: CommitUpload requires a non-nil body") + } + // the client connection's lifetime must not abort a commit already in flight + ctx = context.WithoutCancel(ctx) + nodeID, _, err := d.resolveRef(ctx, ref) + if err != nil { + return err + } + c := d.client(ctx) + fi, err := c.GetFileByID(nodeID) + if err != nil { + return err + } + if err = c.UploadFileVersion(ctx, nodeID, fi.Name, source.Body, source.Length); err != nil { + return err + } + if !source.NodeExisted { + d.deletePlaceholderVersion(c, nodeID) + } + return nil +} + +// deletePlaceholderVersion removes the zero-byte stub version created by TouchFile, +// leaving only the real content at version 1. Best-effort: all errors are logged, +// never returned; the upload already succeeded before this is called. +func (d *Driver) deletePlaceholderVersion(c *kwlib.APIClient, fileID string) { + log := d.log.With().Str("fileID", fileID).Logger() + + versions, err := c.GetFileVersions(fileID) + if err != nil { + log.Warn().Err(err).Msg("placeholder cleanup: could not list versions") + return + } + for _, v := range versions { + if v.VersionNumber > 0 && v.Size == 0 { + if err := c.DeleteFileVersion(fileID, v.ID); err != nil { + ce := kwlib.AsClientError(err) + switch ce.StatusCode { + case http.StatusNotFound: // already gone, fine + case http.StatusUnprocessableEntity: // last version: real bytes not uploaded, leave it + log.Warn().Msg("placeholder cleanup: cannot delete last version, real upload may have failed") + default: + log.Warn().Err(err).Msg("placeholder cleanup: delete version failed") + } + } + return + } + } } func (d *Driver) PrepareUpload(_ context.Context, _ *provider.Reference, _ string, info storage.UploadInfo) (*storage.PrepareUploadResult, error) { return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil } -func (d *Driver) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error { +func (d *Driver) RollbackUpload(ctx context.Context, _ *provider.Reference, _ string, info storage.RollbackInfo) error { + // the prior version stays current, so there is nothing to undo + if info.NodeExisted { + return nil + } + if err := d.client(ctx).DeleteFile(info.NodeID); err != nil { + if kwlib.AsClientError(err).StatusCode == http.StatusNotFound { + return nil + } + return err + } return nil } -func (d *Driver) RestoreRevision(_ context.Context, _ *provider.Reference, _ string) (*storage.RestoreRevisionResult, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") -} - func (d *Driver) RestoreRecycleItem(_ context.Context, _ *provider.Reference, _, _ string, _ *provider.Reference) (*storage.RestoreRecycleItemResult, error) { return nil, errtypes.NotSupported("kiteworks: read-only driver") } @@ -417,16 +738,43 @@ func (d *Driver) UnsetArbitraryMetadata(_ context.Context, _ *provider.Reference return errtypes.NotSupported("kiteworks: read-only driver") } -func (d *Driver) SetLock(_ context.Context, _ *provider.Reference, _ *provider.Lock) (*storage.SetLockResult, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") +func (d *Driver) SetLock(ctx context.Context, ref *provider.Reference, lock *provider.Lock) (*storage.SetLockResult, error) { + nodeID, spaceID, err := d.resolveRef(ctx, ref) + if err != nil { + return nil, err + } + if err := d.client(ctx).LockFile(nodeID); err != nil { + var ce *kwlib.ClientError + if errors.As(err, &ce) && ce.StatusCode == http.StatusForbidden { + return nil, errtypes.PreconditionFailed("kiteworks: file already locked") + } + return nil, err + } + d.locks.Store(nodeID, lock) + return &storage.SetLockResult{SpaceID: spaceID}, nil } +// RefreshLock is a no-op. KW has no lock-refresh endpoint, locks carry no +// expiry, and re-calling the lock endpoint on an already-locked file returns +// 403 even when the caller owns the lock. func (d *Driver) RefreshLock(_ context.Context, _ *provider.Reference, _ *provider.Lock, _ string) error { - return errtypes.NotSupported("kiteworks: read-only driver") + return nil } -func (d *Driver) Unlock(_ context.Context, _ *provider.Reference, _ *provider.Lock) (*storage.UnlockResult, error) { - return nil, errtypes.NotSupported("kiteworks: read-only driver") +func (d *Driver) Unlock(ctx context.Context, ref *provider.Reference, lock *provider.Lock) (*storage.UnlockResult, error) { + nodeID, spaceID, err := d.resolveRef(ctx, ref) + if err != nil { + return nil, err + } + if err := d.client(ctx).UnlockFile(nodeID); err != nil { + var ce *kwlib.ClientError + if errors.As(err, &ce) && ce.StatusCode == http.StatusForbidden { + return nil, errtypes.PreconditionFailed("kiteworks: file not locked or locked by another user") + } + return nil, err + } + d.locks.Delete(nodeID) + return &storage.UnlockResult{SpaceID: spaceID}, nil } func (d *Driver) CreateStorageSpace(_ context.Context, _ *provider.CreateStorageSpaceRequest) (*provider.CreateStorageSpaceResponse, error) { diff --git a/pkg/storage/fs/kiteworks/kiteworks_test.go b/pkg/storage/fs/kiteworks/kiteworks_test.go index 61902b6e70c..4b45917e819 100644 --- a/pkg/storage/fs/kiteworks/kiteworks_test.go +++ b/pkg/storage/fs/kiteworks/kiteworks_test.go @@ -4,8 +4,10 @@ import ( "context" "errors" "io" + "math" "net/http/httptest" "os" + "strings" provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" . "github.com/onsi/ginkgo/v2" @@ -59,7 +61,7 @@ func setupDriver() (storage.FS, *fixture, func()) { ctx := ctxpkg.ContextSetToken(context.Background(), os.Getenv("KITEWORKS_TOKEN")) spaces, err := d.ListStorageSpaces(ctx, nil, false) - Expect(err).ToNot(HaveOccurred(), "real-box ListStorageSpaces failed — check token/endpoint") + Expect(err).ToNot(HaveOccurred(), "real-box ListStorageSpaces failed, check token/endpoint") Expect(spaces).ToNot(BeEmpty(), "real box has no top-level folders") spaceID := spaces[0].Root.OpaqueId @@ -181,6 +183,44 @@ var _ = Describe("kiteworks driver", func() { Expect(path).ToNot(BeEmpty()) }) }) + + Describe("GetQuota", func() { + It("reports the quota of the referenced folder", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}, + } + total, used, remaining, err := d.GetQuota(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(total).To(BeEquivalentTo(1073741824)) + Expect(used).To(BeEquivalentTo(14)) + Expect(remaining).To(BeEquivalentTo(1073741810)) + }) + + It("reports unlimited remaining when no quota is applied", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: "space-2", OpaqueId: "space-2"}, + } + total, used, remaining, err := d.GetQuota(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(total).To(BeEquivalentTo(0)) + Expect(used).To(BeEquivalentTo(99)) + Expect(remaining).To(BeEquivalentTo(uint64(math.MaxUint64))) + }) + + It("reports unlimited rather than failing when the lookup is denied", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: "no-quota-perm-1", OpaqueId: "no-quota-perm-1"}, + } + total, used, remaining, err := d.GetQuota(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(total).To(BeEquivalentTo(0)) + Expect(used).To(BeEquivalentTo(0)) + Expect(remaining).To(BeEquivalentTo(uint64(math.MaxUint64))) + }) + }) }) Context("error propagation", func() { @@ -206,30 +246,436 @@ var _ = Describe("kiteworks driver", func() { }) }) - Context("write rejection", func() { - notSupported := func(err error) bool { return errors.As(err, new(errtypes.NotSupported)) } + Context("write path", func() { + Describe("CreateDir", func() { + It("creates a folder under the space root in mock mode", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}, + Path: "./Documents", + } + result, err := d.CreateDir(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(result.SpaceID).To(Equal(fix.spaceID)) + Expect(result.ResourceID.OpaqueId).To(Equal("new-dir-1")) + Expect(result.ResourceID.StorageId).To(Equal("kiteworks")) + }) + }) - It("rejects CreateDir", func() { - _, err := d.CreateDir(fix.ctx, &provider.Reference{ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}}) - Expect(err).To(Satisfy(notSupported)) + Describe("TouchFile", func() { + It("creates a zero-byte stub and returns a TouchFileResult", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}, + Path: "./newfile.txt", + } + result, err := d.TouchFile(fix.ctx, ref, false, "") + Expect(err).ToNot(HaveOccurred()) + Expect(result.SpaceID).To(Equal(fix.spaceID)) + Expect(result.ResourceID.OpaqueId).To(Equal("touched-1")) + Expect(result.ResourceID.StorageId).To(Equal("kiteworks")) + }) }) - It("rejects TouchFile", func() { - _, err := d.TouchFile(fix.ctx, &provider.Reference{ResourceId: &provider.ResourceId{SpaceId: fix.spaceID}}, false, "") - Expect(err).To(Satisfy(notSupported)) + + Describe("Delete", func() { + It("deletes a folder by ID", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "folder-del-1"}, + } + result, err := d.Delete(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(result.ResourceId.OpaqueId).To(Equal("folder-del-1")) + }) + + It("falls back to file delete when folder delete returns 404", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "file-only-1"}, + } + result, err := d.Delete(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(result.ResourceId.OpaqueId).To(Equal("file-only-1")) + }) }) - It("rejects Delete", func() { - _, err := d.Delete(fix.ctx, &provider.Reference{ResourceId: &provider.ResourceId{SpaceId: fix.spaceID}}) - Expect(err).To(Satisfy(notSupported)) + + Describe("Move", func() { + It("renames a file in place", func() { + skipIfRealBox() + src := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "src-file-1"}, + } + dst := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}, + Path: "./renamed.txt", + } + result, err := d.Move(fix.ctx, src, dst) + Expect(err).ToNot(HaveOccurred()) + Expect(result).ToNot(BeNil()) + }) + + It("moves a file to a different folder", func() { + skipIfRealBox() + src := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "src-file-1"}, + } + dst := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "folder-2"}, + Path: "./src.txt", + } + result, err := d.Move(fix.ctx, src, dst) + Expect(err).ToNot(HaveOccurred()) + Expect(result).ToNot(BeNil()) + }) + + It("renames a folder in place", func() { + skipIfRealBox() + src := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "src-folder-1"}, + } + dst := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}, + Path: "./RenamedFolder", + } + result, err := d.Move(fix.ctx, src, dst) + Expect(err).ToNot(HaveOccurred()) + Expect(result).ToNot(BeNil()) + }) + + It("moves a folder to a different parent", func() { + skipIfRealBox() + src := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "src-folder-1"}, + } + dst := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "folder-2"}, + Path: "./SrcFolder", + } + result, err := d.Move(fix.ctx, src, dst) + Expect(err).ToNot(HaveOccurred()) + Expect(result).ToNot(BeNil()) + }) }) - It("rejects Move", func() { - ref := &provider.Reference{ResourceId: &provider.ResourceId{SpaceId: fix.spaceID}} - _, err := d.Move(fix.ctx, ref, ref) - Expect(err).To(Satisfy(notSupported)) + + Describe("PrepareUpload", func() { + It("returns VersionCreated=false and marks session for new files", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}, + } + result, err := d.PrepareUpload(fix.ctx, ref, "new-sess-1", storage.UploadInfo{NodeExisted: false}) + Expect(err).ToNot(HaveOccurred()) + Expect(result.VersionCreated).To(BeFalse()) + }) + + It("returns VersionCreated=true and does not mark session for existing files", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: fix.spaceID}, + } + result, err := d.PrepareUpload(fix.ctx, ref, "existing-sess-1", storage.UploadInfo{NodeExisted: true}) + Expect(err).ToNot(HaveOccurred()) + Expect(result.VersionCreated).To(BeTrue()) + }) }) - It("rejects SetLock", func() { - _, err := d.SetLock(fix.ctx, &provider.Reference{ResourceId: &provider.ResourceId{SpaceId: fix.spaceID}}, &provider.Lock{LockId: "x"}) - Expect(err).To(Satisfy(notSupported)) + + Describe("CommitUpload", func() { + It("uploads content for an existing file version", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "ver-file-1"}, + } + content := "file content here" + src := storage.UploadSource{ + Body: io.NopCloser(strings.NewReader(content)), + Length: int64(len(content)), + } + err := d.CommitUpload(fix.ctx, ref, "existing-sess-2", src) + Expect(err).ToNot(HaveOccurred()) + }) + + It("uploads content for a new file and cleans up the placeholder version", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "ver-file-1"}, + } + _, err := d.PrepareUpload(fix.ctx, ref, "new-sess-2", storage.UploadInfo{NodeExisted: false}) + Expect(err).ToNot(HaveOccurred()) + + content := "file content here" + src := storage.UploadSource{ + Body: io.NopCloser(strings.NewReader(content)), + Length: int64(len(content)), + } + err = d.CommitUpload(fix.ctx, ref, "new-sess-2", src) + Expect(err).ToNot(HaveOccurred()) + }) + }) + + Describe("RollbackUpload", func() { + rollback := func(nodeID string, nodeExisted bool) error { + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: nodeID}, + } + return d.RollbackUpload(fix.ctx, ref, "rollback-sess-1", storage.RollbackInfo{ + NodeExisted: nodeExisted, + NodeID: nodeID, + ParentID: fix.spaceID, + Filename: "rolled-back.txt", + }) + } + + // error-500 fails every method, so a stray delete would surface as an error + It("issues no request when the node already had content", func() { + skipIfRealBox() + Expect(rollback("error-500", true)).ToNot(HaveOccurred()) + }) + + It("deletes the stub it created for a new file", func() { + skipIfRealBox() + Expect(rollback("file-only-1", false)).ToNot(HaveOccurred()) + }) + + It("treats an already deleted node as success", func() { + skipIfRealBox() + Expect(rollback("rollback-gone-1", false)).ToNot(HaveOccurred()) + }) + + // a TUS terminate before TouchFile ran leaves the placeholder id from initiate + It("tolerates a node id kiteworks never knew", func() { + skipIfRealBox() + Expect(rollback("6b1e9c0e-3f2a-4d5b-8c7d-1a2b3c4d5e6f", false)).ToNot(HaveOccurred()) + }) + + It("propagates errors other than 404", func() { + skipIfRealBox() + Expect(rollback("error-500", false)).To(HaveOccurred()) + }) + + It("does not abort when the calling context is already cancelled", func() { + skipIfRealBox() + ctx, cancel := context.WithCancel(fix.ctx) + cancel() + err := d.RollbackUpload(ctx, &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "file-only-1"}, + }, "rollback-sess-2", storage.RollbackInfo{NodeID: "file-only-1"}) + Expect(err).ToNot(HaveOccurred()) + }) }) + }) + + Context("locking", func() { + precFailed := func(err error) bool { return errors.As(err, new(errtypes.PreconditionFailed)) } + + Describe("GetLock", func() { + It("returns nil for an unlocked file", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "lock-file-1"}, + } + lock, err := d.GetLock(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(lock).To(BeNil()) + }) + + It("returns a synthetic lock for a file locked by another session", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "ext-locked-file-1"}, + } + lock, err := d.GetLock(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(lock).ToNot(BeNil()) + Expect(lock.LockId).To(Equal("kw-ext-locked-file-1")) + }) + + It("returns the stored lock after SetLock", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "lock-file-1"}, + } + stored := &provider.Lock{LockId: "my-lock-id", Type: provider.LockType_LOCK_TYPE_EXCL} + _, err := d.SetLock(fix.ctx, ref, stored) + Expect(err).ToNot(HaveOccurred()) + + got, err := d.GetLock(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(got).ToNot(BeNil()) + Expect(got.LockId).To(Equal("my-lock-id")) + }) + }) + + Describe("SetLock", func() { + It("locks an unlocked file and returns the space ID", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "lock-file-1"}, + } + result, err := d.SetLock(fix.ctx, ref, &provider.Lock{LockId: "lock-1"}) + Expect(err).ToNot(HaveOccurred()) + Expect(result.SpaceID).To(Equal(fix.spaceID)) + }) + + It("returns PreconditionFailed when the file is already locked", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "lock-file-1"}, + } + _, err := d.SetLock(fix.ctx, ref, &provider.Lock{LockId: "lock-1"}) + Expect(err).ToNot(HaveOccurred()) + + _, err = d.SetLock(fix.ctx, ref, &provider.Lock{LockId: "lock-2"}) + Expect(err).To(Satisfy(precFailed)) + }) + }) + + Describe("Unlock", func() { + It("unlocks a locked file and returns the space ID", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "lock-file-1"}, + } + _, err := d.SetLock(fix.ctx, ref, &provider.Lock{LockId: "lock-1"}) + Expect(err).ToNot(HaveOccurred()) + + result, err := d.Unlock(fix.ctx, ref, &provider.Lock{LockId: "lock-1"}) + Expect(err).ToNot(HaveOccurred()) + Expect(result.SpaceID).To(Equal(fix.spaceID)) + }) + + It("returns PreconditionFailed when the file is not locked", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "lock-file-1"}, + } + _, err := d.Unlock(fix.ctx, ref, &provider.Lock{LockId: "lock-1"}) + Expect(err).To(Satisfy(precFailed)) + }) + }) + }) + + Context("versions", func() { + permDenied := func(err error) bool { return errors.As(err, new(errtypes.PermissionDenied)) } + + Describe("ListRevisions", func() { + It("returns versions with fileID@versionID keys", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "versioned-file-1"}, + } + revs, err := d.ListRevisions(fix.ctx, ref) + Expect(err).ToNot(HaveOccurred()) + Expect(revs).To(HaveLen(2)) + Expect(revs[0].Key).To(Equal("versioned-file-1@rev-1")) + Expect(revs[1].Key).To(Equal("versioned-file-1@rev-2")) + Expect(revs[1].Size).To(BeEquivalentTo(20)) + }) + + It("returns PermissionDenied when version_view is not granted", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "file-1"}, + } + _, err := d.ListRevisions(fix.ctx, ref) + Expect(err).To(Satisfy(permDenied)) + }) + }) + + Describe("DownloadRevision", func() { + It("streams version content", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "versioned-file-1"}, + } + ri, rc, err := d.DownloadRevision(fix.ctx, ref, "versioned-file-1@rev-2", func(_ *provider.ResourceInfo) bool { return true }) + Expect(err).ToNot(HaveOccurred()) + Expect(rc).ToNot(BeNil()) + defer rc.Close() + b, err := io.ReadAll(rc) + Expect(err).ToNot(HaveOccurred()) + Expect(string(b)).To(Equal("version 2 content")) + Expect(ri.Size).To(BeEquivalentTo(len("version 2 content"))) + }) + + It("returns nil reader when openReaderFunc returns false", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "versioned-file-1"}, + } + ri, rc, err := d.DownloadRevision(fix.ctx, ref, "versioned-file-1@rev-2", func(_ *provider.ResourceInfo) bool { return false }) + Expect(err).ToNot(HaveOccurred()) + Expect(ri).ToNot(BeNil()) + Expect(rc).To(BeNil()) + }) + + It("returns PermissionDenied when download is not granted", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "file-1"}, + } + _, _, err := d.DownloadRevision(fix.ctx, ref, "file-1@rev-1", func(_ *provider.ResourceInfo) bool { return true }) + Expect(err).To(Satisfy(permDenied)) + }) + }) + + Describe("RestoreRevision", func() { + It("promotes the version successfully", func() { + skipIfRealBox() + result, err := d.RestoreRevision(fix.ctx, nil, "versioned-file-1@rev-2") + Expect(err).ToNot(HaveOccurred()) + Expect(result).ToNot(BeNil()) + }) + + It("returns PermissionDenied when version_promote is not granted", func() { + skipIfRealBox() + _, err := d.RestoreRevision(fix.ctx, nil, "file-1@rev-1") + Expect(err).To(Satisfy(permDenied)) + }) + }) + + Describe("GetMD with version OpaqueId", func() { + It("returns ResourceInfo with the version OpaqueId", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "versioned-file-1@rev-2"}, + } + ri, err := d.GetMD(fix.ctx, ref, nil, nil) + Expect(err).ToNot(HaveOccurred()) + Expect(ri.Type).To(Equal(provider.ResourceType_RESOURCE_TYPE_FILE)) + Expect(ri.Id.OpaqueId).To(Equal("versioned-file-1@rev-2")) + }) + }) + + Describe("Download with version OpaqueId", func() { + It("streams version content via version ref", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "versioned-file-1@rev-2"}, + } + ri, rc, err := d.Download(fix.ctx, ref, func(_ *provider.ResourceInfo) bool { return true }) + Expect(err).ToNot(HaveOccurred()) + Expect(rc).ToNot(BeNil()) + defer rc.Close() + b, err := io.ReadAll(rc) + Expect(err).ToNot(HaveOccurred()) + Expect(string(b)).To(Equal("version 2 content")) + Expect(ri.Size).To(BeEquivalentTo(len("version 2 content"))) + }) + + It("returns PermissionDenied when download is not granted", func() { + skipIfRealBox() + ref := &provider.Reference{ + ResourceId: &provider.ResourceId{SpaceId: fix.spaceID, OpaqueId: "file-1@rev-1"}, + } + _, _, err := d.Download(fix.ctx, ref, func(_ *provider.ResourceInfo) bool { return true }) + Expect(err).To(Satisfy(permDenied)) + }) + }) + }) + + Context("write rejection", func() { + notSupported := func(err error) bool { return errors.As(err, new(errtypes.NotSupported)) } + It("rejects AddGrant", func() { err := d.AddGrant(fix.ctx, &provider.Reference{ResourceId: &provider.ResourceId{SpaceId: fix.spaceID}}, &provider.Grant{}) Expect(err).To(Satisfy(notSupported)) @@ -241,10 +687,17 @@ var _ = Describe("kiteworks driver", func() { }) Context("capabilities", func() { - It("declares an all-false (read-only) set", func() { + It("declares the implemented write, version and lock support", func() { cp, ok := d.(storage.CapabilityProvider) Expect(ok).To(BeTrue()) - Expect(cp.Capabilities(fix.ctx)).To(Equal(storage.Capabilities{})) + Expect(cp.Capabilities(fix.ctx)).To(Equal(storage.Capabilities{ + Upload: true, + CreateContainer: true, + Delete: true, + Move: true, + Versioning: true, + Locking: true, + })) }) }) }) diff --git a/pkg/storage/fs/kiteworks/kwlib/client.go b/pkg/storage/fs/kiteworks/kwlib/client.go index 68811e2c215..2c7e5a95e04 100644 --- a/pkg/storage/fs/kiteworks/kwlib/client.go +++ b/pkg/storage/fs/kiteworks/kwlib/client.go @@ -2,6 +2,7 @@ package kwlib import ( "bytes" + "context" "crypto/tls" "encoding/json" "errors" @@ -18,6 +19,8 @@ import ( "github.com/rs/zerolog" ) +const uploadChunkSize = 8 << 20 // 8 MiB + func NewClientFactory(server, agentString string, insecure bool) *APIClientFactory { transport := &http.Transport{ Proxy: http.ProxyFromEnvironment, @@ -36,28 +39,35 @@ func NewClientFactory(server, agentString string, insecure bool) *APIClientFacto // #nosec transport.TLSClientConfig = &tls.Config{InsecureSkipVerify: true} } + + uploadTransport := transport.Clone() + uploadTransport.ResponseHeaderTimeout = 30 * time.Second + return &APIClientFactory{ - server: server, - agentString: agentString, - httpClient: &http.Client{Transport: transport, Timeout: 15 * time.Second}, + server: server, + agentString: agentString, + httpClient: &http.Client{Transport: transport, Timeout: 15 * time.Second}, + uploadClient: &http.Client{Transport: uploadTransport}, } } type APIClientFactory struct { - server string - agentString string - httpClient *http.Client + server string + agentString string + httpClient *http.Client + uploadClient *http.Client } type APIClient struct { - server string - agentString string - logger *zerolog.Logger - host string - token string - requestId string - remoteAddr string - httpClient *http.Client + server string + agentString string + logger *zerolog.Logger + host string + token string + requestId string + remoteAddr string + httpClient *http.Client + uploadClient *http.Client } func decodeJSON(body io.ReadCloser, out any) error { @@ -67,13 +77,14 @@ func decodeJSON(body io.ReadCloser, out any) error { func (f *APIClientFactory) Build(host, requestId, remoteAddr, token string, l *zerolog.Logger) *APIClient { return &APIClient{ - token: token, - server: f.server, - host: host, - logger: l, - requestId: requestId, - remoteAddr: remoteAddr, - httpClient: f.httpClient, + token: token, + server: f.server, + host: host, + logger: l, + requestId: requestId, + remoteAddr: remoteAddr, + httpClient: f.httpClient, + uploadClient: f.uploadClient, } } @@ -94,7 +105,7 @@ func (c *APIClient) GetTopFolders() (*DirectoryInfo, error) { } func (c *APIClient) GetFolderByID(id string) (*FileInfo, error) { - request, err := c.NewGetRequest(fmt.Sprintf("/rest/folders/%s", id)) + request, err := c.NewGetRequest(fmt.Sprintf("/rest/folders/%s?with=(permissions)", id)) if err != nil { return nil, err } @@ -148,7 +159,7 @@ func (c *APIClient) Search(path string) (*FileInfo, error) { } func (c *APIClient) GetFileByID(id string) (*FileInfo, error) { - request, err := c.NewGetRequest(fmt.Sprintf("/rest/files/%s", id)) + request, err := c.NewGetRequest(fmt.Sprintf("/rest/files/%s?with=(permissions,lockUser)", id)) if err != nil { return nil, err } @@ -194,8 +205,8 @@ func (c *APIClient) GetUser(id string) (*User, error) { return out, nil } -func (c *APIClient) GetQuotaInfo() (*QuotaInfo, error) { - request, err := c.NewGetRequest("/rest/quotas") +func (c *APIClient) GetFolderQuota(folderID string) (*FolderQuota, error) { + request, err := c.NewGetRequest(fmt.Sprintf("/rest/folders/%s/quota", folderID)) if err != nil { return nil, err } @@ -203,7 +214,7 @@ func (c *APIClient) GetQuotaInfo() (*QuotaInfo, error) { if err != nil { return nil, err } - out := &QuotaInfo{} + out := &FolderQuota{} if err := decodeJSON(response.Body, out); err != nil { return nil, err } @@ -251,12 +262,20 @@ func (c *APIClient) CreateFolder(id string, payload CreateDirRequest) (string, e } func (c *APIClient) InitializeUpload(parentID, name string, size int64, numberOfChunks int) (*UploadResult, error) { + return c.initializeUpload(fmt.Sprintf("/rest/folders/%s/actions/initiateUpload", parentID), name, size, numberOfChunks) +} + +func (c *APIClient) InitializeVersionUpload(fileID, name string, size int64, numberOfChunks int) (*UploadResult, error) { + return c.initializeUpload(fmt.Sprintf("/rest/files/%s/actions/initiateUpload", fileID), name, size, numberOfChunks) +} + +func (c *APIClient) initializeUpload(path, name string, size int64, numberOfChunks int) (*UploadResult, error) { payload := InitializeUpload{ FileName: name, TotalSize: size, TotalChunks: numberOfChunks, } - request, err := c.NewPostRequest(fmt.Sprintf("/rest/folders/%s/actions/initiateUpload", parentID), payload) + request, err := c.NewPostRequest(path, payload) if err != nil { return nil, err } @@ -271,7 +290,20 @@ func (c *APIClient) InitializeUpload(parentID, name string, size int64, numberOf return out, nil } -func (c *APIClient) UploadChunk(uploadURI, name string, file io.Reader, chunkIndex int, chunk int64, isLastChunk bool) (*FileInfo, error) { +func (c *APIClient) TerminateUpload(uploadID int64) error { + request, err := c.newRequest("DELETE", fmt.Sprintf("/rest/uploads/%d", uploadID), nil) + if err != nil { + return err + } + response, err := c.SendRequest(request) + if err != nil { + return err + } + response.Body.Close() + return nil +} + +func (c *APIClient) UploadChunk(ctx context.Context, uploadURI, name string, file io.Reader, chunkIndex int, chunk int64, isLastChunk bool) (*FileInfo, error) { body := new(bytes.Buffer) writer := multipart.NewWriter(body) part, err := writer.CreateFormFile("content", name) @@ -296,13 +328,14 @@ func (c *APIClient) UploadChunk(uploadURI, name string, file io.Reader, chunkInd if err != nil { return nil, err } + request = request.WithContext(ctx) request.Header.Set("Content-Type", writer.FormDataContentType()) if isLastChunk { q := request.URL.Query() q.Add("returnEntity", "true") request.URL.RawQuery = q.Encode() } - response, err := c.SendRequest(request) + response, err := c.sendWith(c.uploadClient, request) if err != nil { return nil, err } @@ -313,10 +346,121 @@ func (c *APIClient) UploadChunk(uploadURI, name string, file io.Reader, chunkInd } return out, nil } + // drained so the connection can be reused across chunks + _, _ = io.Copy(io.Discard, response.Body) response.Body.Close() return nil, nil } +func (c *APIClient) MoveFolder(id, destinationFolderID string) error { + request, err := c.NewPostRequest( + fmt.Sprintf("/rest/folders/%s/actions/move", id), + MoveFolderRequest{DestinationFolderID: destinationFolderID}, + ) + if err != nil { + return err + } + _, err = c.SendRequest(request) + return err +} + +// name must be the file's current name: the server requires it even for a version upload. +func (c *APIClient) UploadFileVersion(ctx context.Context, fileID, name string, body io.Reader, length int64) error { + chunks := chunkCount(length) + session, err := c.InitializeVersionUpload(fileID, name, length, chunks) + if err != nil { + return err + } + for i := 0; i < chunks; i++ { + size := int64(uploadChunkSize) + if remaining := length - int64(i)*uploadChunkSize; remaining < size { + size = remaining + } + if _, err := c.UploadChunk(ctx, session.URI, name, body, i, size, i == chunks-1); err != nil { + if termErr := c.TerminateUpload(session.ID); termErr != nil { + c.logger.Warn().Err(termErr).Int64("uploadID", session.ID).Msg("could not terminate kiteworks upload session") + } + return err + } + } + return nil +} + +func chunkCount(length int64) int { + if length <= 0 { + return 1 + } + return int((length + uploadChunkSize - 1) / uploadChunkSize) +} + +func (c *APIClient) GetFileVersions(fileID string) ([]Version, error) { + req, err := c.NewGetRequest(fmt.Sprintf("/rest/files/%s/versions", fileID)) + if err != nil { + return nil, err + } + resp, err := c.SendRequest(req) + if err != nil { + return nil, err + } + out := &VersionList{} + if err := decodeJSON(resp.Body, out); err != nil { + return nil, err + } + return out.Data, nil +} + +func (c *APIClient) DeleteFileVersion(fileID, versionID string) error { + req, err := c.newRequest("DELETE", fmt.Sprintf("/rest/files/%s/versions/%s", fileID, versionID), nil) + if err != nil { + return err + } + _, err = c.SendRequest(req) + return err +} + +func (c *APIClient) PromoteFileVersion(fileID, versionID string) error { + req, err := c.newRequest("POST", fmt.Sprintf("/rest/files/%s/versions/%s/actions/promote", fileID, versionID), nil) + if err != nil { + return err + } + _, err = c.SendRequest(req) + return err +} + +func (c *APIClient) GetVersionContents(fileID, versionID string) (*http.Response, error) { + req, err := c.NewGetRequest(fmt.Sprintf("/rest/files/%s/versions/%s/content", fileID, versionID)) + if err != nil { + return nil, err + } + return c.SendRequest(req) +} + +func (c *APIClient) LockFile(fileID string) error { + req, err := c.newRequest("PATCH", fmt.Sprintf("/rest/files/%s/actions/lock", fileID), nil) + if err != nil { + return err + } + resp, err := c.SendRequest(req) + if err != nil { + return err + } + resp.Body.Close() + return nil +} + +func (c *APIClient) UnlockFile(fileID string) error { + req, err := c.newRequest("PATCH", fmt.Sprintf("/rest/files/%s/actions/unlock", fileID), nil) + if err != nil { + return err + } + resp, err := c.SendRequest(req) + if err != nil { + return err + } + resp.Body.Close() + return nil +} + func (c *APIClient) DeleteFolder(id string) error { request, err := c.newRequest("DELETE", fmt.Sprintf("/rest/folders/%s", id), nil) if err != nil { @@ -432,7 +576,11 @@ func (c *APIClient) newRequest(method, path string, body io.Reader) (*http.Reque } func (c *APIClient) SendRequest(req *http.Request) (*http.Response, error) { - response, err := c.httpClient.Do(req) + return c.sendWith(c.httpClient, req) +} + +func (c *APIClient) sendWith(client *http.Client, req *http.Request) (*http.Response, error) { + response, err := client.Do(req) if err != nil { c.logger.Debug().Str("method", req.Method).Str("path", req.URL.String()).Err(err).Msg("kiteworks API call errored") return nil, err diff --git a/pkg/storage/fs/kiteworks/kwlib/kwlib_test.go b/pkg/storage/fs/kiteworks/kwlib/kwlib_test.go index e0b9d07980e..2f607900fd2 100644 --- a/pkg/storage/fs/kiteworks/kwlib/kwlib_test.go +++ b/pkg/storage/fs/kiteworks/kwlib/kwlib_test.go @@ -1,10 +1,13 @@ package kwlib_test import ( + "context" "encoding/json" "errors" + "io" "net/http" "net/http/httptest" + "strings" "testing" "time" @@ -113,3 +116,367 @@ func TestSendRequest_500_isError(t *testing.T) { t.Fatal("expected error for 500") } } + +// --- Write methods --- + +func TestDeleteFolder_success(t *testing.T) { + srv := serverWith(http.StatusNoContent, "") + defer srv.Close() + c := newClient(t, srv) + if err := c.DeleteFolder("folder-1"); err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestDeleteFolder_error(t *testing.T) { + srv := serverWith(http.StatusForbidden, `{"error":"forbidden"}`) + defer srv.Close() + c := newClient(t, srv) + err := c.DeleteFolder("folder-1") + if err == nil { + t.Fatal("expected error, got nil") + } + var ce *kwlib.ClientError + if !errors.As(err, &ce) || ce.StatusCode != http.StatusForbidden { + t.Fatalf("expected ClientError(403), got %v", err) + } +} + +func TestDeleteFile_success(t *testing.T) { + srv := serverWith(http.StatusNoContent, "") + defer srv.Close() + c := newClient(t, srv) + if err := c.DeleteFile("file-1"); err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestMoveFolder_success(t *testing.T) { + srv := serverWith(http.StatusOK, "") + defer srv.Close() + c := newClient(t, srv) + if err := c.MoveFolder("src-1", "dst-1"); err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestCreateFolder_parsesLocationHeader(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("X-Accellion-Location", "/rest/folders/new-folder-123") + w.WriteHeader(http.StatusCreated) + })) + defer srv.Close() + c := newClient(t, srv) + id, err := c.CreateFolder("parent-1", kwlib.CreateDirRequest{Name: "NewFolder"}) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if id != "new-folder-123" { + t.Fatalf("want new-folder-123, got %q", id) + } +} + +// mirrors the unexported uploadChunkSize in the kwlib package +const testChunkSize = 8 << 20 + +type zeroReader struct{} + +func (zeroReader) Read(p []byte) (int, error) { return len(p), nil } + +type uploadRecorder struct { + totalChunks int + filename string + chunkIndexes []string + chunkSizes []int64 + terminated []string +} + +// versionUploadServer serves an initiateUpload + chunk session for file-1, +// failing every chunk POST with chunkStatus when it is >= 400. +func versionUploadServer(t *testing.T, rec *uploadRecorder, chunkStatus int) *httptest.Server { + t.Helper() + mux := http.NewServeMux() + + mux.HandleFunc("/rest/files/file-1/actions/initiateUpload", func(w http.ResponseWriter, r *http.Request) { + var payload struct { + FileName string `json:"filename"` + TotalChunks int `json:"totalChunks"` + } + if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + t.Errorf("decode initiate payload: %v", err) + } + rec.filename = payload.FileName + rec.totalChunks = payload.TotalChunks + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusCreated) + _, _ = w.Write([]byte(`{"id":42,"uri":"uploads/sess-42"}`)) + }) + + mux.HandleFunc("/uploads/sess-42", func(w http.ResponseWriter, r *http.Request) { + mr, err := r.MultipartReader() + if err != nil { + t.Errorf("multipart reader: %v", err) + return + } + var index string + var size int64 + for { + part, err := mr.NextPart() + if err != nil { + break + } + switch part.FormName() { + case "content": + size, _ = io.Copy(io.Discard, part) + case "index": + b, _ := io.ReadAll(part) + index = string(b) + default: + _, _ = io.Copy(io.Discard, part) + } + } + rec.chunkIndexes = append(rec.chunkIndexes, index) + rec.chunkSizes = append(rec.chunkSizes, size) + + if chunkStatus >= 400 { + w.WriteHeader(chunkStatus) + return + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusCreated) + _, _ = w.Write([]byte(`{"id":"file-1","type":"f","name":"doc.txt"}`)) + }) + + mux.HandleFunc("/rest/uploads/42", func(w http.ResponseWriter, r *http.Request) { + rec.terminated = append(rec.terminated, r.Method) + w.WriteHeader(http.StatusNoContent) + }) + + return httptest.NewServer(mux) +} + +func TestUploadFileVersion_singleChunk(t *testing.T) { + rec := &uploadRecorder{} + srv := versionUploadServer(t, rec, 0) + defer srv.Close() + c := newClient(t, srv) + + if err := c.UploadFileVersion(context.Background(), "file-1", "doc.txt", strings.NewReader("content"), 7); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if rec.filename != "doc.txt" { + t.Fatalf("want filename doc.txt, got %q", rec.filename) + } + if rec.totalChunks != 1 { + t.Fatalf("want totalChunks 1, got %d", rec.totalChunks) + } + if len(rec.chunkSizes) != 1 || rec.chunkSizes[0] != 7 { + t.Fatalf("want one 7-byte chunk, got %v", rec.chunkSizes) + } + if rec.chunkIndexes[0] != "1" { + t.Fatalf("chunk index is 1-based, got %q", rec.chunkIndexes[0]) + } + if len(rec.terminated) != 0 { + t.Fatalf("session should not be terminated on success, got %v", rec.terminated) + } +} + +func TestUploadFileVersion_zeroLengthSendsOneChunk(t *testing.T) { + rec := &uploadRecorder{} + srv := versionUploadServer(t, rec, 0) + defer srv.Close() + c := newClient(t, srv) + + if err := c.UploadFileVersion(context.Background(), "file-1", "empty.txt", strings.NewReader(""), 0); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if rec.totalChunks != 1 { + t.Fatalf("want totalChunks 1 for an empty file, got %d", rec.totalChunks) + } + if len(rec.chunkSizes) != 1 || rec.chunkSizes[0] != 0 { + t.Fatalf("want one 0-byte chunk, got %v", rec.chunkSizes) + } +} + +func TestUploadFileVersion_splitsIntoChunks(t *testing.T) { + rec := &uploadRecorder{} + srv := versionUploadServer(t, rec, 0) + defer srv.Close() + c := newClient(t, srv) + + length := int64(testChunkSize + 100) + body := io.LimitReader(zeroReader{}, length) + if err := c.UploadFileVersion(context.Background(), "file-1", "big.bin", body, length); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if rec.totalChunks != 2 { + t.Fatalf("want totalChunks 2, got %d", rec.totalChunks) + } + want := []int64{testChunkSize, 100} + if len(rec.chunkSizes) != 2 || rec.chunkSizes[0] != want[0] || rec.chunkSizes[1] != want[1] { + t.Fatalf("want chunk sizes %v, got %v", want, rec.chunkSizes) + } + if rec.chunkIndexes[0] != "1" || rec.chunkIndexes[1] != "2" { + t.Fatalf("want indexes [1 2], got %v", rec.chunkIndexes) + } +} + +func TestUploadFileVersion_chunkErrorTerminatesSession(t *testing.T) { + rec := &uploadRecorder{} + srv := versionUploadServer(t, rec, http.StatusForbidden) + defer srv.Close() + c := newClient(t, srv) + + err := c.UploadFileVersion(context.Background(), "file-1", "doc.txt", strings.NewReader("x"), 1) + if err == nil { + t.Fatal("expected error, got nil") + } + var ce *kwlib.ClientError + if !errors.As(err, &ce) || ce.StatusCode != http.StatusForbidden { + t.Fatalf("expected ClientError(403), got %v", err) + } + if len(rec.terminated) != 1 || rec.terminated[0] != http.MethodDelete { + t.Fatalf("want one DELETE to discard the session, got %v", rec.terminated) + } +} + +func TestUploadFileVersion_initiateError(t *testing.T) { + srv := serverWith(http.StatusForbidden, `{"error":"forbidden"}`) + defer srv.Close() + c := newClient(t, srv) + + err := c.UploadFileVersion(context.Background(), "file-1", "doc.txt", strings.NewReader("x"), 1) + var ce *kwlib.ClientError + if !errors.As(err, &ce) || ce.StatusCode != http.StatusForbidden { + t.Fatalf("expected ClientError(403), got %v", err) + } +} + +func TestGetFileVersions_success(t *testing.T) { + body := `{"data":[{"id":"v1","versionNumber":1,"size":0},{"id":"v2","versionNumber":2,"size":100}]}` + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(body)) + })) + defer srv.Close() + c := newClient(t, srv) + versions, err := c.GetFileVersions("file-1") + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if len(versions) != 2 { + t.Fatalf("want 2 versions, got %d", len(versions)) + } + if versions[0].ID != "v1" || versions[0].VersionNumber != 1 || versions[0].Size != 0 { + t.Fatalf("unexpected first version: %+v", versions[0]) + } + if versions[1].ID != "v2" || versions[1].Size != 100 { + t.Fatalf("unexpected second version: %+v", versions[1]) + } +} + +func TestGetFileVersions_decodeError(t *testing.T) { + srv := badJSONServer() + defer srv.Close() + c := newClient(t, srv) + if _, err := c.GetFileVersions("file-1"); err == nil { + t.Fatal("expected error, got nil") + } +} + +func TestDeleteFileVersion_success(t *testing.T) { + srv := serverWith(http.StatusNoContent, "") + defer srv.Close() + c := newClient(t, srv) + if err := c.DeleteFileVersion("file-1", "v1"); err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestDeleteFileVersion_422(t *testing.T) { + srv := serverWith(http.StatusUnprocessableEntity, `{"error":"last version"}`) + defer srv.Close() + c := newClient(t, srv) + err := c.DeleteFileVersion("file-1", "v1") + if err == nil { + t.Fatal("expected error, got nil") + } + var ce *kwlib.ClientError + if !errors.As(err, &ce) || ce.StatusCode != http.StatusUnprocessableEntity { + t.Fatalf("expected ClientError(422), got %v", err) + } +} + +func TestRenameFolder_success(t *testing.T) { + srv := serverWith(http.StatusOK, "") + defer srv.Close() + c := newClient(t, srv) + parentID := "p1" + fi := &kwlib.FileInfo{ID: "f1", Type: kwlib.DirectoryType, Name: "OldName", ParentID: &parentID} + ok, err := c.RenameFolder(fi, "NewName") + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !ok { + t.Fatal("expected ok=true") + } +} + +func TestRenameFile_success(t *testing.T) { + srv := serverWith(http.StatusOK, "") + defer srv.Close() + c := newClient(t, srv) + fi := &kwlib.FileInfo{ID: "f1", Type: kwlib.FileType, Name: "old.txt"} + ok, err := c.RenameFile(fi, "new.txt", false) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !ok { + t.Fatal("expected ok=true") + } +} + +func TestGetFolderQuota_success(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/rest/folders/folder-1/quota" { + t.Errorf("unexpected path: %s", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"storage_quota":1000,"storage_used":250,"storage_available":750}`)) + })) + defer srv.Close() + c := newClient(t, srv) + q, err := c.GetFolderQuota("folder-1") + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if q.StorageQuota != 1000 || q.StorageUsed != 250 || q.StorageAvailable != 750 { + t.Fatalf("unexpected quota: %+v", q) + } +} + +func TestGetFolderQuota_forbidden(t *testing.T) { + srv := serverWith(http.StatusForbidden, `{"error":"forbidden"}`) + defer srv.Close() + c := newClient(t, srv) + _, err := c.GetFolderQuota("folder-1") + var ce *kwlib.ClientError + if !errors.As(err, &ce) || ce.StatusCode != http.StatusForbidden { + t.Fatalf("expected ClientError(403), got %v", err) + } +} + +func TestMoveFile_success(t *testing.T) { + srv := serverWith(http.StatusOK, "") + defer srv.Close() + c := newClient(t, srv) + src := &kwlib.FileInfo{ID: "src-1", Type: kwlib.FileType, Name: "file.txt"} + dst := &kwlib.FileInfo{ID: "dst-folder-1"} + ok, err := c.Move(src, dst, false) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !ok { + t.Fatal("expected ok=true") + } +} diff --git a/pkg/storage/fs/kiteworks/kwlib/types.go b/pkg/storage/fs/kiteworks/kwlib/types.go index 58e2d58bc84..35146eedd8e 100644 --- a/pkg/storage/fs/kiteworks/kwlib/types.go +++ b/pkg/storage/fs/kiteworks/kwlib/types.go @@ -15,7 +15,24 @@ const ( ) const ( - DownloadPermission = "download" + PermDownload = "download" + PermView = "view" + PermPropertiesView = "properties_view" + PermUserView = "user_view" + PermUserAdd = "user_add" + PermUserEdit = "user_edit" + PermUserRemove = "user_remove" + PermRename = "rename" + PermFolderAdd = "folder_add" + PermFileAdd = "file_add" + PermFolderDelete = "folder_delete" + PermFolderMove = "folder_move" + PermVersionView = "version_view" + PermVersionCreate = "version_create" + PermVersionPromote = "version_promote" + PermVersionDelete = "version_delete" + PermFileDelete = "file_delete" + PermFileMove = "file_move" ) type FileSearch struct { @@ -69,6 +86,8 @@ type FileInfo struct { PermaLink string `json:"permalink"` Permissions []Permission `json:"permissions"` Creator User `json:"creator"` + Locked bool `json:"locked"` + LockUser *User `json:"lockUser"` } type FileFingerPrints []FileFingerPrint @@ -143,13 +162,15 @@ func (fi *FileInfo) IsSyncAble() bool { return *fi.SyncAble } -func (fi *FileInfo) HasDownloadPermission() bool { +// HasPermission returns true if the named permission is present and allowed. +// Absent permissions are treated as denied. +func (fi *FileInfo) HasPermission(name string) bool { if fi == nil { return false } for i := range fi.Permissions { - if fi.Permissions[i].Name == DownloadPermission { - return true + if fi.Permissions[i].Name == name { + return fi.Permissions[i].Allowed } } return false @@ -188,9 +209,12 @@ type UploadResult struct { TotalSize int64 `json:"totalSize"` } -type QuotaInfo struct { - FolderQuotaAllowed int64 `json:"folder_quota_allowed"` - FolderQuotaUsed int64 `json:"folder_quota_used"` +// FolderQuota is the response of GET /rest/folders/{id}/quota. +// StorageQuota is 0 when no quota is applied to the folder. +type FolderQuota struct { + StorageQuota int64 `json:"storage_quota"` + StorageUsed int64 `json:"storage_used"` + StorageAvailable int64 `json:"storage_available"` } type User struct { @@ -224,3 +248,19 @@ type FileUpdateRequest struct { type FolderUpdatePutRequest struct { Name string `json:"name,omitempty"` } + +type MoveFolderRequest struct { + DestinationFolderID string `json:"destinationFolderId"` +} + +// Version represents a single version entry from GET /rest/files/{id}/versions. +type Version struct { + ID string `json:"id"` + VersionNumber int `json:"versionNumber"` + Size int64 `json:"size"` + Created Time `json:"created"` +} + +type VersionList struct { + Data []Version `json:"data"` +} diff --git a/pkg/storage/fs/kiteworks/mock_server_test.go b/pkg/storage/fs/kiteworks/mock_server_test.go index 0d846606ef5..d1c052444ea 100644 --- a/pkg/storage/fs/kiteworks/mock_server_test.go +++ b/pkg/storage/fs/kiteworks/mock_server_test.go @@ -1,8 +1,10 @@ package kiteworks_test import ( + "fmt" "net/http" "strings" + "sync/atomic" ) func writeJSON(w http.ResponseWriter, body string) { @@ -29,12 +31,147 @@ func mockKiteworksHandler() http.Handler { w.Header().Set("Content-Type", "text/plain") _, _ = w.Write([]byte("hello kiteworks")) }) + mux.HandleFunc("/rest/folders/space-1/actions/initiateUpload", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusCreated) + writeJSON(w, `{"id":1,"uri":"uploads/touch-mock-1","totalSize":0}`) + }) + mux.HandleFunc("/uploads/touch-mock-1", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusCreated) + writeJSON(w, `{"id":"touched-1","type":"f","name":"newfile.txt","path":"/My Docs/newfile.txt","size":0,"modified":"2024-01-01T00:00:00+0000"}`) + }) + mux.HandleFunc("/rest/folders/space-1/folders", func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("X-Accellion-Location", "/rest/folders/new-dir-1") + w.WriteHeader(http.StatusCreated) + }) mux.HandleFunc("/rest/users/me", func(w http.ResponseWriter, r *http.Request) { writeJSON(w, `{"id":"user-1","name":"Test User","email":"test@example.com"}`) }) - mux.HandleFunc("/rest/quotas", func(w http.ResponseWriter, r *http.Request) { - writeJSON(w, `{"folder_quota_allowed":1073741824,"folder_quota_used":14}`) + mux.HandleFunc("/rest/folders/space-1/quota", func(w http.ResponseWriter, r *http.Request) { + writeJSON(w, `{"storage_quota":1073741824,"storage_used":14,"storage_available":1073741810}`) + }) + // no quota applied to this folder + mux.HandleFunc("/rest/folders/space-2/quota", func(w http.ResponseWriter, r *http.Request) { + writeJSON(w, `{"storage_quota":0,"storage_used":99,"storage_available":0}`) + }) + // viewers lack file_add, so KW rejects the quota lookup + mux.HandleFunc("/rest/folders/no-quota-perm-1/quota", func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusForbidden) + }) + + // Move test nodes + mux.HandleFunc("/rest/folders/src-folder-1", func(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodDelete: + w.WriteHeader(http.StatusNoContent) + case http.MethodPut: + w.WriteHeader(http.StatusOK) + default: + writeJSON(w, `{"id":"src-folder-1","type":"d","name":"SrcFolder","path":"/My Docs/SrcFolder","parentId":"space-1","modified":"2024-01-01T00:00:00+0000"}`) + } + }) + mux.HandleFunc("/rest/folders/src-folder-1/actions/move", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + }) + mux.HandleFunc("/rest/files/src-file-1", func(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodPut: + w.WriteHeader(http.StatusOK) + default: + writeJSON(w, `{"id":"src-file-1","type":"f","name":"src.txt","path":"/My Docs/src.txt","parentId":"space-1","size":10,"modified":"2024-01-01T00:00:00+0000"}`) + } + }) + mux.HandleFunc("/rest/files/actions/move", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + }) + + // Delete test nodes + mux.HandleFunc("/rest/folders/folder-del-1", func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodDelete { + w.WriteHeader(http.StatusNoContent) + } else { + writeJSON(w, `{"id":"folder-del-1","type":"d","name":"ToDelete","path":"/My Docs/ToDelete","parentId":"space-1","modified":"2024-01-01T00:00:00+0000"}`) + } + }) + mux.HandleFunc("/rest/folders/file-only-1", func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNotFound) + }) + mux.HandleFunc("/rest/files/file-only-1", func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodDelete { + w.WriteHeader(http.StatusNoContent) + } else { + writeJSON(w, `{"id":"file-only-1","type":"f","name":"fileonly.txt","path":"/My Docs/fileonly.txt","size":5,"modified":"2024-01-01T00:00:00+0000"}`) + } + }) + mux.HandleFunc("/rest/files/rollback-gone-1", func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNotFound) + }) + + // CommitUpload / version endpoints + mux.HandleFunc("/rest/files/ver-file-1", func(w http.ResponseWriter, r *http.Request) { + writeJSON(w, `{"id":"ver-file-1","type":"f","name":"versionable.txt","path":"/My Docs/versionable.txt","size":17,"modified":"2024-01-01T00:00:00+0000"}`) + }) + mux.HandleFunc("/rest/files/ver-file-1/actions/initiateUpload", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusCreated) + writeJSON(w, `{"id":7,"uri":"uploads/version-mock-1"}`) }) + mux.HandleFunc("/uploads/version-mock-1", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusCreated) + writeJSON(w, `{"id":"ver-file-1","type":"f","name":"versionable.txt","path":"/My Docs/versionable.txt","size":17,"modified":"2024-01-01T00:00:00+0000"}`) + }) + mux.HandleFunc("/rest/uploads/7", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusNoContent) + }) + mux.HandleFunc("/rest/files/ver-file-1/versions", func(w http.ResponseWriter, r *http.Request) { + writeJSON(w, `{"data":[{"id":"ver-1","versionNumber":1,"size":0}]}`) + }) + mux.HandleFunc("/rest/files/ver-file-1/versions/ver-1", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusNoContent) + }) + + var lockState int32 // 0 = unlocked, 1 = locked; accessed via atomic ops + mux.HandleFunc("/rest/files/lock-file-1", func(w http.ResponseWriter, r *http.Request) { + lockedVal := "false" + if atomic.LoadInt32(&lockState) == 1 { + lockedVal = "true" + } + writeJSON(w, `{"id":"lock-file-1","type":"f","name":"lockable.txt","path":"/My Docs/lockable.txt","size":10,"modified":"2024-01-01T00:00:00+0000","locked":`+lockedVal+`}`) + }) + mux.HandleFunc("/rest/files/lock-file-1/actions/lock", func(w http.ResponseWriter, r *http.Request) { + if !atomic.CompareAndSwapInt32(&lockState, 0, 1) { + w.WriteHeader(http.StatusForbidden) + return + } + w.WriteHeader(http.StatusOK) + }) + mux.HandleFunc("/rest/files/lock-file-1/actions/unlock", func(w http.ResponseWriter, r *http.Request) { + if !atomic.CompareAndSwapInt32(&lockState, 1, 0) { + w.WriteHeader(http.StatusForbidden) + return + } + w.WriteHeader(http.StatusOK) + }) + mux.HandleFunc("/rest/files/ext-locked-file-1", func(w http.ResponseWriter, r *http.Request) { + writeJSON(w, `{"id":"ext-locked-file-1","type":"f","name":"ext-locked.txt","path":"/My Docs/ext-locked.txt","size":5,"modified":"2024-01-01T00:00:00+0000","locked":true}`) + }) + + // Versioning test nodes + versionedFileJSON := `{"id":"versioned-file-1","type":"f","name":"versioned.txt","path":"/My Docs/versioned.txt","size":20,"modified":"2024-01-01T00:00:00+0000","permissions":[{"id":1,"name":"version_view","allowed":true},{"id":2,"name":"version_promote","allowed":true},{"id":3,"name":"download","allowed":true}]}` + mux.HandleFunc("/rest/files/versioned-file-1", func(w http.ResponseWriter, r *http.Request) { + writeJSON(w, versionedFileJSON) + }) + mux.HandleFunc("/rest/files/versioned-file-1/versions", func(w http.ResponseWriter, r *http.Request) { + writeJSON(w, `{"data":[{"id":"rev-1","versionNumber":1,"size":10,"created":"2024-01-01T00:00:00+0000"},{"id":"rev-2","versionNumber":2,"size":20,"created":"2024-02-01T00:00:00+0000"}]}`) + }) + mux.HandleFunc("/rest/files/versioned-file-1/versions/rev-2/content", func(w http.ResponseWriter, r *http.Request) { + body := "version 2 content" + w.Header().Set("Content-Type", "text/plain") + w.Header().Set("Content-Length", fmt.Sprintf("%d", len(body))) + _, _ = w.Write([]byte(body)) + }) + mux.HandleFunc("/rest/files/versioned-file-1/versions/rev-2/actions/promote", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + }) + serverError := func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusInternalServerError) _, _ = w.Write([]byte(`{"error":"server error"}`)) @@ -49,4 +186,3 @@ func mockKiteworksHandler() http.Handler { return mux } - diff --git a/pkg/storage/fs/kiteworks/permissions.go b/pkg/storage/fs/kiteworks/permissions.go new file mode 100644 index 00000000000..44943c611ce --- /dev/null +++ b/pkg/storage/fs/kiteworks/permissions.go @@ -0,0 +1,62 @@ +package kiteworks + +import ( + provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" + + "github.com/owncloud/reva/v2/pkg/conversions" + "github.com/owncloud/reva/v2/pkg/storage/fs/kiteworks/kwlib" +) + +// permissionSet returns the CS3 ResourcePermissions for fi based on its KW permissions. +func permissionSet(fi *kwlib.FileInfo) *provider.ResourcePermissions { + if fi.IsDir() { + return folderPermissions(fi) + } + return filePermissions(fi) +} + +func folderPermissions(fi *kwlib.FileInfo) *provider.ResourcePermissions { + canView := fi.HasPermission(kwlib.PermPropertiesView) + return &provider.ResourcePermissions{ + Stat: canView, + GetPath: canView, + ListContainer: canView, + InitiateFileDownload: fi.HasPermission(kwlib.PermDownload), + InitiateFileUpload: fi.HasPermission(kwlib.PermFileAdd), + CreateContainer: fi.HasPermission(kwlib.PermFolderAdd), + Delete: fi.HasPermission(kwlib.PermFolderDelete), + Move: fi.HasPermission(kwlib.PermFolderMove), + ListGrants: fi.HasPermission(kwlib.PermUserView), + AddGrant: fi.HasPermission(kwlib.PermUserAdd), + UpdateGrant: fi.HasPermission(kwlib.PermUserEdit), + RemoveGrant: fi.HasPermission(kwlib.PermUserRemove), + } +} + +func filePermissions(fi *kwlib.FileInfo) *provider.ResourcePermissions { + canView := fi.HasPermission(kwlib.PermView) + return &provider.ResourcePermissions{ + Stat: canView, + GetPath: canView, + InitiateFileDownload: fi.HasPermission(kwlib.PermDownload), + InitiateFileUpload: fi.HasPermission(kwlib.PermVersionCreate), + ListFileVersions: fi.HasPermission(kwlib.PermVersionView), + RestoreFileVersion: fi.HasPermission(kwlib.PermVersionPromote), + Delete: fi.HasPermission(kwlib.PermFileDelete), + Move: fi.HasPermission(kwlib.PermFileMove), + } +} + +// spaceRole returns the CS3 role for the current user on a space root folder. +// Used to build the grants opaque in ListStorageSpaces. +// Precedence: Manager > Editor > Viewer. +func spaceRole(fi *kwlib.FileInfo) *provider.ResourcePermissions { + switch { + case fi.HasPermission(kwlib.PermUserAdd): + return conversions.NewManagerRole().CS3ResourcePermissions() + case fi.HasPermission(kwlib.PermFileAdd): + return conversions.NewEditorRole().CS3ResourcePermissions() + default: + return conversions.NewViewerRole().CS3ResourcePermissions() + } +} diff --git a/pkg/storage/fs/kiteworks/versions.go b/pkg/storage/fs/kiteworks/versions.go new file mode 100644 index 00000000000..9654429cdca --- /dev/null +++ b/pkg/storage/fs/kiteworks/versions.go @@ -0,0 +1,161 @@ +package kiteworks + +import ( + "context" + "io" + "strings" + "time" + + provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" + + "github.com/owncloud/reva/v2/pkg/errtypes" + "github.com/owncloud/reva/v2/pkg/storage" + "github.com/owncloud/reva/v2/pkg/storage/fs/kiteworks/kwlib" +) + +// versionRef identifies a specific KW file version. Version keys are encoded +// as "@" so that GetMD and Download can distinguish them +// from regular file references and look up the right KW resource. +type versionRef struct { + fileID string + versionID string +} + +// parseVersionRef returns a versionRef if ref.OpaqueId encodes a version +// ("fileID@versionID"), and false otherwise. +func parseVersionRef(ref *provider.Reference) (versionRef, bool) { + fileID, versionID, ok := strings.Cut(ref.GetResourceId().GetOpaqueId(), "@") + if !ok { + return versionRef{}, false + } + return versionRef{fileID: fileID, versionID: versionID}, true +} + +// getVersionMD returns a minimal ResourceInfo for a version reference. +// Calls GetFileByID on the parent file to verify access; the returned info +// carries only what handleGet requires (Type=FILE, no "processing" status). +func (d *Driver) getVersionMD(ctx context.Context, ref *provider.Reference, vr versionRef) (*provider.ResourceInfo, error) { + if _, err := d.client(ctx).GetFileByID(vr.fileID); err != nil { + return nil, err + } + return &provider.ResourceInfo{ + Id: ref.GetResourceId(), + Type: provider.ResourceType_RESOURCE_TYPE_FILE, + }, nil +} + +// downloadVersion streams the content of a specific version, used by Download +// when the ref carries a version OpaqueId. +func (d *Driver) downloadVersion(ctx context.Context, ref *provider.Reference, vr versionRef, openReaderFunc func(*provider.ResourceInfo) bool) (*provider.ResourceInfo, io.ReadCloser, error) { + fi, err := d.client(ctx).GetFileByID(vr.fileID) + if err != nil { + return nil, nil, err + } + if !fi.HasPermission(kwlib.PermDownload) { + return nil, nil, errtypes.PermissionDenied(vr.fileID) + } + // Fetch content first to read Content-Length from the response headers. + // This is required because the dataprovider uses ri.Size to set Content-Length, + // and an empty size causes browsers to receive a zero-byte file. + resp, err := d.client(ctx).GetVersionContents(vr.fileID, vr.versionID) + if err != nil { + return nil, nil, err + } + ri := &provider.ResourceInfo{ + Id: ref.GetResourceId(), + Type: provider.ResourceType_RESOURCE_TYPE_FILE, + Name: fi.Name, + Path: fi.Path, + Size: uint64(max(resp.ContentLength, 0)), + } + if !openReaderFunc(ri) { + resp.Body.Close() + return ri, nil, nil + } + return ri, resp.Body, nil +} + +// ListRevisions returns file versions as CS3 FileVersion entries. +// Keys are encoded as "@" so that GetMD and Download +// can route version download requests to the right KW endpoint. +func (d *Driver) ListRevisions(ctx context.Context, ref *provider.Reference) ([]*provider.FileVersion, error) { + nodeID, _, err := d.resolveRef(ctx, ref) + if err != nil { + return nil, err + } + fi, err := d.client(ctx).GetFileByID(nodeID) + if err != nil { + return nil, err + } + if !fi.HasPermission(kwlib.PermVersionView) { + return nil, errtypes.PermissionDenied(nodeID) + } + versions, err := d.client(ctx).GetFileVersions(nodeID) + if err != nil { + return nil, err + } + revs := make([]*provider.FileVersion, 0, len(versions)) + for _, v := range versions { + revs = append(revs, &provider.FileVersion{ + Key: nodeID + "@" + v.ID, + Mtime: uint64(time.Time(v.Created).Unix()), + Size: uint64(v.Size), + }) + } + return revs, nil +} + +// RestoreRevision promotes revisionKey to be the current version of the file. +// revisionKey is "@" as returned by ListRevisions. +// ocdav passes the space-root ResourceId as ref (opaque_id == space_id), so we +// use the fileID embedded in the key rather than resolveRef. +func (d *Driver) RestoreRevision(ctx context.Context, _ *provider.Reference, revisionKey string) (*storage.RestoreRevisionResult, error) { + fileID, versionID, ok := strings.Cut(revisionKey, "@") + if !ok { + return nil, errtypes.BadRequest("kiteworks: invalid revision key: " + revisionKey) + } + fi, err := d.client(ctx).GetFileByID(fileID) + if err != nil { + return nil, err + } + if !fi.HasPermission(kwlib.PermVersionPromote) { + return nil, errtypes.PermissionDenied(fileID) + } + if err := d.client(ctx).PromoteFileVersion(fileID, versionID); err != nil { + return nil, err + } + return &storage.RestoreRevisionResult{}, nil +} + +// DownloadRevision streams the content of the revision identified by revisionKey. +// revisionKey is "@" as returned by ListRevisions. +func (d *Driver) DownloadRevision(ctx context.Context, ref *provider.Reference, revisionKey string, openReaderFunc func(*provider.ResourceInfo) bool) (*provider.ResourceInfo, io.ReadCloser, error) { + nodeID, _, err := d.resolveRef(ctx, ref) + if err != nil { + return nil, nil, err + } + fi, err := d.client(ctx).GetFileByID(nodeID) + if err != nil { + return nil, nil, err + } + if !fi.HasPermission(kwlib.PermDownload) { + return nil, nil, errtypes.PermissionDenied(nodeID) + } + _, versionID, _ := strings.Cut(revisionKey, "@") + resp, err := d.client(ctx).GetVersionContents(nodeID, versionID) + if err != nil { + return nil, nil, err + } + ri := &provider.ResourceInfo{ + Id: ref.GetResourceId(), + Type: provider.ResourceType_RESOURCE_TYPE_FILE, + Name: fi.Name, + Path: fi.Path, + Size: uint64(max(resp.ContentLength, 0)), + } + if !openReaderFunc(ri) { + resp.Body.Close() + return ri, nil, nil + } + return ri, resp.Body, nil +} diff --git a/pkg/storage/storage.go b/pkg/storage/storage.go index ca405033df0..154166b57f3 100644 --- a/pkg/storage/storage.go +++ b/pkg/storage/storage.go @@ -259,8 +259,9 @@ type DeleteStorageSpaceResult struct { // UploadSource carries the staged bytes for a CommitUpload call. type UploadSource struct { - Body io.ReadCloser - Length int64 + Body io.ReadCloser + Length int64 + NodeExisted bool // ScanResult is the antivirus verdict: empty means clean. ScanResult string diff --git a/pkg/upload/coordinator.go b/pkg/upload/coordinator.go index 9af3bab176c..7625e24fd36 100644 --- a/pkg/upload/coordinator.go +++ b/pkg/upload/coordinator.go @@ -547,10 +547,11 @@ func (c *coordinator) commit(ctx context.Context, session Session) (*provider.Re // CommitUpload does not own the body; we opened it, so we close it. err = c.fs.CommitUpload(ctx, &ref, session.ID(), storage.UploadSource{ - Body: f, - Length: session.Size(), - ScanResult: scanResult, - ScanDate: scanDate, + Body: f, + Length: session.Size(), + NodeExisted: session.NodeExists(), + ScanResult: scanResult, + ScanDate: scanDate, }) f.Close() if err != nil {