diff --git a/go/e2e/exec_stop_test.go b/go/e2e/exec_stop_test.go index 0c4169e2b..0ee1f836b 100644 --- a/go/e2e/exec_stop_test.go +++ b/go/e2e/exec_stop_test.go @@ -39,6 +39,47 @@ func TestExecStreamingStopKillsInContainerProcess(t *testing.T) { }) } +// TestExecStreamingStopSweepsDetachedSession verifies that Stop kills work in +// a new session without killing the container's PID 1 keep-alive. The detached +// target holds 300 MiB so it is still dying when the sweep rescans. +func TestExecStreamingStopSweepsDetachedSession(t *testing.T) { + cli, id, name := startKeepAlive(t) + + script := `bun -e 'require("child_process").spawn("bun", ["-e", "globalThis.b = Buffer.alloc(300 << 20, 1); setInterval(() => {}, 1e6)", "hold-3002"], {detached: true, stdio: "ignore"}).unref()'; exec sleep 3000` + stream, err := cli.ExecStreaming(t.Context(), id, agentExecSpec("sh", "-c", script)) + if err != nil { + t.Fatalf("ExecStreaming: %v", err) + } + waitForTop(t, name, func(procs []string) bool { + return hasProcess(procs, "sleep 3000") && hasProcess(procs, "hold-3002") + }) + + if err := stream.Process.Terminate(); err != nil && !isClientKill(err) { + t.Fatalf("Terminate: unexpected error %v", err) + } + waitForTop(t, name, func(procs []string) bool { + return !hasProcess(procs, "sleep 3000") && hasProcess(procs, "hold-3002") + }) + + sweeper, ok := any(cli).(runtime.SessionSweeper) + if !ok { + t.Fatal("PodmanCLI does not implement runtime.SessionSweeper") + } + if err := sweeper.SweepExecSessions(t.Context(), id, strconv.FormatUint(uint64(agentuid.AgentUID), 10)); err != nil { + t.Fatalf("SweepExecSessions: %v", err) + } + waitForTop(t, name, func(procs []string) bool { + return !hasProcess(procs, "sleep 3000") && !hasProcess(procs, "hold-3002") && hasProcess(procs, "sleep infinity") + }) + running, err := cli.Running(t.Context(), name) + if err != nil { + t.Fatalf("Running: %v", err) + } + if !running { + t.Fatal("container stopped during exec-session sweep") + } +} + // TestExecStreamingNaturalExitKeepsStatus pins that an exec which exits on its // own, with stdin still open, returns promptly with its own exit code. func TestExecStreamingNaturalExitKeepsStatus(t *testing.T) { diff --git a/go/internal/runner/host.go b/go/internal/runner/host.go index ca4d33dc6..6fb95e744 100644 --- a/go/internal/runner/host.go +++ b/go/internal/runner/host.go @@ -12,6 +12,7 @@ import ( "log/slog" "maps" "path/filepath" + "strconv" "strings" "sync" "time" @@ -128,6 +129,7 @@ type liveSession struct { sessionID string containerName string containerID runtime.WorkloadID + agentUID uint32 stream *AgentStream state compassv1.AgentSessionState // agentAccountID is the owned agent account this session belongs to, copied @@ -481,6 +483,7 @@ func (h *agentHost) Start(ctx context.Context, req *compassv1.StartAgentSessionR sessionID: sessionID, containerName: name, containerID: handle.ID(), + agentUID: handle.WorkspaceUID(), stream: stream, state: compassv1.AgentSessionState_AGENT_SESSION_STATE_READY, agentAccountID: handle.AgentAccountID(), @@ -516,7 +519,7 @@ func (h *agentHost) Start(ctx context.Context, req *compassv1.StartAgentSessionR // Stop tears a session down. An unknown/already-stopped session succeeds // (idempotent, matching the established StopAgentSession semantics). -func (h *agentHost) Stop(_ context.Context, sessionID string) error { +func (h *agentHost) Stop(ctx context.Context, sessionID string) error { // Resolve session→container under h.mu first, release, then take the // container lock and re-check — the resolve-then-lock protocol (the lock key // is the container, but the caller names a session). A session that vanished @@ -529,6 +532,12 @@ func (h *agentHost) Stop(_ context.Context, sessionID string) error { } unlock := h.lockContainer(s.containerName) defer unlock() + h.mu.Lock() + current, ok := h.sessions[sessionID] + h.mu.Unlock() + if !ok || current != s { + return errSessionUnknown + } h.mu.Lock() s, ok = h.sessions[sessionID] @@ -546,10 +555,16 @@ func (h *agentHost) Stop(_ context.Context, sessionID string) error { if !ok { return nil } - if s.stream == nil { - return nil + if s.stream != nil { + if err := s.stream.Stop(); err != nil { + return err + } } - return s.stream.Stop() + if err := h.sweepExecSessions(ctx, s); err != nil { + h.log.Warn("sweeping detached agent exec sessions after stop", + slog.String("container", s.containerName), slog.String("session_id", sessionID), slog.Any("error", err)) + } + return nil } // Remove tears a container down and everything bound to it: it retires the live @@ -687,6 +702,9 @@ func (h *agentHost) RefreshSecrets(ctx context.Context, sessionID string) error if err != nil { return fmt.Errorf("fetching secrets for session %q: %w", sessionID, err) } + // The container lock keeps a Stop/Reload session sweep from killing these execs. + unlock := h.lockContainer(s.containerName) + defer unlock() if err := h.materializer.Install(ctx, handle.ID(), handle.HomeDir(), handle.WorkspaceUID(), resolved); err != nil { return fmt.Errorf("materializing secrets for session %q: %w", sessionID, err) } @@ -1096,6 +1114,10 @@ func (h *agentHost) reloadLocked(ctx context.Context, sessionID string) error { return err } } + if err := h.sweepExecSessions(ctx, s); err != nil { + h.markErrored(ctx, sessionID, s.containerName, nil) + return fmt.Errorf("sweeping detached agent exec sessions before reload for container %q: %w", s.containerName, err) + } // Hand the control state to the new process before it launches: replay_complete // becomes seq 1 and unacked ops follow, so no concurrent Deliver lands ahead of // the barrier. An ERRORED session was retired at exit, so Bind recreates it. @@ -1258,3 +1280,12 @@ func (h *agentHost) closeSocket(ctx context.Context, containerName string) { h.log.Warn("closing agent socket", slog.String("container", containerName), slog.Any("error", err)) } } + +// sweepExecSessions runs the optional container-backend sweep for this agent. +func (h *agentHost) sweepExecSessions(ctx context.Context, s *liveSession) error { + sweeper, ok := h.engine.(runtime.SessionSweeper) + if !ok { + return nil + } + return sweeper.SweepExecSessions(ctx, s.containerID, strconv.FormatUint(uint64(s.agentUID), 10)) +} diff --git a/go/internal/runner/host_test.go b/go/internal/runner/host_test.go index 42c79dc49..c695995ad 100644 --- a/go/internal/runner/host_test.go +++ b/go/internal/runner/host_test.go @@ -1253,6 +1253,105 @@ func TestStatusIsAnsweredFromLiveSet(t *testing.T) { } } +type recordedSweep struct { + id runtime.WorkloadID + user string +} + +type sessionSweeperRuntime struct { + *stubStreamingRuntime + sweeps []recordedSweep + sweepErr error +} + +func (r *sessionSweeperRuntime) SweepExecSessions(_ context.Context, id runtime.WorkloadID, user string) error { + r.sweeps = append(r.sweeps, recordedSweep{id: id, user: user}) + return r.sweepErr +} + +func TestStopAndReloadSweepExecSessions(t *testing.T) { + engine := &sessionSweeperRuntime{stubStreamingRuntime: newStubStreamingRuntime(t)} + registry := runtime.NewAgentRegistry() + rt := runtime.NewAgentRuntimeWithRegistry(engine, registry) + link := newLink(newRunnerServiceServer(t, newCapturePublish())) + host := NewSessionHost(link, rt, registry, engine, &fakeSpecBuilder{spec: liveSpec()}, AgentHostConfig{RuntimeDir: t.TempDir()}, discardLoggerRunner()) + ctx := t.Context() + t.Cleanup(func() { + if err := host.Stop(ctx, "stop-session"); err != nil { + t.Errorf("Stop after test = %v", err) + } + if err := host.Stop(ctx, "reload-session"); err != nil { + t.Errorf("Stop after test = %v", err) + } + }) + + if _, err := host.Provision(ctx, &compassv1.ProvisionAgentWorkspaceRequest{}, "0123456789abcdef0123456789abcdef"); err != nil { + t.Fatalf("Provision = %v", err) + } + sessionID, err := host.Start(ctx, &compassv1.StartAgentSessionRequest{ContainerName: "cont-1"}, "", "stop-session") + if err != nil { + t.Fatalf("Start = %v", err) + } + if err := host.Stop(ctx, sessionID); err != nil { + t.Fatalf("Stop = %v", err) + } + if len(engine.sweeps) != 1 || engine.sweeps[0] != (recordedSweep{id: "cont-1", user: "1000"}) { + t.Fatalf("sweeps after Stop = %+v, want one sweep for cont-1 as uid 1000", engine.sweeps) + } + + sessionID, err = host.Start(ctx, &compassv1.StartAgentSessionRequest{ContainerName: "cont-1"}, "", "reload-session") + if err != nil { + t.Fatalf("Start before Reload = %v", err) + } + if err := host.Reload(ctx, sessionID); err != nil { + t.Fatalf("Reload = %v", err) + } + want := []recordedSweep{{id: "cont-1", user: "1000"}, {id: "cont-1", user: "1000"}} + if !slices.Equal(engine.sweeps, want) { + t.Fatalf("sweeps after Reload = %+v, want %+v", engine.sweeps, want) + } + if err := host.Stop(ctx, sessionID); err != nil { + t.Fatalf("Stop after Reload = %v", err) + } +} + +// TestReloadSweepFailureMarksErrored pins that a failed sweep after the agent +// was stopped leaves the session ERRORED, never READY with no agent behind it. +func TestReloadSweepFailureMarksErrored(t *testing.T) { + engine := &sessionSweeperRuntime{stubStreamingRuntime: newStubStreamingRuntime(t)} + registry := runtime.NewAgentRegistry() + rt := runtime.NewAgentRuntimeWithRegistry(engine, registry) + link := newLink(newRunnerServiceServer(t, newCapturePublish())) + host := NewSessionHost(link, rt, registry, engine, &fakeSpecBuilder{spec: liveSpec()}, AgentHostConfig{RuntimeDir: t.TempDir()}, discardLoggerRunner()) + ctx := t.Context() + + if _, err := host.Provision(ctx, &compassv1.ProvisionAgentWorkspaceRequest{}, "0123456789abcdef0123456789abcdef"); err != nil { + t.Fatalf("Provision = %v", err) + } + sessionID, err := host.Start(ctx, &compassv1.StartAgentSessionRequest{ContainerName: "cont-1"}, "", "sweep-fail") + if err != nil { + t.Fatalf("Start = %v", err) + } + t.Cleanup(func() { + if err := host.Stop(ctx, sessionID); err != nil { + t.Errorf("Stop after test = %v", err) + } + }) + + sweepErr := errors.New("processes remain") + engine.sweepErr = sweepErr + if err := host.Reload(ctx, sessionID); !errors.Is(err, sweepErr) { + t.Fatalf("Reload = %v, want the sweep error", err) + } + st, err := host.Status(ctx, sessionID) + if err != nil || len(st) != 1 { + t.Fatalf("Status = %+v, %v", st, err) + } + if got := st[0].GetState(); got != compassv1.AgentSessionState_AGENT_SESSION_STATE_ERRORED { + t.Fatalf("state after failed sweep = %v, want ERRORED", got) + } +} + // Reload restarts a session's agent in place, reusing the SAME session id (the // board entry stays continuous). A bug that minted a new id would break board // continuity. diff --git a/go/internal/runtime/applecontainer.go b/go/internal/runtime/applecontainer.go index ae0fe686d..45615da3c 100644 --- a/go/internal/runtime/applecontainer.go +++ b/go/internal/runtime/applecontainer.go @@ -230,6 +230,19 @@ func (a *AppleContainerCLI) Exec(ctx context.Context, id WorkloadID, spec ExecSp }, nil } +// SweepExecSessions removes detached exec processes while preserving the +// container's PID 1 session and keep-alive. +func (a *AppleContainerCLI) SweepExecSessions(ctx context.Context, id WorkloadID, user string) error { + out, err := a.Exec(ctx, id, NewExecSpec("sh", "-s").AsUser(user).WithStdin(sessionSweepScript)) + if err != nil { + return err + } + if out.ExitCode != 0 { + return &CommandError{Summary: "container exec session sweep", ExitCode: out.ExitCode, Stderr: strings.TrimSpace(out.Stderr)} + } + return nil +} + // appleExecArgs assembles the argv for a one-shot `container exec`. Split out so // the argv assembly is unit-testable without spawning the CLI. func appleExecArgs(id WorkloadID, spec ExecSpec) []string { diff --git a/go/internal/runtime/clispawn.go b/go/internal/runtime/clispawn.go index 8ca0e0385..892e9fe3c 100644 --- a/go/internal/runtime/clispawn.go +++ b/go/internal/runtime/clispawn.go @@ -98,6 +98,42 @@ func (e cliEngine) run(ctx context.Context, summary string, args []string) ([]by // The shell notices for those kills go to /dev/null; "$@" keeps the real stderr. const stopWithClientScript = `exec 3>&2; { { sh -c 'echo "$$"; exec cat' && kill -s KILL 0; } | { read -r w; "$@" 2>&3 3>&-; s=$?; kill -s KILL "$w"; exit "$s"; }; } 2>/dev/null` //nolint:gosec // G101: a shell script, not a credential +// sessionSweepScript kills detached exec sessions while preserving PID 1's +// session, which carries the container keep-alive needed for in-place reloads. +// It waits between rounds: a killed process stays visible until the kernel reaps it. +const sessionSweepScript = `IFS=' ' +sid_of() { + stat=$(cat "/proc/$1/stat" 2>/dev/null) || return 1 + # Fields after the last ") " are kernel-written; comm before it may hold anything. + set -- ${stat##*) } + state=$1 + sid=$4 + [ -n "$sid" ] +} +sid_of 1 || { echo "cannot read PID 1 session" >&2; exit 1; } +keep=$sid +sid_of $$ || { echo "cannot read sweep session" >&2; exit 1; } +own=$sid +round=0 +while [ "$round" -lt 100 ]; do + round=$((round + 1)) + live=0 + for d in /proc/[0-9]*; do + pid=${d#/proc/} + sid_of "$pid" || continue + [ "$sid" = "$keep" ] || [ "$sid" = "$own" ] || [ "$state" = Z ] && continue + # A landed kill on a not-yet-reaped process counts until it is gone or a zombie; + # another uid's process (EPERM) is outside this sweep. + if kill -s KILL "$pid" 2>/dev/null; then + live=$((live + 1)) + fi + done + [ "$live" -eq 0 ] && exit 0 + sleep 0.1 +done +echo "processes remain outside PID 1 session after 10s" >&2 +exit 1` + // stopWithClient wraps command so killing the engine client also kills its // in-container process group, which the engine leaves running on its own. // A descendant that called setsid escapes the group kill. diff --git a/go/internal/runtime/podman.go b/go/internal/runtime/podman.go index d7bf52792..cb8e7a12d 100644 --- a/go/internal/runtime/podman.go +++ b/go/internal/runtime/podman.go @@ -410,6 +410,12 @@ type WorkloadRuntime interface { Resize(ctx context.Context, id WorkloadID, limits ResourceLimits) error } +// SessionSweeper is an optional backend capability for removing detached +// exec sessions without stopping the workload's PID 1 keep-alive. +type SessionSweeper interface { + SweepExecSessions(ctx context.Context, id WorkloadID, user string) error +} + // A backend that self-arms egress does NOT grow a verb on WorkloadRuntime: // MicroVMRuntime carries an off-interface marker EgressArmedInGuest(), // and AgentRuntime.provision type-asserts inGuestEgressArmer to skip armEgress @@ -609,6 +615,19 @@ func (p *PodmanCLI) Exec(ctx context.Context, id WorkloadID, spec ExecSpec) (Exe }, nil } +// SweepExecSessions removes detached exec processes while preserving the +// container's PID 1 session and keep-alive. +func (p *PodmanCLI) SweepExecSessions(ctx context.Context, id WorkloadID, user string) error { + out, err := p.Exec(ctx, id, NewExecSpec("sh", "-s").AsUser(user).WithStdin(sessionSweepScript)) + if err != nil { + return err + } + if out.ExitCode != 0 { + return &CommandError{Summary: "podman exec session sweep", ExitCode: out.ExitCode, Stderr: strings.TrimSpace(out.Stderr)} + } + return nil +} + // ExecStreaming starts a streaming `podman exec -i` through the shared // subprocess seam, returning the live pipes plus a kill/wait handle. func (p *PodmanCLI) ExecStreaming(ctx context.Context, id WorkloadID, spec StreamingExecSpec) (*StreamingExec, error) {