From 79ff80ed46ef18fa26319a294eaf482ff387fc54 Mon Sep 17 00:00:00 2001 From: Amp Date: Sat, 26 Sep 2026 08:35:56 +0000 Subject: [PATCH] fix(execution): reserve process capacity before launching Reject starts when all slots are live or reserved, prune only completed history, and release reservations on launch failure. Keep live processes tracked even after failed termination. Add concurrent admission and real-process regressions for #1085. Amp-Thread-ID: https://ampcode.com/threads/T-01a0dcd5-eb93-74c9-a1ac-5b140edaf973 Co-authored-by: Pierre Bruno --- internal/execution/process_manager.go | 50 +++++--- internal/execution/process_manager_test.go | 140 +++++++++++++++++++++ 2 files changed, 174 insertions(+), 16 deletions(-) diff --git a/internal/execution/process_manager.go b/internal/execution/process_manager.go index 29403cbe7..ead61e0bb 100644 --- a/internal/execution/process_manager.go +++ b/internal/execution/process_manager.go @@ -24,11 +24,15 @@ const ( var ( ErrProcessNotFound = errors.New("execution process not found") ErrProcessStdinDisabled = errors.New("execution process does not accept stdin") + ErrProcessCapacity = errors.New("execution process capacity reached; stop a running process before starting another") ) type ProcessManagerOptions struct { CompletedRetention time.Duration - MaxProcesses int + // MaxProcesses bounds retained processes and in-flight launches. Completed + // history is pruned before rejecting a start; live processes are never evicted. + // Zero uses the default limit; a negative value disables the limit. + MaxProcesses int } // ProcessManager owns retained interactive-process identity, transport, @@ -37,6 +41,7 @@ type ProcessManager struct { mu sync.Mutex nextID int processes map[int]*managedProcess + starting int completedRetention time.Duration maxProcesses int startTransport processTransportStarter @@ -120,12 +125,21 @@ func (manager *ProcessManager) Start(ctx context.Context, input ProcessStart, wa if input.Prepared.Command == nil { return ProcessResult{}, errors.New("prepared execution has no command") } + if err := manager.reserve(); err != nil { + if input.Prepared.Cleanup != nil { + input.Prepared.Cleanup() + } + return ProcessResult{}, err + } command := input.Prepared.Command buffer := newProcessOutputBuffer() request := input.Request observer := NewChangeObserver(request.WorkspaceRoots[0]) stdin, tty, transportCleanup, err := manager.startTransport(command, buffer, input.TTY) if err != nil { + manager.mu.Lock() + manager.starting-- + manager.mu.Unlock() if input.Prepared.Cleanup != nil { input.Prepared.Cleanup() } @@ -331,9 +345,13 @@ func (manager *ProcessManager) StopAll() []int { return ids } +// Remove forgets completed history only. A live process must remain tracked +// until completion, even if a termination attempt fails. func (manager *ProcessManager) Remove(id int) { manager.mu.Lock() - delete(manager.processes, id) + if process, ok := manager.processes[id]; ok && process.doneClosed() { + delete(manager.processes, id) + } manager.mu.Unlock() } @@ -358,22 +376,25 @@ func (manager *ProcessManager) get(id int) (*managedProcess, bool) { return process, ok } -func (manager *ProcessManager) store(process *managedProcess) { +func (manager *ProcessManager) reserve() error { manager.mu.Lock() - var evicted *managedProcess - if manager.maxProcesses > 0 && len(manager.processes) >= manager.maxProcesses { - evicted = manager.processToPruneLocked() - if evicted != nil && evicted.doneClosed() { - delete(manager.processes, evicted.id) - evicted = nil + defer manager.mu.Unlock() + if manager.maxProcesses > 0 && len(manager.processes)+manager.starting >= manager.maxProcesses { + completed := manager.processToPruneLocked() + if completed == nil { + return ErrProcessCapacity } + delete(manager.processes, completed.id) } + manager.starting++ + return nil +} + +func (manager *ProcessManager) store(process *managedProcess) { + manager.mu.Lock() + manager.starting-- manager.processes[process.id] = process manager.mu.Unlock() - if evicted != nil { - _, _ = evicted.output.Write([]byte("[zero] session evicted: too many background terminals\n")) - evicted.terminate() - } } func (manager *ProcessManager) processToPruneLocked() *managedProcess { @@ -387,9 +408,6 @@ func (manager *ProcessManager) processToPruneLocked() *managedProcess { return process } } - if len(processes) > 8 { - return processes[0] - } return nil } diff --git a/internal/execution/process_manager_test.go b/internal/execution/process_manager_test.go index 219d74eca..e83103dea 100644 --- a/internal/execution/process_manager_test.go +++ b/internal/execution/process_manager_test.go @@ -3,6 +3,7 @@ package execution import ( "context" "errors" + "fmt" "io" "os" "os/exec" @@ -23,6 +24,145 @@ func processManagerRequest(root string, command *exec.Cmd) Request { } } +func TestProcessManagerReservesCapacityBeforeTransport(t *testing.T) { + root := t.TempDir() + manager := NewProcessManager(ProcessManagerOptions{MaxProcesses: 1}) + entered := make(chan struct{}, 20) + release := make(chan struct{}) + launchErr := errors.New("transport failed") + manager.startTransport = func(*exec.Cmd, io.Writer, bool) (io.WriteCloser, bool, func(), error) { + entered <- struct{}{} + <-release + return nil, false, nil, launchErr + } + start := func() error { + command := exec.Command(os.Args[0]) + cleaned := false + _, err := manager.Start(context.Background(), ProcessStart{ + Prepared: PreparedCommand{Command: command, Cleanup: func() { cleaned = true }}, + Request: processManagerRequest(root, command), + }, 0) + if !cleaned { + t.Error("failed admission or launch did not clean prepared resources") + } + return err + } + first := make(chan error, 1) + go func() { first <- start() }() + <-entered + results := make(chan error, 16) + for range 16 { + go func() { results <- start() }() + } + // Release blocked transports even when testing the unfixed implementation. + select { + case <-entered: + close(release) + <-first + for range 16 { + <-results + } + t.Fatal("additional transport launched while the only slot was reserved") + case err := <-results: + if err == nil || errors.Is(err, launchErr) { + t.Errorf("full manager returned %v, want capacity error", err) + } + } + for range 15 { + if err := <-results; err == nil || errors.Is(err, launchErr) { + t.Errorf("full manager returned %v, want capacity error", err) + } + } + close(release) + if err := <-first; !errors.Is(err, launchErr) { + t.Fatalf("first launch = %v", err) + } + if err := start(); !errors.Is(err, launchErr) { + t.Fatalf("failed launch did not release reservation: %v", err) + } +} + +func TestProcessManagerCapacityDoesNotEvictLiveProcesses(t *testing.T) { + for _, limit := range []int{1, 9} { + t.Run(fmt.Sprint(limit), func(t *testing.T) { + root := t.TempDir() + manager := NewProcessManager(ProcessManagerOptions{MaxProcesses: limit}) + kills := 0 + for id := range limit { + manager.processes[id] = &managedProcess{ + id: id, command: &exec.Cmd{Process: &os.Process{Pid: 123}}, + done: make(chan struct{}), output: newProcessOutputBuffer(), + kill: func(int) error { kills++; return errors.New("kill failed") }, + } + } + manager.Stop(0) // A failed kill must not free capacity. + manager.Remove(0) // Nor may removing a live identity bypass the limit. + launches := 0 + launchErr := errors.New("unexpected launch") + manager.startTransport = func(*exec.Cmd, io.Writer, bool) (io.WriteCloser, bool, func(), error) { + launches++ + return nil, false, nil, launchErr + } + command := exec.Command(os.Args[0]) + input := ProcessStart{Prepared: PreparedCommand{Command: command}, Request: processManagerRequest(root, command)} + if _, err := manager.Start(context.Background(), input, 0); err == nil || errors.Is(err, launchErr) { + t.Fatalf("full manager returned %v, want capacity error before launch", err) + } + if launches != 0 || kills != 1 || manager.Len() != limit { + t.Fatalf("launches=%d kills=%d retained=%d", launches, kills, manager.Len()) + } + manager.processes[0].markDone(nil, 0, AdapterReport{}, nil, nil) + if _, err := manager.Start(context.Background(), input, 0); !errors.Is(err, launchErr) { + t.Fatalf("completed history did not free capacity: %v", err) + } + if launches != 1 || manager.Len() != limit-1 { + t.Fatalf("launches=%d retained=%d after completion", launches, manager.Len()) + } + }) + } +} + +func TestProcessManagerCapacityHelper(t *testing.T) { + if os.Getenv("ZERO_PROCESS_CAPACITY_HELPER") != "1" { + return + } + time.Sleep(time.Minute) + os.Exit(0) +} + +func TestProcessManagerCapacityWithRunningProcess(t *testing.T) { + root := t.TempDir() + manager := NewProcessManager(ProcessManagerOptions{MaxProcesses: 1}) + t.Cleanup(func() { manager.StopAll() }) + start := func() (ProcessResult, *exec.Cmd, error) { + command := exec.Command(os.Args[0], "-test.run=^TestProcessManagerCapacityHelper$") + command.Env = append(os.Environ(), "ZERO_PROCESS_CAPACITY_HELPER=1") + result, err := manager.Start(context.Background(), ProcessStart{ + Prepared: PreparedCommand{Command: command}, Request: processManagerRequest(root, command), + }, 0) + return result, command, err + } + first, _, err := start() + if err != nil || first.Exited { + t.Fatalf("first start = %+v, %v", first, err) + } + if _, command, err := start(); err == nil || command.Process != nil { + t.Fatalf("second start = %v, process=%v; want rejection before OS launch", err, command.Process) + } + if got := len(manager.List()); got != 1 { + t.Fatalf("live processes = %d, want 1", got) + } + stopped, err := manager.Continue(context.Background(), ProcessContinue{ + ProcessID: first.ProcessID, Interrupt: true, Wait: 10 * time.Second, + }) + if err != nil || !stopped.Exited { + t.Fatalf("stop = %+v, %v", stopped, err) + } + if next, _, err := start(); err != nil || next.Exited { + t.Fatalf("start after completion = %+v, %v", next, err) + } +} + func TestProcessManagerRetainsAndContinuesWithStableIdentity(t *testing.T) { if runtime.GOOS == "windows" { t.Skip("test command uses a POSIX shell")