From 76bfbdcbcab4c064a359e57b029a623a215480df Mon Sep 17 00:00:00 2001 From: Samuel K <69881238+skevetter@users.noreply.github.com> Date: Sat, 3 Oct 2026 09:49:45 -0600 Subject: [PATCH 1/3] feat: own runtime process trees with a host-leased supervisor --- README.md | 90 ++++++++- cmd/devsy-runtime-supervisor/main.go | 10 + go.mod | 2 +- internal/processfixture/client.go | 88 ++++++++ internal/processfixture/descendant.go | 73 +++++++ internal/processfixture/main.go | 88 ++++---- internal/processfixture/observer.go | 1 + internal/processprobe/owned_test.go | 154 ++++++++++++++ supervisor/exit_darwin.go | 61 ++++++ supervisor/exit_linux.go | 42 ++++ supervisor/main.go | 116 +++++++++++ supervisor/runner.go | 279 ++++++++++++++++++++++++++ supervisor/runner_test.go | 169 ++++++++++++++++ supervisor/tree_other.go | 21 ++ supervisor/tree_unix.go | 71 +++++++ supervisor/tree_windows.go | 79 ++++++++ 16 files changed, 1294 insertions(+), 50 deletions(-) create mode 100644 cmd/devsy-runtime-supervisor/main.go create mode 100644 internal/processfixture/client.go create mode 100644 internal/processfixture/descendant.go create mode 100644 internal/processprobe/owned_test.go create mode 100644 supervisor/exit_darwin.go create mode 100644 supervisor/exit_linux.go create mode 100644 supervisor/main.go create mode 100644 supervisor/runner.go create mode 100644 supervisor/runner_test.go create mode 100644 supervisor/tree_other.go create mode 100644 supervisor/tree_unix.go create mode 100644 supervisor/tree_windows.go diff --git a/README.md b/README.md index da796e0..70d30c8 100644 --- a/README.md +++ b/README.md @@ -19,6 +19,7 @@ Generated protobuf bindings are included. Consumers do not need protoc or the de - `runtimev1`: protobuf/gRPC bindings, API versions, and `ValidateInfo` for compatibility checks. - `plugin`: shared handshake, plugin name, and gRPC registration/client bridge. - `server`: executable serving entry point. +- `supervisor`: an opt-in go-plugin runner that owns a leased plugin process tree. - `conformance`: reusable protocol behavior suite and structured-error assertions. - `conformance/fake`: persistent fake driver for host integration tests. @@ -262,13 +263,84 @@ CI runs these probes on Linux, macOS, and Windows with the other race tests; | 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. +The last two baseline tests deliberately retain the plain transport's ownership +gaps. The owned-runner regressions add a grandchild and disable the fixture's +response to RPC cancellation. They verify that plugin death, host death, and +explicit client cleanup terminate all three runtime processes without the test +observer killing survivors. A handshake-timeout regression also verifies cleanup +before the runtime becomes ready. The observer remains a fallback only when a +test fails. + +## Owning a runtime process tree + +Build the dedicated supervisor helper alongside your runtime binary: + +```sh +go build -o devsy-runtime-supervisor ./cmd/devsy-runtime-supervisor +``` + +Configure an SDK client with an explicitly verified absolute path for each +executable: + +```go +client := hplugin.NewClient(&hplugin.ClientConfig{ + HandshakeConfig: plugin.Handshake(), + VersionedPlugins: map[int]hplugin.PluginSet{ + plugin.ProtocolVersion: plugin.ClientPlugins(), + }, + AllowedProtocols: []hplugin.Protocol{hplugin.ProtocolGRPC}, + RunnerFunc: supervisor.Runner(supervisor.Options{ + SupervisorBinary: verifiedSupervisorPath, + RuntimeBinary: verifiedRuntimePath, + Args: runtimeArgs, + }), + StartTimeout: 15 * time.Second, +}) +defer client.Kill() +``` + +`hplugin` is `github.com/hashicorp/go-plugin`; `plugin` and `supervisor` are SDK +packages. Use `RunnerFunc` instead of `Cmd`. The host verifies both executables +through its existing binary distribution mechanism before configuring this +runner; go-plugin's `SecureConfig` checks a `Cmd` path and is not applicable to +this runner. An embedding host can also expose `supervisor.Main(args)` as a +dedicated helper command in its own executable, selected by `SupervisorArgs`. +`Main` always exits its process and must only run in that helper process. +The SDK's module releases do not distribute helper executables automatically. + +One supervisor belongs to one go-plugin client. `Client.Kill()` closes the host +lease and waits for the supervisor to be reaped. Abrupt host death closes the +same pipe through the OS, so cleanup continues without a live host. Plugin exit +also triggers cleanup. There is no global process pool or implicit cache. +Runtime arguments and environment configuration travel to the supervisor through +an inherited pipe; the supervisor forwards handshake stdout and diagnostics +stderr without logging the configuration. Buffered diagnostics remain available +after process reaping, until the reader drains them. + +| Platform | Ownership mechanism | +| --- | --- | +| Linux | Dedicated plugin process group; supervisor adopts and reaps orphaned descendants as a subreaper | +| macOS | Dedicated plugin process group; supervisor reaps the plugin and the OS adopts orphaned descendants | +| Windows | Supervisor joins a non-breakaway Job Object before spawning the plugin; descendants inherit membership, and the last handle closes when the supervisor exits | + +Unix observes plugin exit before reaping its group leader, keeping the group ID +reserved while terminating descendants. Linux uses `waitid` with `WNOWAIT` and +macOS uses a process-exit kqueue event. Real permission or ownership-setup errors +remain failures; there is no fallback to unowned launch. The Windows Job Object +uses [kill-on-close semantics](https://learn.microsoft.com/en-us/windows/win32/procthread/job-objects). + +This is ownership for trusted runtime commands, not a sandbox. Unix descendants +must remain in the plugin's process group and retain signalable privileges. +Daemonization, a new session/process group (including a separately created PTY +session), or privilege elevation requires an additional explicit owner. On Unix, +simultaneously killing the host and its supervisor prevents that supervisor from +performing cleanup. Long-lived runtime services and container resources must have +their own resource lifecycle rather than depend on these command processes. + +The default environment remains inherited, with optional `Env` overrides; +client-assigned handshake, certificate, and socket metadata retain precedence. +`Directory` sets the runtime working directory. The additional supervisor start +must be included in startup measurements before selecting session reuse. +Streaming stress and real-runtime trust/environment compatibility remain +separate gates before runtime cutover. diff --git a/cmd/devsy-runtime-supervisor/main.go b/cmd/devsy-runtime-supervisor/main.go new file mode 100644 index 0000000..a83de4f --- /dev/null +++ b/cmd/devsy-runtime-supervisor/main.go @@ -0,0 +1,10 @@ +// devsy-runtime-supervisor owns one plugin process tree for one host lease. +package main + +import ( + "os" + + "github.com/devsy-org/devsy-runtime-sdk/supervisor" +) + +func main() { supervisor.Main(os.Args[1:]) } diff --git a/go.mod b/go.mod index 726000c..0925690 100644 --- a/go.mod +++ b/go.mod @@ -5,6 +5,7 @@ go 1.26.0 require ( github.com/hashicorp/go-hclog v1.6.3 github.com/hashicorp/go-plugin v1.8.0 + golang.org/x/sys v0.47.0 google.golang.org/grpc v1.83.2 google.golang.org/protobuf v1.36.12 ) @@ -17,7 +18,6 @@ require ( github.com/mattn/go-isatty v0.0.17 // indirect github.com/oklog/run v1.1.0 // indirect golang.org/x/net v0.58.0 // indirect - golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.41.0 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect ) diff --git a/internal/processfixture/client.go b/internal/processfixture/client.go new file mode 100644 index 0000000..baa6737 --- /dev/null +++ b/internal/processfixture/client.go @@ -0,0 +1,88 @@ +package main + +import ( + "context" + "fmt" + "os" + "os/exec" + "time" + + sdkplugin "github.com/devsy-org/devsy-runtime-sdk/plugin" + "github.com/devsy-org/devsy-runtime-sdk/supervisor" + "github.com/hashicorp/go-hclog" + hplugin "github.com/hashicorp/go-plugin" +) + +func (f *fixture) childCommand(ctx context.Context) (*exec.Cmd, error) { + if f.uncancelable { + ctx = context.Background() + } + executable, err := os.Executable() + if err != nil { + return nil, 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, + ), + fmt.Sprintf("--grandchild=%t", f.grandchild), + ) + 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 + return cmd, nil +} + +func hostClient(observer string, behavior behavior) (*hplugin.Client, error) { + executable, err := os.Executable() + if err != nil { + return nil, 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", + behavior.ignore, + ), + fmt.Sprintf("--grandchild=%t", behavior.grandchild), + fmt.Sprintf("--uncancelable=%t", behavior.uncancelable), + ) + config := &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, + } + if behavior.owned { + config.Cmd = nil + config.RunnerFunc = supervisor.Runner( + supervisor.Options{ + SupervisorBinary: executable, + SupervisorArgs: []string{"--supervise", observer}, + RuntimeBinary: executable, + Args: cmd.Args[1:], + }, + ) + } + return hplugin.NewClient(config), nil +} diff --git a/internal/processfixture/descendant.go b/internal/processfixture/descendant.go new file mode 100644 index 0000000..8432db6 --- /dev/null +++ b/internal/processfixture/descendant.go @@ -0,0 +1,73 @@ +package main + +import ( + "encoding/json" + "errors" + "net" + "os" + "os/exec" + + "github.com/devsy-org/devsy-runtime-sdk/supervisor" +) + +func superviseFixture(observer string, args []string) { + // #nosec G704 -- Address is supplied by the test-owned loopback observer, not external input. + conn, err := net.Dial("tcp", observer) + if err != nil { + panic(err) + } + events := &reporter{encoder: json.NewEncoder(conn), role: "supervisor", address: observer} + if err := events.Encode( + event{Role: "supervisor", Kind: "ready", PID: os.Getpid()}, + ); err != nil { + panic(err) + } + disconnected := make(chan struct{}) + go events.respond(conn, disconnected) + supervisor.Main(args) +} + +func launchGrandchild(events *reporter) error { + executable, err := os.Executable() + if err != nil { + return err + } + // #nosec G204 -- Fixed self-executable fixture, never a user command. + cmd := exec.Command( + executable, + "--mode", + "grandchild", + "--observer", + events.address, + "--ignore-interrupt=true", + ) + output, err := cmd.StdoutPipe() + if err != nil { + return err + } + cmd.Stderr = os.Stderr + if err := cmd.Start(); err != nil { + _ = output.Close() + return err + } + var ready event + if err := json.NewDecoder(output).Decode(&ready); err != nil { + _ = cmd.Process.Kill() + _ = cmd.Wait() + return err + } + if ready.PID != cmd.Process.Pid { + _ = cmd.Process.Kill() + _ = cmd.Wait() + return errors.New("grandchild readiness mismatch") + } + go func() { _ = cmd.Wait() }() + return nil +} + +func maybeGrandchild(events *reporter, enabled bool) error { + if !enabled { + return nil + } + return launchGrandchild(events) +} diff --git a/internal/processfixture/main.go b/internal/processfixture/main.go index b9468fb..a782043 100644 --- a/internal/processfixture/main.go +++ b/internal/processfixture/main.go @@ -9,7 +9,6 @@ import ( "fmt" "net" "os" - "os/exec" "os/signal" "syscall" "time" @@ -17,8 +16,6 @@ import ( 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" ) @@ -36,51 +33,84 @@ type event struct { type fixture struct { runtimev1.UnimplementedRuntimeDriverServer - observer string - ignore bool - events *reporter + observer string + ignore bool + grandchild bool + uncancelable bool + events *reporter } func main() { + if len(os.Args) > 2 && os.Args[1] == "--supervise" { + superviseFixture(os.Args[2], os.Args[3:]) + } mode := flag.String("mode", pluginRole, "fixture role") observer := flag.String("observer", "", "loopback observer address") ignore := flag.Bool("ignore-interrupt", false, "child ignores normal termination") + owned := flag.Bool("owned", false, "use SDK process supervisor") + grandchild := flag.Bool("grandchild", false, "child starts a blocking grandchild") + uncancelable := flag.Bool("uncancelable", false, "runtime ignores RPC cancellation") + delayed := flag.Bool("delayed", false, "block before handshake") flag.Parse() - if err := run(*mode, *observer, *ignore); err != nil { + behavior := behavior{ + ignore: *ignore, + grandchild: *grandchild, + uncancelable: *uncancelable, + delayed: *delayed, + owned: *owned, + } + if err := run(*mode, *observer, behavior); err != nil { fmt.Fprintln(os.Stderr, err) os.Exit(1) } } -func run(mode, observer string, ignore bool) error { +type behavior struct{ ignore, grandchild, uncancelable, delayed, owned bool } + +func run(mode, observer string, behavior behavior) 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} + events := &reporter{encoder: json.NewEncoder(conn), role: mode, address: observer} 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) + return host(observer, behavior) + case childRole, "grandchild": + return child(disconnected, events, behavior) case pluginRole: - server.Serve(&fixture{observer: observer, ignore: ignore, events: events}) + if behavior.delayed { + <-disconnected + return nil + } + server.Serve( + &fixture{ + observer: observer, + ignore: behavior.ignore, + grandchild: behavior.grandchild, + uncancelable: behavior.uncancelable, + events: events, + }, + ) return nil default: return errors.New("unknown fixture role") } } -func child(disconnected <-chan struct{}, events *reporter, ignore bool) error { +func child(disconnected <-chan struct{}, events *reporter, behavior behavior) error { signals := make(chan os.Signal, 1) signal.Notify(signals, os.Interrupt, syscall.SIGTERM) defer signal.Stop(signals) + if err := maybeGrandchild(events, behavior.grandchild); err != nil { + return err + } // 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 { @@ -91,7 +121,7 @@ func child(disconnected <-chan struct{}, events *reporter, ignore bool) error { case <-disconnected: return nil case <-signals: - if !ignore { + if !behavior.ignore { return nil } if err := events.Encode( @@ -127,17 +157,10 @@ func (f *fixture) Exec( } func (f *fixture) block(ctx context.Context) error { - executable, err := os.Executable() + cmd, err := f.childCommand(ctx) 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 @@ -170,26 +193,11 @@ func (f *fixture) block(ctx context.Context) error { return waitErr } -func host(observer string, ignore bool) error { - executable, err := os.Executable() +func host(observer string, behavior behavior) error { + client, err := hostClient(observer, behavior) 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 { diff --git a/internal/processfixture/observer.go b/internal/processfixture/observer.go index 1ca3f6d..3fa0dbb 100644 --- a/internal/processfixture/observer.go +++ b/internal/processfixture/observer.go @@ -12,6 +12,7 @@ type reporter struct { mu sync.Mutex encoder *json.Encoder role string + address string } func (r *reporter) Encode(e event) error { diff --git a/internal/processprobe/owned_test.go b/internal/processprobe/owned_test.go new file mode 100644 index 0000000..c4c0d78 --- /dev/null +++ b/internal/processprobe/owned_test.go @@ -0,0 +1,154 @@ +package processprobe_test + +import ( + "context" + "os" + "os/exec" + "path/filepath" + "testing" + "time" + + sdkplugin "github.com/devsy-org/devsy-runtime-sdk/plugin" + "github.com/devsy-org/devsy-runtime-sdk/supervisor" + "github.com/hashicorp/go-hclog" + hplugin "github.com/hashicorp/go-plugin" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +func owned(t *testing.T, o *observer, delayed bool) *hplugin.Client { + t.Helper() + args := []string{ + "--mode", + "plugin", + "--observer", + o.listener.Addr().String(), + "--ignore-interrupt=true", + "--grandchild=true", + "--uncancelable=true", + } + timeout := 10 * time.Second + if delayed { + args = append(args, "--delayed=true") + timeout = time.Second + } + client := hplugin.NewClient(&hplugin.ClientConfig{ + HandshakeConfig: sdkplugin.Handshake(), + VersionedPlugins: map[int]hplugin.PluginSet{ + sdkplugin.ProtocolVersion: sdkplugin.ClientPlugins(), + }, + AllowedProtocols: []hplugin.Protocol{hplugin.ProtocolGRPC}, + RunnerFunc: supervisor.Runner(supervisor.Options{ + SupervisorBinary: executable, + SupervisorArgs: []string{"--supervise", o.listener.Addr().String()}, + RuntimeBinary: executable, + Args: args, + Env: append( + fixtureEnvironment(o.directory), + sdkplugin.Handshake().MagicCookieKey+"=stale-cookie", + "PLUGIN_PROTOCOL_VERSIONS=99", + "PLUGIN_UNIX_SOCKET_DIR="+filepath.Join(o.directory, "stale"), + ), + }), + StartTimeout: timeout, + Logger: hclog.NewNullLogger(), + Stderr: os.Stderr, + SyncStderr: os.Stderr, + UnixSocketConfig: &hplugin.UnixSocketConfig{TempDir: o.directory}, + }) + t.Cleanup(client.Kill) + return client +} + +func ownedOperation(t *testing.T, o *observer, client *hplugin.Client) <-chan error { + t.Helper() + driver := connect(t, client) + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + t.Cleanup(cancel) + pending := operation(ctx, t, driver, false) + o.readyChild(t).alive(t) + o.await(t, "grandchild", "ready") + o.lookup(t, "grandchild").alive(t) + return pending +} + +func assertOwnedCleanup(t *testing.T, o *observer) { + t.Helper() + for _, role := range []string{"supervisor", "plugin", "child", "grandchild"} { + o.lookup(t, role).exited(t) + } + entries, err := filepath.Glob(filepath.Join(o.directory, "plugin-dir*")) + if err != nil || len(entries) != 0 { + t.Fatalf("supervisor socket directory leaked: %v (%v)", entries, err) + } +} + +func TestOwnedPluginCrashTerminatesDescendants(t *testing.T) { + o := observe(t) + client := owned(t, o, false) + pending := ownedOperation(t, o, client) + if err := o.lookup(t, "plugin").handle.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("supervisor was not reaped") + } + assertOwnedCleanup(t, o) +} + +func TestOwnedCloseTerminatesUncooperativeDescendants(t *testing.T) { + o := observe(t) + client := owned(t, o, false) + pending := ownedOperation(t, o, client) + client.Kill() + if err := result(t, pending); err == nil { + t.Fatal("terminated Exec reported success") + } + if !client.Exited() { + t.Fatal("supervisor was not reaped") + } + assertOwnedCleanup(t, o) +} + +func TestOwnedHostDeathTerminatesPluginAndDescendants(t *testing.T) { + o := observe(t) + // #nosec G204 -- Fixed fixture built by TestMain; arguments are test-owned. + cmd := exec.Command(executable, "--mode", "host", "--observer", o.listener.Addr().String(), + "--owned=true", "--grandchild=true", "--uncancelable=true", "--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) + } + t.Cleanup(func() { _ = cmd.Process.Kill(); _ = cmd.Wait() }) + o.readyChild(t).alive(t) + o.await(t, "grandchild", "ready") + o.lookup(t, "grandchild").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) + assertOwnedCleanup(t, o) +} + +func TestOwnedHandshakeTimeoutReapsProcesses(t *testing.T) { + o := observe(t) + client := owned(t, o, true) + if _, err := client.Client(); err == nil { + t.Fatal("delayed handshake succeeded") + } + client.Kill() + if !client.Exited() { + t.Fatal("timed-out supervisor was not reaped") + } + o.await(t, "plugin", "ready") + o.lookup(t, "plugin").exited(t) + o.lookup(t, "supervisor").exited(t) +} diff --git a/supervisor/exit_darwin.go b/supervisor/exit_darwin.go new file mode 100644 index 0000000..67f7581 --- /dev/null +++ b/supervisor/exit_darwin.go @@ -0,0 +1,61 @@ +package supervisor + +import ( + "errors" + + "golang.org/x/sys/unix" +) + +func adoptDescendants() error { return nil } +func reapDescendants() error { return nil } + +func watchExit(pid int) (func() error, error) { + queue, err := unix.Kqueue() + if err != nil { + return nil, err + } + unix.CloseOnExec(queue) + event := unix.Kevent_t{Fflags: unix.NOTE_EXIT} + unix.SetKevent(&event, pid, unix.EVFILT_PROC, unix.EV_ADD|unix.EV_ONESHOT) + _, err = unix.Kevent(queue, []unix.Kevent_t{event}, nil, nil) + if err != nil { + _ = unix.Close(queue) + // An already-exited child remains our unreaped zombie; its PID is still reserved. + if errors.Is(err, unix.ESRCH) { + return func() error { return nil }, nil + } + return nil, err + } + return func() error { + defer func() { _ = unix.Close(queue) }() + var events [1]unix.Kevent_t + for { + _, err := unix.Kevent(queue, nil, events[:], nil) + if !errors.Is(err, unix.EINTR) { + return err + } + } + }, nil +} + +func terminateGroup(pid int) error { + err := unix.Kill(-pid, unix.SIGKILL) + if errors.Is(err, unix.ESRCH) { + return nil + } + if !errors.Is(err, unix.EPERM) { + return err + } + // Darwin reports EPERM for a zombie-only group. Preserve real permission failures. + processes, queryErr := unix.SysctlKinfoProcSlice("kern.proc.pgrp", pid) + if queryErr != nil { + return errors.Join(err, queryErr) + } + const zombieState = 5 // SZOMB in Darwin's proc.h. + for _, process := range processes { + if process.Proc.P_stat != zombieState { + return err + } + } + return nil +} diff --git a/supervisor/exit_linux.go b/supervisor/exit_linux.go new file mode 100644 index 0000000..6872b92 --- /dev/null +++ b/supervisor/exit_linux.go @@ -0,0 +1,42 @@ +package supervisor + +import ( + "errors" + + "golang.org/x/sys/unix" +) + +func adoptDescendants() error { return unix.Prctl(unix.PR_SET_CHILD_SUBREAPER, 1, 0, 0, 0) } + +func watchExit(pid int) (func() error, error) { + return func() error { + for { + var info unix.Siginfo + err := unix.Waitid(unix.P_PID, pid, &info, unix.WEXITED|unix.WNOWAIT, nil) + if !errors.Is(err, unix.EINTR) { + return err + } + } + }, nil +} + +func reapDescendants() error { + for { + var status unix.WaitStatus + _, err := unix.Wait4(-1, &status, 0, nil) + if errors.Is(err, unix.ECHILD) { + return nil + } + if err != nil && !errors.Is(err, unix.EINTR) { + return err + } + } +} + +func terminateGroup(pid int) error { + err := unix.Kill(-pid, unix.SIGKILL) + if errors.Is(err, unix.ESRCH) { + return nil + } + return err +} diff --git a/supervisor/main.go b/supervisor/main.go new file mode 100644 index 0000000..ac588f1 --- /dev/null +++ b/supervisor/main.go @@ -0,0 +1,116 @@ +package supervisor + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "io" + "os" + "os/exec" + "os/signal" + "syscall" +) + +// Main is the dedicated supervisor executable entry point; it always exits the +// process. It must never run inside the host or the runtime plugin process. +// On Windows, exiting closes the supervisor's non-inherited Job Object handle. +func Main(args []string) { + flags := flag.NewFlagSet("devsy-runtime-supervisor", flag.ContinueOnError) + leaseFD := flags.Uint64("lease", 0, "inherited host-lifetime pipe") + configFD := flags.Uint64("config", 0, "inherited configuration pipe") + if err := flags.Parse(args); err != nil { + os.Exit(1) + } + if *leaseFD == 0 || *configFD == 0 || flags.NArg() != 0 { + fmt.Fprintln(os.Stderr, "supervisor requires inherited lease and config pipes") + os.Exit(1) + } + lease := os.NewFile(uintptr(*leaseFD), "host-lease") + config := os.NewFile(uintptr(*configFD), "runtime-config") + if lease == nil || config == nil { + fmt.Fprintln(os.Stderr, "invalid supervisor pipe handles") + os.Exit(1) + } + preventLeaseInheritance(lease) + ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + err := supervise(ctx, lease, config) + cancel() + _ = lease.Close() + _ = config.Close() + if err != nil { + fmt.Fprintln(os.Stderr, "runtime supervisor:", err) + os.Exit(1) + } + // Do not close a Windows Job Object containing this process before ExitProcess. + os.Exit(0) +} + +func supervise(ctx context.Context, lease, input *os.File) error { + cfg, err := readConfiguration(input) + _ = input.Close() + if err != nil { + return err + } + defer func() { _ = os.RemoveAll(cfg.SocketDir) }() + stopped := make(chan struct{}) + go func() { _, _ = io.Copy(io.Discard, lease); close(stopped) }() + select { + case <-stopped: + return nil + case <-ctx.Done(): + return ctx.Err() + default: + } + return runTree(ctx, cfg, stopped) +} + +func runTree(ctx context.Context, cfg configuration, stopped <-chan struct{}) error { + tree, err := startTree(cfg) + if err != nil { + return err + } + exited := make(chan error, 1) + go func() { exited <- tree.Wait() }() + select { + case err := <-exited: + return err + case <-stopped: + return stopTree(cfg, tree, exited) + case <-ctx.Done(): + return errors.Join(ctx.Err(), stopTree(cfg, tree, exited)) + } +} + +func stopTree(cfg configuration, tree *processTree, exited <-chan error) error { + // Windows termination kills the supervisor itself, so delete its empty socket directory first. + if err := os.RemoveAll(cfg.SocketDir); err != nil { + fmt.Fprintln(os.Stderr, "supervisor socket cleanup:", err) + } + if err := tree.Kill(); err != nil { + return err + } + <-exited + return nil +} + +func readConfiguration(input *os.File) (configuration, error) { + var cfg configuration + decoder := json.NewDecoder(io.LimitReader(input, 1<<20)) + if err := decoder.Decode(&cfg); err != nil { + return cfg, errors.New("invalid runtime supervisor configuration") + } + if err := cfg.validate(); err != nil { + return cfg, err + } + return cfg, nil +} + +func runtimeCommand(cfg configuration) *exec.Cmd { + // #nosec G204 -- Configuration is supplied by the trusted host over an inherited pipe. + cmd := exec.Command(cfg.Binary, cfg.Args...) + cmd.Env, cmd.Dir = cfg.Env, cfg.Directory + cmd.Stdin, cmd.Stdout, cmd.Stderr = os.Stdin, os.Stdout, os.Stderr + return cmd +} diff --git a/supervisor/runner.go b/supervisor/runner.go new file mode 100644 index 0000000..5f74199 --- /dev/null +++ b/supervisor/runner.go @@ -0,0 +1,279 @@ +// Package supervisor supplies explicit, host-leased runtime process ownership. +package supervisor + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "os/exec" + "path/filepath" + "slices" + "strconv" + "strings" + "sync" + + sdkplugin "github.com/devsy-org/devsy-runtime-sdk/plugin" + "github.com/hashicorp/go-hclog" + "github.com/hashicorp/go-plugin/runner" +) + +// Options selects already verified executables. Env overrides inherited values; +// it must not contain an allowlist unless the caller deliberately wants one. +type Options struct { + SupervisorBinary string + SupervisorArgs []string + RuntimeBinary string + Args []string + Env []string + Directory string +} + +// Runner returns a go-plugin RunnerFunc. Configure RunnerFunc instead of Cmd; +// the go-plugin client explicitly owns the supervisor until Client.Kill returns. +// SupervisorBinary must implement Main's dedicated helper entry point. +func Runner(options Options) func(hclog.Logger, *exec.Cmd, string) (runner.Runner, error) { + options.Args, options.Env = slices.Clone(options.Args), slices.Clone(options.Env) + options.SupervisorArgs = slices.Clone(options.SupervisorArgs) + return func(_ hclog.Logger, spec *exec.Cmd, socketDir string) (runner.Runner, error) { + if err := validateOptions(options); err != nil { + _ = os.RemoveAll(socketDir) + return nil, err + } + cfg := configuration{ + Version: 1, + Binary: options.RuntimeBinary, + Args: options.Args, + Env: runtimeEnvironment(spec.Environ(), options.Env), + Directory: options.Directory, + SocketDir: socketDir, + } + owned, err := newRunner(options.SupervisorBinary, options.SupervisorArgs, spec.Stdin, cfg) + if err != nil { + return nil, errors.Join(err, os.RemoveAll(socketDir)) + } + return owned, nil + } +} + +func validateOptions(options Options) error { + for _, path := range []string{options.SupervisorBinary, options.RuntimeBinary} { + if !filepath.IsAbs(path) { + return errors.New("supervisor and runtime binaries must be absolute paths") + } + info, err := os.Stat(path) + if err != nil { + return err + } + if !info.Mode().IsRegular() { + return errors.New("supervisor and runtime binaries must be regular files") + } + } + if options.Directory != "" && !filepath.IsAbs(options.Directory) { + return errors.New("runtime directory must be absolute") + } + return nil +} + +type ownedRunner struct { + cmd *exec.Cmd + config configuration + leaseRead, leaseWrite *os.File + configRead, configWrite *os.File + stdout, stderr io.ReadCloser + stdoutWrite, stderrWrite *os.File + done chan struct{} + waitErr error + stop sync.Once +} + +func newRunner( + binary string, + args []string, + stdin io.Reader, + cfg configuration, +) (*ownedRunner, error) { + r := &ownedRunner{config: cfg, done: make(chan struct{})} + var err error + r.leaseRead, r.leaseWrite, err = os.Pipe() + if err != nil { + return nil, err + } + r.configRead, r.configWrite, err = os.Pipe() + if err != nil { + r.closeFiles() + return nil, err + } + // #nosec G204 -- Explicitly trusted, absolute supervisor executable; no shell or PATH resolution. + r.cmd = exec.Command(binary, args...) + r.cmd.Stdin = stdin + if err := configurePipes(r.cmd, r.leaseRead, r.configRead); err != nil { + r.closeFiles() + return nil, err + } + if err := r.outputPipes(); err != nil { + r.closeFiles() + return nil, err + } + return r, nil +} + +func (r *ownedRunner) Start(ctx context.Context) error { + if err := ctx.Err(); err != nil { + r.closeFiles() + return err + } + if err := r.cmd.Start(); err != nil { + r.closeFiles() + return err + } + r.closeChildFiles() + go func() { r.waitErr = r.cmd.Wait(); r.closeLease(); close(r.done) }() + encoded := make(chan error, 1) + go func() { + encoded <- json.NewEncoder(r.configWrite).Encode(r.config) + _ = r.configWrite.Close() + }() + select { + case err := <-encoded: + if err == nil { + return nil + } + r.closeLease() + return errors.Join(err, r.Wait(context.Background())) + case <-ctx.Done(): + r.closeLease() + _ = r.configWrite.Close() + return errors.Join(ctx.Err(), r.Wait(context.Background())) + } +} + +func (r *ownedRunner) Kill(ctx context.Context) error { + r.closeLease() + return r.Wait(ctx) +} + +func (r *ownedRunner) Wait(ctx context.Context) error { + select { + case <-r.done: + return r.waitErr + case <-ctx.Done(): + return ctx.Err() + } +} + +func (r *ownedRunner) ID() string { + if r.cmd.Process == nil { + return "" + } + return strconv.Itoa(r.cmd.Process.Pid) +} + +func (r *ownedRunner) Name() string { return r.config.Binary } +func (r *ownedRunner) Stdout() io.ReadCloser { return r.stdout } +func (r *ownedRunner) Stderr() io.ReadCloser { return r.stderr } +func (r *ownedRunner) Diagnose(context.Context) string { + return "runtime supervisor failed; inspect forwarded stderr" +} + +func (*ownedRunner) PluginToHost( + network, address string, +) (translatedNetwork, translatedAddress string, err error) { + return network, address, nil +} + +func (*ownedRunner) HostToPlugin( + network, address string, +) (translatedNetwork, translatedAddress string, err error) { + return network, address, nil +} + +func (r *ownedRunner) closeFiles() { + for _, file := range []*os.File{r.leaseRead, r.leaseWrite, r.configRead, r.configWrite, r.stdoutWrite, r.stderrWrite} { + if file != nil { + _ = file.Close() + } + } + if r.stdout != nil { + _ = r.stdout.Close() + } + if r.stderr != nil { + _ = r.stderr.Close() + } + _ = os.RemoveAll(r.config.SocketDir) +} + +func (r *ownedRunner) closeLease() { r.stop.Do(func() { _ = r.leaseWrite.Close() }) } + +// configuration travels only through an inherited pipe, never process arguments or disk. +type configuration struct { + Version int + Binary string + Args, Env []string + Directory, SocketDir string +} + +func (c configuration) validate() error { + if c.Version != 1 { + return fmt.Errorf("runtime supervisor protocol %d is unsupported; expected 1", c.Version) + } + if !filepath.IsAbs(c.Binary) { + return errors.New("runtime executable must be absolute") + } + if !filepath.IsAbs(c.SocketDir) { + return errors.New("supervisor socket directory must be absolute") + } + return nil +} + +func (r *ownedRunner) outputPipes() error { + output, writeOutput, err := os.Pipe() + if err != nil { + return err + } + r.stdout, r.stdoutWrite = &drainingPipe{File: output}, writeOutput + diagnostic, writeDiagnostic, err := os.Pipe() + if err != nil { + return err + } + r.stderr, r.stderrWrite = &drainingPipe{File: diagnostic}, writeDiagnostic + r.cmd.Stdout, r.cmd.Stderr = writeOutput, writeDiagnostic + return nil +} + +func (r *ownedRunner) closeChildFiles() { + for _, file := range []*os.File{r.leaseRead, r.configRead, r.stdoutWrite, r.stderrWrite} { + _ = file.Close() + } +} + +type drainingPipe struct{ *os.File } + +func (p *drainingPipe) Read(data []byte) (int, error) { + n, err := p.File.Read(data) + if err != nil { + _ = p.Close() + } + return n, err +} + +func runtimeEnvironment(base, overrides []string) []string { + env := append(slices.Clone(base), overrides...) + cookie := sdkplugin.Handshake().MagicCookieKey + for _, value := range base { + key, _, _ := strings.Cut(value, "=") + switch key { + case cookie, + "PLUGIN_MIN_PORT", + "PLUGIN_MAX_PORT", + "PLUGIN_PROTOCOL_VERSIONS", + "PLUGIN_CLIENT_CERT", + "PLUGIN_UNIX_SOCKET_DIR": + // Client-assigned transport metadata must survive inherited or provider environment overrides. + env = append(env, value) + } + } + return env +} diff --git a/supervisor/runner_test.go b/supervisor/runner_test.go new file mode 100644 index 0000000..f8cabe8 --- /dev/null +++ b/supervisor/runner_test.go @@ -0,0 +1,169 @@ +package supervisor + +import ( + "bytes" + "context" + "encoding/json" + "io" + "os" + "os/exec" + "path/filepath" + "testing" + "time" +) + +type inspection struct { + Args []string + Env, Directory string +} + +func TestMain(m *testing.M) { + if len(os.Args) > 1 { + switch os.Args[1] { + case "--helper-role=supervisor": + Main(os.Args[2:]) + case "--helper-role=inspect": + directory, err := os.Getwd() + if err != nil { + panic(err) + } + err = json.NewEncoder(os.Stdout). + Encode(inspection{Args: os.Args[2:], Env: os.Getenv("DEVSY_OWNERSHIP_TEST"), Directory: directory}) + if err != nil { + panic(err) + } + if _, err := os.Stderr.Write( + bytes.Repeat([]byte("diagnostic tail\n"), 64), + ); err != nil { + panic(err) + } + os.Exit(0) + } + } + os.Exit(m.Run()) +} + +func shortDirectory(t *testing.T) string { + t.Helper() + dir, err := os.MkdirTemp("", "dso-") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.RemoveAll(dir) }) + return dir +} + +func fixtureOptions(t *testing.T) Options { + t.Helper() + binary, err := os.Executable() + if err != nil { + t.Fatal(err) + } + return Options{ + SupervisorBinary: binary, SupervisorArgs: []string{"--helper-role=supervisor"}, + RuntimeBinary: binary, Args: []string{"--helper-role=inspect", "space λ argument"}, + Env: []string{"DEVSY_OWNERSHIP_TEST=inherited override"}, Directory: t.TempDir(), + } +} + +func TestRunnerPreservesConfigurationAndDiagnosticTail(t *testing.T) { + t.Setenv("DEVSY_OWNERSHIP_TEST", "host-value") + options := fixtureOptions(t) + factory := Runner(options) + options.Args[1] = "changed after factory construction" + options.Env[0] = "DEVSY_OWNERSHIP_TEST=changed" + dir := shortDirectory(t) + runtime, err := factory(nil, &exec.Cmd{}, dir) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = runtime.Stdout().Close(); _ = runtime.Stderr().Close() }) + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + if err := runtime.Start(ctx); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = runtime.Kill(context.Background()) }) + // Reaping before draining must preserve all buffered output, including stderr's last bytes. + if err := runtime.Wait(ctx); err != nil { + diagnostics, _ := io.ReadAll(runtime.Stderr()) + t.Fatalf("supervisor exit: %v; %s", err, diagnostics) + } + var got inspection + if err := json.NewDecoder(runtime.Stdout()).Decode(&got); err != nil { + t.Fatal(err) + } + if _, err := io.Copy(io.Discard, runtime.Stdout()); err != nil { + t.Fatal(err) + } + assertInspection(t, got, options.Directory) + assertDiagnosticTail(t, runtime.Stderr()) + if _, err := os.Stat(dir); !os.IsNotExist(err) { + t.Fatalf("socket directory survived: %v", err) + } +} + +func TestInvalidOptionsFailBeforeLaunch(t *testing.T) { + for _, change := range []func(*Options){ + func(o *Options) { o.RuntimeBinary = "runtime-from-PATH" }, + func(o *Options) { o.SupervisorBinary = "supervisor-from-PATH" }, + func(o *Options) { o.Directory = "relative-directory" }, + func(o *Options) { o.RuntimeBinary = filepath.Join(t.TempDir(), "missing.exe") }, + } { + options := fixtureOptions(t) + change(&options) + dir := shortDirectory(t) + if _, err := Runner(options)(nil, &exec.Cmd{}, dir); err == nil { + t.Fatal("invalid options accepted") + } + if _, err := os.Stat(dir); !os.IsNotExist(err) { + t.Fatalf("failed launch leaked socket directory: %v", err) + } + } +} + +func TestPreCanceledStartClosesResources(t *testing.T) { + options := fixtureOptions(t) + dir := shortDirectory(t) + runtime, err := Runner(options)(nil, &exec.Cmd{}, dir) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := runtime.Start(ctx); err != context.Canceled { + t.Fatalf("canceled start: %v", err) + } + if runtime.ID() != "" { + t.Fatal("canceled start launched a process") + } + if _, err := os.Stat(dir); !os.IsNotExist(err) { + t.Fatalf("canceled launch leaked directory: %v", err) + } +} + +func assertInspection(t *testing.T, got inspection, directory string) { + t.Helper() + if len(got.Args) != 1 || got.Args[0] != "space λ argument" || got.Env != "inherited override" { + t.Fatalf("configuration changed: %+v", got) + } + expected, err := os.Stat(directory) + if err != nil { + t.Fatal(err) + } + actual, err := os.Stat(got.Directory) + if err != nil || !os.SameFile(expected, actual) { + t.Fatalf("wrong working directory: %s (%v)", got.Directory, err) + } +} + +func assertDiagnosticTail(t *testing.T, input io.Reader) { + t.Helper() + stderr, err := io.ReadAll(input) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(stderr, bytes.Repeat([]byte("diagnostic tail\n"), 64)) { + t.Fatal("diagnostic tail lost") + } +} diff --git a/supervisor/tree_other.go b/supervisor/tree_other.go new file mode 100644 index 0000000..741e715 --- /dev/null +++ b/supervisor/tree_other.go @@ -0,0 +1,21 @@ +//go:build !linux && !darwin && !windows + +package supervisor + +import ( + "fmt" + "os" + "os/exec" + "runtime" +) + +type processTree struct{} + +func unsupported() error { + return fmt.Errorf("runtime process supervision is unsupported on %s", runtime.GOOS) +} +func configurePipes(_ *exec.Cmd, _, _ *os.File) error { return unsupported() } +func preventLeaseInheritance(_ *os.File) {} +func startTree(configuration) (*processTree, error) { return nil, unsupported() } +func (*processTree) Wait() error { return unsupported() } +func (*processTree) Kill() error { return unsupported() } diff --git a/supervisor/tree_unix.go b/supervisor/tree_unix.go new file mode 100644 index 0000000..352c9c0 --- /dev/null +++ b/supervisor/tree_unix.go @@ -0,0 +1,71 @@ +//go:build linux || darwin + +package supervisor + +import ( + "errors" + "os" + "os/exec" + "sync" + "syscall" + + "golang.org/x/sys/unix" +) + +type processTree struct { + mu sync.Mutex + cmd *exec.Cmd + observeExit func() error + reaped bool +} + +func configurePipes(cmd *exec.Cmd, lease, config *os.File) error { + cmd.ExtraFiles = []*os.File{lease, config} + cmd.Args = append(cmd.Args, "--lease=3", "--config=4") + cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} + return nil +} + +func preventLeaseInheritance(file *os.File) { unix.CloseOnExec(int(file.Fd())) } + +func startTree(cfg configuration) (*processTree, error) { + if err := adoptDescendants(); err != nil { + return nil, err + } + cmd := runtimeCommand(cfg) + cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} + if err := cmd.Start(); err != nil { + return nil, err + } + watch, err := watchExit(cmd.Process.Pid) + if err != nil { + _ = unix.Kill(-cmd.Process.Pid, unix.SIGKILL) + _ = cmd.Wait() + return nil, err + } + return &processTree{cmd: cmd, observeExit: watch}, nil +} + +func (t *processTree) Wait() error { + observed := t.observeExit() + t.mu.Lock() + defer t.mu.Unlock() + // Observe exit without reaping: the leader PID remains reserved until group termination. + killed := t.killGroup() + waited := t.cmd.Wait() + t.reaped = true + return errors.Join(observed, killed, waited, reapDescendants()) +} + +func (t *processTree) Kill() error { + t.mu.Lock() + defer t.mu.Unlock() + if t.reaped { + return nil + } + return t.killGroup() +} + +func (t *processTree) killGroup() error { + return terminateGroup(t.cmd.Process.Pid) +} diff --git a/supervisor/tree_windows.go b/supervisor/tree_windows.go new file mode 100644 index 0000000..fe5176b --- /dev/null +++ b/supervisor/tree_windows.go @@ -0,0 +1,79 @@ +package supervisor + +import ( + "os" + "os/exec" + "runtime" + "strconv" + "syscall" + "unsafe" + + "golang.org/x/sys/windows" +) + +type processTree struct { + cmd *exec.Cmd + job windows.Handle +} + +func configurePipes(cmd *exec.Cmd, lease, config *os.File) error { + handles := []syscall.Handle{syscall.Handle(lease.Fd()), syscall.Handle(config.Fd())} + for _, handle := range handles { + if err := windows.SetHandleInformation( + windows.Handle(handle), + windows.HANDLE_FLAG_INHERIT, + windows.HANDLE_FLAG_INHERIT, + ); err != nil { + return err + } + } + cmd.SysProcAttr = &syscall.SysProcAttr{AdditionalInheritedHandles: handles} + cmd.Args = append( + cmd.Args, + "--lease="+strconv.FormatUint(uint64(lease.Fd()), 10), + "--config="+strconv.FormatUint(uint64(config.Fd()), 10), + ) + return nil +} + +func preventLeaseInheritance(file *os.File) { + // exec.Cmd also supplies an explicit handle list; neither lease nor job reaches the plugin. + _ = windows.SetHandleInformation(windows.Handle(file.Fd()), windows.HANDLE_FLAG_INHERIT, 0) +} + +func startTree(cfg configuration) (*processTree, error) { + job, err := windows.CreateJobObject(nil, nil) + if err != nil { + return nil, err + } + limits := windows.JOBOBJECT_EXTENDED_LIMIT_INFORMATION{} + limits.BasicLimitInformation.LimitFlags = windows.JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE + var pinned runtime.Pinner + pinned.Pin(&limits) + defer pinned.Unpin() + _, err = windows.SetInformationJobObject( + job, + windows.JobObjectExtendedLimitInformation, + // #nosec G103 -- Pinned, fixed-size Win32 structure; the native API requires this pointer. + uintptr(unsafe.Pointer(&limits)), + uint32(unsafe.Sizeof(limits)), + ) + if err != nil { + _ = windows.CloseHandle(job) + return nil, err + } + // Join before creating the plugin: CreateProcess descendants inherit membership atomically. + if err := windows.AssignProcessToJobObject(job, windows.CurrentProcess()); err != nil { + _ = windows.CloseHandle(job) + return nil, err + } + // This handle is intentionally held until Main exits. Closing it here would kill this supervisor too. + cmd := runtimeCommand(cfg) + if err := cmd.Start(); err != nil { + return nil, err + } + return &processTree{cmd: cmd, job: job}, nil +} + +func (t *processTree) Wait() error { return t.cmd.Wait() } +func (t *processTree) Kill() error { return windows.TerminateJobObject(t.job, 0) } From b3148ddf27f7e3ce649ec5d7f1c01d190ae7f3c2 Mon Sep 17 00:00:00 2001 From: Samuel K <69881238+skevetter@users.noreply.github.com> Date: Sat, 3 Oct 2026 10:00:55 -0600 Subject: [PATCH 2/3] fix: preserve Unix bytes in supervisor configuration --- supervisor/main.go | 4 ++-- supervisor/runner.go | 4 ++-- supervisor/runner_test.go | 41 +++++++++++++++++++++++++++++++++------ 3 files changed, 39 insertions(+), 10 deletions(-) diff --git a/supervisor/main.go b/supervisor/main.go index ac588f1..94e220e 100644 --- a/supervisor/main.go +++ b/supervisor/main.go @@ -2,7 +2,7 @@ package supervisor import ( "context" - "encoding/json" + "encoding/gob" "errors" "flag" "fmt" @@ -97,7 +97,7 @@ func stopTree(cfg configuration, tree *processTree, exited <-chan error) error { func readConfiguration(input *os.File) (configuration, error) { var cfg configuration - decoder := json.NewDecoder(io.LimitReader(input, 1<<20)) + decoder := gob.NewDecoder(io.LimitReader(input, 1<<20)) if err := decoder.Decode(&cfg); err != nil { return cfg, errors.New("invalid runtime supervisor configuration") } diff --git a/supervisor/runner.go b/supervisor/runner.go index 5f74199..ef62908 100644 --- a/supervisor/runner.go +++ b/supervisor/runner.go @@ -3,7 +3,7 @@ package supervisor import ( "context" - "encoding/json" + "encoding/gob" "errors" "fmt" "io" @@ -133,7 +133,7 @@ func (r *ownedRunner) Start(ctx context.Context) error { go func() { r.waitErr = r.cmd.Wait(); r.closeLease(); close(r.done) }() encoded := make(chan error, 1) go func() { - encoded <- json.NewEncoder(r.configWrite).Encode(r.config) + encoded <- gob.NewEncoder(r.configWrite).Encode(r.config) _ = r.configWrite.Close() }() select { diff --git a/supervisor/runner_test.go b/supervisor/runner_test.go index f8cabe8..676cfc5 100644 --- a/supervisor/runner_test.go +++ b/supervisor/runner_test.go @@ -8,13 +8,15 @@ import ( "os" "os/exec" "path/filepath" + goruntime "runtime" "testing" "time" ) type inspection struct { - Args []string - Env, Directory string + Args [][]byte + Env []byte + Directory string } func TestMain(m *testing.M) { @@ -28,7 +30,11 @@ func TestMain(m *testing.M) { panic(err) } err = json.NewEncoder(os.Stdout). - Encode(inspection{Args: os.Args[2:], Env: os.Getenv("DEVSY_OWNERSHIP_TEST"), Directory: directory}) + Encode(inspection{ + Args: argumentBytes(os.Args[2:]), + Env: []byte(os.Getenv("DEVSY_OWNERSHIP_TEST")), + Directory: directory, + }) if err != nil { panic(err) } @@ -61,8 +67,8 @@ func fixtureOptions(t *testing.T) Options { } return Options{ SupervisorBinary: binary, SupervisorArgs: []string{"--helper-role=supervisor"}, - RuntimeBinary: binary, Args: []string{"--helper-role=inspect", "space λ argument"}, - Env: []string{"DEVSY_OWNERSHIP_TEST=inherited override"}, Directory: t.TempDir(), + RuntimeBinary: binary, Args: []string{"--helper-role=inspect", argumentValue()}, + Env: []string{"DEVSY_OWNERSHIP_TEST=" + environmentValue()}, Directory: t.TempDir(), } } @@ -144,7 +150,8 @@ func TestPreCanceledStartClosesResources(t *testing.T) { func assertInspection(t *testing.T, got inspection, directory string) { t.Helper() - if len(got.Args) != 1 || got.Args[0] != "space λ argument" || got.Env != "inherited override" { + if len(got.Args) != 1 || string(got.Args[0]) != argumentValue() || + !bytes.Equal(got.Env, []byte(environmentValue())) { t.Fatalf("configuration changed: %+v", got) } expected, err := os.Stat(directory) @@ -167,3 +174,25 @@ func assertDiagnosticTail(t *testing.T, input io.Reader) { t.Fatal("diagnostic tail lost") } } + +func argumentBytes(args []string) [][]byte { + data := make([][]byte, len(args)) + for i, arg := range args { + data[i] = []byte(arg) + } + return data +} + +func environmentValue() string { + if goruntime.GOOS == "windows" { + return "inherited override" + } + return "inherited override\xff" +} + +func argumentValue() string { + if goruntime.GOOS == "windows" { + return "space λ argument" + } + return "space λ argument\xff" +} From 293c64be77ed9b0ae230456fa85202b60f6082dc Mon Sep 17 00:00:00 2001 From: Samuel K <69881238+skevetter@users.noreply.github.com> Date: Sat, 3 Oct 2026 10:06:37 -0600 Subject: [PATCH 3/3] fix: retain client-assigned broker and socket metadata --- internal/processprobe/owned_test.go | 12 +++++++----- supervisor/runner.go | 2 ++ 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/internal/processprobe/owned_test.go b/internal/processprobe/owned_test.go index c4c0d78..dbfab7d 100644 --- a/internal/processprobe/owned_test.go +++ b/internal/processprobe/owned_test.go @@ -47,14 +47,16 @@ func owned(t *testing.T, o *observer, delayed bool) *hplugin.Client { fixtureEnvironment(o.directory), sdkplugin.Handshake().MagicCookieKey+"=stale-cookie", "PLUGIN_PROTOCOL_VERSIONS=99", + "PLUGIN_MULTIPLEX_GRPC=false", "PLUGIN_UNIX_SOCKET_DIR="+filepath.Join(o.directory, "stale"), ), }), - StartTimeout: timeout, - Logger: hclog.NewNullLogger(), - Stderr: os.Stderr, - SyncStderr: os.Stderr, - UnixSocketConfig: &hplugin.UnixSocketConfig{TempDir: o.directory}, + GRPCBrokerMultiplex: true, + StartTimeout: timeout, + Logger: hclog.NewNullLogger(), + Stderr: os.Stderr, + SyncStderr: os.Stderr, + UnixSocketConfig: &hplugin.UnixSocketConfig{TempDir: o.directory}, }) t.Cleanup(client.Kill) return client diff --git a/supervisor/runner.go b/supervisor/runner.go index ef62908..df4c4b5 100644 --- a/supervisor/runner.go +++ b/supervisor/runner.go @@ -270,6 +270,8 @@ func runtimeEnvironment(base, overrides []string) []string { "PLUGIN_MAX_PORT", "PLUGIN_PROTOCOL_VERSIONS", "PLUGIN_CLIENT_CERT", + "PLUGIN_MULTIPLEX_GRPC", + "PLUGIN_UNIX_SOCKET_GROUP", "PLUGIN_UNIX_SOCKET_DIR": // Client-assigned transport metadata must survive inherited or provider environment overrides. env = append(env, value)