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 contracts/agents-api/sessions-events.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ A Session stays usable after a Turn fails: new input starts a new Turn. Later re
- **Cancellation.** A queued Turn is cancelled without a live Runtime. A running Turn is cancelled when the Runtime confirms it; completion can win that race. The Turn has stopped when it reads `cancelled`, not when the request returns. A cancellation on an idle Session with no pending input is accepted and has no effect; while an input reservation is pending, it returns 409.
- **Function results.** `turn_id`, `call_id` and `success` are required; `output` and `error` are optional and nullable ([content rules](./message-content.md#function-results)). An identical repeated result returns 202 without another application or event. The result Item appears when the harness applies the result; a result that cancellation prevents from being applied stays stored but produces no Item.
- **Queueing.** A queued Turn starts when a Runtime that supports the Session's harness and configuration is connected and one of Core's [`core.execution_concurrency`](../../docs/configuration.md#settings) work slots is free. A Session stays bound to the Runtime that first ran it.
- **Execution availability.** A service without execution returns 503 `execution_unavailable`, and a Worker that loses execution ownership returns 503. A Session created without a model provider rejects new messages with 400 `model_provider_required` ([model execution](./model-execution.md)).
- **Execution availability.** A service without execution returns 503 `execution_unavailable`, and a Worker that loses execution ownership returns 503.

### Sessions with an Environment

Expand Down
26 changes: 1 addition & 25 deletions contracts/agents-api/v1/items_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,15 +64,6 @@ func TestItemWireFieldsAreExplicitlyNull(t *testing.T) {
expectField(t, got, "output", test.output)
expectField(t, got, "error", test.error)
expectField(t, got, "phase", "")
// Stored payloads keep the original field presence.
stored, err := item.MarshalStored()
if err != nil {
t.Fatal(err)
}
var original, persisted map[string]json.RawMessage
if json.Unmarshal([]byte(test.raw), &original) != nil || json.Unmarshal(stored, &persisted) != nil || len(original) != len(persisted) {
t.Fatalf("stored payload changed: %s", stored)
}
})
}

Expand All @@ -84,14 +75,6 @@ func TestItemWireFieldsAreExplicitlyNull(t *testing.T) {
reasoning := fields(t, Item{ID: "rs", TurnID: "turn", Type: "reasoning"})
expectField(t, reasoning, "status", "null")
expectField(t, reasoning, "summary", "[]")

stored, err := user.MarshalStored()
if err != nil {
t.Fatal(err)
}
if string(stored) != `{"id":"user","turn_id":"turn","type":"message","status":"completed","role":"user","content":[{"type":"input_text","text":"question"}]}` {
t.Fatalf("stored message encoding changed: %s", stored)
}
}

func TestItemEventsCarryNullableOutputIndex(t *testing.T) {
Expand Down Expand Up @@ -156,7 +139,7 @@ func TestReasoningResponsesCarryBothKeys(t *testing.T) {
expectField(t, fields(t, stored["agent"]), "reasoning", `{"effort":"low"}`)
}

func TestStoredSearchItemRoundTripPreservesPayload(t *testing.T) {
func TestSearchItemWireAction(t *testing.T) {
for _, raw := range []string{
`{"id":"search","turn_id":"turn","type":"web_search_call","status":"completed","action":{"type":"search","query":"reference"}}`,
`{"id":"search","turn_id":"turn","type":"web_search_call","status":"completed","action":{"type":"search"}}`,
Expand All @@ -166,13 +149,6 @@ func TestStoredSearchItemRoundTripPreservesPayload(t *testing.T) {
if err := json.Unmarshal([]byte(raw), &item); err != nil {
t.Fatal(err)
}
stored, err := item.MarshalStored()
if err != nil {
t.Fatal(err)
}
if string(stored) != raw {
t.Fatalf("stored replay changed: %s, want %s", stored, raw)
}
wire := fields(t, item)
if item.Action == nil {
expectField(t, wire, "action", "null")
Expand Down
28 changes: 5 additions & 23 deletions contracts/agents-api/v1/subagent_items.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,10 @@ import (
"errors"
)

// MarshalJSON renders the wire shape. Messages always carry content and a
// nullable phase, function results a nullable output and error, and web search
// a nullable action. Other variants use the stored encoding.
// MarshalJSON renders the wire shape, which is also the stored encoding.
// Messages always carry content and a nullable phase, function results a
// nullable output and error, and web search a nullable action. Coordination
// and reasoning variants keep their required fields and nulls.
func (i Item) MarshalJSON() ([]byte, error) {
type wire Item
switch i.Type {
Expand Down Expand Up @@ -36,23 +37,6 @@ func (i Item) MarshalJSON() ([]byte, error) {
Output any `json:"output"`
Error any `json:"error"`
}{wire(i), i.Output, i.Error})
}
return i.MarshalStored()
}

// MarshalStored encodes a persisted Item payload. It keeps the encoding used
// before the wire nulls above, so stored payloads and the byte comparison of
// replayed child Items do not change. Coordination variants keep their
// required fields and nulls in both forms.
func (i Item) MarshalStored() ([]byte, error) {
switch i.Type {
case "web_search_call":
type wire Item
type storedAction WebSearchAction
return json.Marshal(struct {
wire
Action *storedAction `json:"action,omitempty"`
}{wire(i), (*storedAction)(i.Action)})
case "create_subagent_call", "send_subagent_input_call", "agent_message":
content, err := coordinationContent(i.Content)
if err != nil {
Expand Down Expand Up @@ -84,10 +68,8 @@ func (i Item) MarshalStored() ([]byte, error) {
summary = []SummaryText{}
}
return json.Marshal(ReasoningItem{ID: i.ID, TurnID: i.TurnID, Type: i.Type, Status: status, Summary: summary})
default:
type wire Item
return json.Marshal(wire(i))
}
return json.Marshal(wire(i))
}

func coordinationContent(parts []ItemContent) ([]AgentContent, error) {
Expand Down
4 changes: 2 additions & 2 deletions contracts/agents-api/zh/sessions-events.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
---
title: "会话、事件和历史"
source: contracts/agents-api/sessions-events.md
source_hash: d5d0928665f38105167e592f9561c3d3852a13fc11b976f15c1bc4cc2ca9cb14
source_hash: c6141811b426fb0dfb95105b27cc381af5f22e3a7923a171f113f228ae0a5b33
---

本契约涵盖会话(Session)内部发生的事情:发送输入、实时事件流,以及读取轮次(Turn)、条目(Item)和使用量的持久化历史。会话资源本身(创建配置、重试标识、更新、列出和删除)见 [Core 线协议行为](wire-semantics.md)。消息和函数结果内容见[消息内容](message-content.md)。[Agents API 指南](../../../docs/zh/api/public-agent-api.md)展示了使用 SDK 和 HTTP 的调用方式。
Expand Down Expand Up @@ -41,7 +41,7 @@ Turn 失败后会话仍可使用:新输入会启动一个新 Turn。后来预
- **取消。** 排队的 Turn 无需活动 Runtime 即可取消。正在运行的 Turn 只有在 Runtime 确认后才会取消;完成操作可能赢得该竞争。读取到 `cancelled` 时才表示该 Turn 已停止,而不是请求返回时。在没有待处理输入的情况下,对空闲会话执行的取消会被接受且不产生任何效果;而在输入预留待处理期间,取消会返回 409。
- **函数结果。** `turn_id`、`call_id` 和 `success` 为必填项;`output` 和 `error` 为可选项且可为空([内容规则](message-content.md#function-results))。重复提交完全相同的结果会返回 202,不会再次应用或发出事件。harness 应用结果时才会出现结果 Item;如果取消操作导致结果无法应用,结果仍会存储,但不会产生 Item。
- **排队。** 当一个已连接且支持该会话 harness 和配置的 Runtime 接入,并且 Core 的 [`core.execution_concurrency`](../../../docs/zh/configuration.md#settings) 工作槽位有一个空闲时,排队的 Turn 才会启动。会话始终绑定到首次运行它的 Runtime。
- **执行可用性。** 不具备执行能力的服务会返回 503 `execution_unavailable`,失去执行所有权的 Worker 会返回 503。创建时未指定模型提供方的会话会拒绝新消息并返回 400 `model_provider_required`([模型执行](model-execution.md))。
- **执行可用性。** 不具备执行能力的服务会返回 503 `execution_unavailable`,失去执行所有权的 Worker 会返回 503。

### 包含 Environment 的会话 {#sessions-with-an-environment}

Expand Down
2 changes: 1 addition & 1 deletion services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -177,7 +177,7 @@ A streaming Session creation reuses atomic input admission and the live event lo

Public Items read a projection updated in the same Session transaction as admitted messages and journal batches. Item IDs derive from the Turn and source identity, and the first-observation timestamp and tie breakers never change when content or status does. Each new Item's Session position is allocated under the Session lock, preserving observation order for equal timestamps, and each Turn allocates its own zero-based `output_index`, which inputs do not consume; updates and retries keep both.

Item merging never mutates the incoming observation or the previous snapshot: public text delta events read the original fragment after merging, while the Item keeps the accumulated text, and the content slice is copied before its text pointer is replaced. A first observation without its own fragment carries its unchanged text in one delta. Wire-only explicit nulls come from response marshalling, while stored Item payloads keep their original encoding through `Item.MarshalStored`, so replayed child Items compare equal. Structured tool JSON is kept without float conversion, and an unfinished call never becomes a successful result. When a Turn ends, `sessions.EndTurn` makes its unfinished Items incomplete with their partial content and reports them in Session position order before the Turn's event and its settled Session activity; `sessionpg` gives them one shared settlement time. Function results are Session input Items: they emit `item.added` with a null `output_index` and never `item.done`, whose upstream union allows only agent output, and their public output and error come from the saved submission.
Item merging never mutates the incoming observation or the previous snapshot: public text delta events read the original fragment after merging, while the Item keeps the accumulated text, and the content slice is copied before its text pointer is replaced. A first observation without its own fragment carries its unchanged text in one delta. Stored Item payloads use the wire encoding, explicit nulls included, so a replayed child Item compares equal to its stored payload. Structured tool JSON is kept without float conversion, and an unfinished call never becomes a successful result. When a Turn ends, `sessions.EndTurn` makes its unfinished Items incomplete with their partial content and reports them in Session position order before the Turn's event and its settled Session activity; `sessionpg` gives them one shared settlement time. Function results are Session input Items: they emit `item.added` with a null `output_index` and never `item.done`, whose upstream union allows only agent output, and their public output and error come from the saved submission.

## Worker ownership

Expand Down
2 changes: 0 additions & 2 deletions services/core/internal/api/errors_sessions.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,6 @@ func writeSessionsError(w http.ResponseWriter, r *http.Request, err error) {
writeError(w, http.StatusUnauthorized, "installation_authorization_invalid", sessions.ErrInstallationAuthorization.Error())
case errors.Is(err, sessions.ErrExecutorCredentialExists):
writeError(w, http.StatusConflict, "executor_credential_exists", "This executor key ID already exists. Explicitly rotate it to replace the secret.")
case errors.Is(err, execution.ErrModelProviderRequired):
writeError(w, http.StatusBadRequest, "model_provider_required", "This Session was created without a model provider and cannot run. Create a new Session with x_agents_core.model_provider or an Agent that has one saved.")
case errors.Is(err, sessions.ErrHostedEnvironmentFailed):
// Observed official status, type, code, null param and message.
writeError(w, http.StatusConflict, "conflict_error", "the hosted environment failed to provision")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ func TestNativeClassificationPostgresRoundTripAndPublicPrivacy(t *testing.T) {
t.Fatal(err)
}
session := created.Session
receipt := submitMessage(t, pool, tenant, session.ID, "input", json.RawMessage(`{"text":"test"}`))
receipt := submitMessage(t, pool, tenant, session.ID, "input", "test")
transitionTurn(t, pool, tenant, session.ID, receipt.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnQueued, Status: sessions.TurnInProgress})
status := 503
result := execution.Result{ErrorCode: "engine_failed", Error: "Bearer secret-canary https://private.example/key", EngineErrorCode: code, EngineHTTPStatus: &status, Done: proto.DonePayload{Usage: proto.Usage{InputTokens: 7, OutputTokens: 3}, Metadata: map[string]any{proto.DoneMetaAgentSessionID: "native-secret-canary"}}}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"strings"
"testing"

v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/db/sqlc"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/identity"
"github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit"
Expand Down Expand Up @@ -73,13 +74,15 @@ func databaseSessionReads(pool *pgxpool.Pool) func(*Dependencies, *testFakes) {
}
}

// submitMessage admits one message input through the Session service on pool.
func submitMessage(t *testing.T, pool *pgxpool.Pool, tenant, session, key string, payload json.RawMessage) sessions.InputReceipt {
// submitMessage admits one public text message through the Session service on
// pool.
func submitMessage(t *testing.T, pool *pgxpool.Pool, tenant, session, key, text string) sessions.InputReceipt {
t.Helper()
service, err := sessions.NewService(sessionpg.New(pgunit.NewPool(pool), nil), nil)
if err != nil {
t.Fatal(err)
}
payload, _ := json.Marshal(v1.SessionInput{Type: "agent.session.input.message", Input: []v1.InputMessage{{Role: "user", Content: []v1.InputContent{{Type: "input_text", Text: &text}}}}})
receipts, err := service.SubmitInputs(t.Context(), tenant, session, key, []sessions.Input{{Kind: "message", Payload: payload}})
if err != nil {
t.Fatal(err)
Expand Down
2 changes: 1 addition & 1 deletion services/core/internal/api/session_diagnostics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ func TestDiagnosticsCoreHandlerDatabaseBoundary(t *testing.T) {
t.Fatal(err)
}
session := created.Session
receipt := submitMessage(t, pool, tenant, session.ID, "input", json.RawMessage(`{"text":"input-secret-canary"}`))
receipt := submitMessage(t, pool, tenant, session.ID, "input", "input-secret-canary")
transitionTurn(t, pool, tenant, session.ID, receipt.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnQueued, Status: sessions.TurnInProgress})
transitionTurn(t, pool, tenant, session.ID, receipt.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnInProgress, Status: sessions.TurnFailed, Outcome: json.RawMessage(`{"error_code":"device_disconnected","error":"Bearer raw-secret-canary https://private.example/key","done":{"native_id":"secret-native-canary"}}`)})
base := adminSessionsPath + session.ID
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ func TestArchiveWaitingCleanupReceiptBarrier(t *testing.T) {
if err != nil {
t.Fatal(err)
}
inputs, err := service.SubmitInputs(t.Context(), project.TenantID, session.ID, "start", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"run"}`)}})
inputs, err := service.SubmitInputs(t.Context(), project.TenantID, session.ID, "start", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"run"}]}]}`)}})
if err != nil {
t.Fatal(err)
}
Expand Down
7 changes: 0 additions & 7 deletions services/core/internal/execution/environment_admission.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,13 +87,6 @@ func (w *Worker) submitEnvironmentInputs(ctx context.Context, session sessions.S
// Neither kind creates a Turn. The Session lock preserves target and retry identity.
return w.admitInputs(ctx, session.TenantID, session.ID, key, inputs)
}
// Messages start work. A Session from before deployment defaults moved into
// Core may have no frozen provider; reject it here instead of queueing work
// its harness cannot run. Cancellation and results above stay available.
var snapshot Snapshot
if json.Unmarshal(session.Configuration, &snapshot) != nil || !snapshot.ModelProviderConfigured {
return nil, ErrModelProviderRequired
}
changed, unsubscribe := w.dispatcher.notifications.subscribe(session.TenantID, session.ID)
defer unsubscribe()
reserve, cancel := context.WithTimeout(ctx, 5*time.Second)
Expand Down
8 changes: 1 addition & 7 deletions services/core/internal/execution/message_input.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,17 +10,11 @@ import (
)

func messageInput(raw json.RawMessage) (proto.MessageInput, error) {
var input struct {
Text *string `json:"text"`
Input []v1.InputMessage `json:"input"`
}
var input v1.SessionInput
if json.Unmarshal(raw, &input) != nil {
return nil, sessions.ErrInvalidInput
}
var messages proto.MessageInput
if len(input.Input) == 0 && input.Text != nil {
messages = proto.TextInput(*input.Text)
}
for _, message := range input.Input {
if message.Role != "user" {
return nil, sessions.ErrInvalidInput
Expand Down
6 changes: 1 addition & 5 deletions services/core/internal/execution/message_support_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ func TestMessageImageQualificationIsOperationSpecific(t *testing.T) {
}
// Message validation applies even when no function-result validator exists.
raw, _ := json.Marshal(map[string]any{"input": []any{map[string]any{"role": "user", "content": input[0].Content}}})
batch := []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"valid first"}`)}, {Kind: "message", Payload: raw}}
batch := []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"valid first"}]}]}`)}, {Kind: "message", Payload: raw}}
if err := validateProfileInputs(enginetest.Profile(nil), "none", batch); !errors.Is(err, sessions.ErrInvalidInput) {
t.Fatal("image escaped profile validation", err)
}
Expand Down Expand Up @@ -73,10 +73,6 @@ func TestWhitespaceOnlyTextQualificationUsesEngineProfiles(t *testing.T) {
}
}
claude, _ := (engine.Catalog{}).Lookup("claude_sdk")
// Legacy text payloads use the same rule.
if err := validateProfileInputs(claude, "none", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":" \t"}`)}}); !errors.Is(err, ErrWhitespaceOnlyText) {
t.Fatal(err)
}
mixed, _ := json.Marshal(map[string]any{"input": []any{map[string]any{"role": "user", "content": []any{
map[string]any{"type": "input_text", "text": " "}, map[string]any{"type": "input_text", "text": "text"}}},
map[string]any{"role": "user", "content": []any{map[string]any{"type": "input_text", "text": " "}, map[string]any{"type": "input_image", "image_url": url}}}}})
Expand Down
7 changes: 0 additions & 7 deletions services/core/internal/execution/model_execution_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,6 @@ func TestSessionModelExecutionNeverFallsBack(t *testing.T) {
if _, err := d.executionRequest(t.Context(), session, Snapshot{ModelProviderConfigured: true}, runtimedevice.KindCapabilities{}, sessions.ExecutionBinding{}); !errors.Is(err, sessions.ErrNotFound) {
t.Fatal("missing Session credentials fell back", err)
}
// Hosted and self-hosted Runtimes have no model configuration of their own.
for _, environment := range []string{"openai_hosted", "self_hosted"} {
snapshot := Snapshot{Environment: &v1.Environment{Type: environment}}
if _, err := d.executionRequest(t.Context(), sessions.Session{Engine: "codex"}, snapshot, runtimedevice.KindCapabilities{}, sessions.ExecutionBinding{}); !errors.Is(err, ErrModelProviderRequired) {
t.Fatal("provider-free Session dispatched", environment, err)
}
}
// A none device without a frozen provider uses its own provider environment:
// Core sends only the Agent's model and instructions.
instructions := "Keep this instruction."
Expand Down
Loading
Loading