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
4 changes: 2 additions & 2 deletions cmd/rainier/lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ func eligibilityText(verdict, note string) string {
// yet" from "never".
func attachNote(s session) string {
switch s.State {
case "queued", "creating":
case "queued", "creating", "resuming":
return "waits for it to start"
case "suspended_warm", "suspended_cold":
return "resumes it first"
Expand All @@ -141,7 +141,7 @@ func stopNote(s session) string {
return "this build does not know this state; the server decides"
}
switch s.State {
case "queued", "creating":
case "queued", "creating", "resuming":
return "only a running session can be stopped"
case "suspended_warm", "suspended_cold":
return "already stopped"
Expand Down
8 changes: 4 additions & 4 deletions cmd/rainier/sessionstate.go
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ const (
// exited would be the collapse this file exists to prevent.
func lifecycleOf(s session) string {
switch s.State {
case "queued", "creating":
case "queued", "creating", "resuming":
return lifecycleStarting
case "running":
// Regardless of child_exit_code. The sandbox is up.
Expand Down Expand Up @@ -188,7 +188,7 @@ const (
// Everything else is refused by the endpoint.
func canAttach(s session) string {
switch s.State {
case "running", "queued", "creating", "suspended_warm", "suspended_cold":
case "running", "queued", "creating", "resuming", "suspended_warm", "suspended_cold":
return eligibleYes
case "failed":
if s.Reachable {
Expand Down Expand Up @@ -216,7 +216,7 @@ func canStop(s session) string {
switch s.State {
case "running":
return eligibleYes
case "queued", "creating", "suspended_warm", "suspended_cold",
case "queued", "creating", "resuming", "suspended_warm", "suspended_cold",
"failed", "dead", "canceled", "destroyed":
return eligibleNo
default:
Expand All @@ -231,7 +231,7 @@ func canStop(s session) string {
// a live container and deleting it is the cleanup.
func canDelete(s session) string {
switch s.State {
case "running", "queued", "suspended_warm", "suspended_cold",
case "running", "queued", "resuming", "suspended_warm", "suspended_cold",
"failed", "dead", "canceled":
return eligibleYes
case "creating":
Expand Down
2 changes: 2 additions & 0 deletions cmd/rainier/sessionstate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ func TestThreeDimensions(t *testing.T) {
"creating", session{State: "creating", Reachable: true},
lifecycleStarting, processNone, connectionAvailable,
},
{"resuming", session{State: "resuming", Reachable: true}, lifecycleStarting, processNone, connectionAvailable},
{
"running with a live child", session{State: "running", Reachable: true},
lifecycleRunning, processRunning, connectionAvailable,
Expand Down Expand Up @@ -159,6 +160,7 @@ func TestActionEligibilityFollowsRawStates(t *testing.T) {
// it with a conflict because a dispatch may be in flight.
"creating", session{State: "creating"}, eligibleYes, eligibleNo, eligibleNo,
},
{"resuming", session{State: "resuming"}, eligibleYes, eligibleNo, eligibleYes},
{"suspended warm", session{State: "suspended_warm"}, eligibleYes, eligibleNo, eligibleYes},
{"suspended cold", session{State: "suspended_cold"}, eligibleYes, eligibleNo, eligibleYes},
{
Expand Down
2 changes: 1 addition & 1 deletion cmd/rainier/wait.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ func initialAttachGuidance(ctx context.Context, cfg cli.Config, id string, since
} else {
row := resp.Session
switch row.State {
case "queued", "creating":
case "queued", "creating", "resuming":
explanation = "session is " + row.State + " and not ready to attach"
if row.QueueReason != "" {
// Queue text is server-owned prose. Redact known credentials before
Expand Down
19 changes: 11 additions & 8 deletions cmd/runnerd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ func main() {
"bearer token for the controld dial (required when --controld is set; or set RAINIER_RUNNER_TOKEN, which keeps it out of the process list)")
driverFlag := flag.String("driver", envDefault("RAINIER_RUNNER_DRIVER", "docker"),
"execution driver for session sandboxes: docker | microvm")
microvmGuestReconnect := flag.Bool("microvm-guest-reconnect", os.Getenv("RAINIER_MICROVM_GUEST_RECONNECT") == "1",
"opt in to authenticated recovery of surviving microVM guests after control connection or runner restart; requires a matching control plane and guest image")
kernelPath := flag.String("kernel", envDefault("RAINIER_KERNEL_PATH", ""),
"guest vmlinux kernel path for the microvm driver (required when --driver=microvm; the runner refuses to start without a readable one)")
rootfsPath := flag.String("rootfs", envDefault("RAINIER_ROOTFS_PATH", ""),
Expand Down Expand Up @@ -149,14 +151,15 @@ func main() {
log.Fatalf("--driver=microvm: %v", err)
}
mvm, err := driver.NewMicrovm(driver.MicrovmOpts{
Checkpoint: ckpt,
KernelPath: *kernelPath,
BaseRootfs: *rootfsPath,
StateDir: *microvmStateDir,
TotalSlots: *slots,
VCPU: *microvmVCPUs,
MemoryMiB: *microvmMemoryMiB,
ImageSource: imageSource,
GuestReconnect: *microvmGuestReconnect,
Checkpoint: ckpt,
KernelPath: *kernelPath,
BaseRootfs: *rootfsPath,
StateDir: *microvmStateDir,
TotalSlots: *slots,
VCPU: *microvmVCPUs,
MemoryMiB: *microvmMemoryMiB,
ImageSource: imageSource,

SlotGuestCIDR: *microvmGuestCIDR,
SlotUplinkCIDR: *microvmUplinkCIDR,
Expand Down
25 changes: 23 additions & 2 deletions cmd/sessiond/reconnect.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"context"
"crypto/ed25519"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/json"
"errors"
Expand All @@ -29,6 +30,7 @@ type guestReconnectIdentity struct {
key ed25519.PrivateKey
enrolled bool
epoch uint64
launch [32]byte
}

func (b *bootstrapper) enrollGuest(ctx context.Context, c relay.Conn, cfg runner.BootConfig) (map[string]string, error) {
Expand All @@ -38,7 +40,7 @@ func (b *bootstrapper) enrollGuest(ctx context.Context, c relay.Conn, cfg runner
return nil, errGuestReconnect
}
// A failed or unsupported opt-in is sticky: later dials cannot downgrade.
g := &guestReconnectIdentity{session: cfg.SessionID}
g := &guestReconnectIdentity{session: cfg.SessionID, launch: guestLaunchIdentity(cfg)}
b.guest = g
b.mu.Unlock()
if cfg.GuestReconnect != runner.GuestReconnectProtocol || !guestID(cfg.SessionID) {
Expand Down Expand Up @@ -106,7 +108,7 @@ func (b *bootstrapper) reconnectGuest(ctx context.Context, c relay.Conn) error {
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 {
if err != nil || cfg.SessionID != session || cfg.GuestReconnect != 1 || cfg.BootstrapToken != accepted.Token || guestLaunchIdentity(cfg) != g.launch {
return errGuestReconnect
}
b.mu.Lock()
Expand Down Expand Up @@ -335,3 +337,22 @@ func guestReady(ctx context.Context, c relay.Conn, epoch uint64) error {
}
return nil
}

// The live boot retains only a digest of fields whose consumers run once.
// Repositories, git attribution, setup, proxy policy and agent custody paths
// cannot be refreshed merely by changing sessiond's environment. A changed
// launch requires a new boot; ordinary environment values and secret names may
// still refresh. The digest never reaches the host or durable storage.
func guestLaunchIdentity(cfg runner.BootConfig) [32]byte {
cfg.BootstrapToken = ""
cfg.SecretNames = nil
launchEnv := map[string]string{}
for key, value := range cfg.Env {
if key == "RAINIER_AGENTS_B64" || strings.HasPrefix(value, "/rainier/agents/") {
launchEnv[key] = value
}
}
cfg.Env = launchEnv
body, _ := json.Marshal(cfg)
return sha256.Sum256(body)
}
40 changes: 40 additions & 0 deletions cmd/sessiond/reconnect_launch_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
package main

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

func TestReconnectRejectsChangedBootArtifactsBeforeRedemption(t *testing.T) {
for _, field := range []string{"command", "repository", "git", "setup", "proxy", "agent-manifest"} {
t.Run(field, func(t *testing.T) {
b, ch := reconnectFixture(t)
guest, host, ctx := reconnectPair(t)
done := make(chan error, 1)
go func() { done <- reBootstrap(ctx, guest, b); guest.Close() }()
accepted := acceptProof(t, ctx, host, b, ch, 1)
cfg := runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test", BootstrapToken: accepted.Token}
switch field {
case "command":
cfg.Cmd = []string{"different_test"}
case "repository":
cfg.Repos = []runner.RepoSpec{{Owner: "example", Name: "changed"}}
case "git":
cfg.GitAuthorName = "Changed Test"
case "setup":
cfg.Setup = "echo changed_test"
case "proxy":
cfg.ProxyURL = "http://proxy.example.invalid:80"
case "agent-manifest":
cfg.Env = map[string]string{"RAINIER_AGENTS_B64": "changed_test"}
}
sendReconnectConfig(t, ctx, host, cfg)
if raw, err := host.Read(ctx); err == nil || len(raw) != 0 {
t.Fatal("changed launch artifacts reached token redemption")
}
if err := <-done; err != errGuestReconnect {
t.Fatal("changed launch artifacts accepted")
}
})
}
}
2 changes: 1 addition & 1 deletion cmd/sessiond/reconnect_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,7 @@ func reconnectFixture(t *testing.T) (*bootstrapper, runner.GuestReconnectChallen
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))}
return &bootstrapper{guest: &guestReconnectIdentity{session: "session-test", boot: "boot-test", key: key, enrolled: true, launch: guestLaunchIdentity(runner.BootConfig{Protocol: 1, GuestReconnect: 1, SessionID: "session-test"})}}, 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()
Expand Down
7 changes: 4 additions & 3 deletions control/contract_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,7 @@ func TestSessionStateVocabulary(t *testing.T) {
values := map[control.SessionState]string{
control.StateQueued: "queued",
control.StateCreating: "creating",
control.StateResuming: "resuming",
control.StateRunning: "running",
control.StateSuspendedWarm: "suspended_warm",
control.StateSuspendedCold: "suspended_cold",
Expand All @@ -200,7 +201,7 @@ func TestSessionStateVocabulary(t *testing.T) {
}

terminalStates := map[control.SessionState]bool{
control.StateQueued: false, control.StateCreating: false, control.StateRunning: false,
control.StateQueued: false, control.StateCreating: false, control.StateResuming: false, control.StateRunning: false,
control.StateSuspendedWarm: false, control.StateSuspendedCold: false,
control.StateCanceled: true, control.StateFailed: true,
control.StateDead: true, control.StateDestroyed: true,
Expand All @@ -213,7 +214,7 @@ func TestSessionStateVocabulary(t *testing.T) {

slotStates := map[control.SessionState]bool{
control.StateQueued: false,
control.StateCreating: true, control.StateRunning: true,
control.StateCreating: true, control.StateResuming: true, control.StateRunning: true,
control.StateSuspendedWarm: true, control.StateSuspendedCold: false,
control.StateCanceled: false, control.StateFailed: false,
control.StateDead: false, control.StateDestroyed: false,
Expand All @@ -225,7 +226,7 @@ func TestSessionStateVocabulary(t *testing.T) {
}

wantOrder := []control.SessionState{
control.StateQueued, control.StateCreating, control.StateRunning,
control.StateQueued, control.StateCreating, control.StateResuming, control.StateRunning,
control.StateSuspendedWarm, control.StateSuspendedCold,
}
if len(control.NonTerminal) != len(wantOrder) {
Expand Down
64 changes: 48 additions & 16 deletions control/fleet.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,11 @@ type Runner struct {
PoolID PoolID
CapacityUsed int
CapacityTotal int
Connected bool
// CapacityPlacements identifies slots counted in the same CapacityUsed
// observation, by exact placement. Missing entries remain reservations.
CapacityPlacements map[SessionID]uint64

Connected bool
// Generation is monotonic per runner; a stale holder is rejected without
// replacing the event shapes (see RunnerEvent).
Generation uint64
Expand Down Expand Up @@ -69,14 +73,15 @@ type RunnerSession struct {
// and session bindings here rather than accepting a user-created Scope. Every
// field is adapter-derived; none is decoded directly from client JSON.
type RunnerRegistration struct {
WorkspaceID WorkspaceID
PoolID PoolID
RunnerID RunnerID
Generation uint64
CapacityUsed int
CapacityTotal int
Capabilities []string
Sessions []RunnerSession
WorkspaceID WorkspaceID
PoolID PoolID
RunnerID RunnerID
Generation uint64
CapacityUsed int
CapacityTotal int
CapacityPlacements map[SessionID]uint64
Capabilities []string
Sessions []RunnerSession
}

// RunnerRegistrationResult is the application's answer to a registration.
Expand All @@ -92,13 +97,14 @@ type RunnerRegistrationResult struct {
// reconciliation: its identity, generation, capacity, and the sessions it
// holds.
type RunnerSnapshot struct {
WorkspaceID WorkspaceID
PoolID PoolID
RunnerID RunnerID
Generation uint64
CapacityUsed int
CapacityTotal int
Sessions []RunnerSession
WorkspaceID WorkspaceID
PoolID PoolID
RunnerID RunnerID
Generation uint64
CapacityUsed int
CapacityTotal int
CapacityPlacements map[SessionID]uint64
Sessions []RunnerSession
}

// ReconcileResult is the application's answer to a snapshot.
Expand Down Expand Up @@ -199,3 +205,29 @@ type Event struct {
// of which an audit reader needs to know that a command was run.
Command string
}

// ValidateCapacityPlacements rejects an observation that cannot describe its
// aggregate. Keys come only from the authenticated runner's own inventory.
func ValidateCapacityPlacements(used int, placements map[SessionID]uint64) error {
if used < 0 || len(placements) > used {
return ErrInvalid
}
for id, generation := range placements {
if id == "" || generation == 0 || generation > uint64(1<<63-1) {
return ErrInvalid
}
}
return nil
}

// AvailableRunnerSlots subtracts pending create/resume reservations. Reservations already present in this exact aggregate sample occupy one slot,
// not two. An old placement or a legacy sample cannot cancel a reservation.
func AvailableRunnerSlots(host Runner, pending []Session) int {
reserved := 0
for _, row := range pending {
if row.PlacementGeneration == 0 || host.CapacityPlacements[row.ID] != row.PlacementGeneration {
reserved++
}
}
return host.CapacityTotal - host.CapacityUsed - reserved
}
25 changes: 16 additions & 9 deletions control/ports.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,15 +20,22 @@ import (
type Action string

const (
ActionCreate Action = "create"
ActionGet Action = "get"
ActionList Action = "list"
ActionUpdate Action = "update"
ActionDelete Action = "delete"
ActionSuspend Action = "suspend"
ActionResume Action = "resume"
ActionSnapshot Action = "snapshot"
ActionAttach Action = "attach"
// These are audit labels; reconnect is authorized as ActionResume.
ActionGuestEnroll Action = "guest_enroll"
ActionGuestReconnectBegin Action = "guest_reconnect_begin"
ActionGuestReconnectAccept Action = "guest_reconnect_accept"
ActionGuestReconnectConfigure Action = "guest_reconnect_configure"
ActionGuestBootstrapMint Action = "guest_bootstrap_mint"
ActionGuestBootstrapRedeem Action = "guest_bootstrap_redeem"
ActionCreate Action = "create"
ActionGet Action = "get"
ActionList Action = "list"
ActionUpdate Action = "update"
ActionDelete Action = "delete"
ActionSuspend Action = "suspend"
ActionResume Action = "resume"
ActionSnapshot Action = "snapshot"
ActionAttach Action = "attach"
// ActionExec is the audit label for one command run inside a session's
// sandbox. It is deliberately an EVENT label and not an authorization
// verb: exec is authorized as ActionAttach, because it grants nothing a
Expand Down
Loading
Loading