diff --git a/apps/daemon/internal/worldfs/conn_linux.go b/apps/daemon/internal/worldfs/conn_linux.go index 4f9e625f..603b8ded 100644 --- a/apps/daemon/internal/worldfs/conn_linux.go +++ b/apps/daemon/internal/worldfs/conn_linux.go @@ -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. @@ -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{}{}: @@ -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(): @@ -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. diff --git a/apps/daemon/internal/worldfs/dir_linux.go b/apps/daemon/internal/worldfs/dir_linux.go index e8e119fe..cc8bfe05 100644 --- a/apps/daemon/internal/worldfs/dir_linux.go +++ b/apps/daemon/internal/worldfs/dir_linux.go @@ -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 diff --git a/apps/daemon/internal/worldfs/doc.go b/apps/daemon/internal/worldfs/doc.go index e97316f4..ee1ac969 100644 --- a/apps/daemon/internal/worldfs/doc.go +++ b/apps/daemon/internal/worldfs/doc.go @@ -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 @@ -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 @@ -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 diff --git a/apps/daemon/internal/worldfs/drain_linux.go b/apps/daemon/internal/worldfs/drain_linux.go index c6551e16..fb9c6399 100644 --- a/apps/daemon/internal/worldfs/drain_linux.go +++ b/apps/daemon/internal/worldfs/drain_linux.go @@ -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 @@ -95,7 +96,14 @@ 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 @@ -103,13 +111,18 @@ func (f *frontend) release(c cleanup) { 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 { @@ -117,5 +130,5 @@ func (f *frontend) sendRelease(ctx context.Context, c cleanup) bool { } 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() } diff --git a/apps/daemon/internal/worldfs/errors.go b/apps/daemon/internal/worldfs/errors.go index 855eb961..aca4f82a 100644 --- a/apps/daemon/internal/worldfs/errors.go +++ b/apps/daemon/internal/worldfs/errors.go @@ -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") ) diff --git a/apps/daemon/internal/worldfs/files_linux.go b/apps/daemon/internal/worldfs/files_linux.go index 86fda7f3..f2347df3 100644 --- a/apps/daemon/internal/worldfs/files_linux.go +++ b/apps/daemon/internal/worldfs/files_linux.go @@ -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 } @@ -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 @@ -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 } @@ -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() { @@ -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) } diff --git a/apps/daemon/internal/worldfs/frontend_linux_test.go b/apps/daemon/internal/worldfs/frontend_linux_test.go index a0642fa6..5c98b039 100644 --- a/apps/daemon/internal/worldfs/frontend_linux_test.go +++ b/apps/daemon/internal/worldfs/frontend_linux_test.go @@ -26,10 +26,13 @@ import ( // These tests drive the frontend's FUSE methods directly, without a mount, so they need no privileges. -// counts records what the service ran. It refuses Releases while refuse is positive, and when held is set it closes held once a waiting lock request arrives and ends that request Cancelled with EffectPossible once it is cancelled. +// counts records what the service ran. It refuses Releases while refuse is positive, and when held is set it closes held once a waiting lock request arrives and ends that request Cancelled with EffectPossible once it is cancelled. When cut is set, the next Open runs and then breaks every stream, so its response is lost. type counts struct { + srv *fileservicetest.Server refuse atomic.Int32 + opened atomic.Int32 released atomic.Int32 + cut atomic.Bool held chan struct{} mu sync.Mutex @@ -55,6 +58,17 @@ func (c counting) Release(ctx context.Context, a sandboxfs.Attachment, q *sandbo return r, err } +func (c counting) Open(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.OpenRequest) (*sandboxfs.OpenResponse, error) { + r, err := c.Service.Open(ctx, a, q) + if err == nil { + c.opened.Add(1) + } + if c.cut.CompareAndSwap(true, false) { + c.srv.Break() + } + return r, err +} + func (c counting) SetLock(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.SetLockRequest) (*sandboxfs.SetLockResponse, error) { if c.held == nil || !q.Wait { return c.Service.SetLock(ctx, a, q) @@ -77,7 +91,7 @@ func newServer(t *testing.T) (*fileservicetest.Server, *counts, string) { t.Fatal(err) } t.Cleanup(func() { srv.Close() }) - c := &counts{} + c := &counts{srv: srv} srv.Intercept(func(s sandboxfs.Service) sandboxfs.Service { return counting{s, c} }) return srv, c, path } @@ -232,6 +246,152 @@ func TestReleaseQueuedAfterRedial(t *testing.T) { eventually(t, "the queued Release", func() bool { return c.released.Load() == 1 }) } +// An Open whose response is lost while the service is unreachable fails with EIO at once, and its handle ID is released on the next stream once the service is reachable, after the service finished the Open, so the service holds no handle. +func TestLostOpenIsReleased(t *testing.T) { + srv, c, _ := newServer(t) + var down atomic.Bool + reachable := make(chan struct{}) + f := attached(t, func(ctx context.Context) (io.ReadWriteCloser, error) { + if down.Load() { + select { + case <-reachable: + case <-ctx.Done(): + return nil, ctx.Err() + } + } + return srv.Dial(ctx) + }) + draining(t, f) + var e fuse.EntryOut + if st := f.Lookup(nil, &fuse.InHeader{NodeId: f.root.id}, "f", &e); !st.Ok() { + t.Fatalf("Lookup: %v", st) + } + down.Store(true) + c.cut.Store(true) + got := make(chan fuse.Status, 1) + go func() { + var o fuse.OpenOut + got <- f.Open(nil, &fuse.OpenIn{InHeader: fuse.InHeader{NodeId: e.NodeId}, Flags: syscall.O_RDWR}, &o) + }() + wantStatus(t, got, fuse.EIO) + close(reachable) + eventually(t, "the release of the lost handle", func() bool { return c.released.Load() == 1 }) + if c.opened.Load() != 1 { + t.Fatalf("%d opened, want 1", c.opened.Load()) + } +} + +// lateAttach closes attaching when Attach arrives and attaches only once the request was cancelled, as a handler that outlives its stream does. +type lateAttach struct { + sandboxfs.Service + attaching chan struct{} +} + +func (l lateAttach) Attach(ctx context.Context, a sandboxfs.Attachment, q *sandboxfs.AttachRequest) (*sandboxfs.AttachResponse, error) { + close(l.attaching) + <-ctx.Done() + if _, err := l.Service.Attach(context.WithoutCancel(ctx), a, q); err != nil { + return nil, err + } + return nil, ctx.Err() +} + +// A Serve whose context ends while Attach is in doubt detaches on a new stream, which the service serves after the late Attach took effect, so the attachment holds nothing and the failure is not ErrAttachmentDirty. +func TestUncertainAttachIsDetached(t *testing.T) { + srv, err := fileservicetest.New(t.TempDir()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { srv.Close() }) + attaching := make(chan struct{}) + srv.Intercept(func(s sandboxfs.Service) sandboxfs.Service { return lateAttach{s, attaching} }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go func() { + <-attaching + cancel() + }() + if _, _, err := New(fileservicetest.Export, srv.Dial).Serve(ctx, nil, sessionview.WorldMount{}); !errors.Is(err, context.Canceled) || errors.Is(err, ErrAttachmentDirty) { + t.Fatalf("Serve = %v, want the cancellation without ErrAttachmentDirty", err) + } + rw, err := srv.Dial(context.Background()) + if err != nil { + t.Fatal(err) + } + c := sandboxfs.NewClient(rw) + defer c.Close() + var fail *sandboxfs.Failure + if _, err := c.Detach(context.Background(), &sandboxfs.DetachRequest{}); !errors.As(err, &fail) || fail.Code != sandboxfs.CodeStaleAttachment { + t.Fatalf("Detach after Serve: %v, want StaleAttachment", err) + } +} + +// jammed passes reads through and, while jam is set, holds every write and Close until free closes, as a transport whose outgoing traffic is stuck does. +type jammed struct { + io.ReadWriteCloser + jam *atomic.Bool + free <-chan struct{} +} + +func (j jammed) Write(p []byte) (int, error) { + if j.jam.Load() { + <-j.free + } + return j.ReadWriteCloser.Write(p) +} + +func (j jammed) Close() error { + if j.jam.Load() { + <-j.free + } + return j.ReadWriteCloser.Close() +} + +// A Serve whose context ends while Attach is in doubt and the transport holds every write returns ErrAttachmentDirty within detachWait, teardown included. +func TestAbortBoundedWhenTransportBlocks(t *testing.T) { + srv, err := fileservicetest.New(t.TempDir()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { srv.Close() }) + attaching := make(chan struct{}) + srv.Intercept(func(s sandboxfs.Service) sandboxfs.Service { return lateAttach{s, attaching} }) + var jam atomic.Bool + free := make(chan struct{}) + t.Cleanup(func() { close(free) }) + w := New(fileservicetest.Export, func(ctx context.Context) (io.ReadWriteCloser, error) { + rw, err := srv.Dial(ctx) + if err != nil { + return nil, err + } + return jammed{rw, &jam, free}, nil + }) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go func() { + <-attaching + // A request completed after Attach shows that Attach's write has finished, so the cancellation finds Attach in doubt rather than mid-write. + if _, err := w.fs.conn.Describe(context.Background(), &sandboxfs.DescribeRequest{}); err != nil { + t.Errorf("Describe beside Attach: %v", err) + } + jam.Store(true) + cancel() + }() + served := make(chan error, 1) + go func() { + _, _, err := w.Serve(ctx, nil, sessionview.WorldMount{}) + served <- err + }() + select { + case err := <-served: + if !errors.Is(err, ErrAttachmentDirty) { + t.Fatalf("Serve = %v, want ErrAttachmentDirty", err) + } + case <-time.After(detachWait + 2*time.Second): + t.Fatal("Serve still cleaning up 2s after the detach bound") + } +} + // stalled returns a frontend whose stream has failed and whose service is unreachable, with fh open, and a channel closed once a redial waits. func stalled(t *testing.T) (f *frontend, fh uint64, dialing chan struct{}) { srv, _, _ := newServer(t) diff --git a/apps/daemon/internal/worldfs/world_linux.go b/apps/daemon/internal/worldfs/world_linux.go index 41ee6063..2aa19da2 100644 --- a/apps/daemon/internal/worldfs/world_linux.go +++ b/apps/daemon/internal/worldfs/world_linux.go @@ -59,7 +59,7 @@ func New(export sandboxlink.ExportID, dial Dial) *World { // Serve has the signature of [sessionview.World]. Within ctx it connects, attaches the export and presents the mountpoints; it then starts serving dev. A World serves once. // -// A Serve that fails after the export may have been attached detaches it only when nothing can still reach the attachment: when Attach may have taken effect without a response, when a request may still run on a failed stream, or when Detach fails, Serve returns [ErrAttachmentDirty] and the owner of the Link attachment must end it. +// A Serve that fails after the export may have been attached detaches it on a new stream, and returns within 5 more seconds even when ctx has ended or the transport blocks. When that Detach cannot be sent or answered in that time, Serve returns [ErrAttachmentDirty] and the owner of the Link attachment must end it. func (w *World) Serve(ctx context.Context, dev *os.File, mount sessionview.WorldMount) (sessionview.WorldServer, sessionview.Presentation, error) { p, err := w.fs.serve(ctx, dev, mount) if err != nil { @@ -68,7 +68,7 @@ func (w *World) Serve(ctx context.Context, dev *os.File, mount sessionview.World return w, p, nil } -// Stop waits up to 10 seconds for serving to end, which happens once the view's mount namespace is gone. It then detaches, ending every request still waiting on the service after 5 more seconds, so it returns within 15 seconds. It unmounts nothing. After a Serve that failed, or without one, it returns at once. +// Stop waits up to 10 seconds for serving to end, which happens once the view's mount namespace is gone. It then detaches, ending every request still waiting on the service after 5 more seconds, so it returns within 15 seconds even when the transport blocks. It unmounts nothing. After a Serve that failed, or without one, it returns at once. func (w *World) Stop() error { w.fs.stopOnce.Do(func() { w.fs.stopErr = w.fs.stop() }) return w.fs.stopErr @@ -96,7 +96,8 @@ type frontend struct { connTurn chan struct{} // held to use or replace conn; a channel so that waiting for it can be interrupted conn *sandboxfs.Client - gen atomic.Uint64 // counts redials + ids sandboxfs.HandleIDs // the attachment's handle IDs, never reused + gen atomic.Uint64 // counts redials instance sandboxwire.ID service sandboxfs.Identity // the identity the service acts as view sandboxfs.Identity // the identity the view's processes run as @@ -111,7 +112,7 @@ type frontend struct { root *inode born sandboxfs.Timestamp forgets map[sandboxfs.NodeRef]uint64 // references the kernel released, not yet sent - releases []cleanup // handles the kernel released whose Release a failed stream never sent + releases []cleanup // releases for the drainer: of handles the kernel closed whose Release went unanswered, and of acquisitions in doubt kick chan struct{} // wakes the drainer drainCtx context.Context @@ -213,24 +214,18 @@ func (f *frontend) start(dev *os.File, opts *fuse.MountOptions) error { return nil } -// abort undoes a failed Serve and returns its error. Detach on the stream the export was attached on releases what the attachment holds only when nothing else can still reach the attachment, so it is sent only then, and any doubt is [ErrAttachmentDirty]. +// abort undoes a failed Serve and returns its error. When the export may be attached, it drops the stream and detaches on a new one: the service serves that stream only after every request of the earlier ones has finished, so Detach releases whatever Attach, or a request the failure abandoned, left. Detach is settled when it succeeds or the world is lost, since a lost attachment or incarnation holds nothing; anything else, including no answer within detachWait, is [ErrAttachmentDirty]. func (f *frontend) abort(ctx context.Context, err error) error { close(f.served) close(f.drained) defer f.shutdown() - switch { - case f.maybe: - return &Error{Kind: ErrAttachmentDirty, Op: "attach", Err: err} - case !f.attached || f.dead.Load(): + if !f.attached && !f.maybe || f.dead.Load() { return err - case f.gen.Load() != 0: - // A request on the failed stream may still be running. - return &Error{Kind: ErrAttachmentDirty, Op: "present", Err: err} } - // Once ctx has ended, a request it abandoned may still run, and Detach under it fails, so the attachment is dirty. - ctx, cancel := context.WithTimeout(ctx, detachWait) + f.drop() + ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), detachWait) defer cancel() - if _, derr := f.conn.Detach(ctx, &sandboxfs.DetachRequest{}); derr != nil { + if _, derr := call(f, ctx, (*sandboxfs.Client).Detach, &sandboxfs.DetachRequest{}); derr != nil && !f.dead.Load() { return &Error{Kind: ErrAttachmentDirty, Op: "detach", Err: errors.Join(err, derr)} } return err @@ -306,11 +301,19 @@ func (f *frontend) stop() error { return errors.Join(errs...) } +// shutdown ends every request and redial on f.ctx and closes the stream. client returns no stream after it. func (f *frontend) shutdown() { f.cancel() + f.drop() +} + +// drop closes the stream, which never waits for the transport, so the next request redials. +func (f *frontend) drop() { f.connTurn <- struct{}{} - defer func() { <-f.connTurn }() - if f.conn != nil { - f.conn.Close() + c := f.conn + f.conn = nil + <-f.connTurn + if c != nil { + c.Close() } } diff --git a/apps/daemon/internal/worldfs/world_linux_test.go b/apps/daemon/internal/worldfs/world_linux_test.go index 56aa9274..a2a5d6af 100644 --- a/apps/daemon/internal/worldfs/world_linux_test.go +++ b/apps/daemon/internal/worldfs/world_linux_test.go @@ -247,6 +247,42 @@ func TestInstanceChanged(t *testing.T) { } } +// Append is a property of each write: fcntl(F_SETFL) clears and sets O_APPEND on an open descriptor, as on a native file. +func TestAppendFollowsFcntl(t *testing.T) { + requireFUSE(t) + backing := t.TempDir() + writeFile(t, filepath.Join(backing, "f"), "abc") + m, err := serve(t, backing, 0) + if err != nil { + t.Fatalf("Serve: %v", err) + } + fd, err := unix.Open(filepath.Join(m.dir, "f"), unix.O_RDWR|unix.O_APPEND, 0) + if err != nil { + t.Fatal(err) + } + defer unix.Close(fd) + want := func(content string) { + t.Helper() + if b, err := os.ReadFile(filepath.Join(backing, "f")); err != nil || string(b) != content { + t.Fatalf("file = %q, %v; want %q", b, err, content) + } + } + if _, err := unix.FcntlInt(uintptr(fd), unix.F_SETFL, 0); err != nil { + t.Fatal(err) + } + if _, err := unix.Pwrite(fd, []byte("X"), 0); err != nil { + t.Fatal(err) + } + want("Xbc") + if _, err := unix.FcntlInt(uintptr(fd), unix.F_SETFL, unix.O_APPEND); err != nil { + t.Fatal(err) + } + if _, err := unix.Pwrite(fd, []byte("Y"), 0); err != nil { + t.Fatal(err) + } + want("XbcY") +} + // Stop returns within its bound when the service has become unreachable, although Detach must redial. func TestStopUnreachable(t *testing.T) { requireFUSE(t) @@ -283,7 +319,7 @@ func (u unanswered) Attach(ctx context.Context, _ sandboxfs.Attachment, _ *sandb return nil, ctx.Err() } -// A Start that ends while Attach is unanswered keeps the world's report that the attachment must be ended. +// A Start that ends while Attach is unanswered reports the cancellation, and the world's Detach on a new stream leaves nothing for the attachment's owner to end. func TestStartEndsDuringAttach(t *testing.T) { requireFUSE(t) srv, err := fileservicetest.New(t.TempDir()) @@ -304,8 +340,8 @@ func TestStartEndsDuringAttach(t *testing.T) { StagingParent: t.TempDir(), Process: sessionview.Process{Path: "/bin/true", Args: []string{"true"}, Dir: "/", UID: 1000, GID: 1000, Stderr: os.Stderr}, }) - if !errors.Is(err, worldfs.ErrAttachmentDirty) || !errors.Is(err, context.Canceled) { - t.Fatalf("Start = %v, want ErrAttachmentDirty with the cancellation", err) + if errors.Is(err, worldfs.ErrAttachmentDirty) || !errors.Is(err, context.Canceled) { + t.Fatalf("Start = %v, want the cancellation without ErrAttachmentDirty", err) } } diff --git a/apps/sandboxio/fileservicetest/fileservicetest.go b/apps/sandboxio/fileservicetest/fileservicetest.go index a8a5d0ce..d4a3748d 100644 --- a/apps/sandboxio/fileservicetest/fileservicetest.go +++ b/apps/sandboxio/fileservicetest/fileservicetest.go @@ -26,6 +26,8 @@ type Server struct { mu sync.Mutex svc *fileservice.Service + files *sandboxfs.Server // serves svc, as wrapped, to every stream + binds uint64 // stands in for Link's bind sequence: streams are bound in Dial order conns []net.Conn stalled bool wrap func(sandboxfs.Service) sandboxfs.Service @@ -37,7 +39,18 @@ func New(dir string) (*Server, error) { if err != nil { return nil, err } - return &Server{dir: dir, id: sandboxwire.NewID(), svc: svc}, nil + s := &Server{dir: dir, id: sandboxwire.NewID(), svc: svc} + s.serveLocked() + return s, nil +} + +// serveLocked starts serving the current incarnation, as wrapped, to later streams. +func (s *Server) serveLocked() { + var svc sandboxfs.Service = s.svc + if s.wrap != nil { + svc = s.wrap(svc) + } + s.files = sandboxfs.NewServer(svc) } // Dial opens a stream to the current service incarnation. While the server is stalled it waits for ctx to end. @@ -60,19 +73,19 @@ func (s *Server) Dial(ctx context.Context) (io.ReadWriteCloser, error) { Lease: context.Background(), Exports: []sandboxlink.ExportGrant{{ID: Export}}, } - var svc sandboxfs.Service = s.svc - if s.wrap != nil { - svc = s.wrap(svc) - } - go sandboxfs.Serve(context.Background(), server, svc, a) + s.binds++ + go s.files.Serve(context.Background(), server, a, s.binds) return client, nil } -// Intercept wraps the service each later stream serves, so a test can change what it answers. +// Intercept wraps the service each later stream serves, so a test can change what it answers. Streams dialed before and after it are not fenced against each other, so a test calls it before the first Dial. func (s *Server) Intercept(wrap func(sandboxfs.Service) sandboxfs.Service) { s.mu.Lock() defer s.mu.Unlock() s.wrap = wrap + if s.svc != nil { + s.serveLocked() + } } // Break closes every open stream, as a lost transport does. @@ -106,8 +119,13 @@ func (s *Server) Restart() error { s.svc.Close() } svc, err := fileservice.New(s.dir) + if err != nil { + s.svc = nil + return err + } s.svc = svc - return err + s.serveLocked() + return nil } // Close closes every stream and the service. diff --git a/apps/sandboxio/internal/fileservice/files.go b/apps/sandboxio/internal/fileservice/files.go index 594e7121..65d336dd 100644 --- a/apps/sandboxio/internal/fileservice/files.go +++ b/apps/sandboxio/internal/fileservice/files.go @@ -19,7 +19,7 @@ func openFlags(access sandboxfs.AccessMode, flags sandboxfs.OpenFlags) int { f := map[sandboxfs.AccessMode]int{sandboxfs.AccessRead: unix.O_RDONLY, sandboxfs.AccessWrite: unix.O_WRONLY, sandboxfs.AccessReadWrite: unix.O_RDWR}[access] f |= unix.O_NOCTTY for bit, o := range map[sandboxfs.OpenFlags]int{ - sandboxfs.OpenAppend: unix.O_APPEND, sandboxfs.OpenTruncate: unix.O_TRUNC, sandboxfs.OpenNoFollow: unix.O_NOFOLLOW, + sandboxfs.OpenTruncate: unix.O_TRUNC, sandboxfs.OpenNoFollow: unix.O_NOFOLLOW, sandboxfs.OpenSync: unix.O_SYNC, sandboxfs.OpenDataSync: unix.O_DSYNC, } { if flags&bit != 0 { @@ -77,7 +77,7 @@ func (s *Service) Open(_ context.Context, a sandboxfs.Attachment, r *sandboxfs.O if err := openable(n.typ); err != nil { return nil, failure(err, none) } - if err := st.reserve(); err != nil { + if err := st.reserve(r.Handle); err != nil { return nil, err } var fd int @@ -86,14 +86,13 @@ func (s *Service) Open(_ context.Context, a sandboxfs.Attachment, r *sandboxfs.O return err }) if err != nil { - st.unreserve() + st.unreserve(r.Handle) return nil, failure(err, none) } - id, err := st.addHandle(&handle{f: os.NewFile(uintptr(fd), ""), append: r.Flags&sandboxfs.OpenAppend != 0}) - if err != nil { + if err := st.publish(r.Handle, &handle{f: os.NewFile(uintptr(fd), "")}); err != nil { return nil, failure(err, possible) } - return &sandboxfs.OpenResponse{Handle: id}, nil + return &sandboxfs.OpenResponse{}, nil } // Create creates the entry with O_EXCL first, so it never opens a special @@ -108,7 +107,7 @@ func (s *Service) Create(_ context.Context, a sandboxfs.Attachment, r *sandboxfs if err != nil { return nil, err } - if err := st.reserve(); err != nil { + if err := st.reserve(r.Handle); err != nil { return nil, err } flags := openFlags(r.Access, r.Flags) | unix.O_NOFOLLOW @@ -141,20 +140,19 @@ func (s *Service) Create(_ context.Context, a sandboxfs.Attachment, r *sandboxfs return err }) if err != nil { - st.unreserve() + st.unreserve(r.Handle) return nil, failure(err, effect) } n, err := st.addNode(pathFD, &sb) if err != nil { - st.unreserve() + st.unreserve(r.Handle) unix.Close(fd) return nil, failure(err, possible) } - id, err := st.addHandle(&handle{f: os.NewFile(uintptr(fd), ""), append: r.Flags&sandboxfs.OpenAppend != 0}) - if err != nil { + if err := st.publish(r.Handle, &handle{f: os.NewFile(uintptr(fd), "")}); err != nil { return nil, failure(err, possible) } - return &sandboxfs.CreateResponse{Entry: sandboxfs.Entry{Node: n.ref, Attr: st.attr(&sb)}, Handle: id}, nil + return &sandboxfs.CreateResponse{Entry: sandboxfs.Entry{Node: n.ref, Attr: st.attr(&sb)}}, nil } // openExisting opens an existing regular file without following a symlink. @@ -209,8 +207,10 @@ func (s *Service) Read(_ context.Context, a sandboxfs.Attachment, r *sandboxfs.R return &sandboxfs.ReadResponse{Data: buf[:total]}, nil } -// Write writes at the offset, or appends atomically with one write(2) on an -// append handle. It reports the written prefix and the error that stopped it. +// Write writes at the offset, or appends atomically with one write on the +// descriptor with O_APPEND set. It reports the written prefix and the error +// that stopped it; a short append is reported as it is, never continued by +// another append. func (s *Service) Write(_ context.Context, a sandboxfs.Attachment, r *sandboxfs.WriteRequest) (*sandboxfs.WriteResponse, error) { st, err := s.enter(a, r) if err != nil { @@ -220,12 +220,17 @@ func (s *Service) Write(_ context.Context, a sandboxfs.Attachment, r *sandboxfs. if err != nil { return nil, err } - if !h.append && r.Offset > math.MaxInt64-uint64(len(r.Data)) { + if !r.Append && r.Offset > math.MaxInt64-uint64(len(r.Data)) { return nil, sandboxfs.NewFailure(sandboxfs.CodeInvalidArgument, none, "write ends beyond 2^63-1") } + h.writeMu.Lock() + defer h.writeMu.Unlock() total := 0 err = use(h.f, errStaleHandle, func(fd int) error { - if h.append { + if err := setAppend(fd, r.Append); err != nil { + return err + } + if r.Append { n, err := eintr(func() (int, error) { return unix.Write(fd, r.Data) }) total = max(n, 0) return err @@ -251,6 +256,18 @@ func (s *Service) Write(_ context.Context, a sandboxfs.Attachment, r *sandboxfs. return &sandboxfs.WriteResponse{Written: uint32(total), Failure: failure(err, none).(*sandboxfs.Failure)}, nil } +// setAppend sets or clears O_APPEND on fd. An append uses the descriptor's +// flag rather than pwritev2's RWF_APPEND: overlayfs on some kernels drops +// RWF_APPEND and writes at the offset, but every file system honors the flag. +func setAppend(fd int, on bool) error { + fl, err := unix.FcntlInt(uintptr(fd), unix.F_GETFL, 0) + if err != nil || (fl&unix.O_APPEND != 0) == on { + return err + } + _, err = unix.FcntlInt(uintptr(fd), unix.F_SETFL, fl^unix.O_APPEND) + return err +} + // Flush closes a duplicate of the handle's descriptor, which is what a close // of one descriptor does to the file. func (s *Service) Flush(_ context.Context, a sandboxfs.Attachment, r *sandboxfs.FlushRequest) (*sandboxfs.FlushResponse, error) { @@ -322,7 +339,7 @@ func (s *Service) OpenDir(_ context.Context, a sandboxfs.Attachment, r *sandboxf if err != nil { return nil, err } - if err := st.reserve(); err != nil { + if err := st.reserve(r.Handle); err != nil { return nil, err } var fd int @@ -334,14 +351,13 @@ func (s *Service) OpenDir(_ context.Context, a sandboxfs.Attachment, r *sandboxf return err }) if err != nil { - st.unreserve() + st.unreserve(r.Handle) return nil, failure(err, none) } - id, err := st.addHandle(&handle{f: os.NewFile(uintptr(fd), ""), dir: &cursor{dev: uint64(sb.Dev)}}) - if err != nil { + if err := st.publish(r.Handle, &handle{f: os.NewFile(uintptr(fd), ""), dir: &cursor{dev: uint64(sb.Dev)}}); err != nil { return nil, err } - return &sandboxfs.OpenDirResponse{Handle: id}, nil + return &sandboxfs.OpenDirResponse{}, nil } func (s *Service) ReadDir(_ context.Context, a sandboxfs.Attachment, r *sandboxfs.ReadDirRequest) (*sandboxfs.ReadDirResponse, error) { diff --git a/apps/sandboxio/internal/fileservice/service.go b/apps/sandboxio/internal/fileservice/service.go index 33c722ea..27b941fd 100644 --- a/apps/sandboxio/internal/fileservice/service.go +++ b/apps/sandboxio/internal/fileservice/service.go @@ -120,7 +120,7 @@ func probeRename() (noReplace, exchange bool) { } // InstanceID is the service incarnation. Link binds each stream to it, and -// the Attachment a stream passes to sandboxfs.Serve carries it. +// the Attachment each stream's sandboxfs.Server.Serve receives carries it. func (s *Service) InstanceID() sandboxwire.ID { return s.instance } // Close releases every attachment and the export root. Requests still @@ -290,9 +290,9 @@ type state struct { free []int inodes map[inodeKey]*node generation uint64 - handles map[sandboxfs.HandleID]*handle - reserved int - nextHandle uint64 + // handles maps each client-chosen ID to its handle, or to nil while the + // acquisition that reserved the ID runs. + handles map[sandboxfs.HandleID]*handle } // inodeKey identifies a node. The mount ID keeps the same inode reached @@ -349,21 +349,22 @@ type node struct { } type handle struct { - f *os.File - append bool - dir *cursor // set for directory handles + f *os.File + dir *cursor // set for directory handles + + writeMu sync.Mutex // holds a write and the O_APPEND mode it sets on f lockMu sync.Mutex flock sandboxfs.LockMode // held flock mode; zero when none } func newState(s *Service, readOnly bool) *state { - // Random bases make a NodeRef or HandleID from another attachment or + // A random generation base makes a NodeRef from another attachment or // incarnation miss instead of naming a live object. return &state{ svc: s, readOnly: readOnly, done: make(chan struct{}), inodes: map[inodeKey]*node{}, handles: map[sandboxfs.HandleID]*handle{}, - generation: rand.Uint64() >> 2, nextHandle: rand.Uint64() >> 2, + generation: rand.Uint64() >> 2, } } @@ -385,7 +386,9 @@ func (st *state) close() { } } for _, h := range handles { - h.f.Close() + if h != nil { + h.f.Close() + } } } @@ -472,40 +475,44 @@ func (st *state) forget(entries []sandboxfs.ForgetEntry) error { return nil } -// reserve claims a handle slot before a handle is opened, so exhaustion -// fails before anything happens. -func (st *state) reserve() error { +// reserve claims the client-chosen handle ID before an acquisition touches +// the file system, so a duplicate ID or exhaustion fails before anything +// happens. A reserved ID counts toward MaxOpenHandles. +func (st *state) reserve(id sandboxfs.HandleID) error { st.mu.Lock() defer st.mu.Unlock() if st.closed { return errStaleAttachment() } - if len(st.handles)+st.reserved >= int(st.svc.caps.MaxOpenHandles) { + if _, ok := st.handles[id]; ok { + return sandboxfs.NewFailure(sandboxfs.CodeInvalidArgument, sandboxwire.EffectNone, "handle ID is already in use") + } + if len(st.handles) >= int(st.svc.caps.MaxOpenHandles) { return sandboxfs.NewFailure(sandboxfs.CodeResourceExhausted, sandboxwire.EffectNone, "too many open handles") } - st.reserved++ + st.handles[id] = nil return nil } -func (st *state) unreserve() { +// unreserve drops the reservation of an acquisition that failed. +func (st *state) unreserve(id sandboxfs.HandleID) { st.mu.Lock() - st.reserved-- + if !st.closed && st.handles[id] == nil { + delete(st.handles, id) + } st.mu.Unlock() } -// addHandle registers h in a reserved slot. -func (st *state) addHandle(h *handle) (sandboxfs.HandleID, error) { +// publish installs h under its reserved ID. +func (st *state) publish(id sandboxfs.HandleID, h *handle) error { st.mu.Lock() defer st.mu.Unlock() - st.reserved-- if st.closed { h.f.Close() - return 0, errStaleAttachment() + return errStaleAttachment() } - st.nextHandle++ - id := sandboxfs.HandleID(st.nextHandle) st.handles[id] = h - return id, nil + return nil } func (st *state) handle(id sandboxfs.HandleID) (*handle, error) { diff --git a/apps/sandboxio/internal/fileservice/service_test.go b/apps/sandboxio/internal/fileservice/service_test.go index 8c3156e0..cda1b598 100644 --- a/apps/sandboxio/internal/fileservice/service_test.go +++ b/apps/sandboxio/internal/fileservice/service_test.go @@ -5,6 +5,7 @@ package fileservice import ( "bytes" "context" + "encoding/binary" "errors" "fmt" "math" @@ -16,6 +17,7 @@ import ( "strconv" "strings" "sync" + "sync/atomic" "syscall" "testing" "time" @@ -29,7 +31,9 @@ type fixture struct { t *testing.T dir string svc *Service + srv *sandboxfs.Server att sandboxfs.Attachment + ids sandboxfs.HandleIDs // the attachment's handle IDs, across its streams c *sandboxfs.Client root sandboxfs.NodeRef } @@ -47,8 +51,8 @@ func attachRoot(t *testing.T, dir string) *fixture { t.Fatal(err) } t.Cleanup(func() { svc.Close() }) - f := &fixture{t: t, dir: dir, svc: svc, att: attachment(svc, sandboxlink.ExportGrant{ID: "world"})} - f.c, _, _ = f.connect(f.svc, f.att) + f := &fixture{t: t, dir: dir, svc: svc, srv: sandboxfs.NewServer(svc), att: attachment(svc, sandboxlink.ExportGrant{ID: "world"})} + f.c, _, _ = f.connect(f.srv, f.att) resp, err := f.c.Attach(context.Background(), &sandboxfs.AttachRequest{Export: "world"}) if err != nil { t.Fatal(err) @@ -57,12 +61,16 @@ func attachRoot(t *testing.T, dir string) *fixture { return f } -// connect opens a stream to svc for attachment a. It returns the client, the +// binds stands in for Link's bind sequence: each stream is bound after every +// earlier one. +var binds atomic.Uint64 + +// connect opens a stream to srv for attachment a. It returns the client, the // server end of the stream and the channel Serve's result arrives on. -func (f *fixture) connect(svc *Service, a sandboxfs.Attachment) (*sandboxfs.Client, net.Conn, chan error) { +func (f *fixture) connect(srv *sandboxfs.Server, a sandboxfs.Attachment) (*sandboxfs.Client, net.Conn, chan error) { cc, sc := net.Pipe() served := make(chan error, 1) - go func() { served <- sandboxfs.Serve(context.Background(), sc, svc, a) }() + go func() { served <- srv.Serve(context.Background(), sc, a, binds.Add(1)) }() c := sandboxfs.NewClient(cc) f.t.Cleanup(func() { c.Close() }) return c, sc, served @@ -84,21 +92,53 @@ func (f *fixture) lookup(parent sandboxfs.NodeRef, name string) sandboxfs.Entry func (f *fixture) create(name string, access sandboxfs.AccessMode, flags sandboxfs.OpenFlags) (sandboxfs.Entry, sandboxfs.HandleID) { f.t.Helper() - r, err := f.c.Create(context.Background(), &sandboxfs.CreateRequest{Parent: f.root, Name: []byte(name), Mode: 0o640, Access: access, Flags: flags, Exclusive: true}) + h := f.ids.Next() + r, err := f.c.Create(context.Background(), &sandboxfs.CreateRequest{Handle: h, Parent: f.root, Name: []byte(name), Mode: 0o640, Access: access, Flags: flags, Exclusive: true}) if err != nil { f.t.Fatalf("create %q: %v", name, err) } - return r.Entry, r.Handle + return r.Entry, h +} + +// open opens the export's entry name for reading and writing. +func (f *fixture) open(name string) sandboxfs.HandleID { + f.t.Helper() + h := f.ids.Next() + if _, err := f.c.Open(context.Background(), &sandboxfs.OpenRequest{Handle: h, Node: f.lookup(f.root, name).Node, Access: sandboxfs.AccessReadWrite}); err != nil { + f.t.Fatalf("open %q: %v", name, err) + } + return h } func (f *fixture) write(h sandboxfs.HandleID, off uint64, data string) { f.t.Helper() - r, err := f.c.Write(context.Background(), &sandboxfs.WriteRequest{Handle: h, Offset: off, Data: []byte(data)}) - if err != nil || r.Written != uint32(len(data)) || r.Failure != nil { + f.writeRequest(&sandboxfs.WriteRequest{Handle: h, Offset: off, Data: []byte(data)}) +} + +// appendWrite writes data at the end of the file, whatever the offset. +func (f *fixture) appendWrite(h sandboxfs.HandleID, data string) { + f.t.Helper() + f.writeRequest(&sandboxfs.WriteRequest{Handle: h, Offset: 1 << 40, Append: true, Data: []byte(data)}) +} + +func (f *fixture) writeRequest(w *sandboxfs.WriteRequest) { + f.t.Helper() + r, err := f.c.Write(context.Background(), w) + if err != nil || r.Written != uint32(len(w.Data)) || r.Failure != nil { f.t.Fatalf("write: %+v, %v", r, err) } } +// handles counts the attachment's handles, including reserved ones. +func (f *fixture) handles() int { + f.svc.mu.Lock() + st := f.svc.atts[f.att.ID] + f.svc.mu.Unlock() + st.mu.Lock() + defer st.mu.Unlock() + return len(st.handles) +} + func (f *fixture) read(h sandboxfs.HandleID) string { f.t.Helper() r, err := f.c.Read(context.Background(), &sandboxfs.ReadRequest{Handle: h, Size: 4096}) @@ -152,30 +192,52 @@ func TestCreateWriteRead(t *testing.T) { func TestExclusiveCreateConflict(t *testing.T) { f := newFixture(t) f.create("x", sandboxfs.AccessWrite, 0) - _, err := f.c.Create(context.Background(), &sandboxfs.CreateRequest{Parent: f.root, Name: []byte("x"), Mode: 0o600, Access: sandboxfs.AccessWrite, Exclusive: true}) + _, err := f.c.Create(context.Background(), &sandboxfs.CreateRequest{Handle: f.ids.Next(), Parent: f.root, Name: []byte("x"), Mode: 0o600, Access: sandboxfs.AccessWrite, Exclusive: true}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoExists, sandboxwire.EffectNone) - if _, err := f.c.Create(context.Background(), &sandboxfs.CreateRequest{Parent: f.root, Name: []byte("x"), Mode: 0o600, Access: sandboxfs.AccessWrite}); err != nil { + if _, err := f.c.Create(context.Background(), &sandboxfs.CreateRequest{Handle: f.ids.Next(), Parent: f.root, Name: []byte("x"), Mode: 0o600, Access: sandboxfs.AccessWrite}); err != nil { t.Fatalf("non-exclusive create of an existing file: %v", err) } } +// Append is a property of each write, as fcntl(F_SETFL, O_APPEND) makes it +// on a native descriptor, so one handle mixes positioned and append writes. +func TestAppendIsPerWrite(t *testing.T) { + f := newFixture(t) + for name, appendFirst := range map[string]bool{"positioned-then-append": false, "append-then-positioned": true} { + if err := os.WriteFile(filepath.Join(f.dir, name), []byte("abc"), 0o600); err != nil { + t.Fatal(err) + } + h := f.open(name) + if appendFirst { + f.appendWrite(h, "Y") + f.write(h, 0, "X") + } else { + f.write(h, 0, "X") + f.appendWrite(h, "Y") + } + if got := f.read(h); got != "XbcY" { + t.Fatalf("%s: read %q", name, got) + } + } +} + func TestAtomicAppendFromTwoHandles(t *testing.T) { f := newFixture(t) - _, h1 := f.create("log", sandboxfs.AccessWrite, sandboxfs.OpenAppend) - open, err := f.c.Open(context.Background(), &sandboxfs.OpenRequest{Node: f.lookup(f.root, "log").Node, Access: sandboxfs.AccessWrite, Flags: sandboxfs.OpenAppend}) - if err != nil { + _, h1 := f.create("log", sandboxfs.AccessWrite, 0) + h2 := f.ids.Next() + if _, err := f.c.Open(context.Background(), &sandboxfs.OpenRequest{Handle: h2, Node: f.lookup(f.root, "log").Node, Access: sandboxfs.AccessWrite}); err != nil { t.Fatal(err) } const records = 200 var wg sync.WaitGroup - for i, h := range []sandboxfs.HandleID{h1, open.Handle} { + for i, h := range []sandboxfs.HandleID{h1, h2} { wg.Add(1) go func() { defer wg.Done() for j := range records { // Every write names offset zero; append ignores it. rec := fmt.Sprintf("%d:%03d:%s\n", i, j, strings.Repeat("x", 64)) - if r, err := f.c.Write(context.Background(), &sandboxfs.WriteRequest{Handle: h, Data: []byte(rec)}); err != nil || int(r.Written) != len(rec) { + if r, err := f.c.Write(context.Background(), &sandboxfs.WriteRequest{Handle: h, Append: true, Data: []byte(rec)}); err != nil || int(r.Written) != len(rec) { t.Errorf("append: %+v, %v", r, err) return } @@ -257,13 +319,13 @@ func TestReadDirPagesWithCookies(t *testing.T) { t.Fatal(err) } } - dir, err := f.c.OpenDir(context.Background(), &sandboxfs.OpenDirRequest{Node: f.root}) - if err != nil { + dir := f.ids.Next() + if _, err := f.c.OpenDir(context.Background(), &sandboxfs.OpenDirRequest{Handle: dir, Node: f.root}); err != nil { t.Fatal(err) } readAll := func(cookie uint64) (names []string, cookies []uint64) { for { - r, err := f.c.ReadDir(context.Background(), &sandboxfs.ReadDirRequest{Handle: dir.Handle, Cookie: cookie, Limit: 200}) + r, err := f.c.ReadDir(context.Background(), &sandboxfs.ReadDirRequest{Handle: dir, Cookie: cookie, Limit: 200}) if err != nil { t.Fatal(err) } @@ -290,7 +352,7 @@ func TestReadDirPagesWithCookies(t *testing.T) { t.Fatalf("resumed at cookie 10: %v, want %v", again, names[10:]) } - r, err := f.c.ReadDir(context.Background(), &sandboxfs.ReadDirRequest{Handle: dir.Handle, Limit: 1024, WithAttrs: true}) + r, err := f.c.ReadDir(context.Background(), &sandboxfs.ReadDirRequest{Handle: dir, Limit: 1024, WithAttrs: true}) if err != nil || len(r.Entries) == 0 || r.Entries[0].Entry == nil || r.Entries[0].Entry.Attr.Mode&sandboxfs.ModeType != sandboxfs.ModeRegular { t.Fatalf("readdir with attributes: %+v, %v", r, err) } @@ -337,7 +399,7 @@ func TestRootEscapeIsRefused(t *testing.T) { // A client that skips validation gets the same answer from the server. cc, sc := net.Pipe() - go sandboxfs.Serve(context.Background(), sc, f.svc, f.att) + go f.srv.Serve(context.Background(), sc, attachment(f.svc, sandboxlink.ExportGrant{ID: "world"}), binds.Add(1)) defer cc.Close() var e sandboxwire.Encoder e.U64(f.root.ID) @@ -358,9 +420,9 @@ func TestRootEscapeIsRefused(t *testing.T) { link := f.lookup(f.root, "escape").Node _, err = f.c.Lookup(context.Background(), &sandboxfs.LookupRequest{Parent: link, Name: []byte("etc")}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoNotDirectory, sandboxwire.EffectNone) - _, err = f.c.Open(context.Background(), &sandboxfs.OpenRequest{Node: link, Access: sandboxfs.AccessRead}) + _, err = f.c.Open(context.Background(), &sandboxfs.OpenRequest{Handle: f.ids.Next(), Node: link, Access: sandboxfs.AccessRead}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoSymlinkLoop, sandboxwire.EffectNone) - _, err = f.c.OpenDir(context.Background(), &sandboxfs.OpenDirRequest{Node: link}) + _, err = f.c.OpenDir(context.Background(), &sandboxfs.OpenDirRequest{Handle: f.ids.Next(), Node: link}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoNotDirectory, sandboxwire.EffectNone) _, err = f.c.Mkdir(context.Background(), &sandboxfs.MkdirRequest{Parent: link, Name: []byte("x"), Mode: 0o755}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoNotDirectory, sandboxwire.EffectNone) @@ -421,14 +483,15 @@ func TestIncarnationChange(t *testing.T) { defer next.Close() // A stream bound to the old incarnation. - old, _, _ := f.connect(next, f.att) + nextSrv := sandboxfs.NewServer(next) + old, _, _ := f.connect(nextSrv, f.att) _, err = old.Lookup(context.Background(), &sandboxfs.LookupRequest{Parent: f.root, Name: []byte("f")}) wantFailure(t, err, sandboxfs.CodeInstanceChanged, 0, sandboxwire.EffectNone) // The same attachment rebound to the new incarnation holds nothing. a := f.att a.ServerInstanceID = next.InstanceID() - c, _, _ := f.connect(next, a) + c, _, _ := f.connect(nextSrv, a) _, err = c.Lookup(context.Background(), &sandboxfs.LookupRequest{Parent: f.root, Name: []byte("f")}) wantFailure(t, err, sandboxfs.CodeStaleAttachment, 0, sandboxwire.EffectNone) if _, err := c.Attach(context.Background(), &sandboxfs.AttachRequest{Export: "world"}); err != nil { @@ -444,24 +507,24 @@ func TestAttachFollowsExportGrants(t *testing.T) { f := newFixture(t) ctx := context.Background() - c, _, _ := f.connect(f.svc, attachment(f.svc, sandboxlink.ExportGrant{ID: "logs"})) + c, _, _ := f.connect(f.srv, attachment(f.svc, sandboxlink.ExportGrant{ID: "logs"})) _, err := c.Attach(ctx, &sandboxfs.AttachRequest{Export: "world"}) wantFailure(t, err, sandboxfs.CodeUnauthorized, 0, sandboxwire.EffectNone) - c, _, _ = f.connect(f.svc, attachment(f.svc, sandboxlink.ExportGrant{ID: "world", ReadOnly: true})) + c, _, _ = f.connect(f.srv, attachment(f.svc, sandboxlink.ExportGrant{ID: "world", ReadOnly: true})) _, err = c.Attach(ctx, &sandboxfs.AttachRequest{Export: "world"}) wantFailure(t, err, sandboxfs.CodeUnauthorized, 0, sandboxwire.EffectNone) r, err := c.Attach(ctx, &sandboxfs.AttachRequest{Export: "world", ReadOnly: true}) if err != nil { t.Fatal(err) } - _, err = c.Create(ctx, &sandboxfs.CreateRequest{Parent: r.Root.Node, Name: []byte("f"), Mode: 0o644, Access: sandboxfs.AccessWrite}) + _, err = c.Create(ctx, &sandboxfs.CreateRequest{Handle: 1, Parent: r.Root.Node, Name: []byte("f"), Mode: 0o644, Access: sandboxfs.AccessWrite}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoReadOnlyFilesystem, sandboxwire.EffectNone) } func TestDescribeListsGrantedExports(t *testing.T) { f := newFixture(t) - c, _, _ := f.connect(f.svc, attachment(f.svc, sandboxlink.ExportGrant{ID: "logs"})) + c, _, _ := f.connect(f.srv, attachment(f.svc, sandboxlink.ExportGrant{ID: "logs"})) d, err := c.Describe(context.Background(), &sandboxfs.DescribeRequest{}) if err != nil || len(d.Exports) != 0 { t.Fatalf("describe without a world grant: %+v, %v", d, err) @@ -480,9 +543,9 @@ func TestProcMagicLinkIsOpaque(t *testing.T) { cwd := w.Entries[2].Node _, err = f.c.Lookup(ctx, &sandboxfs.LookupRequest{Parent: cwd, Name: []byte("x")}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoNotDirectory, sandboxwire.EffectNone) - _, err = f.c.OpenDir(ctx, &sandboxfs.OpenDirRequest{Node: cwd}) + _, err = f.c.OpenDir(ctx, &sandboxfs.OpenDirRequest{Handle: f.ids.Next(), Node: cwd}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoNotDirectory, sandboxwire.EffectNone) - _, err = f.c.Open(ctx, &sandboxfs.OpenRequest{Node: cwd, Access: sandboxfs.AccessRead}) + _, err = f.c.Open(ctx, &sandboxfs.OpenRequest{Handle: f.ids.Next(), Node: cwd, Access: sandboxfs.AccessRead}) wantFailure(t, err, sandboxfs.CodeErrno, sandboxfs.ErrnoSymlinkLoop, sandboxwire.EffectNone) } @@ -564,7 +627,7 @@ func TestTransportLoss(t *testing.T) { // A second stream of the same attachment reuses its handles. cc, sc := net.Pipe() served := make(chan error, 1) - go func() { served <- sandboxfs.Serve(context.Background(), sc, f.svc, f.att) }() + go func() { served <- f.srv.Serve(context.Background(), sc, f.att, binds.Add(1)) }() wrote := make(chan struct{}, 1) c := sandboxfs.NewClient(wroteConn{cc, wrote}) defer c.Close() @@ -588,3 +651,62 @@ func TestTransportLoss(t *testing.T) { t.Fatal("Serve did not return after the transport closed") } } + +// dropConn closes the stream instead of writing the first frame of type typ, +// as a relay that fails after the service ran the request does. +type dropConn struct { + net.Conn + typ uint16 +} + +func (c dropConn) Write(p []byte) (int, error) { + if len(p) >= sandboxwire.HeaderSize && binary.BigEndian.Uint16(p[4:6]) == c.typ { + c.Conn.Close() + return 0, net.ErrClosed + } + return c.Conn.Write(p) +} + +// A client that lost an Open's reply releases the ID it chose on the next +// stream, and no handle remains. +func TestLostOpenReplyIsReleased(t *testing.T) { + f := newFixture(t) + ctx := context.Background() + if err := os.WriteFile(filepath.Join(f.dir, "f"), nil, 0o600); err != nil { + t.Fatal(err) + } + node := f.lookup(f.root, "f").Node + cc, sc := net.Pipe() + go f.srv.Serve(ctx, dropConn{sc, sandboxwire.ResponseType(uint16(sandboxfs.OpOpen))}, f.att, binds.Add(1)) + c := sandboxfs.NewClient(cc) + defer c.Close() + h := f.ids.Next() + _, err := c.Open(ctx, &sandboxfs.OpenRequest{Handle: h, Node: node, Access: sandboxfs.AccessRead}) + if fail := wantFailure(t, err, sandboxfs.CodeUnknown, 0, sandboxwire.EffectPossible); !errors.Is(fail, sandboxfs.ErrTransport) { + t.Fatalf("%v is not a transport failure", fail) + } + if n := f.handles(); n != 1 { + t.Fatalf("%d handles after the lost reply", n) + } + + resumed, _, _ := f.connect(f.srv, f.att) + if _, err := resumed.Release(ctx, &sandboxfs.ReleaseRequest{Handle: h}); err != nil { + t.Fatalf("release after resuming: %v", err) + } + if n := f.handles(); n != 0 { + t.Fatalf("%d handles after the release", n) + } +} + +// An acquisition that names a live handle ID is refused before it touches the +// file system. +func TestDuplicateHandleIDIsRefused(t *testing.T) { + f := newFixture(t) + _, h := f.create("a", sandboxfs.AccessReadWrite, 0) + _, err := f.c.Create(context.Background(), &sandboxfs.CreateRequest{Handle: h, Parent: f.root, Name: []byte("b"), Mode: 0o600, Access: sandboxfs.AccessWrite}) + wantFailure(t, err, sandboxfs.CodeInvalidArgument, 0, sandboxwire.EffectNone) + if _, err := os.Lstat(filepath.Join(f.dir, "b")); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("the refused Create made its file: %v", err) + } + f.write(h, 0, "still open") +} diff --git a/apps/sandboxio/internal/sandboxio/sandboxio.go b/apps/sandboxio/internal/sandboxio/sandboxio.go index 75e8cfe4..1615235e 100644 --- a/apps/sandboxio/internal/sandboxio/sandboxio.go +++ b/apps/sandboxio/internal/sandboxio/sandboxio.go @@ -85,6 +85,7 @@ func run(ctx context.Context, bootstrapPath string, opt options) error { return &StartupError{StepFileService, err} } defer files.Close() + fileServer := sandboxfs.NewServer(files) // down records that a link attempt failed since the last accepted Hello, // so each drop is logged once rather than on every reconnect attempt. @@ -98,8 +99,8 @@ func run(ctx context.Context, bootstrapPath string, opt options) error { // incarnation, so the link announces that one. ServerInstanceID: files.InstanceID(), Services: []sandboxlink.ServiceHandler{ - {Service: sandboxlink.ServiceFile, Version: sandboxfs.Version, Serve: func(ctx context.Context, b sandboxlink.Bind, _ uint64, s sandboxlink.Stream) { - sandboxfs.Serve(ctx, s, files, sandboxfs.Attachment{ID: b.AttachmentID, ServerInstanceID: b.ExpectedServerInstanceID, Lease: ctx, Exports: b.Exports}) + {Service: sandboxlink.ServiceFile, Version: sandboxfs.Version, Serve: func(ctx context.Context, b sandboxlink.Bind, seq uint64, s sandboxlink.Stream) { + fileServer.Serve(ctx, s, sandboxfs.Attachment{ID: b.AttachmentID, ServerInstanceID: b.ExpectedServerInstanceID, Lease: ctx, Exports: b.Exports}, seq) }}, {Service: sandboxlink.ServiceProcess, Version: sandboxprocess.Version, Serve: func(ctx context.Context, b sandboxlink.Bind, _ uint64, s sandboxlink.Stream) { sandboxprocess.Serve(ctx, s, sandboxprocess.Attachment{ID: b.AttachmentID}, procs) diff --git a/apps/sandboxio/internal/sandboxio/sandboxio_test.go b/apps/sandboxio/internal/sandboxio/sandboxio_test.go index b682b4e3..8a9843de 100644 --- a/apps/sandboxio/internal/sandboxio/sandboxio_test.go +++ b/apps/sandboxio/internal/sandboxio/sandboxio_test.go @@ -153,14 +153,15 @@ func TestServesEachProtocolThroughTheRelay(t *testing.T) { if err != nil { t.Fatal(err) } - created, err := files.Create(ctx, &sandboxfs.CreateRequest{Parent: attached.Root.Node, Name: []byte("note"), Mode: 0o644, Access: sandboxfs.AccessReadWrite, Exclusive: true}) - if err != nil { + var handles sandboxfs.HandleIDs + note := handles.Next() + if _, err := files.Create(ctx, &sandboxfs.CreateRequest{Handle: note, Parent: attached.Root.Node, Name: []byte("note"), Mode: 0o644, Access: sandboxfs.AccessReadWrite, Exclusive: true}); err != nil { t.Fatal(err) } - if w, err := files.Write(ctx, &sandboxfs.WriteRequest{Handle: created.Handle, Data: []byte("hello")}); err != nil || w.Written != 5 { + if w, err := files.Write(ctx, &sandboxfs.WriteRequest{Handle: note, Data: []byte("hello")}); err != nil || w.Written != 5 { t.Fatalf("write: %+v, %v", w, err) } - if r, err := files.Read(ctx, &sandboxfs.ReadRequest{Handle: created.Handle, Size: 64}); err != nil || string(r.Data) != "hello" { + if r, err := files.Read(ctx, &sandboxfs.ReadRequest{Handle: note, Size: 64}); err != nil || string(r.Data) != "hello" { t.Fatalf("read: %+v, %v", r, err) } if b, err := os.ReadFile(filepath.Join(root, "note")); err != nil || string(b) != "hello" { diff --git a/docs/file-access-protocol.md b/docs/file-access-protocol.md index c3123749..5d41a581 100644 --- a/docs/file-access-protocol.md +++ b/docs/file-access-protocol.md @@ -1,6 +1,6 @@ # File access protocol -The File access protocol is how a Runtime reads and changes the files of a sandbox. The Sandbox I/O service in the sandbox serves it, and the Runtime is its client. It is a node and handle protocol shaped like the FUSE low-level operations: lookups acquire node references, opens return handles, reads and writes take offsets, directory reads resume at cookies, and locks work against the sandbox's own processes. Phase 1 serves the [Uncached](#uncached-profile) profile only, with no change stream. +The File access protocol is how a Runtime reads and changes the files of a sandbox. The Sandbox I/O service in the sandbox serves it, and the Runtime is its client. It is a node and handle protocol shaped like the FUSE low-level operations: lookups acquire node references, opens create handles under IDs the client chooses, reads and writes take offsets, directory reads resume at cookies, and locks work against the sandbox's own processes. Phase 1 serves the [Uncached](#uncached-profile) profile only, with no change stream. [`internal/sandboxfs/protocol.go`](../internal/sandboxfs/protocol.go) is the authored definition: message tags, payload layouts, validators and the `Service` interface. The same package holds the generic client and server. [`apps/sandboxio/internal/fileservice`](../apps/sandboxio/internal/fileservice) is the Linux service. Frames use the shared [framing](sandbox-link-protocol.md#framing), and the Link layer supplies the authenticated attachment of each stream. @@ -8,7 +8,7 @@ The File access protocol is how a Runtime reads and changes the files of a sandb - A stream belongs to one attachment, which Link authenticates and hands to the server as a `sandboxfs.Attachment`: its ID, the `ServerInstanceID` the stream was bound to, its lease and the exports it is granted. No request names an attachment, an OS user or a credential. - Request IDs follow the [framing](sandbox-link-protocol.md#framing) rule, and a response carries its request's RequestID. Requests on one stream run concurrently, so responses can arrive in any order. The protocol has no events. -- An attachment attaches to one export. Its node references, handles and locks live in the service until it detaches, its lease ends or the service restarts. A new stream of the same attachment in the same incarnation continues with them. +- An attachment attaches to one export. Its node references, handles and locks live in the service until it detaches, its lease ends or the service restarts. A new stream of the same attachment in the same incarnation continues with them, after the [succession fence](#stream-succession). - A service incarnation is named by `ServerInstanceID`. A request on a stream bound to another incarnation fails with `InstanceChanged`, and nothing is reopened automatically. ## Implement a client @@ -19,10 +19,12 @@ The Go client is `sandboxfs.NewClient(stream)`. It has one method per operation, 2. `Attach` an export and keep the root `NodeRef`. 3. `Lookup`, `Walk`, `Create`, `Mkdir`, `Symlink`, `Link` and `ReadDir` with `WithAttrs` each acquire one reference on every node they return. Release references with `Forget` when the kernel forgets them. 4. `Walk` stops after a symlink. Resolve the link yourself, relative to the view, with `Readlink` and further walks. -5. `Open`, `Create` and `OpenDir` return handles. Send `Flush` on each close of a descriptor for the handle, and `Release` or `ReleaseDir` when the last one closes. -6. Read the `Effect` of every failure. After `EffectPossible`, the request may have taken effect: never replay a mutation automatically. Report the failure, or inspect the state with `GetAttr` or `Lookup` first. +5. `Open`, `Create` and `OpenDir` open a handle under a [handle ID](#handles) the client chooses. Keep one `sandboxfs.HandleIDs` per attachment, across all of its streams, and take each ID from it. Send `Flush` on each close of a descriptor for the handle, and `Release` or `ReleaseDir` when the last one closes. +6. Set `Append` on each `Write` made while the descriptor is in append mode. Append is a property of the write, not of the handle. +7. Read the `Effect` of every failure. After `EffectPossible`, the request may have taken effect: never replay a mutation automatically. Report the failure, or inspect the state with `GetAttr` or `Lookup` first. A failure with a [retryable](#failures) code and `EffectNone` may be resent unchanged. +8. Clean up an acquisition whose outcome is unknown, because its reply was lost or its call was cancelled, with `Release` or `ReleaseDir` of its ID. After a stream fails, resume on a new stream of the same attachment, which the service serves only after the failed stream's requests have finished, and clean up there: an uncertain `Attach` with `Detach`, and each uncertain acquisition with `Release` or `ReleaseDir`. [Handles](#handles) says what the cleanup proves. -Cancelling a call's context returns at once with `Cancelled` or `DeadlineExceeded`. A call cancelled before its request is written fails with `EffectNone`. After the request is written, the client sends `CancelRequest` for it and discards the late response, and the failure is `EffectNone` for a [side-effect-free request](#effects-and-cancellation) and `EffectPossible` for every other one. Cancelling a call while its request is being written fails the stream instead, because a partial frame cannot be withdrawn. A cancel that races the completion of a write may still fail the stream, and the requests in flight then fail with `EffectPossible`. +Cancelling a call's context returns at once. A call cancelled before its request is written fails with `Cancelled` or `DeadlineExceeded` and `EffectNone`. After the request is written, the client sends `CancelRequest` for it and discards the late response, and the call fails with `Cancelled` or `DeadlineExceeded`, with `EffectNone` for a [side-effect-free request](#effects-and-cancellation) and `EffectPossible` for every other one. Cancelling a call while its request is being written fails the stream instead, because a partial frame cannot be withdrawn: the call returns at once, without waiting for the transport to finish the write or close the stream, with the transport failure, `Unknown` and `EffectPossible`, which matches `sandboxfs.ErrTransport`. A cancel that races the completion of a write may still fail the stream, and the requests in flight then fail the same way. To learn what a cancelled request did, interrupt it instead: a call whose context comes from `sandboxfs.WithInterrupt` sends `CancelRequest` when the interrupt channel closes and keeps waiting for the request's own response, so it returns the request's result or the service's failure with its effect. A FUSE frontend uses this for a waiting lock, which may be acquired just before the cancellation arrives. @@ -30,13 +32,15 @@ When the stream fails, every request in flight fails with `Unknown` and `EffectP ## Implement a service -Implement `sandboxfs.Service` and serve each stream with `sandboxfs.Serve(ctx, stream, service, attachment)`. `Serve`: +Implement `sandboxfs.Service`, create one `sandboxfs.NewServer(service)`, and serve every stream with `server.Serve(ctx, stream, attachment, seq)`, where `seq` is the stream's [Link bind sequence](sandbox-link-protocol.md#implement-a-serve-peer). One `Server` serves all of a service's streams, because the [succession fence](#stream-succession) spans them. `Serve`: -- refuses an attachment without an ID, a `ServerInstanceID`, a lease or at least one valid export grant with unique IDs; +- refuses an attachment without an ID, a `ServerInstanceID`, a lease or at least one valid export grant with unique IDs, and returns `sandboxfs.ErrLeaseEnded`, without dispatching anything, when the attachment's lease ended before the stream was admitted; +- admits one stream per `(ServerInstanceID, AttachmentID)` at a time, in bind order: a successor stream reads its requests at once but hands none to the service, without a deadline, until every earlier stream that handed one over has drained. The superseded `Serve` returns `sandboxfs.ErrSuperseded`, at once when it has handed nothing over, and so does the `Serve` of a stream bound before one already admitted, without dispatching any of its requests; - decodes and validates each request, and answers a malformed payload with `InvalidArgument` and `EffectNone`; - ends the stream on a framing violation: an unknown tag, a frame that is not a request, or a RequestID that does not increase; - holds up to `sandboxfs.MaxInFlight` (256) requests, each from admission until its response is written, and answers any more with `ResourceExhausted` and `EffectNone`, so a `CancelRequest` arrives while the client reads responses; - answers `CancelRequest` itself by cancelling the target's context; +- refuses an `Open`, `Create` or `OpenDir` whose handle ID another acquisition on the stream is still using, with `InvalidArgument` and `EffectNone`, and runs a `Release` or `ReleaseDir` of an ID only after the running acquisition of that ID has finished. Together with the succession fence, a service never runs two acquisitions of one ID at once, or a release concurrently with the acquisition of its ID; - returns a method's `*Failure` as the typed failure, and reports any other error, or a response that fails validation, as `Unknown` with `EffectPossible`; - cancels every request's context when the stream ends. @@ -44,6 +48,7 @@ A service must: - generate a new `ServerInstanceID` whenever it loses its node and handle tables, and answer a request whose `Attachment.ServerInstanceID` is not its own with `InstanceChanged`; - answer `StaleAttachment` while the attachment is not attached or after its lease ends, and `StaleNode` or `StaleHandle` for a reference or handle the attachment does not hold. A node ID is reused only with a new generation; +- reserve the handle ID of an `Open`, `Create` or `OpenDir` atomically before any file-system effect, as [Handles](#handles) describes, and publish every state a request creates before its method returns; - call `Capabilities.Admit(request, readOnly)` before running a request and return the failure it reports. Limits that depend on service state, such as `MaxOpenHandles`, stay with the service; - advertise only what it enforces, and advertise locks only when they interoperate with native processes in the sandbox; - act as its own process identity, never as an identity a request supplies, and apply requested permission bits exactly; @@ -57,7 +62,7 @@ A service must: - Each node holds an `O_PATH|O_NOFOLLOW` descriptor. In an attachment a node is one mount ID, device and inode, so hard links share a node while a bind mount and its source stay two. The mount ID comes from `statx` with `STATX_MNT_ID`, or from the `mnt_id` line of `/proc/self/fdinfo/` on kernels older than 5.8. - A lookup opens one component with `openat` and `O_NOFOLLOW` on its parent's descriptor. A symlink, including a proc magic link such as `/proc//cwd`, is a node of its own and is never traversed: `Lookup` and `Readlink` return the link itself, a directory operation on it fails with `Errno` `NotDirectory`, and `Open` fails with `SymlinkLoop`. - Operations on a node, such as opening, truncating, changing its mode or times and linking it, go through the `/proc/self/fd` name of the descriptor the service holds for it. That name resolves to the descriptor's own object, never to a symlink's target, and an open through it first checks the object's file type. -- An open handle holds its own descriptor, so it keeps working after its file is unlinked or renamed. +- An open handle holds its own descriptor, so it keeps working after its file is unlinked or renamed. Files are opened without `O_APPEND`. A handle's writes run one at a time, and each first sets or clears the descriptor's `O_APPEND` to match its `Append`; an append is then one `write` call. The service does not use `pwritev2` with `RWF_APPEND`, which overlayfs on some kernels drops, writing at the offset instead. A short append is reported as it is and never continued by another append. - `Rename` uses `renameat2`. The service declares `RenameNoReplace` and `RenameExchange` only when a probe at start succeeds. - `LockFlock` locks the handle's descriptor, so it interoperates with native `flock`. `POSIXLocks` is false, and `GetLock` and `LockPOSIX` return `Unsupported`. - `ReadDir` cookies are the kernel's directory offsets. `Attr.Ino` combines the device and the inode number as go-fuse's loopback does. @@ -79,14 +84,14 @@ Requests use tags 1 to 30; the response to tag `t` uses `t | 0x8000`. | 6 | `GetAttr` | `Target` | `Attr` | Attributes of a node or open handle | | 7 | `SetAttr` | `Target`, `Set`, selected values | `Attr` | Change the selected attributes | | 8 | `Access` | `Node`, `Mask` | – | Check permissions as the service's identity | -| 9 | `Open` | `Node`, `Access`, `Flags` | `Handle` | Open a regular file | -| 10 | `Create` | `Parent`, `Name`, `Mode`, `Access`, `Flags`, `Exclusive` | `Entry`, `Handle` | Create and open a regular file | +| 9 | `Open` | `Handle`, `Node`, `Access`, `Flags` | – | Open a regular file as `Handle` | +| 10 | `Create` | `Handle`, `Parent`, `Name`, `Mode`, `Access`, `Flags`, `Exclusive` | `Entry` | Create and open a regular file as `Handle` | | 11 | `Read` | `Handle`, `Offset`, `Size` | `Data` | Fewer bytes than asked means end of file | -| 12 | `Write` | `Handle`, `Offset`, `Data` | `Written`, optional `Failure` | Write at `Offset`, or append on an append-opened handle | +| 12 | `Write` | `Handle`, `Offset`, `Append`, `Data` | `Written`, optional `Failure` | Write at `Offset`, or at the end of the file when `Append` is set | | 13 | `Flush` | `Handle`, `Owner` | – | One descriptor for the handle closed | | 14 | `Fsync` | `Handle`, `DataOnly` | – | Make the file's data, and its metadata unless `DataOnly`, durable | | 15 | `Release` | `Handle` | – | Close a file handle | -| 16 | `OpenDir` | `Node` | `Handle` | Open a directory | +| 16 | `OpenDir` | `Handle`, `Node` | – | Open a directory as `Handle` | | 17 | `ReadDir` | `Handle`, `Cookie`, `Limit`, `WithAttrs` | `Entries`, `End` | Read entries after `Cookie` | | 18 | `ReleaseDir` | `Handle` | – | Close a directory handle | | 19 | `Mkdir` | `Parent`, `Name`, `Mode` | `Entry` | Create a directory | @@ -119,19 +124,17 @@ GetAttrResponse Attr SetAttr Target, Set u32, then for each selected bit in order: Size u64, Mode u32, UID u32, GID u32, Atime Timestamp, Mtime Timestamp SetAttrResponse Attr Access Node NodeRef, Mask u32; response (no fields) -Open Node NodeRef, Access enum, Flags u32 -OpenResponse Handle u64 -Create Parent NodeRef, Name bytes, Mode u32, Access enum, Flags u32, Exclusive bool -CreateResponse Entry, Handle u64 +Open Handle u64, Node NodeRef, Access enum, Flags u32; response (no fields) +Create Handle u64, Parent NodeRef, Name bytes, Mode u32, Access enum, Flags u32, Exclusive bool +CreateResponse Entry Read Handle u64, Offset u64, Size u32 ReadResponse Data bytes -Write Handle u64, Offset u64, Data bytes +Write Handle u64, Offset u64, Append bool, Data bytes WriteResponse Written u32, Failure optional Failure Flush Handle u64, Owner u64; response (no fields) Fsync Handle u64, DataOnly bool; response (no fields) Release Handle u64; response (no fields) -OpenDir Node NodeRef -OpenDirResponse Handle u64 +OpenDir Handle u64, Node NodeRef; response (no fields) ReadDir Handle u64, Cookie u64, Limit u32, WithAttrs bool ReadDirResponse Entries count of DirEntry, End bool ReleaseDir Handle u64; response (no fields) @@ -154,13 +157,13 @@ SetLock Handle u64, Kind enum, Owner u64, Lock, Wait bool; response CancelRequest Target u64; response (no fields) ``` -[`testdata`](../internal/sandboxfs/testdata) holds annotated golden frames of `Describe`, `Walk`, `Create`, a short `Write`, `ReadDir` with a cookie, `Rename` and a failure. +[`testdata`](../internal/sandboxfs/testdata) holds annotated golden frames of `Describe`, `Walk`, `Create`, an append `Write`, a short `Write`, `ReadDir` with a cookie, `Rename` and a failure. ### Shared types ```text NodeRef ID u64, Generation u64 // both nonzero -HandleID u64 // nonzero, opaque +HandleID u64 // nonzero, chosen by the client Timestamp Sec i64, Nsec u32 // Nsec below one billion Attr Ino u64, Mode u32, Nlink u32, UID u32, GID u32, Rdev u64, Size u64, Blocks u64, Blksize u32, Atime Timestamp, Mtime Timestamp, Ctime Timestamp Entry Node NodeRef, Attr @@ -192,7 +195,7 @@ Lock Mode enum (LockRead = 1, LockWrite = 2, LockUnlock = 3), Start u64 | `MaxReadDirBytes` | u32 | Largest `ReadDir` limit, 1 to 256 KiB | | `MaxOpenHandles` | u32 | Most open handles per attachment, at least 1 | | `ReadOnly` | bool | The service accepts only read-only attachments | -| `AtomicAppend` | bool | Appends from several handles never interleave within a write | +| `AtomicAppend` | bool | `Write` supports `Append`, and appends from several handles never interleave within a write | | `AtomicRename` | bool | `RenameReplace` replaces the destination atomically | | `RenameNoReplace`, `RenameExchange` | bool | The rename mode is supported | | `HardLinks`, `Symlinks` | bool | `Link` and `Symlink` are supported | @@ -237,14 +240,29 @@ Unknown bits are rejected, `AttrAtime` excludes `AttrAtimeNow` and `AttrMtime` e ### Files - `AccessMode` is `AccessRead` (1), `AccessWrite` (2) or `AccessReadWrite` (3). -- `OpenFlags` are `OpenAppend` (1), `OpenTruncate` (2), `OpenNoFollow` (4), `OpenSync` (8) and `OpenDataSync` (16); unknown flags are rejected. +- `OpenFlags` are `OpenTruncate` (1), `OpenNoFollow` (2), `OpenSync` (4) and `OpenDataSync` (8); unknown flags are rejected. - `Open` opens a regular file node. A node is an object, not a path, so `Open` never follows a symlink: a directory fails with `Errno` `IsDirectory`, a symlink with `SymlinkLoop`, and a special file with `Unsupported`. - `Create` creates a regular file with exactly the requested permission bits; no umask applies. With `Exclusive`, an existing entry fails with `Errno` `Exists`. Without it, an existing regular file is opened, and `OpenTruncate` truncates it. -- `Read` takes an offset up to 2^63−1. A positioned `Write` must end at or before 2^63−1, or it fails with `InvalidArgument`. A handle opened with `OpenAppend` appends each write atomically and ignores `Offset`. +- `Read` takes an offset up to 2^63−1. A positioned `Write` must end at or before 2^63−1, or it fails with `InvalidArgument`. +- A `Write` with `Append` ignores `Offset` and writes its data atomically at the end of the file. Any handle that may write takes both kinds of write, in any order, as a native descriptor does when `fcntl` sets or clears `O_APPEND`. A service without `AtomicAppend` refuses `Append` with `Unsupported`. - A successful `WriteResponse` carries the exact number of bytes written. When a write stops after a nonzero prefix, `Written` counts the prefix and `Failure` says why it stopped; a write that wrote nothing is a failure response. - `Flush` is neither `Fsync` nor `Release`. It reports the errors of closing one descriptor for the handle, and `Owner` names the closing lock owner. - A file operation on a directory handle, or a directory operation on a file handle, fails with `Errno` `IsDirectory`, `NotDirectory` or `BadDescriptor`. +### Handles + +- The client chooses the `HandleID` of each `Open`, `Create` and `OpenDir`: nonzero, and never used before in the attachment, on any of its streams. `sandboxfs.HandleIDs` allocates IDs in increasing order. +- The service reserves the ID before the request has any file-system effect. A reserved ID counts toward `MaxOpenHandles` while its acquisition runs. An acquisition whose ID is reserved or open fails with `InvalidArgument` and `EffectNone`. +- Other requests that name a reserved ID fail with `StaleHandle`, except `Release` and `ReleaseDir`: the server runs them only after the acquisition of that ID has finished, so they close the handle it opened. +- `Release` or `ReleaseDir` of an ID settles an acquisition whose outcome is unknown, whether its reply was lost with the stream or its call was cancelled. Success or `StaleHandle` proves that no handle with that ID remains. It does not prove that the acquisition changed nothing: a `Create` may have created its file, and a truncating `Open` may have truncated it. +- A node reference acquired by a `Create`, `Lookup`, `Walk`, `Mkdir`, `Symlink`, `Link` or `ReadDir` with `WithAttrs` cannot be forgotten when its reply never reaches the client, or reaches a call the client abandoned, because the client never learns the `NodeRef`. It stays held, with the descriptor the service keeps for its node, until the attachment detaches or its lease ends. Their number has no fixed bound: a stream failure loses every reply not yet delivered, including replies the server has written and no longer counts against `MaxInFlight`, and one `Walk` or `ReadDir` acquires a reference for each node it returns. A client leaves no reference to an abandoned call when it ends the context of such a call only before detaching, and otherwise stops it with `sandboxfs.WithInterrupt`, waits for its outcome and forgets what it returns. + +### Stream succession + +For each `(ServerInstanceID, AttachmentID)` the server admits one File stream at a time. Before it dispatches any request of a successor stream, it stops admission on the predecessor, closes and cancels the predecessor's requests, and waits for every admitted handler and state-publication task to finish. The predecessor's replies are discarded. The gate exists before `Attach` creates any state. + +Succession follows Link's bind order, not the order in which streams reach the server: a stream bound before one already admitted is refused before it dispatches anything. The server forgets an attachment's streams only once its lease has ended and every stream it admitted has drained, and refuses every later stream under that lease with `sandboxfs.ErrLeaseEnded`, so the refusal still holds then. A successor reads and admits requests while it waits, but dispatching means handing a request to the service, and it dispatches none until then. A successor superseded while it waits has dispatched nothing, so its `Serve` ends at once, and the next successor waits for the streams it was waiting for. No timeout ends the wait, because a handler may still change state after its cancellation. Effects that completed stay. A request the successor sends therefore sees everything its predecessor's requests did: a `Detach` cleans up an `Attach` whose reply was lost, and a `Release` cleans up a lost acquisition. + ### Directories - `OpenDir` opens a directory node, and `ReadDir` reads its entries after `Cookie`; cookie 0 is the start. Each `DirEntry.Cookie` is the position after that entry and is otherwise opaque. @@ -275,19 +293,21 @@ Failure | Code | Name | Returned when | | --- | --- | --- | -| 1 | `InvalidArgument` | A payload fails validation, a request exceeds a declared limit, the export is unknown, the attachment is already attached, or `Forget` exceeds the references held | +| 1 | `InvalidArgument` | A payload fails validation, a request exceeds a declared limit, the export is unknown, the attachment is already attached, an acquisition names a reserved or open handle ID, or `Forget` exceeds the references held | | 2 | `Unsupported` | The capabilities do not declare the operation or option, or the request needs a feature phase 1 excludes | | 3 | `Unauthorized` | `Attach` names an export the Link binding does not grant, or asks for write access to a read-only grant | | 4 | `StaleAttachment` | The attachment is not attached, has detached or its lease ended | | 5 | `InstanceChanged` | The stream is bound to another service incarnation | | 6 | `StaleNode` | The attachment holds no such `NodeRef` | | 7 | `StaleHandle` | The attachment has no such open handle | -| 8 | `ResourceExhausted` | The stream holds `MaxInFlight` requests, or the attachment has `MaxOpenHandles` handles | +| 8 | `ResourceExhausted` | The stream holds `MaxInFlight` requests, or the attachment has `MaxOpenHandles` handles, counting reserved IDs | | 9 | `Cancelled` | The request was cancelled | | 10 | `DeadlineExceeded` | The caller's deadline passed | | 11 | `Errno` | A file-system call failed; `Errno` says how | | 12 | `Unknown` | The stream failed, or the service failed without a typed error | +`ResourceExhausted` is transient: the same request may succeed later, and `ErrorCode.Retryable` reports it. A client may resend a request that failed with a retryable code and `EffectNone` unchanged, and never resends one that failed with `EffectPossible`. Every other code is final for the request, or reports the caller's own cancellation. + `Errno` is a semantic enum, numbered from 1 in this order: `PermissionDenied`, `OperationNotPermitted`, `NotFound`, `Exists`, `NotDirectory`, `IsDirectory`, `DirectoryNotEmpty`, `InvalidArgument`, `BadDescriptor`, `TooManyOpenFiles`, `NoSpace`, `QuotaExceeded`, `ReadOnlyFilesystem`, `CrossDevice`, `NameTooLong`, `SymlinkLoop`, `FileTooLarge`, `Overflow`, `Busy`, `Again`, `Interrupted`, `IO`, `NoDevice`, `NoSuchDeviceOrAddress`, `BrokenPipe`, `NotSupported`, `NoLocks`, `Deadlock`. Each service converts its native errors; an unknown native error is `IO`, and success is never fabricated. ### Effects and cancellation @@ -298,7 +318,7 @@ Every failure carries `EffectNone`, when the request certainly changed nothing, An ambiguous mutation is never replayed automatically. This includes `Create`, a truncating `Open`, `SetAttr`, `Write` and a lock acquisition. -`CancelRequest` asks the server to stop an outstanding request. Its acknowledgement is not the target's result and never proves that a mutation did not happen; the target's own response still follows. +`CancelRequest` asks the server to stop an outstanding request. Its acknowledgement is not the target's result and never proves that a mutation did not happen; the target's own response still follows. A cancelled acquisition is settled with `Release` or `ReleaseDir` of its ID, as [Handles](#handles) describes. ### Uncached profile @@ -327,4 +347,4 @@ Capabilities may declare smaller limits. ## Verification -`go test ./internal/sandboxfs` covers the golden frames, a round trip of every message, decode rejection and admission while responses go unread. `go test -run '^$' -fuzz FuzzDecode ./internal/sandboxfs` fuzzes the decoder. `go test ./apps/sandboxio/...` runs the Linux service over an in-memory stream against a temporary export, including exclusive create, atomic append, rename modes, directory paging, the root-escape attempts, opaque symlinks and proc magic links, bind-mount aliases, flock against native `flock`, export grants, an incarnation change and transport loss. The bind-mount test reruns itself under `unshare -Urm` and is skipped where unprivileged user namespaces are unavailable. +`go test ./internal/sandboxfs` covers the golden frames, a round trip of every message, decode rejection, the retryable codes, admission while responses go unread, the succession fence for `Attach` then `Detach` and `Open` then `Release`, a waiting successor that ends at once when it is superseded, and a `Release` that follows a cancelled, still running `Open`. `go test -run '^$' -fuzz FuzzDecode ./internal/sandboxfs` fuzzes the decoder. `go test ./apps/sandboxio/...` runs the Linux service over an in-memory stream against a temporary export, including exclusive create, per-write and atomic append, a lost `Open` reply released on the next stream, duplicate handle IDs, rename modes, directory paging, the root-escape attempts, opaque symlinks and proc magic links, bind-mount aliases, flock against native `flock`, export grants, an incarnation change and transport loss. The bind-mount test reruns itself under `unshare -Urm` and is skipped where unprivileged user namespaces are unavailable. diff --git a/internal/sandboxfs/client.go b/internal/sandboxfs/client.go index 6d29da43..51849065 100644 --- a/internal/sandboxfs/client.go +++ b/internal/sandboxfs/client.go @@ -28,9 +28,18 @@ type Client struct { err error done chan struct{} - afterWrite func() // test seam: runs once a frame is recorded as written + afterWrite func() // test seam: runs once a frame is written } +// HandleIDs allocates the handle IDs a client chooses for Open, Create and +// OpenDir: 1, 2, 3 and so on. An attachment keeps one allocator across all of +// its streams, so it never uses an ID twice. Its methods are safe for +// concurrent use. +type HandleIDs struct{ last atomic.Uint64 } + +// Next returns an ID the attachment has not used. +func (h *HandleIDs) Next() HandleID { return HandleID(h.last.Add(1)) } + type call struct { op Op ch chan outcome @@ -49,7 +58,8 @@ func NewClient(conn io.ReadWriteCloser) *Client { return c } -// Close closes the stream. Requests in flight fail with EffectPossible. +// Close closes the stream. Requests in flight fail with EffectPossible. It +// does not wait for the transport to finish closing. func (c *Client) Close() error { c.shutdown(errClientClosed) return nil @@ -84,17 +94,19 @@ func (c *Client) shutdown(cause error) { c.pending, c.abandoned = nil, nil close(c.done) c.mu.Unlock() - c.conn.Close() for _, cl := range pending { cl.ch <- outcome{fail: transportFailure(sandboxwire.EffectPossible, cause)} } + // A transport may hold its close, as yamux does while it cannot queue the + // stream's FIN; no call waits for that. + go c.conn.Close() } // send registers a request and writes it. A request that ends before its // write starts fails with EffectNone. Cancelling ctx during the write fails // the stream, since a partial frame cannot be taken back, and the request -// then fails with EffectPossible. Once the write is recorded as finished, a -// cancellation is left to roundTrip, which sends CancelRequest. +// then fails with EffectPossible. Once the write has finished, a cancellation +// is left to roundTrip, which sends CancelRequest. func (c *Client) send(ctx context.Context, op Op, payload []byte) (uint64, *call, *Failure) { if fail := c.acquire(ctx); fail != nil { return 0, nil, fail @@ -112,32 +124,37 @@ func (c *Client) send(ctx context.Context, op Op, payload []byte) (uint64, *call id, cl := c.seq.Next(), &call{op: op, ch: make(chan outcome, 1)} c.pending[id] = cl c.mu.Unlock() - // The write and the cancellation callback race for state; whichever - // moves it from writing decides. The callback is settled before the - // write slot is released, so it never acts on a later write. - const writing, written, interrupted = 0, 1, 2 - var state atomic.Int32 - settled := make(chan struct{}) - stop := context.AfterFunc(ctx, func() { - defer close(settled) - if state.CompareAndSwap(writing, interrupted) { - c.shutdown(fmt.Errorf("request cancelled while being written: %w", context.Cause(ctx))) - } - }) - err := sandboxwire.WriteFrame(c.conn, sandboxwire.Frame{Type: uint16(op), RequestID: id, Payload: payload}) - state.CompareAndSwap(writing, written) - if c.afterWrite != nil { - c.afterWrite() - } - if !stop() { - <-settled - } - if err != nil { + if err := c.write(ctx, sandboxwire.Frame{Type: uint16(op), RequestID: id, Payload: payload}); err != nil { c.shutdown(err) } return id, cl, nil } +// write writes f for a caller that holds the write slot. It stops waiting +// when ctx ends or the stream fails first and returns why, and the caller +// then fails the stream: a transport can hold a write, and nothing writes on +// a failed stream after it, so releasing the slot cannot interleave frames. +func (c *Client) write(ctx context.Context, f sandboxwire.Frame) error { + written := make(chan error, 1) + go func() { written <- sandboxwire.WriteFrame(c.conn, f) }() + var err error + select { + case err = <-written: + case <-c.done: + return c.Err() + case <-ctx.Done(): + select { + case err = <-written: + default: + return fmt.Errorf("request cancelled while being written: %w", context.Cause(ctx)) + } + } + if err == nil && c.afterWrite != nil { + c.afterWrite() + } + return err +} + // acquire takes the write turn unless ctx ends or the stream fails first. func (c *Client) acquire(ctx context.Context) *Failure { select { diff --git a/internal/sandboxfs/client_test.go b/internal/sandboxfs/client_test.go index 8453bb02..ee92c08b 100644 --- a/internal/sandboxfs/client_test.go +++ b/internal/sandboxfs/client_test.go @@ -4,9 +4,10 @@ import ( "context" "errors" "net" + "sync" "testing" + "time" - "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" ) @@ -14,8 +15,8 @@ import ( // and leaves the stream open. func TestCancelAfterWriteKeepsStream(t *testing.T) { cc, sc := net.Pipe() - a := Attachment{ID: sandboxwire.NewID(), ServerInstanceID: testInstance, Lease: context.Background(), Exports: []sandboxlink.ExportGrant{{ID: "world"}}} - go Serve(context.Background(), sc, &describer{}, a) + a := testAttachment() + go NewServer(&describer{}).Serve(context.Background(), sc, a, 1) c := NewClient(cc) defer c.Close() @@ -32,6 +33,61 @@ func TestCancelAfterWriteKeepsStream(t *testing.T) { } } +// held holds every Write and Close until free closes, as a transport that +// cannot send does. writing closes when the first write starts. +type held struct { + net.Conn + once sync.Once + writing chan struct{} + free <-chan struct{} +} + +func (h *held) Write([]byte) (int, error) { + h.once.Do(func() { close(h.writing) }) + <-h.free + return 0, net.ErrClosed +} + +func (h *held) Close() error { + <-h.free + return h.Conn.Close() +} + +// A call cancelled while the transport holds its write, and then the +// stream's close, returns at once and fails the stream. +func TestCancelDuringHeldWrite(t *testing.T) { + cc, sc := net.Pipe() + defer sc.Close() + free := make(chan struct{}) + defer close(free) + h := &held{Conn: cc, writing: make(chan struct{}), free: free} + c := NewClient(h) + ctx, cancel := context.WithCancel(context.Background()) + go func() { + <-h.writing + cancel() + }() + done := make(chan error, 1) + go func() { + _, err := c.Describe(ctx, &DescribeRequest{}) + done <- err + }() + select { + case err := <-done: + var f *Failure + if !errors.As(err, &f) || f.Effect != sandboxwire.EffectPossible || !errors.Is(err, ErrTransport) { + t.Fatalf("call cancelled mid-write: %v, want a transport failure with EffectPossible", err) + } + case <-time.After(2 * time.Second): + t.Fatal("the cancelled call waits for the transport") + } + select { + case <-c.Done(): + default: + t.Fatal("the stream survived a cancellation mid-write") + } +} + // finisher completes Describe after its request is cancelled, as a lock // acquired just before CancelRequest arrives does. type finisher struct{ Service } @@ -44,8 +100,8 @@ func (finisher) Describe(ctx context.Context, _ Attachment, _ *DescribeRequest) // An interrupt cancels the request and still returns its own outcome. func TestInterruptReturnsOutcome(t *testing.T) { cc, sc := net.Pipe() - a := Attachment{ID: sandboxwire.NewID(), ServerInstanceID: testInstance, Lease: context.Background(), Exports: []sandboxlink.ExportGrant{{ID: "world"}}} - go Serve(context.Background(), sc, finisher{}, a) + a := testAttachment() + go NewServer(finisher{}).Serve(context.Background(), sc, a, 1) c := NewClient(cc) defer c.Close() diff --git a/internal/sandboxfs/protocol.go b/internal/sandboxfs/protocol.go index 1dae79ab..b0b65fe9 100644 --- a/internal/sandboxfs/protocol.go +++ b/internal/sandboxfs/protocol.go @@ -17,7 +17,7 @@ import ( // Version is the protocol version. Link matches it exactly when it opens a // file stream. -const Version uint16 = 1 +const Version uint16 = 2 // Op is a request tag. Each response carries its request's tag with // sandboxwire.ResponseType. The protocol has no events. @@ -137,7 +137,10 @@ type Lease interface { // Service is implemented by a file service, one method per operation. // CancelRequest is not a method: the server answers it by cancelling the // target's context. A method returns a *Failure for a typed failure; the server -// reports any other error as Unknown with EffectPossible. +// reports any other error as Unknown with EffectPossible. A method publishes +// every state it creates before it returns. The server never runs two +// acquisitions (Open, Create, OpenDir) of one handle ID of an attachment at +// once, nor a Release or ReleaseDir while the acquisition of its ID runs. type Service interface { Describe(context.Context, Attachment, *DescribeRequest) (*DescribeResponse, error) Attach(context.Context, Attachment, *AttachRequest) (*AttachResponse, error) @@ -207,6 +210,13 @@ var codeNames = [...]string{ func (c ErrorCode) Valid() bool { return c >= 1 && int(c) < len(codeNames) } func (c ErrorCode) String() string { return enumName(codeNames[:], uint16(c), "ErrorCode") } +// Retryable reports whether the same request may succeed later. +// ResourceExhausted is transient; a client may resend the request unchanged +// after such a failure with EffectNone, and never after EffectPossible. Every +// other code is final for the request, or reports the caller's own +// cancellation. +func (c ErrorCode) Retryable() bool { return c == CodeResourceExhausted } + // Errno is the semantic error of a failed file-system call. Each adapter // converts its native errors; an unknown native error is ErrnoIO. type Errno uint16 @@ -329,13 +339,12 @@ const DurabilityFsyncRequired Durability = 1 type OpenFlags uint32 const ( - OpenAppend OpenFlags = 1 << iota - OpenTruncate + OpenTruncate OpenFlags = 1 << iota OpenNoFollow OpenSync OpenDataSync - openFlagsAll = OpenAppend | OpenTruncate | OpenNoFollow | OpenSync | OpenDataSync + openFlagsAll = OpenTruncate | OpenNoFollow | OpenSync | OpenDataSync ) // AttrMask selects the attributes SetAttr changes. AttrAtimeNow and @@ -396,7 +405,8 @@ type NodeRef struct { Generation uint64 } -// HandleID names an open file or directory handle. It is opaque and nonzero. +// HandleID names an open file or directory handle. The client chooses it when +// it opens the handle: nonzero, and never used twice within an attachment. type HandleID uint64 // LockOwner identifies a lock owner within an attachment. @@ -642,17 +652,21 @@ type AccessRequest struct { type AccessResponse struct{} +// OpenRequest opens Node as Handle, an ID the client chose and never used +// before in the attachment. type OpenRequest struct { + Handle HandleID Node NodeRef Access AccessMode Flags OpenFlags } -type OpenResponse struct { - Handle HandleID -} +type OpenResponse struct{} +// CreateRequest creates and opens a regular file as Handle, an ID the client +// chose and never used before in the attachment. type CreateRequest struct { + Handle HandleID Parent NodeRef Name []byte Mode uint32 // permission bits, applied as given @@ -662,8 +676,7 @@ type CreateRequest struct { } type CreateResponse struct { - Entry Entry - Handle HandleID + Entry Entry } type ReadRequest struct { @@ -677,11 +690,12 @@ type ReadResponse struct { Data []byte } -// WriteRequest writes Data at Offset. An append-opened handle appends -// atomically and ignores Offset. +// WriteRequest writes Data at Offset. An Append write ignores Offset and +// writes Data atomically at the end of the file. type WriteRequest struct { Handle HandleID Offset uint64 + Append bool Data []byte } @@ -714,14 +728,15 @@ type ReleaseRequest struct { type ReleaseResponse struct{} +// OpenDirRequest opens directory Node as Handle, an ID the client chose and +// never used before in the attachment. type OpenDirRequest struct { - Node NodeRef -} - -type OpenDirResponse struct { Handle HandleID + Node NodeRef } +type OpenDirResponse struct{} + // ReadDirRequest reads entries after Cookie; cookie zero is the start. Limit // bounds the entries' WireSize sum. "." and ".." are never returned. type ReadDirRequest struct { @@ -968,6 +983,30 @@ func sideEffectFree(r Request) bool { return false } +// acquires returns the handle ID that r opens. +func acquires(r Request) (HandleID, bool) { + switch r := r.(type) { + case *OpenRequest: + return r.Handle, true + case *CreateRequest: + return r.Handle, true + case *OpenDirRequest: + return r.Handle, true + } + return 0, false +} + +// releases returns the handle ID that r closes. +func releases(r Request) (HandleID, bool) { + switch r := r.(type) { + case *ReleaseRequest: + return r.Handle, true + case *ReleaseDirRequest: + return r.Handle, true + } + return 0, false +} + // modifiesFiles reports whether r changes the file system, which a read-only // attachment refuses. func modifiesFiles(r Request) bool { @@ -1033,6 +1072,9 @@ func (c *Capabilities) Admit(r Request, readOnly bool) *Failure { if tooLong(r.Data, c.MaxWriteBytes) { return NewFailure(CodeInvalidArgument, sandboxwire.EffectNone, "write exceeds MaxWriteBytes") } + if r.Append && !c.AtomicAppend { + return unsupported("an append write") + } case *ReadDirRequest: if r.Limit > c.MaxReadDirBytes { return NewFailure(CodeInvalidArgument, sandboxwire.EffectNone, "limit exceeds MaxReadDirBytes") @@ -1719,26 +1761,29 @@ func (*AccessResponse) decode(*decoder) {} func (*AccessResponse) validate() error { return nil } func (r *OpenRequest) encode(e *sandboxwire.Encoder) { + e.U64(uint64(r.Handle)) r.Node.encode(e) e.Enum(uint16(r.Access)) e.U32(uint32(r.Flags)) } func (r *OpenRequest) decode(d *decoder) { + r.Handle = HandleID(d.u64()) r.Node.decode(d) r.Access = AccessMode(d.u16()) r.Flags = OpenFlags(d.u32()) } func (r *OpenRequest) validate() error { - return errors.Join(r.Node.validate(), validAccess(r.Access), validOpenFlags(r.Flags)) + return errors.Join(r.Handle.validate(), r.Node.validate(), validAccess(r.Access), validOpenFlags(r.Flags)) } -func (r *OpenResponse) encode(e *sandboxwire.Encoder) { e.U64(uint64(r.Handle)) } -func (r *OpenResponse) decode(d *decoder) { r.Handle = HandleID(d.u64()) } -func (r *OpenResponse) validate() error { return r.Handle.validate() } +func (*OpenResponse) encode(*sandboxwire.Encoder) {} +func (*OpenResponse) decode(*decoder) {} +func (*OpenResponse) validate() error { return nil } func (r *CreateRequest) encode(e *sandboxwire.Encoder) { + e.U64(uint64(r.Handle)) r.Parent.encode(e) e.Bytes(r.Name) e.U32(r.Mode) @@ -1748,6 +1793,7 @@ func (r *CreateRequest) encode(e *sandboxwire.Encoder) { } func (r *CreateRequest) decode(d *decoder) { + r.Handle = HandleID(d.u64()) r.Parent.decode(d) r.Name = d.bytes() r.Mode = d.u32() @@ -1757,14 +1803,12 @@ func (r *CreateRequest) decode(d *decoder) { } func (r *CreateRequest) validate() error { - return errors.Join(r.Parent.validate(), validName(r.Name), validPerm(r.Mode), validAccess(r.Access), validOpenFlags(r.Flags)) + return errors.Join(r.Handle.validate(), r.Parent.validate(), validName(r.Name), validPerm(r.Mode), validAccess(r.Access), validOpenFlags(r.Flags)) } -func (r *CreateResponse) encode(e *sandboxwire.Encoder) { r.Entry.encode(e); e.U64(uint64(r.Handle)) } -func (r *CreateResponse) decode(d *decoder) { r.Entry.decode(d); r.Handle = HandleID(d.u64()) } -func (r *CreateResponse) validate() error { - return errors.Join(r.Entry.validate(), r.Handle.validate()) -} +func (r *CreateResponse) encode(e *sandboxwire.Encoder) { r.Entry.encode(e) } +func (r *CreateResponse) decode(d *decoder) { r.Entry.decode(d) } +func (r *CreateResponse) validate() error { return r.Entry.validate() } func (r *ReadRequest) encode(e *sandboxwire.Encoder) { e.U64(uint64(r.Handle)) @@ -1796,11 +1840,13 @@ func (r *ReadResponse) validate() error { func (r *WriteRequest) encode(e *sandboxwire.Encoder) { e.U64(uint64(r.Handle)) e.U64(r.Offset) + e.Bool(r.Append) e.Bytes(r.Data) } func (r *WriteRequest) decode(d *decoder) { r.Handle = HandleID(d.u64()) r.Offset = d.u64() + r.Append = d.bool() r.Data = d.bytes() } @@ -1856,13 +1902,13 @@ func (*ReleaseResponse) encode(*sandboxwire.Encoder) {} func (*ReleaseResponse) decode(*decoder) {} func (*ReleaseResponse) validate() error { return nil } -func (r *OpenDirRequest) encode(e *sandboxwire.Encoder) { r.Node.encode(e) } -func (r *OpenDirRequest) decode(d *decoder) { r.Node.decode(d) } -func (r *OpenDirRequest) validate() error { return r.Node.validate() } +func (r *OpenDirRequest) encode(e *sandboxwire.Encoder) { e.U64(uint64(r.Handle)); r.Node.encode(e) } +func (r *OpenDirRequest) decode(d *decoder) { r.Handle = HandleID(d.u64()); r.Node.decode(d) } +func (r *OpenDirRequest) validate() error { return errors.Join(r.Handle.validate(), r.Node.validate()) } -func (r *OpenDirResponse) encode(e *sandboxwire.Encoder) { e.U64(uint64(r.Handle)) } -func (r *OpenDirResponse) decode(d *decoder) { r.Handle = HandleID(d.u64()) } -func (r *OpenDirResponse) validate() error { return r.Handle.validate() } +func (*OpenDirResponse) encode(*sandboxwire.Encoder) {} +func (*OpenDirResponse) decode(*decoder) {} +func (*OpenDirResponse) validate() error { return nil } func (r *ReadDirRequest) encode(e *sandboxwire.Encoder) { e.U64(uint64(r.Handle)) diff --git a/internal/sandboxfs/protocol_test.go b/internal/sandboxfs/protocol_test.go index ed8ad24a..d984ca07 100644 --- a/internal/sandboxfs/protocol_test.go +++ b/internal/sandboxfs/protocol_test.go @@ -65,7 +65,7 @@ func TestGoldenFixtures(t *testing.T) { {file: "walk_request.hex", op: OpWalk, id: 2, req: &WalkRequest{Parent: testNode, Names: [][]byte{[]byte("link"), []byte("x")}}}, {file: "walk_response.hex", op: OpWalk, id: 2, resp: &WalkResponse{Entries: []Entry{{Node: NodeRef{ID: 2, Generation: 7}, Attr: symlink}}}}, {file: "create_request.hex", op: OpCreate, id: 3, req: &CreateRequest{ - Parent: testNode, Name: []byte("notes.md"), Mode: 0o644, Access: AccessReadWrite, Flags: OpenAppend, Exclusive: true}}, + Handle: 9, Parent: testNode, Name: []byte("notes.md"), Mode: 0o644, Access: AccessReadWrite, Flags: OpenSync, Exclusive: true}}, {file: "write_response.hex", op: OpWrite, id: 4, resp: &WriteResponse{ Written: 4096, Failure: &Failure{Code: CodeErrno, Errno: ErrnoNoSpace, Effect: sandboxwire.EffectNone, Message: "disk full"}}}, {file: "readdir_request.hex", op: OpReadDir, id: 5, req: &ReadDirRequest{Handle: 0x42, Cookie: 0x1c6a3e5f0b9d2471, Limit: 65536}}, @@ -77,6 +77,7 @@ func TestGoldenFixtures(t *testing.T) { Parent: testNode, Name: []byte("old"), NewParent: testNode2, NewName: []byte("new"), Mode: RenameExchange}}, {file: "failure_response.hex", op: OpLookup, id: 7, fail: &Failure{ Code: CodeErrno, Errno: ErrnoNotFound, Effect: sandboxwire.EffectNone, Message: "missing"}}, + {file: "write_request.hex", op: OpWrite, id: 8, req: &WriteRequest{Handle: 9, Append: true, Data: []byte("log\n")}}, } { t.Run(tc.file, func(t *testing.T) { want := readHexFixture(t, tc.file) @@ -137,14 +138,14 @@ func samples() []struct { {&GetAttrRequest{Target: Target{Kind: TargetHandle, Handle: 3}}, &GetAttrResponse{Attr: testAttr}}, {&SetAttrRequest{Target: Target{Kind: TargetNode, Node: testNode}, Set: AttrSize | AttrMode | AttrUID | AttrGID | AttrAtime | AttrMtimeNow, Size: 9, Mode: 0o4755, UID: 5, GID: 6, Atime: Timestamp{7, 8}}, &SetAttrResponse{Attr: testAttr}}, {&AccessRequest{Node: testNode, Mask: MayRead | MayExecute}, &AccessResponse{}}, - {&OpenRequest{Node: testNode, Access: AccessWrite, Flags: OpenTruncate | OpenSync}, &OpenResponse{Handle: 4}}, - {&CreateRequest{Parent: testNode, Name: []byte("n"), Mode: 0o600, Access: AccessRead, Flags: OpenDataSync}, &CreateResponse{Entry: testEntry, Handle: 5}}, + {&OpenRequest{Handle: 4, Node: testNode, Access: AccessWrite, Flags: OpenTruncate | OpenSync}, &OpenResponse{}}, + {&CreateRequest{Handle: 5, Parent: testNode, Name: []byte("n"), Mode: 0o600, Access: AccessRead, Flags: OpenDataSync}, &CreateResponse{Entry: testEntry}}, {&ReadRequest{Handle: 4, Offset: 10, Size: 65536}, &ReadResponse{Data: []byte("data")}}, - {&WriteRequest{Handle: 4, Offset: 10, Data: []byte("data")}, &WriteResponse{Written: 4}}, + {&WriteRequest{Handle: 4, Offset: 10, Append: true, Data: []byte("data")}, &WriteResponse{Written: 4}}, {&FlushRequest{Handle: 4, Owner: 11}, &FlushResponse{}}, {&FsyncRequest{Handle: 4, DataOnly: true}, &FsyncResponse{}}, {&ReleaseRequest{Handle: 4}, &ReleaseResponse{}}, - {&OpenDirRequest{Node: testNode}, &OpenDirResponse{Handle: 6}}, + {&OpenDirRequest{Handle: 6, Node: testNode}, &OpenDirResponse{}}, {&ReadDirRequest{Handle: 6, Cookie: 1, Limit: 4096, WithAttrs: true}, &ReadDirResponse{Entries: []DirEntry{{Name: []byte("f"), Ino: testAttr.Ino, Type: ModeRegular, Cookie: 2, Entry: &testEntry}}, End: true}}, {&ReleaseDirRequest{Handle: 6}, &ReleaseDirResponse{}}, {&MkdirRequest{Parent: testNode, Name: []byte("d"), Mode: 0o1777}, &MkdirResponse{Entry: testEntry}}, @@ -188,6 +189,14 @@ func TestRoundTripEveryMessage(t *testing.T) { } } +func TestRetryableCodes(t *testing.T) { + for c := CodeInvalidArgument; c.Valid(); c++ { + if c.Retryable() != (c == CodeResourceExhausted) { + t.Errorf("%s.Retryable() = %v", c, c.Retryable()) + } + } +} + func TestDecodeRejects(t *testing.T) { encode := func(m message) []byte { var e sandboxwire.Encoder @@ -200,9 +209,10 @@ func TestDecodeRejects(t *testing.T) { }{ "dot-dot name": {OpLookup, encode(&LookupRequest{Parent: testNode, Name: []byte("..")})}, "slash in name": {OpMkdir, encode(&MkdirRequest{Parent: testNode, Name: []byte("a/b")})}, - "zero generation": {OpOpenDir, encode(&OpenDirRequest{Node: NodeRef{ID: 1}})}, - "unknown open flag": {OpOpen, encode(&OpenRequest{Node: testNode, Access: AccessRead, Flags: 1 << 5})}, - "zero access mode": {OpOpen, encode(&OpenRequest{Node: testNode})}, + "zero generation": {OpOpenDir, encode(&OpenDirRequest{Handle: 1, Node: NodeRef{ID: 1}})}, + "zero handle ID": {OpOpen, encode(&OpenRequest{Node: testNode, Access: AccessRead})}, + "unknown open flag": {OpOpen, encode(&OpenRequest{Handle: 1, Node: testNode, Access: AccessRead, Flags: 1 << 4})}, + "zero access mode": {OpOpen, encode(&OpenRequest{Handle: 1, Node: testNode})}, "duplicate forget": {OpForget, encode(&ForgetRequest{Entries: []ForgetEntry{{testNode, 1}, {testNode, 1}}})}, "partial flock range": {OpSetLock, encode(&SetLockRequest{Handle: 1, Kind: LockFlock, Lock: Lock{Mode: LockRead, End: 10}})}, "trailing byte": {OpReleaseDir, append(encode(&ReleaseDirRequest{Handle: 1}), 0)}, diff --git a/internal/sandboxfs/server.go b/internal/sandboxfs/server.go index 52735d53..e431a0cf 100644 --- a/internal/sandboxfs/server.go +++ b/internal/sandboxfs/server.go @@ -6,6 +6,7 @@ import ( "io" "slices" "sync" + "sync/atomic" "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" ) @@ -16,33 +17,84 @@ import ( // always gets through while the client reads responses. const MaxInFlight = 256 -// Serve answers the requests on conn with svc until the stream ends or ctx is -// done. a is the attachment Link authenticated for the stream; Serve refuses -// one without an ID, server instance, lease or valid export grants. Serve owns -// conn and closes it. On return every request context is cancelled and every -// handler has finished. A stream that ends cleanly returns nil. -func Serve(ctx context.Context, conn io.ReadWriteCloser, svc Service, a Attachment) error { +// ErrSuperseded ends a stream that a successor stream of the same attachment +// replaced. +var ErrSuperseded = errors.New("sandboxfs: stream superseded by a successor") + +// ErrLeaseEnded refuses a stream whose attachment's lease had ended before +// the stream was admitted. +var ErrLeaseEnded = errors.New("sandboxfs: attachment lease ended") + +// Server answers File streams for one Service. It admits one stream at a time +// for each (ServerInstanceID, AttachmentID), in the order Link bound them: a +// stream bound before one already admitted is refused, and a successor runs no +// request until every earlier stream that ran one has stopped admitting +// requests and every request it admitted has finished. A successor that has +// run nothing ends at once when it is superseded or its stream or context +// ends. Its methods are safe for concurrent use. +type Server struct { + svc Service + + mu sync.Mutex + // streams holds the newest stream of each attachment, kept after it ends + // until the attachment's lease has ended and the stream has settled. + streams map[streamKey]*stream + + awaitPredecessor func() // test seam: runs when a request of a successor starts waiting +} + +type streamKey struct{ instance, attachment sandboxwire.ID } + +// NewServer returns a Server for svc. A file service uses one Server for all +// of its streams, because the succession fence spans them. +func NewServer(svc Service) *Server { + return &Server{svc: svc, streams: map[streamKey]*stream{}} +} + +// Serve answers the requests on conn until the stream ends, ctx is done or a +// successor stream of the same attachment supersedes it. a is the attachment +// Link authenticated for the stream; Serve refuses one without an ID, server +// instance, lease or valid export grants. seq is the stream's Link bind +// sequence: Serve refuses a stream bound before one of its attachment it +// already admitted with ErrSuperseded, and a stream whose lease has ended with +// ErrLeaseEnded, without running anything. Serve owns conn and closes it. On +// return every request context is cancelled and every handler has finished. A +// stream that ends cleanly returns nil; a superseded one returns ErrSuperseded. +func (s *Server) Serve(ctx context.Context, conn io.ReadWriteCloser, a Attachment, seq uint64) error { if err := a.validate(); err != nil { conn.Close() return err } a.Exports = slices.Clone(a.Exports) ctx, cancel := context.WithCancel(ctx) - s := &server{conn: conn, svc: svc, a: a, inflight: map[uint64]context.CancelFunc{}} + st := &stream{srv: s, key: streamKey{a.ServerInstanceID, a.ID}, conn: conn, svc: s.svc, a: a, bindSeq: seq, cancel: cancel, + drained: make(chan struct{}), inflight: map[uint64]context.CancelFunc{}, acquiring: map[HandleID]chan struct{}{}} + prev, err := s.admit(st) + if err != nil { + cancel() + conn.Close() + return err + } stop := context.AfterFunc(ctx, func() { conn.Close() }) defer func() { cancel() stop() conn.Close() - s.wg.Wait() + st.wg.Wait() + close(st.drained) }() + if prev != nil { + prev.fence() + } for { f, err := sandboxwire.ReadFrame(conn, sandboxwire.MaxPayload) if err == nil { - err = s.dispatch(ctx, f) + err = st.dispatch(ctx, f) } switch { case err == nil: + case st.superseded(): + return ErrSuperseded case ctx.Err() != nil: return ctx.Err() case errors.Is(err, io.EOF): @@ -53,22 +105,128 @@ func Serve(ctx context.Context, conn io.ReadWriteCloser, svc Service, a Attachme } } -type server struct { - conn io.ReadWriteCloser - svc Service - a Attachment - seq sandboxwire.RequestSequence // read loop only - wmu sync.Mutex - wg sync.WaitGroup - - mu sync.Mutex - inflight map[uint64]context.CancelFunc // running handlers, for CancelRequest - held int // admitted requests whose response is not yet written +// admit makes st the newest stream of its attachment and returns the +// predecessor it must fence. It refuses st when its lease has ended, or when a +// stream bound later, or the same stream, was admitted already. st waits for +// the predecessor to drain when the predecessor has run a request; otherwise +// the predecessor never will, and st waits for what the predecessor was +// waiting for. It forgets each attachment whose lease has ended and whose +// newest stream has settled. Every stream of an attachment shares its lease, +// and the lease check and the forgetting happen under one lock, so a stream +// bound before a forgotten one is always refused. +func (s *Server) admit(st *stream) (prev *stream, err error) { + s.mu.Lock() + defer s.mu.Unlock() + if closed(st.a.Lease.Done()) { + return nil, ErrLeaseEnded + } + for k, e := range s.streams { + if closed(e.a.Lease.Done()) && e.settled() { + delete(s.streams, k) + } + } + prev = s.streams[st.key] + if prev != nil && prev.bindSeq >= st.bindSeq { + return nil, ErrSuperseded + } + if prev != nil { + st.after = prev.after + if prev.started.Load() { + st.after = prev.drained + } + } + s.streams[st.key] = st + return prev, nil +} + +// start reports whether st may run a request: once every earlier stream that +// ran one has drained, and only while no successor has superseded st. It +// reports false when st was superseded or ctx ended first; st has then run +// nothing. The check and the mark are one step under the Server's lock, so a +// successor either waits for st or knows that st will run nothing. +func (st *stream) start(ctx context.Context) bool { + if st.started.Load() { + return true + } + if st.after != nil { + if st.srv.awaitPredecessor != nil { + st.srv.awaitPredecessor() + } + // No deadline: a handler an earlier stream admitted may still change + // state this stream's requests depend on. + select { + case <-st.after: + case <-ctx.Done(): + return false + } + } + st.srv.mu.Lock() + defer st.srv.mu.Unlock() + if st.srv.streams[st.key] != st || ctx.Err() != nil { + return false + } + st.started.Store(true) + return true +} + +// settled reports whether st has drained and so has every stream admitted +// before it that ran a request. Until then the attachment's entry carries the +// fence a later stream must wait behind. +func (st *stream) settled() bool { + return closed(st.drained) && (st.after == nil || closed(st.after)) +} + +func closed(c <-chan struct{}) bool { + select { + case <-c: + return true + default: + return false + } +} + +// stream is one served File stream. +type stream struct { + srv *Server + key streamKey + conn io.ReadWriteCloser + svc Service + a Attachment + bindSeq uint64 + cancel context.CancelFunc + drained chan struct{} // closed once the stream has ended and every handler finished + after <-chan struct{} // closed once every earlier stream that ran a request has drained; nil when none did + started atomic.Bool // a request passed the fence; set under the Server's lock + seq sandboxwire.RequestSequence // read loop only + wmu sync.Mutex + wg sync.WaitGroup + + mu sync.Mutex + fenced bool + inflight map[uint64]context.CancelFunc // running handlers, for CancelRequest + held int // admitted requests whose response is not yet written + acquiring map[HandleID]chan struct{} // handle IDs of running acquisitions, closed when each finishes +} + +// fence stops admission on the stream, cancels its requests and closes it, so +// its later responses are discarded. +func (st *stream) fence() { + st.mu.Lock() + st.fenced = true + st.mu.Unlock() + st.cancel() + st.conn.Close() +} + +func (st *stream) superseded() bool { + st.mu.Lock() + defer st.mu.Unlock() + return st.fenced } // dispatch starts one request. A frame that is not a request, or whose // RequestID does not increase, ends the stream. -func (s *server) dispatch(ctx context.Context, f sandboxwire.Frame) error { +func (st *stream) dispatch(ctx context.Context, f sandboxwire.Frame) error { kind, err := tags.Classify(f.Type) if err != nil { return err @@ -76,37 +234,59 @@ func (s *server) dispatch(ctx context.Context, f sandboxwire.Frame) error { if kind != sandboxwire.KindRequest { return malformed("message type %#04x from the client", f.Type) } - if !s.seq.Admit(f.RequestID) { + if !st.seq.Admit(f.RequestID) { return malformed("request ID %d does not increase", f.RequestID) } op := Op(f.Type) req, err := decodeRequest(op, f.Payload) - s.mu.Lock() + st.mu.Lock() + if st.fenced || ctx.Err() != nil { + st.mu.Unlock() + return context.Canceled + } switch { case err != nil: - s.mu.Unlock() - return s.reply(f.RequestID, op, nil, NewFailure(CodeInvalidArgument, sandboxwire.EffectNone, err.Error())) + st.mu.Unlock() + return st.reply(f.RequestID, op, nil, NewFailure(CodeInvalidArgument, sandboxwire.EffectNone, err.Error())) case op == OpCancelRequest: - if cancel := s.inflight[req.(*CancelRequestRequest).Target]; cancel != nil { + if cancel := st.inflight[req.(*CancelRequestRequest).Target]; cancel != nil { cancel() } - s.mu.Unlock() - return s.reply(f.RequestID, op, &CancelRequestResponse{}, nil) - case s.held >= MaxInFlight: - s.mu.Unlock() - return s.reply(f.RequestID, op, nil, NewFailure(CodeResourceExhausted, sandboxwire.EffectNone, "too many requests in flight")) + st.mu.Unlock() + return st.reply(f.RequestID, op, &CancelRequestResponse{}, nil) + case st.held >= MaxInFlight: + st.mu.Unlock() + return st.reply(f.RequestID, op, nil, NewFailure(CodeResourceExhausted, sandboxwire.EffectNone, "too many requests in flight")) + } + // An acquisition owns its handle ID until it finishes; a release of that + // ID waits for it, so it releases whatever the acquisition produced. + var acquired, wait chan struct{} + if id, ok := acquires(req); ok { + if st.acquiring[id] != nil { + st.mu.Unlock() + return st.reply(f.RequestID, op, nil, NewFailure(CodeInvalidArgument, sandboxwire.EffectNone, "handle ID is already in use")) + } + acquired = make(chan struct{}) + st.acquiring[id] = acquired + } else if id, ok := releases(req); ok { + wait = st.acquiring[id] } rctx, cancel := context.WithCancel(ctx) - s.inflight[f.RequestID] = cancel - s.held++ - s.mu.Unlock() - s.wg.Add(1) + st.inflight[f.RequestID] = cancel + st.held++ + st.mu.Unlock() + st.wg.Add(1) go func() { - defer s.wg.Done() - resp, err := opSpecs[op].serve(rctx, s.svc, s.a, req) - s.mu.Lock() - delete(s.inflight, f.RequestID) - s.mu.Unlock() + defer st.wg.Done() + resp, err := st.serve(rctx, op, req, wait) + st.mu.Lock() + delete(st.inflight, f.RequestID) + if acquired != nil { + id, _ := acquires(req) + delete(st.acquiring, id) + close(acquired) + } + st.mu.Unlock() cancel() var fail *Failure if err != nil && !errors.As(err, &fail) { @@ -115,25 +295,45 @@ func (s *server) dispatch(ctx context.Context, f sandboxwire.Frame) error { // The request keeps its slot until its response is written, so a // client that stops reading stops admission instead of piling up // finished requests. - werr := s.reply(f.RequestID, op, resp, fail) - s.mu.Lock() - s.held-- - s.mu.Unlock() + werr := st.reply(f.RequestID, op, resp, fail) + st.mu.Lock() + st.held-- + st.mu.Unlock() if werr != nil { - s.conn.Close() + st.conn.Close() } }() return nil } +// serve runs one request once the stream may run requests, and after the +// acquisition it must follow, if any. +func (st *stream) serve(ctx context.Context, op Op, req Request, after <-chan struct{}) (message, error) { + if !st.start(ctx) { + return nil, NewFailure(CodeCancelled, sandboxwire.EffectNone, "stream superseded or ended before the request ran") + } + if after != nil { + select { + case <-after: + case <-ctx.Done(): + return nil, contextFailure(ctx.Err(), sandboxwire.EffectNone) + } + } + return opSpecs[op].serve(ctx, st.svc, st.a, req) +} + // reply writes a response. A response the service built wrongly becomes -// Unknown with EffectPossible, since the request may have run. -func (s *server) reply(id uint64, op Op, resp message, fail *Failure) error { +// Unknown with EffectPossible, since the request may have run. A fenced +// stream's responses are discarded. +func (st *stream) reply(id uint64, op Op, resp message, fail *Failure) error { payload, err := encodeResponse(resp, fail) if err != nil { payload, _ = encodeResponse(nil, NewFailure(CodeUnknown, sandboxwire.EffectPossible, "service response: "+err.Error())) } - s.wmu.Lock() - defer s.wmu.Unlock() - return sandboxwire.WriteFrame(s.conn, sandboxwire.Frame{Type: sandboxwire.ResponseType(uint16(op)), RequestID: id, Payload: payload}) + st.wmu.Lock() + defer st.wmu.Unlock() + if st.superseded() { + return nil + } + return sandboxwire.WriteFrame(st.conn, sandboxwire.Frame{Type: sandboxwire.ResponseType(uint16(op)), RequestID: id, Payload: payload}) } diff --git a/internal/sandboxfs/server_test.go b/internal/sandboxfs/server_test.go index f397aad7..4b63bd13 100644 --- a/internal/sandboxfs/server_test.go +++ b/internal/sandboxfs/server_test.go @@ -4,6 +4,7 @@ import ( "context" "errors" "net" + "sync" "sync/atomic" "testing" "time" @@ -12,6 +13,10 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" ) +func testAttachment() Attachment { + return Attachment{ID: sandboxwire.NewID(), ServerInstanceID: testInstance, Lease: context.Background(), Exports: []sandboxlink.ExportGrant{{ID: "world"}}} +} + // describer answers Describe and counts the calls. type describer struct { Service @@ -29,8 +34,8 @@ func TestUnreadResponsesStopAdmission(t *testing.T) { cc, sc := net.Pipe() svc := &describer{} served := make(chan error, 1) - a := Attachment{ID: sandboxwire.NewID(), ServerInstanceID: testInstance, Lease: context.Background(), Exports: []sandboxlink.ExportGrant{{ID: "world"}}} - go func() { served <- Serve(context.Background(), sc, svc, a) }() + a := testAttachment() + go func() { served <- NewServer(svc).Serve(context.Background(), sc, a, 1) }() var read atomic.Int32 // net.Pipe completes a write once the server has read it go func() { for id := uint64(1); id <= 2*MaxInFlight; id++ { @@ -60,8 +65,8 @@ func TestRequestIDsMustIncrease(t *testing.T) { defer cc.Close() svc := &describer{} served := make(chan error, 1) - a := Attachment{ID: sandboxwire.NewID(), ServerInstanceID: testInstance, Lease: context.Background(), Exports: []sandboxlink.ExportGrant{{ID: "world"}}} - go func() { served <- Serve(context.Background(), sc, svc, a) }() + a := testAttachment() + go func() { served <- NewServer(svc).Serve(context.Background(), sc, a, 1) }() describe := sandboxwire.Frame{Type: uint16(OpDescribe), RequestID: 5} if err := sandboxwire.WriteFrame(cc, describe); err != nil { t.Fatal(err) @@ -81,3 +86,344 @@ func TestRequestIDsMustIncrease(t *testing.T) { t.Fatalf("a repeated RequestID kept the stream after %d calls", svc.calls.Load()) } } + +// ordered holds Attach and Open until the test lets them finish, whatever +// their context says, as a handler inside a system call does. It records the +// attachment and the handles they create. +type ordered struct { + Service + entered, proceed chan struct{} + + mu sync.Mutex + attached bool + handles map[HandleID]bool +} + +func newOrdered() *ordered { + return &ordered{entered: make(chan struct{}), proceed: make(chan struct{}), handles: map[HandleID]bool{}} +} + +func (o *ordered) hold() { + close(o.entered) + <-o.proceed +} + +func (o *ordered) Describe(context.Context, Attachment, *DescribeRequest) (*DescribeResponse, error) { + return &DescribeResponse{ServerInstanceID: testInstance, Capabilities: testCaps}, nil +} + +func (o *ordered) Attach(context.Context, Attachment, *AttachRequest) (*AttachResponse, error) { + o.hold() + o.mu.Lock() + defer o.mu.Unlock() + o.attached = true + return &AttachResponse{Root: Entry{Node: testNode, Attr: testDirAttr}}, nil +} + +func (o *ordered) Detach(context.Context, Attachment, *DetachRequest) (*DetachResponse, error) { + o.mu.Lock() + defer o.mu.Unlock() + if !o.attached { + return nil, NewFailure(CodeStaleAttachment, sandboxwire.EffectNone, "not attached") + } + o.attached = false + return &DetachResponse{}, nil +} + +func (o *ordered) Open(_ context.Context, _ Attachment, r *OpenRequest) (*OpenResponse, error) { + o.hold() + o.mu.Lock() + defer o.mu.Unlock() + o.handles[r.Handle] = true + return &OpenResponse{}, nil +} + +func (o *ordered) Release(_ context.Context, _ Attachment, r *ReleaseRequest) (*ReleaseResponse, error) { + o.mu.Lock() + defer o.mu.Unlock() + if !o.handles[r.Handle] { + return nil, NewFailure(CodeStaleHandle, sandboxwire.EffectNone, "no such handle") + } + delete(o.handles, r.Handle) + return &ReleaseResponse{}, nil +} + +// leftover reports whether an attachment or handle remains. +func (o *ordered) leftover() bool { + o.mu.Lock() + defer o.mu.Unlock() + return o.attached || len(o.handles) > 0 +} + +func serveStream(t *testing.T, srv *Server, a Attachment, seq uint64) (*Client, chan error) { + cc, sc := net.Pipe() + served := make(chan error, 1) + go func() { served <- srv.Serve(context.Background(), sc, a, seq) }() + c := NewClient(cc) + t.Cleanup(func() { c.Close() }) + return c, served +} + +// A request still running on a failed stream finishes before the successor +// stream of its attachment dispatches the cleanup for it, so the cleanup +// finds what the request created. +func TestSuccessorWaitsForPredecessor(t *testing.T) { + ctx := context.Background() + for _, tc := range []struct { + name string + acquire, settle func(*Client) error + }{ + {"Attach then Detach", + func(c *Client) error { _, err := c.Attach(ctx, &AttachRequest{Export: "world"}); return err }, + func(c *Client) error { _, err := c.Detach(ctx, &DetachRequest{}); return err }}, + {"Open then Release", + func(c *Client) error { + _, err := c.Open(ctx, &OpenRequest{Handle: 7, Node: testNode, Access: AccessRead}) + return err + }, + func(c *Client) error { _, err := c.Release(ctx, &ReleaseRequest{Handle: 7}); return err }}, + } { + t.Run(tc.name, func(t *testing.T) { + svc := newOrdered() + srv := NewServer(svc) + waiting := make(chan struct{}) + srv.awaitPredecessor = func() { close(waiting) } + a := testAttachment() + + first, served := serveStream(t, srv, a, 1) + acquired := make(chan error, 1) + go func() { acquired <- tc.acquire(first) }() + <-svc.entered + + // The client gives up on the first stream and resumes on a second. + second, _ := serveStream(t, srv, a, 2) + settled := make(chan error, 1) + go func() { settled <- tc.settle(second) }() + select { + case <-waiting: + case err := <-settled: + t.Fatalf("the successor served the cleanup while the predecessor ran: %v", err) + } + close(svc.proceed) + + if err := <-settled; err != nil { + t.Fatalf("cleanup on the successor: %v", err) + } + if err := <-acquired; !errors.Is(err, ErrTransport) { + t.Fatalf("the superseded stream answered: %v", err) + } + if err := <-served; !errors.Is(err, ErrSuperseded) { + t.Fatalf("superseded Serve returned %v", err) + } + if svc.leftover() { + t.Fatal("the cleanup left state behind") + } + }) + } +} + +// A successor waiting behind a predecessor whose handler is stuck ends as soon +// as a newer stream supersedes it. The newest stream keeps waiting for the +// stuck handler: its Detach finds the attachment that handler creates. +func TestSupersededWaiterEnds(t *testing.T) { + ctx := context.Background() + svc := newOrdered() + srv := NewServer(svc) + waiting := make(chan struct{}, 2) + srv.awaitPredecessor = func() { waiting <- struct{}{} } + a := testAttachment() + first, _ := serveStream(t, srv, a, 1) + go first.Attach(ctx, &AttachRequest{Export: "world"}) + <-svc.entered + + second, secondServed := serveStream(t, srv, a, 2) + go second.Detach(ctx, &DetachRequest{}) + <-waiting + third, _ := serveStream(t, srv, a, 3) + select { + case err := <-secondServed: + if !errors.Is(err, ErrSuperseded) { + t.Fatalf("superseded Serve returned %v", err) + } + case <-time.After(5 * time.Second): + t.Fatal("the superseded stream still waits for the stuck handler") + } + + settled := make(chan error, 1) + go func() { + _, err := third.Detach(ctx, &DetachRequest{}) + settled <- err + }() + <-waiting + close(svc.proceed) + if err := <-settled; err != nil { + t.Fatalf("Detach on the newest stream: %v, want it to follow the Attach", err) + } + if svc.leftover() { + t.Fatal("the attachment remains") + } +} + +// The end of a waiting successor's stream, and then of the attachment's lease, +// leave the fence in place: a stream of the same attachment ID under a new +// lease, as Link serves an ID bound again after its attachment closed, still +// waits for the stuck handler of the streams before it. The leases and each +// stream's context end separately. +func TestFenceOutlivesLease(t *testing.T) { + ctx := context.Background() + svc := newOrdered() + srv := NewServer(svc) + waiting := make(chan struct{}, 2) + srv.awaitPredecessor = func() { waiting <- struct{}{} } + lease, endLease := context.WithCancel(ctx) + a := testAttachment() + a.Lease = lease + first, _ := serveStream(t, srv, a, 1) + go first.Attach(ctx, &AttachRequest{Export: "world"}) + <-svc.entered + + second, secondServed := serveStream(t, srv, a, 2) + go second.Detach(ctx, &DetachRequest{}) + <-waiting + second.Close() + <-secondServed + endLease() + + renewed := a + renewed.Lease = ctx + third, _ := serveStream(t, srv, renewed, 3) + settled := make(chan error, 1) + go func() { + _, err := third.Detach(ctx, &DetachRequest{}) + settled <- err + }() + select { + case <-waiting: + case err := <-settled: + t.Fatalf("the later stream ran Detach beside the stuck Attach: %v", err) + } + close(svc.proceed) + if err := <-settled; err != nil { + t.Fatalf("Detach: %v, want it to follow the Attach", err) + } +} + +// A stream Link bound before the attachment's newest stream, but whose handler +// reaches Serve later, is refused without dispatching anything, also after the +// newer stream ended. The newer stream keeps serving. +func TestOlderBindIsRefused(t *testing.T) { + ctx := context.Background() + svc := &describer{} + srv := NewServer(svc) + a := testAttachment() + newer, served := serveStream(t, srv, a, 2) + if _, err := newer.Describe(ctx, &DescribeRequest{}); err != nil { + t.Fatal(err) + } + // The older stream's client is gone, so a Serve that admitted it would + // end at once with EOF instead of ErrSuperseded. + late := func() error { + cc, sc := net.Pipe() + cc.Close() + return srv.Serve(ctx, sc, a, 1) + } + if err := late(); !errors.Is(err, ErrSuperseded) { + t.Fatalf("older stream: %v, want ErrSuperseded", err) + } + if _, err := newer.Describe(ctx, &DescribeRequest{}); err != nil { + t.Fatalf("newer stream after the refusal: %v", err) + } + newer.Close() + <-served + if err := late(); !errors.Is(err, ErrSuperseded) || svc.calls.Load() != 2 { + t.Fatalf("older stream after the newer ended: %v, %d calls served", err, svc.calls.Load()) + } +} + +// A stream that reaches the server after its attachment's lease ended, and +// after the newer stream it was bound before has drained and been forgotten, +// is refused. +func TestEndedLeaseIsRefused(t *testing.T) { + ctx := context.Background() + svc := &describer{} + srv := NewServer(svc) + lease, endLease := context.WithCancel(ctx) + a := testAttachment() + a.Lease = lease + newer, served := serveStream(t, srv, a, 2) + if _, err := newer.Describe(ctx, &DescribeRequest{}); err != nil { + t.Fatal(err) + } + newer.Close() + <-served + endLease() + cc, sc := net.Pipe() + cc.Close() + if err := srv.Serve(ctx, sc, a, 1); !errors.Is(err, ErrLeaseEnded) { + t.Fatalf("late stream: %v, want ErrLeaseEnded", err) + } +} + +// On a healthy stream, a Release of a handle whose Open still runs, as after +// the client cancelled the Open, waits and releases what the Open created. A +// second acquisition of the pending ID is refused without effect. +func TestReleaseFollowsPendingAcquisition(t *testing.T) { + svc := newOrdered() + cc, sc := net.Pipe() + defer cc.Close() + go NewServer(svc).Serve(context.Background(), sc, testAttachment(), 1) + responses := make(chan sandboxwire.Frame, 8) + go func() { + for { + f, err := sandboxwire.ReadFrame(cc, sandboxwire.MaxPayload) + if err != nil { + return + } + responses <- f + } + }() + send := func(id uint64, r Request) { + payload, err := encodeRequest(r) + if err == nil { + err = sandboxwire.WriteFrame(cc, sandboxwire.Frame{Type: uint16(r.Op()), RequestID: id, Payload: payload}) + } + if err != nil { + t.Fatal(err) + } + } + outcomes := map[uint64]*Failure{} + receive := func() uint64 { + f := <-responses + _, fail, err := decodeResponse(Op(f.Type^sandboxwire.ResponseType(0)), f.Payload) + if err != nil { + t.Fatal(err) + } + outcomes[f.RequestID] = fail + return f.RequestID + } + + open := &OpenRequest{Handle: 7, Node: testNode, Access: AccessRead} + send(1, open) + <-svc.entered + send(2, &CancelRequestRequest{Target: 1}) + send(3, &ReleaseRequest{Handle: 7}) + send(4, open) + send(5, &DescribeRequest{}) + // The server dispatches in order, so the Describe answer shows that it + // has dispatched the Release. + for receive() != 5 { + } + if _, ok := outcomes[3]; ok { + t.Fatal("Release answered while its Open ran") + } + if f := outcomes[4]; f == nil || f.Code != CodeInvalidArgument || f.Effect != sandboxwire.EffectNone { + t.Fatalf("second Open of a pending ID: %v", f) + } + close(svc.proceed) + for len(outcomes) < 5 { + receive() + } + if f := outcomes[3]; f != nil || svc.leftover() { + t.Fatalf("Release after the Open: %v; handles left: %v", f, svc.leftover()) + } +} diff --git a/internal/sandboxfs/testdata/create_request.hex b/internal/sandboxfs/testdata/create_request.hex index ff011c93..af838b25 100644 --- a/internal/sandboxfs/testdata/create_request.hex +++ b/internal/sandboxfs/testdata/create_request.hex @@ -1,11 +1,12 @@ -# Exclusive Create of an append-mode file. -00000027 # PayloadLength 39 +# Exclusive Create of a synchronous-write file as the client's handle 9. +0000002f # PayloadLength 47 000a # MessageType: Create (10) 0000 # Flags 0000000000000003 # RequestID 3 +0000000000000009 # Handle 9 0000000000000001 0000000000000007 # Parent: ID 1, Generation 7 00000008 6e6f7465732e6d64 # Name "notes.md" 000001a4 # Mode 0644 0003 # Access: ReadWrite -00000001 # Flags: OpenAppend +00000004 # Flags: OpenSync 01 # Exclusive diff --git a/internal/sandboxfs/testdata/write_request.hex b/internal/sandboxfs/testdata/write_request.hex new file mode 100644 index 00000000..cf9dfadd --- /dev/null +++ b/internal/sandboxfs/testdata/write_request.hex @@ -0,0 +1,9 @@ +# Append write: the service ignores Offset and writes at the end of the file. +00000019 # PayloadLength 25 +000c # MessageType: Write (12) +0000 # Flags +0000000000000008 # RequestID 8 +0000000000000009 # Handle 9 +0000000000000000 # Offset +01 # Append +00000004 6c6f670a # Data "log\n"