From 6006bb4d23f425ce5e10abbc149258d1da2cace5 Mon Sep 17 00:00:00 2001 From: Firas Frikha Date: Fri, 14 Aug 2026 18:52:53 +0200 Subject: [PATCH 1/8] feat: route data-path uploads through the coordinator --- .../services/dataprovider/dataprovider.go | 26 ++++++++++++++++--- pkg/rhttp/datatx/datatx.go | 5 +++- pkg/rhttp/datatx/manager/simple/simple.go | 7 ++--- pkg/rhttp/datatx/manager/spaces/spaces.go | 7 ++--- pkg/rhttp/datatx/manager/tus/tus.go | 11 ++++---- 5 files changed, 41 insertions(+), 15 deletions(-) diff --git a/internal/http/services/dataprovider/dataprovider.go b/internal/http/services/dataprovider/dataprovider.go index bffe73bebe1..1ca126d8e2d 100644 --- a/internal/http/services/dataprovider/dataprovider.go +++ b/internal/http/services/dataprovider/dataprovider.go @@ -33,6 +33,7 @@ import ( "github.com/owncloud/reva/v2/pkg/rhttp/router" "github.com/owncloud/reva/v2/pkg/storage" "github.com/owncloud/reva/v2/pkg/storage/fs/registry" + "github.com/owncloud/reva/v2/pkg/upload" ) func init() { @@ -51,6 +52,7 @@ type config struct { NatsEnableTLS bool `mapstructure:"nats_enable_tls"` NatsUsername string `mapstructure:"nats_username"` NatsPassword string `mapstructure:"nats_password"` + UploadDirectory string `mapstructure:"upload_directory" docs:";Local directory for staging upload sessions. Overrides the driver's root. Required for drivers that have no local filesystem root."` } func (c *config) init() { @@ -104,7 +106,12 @@ func New(m map[string]interface{}, log *zerolog.Logger) (global.Service, error) return nil, err } - dataTXs, err := getDataTXs(conf, fs, evstream, log) + coord, err := getCoordinator(conf, fs, evstream, log) + if err != nil { + return nil, err + } + + dataTXs, err := getDataTXs(conf, coord, fs, evstream, log) if err != nil { return nil, err } @@ -126,7 +133,20 @@ func getFS(c *config, stream events.Stream, log *zerolog.Logger) (storage.FS, er return nil, fmt.Errorf("driver not found: %s", c.Driver) } -func getDataTXs(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (map[string]http.Handler, error) { +// getCoordinator builds the coordinator that owns the upload lifecycle for the +// driver this service mounts. +func getCoordinator(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (upload.Coordinator, error) { + store := upload.NewFileStoreFromConfig(c.UploadDirectory, c.Drivers[c.Driver], log) + if store == nil { + return nil, fmt.Errorf("dataprovider: cannot determine the upload directory, set upload_directory") + } + if err := store.Setup(); err != nil { + return nil, fmt.Errorf("dataprovider: upload directory setup failed: %w", err) + } + return upload.NewCoordinator(fs, store, store.UploadDir(), publisher), nil +} + +func getDataTXs(c *config, coord upload.Coordinator, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (map[string]http.Handler, error) { if c.DataTXs == nil { c.DataTXs = make(map[string]map[string]interface{}) } @@ -146,7 +166,7 @@ func getDataTXs(c *config, fs storage.FS, publisher events.Publisher, log *zerol for t := range c.DataTXs { if f, ok := datatxregistry.NewFuncs[t]; ok { if tx, err := f(c.DataTXs[t], publisher, log); err == nil { - if handler, err := tx.Handler(fs); err == nil { + if handler, err := tx.Handler(coord, fs); err == nil { txs[t] = handler } } diff --git a/pkg/rhttp/datatx/datatx.go b/pkg/rhttp/datatx/datatx.go index b770f73b890..6dd36b86eb6 100644 --- a/pkg/rhttp/datatx/datatx.go +++ b/pkg/rhttp/datatx/datatx.go @@ -28,12 +28,15 @@ import ( provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" "github.com/owncloud/reva/v2/pkg/events" "github.com/owncloud/reva/v2/pkg/storage" + "github.com/owncloud/reva/v2/pkg/upload" "github.com/owncloud/reva/v2/pkg/utils" ) // DataTX provides an abstraction around various data transfer protocols. type DataTX interface { - Handler(fs storage.FS) (http.Handler, error) + // Handler serves the protocol's data path. Uploads go through coord, which + // owns the upload lifecycle for every driver; downloads read from driver. + Handler(coord upload.Coordinator, driver storage.FS) (http.Handler, error) } // EmitFileUploadedEvent is a helper function which publishes a FileUploaded event diff --git a/pkg/rhttp/datatx/manager/simple/simple.go b/pkg/rhttp/datatx/manager/simple/simple.go index 39e60e15662..b7141e8035f 100644 --- a/pkg/rhttp/datatx/manager/simple/simple.go +++ b/pkg/rhttp/datatx/manager/simple/simple.go @@ -40,6 +40,7 @@ import ( "github.com/owncloud/reva/v2/pkg/storage" "github.com/owncloud/reva/v2/pkg/storage/cache" "github.com/owncloud/reva/v2/pkg/storagespace" + "github.com/owncloud/reva/v2/pkg/upload" "github.com/owncloud/reva/v2/pkg/utils" ) @@ -78,7 +79,7 @@ func New(m map[string]interface{}, publisher events.Publisher, log *zerolog.Logg }, nil } -func (m *manager) Handler(fs storage.FS) (http.Handler, error) { +func (m *manager) Handler(coord upload.Coordinator, driver storage.FS) (http.Handler, error) { h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { sublog := m.log.With().Str("path", r.URL.Path).Logger() r = r.WithContext(appctx.WithLogger(r.Context(), &sublog)) @@ -92,7 +93,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) { metrics.DownloadsActive.Sub(1) }() } - download.GetOrHeadFile(w, r, fs, "") + download.GetOrHeadFile(w, r, driver, "") case "PUT": metrics.UploadsActive.Add(1) defer func() { @@ -114,7 +115,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) { ctx = ctxpkg.ContextSetLockID(ctx, lockID) } - info, err := fs.Upload(ctx, storage.UploadRequest{ + info, err := coord.Upload(ctx, storage.UploadRequest{ Ref: ref, Body: r.Body, Length: r.ContentLength, diff --git a/pkg/rhttp/datatx/manager/spaces/spaces.go b/pkg/rhttp/datatx/manager/spaces/spaces.go index b0b2f3b6ad8..514bccf1df4 100644 --- a/pkg/rhttp/datatx/manager/spaces/spaces.go +++ b/pkg/rhttp/datatx/manager/spaces/spaces.go @@ -42,6 +42,7 @@ import ( "github.com/owncloud/reva/v2/pkg/storage" "github.com/owncloud/reva/v2/pkg/storage/cache" "github.com/owncloud/reva/v2/pkg/storagespace" + "github.com/owncloud/reva/v2/pkg/upload" "github.com/owncloud/reva/v2/pkg/utils" ) @@ -80,7 +81,7 @@ func New(m map[string]interface{}, publisher events.Publisher, log *zerolog.Logg }, nil } -func (m *manager) Handler(fs storage.FS) (http.Handler, error) { +func (m *manager) Handler(coord upload.Coordinator, driver storage.FS) (http.Handler, error) { h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { var spaceID string spaceID, r.URL.Path = router.ShiftPath(r.URL.Path) @@ -97,7 +98,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) { metrics.DownloadsActive.Sub(1) }() } - download.GetOrHeadFile(w, r, fs, spaceID) + download.GetOrHeadFile(w, r, driver, spaceID) case "PUT": metrics.UploadsActive.Add(1) defer func() { @@ -117,7 +118,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) { Path: fn, } var info *provider.ResourceInfo - info, err = fs.Upload(ctx, storage.UploadRequest{ + info, err = coord.Upload(ctx, storage.UploadRequest{ Ref: ref, Body: r.Body, Length: r.ContentLength, diff --git a/pkg/rhttp/datatx/manager/tus/tus.go b/pkg/rhttp/datatx/manager/tus/tus.go index ac0e037846e..5e364a05a21 100644 --- a/pkg/rhttp/datatx/manager/tus/tus.go +++ b/pkg/rhttp/datatx/manager/tus/tus.go @@ -40,6 +40,7 @@ import ( "github.com/owncloud/reva/v2/pkg/rhttp/datatx/metrics" "github.com/owncloud/reva/v2/pkg/storage" "github.com/owncloud/reva/v2/pkg/storagespace" + "github.com/owncloud/reva/v2/pkg/upload" ) func init() { @@ -87,8 +88,8 @@ func New(m map[string]interface{}, publisher events.Publisher, log *zerolog.Logg }, nil } -func (m *manager) Handler(fs storage.FS) (http.Handler, error) { - composable, ok := fs.(storage.ComposableFS) +func (m *manager) Handler(_ upload.Coordinator, driver storage.FS) (http.Handler, error) { + composable, ok := driver.(storage.ComposableFS) if !ok { return nil, errtypes.NotSupported("file system does not support the tus protocol") } @@ -130,7 +131,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) { return nil, err } - if usl, ok := fs.(storage.UploadSessionLister); ok { + if usl, ok := driver.(storage.UploadSessionLister); ok { // We can currently only send updates if the fs is decomposedfs as we read very specific keys from the storage map of the tus info go func() { for { @@ -174,7 +175,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) { metrics.UploadsActive.Sub(1) }() // set etag, mtime and file id - setHeaders(fs, w, r) + setHeaders(driver, w, r) handler.PostFile(w, r) case "HEAD": handler.HeadFile(w, r) @@ -184,7 +185,7 @@ func (m *manager) Handler(fs storage.FS) (http.Handler, error) { metrics.UploadsActive.Sub(1) }() // set etag, mtime and file id - setHeaders(fs, w, r) + setHeaders(driver, w, r) handler.PatchFile(w, r) case "DELETE": handler.DelFile(w, r) From f0ee164c0770ab779d967ecc56fe19f8ae753bca Mon Sep 17 00:00:00 2001 From: Firas Frikha Date: Fri, 14 Aug 2026 18:56:44 +0200 Subject: [PATCH 2/8] feat: serve the tus protocol from the coordinator --- pkg/rhttp/datatx/manager/tus/tus.go | 74 +++++++++++------------------ 1 file changed, 29 insertions(+), 45 deletions(-) diff --git a/pkg/rhttp/datatx/manager/tus/tus.go b/pkg/rhttp/datatx/manager/tus/tus.go index 5e364a05a21..d846525289e 100644 --- a/pkg/rhttp/datatx/manager/tus/tus.go +++ b/pkg/rhttp/datatx/manager/tus/tus.go @@ -33,7 +33,6 @@ import ( "github.com/owncloud/reva/v2/internal/http/services/owncloud/ocdav/net" "github.com/owncloud/reva/v2/pkg/appctx" - "github.com/owncloud/reva/v2/pkg/errtypes" "github.com/owncloud/reva/v2/pkg/events" "github.com/owncloud/reva/v2/pkg/rhttp/datatx" "github.com/owncloud/reva/v2/pkg/rhttp/datatx/manager/registry" @@ -88,20 +87,14 @@ func New(m map[string]interface{}, publisher events.Publisher, log *zerolog.Logg }, nil } -func (m *manager) Handler(_ upload.Coordinator, driver storage.FS) (http.Handler, error) { - composable, ok := driver.(storage.ComposableFS) - if !ok { - return nil, errtypes.NotSupported("file system does not support the tus protocol") - } - +func (m *manager) Handler(coord upload.Coordinator, _ storage.FS) (http.Handler, error) { // A storage backend for tusd may consist of multiple different parts which // handle upload creation, locking, termination and so on. The composer is a - // place where all those separated pieces are joined together. In this example - // we only use the file store but you may plug in multiple. + // place where all those separated pieces are joined together. composer := tusd.NewStoreComposer() - // let the composable storage tell tus which extensions it supports - composable.UseIn(composer) + // the coordinator serves the tus protocol on behalf of every driver + coord.UseIn(composer) config := tusd.Config{ StoreComposer: composer, @@ -131,33 +124,29 @@ func (m *manager) Handler(_ upload.Coordinator, driver storage.FS) (http.Handler return nil, err } - if usl, ok := driver.(storage.UploadSessionLister); ok { - // We can currently only send updates if the fs is decomposedfs as we read very specific keys from the storage map of the tus info - go func() { - for { - ev := <-handler.CompleteUploads - // We should be able to get the upload progress with fs.GetUploadProgress, but currently tus will erase the info files - // so we create a Progress instance here that is used to read the correct properties - ups, err := usl.ListUploadSessions(context.Background(), storage.UploadSessionFilter{ID: &ev.Upload.ID}) - if err != nil { - appctx.GetLogger(context.Background()).Error().Err(err).Str("session", ev.Upload.ID).Msg("failed to list upload session") - } else { - if len(ups) < 1 { - appctx.GetLogger(context.Background()).Error().Str("session", ev.Upload.ID).Msg("upload session not found") - continue - } - up := ups[0] - executant := up.Executant() - ref := up.Reference() - if m.publisher != nil { - if err := datatx.EmitFileUploadedEvent(up.SpaceOwner(), &executant, &ref, m.publisher); err != nil { - appctx.GetLogger(context.Background()).Error().Err(err).Msg("failed to publish FileUploaded event") - } - } + go func() { + for { + ev := <-handler.CompleteUploads + // tus erases its info files, so read the session back for the event's properties + ups, err := coord.ListUploadSessions(context.Background(), storage.UploadSessionFilter{ID: &ev.Upload.ID}) + if err != nil { + appctx.GetLogger(context.Background()).Error().Err(err).Str("session", ev.Upload.ID).Msg("failed to list upload session") + continue + } + if len(ups) < 1 { + appctx.GetLogger(context.Background()).Error().Str("session", ev.Upload.ID).Msg("upload session not found") + continue + } + up := ups[0] + executant := up.Executant() + ref := up.Reference() + if m.publisher != nil { + if err := datatx.EmitFileUploadedEvent(up.SpaceOwner(), &executant, &ref, m.publisher); err != nil { + appctx.GetLogger(context.Background()).Error().Err(err).Msg("failed to publish FileUploaded event") } } - }() - } + } + }() h := handler.Middleware(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { sublog := m.log.With().Str("uploadid", r.URL.Path).Logger() @@ -175,7 +164,7 @@ func (m *manager) Handler(_ upload.Coordinator, driver storage.FS) (http.Handler metrics.UploadsActive.Sub(1) }() // set etag, mtime and file id - setHeaders(driver, w, r) + setHeaders(coord, w, r) handler.PostFile(w, r) case "HEAD": handler.HeadFile(w, r) @@ -185,7 +174,7 @@ func (m *manager) Handler(_ upload.Coordinator, driver storage.FS) (http.Handler metrics.UploadsActive.Sub(1) }() // set etag, mtime and file id - setHeaders(driver, w, r) + setHeaders(coord, w, r) handler.PatchFile(w, r) case "DELETE": handler.DelFile(w, r) @@ -206,15 +195,10 @@ func (m *manager) Handler(_ upload.Coordinator, driver storage.FS) (http.Handler return h, nil } -func setHeaders(fs storage.FS, w http.ResponseWriter, r *http.Request) { +func setHeaders(coord upload.Coordinator, w http.ResponseWriter, r *http.Request) { ctx := r.Context() id := path.Base(r.URL.Path) - datastore, ok := fs.(tusd.DataStore) - if !ok { - appctx.GetLogger(ctx).Error().Interface("fs", fs).Msg("storage is not a tus datastore") - return - } - upload, err := datastore.GetUpload(ctx, id) + upload, err := coord.GetUpload(ctx, id) if err != nil { appctx.GetLogger(ctx).Error().Err(err).Msg("could not get upload from storage") return From 77f9cff127a167be5e56052c7d4486e775fd9d58 Mon Sep 17 00:00:00 2001 From: Firas Frikha Date: Fri, 14 Aug 2026 23:38:48 +0200 Subject: [PATCH 3/8] feat: consume postprocessing results in the coordinator --- .../storageprovider/storageprovider.go | 43 +- .../services/dataprovider/dataprovider.go | 8 + .../fs/nextcloud/nextcloud_server_mock.go | 13 +- .../utils/decomposedfs/upload_async_test.go | 729 ------------------ .../fixtures/storageprovider-nextcloud.toml | 2 + .../integration/grpc/storageprovider_test.go | 8 +- 6 files changed, 64 insertions(+), 739 deletions(-) delete mode 100644 pkg/storage/utils/decomposedfs/upload_async_test.go diff --git a/internal/grpc/services/storageprovider/storageprovider.go b/internal/grpc/services/storageprovider/storageprovider.go index b23604a275f..39692a1188c 100644 --- a/internal/grpc/services/storageprovider/storageprovider.go +++ b/internal/grpc/services/storageprovider/storageprovider.go @@ -47,6 +47,7 @@ import ( "github.com/owncloud/reva/v2/pkg/storage" "github.com/owncloud/reva/v2/pkg/storage/fs/registry" "github.com/owncloud/reva/v2/pkg/storagespace" + "github.com/owncloud/reva/v2/pkg/upload" "github.com/owncloud/reva/v2/pkg/utils" "github.com/pkg/errors" "github.com/rs/zerolog" @@ -71,6 +72,7 @@ type config struct { MountID string `mapstructure:"mount_id"` UploadExpiration int64 `mapstructure:"upload_expiration" docs:"0;Duration for how long uploads will be valid."` Events eventconfig `mapstructure:"events" docs:"0;Event stream configuration"` + UploadDirectory string `mapstructure:"upload_directory" docs:";Local directory for staging upload sessions. Overrides the driver's root. Required for drivers that have no local filesystem root."` } type eventconfig struct { @@ -106,6 +108,7 @@ func (c *config) init() { type Service struct { conf *config Storage storage.FS + Coordinator upload.Coordinator dataServerURL *url.URL availableXS []*provider.ResourceChecksumPriority } @@ -175,7 +178,14 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc. c.init() - fs, err := getFS(c, log) + // One stream for both the driver and the coordinator: a second one would open a + // second nats connection for the same events. + evstream, err := estreamFromConfig(c.Events) + if err != nil { + return nil, err + } + + fs, err := getFS(c, evstream, log) if err != nil { return nil, err } @@ -202,9 +212,15 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc. return nil, err } + coord, err := getCoordinator(c, fs, evstream, log) + if err != nil { + return nil, err + } + service := &Service{ conf: c, Storage: fs, + Coordinator: coord, dataServerURL: u, availableXS: xsTypes, } @@ -212,6 +228,22 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc. return service, nil } +// getCoordinator builds the coordinator that initiates uploads for the driver +// this service mounts. It stages sessions in the same directory the dataprovider +// appends bytes to, so an upload initiated here can be continued there. +func getCoordinator(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (upload.Coordinator, error) { + store := upload.NewFileStoreFromConfig(c.UploadDirectory, c.Drivers[c.Driver], log) + if store == nil { + return nil, fmt.Errorf("storageprovider: cannot determine the upload directory, set upload_directory") + } + if err := store.Setup(); err != nil { + return nil, fmt.Errorf("storageprovider: upload directory setup failed: %w", err) + } + + // No chunk folder: only the data path assembles chunks. + return upload.NewCoordinator(fs, store, "", publisher), nil +} + func (s *Service) SetArbitraryMetadata(ctx context.Context, req *provider.SetArbitraryMetadataRequest) (*provider.SetArbitraryMetadataResponse, error) { ctx = ctxpkg.ContextSetLockID(ctx, req.LockId) @@ -427,7 +459,7 @@ func (s *Service) InitiateFileUpload(ctx context.Context, req *provider.Initiate metadata["expires"] = strconv.Itoa(int(expirationTimestamp.Seconds)) } - uploadIDs, err := s.Storage.InitiateUpload(ctx, req.Ref, uploadLength, metadata) + uploadIDs, err := s.Coordinator.InitiateUpload(ctx, req.Ref, uploadLength, metadata) if err != nil { var st *rpc.Status switch err.(type) { @@ -1266,12 +1298,7 @@ func (s *Service) addMissingStorageProviderID(resourceID *provider.ResourceId, s } } -func getFS(c *config, log *zerolog.Logger) (storage.FS, error) { - evstream, err := estreamFromConfig(c.Events) - if err != nil { - return nil, err - } - +func getFS(c *config, evstream events.Stream, log *zerolog.Logger) (storage.FS, error) { if f, ok := registry.NewFuncs[c.Driver]; ok { driverConf := c.Drivers[c.Driver] driverConf["mount_id"] = c.MountID // pass the mount id to the driver diff --git a/internal/http/services/dataprovider/dataprovider.go b/internal/http/services/dataprovider/dataprovider.go index 1ca126d8e2d..b138306e8fe 100644 --- a/internal/http/services/dataprovider/dataprovider.go +++ b/internal/http/services/dataprovider/dataprovider.go @@ -111,6 +111,14 @@ func New(m map[string]interface{}, log *zerolog.Logger) (global.Service, error) return nil, err } + // only the data path consumes postprocessing results: one consumer group gets + // one copy of each event, so a second subscriber would take half of them + if ac := upload.AsyncConfFromDriverConf(conf.Drivers[conf.Driver]); ac.Enabled { + if err := coord.StartPostprocessing(evstream, ac.ConsumerGroup, ac.MountID, ac.NumConsumers); err != nil { + return nil, fmt.Errorf("dataprovider: could not start postprocessing: %w", err) + } + } + dataTXs, err := getDataTXs(conf, coord, fs, evstream, log) if err != nil { return nil, err diff --git a/pkg/storage/fs/nextcloud/nextcloud_server_mock.go b/pkg/storage/fs/nextcloud/nextcloud_server_mock.go index 41dfbe210dc..1221575a855 100644 --- a/pkg/storage/fs/nextcloud/nextcloud_server_mock.go +++ b/pkg/storage/fs/nextcloud/nextcloud_server_mock.go @@ -74,6 +74,11 @@ var responses = map[string]Response{ `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/EmptyRecycle `: {200, ``, serverStateEmpty}, + `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetQuota `: {200, `{"totalBytes":456,"usedBytes":123}`, serverStateEmpty}, + + // the parent of an upload target, stated for its permissions and its id + `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetMD {"ref":{"path":"/"},"mdKeys":[]} EMPTY`: {200, `{"opaque":{},"type":2,"id":{"opaque_id":"fileid-/"},"checksum":{},"etag":"deadbeef","mime_type":"httpd/unix-directory","mtime":{"seconds":1234567890},"path":"/","permission_set":{"initiate_file_upload":true,"stat":true,"list_container":true},"size":12345,"canonical_metadata":{},"arbitrary_metadata":{"metadata":{}}}`, serverStateEmpty}, + `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetMD {"ref":{"path":"/"},"mdKeys":null} EMPTY`: {404, ``, serverStateEmpty}, `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetMD {"ref":{"path":"/"},"mdKeys":null} HOME`: {200, `{"opaque":{},"type":1,"id":{"opaque_id":"fileid-/some/path"},"checksum":{},"etag":"deadbeef","mime_type":"text/plain","mtime":{"seconds":1234567890},"path":"/","permission_set":{},"size":12345,"canonical_metadata":{},"arbitrary_metadata":{"metadata":{"da":"ta","some":"arbi","trary":"meta"}}}`, serverStateHome}, @@ -111,7 +116,13 @@ var responses = map[string]Response{ `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetPathByID {"storage_id":"00000000-0000-0000-0000-000000000000","opaque_id":"fileid-/some/path"} EMPTY`: {200, "/subdir", serverStateEmpty}, - `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetMD {"ref":{"path":"/file"},"mdKeys":null}`: {404, ``, serverStateEmpty}, + `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetMD {"ref":{"path":"/file"},"mdKeys":null}`: {404, ``, serverStateEmpty}, + // the coordinator resolves the upload target itself: the file does not exist yet, + // so it stats the parent for the permissions and the parent id + `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetMD {"ref":{"path":"/file"},"mdKeys":[]}`: {404, ``, serverStateEmpty}, + // a zero-length upload finishes at once, so the node is created and read back here + `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/TouchFile {"ref":{"resource_id":{"opaque_id":"fileid-/"},"path":"file"},"markprocessing":false,"mtime":"1234567890"}`: {200, ``, serverStateEmpty}, + `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/GetMD {"ref":{"resource_id":{"opaque_id":"fileid-/"},"path":"file"},"mdKeys":[]}`: {200, `{"opaque":{},"type":1,"id":{"opaque_id":"fileid-/file"},"checksum":{},"etag":"deadbeef","mime_type":"text/plain","mtime":{"seconds":1234567890},"path":"/file","permission_set":{},"size":0,"canonical_metadata":{},"arbitrary_metadata":{"metadata":{}}}`, serverStateEmpty}, `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/InitiateUpload {"ref":{"path":"/file"},"uploadLength":0,"metadata":{"providerID":""}}`: {200, `{"simple": "yes","tus": "yes"}`, serverStateEmpty}, `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/InitiateUpload {"ref":{"resource_id":{"storage_id":"f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c"},"path":"/versionedFile"},"uploadLength":1,"metadata":{}}`: {200, `{"simple": "yes","tus": "yes"}`, serverStateEmpty}, `POST /apps/sciencemesh/~f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c/api/storage/InitiateUpload {"ref":{"resource_id":{"storage_id":"f7fbf8c8-139b-4376-b307-cf0a8c2d0d9c"},"path":"/versionedFile"},"uploadLength":2,"metadata":{}}`: {200, `{"simple": "yes","tus": "yes"}`, serverStateEmpty}, diff --git a/pkg/storage/utils/decomposedfs/upload_async_test.go b/pkg/storage/utils/decomposedfs/upload_async_test.go deleted file mode 100644 index 65a25e78898..00000000000 --- a/pkg/storage/utils/decomposedfs/upload_async_test.go +++ /dev/null @@ -1,729 +0,0 @@ -package decomposedfs - -import ( - "bytes" - "context" - "io" - "os" - "path/filepath" - - userpb "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1" - cs3permissions "github.com/cs3org/go-cs3apis/cs3/permissions/v1beta1" - v1beta11 "github.com/cs3org/go-cs3apis/cs3/rpc/v1beta1" - provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" - "github.com/owncloud/reva/v2/pkg/appctx" - ruser "github.com/owncloud/reva/v2/pkg/ctx" - "github.com/owncloud/reva/v2/pkg/events" - "github.com/owncloud/reva/v2/pkg/events/stream" - "github.com/owncloud/reva/v2/pkg/rgrpc/todo/pool" - "github.com/owncloud/reva/v2/pkg/storage" - "github.com/owncloud/reva/v2/pkg/storage/cache" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/aspects" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/lookup" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/metadata" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/node" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/options" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/permissions" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/permissions/mocks" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/timemanager" - "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/tree" - treemocks "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/tree/mocks" - "github.com/owncloud/reva/v2/pkg/storagespace" - "github.com/owncloud/reva/v2/pkg/store" - "github.com/owncloud/reva/v2/pkg/utils" - "github.com/owncloud/reva/v2/tests/helpers" - "github.com/rs/zerolog" - "github.com/stretchr/testify/mock" - "google.golang.org/grpc" - - . "github.com/onsi/ginkgo/v2" - . "github.com/onsi/gomega" -) - -var _ = Describe("Async file uploads", Ordered, func() { - var ( - ref = &provider.Reference{ - ResourceId: &provider.ResourceId{ - SpaceId: "u-s-e-r-id", - }, - Path: "/foo", - } - - rootRef = &provider.Reference{ - ResourceId: &provider.ResourceId{ - SpaceId: "u-s-e-r-id", - OpaqueId: "u-s-e-r-id", - }, - Path: "/", - } - - user = &userpb.User{ - Id: &userpb.UserId{ - Idp: "idp", - OpaqueId: "u-s-e-r-id", - Type: userpb.UserType_USER_TYPE_PRIMARY, - }, - Username: "username", - } - - firstContent = []byte("0123456789") - secondContent = []byte("01234567890123456789") - - ctx context.Context - - pub chan interface{} - con chan interface{} - uploadID string - - fs storage.FS - o *options.Options - lu *lookup.Lookup - pmock *mocks.PermissionsChecker - cs3permissionsclient *mocks.CS3PermissionsClient - permissionsSelector pool.Selectable[cs3permissions.PermissionsAPIClient] - bs *treemocks.Blobstore - - succeedPostprocessing = func(uploadID string) { - // finish postprocessing - con <- events.PostprocessingFinished{ - UploadID: uploadID, - Outcome: events.PPOutcomeContinue, - } - // wait for upload to be ready - ev, ok := (<-pub).(events.UploadReady) - Expect(ok).To(BeTrue()) - Expect(ev.Failed).To(BeFalse()) - Expect(ev.ResourceID).ToNot(BeNil()) - Expect(ev.ResourceID.OpaqueId).ToNot(BeEmpty()) - Expect(ev.ResourceID.OpaqueId).ToNot(Equal(ev.ResourceID.SpaceId), "ResourceID.OpaqueId should be the file node ID, not the space ID") - } - - failPostprocessing = func(uploadID string, outcome events.PostprocessingOutcome) { - // finish postprocessing - con <- events.PostprocessingFinished{ - UploadID: uploadID, - Outcome: outcome, - } - // wait for upload to be ready - ev, ok := (<-pub).(events.UploadReady) - Expect(ok).To(BeTrue()) - Expect(ev.Failed).To(BeTrue()) - } - - fileStatus = func() (bool, string, int) { - // check processing status - resources, err := fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(BeElementOf([2]int{0, 1}), "should not have more than one child") - - item := resources[0] - Expect(item.Path).To(Equal(ref.Path)) - return len(resources) == 1, utils.ReadPlainFromOpaque(item.Opaque, "status"), int(item.GetSize()) - } - parentSize = func() int { - parentInfo, err := fs.GetMD(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - return int(parentInfo.Size) - } - revisionCount = func() int { - revisions, err := fs.ListRevisions(ctx, ref) - Expect(err).ToNot(HaveOccurred()) - return len(revisions) - } - ) - - BeforeEach(func() { - zl := zerolog.New(os.Stdout).Level(zerolog.DebugLevel) - ctx = appctx.WithLogger(ruser.ContextSetUser(context.Background(), user), &zl) - - // setup test - tmpRoot, err := helpers.TempDir("reva-unit-tests-*-root") - Expect(err).ToNot(HaveOccurred()) - - o, err = options.New(map[string]interface{}{ - "root": tmpRoot, - "asyncfileuploads": true, - "treetime_accounting": true, - "treesize_accounting": true, - }) - Expect(err).ToNot(HaveOccurred()) - - lu = lookup.New(metadata.NewXattrsBackend(o.Root, cache.Config{}), o, &timemanager.Manager{}) - pmock = &mocks.PermissionsChecker{} - - cs3permissionsclient = &mocks.CS3PermissionsClient{} - pool.RemoveSelector("PermissionsSelector" + "any") - permissionsSelector = pool.GetSelector[cs3permissions.PermissionsAPIClient]( - "PermissionsSelector", - "any", - func(cc grpc.ClientConnInterface) cs3permissions.PermissionsAPIClient { - return cs3permissionsclient - }, - ) - bs = &treemocks.Blobstore{} - - // create space uses CheckPermission endpoint - cs3permissionsclient.On("CheckPermission", mock.Anything, mock.Anything, mock.Anything).Return(&cs3permissions.CheckPermissionResponse{ - Status: &v1beta11.Status{Code: v1beta11.Code_CODE_OK}, - }, nil).Times(1) - - // for this test we don't care about permissions - pmock.On("AssemblePermissions", mock.Anything, mock.Anything). - Return(&provider.ResourcePermissions{ - Stat: true, - GetQuota: true, - InitiateFileUpload: true, - ListContainer: true, - ListFileVersions: true, - }, nil) - - // setup fs - pub, con = make(chan interface{}), make(chan interface{}) - tree := tree.New(lu, bs, o, store.Create(), &zerolog.Logger{}) - - aspects := aspects.Aspects{ - Lookup: lu, - Tree: tree, - Permissions: permissions.NewPermissions(pmock, permissionsSelector), - EventStream: stream.Chan{pub, con}, - Trashbin: &DecomposedfsTrashbin{}, - } - fs, err = New(o, aspects, &zerolog.Logger{}) - Expect(err).ToNot(HaveOccurred()) - - resp, err := fs.CreateStorageSpace(ctx, &provider.CreateStorageSpaceRequest{Owner: user, Type: "personal"}) - Expect(err).ToNot(HaveOccurred()) - Expect(resp.Status.Code).To(Equal(v1beta11.Code_CODE_OK)) - resID, err := storagespace.ParseID(resp.StorageSpace.Id.OpaqueId) - Expect(err).ToNot(HaveOccurred()) - ref.ResourceId = &resID - - bs.On("Upload", mock.AnythingOfType("*node.Node"), mock.AnythingOfType("string"), mock.Anything). - Return(nil). - Run(func(args mock.Arguments) { - n := args.Get(0).(*node.Node) - data, err := os.ReadFile(args.Get(1).(string)) - Expect(err).ToNot(HaveOccurred()) - Expect(len(data)).To(Equal(int(n.Blobsize))) - }) - - // start upload of a file - uploadIds, err := fs.InitiateUpload(ctx, ref, 10, map[string]string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(uploadIds)).To(Equal(2)) - Expect(uploadIds["simple"]).ToNot(BeEmpty()) - Expect(uploadIds["tus"]).ToNot(BeEmpty()) - - uploadRef := &provider.Reference{Path: "/" + uploadIds["simple"]} - - _, err = fs.Upload(ctx, storage.UploadRequest{ - Ref: uploadRef, - Body: io.NopCloser(bytes.NewReader(firstContent)), - Length: int64(len(firstContent)), - }, nil) - Expect(err).ToNot(HaveOccurred()) - - uploadID = uploadIds["simple"] - - // wait for bytes received event - _, ok := (<-pub).(events.BytesReceived) - Expect(ok).To(BeTrue()) - - // blobstore not called yet - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 0) - }) - - AfterEach(func() { - if o.Root != "" { - os.RemoveAll(o.Root) - } - close(pub) - close(con) - }) - - When("the uploaded file is new", func() { - It("succeeds eventually", func() { - // node is created - resources, err := fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(Equal(1)) - - item := resources[0] - Expect(item.Path).To(Equal(ref.Path)) - Expect(utils.ReadPlainFromOpaque(item.Opaque, "status")).To(Equal("processing")) - - succeedPostprocessing(uploadID) - - // blobstore called now - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 1) - - // node ready - resources, err = fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(Equal(1)) - - item = resources[0] - Expect(item.Path).To(Equal(ref.Path)) - Expect(utils.ReadPlainFromOpaque(item.Opaque, "status")).To(BeEmpty()) - - }) - - It("deletes node and bytes when instructed", func() { - // node is created - resources, err := fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(Equal(1)) - - item := resources[0] - Expect(item.Path).To(Equal(ref.Path)) - Expect(utils.ReadPlainFromOpaque(item.Opaque, "status")).To(Equal("processing")) - - // bytes are in dedicated path - _, err = os.Stat(filepath.Join(o.Root, "uploads", uploadID)) - Expect(err).To(BeNil()) - - failPostprocessing(uploadID, events.PPOutcomeDelete) - - // blobstore still not called now - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 0) - - // node gone - resources, err = fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(Equal(0)) - - // bytes gone - _, err = os.Stat(filepath.Join(o.Root, "uploads", uploadID)) - Expect(err).ToNot(BeNil()) - }) - - It("releases the quota and removes the node when the node metadata is unreadable", func() { - // node is created and the optimistic size has been propagated - resources, err := fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(Equal(1)) - Expect(parentSize()).To(Equal(len(firstContent))) - - // simulate an orphaned node: the node file is still there but its - // metadata is gone, e.g. because an ancestor was trashed while the - // upload was in flight. Reading the node now fails. Purge instead of - // removing the file directly, so the cached attributes go as well. - nodePath := lu.InternalPath(ref.GetResourceId().GetSpaceId(), resources[0].GetId().GetOpaqueId()) - Expect(lu.MetadataBackend().Purge(ctx, nodePath)).To(Succeed()) - _, err = node.ReadNode(ctx, lu, ref.GetResourceId().GetSpaceId(), resources[0].GetId().GetOpaqueId(), false, nil, true) - Expect(err).To(HaveOccurred(), "node should be unreadable after purging its metadata") - - // No UploadReady event is published for an orphaned session: there is - // no node left to report on. Wait for the bytes to be cleaned up - // instead of for an event that will never arrive. - con <- events.PostprocessingFinished{ - UploadID: uploadID, - Outcome: events.PPOutcomeContinue, - } - Eventually(func() bool { - _, err := os.Stat(filepath.Join(o.Root, "uploads", uploadID)) - return err != nil - }).Should(BeTrue(), "the upload bytes should be cleaned up") - - // the blob was never written - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 0) - - // the orphaned node is gone ... - _, err = os.Stat(nodePath) - Expect(err).ToNot(BeNil()) - - // ... and most importantly the quota has been released - Eventually(parentSize).Should(Equal(0)) - }) - - It("deletes node and keeps the bytes when instructed", func() { - // node is created - resources, err := fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(Equal(1)) - - item := resources[0] - Expect(item.Path).To(Equal(ref.Path)) - Expect(utils.ReadPlainFromOpaque(item.Opaque, "status")).To(Equal("processing")) - - // bytes are in dedicated path - _, err = os.Stat(filepath.Join(o.Root, "uploads", uploadID)) - Expect(err).To(BeNil()) - - failPostprocessing(uploadID, events.PPOutcomeAbort) - - // blobstore still not called now - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 0) - - // node gone - resources, err = fs.ListFolder(ctx, rootRef, []string{}, []string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(resources)).To(Equal(0)) - - // bytes are still here - _, err = os.Stat(filepath.Join(o.Root, "uploads", uploadID)) - Expect(err).To(BeNil()) - }) - }) - - When("the uploaded file creates a new version", func() { - JustBeforeEach(func() { - succeedPostprocessing(uploadID) - - // make sure there is no version yet - revs, err := fs.ListRevisions(ctx, ref) - Expect(err).To(BeNil()) - Expect(len(revs)).To(Equal(0)) - - // upload again - uploadIds, err := fs.InitiateUpload(ctx, ref, 10, map[string]string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(uploadIds)).To(Equal(2)) - Expect(uploadIds["simple"]).ToNot(BeEmpty()) - Expect(uploadIds["tus"]).ToNot(BeEmpty()) - - uploadRef := &provider.Reference{Path: "/" + uploadIds["simple"]} - - _, err = fs.Upload(ctx, storage.UploadRequest{ - Ref: uploadRef, - Body: io.NopCloser(bytes.NewReader(firstContent)), - Length: int64(len(firstContent)), - }, nil) - Expect(err).ToNot(HaveOccurred()) - - uploadID = uploadIds["simple"] - - // wait for bytes received event - _, ok := (<-pub).(events.BytesReceived) - Expect(ok).To(BeTrue()) - - // version already created - revs, err = fs.ListRevisions(ctx, ref) - Expect(err).To(BeNil()) - Expect(len(revs)).To(Equal(1)) - - // at this stage: blobstore called once for the original file - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 1) - - }) - - It("succeeds eventually, creating a new version", func() { - succeedPostprocessing(uploadID) - - // version still existing - revs, err := fs.ListRevisions(ctx, ref) - Expect(err).To(BeNil()) - Expect(len(revs)).To(Equal(1)) - - // blobstore now called twice - for original file and new version - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 2) - - // bytes are gone from upload path - _, err = os.Stat(filepath.Join(o.Root, "uploads", uploadID)) - Expect(err).ToNot(BeNil()) - }) - - It("removes new version and restores old one when instructed", func() { - _, status, _ := fileStatus() - Expect(status).To(Equal("processing")) - - failPostprocessing(uploadID, events.PPOutcomeDelete) - - _, status, _ = fileStatus() - Expect(status).To(Equal("")) - - // version gone now - revs, err := fs.ListRevisions(ctx, ref) - Expect(err).To(BeNil()) - Expect(len(revs)).To(Equal(0)) - - // bytes are removed from upload path - _, err = os.Stat(filepath.Join(o.Root, "uploads", uploadID)) - Expect(err).ToNot(BeNil()) - - // blobstore still called only once for the original file - bs.AssertNumberOfCalls(GinkgoT(), "Upload", 1) - }) - - }) - When("Two uploads are processed in parallel", func() { - var secondUploadID string - - JustBeforeEach(func() { - // upload again - uploadIds, err := fs.InitiateUpload(ctx, ref, 20, map[string]string{}) - Expect(err).ToNot(HaveOccurred()) - Expect(len(uploadIds)).To(Equal(2)) - Expect(uploadIds["simple"]).ToNot(BeEmpty()) - Expect(uploadIds["tus"]).ToNot(BeEmpty()) - - uploadRef := &provider.Reference{Path: "/" + uploadIds["simple"]} - - _, err = fs.Upload(ctx, storage.UploadRequest{ - Ref: uploadRef, - Body: io.NopCloser(bytes.NewReader(secondContent)), - Length: int64(len(secondContent)), - }, nil) - Expect(err).ToNot(HaveOccurred()) - - secondUploadID = uploadIds["simple"] - - // wait for bytes received event - _, ok := (<-pub).(events.BytesReceived) - Expect(ok).To(BeTrue()) - }) - - It("doesn't remove processing status when first upload is finished", func() { - succeedPostprocessing(uploadID) - - _, status, _ := fileStatus() - // check processing status - Expect(status).To(Equal("processing")) - }) - - It("removes processing status when second upload is finished, even if first isn't", func() { - succeedPostprocessing(secondUploadID) - - _, status, _ := fileStatus() - Expect(status).To(Equal("")) - }) - - It("correctly calculates the size when the second upload is finished, even if first is deleted", func() { - succeedPostprocessing(secondUploadID) - - _, status, size := fileStatus() - Expect(status).To(Equal("")) - // size should match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should match second upload as well - Expect(parentSize()).To(Equal(len(secondContent))) - - failPostprocessing(uploadID, events.PPOutcomeDelete) - - // check processing status - _, _, size = fileStatus() - // size should still match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should still match second upload as well - Expect(parentSize()).To(Equal(len(secondContent))) - }) - - It("the first can succeed before the second succeeds", func() { - succeedPostprocessing(uploadID) - - _, status, size := fileStatus() - // check processing status - Expect(status).To(Equal("processing")) - // size should match the second upload - Expect(size).To(Equal((len(secondContent)))) - - // parent size should match the second upload - Expect(parentSize()).To(Equal(len(secondContent))) - - succeedPostprocessing(secondUploadID) - - // check processing status has been removed - _, status, size = fileStatus() - Expect(status).To(Equal("")) - - // size should still match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should still match second upload - Expect(parentSize()).To(Equal(len(secondContent))) - - // file should have one revision - Expect(revisionCount()).To(Equal(1)) - }) - - It("the first can succeed after the second succeeds", func() { - succeedPostprocessing(secondUploadID) - - _, status, size := fileStatus() - // check processing status has been removed because the most recent upload finished and can be downloaded - Expect(status).To(Equal("")) - // size should match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should match second upload as well - Expect(parentSize()).To(Equal(len(secondContent))) - - succeedPostprocessing(uploadID) - - _, status, size = fileStatus() - // check processing status is still unset - Expect(status).To(Equal("")) - // size should still match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should still match second upload - Expect(parentSize()).To(Equal(len(secondContent))) - - // file should have one revision - Expect(revisionCount()).To(Equal(1)) - }) - - It("the first can succeed before the second fails", func() { - succeedPostprocessing(uploadID) - - _, status, size := fileStatus() - // check processing status - Expect(status).To(Equal("processing")) - // size should match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should match the second upload - Expect(parentSize()).To(Equal(len(secondContent))) - - failPostprocessing(secondUploadID, events.PPOutcomeDelete) - - _, status, size = fileStatus() - // check processing status has been removed - Expect(status).To(Equal("")) - // size should match the first upload - Expect(size).To(Equal(len(firstContent))) - - // parent size should match first upload - Expect(parentSize()).To(Equal(len(firstContent))) - - // file should not have any revisions - Expect(revisionCount()).To(Equal(0)) - }) - - It("the first can succeed after the second fails", func() { - failPostprocessing(secondUploadID, events.PPOutcomeDelete) - - _, _, size := fileStatus() - // check processing status has not been unset - // FIXME we need to fall back to the previous processing id - // Expect(status).To(Equal("processing")) - // size should match the first upload - Expect(size).To(Equal(len(firstContent))) - - // parent size should match first upload as well - Expect(parentSize()).To(Equal(len(firstContent))) - - succeedPostprocessing(uploadID) - - _, status, size := fileStatus() - // check processing status is now unset - Expect(status).To(Equal("")) - // size should still match the first upload - Expect(size).To(Equal(len(firstContent))) - - // parent size should still match first upload - Expect(parentSize()).To(Equal(len(firstContent))) - - // file should not have any revisions - Expect(revisionCount()).To(Equal(0)) - }) - - It("the first can fail before the second succeeds", func() { - failPostprocessing(uploadID, events.PPOutcomeDelete) - - _, status, size := fileStatus() - // check processing status - Expect(status).To(Equal("processing")) - // size should match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should match second upload as well - Expect(parentSize()).To(Equal(len(secondContent))) - - succeedPostprocessing(secondUploadID) - - _, status, size = fileStatus() - // check processing status has been removed - Expect(status).To(Equal("")) - // size should still match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should still match second upload - Expect(parentSize()).To(Equal(len(secondContent))) - - // file should not have any revisions - // FIXME we need to delete the revision - // Expect(revisionCount()).To(Equal(0)) - }) - - It("the first can fail after the second succeeds", func() { - succeedPostprocessing(secondUploadID) - - _, status, size := fileStatus() - // check processing status has been removed because the most recent upload finished and can be downloaded - Expect(status).To(Equal("")) - // size should match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should match second upload as well - Expect(parentSize()).To(Equal(len(secondContent))) - - failPostprocessing(uploadID, events.PPOutcomeDelete) - - _, status, size = fileStatus() - // check processing status is still unset - Expect(status).To(Equal("")) - // size should still match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should still match second upload - Expect(parentSize()).To(Equal(len(secondContent))) - - // file should not have any revisions - // FIXME we need to delete the revision - // Expect(revisionCount()).To(Equal(0)) - }) - - It("the first can fail before the second fails", func() { - failPostprocessing(uploadID, events.PPOutcomeDelete) - - _, status, size := fileStatus() - // check processing status - Expect(status).To(Equal("processing")) - // size should match the second upload - Expect(size).To(Equal(len(secondContent))) - - // parent size should match second upload as well - Expect(parentSize()).To(Equal(len(secondContent))) - - failPostprocessing(secondUploadID, events.PPOutcomeDelete) - - // check file has been removed - // if all uploads have been processed with outcome delete -> delete the file - // exists, _, _ := fileStatus() - // FIXME this should be false, but we are not deleting the resource - // Expect(exists).To(BeFalse()) - - // parent size should be 0 - // FIXME we are not correctly reverting the sizediff - // Expect(parentSize()).To(Equal(0)) - }) - - It("the first can fail after the second fails", func() { - failPostprocessing(secondUploadID, events.PPOutcomeDelete) - - _, status, size := fileStatus() - // check processing status has been removed because the most recent upload finished and can be downloaded - Expect(status).To(Equal("")) - // size should match the first upload - Expect(size).To(Equal(len(firstContent))) - - // parent size should match second first as well - Expect(parentSize()).To(Equal(len(firstContent))) - - failPostprocessing(uploadID, events.PPOutcomeDelete) - - // check file has been removed - // if all uploads have been processed with outcome delete -> delete the file - // exists, _, _ := fileStatus() - // FIXME this should be false, but we are not deleting the resource - // Expect(exists).To(BeFalse()) - - // parent size should be 0 - // FIXME we are not correctly reverting the sizediff - // Expect(parentSize()).To(Equal(0)) - }) - }) -}) diff --git a/tests/integration/grpc/fixtures/storageprovider-nextcloud.toml b/tests/integration/grpc/fixtures/storageprovider-nextcloud.toml index ef85bb5d060..bb143c9759f 100644 --- a/tests/integration/grpc/fixtures/storageprovider-nextcloud.toml +++ b/tests/integration/grpc/fixtures/storageprovider-nextcloud.toml @@ -3,6 +3,8 @@ address = "{{grpc_address}}" [grpc.services.storageprovider] driver = "nextcloud" +# the nextcloud driver has no local root, so name the upload directory here +upload_directory = "{{root}}/uploads" [grpc.services.storageprovider.drivers.nextcloud] endpoint = "http://localhost:8080/apps/sciencemesh/" diff --git a/tests/integration/grpc/storageprovider_test.go b/tests/integration/grpc/storageprovider_test.go index 8cfdba19e5f..aea49ffb20a 100644 --- a/tests/integration/grpc/storageprovider_test.go +++ b/tests/integration/grpc/storageprovider_test.go @@ -35,6 +35,7 @@ import ( "github.com/owncloud/reva/v2/pkg/storage/fs/ocis" "github.com/owncloud/reva/v2/pkg/storage/fs/registry" jwt "github.com/owncloud/reva/v2/pkg/token/manager/jwt" + "github.com/owncloud/reva/v2/pkg/utils" "github.com/owncloud/reva/v2/tests/helpers" . "github.com/onsi/ginkgo/v2" @@ -367,7 +368,12 @@ var _ = Describe("storage providers", func() { assertUploads := func(provider string) { It("returns upload URLs for simple and tus", func() { fileRef := ref(provider, filePath) - res, err := providerClient.InitiateFileUpload(ctx, &storagep.InitiateFileUploadRequest{Ref: fileRef}) + // a fixed mtime, so the nextcloud mock's exact-body matching can match the + // TouchFile the coordinator makes + res, err := providerClient.InitiateFileUpload(ctx, &storagep.InitiateFileUploadRequest{ + Ref: fileRef, + Opaque: utils.AppendPlainToOpaque(nil, "X-OC-Mtime", "1234567890"), + }) Expect(err).ToNot(HaveOccurred()) Expect(res.Status.Code).To(Equal(rpcv1beta1.Code_CODE_OK)) Expect(len(res.Protocols)).To(Equal(2)) From 86133faa8173a098bdb6b6e6c01de37a2df0bf63 Mon Sep 17 00:00:00 2001 From: Firas Frikha Date: Mon, 17 Aug 2026 10:19:17 +0200 Subject: [PATCH 4/8] feat: WIP --- .../utils/decomposedfs/decomposedfs.go | 319 +----------------- 1 file changed, 8 insertions(+), 311 deletions(-) diff --git a/pkg/storage/utils/decomposedfs/decomposedfs.go b/pkg/storage/utils/decomposedfs/decomposedfs.go index eab994b780c..cbe89b187eb 100644 --- a/pkg/storage/utils/decomposedfs/decomposedfs.go +++ b/pkg/storage/utils/decomposedfs/decomposedfs.go @@ -28,9 +28,7 @@ import ( "path/filepath" "strconv" "strings" - "time" - user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1" rpcv1beta1 "github.com/cs3org/go-cs3apis/cs3/rpc/v1beta1" provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" "github.com/jellydator/ttlcache/v2" @@ -49,7 +47,6 @@ import ( "github.com/owncloud/reva/v2/pkg/events" "github.com/owncloud/reva/v2/pkg/logger" "github.com/owncloud/reva/v2/pkg/rgrpc/todo/pool" - "github.com/owncloud/reva/v2/pkg/rhttp/datatx/metrics" "github.com/owncloud/reva/v2/pkg/storage" "github.com/owncloud/reva/v2/pkg/storage/utils/chunking" "github.com/owncloud/reva/v2/pkg/storage/utils/decomposedfs/aspects" @@ -81,11 +78,9 @@ const ( var ( tracer trace.Tracer + // the coordinator consumes the postprocessing events; reverting a revision is + // the driver's own business _registeredEvents = []events.Unmarshaller{ - events.PostprocessingFinished{}, - events.PostprocessingStepFinished{}, - events.RestartPostprocessing{}, - events.CleanUpload{}, events.RevertRevision{}, } ) @@ -267,7 +262,9 @@ func New(o *options.Options, aspects aspects.Aspects, log *zerolog.Logger) (stor return nil, errors.New("need nats for async file processing") } - ch, err := events.Consume(fs.stream, o.Events.ConsumerGroup, _registeredEvents...) + // a group of its own: the coordinator holds o.Events.ConsumerGroup, and one + // group gets one copy of each event + ch, err := events.Consume(fs.stream, o.Events.ConsumerGroup+"-revisions", _registeredEvents...) if err != nil { return nil, err } @@ -277,15 +274,15 @@ func New(o *options.Options, aspects aspects.Aspects, log *zerolog.Logger) (stor } for i := 0; i < o.Events.NumConsumers; i++ { - go fs.Postprocessing(ch) + go fs.ConsumeRevisionEvents(ch) } } return fs, nil } -// Postprocessing starts the postprocessing result collector -func (fs *Decomposedfs) Postprocessing(ch <-chan events.Event) { +// ConsumeRevisionEvents starts the revision event collector +func (fs *Decomposedfs) ConsumeRevisionEvents(ch <-chan events.Event) { log := logger.New() for event := range ch { evCtx := context.Background() @@ -299,180 +296,6 @@ func (fs *Decomposedfs) processEvent(evCtx context.Context, event events.Event, defer span.End() switch ev := event.Event.(type) { - case events.PostprocessingFinished: - sublog := log.With().Str("event", "PostprocessingFinished").Str("uploadid", ev.UploadID).Logger() - if ev.ResourceID != nil && ev.ResourceID.GetStorageId() != "" && ev.ResourceID.GetStorageId() != fs.o.MountID { - sublog.Debug().Msg("ignoring event for different storage") - return - } - session, err := fs.sessionStore.Get(ctx, ev.UploadID) - if err != nil { - sublog.Error().Err(err).Msg("Failed to get upload") - return // NOTE: since we can't get the upload, we can't delete the blob - } - - ctx = session.Context(ctx) - - n, err := session.Node(ctx) - if err != nil { - // The node metadata is unreadable, so this upload can never finish: - // the destination cannot be resolved. Clean the session up instead of - // leaving it behind to be retried forever. Cleanup falls back to the - // session metadata to release the quota. - sublog.Error().Err(err).Msg("could not read node, cleaning up orphaned session") - session.Cleanup(true, true, true, false) - return - } - sublog = log.With().Str("spaceid", session.SpaceID()).Str("nodeid", session.NodeID()).Logger() - if !n.Exists { - sublog.Debug().Msg("node no longer exists") - session.Cleanup(false, true, true, false) - return - } - - var ( - failed bool - revertNodeMetadata bool - keepUpload bool - ) - unmarkPostprocessing := true - - switch ev.Outcome { - default: - sublog.Error().Str("outcome", string(ev.Outcome)).Msg("unknown postprocessing outcome - aborting") - fallthrough - case events.PPOutcomeAbort: - failed = true - revertNodeMetadata = true - keepUpload = true - metrics.UploadSessionsAborted.Inc() - case events.PPOutcomeContinue: - if err := session.Finalize(ctx); err != nil { - sublog.Error().Err(err).Msg("could not finalize upload") - failed = true - revertNodeMetadata = false - keepUpload = true - // keep postprocessing status so the upload is not deleted during housekeeping - unmarkPostprocessing = false - } else { - metrics.UploadSessionsFinalized.Inc() - } - case events.PPOutcomeDelete: - failed = true - revertNodeMetadata = true - metrics.UploadSessionsDeleted.Inc() - } - - getParent := func() *node.Node { - p, err := n.Parent(ctx) - if err != nil { - sublog.Error().Err(err).Msg("could not read parent") - return nil - } - return p - } - - now := time.Now() - if failed { - // if no other upload session is in progress (processing id != session id) or has finished (processing id == "") - latestSession, err := n.ProcessingID(ctx) - if err != nil { - sublog.Error().Err(err).Msg("reading node for session failed") - } - if latestSession == session.ID() { - // propagate reverted sizeDiff after failed postprocessing - if err := fs.tp.Propagate(ctx, n, -session.SizeDiff()); err != nil { - sublog.Error().Err(err).Msg("could not propagate tree size change") - } - } - } else if p := getParent(); p != nil { - // update parent tmtime to propagate etag change after successful postprocessing - _ = p.SetTMTime(ctx, &now) - if err := fs.tp.Propagate(ctx, p, 0); err != nil { - sublog.Error().Err(err).Msg("could not propagate etag change") - } - } - - session.Cleanup(revertNodeMetadata, !keepUpload, !keepUpload, unmarkPostprocessing) - - var isVersion bool - if session.NodeExists() { - info, err := session.GetInfo(ctx) - if err == nil && info.MetaData["versionsPath"] != "" { - isVersion = true - } - } - - if err := events.Publish( - ctx, - fs.stream, - events.UploadReady{ - UploadID: ev.UploadID, - Failed: failed, - ExecutingUser: ev.ExecutingUser, - Filename: ev.Filename, - FileRef: &provider.Reference{ - ResourceId: &provider.ResourceId{ - StorageId: session.ProviderID(), - SpaceId: session.SpaceID(), - OpaqueId: session.SpaceID(), - }, - Path: utils.MakeRelativePath(filepath.Join(session.Dir(), session.Filename())), - }, - ResourceID: &provider.ResourceId{ - StorageId: session.ProviderID(), - SpaceId: session.SpaceID(), - OpaqueId: session.NodeID(), - }, - Timestamp: utils.TimeToTS(now), - SpaceOwner: n.SpaceOwnerOrManager(ctx), - IsVersion: isVersion, - ImpersonatingUser: ev.ImpersonatingUser, - }, - ); err != nil { - sublog.Error().Err(err).Msg("Failed to publish UploadReady event") - } - case events.RestartPostprocessing: - sublog := log.With().Str("event", "RestartPostprocessing").Str("uploadid", ev.UploadID).Logger() - session, err := fs.sessionStore.Get(ctx, ev.UploadID) - if err != nil { - sublog.Error().Err(err).Msg("Failed to get upload") - return - } - n, err := session.Node(ctx) - if err != nil { - sublog.Error().Err(err).Msg("could not read node") - return - } - sublog = log.With().Str("spaceid", session.SpaceID()).Str("nodeid", session.NodeID()).Logger() - s, err := session.URL(ctx) - if err != nil { - sublog.Error().Err(err).Msg("could not create url") - return - } - - metrics.UploadSessionsRestarted.Inc() - - // restart postprocessing - if err := events.Publish(ctx, fs.stream, events.BytesReceived{ - UploadID: session.ID(), - URL: s, - SpaceOwner: n.SpaceOwnerOrManager(ctx), - ExecutingUser: &user.User{Id: &user.UserId{OpaqueId: "postprocessing-restart"}}, // send nil instead? - ResourceID: &provider.ResourceId{SpaceId: n.SpaceID, OpaqueId: n.ID}, - Filename: session.Filename(), - Filesize: uint64(session.Size()), - }); err != nil { - sublog.Error().Err(err).Msg("Failed to publish BytesReceived event") - } - case events.CleanUpload: - sublog := log.With().Str("event", "CleanUpload").Str("uploadid", ev.UploadID).Logger() - session, err := fs.sessionStore.Get(ctx, ev.UploadID) - if err != nil { - sublog.Error().Err(err).Msg("Failed to get upload") - return // NOTE: since we can't get the upload, we can't delete the blob - } - session.Cleanup(true, !ev.KeepUpload, !ev.KeepUpload, true) case events.RevertRevision: sublog := log.With().Str("event", "RevertRevision").Interface("nodeid", ev.ResourceID).Logger() if ev.ResourceID != nil && ev.ResourceID.GetStorageId() != "" && ev.ResourceID.GetStorageId() != fs.o.MountID { @@ -489,132 +312,6 @@ func (fs *Decomposedfs) processEvent(evCtx context.Context, event events.Event, sublog.Error().Err(err).Msg("Failed to revert revision") return } - case events.PostprocessingStepFinished: - sublog := log.With().Str("event", "PostprocessingStepFinished").Str("uploadid", ev.UploadID).Logger() - if ev.ResourceID != nil && ev.ResourceID.GetStorageId() != "" && ev.ResourceID.GetStorageId() != fs.o.MountID { - sublog.Debug().Msg("ignoring event for different storage") - return - } - if ev.FinishedStep != events.PPStepAntivirus { - // atm we are only interested in antivirus results - return - } - - res := ev.Result.(events.VirusscanResult) - if res.ErrorMsg != "" { - // scan failed somehow - // Should we handle this here? - return - } - sublog = log.With().Str("scan_description", res.Description).Bool("infected", res.Infected).Logger() - - var n *node.Node - switch ev.UploadID { - case "": - // uploadid is empty -> this was an on-demand scan - /* ON DEMAND SCANNING NOT SUPPORTED ATM - ctx := ctxpkg.ContextSetUser(context.Background(), ev.ExecutingUser) - ref := &provider.Reference{ResourceId: ev.ResourceID} - - no, err := fs.lu.NodeFromResource(ctx, ref) - if err != nil { - log.Error().Err(err).Interface("resourceID", ev.ResourceID).Msg("Failed to get node after scan") - continue - - } - n = no - if ev.Outcome == events.PPOutcomeDelete { - // antivir wants us to delete the file. We must obey and need to - - // check if there a previous versions existing - revs, err := fs.ListRevisions(ctx, ref) - if len(revs) == 0 { - if err != nil { - log.Error().Err(err).Interface("resourceID", ev.ResourceID).Msg("Failed to list revisions. Fallback to delete file") - } - - // no versions -> trash file - err := fs.Delete(ctx, ref) - if err != nil { - log.Error().Err(err).Interface("resourceID", ev.ResourceID).Msg("Failed to delete infected resource") - continue - } - - // now purge it from the recycle bin - if err := fs.PurgeRecycleItem(ctx, &provider.Reference{ResourceId: &provider.ResourceId{SpaceId: n.SpaceID, OpaqueId: n.SpaceID}}, n.ID, "/"); err != nil { - log.Error().Err(err).Interface("resourceID", ev.ResourceID).Msg("Failed to purge infected resource from trash") - } - - // remove cache entry in gateway - fs.cache.RemoveStatContext(ctx, ev.ExecutingUser.GetId(), &provider.ResourceId{SpaceId: n.SpaceID, OpaqueId: n.ID}) - continue - } - - // we have versions - find the newest - versions := make(map[uint64]string) // remember all versions - we need them later - var nv uint64 - for _, v := range revs { - versions[v.Mtime] = v.Key - if v.Mtime > nv { - nv = v.Mtime - } - } - - // restore newest version - if err := fs.RestoreRevision(ctx, ref, versions[nv]); err != nil { - log.Error().Err(err).Interface("resourceID", ev.ResourceID).Str("revision", versions[nv]).Msg("Failed to restore revision") - continue - } - - // now find infected version - revs, err = fs.ListRevisions(ctx, ref) - if err != nil { - log.Error().Err(err).Interface("resourceID", ev.ResourceID).Msg("Error listing revisions after restore") - } - - for _, v := range revs { - // we looking for a version that was previously not there - if _, ok := versions[v.Mtime]; ok { - continue - } - - if err := fs.DeleteRevision(ctx, ref, v.Key); err != nil { - log.Error().Err(err).Interface("resourceID", ev.ResourceID).Str("revision", v.Key).Msg("Failed to delete revision") - } - } - - // remove cache entry in gateway - fs.cache.RemoveStatContext(ctx, ev.ExecutingUser.GetId(), &provider.ResourceId{SpaceId: n.SpaceID, OpaqueId: n.ID}) - continue - } - */ - default: - // uploadid is not empty -> this is an async upload - session, err := fs.sessionStore.Get(ctx, ev.UploadID) - if err != nil { - sublog.Error().Err(err).Msg("Failed to get upload") - return - } - - n, err = session.Node(ctx) - if err != nil { - sublog.Error().Err(err).Msg("Failed to get node after scan") - return - } - sublog = log.With().Str("spaceid", session.SpaceID()).Str("nodeid", session.NodeID()).Logger() - - session.SetScanData(res.Description, res.Scandate) - if err := session.Persist(ctx); err != nil { - sublog.Error().Err(err).Msg("Failed to persist scan results") - } - } - - if err := n.SetScanData(ctx, res.Description, res.Scandate); err != nil { - sublog.Error().Err(err).Msg("Failed to set scan results") - return - } - - metrics.UploadSessionsScanned.Inc() default: log.Error().Interface("event", ev).Msg("Unknown event") } From 6b608d11adeee3c258654bb948c6c2545da48889 Mon Sep 17 00:00:00 2001 From: "lars.jurgensen" Date: Wed, 19 Aug 2026 13:48:54 +0200 Subject: [PATCH 5/8] feat: handle missing node id in tus --- internal/http/services/owncloud/ocdav/tus.go | 3 ++- pkg/rhttp/datatx/manager/tus/tus.go | 4 ++++ 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/internal/http/services/owncloud/ocdav/tus.go b/internal/http/services/owncloud/ocdav/tus.go index 06ad4f71d24..ec8bbaa545c 100644 --- a/internal/http/services/owncloud/ocdav/tus.go +++ b/internal/http/services/owncloud/ocdav/tus.go @@ -319,7 +319,8 @@ func (s *svc) handleTusPost(ctx context.Context, w http.ResponseWriter, r *http. sReq.Ref.Path = uReq.Ref.GetPath() sReq.Ref.ResourceId = nil } else { - if resid, err := storagespace.ParseID(httpRes.Header.Get(net.HeaderOCFileID)); err == nil { + // new files have no node id yet; keep the path-based ref instead + if resid, err := storagespace.ParseID(httpRes.Header.Get(net.HeaderOCFileID)); err == nil && resid.GetOpaqueId() != "" { sReq.Ref = &provider.Reference{ ResourceId: &resid, } diff --git a/pkg/rhttp/datatx/manager/tus/tus.go b/pkg/rhttp/datatx/manager/tus/tus.go index d846525289e..aa07c80cc5d 100644 --- a/pkg/rhttp/datatx/manager/tus/tus.go +++ b/pkg/rhttp/datatx/manager/tus/tus.go @@ -212,6 +212,10 @@ func setHeaders(coord upload.Coordinator, w http.ResponseWriter, r *http.Request if expires != "" { w.Header().Set(net.HeaderTusUploadExpires, expires) } + // the node id is only valid once the upload commits; skip the header for new files + if info.Storage["NodeExists"] != "true" { + return + } resourceid := &provider.ResourceId{ StorageId: info.MetaData["providerID"], SpaceId: info.Storage["SpaceRoot"], From 61dc4ab617592fc247d9b3a785410a02271d3256 Mon Sep 17 00:00:00 2001 From: "lars.jurgensen" Date: Thu, 20 Aug 2026 11:55:22 +0200 Subject: [PATCH 6/8] feat: add changelog --- .../unreleased/feat-upload-coordinator.md | 49 +++++++++++++++++++ internal/http/services/owncloud/ocdav/tus.go | 2 +- pkg/rhttp/datatx/manager/simple/simple.go | 2 +- 3 files changed, 51 insertions(+), 2 deletions(-) create mode 100644 changelog/unreleased/feat-upload-coordinator.md diff --git a/changelog/unreleased/feat-upload-coordinator.md b/changelog/unreleased/feat-upload-coordinator.md new file mode 100644 index 00000000000..841fd93ad82 --- /dev/null +++ b/changelog/unreleased/feat-upload-coordinator.md @@ -0,0 +1,49 @@ +Enhancement: Extract the upload state machine into a driver-agnostic coordinator + +The upload state machine (TUS session management, postprocessing event loop, +antivirus integration, and restart safety) has been extracted from decomposedfs +into a new coordinator in `pkg/upload`. Every storage driver now inherits TUS +chunked uploads, postprocessing, and AV scanning without reimplementing any of +it. + +Drivers integrate by implementing four new methods on the storage interface: + +- `MarkProcessing` sets or clears a "processing" flag on a resource so readers + see a grayed-out placeholder while bytes are in flight. Drivers that do not + need concurrent-upload protection may implement this as a no-op. +- `PrepareUpload` is called after all bytes are received and before + postprocessing begins. Decomposedfs uses this to lock the node, snapshot the + previous version, and propagate the optimistic size change. Drivers with no + such requirements may return immediately. +- `CommitUpload` writes the staged bytes to the resource and receives + pre-computed checksums. +- `RollbackUpload` is the inverse of `PrepareUpload` and is called when + postprocessing fails or is aborted. Drivers that returned immediately from + `PrepareUpload` may return nil. The `RollbackInfo` struct carries the node + identity from the upload session rather than from live node metadata, so a + rollback can still release the quota of a node whose metadata has become + unreadable (e.g. because an ancestor was trashed mid-upload). + +The coordinator owns the upload session files for the decomposedfs driver at +the same on-disk location as before (`/uploads/`), so existing in-flight +uploads continue without interruption and no migration is required. + +**Configuration:** + +Both storageprovider and dataprovider gain an `upload_directory` config key that +sets the local directory where temporary upload session files and staged bytes are stored. +For decomposedfs this is optional; the coordinator falls back to `/uploads/` +inside the driver's own root directory. For drivers that have no local filesystem +root, `upload_directory` must be set explicitly; otherwise the service fails to start. + +The postprocessing consumer settings (`asyncfileuploads`, `consumer_group`, +`numconsumers`, `mount_id`) are read from the driver's own config block, the +same keys decomposedfs already uses. No new top-level config is introduced. + +https://github.com/owncloud/reva/pull/702 +https://github.com/owncloud/reva/pull/703 +https://github.com/owncloud/reva/pull/714 +https://github.com/owncloud/reva/pull/715 +https://github.com/owncloud/reva/pull/717 +https://github.com/owncloud/reva/pull/720 +https://github.com/owncloud/reva/pull/721 diff --git a/internal/http/services/owncloud/ocdav/tus.go b/internal/http/services/owncloud/ocdav/tus.go index ec8bbaa545c..2661b59097c 100644 --- a/internal/http/services/owncloud/ocdav/tus.go +++ b/internal/http/services/owncloud/ocdav/tus.go @@ -320,7 +320,7 @@ func (s *svc) handleTusPost(ctx context.Context, w http.ResponseWriter, r *http. sReq.Ref.ResourceId = nil } else { // new files have no node id yet; keep the path-based ref instead - if resid, err := storagespace.ParseID(httpRes.Header.Get(net.HeaderOCFileID)); err == nil && resid.GetOpaqueId() != "" { + if resid, err := storagespace.ParseID(httpRes.Header.Get(net.HeaderOCFileID)); err == nil && resid.GetOpaqueId() != "" { sReq.Ref = &provider.Reference{ ResourceId: &resid, } diff --git a/pkg/rhttp/datatx/manager/simple/simple.go b/pkg/rhttp/datatx/manager/simple/simple.go index b7141e8035f..057c0d90bd5 100644 --- a/pkg/rhttp/datatx/manager/simple/simple.go +++ b/pkg/rhttp/datatx/manager/simple/simple.go @@ -24,8 +24,8 @@ import ( userpb "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1" provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" - ctxpkg "github.com/owncloud/reva/v2/pkg/ctx" "github.com/mitchellh/mapstructure" + ctxpkg "github.com/owncloud/reva/v2/pkg/ctx" "github.com/pkg/errors" "github.com/rs/zerolog" From 3dc914a14b05b899eda20302db9f7de38fa56f2d Mon Sep 17 00:00:00 2001 From: "lars.jurgensen" Date: Thu, 20 Aug 2026 15:42:11 +0200 Subject: [PATCH 7/8] feat: make TouchFile child-link creation atomic --- pkg/storage/utils/decomposedfs/tree/tree.go | 20 ++++++++------------ 1 file changed, 8 insertions(+), 12 deletions(-) diff --git a/pkg/storage/utils/decomposedfs/tree/tree.go b/pkg/storage/utils/decomposedfs/tree/tree.go index 1b2c397a0e3..3a87fa6db11 100644 --- a/pkg/storage/utils/decomposedfs/tree/tree.go +++ b/pkg/storage/utils/decomposedfs/tree/tree.go @@ -167,20 +167,16 @@ func (t *Tree) TouchFile(ctx context.Context, n *node.Node, markprocessing bool, return err } - // link child name to parent if it is new + // link child name to parent, create-only so a concurrent upload of the same + // new file cannot silently clobber it: the loser gets AlreadyExists and its + // CAS loop re-reads and retries instead of losing the update (OCISDEV-855) childNameLink := filepath.Join(n.ParentPath(), n.Name) - var link string - link, err = os.Readlink(childNameLink) - if err == nil && link != "../"+n.ID { - if err = os.Remove(childNameLink); err != nil { - return errors.Wrap(err, "Decomposedfs: could not remove symlink child entry") - } - } - if errors.Is(err, fs.ErrNotExist) || link != "../"+n.ID { - relativeNodePath := filepath.Join("../../../../../", lookup.Pathify(n.ID, 4, 2)) - if err = os.Symlink(relativeNodePath, childNameLink); err != nil { - return errors.Wrap(err, "Decomposedfs: could not symlink child entry") + relativeNodePath := filepath.Join("../../../../../", lookup.Pathify(n.ID, 4, 2)) + if err = os.Symlink(relativeNodePath, childNameLink); err != nil { + if errors.Is(err, fs.ErrExist) { + return errtypes.AlreadyExists(n.Name) } + return errors.Wrap(err, "Decomposedfs: could not symlink child entry") } return t.Propagate(ctx, n, 0) From 79bb0520701d4774672cdb43a1ebc41ebe08dcbc Mon Sep 17 00:00:00 2001 From: "lars.jurgensen" Date: Fri, 21 Aug 2026 16:43:57 +0200 Subject: [PATCH 8/8] feat: add NewCoordinatorFromConfig --- .../storageprovider/storageprovider.go | 21 +++---------------- .../services/dataprovider/dataprovider.go | 18 +++------------- pkg/upload/coordinator.go | 19 +++++++++++++++++ 3 files changed, 25 insertions(+), 33 deletions(-) diff --git a/internal/grpc/services/storageprovider/storageprovider.go b/internal/grpc/services/storageprovider/storageprovider.go index 39692a1188c..3c69dd3e3ee 100644 --- a/internal/grpc/services/storageprovider/storageprovider.go +++ b/internal/grpc/services/storageprovider/storageprovider.go @@ -212,9 +212,10 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc. return nil, err } - coord, err := getCoordinator(c, fs, evstream, log) + // storageprovider only initiates uploads; the data path assembles chunks, so no chunking here. + coord, err := upload.NewCoordinatorFromConfig(c.UploadDirectory, c.Drivers[c.Driver], fs, evstream, log, false) if err != nil { - return nil, err + return nil, fmt.Errorf("storageprovider: %w", err) } service := &Service{ @@ -228,22 +229,6 @@ func New(m map[string]interface{}, ss *grpc.Server, log *zerolog.Logger) (rgrpc. return service, nil } -// getCoordinator builds the coordinator that initiates uploads for the driver -// this service mounts. It stages sessions in the same directory the dataprovider -// appends bytes to, so an upload initiated here can be continued there. -func getCoordinator(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (upload.Coordinator, error) { - store := upload.NewFileStoreFromConfig(c.UploadDirectory, c.Drivers[c.Driver], log) - if store == nil { - return nil, fmt.Errorf("storageprovider: cannot determine the upload directory, set upload_directory") - } - if err := store.Setup(); err != nil { - return nil, fmt.Errorf("storageprovider: upload directory setup failed: %w", err) - } - - // No chunk folder: only the data path assembles chunks. - return upload.NewCoordinator(fs, store, "", publisher), nil -} - func (s *Service) SetArbitraryMetadata(ctx context.Context, req *provider.SetArbitraryMetadataRequest) (*provider.SetArbitraryMetadataResponse, error) { ctx = ctxpkg.ContextSetLockID(ctx, req.LockId) diff --git a/internal/http/services/dataprovider/dataprovider.go b/internal/http/services/dataprovider/dataprovider.go index b138306e8fe..6b537a6a009 100644 --- a/internal/http/services/dataprovider/dataprovider.go +++ b/internal/http/services/dataprovider/dataprovider.go @@ -106,9 +106,10 @@ func New(m map[string]interface{}, log *zerolog.Logger) (global.Service, error) return nil, err } - coord, err := getCoordinator(conf, fs, evstream, log) + // the data path assembles chunks, so enable chunking + coord, err := upload.NewCoordinatorFromConfig(conf.UploadDirectory, conf.Drivers[conf.Driver], fs, evstream, log, true) if err != nil { - return nil, err + return nil, fmt.Errorf("dataprovider: %w", err) } // only the data path consumes postprocessing results: one consumer group gets @@ -141,19 +142,6 @@ func getFS(c *config, stream events.Stream, log *zerolog.Logger) (storage.FS, er return nil, fmt.Errorf("driver not found: %s", c.Driver) } -// getCoordinator builds the coordinator that owns the upload lifecycle for the -// driver this service mounts. -func getCoordinator(c *config, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (upload.Coordinator, error) { - store := upload.NewFileStoreFromConfig(c.UploadDirectory, c.Drivers[c.Driver], log) - if store == nil { - return nil, fmt.Errorf("dataprovider: cannot determine the upload directory, set upload_directory") - } - if err := store.Setup(); err != nil { - return nil, fmt.Errorf("dataprovider: upload directory setup failed: %w", err) - } - return upload.NewCoordinator(fs, store, store.UploadDir(), publisher), nil -} - func getDataTXs(c *config, coord upload.Coordinator, fs storage.FS, publisher events.Publisher, log *zerolog.Logger) (map[string]http.Handler, error) { if c.DataTXs == nil { c.DataTXs = make(map[string]map[string]interface{}) diff --git a/pkg/upload/coordinator.go b/pkg/upload/coordinator.go index 801856218ee..a0b004e954e 100644 --- a/pkg/upload/coordinator.go +++ b/pkg/upload/coordinator.go @@ -11,6 +11,7 @@ import ( user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1" provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" "github.com/google/uuid" + "github.com/rs/zerolog" tusd "github.com/tus/tusd/v2/pkg/handler" "github.com/owncloud/reva/v2/pkg/appctx" @@ -62,6 +63,24 @@ func NewCoordinator(fs storage.FS, store SessionStore, chunkFolder string, pub e return c } +// NewCoordinatorFromConfig sets up a coordinator and its file store. Pass +// withChunking=false when only initiating uploads: chunk assembly happens on the +// data path, so the storageprovider doesn't need it. +func NewCoordinatorFromConfig(uploadDir string, driverConf map[string]interface{}, fs storage.FS, pub events.Publisher, log *zerolog.Logger, withChunking bool) (Coordinator, error) { + store := NewFileStoreFromConfig(uploadDir, driverConf, log) + if store == nil { + return nil, fmt.Errorf("cannot determine the upload directory, set upload_directory") + } + if err := store.Setup(); err != nil { + return nil, fmt.Errorf("upload directory setup failed: %w", err) + } + chunkFolder := "" + if withChunking { + chunkFolder = store.UploadDir() + } + return NewCoordinator(fs, store, chunkFolder, pub), nil +} + // InitiateUpload resolves the target, then creates and persists the session that // bytes are appended to. func (c *coordinator) InitiateUpload(ctx context.Context, ref *provider.Reference, uploadLength int64, metadata map[string]string) (map[string]string, error) {