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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apps/daemon/internal/gateway/gateway_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
2 changes: 1 addition & 1 deletion apps/sandboxio/internal/netservice/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}

Expand Down
2 changes: 1 addition & 1 deletion apps/sandboxio/internal/netservice/service_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
4 changes: 2 additions & 2 deletions apps/sandboxio/internal/sandboxio/sandboxio.go
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down
2 changes: 1 addition & 1 deletion docs/sandbox-link-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
39 changes: 35 additions & 4 deletions internal/sandboxlink/relay/relay_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
}
Expand All @@ -146,17 +148,18 @@ 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 {
err = s.CloseWrite()
}
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{
Expand Down Expand Up @@ -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)
}
}
26 changes: 18 additions & 8 deletions internal/sandboxlink/serve.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
}
}
Expand All @@ -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 {
Expand All @@ -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
Expand Down
9 changes: 6 additions & 3 deletions internal/sandboxlink/serve_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}
Expand All @@ -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)
Expand Down
Loading