diff --git a/cmd/sessiond/bootconfig.go b/cmd/sessiond/bootconfig.go index fe82111..cbfc014 100644 --- a/cmd/sessiond/bootconfig.go +++ b/cmd/sessiond/bootconfig.go @@ -65,6 +65,10 @@ const ( // each connection and, when that configuration brings a token it has not // already spent, exchanges it for the session's secrets. // +// Opted-in fresh boots instead retain a signing identity and require the +// fail-closed reconnect preamble in reconnect.go. The following describes the +// unchanged legacy mode. +// // It is stateful for one reason, and it is the reason a redial is not a // resume: a bootstrap token is SINGLE-USE. A sessiond that re-exchanged on // every connection would be refused on the first redial inside a live VM and @@ -74,7 +78,11 @@ const ( // exactly "at boot and again after every resume", with no second signal // needed. type bootstrapper struct { - mu sync.Mutex + mu sync.Mutex + guest *guestReconnectIdentity + reconnecting bool + configured []string + configurationApplied func(context.Context) error // attempted is the last token this session sent a request for, whatever // the answer was. // @@ -161,6 +169,9 @@ func readBootConfig(ctx context.Context, conn relay.Conn) (runner.BootConfig, er // dispatcher because the dispatcher is not serving yet: the relay it rides on // is started later, by dialLoop, once there is a session to serve. func (b *bootstrapper) exchange(ctx context.Context, conn relay.Conn, cfg runner.BootConfig) (map[string]string, error) { + if cfg.GuestReconnect != 0 { + return b.enrollGuest(ctx, conn, cfg) + } if len(cfg.SecretNames) == 0 { // A clean boot. "No secrets declared" and "declared and never // arrived" are different facts, and this is the first. @@ -291,11 +302,16 @@ func (b *bootstrapper) forget() int { // // RAINIER_DIAL is deliberately NOT set. There is nothing to dial. func applyBootConfig(cfg runner.BootConfig, secrets map[string]string) error { + return visitBootConfig(cfg, secrets, os.Setenv) +} + +// visitBootConfig shares the legacy precedence and encoding with reconnect. +func visitBootConfig(cfg runner.BootConfig, secrets map[string]string, put func(string, string) error) error { set := func(k, v string) error { if v == "" { return nil } - return os.Setenv(k, v) + return put(k, v) } b64 := func(s string) string { if s == "" { @@ -312,7 +328,7 @@ func applyBootConfig(cfg runner.BootConfig, secrets map[string]string) error { // either. (The plane also drops such a name from SecretNames, so this is // the second of two fences rather than the only one.) for k, v := range secrets { - if err := os.Setenv(k, v); err != nil { + if err := put(k, v); err != nil { return fmt.Errorf("applying this session's environment: %w", err) } } @@ -364,7 +380,7 @@ func applyBootConfig(cfg runner.BootConfig, secrets map[string]string) error { return err } for k, v := range cfg.Env { - if err := os.Setenv(k, v); err != nil { + if err := put(k, v); err != nil { return fmt.Errorf("applying this session's configuration: %w", err) } } @@ -422,7 +438,13 @@ func bootOverVsock(ctx context.Context, dial dialSession, b *bootstrapper) (rela _ = conn.Close() conn = nil } - if err := applyBootConfig(cfg, secrets); err != nil { + apply := applyBootConfig + if cfg.GuestReconnect != 0 { + apply = func(cfg runner.BootConfig, secrets map[string]string) error { + return b.refreshConfiguration(ctx, cfg, secrets) + } + } + if err := apply(cfg, secrets); err != nil { if conn != nil { _ = conn.Close() } @@ -434,7 +456,8 @@ func bootOverVsock(ctx context.Context, dial dialSession, b *bootstrapper) (rela return conn, cfg, failure, nil } -// reBootstrap is the preamble on every LATER connection: read the +// reBootstrap selects authenticated reconnect for an opted-in guest. It never +// downgrades that guest to legacy bootstrap. In legacy mode it reads the // configuration again, and exchange again when it brings a token this // session has not spent — which is what a cold resume brings and a redial // does not. @@ -453,10 +476,20 @@ func bootOverVsock(ctx context.Context, dial dialSession, b *bootstrapper) (rela // which for a spent token is forever. So it is logged and the connection is // served. func reBootstrap(ctx context.Context, conn relay.Conn, b *bootstrapper) error { + b.mu.Lock() + guestMode := b.guest != nil + b.mu.Unlock() + if guestMode { + return b.reconnectGuest(ctx, conn) + } cfg, err := readBootConfig(ctx, conn) if err != nil { return err } + // Capability selection belongs to the first boot, never a legacy redial. + if cfg.GuestReconnect != 0 { + return errGuestReconnect + } secrets, xerr := b.exchange(ctx, conn, cfg) if xerr != nil { log.Printf("this session's secrets were not re-delivered on a new connection (%v); "+ diff --git a/cmd/sessiond/main.go b/cmd/sessiond/main.go index b81617c..af1ac30 100644 --- a/cmd/sessiond/main.go +++ b/cmd/sessiond/main.go @@ -251,6 +251,9 @@ func main() { execs := sandboxexec.NewRunner(workspaceRoot, sandboxexec.SessionEnv(os.Environ(), envAssignments(chainVars)), sandboxexec.NewSpawner().Start) + if overVsock { + boots.bindExecEnvironment(execs, envAssignments(chainVars)) + } // What the RUNNER needs to know about those commands, and the only thing // it needs: how many are running. runnerd's idle auto-stop stops a session diff --git a/cmd/sessiond/reconnect.go b/cmd/sessiond/reconnect.go new file mode 100644 index 0000000..546aaf4 --- /dev/null +++ b/cmd/sessiond/reconnect.go @@ -0,0 +1,315 @@ +package main + +import ( + "bytes" + "context" + "crypto/ed25519" + "crypto/rand" + "encoding/base64" + "encoding/json" + "errors" + "io" + "os" + "strings" + "time" + + "github.com/tokencanopy/rainier/internal/relay" + "github.com/tokencanopy/rainier/internal/sandboxexec" + "github.com/tokencanopy/rainier/protocol/runner" +) + +const guestHandshakeWait = 5 * time.Second + +var errGuestReconnect = errors.New("guest reconnect unavailable") + +// This identity belongs to one sessiond lifetime. Never serialize it, log it, +// export it through the environment or replace it after ambiguous enrollment. +type guestReconnectIdentity struct { + session, boot string + key ed25519.PrivateKey + enrolled bool + epoch uint64 +} + +func (b *bootstrapper) enrollGuest(ctx context.Context, c relay.Conn, cfg runner.BootConfig) (map[string]string, error) { + b.mu.Lock() + if b.guest != nil { + b.mu.Unlock() + return nil, errGuestReconnect + } + // A failed or unsupported opt-in is sticky: later dials cannot downgrade. + g := &guestReconnectIdentity{session: cfg.SessionID} + b.guest = g + b.mu.Unlock() + if cfg.GuestReconnect != runner.GuestReconnectProtocol || !guestID(cfg.SessionID) { + return nil, errGuestReconnect + } + pub, key, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + return nil, errGuestReconnect + } + var boot [32]byte + if _, err = rand.Read(boot[:]); err != nil { + return nil, errGuestReconnect + } + request := runner.GuestReconnectEnrollRequest{Protocol: 1, Token: cfg.BootstrapToken, BootEpoch: base64.RawURLEncoding.EncodeToString(boot[:]), PublicKey: base64.RawURLEncoding.EncodeToString(pub)} + body, _ := json.Marshal(request) + if _, err = runner.DecodeGuestReconnectEnrollRequest(body); err != nil { + return nil, errGuestReconnect + } + env, err := guestEnvironmentExchange(ctx, c, runner.MethodEnrollGuestReconnect, body, cfg.SecretNames) + if err != nil { + return nil, errGuestReconnect + } + b.mu.Lock() + g.key = key + g.boot = request.BootEpoch + g.enrolled = true + b.attempted = cfg.BootstrapToken + b.mu.Unlock() + return env, nil +} + +func guestID(s string) bool { + if len(s) == 0 || len(s) > 256 { + return false + } + for i := range s { + if s[i] < 0x21 || s[i] > 0x7e { + return false + } + } + return true +} + +// reconnectGuest is a preamble only: it never starts setup, init, an agent or a +// relay. The caller may serve its retained session only after this succeeds. +func (b *bootstrapper) reconnectGuest(ctx context.Context, c relay.Conn) error { + b.mu.Lock() + g := b.guest + if g == nil || !g.enrolled || b.reconnecting { + b.mu.Unlock() + return errGuestReconnect + } + b.reconnecting = true + session, boot, key, epoch := g.session, g.boot, g.key, g.epoch + b.mu.Unlock() + defer func() { b.mu.Lock(); b.reconnecting = false; b.mu.Unlock() }() + accepted, err := guestProof(ctx, c, session, boot, key, epoch) + if err != nil { + return errGuestReconnect + } + // Acceptance consumes this epoch even if delivery is subsequently lost. + b.mu.Lock() + g.epoch = accepted.Epoch + b.mu.Unlock() + delivery, cancel := context.WithTimeout(ctx, time.Duration(accepted.ExpiresInSec)*time.Second) + defer cancel() + cfg, err := readBootConfig(delivery, c) + if err != nil || cfg.SessionID != session || cfg.GuestReconnect != 1 || cfg.BootstrapToken != accepted.Token { + return errGuestReconnect + } + b.mu.Lock() + repeated := b.attempted == cfg.BootstrapToken + b.attempted = cfg.BootstrapToken + b.mu.Unlock() + if repeated { + return errGuestReconnect + } + body, _ := json.Marshal(struct { + Protocol int `json:"protocol"` + Token string `json:"token"` + }{1, cfg.BootstrapToken}) + env, err := guestEnvironmentExchange(delivery, c, runner.MethodFetchSessionSecrets, body, cfg.SecretNames) + if err != nil || delivery.Err() != nil { + return errGuestReconnect + } + if err = b.refreshConfiguration(delivery, cfg, env); err != nil { + return errGuestReconnect + } + return nil +} + +func guestProof(ctx context.Context, c relay.Conn, session, boot string, key ed25519.PrivateKey, epoch uint64) (runner.GuestReconnectAcceptResponse, error) { + var zero runner.GuestReconnectAcceptResponse + ctx, cancel := context.WithTimeout(ctx, guestHandshakeWait) + defer cancel() + ev, err := relay.ReadGuestReconnectFrame(ctx, c) + if err != nil || ev.Kind != relay.KindGuestReconnectChallenge { + return zero, errGuestReconnect + } + challenge, err := runner.DecodeGuestReconnectChallenge(ev.Payload) + if err != nil || challenge.SessionID != session || challenge.BootEpoch != boot { + return zero, errGuestReconnect + } + message, err := challenge.SigningMessage() + if err != nil { + return zero, errGuestReconnect + } + proof := runner.GuestReconnectAcceptRequest{Protocol: 1, AttemptID: challenge.AttemptID, Signature: base64.RawURLEncoding.EncodeToString(ed25519.Sign(key, message))} + body, _ := json.Marshal(proof) + if relay.WriteGuestReconnectFrame(ctx, c, relay.KindGuestReconnectProof, body) != nil { + return zero, errGuestReconnect + } + ev, err = relay.ReadGuestReconnectFrame(ctx, c) + if err != nil || ev.Kind != relay.KindGuestReconnectAccepted { + return zero, errGuestReconnect + } + accepted, err := runner.DecodeGuestReconnectAcceptResponse(ev.Payload) + if err != nil || accepted.Epoch <= epoch || ctx.Err() != nil { + return zero, errGuestReconnect + } + return accepted, nil +} + +// Both enrollment and refresh must redeem authority even for an empty declared +// secret set. These are authorized payloads under the normal transport limit, +// not unauthenticated cryptographic frames. Never expose provider/parser text. +func guestEnvironmentExchange(ctx context.Context, c relay.Conn, method string, body []byte, names []string) (map[string]string, error) { + ctx, cancel := context.WithTimeout(ctx, secretsExchangeWait) + defer cancel() + payload, _ := json.Marshal(relay.ControlEvent{Kind: "req:" + method, ID: secretsRequestID, Payload: body}) + frame, _ := relay.Encode(relay.Frame{Type: relay.FrameControl, Payload: payload}) + if c.Write(ctx, frame) != nil { + return nil, errGuestReconnect + } + raw, err := c.Read(ctx) + if err != nil { + return nil, errGuestReconnect + } + f, err := relay.Decode(raw) + if err != nil || f.Type != relay.FrameControl || f.AttachID != 0 { + return nil, errGuestReconnect + } + var ev relay.ControlEvent + if json.Unmarshal(f.Payload, &ev) != nil || ev.Kind != "resp" || ev.ID != secretsRequestID || !ev.OK { + return nil, errGuestReconnect + } + env, err := decodeGuestEnvironment(ev.Payload) + if err != nil || ctx.Err() != nil { + return nil, errGuestReconnect + } + for _, n := range names { + if _, ok := env[n]; !ok { + return nil, errGuestReconnect + } + } + return env, nil +} + +// Require exactly one non-null environment object, with unique string values. +func decodeGuestEnvironment(raw []byte) (map[string]string, error) { + d := json.NewDecoder(bytes.NewReader(raw)) + tok, err := d.Token() + if err != nil || tok != json.Delim('{') { + return nil, errGuestReconnect + } + tok, err = d.Token() + if err != nil || tok != "env" { + return nil, errGuestReconnect + } + tok, err = d.Token() + if err != nil || tok != json.Delim('{') { + return nil, errGuestReconnect + } + env := map[string]string{} + for d.More() { + tok, err = d.Token() + k, ok := tok.(string) + if err != nil || !ok { + return nil, errGuestReconnect + } + if _, exists := env[k]; exists { + return nil, errGuestReconnect + } + tok, err = d.Token() + v, ok := tok.(string) + if err != nil || !ok { + return nil, errGuestReconnect + } + if !validGuestEnv(k, v) { + return nil, errGuestReconnect + } + env[k] = v + } + tok, err = d.Token() + if err != nil || tok != json.Delim('}') { + return nil, errGuestReconnect + } + tok, err = d.Token() + if err != nil || tok != json.Delim('}') { + return nil, errGuestReconnect + } + if _, err = d.Token(); err != io.EOF { + return nil, errGuestReconnect + } + return env, nil +} +func validGuestEnv(k, v string) bool { + return k != "" && !strings.ContainsAny(k, "=\x00") && !strings.ContainsRune(v, '\x00') +} + +// Replace only keys previously owned by this boot configuration. Removing an +// empty proxy/script/secret setting must remove its old process value as well. +// Existing children keep their inherited environments; this updates sessiond +// and future children, not an already running coding agent's credentials. +func (b *bootstrapper) refreshConfiguration(ctx context.Context, cfg runner.BootConfig, secrets map[string]string) error { + env := map[string]string{} + if visitBootConfig(cfg, secrets, func(k, v string) error { + if ctx.Err() != nil || !validGuestEnv(k, v) { + return errGuestReconnect + } + env[k] = v + return nil + }) != nil { + return errGuestReconnect + } + b.mu.Lock() + defer b.mu.Unlock() + if ctx.Err() != nil { + return errGuestReconnect + } + for _, k := range b.configured { + if _, ok := env[k]; !ok { + if os.Unsetenv(k) != nil { + return errGuestReconnect + } + } + } + for k, v := range env { + if os.Setenv(k, v) != nil { + return errGuestReconnect + } + } + b.configured = namesOf(env) + b.delivered = namesOf(secrets) + if b.configurationApplied != nil { + if b.configurationApplied(ctx) != nil { + return errGuestReconnect + } + } + if ctx.Err() != nil { + return errGuestReconnect + } + return nil +} + +// bindExecEnvironment is installed before dialLoop begins. Refresh the runner's +// immutable launch snapshot before the preamble can publish readiness. Existing +// execs and the agent retain their environments; boot-chain exports stay intact. +func (b *bootstrapper) bindExecEnvironment(execs *sandboxexec.Runner, extra []string) { + extra = append([]string(nil), extra...) + b.mu.Lock() + defer b.mu.Unlock() + b.configurationApplied = func(ctx context.Context) error { + if ctx.Err() != nil { + return errGuestReconnect + } + env := sandboxexec.SessionEnv(os.Environ(), extra) + if ctx.Err() != nil { + return errGuestReconnect + } + execs.ReplaceEnvironment(env) + return ctx.Err() + } +} diff --git a/cmd/sessiond/reconnect_process_test.go b/cmd/sessiond/reconnect_process_test.go new file mode 100644 index 0000000..7ea1266 --- /dev/null +++ b/cmd/sessiond/reconnect_process_test.go @@ -0,0 +1,234 @@ +package main + +import ( + "bytes" + "context" + "encoding/base64" + "encoding/json" + "fmt" + "net" + "os" + "os/exec" + "strings" + "testing" + "time" + + "github.com/tokencanopy/rainier/internal/relay" + "github.com/tokencanopy/rainier/internal/sandboxexec" + "github.com/tokencanopy/rainier/internal/session" + "github.com/tokencanopy/rainier/protocol/runner" + "github.com/tokencanopy/rainier/protocol/terminal" +) + +// A separately executed guest uses the production boot/preamble and PTY code +// over real TCP streams. The host is a protocol fixture: shipping host opt-in, +// AF_VSOCK/KVM, policy authorization and relay takeover remain separate gates. +func TestGuestReconnectExecutable(t *testing.T) { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer ln.Close() + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + cmd := exec.CommandContext(ctx, os.Args[0], "-test.run=^TestGuestReconnectProcessGuest$", "-test.v") + cmd.Env = append(os.Environ(), "RAINIER_TEST_GUEST_ADDRESS="+ln.Addr().String()) + var output bytes.Buffer + cmd.Stdout = &output + cmd.Stderr = &output + if err = cmd.Start(); err != nil { + t.Fatal(err) + } + waited := false + defer func() { + if !waited { + _ = cmd.Process.Kill() + _ = cmd.Wait() + } + }() + accept := func() relay.Conn { + t.Helper() + _ = ln.(*net.TCPListener).SetDeadline(time.Now().Add(5 * time.Second)) + c, err := ln.Accept() + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { c.Close() }) + return relay.NetConn(c) + } + initial := accept() + token := base64.RawURLEncoding.EncodeToString(make([]byte, 32)) + cfg := runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "process-test", BootstrapToken: token, Env: map[string]string{"RECONNECT_PROCESS_TEST": "initial", "RECONNECT_REMOVED_TEST": "old"}} + sendReconnectConfig(t, ctx, initial, cfg) + raw, err := initial.Read(ctx) + if err != nil { + t.Fatal(err) + } + f, _ := relay.Decode(raw) + var event relay.ControlEvent + if json.Unmarshal(f.Payload, &event) != nil || event.Kind != "req:"+runner.MethodEnrollGuestReconnect { + t.Fatal("missing initial enrollment") + } + enrolled, err := runner.DecodeGuestReconnectEnrollRequest(event.Payload) + if err != nil { + t.Fatal(err) + } + controlWrite(t, ctx, initial, relay.ControlEvent{Kind: "resp", ID: event.ID, OK: true, Payload: json.RawMessage(`{"env":{}}`)}) + initial.Close() + // First reconnect is explicitly fenced. A subsequent fresh challenge must + // still be usable, and the child shell must not be restarted in either case. + refused := accept() + writeReconnect(t, ctx, refused, relay.KindGuestReconnectRefused, runner.GuestReconnectErrorResponse{Error: "fenced"}) + refused.Close() + host := accept() + challenge := runner.GuestReconnectChallenge{Protocol: 1, SessionID: cfg.SessionID, BootEpoch: enrolled.BootEpoch, HostIncarnation: "host-test", AttemptID: "process-attempt", PlacementGeneration: 1, Challenge: token} + writeReconnect(t, ctx, host, relay.KindGuestReconnectChallenge, challenge) + event, err = relay.ReadGuestReconnectFrame(ctx, host) + if err != nil { + t.Fatal(err) + } + proof, err := runner.DecodeGuestReconnectAcceptRequest(event.Payload) + if err != nil || event.Kind != relay.KindGuestReconnectProof { + t.Fatal("invalid process proof") + } + if err = challenge.VerifyProof(enrolled.PublicKey, proof.Signature); err != nil { + t.Fatal(err) + } + tokenBytes := make([]byte, 32) + tokenBytes[0] = 7 + fresh := base64.RawURLEncoding.EncodeToString(tokenBytes) + writeReconnect(t, ctx, host, relay.KindGuestReconnectAccepted, runner.GuestReconnectAcceptResponse{Epoch: 1, Token: fresh, ExpiresInSec: 120}) + cfg.BootstrapToken = fresh + cfg.Env = map[string]string{"RECONNECT_PROCESS_TEST": "refreshed"} + sendReconnectConfig(t, ctx, host, cfg) + raw, err = host.Read(ctx) + if err != nil { + t.Fatal(err) + } + f, _ = relay.Decode(raw) + if json.Unmarshal(f.Payload, &event) != nil || event.Kind != "req:"+runner.MethodFetchSessionSecrets { + t.Fatal("missing fresh redemption") + } + var request struct { + Token string `json:"token"` + } + _ = json.Unmarshal(event.Payload, &request) + if request.Token != fresh { + t.Fatal("wrong process token") + } + controlWrite(t, ctx, host, relay.ControlEvent{Kind: "resp", ID: event.ID, OK: true, Payload: json.RawMessage(`{"env":{}}`)}) + err = cmd.Wait() + waited = true + if err != nil { + t.Fatalf("guest process: %v\n%s", err, output.String()) + } + if !strings.Contains(output.String(), "retained PTY shell and refreshed guest configuration") { + t.Fatalf("missing process evidence: %s", output.String()) + } + if strings.Contains(output.String(), enrolled.PublicKey) || strings.Contains(output.String(), fresh) { + t.Fatal("guest logged credentials") + } + t.Log("separate guest process: enrollment, refused reconnect, verified proof, fresh redemption, same PTY shell across reconnect; no credential output") +} + +func TestGuestReconnectProcessGuest(t *testing.T) { + address := os.Getenv("RAINIER_TEST_GUEST_ADDRESS") + if address == "" { + t.Skip("subprocess helper") + } + ctx, cancel := context.WithTimeout(context.Background(), 12*time.Second) + defer cancel() + dial := func(ctx context.Context) (relay.Conn, error) { + c, err := (&net.Dialer{}).DialContext(ctx, "tcp", address) + if err != nil { + return nil, err + } + return relay.NetConn(c), nil + } + b := &bootstrapper{} + c, _, failure, err := bootOverVsock(ctx, dial, b) + if err != nil || failure != nil { + t.Fatalf("boot: %v %v", err, failure) + } + c.Close() + chunks := make(chan string, 32) + p, err := session.StartProc([]string{"/bin/sh", "-c", `stty -echo; while read line; do printf 'SHELL:%s:%s:%s\n' "$$" "$RECONNECT_PROCESS_TEST" "$line"; done`}, 80, 24, func(raw []byte) { chunks <- string(raw) }) + if err != nil { + t.Fatal(err) + } + defer func() { p.Stop(); p.Wait() }() + turn := func(label string) string { + t.Helper() + if _, err := p.Write([]byte(label + "\n")); err != nil { + t.Fatal(err) + } + var output string + for !strings.Contains(output, ":"+label) { + select { + case s := <-chunks: + output += s + case <-ctx.Done(): + t.Fatal("PTY turn timeout") + } + } + for _, line := range strings.Split(output, "\n") { + if strings.HasPrefix(line, "SHELL:") { + return strings.TrimSpace(line) + } + } + t.Fatalf("missing shell response: %s", output) + return "" + } + + extra := []string{"CHAIN_TEST=retained"} + execs := sandboxexec.NewRunner(t.TempDir(), sandboxexec.SessionEnv(os.Environ(), extra), sandboxexec.NewSpawner().Start) + b.bindExecEnvironment(execs, extra) + execTurn := func(want string) { + t.Helper() + a := execs.OpenExec(runner.ExecSpec{Argv: []string{"/bin/sh", "-c", `printf '%s:%s:%s' "$RECONNECT_PROCESS_TEST" "$CHAIN_TEST" "${RECONNECT_REMOVED_TEST-unset}"`}}) + defer a.Close() + var output string + for { + select { + case m, ok := <-a.Msgs(): + if !ok { + t.Fatal("exec closed without exit") + } + switch m.Type { + case terminal.TypeExecStdout: + output += string(m.Data) + case terminal.TypeExecError: + t.Fatal(m.Reason) + case terminal.TypeExecExit: + if m.ExitCode != 0 || m.Signal != "" || output != want { + t.Fatalf("future exec environment=%q, want %q", output, want) + } + return + } + case <-ctx.Done(): + t.Fatal("exec turn timeout") + } + } + } + execTurn("initial:retained:old") + before := turn("before") + tr := sessionTransport{dial: dial, preamble: func(ctx context.Context, c relay.Conn) error { return reBootstrap(ctx, c, b) }} + if c, err := tr.connect(ctx); err == nil || c != nil { + t.Fatal("refused connection became ready") + } + execTurn("initial:retained:old") + c, err = tr.connect(ctx) + if err != nil { + t.Fatal(err) + } + defer c.Close() + execTurn("refreshed:retained:unset") + after := turn("after") + if strings.TrimSuffix(before, ":before") != strings.TrimSuffix(after, ":after") { + t.Fatalf("shell restarted or inherited environment changed: %s / %s", before, after) + } + if !strings.Contains(after, ":initial:after") || os.Getenv("RECONNECT_PROCESS_TEST") != "refreshed" { + t.Fatal("guest/child environment contract violated") + } + fmt.Println("retained PTY shell and refreshed guest configuration") +} diff --git a/cmd/sessiond/reconnect_test.go b/cmd/sessiond/reconnect_test.go new file mode 100644 index 0000000..795a959 --- /dev/null +++ b/cmd/sessiond/reconnect_test.go @@ -0,0 +1,427 @@ +package main + +import ( + "context" + "crypto/ed25519" + "crypto/rand" + "encoding/base64" + "encoding/json" + "net" + "os" + "strings" + "testing" + "time" + + "github.com/tokencanopy/rainier/internal/relay" + "github.com/tokencanopy/rainier/protocol/runner" +) + +func TestReconnectEnrollmentWithoutSecrets(t *testing.T) { + a, z := net.Pipe() + defer a.Close() + defer z.Close() + b := &bootstrapper{} + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + done := make(chan error, 1) + go func() { + _, err := b.exchange(ctx, relay.NetConn(a), runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test", BootstrapToken: base64.RawURLEncoding.EncodeToString(make([]byte, 32))}) + done <- err + }() + host := relay.NetConn(z) + raw, err := host.Read(ctx) + if err != nil { + t.Fatalf("no enrollment for secret-free boot: %v", err) + } + f, err := relay.Decode(raw) + if err != nil { + t.Fatal(err) + } + var ev relay.ControlEvent + if err = json.Unmarshal(f.Payload, &ev); err != nil { + t.Fatal(err) + } + if ev.Kind != "req:"+runner.MethodEnrollGuestReconnect { + t.Fatalf("method=%s", ev.Kind) + } + enroll, err := runner.DecodeGuestReconnectEnrollRequest(ev.Payload) + if err != nil { + t.Fatal(err) + } + if enroll.BootEpoch == "" || enroll.PublicKey == "" { + t.Fatal("missing guest identity") + } + body, _ := json.Marshal(relay.ControlEvent{Kind: "resp", ID: ev.ID, OK: true, Payload: json.RawMessage(`{"env":{}}`)}) + frame, _ := relay.Encode(relay.Frame{Type: relay.FrameControl, Payload: body}) + if err = host.Write(ctx, frame); err != nil { + t.Fatal(err) + } + if err = <-done; err != nil { + t.Fatal(err) + } +} + +func reconnectPair(t *testing.T) (relay.Conn, relay.Conn, context.Context) { + t.Helper() + a, z := net.Pipe() + t.Cleanup(func() { a.Close(); z.Close() }) + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + t.Cleanup(cancel) + return relay.NetConn(a), relay.NetConn(z), ctx +} +func reconnectFixture(t *testing.T) (*bootstrapper, runner.GuestReconnectChallenge) { + t.Helper() + _, key, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + return &bootstrapper{guest: &guestReconnectIdentity{session: "session-test", boot: "boot-test", key: key, enrolled: true}}, runner.GuestReconnectChallenge{Protocol: 1, SessionID: "session-test", BootEpoch: "boot-test", HostIncarnation: "host-test", AttemptID: "attempt-test", PlacementGeneration: 1, Challenge: base64.RawURLEncoding.EncodeToString(make([]byte, 32))} +} +func writeReconnect(t *testing.T, ctx context.Context, c relay.Conn, kind string, v any) { + t.Helper() + body, err := json.Marshal(v) + if err != nil { + t.Fatal(err) + } + if err = relay.WriteGuestReconnectFrame(ctx, c, kind, body); err != nil { + t.Fatal(err) + } +} +func controlWrite(t *testing.T, ctx context.Context, c relay.Conn, event relay.ControlEvent) { + t.Helper() + body, _ := json.Marshal(event) + raw, _ := relay.Encode(relay.Frame{Type: relay.FrameControl, Payload: body}) + if err := c.Write(ctx, raw); err != nil { + t.Fatal(err) + } +} +func acceptProof(t *testing.T, ctx context.Context, host relay.Conn, b *bootstrapper, ch runner.GuestReconnectChallenge, epoch uint64) runner.GuestReconnectAcceptResponse { + t.Helper() + writeReconnect(t, ctx, host, relay.KindGuestReconnectChallenge, ch) + ev, err := relay.ReadGuestReconnectFrame(ctx, host) + if err != nil { + t.Fatal(err) + } + if ev.Kind != relay.KindGuestReconnectProof { + t.Fatalf("kind=%s", ev.Kind) + } + proof, err := runner.DecodeGuestReconnectAcceptRequest(ev.Payload) + if err != nil { + t.Fatal(err) + } + if proof.AttemptID != ch.AttemptID { + t.Fatal("wrong attempt") + } + if err = ch.VerifyProof(base64.RawURLEncoding.EncodeToString(b.guest.key.Public().(ed25519.PublicKey)), proof.Signature); err != nil { + t.Fatal(err) + } + tokenBytes := make([]byte, 32) + tokenBytes[0] = byte(epoch) + accepted := runner.GuestReconnectAcceptResponse{Epoch: epoch, Token: base64.RawURLEncoding.EncodeToString(tokenBytes), ExpiresInSec: 120} + writeReconnect(t, ctx, host, relay.KindGuestReconnectAccepted, accepted) + return accepted +} +func sendReconnectConfig(t *testing.T, ctx context.Context, host relay.Conn, cfg runner.BootConfig) { + t.Helper() + body, _ := json.Marshal(cfg) + controlWrite(t, ctx, host, relay.ControlEvent{Kind: relay.KindBootConfig, Payload: body}) +} + +func TestReconnectRefreshEmptySecretsRemovesOldConfiguration(t *testing.T) { + b, ch := reconnectFixture(t) + for _, k := range []string{"RAINIER_SESSION", "HTTP_PROXY", "http_proxy", "HTTPS_PROXY", "https_proxy", "NO_PROXY", "no_proxy", "OLD_TEST_SECRET", "OLD_TEST_CONFIG", "NEW_TEST_CONFIG", "UNOWNED_TEST"} { + t.Setenv(k, "") + } + t.Setenv("UNOWNED_TEST", "keep") + if err := b.refreshConfiguration(context.Background(), runner.BootConfig{SessionID: "session-test", ProxyURL: "http://proxy.invalid", NoProxy: "localhost", Env: map[string]string{"OLD_TEST_CONFIG": "old"}}, map[string]string{"OLD_TEST_SECRET": "synthetic"}); err != nil { + t.Fatal(err) + } + guest, host, ctx := reconnectPair(t) + done := make(chan error, 1) + go func() { done <- reBootstrap(ctx, guest, b) }() + accepted := acceptProof(t, ctx, host, b, ch, 1) + sendReconnectConfig(t, ctx, host, runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test", BootstrapToken: accepted.Token, Env: map[string]string{"NEW_TEST_CONFIG": "new"}}) + raw, err := host.Read(ctx) + if err != nil { + t.Fatal(err) + } + f, _ := relay.Decode(raw) + var request relay.ControlEvent + if json.Unmarshal(f.Payload, &request) != nil || request.Kind != "req:"+runner.MethodFetchSessionSecrets { + t.Fatalf("no fresh redemption: %s", f.Payload) + } + var body struct { + Token string `json:"token"` + } + _ = json.Unmarshal(request.Payload, &body) + if body.Token != accepted.Token { + t.Fatal("stale token") + } + controlWrite(t, ctx, host, relay.ControlEvent{Kind: "resp", ID: request.ID, OK: true, Payload: json.RawMessage(`{"env":{}}`)}) + if err = <-done; err != nil { + t.Fatal(err) + } + for _, k := range []string{"OLD_TEST_SECRET", "OLD_TEST_CONFIG", "HTTP_PROXY", "NO_PROXY"} { + if _, ok := os.LookupEnv(k); ok { + t.Fatalf("stale %s", k) + } + } + if os.Getenv("NEW_TEST_CONFIG") != "new" || os.Getenv("UNOWNED_TEST") != "keep" { + t.Fatal("configuration not refreshed") + } + if b.guest.epoch != 1 { + t.Fatal("epoch not retained") + } +} + +func TestReconnectRefusesUntrustedChallenge(t *testing.T) { + for _, which := range []string{"session", "boot", "protocol", "configuration", "oversize", "cancel"} { + t.Run(which, func(t *testing.T) { + b, ch := reconnectFixture(t) + guest, host, ctx := reconnectPair(t) + if which == "cancel" { + var cancel context.CancelFunc + ctx, cancel = context.WithCancel(ctx) + cancel() + } + done := make(chan error, 1) + go func() { done <- reBootstrap(ctx, guest, b); guest.Close() }() + switch which { + case "session": + ch.SessionID = "other-test" + case "boot": + ch.BootEpoch = "other-boot" + case "protocol": + ch.Protocol = 2 + } + if which == "oversize" { + _ = host.Write(ctx, []byte(strings.Repeat("x", 4097))) + } else if which != "cancel" { + body, _ := json.Marshal(ch) + kind := relay.KindGuestReconnectChallenge + if which == "configuration" { + kind = relay.KindBootConfig + } + payload, _ := json.Marshal(relay.ControlEvent{Kind: kind, Payload: body}) + frame, _ := relay.Encode(relay.Frame{Type: relay.FrameControl, Payload: payload}) + _ = host.Write(ctx, frame) + } + if err := <-done; err != errGuestReconnect { + t.Fatalf("failure=%v", err) + } + raw, err := host.Read(ctx) + if err == nil || len(raw) != 0 { + t.Fatal("refused peer received a proof") + } + if b.guest.epoch != 0 { + t.Fatal("refusal changed epoch") + } + }) + } +} + +func TestReconnectDeliveryFailsClosed(t *testing.T) { + for _, which := range []string{"replay", "refused", "bad-token", "wrong-session", "downgrade", "exchange-refusal", "null-env", "missing-secret", "bad-env"} { + t.Run(which, func(t *testing.T) { + b, ch := reconnectFixture(t) + guest, host, ctx := reconnectPair(t) + if which == "replay" { + b.guest.epoch = 1 + } + done := make(chan error, 1) + go func() { done <- reBootstrap(ctx, guest, b); guest.Close() }() + if which == "refused" { + writeReconnect(t, ctx, host, relay.KindGuestReconnectRefused, runner.GuestReconnectErrorResponse{Error: "fenced"}) + } else { + accepted := acceptProof(t, ctx, host, b, ch, 1) + if which != "replay" { + cfg := runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test", BootstrapToken: accepted.Token} + switch which { + case "bad-token": + cfg.BootstrapToken = "stale" + case "wrong-session": + cfg.SessionID = "other-test" + case "downgrade": + cfg.GuestReconnect = 0 + case "missing-secret": + cfg.SecretNames = []string{"ABSENT_TEST"} + } + sendReconnectConfig(t, ctx, host, cfg) + if which == "exchange-refusal" || which == "null-env" || which == "missing-secret" || which == "bad-env" { + if _, err := host.Read(ctx); err != nil { + t.Fatal(err) + } + payload := json.RawMessage(`{"env":{}}`) + if which == "null-env" { + payload = json.RawMessage(`{"env":null}`) + } + if which == "bad-env" { + payload = json.RawMessage(`{"env":{"BAD=KEY":"synthetic"}}`) + } + controlWrite(t, ctx, host, relay.ControlEvent{Kind: "resp", ID: secretsRequestID, OK: which != "exchange-refusal", Payload: payload}) + } + } + } + if err := <-done; err != errGuestReconnect { + t.Fatalf("accepted invalid delivery: %v", err) + } + if b.reconnecting { + t.Fatal("attempt slot leaked") + } + if which != "refused" && b.guest.epoch != 1 { + t.Fatal("accepted epoch lost after failed delivery") + } + }) + } +} + +func TestGuestEnvironmentStrict(t *testing.T) { + for _, raw := range []string{`null`, `{}`, `{"env":null}`, `{"env":{},"env":{}}`, `{"Env":{}}`, `{"env":{"A":null}}`, `{"env":{"A":"1","A":"2"}}`, `{"env":{}} {}`, `{"env":{},"extra":1}`} { + if _, err := decodeGuestEnvironment([]byte(raw)); err == nil { + t.Fatalf("accepted %s", raw) + } + } + env, err := decodeGuestEnvironment([]byte(`{"env":{"EMPTY":""}}`)) + if err != nil || len(env) != 1 { + t.Fatal("empty string must be valid") + } +} + +func TestGuestConfigurationValidationBeforeMutation(t *testing.T) { + b, _ := reconnectFixture(t) + t.Setenv("RETAIN_TEST", "original") + b.configured = []string{"RETAIN_TEST"} + if err := b.refreshConfiguration(context.Background(), runner.BootConfig{Env: map[string]string{"INVALID=KEY": "x"}}, nil); err == nil { + t.Fatal("accepted invalid environment") + } + if os.Getenv("RETAIN_TEST") != "original" { + t.Fatal("failed validation mutated environment") + } +} + +func TestLegacyReconnectCannotEnrollLater(t *testing.T) { + guest, host, ctx := reconnectPair(t) + b := &bootstrapper{} + done := make(chan error, 1) + go func() { done <- reBootstrap(ctx, guest, b); guest.Close() }() + sendReconnectConfig(t, ctx, host, runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test", BootstrapToken: base64.RawURLEncoding.EncodeToString(make([]byte, 32))}) + if raw, err := host.Read(ctx); err == nil || len(raw) != 0 { + t.Fatal("legacy reconnect attempted enrollment outside a fresh boot") + } + if err := <-done; err == nil { + t.Fatal("late opt-in accepted") + } +} + +func TestEnrollmentFailureIsSticky(t *testing.T) { + guest, host, ctx := reconnectPair(t) + b := &bootstrapper{} + done := make(chan error, 1) + cfg := runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test", BootstrapToken: base64.RawURLEncoding.EncodeToString(make([]byte, 32))} + go func() { _, err := b.exchange(ctx, guest, cfg); done <- err }() + if _, err := host.Read(ctx); err != nil { + t.Fatal(err) + } + controlWrite(t, ctx, host, relay.ControlEvent{Kind: "resp", ID: secretsRequestID, Payload: json.RawMessage(`{"error":"unavailable"}`)}) + if err := <-done; err == nil { + t.Fatal("refused enrollment succeeded") + } + if _, err := b.exchange(ctx, guest, cfg); err == nil { + t.Fatal("ambiguous enrollment retried") + } + if err := reBootstrap(ctx, guest, b); err == nil { + t.Fatal("failed enrollment reconnected") + } + if b.guest.enrolled || len(b.guest.key) != 0 { + t.Fatal("failed enrollment retained usable identity") + } +} + +func TestReconnectAttemptBudgetAndTransportClosure(t *testing.T) { + b, ch := reconnectFixture(t) + guest, host, ctx := reconnectPair(t) + tr := sessionTransport{dial: func(context.Context) (relay.Conn, error) { return guest, nil }, preamble: func(ctx context.Context, c relay.Conn) error { return reBootstrap(ctx, c, b) }} + done := make(chan error, 1) + go func() { + c, err := tr.connect(ctx) + if c != nil { + done <- nil + return + } + done <- err + }() + writeReconnect(t, ctx, host, relay.KindGuestReconnectChallenge, ch) + if _, err := relay.ReadGuestReconnectFrame(ctx, host); err != nil { + t.Fatal(err) + } + other, _, otherCtx := reconnectPair(t) + if err := reBootstrap(otherCtx, other, b); err != errGuestReconnect { + t.Fatal("parallel attempt not refused") + } + writeReconnect(t, ctx, host, relay.KindGuestReconnectRefused, runner.GuestReconnectErrorResponse{Error: "fenced"}) + if err := <-done; err != errGuestReconnect { + t.Fatalf("unready transport: %v", err) + } + if _, err := host.Read(ctx); err == nil { + t.Fatal("failed transport remained open") + } + if b.reconnecting { + t.Fatal("slot leaked") + } +} + +func TestReconnectDeadlineReleasesAttempt(t *testing.T) { + b, ch := reconnectFixture(t) + guest, host, _ := reconnectPair(t) + ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) + defer cancel() + done := make(chan error, 1) + go func() { done <- reBootstrap(ctx, guest, b) }() + writeReconnect(t, ctx, host, relay.KindGuestReconnectChallenge, ch) + if _, err := relay.ReadGuestReconnectFrame(ctx, host); err != nil { + t.Fatal(err) + } + // The peer never supplies acceptance. This same deadline covers both legs. + if err := <-done; err != errGuestReconnect { + t.Fatal("deadline did not fail closed") + } + if b.reconnecting || b.guest.epoch != 0 { + t.Fatal("timeout retained authority or attempt") + } +} + +func TestEnrollmentEnvironmentAboveHandshakeLimit(t *testing.T) { + guest, host, ctx := reconnectPair(t) + b := &bootstrapper{} + done := make(chan error, 1) + cfg := runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test", BootstrapToken: base64.RawURLEncoding.EncodeToString(make([]byte, 32)), SecretNames: []string{"LARGE_TEST"}} + want := strings.Repeat("synthetic", 1024) + go func() { + env, err := b.exchange(ctx, guest, cfg) + if err == nil && env["LARGE_TEST"] != want { + err = errGuestReconnect + } + done <- err + }() + if _, err := host.Read(ctx); err != nil { + t.Fatal(err) + } + payload, _ := json.Marshal(runner.GuestReconnectEnrollResponse{Env: map[string]string{"LARGE_TEST": want}}) + controlWrite(t, ctx, host, relay.ControlEvent{Kind: "resp", ID: secretsRequestID, OK: true, Payload: payload}) + if err := <-done; err != nil { + t.Fatal("authorized environment incorrectly bounded by handshake limit") + } +} + +func TestConfigurationDeadlineAfterApplication(t *testing.T) { + b, _ := reconnectFixture(t) + t.Setenv("DEADLINE_TEST", "") + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + // Deterministically expire while the post-application consumer runs, after + // the pre-application check. A preamble must not publish readiness afterward. + b.configurationApplied = func(context.Context) error { cancel(); return nil } + if err := b.refreshConfiguration(ctx, runner.BootConfig{Env: map[string]string{"DEADLINE_TEST": "new"}}, nil); err != errGuestReconnect { + t.Fatalf("expired delivery returned ready: %v", err) + } +} diff --git a/docs/design/2026-09-30-guest-reconnect-host.md b/docs/design/2026-09-30-guest-reconnect-host.md index 51387ab..38d510e 100644 --- a/docs/design/2026-09-30-guest-reconnect-host.md +++ b/docs/design/2026-09-30-guest-reconnect-host.md @@ -48,8 +48,13 @@ redial. A built execution probe calls the port through real runnerd/RunAgent wit a synthetic driver and proof provider. No shipping CLI or listener invokes this port yet, so this is executable port coverage, not VM or guest continuity evidence. -Depends on OSS PR113. Before enablement: guest key custody/enrollment, negotiated -capability, strict 4 KiB/five-second guest handshake, bounded peer acceptance, +Depends on OSS PR113. Guest key custody and the guest half of the bounded +handshake are supplied by the follow-on slice linked below. Before enablement: +negotiated capability, host handshake integration, bounded peer acceptance, current configuration delivery, cold lifecycle ordering, relay fencing and real PID/PTY/agent qualification. Hosted deployment additionally depends on Cloud 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. diff --git a/docs/design/2026-09-30-guest-reconnect-session.md b/docs/design/2026-09-30-guest-reconnect-session.md new file mode 100644 index 0000000..0406c7d --- /dev/null +++ b/docs/design/2026-09-30-guest-reconnect-session.md @@ -0,0 +1,106 @@ +# Guest signing identity and reconnect preamble + +This slice adds the guest side of authenticated live reconnection to `sessiond`. +It depends on the shared [RPC contracts](2026-09-30-guest-reconnect-rpc.md) and +[host authorization port](2026-09-30-guest-reconnect-host.md). It does not enable +host listeners or advertise a supported reconnect capability. + +## Fresh boot and retained identity + +`BootConfig.guest_reconnect` is an optional uint64. Absent/zero preserves the +existing guest bootstrap behavior. Protocol 1 selects enrollment on a fresh boot +only; unsupported values fail closed. A legacy redial cannot enroll later. +No production host sets this field yet. Before any host sets it, the negotiated +image/runner/control-plane capability, cold lifecycle and relay fencing must be +implemented together. + +A fresh opted-in sessiond generates an Ed25519 key and a random boot epoch in +process memory. The public key and boot epoch ride `enroll_guest_reconnect` +with the single-use bootstrap token. Enrollment is required even when no secrets +are declared. The response must contain a non-null environment object and every +declared secret (an empty string is valid). An ambiguous or refused enrollment +cannot be retried or downgraded; a fresh authorized boot is required. + +The private key is not serialized, logged, stored on disk or exported through +the environment. Losing sessiond loses its reconnect identity. This is continuity +of an enrolled guest, not protection against root within that same sandbox. + +## Stream protocol and readiness + +Later connections use existing newline-delimited `FrameControl` envelopes, with +attachment ID zero and exactly `kind` and `payload` in their control event: + +| Kind | Payload | +| --- | --- | +| `guest_reconnect_challenge` | Shared `GuestReconnectChallenge` | +| `guest_reconnect_proof` | Shared `GuestReconnectAcceptRequest` | +| `guest_reconnect_accepted` | Shared `GuestReconnectAcceptResponse` | +| `guest_reconnect_refused` | Shared fixed-code refusal | + +The encoded outer frame is at most 4096 bytes, excluding its newline. `NetConn` +retains its fixed 64 KiB read buffer but assembles no more than the selected +handshake bound; oversized unterminated frames close the connection. Unknown, +duplicate, missing, null and extra envelope fields fail closed. Shared strict +payload decoders validate challenge and acceptance. Any refusal ends the attempt +without echoing its contents. One guest attempt runs at a time and challenge +through acceptance has one five-second deadline. + +The guest checks its retained session and boot before signing the shared binary +transcript. An accepted epoch must exceed the last observed epoch. The epoch is +remembered before configuration delivery: a lost delivery requires a new +challenge and epoch. Nothing caches a signed proof or replays an accepted token. +The host/control plane remain responsible for durable authority, current placement, +policy, single-use attempts, ownership checks and fencing the old relay. + +After acceptance the guest requires a boot configuration for the same session, +protocol and accepted token. It redeems that fresh token through +`fetch_session_secrets`, including when zero secrets are declared. Configuration +and exchange waits retain their existing 30-second bounds and are additionally +bounded by the accepted token TTL. Failed proof, configuration, redemption or +validation returns a fixed error; `sessionTransport.connect` closes the stream +and does not make it available to the relay or credential dispatcher. + +## Configuration and process continuity + +The new path validates all environment names/values before applying any changes. +It preserves existing precedence (secrets, fixed launch settings, explicit env), +updates non-secret settings even for an empty secret response, and removes keys +previously owned by this configuration when they disappear. Unowned inherited +variables remain untouched. The legacy boot path keeps its existing behavior. + +This updates sessiond and atomically replaces the exec runner environment before +readiness, preserving boot-chain exports. New `rainier exec` children receive the +current settings; removed secrets are absent. A proof or authorization refusal +does not update that snapshot. Expiry during application can leave settings +applied, but the preamble returns failure and does not publish readiness. +The delivery context is checked during validation and after all +configuration consumers complete, so expiry during application refuses readiness. +An already-running child retains its +inherited environment; this is not live credential rotation inside a coding +agent. Reconnect does not execute setup/init, start an agent or replace a PTY. +Repository diff registration, generated git configuration and agent-sync entries +remain boot-time artifacts. This slice does not rebuild them. Future host +configuration delivery must preserve those launch invariants or explicitly +reconcile them before enabling changes to those fields. + +## Verification and remaining gates + +Tests cover secret-free enrollment, sticky refusal, late opt-in rejection, +strict stream envelopes, oversize reads before newline, cross-session/boot proof +refusal, replay epochs, failed/mismatched configuration, missing secrets, +concurrent attempts, connection closure and configuration replacement. + +`TestGuestReconnectExecutable` runs a separately executed guest test artifact +against a real TCP protocol fixture, using the production bootstrap/preamble and +PTY implementation. It verifies enrollment, a refused connection, a signed new +connection, fresh token redemption and the same shell PID on the same PTY before +and after, including the child-versus-parent environment distinction, refreshed +real exec children, removed values and preserved boot-chain exports. This is the +closest executable check for an unenabled guest path; the shipping host has no +caller yet. It is not AF_VSOCK/KVM or a hosted authorization/relay takeover test. + +Still required before enabling support: fresh configuration delivery from the +control plane, original-instance listener ownership, accepted-epoch relay takeover, +cold-resume enrollment/generation ordering, runner recovery integration and real +live-agent PID/PTY/tool-turn qualification. No snapshot or memory persistence is +introduced. User-facing enablement documentation is deferred to that integration. diff --git a/internal/relay/netconn.go b/internal/relay/netconn.go index 1168ef2..4b31850 100644 --- a/internal/relay/netconn.go +++ b/internal/relay/netconn.go @@ -148,3 +148,32 @@ func (n *netConn) Write(ctx context.Context, b []byte) error { } func (n *netConn) Close() error { return n.c.Close() } + +// ReadLimited reads one stream frame with a caller-selected bound before +// decoding it. It preserves buffered bytes for the following normal Read. +// Only one reader may use a Conn at a time. Oversize closes the connection. +func (n *netConn) ReadLimited(ctx context.Context, limit int) ([]byte, error) { + if limit < 1 || limit > maxFrameBytes { + return nil, ErrFrameTooLarge + } + if ctx.Err() != nil { + return nil, ctx.Err() + } + stop := context.AfterFunc(ctx, func() { _ = n.c.Close() }) + defer stop() + line := make([]byte, 0, min(limit, 4096)) + for { + b, err := n.r.ReadByte() + if err != nil { + return nil, err + } + if b == '\n' { + return bytes.TrimRight(line, "\r"), nil + } + if len(line) == limit { + _ = n.c.Close() + return nil, ErrFrameTooLarge + } + line = append(line, b) + } +} diff --git a/internal/relay/reconnect.go b/internal/relay/reconnect.go new file mode 100644 index 0000000..89974f2 --- /dev/null +++ b/internal/relay/reconnect.go @@ -0,0 +1,109 @@ +package relay + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "io" +) + +// Reconnect control kinds run before the ordinary relay is online. They use +// the existing FrameControl envelope with only kind and payload fields. +const ( + KindGuestReconnectChallenge = "guest_reconnect_challenge" + KindGuestReconnectProof = "guest_reconnect_proof" + KindGuestReconnectAccepted = "guest_reconnect_accepted" + KindGuestReconnectRefused = "guest_reconnect_refused" + GuestReconnectFrameLimit = 4096 +) + +var errReconnectFrame = errors.New("invalid guest reconnect frame") + +// ReadGuestReconnectFrame requires a transport with pre-allocation read bounds +// (NetConn today). It rejects unknown, duplicate, null and extra envelope fields, +// non-control/attachment frames and unknown kinds. Payloads need their shared +// protocol decoder as well. The caller owns its handshake context and cleanup. +func ReadGuestReconnectFrame(ctx context.Context, c Conn) (ControlEvent, error) { + var zero ControlEvent + reader, ok := c.(interface { + ReadLimited(context.Context, int) ([]byte, error) + }) + if !ok { + return zero, errReconnectFrame + } + raw, err := reader.ReadLimited(ctx, GuestReconnectFrameLimit) + if err != nil { + return zero, errReconnectFrame + } + var frame Frame + if strictReconnectObject(raw, &frame, "t", "a", "p") != nil || frame.Type != FrameControl || frame.AttachID != 0 { + return zero, errReconnectFrame + } + var event ControlEvent + if strictReconnectObject(frame.Payload, &event, "kind", "payload") != nil || !reconnectKind(event.Kind) { + return zero, errReconnectFrame + } + return event, nil +} + +// WriteGuestReconnectFrame writes one bounded handshake control frame. The +// caller must supply a deadline; no configuration or relay traffic belongs here. +func WriteGuestReconnectFrame(ctx context.Context, c Conn, kind string, payload []byte) error { + if !reconnectKind(kind) { + return errReconnectFrame + } + body, err := json.Marshal(ControlEvent{Kind: kind, Payload: payload}) + if err != nil { + return errReconnectFrame + } + raw, err := Encode(Frame{Type: FrameControl, Payload: body}) + if err != nil || len(raw) > GuestReconnectFrameLimit { + return errReconnectFrame + } + if err := c.Write(ctx, raw); err != nil { + return errReconnectFrame + } + return nil +} +func reconnectKind(kind string) bool { + switch kind { + case KindGuestReconnectChallenge, KindGuestReconnectProof, KindGuestReconnectAccepted, KindGuestReconnectRefused: + return true + } + return false +} +func strictReconnectObject(raw []byte, out any, keys ...string) error { + d := json.NewDecoder(bytes.NewReader(raw)) + token, err := d.Token() + if err != nil || token != json.Delim('{') { + return errReconnectFrame + } + allowed := make(map[string]bool, len(keys)) + for _, k := range keys { + allowed[k] = true + } + seen := make(map[string]bool, len(keys)) + for d.More() { + key, err := d.Token() + name, ok := key.(string) + if err != nil || !ok || !allowed[name] || seen[name] { + return errReconnectFrame + } + var value json.RawMessage + if d.Decode(&value) != nil || bytes.Equal(bytes.TrimSpace(value), []byte("null")) { + return errReconnectFrame + } + seen[name] = true + } + if token, err = d.Token(); err != nil || token != json.Delim('}') || len(seen) != len(keys) { + return errReconnectFrame + } + if _, err = d.Token(); err != io.EOF { + return errReconnectFrame + } + if json.Unmarshal(raw, out) != nil { + return errReconnectFrame + } + return nil +} diff --git a/internal/relay/reconnect_test.go b/internal/relay/reconnect_test.go new file mode 100644 index 0000000..c90af9b --- /dev/null +++ b/internal/relay/reconnect_test.go @@ -0,0 +1,84 @@ +package relay + +import ( + "context" + "encoding/base64" + "errors" + "net" + "strconv" + "strings" + "testing" + "time" +) + +func TestGuestHandshakeReadLimitBeforeNewline(t *testing.T) { + a, b := net.Pipe() + defer a.Close() + defer b.Close() + reader, ok := NetConn(a).(interface { + ReadLimited(context.Context, int) ([]byte, error) + }) + if !ok { + t.Fatal("stream cannot bound an unauthenticated handshake read") + } + go func() { _, _ = b.Write([]byte(strings.Repeat("x", 4097))) }() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + _, err := reader.ReadLimited(ctx, 4096) + if !errors.Is(err, ErrFrameTooLarge) { + t.Fatalf("oversize refusal=%v", err) + } +} + +func TestGuestReconnectFrames(t *testing.T) { + for _, kind := range []string{KindGuestReconnectChallenge, KindGuestReconnectProof, KindGuestReconnectAccepted, KindGuestReconnectRefused} { + t.Run(kind, func(t *testing.T) { + a, b := net.Pipe() + defer a.Close() + defer b.Close() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + writer, reader := NetConn(a), NetConn(b) + done := make(chan error, 1) + go func() { + if err := WriteGuestReconnectFrame(ctx, writer, kind, []byte(`{"protocol":1}`)); err != nil { + done <- err + return + } + done <- writer.Write(ctx, []byte("normal-frame")) + }() + ev, err := ReadGuestReconnectFrame(ctx, reader) + if err != nil || ev.Kind != kind { + t.Fatalf("roundtrip: %v %v", ev, err) + } + raw, err := reader.Read(ctx) + if err != nil || string(raw) != "normal-frame" { + t.Fatal("lost post-handshake bytes") + } + if err = <-done; err != nil { + t.Fatal(err) + } + }) + } +} +func TestGuestReconnectStrictEnvelope(t *testing.T) { + goodEvent := `{"kind":"guest_reconnect_proof","payload":{"protocol":1}}` + frame := func(event string) string { + return `{"t":4,"a":0,"p":"` + base64.StdEncoding.EncodeToString([]byte(event)) + `"}` + } + good := frame(goodEvent) + cases := []string{`null`, good + ` {}`, strings.Replace(good, `"t":4`, `"t":4,"t":4`, 1), strings.Replace(good, `"a":0`, `"a":1`, 1), strings.Replace(good, `"t":4`, `"t":3`, 1), strings.Replace(good, `"t":4`, `"t":null`, 1), strings.Replace(good, `"t":4`, `"t":4,"unknown":0`, 1), frame(`{"kind":"guest_reconnect_proof","payload":null}`), frame(`{"kind":"guest_reconnect_proof","kind":"guest_reconnect_proof","payload":{}}`), frame(`{"kind":"unknown","payload":{}}`), frame(`{"Kind":"guest_reconnect_proof","payload":{}}`), frame(`{"kind":"guest_reconnect_proof","payload":{},"id":1}`)} + for i, raw := range cases { + t.Run(strconv.Itoa(i), func(t *testing.T) { + a, b := net.Pipe() + defer a.Close() + defer b.Close() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + go NetConn(a).Write(ctx, []byte(raw)) + if _, err := ReadGuestReconnectFrame(ctx, NetConn(b)); err == nil { + t.Fatal("accepted invalid envelope") + } + }) + } +} diff --git a/internal/sandboxexec/exec.go b/internal/sandboxexec/exec.go index fdaba1f..b01a8f7 100644 --- a/internal/sandboxexec/exec.go +++ b/internal/sandboxexec/exec.go @@ -243,6 +243,16 @@ func NewRunner(root string, env []string, start Starter) *Runner { live: map[*attachment]struct{}{}} } +// ReplaceEnvironment installs the authorized launch environment for future execs. +// It does not mutate a request already composed or a process already running. +// Callers retain boot-chain exports and fence old relay admission separately. +func (r *Runner) ReplaceEnvironment(env []string) { + snapshot := append([]string(nil), env...) + r.mu.Lock() + r.env = snapshot + r.mu.Unlock() +} + // OpenExec starts one exec and returns its live attachment. It never returns // an error: every refusal is an `exec_error` the caller has to see and map to // an exit code, so a refused exec is an attachment that emits its reason and @@ -746,7 +756,9 @@ func (r *Runner) validate(spec runner.ExecSpec) (Request, string) { // the agent, then TERM under --tty (matching session.StartProc), then the // caller's own additions, each checked by name. func (r *Runner) composeEnv(spec runner.ExecSpec) ([]string, string) { + r.mu.Lock() env := append([]string(nil), r.env...) + r.mu.Unlock() if spec.TTY { env = append(env, "TERM=xterm-256color") } diff --git a/internal/sandboxexec/reconnect_env_test.go b/internal/sandboxexec/reconnect_env_test.go new file mode 100644 index 0000000..01ded1e --- /dev/null +++ b/internal/sandboxexec/reconnect_env_test.go @@ -0,0 +1,62 @@ +package sandboxexec + +import ( + "github.com/tokencanopy/rainier/protocol/runner" + "slices" + "sync" + "testing" +) + +func TestReconnectReplacesFutureExecEnvironment(t *testing.T) { + r := NewRunner(t.TempDir(), []string{"OLD_TEST_SECRET=old", "KEEP_TEST=old"}, nil) + replace, ok := any(r).(interface{ ReplaceEnvironment([]string) }) + if !ok { + t.Fatal("exec runner cannot replace its stale boot environment") + } + before, reason := r.composeEnv(runner.ExecSpec{}) + if reason != "" { + t.Fatal(reason) + } + current := []string{"KEEP_TEST=current", "EMPTY_TEST="} + replace.ReplaceEnvironment(current) + current[0] = "KEEP_TEST=mutated" + after, reason := r.composeEnv(runner.ExecSpec{}) + if reason != "" { + t.Fatal(reason) + } + if !slices.Contains(before, "OLD_TEST_SECRET=old") || !slices.Contains(before, "KEEP_TEST=old") { + t.Fatal("already composed environment changed") + } + if slices.Contains(after, "OLD_TEST_SECRET=old") || !slices.Contains(after, "KEEP_TEST=current") || !slices.Contains(after, "EMPTY_TEST=") { + t.Fatalf("future exec environment is stale or aliased: %v", after) + } + after[0] = "mutated" + again, _ := r.composeEnv(runner.ExecSpec{}) + if !slices.Contains(again, "KEEP_TEST=current") { + t.Fatal("caller mutated stored environment") + } +} + +func TestReconnectEnvironmentReplacementIsAtomic(t *testing.T) { + r := NewRunner(t.TempDir(), []string{"A=old", "B=old"}, nil) + var wg sync.WaitGroup + wg.Add(2) + go func() { + defer wg.Done() + for i := 0; i < 1000; i++ { + r.ReplaceEnvironment([]string{"A=new", "B=new"}) + r.ReplaceEnvironment([]string{"A=old", "B=old"}) + } + }() + go func() { + defer wg.Done() + for i := 0; i < 1000; i++ { + env, _ := r.composeEnv(runner.ExecSpec{}) + if !(slices.Equal(env, []string{"A=old", "B=old"}) || slices.Equal(env, []string{"A=new", "B=new"})) { + t.Error("mixed environment snapshot") + return + } + } + }() + wg.Wait() +} diff --git a/protocol/runner/messages.go b/protocol/runner/messages.go index 888def2..199db05 100644 --- a/protocol/runner/messages.go +++ b/protocol/runner/messages.go @@ -412,6 +412,11 @@ type Spec struct { // writes no part of this to disk, and the values SecretNames names arrive in // the guest over the token exchange, from the control plane, never from here. type BootConfig struct { + // GuestReconnect opts a fresh guest into protocol-1 key enrollment. Zero + // preserves legacy boot. Hosts must leave this absent until end-to-end + // capability negotiation, lifecycle and relay fencing are implemented. + GuestReconnect uint64 `json:"guest_reconnect,omitempty"` + // Protocol is SessionBootstrapProtocolVersion. A guest that reads a // version it does not speak fails its boot chain rather than booting on a // configuration it has half understood.