diff --git a/README.md b/README.md index 70d30c8..a6ab041 100644 --- a/README.md +++ b/README.md @@ -343,4 +343,35 @@ client-assigned handshake, certificate, and socket metadata retain precedence. `Directory` sets the runtime working directory. The additional supervisor start must be included in startup measurements before selecting session reuse. Streaming stress and real-runtime trust/environment compatibility remain -separate gates before runtime cutover. +separate gates before runtime cutover. The streaming probes below cover the +owned transport; real-runtime compatibility remains outstanding. + +## Streaming stress under process ownership + +Run the supervisor-backed transport probes with: + +```sh +mise exec -- go test -race -v ./internal/streamprobe +``` + +Each duplex run transfers 100 MiB of binary stdin and checks 100 MiB on each +output channel, followed by distinct binary tails and exactly one nonzero exit. +Incremental SHA-256 comparisons use fixed-size buffers rather than retaining the +payload. Separate runs pace the host consumer and plugin reader. Each side has +one stream sender; an independent control stream releases the paused reader. + +Backpressure probes pause either the plugin input reader or the host output +consumer. They check that the bulk sender cannot finish, record host and plugin +live Go heap after GC, and enforce a 32 MiB growth budget over the connected +baseline. These are retained-heap checkpoints, not peak RSS measurements or a +memory bound for arbitrary runtime implementations. Cancellation must unblock +and join the sender; a control RPC must still succeed. Additional probes cover +command early exit and an actual plugin process crash while stdin is active. + +The fixture uses the real go-plugin/gRPC transport and the SDK supervisor, with +runtime executable paths containing spaces and Unicode. CI runs these probes +under the race detector on Linux, macOS, and Windows; their output is retained +in the existing `runtime-tests-*` artifacts. The fixture's command outcomes are +synthetic and do not certify a real runtime's child-process or environment +behavior. Real-runtime compatibility and supervisor startup measurements remain +separate integration gates. diff --git a/internal/streamfixture/main.go b/internal/streamfixture/main.go new file mode 100644 index 0000000..5910d33 --- /dev/null +++ b/internal/streamfixture/main.go @@ -0,0 +1,153 @@ +// Stream fixture exercises transport flow control without buffering whole payloads. +package main + +import ( + "encoding/json" + "os" + "runtime" + "sync" + "time" + + "github.com/devsy-org/devsy-runtime-sdk/runtimev1" + "github.com/devsy-org/devsy-runtime-sdk/server" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +type stream = grpc.BidiStreamingServer[runtimev1.ExecClientMessage, runtimev1.ExecServerMessage] + +type fixture struct { + runtimev1.UnimplementedRuntimeDriverServer + gate chan struct{} + release sync.Once +} + +func (f *fixture) Exec(s stream) error { + first, err := s.Recv() + if err != nil { + return err + } + if first.GetStart() == nil || len(first.GetStart().GetArgv()) != 1 { + return status.Error(codes.InvalidArgument, "expected fixture command") + } + return f.run(s, first.GetStart().GetArgv()[0]) +} + +func (f *fixture) run(s stream, command string) error { + switch command { + case "memory": + return memory(s) + case "release": + f.release.Do(func() { close(f.gate) }) + return exit(s, 0) + case "early-exit", "crash": + return interrupt(s, command) + case "paused-input": + return f.paused(s) + case "duplex": + return duplex(s, false) + default: + return status.Error(codes.InvalidArgument, "unknown fixture command") + } +} + +func duplex(s stream, slow bool) error { + var transferred int + for { + message, err := s.Recv() + if err != nil { + return err + } + if message.GetCloseStdin() != nil { + return tail(s) + } + data := message.GetStdin() + if err := echo(s, data); err != nil { + return err + } + transferred += len(data) + if slow && transferred%(1<<20) == 0 { + time.Sleep(time.Millisecond) + } + } +} + +func tail(s stream) error { + if err := output(s, []byte("stdout tail\x00\xff"), false); err != nil { + return err + } + if err := output(s, []byte("stderr tail\xff\x00"), true); err != nil { + return err + } + return exit(s, 7) +} + +func interrupt(s stream, command string) error { + if err := output(s, []byte("ready"), false); err != nil { + return err + } + // A client data frame proves the sender is active before exit or crash. + if _, err := s.Recv(); err != nil { + return err + } + if command == "crash" { + os.Exit(24) + } + return exit(s, 0) +} + +func output(s stream, data []byte, stderr bool) error { + chunk := &runtimev1.OutputChunk{Data: data} + message := &runtimev1.ExecServerMessage{ + Payload: &runtimev1.ExecServerMessage_Stdout{Stdout: chunk}, + } + if stderr { + message.Payload = &runtimev1.ExecServerMessage_Stderr{Stderr: chunk} + } + return s.Send(message) +} + +func exit(s stream, code int32) error { + return s.Send(&runtimev1.ExecServerMessage{Payload: &runtimev1.ExecServerMessage_Exit{ + Exit: &runtimev1.ExecExit{ExitCode: code}, + }}) +} + +func main() { server.Serve(&fixture{gate: make(chan struct{})}) } + +func memory(s stream) error { + runtime.GC() + var stats runtime.MemStats + runtime.ReadMemStats(&stats) + data, err := json.Marshal(stats.HeapAlloc) + if err != nil { + return err + } + if err := output(s, data, false); err != nil { + return err + } + return exit(s, 0) +} + +func (f *fixture) paused(s stream) error { + if err := output(s, []byte("ready"), false); err != nil { + return err + } + select { + case <-f.gate: + return duplex(s, true) + case <-s.Context().Done(): + return status.FromContextError(s.Context().Err()).Err() + } +} + +func echo(s stream, data []byte) error { + if len(data) == 0 || len(data) > runtimev1.ChunkSize { + return status.Error(codes.InvalidArgument, "expected bounded stdin") + } + if err := output(s, data, false); err != nil { + return err + } + return output(s, data, true) +} diff --git a/internal/streamprobe/fixture_test.go b/internal/streamprobe/fixture_test.go new file mode 100644 index 0000000..74d8510 --- /dev/null +++ b/internal/streamprobe/fixture_test.go @@ -0,0 +1,113 @@ +package streamprobe_test + +import ( + "context" + "os" + "os/exec" + "path/filepath" + "testing" + "time" + + sdkplugin "github.com/devsy-org/devsy-runtime-sdk/plugin" + "github.com/devsy-org/devsy-runtime-sdk/runtimev1" + "github.com/devsy-org/devsy-runtime-sdk/supervisor" + "github.com/hashicorp/go-hclog" + hplugin "github.com/hashicorp/go-plugin" +) + +var executable string + +func TestMain(m *testing.M) { + if len(os.Args) > 1 && os.Args[1] == "--stream-supervisor" { + supervisor.Main(os.Args[2:]) + } + dir, err := os.MkdirTemp("", "runtime stream λ ") + if err != nil { + panic(err) + } + executable = filepath.Join(dir, "stream runtime.exe") + // #nosec G204 -- Fixed fixture package, output in a test-owned directory. + cmd := exec.Command("go", "build", "-race", "-o", executable, "../streamfixture") + cmd.Stdout, cmd.Stderr = os.Stdout, os.Stderr + if err := cmd.Run(); err != nil { + _ = os.RemoveAll(dir) + os.Exit(1) + } + code := m.Run() + _ = os.RemoveAll(dir) + os.Exit(code) +} + +func launch(t *testing.T) (*hplugin.Client, runtimev1.RuntimeDriverClient) { + t.Helper() + helper, err := os.Executable() + if err != nil { + t.Fatal(err) + } + socketDir, err := os.MkdirTemp("", "ds-") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + if err := os.RemoveAll(socketDir); err != nil { + t.Errorf("socket cleanup: %v", err) + } + }) + client := hplugin.NewClient(&hplugin.ClientConfig{ + HandshakeConfig: sdkplugin.Handshake(), + VersionedPlugins: map[int]hplugin.PluginSet{ + sdkplugin.ProtocolVersion: sdkplugin.ClientPlugins(), + }, + AllowedProtocols: []hplugin.Protocol{hplugin.ProtocolGRPC}, + RunnerFunc: supervisor.Runner(supervisor.Options{ + SupervisorBinary: helper, + SupervisorArgs: []string{"--stream-supervisor"}, + RuntimeBinary: executable, + }), + StartTimeout: 30 * time.Second, + Logger: hclog.NewNullLogger(), + Stderr: os.Stderr, SyncStderr: os.Stderr, + UnixSocketConfig: &hplugin.UnixSocketConfig{TempDir: socketDir}, + }) + t.Cleanup(client.Kill) + transport, err := client.Client() + if err != nil { + t.Fatal(err) + } + raw, err := transport.Dispense(sdkplugin.Name) + if err != nil { + t.Fatal(err) + } + driver, ok := raw.(runtimev1.RuntimeDriverClient) + if !ok { + t.Fatalf("unexpected driver %T", raw) + } + return client, driver +} + +func deadline(t *testing.T) context.Context { + t.Helper() + ctx, cancel := context.WithTimeout(t.Context(), 90*time.Second) + t.Cleanup(cancel) + return ctx +} + +func start( + ctx context.Context, + t *testing.T, + driver runtimev1.RuntimeDriverClient, + command string, +) runtimev1.RuntimeDriver_ExecClient { + t.Helper() + s, err := driver.Exec(ctx) + if err != nil { + t.Fatal(err) + } + err = s.Send(&runtimev1.ExecClientMessage{Payload: &runtimev1.ExecClientMessage_Start{ + Start: &runtimev1.ExecStart{WorkspaceId: "probe", Argv: []string{command}}, + }}) + if err != nil { + t.Fatal(err) + } + return s +} diff --git a/internal/streamprobe/stress_test.go b/internal/streamprobe/stress_test.go new file mode 100644 index 0000000..4110cde --- /dev/null +++ b/internal/streamprobe/stress_test.go @@ -0,0 +1,195 @@ +package streamprobe_test + +import ( + "bytes" + "context" + "io" + "testing" + + "github.com/devsy-org/devsy-runtime-sdk/runtimev1" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +const ( + pausedInput = "paused-input" + slowHost = "slow-host" + duplexCommand = "duplex" +) + +func TestOwnedDuplex100MiB(t *testing.T) { + for _, scenario := range []string{duplexCommand, slowHost, pausedInput} { + t.Run(scenario, func(t *testing.T) { duplex(t, scenario) }) + } +} + +func duplex(t *testing.T, scenario string) { + t.Helper() + _, driver := launch(t) + ctx, cancel := context.WithCancel(deadline(t)) + defer cancel() + command := scenario + if scenario == slowHost { + command = duplexCommand + } + s := start(ctx, t, driver, command) + if scenario == pausedInput { + ready(t, s) + } + done := sender(s) + defer func() { cancel(); _ = joined(t, done) }() + if scenario == pausedInput { + release(t, driver) + } + if err := receive(s, scenario == slowHost); err != nil { + t.Fatal(err) + } + if err := joined(t, done); err != nil { + t.Fatal(err) + } +} + +func ready(t *testing.T, s runtimev1.RuntimeDriver_ExecClient) { + t.Helper() + frame, err := s.Recv() + if err != nil || string(frame.GetStdout().GetData()) != "ready" { + t.Fatalf("readiness: %v", err) + } +} + +func release(t *testing.T, driver runtimev1.RuntimeDriverClient) { + t.Helper() + s := start(deadline(t), t, driver, "release") + if err := s.CloseSend(); err != nil { + t.Fatal(err) + } + frame, err := s.Recv() + if err != nil || frame.GetExit() == nil { + t.Fatalf("release: %v", err) + } + if _, err := s.Recv(); err != io.EOF { + t.Fatalf("release completion: %v", err) + } +} + +func TestOwnedBackpressureCancellationBoundsRetainedHeap(t *testing.T) { + for _, command := range []string{pausedInput, duplexCommand} { + t.Run(command, func(t *testing.T) { backpressure(t, command) }) + } +} + +func backpressure(t *testing.T, command string) { + t.Helper() + _, driver := launch(t) + pluginBefore, hostBefore := memory(t, driver), hostMemory() + ctx, cancel := context.WithCancel(deadline(t)) + defer cancel() + s := start(ctx, t, driver, command) + prime(t, s, command) + done := sender(s) + defer func() { cancel(); _ = joined(t, done) }() + awaitSender(t, done) + pluginDuring, hostDuring := memory(t, driver), hostMemory() + assertHeapBudget(t, [2]uint64{hostBefore, hostDuring}, [2]uint64{pluginBefore, pluginDuring}) + select { + case <-done.done: + t.Fatalf("sender bypassed paused consumer: %v", done.err) + default: + } + cancel() + canceled(t, s) + if err := joined(t, done); err == nil { + t.Fatal("blocked stdin completed after cancellation") + } + _ = memory(t, driver) +} + +func prime(t *testing.T, s runtimev1.RuntimeDriver_ExecClient, command string) { + t.Helper() + if command == pausedInput { + ready(t, s) + } + if err := s.Send( + &runtimev1.ExecClientMessage{Payload: &runtimev1.ExecClientMessage_Stdin{Stdin: pattern()}}, + ); err != nil { + t.Fatal(err) + } + if command != pausedInput { + frame, err := s.Recv() + if err != nil || !bytes.Equal(frame.GetStdout().GetData(), pattern()) { + t.Fatalf("initial output: %v", err) + } + } +} + +func assertHeapBudget(t *testing.T, host, plugin [2]uint64) { + t.Helper() + const budget = 32 << 20 + t.Logf( + "retained heap baseline/paused: host=%d/%d plugin=%d/%d budget=%d", + host[0], + host[1], + plugin[0], + plugin[1], + budget, + ) + if plugin[1] > plugin[0]+budget || host[1] > host[0]+budget { + t.Fatal("paused stream retained more than the 32 MiB heap-growth budget") + } +} + +func TestOwnedStreamEarlyExitAndCrashJoinSender(t *testing.T) { + for _, command := range []string{"early-exit", "crash"} { + t.Run(command, func(t *testing.T) { interrupted(t, command) }) + } +} + +func interrupted(t *testing.T, command string) { + t.Helper() + client, driver := launch(t) + ctx, cancel := context.WithCancel(deadline(t)) + defer cancel() + s := start(ctx, t, driver, command) + ready(t, s) + done := sender(s) + defer func() { cancel(); _ = joined(t, done) }() + assertInterrupted(t, s, command) + if err := joined(t, done); err == nil { + t.Fatal("stdin completed despite command interruption") + } + client.Kill() + if !client.Exited() { + t.Fatal("supervisor was not reaped") + } +} + +func canceled(t *testing.T, s runtimev1.RuntimeDriver_ExecClient) { + t.Helper() + for { + _, err := s.Recv() + if err == nil { + continue + } + if status.Code(err) != codes.Canceled { + t.Fatalf("cancel: %v", err) + } + return + } +} + +func assertInterrupted(t *testing.T, s runtimev1.RuntimeDriver_ExecClient, command string) { + t.Helper() + frame, err := s.Recv() + if command == "crash" { + if status.Code(err) != codes.Unavailable { + t.Fatalf("crash: %v", err) + } + } else { + if err != nil || frame.GetExit() == nil || frame.GetExit().GetExitCode() != 0 { + t.Fatalf("early exit: %v (%v)", frame, err) + } + if _, err := s.Recv(); err != io.EOF { + t.Fatalf("early completion: %v", err) + } + } +} diff --git a/internal/streamprobe/traffic_test.go b/internal/streamprobe/traffic_test.go new file mode 100644 index 0000000..17a917a --- /dev/null +++ b/internal/streamprobe/traffic_test.go @@ -0,0 +1,210 @@ +package streamprobe_test + +import ( + "bytes" + "crypto/sha256" + "encoding/json" + "errors" + "fmt" + "hash" + "io" + "runtime" + "testing" + "time" + + "github.com/devsy-org/devsy-runtime-sdk/runtimev1" +) + +const trafficBytes = 100 << 20 + +func pattern() []byte { + data := make([]byte, runtimev1.ChunkSize) + for i := range data { + data[i] = byte(i*31 + i/256) + } + return data +} + +func send(s runtimev1.RuntimeDriver_ExecClient, started chan struct{}) error { + data := pattern() + for remaining := trafficBytes; remaining > 0; remaining -= len(data) { + if err := s.Send( + &runtimev1.ExecClientMessage{Payload: &runtimev1.ExecClientMessage_Stdin{Stdin: data}}, + ); err != nil { + return err + } + if remaining == trafficBytes { + close(started) + } + } + if err := s.Send(&runtimev1.ExecClientMessage{Payload: &runtimev1.ExecClientMessage_CloseStdin{ + CloseStdin: &runtimev1.CloseStdin{}, + }}); err != nil { + return err + } + return s.CloseSend() +} + +type upload struct { + done chan struct{} + started chan struct{} + err error +} + +func sender(s runtimev1.RuntimeDriver_ExecClient) *upload { + pending := &upload{done: make(chan struct{}), started: make(chan struct{})} + go func() { pending.err = send(s, pending.started); close(pending.done) }() + return pending +} + +func joined(t *testing.T, pending *upload) error { + t.Helper() + select { + case <-pending.done: + return pending.err + case <-time.After(10 * time.Second): + t.Fatal("stream sender did not stop") + return nil + } +} + +type output struct { + stdout, stderr hash.Hash + outBytes, errBytes int + exited bool + slow bool +} + +func receive(s runtimev1.RuntimeDriver_ExecClient, slow bool) error { + result := output{stdout: sha256.New(), stderr: sha256.New(), slow: slow} + for { + message, err := s.Recv() + if err == io.EOF { + return result.verify() + } + if err != nil { + return err + } + if err := result.accept(message); err != nil { + return err + } + } +} + +func (o *output) accept(message *runtimev1.ExecServerMessage) error { + if o.exited { + return errors.New("frame after terminal exit") + } + return o.channel(message) +} + +func (o *output) channel(message *runtimev1.ExecServerMessage) error { + switch payload := message.Payload.(type) { + case *runtimev1.ExecServerMessage_Exit: + return o.terminal(payload.Exit) + case *runtimev1.ExecServerMessage_Stdout: + return o.writeStdout(payload.Stdout.GetData()) + case *runtimev1.ExecServerMessage_Stderr: + if err := digest(o.stderr, payload.Stderr.GetData()); err != nil { + return err + } + o.errBytes += len(payload.Stderr.GetData()) + default: + return errors.New("unexpected output") + } + return nil +} + +func (o *output) terminal(exit *runtimev1.ExecExit) error { + if exit.GetExitCode() != 7 || exit.GetSignal() != "" { + return errors.New("wrong exit") + } + o.exited = true + return nil +} + +func (o *output) writeStdout(data []byte) error { + if err := digest(o.stdout, data); err != nil { + return err + } + o.outBytes += len(data) + if o.slow && o.outBytes%(1<<20) == 0 { + time.Sleep(time.Millisecond) + } + return nil +} + +func digest(h hash.Hash, data []byte) error { + if len(data) == 0 || len(data) > runtimev1.ChunkSize { + return errors.New("unbounded output frame") + } + _, err := h.Write(data) + return err +} + +func (o *output) verify() error { + if !o.exited { + return errors.New("missing terminal exit") + } + for _, channel := range []struct { + got hash.Hash + size int + tail string + }{ + {o.stdout, o.outBytes, "stdout tail\x00\xff"}, {o.stderr, o.errBytes, "stderr tail\xff\x00"}, + } { + want := sha256.New() + data := pattern() + for remaining := trafficBytes; remaining > 0; remaining -= len(data) { + _, _ = want.Write(data) + } + _, _ = want.Write([]byte(channel.tail)) + if channel.size != trafficBytes+len(channel.tail) || + !bytes.Equal(channel.got.Sum(nil), want.Sum(nil)) { + return fmt.Errorf("binary output or tail lost: %d bytes", channel.size) + } + } + return nil +} + +func memory(t *testing.T, driver runtimev1.RuntimeDriverClient) uint64 { + t.Helper() + s := start(deadline(t), t, driver, "memory") + if err := s.CloseSend(); err != nil { + t.Fatal(err) + } + frame, err := s.Recv() + if err != nil { + t.Fatal(err) + } + var heap uint64 + if err := json.Unmarshal(frame.GetStdout().GetData(), &heap); err != nil { + t.Fatal(err) + } + terminal, err := s.Recv() + if err != nil || terminal.GetExit() == nil { + t.Fatalf("missing memory exit: %v", err) + } + if _, err := s.Recv(); err != io.EOF { + t.Fatalf("memory stream: %v", err) + } + return heap +} + +func hostMemory() uint64 { + runtime.GC() + var stats runtime.MemStats + runtime.ReadMemStats(&stats) + return stats.HeapAlloc +} + +func awaitSender(t *testing.T, pending *upload) { + t.Helper() + select { + case <-pending.started: + case <-pending.done: + t.Fatalf("sender stopped before first data frame: %v", pending.err) + case <-time.After(10 * time.Second): + t.Fatal("sender did not start") + } +}