From 52a2ba2237b61ccf1aaf21877f6d66008ab4d37c Mon Sep 17 00:00:00 2001 From: Ilyas Salikhov Date: Tue, 22 Sep 2026 13:50:40 +0300 Subject: [PATCH 1/3] fix(go): bound command streaming output retention --- README.md | 3 + docs/go-command-streaming.md | 128 +++++ packages/go-sdk/GO_PARITY.md | 3 +- packages/go-sdk/README.md | 83 ++- packages/go-sdk/command_streaming_test.go | 531 ++++++++++++++++++ packages/go-sdk/commands.go | 255 +++++++-- packages/go-sdk/envd_test.go | 2 +- packages/go-sdk/examples_test.go | 43 ++ .../integration/command_streaming_test.go | 88 +++ packages/go-sdk/pty.go | 21 +- reference/manifest.json | 2 +- reference/sdk/go/core.md | 113 +++- 12 files changed, 1216 insertions(+), 56 deletions(-) create mode 100644 docs/go-command-streaming.md create mode 100644 packages/go-sdk/command_streaming_test.go create mode 100644 packages/go-sdk/integration/command_streaming_test.go diff --git a/README.md b/README.md index 45813849..27578946 100644 --- a/README.md +++ b/README.md @@ -93,4 +93,7 @@ The CLI stores local configuration in `~/.agentbox/config.json`; environment var Contributor setup, spec synchronization, verification, KVM testing, versioning, and publication are documented in [RELEASING.md](RELEASING.md). +The [Go command streaming design](docs/go-command-streaming.md) documents the +opt-in output policy for long-running processes. + AgentBox SDK is derived from upstream work described in [UPSTREAM.md](UPSTREAM.md). Licensing notices are in [LICENSE](LICENSE) and [NOTICE](NOTICE). diff --git a/docs/go-command-streaming.md b/docs/go-command-streaming.md new file mode 100644 index 00000000..cbf8c2e5 --- /dev/null +++ b/docs/go-command-streaming.md @@ -0,0 +1,128 @@ +# Go command streaming without output retention + +## Contract and owner + +The handwritten Go command transport in +[`commands.go`](../packages/go-sdk/commands.go) supports explicit streaming on +start and reattachment. [`pty.go`](../packages/go-sdk/pty.go) shares the delivery +implementation. No backend, wire or generated binding changes are required. +The minimum Go version remains 1.24.0. + +A non-nil `CommandOptions.Streaming`, `CommandConnectOptions.Streaming` or +`PTYOptions.Streaming` selects streaming. Use `Commands.ConnectWithOptions` or +`PTY.ConnectWithOptions` for attachment by PID or tag. Existing `Connect` calls +and nil options retain collecting behavior. + +| Concern | Collecting mode | Streaming mode | +| --- | --- | --- | +| Result | Full stdout/stderr and exit metadata | PID, status, exit code and errors; nil stdout/stderr | +| Callbacks | Existing callback and channel delivery | Synchronous callback only for that output | +| Channels | Unbounded queues with a preserved natural-completion tail | Explicitly enabled, unbuffered, at most 32 KiB per delivered slice | +| Unselected output | Captured and queued | Discarded; its public channel is closed | +| Callback plus channel | Existing fan-out | Invalid for the same output | +| `Wait` without reads | Supported | Supported only when no channels are enabled; callbacks must return | +| Local detach | `Close` or attachment context cancellation | Same; no signal or stdin EOF | + +`CommandStreamingOptions` has `StdoutChannel`, `StderrChannel` and `PTYChannel` +flags. Existing start callbacks remain in `CommandOptions`/`PTYOptions`; +`CommandConnectOptions` supplies `OnStdout`, `OnStderr` and `OnPTY`. +`Run` continues to drain channels itself and collects output by default. + +## Memory and delivery bounds + +Streaming creates no queue goroutines or output history. An enabled channel uses +synchronous delivery, splitting events into independent slices of at most 32 KiB. +The SDK retains at most one pending delivery slice plus the current decoded +transport event. A streaming Connect client imposes a 4 MiB message-size limit; +this also bounds messages selected for discard. Decode buffers and HTTP transport +buffers add overhead. This is a bound independent of total emitted bytes, not a +promise of a 4 MiB total heap or RSS. Memory retained by the consumer is outside +the SDK bound. + +A reader that stops blocks subsequent event processing, including exit events and +other outputs. Backpressure can propagate to the remote process. Read every +enabled channel concurrently, or select callbacks/discard for unused outputs. +There is no silent overflow, truncation, or goroutine per chunk. An oversized +transport message fails and closes the local attachment. + +Callbacks execute serially on the receiver. Their byte slices are read-only and +valid only until the callback returns; copy bytes needed later. Channel receivers +own their independent slices. Channel chunk boundaries may differ from transport +boundaries, but byte order and content are preserved. Slow callbacks also apply +backpressure. Callbacks must not call `Wait` on their own handle. + +## Completion, cancellation and errors + +Natural streaming completion hands every enabled-channel byte to a reader before +closing the channels and `Done`; application processing may finish later. +Collecting-mode completion preserves its queued channel tail, even after `Wait`. + +`CommandHandle.Close` cancels the local attachment and releases queued output. +Canceling the context supplied to start/connect does the same while the receiver +is active. Both close the HTTP response without relying on GC and without killing +the process or closing stdin. `Close` is idempotent and does not wait for arbitrary +callback code. `Done`/`Wait` complete when that callback returns and the receiver +exits. Canceling only `Wait`'s context does not detach. + +Local cancellation reports `context.Canceled` or `context.DeadlineExceeded`. +A confirmed process end retains its result, including `CommandExitError` for a +nonzero exit. In streaming mode that error contains no stdout/stderr or server +error message. Transport failures preserve Connect error classification but use +generic diagnostic text, avoiding output-bearing server error details. + +Reattachment adds no stdout replay, retry, restart or exactly-once promise. +Consumers must select the streaming policy on both launch and reconnect. + +## Verification + +Hermetic regressions in +[`command_streaming_test.go`](../packages/go-sdk/command_streaming_test.go) cover: + +- Multiple MiB of ordered stdout/stderr through callbacks and channels, including + inspection of a still-live handle at a synchronized callback barrier. +- No capture or hidden channel queue; unused stderr discard and closed channels. +- Channel blocks larger than 32 KiB, byte ownership and ordered splitting. +- Unread output preventing completion, context cancellation, repeated detach, + response-body closure, and no signal/EOF side effect. +- Start, attach by PID/tag, nonzero exit, transport failure and oversized events. +- Cancellation while a callback is blocked, collecting tail release on `Close`, + PTY callback/channel delivery and PTY detach. + +Existing collecting-mode no-drain, queued-tail, command/PTY and response-closure +regressions remain required. The supported Go matrix is 1.24–1.27; run +`make go-check` and generation/artifact gates described in +[`RELEASING.md`](../RELEASING.md) using the pinned containers. + +Validation passed: `make go-check` (91.2% handwritten statement coverage), builds +and hermetic tests on Go 1.24.13, 1.25.14, 1.26.8 and 1.27.1, race checks on +1.24.13 and 1.27.1, format/lint/type checks, workspace tests, release-artifact +builds and clean installation checks. `make generate` and reference-contract +checks passed; a second generation produced identical tracked files. This does +not replace the full multi-language KVM release suite required before publication. + +`BenchmarkCommandOutputRetention` compares callback consumers at fixed 1 KiB +chunks, retaining the handle until observation and leaving channels unread. +A Go 1.24.13 linux/arm64 sample (`-benchtime=1x -benchmem`) measured: + +| Emitted stdout | Collecting retained output storage | Streaming retained output storage | Streaming allocated bytes | +| --- | --- | --- | --- | +| 1 MiB | 2,161,664 bytes | 0 bytes | 1,904 bytes | +| 16 MiB | 37,428,224 bytes | 0 bytes | 1,904 bytes | + +Retained storage counts capture capacity and queued slices, excluding queue +metadata and any collecting slice already in flight. The benchmark isolates SDK +delivery from transport decoding. Allocation counts are supporting evidence; +structural live-retention assertions are the deterministic gate. + +[`TestCommandStreamingAttachDetachKVM`](../packages/go-sdk/integration/command_streaming_test.go) +starts `cat`, exchanges bytes, detaches/reconnects three times using PID and tag, +exchanges further bytes, then explicitly closes stdin and checks exit metadata. +It creates and deletes only its own sandbox. This runtime smoke has passed. + +## Delivery + +The implementation must follow the coordinated SDK release procedure in +[`RELEASING.md`](../RELEASING.md). No ad-hoc fork or replacement module is needed. +Publication and the consumer dependency upgrade remain release steps; no released +module version is assigned by this source change. The consuming orchestrator must +use the released version and explicitly enable streaming on launch and reconnect. diff --git a/packages/go-sdk/GO_PARITY.md b/packages/go-sdk/GO_PARITY.md index ab621e1f..d176a053 100644 --- a/packages/go-sdk/GO_PARITY.md +++ b/packages/go-sdk/GO_PARITY.md @@ -12,7 +12,8 @@ concurrency, and race coverage. | Network rules, IAM payloads, metrics, structured logs | sandbox network/IAM/metrics/log tests | `sandbox_test.go` | `integration/sdk_test.go` | | Forks, snapshots, signed upload/download URLs | sandbox fork/snapshot/signature tests | `sandbox_test.go` | core KVM lifecycle | | Foreground/background commands, attach/list, stdin/EOF, signals, output streams, exit errors | command and command-handle tests | `envd_test.go` | `integration/sdk_test.go` | -| PTY create/attach/input/resize/kill | PTY tests | `envd_test.go` | KVM command transport | +| Opt-in bounded command delivery, callback/discard/channel policy, detach without EOF | Go-specific explicit memory policy; default behavior remains aligned | `command_streaming_test.go` | `integration/command_streaming_test.go` | +| PTY create/attach/input/resize/kill | PTY tests | `envd_test.go`, `command_streaming_test.go` | KVM command transport | | Text/binary/stream reads and writes, batch writes, list/stat/metadata/exists/mkdir/move/remove/watch | filesystem and watch-handle tests | `envd_test.go` | `integration/sdk_test.go` | | Base images/templates, private registries, Dockerfile parsing, copy, packages, env/user/workdir/start/ready/cache | template builder/parser tests | `template_test.go` | `integration/sdk_test.go` | | Build request/upload/start/poll/log/status, visibility, tags, list/info/delete | template API/build tests | `template_test.go` | `integration/sdk_test.go` | diff --git a/packages/go-sdk/README.md b/packages/go-sdk/README.md index e2b14e78..4f22c77f 100644 --- a/packages/go-sdk/README.md +++ b/packages/go-sdk/README.md @@ -52,11 +52,84 @@ Documentation: [core SDK](https://docs.agentbox.ru/en/sdk/), [templates](https://docs.agentbox.ru/en/sdk/templates/), and [Code Interpreter](https://docs.agentbox.ru/en/sdk/code-interpreter/). -`Sandbox.Kill` returns `false, nil` when the sandbox no longer exists. Command -handles can be waited on without draining their live output channels; `Wait` -always returns complete stdout and stderr collected by the SDK. Output channels -may finish draining and close after `Wait` returns. PTY callers can consume -`CommandHandle.PTY` or use `PTYOptions.OnPTY`. +`Sandbox.Kill` returns `false, nil` when the sandbox no longer exists. + +## Command output + +By default, command handles collect complete stdout/stderr for `Wait`, which can +be called without draining their output channels. Channel tails may finish +draining after `Wait` returns. Callbacks also receive output in this mode. + +For long-lived processes, opt into streaming on both start and reconnect: + +```go +policy := &agentbox.CommandStreamingOptions{} +handle, err := sandbox.Commands.Start(ctx, "cat", &agentbox.CommandOptions{ + Stdin: true, + Streaming: policy, + OnStdout: func(chunk []byte) { consume(chunk) }, + // No stderr callback or channel: discard stderr without retaining it. +}) +if err != nil { + return err +} +defer handle.Close() +pid, err := handle.PID(ctx) +if err != nil { + return err +} +// Exchange messages, then detach locally. The remote process and stdin stay open. +_ = handle.Close() +_, _ = handle.Wait(ctx) // reports context.Canceled for an active attachment +attached, err := sandbox.Commands.ConnectWithOptions(ctx, pid, "", &agentbox.CommandConnectOptions{ + Streaming: policy, + OnStdout: func(chunk []byte) { consume(chunk) }, +}) +if err != nil { + return err +} +defer attached.Close() +``` + +A non-nil `Streaming` selects exactly one delivery path per output: its callback, +its explicitly enabled channel (`StdoutChannel`, `StderrChannel`, `PTYChannel`), +or discard. Selecting both a callback and a channel for the same output returns +`InvalidArgumentError`. Unselected channels are already closed. Callbacks execute +serially on the receiver; their read-only slices are valid until the callback +returns. Copy data that must be retained and return promptly. Cancellation closes +the transport but cannot interrupt callback code; `Done`/`Wait` finish after it +returns. Do not call `Wait` from a callback. + +Streaming channels are unbuffered. Consume every enabled channel concurrently +with `Wait`; a stalled reader applies backpressure and can eventually stall the +remote process. Byte order is preserved, but transport chunks may be split into +slices of at most **32 KiB**. The receiver owns those slices. There is no output +queue or captured result, only one pending channel slice plus one in-flight +transport event. Each encoded/decompressed Connect message is limited to **4 MiB**; +an oversized message fails the attachment with a resource-exhausted error. HTTP +buffers, decoding allocations and caller-retained slices are additional memory. +The limit is independent of total process output, not an exact heap/RSS cap. + +`Wait` returns PID, status, exit code and errors, with nil stdout/stderr in streaming +mode. Successful completion means all enabled channel bytes were handed to their +readers; their processing may still be running. Nonzero exits return +`CommandExitError` with an empty output result and empty server error message. +Transport failures retain their error code but omit server-provided diagnostic +text/details that might contain output. Neither path silently truncates output +and reports success. + +Cancel the context passed to `Start`/`ConnectWithOptions`, or call `Close`, to +release the local attachment without killing the process or sending EOF. +Cancellation returns `context.Canceled`/`context.DeadlineExceeded`, unless a process +end event was already confirmed. `Close` also discards an unread collecting-mode +channel tail; natural completion preserves that tail. Canceling only the context +passed to `Wait` stops waiting without detaching. Reconnecting by PID or tag adds +no stdout replay, delivery retry, process restart or exactly-once guarantee. + +PTY callers can use `PTYOptions.Streaming` with `OnPTY` or `PTYChannel`, and +`PTY.ConnectWithOptions` to reattach with the same policy. Default PTY behavior +remains unchanged. See the [streaming contract](../../docs/go-command-streaming.md) +for bounds and validation, and `examples_test.go` for a complete example. Streaming uploads have no SDK deadline by default. Set `WriteFileOptions.RequestTimeout` to limit a complete upload. A client supplied diff --git a/packages/go-sdk/command_streaming_test.go b/packages/go-sdk/command_streaming_test.go new file mode 100644 index 00000000..e869bc56 --- /dev/null +++ b/packages/go-sdk/command_streaming_test.go @@ -0,0 +1,531 @@ +package agentbox + +import ( + "bytes" + "context" + "errors" + "fmt" + "net/http" + "net/http/httptest" + "runtime" + "strings" + "sync/atomic" + "testing" + "time" + + "connectrpc.com/connect" + api "github.com/abox-dev/sdk/packages/go-sdk/internal/gen/api" + process "github.com/abox-dev/sdk/packages/go-sdk/internal/gen/envd/process" + "github.com/abox-dev/sdk/packages/go-sdk/internal/gen/envd/process/processconnect" +) + +type streamingProcessServer struct { + testProcessServer + chunks int + size int + failure bool + pty bool + exitCode int32 + signals atomic.Int64 + eof atomic.Int64 +} + +func (server *streamingProcessServer) events(send func(*process.ProcessEvent) error) error { + if err := send(&process.ProcessEvent{Event: &process.ProcessEvent_Start{Start: &process.ProcessEvent_StartEvent{Pid: 42}}}); err != nil { + return err + } + for i := range server.chunks { + chunk := bytes.Repeat([]byte{byte(i % 251)}, server.size) + outputs := []*process.ProcessEvent_DataEvent{ + {Output: &process.ProcessEvent_DataEvent_Stdout{Stdout: chunk}}, + {Output: &process.ProcessEvent_DataEvent_Stderr{Stderr: chunk}}, + } + if server.pty { + outputs = []*process.ProcessEvent_DataEvent{{Output: &process.ProcessEvent_DataEvent_Pty{Pty: chunk}}} + } + for _, data := range outputs { + if err := send(&process.ProcessEvent{Event: &process.ProcessEvent_Data{Data: data}}); err != nil { + return err + } + } + } + if server.failure { + return connect.NewError(connect.CodeInternal, errors.New("secret output")) + } + message := "secret output" + return send(&process.ProcessEvent{Event: &process.ProcessEvent_End{End: &process.ProcessEvent_EndEvent{Exited: true, ExitCode: server.exitCode, Status: "exited", Error: &message}}}) +} +func (server *streamingProcessServer) Start(_ context.Context, _ *connect.Request[process.StartRequest], stream *connect.ServerStream[process.StartResponse]) error { + return server.events(func(event *process.ProcessEvent) error { return stream.Send(&process.StartResponse{Event: event}) }) +} +func (server *streamingProcessServer) Connect(_ context.Context, _ *connect.Request[process.ConnectRequest], stream *connect.ServerStream[process.ConnectResponse]) error { + return server.events(func(event *process.ProcessEvent) error { return stream.Send(&process.ConnectResponse{Event: event}) }) +} +func (server *streamingProcessServer) SendSignal(context.Context, *connect.Request[process.SendSignalRequest]) (*connect.Response[process.SendSignalResponse], error) { + server.signals.Add(1) + return connect.NewResponse(&process.SendSignalResponse{}), nil +} +func (server *streamingProcessServer) CloseStdin(context.Context, *connect.Request[process.CloseStdinRequest]) (*connect.Response[process.CloseStdinResponse], error) { + server.eof.Add(1) + return connect.NewResponse(&process.CloseStdinResponse{}), nil +} +func newStreamingTestSandbox(t *testing.T, fixture *streamingProcessServer) (*Sandbox, *closeTrackingTransport) { + t.Helper() + path, handler := processconnect.NewProcessHandler(fixture, connect.WithCodec(tolerantJSONCodec{})) + mux := http.NewServeMux() + mux.Handle(path, handler) + server := httptest.NewServer(mux) + t.Cleanup(server.Close) + transport := &closeTrackingTransport{base: http.DefaultTransport.(*http.Transport).Clone()} + t.Cleanup(func() { transport.base.(*http.Transport).CloseIdleConnections() }) + client, err := NewClient(WithAPIURL(server.URL), WithSandboxURL(server.URL), WithHTTPClient(&http.Client{Transport: transport})) + if err != nil { + t.Fatal(err) + } + return client.sandboxFromAPI(api.Sandbox{SandboxID: "sbx", TemplateID: "base"}), transport +} +func waitStreamingTest(t *testing.T, handle *CommandHandle) (CommandResult, error) { + t.Helper() + ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + defer cancel() + return handle.Wait(ctx) +} +func assertNoRetainedOutput(t *testing.T, handle *CommandHandle) { + t.Helper() + <-handle.Done + if handle.result.Stdout != nil || handle.result.Stderr != nil { + t.Fatal("streaming mode captured output") + } + for _, stream := range []*outputStream{handle.stdout, handle.stderr, handle.pty} { + stream.mu.Lock() + retained := len(stream.queue) + stream.mu.Unlock() + if retained != 0 || !stream.direct || cap(stream.output) != 0 { + t.Fatalf("hidden output queue: %d", retained) + } + } +} + +func TestCommandStreamingCallbacks(t *testing.T) { + fixture := &streamingProcessServer{chunks: 2048, size: 1024} + sandbox, transport := newStreamingTestSandbox(t, fixture) + for _, attach := range []string{"start", "pid", "tag"} { + t.Run(attach, func(t *testing.T) { + var counts [2]int + callback := func(index int) func([]byte) { + return func(chunk []byte) { + want := bytes.Repeat([]byte{byte(counts[index] % 251)}, fixture.size) + if !bytes.Equal(chunk, want) { + t.Errorf("unexpected callback chunk %d", counts[index]) + } + counts[index]++ + } + } + policy := &CommandStreamingOptions{} + var handle *CommandHandle + var err error + switch attach { + case "start": + handle, err = sandbox.Commands.Start(t.Context(), "flood", &CommandOptions{Streaming: policy, OnStdout: callback(0), OnStderr: callback(1)}) + case "pid": + handle, err = sandbox.Commands.ConnectWithOptions(t.Context(), 42, "", &CommandConnectOptions{Streaming: policy, OnStdout: callback(0), OnStderr: callback(1)}) + case "tag": + handle, err = sandbox.Commands.ConnectWithOptions(t.Context(), 0, "rpc", &CommandConnectOptions{Streaming: policy, OnStdout: callback(0), OnStderr: callback(1)}) + } + if err != nil { + t.Fatal(err) + } + defer handle.Close() + result, err := waitStreamingTest(t, handle) + if err != nil || result.PID != 42 || result.Status != "exited" { + t.Fatalf("result: %+v %v", result, err) + } + if counts != [2]int{fixture.chunks, fixture.chunks} { + t.Fatalf("callbacks: %v", counts) + } + assertNoRetainedOutput(t, handle) + for _, channel := range []<-chan []byte{handle.Stdout, handle.Stderr, handle.PTY} { + if _, ok := <-channel; ok { + t.Fatal("unused channel is open") + } + } + }) + } + if transport.closed.Load() != 3 { + t.Fatal("response bodies not released") + } +} + +func TestCommandStreamingChannels(t *testing.T) { + fixture := &streamingProcessServer{chunks: 32, size: 100001} + sandbox, _ := newStreamingTestSandbox(t, fixture) + for _, stderr := range []bool{false, true} { + t.Run(fmt.Sprintf("stderr=%t", stderr), func(t *testing.T) { + handle, err := sandbox.Commands.Start(t.Context(), "flood", &CommandOptions{Streaming: &CommandStreamingOptions{StdoutChannel: true, StderrChannel: stderr}}) + if err != nil { + t.Fatal(err) + } + defer handle.Close() + // One receiver drains both enabled channels without assuming event boundaries. + counts := [2]int{} + stdout, stderrChannel := handle.Stdout, handle.Stderr + for stdout != nil || stderrChannel != nil { + var chunk []byte + var ok bool + index := 0 + select { + case chunk, ok = <-stdout: + if !ok { + stdout = nil + continue + } + case chunk, ok = <-stderrChannel: + index = 1 + if !ok { + stderrChannel = nil + continue + } + case <-time.After(5 * time.Second): + t.Fatal("channel stalled") + } + if len(chunk) > commandStreamChunkBytes { + t.Fatal("unbounded channel slice") + } + for _, value := range chunk { + if value != byte((counts[index]/fixture.size)%251) { + t.Fatal("bytes lost or reordered") + } + counts[index]++ + } + // Retaining and mutating a delivered slice must not affect subsequent data. + clear(chunk) + } + result, err := waitStreamingTest(t, handle) + if err != nil || result.PID != 42 { + t.Fatalf("wait: %+v %v", result, err) + } + want := [2]int{fixture.chunks * fixture.size, 0} + if stderr { + want[1] = want[0] + } + if counts != want { + t.Fatalf("bytes: %v want %v", counts, want) + } + assertNoRetainedOutput(t, handle) + }) + } +} + +func TestCommandStreamingBlockedDeliveryCancellation(t *testing.T) { + fixture := &streamingProcessServer{chunks: 100, size: 1024} + sandbox, transport := newStreamingTestSandbox(t, fixture) + for i := range 30 { + ctx, cancel := context.WithCancel(t.Context()) + handle, err := sandbox.Commands.ConnectWithOptions(ctx, 42, "", &CommandConnectOptions{Streaming: &CommandStreamingOptions{StdoutChannel: true}}) + if err != nil { + cancel() + t.Fatal(err) + } + // Receiving one chunk ensures delivery has started; the next unread chunk + // must prevent the end event from being observed even if the server has exited. + <-handle.Stdout + select { + case <-handle.Done: + t.Fatal("unread channel silently discarded output") + default: + } + if i%2 == 0 { + cancel() + } else { + handle.Close() + } + result, err := waitStreamingTest(t, handle) + cancel() + if !errors.Is(err, context.Canceled) || result.PID != 42 { + t.Fatalf("cancel: %+v %v", result, err) + } + assertNoRetainedOutput(t, handle) + if _, ok := <-handle.Stdout; ok { + t.Fatal("channel retained data after cancellation") + } + handle.Close() + } + if transport.closed.Load() != 30 || fixture.signals.Load() != 0 || fixture.eof.Load() != 0 { + t.Fatal("detach leaked a response or signaled the process") + } +} + +func TestCommandStreamingFailures(t *testing.T) { + for _, kind := range []string{"exit", "transport", "oversized"} { + t.Run(kind, func(t *testing.T) { + fixture := &streamingProcessServer{chunks: 1, size: 1024} + switch kind { + case "exit": + fixture.exitCode = 9 + case "transport": + fixture.failure = true + case "oversized": + fixture.size = commandStreamMessageBytes + } + sandbox, transport := newStreamingTestSandbox(t, fixture) + handle, err := sandbox.Commands.Start(t.Context(), "fail", &CommandOptions{Streaming: &CommandStreamingOptions{}}) + if err != nil { + t.Fatal(err) + } + defer handle.Close() + result, err := waitStreamingTest(t, handle) + if err == nil || result.PID != 42 || strings.Contains(fmt.Sprintf("%+v", err), "secret") { + t.Fatalf("failure: %+v %v", result, err) + } + if kind == "exit" { + var exit *CommandExitError + if !errors.As(err, &exit) || exit.Result.ExitCode != 9 || exit.Result.Stdout != nil || exit.Result.Stderr != nil || exit.Message != "" { + t.Fatalf("exit: %#v", err) + } + } else { + want := connect.CodeInternal + if kind == "oversized" { + want = connect.CodeResourceExhausted + } + if connect.CodeOf(err) != want { + t.Fatalf("code: %s want %s", connect.CodeOf(err), want) + } + + } + assertNoRetainedOutput(t, handle) + if transport.closed.Load() != 1 { + t.Fatal("unclosed response") + } + }) + } +} + +func TestCommandStreamingValidationAndPTY(t *testing.T) { + sandbox, closeServer := newEnvdTestSandbox(t) + defer closeServer() + callback := func([]byte) {} + for _, options := range []*CommandConnectOptions{ + {Streaming: &CommandStreamingOptions{StdoutChannel: true}, OnStdout: callback}, + {Streaming: &CommandStreamingOptions{StderrChannel: true}, OnStderr: callback}, + {Streaming: &CommandStreamingOptions{PTYChannel: true}, OnPTY: callback}, + } { + if _, err := sandbox.Commands.ConnectWithOptions(t.Context(), 7, "", options); err == nil { + t.Fatal("accepted fan-out") + } + } + if _, err := sandbox.Commands.Start(t.Context(), "echo", &CommandOptions{Streaming: &CommandStreamingOptions{StdoutChannel: true}, OnStdout: callback}); err == nil { + t.Fatal("accepted start fan-out") + } + if _, err := sandbox.PTY.Create(t.Context(), "sh", &PTYOptions{Streaming: &CommandStreamingOptions{PTYChannel: true}, OnPTY: callback}); err == nil { + t.Fatal("accepted PTY fan-out") + } + for _, useCallback := range []bool{false, true} { + var output bytes.Buffer + options := &PTYOptions{Streaming: &CommandStreamingOptions{PTYChannel: !useCallback}} + if useCallback { + options.OnPTY = func(chunk []byte) { output.Write(chunk) } + } + handle, err := sandbox.PTY.Create(t.Context(), "sh", options) + if err != nil { + t.Fatal(err) + } + for chunk := range handle.PTY { + output.Write(chunk) + } + if _, err := waitStreamingTest(t, handle); err != nil { + t.Fatal(err) + } + if output.String() != "pty\n" { + t.Fatalf("PTY: %q", output.String()) + } + assertNoRetainedOutput(t, handle) + } + result, err := sandbox.Commands.Run(t.Context(), "flood", &CommandOptions{Streaming: &CommandStreamingOptions{StdoutChannel: true}}) + if err != nil || result.Stdout != nil || result.PID != 7 { + t.Fatalf("streaming Run: %+v %v", result, err) + } + handle, err := sandbox.Commands.ConnectWithOptions(t.Context(), 7, "", nil) + if err != nil { + t.Fatal(err) + } + result, err = waitStreamingTest(t, handle) + if err != nil || result.PID != 7 || string(result.Stderr) != "connected" { + t.Fatalf("collecting Connect: %+v %v", result, err) + } + handle.Close() +} + +func TestCommandCancellationDuringCallback(t *testing.T) { + sandbox, transport := newStreamingTestSandbox(t, &streamingProcessServer{chunks: 2, size: 16}) + entered, release := make(chan struct{}), make(chan struct{}) + defer close(release) + handle, err := sandbox.Commands.Start(t.Context(), "callback", &CommandOptions{Streaming: &CommandStreamingOptions{}, OnStdout: func([]byte) { close(entered); <-release }}) + if err != nil { + t.Fatal(err) + } + <-entered + handle.Close() + deadline := time.Now().Add(5 * time.Second) + for transport.closed.Load() == 0 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + if transport.closed.Load() != 1 { + t.Fatal("callback prevented response closure") + } + select { + case <-handle.Done: + t.Fatal("callback was abandoned") + default: + } + release <- struct{}{} + if _, err := waitStreamingTest(t, handle); !errors.Is(err, context.Canceled) { + t.Fatalf("callback cancellation: %v", err) + } + assertNoRetainedOutput(t, handle) +} + +func BenchmarkCommandOutputRetention(b *testing.B) { + for _, streaming := range []bool{false, true} { + for _, chunks := range []int{1024, 16384} { + b.Run(fmt.Sprintf("streaming=%t/chunks=%d", streaming, chunks), func(b *testing.B) { + var policy *CommandStreamingOptions + if streaming { + policy = &CommandStreamingOptions{} + } + event := &process.ProcessEvent{Event: &process.ProcessEvent_Data{Data: &process.ProcessEvent_DataEvent{Output: &process.ProcessEvent_DataEvent_Stdout{Stdout: make([]byte, 1024)}}}} + b.ReportAllocs() + for b.Loop() { + handle := newCommandHandle(nil, "", policy) + n := 0 + handle.receive(b.Context(), func() (*process.ProcessEvent, bool) { n++; return event, n <= chunks }, func() error { return nil }, func() error { return nil }, outputCallbacks{stdout: func([]byte) {}}) + retained := cap(handle.result.Stdout) + handle.stdout.mu.Lock() + for _, chunk := range handle.stdout.queue { + retained += cap(chunk) + } + handle.stdout.mu.Unlock() + b.ReportMetric(float64(retained), "retained-B") + handle.Close() + runtime.KeepAlive(handle) + } + }) + } + } +} + +func TestCommandStreamingLiveRetention(t *testing.T) { + fixture := &streamingProcessServer{chunks: 4096, size: 1024} + sandbox, _ := newStreamingTestSandbox(t, fixture) + halfway, release := make(chan struct{}), make(chan struct{}) + count := 0 + handle, err := sandbox.Commands.Start(t.Context(), "long", &CommandOptions{Streaming: &CommandStreamingOptions{}, OnStdout: func([]byte) { + count++ + if count == 2048 { + close(halfway) + <-release + } + }}) + if err != nil { + t.Fatal(err) + } + defer handle.Close() + // The callback barrier keeps the handle live and stops all receiver writes. + <-halfway + if handle.result.Stdout != nil || handle.result.Stderr != nil { + t.Error("live result accumulated output") + } + for _, stream := range []*outputStream{handle.stdout, handle.stderr, handle.pty} { + stream.mu.Lock() + if len(stream.queue) != 0 { + t.Error("live callback accumulated queued output") + } + stream.mu.Unlock() + } + close(release) + if _, err := waitStreamingTest(t, handle); err != nil { + t.Fatal(err) + } + if count != fixture.chunks { + t.Fatalf("chunks: %d", count) + } +} + +func TestCommandCollectingCloseReleasesTail(t *testing.T) { + sandbox, closeServer := newEnvdTestSandbox(t) + defer closeServer() + handle, err := sandbox.Commands.Start(t.Context(), "flood", nil) + if err != nil { + t.Fatal(err) + } + if _, err := waitStreamingTest(t, handle); err != nil { + t.Fatal(err) + } + handle.Close() + for _, stream := range []*outputStream{handle.stdout, handle.stderr, handle.pty} { + stream.mu.Lock() + retained := len(stream.queue) + stream.mu.Unlock() + if retained != 0 { + t.Fatal("Close retained collecting tail") + } + // A slice already in flight may win its send concurrently with abort. + for range stream.output { + } + } +} + +func TestPTYStreamingDetach(t *testing.T) { + sandbox, closeServer := newEnvdTestSandbox(t) + defer closeServer() + for _, attach := range []bool{false, true} { + var handle *CommandHandle + var err error + if attach { + handle, err = sandbox.PTY.ConnectWithOptions(t.Context(), 7, "", &CommandConnectOptions{Streaming: &CommandStreamingOptions{StderrChannel: true}}) + } else { + handle, err = sandbox.PTY.Create(t.Context(), "sh", &PTYOptions{Streaming: &CommandStreamingOptions{PTYChannel: true}}) + } + if err != nil { + t.Fatal(err) + } + if _, err := handle.PID(t.Context()); err != nil { + t.Fatal(err) + } + handle.Close() + if _, err := waitStreamingTest(t, handle); !errors.Is(err, context.Canceled) { + t.Fatalf("PTY detach: %v", err) + } + assertNoRetainedOutput(t, handle) + } +} + +func TestPTYStreamingConnectOutput(t *testing.T) { + fixture := &streamingProcessServer{chunks: 40, size: 128, pty: true} + sandbox, _ := newStreamingTestSandbox(t, fixture) + for _, callback := range []bool{false, true} { + var got bytes.Buffer + options := &CommandConnectOptions{Streaming: &CommandStreamingOptions{PTYChannel: !callback}} + if callback { + options.OnPTY = func(chunk []byte) { got.Write(chunk) } + } + handle, err := sandbox.PTY.ConnectWithOptions(t.Context(), 42, "", options) + if err != nil { + t.Fatal(err) + } + for chunk := range handle.PTY { + got.Write(chunk) + } + if _, err := waitStreamingTest(t, handle); err != nil { + t.Fatal(err) + } + var want []byte + for i := range fixture.chunks { + want = append(want, bytes.Repeat([]byte{byte(i)}, fixture.size)...) + } + if !bytes.Equal(got.Bytes(), want) { + t.Fatal("PTY attach lost output") + } + assertNoRetainedOutput(t, handle) + handle.Close() + } +} diff --git a/packages/go-sdk/commands.go b/packages/go-sdk/commands.go index dc10fbd1..90fc2a7c 100644 --- a/packages/go-sdk/commands.go +++ b/packages/go-sdk/commands.go @@ -26,9 +26,42 @@ type CommandOptions struct { Stdin bool OnStdout func([]byte) OnStderr func([]byte) + // Streaming opts out of output capture; nil preserves collecting behavior. + Streaming *CommandStreamingOptions +} + +// CommandStreamingOptions enables output delivery without retaining command output. +// Each output uses its callback, its explicitly enabled channel, or is discarded. +// A callback and channel for the same output are mutually exclusive. +// Callbacks run synchronously and must return promptly. Their slices are read-only +// and valid only until the callback returns; copy bytes that must outlive it. +// Channel slices belong to the receiver and contain at most 32 KiB each. Channels +// are unbuffered: all enabled channels must be consumed concurrently with Wait. +// A stalled reader pauses transport reads and may eventually stall the process. +// The SDK retains no output queue and at most one pending 32 KiB channel slice, +// plus one transport event (limited to 4 MiB encoded and decompressed). Transport +// decoding and HTTP buffers add overhead; caller-retained bytes are not bounded. +// Oversized events fail the local attachment; bytes are never silently dropped. +// Reattachment adds no replay, retry, restart, or exactly-once guarantee. +type CommandStreamingOptions struct { + StdoutChannel bool + StderrChannel bool + PTYChannel bool +} + +// CommandConnectOptions configures output delivery when attaching to a process. +type CommandConnectOptions struct { + OnStdout func([]byte) + OnStderr func([]byte) + OnPTY func([]byte) + // Streaming opts out of output capture; nil preserves collecting behavior. + Streaming *CommandStreamingOptions } -// CommandResult contains collected process output. +const commandStreamChunkBytes = 32 << 10 +const commandStreamMessageBytes = 4 << 20 + +// CommandResult contains process metadata and, in collecting mode, output. type CommandResult struct { PID uint32 ExitCode int @@ -67,8 +100,10 @@ type CommandService struct { client processconnect.ProcessClient } -// CommandHandle represents a streaming process. Wait can be called without -// draining the output channels and always returns the complete collected output. +// CommandHandle represents a process attachment. In collecting mode, Wait returns +// complete output without requiring channel reads. In streaming mode, Wait returns +// only metadata and requires enabled channels to be consumed. Close detaches +// locally without killing the process or closing its stdin. type CommandHandle struct { service *CommandService ready chan struct{} @@ -82,12 +117,14 @@ type CommandHandle struct { PTY <-chan []byte Done <-chan struct{} - stdout *outputStream - stderr *outputStream - pty *outputStream - done chan struct{} - result CommandResult - err error + stdout *outputStream + stderr *outputStream + pty *outputStream + done chan struct{} + result CommandResult + err error + streaming bool + cancel context.CancelFunc } func newCommandService(sandbox *Sandbox) *CommandService { @@ -95,7 +132,8 @@ func newCommandService(sandbox *Sandbox) *CommandService { return &CommandService{sandbox: sandbox, client: client} } -// Run executes a foreground command and collects its output. +// Run executes a foreground command, draining channels and waiting for completion. +// It collects output unless CommandOptions.Streaming is set. func (service *CommandService) Run(ctx context.Context, command string, options *CommandOptions) (CommandResult, error) { handle, err := service.Start(ctx, command, options) if err != nil { @@ -131,6 +169,11 @@ func (service *CommandService) Start(ctx context.Context, command string, option if options == nil { options = &CommandOptions{} } + callbacks := outputCallbacks{stdout: options.OnStdout, stderr: options.OnStderr} + if err := validateCommandStreaming(options.Streaming, callbacks); err != nil { + return nil, err + } + ctx, cancel := context.WithCancel(ctx) config := &process.ProcessConfig{Cmd: command, Args: options.Args, Envs: options.Env} if options.Cwd != "" { config.Cwd = &options.Cwd @@ -141,35 +184,54 @@ func (service *CommandService) Start(ctx context.Context, command string, option request.Msg.Tag = &options.Tag } service.addHeaders(request.Header()) - stream, err := service.client.Start(ctx, request) + stream, err := service.outputClient(options.Streaming).Start(ctx, request) if err != nil { + cancel() return nil, connectError(err) } - handle := newCommandHandle(service, options.Tag) + handle := newCommandHandle(service, options.Tag, options.Streaming) + handle.cancel = cancel go handle.receive(ctx, func() (*process.ProcessEvent, bool) { if !stream.Receive() { return nil, false } return stream.Msg().GetEvent(), true - }, stream.Err, stream.Close, outputCallbacks{stdout: options.OnStdout, stderr: options.OnStderr}) + }, stream.Err, stream.Close, callbacks) return handle, nil } // Connect attaches to an existing process by PID or tag. func (service *CommandService) Connect(ctx context.Context, pid uint32, tag string) (*CommandHandle, error) { + return service.ConnectWithOptions(ctx, pid, tag, nil) +} + +// ConnectWithOptions attaches by PID or tag with an explicit output policy. +// Canceling ctx or calling Close detaches locally without sending a signal or EOF. +func (service *CommandService) ConnectWithOptions(ctx context.Context, pid uint32, tag string, options *CommandConnectOptions) (*CommandHandle, error) { + if options == nil { + options = &CommandConnectOptions{} + } + callbacks := outputCallbacks{stdout: options.OnStdout, stderr: options.OnStderr, pty: options.OnPTY} + if err := validateCommandStreaming(options.Streaming, callbacks); err != nil { + return nil, err + } selector, err := processSelector(pid, tag) if err != nil { return nil, err } request := connect.NewRequest(&process.ConnectRequest{Process: selector}) service.addHeaders(request.Header()) - stream, err := service.client.Connect(ctx, request) + ctx, cancel := context.WithCancel(ctx) + stream, err := service.outputClient(options.Streaming).Connect(ctx, request) if err != nil { + cancel() return nil, connectError(err) } - handle := newCommandHandle(service, tag) + handle := newCommandHandle(service, tag, options.Streaming) + handle.cancel = cancel if pid != 0 { handle.pid = pid + handle.result.PID = pid handle.closeReady() } go handle.receive(ctx, func() (*process.ProcessEvent, bool) { @@ -177,7 +239,7 @@ func (service *CommandService) Connect(ctx context.Context, pid uint32, tag stri return nil, false } return stream.Msg().GetEvent(), true - }, stream.Err, stream.Close, outputCallbacks{}) + }, stream.Err, stream.Close, callbacks) return handle, nil } @@ -235,15 +297,20 @@ func (service *CommandService) addHeaders(header http.Header) { header.Set("Keepalive-Ping-Interval", "50") } -func newCommandHandle(service *CommandService, tag string) *CommandHandle { +func newCommandHandle(service *CommandService, tag string, policy *CommandStreamingOptions) *CommandHandle { ready := make(chan struct{}) - stdout := newOutputStream() - stderr := newOutputStream() - pty := newOutputStream() + var stdout, stderr, pty *outputStream + if policy == nil { + stdout, stderr, pty = newOutputStream(), newOutputStream(), newOutputStream() + } else { + stdout = newDirectOutputStream(policy.StdoutChannel) + stderr = newDirectOutputStream(policy.StderrChannel) + pty = newDirectOutputStream(policy.PTYChannel) + } done := make(chan struct{}) handle := &CommandHandle{ service: service, ready: ready, closeReady: sync.OnceFunc(func() { close(ready) }), - tag: tag, + tag: tag, streaming: policy != nil, Stdout: stdout.output, Stderr: stderr.output, PTY: pty.output, Done: done, stdout: stdout, stderr: stderr, pty: pty, done: done, } @@ -267,7 +334,21 @@ func cleanupCommandOutputs(outputs commandOutputs) { func (handle *CommandHandle) receive(ctx context.Context, next func() (*process.ProcessEvent, bool), streamErr, closeStream func() error, callbacks outputCallbacks) { defer close(handle.done) - defer func() { _ = closeStream() }() + closeResponse := sync.OnceFunc(func() { _ = closeStream() }) + stopCancel := context.AfterFunc(ctx, func() { + cleanupCommandOutputs(commandOutputs{handle.stdout, handle.stderr, handle.pty}) + closeResponse() + }) + defer func() { + stopCancel() + if ctx.Err() != nil { + cleanupCommandOutputs(commandOutputs{handle.stdout, handle.stderr, handle.pty}) + } + closeResponse() + if handle.cancel != nil { + handle.cancel() + } + }() defer handle.stdout.close() defer handle.stderr.close() defer handle.pty.close() @@ -288,27 +369,25 @@ func (handle *CommandHandle) receive(ctx context.Context, next func() (*process. continue } if data := event.GetData(); data != nil { - switch output := data.GetOutput().(type) { + var chunk []byte + var output *outputStream + var callback func([]byte) + switch data := data.GetOutput().(type) { case *process.ProcessEvent_DataEvent_Stdout: - chunk := bytes.Clone(output.Stdout) - handle.result.Stdout = append(handle.result.Stdout, chunk...) - handle.stdout.send(chunk) - if callbacks.stdout != nil { - callbacks.stdout(chunk) + chunk, output, callback = data.Stdout, handle.stdout, callbacks.stdout + if !handle.streaming { + handle.result.Stdout = append(handle.result.Stdout, chunk...) } case *process.ProcessEvent_DataEvent_Stderr: - chunk := bytes.Clone(output.Stderr) - handle.result.Stderr = append(handle.result.Stderr, chunk...) - handle.stderr.send(chunk) - if callbacks.stderr != nil { - callbacks.stderr(chunk) + chunk, output, callback = data.Stderr, handle.stderr, callbacks.stderr + if !handle.streaming { + handle.result.Stderr = append(handle.result.Stderr, chunk...) } case *process.ProcessEvent_DataEvent_Pty: - chunk := bytes.Clone(output.Pty) - handle.pty.send(chunk) - if callbacks.pty != nil { - callbacks.pty(chunk) - } + chunk, output, callback = data.Pty, handle.pty, callbacks.pty + } + if output != nil && !handle.deliver(ctx, output, callback, chunk) { + break } } if end := event.GetEnd(); end != nil { @@ -316,17 +395,89 @@ func (handle *CommandHandle) receive(ctx context.Context, next func() (*process. handle.result.ExitCode = int(end.GetExitCode()) handle.result.Status = end.GetStatus() if end.GetExited() && end.GetExitCode() != 0 { - handle.err = &CommandExitError{Result: handle.result, Message: end.GetError()} + message := end.GetError() + if handle.streaming { + message = "" + } + handle.err = &CommandExitError{Result: handle.result, Message: message} } return } } handle.closeReady() - if err := streamErr(); err != nil && !errors.Is(err, context.Canceled) { + // Closing the response can surface as a transport read error. The attachment + // context is authoritative when local cancellation caused delivery to stop. + if ctx.Err() != nil { + handle.err = ctx.Err() + return + } + if err := streamErr(); err != nil { + if handle.streaming && !errors.Is(err, context.Canceled) && !errors.Is(err, context.DeadlineExceeded) { + // Preserve the error code without retaining server-provided output in errors. + err = connect.NewError(connect.CodeOf(err), errors.New("command attachment failed")) + } handle.err = connectError(err) } } +func validateCommandStreaming(policy *CommandStreamingOptions, callbacks outputCallbacks) error { + if policy != nil && ((policy.StdoutChannel && callbacks.stdout != nil) || + (policy.StderrChannel && callbacks.stderr != nil) || (policy.PTYChannel && callbacks.pty != nil)) { + return &InvalidArgumentError{Message: "streaming output must use either a callback or a channel"} + } + return nil +} + +func (service *CommandService) outputClient(policy *CommandStreamingOptions) processconnect.ProcessClient { + if policy == nil { + return service.client + } + return processconnect.NewProcessClient(service.sandbox.client.envdClient, service.sandbox.envdURL(envdPort, false), + connect.WithCodec(tolerantJSONCodec{}), connect.WithAcceptCompression("gzip", nil, nil), + connect.WithReadMaxBytes(commandStreamMessageBytes)) +} + +func (handle *CommandHandle) deliver(ctx context.Context, output *outputStream, callback func([]byte), chunk []byte) bool { + if !handle.streaming { + chunk = bytes.Clone(chunk) + output.send(chunk) + if callback != nil { + callback(chunk) + } + return true + } + if ctx.Err() != nil { + return false + } + if callback != nil { + callback(chunk) + } else if output.enabled { + for len(chunk) > 0 { + n := min(len(chunk), commandStreamChunkBytes) + part := bytes.Clone(chunk[:n]) + select { + case output.output <- part: + case <-ctx.Done(): + return false + } + chunk = chunk[n:] + } + } + return ctx.Err() == nil +} + +// Close cancels this local attachment and releases queued output, including an +// unread collecting-mode tail. It does not kill the process or send stdin EOF. +// Close does not wait for user callbacks; Done closes after the receiver exits. +// A callback must return before Wait can complete. Close is safe to call repeatedly. +func (handle *CommandHandle) Close() error { + if handle.cancel != nil { + handle.cancel() + } + cleanupCommandOutputs(commandOutputs{handle.stdout, handle.stderr, handle.pty}) + return nil +} + type outputStream struct { output chan []byte wake chan struct{} @@ -335,6 +486,8 @@ type outputStream struct { mu sync.Mutex queue [][]byte closed bool + direct bool + enabled bool } func newOutputStream() *outputStream { @@ -351,6 +504,15 @@ func newOutputStream() *outputStream { return stream } +func newDirectOutputStream(enabled bool) *outputStream { + stream := &outputStream{output: make(chan []byte), direct: true, enabled: enabled} + stream.abortOnce = func() {} // Direct delivery owns no queue or background goroutine. + if !enabled { + close(stream.output) + } + return stream +} + func (stream *outputStream) send(chunk []byte) { stream.mu.Lock() if !stream.closed { @@ -361,6 +523,12 @@ func (stream *outputStream) send(chunk []byte) { } func (stream *outputStream) close() { + if stream.direct { + if stream.enabled { + close(stream.output) + } + return + } stream.mu.Lock() stream.closed = true stream.mu.Unlock() @@ -420,7 +588,12 @@ func (handle *CommandHandle) PID(ctx context.Context) (uint32, error) { } } -// Wait waits for completion and returns collected output. +// Wait waits for receiver completion. Collecting mode preserves the queued channel +// tail and returns full output. Streaming mode returns metadata with nil output; +// all enabled channels must be read concurrently, and are closed before Done. +// Attachment cancellation returns a cancellation error unless a process end event +// was already confirmed. CommandExitError preserves confirmed nonzero exits. +// Canceling only Wait's context stops waiting; use Close to cancel the attachment. func (handle *CommandHandle) Wait(ctx context.Context) (CommandResult, error) { select { case <-ctx.Done(): diff --git a/packages/go-sdk/envd_test.go b/packages/go-sdk/envd_test.go index 910328e9..81610ae6 100644 --- a/packages/go-sdk/envd_test.go +++ b/packages/go-sdk/envd_test.go @@ -651,7 +651,7 @@ func TestCommandHelpersAndFileMappings(t *testing.T) { if _, err := sandbox.PTY.Create(ctx, "", nil); err == nil { t.Fatal("expected empty PTY command") } - handle := newCommandHandle(sandbox.Commands, "") + handle := newCommandHandle(sandbox.Commands, "", nil) canceled, cancel := context.WithCancel(ctx) cancel() if _, err := handle.PID(canceled); err == nil { diff --git a/packages/go-sdk/examples_test.go b/packages/go-sdk/examples_test.go index 1468b63f..65745a21 100644 --- a/packages/go-sdk/examples_test.go +++ b/packages/go-sdk/examples_test.go @@ -24,3 +24,46 @@ func ExampleTemplateBuilder() { template := agentbox.NewTemplate(".").FromPython("3.13").Copy("requirements.txt", "/app/", nil).PipInstall().Workdir("/app") _, _ = template.JSON() } + +func ExampleCommandService_Start_streaming() { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + client, err := agentbox.NewClient() + if err != nil { + log.Fatal(err) + } + sandbox, err := client.Sandboxes.Create(ctx, nil) + if err != nil { + log.Fatal(err) + } + defer sandbox.Kill(context.Background()) + handle, err := sandbox.Commands.Start(ctx, "cat", &agentbox.CommandOptions{ + Stdin: true, + Streaming: &agentbox.CommandStreamingOptions{}, + OnStdout: func(chunk []byte) { log.Printf("output: %s", chunk) }, + // With no callback or channel selected, stderr is discarded. + }) + if err != nil { + log.Fatal(err) + } + defer handle.Close() + pid, err := handle.PID(ctx) + if err != nil { + log.Fatal(err) + } + _, _ = handle.Write(ctx, []byte("hello\n")) + // Close only detaches. The process remains alive with stdin open. + _ = handle.Close() + _, _ = handle.Wait(ctx) + attached, err := sandbox.Commands.ConnectWithOptions(ctx, pid, "", &agentbox.CommandConnectOptions{ + Streaming: &agentbox.CommandStreamingOptions{}, + OnStdout: func(chunk []byte) { log.Printf("output: %s", chunk) }, + }) + if err != nil { + log.Fatal(err) + } + defer attached.Close() + _, _ = attached.Write(ctx, []byte("hello again\n")) + _ = attached.CloseStdin(ctx) + _, _ = attached.Wait(ctx) +} diff --git a/packages/go-sdk/integration/command_streaming_test.go b/packages/go-sdk/integration/command_streaming_test.go new file mode 100644 index 00000000..cc8a0a37 --- /dev/null +++ b/packages/go-sdk/integration/command_streaming_test.go @@ -0,0 +1,88 @@ +//go:build integration + +package integration_test + +import ( + "bytes" + "context" + "errors" + "testing" + "time" + + "github.com/abox-dev/sdk/packages/go-sdk" +) + +func TestCommandStreamingAttachDetachKVM(t *testing.T) { + client, err := agentbox.NewClient() + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Minute) + defer cancel() + sandbox, err := client.Sandboxes.Create(ctx, &agentbox.CreateSandboxOptions{Timeout: 3 * time.Minute, Metadata: map[string]string{"sdk": "go-streaming-smoke"}}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + cleanupCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + if _, err := sandbox.Kill(cleanupCtx); err != nil { + t.Errorf("sandbox cleanup: %v", err) + } + }) + policy := &agentbox.CommandStreamingOptions{StdoutChannel: true} + handle, err := sandbox.Commands.Start(ctx, "cat", &agentbox.CommandOptions{Stdin: true, Tag: "streaming-smoke", Streaming: policy}) + if err != nil { + t.Fatal(err) + } + defer handle.Close() + pid, err := handle.PID(ctx) + if err != nil { + t.Fatal(err) + } + exchange := func(handle *agentbox.CommandHandle, message string) { + t.Helper() + if _, err := handle.Write(ctx, []byte(message)); err != nil { + t.Fatal(err) + } + var got []byte + for !bytes.Contains(got, []byte(message)) { + select { + case chunk, ok := <-handle.Stdout: + if !ok { + t.Fatal("process output closed") + } + got = append(got, chunk...) + case <-ctx.Done(): + t.Fatal(ctx.Err()) + } + } + } + exchange(handle, "before detach\n") + for i := range 3 { + handle.Close() + result, err := handle.Wait(ctx) + if !errors.Is(err, context.Canceled) || result.PID != pid || len(result.Stdout) != 0 || len(result.Stderr) != 0 { + t.Fatalf("detach: %+v %v", result, err) + } + attachPID, tag := pid, "" + if i%2 == 1 { + attachPID, tag = 0, "streaming-smoke" + } + handle, err = sandbox.Commands.ConnectWithOptions(ctx, attachPID, tag, &agentbox.CommandConnectOptions{Streaming: policy}) + if err != nil { + t.Fatal(err) + } + defer handle.Close() + exchange(handle, "after detach\n") + } + if err := handle.CloseStdin(ctx); err != nil { + t.Fatal(err) + } + for range handle.Stdout { + } + result, err := handle.Wait(ctx) + if err != nil || result.ExitCode != 0 || result.PID != pid || result.Stdout != nil || result.Stderr != nil { + t.Fatalf("exit: %+v %v", result, err) + } +} diff --git a/packages/go-sdk/pty.go b/packages/go-sdk/pty.go index bb18137d..af0e534f 100644 --- a/packages/go-sdk/pty.go +++ b/packages/go-sdk/pty.go @@ -16,6 +16,8 @@ type PTYOptions struct { Cols uint32 Rows uint32 OnPTY func([]byte) + // Streaming selects bounded delivery; nil preserves collecting-mode PTY queues. + Streaming *CommandStreamingOptions } // PTYService manages pseudo-terminal processes. @@ -29,6 +31,11 @@ func (service *PTYService) Create(ctx context.Context, command string, options * if options == nil { options = &PTYOptions{} } + callbacks := outputCallbacks{pty: options.OnPTY} + if err := validateCommandStreaming(options.Streaming, callbacks); err != nil { + return nil, err + } + ctx, cancel := context.WithCancel(ctx) cols, rows := options.Cols, options.Rows if cols == 0 { cols = 80 @@ -45,17 +52,19 @@ func (service *PTYService) Create(ctx context.Context, command string, options * request.Msg.Tag = &options.Tag } service.commands.addHeaders(request.Header()) - stream, err := service.commands.client.Start(ctx, request) + stream, err := service.commands.outputClient(options.Streaming).Start(ctx, request) if err != nil { + cancel() return nil, connectError(err) } - handle := newCommandHandle(service.commands, options.Tag) + handle := newCommandHandle(service.commands, options.Tag, options.Streaming) + handle.cancel = cancel go handle.receive(ctx, func() (*process.ProcessEvent, bool) { if !stream.Receive() { return nil, false } return stream.Msg().GetEvent(), true - }, stream.Err, stream.Close, outputCallbacks{pty: options.OnPTY}) + }, stream.Err, stream.Close, callbacks) return handle, nil } @@ -64,6 +73,12 @@ func (service *PTYService) Connect(ctx context.Context, pid uint32, tag string) return service.commands.Connect(ctx, pid, tag) } +// ConnectWithOptions attaches to an existing PTY using the given output policy. +// Use OnPTY or Streaming.PTYChannel for terminal output. +func (service *PTYService) ConnectWithOptions(ctx context.Context, pid uint32, tag string, options *CommandConnectOptions) (*CommandHandle, error) { + return service.commands.ConnectWithOptions(ctx, pid, tag, options) +} + // Input sends terminal input. func (service *PTYService) Input(ctx context.Context, handle *CommandHandle, data []byte) error { pid, err := handle.PID(ctx) diff --git a/reference/manifest.json b/reference/manifest.json index dc0829ca..35bf4ace 100644 --- a/reference/manifest.json +++ b/reference/manifest.json @@ -45,7 +45,7 @@ "sdk/cli/sandbox.md": "16b64c5a4e932eca2452e400e0faf154eba1d3c7373513aa7cd9744e02908313", "sdk/cli/template.md": "f270b22b9ee10a9b954a14f24e04c5449b4d5823bb28002ffde3ed682594cfc5", "sdk/go/code-interpreter.md": "9c56333b8181878d656baa27a4b56b768a0cb01c9d764f7ac8d52d32d8faa70c", - "sdk/go/core.md": "6fafc6064961034202df3c0d78beed0f82cb508b00f5225f37bf3297930fd630", + "sdk/go/core.md": "a74a267c34eba6d43635f92008500929739be68113ade0b6cbf16602d3bd8d7c", "sdk/javascript/code-interpreter/README.md": "a787206eb86f81fc306b373ce20bc6c91464b0ab3c3ec7650eccbb6be5f676c5", "sdk/javascript/code-interpreter/classes/Sandbox.md": "19d9f1f1e8847606a54d93c5a6e5d4c89f936a745866e8d6a2aa06758b605e60", "sdk/javascript/code-interpreter/enumerations/ChartType.md": "cf1ba00cd3258cc88814b7215e101b0d775ceb15dd833636d58ae538d8107ec1", diff --git a/reference/sdk/go/core.md b/reference/sdk/go/core.md index 5160e01b..7a98b7c8 100644 --- a/reference/sdk/go/core.md +++ b/reference/sdk/go/core.md @@ -58,9 +58,11 @@ Create a client, start a sandbox, and run a command: - [func WithProxy\(value string\) ClientOption](<#WithProxy>) - [func WithRequestTimeout\(timeout time.Duration\) ClientOption](<#WithRequestTimeout>) - [func WithSandboxURL\(value string\) ClientOption](<#WithSandboxURL>) +- [type CommandConnectOptions](<#CommandConnectOptions>) - [type CommandExitError](<#CommandExitError>) - [func \(e \*CommandExitError\) Error\(\) string](<#CommandExitError.Error>) - [type CommandHandle](<#CommandHandle>) + - [func \(handle \*CommandHandle\) Close\(\) error](<#CommandHandle.Close>) - [func \(handle \*CommandHandle\) CloseStdin\(ctx context.Context\) error](<#CommandHandle.CloseStdin>) - [func \(handle \*CommandHandle\) Kill\(ctx context.Context\) error](<#CommandHandle.Kill>) - [func \(handle \*CommandHandle\) PID\(ctx context.Context\) \(uint32, error\)](<#CommandHandle.PID>) @@ -70,11 +72,13 @@ Create a client, start a sandbox, and run a command: - [type CommandResult](<#CommandResult>) - [type CommandService](<#CommandService>) - [func \(service \*CommandService\) Connect\(ctx context.Context, pid uint32, tag string\) \(\*CommandHandle, error\)](<#CommandService.Connect>) + - [func \(service \*CommandService\) ConnectWithOptions\(ctx context.Context, pid uint32, tag string, options \*CommandConnectOptions\) \(\*CommandHandle, error\)](<#CommandService.ConnectWithOptions>) - [func \(service \*CommandService\) Kill\(ctx context.Context, pid uint32, tag string\) error](<#CommandService.Kill>) - [func \(service \*CommandService\) List\(ctx context.Context\) \(\[\]ProcessInfo, error\)](<#CommandService.List>) - [func \(service \*CommandService\) Run\(ctx context.Context, command string, options \*CommandOptions\) \(CommandResult, error\)](<#CommandService.Run>) - [func \(service \*CommandService\) Start\(ctx context.Context, command string, options \*CommandOptions\) \(\*CommandHandle, error\)](<#CommandService.Start>) - [func \(service \*CommandService\) Terminate\(ctx context.Context, pid uint32, tag string\) error](<#CommandService.Terminate>) +- [type CommandStreamingOptions](<#CommandStreamingOptions>) - [type ConnectSandboxOptions](<#ConnectSandboxOptions>) - [type CopyOptions](<#CopyOptions>) - [type CreateSandboxOptions](<#CreateSandboxOptions>) @@ -120,6 +124,7 @@ Create a client, start a sandbox, and run a command: - [type PTYOptions](<#PTYOptions>) - [type PTYService](<#PTYService>) - [func \(service \*PTYService\) Connect\(ctx context.Context, pid uint32, tag string\) \(\*CommandHandle, error\)](<#PTYService.Connect>) + - [func \(service \*PTYService\) ConnectWithOptions\(ctx context.Context, pid uint32, tag string, options \*CommandConnectOptions\) \(\*CommandHandle, error\)](<#PTYService.ConnectWithOptions>) - [func \(service \*PTYService\) Create\(ctx context.Context, command string, options \*PTYOptions\) \(\*CommandHandle, error\)](<#PTYService.Create>) - [func \(service \*PTYService\) Input\(ctx context.Context, handle \*CommandHandle, data \[\]byte\) error](<#PTYService.Input>) - [func \(service \*PTYService\) Kill\(ctx context.Context, handle \*CommandHandle\) error](<#PTYService.Kill>) @@ -534,6 +539,19 @@ WithRequestTimeout sets the default unary request timeout. Zero disables it. WithSandboxURL overrides the sandbox proxy URL. + +## type CommandConnectOptions + +CommandConnectOptions configures output delivery when attaching to a process. + + type CommandConnectOptions struct { + OnStdout func([]byte) + OnStderr func([]byte) + OnPTY func([]byte) + // Streaming opts out of output capture; nil preserves collecting behavior. + Streaming *CommandStreamingOptions + } + ## type CommandExitError @@ -554,7 +572,7 @@ Error describes the non\-zero command exit code. ## type CommandHandle -CommandHandle represents a streaming process. Wait can be called without draining the output channels and always returns the complete collected output. +CommandHandle represents a process attachment. In collecting mode, Wait returns complete output without requiring channel reads. In streaming mode, Wait returns only metadata and requires enabled channels to be consumed. Close detaches locally without killing the process or closing its stdin. type CommandHandle struct { Stdout <-chan []byte @@ -564,6 +582,13 @@ CommandHandle represents a streaming process. Wait can be called without drainin // contains filtered or unexported fields } + +### func \(\*CommandHandle\) Close + + func (handle *CommandHandle) Close() error + +Close cancels this local attachment and releases queued output, including an unread collecting\-mode tail. It does not kill the process or send stdin EOF. Close does not wait for user callbacks; Done closes after the receiver exits. A callback must return before Wait can complete. Close is safe to call repeatedly. + ### func \(\*CommandHandle\) CloseStdin @@ -590,7 +615,7 @@ PID waits for and returns the process identifier. func (handle *CommandHandle) Wait(ctx context.Context) (CommandResult, error) -Wait waits for completion and returns collected output. +Wait waits for receiver completion. Collecting mode preserves the queued channel tail and returns full output. Streaming mode returns metadata with nil output; all enabled channels must be read concurrently, and are closed before Done. Attachment cancellation returns a cancellation error unless a process end event was already confirmed. CommandExitError preserves confirmed nonzero exits. Canceling only Wait's context stops waiting; use Close to cancel the attachment. ### func \(\*CommandHandle\) Write @@ -612,12 +637,14 @@ CommandOptions configures a command process. Stdin bool OnStdout func([]byte) OnStderr func([]byte) + // Streaming opts out of output capture; nil preserves collecting behavior. + Streaming *CommandStreamingOptions } ## type CommandResult -CommandResult contains collected process output. +CommandResult contains process metadata and, in collecting mode, output. type CommandResult struct { PID uint32 @@ -643,6 +670,13 @@ CommandService executes and manages sandbox processes. Connect attaches to an existing process by PID or tag. + +### func \(\*CommandService\) ConnectWithOptions + + func (service *CommandService) ConnectWithOptions(ctx context.Context, pid uint32, tag string, options *CommandConnectOptions) (*CommandHandle, error) + +ConnectWithOptions attaches by PID or tag with an explicit output policy. Canceling ctx or calling Close detaches locally without sending a signal or EOF. + ### func \(\*CommandService\) Kill @@ -662,7 +696,7 @@ List returns currently running processes. func (service *CommandService) Run(ctx context.Context, command string, options *CommandOptions) (CommandResult, error) -Run executes a foreground command and collects its output. +Run executes a foreground command, draining channels and waiting for completion. It collects output unless CommandOptions.Streaming is set. ### func \(\*CommandService\) Start @@ -671,6 +705,57 @@ Run executes a foreground command and collects its output. Start starts a process and streams output through the returned handle. +###### Example (Streaming) + + + + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + client, err := agentbox.NewClient() + if err != nil { + log.Fatal(err) + } + sandbox, err := client.Sandboxes.Create(ctx, nil) + if err != nil { + log.Fatal(err) + } + defer sandbox.Kill(context.Background()) + handle, err := sandbox.Commands.Start(ctx, "cat", &agentbox.CommandOptions{ + Stdin: true, + Streaming: &agentbox.CommandStreamingOptions{}, + OnStdout: func(chunk []byte) { log.Printf("output: %s", chunk) }, + // With no callback or channel selected, stderr is discarded. + }) + if err != nil { + log.Fatal(err) + } + defer handle.Close() + pid, err := handle.PID(ctx) + if err != nil { + log.Fatal(err) + } + _, _ = handle.Write(ctx, []byte("hello\n")) + // Close only detaches. The process remains alive with stdin open. + _ = handle.Close() + _, _ = handle.Wait(ctx) + attached, err := sandbox.Commands.ConnectWithOptions(ctx, pid, "", &agentbox.CommandConnectOptions{ + Streaming: &agentbox.CommandStreamingOptions{}, + OnStdout: func(chunk []byte) { log.Printf("output: %s", chunk) }, + }) + if err != nil { + log.Fatal(err) + } + defer attached.Close() + _, _ = attached.Write(ctx, []byte("hello again\n")) + _ = attached.CloseStdin(ctx) + _, _ = attached.Wait(ctx) + + + + + + ### func \(\*CommandService\) Terminate @@ -678,6 +763,17 @@ Start starts a process and streams output through the returned handle. Terminate sends SIGTERM to a process. + +## type CommandStreamingOptions + +CommandStreamingOptions enables output delivery without retaining command output. Each output uses its callback, its explicitly enabled channel, or is discarded. A callback and channel for the same output are mutually exclusive. Callbacks run synchronously and must return promptly. Their slices are read\-only and valid only until the callback returns; copy bytes that must outlive it. Channel slices belong to the receiver and contain at most 32 KiB each. Channels are unbuffered: all enabled channels must be consumed concurrently with Wait. A stalled reader pauses transport reads and may eventually stall the process. The SDK retains no output queue and at most one pending 32 KiB channel slice, plus one transport event \(limited to 4 MiB encoded and decompressed\). Transport decoding and HTTP buffers add overhead; caller\-retained bytes are not bounded. Oversized events fail the local attachment; bytes are never silently dropped. Reattachment adds no replay, retry, restart, or exactly\-once guarantee. + + type CommandStreamingOptions struct { + StdoutChannel bool + StderrChannel bool + PTYChannel bool + } + ## type ConnectSandboxOptions @@ -1061,6 +1157,8 @@ PTYOptions configures an interactive terminal. Cols uint32 Rows uint32 OnPTY func([]byte) + // Streaming selects bounded delivery; nil preserves collecting-mode PTY queues. + Streaming *CommandStreamingOptions } @@ -1079,6 +1177,13 @@ PTYService manages pseudo\-terminal processes. Connect attaches to an existing PTY process. + +### func \(\*PTYService\) ConnectWithOptions + + func (service *PTYService) ConnectWithOptions(ctx context.Context, pid uint32, tag string, options *CommandConnectOptions) (*CommandHandle, error) + +ConnectWithOptions attaches to an existing PTY using the given output policy. Use OnPTY or Streaming.PTYChannel for terminal output. + ### func \(\*PTYService\) Create From 6b213d55d10e58f41150022ed9321231b39e2163 Mon Sep 17 00:00:00 2001 From: Ilyas Salikhov Date: Tue, 22 Sep 2026 13:54:15 +0300 Subject: [PATCH 2/3] chore: release AgentBox SDK v0.1.8 --- docs/go-command-streaming.md | 8 ++++---- packages/cli/package.json | 2 +- packages/code-interpreter-js/package.json | 2 +- packages/code-interpreter-python/package.json | 2 +- packages/code-interpreter-python/pyproject.toml | 2 +- packages/code-interpreter-python/uv.lock | 4 ++-- packages/go-sdk/README.md | 4 ++-- packages/go-sdk/version.go | 2 +- packages/js-sdk/package.json | 2 +- packages/python-sdk/package.json | 2 +- packages/python-sdk/pyproject.toml | 2 +- packages/python-sdk/uv.lock | 2 +- reference/manifest.json | 14 +++++++------- reference/sdk/go/core.md | 2 +- 14 files changed, 25 insertions(+), 25 deletions(-) diff --git a/docs/go-command-streaming.md b/docs/go-command-streaming.md index cbf8c2e5..603c8700 100644 --- a/docs/go-command-streaming.md +++ b/docs/go-command-streaming.md @@ -93,7 +93,7 @@ regressions remain required. The supported Go matrix is 1.24–1.27; run `make go-check` and generation/artifact gates described in [`RELEASING.md`](../RELEASING.md) using the pinned containers. -Validation passed: `make go-check` (91.2% handwritten statement coverage), builds +Validation passed: `make go-check` (91.1% handwritten statement coverage), builds and hermetic tests on Go 1.24.13, 1.25.14, 1.26.8 and 1.27.1, race checks on 1.24.13 and 1.27.1, format/lint/type checks, workspace tests, release-artifact builds and clean installation checks. `make generate` and reference-contract @@ -123,6 +123,6 @@ It creates and deletes only its own sandbox. This runtime smoke has passed. The implementation must follow the coordinated SDK release procedure in [`RELEASING.md`](../RELEASING.md). No ad-hoc fork or replacement module is needed. -Publication and the consumer dependency upgrade remain release steps; no released -module version is assigned by this source change. The consuming orchestrator must -use the released version and explicitly enable streaming on launch and reconnect. +The coordinated release version is 0.1.8. Consumers must upgrade to +`github.com/abox-dev/sdk/packages/go-sdk@v0.1.8` and explicitly enable streaming +on both launch and reconnect. diff --git a/packages/cli/package.json b/packages/cli/package.json index 6ca67883..18c875d9 100644 --- a/packages/cli/package.json +++ b/packages/cli/package.json @@ -1,6 +1,6 @@ { "name": "@abox-dev/cli", - "version": "0.1.7", + "version": "0.1.8", "description": "CLI for AgentBox sandboxes and templates", "homepage": "https://docs.agentbox.ru/en/cli/", "license": "MIT", diff --git a/packages/code-interpreter-js/package.json b/packages/code-interpreter-js/package.json index b62cc8a9..b93cf0ec 100644 --- a/packages/code-interpreter-js/package.json +++ b/packages/code-interpreter-js/package.json @@ -1,6 +1,6 @@ { "name": "@abox-dev/code-interpreter", - "version": "0.1.7", + "version": "0.1.8", "packageManager": "pnpm@10.34.5", "description": "AgentBox Code Interpreter - Stateful code execution", "homepage": "https://docs.agentbox.ru/en/sdk/code-interpreter/", diff --git a/packages/code-interpreter-python/package.json b/packages/code-interpreter-python/package.json index a1a9cade..4cd53c63 100644 --- a/packages/code-interpreter-python/package.json +++ b/packages/code-interpreter-python/package.json @@ -1,7 +1,7 @@ { "name": "@abox-dev/code-interpreter-python", "private": true, - "version": "0.1.7", + "version": "0.1.8", "scripts": { "test": "uv run pytest -n 2 --verbose -x tests/test_sandbox_url.py", "test:integration": "uv run pytest -n 2 --verbose -x", diff --git a/packages/code-interpreter-python/pyproject.toml b/packages/code-interpreter-python/pyproject.toml index 3a86de1d..1bb86bc8 100644 --- a/packages/code-interpreter-python/pyproject.toml +++ b/packages/code-interpreter-python/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "abox-code-interpreter" -version = "0.1.7" +version = "0.1.8" description = "AgentBox Code Interpreter - Stateful code execution" authors = [{ name = "RetailDriver LLC" }] license = "MIT" diff --git a/packages/code-interpreter-python/uv.lock b/packages/code-interpreter-python/uv.lock index 4de9fd00..63963386 100644 --- a/packages/code-interpreter-python/uv.lock +++ b/packages/code-interpreter-python/uv.lock @@ -16,7 +16,7 @@ members = [ [[package]] name = "abox-code-interpreter" -version = "0.1.7" +version = "0.1.8" source = { editable = "." } dependencies = [ { name = "abox-sdk" }, @@ -58,7 +58,7 @@ dev = [ [[package]] name = "abox-sdk" -version = "0.1.7" +version = "0.1.8" source = { editable = "../python-sdk" } dependencies = [ { name = "attrs" }, diff --git a/packages/go-sdk/README.md b/packages/go-sdk/README.md index 4f22c77f..ff297060 100644 --- a/packages/go-sdk/README.md +++ b/packages/go-sdk/README.md @@ -44,8 +44,8 @@ func main() { Code Interpreter is available from `github.com/abox-dev/sdk/packages/go-sdk/codeinterpreter`. -API reference for this release: [core SDK on pkg.go.dev](https://pkg.go.dev/github.com/abox-dev/sdk/packages/go-sdk@v0.1.7) and -[Code Interpreter on pkg.go.dev](https://pkg.go.dev/github.com/abox-dev/sdk/packages/go-sdk/codeinterpreter@v0.1.7). +API reference for this release: [core SDK on pkg.go.dev](https://pkg.go.dev/github.com/abox-dev/sdk/packages/go-sdk@v0.1.8) and +[Code Interpreter on pkg.go.dev](https://pkg.go.dev/github.com/abox-dev/sdk/packages/go-sdk/codeinterpreter@v0.1.8). Documentation: [core SDK](https://docs.agentbox.ru/en/sdk/), [sandboxes](https://docs.agentbox.ru/en/sdk/sandboxes/), diff --git a/packages/go-sdk/version.go b/packages/go-sdk/version.go index b3737d13..926be685 100644 --- a/packages/go-sdk/version.go +++ b/packages/go-sdk/version.go @@ -1,4 +1,4 @@ package agentbox // Version is the AgentBox SDK release version. -const Version = "0.1.7" +const Version = "0.1.8" diff --git a/packages/js-sdk/package.json b/packages/js-sdk/package.json index 62d1c54b..7d587f33 100644 --- a/packages/js-sdk/package.json +++ b/packages/js-sdk/package.json @@ -1,6 +1,6 @@ { "name": "@abox-dev/sdk", - "version": "0.1.7", + "version": "0.1.8", "description": "AgentBox SDK for secure cloud sandboxes", "homepage": "https://docs.agentbox.ru/en/sdk/", "license": "MIT", diff --git a/packages/python-sdk/package.json b/packages/python-sdk/package.json index 68286d34..61506ba1 100644 --- a/packages/python-sdk/package.json +++ b/packages/python-sdk/package.json @@ -1,7 +1,7 @@ { "name": "@abox-dev/python-sdk", "private": true, - "version": "0.1.7", + "version": "0.1.8", "scripts": { "test": "uv run pytest -n 4 --verbose -x tests/test_*.py tests/shared", "test:integration": "uv run pytest -n 4 --verbose -x tests/sync tests/async", diff --git a/packages/python-sdk/pyproject.toml b/packages/python-sdk/pyproject.toml index 9ea3bd89..ab9d5ebe 100644 --- a/packages/python-sdk/pyproject.toml +++ b/packages/python-sdk/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "abox-sdk" -version = "0.1.7" +version = "0.1.8" description = "AgentBox SDK for secure cloud sandboxes" authors = [{ name = "RetailDriver LLC" }] license = "MIT" diff --git a/packages/python-sdk/uv.lock b/packages/python-sdk/uv.lock index 615a6716..92a38cc9 100644 --- a/packages/python-sdk/uv.lock +++ b/packages/python-sdk/uv.lock @@ -8,7 +8,7 @@ resolution-markers = [ [[package]] name = "abox-sdk" -version = "0.1.7" +version = "0.1.8" source = { editable = "." } dependencies = [ { name = "attrs" }, diff --git a/reference/manifest.json b/reference/manifest.json index 35bf4ace..1d2e6204 100644 --- a/reference/manifest.json +++ b/reference/manifest.json @@ -45,7 +45,7 @@ "sdk/cli/sandbox.md": "16b64c5a4e932eca2452e400e0faf154eba1d3c7373513aa7cd9744e02908313", "sdk/cli/template.md": "f270b22b9ee10a9b954a14f24e04c5449b4d5823bb28002ffde3ed682594cfc5", "sdk/go/code-interpreter.md": "9c56333b8181878d656baa27a4b56b768a0cb01c9d764f7ac8d52d32d8faa70c", - "sdk/go/core.md": "a74a267c34eba6d43635f92008500929739be68113ade0b6cbf16602d3bd8d7c", + "sdk/go/core.md": "5f53d0230c06cde85d119df4f21cfeb6c69ce5f84ef65d3b86776b76ad7bc2ea", "sdk/javascript/code-interpreter/README.md": "a787206eb86f81fc306b373ce20bc6c91464b0ab3c3ec7650eccbb6be5f676c5", "sdk/javascript/code-interpreter/classes/Sandbox.md": "19d9f1f1e8847606a54d93c5a6e5d4c89f936a745866e8d6a2aa06758b605e60", "sdk/javascript/code-interpreter/enumerations/ChartType.md": "cf1ba00cd3258cc88814b7215e101b0d775ceb15dd833636d58ae538d8107ec1", @@ -185,12 +185,12 @@ }, "monoRevision": "a7d59d50e3dd84faee59aca08df1b5bb175266e6", "packages": { - "@abox-dev/cli": "0.1.7", - "@abox-dev/code-interpreter": "0.1.7", - "@abox-dev/sdk": "0.1.7", - "abox-code-interpreter": "0.1.7", - "abox-sdk": "0.1.7", - "github.com/abox-dev/sdk/packages/go-sdk": "0.1.7" + "@abox-dev/cli": "0.1.8", + "@abox-dev/code-interpreter": "0.1.8", + "@abox-dev/sdk": "0.1.8", + "abox-code-interpreter": "0.1.8", + "abox-sdk": "0.1.8", + "github.com/abox-dev/sdk/packages/go-sdk": "0.1.8" }, "schemaVersion": 1 } diff --git a/reference/sdk/go/core.md b/reference/sdk/go/core.md index 7a98b7c8..2ce9a329 100644 --- a/reference/sdk/go/core.md +++ b/reference/sdk/go/core.md @@ -266,7 +266,7 @@ Create a client, start a sandbox, and run a command: Version is the AgentBox SDK release version. - const Version = "0.1.7" + const Version = "0.1.8" ## func IAMTokenPlaceholder From 8c09a8d4a92bc92073453049e2b34bb707400cbb Mon Sep 17 00:00:00 2001 From: Ilyas Salikhov Date: Tue, 22 Sep 2026 13:58:36 +0300 Subject: [PATCH 3/3] test(go): allow CI headroom for streaming flood fixtures --- packages/go-sdk/command_streaming_test.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/packages/go-sdk/command_streaming_test.go b/packages/go-sdk/command_streaming_test.go index e869bc56..23afaad1 100644 --- a/packages/go-sdk/command_streaming_test.go +++ b/packages/go-sdk/command_streaming_test.go @@ -86,7 +86,9 @@ func newStreamingTestSandbox(t *testing.T, fixture *streamingProcessServer) (*Sa } func waitStreamingTest(t *testing.T, handle *CommandHandle) (CommandResult, error) { t.Helper() - ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + // Flood fixtures decode several MiB through HTTP; allow shared CI runners + // headroom while keeping an explicit deadline for a stuck receiver. + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) defer cancel() return handle.Wait(ctx) }