diff --git a/.changeset/session-inbound-events.md b/.changeset/session-inbound-events.md new file mode 100644 index 000000000..de1b7ec2d --- /dev/null +++ b/.changeset/session-inbound-events.md @@ -0,0 +1,6 @@ +--- +"@truefoundry/trueforge-core": patch +"@truefoundry/trueforge": patch +--- + +Add `session_inbound_events` store API for durable tip HITL send-event inbox (insert / list unconsumed / mark consumed), with Postgres and SQLite migrations. diff --git a/packages/trueforge-core/src/agent-session/index.ts b/packages/trueforge-core/src/agent-session/index.ts index 3b9b5d66c..d8b492690 100644 --- a/packages/trueforge-core/src/agent-session/index.ts +++ b/packages/trueforge-core/src/agent-session/index.ts @@ -21,6 +21,9 @@ export { } from './schemas/turn'; export type { TerminalTurnState, Turn, TurnInputItem, TurnMetrics, TurnState } from './schemas/turn'; +export { SessionInboundEventItemSchema } from './schemas/sendEvent'; +export type { SessionInboundEventItem } from './schemas/sendEvent'; + export { SessionMetadataSchema, SessionMetricsSchema, @@ -80,16 +83,20 @@ export type { GetSessionInput, GetTurnInput, ISessionStore, + InsertSessionInboundEventsInput, ListSessionEventsInput, ListSessionsInput, ListTurnEventsInput, ListTurnsInput, + ListUnconsumedSessionInboundEventsInput, + MarkSessionInboundEventsConsumedInput, NewThreadInit, OverwriteThreadContextInput, PatchMCPServersInput, PatchSandboxInfoInput, PatchThreadCapabilityStateInput, RemoveThreadsInput, + SessionInboundEventRecord, TurnContextAppend, TurnRecordWithoutSnapshot, UpdateSessionInput, @@ -100,6 +107,7 @@ export { PreviousTurnRunningError, SessionAlreadyExistsError, SessionExternalIdConflictError, + SessionInboundEventAlreadyExistsError, SessionNotFoundError, SessionStoreConflictError, SessionStoreInvariantError, diff --git a/packages/trueforge-core/src/agent-session/schemas/sendEvent.ts b/packages/trueforge-core/src/agent-session/schemas/sendEvent.ts new file mode 100644 index 000000000..f25a46587 --- /dev/null +++ b/packages/trueforge-core/src/agent-session/schemas/sendEvent.ts @@ -0,0 +1,17 @@ +/** + * Inbound send-event payloads for tip HITL (client → harness), distinct from the + * stream log ({@link PersistedTurnEvent} / session_event). + * + * Public send is session-scoped (`POST …/sessions/{id}/events`) with required + * body `turn_id` (one batch → one tip) plus `SessionInboundEventItem`s; rows stamp that + * tip id. v1 union is tip-only; approval policies may relax `turn_id` later. + * `user.message` stays on createTurn / steer. + */ +import { z } from '@hono/zod-openapi'; +import { UserToolApprovalMessageSchema, UserToolResponseMessageSchema } from '../../core/events/schema'; + +export const SessionInboundEventItemSchema = z + .discriminatedUnion('type', [UserToolApprovalMessageSchema, UserToolResponseMessageSchema]) + .openapi('SessionInboundEventItem'); + +export type SessionInboundEventItem = z.infer; diff --git a/packages/trueforge-core/src/agent-session/store/ISessionStore.ts b/packages/trueforge-core/src/agent-session/store/ISessionStore.ts index b11a7ca1d..69b74aa2f 100644 --- a/packages/trueforge-core/src/agent-session/store/ISessionStore.ts +++ b/packages/trueforge-core/src/agent-session/store/ISessionStore.ts @@ -11,6 +11,7 @@ import type { SessionRecord } from '../models/SessionRecord'; import type { TurnRecord } from '../models/TurnRecord'; import type { PersistedTurnEvent, SessionEventItem } from '../schemas/events'; import type { TokenPagination } from '../schemas/pagination'; +import type { SessionInboundEventItem } from '../schemas/sendEvent'; import type { SessionMetadata } from '../schemas/session'; import type { CancellationReason, TerminalTurnState } from '../schemas/turn'; @@ -170,6 +171,52 @@ export interface AppendToEventsInput { events: PersistedTurnEvent[]; } +/** One durable inbound send-event row (tip HITL and/or session-scoped). */ +export interface SessionInboundEventRecord { + event_id: string; + /** Tip id when tip-scoped; null for session-only (e.g. future policies). */ + turn_id: string | null; + /** Validated {@link SessionInboundEventItem} body (widens when policy lands). */ + payload: SessionInboundEventItem; + /** ISO-8601; copied from insert input. Ordering uses `event_id`. */ + created_at: string; +} + +export interface InsertSessionInboundEventsInput { + session_id: string; + /** + * Tip that receives this batch (v1 required). One send = one tip; stamp every + * row with this id. Relax to optional/null when session-scoped policies land. + */ + turn_id: string; + /** + * Caller mints `event_id` (monotonic ULID) — same contract as session_event. + * Empty array is a no-op. + */ + events: Array<{ + event_id: string; + payload: SessionInboundEventItem; + created_at: string; + }>; +} + +export interface ListUnconsumedSessionInboundEventsInput { + session_id: string; + /** + * Three-way filter — pass the key explicitly (do not omit): + * - `undefined` — all unconsumed for the session + * - `string` — unconsumed for that turn only + * - `null` — session-scoped rows only (`turn_id` IS NULL; empty until + * policies allow null inserts) + */ + turn_id: string | null | undefined; +} + +export interface MarkSessionInboundEventsConsumedInput { + session_id: string; + event_ids: string[]; +} + export interface AddThreadsInput { session_id: string; turn_id: string; @@ -360,6 +407,31 @@ export interface ISessionStore< */ appendToEvents(input: AppendToEventsInput): Promise; + /** + * Durable inbound send-event inbox for the session. Column stays nullable for later + * session-scoped policies. Tip must be non-terminal (v1: `running`; `paused` + * when that status lands) — terminal tip → {@link TurnNotRunningError}. + * Missing session → {@link SessionNotFoundError}; unknown turn → + * {@link TurnNotFoundError}. Duplicate `event_id` → + * {@link SessionInboundEventAlreadyExistsError}. + */ + insertSessionInboundEvents(input: InsertSessionInboundEventsInput): Promise; + + /** + * Unconsumed inbox rows, ordered by monotonic `event_id` ascending. + * See {@link ListUnconsumedSessionInboundEventsInput.turn_id} for filtering. + * Missing session → {@link SessionNotFoundError}. + */ + listUnconsumedSessionInboundEvents( + input: ListUnconsumedSessionInboundEventsInput, + ): Promise; + + /** + * Marks inbox rows consumed. Already-consumed or unknown ids are ignored. + * Empty `event_ids` is a no-op. Missing session → {@link SessionNotFoundError}. + */ + markSessionInboundEventsConsumed(input: MarkSessionInboundEventsConsumedInput): Promise; + /** Adds thread snapshots to the turn (sub-agent spawns). */ addThreads(input: AddThreadsInput): Promise; diff --git a/packages/trueforge-core/src/agent-session/store/InMemorySessionStore.ts b/packages/trueforge-core/src/agent-session/store/InMemorySessionStore.ts index a2e114b61..7684ec873 100644 --- a/packages/trueforge-core/src/agent-session/store/InMemorySessionStore.ts +++ b/packages/trueforge-core/src/agent-session/store/InMemorySessionStore.ts @@ -4,6 +4,7 @@ import type { SessionRecord } from '../models/SessionRecord'; import type { TurnRecord, TurnSnapshot } from '../models/TurnRecord'; import type { PersistedTurnEvent, SessionEventItem } from '../schemas/events'; import type { TokenPagination } from '../schemas/pagination'; +import type { SessionInboundEventItem } from '../schemas/sendEvent'; import type { TerminalTurnState } from '../schemas/turn'; import { assertCreateTurnThreadDelta } from './assertCreateTurnThreadDelta'; import type { @@ -18,17 +19,21 @@ import type { GetSessionByExternalIdInput, GetSessionInput, GetTurnInput, + InsertSessionInboundEventsInput, ISessionStore, ListSessionEventsInput, ListSessionsInput, ListTurnEventsInput, ListTurnsInput, + ListUnconsumedSessionInboundEventsInput, + MarkSessionInboundEventsConsumedInput, NewThreadInit, OverwriteThreadContextInput, PatchMCPServersInput, PatchSandboxInfoInput, PatchThreadCapabilityStateInput, RemoveThreadsInput, + SessionInboundEventRecord, TurnContextAppend, TurnRecordWithoutSnapshot, UpdateSessionInput, @@ -45,6 +50,7 @@ import { PreviousTurnRunningError, SessionAlreadyExistsError, SessionExternalIdConflictError, + SessionInboundEventAlreadyExistsError, SessionNotFoundError, SessionStoreInvariantError, TurnAlreadyExistsError, @@ -56,6 +62,14 @@ import { type StoredEvent = PersistedTurnEvent; +interface StoredInboundEvent { + event_id: string; + turn_id: string | null; + payload: SessionInboundEventItem; + created_at: string; + consumed: boolean; +} + interface StoredSession { record: SessionRecord; turnIds: string[]; @@ -170,6 +184,8 @@ export class InMemorySessionStore< private readonly sessions = new Map>(); private readonly turns = new Map>(); private readonly events = new Map(); + /** session_id → inbound send-event inbox */ + private readonly inboundEvents = new Map(); async createSession(input: CreateSessionInput): Promise { const key = sessionKey(input.session_id); @@ -219,6 +235,7 @@ export class InMemorySessionStore< this.turns.delete(tKey); this.events.delete(tKey); } + this.inboundEvents.delete(sessionKey(input.session_id)); this.sessions.delete(sKey); } @@ -485,6 +502,81 @@ export class InMemorySessionStore< return; } + async insertSessionInboundEvents(input: InsertSessionInboundEventsInput): Promise { + if (input.events.length === 0) { + return; + } + this.requireSession(input.session_id); + this.requireRunningTurn(input.session_id, input.turn_id); + const sKey = sessionKey(input.session_id); + let list = this.inboundEvents.get(sKey); + if (!list) { + list = []; + this.inboundEvents.set(sKey, list); + } + const existing = new Set(list.map(row => row.event_id)); + for (const event of input.events) { + if (existing.has(event.event_id)) { + throw new SessionInboundEventAlreadyExistsError(input.session_id, event.event_id); + } + existing.add(event.event_id); + } + for (const event of input.events) { + list.push({ + event_id: event.event_id, + turn_id: input.turn_id, + payload: deepCopy(event.payload), + created_at: event.created_at, + consumed: false, + }); + } + } + + async listUnconsumedSessionInboundEvents( + input: ListUnconsumedSessionInboundEventsInput, + ): Promise { + this.requireSession(input.session_id); + const list = this.inboundEvents.get(sessionKey(input.session_id)) ?? []; + return list + .filter(row => { + if (row.consumed) { + return false; + } + if (input.turn_id === undefined) { + return true; + } + if (input.turn_id === null) { + return row.turn_id === null; + } + return row.turn_id === input.turn_id; + }) + .slice() + .sort((a, b) => (a.event_id < b.event_id ? -1 : a.event_id > b.event_id ? 1 : 0)) + .map(row => ({ + event_id: row.event_id, + turn_id: row.turn_id, + payload: deepCopy(row.payload), + created_at: row.created_at, + })); + } + + async markSessionInboundEventsConsumed(input: MarkSessionInboundEventsConsumedInput): Promise { + if (input.event_ids.length === 0) { + return; + } + this.requireSession(input.session_id); + const list = this.inboundEvents.get(sessionKey(input.session_id)); + if (!list) { + return; + } + const wanted = new Set(input.event_ids); + for (const row of list) { + if (wanted.has(row.event_id)) { + row.consumed = true; + } + } + } + /** Cost from turn metrics when present; duration is completed_at − created_at, floored at 0. */ private addTerminalSessionMetrics(sessionId: string, created_at: Date, state: TerminalTurnState): void { const stored = this.sessions.get(sessionKey(sessionId)); @@ -499,6 +591,14 @@ export class InMemorySessionStore< stored.record.metrics.total_duration_ms += elapsed_ms > 0 ? Math.trunc(elapsed_ms) : 0; } + private requireSession(sessionId: string): StoredSession { + const stored = this.sessions.get(sessionKey(sessionId)); + if (!stored) { + throw new SessionNotFoundError(sessionId); + } + return stored; + } + private requireTurn(sessionId: string, turnId: string): TurnRecord { const turn = this.turns.get(turnKey({ session_id: sessionId, turn_id: turnId })); if (!turn) { diff --git a/packages/trueforge-core/src/agent-session/store/SessionStoreErrors.ts b/packages/trueforge-core/src/agent-session/store/SessionStoreErrors.ts index 89b31ee90..b3829cac6 100644 --- a/packages/trueforge-core/src/agent-session/store/SessionStoreErrors.ts +++ b/packages/trueforge-core/src/agent-session/store/SessionStoreErrors.ts @@ -74,6 +74,18 @@ export class TurnAlreadyExistsError extends SessionStoreConflictError { } } +export class SessionInboundEventAlreadyExistsError extends SessionStoreConflictError { + readonly session_id: string; + readonly event_id: string; + + constructor(session_id: string, event_id: string, options?: ErrorOptions) { + super(`Session inbound event already exists: ${session_id}/${event_id}`, options); + this.name = 'SessionInboundEventAlreadyExistsError'; + this.session_id = session_id; + this.event_id = event_id; + } +} + export class PreviousTurnRunningError extends SessionStoreConflictError { readonly previous_turn_id: string; diff --git a/packages/trueforge-core/tests/agent-session/store/storeContractSuite.ts b/packages/trueforge-core/tests/agent-session/store/storeContractSuite.ts index 28a85e43f..1cc25afc5 100644 --- a/packages/trueforge-core/tests/agent-session/store/storeContractSuite.ts +++ b/packages/trueforge-core/tests/agent-session/store/storeContractSuite.ts @@ -687,6 +687,27 @@ export function runStoreContractSuite(createStore: () => ISessionStore) { order: undefined, }), ).rejects.toBeInstanceOf(TurnNotFoundError); + await expect( + store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + events: [ + { + event_id: newEventId(), + payload: { + type: 'user.tool_approval', + thread_id: 'main', + tool_call_id: 'tc-1', + approval: { status: 'allow' }, + }, + created_at: new Date().toISOString(), + }, + ], + }), + ).rejects.toBeInstanceOf(SessionNotFoundError); + await expect( + store.listUnconsumedSessionInboundEvents({ session_id: sessionId, turn_id: undefined }), + ).rejects.toBeInstanceOf(SessionNotFoundError); await expect( store.listSessionEvents({ session_id: sessionId, @@ -2208,6 +2229,269 @@ export function runStoreContractSuite(createStore: () => ISessionStore) { expect(data.map(e => e.id)).toEqual([created.id, model.id]); }); + it('session_inbound_events: insert, list unconsumed, mark consumed, duplicate id', async () => { + const store = createStore(); + await seedSession(store); + await store.createTurn(makeCreateTurnInput({ sessionId, turnId: 'turn-1' })); + + const earlier = { + event_id: 'evt-a', + payload: { + type: 'user.tool_approval' as const, + thread_id: 'main', + tool_call_id: 'tc-1', + approval: { status: 'allow' as const }, + }, + created_at: new Date().toISOString(), + }; + const later = { + event_id: 'evt-b', + payload: { + type: 'user.tool_approval' as const, + thread_id: 'main', + tool_call_id: 'tc-2', + approval: { status: 'deny' as const, reason: 'nope' }, + }, + created_at: new Date().toISOString(), + }; + + await store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + events: [later, earlier], + }); + + let pending = await store.listUnconsumedSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + }); + expect(pending.map(e => e.event_id)).toEqual([earlier.event_id, later.event_id]); + expect(pending[0]?.payload).toEqual(earlier.payload); + expect(pending[0]?.turn_id).toBe('turn-1'); + + await store.markSessionInboundEventsConsumed({ + session_id: sessionId, + event_ids: [earlier.event_id], + }); + pending = await store.listUnconsumedSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + }); + expect(pending.map(e => e.event_id)).toEqual([later.event_id]); + + await expect( + store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + events: [later], + }), + ).rejects.toMatchObject({ + name: 'SessionInboundEventAlreadyExistsError', + event_id: later.event_id, + }); + + // Later id in the batch collides — error must name that id. + const fresh = { + event_id: 'evt-fresh', + payload: { + type: 'user.tool_approval' as const, + thread_id: 'main', + tool_call_id: 'tc-fresh', + approval: { status: 'allow' as const }, + }, + created_at: new Date().toISOString(), + }; + await expect( + store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + events: [fresh, later], + }), + ).rejects.toMatchObject({ + name: 'SessionInboundEventAlreadyExistsError', + event_id: later.event_id, + }); + expect( + (await store.listUnconsumedSessionInboundEvents({ session_id: sessionId, turn_id: 'turn-1' })).map( + e => e.event_id, + ), + ).toEqual([later.event_id]); + + const dupId = 'evt-dup'; + await expect( + store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + events: [ + { + event_id: dupId, + payload: { + type: 'user.tool_approval' as const, + thread_id: 'main', + tool_call_id: 'tc-dup', + approval: { status: 'allow' as const }, + }, + created_at: new Date().toISOString(), + }, + { + event_id: dupId, + payload: { + type: 'user.tool_approval' as const, + thread_id: 'main', + tool_call_id: 'tc-dup-2', + approval: { status: 'deny' as const, reason: 'dup' }, + }, + created_at: new Date().toISOString(), + }, + ], + }), + ).rejects.toMatchObject({ + name: 'SessionInboundEventAlreadyExistsError', + event_id: dupId, + }); + // Failed batch must not leave a partial row (SQL PK is all-or-nothing). + expect( + (await store.listUnconsumedSessionInboundEvents({ session_id: sessionId, turn_id: 'turn-1' })).map( + e => e.event_id, + ), + ).toEqual([later.event_id]); + + // Terminal tip rejects inbox writes. + await finishTurn(store, 'turn-1'); + await expect( + store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + events: [ + { + event_id: 'evt-after-done', + payload: { + type: 'user.tool_response' as const, + thread_id: 'main', + tool_call_id: 'tc-3', + content: 'client result', + }, + created_at: new Date().toISOString(), + }, + ], + }), + ).rejects.toBeInstanceOf(TurnNotRunningError); + expect( + (await store.listUnconsumedSessionInboundEvents({ session_id: sessionId, turn_id: 'turn-1' })).map( + e => e.event_id, + ), + ).toEqual([later.event_id]); + }); + + it('session_inbound_events: list filter turn_id string | null | undefined', async () => { + const store = createStore(); + await seedSession(store); + await store.createTurn(makeCreateTurnInput({ sessionId, turnId: 'turn-a' })); + + const forA = { + event_id: 'evt-a', + payload: { + type: 'user.tool_approval' as const, + thread_id: 'main', + tool_call_id: 'tc-a', + approval: { status: 'allow' as const }, + }, + created_at: new Date().toISOString(), + }; + await store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-a', + events: [forA], + }); + + await finishTurn(store, 'turn-a'); + await store.createTurn( + makeCreateTurnInput({ sessionId, turnId: 'turn-b', previousTurnId: 'turn-a', firstTurnId: 'turn-a' }), + ); + + const forB = { + event_id: 'evt-b', + payload: { + type: 'user.tool_approval' as const, + thread_id: 'main', + tool_call_id: 'tc-b', + approval: { status: 'allow' as const }, + }, + created_at: new Date().toISOString(), + }; + + await store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-b', + events: [forB], + }); + + // string — that turn only + expect( + ( + await store.listUnconsumedSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-a', + }) + ).map(e => e.event_id), + ).toEqual([forA.event_id]); + expect( + ( + await store.listUnconsumedSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-b', + }) + ).map(e => e.event_id), + ).toEqual([forB.event_id]); + + // null — session-scoped only (v1 insert always sets turn_id; no such rows yet) + expect( + ( + await store.listUnconsumedSessionInboundEvents({ + session_id: sessionId, + turn_id: null, + }) + ).map(e => e.event_id), + ).toEqual([]); + + // omitted — all unconsumed, ordered by event_id + // undefined — all unconsumed, ordered by event_id + expect( + ( + await store.listUnconsumedSessionInboundEvents({ + session_id: sessionId, + turn_id: undefined, + }) + ).map(e => e.event_id), + ).toEqual([forA.event_id, forB.event_id]); + }); + + it('session_inbound_events cascade away with deleteSession', async () => { + const store = createStore(); + await seedSession(store); + await store.createTurn(makeCreateTurnInput({ sessionId, turnId: 'turn-1' })); + await store.insertSessionInboundEvents({ + session_id: sessionId, + turn_id: 'turn-1', + events: [ + { + event_id: newEventId(), + payload: { + type: 'user.tool_approval', + thread_id: 'main', + tool_call_id: 'tc-x', + approval: { status: 'allow' }, + }, + created_at: new Date().toISOString(), + }, + ], + }); + await store.deleteSession({ tenant_id: tenant, session_id: sessionId }); + await expect( + store.listUnconsumedSessionInboundEvents({ session_id: sessionId, turn_id: undefined }), + ).rejects.toBeInstanceOf(SessionNotFoundError); + }); + it('add/remove threads and append/overwrite context', async () => { const store = createStore(); await seedSession(store); diff --git a/packages/trueforge/src/db/postgres/migrations/20260918_000002_session_inbound_events.ts b/packages/trueforge/src/db/postgres/migrations/20260918_000002_session_inbound_events.ts new file mode 100644 index 000000000..1bf4f1833 --- /dev/null +++ b/packages/trueforge/src/db/postgres/migrations/20260918_000002_session_inbound_events.ts @@ -0,0 +1,37 @@ +import { sql, type Kysely } from 'kysely'; + +/** + * Session inbound send-event inbox (tip HITL + future session-scoped payloads). + * `turn_id` column nullable (v1 insert always sets it; null reserved for session-only policies later). + */ +export async function up(db: Kysely): Promise { + await sql`SET LOCAL lock_timeout = '5s'`.execute(db); + + await db.schema + .createTable('session_inbound_events') + .addColumn('session_id', 'text', col => col.notNull()) + .addColumn('event_id', 'text', col => col.notNull()) + .addColumn('turn_id', 'text') + .addColumn('payload', 'jsonb', col => col.notNull()) + .addColumn('consumed', 'boolean', col => col.notNull().defaultTo(false)) + .addColumn('created_at', 'timestamptz', col => col.notNull()) + .addPrimaryKeyConstraint('session_inbound_events_pkey', ['session_id', 'event_id']) + .execute(); + + await db.schema + .alterTable('session_inbound_events') + .addForeignKeyConstraint('session_inbound_events_session_fkey', ['session_id'], 'session', ['session_id']) + .onDelete('cascade') + .execute(); + + await sql` + CREATE INDEX session_inbound_events_unconsumed_idx + ON session_inbound_events (session_id, turn_id, event_id) + WHERE consumed = false + `.execute(db); +} + +export async function down(db: Kysely): Promise { + await sql`SET LOCAL lock_timeout = '5s'`.execute(db); + await db.schema.dropTable('session_inbound_events').ifExists().cascade().execute(); +} diff --git a/packages/trueforge/src/db/postgres/session-store/PostgresSessionStore.ts b/packages/trueforge/src/db/postgres/session-store/PostgresSessionStore.ts index 76b6ce66b..8a254caa1 100644 --- a/packages/trueforge/src/db/postgres/session-store/PostgresSessionStore.ts +++ b/packages/trueforge/src/db/postgres/session-store/PostgresSessionStore.ts @@ -21,16 +21,20 @@ import type { GetSessionByExternalIdInput, GetSessionInput, GetTurnInput, + InsertSessionInboundEventsInput, ISessionStore, ListSessionEventsInput, ListSessionsInput, ListTurnEventsInput, ListTurnsInput, + ListUnconsumedSessionInboundEventsInput, + MarkSessionInboundEventsConsumedInput, OverwriteThreadContextInput, PatchMCPServersInput, PatchSandboxInfoInput, PatchThreadCapabilityStateInput, RemoveThreadsInput, + SessionInboundEventRecord, TurnRecordWithoutSnapshot, UpdateSessionInput, UpdateTurnStateInput, @@ -52,6 +56,11 @@ import { listSessionEvents as listSessionEventsQuery, listTurnEvents as listTurnEventsQuery, } from './queries/events'; +import { + insertSessionInboundEvents as insertSessionInboundEventsQuery, + listUnconsumedSessionInboundEvents as listUnconsumedSessionInboundEventsQuery, + markSessionInboundEventsConsumed as markSessionInboundEventsConsumedQuery, +} from './queries/inboundEvents'; import { createSession as createSessionQuery, deleteSession as deleteSessionQuery, @@ -216,6 +225,20 @@ export class PostgresSessionStore implements ISessionStore { + return insertSessionInboundEventsQuery(this.db, input); + } + + listUnconsumedSessionInboundEvents( + input: ListUnconsumedSessionInboundEventsInput, + ): Promise { + return listUnconsumedSessionInboundEventsQuery(this.db, input); + } + + markSessionInboundEventsConsumed(input: MarkSessionInboundEventsConsumedInput): Promise { + return markSessionInboundEventsConsumedQuery(this.db, input); + } + addThreads(input: AddThreadsInput): Promise { return addThreadsQuery(this.db, input); } diff --git a/packages/trueforge/src/db/postgres/session-store/queries/inboundEvents.ts b/packages/trueforge/src/db/postgres/session-store/queries/inboundEvents.ts new file mode 100644 index 000000000..da3a6f1a3 --- /dev/null +++ b/packages/trueforge/src/db/postgres/session-store/queries/inboundEvents.ts @@ -0,0 +1,148 @@ +import type { + InsertSessionInboundEventsInput, + ListUnconsumedSessionInboundEventsInput, + MarkSessionInboundEventsConsumedInput, + SessionInboundEventRecord, +} from '@truefoundry/trueforge-core/agent-session/store/ISessionStore'; +import { + SessionInboundEventAlreadyExistsError, + SessionNotFoundError, + TurnNotFoundError, + TurnNotRunningError, +} from '@truefoundry/trueforge-core/agent-session/store/SessionStoreErrors'; +import type { Kysely } from 'kysely'; +import { sql } from 'kysely'; +import { firstCollidingEventId, firstDuplicateEventIdInBatch } from '../../../sessionInboundEvents'; +import { isUniqueViolation } from '../../client'; +import { json } from '../../sqlExpressions'; +import type { Database } from '../../types'; + +async function requireSession(db: Kysely, sessionId: string): Promise { + const row = await db + .selectFrom('session') + .select('session_id') + .where('session_id', '=', sessionId) + .executeTakeFirst(); + if (!row) { + throw new SessionNotFoundError(sessionId); + } +} + +/** Exists + non-terminal (v1: `running` only; `paused` will be allowed when that status lands). */ +async function requireTurn(db: Kysely, sessionId: string, turnId: string): Promise { + const row = await db + .selectFrom('turn') + .select(['turn_id', 'state']) + .where('session_id', '=', sessionId) + .where('turn_id', '=', turnId) + .executeTakeFirst(); + if (!row) { + throw new TurnNotFoundError(turnId); + } + if (row.state.status !== 'running') { + throw new TurnNotRunningError(turnId, row.state); + } +} + +async function resolveCollidingEventId( + db: Kysely, + sessionId: string, + events: InsertSessionInboundEventsInput['events'], +): Promise { + const ids = [...new Set(events.map(e => e.event_id))]; + if (ids.length === 0) { + return ''; + } + const rows = await db + .selectFrom('session_inbound_events') + .select('event_id') + .where('session_id', '=', sessionId) + .where('event_id', 'in', ids) + .execute(); + return firstCollidingEventId(events, new Set(rows.map(r => r.event_id))); +} + +export async function insertSessionInboundEvents( + db: Kysely, + input: InsertSessionInboundEventsInput, +): Promise { + if (input.events.length === 0) { + return; + } + await requireSession(db, input.session_id); + await requireTurn(db, input.session_id, input.turn_id); + + const duplicateInBatch = firstDuplicateEventIdInBatch(input.events); + if (duplicateInBatch !== undefined) { + throw new SessionInboundEventAlreadyExistsError(input.session_id, duplicateInBatch); + } + + try { + await db + .insertInto('session_inbound_events') + .values( + input.events.map(event => ({ + session_id: input.session_id, + event_id: event.event_id, + turn_id: input.turn_id, + payload: json(event.payload), + consumed: false, + created_at: sql`${event.created_at}::timestamptz`, + })), + ) + .execute(); + } catch (error) { + if (isUniqueViolation(error)) { + const eventId = await resolveCollidingEventId(db, input.session_id, input.events); + throw new SessionInboundEventAlreadyExistsError(input.session_id, eventId, { + cause: error, + }); + } + throw error; + } +} + +export async function listUnconsumedSessionInboundEvents( + db: Kysely, + input: ListUnconsumedSessionInboundEventsInput, +): Promise { + await requireSession(db, input.session_id); + + let query = db + .selectFrom('session_inbound_events') + .select(['event_id', 'turn_id', 'payload', 'created_at']) + .where('session_id', '=', input.session_id) + .where('consumed', '=', false); + + if (input.turn_id === null) { + query = query.where('turn_id', 'is', null); + } else if (input.turn_id !== undefined) { + query = query.where('turn_id', '=', input.turn_id); + } + + const rows = await query.orderBy('event_id', 'asc').execute(); + + return rows.map(row => ({ + event_id: row.event_id, + turn_id: row.turn_id, + payload: row.payload, + created_at: new Date(row.created_at).toISOString(), + })); +} + +export async function markSessionInboundEventsConsumed( + db: Kysely, + input: MarkSessionInboundEventsConsumedInput, +): Promise { + if (input.event_ids.length === 0) { + return; + } + await requireSession(db, input.session_id); + + await db + .updateTable('session_inbound_events') + .set({ consumed: true }) + .where('session_id', '=', input.session_id) + .where('event_id', 'in', input.event_ids) + .execute(); +} diff --git a/packages/trueforge/src/db/postgres/types.ts b/packages/trueforge/src/db/postgres/types.ts index 41d78f462..342f6f376 100644 --- a/packages/trueforge/src/db/postgres/types.ts +++ b/packages/trueforge/src/db/postgres/types.ts @@ -6,6 +6,7 @@ import type { AgentSpec, CreatedBySubject, PersistedTurnEvent, + SessionInboundEventItem, SessionMetadata, SessionMetrics, SessionSource, @@ -249,6 +250,19 @@ export interface SessionEventTable { created_at: Date; } +/** + * Session inbound send-event inbox (tip HITL + future session-scoped payloads). + * PRIMARY KEY (session_id, event_id). `turn_id` nullable. + */ +export interface SessionInboundEventsTable { + session_id: string; + event_id: string; + turn_id: string | null; + payload: JSONColumnType; + consumed: boolean; + created_at: Date; +} + /** * pure immutable CONTENT; no state → no checkpoint field * PRIMARY KEY (session_id, thread_id, append_id) @@ -520,6 +534,7 @@ export interface Database { turn: TurnTable; turn_thread: TurnThreadTable; session_event: SessionEventTable; + session_inbound_events: SessionInboundEventsTable; thread_context_log: ThreadContextLogTable; thread_capability_state: ThreadCapabilityStateTable; model_provider: ModelProviderTable; diff --git a/packages/trueforge/src/db/sessionInboundEvents.ts b/packages/trueforge/src/db/sessionInboundEvents.ts new file mode 100644 index 000000000..2e598acf5 --- /dev/null +++ b/packages/trueforge/src/db/sessionInboundEvents.ts @@ -0,0 +1,28 @@ +import type { InsertSessionInboundEventsInput } from '@truefoundry/trueforge-core/agent-session/store/ISessionStore'; + +type InboundInsertEvent = InsertSessionInboundEventsInput['events'][number]; + +/** First repeated `event_id` in the batch (input order), if any. */ +export function firstDuplicateEventIdInBatch( + events: readonly Pick[], +): string | undefined { + const seen = new Set(); + for (const event of events) { + if (seen.has(event.event_id)) { + return event.event_id; + } + seen.add(event.event_id); + } + return undefined; +} + +/** + * After a unique/PK violation, pick the colliding id: first input `event_id` + * that already exists. + */ +export function firstCollidingEventId( + events: readonly Pick[], + existingEventIds: ReadonlySet, +): string { + return events.find(e => existingEventIds.has(e.event_id))?.event_id ?? events[0]?.event_id ?? ''; +} diff --git a/packages/trueforge/src/db/sqlite/client.ts b/packages/trueforge/src/db/sqlite/client.ts index 8c29149aa..ce6dc9407 100644 --- a/packages/trueforge/src/db/sqlite/client.ts +++ b/packages/trueforge/src/db/sqlite/client.ts @@ -141,6 +141,7 @@ const JSON_RESULT_COLUMNS = new Set([ 'turn_state', 'thread_checkpoint', 'event', + 'payload', 'manifest', 'metadata', 'build_metadata', diff --git a/packages/trueforge/src/db/sqlite/migrations/20260918_000002_session_inbound_events.ts b/packages/trueforge/src/db/sqlite/migrations/20260918_000002_session_inbound_events.ts new file mode 100644 index 000000000..e7fc8ddcc --- /dev/null +++ b/packages/trueforge/src/db/sqlite/migrations/20260918_000002_session_inbound_events.ts @@ -0,0 +1,30 @@ +import { type Kysely, sql } from 'kysely'; + +/** + * Session inbound send-event inbox (tip HITL + future session-scoped payloads). + * `turn_id` column nullable (v1 insert always sets it; null reserved for session-only policies later). + */ +export async function up(db: Kysely): Promise { + await sql` + CREATE TABLE session_inbound_events ( + session_id TEXT NOT NULL REFERENCES session(session_id) ON DELETE CASCADE, + event_id TEXT NOT NULL, + turn_id TEXT, + payload BLOB NOT NULL, + consumed INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL, + PRIMARY KEY (session_id, event_id) + ) STRICT + `.execute(db); + + await sql` + CREATE INDEX session_inbound_events_unconsumed_idx + ON session_inbound_events (session_id, turn_id, event_id) + WHERE consumed = 0 + `.execute(db); +} + +export async function down(db: Kysely): Promise { + await sql`DROP INDEX IF EXISTS session_inbound_events_unconsumed_idx`.execute(db); + await sql`DROP TABLE IF EXISTS session_inbound_events`.execute(db); +} diff --git a/packages/trueforge/src/db/sqlite/session-store/SqliteSessionStore.ts b/packages/trueforge/src/db/sqlite/session-store/SqliteSessionStore.ts index f0c597a68..e4efa84b3 100644 --- a/packages/trueforge/src/db/sqlite/session-store/SqliteSessionStore.ts +++ b/packages/trueforge/src/db/sqlite/session-store/SqliteSessionStore.ts @@ -14,16 +14,20 @@ import type { GetSessionByExternalIdInput, GetSessionInput, GetTurnInput, + InsertSessionInboundEventsInput, ISessionStore, ListSessionEventsInput, ListSessionsInput, ListTurnEventsInput, ListTurnsInput, + ListUnconsumedSessionInboundEventsInput, + MarkSessionInboundEventsConsumedInput, OverwriteThreadContextInput, PatchMCPServersInput, PatchSandboxInfoInput, PatchThreadCapabilityStateInput, RemoveThreadsInput, + SessionInboundEventRecord, TurnRecordWithoutSnapshot, UpdateSessionInput, UpdateTurnStateInput, @@ -40,6 +44,11 @@ import { listSessionEvents as listSessionEventsQuery, listTurnEvents as listTurnEventsQuery, } from './queries/events'; +import { + insertSessionInboundEvents as insertSessionInboundEventsQuery, + listUnconsumedSessionInboundEvents as listUnconsumedSessionInboundEventsQuery, + markSessionInboundEventsConsumed as markSessionInboundEventsConsumedQuery, +} from './queries/inboundEvents'; import { createSession as createSessionQuery, deleteSession as deleteSessionQuery, @@ -191,6 +200,20 @@ export class SqliteSessionStore implements ISessionStore { + return insertSessionInboundEventsQuery(this.db, input); + } + + listUnconsumedSessionInboundEvents( + input: ListUnconsumedSessionInboundEventsInput, + ): Promise { + return listUnconsumedSessionInboundEventsQuery(this.db, input); + } + + markSessionInboundEventsConsumed(input: MarkSessionInboundEventsConsumedInput): Promise { + return markSessionInboundEventsConsumedQuery(this.db, input); + } + addThreads(input: AddThreadsInput): Promise { return addThreadsQuery(this.db, input); } diff --git a/packages/trueforge/src/db/sqlite/session-store/queries/inboundEvents.ts b/packages/trueforge/src/db/sqlite/session-store/queries/inboundEvents.ts new file mode 100644 index 000000000..108a965bf --- /dev/null +++ b/packages/trueforge/src/db/sqlite/session-store/queries/inboundEvents.ts @@ -0,0 +1,149 @@ +import type { SessionInboundEventItem, TurnState } from '@truefoundry/trueforge-core/agent-session'; +import type { + InsertSessionInboundEventsInput, + ListUnconsumedSessionInboundEventsInput, + MarkSessionInboundEventsConsumedInput, + SessionInboundEventRecord, +} from '@truefoundry/trueforge-core/agent-session/store/ISessionStore'; +import { + SessionInboundEventAlreadyExistsError, + SessionNotFoundError, + TurnNotFoundError, + TurnNotRunningError, +} from '@truefoundry/trueforge-core/agent-session/store/SessionStoreErrors'; +import type { Kysely } from 'kysely'; +import { sql } from 'kysely'; +import { firstCollidingEventId, firstDuplicateEventIdInBatch } from '../../../sessionInboundEvents'; +import { isUniqueViolation } from '../../client'; +import { jsonbBind, jsonText } from '../../sqlExpressions'; +import type { Database } from '../../types'; + +async function requireSession(db: Kysely, sessionId: string): Promise { + const row = await db + .selectFrom('session') + .select('session_id') + .where('session_id', '=', sessionId) + .executeTakeFirst(); + if (!row) { + throw new SessionNotFoundError(sessionId); + } +} + +/** Exists + non-terminal (v1: `running` only; `paused` will be allowed when that status lands). */ +async function requireTurn(db: Kysely, sessionId: string, turnId: string): Promise { + const row = await db + .selectFrom('turn') + .select(['turn_id', jsonText(sql.ref('state')).as('state')]) + .where('session_id', '=', sessionId) + .where('turn_id', '=', turnId) + .executeTakeFirst(); + if (!row) { + throw new TurnNotFoundError(turnId); + } + if (row.state.status !== 'running') { + throw new TurnNotRunningError(turnId, row.state); + } +} + +async function resolveCollidingEventId( + db: Kysely, + sessionId: string, + events: InsertSessionInboundEventsInput['events'], +): Promise { + const ids = [...new Set(events.map(e => e.event_id))]; + if (ids.length === 0) { + return ''; + } + const rows = await db + .selectFrom('session_inbound_events') + .select('event_id') + .where('session_id', '=', sessionId) + .where('event_id', 'in', ids) + .execute(); + return firstCollidingEventId(events, new Set(rows.map(r => r.event_id))); +} + +export async function insertSessionInboundEvents( + db: Kysely, + input: InsertSessionInboundEventsInput, +): Promise { + if (input.events.length === 0) { + return; + } + await requireSession(db, input.session_id); + await requireTurn(db, input.session_id, input.turn_id); + + const duplicateInBatch = firstDuplicateEventIdInBatch(input.events); + if (duplicateInBatch !== undefined) { + throw new SessionInboundEventAlreadyExistsError(input.session_id, duplicateInBatch); + } + + try { + await db + .insertInto('session_inbound_events') + .values( + input.events.map(event => ({ + session_id: input.session_id, + event_id: event.event_id, + turn_id: input.turn_id, + payload: jsonbBind(event.payload), + consumed: 0, + created_at: event.created_at, + })), + ) + .execute(); + } catch (error) { + if (isUniqueViolation(error)) { + const eventId = await resolveCollidingEventId(db, input.session_id, input.events); + throw new SessionInboundEventAlreadyExistsError(input.session_id, eventId, { + cause: error, + }); + } + throw error; + } +} + +export async function listUnconsumedSessionInboundEvents( + db: Kysely, + input: ListUnconsumedSessionInboundEventsInput, +): Promise { + await requireSession(db, input.session_id); + + let query = db + .selectFrom('session_inbound_events') + .select(['event_id', 'turn_id', 'created_at', jsonText(sql.ref('payload')).as('payload')]) + .where('session_id', '=', input.session_id) + .where('consumed', '=', 0); + + if (input.turn_id === null) { + query = query.where('turn_id', 'is', null); + } else if (input.turn_id !== undefined) { + query = query.where('turn_id', '=', input.turn_id); + } + + const rows = await query.orderBy('event_id', 'asc').execute(); + + return rows.map(row => ({ + event_id: row.event_id, + turn_id: row.turn_id, + payload: row.payload, + created_at: row.created_at, + })); +} + +export async function markSessionInboundEventsConsumed( + db: Kysely, + input: MarkSessionInboundEventsConsumedInput, +): Promise { + if (input.event_ids.length === 0) { + return; + } + await requireSession(db, input.session_id); + + await db + .updateTable('session_inbound_events') + .set({ consumed: 1 }) + .where('session_id', '=', input.session_id) + .where('event_id', 'in', input.event_ids) + .execute(); +} diff --git a/packages/trueforge/src/db/sqlite/types.ts b/packages/trueforge/src/db/sqlite/types.ts index 9c2b62303..befe4af55 100644 --- a/packages/trueforge/src/db/sqlite/types.ts +++ b/packages/trueforge/src/db/sqlite/types.ts @@ -9,6 +9,7 @@ import type { AgentSpec, CreatedBySubject, PersistedTurnEvent, + SessionInboundEventItem, SessionMetadata, SessionMetrics, SessionSource, @@ -147,6 +148,20 @@ export interface SessionEventTable { created_at: string; } +/** + * Session inbound send-event inbox (tip HITL + future session-scoped payloads). + * PRIMARY KEY (session_id, event_id). `turn_id` nullable. + * `consumed` is INTEGER 0/1 (STRICT has no boolean). + */ +export interface SessionInboundEventsTable { + session_id: string; + event_id: string; + turn_id: string | null; + payload: JsonbColumn; + consumed: number; + created_at: string; +} + /** * Pure immutable content; no state → no checkpoint field. * PRIMARY KEY (append_id) AUTOINCREMENT @@ -340,6 +355,7 @@ export interface Database { turn_thread: TurnThreadTable; turn_thread_context: TurnThreadContextTable; session_event: SessionEventTable; + session_inbound_events: SessionInboundEventsTable; thread_context_log: ThreadContextLogTable; thread_capability_state: ThreadCapabilityStateTable; model_provider: ModelProviderTable;