Skip to content

Commit f26b856

Browse files
Ocisdev 900 pr6d orphan rollback (#720)
* fix: release the quota of an orphaned upload through the driver seam * feat: cleanup * feat: apply review comments --------- Co-authored-by: Firas Frikha <firas.frikha@kiteworks.com>
1 parent d37b513 commit f26b856

20 files changed

Lines changed: 263 additions & 38 deletions

File tree

‎pkg/ocm/storage/received/ocm.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -716,6 +716,6 @@ func (d *driver) PrepareUpload(_ context.Context, _ *provider.Reference, _ strin
716716
return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil
717717
}
718718

719-
func (d *driver) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ bool, _ int64) error {
719+
func (d *driver) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error {
720720
return nil
721721
}

‎pkg/storage/fs/cephfs/upload.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,7 @@ func (fs *cephfs) PrepareUpload(_ context.Context, _ *provider.Reference, _ stri
159159
return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil
160160
}
161161

162-
func (fs *cephfs) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ bool, _ int64) error {
162+
func (fs *cephfs) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error {
163163
return nil
164164
}
165165

‎pkg/storage/fs/hello/unimplemented.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -98,7 +98,7 @@ func (fs *hellofs) PrepareUpload(_ context.Context, _ *provider.Reference, _ str
9898
return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil
9999
}
100100

101-
func (fs *hellofs) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ bool, _ int64) error {
101+
func (fs *hellofs) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error {
102102
return nil
103103
}
104104

‎pkg/storage/fs/kiteworks/kiteworks.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -279,7 +279,7 @@ func (d *Driver) PrepareUpload(_ context.Context, _ *provider.Reference, _ strin
279279
return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil
280280
}
281281

282-
func (d *Driver) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ bool, _ int64) error {
282+
func (d *Driver) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error {
283283
return nil
284284
}
285285

‎pkg/storage/fs/nextcloud/nextcloud.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -453,7 +453,7 @@ func (nc *StorageDriver) PrepareUpload(_ context.Context, _ *provider.Reference,
453453
return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil
454454
}
455455

456-
func (nc *StorageDriver) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ bool, _ int64) error {
456+
func (nc *StorageDriver) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error {
457457
return nil
458458
}
459459

‎pkg/storage/fs/owncloudsql/upload.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -181,7 +181,7 @@ func (fs *owncloudsqlfs) PrepareUpload(_ context.Context, _ *provider.Reference,
181181
return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil
182182
}
183183

184-
func (fs *owncloudsqlfs) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ bool, _ int64) error {
184+
func (fs *owncloudsqlfs) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error {
185185
return nil
186186
}
187187

‎pkg/storage/fs/s3/upload.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -86,6 +86,6 @@ func (fs *s3FS) PrepareUpload(_ context.Context, _ *provider.Reference, _ string
8686
return &storage.PrepareUploadResult{VersionCreated: info.NodeExisted}, nil
8787
}
8888

89-
func (fs *s3FS) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ bool, _ int64) error {
89+
func (fs *s3FS) RollbackUpload(_ context.Context, _ *provider.Reference, _ string, _ storage.RollbackInfo) error {
9090
return nil
9191
}

‎pkg/storage/storage.go‎

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,18 @@ type PrepareUploadResult struct {
9090
SizeDiff int64
9191
}
9292

93+
// RollbackInfo carries what a driver needs to undo PrepareUpload. NodeID and
94+
// ParentID come from the upload session rather than the node, so a rollback can
95+
// still release the quota of a node whose own metadata has become unreadable.
96+
type RollbackInfo struct {
97+
NodeExisted bool // true when the target node had a prior version; drivers with nothing to undo for new nodes may no-op
98+
SizeDiff int64
99+
NodeID string
100+
ParentID string
101+
Filename string
102+
Size int64
103+
}
104+
93105
// FS is the interface to implement access to the storage.
94106
type FS interface {
95107
// Minimal set for a readonly storage driver
@@ -149,9 +161,9 @@ type FS interface {
149161
// It is the inverse of PrepareUpload: restores previous metadata and reverts the optimistic
150162
// size propagation. The caller (coordinator) is responsible for unmarking the processing flag
151163
// and deleting the upload session files. Drivers that performed no work in PrepareUpload may return nil.
152-
// nodeExisted indicates whether the target node had a prior version; drivers that have nothing
153-
// to undo for new nodes may no-op when nodeExisted is false.
154-
RollbackUpload(ctx context.Context, ref *provider.Reference, sessionID string, nodeExisted bool, sizeDiff int64) error
164+
// Callers must keep the session on a returned error: the rollback is retryable and info
165+
// carries state the driver cannot recover once the session files are gone.
166+
RollbackUpload(ctx context.Context, ref *provider.Reference, sessionID string, info RollbackInfo) error
155167

156168
// Revisions
157169

‎pkg/storage/utils/decomposedfs/rollback_upload_test.go‎

Lines changed: 115 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
package decomposedfs_test
22

33
import (
4+
"os"
5+
46
provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
57
. "github.com/onsi/ginkgo/v2"
68
. "github.com/onsi/gomega"
@@ -68,7 +70,7 @@ var _ = Describe("RollbackUpload", func() {
6870
Expect(revsBefore).To(HaveLen(1))
6971

7072
Expect(env.Fs.MarkProcessing(env.Ctx, ref, true, "session-overwrite")).To(Succeed())
71-
Expect(env.Fs.RollbackUpload(env.Ctx, ref, "session-overwrite", true, result.SizeDiff)).To(Succeed())
73+
Expect(env.Fs.RollbackUpload(env.Ctx, ref, "session-overwrite", storage.RollbackInfo{NodeExisted: true, SizeDiff: result.SizeDiff})).To(Succeed())
7274

7375
revsAfter, err := env.Fs.ListRevisions(env.Ctx, ref)
7476
Expect(err).ToNot(HaveOccurred())
@@ -104,7 +106,7 @@ var _ = Describe("RollbackUpload", func() {
104106
Expect(err).ToNot(HaveOccurred())
105107

106108
Expect(env.Fs.MarkProcessing(env.Ctx, ref, true, "session-b")).To(Succeed())
107-
Expect(env.Fs.RollbackUpload(env.Ctx, ref, "session-a", true, 10)).To(Succeed())
109+
Expect(env.Fs.RollbackUpload(env.Ctx, ref, "session-a", storage.RollbackInfo{NodeExisted: true, SizeDiff: 10})).To(Succeed())
108110

109111
n, err := env.Lookup.NodeFromResource(env.Ctx, ref)
110112
Expect(err).ToNot(HaveOccurred())
@@ -120,7 +122,117 @@ var _ = Describe("RollbackUpload", func() {
120122
ResourceId: env.SpaceRootRes,
121123
Path: "/dir1/does-not-exist.txt",
122124
}
123-
Expect(env.Fs.RollbackUpload(env.Ctx, missingRef, "any-session-id", false, 0)).To(Succeed())
125+
Expect(env.Fs.RollbackUpload(env.Ctx, missingRef, "any-session-id", storage.RollbackInfo{})).To(Succeed())
126+
})
127+
})
128+
129+
// A node whose metadata is gone but whose file remains, e.g. because an
130+
// ancestor was trashed while the upload was in flight. ReadNode fails on the
131+
// missing parent id, so the rollback has only the session to go on.
132+
Context("orphaned node: the metadata is unreadable", func() {
133+
var (
134+
idRef *provider.Reference
135+
nodePath string
136+
dir1Ref *provider.Reference
137+
info storage.RollbackInfo
138+
)
139+
140+
// parentSize is the size dir1 accounts for, which the propagation moves.
141+
parentSize := func() uint64 {
142+
parentInfo, err := env.Fs.GetMD(env.Ctx, dir1Ref, []string{}, []string{})
143+
Expect(err).ToNot(HaveOccurred())
144+
return parentInfo.Size
145+
}
146+
147+
// orphan prepares an upload, then purges the target's metadata. It returns
148+
// the parent size from before the optimistic propagation.
149+
orphan := func() uint64 {
150+
dir1Ref = &provider.Reference{ResourceId: env.SpaceRootRes, Path: "/dir1"}
151+
sizeBefore := parentSize()
152+
153+
_, err := env.Fs.TouchFile(env.Ctx, ref, false, "")
154+
Expect(err).ToNot(HaveOccurred())
155+
156+
n, err := env.Lookup.NodeFromResource(env.Ctx, ref)
157+
Expect(err).ToNot(HaveOccurred())
158+
nodePath = n.InternalPath()
159+
// Address the node by id: the path walk needs the parent metadata the
160+
// purge below destroys.
161+
idRef = &provider.Reference{ResourceId: &provider.ResourceId{
162+
StorageId: env.SpaceRootRes.GetStorageId(),
163+
SpaceId: env.SpaceRootRes.GetSpaceId(),
164+
OpaqueId: n.ID,
165+
}}
166+
167+
result, err := env.Fs.PrepareUpload(env.Ctx, ref, "session-orphan", storage.UploadInfo{
168+
NodeExisted: false,
169+
Size: 30,
170+
})
171+
Expect(err).ToNot(HaveOccurred())
172+
Expect(env.Fs.MarkProcessing(env.Ctx, ref, true, "session-orphan")).To(Succeed())
173+
Expect(parentSize()).To(Equal(sizeBefore+30), "the optimistic size should be propagated")
174+
175+
// Purge, not remove: the cached attributes have to go too.
176+
Expect(env.Lookup.MetadataBackend().Purge(env.Ctx, nodePath)).To(Succeed())
177+
_, err = env.Lookup.NodeFromResource(env.Ctx, idRef)
178+
Expect(err).To(HaveOccurred(), "the node should be unreadable")
179+
180+
info = storage.RollbackInfo{
181+
SizeDiff: result.SizeDiff,
182+
NodeID: idRef.GetResourceId().GetOpaqueId(),
183+
ParentID: n.ParentID,
184+
Filename: "rollback-target.txt",
185+
Size: 30,
186+
}
187+
return sizeBefore
188+
}
189+
190+
It("releases the quota and removes the unreachable node", func() {
191+
sizeBefore := orphan()
192+
193+
Expect(env.Fs.RollbackUpload(env.Ctx, idRef, "session-orphan", info)).To(Succeed())
194+
195+
Expect(parentSize()).To(Equal(sizeBefore))
196+
197+
_, err := os.Stat(nodePath)
198+
Expect(err).To(HaveOccurred(), "the orphaned node should be gone")
199+
})
200+
201+
// Without them there is no way to reach the parent, so report the failure
202+
// rather than pretending the upload was rolled back.
203+
It("reports the lookup failure when the session carries no ids", func() {
204+
orphan()
205+
206+
err := env.Fs.RollbackUpload(env.Ctx, idRef, "session-orphan", storage.RollbackInfo{SizeDiff: 30})
207+
208+
Expect(err).To(MatchError(ContainSubstring("node lookup failed")))
209+
_, statErr := os.Stat(nodePath)
210+
Expect(statErr).ToNot(HaveOccurred(), "the node should be left for a retry")
211+
})
212+
213+
// The space is what the propagation walks up to, so there is nothing to
214+
// walk without it.
215+
It("reports the lookup failure when the reference names no space", func() {
216+
orphan()
217+
218+
err := env.Fs.RollbackUpload(env.Ctx, &provider.Reference{Path: "/dir1/rollback-target.txt"}, "session-orphan", info)
219+
220+
Expect(err).To(MatchError(ContainSubstring("node lookup failed")))
221+
_, statErr := os.Stat(nodePath)
222+
Expect(statErr).ToNot(HaveOccurred(), "the node should be left for a retry")
223+
})
224+
225+
// Removing the node first would destroy the parent id the retry needs.
226+
It("keeps the node when the quota cannot be released", func() {
227+
orphan()
228+
// A parent id pointing nowhere fails the walk up the tree.
229+
info.ParentID = "does-not-exist"
230+
231+
err := env.Fs.RollbackUpload(env.Ctx, idRef, "session-orphan", info)
232+
233+
Expect(err).To(MatchError(ContainSubstring("could not revert propagate")))
234+
_, statErr := os.Stat(nodePath)
235+
Expect(statErr).ToNot(HaveOccurred(), "the node should be left for a retry")
124236
})
125237
})
126238
})

‎pkg/storage/utils/decomposedfs/upload.go‎

Lines changed: 47 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ package decomposedfs
2121
import (
2222
"context"
2323
"fmt"
24+
iofs "io/fs"
2425
"os"
2526
"path/filepath"
2627
"strings"
@@ -648,10 +649,12 @@ func (fs *Decomposedfs) PrepareUpload(ctx context.Context, ref *provider.Referen
648649
// RollbackUpload reverts the node state written by PrepareUpload after a failed or aborted
649650
// postprocessing run. It restores the previous revision (or purges the node if versioning is
650651
// disabled and no prior version exists) and reverts the optimistic size propagation.
651-
func (fs *Decomposedfs) RollbackUpload(ctx context.Context, ref *provider.Reference, sessionID string, nodeExisted bool, sizeDiff int64) error {
652+
func (fs *Decomposedfs) RollbackUpload(ctx context.Context, ref *provider.Reference, sessionID string, info storage.RollbackInfo) error {
652653
n, err := fs.lu.NodeFromResource(ctx, ref)
653654
if err != nil {
654-
return fmt.Errorf("RollbackUpload: node lookup failed: %w", err)
655+
// The node metadata is unreadable, so the upload can never finish and its
656+
// quota would stay consumed forever.
657+
return fs.rollbackOrphaned(ctx, ref, sessionID, info, err)
655658
}
656659
if !n.Exists {
657660
return nil // nothing was written yet
@@ -669,7 +672,7 @@ func (fs *Decomposedfs) RollbackUpload(ctx context.Context, ref *provider.Refere
669672
return nil
670673
}
671674

672-
if nodeExisted {
675+
if info.NodeExisted {
673676
if err := n.RevertCurrentRevision(ctx, false); err != nil {
674677
return err
675678
}
@@ -681,14 +684,53 @@ func (fs *Decomposedfs) RollbackUpload(ctx context.Context, ref *provider.Refere
681684
}
682685
}
683686

684-
if sizeDiff != 0 {
685-
if err := fs.tp.Propagate(ctx, n, -sizeDiff); err != nil {
687+
if info.SizeDiff != 0 {
688+
if err := fs.tp.Propagate(ctx, n, -info.SizeDiff); err != nil {
686689
appctx.GetLogger(ctx).Error().Err(err).Msg("RollbackUpload: could not revert propagate")
687690
}
688691
}
689692
return nil
690693
}
691694

695+
// node metadata corrupt, use session ids to release quota
696+
func (fs *Decomposedfs) rollbackOrphaned(ctx context.Context, ref *provider.Reference, sessionID string, info storage.RollbackInfo, lookupErr error) error {
697+
if info.NodeID == "" || info.ParentID == "" {
698+
return fmt.Errorf("RollbackUpload: node lookup failed: %w", lookupErr)
699+
}
700+
spaceID := ref.GetResourceId().GetSpaceId()
701+
if spaceID == "" {
702+
return fmt.Errorf("RollbackUpload: node lookup failed: %w", lookupErr)
703+
}
704+
appctx.GetLogger(ctx).Info().Err(lookupErr).Str("sessionid", sessionID).Str("nodeid", info.NodeID).
705+
Msg("node unreadable, rolling back orphaned upload")
706+
n := node.New(spaceID, info.NodeID, info.ParentID, info.Filename, info.Size, sessionID,
707+
provider.ResourceType_RESOURCE_TYPE_FILE, nil, fs.lu)
708+
spaceRoot, err := node.ReadNode(ctx, fs.lu, spaceID, spaceID, false, nil, false)
709+
if err != nil {
710+
return fmt.Errorf("RollbackUpload: space root lookup failed: %w", err)
711+
}
712+
n.SpaceRoot = spaceRoot
713+
if info.SizeDiff != 0 {
714+
// return error so caller keeps session for retry
715+
if err := fs.tp.Propagate(ctx, n, -info.SizeDiff); err != nil {
716+
return fmt.Errorf("RollbackUpload: could not revert propagate: %w", err)
717+
}
718+
}
719+
nodePath := n.InternalPath()
720+
if err := utils.RemoveItem(nodePath); err != nil && !errors.Is(err, iofs.ErrNotExist) {
721+
appctx.GetLogger(ctx).Error().Err(err).Str("nodepath", nodePath).Msg("RollbackUpload: removing orphaned node failed")
722+
}
723+
if err := fs.lu.MetadataBackend().Purge(ctx, nodePath); err != nil && !errors.Is(err, iofs.ErrNotExist) {
724+
appctx.GetLogger(ctx).Error().Err(err).Str("nodepath", nodePath).Msg("RollbackUpload: purging orphaned node metadata failed")
725+
}
726+
// parent holds a child entry pointing to the now-removed node
727+
childEntry := filepath.Join(n.ParentPath(), n.Name)
728+
if err := os.Remove(childEntry); err != nil && !errors.Is(err, iofs.ErrNotExist) {
729+
appctx.GetLogger(ctx).Error().Err(err).Str("path", childEntry).Msg("RollbackUpload: removing orphaned child entry failed")
730+
}
731+
return nil
732+
}
733+
692734
func validateChecksums(ctx context.Context, lu node.PathLookup, n *node.Node, versionPath string) error {
693735
for _, t := range []string{"md5", "sha1", "adler32"} {
694736
key := prefixes.ChecksumPrefix + t

0 commit comments

Comments
 (0)