From c8b303b2c357468ecdf159520316a683a8a740a1 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 17:24:59 +0000 Subject: [PATCH] Delete stored compatibility encodings Session message inputs are stored only as the public agent.session.input.message event, so the {"text"} branches in Item projection and Runtime dispatch go, and the recovery fixture and tests submit the public shape. Items are stored with their wire encoding: MarshalStored folds into MarshalJSON, and replayed child Items compare against the same bytes. Sessions always freeze a model provider at creation, so the execution guards for Sessions stored without one, ErrModelProviderRequired and its handlers are deleted. --- contracts/agents-api/sessions-events.md | 2 +- contracts/agents-api/v1/items_test.go | 26 +------- contracts/agents-api/v1/subagent_items.go | 28 ++------- contracts/agents-api/zh/sessions-events.md | 4 +- services/core/IMPLEMENTATION.md | 2 +- services/core/internal/api/errors_sessions.go | 2 - .../native_classification_integration_test.go | 2 +- .../session_diagnostics_public_compat_test.go | 7 ++- .../internal/api/session_diagnostics_test.go | 2 +- .../archive_cancellation_cleanup_test.go | 2 +- .../execution/environment_admission.go | 7 --- .../core/internal/execution/message_input.go | 8 +-- .../execution/message_support_test.go | 6 +- .../execution/model_execution_test.go | 7 --- .../internal/execution/prepared_dispatch.go | 5 -- services/core/internal/execution/request.go | 9 --- services/core/internal/execution/worker.go | 3 - .../internal/execution/worker_schedule.go | 3 - services/core/internal/items/changes_test.go | 2 +- services/core/internal/items/inputs.go | 6 +- .../sessionpg/execution_inputs_test.go | 2 +- .../postgres/sessionpg/inputs_test.go | 11 ++-- .../persistence/postgres/sessionpg/items.go | 2 +- .../core/internal/sessions/creation_test.go | 2 +- .../sessions/environment_inputs_test.go | 2 +- .../core/internal/sessions/inputs_test.go | 16 +++-- services/core/internal/sessions/subagents.go | 2 +- services/core/tests/fixtures/main.go | 2 +- .../admin_session_archive_race_test.go | 3 +- .../integration/admin_session_archive_test.go | 4 +- .../admin_session_archive_worker_http_test.go | 2 +- .../integration/archive_cancellation_test.go | 2 +- .../integration/claude_execution_test.go | 2 +- .../tests/integration/command_output_test.go | 2 +- .../creation_stream_settlement_public_test.go | 2 +- .../deployment_model_providers_http_test.go | 60 ------------------- .../core/tests/integration/dispatch_test.go | 3 +- .../integration/environment_admission_test.go | 2 +- .../environment_directory_active_test.go | 2 +- .../environment_expiry_dispatch_test.go | 5 +- .../environment_expiry_worker_test.go | 2 +- .../environment_initial_input_test.go | 8 ++- .../environment_initial_public_test.go | 2 +- .../environment_worker_helpers_test.go | 2 +- .../tests/integration/environments_test.go | 2 +- .../core/tests/integration/execution_test.go | 2 +- .../function_input_execution_test.go | 2 +- .../function_inputs_public_test.go | 4 +- .../tests/integration/function_inputs_test.go | 6 +- .../integration/function_item_events_test.go | 15 +---- .../integration/function_state_public_test.go | 2 +- ...sted_initialization_failure_public_test.go | 4 +- .../input_conflicts_public_test.go | 4 +- .../core/tests/integration/inputs_test.go | 11 +++- .../core/tests/integration/item_order_test.go | 6 +- .../core/tests/integration/item_reads_test.go | 8 +-- .../integration/local_artifact_export_test.go | 3 +- .../local_environment_file_write_test.go | 4 +- .../local_environment_worker_test.go | 3 +- .../message_image_admission_test.go | 2 +- .../integration/prepared_dispatch_test.go | 2 +- .../integration/public_execution_test.go | 2 +- .../runtime_input_admission_test.go | 2 +- .../tests/integration/runtime_pending_test.go | 2 +- .../runtime_wake_hint_integration_test.go | 3 +- .../runtime_worker_recovery_test.go | 5 +- .../sandbox_deployment_switch_test.go | 2 +- .../self_hosted_cancel_public_test.go | 4 +- .../session_deletion_lifecycle_public_test.go | 6 +- .../integration/session_deletion_test.go | 5 +- .../tests/integration/session_events_test.go | 6 +- .../integration/session_metadata_test.go | 2 +- .../integration/stream_authority_http_test.go | 2 +- .../token_usage_integration_test.go | 6 +- .../tests/integration/turn_events_test.go | 2 +- .../worker_preparation_failure_test.go | 11 ++-- .../tests/integration/worker_wakeup_test.go | 5 +- 77 files changed, 133 insertions(+), 287 deletions(-) diff --git a/contracts/agents-api/sessions-events.md b/contracts/agents-api/sessions-events.md index 1ce4367df..7615b7384 100644 --- a/contracts/agents-api/sessions-events.md +++ b/contracts/agents-api/sessions-events.md @@ -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 diff --git a/contracts/agents-api/v1/items_test.go b/contracts/agents-api/v1/items_test.go index 5e5c154e9..a507c1319 100644 --- a/contracts/agents-api/v1/items_test.go +++ b/contracts/agents-api/v1/items_test.go @@ -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) - } }) } @@ -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) { @@ -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"}}`, @@ -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") diff --git a/contracts/agents-api/v1/subagent_items.go b/contracts/agents-api/v1/subagent_items.go index b9c4a1533..0c21e078b 100644 --- a/contracts/agents-api/v1/subagent_items.go +++ b/contracts/agents-api/v1/subagent_items.go @@ -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 { @@ -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 { @@ -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) { diff --git a/contracts/agents-api/zh/sessions-events.md b/contracts/agents-api/zh/sessions-events.md index 519727458..cb7230f99 100644 --- a/contracts/agents-api/zh/sessions-events.md +++ b/contracts/agents-api/zh/sessions-events.md @@ -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 的调用方式。 @@ -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} diff --git a/services/core/IMPLEMENTATION.md b/services/core/IMPLEMENTATION.md index e80ea729c..47031e10a 100644 --- a/services/core/IMPLEMENTATION.md +++ b/services/core/IMPLEMENTATION.md @@ -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 diff --git a/services/core/internal/api/errors_sessions.go b/services/core/internal/api/errors_sessions.go index 8806156ee..64ed5b977 100644 --- a/services/core/internal/api/errors_sessions.go +++ b/services/core/internal/api/errors_sessions.go @@ -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") diff --git a/services/core/internal/api/native_classification_integration_test.go b/services/core/internal/api/native_classification_integration_test.go index 47a621d6a..0a13c9e78 100644 --- a/services/core/internal/api/native_classification_integration_test.go +++ b/services/core/internal/api/native_classification_integration_test.go @@ -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"}}} diff --git a/services/core/internal/api/session_diagnostics_public_compat_test.go b/services/core/internal/api/session_diagnostics_public_compat_test.go index bb5fa50eb..4394dd0b2 100644 --- a/services/core/internal/api/session_diagnostics_public_compat_test.go +++ b/services/core/internal/api/session_diagnostics_public_compat_test.go @@ -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" @@ -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) diff --git a/services/core/internal/api/session_diagnostics_test.go b/services/core/internal/api/session_diagnostics_test.go index a9365ba25..21ad01899 100644 --- a/services/core/internal/api/session_diagnostics_test.go +++ b/services/core/internal/api/session_diagnostics_test.go @@ -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 diff --git a/services/core/internal/execution/archive_cancellation_cleanup_test.go b/services/core/internal/execution/archive_cancellation_cleanup_test.go index a0cde3cd0..95353a660 100644 --- a/services/core/internal/execution/archive_cancellation_cleanup_test.go +++ b/services/core/internal/execution/archive_cancellation_cleanup_test.go @@ -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) } diff --git a/services/core/internal/execution/environment_admission.go b/services/core/internal/execution/environment_admission.go index d078186e7..52d1c3b05 100644 --- a/services/core/internal/execution/environment_admission.go +++ b/services/core/internal/execution/environment_admission.go @@ -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) diff --git a/services/core/internal/execution/message_input.go b/services/core/internal/execution/message_input.go index 328567f6f..c0dd9e8db 100644 --- a/services/core/internal/execution/message_input.go +++ b/services/core/internal/execution/message_input.go @@ -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 diff --git a/services/core/internal/execution/message_support_test.go b/services/core/internal/execution/message_support_test.go index 773893fe6..edccfaa8c 100644 --- a/services/core/internal/execution/message_support_test.go +++ b/services/core/internal/execution/message_support_test.go @@ -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) } @@ -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}}}}}) diff --git a/services/core/internal/execution/model_execution_test.go b/services/core/internal/execution/model_execution_test.go index 627a3161f..d4964e10b 100644 --- a/services/core/internal/execution/model_execution_test.go +++ b/services/core/internal/execution/model_execution_test.go @@ -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." diff --git a/services/core/internal/execution/prepared_dispatch.go b/services/core/internal/execution/prepared_dispatch.go index 751578383..fa2b11b0b 100644 --- a/services/core/internal/execution/prepared_dispatch.go +++ b/services/core/internal/execution/prepared_dispatch.go @@ -7,7 +7,6 @@ import ( "strings" "time" - v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -42,10 +41,6 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t if json.Unmarshal(session.Configuration, &snapshot) != nil || strings.TrimSpace(snapshot.Agent.Model) == "" { return run, sessions.ErrInvalidInput } - if !snapshot.ModelProviderConfigured && snapshot.Environment != nil && v1.ModelProviderRequired(snapshot.Environment.Type) { - // Reserved before providers were required; the caller settles it as failed. - return run, ErrModelProviderRequired - } bound, err := d.SessionsReader.GetSessionExecutionBinding(ctx, tenantID, sessionID) if err != nil { return run, err diff --git a/services/core/internal/execution/request.go b/services/core/internal/execution/request.go index 99b06d1f1..621bc16ad 100644 --- a/services/core/internal/execution/request.go +++ b/services/core/internal/execution/request.go @@ -5,17 +5,12 @@ import ( "encoding/json" "errors" - v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/vaults" ) -// ErrModelProviderRequired reports a hosted or self-hosted Session that has no -// frozen model provider and therefore cannot run. -var ErrModelProviderRequired = errors.New("the Session has no model provider") - func (d *Dispatcher) executionRequest(ctx context.Context, session sessions.Session, snapshot Snapshot, caps runtimedevice.KindCapabilities, bound sessions.ExecutionBinding) (proto.PromptRequestPayload, error) { recoverNativeSession := bound.HasStartedTurn && bound.NativeSessionID == "" if recoverNativeSession && !caps.NativeSessionRecovery { @@ -31,10 +26,6 @@ func (d *Dispatcher) executionRequest(ctx context.Context, session sessions.Sess if err != nil { return proto.PromptRequestPayload{}, err } - } else if snapshot.Environment != nil && v1.ModelProviderRequired(snapshot.Environment.Type) { - // Require the frozen bundle before dispatch so the harness cannot - // select an implicit provider endpoint. - return proto.PromptRequestPayload{}, ErrModelProviderRequired } options["model"], options["system_prompt"] = snapshot.Agent.Model, snapshot.Agent.Instructions if snapshot.Agent.XAgentsCore != nil && len(snapshot.Agent.XAgentsCore.HarnessConfig) > 0 { diff --git a/services/core/internal/execution/worker.go b/services/core/internal/execution/worker.go index 2ebb2c20f..137dd1a82 100644 --- a/services/core/internal/execution/worker.go +++ b/services/core/internal/execution/worker.go @@ -365,9 +365,6 @@ func (w *Worker) runClaim(ctx context.Context, item sessions.ExecutionWork) erro var rejection *preparationRejection capacityRejected := errors.As(err, &rejection) && rejection.operation == proto.TypeExecutionPrepare && rejection.code == "preparation_capacity" outcome := json.RawMessage(`{"error_code":"execution_unavailable"}`) - if errors.Is(err, ErrModelProviderRequired) { - outcome = json.RawMessage(`{"error_code":"model_provider_required"}`) - } finish, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() turn, err := w.dispatcher.SessionsReader.GetTurn(finish, item.TenantID, item.SessionID, item.TurnID) diff --git a/services/core/internal/execution/worker_schedule.go b/services/core/internal/execution/worker_schedule.go index de5ac027a..bf3a821cc 100644 --- a/services/core/internal/execution/worker_schedule.go +++ b/services/core/internal/execution/worker_schedule.go @@ -93,9 +93,6 @@ func (w *Worker) runEnvironmentInput(ctx context.Context, item scheduledWork) er if err == nil { return nil } - if errors.Is(err, ErrModelProviderRequired) { - return w.dispatcher.sessionExecution.FailEnvironmentInput(ctx, item.TenantID, item.SessionID, item.reservationID, "model_provider_required") - } if errors.Is(err, errPreparationFailed) && run.Reservation.State == sessions.EnvironmentInputPending { return w.dispatcher.sessionExecution.FailEnvironmentInput(ctx, item.TenantID, item.SessionID, item.reservationID, "runtime_preparation_failed") } diff --git a/services/core/internal/items/changes_test.go b/services/core/internal/items/changes_test.go index 38cf514ca..96f98766a 100644 --- a/services/core/internal/items/changes_test.go +++ b/services/core/internal/items/changes_test.go @@ -57,7 +57,7 @@ func TestObserveDecidesItemChanges(t *testing.T) { }, { name: "an input message takes no output index", kind: "message", - update: project("message", `{"text":"hello"}`), + update: project("message", `{"type":"agent.session.input.message","input":[{"role":"user","content":[{"type":"input_text","text":"hello"}]}]}`), check: func(c Change) bool { return !c.Output && c.Item.Role == "user" }, }, } { diff --git a/services/core/internal/items/inputs.go b/services/core/internal/items/inputs.go index 62d0efe8d..67bd00d4d 100644 --- a/services/core/internal/items/inputs.go +++ b/services/core/internal/items/inputs.go @@ -9,20 +9,16 @@ import ( func inputMessages(turn string, sequence int64, raw json.RawMessage) []Update { var p struct { - Text *string `json:"text"` Input []struct { Role string `json:"role"` Content []v1.ItemContent `json:"content"` } `json:"input"` } - // Internal admission predates the public schema and accepts arbitrary objects. + // Session admission stores any JSON object; only public messages project. if err := json.Unmarshal(raw, &p); err != nil { return nil } key := "input:" + strconv.FormatInt(sequence, 10) - if p.Text != nil { - return []Update{{Item: message(turn, key, "user", *p.Text, "completed")}} - } var updates []Update for i, input := range p.Input { if input.Role != "user" || len(input.Content) == 0 { diff --git a/services/core/internal/persistence/postgres/sessionpg/execution_inputs_test.go b/services/core/internal/persistence/postgres/sessionpg/execution_inputs_test.go index 22807a559..80af513a5 100644 --- a/services/core/internal/persistence/postgres/sessionpg/execution_inputs_test.go +++ b/services/core/internal/persistence/postgres/sessionpg/execution_inputs_test.go @@ -308,7 +308,7 @@ func TestEnvironmentInputPromotionRollsBackHistoryAndSettlement(t *testing.T) { ctx := t.Context() pending := reserve(t, service, tenant, session, "pending") name := "reservation_failure_" + strings.ReplaceAll(uuid.NewString(), "-", "") - table, expression := "turn_inputs", "session_id <> '"+text(session)+"'::uuid OR payload->>'text' <> 'second'" + table, expression := "turn_inputs", "session_id <> '"+text(session)+"'::uuid OR payload#>>'{input,0,content,0,text}' <> 'second'" switch phase { case "settlement": table, expression = "environment_input_reservations", "id <> '"+pending.ID+"'::uuid OR state <> 'admitted'" diff --git a/services/core/internal/persistence/postgres/sessionpg/inputs_test.go b/services/core/internal/persistence/postgres/sessionpg/inputs_test.go index e2f9e1b2f..450660b4b 100644 --- a/services/core/internal/persistence/postgres/sessionpg/inputs_test.go +++ b/services/core/internal/persistence/postgres/sessionpg/inputs_test.go @@ -14,6 +14,7 @@ import ( "github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgxpool" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -22,8 +23,10 @@ var messagePayload = json.RawMessage(`{"input":[{"role":"user","content":[{"type var cancelInput = sessions.Input{Kind: "cancel", Payload: json.RawMessage(`{}`)} +// messageInput is a public message event with one text part, as the events +// route stores it. func messageInput(text string) sessions.Input { - payload, _ := json.Marshal(map[string]string{"text": text}) + payload, _ := json.Marshal(v1.SessionInput{Type: "agent.session.input.message", Input: []v1.InputMessage{{Role: "user", Content: []v1.InputContent{{Type: "input_text", Text: &text}}}}}) return sessions.Input{Kind: "message", Payload: payload} } @@ -314,7 +317,7 @@ func TestInputBatchesAreOrderedAndIdempotentAcrossConnections(t *testing.T) { t.Fatalf("inputs=%d err=%v", len(inputs), err) } for i, input := range inputs { - var payload map[string]string + var payload v1.SessionInput if err := json.Unmarshal(input.Payload, &payload); err != nil { t.Fatal(err) } @@ -322,7 +325,7 @@ func TestInputBatchesAreOrderedAndIdempotentAcrossConnections(t *testing.T) { if i%2 == 1 { want = "second" } - if payload["text"] != want { + if *payload.Input[0].Content[0].Text != want { t.Fatalf("batch interleaved at %d: %v", i, payload) } } @@ -355,7 +358,7 @@ func TestBatchRetriesCompareTheWholeRequestAndRetainTargets(t *testing.T) { t.Fatalf("changed batch accepted: %v", err) } } - batch[1].Payload = json.RawMessage(`{ "text" : "one" }`) + batch[1].Payload = json.RawMessage(`{ "input" : [ { "content" : [ { "text" : "one", "type" : "input_text" } ], "role" : "user" } ], "type" : "agent.session.input.message" }`) pool.Close() restartedStore, restarted := stagingService(t, pgtest.Open(t)) retry, err := restarted.SubmitInputs(ctx, text(tenant), text(session), "mixed", batch) diff --git a/services/core/internal/persistence/postgres/sessionpg/items.go b/services/core/internal/persistence/postgres/sessionpg/items.go index 3e818822f..9c8072bc4 100644 --- a/services/core/internal/persistence/postgres/sessionpg/items.go +++ b/services/core/internal/persistence/postgres/sessionpg/items.go @@ -64,7 +64,7 @@ func (t *SessionTx) PutItem(ctx context.Context, turnID string, created time.Tim if err != nil { return nil, err } - payload, err := change.Item.MarshalStored() + payload, err := json.Marshal(change.Item) if err != nil { return nil, err } diff --git a/services/core/internal/sessions/creation_test.go b/services/core/internal/sessions/creation_test.go index 2fa58c0da..865fee9e1 100644 --- a/services/core/internal/sessions/creation_test.go +++ b/services/core/internal/sessions/creation_test.go @@ -423,7 +423,7 @@ func TestCreateSession(t *testing.T) { } want := []string{"UpsertSession create", "LockSkills skill_a", "ReadSkillVersion skill_a 3", "SaveModelExecution https://model.example/v1", "SaveExecutionConfiguration available " + uuid.Nil.String(), "SaveInitialFiles 1", "SaveSetup skill_a@3", "CreateEnvironment", - "LoadEnvironmentInput", `CreateInputReservation [{"kind":"message","payload":{"text":"hi"}}] initial`, "LoadEnvironmentInput", + "LoadEnvironmentInput", `CreateInputReservation [{"kind":"message","payload":` + hi + `}] initial`, "LoadEnvironmentInput", "PruneChanges", "AuditCreation session:session environment:environment:session", "LoadSession"} if strings.Join(calls, "\n") != strings.Join(want, "\n") { t.Fatalf("calls:\n%s\nwant:\n%s", strings.Join(calls, "\n"), strings.Join(want, "\n")) diff --git a/services/core/internal/sessions/environment_inputs_test.go b/services/core/internal/sessions/environment_inputs_test.go index f49fd2aaf..46b8cc9f7 100644 --- a/services/core/internal/sessions/environment_inputs_test.go +++ b/services/core/internal/sessions/environment_inputs_test.go @@ -51,7 +51,7 @@ func TestJoinsActiveTurn(t *testing.T) { } func TestReserveEnvironmentInput(t *testing.T) { - const batch = `[{"kind":"message","payload":{"text":"hi"}}]` + const batch = `[{"kind":"message","payload":` + hi + `}]` find := "FindInputReservation request " + batch reserve := func(t *testing.T, tx *fakeInputTx) (EnvironmentInputReservation, error) { tx.loadEnvironmentInput = inputs() diff --git a/services/core/internal/sessions/inputs_test.go b/services/core/internal/sessions/inputs_test.go index 22930d983..d847b3f82 100644 --- a/services/core/internal/sessions/inputs_test.go +++ b/services/core/internal/sessions/inputs_test.go @@ -10,6 +10,7 @@ import ( "testing" "time" + v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/items" ) @@ -104,11 +105,16 @@ func (s *fakeStorage) WithInputs(ctx context.Context, tenant, session string, ap // inputKey is the idempotency key of the input batches under test. const inputKey = "request" +// messageInput is a public message event with one text part, as the events +// route stores it. func messageInput(text string) Input { - payload, _ := json.Marshal(map[string]string{"text": text}) + payload, _ := json.Marshal(v1.SessionInput{Type: "agent.session.input.message", Input: []v1.InputMessage{{Role: "user", Content: []v1.InputContent{{Type: "input_text", Text: &text}}}}}) return Input{Kind: "message", Payload: payload} } +// hi is the payload of messageInput("hi") as admission normalizes it. +const hi = `{"input":[{"content":[{"text":"hi","type":"input_text"}],"role":"user"}],"type":"agent.session.input.message"}` + var cancelInput = Input{Kind: "cancel", Payload: json.RawMessage(`{}`)} // sequences is a fake CreateTurnInput that allocates sequences from first. @@ -144,7 +150,7 @@ func TestValidateInputs(t *testing.T) { } func TestValidateMessageInputs(t *testing.T) { - if _, encoded, err := validateMessageInputs([]Input{messageInput("hi")}); err != nil || string(encoded) != `[{"kind":"message","payload":{"text":"hi"}}]` { + if _, encoded, err := validateMessageInputs([]Input{messageInput("hi")}); err != nil || string(encoded) != `[{"kind":"message","payload":`+hi+`}]` { t.Fatalf("batch %s, %v", encoded, err) } for name, inputs := range map[string][]Input{"cancel": {messageInput("hi"), cancelInput}, "no inputs": nil} { @@ -205,13 +211,13 @@ func TestAdmitInput(t *testing.T) { tx := newInputTx(t) tx.loadActiveTurn, tx.createTurn, tx.appendChanges = activeTurn(nil), returns(turnWith(TurnQueued)), collect(&changes) tx.createTurnInput, tx.loadUsage = sequences(7), returns(json.RawMessage(`{}`)) - tx.loadInputSource = returns(Source{Turn: testTurn, Kind: "message", Sequence: 7, Payload: json.RawMessage(`{"text":"hi"}`)}) + tx.loadInputSource = returns(Source{Turn: testTurn, Kind: "message", Sequence: 7, Payload: messageInput("hi").Payload}) tx.loadItem, tx.putItem = returns(items.Stored{}), func(items.Change) (*int32, error) { return nil, nil } receipt, err := admitInput(t.Context(), tx, inputKey, 0, messageInput("hi")) if err != nil || receipt != (InputReceipt{Sequence: 7, TurnID: testTurn}) { t.Fatalf("receipt %+v, %v", receipt, err) } - item := items.Identity(testTurn, "input:7") + item := items.Identity(testTurn, "input:7:0") assertCalls(t, tx.fakeTx, "LoadActiveTurn", "CreateTurn", "AppendChanges agent.session.turn.created", "CreateTurnInput "+testTurn+" request 0 message", "LoadInputSource 7", "LoadItem "+testTurn+" "+item, "PutItem "+testTurn+" "+item, "AppendChanges agent.session.turn.item.added", "LoadUsage", "AppendChanges agent.session.in_progress") @@ -255,7 +261,7 @@ func TestAdmitInput(t *testing.T) { } func TestSubmitInputs(t *testing.T) { - const batch = `[{"kind":"message","payload":{"text":"hi"}},{"kind":"cancel","payload":{}}]` + const batch = `[{"kind":"message","payload":` + hi + `},{"kind":"cancel","payload":{}}]` inputs := []Input{messageInput("hi"), cancelInput} running := turnWith(TurnInProgress) submit := func(t *testing.T, tx *fakeInputTx, inputs []Input) ([]InputReceipt, error) { diff --git a/services/core/internal/sessions/subagents.go b/services/core/internal/sessions/subagents.go index ba982ab0b..14327ff31 100644 --- a/services/core/internal/sessions/subagents.go +++ b/services/core/internal/sessions/subagents.go @@ -346,7 +346,7 @@ func projectSubagentItem(ctx context.Context, tx ProjectionTx, raw json.RawMessa // the Session stream carries root work, and child history is read through the // Subagent routes. func putChildItem(ctx context.Context, tx SubagentProjectionTx, turn ChildTurn, position int32, item v1.Item) error { - payload, err := item.MarshalStored() + payload, err := json.Marshal(item) if err != nil { return err } diff --git a/services/core/tests/fixtures/main.go b/services/core/tests/fixtures/main.go index b55810eca..c97d6f33b 100644 --- a/services/core/tests/fixtures/main.go +++ b/services/core/tests/fixtures/main.go @@ -58,7 +58,7 @@ func seed() error { return err } for _, status := range []string{sessions.TurnCompleted, sessions.TurnFailed, sessions.TurnCancelled, sessions.TurnInProgress} { - receipts, err := service.SubmitInputs(ctx, f.Tenant, f.Session, uuid.NewString(), []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"recovery fixture"}`)}}) + receipts, err := service.SubmitInputs(ctx, f.Tenant, f.Session, uuid.NewString(), []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"type":"agent.session.input.message","input":[{"role":"user","content":[{"type":"input_text","text":"recovery fixture"}]}]}`)}}) if err != nil { return err } diff --git a/services/core/tests/integration/admin_session_archive_race_test.go b/services/core/tests/integration/admin_session_archive_race_test.go index 2bd929df6..e973a7c72 100644 --- a/services/core/tests/integration/admin_session_archive_race_test.go +++ b/services/core/tests/integration/admin_session_archive_race_test.go @@ -1,7 +1,6 @@ package integration import ( - "encoding/json" "errors" "strings" "sync" @@ -59,7 +58,7 @@ func TestManagedSessionArchiveOrdersConcurrentInput(t *testing.T) { go func() { defer wg.Done() <-start - _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "racing-input", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"racing"}`)}}) + _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "racing-input", []sessions.Input{messageInput("racing")}) if err != nil && !errors.Is(err, sessions.ErrEnvironmentUnavailable) { t.Error(err) } diff --git a/services/core/tests/integration/admin_session_archive_test.go b/services/core/tests/integration/admin_session_archive_test.go index df75725eb..846c69b12 100644 --- a/services/core/tests/integration/admin_session_archive_test.go +++ b/services/core/tests/integration/admin_session_archive_test.go @@ -85,7 +85,7 @@ func archiveAllocation(t *testing.T, w *Store, tenant string, session sessions.S func TestManagedSessionArchiveUnallocatedAndGuards(t *testing.T) { s, w, installation := managedArchiveFixture(t) input := managerSessionInput(uuid.NewString()) - input.InitialInputs = []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"waiting"}`)}} + input.InitialInputs = []sessions.Input{messageInput("waiting")} tenant, session := managedArchiveSession(t, s, input) ctx := adminDeleteContext(t.Context(), tenant, uuid.NewString()) active, err := sessionAdapter(s).GetManagedSessionArchive(t.Context(), tenant, session.ID) @@ -135,7 +135,7 @@ func TestManagedSessionArchiveUnallocatedAndGuards(t *testing.T) { if _, err := deploymentExecution(t, w).ReserveAllocation(t.Context(), deployment.AllocationKey{TenantID: tenant, EnvironmentID: session.Environment.ID}, installation, runtimedevice.HashCredential(uuid.NewString())); !errors.Is(err, deployment.ErrInvalidInput) { t.Fatal("archived Environment allocated after archive", err) } - if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "later", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"later"}`)}}); !errors.Is(err, sessions.ErrEnvironmentUnavailable) { + if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "later", []sessions.Input{messageInput("later")}); !errors.Is(err, sessions.ErrEnvironmentUnavailable) { t.Fatal("archived Environment accepted new input", err) } view, err := deploymentService(t, s).View(t.Context()) diff --git a/services/core/tests/integration/admin_session_archive_worker_http_test.go b/services/core/tests/integration/admin_session_archive_worker_http_test.go index 134fd92d0..00f830760 100644 --- a/services/core/tests/integration/admin_session_archive_worker_http_test.go +++ b/services/core/tests/integration/admin_session_archive_worker_http_test.go @@ -94,7 +94,7 @@ func TestAdminSessionArchiveWorkerHTTPPostgres(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(t.Context(), s, project.TenantID, active.ID, "pending-turn", json.RawMessage(`{"text":"pending"}`)) + input, err := sendMessage(t.Context(), s, project.TenantID, active.ID, "pending-turn", messageText("pending")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/archive_cancellation_test.go b/services/core/tests/integration/archive_cancellation_test.go index 0b64fef50..820ffb596 100644 --- a/services/core/tests/integration/archive_cancellation_test.go +++ b/services/core/tests/integration/archive_cancellation_test.go @@ -123,7 +123,7 @@ func TestArchiveWaitingCancellationReceipts(t *testing.T) { } time.Sleep(time.Millisecond) } - pending, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), h.tenant, session.ID, "pending", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}, {Kind: "message", Payload: json.RawMessage(`{"text":"second"}`)}}) + pending, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), h.tenant, session.ID, "pending", []sessions.Input{messageInput("first"), messageInput("second")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/claude_execution_test.go b/services/core/tests/integration/claude_execution_test.go index 648969f7c..c5c718d74 100644 --- a/services/core/tests/integration/claude_execution_test.go +++ b/services/core/tests/integration/claude_execution_test.go @@ -139,7 +139,7 @@ func TestClaudeInvalidImageResultRejectsWholeBatchBeforePersistence(t *testing.T } return sessions.Input{Kind: "tool_result", Payload: payload} } - batch := []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"Follow up"}`)}, result(`{"success":true,"output":[{"type":"input_image","image_url":"data:image/png;base64,AA=="}]}`), {Kind: "cancel", Payload: json.RawMessage(`{}`)}} + batch := []sessions.Input{messageInput("Follow up"), result(`{"success":true,"output":[{"type":"input_image","image_url":"data:image/png;base64,AA=="}]}`), {Kind: "cancel", Payload: json.RawMessage(`{}`)}} if _, err := worker.SubmitInputs(t.Context(), h.tenant, h.session.ID, "batch", batch); !errors.Is(err, sessions.ErrInvalidInput) { t.Fatal(err) } diff --git a/services/core/tests/integration/command_output_test.go b/services/core/tests/integration/command_output_test.go index 01406737e..647d05650 100644 --- a/services/core/tests/integration/command_output_test.go +++ b/services/core/tests/integration/command_output_test.go @@ -21,7 +21,7 @@ func TestCommandOutputCommitsFragmentsSnapshotsAndRecovery(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "start", json.RawMessage(`{"text":"run commands"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "start", messageText("run commands")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/creation_stream_settlement_public_test.go b/services/core/tests/integration/creation_stream_settlement_public_test.go index 63e64c863..954300f88 100644 --- a/services/core/tests/integration/creation_stream_settlement_public_test.go +++ b/services/core/tests/integration/creation_stream_settlement_public_test.go @@ -170,7 +170,7 @@ func TestCreationStreamPublicLifetimes(t *testing.T) { created.ended(t, 5*time.Second) connect(first.Session.Environment.ID) - if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, first.Session.ID, "later", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"later"}`)}}); err != nil { + if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, first.Session.ID, "later", []sessions.Input{messageInput("later")}); err != nil { t.Fatal(err) } if current, err := sessionAdapter(s).GetSession(t.Context(), tenant, first.Session.ID); err != nil || !current.PendingInput { diff --git a/services/core/tests/integration/deployment_model_providers_http_test.go b/services/core/tests/integration/deployment_model_providers_http_test.go index b796d6ebf..8dc8687af 100644 --- a/services/core/tests/integration/deployment_model_providers_http_test.go +++ b/services/core/tests/integration/deployment_model_providers_http_test.go @@ -4,24 +4,19 @@ import ( "bytes" "context" "encoding/json" - "errors" "net/http/httptest" "strings" "testing" - "time" v1 "github.com/MiniMax-AI/OpenAgentCore/contracts/agents-api/v1" - "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/adminaudit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/api" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/credentialcrypto" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/modelconfiguration" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/auditpg" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/modelconfigurationpg" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" "github.com/google/uuid" ) @@ -230,61 +225,6 @@ func TestDeploymentModelProvidersHTTP(t *testing.T) { } } -// A hosted or self-hosted Session created before providers were required has -// no frozen provider: new work is rejected before anything is queued, and input -// reserved before the upgrade fails with that reason instead of waiting. -func TestLegacySessionWithoutProviderCannotStartWork(t *testing.T) { - h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), true) - legacy, err := h.s.CreateSession(t.Context(), h.tenant, sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(), - Configuration: []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`)}) - if err != nil { - t.Fatal(err) - } - executor := connectFixtureRuntime(t, h, legacy) - // Reserved directly, as a pre-upgrade Core did. - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, legacy.ID, "before-upgrade", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"old"}`)}}) - if err != nil { - t.Fatal(err) - } - worker, stop := startEnvironmentExpiryWorker(t, h.s, h.d) - defer stop() - _, pool := testStore(t) - reservations := func() int { - t.Helper() - var count int - if err := pool.QueryRow(t.Context(), "SELECT count(*) FROM environment_input_reservations WHERE session_id=$1", legacy.ID).Scan(&count); err != nil { - t.Fatal(err) - } - return count - } - before := reservations() - message := []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"start"}`)}} - if _, err := worker.SubmitInputs(t.Context(), h.tenant, legacy.ID, uuid.NewString(), message); !errors.Is(err, execution.ErrModelProviderRequired) { - t.Fatal("provider-free Session accepted work", err) - } - if after := reservations(); after != before { - t.Fatal("rejected work was queued", before, after) - } - awaitDaemonRemoteCondition(t, t.Context(), 5*time.Second, "legacy reservation settled", func() bool { - got, err := sessionAdapter(h.s).GetEnvironmentInputReservation(t.Context(), h.tenant, legacy.ID, pending.ID) - return err == nil && got.State == sessions.EnvironmentInputFailed - }) - session, err := sessionAdapter(h.s).GetSession(t.Context(), h.tenant, legacy.ID) - if err != nil || session.EnvironmentInputActivity == nil || session.EnvironmentInputActivity.Failure != "model_provider_required" { - t.Fatal("legacy reservation did not fail with its reason", session.EnvironmentInputActivity, err) - } - _ = executor.conn.SetReadDeadline(time.Now().Add(500 * time.Millisecond)) - for { - var frame proto.Envelope - if executor.conn.ReadJSON(&frame) != nil { - break - } - if frame.Type == proto.TypeExecutionPrepare { - t.Fatal("provider-free work reached the executor", frame.Type) - } - } -} - // A none Session may freeze the deployment default, so its caller intent is // recorded first: a same-key retry returns the committed Session after the // default was replaced or removed. diff --git a/services/core/tests/integration/dispatch_test.go b/services/core/tests/integration/dispatch_test.go index 36141beec..7e733aad4 100644 --- a/services/core/tests/integration/dispatch_test.go +++ b/services/core/tests/integration/dispatch_test.go @@ -125,8 +125,7 @@ func newDispatchHarnessForSession(t *testing.T, configuration []byte, local bool func (h *dispatchHarness) message(key, text string) sessions.InputReceipt { h.t.Helper() - body, _ := json.Marshal(map[string]string{"text": text}) - r, err := sendMessage(context.Background(), h.s, h.tenant, h.session.ID, key, body) + r, err := sendMessage(context.Background(), h.s, h.tenant, h.session.ID, key, messageText(text)) if err != nil { h.t.Fatal(err) } diff --git a/services/core/tests/integration/environment_admission_test.go b/services/core/tests/integration/environment_admission_test.go index 27e191352..cfbc778fb 100644 --- a/services/core/tests/integration/environment_admission_test.go +++ b/services/core/tests/integration/environment_admission_test.go @@ -51,7 +51,7 @@ func submitEnvironmentAdmission(ctx context.Context, h *dispatchHarness, worker } func environmentAdmissionInputs() []sessions.Input { - return []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}, {Kind: "message", Payload: json.RawMessage(`{"text":"second"}`)}} + return []sessions.Input{messageInput("first"), messageInput("second")} } func awaitEnvironmentAdmission(t *testing.T, result <-chan environmentAdmissionResult) environmentAdmissionResult { diff --git a/services/core/tests/integration/environment_directory_active_test.go b/services/core/tests/integration/environment_directory_active_test.go index 96696846f..6d63a0c74 100644 --- a/services/core/tests/integration/environment_directory_active_test.go +++ b/services/core/tests/integration/environment_directory_active_test.go @@ -10,7 +10,7 @@ import ( func TestEnvironmentDirectoryActiveRunUsesExistingOwner(t *testing.T) { h, w, environment := directoryWorker(t) awaitFixtureCapabilities(t, h, workerEnvironmentCapabilities()) - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "execute", []sessions.Input{{Kind: "message", Payload: []byte(`{"text":"work"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "execute", []sessions.Input{messageInput("work")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/environment_expiry_dispatch_test.go b/services/core/tests/integration/environment_expiry_dispatch_test.go index 623fe0ea5..4bc66594b 100644 --- a/services/core/tests/integration/environment_expiry_dispatch_test.go +++ b/services/core/tests/integration/environment_expiry_dispatch_test.go @@ -1,7 +1,6 @@ package integration import ( - "encoding/json" "testing" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -24,7 +23,7 @@ func TestWorkerEnvironmentExpiryAtFullExecutionCapacity(t *testing.T) { var active []sessions.Session for _, key := range []string{"one", "two", "three", "four"} { session := publicSession(t, h, key) - if _, err := worker.SubmitInputs(t.Context(), h.tenant, session.ID, key, []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"remain active"}`)}}); err != nil { + if _, err := worker.SubmitInputs(t.Context(), h.tenant, session.ID, key, []sessions.Input{messageInput("remain active")}); err != nil { t.Fatal(err) } requests = append(requests, h.read(testExecutionRequest)) @@ -66,7 +65,7 @@ func TestWorkerEnvironmentExpirySkipsBusySessionAndAllowsDispatch(t *testing.T) } worker, stop := startEnvironmentExpiryWorker(t, h.s, h.d) h.session = publicSession(t, h, "unrelated") - receipt, err := worker.SubmitInputs(t.Context(), h.tenant, h.session.ID, "work", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"make normal progress"}`)}}) + receipt, err := worker.SubmitInputs(t.Context(), h.tenant, h.session.ID, "work", []sessions.Input{messageInput("make normal progress")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/environment_expiry_worker_test.go b/services/core/tests/integration/environment_expiry_worker_test.go index a7d8b4fdd..04838a170 100644 --- a/services/core/tests/integration/environment_expiry_worker_test.go +++ b/services/core/tests/integration/environment_expiry_worker_test.go @@ -25,7 +25,7 @@ func newEnvironmentExpiryReservation(t *testing.T, s *Store) (string, sessions.E if err != nil { t.Fatal(err) } - pending, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "pending", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"wait for the environment"}`)}}) + pending, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "pending", []sessions.Input{messageInput("wait for the environment")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/environment_initial_input_test.go b/services/core/tests/integration/environment_initial_input_test.go index 5bd57dfff..da0f37d64 100644 --- a/services/core/tests/integration/environment_initial_input_test.go +++ b/services/core/tests/integration/environment_initial_input_test.go @@ -105,9 +105,11 @@ func TestEnvironmentInitialInputCreationRetainsCursorIdentityAndPromotion(t *tes t.Fatal("creation snapshot differs from the committed projection", session.EnvironmentInputActivity, session.PendingInput) } reservation := initialEnvironmentReservation(t, s, pool, tenant, session.ID) - storedBatch, marshalErr := json.Marshal(reservation.Inputs) - originalBatch, _ := json.Marshal(input.InitialInputs) - if marshalErr != nil || reservation.State != sessions.EnvironmentInputPending || reservation.Deadline.Sub(reservation.CreatedAt) != 5*time.Minute || string(storedBatch) != string(originalBatch) { + // jsonb keeps the batch's JSON value, not its key order. + var storedBatch, originalBatch any + stored, marshalErr := json.Marshal(reservation.Inputs) + original, _ := json.Marshal(input.InitialInputs) + if marshalErr != nil || json.Unmarshal(stored, &storedBatch) != nil || json.Unmarshal(original, &originalBatch) != nil || reservation.State != sessions.EnvironmentInputPending || reservation.Deadline.Sub(reservation.CreatedAt) != 5*time.Minute || !reflect.DeepEqual(storedBatch, originalBatch) { t.Fatal("initial batch/deadline changed", reservation) } environmentInputHistory(t, pool, session.ID, 0, 0) diff --git a/services/core/tests/integration/environment_initial_public_test.go b/services/core/tests/integration/environment_initial_public_test.go index a3b47fa76..30559672c 100644 --- a/services/core/tests/integration/environment_initial_public_test.go +++ b/services/core/tests/integration/environment_initial_public_test.go @@ -31,7 +31,7 @@ func TestEnvironmentInitialFailureOfficialClient(t *testing.T) { configuration := json.RawMessage(`{"agent":{"id":"agent_initial_failure","model":"fixture","tools":[],"multi_agent":{"enabled":false,"max_concurrent_subagents":null},"reasoning":{},"service_tier":"auto","text":{"format":{"type":"text"},"verbosity":"medium"}},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`) session, err := s.CreateSession(t.Context(), tenant, sessions.CreateSession{ Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: "initial", Configuration: configuration, - InitialInputs: []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"private-input-marker"}`)}}, + InitialInputs: []sessions.Input{messageInput("private-input-marker")}, }) if err != nil { t.Fatal(err) diff --git a/services/core/tests/integration/environment_worker_helpers_test.go b/services/core/tests/integration/environment_worker_helpers_test.go index d587a3ca4..3ddd0529c 100644 --- a/services/core/tests/integration/environment_worker_helpers_test.go +++ b/services/core/tests/integration/environment_worker_helpers_test.go @@ -47,7 +47,7 @@ func unboundWorkerEnvironmentReservation(t *testing.T, h *dispatchHarness) sessi if err != nil { t.Fatal(err) } - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, session.ID, "work", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, session.ID, "work", []sessions.Input{messageInput("first")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/environments_test.go b/services/core/tests/integration/environments_test.go index ee0c7a98f..430b829a1 100644 --- a/services/core/tests/integration/environments_test.go +++ b/services/core/tests/integration/environments_test.go @@ -189,7 +189,7 @@ func TestEnvironmentCreationFailureRollsBackAllResources(t *testing.T) { constraint := "environment_failure_" + strings.ReplaceAll(marker, "-", "") table, expression := "environments", "status <> 'pending'" if phase == "input" { - table, expression = "environment_input_reservations", "NOT (batch @> '[{\"payload\":{\"text\":\""+marker+"\"}}]'::jsonb)" + table, expression = "environment_input_reservations", "NOT (batch @> '[{\"payload\":{\"input\":[{\"content\":[{\"text\":\""+marker+"\"}]}]}}]'::jsonb)" } if phase == "activity" { table, expression = "session_events", "NOT (payload ? 'environment_input_activity')" diff --git a/services/core/tests/integration/execution_test.go b/services/core/tests/integration/execution_test.go index 0be73a2e7..8e73078db 100644 --- a/services/core/tests/integration/execution_test.go +++ b/services/core/tests/integration/execution_test.go @@ -276,7 +276,7 @@ func TestExecutionWriterSerializesWritesOnItsLease(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), time.Second) defer cancel() task := tasks[0] - _, err := sendMessage(ctx, s, task.tenant, task.session, "public", json.RawMessage(`{"text":"additional"}`)) + _, err := sendMessage(ctx, s, task.tenant, task.session, "public", messageText("additional")) close(release) if err != nil { t.Fatal("public admission used owner gate", err) diff --git a/services/core/tests/integration/function_input_execution_test.go b/services/core/tests/integration/function_input_execution_test.go index 63dbe8d17..3d5aeeafa 100644 --- a/services/core/tests/integration/function_input_execution_test.go +++ b/services/core/tests/integration/function_input_execution_test.go @@ -18,7 +18,7 @@ func TestExecutionFunctionInputBatchStillSteersMessages(t *testing.T) { h.write(input.TurnID, proto.TypeFunctionCall, proto.FunctionCallPayload{CallID: "a", Name: "lookup_ticket", Arguments: json.RawMessage(`{}`)}) state := functionState(t, h, 1) raw, _ := json.Marshal(sessions.FunctionResultInput{TurnID: input.TurnID, CallID: state.RequiredActions[0].CallID, Result: json.RawMessage(`{"success":true,"output":"answer"}`)}) - batch := []sessions.Input{{Kind: "tool_result", Payload: raw}, {Kind: "message", Payload: json.RawMessage(`{"text":"Follow up"}`)}} + batch := []sessions.Input{{Kind: "tool_result", Payload: raw}, messageInput("Follow up")} receipts, err := submitInputs(t.Context(), h.s, h.tenant, h.session.ID, "mixed", batch) if err != nil { t.Fatal(err) diff --git a/services/core/tests/integration/function_inputs_public_test.go b/services/core/tests/integration/function_inputs_public_test.go index 486a53276..5b512cffe 100644 --- a/services/core/tests/integration/function_inputs_public_test.go +++ b/services/core/tests/integration/function_inputs_public_test.go @@ -33,7 +33,7 @@ func TestFunctionInputsOfficialClientAtomicAdmission(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "start", json.RawMessage(`{"text":"fixture"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "start", messageText("fixture")) if err != nil { t.Fatal(err) } @@ -110,7 +110,7 @@ func TestFunctionInputsOfficialClientAtomicAdmission(t *testing.T) { if _, err := transitionTurn(ctx, s, tenant, session.ID, input.TurnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnWaiting, Status: sessions.TurnFailed}); err != nil { t.Fatal(err) } - next, err := sendMessage(ctx, s, tenant, session.ID, "next", json.RawMessage(`{"text":"next"}`)) + next, err := sendMessage(ctx, s, tenant, session.ID, "next", messageText("next")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/function_inputs_test.go b/services/core/tests/integration/function_inputs_test.go index 435e1579f..8046aa039 100644 --- a/services/core/tests/integration/function_inputs_test.go +++ b/services/core/tests/integration/function_inputs_test.go @@ -39,7 +39,7 @@ func TestFunctionInputBatchesPersistAndReplayWithoutRetargeting(t *testing.T) { s, pool := testStore(t) tenant, session, turn := functionInputFixture(t, s, functionExecution(t)) full := `{"success":false,"output":[{"type":"input_text","text":""},{"type":"input_image","image_url":"data:image/png;base64,AA=="},{"type":"input_text","text":"after"}],"error":"failed"}` - batch := []sessions.Input{resultInput(t, turn, "a", full), {Kind: "message", Payload: json.RawMessage(`{"text":"Follow up"}`)}, resultInput(t, turn, "b", `{"success":true,"output":null,"error":null}`), {Kind: "cancel", Payload: json.RawMessage(`{}`)}} + batch := []sessions.Input{resultInput(t, turn, "a", full), messageInput("Follow up"), resultInput(t, turn, "b", `{"success":true,"output":null,"error":null}`), {Kind: "cancel", Payload: json.RawMessage(`{}`)}} receipts, err := submitInputs(t.Context(), s, tenant, session.ID, "batch", batch) if err != nil || len(receipts) != 4 { t.Fatal(receipts, err) @@ -105,7 +105,7 @@ func TestFunctionInputBatchFailureRollsBackEveryWrite(t *testing.T) { s, _ := testStore(t) functions := functionExecution(t) tenant, session, turn := functionInputFixture(t, s, functions) - message := sessions.Input{Kind: "message", Payload: json.RawMessage(`{"text":"Must roll back"}`)} + message := messageInput("Must roll back") cancel := sessions.Input{Kind: "cancel", Payload: json.RawMessage(`{}`)} first := resultInput(t, turn, "a", `{"success":true}`) batch := []sessions.Input{message, first, cancel} @@ -165,7 +165,7 @@ func TestFunctionInputConcurrentBatchesSelectOneResult(t *testing.T) { var wg sync.WaitGroup results := make(chan error, 2) for i := range 2 { - batch := []sessions.Input{{Kind: "message", Payload: json.RawMessage(fmt.Sprintf(`{"text":"message-%d"}`, i))}, resultInput(t, turn, "a", fmt.Sprintf(`{"success":true,"output":"%d"}`, i))} + batch := []sessions.Input{{Kind: "message", Payload: messageText(fmt.Sprintf("message-%d", i))}, resultInput(t, turn, "a", fmt.Sprintf(`{"success":true,"output":"%d"}`, i))} wg.Add(1) go func() { defer wg.Done() diff --git a/services/core/tests/integration/function_item_events_test.go b/services/core/tests/integration/function_item_events_test.go index 7fca61f85..edd9d16d7 100644 --- a/services/core/tests/integration/function_item_events_test.go +++ b/services/core/tests/integration/function_item_events_test.go @@ -64,7 +64,7 @@ func TestFunctionResultItemsRetainSubmittedFields(t *testing.T) { `{"success":false,"output":[{"type":"input_text","text":"before"},{"type":"input_image","image_url":"data:image/png;base64,AA=="}],"error":"failure"}`, } { t.Run(raw, func(t *testing.T) { - s, pool := testStore(t) + s, _ := testStore(t) functions := functionExecution(t) tenant, session := newTurnSession(t, s) turn := submitMessage(t, s, tenant, session.ID, "start").TurnID @@ -119,19 +119,6 @@ func TestFunctionResultItemsRetainSubmittedFields(t *testing.T) { results++ } } - // The stored payload keeps the submitted field presence. - var stored map[string]any - if err := pool.QueryRow(t.Context(), `SELECT payload FROM session_items WHERE turn_id = $1 AND payload->>'type' = 'function_call_output'`, turn).Scan(&stored); err != nil { - t.Fatal(err) - } - var submitted map[string]any - _ = json.Unmarshal([]byte(raw), &submitted) - for _, field := range []string{"output", "error"} { - _, present := submitted[field] - if _, exists := stored[field]; present != exists { - t.Fatalf("stored %s presence changed: %v", field, stored) - } - } if results != 2 { t.Fatal("missing saved or streamed result", results) } diff --git a/services/core/tests/integration/function_state_public_test.go b/services/core/tests/integration/function_state_public_test.go index 7c1b84cd7..538ccda86 100644 --- a/services/core/tests/integration/function_state_public_test.go +++ b/services/core/tests/integration/function_state_public_test.go @@ -30,7 +30,7 @@ func TestFunctionStateOfficialClientReadsAndLiveEvents(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "start", json.RawMessage(`{"text":"fixture"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "start", messageText("fixture")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/hosted_initialization_failure_public_test.go b/services/core/tests/integration/hosted_initialization_failure_public_test.go index ee461f4e5..543f11500 100644 --- a/services/core/tests/integration/hosted_initialization_failure_public_test.go +++ b/services/core/tests/integration/hosted_initialization_failure_public_test.go @@ -220,7 +220,7 @@ func TestHostedInitializationFailureRecordsSafeSessionFailure(t *testing.T) { !last.EnvironmentFailure.FailedAt.Equal(read.EnvironmentFailure.FailedAt) || last.EnvironmentInputActivity != nil || !last.Settled { t.Fatal("failed snapshot", last) } - if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "later", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"later"}`)}}); !errors.Is(err, sessions.ErrHostedEnvironmentFailed) { + if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, session.ID, "later", []sessions.Input{messageInput("later")}); !errors.Is(err, sessions.ErrHostedEnvironmentFailed) { t.Fatal("failed hosted Environment admitted input", err) } raw, _ := json.Marshal(events) @@ -249,7 +249,7 @@ func TestHostedInitializationFailureSettlesPendingInitialInput(t *testing.T) { tenant := uuid.NewString() session, environment := hostedFailureSession(t, s, tenant, sessions.CreateSession{ Initialization: environmentconfig.Setup{Commands: []environmentconfig.SetupCommand{{Command: "exit 3"}}}, - InitialInputs: []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"initial"}`)}}, + InitialInputs: []sessions.Input{messageInput("initial")}, }) p := &hostedFailureProvider{lifecycleProvider: lifecycleProvider{resources: map[string]sandbox.Info{}}, fail: "setup", result: failedInitialization(3)} diff --git a/services/core/tests/integration/input_conflicts_public_test.go b/services/core/tests/integration/input_conflicts_public_test.go index 41578612f..f390b0419 100644 --- a/services/core/tests/integration/input_conflicts_public_test.go +++ b/services/core/tests/integration/input_conflicts_public_test.go @@ -58,7 +58,7 @@ func TestSessionInputConflictsAndResultTargetsPostgres(t *testing.T) { input := sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(), Configuration: json.RawMessage(`{` + conflictAgent + `,"environment":` + environment + `}`)} if initial { - input.InitialInputs = []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"reserved"}`)}} + input.InitialInputs = []sessions.Input{messageInput("reserved")} } session, err := s.CreateSession(ctx, tenant, input) if err != nil { @@ -69,7 +69,7 @@ func TestSessionInputConflictsAndResultTargetsPostgres(t *testing.T) { // waiting starts a Turn that waits for one function result. waiting := func(session, key, call string) sessions.InputReceipt { t.Helper() - receipt, err := sendMessage(ctx, s, tenant, session, key, json.RawMessage(`{"text":"work"}`)) + receipt, err := sendMessage(ctx, s, tenant, session, key, messageText("work")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/inputs_test.go b/services/core/tests/integration/inputs_test.go index 8dd0279e2..36dd59666 100644 --- a/services/core/tests/integration/inputs_test.go +++ b/services/core/tests/integration/inputs_test.go @@ -9,6 +9,7 @@ import ( "github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgxpool" + 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/persistence/postgres/sessionpg" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" @@ -16,9 +17,15 @@ import ( var messagePayload = json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_text","text":"hello"}]}]}`) +// messageText is a public message event with one text part, as the events +// route stores it. +func messageText(text string) json.RawMessage { + payload, _ := json.Marshal(v1.SessionInput{Type: "agent.session.input.message", Input: []v1.InputMessage{{Role: "user", Content: []v1.InputContent{{Type: "input_text", Text: &text}}}}}) + return payload +} + func messageInput(text string) sessions.Input { - payload, _ := json.Marshal(map[string]string{"text": text}) - return sessions.Input{Kind: "message", Payload: payload} + return sessions.Input{Kind: "message", Payload: messageText(text)} } func newTurnSession(t *testing.T, s *Store) (string, sessions.Session) { diff --git a/services/core/tests/integration/item_order_test.go b/services/core/tests/integration/item_order_test.go index d7c152ac7..38b6c24d7 100644 --- a/services/core/tests/integration/item_order_test.go +++ b/services/core/tests/integration/item_order_test.go @@ -22,7 +22,7 @@ func TestItemObservationOrderSurvivesTiesUpdatesRetriesAndRecovery(t *testing.T) if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "first", json.RawMessage(`{"text":"question"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "first", messageText("question")) if err != nil { t.Fatal(err) } @@ -51,7 +51,7 @@ func TestItemObservationOrderSurvivesTiesUpdatesRetriesAndRecovery(t *testing.T) t.Fatal(err) } } - if _, err = sendMessage(ctx, s, tenant, session.ID, "steer", json.RawMessage(`{"text":"continue"}`)); err != nil { + if _, err = sendMessage(ctx, s, tenant, session.ID, "steer", messageText("continue")); err != nil { t.Fatal(err) } page, err = sessionAdapter(s).ListItems(ctx, tenant, session.ID, "", 100, true) @@ -131,7 +131,7 @@ func TestItemObservationOrderSurvivesTiesUpdatesRetriesAndRecovery(t *testing.T) t.Fatal(err) } checkOrder() - next, err := sendMessage(ctx, s, tenant, session.ID, "next-turn", json.RawMessage(`{"text":"new turn"}`)) + next, err := sendMessage(ctx, s, tenant, session.ID, "next-turn", messageText("new turn")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/item_reads_test.go b/services/core/tests/integration/item_reads_test.go index 31b01fa3b..27585bf45 100644 --- a/services/core/tests/integration/item_reads_test.go +++ b/services/core/tests/integration/item_reads_test.go @@ -22,7 +22,7 @@ func TestItemsRecoverSnapshotsPartialResultsPaginationAndIsolation(t *testing.T) if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "first", json.RawMessage(`{"text":"question"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "first", messageText("question")) if err != nil { t.Fatal(err) } @@ -122,7 +122,7 @@ func TestItemProjectionFailureRollsBackJournalAndAggregateRecovers(t *testing.T) journal := executionOwner(t, s).Sessions tenant := uuid.NewString() session, _ := s.CreateSession(ctx, tenant, sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: "legacy"}) - input, err := sendMessage(ctx, s, tenant, session.ID, "input", json.RawMessage(`{"text":"test"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "input", messageText("test")) if err != nil { t.Fatal(err) } @@ -159,7 +159,7 @@ func TestReceiptOnlyTextRecoversWithoutInventingCompletion(t *testing.T) { tenant := uuid.NewString() for _, receiptOnly := range []bool{true, false} { session, _ := s.CreateSession(ctx, tenant, sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString()}) - input, err := sendMessage(ctx, s, tenant, session.ID, "first", json.RawMessage(`{"text":"test"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "first", messageText("test")) if err != nil { t.Fatal(err) } @@ -193,7 +193,7 @@ func TestLegacyFailureRetainsPartialAnswerAcrossRecovery(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "first", json.RawMessage(`{"text":"question"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "first", messageText("question")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/local_artifact_export_test.go b/services/core/tests/integration/local_artifact_export_test.go index 06775a674..056591d51 100644 --- a/services/core/tests/integration/local_artifact_export_test.go +++ b/services/core/tests/integration/local_artifact_export_test.go @@ -3,7 +3,6 @@ package integration import ( "archive/tar" "bytes" - "encoding/json" "testing" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" @@ -46,7 +45,7 @@ func completeLocalArtifactExport(t *testing.T, h *dispatchHarness, worker *execu t.Fatal("capture published before native completion", err) } completeCaptureDirectoryRead(t, h, worker, environment) - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "during-artifact-capture", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"run after the completed native execution"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "during-artifact-capture", []sessions.Input{messageInput("run after the completed native execution")}) if err != nil || pending.State != sessions.EnvironmentInputPending || len(pending.Receipts) != 0 { t.Fatalf("input during artifact capture was assigned to the finished executor: %+v %v", pending, err) } diff --git a/services/core/tests/integration/local_environment_file_write_test.go b/services/core/tests/integration/local_environment_file_write_test.go index 0b802a30f..89d87e926 100644 --- a/services/core/tests/integration/local_environment_file_write_test.go +++ b/services/core/tests/integration/local_environment_file_write_test.go @@ -56,7 +56,7 @@ func TestLocalEnvironmentFileWriteOwnsMutationBeforeDispatch(t *testing.T) { if err != nil || intent.State != "pending" || intent.Identity.DeviceID != h.device.ID { t.Fatal("dispatch preceded durable ownership", intent, err) } - if _, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "concurrent", []sessions.Input{{Kind: "message", Payload: []byte(`{"text":"work"}`)}}); !errors.Is(err, sessions.ErrTurnConflict) { + if _, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "concurrent", []sessions.Input{messageInput("work")}); !errors.Is(err, sessions.ErrTurnConflict) { t.Fatal("upload admitted concurrent execution", err) } cancel() @@ -98,7 +98,7 @@ func TestLocalEnvironmentFileWriteLostReceiptRemainsPending(t *testing.T) { if err != nil || intent.State != "pending" { t.Fatal("disconnect guessed rejection", intent, err) } - if _, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "after-loss", []sessions.Input{{Kind: "message", Payload: []byte(`{"text":"work"}`)}}); !errors.Is(err, sessions.ErrTurnConflict) { + if _, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "after-loss", []sessions.Input{messageInput("work")}); !errors.Is(err, sessions.ErrTurnConflict) { t.Fatal("unknown upload admitted execution", err) } } diff --git a/services/core/tests/integration/local_environment_worker_test.go b/services/core/tests/integration/local_environment_worker_test.go index 494efc624..5d9b3214c 100644 --- a/services/core/tests/integration/local_environment_worker_test.go +++ b/services/core/tests/integration/local_environment_worker_test.go @@ -2,7 +2,6 @@ package integration import ( "context" - "encoding/json" "errors" "testing" "time" @@ -109,7 +108,7 @@ func TestLocalEnvironmentWorkerRejectsGeneralDeviceDespiteCapability(t *testing. func TestLocalEnvironmentWorkerSchedulesPreparationWithoutRemoteResolver(t *testing.T) { h, worker, environment := localWorker(t, true, true) - reservation, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "local-input", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + reservation, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "local-input", []sessions.Input{messageInput("first")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/message_image_admission_test.go b/services/core/tests/integration/message_image_admission_test.go index 02024727a..1cf450a41 100644 --- a/services/core/tests/integration/message_image_admission_test.go +++ b/services/core/tests/integration/message_image_admission_test.go @@ -14,7 +14,7 @@ import ( func imageAdmissionBatch() []sessions.Input { return []sessions.Input{ - {Kind: "message", Payload: json.RawMessage(`{"text":"do not partially admit"}`)}, + messageInput("do not partially admit"), {Kind: "message", Payload: json.RawMessage(`{"input":[{"role":"user","content":[{"type":"input_image","image_url":"data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAwMCAO+aXioAAAAASUVORK5CYII="}]}]}`)}, } } diff --git a/services/core/tests/integration/prepared_dispatch_test.go b/services/core/tests/integration/prepared_dispatch_test.go index bbb1f066f..08cfbdf04 100644 --- a/services/core/tests/integration/prepared_dispatch_test.go +++ b/services/core/tests/integration/prepared_dispatch_test.go @@ -24,7 +24,7 @@ func preparedDispatchHarness(t *testing.T) (*dispatchHarness, sessions.Environme assertNoRuntimeAllocation(t, h) h.d, h.lease = h.bound(), h.owner().Lease enableWorkerEnvironment(t, h) - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "pending", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}, {Kind: "message", Payload: json.RawMessage(`{"text":"second"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "pending", []sessions.Input{messageInput("first"), messageInput("second")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/public_execution_test.go b/services/core/tests/integration/public_execution_test.go index 879c09a1d..b327607b8 100644 --- a/services/core/tests/integration/public_execution_test.go +++ b/services/core/tests/integration/public_execution_test.go @@ -144,7 +144,7 @@ func TestWorkerRestartReconcilesClaimedButPreservesQueuedWork(t *testing.T) { } checkMeasurement(false) queued := publicSession(t, h, "queued") - if _, err := sendMessage(ctx, h.s, h.tenant, queued.ID, "first", json.RawMessage(`{"text":"Not sent"}`)); err != nil { + if _, err := sendMessage(ctx, h.s, h.tenant, queued.ID, "first", messageText("Not sent")); err != nil { t.Fatal(err) } worker := startOwnedWorker(t, ctx, h.s, h.d, h.owner()) diff --git a/services/core/tests/integration/runtime_input_admission_test.go b/services/core/tests/integration/runtime_input_admission_test.go index 66617e3db..8c45dac30 100644 --- a/services/core/tests/integration/runtime_input_admission_test.go +++ b/services/core/tests/integration/runtime_input_admission_test.go @@ -15,7 +15,7 @@ import ( func TestManagedRuntimeMaintenancePreservesCancelAndRetry(t *testing.T) { s, _ := newManagedTestStore(t) tenant, session, _ := managedSession(t, s) - inputs := []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"accepted work"}`)}} + inputs := []sessions.Input{messageInput("accepted work")} accepted, err := submitInputs(t.Context(), s, tenant, session.ID, "work", inputs) if err != nil { t.Fatal(err) diff --git a/services/core/tests/integration/runtime_pending_test.go b/services/core/tests/integration/runtime_pending_test.go index 65128f89c..8245750f5 100644 --- a/services/core/tests/integration/runtime_pending_test.go +++ b/services/core/tests/integration/runtime_pending_test.go @@ -16,7 +16,7 @@ import ( func TestManagedRuntimeAutomaticBootstrapRecoversCommittedSessions(t *testing.T) { s, _ := newManagedTestStore(t) tenant, idle, idleEnvironment := managedSession(t, s) - initial, err := s.CreateSession(t.Context(), tenant, sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(), Configuration: json.RawMessage(`{"agent":{"model":"test"},"environment":{"type":"openai_hosted"}}`), InitialInputs: []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"hello"}`)}}}) + initial, err := s.CreateSession(t.Context(), tenant, sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(), Configuration: json.RawMessage(`{"agent":{"model":"test"},"environment":{"type":"openai_hosted"}}`), InitialInputs: []sessions.Input{messageInput("hello")}}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/runtime_wake_hint_integration_test.go b/services/core/tests/integration/runtime_wake_hint_integration_test.go index f31a2e4a7..0d29a6720 100644 --- a/services/core/tests/integration/runtime_wake_hint_integration_test.go +++ b/services/core/tests/integration/runtime_wake_hint_integration_test.go @@ -113,8 +113,7 @@ func newWakeHintIntegration(t *testing.T) *wakeHintIntegration { } func wakeHintInput(text string) []sessions.Input { - payload, _ := json.Marshal(map[string]string{"text": text}) - return []sessions.Input{{Kind: "message", Payload: payload}} + return []sessions.Input{messageInput(text)} } func (f *wakeHintIntegration) pending(t *testing.T, target wakeHintIntegrationTarget, key string) sessions.EnvironmentInputReservation { diff --git a/services/core/tests/integration/runtime_worker_recovery_test.go b/services/core/tests/integration/runtime_worker_recovery_test.go index 654f8a444..0a8b08e51 100644 --- a/services/core/tests/integration/runtime_worker_recovery_test.go +++ b/services/core/tests/integration/runtime_worker_recovery_test.go @@ -2,7 +2,6 @@ package integration import ( "context" - "encoding/json" "errors" "testing" "time" @@ -35,7 +34,7 @@ func runtimeWorkerHarness(t *testing.T) (*dispatchHarness, *pgxpool.Pool) { func TestPreparedDispatchKeepsPendingReservationAfterComputeConflict(t *testing.T) { h, _ := runtimeWorkerHarness(t) h.d, h.lease = h.bound(), h.owner().Lease - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "pending", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "pending", []sessions.Input{messageInput("first")}) if err != nil { t.Fatal(err) } @@ -56,7 +55,7 @@ func TestPreparedDispatchKeepsPendingReservationAfterComputeConflict(t *testing. func TestWorkerWaitsForComputeAndSurvivesPromotionConflict(t *testing.T) { h, pool := runtimeWorkerHarness(t) - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "pending", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "pending", []sessions.Input{messageInput("first")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/sandbox_deployment_switch_test.go b/services/core/tests/integration/sandbox_deployment_switch_test.go index 49ee8ba6a..7c88d88c0 100644 --- a/services/core/tests/integration/sandbox_deployment_switch_test.go +++ b/services/core/tests/integration/sandbox_deployment_switch_test.go @@ -311,7 +311,7 @@ func TestSandboxSwitchPreservesReleasedAllocationAndItemHistory(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(t.Context(), s, tenant, history.ID, "history", json.RawMessage(`{"text":"retained request"}`)) + input, err := sendMessage(t.Context(), s, tenant, history.ID, "history", messageText("retained request")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/self_hosted_cancel_public_test.go b/services/core/tests/integration/self_hosted_cancel_public_test.go index 31f6f8989..bf9551100 100644 --- a/services/core/tests/integration/self_hosted_cancel_public_test.go +++ b/services/core/tests/integration/self_hosted_cancel_public_test.go @@ -117,7 +117,7 @@ func TestSelfHostedCancellationOfficialClient(t *testing.T) { return value } idleReceipts := receipts(created.IdleKey, "") - later, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, created.LaterID, "controlled-later-input", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"Retain pending input."}`)}}) + later, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), tenant, created.LaterID, "controlled-later-input", []sessions.Input{messageInput("Retain pending input.")}) if err != nil || later.State != sessions.EnvironmentInputPending || later.IsInitial { t.Fatal("could not establish controlled later reservation", err) } @@ -131,7 +131,7 @@ func TestSelfHostedCancellationOfficialClient(t *testing.T) { start := func() string { t.Helper() // Controlled callbacks isolate HTTP admission; no daemon or model runs in this fixture. - input, err := sendMessage(t.Context(), s, tenant, created.ID, uuid.NewString(), json.RawMessage(`{"text":"Controlled active work."}`)) + input, err := sendMessage(t.Context(), s, tenant, created.ID, uuid.NewString(), messageText("Controlled active work.")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/session_deletion_lifecycle_public_test.go b/services/core/tests/integration/session_deletion_lifecycle_public_test.go index 57216d5a3..440acbe68 100644 --- a/services/core/tests/integration/session_deletion_lifecycle_public_test.go +++ b/services/core/tests/integration/session_deletion_lifecycle_public_test.go @@ -45,7 +45,7 @@ func TestSessionDeletionLifecyclePostgres(t *testing.T) { input := sessions.CreateSession{Creator: FixtureCreator(), Engine: "codex", IdempotencyKey: uuid.NewString(), Configuration: json.RawMessage(`{` + deletionAgent + `,"environment":` + environment + `}`)} if initial { - input.InitialInputs = []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"reserved"}`)}} + input.InitialInputs = []sessions.Input{messageInput("reserved")} } session, err := s.CreateSession(ctx, tenant, input) if err != nil { @@ -59,7 +59,7 @@ func TestSessionDeletionLifecyclePostgres(t *testing.T) { turn := func(to ...string) string { t.Helper() session := create(none, false) - receipt, err := sendMessage(ctx, s, tenant, session.ID, "input", json.RawMessage(`{"text":"work"}`)) + receipt, err := sendMessage(ctx, s, tenant, session.ID, "input", messageText("work")) if err != nil { t.Fatal(err) } @@ -84,7 +84,7 @@ func TestSessionDeletionLifecyclePostgres(t *testing.T) { } reserve := func(session sessions.Session) sessions.EnvironmentInputReservation { t.Helper() - reservation, err := sessionService(t, s).ReserveEnvironmentInput(ctx, tenant, session.ID, "later", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"later"}`)}}) + reservation, err := sessionService(t, s).ReserveEnvironmentInput(ctx, tenant, session.ID, "later", []sessions.Input{messageInput("later")}) if err != nil || reservation.State != sessions.EnvironmentInputPending { t.Fatal(reservation, err) } diff --git a/services/core/tests/integration/session_deletion_test.go b/services/core/tests/integration/session_deletion_test.go index 1484c5553..b62dc8629 100644 --- a/services/core/tests/integration/session_deletion_test.go +++ b/services/core/tests/integration/session_deletion_test.go @@ -2,7 +2,6 @@ package integration import ( "context" - "encoding/json" "errors" "strings" "sync" @@ -62,7 +61,7 @@ func TestSessionDeletionWaitsForSettledTurnAndRejectsAdmission(t *testing.T) { if err != nil { t.Fatal(err) } - receipt, err := sendMessage(ctx, s, tenant, session.ID, "input", json.RawMessage(`{"text":"retained"}`)) + receipt, err := sendMessage(ctx, s, tenant, session.ID, "input", messageText("retained")) if err != nil { t.Fatal(err) } @@ -139,7 +138,7 @@ func TestSessionDeletionWaitsForSettledTurnAndRejectsAdmission(t *testing.T) { if _, err := createSession(ctx, fresh, tenant, input); !errors.Is(err, sessions.ErrIdempotencyConflict) { t.Fatal(err) } - if _, err := sendMessage(ctx, fresh, tenant, session.ID, "input", json.RawMessage(`{"text":"retained"}`)); !errors.Is(err, sessions.ErrNotFound) { + if _, err := sendMessage(ctx, fresh, tenant, session.ID, "input", messageText("retained")); !errors.Is(err, sessions.ErrNotFound) { t.Fatal(err) } if _, err := requestCancel(ctx, fresh, tenant, session.ID, "late-cancel"); !errors.Is(err, sessions.ErrNotFound) { diff --git a/services/core/tests/integration/session_events_test.go b/services/core/tests/integration/session_events_test.go index b68308acb..210b755ab 100644 --- a/services/core/tests/integration/session_events_test.go +++ b/services/core/tests/integration/session_events_test.go @@ -56,7 +56,7 @@ func TestSessionEventsCommitSnapshotsRetriesAndIsolation(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "start", json.RawMessage(`{"text":"question"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "start", messageText("question")) if err != nil { t.Fatal(err) } @@ -64,7 +64,7 @@ func TestSessionEventsCommitSnapshotsRetriesAndIsolation(t *testing.T) { if err != nil { t.Fatal(err) } - if _, err = sendMessage(ctx, s, tenant, session.ID, "start", json.RawMessage(`{"text":"question"}`)); err != nil { + if _, err = sendMessage(ctx, s, tenant, session.ID, "start", messageText("question")); err != nil { t.Fatal(err) } after, _ := sessionAdapter(s).SessionEventCursor(ctx, tenant, session.ID) @@ -166,7 +166,7 @@ func TestSessionEventsRetentionAndQueuedCancellation(t *testing.T) { }) inputs := make([]sessions.Input, 64) for i := range inputs { - inputs[i] = sessions.Input{Kind: "message", Payload: json.RawMessage(`{"text":"input"}`)} + inputs[i] = messageInput("input") } for range 5 { if _, err = submitInputs(ctx, s, tenant, session.ID, uuid.NewString(), inputs); err != nil { diff --git a/services/core/tests/integration/session_metadata_test.go b/services/core/tests/integration/session_metadata_test.go index 021a87b0c..2dad9bbb9 100644 --- a/services/core/tests/integration/session_metadata_test.go +++ b/services/core/tests/integration/session_metadata_test.go @@ -114,7 +114,7 @@ func TestSessionMetadataPreservesTerminalActivity(t *testing.T) { t.Fatal(err) } for _, status := range []string{sessions.TurnCompleted, sessions.TurnFailed, sessions.TurnCancelled} { - receipt, err := sendMessage(ctx, s, tenant, session.ID, uuid.NewString(), []byte(`{"text":"metadata fixture"}`)) + receipt, err := sendMessage(ctx, s, tenant, session.ID, uuid.NewString(), messageText("metadata fixture")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/stream_authority_http_test.go b/services/core/tests/integration/stream_authority_http_test.go index 33dd4666a..3ca7984c7 100644 --- a/services/core/tests/integration/stream_authority_http_test.go +++ b/services/core/tests/integration/stream_authority_http_test.go @@ -100,7 +100,7 @@ func TestLiveStreamClosesAfterKeyRevocationOrProjectArchive(t *testing.T) { t.Fatal("revocation affected peer", valid.StatusCode) } // New events remain available to valid callers after the reader has closed. - if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), project.TenantID, session.ID, "after-revocation", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"new event"}`)}}); err != nil { + if _, err := sessionService(t, s).ReserveEnvironmentInput(t.Context(), project.TenantID, session.ID, "after-revocation", []sessions.Input{messageInput("new event")}); err != nil { t.Fatal(err) } } diff --git a/services/core/tests/integration/token_usage_integration_test.go b/services/core/tests/integration/token_usage_integration_test.go index 6b533dc6b..434b4f869 100644 --- a/services/core/tests/integration/token_usage_integration_test.go +++ b/services/core/tests/integration/token_usage_integration_test.go @@ -33,7 +33,7 @@ func TestTokenUsageDurableSnapshotsAndSessionTotals(t *testing.T) { } } for n, status := range []string{sessions.TurnFailed, sessions.TurnCancelled} { - admission, err := sendMessage(ctx, s, tenant, session.ID, fmt.Sprint(n), json.RawMessage(`{"text":"measure"}`)) + admission, err := sendMessage(ctx, s, tenant, session.ID, fmt.Sprint(n), messageText("measure")) if err != nil { t.Fatal(err) } @@ -113,7 +113,7 @@ func TestCancellationReceiptUsageSurvivesRecovery(t *testing.T) { if err != nil { t.Fatal(err) } - admission, err := sendMessage(ctx, s, tenant, session.ID, "start", json.RawMessage(`{"text":"measure"}`)) + admission, err := sendMessage(ctx, s, tenant, session.ID, "start", messageText("measure")) if err != nil { t.Fatal(err) } @@ -187,7 +187,7 @@ func TestSessionUsageRequiresEveryRootTurnEndedAndMeasured(t *testing.T) { } submit := func(key string) sessions.InputReceipt { t.Helper() - admission, err := sendMessage(ctx, s, tenant, session.ID, key, json.RawMessage(`{"text":"measure"}`)) + admission, err := sendMessage(ctx, s, tenant, session.ID, key, messageText("measure")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/turn_events_test.go b/services/core/tests/integration/turn_events_test.go index e95bf1a40..5591fd3ac 100644 --- a/services/core/tests/integration/turn_events_test.go +++ b/services/core/tests/integration/turn_events_test.go @@ -20,7 +20,7 @@ func TestTurnEventBatchesAreOrderedIsolatedAndDurable(t *testing.T) { if err != nil { t.Fatal(err) } - input, err := sendMessage(ctx, s, tenant, session.ID, "start", json.RawMessage(`{"text":"test"}`)) + input, err := sendMessage(ctx, s, tenant, session.ID, "start", messageText("test")) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/worker_preparation_failure_test.go b/services/core/tests/integration/worker_preparation_failure_test.go index d1d266352..9d1223560 100644 --- a/services/core/tests/integration/worker_preparation_failure_test.go +++ b/services/core/tests/integration/worker_preparation_failure_test.go @@ -1,7 +1,6 @@ package integration import ( - "encoding/json" "testing" "time" @@ -15,7 +14,7 @@ func TestWorkerSettlesConfirmedPreparationFailureAndAcceptsNewInput(t *testing.T h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), false) enableWorkerEnvironment(t, h) frames := workerFrames(t, h) - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "first", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "first", []sessions.Input{messageInput("first")}) if err != nil { t.Fatal(err) } @@ -58,7 +57,7 @@ func TestWorkerSettlesConfirmedPreparationFailureAndAcceptsNewInput(t *testing.T t.Fatal("failed input retried", frame.Type) case <-time.After(1200 * time.Millisecond): } - next, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "next", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"next"}`)}}) + next, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "next", []sessions.Input{messageInput("next")}) if err != nil { t.Fatal("new input remained blocked", err) } @@ -98,7 +97,7 @@ func TestWorkerRetriesUncertainPreparationFailure(t *testing.T) { h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), false) enableWorkerEnvironment(t, h) frames := workerFrames(t, h) - pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "retry", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"retry"}`)}}) + pending, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "retry", []sessions.Input{messageInput("retry")}) if err != nil { t.Fatal(err) } @@ -130,7 +129,7 @@ func TestWorkerPreparationRejectionPreservesCancellationAndNewerInput(t *testing h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), false) enableWorkerEnvironment(t, h) frames := workerFrames(t, h) - first, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "first", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + first, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "first", []sessions.Input{messageInput("first")}) if err != nil { t.Fatal(err) } @@ -140,7 +139,7 @@ func TestWorkerPreparationRejectionPreservesCancellationAndNewerInput(t *testing if _, err := cancelEnvironmentInput(t.Context(), h.s, h.tenant, h.session.ID, first.ID); err != nil { t.Fatal(err) } - next, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "next", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"next"}`)}}) + next, err := sessionService(t, h.s).ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "next", []sessions.Input{messageInput("next")}) if err != nil { t.Fatal(err) } diff --git a/services/core/tests/integration/worker_wakeup_test.go b/services/core/tests/integration/worker_wakeup_test.go index 3596bbff1..3f4543fb6 100644 --- a/services/core/tests/integration/worker_wakeup_test.go +++ b/services/core/tests/integration/worker_wakeup_test.go @@ -2,7 +2,6 @@ package integration import ( "context" - "encoding/json" "strings" "sync/atomic" "testing" @@ -60,7 +59,7 @@ func TestWorkerSchedulerCommittedAdmissionWakesBeforeMaintenance(t *testing.T) { t.Error("worker did not stop") } }() - input := []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"wake"}`)}} + input := []sessions.Input{messageInput("wake")} switch operation { case "submit": _, err = worker.SubmitInputs(ctx, h.tenant, h.session.ID, "wake", input) @@ -112,7 +111,7 @@ func TestWorkerSchedulerHintBypassesEnvironmentScanThrottle(t *testing.T) { admitted := make(chan error, 1) started := time.Now() go func() { - _, err := worker.SubmitInputs(ctx, h.tenant, h.session.ID, "wake", []sessions.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"wake"}`)}}) + _, err := worker.SubmitInputs(ctx, h.tenant, h.session.ID, "wake", []sessions.Input{messageInput("wake")}) admitted <- err }() prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare)