From b72535eb305605ba7130817fb024cbc2377071e9 Mon Sep 17 00:00:00 2001 From: chris lee Date: Thu, 8 Oct 2026 15:04:00 -0400 Subject: [PATCH 01/15] Add Darwin transport and OS operations to shared guest service --- docs/macos-guest-agent.md | 76 ++++++++ lib/system/guest_agent/darwin_test.go | 60 ++++++ lib/system/guest_agent/exec.go | 113 +++++------- lib/system/guest_agent/exec_test.go | 4 +- lib/system/guest_agent/main.go | 11 +- lib/system/guest_agent/network.go | 2 + lib/system/guest_agent/network_darwin.go | 15 ++ lib/system/guest_agent/service_test.go | 172 ++++++++++++++++++ lib/system/guest_agent/shutdown.go | 2 + lib/system/guest_agent/shutdown_darwin.go | 37 ++++ lib/system/guest_agent/vsock_darwin.go | 113 ++++++++++++ lib/system/guest_agent/vsock_darwin_nocgo.go | 14 ++ .../guest_agent/vsock_darwin_nocgo_test.go | 15 ++ lib/system/guest_agent/vsock_darwin_test.go | 35 ++++ lib/system/guest_agent/vsock_linux.go | 15 ++ 15 files changed, 610 insertions(+), 74 deletions(-) create mode 100644 docs/macos-guest-agent.md create mode 100644 lib/system/guest_agent/darwin_test.go create mode 100644 lib/system/guest_agent/network_darwin.go create mode 100644 lib/system/guest_agent/service_test.go create mode 100644 lib/system/guest_agent/shutdown_darwin.go create mode 100644 lib/system/guest_agent/vsock_darwin.go create mode 100644 lib/system/guest_agent/vsock_darwin_nocgo.go create mode 100644 lib/system/guest_agent/vsock_darwin_nocgo_test.go create mode 100644 lib/system/guest_agent/vsock_darwin_test.go create mode 100644 lib/system/guest_agent/vsock_linux.go diff --git a/docs/macos-guest-agent.md b/docs/macos-guest-agent.md new file mode 100644 index 000000000..6e552c0af --- /dev/null +++ b/docs/macos-guest-agent.md @@ -0,0 +1,76 @@ +# Experimental Darwin GuestService + +The existing guest-agent executable now builds for macOS and serves the same +`guest.GuestService` gRPC contract as Linux on vsock port **2222**. It reuses +exec, copy-to/from-guest, and stat implementations; this is not the browser +prototype's custom HTTP protocol or its ports. + +## Build + +Build on a macOS development machine with the macOS SDK and cgo enabled: + +```sh +CGO_ENABLED=1 GOOS=darwin GOARCH=arm64 go build \ + -o guest-agent-darwin-arm64 ./lib/system/guest_agent +``` + +Darwin native AF_VSOCK requires cgo. A cgo-disabled Darwin executable builds but +refuses to listen with an explicit error. Linux keeps its existing Go vsock +listener and cgo-disabled build path. The Darwin listener accepts only host +CID 2, marks descriptors close-on-exec, and provides net.Conn deadlines for +gRPC. Do not start the executable on the development host to test guest services. + +## Provisioning boundary + +Install the binary **inside a stopped-template provisioning guest**, not on the +host. System operations need a root LaunchDaemon. Use a root-owned, non-writable +binary location and launchd configuration; provision permissions and logs +explicitly. The default readiness file is `/var/run/hypeman/guest-agent-ready` +(`HYPEMAN_AGENT_READY_FILE` can override it). + +A listening system agent does not establish autologin, desktop readiness, TCC +permissions, or browser readiness. Root and user desktop agents need a reviewed +handoff and explicit session selection before desktop execution is supported. +No desktop agent installer or public session-selector extension is introduced +by this initial patch. + +The host API and hypervisor still enforce authorization. The vsock host-CID +check is transport admission, not a replacement for instance authority checks. +Do not expose this privileged service through unauthenticated host forwarding. + +## OS-specific operations + +- **Shutdown:** root-only `/sbin/shutdown -h now`, not a signal to launchd/PID 1. + Signal 0 (default) and SIGTERM mean orderly shutdown; other signals are rejected. + Permission, cancellation, and command errors are reported. An RPC response is + not proof the VMM has exited: the host must wait for teardown. +- **Network identity reconfiguration:** returns gRPC `Unimplemented` on Darwin. + Current VZ NAT uses guest DHCP; no static-address, MAC-rekey, ingress, or network + policy parity is claimed. +- **GPU status:** Linux NVIDIA initialization reporting is not macOS graphics + readiness. A Darwin guest without that device reports the existing unknown + state. + +## Validation and remaining integration + +In-process gRPC tests exercise exec stdout/stderr/exit/env, disconnect +cancellation, and file copy/stat round-trips. Darwin-specific tests exercise +shutdown policy without executing a real shutdown, explicit network rejection, +and transport deadline errors using a local socket pair. These do not prove a +live guest AF_VSOCK handshake for this executable. + +Remaining draft gates: + +- Provision in a test guest and exercise real host GuestService connectivity. +- Connect normal API exec/files, system readiness, graceful stop/recovery, and + guest-agent version compatibility; do not silently redefine `Running`. +- Root/desktop session authorization and image provisioning. +- Broader backpressure/large-output validation of bounded non-TTY streaming, + PTY/disconnect and descendant-process cleanup, transfer failure/size handling, + and privilege/logging security review. +- Linux test execution on an appropriate runner, and independent authenticated + review. + +The live macOS benchmark guest and its prototype agent are unchanged by this +source patch. Native gRPC integration remains experimental until those gates +are satisfied. diff --git a/lib/system/guest_agent/darwin_test.go b/lib/system/guest_agent/darwin_test.go new file mode 100644 index 000000000..d813c2a98 --- /dev/null +++ b/lib/system/guest_agent/darwin_test.go @@ -0,0 +1,60 @@ +//go:build darwin + +package main + +import ( + "context" + "errors" + "syscall" + "testing" + + pb "github.com/kernel/hypeman/lib/guest" + "github.com/stretchr/testify/require" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +func TestDarwinShutdownPolicy(t *testing.T) { + for _, tc := range []struct { + name string + signal int32 + euid int + canceled bool + commandError error + code codes.Code + called bool + }{ + {name: "default orderly shutdown", called: true}, + {name: "explicit orderly shutdown", signal: int32(syscall.SIGTERM), called: true}, + {name: "reject arbitrary signal", signal: int32(syscall.SIGKILL), code: codes.InvalidArgument}, + {name: "reject desktop agent", euid: 501, code: codes.PermissionDenied}, + {name: "canceled before command", canceled: true, code: codes.Canceled}, + {name: "command failure", commandError: errors.New("shutdown refused"), code: codes.Internal, called: true}, + } { + t.Run(tc.name, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + if tc.canceled { + cancel() + } + called := false + response, err := requestDarwinShutdown(ctx, &pb.ShutdownRequest{Signal: tc.signal}, tc.euid, func(context.Context) error { + called = true + return tc.commandError + }) + require.Equal(t, tc.called, called) + require.Equal(t, tc.code, status.Code(err)) + if tc.code == codes.OK { + require.NotNil(t, response) + } else { + require.Nil(t, response) + } + }) + } +} + +func TestDarwinNetworkReconfigurationUnsupported(t *testing.T) { + response, err := (&guestServer{}).ReconfigureNetwork(context.Background(), &pb.ReconfigureNetworkRequest{}) + require.Nil(t, response) + require.Equal(t, codes.Unimplemented, status.Code(err)) +} diff --git a/lib/system/guest_agent/exec.go b/lib/system/guest_agent/exec.go index 409754c65..98f9175e1 100644 --- a/lib/system/guest_agent/exec.go +++ b/lib/system/guest_agent/exec.go @@ -3,7 +3,6 @@ package main import ( "context" "fmt" - "io" "log" "os" "os/exec" @@ -30,16 +29,16 @@ func (s *guestServer) Exec(stream pb.GuestService_ExecServer) error { return fmt.Errorf("first message must be ExecStart") } - command := start.Command - if len(command) == 0 { - command = []string{"/bin/sh"} + if len(start.Command) == 0 { + start.Command = []string{"/bin/sh"} } + command := start.Command log.Printf("[guest-agent] exec: command=%v tty=%v cwd=%s timeout=%d", command, start.Tty, start.Cwd, start.TimeoutSeconds) // Create context with timeout if specified - ctx := context.Background() + ctx := stream.Context() if start.TimeoutSeconds > 0 { var cancel context.CancelFunc ctx, cancel = context.WithTimeout(ctx, time.Duration(start.TimeoutSeconds)*time.Second) @@ -59,31 +58,24 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E return fmt.Errorf("empty command") } - cmd := exec.CommandContext(ctx, start.Command[0], start.Command[1:]...) - - // Set up environment (no TTY defaults for non-TTY mode) + outputCtx, cancel := context.WithCancel(ctx) + defer cancel() + // CommandContext must observe both caller cancellation and stream-send errors. + cmd := exec.CommandContext(outputCtx, start.Command[0], start.Command[1:]...) cmd.Env = s.buildEnv(start.Env, false) - - // Set up working directory - if start.Cwd != "" { - cmd.Dir = start.Cwd + cmd.Dir = start.Cwd + cmd.WaitDelay = 2 * time.Second + var sendMu sync.Mutex + cmd.Stdout = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel} + cmd.Stderr = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel, stderr: true} + stdin, err := cmd.StdinPipe() + if err != nil { + return fmt.Errorf("open command stdin: %w", err) } - - stdin, _ := cmd.StdinPipe() - stdout, _ := cmd.StdoutPipe() - stderr, _ := cmd.StderrPipe() - if err := cmd.Start(); err != nil { return fmt.Errorf("start command: %w", err) } - // Mutex to protect concurrent stream.Send calls (gRPC streams are not thread-safe) - var sendMu sync.Mutex - - // Use WaitGroup to ensure all output is read before sending - var wg sync.WaitGroup - var stdoutData, stderrData []byte - // Handle stdin in background go func() { defer stdin.Close() @@ -98,50 +90,12 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E } }() - // Read all stdout/stderr BEFORE calling Wait() - Wait() closes the pipes! - wg.Add(1) - go func() { - defer wg.Done() - data, _ := io.ReadAll(stdout) - stdoutData = data - }() - - wg.Add(1) - go func() { - defer wg.Done() - data, _ := io.ReadAll(stderr) - stderrData = data - }() - - // Wait for all reads to complete FIRST (before Wait closes pipes) - wg.Wait() - - // Now safe to call Wait - pipes are fully drained + // os/exec drains stdout/stderr through bounded writers before Wait returns. waitErr := cmd.Wait() - - // Now stream output in chunks (streaming compatible) - const chunkSize = 32 * 1024 - for i := 0; i < len(stdoutData); i += chunkSize { - end := i + chunkSize - if end > len(stdoutData) { - end = len(stdoutData) - } - sendMu.Lock() - stream.Send(&pb.ExecResponse{ - Response: &pb.ExecResponse_Stdout{Stdout: stdoutData[i:end]}, - }) - sendMu.Unlock() - } - for i := 0; i < len(stderrData); i += chunkSize { - end := i + chunkSize - if end > len(stderrData) { - end = len(stderrData) + if waitErr != nil { + if _, exited := waitErr.(*exec.ExitError); !exited { + return fmt.Errorf("stream command output: %w", waitErr) } - sendMu.Lock() - stream.Send(&pb.ExecResponse{ - Response: &pb.ExecResponse_Stderr{Stderr: stderrData[i:end]}, - }) - sendMu.Unlock() } exitCode := int32(0) @@ -160,6 +114,33 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E }) } +type execStreamWriter struct { + stream pb.GuestService_ExecServer + mu *sync.Mutex + cancel context.CancelFunc + stderr bool +} + +func (w *execStreamWriter) Write(data []byte) (int, error) { + w.mu.Lock() + defer w.mu.Unlock() + written := 0 + for len(data) > 0 { + n := min(len(data), 32*1024) + response := &pb.ExecResponse{Response: &pb.ExecResponse_Stdout{Stdout: data[:n]}} + if w.stderr { + response.Response = &pb.ExecResponse_Stderr{Stderr: data[:n]} + } + if err := w.stream.Send(response); err != nil { + w.cancel() + return written, err + } + written += n + data = data[n:] + } + return written, nil +} + // executeTTY executes command with TTY func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_ExecServer, start *pb.ExecStart) error { // Run command directly with PTY - guest-agent is already running in container namespace diff --git a/lib/system/guest_agent/exec_test.go b/lib/system/guest_agent/exec_test.go index ca3440536..887a72d7b 100644 --- a/lib/system/guest_agent/exec_test.go +++ b/lib/system/guest_agent/exec_test.go @@ -16,9 +16,9 @@ func TestBuildEnv(t *testing.T) { }) t.Run("non-TTY session does not add xterm-256color", func(t *testing.T) { + t.Setenv("TERM", "guest-agent-test") env := s.buildEnv(nil, false) - // Non-TTY should not add our default TERM - // (host environment TERM may still be present, that's fine) + assert.Contains(t, env, "TERM=guest-agent-test", "non-TTY preserves the inherited TERM") assert.NotContains(t, env, "TERM=xterm-256color", "non-TTY should not add xterm-256color default") }) diff --git a/lib/system/guest_agent/main.go b/lib/system/guest_agent/main.go index cc8f5d049..33db3e889 100644 --- a/lib/system/guest_agent/main.go +++ b/lib/system/guest_agent/main.go @@ -3,20 +3,19 @@ package main import ( "fmt" "log" + "net" "os" "path/filepath" "strconv" "time" pb "github.com/kernel/hypeman/lib/guest" - "github.com/mdlayher/vsock" "google.golang.org/grpc" ) const ( - readySentinelPrefix = "HYPEMAN-AGENT-READY" - defaultReadyFilePath = "/run/hypeman/guest-agent-ready" - readyFDEnv = "HYPEMAN_AGENT_READY_FD" + readySentinelPrefix = "HYPEMAN-AGENT-READY" + readyFDEnv = "HYPEMAN_AGENT_READY_FD" ) // guestServer implements the gRPC GuestService @@ -27,11 +26,11 @@ type guestServer struct { func main() { // Listen on vsock port 2222 with retries - var l *vsock.Listener + var l net.Listener var err error for i := 0; i < 10; i++ { - l, err = vsock.Listen(2222, nil) + l, err = listenVsock(2222) if err == nil { break } diff --git a/lib/system/guest_agent/network.go b/lib/system/guest_agent/network.go index 66b71bad8..3b01f5e53 100644 --- a/lib/system/guest_agent/network.go +++ b/lib/system/guest_agent/network.go @@ -1,3 +1,5 @@ +//go:build linux + package main import ( diff --git a/lib/system/guest_agent/network_darwin.go b/lib/system/guest_agent/network_darwin.go new file mode 100644 index 000000000..38a2370ce --- /dev/null +++ b/lib/system/guest_agent/network_darwin.go @@ -0,0 +1,15 @@ +//go:build darwin + +package main + +import ( + "context" + + pb "github.com/kernel/hypeman/lib/guest" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +func (s *guestServer) ReconfigureNetwork(context.Context, *pb.ReconfigureNetworkRequest) (*pb.ReconfigureNetworkResponse, error) { + return nil, status.Error(codes.Unimplemented, "Darwin network identity reconfiguration is not supported; VZ NAT uses guest DHCP") +} diff --git a/lib/system/guest_agent/service_test.go b/lib/system/guest_agent/service_test.go new file mode 100644 index 000000000..9171ab672 --- /dev/null +++ b/lib/system/guest_agent/service_test.go @@ -0,0 +1,172 @@ +package main + +import ( + "context" + "io" + "net" + "os" + "path/filepath" + "strconv" + "strings" + "syscall" + "testing" + "time" + + pb "github.com/kernel/hypeman/lib/guest" + "github.com/stretchr/testify/require" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/test/bufconn" +) + +// Exercise the shared RPC contract without starting a VM or listening on vsock. +func testGuestClient(t *testing.T) pb.GuestServiceClient { + t.Helper() + listener := bufconn.Listen(64 << 10) + server := grpc.NewServer() + pb.RegisterGuestServiceServer(server, &guestServer{}) + go server.Serve(listener) + t.Cleanup(server.Stop) + t.Cleanup(func() { listener.Close() }) + conn, err := grpc.NewClient("passthrough:///guest-test", grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithContextDialer(func(context.Context, string) (net.Conn, error) { + return listener.Dial() + })) + require.NoError(t, err) + t.Cleanup(func() { conn.Close() }) + return pb.NewGuestServiceClient(conn) +} + +func TestGuestServiceExecRoundTrip(t *testing.T) { + client := testGuestClient(t) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + stream, err := client.Exec(ctx) + require.NoError(t, err) + require.NoError(t, stream.Send(&pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", "printf '%s' \"$HYPEMAN_EXEC_TEST\"; printf 'stderr' >&2; exit 7"}, + Env: map[string]string{"HYPEMAN_EXEC_TEST": "stdout"}, + }}})) + require.NoError(t, stream.CloseSend()) + var stdout, stderr strings.Builder + var exitCode *int32 + for { + response, err := stream.Recv() + if err == io.EOF { + break + } + require.NoError(t, err) + switch value := response.Response.(type) { + case *pb.ExecResponse_Stdout: + stdout.Write(value.Stdout) + case *pb.ExecResponse_Stderr: + stderr.Write(value.Stderr) + case *pb.ExecResponse_ExitCode: + code := value.ExitCode + exitCode = &code + } + } + require.Equal(t, "stdout", stdout.String()) + require.Equal(t, "stderr", stderr.String()) + require.NotNil(t, exitCode) + require.Equal(t, int32(7), *exitCode) +} + +func TestGuestServiceExecStreamsBeforeExit(t *testing.T) { + client := testGuestClient(t) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + stream, err := client.Exec(ctx) + require.NoError(t, err) + require.NoError(t, stream.Send(&pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", "printf ready; read input; printf '%s' \"$input\""}, + }}})) + first, err := stream.Recv() + require.NoError(t, err, "output must arrive while the command is still waiting for stdin") + require.Equal(t, "ready", string(first.GetStdout())) + require.NoError(t, stream.Send(&pb.ExecRequest{Request: &pb.ExecRequest_Stdin{Stdin: []byte("finish\n")}})) + require.NoError(t, stream.CloseSend()) + var output strings.Builder + exited := false + for { + response, err := stream.Recv() + if err == io.EOF { + break + } + require.NoError(t, err) + output.Write(response.GetStdout()) + if code, ok := response.Response.(*pb.ExecResponse_ExitCode); ok { + require.Equal(t, int32(0), code.ExitCode) + exited = true + } + } + require.True(t, exited) + require.Equal(t, "finish", output.String()) +} + +func TestGuestServiceExecCancellation(t *testing.T) { + client := testGuestClient(t) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + pidPath := filepath.Join(t.TempDir(), "pid") + stream, err := client.Exec(ctx) + require.NoError(t, err) + require.NoError(t, stream.Send(&pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", "printf '%s' \"$$\" > \"$HYPEMAN_PID_FILE\"; exec /bin/sleep 30"}, + Env: map[string]string{"HYPEMAN_PID_FILE": pidPath}, + }}})) + require.NoError(t, stream.CloseSend()) + var pid int + require.Eventually(t, func() bool { + data, err := os.ReadFile(pidPath) + if err != nil { + return false + } + pid, err = strconv.Atoi(string(data)) + return err == nil && pid > 0 + }, 5*time.Second, 10*time.Millisecond) + defer syscall.Kill(pid, syscall.SIGKILL) + cancel() + require.Eventually(t, func() bool { + return syscall.Kill(pid, 0) == syscall.ESRCH + }, 5*time.Second, 10*time.Millisecond, "disconnect must cancel the command, not leave it running") +} + +func TestGuestServiceFileRoundTrip(t *testing.T) { + client := testGuestClient(t) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + path := filepath.Join(t.TempDir(), "copied") + payload := "shared guest service file payload" + stream, err := client.CopyToGuest(ctx) + require.NoError(t, err) + require.NoError(t, stream.Send(&pb.CopyToGuestRequest{Request: &pb.CopyToGuestRequest_Start{Start: &pb.CopyToGuestStart{Path: path, Mode: 0600, Size: int64(len(payload))}}})) + require.NoError(t, stream.Send(&pb.CopyToGuestRequest{Request: &pb.CopyToGuestRequest_Data{Data: []byte(payload)}})) + require.NoError(t, stream.Send(&pb.CopyToGuestRequest{Request: &pb.CopyToGuestRequest_End{End: &pb.CopyToGuestEnd{}}})) + result, err := stream.CloseAndRecv() + require.NoError(t, err) + require.True(t, result.Success, result.Error) + require.Equal(t, int64(len(payload)), result.BytesWritten) + info, err := client.StatPath(ctx, &pb.StatPathRequest{Path: path}) + require.NoError(t, err) + require.True(t, info.Exists && info.IsFile) + require.Equal(t, uint32(0600), info.Mode&0777) + require.Equal(t, int64(len(payload)), info.Size) + read, err := client.CopyFromGuest(ctx, &pb.CopyFromGuestRequest{Path: path}) + require.NoError(t, err) + var content strings.Builder + final := false + for { + response, err := read.Recv() + if err == io.EOF { + break + } + require.NoError(t, err) + require.Nil(t, response.GetError()) + content.Write(response.GetData()) + if response.GetEnd() != nil && response.GetEnd().Final { + final = true + } + } + require.True(t, final) + require.Equal(t, payload, content.String()) +} diff --git a/lib/system/guest_agent/shutdown.go b/lib/system/guest_agent/shutdown.go index e29690aef..5eed220cb 100644 --- a/lib/system/guest_agent/shutdown.go +++ b/lib/system/guest_agent/shutdown.go @@ -1,3 +1,5 @@ +//go:build linux + package main import ( diff --git a/lib/system/guest_agent/shutdown_darwin.go b/lib/system/guest_agent/shutdown_darwin.go new file mode 100644 index 000000000..84a5ff99d --- /dev/null +++ b/lib/system/guest_agent/shutdown_darwin.go @@ -0,0 +1,37 @@ +//go:build darwin + +package main + +import ( + "context" + "os" + "os/exec" + "syscall" + + pb "github.com/kernel/hypeman/lib/guest" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// Darwin's launchd is not Linux init: sending it SIGTERM is not a shutdown API. +func (s *guestServer) Shutdown(ctx context.Context, req *pb.ShutdownRequest) (*pb.ShutdownResponse, error) { + return requestDarwinShutdown(ctx, req, os.Geteuid(), func(ctx context.Context) error { + return exec.CommandContext(ctx, "/sbin/shutdown", "-h", "now").Run() + }) +} + +func requestDarwinShutdown(ctx context.Context, req *pb.ShutdownRequest, euid int, run func(context.Context) error) (*pb.ShutdownResponse, error) { + if req.Signal != 0 && req.Signal != int32(syscall.SIGTERM) { + return nil, status.Error(codes.InvalidArgument, "Darwin supports only an orderly shutdown, not arbitrary init signals") + } + if euid != 0 { + return nil, status.Error(codes.PermissionDenied, "Darwin shutdown requires the system guest agent running as root") + } + if err := ctx.Err(); err != nil { + return nil, status.FromContextError(err).Err() + } + if err := run(ctx); err != nil { + return nil, status.Errorf(codes.Internal, "Darwin shutdown failed: %v", err) + } + return &pb.ShutdownResponse{}, nil +} diff --git a/lib/system/guest_agent/vsock_darwin.go b/lib/system/guest_agent/vsock_darwin.go new file mode 100644 index 000000000..b3c47ce4a --- /dev/null +++ b/lib/system/guest_agent/vsock_darwin.go @@ -0,0 +1,113 @@ +//go:build darwin && cgo + +package main + +/* +#include +#include +#include +#include +#include +static int listen_vm(unsigned port) { + int fd=socket(AF_VSOCK,SOCK_STREAM,0);if(fd<0)return -1; + struct sockaddr_vm a={0};a.svm_len=sizeof(a);a.svm_family=AF_VSOCK;a.svm_cid=VMADDR_CID_ANY;a.svm_port=port; + if(bind(fd,(struct sockaddr*)&a,sizeof(a))<0 || listen(fd,16)<0 || fcntl(fd,F_SETFL,O_NONBLOCK)<0){int e=errno;close(fd);errno=e;return -1;} + if(fcntl(fd,F_SETFD,FD_CLOEXEC)<0){int e=errno;close(fd);errno=e;return -1;}return fd; +} +static int accept_vm(int fd) { + struct sockaddr_vm a={0};socklen_t n=sizeof(a);int c=accept(fd,(struct sockaddr*)&a,&n);if(c<0)return -1; + if(a.svm_cid!=VMADDR_CID_HOST){close(c);errno=EAGAIN;return -1;} + if(fcntl(c,F_SETFL,O_NONBLOCK)<0){int e=errno;close(c);errno=e;return -1;} + if(fcntl(c,F_SETFD,FD_CLOEXEC)<0){int e=errno;close(c);errno=e;return -1;}return c; +} +*/ +import "C" +import ( + "errors" + "fmt" + "net" + "os" + "sync" + "syscall" + "time" +) + +const defaultReadyFilePath = "/var/run/hypeman/guest-agent-ready" + +type vmAddr string + +func (a vmAddr) Network() string { return "vsock" } +func (a vmAddr) String() string { return string(a) } + +// Do not promote os.File's zero-copy methods: AF_VSOCK is not a TCP/file +// descriptor and must use ordinary reads/writes for gRPC transport. +type vmConn struct { + file *os.File + local net.Addr +} + +func (c *vmConn) Read(p []byte) (int, error) { + n, err := c.file.Read(p) + if err != nil { + return n, &net.OpError{Op: "read", Net: "vsock", Err: err} + } + return n, err +} +func (c *vmConn) Write(p []byte) (int, error) { + n, err := c.file.Write(p) + if err != nil { + return n, &net.OpError{Op: "write", Net: "vsock", Err: err} + } + return n, err +} +func (c *vmConn) Close() error { return c.file.Close() } +func (c *vmConn) SetDeadline(t time.Time) error { return c.file.SetDeadline(t) } +func (c *vmConn) SetReadDeadline(t time.Time) error { return c.file.SetReadDeadline(t) } +func (c *vmConn) SetWriteDeadline(t time.Time) error { return c.file.SetWriteDeadline(t) } + +func (c *vmConn) LocalAddr() net.Addr { return c.local } +func (c *vmConn) RemoteAddr() net.Addr { return vmAddr("host:2") } + +type vmListener struct { + mu sync.Mutex + fd C.int + port uint32 + closed bool +} + +func listenVsock(port uint32) (net.Listener, error) { + fd, err := C.listen_vm(C.uint(port)) + if fd < 0 { + return nil, fmt.Errorf("listen vsock: %w", err) + } + return &vmListener{fd: fd, port: port}, nil +} +func (l *vmListener) Addr() net.Addr { return vmAddr(fmt.Sprintf("guest:%d", l.port)) } +func (l *vmListener) Close() error { + l.mu.Lock() + defer l.mu.Unlock() + if l.closed { + return net.ErrClosed + } + l.closed = true + C.close(l.fd) + return nil +} +func (l *vmListener) Accept() (net.Conn, error) { + for { + l.mu.Lock() + if l.closed { + l.mu.Unlock() + return nil, net.ErrClosed + } + fd, err := C.accept_vm(l.fd) + l.mu.Unlock() + if fd >= 0 { + return &vmConn{file: os.NewFile(uintptr(fd), "guest-vsock"), local: l.Addr()}, nil + } + if !errors.Is(err, syscall.EAGAIN) && !errors.Is(err, syscall.EINTR) { + return nil, fmt.Errorf("accept vsock: %w", err) + } + time.Sleep(10 * time.Millisecond) + } +} diff --git a/lib/system/guest_agent/vsock_darwin_nocgo.go b/lib/system/guest_agent/vsock_darwin_nocgo.go new file mode 100644 index 000000000..cd4201845 --- /dev/null +++ b/lib/system/guest_agent/vsock_darwin_nocgo.go @@ -0,0 +1,14 @@ +//go:build darwin && !cgo + +package main + +import ( + "fmt" + "net" +) + +const defaultReadyFilePath = "/var/run/hypeman/guest-agent-ready" + +func listenVsock(uint32) (net.Listener, error) { + return nil, fmt.Errorf("Darwin guest vsock requires a cgo-enabled build") +} diff --git a/lib/system/guest_agent/vsock_darwin_nocgo_test.go b/lib/system/guest_agent/vsock_darwin_nocgo_test.go new file mode 100644 index 000000000..630f8b5d8 --- /dev/null +++ b/lib/system/guest_agent/vsock_darwin_nocgo_test.go @@ -0,0 +1,15 @@ +//go:build darwin && !cgo + +package main + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestDarwinVsockRequiresCGO(t *testing.T) { + listener, err := listenVsock(2222) + require.Nil(t, listener) + require.ErrorContains(t, err, "cgo-enabled build") +} diff --git a/lib/system/guest_agent/vsock_darwin_test.go b/lib/system/guest_agent/vsock_darwin_test.go new file mode 100644 index 000000000..732bf71e8 --- /dev/null +++ b/lib/system/guest_agent/vsock_darwin_test.go @@ -0,0 +1,35 @@ +//go:build darwin && cgo + +package main + +import ( + "errors" + "net" + "os" + "testing" + "time" + + "golang.org/x/sys/unix" +) + +func TestVsockReadTimeoutImplementsNetError(t *testing.T) { + fds, err := unix.Socketpair(unix.AF_UNIX, unix.SOCK_STREAM, 0) + if err != nil { + t.Fatal(err) + } + defer unix.Close(fds[1]) + if err = unix.SetNonblock(fds[0], true); err != nil { + unix.Close(fds[0]) + t.Fatal(err) + } + conn := &vmConn{file: os.NewFile(uintptr(fds[0]), "test-vsock"), local: vmAddr("guest:2222")} + defer conn.Close() + if err = conn.SetReadDeadline(time.Now().Add(-time.Second)); err != nil { + t.Fatal(err) + } + _, err = conn.Read(make([]byte, 1)) + var ne net.Error + if !errors.As(err, &ne) || !ne.Timeout() { + t.Fatalf("transport timeout must implement net.Error: %v", err) + } +} diff --git a/lib/system/guest_agent/vsock_linux.go b/lib/system/guest_agent/vsock_linux.go new file mode 100644 index 000000000..2b821da46 --- /dev/null +++ b/lib/system/guest_agent/vsock_linux.go @@ -0,0 +1,15 @@ +//go:build linux + +package main + +import ( + "net" + + "github.com/mdlayher/vsock" +) + +const defaultReadyFilePath = "/run/hypeman/guest-agent-ready" + +func listenVsock(port uint32) (net.Listener, error) { + return vsock.Listen(port, nil) +} From 5ef2b068697f80e65860d2ce045364d91012a0e4 Mon Sep 17 00:00:00 2001 From: chris lee Date: Thu, 8 Oct 2026 15:35:27 -0400 Subject: [PATCH 02/15] Wire opt-in macOS guest service into API, readiness and shutdown --- cmd/api/api/cp.go | 5 +++ cmd/api/api/exec.go | 16 +++++-- cmd/api/api/exec_disconnect_test.go | 45 +++++++++++++++++++ cmd/api/api/macos_test.go | 30 +++++++++++++ docs/macos-experimental.md | 16 ++++--- docs/macos-guest-agent.md | 25 ++++++++++- lib/guest/client.go | 25 +++++++---- lib/images/types.go | 3 ++ lib/instances/macos.go | 10 ++++- lib/instances/macos_test.go | 69 +++++++++++++++++++++++++++++ lib/instances/query.go | 28 +++++++----- lib/instances/stop.go | 4 +- lib/instances/vsock.go | 4 +- 13 files changed, 244 insertions(+), 36 deletions(-) create mode 100644 cmd/api/api/exec_disconnect_test.go diff --git a/cmd/api/api/cp.go b/cmd/api/api/cp.go index bfa64e45c..ac4812b27 100644 --- a/cmd/api/api/cp.go +++ b/cmd/api/api/cp.go @@ -104,6 +104,11 @@ func (s *ApiService) CpHandler(w http.ResponseWriter, r *http.Request) { return } + if inst.MacOS != nil && (!inst.MacOS.GuestAgent || inst.SkipGuestAgent) { + http.Error(w, `{"code":"unsupported","message":"file copy requires the shared macOS guest agent to be enabled"}`, http.StatusNotImplemented) + return + } + if inst.State != instances.StateRunning { http.Error(w, fmt.Sprintf(`{"code":"invalid_state","message":"instance must be running (current state: %s)"}`, inst.State), http.StatusConflict) return diff --git a/cmd/api/api/exec.go b/cmd/api/api/exec.go index 9313e5372..f1d7ed2fe 100644 --- a/cmd/api/api/exec.go +++ b/cmd/api/api/exec.go @@ -71,8 +71,8 @@ func (s *ApiService) ExecHandler(w http.ResponseWriter, r *http.Request) { return } - if inst.MacOS != nil { - http.Error(w, `{"code":"unsupported","message":"exec is not implemented for experimental macOS instances"}`, http.StatusNotImplemented) + if inst.MacOS != nil && (!inst.MacOS.GuestAgent || inst.SkipGuestAgent) { + http.Error(w, `{"code":"unsupported","message":"exec is not implemented for macOS images without the shared guest agent enabled"}`, http.StatusNotImplemented) return } @@ -128,6 +128,8 @@ func (s *ApiService) ExecHandler(w http.ResponseWriter, r *http.Request) { tracer := otel.Tracer("hypeman/exec") ctx, span := tracer.Start(ctx, "exec.session", trace.WithAttributes(execSpanAttributes(inst.Id, execReq.TTY)...)) defer span.End() + ctx, cancel := context.WithCancel(ctx) + defer cancel() // Audit log: exec session started log.InfoContext(ctx, "exec session started", @@ -146,11 +148,10 @@ func (s *ApiService) ExecHandler(w http.ResponseWriter, r *http.Request) { var resizeChan chan *guest.WindowSize if execReq.TTY { resizeChan = make(chan *guest.WindowSize, 10) - defer close(resizeChan) } // Create WebSocket read/writer wrapper that handles resize messages - wsConn := &wsReadWriter{ws: ws, ctx: ctx, resizeChan: resizeChan} + wsConn := &wsReadWriter{ws: ws, ctx: ctx, resizeChan: resizeChan, cancel: cancel} dialer, err := s.InstanceManager.GetVsockDialer(ctx, inst.Id) if err != nil { @@ -221,6 +222,7 @@ type wsReadWriter struct { reader io.Reader mu sync.Mutex resizeChan chan<- *guest.WindowSize // Channel to send resize events (nil if not TTY) + cancel context.CancelFunc // Exec session cancellation on disconnect (nil for other users). } func (w *wsReadWriter) Read(p []byte) (n int, err error) { @@ -241,6 +243,9 @@ func (w *wsReadWriter) Read(p []byte) (n int, err error) { // Read next WebSocket message messageType, data, err := w.ws.ReadMessage() if err != nil { + if w.cancel != nil { + w.cancel() + } return 0, err } @@ -273,6 +278,9 @@ func (w *wsReadWriter) Read(p []byte) (n int, err error) { func (w *wsReadWriter) Write(p []byte) (n int, err error) { if err := w.ws.WriteMessage(websocket.BinaryMessage, p); err != nil { + if w.cancel != nil { + w.cancel() + } return 0, err } return len(p), nil diff --git a/cmd/api/api/exec_disconnect_test.go b/cmd/api/api/exec_disconnect_test.go new file mode 100644 index 000000000..ad16654f5 --- /dev/null +++ b/cmd/api/api/exec_disconnect_test.go @@ -0,0 +1,45 @@ +package api + +import ( + "context" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/gorilla/websocket" + "github.com/stretchr/testify/require" +) + +func TestExecWebsocketDisconnectCancelsSession(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + readDone := make(chan error, 1) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + readDone <- err + return + } + defer conn.Close() + wrapped := &wsReadWriter{ws: conn, ctx: ctx, cancel: cancel} + _, err = wrapped.Read(make([]byte, 1)) + readDone <- err + })) + defer server.Close() + conn, _, err := websocket.DefaultDialer.Dial("ws"+strings.TrimPrefix(server.URL, "http"), nil) + require.NoError(t, err) + require.NoError(t, conn.Close()) + select { + case <-ctx.Done(): + case <-time.After(5 * time.Second): + t.Fatal("WebSocket disconnect did not cancel exec session") + } + select { + case err := <-readDone: + require.Error(t, err) + case <-time.After(5 * time.Second): + t.Fatal("WebSocket reader did not exit") + } +} diff --git a/cmd/api/api/macos_test.go b/cmd/api/api/macos_test.go index a7de011ac..cc8ced902 100644 --- a/cmd/api/api/macos_test.go +++ b/cmd/api/api/macos_test.go @@ -46,6 +46,36 @@ func TestMacOSStatRejectedBeforeGuestDial(t *testing.T) { require.Equal(t, "unsupported", unsupported.Code) } +func TestMacOSGuestAgentAdmissionBeforeUpgrade(t *testing.T) { + for _, tc := range []struct { + name string + declared, skipped bool + status int + }{ + {"unmanaged", false, false, http.StatusNotImplemented}, + {"disabled", true, true, http.StatusNotImplemented}, + {"enabled but stopped", true, false, http.StatusConflict}, + } { + for _, operation := range []string{"exec", "cp"} { + t.Run(tc.name+"/"+operation, func(t *testing.T) { + s := &ApiService{} + inst := &instances.Instance{StoredMetadata: instances.StoredMetadata{ + MacOS: &images.MacOSImage{GuestAgent: tc.declared}, SkipGuestAgent: tc.skipped, + }, State: instances.StateStopped} + ctx := mw.WithResolvedInstance(context.Background(), "test", inst) + r := httptest.NewRequest(http.MethodGet, "/instances/test/"+operation, nil).WithContext(ctx) + w := httptest.NewRecorder() + if operation == "exec" { + s.ExecHandler(w, r) + } else { + s.CpHandler(w, r) + } + require.Equal(t, tc.status, w.Code) + }) + } + } +} + func TestMacOSSchemaDefersTemplateDefaults(t *testing.T) { spec, err := oapi.GetSwagger() require.NoError(t, err) diff --git a/docs/macos-experimental.md b/docs/macos-experimental.md index 187542f00..f1378645b 100644 --- a/docs/macos-experimental.md +++ b/docs/macos-experimental.md @@ -72,14 +72,19 @@ attached. NAT uses the preserved MAC. The IP is observed from the host's VZ DHCP leases, which may retain an old lease before the guest is ready. **`Running` means the VMM is running, not that SSH, the desktop or Chrome is -ready.** Readiness is currently tested over SSH. Startup does not launch Chrome. +ready.** Unmanaged templates still need an external readiness check. Templates +provisioned with the shared Darwin GuestService can opt in with `guest_agent: true` +in their platform configuration; system readiness is observed separately through +`guest_agent_ready_at`. See [Darwin GuestService](macos-guest-agent.md). Startup +does not launch Chrome. Guest identity, host keys, user secrets and MAC are preserved. Only one active instance with a given Mac identifier is allowed. Do not run the source bundle or an external clone concurrently. Production provisioning needs identity and credential rekeying; this spike is not a multi-tenant image format. -`POST /instances/{id}/stop` shuts down the VMM; without a Darwin guest agent, -it does **not** guarantee an orderly guest OS/application shutdown. Use a human +`POST /instances/{id}/stop` attempts orderly GuestService shutdown for opted-in +images, then falls back to VMM shutdown if necessary. Without an enabled Darwin +guest agent it does **not** guarantee an orderly guest OS/application shutdown. Use a human or an authorized guest shutdown workflow before destructive operations when application consistency matters. Start cold-boots the instance's existing disk. Delete removes instance storage, not the imported image. Deleting the imported @@ -87,8 +92,9 @@ image does not affect existing instances: each owns independent copies and starts without the template. Unsupported instance operations reject requests: snapshot/fork/standby/restore, -updates, volumes, env/commands/credential brokering, Linux guest-agent exec and -vsock operations, health/restart/auto-standby policies, passthrough and shaping. +updates, volumes, startup env/commands/credential brokering, +health/restart/auto-standby policies, passthrough and shaping. Exec/file-copy over +vsock require an opted-in shared agent; unmanaged images reject those requests. The shim's standalone save/restore proof does not make API snapshot semantics safe; consistent disk+aux+state bundle handling remains separate work. diff --git a/docs/macos-guest-agent.md b/docs/macos-guest-agent.md index 6e552c0af..51e5eef67 100644 --- a/docs/macos-guest-agent.md +++ b/docs/macos-guest-agent.md @@ -38,6 +38,27 @@ The host API and hypervisor still enforce authorization. The vsock host-CID check is transport admission, not a replacement for instance authority checks. Do not expose this privileged service through unauthenticated host forwarding. +## Image declaration and host integration + +Set `"guest_agent": true` in the image's experimental macOS platform configuration +only after provisioning the root shared agent on vsock 2222. This field is omitted +by default, so existing local templates remain unmanaged. The normal instance +`skip_guest_agent` option can disable agent integration even for a declared image. + +For declared/enabled instances, the normal exec and file-copy handlers use the +shared GuestService and stop attempts its Shutdown RPC before the existing +forced-stop fallback. An unavailable agent still produces a timeout/failure, +not an assertion that it is ready. Readiness probing records the existing +`guest_agent_ready_at` marker without requiring or inventing a Linux workload +start marker. `Running` remains VMM-running for macOS; it does not certify system, +desktop or browser readiness. This readiness timestamp is a boot observation, +not continuous agent health. Start clears the prior boot's readiness marker. + +Undeclared/disabled images reject exec/copy before WebSocket upgrade and skip +agent shutdown/probes. The declaration is an image capability claim, not a +verified live guest handshake or a security credential. Runtime connectivity +and privileged provisioning still require the validation below. + ## OS-specific operations - **Shutdown:** root-only `/sbin/shutdown -h now`, not a signal to launchd/PID 1. @@ -62,8 +83,8 @@ live guest AF_VSOCK handshake for this executable. Remaining draft gates: - Provision in a test guest and exercise real host GuestService connectivity. -- Connect normal API exec/files, system readiness, graceful stop/recovery, and - guest-agent version compatibility; do not silently redefine `Running`. +- Live normal API exec/files, readiness and graceful stop/recovery validation; + guest-agent version compatibility and bounded readiness wait semantics. - Root/desktop session authorization and image provisioning. - Broader backpressure/large-output validation of bounded non-TTY streaming, PTY/disconnect and descendant-process cleanup, transfer failure/size handling, diff --git a/lib/guest/client.go b/lib/guest/client.go index 61376bf02..38b2b26a5 100644 --- a/lib/guest/client.go +++ b/lib/guest/client.go @@ -442,6 +442,8 @@ func retryableConnectionErrorType(err error) string { // execIntoInstanceOnce executes command in instance via vsock using gRPC (single attempt). func execIntoInstanceOnce(ctx context.Context, dialer hypervisor.VsockDialer, opts ExecOptions) (*ExitStatus, error) { + ctx, cancel := context.WithCancel(ctx) + defer cancel() start := time.Now() var bytesSent int64 @@ -515,14 +517,21 @@ func execIntoInstanceOnce(ctx context.Context, dialer hypervisor.VsockDialer, op // Handle resize events in background (if channel provided) if opts.ResizeChan != nil { go func() { - for resize := range opts.ResizeChan { - streamMu.Lock() - stream.Send(&ExecRequest{ - Request: &ExecRequest_Resize{ - Resize: resize, - }, - }) - streamMu.Unlock() + for { + select { + case <-ctx.Done(): + return + case resize, ok := <-opts.ResizeChan: + if !ok { + return + } + streamMu.Lock() + err := stream.Send(&ExecRequest{Request: &ExecRequest_Resize{Resize: resize}}) + streamMu.Unlock() + if err != nil { + return + } + } } }() } diff --git a/lib/images/types.go b/lib/images/types.go index b2c8ac272..f49437ef9 100644 --- a/lib/images/types.go +++ b/lib/images/types.go @@ -36,6 +36,9 @@ type MacOSImage struct { MAC string `json:"mac"` CPUs uint `json:"cpus"` Memory uint64 `json:"memory"` + // GuestAgent declares a provisioned system GuestService on vsock 2222. + // Readiness is probed separately; old templates remain unmanaged. + GuestAgent bool `json:"guest_agent,omitempty"` } // Validate checks the platform fields every macOS bundle must carry, whether it diff --git a/lib/instances/macos.go b/lib/instances/macos.go index c02eb32d0..cc2be6c02 100644 --- a/lib/instances/macos.go +++ b/lib/instances/macos.go @@ -29,7 +29,7 @@ func prepareMacOSCreate(req *CreateInstanceRequest, img *images.Image, caps hype return fmt.Errorf("%w: incomplete macOS image", ErrImageNotReady) } if req.HotplugSize != 0 || req.OverlaySize != 0 || len(req.Volumes) != 0 || len(req.Devices) != 0 || req.GPU != nil || len(req.Env) != 0 || len(req.Entrypoint) != 0 || len(req.Cmd) != 0 || req.NetworkEgress != nil || len(req.Credentials) != 0 || req.AutoStandby != nil || req.HealthCheck != nil || req.RestartPolicy != nil || req.SnapshotPolicy != nil || req.DiskIOBps != 0 || req.NetworkBandwidthDownload != 0 || req.NetworkBandwidthUpload != 0 { - return fmt.Errorf("%w: experimental macOS supports local disk clone, CPU/RAM, tags, expiration and NAT only; Linux commands/env/volumes, overlays, agents, policies and I/O shaping are unsupported", ErrInvalidRequest) + return fmt.Errorf("%w: experimental macOS supports local disk clone, CPU/RAM, tags, expiration and NAT only; Linux startup commands/env/volumes, overlays, policies and I/O shaping are unsupported", ErrInvalidRequest) } size, vcpus := req.Size, req.Vcpus if size == 0 { @@ -43,7 +43,7 @@ func prepareMacOSCreate(req *CreateInstanceRequest, img *images.Image, caps hype } req.Size, req.Vcpus = size, vcpus req.OverlaySize = *img.SizeBytes // Reserve the writable boot disk, not a Linux overlay. - req.SkipGuestAgent = true + req.SkipGuestAgent = req.SkipGuestAgent || !img.MacOS.GuestAgent req.SkipKernelHeaders = true return nil } @@ -100,6 +100,12 @@ func (m *manager) checkMacOSIdentityAvailable(ctx context.Context, stored *Store return nil } +// guestAgentEnabled keeps old macOS templates unmanaged even if older or +// hand-written metadata did not set SkipGuestAgent explicitly. +func guestAgentEnabled(stored *StoredMetadata) bool { + return !stored.SkipGuestAgent && (stored.MacOS == nil || stored.MacOS.GuestAgent) +} + // rejectMacOS refuses operations the experimental macOS guest does not implement. // Callers check it right after loading the record they already hold under the lock. func (s *StoredMetadata) rejectMacOS(operation string) error { diff --git a/lib/instances/macos_test.go b/lib/instances/macos_test.go index ba44e5b0b..5e020ebf3 100644 --- a/lib/instances/macos_test.go +++ b/lib/instances/macos_test.go @@ -35,6 +35,75 @@ func TestMacOSRequestDefaultsAndRejections(t *testing.T) { } require.ErrorIs(t, prepareMacOSCreate(&CreateInstanceRequest{}, testMacImage(), hypervisor.Capabilities{}), ErrInvalidRequest) } +func TestMacOSGuestAgentOptIn(t *testing.T) { + if runtime.GOOS != "darwin" || runtime.GOARCH != "arm64" { + t.Skip("Mac request acceptance requires Apple silicon") + } + image := testMacImage() + image.MacOS.GuestAgent = true + request := CreateInstanceRequest{} + require.NoError(t, prepareMacOSRequest(&request, image, hypervisor.TypeVZ)) + require.False(t, request.SkipGuestAgent) + request = CreateInstanceRequest{SkipGuestAgent: true} + require.NoError(t, prepareMacOSRequest(&request, image, hypervisor.TypeVZ)) + require.True(t, request.SkipGuestAgent) +} + +func TestMacOSAgentReadinessSeparateFromRunning(t *testing.T) { + for _, tc := range []struct { + name string + declared, skipped, ready bool + }{ + {"legacy image", false, false, false}, + {"disabled", true, true, false}, + {"not ready", true, false, false}, + {"ready", true, false, true}, + } { + t.Run(tc.name, func(t *testing.T) { + image := testMacImage() + image.MacOS.GuestAgent = tc.declared + stored := StoredMetadata{Id: "mac", MacOS: image.MacOS, SkipGuestAgent: tc.skipped} + now := time.Now().UTC() + calls := 0 + m := &manager{now: func() time.Time { return now }, guestAgentReadyProbe: func(context.Context, *StoredMetadata) bool { + calls++ + return tc.ready + }} + require.Equal(t, StateRunning, deriveRunningState(&stored)) + hydrated := m.hydrateBootMarkersFromLogs(context.Background(), &stored) + require.Equal(t, tc.ready, hydrated) + require.Equal(t, StateRunning, deriveRunningState(&stored)) + require.Nil(t, stored.ProgramStartedAt, "macOS does not invent a Linux workload marker") + if tc.ready { + require.Equal(t, &now, stored.GuestAgentReadyAt) + } else { + require.Nil(t, stored.GuestAgentReadyAt) + } + if tc.declared && !tc.skipped { + require.Equal(t, 1, calls) + } else { + require.Zero(t, calls) + } + }) + } +} + +func TestMacOSAgentReadyMarkerPersisted(t *testing.T) { + p := paths.New(t.TempDir()) + now := time.Now().UTC() + m := &manager{paths: p, now: func() time.Time { return now }, guestAgentReadyProbe: func(context.Context, *StoredMetadata) bool { return true }} + require.NoError(t, m.ensureDirectories("mac")) + image := testMacImage() + image.MacOS.GuestAgent = true + require.NoError(t, m.saveMetadata(&metadata{StoredMetadata: StoredMetadata{Id: "mac", DataDir: p.InstanceDir("mac"), MacOS: image.MacOS}})) + m.persistBootMarkers(context.Background(), "mac") + meta, err := m.loadMetadata("mac") + require.NoError(t, err) + require.Equal(t, &now, meta.GuestAgentReadyAt) + require.Nil(t, meta.ProgramStartedAt) + require.True(t, meta.MacOS.GuestAgent) +} + func TestMacOSBootConfigNoLinuxDevices(t *testing.T) { p := paths.New(t.TempDir()) m := &manager{paths: p} diff --git a/lib/instances/query.go b/lib/instances/query.go index 02ec7ed49..fbf739ccd 100644 --- a/lib/instances/query.go +++ b/lib/instances/query.go @@ -219,7 +219,7 @@ func (m *manager) updateCachedHypervisorStateFromInstance(inst *Instance) { } func deriveRunningState(stored *StoredMetadata) State { - // For unmanaged macOS desktops, Running means VMM running, not app/agent ready. + // For macOS, Running means VMM running. Agent readiness is a separate marker. if stored.MacOS != nil { return StateRunning } @@ -293,8 +293,8 @@ func advancePhaseIfRunning(stored *StoredMetadata) { // services do not need to forward stdout/stderr to the serial console. // Returns true when at least one missing marker was found and populated. func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *StoredMetadata) bool { - needProgram := stored.ProgramStartedAt == nil - needAgent := !stored.SkipGuestAgent && stored.GuestAgentReadyAt == nil + needProgram := stored.MacOS == nil && stored.ProgramStartedAt == nil + needAgent := guestAgentEnabled(stored) && stored.GuestAgentReadyAt == nil if !needProgram && !needAgent { m.clearBootMarkerRescan(stored.Id) return false @@ -308,7 +308,10 @@ func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *Stored ) defer span.End() - programStartedAt, guestAgentReadyAt := m.parseBootMarkers(ctx, stored.Id, needProgram, needAgent, stored.StartedAt) + var programStartedAt, guestAgentReadyAt *time.Time + if stored.MacOS == nil { + programStartedAt, guestAgentReadyAt = m.parseBootMarkers(ctx, stored.Id, needProgram, needAgent, stored.StartedAt) + } hydrated := false if needProgram && programStartedAt != nil { stored.ProgramStartedAt = programStartedAt @@ -318,7 +321,7 @@ func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *Stored stored.GuestAgentReadyAt = guestAgentReadyAt hydrated = true } - if needAgent && stored.GuestAgentReadyAt == nil && stored.ProgramStartedAt != nil && m.hydrateGuestAgentReadyFromProbe(ctx, stored) { + if needAgent && stored.GuestAgentReadyAt == nil && (stored.MacOS != nil || stored.ProgramStartedAt != nil) && m.hydrateGuestAgentReadyFromProbe(ctx, stored) { hydrated = true } if hydrated { @@ -421,7 +424,7 @@ func (m *manager) nowUTC() time.Time { } func (m *manager) hydrateGuestAgentReadyFromProbe(ctx context.Context, stored *StoredMetadata) bool { - if stored == nil || stored.SkipGuestAgent || stored.GuestAgentReadyAt != nil { + if stored == nil || !guestAgentEnabled(stored) || stored.GuestAgentReadyAt != nil { return false } probe := m.guestAgentReadyProbe @@ -437,7 +440,7 @@ func (m *manager) hydrateGuestAgentReadyFromProbe(ctx context.Context, stored *S } func probeGuestAgentReady(ctx context.Context, stored *StoredMetadata) bool { - if stored == nil || stored.SkipGuestAgent { + if stored == nil || !guestAgentEnabled(stored) { return false } dialer, err := hypervisor.NewVsockDialer(stored.HypervisorType, stored.VsockSocket, stored.VsockCID) @@ -643,13 +646,16 @@ func (m *manager) persistBootMarkers(ctx context.Context, id string) { return } - needProgram := meta.ProgramStartedAt == nil - needAgent := !meta.SkipGuestAgent && meta.GuestAgentReadyAt == nil + needProgram := meta.MacOS == nil && meta.ProgramStartedAt == nil + needAgent := guestAgentEnabled(&meta.StoredMetadata) && meta.GuestAgentReadyAt == nil if !needProgram && !needAgent { return } - programStartedAt, guestAgentReadyAt := m.parseBootMarkers(ctx, id, needProgram, needAgent, meta.StartedAt) + var programStartedAt, guestAgentReadyAt *time.Time + if meta.MacOS == nil { + programStartedAt, guestAgentReadyAt = m.parseBootMarkers(ctx, id, needProgram, needAgent, meta.StartedAt) + } updated := false if needProgram && programStartedAt != nil { meta.ProgramStartedAt = programStartedAt @@ -659,7 +665,7 @@ func (m *manager) persistBootMarkers(ctx context.Context, id string) { meta.GuestAgentReadyAt = guestAgentReadyAt updated = true } - if needAgent && meta.GuestAgentReadyAt == nil && meta.ProgramStartedAt != nil && m.hydrateGuestAgentReadyFromProbe(ctx, &meta.StoredMetadata) { + if needAgent && meta.GuestAgentReadyAt == nil && (meta.MacOS != nil || meta.ProgramStartedAt != nil) && m.hydrateGuestAgentReadyFromProbe(ctx, &meta.StoredMetadata) { updated = true } if !updated { diff --git a/lib/instances/stop.go b/lib/instances/stop.go index 603b9a75e..88b01d2bb 100644 --- a/lib/instances/stop.go +++ b/lib/instances/stop.go @@ -39,12 +39,12 @@ func hasVFIODevices(stored *StoredMetadata) bool { return storedVGPUDevicePath(stored) != "" || len(stored.Devices) > 0 } -// tryGracefulGuestShutdown asks guest init to shut down and waits for the +// tryGracefulGuestShutdown asks GuestService to shut down and waits for the // hypervisor process to exit. Returns true if the process exited in time. func (m *manager) tryGracefulGuestShutdown(ctx context.Context, inst *Instance, stopTimeout int) bool { log := logger.FromContext(ctx) - if inst.SkipGuestAgent { + if !guestAgentEnabled(&inst.StoredMetadata) { log.DebugContext(ctx, "guest-agent disabled, skipping graceful guest shutdown", "instance_id", inst.Id) return false } diff --git a/lib/instances/vsock.go b/lib/instances/vsock.go index f4060950e..352f0ea7d 100644 --- a/lib/instances/vsock.go +++ b/lib/instances/vsock.go @@ -14,8 +14,8 @@ func (m *manager) GetVsockDialer(ctx context.Context, instanceID string) (hyperv return nil, err } - if inst.MacOS != nil { - return nil, fmt.Errorf("%w: macOS guest-agent exec/stat is not implemented", ErrInvalidRequest) + if inst.MacOS != nil && !guestAgentEnabled(&inst.StoredMetadata) { + return nil, fmt.Errorf("%w: macOS image does not enable the shared guest agent", ErrInvalidRequest) } return hypervisor.NewVsockDialer(hypervisor.Type(inst.HypervisorType), inst.VsockSocket, inst.VsockCID) } From bc563782958aed70d93321d1f2d9488adaba347b Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 14:43:20 +0000 Subject: [PATCH 03/15] Bound exec drain after cancellation and kill the process group - Persist macOS guest-agent readiness on the public read path by hydrating boot markers for macOS instances too. - Kill the command's process group on cancellation or timeout so shell descendants do not outlive the command. - Bound how long a cancelled exec waits for output to drain, so a client that stopped reading cannot hold the handler open. - Hold the fork lock across Darwin vsock accept so an accepted descriptor is not inherited by a concurrent fork before it is marked close-on-exec. --- lib/instances/macos_test.go | 21 ++++++ lib/instances/query.go | 2 +- lib/system/guest_agent/exec.go | 68 ++++++++++++++++++-- lib/system/guest_agent/service_test.go | 89 ++++++++++++++++++++++++++ lib/system/guest_agent/vsock_darwin.go | 4 ++ 5 files changed, 178 insertions(+), 6 deletions(-) diff --git a/lib/instances/macos_test.go b/lib/instances/macos_test.go index 5e020ebf3..16baa808e 100644 --- a/lib/instances/macos_test.go +++ b/lib/instances/macos_test.go @@ -104,6 +104,27 @@ func TestMacOSAgentReadyMarkerPersisted(t *testing.T) { require.True(t, meta.MacOS.GuestAgent) } +func TestMacOSAgentReadyMarkerPersistedByPublicRead(t *testing.T) { + p := paths.New(t.TempDir()) + now := time.Now().UTC() + m := &manager{paths: p, now: func() time.Time { return now }, guestAgentReadyProbe: func(context.Context, *StoredMetadata) bool { return true }} + require.NoError(t, m.ensureDirectories("mac")) + socket := filepath.Join(p.InstanceDir("mac"), "vz.sock") + require.NoError(t, os.WriteFile(socket, nil, 0600)) + image := testMacImage() + image.MacOS.GuestAgent = true + stored := StoredMetadata{Id: "mac", DataDir: p.InstanceDir("mac"), SocketPath: socket, MacOS: image.MacOS, HypervisorType: hypervisor.TypeVZ, CreatedAt: now} + require.NoError(t, m.saveMetadata(&metadata{StoredMetadata: stored})) + m.storeCachedHypervisorState("mac", hypervisor.StateRunning) + + inst, err := m.GetInstance(context.Background(), "mac") + require.NoError(t, err) + require.NotNil(t, inst.GuestAgentReadyAt) + meta, err := m.loadMetadata("mac") + require.NoError(t, err) + require.NotNil(t, meta.GuestAgentReadyAt, "readiness must reach metadata through the normal read path") +} + func TestMacOSBootConfigNoLinuxDevices(t *testing.T) { p := paths.New(t.TempDir()) m := &manager{paths: p} diff --git a/lib/instances/query.go b/lib/instances/query.go index fbf739ccd..812d52067 100644 --- a/lib/instances/query.go +++ b/lib/instances/query.go @@ -137,7 +137,7 @@ func (m *manager) deriveStateWithOptions(ctx context.Context, stored *StoredMeta return stateResult{State: StateCreated} case hypervisor.StateRunning: hydrated := false - if hydrateBootMarkers && stored.MacOS == nil { + if hydrateBootMarkers { hydrated = m.hydrateBootMarkersFromLogs(ctx, stored) } return stateResult{ diff --git a/lib/system/guest_agent/exec.go b/lib/system/guest_agent/exec.go index 98f9175e1..07ebc5f30 100644 --- a/lib/system/guest_agent/exec.go +++ b/lib/system/guest_agent/exec.go @@ -2,12 +2,14 @@ package main import ( "context" + "errors" "fmt" "log" "os" "os/exec" "strings" "sync" + "syscall" "time" "github.com/creack/pty" @@ -65,6 +67,8 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E cmd.Env = s.buildEnv(start.Env, false) cmd.Dir = start.Cwd cmd.WaitDelay = 2 * time.Second + cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} + cmd.Cancel = func() error { return killProcessGroup(cmd) } var sendMu sync.Mutex cmd.Stdout = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel} cmd.Stderr = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel, stderr: true} @@ -90,8 +94,16 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E } }() - // os/exec drains stdout/stderr through bounded writers before Wait returns. - waitErr := cmd.Wait() + // cmd.Wait also waits for the stdout/stderr copies, which can block in stream.Send. + waitDone := make(chan struct{}) + var waitErr error + go func() { + waitErr = cmd.Wait() + close(waitDone) + }() + if err := awaitDrain(ctx, waitDone); err != nil { + return err + } if waitErr != nil { if _, exited := waitErr.(*exec.ExitError); !exited { return fmt.Errorf("stream command output: %w", waitErr) @@ -114,6 +126,37 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E }) } +// execDrainGrace bounds how long a cancelled command may take to finish. A descendant that +// holds the terminal or output, or a client that stopped reading, must not block the handler. +var execDrainGrace = 5 * time.Second + +// awaitDrain waits for done. After ctx ends it allows execDrainGrace more, then returns an +// error so the handler returns, ending the RPC and releasing any output send still blocked on it. +func awaitDrain(ctx context.Context, done <-chan struct{}) error { + select { + case <-done: + return nil + case <-ctx.Done(): + } + timer := time.NewTimer(execDrainGrace) + defer timer.Stop() + select { + case <-done: + return nil + case <-timer.C: + return fmt.Errorf("command output did not drain after cancellation: %w", ctx.Err()) + } +} + +// killProcessGroup kills the command and its descendants. Both exec paths start the command +// as a group leader (Setpgid for no-TTY, Setsid for TTY), so the group id is its pid. +func killProcessGroup(cmd *exec.Cmd) error { + if err := syscall.Kill(-cmd.Process.Pid, syscall.SIGKILL); err != nil && !errors.Is(err, syscall.ESRCH) { + return err + } + return nil +} + type execStreamWriter struct { stream pb.GuestService_ExecServer mu *sync.Mutex @@ -150,6 +193,7 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe } cmd := exec.CommandContext(ctx, start.Command[0], start.Command[1:]...) + cmd.Cancel = func() error { return killProcessGroup(cmd) } // Set up environment (TTY mode adds TERM default) cmd.Env = s.buildEnv(start.Env, true) @@ -226,11 +270,25 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe } }() - // Wait for command or context cancellation - waitErr := cmd.Wait() + waitDone := make(chan struct{}) + var waitErr error + go func() { + waitErr = cmd.Wait() + close(waitDone) + }() + if err := awaitDrain(ctx, waitDone); err != nil { + return err + } // Wait for all output to be sent - wg.Wait() + outputDone := make(chan struct{}) + go func() { + wg.Wait() + close(outputDone) + }() + if err := awaitDrain(ctx, outputDone); err != nil { + return err + } exitCode := int32(0) if cmd.ProcessState != nil { diff --git a/lib/system/guest_agent/service_test.go b/lib/system/guest_agent/service_test.go index 9171ab672..1f9c6256a 100644 --- a/lib/system/guest_agent/service_test.go +++ b/lib/system/guest_agent/service_test.go @@ -2,9 +2,11 @@ package main import ( "context" + "errors" "io" "net" "os" + "os/exec" "path/filepath" "strconv" "strings" @@ -131,6 +133,93 @@ func TestGuestServiceExecCancellation(t *testing.T) { }, 5*time.Second, 10*time.Millisecond, "disconnect must cancel the command, not leave it running") } +func TestGuestServiceExecCancellationKillsDescendants(t *testing.T) { + for _, tty := range []bool{false, true} { + t.Run(map[bool]string{false: "pipes", true: "tty"}[tty], func(t *testing.T) { + client := testGuestClient(t) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + pidPath := filepath.Join(t.TempDir(), "pid") + stream, err := client.Exec(ctx) + require.NoError(t, err) + require.NoError(t, stream.Send(&pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", "sleep 30 & printf '%s' \"$!\" > \"$HYPEMAN_PID_FILE\"; wait"}, + Env: map[string]string{"HYPEMAN_PID_FILE": pidPath}, + Tty: tty, + }}})) + require.NoError(t, stream.CloseSend()) + var pid int + require.Eventually(t, func() bool { + data, err := os.ReadFile(pidPath) + if err != nil { + return false + } + pid, err = strconv.Atoi(string(data)) + return err == nil && pid > 0 + }, 5*time.Second, 10*time.Millisecond) + defer syscall.Kill(pid, syscall.SIGKILL) + cancel() + require.Eventually(t, func() bool { + return processGone(pid) + }, 5*time.Second, 10*time.Millisecond, "cancellation must reach descendants, not only the shell") + }) + } +} + +// processGone reports whether pid no longer runs. A zombie counts as gone. +func processGone(pid int) bool { + if errors.Is(syscall.Kill(pid, 0), syscall.ESRCH) { + return true + } + out, err := exec.Command("ps", "-o", "stat=", "-p", strconv.Itoa(pid)).Output() + return err != nil || strings.HasPrefix(strings.TrimSpace(string(out)), "Z") +} + +// stalledExecStream models a client that stopped reading: Send blocks until release closes. +type stalledExecStream struct { + grpc.ServerStream + ctx context.Context + start *pb.ExecRequest + release chan struct{} +} + +func (s *stalledExecStream) Context() context.Context { return s.ctx } +func (s *stalledExecStream) Recv() (*pb.ExecRequest, error) { + if s.start != nil { + req := s.start + s.start = nil + return req, nil + } + <-s.ctx.Done() + return nil, s.ctx.Err() +} +func (s *stalledExecStream) Send(*pb.ExecResponse) error { + <-s.release + return nil +} + +func TestGuestServiceExecTimeoutFinishesWhenClientStopsReading(t *testing.T) { + restore := execDrainGrace + execDrainGrace = 100 * time.Millisecond + t.Cleanup(func() { execDrainGrace = restore }) + ctx, cancel := context.WithCancel(context.Background()) + t.Cleanup(cancel) + release := make(chan struct{}) + t.Cleanup(func() { close(release) }) + stream := &stalledExecStream{ctx: ctx, release: release, start: &pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", "yes"}, + TimeoutSeconds: 1, + }}}} + done := make(chan error, 1) + go func() { done <- (&guestServer{}).Exec(stream) }() + select { + case err := <-done: + require.ErrorContains(t, err, "did not drain") + case <-time.After(10 * time.Second): + t.Fatal("exec stayed blocked behind a stalled client after its timeout") + } +} + func TestGuestServiceFileRoundTrip(t *testing.T) { client := testGuestClient(t) ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) diff --git a/lib/system/guest_agent/vsock_darwin.go b/lib/system/guest_agent/vsock_darwin.go index b3c47ce4a..419de01b4 100644 --- a/lib/system/guest_agent/vsock_darwin.go +++ b/lib/system/guest_agent/vsock_darwin.go @@ -100,7 +100,11 @@ func (l *vmListener) Accept() (net.Conn, error) { l.mu.Unlock() return nil, net.ErrClosed } + // Darwin has no accept4: the accepted fd is inheritable until accept_vm sets FD_CLOEXEC. + // Holding the fork lock across that window keeps a concurrent fork from copying it. + syscall.ForkLock.RLock() fd, err := C.accept_vm(l.fd) + syscall.ForkLock.RUnlock() l.mu.Unlock() if fd >= 0 { return &vmConn{file: os.NewFile(uintptr(fd), "guest-vsock"), local: l.Addr()}, nil From 1532a55fe18b6d3aa8602bbd1d498a5ff14b9c49 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 14:44:26 +0000 Subject: [PATCH 04/15] Document root exec authority and exec timeout scope The guest agent runs commands as root, as the Linux agent does, and the host API is the only authorization boundary. Exec timeout ends the command; it does not bound a healthy output stream. --- docs/macos-guest-agent.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/docs/macos-guest-agent.md b/docs/macos-guest-agent.md index 51e5eef67..03f523c80 100644 --- a/docs/macos-guest-agent.md +++ b/docs/macos-guest-agent.md @@ -36,6 +36,9 @@ by this initial patch. The host API and hypervisor still enforce authorization. The vsock host-CID check is transport admission, not a replacement for instance authority checks. +As with the Linux guest agent, commands and file operations run as root with no +per-caller credential switch; the host API alone decides who may reach them. +An exec `timeout_seconds` ends the command; it does not bound a healthy stream. Do not expose this privileged service through unauthenticated host forwarding. ## Image declaration and host integration From 03d413cee491a70a12b8e1576364fbb5d33214e2 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 14:58:04 +0000 Subject: [PATCH 05/15] Bound the exit-code send and kill exec descendants after the shell exits - Send the final exit code under the same bound as command output, so a client that stopped reading cannot hold the handler past a timeout. - Kill the command's process group from the context rather than only through exec.Cmd.Cancel, which is not called once the direct child has exited. - Make the drain bound a server field instead of a package variable, so tests do not mutate shared state read by handlers that outlive them. --- lib/system/guest_agent/exec.go | 71 ++++++++++++++++------ lib/system/guest_agent/main.go | 2 + lib/system/guest_agent/service_test.go | 84 +++++++++++++++++++++++--- 3 files changed, 130 insertions(+), 27 deletions(-) diff --git a/lib/system/guest_agent/exec.go b/lib/system/guest_agent/exec.go index 07ebc5f30..ffa0581a0 100644 --- a/lib/system/guest_agent/exec.go +++ b/lib/system/guest_agent/exec.go @@ -79,6 +79,11 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E if err := cmd.Start(); err != nil { return fmt.Errorf("start command: %w", err) } + // Cancel is not called once the direct child has exited, so a descendant that + // still holds the pipes would outlive the timeout. Kill the group from the context. + finished := make(chan struct{}) + defer close(finished) + go killGroupOnDone(ctx, cmd, finished) // Handle stdin in background go func() { @@ -101,7 +106,7 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E waitErr = cmd.Wait() close(waitDone) }() - if err := awaitDrain(ctx, waitDone); err != nil { + if err := awaitBounded(ctx, waitDone, s.drainBound()); err != nil { return err } if waitErr != nil { @@ -120,32 +125,62 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E log.Printf("[guest-agent] command finished with exit code: %d", exitCode) - // Send exit code - return stream.Send(&pb.ExecResponse{ - Response: &pb.ExecResponse_ExitCode{ExitCode: exitCode}, - }) + return s.sendExitCode(ctx, stream, exitCode) } -// execDrainGrace bounds how long a cancelled command may take to finish. A descendant that +// defaultDrainGrace bounds how long a cancelled command may take to finish. A descendant that // holds the terminal or output, or a client that stopped reading, must not block the handler. -var execDrainGrace = 5 * time.Second +const defaultDrainGrace = 5 * time.Second -// awaitDrain waits for done. After ctx ends it allows execDrainGrace more, then returns an -// error so the handler returns, ending the RPC and releasing any output send still blocked on it. -func awaitDrain(ctx context.Context, done <-chan struct{}) error { +func (s *guestServer) drainBound() time.Duration { + if s.drainGrace > 0 { + return s.drainGrace + } + return defaultDrainGrace +} + +// awaitBounded waits for done. After ctx ends it allows grace more, then returns an +// error so the handler returns, ending the RPC and releasing any send still blocked on it. +func awaitBounded(ctx context.Context, done <-chan struct{}, grace time.Duration) error { select { case <-done: return nil case <-ctx.Done(): } - timer := time.NewTimer(execDrainGrace) + timer := time.NewTimer(grace) defer timer.Stop() select { case <-done: return nil case <-timer.C: - return fmt.Errorf("command output did not drain after cancellation: %w", ctx.Err()) + return fmt.Errorf("command did not finish after cancellation: %w", ctx.Err()) + } +} + +// killGroupOnDone kills cmd's process group once ctx ends, unless finished closes first. +func killGroupOnDone(ctx context.Context, cmd *exec.Cmd, finished <-chan struct{}) { + select { + case <-ctx.Done(): + _ = killProcessGroup(cmd) + case <-finished: + } +} + +// sendExitCode delivers the final status under the same bound as output, so a client +// that stopped reading cannot hold the handler past cancellation. +func (s *guestServer) sendExitCode(ctx context.Context, stream pb.GuestService_ExecServer, exitCode int32) error { + sent := make(chan struct{}) + var err error + go func() { + err = stream.Send(&pb.ExecResponse{ + Response: &pb.ExecResponse_ExitCode{ExitCode: exitCode}, + }) + close(sent) + }() + if boundErr := awaitBounded(ctx, sent, s.drainBound()); boundErr != nil { + return boundErr } + return err } // killProcessGroup kills the command and its descendants. Both exec paths start the command @@ -221,6 +256,9 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe return fmt.Errorf("start pty: %w", err) } defer ptmx.Close() + finished := make(chan struct{}) + defer close(finished) + go killGroupOnDone(ctx, cmd, finished) // Mutex to protect concurrent stream.Send calls (gRPC streams are not thread-safe) var sendMu sync.Mutex @@ -276,7 +314,7 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe waitErr = cmd.Wait() close(waitDone) }() - if err := awaitDrain(ctx, waitDone); err != nil { + if err := awaitBounded(ctx, waitDone, s.drainBound()); err != nil { return err } @@ -286,7 +324,7 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe wg.Wait() close(outputDone) }() - if err := awaitDrain(ctx, outputDone); err != nil { + if err := awaitBounded(ctx, outputDone, s.drainBound()); err != nil { return err } @@ -300,10 +338,7 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe log.Printf("[guest-agent] TTY command finished with exit code: %d", exitCode) - // Send exit code - return stream.Send(&pb.ExecResponse{ - Response: &pb.ExecResponse_ExitCode{ExitCode: exitCode}, - }) + return s.sendExitCode(ctx, stream, exitCode) } // buildEnv constructs environment variables by merging provided env with defaults. diff --git a/lib/system/guest_agent/main.go b/lib/system/guest_agent/main.go index 33db3e889..5f8159fc6 100644 --- a/lib/system/guest_agent/main.go +++ b/lib/system/guest_agent/main.go @@ -22,6 +22,8 @@ const ( type guestServer struct { pb.UnimplementedGuestServiceServer gpuReporter *gpuInitReporter + // drainGrace bounds how long a cancelled exec may take to finish; zero means defaultDrainGrace. + drainGrace time.Duration } func main() { diff --git a/lib/system/guest_agent/service_test.go b/lib/system/guest_agent/service_test.go index 1f9c6256a..d33506099 100644 --- a/lib/system/guest_agent/service_test.go +++ b/lib/system/guest_agent/service_test.go @@ -176,11 +176,13 @@ func processGone(pid int) bool { } // stalledExecStream models a client that stopped reading: Send blocks until release closes. +// With stallExitOnly, only the final exit-code send blocks. type stalledExecStream struct { grpc.ServerStream - ctx context.Context - start *pb.ExecRequest - release chan struct{} + ctx context.Context + start *pb.ExecRequest + release chan struct{} + stallExitOnly bool } func (s *stalledExecStream) Context() context.Context { return s.ctx } @@ -193,15 +195,15 @@ func (s *stalledExecStream) Recv() (*pb.ExecRequest, error) { <-s.ctx.Done() return nil, s.ctx.Err() } -func (s *stalledExecStream) Send(*pb.ExecResponse) error { +func (s *stalledExecStream) Send(resp *pb.ExecResponse) error { + if _, exit := resp.Response.(*pb.ExecResponse_ExitCode); s.stallExitOnly && !exit { + return nil + } <-s.release return nil } func TestGuestServiceExecTimeoutFinishesWhenClientStopsReading(t *testing.T) { - restore := execDrainGrace - execDrainGrace = 100 * time.Millisecond - t.Cleanup(func() { execDrainGrace = restore }) ctx, cancel := context.WithCancel(context.Background()) t.Cleanup(cancel) release := make(chan struct{}) @@ -211,15 +213,79 @@ func TestGuestServiceExecTimeoutFinishesWhenClientStopsReading(t *testing.T) { TimeoutSeconds: 1, }}}} done := make(chan error, 1) - go func() { done <- (&guestServer{}).Exec(stream) }() + go func() { done <- (&guestServer{drainGrace: 100 * time.Millisecond}).Exec(stream) }() select { case err := <-done: - require.ErrorContains(t, err, "did not drain") + require.ErrorContains(t, err, "did not finish") case <-time.After(10 * time.Second): t.Fatal("exec stayed blocked behind a stalled client after its timeout") } } +func TestGuestServiceExecTimeoutBoundsExitCodeSend(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + t.Cleanup(cancel) + release := make(chan struct{}) + t.Cleanup(func() { close(release) }) + stream := &stalledExecStream{ctx: ctx, release: release, stallExitOnly: true, start: &pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", "echo started; sleep 30"}, + TimeoutSeconds: 1, + }}}} + done := make(chan error, 1) + go func() { done <- (&guestServer{drainGrace: 100 * time.Millisecond}).Exec(stream) }() + select { + case err := <-done: + require.ErrorContains(t, err, "did not finish") + case <-time.After(10 * time.Second): + t.Fatal("exit-code send stayed blocked after the command timed out") + } +} + +func TestGuestServiceExecTimeoutKillsDescendantAfterShellExits(t *testing.T) { + for _, tty := range []bool{false, true} { + t.Run(map[bool]string{false: "pipes", true: "tty"}[tty], func(t *testing.T) { + client := testGuestClient(t) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + pidPath := filepath.Join(t.TempDir(), "pid") + stream, err := client.Exec(ctx) + require.NoError(t, err) + require.NoError(t, stream.Send(&pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", "sleep 30 & printf '%s' \"$!\" > \"$HYPEMAN_PID_FILE\"; exit 0"}, + Env: map[string]string{"HYPEMAN_PID_FILE": pidPath}, + Tty: tty, + TimeoutSeconds: 1, + }}})) + require.NoError(t, stream.CloseSend()) + var pid int + require.Eventually(t, func() bool { + data, err := os.ReadFile(pidPath) + if err != nil { + return false + } + pid, err = strconv.Atoi(string(data)) + return err == nil && pid > 0 + }, 5*time.Second, 10*time.Millisecond) + defer syscall.Kill(pid, syscall.SIGKILL) + finished := make(chan error, 1) + go func() { + for { + if _, err := stream.Recv(); err != nil { + finished <- err + return + } + } + }() + select { + case <-finished: + case <-time.After(10 * time.Second): + t.Fatal("RPC stayed open while a descendant held the command's output") + } + require.True(t, processGone(pid), "timeout must kill the descendant even after the shell exited") + }) + } +} + func TestGuestServiceFileRoundTrip(t *testing.T) { client := testGuestClient(t) ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) From 1830db3c8ec10a48e213e45f9bcecc353c7e227f Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 15:07:32 +0000 Subject: [PATCH 06/15] Share boot marker discovery between hydration and persistence Hydration and persistence duplicated the missing-marker checks, the serial log parse, marker assignment and the readiness probe. Extract them into one applyBootMarkers routine. Hydration keeps its scan throttling and rescan bookkeeping; persistence keeps its save and metrics. --- lib/instances/query.go | 60 ++++++++++++++++++++---------------------- 1 file changed, 28 insertions(+), 32 deletions(-) diff --git a/lib/instances/query.go b/lib/instances/query.go index 812d52067..038deaa4b 100644 --- a/lib/instances/query.go +++ b/lib/instances/query.go @@ -293,8 +293,7 @@ func advancePhaseIfRunning(stored *StoredMetadata) { // services do not need to forward stdout/stderr to the serial console. // Returns true when at least one missing marker was found and populated. func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *StoredMetadata) bool { - needProgram := stored.MacOS == nil && stored.ProgramStartedAt == nil - needAgent := guestAgentEnabled(stored) && stored.GuestAgentReadyAt == nil + needProgram, needAgent := bootMarkersMissing(stored) if !needProgram && !needAgent { m.clearBootMarkerRescan(stored.Id) return false @@ -308,6 +307,32 @@ func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *Stored ) defer span.End() + hydrated := m.applyBootMarkers(ctx, stored) + if hydrated { + advancePhaseIfRunning(stored) + m.clearBootMarkerRescan(stored.Id) + } else { + m.deferBootMarkerRescan(stored.Id) + } + return hydrated +} + +// bootMarkersMissing reports which boot markers stored still lacks. +func bootMarkersMissing(stored *StoredMetadata) (needProgram, needAgent bool) { + needProgram = stored.MacOS == nil && stored.ProgramStartedAt == nil + needAgent = guestAgentEnabled(stored) && stored.GuestAgentReadyAt == nil + return needProgram, needAgent +} + +// applyBootMarkers fills missing boot markers on stored from serial logs, then from +// the guest-agent probe, and reports whether it filled any. Callers own locking, +// phase advancement and persistence. +func (m *manager) applyBootMarkers(ctx context.Context, stored *StoredMetadata) bool { + needProgram, needAgent := bootMarkersMissing(stored) + if !needProgram && !needAgent { + return false + } + var programStartedAt, guestAgentReadyAt *time.Time if stored.MacOS == nil { programStartedAt, guestAgentReadyAt = m.parseBootMarkers(ctx, stored.Id, needProgram, needAgent, stored.StartedAt) @@ -324,12 +349,6 @@ func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *Stored if needAgent && stored.GuestAgentReadyAt == nil && (stored.MacOS != nil || stored.ProgramStartedAt != nil) && m.hydrateGuestAgentReadyFromProbe(ctx, stored) { hydrated = true } - if hydrated { - advancePhaseIfRunning(stored) - m.clearBootMarkerRescan(stored.Id) - } else { - m.deferBootMarkerRescan(stored.Id) - } return hydrated } @@ -645,30 +664,7 @@ func (m *manager) persistBootMarkers(ctx context.Context, id string) { if err != nil { return } - - needProgram := meta.MacOS == nil && meta.ProgramStartedAt == nil - needAgent := guestAgentEnabled(&meta.StoredMetadata) && meta.GuestAgentReadyAt == nil - if !needProgram && !needAgent { - return - } - - var programStartedAt, guestAgentReadyAt *time.Time - if meta.MacOS == nil { - programStartedAt, guestAgentReadyAt = m.parseBootMarkers(ctx, id, needProgram, needAgent, meta.StartedAt) - } - updated := false - if needProgram && programStartedAt != nil { - meta.ProgramStartedAt = programStartedAt - updated = true - } - if needAgent && guestAgentReadyAt != nil { - meta.GuestAgentReadyAt = guestAgentReadyAt - updated = true - } - if needAgent && meta.GuestAgentReadyAt == nil && (meta.MacOS != nil || meta.ProgramStartedAt != nil) && m.hydrateGuestAgentReadyFromProbe(ctx, &meta.StoredMetadata) { - updated = true - } - if !updated { + if !m.applyBootMarkers(ctx, &meta.StoredMetadata) { return } From 74389b61f6844f45aa4858d2d7f7874033331f90 Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 19:10:25 +0000 Subject: [PATCH 07/15] Report exec exit status correctly and share one guest-agent predicate - A command that exits 0 while a background child still holds its output returns exec.ErrWaitDelay after the wait delay. That is a normal exit, so report its status instead of failing the RPC. - A command that hits its deadline reports 124 (GNU timeout convention). The group kill leaves a process state behind, so the old check never fired. - Share exit-code selection between the TTY and non-TTY paths. - Define StoredMetadata.GuestAgentEnabled once and use it in exec, cp, stop, vsock and boot-marker code instead of five copies of the expression. - Make the websocket exec cancel function required and drop the nil checks. - Remove the unreachable empty-command checks and the redundant Cwd guard. --- cmd/api/api/cp.go | 2 +- cmd/api/api/exec.go | 12 +++---- lib/instances/macos.go | 9 ++--- lib/instances/query.go | 6 ++-- lib/instances/stop.go | 2 +- lib/instances/vsock.go | 2 +- lib/system/guest_agent/exec.go | 49 ++++++++++++-------------- lib/system/guest_agent/service_test.go | 37 +++++++++++++++++++ 8 files changed, 74 insertions(+), 45 deletions(-) diff --git a/cmd/api/api/cp.go b/cmd/api/api/cp.go index ac4812b27..67d054efc 100644 --- a/cmd/api/api/cp.go +++ b/cmd/api/api/cp.go @@ -104,7 +104,7 @@ func (s *ApiService) CpHandler(w http.ResponseWriter, r *http.Request) { return } - if inst.MacOS != nil && (!inst.MacOS.GuestAgent || inst.SkipGuestAgent) { + if inst.MacOS != nil && !inst.GuestAgentEnabled() { http.Error(w, `{"code":"unsupported","message":"file copy requires the shared macOS guest agent to be enabled"}`, http.StatusNotImplemented) return } diff --git a/cmd/api/api/exec.go b/cmd/api/api/exec.go index f1d7ed2fe..60a3827fb 100644 --- a/cmd/api/api/exec.go +++ b/cmd/api/api/exec.go @@ -71,7 +71,7 @@ func (s *ApiService) ExecHandler(w http.ResponseWriter, r *http.Request) { return } - if inst.MacOS != nil && (!inst.MacOS.GuestAgent || inst.SkipGuestAgent) { + if inst.MacOS != nil && !inst.GuestAgentEnabled() { http.Error(w, `{"code":"unsupported","message":"exec is not implemented for macOS images without the shared guest agent enabled"}`, http.StatusNotImplemented) return } @@ -222,7 +222,7 @@ type wsReadWriter struct { reader io.Reader mu sync.Mutex resizeChan chan<- *guest.WindowSize // Channel to send resize events (nil if not TTY) - cancel context.CancelFunc // Exec session cancellation on disconnect (nil for other users). + cancel context.CancelFunc // ends the exec session when the websocket side fails } func (w *wsReadWriter) Read(p []byte) (n int, err error) { @@ -243,9 +243,7 @@ func (w *wsReadWriter) Read(p []byte) (n int, err error) { // Read next WebSocket message messageType, data, err := w.ws.ReadMessage() if err != nil { - if w.cancel != nil { - w.cancel() - } + w.cancel() return 0, err } @@ -278,9 +276,7 @@ func (w *wsReadWriter) Read(p []byte) (n int, err error) { func (w *wsReadWriter) Write(p []byte) (n int, err error) { if err := w.ws.WriteMessage(websocket.BinaryMessage, p); err != nil { - if w.cancel != nil { - w.cancel() - } + w.cancel() return 0, err } return len(p), nil diff --git a/lib/instances/macos.go b/lib/instances/macos.go index cc2be6c02..9e42e61d7 100644 --- a/lib/instances/macos.go +++ b/lib/instances/macos.go @@ -100,10 +100,11 @@ func (m *manager) checkMacOSIdentityAvailable(ctx context.Context, stored *Store return nil } -// guestAgentEnabled keeps old macOS templates unmanaged even if older or -// hand-written metadata did not set SkipGuestAgent explicitly. -func guestAgentEnabled(stored *StoredMetadata) bool { - return !stored.SkipGuestAgent && (stored.MacOS == nil || stored.MacOS.GuestAgent) +// GuestAgentEnabled reports whether exec, copy, readiness and shutdown may use the +// shared guest agent. A macOS image must declare it; hand-written metadata that +// leaves SkipGuestAgent unset does not enable it. +func (s *StoredMetadata) GuestAgentEnabled() bool { + return !s.SkipGuestAgent && (s.MacOS == nil || s.MacOS.GuestAgent) } // rejectMacOS refuses operations the experimental macOS guest does not implement. diff --git a/lib/instances/query.go b/lib/instances/query.go index 038deaa4b..1ad989ec3 100644 --- a/lib/instances/query.go +++ b/lib/instances/query.go @@ -320,7 +320,7 @@ func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *Stored // bootMarkersMissing reports which boot markers stored still lacks. func bootMarkersMissing(stored *StoredMetadata) (needProgram, needAgent bool) { needProgram = stored.MacOS == nil && stored.ProgramStartedAt == nil - needAgent = guestAgentEnabled(stored) && stored.GuestAgentReadyAt == nil + needAgent = stored.GuestAgentEnabled() && stored.GuestAgentReadyAt == nil return needProgram, needAgent } @@ -443,7 +443,7 @@ func (m *manager) nowUTC() time.Time { } func (m *manager) hydrateGuestAgentReadyFromProbe(ctx context.Context, stored *StoredMetadata) bool { - if stored == nil || !guestAgentEnabled(stored) || stored.GuestAgentReadyAt != nil { + if stored == nil || !stored.GuestAgentEnabled() || stored.GuestAgentReadyAt != nil { return false } probe := m.guestAgentReadyProbe @@ -459,7 +459,7 @@ func (m *manager) hydrateGuestAgentReadyFromProbe(ctx context.Context, stored *S } func probeGuestAgentReady(ctx context.Context, stored *StoredMetadata) bool { - if stored == nil || !guestAgentEnabled(stored) { + if stored == nil || !stored.GuestAgentEnabled() { return false } dialer, err := hypervisor.NewVsockDialer(stored.HypervisorType, stored.VsockSocket, stored.VsockCID) diff --git a/lib/instances/stop.go b/lib/instances/stop.go index 88b01d2bb..81103fdb8 100644 --- a/lib/instances/stop.go +++ b/lib/instances/stop.go @@ -44,7 +44,7 @@ func hasVFIODevices(stored *StoredMetadata) bool { func (m *manager) tryGracefulGuestShutdown(ctx context.Context, inst *Instance, stopTimeout int) bool { log := logger.FromContext(ctx) - if !guestAgentEnabled(&inst.StoredMetadata) { + if !inst.StoredMetadata.GuestAgentEnabled() { log.DebugContext(ctx, "guest-agent disabled, skipping graceful guest shutdown", "instance_id", inst.Id) return false } diff --git a/lib/instances/vsock.go b/lib/instances/vsock.go index 352f0ea7d..8cbefe4e0 100644 --- a/lib/instances/vsock.go +++ b/lib/instances/vsock.go @@ -14,7 +14,7 @@ func (m *manager) GetVsockDialer(ctx context.Context, instanceID string) (hyperv return nil, err } - if inst.MacOS != nil && !guestAgentEnabled(&inst.StoredMetadata) { + if inst.MacOS != nil && !inst.StoredMetadata.GuestAgentEnabled() { return nil, fmt.Errorf("%w: macOS image does not enable the shared guest agent", ErrInvalidRequest) } return hypervisor.NewVsockDialer(hypervisor.Type(inst.HypervisorType), inst.VsockSocket, inst.VsockCID) diff --git a/lib/system/guest_agent/exec.go b/lib/system/guest_agent/exec.go index ffa0581a0..ecb40d394 100644 --- a/lib/system/guest_agent/exec.go +++ b/lib/system/guest_agent/exec.go @@ -56,10 +56,6 @@ func (s *guestServer) Exec(stream pb.GuestService_ExecServer) error { // executeNoTTY executes command without TTY func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_ExecServer, start *pb.ExecStart) error { // Run command directly - guest-agent is already running in container namespace - if len(start.Command) == 0 { - return fmt.Errorf("empty command") - } - outputCtx, cancel := context.WithCancel(ctx) defer cancel() // CommandContext must observe both caller cancellation and stream-send errors. @@ -109,19 +105,16 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E if err := awaitBounded(ctx, waitDone, s.drainBound()); err != nil { return err } + // ErrWaitDelay means the command exited but a descendant still held its output. + // The status is already known, so it is reported as a normal exit. if waitErr != nil { - if _, exited := waitErr.(*exec.ExitError); !exited { + _, exited := waitErr.(*exec.ExitError) + if !exited && !errors.Is(waitErr, exec.ErrWaitDelay) { return fmt.Errorf("stream command output: %w", waitErr) } } - exitCode := int32(0) - if cmd.ProcessState != nil { - exitCode = int32(cmd.ProcessState.ExitCode()) - } else if waitErr != nil { - // If killed by timeout, exit with 124 (GNU timeout convention) - exitCode = 124 - } + exitCode := exitCodeOf(ctx, cmd, waitErr) log.Printf("[guest-agent] command finished with exit code: %d", exitCode) @@ -183,6 +176,21 @@ func (s *guestServer) sendExitCode(ctx context.Context, stream pb.GuestService_E return err } +// exitCodeOf reports how a command ended. A command that hit its deadline exits 124 +// (GNU timeout convention) even though the group kill leaves a process state behind. +func exitCodeOf(ctx context.Context, cmd *exec.Cmd, waitErr error) int32 { + if errors.Is(ctx.Err(), context.DeadlineExceeded) { + return 124 + } + if cmd.ProcessState != nil { + return int32(cmd.ProcessState.ExitCode()) + } + if waitErr != nil { + return 124 + } + return 0 +} + // killProcessGroup kills the command and its descendants. Both exec paths start the command // as a group leader (Setpgid for no-TTY, Setsid for TTY), so the group id is its pid. func killProcessGroup(cmd *exec.Cmd) error { @@ -223,20 +231,13 @@ func (w *execStreamWriter) Write(data []byte) (int, error) { func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_ExecServer, start *pb.ExecStart) error { // Run command directly with PTY - guest-agent is already running in container namespace // This ensures PTY and shell are in the same namespace, fixing Ctrl+C signal handling - if len(start.Command) == 0 { - return fmt.Errorf("empty command") - } - cmd := exec.CommandContext(ctx, start.Command[0], start.Command[1:]...) cmd.Cancel = func() error { return killProcessGroup(cmd) } // Set up environment (TTY mode adds TERM default) cmd.Env = s.buildEnv(start.Env, true) - // Set up working directory - if start.Cwd != "" { - cmd.Dir = start.Cwd - } + cmd.Dir = start.Cwd // Set up initial window size (use defaults if not specified) ws := &pty.Winsize{ @@ -328,13 +329,7 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe return err } - exitCode := int32(0) - if cmd.ProcessState != nil { - exitCode = int32(cmd.ProcessState.ExitCode()) - } else if waitErr != nil { - // If killed by timeout, exit with 124 (GNU timeout convention) - exitCode = 124 - } + exitCode := exitCodeOf(ctx, cmd, waitErr) log.Printf("[guest-agent] TTY command finished with exit code: %d", exitCode) diff --git a/lib/system/guest_agent/service_test.go b/lib/system/guest_agent/service_test.go index d33506099..dae9d4fcd 100644 --- a/lib/system/guest_agent/service_test.go +++ b/lib/system/guest_agent/service_test.go @@ -286,6 +286,43 @@ func TestGuestServiceExecTimeoutKillsDescendantAfterShellExits(t *testing.T) { } } +// runExec drives one exec to completion and returns its stdout and exit code. +func runExec(t *testing.T, start *pb.ExecStart) (string, int32) { + t.Helper() + client := testGuestClient(t) + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) + defer cancel() + stream, err := client.Exec(ctx) + require.NoError(t, err) + require.NoError(t, stream.Send(&pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: start}})) + require.NoError(t, stream.CloseSend()) + var stdout strings.Builder + exit := int32(-1) + for { + resp, err := stream.Recv() + if errors.Is(err, io.EOF) { + break + } + require.NoError(t, err) + stdout.Write(resp.GetStdout()) + if code, ok := resp.Response.(*pb.ExecResponse_ExitCode); ok { + exit = code.ExitCode + } + } + return stdout.String(), exit +} + +func TestGuestServiceExecExitZeroWhenDescendantHoldsOutput(t *testing.T) { + stdout, exit := runExec(t, &pb.ExecStart{Command: []string{"/bin/sh", "-c", "sleep 5 & echo hi"}}) + require.Equal(t, "hi\n", stdout) + require.Equal(t, int32(0), exit, "a command that exits 0 must report 0, not a stream error") +} + +func TestGuestServiceExecTimeoutReports124(t *testing.T) { + _, exit := runExec(t, &pb.ExecStart{Command: []string{"/bin/sleep", "30"}, TimeoutSeconds: 1}) + require.Equal(t, int32(124), exit) +} + func TestGuestServiceFileRoundTrip(t *testing.T) { client := testGuestClient(t) ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) From b99e4a23f814bf9ca5dd386cab362ca15df073dc Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 19:16:21 +0000 Subject: [PATCH 08/15] Use one cancellation context for exec and consolidate exec constants - Run exec under one cancellable context. The stream-send error path and caller disconnect cancel the same context, and killGroupOnDone is the only place the process group is killed. This drops the second context and the command cancel hooks, which both killed the same group. - Keep the already-cancelled error before start, as before. - Define the default ready-file path once per OS in ready_linux.go and ready_darwin.go instead of three copies across transport files. - Remove PR-status wording from the guest-agent doc. It described this patch rather than the code and would rot after merge. --- docs/macos-guest-agent.md | 10 +++------- lib/system/guest_agent/exec.go | 16 +++++++++------- lib/system/guest_agent/ready_darwin.go | 7 +++++++ lib/system/guest_agent/ready_linux.go | 7 +++++++ lib/system/guest_agent/vsock_darwin.go | 2 -- lib/system/guest_agent/vsock_darwin_nocgo.go | 2 -- lib/system/guest_agent/vsock_linux.go | 2 -- 7 files changed, 26 insertions(+), 20 deletions(-) create mode 100644 lib/system/guest_agent/ready_darwin.go create mode 100644 lib/system/guest_agent/ready_linux.go diff --git a/docs/macos-guest-agent.md b/docs/macos-guest-agent.md index 03f523c80..0292d2225 100644 --- a/docs/macos-guest-agent.md +++ b/docs/macos-guest-agent.md @@ -31,8 +31,8 @@ explicitly. The default readiness file is `/var/run/hypeman/guest-agent-ready` A listening system agent does not establish autologin, desktop readiness, TCC permissions, or browser readiness. Root and user desktop agents need a reviewed handoff and explicit session selection before desktop execution is supported. -No desktop agent installer or public session-selector extension is introduced -by this initial patch. +No desktop agent installer or public session-selector extension is part of this +slice. The host API and hypervisor still enforce authorization. The vsock host-CID check is transport admission, not a replacement for instance authority checks. @@ -83,7 +83,7 @@ shutdown policy without executing a real shutdown, explicit network rejection, and transport deadline errors using a local socket pair. These do not prove a live guest AF_VSOCK handshake for this executable. -Remaining draft gates: +Remaining gates: - Provision in a test guest and exercise real host GuestService connectivity. - Live normal API exec/files, readiness and graceful stop/recovery validation; @@ -94,7 +94,3 @@ Remaining draft gates: and privilege/logging security review. - Linux test execution on an appropriate runner, and independent authenticated review. - -The live macOS benchmark guest and its prototype agent are unchanged by this -source patch. Native gRPC integration remains experimental until those gates -are satisfied. diff --git a/lib/system/guest_agent/exec.go b/lib/system/guest_agent/exec.go index ecb40d394..02127d983 100644 --- a/lib/system/guest_agent/exec.go +++ b/lib/system/guest_agent/exec.go @@ -55,16 +55,19 @@ func (s *guestServer) Exec(stream pb.GuestService_ExecServer) error { // executeNoTTY executes command without TTY func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_ExecServer, start *pb.ExecStart) error { - // Run command directly - guest-agent is already running in container namespace - outputCtx, cancel := context.WithCancel(ctx) + // Run command directly - guest-agent is already running in container namespace. + // One cancellable context covers caller cancellation and stream-send errors; + // killGroupOnDone is the only place the command is killed. + ctx, cancel := context.WithCancel(ctx) defer cancel() - // CommandContext must observe both caller cancellation and stream-send errors. - cmd := exec.CommandContext(outputCtx, start.Command[0], start.Command[1:]...) + if err := ctx.Err(); err != nil { + return fmt.Errorf("start command: %w", err) + } + cmd := exec.Command(start.Command[0], start.Command[1:]...) cmd.Env = s.buildEnv(start.Env, false) cmd.Dir = start.Cwd cmd.WaitDelay = 2 * time.Second cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} - cmd.Cancel = func() error { return killProcessGroup(cmd) } var sendMu sync.Mutex cmd.Stdout = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel} cmd.Stderr = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel, stderr: true} @@ -231,8 +234,7 @@ func (w *execStreamWriter) Write(data []byte) (int, error) { func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_ExecServer, start *pb.ExecStart) error { // Run command directly with PTY - guest-agent is already running in container namespace // This ensures PTY and shell are in the same namespace, fixing Ctrl+C signal handling - cmd := exec.CommandContext(ctx, start.Command[0], start.Command[1:]...) - cmd.Cancel = func() error { return killProcessGroup(cmd) } + cmd := exec.Command(start.Command[0], start.Command[1:]...) // Set up environment (TTY mode adds TERM default) cmd.Env = s.buildEnv(start.Env, true) diff --git a/lib/system/guest_agent/ready_darwin.go b/lib/system/guest_agent/ready_darwin.go new file mode 100644 index 000000000..5504a5ae8 --- /dev/null +++ b/lib/system/guest_agent/ready_darwin.go @@ -0,0 +1,7 @@ +//go:build darwin + +package main + +// defaultReadyFilePath is where the guest agent records readiness when +// HYPEMAN_AGENT_READY_FILE is unset. +const defaultReadyFilePath = "/var/run/hypeman/guest-agent-ready" diff --git a/lib/system/guest_agent/ready_linux.go b/lib/system/guest_agent/ready_linux.go new file mode 100644 index 000000000..724c7c0fc --- /dev/null +++ b/lib/system/guest_agent/ready_linux.go @@ -0,0 +1,7 @@ +//go:build linux + +package main + +// defaultReadyFilePath is where the guest agent records readiness when +// HYPEMAN_AGENT_READY_FILE is unset. +const defaultReadyFilePath = "/run/hypeman/guest-agent-ready" diff --git a/lib/system/guest_agent/vsock_darwin.go b/lib/system/guest_agent/vsock_darwin.go index 419de01b4..c0221e37e 100644 --- a/lib/system/guest_agent/vsock_darwin.go +++ b/lib/system/guest_agent/vsock_darwin.go @@ -32,8 +32,6 @@ import ( "time" ) -const defaultReadyFilePath = "/var/run/hypeman/guest-agent-ready" - type vmAddr string func (a vmAddr) Network() string { return "vsock" } diff --git a/lib/system/guest_agent/vsock_darwin_nocgo.go b/lib/system/guest_agent/vsock_darwin_nocgo.go index cd4201845..7fae1f3b7 100644 --- a/lib/system/guest_agent/vsock_darwin_nocgo.go +++ b/lib/system/guest_agent/vsock_darwin_nocgo.go @@ -7,8 +7,6 @@ import ( "net" ) -const defaultReadyFilePath = "/var/run/hypeman/guest-agent-ready" - func listenVsock(uint32) (net.Listener, error) { return nil, fmt.Errorf("Darwin guest vsock requires a cgo-enabled build") } diff --git a/lib/system/guest_agent/vsock_linux.go b/lib/system/guest_agent/vsock_linux.go index 2b821da46..3de52ddc5 100644 --- a/lib/system/guest_agent/vsock_linux.go +++ b/lib/system/guest_agent/vsock_linux.go @@ -8,8 +8,6 @@ import ( "github.com/mdlayher/vsock" ) -const defaultReadyFilePath = "/run/hypeman/guest-agent-ready" - func listenVsock(port uint32) (net.Listener, error) { return vsock.Listen(port, nil) } From 692e96a814d9e3611b02ed6ae6b5d9d2f0fabb3b Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 19:38:33 +0000 Subject: [PATCH 09/15] Test the guest-agent opt-in through the validator and defaults split Runs on every host with fixture capabilities instead of skipping off Apple silicon. --- lib/instances/macos_test.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/lib/instances/macos_test.go b/lib/instances/macos_test.go index 16baa808e..800c57c86 100644 --- a/lib/instances/macos_test.go +++ b/lib/instances/macos_test.go @@ -36,16 +36,15 @@ func TestMacOSRequestDefaultsAndRejections(t *testing.T) { require.ErrorIs(t, prepareMacOSCreate(&CreateInstanceRequest{}, testMacImage(), hypervisor.Capabilities{}), ErrInvalidRequest) } func TestMacOSGuestAgentOptIn(t *testing.T) { - if runtime.GOOS != "darwin" || runtime.GOARCH != "arm64" { - t.Skip("Mac request acceptance requires Apple silicon") - } + caps := hypervisor.Capabilities{SupportsMacOSBoot: true} image := testMacImage() image.MacOS.GuestAgent = true request := CreateInstanceRequest{} - require.NoError(t, prepareMacOSRequest(&request, image, hypervisor.TypeVZ)) + require.NoError(t, validateMacOSCreate(request, image, caps)) + applyMacOSDefaults(&request, image) require.False(t, request.SkipGuestAgent) request = CreateInstanceRequest{SkipGuestAgent: true} - require.NoError(t, prepareMacOSRequest(&request, image, hypervisor.TypeVZ)) + applyMacOSDefaults(&request, image) require.True(t, request.SkipGuestAgent) } From 84d4b16ea49c3ef0a607a69f978e6b7ba93fe3de Mon Sep 17 00:00:00 2001 From: chruffins <23645059+chruffins@users.noreply.github.com> Date: Fri, 9 Oct 2026 19:49:10 +0000 Subject: [PATCH 10/15] Share the bounded-wait helper and the boot-marker and shutdown policy rules - runBounded runs a call on its own goroutine under the drain bound. The command wait, output drain and exit-code send each repeated the goroutine, channel and bound by hand. - StoredMetadata.programMarkerSettled states the macOS rule that no workload start marker is awaited. Hydration and the agent probe gate use it. - Boot markers and the metrics readiness check read GuestAgentEnabled instead of the raw skip flag. For Linux records these are identical. - requestDarwinShutdown and its policy test move to untagged files, so the shutdown policy runs on Linux CI. Only the euid and command stay darwin-only. - The two stalled-client exec tests share one helper. --- lib/instances/metrics.go | 2 +- lib/instances/query.go | 10 +++- lib/system/guest_agent/darwin_test.go | 41 --------------- lib/system/guest_agent/exec.go | 39 ++++++-------- lib/system/guest_agent/service_test.go | 30 ++++------- lib/system/guest_agent/shutdown_darwin.go | 19 ------- lib/system/guest_agent/shutdown_policy.go | 28 ++++++++++ .../guest_agent/shutdown_policy_test.go | 52 +++++++++++++++++++ 8 files changed, 116 insertions(+), 105 deletions(-) create mode 100644 lib/system/guest_agent/shutdown_policy.go create mode 100644 lib/system/guest_agent/shutdown_policy_test.go diff --git a/lib/instances/metrics.go b/lib/instances/metrics.go index 8d86b6cc7..d1ad9a9fa 100644 --- a/lib/instances/metrics.go +++ b/lib/instances/metrics.go @@ -603,7 +603,7 @@ func timeToRunningReadyAt(stored *StoredMetadata) *time.Time { if stored == nil || stored.ProgramStartedAt == nil { return nil } - if stored.SkipGuestAgent || stored.GuestAgentReadyAt == nil { + if !stored.GuestAgentEnabled() || stored.GuestAgentReadyAt == nil { return stored.ProgramStartedAt } if stored.GuestAgentReadyAt.After(*stored.ProgramStartedAt) { diff --git a/lib/instances/query.go b/lib/instances/query.go index 1ad989ec3..c9917bcae 100644 --- a/lib/instances/query.go +++ b/lib/instances/query.go @@ -317,9 +317,15 @@ func (m *manager) hydrateBootMarkersFromLogs(ctx context.Context, stored *Stored return hydrated } +// programMarkerSettled reports that no workload-start marker is awaited or already +// recorded. macOS guests have no Linux workload marker, so it is never awaited. +func (s *StoredMetadata) programMarkerSettled() bool { + return s.MacOS != nil || s.ProgramStartedAt != nil +} + // bootMarkersMissing reports which boot markers stored still lacks. func bootMarkersMissing(stored *StoredMetadata) (needProgram, needAgent bool) { - needProgram = stored.MacOS == nil && stored.ProgramStartedAt == nil + needProgram = !stored.programMarkerSettled() needAgent = stored.GuestAgentEnabled() && stored.GuestAgentReadyAt == nil return needProgram, needAgent } @@ -346,7 +352,7 @@ func (m *manager) applyBootMarkers(ctx context.Context, stored *StoredMetadata) stored.GuestAgentReadyAt = guestAgentReadyAt hydrated = true } - if needAgent && stored.GuestAgentReadyAt == nil && (stored.MacOS != nil || stored.ProgramStartedAt != nil) && m.hydrateGuestAgentReadyFromProbe(ctx, stored) { + if needAgent && stored.GuestAgentReadyAt == nil && stored.programMarkerSettled() && m.hydrateGuestAgentReadyFromProbe(ctx, stored) { hydrated = true } return hydrated diff --git a/lib/system/guest_agent/darwin_test.go b/lib/system/guest_agent/darwin_test.go index d813c2a98..d7582ea53 100644 --- a/lib/system/guest_agent/darwin_test.go +++ b/lib/system/guest_agent/darwin_test.go @@ -4,8 +4,6 @@ package main import ( "context" - "errors" - "syscall" "testing" pb "github.com/kernel/hypeman/lib/guest" @@ -14,45 +12,6 @@ import ( "google.golang.org/grpc/status" ) -func TestDarwinShutdownPolicy(t *testing.T) { - for _, tc := range []struct { - name string - signal int32 - euid int - canceled bool - commandError error - code codes.Code - called bool - }{ - {name: "default orderly shutdown", called: true}, - {name: "explicit orderly shutdown", signal: int32(syscall.SIGTERM), called: true}, - {name: "reject arbitrary signal", signal: int32(syscall.SIGKILL), code: codes.InvalidArgument}, - {name: "reject desktop agent", euid: 501, code: codes.PermissionDenied}, - {name: "canceled before command", canceled: true, code: codes.Canceled}, - {name: "command failure", commandError: errors.New("shutdown refused"), code: codes.Internal, called: true}, - } { - t.Run(tc.name, func(t *testing.T) { - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - if tc.canceled { - cancel() - } - called := false - response, err := requestDarwinShutdown(ctx, &pb.ShutdownRequest{Signal: tc.signal}, tc.euid, func(context.Context) error { - called = true - return tc.commandError - }) - require.Equal(t, tc.called, called) - require.Equal(t, tc.code, status.Code(err)) - if tc.code == codes.OK { - require.NotNil(t, response) - } else { - require.Nil(t, response) - } - }) - } -} - func TestDarwinNetworkReconfigurationUnsupported(t *testing.T) { response, err := (&guestServer{}).ReconfigureNetwork(context.Background(), &pb.ReconfigureNetworkRequest{}) require.Nil(t, response) diff --git a/lib/system/guest_agent/exec.go b/lib/system/guest_agent/exec.go index 02127d983..4d8a94dfa 100644 --- a/lib/system/guest_agent/exec.go +++ b/lib/system/guest_agent/exec.go @@ -99,13 +99,8 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E }() // cmd.Wait also waits for the stdout/stderr copies, which can block in stream.Send. - waitDone := make(chan struct{}) var waitErr error - go func() { - waitErr = cmd.Wait() - close(waitDone) - }() - if err := awaitBounded(ctx, waitDone, s.drainBound()); err != nil { + if err := s.runBounded(ctx, func() { waitErr = cmd.Wait() }); err != nil { return err } // ErrWaitDelay means the command exited but a descendant still held its output. @@ -165,20 +160,28 @@ func killGroupOnDone(ctx context.Context, cmd *exec.Cmd, finished <-chan struct{ // sendExitCode delivers the final status under the same bound as output, so a client // that stopped reading cannot hold the handler past cancellation. func (s *guestServer) sendExitCode(ctx context.Context, stream pb.GuestService_ExecServer, exitCode int32) error { - sent := make(chan struct{}) var err error - go func() { + if boundErr := s.runBounded(ctx, func() { err = stream.Send(&pb.ExecResponse{ Response: &pb.ExecResponse_ExitCode{ExitCode: exitCode}, }) - close(sent) - }() - if boundErr := awaitBounded(ctx, sent, s.drainBound()); boundErr != nil { + }); boundErr != nil { return boundErr } return err } +// runBounded runs fn on its own goroutine and waits for it under the drain bound. +// If the bound expires, fn keeps running and must not rely on the handler still waiting. +func (s *guestServer) runBounded(ctx context.Context, fn func()) error { + done := make(chan struct{}) + go func() { + fn() + close(done) + }() + return awaitBounded(ctx, done, s.drainBound()) +} + // exitCodeOf reports how a command ended. A command that hit its deadline exits 124 // (GNU timeout convention) even though the group kill leaves a process state behind. func exitCodeOf(ctx context.Context, cmd *exec.Cmd, waitErr error) int32 { @@ -311,23 +314,13 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe } }() - waitDone := make(chan struct{}) var waitErr error - go func() { - waitErr = cmd.Wait() - close(waitDone) - }() - if err := awaitBounded(ctx, waitDone, s.drainBound()); err != nil { + if err := s.runBounded(ctx, func() { waitErr = cmd.Wait() }); err != nil { return err } // Wait for all output to be sent - outputDone := make(chan struct{}) - go func() { - wg.Wait() - close(outputDone) - }() - if err := awaitBounded(ctx, outputDone, s.drainBound()); err != nil { + if err := s.runBounded(ctx, wg.Wait); err != nil { return err } diff --git a/lib/system/guest_agent/service_test.go b/lib/system/guest_agent/service_test.go index dae9d4fcd..b925dfa4a 100644 --- a/lib/system/guest_agent/service_test.go +++ b/lib/system/guest_agent/service_test.go @@ -203,13 +203,16 @@ func (s *stalledExecStream) Send(resp *pb.ExecResponse) error { return nil } -func TestGuestServiceExecTimeoutFinishesWhenClientStopsReading(t *testing.T) { +// runStalledExec runs a timed-out command against a client that stopped reading and +// requires the handler to return the bounded-drain error instead of hanging. +func runStalledExec(t *testing.T, command string, exitOnly bool) { + t.Helper() ctx, cancel := context.WithCancel(context.Background()) t.Cleanup(cancel) release := make(chan struct{}) t.Cleanup(func() { close(release) }) - stream := &stalledExecStream{ctx: ctx, release: release, start: &pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ - Command: []string{"/bin/sh", "-c", "yes"}, + stream := &stalledExecStream{ctx: ctx, release: release, stallExitOnly: exitOnly, start: &pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ + Command: []string{"/bin/sh", "-c", command}, TimeoutSeconds: 1, }}}} done := make(chan error, 1) @@ -222,23 +225,12 @@ func TestGuestServiceExecTimeoutFinishesWhenClientStopsReading(t *testing.T) { } } +func TestGuestServiceExecTimeoutFinishesWhenClientStopsReading(t *testing.T) { + runStalledExec(t, "yes", false) +} + func TestGuestServiceExecTimeoutBoundsExitCodeSend(t *testing.T) { - ctx, cancel := context.WithCancel(context.Background()) - t.Cleanup(cancel) - release := make(chan struct{}) - t.Cleanup(func() { close(release) }) - stream := &stalledExecStream{ctx: ctx, release: release, stallExitOnly: true, start: &pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{ - Command: []string{"/bin/sh", "-c", "echo started; sleep 30"}, - TimeoutSeconds: 1, - }}}} - done := make(chan error, 1) - go func() { done <- (&guestServer{drainGrace: 100 * time.Millisecond}).Exec(stream) }() - select { - case err := <-done: - require.ErrorContains(t, err, "did not finish") - case <-time.After(10 * time.Second): - t.Fatal("exit-code send stayed blocked after the command timed out") - } + runStalledExec(t, "echo started; sleep 30", true) } func TestGuestServiceExecTimeoutKillsDescendantAfterShellExits(t *testing.T) { diff --git a/lib/system/guest_agent/shutdown_darwin.go b/lib/system/guest_agent/shutdown_darwin.go index 84a5ff99d..1b8199c66 100644 --- a/lib/system/guest_agent/shutdown_darwin.go +++ b/lib/system/guest_agent/shutdown_darwin.go @@ -6,11 +6,8 @@ import ( "context" "os" "os/exec" - "syscall" pb "github.com/kernel/hypeman/lib/guest" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" ) // Darwin's launchd is not Linux init: sending it SIGTERM is not a shutdown API. @@ -19,19 +16,3 @@ func (s *guestServer) Shutdown(ctx context.Context, req *pb.ShutdownRequest) (*p return exec.CommandContext(ctx, "/sbin/shutdown", "-h", "now").Run() }) } - -func requestDarwinShutdown(ctx context.Context, req *pb.ShutdownRequest, euid int, run func(context.Context) error) (*pb.ShutdownResponse, error) { - if req.Signal != 0 && req.Signal != int32(syscall.SIGTERM) { - return nil, status.Error(codes.InvalidArgument, "Darwin supports only an orderly shutdown, not arbitrary init signals") - } - if euid != 0 { - return nil, status.Error(codes.PermissionDenied, "Darwin shutdown requires the system guest agent running as root") - } - if err := ctx.Err(); err != nil { - return nil, status.FromContextError(err).Err() - } - if err := run(ctx); err != nil { - return nil, status.Errorf(codes.Internal, "Darwin shutdown failed: %v", err) - } - return &pb.ShutdownResponse{}, nil -} diff --git a/lib/system/guest_agent/shutdown_policy.go b/lib/system/guest_agent/shutdown_policy.go new file mode 100644 index 000000000..710a430a0 --- /dev/null +++ b/lib/system/guest_agent/shutdown_policy.go @@ -0,0 +1,28 @@ +package main + +import ( + "context" + "syscall" + + pb "github.com/kernel/hypeman/lib/guest" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// requestDarwinShutdown applies the Darwin shutdown policy. It takes the effective +// uid and the shutdown command as arguments, so the policy runs on every host. +func requestDarwinShutdown(ctx context.Context, req *pb.ShutdownRequest, euid int, run func(context.Context) error) (*pb.ShutdownResponse, error) { + if req.Signal != 0 && req.Signal != int32(syscall.SIGTERM) { + return nil, status.Error(codes.InvalidArgument, "Darwin supports only an orderly shutdown, not arbitrary init signals") + } + if euid != 0 { + return nil, status.Error(codes.PermissionDenied, "Darwin shutdown requires the system guest agent running as root") + } + if err := ctx.Err(); err != nil { + return nil, status.FromContextError(err).Err() + } + if err := run(ctx); err != nil { + return nil, status.Errorf(codes.Internal, "Darwin shutdown failed: %v", err) + } + return &pb.ShutdownResponse{}, nil +} diff --git a/lib/system/guest_agent/shutdown_policy_test.go b/lib/system/guest_agent/shutdown_policy_test.go new file mode 100644 index 000000000..c2c5873cd --- /dev/null +++ b/lib/system/guest_agent/shutdown_policy_test.go @@ -0,0 +1,52 @@ +package main + +import ( + "context" + "errors" + "syscall" + "testing" + + pb "github.com/kernel/hypeman/lib/guest" + "github.com/stretchr/testify/require" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +func TestDarwinShutdownPolicy(t *testing.T) { + for _, tc := range []struct { + name string + signal int32 + euid int + canceled bool + commandError error + code codes.Code + called bool + }{ + {name: "default orderly shutdown", called: true}, + {name: "explicit orderly shutdown", signal: int32(syscall.SIGTERM), called: true}, + {name: "reject arbitrary signal", signal: int32(syscall.SIGKILL), code: codes.InvalidArgument}, + {name: "reject desktop agent", euid: 501, code: codes.PermissionDenied}, + {name: "canceled before command", canceled: true, code: codes.Canceled}, + {name: "command failure", commandError: errors.New("shutdown refused"), code: codes.Internal, called: true}, + } { + t.Run(tc.name, func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + if tc.canceled { + cancel() + } + called := false + response, err := requestDarwinShutdown(ctx, &pb.ShutdownRequest{Signal: tc.signal}, tc.euid, func(context.Context) error { + called = true + return tc.commandError + }) + require.Equal(t, tc.called, called) + require.Equal(t, tc.code, status.Code(err)) + if tc.code == codes.OK { + require.NotNil(t, response) + } else { + require.Nil(t, response) + } + }) + } +} From 912c28c79bc301b7bc2dd5f592e59b28625a9200 Mon Sep 17 00:00:00 2001 From: chris lee Date: Fri, 9 Oct 2026 16:23:34 -0400 Subject: [PATCH 11/15] Share exec setup and bounded completion across TTY modes --- lib/system/guest_agent/exec.go | 94 +++++++++---------- .../guest_agent/exec_start_cancel_test.go | 23 +++++ 2 files changed, 67 insertions(+), 50 deletions(-) create mode 100644 lib/system/guest_agent/exec_start_cancel_test.go diff --git a/lib/system/guest_agent/exec.go b/lib/system/guest_agent/exec.go index 4d8a94dfa..2ebb6e6d3 100644 --- a/lib/system/guest_agent/exec.go +++ b/lib/system/guest_agent/exec.go @@ -47,27 +47,20 @@ func (s *guestServer) Exec(stream pb.GuestService_ExecServer) error { defer cancel() } + ctx, cancel := context.WithCancel(ctx) + defer cancel() + cmd, err := s.execCommand(ctx, start) + if err != nil { + return err + } if start.Tty { - return s.executeTTY(ctx, stream, start) + return s.executeTTY(ctx, cancel, stream, start, cmd) } - return s.executeNoTTY(ctx, stream, start) + return s.executeNoTTY(ctx, cancel, stream, cmd) } // executeNoTTY executes command without TTY -func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_ExecServer, start *pb.ExecStart) error { - // Run command directly - guest-agent is already running in container namespace. - // One cancellable context covers caller cancellation and stream-send errors; - // killGroupOnDone is the only place the command is killed. - ctx, cancel := context.WithCancel(ctx) - defer cancel() - if err := ctx.Err(); err != nil { - return fmt.Errorf("start command: %w", err) - } - cmd := exec.Command(start.Command[0], start.Command[1:]...) - cmd.Env = s.buildEnv(start.Env, false) - cmd.Dir = start.Cwd - cmd.WaitDelay = 2 * time.Second - cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} +func (s *guestServer) executeNoTTY(ctx context.Context, cancel context.CancelFunc, stream pb.GuestService_ExecServer, cmd *exec.Cmd) error { var sendMu sync.Mutex cmd.Stdout = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel} cmd.Stderr = &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel, stderr: true} @@ -98,7 +91,28 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E } }() - // cmd.Wait also waits for the stdout/stderr copies, which can block in stream.Send. + return s.finishExec(ctx, stream, cmd, nil) +} + +// execCommand owns setup and pre-start cancellation for both I/O modes. +func (s *guestServer) execCommand(ctx context.Context, start *pb.ExecStart) (*exec.Cmd, error) { + if err := ctx.Err(); err != nil { + return nil, fmt.Errorf("start command: %w", err) + } + cmd := exec.CommandContext(ctx, start.Command[0], start.Command[1:]...) + cmd.Cancel = func() error { return killProcessGroup(cmd) } + cmd.WaitDelay = 2 * time.Second + cmd.Env = s.buildEnv(start.Env, start.Tty) + cmd.Dir = start.Cwd + if !start.Tty { + cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} + } + return cmd, nil +} + +// finishExec owns bounded wait, output drain and final status for both I/O modes. +func (s *guestServer) finishExec(ctx context.Context, stream pb.GuestService_ExecServer, cmd *exec.Cmd, outputDone <-chan struct{}) error { + // Non-TTY cmd.Wait also waits for copies, which can block in stream.Send. var waitErr error if err := s.runBounded(ctx, func() { waitErr = cmd.Wait() }); err != nil { return err @@ -112,6 +126,11 @@ func (s *guestServer) executeNoTTY(ctx context.Context, stream pb.GuestService_E } } + if outputDone != nil { + if err := awaitBounded(ctx, outputDone, s.drainBound()); err != nil { + return err + } + } exitCode := exitCodeOf(ctx, cmd, waitErr) log.Printf("[guest-agent] command finished with exit code: %d", exitCode) @@ -234,15 +253,7 @@ func (w *execStreamWriter) Write(data []byte) (int, error) { } // executeTTY executes command with TTY -func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_ExecServer, start *pb.ExecStart) error { - // Run command directly with PTY - guest-agent is already running in container namespace - // This ensures PTY and shell are in the same namespace, fixing Ctrl+C signal handling - cmd := exec.Command(start.Command[0], start.Command[1:]...) - - // Set up environment (TTY mode adds TERM default) - cmd.Env = s.buildEnv(start.Env, true) - - cmd.Dir = start.Cwd +func (s *guestServer) executeTTY(ctx context.Context, cancel context.CancelFunc, stream pb.GuestService_ExecServer, start *pb.ExecStart, cmd *exec.Cmd) error { // Set up initial window size (use defaults if not specified) ws := &pty.Winsize{ @@ -269,8 +280,8 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe // Mutex to protect concurrent stream.Send calls (gRPC streams are not thread-safe) var sendMu sync.Mutex - // Use WaitGroup to ensure all output is sent before exit code - var wg sync.WaitGroup + outputDone := make(chan struct{}) + output := &execStreamWriter{stream: stream, mu: &sendMu, cancel: cancel} // Handle stdin and resize in background go func() { @@ -295,18 +306,15 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe }() // Stream output - wg.Add(1) go func() { - defer wg.Done() + defer close(outputDone) buf := make([]byte, 32*1024) for { n, err := ptmx.Read(buf) if n > 0 { - sendMu.Lock() - stream.Send(&pb.ExecResponse{ - Response: &pb.ExecResponse_Stdout{Stdout: buf[:n]}, - }) - sendMu.Unlock() + if _, sendErr := output.Write(buf[:n]); sendErr != nil { + return + } } if err != nil { return @@ -314,21 +322,7 @@ func (s *guestServer) executeTTY(ctx context.Context, stream pb.GuestService_Exe } }() - var waitErr error - if err := s.runBounded(ctx, func() { waitErr = cmd.Wait() }); err != nil { - return err - } - - // Wait for all output to be sent - if err := s.runBounded(ctx, wg.Wait); err != nil { - return err - } - - exitCode := exitCodeOf(ctx, cmd, waitErr) - - log.Printf("[guest-agent] TTY command finished with exit code: %d", exitCode) - - return s.sendExitCode(ctx, stream, exitCode) + return s.finishExec(ctx, stream, cmd, outputDone) } // buildEnv constructs environment variables by merging provided env with defaults. diff --git a/lib/system/guest_agent/exec_start_cancel_test.go b/lib/system/guest_agent/exec_start_cancel_test.go new file mode 100644 index 000000000..96b5398ae --- /dev/null +++ b/lib/system/guest_agent/exec_start_cancel_test.go @@ -0,0 +1,23 @@ +package main + +import ( + "context" + "fmt" + "path/filepath" + "testing" + + pb "github.com/kernel/hypeman/lib/guest" + "github.com/stretchr/testify/require" +) + +func TestCancelledExecDoesNotAttemptStart(t *testing.T) { + for _, tty := range []bool{false, true} { + t.Run(fmt.Sprintf("tty=%t", tty), func(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + stream := &stalledExecStream{ctx: ctx, start: &pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{Command: []string{filepath.Join(t.TempDir(), "nonexistent-executable")}, Tty: tty}}}} + err := (&guestServer{}).Exec(stream) + require.ErrorIs(t, err, context.Canceled, "a canceled request must be rejected before trying to start the command") + }) + } +} From c9469c24b11bf0bced8356f9282ae4aadd7c9c0f Mon Sep 17 00:00:00 2001 From: chris lee Date: Fri, 9 Oct 2026 17:45:01 -0400 Subject: [PATCH 12/15] Test exec output-send cancellation for both I/O modes --- .../guest_agent/exec_send_cancel_test.go | 48 +++++++++++++++++++ 1 file changed, 48 insertions(+) create mode 100644 lib/system/guest_agent/exec_send_cancel_test.go diff --git a/lib/system/guest_agent/exec_send_cancel_test.go b/lib/system/guest_agent/exec_send_cancel_test.go new file mode 100644 index 000000000..9419f3447 --- /dev/null +++ b/lib/system/guest_agent/exec_send_cancel_test.go @@ -0,0 +1,48 @@ +package main + +import ( + "context" + "errors" + "os" + "path/filepath" + "strconv" + "strings" + "syscall" + "testing" + "time" + + pb "github.com/kernel/hypeman/lib/guest" + "github.com/stretchr/testify/require" +) + +type failedOutputExecStream struct { + stalledExecStream + failure error +} + +func (s *failedOutputExecStream) Send(*pb.ExecResponse) error { return s.failure } + +func TestExecSendFailureCancelsProcessInBothModes(t *testing.T) { + for _, tty := range []bool{false, true} { + name := "stream" + if tty { + name = "tty" + } + t.Run(name, func(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + release := make(chan struct{}) + defer close(release) + pidFile := filepath.Join(t.TempDir(), "pid") + failure := errors.New("synthetic send failure") + stream := &failedOutputExecStream{stalledExecStream: stalledExecStream{ctx: ctx, release: release, start: &pb.ExecRequest{Request: &pb.ExecRequest_Start{Start: &pb.ExecStart{Tty: tty, Command: []string{"/bin/sh", "-c", `echo $$ > "$1"; printf ready; exec /bin/sleep 30`, "sh", pidFile}}}}}, failure: failure} + require.Error(t, (&guestServer{}).Exec(stream)) + require.NoError(t, ctx.Err(), "send failure must cancel without waiting for the caller deadline") + data, err := os.ReadFile(pidFile) + require.NoError(t, err) + pid, err := strconv.Atoi(strings.TrimSpace(string(data))) + require.NoError(t, err) + require.Eventually(t, func() bool { return errors.Is(syscall.Kill(pid, 0), syscall.ESRCH) }, time.Second, 10*time.Millisecond, "failed output must not leave the command alive") + }) + } +} From 64a8c137c9e14846c5230fc008e8d0b26544790e Mon Sep 17 00:00:00 2001 From: chris lee Date: Sat, 10 Oct 2026 22:20:29 -0400 Subject: [PATCH 13/15] fix: validate macOS root readiness and retain shutdown grace --- docs/macos-guest-agent.md | 27 ++++++++++++++++--- .../guest_agent_ready_command_test.go | 13 +++++++++ lib/instances/macos_test.go | 5 ++-- lib/instances/query.go | 10 ++++++- lib/instances/shutdown_wait_test.go | 18 +++++++++++++ lib/instances/stop.go | 18 ++++++++----- 6 files changed, 78 insertions(+), 13 deletions(-) create mode 100644 lib/instances/guest_agent_ready_command_test.go create mode 100644 lib/instances/shutdown_wait_test.go diff --git a/docs/macos-guest-agent.md b/docs/macos-guest-agent.md index 0292d2225..61f8a3a68 100644 --- a/docs/macos-guest-agent.md +++ b/docs/macos-guest-agent.md @@ -83,11 +83,32 @@ shutdown policy without executing a real shutdown, explicit network rejection, and transport deadline errors using a local socket pair. These do not prove a live guest AF_VSOCK handshake for this executable. +## Isolated live validation (2026-10-10) + +A manually provisioned root LaunchDaemon in the isolated QA guest survives normal +API cold boots. Authenticated normal API exec reports UID0 in both TTY modes; +private root-owned copy-to/from round trips and command/descendant timeout cleanup +pass repeated race-instrumented client tests. The readiness probe uses +`/usr/bin/true` on Darwin (`/bin/true` remains the Linux command); a live regression +first failed with the Linux path, then passed after this correction. A normal GET +persisted `GuestAgentReadyAt` while leaving `ProgramStartedAt` absent. + +Darwin can disconnect vsock while shutdown is in progress, before the RPC reply +arrives. The host therefore retains the configured grace period on Darwin even +when that reply is unavailable, instead of Linux's short failure fallback. Success +still requires the owned VMM process to exit. A live normal API stop confirmed +shutdown RPC acceptance, VMM exit and closed storage without forced fallback after +the correction. HTTP stop200 or `Stopped` alone is not a build/export receipt. + +This is manual provisioning of one QA instance, not automatic image installation, +public macOS builds, output publication or reusable-image sanitation proof. The +immutable imported base remains unchanged and does not acquire this capability. + Remaining gates: -- Provision in a test guest and exercise real host GuestService connectivity. -- Live normal API exec/files, readiness and graceful stop/recovery validation; - guest-agent version compatibility and bounded readiness wait semantics. +- Reusable image provisioning and guest-agent version compatibility. +- Strict build/export stop receipts and API recovery with the root deployment; + bounded readiness wait semantics. - Root/desktop session authorization and image provisioning. - Broader backpressure/large-output validation of bounded non-TTY streaming, PTY/disconnect and descendant-process cleanup, transfer failure/size handling, diff --git a/lib/instances/guest_agent_ready_command_test.go b/lib/instances/guest_agent_ready_command_test.go new file mode 100644 index 000000000..ef9ad91e1 --- /dev/null +++ b/lib/instances/guest_agent_ready_command_test.go @@ -0,0 +1,13 @@ +package instances + +import ( + "testing" + + "github.com/kernel/hypeman/lib/images" + "github.com/stretchr/testify/require" +) + +func TestGuestAgentReadyExecutableMatchesGuestOS(t *testing.T) { + require.Equal(t, "/bin/true", guestAgentReadyExecutable(&StoredMetadata{})) + require.Equal(t, "/usr/bin/true", guestAgentReadyExecutable(&StoredMetadata{MacOS: &images.MacOSImage{}})) +} diff --git a/lib/instances/macos_test.go b/lib/instances/macos_test.go index 800c57c86..566958b15 100644 --- a/lib/instances/macos_test.go +++ b/lib/instances/macos_test.go @@ -40,11 +40,10 @@ func TestMacOSGuestAgentOptIn(t *testing.T) { image := testMacImage() image.MacOS.GuestAgent = true request := CreateInstanceRequest{} - require.NoError(t, validateMacOSCreate(request, image, caps)) - applyMacOSDefaults(&request, image) + require.NoError(t, prepareMacOSCreate(&request, image, caps)) require.False(t, request.SkipGuestAgent) request = CreateInstanceRequest{SkipGuestAgent: true} - applyMacOSDefaults(&request, image) + require.NoError(t, prepareMacOSCreate(&request, image, caps)) require.True(t, request.SkipGuestAgent) } diff --git a/lib/instances/query.go b/lib/instances/query.go index c9917bcae..27a0afeb1 100644 --- a/lib/instances/query.go +++ b/lib/instances/query.go @@ -477,13 +477,21 @@ func probeGuestAgentReady(ctx context.Context, stored *StoredMetadata) bool { defer cancel() exit, err := guest.ExecIntoInstance(probeCtx, dialer, guest.ExecOptions{ - Command: []string{"/bin/true"}, + Command: []string{guestAgentReadyExecutable(stored)}, Timeout: int32(guestAgentReadyProbeTimeout / time.Second), WaitForAgent: guestAgentReadyProbeWait, }) return err == nil && exit != nil && exit.Code == 0 } +func guestAgentReadyExecutable(stored *StoredMetadata) string { + // Darwin provides true in /usr/bin, not /bin. Keep the Linux probe unchanged. + if stored.MacOS != nil { + return "/usr/bin/true" + } + return "/bin/true" +} + // appLogPathsForMarkerScan returns app log paths in chronological order // (oldest rotated file to newest active file). func (m *manager) appLogPathsForMarkerScan(id string) []string { diff --git a/lib/instances/shutdown_wait_test.go b/lib/instances/shutdown_wait_test.go new file mode 100644 index 000000000..1f9449680 --- /dev/null +++ b/lib/instances/shutdown_wait_test.go @@ -0,0 +1,18 @@ +package instances + +import ( + "testing" + "time" + + "github.com/kernel/hypeman/lib/images" + "github.com/stretchr/testify/require" +) + +func TestGuestShutdownWaitPreservesDarwinGraceOnLostReply(t *testing.T) { + linux := &Instance{} + mac := &Instance{StoredMetadata: StoredMetadata{MacOS: &images.MacOSImage{}}} + require.Equal(t, 500*time.Millisecond, guestShutdownWaitTimeout(linux, 30, false)) + require.Equal(t, 30*time.Second, guestShutdownWaitTimeout(linux, 30, true)) + require.Equal(t, 30*time.Second, guestShutdownWaitTimeout(mac, 30, false)) + require.Equal(t, 30*time.Second, guestShutdownWaitTimeout(mac, 30, true)) +} diff --git a/lib/instances/stop.go b/lib/instances/stop.go index 81103fdb8..3a2973fdf 100644 --- a/lib/instances/stop.go +++ b/lib/instances/stop.go @@ -84,7 +84,7 @@ func (m *manager) tryGracefulGuestShutdown(ctx context.Context, inst *Instance, // Drop potentially stale pooled connection and retry once. guest.CloseConn(dialer.Key()) if retryErr := sendShutdown(); retryErr != nil { - log.WarnContext(ctx, "shutdown RPC failed; falling back to hypervisor shutdown", "instance_id", inst.Id, "error", retryErr) + log.WarnContext(ctx, "shutdown RPC response unavailable; awaiting guest power-off", "instance_id", inst.Id, "error", retryErr) } else { shutdownSent = true } @@ -92,11 +92,7 @@ func (m *manager) tryGracefulGuestShutdown(ctx context.Context, inst *Instance, shutdownSent = true } - waitTimeout := time.Duration(stopTimeout) * time.Second - if !shutdownSent && waitTimeout > shutdownFailureFallbackWait { - // If we couldn't signal the guest, don't burn the full graceful timeout. - waitTimeout = shutdownFailureFallbackWait - } + waitTimeout := guestShutdownWaitTimeout(inst, stopTimeout, shutdownSent) if WaitForProcessExit(pid, waitTimeout) { log.DebugContext(ctx, "VM shut down gracefully", "instance_id", inst.Id) @@ -107,6 +103,16 @@ func (m *manager) tryGracefulGuestShutdown(ctx context.Context, inst *Instance, return false } +func guestShutdownWaitTimeout(inst *Instance, stopTimeout int, shutdownSent bool) time.Duration { + wait := time.Duration(stopTimeout) * time.Second + // Darwin may close vsock before the reply while power-off is in progress. + // Preserve its grace period; success still requires actual owned VMM exit. + if !shutdownSent && inst.MacOS == nil && wait > shutdownFailureFallbackWait { + return shutdownFailureFallbackWait + } + return wait +} + // stopInstance gracefully stops an active instance. // Flow: send Shutdown RPC -> wait for VM to power off -> // fall back to hypervisor shutdown -> final SIGKILL if still alive. From 05bd2fff088de610735017458058aad7aac8ee3a Mon Sep 17 00:00:00 2001 From: chris lee Date: Sun, 11 Oct 2026 00:38:25 -0400 Subject: [PATCH 14/15] Schedule Darwin power-off after shutdown acknowledgement --- lib/system/guest_agent/shutdown_darwin.go | 25 ++++++++++++++++++- .../shutdown_schedule_darwin_test.go | 19 ++++++++++++++ 2 files changed, 43 insertions(+), 1 deletion(-) create mode 100644 lib/system/guest_agent/shutdown_schedule_darwin_test.go diff --git a/lib/system/guest_agent/shutdown_darwin.go b/lib/system/guest_agent/shutdown_darwin.go index 1b8199c66..fcca21846 100644 --- a/lib/system/guest_agent/shutdown_darwin.go +++ b/lib/system/guest_agent/shutdown_darwin.go @@ -4,8 +4,10 @@ package main import ( "context" + "log" "os" "os/exec" + "syscall" pb "github.com/kernel/hypeman/lib/guest" ) @@ -13,6 +15,27 @@ import ( // Darwin's launchd is not Linux init: sending it SIGTERM is not a shutdown API. func (s *guestServer) Shutdown(ctx context.Context, req *pb.ShutdownRequest) (*pb.ShutdownResponse, error) { return requestDarwinShutdown(ctx, req, os.Geteuid(), func(ctx context.Context) error { - return exec.CommandContext(ctx, "/sbin/shutdown", "-h", "now").Run() + if err := ctx.Err(); err != nil { + return err + } + command := darwinShutdownCommand() + if err := command.Start(); err != nil { + return err + } + go func() { + if err := command.Wait(); err != nil { + log.Printf("scheduled Darwin shutdown failed: %v", err) + } + }() + return nil }) } + +func darwinShutdownCommand() *exec.Cmd { + // A short detached delay lets gRPC send the acceptance reply before macOS + // tears down vsock. The host must still independently confirm owned VMM exit. + // Once accepted, shutdown is not cancelled by the RPC connection disappearing. + command := exec.Command("/bin/sh", "-c", "/bin/sleep 1; exec /sbin/shutdown -h now") + command.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} + return command +} diff --git a/lib/system/guest_agent/shutdown_schedule_darwin_test.go b/lib/system/guest_agent/shutdown_schedule_darwin_test.go new file mode 100644 index 000000000..017c5d352 --- /dev/null +++ b/lib/system/guest_agent/shutdown_schedule_darwin_test.go @@ -0,0 +1,19 @@ +//go:build darwin + +package main + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestDarwinShutdownDefersPowerOffUntilReply(t *testing.T) { + // Inspect only: never execute a host shutdown command in a test. + command := darwinShutdownCommand() + require.Equal(t, "/bin/sh", command.Path) + require.Equal(t, []string{"/bin/sh", "-c", "/bin/sleep 1; exec /sbin/shutdown -h now"}, command.Args) + require.NotNil(t, command.SysProcAttr) + require.True(t, command.SysProcAttr.Setpgid) + require.Nil(t, command.Cancel, "accepted shutdown must survive RPC cancellation") +} From c336b304e6d0f60ba56eeffa293cd9e2d75b2fb3 Mon Sep 17 00:00:00 2001 From: chris lee Date: Sun, 11 Oct 2026 00:38:25 -0400 Subject: [PATCH 15/15] Keep macOS vsock accept errors retryable for gRPC --- lib/system/guest_agent/vsock_darwin.go | 9 ++++++++- lib/system/guest_agent/vsock_darwin_test.go | 13 +++++++++++++ 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/lib/system/guest_agent/vsock_darwin.go b/lib/system/guest_agent/vsock_darwin.go index c0221e37e..048267d24 100644 --- a/lib/system/guest_agent/vsock_darwin.go +++ b/lib/system/guest_agent/vsock_darwin.go @@ -108,8 +108,15 @@ func (l *vmListener) Accept() (net.Conn, error) { return &vmConn{file: os.NewFile(uintptr(fd), "guest-vsock"), local: l.Addr()}, nil } if !errors.Is(err, syscall.EAGAIN) && !errors.Is(err, syscall.EINTR) { - return nil, fmt.Errorf("accept vsock: %w", err) + return nil, acceptError(err) } time.Sleep(10 * time.Millisecond) } } + +// acceptError keeps the net.OpError Temporary contract. grpc-go retries only a +// direct Temporary() assertion, so wrapping with fmt.Errorf would stop the server +// on transient EMFILE/ENFILE/ECONNABORTED. +func acceptError(err error) error { + return &net.OpError{Op: "accept", Net: "vsock", Err: err} +} diff --git a/lib/system/guest_agent/vsock_darwin_test.go b/lib/system/guest_agent/vsock_darwin_test.go index 732bf71e8..4fc2829f4 100644 --- a/lib/system/guest_agent/vsock_darwin_test.go +++ b/lib/system/guest_agent/vsock_darwin_test.go @@ -6,6 +6,7 @@ import ( "errors" "net" "os" + "syscall" "testing" "time" @@ -33,3 +34,15 @@ func TestVsockReadTimeoutImplementsNetError(t *testing.T) { t.Fatalf("transport timeout must implement net.Error: %v", err) } } + +func TestVsockAcceptErrorsRemainGRPCTemporary(t *testing.T) { + for _, errno := range []syscall.Errno{syscall.EMFILE, syscall.ENFILE, syscall.ECONNABORTED} { + err := acceptError(errno) + if ne, ok := err.(interface{ Temporary() bool }); !ok || !ne.Temporary() { + t.Fatalf("%v: grpc direct Temporary assertion failed for %T", errno, err) + } + if !errors.Is(err, errno) { + t.Fatalf("%v: lost errno", errno) + } + } +}