Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
90 changes: 78 additions & 12 deletions pkg/storage/fs/posix/tree/revisions.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,13 +26,16 @@ import (
"path/filepath"
"strconv"
"strings"
"syscall"
"time"

provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
"github.com/pkg/errors"
"github.com/pkg/xattr"

"github.com/opencloud-eu/reva/v2/pkg/appctx"
"github.com/opencloud-eu/reva/v2/pkg/errtypes"
"github.com/opencloud-eu/reva/v2/pkg/storage/fs/posix/blobstore"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/metadata"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/metadata/prefixes"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node"
Expand Down Expand Up @@ -280,39 +283,102 @@ func (tp *Tree) RestoreRevision(ctx context.Context, srcNode, targetNode metadat
}
defer rf.Close()

wf, err := os.OpenFile(target, os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
// Write the revision to a temp file and rename it over the target, like the blobstore does for
// uploads. Readers then see either the old or the restored content, never a mix, and a failed
// copy leaves the target untouched.
tmpDir := filepath.Join(tp.lookup.InternalSpaceRoot(targetNode.GetSpaceID()), blobstore.TMPDir)
if err := os.MkdirAll(tmpDir, 0700); err != nil {
return err
}
defer wf.Close()
err = wf.Truncate(0)
wf, err := os.CreateTemp(tmpDir, filepath.Base(target)+".restore-*")
if err != nil {
return err
}
tempName := wf.Name()
renamed := false
defer func() {
_ = wf.Close()
if !renamed {
_ = os.Remove(tempName)
}
}()

if _, err := io.Copy(wf, rf); err != nil {
return err
}
if err := wf.Sync(); err != nil {
return err
}

err = tp.lookup.CopyMetadata(ctx, srcNode, targetNode, 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
})
// keep the mode and owner of the target, the in-place write did not change them
fi, err := os.Stat(target)
if err != nil {
return err
}
if err := wf.Chmod(fi.Mode().Perm()); err != nil {
return err
}
if st, ok := fi.Sys().(*syscall.Stat_t); ok {
if err := wf.Chown(int(st.Uid), int(st.Gid)); err != nil {
tp.log.Warn().Err(err).Str("target", target).Msg("could not keep the owner of the restored file")
}
}
if err := wf.Close(); err != nil {
return err
}

// the node metadata lives in xattrs on the file itself, so copy it to the temp file before the
// rename replaces the target
names, err := xattr.List(target)
if err != nil {
return err
}
for _, name := range names {
if !strings.HasPrefix(name, prefixes.OcPrefix) {
continue
}
value, err := xattr.Get(target, name)
if err != nil {
return err
}
if err := xattr.Set(tempName, name, value); err != nil {
return err
}
}

revisionAttrs, err := tp.lookup.MetadataBackend().All(ctx, srcNode)
if err != nil {
return errtypes.InternalError("failed to copy blob xattrs to old revision to node: " + err.Error())
return errtypes.InternalError("failed to read blob xattrs of old revision: " + err.Error())
}
for name, value := range revisionAttrs {
if strings.HasPrefix(name, prefixes.ChecksumPrefix) ||
name == prefixes.TypeAttr ||
name == prefixes.BlobIDAttr ||
name == prefixes.BlobsizeAttr {
if err := xattr.Set(tempName, name, value); err != nil {
return errtypes.InternalError("failed to copy blob xattrs to old revision to node: " + err.Error())
}
}
}

// set the node mtime to the current time if no mtime was provided
if mtime.IsZero() {
mtime = time.Now()
}
err = os.Chtimes(target, mtime, mtime)
if err := xattr.Set(tempName, prefixes.MTimeAttr, []byte(mtime.UTC().Format(time.RFC3339Nano))); err != nil {
return errtypes.InternalError("failed to set mtime attribute on node: " + err.Error())
}
err = os.Chtimes(tempName, mtime, mtime)
if err != nil {
return errtypes.InternalError("failed to update times:" + err.Error())
}

if err := os.Rename(tempName, target); err != nil {
return err
}
renamed = true

// set the mtime again through the metadata backend so that it updates its cache
err = tp.lookup.MetadataBackend().SetMultiple(ctx, targetNode,
map[string][]byte{
prefixes.MTimeAttr: []byte(mtime.UTC().Format(time.RFC3339Nano)),
Expand Down
148 changes: 148 additions & 0 deletions pkg/storage/fs/posix/tree/revisions_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,148 @@
package tree_test

import (
"bytes"
"io"
"os"
"path/filepath"
"sync"
"sync/atomic"
"time"

. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"

"github.com/opencloud-eu/reva/v2/pkg/storage/fs/posix/blobstore"
helpers "github.com/opencloud-eu/reva/v2/pkg/storage/fs/posix/testhelpers"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/metadata"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/metadata/prefixes"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/node"
)

var _ = Describe("Revisions", func() {
var (
renv *helpers.TestEnv
)

BeforeEach(func() {
var err error
renv, err = helpers.NewTestEnv(map[string]any{"watch_fs": false})
Expect(err).ToNot(HaveOccurred())
})

AfterEach(func() {
if renv != nil {
renv.Cleanup()
}
})

// createFile creates a file in the space root with the given content
createFile := func(name string, content []byte) *node.Node {
n, err := renv.CreateTestFile(name, "blobid-"+name, renv.SpaceRootRes.OpaqueId, renv.SpaceRootRes.SpaceId, int64(len(content)))
Expect(err).ToNot(HaveOccurred())
Expect(os.WriteFile(n.InternalPath(), content, 0600)).To(Succeed())
return n
}

// createRevision stores the current content of n as a revision and returns the revision node
createRevision := func(n *node.Node, version string) metadata.MetadataNode {
_, err := renv.Tree.CreateRevision(renv.Ctx, n, version)
Expect(err).ToNot(HaveOccurred())
return node.NewBaseNode(n.SpaceID, n.ID+node.RevisionIDDelimiter+version, renv.Lookup)
}

tmpDirEntries := func() []os.DirEntry {
entries, _ := os.ReadDir(filepath.Join(renv.Lookup.InternalSpaceRoot(renv.SpaceRootRes.SpaceId), blobstore.TMPDir))
return entries
}

Describe("RestoreRevision", func() {
It("restores the content and the blob metadata of the revision", func() {
n := createFile("file.txt", []byte("old content"))
rev := createRevision(n, "2020-01-01T00:00:00Z")
Expect(os.WriteFile(n.InternalPath(), []byte("new content, longer"), 0600)).To(Succeed())
Expect(renv.Lookup.MetadataBackend().Set(renv.Ctx, n, prefixes.BlobsizeAttr, []byte("19"))).To(Succeed())

Expect(renv.Tree.RestoreRevision(renv.Ctx, rev, n, time.Now())).To(Succeed())

Expect(os.ReadFile(n.InternalPath())).To(Equal([]byte("old content")))
blobsize, err := renv.Lookup.MetadataBackend().GetInt64(renv.Ctx, n, prefixes.BlobsizeAttr)
Expect(err).ToNot(HaveOccurred())
Expect(blobsize).To(Equal(int64(len("old content"))))
// the node metadata survives the restore
id, err := renv.Lookup.MetadataBackend().Get(renv.Ctx, n, prefixes.IDAttr)
Expect(err).ToNot(HaveOccurred())
Expect(string(id)).To(Equal(n.ID))
Expect(tmpDirEntries()).To(BeEmpty())
})

It("leaves the target untouched when copying the revision fails", func() {
n := createFile("file.txt", []byte("current content"))
// a revision that cannot be read: its path is a directory, so reading it fails
broken := node.NewBaseNode(n.SpaceID, n.ID+node.RevisionIDDelimiter+"2020-01-01T00:00:00Z", renv.Lookup)
Expect(os.MkdirAll(broken.InternalPath(), 0700)).To(Succeed())

Expect(renv.Tree.RestoreRevision(renv.Ctx, broken, n, time.Now())).ToNot(Succeed())

Expect(os.ReadFile(n.InternalPath())).To(Equal([]byte("current content")))
id, err := renv.Lookup.MetadataBackend().Get(renv.Ctx, n, prefixes.IDAttr)
Expect(err).ToNot(HaveOccurred())
Expect(string(id)).To(Equal(n.ID))
Expect(tmpDirEntries()).To(BeEmpty())
})

It("never lets a concurrent reader see a partially restored file", func() {
const size = 8 << 20
contentA := bytes.Repeat([]byte("a"), size)
contentB := bytes.Repeat([]byte("b"), size)

n := createFile("big.bin", contentA)
revA := createRevision(n, "2020-01-01T00:00:00Z")
Expect(os.WriteFile(n.InternalPath(), contentB, 0600)).To(Succeed())
revB := createRevision(n, "2020-01-02T00:00:00Z")

var (
stop atomic.Bool
bad atomic.Int64
reads atomic.Int64
wg sync.WaitGroup
badSize atomic.Int64
)
for i := 0; i < 4; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for !stop.Load() {
f, err := os.Open(n.InternalPath())
if err != nil {
continue
}
data, err := io.ReadAll(f)
_ = f.Close()
if err != nil {
continue
}
reads.Add(1)
if !bytes.Equal(data, contentA) && !bytes.Equal(data, contentB) {
bad.Add(1)
badSize.Store(int64(len(data)))
}
}
}()
}

for i := 0; i < 10; i++ {
rev := revA
if i%2 == 1 {
rev = revB
}
Expect(renv.Tree.RestoreRevision(renv.Ctx, rev, n, time.Now())).To(Succeed())
}
stop.Store(true)
wg.Wait()

Expect(reads.Load()).To(BeNumerically(">", 0))
Expect(bad.Load()).To(BeZero(), "%d of %d reads saw neither the old nor the new content (e.g. %d bytes)", bad.Load(), reads.Load(), badSize.Load())
})
})
})
35 changes: 32 additions & 3 deletions pkg/storage/pkg/decomposedfs/upload/upload.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ import (
"github.com/opencloud-eu/reva/v2/pkg/rhttp/datatx/metrics"
"github.com/opencloud-eu/reva/v2/pkg/rhttp/datatx/utils/download"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/disk"
"github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/metadata"
"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"
Expand Down Expand Up @@ -453,9 +454,37 @@ func (session *DecomposedFsSession) Cleanup(revertNodeMetadata, cleanBin, cleanI
} else {
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")
if err := n.RevertUpload(ctx, versionID); err != nil {
sublog.Error().Err(err).Str("versionID", versionID).Msg("reverting node metadata failed")
ok := func() bool {
// lock the node so that no other upload can finish while we revert. RevertUpload
// takes the same lock, which this node object then already holds.
unlock, err := session.store.lu.MetadataBackend().Lock(n)
if err != nil {
sublog.Error().Err(err).Str("versionID", versionID).Msg("locking node failed")
return false
}
defer func() { _ = unlock() }()

// only roll back if this session still owns the node. Otherwise a newer upload has
// replaced the content since, and reverting would overwrite it with the old version.
// Read the status through the backend, the node may have cached it before we locked.
status, err := session.store.lu.MetadataBackend().Get(ctx, n, prefixes.StatusPrefix)
if err != nil && !metadata.IsAttrUnset(err) {
sublog.Error().Err(err).Str("versionID", versionID).Msg("reading processingid for session failed")
return false
}
if string(status) != node.ProcessingStatus+session.ID() {
sublog.Info().Str("versionID", versionID).Str("status", string(status)).Msg("node was changed by another upload, keeping it and the revision")
return true
}

sublog.Debug().Str("nodepath", n.InternalPath()).Str("versionID", versionID).Msg("restoring revision")
if err := n.RevertUpload(ctx, versionID); err != nil {
sublog.Error().Err(err).Str("versionID", versionID).Msg("reverting node metadata failed")
return false
}
return true
}()
if !ok {
return
}
} else {
Expand Down
54 changes: 54 additions & 0 deletions pkg/storage/pkg/decomposedfs/upload_async_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -685,4 +685,58 @@ var _ = Describe("Async file uploads", Ordered, func() {
// Expect(parentSize()).To(Equal(0))
})
})

When("two uploads overwrite an existing file in parallel", func() {
var (
olderUploadID string
newerUploadID string
olderContent = []byte("012345678901234")
)

JustBeforeEach(func() {
// the file exists with firstContent
succeedPostprocessing(uploadID)

upload := func(content []byte) string {
uploadIds, err := fs.InitiateUpload(ctx, ref, int64(len(content)), map[string]string{})
Expect(err).ToNot(HaveOccurred())
_, err = fs.Upload(ctx, storage.UploadRequest{
Ref: &provider.Reference{Path: "/" + uploadIds["simple"]},
Body: io.NopCloser(bytes.NewReader(content)),
Length: int64(len(content)),
}, nil)
Expect(err).ToNot(HaveOccurred())
_, ok := (<-pub).(events.BytesReceived)
Expect(ok).To(BeTrue())
return uploadIds["simple"]
}
olderUploadID = upload(olderContent)
newerUploadID = upload(secondContent)
})

It("keeps the newer upload when the older one fails after the newer one succeeded", func() {
succeedPostprocessing(newerUploadID)

_, status, size := fileStatus()
Expect(status).To(Equal(""))
Expect(size).To(Equal(len(secondContent)))

failPostprocessing(olderUploadID, events.PPOutcomeDelete)

// rolling back the older upload must not replace the newer content with the old version
_, status, size = fileStatus()
Expect(status).To(Equal(""))
Expect(size).To(Equal(len(secondContent)))
Expect(parentSize()).To(Equal(len(secondContent)))

// the version of the original file is kept
revisions, err := fs.ListRevisions(ctx, ref)
Expect(err).ToNot(HaveOccurred())
sizes := []uint64{}
for _, r := range revisions {
sizes = append(sizes, r.Size)
}
Expect(sizes).To(ContainElement(uint64(len(firstContent))))
})
})
})