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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 55 additions & 0 deletions docs/design/2026-09-30-guest-reconnect-host.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
# Runner reconnect authorization callback

This dependent slice follows the shared [RPC contract](2026-09-30-guest-reconnect-rpc.md).
It implements the optional `driver.GuestReconnectHost` port on runnerd. It does
not advertise a capability or enable the guest listener. The next driver handshake
can call `host.AuthorizeGuestReconnect(ctx, sessionID, prove)`; `prove` exchanges
only the validated challenge for a signature on its bounded guest connection.

One call captures the accepted agent connection, local placement generation and
registry boot identity. Both RPCs use that connection's existing outbound queue
and runner-origin pending table. Redial never moves an old proof to the new queue.
Before the guest callback, before accept, and before returning authority, check the
original connection and live local placement again. The control plane remains the
source of durable placement, expiry, membership, policy and single-use decisions.
The driver supplies session identity; guest data cannot choose workload scope.
Unknown/zero placement and unaccepted connections refuse. Current cold-resume
placement bookkeeping is insufficient for this guard and must be reconciled before
enabling recovery; no permissive fallback is introduced.

The whole call has a five-second context budget, at most one in-flight operation
per session and 64 across a runner. Admission fails immediately when full. Slots
and pending RPC entries are released on return. Guest proof callbacks must honor
the context and bound peer/frame resources; the callback runs inline, so a broken
implementation that ignores cancellation retains its slot rather than spawning
unbounded detached work. It must not send configuration or replace a relay.

Strict shared decoders validate challenge, proof and acceptance. The challenge
must match the expected session, placement and original connection incarnation.
Refusals return a zero acceptance and only invalid/expired/fenced/unavailable;
upstream and guest error text is discarded. A failure after committed acceptance
can strand that attempt: start a fresh challenge, never replay or cache authority.
Successful return still requires fresh configuration and old-relay fencing by the
caller. No token, proof or challenge is logged.

The generic guest relay now explicitly refuses begin/accept by method name, for
both Docker and microVM guests, in addition to its existing high-bit request and
response fence. Enrollment remains guest-originated. Legacy bootstrap and other
session RPCs preserve their behavior.

A pair of separately callable begin/accept methods was rejected because callers
could accidentally migrate an attempt across agent reconnection, omit validation
or omit cleanup. A new public RPC/HTTP route was rejected: the existing agent
session-RPC queue and driver host callback already own this boundary.

Verification covers real WebSocket round trips, wrong scope, malformed replies,
fixed refusals, cancellation, deadline, admission, local replacement and control
redial. A built execution probe calls the port through real runnerd/RunAgent with
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,
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.
3 changes: 3 additions & 0 deletions docs/design/2026-09-30-guest-reconnect-rpc.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,3 +64,6 @@ initial and cold-resume enrollment ordering, bounded guest peer handling,
configuration refresh, and relay takeover. It must retain the five-second
challenge/handshake budget and prove original-socket fencing, replay refusal,
tenant isolation and real PID/PTY/agent continuity before advertising support.

The follow-on [runner host callback](2026-09-30-guest-reconnect-host.md) now supplies
connection-bound begin/proof/accept orchestration without enabling the listener.
23 changes: 23 additions & 0 deletions internal/driver/reconnect.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
package driver

import (
"context"
"github.com/tokencanopy/rainier/protocol/runner"
)

// GuestReconnectProof exchanges a validated challenge with the guest and returns
// its canonical signature. It must honor ctx, bound unauthenticated frames and
// never send configuration or replace the active relay. Errors are not exposed.
type GuestReconnectProof func(context.Context, runner.GuestReconnectChallenge) (string, error)

// GuestReconnectHost is an optional authorization port for a future reconnect
// listener. Implementing it does not advertise or enable reconnect capability.
type GuestReconnectHost interface {
// AuthorizeGuestReconnect performs challenge and acceptance on one accepted
// control connection, within a single five-second budget. sessionID is the
// driver's assignment, never guest input. Proof is called only after challenge
// scope is checked. Failures return a zero response and one fixed error code:
// invalid, expired, fenced or unavailable. Success is committed authority, not
// permission to skip fresh configuration or old-relay fencing.
AuthorizeGuestReconnect(context.Context, string, GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error)
}
8 changes: 6 additions & 2 deletions internal/runnerd/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -258,6 +258,10 @@ func (s *Server) agentSession(ctx context.Context, cfg AgentConfig) (established

ag := &agentSessionState{}
out := make(chan runner.FromRunner, 64)
rc := &reconnectControl{ctx: connCtx, state: ag, out: out}
s.reconnectControl.Store(rc)
defer s.reconnectControl.CompareAndSwap(rc, nil)

send := func(m runner.FromRunner) {
cctx, ccancel := context.WithTimeout(ctx, 5*time.Second)
m.Used, m.Total, _ = s.drv.Capacity(cctx) // best-effort; piggybacked on every message
Expand Down Expand Up @@ -607,8 +611,8 @@ func (s *Server) forwardSessionRPC(m runner.ToRunner, send func(runner.FromRunne
return
}
// A response to a request this RUNNER originated stops here: the runner
// is a pure forwarder for everything a sandbox asked, and the one thing
// it asks for itself (a bootstrap token on a cold resume) is answered to
// is a pure forwarder for everything a sandbox asked; cold bootstrap and
// reconnect authorization requested by the runner are answered to
// a caller inside this process, not to the guest. The id spaces are
// disjoint so the two can share one connection — see
// runnerOriginatedIDBase.
Expand Down
12 changes: 8 additions & 4 deletions internal/runnerd/microvm.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,9 @@ func isRunnerOriginated(id uint64) bool { return id&runnerOriginatedIDBase != 0
// rather than answering it there. `method` is the envelope's method, which
// is "resp" for a response, so both arms consult the same fence.
//
// Two things are refused, and only this hop can refuse either of them.
// Runner-only methods and runner-owned IDs are refused at this hop.
// Reconnect begin/accept are also runner-only: a guest may enroll its fresh
// boot but may not originate the host's later authorization flow.
//
// A sandbox may not mint its own bootstrap token. The design's whole point
// is that the token is SINGLE-USE and lives 120 seconds: a guest that could
Expand All @@ -103,6 +105,8 @@ func isRunnerOriginated(id uint64) bool { return id&runnerOriginatedIDBase != 0
// waiting on that number.
func refuseSandboxOrigin(method string, id uint64) string {
switch {
case method == runner.MethodBeginGuestReconnect || method == runner.MethodAcceptGuestReconnect:
return "fenced"
case method == runner.MethodMintSessionBootstrap:
return "a sandbox may not mint its own bootstrap token; a fresh one is minted by the runner on a cold resume"
case isRunnerOriginated(id):
Expand All @@ -114,8 +118,8 @@ func refuseSandboxOrigin(method string, id uint64) string {
// runnerRPCTable is the pending table for requests this runner originated.
//
// It is deliberately tiny and deliberately separate from the sandbox's: the
// runner asks the control plane exactly one thing (a bootstrap token, on a
// cold resume), and a table shared with the forwarding path would make the
// runner asks for cold bootstrap and guest reconnect authorization, and a
// table shared with the forwarding path would make the
// forwarder a participant in the conversation it is supposed to be a pipe
// for.
type runnerRPCTable struct {
Expand Down Expand Up @@ -191,7 +195,7 @@ const mintBootstrapTimeout = 10 * time.Second
// MintSessionBootstrap asks the control plane for a fresh single-use token
// for sessionID, on the runner's own control connection.
//
// It is the one request this runner originates on the session RPC. The
// This cold-resume request originates inside the runner. The
// message carries no session id of its own — FromRunner.Session is the id,
// and the control plane answers from the row its placement guard read — which
// is what stops a runner holding session A from minting for session B.
Expand Down
10 changes: 10 additions & 0 deletions internal/runnerd/microvm_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -264,6 +264,16 @@ func TestASandboxMayNotMintItsOwnBootstrapToken(t *testing.T) {
ev relay.ControlEvent
wantReason string
}{
{
name: "guest begin reconnect",
ev: relay.ControlEvent{Kind: "req:" + runner.MethodBeginGuestReconnect, ID: 5, Payload: []byte(`{"protocol":1}`)},
wantReason: "fenced",
},
{
name: "guest accept reconnect",
ev: relay.ControlEvent{Kind: "req:" + runner.MethodAcceptGuestReconnect, ID: 6, Payload: []byte(`{"protocol":1}`)},
wantReason: "fenced",
},
{
name: "the mint method itself",
ev: relay.ControlEvent{Kind: "req:" + runner.MethodMintSessionBootstrap, ID: 7, Payload: []byte(`{"protocol":1}`)},
Expand Down
166 changes: 166 additions & 0 deletions internal/runnerd/reconnect.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,166 @@
package runnerd

import (
"context"
"encoding/json"
"errors"
"strconv"
"time"

"github.com/tokencanopy/rainier/internal/driver"
"github.com/tokencanopy/rainier/protocol/runner"
)

const reconnectBudget = 5 * time.Second
const maxGuestReconnects = 64

var (
errReconnectInvalid = errors.New("invalid")
errReconnectExpired = errors.New("expired")
errReconnectFenced = errors.New("fenced")
errReconnectUnavailable = errors.New("unavailable")
)

// A snapshot of one agent connection. out belongs only to this connection;
// a late request can never land on the replacement's queue.
type reconnectControl struct {
ctx context.Context
state *agentSessionState
out chan<- runner.FromRunner
}

var _ driver.GuestReconnectHost = (*Server)(nil)

// AuthorizeGuestReconnect implements driver.GuestReconnectHost. It does not
// deliver configuration, attach a guest, or advertise reconnect capability.
func (s *Server) AuthorizeGuestReconnect(ctx context.Context, id string, prove driver.GuestReconnectProof) (runner.GuestReconnectAcceptResponse, error) {
var zero runner.GuestReconnectAcceptResponse
if prove == nil {
return zero, errReconnectInvalid
}
row, ok := s.reg.snapshot(id)
if !ok || row.placementGen == 0 || (row.state != "running" && row.state != "starting") {
return zero, errReconnectFenced
}
rc := s.reconnectControl.Load()
if rc == nil || rc.ctx.Err() != nil || rc.state.generation.Load() == 0 {
return zero, errReconnectUnavailable
}
generation := rc.state.generation.Load()
if !s.claimReconnect(id) {
return zero, errReconnectUnavailable
}
defer s.releaseReconnect(id)
ctx, cancel := context.WithTimeout(ctx, reconnectBudget)
defer cancel()
stop := context.AfterFunc(rc.ctx, cancel)
defer stop()
valid := func() bool {
current, exists := s.reg.snapshot(id)
return ctx.Err() == nil && rc.ctx.Err() == nil && s.reconnectControl.Load() == rc && rc.state.generation.Load() == generation && exists && current.placementGen == row.placementGen && current.boot == row.boot && current.handle == row.handle && (current.state == "running" || current.state == "starting")
}
if !valid() {
return zero, errReconnectFenced
}
payload, err := s.reconnectCall(ctx, rc, id, runner.MethodBeginGuestReconnect, runner.GuestReconnectBeginRequest{Protocol: runner.GuestReconnectProtocol})
if err != nil {
return zero, err
}
if !valid() {
return zero, errReconnectFenced
}
challenge, err := runner.DecodeGuestReconnectChallenge(payload)
if err != nil || challenge.SessionID != id || challenge.PlacementGeneration != row.placementGen || challenge.HostIncarnation != strconv.FormatUint(generation, 10) {
return zero, errReconnectInvalid
}
signature, err := prove(ctx, challenge)
if err != nil {
return zero, errReconnectUnavailable
}
if !valid() {
return zero, errReconnectFenced
}
request := runner.GuestReconnectAcceptRequest{Protocol: runner.GuestReconnectProtocol, AttemptID: challenge.AttemptID, Signature: signature}
body, _ := json.Marshal(request)
if _, err := runner.DecodeGuestReconnectAcceptRequest(body); err != nil {
return zero, errReconnectInvalid
}
payload, err = s.reconnectCall(ctx, rc, id, runner.MethodAcceptGuestReconnect, request)
if err != nil {
return zero, err
}
if !valid() {
return zero, errReconnectFenced
}
accepted, err := runner.DecodeGuestReconnectAcceptResponse(payload)
if err != nil {
return zero, errReconnectInvalid
}
return accepted, nil
}

func (s *Server) claimReconnect(id string) bool {
s.reconnectMu.Lock()
defer s.reconnectMu.Unlock()
if _, ok := s.reconnecting[id]; ok || len(s.reconnecting) >= maxGuestReconnects {
return false
}
if s.reconnecting == nil {
s.reconnecting = make(map[string]struct{})
}
s.reconnecting[id] = struct{}{}
return true
}
func (s *Server) releaseReconnect(id string) {
s.reconnectMu.Lock()
defer s.reconnectMu.Unlock()
delete(s.reconnecting, id)
}

func (s *Server) reconnectCall(ctx context.Context, rc *reconnectControl, session, method string, request any) ([]byte, error) {
if ctx.Err() != nil || rc.ctx.Err() != nil {
return nil, errReconnectUnavailable
}
body, err := json.Marshal(request)
if err != nil {
return nil, errReconnectInvalid
}
id, ch := s.runnerRPC.begin(session)
defer s.runnerRPC.end(id)
msg := runner.FromRunner{Type: "session_req", Session: session, RPC: &runner.RPCEnvelope{ID: id, Method: method, Payload: body}}
msg.Used, msg.Total, _ = s.drv.Capacity(ctx)
msg.Active, msg.IdleExited = s.reg.counts()
if ctx.Err() != nil || rc.ctx.Err() != nil {
return nil, errReconnectUnavailable
}
select {
case rc.out <- msg:
default:
return nil, errReconnectUnavailable
}
select {
case <-ctx.Done():
return nil, errReconnectUnavailable
case answer := <-ch:
if ctx.Err() != nil || rc.ctx.Err() != nil {
return nil, errReconnectUnavailable
}
if answer.OK {
return answer.Payload, nil
}
refusal, err := runner.DecodeGuestReconnectErrorResponse(answer.Payload)
if err != nil {
return nil, errReconnectUnavailable
}
switch refusal.Error {
case "invalid":
return nil, errReconnectInvalid
case "expired":
return nil, errReconnectExpired
case "fenced":
return nil, errReconnectFenced
default:
return nil, errReconnectUnavailable
}
}
}
Loading
Loading