From 2a05417d7a2b3147e9372d0b89b31d4bd0241013 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?J=C3=B6rn=20Friedrich=20Dreyer?= Date: Mon, 5 Oct 2026 13:06:09 +0200 Subject: [PATCH 1/4] chore: bump reva to 2.51.0 --- go.mod | 2 +- go.sum | 4 +- .../services/gateway/storageprovidercache.go | 13 +- .../reva/v2/pkg/events/postprocessing.go | 17 ++ .../storage/pkg/decomposedfs/decomposedfs.go | 26 +++ .../pkg/storage/pkg/decomposedfs/node/node.go | 174 +++++++++++++++++- .../pkg/storage/pkg/decomposedfs/revisions.go | 23 ++- .../pkg/decomposedfs/tree/revisions.go | 56 ------ .../storage/pkg/decomposedfs/upload/upload.go | 80 +++----- .../utils/decomposedfs/decomposedfs.go | 26 +++ .../storage/utils/decomposedfs/node/node.go | 168 ++++++++++++++++- .../storage/utils/decomposedfs/revisions.go | 23 ++- .../utils/decomposedfs/upload/upload.go | 76 +++----- vendor/modules.txt | 2 +- 14 files changed, 522 insertions(+), 168 deletions(-) diff --git a/go.mod b/go.mod index 271b78b785..44cfe9ee4b 100644 --- a/go.mod +++ b/go.mod @@ -64,7 +64,7 @@ require ( github.com/open-policy-agent/opa v1.19.1 github.com/opencloud-eu/icap-client v0.0.0-20250930132611-28a2afe62d89 github.com/opencloud-eu/libre-graph-api-go v1.0.8 - github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002054620-6fb86feb3aa2 + github.com/opencloud-eu/reva/v2 v2.51.0 github.com/opensearch-project/opensearch-go/v4 v4.7.3 github.com/orcaman/concurrent-map v1.0.0 github.com/pkg/errors v0.9.1 diff --git a/go.sum b/go.sum index bf6466df2f..c30cad60d3 100644 --- a/go.sum +++ b/go.sum @@ -761,8 +761,8 @@ github.com/opencloud-eu/icap-client v0.0.0-20250930132611-28a2afe62d89 h1:W1ms+l github.com/opencloud-eu/icap-client v0.0.0-20250930132611-28a2afe62d89/go.mod h1:vigJkNss1N2QEceCuNw/ullDehncuJNFB6mEnzfq9UI= github.com/opencloud-eu/libre-graph-api-go v1.0.8 h1:VTkFLkLNna+T4RN6JUZnVQ1bqVAfPtMeORXbFzEYGKs= github.com/opencloud-eu/libre-graph-api-go v1.0.8/go.mod h1:lTM8JeGblNpoMySTW7Lui2+c5TTLI95mwxtdUIHHrhU= -github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002054620-6fb86feb3aa2 h1:ROTURL0i5+PSIjHB3nW59oroSfLM/bFcJEZq9Hb9GWA= -github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002054620-6fb86feb3aa2/go.mod h1:ngLsDakKnt8zCFCeULhuSZFESGy5NIuSq5/JGeOhPdU= +github.com/opencloud-eu/reva/v2 v2.51.0 h1:PbvUtLlbCpS5em48cSGk2sHRdjVA7htRQVKjpckI2w4= +github.com/opencloud-eu/reva/v2 v2.51.0/go.mod h1:ngLsDakKnt8zCFCeULhuSZFESGy5NIuSq5/JGeOhPdU= github.com/opencloud-eu/secure v0.0.0-20260312082735-b6f5cb2244e4 h1:l2oB/RctH+t8r7QBj5p8thfEHCM/jF35aAY3WQ3hADI= github.com/opencloud-eu/secure v0.0.0-20260312082735-b6f5cb2244e4/go.mod h1:BmF5hyM6tXczk3MpQkFf1hpKSRqCyhqcbiQtiAF7+40= github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= diff --git a/vendor/github.com/opencloud-eu/reva/v2/internal/grpc/services/gateway/storageprovidercache.go b/vendor/github.com/opencloud-eu/reva/v2/internal/grpc/services/gateway/storageprovidercache.go index f427e1f3c0..91efba06d6 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/internal/grpc/services/gateway/storageprovidercache.go +++ b/vendor/github.com/opencloud-eu/reva/v2/internal/grpc/services/gateway/storageprovidercache.go @@ -27,6 +27,7 @@ import ( ctxpkg "github.com/opencloud-eu/reva/v2/pkg/ctx" sdk "github.com/opencloud-eu/reva/v2/pkg/sdk/common" "github.com/opencloud-eu/reva/v2/pkg/storage/cache" + "github.com/opencloud-eu/reva/v2/pkg/storagespace" "github.com/opencloud-eu/reva/v2/pkg/utils" "github.com/pkg/errors" "google.golang.org/grpc" @@ -129,7 +130,17 @@ func (c *cachedSpacesAPIClient) UpdateStorageSpace(ctx context.Context, in *prov return c.c.UpdateStorageSpace(ctx, in, opts...) } func (c *cachedSpacesAPIClient) DeleteStorageSpace(ctx context.Context, in *provider.DeleteStorageSpaceRequest, opts ...grpc.CallOption) (*provider.DeleteStorageSpaceResponse, error) { - return c.c.DeleteStorageSpace(ctx, in, opts...) + resp, err := c.c.DeleteStorageSpace(ctx, in, opts...) + switch { + case err != nil: + return nil, err + case resp.Status.Code != rpc.Code_CODE_OK: + return resp, nil + default: + _, spaceid, _, _ := storagespace.SplitID(in.GetId().GetOpaqueId()) + _ = c.createPersonalSpaceCache.Delete(spaceid) + return resp, nil + } } /* diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/events/postprocessing.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/events/postprocessing.go index 48e29e5cd5..b2694dd524 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/events/postprocessing.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/events/postprocessing.go @@ -238,3 +238,20 @@ func (CleanUpload) Unmarshal(v []byte) (interface{}, error) { err := json.Unmarshal(v, &e) return e, err } + +// DeleteRevision can be emitted to delete a revision of a node. If Timestamp +// is set exactly that revision is deleted. If Timestamp is nil the node's +// current version is targeted and the node reverts to the pre-upload state, +// unstalling it - only while the node is actually stuck in processing, so a +// redelivered event is a no-op. +type DeleteRevision struct { + ResourceID *provider.ResourceId + Timestamp *types.Timestamp +} + +// Unmarshal to fulfill umarshaller interface +func (DeleteRevision) Unmarshal(v []byte) (interface{}, error) { + e := DeleteRevision{} + err := json.Unmarshal(v, &e) + return e, err +} diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/decomposedfs.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/decomposedfs.go index ae56d70a5a..78c8afa7b9 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/decomposedfs.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/decomposedfs.go @@ -83,6 +83,7 @@ var ( events.RestartPostprocessing{}, events.StartPostprocessingStep{}, events.CleanUpload{}, + events.DeleteRevision{}, } ) @@ -464,6 +465,31 @@ func (fs *Decomposedfs) handlePostprocessingEvent(ctx context.Context, event eve 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.DeleteRevision: + sublog := log.With().Str("event", "DeleteRevision").Interface("nodeid", ev.ResourceID).Logger() + n, err := fs.lu.NodeFromID(ctx, ev.ResourceID) + if err != nil { + sublog.Error().Err(err).Msg("Failed to get node") + return + } + + var deleteErr error + if ev.Timestamp == nil { + // the node's current revision is targeted - revert to the pre-upload + // state. Only stuck nodes are targeted: a reverted node has its + // processing flag removed, so a redelivered event is a no-op. + if !n.IsProcessing(ctx) { + sublog.Debug().Msg("node is not stuck, ignoring") + return + } + _, deleteErr = n.DeleteRevision(ctx, "") + } else { + versionID := time.Unix(int64(ev.Timestamp.Seconds), int64(ev.Timestamp.Nanos)).UTC().Format(time.RFC3339Nano) + deleteErr = fs.deleteRevisionFile(ctx, n, n.ID+node.RevisionIDDelimiter+versionID) + } + if deleteErr != nil { + sublog.Error().Err(deleteErr).Msg("Failed to delete revision") + } case events.StartPostprocessingStep: sublog := log.With().Str("event", "StartPostprocessingStep").Str("uploadid", ev.UploadID).Logger() if ev.UploadID == "" { diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node/node.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node/node.go index 33516a6a64..6aa8f68827 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node/node.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node/node.go @@ -1324,7 +1324,12 @@ func (n *Node) DeleteGrant(ctx context.Context, g *provider.Grant) (err error) { // Purge removes a node from disk. It does not move it to the trash func (n *Node) Purge(ctx context.Context) error { - return n.lu.PurgeNode(n) + if err := n.lu.PurgeNode(n); err != nil { + return err + } + + // remove .mpk and .mlock files + return n.lu.MetadataBackend().Purge(ctx, n) } // ListGrants lists all grants of the current node. @@ -1551,6 +1556,173 @@ func (n *Node) SetDTime(ctx context.Context, t *time.Time) (err error) { return n.lu.TimeManager().SetDTime(ctx, n, t) } +// RevertUpload reverts the upload that created the given revision: it +// restores the revision onto the node and removes the revision, including its +// metadata sidecars. The node's metadata lock is held for the duration of the +// operation. +func (n *Node) RevertUpload(ctx context.Context, versionID string) error { + if versionID == "" { + return errors.New("empty versionID") + } + + unlock, err := n.lu.MetadataBackend().Lock(n) + if err != nil { + return err + } + defer func() { + _ = unlock() + }() + + revisionNode := NewBaseNode(n.SpaceID, n.ID+RevisionIDDelimiter+versionID, n.lu) + revisionPath := revisionNode.InternalPath() + if _, err := os.Stat(revisionPath); err != nil { + appctx.GetLogger(ctx).Error().Str("versionpath", revisionPath).Err(err).Msg("revision does not exist") + return err + } + + if err := n.lu.CopyMetadata(ctx, revisionNode, n, func(attributeName string, value []byte) (newValue []byte, copy bool) { + return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) || + attributeName == prefixes.TypeAttr || + attributeName == prefixes.BlobIDAttr || + attributeName == prefixes.BlobsizeAttr || + attributeName == prefixes.MTimeAttr + }); err != nil { + appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("restoring revision metadata failed") + return err + } + + if err := os.RemoveAll(revisionPath); err != nil { + appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("error removing version") + return err + } + + // remove the revision's metadata sidecars + if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionNode)); err != nil { + appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata, continuing") + } + if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionNode)); err != nil { + appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata lockfile, continuing") + } + if err := n.lu.MetadataBackend().Purge(ctx, revisionNode); err != nil { + appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not purge revision from cache, continuing") + } + return nil +} + +// DeleteRevision deletes a revision of the node and returns the blob id of the +// deleted revision, so that the caller can delete the blob. If versionID is +// empty the node's current revision is deleted: the latest stored revision is +// restored onto the node (or the node is purged when there is none) and the +// processing flag is removed, unstalling the node. No blob is returned in that +// case. +func (n *Node) DeleteRevision(ctx context.Context, versionID string) (string, error) { + if versionID != "" { + return n.deleteRevision(ctx, n.ID+RevisionIDDelimiter+versionID) + } + + revisionPath, err := n.getLatestRevision(ctx) + if err != nil { + return "", err + } + + if revisionPath == "" { + // there is no revision - delete the node + unlock, err := n.lu.MetadataBackend().Lock(n) + if err != nil { + return "", err + } + defer func() { + _ = unlock() + }() + if err := n.Purge(ctx); err != nil { + appctx.GetLogger(ctx).Info().Str("nodepath", n.InternalPath()).Err(err).Msg("error purging node") + return "", err + } + return "", nil + } + + latestID := strings.TrimPrefix(revisionPath, n.lu.VersionPath(n.SpaceID, n.ID, "")) + if err := n.RevertUpload(ctx, latestID); err != nil { + return "", err + } + + // we just deleted the current revision - remove processing flag if set + if uploadid, err := n.ProcessingID(ctx); err == nil { + return "", n.UnmarkProcessing(ctx, uploadid) + } + return "", nil +} + +// deleteRevision deletes the revision node identified by revisionID (the node's +// ID + RevisionIDDelimiter + timestamp) together with its metadata sidecars +// and returns its blob id. It is a no-op if the revision does not exist. +func (n *Node) deleteRevision(ctx context.Context, revisionID string) (string, error) { + log := appctx.GetLogger(ctx) + + revisionNode := NewBaseNode(n.SpaceID, revisionID, n.lu) + revisionPath := revisionNode.InternalPath() + if _, err := os.Stat(revisionPath); err != nil { + log.Warn().Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("revision does not exist, nothing to delete") + return "", nil + } + + unlock, err := n.lu.MetadataBackend().Lock(n) + if err != nil { + return "", err + } + defer func() { + _ = unlock() + }() + + blobID, _, err := n.lu.ReadBlobIDAndSizeAttr(ctx, revisionNode, nil) + if err != nil { + return "", err + } + + if err := os.RemoveAll(revisionPath); err != nil { + return "", err + } + + if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionNode)); err != nil { + log.Warn().Err(err).Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("could not delete revision metadata, continuing") + } + if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionNode)); err != nil { + log.Warn().Err(err).Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("could not delete revision metadata lockfile, continuing") + } + if err := n.lu.MetadataBackend().Purge(ctx, revisionNode); err != nil { + log.Warn().Err(err).Str("nodeid", n.ID).Str("revisionid", revisionID).Msg("could not purge revision from cache, continuing") + } + + return blobID, nil +} + +func (n *Node) getLatestRevision(ctx context.Context) (string, error) { + revPrefix := n.lu.VersionPath(n.SpaceID, n.ID, "") + revisions, err := filepath.Glob(n.lu.VersionPath(n.SpaceID, n.ID, "*")) + if err != nil { + appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error reading revisions") + return "", err + } + + revPath, latest := "", time.Time{} + for _, rev := range revisions { + if strings.HasSuffix(rev, ".mpk") || strings.HasSuffix(rev, ".mlock") { + continue + } + revDate, err := time.Parse(time.RFC3339Nano, strings.TrimPrefix(rev, revPrefix)) + if err != nil { + appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error parsing revision date") + continue + } + + if revDate.After(latest) { + latest = revDate + revPath = rev + } + } + return revPath, nil +} + // ReadChildNodeFromLink reads the child node id from a link func ReadChildNodeFromLink(ctx context.Context, path string) (string, error) { _, span := tracer.Start(ctx, "readChildNodeFromLink") diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/revisions.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/revisions.go index 7b7e458444..c0c0791660 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/revisions.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/revisions.go @@ -162,11 +162,28 @@ func (fs *Decomposedfs) DeleteRevision(ctx context.Context, ref *provider.Refere return err } - if err := os.RemoveAll(fs.lu.InternalPath(n.SpaceID, revisionKey)); err != nil { - return err + return fs.deleteRevisionFile(ctx, n, revisionKey) +} + +// deleteRevisionFile deletes the revision node identified by revisionKey +// (nodeID + RevisionIDDelimiter + timestamp) together with its metadata +// sidecars and its own blob. It is a no-op if the revision does not exist. +func (fs *Decomposedfs) deleteRevisionFile(ctx context.Context, n *node.Node, revisionKey string) error { + kp := strings.SplitN(revisionKey, node.RevisionIDDelimiter, 2) + if len(kp) != 2 { + return errtypes.NotFound(revisionKey) } - return fs.tp.DeleteBlob(n) + blobID, err := n.DeleteRevision(ctx, kp[1]) + if err != nil { + return err + } + if blobID == "" { + // no blob to delete (0-byte file or current-revision revert) + return nil + } + + return fs.tp.DeleteBlob(&node.Node{BaseNode: node.BaseNode{SpaceID: n.SpaceID}, BlobID: blobID}) } func (fs *Decomposedfs) getRevisionNode(ctx context.Context, ref *provider.Reference, revisionKey string, hasPermission func(*provider.ResourcePermissions) bool) (*node.Node, error) { diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/tree/revisions.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/tree/revisions.go index 09de801621..33f731a7a8 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/tree/revisions.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/tree/revisions.go @@ -275,62 +275,6 @@ func (tp *Tree) DownloadRevision(ctx context.Context, ref *provider.Reference, r return ri, reader, nil } -// DeleteRevision deletes the specified revision of the resource -func (tp *Tree) DeleteRevision(ctx context.Context, ref *provider.Reference, revisionKey string) error { - _, span := tracer.Start(ctx, "DeleteRevision") - defer span.End() - n, err := tp.getRevisionNode(ctx, ref, revisionKey, func(rp *provider.ResourcePermissions) bool { - return rp.RestoreFileVersion - }) - if err != nil { - return err - } - - if err := os.RemoveAll(tp.lookup.InternalPath(n.SpaceID, revisionKey)); err != nil { - return err - } - - return tp.DeleteBlob(n) -} - -func (tp *Tree) getRevisionNode(ctx context.Context, ref *provider.Reference, revisionKey string, hasPermission func(*provider.ResourcePermissions) bool) (*node.Node, error) { - _, span := tracer.Start(ctx, "getRevisionNode") - defer span.End() - log := appctx.GetLogger(ctx) - - // verify revision key format - kp := strings.SplitN(revisionKey, node.RevisionIDDelimiter, 2) - if len(kp) != 2 { - log.Error().Str("revisionKey", revisionKey).Msg("malformed revisionKey") - return nil, errtypes.NotFound(revisionKey) - } - log.Debug().Str("revisionKey", revisionKey).Msg("DownloadRevision") - - spaceID := ref.ResourceId.SpaceId - // check if the node is available and has not been deleted - n, err := node.ReadNode(ctx, tp.lookup, spaceID, kp[0], "", false, nil, false) - if err != nil { - return nil, err - } - if !n.Exists { - err = errtypes.NotFound(filepath.Join(n.ParentID, n.Name)) - return nil, err - } - - p, err := tp.permissions.AssemblePermissions(ctx, n) - switch { - case err != nil: - return nil, err - case !hasPermission(p): - return nil, errtypes.PermissionDenied(filepath.Join(n.ParentID, n.Name)) - } - - // Set space owner in context - storagespace.ContextSendSpaceOwnerID(ctx, n.SpaceOwnerOrManager(ctx)) - - return n, nil -} - func (tp *Tree) RestoreRevision(ctx context.Context, sourceNode, targetNode metadata.MetadataNode, mtime time.Time) error { err := tp.lookup.CopyMetadata(ctx, sourceNode, targetNode, func(attributeName string, value []byte) (newValue []byte, copy bool) { return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) || diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/upload/upload.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/upload/upload.go index a3e7b79868..e67cfcd031 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/upload/upload.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/upload/upload.go @@ -353,22 +353,29 @@ func (session *DecomposedFsSession) Finalize(ctx context.Context) (err error) { // another upload on this node is in progress or has finished since we started if !isProcessing || processingID != session.ID() { - versionID := n.ID + node.RevisionIDDelimiter + session.MTime().UTC().Format(time.RFC3339Nano) - // There should be a revision node (created by the other upload that finished before us), read it and upload our blob there. - existingRevisionNode, revisionNodeUnlock, err := node.LockAndReadNode(ctx, session.store.lu, session.SpaceID(), versionID, "", false, spaceRoot, false) - if err != nil || !existingRevisionNode.Exists { - // The revision node has not been created. Likely because the file on disk was modified externally and re-assilimated (watchfs == true) - // Let's create the revision node now and upload the blob to it. - n, revisionNodeUnlock, err = session.createRevisionNodeForUpload(ctx, n, session.MTime().UTC().Format(time.RFC3339Nano)) - if err != nil { - appctx.GetLogger(ctx).Debug().Err(err).Str("versionID", session.MTime().UTC().Format(time.RFC3339Nano)).Msg("failed to create revision node for upload finalization") - return err - } + // The node's current content is already this upload's content (e.g. a + // newer upload was reverted in the meantime) - no revision is needed, + // upload the blob to the node itself. + if attribs.String(prefixes.BlobIDAttr) == session.ID() { + appctx.GetLogger(ctx).Debug().Str("nodepath", n.InternalPath()).Msg("node already contains this upload's content, no revision needed") } else { - n = existingRevisionNode + versionID := n.ID + node.RevisionIDDelimiter + session.MTime().UTC().Format(time.RFC3339Nano) + // There should be a revision node (created by the other upload that finished before us), read it and upload our blob there. + existingRevisionNode, revisionNodeUnlock, err := node.LockAndReadNode(ctx, session.store.lu, session.SpaceID(), versionID, "", false, spaceRoot, false) + if err != nil || !existingRevisionNode.Exists { + // The revision node has not been created. Likely because the file on disk was modified externally and re-assilimated (watchfs == true) + // Let's create the revision node now and upload the blob to it. + n, revisionNodeUnlock, err = session.createRevisionNodeForUpload(ctx, n, session.MTime().UTC().Format(time.RFC3339Nano)) + if err != nil { + appctx.GetLogger(ctx).Debug().Err(err).Str("versionID", session.MTime().UTC().Format(time.RFC3339Nano)).Msg("failed to create revision node for upload finalization") + return err + } + } else { + n = existingRevisionNode + } + appctx.GetLogger(ctx).Debug().Str("new nodepath", n.InternalPath()).Msg("uploading to revision node, that was created for us by another upload") + defer func() { _ = revisionNodeUnlock() }() } - appctx.GetLogger(ctx).Debug().Str("new nodepath", n.InternalPath()).Msg("uploading to revision node, that was created for us by another upload") - defer func() { _ = revisionNodeUnlock() }() } // upload the data to the blobstore @@ -433,17 +440,6 @@ func checkHash(expected string, h hash.Hash) error { return nil } -func (session *DecomposedFsSession) removeNode(ctx context.Context) { - n, err := session.Node(ctx) - if err != nil { - appctx.GetLogger(ctx).Error().Str("session", session.ID()).Err(err).Msg("getting node from session failed") - return - } - if err := n.Purge(ctx); err != nil { - appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("purging node failed") - } -} - // cleanup cleans up after the upload is finished func (session *DecomposedFsSession) Cleanup(revertNodeMetadata, cleanBin, cleanInfo, unmarkPostprocessing bool) { ctx := session.Context(context.Background()) @@ -455,35 +451,13 @@ func (session *DecomposedFsSession) Cleanup(revertNodeMetadata, cleanBin, cleanI if err != nil { sublog.Error().Err(err).Msg("reading node for session failed") } else { - if session.NodeExists() && session.info.MetaData["versionID"] != "" { - versionID := session.info.MetaData["versionID"] + versionID := strings.TrimPrefix(session.info.MetaData["versionID"], n.ID+node.RevisionIDDelimiter) + if session.NodeExists() && versionID != "" { sublog.Debug().Str("nodepath", n.InternalPath()).Str("versionID", versionID).Msg("restoring revision") - revisionNode, err := node.ReadNode(ctx, session.store.lu, session.SpaceID(), versionID, "", false, n.SpaceRoot, false) - if err != nil { - sublog.Error().Err(err).Str("versionID", versionID).Msg("reading revision node failed") + if err := n.RevertUpload(ctx, versionID); err != nil { + sublog.Error().Err(err).Str("versionID", versionID).Msg("reverting node metadata failed") return } - - if !revisionNode.Exists { - sublog.Error().Str("versionID", versionID).Msg("revision node does not exist") - return - } - - // restore the revision - mtime, err := revisionNode.GetMTime(ctx) - if err != nil { - sublog.Error().Err(err).Str("versionID", versionID).Msg("getting mtime of revision node failed") - mtime = time.Now() - } - - if err := session.store.tp.RestoreRevision(ctx, revisionNode, n, mtime); err != nil { - sublog.Error().Err(err).Str("versionID", versionID).Msg("restoring revision node failed") - return - } - - if err := os.RemoveAll(revisionNode.InternalPath()); err != nil { - sublog.Error().Err(err).Str("revisionpath", revisionNode.InternalPath()).Msg("removing restored revision file failed") - } } else { // if no other upload session is in progress (processing id != session id) or has finished (processing id == "") latestSession, err := n.ProcessingID(ctx) @@ -492,7 +466,9 @@ func (session *DecomposedFsSession) Cleanup(revertNodeMetadata, cleanBin, cleanI } if latestSession == session.ID() { // actually delete the node - session.removeNode(ctx) + if err := n.Purge(ctx); err != nil { + sublog.Error().Err(err).Str("nodepath", n.InternalPath()).Msg("purging node failed") + } } // FIXME else if the upload has become a revision, delete the revision, or if it is the last one, delete the node } diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/decomposedfs.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/decomposedfs.go index 13c345e347..fd79d52f23 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/decomposedfs.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/decomposedfs.go @@ -83,6 +83,7 @@ var ( events.PostprocessingStepFinished{}, events.RestartPostprocessing{}, events.CleanUpload{}, + events.DeleteRevision{}, } ) @@ -452,6 +453,31 @@ func (fs *Decomposedfs) handlePostprocessingEvent(ctx context.Context, event eve 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.DeleteRevision: + sublog := log.With().Str("event", "DeleteRevision").Interface("nodeid", ev.ResourceID).Logger() + n, err := fs.lu.NodeFromID(ctx, ev.ResourceID) + if err != nil { + sublog.Error().Err(err).Msg("Failed to get node") + return + } + + var deleteErr error + if ev.Timestamp == nil { + // the node's current revision is targeted - revert to the pre-upload + // state. Only stuck nodes are targeted: a reverted node has its + // processing flag removed, so a redelivered event is a no-op. + if !n.IsProcessing(ctx) { + sublog.Debug().Msg("node is not stuck, ignoring") + return + } + _, deleteErr = n.DeleteRevision(ctx, "") + } else { + versionID := time.Unix(int64(ev.Timestamp.Seconds), int64(ev.Timestamp.Nanos)).UTC().Format(time.RFC3339Nano) + deleteErr = fs.deleteRevisionFile(ctx, n, n.ID+node.RevisionIDDelimiter+versionID) + } + if deleteErr != nil { + sublog.Error().Err(deleteErr).Msg("Failed to delete revision") + } case events.PostprocessingStepFinished: sublog := log.With().Str("event", "PostprocessingStepFinished").Str("uploadid", ev.UploadID).Logger() if ev.FinishedStep != events.PPStepAntivirus { diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/node/node.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/node/node.go index fdb12817ff..d3f866a8c4 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/node/node.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/node/node.go @@ -1165,7 +1165,12 @@ func (n *Node) Purge(ctx context.Context) error { // remove child entry in parent src := filepath.Join(n.ParentPath(), n.Name) - return os.Remove(src) + if err := os.Remove(src); err != nil { + return err + } + + // remove .mpk and .mlock files + return n.lu.MetadataBackend().Purge(ctx, n.InternalPath()) } // ListGrants lists all grants of the current node. @@ -1391,3 +1396,164 @@ func (n *Node) GetDTime(ctx context.Context) (time.Time, error) { func (n *Node) SetDTime(ctx context.Context, t *time.Time) (err error) { return n.lu.TimeManager().SetDTime(ctx, n, t) } + +// RevertUpload reverts the upload that created the given revision: it +// restores the revision onto the node and removes the revision, including its +// metadata sidecars. The node's metadata lock is held for the duration of the +// operation. +func (n *Node) RevertUpload(ctx context.Context, versionID string) error { + if versionID == "" { + return errors.New("empty versionID") + } + + lock, err := lockedfile.OpenFile(n.lu.MetadataBackend().LockfilePath(n.InternalPath()), os.O_CREATE|os.O_WRONLY, 0600) + if err != nil { + return err + } + defer lock.Close() + + revisionPath := n.InternalPath() + RevisionIDDelimiter + versionID + if _, err := os.Stat(revisionPath); err != nil { + appctx.GetLogger(ctx).Error().Str("versionpath", revisionPath).Err(err).Msg("revision does not exist") + return err + } + + if err := n.lu.CopyMetadata(ctx, revisionPath, n.InternalPath(), func(attributeName string, value []byte) (newValue []byte, copy bool) { + return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) || + attributeName == prefixes.TypeAttr || + attributeName == prefixes.BlobIDAttr || + attributeName == prefixes.BlobsizeAttr || + attributeName == prefixes.MTimeAttr + }, false); err != nil { + appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("restoring revision metadata failed") + return err + } + + if err := os.RemoveAll(revisionPath); err != nil { + appctx.GetLogger(ctx).Info().Str("versionpath", revisionPath).Str("nodepath", n.InternalPath()).Err(err).Msg("error removing version") + return err + } + + // remove the revision's metadata sidecars + if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionPath)); err != nil { + appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata, continuing") + } + if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionPath)); err != nil { + appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not delete revision metadata lockfile, continuing") + } + if err := n.lu.MetadataBackend().Purge(ctx, revisionPath); err != nil { + appctx.GetLogger(ctx).Warn().Err(err).Str("versionpath", revisionPath).Msg("could not purge revision from cache, continuing") + } + return nil +} + +// DeleteRevision deletes a revision of the node and returns the blob id of the +// deleted revision, so that the caller can delete the blob. If versionID is +// empty the node's current revision is deleted: the latest stored revision is +// restored onto the node (or the node is purged when there is none) and the +// processing flag is removed, unstalling the node. No blob is returned in that +// case. +func (n *Node) DeleteRevision(ctx context.Context, versionID string) (string, error) { + if versionID != "" { + return n.deleteRevision(ctx, versionID) + } + + revisionPath, err := n.getLatestRevision(ctx) + if err != nil { + return "", err + } + + if revisionPath == "" { + // there is no revision - delete the node + lock, err := lockedfile.OpenFile(n.lu.MetadataBackend().LockfilePath(n.InternalPath()), os.O_CREATE|os.O_WRONLY, 0600) + if err != nil { + return "", err + } + defer lock.Close() + if err := n.Purge(ctx); err != nil { + appctx.GetLogger(ctx).Info().Str("nodepath", n.InternalPath()).Err(err).Msg("error purging node") + return "", err + } + return "", nil + } + + latestID := strings.TrimPrefix(revisionPath, n.InternalPath()+RevisionIDDelimiter) + if err := n.RevertUpload(ctx, latestID); err != nil { + return "", err + } + + // we just deleted the current revision - remove processing flag if set + if uploadid, err := n.ProcessingID(ctx); err == nil { + return "", n.UnmarkProcessing(ctx, uploadid) + } + return "", nil +} + +// deleteRevision deletes the revision identified by versionID (the timestamp +// following the node's InternalPath + RevisionIDDelimiter) together with its +// metadata sidecars and returns its blob id. It is a no-op if the revision +// does not exist. +func (n *Node) deleteRevision(ctx context.Context, versionID string) (string, error) { + log := appctx.GetLogger(ctx) + + revisionPath := n.InternalPath() + RevisionIDDelimiter + versionID + if _, err := os.Stat(revisionPath); err != nil { + log.Warn().Str("nodeid", n.ID).Str("versionid", versionID).Msg("revision does not exist, nothing to delete") + return "", nil + } + + lock, err := lockedfile.OpenFile(n.lu.MetadataBackend().LockfilePath(n.InternalPath()), os.O_CREATE|os.O_WRONLY, 0600) + if err != nil { + return "", err + } + defer lock.Close() + + blobID, _, err := n.lu.ReadBlobIDAndSizeAttr(ctx, revisionPath, nil) + if err != nil { + return "", err + } + + if err := os.RemoveAll(revisionPath); err != nil { + return "", err + } + + if err := os.Remove(n.lu.MetadataBackend().MetadataPath(revisionPath)); err != nil { + log.Warn().Err(err).Str("nodeid", n.ID).Str("versionid", versionID).Msg("could not delete revision metadata, continuing") + } + if err := os.Remove(n.lu.MetadataBackend().LockfilePath(revisionPath)); err != nil { + log.Warn().Err(err).Str("nodeid", n.ID).Str("versionid", versionID).Msg("could not delete revision metadata lockfile, continuing") + } + if err := n.lu.MetadataBackend().Purge(ctx, revisionPath); err != nil { + log.Warn().Err(err).Str("nodeid", n.ID).Str("versionid", versionID).Msg("could not purge revision from cache, continuing") + } + + return blobID, nil +} + +func (n *Node) getLatestRevision(ctx context.Context) (string, error) { + revPrefix := n.InternalPath() + RevisionIDDelimiter + revisions, err := filepath.Glob(revPrefix + "*") + if err != nil { + appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error reading revisions") + return "", err + } + + revPath, latest := "", time.Time{} + for _, rev := range revisions { + if strings.HasSuffix(rev, ".mpk") || strings.HasSuffix(rev, ".mlock") { + continue + } + revDate, err := time.Parse(time.RFC3339Nano, strings.TrimPrefix(rev, revPrefix)) + if err != nil { + appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("error parsing revision date") + continue + } + + appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Str("revPath", revPath).Interface("time", revDate).Err(err).Msg("error parsing revision date") + if revDate.After(latest) { + latest = revDate + revPath = rev + } + } + return revPath, nil +} diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/revisions.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/revisions.go index 5833f23763..b351623be8 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/revisions.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/revisions.go @@ -345,11 +345,28 @@ func (fs *Decomposedfs) DeleteRevision(ctx context.Context, ref *provider.Refere return err } - if err := os.RemoveAll(fs.lu.InternalPath(n.SpaceID, revisionKey)); err != nil { - return err + return fs.deleteRevisionFile(ctx, n, revisionKey) +} + +// deleteRevisionFile deletes the revision node identified by revisionKey +// (nodeID + RevisionIDDelimiter + timestamp) together with its metadata +// sidecars and its own blob. It is a no-op if the revision does not exist. +func (fs *Decomposedfs) deleteRevisionFile(ctx context.Context, n *node.Node, revisionKey string) error { + kp := strings.SplitN(revisionKey, node.RevisionIDDelimiter, 2) + if len(kp) != 2 { + return errtypes.NotFound(revisionKey) } - return fs.tp.DeleteBlob(n) + blobID, err := n.DeleteRevision(ctx, kp[1]) + if err != nil { + return err + } + if blobID == "" { + // no blob to delete (0-byte file or current-revision revert) + return nil + } + + return fs.tp.DeleteBlob(&node.Node{SpaceID: n.SpaceID, BlobID: blobID}) } func (fs *Decomposedfs) getRevisionNode(ctx context.Context, ref *provider.Reference, revisionKey string, hasPermission func(*provider.ResourcePermissions) bool) (*node.Node, error) { diff --git a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/upload/upload.go b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/upload/upload.go index 86c8676581..1087a2c100 100644 --- a/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/upload/upload.go +++ b/vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/upload/upload.go @@ -316,57 +316,10 @@ func checkHash(expected string, h hash.Hash) error { return nil } -func (session *OcisSession) removeNode(ctx context.Context) { - n, err := session.Node(ctx) - if err != nil { - appctx.GetLogger(ctx).Error().Str("session", session.ID()).Err(err).Msg("getting node from session failed") - return - } - if err := n.Purge(ctx); err != nil { - appctx.GetLogger(ctx).Error().Str("nodepath", n.InternalPath()).Err(err).Msg("purging node failed") - } -} - // cleanup cleans up after the upload is finished func (session *OcisSession) Cleanup(revertNodeMetadata, cleanBin, cleanInfo, unmarkPostprocessing bool) { ctx := session.Context(context.Background()) - if revertNodeMetadata { - n, err := session.Node(ctx) - if err != nil { - appctx.GetLogger(ctx).Error().Err(err).Str("sessionid", session.ID()).Msg("reading node for session failed") - } else { - if session.NodeExists() && session.info.MetaData["versionsPath"] != "" { - p := session.info.MetaData["versionsPath"] - if err := session.store.lu.CopyMetadata(ctx, p, n.InternalPath(), func(attributeName string, value []byte) (newValue []byte, copy bool) { - return value, strings.HasPrefix(attributeName, prefixes.ChecksumPrefix) || - attributeName == prefixes.TypeAttr || - attributeName == prefixes.BlobIDAttr || - attributeName == prefixes.BlobsizeAttr || - attributeName == prefixes.MTimeAttr - }, true); err != nil { - appctx.GetLogger(ctx).Info().Str("versionpath", p).Str("nodepath", n.InternalPath()).Err(err).Msg("renaming version node failed") - } - - if err := os.RemoveAll(p); err != nil { - appctx.GetLogger(ctx).Info().Str("versionpath", p).Str("nodepath", n.InternalPath()).Err(err).Msg("error removing version") - } - - } else { - // 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 { - appctx.GetLogger(ctx).Error().Err(err).Str("spaceid", n.SpaceID).Str("nodeid", n.ID).Str("uploadid", session.ID()).Msg("reading processingid for session failed") - } - if latestSession == session.ID() { - // actually delete the node - session.removeNode(ctx) - } - // FIXME else if the upload has become a revision, delete the revision, or if it is the last one, delete the node - } - } - } - if cleanBin { if err := os.Remove(session.binPath()); err != nil && !errors.Is(err, fs.ErrNotExist) { appctx.GetLogger(ctx).Error().Str("path", session.binPath()).Err(err).Msg("removing upload failed") @@ -376,8 +329,37 @@ func (session *OcisSession) Cleanup(revertNodeMetadata, cleanBin, cleanInfo, unm if cleanInfo { if err := os.Remove(session.infoPath()); err != nil { appctx.GetLogger(ctx).Error().Err(err).Str("session", session.ID()).Msg("removing upload info failed") + } + } + + if revertNodeMetadata { + n, err := session.Node(ctx) + if err != nil { + appctx.GetLogger(ctx).Error().Err(err).Str("sessionid", session.ID()).Msg("reading node for session failed") return } + + versionID := strings.TrimPrefix(session.info.MetaData["versionsPath"], n.InternalPath()+node.RevisionIDDelimiter) + if session.NodeExists() && versionID != "" { + if err := n.RevertUpload(ctx, versionID); err != nil { + appctx.GetLogger(ctx).Error().Err(err).Str("nodepath", n.InternalPath()).Msg("reverting node metadata failed") + return + } + } else { + // 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 { + appctx.GetLogger(ctx).Error().Err(err).Str("spaceid", n.SpaceID).Str("nodeid", n.ID).Str("uploadid", session.ID()).Msg("reading processingid for session failed") + } + if latestSession == session.ID() { + // actually delete the node + if err := n.Purge(ctx); err != nil { + appctx.GetLogger(ctx).Error().Err(err).Str("nodepath", n.InternalPath()).Msg("purging node failed") + return + } + } + // FIXME else if the upload has become a revision, delete the revision, or if it is the last one, delete the node + } } if unmarkPostprocessing { diff --git a/vendor/modules.txt b/vendor/modules.txt index bd32502f4e..909c2e43fd 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -1362,7 +1362,7 @@ github.com/opencloud-eu/icap-client # github.com/opencloud-eu/libre-graph-api-go v1.0.8 ## explicit; go 1.23 github.com/opencloud-eu/libre-graph-api-go -# github.com/opencloud-eu/reva/v2 v2.50.1-0.20261002054620-6fb86feb3aa2 +# github.com/opencloud-eu/reva/v2 v2.51.0 ## explicit; go 1.26.0 github.com/opencloud-eu/reva/v2/cmd/revad/internal/grace github.com/opencloud-eu/reva/v2/cmd/revad/runtime From ad6ff3ee8603c3d1fb7e6d9996cbd776e5376d2d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?J=C3=B6rn=20Friedrich=20Dreyer?= Date: Thu, 24 Sep 2026 13:51:31 +0200 Subject: [PATCH 2/4] fix(storage-users): add delete-stale-nodes command MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Jörn Friedrich Dreyer --- services/storage-users/pkg/command/uploads.go | 232 ++++++++++++++++++ 1 file changed, 232 insertions(+) diff --git a/services/storage-users/pkg/command/uploads.go b/services/storage-users/pkg/command/uploads.go index 3c320d925f..ce789cc662 100644 --- a/services/storage-users/pkg/command/uploads.go +++ b/services/storage-users/pkg/command/uploads.go @@ -4,16 +4,21 @@ import ( "context" "encoding/json" "fmt" + "io/fs" + "log" "os" + "path/filepath" "strconv" "strings" "time" "github.com/olekukonko/tablewriter" "github.com/olekukonko/tablewriter/tw" + "github.com/shamaton/msgpack/v2" "github.com/spf13/cobra" userpb "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1" + provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1" "github.com/opencloud-eu/opencloud/pkg/config/configlog" "github.com/opencloud-eu/opencloud/services/storage-users/pkg/config" "github.com/opencloud-eu/opencloud/services/storage-users/pkg/config/parser" @@ -22,9 +27,19 @@ import ( "github.com/opencloud-eu/reva/v2/pkg/events" "github.com/opencloud-eu/reva/v2/pkg/storage" "github.com/opencloud-eu/reva/v2/pkg/storage/fs/registry" + "github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/lookup" + "github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node" "github.com/opencloud-eu/reva/v2/pkg/utils" ) +const ( + MSGPACK_KEY_USER_OCIS_NODESTATUS = "user.ocis.nodestatus" + + // Log indentation levels + LOG_INDENT_L1 = " " // 2 spaces + LOG_INDENT_L2 = LOG_INDENT_L1 + LOG_INDENT_L1 +) + // Session contains the information of an upload session type Session struct { ID string `json:"id"` @@ -50,6 +65,7 @@ func Uploads(cfg *config.Config) *cobra.Command { } uploadsCmd.AddCommand([]*cobra.Command{ ListUploadSessions(cfg), + DeleteStaleProcessingNodes(cfg), }...) return uploadsCmd @@ -318,3 +334,219 @@ func buildInfo(filter storage.UploadSessionFilter) string { b.WriteString(":") return b.String() } + +// DeleteStaleProcessingNodes is the entry point for the delete-stale-nodes command +func DeleteStaleProcessingNodes(cfg *config.Config) *cobra.Command { + deleteStaleNodesCmd := &cobra.Command{ + Use: "delete-stale-nodes", + Short: "Delete all nodes in processing state that are not referenced by any upload session", + PreRunE: func(cmd *cobra.Command, args []string) error { + return configlog.ReturnFatal(parser.ParseConfig(cfg)) + }, + RunE: func(cmd *cobra.Command, args []string) error { + spaceIDs := []string{} + dryRun, _ := cmd.Flags().GetBool("dry-run") + verbose, _ := cmd.Flags().GetBool("verbose") + start := time.Now() + + // Check if specific space ID provided + if cmd.Flags().Changed("spaceid") { + spaceID, _ := cmd.Flags().GetString("spaceid") + spaceIDs = append(spaceIDs, spaceID) + } else { + fmt.Println("Scanning all spaces for stale processing nodes...") + spaceIDs = globSpaceIDs(cfg) + } + + if verbose { + fmt.Printf("Spaces to cleanup: %d\n", len(spaceIDs)) + for _, spaceID := range spaceIDs { + fmt.Printf(" - %s\n", spaceID) + } + } + + var stream events.Stream + if !dryRun { + s, err := event.NewStream(cfg) + if err != nil { + log.Fatalf("Failed to create event stream: %v", err) + } + stream = s + } + + staleCount := 0 + for _, spaceID := range spaceIDs { + staleCount += deleteStaleUploads(cfg, spaceID, dryRun, verbose, stream) + } + + if verbose { + fmt.Printf("Took %ds\n", int(time.Since(start).Seconds())) + } + fmt.Printf("Total stale nodes: %d\n", staleCount) + + return nil + }, + } + deleteStaleNodesCmd.Flags().String("spaceid", "", "Space ID to check for processing nodes (omit to check all spaces)") + deleteStaleNodesCmd.Flags().Bool("dry-run", true, "Only show what would be deleted without actually deleting") + deleteStaleNodesCmd.Flags().Bool("verbose", false, "Enable verbose logging") + return deleteStaleNodesCmd +} + +// globSpaceIDs returns a list of all space IDs in the storage root +func globSpaceIDs(cfg *config.Config) []string { + fsys := os.DirFS(cfg.Drivers.Decomposed.Root) + dirs, err := fs.Glob(fsys, "spaces/*/*/nodes") + if err != nil { + fmt.Fprintf(os.Stderr, "Error globbing spaces root directory %s: %v\n", cfg.Drivers.Decomposed.Root, err) + return []string{} + } + + spaceIDs := []string{} + for _, dir := range dirs { + // For dir i.e. spaces/9d/408cec-8f0a-4d33-8715-89df1217a10c/nodes + // spaceID is 9d408cec-8f0a-4d33-8715-89df1217a10c + spaceIDs = append(spaceIDs, strings.ReplaceAll(strings.TrimSuffix(strings.TrimPrefix(dir, "spaces/"), "/nodes"), "/", "")) + } + return spaceIDs +} + +// delete stale processing nodes for a given spaceID +func deleteStaleUploads(cfg *config.Config, spaceID string, dryRun bool, verbose bool, stream events.Stream) int { + if verbose { + fmt.Printf("\nDeleting stale processing nodes for space: %s\n", spaceID) + } + + // Find .mpk files in space directory + spaceRoot := filepath.Join(cfg.Drivers.Decomposed.Root, "spaces", lookup.Pathify(spaceID, 1, 2)) + mpkFiles := []string{} + err := filepath.Walk(spaceRoot, func(path string, info os.FileInfo, err error) error { + if err != nil { + fmt.Fprintf(os.Stderr, "Error accessing path %s: %s\n", path, err) + return filepath.SkipDir + } + if !info.IsDir() && strings.HasSuffix(path, ".mpk") { + mpkFiles = append(mpkFiles, path) + } + return nil + }) + + if err != nil { + fmt.Fprintf(os.Stderr, "Error walking space directory %s: %s\n", spaceRoot, err) + return 0 + } + + if verbose { + fmt.Printf("%sFound total %d .mpk files\n", LOG_INDENT_L1, len(mpkFiles)) + } + + staleCount := 0 + for _, path := range mpkFiles { + staleCount += deleteStaleNode(cfg, path, dryRun, verbose, stream) + } + + if verbose { + fmt.Printf("%sFound total %d stale nodes\n", LOG_INDENT_L1, staleCount) + } + + return staleCount +} + +// deleteStaleNode deletes a stale node: if it is not referenced by any upload session +// returns 1 if the node stale node was detected for deletion, 0 otherwise, for counting purposes +func deleteStaleNode(cfg *config.Config, path string, dryRun bool, verbose bool, stream events.Stream) int { + nodeDir := filepath.Dir(path) + + // Read .mpk file to get processing info + b, err := os.ReadFile(path) + if err != nil { + fmt.Fprintf(os.Stderr, "Error reading file %s: %s\n", path, err) + return 0 + } + var mpkData map[string]any + if err := msgpack.Unmarshal(b, &mpkData); err != nil { + fmt.Fprintf(os.Stderr, "Error unmarshaling file %s: %s\n", path, err) + return 0 + } + + processingID := extractProcessingID(mpkData) + if processingID == "" { + return 0 + } + + // Construct path to upload info file: + // i.e. ~/.ocis/storage/users/uploads/5329c14b-b786-4b27-8f7d-7429f03009d7.info + // And pass only the .info file not exists: err is ErrNotExist + pathUploadInfo := filepath.Join(cfg.Drivers.Decomposed.Root, "uploads", processingID) + ".info" + _, infoStatErr := os.Stat(pathUploadInfo) + if infoStatErr == nil { + return 0 + } + if !os.IsNotExist(infoStatErr) { + // Tere was an error other than file not existing, log and return + fmt.Fprintf(os.Stderr, "Error checking upload info %s: %s\n", pathUploadInfo, infoStatErr) + return 0 + } + + if verbose { + fmt.Printf("%sFound stale upload at %s (Processing ID: %s)\n", LOG_INDENT_L1, path, processingID) + fmt.Printf("%sUpload info missing at: %s\n", LOG_INDENT_L2, pathUploadInfo) + } + + if dryRun { + return 1 + } + + rid := extractResourceID(strings.TrimSuffix(path, ".mpk")) + if rid == nil { + fmt.Fprintf(os.Stderr, "Failed to extract resource ID from path %s\n", path) + return 0 + } + + // A nil Timestamp targets the node's current revision: the driver reverts the + // node to its previous version, or purges it if there is none, and clears the + // processing flag. A non-nil timestamp would instead delete the revision with + // that exact timestamp, which never matches the one stuck in processing. + if err := events.Publish(context.Background(), stream, events.DeleteRevision{ + ResourceID: rid, + }); err != nil { + // if publishing fails there is no need to try publishing other events - they will fail too. + log.Fatalf("Failed to send delete revision event for node '%s'\n", path) + } + + if verbose { + fmt.Printf("%sDeleted stale node: %s\n", LOG_INDENT_L2, nodeDir) + } + + return 1 +} + +func extractProcessingID(mpkData map[string]any) string { + processingID := "" + for k, v := range mpkData { + vStr := string(v.([]byte)) + if k == MSGPACK_KEY_USER_OCIS_NODESTATUS && strings.Contains(vStr, node.ProcessingStatus) { + processingID = strings.Split(vStr, ":")[1] + break + } + } + return processingID +} + +func extractResourceID(path string) *provider.ResourceId { + // path looks like /.../storage/users/spaces/f2/06bccf-0f10-4070-9e63-40943f060667/nodes/5b/ba/1e/a7/-f185-4f31-8342-ed4b5743f096 + parts := strings.Split(path, "spaces") + if len(parts) < 2 { + return nil + } + + spaceParts := strings.Split(parts[1], "nodes") + if len(spaceParts) < 2 { + return nil + } + + return &provider.ResourceId{ + SpaceId: strings.ReplaceAll(spaceParts[0], "/", ""), + OpaqueId: strings.ReplaceAll(spaceParts[1], "/", ""), + } +} From fabf82850e2b30ab4f720ffec451840ef4a2e28a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?J=C3=B6rn=20Friedrich=20Dreyer?= Date: Thu, 24 Sep 2026 15:13:33 +0200 Subject: [PATCH 3/4] fix(storage-users): use user.oc. metadata prefix for stale node detection MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Jörn Friedrich Dreyer --- services/storage-users/pkg/command/uploads.go | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/services/storage-users/pkg/command/uploads.go b/services/storage-users/pkg/command/uploads.go index ce789cc662..7649294eaf 100644 --- a/services/storage-users/pkg/command/uploads.go +++ b/services/storage-users/pkg/command/uploads.go @@ -28,13 +28,12 @@ import ( "github.com/opencloud-eu/reva/v2/pkg/storage" "github.com/opencloud-eu/reva/v2/pkg/storage/fs/registry" "github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/lookup" + "github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/metadata/prefixes" "github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node" "github.com/opencloud-eu/reva/v2/pkg/utils" ) const ( - MSGPACK_KEY_USER_OCIS_NODESTATUS = "user.ocis.nodestatus" - // Log indentation levels LOG_INDENT_L1 = " " // 2 spaces LOG_INDENT_L2 = LOG_INDENT_L1 + LOG_INDENT_L1 @@ -525,7 +524,7 @@ func extractProcessingID(mpkData map[string]any) string { processingID := "" for k, v := range mpkData { vStr := string(v.([]byte)) - if k == MSGPACK_KEY_USER_OCIS_NODESTATUS && strings.Contains(vStr, node.ProcessingStatus) { + if k == prefixes.StatusPrefix && strings.Contains(vStr, node.ProcessingStatus) { processingID = strings.Split(vStr, ":")[1] break } From 5ff51c1856a8b7fb039c9b6f5e7d7bc4a345834a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?J=C3=B6rn=20Friedrich=20Dreyer?= Date: Thu, 24 Sep 2026 15:33:29 +0200 Subject: [PATCH 4/4] docs(storage-users): document the delete-stale-nodes command MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Jörn Friedrich Dreyer --- services/storage-users/README.md | 47 ++++++++++++++++++- services/storage-users/pkg/command/uploads.go | 2 +- 2 files changed, 47 insertions(+), 2 deletions(-) diff --git a/services/storage-users/README.md b/services/storage-users/README.md index 0289db18f0..ae2e550e8e 100644 --- a/services/storage-users/README.md +++ b/services/storage-users/README.md @@ -55,7 +55,8 @@ opencloud storage-users uploads ```plaintext COMMANDS: - sessions Print a list of upload sessions + delete-stale-nodes Delete (or revert) all nodes in processing state that are not referenced by any upload session + sessions Print a list of upload sessions ``` #### Sessions command @@ -140,6 +141,50 @@ opencloud storage-users uploads sessions --expired=true --clean opencloud storage-users uploads sessions --processing=false --has-virus=false --resume ``` + +#### Delete Stale Nodes command + +This command allows to remove (or revert) nodes that are stale, meaning they are in postprocessing but their upload session is gone. +It will check all nodes that are in postprocessing and find those without an upload session. Then it will delete the node if there are no other versions. If there are other versions, it will instead revert the node to the previous version. +The command reads the decomposed storage directly, so it only finds nodes when the msgpack metadata backend (`.mpk` files) is used. The stale nodes are cleaned up by `DeleteRevision` events, which are consumed by the running `storage-users` service: it needs to be up for the cleanup to actually happen. + +```bash + opencloud storage-users uploads delete-stale-nodes +``` +```plaintext +Delete (or revert) all nodes in processing state that are not referenced by any upload session + +Usage: + opencloud storage-users uploads delete-stale-nodes [flags] + +Flags: + --dry-run Only show what would be deleted without actually deleting (default true) + -h, --help help for delete-stale-nodes + --spaceid string Space ID to check for processing nodes (omit to check all spaces) + --verbose Enable verbose logging +``` + +#### Command Examples + +Dry run to see what would be deleted (recommended first step) + +```bash +opencloud storage-users uploads delete-stale-nodes +``` + +Set `--dry-run=false` to actually delete the stale nodes + +```bash +opencloud storage-users uploads delete-stale-nodes --dry-run=false +``` + +Use `--verbose` to get more information about what is happening + +```bash +opencloud storage-users uploads delete-stale-nodes --dry-run=false --verbose +``` + + ### Manage Trash-Bin Items This command set provides commands to get an overview of trash-bin items, restore items and purge old items of `personal` spaces and `project` spaces (spaces that have been created manually). `trash-bin` commands require a `spaceID` as parameter. diff --git a/services/storage-users/pkg/command/uploads.go b/services/storage-users/pkg/command/uploads.go index 7649294eaf..866466f207 100644 --- a/services/storage-users/pkg/command/uploads.go +++ b/services/storage-users/pkg/command/uploads.go @@ -338,7 +338,7 @@ func buildInfo(filter storage.UploadSessionFilter) string { func DeleteStaleProcessingNodes(cfg *config.Config) *cobra.Command { deleteStaleNodesCmd := &cobra.Command{ Use: "delete-stale-nodes", - Short: "Delete all nodes in processing state that are not referenced by any upload session", + Short: "Delete (or revert) all nodes in processing state that are not referenced by any upload session", PreRunE: func(cmd *cobra.Command, args []string) error { return configlog.ReturnFatal(parser.ParseConfig(cfg)) },