Skip to content
Merged
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
41 changes: 41 additions & 0 deletions go/e2e/exec_stop_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
39 changes: 35 additions & 4 deletions go/internal/runner/host.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (
"log/slog"
"maps"
"path/filepath"
"strconv"
"strings"
"sync"
"time"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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
Expand All @@ -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]
Expand All @@ -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
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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))
}
99 changes: 99 additions & 0 deletions go/internal/runner/host_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
13 changes: 13 additions & 0 deletions go/internal/runtime/applecontainer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
36 changes: 36 additions & 0 deletions go/internal/runtime/clispawn.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
19 changes: 19 additions & 0 deletions go/internal/runtime/podman.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down
Loading