diff --git a/docs/CLI-ADAPTER.md b/docs/CLI-ADAPTER.md index fee63bc..9ef4aef 100644 --- a/docs/CLI-ADAPTER.md +++ b/docs/CLI-ADAPTER.md @@ -48,9 +48,32 @@ The generated `exec.go` is defensive by default: OOM the adapter; truncation is flagged in the reply. - **Structured failures** — a non-zero exit is returned as `{"stdout","stderr","exit","truncated"}` rather than an opaque error, so the - caller sees everything the CLI produced. Only spawn failures (binary missing) - and timeouts surface as IPC errors; the per-method `timeout`/`duration` bounds - the run and the child is killed on cancel. + caller sees everything the CLI produced. Spawn failures (binary missing), + timeouts and a child ended by a signal surface as IPC errors, never as + `{"exit":-1}`; the per-method `timeout`/`duration` bounds the run and the + child is killed on cancel (SIGTERM, then SIGKILL after 5 s). +- **Children never outlive the adapter** — a child still running when the + adapter goes away is stopped, however it goes: a SIGTERM waits for in-flight + calls (and so their children); a dead supervisor makes the adapter shut down + (it watches its parent; macOS has no parent-death signal); a SIGKILLed + adapter's children get the Linux parent-death signal and, on every OS, are + taken down by the child guard (`internal/backend/childguard.go`, a copy of + the adapter that exists only while a child is in flight). Children stay in + the adapter's process group, so a supervisor group stop reaches them too. A + server a CLI deliberately daemonizes into its own session is outside all of + this. +- **Staging leaves nothing behind** — native-binary downloads go to + `$APP/.staged/tmp`, not the shared `TMPDIR`; a stop during the first-spawn + download removes the partial file, and a SIGKILLed start's leftover is swept + by the next start. +- **Ready before the binary** — an app that ships its binary (`assets`) listens + at once and stages the binary in the background: the supervisor waits only a + few seconds for the socket (and never marks a later one ready), and + `.help` needs no binary. A call that needs it waits within its own + deadline. A failed download is reported per call and retried by a later call + (backing off from 2 s to 5 min) instead of exiting, so a first start without + network recovers without a respawn. `tar` and install steps get the same + parent-death signal as CLI calls. ## Why HTTP works today and CLI doesn't diff --git a/docs/R2-ARTIFACT-REGISTRY.md b/docs/R2-ARTIFACT-REGISTRY.md index c4785c3..3697ffc 100644 --- a/docs/R2-ARTIFACT-REGISTRY.md +++ b/docs/R2-ARTIFACT-REGISTRY.md @@ -39,11 +39,14 @@ set install order + args fold into the bundle tarball fetc (`proc.exec`, `fs.write $APP`, `net.dial `). The whole tarball is sha-pinned in the catalogue, so `install.json` (and the expected asset shas) can't be altered undetected. -4. **Install + call** (host). The generated cli adapter calls `StageAssets($APP)` - on first spawn (`internal/backend/stage.go`): read `install.json` → select the - asset(s) for `runtime.GOOS/GOARCH` → in ascending `order`, fetch from R2, - verify sha256, stage under `$APP` (single file, or `tar.gz` extracted via the - host `tar`), run any install `args` — then exec the staged `exec_path` per call. +4. **Install + call** (host). The generated cli adapter stages its assets on + first spawn (`internal/backend/stage.go`, `NewStager` → `StageAssetsContext`), + in the background while it already serves: read `install.json` → select the + asset(s) for `runtime.GOOS/GOARCH` → in ascending `order`, fetch from R2 into + `$APP/.staged/tmp`, verify sha256, stage under `$APP` (single file, or `tar.gz` + extracted via the host `tar`), run any install `args` — then exec the staged + `exec_path` per call. A call that arrives before staging is done waits for it + within its own deadline; a failed download is retried by a later call. ## R2 layout diff --git a/internal/scaffold/cli_lifecycle_e2e_test.go b/internal/scaffold/cli_lifecycle_e2e_test.go new file mode 100644 index 0000000..3bfb03c --- /dev/null +++ b/internal/scaffold/cli_lifecycle_e2e_test.go @@ -0,0 +1,256 @@ +//go:build !windows + +package scaffold + +import ( + "bufio" + "encoding/json" + "errors" + "net" + "os" + "os/exec" + "path/filepath" + "runtime" + "strconv" + "strings" + "syscall" + "testing" + "time" + + "github.com/pilot-protocol/app-store/pkg/ipc" +) + +// buildCLIAdapter generates the adapter for spec and builds it, the way +// the other e2e tests do (go.sum seeded from this module, so the build needs no +// network). Returns the binary and the generated project dir. +func buildCLIAdapter(t *testing.T, spec string) (bin, proj string) { + t.Helper() + if testing.Short() { + t.Skip("builds and runs a real adapter binary; skipped under -short") + } + if _, err := exec.LookPath("go"); err != nil { + t.Skip("go toolchain not available") + } + root := t.TempDir() + cfg := parseSpec(t, spec) + proj = filepath.Join(root, "proj") + if _, err := Generate(cfg, proj); err != nil { + t.Fatalf("generate: %v", err) + } + if sum, err := os.ReadFile(filepath.Join("..", "..", "go.sum")); err == nil { + _ = os.WriteFile(filepath.Join(proj, "go.sum"), sum, 0o644) + } + bin = filepath.Join(root, "adapter") + build := exec.Command("go", "build", "-o", bin, "./cmd/"+cfg.BinaryName) + build.Dir = proj + build.Env = append(os.Environ(), "GOFLAGS=-mod=mod") + if out, err := build.CombinedOutput(); err != nil { + t.Fatalf("build adapter: %v\n%s", err, out) + } + return bin, proj +} + +// shortTempDir is a per-test dir with a short path: a t.TempDir() for a +// subtest can exceed sun_path (104 bytes on darwin) once app.sock is appended. +func shortTempDir(t *testing.T, prefix string) string { + t.Helper() + dir, err := os.MkdirTemp("", prefix) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { os.RemoveAll(dir) }) + return dir +} + +// waitForSocket waits until the adapter accepts connections on sock (the file +// appears at bind, a moment before listen). +func waitForSocket(t *testing.T, sock string, within time.Duration) { + t.Helper() + for deadline := time.Now().Add(within); time.Now().Before(deadline); time.Sleep(10 * time.Millisecond) { + if c, err := net.DialTimeout("unix", sock, time.Second); err == nil { + _ = c.Close() + return + } + } + t.Fatalf("adapter socket %s did not accept connections within %s", sock, within) +} + +// ipcCall dials the adapter once per call, like the supervisor does. +func ipcCall(sock, method, args string, deadline time.Duration) (json.RawMessage, error) { + conn, err := net.DialTimeout("unix", sock, 3*time.Second) + if err != nil { + return nil, err + } + defer conn.Close() + _ = conn.SetDeadline(time.Now().Add(deadline)) + var out json.RawMessage + err = ipc.Call(conn, method, json.RawMessage(args), &out) + return out, err +} + +// processGone reports whether pid has exited within the given time (a zombie +// counts as alive: its parent has not reaped it yet). It SIGKILLs a survivor so +// a failing test leaves nothing behind. +func processGone(pid int, within time.Duration) bool { + for deadline := time.Now().Add(within); ; time.Sleep(20 * time.Millisecond) { + if err := syscall.Kill(pid, 0); errors.Is(err, syscall.ESRCH) { + return true + } + if time.Now().After(deadline) { + _ = syscall.Kill(pid, syscall.SIGKILL) + return false + } + } +} + +// TestCLIAdapterChildLifecycleE2E: a CLI child still running when its call +// times out, when the adapter is stopped, or when the process that spawned the +// adapter dies must not outlive the call or the adapter, and a timed-out call +// must fail rather than report a successful {"exit":-1}. io.pilot.miren +// 0.1.0: `miren login` started over miren.exec was left running, reparented to +// init, after the supervisor SIGKILLed the adapter and after the daemon died +// (Linux: the adapter's Pdeathsig is not inherited by its children; macOS: the +// adapter itself kept serving as an orphan), and a call cut off by its 60s +// deadline replied {"exit":-1} as a success. +func TestCLIAdapterChildLifecycleE2E(t *testing.T) { + root := t.TempDir() + // `hang ` records its pid, then blocks without writing anything + // (so SIGPIPE cannot end it early once the adapter is gone). + tool := filepath.Join(root, "faketool") + script := "#!/bin/sh\n" + + "case \"$1\" in\n" + + " hang) echo $$ > \"$2\"; exec sleep 300;;\n" + + " *) exit 2;;\n" + + "esac\n" + if err := os.WriteFile(tool, []byte(script), 0o755); err != nil { + t.Fatal(err) + } + bin, proj := buildCLIAdapter(t, ` +id: io.pilot.faketool +app_version: 0.1.0 +description: "Fronts faketool." +namespace: faketool +backend: + type: cli + command: ["`+tool+`"] +methods: + - name: faketool.hang + summary: "Blocks until killed." + timeout: 4s # room for the child to start on a loaded host before the deadline + cli: {args: ["hang", "${pidfile}"]} + - name: faketool.run + summary: "Passthrough." + cli: {passthrough: true} +`) + manifest := filepath.Join(proj, "manifest.json") + + start := func(t *testing.T) (*exec.Cmd, string) { + t.Helper() + sock := filepath.Join(shortTempDir(t, "cla"), "app.sock") + a := exec.Command(bin, "--socket", sock, "--manifest", manifest) + a.Stderr = os.Stderr + a.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} // like the supervisor + if err := a.Start(); err != nil { + t.Fatalf("start adapter: %v", err) + } + t.Cleanup(func() { _ = a.Process.Kill(); _, _ = a.Process.Wait() }) + waitForSocket(t, sock, 10*time.Second) + return a, sock + } + childPID := func(t *testing.T, pidfile string) int { + t.Helper() + for deadline := time.Now().Add(10 * time.Second); time.Now().Before(deadline); time.Sleep(20 * time.Millisecond) { + if b, err := os.ReadFile(pidfile); err == nil && len(b) > 0 { + if pid, err := strconv.Atoi(strings.TrimSpace(string(b))); err == nil { + return pid + } + } + } + t.Fatal("child never started") + return 0 + } + hangInFlight := func(t *testing.T, sock string) int { + t.Helper() + pidfile := filepath.Join(t.TempDir(), "pid") + go func() { _, _ = ipcCall(sock, "faketool.run", `{"args":["hang","`+pidfile+`"]}`, 30*time.Second) }() + return childPID(t, pidfile) + } + + t.Run("timeout is an error and kills the child", func(t *testing.T) { + _, sock := start(t) + pidfile := filepath.Join(t.TempDir(), "pid") + out, err := ipcCall(sock, "faketool.hang", `{"pidfile":"`+pidfile+`"}`, 30*time.Second) + if err == nil || !strings.Contains(err.Error(), "deadline exceeded") { + t.Fatalf("timed-out call: out=%s err=%v, want an IPC error naming the deadline (not a {\"exit\":-1} result)", out, err) + } + if pid := childPID(t, pidfile); !processGone(pid, 3*time.Second) { + t.Errorf("child pid %d still running after its call timed out", pid) + } + }) + + t.Run("SIGTERM stops the in-flight child before the adapter exits", func(t *testing.T) { + a, sock := start(t) + pid := hangInFlight(t, sock) + _ = a.Process.Signal(syscall.SIGTERM) + if _, err := a.Process.Wait(); err != nil { + t.Fatal(err) + } + if !processGone(pid, 100*time.Millisecond) { + t.Errorf("in-flight child pid %d outlived the adapter's SIGTERM exit", pid) + } + if _, err := os.Stat(sock); !os.IsNotExist(err) { + t.Errorf("socket left after a clean exit: %v", err) + } + }) + + t.Run("SIGKILL of the adapter takes the in-flight child down", func(t *testing.T) { + // Linux: the child's parent-death signal. Every OS: the child guard + // (childguard.go), which needs the adapter to lead its process group. + a, sock := start(t) + pid := hangInFlight(t, sock) + _ = a.Process.Kill() + _, _ = a.Process.Wait() + if !processGone(pid, 3*time.Second) { + t.Errorf("in-flight child pid %d outlived the SIGKILLed adapter (%s)", pid, runtime.GOOS) + } + }) + + t.Run("parent death stops the adapter and its in-flight child", func(t *testing.T) { + // The adapter's parent is a shell standing in for the supervisor. It is + // SIGKILLed; the adapter (no Pdeathsig here, as on macOS) must notice, + // cancel the running call, reap its child and exit. + sock := filepath.Join(shortTempDir(t, "clp"), "app.sock") + parent := exec.Command("/bin/sh", "-c", `"$0" "$@" & echo $!; wait`, bin, "--socket", sock, "--manifest", manifest) + parent.Stderr = os.Stderr + stdout, err := parent.StdoutPipe() + if err != nil { + t.Fatal(err) + } + if err := parent.Start(); err != nil { + t.Fatal(err) + } + line, err := bufio.NewReader(stdout).ReadString('\n') + if err != nil { + t.Fatalf("read adapter pid: %v", err) + } + adapterPID, err := strconv.Atoi(strings.TrimSpace(line)) + if err != nil { + t.Fatalf("adapter pid %q: %v", line, err) + } + t.Cleanup(func() { _ = syscall.Kill(adapterPID, syscall.SIGKILL) }) + waitForSocket(t, sock, 10*time.Second) + pid := hangInFlight(t, sock) + + _ = parent.Process.Kill() + _ = parent.Wait() + if !processGone(adapterPID, 5*time.Second) { + t.Fatalf("adapter pid %d kept running after its parent died", adapterPID) + } + if !processGone(pid, 500*time.Millisecond) { + t.Errorf("in-flight child pid %d outlived the orphaned adapter", pid) + } + if _, err := os.Stat(sock); !os.IsNotExist(err) { + t.Errorf("socket left after the adapter shut down: %v", err) + } + }) +} diff --git a/internal/scaffold/scaffold.go b/internal/scaffold/scaffold.go index d00a8bf..22abf52 100644 --- a/internal/scaffold/scaffold.go +++ b/internal/scaffold/scaffold.go @@ -21,6 +21,17 @@ type file struct { tmpl string } +// childProcFiles are the per-OS process attributes the exec runner applies to +// every CLI child (Linux: parent-death signal), plus the guard that takes an +// in-flight child down with a SIGKILLed adapter; emitted with client_cli.go.tmpl. +func childProcFiles() []file { + return []file{ + {filepath.Join("internal", "backend", "childproc_linux.go"), "childproc_linux.go.tmpl"}, + {filepath.Join("internal", "backend", "childproc_other.go"), "childproc_other.go.tmpl"}, + {filepath.Join("internal", "backend", "childguard.go"), "childguard.go.tmpl"}, + } +} + // Generate renders a full adapter project for cfg into outDir. cfg must already // be Resolve()d and Validate()d. Returns the list of written paths. func Generate(cfg *Config, outDir string) ([]string, error) { @@ -69,6 +80,7 @@ func Generate(cfg *Config, outDir string) ([]string, error) { } case "cli": files = append(files, file{filepath.Join("internal", "backend", "exec.go"), "client_cli.go.tmpl"}) + files = append(files, childProcFiles()...) // Native-binary delivery: emit the staging runtime only when the app // actually ships assets (an already-installed cli needs no stager). if cfg.HasAssets() { @@ -83,6 +95,7 @@ func Generate(cfg *Config, outDir string) ([]string, error) { file{filepath.Join("internal", "backend", "cloud.go"), "cloud.go.tmpl"}, file{filepath.Join("internal", "backend", "signer.go"), "signer.go.tmpl"}, ) + files = append(files, childProcFiles()...) if cfg.HasAssets() { files = append(files, file{filepath.Join("internal", "backend", "stage.go"), "stage.go.tmpl"}) } diff --git a/internal/scaffold/stage_lifecycle_e2e_test.go b/internal/scaffold/stage_lifecycle_e2e_test.go new file mode 100644 index 0000000..a36b5a0 --- /dev/null +++ b/internal/scaffold/stage_lifecycle_e2e_test.go @@ -0,0 +1,273 @@ +//go:build !windows + +package scaffold + +import ( + "crypto/sha256" + "encoding/hex" + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "os/exec" + "path/filepath" + "runtime" + "strings" + "sync" + "sync/atomic" + "syscall" + "testing" + "time" +) + +// TestStageLifecycleE2E covers first-start asset staging for a cli app that +// ships its binary (io.pilot.miren 0.1.0 stages a 40-60 MB CLI tarball): +// +// - The socket must appear, and help must answer, while the binary is still +// downloading. The supervisor gives a new app 3s to create its socket and +// otherwise never marks it ready, so a download that ran before the socket +// (6-25 s for miren) left the app answering "app not ready" until its next +// respawn. +// - A stop during the download is a clean exit that removes the partial file; +// a SIGKILLed download's leftover is swept by the next start. Neither may +// land in the shared TMPDIR (miren stranded a pilot-asset-* there per +// interrupted start). +// - A failed download (no network on first start) is not fatal: help keeps +// answering, a call reports the failure, and a later call retries. +func TestStageLifecycleE2E(t *testing.T) { + bin, proj := buildCLIAdapter(t, cliAssetsSpec) + manifest, err := os.ReadFile(filepath.Join(proj, "manifest.json")) + if err != nil { + t.Fatal(err) + } + + // A ~2 MB "binary". Served at once; or the first 64 KiB, then nothing until + // release is closed (mode "stall") or until the client goes away (mode + // "hang"); or failing with HTTP 500 while fail is set. + body := []byte("#!/bin/sh\necho toolx 1.0\n" + strings.Repeat("#"+strings.Repeat("x", 1022)+"\n", 2048)) + sum := sha256.Sum256(body) + var ( + fail atomic.Bool + releaseM sync.Mutex + release = make(chan struct{}) + ) + releaseSlow := func() { + releaseM.Lock() + defer releaseM.Unlock() + select { + case <-release: + default: + close(release) + } + } + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if fail.Load() { + http.Error(w, "registry down", http.StatusInternalServerError) + return + } + mode := r.URL.Query().Get("mode") + if mode == "" { + _, _ = w.Write(body) + return + } + _, _ = w.Write(body[:64<<10]) + w.(http.Flusher).Flush() + wait := release + if mode == "hang" { + wait = nil // blocks until the client disconnects + } + select { + case <-r.Context().Done(): + return + case <-wait: + } + _, _ = w.Write(body[64<<10:]) + })) + defer srv.Close() + defer releaseSlow() + + type env struct{ app, tmp, sock string } + // install.json is read at runtime, so point this host's asset at the local + // server (the generated spec's registry URLs are placeholders). + writeSpec := func(t *testing.T, e env, mode string) { + t.Helper() + url := srv.URL + "/toolx" + if mode != "" { + url += "?mode=" + mode + } + spec := map[string]any{"schema": 1, "app": "io.pilot.toolx", "version": "0.2.0", "command": "toolx", + "assets": []map[string]any{{"role": "binary", "os": runtime.GOOS, "arch": runtime.GOARCH, + "url": url, "sha256": hex.EncodeToString(sum[:]), "exec_path": "bin/toolx", "order": 1}}} + raw, _ := json.Marshal(spec) + if err := os.WriteFile(filepath.Join(e.app, "install.json"), raw, 0o644); err != nil { + t.Fatal(err) + } + } + setup := func(t *testing.T, mode string) env { + t.Helper() + dir := shortTempDir(t, "stg") + e := env{app: filepath.Join(dir, "app"), tmp: filepath.Join(dir, "tmp"), sock: filepath.Join(dir, "app", "app.sock")} + for _, d := range []string{e.app, e.tmp} { + if err := os.MkdirAll(d, 0o755); err != nil { + t.Fatal(err) + } + } + if err := os.WriteFile(filepath.Join(e.app, "manifest.json"), manifest, 0o644); err != nil { + t.Fatal(err) + } + writeSpec(t, e, mode) + return e + } + start := func(t *testing.T, e env) *exec.Cmd { + t.Helper() + a := exec.Command(bin, "--socket", e.sock, "--manifest", filepath.Join(e.app, "manifest.json")) + a.Env = append(os.Environ(), "TMPDIR="+e.tmp) + a.Stderr = os.Stderr + a.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} // like the supervisor + if err := a.Start(); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = a.Process.Kill(); _, _ = a.Process.Wait() }) + return a + } + waitPartial := func(t *testing.T, e env) { + t.Helper() + for deadline := time.Now().Add(10 * time.Second); time.Now().Before(deadline); time.Sleep(20 * time.Millisecond) { + m, _ := filepath.Glob(filepath.Join(e.app, ".staged", "tmp", "pilot-asset-*")) + for _, f := range m { + if fi, err := os.Stat(f); err == nil && fi.Size() > 0 { + return + } + } + } + t.Fatal("download never started") + } + noTmpLeft := func(t *testing.T, e env) { + t.Helper() + if ents, _ := os.ReadDir(e.tmp); len(ents) > 0 { + t.Errorf("TMPDIR not empty: %v", ents) + } + } + version := func(t *testing.T, e env) (string, error) { + t.Helper() + out, err := ipcCall(e.sock, "toolx.version", `{}`, 30*time.Second) + return string(out), err + } + + t.Run("socket and help are up while the binary is still downloading", func(t *testing.T) { + e := setup(t, "stall") + started := time.Now() + start(t, e) + // The supervisor's readiness window is 3s; the download is stalled. + waitForSocket(t, e.sock, 2500*time.Millisecond) + if _, err := ipcCall(e.sock, "toolx.help", `{}`, 5*time.Second); err != nil { + t.Fatalf("help while staging: %v", err) + } + t.Logf("socket + help answered %s after start, download still stalled", time.Since(started).Round(time.Millisecond)) + waitPartial(t, e) + + // A call that needs the binary waits for it, then runs it. + type res struct { + out string + err error + } + got := make(chan res, 1) + go func() { o, err := version(t, e); got <- res{o, err} }() + select { + case r := <-got: + t.Fatalf("call returned before the binary was installed: %q %v", r.out, r.err) + case <-time.After(300 * time.Millisecond): + } + releaseSlow() + r := <-got + if r.err != nil || !strings.Contains(r.out, "toolx 1.0") { + t.Fatalf("call after staging: out=%q err=%v", r.out, r.err) + } + if _, err := os.Stat(filepath.Join(e.app, ".staged", "tmp")); !os.IsNotExist(err) { + t.Errorf("$APP/.staged/tmp left after staging: %v", err) + } + noTmpLeft(t, e) + }) + + t.Run("SIGTERM mid-download exits cleanly and removes the partial", func(t *testing.T) { + e := setup(t, "hang") + a := start(t, e) + waitPartial(t, e) + _ = a.Process.Signal(syscall.SIGTERM) + done := make(chan error, 1) + go func() { done <- a.Wait() }() + select { + case err := <-done: + if err != nil { + t.Errorf("adapter stopped mid-staging: %v (want a clean exit)", err) + } + case <-time.After(5 * time.Second): + t.Fatal("adapter did not exit within 5s of SIGTERM while staging") + } + if _, err := os.Stat(filepath.Join(e.app, ".staged", "tmp")); !os.IsNotExist(err) { + t.Errorf("$APP/.staged/tmp left behind: %v", err) + } + if _, err := os.Stat(e.sock); !os.IsNotExist(err) { + t.Errorf("socket left behind: %v", err) + } + noTmpLeft(t, e) + }) + + t.Run("SIGKILL mid-download: next start sweeps the leftover", func(t *testing.T) { + e := setup(t, "hang") + a := start(t, e) + waitPartial(t, e) + _ = a.Process.Kill() + _, _ = a.Process.Wait() + noTmpLeft(t, e) + writeSpec(t, e, "") + // A SIGKILLed adapter leaves app.sock; the supervisor removes a stale + // socket before it spawns the next instance. + _ = os.Remove(e.sock) + start(t, e) + waitForSocket(t, e.sock, 10*time.Second) + if out, err := version(t, e); err != nil || !strings.Contains(out, "toolx 1.0") { + t.Fatalf("call after restart: out=%q err=%v", out, err) + } + if _, err := os.Stat(filepath.Join(e.app, ".staged", "tmp")); !os.IsNotExist(err) { + t.Errorf("$APP/.staged/tmp not swept by the next start: %v", err) + } + noTmpLeft(t, e) + }) + + t.Run("a failed download keeps serving and a later call retries", func(t *testing.T) { + fail.Store(true) + defer fail.Store(false) + e := setup(t, "") + a := start(t, e) + waitForSocket(t, e.sock, 2500*time.Millisecond) + // The first attempt fails; the call reports it as an error. + var err error + for deadline := time.Now().Add(5 * time.Second); time.Now().Before(deadline); time.Sleep(50 * time.Millisecond) { + if _, err = version(t, e); err != nil && strings.Contains(err.Error(), "HTTP 500") { + break + } + } + if err == nil || !strings.Contains(err.Error(), "install assets") || !strings.Contains(err.Error(), "HTTP 500") { + t.Fatalf("call with the registry down: err=%v, want an install-assets error naming HTTP 500", err) + } + if _, err := ipcCall(e.sock, "toolx.help", `{}`, 5*time.Second); err != nil { + t.Fatalf("help after a failed download: %v", err) + } + if err := syscall.Kill(a.Process.Pid, 0); err != nil { + t.Fatalf("adapter exited after a failed download: %v", err) + } + // Registry back: a call after the retry delay installs and runs it. + fail.Store(false) + var out string + for deadline := time.Now().Add(15 * time.Second); time.Now().Before(deadline); time.Sleep(250 * time.Millisecond) { + if out, err = version(t, e); err == nil { + break + } + } + if err != nil || !strings.Contains(out, "toolx 1.0") { + t.Fatalf("call once the registry is back: out=%q err=%v", out, err) + } + noTmpLeft(t, e) + }) +} diff --git a/internal/scaffold/templates/childguard.go.tmpl b/internal/scaffold/templates/childguard.go.tmpl new file mode 100644 index 0000000..064402d --- /dev/null +++ b/internal/scaffold/templates/childguard.go.tmpl @@ -0,0 +1,206 @@ +// Child guard for {{.ID}}. GENERATED by pilot-app; edit pilot.app.yaml and +// re-generate. +// +// A CLI child that is still running when this adapter dies must die with it. +// On Linux the kernel does that for the child itself (Pdeathsig, see +// childproc_linux.go). macOS has no parent-death signal, and on neither OS +// does it reach what the child started. When the adapter is SIGKILLed (the +// app-store supervisor's stop path through v1.0.3 signals only the adapter's +// pid; also the OOM killer, an operator, a crash) nothing else would ever stop +// that child: it is reparented to launchd/init, keeps whatever it holds (a +// listening port, a download), and the supervisor's stale-instance reaper only +// matches the adapter's own argv. Found with io.pilot.otto: an in-flight +// `otto.exec ["mcp","serve-http",...]` kept its port open after the adapter +// was SIGKILLed on macOS. +// +// So while a CLI child is in flight the adapter keeps a guard: this same binary +// re-executed with guardArg, in the adapter's process group, reading a pipe +// whose write end only the adapter holds (close-on-exec, so no CLI child +// inherits it). The kernel closes that pipe however the adapter goes away. On +// EOF the guard waits until it has been reparented, which proves the adapter +// is gone rather than merely done with it, and then SIGKILLs the process +// group: every CLI child and whatever they started in it, the guard included. +// A process group id is not reused while any member is alive, and the guard is +// a member, so the signal cannot reach an unrelated group. An adapter with no +// child in flight for guardIdle tells its guard to go ("release") and the next +// call starts a new one. The guard runs only when the adapter leads its own +// process group, as the supervisor spawns it, so it never signals a group the +// adapter does not own. +package backend + +import ( + "bufio" + "log" + "os" + "os/exec" + "os/signal" + "strconv" + "sync" + "syscall" + "time" +) + +// guardArg is the argv[1] that turns this binary into the child guard; argv[2] +// is the adapter's pid. +const guardArg = "--pilot-child-guard" + +// guardIdle is how long the guard outlives the last in-flight child, so a +// burst of back-to-back calls shares one guard instead of starting one per +// call. Short: an idle adapter should not keep an extra process around. +const guardIdle = time.Second + +type childGuard struct { + mu sync.Mutex + w *os.File // write end of the guard's stdin; nil = no guard running + done chan struct{} + inflight int + idle *time.Timer + off bool // this adapter does not lead its process group +} + +var guard childGuard + +// acquire is called before a CLI child starts: it makes sure a guard is +// running. Failing to start one is logged, not fatal: the call still runs, +// and on Linux Pdeathsig still covers the child itself. +func (g *childGuard) acquire() { + g.mu.Lock() + defer g.mu.Unlock() + g.inflight++ + if g.idle != nil { + g.idle.Stop() + g.idle = nil + } + if g.w != nil || g.off { + return + } + if syscall.Getpgrp() != os.Getpid() { + g.off = true // not our group to take down + return + } + if err := g.startLocked(); err != nil { + log.Printf("child guard: %v (continuing without it)", err) + } +} + +// release is called once a CLI child has been reaped. +func (g *childGuard) release() { + g.mu.Lock() + defer g.mu.Unlock() + g.inflight-- + if g.inflight > 0 || g.w == nil { + return + } + g.idle = time.AfterFunc(guardIdle, func() { + g.mu.Lock() + defer g.mu.Unlock() + if g.inflight == 0 { + g.retireLocked() + } + }) +} + +// retireLocked tells the guard to exit without acting and closes its pipe. +func (g *childGuard) retireLocked() { + if g.w == nil { + return + } + _, _ = g.w.WriteString("release\n") + _ = g.w.Close() + g.w = nil +} + +func (g *childGuard) startLocked() error { + exe, err := os.Executable() + if err != nil { + return err + } + r, w, err := os.Pipe() + if err != nil { + return err + } + // The adapter's pid is passed in, not read with getppid() in the guard: a + // guard that is slow to start could otherwise read the pid it was already + // reparented to and never act. + cmd := exec.Command(exe, guardArg, strconv.Itoa(os.Getpid())) + cmd.Stdin = r + cmd.Stderr = os.Stderr + cmd.Env = []string{} + // No Setpgid: the guard must be in the group it guards. No Pdeathsig: it + // has to outlive the adapter. + if err := cmd.Start(); err != nil { + _ = r.Close() + _ = w.Close() + return err + } + _ = r.Close() + done := make(chan struct{}) + g.w, g.done = w, done + go func() { + _ = cmd.Wait() + close(done) + g.mu.Lock() + defer g.mu.Unlock() + if g.w == w { // it died on its own: start a new one on the next call + _ = w.Close() + g.w = nil + } + }() + return nil +} + +// CloseChildGuard is main's last step on a clean shutdown. With nothing in +// flight the guard is released and waited for, so the adapter leaves an empty +// process group behind. If a child is somehow still running, the guard is left +// in place to take the group down once the adapter has exited. +func CloseChildGuard(timeout time.Duration) { + g := &guard + g.mu.Lock() + if g.inflight > 0 || g.w == nil { + g.mu.Unlock() + return + } + if g.idle != nil { + g.idle.Stop() + g.idle = nil + } + done := g.done + g.retireLocked() + g.mu.Unlock() + select { + case <-done: + case <-time.After(timeout): + } +} + +// IsChildGuard reports whether this process was started as the child guard. +func IsChildGuard(args []string) bool { return len(args) == 3 && args[1] == guardArg } + +// ChildGuardMain is the guard process. It ignores the usual stop signals (a +// supervisor may signal the whole group; the guard then waits for the adapter +// to finish) and ends on its own. +func ChildGuardMain() { + signal.Ignore(syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP, syscall.SIGPIPE) + adapter, err := strconv.Atoi(os.Args[2]) + if err != nil || adapter <= 1 { + return + } + released := false + sc := bufio.NewScanner(os.Stdin) + for sc.Scan() { + if sc.Text() == "release" { + released = true + } + } + if released { + return + } + // The kernel closes a dying process's descriptors before it reparents + // its children: wait for the reparent before acting on the EOF. + for deadline := time.Now().Add(5 * time.Second); os.Getppid() == adapter; time.Sleep(5 * time.Millisecond) { + if time.Now().After(deadline) { + return // the adapter is alive and closed the pipe itself + } + } + _ = syscall.Kill(-syscall.Getpgrp(), syscall.SIGKILL) +} diff --git a/internal/scaffold/templates/childproc_linux.go.tmpl b/internal/scaffold/templates/childproc_linux.go.tmpl new file mode 100644 index 0000000..7d2e43f --- /dev/null +++ b/internal/scaffold/templates/childproc_linux.go.tmpl @@ -0,0 +1,26 @@ +//go:build linux + +package backend + +import ( + "runtime" + "syscall" +) + +// childProcAttr keeps every CLI child in the adapter's process group (so the +// supervisor's group-wide stop and its orphan reaper reach it) and asks the +// kernel to SIGKILL the child when the adapter dies. Without Pdeathsig an +// in-flight child outlives an adapter that is SIGKILLed — the supervisor's +// stop path and, via the adapter's own Pdeathsig, a daemon crash — and is +// reparented to init. GENERATED by pilot-app. +func childProcAttr() *syscall.SysProcAttr { + return &syscall.SysProcAttr{Pdeathsig: syscall.SIGKILL} +} + +// lockForChild pins the goroutine that forks a CLI child to its OS thread for +// the child's lifetime: Linux delivers Pdeathsig when the forking *thread* +// exits, and the Go runtime may retire an idle thread. Returns the unlock. +func lockForChild() func() { + runtime.LockOSThread() + return runtime.UnlockOSThread +} diff --git a/internal/scaffold/templates/childproc_other.go.tmpl b/internal/scaffold/templates/childproc_other.go.tmpl new file mode 100644 index 0000000..b62d4cc --- /dev/null +++ b/internal/scaffold/templates/childproc_other.go.tmpl @@ -0,0 +1,12 @@ +//go:build !linux + +package backend + +import "syscall" + +// childProcAttr: there is no parent-death signal outside Linux. The child stays +// in the adapter's process group, so the supervisor's group-wide stop and its +// stale-instance reaper (kill(-pgid)) still reach it. GENERATED by pilot-app. +func childProcAttr() *syscall.SysProcAttr { return nil } + +func lockForChild() func() { return func() {} } diff --git a/internal/scaffold/templates/client_cli.go.tmpl b/internal/scaffold/templates/client_cli.go.tmpl index 9edcb5f..e6b276b 100644 --- a/internal/scaffold/templates/client_cli.go.tmpl +++ b/internal/scaffold/templates/client_cli.go.tmpl @@ -41,6 +41,7 @@ import ( "regexp" "sort" "strings" + "syscall" "time" ) @@ -48,8 +49,8 @@ import ( // the adapter's memory. Output past the cap is dropped and flagged truncated. const maxOutputBytes = 4 << 20 // 4 MiB -// killGrace is how long the child has to exit after its context is cancelled -// before Wait abandons it; pairs with the SIGKILL exec sends on cancel. +// killGrace is how long the child has, after the SIGTERM sent on cancel, to stop +// before exec escalates to SIGKILL and Wait stops waiting on its pipes. const killGrace = 5 * time.Second // baseEnvKeys is the minimal environment every fronted CLI may see. Anything an @@ -62,6 +63,10 @@ var placeholderRE = regexp.MustCompile(`\$\{([^}]+)\}`) type Runner struct { base []string // e.g. ["weathercli"] or ["python", "-m", "tool"] envKeys []string // extra host env vars forwarded to the child (allowlist) + // resolve, when set, returns the path of the command binary for a call + // (the staged asset, which may still be downloading on first start). "" + // keeps base[0]. + resolve func(context.Context) (string, error) } // Spec is one method's invocation shape, baked in at generation time. @@ -85,6 +90,15 @@ func NewRunner(base []string, envKeys ...string) *Runner { return &Runner{base: base, envKeys: envKeys} } +// SetCommandResolver makes every call resolve its command binary through +// resolve instead of using base[0] as-is: main wires it to the asset stager, so +// the adapter can listen (and serve help) while the binary is still being +// installed, and a call waits for it, bounded by the call's own deadline. An +// error from resolve fails the call. +func (r *Runner) SetCommandResolver(resolve func(context.Context) (string, error)) { + r.resolve = resolve +} + // Run executes the command for one call. On clean exit with pure-JSON stdout the // JSON is passed through verbatim; otherwise the reply is wrapped as // {"stdout","stderr","exit","truncated"} so the IPC reply is always a JSON @@ -112,9 +126,28 @@ func (r *Runner) Run(ctx context.Context, spec Spec, call Call) (json.RawMessage } } - cmd := exec.CommandContext(ctx, r.base[0], args...) + argv0 := r.base[0] + if r.resolve != nil { + p, err := r.resolve(ctx) + if err != nil { + return nil, fmt.Errorf("backend: %s: %w", r.base[0], err) + } + if p != "" { + argv0 = p + } + } + + cmd := exec.CommandContext(ctx, argv0, args...) cmd.Env = r.childEnv() + // On timeout/cancel, ask the CLI to stop (SIGTERM) before exec's WaitDelay + // escalates to SIGKILL. A bare SIGKILL gives the CLI no chance to stop what + // it started itself: aegis's model download (curl) and its L2 llama-server + // kept running, reparented to init, after every timed-out call. + cmd.Cancel = func() error { return cmd.Process.Signal(syscall.SIGTERM) } cmd.WaitDelay = killGrace + // Same process group as the adapter (so a group-wide stop reaches it) and, + // on Linux, killed by the kernel if the adapter dies mid-call. + cmd.SysProcAttr = childProcAttr() if call.Stdin != "" { cmd.Stdin = strings.NewReader(call.Stdin) } @@ -123,17 +156,29 @@ func (r *Runner) Run(ctx context.Context, spec Spec, call Call) (json.RawMessage cmd.Stdout = stdout cmd.Stderr = stderr + // While the child runs, a guard process stands by to take it down if this + // adapter is SIGKILLed (childguard.go). + guard.acquire() + unlock := lockForChild() runErr := cmd.Run() + unlock() + guard.release() exitCode := 0 if runErr != nil { var ee *exec.ExitError - if errors.As(runErr, &ee) { + if ctxErr := ctx.Err(); ctxErr != nil { + // exec killed the child because the call's deadline passed or the + // adapter is shutting down. That is not a CLI result: an ExitError + // here only says "signal: killed" (exit -1). + return nil, fmt.Errorf("backend: %s: %w: %s", argv0, ctxErr, strings.TrimSpace(stderr.String())) + } + if errors.As(runErr, &ee) && ee.Exited() { exitCode = ee.ExitCode() // ran, exited non-zero — a normal CLI result } else { - // Failed to start (binary missing) or killed (timeout/cancel): the - // command never produced a meaningful result. Surface as an error. - return nil, fmt.Errorf("backend: %s: %w: %s", r.base[0], runErr, strings.TrimSpace(stderr.String())) + // Failed to start (binary missing) or ended by a signal: the command + // never produced a meaningful result. Surface as an error. + return nil, fmt.Errorf("backend: %s: %w: %s", argv0, runErr, strings.TrimSpace(stderr.String())) } } diff --git a/internal/scaffold/templates/main.go.tmpl b/internal/scaffold/templates/main.go.tmpl index 7fbf31f..f24a042 100644 --- a/internal/scaffold/templates/main.go.tmpl +++ b/internal/scaffold/templates/main.go.tmpl @@ -18,8 +18,10 @@ import ( {{- end}} "context" "encoding/json" + "errors" "flag" "fmt" + "io" "log" "net" {{- if eq .Backend.Type "http"}} @@ -31,6 +33,7 @@ import ( {{- if or .Backend.Headers (eq .Backend.Type "http")}} "strings" {{- end}} + "sync" "syscall" "time" @@ -39,6 +42,17 @@ import ( ) func main() { +{{- if or (eq .Backend.Type "cli") (eq .Backend.Type "hybrid")}} + // This binary doubles as the guard that stops in-flight CLI children if + // the adapter is killed (backend/childguard.go); that mode takes no flags. + if backend.IsChildGuard(os.Args) { + backend.ChildGuardMain() + return + } +{{- end}} + // Captured first, before anything slow (asset staging), so a parent that + // dies while we start up is still noticed by watchParent. + parentPID := os.Getppid() // The pilot daemon spawns every app with the full standard lifecycle flag // set. We MUST define all of them or Go's flag package hard-fails ("flag // provided but not defined") and the supervisor sees a crash-loop. This @@ -63,6 +77,16 @@ func main() { log.Fatal("{{.BinaryName}}: --socket is required (the pilot daemon supplies it)") } + // Installed before any startup work (asset staging can take a while on + // first spawn) so a stop during startup is handled too: staging is + // cancelled and cleans up instead of the default SIGTERM action killing + // the process mid-download. + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stop() + ctx, parentGone := context.WithCancel(ctx) + defer parentGone() + go watchParent(ctx, parentPID, parentGone) + {{- if eq .Backend.Type "http"}} cfg := resolveConfig(*manifestPath) {{- if or .Managed .HasBrokerSignup}} @@ -84,15 +108,12 @@ func main() { appDir = filepath.Dir(*manifestPath) } {{- if .HasAssets}} - stagedCmd, err := backend.StageAssets(appDir) - if err != nil { - log.Fatalf("{{.BinaryName}}: install assets: %v", err) - } - base := {{printf "%#v" .Backend.Command}} - if stagedCmd != "" { - base[0] = stagedCmd - } - runner := backend.NewRunner(base{{range .Backend.EnvPassthrough}}, {{printf "%q" .}}{{end}}) + // Assets are staged in the background (see backend/stage.go): the socket + // must appear within the supervisor's readiness window, which a first-start + // download does not fit in. Local calls wait for the binary. + stager := backend.NewStager(ctx, appDir) + runner := backend.NewRunner({{printf "%#v" .Backend.Command}}{{range .Backend.EnvPassthrough}}, {{printf "%q" .}}{{end}}) + runner.SetCommandResolver(stager.Command) {{- else}} runner := backend.NewRunner({{printf "%#v" .Backend.Command}}{{range .Backend.EnvPassthrough}}, {{printf "%q" .}}{{end}}) {{- end}} @@ -104,20 +125,17 @@ func main() { {{- if .HasAssets}} // Native delivery: fetch this host's binaries from the Pilot R2 artifact // registry (verify sha → stage under $APP → run ordered install args) and - // exec the staged path, not an assumed-installed command. See backend/stage.go. + // exec the staged path, not an assumed-installed command. Staging runs in the + // background (see backend/stage.go): the socket must appear within the + // supervisor's readiness window, which a first-start download does not fit + // in. Help answers at once; a call that needs the binary waits for it. appDir := os.Getenv("APP") if *manifestPath != "" { appDir = filepath.Dir(*manifestPath) } - stagedCmd, err := backend.StageAssets(appDir) - if err != nil { - log.Fatalf("{{.BinaryName}}: install assets: %v", err) - } - base := {{printf "%#v" .Backend.Command}} - if stagedCmd != "" { - base[0] = stagedCmd - } - runner := backend.NewRunner(base{{range .Backend.EnvPassthrough}}, {{printf "%q" .}}{{end}}) + stager := backend.NewStager(ctx, appDir) + runner := backend.NewRunner({{printf "%#v" .Backend.Command}}{{range .Backend.EnvPassthrough}}, {{printf "%q" .}}{{end}}) + runner.SetCommandResolver(stager.Command) {{- else}} runner := backend.NewRunner({{printf "%#v" .Backend.Command}}{{range .Backend.EnvPassthrough}}, {{printf "%q" .}}{{end}}) {{- end}} @@ -132,49 +150,188 @@ func main() { registerHandlers(d, runner, readAppVersion(*manifestPath)) {{- end}} - ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) - defer stop() + serveErr := serve(ctx, *socket, d) +{{- if and (ne .Backend.Type "http") .HasAssets}} + // serve also returns on a listen error: stop staging either way, and let an + // interrupted download remove its partial file before the process exits. + parentGone() + stager.Wait(5 * time.Second) +{{- end}} +{{- if or (eq .Backend.Type "cli") (eq .Backend.Type "hybrid")}} + backend.CloseChildGuard(2 * time.Second) +{{- end}} + if serveErr != nil { + log.Fatalf("{{.BinaryName}}: serve: %v", serveErr) + } +} - if err := serve(ctx, *socket, d); err != nil { - log.Fatalf("{{.BinaryName}}: serve: %v", err) +// watchParent shuts the app down when the process that spawned it (the pilot +// daemon's supervisor) is gone. Linux supervisors also set Pdeathsig, but macOS +// has no parent-death signal and older supervisors set none: without this an app +// orphaned by a daemon crash keeps serving, reparented to launchd/init, and the +// next daemon spawns a second copy beside it. A pid of 1 at start means we were +// launched by init itself (e.g. a container's PID 1): nothing to watch. +func watchParent(ctx context.Context, parentPID int, cancel context.CancelFunc) { + if parentPID <= 1 { + return + } + t := time.NewTicker(250 * time.Millisecond) + defer t.Stop() + for { + if os.Getppid() != parentPID { + log.Printf("{{.BinaryName}}: parent pid %d is gone; shutting down", parentPID) + cancel() + return + } + select { + case <-ctx.Done(): + return + case <-t.C: + } } } -// serve listens on the unix socket and runs ipc.Serve per connection. The -// daemon dials the socket per call, so connections are short-lived; one Serve -// goroutine per accepted conn keeps concurrent callers independent. +// serve listens on the unix socket and serves each connection in its own +// goroutine. The daemon dials the socket per call, so connections are +// short-lived and concurrent callers stay independent. func serve(ctx context.Context, socketPath string, d *ipc.Dispatcher) error { _ = os.Remove(socketPath) // a stale socket from a crashed prior instance would block Listen - ln, err := net.Listen("unix", socketPath) + ln, err := net.ListenUnix("unix", &net.UnixAddr{Name: socketPath, Net: "unix"}) if err != nil { return fmt.Errorf("listen %s: %w", socketPath, err) } - defer ln.Close() - defer os.Remove(socketPath) + // Go unlinks a unix listener's path on Close, whatever that path names by + // then. If a newer instance has replaced app.sock (the daemon respawned the + // app while this copy was orphaned), that would delete the live instance's + // socket and take it offline. Remove only the socket this process bound. + ln.SetUnlinkOnClose(false) + owned, ownedErr := os.Stat(socketPath) + defer func() { + _ = ln.Close() + if ownedErr != nil { + return + } + if cur, err := os.Stat(socketPath); err == nil && os.SameFile(cur, owned) { + _ = os.Remove(socketPath) + } + }() go func() { <-ctx.Done() _ = ln.Close() // unblocks Accept }() + // Live connections, so shutdown can wait for in-flight calls: their ctx is + // cancelled with ours, which makes the backend kill any CLI child it is + // running. Returning before those handlers finish would exit the process + // with the child still alive (reparented to init). + var ( + mu sync.Mutex + conns = map[net.Conn]struct{}{} + wg sync.WaitGroup + ) + log.Printf("{{.BinaryName}}: listening on %s", socketPath) + var backoff time.Duration for { conn, err := ln.Accept() if err != nil { if ctx.Err() != nil { + drain(&mu, conns, &wg) return nil // clean shutdown } + // Running out of descriptors is transient (in-flight calls give + // theirs back as they finish): back off and keep serving, as + // net/http does, instead of exiting and dropping every call. + if errors.Is(err, syscall.EMFILE) || errors.Is(err, syscall.ENFILE) || errors.Is(err, syscall.ECONNABORTED) { + backoff = acceptBackoff(backoff) + log.Printf("{{.BinaryName}}: accept: %v; retrying in %s", err, backoff) + select { + case <-ctx.Done(): + drain(&mu, conns, &wg) + return nil + case <-time.After(backoff): + } + continue + } return fmt.Errorf("accept: %w", err) } + backoff = 0 + mu.Lock() + conns[conn] = struct{}{} + mu.Unlock() + wg.Add(1) go func() { - defer conn.Close() - if err := ipc.Serve(ctx, conn, d); err != nil { - log.Printf("{{.BinaryName}}: conn closed with error: %v", err) - } + defer wg.Done() + defer func() { + mu.Lock() + delete(conns, conn) + mu.Unlock() + }() + serveConn(ctx, conn, d) }() } } +// acceptBackoff doubles the accept retry delay: 5ms first, capped at 1s. +func acceptBackoff(prev time.Duration) time.Duration { + if prev < 5*time.Millisecond { + return 5 * time.Millisecond + } + if next := 2 * prev; next < time.Second { + return next + } + return time.Second +} + +// serveConn serves one connection. The caller hanging up cancels the context +// its call runs under, so an abandoned call (the caller timed out or was +// interrupted) stops its backend request or CLI child at once instead of +// holding that work, its connection and its descriptors until the method's own +// timeout. A goroutine is the connection's only reader: it feeds ipc.Serve +// through a pipe and so sees the hang-up even while a handler is running. +func serveConn(ctx context.Context, conn net.Conn, d *ipc.Dispatcher) { + ctx, cancel := context.WithCancel(ctx) + defer cancel() + defer conn.Close() + pr, pw := io.Pipe() + go func() { + _, err := io.Copy(pw, conn) + cancel() + _ = pw.CloseWithError(err) // nil → io.EOF: a clean end for ipc.Serve + }() + rw := struct { + io.Reader + io.Writer + }{pr, conn} + if err := ipc.Serve(ctx, rw, d); err != nil && ctx.Err() == nil { + log.Printf("{{.BinaryName}}: conn closed with error: %v", err) + } + _ = pr.Close() // Serve returned first (shutdown): release the copier +} + +// shutdownGrace bounds how long shutdown waits for in-flight calls to unwind +// (a cancelled CLI child is SIGKILLed at once; its pipes get backend.killGrace). +const shutdownGrace = 8 * time.Second + +// drain closes every open connection (unblocking idle reads; a handler that is +// still running keeps going until its cancelled work returns) and waits, up to +// shutdownGrace, for all connection goroutines to finish. +func drain(mu *sync.Mutex, conns map[net.Conn]struct{}, wg *sync.WaitGroup) { + mu.Lock() + for c := range conns { + _ = c.Close() + } + mu.Unlock() + done := make(chan struct{}) + go func() { wg.Wait(); close(done) }() + select { + case <-done: + case <-time.After(shutdownGrace): + log.Printf("{{.BinaryName}}: shutdown: in-flight calls still running after %s", shutdownGrace) + } +} + // registerHandlers wires each IPC method to its backend call. Method names MUST // match the manifest's `exposes` list, or the daemon won't broker them. {{- if eq .Backend.Type "http"}} diff --git a/internal/scaffold/templates/stage.go.tmpl b/internal/scaffold/templates/stage.go.tmpl index 4f0451b..641314e 100644 --- a/internal/scaffold/templates/stage.go.tmpl +++ b/internal/scaffold/templates/stage.go.tmpl @@ -2,13 +2,21 @@ // registry. GENERATED by pilot-app (only for cli apps that ship assets); edit // pilot.app.yaml and re-generate. // -// At startup the adapter calls StageAssets($APP). It reads $APP/install.json -// (shipped in the bundle), selects the asset(s) matching this host's os/arch, -// and for each — in ascending install order — fetches it from its R2 URL, -// verifies its sha256 against the (tamper-pinned) install spec, stages it under -// $APP (a single file at exec_path, or a tar.gz extracted in place), and runs -// any install args. The fronted command then execs the staged exec_path instead -// of an assumed-installed binary. +// At startup the adapter starts a Stager for $APP, which runs StageAssetsContext +// in the background. It reads $APP/install.json (shipped in the bundle), selects +// the asset(s) matching this host's os/arch, and for each — in ascending install +// order — fetches it from its R2 URL, verifies its sha256 against the +// (tamper-pinned) install spec, stages it under $APP (a single file at +// exec_path, or a tar.gz extracted in place), and runs any install args. The +// fronted command then execs the staged exec_path instead of an +// assumed-installed binary. +// +// Staging runs beside the IPC server, not before it: a first-start download +// (tens of MB) takes longer than the few seconds the supervisor waits for the +// socket before it gives up on the spawn, and help needs no binary. A call that +// needs the binary waits for staging within its own deadline; a failed attempt +// (e.g. no network on first start) is retried by a later call instead of +// exiting, so the app recovers once the registry is reachable. // // Integrity: each asset's sha256 is checked after download; the whole bundle // tarball is itself sha-pinned in the catalogue, so install.json (and thus the @@ -24,6 +32,7 @@ import ( "errors" "fmt" "io" + "log" "net/http" "os" "os/exec" @@ -32,6 +41,7 @@ import ( "runtime" "sort" "strings" + "sync" "time" ) @@ -43,6 +53,22 @@ const fetchTimeout = 10 * time.Minute // " init" step). const installStepTimeout = 2 * time.Minute +// stagingDirName ($APP/.staged/tmp) is where downloads land before they are +// verified and staged: inside $APP (never the shared TMPDIR), so a download +// interrupted by a stop, a crash or the daemon's death leaves nothing behind +// that the next start does not sweep, and an uninstall reclaims it. A 40-60 MB +// CLI tarball per interrupted first start used to pile up in $TMPDIR as +// pilot-asset-*. (Same location as fix/cli-server-teardown-parent-watch.) +var stagingDirName = filepath.Join(".staged", "tmp") + +// Retry pacing for a failed staging attempt: the first retry may start +// retryMin after the failure, doubling up to retryMax, so a call loop against +// an unreachable registry or a bad artifact does not refetch on every call. +const ( + retryMin = 2 * time.Second + retryMax = 5 * time.Minute +) + // maxAssetBytes caps a single download / extracted file so a malicious or // corrupt artifact can't fill the disk. const maxAssetBytes = 2 << 30 // 2 GiB @@ -67,11 +93,111 @@ type installAsset struct { Args []string `json:"args"` } +// Stager stages the app's assets in the background and hands the staged +// command path to calls that need it (Runner.SetCommandResolver). +type Stager struct { + appDir string + ctx context.Context // the adapter's lifetime: cancelling it aborts staging + + mu sync.Mutex + cur *stageRun // the attempt in progress, or the last one + backoff time.Duration + notAfter time.Time // no new attempt before this (after a failure) +} + +type stageRun struct { + done chan struct{} + cmd string + err error +} + +// NewStager starts staging appDir's assets in the background. Cancelling ctx +// (the adapter is stopping) aborts the attempt in progress, which removes its +// partial download. +func NewStager(ctx context.Context, appDir string) *Stager { + s := &Stager{appDir: appDir, ctx: ctx} + s.mu.Lock() + s.startLocked() + s.mu.Unlock() + return s +} + +func (s *Stager) startLocked() *stageRun { + r := &stageRun{done: make(chan struct{})} + s.cur = r + go func() { + defer close(r.done) + start := time.Now() + r.cmd, r.err = StageAssetsContext(s.ctx, s.appDir) + s.mu.Lock() + defer s.mu.Unlock() + switch { + case r.err == nil: + s.backoff = 0 + if d := time.Since(start); d > time.Second { + log.Printf("stage: assets installed in %s", d.Round(time.Millisecond)) + } + case s.ctx.Err() != nil: + // Stopping: not a failure to report. + default: + s.backoff = min(max(2*s.backoff, retryMin), retryMax) + s.notAfter = time.Now().Add(s.backoff) + log.Printf("stage: install assets: %v (a call retries after %s)", r.err, s.backoff) + } + }() + return r +} + +// Command returns the staged command path, waiting — at most until ctx is done +// — for an attempt in progress. After a failed attempt, a call past the retry +// delay starts a new one; before it, the call fails with the last error. +func (s *Stager) Command(ctx context.Context) (string, error) { + s.mu.Lock() + r := s.cur + select { + case <-r.done: + if r.err != nil && s.ctx.Err() == nil && !time.Now().Before(s.notAfter) { + r = s.startLocked() + } + default: + } + s.mu.Unlock() + select { + case <-r.done: + if r.err != nil { + return "", fmt.Errorf("install assets: %w", r.err) + } + return r.cmd, nil + case <-ctx.Done(): + return "", fmt.Errorf("install assets: still installing: %w", ctx.Err()) + } +} + +// Wait waits, up to timeout, for the attempt in progress to return. main calls +// it on shutdown, after cancelling the Stager's ctx, so that an interrupted +// download has removed its partial file before the process exits. +func (s *Stager) Wait(timeout time.Duration) { + s.mu.Lock() + r := s.cur + s.mu.Unlock() + select { + case <-r.done: + case <-time.After(timeout): + } +} + // StageAssets materializes the registry assets for this host and returns the // absolute path of the staged command binary (the asset whose exec_path matches // install.json's "command"). When there is no install.json the app ships no // assets, so it returns ("", nil) and the caller keeps the command as-is. func StageAssets(appDir string) (string, error) { + return StageAssetsContext(context.Background(), appDir) +} + +// StageAssetsContext is StageAssets bounded by ctx: cancelling it (the adapter +// is being stopped) aborts the download / extraction and removes the partial +// download. +func StageAssetsContext(ctx context.Context, appDir string) (string, error) { raw, err := os.ReadFile(filepath.Join(appDir, "install.json")) if errors.Is(err, os.ErrNotExist) { return "", nil @@ -95,9 +221,17 @@ func StageAssets(appDir string) (string, error) { } sort.SliceStable(host, func(i, j int) bool { return host[i].Order < host[j].Order }) + // Whatever an interrupted earlier start left in the staging dir is garbage. + staging := filepath.Join(appDir, stagingDirName) + _ = os.RemoveAll(staging) + defer os.RemoveAll(staging) + // Staging is a one-off: do not keep its registry connection open (idle in + // the keep-alive pool) for the adapter's lifetime. + defer http.DefaultClient.CloseIdleConnections() + var cmdPath string for _, a := range host { - execPath, err := stageOne(appDir, a) + execPath, err := stageOne(ctx, appDir, staging, a) if err != nil { return "", err } @@ -105,7 +239,7 @@ func StageAssets(appDir string) (string, error) { cmdPath = execPath } if len(a.Args) > 0 { - if err := runInstallStep(execPath, a.Args); err != nil { + if err := runInstallStep(ctx, execPath, a.Args); err != nil { return "", fmt.Errorf("stage: install step for %q failed: %w", a.ExecPath, err) } } @@ -118,7 +252,7 @@ func StageAssets(appDir string) (string, error) { // tar.gz asset is extracted in place and exec_path names a file inside the // extracted tree. Staging is idempotent: a sha-stamped marker skips re-work on // re-spawn. -func stageOne(appDir string, a installAsset) (string, error) { +func stageOne(ctx context.Context, appDir, staging string, a installAsset) (string, error) { execAbs := filepath.Join(appDir, filepath.FromSlash(a.ExecPath)) marker := filepath.Join(appDir, ".staged", a.SHA256) if _, err := os.Stat(marker); err == nil { @@ -127,7 +261,7 @@ func stageOne(appDir string, a installAsset) (string, error) { } } - body, err := download(a.URL, a.SHA256) + body, err := download(ctx, staging, a.URL, a.SHA256) if err != nil { return "", err } @@ -135,7 +269,7 @@ func stageOne(appDir string, a installAsset) (string, error) { switch a.Unpack { case "tar.gz": - if err := extractTarGz(body, appDir); err != nil { + if err := extractTarGz(ctx, body, appDir); err != nil { return "", fmt.Errorf("stage: extract %q: %w", a.ExecPath, err) } default: @@ -155,8 +289,8 @@ func stageOne(appDir string, a installAsset) (string, error) { // download fetches url to a temp file, verifying its sha256 streams-as-it-goes. // Returns the temp file path (caller removes it). A mismatch is fatal so a // tampered or wrong artifact is never installed. -func download(url, wantSHA string) (string, error) { - ctx, cancel := context.WithTimeout(context.Background(), fetchTimeout) +func download(ctx context.Context, staging, url, wantSHA string) (string, error) { + ctx, cancel := context.WithTimeout(ctx, fetchTimeout) defer cancel() req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) if err != nil { @@ -171,7 +305,10 @@ func download(url, wantSHA string) (string, error) { return "", fmt.Errorf("stage: fetch %s: HTTP %d", url, resp.StatusCode) } - tmp, err := os.CreateTemp("", "pilot-asset-*") + if err := os.MkdirAll(staging, 0o700); err != nil { + return "", fmt.Errorf("stage: staging dir: %w", err) + } + tmp, err := os.CreateTemp(staging, "pilot-asset-*") if err != nil { return "", fmt.Errorf("stage: temp file: %w", err) } @@ -217,16 +354,19 @@ func installFile(src, dest, role string) error { // is present on every linux/darwin host and handles them. Before extracting, // every entry name is scanned and any absolute path or "../" traversal is // rejected (zip-slip defence), since tar's own stripping is not relied upon. -func extractTarGz(archive, dir string) error { - if err := assertSafeArchive(archive); err != nil { +func extractTarGz(ctx context.Context, archive, dir string) error { + if err := assertSafeArchive(ctx, archive); err != nil { return err } if err := os.MkdirAll(dir, 0o755); err != nil { return err } - ctx, cancel := context.WithTimeout(context.Background(), installStepTimeout) + ctx, cancel := context.WithTimeout(ctx, installStepTimeout) defer cancel() cmd := exec.CommandContext(ctx, "tar", "-xzf", archive, "-C", dir) + cmd.SysProcAttr = childProcAttr() // dies with the adapter (Linux) + unlock := lockForChild() + defer unlock() if out, err := cmd.CombinedOutput(); err != nil { return fmt.Errorf("tar -xzf: %w: %s", err, strings.TrimSpace(string(out))) } @@ -235,11 +375,14 @@ func extractTarGz(archive, dir string) error { // assertSafeArchive lists the archive (tar -tzf) and rejects any member that is // an absolute path or escapes the extraction root via "..". -func assertSafeArchive(archive string) error { - ctx, cancel := context.WithTimeout(context.Background(), installStepTimeout) +func assertSafeArchive(ctx context.Context, archive string) error { + ctx, cancel := context.WithTimeout(ctx, installStepTimeout) defer cancel() cmd := exec.CommandContext(ctx, "tar", "-tzf", archive) + cmd.SysProcAttr = childProcAttr() + unlock := lockForChild() out, err := cmd.Output() + unlock() if err != nil { return fmt.Errorf("tar -tzf (list): %w", err) } @@ -282,12 +425,15 @@ func copyFile(src, dest string, mode os.FileMode) error { // runInstallStep runs a one-time post-stage command (the publisher's optional // install args), e.g. " init". A non-zero exit fails the install so a // broken setup never silently serves. -func runInstallStep(execPath string, args []string) error { - ctx, cancel := context.WithTimeout(context.Background(), installStepTimeout) +func runInstallStep(ctx context.Context, execPath string, args []string) error { + ctx, cancel := context.WithTimeout(ctx, installStepTimeout) defer cancel() cmd := exec.CommandContext(ctx, execPath, args...) cmd.Env = os.Environ() + cmd.SysProcAttr = childProcAttr() + unlock := lockForChild() out, err := cmd.CombinedOutput() + unlock() if err != nil { return fmt.Errorf("%s %s: %w: %s", execPath, strings.Join(args, " "), err, strings.TrimSpace(string(out))) } diff --git a/internal/scaffold/zz_adapter_lifecycle_e2e_test.go b/internal/scaffold/zz_adapter_lifecycle_e2e_test.go new file mode 100644 index 0000000..4e3f6b8 --- /dev/null +++ b/internal/scaffold/zz_adapter_lifecycle_e2e_test.go @@ -0,0 +1,348 @@ +//go:build !windows + +package scaffold + +import ( + "bufio" + "bytes" + "encoding/json" + "fmt" + "io" + "net" + "net/http" + "net/http/httptest" + "os" + "os/exec" + "path/filepath" + "strconv" + "strings" + "sync" + "sync/atomic" + "syscall" + "testing" + "time" + + "github.com/pilot-protocol/app-store/pkg/ipc" +) + +// lifecycleSpec is a byo http app with one fast and one slow method. The slow +// backend route answers only when its request is cancelled, so a test can see +// whether the adapter gives up on a call its caller has abandoned. +const lifecycleSpec = ` +id: io.pilot.lifecyclex +app_version: 0.1.0 +description: "App exercising the adapter process lifecycle." +namespace: lifecyclex +backend: + base_url: https://placeholder.invalid +methods: + - name: lifecyclex.fast + summary: "Answers at once." + http: { verb: GET, path: "/fast" } + - name: lifecyclex.slow + summary: "Answers only when cancelled." + http: { verb: POST, path: "/slow" } +` + +// lifecycleAdapter is a built lifecyclex adapter plus the backend it talks to. +type lifecycleAdapter struct { + bin string + backend *httptest.Server + slowOpen atomic.Int64 // /slow requests still waiting on the backend + cancelled atomic.Int64 // /slow requests the adapter abandoned +} + +func buildLifecycleAdapter(t *testing.T) *lifecycleAdapter { + t.Helper() + if testing.Short() { + t.Skip("builds and runs a real adapter binary; skipped under -short") + } + if _, err := exec.LookPath("go"); err != nil { + t.Skip("go toolchain not available") + } + la := &lifecycleAdapter{} + la.backend = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/slow" { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"ok":true}`)) + return + } + // Drain the body first: net/http only watches an HTTP/1 connection + // for the client going away once the request body has been read. + _, _ = io.Copy(io.Discard, r.Body) + la.slowOpen.Add(1) + defer la.slowOpen.Add(-1) + select { + case <-r.Context().Done(): + la.cancelled.Add(1) + case <-time.After(30 * time.Second): + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"late":true}`)) + } + })) + t.Cleanup(la.backend.Close) + + root := t.TempDir() + cfg := parseSpec(t, lifecycleSpec) + proj := filepath.Join(root, "proj") + if _, err := Generate(cfg, proj); err != nil { + t.Fatalf("generate: %v", err) + } + if sum, err := os.ReadFile(filepath.Join("..", "..", "go.sum")); err == nil { + _ = os.WriteFile(filepath.Join(proj, "go.sum"), sum, 0o644) + } + la.bin = filepath.Join(root, "adapter") + build := exec.Command("go", "build", "-o", la.bin, "./cmd/"+cfg.BinaryName) + build.Dir = proj + build.Env = append(os.Environ(), "GOFLAGS=-mod=mod") + if out, err := build.CombinedOutput(); err != nil { + t.Fatalf("build adapter: %v\n%s", err, out) + } + return la +} + +// sockPath returns a socket path short enough for sun_path (~104 on darwin). +func sockPath(t *testing.T) string { + t.Helper() + dir, err := os.MkdirTemp("", "lcx") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.RemoveAll(dir) }) + return filepath.Join(dir, "a.sock") +} + +func (la *lifecycleAdapter) env() []string { + return append(os.Environ(), "LIFECYCLEX_BACKEND_URL="+la.backend.URL) +} + +// start runs the adapter as a direct child of the test process. +func (la *lifecycleAdapter) start(t *testing.T, sock string) *exec.Cmd { + t.Helper() + cmd := exec.Command(la.bin, "--socket", sock) + cmd.Env = la.env() + cmd.Stderr = os.Stderr + if err := cmd.Start(); err != nil { + t.Fatalf("start adapter: %v", err) + } + t.Cleanup(func() { _ = cmd.Process.Kill(); _, _ = cmd.Process.Wait() }) + return cmd +} + +func waitSocket(t *testing.T, sock string, notInode uint64) { + t.Helper() + for deadline := time.Now().Add(10 * time.Second); time.Now().Before(deadline); time.Sleep(20 * time.Millisecond) { + if fi, err := os.Stat(sock); err == nil && inode(fi) != notInode { + return + } + } + t.Fatalf("socket %s never appeared", sock) +} + +func inode(fi os.FileInfo) uint64 { + if st, ok := fi.Sys().(*syscall.Stat_t); ok { + return uint64(st.Ino) + } + return 0 +} + +func callFast(sock string) error { + conn, err := net.DialTimeout("unix", sock, 3*time.Second) + if err != nil { + return err + } + defer conn.Close() + _ = conn.SetDeadline(time.Now().Add(5 * time.Second)) + var out json.RawMessage + return ipc.Call(conn, "lifecyclex.fast", json.RawMessage(`{}`), &out) +} + +// sendSlow opens a connection and sends one lifecyclex.slow request without +// waiting for the reply; the caller decides when to hang up. +func sendSlow(t *testing.T, sock string, i int) net.Conn { + t.Helper() + conn, err := trySendSlow(sock, i) + if err != nil { + t.Fatal(err) + } + return conn +} + +func trySendSlow(sock string, i int) (net.Conn, error) { + conn, err := net.DialTimeout("unix", sock, 3*time.Second) + if err != nil { + return nil, fmt.Errorf("dial %d: %w", i, err) + } + req := &ipc.Envelope{Type: ipc.EnvReq, ReqID: fmt.Sprintf("r%d", i), Method: "lifecyclex.slow", Payload: json.RawMessage(`{}`)} + if err := ipc.WriteFrame(conn, req); err != nil { + _ = conn.Close() + return nil, fmt.Errorf("write %d: %w", i, err) + } + return conn, nil +} + +// syncBuffer is a bytes.Buffer safe to read while exec's copier writes to it. +type syncBuffer struct { + mu sync.Mutex + buf bytes.Buffer +} + +func (b *syncBuffer) Write(p []byte) (int, error) { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.Write(p) +} + +func (b *syncBuffer) String() string { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.String() +} + +func alive(pid int) bool { return syscall.Kill(pid, 0) == nil } + +func eventually(d time.Duration, cond func() bool) bool { + for deadline := time.Now().Add(d); time.Now().Before(deadline); time.Sleep(25 * time.Millisecond) { + if cond() { + return true + } + } + return cond() +} + +// TestGeneratedAdapterLifecycleE2E runs a real generated adapter binary and +// checks it never outlives, or disturbs, what the daemon expects: +// +// - it exits when the process that spawned it dies (macOS has no Pdeathsig); +// - an old instance exiting leaves a newer instance's socket in place; +// - a caller hanging up cancels the call's backend request; +// - running out of descriptors is survived, not fatal. +func TestGeneratedAdapterLifecycleE2E(t *testing.T) { + la := buildLifecycleAdapter(t) + + t.Run("exits when its parent dies", func(t *testing.T) { + sock := sockPath(t) + // A shell stands in for the daemon: it starts the adapter, prints its + // pid, and is then SIGKILLed the way a crashed daemon disappears. + parent := exec.Command("sh", "-c", `"$0" --socket "$1" & echo $!; exec sleep 60`, la.bin, sock) + parent.Env = la.env() + parent.Stderr = os.Stderr + out, err := parent.StdoutPipe() + if err != nil { + t.Fatal(err) + } + if err := parent.Start(); err != nil { + t.Fatal(err) + } + line, err := bufio.NewReader(out).ReadString('\n') + if err != nil { + t.Fatalf("read adapter pid: %v", err) + } + pid, err := strconv.Atoi(strings.TrimSpace(line)) + if err != nil { + t.Fatalf("adapter pid %q: %v", line, err) + } + t.Cleanup(func() { _ = syscall.Kill(pid, syscall.SIGKILL) }) + waitSocket(t, sock, 0) + if err := callFast(sock); err != nil { + t.Fatalf("call before parent death: %v", err) + } + + _ = parent.Process.Kill() + _, _ = parent.Process.Wait() + if !eventually(3*time.Second, func() bool { return !alive(pid) }) { + t.Fatalf("adapter pid %d still running 3s after its parent was killed (orphaned)", pid) + } + if _, err := os.Stat(sock); !os.IsNotExist(err) { + t.Errorf("socket %s left behind after the orphaned adapter exited (stat err=%v)", sock, err) + } + }) + + t.Run("old instance exiting keeps the newer socket", func(t *testing.T) { + sock := sockPath(t) + old := la.start(t, sock) + waitSocket(t, sock, 0) + fi, err := os.Stat(sock) + if err != nil { + t.Fatal(err) + } + // A respawn next to an orphan: the new instance replaces app.sock. + la.start(t, sock) + waitSocket(t, sock, inode(fi)) + + _ = old.Process.Signal(syscall.SIGTERM) + _, _ = old.Process.Wait() + if _, err := os.Stat(sock); err != nil { + t.Fatalf("old instance's exit removed the new instance's socket: %v", err) + } + if err := callFast(sock); err != nil { + t.Errorf("new instance unreachable after the old one exited: %v", err) + } + }) + + t.Run("caller hang-up cancels the backend request", func(t *testing.T) { + sock := sockPath(t) + la.start(t, sock) + waitSocket(t, sock, 0) + before := la.cancelled.Load() + const n = 5 + for i := 0; i < n; i++ { + conn := sendSlow(t, sock, i) + if !eventually(3*time.Second, func() bool { return la.slowOpen.Load() > 0 }) { + t.Fatalf("call %d never reached the backend", i) + } + _ = conn.Close() // the caller gives up + if !eventually(3*time.Second, func() bool { return la.cancelled.Load() == before+int64(i)+1 }) { + t.Fatalf("call %d: backend request still open 3s after its caller hung up (cancelled=%d)", i, la.cancelled.Load()-before) + } + } + if err := callFast(sock); err != nil { + t.Errorf("adapter unhealthy after abandoned calls: %v", err) + } + }) + + t.Run("survives running out of descriptors", func(t *testing.T) { + sock := sockPath(t) + // Few descriptors, so a handful of waiting callers exhaust them. + cmd := exec.Command("sh", "-c", `ulimit -n 32 && exec "$0" --socket "$1"`, la.bin, sock) + cmd.Env = la.env() + var logs syncBuffer + cmd.Stderr = io.MultiWriter(os.Stderr, &logs) + if err := cmd.Start(); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = cmd.Process.Kill(); _, _ = cmd.Process.Wait() }) + exited := make(chan error, 1) + go func() { exited <- cmd.Wait() }() + waitSocket(t, sock, 0) + + var conns []net.Conn + for i := 0; i < 40; i++ { // two descriptors each: more than 32 in total + c, err := trySendSlow(sock, i) + if err != nil { + // A caller may be turned away here: when accept(2) fails with + // EMFILE, macOS drops the pending connection (the peer sees + // EPIPE/ECONNRESET) where Linux leaves it queued. Either way + // the adapter itself must keep running. + t.Logf("caller %d turned away while the adapter is out of descriptors: %v", i, err) + continue + } + conns = append(conns, c) + } + time.Sleep(1500 * time.Millisecond) + select { + case err := <-exited: + t.Fatalf("adapter exited while out of descriptors: %v", err) + default: + } + if !strings.Contains(logs.String(), "too many open files") { + t.Fatalf("adapter never ran out of descriptors, so this proves nothing; its log:\n%s", logs.String()) + } + for _, c := range conns { + _ = c.Close() + } + if !eventually(5*time.Second, func() bool { return callFast(sock) == nil }) { + t.Fatalf("adapter did not recover after callers hung up: %v", callFast(sock)) + } + }) +} diff --git a/internal/scaffold/zz_cli_inflight_server_e2e_test.go b/internal/scaffold/zz_cli_inflight_server_e2e_test.go new file mode 100644 index 0000000..f1049c2 --- /dev/null +++ b/internal/scaffold/zz_cli_inflight_server_e2e_test.go @@ -0,0 +1,409 @@ +//go:build !windows + +package scaffold + +import ( + "bufio" + "encoding/json" + "errors" + "net" + "os" + "os/exec" + "path/filepath" + "runtime" + "strconv" + "strings" + "syscall" + "testing" + "time" + + "github.com/pilot-protocol/app-store/pkg/ipc" +) + +// miniSupEnv switches the test binary into a stand-in for the pilot daemon's +// supervisor (TestHelperMiniSupervisor): it spawns the adapter named in the +// variable with the supervisor's process attributes, prints the adapter's pid +// and then waits to be killed. +const miniSupEnv = "SCAFFOLD_TEST_MINI_SUPERVISOR" + +// TestHelperMiniSupervisor is not a test: it is the helper process for +// TestCLIInflightServerE2E and skips unless miniSupEnv is set. +func TestHelperMiniSupervisor(t *testing.T) { + spec := os.Getenv(miniSupEnv) + if spec == "" { + t.Skip("helper process for TestCLIInflightServerE2E") + } + var argv []string + if err := json.Unmarshal([]byte(spec), &argv); err != nil || len(argv) == 0 { + t.Fatalf("bad %s: %q", miniSupEnv, spec) + } + // app-store spawn(): own process group everywhere, plus Pdeathsig=SIGKILL + // on Linux from a thread that stays alive as long as this process does. + runtime.LockOSThread() + cmd := exec.Command(argv[0], argv[1:]...) + cmd.Stderr = os.Stderr + cmd.SysProcAttr = superviseAttr() + if err := cmd.Start(); err != nil { + t.Fatalf("spawn adapter: %v", err) + } + os.Stdout.WriteString(strconv.Itoa(cmd.Process.Pid) + "\n") + time.Sleep(time.Hour) // until the test SIGKILLs us, as a crashed daemon disappears +} + +// TestCLIInflightServerE2E: the io.pilot.otto shape. A passthrough call starts +// a long-running server (`otto.exec ["mcp","serve-http","--port",P]`) that is +// still running when the app goes away. In otto 0.20.0 the child kept its port +// open, reparented to init/launchd, after the daemon died (Linux: the +// adapter's Pdeathsig killed the adapter, not the child; macOS: the adapter +// itself was orphaned and kept serving) and, on macOS, after the adapter was +// SIGKILLed (the supervisor's stop path). The child guard (childguard.go) +// covers the SIGKILL case on every OS; these subtests pin its lifecycle too. +func TestCLIInflightServerE2E(t *testing.T) { + if testing.Short() { + t.Skip("builds and runs a real adapter binary; skipped under -short") + } + if _, err := exec.LookPath("go"); err != nil { + t.Skip("go toolchain not available") + } + root := t.TempDir() + tool := filepath.Join(root, "srvtool") + // `serve ` records its pid and then runs until killed, without + // writing anything (so a closed pipe cannot end it early). + script := "#!/bin/sh\n" + + "case \"$1\" in\n" + + " serve) echo $$ > \"$2\"; exec sleep 300;;\n" + + " quick) echo '{\"ok\":true}';;\n" + + " *) echo '{}';;\n" + + "esac\n" + if err := os.WriteFile(tool, []byte(script), 0o755); err != nil { + t.Fatal(err) + } + cfg := parseSpec(t, ` +id: io.pilot.srvtool +app_version: 0.1.0 +description: "Fronts srvtool." +namespace: srvtool +backend: + type: cli + command: ["`+tool+`"] +methods: + - name: srvtool.exec + summary: "Passthrough." + timeout: 120s + cli: {passthrough: true} +`) + proj := filepath.Join(root, "proj") + if _, err := Generate(cfg, proj); err != nil { + t.Fatalf("generate: %v", err) + } + if sum, err := os.ReadFile(filepath.Join("..", "..", "go.sum")); err == nil { + _ = os.WriteFile(filepath.Join(proj, "go.sum"), sum, 0o644) + } + bin := filepath.Join(root, "adapter") + build := exec.Command("go", "build", "-o", bin, "./cmd/"+cfg.BinaryName) + build.Dir = proj + build.Env = append(os.Environ(), "GOFLAGS=-mod=mod") + if out, err := build.CombinedOutput(); err != nil { + t.Fatalf("build adapter: %v\n%s", err, out) + } + + shortDir := func(t *testing.T) string { + t.Helper() + dir, err := os.MkdirTemp("", "cis") // short: sun_path is ~104 bytes on darwin + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.RemoveAll(dir) }) + return dir + } + waitSock := func(t *testing.T, sock string) { + t.Helper() + for deadline := time.Now().Add(10 * time.Second); time.Now().Before(deadline); time.Sleep(20 * time.Millisecond) { + if _, err := os.Stat(sock); err == nil { + return + } + } + t.Fatal("adapter socket never appeared") + } + // startServe issues the long-running passthrough call without waiting for + // it and returns the server child's pid. + startServe := func(t *testing.T, sock, dir string) int { + t.Helper() + pidfile := filepath.Join(dir, "child.pid") + go func() { + conn, err := net.DialTimeout("unix", sock, 3*time.Second) + if err != nil { + return + } + defer conn.Close() + var out json.RawMessage + _ = ipc.Call(conn, "srvtool.exec", map[string]any{"args": []string{"serve", pidfile}}, &out) + }() + for deadline := time.Now().Add(10 * time.Second); time.Now().Before(deadline); time.Sleep(20 * time.Millisecond) { + if b, err := os.ReadFile(pidfile); err == nil { + if pid, err := strconv.Atoi(strings.TrimSpace(string(b))); err == nil && pid > 0 { + t.Cleanup(func() { _ = syscall.Kill(pid, syscall.SIGKILL) }) + return pid + } + } + } + t.Fatal("server child never started") + return 0 + } + gone := func(pid int, within time.Duration) bool { + for deadline := time.Now().Add(within); ; time.Sleep(20 * time.Millisecond) { + if !procAlive(pid) { + return true + } + if time.Now().After(deadline) { + return false + } + } + } + + t.Run("daemon dies with a server call in flight", func(t *testing.T) { + dir := shortDir(t) + sock := filepath.Join(dir, "app.sock") + argv, _ := json.Marshal([]string{bin, "--socket", sock, "--manifest", filepath.Join(proj, "manifest.json")}) + sup := exec.Command(os.Args[0], "-test.run=^TestHelperMiniSupervisor$", "-test.v=false") + sup.Env = append(os.Environ(), miniSupEnv+"="+string(argv)) + sup.Stderr = os.Stderr + out, err := sup.StdoutPipe() + if err != nil { + t.Fatal(err) + } + if err := sup.Start(); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = sup.Process.Kill(); _, _ = sup.Process.Wait() }) + line, err := bufio.NewReader(out).ReadString('\n') + if err != nil { + t.Fatalf("read adapter pid from mini supervisor: %v", err) + } + apid, err := strconv.Atoi(strings.TrimSpace(line)) + if err != nil { + t.Fatalf("adapter pid %q: %v", line, err) + } + t.Cleanup(func() { _ = syscall.Kill(-apid, syscall.SIGKILL) }) + waitSock(t, sock) + child := startServe(t, sock, dir) + + _ = sup.Process.Kill() // the daemon crashes + _, _ = sup.Process.Wait() + // Linux: the adapter's Pdeathsig kills it and the child's own + // Pdeathsig kills the child. macOS: the adapter sees it was reparented, + // shuts down and reaps the child before exiting. + if !gone(apid, 5*time.Second) { + t.Errorf("adapter pid %d outlived its supervisor", apid) + } + if !gone(child, 5*time.Second) { + t.Errorf("server child pid %d outlived the dead supervisor and adapter (reparented to init/launchd)", child) + } + }) + + // startAdapter runs the adapter the way the supervisor does: leader of its + // own process group. + startAdapter := func(t *testing.T, dir string) (*exec.Cmd, string) { + t.Helper() + sock := filepath.Join(dir, "app.sock") + a := exec.Command(bin, "--socket", sock, "--manifest", filepath.Join(proj, "manifest.json")) + a.Stderr = os.Stderr + a.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} + if err := a.Start(); err != nil { + t.Fatal(err) + } + apid := a.Process.Pid + t.Cleanup(func() { _ = syscall.Kill(-apid, syscall.SIGKILL); _, _ = a.Process.Wait() }) + waitSock(t, sock) + return a, sock + } + + t.Run("adapter SIGKILLed with a server call in flight", func(t *testing.T) { + // The app-store supervisor's stop path through v1.0.3: SIGKILL to the + // adapter's pid only. Linux: the child's Pdeathsig. Everywhere + // (macOS has no Pdeathsig): the child guard takes the group down. + dir := shortDir(t) + a, sock := startAdapter(t, dir) + apid := a.Process.Pid + child := startServe(t, sock, dir) + _ = syscall.Kill(apid, syscall.SIGKILL) + _, _ = a.Process.Wait() + if !gone(child, 3*time.Second) { + t.Errorf("server child pid %d outlived the SIGKILLed adapter", child) + } + if left := groupLeft(apid, 3*time.Second); len(left) > 0 { + t.Errorf("adapter's process group %d not empty after it was SIGKILLed: %v", apid, left) + } + }) + + t.Run("clean stop leaves an empty process group", func(t *testing.T) { + dir := shortDir(t) + a, sock := startAdapter(t, dir) + apid := a.Process.Pid + pidfile := filepath.Join(dir, "quick.pid") + conn, err := net.DialTimeout("unix", sock, 3*time.Second) + if err != nil { + t.Fatal(err) + } + var out json.RawMessage + if err := ipc.Call(conn, "srvtool.exec", map[string]any{"args": []string{"quick", pidfile}}, &out); err != nil { + t.Fatalf("quick call: %v", err) + } + conn.Close() + // The guard outlives the call briefly (guardIdle), so the group is not + // empty yet; a SIGTERM now must still leave nothing behind. + if left := groupMembers(apid); len(left) < 2 { + t.Fatalf("expected the adapter and its idle child guard in group %d, got %v", apid, left) + } + _ = a.Process.Signal(syscall.SIGTERM) + _, _ = a.Process.Wait() + // CloseChildGuard released and reaped the guard before the adapter exited. + if left := groupMembers(apid); len(left) > 0 { + t.Errorf("process group %d not empty after a clean SIGTERM: %v", apid, left) + } + }) + + t.Run("idle adapter retires its guard", func(t *testing.T) { + dir := shortDir(t) + a, sock := startAdapter(t, dir) + apid := a.Process.Pid + conn, err := net.DialTimeout("unix", sock, 3*time.Second) + if err != nil { + t.Fatal(err) + } + var out json.RawMessage + if err := ipc.Call(conn, "srvtool.exec", map[string]any{"args": []string{"quick"}}, &out); err != nil { + t.Fatalf("quick call: %v", err) + } + conn.Close() + deadline := time.Now().Add(6 * time.Second) + for len(groupMembers(apid)) > 1 && time.Now().Before(deadline) { + time.Sleep(50 * time.Millisecond) + } + if left := groupMembers(apid); len(left) != 1 || left[0] != apid { + t.Fatalf("idle adapter group %d = %v, want only the adapter (guard retired)", apid, left) + } + if !procAlive(apid) { + t.Fatal("adapter died when its guard retired") + } + }) + + t.Run("adapter keeps serving when its guard dies", func(t *testing.T) { + dir := shortDir(t) + a, sock := startAdapter(t, dir) + apid := a.Process.Pid + child := startServe(t, sock, dir) + var guardPID int + for _, pid := range groupMembers(apid) { + if pid != apid && pid != child { + guardPID = pid + } + } + if guardPID == 0 { + t.Fatalf("no child guard in group %d: %v", apid, groupMembers(apid)) + } + _ = syscall.Kill(guardPID, syscall.SIGKILL) + if !gone(guardPID, 3*time.Second) { + t.Fatalf("guard %d did not die", guardPID) + } + if !procAlive(apid) || !procAlive(child) { + t.Fatalf("killing the guard took down adapter (alive=%v) or child (alive=%v)", procAlive(apid), procAlive(child)) + } + // The next call starts a new guard, which still covers the old child. + second := startServe(t, sock, shortDir(t)) + _ = syscall.Kill(apid, syscall.SIGKILL) + _, _ = a.Process.Wait() + for _, pid := range []int{child, second} { + if !gone(pid, 3*time.Second) { + t.Errorf("child %d outlived the SIGKILLed adapter after its first guard died", pid) + } + } + }) + + t.Run("supervisor group stop reaches the server child", func(t *testing.T) { + dir := shortDir(t) + sock := filepath.Join(dir, "app.sock") + a := exec.Command(bin, "--socket", sock, "--manifest", filepath.Join(proj, "manifest.json")) + a.Stderr = os.Stderr + a.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} // as the supervisor spawns it + if err := a.Start(); err != nil { + t.Fatal(err) + } + apid := a.Process.Pid + t.Cleanup(func() { _ = syscall.Kill(-apid, syscall.SIGKILL); _, _ = a.Process.Wait() }) + waitSock(t, sock) + child := startServe(t, sock, dir) + + // The child must stay in the adapter's process group: on macOS (no + // Pdeathsig) a group-wide signal from the supervisor is what reaches it + // when the adapter is SIGKILLed. + if pg, err := syscall.Getpgid(child); err != nil || pg != apid { + t.Fatalf("server child pgid = %d (err %v), want the adapter's group %d", pg, err, apid) + } + if err := syscall.Kill(-apid, syscall.SIGKILL); err != nil { + t.Fatalf("kill group: %v", err) + } + _, _ = a.Process.Wait() + if !gone(child, 3*time.Second) { + t.Errorf("server child pid %d survived a SIGKILL of the adapter's process group", child) + } + }) +} + +// procAlive reports whether pid is a live (non-zombie) process. A zombie is +// dead for this purpose: in a container nothing may ever reap it. +func procAlive(pid int) bool { + if err := syscall.Kill(pid, 0); errors.Is(err, syscall.ESRCH) { + return false + } + if b, err := os.ReadFile("/proc/" + strconv.Itoa(pid) + "/stat"); err == nil { + s := string(b) + if i := strings.LastIndexByte(s, ')'); i >= 0 && i+2 < len(s) && (s[i+2] == 'Z' || s[i+2] == 'X') { + return false + } + return true + } + if runtime.GOOS == "darwin" { + out, err := exec.Command("ps", "-o", "stat=", "-p", strconv.Itoa(pid)).Output() + if err != nil { + return false + } + return !strings.HasPrefix(strings.TrimSpace(string(out)), "Z") + } + return true +} + +// groupLeft waits up to d for process group pgid to empty and returns what is +// still in it. +func groupLeft(pgid int, d time.Duration) []int { + deadline := time.Now().Add(d) + for { + left := groupMembers(pgid) + if len(left) == 0 || time.Now().After(deadline) { + return left + } + time.Sleep(20 * time.Millisecond) + } +} + +// groupMembers lists the live (non-zombie) members of process group pgid. +func groupMembers(pgid int) []int { + // ps is on every darwin host and in the Linux test images (procps). + out, err := exec.Command("ps", "-A", "-o", "pid=,pgid=,stat=").Output() + if err != nil { + return nil + } + var pids []int + for _, ln := range strings.Split(string(out), "\n") { + f := strings.Fields(ln) + if len(f) < 3 { + continue + } + pid, _ := strconv.Atoi(f[0]) + pg, _ := strconv.Atoi(f[1]) + if pg == pgid && !strings.HasPrefix(f[2], "Z") { + pids = append(pids, pid) + } + } + return pids +} diff --git a/internal/scaffold/zz_cli_inflight_server_linux_test.go b/internal/scaffold/zz_cli_inflight_server_linux_test.go new file mode 100644 index 0000000..74c58fe --- /dev/null +++ b/internal/scaffold/zz_cli_inflight_server_linux_test.go @@ -0,0 +1,11 @@ +//go:build linux + +package scaffold + +import "syscall" + +// superviseAttr mirrors the app-store supervisor's spawn attributes on Linux: +// its own process group, and SIGKILL when the supervisor dies. +func superviseAttr() *syscall.SysProcAttr { + return &syscall.SysProcAttr{Setpgid: true, Pdeathsig: syscall.SIGKILL} +} diff --git a/internal/scaffold/zz_cli_inflight_server_other_test.go b/internal/scaffold/zz_cli_inflight_server_other_test.go new file mode 100644 index 0000000..674efcf --- /dev/null +++ b/internal/scaffold/zz_cli_inflight_server_other_test.go @@ -0,0 +1,11 @@ +//go:build !linux && !windows + +package scaffold + +import "syscall" + +// superviseAttr mirrors the app-store supervisor's spawn attributes outside +// Linux: its own process group only (there is no parent-death signal). +func superviseAttr() *syscall.SysProcAttr { + return &syscall.SysProcAttr{Setpgid: true} +}