From 34996d630f07a9966217f09e9ab4e06fbba32cbd Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 12:47:43 -0400 Subject: [PATCH] Give Ask Fabro sessions their own event log Step 4 of the legacy executor deletion, second commit. Ask Fabro's sessions were the last writer of `run_events`: a session's creation, its turns and their messages, tool calls and endings went into the run's legacy event log, keyed by the run's sequence. They now have a log of their own. - `run_session_events` (migration `2026091802`): one row per session event, numbered per session from 1, with the owning run, the turn, the event name and its properties. `RunSessionEventStore` appends under the write lock, lists a session from a sequence, names a session's owner from its creation event, deletes a run's sessions with the run, and publishes each committed event to its subscribers. - `fabro_types::SessionEvent`: `seq`, `session_id`, `run_id`, `ts` and a flattened body (`event` naming the kind, `properties` its fields), with the same event names and property shapes the legacy events carried, so the web app and the CLI read the same JSON. The property structs move to `session_event`; `run_event::session` re-exports them under their old names until the legacy event log goes. - The API: `GET /sessions/{id}/events` pages `PaginatedSessionEventList` by the session's own sequence, `GET /sessions/{id}/attach` replays and streams `SessionEvent` frames (subscribed before the replay, so no event falls between the two), the turn stream carries the same frames, and an interrupt answers with the recorded event. The session projection folds `SessionEvent`s; the legacy `find_session_owner` over `run_events` is gone. - The CLI's `run ask` and the web app's session stream read `SessionEvent`; the web runtime no longer accepts the nested legacy envelope shape. The two session resume tests in the server keep failing for a reason this commit does not touch: Ask Fabro reconnects to the run's sandbox from the projection's sandbox instance, which the Petri projection does not carry yet (`VIEWS.md`, the `scope.acquired` gap). Co-Authored-By: Claude Fable 5.1 --- .../app/lib/ask-fabro-runtime.test.ts | 33 +- apps/fabro-web/app/lib/ask-fabro-runtime.ts | 39 +- apps/fabro-web/app/lib/session-stream.ts | 4 +- docs/public/api-reference/fabro-api.yaml | 91 +++- lib/apps/fabro-cli/src/commands/run/ask.rs | 40 +- lib/apps/fabro-server/src/server.rs | 14 +- .../src/server/handler/sessions.rs | 475 ++++++++---------- .../fabro-server/tests/it/api/sessions.rs | 63 +-- lib/components/fabro-store/src/keys.rs | 1 + lib/components/fabro-store/src/lib.rs | 2 + .../src/run_session_event_store.rs | 351 +++++++++++++ .../fabro-store/src/run_sessions.rs | 176 ++++--- .../fabro-store/src/run_summary_store.rs | 156 +----- lib/components/fabro-store/src/slate/mod.rs | 113 +---- lib/foundation/fabro-api/build.rs | 1 + lib/foundation/fabro-api/src/lib.rs | 15 +- .../tests/session_contract_round_trip.rs | 37 +- lib/foundation/fabro-client/src/client.rs | 8 +- .../2026091802_run_session_events.sql | 19 + lib/foundation/fabro-db/src/lib.rs | 6 + lib/foundation/fabro-types/src/lib.rs | 2 + .../fabro-types/src/run_event/session.rs | 168 +------ .../fabro-types/src/session_event.rs | 292 +++++++++++ .../src/.openapi-generator/FILES | 3 + .../fabro-api-client/src/api/sessions-api.ts | 38 +- .../fabro-api-client/src/models/index.ts | 3 + .../models/paginated-session-event-list.ts | 29 ++ .../src/models/session-detail.ts | 2 +- .../src/models/session-event-name.ts | 34 ++ .../src/models/session-event.ts | 33 ++ 30 files changed, 1281 insertions(+), 967 deletions(-) create mode 100644 lib/components/fabro-store/src/run_session_event_store.rs create mode 100644 lib/foundation/fabro-db/migrations/2026091802_run_session_events.sql create mode 100644 lib/foundation/fabro-types/src/session_event.rs create mode 100644 lib/packages/fabro-api-client/src/models/paginated-session-event-list.ts create mode 100644 lib/packages/fabro-api-client/src/models/session-event-name.ts create mode 100644 lib/packages/fabro-api-client/src/models/session-event.ts diff --git a/apps/fabro-web/app/lib/ask-fabro-runtime.test.ts b/apps/fabro-web/app/lib/ask-fabro-runtime.test.ts index fdb360045..d99103324 100644 --- a/apps/fabro-web/app/lib/ask-fabro-runtime.test.ts +++ b/apps/fabro-web/app/lib/ask-fabro-runtime.test.ts @@ -8,43 +8,16 @@ import type { SessionStreamEvent } from "./session-stream"; function event(name: string, properties: Record): SessionStreamEvent { return { - seq: 0, - event: { event: name, properties }, - } as unknown as SessionStreamEvent; -} - -function flattenedEvent( - name: string, - properties: Record, -): SessionStreamEvent { - return { - seq: 0, - id: "evt_1", - ts: "2026-05-22T16:25:34.940200Z", + seq: 1, + session_id: "01HZX6M0P7SE4VJ9Y3X2B8E9QF", run_id: "run_1", + ts: "2026-05-22T16:25:34.940200Z", event: name, properties, } as unknown as SessionStreamEvent; } describe("applyTurnEvent", () => { - test("appends assistant deltas from flattened SSE event envelopes", () => { - const acc = { - activeTextIndex: null, - parts: [], - toolCallIndex: new Map(), - } as Parameters[0]; - - expect( - applyTurnEvent( - acc, - flattenedEvent("run.session.assistant_delta", { delta: "Hello" }), - ), - ).toBe(true); - - expect(acc.parts).toEqual([{ type: "text", text: "Hello" }]); - }); - test("appends assistant deltas into a single streaming text part", () => { const acc = { activeTextIndex: null, diff --git a/apps/fabro-web/app/lib/ask-fabro-runtime.ts b/apps/fabro-web/app/lib/ask-fabro-runtime.ts index 150f30a69..8e4868645 100644 --- a/apps/fabro-web/app/lib/ask-fabro-runtime.ts +++ b/apps/fabro-web/app/lib/ask-fabro-runtime.ts @@ -91,43 +91,22 @@ function snapshot(acc: TurnAccumulator): ChatModelRunResult { return { content: acc.parts.slice() }; } -/** - * Apply a single `EventEnvelope` to the accumulator. Returns true if the - * accumulator changed and a fresh `ChatModelRunResult` should be yielded. - */ -interface NestedRunEvent { - event?: string; - properties?: Record; -} - function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } -function eventPayload(envelope: SessionStreamEvent): { - eventName: string; - props: Record; -} { - const raw = envelope as unknown as Record; - if (typeof raw.event === "string") { - return { - eventName: raw.event, - props: isRecord(raw.properties) ? raw.properties : {}, - }; - } - - const nested = isRecord(raw.event) ? (raw.event as NestedRunEvent) : {}; - return { - eventName: nested.event ?? "", - props: isRecord(nested.properties) ? nested.properties : {}, - }; -} - +/** + * Apply a single `SessionEvent` to the accumulator. Returns true if the + * accumulator changed and a fresh `ChatModelRunResult` should be yielded. + */ export function applyTurnEvent( acc: TurnAccumulator, - envelope: SessionStreamEvent, + event: SessionStreamEvent, ): boolean { - const { eventName, props } = eventPayload(envelope); + const eventName = event.event; + const props: Record = isRecord(event.properties) + ? event.properties + : {}; if (eventName === "run.session.assistant_delta") { const delta = typeof props.delta === "string" ? props.delta : ""; diff --git a/apps/fabro-web/app/lib/session-stream.ts b/apps/fabro-web/app/lib/session-stream.ts index 833391b50..d49096194 100644 --- a/apps/fabro-web/app/lib/session-stream.ts +++ b/apps/fabro-web/app/lib/session-stream.ts @@ -1,6 +1,6 @@ import { SessionsApiAxiosParamCreator, - type EventEnvelope, + type SessionEvent, type SubmitTurnRequest, } from "@qltysh/fabro-api-client"; @@ -9,7 +9,7 @@ import { generatedApiConfiguration, } from "./api-client"; -export type SessionStreamEvent = EventEnvelope; +export type SessionStreamEvent = SessionEvent; type FetchLike = ( input: string, diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index dcbf5da6c..62e5c45ca 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -905,7 +905,9 @@ paths: operationId: listSessionEvents tags: [Sessions] summary: List session events - description: Returns run event envelopes filtered to this session's durable `run.session.*` events. `since_seq` uses the owning run event sequence. + description: >- + Returns this session's events in order. Events are numbered per + session from 1; `since_seq` is the first sequence number to include. parameters: - name: since_seq in: query @@ -922,11 +924,11 @@ paths: maximum: 1000 responses: "200": - description: Session-scoped run events + description: The session's events content: application/json: schema: - $ref: "#/components/schemas/PaginatedEventList" + $ref: "#/components/schemas/PaginatedSessionEventList" "404": description: Session not found headers: @@ -948,7 +950,11 @@ paths: operationId: attachSessionEvents tags: [Sessions] summary: Attach to session events - description: Replays and streams this session's durable `run.session.*` events from the owning run event log. The stream remains open until the client disconnects or the server shuts down. + description: >- + Replays this session's events from `since_seq` (the next unseen event + when omitted) as `SessionEvent` frames, then streams new ones as they + are recorded. The stream remains open until the client disconnects or + the server shuts down. parameters: - name: since_seq in: query @@ -957,7 +963,7 @@ paths: minimum: 1 responses: "200": - description: Streamed session-scoped run events + description: Streamed session events content: text/event-stream: schema: @@ -983,7 +989,11 @@ paths: operationId: submitSessionTurn tags: [Sessions] summary: Submit a session turn - description: Starts a streamed turn immediately. Background turns are not supported in this API version. + description: >- + Starts a streamed turn immediately. The stream carries the turn's + `SessionEvent` frames, from `run.session.turn.started` to the event + that ends the turn. Background turns are not supported in this API + version. requestBody: required: true content: @@ -1056,7 +1066,7 @@ paths: content: application/json: schema: - $ref: "#/components/schemas/EventEnvelope" + $ref: "#/components/schemas/SessionEvent" "404": description: Session not found headers: @@ -8268,9 +8278,9 @@ components: SessionDetail: description: >- - Session metadata plus the highest run event sequence the session's - event stream has reached. The conversation itself is held by the - server's durable session record and is not returned over the API. + Session metadata plus the sequence number of the session's latest + event. The conversation itself is held by the server's durable + session record and is not returned over the API. type: object required: - id @@ -8354,6 +8364,67 @@ components: meta: $ref: "#/components/schemas/PaginationMeta" + SessionEventName: + description: What a session event records. + type: string + enum: + - run.session.created + - run.session.turn.started + - run.session.user_message + - run.session.assistant_delta + - run.session.assistant_message + - run.session.tool_call.started + - run.session.tool_call.completed + - run.session.turn.succeeded + - run.session.turn.failed + - run.session.turn.interrupted + + SessionEvent: + description: >- + One recorded event of an Ask Fabro session: its position in the + session (`seq`, from 1), the session and run it belongs to, when it + was recorded, the event's name and its properties. Every event but + `run.session.created` carries the `turn_id` it belongs to in its + properties. + type: object + required: + - seq + - session_id + - run_id + - ts + - event + - properties + properties: + seq: + type: integer + minimum: 1 + session_id: + $ref: "#/components/schemas/SessionId" + run_id: + type: string + ts: + type: string + format: date-time + event: + $ref: "#/components/schemas/SessionEventName" + properties: + type: object + additionalProperties: true + + PaginatedSessionEventList: + description: A page of a session's events, in sequence order. + type: object + required: + - data + - meta + properties: + data: + type: array + items: + $ref: "#/components/schemas/SessionEvent" + meta: + $ref: "#/components/schemas/PaginationMeta" + WorkflowScheduleSummary: description: Workflow schedule summary shown in workflow lists. type: object diff --git a/lib/apps/fabro-cli/src/commands/run/ask.rs b/lib/apps/fabro-cli/src/commands/run/ask.rs index 2c7819186..2e451875e 100644 --- a/lib/apps/fabro-cli/src/commands/run/ask.rs +++ b/lib/apps/fabro-cli/src/commands/run/ask.rs @@ -1,6 +1,6 @@ use anyhow::{Result, bail}; use fabro_api::types::CreateRunSessionRequest; -use fabro_store::EventEnvelope; +use fabro_types::{SessionEvent, SessionEventBody}; use crate::args::AskArgs; use crate::command_context::CommandContext; @@ -24,21 +24,13 @@ pub(crate) async fn run(args: AskArgs, base_ctx: &CommandContext) -> Result<()> let mut saw_terminal = false; while let Some(event) = stream.next_event().await? { render_event(&event, ctx.json_output())?; - match event.event.event_name() { - "run.session.turn.succeeded" | "run.session.turn.interrupted" => { + match &event.body { + SessionEventBody::TurnSucceeded(_) | SessionEventBody::TurnInterrupted(_) => { saw_terminal = true; } - "run.session.turn.failed" => { + SessionEventBody::TurnFailed(props) => { saw_terminal = true; - terminal_error = Some( - event - .event - .properties()? - .get("error") - .and_then(serde_json::Value::as_str) - .unwrap_or("session turn failed") - .to_string(), - ); + terminal_error = Some(props.error.clone()); } _ => {} } @@ -68,28 +60,18 @@ fn session_title(prompt: &str) -> String { clippy::print_stdout, reason = "The ask command streams assistant output and JSON events to stdout." )] -fn render_event(event: &EventEnvelope, json_output: bool) -> Result<()> { +fn render_event(event: &SessionEvent, json_output: bool) -> Result<()> { if json_output { println!("{}", serde_json::to_string(event)?); return Ok(()); } - match event.event.event_name() { - "run.session.assistant_delta" => { - let properties = event.event.properties()?; - if let Some(delta) = properties.get("delta").and_then(serde_json::Value::as_str) { - print!("{delta}"); - } + match &event.body { + SessionEventBody::AssistantDelta(props) => { + print!("{}", props.delta); } - "run.session.assistant_message" => { - let properties = event.event.properties()?; - if let Some(text) = properties - .get("text") - .and_then(serde_json::Value::as_str) - .filter(|text| !text.is_empty()) - { - println!("{text}"); - } + SessionEventBody::AssistantMessage(props) if !props.text.is_empty() => { + println!("{}", props.text); } _ => {} } diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index b31a0b153..73c875001 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -81,8 +81,8 @@ use fabro_store::platform_records::{ }; use fabro_store::{ ArtifactKey, ArtifactStore, AuthCodeStore, AuthSessionStore, Database, KeyedMutex, - NodeArtifact, PendingInterviewRecord, RunSessionRecordStore, RunSummaryStore, - StageArtifactEntry, StageId, + NodeArtifact, PendingInterviewRecord, RunSessionEventStore, RunSessionRecordStore, + RunSummaryStore, StageArtifactEntry, StageId, }; #[cfg(test)] use fabro_types::BlockedReason; @@ -1056,6 +1056,8 @@ pub(crate) struct AppStores { pub(crate) run_summaries: Arc, /// Ask Fabro conversations, keyed by session id. pub(crate) session_records: Arc, + /// The events of Ask Fabro sessions, numbered per session. + pub(crate) session_events: Arc, pub(crate) auth_codes: Arc, pub(crate) auth_sessions: Arc, pub(crate) automations: Arc, @@ -2426,6 +2428,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result anyhow::Result Ok(DeleteRunOutcome::Preserved(response)), SandboxDeleteOutcome::Cleaned => Ok(DeleteRunOutcome::Deleted), diff --git a/lib/apps/fabro-server/src/server/handler/sessions.rs b/lib/apps/fabro-server/src/server/handler/sessions.rs index 213774d8a..f978d7562 100644 --- a/lib/apps/fabro-server/src/server/handler/sessions.rs +++ b/lib/apps/fabro-server/src/server/handler/sessions.rs @@ -11,24 +11,22 @@ use axum::routing::{get, post}; use axum::{Json, Router}; use chrono::{DateTime, Utc}; use fabro_api::types::{ - CreateRunSessionRequest, PaginatedEventList, PaginationMeta, SubmitTurnRequest, + CreateRunSessionRequest, PaginatedSessionEventList, PaginationMeta, SubmitTurnRequest, }; use fabro_llm::lithos_catalog::Catalog; use fabro_llm::{FabroClient, ModelSelectionError, selection}; use fabro_sandbox::SecretRedactor; use fabro_sandbox::reconnect::reconnect_for_run; -use fabro_store::{ - EventPayload, ProjectedRunSession, RunDatabase, project_run_session, project_run_sessions, -}; +use fabro_store::{ProjectedRunSession, project_run_session, project_run_sessions}; use fabro_tool::fabro_client::ClientBackend; -use fabro_types::run_event::{ - RunSessionAssistantDeltaProps, RunSessionAssistantMessageProps, RunSessionCreatedProps, - RunSessionToolCallCompletedProps, RunSessionToolCallStartedProps, RunSessionTurnFailedCode, - RunSessionTurnFailedProps, RunSessionTurnInterruptedProps, RunSessionTurnStartedProps, - RunSessionTurnSucceededProps, RunSessionUserMessageProps, +use fabro_types::session_event::{ + SessionAssistantDeltaProps, SessionAssistantMessageProps, SessionCreatedProps, + SessionToolCallCompletedProps, SessionToolCallStartedProps, SessionTurnFailedCode, + SessionTurnFailedProps, SessionTurnInterruptedProps, SessionTurnStartedProps, + SessionTurnSucceededProps, SessionUserMessageProps, }; use fabro_types::settings::ModelRef as SettingsModelRef; -use fabro_types::{EventBody, EventEnvelope, RunEvent, RunId, SessionDetail, SessionId, TurnId}; +use fabro_types::{RunId, SessionDetail, SessionEvent, SessionEventBody, SessionId, TurnId}; use fabro_workflow::run_tools::register_named_fabro_run_tools; use fabro_workflow::services::FabroRunToolServices; use lithos_llm::catalog::ProviderId; @@ -44,15 +42,12 @@ use pebble_coding_agent::{CodingAgent, CodingAgentOptions, Error as AgentError, use serde_json::Value; use tokio::sync::broadcast::error::RecvError; use tokio::sync::mpsc; -use tokio_stream::StreamExt; use tokio_stream::wrappers::ReceiverStream; use tokio_util::sync::CancellationToken; use tracing::{error, warn}; use super::super::session_runtime::{InterruptTurnError, SessionTurnLease, StartTurnError}; -use super::super::{ - AppState, EventListParams, PaginationParams, paginate_items, parse_run_id_path, -}; +use super::super::{AppState, PaginationParams, paginate_items, parse_run_id_path}; use crate::error::ApiError; use crate::principal_middleware::RequiredUser; use crate::worker_token::issue_worker_token; @@ -116,13 +111,12 @@ async fn list_run_sessions( Ok(id) => id, Err(response) => return response, }; - let run_store = match open_run_reader(&state, run_id).await { - Ok(store) => store, - Err(response) => return response, - }; - match run_store.list_events().await { + if let Err(response) = ensure_run(&state, run_id).await { + return response; + } + match state.stores.session_events.list_for_run(run_id).await { Ok(events) => { - let mut sessions = project_run_sessions(run_id, &events); + let mut sessions = project_run_sessions(&events); match params.order { RunSessionListOrder::UpdatedDesc => sessions.sort_by(|left, right| { right @@ -159,10 +153,9 @@ async fn create_run_session( Ok(id) => id, Err(response) => return response, }; - let run_store = match open_run(&state, run_id).await { - Ok(store) => store, - Err(response) => return response, - }; + if let Err(response) = ensure_run(&state, run_id).await { + return response; + } let llm_result = match state.resolve_llm_client().await { Ok(result) => result, Err(err) => { @@ -187,10 +180,10 @@ async fn create_run_session( let session_id = SessionId::new(); let now = Utc::now(); let event = match append_run_session_event( - &run_store, + &state, run_id, session_id, - EventBody::RunSessionCreated(RunSessionCreatedProps { + SessionEventBody::Created(SessionCreatedProps { title: request.title, model: Some(model), provider: Some(provider), @@ -203,8 +196,7 @@ async fn create_run_session( Err(err) => return store_error(&err).into_response(), }; - let events = vec![event]; - match project_run_session(run_id, session_id, &events) { + match project_run_session(session_id, &[event]) { Some(session) => (StatusCode::CREATED, Json(session.record)).into_response(), None => ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, @@ -223,7 +215,7 @@ async fn get_session( Ok(id) => id, Err(err) => return err.into_response(), }; - let (_, session) = match load_session_read(&state, session_id).await { + let (_, session) = match load_session(&state, session_id).await { Ok(context) => context, Err(response) => return response, }; @@ -234,29 +226,50 @@ async fn session_method_not_found() -> Response { StatusCode::NOT_FOUND.into_response() } +/// Query parameters for `/sessions/{id}/events`: the first sequence number +/// to include and the page size. +#[derive(serde::Deserialize)] +struct SessionEventListParams { + #[serde(default)] + since_seq: Option, + #[serde(default)] + limit: Option, +} + +impl SessionEventListParams { + fn since_seq(&self) -> u32 { + self.since_seq.unwrap_or(1).max(1) + } + + fn limit(&self) -> usize { + self.limit.unwrap_or(100).clamp(1, 1000) + } +} + async fn list_session_events( _auth: RequiredUser, State(state): State>, Path(id): Path, - Query(params): Query, + Query(params): Query, ) -> Response { let session_id = match parse_session_id(&id) { Ok(id) => id, Err(err) => return err.into_response(), }; - let (_, run_store) = match load_session_run_reader(&state, session_id).await { - Ok(context) => context, - Err(response) => return response, - }; - match run_store - .list_events_for_session_from_with_limit(session_id, params.since_seq(), params.limit()) + if let Err(response) = load_session_owner(&state, session_id).await { + return response; + } + let limit = params.limit(); + match state + .stores + .session_events + .list_from(session_id, params.since_seq(), limit.saturating_add(1)) .await { Ok(mut data) => { - let limit = params.limit(); let has_more = data.len() > limit; data.truncate(limit); - Json(PaginatedEventList { + Json(PaginatedSessionEventList { data, meta: PaginationMeta { has_more, @@ -287,13 +300,12 @@ async fn attach_session_events( Ok(id) => id, Err(err) => return err.into_response(), }; - let (_, run_store) = match load_session_run_reader(&state, session_id).await { - Ok(context) => context, - Err(response) => return response, - }; + if let Err(response) = load_session_owner(&state, session_id).await { + return response; + } let start_seq = match params.since_seq { Some(seq) => seq.max(1), - None => match run_store.last_event_seq().await { + None => match state.stores.session_events.last_seq(session_id).await { Ok(last_seq) => last_seq.map_or(1, |seq| seq.saturating_add(1)), Err(err) => return store_error(&err).into_response(), }, @@ -301,58 +313,57 @@ async fn attach_session_events( let shutdown = state.shutdown_token(); let (sender, receiver) = mpsc::channel(SESSION_SSE_BUFFER_CAPACITY); tokio::spawn(async move { + let events = &state.stores.session_events; let mut next_seq = start_seq; - - loop { - let Ok(replay_batch) = run_store - .list_events_for_session_from_with_limit( - session_id, - next_seq, - ATTACH_REPLAY_BATCH_LIMIT, - ) - .await - else { - return; - }; - let replay_has_more = replay_batch.len() > ATTACH_REPLAY_BATCH_LIMIT; - - for event in replay_batch.into_iter().take(ATTACH_REPLAY_BATCH_LIMIT) { - next_seq = event.seq.saturating_add(1); - if let Some(sse_event) = session_sse_event(&event) { - if !send_attach_sse_event(&sender, &shutdown, sse_event).await { + // Subscribed before the replay, so an event committed between the + // replay's last page and the live loop is not missed: the live loop + // skips what the replay already sent by sequence number. + let mut live = events.subscribe(); + 'replay: loop { + loop { + let Ok(batch) = events + .list_from(session_id, next_seq, ATTACH_REPLAY_BATCH_LIMIT + 1) + .await + else { + return; + }; + let has_more = batch.len() > ATTACH_REPLAY_BATCH_LIMIT; + for event in batch.into_iter().take(ATTACH_REPLAY_BATCH_LIMIT) { + next_seq = event.seq.saturating_add(1); + if !send_attach_sse_event(&sender, &shutdown, session_sse_event(&event)).await { return; } } + if !has_more { + break; + } } - if replay_has_more { - continue; - } - break; - } - - let Ok(mut live_stream) = run_store.watch_events_from(next_seq) else { - return; - }; - let session_id_string = session_id.to_string(); - loop { - tokio::select! { - biased; - () = shutdown.cancelled() => break, - () = sender.closed() => break, - next = live_stream.next() => { - let Some(result) = next else { - return; - }; - let Ok(event) = result else { - return; - }; - if event_matches_session(&event, &session_id_string) { - if let Some(sse_event) = session_sse_event(&event) { - if !send_attach_sse_event(&sender, &shutdown, sse_event).await { + loop { + tokio::select! { + biased; + () = shutdown.cancelled() => return, + () = sender.closed() => return, + next = live.recv() => match next { + Ok(event) => { + if event.session_id != session_id || event.seq < next_seq { + continue; + } + next_seq = event.seq.saturating_add(1); + if !send_attach_sse_event( + &sender, + &shutdown, + session_sse_event(&event), + ) + .await + { return; } } + // The subscriber fell behind the store's buffer; the + // table holds what it missed. + Err(RecvError::Lagged(_)) => continue 'replay, + Err(RecvError::Closed) => return, } } } @@ -374,7 +385,7 @@ async fn submit_turn( Ok(id) => id, Err(err) => return err.into_response(), }; - let (run_id, run_store, session) = match load_session(&state, session_id).await { + let (run_id, session) = match load_session(&state, session_id).await { Ok(context) => context, Err(response) => return response, }; @@ -405,16 +416,16 @@ async fn submit_turn( let (sender, receiver) = mpsc::channel(SESSION_SSE_BUFFER_CAPACITY); let now = Utc::now(); for body in [ - EventBody::RunSessionTurnStarted(RunSessionTurnStartedProps { + SessionEventBody::TurnStarted(SessionTurnStartedProps { turn_id, input: input.clone(), }), - EventBody::RunSessionUserMessage(RunSessionUserMessageProps { + SessionEventBody::UserMessage(SessionUserMessageProps { turn_id, text: input.clone(), }), ] { - match append_and_send_event(&run_store, &sender, run_id, session_id, body, now).await { + match append_and_send_event(&state, &sender, run_id, session_id, body, now).await { Ok(()) => {} Err(err) => { drop(turn_lease); @@ -424,7 +435,7 @@ async fn submit_turn( } tokio::spawn(run_streaming_turn( - state, run_id, run_store, session, turn_id, input, sender, turn_lease, + state, run_id, session, turn_id, input, sender, turn_lease, )); let mut response = Sse::new(ReceiverStream::new(receiver)) .keep_alive(KeepAlive::default()) @@ -448,7 +459,7 @@ async fn interrupt_turn( Ok(id) => id, Err(err) => return err.into_response(), }; - let (run_id, run_store, _) = match load_session(&state, session_id).await { + let (run_id, _) = match load_session(&state, session_id).await { Ok(context) => context, Err(response) => return response, }; @@ -463,10 +474,10 @@ async fn interrupt_turn( } }; match append_run_session_event( - &run_store, + &state, run_id, session_id, - EventBody::RunSessionTurnInterrupted(RunSessionTurnInterruptedProps { + SessionEventBody::TurnInterrupted(SessionTurnInterruptedProps { turn_id, error: Some("Interrupted.".to_string()), }), @@ -488,7 +499,6 @@ async fn interrupt_turn( async fn run_streaming_turn( state: Arc, run_id: RunId, - run_store: RunDatabase, session: ProjectedRunSession, turn_id: TurnId, input: String, @@ -498,11 +508,11 @@ async fn run_streaming_turn( let session_id = session.record.id; if turn_lease.interrupt_requested() { let _ = append_and_send_event( - &run_store, + &state, &sender, run_id, session_id, - EventBody::RunSessionTurnInterrupted(RunSessionTurnInterruptedProps { + SessionEventBody::TurnInterrupted(SessionTurnInterruptedProps { turn_id, error: Some("Interrupted.".to_string()), }), @@ -516,14 +526,14 @@ async fn run_streaming_turn( let runtime_entry = turn_lease.entry(); let mut agent_slot = runtime_entry.lock_agent().await; if agent_slot.is_none() { - match build_agent(&state, run_id, &run_store, &session).await { + match build_agent(&state, run_id, &session).await { Ok(agent) => { *agent_slot = Some(agent); } Err(err) => { error!(error = ?err, session_id = %session_id, turn_id = %turn_id, "Failed to build run-backed session runtime"); let _ = append_and_send_event( - &run_store, + &state, &sender, run_id, session_id, @@ -546,14 +556,14 @@ async fn run_streaming_turn( .expect("session runtime slot should be loaded"); let cancel_token = CancellationToken::new(); turn_lease.attach_cancel_token(&cancel_token); - let model_input = match run_store.state().await { + let model_input = match state.load_run_projection(&run_id).await { Ok(projection) => { let snapshot = build_ask_fabro_run_snapshot(&projection, run_id); build_ask_fabro_turn_input(&input, &snapshot) } Err(err) => { warn!( - error = %err, + error = err.detail(), session_id = %session_id, turn_id = %turn_id, "Failed to build Ask Fabro run snapshot" @@ -566,7 +576,7 @@ async fn run_streaming_turn( }; let mut output = None; let result = Box::pin(drive_agent( - &run_store, + &state, agent, run_id, session_id, @@ -596,11 +606,11 @@ async fn run_streaming_turn( match outcome.result { Ok(Ok(())) => { let _ = append_and_send_event( - &run_store, + &state, &sender, run_id, session_id, - EventBody::RunSessionTurnSucceeded(RunSessionTurnSucceededProps { + SessionEventBody::TurnSucceeded(SessionTurnSucceededProps { turn_id, output: outcome.output, }), @@ -611,7 +621,7 @@ async fn run_streaming_turn( Ok(Err(err)) => { turn_lease.entry().clear_agent().await; let body = if matches!(err, AgentError::Interrupted(_)) { - EventBody::RunSessionTurnInterrupted(RunSessionTurnInterruptedProps { + SessionEventBody::TurnInterrupted(SessionTurnInterruptedProps { turn_id, error: Some(err.to_string()), }) @@ -620,13 +630,12 @@ async fn run_streaming_turn( turn_failed_body(turn_id, err.to_string(), outcome.output, code, false) }; let _ = - append_and_send_event(&run_store, &sender, run_id, session_id, body, Utc::now()) - .await; + append_and_send_event(&state, &sender, run_id, session_id, body, Utc::now()).await; } Err(err) => { turn_lease.entry().clear_agent().await; let _ = append_and_send_event( - &run_store, + &state, &sender, run_id, session_id, @@ -634,7 +643,7 @@ async fn run_streaming_turn( turn_id, err.to_string(), outcome.output, - RunSessionTurnFailedCode::AgentError, + SessionTurnFailedCode::AgentError, false, ), Utc::now(), @@ -664,13 +673,13 @@ enum AskFabroBuildError { } impl AskFabroBuildError { - fn code(&self) -> RunSessionTurnFailedCode { + fn code(&self) -> SessionTurnFailedCode { match self { - Self::NoSandbox => RunSessionTurnFailedCode::NoSandbox, - Self::SandboxUnavailable(_) => RunSessionTurnFailedCode::SandboxUnavailable, - Self::LlmUnconfigured(_) => RunSessionTurnFailedCode::LlmUnconfigured, - Self::ModelUnavailable(_) => RunSessionTurnFailedCode::ModelUnavailable, - Self::Agent(_) => RunSessionTurnFailedCode::AgentError, + Self::NoSandbox => SessionTurnFailedCode::NoSandbox, + Self::SandboxUnavailable(_) => SessionTurnFailedCode::SandboxUnavailable, + Self::LlmUnconfigured(_) => SessionTurnFailedCode::LlmUnconfigured, + Self::ModelUnavailable(_) => SessionTurnFailedCode::ModelUnavailable, + Self::Agent(_) => SessionTurnFailedCode::AgentError, } } @@ -684,7 +693,6 @@ impl AskFabroBuildError { async fn build_agent( state: &AppState, run_id: RunId, - run_store: &RunDatabase, session: &ProjectedRunSession, ) -> Result { let catalog = state.catalog(); @@ -707,10 +715,10 @@ async fn build_agent( }; } - let projection = run_store - .state() + let projection = state + .load_run_projection(&run_id) .await - .map_err(|err| AskFabroBuildError::Agent(anyhow::Error::new(err)))?; + .map_err(|err| AskFabroBuildError::Agent(anyhow::anyhow!("{}", err.detail())))?; let sandbox_record = projection .sandbox .as_ref() @@ -756,8 +764,8 @@ async fn build_agent( // A resumed session continues its stored conversation on the model it // recorded; a record whose events outran it (a crash between the event - // log and the record write) is moved past the log's last sequence so the - // stream never reuses a number. + // append and the record write) is moved past the session's last + // sequence so the stream never reuses a number. let stored = state .stores .session_records @@ -767,7 +775,12 @@ async fn build_agent( let builder = match stored { Some(stored) => { let mut record = stored.record; - if let Ok(Some(last_seq)) = run_store.last_event_seq().await { + if let Ok(Some(last_seq)) = state + .stores + .session_events + .last_seq(session.record.id) + .await + { record.resume_after(u64::from(last_seq)); } CodingAgent::resume( @@ -1151,7 +1164,7 @@ User question: } async fn drive_agent( - run_store: &RunDatabase, + state: &AppState, agent: &mut CodingAgent, run_id: RunId, session_id: SessionId, @@ -1171,7 +1184,7 @@ async fn drive_agent( while let Ok(event) = receiver.try_recv() { record_turn_output(output, &event); Box::pin(persist_agent_event( - run_store, run_id, session_id, turn_id, event, sender, + state, run_id, session_id, turn_id, event, sender, )) .await?; } @@ -1182,7 +1195,7 @@ async fn drive_agent( Ok(event) => { record_turn_output(output, &event); Box::pin(persist_agent_event( - run_store, run_id, session_id, turn_id, event, sender, + state, run_id, session_id, turn_id, event, sender, )) .await?; } @@ -1203,10 +1216,10 @@ fn turn_failed_body( turn_id: TurnId, error: String, output: Option, - code: RunSessionTurnFailedCode, + code: SessionTurnFailedCode, retryable: bool, -) -> EventBody { - EventBody::RunSessionTurnFailed(RunSessionTurnFailedProps { +) -> SessionEventBody { + SessionEventBody::TurnFailed(SessionTurnFailedProps { turn_id, error, output, @@ -1215,19 +1228,19 @@ fn turn_failed_body( }) } -fn agent_failure_code(err: &AgentError) -> RunSessionTurnFailedCode { +fn agent_failure_code(err: &AgentError) -> SessionTurnFailedCode { match err { AgentError::ToolExecution(message) if message.contains("denied") || message.contains("not allowed") => { - RunSessionTurnFailedCode::ToolDenied + SessionTurnFailedCode::ToolDenied } - _ => RunSessionTurnFailedCode::AgentError, + _ => SessionTurnFailedCode::AgentError, } } async fn persist_agent_event( - run_store: &RunDatabase, + state: &AppState, run_id: RunId, session_id: SessionId, turn_id: TurnId, @@ -1238,25 +1251,25 @@ async fn persist_agent_event( let Some(body) = agent_event_payload(turn_id, event.event) else { return Ok(()); }; - append_and_send_event(run_store, sender, run_id, session_id, body, ts) + append_and_send_event(state, sender, run_id, session_id, body, ts) .await .map_err(Into::into) } -fn agent_event_payload(event_turn_id: TurnId, event: CodingEvent) -> Option { +fn agent_event_payload(event_turn_id: TurnId, event: CodingEvent) -> Option { match event { CodingEvent::AssistantMessage { text, model, usage, .. - } => Some(EventBody::RunSessionAssistantMessage( - RunSessionAssistantMessageProps { + } => Some(SessionEventBody::AssistantMessage( + SessionAssistantMessageProps { turn_id: event_turn_id, text, model: Some(model), usage: serde_json::to_value(usage).unwrap_or(Value::Null), }, )), - CodingEvent::TextDelta { delta } => Some(EventBody::RunSessionAssistantDelta( - RunSessionAssistantDeltaProps { + CodingEvent::TextDelta { delta } => Some(SessionEventBody::AssistantDelta( + SessionAssistantDeltaProps { turn_id: event_turn_id, delta, }, @@ -1265,8 +1278,8 @@ fn agent_event_payload(event_turn_id: TurnId, event: CodingEvent) -> Option Some(EventBody::RunSessionToolCallStarted( - RunSessionToolCallStartedProps { + } => Some(SessionEventBody::ToolCallStarted( + SessionToolCallStartedProps { turn_id: event_turn_id, tool_name, tool_call_id, @@ -1282,8 +1295,8 @@ fn agent_event_payload(event_turn_id: TurnId, event: CodingEvent) -> Option Some(EventBody::RunSessionToolCallCompleted( - RunSessionToolCallCompletedProps { + } => Some(SessionEventBody::ToolCallCompleted( + SessionToolCallCompletedProps { turn_id: event_turn_id, tool_name, tool_call_id, @@ -1303,65 +1316,42 @@ fn agent_event_payload(event_turn_id: TurnId, event: CodingEvent) -> Option, ) -> fabro_store::Result<()> { - let event = append_run_session_event(run_store, run_id, session_id, body, ts).await?; + let event = append_run_session_event(state, run_id, session_id, body, ts).await?; send_sse_event(sender, &event).await; Ok(()) } async fn append_run_session_event( - run_store: &RunDatabase, + state: &AppState, run_id: RunId, session_id: SessionId, - body: EventBody, + body: SessionEventBody, ts: DateTime, -) -> fabro_store::Result { - let event = RunEvent { - id: format!("evt_{}", ulid::Ulid::new()), - ts, - run_id, - node_id: None, - node_label: None, - stage_id: None, - parallel_group_id: None, - parallel_branch_id: None, - session_id: Some(session_id.to_string()), - parent_session_id: None, - tool_call_id: None, - actor: None, - body, - }; - let payload = EventPayload::new(event.to_value()?, &run_id)?; - run_store.append_event_envelope(&payload).await -} - -async fn send_sse_event(sender: &SessionSseSender, event: &EventEnvelope) -> bool { - let Ok(data) = serde_json::to_string(event) else { - return true; - }; - sender - .send(Ok(Event::default() - .id(event.seq.to_string()) - .event(event.event.event_name()) - .data(data))) +) -> fabro_store::Result { + state + .stores + .session_events + .append(run_id, session_id, body, ts) .await - .is_ok() } -fn session_sse_event(event: &EventEnvelope) -> Option { - let data = serde_json::to_string(event).ok()?; - Some( - Event::default() - .id(event.seq.to_string()) - .event(event.event.event_name()) - .data(data), - ) +async fn send_sse_event(sender: &SessionSseSender, event: &SessionEvent) -> bool { + sender.send(Ok(session_sse_event(event))).await.is_ok() +} + +fn session_sse_event(event: &SessionEvent) -> Event { + let data = serde_json::to_string(event).unwrap_or_else(|_| "null".to_string()); + Event::default() + .id(event.seq.to_string()) + .event(event.event_name()) + .data(data) } async fn send_attach_sse_event( @@ -1377,100 +1367,39 @@ async fn send_attach_sse_event( } } -fn event_matches_session(event: &EventEnvelope, session_id: &str) -> bool { - event - .event - .session_id - .as_deref() - .is_some_and(|id| id == session_id) - && event.event.body.is_run_session_event() -} - +/// The session as its events describe it, with the run that owns it. async fn load_session( state: &AppState, session_id: SessionId, -) -> Result<(RunId, RunDatabase, ProjectedRunSession), Response> { - let run_id = match state.store_ref().find_session_owner(&session_id).await { - Ok(Some(run_id)) => run_id, - Ok(None) => return Err(ApiError::not_found("Session not found.").into_response()), - Err(err) => return Err(store_error(&err).into_response()), - }; - let run_store = open_run(state, run_id).await?; - let events = match run_store.list_events().await { - Ok(events) => events, - Err(err) => return Err(store_error(&err).into_response()), - }; - match project_run_session(run_id, session_id, &events) { - Some(session) => Ok((run_id, run_store, session)), - None => Err(ApiError::not_found("Session not found.").into_response()), - } -} - -async fn load_session_read( - state: &AppState, - session_id: SessionId, ) -> Result<(RunId, ProjectedRunSession), Response> { - let run_id = match state.store_ref().find_session_owner(&session_id).await { - Ok(Some(run_id)) => run_id, - Ok(None) => return Err(ApiError::not_found("Session not found.").into_response()), - Err(err) => return Err(store_error(&err).into_response()), - }; - let run_store = open_run_reader(state, run_id).await?; - let events = match run_store.list_events().await { - Ok(events) => events, - Err(err) => return Err(store_error(&err).into_response()), - }; - match project_run_session(run_id, session_id, &events) { + let run_id = load_session_owner(state, session_id).await?; + let events = state + .stores + .session_events + .list_from(session_id, 1, usize::MAX) + .await + .map_err(|err| store_error(&err).into_response())?; + match project_run_session(session_id, &events) { Some(session) => Ok((run_id, session)), None => Err(ApiError::not_found("Session not found.").into_response()), } } -async fn load_session_run_reader( - state: &AppState, - session_id: SessionId, -) -> Result<(RunId, RunDatabase), Response> { - let run_id = match state.store_ref().find_session_owner(&session_id).await { - Ok(Some(run_id)) => run_id, - Ok(None) => return Err(ApiError::not_found("Session not found.").into_response()), - Err(err) => return Err(store_error(&err).into_response()), - }; - let run_store = open_run_reader(state, run_id).await?; - let events = match run_store - .list_events_for_session_from_with_limit(session_id, 1, 0) - .await - { - Ok(events) => events, - Err(err) => return Err(store_error(&err).into_response()), - }; - if events.is_empty() { - return Err(ApiError::not_found("Session not found.").into_response()); +/// The run that owns `session_id`, from the session's creation event. +async fn load_session_owner(state: &AppState, session_id: SessionId) -> Result { + match state.stores.session_events.owner(session_id).await { + Ok(Some(run_id)) => Ok(run_id), + Ok(None) => Err(ApiError::not_found("Session not found.").into_response()), + Err(err) => Err(store_error(&err).into_response()), } - Ok((run_id, run_store)) } -async fn open_run(state: &AppState, run_id: RunId) -> Result { - state.store_ref().open_run(&run_id).await.map_err(|err| { - if matches!(err, fabro_store::Error::RunNotFound(_)) { - ApiError::not_found("Run not found.").into_response() - } else { - store_error(&err).into_response() - } - }) -} - -async fn open_run_reader(state: &AppState, run_id: RunId) -> Result { - state - .store_ref() - .open_run_reader(&run_id) - .await - .map_err(|err| { - if matches!(err, fabro_store::Error::RunNotFound(_)) { - ApiError::not_found("Run not found.").into_response() - } else { - store_error(&err).into_response() - } - }) +async fn ensure_run(state: &AppState, run_id: RunId) -> Result<(), Response> { + match state.stores.run_summaries.contains(&run_id).await { + Ok(true) => Ok(()), + Ok(false) => Err(ApiError::not_found("Run not found.").into_response()), + Err(err) => Err(store_error(&err).into_response()), + } } fn store_error(err: &fabro_store::Error) -> ApiError { @@ -1728,7 +1657,7 @@ enabled = true }); match body { - Some(EventBody::RunSessionAssistantDelta(props)) => { + Some(SessionEventBody::AssistantDelta(props)) => { assert_eq!(props.turn_id, turn_id); assert_eq!(props.delta, "Hello"); } @@ -2295,14 +2224,12 @@ mod resume_tests { .clear_agent() .await; let log_head_before_resume = state - .store_ref() - .open_run_reader(&run_id) + .stores + .session_events + .last_seq(session_id) .await .unwrap() - .last_event_seq() - .await - .unwrap() - .expect("the run has events"); + .expect("the session has events"); turn(&app, session_id, "Second question").await; diff --git a/lib/apps/fabro-server/tests/it/api/sessions.rs b/lib/apps/fabro-server/tests/it/api/sessions.rs index 936b4545f..0605d8989 100644 --- a/lib/apps/fabro-server/tests/it/api/sessions.rs +++ b/lib/apps/fabro-server/tests/it/api/sessions.rs @@ -79,7 +79,7 @@ async fn create_session_response( } #[tokio::test] -async fn run_bound_session_is_created_as_run_event_and_resolves_by_flat_id() { +async fn run_bound_session_is_created_with_a_session_event_and_resolves_by_flat_id() { let app = fabro_server::test_support::build_test_router(test_app_state()); let (run_id, _run_id_workspace) = create_run(&app).await; @@ -138,65 +138,6 @@ async fn run_bound_session_is_created_as_run_event_and_resolves_by_flat_id() { ); } -#[tokio::test] -async fn generic_session_creation_requires_the_dedicated_operation_without_advancing_history() { - let app = fabro_server::test_support::build_test_router(test_app_state()); - let (run_id, _run_id_workspace) = create_run(&app).await; - let before_request = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}/events"))) - .body(Body::empty()) - .expect("run-events request should build"); - let before = response_json( - app.clone().oneshot(before_request).await.unwrap(), - StatusCode::OK, - format!("GET /api/v1/runs/{run_id}/events before rejected append"), - ) - .await; - - let session_id = fabro_types::SessionId::new(); - let request = Request::builder() - .method("POST") - .uri(api(&format!("/runs/{run_id}/events"))) - .header("content-type", "application/json") - .body(Body::from( - serde_json::to_string(&serde_json::json!({ - "id": ulid::Ulid::new().to_string(), - "ts": "2026-08-31T12:00:00Z", - "run_id": run_id, - "event": "run.session.created", - "session_id": session_id, - "properties": { "title": "Injected session" }, - })) - .expect("session creation event should serialize"), - )) - .expect("append-event request should build"); - let rejected = response_json( - app.clone().oneshot(request).await.unwrap(), - StatusCode::BAD_REQUEST, - format!("POST /api/v1/runs/{run_id}/events with session creation"), - ) - .await; - let detail = rejected["errors"][0]["detail"] - .as_str() - .expect("error response should include detail"); - assert!(detail.contains("dedicated operation endpoint")); - assert!(detail.contains("run.session.created")); - - let after_request = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}/events"))) - .body(Body::empty()) - .expect("run-events request should build"); - let after = response_json( - app.clone().oneshot(after_request).await.unwrap(), - StatusCode::OK, - format!("GET /api/v1/runs/{run_id}/events after rejected append"), - ) - .await; - assert_eq!(after, before); -} - #[tokio::test] async fn sessions_are_listed_only_under_their_owning_run() { let app = fabro_server::test_support::build_test_router(test_app_state()); @@ -420,7 +361,7 @@ async fn unsupported_derived_turn_read_routes_are_removed() { } #[tokio::test] -async fn session_events_are_filtered_by_session_and_paginated_by_run_sequence() { +async fn session_events_are_listed_by_session_and_paginated_by_session_sequence() { let app = fabro_server::test_support::build_test_router(test_app_state()); let (run_id, _run_id_workspace) = create_run(&app).await; let first = create_session(&app, &run_id, "First").await; diff --git a/lib/components/fabro-store/src/keys.rs b/lib/components/fabro-store/src/keys.rs index 28adb69ba..146082a88 100644 --- a/lib/components/fabro-store/src/keys.rs +++ b/lib/components/fabro-store/src/keys.rs @@ -58,6 +58,7 @@ impl AsRef<[u8]> for SlateKey { // matches numeric seq order through `MAX_EVENT_SEQ`. Seek-based event listing // (`run_events_range`) depends on this invariant, so event allocation rejects // larger sequences. +#[cfg(any(test, feature = "test-support"))] pub(crate) fn run_event_key(run_id: &RunId, seq: u32, epoch_ms: i64) -> SlateKey { SlateKey::new("runs") .with(run_id) diff --git a/lib/components/fabro-store/src/lib.rs b/lib/components/fabro-store/src/lib.rs index b0eca994b..0adf5768b 100644 --- a/lib/components/fabro-store/src/lib.rs +++ b/lib/components/fabro-store/src/lib.rs @@ -8,6 +8,7 @@ mod keys; pub mod platform_records; #[cfg(test)] mod record; +mod run_session_event_store; mod run_session_record_store; mod run_sessions; mod run_state; @@ -37,6 +38,7 @@ pub use platform_records::{ PlatformRecord, PlatformRecordHook, PlatformRecordKind, PlatformRecordStore, StagePosition, StoredPlatformRecord, }; +pub use run_session_event_store::RunSessionEventStore; pub use run_session_record_store::{RunSessionRecordStore, StoredSessionRecord}; pub use run_sessions::{ProjectedRunSession, project_run_session, project_run_sessions}; pub use run_state::{RunProjectionReducer, build_summary, projected_usage}; diff --git a/lib/components/fabro-store/src/run_session_event_store.rs b/lib/components/fabro-store/src/run_session_event_store.rs new file mode 100644 index 000000000..412613174 --- /dev/null +++ b/lib/components/fabro-store/src/run_session_event_store.rs @@ -0,0 +1,351 @@ +//! SQLite storage for the events of Ask Fabro sessions. +//! +//! A session's events are numbered per session, from 1, in the order they +//! are appended; an append allocates the next number under the write lock, +//! so two writers of one session never share a number. Every committed +//! event is also published to the store's subscribers, which is how a +//! client attached to a session sees a turn as it runs. + +use chrono::{DateTime, Utc}; +use fabro_types::{RunId, SessionEvent, SessionEventBody, SessionId}; +use sqlx::sqlite::SqliteRow; +use sqlx::{Row as _, SqlitePool}; +use tokio::sync::broadcast; + +use crate::{Error, Result, sqlite_row}; + +const RECORD_NAME: &str = "run session event"; + +/// How many committed events a slow subscriber may fall behind before it +/// is told it lagged and has to replay from the table. +const LIVE_CAPACITY: usize = 1024; + +/// Reads and writes session events in SQLite and publishes each committed +/// one. +pub struct RunSessionEventStore { + pool: SqlitePool, + live: broadcast::Sender, +} + +impl std::fmt::Debug for RunSessionEventStore { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("RunSessionEventStore") + .finish_non_exhaustive() + } +} + +impl RunSessionEventStore { + #[must_use] + pub fn new(pool: SqlitePool) -> Self { + let (live, _) = broadcast::channel(LIVE_CAPACITY); + Self { pool, live } + } + + /// Record `body` as the session's next event and publish it. + pub async fn append( + &self, + run_id: RunId, + session_id: SessionId, + body: SessionEventBody, + ts: DateTime, + ) -> Result { + let properties = serde_json::to_value(&body)?; + let properties_json = + serde_json::to_string(properties.get("properties").ok_or_else(|| { + Error::InvalidEvent("session event body has no properties".into()) + })?)?; + let mut transaction = self.pool.begin_with("BEGIN IMMEDIATE").await?; + let last_seq: Option = + sqlx::query_scalar("SELECT MAX(seq) FROM run_session_events WHERE session_id = ?") + .bind(session_id.to_string()) + .fetch_one(&mut *transaction) + .await?; + let seq = u32::try_from(last_seq.unwrap_or(0)) + .ok() + .and_then(|last| last.checked_add(1)) + .ok_or(Error::EventSequenceExhausted { max_seq: u32::MAX })?; + sqlx::query( + r" +INSERT INTO run_session_events ( + session_id, seq, run_id, turn_id, event_name, recorded_at_ms, properties_json +) VALUES (?, ?, ?, ?, ?, ?, ?) +", + ) + .bind(session_id.to_string()) + .bind(i64::from(seq)) + .bind(run_id.to_string()) + .bind(body.turn_id().map(|turn_id| turn_id.to_string())) + .bind(body.event_name()) + .bind(ts.timestamp_millis()) + .bind(properties_json) + .execute(&mut *transaction) + .await?; + transaction.commit().await?; + let event = SessionEvent { + seq, + session_id, + run_id, + ts, + body, + }; + // No subscriber is not an error: the event is committed either way. + let _ = self.live.send(event.clone()); + Ok(event) + } + + /// The session's events from `since_seq` on, at most `limit` of them, + /// in sequence order. + pub async fn list_from( + &self, + session_id: SessionId, + since_seq: u32, + limit: usize, + ) -> Result> { + let rows = sqlx::query( + r" +SELECT session_id, seq, run_id, event_name, recorded_at_ms, properties_json +FROM run_session_events +WHERE session_id = ? AND seq >= ? +ORDER BY seq ASC +LIMIT ? +", + ) + .bind(session_id.to_string()) + .bind(i64::from(since_seq)) + .bind(i64::try_from(limit).unwrap_or(i64::MAX)) + .fetch_all(&self.pool) + .await?; + rows.iter().map(event_from_row).collect() + } + + /// Every event of every session of `run_id`, session by session in + /// sequence order. + pub async fn list_for_run(&self, run_id: RunId) -> Result> { + let rows = sqlx::query( + r" +SELECT session_id, seq, run_id, event_name, recorded_at_ms, properties_json +FROM run_session_events +WHERE run_id = ? +ORDER BY session_id ASC, seq ASC +", + ) + .bind(run_id.to_string()) + .fetch_all(&self.pool) + .await?; + rows.iter().map(event_from_row).collect() + } + + /// The sequence number of the session's latest event, if it has any. + pub async fn last_seq(&self, session_id: SessionId) -> Result> { + let last_seq: Option = + sqlx::query_scalar("SELECT MAX(seq) FROM run_session_events WHERE session_id = ?") + .bind(session_id.to_string()) + .fetch_one(&self.pool) + .await?; + last_seq + .map(|seq| { + u32::try_from(seq).map_err(|_| { + Error::InvalidEvent(format!("stored {RECORD_NAME} sequence {seq}")) + }) + }) + .transpose() + } + + /// The run that owns `session_id`, from the session's creation event. + pub async fn owner(&self, session_id: SessionId) -> Result> { + let run_id: Option = sqlx::query_scalar( + "SELECT run_id FROM run_session_events WHERE session_id = ? AND seq = 1", + ) + .bind(session_id.to_string()) + .fetch_optional(&self.pool) + .await?; + run_id + .map(|run_id| { + run_id.parse::().map_err(|error| { + Error::InvalidEvent(format!("stored {RECORD_NAME} run id: {error}")) + }) + }) + .transpose() + } + + /// Forget every session event of `run_id`. + pub async fn delete_for_run(&self, run_id: RunId) -> Result { + let result = sqlx::query("DELETE FROM run_session_events WHERE run_id = ?") + .bind(run_id.to_string()) + .execute(&self.pool) + .await?; + Ok(result.rows_affected()) + } + + /// Every event committed from now on, across all sessions. A receiver + /// that falls more than the channel's capacity behind gets + /// `RecvError::Lagged` and replays from `list_from`. + #[must_use] + pub fn subscribe(&self) -> broadcast::Receiver { + self.live.subscribe() + } +} + +fn event_from_row(row: &SqliteRow) -> Result { + let session_id: String = row.try_get("session_id")?; + let session_id = session_id.parse::().map_err(|error| { + Error::InvalidEvent(format!("stored {RECORD_NAME} session id: {error}")) + })?; + let seq: i64 = row.try_get("seq")?; + let seq = u32::try_from(seq) + .map_err(|_| Error::InvalidEvent(format!("stored {RECORD_NAME} sequence {seq}")))?; + let run_id: String = row.try_get("run_id")?; + let run_id = run_id + .parse::() + .map_err(|error| Error::InvalidEvent(format!("stored {RECORD_NAME} run id: {error}")))?; + let ts = sqlite_row::timestamp_from_row(row, RECORD_NAME, "recorded_at_ms")?; + let event_name: String = row.try_get("event_name")?; + let properties_json: String = row.try_get("properties_json")?; + let properties: serde_json::Value = serde_json::from_str(&properties_json)?; + let body: SessionEventBody = serde_json::from_value(serde_json::json!({ + "event": event_name, + "properties": properties, + }))?; + Ok(SessionEvent { + seq, + session_id, + run_id, + ts, + body, + }) +} + +#[cfg(test)] +mod tests { + use chrono::TimeZone; + use fabro_types::session_event::{ + SessionAssistantDeltaProps, SessionCreatedProps, SessionTurnStartedProps, + }; + use fabro_types::{TurnId, fixtures}; + + use super::*; + use crate::test_support; + + fn store() -> RunSessionEventStore { + RunSessionEventStore::new(test_support::in_memory_pool_with(&[ + fabro_db::RUN_SESSION_EVENTS_MIGRATION_SQL, + ])) + } + + fn created() -> SessionEventBody { + SessionEventBody::Created(SessionCreatedProps { + title: Some("Ask Fabro".to_string()), + model: Some("test-model".to_string()), + provider: None, + }) + } + + fn turn_started(turn_id: TurnId) -> SessionEventBody { + SessionEventBody::TurnStarted(SessionTurnStartedProps { + turn_id, + input: "What happened?".to_string(), + }) + } + + #[tokio::test] + async fn appends_number_a_session_from_one_and_read_back_in_order() { + let store = store(); + let session_id = SessionId::new(); + let turn_id = TurnId::new(); + let ts = Utc.with_ymd_and_hms(2026, 9, 18, 12, 0, 0).unwrap(); + + let first = store + .append(fixtures::RUN_1, session_id, created(), ts) + .await + .unwrap(); + let second = store + .append(fixtures::RUN_1, session_id, turn_started(turn_id), ts) + .await + .unwrap(); + assert_eq!(first.seq, 1); + assert_eq!(second.seq, 2); + assert_eq!(second.body.turn_id(), Some(turn_id)); + + let events = store.list_from(session_id, 1, 10).await.unwrap(); + assert_eq!(events, vec![first.clone(), second.clone()]); + assert_eq!(store.list_from(session_id, 2, 10).await.unwrap(), vec![ + second + ]); + assert_eq!(store.list_from(session_id, 1, 1).await.unwrap(), vec![ + first + ]); + assert_eq!(store.last_seq(session_id).await.unwrap(), Some(2)); + assert_eq!(store.last_seq(SessionId::new()).await.unwrap(), None); + } + + #[tokio::test] + async fn sessions_are_numbered_apart_and_owned_by_their_run() { + let store = store(); + let first = SessionId::new(); + let second = SessionId::new(); + let ts = Utc::now(); + store + .append(fixtures::RUN_1, first, created(), ts) + .await + .unwrap(); + store + .append(fixtures::RUN_1, first, turn_started(TurnId::new()), ts) + .await + .unwrap(); + let other = store + .append(fixtures::RUN_2, second, created(), ts) + .await + .unwrap(); + assert_eq!(other.seq, 1); + + assert_eq!(store.owner(first).await.unwrap(), Some(fixtures::RUN_1)); + assert_eq!(store.owner(second).await.unwrap(), Some(fixtures::RUN_2)); + assert_eq!(store.owner(SessionId::new()).await.unwrap(), None); + + let run_events = store.list_for_run(fixtures::RUN_1).await.unwrap(); + assert_eq!(run_events.len(), 2); + assert!( + run_events + .iter() + .all(|event| event.session_id == first && event.run_id == fixtures::RUN_1) + ); + + assert_eq!(store.delete_for_run(fixtures::RUN_1).await.unwrap(), 2); + assert!( + store + .list_for_run(fixtures::RUN_1) + .await + .unwrap() + .is_empty() + ); + assert_eq!(store.owner(first).await.unwrap(), None); + assert_eq!(store.owner(second).await.unwrap(), Some(fixtures::RUN_2)); + } + + #[tokio::test] + async fn a_subscriber_sees_each_committed_event() { + let store = store(); + let session_id = SessionId::new(); + let mut live = store.subscribe(); + let ts = Utc::now(); + let appended = store + .append(fixtures::RUN_1, session_id, created(), ts) + .await + .unwrap(); + let delta = store + .append( + fixtures::RUN_1, + session_id, + SessionEventBody::AssistantDelta(SessionAssistantDeltaProps { + turn_id: TurnId::new(), + delta: "Hel".to_string(), + }), + ts, + ) + .await + .unwrap(); + + assert_eq!(live.recv().await.unwrap(), appended); + assert_eq!(live.recv().await.unwrap(), delta); + } +} diff --git a/lib/components/fabro-store/src/run_sessions.rs b/lib/components/fabro-store/src/run_sessions.rs index 1a7479f69..6c5389ef8 100644 --- a/lib/components/fabro-store/src/run_sessions.rs +++ b/lib/components/fabro-store/src/run_sessions.rs @@ -1,23 +1,25 @@ use std::collections::BTreeMap; use fabro_types::{ - EventBody, EventEnvelope, RunId, RunSessionMetadata, SessionId, SessionStatus, SessionSummary, + RunSessionMetadata, SessionEvent, SessionEventBody, SessionId, SessionStatus, SessionSummary, SessionTurn, }; -/// Ask Fabro session metadata at the event-log position it was read at. +/// Ask Fabro session metadata at the session's event position it was read +/// at. /// -/// The transcript is not projected from run events: pebble's session record -/// holds the durable history, and the `run.session.*` events stream it live. +/// The transcript is not projected from the events: pebble's session record +/// holds the durable history, and the events stream it live. #[derive(Debug, Clone, PartialEq)] pub struct ProjectedRunSession { pub record: RunSessionMetadata, pub last_seq: u32, } -pub fn project_run_sessions(run_id: RunId, events: &[EventEnvelope]) -> Vec { +/// Every session the events describe, in session id order. +pub fn project_run_sessions(events: &[SessionEvent]) -> Vec { let mut projection = RunSessionProjection::default(); - projection.apply(run_id, events); + projection.apply(events); projection .sessions .values() @@ -25,13 +27,14 @@ pub fn project_run_sessions(run_id: RunId, events: &[EventEnvelope]) -> Vec Option { let mut projection = RunSessionProjection::default(); - projection.apply(run_id, events); + projection.apply(events.iter().filter(|event| event.session_id == session_id)); projection.sessions.remove(&session_id) } @@ -41,51 +44,48 @@ struct RunSessionProjection { } impl RunSessionProjection { - fn apply(&mut self, run_id: RunId, events: &[EventEnvelope]) { - for envelope in events { - let Some(session_id) = event_session_id(envelope) else { - continue; - }; - match &envelope.event.body { - EventBody::RunSessionCreated(props) => { - let mut record = RunSessionMetadata::new(session_id, run_id, envelope.event.ts); + fn apply<'a>(&mut self, events: impl IntoIterator) { + for event in events { + let session_id = event.session_id; + match &event.body { + SessionEventBody::Created(props) => { + let mut record = RunSessionMetadata::new(session_id, event.run_id, event.ts); record.title.clone_from(&props.title); record.model.clone_from(&props.model); record.provider.clone_from(&props.provider); self.sessions.insert(session_id, ProjectedRunSession { record, - last_seq: envelope.seq, + last_seq: event.seq, }); } - EventBody::RunSessionTurnStarted(props) => { + SessionEventBody::TurnStarted(props) => { if let Some(session) = self.sessions.get_mut(&session_id) { - session.last_seq = envelope.seq; + session.last_seq = event.seq; session.record.status = SessionStatus::Running; session.record.active_turn = Some(SessionTurn { id: props.turn_id, - started_at: envelope.event.ts, + started_at: event.ts, input: props.input.clone(), }); - session.record.updated_at = envelope.event.ts; + session.record.updated_at = event.ts; } } - EventBody::RunSessionUserMessage(_) - | EventBody::RunSessionAssistantMessage(_) - | EventBody::RunSessionAssistantDelta(_) - | EventBody::RunSessionToolCallStarted(_) - | EventBody::RunSessionToolCallCompleted(_) => { + SessionEventBody::UserMessage(_) + | SessionEventBody::AssistantMessage(_) + | SessionEventBody::AssistantDelta(_) + | SessionEventBody::ToolCallStarted(_) + | SessionEventBody::ToolCallCompleted(_) => { if let Some(session) = self.sessions.get_mut(&session_id) { - session.last_seq = envelope.seq; - session.record.updated_at = envelope.event.ts; + session.last_seq = event.seq; + session.record.updated_at = event.ts; } } - EventBody::RunSessionTurnFailed(_) => { - self.finish_turn(session_id, true, envelope.event.ts, envelope.seq); + SessionEventBody::TurnFailed(_) => { + self.finish_turn(session_id, true, event.ts, event.seq); } - EventBody::RunSessionTurnSucceeded(_) | EventBody::RunSessionTurnInterrupted(_) => { - self.finish_turn(session_id, false, envelope.event.ts, envelope.seq); + SessionEventBody::TurnSucceeded(_) | SessionEventBody::TurnInterrupted(_) => { + self.finish_turn(session_id, false, event.ts, event.seq); } - _ => {} } } } @@ -110,23 +110,15 @@ impl RunSessionProjection { } } -fn event_session_id(envelope: &EventEnvelope) -> Option { - envelope - .event - .session_id - .as_deref() - .and_then(|id| id.parse().ok()) -} - #[cfg(test)] mod tests { use chrono::{TimeZone, Utc}; - use fabro_types::run_event::{ - RunSessionAssistantMessageProps, RunSessionCreatedProps, RunSessionTurnFailedCode, - RunSessionTurnFailedProps, RunSessionTurnStartedProps, RunSessionTurnSucceededProps, - RunSessionUserMessageProps, + use fabro_types::session_event::{ + SessionAssistantMessageProps, SessionCreatedProps, SessionTurnFailedCode, + SessionTurnFailedProps, SessionTurnStartedProps, SessionTurnSucceededProps, + SessionUserMessageProps, }; - use fabro_types::{EventBody, EventEnvelope, RunEvent, TurnId, fixtures}; + use fabro_types::{SessionEvent, SessionEventBody, TurnId, fixtures}; use serde_json::json; use super::{project_run_session, project_run_sessions}; @@ -139,7 +131,7 @@ mod tests { event( 1, session_id, - EventBody::RunSessionCreated(RunSessionCreatedProps { + SessionEventBody::Created(SessionCreatedProps { title: Some("Ask".to_string()), model: Some("test-model".to_string()), provider: None, @@ -148,7 +140,7 @@ mod tests { event( 2, session_id, - EventBody::RunSessionTurnStarted(RunSessionTurnStartedProps { + SessionEventBody::TurnStarted(SessionTurnStartedProps { turn_id, input: "What happened?".to_string(), }), @@ -156,15 +148,15 @@ mod tests { event( 3, session_id, - EventBody::RunSessionUserMessage(RunSessionUserMessageProps { + SessionEventBody::UserMessage(SessionUserMessageProps { turn_id, text: "What happened?".to_string(), }), ), ]; - let running = project_run_session(fixtures::RUN_1, session_id, &events) - .expect("session should project from run events"); + let running = project_run_session(session_id, &events) + .expect("session should project from its events"); assert_eq!(running.record.status, fabro_types::SessionStatus::Running); assert_eq!( running.record.active_turn.as_ref().map(|turn| turn.id), @@ -176,7 +168,7 @@ mod tests { events.push(event( 4, session_id, - EventBody::RunSessionAssistantMessage(RunSessionAssistantMessageProps { + SessionEventBody::AssistantMessage(SessionAssistantMessageProps { turn_id, text: "The run finished.".to_string(), model: Some("test-model".to_string()), @@ -186,13 +178,13 @@ mod tests { events.push(event( 5, session_id, - EventBody::RunSessionTurnSucceeded(RunSessionTurnSucceededProps { + SessionEventBody::TurnSucceeded(SessionTurnSucceededProps { turn_id, output: Some("The run finished.".to_string()), }), )); - let idle = project_run_session(fixtures::RUN_1, session_id, &events).unwrap(); + let idle = project_run_session(session_id, &events).unwrap(); assert_eq!(idle.record.status, fabro_types::SessionStatus::Idle); assert!(idle.record.active_turn.is_none()); assert_eq!(idle.record.model.as_deref(), Some("test-model")); @@ -207,7 +199,7 @@ mod tests { event( 1, session_id, - EventBody::RunSessionCreated(RunSessionCreatedProps { + SessionEventBody::Created(SessionCreatedProps { title: None, model: None, provider: None, @@ -216,7 +208,7 @@ mod tests { event( 2, session_id, - EventBody::RunSessionTurnStarted(RunSessionTurnStartedProps { + SessionEventBody::TurnStarted(SessionTurnStartedProps { turn_id, input: "hi".to_string(), }), @@ -224,40 +216,70 @@ mod tests { event( 3, session_id, - EventBody::RunSessionTurnFailed(RunSessionTurnFailedProps { + SessionEventBody::TurnFailed(SessionTurnFailedProps { turn_id, error: "boom".to_string(), output: None, - code: RunSessionTurnFailedCode::AgentError, + code: SessionTurnFailedCode::AgentError, retryable: false, }), ), ]; - let summaries = project_run_sessions(fixtures::RUN_1, &events); + let summaries = project_run_sessions(&events); assert_eq!(summaries.len(), 1); assert_eq!(summaries[0].status, fabro_types::SessionStatus::Failed); assert!(summaries[0].active_turn.is_none()); } - fn event(seq: u32, session_id: fabro_types::SessionId, body: EventBody) -> EventEnvelope { - EventEnvelope { + #[test] + fn a_session_projects_only_from_its_own_events() { + let first = fabro_types::SessionId::new(); + let second = fabro_types::SessionId::new(); + let events = vec![ + event( + 1, + first, + SessionEventBody::Created(SessionCreatedProps { + title: Some("First".to_string()), + model: None, + provider: None, + }), + ), + event( + 1, + second, + SessionEventBody::Created(SessionCreatedProps { + title: Some("Second".to_string()), + model: None, + provider: None, + }), + ), + event( + 2, + second, + SessionEventBody::TurnStarted(SessionTurnStartedProps { + turn_id: TurnId::new(), + input: "hi".to_string(), + }), + ), + ]; + + let projected = project_run_session(first, &events).unwrap(); + assert_eq!(projected.record.title.as_deref(), Some("First")); + assert_eq!(projected.record.status, fabro_types::SessionStatus::Idle); + assert_eq!(projected.last_seq, 1); + assert_eq!(project_run_sessions(&events).len(), 2); + assert!(project_run_session(fabro_types::SessionId::new(), &events).is_none()); + } + + fn event(seq: u32, session_id: fabro_types::SessionId, body: SessionEventBody) -> SessionEvent { + SessionEvent { seq, - event: RunEvent { - id: format!("evt-{seq}"), - ts: Utc.with_ymd_and_hms(2026, 5, 20, 12, 0, seq).unwrap(), - run_id: fixtures::RUN_1, - node_id: None, - node_label: None, - stage_id: None, - parallel_group_id: None, - parallel_branch_id: None, - session_id: Some(session_id.to_string()), - parent_session_id: None, - tool_call_id: None, - actor: None, - body, - }, + session_id, + run_id: fixtures::RUN_1, + ts: Utc.with_ymd_and_hms(2026, 5, 20, 12, 0, seq).unwrap(), + body, } } } diff --git a/lib/components/fabro-store/src/run_summary_store.rs b/lib/components/fabro-store/src/run_summary_store.rs index fe6ec7c04..193644871 100644 --- a/lib/components/fabro-store/src/run_summary_store.rs +++ b/lib/components/fabro-store/src/run_summary_store.rs @@ -363,7 +363,8 @@ impl RunSummaryStore { Ok(self.pool.begin().await?) } - pub(crate) async fn contains(&self, run_id: &RunId) -> Result { + /// Whether a run with `run_id` is stored. + pub async fn contains(&self, run_id: &RunId) -> Result { Ok( sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM runs WHERE id = ?)") .bind(run_id.to_string()) @@ -520,32 +521,6 @@ impl RunSummaryStore { .await } - pub(crate) async fn find_session_owner(&self, session_id: &SessionId) -> Result> { - let mut query = QueryBuilder::::new(SELECT_EVENT_COLUMNS); - query - .push(" WHERE session_id = ") - .push_bind(session_id.to_string()) - .push(" AND event_name = 'run.session.created'"); - let row = query.build().fetch_optional(&self.pool).await?; - let Some(row) = row else { - return Ok(None); - }; - - let stored_run_id: String = row.try_get("run_id")?; - let run_id = stored_run_id - .parse::() - .map_err(|_| Error::RunEventMismatch { - run_id: stored_run_id.clone(), - seq: 0, - field: "run_id", - })?; - // The WHERE clause pins the row's session_id and event_name columns - // to the requested values, and decoding verifies the envelope against - // every stored column, so a successful decode proves ownership. - decode_event_row(&row, &run_id, &stored_run_id)?; - Ok(Some(run_id)) - } - pub(crate) async fn delete_canonical(&self, run_id: &RunId) -> Result<()> { // Reserve the write lock up front: a deferred transaction upgraded // while another writer is active can fail at once, bypassing the @@ -1565,17 +1540,6 @@ mod tests { ) } - fn session_created_payload(run_id: &RunId, session_id: &SessionId) -> EventPayload { - sql_event_payload( - run_id, - "run.session.created", - None, - None, - Some(session_id), - serde_json::json!({ "title": "Owned session" }), - ) - } - async fn seed_sql_event( store: &RunSummaryStore, run_id: &RunId, @@ -1929,122 +1893,6 @@ mod tests { )); } - #[tokio::test] - async fn session_owner_lookup_resolves_only_typed_creation_events() { - let (_directory, store) = store().await; - let created_at = dt("2026-08-27T12:00:00Z"); - let id = run_id(created_at.timestamp_millis().cast_unsigned(), 31); - store - .upsert_projection(&entry(projection(id, "owner", created_at), 3)) - .await - .unwrap(); - - let owner_session = SessionId::new(); - seed_sql_event( - &store, - &id, - 2, - &session_created_payload(&id, &owner_session), - ) - .await; - assert_eq!( - store.find_session_owner(&owner_session).await.unwrap(), - Some(id) - ); - - let non_owner_session = SessionId::new(); - let non_creation = sql_event_payload( - &id, - "run.session.future", - None, - None, - Some(&non_owner_session), - serde_json::json!({ "kind": "future" }), - ); - seed_sql_event(&store, &id, 3, &non_creation).await; - assert_eq!( - store.find_session_owner(&non_owner_session).await.unwrap(), - None - ); - assert_eq!( - store.find_session_owner(&SessionId::new()).await.unwrap(), - None - ); - } - - #[tokio::test] - async fn session_owner_lookup_rejects_corrupt_selected_events() { - let (_directory, store) = store().await; - let created_at = dt("2026-08-27T12:00:00Z"); - let id = run_id(created_at.timestamp_millis().cast_unsigned(), 32); - let other_id = run_id(created_at.timestamp_millis().cast_unsigned() + 1, 33); - store - .upsert_projection(&entry(projection(id, "owner", created_at), 2)) - .await - .unwrap(); - let session_id = SessionId::new(); - let payload = session_created_payload(&id, &session_id); - seed_sql_event(&store, &id, 2, &payload).await; - - let mut wrong_event_name = payload.as_value().clone(); - wrong_event_name["event"] = "run.session.future".into(); - sqlx::query("UPDATE run_events SET event_json = ? WHERE run_id = ? AND seq = 2") - .bind(serde_json::to_string(&wrong_event_name).unwrap()) - .bind(id.to_string()) - .execute(&store.pool) - .await - .unwrap(); - assert!(matches!( - store.find_session_owner(&session_id).await.unwrap_err(), - Error::RunEventMismatch { - field: "event_name", - .. - } - )); - - seed_sql_event_restore(&store, &id, 2, &payload).await; - let mut wrong_run = payload.as_value().clone(); - wrong_run["run_id"] = other_id.to_string().into(); - sqlx::query("UPDATE run_events SET event_json = ? WHERE run_id = ? AND seq = 2") - .bind(serde_json::to_string(&wrong_run).unwrap()) - .bind(id.to_string()) - .execute(&store.pool) - .await - .unwrap(); - assert!(matches!( - store.find_session_owner(&session_id).await.unwrap_err(), - Error::RunEventMismatch { - field: "run_id", - .. - } - )); - - seed_sql_event_restore(&store, &id, 2, &payload).await; - let mut wrong_session = payload.as_value().clone(); - wrong_session["session_id"] = SessionId::new().to_string().into(); - sqlx::query("UPDATE run_events SET event_json = ? WHERE run_id = ? AND seq = 2") - .bind(serde_json::to_string(&wrong_session).unwrap()) - .bind(id.to_string()) - .execute(&store.pool) - .await - .unwrap(); - assert!(matches!( - store.find_session_owner(&session_id).await.unwrap_err(), - Error::RunEventMismatch { - field: "session_id", - .. - } - )); - - seed_sql_event_restore(&store, &id, 2, &payload).await; - sqlx::query("UPDATE run_events SET event_json = '{}' WHERE run_id = ? AND seq = 2") - .bind(id.to_string()) - .execute(&store.pool) - .await - .unwrap(); - assert!(store.find_session_owner(&session_id).await.is_err()); - } - #[tokio::test] async fn canonical_delete_waits_for_a_concurrent_writer() { let (_directory, store) = store().await; diff --git a/lib/components/fabro-store/src/slate/mod.rs b/lib/components/fabro-store/src/slate/mod.rs index 7206c7c6d..c2ca329d4 100644 --- a/lib/components/fabro-store/src/slate/mod.rs +++ b/lib/components/fabro-store/src/slate/mod.rs @@ -6,7 +6,7 @@ use std::sync::Arc; use std::time::Duration; use chrono::{DateTime, Utc}; -use fabro_types::{RunId, SessionId}; +use fabro_types::RunId; use object_store::ObjectStore; pub use run_store::RunDatabase; use run_store::RunDatabaseInner; @@ -251,12 +251,6 @@ impl Database { self.run_summary_store.set_platform_record_hook(hook); } - /// Resolves the run that owns `session_id` from the canonical typed - /// creation event stored in SQLite. - pub async fn find_session_owner(&self, session_id: &SessionId) -> Result> { - self.run_summary_store.find_session_owner(session_id).await - } - pub async fn delete_run(&self, run_id: &RunId) -> Result<()> { let mut active_runs = self.active_runs.lock().await; let active = active_runs.get(run_id).cloned(); @@ -516,21 +510,6 @@ mod tests { .unwrap() } - fn session_created_payload(label: &str, session_id: &SessionId) -> EventPayload { - EventPayload::new( - serde_json::json!({ - "id": format!("evt-{label}-session-created"), - "ts": "2026-03-27T12:00:05Z", - "run_id": test_run_id(label).to_string(), - "event": "run.session.created", - "session_id": session_id, - "properties": { "title": "Owned session" }, - }), - &test_run_id(label), - ) - .unwrap() - } - async fn append_created(run: &RunDatabase, label: &str, created_at: DateTime) { let run_spec = sample_run_spec(label); run.append_event(&event_payload( @@ -742,96 +721,6 @@ mod tests { ); } - #[tokio::test] - async fn session_owner_claims_are_atomic_durable_and_ignore_legacy_reverse_rows() { - let directory = tempfile::tempdir().unwrap(); - let object_store: Arc = Arc::new(InMemory::new()); - let store = store_test_support::test_database_at( - Arc::clone(&object_store), - "session-owner", - Duration::from_millis(1), - None, - directory.path(), - ); - let first_id = test_run_id("run-1"); - let second_id = test_run_id("run-2"); - let first = store.create_run(&first_id).await.unwrap(); - let second = store.create_run(&second_id).await.unwrap(); - append_created(&first, "run-1", dt("2026-03-27T12:00:00Z")).await; - append_created(&second, "run-2", dt("2026-03-27T12:00:10Z")).await; - - let session_id = SessionId::new(); - assert_eq!( - first - .append_event(&session_created_payload("run-1", &session_id)) - .await - .unwrap(), - 2 - ); - assert_eq!( - store.find_session_owner(&session_id).await.unwrap(), - Some(first_id) - ); - - let legacy_key = keys::session_by_id_key(&session_id).as_ref().to_vec(); - let legacy = store.open_db().await.unwrap(); - assert!(legacy.get(&legacy_key).await.unwrap().is_none()); - legacy - .put( - &legacy_key, - serde_json::to_vec(&serde_json::json!({ "run_id": second_id })).unwrap(), - ) - .await - .unwrap(); - legacy.flush().await.unwrap(); - assert_eq!( - store.find_session_owner(&session_id).await.unwrap(), - Some(first_id), - "legacy reverse rows must not influence ownership" - ); - - assert!( - first - .append_event(&session_created_payload("run-1", &session_id)) - .await - .is_err() - ); - assert_eq!(first.last_event_seq().await.unwrap(), Some(2)); - assert!( - second - .append_event(&session_created_payload("run-2", &session_id)) - .await - .is_err() - ); - assert_eq!(second.last_event_seq().await.unwrap(), Some(1)); - assert_eq!( - store.find_session_owner(&session_id).await.unwrap(), - Some(first_id) - ); - - let reopened = store_test_support::test_database_at( - object_store, - "session-owner", - Duration::from_millis(1), - None, - directory.path(), - ); - assert_eq!( - reopened.find_session_owner(&session_id).await.unwrap(), - Some(first_id) - ); - - reopened.delete_run(&first_id).await.unwrap(); - assert_eq!( - reopened.find_session_owner(&session_id).await.unwrap(), - None - ); - assert!( - legacy.get(&legacy_key).await.unwrap().is_some(), - "legacy reverse rows remain diagnostic-only during the support window" - ); - } - #[tokio::test] async fn delete_run_keeps_global_cas_blobs() { let (_object_store, store) = make_store(); diff --git a/lib/foundation/fabro-api/build.rs b/lib/foundation/fabro-api/build.rs index 9c527af4e..aa0d94802 100644 --- a/lib/foundation/fabro-api/build.rs +++ b/lib/foundation/fabro-api/build.rs @@ -658,6 +658,7 @@ fn main() { &[], ), ("EventEnvelope", "fabro_types::EventEnvelope", &[]), + ("SessionEvent", "fabro_types::SessionEvent", &[]), ("RunStreamItem", "fabro_types::RunStreamItem", &[]), ("RunStreamItemKind", "fabro_types::RunStreamItemKind", &[]), ("PetriAdmission", "fabro_types::PetriAdmission", &[]), diff --git a/lib/foundation/fabro-api/src/lib.rs b/lib/foundation/fabro-api/src/lib.rs index ac4621284..2ea56ac32 100644 --- a/lib/foundation/fabro-api/src/lib.rs +++ b/lib/foundation/fabro-api/src/lib.rs @@ -64,13 +64,14 @@ pub mod types { RunTarget, SandboxDetails, SandboxInfo, SandboxListMeta, SandboxListResponse, SandboxProviderKind, SandboxProviderLookupError, SandboxService, SandboxServiceListResponse, SecretMetadata, SecretType, ServerSettings, SessionDetail, - SessionId, SessionStatus, SessionSummary, SessionTurn, SkillActivationSource, SkillSummary, - StageCompletion, StageContextWindow, StageContextWindowUnavailableReason, StageHandler, - StageId, StageInferenceProjection, StageModelUsage, StageOutcome, StageProjection, - StageState, StageToolBatchProjection, SystemActorKind, SystemIntegrationStatus, - SystemIntegrationsResponse, TodoListProjection, ToolCategory, ToolSource, ToolSummary, - TurnId, UpdateVariableRequest, UserPrincipal, Variable, VariableListResponse, WorkflowPath, - WorkflowSettings, WorkflowVersion, WorkflowVersionId, + SessionEvent, SessionEventBody, SessionId, SessionStatus, SessionSummary, SessionTurn, + SkillActivationSource, SkillSummary, StageCompletion, StageContextWindow, + StageContextWindowUnavailableReason, StageHandler, StageId, StageInferenceProjection, + StageModelUsage, StageOutcome, StageProjection, StageState, StageToolBatchProjection, + SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection, + ToolCategory, ToolSource, ToolSummary, TurnId, UpdateVariableRequest, UserPrincipal, + Variable, VariableListResponse, WorkflowPath, WorkflowSettings, WorkflowVersion, + WorkflowVersionId, }; pub use lithos_llm::catalog::{ModelHandle, ProviderId}; pub use lithos_llm::types::{ diff --git a/lib/foundation/fabro-api/tests/session_contract_round_trip.rs b/lib/foundation/fabro-api/tests/session_contract_round_trip.rs index 38a0fcfd6..13eb627ae 100644 --- a/lib/foundation/fabro-api/tests/session_contract_round_trip.rs +++ b/lib/foundation/fabro-api/tests/session_contract_round_trip.rs @@ -2,12 +2,14 @@ use std::any::{TypeId, type_name}; use chrono::{TimeZone, Utc}; use fabro_api::types::{ - RunSessionMetadata as ApiRunSessionMetadata, SessionDetail as ApiSessionDetail, + PaginatedSessionEventList, RunSessionMetadata as ApiRunSessionMetadata, + SessionDetail as ApiSessionDetail, SessionEvent as ApiSessionEvent, SessionSummary as ApiSessionSummary, SessionTurn as ApiSessionTurn, SubmitTurnRequest, }; +use fabro_types::session_event::SessionTurnStartedProps; use fabro_types::{ - RunSessionMetadata, SessionDetail, SessionId, SessionStatus, SessionSummary, SessionTurn, - TurnId, fixtures, + RunSessionMetadata, SessionDetail, SessionEvent, SessionEventBody, SessionId, SessionStatus, + SessionSummary, SessionTurn, TurnId, fixtures, }; use serde_json::json; @@ -17,6 +19,35 @@ fn session_contract_reuses_domain_types() { assert_same_type::(); assert_same_type::(); assert_same_type::(); + assert_same_type::(); +} + +#[test] +fn a_session_event_page_round_trips_the_flattened_event() { + let session_id = SessionId::new(); + let turn_id = TurnId::new(); + let value = json!({ + "data": [{ + "seq": 2, + "session_id": session_id.to_string(), + "run_id": fixtures::RUN_1, + "ts": "2026-05-20T12:00:01Z", + "event": "run.session.turn.started", + "properties": { "turn_id": turn_id.to_string(), "input": "What changed?" } + }], + "meta": { "has_more": false } + }); + + let page: PaginatedSessionEventList = + serde_json::from_value(value.clone()).expect("page should deserialize"); + assert_eq!( + page.data[0].body, + SessionEventBody::TurnStarted(SessionTurnStartedProps { + turn_id, + input: "What changed?".to_string(), + }) + ); + assert_eq!(serde_json::to_value(&page).unwrap(), value); } #[test] diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index bbacc5862..09a0ffcf0 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -15,7 +15,7 @@ use fabro_types::{ ArtifactUpload, BlobHash, EventEnvelope, Model, ModelTestMode, PairId, PairMessageRecord, PairMessageRequest, PairRecord, PairStartRequest, PairTranscriptResponse, Run, RunEvent, RunEventDetailResponse, RunId, RunPairStatusResponse, RunProjection, RunSessionMetadata, - RunStreamItem, SessionId, StageId, WorkflowVersion, WorkflowVersionId, + RunStreamItem, SessionEvent, SessionId, StageId, WorkflowVersion, WorkflowVersionId, }; use fabro_util::exit::{ErrorExt, ExitClass}; use futures::future::BoxFuture; @@ -76,10 +76,12 @@ pub struct RunStreamPage { type HttpByteStream = Pin> + Send>>; +/// The live stream of an Ask Fabro turn, as `POST /sessions/{id}/turns` +/// serves it: one `SessionEvent` per `data:` frame, in `seq` order. pub struct SessionEventStream { stream: HttpByteStream, pending_bytes: Vec, - buffered_events: VecDeque, + buffered_events: VecDeque, } #[derive(Default)] @@ -245,7 +247,7 @@ impl SessionEventStream { } } - pub async fn next_event(&mut self) -> Result> { + pub async fn next_event(&mut self) -> Result> { loop { if let Some(event) = self.buffered_events.pop_front() { return Ok(Some(event)); diff --git a/lib/foundation/fabro-db/migrations/2026091802_run_session_events.sql b/lib/foundation/fabro-db/migrations/2026091802_run_session_events.sql new file mode 100644 index 000000000..84cf87381 --- /dev/null +++ b/lib/foundation/fabro-db/migrations/2026091802_run_session_events.sql @@ -0,0 +1,19 @@ +-- The events of an Ask Fabro session: the turns and their messages, tool +-- calls and endings, numbered per session from 1. They stream a turn live +-- and project the session's metadata; the conversation itself is the +-- session's record in `run_session_records`. +CREATE TABLE run_session_events ( + session_id TEXT NOT NULL, + seq INTEGER NOT NULL, + run_id TEXT NOT NULL, + turn_id TEXT, + event_name TEXT NOT NULL, + recorded_at_ms INTEGER NOT NULL, + properties_json TEXT NOT NULL, + PRIMARY KEY (session_id, seq), + CHECK (seq >= 1), + CHECK (json_valid(properties_json)) +) WITHOUT ROWID; + +CREATE INDEX run_session_events_by_run +ON run_session_events(run_id, session_id, seq); diff --git a/lib/foundation/fabro-db/src/lib.rs b/lib/foundation/fabro-db/src/lib.rs index 43bfaee1b..d710ae1cd 100644 --- a/lib/foundation/fabro-db/src/lib.rs +++ b/lib/foundation/fabro-db/src/lib.rs @@ -41,6 +41,12 @@ pub const RUN_EVENT_SESSION_OWNER_MIGRATION_SQL: &str = pub const RUN_SESSION_RECORDS_MIGRATION_SQL: &str = include_str!("../migrations/2026091101_run_session_records.sql"); +/// The Ask Fabro session events migration (`run_session_events`), exposed so +/// fixtures in other crates can install the production schema without a +/// filesystem path into this crate. +pub const RUN_SESSION_EVENTS_MIGRATION_SQL: &str = + include_str!("../migrations/2026091802_run_session_events.sql"); + /// The Petri run record migration (`petri_runs`, `petri_records`), exposed /// so fixtures in other crates can install the production schema without a /// filesystem path into this crate. diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index 2f3ad48bd..55367f476 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -45,6 +45,7 @@ pub mod sandbox_provider; pub mod sandbox_services; pub mod secret; pub mod session; +pub mod session_event; pub mod settings; pub mod stage_completion; pub mod stage_handler; @@ -176,6 +177,7 @@ pub use session::{ RunSessionMetadata, SessionDetail, SessionId, SessionStatus, SessionSummary, SessionTurn, TurnId, }; +pub use session_event::{SessionEvent, SessionEventBody}; pub use stage_completion::StageCompletion; pub use stage_handler::StageHandler; pub use stage_id::{InvalidStageVisit, ParallelBranchId, StageId}; diff --git a/lib/foundation/fabro-types/src/run_event/session.rs b/lib/foundation/fabro-types/src/run_event/session.rs index 6cb2766c9..63e4d26d7 100644 --- a/lib/foundation/fabro-types/src/run_event/session.rs +++ b/lib/foundation/fabro-types/src/run_event/session.rs @@ -1,154 +1,16 @@ -use lithos_llm::catalog::ProviderId; -use serde::{Deserialize, Serialize}; -use serde_json::Value; +//! The legacy event log's names for the session event properties, which +//! live in `crate::session_event`. -use crate::TurnId; - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionCreatedProps { - #[serde(default, skip_serializing_if = "Option::is_none")] - pub title: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub model: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub provider: Option, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionTurnStartedProps { - pub turn_id: TurnId, - pub input: String, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionUserMessageProps { - pub turn_id: TurnId, - pub text: String, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionAssistantDeltaProps { - pub turn_id: TurnId, - pub delta: String, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionAssistantMessageProps { - pub turn_id: TurnId, - pub text: String, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub model: Option, - #[serde(default)] - pub usage: Value, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionToolCallStartedProps { - pub turn_id: TurnId, - pub tool_name: String, - pub tool_call_id: String, - pub arguments: Value, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionToolCallCompletedProps { - pub turn_id: TurnId, - pub tool_name: String, - pub tool_call_id: String, - pub output: Value, - pub is_error: bool, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub output_bytes_observed: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub output_bytes_retained: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub output_bytes_omitted: Option, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionTurnSucceededProps { - pub turn_id: TurnId, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub output: Option, -} - -#[derive( - Debug, - Clone, - Copy, - PartialEq, - Eq, - Hash, - Default, - Serialize, - Deserialize, - strum::Display, - strum::EnumString, - strum::IntoStaticStr, -)] -#[serde(rename_all = "snake_case")] -#[strum(serialize_all = "snake_case")] -pub enum RunSessionTurnFailedCode { - NoSandbox, - SandboxUnavailable, - LlmUnconfigured, - ModelUnavailable, - ToolDenied, - #[default] - AgentError, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionTurnFailedProps { - pub turn_id: TurnId, - pub error: String, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub output: Option, - #[serde(default)] - pub code: RunSessionTurnFailedCode, - #[serde(default)] - pub retryable: bool, -} - -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct RunSessionTurnInterruptedProps { - pub turn_id: TurnId, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub error: Option, -} - -#[cfg(test)] -mod tests { - use serde_json::json; - - use super::{RunSessionCreatedProps, RunSessionToolCallCompletedProps}; - use crate::TurnId; - - #[test] - fn session_created_deserializes_legacy_payload_without_provider() { - let props: RunSessionCreatedProps = serde_json::from_value(json!({ - "title": "Legacy session", - "model": "gpt-5.4" - })) - .unwrap(); - - assert_eq!(props.model.as_deref(), Some("gpt-5.4")); - assert_eq!(props.provider, None); - } - - #[test] - fn tool_completion_deserializes_without_output_byte_counts() { - let props: RunSessionToolCallCompletedProps = serde_json::from_value(json!({ - "turn_id": TurnId::new(), - "tool_name": "shell", - "tool_call_id": "call_1", - "output": "ok", - "is_error": false - })) - .unwrap(); - - assert!(props.output_bytes_observed.is_none()); - assert!(props.output_bytes_retained.is_none()); - assert!(props.output_bytes_omitted.is_none()); - } -} +pub use crate::session_event::{ + SessionAssistantDeltaProps as RunSessionAssistantDeltaProps, + SessionAssistantMessageProps as RunSessionAssistantMessageProps, + SessionCreatedProps as RunSessionCreatedProps, + SessionToolCallCompletedProps as RunSessionToolCallCompletedProps, + SessionToolCallStartedProps as RunSessionToolCallStartedProps, + SessionTurnFailedCode as RunSessionTurnFailedCode, + SessionTurnFailedProps as RunSessionTurnFailedProps, + SessionTurnInterruptedProps as RunSessionTurnInterruptedProps, + SessionTurnStartedProps as RunSessionTurnStartedProps, + SessionTurnSucceededProps as RunSessionTurnSucceededProps, + SessionUserMessageProps as RunSessionUserMessageProps, +}; diff --git a/lib/foundation/fabro-types/src/session_event.rs b/lib/foundation/fabro-types/src/session_event.rs new file mode 100644 index 000000000..9c4f3d0aa --- /dev/null +++ b/lib/foundation/fabro-types/src/session_event.rs @@ -0,0 +1,292 @@ +//! The events of an Ask Fabro session. +//! +//! A session's conversation is pebble's session record; its events are the +//! live view of a turn: the turn starting, the user's message, the +//! assistant's deltas and messages, the tool calls, and the turn's end. They +//! are numbered per session, from 1, and stream to the session's clients as +//! they are recorded. + +use chrono::{DateTime, Utc}; +use lithos_llm::catalog::ProviderId; +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +use crate::{RunId, SessionId, TurnId}; + +/// One recorded event of a session. +/// +/// On the wire the body is flattened: `event` names the kind and +/// `properties` holds its fields, beside `seq`, `session_id`, `run_id` and +/// `ts`. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionEvent { + /// The event's position in its session, from 1. + pub seq: u32, + pub session_id: SessionId, + pub run_id: RunId, + pub ts: DateTime, + #[serde(flatten)] + pub body: SessionEventBody, +} + +impl SessionEvent { + #[must_use] + pub fn event_name(&self) -> &'static str { + self.body.event_name() + } +} + +/// What a session event records. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(tag = "event", content = "properties")] +pub enum SessionEventBody { + #[serde(rename = "run.session.created")] + Created(SessionCreatedProps), + #[serde(rename = "run.session.turn.started")] + TurnStarted(SessionTurnStartedProps), + #[serde(rename = "run.session.user_message")] + UserMessage(SessionUserMessageProps), + #[serde(rename = "run.session.assistant_delta")] + AssistantDelta(SessionAssistantDeltaProps), + #[serde(rename = "run.session.assistant_message")] + AssistantMessage(SessionAssistantMessageProps), + #[serde(rename = "run.session.tool_call.started")] + ToolCallStarted(SessionToolCallStartedProps), + #[serde(rename = "run.session.tool_call.completed")] + ToolCallCompleted(SessionToolCallCompletedProps), + #[serde(rename = "run.session.turn.succeeded")] + TurnSucceeded(SessionTurnSucceededProps), + #[serde(rename = "run.session.turn.failed")] + TurnFailed(SessionTurnFailedProps), + #[serde(rename = "run.session.turn.interrupted")] + TurnInterrupted(SessionTurnInterruptedProps), +} + +impl SessionEventBody { + #[must_use] + pub fn event_name(&self) -> &'static str { + match self { + Self::Created(_) => "run.session.created", + Self::TurnStarted(_) => "run.session.turn.started", + Self::UserMessage(_) => "run.session.user_message", + Self::AssistantDelta(_) => "run.session.assistant_delta", + Self::AssistantMessage(_) => "run.session.assistant_message", + Self::ToolCallStarted(_) => "run.session.tool_call.started", + Self::ToolCallCompleted(_) => "run.session.tool_call.completed", + Self::TurnSucceeded(_) => "run.session.turn.succeeded", + Self::TurnFailed(_) => "run.session.turn.failed", + Self::TurnInterrupted(_) => "run.session.turn.interrupted", + } + } + + /// The turn the event belongs to; `None` for the session's creation. + #[must_use] + pub fn turn_id(&self) -> Option { + match self { + Self::Created(_) => None, + Self::TurnStarted(props) => Some(props.turn_id), + Self::UserMessage(props) => Some(props.turn_id), + Self::AssistantDelta(props) => Some(props.turn_id), + Self::AssistantMessage(props) => Some(props.turn_id), + Self::ToolCallStarted(props) => Some(props.turn_id), + Self::ToolCallCompleted(props) => Some(props.turn_id), + Self::TurnSucceeded(props) => Some(props.turn_id), + Self::TurnFailed(props) => Some(props.turn_id), + Self::TurnInterrupted(props) => Some(props.turn_id), + } + } + + /// Whether the event ends a turn. + #[must_use] + pub fn ends_turn(&self) -> bool { + matches!( + self, + Self::TurnSucceeded(_) | Self::TurnFailed(_) | Self::TurnInterrupted(_) + ) + } +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionCreatedProps { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub title: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub model: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub provider: Option, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionTurnStartedProps { + pub turn_id: TurnId, + pub input: String, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionUserMessageProps { + pub turn_id: TurnId, + pub text: String, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionAssistantDeltaProps { + pub turn_id: TurnId, + pub delta: String, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionAssistantMessageProps { + pub turn_id: TurnId, + pub text: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub model: Option, + #[serde(default)] + pub usage: Value, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionToolCallStartedProps { + pub turn_id: TurnId, + pub tool_name: String, + pub tool_call_id: String, + pub arguments: Value, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionToolCallCompletedProps { + pub turn_id: TurnId, + pub tool_name: String, + pub tool_call_id: String, + pub output: Value, + pub is_error: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_observed: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_retained: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_omitted: Option, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionTurnSucceededProps { + pub turn_id: TurnId, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output: Option, +} + +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Hash, + Default, + Serialize, + Deserialize, + strum::Display, + strum::EnumString, + strum::IntoStaticStr, +)] +#[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] +pub enum SessionTurnFailedCode { + NoSandbox, + SandboxUnavailable, + LlmUnconfigured, + ModelUnavailable, + ToolDenied, + #[default] + AgentError, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionTurnFailedProps { + pub turn_id: TurnId, + pub error: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output: Option, + #[serde(default)] + pub code: SessionTurnFailedCode, + #[serde(default)] + pub retryable: bool, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionTurnInterruptedProps { + pub turn_id: TurnId, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +#[cfg(test)] +mod tests { + use chrono::{TimeZone, Utc}; + use serde_json::json; + + use super::{ + SessionCreatedProps, SessionEvent, SessionEventBody, SessionToolCallCompletedProps, + SessionTurnStartedProps, + }; + use crate::{SessionId, TurnId, fixtures}; + + #[test] + fn a_session_event_flattens_its_body_on_the_wire() { + let session_id = SessionId::new(); + let turn_id = TurnId::new(); + let event = SessionEvent { + seq: 2, + session_id, + run_id: fixtures::RUN_1, + ts: Utc.with_ymd_and_hms(2026, 5, 20, 12, 0, 0).unwrap(), + body: SessionEventBody::TurnStarted(SessionTurnStartedProps { + turn_id, + input: "What happened?".to_string(), + }), + }; + + let value = serde_json::to_value(&event).unwrap(); + assert_eq!( + value, + json!({ + "seq": 2, + "session_id": session_id.to_string(), + "run_id": fixtures::RUN_1, + "ts": "2026-05-20T12:00:00Z", + "event": "run.session.turn.started", + "properties": { "turn_id": turn_id.to_string(), "input": "What happened?" } + }) + ); + let round_trip: SessionEvent = serde_json::from_value(value).unwrap(); + assert_eq!(round_trip, event); + assert_eq!(round_trip.body.turn_id(), Some(turn_id)); + assert!(!round_trip.body.ends_turn()); + } + + #[test] + fn a_created_event_has_no_turn() { + let body = SessionEventBody::Created(SessionCreatedProps { + title: Some("Ask".to_string()), + model: None, + provider: None, + }); + assert_eq!(body.event_name(), "run.session.created"); + assert_eq!(body.turn_id(), None); + } + + #[test] + fn tool_completion_deserializes_without_output_byte_counts() { + let props: SessionToolCallCompletedProps = serde_json::from_value(json!({ + "turn_id": TurnId::new(), + "tool_name": "shell", + "tool_call_id": "call_1", + "output": "ok", + "is_error": false + })) + .unwrap(); + + assert!(props.output_bytes_observed.is_none()); + assert!(props.output_bytes_retained.is_none()); + assert!(props.output_bytes_omitted.is_none()); + } +} diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 9e6dd0fff..0d57f88d0 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -272,6 +272,7 @@ models/paginated-run-list.ts models/paginated-run-stage-list.ts models/paginated-run-stream-list.ts models/paginated-saved-query-list.ts +models/paginated-session-event-list.ts models/paginated-session-list.ts models/paginated-workflow-list-response.ts models/pagination-meta.ts @@ -499,6 +500,8 @@ models/server-slate-db-settings.ts models/server-storage-settings.ts models/server-web-settings.ts models/session-detail.ts +models/session-event-name.ts +models/session-event.ts models/session-status.ts models/session-summary.ts models/session-turn.ts diff --git a/lib/packages/fabro-api-client/src/api/sessions-api.ts b/lib/packages/fabro-api-client/src/api/sessions-api.ts index c6474faf8..677dac92d 100644 --- a/lib/packages/fabro-api-client/src/api/sessions-api.ts +++ b/lib/packages/fabro-api-client/src/api/sessions-api.ts @@ -26,9 +26,7 @@ import type { CreateRunSessionRequest } from '../models'; // @ts-ignore import type { ErrorResponse } from '../models'; // @ts-ignore -import type { EventEnvelope } from '../models'; -// @ts-ignore -import type { PaginatedEventList } from '../models'; +import type { PaginatedSessionEventList } from '../models'; // @ts-ignore import type { PaginatedSessionList } from '../models'; // @ts-ignore @@ -36,6 +34,8 @@ import type { RunSessionMetadata } from '../models'; // @ts-ignore import type { SessionDetail } from '../models'; // @ts-ignore +import type { SessionEvent } from '../models'; +// @ts-ignore import type { SubmitTurnRequest } from '../models'; /** * SessionsApi - axios parameter creator @@ -43,7 +43,7 @@ import type { SubmitTurnRequest } from '../models'; export const SessionsApiAxiosParamCreator = function (configuration?: Configuration) { return { /** - * Replays and streams this session\'s durable `run.session.*` events from the owning run event log. The stream remains open until the client disconnects or the server shuts down. + * Replays this session\'s events from `since_seq` (the next unseen event when omitted) as `SessionEvent` frames, then streams new ones as they are recorded. The stream remains open until the client disconnects or the server shuts down. * @summary Attach to session events * @param {string} id * @param {number} [sinceSeq] @@ -272,7 +272,7 @@ export const SessionsApiAxiosParamCreator = function (configuration?: Configurat }; }, /** - * Returns run event envelopes filtered to this session\'s durable `run.session.*` events. `since_seq` uses the owning run event sequence. + * Returns this session\'s events in order. Events are numbered per session from 1; `since_seq` is the first sequence number to include. * @summary List session events * @param {string} id * @param {number} [sinceSeq] @@ -322,7 +322,7 @@ export const SessionsApiAxiosParamCreator = function (configuration?: Configurat }; }, /** - * Starts a streamed turn immediately. Background turns are not supported in this API version. + * Starts a streamed turn immediately. The stream carries the turn\'s `SessionEvent` frames, from `run.session.turn.started` to the event that ends the turn. Background turns are not supported in this API version. * @summary Submit a session turn * @param {string} id * @param {SubmitTurnRequest} submitTurnRequest @@ -376,7 +376,7 @@ export const SessionsApiFp = function(configuration?: Configuration) { const localVarAxiosParamCreator = SessionsApiAxiosParamCreator(configuration) return { /** - * Replays and streams this session\'s durable `run.session.*` events from the owning run event log. The stream remains open until the client disconnects or the server shuts down. + * Replays this session\'s events from `since_seq` (the next unseen event when omitted) as `SessionEvent` frames, then streams new ones as they are recorded. The stream remains open until the client disconnects or the server shuts down. * @summary Attach to session events * @param {string} id * @param {number} [sinceSeq] @@ -424,7 +424,7 @@ export const SessionsApiFp = function(configuration?: Configuration) { * @param {*} [options] Override http request option. * @throws {RequiredError} */ - async interruptSessionTurn(id: string, turnId: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + async interruptSessionTurn(id: string, turnId: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { const localVarAxiosArgs = await localVarAxiosParamCreator.interruptSessionTurn(id, turnId, options); const localVarOperationServerIndex = configuration?.serverIndex ?? 0; const localVarOperationServerBasePath = operationServerMap['SessionsApi.interruptSessionTurn']?.[localVarOperationServerIndex]?.url; @@ -447,7 +447,7 @@ export const SessionsApiFp = function(configuration?: Configuration) { return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, /** - * Returns run event envelopes filtered to this session\'s durable `run.session.*` events. `since_seq` uses the owning run event sequence. + * Returns this session\'s events in order. Events are numbered per session from 1; `since_seq` is the first sequence number to include. * @summary List session events * @param {string} id * @param {number} [sinceSeq] @@ -455,14 +455,14 @@ export const SessionsApiFp = function(configuration?: Configuration) { * @param {*} [options] Override http request option. * @throws {RequiredError} */ - async listSessionEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + async listSessionEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { const localVarAxiosArgs = await localVarAxiosParamCreator.listSessionEvents(id, sinceSeq, limit, options); const localVarOperationServerIndex = configuration?.serverIndex ?? 0; const localVarOperationServerBasePath = operationServerMap['SessionsApi.listSessionEvents']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, /** - * Starts a streamed turn immediately. Background turns are not supported in this API version. + * Starts a streamed turn immediately. The stream carries the turn\'s `SessionEvent` frames, from `run.session.turn.started` to the event that ends the turn. Background turns are not supported in this API version. * @summary Submit a session turn * @param {string} id * @param {SubmitTurnRequest} submitTurnRequest @@ -485,7 +485,7 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP const localVarFp = SessionsApiFp(configuration) return { /** - * Replays and streams this session\'s durable `run.session.*` events from the owning run event log. The stream remains open until the client disconnects or the server shuts down. + * Replays this session\'s events from `since_seq` (the next unseen event when omitted) as `SessionEvent` frames, then streams new ones as they are recorded. The stream remains open until the client disconnects or the server shuts down. * @summary Attach to session events * @param {string} id * @param {number} [sinceSeq] @@ -524,7 +524,7 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP * @param {*} [options] Override http request option. * @throws {RequiredError} */ - interruptSessionTurn(id: string, turnId: string, options?: RawAxiosRequestConfig): AxiosPromise { + interruptSessionTurn(id: string, turnId: string, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.interruptSessionTurn(id, turnId, options).then((request) => request(axios, basePath)); }, /** @@ -541,7 +541,7 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP return localVarFp.listRunSessions(id, pageLimit, pageOffset, order, options).then((request) => request(axios, basePath)); }, /** - * Returns run event envelopes filtered to this session\'s durable `run.session.*` events. `since_seq` uses the owning run event sequence. + * Returns this session\'s events in order. Events are numbered per session from 1; `since_seq` is the first sequence number to include. * @summary List session events * @param {string} id * @param {number} [sinceSeq] @@ -549,11 +549,11 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP * @param {*} [options] Override http request option. * @throws {RequiredError} */ - listSessionEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig): AxiosPromise { + listSessionEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.listSessionEvents(id, sinceSeq, limit, options).then((request) => request(axios, basePath)); }, /** - * Starts a streamed turn immediately. Background turns are not supported in this API version. + * Starts a streamed turn immediately. The stream carries the turn\'s `SessionEvent` frames, from `run.session.turn.started` to the event that ends the turn. Background turns are not supported in this API version. * @summary Submit a session turn * @param {string} id * @param {SubmitTurnRequest} submitTurnRequest @@ -571,7 +571,7 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP */ export class SessionsApi extends BaseAPI { /** - * Replays and streams this session\'s durable `run.session.*` events from the owning run event log. The stream remains open until the client disconnects or the server shuts down. + * Replays this session\'s events from `since_seq` (the next unseen event when omitted) as `SessionEvent` frames, then streams new ones as they are recorded. The stream remains open until the client disconnects or the server shuts down. * @summary Attach to session events * @param {string} id * @param {number} [sinceSeq] @@ -632,7 +632,7 @@ export class SessionsApi extends BaseAPI { } /** - * Returns run event envelopes filtered to this session\'s durable `run.session.*` events. `since_seq` uses the owning run event sequence. + * Returns this session\'s events in order. Events are numbered per session from 1; `since_seq` is the first sequence number to include. * @summary List session events * @param {string} id * @param {number} [sinceSeq] @@ -645,7 +645,7 @@ export class SessionsApi extends BaseAPI { } /** - * Starts a streamed turn immediately. Background turns are not supported in this API version. + * Starts a streamed turn immediately. The stream carries the turn\'s `SessionEvent` frames, from `run.session.turn.started` to the event that ends the turn. Background turns are not supported in this API version. * @summary Submit a session turn * @param {string} id * @param {SubmitTurnRequest} submitTurnRequest diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 5ed8c356e..876c38994 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -242,6 +242,7 @@ export * from './paginated-run-list'; export * from './paginated-run-stage-list'; export * from './paginated-run-stream-list'; export * from './paginated-saved-query-list'; +export * from './paginated-session-event-list'; export * from './paginated-session-list'; export * from './paginated-workflow-list-response'; export * from './pagination-meta'; @@ -469,6 +470,8 @@ export * from './server-slate-db-settings'; export * from './server-storage-settings'; export * from './server-web-settings'; export * from './session-detail'; +export * from './session-event'; +export * from './session-event-name'; export * from './session-status'; export * from './session-summary'; export * from './session-turn'; diff --git a/lib/packages/fabro-api-client/src/models/paginated-session-event-list.ts b/lib/packages/fabro-api-client/src/models/paginated-session-event-list.ts new file mode 100644 index 000000000..3718a4042 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/paginated-session-event-list.ts @@ -0,0 +1,29 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.2.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + +// May contain unused imports in some cases +// @ts-ignore +import type { PaginationMeta } from './pagination-meta'; +// May contain unused imports in some cases +// @ts-ignore +import type { SessionEvent } from './session-event'; + +/** + * A page of a session\'s events, in sequence order. + */ +export interface PaginatedSessionEventList { + 'data': Array; + 'meta': PaginationMeta; +} diff --git a/lib/packages/fabro-api-client/src/models/session-detail.ts b/lib/packages/fabro-api-client/src/models/session-detail.ts index 27d8be07d..1e4855a0d 100644 --- a/lib/packages/fabro-api-client/src/models/session-detail.ts +++ b/lib/packages/fabro-api-client/src/models/session-detail.ts @@ -21,7 +21,7 @@ import type { SessionStatus } from './session-status'; import type { SessionTurn } from './session-turn'; /** - * Session metadata plus the highest run event sequence the session\'s event stream has reached. The conversation itself is held by the server\'s durable session record and is not returned over the API. + * Session metadata plus the sequence number of the session\'s latest event. The conversation itself is held by the server\'s durable session record and is not returned over the API. */ export interface SessionDetail { /** diff --git a/lib/packages/fabro-api-client/src/models/session-event-name.ts b/lib/packages/fabro-api-client/src/models/session-event-name.ts new file mode 100644 index 000000000..39f8fca1e --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/session-event-name.ts @@ -0,0 +1,34 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.2.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * What a session event records. + */ + +export const SessionEventName = { + RUN_SESSION_CREATED: 'run.session.created', + RUN_SESSION_TURN_STARTED: 'run.session.turn.started', + RUN_SESSION_USER_MESSAGE: 'run.session.user_message', + RUN_SESSION_ASSISTANT_DELTA: 'run.session.assistant_delta', + RUN_SESSION_ASSISTANT_MESSAGE: 'run.session.assistant_message', + RUN_SESSION_TOOL_CALL_STARTED: 'run.session.tool_call.started', + RUN_SESSION_TOOL_CALL_COMPLETED: 'run.session.tool_call.completed', + RUN_SESSION_TURN_SUCCEEDED: 'run.session.turn.succeeded', + RUN_SESSION_TURN_FAILED: 'run.session.turn.failed', + RUN_SESSION_TURN_INTERRUPTED: 'run.session.turn.interrupted' +} as const; + +export type SessionEventName = typeof SessionEventName[keyof typeof SessionEventName]; diff --git a/lib/packages/fabro-api-client/src/models/session-event.ts b/lib/packages/fabro-api-client/src/models/session-event.ts new file mode 100644 index 000000000..c1f04e2e0 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/session-event.ts @@ -0,0 +1,33 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.2.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + +// May contain unused imports in some cases +// @ts-ignore +import type { SessionEventName } from './session-event-name'; + +/** + * One recorded event of an Ask Fabro session: its position in the session (`seq`, from 1), the session and run it belongs to, when it was recorded, the event\'s name and its properties. Every event but `run.session.created` carries the `turn_id` it belongs to in its properties. + */ +export interface SessionEvent { + 'seq': number; + /** + * Durable session identifier. + */ + 'session_id': string; + 'run_id': string; + 'ts': string; + 'event': SessionEventName; + 'properties': { [key: string]: any; }; +}