diff --git a/docs/CLI-ADAPTER.md b/docs/CLI-ADAPTER.md index 9a62a9c..9ef4aef 100644 --- a/docs/CLI-ADAPTER.md +++ b/docs/CLI-ADAPTER.md @@ -48,9 +48,10 @@ 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 (SIGTERM, then SIGKILL after 5 s). + 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 @@ -65,6 +66,14 @@ The generated `exec.go` is defensive by default: `$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 index 249ee4b..3bfb03c 100644 --- a/internal/scaffold/cli_lifecycle_e2e_test.go +++ b/internal/scaffold/cli_lifecycle_e2e_test.go @@ -3,6 +3,7 @@ package scaffold import ( + "bufio" "encoding/json" "errors" "net" @@ -19,12 +20,11 @@ import ( "github.com/pilot-protocol/app-store/pkg/ipc" ) -// TestCLIAdapterChildLifecycleE2E: a CLI child that is still running when its -// call times out, or when the adapter is stopped, must not outlive the call or -// the adapter (io.pilot.miren: `miren login` over miren.exec was left running, -// reparented to init, after the supervisor SIGKILLed the adapter), and a -// timed-out call must fail rather than report a successful {"exit":-1}. -func TestCLIAdapterChildLifecycleE2E(t *testing.T) { +// 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") } @@ -32,7 +32,88 @@ func TestCLIAdapterChildLifecycleE2E(t *testing.T) { 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") @@ -44,72 +125,37 @@ func TestCLIAdapterChildLifecycleE2E(t *testing.T) { if err := os.WriteFile(tool, []byte(script), 0o755); err != nil { t.Fatal(err) } - spec := ` + bin, proj := buildCLIAdapter(t, ` id: io.pilot.faketool app_version: 0.1.0 description: "Fronts faketool." namespace: faketool backend: type: cli - command: ["` + tool + `"] + command: ["`+tool+`"] methods: - name: faketool.hang summary: "Blocks until killed." - timeout: 2s + 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} -` - 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) - } +`) + manifest := filepath.Join(proj, "manifest.json") start := func(t *testing.T) (*exec.Cmd, string) { t.Helper() - // Short path: t.TempDir() for a subtest can exceed sun_path (104 on darwin). - dir, err := os.MkdirTemp("", "cla") - if err != nil { - t.Fatal(err) - } - t.Cleanup(func() { os.RemoveAll(dir) }) - sock := filepath.Join(dir, "app.sock") - a := exec.Command(bin, "--socket", sock, "--manifest", filepath.Join(proj, "manifest.json")) + 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() }) - 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 a, sock - } - } - t.Fatal("adapter socket never appeared") - return nil, "" - } - call := func(sock, method, args string) error { - conn, err := net.DialTimeout("unix", sock, 3*time.Second) - if err != nil { - return err - } - defer conn.Close() - _ = conn.SetDeadline(time.Now().Add(30 * time.Second)) - var out json.RawMessage - return ipc.Call(conn, method, json.RawMessage(args), &out) + waitForSocket(t, sock, 10*time.Second) + return a, sock } childPID := func(t *testing.T, pidfile string) int { t.Helper() @@ -123,55 +169,88 @@ methods: t.Fatal("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 err := syscall.Kill(pid, 0); errors.Is(err, syscall.ESRCH) { - return true - } - if time.Now().After(deadline) { - _ = syscall.Kill(pid, syscall.SIGKILL) - return false - } - } + 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") - err := call(sock, "faketool.hang", `{"pidfile":"`+pidfile+`"}`) + 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: err = %v, want an IPC error mentioning the deadline (not a {\"exit\":-1} result)", err) + 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); !gone(pid, 3*time.Second) { + if pid := childPID(t, pidfile); !processGone(pid, 3*time.Second) { t.Errorf("child pid %d still running after its call timed out", pid) } }) - stopWith := func(sig syscall.Signal) func(t *testing.T) { - return func(t *testing.T) { - a, sock := start(t) - pidfile := filepath.Join(t.TempDir(), "pid") - go func() { _ = call(sock, "faketool.run", `{"args":["hang","`+pidfile+`"]}`) }() - pid := childPID(t, pidfile) - _ = a.Process.Signal(sig) - _, _ = a.Process.Wait() - // SIGTERM: the adapter must not exit before its child is gone. - // SIGKILL (the supervisor's stop path; also what Pdeathsig does to the - // adapter when the daemon dies): the kernel's parent-death signal. - within := 100 * time.Millisecond - if sig == syscall.SIGKILL { - within = 2 * time.Second - } - if !gone(pid, within) { - t.Errorf("in-flight child pid %d outlived the adapter (%v)", pid, sig) - } + 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) } - } - t.Run("SIGTERM stops the in-flight child before exit", stopWith(syscall.SIGTERM)) + 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) { - if runtime.GOOS != "linux" { - t.Skip("no parent-death signal outside Linux; there the supervisor's process-group stop reaches the child") + // 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) } - stopWith(syscall.SIGKILL)(t) }) } diff --git a/internal/scaffold/stage_lifecycle_e2e_test.go b/internal/scaffold/stage_lifecycle_e2e_test.go index 3f805fb..a36b5a0 100644 --- a/internal/scaffold/stage_lifecycle_e2e_test.go +++ b/internal/scaffold/stage_lifecycle_e2e_test.go @@ -13,97 +13,109 @@ import ( "path/filepath" "runtime" "strings" + "sync" + "sync/atomic" "syscall" "testing" "time" ) -// TestStageInterruptedLeavesNothingE2E: stopping an adapter while it is still -// downloading its native binary on first start must not leave the partial -// download behind in the shared TMPDIR (io.pilot.miren: a 40-60 MB -// pilot-asset-* per interrupted start), SIGTERM must be handled (clean exit, -// partial removed), and a SIGKILLed start's leftover must be swept by the next -// start. -func TestStageInterruptedLeavesNothingE2E(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() - cfg := parseSpec(t, cliAssetsSpec) - 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) +// 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 either slowly (64 KiB per 100 ms) or at once. + // 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 r.URL.Query().Get("slow") == "" { + if fail.Load() { + http.Error(w, "registry down", http.StatusInternalServerError) + return + } + mode := r.URL.Query().Get("mode") + if mode == "" { _, _ = w.Write(body) return } - for off := 0; off < len(body); off += 64 << 10 { - end := min(off+64<<10, len(body)) - if _, err := w.Write(body[off:end]); err != nil { - return - } - w.(http.Flusher).Flush() - select { - case <-r.Context().Done(): - return - case <-time.After(100 * time.Millisecond): - } + _, _ = 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 R2 URLs are placeholders). - writeSpec := func(t *testing.T, e env, slow bool) { + // 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 slow { - url += "?slow=1" + 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{{"name": "toolx", "role": "binary", "os": runtime.GOOS, "arch": runtime.GOARCH, + "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, slow bool) env { + setup := func(t *testing.T, mode string) env { t.Helper() - dir, err := os.MkdirTemp("", "stg") // short: sun_path - if err != nil { - t.Fatal(err) - } - t.Cleanup(func() { os.RemoveAll(dir) }) + 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) } } - mf, _ := os.ReadFile(filepath.Join(proj, "manifest.json")) - _ = os.WriteFile(filepath.Join(e.app, "manifest.json"), mf, 0o644) - writeSpec(t, e, slow) + 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 { @@ -111,6 +123,7 @@ func TestStageInterruptedLeavesNothingE2E(t *testing.T) { 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) } @@ -121,8 +134,7 @@ func TestStageInterruptedLeavesNothingE2E(t *testing.T) { 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-*")) - tm, _ := filepath.Glob(filepath.Join(e.tmp, "pilot-asset-*")) - for _, f := range append(m, tm...) { + for _, f := range m { if fi, err := os.Stat(f); err == nil && fi.Size() > 0 { return } @@ -136,9 +148,49 @@ func TestStageInterruptedLeavesNothingE2E(t *testing.T) { 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, true) + e := setup(t, "hang") a := start(t, e) waitPartial(t, e) _ = a.Process.Signal(syscall.SIGTERM) @@ -155,34 +207,67 @@ func TestStageInterruptedLeavesNothingE2E(t *testing.T) { 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, true) + e := setup(t, "hang") a := start(t, e) waitPartial(t, e) _ = a.Process.Kill() _, _ = a.Process.Wait() noTmpLeft(t, e) - writeSpec(t, e, false) - b := start(t, e) - for deadline := time.Now().Add(15 * time.Second); ; time.Sleep(20 * time.Millisecond) { - if _, err := os.Stat(e.sock); err == nil { - break - } - if time.Now().After(deadline) { - t.Fatal("second start never listened") - } + 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) } - if _, err := os.Stat(filepath.Join(e.app, "bin", "toolx")); err != nil { - t.Errorf("asset not staged: %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) - _ = b.Process.Signal(syscall.SIGTERM) - _, _ = b.Process.Wait() }) } diff --git a/internal/scaffold/templates/client_cli.go.tmpl b/internal/scaffold/templates/client_cli.go.tmpl index c6f64c2..e6b276b 100644 --- a/internal/scaffold/templates/client_cli.go.tmpl +++ b/internal/scaffold/templates/client_cli.go.tmpl @@ -63,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. @@ -86,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 @@ -113,7 +126,18 @@ 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 @@ -147,14 +171,14 @@ func (r *Runner) Run(ctx context.Context, spec Spec, call Call) (json.RawMessage // 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", r.base[0], ctxErr, strings.TrimSpace(stderr.String())) + 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 7180ac4..5908ff8 100644 --- a/internal/scaffold/templates/main.go.tmpl +++ b/internal/scaffold/templates/main.go.tmpl @@ -108,19 +108,12 @@ func main() { appDir = filepath.Dir(*manifestPath) } {{- if .HasAssets}} - stagedCmd, err := backend.StageAssetsContext(ctx, appDir) - if err != nil { - if ctx.Err() != nil { - log.Printf("{{.BinaryName}}: stopped while installing assets: %v", err) - return - } - 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}} @@ -132,24 +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.StageAssetsContext(ctx, appDir) - if err != nil { - if ctx.Err() != nil { - log.Printf("{{.BinaryName}}: stopped while installing assets: %v", err) - return - } - 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}} @@ -165,6 +151,12 @@ func main() { {{- end}} 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}} diff --git a/internal/scaffold/templates/stage.go.tmpl b/internal/scaffold/templates/stage.go.tmpl index c235128..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 StageAssetsContext(ctx, $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" ) @@ -51,6 +61,14 @@ const installStepTimeout = 2 * time.Minute // 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 @@ -75,6 +93,99 @@ 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 @@ -114,6 +225,9 @@ func StageAssetsContext(ctx context.Context, appDir string) (string, error) { 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 { @@ -250,6 +364,9 @@ func extractTarGz(ctx context.Context, archive, dir string) error { 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))) } @@ -262,7 +379,10 @@ 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) } @@ -310,7 +430,10 @@ func runInstallStep(ctx context.Context, execPath string, args []string) error { 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))) }