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
140 changes: 140 additions & 0 deletions go/e2e/exec_stop_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
//go:build podman

package e2e

import (
"errors"
"os/exec"
"slices"
"strconv"
"strings"
"testing"
"time"

"github.com/RigelBuild/compass/go/internal/agentuid"
"github.com/RigelBuild/compass/go/internal/runtime"
)

// TestExecStreamingStopKillsInContainerProcess pins that terminating a podman
// streaming exec ends the processes inside the container, not only the host-side
// `podman exec -i` client. The exec forks a background child, so a fix that
// kills only the direct exec process still leaves a survivor.
func TestExecStreamingStopKillsInContainerProcess(t *testing.T) {
cli, id, name := startKeepAlive(t)

script := "trap '' TERM HUP; sleep 3001 & 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, "sleep 3001")
})

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, "sleep 3001")
})
}

// 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) {
cli, id, _ := startKeepAlive(t)

stream, err := cli.ExecStreaming(t.Context(), id, agentExecSpec("sh", "-c", "exit 3"))
if err != nil {
t.Fatalf("ExecStreaming: %v", err)
}
done := make(chan error, 1)
go func() { done <- stream.Process.Wait() }()
select {
case err = <-done:
case <-time.After(30 * time.Second):
t.Fatal("exec did not return after its command exited")
}
exitErr, ok := errors.AsType[*exec.ExitError](err)
if !ok || exitErr.ExitCode() != 3 {
t.Fatalf("Wait = %v, want exit status 3", err)
}
}

func agentExecSpec(command ...string) runtime.StreamingExecSpec {
return runtime.NewStreamingExecSpec(command...).AsUser(strconv.FormatUint(uint64(agentuid.AgentUID), 10))
}

// startKeepAlive starts an agent-image container running the production
// `sleep infinity` keep-alive and removes it at cleanup.
func startKeepAlive(t *testing.T) (*runtime.PodmanCLI, runtime.WorkloadID, string) {
t.Helper()
if !podmanUsable() {
t.Skip("rootless podman cannot run compass-agent:latest here; skipping the real-stack e2e")
}
cli := runtime.NewPodmanCLI()
name := "compass-e2e-exec-" + strconv.FormatInt(time.Now().UnixNano(), 36)
t.Cleanup(func() {
if out, err := exec.Command("podman", "rm", "--force", "--time", "0", name).CombinedOutput(); err != nil {
t.Logf("cleanup podman rm %s: %v: %s", name, err, out)
}
})
id, err := cli.Create(t.Context(), runtime.WorkloadSpec{
Image: agentImage,
Name: name,
UID: agentuid.AgentUID,
Command: []string{"sleep", "infinity"},
})
if err != nil {
t.Fatalf("Create: %v", err)
}
if err := cli.Start(t.Context(), id); err != nil {
t.Fatalf("Start: %v", err)
}
return cli, id, name
}

// isClientKill reports the deliberate SIGKILL of the host-side exec client.
func isClientKill(err error) bool {
exitErr, ok := errors.AsType[*exec.ExitError](err)
return ok && exitErr.ExitCode() == -1
}

func hasProcess(procs []string, args string) bool {
return slices.ContainsFunc(procs, func(p string) bool { return strings.Contains(p, args) })
}

// listProcs prints "pid state args" per container process from /proc. The
// agent image has no ps, and `podman top` exits 125 in the CI e2e container.
const listProcs = `for d in /proc/[0-9]*; do read -r s < "$d/stat" || continue; set -f -- $s; echo "${d#/proc/} $3 $(tr '\0' ' ' < "$d/cmdline")"; done 2>/dev/null`

// waitForTop polls the container's process list until done accepts the live
// processes. Zombies are excluded: the keep-alive never reaps an orphan.
func waitForTop(t *testing.T, name string, done func([]string) bool) {
t.Helper()
var procs []string
for deadline := time.Now().Add(15 * time.Second); time.Now().Before(deadline); {
out, err := exec.Command("podman", "exec", name, "sh", "-c", listProcs).CombinedOutput()
if err != nil {
t.Fatalf("list processes in %s: %v: %s", name, err, out)
}
procs = liveProcesses(string(out))
if done(procs) {
return
}
time.Sleep(50 * time.Millisecond) //nolint:forbidigo // bounded poll tick on the process list with a deadline (rule://go-no-sleep-in-test poll-until exemption)
}
t.Fatalf("container %s: live processes %q never reached the expected state", name, procs)
}

func liveProcesses(list string) []string {
var procs []string
for line := range strings.Lines(list) {
fields := strings.Fields(line)
if len(fields) < 3 || fields[1] == "Z" {
continue
}
procs = append(procs, strings.TrimSpace(line))
}
return procs
}
2 changes: 1 addition & 1 deletion go/internal/runtime/applecontainer.go
Original file line number Diff line number Diff line change
Expand Up @@ -275,7 +275,7 @@ func appleExecStreamingArgs(id WorkloadID, spec StreamingExecSpec) []string {
args = append(args, "--env", kv.key+"="+kv.value)
}
args = append(args, id.String())
args = append(args, spec.Command...)
args = append(args, stopWithClient(spec.Command)...)
return args
}

Expand Down
3 changes: 2 additions & 1 deletion go/internal/runtime/applecontainer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,8 @@ func TestAppleExecStreamingArgsAssemblesInteractiveExec(t *testing.T) {
"--workdir", "/work",
"--env", "COMPASS_MODEL=test-model",
"--env", "HOME=/home/agent",
"ctr123", "compass-agent",
"ctr123",
"sh", "-c", stopWithClientScript, "sh", "compass-agent",
}
if !slices.Equal(args, want) {
t.Fatalf("appleExecStreamingArgs = %q, want %q", args, want)
Expand Down
21 changes: 17 additions & 4 deletions go/internal/runtime/clispawn.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,12 +92,25 @@ func (e cliEngine) run(ctx context.Context, summary string, args []string) ([]by
return stdout, nil
}

// stopWithClientScript runs "$@" with stdin relayed by a watcher in the same
// process group. Killing the engine client closes that stdin, and the watcher
// then SIGKILLs the group; a natural exit kills the watcher and keeps the status.
// 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

// 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.
func stopWithClient(command []string) []string {
return append([]string{"sh", "-c", stopWithClientScript, "sh"}, command...)
}

// spawnStreaming starts `<program> <args>` streaming, returning the live pipes
// plus a kill/wait handle. The process is bound to a cancellable child of ctx:
// its Cancel SIGKILLs the process and WaitDelay bounds the reap, so cancelling
// the parent context or calling ChildHandle.Kill terminates the in-container
// agent even without a Go Drop. stdout/stderr are caller-owned os.Pipes, so
// Wait reaps at exit without closing them under output still in the pipe.
// its Cancel SIGKILLs the engine client and WaitDelay bounds the reap. The
// in-container process dies with the client only because exec argv builders
// wrap the command in stopWithClient. stdout/stderr are caller-owned os.Pipes,
// so Wait reaps at exit without closing them under output still in the pipe.
func (e cliEngine) spawnStreaming(ctx context.Context, args []string) (*StreamingExec, error) {
execCtx, cancel := context.WithCancel(ctx)
//nolint:gosec // G204: the container-engine seam — see spawnCapture. The
Expand Down
2 changes: 1 addition & 1 deletion go/internal/runtime/podman.go
Original file line number Diff line number Diff line change
Expand Up @@ -782,7 +782,7 @@ func execStreamingArgs(id WorkloadID, spec StreamingExecSpec) []string {
args = append(args, "-e", kv.key+"="+kv.value)
}
args = append(args, id.String())
args = append(args, spec.Command...)
args = append(args, stopWithClient(spec.Command)...)
return args
}

Expand Down
8 changes: 5 additions & 3 deletions go/internal/runtime/podman_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,8 @@ func TestExecStreamingArgsAssemblesInteractiveExec(t *testing.T) {
"-e", "COMPASS_MODEL=test-model",
"-e", "COMPASS_WORKDIR=/work",
"-e", "HOME=/home/agent",
"ctr123", "compass-agent",
"ctr123",
"sh", "-c", stopWithClientScript, "sh", "compass-agent",
}
if !slices.Equal(args, want) {
t.Fatalf("execStreamingArgs = %q, want %q", args, want)
Expand Down Expand Up @@ -237,7 +238,7 @@ func TestExecStreamingArgsMinimalOmitsUserAndWorkdir(t *testing.T) {

args := execStreamingArgs(WorkloadID("c"), spec)

want := []string{"exec", "--interactive", "c", "compass-agent"}
want := []string{"exec", "--interactive", "c", "sh", "-c", stopWithClientScript, "sh", "compass-agent"}
if !slices.Equal(args, want) {
t.Fatalf("execStreamingArgs = %q, want %q", args, want)
}
Expand Down Expand Up @@ -296,7 +297,8 @@ func TestExecStreamingArgsCarriesInlineEnvNotEnvFile(t *testing.T) {
"--user", "1000",
"--workdir", "/work",
"-e", "HOME=/home/agent",
"ctr123", "compass-agent",
"ctr123",
"sh", "-c", stopWithClientScript, "sh", "compass-agent",
}
if !slices.Equal(args, want) {
t.Fatalf("execStreamingArgs = %q, want %q", args, want)
Expand Down
Loading