diff --git a/apps/daemon/internal/gateway/gateway_test.go b/apps/daemon/internal/gateway/gateway_test.go index b2f319c9..266b978b 100644 --- a/apps/daemon/internal/gateway/gateway_test.go +++ b/apps/daemon/internal/gateway/gateway_test.go @@ -117,7 +117,7 @@ func startSandbox(t *testing.T) *sandbox { sandboxlink.Serve(ctx, sandboxlink.ServeConfig{URL: srv.URL, TLS: srv.TLS, Credential: []byte("serve credential"), Resource: resource, ServerInstanceID: sandboxwire.NewID(), Services: []sandboxlink.ServiceHandler{{Service: sandboxlink.ServiceNetwork, Version: sandboxnet.Version, - Serve: func(ctx context.Context, b sandboxlink.Bind, s sandboxlink.Stream) { + Serve: func(ctx context.Context, b sandboxlink.Bind, _ uint64, s sandboxlink.Stream) { sandboxnet.Serve(ctx, s, b.Egress, sb) }}}, OnConnected: func(sandboxlink.HelloAccepted) { diff --git a/apps/sandboxio/internal/netservice/service.go b/apps/sandboxio/internal/netservice/service.go index 75a8a433..9c980382 100644 --- a/apps/sandboxio/internal/netservice/service.go +++ b/apps/sandboxio/internal/netservice/service.go @@ -32,7 +32,7 @@ func New() *Service { return &Service{resolver: net.DefaultResolver} } // Handle serves one Network stream under the egress its Bind carries. It is // the Serve function of the sandboxlink.ServiceNetwork handler. -func (s *Service) Handle(ctx context.Context, b sandboxlink.Bind, st sandboxlink.Stream) { +func (s *Service) Handle(ctx context.Context, b sandboxlink.Bind, _ uint64, st sandboxlink.Stream) { sandboxnet.Serve(ctx, st, b.Egress, s) } diff --git a/apps/sandboxio/internal/netservice/service_test.go b/apps/sandboxio/internal/netservice/service_test.go index d7df57c3..1925315c 100644 --- a/apps/sandboxio/internal/netservice/service_test.go +++ b/apps/sandboxio/internal/netservice/service_test.go @@ -78,7 +78,7 @@ func newFixture(t *testing.T) *fixture { sandboxlink.Serve(ctx, sandboxlink.ServeConfig{URL: f.srv.URL, TLS: f.srv.TLS, Credential: []byte("serve credential"), Resource: resource, ServerInstanceID: sandboxwire.NewID(), Services: []sandboxlink.ServiceHandler{{Service: sandboxlink.ServiceNetwork, Version: sandboxnet.Version, - Serve: func(ctx context.Context, b sandboxlink.Bind, s sandboxlink.Stream) { + Serve: func(ctx context.Context, b sandboxlink.Bind, _ uint64, s sandboxlink.Stream) { f.served <- sandboxnet.Serve(ctx, s, b.Egress, svc) }}}, OnConnected: func(sandboxlink.HelloAccepted) { diff --git a/apps/sandboxio/internal/sandboxio/sandboxio.go b/apps/sandboxio/internal/sandboxio/sandboxio.go index 3a7dfe6f..75e8cfe4 100644 --- a/apps/sandboxio/internal/sandboxio/sandboxio.go +++ b/apps/sandboxio/internal/sandboxio/sandboxio.go @@ -98,10 +98,10 @@ 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, s sandboxlink.Stream) { + {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.ServiceProcess, Version: sandboxprocess.Version, Serve: func(ctx context.Context, b sandboxlink.Bind, s sandboxlink.Stream) { + {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) }}, {Service: sandboxlink.ServiceNetwork, Version: sandboxnet.Version, Serve: netservice.New().Handle}, diff --git a/docs/sandbox-link-protocol.md b/docs/sandbox-link-protocol.md index b89a27b7..47177bf5 100644 --- a/docs/sandbox-link-protocol.md +++ b/docs/sandbox-link-protocol.md @@ -19,7 +19,7 @@ A stream ends in one of two ways, and the relay keeps them apart end to end. An The Sandbox I/O service runs `sandboxlink.Serve` with a `ServeConfig`: - `URL`, `Credential` and `Resource` come from the [bootstrap input](sandbox-bootstrap.md). The credential identifies the service, so the Hello carries no peer ID. `ServerInstanceID` is a new ID whenever the service starts without its operation and handle registries. -- `Services` holds one handler per offered service and version. A handler receives the `Bind`, which carries the authorized binding including a File stream's exports, and the stream. It owns the stream and returns when it is done with it. Its context ends when the attachment closes or `Serve` returns. +- `Services` holds one handler per offered service and version. A handler receives the `Bind`, which carries the authorized binding including a File stream's exports, a bind sequence and the stream. It owns the stream and returns when it is done with it. Its context ends when the attachment closes or `Serve` returns. Bind order is the order in which `Serve` assigns bind sequences, under the lock that tracks attachments and before it answers `Bound`; concurrent `Bound` and `Opened` replies can reach the opener in another order, and a handler that runs late keeps its stream's place. The sequence belongs to one `Serve` call and survives reconnects; it increases strictly but has gaps, since all attachments share it, and a service may rely on it to fence succession. - `Serve` reconnects with jittered exponential backoff, sending the same `ServerInstanceID`, whenever the link drops. It returns when its context ends or when the relay refuses the Hello with a failure that is not [retryable](#failures), for example `AuthenticationFailed` after the credential is withdrawn or `StaleGeneration` after a newer sandbox took over the resource. Before returning it cancels every handler's context and waits for the handlers. - `OnAttachmentLost` fires when an attachment's last open stream ends while the attachment is still open, such as when the link drops. `OnAttachmentRestored` fires when a stream binds a lost attachment again. `OnAttachmentClosed` fires with the reason when the relay reports `AttachmentClosed`. Losing a socket is not closing an attachment: the service keeps an attachment's state until it is closed. - Closing is final. A `Bind` the relay sent before a close can arrive after the `AttachmentClosed`, so `Serve` refuses a `Bind` for an attachment closed within the last `sandboxlink.HandshakeTimeout` with `LeaseExpired`. A stream of a closed attachment that ends later never marks another attachment lost. diff --git a/internal/sandboxlink/relay/relay_test.go b/internal/sandboxlink/relay/relay_test.go index 2972697d..58bfa82f 100644 --- a/internal/sandboxlink/relay/relay_test.go +++ b/internal/sandboxlink/relay/relay_test.go @@ -113,7 +113,8 @@ func newFixture(t *testing.T) *fixture { } // servePeer is a serve peer whose File handler echoes until EOF and whose -// Process handler resets its stream after one byte. +// Process handler reports its bind sequence after one byte and then resets its +// stream. type servePeer struct { instance sandboxwire.ID connected chan sandboxlink.HelloAccepted @@ -123,6 +124,7 @@ type servePeer struct { closed chan sandboxlink.CloseReason binds chan sandboxlink.Bind echoed chan error + seqs chan uint64 done chan struct{} err error // Serve's result, once done is closed } @@ -146,8 +148,8 @@ func (f *fixture) startServe(generation uint64) *servePeer { f.auth.AddServe(credential, sandboxlink.ServePeer{PeerID: sandboxwire.NewID(), Resource: resource(generation)}) p := &servePeer{instance: sandboxwire.NewID(), connected: make(chan sandboxlink.HelloAccepted, 16), conns: make(chan net.Conn, 16), lost: make(chan sandboxwire.ID, 16), restored: make(chan sandboxwire.ID, 16), closed: make(chan sandboxlink.CloseReason, 128), - binds: make(chan sandboxlink.Bind, 16), echoed: make(chan error, 16), done: make(chan struct{})} - echo := func(_ context.Context, b sandboxlink.Bind, s sandboxlink.Stream) { + binds: make(chan sandboxlink.Bind, 16), echoed: make(chan error, 16), seqs: make(chan uint64, 16), done: make(chan struct{})} + echo := func(_ context.Context, b sandboxlink.Bind, _ uint64, s sandboxlink.Stream) { put(p.binds, b) _, err := io.Copy(s, s) if err == nil { @@ -155,8 +157,9 @@ func (f *fixture) startServe(generation uint64) *servePeer { } put(p.echoed, err) } - resetAfterOne := func(_ context.Context, _ sandboxlink.Bind, s sandboxlink.Stream) { + resetAfterOne := func(_ context.Context, _ sandboxlink.Bind, seq uint64, s sandboxlink.Stream) { s.Read(make([]byte, 1)) + put(p.seqs, seq) s.Reset() } cfg := sandboxlink.ServeConfig{ @@ -487,3 +490,31 @@ func TestServeReconnect(t *testing.T) { t.Fatalf("restored %s, want %s", id, attachment) } } + +// Handlers receive bind sequences in bind order however late they run, and the +// sequence keeps increasing after the serve peer reconnects. +func TestBindSequence(t *testing.T) { + f := newFixture(t) + p := f.serve(1) + attachment := sandboxwire.NewID() + // A Process handler reports its sequence once its stream's first byte + // arrives, so the test decides which handler runs first. + seq := func(s sandboxlink.Stream) uint64 { + t.Helper() + if _, err := s.Write([]byte{0}); err != nil { + t.Fatal(err) + } + return recv(t, p.seqs) + } + first, _ := f.mustOpen(sandboxlink.ServiceProcess, attachment, 1) + second, _ := f.mustOpen(sandboxlink.ServiceProcess, attachment, 1) + later := seq(second) + earlier := seq(first) + recv(t, p.conns).Close() + recv(t, p.connected) + third, _ := f.mustOpen(sandboxlink.ServiceProcess, attachment, 1) + reconnected := seq(third) + if !(0 < earlier && earlier < later && later < reconnected) { + t.Fatalf("bind sequences %d, %d, then %d after a reconnect; want them to increase from above zero", earlier, later, reconnected) + } +} diff --git a/internal/sandboxlink/serve.go b/internal/sandboxlink/serve.go index 98964e07..726a008c 100644 --- a/internal/sandboxlink/serve.go +++ b/internal/sandboxlink/serve.go @@ -15,10 +15,17 @@ import ( // ServiceHandler serves one service. Serve owns the stream and returns when // it is done with it. ctx ends when the attachment closes or Serve returns. // The Bind carries the authorized binding, including a File stream's exports. +// seq is the stream's bind sequence. Bind order is the order in which the +// serve peer assigns it, under the lock that tracks attachments and before it +// answers Bound; concurrent Bound and Opened replies may reach the opener in +// another order, and a handler that runs late keeps its stream's place. The +// sequence belongs to one Serve call and survives reconnects; it increases +// strictly but has gaps, since all attachments share it. A service may rely +// on it to fence a stream's successor. type ServiceHandler struct { Service Service Version uint16 - Serve func(ctx context.Context, b Bind, s Stream) + Serve func(ctx context.Context, b Bind, seq uint64, s Stream) } // ServeConfig configures a serve peer. Dial replaces the WebSocket dial to URL @@ -119,6 +126,7 @@ type server struct { mu sync.Mutex attachments map[sandboxwire.ID]*served + binds uint64 // bind sequence of the last stream bound // closedIDs holds recently closed attachment IDs until the time given. A // Bind the relay sent before a close can reach bind after the close. closedIDs map[sandboxwire.ID]time.Time @@ -199,6 +207,7 @@ func (s *server) bind(ctx context.Context, st *yamux.Stream) { } var refuse Code var a *served + var seq uint64 switch { case err != nil || !ok: refuse = ProtocolViolation @@ -209,7 +218,7 @@ func (s *server) bind(ctx context.Context, st *yamux.Stream) { case b.ExpectedServerInstanceID != s.cfg.ServerInstanceID: refuse = InstanceChanged default: - if a = s.track(ctx, b.AttachmentID); a == nil { + if a, seq = s.track(ctx, b.AttachmentID); a == nil { refuse = LeaseExpired } } @@ -228,16 +237,16 @@ func (s *server) bind(ctx context.Context, st *yamux.Stream) { return } st.SetDeadline(time.Time{}) - h.Serve(a.ctx, b, st) + h.Serve(a.ctx, b, seq, st) } -// track counts a bound stream of attachment id. It returns nil for a recently -// closed attachment. -func (s *server) track(ctx context.Context, id sandboxwire.ID) *served { +// track counts a bound stream of attachment id and returns its bind sequence. +// It returns nil for a recently closed attachment. +func (s *server) track(ctx context.Context, id sandboxwire.ID) (*served, uint64) { s.mu.Lock() defer s.mu.Unlock() if until, closed := s.closedIDs[id]; closed && time.Now().Before(until) { - return nil + return nil, 0 } a := s.attachments[id] if a == nil { @@ -252,7 +261,8 @@ func (s *server) track(ctx context.Context, id sandboxwire.ID) *served { s.cfg.OnAttachmentRestored(id) } } - return a + s.binds++ + return a, s.binds } // release ends one bound stream of a. Only the current attachment of its ID diff --git a/internal/sandboxlink/serve_test.go b/internal/sandboxlink/serve_test.go index 70c23cdf..5c0d40bd 100644 --- a/internal/sandboxlink/serve_test.go +++ b/internal/sandboxlink/serve_test.go @@ -19,7 +19,7 @@ func TestCloseBeforeBind(t *testing.T) { s := testServer(&lost) id := sandboxwire.NewID() s.closed(AttachmentClosed{AttachmentID: id, Reason: CloseRevoked}) - if a := s.track(context.Background(), id); a != nil { + if a, _ := s.track(context.Background(), id); a != nil { t.Fatal("a Bind after AttachmentClosed was tracked") } } @@ -30,10 +30,13 @@ func TestReleaseOfClosedAttachment(t *testing.T) { var lost []sandboxwire.ID s := testServer(&lost) id := sandboxwire.NewID() - old := s.track(context.Background(), id) + old, oldSeq := s.track(context.Background(), id) s.closed(AttachmentClosed{AttachmentID: id, Reason: CloseRequested}) s.closedIDs[id] = time.Now() // let the tombstone lapse - current := s.track(context.Background(), id) + current, seq := s.track(context.Background(), id) + if seq <= oldSeq { + t.Fatalf("bind sequence %d after %d; want it to increase", seq, oldSeq) + } s.release(old) if current.streams != 1 || len(lost) != 0 { t.Fatalf("after the old stream ended: %d streams, lost %v; want 1 stream and none lost", current.streams, lost)