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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apps/daemon/internal/transport/ws.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions docs/runtime-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
4 changes: 3 additions & 1 deletion docs/zh/runtime-protocol.md
Original file line number Diff line number Diff line change
@@ -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)。
Expand Down
3 changes: 3 additions & 0 deletions internal/agentdaemon/proto/envelope.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/execution/delivery.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
6 changes: 5 additions & 1 deletion services/core/internal/execution/journal.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
)

type journal struct {
ctx context.Context
writer eventWriter
tenant, session, turn string
next int32
Expand Down Expand Up @@ -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 {
Expand All @@ -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})
Expand Down Expand Up @@ -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
Expand All @@ -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 {
Expand Down
63 changes: 63 additions & 0 deletions services/core/internal/execution/journal_failure.go
Original file line number Diff line number Diff line change
@@ -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", ""
}
49 changes: 49 additions & 0 deletions services/core/internal/execution/journal_failure_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
}
58 changes: 58 additions & 0 deletions services/core/internal/execution/journal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}
2 changes: 1 addition & 1 deletion services/core/internal/runtimegateway/session.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 7 additions & 5 deletions services/core/internal/sessions/journal.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
22 changes: 22 additions & 0 deletions services/core/internal/sessions/journal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
Loading
Loading