Skip to content
Merged
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
10 changes: 7 additions & 3 deletions apps/daemon/internal/worldfs/conn_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
var (
errDead = errors.New("worldfs: the world is lost")
errInterrupted = errors.New("worldfs: interrupted before the request was sent")
errStopped = errors.New("worldfs: the world stopped")
)

// connect opens a stream and describes the service within ctx.
Expand All @@ -30,7 +31,7 @@ func (f *frontend) connect(ctx context.Context) (*sandboxfs.Client, *sandboxfs.D
return c, d, nil
}

// client returns the stream's client. After the stream failed it redials within ctx and continues only with the same service instance. A failed redial is an [ErrConnect] error, and an interrupt while waiting for the stream or redialing is errInterrupted: either way the request was never sent.
// client returns the stream's client. After the stream failed it redials within ctx and continues only with the same service instance. Once shutdown began it returns no stream. A failed redial is an [ErrConnect] error, and an interrupt while waiting for the stream or redialing is errInterrupted: either way the request was never sent.
func (f *frontend) client(ctx context.Context, interrupt <-chan struct{}) (*sandboxfs.Client, error) {
select {
case f.connTurn <- struct{}{}:
Expand All @@ -43,6 +44,9 @@ func (f *frontend) client(ctx context.Context, interrupt <-chan struct{}) (*sand
if f.dead.Load() {
return nil, errDead
}
if f.ctx.Err() != nil {
return nil, &Error{Kind: ErrConnect, Op: "reconnect", Err: errStopped}
}
if c := f.conn; c != nil {
select {
case <-c.Done():
Expand Down Expand Up @@ -126,10 +130,10 @@ func unsent(err error) bool {
return errors.Is(err, ErrConnect) || errors.As(err, &fail) && fail.Effect == sandboxwire.EffectNone && errors.Is(err, sandboxfs.ErrTransport)
}

// retryable reports whether a request certainly did nothing and may succeed when sent again: it was never sent, or the service refused it for now with ResourceExhausted and EffectNone.
// retryable reports whether a request certainly did nothing and may succeed when sent again: it was never sent, or the service refused it for now with a retryable code and EffectNone.
func retryable(err error) bool {
var fail *sandboxfs.Failure
return unsent(err) || errors.As(err, &fail) && fail.Code == sandboxfs.CodeResourceExhausted && fail.Effect == sandboxwire.EffectNone
return unsent(err) || errors.As(err, &fail) && fail.Code.Retryable() && fail.Effect == sandboxwire.EffectNone
}

// noEffect reports whether a request certainly changed nothing.
Expand Down
7 changes: 4 additions & 3 deletions apps/daemon/internal/worldfs/dir_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,11 +14,12 @@ func (f *frontend) OpenDir(_ <-chan struct{}, in *fuse.OpenIn, out *fuse.OpenOut
}
h := &handle{node: n}
if !n.synthetic() {
r, err := call(f, f.ctx, (*sandboxfs.Client).OpenDir, &sandboxfs.OpenDirRequest{Node: n.ref})
if err != nil {
id := f.ids.Next()
if _, err := call(f, f.ctx, (*sandboxfs.Client).OpenDir, &sandboxfs.OpenDirRequest{Handle: id, Node: n.ref}); err != nil {
f.settle(cleanup{handle: id, dir: true}, err)
return status(err)
}
h.server = r.Handle
h.server = id
}
*out = fuse.OpenOut{Fh: f.newHandle(h)}
return fuse.OK
Expand Down
13 changes: 9 additions & 4 deletions apps/daemon/internal/worldfs/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
//
// # Mapping
//
// A kernel node ID names a NodeRef, and each kernel lookup holds one server lookup reference; FORGET and BATCH_FORGET release the same counts with Forget. A kernel file handle names a HandleID, and a kernel lock owner is the attachment's LockOwner.
// A kernel node ID names a NodeRef, and each kernel lookup holds one server lookup reference; FORGET and BATCH_FORGET release the same counts with Forget. A kernel file handle names a HandleID the frontend chose with [sandboxfs.HandleIDs] and never reuses within the attachment, and a kernel lock owner is the attachment's LockOwner.
//
// LOOKUP Lookup
// FORGET, BATCH_FORGET Forget, batched
Expand All @@ -18,7 +18,8 @@
// LINK, SYMLINK Link, Symlink
// CREATE Create; O_EXCL is Exclusive
// OPEN Open
// READ, WRITE Read, Write; a write that stopped after a prefix is a short write
// READ Read
// WRITE Write; O_APPEND in the write's flags, as open or fcntl(F_SETFL) last set it, is Append; a write that stopped after a prefix is a short write
// FLUSH Flush with the closing lock owner
// FSYNC, FSYNCDIR Fsync of the file or directory handle
// RELEASE, RELEASEDIR Release, ReleaseDir
Expand Down Expand Up @@ -69,7 +70,11 @@
//
// # Link loss
//
// When the stream fails, requests in flight fail with EIO. The next request redials, sends Describe and continues with the same node and handle tables only when ServerInstanceID is unchanged; a request is sent again only when the stream never sent it. A Forget, Release or ReleaseDir that certainly did nothing, because the stream never sent it or the service refused it with ResourceExhausted and EffectNone, is queued and sent again after the next redial, or after a backoff of up to 2 seconds; one that may have reached the service is never sent again. A new incarnation or an ended attachment marks the world lost: every later request fails with EIO and [World.Lost] closes. The attachment has ended when the service answers StaleAttachment, or when the redial fails with a Link failure that is not retryable (sandboxlink.Code.Retryable), such as LeaseExpired or StaleGeneration; a redial refused with InstanceChanged is a new incarnation.
// When the stream fails, requests in flight fail with EIO. The next request redials, sends Describe and continues with the same node and handle tables only when ServerInstanceID is unchanged; a request is sent again only when the stream never sent it. The service serves the new stream only after every request of the failed one has finished, so a request on it sees everything the failed stream's requests did.
//
// Detach cannot overtake a request still running on a failed stream, so a Serve that fails detaches only when the export was attached on the stream still in use and nothing else is in doubt. When Attach may have taken effect without a response, when a request may still run on a failed stream, or when Detach fails, Serve fails with [ErrAttachmentDirty], and the owner of the Link attachment ends it; the service releases everything the attachment holds when its lease ends. [World.Stop] detaches without sending what is queued, since Detach drops every reference and handle the attachment holds.
// A Forget that certainly did nothing, because the stream never sent it or the service refused it with a retryable code (sandboxfs.ErrorCode.Retryable) and EffectNone, is queued and sent again after the next redial, or after a backoff of up to 2 seconds; one that may have reached the service is never sent again. A Release or ReleaseDir is queued and sent again in the same way, and also when the stream failed before its answer arrived: handle IDs are never reused, so a repeat of one that ran finds StaleHandle, which settles it as success does. An Open, Create or OpenDir that failed after it may have taken effect returns its error at once and is never sent again; the frontend queues a release of its handle ID, sent the same way, which removes any handle it left, though not a file it created or truncated.
//
// A new incarnation or an ended attachment marks the world lost: every later request fails with EIO and [World.Lost] closes. The attachment has ended when the service answers StaleAttachment, or when the redial fails with a Link failure that is not retryable (sandboxlink.Code.Retryable), such as LeaseExpired or StaleGeneration; a redial refused with InstanceChanged is a new incarnation.
//
// A Serve that fails after Attach was sent closes the stream and sends Detach on a new one, and returns within 5 seconds even when Start's context has ended or the transport blocks; closing a stream never waits for the transport. Detach then follows everything sent before, including an Attach whose response was lost. Serve fails with [ErrAttachmentDirty], and the owner of the Link attachment ends it, only when that Detach cannot be sent or answered in that time: the service is unreachable, the redial is refused with a retryable Link failure, the transport blocks, or Detach fails or its stream fails before the answer. The service releases everything the attachment holds when its lease ends. [World.Stop] detaches without sending what is queued, since Detach drops every reference and handle the attachment holds.
package worldfs
21 changes: 17 additions & 4 deletions apps/daemon/internal/worldfs/drain_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,13 @@ package worldfs

import (
"context"
"errors"
"time"

"github.com/MiniMax-AI/OpenAgentCore/internal/sandboxfs"
)

// cleanup is a Release or ReleaseDir of a handle the kernel closed.
// cleanup is a Release or ReleaseDir of a handle the kernel closed, or of a handle ID whose acquisition may have taken effect without a response.
type cleanup struct {
handle sandboxfs.HandleID
dir bool
Expand Down Expand Up @@ -95,27 +96,39 @@ func (f *frontend) sendForgets() bool {
return true
}

// release sends a Release or ReleaseDir for a handle the kernel closed. One that certainly did nothing and may succeed later is queued for the drainer; one that may have reached the service is never sent again.
// settle queues a release of the handle ID of an acquisition that failed after it may have taken effect, so the kernel request returns its error without waiting for the service. The drainer sends it after the acquisition, on its stream or on a successor the service fences behind it, so the Release finds whatever the acquisition left. The acquisition itself is never sent again.
func (f *frontend) settle(c cleanup, err error) {
if !noEffect(err) {
f.queue(c)
}
}

// release sends a Release or ReleaseDir for a handle the kernel closed. One the service has not answered is queued for the drainer.
func (f *frontend) release(c cleanup) {
if f.sendRelease(f.ctx, c) {
return
}
if f.seams.queue != nil {
f.seams.queue()
}
f.queue(c)
}

// queue hands c to the drainer.
func (f *frontend) queue(c cleanup) {
f.mu.Lock()
f.releases = append(f.releases, c)
f.wake()
f.mu.Unlock()
}

// sendRelease sends c and reports false when it must be sent again.
// sendRelease sends c and reports false when it must be sent again: when it certainly did nothing and may succeed later, or when the stream failed before its answer arrived. The frontend never reuses a handle ID, so a repeat of a Release that ran finds StaleHandle, which settles it as success does.
func (f *frontend) sendRelease(ctx context.Context, c cleanup) bool {
var err error
if c.dir {
_, err = call(f, ctx, (*sandboxfs.Client).ReleaseDir, &sandboxfs.ReleaseDirRequest{Handle: c.handle})
} else {
_, err = call(f, ctx, (*sandboxfs.Client).Release, &sandboxfs.ReleaseRequest{Handle: c.handle})
}
return err == nil || !retryable(err) || f.dead.Load() || f.closed.Load()
return err == nil || !retryable(err) && !errors.Is(err, sandboxfs.ErrTransport) || f.dead.Load() || f.closed.Load()
}
2 changes: 1 addition & 1 deletion apps/daemon/internal/worldfs/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ var (
ErrInstanceChanged = errors.New("worldfs: file service instance changed")
ErrAttachmentLost = errors.New("worldfs: attachment ended")
ErrTopologyChanged = errors.New("worldfs: pinned topology changed")
// ErrAttachmentDirty is a failed Serve that cannot show the attachment holds nothing: the owner of the Link attachment must end it.
// ErrAttachmentDirty is a failed Serve whose cleanup Detach could not be sent or answered, so it cannot show the attachment holds nothing: the owner of the Link attachment must end it.
ErrAttachmentDirty = errors.New("worldfs: the attachment may still hold state")
)

Expand Down
19 changes: 11 additions & 8 deletions apps/daemon/internal/worldfs/files_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ func openFlags(flags uint32) (acc sandboxfs.AccessMode, of sandboxfs.OpenFlags,
for _, m := range []struct {
bit uint32
flag sandboxfs.OpenFlags
}{{syscall.O_APPEND, sandboxfs.OpenAppend}, {syscall.O_TRUNC, sandboxfs.OpenTruncate}, {syscall.O_NOFOLLOW, sandboxfs.OpenNoFollow}} {
}{{syscall.O_TRUNC, sandboxfs.OpenTruncate}, {syscall.O_NOFOLLOW, sandboxfs.OpenNoFollow}} {
if flags&m.bit != 0 {
of |= m.flag
}
Expand Down Expand Up @@ -53,11 +53,12 @@ func (f *frontend) Open(_ <-chan struct{}, in *fuse.OpenIn, out *fuse.OpenOut) f
return fuse.EPERM
}
} else {
r, err := call(f, f.ctx, (*sandboxfs.Client).Open, &sandboxfs.OpenRequest{Node: n.ref, Access: acc, Flags: of})
if err != nil {
id := f.ids.Next()
if _, err := call(f, f.ctx, (*sandboxfs.Client).Open, &sandboxfs.OpenRequest{Handle: id, Node: n.ref, Access: acc, Flags: of}); err != nil {
f.settle(cleanup{handle: id}, err)
return status(err)
}
h.server = r.Handle
h.server = id
}
*out = fuse.OpenOut{Fh: f.newHandle(h), OpenFlags: fuse.FOPEN_DIRECT_IO}
return fuse.OK
Expand All @@ -72,14 +73,16 @@ func (f *frontend) Create(_ <-chan struct{}, in *fuse.CreateIn, name string, out
if !ok {
return fuse.EINVAL
}
id := f.ids.Next()
r, err := call(f, f.ctx, (*sandboxfs.Client).Create, &sandboxfs.CreateRequest{
Parent: p.ref, Name: []byte(name), Mode: in.Mode & sandboxfs.ModePerm, Access: acc, Flags: of, Exclusive: in.Flags&syscall.O_EXCL != 0,
Handle: id, Parent: p.ref, Name: []byte(name), Mode: in.Mode & sandboxfs.ModePerm, Access: acc, Flags: of, Exclusive: in.Flags&syscall.O_EXCL != 0,
})
if err != nil {
f.settle(cleanup{handle: id}, err)
return status(err)
}
n := f.adopt(r.Entry, &out.EntryOut)
out.OpenOut = fuse.OpenOut{Fh: f.newHandle(&handle{node: n, server: r.Handle}), OpenFlags: fuse.FOPEN_DIRECT_IO}
out.OpenOut = fuse.OpenOut{Fh: f.newHandle(&handle{node: n, server: id}), OpenFlags: fuse.FOPEN_DIRECT_IO}
return fuse.OK
}

Expand All @@ -98,7 +101,7 @@ func (f *frontend) Read(_ <-chan struct{}, in *fuse.ReadIn, _ []byte) (fuse.Read
return fuse.ReadResultData(r.Data), fuse.OK
}

// Write reports a short write when the service stopped after a prefix.
// Write appends when the write's flags hold O_APPEND, which fcntl(F_SETFL) may have set or cleared since the open. It reports a short write when the service stopped after a prefix.
func (f *frontend) Write(_ <-chan struct{}, in *fuse.WriteIn, data []byte) (uint32, fuse.Status) {
h, st := f.handle(in.Fh)
if !st.Ok() {
Expand All @@ -108,7 +111,7 @@ func (f *frontend) Write(_ <-chan struct{}, in *fuse.WriteIn, data []byte) (uint
return 0, fuse.EBADF
}
data = data[:min(uint32(len(data)), f.caps.MaxWriteBytes)]
r, err := call(f, f.ctx, (*sandboxfs.Client).Write, &sandboxfs.WriteRequest{Handle: h.server, Offset: in.Offset, Data: data})
r, err := call(f, f.ctx, (*sandboxfs.Client).Write, &sandboxfs.WriteRequest{Handle: h.server, Offset: in.Offset, Append: in.Flags&syscall.O_APPEND != 0, Data: data})
if err != nil {
return 0, status(err)
}
Expand Down
Loading
Loading