From 523a84b98e2bdc66bd0bedd0f72bc168a3bf65cc Mon Sep 17 00:00:00 2001 From: sam Date: Thu, 8 Oct 2026 16:23:01 +0800 Subject: [PATCH] Align execution journal with Runtime frame limits --- apps/daemon/internal/transport/ws.go | 2 +- docs/runtime-protocol.md | 2 + docs/zh/runtime-protocol.md | 4 +- internal/agentdaemon/proto/envelope.go | 3 + services/core/IMPLEMENTATION.md | 2 +- services/core/internal/execution/delivery.go | 2 +- services/core/internal/execution/journal.go | 6 +- .../internal/execution/journal_failure.go | 63 ++++++++++++++ .../execution/journal_failure_test.go | 49 +++++++++++ .../core/internal/execution/journal_test.go | 58 +++++++++++++ .../core/internal/runtimegateway/session.go | 2 +- services/core/internal/sessions/journal.go | 12 +-- .../core/internal/sessions/journal_test.go | 22 +++++ .../internal/store/command_output_test.go | 86 +++++++++++++++++++ 14 files changed, 302 insertions(+), 11 deletions(-) create mode 100644 services/core/internal/execution/journal_failure.go create mode 100644 services/core/internal/execution/journal_failure_test.go diff --git a/apps/daemon/internal/transport/ws.go b/apps/daemon/internal/transport/ws.go index feeb2a7a5..bf1bf1f50 100644 --- a/apps/daemon/internal/transport/ws.go +++ b/apps/daemon/internal/transport/ws.go @@ -122,7 +122,7 @@ func Dial(ctx context.Context, opts DialOptions) (*Conn, error) { } // Bound a single inbound frame so a misbehaving server can't OOM // us. Matches the gateway's 4 MiB outbound ceiling. - wsConn.SetReadLimit(4 * 1024 * 1024) + wsConn.SetReadLimit(proto.MaxFrameBytes) c := &Conn{ ws: wsConn, diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index aa2cb8441..a9607ab4e 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -6,6 +6,8 @@ This protocol connects Core to a Runtime daemon after the daemon has its machine Hosted and self-hosted Runtimes use the same protocol. A Harness joins through the [Harness adapter contract](../contracts/agents-api/harness-onboarding.md), which owns the Executor and Turn lifecycle obligations behind the Runtime registry. +Connection frames are limited to 4 MiB by `proto.MaxFrameBytes`. Core journals a transport-valid observation without truncating its payload. Ordinary journal batches hold up to 64 observations and 1 MiB; a larger observation is stored alone. Each Turn remains limited to 65,536 observations and 32 MiB, with one reserved terminal outcome beyond those limits. Journal failures fail the Turn and retain committed Items. The `execution journal failed` log records Project, Session and Turn IDs, trace context, stage, event kind, byte count, journal position and a bounded failure category or SQLSTATE; it never records event payloads or raw error text. + ## Ownership and connection Core owns durable Session, Turn, input and Environment records, scheduling and reconciliation. Runtime owns native Executors, active Turns, transfer state and cleanup until settlement. A Sandbox Provider owns placement and the surrounding compute. Releasing an execution admission or closing an Executor never deletes, suspends or reclaims a sandbox. The daemon is not an isolation boundary; see [Runtime and outer isolation](./concepts.md#runtime-and-outer-isolation). diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 8a3077ceb..96a80ffac 100644 --- a/docs/zh/runtime-protocol.md +++ b/docs/zh/runtime-protocol.md @@ -1,13 +1,15 @@ --- title: "Core–Runtime 协议" source: docs/runtime-protocol.md -source_hash: cbc3c6419e2d4991df82d5bbfe3b556d35cccb57d1e4e7437facca7487caabcc +source_hash: 01a9c4ddb5066c7f1aef74f5fc2a6833b389cdc16168d179221b95134dae1978 --- 此协议在 Runtime daemon 获取机器凭据后连接 Core 与 daemon,定义 daemon 连接上消息的含义和顺序。wire 类型、限制和验证器仅在 [`internal/agentdaemon/proto`](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/internal/agentdaemon/proto) 中定义一次;Core 的 [gateway](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/services/core/internal/runtimegateway) 与参考 Runtime 的 [dispatcher](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/apps/daemon/internal/dispatch) 都使用它们,因此无需同步第二套 payload schema。签发凭据和打开连接的 HTTP 路由见[机器连接 API](../../contracts/agents-api/zh/machine-api.md)。 托管和自托管 Runtime 使用同一协议。Harness 通过 [Harness adapter 契约](../../contracts/agents-api/zh/harness-onboarding.md)接入,该契约负责 Runtime registry 后的 Executor 和 Turn 生命周期义务。 +连接帧受 `proto.MaxFrameBytes` 的 4 MiB 上限约束。Core 会完整记录传输范围内的事件,不截断 payload。普通 journal 批次最多包含 64 个事件、合计 1 MiB;更大的单个事件独立存储。每个 Turn 仍限制为 65,536 个事件和 32 MiB,并在限制之外预留一条终态记录。journal 失败会使 Turn 失败并保留已提交的 Items。`execution journal failed` 日志记录 Project、Session、Turn ID、追踪上下文、阶段、事件类型、字节数、journal 位置以及有限的失败分类或 SQLSTATE;不记录事件 payload 或原始错误文本。 + ## 所有权与连接 {#ownership-and-connection} Core 拥有持久化的 Session、Turn、input 和 Environment 记录、调度与协调。Runtime 拥有原生 Executor、活动 Turn、传输状态和清理责任,直到完成结算。Sandbox Provider 拥有执行位置与外围计算资源。释放执行准入或关闭 Executor 不会删除、暂停或回收沙箱。daemon 不是隔离边界;参见 [Runtime 与外层隔离](concepts.md#runtime-and-outer-isolation)。 diff --git a/internal/agentdaemon/proto/envelope.go b/internal/agentdaemon/proto/envelope.go index da98ad3df..b827fb31d 100644 --- a/internal/agentdaemon/proto/envelope.go +++ b/internal/agentdaemon/proto/envelope.go @@ -27,6 +27,9 @@ import ( "fmt" ) +// MaxFrameBytes bounds a complete JSON frame on the Core–Runtime connection. +const MaxFrameBytes = 4 * 1024 * 1024 + // Envelope is the outer JSON frame. Payload is held as raw JSON so the // routing layer can dispatch by Type before paying a per-event decode. type Envelope struct { diff --git a/services/core/IMPLEMENTATION.md b/services/core/IMPLEMENTATION.md index 9231be3fa..352af38fb 100644 --- a/services/core/IMPLEMENTATION.md +++ b/services/core/IMPLEMENTATION.md @@ -163,7 +163,7 @@ Neutral tool and message observations go into the journal before `sessions.Proje ## Journal and live events -Execution observations are written to tenant-scoped `turn_events` in ordered, idempotent batches through `sessions.AppendTurnEvents` before they back recovery or publication, with daemon payloads intact. Flush at least every 100 ms while consuming events and before terminal persistence; the terminal outcome, its journal entry and native continuity commit together. `sessions` enforces the journal limits: 64 observations, 512 KiB per payload and 1 MiB per batch, and 65,536 observations and 32 MiB per Turn, with one extra entry reserved for the terminal outcome. Never infer a successful completion after a persistence error or stream overflow. +Execution observations are written to tenant-scoped `turn_events` in ordered, idempotent batches through `sessions.AppendTurnEvents` before they back recovery or publication, with daemon payloads intact. Flush at least every 100 ms while consuming events and before terminal persistence; the terminal outcome, its journal entry and native continuity commit together. `sessions` enforces the [journal limits](../../docs/runtime-protocol.md), including the reserved terminal entry. Never infer a successful completion after a persistence error or stream overflow. Live Session SSE reads `session_events` committed with the corresponding input, Item or lifecycle change under the Session lock, as immutable transition snapshots. The notification buffer keeps at most 256 events and 64 MiB per Session after each transaction (the `sessions` retention bounds, which `sessionpg` applies), and read batches at most 32 events or 1 MiB, each keeping a single oversized event. GET polls committed events every 100 ms from the committed high-water mark; a missing sequence position ends the stream with a safe error. Socket writes have a five-second deadline and hold no database connection. Rebuilding historical indexes emits no live events. diff --git a/services/core/internal/execution/delivery.go b/services/core/internal/execution/delivery.go index babd31a2a..11a0d3a19 100644 --- a/services/core/internal/execution/delivery.go +++ b/services/core/internal/execution/delivery.go @@ -73,7 +73,7 @@ func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, pe abort(peer, request.RunID) } }() - journal := &journal{writer: d.sessionExecution, tenant: tenantID, session: sessionID, turn: request.RunID, next: 1, + journal := &journal{ctx: ctx, writer: d.sessionExecution, tenant: tenantID, session: sessionID, turn: request.RunID, next: 1, observeSubagents: request.ObserveSubagentIdentities} defer func() { finishCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) diff --git a/services/core/internal/execution/journal.go b/services/core/internal/execution/journal.go index cf3a15438..45297eeca 100644 --- a/services/core/internal/execution/journal.go +++ b/services/core/internal/execution/journal.go @@ -10,6 +10,7 @@ import ( ) type journal struct { + ctx context.Context writer eventWriter tenant, session, turn string next int32 @@ -45,6 +46,7 @@ func recordCancellation(ctx context.Context, journal *journal, reply cancellatio func (j *journal) observe(ctx context.Context, env proto.Envelope) error { if err := j.enqueue(env); err != nil { + j.reportFailure("enqueue", env.Type, len(env.Payload), err) return err } if j.bytes > 768*1024 || len(j.batch) >= 64 { @@ -63,7 +65,7 @@ func (j *journal) enqueue(env proto.Envelope) error { default: return nil } - if len(env.Payload) > 512*1024 { + if len(env.Payload) > proto.MaxFrameBytes { return sessions.ErrEventLimit } j.batch = append(j.batch, sessions.ExecutionEvent{Kind: env.Type, Payload: env.Payload}) @@ -94,6 +96,7 @@ func (j *journal) flush(ctx context.Context) error { // An uncertain commit must retry the same batch even after more frames arrive. j.pendingCount = count if err := j.writer.AppendTurnEvents(ctx, j.tenant, j.session, j.turn, j.next, j.batch[:count]); err != nil { + j.reportFailure("flush", j.batch[0].Kind, size, err) return err } j.pendingCount = 0 @@ -115,6 +118,7 @@ func (j *journal) drain(upstream <-chan proto.Envelope, result *Result) error { return observedErr } if err := j.enqueue(env); err != nil { + j.reportFailure("drain", env.Type, len(env.Payload), err) observedErr = err } if err := result.mergeObservation(env); err != nil { diff --git a/services/core/internal/execution/journal_failure.go b/services/core/internal/execution/journal_failure.go new file mode 100644 index 000000000..c3dde0c1d --- /dev/null +++ b/services/core/internal/execution/journal_failure.go @@ -0,0 +1,63 @@ +package execution + +import ( + "context" + "errors" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" + "github.com/jackc/pgx/v5/pgconn" +) + +// Journal failures can contain SQL, model output and credentials. Record only +// bounded categories and sizes; the error text and event payload stay private. +func (j *journal) reportFailure(stage, kind string, size int, err error) { + ctx := j.ctx + if ctx == nil { + ctx = context.Background() + } + switch kind { + case proto.TypeDelta, proto.TypeOutputMessage, proto.TypeThinking, proto.TypeToolCall, proto.TypeCommandOutput, proto.TypeUsage, + proto.TypeError, proto.TypeDone, proto.TypePromptSteerAck, proto.TypeSubagentIdentity, proto.TypeSubagentLifecycle, proto.TypeSubagentTurn, proto.TypeSubagentItem, proto.TypeSubagentCoordination, "cancel_receipt": + default: + kind = "unknown" + } + reason, state := journalFailureCategory(err) + obslog.Ctx(ctx).Error("execution journal failed", "project_id", j.tenant, "session_id", j.session, "turn_id", j.turn, + "stage", stage, "event_kind", kind, "payload_bytes", size, "next_ordinal", j.next, + "pending_events", len(j.batch), "reason", reason, "sqlstate", state) +} + +func journalFailureCategory(err error) (string, string) { + switch { + case errors.Is(err, context.DeadlineExceeded): + return "deadline_exceeded", "" + case errors.Is(err, context.Canceled): + return "cancelled", "" + case errors.Is(err, sessions.ErrEventLimit): + return "event_limit", "" + case errors.Is(err, sessions.ErrInvalidInput): + return "invalid_event", "" + case errors.Is(err, sessions.ErrIdempotencyConflict): + return "idempotency_conflict", "" + case errors.Is(err, sessions.ErrTurnConflict): + return "turn_conflict", "" + } + var pg *pgconn.PgError + if errors.As(err, &pg) { + if len(pg.Code) == 5 { + valid := true + for _, c := range pg.Code { + if !(c >= '0' && c <= '9' || c >= 'A' && c <= 'Z') { + valid = false + } + } + if valid { + return "database", pg.Code + } + } + return "database", "" + } + return "unknown", "" +} diff --git a/services/core/internal/execution/journal_failure_test.go b/services/core/internal/execution/journal_failure_test.go new file mode 100644 index 000000000..4a05c415f --- /dev/null +++ b/services/core/internal/execution/journal_failure_test.go @@ -0,0 +1,49 @@ +package execution + +import ( + "context" + "errors" + "fmt" + "strings" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" + "github.com/jackc/pgx/v5/pgconn" +) + +func TestJournalFailureLogsSafeCategoriesAndCorrelation(t *testing.T) { + out := captureExecutionLogs(t) + j := journal{ctx: obslog.WithRequestID(reservationTrace(t.Context(), "reservation"), "request-id"), tenant: "project-id", session: "session-id", turn: "turn-id", next: 112} + for _, tc := range []struct { + err error + reason, state string + }{ + {fmt.Errorf("secret-canary: %w", sessions.ErrEventLimit), "event_limit", ""}, + {context.DeadlineExceeded, "deadline_exceeded", ""}, + {context.Canceled, "cancelled", ""}, + {sessions.ErrInvalidInput, "invalid_event", ""}, + {sessions.ErrIdempotencyConflict, "idempotency_conflict", ""}, + {sessions.ErrTurnConflict, "turn_conflict", ""}, + {&pgconn.PgError{Code: "22P05", Message: "secret-canary", Detail: "secret-canary"}, "database", "22P05"}, + {&pgconn.PgError{Code: "secret-canary"}, "database", ""}, + {errors.New("secret-canary"), "unknown", ""}, + } { + reason, state := journalFailureCategory(tc.err) + if reason != tc.reason || state != tc.state { + t.Fatalf("category %q %q", reason, state) + } + j.reportFailure("flush", proto.TypeToolCall, 536023, tc.err) + } + j.reportFailure("enqueue", "secret-canary", 1, errors.New("secret-canary")) + logs := out.String() + if strings.Contains(logs, "secret-canary") { + t.Fatal("raw diagnostic leaked") + } + for _, want := range []string{"project-id", "session-id", "turn-id", "request-id", "trace_id", "536023", "112", "event_limit", "22P05"} { + if !strings.Contains(logs, want) { + t.Fatalf("missing correlation/category %q", want) + } + } +} diff --git a/services/core/internal/execution/journal_test.go b/services/core/internal/execution/journal_test.go index 1955836aa..7cf495adb 100644 --- a/services/core/internal/execution/journal_test.go +++ b/services/core/internal/execution/journal_test.go @@ -159,3 +159,61 @@ func TestCancellationReceiptContinuitySurvivesFlushFailure(t *testing.T) { t.Fatal("receipt or partial text lost") } } + +// Transport-valid terminal snapshots must survive the same journal as their +// streamed output. The final snapshot is authoritative, not another delta. +func TestJournalAcceptsTransportSizedCommandSnapshot(t *testing.T) { + for _, size := range []int{528382, 2 * 1024 * 1024} { + t.Run(fmt.Sprint(size), func(t *testing.T) { + writer := &validatingWriter{} + j := journal{writer: writer, next: 1, turn: "c8b7fc8b-23e1-49f5-aaec-6d81398d218a"} + for _, p := range []proto.ToolCallPayload{ + {ID: "cmd", Stage: "before", Observation: &proto.ToolObservation{Kind: "command", Command: "fixture", Status: "in_progress"}}, + {ID: "cmd", Stage: "after", Observation: &proto.ToolObservation{Kind: "command", Command: "fixture", Status: "completed", Output: []byte(`"` + strings.Repeat("x", size) + `"`)}}, + } { + env, err := proto.NewEnvelope(proto.TypeToolCall, j.turn, p) + if err != nil { + t.Fatal(err) + } + if err = j.observe(t.Context(), env); err != nil { + t.Fatal(err) + } + } + if err := j.flush(t.Context()); err != nil { + t.Fatal(err) + } + if len(writer.events) != 2 || j.next != 3 { + t.Fatal("lost terminal snapshot") + } + }) + } +} + +type validatingWriter struct{ events []sessions.ExecutionEvent } + +func (w *validatingWriter) AppendTurnEvents(_ context.Context, _, _, turn string, first int32, events []sessions.ExecutionEvent) error { + if _, err := sessions.NewJournalBatch(turn, first, events); err != nil { + return err + } + w.events = append(w.events, events...) + return nil +} + +func TestJournalRejectsOversizedEventAndRetainsHistory(t *testing.T) { + writer := &validatingWriter{} + j := journal{writer: writer, next: 1, turn: "c8b7fc8b-23e1-49f5-aaec-6d81398d218a"} + before, _ := proto.NewEnvelope(proto.TypeDelta, j.turn, proto.DeltaPayload{Delta: "retained"}) + if err := j.observe(t.Context(), before); err != nil { + t.Fatal(err) + } + oversized, _ := proto.NewEnvelope(proto.TypeDelta, j.turn, proto.DeltaPayload{Delta: strings.Repeat("x", proto.MaxFrameBytes)}) + if err := j.observe(t.Context(), oversized); !errors.Is(err, sessions.ErrEventLimit) { + t.Fatal(err) + } + if err := j.flush(t.Context()); err != nil { + t.Fatal(err) + } + if len(writer.events) != 1 || string(writer.events[0].Payload) != string(before.Payload) { + t.Fatal("history changed after rejection") + } +} diff --git a/services/core/internal/runtimegateway/session.go b/services/core/internal/runtimegateway/session.go index a1b3f2389..8d4cc0035 100644 --- a/services/core/internal/runtimegateway/session.go +++ b/services/core/internal/runtimegateway/session.go @@ -41,7 +41,7 @@ var ( // ReadLimit caps a single inbound frame at 4 MiB. tool_call // results can be large but anything past this is almost certainly // a misbehaving daemon (or hostile input). - ReadLimit int64 = 4 * 1024 * 1024 + ReadLimit int64 = proto.MaxFrameBytes // CloseRuntimeDeleted is a custom WS close code (4001) sent when // a heartbeat discovers the runtime has been deleted. The daemon diff --git a/services/core/internal/sessions/journal.go b/services/core/internal/sessions/journal.go index a2af177fd..e083458c4 100644 --- a/services/core/internal/sessions/journal.go +++ b/services/core/internal/sessions/journal.go @@ -6,16 +6,18 @@ import ( "github.com/google/uuid" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/jsonobject" ) // A Turn's journal holds its execution observations in order. One batch holds -// at most 64 observations of at most 512 KiB each and 1 MiB together, and one -// Turn's journal at most 65,536 observations and 32 MiB. The terminal outcome -// is recorded beyond those limits, in the one entry they reserve for it. +// at most 64 observations and 1 MiB together, or one observation up to the +// transport frame ceiling. A Turn holds at most 65,536 observations and +// 32 MiB. The terminal outcome is recorded beyond those limits, in the one +// entry they reserve for it. const ( journalBatchEvents = 64 - journalPayloadBytes = 512 * 1024 + journalPayloadBytes = proto.MaxFrameBytes journalBatchBytes = 1024 * 1024 journalTurnEvents = 65536 journalTurnBytes = 32 * 1024 * 1024 @@ -51,7 +53,7 @@ func NewJournalBatch(turn string, first int32, events []ExecutionEvent) (Journal batch.events[i] = ExecutionEvent{Kind: event.Kind, Payload: payload} batch.bytes += int64(len(payload)) } - if batch.bytes > journalBatchBytes { + if len(batch.events) > 1 && batch.bytes > journalBatchBytes { return JournalBatch{}, ErrEventLimit } return batch, nil diff --git a/services/core/internal/sessions/journal_test.go b/services/core/internal/sessions/journal_test.go index 40c0e08fe..7b4213fb1 100644 --- a/services/core/internal/sessions/journal_test.go +++ b/services/core/internal/sessions/journal_test.go @@ -163,3 +163,25 @@ func TestExecutionOperationsValidationOrder(t *testing.T) { } } } + +func TestJournalSingleLargeObservationKeepsBatchAndTurnBudgets(t *testing.T) { + payload := json.RawMessage(`{"text":"` + strings.Repeat("x", journalPayloadBytes-len(`{"text":""}`)) + `"}`) + large := ExecutionEvent{Kind: "delta", Payload: payload} + batch, err := NewJournalBatch(testTurn, 1, []ExecutionEvent{large}) + if err != nil || batch.bytes != journalPayloadBytes { + t.Fatalf("maximum event: bytes=%d err=%v", batch.bytes, err) + } + if err := batch.admit(JournalTurn{Status: TurnInProgress, EventBytes: journalTurnBytes - batch.bytes}); err != nil { + t.Fatal(err) + } + if err := batch.admit(JournalTurn{Status: TurnInProgress, EventBytes: journalTurnBytes - batch.bytes + 1}); !errors.Is(err, ErrEventLimit) { + t.Fatal("turn budget bypass", err) + } + if _, err := NewJournalBatch(testTurn, 1, []ExecutionEvent{large, event()}); !errors.Is(err, ErrEventLimit) { + t.Fatal("oversized multi-event batch accepted", err) + } + large.Payload = append([]byte(" "), payload...) + if _, err := NewJournalBatch(testTurn, 1, []ExecutionEvent{large}); !errors.Is(err, ErrInvalidInput) { + t.Fatal("oversized event accepted", err) + } +} diff --git a/services/core/internal/store/command_output_test.go b/services/core/internal/store/command_output_test.go index 8ae0b4ba5..25a46bfff 100644 --- a/services/core/internal/store/command_output_test.go +++ b/services/core/internal/store/command_output_test.go @@ -4,8 +4,11 @@ import ( "context" "encoding/json" "errors" + "fmt" "reflect" + "strings" "testing" + "time" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" @@ -150,3 +153,86 @@ func TestExecutionJournalsCommandOutputBeforeCancellation(t *testing.T) { t.Fatalf("journal/cancellation lost partial output: %+v %v", page, err) } } + +func TestExecutionSecondTurnCommandSnapshotAndStorageFailure(t *testing.T) { + for _, failStorage := range []bool{false, true} { + t.Run(fmt.Sprint(failStorage), func(t *testing.T) { + h := newDispatchHarness(t) + ctx := t.Context() + first := h.message("first", "first input") + result := h.run(ctx, first.TurnID) + h.read(testExecutionRequest) + h.write(first.TurnID, proto.TypeDone, proto.DonePayload{Content: "retained first reply"}) + h.finished(result, sessions.TurnCompleted) + history, err := sessionReads(h.db.pool).ListItems(ctx, h.tenant, h.session.ID, "", 100, true) + if err != nil { + t.Fatal(err) + } + + second := h.message("second", "second input") + result = h.run(ctx, second.TurnID) + h.read(testExecutionRequest) + h.write(second.TurnID, proto.TypeToolCall, proto.ToolCallPayload{ID: "cmd", Stage: "before", Observation: &proto.ToolObservation{Kind: "command", Command: "fixture", Status: "in_progress"}}) + h.write(second.TurnID, proto.TypeCommandOutput, proto.CommandOutputPayload{ID: "cmd", Delta: "retained stream"}) + deadline := time.Now().Add(5 * time.Second) + for { + events, err := h.s.ListTurnEvents(ctx, h.tenant, h.session.ID, second.TurnID, 0, 100) + if err != nil { + t.Fatal(err) + } + if len(events) >= 2 { + break + } + if time.Now().After(deadline) { + t.Fatal("stream not persisted") + } + time.Sleep(10 * time.Millisecond) + } + // A full snapshot exceeds the old journal limit. A NUL is instead a + // deterministic PostgreSQL JSONB failure, not a native engine failure. + output := strings.Repeat("x", 2*1024*1024) + expectedStatus, expectedOutput := "completed", output + turnStatus := sessions.TurnCompleted + if failStorage { + output = "\x00" + expectedStatus, expectedOutput = "incomplete", "retained stream" + turnStatus = sessions.TurnFailed + } + raw, _ := json.Marshal(output) + h.write(second.TurnID, proto.TypeToolCall, proto.ToolCallPayload{ID: "cmd", Stage: "after", Observation: &proto.ToolObservation{Kind: "command", Command: "fixture", Status: "completed", Output: raw}}) + h.write(second.TurnID, proto.TypeDone, proto.DonePayload{Content: "second reply"}) + turn := h.finished(result, turnStatus) + if failStorage { + var outcome struct { + ErrorCode string `json:"error_code"` + } + if json.Unmarshal(turn.Outcome, &outcome) != nil || outcome.ErrorCode != "event_persistence_failed" { + t.Fatal("wrong failure outcome") + } + } + page, err := sessionReads(h.db.pool).ListItems(ctx, h.tenant, h.session.ID, "", 100, true) + if err != nil || len(page.Items) < len(history.Items)+2 { + t.Fatal("missing Items", err) + } + if !reflect.DeepEqual(page.Items[:len(history.Items)], history.Items) { + t.Fatal("prior Turn Items changed") + } + found := false + for _, item := range page.Items { + if item.TurnID == second.TurnID && item.Type == "command_execution" { + found = true + if item.Status != expectedStatus || item.Output != expectedOutput { + t.Fatal("command snapshot/terminal state mismatch") + } + } + } + if !found { + t.Fatal("missing command Item") + } + events, err := h.s.ListTurnEvents(ctx, h.tenant, h.session.ID, second.TurnID, 0, 100) + if err != nil || len(events) == 0 || events[len(events)-1].Kind != "execution_"+turnStatus { + t.Fatal("journal terminal mismatch", err) + } + }) + } +}