From 8bdad289337c22a776a6cfc39333985be2369fac Mon Sep 17 00:00:00 2001 From: Samuel K <69881238+skevetter@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:32:20 -0600 Subject: [PATCH 1/4] test: probe runtime cancellation and process ownership --- .github/workflows/ci.yml | 12 +- README.md | 33 ++++ internal/processfixture/main.go | 217 +++++++++++++++++++++++++ internal/processfixture/observer.go | 38 +++++ internal/processprobe/observer_test.go | 212 ++++++++++++++++++++++++ internal/processprobe/process_test.go | 216 ++++++++++++++++++++++++ 6 files changed, 727 insertions(+), 1 deletion(-) create mode 100644 internal/processfixture/main.go create mode 100644 internal/processfixture/observer.go create mode 100644 internal/processprobe/observer_test.go create mode 100644 internal/processprobe/process_test.go diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 0861842..dca71bc 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -21,7 +21,17 @@ jobs: - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7 with: go-version: '1.26.8' - - run: go test -race ./... + - name: Run race tests and record process-ownership evidence + shell: bash + run: | + set -o pipefail + go test -race -json ./... | tee "$RUNNER_TEMP/test-report.json" + - uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7 + if: always() + with: + name: runtime-tests-${{ matrix.os }} + path: ${{ runner.temp }}/test-report.json + if-no-files-found: error - run: go vet ./... generated: name: Generated bindings diff --git a/README.md b/README.md index 4971360..da796e0 100644 --- a/README.md +++ b/README.md @@ -239,3 +239,36 @@ spaces and non-ASCII characters, and publishes `runtime-spawn-*` JSON artifacts. These jobs gate release automation alongside the existing quality checks. Cross-platform cancellation and descendant-process cleanup are the next host hardening gates; successful startup measurements do not certify those behaviors. + +## Process ownership experiments + +Run the real-process cancellation and crash probes with: + +```sh +mise exec -- go test -race -v ./internal/processprobe +``` + +The fixture launches a host, a plugin, and a blocking runtime child from an +executable path containing spaces and Unicode. Readiness acknowledgements +precede cancellation, and independent observer connections answer liveness +checks. The plugin acknowledges child reaping only after `exec.Cmd.Wait` returns. +CI runs these probes on Linux, macOS, and Windows with the other race tests; +`runtime-tests-*` artifacts retain their JSON test output. + +| Scenario | Behavior asserted by the spike | +| --- | --- | +| Cancel unary RPC or Exec | A cooperative plugin cancels and reaps its child | +| Child ignores interruption | A bounded forced-kill fallback reaps the child; Unix also acknowledges the ignored signal | +| Abrupt plugin death | The runtime child survives until the independent test observer kills it | +| Abrupt host death | Transport loss cancels the RPC and reaps the child, but the plugin survives until the observer kills it | + +The last two tests deliberately record ownership gaps in the current transport. +A passing probe suite does not mean abrupt process-tree cleanup is implemented. +`exec.CommandContext` and `Client.Kill()` alone do not establish a complete +process-tree policy. Windows cannot deliver `os.Interrupt` through +`os.Process.Signal`, so cancellation exercises the forced-kill fallback there. + +These results block runtime cutover until explicit host/plugin/descendant +ownership is designed and tested on all supported platforms. The fixture is an +experiment, not an exported process supervisor. Streaming stress and executable +trust/environment compatibility remain separate host-hardening work. diff --git a/internal/processfixture/main.go b/internal/processfixture/main.go new file mode 100644 index 0000000..b9468fb --- /dev/null +++ b/internal/processfixture/main.go @@ -0,0 +1,217 @@ +// Process fixture only: exercises ownership failures through the real SDK transport. +package main + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "net" + "os" + "os/exec" + "os/signal" + "syscall" + "time" + + sdkplugin "github.com/devsy-org/devsy-runtime-sdk/plugin" + "github.com/devsy-org/devsy-runtime-sdk/runtimev1" + "github.com/devsy-org/devsy-runtime-sdk/server" + "github.com/hashicorp/go-hclog" + hplugin "github.com/hashicorp/go-plugin" + "google.golang.org/grpc" + "google.golang.org/grpc/status" +) + +const ( + childRole = "child" + pluginRole = "plugin" +) + +type event struct { + Role string `json:"role"` + Kind string `json:"kind"` + PID int `json:"pid"` +} + +type fixture struct { + runtimev1.UnimplementedRuntimeDriverServer + observer string + ignore bool + events *reporter +} + +func main() { + mode := flag.String("mode", pluginRole, "fixture role") + observer := flag.String("observer", "", "loopback observer address") + ignore := flag.Bool("ignore-interrupt", false, "child ignores normal termination") + flag.Parse() + if err := run(*mode, *observer, *ignore); err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } +} + +func run(mode, observer string, ignore bool) error { + conn, err := net.DialTimeout("tcp", observer, 5*time.Second) + if err != nil { + return err + } + defer func() { _ = conn.Close() }() + events := &reporter{encoder: json.NewEncoder(conn), role: mode} + if err := events.Encode(event{Role: mode, Kind: "ready", PID: os.Getpid()}); err != nil { + return err + } + disconnected := make(chan struct{}) + go events.respond(conn, disconnected) + switch mode { + case childRole: + return child(disconnected, events, ignore) + case "host": + return host(observer, ignore) + case pluginRole: + server.Serve(&fixture{observer: observer, ignore: ignore, events: events}) + return nil + default: + return errors.New("unknown fixture role") + } +} + +func child(disconnected <-chan struct{}, events *reporter, ignore bool) error { + signals := make(chan os.Signal, 1) + signal.Notify(signals, os.Interrupt, syscall.SIGTERM) + defer signal.Stop(signals) + // Publish readiness only after the signal handler and observer connection exist. + if err := json.NewEncoder(os.Stdout). + Encode(event{Role: childRole, Kind: "armed", PID: os.Getpid()}); err != nil { + return err + } + for { + select { + case <-disconnected: + return nil + case <-signals: + if !ignore { + return nil + } + if err := events.Encode( + event{Role: childRole, Kind: "ignored", PID: os.Getpid()}, + ); err != nil { + return err + } + } + } +} + +func (f *fixture) Info( + ctx context.Context, + _ *runtimev1.InfoRequest, +) (*runtimev1.InfoResponse, error) { + if err := f.block(ctx); err != nil { + return nil, err + } + return nil, status.FromContextError(ctx.Err()).Err() +} + +func (f *fixture) Exec( + stream grpc.BidiStreamingServer[runtimev1.ExecClientMessage, runtimev1.ExecServerMessage], +) error { + first, err := stream.Recv() + if err != nil { + return err + } + if first.GetStart() == nil { + return errors.New("missing Exec start") + } + return f.block(stream.Context()) +} + +func (f *fixture) block(ctx context.Context) error { + executable, err := os.Executable() + if err != nil { + return err + } + // #nosec G204 -- Self-executable fixture, with arguments controlled by the test harness. + cmd := exec.CommandContext(ctx, executable, "--mode", childRole, "--observer", f.observer, + fmt.Sprintf("--ignore-interrupt=%t", f.ignore)) + cmd.Cancel = func() error { return cmd.Process.Signal(os.Interrupt) } + // Windows cannot deliver os.Interrupt to a child; WaitDelay also bounds this fallback. + cmd.WaitDelay = 100 * time.Millisecond + cmd.Stderr = os.Stderr + out, err := cmd.StdoutPipe() + if err != nil { + return err + } + if err := cmd.Start(); err != nil { + _ = out.Close() + return err + } + var ready event + if err := json.NewDecoder(out).Decode(&ready); err != nil { + _ = cmd.Process.Kill() + _ = cmd.Wait() + return err + } + if err := f.events.Encode( + event{Role: pluginRole, Kind: "child-ready", PID: ready.PID}, + ); err != nil { + _ = cmd.Process.Kill() + _ = cmd.Wait() + return err + } + // Wait is required even when cancellation or forced termination has occurred. + waitErr := cmd.Wait() + if err := f.events.Encode(event{Role: pluginRole, Kind: "reaped", PID: ready.PID}); err != nil { + return err + } + if ctx.Err() != nil { + return status.FromContextError(ctx.Err()).Err() + } + return waitErr +} + +func host(observer string, ignore bool) error { + executable, err := os.Executable() + if err != nil { + return err + } + // #nosec G204 -- Self-executable fixture, with arguments controlled by the test harness. + cmd := exec.Command(executable, "--mode", pluginRole, "--observer", observer, + fmt.Sprintf("--ignore-interrupt=%t", ignore)) + client := hplugin.NewClient(&hplugin.ClientConfig{ + HandshakeConfig: sdkplugin.Handshake(), + VersionedPlugins: map[int]hplugin.PluginSet{ + sdkplugin.ProtocolVersion: sdkplugin.ClientPlugins(), + }, + AllowedProtocols: []hplugin.Protocol{hplugin.ProtocolGRPC}, + Cmd: cmd, + StartTimeout: 10 * time.Second, + Logger: hclog.NewNullLogger(), + Stderr: os.Stderr, + SyncStderr: os.Stderr, + }) + defer client.Kill() + rpc, err := client.Client() + if err != nil { + return err + } + raw, err := rpc.Dispense(sdkplugin.Name) + if err != nil { + return err + } + driver, ok := raw.(runtimev1.RuntimeDriverClient) + if !ok { + return errors.New("unexpected runtime client") + } + stream, err := driver.Exec(context.Background()) + if err != nil { + return err + } + if err := stream.Send(&runtimev1.ExecClientMessage{Payload: &runtimev1.ExecClientMessage_Start{ + Start: &runtimev1.ExecStart{WorkspaceId: "process-probe", Argv: []string{"block"}}, + }}); err != nil { + return err + } + _, err = stream.Recv() + return err +} diff --git a/internal/processfixture/observer.go b/internal/processfixture/observer.go new file mode 100644 index 0000000..1ca3f6d --- /dev/null +++ b/internal/processfixture/observer.go @@ -0,0 +1,38 @@ +package main + +import ( + "encoding/json" + "net" + "os" + "sync" +) + +// One writer serializes RPC lifecycle events and independent liveness replies. +type reporter struct { + mu sync.Mutex + encoder *json.Encoder + role string +} + +func (r *reporter) Encode(e event) error { + r.mu.Lock() + defer r.mu.Unlock() + return r.encoder.Encode(e) +} + +func (r *reporter) respond(conn net.Conn, disconnected chan<- struct{}) { + defer close(disconnected) + decoder := json.NewDecoder(conn) + for { + var e event + if err := decoder.Decode(&e); err != nil { + return + } + if e.Kind != "ping" { + return + } + if err := r.Encode(event{Role: r.role, Kind: "alive", PID: os.Getpid()}); err != nil { + return + } + } +} diff --git a/internal/processprobe/observer_test.go b/internal/processprobe/observer_test.go new file mode 100644 index 0000000..993a8c3 --- /dev/null +++ b/internal/processprobe/observer_test.go @@ -0,0 +1,212 @@ +package processprobe_test + +import ( + "encoding/json" + "errors" + "net" + "os" + "sync" + "testing" + "time" +) + +type event struct { + Role string `json:"role"` + Kind string `json:"kind"` + PID int `json:"pid"` +} + +type process struct { + handle *os.Process + pid int + role string + pong chan struct{} + conn net.Conn + closed chan struct{} +} + +type observer struct { + listener net.Listener + pending []event + events chan event + mu sync.Mutex + processes map[string]*process +} + +func observe(t *testing.T) *observer { + t.Helper() + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + o := &observer{ + listener: listener, + events: make(chan event, 32), + processes: make(map[string]*process), + } + t.Cleanup(func() { o.cleanup(t) }) + go o.accept() + return o +} + +func (o *observer) accept() { + for { + conn, err := o.listener.Accept() + if err != nil { + return + } + go o.read(conn) + } +} + +func (o *observer) read(conn net.Conn) { + defer func() { _ = conn.Close() }() + decoder := json.NewDecoder(conn) + var ready event + if err := decoder.Decode(&ready); err != nil { + return + } + handle, err := os.FindProcess(ready.PID) + if err != nil { + return + } + p := &process{ + handle: handle, + pid: ready.PID, + conn: conn, + closed: make(chan struct{}), + role: ready.Role, + pong: make(chan struct{}, 1), + } + defer close(p.closed) + o.mu.Lock() + o.processes[ready.Role] = p + o.mu.Unlock() + o.events <- ready + for { + var next event + if err := decoder.Decode(&next); err != nil { + return + } + if next.Kind == "alive" { + p.pong <- struct{}{} + continue + } + o.events <- next + } +} + +func (o *observer) await(t *testing.T, role, kind string) event { + t.Helper() + if e, ok := o.take(role, kind); ok { + return e + } + timer := time.NewTimer(15 * time.Second) + defer timer.Stop() + for { + select { + case e := <-o.events: + if e.Role == role && e.Kind == kind { + return e + } + o.pending = append(o.pending, e) + case <-timer.C: + t.Fatalf("timed out waiting for %s/%s", role, kind) + return event{} + } + } +} + +func (o *observer) lookup(t *testing.T, role string) *process { + t.Helper() + o.mu.Lock() + defer o.mu.Unlock() + p := o.processes[role] + if p == nil || p.pid <= 0 { + t.Fatalf("missing %s process", role) + } + return p +} + +func (p *process) exited(t *testing.T) { + t.Helper() + select { + case <-p.closed: + case <-time.After(10 * time.Second): + t.Fatalf("process %d still holds its observer connection", p.pid) + } +} + +func (p *process) alive(t *testing.T) { + t.Helper() + select { + case <-p.closed: + t.Fatalf("process %d unexpectedly closed its observer connection", p.pid) + default: + } + if err := p.conn.SetWriteDeadline(time.Now().Add(5 * time.Second)); err != nil { + t.Fatal(err) + } + if err := json.NewEncoder(p.conn).Encode(event{Role: p.role, Kind: "ping"}); err != nil { + t.Fatal(err) + } + select { + case <-p.pong: + case <-p.closed: + t.Fatalf("process %d exited during liveness check", p.pid) + case <-time.After(5 * time.Second): + t.Fatalf("process %d did not answer liveness check", p.pid) + } +} + +func (p *process) kill(t *testing.T) { + t.Helper() + if err := p.handle.Kill(); err != nil { + t.Fatal(err) + } + p.exited(t) +} + +func (o *observer) cleanup(t *testing.T) { + t.Helper() + _ = o.listener.Close() + o.mu.Lock() + defer o.mu.Unlock() + // Retain process handles from readiness through cleanup, rather than rediscovering a PID after death. + for _, p := range o.processes { + select { + case <-p.closed: + default: + if err := p.handle.Kill(); err != nil && !errors.Is(err, os.ErrProcessDone) { + t.Errorf("kill fixture %d: %v", p.pid, err) + } + select { + case <-p.closed: + case <-time.After(10 * time.Second): + t.Errorf("fixture %d did not disconnect after cleanup", p.pid) + } + } + _ = p.conn.Close() + _ = p.handle.Release() + } +} + +func (o *observer) take(role, kind string) (event, bool) { + for i, e := range o.pending { + if e.Role == role && e.Kind == kind { + o.pending = append(o.pending[:i], o.pending[i+1:]...) + return e, true + } + } + return event{}, false +} + +func (o *observer) readyChild(t *testing.T) *process { + t.Helper() + armed := o.await(t, "plugin", "child-ready") + ready := o.await(t, "child", "ready") + if armed.PID != ready.PID { + t.Fatal("child readiness PID mismatch") + } + return o.lookup(t, "child") +} diff --git a/internal/processprobe/process_test.go b/internal/processprobe/process_test.go new file mode 100644 index 0000000..be17881 --- /dev/null +++ b/internal/processprobe/process_test.go @@ -0,0 +1,216 @@ +// Package processprobe_test records the limits of plain go-plugin process ownership. +package processprobe_test + +import ( + "context" + "fmt" + "os" + "os/exec" + "path/filepath" + "runtime" + "testing" + "time" + + sdkplugin "github.com/devsy-org/devsy-runtime-sdk/plugin" + "github.com/devsy-org/devsy-runtime-sdk/runtimev1" + "github.com/hashicorp/go-hclog" + hplugin "github.com/hashicorp/go-plugin" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +var executable string + +func TestMain(m *testing.M) { + dir, err := os.MkdirTemp("", "devsy-process-probe-") + if err != nil { + panic(err) + } + executable = filepath.Join(dir, "process λ fixture.exe") + // #nosec G204 -- Fixed fixture package; output path is owned by the test harness. + cmd := exec.Command("go", "build", "-race", "-o", executable, "../processfixture") + cmd.Stdout, cmd.Stderr = os.Stdout, os.Stderr + if err := cmd.Run(); err != nil { + _ = os.RemoveAll(dir) + os.Exit(1) + } + code := m.Run() + _ = os.RemoveAll(dir) + os.Exit(code) +} + +func launch(t *testing.T, o *observer, ignore bool) (*hplugin.Client, *exec.Cmd) { + t.Helper() + // #nosec G204 -- Fixed fixture built by TestMain; no user-controlled executable. + cmd := exec.Command(executable, "--mode", "plugin", "--observer", o.listener.Addr().String(), + fmt.Sprintf("--ignore-interrupt=%t", ignore)) + client := hplugin.NewClient(&hplugin.ClientConfig{ + HandshakeConfig: sdkplugin.Handshake(), + VersionedPlugins: map[int]hplugin.PluginSet{ + sdkplugin.ProtocolVersion: sdkplugin.ClientPlugins(), + }, + AllowedProtocols: []hplugin.Protocol{hplugin.ProtocolGRPC}, + Cmd: cmd, + StartTimeout: 10 * time.Second, + Logger: hclog.NewNullLogger(), + Stderr: os.Stderr, + SyncStderr: os.Stderr, + }) + t.Cleanup(client.Kill) + return client, cmd +} + +func connect(t *testing.T, client *hplugin.Client) runtimev1.RuntimeDriverClient { + t.Helper() + rpc, err := client.Client() + if err != nil { + t.Fatal(err) + } + raw, err := rpc.Dispense(sdkplugin.Name) + if err != nil { + t.Fatal(err) + } + driver, ok := raw.(runtimev1.RuntimeDriverClient) + if !ok { + t.Fatalf("unexpected client %T", raw) + } + return driver +} + +func operation( + ctx context.Context, + t *testing.T, + driver runtimev1.RuntimeDriverClient, + unary bool, +) <-chan error { + t.Helper() + result := make(chan error, 1) + if unary { + go func() { _, err := driver.Info(ctx, &runtimev1.InfoRequest{}); result <- err }() + return result + } + stream, err := driver.Exec(ctx) + if err != nil { + t.Fatal(err) + } + if err := stream.Send(&runtimev1.ExecClientMessage{Payload: &runtimev1.ExecClientMessage_Start{ + Start: &runtimev1.ExecStart{WorkspaceId: "process-probe", Argv: []string{"block"}}, + }}); err != nil { + t.Fatal(err) + } + go func() { _, err := stream.Recv(); result <- err }() + return result +} + +func result(t *testing.T, errors <-chan error) error { + t.Helper() + select { + case err := <-errors: + return err + case <-time.After(15 * time.Second): + t.Fatal("operation did not terminate") + return nil + } +} + +func TestCancellationReapsRuntimeChild(t *testing.T) { + for _, unary := range []bool{true, false} { + for _, ignore := range []bool{false, true} { + t.Run(fmt.Sprintf("unary=%t/ignore-interrupt=%t", unary, ignore), func(t *testing.T) { + o := observe(t) + client, _ := launch(t, o, ignore) + driver := connect(t, client) + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + pending := operation(ctx, t, driver, unary) + child := o.readyChild(t) + child.alive(t) + cancel() + if err := result(t, pending); status.Code(err) != codes.Canceled { + t.Fatalf("cancel: %v", err) + } + // A canceled RPC can return before the plugin finishes cleanup: wait for the actual Wait acknowledgement. + awaitIgnored(t, o, ignore) + reaped := o.await(t, "plugin", "reaped") + if reaped.PID != child.pid { + t.Fatal("wrong process was reaped") + } + child.exited(t) + client.Kill() + if !client.Exited() { + t.Fatal("plugin was not reaped") + } + o.lookup(t, "plugin").exited(t) + }) + } + } +} + +func TestPluginCrashLeavesRuntimeChild(t *testing.T) { + o := observe(t) + client, cmd := launch(t, o, true) + driver := connect(t, client) + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + pending := operation(ctx, t, driver, false) + child := o.readyChild(t) + child.alive(t) + if err := cmd.Process.Kill(); err != nil { + t.Fatal(err) + } + if err := result(t, pending); status.Code(err) != codes.Unavailable { + t.Fatalf("crash: %v", err) + } + client.Kill() + if !client.Exited() { + t.Fatal("plugin was not reaped") + } + o.lookup(t, "plugin").exited(t) + // This is a measured ownership gap, not a production guarantee. The observer prevents a test orphan. + child.alive(t) + t.Log("plain go-plugin cannot clean a runtime child after abrupt plugin death") + child.kill(t) +} + +func TestHostDeathLeavesPlugin(t *testing.T) { + o := observe(t) + // #nosec G204 -- Fixed fixture built by TestMain, with a test-owned observer. + cmd := exec.Command( + executable, + "--mode", + "host", + "--observer", + o.listener.Addr().String(), + "--ignore-interrupt=true", + ) + cmd.Stdout, cmd.Stderr = os.Stdout, os.Stderr + if err := cmd.Start(); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = cmd.Process.Kill(); _ = cmd.Wait() }) + child := o.readyChild(t) + plugin := o.lookup(t, "plugin") + child.alive(t) + if err := cmd.Process.Kill(); err != nil { + t.Fatal(err) + } + if err := cmd.Wait(); err == nil { + t.Fatal("killed host reported success") + } + o.lookup(t, "host").exited(t) + // Transport loss cancels the RPC, allowing a cooperative plugin to reap its child. + o.await(t, "plugin", "reaped") + child.exited(t) + plugin.alive(t) + t.Log( + "transport cancellation cleans the child, but abrupt host death leaves the plugin serving", + ) + plugin.kill(t) +} + +func awaitIgnored(t *testing.T, o *observer, ignore bool) { + t.Helper() + if ignore && runtime.GOOS != "windows" { + o.await(t, "child", "ignored") + } +} From 1a397d7fb504b9698acfbd85960f17f030ec4c1f Mon Sep 17 00:00:00 2001 From: Samuel K <69881238+skevetter@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:37:05 -0600 Subject: [PATCH 2/4] test: isolate crash probe temporary files --- internal/processprobe/observer_test.go | 12 ++++++++++++ internal/processprobe/process_test.go | 8 ++++++++ 2 files changed, 20 insertions(+) diff --git a/internal/processprobe/observer_test.go b/internal/processprobe/observer_test.go index 993a8c3..e629699 100644 --- a/internal/processprobe/observer_test.go +++ b/internal/processprobe/observer_test.go @@ -26,6 +26,7 @@ type process struct { } type observer struct { + directory string listener net.Listener pending []event events chan event @@ -35,11 +36,21 @@ type observer struct { func observe(t *testing.T) *observer { t.Helper() + directory, err := os.MkdirTemp("", "dpc-") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + if err := os.RemoveAll(directory); err != nil { + t.Errorf("remove fixture directory: %v", err) + } + }) listener, err := net.Listen("tcp", "127.0.0.1:0") if err != nil { t.Fatal(err) } o := &observer{ + directory: directory, listener: listener, events: make(chan event, 32), processes: make(map[string]*process), @@ -169,6 +180,7 @@ func (p *process) kill(t *testing.T) { func (o *observer) cleanup(t *testing.T) { t.Helper() + _ = o.listener.Close() o.mu.Lock() defer o.mu.Unlock() diff --git a/internal/processprobe/process_test.go b/internal/processprobe/process_test.go index be17881..9da3009 100644 --- a/internal/processprobe/process_test.go +++ b/internal/processprobe/process_test.go @@ -44,8 +44,10 @@ func launch(t *testing.T, o *observer, ignore bool) (*hplugin.Client, *exec.Cmd) // #nosec G204 -- Fixed fixture built by TestMain; no user-controlled executable. cmd := exec.Command(executable, "--mode", "plugin", "--observer", o.listener.Addr().String(), fmt.Sprintf("--ignore-interrupt=%t", ignore)) + cmd.Env = fixtureEnvironment(o.directory) client := hplugin.NewClient(&hplugin.ClientConfig{ HandshakeConfig: sdkplugin.Handshake(), + SkipHostEnv: true, VersionedPlugins: map[int]hplugin.PluginSet{ sdkplugin.ProtocolVersion: sdkplugin.ClientPlugins(), }, @@ -183,6 +185,7 @@ func TestHostDeathLeavesPlugin(t *testing.T) { o.listener.Addr().String(), "--ignore-interrupt=true", ) + cmd.Env = fixtureEnvironment(o.directory) cmd.Stdout, cmd.Stderr = os.Stdout, os.Stderr if err := cmd.Start(); err != nil { t.Fatal(err) @@ -214,3 +217,8 @@ func awaitIgnored(t *testing.T, o *observer, ignore bool) { o.await(t, "child", "ignored") } } + +func fixtureEnvironment(directory string) []string { + // Keep inherited compatibility, while placing crash-left socket files in a short, test-owned directory. + return append(os.Environ(), "TMPDIR="+directory, "TEMP="+directory, "TMP="+directory) +} From 8e6ab37fbf9f3b116c7484101671c90e56a9d80e Mon Sep 17 00:00:00 2001 From: Samuel K <69881238+skevetter@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:41:22 -0600 Subject: [PATCH 3/4] test: unblock process observer during teardown --- internal/processprobe/observer_test.go | 41 ++++++++++++++++++++++---- 1 file changed, 35 insertions(+), 6 deletions(-) diff --git a/internal/processprobe/observer_test.go b/internal/processprobe/observer_test.go index e629699..46e0fee 100644 --- a/internal/processprobe/observer_test.go +++ b/internal/processprobe/observer_test.go @@ -30,6 +30,7 @@ type observer struct { listener net.Listener pending []event events chan event + done chan struct{} mu sync.Mutex processes map[string]*process } @@ -53,6 +54,7 @@ func observe(t *testing.T) *observer { directory: directory, listener: listener, events: make(chan event, 32), + done: make(chan struct{}), processes: make(map[string]*process), } t.Cleanup(func() { o.cleanup(t) }) @@ -90,20 +92,23 @@ func (o *observer) read(conn net.Conn) { pong: make(chan struct{}, 1), } defer close(p.closed) - o.mu.Lock() - o.processes[ready.Role] = p - o.mu.Unlock() - o.events <- ready + if !o.register(ready.Role, p) { + return + } + o.forward(ready) for { var next event if err := decoder.Decode(&next); err != nil { return } if next.Kind == "alive" { - p.pong <- struct{}{} + select { + case p.pong <- struct{}{}: + case <-o.done: + } continue } - o.events <- next + o.forward(next) } } @@ -181,6 +186,7 @@ func (p *process) kill(t *testing.T) { func (o *observer) cleanup(t *testing.T) { t.Helper() + close(o.done) _ = o.listener.Close() o.mu.Lock() defer o.mu.Unlock() @@ -222,3 +228,26 @@ func (o *observer) readyChild(t *testing.T) *process { } return o.lookup(t, "child") } + +func (o *observer) forward(e event) { + // Once teardown starts, stop queueing evidence while still reading until the process disconnects. + select { + case o.events <- e: + case <-o.done: + } +} + +func (o *observer) register(role string, p *process) bool { + o.mu.Lock() + defer o.mu.Unlock() + select { + case <-o.done: + // A readiness message may race with teardown; do not leave this late fixture untracked. + _ = p.handle.Kill() + _ = p.handle.Release() + return false + default: + o.processes[role] = p + return true + } +} From de4c11e6a955ef93e2511c9ec860a7886192888e Mon Sep 17 00:00:00 2001 From: Samuel K <69881238+skevetter@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:42:51 -0600 Subject: [PATCH 4/4] refactor: separate process observer event routing --- internal/processprobe/observer_test.go | 20 ++++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/internal/processprobe/observer_test.go b/internal/processprobe/observer_test.go index 46e0fee..949d34e 100644 --- a/internal/processprobe/observer_test.go +++ b/internal/processprobe/observer_test.go @@ -101,14 +101,7 @@ func (o *observer) read(conn net.Conn) { if err := decoder.Decode(&next); err != nil { return } - if next.Kind == "alive" { - select { - case p.pong <- struct{}{}: - case <-o.done: - } - continue - } - o.forward(next) + o.deliver(p, next) } } @@ -251,3 +244,14 @@ func (o *observer) register(role string, p *process) bool { return true } } + +func (o *observer) deliver(p *process, e event) { + if e.Kind == "alive" { + select { + case p.pong <- struct{}{}: + case <-o.done: + } + return + } + o.forward(e) +}