diff --git a/backends/example-platform/apps/service/bin/production-server.ts b/backends/example-platform/apps/service/bin/production-server.ts index f7b790ee6a9..c8da5fad003 100644 --- a/backends/example-platform/apps/service/bin/production-server.ts +++ b/backends/example-platform/apps/service/bin/production-server.ts @@ -88,6 +88,8 @@ export async function startProductionServer( mcp_handler: () => Response.json({ error: "unavailable" }, { status: 503, headers: { "cache-control": "no-store" } }), tasks: { authorization, codecRootSecret: config.codecKey, cursorSigningKeyset }, conversations: { authorization, codecRootSecret: config.codecKey, cursorSigningKeyset }, + chat: { authorization, codecRootSecret: config.codecKey, cursorSigningKeyset }, + settings: { authorization }, device_sessions: authorization, device_ownership_key: config.codecKey, transcription_source: createDeepgramTranscriptionSource({ @@ -186,7 +188,7 @@ export async function runProductionServer( try { const running = await startProductionServer(env, factories, controller.signal); clearTimeout(startupDeadline); - if (!controller.signal.aborted) console.info("omi-platform ready: memories.read, tasks, device audio uploads"); + if (!controller.signal.aborted) console.info("omi-platform ready: memories.read, tasks, device audio uploads, conversations.read, chat.read, settings"); await requested; await running.stop(); return 0; diff --git a/backends/example-platform/apps/service/composition/chat-conversation-sessions.test.ts b/backends/example-platform/apps/service/composition/chat-conversation-sessions.test.ts new file mode 100644 index 00000000000..df44cdd5db3 --- /dev/null +++ b/backends/example-platform/apps/service/composition/chat-conversation-sessions.test.ts @@ -0,0 +1,95 @@ +import { describe, expect, test } from "bun:test"; +import { + MAIN_CHAT_CONVERSATION_ID, + composeChatSessionsIntoConversationPage, + type ChatConversationSessionItem, +} from "./chat-conversation-sessions"; + +const session = ( + overrides: Partial = {}, +): ChatConversationSessionItem => Object.freeze({ + id: MAIN_CHAT_CONVERSATION_ID, + title: "hello", + overview: "answer", + createdAt: 1000, + updatedAt: 2000, + startedAt: 1000, + finishedAt: null, + source: "chat", + status: "in_progress", + discarded: false, + starred: false, + visibility: "private", + isLocked: false, + folderId: null, + revision: null, + ...overrides, +}); + +const page = ( + items: ReadonlyArray>, + extras: Record = {}, +) => Object.freeze({ + contractVersion: "1.0.0", + items, + window: Object.freeze({ + status: "complete", + complete: true, + hasMore: false, + nextCursor: extras.nextCursor ?? null, + }), + completeness: Object.freeze({ + version: "conversations-completeness-v1", + status: "complete", + reasons: Object.freeze([]), + }), + absence: items.length === 0 ? Object.freeze({ kind: "query_gap" }) : null, +}); + +describe("chat conversation composition", () => { + test("does not invent chat:chat-main when no granted sessions exist", () => { + expect(composeChatSessionsIntoConversationPage(page([]), [])).toBeNull(); + expect(composeChatSessionsIntoConversationPage(page([{ + id: "recording:one", + updatedAt: 3000, + title: "Recording", + }]), [])).toBeNull(); + }); + + test("merges a persisted main chat session onto the listen page without changing the cursor", () => { + const listen = page([ + { id: "recording:newer", updatedAt: 4000, title: "Newer" }, + { id: "recording:older", updatedAt: 500, title: "Older" }, + ], { nextCursor: "listen-cursor" }); + const composed = composeChatSessionsIntoConversationPage(listen, [session()]); + expect(composed?.window).toEqual(listen.window); + expect(composed?.absence).toBeNull(); + expect(composed?.items.map((item) => item.id)).toEqual([ + "recording:newer", + MAIN_CHAT_CONVERSATION_ID, + "recording:older", + ]); + }); + + test("replaces a listen-claimed chat id with the granted chat session and clears an empty-page gap", () => { + const composed = composeChatSessionsIntoConversationPage( + page([]), + [session({ title: "saved prompt" })], + ); + expect(composed?.absence).toBeNull(); + expect(composed?.items).toEqual([session({ title: "saved prompt" })]); + expect(composeChatSessionsIntoConversationPage( + page([{ id: MAIN_CHAT_CONVERSATION_ID, updatedAt: 1, title: "stale" }]), + [session()], + )?.items).toEqual([session()]); + }); + + test("rejects a malformed listen envelope instead of dropping or inventing rows", () => { + expect(composeChatSessionsIntoConversationPage(null, [session()])).toBeNull(); + expect(composeChatSessionsIntoConversationPage({ items: "rows" }, [session()])).toBeNull(); + expect(composeChatSessionsIntoConversationPage( + page([{ id: "recording:one", updatedAt: "later" }]), + [session()], + )).toBeNull(); + }); +}); diff --git a/backends/example-platform/apps/service/composition/chat-conversation-sessions.ts b/backends/example-platform/apps/service/composition/chat-conversation-sessions.ts new file mode 100644 index 00000000000..fddbe1b2ef3 --- /dev/null +++ b/backends/example-platform/apps/service/composition/chat-conversation-sessions.ts @@ -0,0 +1,75 @@ +export const MAIN_CHAT_CONVERSATION_ID = "chat:chat-main"; + +export type ChatConversationSessionItem = { + readonly id: typeof MAIN_CHAT_CONVERSATION_ID; + readonly title: string; + readonly overview: string; + readonly createdAt: number; + readonly updatedAt: number; + readonly startedAt: number; + readonly finishedAt: number | null; + readonly source: "chat"; + readonly status: "completed" | "in_progress"; + readonly discarded: false; + readonly starred: false; + readonly visibility: "private"; + readonly isLocked: false; + readonly folderId: null; + readonly revision: null; +}; + +export type ConversationEnvelopePage = { + readonly contractVersion: unknown; + readonly items: readonly Record[]; + readonly window: unknown; + readonly completeness: unknown; + readonly absence: { readonly kind: "query_gap" } | null; +}; + +const record = (value: unknown): Record | null => + value !== null && typeof value === "object" && !Array.isArray(value) + ? value as Record + : null; + +const compareItems = ( + left: Record, + right: Record, +): number => { + const leftUpdated = left.updatedAt; + const rightUpdated = right.updatedAt; + if (typeof leftUpdated === "number" && typeof rightUpdated === "number" + && leftUpdated !== rightUpdated) { + return rightUpdated - leftUpdated; + } + const leftId = typeof left.id === "string" ? left.id : ""; + const rightId = typeof right.id === "string" ? right.id : ""; + return leftId < rightId ? -1 : leftId > rightId ? 1 : 0; +}; + +export const composeChatSessionsIntoConversationPage = ( + page: unknown, + sessions: readonly ChatConversationSessionItem[], +): ConversationEnvelopePage | null => { + const envelope = record(page); + if (envelope === null || !Array.isArray(envelope.items) || sessions.length === 0) { + return null; + } + const items: Record[] = []; + const sessionIds = new Set(sessions.map((session) => session.id)); + for (const item of envelope.items) { + const row = record(item); + if (row === null || typeof row.id !== "string" || typeof row.updatedAt !== "number") { + return null; + } + if (!sessionIds.has(row.id)) items.push(row); + } + items.push(...sessions); + items.sort(compareItems); + return Object.freeze({ + contractVersion: envelope.contractVersion, + items: Object.freeze(items), + window: envelope.window, + completeness: envelope.completeness, + absence: items.length === 0 ? { kind: "query_gap" } : null, + }); +}; diff --git a/backends/example-platform/apps/service/memory-service-app.ts b/backends/example-platform/apps/service/memory-service-app.ts index 8ec25a88744..f1f7c7ff818 100644 --- a/backends/example-platform/apps/service/memory-service-app.ts +++ b/backends/example-platform/apps/service/memory-service-app.ts @@ -24,6 +24,8 @@ export const createMemoryServiceApp = ( tasks?: { readonly executeRequest: (request: Request) => Promise }, deviceSessions?: { readonly fetch: (request: Request) => Promise }, conversations?: {readonly executeRequest: (request: Request) => Promise}, + chat?: {readonly executeRequest: (request: Request) => Promise}, + settings?: {readonly executeRequest: (request: Request) => Promise}, ): Hono => { const app = createServiceApp(mcpHandler, observability); registerMemoryRoutes(app, memoryRoutes); @@ -40,5 +42,7 @@ export const createMemoryServiceApp = ( app.get("/v1/device-sessions/:id/transcript", context => deviceSessions.fetch(context.req.raw)); } if (conversations) app.get("/v1/conversations", context => conversations.executeRequest(context.req.raw)); + if (chat) app.get("/v1/chat-messages", context => chat.executeRequest(context.req.raw)); + if (settings) app.get("/v1/settings", context => settings.executeRequest(context.req.raw)); return app; }; diff --git a/backends/example-platform/apps/service/routes/chat-messages.ts b/backends/example-platform/apps/service/routes/chat-messages.ts index 56c5e13bc3f..2e1ef11d31b 100644 --- a/backends/example-platform/apps/service/routes/chat-messages.ts +++ b/backends/example-platform/apps/service/routes/chat-messages.ts @@ -38,6 +38,7 @@ import { import type { ChatMessageRecord, ChatMessagesStore, + StoredChatMessage, WritableChatMessageType, } from "../stores/chat-messages-store"; import type { @@ -302,7 +303,7 @@ export const chatMessagePayloadHash = (create: ParsedCreate): string => { return `sha256:${createHash("sha256").update(canonicalJson(subject), "utf8").digest("hex")}`; }; -const parseHistoryQuery = (request: Request): { +export const parseHistoryQuery = (request: Request): { readonly limit: number; readonly olderCursor: string | null; } | null => { @@ -328,7 +329,7 @@ const TERMINAL_KINDS = new Set(["done", "failed", "cancelled"]); const isTerminal = (event: ChatGenerationEvent): boolean => TERMINAL_KINDS.has(event.frame.kind); type ChatGenerationOutcome = "completed" | "cancelled" | null; -type ChatWireMessage = ChatMessageRecord & { +export type ChatWireMessage = ChatMessageRecord & { readonly generationOutcome: ChatGenerationOutcome; }; @@ -345,19 +346,16 @@ const sameCanonicalMessage = (left: ChatMessageRecord, right: ChatMessageRecord) * row deliberately does not duplicate that state, so an orphan or a mismatched * terminal fails closed instead of silently becoming a completed answer. */ -const projectHistoryMessage = ( - accountId: string, +export const projectLoadedHistoryMessage = ( message: ChatMessageRecord, - messages: ChatMessagesStore, - events: ChatGenerationEventsStore, + stored: StoredChatMessage | null, + generationEvents: readonly ChatGenerationEvent[] | null, ): ChatWireMessage => { if (message.sender !== "ai") return withGenerationOutcome(message, null); - const stored = messages.readMessage(accountId, message.id); if (stored === null || stored.generationId === null || !sameCanonicalMessage(stored.message, message)) { throw new TypeError("canonical assistant has no matching generation identity"); } - const generationEvents = events.listAfter(accountId, stored.generationId, null); const terminals = generationEvents?.filter(isTerminal) ?? []; if (terminals.length !== 1) { throw new TypeError("canonical assistant has no unique terminal event"); @@ -374,6 +372,20 @@ const projectHistoryMessage = ( ); }; +export const projectHistoryMessage = ( + accountId: string, + message: ChatMessageRecord, + messages: Pick, + events: Pick, +): ChatWireMessage => { + if (message.sender !== "ai") return projectLoadedHistoryMessage(message, null, null); + const stored = messages.readMessage(accountId, message.id); + const generationEvents = stored?.generationId === undefined || stored.generationId === null + ? null + : events.listAfter(accountId, stored.generationId, null); + return projectLoadedHistoryMessage(message, stored, generationEvents); +}; + type ExternalChatGenerationEvent = ChatGenerationEvent & { readonly frame: Exclude; }; diff --git a/backends/example-platform/docs/memory-productionization/chat-deployed.md b/backends/example-platform/docs/memory-productionization/chat-deployed.md new file mode 100644 index 00000000000..65783c36aa0 --- /dev/null +++ b/backends/example-platform/docs/memory-productionization/chat-deployed.md @@ -0,0 +1,33 @@ +# Persisted chat history reads + +The deployed service mounts `GET /v1/chat-messages` through the same admission, +readiness and drain boundary as the other REST routes. It verifies the original +Firebase identity and requires a registered active credential with the exact +`chat.read` grant. A `memories.read`, `tasks.read` or `conversations.read` grant +does not confer this read permission. Missing or revoked grants return 403; they +never become an empty successful transcript. Deployment still needs the +authoritative account, control, credential and grant records described in +`deployed-entry.md`. + +Migration 0055 adds account-owned chat message rows and generation events. History +uses the insertion snapshot and opaque HMAC cursor already used by the local +service. An empty granted account is an honest empty page with the existing +attachment capability advertisement. Assistant rows require a unique terminal +generation event; an orphan or mismatched terminal is 503 rather than a completed +answer. Human rows keep `generationOutcome: null`. + +`POST /v1/chat-messages` and generation SSE are not mounted. Unmounted writes stay +404 `{error:"not_found"}`. Do not invent chat quotas or mount admission until a +real entitlement producer exists. + +Verification uses `bun run check:deployed` for grant denial, empty-page shape, +projection fail-closed behavior, string generation frames, route pairing and the +production import closure, and `bun run test:postgres` for actual application-role +reads, account isolation, unique-terminal assistant outcomes, grant revocation +and conversation-list composition of `chat:chat-main`. Docker is +required for that real PostgreSQL 18.4 gate. These tests use isolated synthetic +identities; they do not activate a deployed user or prove live generation. +Do not apply migrations 55-56 or deploy this entry until the existing operator +migration sequence can run against based-hardware-dev. A process built from this +manifest will not become ready against a database that still has only +migrations 1–54. diff --git a/backends/example-platform/docs/memory-productionization/conversations-deployed.md b/backends/example-platform/docs/memory-productionization/conversations-deployed.md index 6ce72607b03..05990e930b4 100644 --- a/backends/example-platform/docs/memory-productionization/conversations-deployed.md +++ b/backends/example-platform/docs/memory-productionization/conversations-deployed.md @@ -45,10 +45,14 @@ Migration 0053 removes the old whole-account function and introduces a distinct metadata read; an old process calling the retired function fails unavailable rather than returning an empty successful history during a mixed-revision rollout. -This is the persisted Listen/recording list, not a production chat history store. -Chat conversations, editable metadata, folders and star mutations still need their -own actual persisted domain composition. Full transcript data remains on the -existing account-scoped device-session transcript route. +This is the persisted Listen/recording list, plus granted main-chat sessions. +First-page envelope reads that also hold `chat.read` include `chat:chat-main` when +that account actually has main-session messages. Missing `chat.read`, revoked +grants, and empty chat history never invent that row. Later cursor pages keep the +Listen sequence and do not repeat the chat session. Chat writes, editable +metadata, folders and star mutations still need their own persisted domain +composition. Full recording transcript data remains on the existing +account-scoped device-session transcript route. Verification uses `bun run check:deployed` for projection, expiry/cancellation, route and shell contracts, and `bun run test:postgres` for actual application-role diff --git a/backends/example-platform/docs/memory-productionization/deployed-entry.md b/backends/example-platform/docs/memory-productionization/deployed-entry.md index a717509ebc0..a02d1dbcfe3 100644 --- a/backends/example-platform/docs/memory-productionization/deployed-entry.md +++ b/backends/example-platform/docs/memory-productionization/deployed-entry.md @@ -16,11 +16,14 @@ Build from this backend directory with `docker build --platform linux/amd64 -t omi-platform-dev .`. This increment serves authenticated canonical memory reads, task reads and -mutations, and indexed device audio uploads. All three use the same database -generation and Firebase authorization configuration inside the readiness and -shutdown boundary. Audio upload completion does not certify transcription or -conversation formation. Chat and authenticated MCP remain unavailable; MCP -returns 503. This is not full backend parity or production qualification. +mutations, indexed device audio uploads, conversation reads, chat history +reads, granted main-chat session composition on the first conversation page, and Settings GET. All of them use the same database generation and Firebase authorization +configuration inside the readiness and shutdown boundary. Audio upload completion +does not certify transcription or conversation formation. Chat writes, generation +SSE, Settings identity/entitlement producers, attachments and authenticated MCP +remain unavailable; MCP returns 503. Missing `chat.read` is 403, not an empty +successful transcript. Signed-in Settings without a producer is 503, not a 200 +profile invented from the Firebase token. This is not full backend parity or production qualification. Required configuration: diff --git a/backends/example-platform/docs/memory-productionization/settings-deployed.md b/backends/example-platform/docs/memory-productionization/settings-deployed.md new file mode 100644 index 00000000000..f6a6f82ca7b --- /dev/null +++ b/backends/example-platform/docs/memory-productionization/settings-deployed.md @@ -0,0 +1,17 @@ +# Production Settings reads + +The deployed service mounts `GET /v1/settings` through the same admission, +readiness and drain boundary as the other REST routes. Absent credentials return +the signed-out envelope `{identity:null,entitlement:null}`. Present invalid +credentials are 401. A verified Firebase identity without an owner-backed profile +and entitlement producer is 503 `{error:"service_unavailable"}` with +`retry-after: 60`. Token claims are never projected as display names, emails, or +plans. PostgreSQL account/grant rows are not a Settings producer. + +Mutations stay unmounted. Unmounted writes stay 404 `{error:"not_found"}`. Do not +invent identity, billing, or usage to satisfy the page. + +Verification uses `bun run check:deployed` for signed-out, unauthorized, verified +unavailable, grammar, pairing and the production import closure. These tests use +isolated synthetic identities; they do not activate a deployed user or prove a +billing producer. diff --git a/backends/example-platform/drivers/postgres/chat-messages.real.test.ts b/backends/example-platform/drivers/postgres/chat-messages.real.test.ts new file mode 100644 index 00000000000..a7d0a7b1c30 --- /dev/null +++ b/backends/example-platform/drivers/postgres/chat-messages.real.test.ts @@ -0,0 +1,192 @@ +import { expect, test } from "bun:test"; +import postgres from "postgres"; +import { createHash, randomUUID } from "node:crypto"; +import { runPostgresMigrations } from "./migrations/runner"; +import { createPostgresJsTransactionPool } from "./postgresjs"; +import type { PostgresTransactionPool } from "./connection"; +import { createPostgresFirebaseChatReadRuntime } from "./firebase-chat-read-runtime"; +import { seedProdLocalFirebaseAuthorizationSql } from "../../scripts/prod-local-identity-seed"; +import { CHAT_CAPABILITIES } from "../../apps/service/routes/chat-messages"; + +const url = process.env.OMI_TEST_POSTGRES_URL; +const realTest = url ? test : test.skip; + +realTest("real chat reads require chat.read and never invent empty success", async () => { + const endpoint = new URL(url!); + if (endpoint.hostname !== "127.0.0.1" || endpoint.protocol !== "postgres:") { + throw Error("postgres_test_not_loopback_only"); + } + const owner = postgres(url!, { max: 1 }); + const pool = createPostgresJsTransactionPool({ + connectionString: url!, + maxConnections: 2, + }); + const suffix = randomUUID(); + const generation = createHash("sha256").update(suffix).digest("hex"); + const now = () => Math.floor(Date.now() / 1000); + const project = "synthetic-chat-project"; + const app = "synthetic-chat-app"; + const uid = `uid-${suffix}`; + const account = `account-${suffix}`; + const principal = `principal-${suffix}`; + const credential = `credential-${suffix}`; + const humanId = "11111111-1111-4111-8111-111111111111"; + const aiId = "22222222-2222-4222-8222-222222222222"; + const generationId = "gen_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; + try { + await owner.unsafe(`DO $roles$ BEGIN + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_application') THEN CREATE ROLE omi_platform_application NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_cleanup') THEN CREATE ROLE omi_platform_cleanup NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_restore') THEN CREATE ROLE omi_platform_restore NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_restore_operator') THEN CREATE ROLE omi_platform_restore_operator NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + END $roles$;`); + await runPostgresMigrations(owner); + for (const target of [account, `other-${suffix}`]) { + for (const statement of seedProdLocalFirebaseAuthorizationSql({ + firebase_project_id: project, + firebase_uid: target === account ? uid : `other-${uid}`, + application_id: app, + account_id: target, + principal_id: principal, + credential_id: credential, + grant_id: `memory-${suffix}`, + }, now())) { + await owner.unsafe(statement.text, [...statement.values]); + } + await owner.unsafe( + `INSERT INTO omi_memory.application_grant_revisions(account_id,application_id,credential_id,credential_generation,capability,grant_id,grant_version,lifecycle,enabled,scopes,record_schema_version,record_json,content_hash) VALUES($1,$2,$3,1,$4,$5,1,'active',true,'[]','grant-v1','{}',$6)`, + [target, app, credential, "chat.read", `chat.read-${suffix}`, "3".repeat(64)], + ); + await owner.unsafe( + `INSERT INTO omi_memory.application_grant_heads(account_id,application_id,credential_id,credential_generation,capability,grant_id,grant_version) VALUES($1,$2,$3,1,$4,$5,1)`, + [target, app, credential, "chat.read", `chat.read-${suffix}`], + ); + } + await owner.unsafe( + `INSERT INTO omi_memory.postgres_restore_admission_revisions(database_generation_digest,release_revision,state,restore_id,restored_snapshot_digest,checkpoint_candidate_digest,checkpoint_evidence_digest,first_approval_subject_digest,first_approval_receipt_digest,second_approval_subject_digest,second_approval_receipt_digest,manual_release_receipt_digest,previous_release_revision,content_hash) VALUES($1,1,'released',$2,$3,$3,$3,$4,$5,$6,$7,$3,NULL,$3)`, + [generation, `synthetic-${suffix}`, "9".repeat(64), "4".repeat(64), "5".repeat(64), "6".repeat(64), "7".repeat(64)], + ); + await owner.unsafe( + "INSERT INTO omi_memory.postgres_restore_admission_heads(database_generation_digest,release_revision) VALUES($1,1)", + [generation], + ); + const appPool: PostgresTransactionPool = { + withTransaction: (options, callback) => pool.withTransaction(options, async (connection) => { + await connection.query({ + name: "chat_test.role", + text: "SET LOCAL ROLE omi_platform_application", + values: [], + }); + return callback(connection); + }), + }; + const runtime = createPostgresFirebaseChatReadRuntime({ + authorization: { + pool: appPool, + project_id: project, + application_id: app, + runtime_mode: "deployed", + context_ttl_seconds: 60, + database_generation_digest: generation, + id_token_adapter: { + verification_source: "firebase_production", + async verifyIdToken(token) { + const selected = token.startsWith("other") ? `other-${uid}` : uid; + return { + aud: project, + iss: `https://securetoken.google.com/${project}`, + sub: selected, + uid: selected, + iat: now() - 10, + auth_time: now() - 10, + exp: now() + 600, + }; + }, + }, + }, + codecRootSecret: new Uint8Array(32).fill(7), + cursorSigningKeyset: { + active_key_id: "test", + keys: [{ key_id: "test", secret: new Uint8Array(32).fill(8) }], + }, + }); + const call = (query = "?limit=50", token = "header.payload.signature", method = "GET") => + runtime.executeRequest(new Request(`https://chat.example/v1/chat-messages${query}`, { + method, + headers: { authorization: `Bearer ${token}` }, + ...(method === "POST" ? { body: "{}" } : {}), + })); + const empty = await call(); + expect(empty.status).toBe(200); + expect(await empty.json()).toEqual({ + messages: [], + page: { olderCursor: null, hasOlder: false }, + capabilities: CHAT_CAPABILITIES, + }); + await owner.unsafe( + `INSERT INTO omi_memory.chat_messages(account_id,id,text,sender,message_type,created_at,updated_at,chat_session_id,app_id,journal_revision,payload_hash,message_source,rating,reported,server_revision,attachments_json,generation_id) VALUES($1,$2,'hello','human','text',1000,1000,NULL,NULL,0,'sha256:human','desktop_chat',NULL,false,'rev-human','[]'::jsonb,'gen_human')`, + [account, humanId], + ); + const human = await call(); + expect(human.status).toBe(200); + const humanPage = await human.json() as { + messages: Array<{ sender: string; generationOutcome: unknown; attachments: unknown }>; + page: { olderCursor: string | null; hasOlder: boolean }; + }; + expect(humanPage.messages).toHaveLength(1); + expect(humanPage.messages[0]?.sender).toBe("human"); + expect(humanPage.messages[0]?.generationOutcome).toBeNull(); + expect(humanPage.messages[0]?.attachments).toEqual([]); + expect(humanPage.page).toEqual({ olderCursor: null, hasOlder: false }); + await owner.unsafe( + `INSERT INTO omi_memory.chat_messages(account_id,id,text,sender,message_type,created_at,updated_at,chat_session_id,app_id,journal_revision,payload_hash,message_source,rating,reported,server_revision,attachments_json,generation_id) VALUES($1,$2,'answer','ai','text',2000,2000,NULL,NULL,0,'sha256:ai','desktop_chat',NULL,false,'rev-ai','[]'::jsonb,$3)`, + [account, aiId, generationId], + ); + expect((await call()).status).toBe(503); + const assistant = { + id: aiId, + text: "answer", + sender: "ai", + type: "text", + createdAt: 2000, + updatedAt: 2000, + chatSessionId: null, + appId: null, + journalRevision: 0, + payloadHash: "sha256:ai", + messageSource: "desktop_chat", + rating: null, + reported: false, + revision: "rev-ai", + attachments: [], + }; + await owner.unsafe( + `INSERT INTO omi_memory.chat_generation_events(account_id,generation_id,sequence,event_id,created_at,frame_json) VALUES($1,$2,1,'evt-done',2000,$3::text::jsonb)`, + [account, generationId, JSON.stringify({ kind: "done", message: assistant })], + ); + const completed = await call(); + expect(completed.status).toBe(200); + const completedPage = await completed.json() as { + messages: Array<{ id: string; sender: string; generationOutcome: unknown }>; + }; + expect(completedPage.messages.map((row) => [row.id, row.sender, row.generationOutcome])).toEqual([ + [humanId, "human", null], + [aiId, "ai", "completed"], + ]); + expect((await call("", "other.payload.signature")).status).toBe(200); + expect(await (await call("", "other.payload.signature")).json()).toEqual({ + messages: [], + page: { olderCursor: null, hasOlder: false }, + capabilities: CHAT_CAPABILITIES, + }); + await owner.unsafe( + "DELETE FROM omi_memory.application_grant_heads WHERE account_id=$1 AND capability='chat.read'", + [account], + ); + expect((await call()).status).toBe(403); + expect((await call("", "header.payload.signature", "POST")).status).toBe(404); + } finally { + await pool.close(); + await owner.end(); + } +}); diff --git a/backends/example-platform/drivers/postgres/chat-read-repository.test.ts b/backends/example-platform/drivers/postgres/chat-read-repository.test.ts new file mode 100644 index 00000000000..77f672bcc1a --- /dev/null +++ b/backends/example-platform/drivers/postgres/chat-read-repository.test.ts @@ -0,0 +1,264 @@ +import { expect, test } from "bun:test"; +import { createAuthorizedLedgerWriteContextIssuer } from "../../apps/service/auth/authorized-context-internal"; +import { + authorizationStateDigest, + type AuthorityStateRow, +} from "./transaction"; +import type { + PostgresTransactionPool, + CheckedOutPostgresConnection, +} from "./connection"; +import { withAuthorizedChatRead } from "./chat-read-repository"; + +const hash = (character: string): string => character.repeat(64); +const account = "account:alice"; + +const authorityRow = (): AuthorityStateRow => ({ + account_id: account, + principal_id: "principal:chat", + application_id: "app:chat", + credential_id: "credential:chat", + credential_generation: 1, + capability: "chat.read", + grant_id: "grant:chat", + grant_version: 1, + account_epoch: 2, + control_conflict_reason: null, + control_conflict_at_revision: null, + destination_activation_epoch: 2, + destination_activation_revision: 3, + lifecycle_state: "active", + deletion_epoch: null, + account_generation: "new", + credential_lifecycle: "active", + grant_lifecycle: "active", + grant_enabled: true, + authentication_strength: "service-workload", + credential_expires_at_epoch_seconds: 10_000, + control_revision: 3, + control_content_hash: hash("1"), + credential_content_hash: hash("2"), + grant_content_hash: hash("3"), + db_now_epoch_seconds: 100, +}); + +const context = (capability = "chat.read") => { + const row = authorityRow(); + return createAuthorizedLedgerWriteContextIssuer().issue( + { + context_version: "authorized-ledger-write-context-v1", + principal_id: row.principal_id, + account_id: account, + application_id: row.application_id, + credential_id: row.credential_id, + credential_generation: row.credential_generation, + capability, + grant_id: row.grant_id, + grant_version: row.grant_version, + account_epoch: 2, + destination_activation_revision: 3, + lifecycle_state: "active", + deletion_epoch: null, + authentication_strength: row.authentication_strength, + issued_at_epoch_seconds: 50, + expires_at_epoch_seconds: 500, + authorization_state_digest: authorizationStateDigest(row), + }, + 100, + ); +}; + +test("chat history remains private until final clock and cancellation checks pass", async () => { + for (const mode of ["success", "expired", "aborted"]) { + const controller = new AbortController(); + let committed = false; + let released!: () => void; + let began!: () => void; + const wait = new Promise((resolve) => { + released = resolve; + }); + const started = new Promise((resolve) => { + began = resolve; + }); + const connection: CheckedOutPostgresConnection = { + connectionIdentity: {}, + async execute() { + return { rowCount: 0 }; + }, + async query(statement) { + const rows = statement.name === "authority.lock_and_revalidate" + ? [authorityRow()] + : statement.name === "chat.read_snapshot" + ? [{ sequence: 0 }] + : statement.name === "chat.final_clock" + ? [{ now: mode === "expired" ? 500 : 499 }] + : []; + return rows as never; + }, + }; + const pool: PostgresTransactionPool = { + async withTransaction(options, operation) { + expect(options.signal).toBe(controller.signal); + const result = await operation(connection); + committed = true; + return result; + }, + }; + const pending = withAuthorizedChatRead( + pool, + context(), + controller.signal, + async (storage) => { + began(); + await wait; + return storage.readSnapshotSequence(); + }, + ); + await started; + expect(committed).toBe(false); + if (mode === "aborted") controller.abort(); + released(); + if (mode === "success") await expect(pending).resolves.toBe(0); + else if (mode === "expired") await expect(pending).rejects.toMatchObject({ code: "expired_context" }); + else await expect(pending).rejects.toThrow(); + expect(committed).toBe(mode === "success"); + } +}); + +test("an unrelated capability never checks out a chat connection", async () => { + const pool: PostgresTransactionPool = { + async withTransaction() { + throw Error("must not reach pool"); + }, + }; + await expect(withAuthorizedChatRead( + pool, + context("memories.read"), + new AbortController().signal, + () => null, + )).rejects.toMatchObject({ code: "capability_denied" }); +}); + +test("unreadable stored chat history fails closed instead of inventing an empty page", async () => { + const connection: CheckedOutPostgresConnection = { + connectionIdentity: {}, + async execute() { + return { rowCount: 0 }; + }, + async query(statement) { + const rows = statement.name === "authority.lock_and_revalidate" + ? [authorityRow()] + : statement.name === "chat.read_history" + ? [{ page: { hasOlder: false, messages: "corrupt" } }] + : statement.name === "chat.final_clock" + ? [{ now: 100 }] + : []; + return rows as never; + }, + }; + const pool: PostgresTransactionPool = { + async withTransaction(_options, operation) { + return operation(connection); + }, + }; + await expect(withAuthorizedChatRead( + pool, + context(), + new AbortController().signal, + (storage) => storage.listHistory({ + limit: 50, + snapshotSequence: 0, + olderThan: null, + }), + )).rejects.toMatchObject({ code: "persistence_failed" }); +}); + +test("a string generation frame fails closed instead of completing an assistant row", async () => { + const connection: CheckedOutPostgresConnection = { + connectionIdentity: {}, + async execute() { + return { rowCount: 0 }; + }, + async query(statement) { + const rows = statement.name === "authority.lock_and_revalidate" + ? [authorityRow()] + : statement.name === "chat.read_generation_events" + ? [{ events: [{ id: "evt-done", generationId: "gen", sequence: 1, createdAt: 1, frame: "{\"kind\":\"done\"}" }] }] + : statement.name === "chat.final_clock" + ? [{ now: 100 }] + : []; + return rows as never; + }, + }; + const pool: PostgresTransactionPool = { + async withTransaction(_options, operation) { + return operation(connection); + }, + }; + await expect(withAuthorizedChatRead( + pool, + context(), + new AbortController().signal, + (storage) => storage.listGenerationEvents("gen"), + )).rejects.toMatchObject({ code: "persistence_failed" }); +}); + +test("unreadable chat conversation sessions fail closed instead of inventing chat:chat-main", async () => { + const connection: CheckedOutPostgresConnection = { + connectionIdentity: {}, + async execute() { + return { rowCount: 0 }; + }, + async query(statement) { + const rows = statement.name === "authority.lock_and_revalidate" + ? [authorityRow()] + : statement.name === "chat.read_conversation_sessions" + ? [{ sessions: { id: "chat:chat-main" } }] + : statement.name === "chat.final_clock" + ? [{ now: 100 }] + : []; + return rows as never; + }, + }; + const pool: PostgresTransactionPool = { + async withTransaction(_options, operation) { + return operation(connection); + }, + }; + await expect(withAuthorizedChatRead( + pool, + context(), + new AbortController().signal, + (storage) => storage.listConversationSessions(), + )).rejects.toMatchObject({ code: "persistence_failed" }); +}); + +test("granted empty chat conversation sessions stay an empty list", async () => { + const connection: CheckedOutPostgresConnection = { + connectionIdentity: {}, + async execute() { + return { rowCount: 0 }; + }, + async query(statement) { + const rows = statement.name === "authority.lock_and_revalidate" + ? [authorityRow()] + : statement.name === "chat.read_conversation_sessions" + ? [{ sessions: [] }] + : statement.name === "chat.final_clock" + ? [{ now: 100 }] + : []; + return rows as never; + }, + }; + const pool: PostgresTransactionPool = { + async withTransaction(_options, operation) { + return operation(connection); + }, + }; + await expect(withAuthorizedChatRead( + pool, + context(), + new AbortController().signal, + (storage) => storage.listConversationSessions(), + )).resolves.toEqual([]); +}); diff --git a/backends/example-platform/drivers/postgres/chat-read-repository.ts b/backends/example-platform/drivers/postgres/chat-read-repository.ts new file mode 100644 index 00000000000..c34f8386376 --- /dev/null +++ b/backends/example-platform/drivers/postgres/chat-read-repository.ts @@ -0,0 +1,266 @@ +// domain-pending(DIV-CHAT-SENDER-001) +// domain-pending(DIV-CHAT-TYPE-001) +// domain-pending(DIV-CHAT-SESSION-001) +// domain-pending(DIV-CHAT-REV-001) +// domain-pending(DIV-CHAT-HASH-001) +// domain-pending(DIV-CHAT-SOURCE-001) + +import { + assertAuthorizedLedgerWriteContextCurrentAt, + type AuthorizedLedgerWriteContext, +} from "../../apps/service/auth/authorized-context"; +import type { + ChatGenerationEvent, + ChatGenerationFrame, +} from "../../apps/service/stores/chat-generation-events-store"; +import { + compareChatHistoryKeys, + detachChatMessage, + type ChatHistoryQuery, + type ChatHistoryStorePage, + type ChatMessageRecord, + type StoredChatMessage, +} from "../../apps/service/stores/chat-messages-store"; +import { + MAIN_CHAT_CONVERSATION_ID, + type ChatConversationSessionItem, +} from "../../apps/service/composition/chat-conversation-sessions"; +import type { PostgresTransactionPool } from "./connection"; +import { + PostgresRepositoryError, + withAuthorizedSerializableConnectionTransaction, +} from "./transaction"; + +export interface ChatReadStorage { + readSnapshotSequence(): Promise; + listHistory(query: ChatHistoryQuery): Promise; + readMessage(messageId: string): Promise; + listGenerationEvents(generationId: string): Promise; + listConversationSessions(): Promise; +} + +const fail = (): never => { + throw new PostgresRepositoryError("persistence_failed"); +}; + +const integer = (value: unknown): number | null => { + if (typeof value === "number" && Number.isSafeInteger(value) && value >= 0) return value; + if (typeof value === "string" && /^(0|[1-9][0-9]*)$/.test(value)) { + const parsed = Number(value); + return Number.isSafeInteger(parsed) && parsed >= 0 ? parsed : null; + } + return null; +}; + +const record = (value: unknown): Record | null => + value !== null && typeof value === "object" && !Array.isArray(value) + ? value as Record + : null; + +const parseStored = (value: unknown): StoredChatMessage => { + const row = record(value); + if (row === null) return fail(); + const generationId = row.generationId; + if (!(generationId === null || typeof generationId === "string")) return fail(); + let message: ChatMessageRecord; + try { + message = detachChatMessage({ + id: row.id as never, + text: row.text as never, + sender: row.sender as never, + type: row.type as never, + createdAt: row.createdAt as never, + updatedAt: row.updatedAt as never, + chatSessionId: row.chatSessionId as never, + appId: row.appId as never, + journalRevision: row.journalRevision as never, + payloadHash: row.payloadHash as never, + messageSource: row.messageSource as never, + rating: row.rating as never, + reported: row.reported as never, + revision: row.revision as never, + attachments: row.attachments as never, + }); + } catch { + return fail(); + } + return Object.freeze({ message, generationId }); +}; + +const parseFrame = (value: unknown): ChatGenerationFrame => { + const frame = record(value); + if (frame === null || typeof frame.kind !== "string" || frame.kind.length === 0) return fail(); + return frame as unknown as ChatGenerationFrame; +}; + +const parseEvent = (value: unknown): ChatGenerationEvent => { + const row = record(value); + if (row === null) return fail(); + const sequence = integer(row.sequence); + const createdAt = integer(row.createdAt); + if (typeof row.id !== "string" || row.id.length === 0 + || typeof row.generationId !== "string" || row.generationId.length === 0 + || sequence === null || sequence < 1 || createdAt === null) return fail(); + return Object.freeze({ + id: row.id, + generationId: row.generationId, + sequence, + createdAt, + frame: parseFrame(row.frame), + }); +}; + +const parseConversationSession = (value: unknown): ChatConversationSessionItem => { + const row = record(value); + const createdAt = integer(row?.createdAt); + const updatedAt = integer(row?.updatedAt); + const startedAt = integer(row?.startedAt); + const finishedAt = row?.finishedAt === null ? null : integer(row?.finishedAt); + if (row === null + || row.id !== MAIN_CHAT_CONVERSATION_ID + || typeof row.title !== "string" || row.title.length === 0 || row.title.length > 240 + || typeof row.overview !== "string" || row.overview.length === 0 || row.overview.length > 240 + || createdAt === null || updatedAt === null || startedAt === null + || updatedAt < createdAt || startedAt !== createdAt + || (finishedAt !== null && finishedAt < createdAt) + || row.source !== "chat" + || (row.status !== "completed" && row.status !== "in_progress") + || row.discarded !== false || row.starred !== false + || row.visibility !== "private" || row.isLocked !== false + || row.folderId !== null || row.revision !== null) { + return fail(); + } + return Object.freeze({ + id: MAIN_CHAT_CONVERSATION_ID, + title: row.title, + overview: row.overview, + createdAt, + updatedAt, + startedAt, + finishedAt, + source: "chat", + status: row.status, + discarded: false, + starred: false, + visibility: "private", + isLocked: false, + folderId: null, + revision: null, + }); +}; + +const parseConversationSessions = (value: unknown): readonly ChatConversationSessionItem[] => { + if (!Array.isArray(value) || value.length > 1) return fail(); + return Object.freeze(value.map(parseConversationSession)); +}; + +const parseHistoryPage = (value: unknown): ChatHistoryStorePage => { + const page = record(value); + if (page === null || typeof page.hasOlder !== "boolean" || !Array.isArray(page.messages)) { + return fail(); + } + const stored = page.messages.map(parseStored); + const messages = stored.map((row) => row.message).sort((left, right) => compareChatHistoryKeys( + { createdAt: left.createdAt, id: left.id }, + { createdAt: right.createdAt, id: right.id }, + )); + return Object.freeze({ + messages: Object.freeze(messages), + hasOlder: page.hasOlder, + }); +}; + +export async function withAuthorizedChatRead( + pool: PostgresTransactionPool, + authority: AuthorizedLedgerWriteContext, + signal: AbortSignal, + project: (storage: ChatReadStorage) => Result | Promise, +): Promise { + if (authority.capability !== "chat.read") { + throw new PostgresRepositoryError("capability_denied"); + } + signal.throwIfAborted(); + const boundedPool: PostgresTransactionPool = { + withTransaction: (options, operation) => + pool.withTransaction({ ...options, signal }, operation), + }; + return withAuthorizedSerializableConnectionTransaction( + boundedPool, + authority, + async ({ connection, dbNowEpochSeconds }) => { + const storage: ChatReadStorage = { + async readSnapshotSequence() { + const rows = await connection.query({ + name: "chat.read_snapshot", + text: "SELECT omi_memory.read_chat_snapshot_sequence() AS sequence", + values: [], + }); + if (rows.length !== 1) return fail(); + const sequence = integer(rows[0]!.sequence); + if (sequence === null) return fail(); + return sequence; + }, + async listHistory(query: ChatHistoryQuery) { + const rows = await connection.query({ + name: "chat.read_history", + text: "SELECT omi_memory.read_chat_history($1,$2,$3,$4) AS page", + values: [ + query.limit, + query.snapshotSequence, + query.olderThan?.createdAt ?? null, + query.olderThan?.id ?? null, + ], + }); + if (rows.length !== 1) return fail(); + return parseHistoryPage(rows[0]!.page); + }, + async readMessage(messageId: string) { + const rows = await connection.query({ + name: "chat.read_message", + text: "SELECT omi_memory.read_chat_message($1) AS message", + values: [messageId], + }); + if (rows.length !== 1) return fail(); + if (rows[0]!.message === null) return null; + return parseStored(rows[0]!.message); + }, + async listGenerationEvents(generationId: string) { + const rows = await connection.query({ + name: "chat.read_generation_events", + text: "SELECT omi_memory.read_chat_generation_events($1) AS events", + values: [generationId], + }); + if (rows.length !== 1) return fail(); + const events = rows[0]!.events; + if (!Array.isArray(events)) return fail(); + return Object.freeze(events.map(parseEvent)); + }, + async listConversationSessions() { + const rows = await connection.query({ + name: "chat.read_conversation_sessions", + text: "SELECT omi_memory.read_chat_conversation_sessions() AS sessions", + values: [], + }); + if (rows.length !== 1) return fail(); + return parseConversationSessions(rows[0]!.sessions); + }, + }; + const result = await project(storage); + const clocks = await connection.query({ + name: "chat.final_clock", + text: "SELECT floor(extract(epoch FROM clock_timestamp()))::bigint AS now", + values: [], + }); + const now = integer(clocks[0]?.now); + if (clocks.length !== 1 || now === null) return fail(); + void dbNowEpochSeconds; + signal.throwIfAborted(); + try { + assertAuthorizedLedgerWriteContextCurrentAt(authority, now); + } catch { + throw new PostgresRepositoryError("expired_context"); + } + return result; + }, + ); +} diff --git a/backends/example-platform/drivers/postgres/conversations.real.test.ts b/backends/example-platform/drivers/postgres/conversations.real.test.ts index 05041a9bda8..360d153d581 100644 --- a/backends/example-platform/drivers/postgres/conversations.real.test.ts +++ b/backends/example-platform/drivers/postgres/conversations.real.test.ts @@ -392,3 +392,189 @@ realTest( }, 120000 ); + +realTest( + "real conversation reads compose granted chat sessions without inventing empty chat:chat-main", + async () => { + const endpoint = new URL(url!); + if (endpoint.hostname !== "127.0.0.1" || endpoint.protocol !== "postgres:") + throw Error("postgres_test_not_loopback_only"); + const owner = postgres(url!, { max: 1 }); + const pool = createPostgresJsTransactionPool({ + connectionString: url!, + maxConnections: 2, + }); + const suffix = randomUUID(), + generation = createHash("sha256").update(suffix).digest("hex"), + now = () => Math.floor(Date.now() / 1000); + const project = "synthetic-chat-conversation-project", + app = "synthetic-chat-conversation-app", + uid = `uid-${suffix}`, + account = `account-${suffix}`, + principal = `principal-${suffix}`, + credential = `credential-${suffix}`; + try { + await owner.unsafe(`DO $roles$ BEGIN + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_application') THEN CREATE ROLE omi_platform_application NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_cleanup') THEN CREATE ROLE omi_platform_cleanup NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_restore') THEN CREATE ROLE omi_platform_restore NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + IF NOT EXISTS(SELECT FROM pg_roles WHERE rolname='omi_platform_restore_operator') THEN CREATE ROLE omi_platform_restore_operator NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT; END IF; + END $roles$;`); + await runPostgresMigrations(owner); + for (const statement of seedProdLocalFirebaseAuthorizationSql( + { + firebase_project_id: project, + firebase_uid: uid, + application_id: app, + account_id: account, + principal_id: principal, + credential_id: credential, + grant_id: `memory-${suffix}`, + }, + now() + )) + await owner.unsafe(statement.text, [...statement.values]); + for (const capability of ["conversations.read"]) { + await owner.unsafe( + `INSERT INTO omi_memory.application_grant_revisions(account_id,application_id,credential_id,credential_generation,capability,grant_id,grant_version,lifecycle,enabled,scopes,record_schema_version,record_json,content_hash) VALUES($1,$2,$3,1,$4,$5,1,'active',true,'[]','grant-v1','{}',$6)`, + [account, app, credential, capability, `${capability}-${suffix}`, "3".repeat(64)] + ); + await owner.unsafe( + `INSERT INTO omi_memory.application_grant_heads(account_id,application_id,credential_id,credential_generation,capability,grant_id,grant_version) VALUES($1,$2,$3,1,$4,$5,1)`, + [account, app, credential, capability, `${capability}-${suffix}`] + ); + } + await owner.unsafe( + `INSERT INTO omi_memory.postgres_restore_admission_revisions(database_generation_digest,release_revision,state,restore_id,restored_snapshot_digest,checkpoint_candidate_digest,checkpoint_evidence_digest,first_approval_subject_digest,first_approval_receipt_digest,second_approval_subject_digest,second_approval_receipt_digest,manual_release_receipt_digest,previous_release_revision,content_hash) VALUES($1,1,'released',$2,$3,$3,$3,$4,$5,$6,$7,$3,NULL,$3)`, + [generation, `synthetic-${suffix}`, "9".repeat(64), "4".repeat(64), "5".repeat(64), "6".repeat(64), "7".repeat(64)] + ); + await owner.unsafe( + "INSERT INTO omi_memory.postgres_restore_admission_heads(database_generation_digest,release_revision) VALUES($1,1)", + [generation] + ); + const appPool: PostgresTransactionPool = { + withTransaction: (options, callback) => + pool.withTransaction(options, async (connection) => { + await connection.query({ + name: "conversation_chat_test.role", + text: "SET LOCAL ROLE omi_platform_application", + values: [], + }); + return callback(connection); + }), + }; + const runtime = createPostgresFirebaseConversationReadRuntime({ + authorization: { + pool: appPool, + project_id: project, + application_id: app, + runtime_mode: "deployed", + context_ttl_seconds: 60, + database_generation_digest: generation, + id_token_adapter: { + verification_source: "firebase_production", + async verifyIdToken() { + return { + aud: project, + iss: `https://securetoken.google.com/${project}`, + sub: uid, + uid, + iat: now() - 10, + auth_time: now() - 10, + exp: now() + 600, + }; + }, + }, + }, + codecRootSecret: new Uint8Array(32).fill(7), + cursorSigningKeyset: { + active_key_id: "test", + keys: [{ key_id: "test", secret: new Uint8Array(32).fill(8) }], + }, + }); + const call = (query = "") => + runtime.executeRequest( + new Request(`https://conversations.example/v1/conversations${query}`, { + headers: { authorization: "Bearer header.payload.signature" }, + }) + ); + const ids = (body: { items: Array<{ id: string }> }) => body.items.map((item) => item.id); + await owner.unsafe( + `INSERT INTO omi_memory.chat_messages(account_id,id,text,sender,message_type,created_at,updated_at,chat_session_id,app_id,journal_revision,payload_hash,message_source,rating,reported,server_revision,attachments_json,generation_id) VALUES($1,$2,'saved prompt','human','text',1000,1000,NULL,NULL,0,'sha256:human','desktop_chat',NULL,false,'rev-human','[]'::jsonb,'gen_human')`, + [account, "11111111-1111-4111-8111-111111111111"] + ); + const withoutGrant = (await (await call()).json()) as { items: Array<{ id: string }> }; + expect(ids(withoutGrant)).not.toContain("chat:chat-main"); + await owner.unsafe( + `INSERT INTO omi_memory.application_grant_revisions(account_id,application_id,credential_id,credential_generation,capability,grant_id,grant_version,lifecycle,enabled,scopes,record_schema_version,record_json,content_hash) VALUES($1,$2,$3,1,$4,$5,1,'active',true,'[]','grant-v1','{}',$6)`, + [account, app, credential, "chat.read", `chat.read-${suffix}`, "3".repeat(64)] + ); + await owner.unsafe( + `INSERT INTO omi_memory.application_grant_heads(account_id,application_id,credential_id,credential_generation,capability,grant_id,grant_version) VALUES($1,$2,$3,1,$4,$5,1)`, + [account, app, credential, "chat.read", `chat.read-${suffix}`] + ); + const emptyListen = (await (await call()).json()) as { + items: Array<{ id: string; title: string; source: string; status: string }>; + absence: unknown; + window: { nextCursor: string | null }; + }; + expect(emptyListen.absence).toBeNull(); + expect(emptyListen.items).toEqual([ + expect.objectContaining({ + id: "chat:chat-main", + title: "saved prompt", + overview: "saved prompt", + source: "chat", + status: "in_progress", + }), + ]); + const first = randomUUID(); + const second = randomUUID(); + for (const id of [first, second]) { + await owner.unsafe( + "INSERT INTO omi_memory.listen_capture_sessions(account_id,session_id,conversation_id,started_at,source,codec,sample_rate,channels,content_hash) VALUES($1,$2,$3,clock_timestamp()-interval '2 seconds','omi','21',16000,1,$4)", + [account, id, `conversation:${id}`, "1".repeat(64)] + ); + await owner.unsafe( + "INSERT INTO omi_memory.listen_capture_audio_uploads(account_id,session_id,capture_id,device_id,codec_id,upload_completed_at) VALUES($1,$2::text,$2::uuid,'synthetic-device',21,clock_timestamp())", + [account, id] + ); + await owner.unsafe( + "INSERT INTO omi_memory.listen_audio_transcriptions(account_id,session_id,state,attempts,available_at,updated_at,provider_result) VALUES($1,$2,'completed',1,clock_timestamp(),clock_timestamp(),$3::text::jsonb)", + [ + account, + id, + JSON.stringify({ + durationSeconds: 1, + segments: [{ text: `recording ${id}`, start: 0, end: 1, speaker: 0 }], + }), + ] + ); + } + const firstPage = (await (await call("?limit=1")).json()) as { + items: Array<{ id: string }>; + window: { nextCursor: string | null; hasMore: boolean }; + }; + expect(ids(firstPage)).toContain("chat:chat-main"); + expect(firstPage.items).toHaveLength(2); + expect(firstPage.window.hasMore).toBe(true); + expect(firstPage.window.nextCursor).not.toBeNull(); + const secondPage = (await ( + await call(`?limit=1&cursor=${encodeURIComponent(firstPage.window.nextCursor!)}`) + ).json()) as { items: Array<{ id: string }> }; + expect(ids(secondPage)).not.toContain("chat:chat-main"); + expect(secondPage.items).toHaveLength(1); + await owner.unsafe( + "DELETE FROM omi_memory.application_grant_heads WHERE account_id=$1 AND capability='chat.read'", + [account] + ); + const revoked = (await (await call("?limit=1")).json()) as { items: Array<{ id: string }> }; + expect(ids(revoked)).not.toContain("chat:chat-main"); + expect(revoked.items).toHaveLength(1); + } finally { + await pool.close(); + await owner.end(); + } + }, + 60000 +); diff --git a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.test.ts b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.test.ts index cc1d8dfcc6e..6e4fac825b4 100644 --- a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.test.ts +++ b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.test.ts @@ -148,3 +148,42 @@ test("deployed conversation collection shares Firebase admission and does not mo expect(await response.text()).toBe('{"error":"unauthorized"}'); expect((await app.request("/v1/conversations/recording/title",{method:"PATCH"})).status).toBe(404); }); + +test("deployed chat history shares Firebase admission and does not mount writes", async () => { + const base = options(); + const app = createPostgresFirebaseAuthorizedMemoryServiceApp({ + ...base, + chat: { + authorization: base.memory_read.authorization, + codecRootSecret: new Uint8Array(32).fill(7), + cursorSigningKeyset: { active_key_id: "test", keys: [{ key_id: "test", secret: new Uint8Array(32).fill(8) }] }, + }, + }); + const response = await app.request("/v1/chat-messages?limit=50", { + headers: { authorization: "Bearer invalid.token" }, + }); + expect(response.status).toBe(401); + expect(await response.json()).toEqual({ + error: { code: "unauthorized", retryable: false, action: "reauthenticate" }, + }); + expect((await app.request("/v1/chat-messages", { method: "POST", body: "{}" })).status).toBe(404); + expect(await (await app.request("/v1/chat-messages", { method: "POST", body: "{}" })).text()) + .toBe('{"error":"not_found"}'); +}); + +test("deployed settings verify identity without inventing a signed-in profile", async () => { + const base = options(); + const app = createPostgresFirebaseAuthorizedMemoryServiceApp({ + ...base, + settings: { authorization: base.memory_read.authorization }, + }); + const signedOut = await app.request("/v1/settings"); + expect(signedOut.status).toBe(200); + expect(await signedOut.json()).toEqual({ identity: null, entitlement: null }); + const denied = await app.request("/v1/settings", { + headers: { authorization: "Bearer invalid.token" }, + }); + expect(denied.status).toBe(401); + expect(await denied.text()).toBe('{"error":"unauthorized"}'); + expect((await app.request("/v1/settings", { method: "PATCH", body: "{}" })).status).toBe(404); +}); diff --git a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.ts b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.ts index b1b03e485e8..629d3650a27 100644 --- a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.ts +++ b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-app.ts @@ -1,4 +1,6 @@ import {createPostgresFirebaseConversationReadRuntime, type PostgresFirebaseConversationReadOptions} from "./firebase-conversation-read-runtime"; +import { createPostgresFirebaseChatReadRuntime, type PostgresFirebaseChatReadOptions } from "./firebase-chat-read-runtime"; +import { createPostgresFirebaseSettingsRuntime, type PostgresFirebaseSettingsOptions } from "./firebase-settings-runtime"; import { createPostgresFirebaseDeviceSessionRuntime } from "./firebase-device-session-runtime"; import type { PrerecordedTranscriptionSource } from "../../apps/service/listen/prerecorded-transcription"; import type { PostgresFirebaseAuthorizationRuntimeOptions } from "./firebase-authorized-runtime-support"; @@ -29,6 +31,8 @@ export interface PostgresFirebaseAuthorizedMemoryServiceAppOptions { readonly observability?: ServiceAppObservability; readonly tasks?: PostgresFirebaseTasksOptions; readonly conversations?: PostgresFirebaseConversationReadOptions; + readonly chat?: PostgresFirebaseChatReadOptions; + readonly settings?: PostgresFirebaseSettingsOptions; readonly device_sessions?: PostgresFirebaseAuthorizationRuntimeOptions; readonly device_ownership_key?: Uint8Array; readonly transcription_source?: PrerecordedTranscriptionSource; @@ -48,7 +52,7 @@ export const createPostgresFirebaseAuthorizedMemoryServiceApp = ( const descriptors = Object.getOwnPropertyDescriptors(options); const required = ["mcp_handler", "memory_read", "now_epoch_seconds", "counter"] as const; if (Reflect.ownKeys(descriptors).some((key) => - typeof key !== "string" || ![...required, "observability", "tasks", "device_sessions", "transcription_source", "device_ownership_key", "conversations"].includes(key)) + typeof key !== "string" || ![...required, "observability", "tasks", "device_sessions", "transcription_source", "device_ownership_key", "conversations", "chat", "settings"].includes(key)) || required.some((key) => !Object.hasOwn(descriptors, key)) || Object.values(descriptors).some((entry) => !entry.enumerable || !("value" in entry))) { throw new TypeError("invalid PostgreSQL Firebase memory service options"); @@ -73,5 +77,7 @@ export const createPostgresFirebaseAuthorizedMemoryServiceApp = ( descriptors.tasks ? createPostgresFirebaseTasksRuntime(descriptors.tasks.value as PostgresFirebaseTasksOptions) : undefined, descriptors.device_sessions ? createPostgresFirebaseDeviceSessionRuntime(descriptors.device_sessions.value as PostgresFirebaseAuthorizationRuntimeOptions, descriptors.transcription_source?.value as PrerecordedTranscriptionSource | undefined, descriptors.device_ownership_key?.value as Uint8Array | undefined) : undefined, descriptors.conversations ? createPostgresFirebaseConversationReadRuntime(descriptors.conversations.value as PostgresFirebaseConversationReadOptions) : undefined, + descriptors.chat ? createPostgresFirebaseChatReadRuntime(descriptors.chat.value as PostgresFirebaseChatReadOptions) : undefined, + descriptors.settings ? createPostgresFirebaseSettingsRuntime(descriptors.settings.value as PostgresFirebaseSettingsOptions) : undefined, ); }; diff --git a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.test.ts b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.test.ts index c9108b3216d..d2a960ec4be 100644 --- a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.test.ts +++ b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.test.ts @@ -350,3 +350,23 @@ test("task routes cannot use a pool outside the process readiness proof", () => } }, })).toThrow("invalid PostgreSQL Firebase memory service process options"); }); + +test("chat routes cannot use a pool outside the process readiness proof", () => { + const base = fixture(); + expect(() => createPostgresFirebaseAuthorizedMemoryServiceProcess({ ...base.options, + service_options: { ...base.service_options, chat: { + authorization: { ...base.service_options.memory_read.authorization, pool: fixture().pool }, + codecRootSecret: new Uint8Array(32).fill(1), + cursorSigningKeyset: { active_key_id: "test", keys: [{ key_id: "test", secret: new Uint8Array(32).fill(2) }] }, + } }, + })).toThrow("invalid PostgreSQL Firebase memory service process options"); +}); + +test("settings routes cannot use a pool outside the process readiness proof", () => { + const base = fixture(); + expect(() => createPostgresFirebaseAuthorizedMemoryServiceProcess({ ...base.options, + service_options: { ...base.service_options, settings: { + authorization: { ...base.service_options.memory_read.authorization, pool: fixture().pool }, + } }, + })).toThrow("invalid PostgreSQL Firebase memory service process options"); +}); diff --git a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.ts b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.ts index 1070c6d899f..325548d93f7 100644 --- a/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.ts +++ b/backends/example-platform/drivers/postgres/firebase-authorized-memory-service-process.ts @@ -113,17 +113,22 @@ const nestedAuthorization = (serviceOptions: unknown): Readonly<{ }> => { const service = exactRecordWithOptional(serviceOptions, [ "counter", "mcp_handler", "memory_read", "now_epoch_seconds", - ], ["observability", "tasks", "device_sessions", "transcription_source", "device_ownership_key", "conversations"]); + ], ["observability", "tasks", "device_sessions", "transcription_source", "device_ownership_key", "conversations", "chat", "settings"]); const memoryRead = exactRecord(service["memory_read"], ["authorization", "product"]); const authorization = exactRecord(memoryRead["authorization"], [ "application_id", "context_ttl_seconds", "database_generation_digest", "id_token_adapter", "pool", "project_id", "runtime_mode", ]); - for (const key of ["tasks", "device_sessions", "conversations"]) { + for (const key of ["tasks", "device_sessions", "conversations", "chat", "settings"]) { if (!Object.hasOwn(service, key)) continue; - const candidate = key !== "device_sessions" - ? exactRecord(service[key], ["authorization", "codecRootSecret", "cursorSigningKeyset"])["authorization"] - : service[key]; + const candidate = key === "device_sessions" + ? service[key] + : exactRecord( + service[key], + key === "settings" + ? ["authorization"] + : ["authorization", "codecRootSecret", "cursorSigningKeyset"], + )["authorization"]; const paired = exactRecord(candidate, Object.keys(authorization)); if (Object.keys(authorization).some(field => paired[field] !== authorization[field])) fail(); } diff --git a/backends/example-platform/drivers/postgres/firebase-authorized-runtime-support.ts b/backends/example-platform/drivers/postgres/firebase-authorized-runtime-support.ts index 6f4eb250840..bb73c92b7f1 100644 --- a/backends/example-platform/drivers/postgres/firebase-authorized-runtime-support.ts +++ b/backends/example-platform/drivers/postgres/firebase-authorized-runtime-support.ts @@ -60,11 +60,12 @@ const exactOptions = ( /** @internal Shared fixed-query authorization construction for PG read/write runtimes. */ export const createPostgresFirebaseAuthorizationRuntime = ( optionsValue: PostgresFirebaseAuthorizationRuntimeOptions, - capability: "memories.read" | "memories.write" | "memories.export" | "tasks.read" | "tasks.write" | "listen.capture.write" | "conversations.read", + capability: "memories.read" | "memories.write" | "memories.export" | "tasks.read" | "tasks.write" | "listen.capture.write" | "conversations.read" | "chat.read", ): PostgresFirebaseAuthorizationRuntimeBinding => { if (capability !== "memories.read" && capability !== "memories.write" && capability !== "memories.export" && capability !== "tasks.read" - && capability !== "tasks.write" && capability !== "listen.capture.write" && capability !== "conversations.read") { + && capability !== "tasks.write" && capability !== "listen.capture.write" + && capability !== "conversations.read" && capability !== "chat.read") { throw new TypeError("invalid PostgreSQL Firebase runtime capability"); } const options = exactOptions(optionsValue); diff --git a/backends/example-platform/drivers/postgres/firebase-chat-read-runtime.test.ts b/backends/example-platform/drivers/postgres/firebase-chat-read-runtime.test.ts new file mode 100644 index 00000000000..00724ef2c2b --- /dev/null +++ b/backends/example-platform/drivers/postgres/firebase-chat-read-runtime.test.ts @@ -0,0 +1,79 @@ +import { expect, test } from "bun:test"; +import { CHAT_CAPABILITIES } from "../../apps/service/routes/chat-messages"; +import type { PostgresTransactionPool } from "./connection"; +import { createPostgresFirebaseChatReadRuntime } from "./firebase-chat-read-runtime"; + +const authorization = (pool: PostgresTransactionPool) => ({ + pool, + project_id: "qa-project", + runtime_mode: "deployed" as const, + application_id: "qa-app", + context_ttl_seconds: 60, + database_generation_digest: "a".repeat(64), + id_token_adapter: { + verification_source: "firebase_production" as const, + verifyIdToken: async () => { + const now = Math.floor(Date.now() / 1000); + return { + aud: "qa-project", + iss: "https://securetoken.google.com/qa-project", + sub: "unmigrated-uid", + uid: "unmigrated-uid", + iat: now - 60, + auth_time: now - 60, + exp: now + 3600, + }; + }, + }, +}); + +const runtimeFor = (pool: PostgresTransactionPool) => createPostgresFirebaseChatReadRuntime({ + authorization: authorization(pool), + codecRootSecret: new Uint8Array(32).fill(7), + cursorSigningKeyset: { + active_key_id: "test", + keys: [{ key_id: "test", secret: new Uint8Array(32).fill(8) }], + }, +}); + +test("chat GET denies missing tokens and missing grants instead of returning empty history", async () => { + let queries = 0; + const pool: PostgresTransactionPool = { + withTransaction: async (_options, callback) => callback({ + connectionIdentity: {}, + query: async () => { + queries += 1; + return []; + }, + execute: async () => ({ rowCount: 0 }), + }), + }; + const runtime = runtimeFor(pool); + const request = (headers: HeadersInit = {}) => new Request( + "https://service.example/v1/chat-messages?limit=50", + { headers }, + ); + const missing = await runtime.executeRequest(request()); + expect(missing.status).toBe(401); + expect(await missing.json()).toEqual({ + error: { code: "unauthorized", retryable: false, action: "reauthenticate" }, + }); + expect(queries).toBe(0); + const denied = await runtime.executeRequest(request({ + authorization: "Bearer header.payload.signature", + })); + expect(denied.status).toBe(403); + expect(await denied.json()).toEqual({ + error: { code: "forbidden", retryable: false, action: "none" }, + }); + expect(queries).toBeGreaterThan(0); + expect((await runtime.executeRequest(new Request( + "https://service.example/v1/chat-messages", + { method: "POST", body: "{}" }, + ))).status).toBe(404); + expect(CHAT_CAPABILITIES.maxAttachmentsPerMessage).toBe(4); + const cancelled = new AbortController(); + cancelled.abort(); + expect((await runtime.executeRequest(new Request(request(), { signal: cancelled.signal }))).status) + .toBe(503); +}); diff --git a/backends/example-platform/drivers/postgres/firebase-chat-read-runtime.ts b/backends/example-platform/drivers/postgres/firebase-chat-read-runtime.ts new file mode 100644 index 00000000000..58c3c886737 --- /dev/null +++ b/backends/example-platform/drivers/postgres/firebase-chat-read-runtime.ts @@ -0,0 +1,172 @@ +// domain-pending(DIV-CHAT-SENDER-001) +// domain-pending(DIV-CHAT-TYPE-001) +// domain-pending(DIV-CHAT-SESSION-001) +// domain-pending(DIV-CHAT-REV-001) +// domain-pending(DIV-CHAT-HASH-001) +// domain-pending(DIV-CHAT-SOURCE-001) + +import { + ExpiredChatHistoryCursorError, + InvalidChatHistoryCursorError, + createChatHistoryCursorCodec, +} from "../../apps/service/chat/history-cursor"; +import { createServedCounter } from "../../apps/service/observability/served-count"; +import { + CHAT_CAPABILITIES, + parseHistoryQuery, + projectLoadedHistoryMessage, +} from "../../apps/service/routes/chat-messages"; +import type { McpCursorSigningKeyset } from "../../apps/mcp/cursor"; +import { + createPostgresFirebaseAuthorizationRuntime, + type PostgresFirebaseAuthorizationRuntimeOptions, +} from "./firebase-authorized-runtime-support"; +import { withAuthorizedChatRead } from "./chat-read-repository"; + +export interface PostgresFirebaseChatReadOptions { + readonly authorization: PostgresFirebaseAuthorizationRuntimeOptions; + readonly codecRootSecret: Uint8Array; + readonly cursorSigningKeyset: McpCursorSigningKeyset; +} + +const CHAT_CURSOR_TTL_SECONDS = 3_600; +const SERVICE_UNAVAILABLE_RETRY_AFTER_SECONDS = 60; +const JSON_HEADERS = Object.freeze({ + "cache-control": "no-store", + "content-type": "application/json", +}); + +const json = ( + body: unknown, + status: number, + extra: Readonly> = {}, +): Response => new Response( + typeof body === "string" ? body : JSON.stringify(body), + { status, headers: { ...JSON_HEADERS, ...extra } }, +); + +const errorResponse = ( + status: number, + code: string, + action: string, + retryable = false, + extra: Readonly> = {}, +): Response => json({ error: { code, retryable, action } }, status, extra); + +const unavailable = (): Response => errorResponse( + 503, + "service_unavailable", + "retry", + true, + { "retry-after": String(SERVICE_UNAVAILABLE_RETRY_AFTER_SECONDS) }, +); + +export function createPostgresFirebaseChatReadRuntime( + options: PostgresFirebaseChatReadOptions, +) { + const runtime = createPostgresFirebaseAuthorizationRuntime( + options.authorization, + "chat.read", + ); + void options.codecRootSecret; + const cursor = createChatHistoryCursorCodec({ + activeId: options.cursorSigningKeyset.active_key_id, + keys: options.cursorSigningKeyset.keys.map((key) => ({ + id: key.key_id, + secret: new Uint8Array(key.secret), + })), + }); + const counter = createServedCounter(); + return Object.freeze({ + async executeRequest(request: Request): Promise { + if (request.method !== "GET" || new URL(request.url).pathname !== "/v1/chat-messages") { + return json({ error: "not_found" }, 404); + } + try { + request.signal.throwIfAborted(); + const token = request.headers.get("authorization")?.match(/^Bearer (\S+)$/)?.[1] ?? ""; + const authorization = await runtime.authorizer.authorize( + token, + Math.floor(Date.now() / 1000), + ); + if (!authorization.authorized) { + if (authorization.outcome === "authentication") { + return errorResponse(401, "unauthorized", "reauthenticate"); + } + return authorization.outcome === "unavailable" + ? unavailable() + : errorResponse(403, "forbidden", "none"); + } + const query = parseHistoryQuery(request); + if (query === null) { + counter.recordDomainRead("denied"); + return errorResponse(400, "bad_request", "edit_request"); + } + const authority = authorization.context; + return await withAuthorizedChatRead( + runtime.pool, + authority, + request.signal, + async (storage) => { + const accountEpoch = authority.account_epoch; + const nowEpochSeconds = Math.floor(Date.now() / 1000); + try { + const claims = query.olderCursor === null ? null : cursor.verify(query.olderCursor, { + accountId: authority.account_id, + accountEpoch, + nowEpochSeconds, + }); + const snapshotSequence = claims?.snapshotSequence + ?? await storage.readSnapshotSequence(); + const page = await storage.listHistory({ + limit: query.limit, + snapshotSequence, + olderThan: claims?.olderThan ?? null, + }); + const messages = []; + for (const message of page.messages) { + const stored = message.sender === "ai" + ? await storage.readMessage(message.id) + : null; + const generationEvents = stored?.generationId + ? await storage.listGenerationEvents(stored.generationId) + : null; + messages.push(projectLoadedHistoryMessage(message, stored, generationEvents)); + } + const oldest = messages[0]; + const olderCursor = page.hasOlder && oldest !== undefined + ? cursor.issue({ + accountId: authority.account_id, + accountEpoch, + snapshotSequence, + olderThan: { createdAt: oldest.createdAt, id: oldest.id }, + issuedAtEpochSeconds: claims?.issuedAtEpochSeconds ?? nowEpochSeconds, + ttlSeconds: CHAT_CURSOR_TTL_SECONDS, + }) + : null; + counter.recordDomainRead("served"); + return json({ + messages, + page: { olderCursor, hasOlder: page.hasOlder }, + capabilities: CHAT_CAPABILITIES, + }, 200); + } catch (error) { + if (error instanceof ExpiredChatHistoryCursorError) { + counter.recordDomainRead("denied"); + return errorResponse(410, "cursor_expired", "refresh_history"); + } + if (error instanceof InvalidChatHistoryCursorError) { + counter.recordDomainRead("denied"); + return errorResponse(400, "bad_request", "refresh_history"); + } + throw error; + } + }, + ); + } catch { + counter.recordDomainRead("failed"); + return unavailable(); + } + }, + }); +} diff --git a/backends/example-platform/drivers/postgres/firebase-conversation-read-runtime.ts b/backends/example-platform/drivers/postgres/firebase-conversation-read-runtime.ts index 34611951d0f..90e11607540 100644 --- a/backends/example-platform/drivers/postgres/firebase-conversation-read-runtime.ts +++ b/backends/example-platform/drivers/postgres/firebase-conversation-read-runtime.ts @@ -1,6 +1,7 @@ import { createHash } from "node:crypto"; import { Hono } from "hono"; import { prepareConversationsRead } from "../../apps/service/composition/conversations-read"; +import { composeChatSessionsIntoConversationPage } from "../../apps/service/composition/chat-conversation-sessions"; import { parseConversationReadWindow, registerConversationReadRoutes, @@ -12,6 +13,7 @@ import { type PostgresFirebaseAuthorizationRuntimeOptions, } from "./firebase-authorized-runtime-support"; import { withAuthorizedConversationRead } from "./conversation-read-repository"; +import { withAuthorizedChatRead } from "./chat-read-repository"; export interface PostgresFirebaseConversationReadOptions { readonly authorization: PostgresFirebaseAuthorizationRuntimeOptions; @@ -30,6 +32,10 @@ export function createPostgresFirebaseConversationReadRuntime( options.authorization, "conversations.read" ); + const chatRuntime = createPostgresFirebaseAuthorizationRuntime( + options.authorization, + "chat.read" + ); const secret = new Uint8Array(options.codecRootSecret); const keys = { active_key_id: options.cursorSigningKeyset.active_key_id, @@ -65,7 +71,7 @@ export function createPostgresFirebaseConversationReadRuntime( const authority = authorization.context; const window = parseConversationReadWindow(request); if (window === null) return failed(400, "bad_request"); - return await withAuthorizedConversationRead( + const listen = await withAuthorizedConversationRead( runtime.pool, authority, request.signal, @@ -174,6 +180,38 @@ export function createPostgresFirebaseConversationReadRuntime( return response; } ); + if (listen.status !== 200 || window.legacy || window.cursor !== null) { + return listen; + } + const chatAuthorization = await chatRuntime.authorizer.authorize( + token, + Math.floor(Date.now() / 1000) + ); + if (!chatAuthorization.authorized) { + return chatAuthorization.outcome === "authorization" + ? listen + : failed(503, "unavailable"); + } + const sessions = await withAuthorizedChatRead( + chatRuntime.pool, + chatAuthorization.context, + request.signal, + (storage) => storage.listConversationSessions() + ); + if (sessions.length === 0) return listen; + const composed = composeChatSessionsIntoConversationPage( + await listen.json(), + sessions + ); + if (composed === null) return failed(503, "unavailable"); + return new Response(JSON.stringify(composed), { + status: 200, + headers: { + "cache-control": "no-store", + "content-type": + listen.headers.get("content-type") ?? "application/json", + }, + }); } catch { return failed(503, "unavailable"); } diff --git a/backends/example-platform/drivers/postgres/firebase-settings-runtime.test.ts b/backends/example-platform/drivers/postgres/firebase-settings-runtime.test.ts new file mode 100644 index 00000000000..296ec3b782a --- /dev/null +++ b/backends/example-platform/drivers/postgres/firebase-settings-runtime.test.ts @@ -0,0 +1,92 @@ +import { expect, test } from "bun:test"; +import { FIREBASE_ID_TOKEN_VERIFICATION_UNAVAILABLE } from "../../apps/service/auth/firebase-identity"; +import type { PostgresTransactionPool } from "./connection"; +import { createPostgresFirebaseSettingsRuntime } from "./firebase-settings-runtime"; + +const PROJECT = "qa-project"; +const SIGNED_TOKEN = "eyJhbGciOiJSUzI1NiJ9.eyJzdWIiOiJ1c2VyIn0.signature"; + +const claims = (now: number) => ({ + aud: PROJECT, + iss: `https://securetoken.google.com/${PROJECT}`, + sub: "firebase-user-1", + uid: "firebase-user-1", + exp: now + 3_600, + iat: now - 60, + auth_time: now - 120, + email: "private@example.invalid", + name: "Invented Profile", +}); + +const authorization = ( + verify: (token: string, checkRevoked: boolean) => Promise, + pool: PostgresTransactionPool, +) => ({ + pool, + project_id: PROJECT, + runtime_mode: "deployed" as const, + application_id: "qa-app", + context_ttl_seconds: 60, + database_generation_digest: "a".repeat(64), + id_token_adapter: { + verification_source: "firebase_production" as const, + verifyIdToken: verify, + }, +}); + +const unusedPool: PostgresTransactionPool = { + withTransaction: async () => { + throw new Error("settings must not query PostgreSQL until a producer exists"); + }, +}; + +test("signed-out settings stay null and verified identity stays unavailable without inventing a profile", async () => { + const calls: boolean[] = []; + const runtime = createPostgresFirebaseSettingsRuntime({ + authorization: authorization(async (_token, checkRevoked) => { + calls.push(checkRevoked); + return claims(Math.floor(Date.now() / 1000)); + }, unusedPool), + }); + const absent = await runtime.executeRequest(new Request("https://service.example/v1/settings")); + expect(absent.status).toBe(200); + expect(await absent.json()).toEqual({ identity: null, entitlement: null }); + expect(calls).toEqual([]); + + const verified = await runtime.executeRequest(new Request("https://service.example/v1/settings", { + headers: { authorization: `Bearer ${SIGNED_TOKEN}` }, + })); + const body = await verified.text(); + expect(verified.status).toBe(503); + expect(verified.headers.get("retry-after")).toBe("60"); + expect(body).toBe('{"error":"service_unavailable"}'); + expect(body).not.toContain("private@example.invalid"); + expect(body).not.toContain("Invented Profile"); + expect(calls).toEqual([false, true]); + + expect((await runtime.executeRequest(new Request( + "https://service.example/v1/settings", + { headers: { authorization: "Bearer invalid.token" } }, + ))).status).toBe(401); + expect((await runtime.executeRequest(new Request( + "https://service.example/v1/settings?appearance=dark", + { headers: { authorization: `Bearer ${SIGNED_TOKEN}` } }, + ))).status).toBe(400); + expect((await runtime.executeRequest(new Request( + "https://service.example/v1/settings", + { method: "POST", body: "{}" }, + ))).status).toBe(404); +}); + +test("Firebase verification outage is unavailable rather than a signed-out envelope", async () => { + const runtime = createPostgresFirebaseSettingsRuntime({ + authorization: authorization(async () => { + throw new Error(FIREBASE_ID_TOKEN_VERIFICATION_UNAVAILABLE); + }, unusedPool), + }); + const result = await runtime.executeRequest(new Request("https://service.example/v1/settings", { + headers: { authorization: `Bearer ${SIGNED_TOKEN}` }, + })); + expect(result.status).toBe(503); + expect(await result.text()).toBe('{"error":"service_unavailable"}'); +}); diff --git a/backends/example-platform/drivers/postgres/firebase-settings-runtime.ts b/backends/example-platform/drivers/postgres/firebase-settings-runtime.ts new file mode 100644 index 00000000000..cd1bb3e489c --- /dev/null +++ b/backends/example-platform/drivers/postgres/firebase-settings-runtime.ts @@ -0,0 +1,90 @@ +import { + createFirebaseIdentityVerifier, + isFirebaseIdentityRefreshUnavailable, +} from "../../apps/service/auth/firebase-identity"; +import type { PostgresFirebaseAuthorizationRuntimeOptions } from "./firebase-authorized-runtime-support"; + +export interface PostgresFirebaseSettingsOptions { + readonly authorization: PostgresFirebaseAuthorizationRuntimeOptions; +} + +const JSON_HEADERS = Object.freeze({ + "cache-control": "no-store", + "content-type": "application/json", +}); +const UNAVAILABLE_HEADERS = Object.freeze({ + ...JSON_HEADERS, + "retry-after": "60", +}); +const SIGNED_OUT_BODY = '{"identity":null,"entitlement":null}'; +const UNAUTHORIZED_BODY = '{"error":"unauthorized"}'; +const BAD_REQUEST_BODY = '{"error":"bad_request"}'; +const UNAVAILABLE_BODY = '{"error":"service_unavailable"}'; +const NOT_FOUND_BODY = '{"error":"not_found"}'; + +const response = ( + body: string, + status: number, + headers: Readonly> = JSON_HEADERS, +): Response => new Response(body, { status, headers }); + +const bearerToken = (presentHeader: string): string | null => { + if (!presentHeader.startsWith("Bearer ")) return null; + const token = presentHeader.slice("Bearer ".length); + return token.length > 0 ? token : null; +}; + +const hasInvalidRequestGrammar = (request: Request): boolean => { + let url: URL; + try { + url = new URL(request.url); + } catch { + return true; + } + if ([...url.searchParams].length > 0) return true; + const contentLength = request.headers.get("content-length"); + if (contentLength !== null && contentLength !== "0") return true; + return request.headers.has("transfer-encoding"); +}; + +export function createPostgresFirebaseSettingsRuntime( + options: PostgresFirebaseSettingsOptions, +) { + const authorization = options.authorization; + const verifier = createFirebaseIdentityVerifier({ + project_id: authorization.project_id, + runtime_mode: authorization.runtime_mode, + adapter: authorization.id_token_adapter, + }); + return Object.freeze({ + async executeRequest(request: Request): Promise { + if (request.method !== "GET" || new URL(request.url).pathname !== "/v1/settings") { + return response(NOT_FOUND_BODY, 404); + } + if (hasInvalidRequestGrammar(request)) { + return response(BAD_REQUEST_BODY, 400); + } + const header = request.headers.get("authorization"); + if (header === null) { + return response(SIGNED_OUT_BODY, 200); + } + const token = bearerToken(header); + if (token === null) { + return response(UNAUTHORIZED_BODY, 401); + } + try { + request.signal.throwIfAborted(); + const identity = await verifier.resolve(token, Math.floor(Date.now() / 1000)); + if (identity === null) { + return response(UNAUTHORIZED_BODY, 401); + } + if (isFirebaseIdentityRefreshUnavailable(identity)) { + return response(UNAVAILABLE_BODY, 503, UNAVAILABLE_HEADERS); + } + return response(UNAVAILABLE_BODY, 503, UNAVAILABLE_HEADERS); + } catch { + return response(UNAVAILABLE_BODY, 503, UNAVAILABLE_HEADERS); + } + }, + }); +} diff --git a/backends/example-platform/drivers/postgres/migrations/0055-chat-messages.sql b/backends/example-platform/drivers/postgres/migrations/0055-chat-messages.sql new file mode 100644 index 00000000000..2387ef40c00 --- /dev/null +++ b/backends/example-platform/drivers/postgres/migrations/0055-chat-messages.sql @@ -0,0 +1,273 @@ +CREATE TABLE omi_memory.chat_messages ( + sequence bigint GENERATED ALWAYS AS IDENTITY (INCREMENT BY 1 MINVALUE 1 MAXVALUE 9007199254740991), + account_id text NOT NULL REFERENCES omi_memory.platform_accounts(account_id), + id text NOT NULL CHECK (id <> ''), + text text NOT NULL, + sender text NOT NULL CHECK (sender <> ''), + message_type text NOT NULL CHECK (message_type <> ''), + created_at bigint NOT NULL CHECK (created_at >= 0), + updated_at bigint NOT NULL CHECK (updated_at >= 0), + chat_session_id text, + app_id text, + journal_revision bigint NOT NULL CHECK (journal_revision >= 0), + payload_hash text NOT NULL CHECK (payload_hash <> ''), + message_source text NOT NULL CHECK (message_source <> ''), + rating double precision, + reported boolean NOT NULL, + server_revision text, + attachments_json jsonb CHECK (attachments_json IS NULL OR jsonb_typeof(attachments_json)='array'), + generation_id text, + PRIMARY KEY (account_id, id), + UNIQUE (account_id, sequence) +); +CREATE INDEX chat_messages_history ON omi_memory.chat_messages ( + account_id, app_id, chat_session_id, created_at DESC, id DESC, sequence +); +CREATE TABLE omi_memory.chat_generation_events ( + account_id text NOT NULL REFERENCES omi_memory.platform_accounts(account_id), + generation_id text NOT NULL CHECK (generation_id <> ''), + sequence bigint NOT NULL CHECK (sequence > 0 AND sequence <= 9007199254740991), + event_id text NOT NULL CHECK (event_id <> ''), + created_at bigint NOT NULL CHECK (created_at >= 0), + frame_json jsonb NOT NULL, + PRIMARY KEY (account_id, generation_id, sequence), + UNIQUE (account_id, generation_id, event_id) +); +REVOKE ALL ON omi_memory.chat_messages FROM PUBLIC, omi_platform_application; +REVOKE ALL ON omi_memory.chat_generation_events FROM PUBLIC, omi_platform_application; +ALTER TABLE omi_memory.chat_messages ENABLE ROW LEVEL SECURITY; +ALTER TABLE omi_memory.chat_generation_events ENABLE ROW LEVEL SECURITY; +CREATE POLICY chat_messages_select ON omi_memory.chat_messages FOR SELECT TO omi_platform_application + USING (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true) IN ('chat.read','chat.write')); +CREATE POLICY chat_messages_insert ON omi_memory.chat_messages FOR INSERT TO omi_platform_application + WITH CHECK (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true)='chat.write'); +CREATE POLICY chat_messages_update ON omi_memory.chat_messages FOR UPDATE TO omi_platform_application + USING (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true)='chat.write') + WITH CHECK (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true)='chat.write'); +CREATE POLICY chat_generation_events_select ON omi_memory.chat_generation_events FOR SELECT TO omi_platform_application + USING (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true) IN ('chat.read','chat.write')); +CREATE POLICY chat_generation_events_insert ON omi_memory.chat_generation_events FOR INSERT TO omi_platform_application + WITH CHECK (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true)='chat.write'); +CREATE POLICY chat_generation_events_update ON omi_memory.chat_generation_events FOR UPDATE TO omi_platform_application + USING (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true)='chat.write') + WITH CHECK (account_id=current_setting('omi.account_id',true) AND current_setting('omi.capability',true)='chat.write'); + +CREATE FUNCTION omi_memory.read_chat_snapshot_sequence() +RETURNS bigint LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,omi_memory AS $function$ +DECLARE v_account text:=nullif(current_setting('omi.account_id',true),''); v_sequence bigint; +BEGIN + IF v_account IS NULL OR nullif(current_setting('omi.principal_id',true),'') IS NULL + OR current_setting('omi.capability',true) IS DISTINCT FROM 'chat.read' THEN + RAISE EXCEPTION USING ERRCODE='P1005',MESSAGE='chat_authority_denied'; + END IF; + SELECT coalesce(max(sequence),0) INTO v_sequence FROM omi_memory.chat_messages WHERE account_id=v_account; + RETURN v_sequence; +END $function$; +REVOKE ALL ON FUNCTION omi_memory.read_chat_snapshot_sequence() FROM PUBLIC; +GRANT EXECUTE ON FUNCTION omi_memory.read_chat_snapshot_sequence() TO omi_platform_application; + +CREATE FUNCTION omi_memory.read_chat_history(p_limit integer,p_snapshot_sequence bigint,p_older_created_at bigint,p_older_id text) +RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,omi_memory AS $function$ +DECLARE v_account text:=nullif(current_setting('omi.account_id',true),''); v_page jsonb; +BEGIN + IF v_account IS NULL OR nullif(current_setting('omi.principal_id',true),'') IS NULL + OR current_setting('omi.capability',true) IS DISTINCT FROM 'chat.read' THEN + RAISE EXCEPTION USING ERRCODE='P1005',MESSAGE='chat_authority_denied'; + END IF; + IF p_limit IS NULL OR p_limit<1 OR p_limit>100 + OR p_snapshot_sequence IS NULL OR p_snapshot_sequence<0 + OR (p_older_created_at IS NULL) IS DISTINCT FROM (p_older_id IS NULL) + OR (p_older_id IS NOT NULL AND p_older_id='') + OR (p_older_created_at IS NOT NULL AND p_older_created_at<0) THEN + RAISE EXCEPTION USING ERRCODE='P1002',MESSAGE='chat_history_invalid'; + END IF; + WITH selected AS MATERIALIZED ( + SELECT m.id,m.text,m.sender,m.message_type,m.created_at,m.updated_at,m.chat_session_id,m.app_id, + m.journal_revision,m.payload_hash,m.message_source,m.rating,m.reported,m.server_revision,m.attachments_json,m.generation_id + FROM omi_memory.chat_messages m + WHERE m.account_id=v_account AND m.app_id IS NULL AND m.chat_session_id IS NULL + AND m.sequence<=p_snapshot_sequence + AND (p_older_created_at IS NULL OR m.created_atp_limit, + 'messages',coalesce(( + SELECT jsonb_agg(jsonb_build_object( + 'id',s.id,'text',s.text,'sender',s.sender,'type',s.message_type,'createdAt',s.created_at,'updatedAt',s.updated_at, + 'chatSessionId',s.chat_session_id,'appId',s.app_id,'journalRevision',s.journal_revision,'payloadHash',s.payload_hash, + 'messageSource',s.message_source,'rating',s.rating,'reported',s.reported,'revision',s.server_revision, + 'attachments',coalesce(s.attachments_json,'[]'::jsonb),'generationId',s.generation_id + ) ORDER BY s.created_at DESC, s.id DESC) + FROM (SELECT * FROM selected ORDER BY created_at DESC, id DESC LIMIT p_limit) s + ),'[]'::jsonb) + ) INTO v_page; + RETURN v_page; +END $function$; +REVOKE ALL ON FUNCTION omi_memory.read_chat_history(integer,bigint,bigint,text) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION omi_memory.read_chat_history(integer,bigint,bigint,text) TO omi_platform_application; + +CREATE FUNCTION omi_memory.read_chat_message(p_id text) +RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,omi_memory AS $function$ +DECLARE v_account text:=nullif(current_setting('omi.account_id',true),''); v_row jsonb; +BEGIN + IF v_account IS NULL OR nullif(current_setting('omi.principal_id',true),'') IS NULL + OR current_setting('omi.capability',true) IS DISTINCT FROM 'chat.read' THEN + RAISE EXCEPTION USING ERRCODE='P1005',MESSAGE='chat_authority_denied'; + END IF; + IF p_id IS NULL OR p_id='' THEN + RAISE EXCEPTION USING ERRCODE='P1002',MESSAGE='chat_message_invalid'; + END IF; + SELECT jsonb_build_object( + 'id',m.id,'text',m.text,'sender',m.sender,'type',m.message_type,'createdAt',m.created_at,'updatedAt',m.updated_at, + 'chatSessionId',m.chat_session_id,'appId',m.app_id,'journalRevision',m.journal_revision,'payloadHash',m.payload_hash, + 'messageSource',m.message_source,'rating',m.rating,'reported',m.reported,'revision',m.server_revision, + 'attachments',coalesce(m.attachments_json,'[]'::jsonb),'generationId',m.generation_id + ) INTO v_row + FROM omi_memory.chat_messages m + WHERE m.account_id=v_account AND m.id=p_id; + RETURN v_row; +END $function$; +REVOKE ALL ON FUNCTION omi_memory.read_chat_message(text) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION omi_memory.read_chat_message(text) TO omi_platform_application; + +CREATE FUNCTION omi_memory.read_chat_generation_events(p_generation_id text) +RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,omi_memory AS $function$ +DECLARE v_account text:=nullif(current_setting('omi.account_id',true),''); v_rows jsonb; +BEGIN + IF v_account IS NULL OR nullif(current_setting('omi.principal_id',true),'') IS NULL + OR current_setting('omi.capability',true) IS DISTINCT FROM 'chat.read' THEN + RAISE EXCEPTION USING ERRCODE='P1005',MESSAGE='chat_authority_denied'; + END IF; + IF p_generation_id IS NULL OR p_generation_id='' THEN + RAISE EXCEPTION USING ERRCODE='P1002',MESSAGE='chat_generation_invalid'; + END IF; + SELECT coalesce(jsonb_agg(jsonb_build_object( + 'id',e.event_id,'generationId',e.generation_id,'sequence',e.sequence,'createdAt',e.created_at,'frame',e.frame_json + ) ORDER BY e.sequence),'[]'::jsonb) INTO v_rows + FROM omi_memory.chat_generation_events e + WHERE e.account_id=v_account AND e.generation_id=p_generation_id; + RETURN v_rows; +END $function$; +REVOKE ALL ON FUNCTION omi_memory.read_chat_generation_events(text) FROM PUBLIC; +GRANT EXECUTE ON FUNCTION omi_memory.read_chat_generation_events(text) TO omi_platform_application; + +CREATE OR REPLACE FUNCTION omi_memory.cleanup_surface_tables(p_surface text) +RETURNS TABLE(table_name text) +LANGUAGE sql +IMMUTABLE +SECURITY DEFINER +SET search_path = pg_catalog, omi_memory +AS $function$ + SELECT mapping.table_name + FROM (VALUES + ('product_projections', 'listen_conversation_cursor_positions'), + ('product_projections', 'listen_conversation_read_revisions'), + ('product_projections', 'chat_messages'), + ('product_projections', 'chat_generation_events'), + ('staged_results', 'listen_audio_transcriptions'), + ('staged_results', 'listen_capture_audio_uploads'), + ('staged_results', 'listen_capture_audio_chunks'), + ('product_projections', 'task_records'), + ('product_projections', 'task_sequences'), + ('product_projections', 'task_write_receipts'), + ('product_projections', 'task_stragglers'), + ('product_projections', 'memory_render_responses'), + ('durable_work', 'memory_work_acceptances'), + ('durable_work', 'memory_work_execution_policies'), + ('durable_work', 'memory_work_heads'), + ('durable_work', 'memory_work_input_manifest'), + ('durable_work', 'memory_work_outbox_events'), + ('durable_work', 'memory_work_state_revisions'), + ('durable_work', 'memory_work_success_results'), + ('staged_results', 'memory_work_staged_results'), + ('staged_results', 'memory_formation_work_inputs'), + ('staged_results', 'memory_predicate_batch_work_inputs'), + ('staged_results', 'memory_derived_group_dream_work_inputs'), + ('staged_results', 'memory_query_evaluation_inputs'), + ('staged_results', 'memory_listen_attribution_belief_inputs'), + ('staged_results', 'memory_candidate_derivation_artifacts'), + ('staged_results', 'listen_capture_sessions'), + ('staged_results', 'listen_capture_session_state_revisions'), + ('staged_results', 'listen_capture_segments'), + ('staged_results', 'listen_formation_finalizations'), + ('staged_results', 'listen_conversation_finalization_intents'), + ('staged_results', 'listen_formation_outbox'), + ('staged_results', 'listen_formation_delivery_revisions'), + ('staged_results', 'listen_formation_delivery_heads'), + ('authoritative_memory', 'memory_attribution_belief_revisions'), + ('authoritative_memory', 'memory_claim_evidence_refs'), + ('authoritative_memory', 'memory_claim_lineages'), + ('authoritative_memory', 'memory_claim_liveness_fences'), + ('authoritative_memory', 'memory_claim_predicate_refs'), + ('authoritative_memory', 'memory_claim_revisions'), + ('authoritative_memory', 'memory_claim_source_provisionals'), + ('authoritative_memory', 'memory_claim_supersessions'), + ('authoritative_memory', 'memory_consumed_markers'), + ('authoritative_memory', 'memory_coreference_support_evidence_refs'), + ('authoritative_memory', 'memory_coreference_support_revisions'), + ('authoritative_memory', 'memory_derivation_attempts'), + ('authoritative_memory', 'memory_derivation_commits'), + ('authoritative_memory', 'memory_derivation_inputs'), + ('authoritative_memory', 'memory_entity_identities'), + ('authoritative_memory', 'memory_entity_revisions'), + ('authoritative_memory', 'memory_event_identities'), + ('authoritative_memory', 'memory_event_revisions'), + ('authoritative_memory', 'memory_evidence_identities'), + ('authoritative_memory', 'memory_evidence_revisions'), + ('authoritative_memory', 'memory_formation_extraction_evidence'), + ('authoritative_memory', 'memory_formation_extraction_outcomes'), + ('authoritative_memory', 'memory_formation_placement_outcomes'), + ('authoritative_memory', 'memory_formation_outcomes'), + ('authoritative_memory', 'memory_generated_adjacency'), + ('authoritative_memory', 'memory_graph_heads'), + ('authoritative_memory', 'memory_idempotency_receipts'), + ('authoritative_memory', 'memory_identity_authorization_identities'), + ('authoritative_memory', 'memory_identity_authorization_entity_endpoints'), + ('authoritative_memory', 'memory_identity_authorization_revisions'), + ('authoritative_memory', 'memory_identity_authorization_support'), + ('authoritative_memory', 'memory_identity_constraint_entity_endpoints'), + ('authoritative_memory', 'memory_identity_revisions'), + ('authoritative_memory', 'memory_identity_support'), + ('authoritative_memory', 'memory_mention_revisions'), + ('authoritative_memory', 'memory_placement_artifacts'), + ('authoritative_memory', 'memory_predicate_assertion_revisions'), + ('authoritative_memory', 'memory_predicate_identities'), + ('authoritative_memory', 'memory_predicate_revisions'), + ('authoritative_memory', 'memory_revisions'), + ('authoritative_memory', 'memory_source_local_claim_roles'), + ('account_access', 'application_credential_heads'), + ('account_access', 'application_credential_revisions'), + ('account_access', 'application_grant_heads'), + ('account_access', 'application_grant_revisions'), + ('account_access', 'firebase_application_credential_bindings'), + ('account_access', 'firebase_identity_bindings'), + ('experiment_results', 'memory_strategy_assignment_bundles'), + ('experiment_results', 'memory_strategy_assignment_policies'), + ('experiment_results', 'memory_strategy_baseline_read_groundings'), + ('experiment_results', 'memory_strategy_candidate_read_groundings'), + ('experiment_results', 'memory_strategy_definitions'), + ('experiment_results', 'memory_strategy_evaluation_baselines'), + ('experiment_results', 'memory_strategy_evaluation_pairs'), + ('experiment_results', 'memory_strategy_policy_shadows'), + ('experiment_results', 'memory_strategy_shadow_assignments'), + ('experiment_results', 'memory_strategy_shadow_results'), + ('product_projections', 'memory_product_membership_claim_lineages'), + ('product_projections', 'memory_product_membership_revisions'), + ('product_projections', 'memory_product_operation_receipts'), + ('product_projections', 'memory_product_projection_citation_evidence_refs'), + ('product_projections', 'memory_product_projection_citations'), + ('product_projections', 'memory_product_projection_payloads'), + ('product_projections', 'memory_product_projection_revisions'), + ('product_projections', 'memory_product_propositions'), + ('product_projections', 'memory_product_redirect_successors'), + ('product_projections', 'memory_product_redirects'), + ('rebuildable_groups_indexes', 'memory_product_group_members'), + ('rebuildable_groups_indexes', 'memory_product_group_projections'), + ('migration_state', 'memory_legacy_proposition_mappings'), + ('migration_state', 'memory_migration_item_tombstones') + ) AS mapping(surface, table_name) + WHERE mapping.surface = p_surface + ORDER BY mapping.table_name +$function$; diff --git a/backends/example-platform/drivers/postgres/migrations/0056-chat-conversation-sessions.sql b/backends/example-platform/drivers/postgres/migrations/0056-chat-conversation-sessions.sql new file mode 100644 index 00000000000..9de798fb22a --- /dev/null +++ b/backends/example-platform/drivers/postgres/migrations/0056-chat-conversation-sessions.sql @@ -0,0 +1,51 @@ +CREATE FUNCTION omi_memory.read_chat_conversation_sessions() +RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,omi_memory AS $function$ +DECLARE v_account text:=nullif(current_setting('omi.account_id',true),''); v_sessions jsonb; +BEGIN + IF v_account IS NULL OR nullif(current_setting('omi.principal_id',true),'') IS NULL + OR current_setting('omi.capability',true) IS DISTINCT FROM 'chat.read' THEN + RAISE EXCEPTION USING ERRCODE='P1005',MESSAGE='chat_authority_denied'; + END IF; + SELECT coalesce(jsonb_agg(to_jsonb(session) ORDER BY session."updatedAt" DESC, session.id),'[]'::jsonb) + INTO v_sessions + FROM ( + SELECT + 'chat:chat-main'::text AS id, + CASE + WHEN length(btrim(title_text))=0 THEN 'Chat' + WHEN char_length(btrim(title_text))>240 THEN left(btrim(title_text),237)||'...' + ELSE btrim(title_text) + END AS title, + CASE + WHEN length(btrim(last_text))=0 THEN 'Chat' + WHEN char_length(btrim(last_text))>240 THEN left(btrim(last_text),237)||'...' + ELSE btrim(last_text) + END AS overview, + created_at AS "createdAt", + updated_at AS "updatedAt", + created_at AS "startedAt", + NULL::bigint AS "finishedAt", + 'chat'::text AS source, + 'in_progress'::text AS status, + false AS discarded, + false AS starred, + 'private'::text AS visibility, + false AS "isLocked", + NULL::text AS "folderId", + NULL::text AS revision + FROM ( + SELECT + min(m.created_at) AS created_at, + max(m.created_at) AS updated_at, + (array_agg(m.text ORDER BY CASE WHEN m.sender='human' THEN 0 ELSE 1 END, m.created_at, m.id))[1] AS title_text, + (array_agg(m.text ORDER BY m.created_at DESC, m.id DESC))[1] AS last_text, + count(*) AS n + FROM omi_memory.chat_messages m + WHERE m.account_id=v_account AND m.app_id IS NULL AND m.chat_session_id IS NULL + ) AS aggregated + WHERE n>0 + ) AS session; + RETURN v_sessions; +END $function$; +REVOKE ALL ON FUNCTION omi_memory.read_chat_conversation_sessions() FROM PUBLIC; +GRANT EXECUTE ON FUNCTION omi_memory.read_chat_conversation_sessions() TO omi_platform_application; diff --git a/backends/example-platform/drivers/postgres/migrations/manifest.test.ts b/backends/example-platform/drivers/postgres/migrations/manifest.test.ts index ab5f522af6b..12e6b44476d 100644 --- a/backends/example-platform/drivers/postgres/migrations/manifest.test.ts +++ b/backends/example-platform/drivers/postgres/migrations/manifest.test.ts @@ -14,7 +14,7 @@ describe("PostgreSQL migration manifest", () => { .toEqual([ 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20, 21, 22, 23, 24, 25, 26, 27, 28, 29, 30, 31, 32, 33, 34, 35, 36, 37, 38, 39, 40, 41, 42, - 43, 44, 45, 46, 47, 48, 49, 50, 51, 52, 53, 54, + 43, 44, 45, 46, 47, 48, 49, 50, 51, 52, 53, 54, 55, 56, ]); expect(new Set(POSTGRES_MIGRATIONS.map((migration) => migration.name)).size) .toBe(POSTGRES_MIGRATIONS.length); diff --git a/backends/example-platform/drivers/postgres/migrations/manifest.ts b/backends/example-platform/drivers/postgres/migrations/manifest.ts index 3878da40e9c..4a45fa40bef 100644 --- a/backends/example-platform/drivers/postgres/migrations/manifest.ts +++ b/backends/example-platform/drivers/postgres/migrations/manifest.ts @@ -286,4 +286,6 @@ export const POSTGRES_MIGRATIONS: readonly PostgresMigrationManifestEntry[] = Ob Object.freeze({version:52,name:"listen-audio-batches",fileName:"0052-listen-audio-batches.sql",sha256:"8c26e0de79a446a2f3ffe097c30bed1422d3e760c7b0c28faeffee4b2cae2a08"}), Object.freeze({version:53,name:"listen-conversation-pagination",fileName:"0053-listen-conversation-pagination.sql",sha256:"e4be81ab2cc431e224616146189bcdcc603f09394e5d4c7013004aa464f467f1"}), Object.freeze({version:54,name:"listen-capture-time",fileName:"0054-listen-capture-time.sql",sha256:"41ea68570b2de03e25e0cfa82cc6685e4e7ab14fb2a2111880a46e6f4567978e"}), + Object.freeze({version:55,name:"chat-messages",fileName:"0055-chat-messages.sql",sha256:"c5b853f1bd658e6060a83024257badfad888a0a37eab7adb8367f427fe4626f5"}), + Object.freeze({version:56,name:"chat-conversation-sessions",fileName:"0056-chat-conversation-sessions.sql",sha256:"c9b7fd96f98baf7ad6e9049dd89dd63cfbd9dcc2c87c77f6c1a5fdb894f4ba92"}), ]); diff --git a/backends/example-platform/drivers/postgres/migrations/schema-static.test.ts b/backends/example-platform/drivers/postgres/migrations/schema-static.test.ts index e241d232be8..d296fa43d11 100644 --- a/backends/example-platform/drivers/postgres/migrations/schema-static.test.ts +++ b/backends/example-platform/drivers/postgres/migrations/schema-static.test.ts @@ -155,6 +155,8 @@ const expectedTables = [ "listen_capture_audio_uploads", "listen_conversation_read_revisions", "listen_conversation_cursor_positions", + "chat_generation_events", + "chat_messages", "memory_render_responses", "task_records", "task_sequences", diff --git a/backends/example-platform/migration/postgres/deletion-surface-registry.ts b/backends/example-platform/migration/postgres/deletion-surface-registry.ts index 5240e449283..6ee0631c22e 100644 --- a/backends/example-platform/migration/postgres/deletion-surface-registry.ts +++ b/backends/example-platform/migration/postgres/deletion-surface-registry.ts @@ -56,6 +56,7 @@ export const POSTGRES_DELETION_SURFACE_TABLES = Object.freeze({ ]), product_projections: Object.freeze([ "memory_render_responses", "listen_conversation_read_revisions", "listen_conversation_cursor_positions", + "chat_messages", "chat_generation_events", "task_records", "task_sequences", "task_write_receipts", "task_stragglers", "memory_product_membership_claim_lineages", "memory_product_membership_revisions", "memory_product_operation_receipts", "memory_product_projection_citation_evidence_refs", diff --git a/backends/example-platform/package.json b/backends/example-platform/package.json index f123c9210b3..e24238201d4 100644 --- a/backends/example-platform/package.json +++ b/backends/example-platform/package.json @@ -28,7 +28,7 @@ "build:production-deps": "bun run scripts/build-production-dependency-artifact.ts", "app": "bun run apps/service/bin/dev-server.ts", "logs": "bun integration/dev-logs.ts", - "check:deployed": "bun test scripts/prod-local.test.ts scripts/prod-local-identity.test.ts scripts/prod-local-identity-e2e.test.ts scripts/prod-local-identity-seed.test.ts integration/lib/supervise.test.ts scripts/lint-import-graph.test.ts drivers/sqlite/service-stores/account-lifecycle.test.ts migration/postgres/deletion-surface-registry.schema.test.ts migration/postgres/account-deletion-cleanup-participant.test.ts apps/service/codecs/capture-ownership.test.ts drivers/postgres/migrations/schema-static.test.ts drivers/postgres/firebase-chat-generation-context-source.test.ts apps/service/chat/generation-context.test.ts core/retrieve/temporal.test.ts apps/service/workers/formation-work-producer.test.ts drivers/postgres/conversation-read-repository.test.ts drivers/postgres/conversation-read-projection.test.ts apps/service/listen/prerecorded-transcription.test.ts apps/service/listen/device-transcription.test.ts apps/service/bin/production-server.test.ts drivers/model/device-audio.test.ts drivers/model/deepgram-transcription.test.ts drivers/postgres/firebase-authorized-memory-service-app.test.ts apps/service/stores/device-session-upload.test.ts drivers/postgres/firebase-device-session-runtime.test.ts drivers/postgres/listen-finalization-repository.test.ts drivers/postgres/postgresjs.test.ts deploy/gcp/dev/release-candidate.test.ts drivers/postgres/migrations/manifest.test.ts scripts/postgres-test-lifecycle.test.ts scripts/postgres-test-resources.test.ts drivers/model/persisted-render.test.ts scripts/bind-firebase-identity.test.ts drivers/model/http-render.test.ts apps/service/composition/deployed-config.test.ts drivers/postgres/firebase-authorized-memory-service-process.test.ts drivers/postgres/firebase-authorized-memory-read-runtime.test.ts drivers/postgres/firebase-memory-route-read-port.test.ts apps/service/memory-service-app.test.ts && bun run lint:imports && bun scripts/trace-value-imports.ts apps/service/bin/production-server.ts --forbid apps/qa --forbid drivers/sqlite --forbid drivers/model/glm --forbid integration/local-test-gateway --forbid harness/ --forbid spikes/ --forbid migration/", + "check:deployed": "bun test scripts/prod-local.test.ts scripts/prod-local-identity.test.ts scripts/prod-local-identity-e2e.test.ts scripts/prod-local-identity-seed.test.ts integration/lib/supervise.test.ts scripts/lint-import-graph.test.ts drivers/sqlite/service-stores/account-lifecycle.test.ts migration/postgres/deletion-surface-registry.schema.test.ts migration/postgres/account-deletion-cleanup-participant.test.ts apps/service/codecs/capture-ownership.test.ts drivers/postgres/migrations/schema-static.test.ts drivers/postgres/firebase-chat-generation-context-source.test.ts apps/service/chat/generation-context.test.ts core/retrieve/temporal.test.ts apps/service/workers/formation-work-producer.test.ts drivers/postgres/conversation-read-repository.test.ts drivers/postgres/chat-read-repository.test.ts drivers/postgres/firebase-chat-read-runtime.test.ts drivers/postgres/firebase-settings-runtime.test.ts drivers/postgres/chat-messages.real.test.ts apps/service/composition/chat-conversation-sessions.test.ts drivers/postgres/conversation-read-projection.test.ts apps/service/listen/prerecorded-transcription.test.ts apps/service/listen/device-transcription.test.ts apps/service/bin/production-server.test.ts drivers/model/device-audio.test.ts drivers/model/deepgram-transcription.test.ts drivers/postgres/firebase-authorized-memory-service-app.test.ts apps/service/stores/device-session-upload.test.ts drivers/postgres/firebase-device-session-runtime.test.ts drivers/postgres/listen-finalization-repository.test.ts drivers/postgres/postgresjs.test.ts deploy/gcp/dev/release-candidate.test.ts drivers/postgres/migrations/manifest.test.ts scripts/postgres-test-lifecycle.test.ts scripts/postgres-test-resources.test.ts drivers/model/persisted-render.test.ts scripts/bind-firebase-identity.test.ts drivers/model/http-render.test.ts apps/service/composition/deployed-config.test.ts drivers/postgres/firebase-authorized-memory-service-process.test.ts drivers/postgres/firebase-authorized-memory-read-runtime.test.ts drivers/postgres/firebase-memory-route-read-port.test.ts apps/service/memory-service-app.test.ts && bun run lint:imports && bun scripts/trace-value-imports.ts apps/service/bin/production-server.ts --forbid apps/qa --forbid drivers/sqlite --forbid drivers/model/glm --forbid integration/local-test-gateway --forbid harness/ --forbid spikes/ --forbid migration/", "start:deployed": "bun apps/service/bin/production-server.ts", "bind:firebase": "bun scripts/bind-firebase-identity.ts" }, diff --git a/backends/example-platform/scripts/lint-import-graph.ts b/backends/example-platform/scripts/lint-import-graph.ts index a60b1d1fb4e..0c6e69487b5 100644 --- a/backends/example-platform/scripts/lint-import-graph.ts +++ b/backends/example-platform/scripts/lint-import-graph.ts @@ -614,6 +614,7 @@ for (const file of files(root, new Set(["frontend"]))) { && shown !== "drivers/postgres/product-projection-repository.ts" && shown !== "drivers/postgres/tasks-repository.ts" && shown !== "drivers/postgres/conversation-read-repository.ts" + && shown !== "drivers/postgres/chat-read-repository.ts" && shown !== "drivers/postgres/legacy-proposition-migration-repository.ts" && shown !== "drivers/postgres/memory-experiment-repository.ts" && shown !== "drivers/postgres/memory-query-evaluation-source.ts" diff --git a/backends/example-platform/scripts/postgres-test.ts b/backends/example-platform/scripts/postgres-test.ts index 87f3b4b7315..8c74d285ab2 100644 --- a/backends/example-platform/scripts/postgres-test.ts +++ b/backends/example-platform/scripts/postgres-test.ts @@ -389,7 +389,7 @@ const testPostgres = async (preserve: boolean): Promise => { childEnvironment["OMI_TEST_POSTGRES_URL"] = postgresTestConnectionString(prepared.state, passwordFrom(prepared.state)); childEnvironment["OMI_TEST_POSTGRES_IMAGE"] = prepared.state.image; try { - const result = command(["bun", "test", "drivers/postgres/postgresjs.real.test.ts", "drivers/postgres/derived-group-dream-work-input.real.test.ts", "drivers/postgres/derived-group-dream-success.real.test.ts", "drivers/postgres/derived-group-dream-one-shot-runtime.real.test.ts", "drivers/postgres/firebase-binding.real.test.ts", "drivers/postgres/tasks.real.test.ts", "drivers/postgres/conversations.real.test.ts", "migration/postgres/deletion-receipts.real.test.ts"], { + const result = command(["bun", "test", "drivers/postgres/postgresjs.real.test.ts", "drivers/postgres/derived-group-dream-work-input.real.test.ts", "drivers/postgres/derived-group-dream-success.real.test.ts", "drivers/postgres/derived-group-dream-one-shot-runtime.real.test.ts", "drivers/postgres/firebase-binding.real.test.ts", "drivers/postgres/tasks.real.test.ts", "drivers/postgres/conversations.real.test.ts", "drivers/postgres/chat-messages.real.test.ts", "migration/postgres/deletion-receipts.real.test.ts"], { env: childEnvironment, inherit: true, }); if (result.exitCode !== 0) process.exitCode = result.exitCode; diff --git a/docs/provider-agnostic-handoff.md b/docs/provider-agnostic-handoff.md index 1015e1ac138..3e33b259dac 100644 --- a/docs/provider-agnostic-handoff.md +++ b/docs/provider-agnostic-handoff.md @@ -22,9 +22,13 @@ GitHub `Validate backend worker` has been red since ff4738ee63: the PostgreSQL L Cloud Linux after fa6c525055: GitHub `Validate backend worker` passed on that commit, including real PostgreSQL 18.4 Listen qualification with `errorCode: null`. Worker transcript GET/POST now projects the same client fields as `parseRecordingTranscript` / `deviceTranscriptionProjection`. Unreadable stored segments stay 503 instead of a 200 body the app would treat as an unacknowledged transcript. Root `bun run check` passed locally after that projection. Apple Debug builds, live ScreenCaptureKit, and physical BLE/iPad remain unverified on Linux. +Cloud Linux after the chat GET slice: production mounts `GET /v1/chat-messages` under an explicit `chat.read` grant. Missing or revoked grants are 403. Granted empty history is an honest empty page with the existing chat attachment capabilities. Assistant rows without a unique terminal generation event are 503. Writes and SSE stay unmounted 404. Migration 0055 is in the checksummed manifest and is not applied to DEV from this VM. `bun run check:deployed` passed locally (real PostgreSQL 18.4 skipped without Docker). Apple Debug builds, live ScreenCaptureKit, and physical BLE/iPad remain unverified on Linux. + +Cloud Linux after the Settings GET slice: production mounts `GET /v1/settings`. Signed-out remains `{identity:null,entitlement:null}`. A verified Firebase identity without an owner-backed profile/entitlement producer is 503, not a 200 invented from token claims. Chat 403 and conversation/memory/task grant denials surface typed copy instead of empty libraries. First-page `GET /v1/conversations` now composes `chat:chat-main` when `chat.read` is granted and main-session messages exist; it never invents that row for empty or ungranted history. Apple Debug builds, live ScreenCaptureKit, and physical BLE/iPad remain unverified on Linux. + ## Verified and pushed -Backend slices include authorized PostgreSQL tasks/audio (a64e477a56), transcription with durable paid-response recovery (f22f1d8572), conversation reads with revision-fenced cursors (8546a0e574), and current trusted chat context packets (01bdb648bb). Chat persistence and deployed gateway identity composition remain missing. +Backend slices include authorized PostgreSQL tasks/audio (a64e477a56), transcription with durable paid-response recovery (f22f1d8572), conversation reads with revision-fenced cursors (8546a0e574), trusted chat context packets (01bdb648bb), GET-only PostgreSQL chat history under an explicit `chat.read` grant, and GET-only Settings that stay unavailable until an owner-backed identity/entitlement producer exists. Chat writes, generation SSE, Settings identity/entitlement producers, attachments and deployed gateway identity composition remain missing. Native slices include capability-gated LED/microphone controls (9ea1065492), account-session device retirement (9af6d6c248), bounded reconnect plus truthful connecting UI (d2bd77e8b3), and acknowledged Find Device commands (5dd93b4ecd). Native retries use 1/2/4-second delays only after a previously ready connection; manual disconnect, account retirement and radio-off cancel them. Mind Map remains collapsed on Home under the existing design decision. Canonical capture ownership receipts are pushed in b4b989025e; the deletion inventory and real cleanup gate are repaired in 0748acd29b. diff --git a/docs/provider-agnostic-parity.md b/docs/provider-agnostic-parity.md index c6213dc2a18..9da4f6d105d 100644 --- a/docs/provider-agnostic-parity.md +++ b/docs/provider-agnostic-parity.md @@ -26,21 +26,21 @@ Working acceptance checklist, inspected 2026-09-07 against the current v5 workin | Capability / app owner | Existing implementation and wire | Missing implementation or integration | Required acceptance evidence | | ------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | Authentication, setup and account binding — native auth + backend authorization | Native Google/Firebase flow, PKCE, secure session storage, refresh/cancel/sign-out; shared RN setup revision. Production memory authorization composes verified Firebase identity → application grant → account control. [Authorization composition](../backends/example-platform/apps/service/composition/firebase-memory-authorization.ts), [Postgres source](../backends/example-platform/drivers/postgres/firebase-application-authorization-source.ts). | A verified token alone does not create an application grant or account binding. Production onboarding/provisioning must establish the existing authority records and lifecycle, without treating a staging bearer or user-supplied account ID as a grant. Installed mobile builds need explicit intended backend configuration. PWA production sign-in is not implemented by the local proxy. | Fresh and existing users complete real native sign-in; correct account/grant admits access; missing, revoked, wrong-project and wrong-account grants deny correctly; refresh, cancel, sign-out and account changes retire old work. | -| Chat and SSE — `chatClient`, native generation streams | `GET/POST /v1/chat-messages`, `GET /v1/chat-generations/:id/events`, cancellation. Worker persists generation events; local portable service has [chat routes](../backends/example-platform/apps/service/routes/chat-messages.ts). | Production entry does not compose chat stores, admission, idempotency, streaming lifecycle or model execution. Worker Firebase accounts still use a staging lifetime admission limit, not a production subscription-period authority. | Actual deployed app sends a prompt, observes incremental and terminal events, reloads history, reconnects with event cursor, cancels, retries an uncertain admission exactly once, and shows provider failure honestly. Exercise entitlement exhaustion and a second isolated account. | +| Chat and SSE — `chatClient`, native generation streams | `GET/POST /v1/chat-messages`, `GET /v1/chat-generations/:id/events`, cancellation. Worker persists generation events; local portable service has [chat routes](../backends/example-platform/apps/service/routes/chat-messages.ts). Production entry now mounts `GET /v1/chat-messages` under the exact `chat.read` grant; missing grants are 403, granted empty history is an honest empty page, and assistant rows without a unique terminal event are 503. [Chat deployed read](../backends/example-platform/docs/memory-productionization/chat-deployed.md). | Production entry does not compose chat admission, POST, SSE, cancellation, attachments, or a source-owned entitlement producer. Worker Firebase accounts still use a staging lifetime admission limit, not a production subscription-period authority. Migration 55 is not yet applied to DEV. | Actual deployed app sends a prompt, observes incremental and terminal events, reloads history, reconnects with event cursor, cancels, retries an uncertain admission exactly once, and shows provider failure honestly. Exercise entitlement exhaustion and a second isolated account. | | Recording/transcription — `useNativeDevices`, `deviceSessionClient` | Worker `POST /v1/device-sessions`, indexed `/audio`, `/complete`, and `GET /:id/transcript`; durable transcription jobs feed server-created conversation records. Portable local service instead has [WebSocket `/v4/listen`](../backends/example-platform/apps/service/routes/listen.ts). | The production entry stores indexed audio chunks and upload completion in PostgreSQL using the canonical Listen authority. Transcription processing and transcript reads now use persisted claims, saved provider results and canonical batched Listen publication. Retries resume saved results without another provider call. Conversation listing is now mounted; downstream formation workers remain unwired; upload completion is not transcript completion. Phone-microphone recording is not implemented by the wearable path; encrypted journal recovery is implemented with physical offline/restart verification still open. | Physical Omi audio survives lost acknowledgements and reconnects without duplicate bytes; completion leads to an actual stored transcript in the mobile conversation detail; queued/running jobs can resume; terminal failures are truthful and refresh without restarting paid work; wrong-account reads fail. Test real background/foreground behavior and process death separately. | -| Conversations — shared `ConversationsPage` | Worker projects namespaced `chat:` and `recording:` IDs, bounded metadata and transcript detail. Local portable composition has conversation projection routes. Mobile reuses the full page rather than only recap cards. | Production entry now reads recorded and canonical finalized Listen sessions under the exact conversations.read grant. Persisted account revisions fence cursors when visibility or state changes. Chat conversation composition and metadata mutations remain missing. Conversation pages now select bounded PostgreSQL rows before transcript projection; signed cursor positions persist for 900 seconds with a 10,000-position account metadata ceiling. | Persisted chat and recordings appear after a fresh app launch, account isolation holds, no source-ID collision occurs, transcript/detail works, and pagination/search scope is explicit. No duplicated or skipped records at cursor boundaries. | +| Conversations — shared `ConversationsPage` | Worker projects namespaced `chat:` and `recording:` IDs, bounded metadata and transcript detail. Local portable composition has conversation projection routes. Mobile reuses the full page rather than only recap cards. | Production entry now reads recorded and canonical finalized Listen sessions under the exact conversations.read grant. Persisted account revisions fence cursors when visibility or state changes. First-page envelope reads also compose `chat:chat-main` when the same credential holds `chat.read` and main-session messages exist; missing grants and empty history never invent that row. Metadata mutations remain missing. Conversation pages now select bounded PostgreSQL rows before transcript projection; signed cursor positions persist for 900 seconds with a 10,000-position account metadata ceiling. | Persisted chat and recordings appear after a fresh app launch, account isolation holds, no source-ID collision occurs, transcript/detail works, and pagination/search scope is explicit. No duplicated or skipped records at cursor boundaries. | | Tasks — shared task client/UI + canonical task authority | App `GET /v1/tasks`, `POST /v1/tasks/ops`; complete/reopen/edit uses account epoch and revision, retained operation identity, failure retry/dismiss, authoritative refresh. Local SQLite app path is exercised; Worker delegates canonical reads/writes through the same service binding when configured. | Production entry now composes verified Firebase task grants, persisted control/epoch fences and transactional task/receipt storage. Remaining limits include a 10,000-distinct-record lifetime ceiling including tombstones and an undeployed straggler retention scheduler. Real account and grant issuance remains required. Task generation from real recordings also needs its actual producer, not a seeded task list. | Signed-in app edits a real task, restarts, sees saved state; duplicate op is acknowledged once; stale revision/epoch blocks writes; conflict refresh preserves draft; lost acknowledgement retries the same op; another account cannot read/write it. | | Memories — ratified read adapter + Firebase/Postgres reader | App `GET /v1/memories`; production entry composes a real authorized reader and HTTP render-model port. Worker [memory adapter](../apps/backend-worker/src/memory-service.ts) forwards original Firebase credentials via [service binding](../apps/backend-worker/src/canonical-service.ts), never staging credentials. | Production data ingestion, ownership/grants, render model configuration and actual accepted-work/STM coverage must be present. The current entry intentionally declares coverage unavailable. Memory-read deployment does not establish the formation/write pipeline. | Authorized non-seeded source data yields grounded renders with valid citations/provenance and honest completeness. Missing coverage stays incomplete, cross-account reads fail, revocation works, and the mobile page consumes the deployed response without a legacy-memory mapping. | -| Settings/account/privacy — `SettingsPage`, `desktopCloudClient` | Native account/subscription/recording-storage/training/private-sync/webhook reads and supported privacy mutations use existing primary-cloud APIs. Web `GET /v1/settings` preserves nullable identity and shows chat requests or transcription seconds with their actual units; policy links and native OS permissions are real actions. Settings is a mobile bottom destination. | Worker lacks legacy `/v1/users/*` settings APIs. Local portable `/v1/settings` is allowlisted for GET by the development proxy but is not mounted by production entry. Production identity, entitlements and privacy mutation authorities must be composed; do not present staging labels as verified profile/billing. Name/language editing, notification preferences and storage controls need real persistence before adding forms. | Mobile Settings remains selected in bottom navigation; authenticated profile/usage reflect the actual user; a supported privacy change survives reload and affects the authoritative policy; unavailable settings show retryable user-facing copy without request IDs. | +| Settings/account/privacy — `SettingsPage`, `desktopCloudClient` | Native account/subscription/recording-storage/training/private-sync/webhook reads and supported privacy mutations use existing primary-cloud APIs. Web `GET /v1/settings` preserves nullable identity and shows chat requests or transcription seconds with their actual units; policy links and native OS permissions are real actions. Settings is a mobile bottom destination. Production entry now mounts `GET /v1/settings`: signed-out is `{identity:null,entitlement:null}`, invalid credentials are 401, and a verified Firebase identity without an owner-backed producer is 503. [Settings deployed read](../backends/example-platform/docs/memory-productionization/settings-deployed.md). | Worker lacks legacy `/v1/users/*` settings APIs. Production identity, entitlements and privacy mutation authorities must still be composed; do not present staging labels as verified profile/billing. Name/language editing, notification preferences and storage controls need real persistence before adding forms. | Mobile Settings remains selected in bottom navigation; authenticated profile/usage reflect the actual user; a supported privacy change survives reload and affects the authoritative policy; unavailable settings show retryable user-facing copy without request IDs. | | Connectors/apps — `ConnectorsPage` | `/v1/apps`, `/v1/apps/enabled`, profile, enable/disable calls exist in [cloud client](../react-native/src/desktopCloudClient.ts); native transport can use primary cloud APIs. Apps tab is retained. | Worker and production portable entry do not expose connector catalog/ownership/enablement APIs. Production browser routing/auth is missing. A connector toggle is not proof that its external authorization or ingestion works. | Real catalog loads, enabled state survives reload, ownership is enforced, external authorization cancellation/errors are handled, and app data arrives through an authorized producer. | | Attachments — native transport + backend attachment authority | Worker upload admission/completion, checksums, ownership and chat attachment admission exist in [attachments](../apps/backend-worker/src/attachments.ts). Local portable service has [attachment routes](../backends/example-platform/apps/service/routes/chat-attachments.ts). | Production portable entry does not compose attachment storage/signing/completion. End-user native picker/upload and signed object-store integration need their own evidence; backend tests alone do not establish them. | Attach, upload, complete and send through the actual app; reject mismatched checksum/size/owner, expire access and recover uncertain completion. Do not expose provider credentials in JS or public object URLs. | | Offline/process recovery — native storage owner | Capture uses an encrypted native recording journal, with bounded fsynced entries, native owner receipts and indexed batch acknowledgement before replay cleanup. [Journal implementation and limits](native-recording-journal.md) describe the account/origin/login partition and recovery contract; the separate [sync outbox](../packages/sync/src/outbox.ts) is not the recording journal. | Physical offline/restart replay and suspended background packet delivery remain unverified. Cached ownership is permitted only after classified offline transport failures; every replay requires fresh matching authority. Process recovery protects already persisted packets and does not imply continuous recording after OS termination. | Capture while offline, terminate process, relaunch/reconnect and upload once; sign-out/account-switch cannot replay old account audio; corrupt/partial journal entries fail visibly and storage limits are enforced. | -Settings authority boundary: the canonical [settings projection](../backends/example-platform/apps/service/control/settings-projection.ts) requires identity and shares its entitlement projection with enforcement. Its implementations are [local SQLite](../backends/example-platform/drivers/sqlite/service-stores/settings-projection.ts) and in-memory; [app-facing producers](../backends/example-platform/apps/service/app-facing.ts) seed QA identity and optional unmetered entitlement. PostgreSQL account/grant records do not supply profile, billing or consumed usage, and [verified Firebase identity](../backends/example-platform/apps/service/auth/firebase-identity.ts) supplies only project, UID, authentication strength and expiry. A deployed settings reader needs an owner-backed profile projection and the same persisted entitlement/usage producer used by enforcement. Meanwhile, native `/v1/users/profile` and `/v1/users/me/subscription` continue through the primary cloud [client](../react-native/src/desktopCloudClient.ts) and [route policy](../react-native/src/v5BackendOrigin.ts); they are not portable-backend implementations. Web service settings read the existing local projection, preserve absent identity/entitlement, and label chat requests and transcription seconds separately. Unknown allowance units remain unavailable; this local read does not create a deployed producer. +Settings authority boundary: the canonical [settings projection](../backends/example-platform/apps/service/control/settings-projection.ts) requires identity and shares its entitlement projection with enforcement. Its implementations are [local SQLite](../backends/example-platform/drivers/sqlite/service-stores/settings-projection.ts) and in-memory; [app-facing producers](../backends/example-platform/apps/service/app-facing.ts) seed QA identity and optional unmetered entitlement. PostgreSQL account/grant records do not supply profile, billing or consumed usage, and [verified Firebase identity](../backends/example-platform/apps/service/auth/firebase-identity.ts) supplies only project, UID, authentication strength and expiry. Production now mounts `GET /v1/settings` and returns 503 after a verified identity until that owner-backed producer exists; it does not project token claims as a profile. Native `/v1/users/profile` and `/v1/users/me/subscription` continue through the primary cloud [client](../react-native/src/desktopCloudClient.ts) and [route policy](../react-native/src/v5BackendOrigin.ts); they are not portable-backend implementations. Web service settings read the existing local projection, preserve absent identity/entitlement, and label chat requests and transcription seconds separately. Unknown allowance units remain unavailable; this local read does not create a deployed producer. -Chat identity boundary: the PostgreSQL [memory context adapter](../backends/example-platform/drivers/postgres/firebase-chat-generation-context-source.ts) now produces the trusted context packet consumed by the current supervisor, using the authenticated account-bound memory read. This does not authorize chat generation. The unmounted [gateway source](../backends/example-platform/apps/service/chat/generation-source.ts) still sends the opaque context account ID as `x-omi-user-uid` and defaults its service caller to `platform`; neither establishes the verified Firebase UID or a gateway-allowlisted caller. Production chat must compose the actual verified Firebase identity, account binding and explicit caller configuration before mounting this source. Its PostgreSQL admission, message/event persistence and entitlement authority are also missing; the local QA principal and memory-read grant cannot substitute for them. +Chat identity boundary: the PostgreSQL [memory context adapter](../backends/example-platform/drivers/postgres/firebase-chat-generation-context-source.ts) now produces the trusted context packet consumed by the current supervisor, using the authenticated account-bound memory read. This does not authorize chat generation. Production now persists and reads chat history under `chat.read`; it still does not mount POST/SSE. The unmounted [gateway source](../backends/example-platform/apps/service/chat/generation-source.ts) still sends the opaque context account ID as `x-omi-user-uid` and defaults its service caller to `platform`; neither establishes the verified Firebase UID or a gateway-allowlisted caller. Production chat writes must compose the actual verified Firebase identity, account binding, explicit caller configuration and a source-owned entitlement producer before mounting admission. The local QA principal and memory-read grant cannot substitute for `chat.read` or `chat.write`. -Chat entitlement integration prerequisite: the legacy origin owns `backend/utils/subscription.py:get_chat_quota_snapshot`, catalog revision 2 in `backend/config/plan_catalog.json`, and UTC monthly consumption in `backend/database/user_usage.py:get_monthly_chat_usage`. The `/v1/users/me/usage-quota` response is a read snapshot, not an atomic reservation; the gateway's `gateway/jit_budget.py` reservation is restricted to JIT QA run budgets and cannot authorize subscriber chat. The required source-owned admission operation must bind the verified Firebase project/UID, canonical account and epoch, client-message ID and payload hash to one durable reservation. Exact replay must return its existing receipt without consuming again; changed payload must conflict; missing source state must be unavailable; exhaustion must remain an entitlement refusal. The receipt needs the source catalog/subscription revision, UTC usage period and billing unit, and must be persisted with PostgreSQL message admission and its first generation event. Completion must settle that same reservation once from gateway-owned provider-attempt accounting; interrupted or unknown-cost attempts must not be silently released. A same-message retry after an uncertain response must reconcile the reservation rather than allocate another. These are integration requirements, not an implemented public endpoint: the legacy source producer and its authenticated transport do not yet exist. PostgreSQL `chat.read` and `chat.write` would require explicit grants through the existing application authorization lifecycle; neither Firebase identity alone nor a memory grant supplies them. The current token-only `ChatGenerationUsage` shape also lacks the cost and settlement evidence needed for monetary plans. +Chat entitlement integration prerequisite: the legacy origin owns `backend/utils/subscription.py:get_chat_quota_snapshot`, catalog revision 2 in `backend/config/plan_catalog.json`, and UTC monthly consumption in `backend/database/user_usage.py:get_monthly_chat_usage`. The `/v1/users/me/usage-quota` response is a read snapshot, not an atomic reservation; the gateway's `gateway/jit_budget.py` reservation is restricted to JIT QA run budgets and cannot authorize subscriber chat. The required source-owned admission operation must bind the verified Firebase project/UID, canonical account and epoch, client-message ID and payload hash to one durable reservation. Exact replay must return its existing receipt without consuming again; changed payload must conflict; missing source state must be unavailable; exhaustion must remain an entitlement refusal. The receipt needs the source catalog/subscription revision, UTC usage period and billing unit, and must be persisted with PostgreSQL message admission and its first generation event. Completion must settle that same reservation once from gateway-owned provider-attempt accounting; interrupted or unknown-cost attempts must not be silently released. A same-message retry after an uncertain response must reconcile the reservation rather than allocate another. These are integration requirements, not an implemented public endpoint: the legacy source producer and its authenticated transport do not yet exist. PostgreSQL `chat.read` now exists as an explicit application grant for history GET; `chat.write` still requires that same lifecycle and is not conferred by Firebase identity or a memory grant. The current token-only `ChatGenerationUsage` shape also lacks the cost and settlement evidence needed for monetary plans. ## Acceptance checklist for subsequent work diff --git a/react-native/__tests__/chatClient.test.ts b/react-native/__tests__/chatClient.test.ts index 626bd9bcc9f..73401dc0c2b 100644 --- a/react-native/__tests__/chatClient.test.ts +++ b/react-native/__tests__/chatClient.test.ts @@ -387,6 +387,14 @@ test('maps ratified public recovery without automatically retrying', () => { expect( chatErrorCopy(new ChatBackendError(404, 'not_found', false, 'none', null)), ).toBe('This request cannot be completed.'); + expect( + chatErrorCopy(new ChatBackendError(403, 'forbidden', false, 'none', null)), + ).toBe('Chat is not available for this account.'); + expect( + chatHistoryErrorCopy( + new ChatBackendError(403, 'forbidden', false, 'none', null), + ), + ).toBe('Chat is not available for this account.'); expect( chatHistoryErrorCopy( new ChatBackendError(401, 'unauthorized', false, 'reauthenticate', null), diff --git a/react-native/__tests__/desktopReadClient.test.ts b/react-native/__tests__/desktopReadClient.test.ts index 47359739364..fa3c63721d5 100644 --- a/react-native/__tests__/desktopReadClient.test.ts +++ b/react-native/__tests__/desktopReadClient.test.ts @@ -4,6 +4,7 @@ import { conversationGroupLabel, desktopBackendConfigurationCopy, desktopBackendUnauthorizedCopy, + desktopBackendForbiddenCopy, desktopCloudBaseURL, desktopLocalBackendServiceCopy, desktopProjectionUnavailableCopy, @@ -257,6 +258,9 @@ test('maps native cloud-first backend failures to actionable, credential-safe co expect( desktopReadErrorCopy(new Error('Conversations response is malformed')), ).toBe('This saved data could not be loaded. Retry without changing it.'); + expect(desktopReadErrorCopy(new Error(desktopBackendForbiddenCopy))).toBe( + desktopBackendForbiddenCopy, + ); }); describe('desktopRecoveryCopy', () => { @@ -857,6 +861,26 @@ test('maps a cloud 401 to typed unauthorized copy without fabricating rows', asy }); }); +test('maps a cloud 403 to typed grant-denied copy without treating an empty body as success', async () => { + const backend = backendFor(() => ({ + status: 403, + body: JSON.stringify({items: []}), + })); + const result = await loadDesktopReads(backend); + expect(result.conversations).toEqual({ + status: 'error', + error: desktopBackendForbiddenCopy, + }); + expect(result.memories).toEqual({ + status: 'error', + error: desktopBackendForbiddenCopy, + }); + expect(result.tasks).toEqual({ + status: 'error', + error: desktopBackendForbiddenCopy, + }); +}); + test('parses catalogue, enabled, owned, and service app records without inventing rows', () => { const app = parseCloudApp( { diff --git a/react-native/src/chatClient.ts b/react-native/src/chatClient.ts index cd3a674399f..d705d3c1198 100644 --- a/react-native/src/chatClient.ts +++ b/react-native/src/chatClient.ts @@ -57,6 +57,9 @@ export function chatErrorCopy(error: unknown): string { if (error.action === 'reauthenticate' || error.status === 401) { return 'Sign in again to continue.'; } + if (error.status === 403 || error.backendCode === 'forbidden') { + return 'Chat is not available for this account.'; + } if (error.status === 429) { return error.retryAfterSeconds === null ? 'Too many requests. Try again shortly.' diff --git a/react-native/src/chatConversationHistory.test.tsx b/react-native/src/chatConversationHistory.test.tsx new file mode 100644 index 00000000000..5a6720d9053 --- /dev/null +++ b/react-native/src/chatConversationHistory.test.tsx @@ -0,0 +1,219 @@ +import React from 'react'; +import ReactTestRenderer, {act} from 'react-test-renderer'; +import {Text} from 'react-native'; +import type {NativeHttpResponse} from './omiNativeTypes'; +import type { + ConversationProjection, + DomainReadOutcome, +} from './desktopReadClient'; +import {ChatBackendError} from './chatClient'; + +const mockRequest = jest.fn(); +let mockInvalidated: (() => void) | undefined; +jest.mock('./omiNative', () => ({ + omiBackend: {request: (request: unknown) => mockRequest(request)}, + subscribeOmiBackendSessionInvalidated: (listener: () => void) => { + mockInvalidated = listener; + return () => { + mockInvalidated = undefined; + }; + }, +})); + +const {ConversationsPage} = require('./pages/Conversations'); +const {MAIN_CHAT_CONVERSATION_ID} = require('./chatConversationHistory'); + +function historyResponse( + messages: Array<{id: string; text: string; sender: 'human' | 'ai'}>, +): NativeHttpResponse { + return { + id: 'chat-history', + status: 200, + body: JSON.stringify({ + messages: messages.map(message => ({ + ...message, + createdAt: 1_000, + generationOutcome: message.sender === 'ai' ? 'completed' : null, + type: 'text', + updatedAt: 1_000, + chatSessionId: null, + appId: null, + journalRevision: 1, + payloadHash: 'sha256:test', + messageSource: 'desktop_chat', + rating: null, + reported: false, + revision: '1', + attachments: [], + })), + page: {olderCursor: null, hasOlder: false}, + capabilities: { + maxAttachmentsPerMessage: 4, + maxAttachmentBytes: 52_428_800, + allowedAttachmentMimeTypes: ['text/plain'], + }, + }), + }; +} + +function textOf(renderer: ReactTestRenderer.ReactTestRenderer) { + return renderer.root + .findAllByType(Text) + .map(node => node.props.children) + .flat() + .join(' '); +} + +function conversation( + overrides: Partial, +): ConversationProjection { + return { + kind: 'conversation', + id: MAIN_CHAT_CONVERSATION_ID, + title: 'saved prompt', + summary: 'saved prompt', + searchableText: 'saved prompt', + createdAt: '2026-09-07T00:00:00Z', + updatedAt: '2026-09-07T00:01:00Z', + startedAt: '2026-09-07T00:00:00Z', + finishedAt: null, + starred: false, + status: 'in_progress', + source: 'chat', + visibility: 'private', + folderId: null, + locked: false, + discarded: false, + ...overrides, + }; +} + +function outcome( + items: ConversationProjection[], +): DomainReadOutcome { + return { + status: 'success', + value: { + items, + page: { + windowStatus: 'complete', + complete: true, + hasMore: false, + nextCursor: null, + completenessStatus: 'complete', + reasons: [], + }, + }, + }; +} + +const renderers: ReactTestRenderer.ReactTestRenderer[] = []; + +afterEach(() => { + act(() => renderers.splice(0).forEach(renderer => renderer.unmount())); + mockRequest.mockReset(); +}); + +async function renderPage(items: ConversationProjection[]) { + const native = require('react-native'); + jest.spyOn(native, 'useWindowDimensions').mockReturnValue({ + width: 390, + height: 844, + scale: 1, + fontScale: 1, + }); + let renderer!: ReactTestRenderer.ReactTestRenderer; + await act(async () => { + renderer = ReactTestRenderer.create( + , + ); + }); + renderers.push(renderer); + return renderer; +} + +test('opens chat:chat-main into persisted messages instead of a title-only detail', async () => { + mockRequest.mockResolvedValue( + historyResponse([ + {id: 'human-1', text: 'saved prompt', sender: 'human'}, + {id: 'ai-1', text: 'saved answer', sender: 'ai'}, + ]), + ); + const renderer = await renderPage([conversation({})]); + expect(mockRequest).not.toHaveBeenCalled(); + await act(async () => + renderer.root + .findAll( + node => + node.props.accessibilityLabel === 'Open conversation saved prompt', + )[0]! + .props.onPress(), + ); + expect(mockRequest).toHaveBeenCalledWith({ + id: 'chat-history', + method: 'GET', + expectedApiContract: 'canonical', + path: '/v1/chat-messages?limit=50', + }); + expect(textOf(renderer)).toContain('You · saved prompt'); + expect(textOf(renderer)).toContain('Omi · saved answer'); +}); + +test('does not load main chat history for another chat session id', async () => { + const renderer = await renderPage([ + conversation({id: 'chat:session-alpha', title: 'Other session'}), + ]); + await act(async () => + renderer.root + .findAll( + node => + node.props.accessibilityLabel === 'Open conversation Other session', + )[0]! + .props.onPress(), + ); + expect(mockRequest).not.toHaveBeenCalled(); + expect(textOf(renderer)).toContain( + 'Chat history for this conversation is not available here.', + ); + expect(textOf(renderer)).not.toContain('You ·'); +}); + +test('shows typed chat grant denial instead of an empty message list', async () => { + mockRequest.mockRejectedValue( + new ChatBackendError(403, 'forbidden', false, 'none', null), + ); + const renderer = await renderPage([conversation({})]); + await act(async () => + renderer.root + .findAll( + node => + node.props.accessibilityLabel === 'Open conversation saved prompt', + )[0]! + .props.onPress(), + ); + expect(textOf(renderer)).toContain( + 'Chat is not available for this account.', + ); + expect(textOf(renderer)).not.toContain('No messages in this chat yet.'); + expect(textOf(renderer)).not.toContain('You ·'); +}); + +test('keeps an honest empty chat page instead of inventing a completed answer', async () => { + mockRequest.mockResolvedValue(historyResponse([])); + const renderer = await renderPage([conversation({})]); + await act(async () => + renderer.root + .findAll( + node => + node.props.accessibilityLabel === 'Open conversation saved prompt', + )[0]! + .props.onPress(), + ); + expect(textOf(renderer)).toContain('No messages in this chat yet.'); + expect(textOf(renderer)).not.toContain('You ·'); + expect(textOf(renderer)).not.toContain('Omi ·'); +}); diff --git a/react-native/src/chatConversationHistory.ts b/react-native/src/chatConversationHistory.ts new file mode 100644 index 00000000000..a773f17924f --- /dev/null +++ b/react-native/src/chatConversationHistory.ts @@ -0,0 +1,73 @@ +import {useEffect, useRef, useState} from 'react'; +import { + chatHistoryErrorCopy, + loadChatHistory, + type ChatMessage, +} from './chatClient'; +import {omiBackend, subscribeOmiBackendSessionInvalidated} from './omiNative'; + +export const MAIN_CHAT_CONVERSATION_ID = 'chat:chat-main'; + +type ChatHistoryRead = + | {status: 'idle'} + | {status: 'loading'} + | {status: 'error'; error: string} + | {status: 'loaded'; messages: ChatMessage[]}; + +export function useChatConversationHistory(active: boolean) { + const [result, setResult] = useState({status: 'idle'}); + const [reload, setReload] = useState(0); + const epoch = useRef(0); + useEffect( + () => + subscribeOmiBackendSessionInvalidated(() => { + epoch.current++; + setResult(previous => + previous.status === 'idle' + ? previous + : { + status: 'error', + error: chatHistoryErrorCopy({ + code: 'OMI_HTTP_UNAUTHORIZED', + }), + }, + ); + }), + [], + ); + useEffect(() => { + const current = ++epoch.current; + let alive = true; + if (!active) { + setResult({status: 'idle'}); + return () => { + alive = false; + }; + } + setResult({status: 'loading'}); + const load = async () => { + try { + if (omiBackend == null) { + throw new Error('Native transport unavailable'); + } + const messages = await loadChatHistory(omiBackend); + if (!alive || epoch.current !== current) { + return; + } + setResult({status: 'loaded', messages}); + } catch (error) { + if (alive && epoch.current === current) { + setResult({status: 'error', error: chatHistoryErrorCopy(error)}); + } + } + }; + void load(); + return () => { + alive = false; + }; + }, [active, reload]); + return { + result, + reload: () => setReload(value => value + 1), + }; +} diff --git a/react-native/src/desktopReadClient.ts b/react-native/src/desktopReadClient.ts index fd3d6a55dde..4bad4f59392 100644 --- a/react-native/src/desktopReadClient.ts +++ b/react-native/src/desktopReadClient.ts @@ -197,6 +197,8 @@ export const desktopLocalBackendServiceCopy = 'The configured local Omi service is unavailable. Check its connection, then retry.'; export const desktopProjectionUnavailableCopy = 'This saved data is not available from the selected Omi service yet. Retry after its persisted projection is connected.'; +export const desktopBackendForbiddenCopy = + 'This saved data is not available for this account.'; const desktopReadFailureCopy = 'This saved data could not be loaded. Retry without changing it.'; const desktopRecoveryGenericCopy = @@ -213,7 +215,8 @@ export function desktopRecoveryCopy( outcome.error === desktopBackendUnauthorizedCopy || outcome.error === desktopBackendServiceCopy || outcome.error === desktopLocalBackendServiceCopy || - outcome.error === desktopProjectionUnavailableCopy) + outcome.error === desktopProjectionUnavailableCopy || + outcome.error === desktopBackendForbiddenCopy) ) { return outcome.error; } @@ -259,6 +262,7 @@ export function desktopReadErrorCopy(error: unknown): string { desktopBackendServiceCopy, desktopLocalBackendServiceCopy, desktopProjectionUnavailableCopy, + desktopBackendForbiddenCopy, desktopReadFailureCopy, ].includes(message) ? message @@ -376,6 +380,9 @@ async function read( unauthorized.code = 'unauthorized'; throw unauthorized; } + if (response.status === 403) { + throw new Error(desktopBackendForbiddenCopy); + } if (response.status === 503 && response.body !== null) { try { const body = object(JSON.parse(response.body), `${id} error`); diff --git a/react-native/src/pages/Conversations.test.tsx b/react-native/src/pages/Conversations.test.tsx new file mode 100644 index 00000000000..8b05137ff87 --- /dev/null +++ b/react-native/src/pages/Conversations.test.tsx @@ -0,0 +1,39 @@ +import React from 'react'; +import ReactTestRenderer, {act} from 'react-test-renderer'; +import {Text} from 'react-native'; +import {ConversationsPage} from './Conversations'; + +function textOf(renderer: ReactTestRenderer.ReactTestRenderer): string { + return renderer.root + .findAllByType(Text) + .flatMap(node => + Array.isArray(node.props.children) + ? node.props.children + : [node.props.children], + ) + .filter( + (value): value is string | number => + typeof value === 'string' || typeof value === 'number', + ) + .join(' '); +} + +test('conversation grant denial shows the typed error instead of an empty library', () => { + let renderer!: ReactTestRenderer.ReactTestRenderer; + act(() => { + renderer = ReactTestRenderer.create( + , + ); + }); + expect(textOf(renderer)).toContain( + 'This saved data is not available for this account.', + ); + expect(textOf(renderer)).not.toContain('No conversations yet.'); + expect(textOf(renderer)).not.toContain('Conversations could not be loaded.'); +}); diff --git a/react-native/src/pages/Conversations.tsx b/react-native/src/pages/Conversations.tsx index 2257b99ba9f..adcafaf34a8 100644 --- a/react-native/src/pages/Conversations.tsx +++ b/react-native/src/pages/Conversations.tsx @@ -18,6 +18,8 @@ import { } from '../desktopReadClient'; import {FocusPressable} from '../ui/Pressable'; import {RecordingTranscript} from '../ui/RecordingTranscript'; +import {ChatConversationHistory} from '../ui/ChatConversationHistory'; +import {MAIN_CHAT_CONVERSATION_ID} from '../chatConversationHistory'; import {ReadStatus} from '../ui/ReadStatus'; import {styles} from '../ui/styles'; @@ -260,9 +262,7 @@ export function ConversationsPage({ Conversations unavailable - - Conversations could not be loaded. - + {error} ) : grouped.length === 0 ? ( @@ -388,6 +388,16 @@ export function ConversationsPage({ revision={selected.updatedAt ?? undefined} /> )} + {selected.source === 'chat' && + selected.id === MAIN_CHAT_CONVERSATION_ID && ( + + )} + {selected.source === 'chat' && + selected.id !== MAIN_CHAT_CONVERSATION_ID && ( + + Chat history for this conversation is not available here. + + )} )} diff --git a/react-native/src/pages/Tasks.test.tsx b/react-native/src/pages/Tasks.test.tsx index a19c756c646..0a71d874213 100644 --- a/react-native/src/pages/Tasks.test.tsx +++ b/react-native/src/pages/Tasks.test.tsx @@ -1,5 +1,6 @@ import React from 'react'; import ReactTestRenderer, {act} from 'react-test-renderer'; +import {Text} from 'react-native'; import {TasksPage} from './Tasks'; import {TaskPagination} from '../ui/TaskPagination'; import type {TaskMutationProps} from '../ui/TaskEditor'; @@ -170,6 +171,29 @@ test('conflict refresh preserves dirty description while untouched descriptions ); }); +test('task grant denial shows the typed error instead of an empty library', () => { + let renderer!: ReactTestRenderer.ReactTestRenderer; + act(() => { + renderer = ReactTestRenderer.create( + , + ); + }); + const copy = renderer.root + .findAllByType(Text) + .map(node => node.props.children) + .flat() + .join(' '); + expect(copy).toContain('This saved data is not available for this account.'); + expect(copy).not.toContain('No tasks yet.'); + expect(copy).not.toContain('Saved tasks could not be loaded.'); +}); + test('task pagination stays available when loaded task search has no matches', () => { const onLoadMore = jest.fn(); let renderer!: ReactTestRenderer.ReactTestRenderer; diff --git a/react-native/src/pages/Tasks.tsx b/react-native/src/pages/Tasks.tsx index e8e6dacee4a..25e82c517cd 100644 --- a/react-native/src/pages/Tasks.tsx +++ b/react-native/src/pages/Tasks.tsx @@ -119,9 +119,7 @@ export function TasksPage({ ) : error !== null ? ( Tasks unavailable - - Saved tasks could not be loaded. - + {error} ) : filtered.length === 0 ? ( diff --git a/react-native/src/ui/ChatConversationHistory.tsx b/react-native/src/ui/ChatConversationHistory.tsx new file mode 100644 index 00000000000..3ba20b8134b --- /dev/null +++ b/react-native/src/ui/ChatConversationHistory.tsx @@ -0,0 +1,50 @@ +import React from 'react'; +import {ActivityIndicator, Text, View} from 'react-native'; +import {useChatConversationHistory} from '../chatConversationHistory'; +import {FocusPressable} from './Pressable'; +import {styles} from './styles'; + +export function ChatConversationHistory() { + const {result, reload} = useChatConversationHistory(true); + return ( + + + Messages + + {result.status === 'loading' || result.status === 'idle' ? ( + + + Loading chat… + + ) : result.status === 'error' ? ( + <> + + {result.error} + + + Check again + + + ) : result.messages.length === 0 ? ( + + No messages in this chat yet. + + ) : ( + result.messages.map(message => ( + + {`${message.sender === 'human' ? 'You' : 'Omi'} · ${message.text}`} + + )) + )} + + ); +}