chore: bump reva to 2.51.0

This commit is contained in:
Jörn Friedrich Dreyer committed 2026-10-05 14:12:21 +02:00
1 parent 4dec19089d
commit 2a05417d7a
14 files changed
+522 -168

No files matched your search

+1 -1
View File
@@ -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
+2 -2
View File
@@ -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=
@@ -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
}
}
/*
+17
View File
@@ -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
}
@@ -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 == "" {
@@ -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")
@@ -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) {
@@ -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) ||
@@ -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
}
@@ -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 {
@@ -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
}
@@ -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) {
@@ -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 {
+1 -1
View File
@@ -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