From 7fa5fb4de3940faf9833d14e36f02c27df401324 Mon Sep 17 00:00:00 2001 From: jiashuoz Date: Thu, 1 Oct 2026 12:25:40 +0800 Subject: [PATCH] feat(driver): bound guest admission and bridge reconnect proofs --- .../design/2026-09-30-guest-reconnect-host.md | 4 + .../2026-10-01-guest-reconnect-transport.md | 79 ++++ internal/driver/microvm_vsock.go | 29 +- internal/driver/reconnect_transport.go | 112 +++++ internal/driver/reconnect_transport_test.go | 407 ++++++++++++++++++ internal/runnerd/reconnect_process_test.go | 16 +- .../runnerd/testdata/reconnect-host/main.go | 86 +++- 7 files changed, 698 insertions(+), 35 deletions(-) create mode 100644 docs/design/2026-10-01-guest-reconnect-transport.md create mode 100644 internal/driver/reconnect_transport.go create mode 100644 internal/driver/reconnect_transport_test.go diff --git a/docs/design/2026-09-30-guest-reconnect-host.md b/docs/design/2026-09-30-guest-reconnect-host.md index 38d510e..abb79b5 100644 --- a/docs/design/2026-09-30-guest-reconnect-host.md +++ b/docs/design/2026-09-30-guest-reconnect-host.md @@ -58,3 +58,7 @@ PR146/147/148. No RAM persistence is introduced. Guest identity and preamble implementation are described in [the guest-side slice](2026-09-30-guest-reconnect-session.md). Host listener integration and capability enablement remain separate gates. + +The [host stream transport](2026-10-01-guest-reconnect-transport.md) supplies +the bounded proof callback and moves legacy listener admission before worker +creation. Live listener integration remains gated on the complete handoff. diff --git a/docs/design/2026-10-01-guest-reconnect-transport.md b/docs/design/2026-10-01-guest-reconnect-transport.md new file mode 100644 index 0000000..db76035 --- /dev/null +++ b/docs/design/2026-10-01-guest-reconnect-transport.md @@ -0,0 +1,79 @@ +# Host guest-proof transport and listener admission + +This dependent slice connects the [runner authorization port](2026-09-30-guest-reconnect-host.md) +to the [guest proof frames](2026-09-30-guest-reconnect-session.md). It is a transport +prerequisite, not enablement of live reconnection. + +## Admission before work + +The existing microVM listener now claims its one boot connection inline before +starting a serving goroutine. Later peers are closed inline, without a worker or +per-peer log entry. A guest that holds the first configuration write open cannot +cause an unbounded number of rejected-peer goroutines or log entries. Teardown +still owns and closes the claimed connection; a failed first write does not +release the single-use claim. Boot configuration and legacy behavior are otherwise +unchanged. The listener still refuses every later connection without sending data. + +## One admitted stream, one proof + +`driver.AuthorizeGuestConnection` supplies the real bounded stream exchange to +`GuestReconnectHost.AuthorizeGuestReconnect`. The driver supplies the session +assignment; guest bytes cannot select it. The helper validates the challenge and +expected session before writing, sends only that challenge, reads one strict +4096-byte outer frame, requires the proof kind and matching attempt ID, and uses +the shared strict proof decoder before returning its signature to the host. +The control plane, not this transport, verifies the signature and authority. + +The entire call has a five-second ceiling. The caller deadline and the host's +proof context can shorten it. Cancellation closes blocked peer I/O; failures +close the connection and return zero acceptance with only a fixed refusal code. +Unknown host/provider text becomes `unavailable`. Duplicate proof callbacks, +acceptance without a completed proof, malformed acceptance and late success are +refused. The call remains synchronous; a host violating its cancellation contract +retains its admitted slot rather than creating detached work. + +Success returns committed acceptance to the caller and leaves the same connection +open, including any buffered bytes. It does not send the accepted token, a boot +configuration or a refusal payload, and does not attach or replace a relay. This +separates the proof transport from the guarded installation that still must be +implemented. Rebuilding a reader around the raw socket would lose buffered bytes; +the caller must retain the same `relay.Conn`. + +## Integration limits + +No shipping listener invokes this helper yet and no boot opt-in/capability is +advertised. Relaxing `guestChannel.served` without the complete handoff would +replay configuration or let an unproven peer evict a healthy guest, so that guard +remains intact. + +The next integration needs an instance-bound handoff that retains original runner +connection, local handle/boot, placement and accepted epoch through current +configuration delivery and relay installation. It must fence the old relay before +new input/RPC, bound peer admission before launching work, validate recovered +VMM/jail/socket ownership, and resolve cold enrollment/generation ordering. +Guest configuration redemption/readiness must complete before attachments can +interleave with the guest preamble; authorization success alone is insufficient. +Current launch configuration must come through the authenticated control-plane +path, not the cached initial boot token/configuration. Boot-time artifact invariants +and existing-child environment limits from the guest design still apply. + +An all-in-one listener change using cached `g.boot` was rejected because it would +violate fresh configuration and single-use token rules. Spawning a goroutine for +all peers and then checking the claim was rejected because admission no longer +bounds allocated work. Keeping stream encoding in callers was rejected because +it would duplicate strict framing, cancellation and attempt correlation rules. + +## Verification + +Tests pin admission before the next accept, unchanged first-boot behavior, +challenge/proof round trips, invalid scope/attempt/kind/schema, oversized data +without newline, duplicate callbacks, zero authority on refusal, fixed errors, +caller and host cancellation, deadline, and buffered-reader continuity. + +The built runner probe now runs the production host port and this transport over +real agent WebSocket and guest TCP streams. Its synthetic guest signs the +challenge; the control-plane fixture verifies that exact signature. Wrong scope, +revoked authority, wrong attempt and oversized proof are refused. It asserts that +no acceptance/configuration bytes reach the guest from this helper. This is +executable transport coverage, not a shipping listener, AF_VSOCK/KVM, enrolled +guest recovery, relay takeover or real coding-agent qualification result. diff --git a/internal/driver/microvm_vsock.go b/internal/driver/microvm_vsock.go index 3017c61..143cdfa 100644 --- a/internal/driver/microvm_vsock.go +++ b/internal/driver/microvm_vsock.go @@ -368,15 +368,10 @@ func (m *Microvm) openGuestChannel(sessionID, udsPath, listenPath string, cfg ru return g, nil } -// acceptGuests keeps accepting for the life of the instance. -// -// It keeps accepting even though exactly one connection is ever SERVED, -// because the alternative is worse in both directions: an accept loop that -// stopped would leave later dials queued in the kernel with nobody to refuse -// them, and a loop that served inline would be wedged for good by the first -// guest that connected and did not read (the boot config can exceed a -// megabyte — see bootConfigWriteTimeout). So each connection gets a goroutine -// of its own, and all but the first are refused in it. +// acceptGuests keeps accepting for the life of the instance. Admission happens +// inline before starting a worker: only the one claimed boot connection may +// allocate a serving goroutine. Refused peers receive no bytes and no per-peer +// log line. A blocked boot write cannot prevent rejection of later peers. func (m *Microvm) acceptGuests(sessionID string, g *guestChannel) { for { c, err := g.listener.Accept() @@ -386,11 +381,16 @@ func (m *Microvm) acceptGuests(sessionID string, g *guestChannel) { } return } + if !g.claim(c) { + _ = c.Close() + continue + } go m.serveGuest(sessionID, g, c) } } -// serveGuest hands ONE guest connection its configuration and then hands the +// serveGuest handles the ONE connection already claimed by acceptGuests. +// It sends that guest its configuration and then hands the // connection to the runner. // // The boot configuration is the FIRST frame on the conn, written here before @@ -407,15 +407,6 @@ func (m *Microvm) acceptGuests(sessionID string, g *guestChannel) { // holding it. There is nothing to attach it to, and the claim stays spent: // this boot has had its one connection. func (m *Microvm) serveGuest(sessionID string, g *guestChannel, c net.Conn) { - if !g.claim(c) { - // Not an error the operator can act on and not a rarity worth a - // line per occurrence — but it IS somebody in the guest dialling a - // socket that is not theirs, so it is said once per attempt and - // names nothing about the session but its id. - log.Printf("microvm: session %s: refusing a second connection on a control channel that serves one guest per boot", sessionID) - _ = c.Close() - return - } conn := relay.NetConn(c) ctx, cancel := context.WithTimeout(context.Background(), bootConfigWriteTimeout) defer cancel() diff --git a/internal/driver/reconnect_transport.go b/internal/driver/reconnect_transport.go new file mode 100644 index 0000000..c75a883 --- /dev/null +++ b/internal/driver/reconnect_transport.go @@ -0,0 +1,112 @@ +package driver + +import ( + "context" + "encoding/json" + "errors" + "sync/atomic" + "time" + + "github.com/tokencanopy/rainier/internal/relay" + "github.com/tokencanopy/rainier/protocol/runner" +) + +var errGuestStream = errors.New("unavailable") + +// AuthorizeGuestConnection exchanges one bounded challenge/proof on an already +// admitted guest stream using the runner's original-connection authorization +// port. session is the driver's assignment, never data read from the guest. +// Failure closes the stream and returns zero authority with a fixed code. +// +// Success does NOT write accepted/configuration, replace a relay, or transfer +// ownership of c. The caller must fence the original instance/placement/epoch, +// acquire fresh configuration and install the relay before publishing readiness. +// Keep the same Conn for that handoff: it may buffer bytes after the proof. +// There is deliberately no shipping listener caller until those gates exist. +// +// The caller must bound peer admission before invoking this synchronously. Like +// GuestReconnectHost, this helper does not detach a host that ignores context. +func AuthorizeGuestConnection(ctx context.Context, host GuestReconnectHost, session string, c relay.Conn) (runner.GuestReconnectAcceptResponse, error) { + var zero runner.GuestReconnectAcceptResponse + if c == nil { + return zero, errGuestStream + } + success := false + defer func() { + if !success { + _ = c.Close() + } + }() + if host == nil || !guestStreamSession(session) || ctx.Err() != nil { + return zero, errGuestStream + } + ctx, cancel := context.WithTimeout(ctx, 5*time.Second) + defer cancel() + stop := context.AfterFunc(ctx, func() { _ = c.Close() }) + defer stop() + var invoked, proved, invalid atomic.Bool + accepted, err := host.AuthorizeGuestReconnect(ctx, session, func(proofCtx context.Context, challenge runner.GuestReconnectChallenge) (string, error) { + if !invoked.CompareAndSwap(false, true) { + invalid.Store(true) + return "", errGuestStream + } + body, _ := json.Marshal(challenge) + if _, err := runner.DecodeGuestReconnectChallenge(body); err != nil || challenge.SessionID != session { + invalid.Store(true) + return "", errGuestStream + } + // Both the outer transport budget and the host's connection-bound proof + // context must remain live. Neither may extend the other. + pctx, pcancel := context.WithCancel(ctx) + defer pcancel() + pstop := context.AfterFunc(proofCtx, pcancel) + defer pstop() + if proofCtx.Err() != nil { + return "", errGuestStream + } + if relay.WriteGuestReconnectFrame(pctx, c, relay.KindGuestReconnectChallenge, body) != nil { + return "", errGuestStream + } + event, err := relay.ReadGuestReconnectFrame(pctx, c) + if err != nil || event.Kind != relay.KindGuestReconnectProof { + return "", errGuestStream + } + proof, err := runner.DecodeGuestReconnectAcceptRequest(event.Payload) + if err != nil || proof.AttemptID != challenge.AttemptID || pctx.Err() != nil || proofCtx.Err() != nil { + return "", errGuestStream + } + proved.Store(true) + return proof.Signature, nil + }) + if err != nil || !proved.Load() || invalid.Load() || ctx.Err() != nil { + return zero, guestStreamError(err) + } + body, _ := json.Marshal(accepted) + accepted, err = runner.DecodeGuestReconnectAcceptResponse(body) + if err != nil || ctx.Err() != nil { + return zero, errGuestStream + } + success = true + return accepted, nil +} + +func guestStreamError(err error) error { + if err != nil { + switch err.Error() { + case "invalid", "expired", "fenced": + return errors.New(err.Error()) + } + } + return errGuestStream +} +func guestStreamSession(id string) bool { + if len(id) == 0 || len(id) > 256 { + return false + } + for i := range id { + if id[i] < 0x21 || id[i] > 0x7e { + return false + } + } + return true +} diff --git a/internal/driver/reconnect_transport_test.go b/internal/driver/reconnect_transport_test.go new file mode 100644 index 0000000..80555fe --- /dev/null +++ b/internal/driver/reconnect_transport_test.go @@ -0,0 +1,407 @@ +package driver + +import ( + "context" + "crypto/ed25519" + "crypto/rand" + "encoding/base64" + "encoding/json" + "errors" + "github.com/tokencanopy/rainier/internal/relay" + "net" + "runtime" + "strings" + "testing" + "time" + + "github.com/tokencanopy/rainier/protocol/runner" +) + +type admissionListener struct { + g *guestChannel + conn net.Conn + called bool + claimed bool +} + +func (l *admissionListener) Accept() (net.Conn, error) { + if !l.called { + l.called = true + return l.conn, nil + } + l.g.mu.Lock() + l.claimed = l.g.served + l.g.mu.Unlock() + return nil, net.ErrClosed +} +func (*admissionListener) Close() error { return nil } +func (*admissionListener) Addr() net.Addr { return &net.UnixAddr{Name: "admission-test", Net: "unix"} } + +func TestGuestAdmissionPrecedesWorker(t *testing.T) { + // The next Accept observes admission synchronously, before the worker can + // run. A busy listener must not spawn one goroutine per refused peer. + previous := runtime.GOMAXPROCS(1) + defer runtime.GOMAXPROCS(previous) + host, peer := net.Pipe() + defer peer.Close() + g := &guestChannel{boot: runner.BootConfig{Protocol: 1, SessionID: "admission-test"}} + listener := &admissionListener{g: g, conn: host} + g.listener = listener + defer g.close() + m := &Microvm{} + m.acceptGuests("admission-test", g) + if !listener.claimed { + t.Fatal("accept loop admitted another peer before claiming the worker's connection") + } +} + +type streamAuthorizer func(context.Context, string, GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) + +func (f streamAuthorizer) AuthorizeGuestReconnect(ctx context.Context, id string, prove GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + return f(ctx, id, prove) +} +func streamChallenge() runner.GuestReconnectChallenge { + return runner.GuestReconnectChallenge{Protocol: 1, SessionID: "stream-test", BootEpoch: "boot-test", HostIncarnation: "host-test", AttemptID: "attempt-test", PlacementGeneration: 1, Challenge: base64.RawURLEncoding.EncodeToString(make([]byte, 32))} +} +func streamAcceptance() runner.GuestReconnectAcceptResponse { + return runner.GuestReconnectAcceptResponse{Epoch: 1, Token: base64.RawURLEncoding.EncodeToString(make([]byte, 32)), ExpiresInSec: 120} +} +func TestGuestStreamProofReachesHostAuthority(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + pub, key, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + challenge := streamChallenge() + accepted := streamAcceptance() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + authority := streamAuthorizer(func(ctx context.Context, id string, prove GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + if id != "stream-test" { + return runner.GuestReconnectAcceptResponse{}, net.ErrClosed + } + signature, err := prove(ctx, challenge) + if err != nil { + return runner.GuestReconnectAcceptResponse{}, err + } + if err = challenge.VerifyProof(base64.RawURLEncoding.EncodeToString(pub), signature); err != nil { + return runner.GuestReconnectAcceptResponse{}, err + } + return accepted, nil + }) + done := make(chan error, 1) + go func() { + got, err := AuthorizeGuestConnection(ctx, authority, "stream-test", relay.NetConn(a)) + if err == nil && got != accepted { + err = net.ErrClosed + } + done <- err + }() + guest := relay.NetConn(z) + event, err := relay.ReadGuestReconnectFrame(ctx, guest) + if err != nil { + t.Fatalf("host did not deliver a bounded challenge: %v", err) + } + if event.Kind != relay.KindGuestReconnectChallenge { + t.Fatal("host sent authority or configuration before proof") + } + decoded, err := runner.DecodeGuestReconnectChallenge(event.Payload) + if err != nil || decoded != challenge { + t.Fatal("wrong challenge") + } + message, _ := decoded.SigningMessage() + body, _ := json.Marshal(runner.GuestReconnectAcceptRequest{Protocol: 1, AttemptID: decoded.AttemptID, Signature: base64.RawURLEncoding.EncodeToString(ed25519.Sign(key, message))}) + if err = relay.WriteGuestReconnectFrame(ctx, guest, relay.KindGuestReconnectProof, body); err != nil { + t.Fatal(err) + } + if err = <-done; err != nil { + t.Fatal(err) + } + // Authorization alone must not publish accepted/config or transfer a relay. + // The caller still owns that guarded handoff. The stream remains usable. + written := make(chan error, 1) + go func() { written <- relay.NetConn(a).Write(ctx, []byte("caller-owned")) }() + raw, err := guest.Read(ctx) + if err != nil || string(raw) != "caller-owned" { + t.Fatal("stream was closed or helper published authority") + } + if err = <-written; err != nil { + t.Fatal(err) + } +} + +func TestGuestStreamRejectsInvalidProof(t *testing.T) { + for _, which := range []string{"attempt", "kind", "unknown", "duplicate", "oversize", "null", "bad-signature", "malformed-acceptance"} { + t.Run(which, func(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + challenge := streamChallenge() + authority := streamAuthorizer(func(ctx context.Context, _ string, prove GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + _, err := prove(ctx, challenge) + if err != nil { + return runner.GuestReconnectAcceptResponse{}, err + } + accepted := streamAcceptance() + if which == "malformed-acceptance" { + accepted.Token = "invalid" + } + return accepted, nil + }) + done := make(chan error, 1) + go func() { + got, err := AuthorizeGuestConnection(ctx, authority, "stream-test", relay.NetConn(a)) + if got != (runner.GuestReconnectAcceptResponse{}) { + done <- errors.New("nonzero refused authority") + return + } + done <- err + }() + guest := relay.NetConn(z) + if _, err := relay.ReadGuestReconnectFrame(ctx, guest); err != nil { + t.Fatal(err) + } + if which == "oversize" { + _, _ = z.Write([]byte(strings.Repeat("x", 4097))) + } else { + proof := runner.GuestReconnectAcceptRequest{Protocol: 1, AttemptID: challenge.AttemptID, Signature: base64.RawURLEncoding.EncodeToString(make([]byte, 64))} + if which == "attempt" { + proof.AttemptID = "other-attempt" + } + if which == "bad-signature" { + proof.Signature = "not-canonical" + } + body, _ := json.Marshal(proof) + kind := relay.KindGuestReconnectProof + switch which { + case "kind": + kind = relay.KindGuestReconnectAccepted + case "unknown": + body = append(body[:len(body)-1], []byte(`,"extra":1}`)...) + case "duplicate": + body = append(body[:len(body)-1], []byte(`,"protocol":1}`)...) + case "null": + body = []byte(`null`) + } + _ = relay.WriteGuestReconnectFrame(ctx, guest, kind, body) + } + if err := <-done; err == nil || err.Error() != "unavailable" { + t.Fatalf("refusal=%v", err) + } + if raw, err := guest.Read(ctx); err == nil || len(raw) > 0 { + t.Fatal("failed peer received authority or remained connected") + } + }) + } +} + +func TestGuestStreamHostRefusalSendsNothing(t *testing.T) { + for _, reason := range []string{"invalid", "expired", "fenced", "provider-secret-test"} { + t.Run(reason, func(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + authority := streamAuthorizer(func(context.Context, string, GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + return streamAcceptance(), errors.New(reason) + }) + got, err := AuthorizeGuestConnection(ctx, authority, "stream-test", relay.NetConn(a)) + want := reason + if reason == "provider-secret-test" { + want = "unavailable" + } + if got != (runner.GuestReconnectAcceptResponse{}) || err == nil || err.Error() != want { + t.Fatalf("bad refusal %v", err) + } + if raw, err := relay.NetConn(z).Read(ctx); err == nil || len(raw) > 0 { + t.Fatal("refused host sent guest data") + } + }) + } +} + +func TestGuestStreamRequiresProofAndCorrectScope(t *testing.T) { + for _, which := range []string{"no-proof", "wrong-session", "invalid-challenge", "second-proof"} { + t.Run(which, func(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + authority := streamAuthorizer(func(ctx context.Context, _ string, prove GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + if which != "no-proof" { + ch := streamChallenge() + if which == "wrong-session" { + ch.SessionID = "other-test" + } + if which == "invalid-challenge" { + ch.Protocol = 2 + } + _, _ = prove(ctx, ch) + if which == "second-proof" { + _, _ = prove(ctx, ch) + } + } + // Even a misbehaving host cannot report acceptance without a valid exchange. + return streamAcceptance(), nil + }) + done := make(chan error, 1) + go func() { + _, err := AuthorizeGuestConnection(ctx, authority, "stream-test", relay.NetConn(a)) + done <- err + }() + if which == "second-proof" { + guest := relay.NetConn(z) + ev, err := relay.ReadGuestReconnectFrame(ctx, guest) + if err != nil { + t.Fatal(err) + } + ch, _ := runner.DecodeGuestReconnectChallenge(ev.Payload) + body, _ := json.Marshal(runner.GuestReconnectAcceptRequest{Protocol: 1, AttemptID: ch.AttemptID, Signature: base64.RawURLEncoding.EncodeToString(make([]byte, 64))}) + _ = relay.WriteGuestReconnectFrame(ctx, guest, relay.KindGuestReconnectProof, body) + } + if err := <-done; err == nil { + t.Fatal("accepted invalid authority choreography") + } + if raw, err := relay.NetConn(z).Read(ctx); err == nil || len(raw) > 0 { + t.Fatal("invalid challenge or acceptance escaped to guest") + } + }) + } +} + +func TestGuestStreamCancellationClosesIO(t *testing.T) { + for _, which := range []string{"outer", "host", "deadline"} { + t.Run(which, func(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + if which == "deadline" { + var timed context.CancelFunc + ctx, timed = context.WithTimeout(ctx, 150*time.Millisecond) + defer timed() + } + cancelProof := make(chan context.CancelFunc, 1) + authority := streamAuthorizer(func(ctx context.Context, _ string, prove GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + pctx, pcancel := context.WithCancel(ctx) + defer pcancel() + cancelProof <- pcancel + _, err := prove(pctx, streamChallenge()) + return runner.GuestReconnectAcceptResponse{}, err + }) + done := make(chan error, 1) + go func() { + _, err := AuthorizeGuestConnection(ctx, authority, "stream-test", relay.NetConn(a)) + done <- err + }() + guest := relay.NetConn(z) + readCtx, readCancel := context.WithTimeout(context.Background(), time.Second) + defer readCancel() + if _, err := relay.ReadGuestReconnectFrame(readCtx, guest); err != nil { + t.Fatal(err) + } + hostCancel := <-cancelProof + if which == "outer" { + cancel() + } + if which == "host" { + hostCancel() + } + select { + case err := <-done: + if err == nil { + t.Fatal("canceled proof succeeded") + } + case <-readCtx.Done(): + t.Fatal("cancellation left proof read blocked") + } + if raw, err := guest.Read(readCtx); err == nil || len(raw) > 0 { + t.Fatal("canceled peer received data") + } + }) + } +} + +func TestGuestStreamPreservesBufferedBytesForCaller(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + hostConn := relay.NetConn(a) + guest := relay.NetConn(z) + authority := streamAuthorizer(func(ctx context.Context, _ string, prove GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + _, err := prove(ctx, streamChallenge()) + return streamAcceptance(), err + }) + done := make(chan error, 1) + go func() { _, err := AuthorizeGuestConnection(ctx, authority, "stream-test", hostConn); done <- err }() + event, err := relay.ReadGuestReconnectFrame(ctx, guest) + if err != nil { + t.Fatal(err) + } + challenge, _ := runner.DecodeGuestReconnectChallenge(event.Payload) + body, _ := json.Marshal(runner.GuestReconnectAcceptRequest{Protocol: 1, AttemptID: challenge.AttemptID, Signature: base64.RawURLEncoding.EncodeToString(make([]byte, 64))}) + payload, _ := json.Marshal(relay.ControlEvent{Kind: relay.KindGuestReconnectProof, Payload: body}) + frame, _ := relay.Encode(relay.Frame{Type: relay.FrameControl, Payload: payload}) + written := make(chan error, 1) + go func() { _, err := z.Write(append(frame, []byte("\nnext-frame\n")...)); written <- err }() + if err = <-done; err != nil { + t.Fatal(err) + } + raw, err := hostConn.Read(ctx) + if err != nil || string(raw) != "next-frame" { + t.Fatal("proof reader lost subsequent bytes") + } + if err = <-written; err != nil { + t.Fatal(err) + } +} + +func TestGuestStreamRefusesLateHostSuccess(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + authority := streamAuthorizer(func(ctx context.Context, _ string, prove GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) { + _, err := prove(ctx, streamChallenge()) + if err != nil { + return runner.GuestReconnectAcceptResponse{}, err + } + cancel() + return streamAcceptance(), nil + }) + done := make(chan error, 1) + go func() { + got, err := AuthorizeGuestConnection(ctx, authority, "stream-test", relay.NetConn(a)) + if got != (runner.GuestReconnectAcceptResponse{}) { + done <- errors.New("late nonzero authority") + return + } + done <- err + }() + guest := relay.NetConn(z) + readCtx, readCancel := context.WithTimeout(context.Background(), time.Second) + defer readCancel() + event, err := relay.ReadGuestReconnectFrame(readCtx, guest) + if err != nil { + t.Fatal(err) + } + ch, _ := runner.DecodeGuestReconnectChallenge(event.Payload) + body, _ := json.Marshal(runner.GuestReconnectAcceptRequest{Protocol: 1, AttemptID: ch.AttemptID, Signature: base64.RawURLEncoding.EncodeToString(make([]byte, 64))}) + _ = relay.WriteGuestReconnectFrame(readCtx, guest, relay.KindGuestReconnectProof, body) + if err = <-done; err == nil || err.Error() != "unavailable" { + t.Fatalf("late acceptance: %v", err) + } + if raw, err := guest.Read(readCtx); err == nil || len(raw) > 0 { + t.Fatal("late authority reached guest") + } +} diff --git a/internal/runnerd/reconnect_process_test.go b/internal/runnerd/reconnect_process_test.go index 9c48c46..44e19d3 100644 --- a/internal/runnerd/reconnect_process_test.go +++ b/internal/runnerd/reconnect_process_test.go @@ -16,8 +16,8 @@ import ( ) // No shipping command invokes the optional callback yet. Build its public call -// site and exercise that artifact across a real agent WebSocket; do not imply -// that a live VM, guest key owner or relay takeover was exercised. +// site and bounded guest transport across real agent WebSocket and guest TCP +// streams. No live VM, enrolled guest key owner or relay takeover is exercised. func TestReconnectHostBuiltExecution(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), time.Minute) t.Cleanup(cancel) @@ -74,9 +74,13 @@ func TestReconnectHostBuiltExecution(t *testing.T) { if result := conn.readMsg(t); result.Type != "result" || !result.OK { t.Fatal("probe placement failed") } - for _, tc := range []struct{ name, want string }{{"valid", "accepted"}, {"wrong_scope", "invalid"}, {"revoked", "fenced"}} { + for _, tc := range []struct{ name, want string }{{"valid", "accepted"}, {"wrong_scope", "invalid"}, {"revoked", "fenced"}, {"wrong_attempt", "unavailable"}, {"oversize", "unavailable"}} { t.Run(tc.name, func(t *testing.T) { - if _, err := io.WriteString(input, "authorize\n"); err != nil { + command := "authorize" + if tc.name == "wrong_attempt" || tc.name == "oversize" { + command = tc.name + } + if _, err := io.WriteString(input, command+"\n"); err != nil { t.Fatal(err) } req := conn.readMsg(t) @@ -86,6 +90,8 @@ func TestReconnectHostBuiltExecution(t *testing.T) { switch tc.name { case "revoked": replyReconnect(t, conn, req, false, []byte(`{"error":"fenced"}`)) + case "wrong_attempt", "oversize": + replyReconnect(t, conn, req, true, reconnectChallengeJSON()) case "wrong_scope": replyReconnect(t, conn, req, true, []byte(strings.Replace(string(reconnectChallengeJSON()), "session_test", "wrong_test", 1))) default: @@ -109,5 +115,5 @@ func TestReconnectHostBuiltExecution(t *testing.T) { } }) } - t.Log("built host callback: accepted proof, wrong-scope and revoked authority paths passed; synthetic driver/guest, real agent WebSocket") + t.Log("built host/stream bridge: accepted proof, wrong-scope, revoked, wrong-attempt and oversized guest paths passed; synthetic driver/guest, real WebSocket and TCP streams") } diff --git a/internal/runnerd/testdata/reconnect-host/main.go b/internal/runnerd/testdata/reconnect-host/main.go index caa5711..6ff1dad 100644 --- a/internal/runnerd/testdata/reconnect-host/main.go +++ b/internal/runnerd/testdata/reconnect-host/main.go @@ -1,6 +1,7 @@ // A built execution probe for the optional host callback, not a shipping command. -// The driver and guest proof provider are synthetic; RunAgent and the host callback -// are real. No capability or local RPC endpoint is enabled by this probe. +// The driver and signing guest are synthetic. RunAgent, host authorization and +// bounded guest stream transport run over real WebSocket and TCP connections. +// No capability or shipping listener is enabled by this probe. package main import ( @@ -8,11 +9,17 @@ import ( "context" "crypto/ed25519" "encoding/base64" + "encoding/json" + "errors" "fmt" "github.com/tokencanopy/rainier/internal/driver" + "github.com/tokencanopy/rainier/internal/relay" "github.com/tokencanopy/rainier/internal/runnerd" "github.com/tokencanopy/rainier/protocol/runner" + "net" "os" + "strings" + "time" ) func main() { @@ -24,19 +31,13 @@ func main() { defer close(done) _ = s.RunAgent(ctx, runnerd.AgentConfig{ControldURL: os.Args[1], Token: "testtoken", RunnerName: "runner_process_test"}) }() - private := ed25519.NewKeyFromSeed(make([]byte, 32)) input := bufio.NewScanner(os.Stdin) for input.Scan() { - if input.Text() != "authorize" { + mode := input.Text() + if mode != "authorize" && mode != "wrong_attempt" && mode != "oversize" { break } - out, err := s.AuthorizeGuestReconnect(ctx, "session_test", func(ctx context.Context, c runner.GuestReconnectChallenge) (string, error) { - message, err := c.SigningMessage() - if err != nil { - return "", err - } - return base64.RawURLEncoding.EncodeToString(ed25519.Sign(private, message)), nil - }) + out, err := authorizeStream(ctx, s, mode) if err != nil { fmt.Println(err.Error()) } else if out.Epoch == 2 && out.ExpiresInSec == 120 { @@ -48,3 +49,66 @@ func main() { cancel() <-done } + +func authorizeStream(parent context.Context, s *runnerd.Server, mode string) (runner.GuestReconnectAcceptResponse, error) { + var zero runner.GuestReconnectAcceptResponse + ctx, cancel := context.WithTimeout(parent, 6*time.Second) + defer cancel() + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + return zero, errors.New("unavailable") + } + defer listener.Close() + peer, err := (&net.Dialer{}).DialContext(ctx, "tcp", listener.Addr().String()) + if err != nil { + return zero, errors.New("unavailable") + } + defer peer.Close() + host, err := listener.Accept() + if err != nil { + return zero, errors.New("unavailable") + } + defer host.Close() + guestDone := make(chan error, 1) + go func() { + conn := relay.NetConn(peer) + event, err := relay.ReadGuestReconnectFrame(ctx, conn) + if err != nil { + guestDone <- nil + return + } // authority refused before a challenge + if event.Kind != relay.KindGuestReconnectChallenge { + guestDone <- errors.New("unexpected") + return + } + challenge, err := runner.DecodeGuestReconnectChallenge(event.Payload) + if err != nil { + guestDone <- errors.New("unexpected") + return + } + if mode == "oversize" { + _, _ = peer.Write([]byte(strings.Repeat("x", 4097))) + } else { + message, _ := challenge.SigningMessage() + private := ed25519.NewKeyFromSeed(make([]byte, 32)) + proof := runner.GuestReconnectAcceptRequest{Protocol: 1, AttemptID: challenge.AttemptID, Signature: base64.RawURLEncoding.EncodeToString(ed25519.Sign(private, message))} + if mode == "wrong_attempt" { + proof.AttemptID = "other_test" + } + body, _ := json.Marshal(proof) + _ = relay.WriteGuestReconnectFrame(ctx, conn, relay.KindGuestReconnectProof, body) + } + // Neither acceptance nor configuration belongs to this transport helper. + if raw, err := conn.Read(ctx); err == nil || len(raw) > 0 { + guestDone <- errors.New("unexpected") + return + } + guestDone <- nil + }() + accepted, err := driver.AuthorizeGuestConnection(ctx, s, "session_test", relay.NetConn(host)) + _ = host.Close() + if guestErr := <-guestDone; guestErr != nil { + return zero, guestErr + } + return accepted, err +}