From f633d2cae7a4b2beb6b8a0481253e145e77dad0e Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Wed, 23 Sep 2026 12:01:36 +0000 Subject: [PATCH] Settle Session creation streams like the official service Close fresh creation streams right after the first settling idle or failed event, with a same-snapshot fallback for settlements that record no event; end same-key stream retries immediately; send the committed JSON 201 projection as the created snapshot; carry Turn usage on terminal Turn events; and record new Turns as turn.created, user items, session.in_progress, turn.in_progress. Squashed from the reviewed branch head 80f7a1703aa6021e7ff77545d7ce2dfaf5e82305 onto main 4f3822a. --- CONTRIBUTING.md | 58 +- apps/web/e2e/fixture-core.mjs | 24 +- contracts/agents-api/README.md | 55 +- contracts/agents-api/environments.md | 12 +- contracts/agents-api/history-events-usage.md | 91 +++ contracts/agents-api/openapi.yaml | 109 +-- contracts/agents-api/operation-evidence.md | 7 +- contracts/agents-api/v1/events.go | 27 + contracts/agents-api/v1/events_test.go | 44 ++ packages/agents-client/src/client.test.ts | 145 +++- packages/agents-client/src/client.ts | 35 +- packages/agents-client/src/index.ts | 2 +- packages/agents-client/src/types.ts | 5 + services/agents-api/internal/api/handler.go | 2 +- .../internal/api/session_creation_identity.go | 3 +- .../internal/api/session_creation_stream.go | 61 +- .../api/session_creation_stream_test.go | 726 ++++++++++++++++++ services/agents-api/internal/api/stream.go | 144 +++- .../agents-api/internal/api/stream_test.go | 49 ++ .../creation_stream_settlement_public_test.go | 242 ++++++ .../store/environment_initial_input_test.go | 35 +- .../store/environment_input_activity.go | 24 +- .../store/environment_input_activity_test.go | 18 +- .../internal/store/function_state.go | 4 +- .../agents-api/internal/store/scheduling.go | 93 ++- .../internal/store/session_creation_stream.go | 20 +- .../store/session_creation_stream_test.go | 59 +- .../session_environment_snapshot_test.go | 6 +- .../internal/store/session_events.go | 9 +- .../agents-api/internal/store/sessions.go | 4 + .../agents-api/internal/store/turn_inputs.go | 9 + .../tests/official_agent_reference_retry.py | 25 +- .../agents-api/tests/official_execution.py | 3 + .../tests/official_self_hosted_initial.py | 23 +- .../tests/official_session_creation_stream.py | 110 ++- 35 files changed, 2026 insertions(+), 257 deletions(-) create mode 100644 contracts/agents-api/v1/events_test.go create mode 100644 services/agents-api/internal/api/session_creation_stream_test.go create mode 100644 services/agents-api/internal/store/creation_stream_settlement_public_test.go diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 6884a915e..26f9bc3cf 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -973,10 +973,11 @@ lock, and return terminal storage outcomes without rolling their transaction bac Initial messages for a newly created Environment-bearing Session use that same reservation in the creation transaction, including its connection-action event. -The creation winner alone inserts it; the original pre-work snapshot and stream -cursor remain unchanged. A durable initial/later flag defaults historical rows to -later input without inferring origin. Initial expiry projects a failed Session and -safe error before any Turn exists; later expiry retains idle semantics. Failure +The creation winner alone inserts it; the stream cursor still precedes that +event, and creation retries never re-insert it. A durable initial/later flag +defaults historical rows to later input without inferring origin. Initial expiry +projects a failed Session and safe error before any Turn exists; later expiry +retains idle semantics. Failure events capture the settled activity and Usage atomically. Late connections and creation retries cannot reset or replay expired input, and newer work supersedes old activity without changing its event snapshots. The Environment itself is not @@ -987,8 +988,9 @@ immediate Turn admission. Cancellation/deletion keep their existing semantics. Ordinary and streamed public self-hosted creation accept initial text through this transaction after configuration and new-work lease checks. They return the owned Environment ID and executor URL while offline, without waiting for admission. -The creation stream sends its original pre-work `created` snapshot before the -committed connection action. A disconnected observer leaves committed input intact; +The creation stream's `created` snapshot is the committed JSON 201 projection, +including the connection action; the committed action event then follows from the +creation cursor. A disconnected observer leaves committed input intact; only the existing Worker prepares, promotes and starts it. Saved-Agent retries with recorded intent recover before fresh execution admission or source resolution; inline retries keep their existing resolved-snapshot validation. @@ -2073,19 +2075,45 @@ replaced; do not carry obsolete compatibility code forward to satisfy this secti Missing sequence positions produce a safe stream error and close; recover via Session/Turn/Items queries. Socket writes have a five-second deadline and hold no database connection. Client disconnect releases the handler; comments keep - idle connections alive. SSE does not close merely because one Turn finishes. + idle connections alive. GET SSE does not close merely because one Turn finishes; + only creation responses end on settlement (below). - Session creation with `stream=true` reuses atomic input admission and the live event loop. The upsert returns its cursor under the Session lock, before initial inputs; never replace it with a post-commit cursor lookup. A new response emits - one request-local `agent.session.created` with the pre-input resource snapshot, - then committed changes from that cursor. The local creation retry key excludes - response mode: retries observe only later events and admit no work again. Retry - the same request/key with `stream=false` to recover a lost Session ID. GET event - streams retain their current live-only start. Disconnect never cancels admitted - work. Exact upstream created-snapshot timing, POST stream lifetime and creation - retry response semantics remain unverified; the separate SDK one-Turn helper - does not define this endpoint. Do not present local retry behavior as replay. + one request-local `agent.session.created` with the committed Session projection + that the JSON 201 response returns (read after the commit), then committed + changes from that cursor exactly once. A fresh creation stream ends right + after the first `agent.session.idle` recorded when a Turn ends or an input + reservation stops being pending (expired, cancelled or failed), or any + `agent.session.failed`, and never sends the events after it. A self-hosted + connection clearing pending input to idle, `requires_action`, function results + and resumed work keep it open. A creation that admitted nothing (no Turn or + reservation) ends right after `created`. Settlements that record no event use + a fallback: after an empty drain the stream reads the JSON-path projection and + the event cursor in one database snapshot and, if the Session is idle or failed + with no queued, running or waiting Turn and no pending reservation, sends only + events up to that cursor, then ends. Accepted follow-ups: another client's work + drained before that read can still be sent, and idles recorded by an older + binary during a rolling deploy carry no settled marker and rely on the + fallback. An input reservation made while the ending Turn captured Artifacts + can start a later Turn that the stream does not follow. The settled marker and + pending-input flag are Store-internal, never wire fields, and add no events. + Re-read the projection after a sent Session status event and otherwise at most + once a second. The local creation retry key excludes response mode; a same-key + `stream=true` retry of an existing creation returns 201 with only the + connection comment and ends at once, admitting nothing and following no work, + because official same-key requests create distinct Sessions. Retry the same + request/key with `stream=false`, or use the GET events stream, to recover. GET + event streams keep their live-only start and never end on settlement. + Disconnect never cancels admitted work. Official observations cover `none` + creation; self-hosted, hosted and no-input stream lifetimes and the retry + behavior are local choices, and the separate SDK one-Turn helper does not + define this endpoint. Do + not present local retry behavior as replay. A new Turn records `turn.created`, + its user input Items, then Session activity in one transaction. Terminal Turn + events carry top-level `usage` copied from their Turn snapshot, null when + unknown; never derive or sum it. - Public Turn retrieve/list project persisted execution state and the immutable Session Agent identity. Scope both resources and pagination cursors to the diff --git a/apps/web/e2e/fixture-core.mjs b/apps/web/e2e/fixture-core.mjs index 4f5b1a67f..44fb8f85a 100644 --- a/apps/web/e2e/fixture-core.mjs +++ b/apps/web/e2e/fixture-core.mjs @@ -592,6 +592,7 @@ function emitTurnLifecycle(status) { session_id: "session_snapshot", turn_id: terminal.id, turn: terminal, + usage: terminal.usage, })}\n\n`; for (const stream of streamResponses.keys()) stream.write(event); return true; @@ -1065,8 +1066,9 @@ const server = http.createServer(async (request, response) => { "cache-control": "no-cache, no-transform", connection: "keep-alive", }); - response.write(": fixture creation retry observes only future events\n\n"); - setTimeout(() => response.end(), state.controls.sessionCreateStreamCloseDelayMs); + // Like Core, a same-key stream retry sends no events and ends at once; + // clients recover the Session with stream=false. + response.end(": connected\n\n"); return; } return sendJson(response, receipt.session, 201); @@ -1121,7 +1123,6 @@ const server = http.createServer(async (request, response) => { created_at: baseline + state.sequence, last_active_at: baseline + state.sequence, }; - const createdSnapshot = structuredClone(created); let initialTurn = null; let initialItems = []; if (hasInitialInput && initialInputMessages) { @@ -1153,6 +1154,11 @@ const server = http.createServer(async (request, response) => { created.status = "in_progress"; created.last_active_at = baseline + state.sequence; } + // As in Core, the created event repeats this fixture's JSON 201 body. The + // fixture queues a Turn for any initial input, so both show in_progress; + // Core instead shows requires_action for self_hosted and idle while an + // openai_hosted Environment provisions. + const createdSnapshot = structuredClone(created); state.sessions.unshift(created); if (typeof idempotencyKey === "string") { state.sessionCreateReceipts.set(idempotencyKey, { fingerprint, session: created }); @@ -1188,12 +1194,6 @@ const server = http.createServer(async (request, response) => { turn_id: turn.id, turn, })}\n\n`); - response.write(`event: agent.session.in_progress\nid: progress_${state.sequence}\ndata: ${JSON.stringify({ - type: "agent.session.in_progress", - event_id: `progress_${state.sequence}`, - session_id: created.id, - session: created, - })}\n\n`); for (const [index, item] of items.entries()) { response.write(`event: agent.session.turn.item.added\nid: item_${state.sequence}_${index + 1}\ndata: ${JSON.stringify({ type: "agent.session.turn.item.added", @@ -1203,6 +1203,12 @@ const server = http.createServer(async (request, response) => { item, })}\n\n`); } + response.write(`event: agent.session.in_progress\nid: progress_${state.sequence}\ndata: ${JSON.stringify({ + type: "agent.session.in_progress", + event_id: `progress_${state.sequence}`, + session_id: created.id, + session: created, + })}\n\n`); } setTimeout(() => response.end(), state.controls.sessionCreateStreamCloseDelayMs); return; diff --git a/contracts/agents-api/README.md b/contracts/agents-api/README.md index 68a02be7e..5d7c1661a 100644 --- a/contracts/agents-api/README.md +++ b/contracts/agents-api/README.md @@ -776,27 +776,50 @@ express the string/array union, so input is unconstrained with a type descriptio ### Session creation streaming `POST /v1/agents/sessions` also accepts `stream=true` for the supported creation -inputs. Fresh creation sends `agent.session.created` with the pre-input Session, -then its committed activity/Turn/Item/output events. Self-hosted initial creation -first requests the Environment connection, before native readiness and a Turn. -The cursor comes -from the atomic creation upsert, so rapid initial execution cannot move the start -past its own events. The ordinary bounded-buffer/gap policy still applies. +inputs. Fresh creation sends `agent.session.created` with the committed Session, +the same projection as the JSON 201 body (`in_progress` after `none` initial +input), then its committed activity/Turn/Item/output events. A new Turn publishes +`turn.created`, its user input `item.added`, `agent.session.in_progress`, then +`turn.in_progress`. Self-hosted initial creation shows and then emits the +Environment connection request, before native readiness and a Turn. The cursor +comes from the atomic creation upsert, so rapid initial execution cannot move the +start past its own events. The ordinary bounded-buffer/gap policy still applies. + +A fresh creation stream ends right after the first `agent.session.idle` recorded +when a Turn ends or an input reservation stops being pending (expired, cancelled +or failed), or any `agent.session.failed`, and never sends the events after it. A +pinned-SDK loop over `sessions.create(..., stream=True)` therefore ends right +after the initial Turn's idle. `requires_action`, function results, resumed work +and a self-hosted connection clearing pending input keep it open, and a +provisioning or offline reservation keeps it open until a Turn settles or the +reservation expires or fails. A creation that admitted nothing ends right after +`created`. A settlement that records no event ends the stream after the events +up to the cursor read with a settled projection in one snapshot; another client's +work drained before that read can still be sent. Observe later Turns with the GET +event stream, which never ends on its own. Terminal Turn events carry the Turn +snapshot's `usage` at the top level, null when unknown. The local `Idempotency-Key` creation extension shares identity across response -modes. Retrying creation streams only future changes and never resubmits input or -replays old events. Recover a lost Session ID by repeating the same request/key -with `stream=false`, then use Session/Turn/Items reads. Disconnect only stops the -HTTP observer; committed reservations and admitted execution continue. Idle streams remain open for later -Turns. Pinned SDK3.13.0 proves the creation stream and created-event schema; exact -upstream initial snapshot/order, POST stream lifetime and retry behavior have not -been compared with the hosted service. These choices are not full conformance. +modes. A same-key `stream=true` retry of an existing creation returns 201 with +only the connection comment and ends at once: it replays nothing, resubmits no +input and follows no work, since official same-key requests create distinct +Sessions. Recover a lost +Session ID by repeating the same request/key with `stream=false`, then use +Session/Turn/Items reads. Disconnect only stops the HTTP observer; committed +reservations and admitted execution continue. September 23 official `none` +observations match the created snapshot, the end at idle, the start order and the +terminal usage field ([evidence](history-events-usage.md#creation-stream-settlement-2026-09-23)). +Self-hosted, hosted and no-input creation stream lifetimes and the stream retry +behavior are local choices. +These choices are not full conformance. `official_session_creation_stream.py` covers the pinned client and raw HTTP on -real PostgreSQL: idle/initial text and saved Agents, first snapshots and ordered -Items, retries across response modes, later Turns, disconnect recovery, isolation +real PostgreSQL: initial text and saved Agents, created snapshots equal to the JSON +201 body, Turn start order, terminal usage, the end at idle, GET continuation, +immediately ending stream retries and JSON retries, disconnect recovery, isolation and errors before stream headers. Store tests cover concurrent upsert ownership, -pre-admission cursors and observers draining after execution has completed. +pre-admission cursors, post-admission projections and observers draining after +execution has completed; API tests cover the stream lifetimes. ## Acceptance evidence and remaining scope diff --git a/contracts/agents-api/environments.md b/contracts/agents-api/environments.md index f515682ff..210f98f35 100644 --- a/contracts/agents-api/environments.md +++ b/contracts/agents-api/environments.md @@ -82,9 +82,15 @@ the exact local profile are validated before persistence. Supported optional functions remain engine-specific. Initial text commits a reservation and connection action, then returns the Session -and Environment connection target while offline. Streamed creation sends its -original `created` snapshot before the connection action. The existing Worker -prepares and admits the input; closing the stream leaves committed work intact. +and Environment connection target while offline. Streamed creation sends the same +committed projection as its `created` snapshot, already showing `requires_action` +and the connection action, then the committed `requires_action` event. Like +every fresh creation stream it ends right after the idle recorded when the +admitted Turn ends or the reservation stops being pending, or after a failure; +the connection alone clearing the action does not end it. A creation without +input ends right after its created snapshot, and a same-key stream retry ends at +once without events. The existing Worker prepares and admits the input; closing the stream +leaves committed work intact. Initial expiry leaves a failed Session, safe error and empty actions without a Turn or an Environment failure. Creation retries preserve the original identity, deadline and input. Later live observers do not replay creation events. diff --git a/contracts/agents-api/history-events-usage.md b/contracts/agents-api/history-events-usage.md index 9e1a42451..d116fba88 100644 --- a/contracts/agents-api/history-events-usage.md +++ b/contracts/agents-api/history-events-usage.md @@ -116,3 +116,94 @@ Sources: [Sessions](https://developers.openai.com/api/docs/guides/agents-api/ses Sanitized request evidence is retained privately under `~/.parsar/remediation/20260922/history-events-usage/official/`; credentials are excluded from source and evidence. + +## Creation stream settlement, 2026-09-23 + +Evidence: the second official-semantics campaign scan compared four owned +`environment:none` official Sessions with Core (private +`~/.parsar/remediation/20260923/campaign-scan-2/events-tools/`, `findings.json` +EVT-01..24 with raw frames under `official/`). All 1091 official events passed +strict validation against the pinned types. Two independent official creation +streams (structured output, and a function call with its result) were closed by +the server; they, a third creation stream that the client closed while a call +was pending, and the 2026-09-22 probe above agree on the snapshot, order and +terminal-event observations below. This batch changes only the four Core-owned +stream differences EVT-01..04; the plan is +`~/.parsar/remediation/20260923/creation-stream-settlement/PLAN.md`. + +- **Creation stream lifetime (EVT-01).** The official service closed both creation + streams right after the first `agent.session.idle`, without `[DONE]` or an error + frame, and kept the function stream open through `requires_action` and the + result. A fresh Core creation stream ends right after the first + `agent.session.idle` recorded when a Turn ends or an input reservation stops + being pending (expired, cancelled or failed), or any `agent.session.failed`, + and never sends the events after it. A self-hosted connection that clears + pending input to idle does not end it. A creation that admitted nothing ends + right after `created`. Settlements that record no event, such as a reservation + cancelled while its Session is already idle, use a fallback: after an empty + drain the stream reads the JSON-path projection and the event cursor in one + database snapshot and, if settled, sends only events up to that cursor, then + ends. The settled marker and pending-input flag are Store-internal and add no + events. A same-key `stream=true` retry of an existing creation returns 201 with + only the connection comment and ends at once; official same-key requests create + distinct Sessions, so there is no retry stream to follow, and recovery uses + `stream=false` or GET. The TypeScript client reports that empty creation stream + as `CreationStreamRetryError`. Accepted follow-ups: another client's work + drained before the fallback read can still be sent after a silent settlement; + idles recorded by an older binary during a rolling deploy carry no settled + marker and rely on the fallback; and an input reservation made while the ending + Turn captured Artifacts can start a later Turn that the stream does not follow. + Four review rounds replaced event-only, projection-only and retry-following + designs before merge. GET event streams are unchanged: live-only, no replay, + and they never end on their own. Session deletion still ends both. +- **Created snapshot (EVT-02).** Official `agent.session.created` carried the + post-admission Session (`in_progress`, no actions, null usage), like the JSON + 201 body. Core now sends the committed projection that JSON 201 returns, read + after the creation commit, while the stream still starts at the creation + upsert cursor, so every initial Turn and Item event follows exactly once. The + snapshot is read after the commit, so it can already show a later state than + the events that follow it; the JSON 201 body has the same race. For + self-hosted input the snapshot already requests the Environment connection and + the committed `requires_action` event follows. Hosted initial input remains + `idle` while it provisions. +- **Terminal usage (EVT-03).** Official `agent.session.turn.completed` and + `.cancelled` (8/8) carried a top-level `usage`, null at emission even when later + reads were measured. Core terminal Turn events (`completed`, `failed`, + `cancelled`), root and child, now carry `usage` copied from the rendered Turn + snapshot, with explicit null when unknown. Other events omit it. Codex can + therefore publish measured counters at settlement, while Claude and MiniMax + stay null; no counter is derived or summed. The TypeScript client accepts the + field on terminal Turn events only and still accepts older events without it. +- **Turn start order (EVT-04).** Official new Turns published `turn.created`, the + user `item.added` (`output_index` null), `agent.session.in_progress`, then + `turn.in_progress`. Core now records the Session activity after the admitting + input's Items in the same transaction, for creation, events.create and + reservation promotion. When one batch holds several message events, later + messages follow that activity; their official order was not observed. + +Deferred, with evidence retained in `findings.json`: + +- EVT-05: Core emits `turn.in_progress` when a function result resumes a waiting + Turn and publishes the result Item at native application; the latter is the + known INTERACTION-PUBLICATION-001 receipt boundary. +- EVT-06: Core emits an interim `agent.session.in_progress` when cancelling a + Turn that waits on a function result. Statuses match. +- EVT-07: official mid-Turn attach sent catch-up Item snapshots, but not + deterministically; one more official sample is needed before designing. +- EVT-08: official Items omit in-progress and incomplete output Items; Core keeps + them under native history ownership. +- EVT-09 and EVT-10: null-valued `output_index`, `phase` and `error` fields and + the initial assistant content belong to a separate serialization batch. +- EVT-11 and EVT-12: unknown call/Turn result and conflict error codes belong to + ERROR-PROTOCOL-001. +- EVT-13: official Session usage became null when any root Turn usage was + unknown; Core sums the known Turns. The in-progress case is unverified. +- EVT-19: the Core terminal sequence for a Turn cancelled mid-text is recorded by + this batch's live acceptance, not changed by it. + +Batch validation used targeted Go API, store and contract tests on a dedicated +PostgreSQL database, including the pinned Python SDK 3.13.0 creation-stream, +initial-input, self-hosted and initial-failure scripts, plus the TypeScript +client and Core Web unit tests. Resource-level replay, live model acceptance and +the full gate are recorded separately; retry, self-hosted, hosted and no-input +creation stream lifetimes have no official observation. diff --git a/contracts/agents-api/openapi.yaml b/contracts/agents-api/openapi.yaml index 3291782b6..6105d895c 100644 --- a/contracts/agents-api/openapi.yaml +++ b/contracts/agents-api/openapi.yaml @@ -1508,6 +1508,13 @@ definitions: type: string type: type: string + usage: + allOf: + - $ref: '#/definitions/v1.TokenUsage' + description: |- + Usage is present only on terminal Turn events, where it mirrors the Turn + snapshot and is null when unknown. Other events omit it. + x-nullable: true required: - event_id - type @@ -3028,54 +3035,60 @@ paths: initial timeout. Initial input is required for none and for streamed creation outside self_hosted. Omitted/null input remains valid for non-streaming hosted and self_hosted creation. With stream=true, returns live Session events starting - at creation; disconnect does not cancel execution. New Sessions retain their - authenticated creator; all creation retries require the same typed subject, - including across key rotation. Saved-Agent retries and inline requests using - Vault attachments or credential references retain caller intent independently - of later resource changes; unrelated inline retries preserve resolved/default - equivalences. Unknown historical creators reject retries; known creators without - recorded intent retain resolved-snapshot retry rules. These conflict policies - are local and not verified hosted parity. Creation retries observe future - events without replay; retry with stream=false to retrieve the Session. Claude - SDK on none and Core-managed Docker openai_hosted supports qualified object-root - json_schema output with medium verbosity, single-Agent execution and ordinary - functions. Hosted execution reuses native workspace tools and Files/Artifacts; - Skills, Plugins, capability directories, HTTP MCP, Subagent and tool_search - combinations remain unqualified, including inherited template contents. Other - non-text initial input remains unsupported. Basic Codex and Claude SDK openai_hosted - creation requires an explicitly configured managed provider. The Claude workspace - profile supports non-deferred function tools with text or successful inline - PNG/JPEG results alongside native workspace tools; HTTP MCP remains unsupported. - Idle Sessions provision automatically; initial provisioning has no caller - connection action. Network defaults to enabled; disabled and restricted exact - ASCII hostnames are supported. Restricted policy requires 1–100 allowed domains. - Unsupported hostname forms and startup installations are rejected. Confidential - env, system/npm/Python packages and ordered setup commands use the shared - initialization lifecycle; requested network applies after setup. Initial inline - and tenant-owned file_id files freeze encrypted bytes before provisioning, - then install through the common Core lifecycle before native execution or - live Files access. With a template reference, omitted/null files, env, packages - and setup_commands inherit. Non-null files and command lists replace; env - overlays by key; each package manager inherits on omission/null and otherwise - replaces its list. Empty lists clear their selected field. Tenant-owned environment_template_id - references inherit omitted/null network and allow only narrowing overrides. - Inline hosted network:null retains the enabled default; updating a Template - with network:null resets its saved policy to enabled. Core freezes effective - configuration; template updates/deletion do not alter Session snapshots or - same-intent creation retries. Inline or tenant-owned skill_reference Skills - share initialization. Templates preserve default/latest/explicit selectors; - Session creation freezes concrete metadata and encrypted content atomically. - Skill, Plugin and capability-directory list omission/null inherit; a non-null - list replaces, including empty-list clearing. Omitted/null Skill version selectors - resolve the default version. Source deletion/default updates cannot change - committed Session Skill contents. Deferred function discovery uses type-only - tool_search and per-function defer_loading in the qualified single-agent Claude - environment:none function profile, including qualified inline image messages - and text results. Explicit web_search mode disabled and programmatic_tool_calling - enabled false use frozen common Runtime controls. Enabled forms remain unqualified. - Omitted programmatic configuration preserves native behavior, a documented - difference from the official default-on behavior. Other combinations remain - unqualified; see the operation coverage. + with the committed creation snapshot and closes right after the first agent.session.idle + recorded when a Turn ends or an input reservation stops being pending, or + any agent.session.failed, without sending later events. A creation that admitted + nothing closes after the snapshot; a settlement that records no event closes + after events up to the cursor read with a settled Session projection. Required + actions keep it open; disconnect does not cancel execution. The GET events + stream remains live-only. New Sessions retain their authenticated creator; + all creation retries require the same typed subject, including across key + rotation. Saved-Agent retries and inline requests using Vault attachments + or credential references retain caller intent independently of later resource + changes; unrelated inline retries preserve resolved/default equivalences. + Unknown historical creators reject retries; known creators without recorded + intent retain resolved-snapshot retry rules. These conflict policies are local + and not verified hosted parity. A same-key stream=true retry of an existing + creation returns 201 with no events and closes at once; retry with stream=false + or use the GET events stream to recover. Claude SDK on none and Core-managed + Docker openai_hosted supports qualified object-root json_schema output with + medium verbosity, single-Agent execution and ordinary functions. Hosted execution + reuses native workspace tools and Files/Artifacts; Skills, Plugins, capability + directories, HTTP MCP, Subagent and tool_search combinations remain unqualified, + including inherited template contents. Other non-text initial input remains + unsupported. Basic Codex and Claude SDK openai_hosted creation requires an + explicitly configured managed provider. The Claude workspace profile supports + non-deferred function tools with text or successful inline PNG/JPEG results + alongside native workspace tools; HTTP MCP remains unsupported. Idle Sessions + provision automatically; initial provisioning has no caller connection action. + Network defaults to enabled; disabled and restricted exact ASCII hostnames + are supported. Restricted policy requires 1–100 allowed domains. Unsupported + hostname forms and startup installations are rejected. Confidential env, system/npm/Python + packages and ordered setup commands use the shared initialization lifecycle; + requested network applies after setup. Initial inline and tenant-owned file_id + files freeze encrypted bytes before provisioning, then install through the + common Core lifecycle before native execution or live Files access. With a + template reference, omitted/null files, env, packages and setup_commands inherit. + Non-null files and command lists replace; env overlays by key; each package + manager inherits on omission/null and otherwise replaces its list. Empty lists + clear their selected field. Tenant-owned environment_template_id references + inherit omitted/null network and allow only narrowing overrides. Inline hosted + network:null retains the enabled default; updating a Template with network:null + resets its saved policy to enabled. Core freezes effective configuration; + template updates/deletion do not alter Session snapshots or same-intent creation + retries. Inline or tenant-owned skill_reference Skills share initialization. + Templates preserve default/latest/explicit selectors; Session creation freezes + concrete metadata and encrypted content atomically. Skill, Plugin and capability-directory + list omission/null inherit; a non-null list replaces, including empty-list + clearing. Omitted/null Skill version selectors resolve the default version. + Source deletion/default updates cannot change committed Session Skill contents. + Deferred function discovery uses type-only tool_search and per-function defer_loading + in the qualified single-agent Claude environment:none function profile, including + qualified inline image messages and text results. Explicit web_search mode + disabled and programmatic_tool_calling enabled false use frozen common Runtime + controls. Enabled forms remain unqualified. Omitted programmatic configuration + preserves native behavior, a documented difference from the official default-on + behavior. Other combinations remain unqualified; see the operation coverage. parameters: - description: agents=v1 in: header diff --git a/contracts/agents-api/operation-evidence.md b/contracts/agents-api/operation-evidence.md index a8e104e55..eb31f96f5 100644 --- a/contracts/agents-api/operation-evidence.md +++ b/contracts/agents-api/operation-evidence.md @@ -1,6 +1,6 @@ # Pinned operation evidence inventory — 2026-09-23 -Baseline inventory of main `b5715912f09333e2b4449ec6f0eecaabce44c9b7`. The Session admission batch below updates creation and metadata validation, the list query tolerance batch (L) updates list and resource query handling, the validation error batch (X) updates field error codes/params, malformed path IDs and U+0000 handling, and the Environment Files wire batch (G) updates Files.create/list status, envelope, query, path and empty-page behavior; historical evidence retains its original revision and scope. This inventory guides repeated qualification and does not assert complete compatibility. +Baseline inventory of main `b5715912f09333e2b4449ec6f0eecaabce44c9b7`. The Session admission batch below updates creation and metadata validation, the list query tolerance batch (L) updates list and resource query handling, the validation error batch (X) updates field error codes/params, malformed path IDs and U+0000 handling, the Environment Files wire batch (G) updates Files.create/list status, envelope, query, path and empty-page behavior, and the creation stream settlement batch (J) updates creation SSE lifetime/snapshot, terminal Turn usage and Turn start order; historical evidence retains its original revision and scope. This inventory guides repeated qualification and does not assert complete compatibility. Baseline: `contracts/agents-api/upstream.json`, SDK **3.13.0**, upstream commit **d7c41efee1b0802b79f3f88a678ef2052b06e9ce**, `OpenAI-Beta: agents=v1`. AGENTS.md and relevant CONTRIBUTING.md compatibility, ownership and evidence rules govern this inventory. @@ -37,6 +37,7 @@ Repository paths below are relative to the inspected worktree; private evidence | V | [Resource selector and error qualification](resource-selector-semantics.md); private September23 `skill-version-alignment/official/` and `files-error-alignment/official-probe.json`: owned Skill/Template/Session selector observations and seven missing File/cursor requests. Preserve each probe phase and distinguish actual hosted execution from metadata-only reads. | | W | `~/.parsar/remediation/20260923/file-resource-semantics/official-files/` and `official-skills/`: 56 owned resource requests, no Sessions/models; [qualified observations and limitations](file-resource-semantics.md). | | X | [Validation error fields](official-semantics-alignment.md#validation-error-fields--september-23); private `~/.parsar/remediation/20260923/campaign-scan-1/{vaults-agents,sessions,skills-files-templates}/findings.json` VA-07/08/09/10, SES-28, SFT-20 and the September 22 Session `metadata.a` null observation. Metadata/name field errors, U+0000 local limit, malformed path IDs and Template network codes; Go handler and real-PostgreSQL route/no-write tests, no model execution. | +| J | [Creation stream settlement](history-events-usage.md#creation-stream-settlement-2026-09-23); private `~/.parsar/remediation/20260923/campaign-scan-2/events-tools/findings.json` EVT-01..04 with raw frames under `official/`: four owned official `none` Sessions (text, structured output, function result, cancellation), 1091 events validated against the pinned types. Core-side changes carry API/store/contract DB tests; live acceptance is recorded with the batch. EVT-05..13 and EVT-19 stay deferred or record-only. | | L | [List query tolerance](list-query-semantics.md#list-query-tolerance--september-23-2026); private `~/.parsar/remediation/20260923/campaign-scan-1/{vaults-agents,sessions,skills-files-templates}/findings.json`: owned-collection unknown/repeated keys, limit bounds, Vault status union and Files empty purpose, plus unknown keys on a deleted Vault read and Agent delete. Rows A1–D2 of that section; no model execution. | | G | [Environment Files wire alignment](environment-files.md#wire-alignment--september-23-2026); private `~/.parsar/remediation/20260923/campaign-scan-2/hosted-env/findings.json` HE-10, 16, 18, 32, 34–39 with raw records under `official/` (labels `fc01`–`fc17`, `fl01`–`fl16`, `files-*-pending`) and `run1/`: three owned hosted Sessions, all deleted; first official Environment Files observations. Rows F1–F9 of that section. Go handler, real-PostgreSQL Worker, gateway/daemon and Rust helper tests without a model; live acceptance is recorded with the batch. | | Y | [Artifact capture and listing](official-semantics-alignment.md#artifact-capture-and-listing--september-23); private `~/.parsar/remediation/20260923/campaign-scan-2/hosted-env/findings.json` HE-50..62 with raw records under `official/` (labels `al01`–`al09`, `ar01`–`ar04`, `ac01`–`ac05`, `ad01`/`ad02`): three owned Sessions and two tiny Turns, all deleted; first official Artifact observations. Rows A1–A4 of that section: symlink skip, republication, list envelope and malformed filter. Rust link tests, real-PostgreSQL store/HTTP and pinned-SDK tests without a model; live acceptance is recorded with the batch. | @@ -52,13 +53,13 @@ Paths in the appendix include `/v1`. SDK names here omit `client.`. `P` means pa | 3 | beta.agents.update | P: atomic replacements; empty body touches timestamp; metadata/name errors use official code and param | R `agent-patch-metadata`, `agent-null-fields`, `agent-noop`, `agent-nested-reasoning`, rejection labels; X VA-07/08/09/10 | C DB no-op/unchanged snapshot; resource implementation-validation.md actual PostgreSQL SDK update tests | Model-dependent default recomputation; uncommon nested/null/error variants | | 4 | beta.agents.list | P: scoped cursor list; unknown keys ignored, limit 0/above 100 clamp | R `agent-list-empty-scoped`, `agent-list-limit101`; L VA-01/02/03/04/18 | L DB `official_list_query.py` tenant A/B | Core page capacity 100; no inferred official cap. Overflowing limits unsampled | | 5 | beta.agents.delete | P: resource deletion | R `cleanup-agent`, subsequent 404 | C DB delete/post-delete | Referenced/in-flight/repeated-delete exact parity | -| 6 | beta.agents.sessions.create | P: JSON/live SSE 201, saved/inline frozen config, initial messages, native profiles | S `create-1/2.json`, `omitted-input.json`, `null-input.json`, `empty-array-input.json`, `retry-status-original/repeat.json`; H stream | N Live none admission and retry; C Live three hosted profiles; D/T/K/I recorded additional workflows | Session admission batch removes idle `none` creation; local idempotent create still differs from two official IDs. Many input/tool/environment combinations restricted | +| 6 | beta.agents.sessions.create | P: JSON/live SSE 201, saved/inline frozen config, initial messages, native profiles; fresh creation SSE sends the JSON 201 projection, then ends right after the first idle recorded when a Turn ends or an input reservation stops being pending, or any failed, never sending later events; nothing admitted ends after `created`; a silent settlement ends after events up to the cursor read with a settled projection in one snapshot. A same-key stream retry returns 201, sends no events and ends at once | S `create-1/2.json`, `omitted-input.json`, `null-input.json`, `empty-array-input.json`, `retry-status-original/repeat.json`; H stream; J creation streams closed after idle, open through requires_action | N Live none admission and retry; C Live three hosted profiles; D/T/K/I recorded additional workflows; J DB creation-stream lifetime/snapshot/retry | Session admission batch removes idle `none` creation; local idempotent create still differs from two official IDs. Stream retry, self-hosted, hosted and no-input stream lifetimes are local; work drained before a silent-settlement read can still be sent. Many input/tool/environment combinations restricted | | 7 | beta.agents.sessions.retrieve | P: persisted state, required actions, usage; malformed ID equals missing | S `retrieve-1.json`, `session-after-1.json`; H recovered state; X SES-28 | C Live history; T pending actions; H Core acceptance recorded | Complete statuses/actions/lifecycle timing; Claude/MiniMax public usage remains null | | 8 | beta.agents.sessions.update | P: metadata-only replacement/clear; metadata errors use official code and `metadata`/`metadata.` param | S `metadata-replace/null/empty/omit/invalid-value.json`; `update-agent.json` uses newer unpinned field | N Live completed Session metadata rejection/clear/isolation; recorded active controlled metadata coverage | Session admission batch changes empty update to observed official 400; Session agent update belongs to baseline upgrade, not fixed-pin operation gap | | 9 | beta.agents.sessions.list | P: Agent filter, full envelope, cursor paging; unknown keys ignored, limit 0/above 100 clamp | S `list-filter.json`, `list-empty-after.json`, `list-owned-cross-filter-cursor.json`, limit/order/unknown-query samples; L SES-10/11/15/16/17 | C Live order/cursors/empty/tenant checks; L DB tenant A/B | Core page capacity 100; official cap unknown. Eventual visibility sample is not a required delay | | 10 | beta.agents.sessions.delete | P: public deletion, owned managed cleanup, user compute retained | S cleanup files 200/deleted; retry-session active cleanup initially 409 | D recorded real cleanup; C cleanup separately recorded | Physical purge/retention and all active/unknown-effect races; caller compute ownership preserved | | 11 | beta.agents.sessions.events.create | P: 202/empty body, empty-array authenticated no-op, text/cancel/function admission | S `second-turn-create.json`, `events-empty/null.json`; H second-input | C Live real continuation/no-op; T qualified message/result/cancel workflows | Mixed prepared-environment batches, native receipt vs durable acceptance, cancel-before-result-publication timing; unqualified content/tools | -| 12 | beta.agents.sessions.events.stream | P: live-only SSE, typed persisted projections | H create/reconnect frames; no historical frames in sampled idle interval | C Live; H recorded three-harness disconnect/recovery; T pending actions | Full SSE/Item variants/order; child deltas settle late; no replay guarantee or observer-disconnect proof for every state | +| 12 | beta.agents.sessions.events.stream | P: live-only SSE, typed persisted projections; terminal Turn events carry top-level Turn usage (null when unknown); new Turns start `turn.created`, user `item.added`, `session.in_progress`, `turn.in_progress` | H create/reconnect frames; no historical frames in sampled idle interval; J GET streams never ended, terminal `usage` present | C Live; H recorded three-harness disconnect/recovery; T pending actions; J DB order/usage | Full SSE/Item variants/order (EVT-05..10 deferred); child deltas settle late; no replay guarantee or observer-disconnect proof for every state | | 13 | beta.agents.sessions.turns.retrieve | P: persisted root/child Turn identity; malformed ID equals missing | H/S contain Turn list payloads; no isolated positive retrieve raw request identified in this set; X SES-28 official `turn_` 404 | Recorded H/B scoped Turn recovery; C history uses list | Distinguish list-shape evidence from retrieve wire qualification; full lifecycle/usage | | 14 | beta.agents.sessions.turns.list | P: full envelope, ordered root+child history; limit outside 1–100 rejects with the Beta code | S `turns-final/empty-page/limit-high/order-empty.json`; H both directions; L SES-14/15/16/17; X SES-28 | C Live paging; H/B real child/root identity recorded; L DB tenant A/B | All interleavings, same-timestamp paging, interim/failed usage | | 15 | beta.agents.sessions.items.list | P: scoped root Items, full envelope; limit 0/above 100 clamp within the pinned 1–100 page | S `items-final/empty-page/limit-high.json`; H both directions; L SES-12/13/15/16/17 | C Live; H/T recorded content/coordination/result variants; L DB tenant A/B | Full Item union. Newer turn_id filter excluded from pin | diff --git a/contracts/agents-api/v1/events.go b/contracts/agents-api/v1/events.go index 8d6d20ed9..9ddc64483 100644 --- a/contracts/agents-api/v1/events.go +++ b/contracts/agents-api/v1/events.go @@ -1,5 +1,7 @@ package v1 +import "encoding/json" + // SessionEvent contains the supported live event variants of the pinned protocol. type SessionEvent struct { Subagent *Subagent `json:"subagent,omitempty"` @@ -18,6 +20,9 @@ type SessionEvent struct { Text *string `json:"text,omitempty"` Error *StreamError `json:"error,omitempty"` Environment *SessionEnvironmentState `json:"environment,omitempty"` + // Usage is present only on terminal Turn events, where it mirrors the Turn + // snapshot and is null when unknown. Other events omit it. + Usage *TokenUsage `json:"usage,omitempty" extensions:"x-nullable"` } type StreamError struct { @@ -25,3 +30,25 @@ type StreamError struct { Type string `json:"type"` Message string `json:"message"` } + +// TerminalTurnEvent reports whether an event type settles a Turn. +func TerminalTurnEvent(eventType string) bool { + switch eventType { + case "agent.session.turn.completed", "agent.session.turn.failed", "agent.session.turn.cancelled": + return true + } + return false +} + +// MarshalJSON keeps the nullable top-level usage on terminal Turn events only. +func (e SessionEvent) MarshalJSON() ([]byte, error) { + type wire SessionEvent + if !TerminalTurnEvent(e.Type) { + e.Usage = nil + return json.Marshal(wire(e)) + } + return json.Marshal(struct { + wire + Usage *TokenUsage `json:"usage"` + }{wire(e), e.Usage}) +} diff --git a/contracts/agents-api/v1/events_test.go b/contracts/agents-api/v1/events_test.go new file mode 100644 index 000000000..5d0891072 --- /dev/null +++ b/contracts/agents-api/v1/events_test.go @@ -0,0 +1,44 @@ +package v1 + +import ( + "encoding/json" + "testing" +) + +func TestSessionEventUsageOnlyOnTerminalTurnEvents(t *testing.T) { + usage := &TokenUsage{InputTokens: 7, InputTokensDetails: InputTokenDetails{CachedTokens: 2}, OutputTokens: 3, OutputTokensDetails: OutputTokenDetails{ReasoningTokens: 1}, TotalTokens: 10} + measured := `{"input_tokens":7,"input_tokens_details":{"cached_tokens":2},"output_tokens":3,"output_tokens_details":{"reasoning_tokens":1},"total_tokens":10}` + for _, test := range []struct { + event string + usage *TokenUsage + want string + }{ + {"agent.session.turn.completed", usage, measured}, + {"agent.session.turn.completed", nil, "null"}, + {"agent.session.turn.failed", nil, "null"}, + {"agent.session.turn.cancelled", usage, measured}, + {"agent.session.turn.in_progress", usage, ""}, + {"agent.session.turn.created", nil, ""}, + {"agent.session.idle", usage, ""}, + } { + raw, err := json.Marshal(SessionEvent{Type: test.event, EventID: "event", TurnID: "turn", Turn: &Turn{ID: "turn", Usage: test.usage}, Usage: test.usage}) + if err != nil { + t.Fatal(err) + } + var fields map[string]json.RawMessage + if err := json.Unmarshal(raw, &fields); err != nil { + t.Fatal(err) + } + value, present := fields["usage"] + if present != (test.want != "") || (present && string(value) != test.want) { + t.Fatal("unexpected top-level usage", test.event, string(raw)) + } + if fields["type"] == nil || fields["event_id"] == nil || fields["turn"] == nil || fields["turn_id"] == nil { + t.Fatal("terminal usage changed the other event fields", string(raw)) + } + var decoded SessionEvent + if err := json.Unmarshal(raw, &decoded); err != nil || (present && string(value) != "null" && decoded.Usage == nil) { + t.Fatal("event did not round-trip", string(raw), err) + } + } +} diff --git a/packages/agents-client/src/client.test.ts b/packages/agents-client/src/client.test.ts index 48d0e2746..c5dc9b8cb 100644 --- a/packages/agents-client/src/client.test.ts +++ b/packages/agents-client/src/client.test.ts @@ -1,6 +1,6 @@ import { afterEach, describe, expect, it, vi } from "vitest"; -import { AgentCoreError, createIdempotencyKey, OpenAIAgentsClient } from "./client"; +import { AgentCoreError, CreationStreamRetryError, createIdempotencyKey, OpenAIAgentsClient } from "./client"; import hostedDadf64 from "./fixtures/parsar-dadf64a7/openai-hosted.json"; import eventBatchDadf64 from "./fixtures/parsar-dadf64a7/session-event-batch.json"; import type { @@ -330,6 +330,36 @@ describe("OpenAIAgentsClient", () => { ]); }); + it("completes a creation stream that ends right after its initial Turn settles", async () => { + const admitted = { ...sessionResource(), status: "in_progress", last_active_at: 2 }; + const cancelled = turnResource({ status: "cancelled", completed_at: 3 }); + const frames = [ + { type: "agent.session.created", event_id: "created", session: admitted }, + { type: "agent.session.turn.created", event_id: "turn", session_id: "session", turn_id: "turn_1", turn: turnResource() }, + { + type: "agent.session.turn.item.added", event_id: "input", session_id: "session", turn_id: "turn_1", + item: messageItem({ status: "completed", role: "user", content: [{ type: "input_text", text: "First" }] }), + }, + { type: "agent.session.in_progress", event_id: "progress", session: admitted }, + { type: "agent.session.turn.cancelled", event_id: "done", session_id: "session", turn_id: "turn_1", turn: cancelled, usage: null }, + { type: "agent.session.idle", event_id: "idle", session: { ...sessionResource(), last_active_at: 3 } }, + ]; + const onSession = vi.fn(); + const onEvent = vi.fn(); + const client = new OpenAIAgentsClient({ + fetch: recordingFetch(streamResponse([ + ": connected\n\n", + ...frames.map((frame) => `event: ${frame.type}\ndata: ${JSON.stringify(frame)}\n\n`), + ], 201), []), + }); + + await client.createSessionStream({ environment: { type: "none" }, input: "First" }, "create-key", { onSession, onEvent }); + + expect(onSession).toHaveBeenCalledWith(expect.objectContaining({ id: "session", status: "in_progress" })); + expect(onEvent.mock.calls.map((call) => call[0]?.type)).toEqual(frames.map((frame) => frame.type)); + expect(onEvent.mock.calls[4]?.[0]).toMatchObject({ usage: null, turn: { status: "cancelled", usage: null } }); + }); + it.each([ ["default", { type: "openai_hosted" }], ["enabled", { type: "openai_hosted", network: { access: "enabled" } }], @@ -774,20 +804,45 @@ describe("OpenAIAgentsClient", () => { expect(cancelled).toBe(true); }); - it.each([ - ["null body", () => new Response(null), 0], - ["comments only", () => streamResponse([": connected\n\n: keepalive\n\n"]), 1], - ["DONE only", () => streamResponse(["data: [DONE]\n\n"]), 1], - ])("turns a creation stream with %s into empty_stream", async (_label, response, openCalls) => { + it("turns a creation response without a stream body into empty_stream", async () => { const onOpen = vi.fn(); - const client = new OpenAIAgentsClient({ fetch: recordingFetch(response(), []) }); + const client = new OpenAIAgentsClient({ fetch: recordingFetch(new Response(null, { status: 201 }), []) }); await expect(client.createSessionStream( { environment: { type: "none" } }, "create-key", { onOpen, onSession: vi.fn(), onEvent: vi.fn() }, )).rejects.toMatchObject({ status: 502, code: "empty_stream" }); - expect(onOpen).toHaveBeenCalledTimes(openCalls); + expect(onOpen).not.toHaveBeenCalled(); + }); + + it.each([ + ["the connection comment only", [": connected\n\n"]], + ["comments only", [": connected\n\n: keepalive\n\n"]], + ["DONE only", ["data: [DONE]\n\n"]], + ])("reports a creation stream with %s as a same-key retry to recover with stream=false", async (_label, chunks) => { + const onOpen = vi.fn(); + const onSession = vi.fn(); + const onEvent = vi.fn(); + const client = new OpenAIAgentsClient({ fetch: recordingFetch(streamResponse(chunks, 201), []) }); + + const failure = client.createSessionStream( + { environment: { type: "none" }, input: "First" }, + "create-key", + { onOpen, onSession, onEvent }, + ); + await expect(failure).rejects.toBeInstanceOf(CreationStreamRetryError); + await expect(failure).rejects.toMatchObject({ status: 409, code: "creation_stream_retry", message: expect.stringContaining("stream=false") }); + expect(onOpen).toHaveBeenCalledTimes(1); + expect(onSession).not.toHaveBeenCalled(); + expect(onEvent).not.toHaveBeenCalled(); + }); + + it("keeps empty_stream for a live events stream that ends without events", async () => { + const client = new OpenAIAgentsClient({ fetch: recordingFetch(streamResponse([": connected\n\n"]), []) }); + + await expect(client.streamEvents("session", { onEvent: vi.fn() })) + .rejects.toMatchObject({ status: 502, code: "empty_stream" }); }); it.each([ @@ -932,6 +987,80 @@ describe("OpenAIAgentsClient", () => { expect(onEvent.mock.calls[2]?.[0]).toMatchObject({ item_id: "item_1", item: { id: "item_1", turn_id: "turn_1" } }); }); + const measuredUsage = { + input_tokens: 7, + input_tokens_details: { cached_tokens: 2 }, + output_tokens: 3, + output_tokens_details: { reasoning_tokens: 1 }, + total_tokens: 10, + }; + + function terminalTurnEvent(status: string, fields: Record = {}): Record { + return { + type: `agent.session.turn.${status}`, + event_id: `turn-${status}`, + session_id: "session", + turn_id: "turn_1", + turn: turnResource({ + status, + started_at: 1, + completed_at: 2, + error: status === "failed" ? { code: "internal_error", message: "The execution could not complete." } : null, + }), + ...fields, + }; + } + + it.each([ + ["measured", "completed", measuredUsage], + ["unknown", "cancelled", null], + ["unknown failed", "failed", null], + ])("preserves %s top-level usage on a terminal Turn event", async (_label, status, usage) => { + const event = terminalTurnEvent(status, { turn: { ...terminalTurnEvent(status).turn as object, usage }, usage }); + const onEvent = vi.fn(); + const client = new OpenAIAgentsClient({ + fetch: recordingFetch(streamResponse([`event: ${event.type}\ndata: ${JSON.stringify(event)}\n\n`]), []), + }); + + await client.streamEvents("session", { onEvent }); + + expect(onEvent).toHaveBeenCalledTimes(1); + const projected = onEvent.mock.calls[0]?.[0] as SessionEvent; + expect(projected).toMatchObject({ type: event.type, turn_id: "turn_1", turn: { status, usage } }); + expect(Object.prototype.hasOwnProperty.call(projected, "usage")).toBe(true); + expect(projected.usage).toEqual(usage); + }); + + it("accepts a terminal Turn event without top-level usage", async () => { + const onEvent = vi.fn(); + const client = new OpenAIAgentsClient({ + fetch: recordingFetch(streamResponse([`data: ${JSON.stringify(terminalTurnEvent("completed"))}\n\n`]), []), + }); + + await client.streamEvents("session", { onEvent }); + + expect(onEvent).toHaveBeenCalledTimes(1); + expect(Object.prototype.hasOwnProperty.call(onEvent.mock.calls[0]?.[0], "usage")).toBe(false); + }); + + it.each([ + ["a non-terminal Turn event", { + type: "agent.session.turn.in_progress", event_id: "event", session_id: "session", turn_id: "turn_1", + turn: turnResource({ status: "in_progress", started_at: 1 }), usage: null, + }], + ["a malformed terminal value", terminalTurnEvent("completed", { usage: { input_tokens: 1 } })], + ["an estimated terminal value", terminalTurnEvent("completed", { usage: 0 })], + ])("rejects top-level usage on %s", async (_label, event) => { + const onEvent = vi.fn(); + const client = new OpenAIAgentsClient({ + fetch: recordingFetch(streamResponse([`data: ${JSON.stringify(event)}\n\n`]), []), + }); + + await expect(client.streamEvents("session", { onEvent })) + .rejects.toMatchObject({ status: 502, code: "invalid_stream_event" }); + expect(onEvent).not.toHaveBeenCalled(); + }); + it.each([ ["Turn Session", { type: "agent.session.turn.created", event_id: "event", session_id: "session", turn_id: "turn_1", diff --git a/packages/agents-client/src/client.ts b/packages/agents-client/src/client.ts index c41af13bf..5c46ca68f 100644 --- a/packages/agents-client/src/client.ts +++ b/packages/agents-client/src/client.ts @@ -96,6 +96,22 @@ export class AgentCoreError extends Error { } } +/** + * A creation stream ended without events. Core sends no events on a same-key + * retry of an existing creation; repeat the same request and Idempotency-Key + * with `stream=false` to retrieve the Session. + */ +export class CreationStreamRetryError extends AgentCoreError { + constructor() { + super( + "Agent Core already recorded this Session creation, so its stream sends no events. Repeat the same request and Idempotency-Key with stream=false to retrieve the Session.", + 409, + "creation_stream_retry", + ); + this.name = "CreationStreamRetryError"; + } +} + function trimTrailingSlash(value: string): string { return value.replace(/\/+$/, ""); } @@ -282,6 +298,11 @@ const streamErrorFields = new Set(["code", "type", "message"]); const environmentStateFields = new Set(["id", "type", "status", "error"]); const snapshotEventFields = new Set(["type", "event_id", "session_id", "session"]); const turnEventFields = new Set(["type", "event_id", "session_id", "turn_id", "turn"]); +// Terminal Turn events also carry top-level usage mirroring the Turn snapshot. +const terminalTurnEventFields = new Set([...turnEventFields, "usage"]); +const terminalTurnEventTypes = new Set([ + "agent.session.turn.completed", "agent.session.turn.failed", "agent.session.turn.cancelled", +]); const itemEventFields = new Set(["type", "event_id", "session_id", "turn_id", "item_id", "output_index", "item"]); const contentPartEventFields = new Set([ "type", "event_id", "session_id", "turn_id", "item_id", "output_index", "content_index", "part", @@ -1127,17 +1148,22 @@ function projectStreamEventSession( "agent.session.turn.cancelled": "cancelled", }; if (hasOwn(turnStatusByEvent, event.type)) { - if (!exactFields(value, turnEventFields) || event.turn === undefined) return invalidStreamEvent(); + // Older Core releases omit terminal usage; no other Turn event may carry it. + const withUsage = terminalTurnEventTypes.has(event.type) && hasOwn(value, "usage"); + if (!exactFields(value, withUsage ? terminalTurnEventFields : turnEventFields) || event.turn === undefined) { + return invalidStreamEvent(); + } const sessionId = eventSessionId(value, expectedSessionId, true)!; const turnId = requiredEventString(value, "turn_id"); const turn = projectAgentTurn(event.turn, expectedSessionId, invalidStreamEvent); + const usage = withUsage ? { usage: projectTokenUsage(value.usage, invalidStreamEvent) } : {}; if ( !sameResourceId(turn.id, turnId) || // Native child history can first appear as a completed created snapshot. (!(event.type === "agent.session.turn.created" && turn.subagent_id != null) && turn.status !== turnStatusByEvent[event.type]) || (immutable !== undefined && turn.subagent_id == null && !sameResourceId(turn.agent_id, immutable.agent.id)) ) return invalidStreamEvent(); - return { ...base, session_id: sessionId, turn_id: turnId, turn } as SessionEvent; + return { ...base, session_id: sessionId, turn_id: turnId, turn, ...usage } as SessionEvent; } if (event.type === "agent.session.turn.item.added" || event.type === "agent.session.turn.item.done") { @@ -1249,6 +1275,8 @@ function projectUnknownStreamEvent( interface EventStreamConsumerOptions extends StreamOptions { onParsedEvent: (event: SessionEvent) => void; expectedSessionId?: () => string | undefined; + /** Replaces the generic error when a present stream body ends without events. */ + emptyStreamError?: () => AgentCoreError; } async function consumeEventStream( @@ -1303,7 +1331,7 @@ async function consumeEventStream( decoder.finish(); options.signal?.throwIfAborted(); if (!sawEvent) { - throw new AgentCoreError("Agent core returned an empty event stream.", 502, "empty_stream"); + throw options.emptyStreamError?.() ?? new AgentCoreError("Agent core returned an empty event stream.", 502, "empty_stream"); } } catch (error) { await reader.cancel(error).catch(() => undefined); @@ -1932,6 +1960,7 @@ export class OpenAIAgentsClient implements AgentCore { let immutableSession: ImmutableSessionProjection | undefined; await consumeEventStream(response.body, { ...options, + emptyStreamError: () => new CreationStreamRetryError(), expectedSessionId: () => createdSessionId, onParsedEvent: (event) => { if (createdSessionId === undefined) { diff --git a/packages/agents-client/src/index.ts b/packages/agents-client/src/index.ts index 1e54698a2..d94a0142a 100644 --- a/packages/agents-client/src/index.ts +++ b/packages/agents-client/src/index.ts @@ -1,4 +1,4 @@ -export { AgentCoreError, createIdempotencyKey, OpenAIAgentsClient } from "./client"; +export { AgentCoreError, CreationStreamRetryError, createIdempotencyKey, OpenAIAgentsClient } from "./client"; export type { OpenAIAgentsClientOptions } from "./client"; export { createSSEDecoder } from "./sse"; export type { SSEDecoder, SSEMessage } from "./sse"; diff --git a/packages/agents-client/src/types.ts b/packages/agents-client/src/types.ts index 689bd0753..5e6a5696b 100644 --- a/packages/agents-client/src/types.ts +++ b/packages/agents-client/src/types.ts @@ -617,6 +617,11 @@ export interface SessionEventBase { delta?: string; text?: string; error?: StreamError; + /** + * Present on terminal Turn events only. It mirrors that Turn snapshot's usage + * and is null when unknown; a later Turn read can still report measured usage. + */ + usage?: TokenUsage | null; } export type AgentSessionEnvironmentEvent = { diff --git a/services/agents-api/internal/api/handler.go b/services/agents-api/internal/api/handler.go index 099211726..2d25bb727 100644 --- a/services/agents-api/internal/api/handler.go +++ b/services/agents-api/internal/api/handler.go @@ -127,7 +127,7 @@ func NewHandler(s ResourceStore, auth *Authenticator, engine string, options ... // createSession atomically reserves or admits initial text with the Session. // @Summary Create an execution Session -// @Description Supports inline configuration or a tenant-owned saved agent_id with per-Session field replacements. Execution supports model/instructions, text verbosity, non-deferred function tools, adapter-qualified multi_agent with persisted Subagent reads, implicit reasoning, service tier auto and environment type none, subject to the configured engine. Codex additionally supports HTTP MCP with explicit service origin, native allowed_tools and boolean required defaulting to false. Session vault_ids attach only project-owned Vaults; credential_id selects an attached static bearer credential for the exact HTTPS URL, while null/omission selects a unique match or remains anonymous. Ambiguous selection rejects creation. Frozen private selections never populate an omitted public credential_id; missing decryption configuration fails dispatch without anonymous fallback. Required initialization uses native startup before the first native Turn, including cold resume, and requires a separately advertised capability; exact hosted creation timing and error parity remain unverified. Other MCP origins and OAuth remain unsupported. The self_hosted profile requires Codex, an absolute workspace_directory and empty capability_directories, with optional non-deferred function tools and HTTP MCP using explicit service origin, optionally authenticated by the attached Vault rules. Remote MCP and remote Bearer authentication each require separately advertised combination support; old peers cannot receive unsupported work. Omitted/null capability_directories use the empty-list default; self_hosted requires configured execution plus executor registry. Claude SDK currently requires medium verbosity and object-root function schemas. It supports anonymous or attached static-bearer service-origin HTTP MCP on none with boolean required and separately advertised MCP/bearer/required runtime support. Required servers must be connected before the first native input is released; pending or failed startup rejects execution. The shared Vault selection and immutable binding rules apply; unsupported native labels/tool names reject before persistence. An attached Vault with no matching credential may remain anonymous; missing keys or failed credential lookup/decryption never fall back to anonymous execution. Omitted stream defaults to false; stream and agent_id cannot be null. Metadata may be null; non-string values and limit violations return invalid_request_error with a metadata or metadata. param. Hosted network policy rejections return invalid_request_error with a null param. Initial input accepts a string or ordered user-message array. Codex and Claude SDK on none and qualified openai_hosted also accept inline PNG/JPEG image content; other image combinations and remote URLs are unsupported. None initial input atomically starts a Turn; self_hosted initial input is reserved while returning its Environment connection target, with execution deferred to native readiness and Session failure on initial timeout. Initial input is required for none and for streamed creation outside self_hosted. Omitted/null input remains valid for non-streaming hosted and self_hosted creation. With stream=true, returns live Session events starting at creation; disconnect does not cancel execution. New Sessions retain their authenticated creator; all creation retries require the same typed subject, including across key rotation. Saved-Agent retries and inline requests using Vault attachments or credential references retain caller intent independently of later resource changes; unrelated inline retries preserve resolved/default equivalences. Unknown historical creators reject retries; known creators without recorded intent retain resolved-snapshot retry rules. These conflict policies are local and not verified hosted parity. Creation retries observe future events without replay; retry with stream=false to retrieve the Session. Claude SDK on none and Core-managed Docker openai_hosted supports qualified object-root json_schema output with medium verbosity, single-Agent execution and ordinary functions. Hosted execution reuses native workspace tools and Files/Artifacts; Skills, Plugins, capability directories, HTTP MCP, Subagent and tool_search combinations remain unqualified, including inherited template contents. Other non-text initial input remains unsupported. Basic Codex and Claude SDK openai_hosted creation requires an explicitly configured managed provider. The Claude workspace profile supports non-deferred function tools with text or successful inline PNG/JPEG results alongside native workspace tools; HTTP MCP remains unsupported. Idle Sessions provision automatically; initial provisioning has no caller connection action. Network defaults to enabled; disabled and restricted exact ASCII hostnames are supported. Restricted policy requires 1–100 allowed domains. Unsupported hostname forms and startup installations are rejected. Confidential env, system/npm/Python packages and ordered setup commands use the shared initialization lifecycle; requested network applies after setup. Initial inline and tenant-owned file_id files freeze encrypted bytes before provisioning, then install through the common Core lifecycle before native execution or live Files access. With a template reference, omitted/null files, env, packages and setup_commands inherit. Non-null files and command lists replace; env overlays by key; each package manager inherits on omission/null and otherwise replaces its list. Empty lists clear their selected field. Tenant-owned environment_template_id references inherit omitted/null network and allow only narrowing overrides. Inline hosted network:null retains the enabled default; updating a Template with network:null resets its saved policy to enabled. Core freezes effective configuration; template updates/deletion do not alter Session snapshots or same-intent creation retries. Inline or tenant-owned skill_reference Skills share initialization. Templates preserve default/latest/explicit selectors; Session creation freezes concrete metadata and encrypted content atomically. Skill, Plugin and capability-directory list omission/null inherit; a non-null list replaces, including empty-list clearing. Omitted/null Skill version selectors resolve the default version. Source deletion/default updates cannot change committed Session Skill contents. Deferred function discovery uses type-only tool_search and per-function defer_loading in the qualified single-agent Claude environment:none function profile, including qualified inline image messages and text results. Explicit web_search mode disabled and programmatic_tool_calling enabled false use frozen common Runtime controls. Enabled forms remain unqualified. Omitted programmatic configuration preserves native behavior, a documented difference from the official default-on behavior. Other combinations remain unqualified; see the operation coverage. +// @Description Supports inline configuration or a tenant-owned saved agent_id with per-Session field replacements. Execution supports model/instructions, text verbosity, non-deferred function tools, adapter-qualified multi_agent with persisted Subagent reads, implicit reasoning, service tier auto and environment type none, subject to the configured engine. Codex additionally supports HTTP MCP with explicit service origin, native allowed_tools and boolean required defaulting to false. Session vault_ids attach only project-owned Vaults; credential_id selects an attached static bearer credential for the exact HTTPS URL, while null/omission selects a unique match or remains anonymous. Ambiguous selection rejects creation. Frozen private selections never populate an omitted public credential_id; missing decryption configuration fails dispatch without anonymous fallback. Required initialization uses native startup before the first native Turn, including cold resume, and requires a separately advertised capability; exact hosted creation timing and error parity remain unverified. Other MCP origins and OAuth remain unsupported. The self_hosted profile requires Codex, an absolute workspace_directory and empty capability_directories, with optional non-deferred function tools and HTTP MCP using explicit service origin, optionally authenticated by the attached Vault rules. Remote MCP and remote Bearer authentication each require separately advertised combination support; old peers cannot receive unsupported work. Omitted/null capability_directories use the empty-list default; self_hosted requires configured execution plus executor registry. Claude SDK currently requires medium verbosity and object-root function schemas. It supports anonymous or attached static-bearer service-origin HTTP MCP on none with boolean required and separately advertised MCP/bearer/required runtime support. Required servers must be connected before the first native input is released; pending or failed startup rejects execution. The shared Vault selection and immutable binding rules apply; unsupported native labels/tool names reject before persistence. An attached Vault with no matching credential may remain anonymous; missing keys or failed credential lookup/decryption never fall back to anonymous execution. Omitted stream defaults to false; stream and agent_id cannot be null. Metadata may be null; non-string values and limit violations return invalid_request_error with a metadata or metadata. param. Hosted network policy rejections return invalid_request_error with a null param. Initial input accepts a string or ordered user-message array. Codex and Claude SDK on none and qualified openai_hosted also accept inline PNG/JPEG image content; other image combinations and remote URLs are unsupported. None initial input atomically starts a Turn; self_hosted initial input is reserved while returning its Environment connection target, with execution deferred to native readiness and Session failure on initial timeout. Initial input is required for none and for streamed creation outside self_hosted. Omitted/null input remains valid for non-streaming hosted and self_hosted creation. With stream=true, returns live Session events starting with the committed creation snapshot and closes right after the first agent.session.idle recorded when a Turn ends or an input reservation stops being pending, or any agent.session.failed, without sending later events. A creation that admitted nothing closes after the snapshot; a settlement that records no event closes after events up to the cursor read with a settled Session projection. Required actions keep it open; disconnect does not cancel execution. The GET events stream remains live-only. New Sessions retain their authenticated creator; all creation retries require the same typed subject, including across key rotation. Saved-Agent retries and inline requests using Vault attachments or credential references retain caller intent independently of later resource changes; unrelated inline retries preserve resolved/default equivalences. Unknown historical creators reject retries; known creators without recorded intent retain resolved-snapshot retry rules. These conflict policies are local and not verified hosted parity. A same-key stream=true retry of an existing creation returns 201 with no events and closes at once; retry with stream=false or use the GET events stream to recover. Claude SDK on none and Core-managed Docker openai_hosted supports qualified object-root json_schema output with medium verbosity, single-Agent execution and ordinary functions. Hosted execution reuses native workspace tools and Files/Artifacts; Skills, Plugins, capability directories, HTTP MCP, Subagent and tool_search combinations remain unqualified, including inherited template contents. Other non-text initial input remains unsupported. Basic Codex and Claude SDK openai_hosted creation requires an explicitly configured managed provider. The Claude workspace profile supports non-deferred function tools with text or successful inline PNG/JPEG results alongside native workspace tools; HTTP MCP remains unsupported. Idle Sessions provision automatically; initial provisioning has no caller connection action. Network defaults to enabled; disabled and restricted exact ASCII hostnames are supported. Restricted policy requires 1–100 allowed domains. Unsupported hostname forms and startup installations are rejected. Confidential env, system/npm/Python packages and ordered setup commands use the shared initialization lifecycle; requested network applies after setup. Initial inline and tenant-owned file_id files freeze encrypted bytes before provisioning, then install through the common Core lifecycle before native execution or live Files access. With a template reference, omitted/null files, env, packages and setup_commands inherit. Non-null files and command lists replace; env overlays by key; each package manager inherits on omission/null and otherwise replaces its list. Empty lists clear their selected field. Tenant-owned environment_template_id references inherit omitted/null network and allow only narrowing overrides. Inline hosted network:null retains the enabled default; updating a Template with network:null resets its saved policy to enabled. Core freezes effective configuration; template updates/deletion do not alter Session snapshots or same-intent creation retries. Inline or tenant-owned skill_reference Skills share initialization. Templates preserve default/latest/explicit selectors; Session creation freezes concrete metadata and encrypted content atomically. Skill, Plugin and capability-directory list omission/null inherit; a non-null list replaces, including empty-list clearing. Omitted/null Skill version selectors resolve the default version. Source deletion/default updates cannot change committed Session Skill contents. Deferred function discovery uses type-only tool_search and per-function defer_loading in the qualified single-agent Claude environment:none function profile, including qualified inline image messages and text results. Explicit web_search mode disabled and programmatic_tool_calling enabled false use frozen common Runtime controls. Enabled forms remain unqualified. Omitted programmatic configuration preserves native behavior, a documented difference from the official default-on behavior. Other combinations remain unqualified; see the operation coverage. // @Tags Sessions // @Accept json // @Produce json,text/event-stream diff --git a/services/agents-api/internal/api/session_creation_identity.go b/services/agents-api/internal/api/session_creation_identity.go index 592ecff2e..320869c02 100644 --- a/services/agents-api/internal/api/session_creation_identity.go +++ b/services/agents-api/internal/api/session_creation_identity.go @@ -54,7 +54,8 @@ func (h *Handler) recoverSessionCreation(w http.ResponseWriter, r *http.Request, writeError(w, http.StatusServiceUnavailable, "stream_unavailable", "Live events are unavailable.") return true } - h.respondSessionCreationStream(w, r, events, result) + // Recorded-intent lookup finds an existing creation, which sends no events. + h.respondSessionCreationStream(w, r, events, nil, result) } else { session, err := h.store.GetSession(r.Context(), tenantID(r), result.Session.ID) if err != nil { diff --git a/services/agents-api/internal/api/session_creation_stream.go b/services/agents-api/internal/api/session_creation_stream.go index b5a8db8ee..dc4c8e02c 100644 --- a/services/agents-api/internal/api/session_creation_stream.go +++ b/services/agents-api/internal/api/session_creation_stream.go @@ -15,8 +15,10 @@ type sessionStreamCreator interface { } func (h *Handler) createSessionStream(w http.ResponseWriter, r *http.Request, input store.CreateSessionInput) { + // Check every stream capability before the creation can commit. events, ok := h.store.(eventStore) - if !ok { + snapshots, snapshotted := h.store.(sessionSnapshotStore) + if !ok || !snapshotted { writeError(w, http.StatusServiceUnavailable, "stream_unavailable", "Live events are unavailable.") return } @@ -43,18 +45,63 @@ func (h *Handler) createSessionStream(w http.ResponseWriter, r *http.Request, in writeStoreError(w, r, err) return } - h.respondSessionCreationStream(w, r, events, result) + h.respondSessionCreationStream(w, r, events, snapshots, result) +} + +type sessionSnapshotStore interface { + SessionStreamSnapshot(context.Context, string, string) (store.Session, int64, error) } -func (h *Handler) respondSessionCreationStream(w http.ResponseWriter, r *http.Request, events eventStore, result store.SessionCreation) { +// respondSessionCreationStream renders result.Session, the committed Session +// projection read after result.Cursor, which is also the JSON 201 body. A fresh +// creation sends it as agent.session.created, then streams changes after the +// cursor until the Session settles, or ends at once when nothing was admitted. +// A same-key retry of an existing creation admits nothing and sends no events; +// official same-key requests create distinct Sessions, so there is no retry +// stream to follow. Recover with stream=false or the GET events stream. Only a +// fresh creation uses snapshots. +func (h *Handler) respondSessionCreationStream(w http.ResponseWriter, r *http.Request, events eventStore, snapshots sessionSnapshotStore, result store.SessionCreation) { + if !result.Created { + openEventStream(w, http.StatusCreated) + return + } response, err := sessionResponse(result.Session, h.executorURL) if err != nil { writeStoreError(w, r, err) return } - var initial *v1.SessionEvent - if result.Created { - initial = &v1.SessionEvent{Type: "agent.session.created", EventID: uuid.NewString(), Session: &response} + created := v1.SessionEvent{Type: "agent.session.created", EventID: uuid.NewString(), Session: &response} + if session := result.Session; sessionSettled(session, response) && session.LastTurn == nil && session.EnvironmentInputActivity == nil { + // Nothing was admitted, e.g. self_hosted creation without input. + if write := openEventStream(w, http.StatusCreated); write != nil { + _ = emitSessionEvent(write, session.ID, created) + } + return + } + tenant, id := tenantID(r), result.Session.ID + settlement := func(ctx context.Context) (bool, int64, error) { + session, cursor, err := snapshots.SessionStreamSnapshot(ctx, tenant, id) + if err != nil { + return false, 0, err + } + response, err := sessionResponse(session, h.executorURL) + return err == nil && sessionSettled(session, response), cursor, err + } + h.serveSessionEvents(w, r, events, result.Session, result.Cursor, &created, http.StatusCreated, settlement) +} + +// sessionSettled reports that a committed projection has no admitted work left: +// the Session is idle or failed, its latest Turn is not queued, running or +// waiting, and its latest input reservation is not pending. +func sessionSettled(session store.Session, response v1.Session) bool { + if response.Status != "idle" && response.Status != "failed" { + return false + } + if turn := session.LastTurn; turn != nil { + switch turn.Status { + case store.TurnQueued, store.TurnInProgress, store.TurnWaiting: + return false + } } - h.serveSessionEvents(w, r, events, result.Session, result.Cursor, initial, http.StatusCreated) + return !session.PendingInput } diff --git a/services/agents-api/internal/api/session_creation_stream_test.go b/services/agents-api/internal/api/session_creation_stream_test.go new file mode 100644 index 000000000..dce8d8d57 --- /dev/null +++ b/services/agents-api/internal/api/session_creation_stream_test.go @@ -0,0 +1,726 @@ +package api + +import ( + "bufio" + "context" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "reflect" + "strings" + "sync/atomic" + "testing" + "time" + + v1 "github.com/MiniMax-AI-Dev/parsar/contracts/agents-api/v1" + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/identity" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" + "github.com/google/uuid" +) + +// creationStreamFixture serves one Session through both creation paths: the +// execution upsert (CreateSessionStream) and recorded-intent retry lookup. Like +// the Store, it commits events with ordered sequences together with the state +// they describe. +type creationStreamFixture struct { + streamFixture + creation store.SessionCreation + recorded bool + deleted bool + sequence int64 + snapshots int + // beforeSnapshot and afterSnapshot commit racing work around the next read. + beforeSnapshot, afterSnapshot func() +} + +// GetSession returns the committed projection that commit changes atomically +// with its events. +func (f *creationStreamFixture) GetSession(_ context.Context, tenant, id string) (store.Session, error) { + f.mu.Lock() + defer f.mu.Unlock() + if f.deleted || tenant != f.session.TenantID || id != f.session.ID { + return store.Session{}, store.ErrNotFound + } + return f.session, nil +} + +func (f *creationStreamFixture) SessionEventCursor(context.Context, string, string) (int64, error) { + f.mu.Lock() + defer f.mu.Unlock() + return f.sequence, nil +} + +func (f *creationStreamFixture) ListSessionEvents(_ context.Context, _, _ string, cursor int64) ([]store.SessionChange, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.cursors = append(f.cursors, cursor) + var changes []store.SessionChange + for _, change := range f.changes { + if change.Sequence > cursor { + changes = append(changes, change) + } + } + return changes, nil +} + +// SessionStreamSnapshot returns the projection and cursor from one snapshot, +// running the one-shot race hooks before and after that read. +func (f *creationStreamFixture) SessionStreamSnapshot(_ context.Context, tenant, id string) (store.Session, int64, error) { + f.mu.Lock() + before := f.beforeSnapshot + f.beforeSnapshot = nil + f.mu.Unlock() + if before != nil { + before() + } + f.mu.Lock() + f.snapshots++ + session, cursor, deleted, race := f.session, f.sequence, f.deleted, f.afterSnapshot + f.afterSnapshot = nil + f.mu.Unlock() + if race != nil { + race() + } + if deleted || tenant != session.TenantID || id != session.ID { + return store.Session{}, 0, store.ErrNotFound + } + return session, cursor, nil +} + +func (f *creationStreamFixture) FindSessionCreation(context.Context, string, string, json.RawMessage, identity.Subject) (store.SessionCreation, error) { + if !f.recorded { + return store.SessionCreation{}, store.ErrNotFound + } + f.mu.Lock() + defer f.mu.Unlock() + // Recorded-intent lookup returns only the resource row and its cursor. + row := f.session + row.LastTurn, row.Usage, row.RequiredActions = nil, nil, nil + return store.SessionCreation{Session: row, Cursor: 10}, nil +} + +type creationStreamAdmission struct { + inputRecorder + fixture *creationStreamFixture +} + +func (a *creationStreamAdmission) CreateSessionStream(context.Context, string, store.CreateSessionInput) (store.SessionCreation, error) { + return a.fixture.creation, nil +} + +type sseFrame struct { + Type string + Data map[string]json.RawMessage +} + +func (f sseFrame) event(t *testing.T) v1.SessionEvent { + t.Helper() + raw, _ := json.Marshal(f.Data) + var event v1.SessionEvent + if err := json.Unmarshal(raw, &event); err != nil { + t.Fatal(err) + } + return event +} + +type creationStreamHarness struct { + t *testing.T + fixture *creationStreamFixture + server *httptest.Server + active atomic.Int32 +} + +func newCreationStreamHarness(t *testing.T) *creationStreamHarness { + t.Helper() + tenant := uuid.NewString() + fixture := &creationStreamFixture{sequence: 10, streamFixture: streamFixture{session: store.Session{ + ID: uuid.NewString(), TenantID: tenant, CreatedAt: time.Unix(1700000000, 0), Metadata: map[string]string{}, + Configuration: json.RawMessage(`{"agent":{"id":"agent_test","model":"model","tools":[]},"environment":{"type":"none"}}`), + }}} + auth, err := NewAuthenticator([]APIKey{{ + OrganizationID: "test-org", ProjectID: uuid.NewString(), SubjectKind: "service_account", SubjectID: "test-runner", + TokenSHA256: device.HashCredential("key"), TenantID: tenant, + }}) + if err != nil { + t.Fatal(err) + } + handler, err := NewHandler(fixture, auth, "codex", WithExecution(&creationStreamAdmission{fixture: fixture})) + if err != nil { + t.Fatal(err) + } + h := &creationStreamHarness{t: t, fixture: fixture} + h.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + h.active.Add(1) + defer h.active.Add(-1) + handler.ServeHTTP(w, r) + })) + t.Cleanup(h.server.Close) + // Runs before server.Close: every stream handler must have returned. + t.Cleanup(func() { + deadline := time.Now().Add(3 * time.Second) + for h.active.Load() != 0 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + if h.active.Load() != 0 { + t.Error("stream handler leaked after the stream ended or disconnected") + } + }) + return h +} + +// turn returns a Turn snapshot with the given status for the fixture Session. +func (h *creationStreamHarness) turn(status string, usage string) *store.Turn { + turn := &store.Turn{ID: "turn_1", SessionID: h.fixture.session.ID, Status: status, CreatedAt: time.Unix(1700000001, 0)} + if usage != "" { + turn.Usage = json.RawMessage(usage) + } + return turn +} + +// commit applies a projection change and its events atomically, as the Store +// commits Session events with the state they describe. +func (h *creationStreamHarness) commit(update func(*store.Session), changes ...store.SessionChange) { + h.fixture.mu.Lock() + defer h.fixture.mu.Unlock() + if update != nil { + update(&h.fixture.session) + } + for _, change := range changes { + h.fixture.sequence++ + change.Sequence = h.fixture.sequence + h.fixture.changes = append(h.fixture.changes, change) + } +} + +func (h *creationStreamHarness) snapshots() int { + h.fixture.mu.Lock() + defer h.fixture.mu.Unlock() + return h.fixture.snapshots +} + +func (h *creationStreamHarness) open(method, path, body string) (<-chan sseFrame, func()) { + h.t.Helper() + ctx, cancel := context.WithCancel(h.t.Context()) + request, err := http.NewRequestWithContext(ctx, method, h.server.URL+path, strings.NewReader(body)) + if err != nil { + h.t.Fatal(err) + } + request.Header.Set("Authorization", "Bearer key") + request.Header.Set("OpenAI-Beta", "agents=v1") + request.Header.Set("Idempotency-Key", "creation") + response, err := h.server.Client().Do(request) + if err != nil { + h.t.Fatal(err) + } + want := http.StatusCreated + if method == http.MethodGet { + want = http.StatusOK + } + if response.StatusCode != want || response.Header.Get("Content-Type") != "text/event-stream" { + h.t.Fatal("stream was not opened", response.StatusCode, response.Header) + } + frames := make(chan sseFrame, 64) + go func() { + defer close(frames) + defer response.Body.Close() + scanner := bufio.NewScanner(response.Body) + scanner.Buffer(make([]byte, 1<<20), 1<<20) + var frame sseFrame + for scanner.Scan() { + line := scanner.Text() + if value, ok := strings.CutPrefix(line, "event: "); ok { + frame.Type = value + } else if value, ok := strings.CutPrefix(line, "data: "); ok { + _ = json.Unmarshal([]byte(value), &frame.Data) + } else if line == "" && frame.Type != "" { + frames <- frame + frame = sseFrame{} + } + } + }() + return frames, cancel +} + +func nextFrame(t *testing.T, frames <-chan sseFrame) sseFrame { + t.Helper() + select { + case frame, ok := <-frames: + if !ok { + t.Fatal("stream ended early") + } + return frame + case <-time.After(5 * time.Second): + t.Fatal("stream event timed out") + } + return sseFrame{} +} + +func expectEnded(t *testing.T, frames <-chan sseFrame) { + t.Helper() + select { + case frame, ok := <-frames: + if ok { + t.Fatal("stream continued after settlement", frame.Type) + } + case <-time.After(5 * time.Second): + t.Fatal("stream did not end after settlement") + } +} + +// expectOpen outlasts the periodic projection re-read. +func expectOpen(t *testing.T, frames <-chan sseFrame) { + t.Helper() + select { + case frame, ok := <-frames: + t.Fatal("stream changed while it should wait", frame.Type, ok) + case <-time.After(1200 * time.Millisecond): + } +} + +func expectTypes(t *testing.T, frames <-chan sseFrame, types ...string) []sseFrame { + t.Helper() + observed := make([]sseFrame, 0, len(types)) + for _, want := range types { + frame := nextFrame(t, frames) + if frame.Type != want || string(frame.Data["type"]) != `"`+want+`"` { + t.Fatal("unexpected event order", frame.Type, "want", want) + } + observed = append(observed, frame) + } + return observed +} + +func sameJSON(t *testing.T, raw json.RawMessage, value any) bool { + t.Helper() + encoded, err := json.Marshal(value) + if err != nil { + t.Fatal(err) + } + var left, right any + if json.Unmarshal(raw, &left) != nil || json.Unmarshal(encoded, &right) != nil { + return false + } + return reflect.DeepEqual(left, right) +} + +func turnChange(kind string, turn *store.Turn) store.SessionChange { + return store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.turn." + kind, EventID: uuid.NewString(), SessionID: turn.SessionID, TurnID: turn.ID}, Turn: turn} +} + +// sessionChange records a Turn-owned Session status snapshot, settled when its +// Turn is terminal, as the Store records it. +func sessionChange(status string, turn *store.Turn, actions []v1.FunctionCallAction) store.SessionChange { + settled := turn != nil && (turn.Status == store.TurnCompleted || turn.Status == store.TurnFailed || turn.Status == store.TurnCancelled) + return store.SessionChange{Event: v1.SessionEvent{Type: "agent.session." + status, EventID: uuid.NewString(), SessionID: "session"}, Turn: turn, RequiredActions: actions, Settled: settled} +} + +func activityChange(activity *store.EnvironmentInputActivity, settled bool) store.SessionChange { + return store.SessionChange{Event: v1.SessionEvent{Type: "agent.session." + activity.Status, EventID: uuid.NewString()}, EnvironmentInputActivity: activity, Settled: settled} +} + +func userItemChange(turn *store.Turn) store.SessionChange { + text := "First" + return store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.turn.item.added", EventID: uuid.NewString(), SessionID: turn.SessionID, TurnID: turn.ID, + Item: &v1.Item{ID: "item_1", TurnID: turn.ID, Type: "message", Status: "completed", Role: "user", Content: []v1.ItemContent{{Type: "input_text", Text: &text}}}}} +} + +const measuredUsage = `{"input_tokens":7,"input_tokens_details":{"cached_tokens":2},"output_tokens":3,"output_tokens_details":{"reasoning_tokens":1},"total_tokens":10}` + +const creationBody = `{"agent":{"model":"model"},"environment":{"type":"none"},"input":"First","stream":true}` + +const savedBody = `{"agent_id":"agent_test","environment":{"type":"none"},"input":"First","stream":true}` + +func setTurn(turn *store.Turn) func(*store.Session) { + return func(session *store.Session) { session.LastTurn, session.RequiredActions = turn, nil } +} + +func TestCreationStreamClosesOnceSettled(t *testing.T) { + for _, test := range []struct{ terminal, session, usage string }{ + {"completed", "idle", measuredUsage}, + {"cancelled", "idle", ""}, + {"failed", "failed", ""}, + } { + t.Run(test.terminal, func(t *testing.T) { + h := newCreationStreamHarness(t) + // The admission result is the committed post-input projection. + queued := h.turn(store.TurnQueued, "") + h.commit(setTurn(queued), turnChange("created", queued), userItemChange(queued), sessionChange("in_progress", queued, nil)) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: 10} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + observed := expectTypes(t, frames, "agent.session.created", "agent.session.turn.created", "agent.session.turn.item.added", "agent.session.in_progress") + created := observed[0].event(t) + projection, err := sessionResponse(h.fixture.creation.Session, "") + if err != nil || created.Session == nil || created.Session.Status != "in_progress" || !sameJSON(t, observed[0].Data["session"], projection) { + t.Fatal("created snapshot differs from the JSON projection", string(observed[0].Data["session"]), err) + } + running := h.turn(store.TurnInProgress, "") + h.commit(setTurn(running), turnChange("in_progress", running)) + expectTypes(t, frames, "agent.session.turn.in_progress") + expectOpen(t, frames) + + settled := h.turn(test.terminal, test.usage) + h.commit(setTurn(settled), turnChange(test.terminal, settled), sessionChange(test.session, settled, nil)) + terminal := expectTypes(t, frames, "agent.session.turn."+test.terminal, "agent.session."+test.session) + expectEnded(t, frames) + want := "null" + if test.usage != "" { + want = test.usage + } + var turn map[string]json.RawMessage + if json.Unmarshal(terminal[0].Data["turn"], &turn) != nil || string(terminal[0].Data["usage"]) != want || string(turn["usage"]) != want { + t.Fatal("terminal usage does not mirror the Turn snapshot", string(terminal[0].Data["usage"]), string(turn["usage"])) + } + for _, frame := range append(observed, terminal[1]) { + if _, present := frame.Data["usage"]; present { + t.Fatal("usage is limited to terminal Turn events", frame.Type) + } + } + }) + } +} + +func TestCreationStreamStaysOpenAcrossRequiredAction(t *testing.T) { + h := newCreationStreamHarness(t) + queued := h.turn(store.TurnQueued, "") + h.commit(setTurn(queued), turnChange("created", queued), userItemChange(queued), sessionChange("in_progress", queued, nil)) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: 10} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + expectTypes(t, frames, "agent.session.created", "agent.session.turn.created", "agent.session.turn.item.added", "agent.session.in_progress") + + waiting := h.turn(store.TurnWaiting, "") + actions := []v1.FunctionCallAction{{Type: "function_call", CallID: "call_1", Name: "lookup", TurnID: waiting.ID, Arguments: json.RawMessage(`{}`)}} + h.commit(func(session *store.Session) { session.LastTurn, session.RequiredActions = waiting, actions }, + sessionChange("requires_action", waiting, actions)) + expectTypes(t, frames, "agent.session.requires_action") + expectOpen(t, frames) + + // The function result resumes the Turn; the stream ends once it settles. + resumed := h.turn(store.TurnInProgress, "") + h.commit(setTurn(resumed), sessionChange("in_progress", resumed, nil)) + expectTypes(t, frames, "agent.session.in_progress") + expectOpen(t, frames) + completed := h.turn(store.TurnCompleted, measuredUsage) + h.commit(setTurn(completed), turnChange("completed", completed), sessionChange("idle", completed, nil)) + expectTypes(t, frames, "agent.session.turn.completed", "agent.session.idle") + expectEnded(t, frames) +} + +func TestCreationStreamWithoutAdmittedWorkClosesAfterCreated(t *testing.T) { + h := newCreationStreamHarness(t) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: 10} + // Work committed after the creation is not followed. + queued := h.turn(store.TurnQueued, "") + h.commit(nil, turnChange("created", queued)) + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + expectTypes(t, frames, "agent.session.created") + expectEnded(t, frames) + h.fixture.mu.Lock() + defer h.fixture.mu.Unlock() + if h.fixture.snapshots != 0 || len(h.fixture.cursors) != 0 { + t.Fatal("a creation with nothing admitted read its projection or events", h.fixture.snapshots, h.fixture.cursors) + } +} + +func TestCreationStreamWaitsForPendingProvisioning(t *testing.T) { + for _, outcome := range []string{"failure", "expiry without an event"} { + t.Run(outcome, func(t *testing.T) { + h := newCreationStreamHarness(t) + // Hosted initial input is idle, without a Turn or public action, while it provisions. + h.commit(func(session *store.Session) { session.PendingInput = true }) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: 10} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + if created := expectTypes(t, frames, "agent.session.created")[0].event(t); created.Session.Status != "idle" { + t.Fatal("provisioning snapshot changed", created.Session.Status) + } + expectOpen(t, frames) + if outcome == "failure" { + failed := &store.EnvironmentInputActivity{Status: "failed", Failure: "environment_unavailable", LastActiveAt: time.Unix(1700000005, 0)} + h.commit(func(session *store.Session) { session.PendingInput, session.EnvironmentInputActivity = false, failed }, + activityChange(failed, true)) + expectTypes(t, frames, "agent.session.failed") + } else { + // A reservation can settle without recording an event. + h.commit(func(session *store.Session) { session.PendingInput = false }) + } + expectEnded(t, frames) + }) + } +} + +func TestCreationStreamEndsWhenSessionIsDeleted(t *testing.T) { + h := newCreationStreamHarness(t) + queued := h.turn(store.TurnQueued, "") + h.commit(setTurn(queued)) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: 10} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + expectTypes(t, frames, "agent.session.created") + h.commit(func(*store.Session) { h.fixture.deleted = true }) + expectEnded(t, frames) +} + +func TestGetStreamStaysOpenAfterSettlement(t *testing.T) { + h := newCreationStreamHarness(t) + h.commit(setTurn(h.turn(store.TurnInProgress, ""))) + frames, cancel := h.open(http.MethodGet, "/v1/agents/sessions/"+h.fixture.session.ID+"/events", "") + defer cancel() + completed := h.turn(store.TurnCompleted, "") + h.commit(setTurn(completed), turnChange("completed", completed), sessionChange("idle", completed, nil)) + expectTypes(t, frames, "agent.session.turn.completed", "agent.session.idle") + expectOpen(t, frames) + later := h.turn(store.TurnQueued, "") + failed := h.turn(store.TurnFailed, "") + h.commit(setTurn(failed), turnChange("created", later), sessionChange("failed", failed, nil)) + expectTypes(t, frames, "agent.session.turn.created", "agent.session.failed") + expectOpen(t, frames) +} + +// The creation stream never runs into another client's Turn, whether that work +// commits in the same drained batch as the settling idle or between the +// fallback's settled snapshot and its drain. +func TestCreationStreamStopsBeforeLaterWork(t *testing.T) { + later := func(h *creationStreamHarness) []store.SessionChange { + queued := h.turn(store.TurnQueued, "") + queued.ID = "turn_b" + return []store.SessionChange{turnChange("created", queued), userItemChange(queued), sessionChange("in_progress", queued, nil)} + } + t.Run("same batch as the settling idle", func(t *testing.T) { + h := newCreationStreamHarness(t) + h.commit(setTurn(h.turn(store.TurnInProgress, ""))) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: h.fixture.sequence} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + expectTypes(t, frames, "agent.session.created") + // The settling idle and B's new Turn commit before the next poll drains them together. + completed := h.turn(store.TurnCompleted, "") + changes := append([]store.SessionChange{turnChange("completed", completed), sessionChange("idle", completed, nil)}, later(h)...) + h.commit(func(session *store.Session) { session.LastTurn = changes[2].Turn }, changes...) + expectTypes(t, frames, "agent.session.turn.completed", "agent.session.idle") + expectEnded(t, frames) + }) + t.Run("between the settled snapshot and its drain", func(t *testing.T) { + h := newCreationStreamHarness(t) + h.commit(func(session *store.Session) { session.PendingInput = true }) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: h.fixture.sequence} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + expectTypes(t, frames, "agent.session.created") + expectOpen(t, frames) + h.fixture.mu.Lock() + // After an empty drain, the reservation settles without a status event, + // alongside an unsent environment event, before the fallback read. + h.fixture.beforeSnapshot = func() { + h.commit(func(session *store.Session) { session.PendingInput = false }, + store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.environment.disconnected", EventID: uuid.NewString(), SessionID: h.fixture.session.ID, + Environment: &v1.SessionEnvironmentState{ID: "environment", Type: "self_hosted", Status: "disconnected"}}}) + } + // B's Turn commits after that settled read and before its drain. + h.fixture.afterSnapshot = func() { + changes := later(h) + h.commit(func(session *store.Session) { session.LastTurn = changes[0].Turn }, changes...) + } + h.fixture.mu.Unlock() + expectTypes(t, frames, "agent.session.environment.disconnected") + expectEnded(t, frames) + }) +} + +func TestCreationStreamIgnoresConnectionIdle(t *testing.T) { + h := newCreationStreamHarness(t) + waiting := &store.EnvironmentInputActivity{Status: "requires_action", LastActiveAt: time.Unix(1700000002, 0)} + h.commit(func(session *store.Session) { session.PendingInput = true }) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: h.fixture.sequence} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + expectTypes(t, frames, "agent.session.created") + // A self-hosted connection clears the action; the input is still pending. + connected := &store.EnvironmentInputActivity{Status: "idle", LastActiveAt: waiting.LastActiveAt} + h.commit(func(session *store.Session) { session.EnvironmentInputActivity = connected }, activityChange(connected, false)) + expectTypes(t, frames, "agent.session.idle") + expectOpen(t, frames) + // Its later expiry records no event and ends the stream through the projection. + h.commit(func(session *store.Session) { session.PendingInput = false }) + expectEnded(t, frames) +} + +func TestCreationStreamBoundsProjectionReads(t *testing.T) { + h := newCreationStreamHarness(t) + running := h.turn(store.TurnInProgress, "") + h.commit(setTurn(running)) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Created: true, Cursor: h.fixture.sequence} + frames, cancel := h.open(http.MethodPost, "/v1/agents/sessions", creationBody) + defer cancel() + expectTypes(t, frames, "agent.session.created") + start := h.snapshots() + deadline := time.Now().Add(2200 * time.Millisecond) + count := 0 + for time.Now().Before(deadline) { + text := "x" + h.commit(nil, store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.turn.output_text.delta", EventID: uuid.NewString(), SessionID: running.SessionID, TurnID: running.ID, Delta: &text}}) + count++ + time.Sleep(20 * time.Millisecond) + } + for range count { + nextFrame(t, frames) + } + // Deltas carry no Session status, so reads stay on the once-a-second timer. + if reads := h.snapshots() - start; reads > 4 { + t.Fatal("chatty stream re-read its projection too often", reads, count) + } + h.commit(nil, sessionChange("in_progress", running, nil)) + expectTypes(t, frames, "agent.session.in_progress") + expectOpen(t, frames) +} + +// A same-key stream retry of an existing creation sends the connection comment +// and ends, whatever the Session is doing: nothing is replayed or followed. +func TestCreationRetryStreamEndsImmediately(t *testing.T) { + for _, test := range []struct { + name string + setup func(h *creationStreamHarness) + upsert bool + }{ + {"initial work running", func(h *creationStreamHarness) { + running := h.turn(store.TurnInProgress, "") + h.commit(setTurn(running), turnChange("in_progress", running)) + }, false}, + {"initial work settled", func(h *creationStreamHarness) { + completed := h.turn(store.TurnCompleted, "") + h.commit(setTurn(completed), turnChange("completed", completed), sessionChange("idle", completed, nil)) + }, false}, + {"superseded by another client's Turn", func(h *creationStreamHarness) { + later := h.turn(store.TurnInProgress, "") + later.ID = "turn_b" + h.commit(setTurn(later), turnChange("created", later), userItemChange(later), sessionChange("in_progress", later, nil)) + }, false}, + {"upsert retry", func(h *creationStreamHarness) { + h.commit(setTurn(h.turn(store.TurnInProgress, ""))) + h.fixture.creation = store.SessionCreation{Session: h.fixture.session, Cursor: 10} + }, true}, + } { + t.Run(test.name, func(t *testing.T) { + h := newCreationStreamHarness(t) + h.fixture.recorded = !test.upsert + test.setup(h) + body := savedBody + if test.upsert { + body = creationBody + } + request, err := http.NewRequestWithContext(t.Context(), http.MethodPost, h.server.URL+"/v1/agents/sessions", strings.NewReader(body)) + if err != nil { + t.Fatal(err) + } + request.Header.Set("Authorization", "Bearer key") + request.Header.Set("OpenAI-Beta", "agents=v1") + request.Header.Set("Idempotency-Key", "creation") + response, err := h.server.Client().Do(request) + if err != nil { + t.Fatal(err) + } + raw, err := io.ReadAll(response.Body) + _ = response.Body.Close() + if err != nil || response.StatusCode != http.StatusCreated || response.Header.Get("Content-Type") != "text/event-stream" || string(raw) != ": connected\n\n" { + t.Fatal("retry stream did not end at once", response.StatusCode, string(raw), err) + } + h.fixture.mu.Lock() + defer h.fixture.mu.Unlock() + if h.fixture.snapshots != 0 || len(h.fixture.cursors) != 0 { + t.Fatal("retry stream read Session state or events", h.fixture.snapshots, h.fixture.cursors) + } + }) + } +} + +func TestSettlingEvents(t *testing.T) { + for _, test := range []struct { + change store.SessionChange + settles bool + }{ + {sessionChange("idle", &store.Turn{Status: store.TurnCompleted}, nil), true}, + {activityChange(&store.EnvironmentInputActivity{Status: "idle"}, true), true}, + {activityChange(&store.EnvironmentInputActivity{Status: "idle"}, false), false}, + {activityChange(&store.EnvironmentInputActivity{Status: "failed"}, true), true}, + {sessionChange("failed", &store.Turn{Status: store.TurnFailed}, nil), true}, + {sessionChange("requires_action", &store.Turn{Status: store.TurnWaiting}, nil), false}, + {sessionChange("in_progress", &store.Turn{Status: store.TurnQueued}, nil), false}, + {turnChange("completed", &store.Turn{Status: store.TurnCompleted}), false}, + } { + if settlingEvent(test.change) != test.settles { + t.Fatal("unexpected settling event", test.change.Event.Type, test.change.Settled) + } + } +} + +func TestSessionSettledProjection(t *testing.T) { + for _, test := range []struct { + status string + turn string + pending bool + settled bool + }{ + {"idle", store.TurnCompleted, false, true}, + {"failed", store.TurnFailed, false, true}, + {"idle", "", false, true}, + {"failed", "", false, true}, + {"idle", "", true, false}, + {"idle", store.TurnCancelled, true, false}, + {"in_progress", store.TurnQueued, false, false}, + {"requires_action", store.TurnWaiting, false, false}, + {"requires_action", "", true, false}, + } { + session := store.Session{PendingInput: test.pending} + if test.turn != "" { + session.LastTurn = &store.Turn{Status: test.turn} + } + if sessionSettled(session, v1.Session{Status: test.status}) != test.settled { + t.Fatal("unexpected settlement", test) + } + } +} + +// eventOnlyStore streams events but cannot read the same-snapshot projection. +type eventOnlyStore struct{ ResourceStore } + +func (eventOnlyStore) SessionEventCursor(context.Context, string, string) (int64, error) { + return 0, nil +} + +func (eventOnlyStore) ListSessionEvents(context.Context, string, string, int64) ([]store.SessionChange, error) { + return nil, nil +} + +type countingStreamAdmission struct { + inputRecorder + calls atomic.Int32 +} + +func (a *countingStreamAdmission) CreateSessionStream(context.Context, string, store.CreateSessionInput) (store.SessionCreation, error) { + a.calls.Add(1) + return store.SessionCreation{}, store.ErrInvalidInput +} + +func TestCreationStreamCapabilityIsCheckedBeforeCreation(t *testing.T) { + auth, err := NewAuthenticator([]APIKey{{OrganizationID: "test-org", ProjectID: uuid.NewString(), SubjectKind: "service_account", SubjectID: "test-runner", TokenSHA256: device.HashCredential("key"), TenantID: uuid.NewString()}}) + if err != nil { + t.Fatal(err) + } + admission := &countingStreamAdmission{} + handler, err := NewHandler(eventOnlyStore{}, auth, "codex", WithExecution(admission)) + if err != nil { + t.Fatal(err) + } + request := httptest.NewRequest(http.MethodPost, "/v1/agents/sessions", strings.NewReader(creationBody)) + request.Header.Set("Authorization", "Bearer key") + request.Header.Set("OpenAI-Beta", "agents=v1") + response := httptest.NewRecorder() + handler.ServeHTTP(response, request) + if response.Code != http.StatusServiceUnavailable || admission.calls.Load() != 0 { + t.Fatal("unsupported stream store reached creation", response.Code, admission.calls.Load(), response.Body.String()) + } +} diff --git a/services/agents-api/internal/api/stream.go b/services/agents-api/internal/api/stream.go index 175f75930..532d6af15 100644 --- a/services/agents-api/internal/api/stream.go +++ b/services/agents-api/internal/api/stream.go @@ -51,36 +51,32 @@ func (h *Handler) streamEvents(w http.ResponseWriter, r *http.Request) { writeStoreError(w, r, err) return } - h.serveSessionEvents(w, r, events, session, cursor, nil, http.StatusOK) + h.serveSessionEvents(w, r, events, session, cursor, nil, http.StatusOK, nil) } -func (h *Handler) serveSessionEvents(w http.ResponseWriter, r *http.Request, events eventStore, session store.Session, cursor int64, initial *v1.SessionEvent, status int) { +// streamSettlement reads a creation stream's committed Session projection and +// the Session event cursor from one snapshot, reporting whether it is settled. +type streamSettlement func(context.Context) (settled bool, cursor int64, err error) + +// serveSessionEvents streams committed changes after cursor. GET streams pass a +// nil settlement and stay live-only until disconnect, deletion or failure. +// +// A fresh creation response ends right after it sends a settling Session event +// and never sends later events: an agent.session.idle recorded when a Turn ends +// or an input reservation stops being pending, or any agent.session.failed. +// Settlements that record no event, such as a reservation cancelled while its +// Session is already idle, use the projection: after an empty drain the stream +// reads the projection and the event cursor from one snapshot and, if settled, +// sends only events up to that cursor before ending. Later work drained before +// that read can still be sent. The projection is re-read after a sent Session +// status event and otherwise at most once a second. +func (h *Handler) serveSessionEvents(w http.ResponseWriter, r *http.Request, events eventStore, session store.Session, cursor int64, initial *v1.SessionEvent, status int, settlement streamSettlement) { id, tenant := session.ID, tenantID(r) - w.Header().Set("Content-Type", "text/event-stream") - w.Header().Set("Cache-Control", "no-cache") - w.Header().Set("X-Accel-Buffering", "no") - w.WriteHeader(status) - controller := http.NewResponseController(w) - write := func(data []byte) error { - if err := controller.SetWriteDeadline(time.Now().Add(5 * time.Second)); err != nil { - return err - } - if _, err := w.Write(data); err != nil { - return err - } - return controller.Flush() - } - if err := write([]byte(": connected\n\n")); err != nil { + write := openEventStream(w, status) + if write == nil { return } - emit := func(event v1.SessionEvent) error { - payload, err := json.Marshal(event) - if err != nil { - writeStreamFailure(write, id) - return err - } - return write([]byte(fmt.Sprintf("event: %s\ndata: %s\n\n", event.Type, payload))) - } + emit := func(event v1.SessionEvent) error { return emitSessionEvent(write, id, event) } if initial != nil { if err := emit(*initial); err != nil { return @@ -90,6 +86,9 @@ func (h *Handler) serveSessionEvents(w http.ResponseWriter, r *http.Request, eve defer poll.Stop() heartbeat := time.NewTicker(15 * time.Second) defer heartbeat.Stop() + // limit is the snapshot cursor of a settled projection; nothing after it is sent. + recheck, limit := true, int64(-1) + var checked time.Time for { changes, err := events.ListSessionEvents(r.Context(), tenant, id, cursor) if errors.Is(err, store.ErrNotFound) { @@ -100,6 +99,9 @@ func (h *Handler) serveSessionEvents(w http.ResponseWriter, r *http.Request, eve return } for _, change := range changes { + if limit >= 0 && change.Sequence > limit { + return + } event, err := streamResponse(session, change, h.executorURL) if err != nil { writeStreamFailure(write, id) @@ -109,10 +111,35 @@ func (h *Handler) serveSessionEvents(w http.ResponseWriter, r *http.Request, eve return } cursor = change.Sequence + if settlement != nil && settlingEvent(change) { + return + } + recheck = recheck || sessionStatusEvent(change.Event.Type) + } + if limit >= 0 && (cursor >= limit || len(changes) == 0) { + return } if len(changes) > 0 { continue } + if settlement != nil && (recheck || time.Since(checked) >= time.Second) { + recheck, checked = false, time.Now() + settled, snapshot, err := settlement(r.Context()) + if errors.Is(err, store.ErrNotFound) || r.Context().Err() != nil { + return + } + if err != nil { + writeStreamFailure(write, id) + return + } + if settled { + if cursor >= snapshot { + return + } + limit = snapshot + continue + } + } select { case <-r.Context().Done(): return @@ -125,15 +152,68 @@ func (h *Handler) serveSessionEvents(w http.ResponseWriter, r *http.Request, eve } } +// openEventStream writes the event-stream headers and the connection comment. It +// returns nil when that write fails. +func openEventStream(w http.ResponseWriter, status int) func([]byte) error { + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-cache") + w.Header().Set("X-Accel-Buffering", "no") + w.WriteHeader(status) + controller := http.NewResponseController(w) + write := func(data []byte) error { + if err := controller.SetWriteDeadline(time.Now().Add(5 * time.Second)); err != nil { + return err + } + if _, err := w.Write(data); err != nil { + return err + } + return controller.Flush() + } + if write([]byte(": connected\n\n")) != nil { + return nil + } + return write +} + +func emitSessionEvent(write func([]byte) error, session string, event v1.SessionEvent) error { + payload, err := json.Marshal(event) + if err != nil { + writeStreamFailure(write, session) + return err + } + return write([]byte(fmt.Sprintf("event: %s\ndata: %s\n\n", event.Type, payload))) +} + +// settlingEvent reports a recorded event after which a creation stream ends: an +// idle marked settled when recorded, or any failure. A self-hosted connection +// that clears a pending input's action to idle is not marked. +func settlingEvent(change store.SessionChange) bool { + switch change.Event.Type { + case "agent.session.failed": + return true + case "agent.session.idle": + return change.Settled + } + return false +} + +func sessionStatusEvent(eventType string) bool { + switch eventType { + case "agent.session.in_progress", "agent.session.requires_action", "agent.session.idle", "agent.session.failed": + return true + } + return false +} + func streamResponse(session store.Session, change store.SessionChange, executorURL string) (v1.SessionEvent, error) { event := change.Event if change.Turn == nil && change.EnvironmentInputActivity == nil { - return event, nil + return withTurnUsage(event), nil } if change.Turn != nil && strings.HasPrefix(event.Type, "agent.session.turn.") { turn, err := turnResponse(session, *change.Turn) event.Turn = &turn - return event, err + return withTurnUsage(event), err } event.SessionID = "" session.RequiredActions = change.RequiredActions @@ -144,6 +224,16 @@ func streamResponse(session store.Session, change store.SessionChange, executorU return event, err } +// withTurnUsage mirrors a terminal Turn snapshot's usage at the event's top +// level, including null when unknown. It never derives or sums counters. +func withTurnUsage(event v1.SessionEvent) v1.SessionEvent { + event.Usage = nil + if v1.TerminalTurnEvent(event.Type) && event.Turn != nil { + event.Usage = event.Turn.Usage + } + return event +} + func writeStreamFailure(write func([]byte) error, session string) { event := v1.SessionEvent{Type: "error", EventID: uuid.NewString(), SessionID: session, Error: &v1.StreamError{Code: "stream_interrupted", Type: "server_error", Message: "The live stream was interrupted. Reconnect and retrieve the Session and its saved Items to recover."}} diff --git a/services/agents-api/internal/api/stream_test.go b/services/agents-api/internal/api/stream_test.go index 8f3338b2a..75b4c2020 100644 --- a/services/agents-api/internal/api/stream_test.go +++ b/services/agents-api/internal/api/stream_test.go @@ -38,6 +38,15 @@ func (f *streamFixture) SessionEventCursor(context.Context, string, string) (int return 10, nil } +func (f *streamFixture) SessionStreamSnapshot(_ context.Context, tenant, id string) (store.Session, int64, error) { + f.mu.Lock() + defer f.mu.Unlock() + if tenant != f.session.TenantID || id != f.session.ID { + return store.Session{}, 0, store.ErrNotFound + } + return f.session, 10, nil +} + func (f *streamFixture) ListSessionEvents(_ context.Context, _, _ string, cursor int64) ([]store.SessionChange, error) { f.mu.Lock() defer f.mu.Unlock() @@ -135,3 +144,43 @@ func TestLiveStreamAuthDisconnectRecoveryAndServerDeadline(t *testing.T) { } _ = response.Body.Close() } + +func TestTerminalTurnEventsMirrorTurnUsage(t *testing.T) { + session := store.Session{ID: "session", Configuration: json.RawMessage(`{"agent":{"id":"agent_test","model":"model","tools":[]},"environment":{"type":"none"}}`)} + measured := json.RawMessage(`{"input_tokens":7,"input_tokens_details":{"cached_tokens":2},"output_tokens":3,"output_tokens_details":{"reasoning_tokens":1},"total_tokens":10}`) + child := &v1.Turn{ID: "child", Status: "cancelled"} + for _, test := range []struct { + change store.SessionChange + want string + }{ + {store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.turn.completed"}, Turn: &store.Turn{ID: "turn", Status: store.TurnCompleted, Usage: measured}}, string(measured)}, + {store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.turn.failed"}, Turn: &store.Turn{ID: "turn", Status: store.TurnFailed}}, "null"}, + // Child Turn snapshots are rendered when recorded. + {store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.turn.cancelled", Turn: child}}, "null"}, + {store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.turn.in_progress"}, Turn: &store.Turn{ID: "turn", Status: store.TurnInProgress, Usage: measured}}, ""}, + {store.SessionChange{Event: v1.SessionEvent{Type: "agent.session.idle"}, Turn: &store.Turn{ID: "turn", Status: store.TurnCompleted, Usage: measured}, SessionUsage: measured}, ""}, + } { + event, err := streamResponse(session, test.change, "") + if err != nil { + t.Fatal(err) + } + raw, err := json.Marshal(event) + if err != nil { + t.Fatal(err) + } + var fields map[string]json.RawMessage + if err := json.Unmarshal(raw, &fields); err != nil { + t.Fatal(err) + } + usage, present := fields["usage"] + if present != (test.want != "") || (present && string(usage) != test.want) { + t.Fatal("terminal usage does not mirror the Turn snapshot", test.change.Event.Type, string(raw)) + } + if present { + var turn map[string]json.RawMessage + if json.Unmarshal(fields["turn"], &turn) != nil || string(turn["usage"]) != string(usage) { + t.Fatal("top-level usage differs from the rendered Turn", string(raw)) + } + } + } +} diff --git a/services/agents-api/internal/store/creation_stream_settlement_public_test.go b/services/agents-api/internal/store/creation_stream_settlement_public_test.go new file mode 100644 index 000000000..4d135d04a --- /dev/null +++ b/services/agents-api/internal/store/creation_stream_settlement_public_test.go @@ -0,0 +1,242 @@ +package store_test + +import ( + "bufio" + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/api" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" + "github.com/google/uuid" +) + +// storeAdmission admits creation input through the Store without a Worker, so +// reservations stay pending until the test settles them. +type storeAdmission struct{ *store.Store } + +type sseLines struct { + lines chan string + stop context.CancelFunc +} + +// openStream reads non-empty SSE lines until the server ends the response. +func openStream(t *testing.T, server *httptest.Server, token, method, path, body, key string) sseLines { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + request, err := http.NewRequestWithContext(ctx, method, server.URL+path, strings.NewReader(body)) + if err != nil { + t.Fatal(err) + } + request.Header.Set("Authorization", "Bearer "+token) + request.Header.Set("OpenAI-Beta", "agents=v1") + if key != "" { + request.Header.Set("Idempotency-Key", key) + } + response, err := server.Client().Do(request) + if err != nil { + t.Fatal(err) + } + if response.StatusCode/100 != 2 || response.Header.Get("Content-Type") != "text/event-stream" { + t.Fatal("stream was not opened", response.StatusCode) + } + lines := make(chan string, 64) + go func() { + defer close(lines) + defer response.Body.Close() + scanner := bufio.NewScanner(response.Body) + for scanner.Scan() { + if line := scanner.Text(); line != "" { + lines <- line + } + } + }() + return sseLines{lines: lines, stop: cancel} +} + +func (s sseLines) next(t *testing.T) string { + t.Helper() + select { + case line, ok := <-s.lines: + if !ok { + t.Fatal("stream ended early") + } + return line + case <-time.After(5 * time.Second): + t.Fatal("stream line timed out") + } + return "" +} + +func (s sseLines) ended(t *testing.T, within time.Duration) { + t.Helper() + select { + case line, ok := <-s.lines: + if ok { + t.Fatal("stream sent more after settlement", line) + } + case <-time.After(within): + t.Fatal("stream did not end once the Session settled") + } +} + +// drain reads lines until the stream is quiet, failing if it ends. +func (s sseLines) drain(t *testing.T) map[string]bool { + t.Helper() + seen := map[string]bool{} + for { + select { + case line, ok := <-s.lines: + if !ok { + t.Fatal("stream ended early") + } + seen[line] = true + case <-time.After(500 * time.Millisecond): + return seen + } + } +} + +func (s sseLines) open(t *testing.T) { + t.Helper() + select { + case line, ok := <-s.lines: + t.Fatal("stream changed while it should wait", line, ok) + case <-time.After(1500 * time.Millisecond): + } +} + +// A self_hosted creation without input ends after its created snapshot, and a +// same-key stream retry ends at once even while later input is pending. A fresh +// creation stream whose initial reservation is cancelled without a Session event +// ends through the committed projection, while GET stays open. +func TestCreationStreamPublicLifetimes(t *testing.T) { + s, pool := store.NewTestStore(t) + tenant, token := uuid.NewString(), uuid.NewString() + auth, err := api.NewAuthenticator([]api.APIKey{{OrganizationID: "test-org", ProjectID: tenant, SubjectKind: "service_account", SubjectID: "test-runner", TokenSHA256: device.HashCredential(token), TenantID: tenant}}) + if err != nil { + t.Fatal(err) + } + handler, err := api.NewHandler(s, auth, "codex", api.WithEnvironmentRemoteURL("https://offline-executor.example"), api.WithExecution(storeAdmission{s})) + if err != nil { + t.Fatal(err) + } + var active atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + active.Add(1) + defer active.Add(-1) + handler.ServeHTTP(w, r) + })) + defer server.Close() + lease, err := s.AcquireExecutionLease(t.Context()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = lease.Close(context.Background()) }() + writer := lease.Store() + connect := func(environment string) { + t.Helper() + generation := uuid.NewString() + if err := writer.ReplaceEnvironmentConnection(t.Context(), tenant, environment, generation); err != nil { + t.Fatal(err) + } + if err := writer.ObserveEnvironmentConnection(t.Context(), tenant, environment, generation, 1, true); err != nil { + t.Fatal(err) + } + } + type createdEvent struct { + Session struct { + ID string `json:"id"` + Status string `json:"status"` + Environment struct { + ID string `json:"id"` + } `json:"environment"` + } `json:"session"` + } + readCreated := func(stream sseLines, status string) createdEvent { + t.Helper() + if line := stream.next(t); line != ": connected" { + t.Fatal(line) + } + if line := stream.next(t); line != "event: agent.session.created" { + t.Fatal(line) + } + var event createdEvent + if data, ok := strings.CutPrefix(stream.next(t), "data: "); !ok || json.Unmarshal([]byte(data), &event) != nil || event.Session.Status != status { + t.Fatal("invalid created snapshot", data) + } + return event + } + const idle = `{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"},"stream":true}` + const initial = `{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"},"input":"initial","stream":true}` + + created := openStream(t, server, token, http.MethodPost, "/v1/agents/sessions", idle, "no-input") + defer created.stop() + first := readCreated(created, "idle") + created.ended(t, 5*time.Second) + + connect(first.Session.Environment.ID) + if _, err := s.ReserveEnvironmentInput(t.Context(), tenant, first.Session.ID, "later", []store.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"later"}`)}}); err != nil { + t.Fatal(err) + } + if current, err := s.GetSession(t.Context(), tenant, first.Session.ID); err != nil || !current.PendingInput { + t.Fatal("later input is not pending", err) + } + retry := openStream(t, server, token, http.MethodPost, "/v1/agents/sessions", idle, "no-input") + defer retry.stop() + if line := retry.next(t); line != ": connected" { + t.Fatal(line) + } + retry.ended(t, 5*time.Second) + + fresh := openStream(t, server, token, http.MethodPost, "/v1/agents/sessions", initial, "initial-input") + defer fresh.stop() + second := readCreated(fresh, "requires_action") + session := second.Session.ID + if line := fresh.next(t); line != "event: agent.session.requires_action" { + t.Fatal(line) + } + fresh.next(t) + live := openStream(t, server, token, http.MethodGet, "/v1/agents/sessions/"+session+"/events", "", "") + defer live.stop() + if line := live.next(t); line != ": connected" { + t.Fatal(line) + } + // The connection clears the action to idle; the initial input is still pending. + connect(second.Session.Environment.ID) + if seen := fresh.drain(t); !seen["event: agent.session.idle"] || !seen["event: agent.session.environment.connected"] || seen["event: agent.session.failed"] { + t.Fatal("unexpected connection events", seen) + } + fresh.open(t) + var reservation string + if err := pool.QueryRow(t.Context(), "SELECT id FROM environment_input_reservations WHERE session_id=$1 AND is_initial", session).Scan(&reservation); err != nil { + t.Fatal(err) + } + cursor, err := s.SessionEventCursor(t.Context(), tenant, session) + if err != nil { + t.Fatal(err) + } + if settled, err := s.CancelEnvironmentInput(t.Context(), tenant, session, reservation); err != nil || settled.State != store.EnvironmentInputCancelled { + t.Fatal(settled.State, err) + } + if after, err := s.SessionEventCursor(t.Context(), tenant, session); err != nil || after != cursor { + t.Fatal("cancellation recorded a Session event; this case needs a silent settlement", after, cursor, err) + } + fresh.ended(t, 5*time.Second) + live.drain(t) + live.open(t) + live.stop() + deadline := time.Now().Add(3 * time.Second) + for active.Load() != 0 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + if active.Load() != 0 { + t.Fatal("stream handler leaked") + } +} diff --git a/services/agents-api/internal/store/environment_initial_input_test.go b/services/agents-api/internal/store/environment_initial_input_test.go index 73fdcd130..041d7b68b 100644 --- a/services/agents-api/internal/store/environment_initial_input_test.go +++ b/services/agents-api/internal/store/environment_initial_input_test.go @@ -87,10 +87,21 @@ func TestEnvironmentInitialInputCreationRetainsCursorIdentityAndPromotion(t *tes input := environmentInput("initial", kind, "/workspace") input.InitialInputs = []Input{messageInput("first"), messageInput("second")} creation, err := s.CreateSessionStream(t.Context(), tenant, input) - if err != nil || !creation.Created || creation.Cursor != 0 || creation.Session.LastTurn != nil || creation.Session.EnvironmentInputActivity != nil { - t.Fatal("creation snapshot borrowed initial work", creation, err) + if err != nil || !creation.Created || creation.Cursor != 0 || creation.Session.LastTurn != nil { + t.Fatal("creation started initial work before native readiness", creation, err) } session := creation.Session + // The snapshot is the committed JSON projection: self-hosted input requests + // its connection, while hosted initial provisioning has no caller action. + activityStatus := func(value Session) string { + if value.EnvironmentInputActivity == nil { + return "" + } + return value.EnvironmentInputActivity.Status + } + if want := map[string]string{"self_hosted": "requires_action", "openai_hosted": ""}[kind]; activityStatus(session) != want || !session.PendingInput { + 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) @@ -113,7 +124,7 @@ func TestEnvironmentInitialInputCreationRetainsCursorIdentityAndPromotion(t *tes } other, _ := testStore(t) retry, err := other.CreateSessionStream(t.Context(), tenant, input) - if err != nil || retry.Created || retry.Cursor != int64(expectedEvents) || retry.Session.ID != session.ID || retry.Session.EnvironmentInputActivity != nil { + if err != nil || retry.Created || retry.Cursor != int64(expectedEvents) || retry.Session.ID != session.ID || retry.Session.EnvironmentInputActivity != nil || retry.Session.LastTurn != nil { t.Fatal("retry changed creation cursor/snapshot", retry, err) } if got := initialEnvironmentReservation(t, other, pool, tenant, session.ID); !reflect.DeepEqual(got, reservation) { @@ -148,7 +159,7 @@ func TestEnvironmentInitialInputCreationRetainsCursorIdentityAndPromotion(t *tes t.Fatal("initial batch did not promote", promoted, err) } active := requireEnvironmentInputActivity(t, s, tenant, session.ID, "", "") - if active.LastTurn == nil || active.LastTurn.Status != TurnInProgress { + if active.LastTurn == nil || active.LastTurn.Status != TurnInProgress || active.PendingInput { t.Fatal("promotion did not claim its Turn", active.LastTurn) } replay, err := writer.PromoteEnvironmentInput(t.Context(), tenant, session.ID, reservation.ID) @@ -163,6 +174,20 @@ func TestEnvironmentInitialInputCreationRetainsCursorIdentityAndPromotion(t *tes if err != nil || !reflect.DeepEqual(events, after[:len(events)]) { t.Fatal("initial snapshot changed after promotion", err) } + // Promotion starts the Turn in the official order: turn.created, the + // user input Item, Session activity, then turn.in_progress. + first := map[string]int{} + for index, change := range after[len(events):] { + if _, seen := first[change.Event.Type]; !seen { + first[change.Event.Type] = index + 1 + } + } + order := []string{"agent.session.turn.created", "agent.session.turn.item.added", "agent.session.in_progress", "agent.session.turn.in_progress"} + for index := range order { + if first[order[index]] == 0 || (index > 0 && first[order[index]] < first[order[index-1]]) { + t.Fatal("promoted Turn events are out of order", order[index], first) + } + } }) } } @@ -192,7 +217,7 @@ func TestEnvironmentInitialInputExpiryHasNoTurnAndCannotReplay(t *testing.T) { reservation = initialEnvironmentReservation(t, s, pool, tenant, session.ID) } failed := requireEnvironmentInputActivity(t, s, tenant, session.ID, "failed", "") - if failed.Environment.Status != "pending" || failed.LastTurn != nil || !failed.EnvironmentInputActivity.LastActiveAt.Equal(*reservation.SettledAt) { + if failed.Environment.Status != "pending" || failed.LastTurn != nil || failed.PendingInput || !failed.EnvironmentInputActivity.LastActiveAt.Equal(*reservation.SettledAt) { t.Fatal("input expiry changed Environment/Turn", failed) } events, err := s.ListSessionEvents(t.Context(), tenant, session.ID, 0) diff --git a/services/agents-api/internal/store/environment_input_activity.go b/services/agents-api/internal/store/environment_input_activity.go index 9a3974ebe..2750709d7 100644 --- a/services/agents-api/internal/store/environment_input_activity.go +++ b/services/agents-api/internal/store/environment_input_activity.go @@ -21,13 +21,23 @@ type EnvironmentInputActivity struct { } func environmentInputActivity(ctx context.Context, q *sqlc.Queries, session pgtype.UUID) (*EnvironmentInputActivity, error) { + activity, _, err := environmentInputState(ctx, q, session) + return activity, err +} + +// environmentInputState also reports whether the latest reservation is still +// pending and can start a Turn, including a provisioning hosted initial input +// that has no public activity. While a Turn is active or newer than it, the +// reservation is not reported. +func environmentInputState(ctx context.Context, q *sqlc.Queries, session pgtype.UUID) (*EnvironmentInputActivity, bool, error) { row, err := q.GetEnvironmentInputActivity(ctx, session) if errors.Is(err, pgx.ErrNoRows) { - return nil, nil + return nil, false, nil } if err != nil { - return nil, err + return nil, false, err } + pending := row.State == EnvironmentInputPending activity := &EnvironmentInputActivity{Status: "idle", LastActiveAt: row.CreatedAt.Time} if row.SettledAt.Valid { activity.LastActiveAt = row.SettledAt.Time @@ -43,13 +53,13 @@ func environmentInputActivity(ctx context.Context, q *sqlc.Queries, session pgty // No Turn has started. The pinned Session contract permits idle while a // hosted Environment provisions; neither a caller action nor an invented // in-progress/idle transition is appropriate here. - return nil, nil + return nil, pending, nil } - if row.State == EnvironmentInputPending && row.EnvironmentType != "openai_hosted" && row.ConnectionStatus != "connected" { + if pending && row.EnvironmentType != "openai_hosted" && row.ConnectionStatus != "connected" { activity.Status = "requires_action" activity.EnvironmentID = uuid.UUID(row.EnvironmentID.Bytes).String() } - return activity, nil + return activity, pending, nil } func withEnvironmentInputActivity(ctx context.Context, q *sqlc.Queries, session pgtype.UUID, apply func() error) error { @@ -60,7 +70,7 @@ func withEnvironmentInputActivity(ctx context.Context, q *sqlc.Queries, session if err := apply(); err != nil { return err } - after, err := environmentInputActivity(ctx, q, session) + after, pending, err := environmentInputState(ctx, q, session) if err != nil || after == nil { // Admitted input is represented by the normal Turn and Session events. return err @@ -74,7 +84,7 @@ func withEnvironmentInputActivity(ctx context.Context, q *sqlc.Queries, session } return recordSessionChange(ctx, q, session, SessionChange{ Event: v1.SessionEvent{Type: "agent.session." + after.Status}, - EnvironmentInputActivity: after, SessionUsage: usage, + EnvironmentInputActivity: after, SessionUsage: usage, Settled: !pending, }) } diff --git a/services/agents-api/internal/store/environment_input_activity_test.go b/services/agents-api/internal/store/environment_input_activity_test.go index de9691ad4..c0748dfc9 100644 --- a/services/agents-api/internal/store/environment_input_activity_test.go +++ b/services/agents-api/internal/store/environment_input_activity_test.go @@ -69,6 +69,10 @@ func TestEnvironmentInputActivityWaitsBeforeTurnAndClearsOnConnection(t *testing if !reflect.DeepEqual(first[0], changes[0]) || changes[2].Turn != nil || changes[2].EnvironmentInputActivity.Status != "idle" { t.Fatal("activity snapshot changed or borrowed a Turn") } + // A connection clears the action, but the input is still pending admission. + if !idle.PendingInput || changes[2].Settled || changes[0].Settled { + t.Fatal("pending input reported as settled") + } if err := writer.ObserveEnvironmentConnection(t.Context(), tenant, environment, generation, 1, false); err != nil { t.Fatal(err) } @@ -106,9 +110,13 @@ func TestEnvironmentInputActivitySettlementAndNewerWork(t *testing.T) { } reservation := reserveEnvironmentInput(t, s, tenant, session.ID, "waiting") value, err := s.GetSession(t.Context(), tenant, session.ID) - if err != nil || value.LastTurn.Status != TurnFailed || value.EnvironmentInputActivity.Status != "requires_action" { + if err != nil || value.LastTurn.Status != TurnFailed || value.EnvironmentInputActivity.Status != "requires_action" || !value.PendingInput { t.Fatal("prior failure hid waiting input", err) } + waitingCursor, err := s.SessionEventCursor(t.Context(), tenant, session.ID) + if err != nil { + t.Fatal(err) + } if state == EnvironmentInputCancelled { _, err = s.CancelEnvironmentInput(t.Context(), tenant, session.ID, reservation.ID) } else { @@ -130,7 +138,13 @@ func TestEnvironmentInputActivitySettlementAndNewerWork(t *testing.T) { if err != nil { t.Fatal(err) } - requireEnvironmentInputActivity(t, s, tenant, session.ID, "idle", "") + // The settled later reservation no longer counts as pending input, so + // creation streams can end instead of waiting for work that cannot start. + settled := requireEnvironmentInputActivity(t, s, tenant, session.ID, "idle", "") + events, err := s.ListSessionEvents(t.Context(), tenant, session.ID, waitingCursor) + if err != nil || settled.PendingInput || len(events) != 1 || events[0].Event.Type != "agent.session.idle" || events[0].Turn != nil || !events[0].Settled { + t.Fatal("settled reservation remained pending", settled.EnvironmentInputActivity, events, err) + } cursor, err := s.SessionEventCursor(t.Context(), tenant, session.ID) if err != nil { t.Fatal(err) diff --git a/services/agents-api/internal/store/function_state.go b/services/agents-api/internal/store/function_state.go index cec067d78..c8e759028 100644 --- a/services/agents-api/internal/store/function_state.go +++ b/services/agents-api/internal/store/function_state.go @@ -61,8 +61,10 @@ func recordSessionActivity(ctx context.Context, q *sqlc.Queries, row sqlc.Turn, if err != nil { return err } + // Mark the idle or failure of an ending Turn; a reservation made during its + // Artifact capture is newer work. return recordSessionChange(ctx, q, row.SessionID, SessionChange{ Event: v1.SessionEvent{Type: "agent.session." + status}, Turn: &turn, - SessionUsage: usage, RequiredActions: actions, + SessionUsage: usage, RequiredActions: actions, Settled: terminalStatus(row.Status), }) } diff --git a/services/agents-api/internal/store/scheduling.go b/services/agents-api/internal/store/scheduling.go index a98c7fb1d..c6fc8442c 100644 --- a/services/agents-api/internal/store/scheduling.go +++ b/services/agents-api/internal/store/scheduling.go @@ -3,6 +3,7 @@ package store import ( "context" "errors" + "fmt" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/db/sqlc" "github.com/google/uuid" @@ -86,39 +87,81 @@ func (s *Store) sessionActivity(ctx context.Context, session Session, err error) if err != nil { return Session{}, err } - id, _ := parseID(session.ID) - tenant, _ := parseID(session.TenantID) + err = pgx.BeginTxFunc(ctx, s.pool, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}, func(tx pgx.Tx) error { + var err error + session, err = readSessionActivity(ctx, s.queries.WithTx(tx), session) + return err + }) + return session, err +} + +// SessionStreamSnapshot reads the projection that GetSession returns and the +// committed Session event cursor from one database snapshot. +func (s *Store) SessionStreamSnapshot(ctx context.Context, tenantID, sessionID string) (Session, int64, error) { + tenant, err := parseID(tenantID) + if err != nil { + return Session{}, 0, err + } + id, err := parseID(sessionID) + if err != nil { + return Session{}, 0, err + } + var session Session + var cursor int64 err = pgx.BeginTxFunc(ctx, s.pool, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}, func(tx pgx.Tx) error { q := s.queries.WithTx(tx) - environment, err := q.GetSessionEnvironment(ctx, sqlc.GetSessionEnvironmentParams{TenantID: tenant, ID: id}) - if err == nil { - value, err := environmentFromRow(environment.Environment, environment.TenantID, environment.Configuration, nil) - if err != nil { - return err - } - session.Environment = &value - session.EnvironmentInputActivity, err = environmentInputActivity(ctx, q, id) - if err != nil { - return err - } - } else if !errors.Is(err, pgx.ErrNoRows) { - return err - } - row, err := q.GetLatestSessionTurn(ctx, id) - if errors.Is(err, pgx.ErrNoRows) { - return nil - } + row, err := q.GetSession(ctx, sqlc.GetSessionParams{TenantID: tenant, ID: id}) if err != nil { return err } - turn := turnFromRow(row) - session.LastTurn = &turn - session.RequiredActions, err = functionActions(ctx, q, row) - if err != nil { + if session, err = sessionFromRow(row); err != nil { return err } - session.Usage, err = q.SessionTokenUsage(ctx, id) + cursor = row.EventSequence + session, err = readSessionActivity(ctx, q, session) return err }) + if errors.Is(err, pgx.ErrNoRows) { + return Session{}, 0, ErrNotFound + } + if err != nil { + return Session{}, 0, fmt.Errorf("read session stream snapshot: %w", err) + } + return session, cursor, nil +} + +// readSessionActivity adds the Environment, reservation activity and latest Turn +// projection within the caller's snapshot. +func readSessionActivity(ctx context.Context, q *sqlc.Queries, session Session) (Session, error) { + id, _ := parseID(session.ID) + tenant, _ := parseID(session.TenantID) + environment, err := q.GetSessionEnvironment(ctx, sqlc.GetSessionEnvironmentParams{TenantID: tenant, ID: id}) + if err == nil { + value, err := environmentFromRow(environment.Environment, environment.TenantID, environment.Configuration, nil) + if err != nil { + return session, err + } + session.Environment = &value + session.EnvironmentInputActivity, session.PendingInput, err = environmentInputState(ctx, q, id) + if err != nil { + return session, err + } + } else if !errors.Is(err, pgx.ErrNoRows) { + return session, err + } + row, err := q.GetLatestSessionTurn(ctx, id) + if errors.Is(err, pgx.ErrNoRows) { + return session, nil + } + if err != nil { + return session, err + } + turn := turnFromRow(row) + session.LastTurn = &turn + session.RequiredActions, err = functionActions(ctx, q, row) + if err != nil { + return session, err + } + session.Usage, err = q.SessionTokenUsage(ctx, id) return session, err } diff --git a/services/agents-api/internal/store/session_creation_stream.go b/services/agents-api/internal/store/session_creation_stream.go index 757a2aa2d..dc447988f 100644 --- a/services/agents-api/internal/store/session_creation_stream.go +++ b/services/agents-api/internal/store/session_creation_stream.go @@ -2,9 +2,12 @@ package store import "context" -// SessionCreation starts observation at the Session upsert. Session contains the -// resource fields before initial work, not a later activity snapshot. Only a new -// creation emits that snapshot; retries observe future changes from Cursor. +// SessionCreation starts observation at the Session upsert. For a new creation, +// CreateSessionStream returns the committed Session projection that +// CreateSession returns, read after the creation commits; Cursor still precedes +// the initial inputs, so their events remain observable exactly once. Retries +// and FindSessionCreation return only the resource row and its cursor. Only a new creation emits a created snapshot and +// streams from Cursor; a stream retry of an existing creation sends no events. type SessionCreation struct { Session Session Created bool @@ -15,5 +18,14 @@ type SessionCreation struct { // The upsert returns its event cursor while holding the Session write lock, before // initial inputs commit. No post-commit cursor lookup may skip those inputs. func (s *Store) CreateSessionStream(ctx context.Context, tenant string, input CreateSessionInput) (SessionCreation, error) { - return s.createSession(ctx, tenant, input) + result, err := s.createSession(ctx, tenant, input) + if err != nil || !result.Created { + // A stream retry sends no events, so it needs no projection read. + return result, err + } + result.Session, err = s.sessionActivity(ctx, result.Session, nil) + if err != nil { + return SessionCreation{}, err + } + return result, nil } diff --git a/services/agents-api/internal/store/session_creation_stream_test.go b/services/agents-api/internal/store/session_creation_stream_test.go index d3cc70edf..dd50dcf51 100644 --- a/services/agents-api/internal/store/session_creation_stream_test.go +++ b/services/agents-api/internal/store/session_creation_stream_test.go @@ -2,6 +2,8 @@ package store import ( "context" + "encoding/json" + "errors" "sync" "testing" @@ -44,18 +46,43 @@ func TestCreationStreamStartsBeforeOwnInputsAndRetriesAtUpsertCursor(t *testing. retries = append(retries, result) } } - if !created.Created || created.Cursor != 0 || created.Session.LastTurn != nil || len(retries) != 7 { - t.Fatal("invalid pre-input creation snapshot", created, retries) + // The snapshot is the committed post-admission projection; the cursor still + // precedes the initial input's events. + if !created.Created || created.Cursor != 0 || created.Session.LastTurn == nil || created.Session.LastTurn.Status != TurnQueued || len(retries) != 7 { + t.Fatal("invalid post-admission creation snapshot", created, retries) } id := created.Session.ID initial, err := s.ListSessionEvents(ctx, tenant, id, created.Cursor) - if err != nil || len(initial) != 3 || initial[0].Event.Type != "agent.session.turn.created" { + if err != nil || len(initial) != 3 { t.Fatal("lost initial events", initial, err) } + for i, kind := range []string{"agent.session.turn.created", "agent.session.turn.item.added", "agent.session.in_progress"} { + if initial[i].Event.Type != kind { + t.Fatal("new Turn events are out of order", i, initial[i].Event.Type) + } + } + if initial[1].Event.Item == nil || initial[1].Event.Item.Role != "user" || initial[2].Turn == nil || initial[2].Turn.Status != TurnQueued { + t.Fatal("invalid initial input or activity snapshot", initial) + } + encoded := func(value Session) string { + raw, err := json.Marshal(value) + if err != nil { + t.Fatal(err) + } + return string(raw) + } + read, err := s.GetSession(ctx, tenant, id) + if err != nil { + t.Fatal(err) + } for _, retry := range retries { if retry.Session.ID != id || retry.Cursor != initial[len(initial)-1].Sequence { t.Fatal(retry) } + // Upsert retries return only the row, without a projection read or input. + if retry.Session.LastTurn != nil || retry.Session.EnvironmentInputActivity != nil || retry.Session.Usage != nil { + t.Fatal("retry read the Session projection", retry.Session) + } if events, err := s.ListSessionEvents(ctx, tenant, id, retry.Cursor); err != nil || len(events) != 0 { t.Fatal("retry replayed initial events", events, err) } @@ -64,6 +91,14 @@ func TestCreationStreamStartsBeforeOwnInputsAndRetriesAtUpsertCursor(t *testing. if err != nil || ordinary.ID != id || ordinary.LastTurn == nil { t.Fatal(ordinary, err) } + // The streamed snapshot is the projection that the JSON response returns. + if encoded(created.Session) != encoded(ordinary) || encoded(ordinary) != encoded(read) { + t.Fatal("streamed and JSON creation projections differ", created.Session, ordinary) + } + snapshot, cursor, err := s.SessionStreamSnapshot(ctx, tenant, id) + if err != nil || encoded(snapshot) != encoded(read) || cursor != initial[len(initial)-1].Sequence || initial[2].Settled { + t.Fatal("stream snapshot differs from the Session read and its cursor", cursor, err) + } transition(t, s, tenant, id, ordinary.LastTurn.ID, TurnQueued, TurnInProgress) transition(t, s, tenant, id, ordinary.LastTurn.ID, TurnInProgress, TurnCompleted) // Completing before the HTTP observer drains does not change its start point. @@ -71,6 +106,16 @@ func TestCreationStreamStartsBeforeOwnInputsAndRetriesAtUpsertCursor(t *testing. if err != nil || len(all) <= len(initial) || all[0].Event.EventID != initial[0].Event.EventID { t.Fatal(all, err) } + // The terminal Turn's idle snapshot is recorded as settled for creation streams. + if last := all[len(all)-1]; last.Event.Type != "agent.session.idle" || !last.Settled { + t.Fatal("terminal idle is not recorded as settled", last.Event.Type, last.Settled) + } + if _, cursor, err := s.SessionStreamSnapshot(ctx, tenant, id); err != nil || cursor != all[len(all)-1].Sequence { + t.Fatal("stream snapshot cursor", cursor, err) + } + if _, _, err := s.SessionStreamSnapshot(ctx, uuid.NewString(), id); !errors.Is(err, ErrNotFound) { + t.Fatal("foreign stream snapshot", err) + } late, err := s.CreateSessionStream(ctx, tenant, input) if err != nil || late.Created || late.Cursor != all[len(all)-1].Sequence { t.Fatal(late, err) @@ -79,10 +124,18 @@ func TestCreationStreamStartsBeforeOwnInputsAndRetriesAtUpsertCursor(t *testing. if err != nil { t.Fatal(err) } + if late.Session.ID != id || late.Session.LastTurn != nil { + t.Fatal("late retry read the Session projection", late.Session.LastTurn) + } future, err := s.ListSessionEvents(ctx, tenant, id, late.Cursor) if err != nil || len(future) != 3 || future[0].Turn.ID != next[0].TurnID { t.Fatal(future, err) } + for i, kind := range []string{"agent.session.turn.created", "agent.session.turn.item.added", "agent.session.in_progress"} { + if future[i].Event.Type != kind { + t.Fatal("later Turn events are out of order", i, future[i].Event.Type) + } + } turns, err := s.ListTurns(ctx, tenant, id, "", 100, true) if err != nil || len(turns.Turns) != 2 { t.Fatal(turns, err) diff --git a/services/agents-api/internal/store/session_environment_snapshot_test.go b/services/agents-api/internal/store/session_environment_snapshot_test.go index 6785a0790..400a0e3d6 100644 --- a/services/agents-api/internal/store/session_environment_snapshot_test.go +++ b/services/agents-api/internal/store/session_environment_snapshot_test.go @@ -44,8 +44,12 @@ func TestSelfHostedCreationSnapshotRetainsEnvironmentAndCursor(t *testing.T) { if value.Created || value.Cursor != events[0].Sequence || snapshot.Environment == nil || snapshot.Environment.ID != environment.ID || string(snapshot.Environment.Configuration) != string(environment.Configuration) { t.Fatal("retry changed Environment or cursor", value) } + } + // Recorded-intent lookup and upsert retries return only the row; they read + // no projection because a stream retry sends no events. + for _, snapshot := range []Session{recovered.Session, retry.Session} { if snapshot.LastTurn != nil || snapshot.EnvironmentInputActivity != nil || snapshot.Usage != nil { - t.Fatal("creation snapshot borrowed later activity", snapshot) + t.Fatal("creation retry borrowed later activity", snapshot) } } changedCreator := input.Creator diff --git a/services/agents-api/internal/store/session_events.go b/services/agents-api/internal/store/session_events.go index 4618af6a2..518ce1ca2 100644 --- a/services/agents-api/internal/store/session_events.go +++ b/services/agents-api/internal/store/session_events.go @@ -24,6 +24,12 @@ type SessionChange struct { SessionUsage json.RawMessage `json:"session_usage,omitempty"` RequiredActions []v1.FunctionCallAction `json:"required_actions,omitempty"` EnvironmentInputActivity *EnvironmentInputActivity `json:"environment_input_activity,omitempty"` + // Settled marks an idle or failed snapshot recorded when a Turn ends, or when + // the latest input reservation stops being pending (expired, cancelled or + // failed). A reservation made while the ending Turn captured Artifacts can + // still be pending and start a later Turn. It is internal, never a wire + // field; snapshots recorded without it read as unsettled. + Settled bool `json:"settled,omitempty"` } func recordSessionChange(ctx context.Context, q *sqlc.Queries, session pgtype.UUID, change SessionChange) error { @@ -71,7 +77,8 @@ func recordTurnChange(ctx context.Context, q *sqlc.Queries, row sqlc.Turn, creat if err := recordSessionChange(ctx, q, row.SessionID, change); err != nil { return err } - if !created && !terminalStatus(row.Status) { + // The admitting input records a new Turn's Session activity after its Items. + if created || !terminalStatus(row.Status) { return nil } return recordSessionActivity(ctx, q, row, nil) diff --git a/services/agents-api/internal/store/sessions.go b/services/agents-api/internal/store/sessions.go index e0a2dc2db..0fe33fe9d 100644 --- a/services/agents-api/internal/store/sessions.go +++ b/services/agents-api/internal/store/sessions.go @@ -47,6 +47,10 @@ type Session struct { RequiredActions []v1.FunctionCallAction Environment *Environment EnvironmentInputActivity *EnvironmentInputActivity + // PendingInput reports that the latest input reservation, read once no Turn + // is active or newer, can still start a Turn. It only supports settlement + // checks and is never rendered. + PendingInput bool } type CreateSessionInput struct { diff --git a/services/agents-api/internal/store/turn_inputs.go b/services/agents-api/internal/store/turn_inputs.go index ba9741446..241dc460b 100644 --- a/services/agents-api/internal/store/turn_inputs.go +++ b/services/agents-api/internal/store/turn_inputs.go @@ -153,11 +153,13 @@ func admitInput(ctx context.Context, q *sqlc.Queries, tenantID string, session p if input.Kind == "tool_result" { return admitFunctionResult(ctx, q, tenantID, session, key, position, input) } + created := false turn, err := q.GetActiveTurn(ctx, session) if errors.Is(err, pgx.ErrNoRows) { if input.Kind == "message" { turn, err = q.CreateTurn(ctx, sqlc.CreateTurnParams{ID: pgtype.UUID{Bytes: uuid.New(), Valid: true}, SessionID: session}) if err == nil { + created = true err = recordTurnChange(ctx, q, turn, true) } } else { @@ -181,6 +183,13 @@ func admitInput(ctx context.Context, q *sqlc.Queries, tenantID string, session p if err := indexInput(ctx, q, session, sequence); err != nil { return InputReceipt{}, err } + if created { + // A new Turn publishes turn.created, then its user input Items, then the + // Session activity, within this transaction. + if err := recordSessionActivity(ctx, q, turn, nil); err != nil { + return InputReceipt{}, err + } + } return inputReceipt(sequence, turn.ID, false), nil } diff --git a/services/agents-api/tests/official_agent_reference_retry.py b/services/agents-api/tests/official_agent_reference_retry.py index 701c388c4..112352452 100644 --- a/services/agents-api/tests/official_agent_reference_retry.py +++ b/services/agents-api/tests/official_agent_reference_retry.py @@ -1,9 +1,7 @@ """Saved-reference retry identity using real PostgreSQL and controlled source mutations.""" import concurrent.futures -import json import sys -import time import httpx2 from openai import OpenAI, ConflictError, NotFoundError, BadRequestError @@ -65,27 +63,14 @@ def mutate(**values): assert transport.post(endpoint, headers=auth, json=spec | bad).status_code == 400 equivalent = spec | {"input":[{"role":"user","content":[{"type":"input_text","text":"initial"}]}], "metadata":{}, "stream":False} assert transport.post(endpoint, headers=auth, json=equivalent).json()["id"] == first.id + # A same-key stream retry sends no events, replays nothing and ends at once. with transport.stream("POST", endpoint, headers=auth, json=spec | {"stream":True}) as stream: - assert stream.status_code == 201 - sessions.events.create(first.id, events=[{"type":"agent.session.input.message","input":[{"role":"user","content":[{"type":"input_text","text":"future-after-retry"}]}]}]) - found = False - deadline = time.monotonic() + 15 - for line in stream.iter_lines(): - assert time.monotonic() < deadline, "future event was not observed" - if not line.startswith("data:"): - continue - event = json.loads(line[5:]) - assert event["type"] != "agent.session.created" - if event["type"] == "agent.session.turn.item.added": - text = json.dumps(event) - assert "initial" not in text, "retry replayed initial work" - if "future-after-retry" in text: - found = True - break - assert found + assert stream.status_code == 201 and stream.headers["content-type"] == "text/event-stream" + assert [line for line in stream.iter_lines() if line] == [": connected"] + sessions.events.create(first.id, events=[{"type":"agent.session.input.message","input":[{"role":"user","content":[{"type":"input_text","text":"future-after-retry"}]}]}]) assert len(list(sessions.turns.list(first.id))) == 1 assert len(list(sessions.items.list(first.id))) == 2 - print("Saved-reference retries: SDK/raw HTTP mutation/deletion, concurrent recovery, initial-input-once, current metadata, restart and future-only SSE passed.") + print("Saved-reference retries: SDK/raw HTTP mutation/deletion, concurrent recovery, initial-input-once, current metadata, restart and eventless stream retries passed.") if __name__ == "__main__": diff --git a/services/agents-api/tests/official_execution.py b/services/agents-api/tests/official_execution.py index 6f670f367..6f88e109b 100644 --- a/services/agents-api/tests/official_execution.py +++ b/services/agents-api/tests/official_execution.py @@ -77,6 +77,9 @@ def until_idle(stream): assert "agent.session.turn." + kind in types, types terminal = next(value for value in first_events if value.type == "agent.session.turn.completed") assert terminal.turn.usage.model_dump() == expected_usage + # Top-level terminal usage mirrors the Turn snapshot; other events omit it. + assert terminal.usage.model_dump() == expected_usage + assert all("usage" not in value.to_dict() for value in first_events if value is not terminal) assert first_events[-1].session.usage.model_dump() == expected_usage text_events = [value for value in first_events if value.type.startswith("agent.session.turn.output_text.")] assert all(value.item_id == answers[0].id and value.output_index == 0 and value.content_index == 0 for value in text_events) diff --git a/services/agents-api/tests/official_self_hosted_initial.py b/services/agents-api/tests/official_self_hosted_initial.py index fe3630abb..82835fb13 100644 --- a/services/agents-api/tests/official_self_hosted_initial.py +++ b/services/agents-api/tests/official_self_hosted_initial.py @@ -94,17 +94,19 @@ def create_once(index): assert created["type"] == "agent.session.created" and waiting["type"] == "agent.session.requires_action" assert set(created) == set(waiting) == {"type", "event_id", "session"} assert created["event_id"] != waiting["event_id"] - check(created["session"], "idle") + # The created snapshot is the committed JSON 201 projection; the + # committed connection action still follows from the cursor. + check(created["session"], "requires_action") value = waiting["session"] - assert created["session"]["id"] == value["id"] and created["session"]["environment"] == value["environment"] - assert created["session"]["created_at"] == created["session"]["last_active_at"] + assert created["session"] == value elif mode == "raw_disconnect": with raw.stream("POST", endpoint, json={**request, "stream": True}, headers={"Idempotency-Key": key}) as response: assert response.status_code == 201 and response.headers["content-type"] == "text/event-stream" created = next(json.loads(line[6:]) for line in response.iter_lines() if line.startswith("data: ")) assert created["type"] == "agent.session.created" - check(created["session"], "idle") + check(created["session"], "requires_action") value = sessions.retrieve(created["session"]["id"]).to_dict() + assert created["session"] == value else: value = sessions.create(**request, extra_headers={"Idempotency-Key": key}).to_dict() assert time.monotonic() - began < 8, "creation waited for offline execution" @@ -136,7 +138,11 @@ def create_once(index): elif phase == "expire": case = cases[0] observations, failures = {}, [] - ready = {name: threading.Event() for name in ("sdk", "raw", "creation_retry")} + # A same-key stream retry of the pending creation ends at once without events. + with client() as observer: + with observer.beta.agents.sessions.create(**case["request"], stream=True, extra_headers={"Idempotency-Key": case["key"]}) as stream: + assert list(stream) == [] + ready = {name: threading.Event() for name in ("sdk", "raw")} def observe(name): try: @@ -148,10 +154,7 @@ def observe(name): event = next(json.loads(line[6:]) for line in response.iter_lines() if line.startswith("data: ")) else: with client() as observer: - events = observer.beta.agents.sessions - stream = (events.create(**case["request"], stream=True, extra_headers={"Idempotency-Key": case["key"]}) - if name == "creation_retry" else events.events.stream(case["id"])) - with stream: + with observer.beta.agents.sessions.events.stream(case["id"]) as stream: ready[name].set() event = next(iter(stream)).to_dict() assert event["type"] == "agent.session.failed" and set(event) == {"type", "event_id", "session"} @@ -171,7 +174,7 @@ def observe(name): assert not worker.is_alive(), "public initial failure observer timed out" if failures: raise AssertionError("public initial failure observer failed") from failures[0] - assert observations["sdk"] == observations["raw"] == observations["creation_retry"] + assert observations["sdk"] == observations["raw"] assert retry(case, "failed") == observations["sdk"]["session"] assert api.beta.agents.environments.retrieve(case["environment_id"]).status == "pending" for retained in cases[1:]: diff --git a/services/agents-api/tests/official_session_creation_stream.py b/services/agents-api/tests/official_session_creation_stream.py index 7fe52bae5..d42966f50 100644 --- a/services/agents-api/tests/official_session_creation_stream.py +++ b/services/agents-api/tests/official_session_creation_stream.py @@ -16,50 +16,86 @@ def verify_creation_streams(client, raw, base, headers, foreign, unsupported): saved = client.beta.agents.create(model="test-model", instructions="Saved stream configuration.") forms = ["First", [{"role": "user", "content": [{"type": "input_text", "text": "First"}]}, {"role": "user", "content": [{"type": "input_text", "text": "Second"}]}]] + cancel = [{"type": "agent.session.input.cancel"}] + + def message(text): + return [{"type": "agent.session.input.message", "input": [{"role": "user", "content": [{"type": "input_text", "text": text}]}]}] + for index, initial in enumerate(forms): config = spec if index != 1 else {"agent_id": saved.id, "environment": {"type": "none"}} request = {**config, "input": initial} key = {"Idempotency-Key": str(uuid.uuid4())} + count = 2 if index == 1 else 1 + observed = [] + # The pinned-SDK loop ends because Core closes the creation stream right + # after the initial Turn settles at idle. with sessions.create(**request, stream=True, extra_headers=key) as stream: - first = next(stream) - assert first.type == "agent.session.created" - assert set(first.to_dict()) == {"type", "event_id", "session"} - session = first.session - assert session.status == "idle" and session.last_active_at == session.created_at - if index == 1: - assert session.agent.id == saved.id and session.agent.instructions == saved.instructions - count = 2 if index == 1 else 1 - events = [next(stream) for _ in range(2 + count)] - assert [event.type for event in events] == ["agent.session.turn.created", "agent.session.in_progress"] + ["agent.session.turn.item.added"] * count - assert len({event.event_id for event in [first, *events]}) == 3 + count - assert events[0].turn.id == list(sessions.turns.list(session.id))[0].id - assert events[1].session.status == "in_progress" - items = list(sessions.items.list(session.id, order="asc")) - assert [event.item.to_dict() for event in events[2:]] == [item.to_dict() for item in items] - assert [item.content[0].text for item in items] == (["First", "Second"] if count == 2 else ["First"]) - sessions.events.create(session.id, events=[{"type": "agent.session.input.cancel"}]) - assert [next(stream).type, next(stream).type] == ["agent.session.turn.cancelled", "agent.session.idle"] - # An idle creation stream stays available for a later Turn. - sessions.events.create(session.id, events=[{"type": "agent.session.input.message", "input": [{"role": "user", "content": [{"type": "input_text", "text": "Next"}]}]}]) - assert next(stream).type == "agent.session.turn.created" + for event in stream: + observed.append(event) + if len(observed) == 1: + assert event.type == "agent.session.created" + assert set(event.to_dict()) == {"type", "event_id", "session"} + session = event.session + # The snapshot is the committed post-admission Session, identical + # to the JSON 201 body that a same-key retry recovers. + assert session.status == "in_progress" and session.required_actions == [] and session.usage is None + reply = raw.post(base + "/v1/agents/sessions", headers={**headers, **key}, json=request) + assert reply.status_code == 201 and reply.json() == event.to_dict()["session"] + if index == 1: + assert session.agent.id == saved.id and session.agent.instructions == saved.instructions + elif event.type == "agent.session.in_progress": + sessions.events.create(session.id, events=cancel) + assert len(observed) < 32, "creation stream did not end at idle" + assert [event.type for event in observed] == (["agent.session.created", "agent.session.turn.created"] + ["agent.session.turn.item.added"] * count + + ["agent.session.in_progress", "agent.session.turn.cancelled", "agent.session.idle"]) + assert len({event.event_id for event in observed}) == len(observed) + turns = list(sessions.turns.list(session.id)) + assert len(turns) == 1 and observed[1].turn.id == turns[0].id + assert observed[2 + count].session.status == "in_progress" + items = list(sessions.items.list(session.id, order="asc")) + assert [event.item.to_dict() for event in observed[2:2 + count]] == [item.to_dict() for item in items] + assert [item.content[0].text for item in items] == (["First", "Second"] if count == 2 else ["First"]) + # Terminal Turn events carry the Turn snapshot's usage, null when unknown. + terminal = observed[-2].to_dict() + assert set(terminal) == {"type", "event_id", "session_id", "turn_id", "turn", "usage"} + assert terminal["usage"] is None and terminal["usage"] == terminal["turn"]["usage"] + assert all("usage" not in event.to_dict() for event in observed if event is not observed[-2]) + assert observed[-1].session.status == "idle" + + # GET stays open after idle and observes later Turns; disconnect never cancels them. + with sessions.events.stream(session.id) as live: + sessions.events.create(session.id, events=message("Next")) + assert [next(live).type for _ in range(3)] == ["agent.session.turn.created", "agent.session.turn.item.added", "agent.session.in_progress"] + sessions.events.create(session.id, events=cancel) + assert [next(live).type for _ in range(2)] == ["agent.session.turn.cancelled", "agent.session.idle"] + sessions.events.create(session.id, events=message("After idle")) + assert next(live).type == "agent.session.turn.created" current = sessions.retrieve(session.id) assert current.status == "in_progress", "disconnect cancelled admitted work" assert sessions.create(**request, stream=False, extra_headers=key).id == session.id - assert len(list(sessions.turns.list(session.id))) == 2 - sessions.events.create(session.id, events=[{"type": "agent.session.input.cancel"}]) - # A creation retry observes from the upsert cursor, with no old created/Turn/Item replay. - with raw.stream("POST", base + "/v1/agents/sessions", headers={**headers, **key, "Last-Event-ID": first.event_id}, - json={**request, "stream": True}) as response: - assert response.status_code == 201 and response.headers["content-type"] == "text/event-stream" - lines = response.iter_lines() - assert next(lines) == ": connected" - previous_turns = {turn.id for turn in sessions.turns.list(session.id)} - sessions.events.create(session.id, events=[{"type": "agent.session.input.message", "input": [{"role": "user", "content": [{"type": "input_text", "text": "After retry"}]}]}]) - event = event_data(lines) - assert event["type"] == "agent.session.turn.created" and event["event_id"] != events[0].event_id - assert event["turn"]["id"] not in previous_turns assert len(list(sessions.turns.list(session.id))) == 3 - sessions.events.create(session.id, events=[{"type": "agent.session.input.cancel"}]) + sessions.events.create(session.id, events=cancel) + + # A same-key stream retry of an existing creation admits nothing and ends + # at once, whether the Session is settled or running another client's Turn: + # no created snapshot, Turn, Item or later event is sent. + def retry_stream(): + with raw.stream("POST", base + "/v1/agents/sessions", headers={**headers, **key, "Last-Event-ID": observed[0].event_id}, + json={**request, "stream": True}) as response: + assert response.status_code == 201 and response.headers["content-type"] == "text/event-stream" + assert [line for line in response.iter_lines() if line] == [": connected"] + + retry_stream() + assert len(list(sessions.turns.list(session.id))) == 3 + sessions.events.create(session.id, events=message("After retry")) + assert sessions.retrieve(session.id).status == "in_progress" + retry_stream() + assert len(list(sessions.turns.list(session.id))) == 4 + sessions.events.create(session.id, events=cancel) + # The pinned SDK sees the empty retry stream end without events. + with sessions.create(**request, stream=True, extra_headers=key) as stream: + assert list(stream) == [] + assert len(list(sessions.turns.list(session.id))) == 4 response = raw.get(base + "/v1/agents/sessions/" + session.id, headers={**headers, "Authorization": "Bearer " + foreign}) assert response.status_code == 404 @@ -71,6 +107,8 @@ def verify_creation_streams(client, raw, base, headers, foreign, unsupported): assert set(event) == {"type", "event_id", "session"} and event["type"] == "agent.session.created" recovered = sessions.create(**request, extra_headers=key) assert recovered.id == event["session"]["id"] and recovered.status == "in_progress" + reply = raw.post(base + "/v1/agents/sessions", headers={**headers, **key}, json=request) + assert reply.status_code == 201 and reply.json() == event["session"] assert len(list(sessions.turns.list(recovered.id))) == 1 sessions.events.create(recovered.id, events=[{"type": "agent.session.input.cancel"}]) @@ -86,4 +124,4 @@ def verify_creation_streams(client, raw, base, headers, foreign, unsupported): response = raw.post(unsupported + "/v1/agents/sessions", headers=headers, json={**request, "stream": True}) assert response.status_code == 400 and response.headers["content-type"].startswith("application/json") assert {session.id for session in sessions.list()} == before - print("Creation streams: fixed SDK/raw HTTP, initial admission and later idle continuation, snapshots/order, saved Agents, safe retries, later Turns, disconnect recovery and pre-stream errors passed.") + print("Creation streams: fixed SDK/raw HTTP, initial admission, post-admission snapshots, Turn start order, terminal usage, settlement closure, live GET continuation, saved Agents, immediately ending retries, disconnect recovery and pre-stream errors passed.")