Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 32 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
153 changes: 153 additions & 0 deletions internal/streamfixture/main.go
Original file line number Diff line number Diff line change
@@ -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)
}
113 changes: 113 additions & 0 deletions internal/streamprobe/fixture_test.go
Original file line number Diff line number Diff line change
@@ -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
}
Loading
Loading