diff --git a/apps/fabro-web/app/lib/api-client.ts b/apps/fabro-web/app/lib/api-client.ts index 1dc2d47b1..334c4411a 100644 --- a/apps/fabro-web/app/lib/api-client.ts +++ b/apps/fabro-web/app/lib/api-client.ts @@ -16,6 +16,7 @@ import { RunOutputsApi, RunsApi, SecretsApi, + SessionsApi, SettingsApi, SystemApi, WorkflowsApi, @@ -114,6 +115,11 @@ export const secretsApi = new SecretsApi( "", generatedAxios, ); +export const sessionsApi = new SessionsApi( + generatedApiConfiguration, + "", + generatedAxios, +); export const settingsApi = new SettingsApi( generatedApiConfiguration, "", @@ -194,6 +200,34 @@ function apiErrorFromAxios(error: unknown): ApiError | null { }); } +export async function apiErrorFromFetchResponse(response: Response): Promise { + if (response.ok) return null; + + const body = await readFetchErrorBody(response); + const requestId = requestIdFromHeaders(response.headers) ?? extractRequestId(body); + return new ApiError({ + status: response.status, + message: extractErrorDetail(body) ?? (response.statusText || `HTTP ${response.status}`), + requestId, + body, + }); +} + +async function readFetchErrorBody(response: Response): Promise { + const contentType = response.headers.get("content-type") ?? ""; + if (contentType.includes("application/json")) { + return response.json().catch(() => null); + } + + const text = await response.text().catch(() => ""); + if (!text) return null; + try { + return JSON.parse(text); + } catch { + return text; + } +} + function extractErrorDetail(body: unknown): string | null { if (!body || typeof body !== "object") return null; const errors = (body as Record).errors; diff --git a/apps/fabro-web/app/lib/session-stream.test.ts b/apps/fabro-web/app/lib/session-stream.test.ts new file mode 100644 index 000000000..15b88a8db --- /dev/null +++ b/apps/fabro-web/app/lib/session-stream.test.ts @@ -0,0 +1,146 @@ +import { afterEach, describe, expect, mock, test } from "bun:test"; + +import { ApiError } from "./api-client"; +import { + attachSessionEvents, + streamSessionTurn, + type SessionStreamEvent, +} from "./session-stream"; + +const encoder = new TextEncoder(); + +afterEach(() => { + mock.restore(); +}); + +function streamResponse(chunks: string[], status = 200, headers: HeadersInit = {}) { + return new Response( + new ReadableStream({ + start(controller) { + for (const chunk of chunks) { + controller.enqueue(encoder.encode(chunk)); + } + controller.close(); + }, + }), + { + status, + headers: { + "content-type": "text/event-stream", + ...headers, + }, + }, + ); +} + +describe("session stream helpers", () => { + test("posts a turn and parses chunked SSE event envelopes", async () => { + const events: SessionStreamEvent[] = []; + const fetchMock = mock(() => + Promise.resolve( + streamResponse( + [ + "id: 3\nevent: run.session.turn.started\n", + 'data: {"seq":3,"event":{"event":"run.session.turn.started","properties":{"turn_id":"turn_1"}}}\n\n', + ], + 200, + { "x-fabro-turn-id": "turn_1" }, + ), + ), + ); + + const result = await streamSessionTurn({ + sessionId: "ses_1", + input: "Summarize", + turnId: "turn_1", + fetchImpl: fetchMock, + onEvent: (event) => events.push(event), + }); + + expect(result.turnId).toBe("turn_1"); + expect(fetchMock.mock.calls[0]?.[0]).toBe("/api/v1/sessions/ses_1/turns"); + expect(JSON.parse(fetchMock.mock.calls[0]?.[1]?.body as string)).toEqual({ + input: "Summarize", + turn_id: "turn_1", + }); + expect(events).toHaveLength(1); + expect(events[0]?.seq).toBe(3); + expect(events[0]?.event.event).toBe("run.session.turn.started"); + }); + + test("attaches to session events from a run sequence", async () => { + const events: SessionStreamEvent[] = []; + const fetchMock = mock(() => + Promise.resolve( + streamResponse([ + 'data: {"seq":7,"event":{"event":"run.session.assistant_message","properties":{}}}\n\n', + ]), + ), + ); + + await attachSessionEvents({ + sessionId: "ses_1", + sinceSeq: 7, + fetchImpl: fetchMock, + onEvent: (event) => events.push(event), + }); + + expect(fetchMock.mock.calls[0]?.[0]).toBe( + "/api/v1/sessions/ses_1/attach?since_seq=7", + ); + expect(events[0]?.seq).toBe(7); + }); + + test("parses CRLF-delimited SSE frames", async () => { + const events: SessionStreamEvent[] = []; + const fetchMock = mock(() => + Promise.resolve( + streamResponse([ + 'data: {"seq":8,"event":{"event":"run.session.assistant_message","properties":{}}}\r\n\r\n', + ]), + ), + ); + + await attachSessionEvents({ + sessionId: "ses_1", + fetchImpl: fetchMock, + onEvent: (event) => events.push(event), + }); + + expect(events[0]?.seq).toBe(8); + }); + + test("converts non-2xx responses to ApiError", async () => { + const fetchMock = mock(() => + Promise.resolve( + new Response( + JSON.stringify({ + errors: [{ + status: "409", + title: "Conflict", + detail: "Session already has an active turn.", + code: "session_active_turn", + }], + }), + { + status: 409, + headers: { "x-request-id": "req_1" }, + }, + ), + ), + ); + + await expect( + streamSessionTurn({ + sessionId: "ses_1", + input: "Summarize", + fetchImpl: fetchMock, + onEvent: () => {}, + }), + ).rejects.toMatchObject({ + status: 409, + requestId: "req_1", + message: "Session already has an active turn.", + } satisfies Partial); + }); +}); diff --git a/apps/fabro-web/app/lib/session-stream.ts b/apps/fabro-web/app/lib/session-stream.ts new file mode 100644 index 000000000..5747239fe --- /dev/null +++ b/apps/fabro-web/app/lib/session-stream.ts @@ -0,0 +1,141 @@ +import { + SessionsApiAxiosParamCreator, + type EventEnvelope, + type SubmitTurnRequest, +} from "@qltysh/fabro-api-client"; + +import { + apiErrorFromFetchResponse, + generatedApiConfiguration, +} from "./api-client"; + +export type SessionStreamEvent = EventEnvelope; + +type FetchLike = ( + input: string, + init?: RequestInit, +) => Promise; + +interface SessionStreamOptions { + sessionId: string; + signal?: AbortSignal; + fetchImpl?: FetchLike; + onEvent: (event: SessionStreamEvent) => void; +} + +export interface StreamSessionTurnOptions extends SessionStreamOptions { + input: string; + turnId?: string; +} + +export interface StreamSessionTurnResult { + turnId: string | null; +} + +export interface AttachSessionEventsOptions extends SessionStreamOptions { + sinceSeq?: number; +} + +export async function streamSessionTurn({ + sessionId, + input, + turnId, + signal, + fetchImpl = fetch, + onEvent, +}: StreamSessionTurnOptions): Promise { + const body: SubmitTurnRequest = { input }; + if (turnId) body.turn_id = turnId; + + const request = await SessionsApiAxiosParamCreator( + generatedApiConfiguration, + ).submitSessionTurn(sessionId, body, { signal }); + const response = await fetchImpl(request.url, fetchInitFromAxiosRequest(request.options)); + await throwIfApiError(response); + + await readEventStream(response, onEvent); + return { turnId: response.headers.get("x-fabro-turn-id") }; +} + +export async function attachSessionEvents({ + sessionId, + sinceSeq, + signal, + fetchImpl = fetch, + onEvent, +}: AttachSessionEventsOptions): Promise { + const request = await SessionsApiAxiosParamCreator( + generatedApiConfiguration, + ).attachSessionEvents(sessionId, sinceSeq, { signal }); + const response = await fetchImpl(request.url, fetchInitFromAxiosRequest(request.options)); + await throwIfApiError(response); + + await readEventStream(response, onEvent); +} + +async function throwIfApiError(response: Response): Promise { + const error = await apiErrorFromFetchResponse(response); + if (error) throw error; +} + +function fetchInitFromAxiosRequest(options: { + method?: string; + headers?: unknown; + data?: unknown; + signal?: unknown; +}): RequestInit { + const init: RequestInit = { + method: options.method, + credentials: "same-origin", + headers: options.headers as HeadersInit, + signal: options.signal as AbortSignal | undefined, + }; + if (options.data !== undefined) { + init.body = typeof options.data === "string" + ? options.data + : JSON.stringify(options.data); + } + return init; +} + +async function readEventStream( + response: Response, + onEvent: (event: SessionStreamEvent) => void, +): Promise { + if (!response.body) return; + + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + + while (true) { + const { value, done } = await reader.read(); + if (done) break; + buffer += decoder.decode(value, { stream: true }); + buffer = drainSseBuffer(buffer, onEvent); + } + + buffer += decoder.decode(); + drainSseBuffer(`${buffer}\n\n`, onEvent); +} + +function drainSseBuffer( + buffer: string, + onEvent: (event: SessionStreamEvent) => void, +): string { + let cursor = 0; + while (true) { + const match = /\r?\n\r?\n/g.exec(buffer.slice(cursor)); + if (!match) return buffer.slice(cursor); + const next = cursor + match.index; + const frame = buffer.slice(cursor, next); + cursor = next + match[0].length; + const data = frame + .split(/\r?\n/) + .filter((line) => line.startsWith("data:")) + .map((line) => line.slice("data:".length).trimStart()) + .join("\n"); + if (!data) continue; + onEvent(JSON.parse(data) as SessionStreamEvent); + } +} diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 417acf422..e1e3b1498 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -822,6 +822,26 @@ paths: operationId: listRunSessions tags: [Sessions] summary: List run sessions + parameters: + - name: page[limit] + in: query + schema: + type: integer + default: 20 + minimum: 1 + maximum: 100 + - name: page[offset] + in: query + schema: + type: integer + default: 0 + minimum: 0 + - name: order + in: query + schema: + type: string + enum: [updated_desc, created_desc] + default: updated_desc responses: "200": description: Ask Fabro sessions for the run @@ -870,11 +890,11 @@ paths: summary: Get session responses: "200": - description: Session record + description: Session detail content: application/json: schema: - $ref: "#/components/schemas/SessionRecord" + $ref: "#/components/schemas/SessionDetail" "404": description: Session not found headers: @@ -884,6 +904,84 @@ paths: application/json: schema: $ref: "#/components/schemas/ErrorResponse" + /api/v1/sessions/{id}/events: + parameters: + - name: id + in: path + required: true + schema: + $ref: "#/components/schemas/SessionId" + get: + 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. + parameters: + - name: since_seq + in: query + schema: + type: integer + default: 1 + minimum: 1 + - name: limit + in: query + schema: + type: integer + default: 100 + minimum: 1 + maximum: 1000 + responses: + "200": + description: Session-scoped run events + content: + application/json: + schema: + $ref: "#/components/schemas/PaginatedEventList" + "404": + description: Session not found + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + + /api/v1/sessions/{id}/attach: + parameters: + - name: id + in: path + required: true + schema: + $ref: "#/components/schemas/SessionId" + get: + 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. + parameters: + - name: since_seq + in: query + schema: + type: integer + minimum: 1 + responses: + "200": + description: Streamed session-scoped run events + content: + text/event-stream: + schema: + type: string + "404": + description: Session not found + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + /api/v1/sessions/{id}/turns: parameters: - name: id @@ -905,6 +1003,11 @@ paths: responses: "200": description: Streamed session events + headers: + x-fabro-turn-id: + description: Durable turn id accepted for this streamed turn. + schema: + $ref: "#/components/schemas/TurnId" content: text/event-stream: schema: @@ -932,6 +1035,10 @@ paths: headers: x-request-id: $ref: "#/components/headers/XRequestId" + x-fabro-active-turn-id: + description: Durable id of the currently active turn. + schema: + $ref: "#/components/schemas/TurnId" content: application/json: schema: @@ -5338,6 +5445,22 @@ components: type: string enum: [idle, running, failed] + SessionTurn: + description: Currently active durable session turn. + type: object + required: + - id + - started_at + - input + properties: + id: + $ref: "#/components/schemas/TurnId" + started_at: + type: string + format: date-time + input: + type: string + SessionMessage: description: Persisted full-fidelity session transcript message. type: object @@ -5373,6 +5496,7 @@ components: - id - run_id - status + - active_turn - created_at - updated_at properties: @@ -5386,6 +5510,10 @@ components: $ref: "#/components/schemas/SessionStatus" model: type: ["string", "null"] + active_turn: + oneOf: + - $ref: "#/components/schemas/SessionTurn" + - type: "null" created_at: type: string format: date-time @@ -5400,6 +5528,7 @@ components: - id - run_id - status + - active_turn - created_at - updated_at properties: @@ -5413,6 +5542,10 @@ components: $ref: "#/components/schemas/SessionStatus" model: type: ["string", "null"] + active_turn: + oneOf: + - $ref: "#/components/schemas/SessionTurn" + - type: "null" created_at: type: string format: date-time @@ -5420,6 +5553,47 @@ components: type: string format: date-time + SessionDetail: + description: Session metadata plus durable transcript projection. + type: object + required: + - id + - run_id + - status + - active_turn + - created_at + - updated_at + - messages + - last_seq + properties: + id: + $ref: "#/components/schemas/SessionId" + run_id: + type: string + title: + type: ["string", "null"] + status: + $ref: "#/components/schemas/SessionStatus" + model: + type: ["string", "null"] + active_turn: + oneOf: + - $ref: "#/components/schemas/SessionTurn" + - type: "null" + created_at: + type: string + format: date-time + updated_at: + type: string + format: date-time + messages: + type: array + items: + $ref: "#/components/schemas/SessionMessage" + last_seq: + type: integer + minimum: 0 + CreateRunSessionRequest: type: object properties: @@ -5436,6 +5610,8 @@ components: properties: input: type: string + turn_id: + $ref: "#/components/schemas/TurnId" PaginatedSessionList: description: Paginated list of sessions. @@ -7908,6 +8084,7 @@ components: - timestamps - timing - billing + - ask_fabro - diff - pull_request - current_question @@ -7974,6 +8151,8 @@ components: oneOf: - $ref: "#/components/schemas/RunBillingSummary" - type: "null" + ask_fabro: + $ref: "#/components/schemas/AskFabro" diff: oneOf: - $ref: "#/components/schemas/DiffSummary" @@ -7991,6 +8170,27 @@ components: links: $ref: "#/components/schemas/RunLinks" + AskFabro: + description: Readiness and defaults for starting an Ask Fabro session on this run. + type: object + required: + - available + - unavailable_reason + - default_model + properties: + available: + type: boolean + unavailable_reason: + type: ["string", "null"] + enum: + - feature_disabled + - no_sandbox + - sandbox_not_ready + - llm_unconfigured + - null + default_model: + type: ["string", "null"] + WorkflowRef: type: object required: [slug, name, graph_name, node_count, edge_count] diff --git a/lib/crates/fabro-api/build.rs b/lib/crates/fabro-api/build.rs index 9a5324342..e1471cf35 100644 --- a/lib/crates/fabro-api/build.rs +++ b/lib/crates/fabro-api/build.rs @@ -475,12 +475,15 @@ fn main() { ("SandboxState", "fabro_types::SandboxState", &[]), ("SandboxResources", "fabro_types::SandboxResources", &[]), ("SandboxTimestamps", "fabro_types::SandboxTimestamps", &[]), + ("AskFabro", "fabro_types::AskFabro", &[]), ("SessionId", "fabro_types::SessionId", &[]), ("TurnId", "fabro_types::TurnId", &[]), ("SessionStatus", "fabro_types::SessionStatus", &[]), + ("SessionTurn", "fabro_types::SessionTurn", &[]), ("SessionMessage", "fabro_types::SessionMessage", &[]), ("SessionRecord", "fabro_types::SessionRecord", &[]), ("SessionSummary", "fabro_types::SessionSummary", &[]), + ("SessionDetail", "fabro_types::SessionDetail", &[]), ]; for (name, path, impls) in replacements { settings.with_replacement(*name, *path, impls.iter().copied()); diff --git a/lib/crates/fabro-api/src/lib.rs b/lib/crates/fabro-api/src/lib.rs index e20c484cb..f1ff83b16 100644 --- a/lib/crates/fabro-api/src/lib.rs +++ b/lib/crates/fabro-api/src/lib.rs @@ -32,8 +32,8 @@ pub mod types { BlockedReason, FailureReason, RunControlAction, RunStatus, SuccessReason, }; pub use fabro_types::{ - AuthMethod, BilledTokenCounts, CommandTermination, Conclusion, DiffStats, DiffSummary, - DirtyStatus, EventEnvelope, ExecOutputTail, FailureCategory, FailureDetail, + AskFabro, AuthMethod, BilledTokenCounts, CommandTermination, Conclusion, DiffStats, + DiffSummary, DirtyStatus, EventEnvelope, ExecOutputTail, FailureCategory, FailureDetail, FailureSignature, GitContext, IdpIdentity, InterviewOption, InterviewQuestionRecord, PairId, PairMessageId, PairMessageRecord, PairMessageRequest, PairRecord, PairStartRequest, PairStatus, PairTarget, PairTargetSelector, PairTranscriptEntry, PairTranscriptResponse, @@ -45,9 +45,10 @@ pub mod types { RunSandboxRuntime, RunServerProvenance, SandboxDetails, SandboxNetwork, SandboxNetworkPolicy, SandboxNetworkPolicyMode, SandboxProvider, SandboxResources, SandboxService, SandboxServiceListResponse, SandboxState, SandboxTimestamps, - SecretMetadata, SecretType, ServerSettings, SessionId, SessionMessage, SessionRecord, - SessionStatus, SessionSummary, StageCompletion, StageHandler, StageOutcome, - StageProjection, StageState, SystemActorKind, TurnId, UserPrincipal, WorkflowSettings, + SecretMetadata, SecretType, ServerSettings, SessionDetail, SessionId, SessionMessage, + SessionRecord, SessionStatus, SessionSummary, SessionTurn, StageCompletion, StageHandler, + StageOutcome, StageProjection, StageState, SystemActorKind, TurnId, UserPrincipal, + WorkflowSettings, }; pub use crate::generated::types::*; diff --git a/lib/crates/fabro-api/tests/run_event_round_trip.rs b/lib/crates/fabro-api/tests/run_event_round_trip.rs index 14a2e7dd9..0ce55ca42 100644 --- a/lib/crates/fabro-api/tests/run_event_round_trip.rs +++ b/lib/crates/fabro-api/tests/run_event_round_trip.rs @@ -169,6 +169,26 @@ fn run_event_round_trips_agent_interrupt_injected() { assert_run_event_round_trip(value); } +#[test] +fn run_event_turn_failed_defaults_code_for_legacy_payloads() { + let value = json!({ + "id": "evt_session_failed", + "ts": "2026-05-20T12:00:00Z", + "run_id": fixtures::RUN_1, + "event": "run.session.turn.failed", + "session_id": "01HZX6M0P7SE4VJ9Y3X2B8E9QF", + "properties": { + "turn_id": "01HZX6M29F1CD5YYMHT1F5D7WQ", + "error": "provider unavailable" + } + }); + + let event: ApiRunEvent = serde_json::from_value(value).unwrap(); + let round_trip = serde_json::to_value(event).unwrap(); + assert_eq!(round_trip["properties"]["code"], "agent_error"); + assert_eq!(round_trip["properties"]["retryable"], false); +} + #[test] fn run_event_round_trips_stage_started() { let value = json!({ diff --git a/lib/crates/fabro-api/tests/run_summary_round_trip.rs b/lib/crates/fabro-api/tests/run_summary_round_trip.rs index 3b7887e72..d310dd7a5 100644 --- a/lib/crates/fabro-api/tests/run_summary_round_trip.rs +++ b/lib/crates/fabro-api/tests/run_summary_round_trip.rs @@ -5,8 +5,9 @@ use chrono::{TimeZone, Utc}; use fabro_api::types::{RepositoryRef as ApiRepositoryRef, Run as ApiRun}; use fabro_types::status::{RunStatus, SuccessReason}; use fabro_types::{ - DiffSummary, PullRequestLink, RepositoryProvider, RepositoryRef, Run, RunBillingSummary, RunId, - RunLifecycle, RunLinks, RunOrigin, RunTimestamps, RunTiming, WorkflowRef, + AskFabro, AskFabroUnavailableReason, DiffSummary, PullRequestLink, RepositoryProvider, + RepositoryRef, Run, RunBillingSummary, RunId, RunLifecycle, RunLinks, RunOrigin, RunTimestamps, + RunTiming, WorkflowRef, }; use serde_json::json; @@ -67,6 +68,11 @@ fn run_summary_json_matches_openapi_shape() { billing: Some(RunBillingSummary { total_usd_micros: Some(123), }), + ask_fabro: AskFabro { + available: false, + unavailable_reason: Some(AskFabroUnavailableReason::SandboxNotReady), + default_model: Some("gpt-5.4".to_string()), + }, diff: Some(DiffSummary { files_changed: 3, additions: 12, @@ -138,6 +144,11 @@ fn run_summary_json_matches_openapi_shape() { "billing": { "total_usd_micros": 123 }, + "ask_fabro": { + "available": false, + "unavailable_reason": "sandbox_not_ready", + "default_model": "gpt-5.4" + }, "diff": { "files_changed": 3, "additions": 12, @@ -225,6 +236,7 @@ fn run_summary_deserializes_when_optional_fields_are_absent() { assert_eq!(summary.lifecycle.pending_control, None); assert_eq!(summary.timing.map(|t| t.wall_time_ms), None); assert_eq!(summary.billing, None); + assert_eq!(summary.ask_fabro, AskFabro::default()); assert_eq!(summary.superseded_by, None); assert_eq!(summary.diff, None); assert_eq!(summary.pull_request, None); diff --git a/lib/crates/fabro-api/tests/session_contract_round_trip.rs b/lib/crates/fabro-api/tests/session_contract_round_trip.rs new file mode 100644 index 000000000..c8b26fb7a --- /dev/null +++ b/lib/crates/fabro-api/tests/session_contract_round_trip.rs @@ -0,0 +1,79 @@ +use std::any::{TypeId, type_name}; + +use chrono::{TimeZone, Utc}; +use fabro_api::types::{ + SessionDetail as ApiSessionDetail, SessionRecord as ApiSessionRecord, + SessionSummary as ApiSessionSummary, SessionTurn as ApiSessionTurn, SubmitTurnRequest, +}; +use fabro_types::{ + SessionDetail, SessionId, SessionMessage, SessionRecord, SessionStatus, SessionSummary, + SessionTurn, TurnId, fixtures, +}; +use serde_json::json; + +#[test] +fn session_contract_reuses_domain_types() { + assert_same_type::(); + assert_same_type::(); + assert_same_type::(); + assert_same_type::(); +} + +#[test] +fn session_detail_round_trips_messages_active_turn_and_last_seq() { + let created_at = Utc.with_ymd_and_hms(2026, 5, 20, 12, 0, 0).unwrap(); + let turn_started_at = Utc.with_ymd_and_hms(2026, 5, 20, 12, 0, 1).unwrap(); + let updated_at = Utc.with_ymd_and_hms(2026, 5, 20, 12, 0, 2).unwrap(); + let session_id = SessionId::new(); + let turn_id = TurnId::new(); + let detail = SessionDetail::new( + SessionRecord { + id: session_id, + run_id: fixtures::RUN_1, + title: Some("Ask Fabro".to_string()), + status: SessionStatus::Running, + model: Some("gpt-5.4".to_string()), + active_turn: Some(SessionTurn { + id: turn_id, + started_at: turn_started_at, + input: "What changed?".to_string(), + }), + created_at, + updated_at, + }, + vec![SessionMessage::user("What changed?", updated_at)], + 7, + ); + + let value = serde_json::to_value(&detail).expect("detail should serialize"); + assert_eq!(value["active_turn"]["id"], turn_id.to_string()); + assert_eq!(value["messages"][0]["kind"], "user"); + assert_eq!(value["last_seq"], 7); + + let round_trip: ApiSessionDetail = + serde_json::from_value(value.clone()).expect("detail should deserialize"); + assert_eq!(serde_json::to_value(round_trip).unwrap(), value); +} + +#[test] +fn submit_turn_request_accepts_client_turn_id() { + let turn_id = TurnId::new(); + let request: SubmitTurnRequest = serde_json::from_value(json!({ + "input": "Summarize this run", + "turn_id": turn_id.to_string() + })) + .expect("submit turn request should deserialize"); + + assert_eq!(request.input, "Summarize this run"); + assert_eq!(request.turn_id, Some(turn_id)); +} + +fn assert_same_type() { + assert_eq!( + TypeId::of::(), + TypeId::of::(), + "{} should be the same type as {}", + type_name::(), + type_name::() + ); +} diff --git a/lib/crates/fabro-client/src/client.rs b/lib/crates/fabro-client/src/client.rs index 48245a409..d8b3e8570 100644 --- a/lib/crates/fabro-client/src/client.rs +++ b/lib/crates/fabro-client/src/client.rs @@ -650,7 +650,8 @@ impl Client { .map_err(|()| anyhow!("server base URL cannot accept path segments"))? .extend(["api", "v1", "sessions", &session_id.to_string(), "turns"]); let body = types::SubmitTurnRequest { - input: input.into(), + input: input.into(), + turn_id: None, }; let response = self .send_http(|http_client| { diff --git a/lib/crates/fabro-server/src/demo/mod.rs b/lib/crates/fabro-server/src/demo/mod.rs index d3f2463db..e2210a272 100644 --- a/lib/crates/fabro-server/src/demo/mod.rs +++ b/lib/crates/fabro-server/src/demo/mod.rs @@ -1151,6 +1151,7 @@ mod runs { billing: total_usd_micros.map(|total_usd_micros| RunBillingSummary { total_usd_micros: Some(total_usd_micros), }), + ask_fabro: Default::default(), diff: None, pull_request: None, current_question: None, diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index d016c9755..7f50b4587 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -82,9 +82,9 @@ use fabro_types::settings::server::{ }; use fabro_types::settings::{InterpString, RunNamespace}; use fabro_types::{ - AgentBackend, EventBody, InterviewQuestionRecord, PairId, PairMessageId, PairTarget, Principal, - PullRequestLink, QuestionType, RunBlobId, RunControlAction, RunEvent, RunId, ServerSettings, - SessionCapability, + AgentBackend, AskFabro, AskFabroUnavailableReason, EventBody, InterviewQuestionRecord, PairId, + PairMessageId, PairTarget, Principal, PullRequestLink, QuestionType, RunBlobId, + RunControlAction, RunEvent, RunId, ServerSettings, SessionCapability, }; use fabro_util::error::{ SharedError, collect_causes, render_compact_with_causes, render_with_causes, @@ -175,6 +175,15 @@ pub struct PaginationParams { pub offset: u32, } +pub(crate) fn paginate_items(items: Vec, pagination: &PaginationParams) -> (Vec, bool) { + let limit = pagination.limit.clamp(1, 100) as usize; + let offset = pagination.offset.min(MAX_PAGE_OFFSET) as usize; + let mut data: Vec<_> = items.into_iter().skip(offset).take(limit + 1).collect(); + let has_more = data.len() > limit; + data.truncate(limit); + (data, has_more) +} + #[derive(serde::Deserialize)] pub(crate) struct DfParams { #[serde(default)] @@ -631,6 +640,43 @@ pub struct AppState { type PullRequestCreateLocks = Arc>>>>; +struct AskFabroReadiness { + feature_enabled: bool, + default_model: Option, +} + +impl AskFabroReadiness { + fn decorate(&self, mut run: fabro_types::Run) -> fabro_types::Run { + run.ask_fabro = self.ask_fabro_for(&run); + run + } + + fn ask_fabro_for(&self, run: &fabro_types::Run) -> AskFabro { + let unavailable_reason = if !self.feature_enabled { + Some(AskFabroUnavailableReason::FeatureDisabled) + } else if run.sandbox.is_none() { + Some(AskFabroUnavailableReason::NoSandbox) + } else if run + .sandbox + .as_ref() + .and_then(|sandbox| sandbox.runtime.as_ref()) + .is_none() + { + Some(AskFabroUnavailableReason::SandboxNotReady) + } else if self.default_model.is_none() { + Some(AskFabroUnavailableReason::LlmUnconfigured) + } else { + None + }; + + AskFabro { + available: unavailable_reason.is_none(), + unavailable_reason, + default_model: self.default_model.clone(), + } + } +} + struct PullRequestCreateGuard { locks: PullRequestCreateLocks, run_id: RunId, @@ -810,6 +856,46 @@ impl AppState { } } + pub(crate) async fn decorate_run_summary(&self, run: fabro_types::Run) -> fabro_types::Run { + self.ask_fabro_readiness().await.decorate(run) + } + + pub(crate) async fn decorate_run_summaries( + &self, + runs: Vec, + ) -> Vec { + let readiness = self.ask_fabro_readiness().await; + runs.into_iter() + .map(|run| readiness.decorate(run)) + .collect() + } + + async fn ask_fabro_readiness(&self) -> AskFabroReadiness { + let feature_enabled = self.server_settings().features.session_sandboxes; + if !feature_enabled { + return AskFabroReadiness { + feature_enabled, + default_model: None, + }; + } + + let provider_ids = self.ready_llm_provider_ids().await; + let default_model = if provider_ids.is_empty() { + None + } else { + Some( + self.catalog() + .default_for_configured_ids(&provider_ids) + .id + .clone(), + ) + }; + AskFabroReadiness { + feature_enabled, + default_model, + } + } + pub(crate) fn vault_or_env(&self, name: &str) -> Option { process_env_var(name).or_else(|| { self.vault diff --git a/lib/crates/fabro-server/src/server/handler/lifecycle.rs b/lib/crates/fabro-server/src/server/handler/lifecycle.rs index 4bb1b4bb9..8360dddf3 100644 --- a/lib/crates/fabro-server/src/server/handler/lifecycle.rs +++ b/lib/crates/fabro-server/src/server/handler/lifecycle.rs @@ -25,7 +25,9 @@ pub(super) fn routes() -> Router> { async fn run_response(state: &AppState, id: RunId, status: StatusCode) -> Response { match state.store.get_cached_summary(&id).await { - Ok(Some(summary)) => (status, Json(summary)).into_response(), + Ok(Some(summary)) => { + (status, Json(state.decorate_run_summary(summary).await)).into_response() + } Ok(None) => ApiError::not_found("Run not found.").into_response(), Err(err) => { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() diff --git a/lib/crates/fabro-server/src/server/handler/runs.rs b/lib/crates/fabro-server/src/server/handler/runs.rs index 8fc6b15ec..99b2b94df 100644 --- a/lib/crates/fabro-server/src/server/handler/runs.rs +++ b/lib/crates/fabro-server/src/server/handler/runs.rs @@ -28,10 +28,10 @@ use tokio::fs; use tracing::info; use super::super::{ - AppState, ListResponse, MAX_PAGE_OFFSET, PaginationParams, RunExecutionMode, - answer_from_request, api_question_from_pending_interview, default_page_limit, - delete_run_internal, load_pending_interview, managed_run, parse_run_id_path, - reject_if_archived, resolve_interp_string, submit_pending_interview_answer, workflow_event, + AppState, ListResponse, PaginationParams, RunExecutionMode, answer_from_request, + api_question_from_pending_interview, default_page_limit, delete_run_internal, + load_pending_interview, managed_run, paginate_items, parse_run_id_path, reject_if_archived, + resolve_interp_string, submit_pending_interview_answer, workflow_event, }; use crate::error::ApiError; use crate::principal_middleware::{ @@ -147,15 +147,6 @@ pub(crate) fn board_columns(include_archived: bool) -> Vec(items: Vec, pagination: &PaginationParams) -> (Vec, bool) { - let limit = pagination.limit.clamp(1, 100) as usize; - let offset = pagination.offset.min(MAX_PAGE_OFFSET) as usize; - let mut data: Vec<_> = items.into_iter().skip(offset).take(limit + 1).collect(); - let has_more = data.len() > limit; - data.truncate(limit); - (data, has_more) -} - async fn list_board_runs( _auth: RequiredUser, State(state): State>, @@ -190,15 +181,20 @@ async fn list_board_runs( }) .collect(); let (page_summaries, has_more) = paginate_items(board_summaries, ¶ms.pagination()); + let data = state + .decorate_run_summaries( + page_summaries + .into_iter() + .map(|entry| entry.summary) + .collect(), + ) + .await; ( StatusCode::OK, Json(serde_json::json!({ "columns": board_columns(include_archived), - "data": page_summaries - .into_iter() - .map(|entry| entry.summary) - .collect::>(), + "data": data, "meta": { "has_more": has_more } })), ) @@ -232,7 +228,11 @@ async fn link_run_parent( return err.into_response(); } if child.parent_id == Some(parent_id) { - return (StatusCode::OK, Json(child)).into_response(); + return ( + StatusCode::OK, + Json(state.decorate_run_summary(child).await), + ) + .into_response(); } let Ok(run_store) = state.store.open_run(&child_id).await else { @@ -268,7 +268,11 @@ async fn unlink_run_parent( } }; let Some(previous_parent_id) = child.parent_id else { - return (StatusCode::OK, Json(child)).into_response(); + return ( + StatusCode::OK, + Json(state.decorate_run_summary(child).await), + ) + .into_response(); }; let Ok(run_store) = state.store.open_run(&child_id).await else { @@ -321,7 +325,11 @@ async fn validate_parent_link( async fn updated_run_response(state: &AppState, run_id: &RunId) -> Response { match state.store.get_cached_summary(run_id).await { - Ok(Some(summary)) => (StatusCode::OK, Json(summary)).into_response(), + Ok(Some(summary)) => ( + StatusCode::OK, + Json(state.decorate_run_summary(summary).await), + ) + .into_response(), Ok(None) => ApiError::not_found("Run not found.").into_response(), Err(err) => { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() @@ -350,6 +358,7 @@ async fn list_runs( .filter(|summary| include_archived || !summary.lifecycle.archived) .collect::>(); let (data, has_more) = paginate_items(items, ¶ms.pagination()); + let data = state.decorate_run_summaries(data).await; ( StatusCode::OK, Json(serde_json::json!({ @@ -430,7 +439,10 @@ async fn resolve_run( .and_then(|repository| repository.origin_url.clone()) }, ) { - Ok(run) => (StatusCode::OK, Json(run.clone())).into_response(), + Ok(run) => { + let run = state.decorate_run_summary(run.clone()).await; + (StatusCode::OK, Json(run)).into_response() + } Err(err @ (ResolveRunError::InvalidSelector | ResolveRunError::AmbiguousPrefix { .. })) => { ApiError::bad_request(err.to_string()).into_response() } @@ -487,7 +499,11 @@ async fn update_run( } }; if current.title == title { - return (StatusCode::OK, Json(current)).into_response(); + return ( + StatusCode::OK, + Json(state.decorate_run_summary(current).await), + ) + .into_response(); } let run_store = match state.store.open_run(&id).await { @@ -508,7 +524,11 @@ async fn update_run( } match state.store.get_cached_summary(&id).await { - Ok(Some(summary)) => (StatusCode::OK, Json(summary)).into_response(), + Ok(Some(summary)) => ( + StatusCode::OK, + Json(state.decorate_run_summary(summary).await), + ) + .into_response(), Ok(None) => ApiError::not_found("Run not found.").into_response(), Err(err) => { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() @@ -605,7 +625,11 @@ async fn create_run( ); } - (StatusCode::CREATED, Json(summary)).into_response() + ( + StatusCode::CREATED, + Json(state.decorate_run_summary(summary).await), + ) + .into_response() } fn run_provenance(headers: &HeaderMap, subject: &Principal) -> RunProvenance { @@ -703,7 +727,9 @@ async fn get_run_status( State(state): State>, ) -> Response { match state.store.get_cached_summary(&id).await { - Ok(Some(run)) => (StatusCode::OK, Json(run)).into_response(), + Ok(Some(run)) => { + (StatusCode::OK, Json(state.decorate_run_summary(run).await)).into_response() + } Ok(None) => ApiError::not_found("Run not found.").into_response(), Err(err) => { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() diff --git a/lib/crates/fabro-server/src/server/handler/sessions.rs b/lib/crates/fabro-server/src/server/handler/sessions.rs index 491974f7f..1e27fbf62 100644 --- a/lib/crates/fabro-server/src/server/handler/sessions.rs +++ b/lib/crates/fabro-server/src/server/handler/sessions.rs @@ -1,8 +1,8 @@ use std::convert::Infallible; use std::sync::Arc; -use axum::extract::{Path, State}; -use axum::http::StatusCode; +use axum::extract::{Path, Query, State}; +use axum::http::{HeaderValue, StatusCode}; use axum::response::sse::{Event, KeepAlive, Sse}; use axum::response::{IntoResponse, Response}; use axum::routing::{get, post}; @@ -14,7 +14,9 @@ use fabro_agent::{ AgentEvent, AgentProfile, AnthropicProfile, Error as AgentError, GeminiProfile, OpenAiProfile, Session, SessionEvent, SessionOptions, ToolApprovalAdapter, WebFetchSummarizer, }; -use fabro_api::types::{CreateRunSessionRequest, SubmitTurnRequest}; +use fabro_api::types::{ + CreateRunSessionRequest, PaginatedEventList, PaginationMeta, SubmitTurnRequest, +}; use fabro_llm::client::Client as LlmClient; use fabro_model::{AgentProfileKind, Catalog, ModelHandle, ProviderId}; use fabro_sandbox::reconnect::reconnect_for_run; @@ -23,22 +25,26 @@ use fabro_store::{ }; use fabro_types::run_event::{ RunSessionAssistantDeltaProps, RunSessionAssistantMessageProps, RunSessionCreatedProps, - RunSessionToolCallCompletedProps, RunSessionToolCallStartedProps, RunSessionTurnFailedProps, - RunSessionTurnInterruptedProps, RunSessionTurnStartedProps, RunSessionTurnSucceededProps, - RunSessionUserMessageProps, + RunSessionToolCallCompletedProps, RunSessionToolCallStartedProps, RunSessionTurnFailedCode, + RunSessionTurnFailedProps, RunSessionTurnInterruptedProps, RunSessionTurnStartedProps, + RunSessionTurnSucceededProps, RunSessionUserMessageProps, }; use fabro_types::settings::{ModelRef as SettingsModelRef, ModelRegistry, ResolvedModelRef}; use fabro_types::{ - EventBody, EventEnvelope, PermissionLevel, RunEvent, RunId, SessionId, SessionRecord, TurnId, + EventBody, EventEnvelope, PermissionLevel, RunEvent, RunId, SessionDetail, SessionId, TurnId, }; 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, ListResponse}; +use super::super::{ + AppState, EventListParams, PaginationParams, paginate_items, parse_run_id_path, +}; use crate::error::ApiError; use crate::principal_middleware::RequiredUser; use crate::server_secrets::LlmClientResult; @@ -57,6 +63,8 @@ pub(super) fn routes() -> Router> { "/sessions/{id}", get(get_session).fallback(session_method_not_found), ) + .route("/sessions/{id}/events", get(list_session_events)) + .route("/sessions/{id}/attach", get(attach_session_events)) .route( "/sessions/{id}/turns", post(submit_turn).fallback(session_method_not_found), @@ -67,14 +75,31 @@ pub(super) fn routes() -> Router> { ) } +#[derive(Debug, Clone, Copy, Default, serde::Deserialize)] +#[serde(rename_all = "snake_case")] +enum RunSessionListOrder { + #[default] + UpdatedDesc, + CreatedDesc, +} + +#[derive(serde::Deserialize)] +struct ListRunSessionsParams { + #[serde(flatten)] + pagination: PaginationParams, + #[serde(default)] + order: RunSessionListOrder, +} + async fn list_run_sessions( _auth: RequiredUser, State(state): State>, Path(run_id): Path, + Query(params): Query, ) -> Response { - let run_id = match parse_run_id(&run_id) { + let run_id = match parse_run_id_path(&run_id) { Ok(id) => id, - Err(err) => return err.into_response(), + Err(response) => return response, }; let run_store = match open_run_reader(&state, run_id).await { Ok(store) => store, @@ -82,7 +107,28 @@ async fn list_run_sessions( }; match run_store.list_events().await { Ok(events) => { - Json(ListResponse::new(project_run_sessions(run_id, &events))).into_response() + let mut sessions = project_run_sessions(run_id, &events); + match params.order { + RunSessionListOrder::UpdatedDesc => sessions.sort_by(|left, right| { + right + .updated_at + .cmp(&left.updated_at) + .then_with(|| right.created_at.cmp(&left.created_at)) + .then_with(|| right.id.cmp(&left.id)) + }), + RunSessionListOrder::CreatedDesc => sessions.sort_by(|left, right| { + right + .created_at + .cmp(&left.created_at) + .then_with(|| right.id.cmp(&left.id)) + }), + } + let (data, has_more) = paginate_items(sessions, ¶ms.pagination); + Json(serde_json::json!({ + "data": data, + "meta": { "has_more": has_more } + })) + .into_response() } Err(err) => store_error(&err).into_response(), } @@ -94,9 +140,9 @@ async fn create_run_session( Path(run_id): Path, Json(request): Json, ) -> Response { - let run_id = match parse_run_id(&run_id) { + let run_id = match parse_run_id_path(&run_id) { Ok(id) => id, - Err(err) => return err.into_response(), + Err(response) => return response, }; let run_store = match open_run(&state, run_id).await { Ok(store) => store, @@ -157,13 +203,145 @@ async fn get_session( Ok(context) => context, Err(response) => return response, }; - Json(session).into_response() + Json(SessionDetail::new( + session.record, + session.runtime_context, + session.last_seq, + )) + .into_response() } async fn session_method_not_found() -> Response { StatusCode::NOT_FOUND.into_response() } +async fn list_session_events( + _auth: RequiredUser, + State(state): State>, + Path(id): Path, + 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()) + .await + { + Ok(mut data) => { + let limit = params.limit(); + let has_more = data.len() > limit; + data.truncate(limit); + Json(PaginatedEventList { + data, + meta: PaginationMeta { has_more }, + }) + .into_response() + } + Err(err) => store_error(&err).into_response(), + } +} + +#[derive(serde::Deserialize)] +struct AttachSessionParams { + #[serde(default)] + since_seq: Option, +} + +async fn attach_session_events( + _auth: RequiredUser, + State(state): State>, + Path(id): Path, + Query(params): Query, +) -> Response { + const ATTACH_REPLAY_BATCH_LIMIT: usize = 256; + + 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, + }; + let start_seq = match params.since_seq { + Some(seq) => seq.max(1), + None => match run_store.list_events().await { + Ok(events) => events.last().map_or(1, |event| event.seq.saturating_add(1)), + Err(err) => return store_error(&err).into_response(), + }, + }; + let shutdown = state.shutdown_token(); + let (sender, receiver) = mpsc::channel(SESSION_SSE_BUFFER_CAPACITY); + tokio::spawn(async move { + 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 { + return; + } + } + } + + 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 { + return; + } + } + } + } + } + } + }); + + Sse::new(ReceiverStream::new(receiver)) + .keep_alive(KeepAlive::default()) + .into_response() +} + async fn submit_turn( _auth: RequiredUser, State(state): State>, @@ -180,12 +358,25 @@ async fn submit_turn( }; let input = request.input; - let turn_id = TurnId::new(); + let turn_id = match request.turn_id { + Some(turn_id) => turn_id, + None => TurnId::new(), + }; let turn_lease = match state.session_runtimes().reserve_turn(session_id, turn_id) { Ok(lease) => lease, - Err(StartTurnError::ActiveTurn) => { - return ApiError::new(StatusCode::CONFLICT, "Session already has an active turn.") - .into_response(); + Err(StartTurnError::ActiveTurn { turn_id }) => { + let mut response = ApiError::with_code( + StatusCode::CONFLICT, + "Session already has an active turn.", + "session_active_turn", + ) + .into_response(); + if let Ok(value) = HeaderValue::from_str(&turn_id.to_string()) { + response + .headers_mut() + .insert("x-fabro-active-turn-id", value); + } + return response; } }; @@ -213,9 +404,13 @@ async fn submit_turn( tokio::spawn(run_streaming_turn( state, run_id, run_store, session, turn_id, input, sender, turn_lease, )); - Sse::new(ReceiverStream::new(receiver)) + let mut response = Sse::new(ReceiverStream::new(receiver)) .keep_alive(KeepAlive::default()) - .into_response() + .into_response(); + if let Ok(value) = HeaderValue::from_str(&turn_id.to_string()) { + response.headers_mut().insert("x-fabro-turn-id", value); + } + response } async fn interrupt_turn( @@ -310,11 +505,13 @@ async fn run_streaming_turn( &sender, run_id, session_id, - EventBody::RunSessionTurnFailed(RunSessionTurnFailedProps { + turn_failed_body( turn_id, - error: err.to_string(), - output: None, - }), + err.to_string(), + None, + err.code(), + err.retryable(), + ), Utc::now(), ) .await; @@ -370,11 +567,8 @@ async fn run_streaming_turn( error: Some(err.to_string()), }) } else { - EventBody::RunSessionTurnFailed(RunSessionTurnFailedProps { - turn_id, - error: err.to_string(), - output: outcome.output, - }) + let code = agent_failure_code(&err); + 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()) @@ -387,11 +581,13 @@ async fn run_streaming_turn( &sender, run_id, session_id, - EventBody::RunSessionTurnFailed(RunSessionTurnFailedProps { + turn_failed_body( turn_id, - error: err.to_string(), - output: outcome.output, - }), + err.to_string(), + outcome.output, + RunSessionTurnFailedCode::AgentError, + false, + ), Utc::now(), ) .await; @@ -404,13 +600,45 @@ struct TurnExecutionOutcome { output: Option, } +#[derive(Debug, thiserror::Error)] +enum AskFabroBuildError { + #[error("{0}")] + LlmUnconfigured(String), + #[error("{0}")] + ModelUnavailable(String), + #[error("run has no sandbox available for Ask Fabro")] + NoSandbox, + #[error("run sandbox is unavailable for Ask Fabro: {0}")] + SandboxUnavailable(#[source] anyhow::Error), + #[error("failed to create Ask Fabro agent session: {0}")] + Agent(#[source] anyhow::Error), +} + +impl AskFabroBuildError { + fn code(&self) -> RunSessionTurnFailedCode { + match self { + Self::NoSandbox => RunSessionTurnFailedCode::NoSandbox, + Self::SandboxUnavailable(_) => RunSessionTurnFailedCode::SandboxUnavailable, + Self::LlmUnconfigured(_) => RunSessionTurnFailedCode::LlmUnconfigured, + Self::ModelUnavailable(_) => RunSessionTurnFailedCode::ModelUnavailable, + Self::Agent(_) => RunSessionTurnFailedCode::AgentError, + } + } + + fn retryable(&self) -> bool { + matches!(self, Self::SandboxUnavailable(_)) + } +} + async fn build_agent_session( state: &AppState, run_id: RunId, session: &ProjectedRunSession, -) -> anyhow::Result { +) -> Result { let catalog = state.catalog(); - let llm_result = state.resolve_llm_client().await?; + let llm_result = state.resolve_llm_client().await.map_err(|err| { + AskFabroBuildError::LlmUnconfigured(format!("LLM credentials are not configured: {err}")) + })?; for (provider, issue) in &llm_result.auth_issues { warn!(provider = %provider, error = %issue, "LLM provider unavailable due to auth issue"); } @@ -420,21 +648,39 @@ async fn build_agent_session( let (provider_id, model, profile_kind) = selected_session_model(&catalog, &llm_result, session)?; if !llm_result.client.has_provider(provider_id.as_str()) { - anyhow::bail!("LLM credentials not configured for provider '{provider_id}'"); + let message = format!("LLM credentials not configured for provider '{provider_id}'"); + return if session.record.model.is_some() { + Err(AskFabroBuildError::ModelUnavailable(message)) + } else { + Err(AskFabroBuildError::LlmUnconfigured(message)) + }; } - let run_store = state.store_ref().open_run_reader(&run_id).await?; - let projection = run_store.state().await?; + let run_store = state + .store_ref() + .open_run_reader(&run_id) + .await + .map_err(|err| AskFabroBuildError::Agent(anyhow::Error::new(err)))?; + let projection = run_store + .state() + .await + .map_err(|err| AskFabroBuildError::Agent(anyhow::Error::new(err)))?; let sandbox_record = projection .sandbox .as_ref() - .ok_or_else(|| anyhow::anyhow!("run has no sandbox available for Ask Fabro"))?; + .ok_or(AskFabroBuildError::NoSandbox)?; + if sandbox_record.runtime.is_none() { + return Err(AskFabroBuildError::SandboxUnavailable(anyhow::anyhow!( + "run sandbox runtime is not ready" + ))); + } let sandbox = reconnect_for_run( sandbox_record, state.vault_or_env("DAYTONA_API_KEY"), Some(run_id), ) - .await?; + .await + .map_err(AskFabroBuildError::SandboxUnavailable)?; let sandbox: Arc = Arc::from(sandbox); let profile = build_profile( provider_id, @@ -459,26 +705,32 @@ async fn build_agent_session( config, None, ) - .map_err(Into::into) + .map_err(|err| AskFabroBuildError::Agent(anyhow::Error::new(err))) } fn selected_session_model( catalog: &Catalog, llm_result: &LlmClientResult, session: &ProjectedRunSession, -) -> anyhow::Result<(ProviderId, String, AgentProfileKind)> { +) -> Result<(ProviderId, String, AgentProfileKind), AskFabroBuildError> { let configured_provider_ids = llm_result.provider_ids(); let selected = match session.record.model.as_deref() { - Some(model_id) => catalog - .get(model_id) - .ok_or_else(|| anyhow::anyhow!("session model '{model_id}' is not in the catalog"))?, + Some(model_id) => catalog.get(model_id).ok_or_else(|| { + AskFabroBuildError::ModelUnavailable(format!( + "session model '{model_id}' is not in the catalog" + )) + })?, None => catalog.default_for_configured_ids(&configured_provider_ids), }; let provider_id = selected.provider.clone(); let model = selected.id.clone(); let profile_kind = catalog .effective_agent_profile(&provider_id, Some(&model)) - .ok_or_else(|| anyhow::anyhow!("provider '{provider_id}' is not configured"))?; + .ok_or_else(|| { + AskFabroBuildError::ModelUnavailable(format!( + "provider '{provider_id}' is not configured" + )) + })?; Ok((provider_id, model, profile_kind)) } @@ -672,6 +924,33 @@ fn record_turn_output(output: &mut Option, event: &SessionEvent) { } } +fn turn_failed_body( + turn_id: TurnId, + error: String, + output: Option, + code: RunSessionTurnFailedCode, + retryable: bool, +) -> EventBody { + EventBody::RunSessionTurnFailed(RunSessionTurnFailedProps { + turn_id, + error, + output, + code, + retryable, + }) +} + +fn agent_failure_code(err: &AgentError) -> RunSessionTurnFailedCode { + match err { + AgentError::ToolExecution(message) + if message.contains("denied") || message.contains("not allowed") => + { + RunSessionTurnFailedCode::ToolDenied + } + _ => RunSessionTurnFailedCode::AgentError, + } +} + async fn persist_agent_event( run_store: &RunDatabase, run_id: RunId, @@ -789,6 +1068,38 @@ async fn send_sse_event(sender: &SessionSseSender, event: &EventEnvelope) -> boo .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_attach_sse_event( + sender: &SessionSseSender, + shutdown: &CancellationToken, + event: Event, +) -> bool { + tokio::select! { + biased; + () = shutdown.cancelled() => false, + () = sender.closed() => false, + result = sender.send(Ok(event)) => result.is_ok(), + } +} + +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() +} + async fn load_session( state: &AppState, session_id: SessionId, @@ -812,7 +1123,7 @@ async fn load_session( async fn load_session_read( state: &AppState, session_id: SessionId, -) -> Result<(RunId, SessionRecord), Response> { +) -> Result<(RunId, ProjectedRunSession), Response> { let run_id = match state.store_ref().get_session_run_id(&session_id).await { Ok(Some(run_id)) => run_id, Ok(None) => return Err(ApiError::not_found("Session not found.").into_response()), @@ -823,12 +1134,35 @@ async fn load_session_read( Ok(events) => events, Err(err) => return Err(store_error(&err).into_response()), }; - match project_run_session(run_id, session_id, &events) { + match fabro_store::project_run_session_with_context(run_id, 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().get_session_run_id(&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()); + } + 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(_)) { @@ -857,12 +1191,6 @@ fn store_error(err: &fabro_store::Error) -> ApiError { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) } -fn parse_run_id(value: &str) -> Result { - value - .parse() - .map_err(|err| ApiError::bad_request(format!("Invalid run ID: {err}"))) -} - fn parse_session_id(value: &str) -> Result { value .parse() diff --git a/lib/crates/fabro-server/src/server/session_runtime.rs b/lib/crates/fabro-server/src/server/session_runtime.rs index 019c197e1..922bb482d 100644 --- a/lib/crates/fabro-server/src/server/session_runtime.rs +++ b/lib/crates/fabro-server/src/server/session_runtime.rs @@ -31,8 +31,10 @@ impl SessionRuntimeManager { .active_turn .lock() .expect("session active turn lock poisoned"); - if active.is_some() { - return Err(StartTurnError::ActiveTurn); + if let Some(active) = active.as_ref() { + return Err(StartTurnError::ActiveTurn { + turn_id: active.turn_id, + }); } *active = Some(ActiveTurn { turn_id, @@ -134,7 +136,7 @@ struct ActiveTurn { #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum StartTurnError { - ActiveTurn, + ActiveTurn { turn_id: TurnId }, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] diff --git a/lib/crates/fabro-server/tests/it/api/runs.rs b/lib/crates/fabro-server/tests/it/api/runs.rs index 9682a8c6e..b37260645 100644 --- a/lib/crates/fabro-server/tests/it/api/runs.rs +++ b/lib/crates/fabro-server/tests/it/api/runs.rs @@ -52,6 +52,58 @@ async fn request_json( .await } +#[tokio::test] +async fn run_responses_include_ask_fabro_affordance() { + let settings = settings_from_toml( + r" +_version = 1 + +[features] +session_sandboxes = true +", + ); + let state = fabro_server::test_support::TestAppStateBuilder::new() + .runtime_settings(settings.server_settings, settings.manifest_run_defaults) + .env_lookup(|name| (name == "OPENAI_API_KEY").then(|| "test-key".to_string())) + .build(); + let app = fabro_server::test_support::build_test_router(state); + let created = create_run(&app, minimal_manifest_json(MINIMAL_DOT)).await; + let run_id = created["id"].as_str().unwrap(); + + assert_eq!(created["ask_fabro"]["available"], false); + assert_eq!( + created["ask_fabro"]["unavailable_reason"], + "sandbox_not_ready" + ); + assert_eq!(created["ask_fabro"]["default_model"], "gpt-5.4"); + + let get_request = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}"))) + .body(Body::empty()) + .unwrap(); + let fetched = response_json( + app.clone().oneshot(get_request).await.unwrap(), + StatusCode::OK, + format!("GET /api/v1/runs/{run_id}"), + ) + .await; + assert_eq!(fetched["ask_fabro"], created["ask_fabro"]); + + let list_request = Request::builder() + .method("GET") + .uri(api("/runs")) + .body(Body::empty()) + .unwrap(); + let list = response_json( + app.clone().oneshot(list_request).await.unwrap(), + StatusCode::OK, + "GET /api/v1/runs", + ) + .await; + assert_eq!(list["data"][0]["ask_fabro"], created["ask_fabro"]); +} + #[tokio::test] async fn retrieve_run_settings_returns_dense_snapshot() { let storage_dir = tempfile::tempdir().unwrap(); diff --git a/lib/crates/fabro-server/tests/it/api/sessions.rs b/lib/crates/fabro-server/tests/it/api/sessions.rs index 15554a223..dc8217198 100644 --- a/lib/crates/fabro-server/tests/it/api/sessions.rs +++ b/lib/crates/fabro-server/tests/it/api/sessions.rs @@ -104,6 +104,8 @@ async fn run_bound_session_is_created_as_run_event_and_resolves_by_flat_id() { assert_eq!(fetched["id"], session_id); assert_eq!(fetched["run_id"], run_id); assert_session_metadata_only(&fetched); + assert_eq!(fetched["messages"].as_array().unwrap().len(), 0); + assert!(fetched["active_turn"].is_null()); let events_request = Request::builder() .method("GET") @@ -123,6 +125,7 @@ async fn run_bound_session_is_created_as_run_event_and_resolves_by_flat_id() { .filter(|event| event["session_id"] == session_id) .collect(); assert_eq!(session_events.len(), 1); + assert_eq!(fetched["last_seq"], session_events[0]["seq"]); assert_eq!(session_events[0]["event"], "run.session.created"); assert!( session_events[0]["properties"].get("permissions").is_none(), @@ -273,6 +276,10 @@ async fn session_turn_fails_when_selected_model_provider_is_unconfigured() { .expect("submit-turn request should build"); let response = app.clone().oneshot(request).await.unwrap(); assert_eq!(response.status(), StatusCode::OK); + assert!( + response.headers().contains_key("x-fabro-turn-id"), + "submit turn should return the generated turn id" + ); let events = session_sse_events(response).await; let failed = events @@ -286,6 +293,8 @@ async fn session_turn_fails_when_selected_model_provider_is_unconfigured() { .contains("provider 'openai'"), "failure should be for the selected model provider: {failed:?}" ); + assert_eq!(failed["properties"]["code"], "model_unavailable"); + assert_eq!(failed["properties"]["retryable"], false); } #[tokio::test] @@ -312,7 +321,7 @@ async fn session_metadata_patch_route_is_removed() { } #[tokio::test] -async fn derived_session_read_routes_are_removed() { +async fn unsupported_derived_turn_read_routes_are_removed() { let app = fabro_server::test_support::build_test_router(test_app_state()); let run_id = create_run(&app).await; let created = create_session(&app, &run_id, "Ask Fabro").await; @@ -321,11 +330,7 @@ async fn derived_session_read_routes_are_removed() { .expect("session response should include an id"); let turn_id = fabro_types::TurnId::new(); - for path in [ - format!("/sessions/{session_id}/turns"), - format!("/sessions/{session_id}/turns/{turn_id}"), - format!("/sessions/{session_id}/events"), - ] { + for path in [format!("/sessions/{session_id}/turns/{turn_id}")] { let request = Request::builder() .method("GET") .uri(api(&path)) @@ -340,6 +345,68 @@ async fn derived_session_read_routes_are_removed() { } } +#[tokio::test] +async fn session_events_are_filtered_by_session_and_paginated_by_run_sequence() { + let app = test_app_with_no_providers(); + let run_id = create_run(&app).await; + let first = create_session(&app, &run_id, "First").await; + let second = create_session(&app, &run_id, "Second").await; + let first_id = first["id"].as_str().unwrap(); + let second_id = second["id"].as_str().unwrap(); + let get_first = Request::builder() + .method("GET") + .uri(api(&format!("/sessions/{first_id}"))) + .body(Body::empty()) + .expect("get-session request should build"); + let first_detail = response_json( + app.clone().oneshot(get_first).await.unwrap(), + StatusCode::OK, + format!("GET /api/v1/sessions/{first_id}"), + ) + .await; + let after_first_created_seq = first_detail["last_seq"].as_u64().unwrap() + 1; + + let turn_id = fabro_types::TurnId::new(); + let submit = Request::builder() + .method("POST") + .uri(api(&format!("/sessions/{first_id}/turns"))) + .header("content-type", "application/json") + .body(Body::from(format!( + r#"{{"turn_id":"{turn_id}","input":"Which provider?"}}"# + ))) + .expect("submit-turn request should build"); + let response = app.clone().oneshot(submit).await.unwrap(); + assert_eq!( + response.headers().get("x-fabro-turn-id").unwrap(), + turn_id.to_string().as_str() + ); + let _ = session_sse_events(response).await; + + let request = Request::builder() + .method("GET") + .uri(api(&format!( + "/sessions/{first_id}/events?since_seq={after_first_created_seq}&limit=1" + ))) + .body(Body::empty()) + .expect("session events request should build"); + let page = response_json( + app.clone().oneshot(request).await.unwrap(), + StatusCode::OK, + format!("GET /api/v1/sessions/{first_id}/events"), + ) + .await; + + assert_eq!(page["data"].as_array().unwrap().len(), 1); + assert_eq!(page["data"][0]["session_id"], first_id); + assert_ne!(page["data"][0]["session_id"], second_id); + assert_eq!(page["data"][0]["event"], "run.session.turn.started"); + assert_eq!( + page["data"][0]["properties"]["turn_id"], + turn_id.to_string() + ); + assert_eq!(page["meta"]["has_more"], true); +} + #[tokio::test] async fn inactive_turn_interrupt_returns_conflict() { let app = fabro_server::test_support::build_test_router(test_app_state()); diff --git a/lib/crates/fabro-store/src/run_sessions.rs b/lib/crates/fabro-store/src/run_sessions.rs index b40bfe4a9..4fe396c72 100644 --- a/lib/crates/fabro-store/src/run_sessions.rs +++ b/lib/crates/fabro-store/src/run_sessions.rs @@ -3,7 +3,7 @@ use std::collections::BTreeMap; use fabro_types::run_event::{RunSessionToolCallCompletedProps, RunSessionToolCallStartedProps}; use fabro_types::{ EventBody, EventEnvelope, RunId, SessionId, SessionMessage, SessionRecord, SessionStatus, - SessionSummary, + SessionSummary, SessionTurn, }; use serde_json::json; @@ -11,6 +11,7 @@ use serde_json::json; pub struct ProjectedRunSession { pub record: SessionRecord, pub runtime_context: Vec, + pub last_seq: u32, } pub fn project_run_sessions(run_id: RunId, events: &[EventEnvelope]) -> Vec { @@ -79,18 +80,26 @@ impl RunSessionProjection { let projected = ProjectedRunSession { record, runtime_context: Vec::new(), + last_seq: envelope.seq, }; self.sessions.insert(session_id, projected); } - EventBody::RunSessionTurnStarted(_) => { + EventBody::RunSessionTurnStarted(props) => { if let Some(session) = self.sessions.get_mut(&session_id) { + session.last_seq = envelope.seq; session.record.status = SessionStatus::Running; + session.record.active_turn = Some(SessionTurn { + id: props.turn_id, + started_at: envelope.event.ts, + input: props.input.clone(), + }); session.record.updated_at = envelope.event.ts; } } EventBody::RunSessionUserMessage(props) => { let project_context = self.should_project_context(session_id); if let Some(session) = self.sessions.get_mut(&session_id) { + session.last_seq = envelope.seq; if project_context { session .runtime_context @@ -102,6 +111,7 @@ impl RunSessionProjection { EventBody::RunSessionAssistantMessage(props) => { let project_context = self.should_project_context(session_id); if let Some(session) = self.sessions.get_mut(&session_id) { + session.last_seq = envelope.seq; if project_context { session.runtime_context.push(SessionMessage::Assistant { content: props.text.clone(), @@ -115,9 +125,16 @@ impl RunSessionProjection { session.record.updated_at = envelope.event.ts; } } + EventBody::RunSessionAssistantDelta(_) => { + if let Some(session) = self.sessions.get_mut(&session_id) { + session.last_seq = envelope.seq; + session.record.updated_at = envelope.event.ts; + } + } EventBody::RunSessionToolCallStarted(props) => { let project_context = self.should_project_context(session_id); if let Some(session) = self.sessions.get_mut(&session_id) { + session.last_seq = envelope.seq; if project_context { append_tool_call(session, props); } @@ -127,6 +144,7 @@ impl RunSessionProjection { EventBody::RunSessionToolCallCompleted(props) => { let project_context = self.should_project_context(session_id); if let Some(session) = self.sessions.get_mut(&session_id) { + session.last_seq = envelope.seq; if project_context { append_tool_result(session, props, envelope.event.ts); } @@ -134,10 +152,10 @@ impl RunSessionProjection { } } EventBody::RunSessionTurnFailed(_) => { - self.finish_turn(session_id, true, envelope.event.ts); + self.finish_turn(session_id, true, envelope.event.ts, envelope.seq); } EventBody::RunSessionTurnSucceeded(_) | EventBody::RunSessionTurnInterrupted(_) => { - self.finish_turn(session_id, false, envelope.event.ts); + self.finish_turn(session_id, false, envelope.event.ts, envelope.seq); } _ => {} } @@ -149,13 +167,16 @@ impl RunSessionProjection { session_id: SessionId, failed: bool, timestamp: chrono::DateTime, + seq: u32, ) { if let Some(session) = self.sessions.get_mut(&session_id) { + session.last_seq = seq; session.record.status = if failed { SessionStatus::Failed } else { SessionStatus::Idle }; + session.record.active_turn = None; session.record.updated_at = timestamp; } } @@ -216,8 +237,8 @@ mod tests { use chrono::{TimeZone, Utc}; use fabro_types::run_event::{ RunSessionAssistantMessageProps, RunSessionCreatedProps, RunSessionToolCallCompletedProps, - RunSessionToolCallStartedProps, RunSessionTurnStartedProps, RunSessionTurnSucceededProps, - RunSessionUserMessageProps, + RunSessionToolCallStartedProps, RunSessionTurnFailedCode, RunSessionTurnFailedProps, + RunSessionTurnStartedProps, RunSessionTurnSucceededProps, RunSessionUserMessageProps, }; use fabro_types::{EventBody, EventEnvelope, RunEvent, SessionMessage, TurnId, fixtures}; use serde_json::json; @@ -392,6 +413,94 @@ mod tests { assert!(value.get("deleted_at").is_none()); } + #[test] + fn projection_tracks_active_turn_and_last_matching_sequence() { + let session_id = fabro_types::SessionId::new(); + let other_session_id = fabro_types::SessionId::new(); + let turn_id = TurnId::new(); + let events = vec![ + event( + 1, + session_id, + EventBody::RunSessionCreated(RunSessionCreatedProps { + title: None, + model: None, + }), + ), + event( + 2, + session_id, + EventBody::RunSessionTurnStarted(RunSessionTurnStartedProps { + turn_id, + input: "Summarize".to_string(), + }), + ), + event( + 3, + other_session_id, + EventBody::RunSessionCreated(RunSessionCreatedProps { + title: Some("Other".to_string()), + model: None, + }), + ), + ]; + + let session = project_run_session_with_context(fixtures::RUN_1, session_id, &events) + .expect("session should project from run events"); + + assert_eq!(session.last_seq, 2); + let active = session.record.active_turn.expect("turn should be active"); + assert_eq!(active.id, turn_id); + assert_eq!(active.started_at, events[1].event.ts); + assert_eq!(active.input, "Summarize"); + } + + #[test] + fn projection_clears_active_turn_when_turn_finishes() { + let session_id = fabro_types::SessionId::new(); + let turn_id = TurnId::new(); + + for body in [ + EventBody::RunSessionTurnSucceeded(RunSessionTurnSucceededProps { + turn_id, + output: None, + }), + EventBody::RunSessionTurnFailed(RunSessionTurnFailedProps { + turn_id, + error: "no sandbox".to_string(), + output: None, + code: RunSessionTurnFailedCode::default(), + retryable: false, + }), + ] { + let events = vec![ + event( + 1, + session_id, + EventBody::RunSessionCreated(RunSessionCreatedProps { + title: None, + model: None, + }), + ), + event( + 2, + session_id, + EventBody::RunSessionTurnStarted(RunSessionTurnStartedProps { + turn_id, + input: "Summarize".to_string(), + }), + ), + event(3, session_id, body), + ]; + + let session = project_run_session_with_context(fixtures::RUN_1, session_id, &events) + .expect("session should project from run events"); + + assert_eq!(session.last_seq, 3); + assert_eq!(session.record.active_turn, None); + } + } + fn event(seq: u32, session_id: fabro_types::SessionId, body: EventBody) -> EventEnvelope { let event = RunEvent { id: format!("evt-{seq}"), diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index c9916f64b..4b854d81f 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -8,13 +8,13 @@ use fabro_types::run_event::{ }; use fabro_types::settings::run::RunSandboxSettings; use fabro_types::{ - AgentBackend, BilledModelUsage, Checkpoint, CheckpointRecord, CommandTermination, Conclusion, - EventBody, FailureSignature, InterviewQuestionRecord, Outcome, PendingInterviewRecord, - PullRequestLink, RepositoryRef, Run, RunBillingSummary, RunControlAction, RunDiff, RunEvent, - RunId, RunLifecycle, RunLinks, RunModel, RunOrigin, RunProjection, RunSandbox, - RunSandboxRuntime, RunSpec, RunStatus, RunTimestamps, SandboxProvider, StageCompletion, - StageHandler, StageId, StageOutcome, StageProjection, StageState, StartRecord, WorkflowRef, - first_event_seq, + AgentBackend, AskFabro, BilledModelUsage, Checkpoint, CheckpointRecord, CommandTermination, + Conclusion, EventBody, FailureSignature, InterviewQuestionRecord, Outcome, + PendingInterviewRecord, PullRequestLink, RepositoryRef, Run, RunBillingSummary, + RunControlAction, RunDiff, RunEvent, RunId, RunLifecycle, RunLinks, RunModel, RunOrigin, + RunProjection, RunSandbox, RunSandboxRuntime, RunSpec, RunStatus, RunTimestamps, + SandboxProvider, StageCompletion, StageHandler, StageId, StageOutcome, StageProjection, + StageState, StartRecord, WorkflowRef, first_event_seq, }; use fabro_util::error::render_compact_with_causes; use serde_json::Value; @@ -674,6 +674,7 @@ pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> Run { billing: total_usd_micros.map(|total_usd_micros| RunBillingSummary { total_usd_micros: Some(total_usd_micros), }), + ask_fabro: AskFabro::default(), diff: diff_summary, pull_request: state.pull_request.clone(), current_question, diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index a72b7d1de..1718cb662 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -4,7 +4,7 @@ use std::sync::atomic::{AtomicU32, Ordering}; use bytes::Bytes; use chrono::Utc; -use fabro_types::{RunBlobId, RunEvent, RunId}; +use fabro_types::{RunBlobId, RunEvent, RunId, SessionId}; use futures::Stream; use slatedb::{Db, DbRead}; use tokio::sync::{Mutex, broadcast, mpsc}; @@ -339,6 +339,25 @@ impl RunDatabase { .await } + /// Returns up to `limit + 1` durable Ask Fabro session events for the given + /// session, starting at `start_seq`. The extra item lets callers compute + /// `has_more` without a second read. + pub async fn list_events_for_session_from_with_limit( + &self, + session_id: SessionId, + start_seq: u32, + limit: usize, + ) -> Result> { + list_events_for_session_from_with_limit( + &self.inner.db, + &self.inner.run_id, + session_id, + start_seq, + limit, + ) + .await + } + pub fn watch_events_from( &self, seq: u32, @@ -570,6 +589,57 @@ where Ok(events) } +async fn list_events_for_session_from_with_limit( + db: &R, + run_id: &RunId, + session_id: SessionId, + start_seq: u32, + limit: usize, +) -> Result> +where + R: DbRead + Sync, +{ + #[derive(serde::Deserialize)] + struct SessionEventProbe<'a> { + #[serde(default, borrow)] + session_id: Option<&'a str>, + #[serde(rename = "event", default, borrow)] + event_name: Option<&'a str>, + } + + let session_id_string = session_id.to_string(); + let max_events = limit.saturating_add(1); + let mut iter = db.scan_prefix(keys::run_events_prefix(run_id)).await?; + let mut events = Vec::new(); + while let Some(entry) = iter.next().await? { + let key = key_to_string(&entry.key)?; + let Some(seq) = keys::parse_event_seq(&key) else { + continue; + }; + if seq < start_seq { + continue; + } + + let probe: SessionEventProbe = serde_json::from_slice(&entry.value)?; + if probe.session_id != Some(session_id_string.as_str()) + || !probe + .event_name + .is_some_and(|name| name.starts_with("run.session.")) + { + continue; + } + + let event: RunEvent = serde_json::from_slice(&entry.value)?; + if event.body.is_run_session_event() { + events.push(EventEnvelope { seq, event }); + if events.len() >= max_events { + break; + } + } + } + Ok(events) +} + async fn list_blobs(db: &R) -> Result> where R: DbRead + Sync, @@ -597,7 +667,7 @@ mod tests { use std::sync::Arc; use std::time::Duration; - use fabro_types::{Graph, RunId, StageId, WorkflowSettings}; + use fabro_types::{Graph, RunId, SessionId, StageId, WorkflowSettings}; use object_store::memory::InMemory; use serde_json::json; @@ -624,6 +694,24 @@ mod tests { stage_prompt_payload_for_stage(run_id, idx, node_id, None) } + fn session_message_payload(run_id: &RunId, idx: u32, session_id: SessionId) -> EventPayload { + EventPayload::new( + json!({ + "id": format!("evt-session-{idx}"), + "ts": "2026-04-09T12:00:00Z", + "run_id": run_id.to_string(), + "session_id": session_id.to_string(), + "event": "run.session.user_message", + "properties": { + "turn_id": fabro_types::TurnId::new().to_string(), + "text": format!("message {idx}"), + }, + }), + run_id, + ) + .unwrap() + } + fn run_created_payload(run_id: &RunId) -> EventPayload { EventPayload::new( json!({ @@ -827,4 +915,51 @@ mod tests { let seqs: Vec = events.iter().map(|e| e.seq).collect(); assert_eq!(seqs, vec![3]); } + + #[tokio::test] + async fn list_events_for_session_returns_only_matching_run_session_events() { + let run = fresh_run().await; + let run_id = run.run_id(); + let session_id = SessionId::new(); + let other_session_id = SessionId::new(); + run.append_event(&stage_prompt_payload(&run_id, 1, Some("noise"))) + .await + .unwrap(); + run.append_event(&session_message_payload(&run_id, 2, session_id)) + .await + .unwrap(); + run.append_event(&session_message_payload(&run_id, 3, other_session_id)) + .await + .unwrap(); + run.append_event(&session_message_payload(&run_id, 4, session_id)) + .await + .unwrap(); + + let events = run + .list_events_for_session_from_with_limit(session_id, 1, 100) + .await + .unwrap(); + + let seqs: Vec = events.iter().map(|e| e.seq).collect(); + assert_eq!(seqs, vec![3, 5]); + } + + #[tokio::test] + async fn list_events_for_session_returns_limit_plus_one_for_has_more_signal() { + let run = fresh_run().await; + let run_id = run.run_id(); + let session_id = SessionId::new(); + for idx in 1..=5 { + run.append_event(&session_message_payload(&run_id, idx, session_id)) + .await + .unwrap(); + } + + let events = run + .list_events_for_session_from_with_limit(session_id, 1, 2) + .await + .unwrap(); + + assert_eq!(events.len(), 3); + } } diff --git a/lib/crates/fabro-tool/src/common.rs b/lib/crates/fabro-tool/src/common.rs index d8fca2a94..6043c9a7e 100644 --- a/lib/crates/fabro-tool/src/common.rs +++ b/lib/crates/fabro-tool/src/common.rs @@ -297,6 +297,7 @@ mod tests { }, timing: None, billing: None, + ask_fabro: fabro_types::AskFabro::default(), diff: None, pull_request: None, current_question: None, diff --git a/lib/crates/fabro-tool/src/create.rs b/lib/crates/fabro-tool/src/create.rs index 882b02153..85b8ae69f 100644 --- a/lib/crates/fabro-tool/src/create.rs +++ b/lib/crates/fabro-tool/src/create.rs @@ -513,6 +513,7 @@ mod tests { }, timing: None, billing: None, + ask_fabro: fabro_types::AskFabro::default(), diff: None, pull_request: None, current_question: None, diff --git a/lib/crates/fabro-tool/src/search.rs b/lib/crates/fabro-tool/src/search.rs index 7a57a4299..6d61b524c 100644 --- a/lib/crates/fabro-tool/src/search.rs +++ b/lib/crates/fabro-tool/src/search.rs @@ -466,6 +466,7 @@ mod tests { }, timing: None, billing: None, + ask_fabro: fabro_types::AskFabro::default(), diff: None, pull_request: None, current_question: None, diff --git a/lib/crates/fabro-types/src/lib.rs b/lib/crates/fabro-types/src/lib.rs index 46250fa2a..3f104088f 100644 --- a/lib/crates/fabro-types/src/lib.rs +++ b/lib/crates/fabro-types/src/lib.rs @@ -104,8 +104,8 @@ pub use run_projection::{ }; pub use run_sandbox::{RunSandbox, RunSandboxRuntime}; pub use run_summary::{ - AutomationRef, Run, RunBillingSummary, RunError, RunLifecycle, RunLinks, RunModel, RunOrigin, - RunOriginKind, RunTimestamps, WorkflowRef, + AskFabro, AskFabroUnavailableReason, AutomationRef, Run, RunBillingSummary, RunError, + RunLifecycle, RunLinks, RunModel, RunOrigin, RunOriginKind, RunTimestamps, WorkflowRef, }; pub use run_title::{RunTitleError, infer_run_title, normalize_explicit_run_title}; pub use sandbox_details::{ @@ -119,8 +119,8 @@ pub use sandbox_services::{ }; pub use secret::{SecretMetadata, SecretType}; pub use session::{ - PermissionLevel, SessionId, SessionMessage, SessionRecord, SessionStatus, SessionSummary, - TurnId, + PermissionLevel, SessionDetail, SessionId, SessionMessage, SessionRecord, SessionStatus, + SessionSummary, SessionTurn, TurnId, }; pub use stage_completion::StageCompletion; pub use stage_handler::StageHandler; diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index 27bba053b..8a6ba95ce 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -559,6 +559,10 @@ impl EventBody { } } + pub fn is_run_session_event(&self) -> bool { + self.event_name().starts_with("run.session.") + } + fn properties_value(&self) -> serde_json::Result { if let Self::Unknown { properties, .. } = self { return Ok(properties.clone()); diff --git a/lib/crates/fabro-types/src/run_event/session.rs b/lib/crates/fabro-types/src/run_event/session.rs index 53e4773b8..a44c04003 100644 --- a/lib/crates/fabro-types/src/run_event/session.rs +++ b/lib/crates/fabro-types/src/run_event/session.rs @@ -63,12 +63,42 @@ pub struct RunSessionTurnSucceededProps { 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, + pub turn_id: TurnId, + pub error: String, #[serde(default, skip_serializing_if = "Option::is_none")] - pub output: Option, + pub output: Option, + #[serde(default)] + pub code: RunSessionTurnFailedCode, + #[serde(default)] + pub retryable: bool, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] diff --git a/lib/crates/fabro-types/src/run_summary.rs b/lib/crates/fabro-types/src/run_summary.rs index f3e66b7be..b638613dd 100644 --- a/lib/crates/fabro-types/src/run_summary.rs +++ b/lib/crates/fabro-types/src/run_summary.rs @@ -8,6 +8,37 @@ use crate::{ RunControlAction, RunId, RunSandbox, RunStatus, RunTiming, }; +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct AskFabro { + pub available: bool, + #[serde(default)] + pub unavailable_reason: Option, + #[serde(default)] + pub default_model: Option, +} + +#[derive( + Debug, + Clone, + Copy, + PartialEq, + Eq, + Hash, + Serialize, + Deserialize, + strum::Display, + strum::EnumString, + strum::IntoStaticStr, +)] +#[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] +pub enum AskFabroUnavailableReason { + FeatureDisabled, + NoSandbox, + SandboxNotReady, + LlmUnconfigured, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct Run { pub id: RunId, @@ -40,6 +71,8 @@ pub struct Run { #[serde(default)] pub billing: Option, #[serde(default)] + pub ask_fabro: AskFabro, + #[serde(default)] pub diff: Option, #[serde(default)] pub pull_request: Option, diff --git a/lib/crates/fabro-types/src/session.rs b/lib/crates/fabro-types/src/session.rs index e97101ae5..4e704e239 100644 --- a/lib/crates/fabro-types/src/session.rs +++ b/lib/crates/fabro-types/src/session.rs @@ -48,15 +48,24 @@ impl SessionStatus { } } +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionTurn { + pub id: TurnId, + pub started_at: DateTime, + pub input: String, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct SessionRecord { - pub id: SessionId, - pub run_id: RunId, - pub title: Option, - pub status: SessionStatus, - pub model: Option, - pub created_at: DateTime, - pub updated_at: DateTime, + pub id: SessionId, + pub run_id: RunId, + pub title: Option, + pub status: SessionStatus, + pub model: Option, + #[serde(default)] + pub active_turn: Option, + pub created_at: DateTime, + pub updated_at: DateTime, } impl SessionRecord { @@ -67,6 +76,7 @@ impl SessionRecord { title: None, status: SessionStatus::Idle, model: None, + active_turn: None, created_at: now, updated_at: now, } @@ -75,25 +85,47 @@ impl SessionRecord { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct SessionSummary { - pub id: SessionId, - pub run_id: RunId, - pub title: Option, - pub status: SessionStatus, - pub model: Option, - pub created_at: DateTime, - pub updated_at: DateTime, + pub id: SessionId, + pub run_id: RunId, + pub title: Option, + pub status: SessionStatus, + pub model: Option, + #[serde(default)] + pub active_turn: Option, + pub created_at: DateTime, + pub updated_at: DateTime, } impl From<&SessionRecord> for SessionSummary { fn from(record: &SessionRecord) -> Self { Self { - id: record.id, - run_id: record.run_id, - title: record.title.clone(), - status: record.status, - model: record.model.clone(), - created_at: record.created_at, - updated_at: record.updated_at, + id: record.id, + run_id: record.run_id, + title: record.title.clone(), + status: record.status, + model: record.model.clone(), + active_turn: record.active_turn.clone(), + created_at: record.created_at, + updated_at: record.updated_at, + } + } +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct SessionDetail { + #[serde(flatten)] + pub record: SessionRecord, + #[serde(default)] + pub messages: Vec, + pub last_seq: u32, +} + +impl SessionDetail { + pub fn new(record: SessionRecord, messages: Vec, last_seq: u32) -> Self { + Self { + record, + messages, + last_seq, } } } diff --git a/lib/crates/fabro-types/tests/run_event_serde.rs b/lib/crates/fabro-types/tests/run_event_serde.rs index 8f96a0fc3..7eeb4e5ab 100644 --- a/lib/crates/fabro-types/tests/run_event_serde.rs +++ b/lib/crates/fabro-types/tests/run_event_serde.rs @@ -3,9 +3,10 @@ use std::collections::BTreeMap; use fabro_types::graph::Graph; use fabro_types::run::{DirtyStatus, ForkSourceRef, GitContext, PreRunPushOutcome}; use fabro_types::run_event::run::{RunCreatedProps, RunParentLinkedProps, RunParentUnlinkedProps}; +use fabro_types::run_event::{RunSessionTurnFailedCode, RunSessionTurnFailedProps}; use fabro_types::settings::InterpString; use fabro_types::settings::run::RunGoal; -use fabro_types::{EventBody, WorkflowSettings, fixtures}; +use fabro_types::{EventBody, TurnId, WorkflowSettings, fixtures}; fn templated_settings() -> WorkflowSettings { let mut settings = WorkflowSettings::default(); @@ -142,3 +143,20 @@ fn run_parent_events_round_trip_parent_ids() { serde_json::from_value(unlinked_json).expect("unlinked event should deserialize"); assert_eq!(unlinked_round_trip.event_name(), "run.parent.unlinked"); } + +#[test] +fn run_session_turn_failed_defaults_code_for_old_events() { + let turn_id = TurnId::new(); + let props: RunSessionTurnFailedProps = serde_json::from_value(serde_json::json!({ + "turn_id": turn_id, + "error": "legacy failure" + })) + .expect("legacy failed props should deserialize"); + + assert_eq!(props.code, RunSessionTurnFailedCode::AgentError); + assert!(!props.retryable); + + let json = serde_json::to_value(props).expect("props should serialize"); + assert_eq!(json["code"], "agent_error"); + assert_eq!(json["retryable"], false); +} diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index 42313eda2..ed8643ef4 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -1638,6 +1638,7 @@ mod tests { }, timing: None, billing: None, + ask_fabro: fabro_types::AskFabro::default(), diff: None, pull_request: None, current_question: None, diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 378a74e29..3afce3e32 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -33,6 +33,7 @@ models/artifact-batch-upload-manifest.ts models/artifact-entry.ts models/artifact-list-response.ts models/artifacts-settings.ts +models/ask-fabro.ts models/auth-config-response.ts models/auth-me-response.ts models/auth-method.ts @@ -371,10 +372,12 @@ models/server-settings.ts models/server-slate-db-settings.ts models/server-storage-settings.ts models/server-web-settings.ts +models/session-detail.ts models/session-message.ts models/session-record.ts models/session-status.ts models/session-summary.ts +models/session-turn.ts models/slack-integration-settings.ts models/ssh-access-request.ts models/ssh-access-response.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 404b0e16d..904647177 100644 --- a/lib/packages/fabro-api-client/src/api/sessions-api.ts +++ b/lib/packages/fabro-api-client/src/api/sessions-api.ts @@ -28,8 +28,12 @@ import type { ErrorResponse } from '../models'; // @ts-ignore import type { EventEnvelope } from '../models'; // @ts-ignore +import type { PaginatedEventList } from '../models'; +// @ts-ignore import type { PaginatedSessionList } from '../models'; // @ts-ignore +import type { SessionDetail } from '../models'; +// @ts-ignore import type { SessionRecord } from '../models'; // @ts-ignore import type { SubmitTurnRequest } from '../models'; @@ -38,6 +42,51 @@ 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. + * @summary Attach to session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + attachSessionEvents: async (id: string, sinceSeq?: number, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('attachSessionEvents', 'id', id) + const localVarPath = `/api/v1/sessions/{id}/attach` + .replace(`{${"id"}}`, encodeURIComponent(String(id))); + // use dummy base URL string because the URL constructor only accepts absolute URLs. + const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL); + let baseOptions; + if (configuration) { + baseOptions = configuration.baseOptions; + } + + const localVarRequestOptions = { method: 'GET', ...baseOptions, ...options}; + const localVarHeaderParameter = {} as any; + const localVarQueryParameter = {} as any; + + // authentication SessionCookie required + + // authentication BearerAuth required + // http bearer authentication required + await setBearerAuthToObject(localVarHeaderParameter, configuration) + + if (sinceSeq !== undefined) { + localVarQueryParameter['since_seq'] = sinceSeq; + } + + localVarHeaderParameter['Accept'] = 'text/event-stream,application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, /** * Creates a read-only Ask Fabro session bound to the run. * @summary Create run session @@ -171,10 +220,13 @@ export const SessionsApiAxiosParamCreator = function (configuration?: Configurat * * @summary List run sessions * @param {string} id Unique run identifier (ULID). + * @param {number} [pageLimit] + * @param {number} [pageOffset] + * @param {ListRunSessionsOrderEnum} [order] * @param {*} [options] Override http request option. * @throws {RequiredError} */ - listRunSessions: async (id: string, options: RawAxiosRequestConfig = {}): Promise => { + listRunSessions: async (id: string, pageLimit?: number, pageOffset?: number, order?: ListRunSessionsOrderEnum, options: RawAxiosRequestConfig = {}): Promise => { // verify required parameter 'id' is not null or undefined assertParamExists('listRunSessions', 'id', id) const localVarPath = `/api/v1/runs/{id}/sessions` @@ -196,6 +248,68 @@ export const SessionsApiAxiosParamCreator = function (configuration?: Configurat // http bearer authentication required await setBearerAuthToObject(localVarHeaderParameter, configuration) + if (pageLimit !== undefined) { + localVarQueryParameter['page[limit]'] = pageLimit; + } + + if (pageOffset !== undefined) { + localVarQueryParameter['page[offset]'] = pageOffset; + } + + if (order !== undefined) { + localVarQueryParameter['order'] = order; + } + + localVarHeaderParameter['Accept'] = 'application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, + /** + * Returns run event envelopes filtered to this session\'s durable `run.session.*` events. `since_seq` uses the owning run event sequence. + * @summary List session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {number} [limit] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + listSessionEvents: async (id: string, sinceSeq?: number, limit?: number, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('listSessionEvents', 'id', id) + const localVarPath = `/api/v1/sessions/{id}/events` + .replace(`{${"id"}}`, encodeURIComponent(String(id))); + // use dummy base URL string because the URL constructor only accepts absolute URLs. + const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL); + let baseOptions; + if (configuration) { + baseOptions = configuration.baseOptions; + } + + const localVarRequestOptions = { method: 'GET', ...baseOptions, ...options}; + const localVarHeaderParameter = {} as any; + const localVarQueryParameter = {} as any; + + // authentication SessionCookie required + + // authentication BearerAuth required + // http bearer authentication required + await setBearerAuthToObject(localVarHeaderParameter, configuration) + + if (sinceSeq !== undefined) { + localVarQueryParameter['since_seq'] = sinceSeq; + } + + if (limit !== undefined) { + localVarQueryParameter['limit'] = limit; + } + localVarHeaderParameter['Accept'] = 'application/json'; setSearchParams(localVarUrlObj, localVarQueryParameter); @@ -261,6 +375,20 @@ export const SessionsApiAxiosParamCreator = function (configuration?: Configurat 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. + * @summary Attach to session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async attachSessionEvents(id: string, sinceSeq?: number, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.attachSessionEvents(id, sinceSeq, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['SessionsApi.attachSessionEvents']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, /** * Creates a read-only Ask Fabro session bound to the run. * @summary Create run session @@ -282,7 +410,7 @@ export const SessionsApiFp = function(configuration?: Configuration) { * @param {*} [options] Override http request option. * @throws {RequiredError} */ - async getSession(id: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + async getSession(id: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { const localVarAxiosArgs = await localVarAxiosParamCreator.getSession(id, options); const localVarOperationServerIndex = configuration?.serverIndex ?? 0; const localVarOperationServerBasePath = operationServerMap['SessionsApi.getSession']?.[localVarOperationServerIndex]?.url; @@ -306,15 +434,33 @@ export const SessionsApiFp = function(configuration?: Configuration) { * * @summary List run sessions * @param {string} id Unique run identifier (ULID). + * @param {number} [pageLimit] + * @param {number} [pageOffset] + * @param {ListRunSessionsOrderEnum} [order] * @param {*} [options] Override http request option. * @throws {RequiredError} */ - async listRunSessions(id: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { - const localVarAxiosArgs = await localVarAxiosParamCreator.listRunSessions(id, options); + async listRunSessions(id: string, pageLimit?: number, pageOffset?: number, order?: ListRunSessionsOrderEnum, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.listRunSessions(id, pageLimit, pageOffset, order, options); const localVarOperationServerIndex = configuration?.serverIndex ?? 0; const localVarOperationServerBasePath = operationServerMap['SessionsApi.listRunSessions']?.[localVarOperationServerIndex]?.url; 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. + * @summary List session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {number} [limit] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + 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. * @summary Submit a session turn @@ -338,6 +484,17 @@ export const SessionsApiFp = function(configuration?: Configuration) { export const SessionsApiFactory = function (configuration?: Configuration, basePath?: string, axios?: AxiosInstance) { 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. + * @summary Attach to session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + attachSessionEvents(id: string, sinceSeq?: number, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.attachSessionEvents(id, sinceSeq, options).then((request) => request(axios, basePath)); + }, /** * Creates a read-only Ask Fabro session bound to the run. * @summary Create run session @@ -356,7 +513,7 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP * @param {*} [options] Override http request option. * @throws {RequiredError} */ - getSession(id: string, options?: RawAxiosRequestConfig): AxiosPromise { + getSession(id: string, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.getSession(id, options).then((request) => request(axios, basePath)); }, /** @@ -374,11 +531,26 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP * * @summary List run sessions * @param {string} id Unique run identifier (ULID). + * @param {number} [pageLimit] + * @param {number} [pageOffset] + * @param {ListRunSessionsOrderEnum} [order] * @param {*} [options] Override http request option. * @throws {RequiredError} */ - listRunSessions(id: string, options?: RawAxiosRequestConfig): AxiosPromise { - return localVarFp.listRunSessions(id, options).then((request) => request(axios, basePath)); + listRunSessions(id: string, pageLimit?: number, pageOffset?: number, order?: ListRunSessionsOrderEnum, options?: RawAxiosRequestConfig): AxiosPromise { + 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. + * @summary List session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {number} [limit] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + 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. @@ -398,6 +570,18 @@ export const SessionsApiFactory = function (configuration?: Configuration, baseP * SessionsApi - object-oriented interface */ 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. + * @summary Attach to session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public attachSessionEvents(id: string, sinceSeq?: number, options?: RawAxiosRequestConfig) { + return SessionsApiFp(this.configuration).attachSessionEvents(id, sinceSeq, options).then((request) => request(this.axios, this.basePath)); + } + /** * Creates a read-only Ask Fabro session bound to the run. * @summary Create run session @@ -437,11 +621,27 @@ export class SessionsApi extends BaseAPI { * * @summary List run sessions * @param {string} id Unique run identifier (ULID). + * @param {number} [pageLimit] + * @param {number} [pageOffset] + * @param {ListRunSessionsOrderEnum} [order] * @param {*} [options] Override http request option. * @throws {RequiredError} */ - public listRunSessions(id: string, options?: RawAxiosRequestConfig) { - return SessionsApiFp(this.configuration).listRunSessions(id, options).then((request) => request(this.axios, this.basePath)); + public listRunSessions(id: string, pageLimit?: number, pageOffset?: number, order?: ListRunSessionsOrderEnum, options?: RawAxiosRequestConfig) { + return SessionsApiFp(this.configuration).listRunSessions(id, pageLimit, pageOffset, order, options).then((request) => request(this.axios, this.basePath)); + } + + /** + * Returns run event envelopes filtered to this session\'s durable `run.session.*` events. `since_seq` uses the owning run event sequence. + * @summary List session events + * @param {string} id + * @param {number} [sinceSeq] + * @param {number} [limit] + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public listSessionEvents(id: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig) { + return SessionsApiFp(this.configuration).listSessionEvents(id, sinceSeq, limit, options).then((request) => request(this.axios, this.basePath)); } /** @@ -456,3 +656,9 @@ export class SessionsApi extends BaseAPI { return SessionsApiFp(this.configuration).submitSessionTurn(id, submitTurnRequest, options).then((request) => request(this.axios, this.basePath)); } } + +export const ListRunSessionsOrderEnum = { + UPDATED_DESC: 'updated_desc', + CREATED_DESC: 'created_desc' +} as const; +export type ListRunSessionsOrderEnum = typeof ListRunSessionsOrderEnum[keyof typeof ListRunSessionsOrderEnum]; diff --git a/lib/packages/fabro-api-client/src/models/ask-fabro.ts b/lib/packages/fabro-api-client/src/models/ask-fabro.ts new file mode 100644 index 000000000..cc1b1570d --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/ask-fabro.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.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * Readiness and defaults for starting an Ask Fabro session on this run. + */ +export interface AskFabro { + 'available': boolean; + 'unavailable_reason': AskFabroUnavailableReasonEnum | null; + 'default_model': string | null; +} + +export const AskFabroUnavailableReasonEnum = { + FEATURE_DISABLED: 'feature_disabled', + NO_SANDBOX: 'no_sandbox', + SANDBOX_NOT_READY: 'sandbox_not_ready', + LLM_UNCONFIGURED: 'llm_unconfigured' +} as const; + +export type AskFabroUnavailableReasonEnum = typeof AskFabroUnavailableReasonEnum[keyof typeof AskFabroUnavailableReasonEnum]; diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 9c5cd2123..6bf39c541 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -10,6 +10,7 @@ export * from './artifact-batch-upload-manifest'; export * from './artifact-entry'; export * from './artifact-list-response'; export * from './artifacts-settings'; +export * from './ask-fabro'; export * from './auth-config-response'; export * from './auth-me-response'; export * from './auth-method'; @@ -347,10 +348,12 @@ export * from './server-settings'; export * from './server-slate-db-settings'; export * from './server-storage-settings'; export * from './server-web-settings'; +export * from './session-detail'; export * from './session-message'; export * from './session-record'; export * from './session-status'; export * from './session-summary'; +export * from './session-turn'; export * from './slack-integration-settings'; export * from './ssh-access-request'; export * from './ssh-access-response'; diff --git a/lib/packages/fabro-api-client/src/models/run.ts b/lib/packages/fabro-api-client/src/models/run.ts index ee2b63081..cf6c4c51e 100644 --- a/lib/packages/fabro-api-client/src/models/run.ts +++ b/lib/packages/fabro-api-client/src/models/run.ts @@ -13,6 +13,9 @@ */ +// May contain unused imports in some cases +// @ts-ignore +import type { AskFabro } from './ask-fabro'; // May contain unused imports in some cases // @ts-ignore import type { AutomationRef } from './automation-ref'; @@ -87,6 +90,7 @@ export interface Run { 'timestamps': RunTimestamps; 'timing': RunTiming | null; 'billing': RunBillingSummary | null; + 'ask_fabro': AskFabro; 'diff': DiffSummary | null; 'pull_request': PullRequestLink | null; 'current_question': RunQuestion | null; diff --git a/lib/packages/fabro-api-client/src/models/session-detail.ts b/lib/packages/fabro-api-client/src/models/session-detail.ts new file mode 100644 index 000000000..d2f072227 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/session-detail.ts @@ -0,0 +1,43 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.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 { SessionMessage } from './session-message'; +// May contain unused imports in some cases +// @ts-ignore +import type { SessionStatus } from './session-status'; +// May contain unused imports in some cases +// @ts-ignore +import type { SessionTurn } from './session-turn'; + +/** + * Session metadata plus durable transcript projection. + */ +export interface SessionDetail { + /** + * Durable session identifier. + */ + 'id': string; + 'run_id': string; + 'title'?: string | null; + 'status': SessionStatus; + 'model'?: string | null; + 'active_turn': SessionTurn | null; + 'created_at': string; + 'updated_at': string; + 'messages': Array; + 'last_seq': number; +} diff --git a/lib/packages/fabro-api-client/src/models/session-record.ts b/lib/packages/fabro-api-client/src/models/session-record.ts index 7c89e9047..94aa77aa7 100644 --- a/lib/packages/fabro-api-client/src/models/session-record.ts +++ b/lib/packages/fabro-api-client/src/models/session-record.ts @@ -16,6 +16,9 @@ // May contain unused imports in some cases // @ts-ignore import type { SessionStatus } from './session-status'; +// May contain unused imports in some cases +// @ts-ignore +import type { SessionTurn } from './session-turn'; /** * Ask Fabro session metadata derived from the owning run event stream. @@ -29,6 +32,7 @@ export interface SessionRecord { 'title'?: string | null; 'status': SessionStatus; 'model'?: string | null; + 'active_turn': SessionTurn | null; 'created_at': string; 'updated_at': string; } diff --git a/lib/packages/fabro-api-client/src/models/session-summary.ts b/lib/packages/fabro-api-client/src/models/session-summary.ts index 207c6d663..4c5e3c75e 100644 --- a/lib/packages/fabro-api-client/src/models/session-summary.ts +++ b/lib/packages/fabro-api-client/src/models/session-summary.ts @@ -16,6 +16,9 @@ // May contain unused imports in some cases // @ts-ignore import type { SessionStatus } from './session-status'; +// May contain unused imports in some cases +// @ts-ignore +import type { SessionTurn } from './session-turn'; /** * List projection of an Ask Fabro session. @@ -29,6 +32,7 @@ export interface SessionSummary { 'title'?: string | null; 'status': SessionStatus; 'model'?: string | null; + 'active_turn': SessionTurn | null; 'created_at': string; 'updated_at': string; } diff --git a/lib/packages/fabro-api-client/src/models/session-turn.ts b/lib/packages/fabro-api-client/src/models/session-turn.ts new file mode 100644 index 000000000..65950d477 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/session-turn.ts @@ -0,0 +1,27 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * Currently active durable session turn. + */ +export interface SessionTurn { + /** + * Durable session turn identifier. + */ + 'id': string; + 'started_at': string; + 'input': string; +} diff --git a/lib/packages/fabro-api-client/src/models/submit-turn-request.ts b/lib/packages/fabro-api-client/src/models/submit-turn-request.ts index ff6e824b0..06375abf8 100644 --- a/lib/packages/fabro-api-client/src/models/submit-turn-request.ts +++ b/lib/packages/fabro-api-client/src/models/submit-turn-request.ts @@ -16,4 +16,8 @@ export interface SubmitTurnRequest { 'input': string; + /** + * Durable session turn identifier. + */ + 'turn_id'?: string; }