fix(agent-runtime): allow one owner material to be bound to multiple sessions (#1500)

* fix(agent-runtime): allow one owner material to be bound to multiple sessions

`agent_session_materials.id` is a global primary key, but
`bindOwnerMaterialsToSession` de-duplicated on `(sessionId, id)` and then
inserted the owner-side material id as the row id. Binding the same owner
material to a second session therefore raised a duplicate-key error
(23505) that the route surfaced as HTTP 500, so an owner material could
only ever be used by one course.

Keep row ids globally unique (every extraction method keys on `id` alone)
and store the shared owner id as metadata instead: the binder now mints a
fresh session row id, records `owner_material_id`, and looks that column
up for idempotent rebinding; a partial unique index on
`(session_id, owner_material_id)` adjudicates concurrent rebinds and the
loser adopts the winner's row. The schema change is additive and
idempotent (`ADD COLUMN IF NOT EXISTS` + `CREATE UNIQUE INDEX IF NOT
EXISTS`), so existing databases upgrade in place with `NULL` for legacy
rows. No FK from `owner_material_id` to the owner library, to keep
owner-library deletion decoupled from session rows.

Tests: backend-neutral contract test (same owner material bound to two
sessions, both readable), PGlite in-place upgrade test, pinned-schema
update, and a host-level test covering the route's 202 path.

* fix(agent-runtime): reuse legacy owner-material bindings and clean up losing uploads on concurrent bind

Rows written by the previous binder use the owner upload id as the session
row id and leave `owner_material_id` NULL, so the new owner-id lookup missed
them and a rebind created a duplicate row and copied the bytes again. The
binder now falls back to the session row id and adopts it only when it is
unambiguously that owner upload: source kind, matching title, no source URL
or derivative, no text, and the deterministic legacy object key with a byte
length that matches the owner record. Adoption stamps `owner_material_id`
through a conditional `backfillOwnerMaterialId` update that leaves extraction
columns untouched; a lost race re-reads the fast path. A different session is
still a fresh row.

The concurrent-bind loser also left its uploaded object behind after
adopting the winner. It now removes that object when the winner references a
different key, best-effort so a failed cleanup cannot fail a bind that
succeeded.

Tests: host-level PGlite regressions for the legacy rebind (id, row count,
byte copies, extraction state, and backfill) plus a different-session bind,
and a controllable byte store that parks the loser after its upload so the
winner commits first, asserting one row and only the winner's object remain.
A storage unit test pins the backfill contract. Bumps @openmaic/storage to
0.30.1 for the new public store method.
This commit is contained in:
wyuc
2026-09-14 18:43:31 +02:00
committed by GitHub
parent cfae106b89
commit ee7a7b64df
8 changed files with 643 additions and 33 deletions
+137 -28
View File
@@ -19,7 +19,10 @@ import {
type AgentSessionMeta,
type ListAgentSessionMaterialsOptions,
} from '@openmaic/storage';
import { getReadyOwnerMaterials } from '@/lib/persistence/owner-materials';
import {
getReadyOwnerMaterials,
type OwnerMaterialRecord,
} from '@/lib/persistence/owner-materials';
import { getServerPersistenceProvider } from '@/lib/persistence/server-provider';
import { getMaterialByteStore } from '@/lib/server/materials/bytes';
@@ -208,7 +211,12 @@ export async function createSourceMaterial(
/**
* Bind owner-library uploads to a session by copying their private bytes into
* the session byte prefix and creating the material rows the agent reads.
* Rebinding the same id is idempotent.
*
* Row ids are globally unique (the extraction pipeline keys on them), so the
* session row is minted with its own fresh `mat_` id and the owner upload id is
* recorded in `ownerMaterialId`. That lets one owner material back several
* sessions, and a unique `(session_id, owner_material_id)` index makes
* rebinding the same upload into the same session idempotent.
*/
export async function bindOwnerMaterialsToSession(
sessionId: string,
@@ -228,33 +236,9 @@ export async function bindOwnerMaterialsToSession(
const bound = [];
for (const id of materialIds) {
const record = byId.get(id)!;
if (!(await store.getMaterial(sessionId, id))) {
let source: Buffer;
try {
source = await byteStore.get(record.ossKey);
} catch {
throw new SessionMaterialBindingError(`material ${id} bytes are unavailable`);
}
const mime = record.mime ?? 'application/octet-stream';
const rawObjectKey = sessionMaterialKey(sessionId, id, rawObjectName(mime));
await byteStore.put(rawObjectKey, source, mime);
try {
await store.createMaterial(sessionId, {
id,
kind: 'source',
title: record.originalName ?? id,
rawAssetId: rawObjectKey,
textChars: 0,
});
} catch (error) {
if (!(await store.getMaterial(sessionId, id))) {
await byteStore.delete(rawObjectKey).catch(() => undefined);
throw error;
}
}
}
const material = await bindOwnerMaterial(store, byteStore, sessionId, id, record);
bound.push({
materialId: id,
materialId: material.id,
...(record.originalName ? { originalName: record.originalName } : {}),
...(record.mime ? { mime: record.mime } : {}),
bytes: record.bytes,
@@ -263,6 +247,131 @@ export async function bindOwnerMaterialsToSession(
return bound;
}
/**
* The deterministic object key the pre-upgrade binder copied an owner upload
* to: the owner id was used as the session row id, so the key follows from the
* session, the owner id, and the MIME type.
*/
function legacyOwnerMaterialKey(
sessionId: string,
ownerMaterialId: string,
mime: string | null,
): string {
return sessionMaterialKey(
sessionId,
ownerMaterialId,
rawObjectName(mime ?? 'application/octet-stream'),
);
}
/**
* Whether a row written before `owner_material_id` existed is the legacy
* binding of exactly this owner upload.
*
* The old binder keyed the row on the owner id, left `owner_material_id` NULL,
* and copied the bytes to the deterministic key above, so a match on the row
* id, kind, title, provenance, and that key plus the stored byte length is
* unambiguous. A row that fails any check stays a distinct row and is never
* silently adopted or overwritten.
*/
async function isLegacyOwnerMaterialBinding(
byteStore: ReturnType<typeof getMaterialByteStore>,
legacy: AgentSessionMaterial,
sessionId: string,
ownerMaterialId: string,
record: Pick<OwnerMaterialRecord, 'mime' | 'originalName' | 'bytes'>,
): Promise<boolean> {
if (legacy.id !== ownerMaterialId) return false;
if (legacy.ownerMaterialId !== null) return false;
if (legacy.kind !== 'source') return false;
if (legacy.title !== (record.originalName ?? ownerMaterialId)) return false;
// A legacy source binding carries copied raw bytes only: no fetch URL, no
// derivative, and no extracted text.
if (legacy.sourceUrl !== null || legacy.derivedFrom !== null) return false;
if (legacy.textAssetId !== null || legacy.textChars !== 0) return false;
const expectedKey = legacyOwnerMaterialKey(sessionId, ownerMaterialId, record.mime);
if (!isSessionMaterialKey(sessionId, expectedKey)) return false;
if (legacy.rawAssetId !== expectedKey) return false;
try {
return (await byteStore.get(expectedKey)).length === record.bytes;
} catch {
return false;
}
}
/**
* Bind one owner upload to a session. The `(session_id, owner_material_id)`
* unique index adjudicates concurrent binds of the same upload into the same
* session, so the loser adopts the winner's row instead of failing.
*/
async function bindOwnerMaterial(
store: PgAgentSessionMaterialStore,
byteStore: ReturnType<typeof getMaterialByteStore>,
sessionId: string,
ownerMaterialId: string,
record: Pick<OwnerMaterialRecord, 'ossKey' | 'mime' | 'originalName' | 'bytes'>,
): Promise<AgentSessionMaterial> {
const existing = await store.getMaterialByOwnerMaterialId(sessionId, ownerMaterialId);
if (existing) return existing;
// Upgrade path: rows written before `owner_material_id` existed use the
// owner upload id as the session row id. Reuse and backfill one instead of
// minting a duplicate row and copying the bytes a second time.
const legacy = await store.getMaterial(sessionId, ownerMaterialId);
if (
legacy &&
(await isLegacyOwnerMaterialBinding(byteStore, legacy, sessionId, ownerMaterialId, record))
) {
const adopted = await store.backfillOwnerMaterialId(sessionId, legacy.id, ownerMaterialId);
if (adopted) return adopted;
// A concurrent bind upgraded or claimed the owner id first; take its row.
const winner = await store.getMaterialByOwnerMaterialId(sessionId, ownerMaterialId);
if (winner) return winner;
}
let source: Buffer;
try {
source = await byteStore.get(record.ossKey);
} catch {
throw new SessionMaterialBindingError(`material ${ownerMaterialId} bytes are unavailable`);
}
const mime = record.mime ?? 'application/octet-stream';
const sessionMaterialId = createMaterialId();
const rawObjectKey = sessionMaterialKey(sessionId, sessionMaterialId, rawObjectName(mime));
await byteStore.put(rawObjectKey, source, mime);
try {
return await store.createMaterial(sessionId, {
id: sessionMaterialId,
ownerMaterialId,
kind: 'source',
title: record.originalName ?? ownerMaterialId,
rawAssetId: rawObjectKey,
textChars: 0,
});
} catch (error) {
// A concurrent bind may have committed the row first (the unique index
// rejected ours); adopt it, and drop the object we just stored when the
// winner does not reference it. Cleanup is best-effort: a failed delete
// must not fail a bind that already succeeded.
const winner = await store
.getMaterialByOwnerMaterialId(sessionId, ownerMaterialId)
.catch(() => null);
if (winner) {
if (winner.rawAssetId !== rawObjectKey) {
await byteStore.delete(rawObjectKey).catch((cleanupError) => {
console.warn(
`[session-materials] failed to delete losing bind object ${rawObjectKey}:`,
cleanupError,
);
});
}
return winner;
}
await byteStore.delete(rawObjectKey).catch(() => undefined);
throw error;
}
}
/**
* The HTTP-visible projection of one material row — the same shape the
* `list_materials` agent tool exposes. Object keys stay off the wire.
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@openmaic/storage",
"version": "0.29.2",
"version": "0.30.1",
"description": "The MAIC pluggable persistence layer: document / runtime / KV / asset primitives with browser and HTTP backends, depending only on @openmaic/dsl.",
"type": "module",
"main": "./dist/index.js",
+83 -4
View File
@@ -59,6 +59,7 @@ CREATE TABLE IF NOT EXISTS agent_session_materials (
session_id TEXT NOT NULL REFERENCES agent_sessions(id) ON DELETE CASCADE,
kind TEXT NOT NULL,
title TEXT,
owner_material_id TEXT,
source_url TEXT,
text_asset_id TEXT,
raw_asset_id TEXT,
@@ -81,9 +82,15 @@ CREATE TABLE IF NOT EXISTS agent_session_materials (
,CONSTRAINT agent_session_materials_extraction_attempts_nonnegative CHECK (extraction_attempts >= 0)
);
ALTER TABLE agent_session_materials ADD COLUMN IF NOT EXISTS owner_material_id TEXT;
CREATE INDEX IF NOT EXISTS agent_session_materials_session_created_idx
ON agent_session_materials (session_id, created_at);
CREATE UNIQUE INDEX IF NOT EXISTS agent_session_materials_session_owner_material_idx
ON agent_session_materials (session_id, owner_material_id)
WHERE owner_material_id IS NOT NULL;
CREATE INDEX IF NOT EXISTS agent_session_materials_extraction_queue_idx
ON agent_session_materials (created_at)
WHERE kind = 'source' AND extraction_status IN ('pending','running');
@@ -150,6 +157,7 @@ interface MaterialRow extends Record<string, unknown> {
session_id: string;
kind: string;
title: string | null;
owner_material_id: string | null;
source_url: string | null;
text_asset_id: string | null;
raw_asset_id: string | null;
@@ -172,6 +180,7 @@ function mapRow(row: MaterialRow): AgentSessionMaterial {
sessionId: row.session_id,
kind: row.kind as AgentSessionMaterialKind,
title: row.title,
ownerMaterialId: row.owner_material_id,
sourceUrl: row.source_url,
textAssetId: row.text_asset_id,
rawAssetId: row.raw_asset_id,
@@ -200,6 +209,15 @@ function isForeignKeyViolation(error: unknown): boolean {
);
}
function isUniqueViolation(error: unknown): boolean {
return (
typeof error === 'object' &&
error !== null &&
'code' in error &&
(error as { code?: unknown }).code === '23505'
);
}
export class PgAgentSessionMaterialStore implements AgentSessionMaterialStore {
private readonly queryable: Queryable;
private readonly tableNames: AgentSessionMaterialTableNames;
@@ -249,10 +267,11 @@ export class PgAgentSessionMaterialStore implements AgentSessionMaterialStore {
try {
const { rows } = await this.queryable.query<MaterialRow>(
`INSERT INTO ${this.table}
(id, session_id, kind, title, source_url, text_asset_id, raw_asset_id,
text_chars, derived_from, extraction_status, created_at)
SELECT $1, $2, $3, $4, $5, $6, $7, $8, $9,
CASE WHEN $3 = 'source' THEN 'idle' ELSE 'done' END, $10
(id, session_id, kind, title, owner_material_id, source_url,
text_asset_id, raw_asset_id, text_chars, derived_from,
extraction_status, created_at)
SELECT $1, $2, $3, $4, $5, $6, $7, $8, $9, $10,
CASE WHEN $3 = 'source' THEN 'idle' ELSE 'done' END, $11
FROM agent_sessions AS session
WHERE session.id = $2 AND session.deleted_at IS NULL
RETURNING *`,
@@ -261,6 +280,7 @@ export class PgAgentSessionMaterialStore implements AgentSessionMaterialStore {
sessionId,
input.kind,
input.title ?? null,
input.ownerMaterialId ?? null,
input.sourceUrl ?? null,
input.textAssetId ?? null,
input.rawAssetId ?? null,
@@ -329,6 +349,65 @@ export class PgAgentSessionMaterialStore implements AgentSessionMaterialStore {
return result.rows[0] ? mapRow(result.rows[0]) : null;
}
/**
* Resolve the session's row for one owner-library upload, or `null`.
*
* This is intentionally not part of the `AgentSessionMaterialStore`
* interface: only the host's owner-material binder needs it, to keep a
* rebind idempotent while row ids stay globally unique. The unique
* `(session_id, owner_material_id)` index makes at most one row match.
*/
async getMaterialByOwnerMaterialId(
sessionId: string,
ownerMaterialId: string,
): Promise<AgentSessionMaterial | null> {
const result = await this.queryable.query<MaterialRow>(
`SELECT material.* FROM ${this.table} AS material
INNER JOIN agent_sessions AS session ON session.id = material.session_id
WHERE material.owner_material_id = $1 AND material.session_id = $2
AND session.deleted_at IS NULL
LIMIT 1`,
[ownerMaterialId, sessionId],
);
return result.rows[0] ? mapRow(result.rows[0]) : null;
}
/**
* Backfill `owner_material_id` on a row the pre-upgrade binder wrote with
* `id = ownerMaterialId` and no `owner_material_id`. The host validates the
* form of the legacy row against its owner record before calling this, so
* only an already-verified binding is stamped.
*
* The partial unique index on `(session_id, owner_material_id)` can still
* reject the stamp when a concurrent bind already claimed the owner id for
* this session; that is an expected race, reported as `null` rather than an
* error so the caller can adopt the winner. Returns the updated row, or
* `null` when the row no longer matches (already upgraded, claimed, or
* gone). Extraction columns are deliberately untouched.
*/
async backfillOwnerMaterialId(
sessionId: string,
materialId: string,
ownerMaterialId: string,
): Promise<AgentSessionMaterial | null> {
try {
const result = await this.queryable.query<MaterialRow>(
`UPDATE ${this.table} AS material
SET owner_material_id = $3
FROM agent_sessions AS session
WHERE material.id = $1 AND material.session_id = $2
AND material.owner_material_id IS NULL
AND material.session_id = session.id AND session.deleted_at IS NULL
RETURNING material.*`,
[materialId, sessionId, ownerMaterialId],
);
return result.rows[0] ? mapRow(result.rows[0]) : null;
} catch (error) {
if (isUniqueViolation(error)) return null;
throw error;
}
}
async enqueueExtraction(sessionId: string, materialId: string): Promise<boolean> {
const result = await this.queryable.query(
`UPDATE ${this.table} AS material
@@ -92,6 +92,13 @@ export interface AgentSessionMaterial {
sessionId: string;
kind: AgentSessionMaterialKind;
title: string | null;
/**
* The owner-library upload this row was bound from, when it came from the
* owner material library (null for fetched/derived rows). The row `id` is
* minted per session, so the same owner upload bound to several sessions
* yields several rows that all record the same `ownerMaterialId`.
*/
ownerMaterialId?: string | null;
/** The fetch's source URL; never a model-invented target. */
sourceUrl: string | null;
/** Asset id (registry) of the extracted text/markdown bytes. */
@@ -110,6 +117,12 @@ export interface AgentSessionMaterial {
export interface CreateAgentSessionMaterialInput {
/** Caller-minted stable id; defaults to a fresh `mat_` id. */
id?: string;
/**
* The owner-library upload this row is bound from, when applicable. A
* unique `(session_id, owner_material_id)` index makes a rebind idempotent
* without requiring the shared owner id to be globally unique.
*/
ownerMaterialId?: string;
kind: AgentSessionMaterialKind;
title?: string;
sourceUrl?: string;
@@ -88,6 +88,49 @@ export function runAgentSessionMaterialContract(
).rejects.toThrow();
});
test('same owner material bound to two sessions succeeds and both sessions can read it', async () => {
const store = makeStore();
await store.createSession({ id: 'session-1', ownerId: 'owner-a', prompt: 'p' });
await store.createSession({ id: 'session-2', ownerId: 'owner-a', prompt: 'p' });
// One owner-library upload reused across sessions: each session gets its
// own globally unique row id, so the shared owner id is metadata rather
// than the primary key.
const first = await store.createMaterial('session-1', {
kind: 'source',
title: 'textbook.pdf',
ownerMaterialId: 'mat_owner',
});
const second = await store.createMaterial('session-2', {
kind: 'source',
title: 'textbook.pdf',
ownerMaterialId: 'mat_owner',
});
expect(first.id).not.toBe(second.id);
expect(first.ownerMaterialId).toBe('mat_owner');
expect(second.ownerMaterialId).toBe('mat_owner');
// Both sessions read their own binding, and neither leaks the other's.
await expect(store.getMaterial('session-1', first.id)).resolves.toMatchObject({
id: first.id,
ownerMaterialId: 'mat_owner',
});
await expect(store.getMaterial('session-2', second.id)).resolves.toMatchObject({
id: second.id,
ownerMaterialId: 'mat_owner',
});
await expect(store.getMaterial('session-1', second.id)).resolves.toBeNull();
// Rebinding the same owner upload into one session is refused by the
// unique (session_id, owner_material_id) index.
await expect(
store.createMaterial('session-1', {
kind: 'source',
title: 'textbook.pdf',
ownerMaterialId: 'mat_owner',
}),
).rejects.toThrow();
});
test('lists newest-first and pages with a keyset before cursor', async () => {
const store = makeStore();
await store.createSession({ id: 'session-1', ownerId: 'owner-a', prompt: 'p' });
@@ -70,6 +70,70 @@ describe('PgAgentSessionMaterialStore with PGlite', () => {
);
});
test('adds owner_material_id to a table provisioned before the column existed', async () => {
// Simulate a 1.0.2 database: the table predates the owner-material column,
// so `CREATE TABLE IF NOT EXISTS` cannot add it. CASCADE also drops the
// partial unique index that depends on the column.
await db.query('ALTER TABLE agent_session_materials DROP COLUMN owner_material_id CASCADE');
const sessions = new PgAgentSessionStore(db, {
withTransaction: (body) => db.transaction((tx: Queryable) => body(tx)),
});
await sessions.createSession({ id: 'session-legacy', ownerId: 'owner-a', prompt: 'p' });
await db.query(
`INSERT INTO agent_session_materials
(id, session_id, kind, title, text_chars, extraction_status, extraction_attempts, created_at)
VALUES ('mat_legacy', 'session-legacy', 'web', 'legacy', 0, 'done', 0, now())`,
);
// Re-running the initializer is the in-place upgrade.
await ensureAgentSessionMaterialSchema(db);
const legacy = await store.getMaterial('session-legacy', 'mat_legacy');
expect(legacy).toMatchObject({ id: 'mat_legacy', ownerMaterialId: null });
// The upgraded table now accepts two sessions binding one owner upload.
await sessions.createSession({ id: 'session-a', ownerId: 'owner-a', prompt: 'p' });
await sessions.createSession({ id: 'session-b', ownerId: 'owner-a', prompt: 'p' });
const first = await store.createMaterial('session-a', {
kind: 'source',
title: 'textbook.pdf',
ownerMaterialId: 'mat_owner',
});
const second = await store.createMaterial('session-b', {
kind: 'source',
title: 'textbook.pdf',
ownerMaterialId: 'mat_owner',
});
expect(first.id).not.toBe(second.id);
expect(first.ownerMaterialId).toBe('mat_owner');
expect(second.ownerMaterialId).toBe('mat_owner');
});
test('backfills owner_material_id onto a legacy row without touching extraction', async () => {
await new PgAgentSessionStore(db, {
withTransaction: (body) => db.transaction((tx: Queryable) => body(tx)),
}).createSession({ id: 'session-1', ownerId: 'owner-a', prompt: 'p' });
// The pre-upgrade binder's shape: id = owner id, owner_material_id NULL.
await db.query(
`INSERT INTO agent_session_materials
(id, session_id, kind, title, raw_asset_id, text_chars, extraction_status,
extraction_attempts, extractor_version, created_at)
VALUES ('mat_owner', 'session-1', 'source', 'textbook.pdf',
'materials/session-1/mat_owner/raw.x', 0, 'done', 2, 'pdf@1', now())`,
);
const backfilled = await store.backfillOwnerMaterialId('session-1', 'mat_owner', 'mat_owner');
expect(backfilled).toMatchObject({
id: 'mat_owner',
ownerMaterialId: 'mat_owner',
rawAssetId: 'materials/session-1/mat_owner/raw.x',
extraction: { status: 'done', attempts: 2, extractorVersion: 'pdf@1' },
});
// The NULL predicate means a repeat is a no-op, not an error.
expect(await store.backfillOwnerMaterialId('session-1', 'mat_owner', 'mat_owner')).toBeNull();
expect(await store.backfillOwnerMaterialId('session-1', 'mat_absent', 'mat_absent')).toBeNull();
});
test('cascades material rows away when the session row is hard-deleted', async () => {
await new PgAgentSessionStore(db, {
withTransaction: (body) => db.transaction((tx: Queryable) => body(tx)),
@@ -476,6 +476,7 @@ CREATE TABLE IF NOT EXISTS agent_session_materials (
session_id TEXT NOT NULL REFERENCES agent_sessions(id) ON DELETE CASCADE,
kind TEXT NOT NULL,
title TEXT,
owner_material_id TEXT,
source_url TEXT,
text_asset_id TEXT,
raw_asset_id TEXT,
@@ -498,9 +499,15 @@ CREATE TABLE IF NOT EXISTS agent_session_materials (
,CONSTRAINT agent_session_materials_extraction_attempts_nonnegative CHECK (extraction_attempts >= 0)
);
ALTER TABLE agent_session_materials ADD COLUMN IF NOT EXISTS owner_material_id TEXT;
CREATE INDEX IF NOT EXISTS agent_session_materials_session_created_idx
ON agent_session_materials (session_id, created_at);
CREATE UNIQUE INDEX IF NOT EXISTS agent_session_materials_session_owner_material_idx
ON agent_session_materials (session_id, owner_material_id)
WHERE owner_material_id IS NOT NULL;
CREATE INDEX IF NOT EXISTS agent_session_materials_extraction_queue_idx
ON agent_session_materials (created_at)
WHERE kind = 'source' AND extraction_status IN ('pending','running');
@@ -0,0 +1,295 @@
/**
* Owner-material binding integration — the issue #1494 regression.
*
* Drives `bindOwnerMaterialsToSession` and the real session-creation route over
* a PGlite-backed durable store. Before the fix the second session's bind hit
* the global `agent_session_materials` primary key and the route answered 500;
* now each session gets its own row id while the shared owner upload id is
* recorded for idempotency, so both sessions bind and read their own row.
*/
import { PGlite } from '@electric-sql/pglite';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { NextRequest } from 'next/server';
import { PgAgentSessionStore, ensureAgentSessionSchema } from '@openmaic/storage/agent-session/pg';
import { ensureAgentSessionMaterialSchema } from '@openmaic/storage/material/pg';
import type { Queryable } from '@openmaic/storage/asset/pg';
import { setMaterialByteStoreForTests } from '@/lib/server/materials/bytes';
import { ensureOwnerMaterialSchema } from '@/lib/persistence/owner-materials';
const mocks = vi.hoisted(() => ({
getAgentSessionStore: vi.fn(),
getServerPersistenceProvider: vi.fn(),
resolveRequestOwnerId: vi.fn(),
scheduleConversationTitle: vi.fn(),
}));
vi.mock('@/lib/config/feature-flags', () => ({
isAgentRuntimeEnabled: () => true,
isAgentRuntimeConfigured: () => true,
}));
vi.mock('@/lib/server/agent-runtime/owner', () => ({
resolveRequestOwnerId: mocks.resolveRequestOwnerId,
}));
vi.mock('@/lib/server/agent-runtime/skills', () => ({
listSkills: async () => [],
findSkill: async () => null,
inferSkillIdFromPrompt: async () => undefined,
}));
vi.mock('@/lib/server/agent-runtime/store', () => ({
getAgentSessionStore: mocks.getAgentSessionStore,
}));
vi.mock('@/lib/server/agent-runtime/conversation-title-task', () => ({
scheduleConversationTitle: mocks.scheduleConversationTitle,
}));
vi.mock('@/lib/persistence/server-provider', () => ({
getServerPersistenceProvider: mocks.getServerPersistenceProvider,
}));
import { POST } from '@/app/api/agent/sessions/route';
import {
bindOwnerMaterialsToSession,
getSessionMaterial,
listSessionMaterials,
} from '@/lib/server/agent-runtime/session-materials';
let dbCounter = 0;
let db: PGlite | undefined;
async function makeHost() {
const instance = new PGlite();
await instance.waitReady;
await ensureAgentSessionSchema(instance);
await ensureOwnerMaterialSchema(instance);
await ensureAgentSessionMaterialSchema(instance);
const bytes = new Map<string, Buffer>();
const puts: string[] = [];
setMaterialByteStoreForTests({
put: async (key, body) => {
bytes.set(key, Buffer.from(body as Uint8Array));
puts.push(key);
},
get: async (key) => {
const value = bytes.get(key);
if (!value) throw new Error(`missing material bytes: ${key}`);
return value;
},
delete: async (key) => void bytes.delete(key),
});
const sessionStore = new PgAgentSessionStore(instance, {
withTransaction: (body) => instance.transaction((tx: Queryable) => body(tx)),
});
dbCounter += 1;
vi.stubEnv('DATABASE_URL', `postgres://binding-${dbCounter}`);
mocks.getAgentSessionStore.mockResolvedValue(sessionStore);
mocks.getServerPersistenceProvider.mockResolvedValue({ pool: instance });
mocks.resolveRequestOwnerId.mockImplementation((_request: NextRequest, headers: Headers) => {
headers.append('Set-Cookie', 'anonymous_id=test; Path=/; HttpOnly');
return 'owner-1';
});
db = instance;
return { db: instance, bytes, puts, sessionStore };
}
async function seedOwnerMaterial(instance: PGlite, id: string) {
await instance.query(
`INSERT INTO owner_material
(id, owner_id, kind, mime, bytes, original_name, oss_key, status, extraction, created_at)
VALUES ($1, 'owner-1', 'source', 'application/pdf', 3, 'textbook.pdf', $2, 'ready', NULL, $3)`,
[id, `owner/${id}/raw`, Date.now()],
);
}
/** The deterministic object key the pre-upgrade binder copied owner bytes to. */
function legacyRawKey(
sessionId: string,
ownerMaterialId: string,
mime = 'application/pdf',
): string {
return `materials/${sessionId}/${ownerMaterialId}/raw.${Buffer.from(mime, 'utf8').toString(
'base64url',
)}`;
}
function post(body: unknown) {
return POST(
new NextRequest('http://localhost/api/agent/sessions', {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify(body),
}),
);
}
beforeEach(() => {
vi.clearAllMocks();
vi.unstubAllEnvs();
setMaterialByteStoreForTests(null);
});
afterEach(async () => {
await db?.close();
db = undefined;
});
describe('owner-material binding across sessions', () => {
it('binds one owner upload to two sessions and both can read their own row', async () => {
const { db: instance, bytes, sessionStore } = await makeHost();
await sessionStore.createSession({ id: 'session-a', ownerId: 'owner-1', prompt: 'p' });
await sessionStore.createSession({ id: 'session-b', ownerId: 'owner-1', prompt: 'p' });
await seedOwnerMaterial(instance, 'mat_owner');
bytes.set('owner/mat_owner/raw', Buffer.from('PDF'));
const first = await bindOwnerMaterialsToSession('session-a', 'owner-1', ['mat_owner']);
const second = await bindOwnerMaterialsToSession('session-b', 'owner-1', ['mat_owner']);
expect(first).toHaveLength(1);
expect(second).toHaveLength(1);
// Both sessions got distinct session-side rows for the same owner upload.
expect(first[0]!.materialId).not.toBe(second[0]!.materialId);
const rowA = await getSessionMaterial('session-a', first[0]!.materialId);
const rowB = await getSessionMaterial('session-b', second[0]!.materialId);
expect(rowA).toMatchObject({
id: first[0]!.materialId,
sessionId: 'session-a',
ownerMaterialId: 'mat_owner',
title: 'textbook.pdf',
});
expect(rowB).toMatchObject({
id: second[0]!.materialId,
sessionId: 'session-b',
ownerMaterialId: 'mat_owner',
title: 'textbook.pdf',
});
// Reads stay session-scoped: neither session can read the other's row.
expect(await getSessionMaterial('session-a', second[0]!.materialId)).toBeNull();
expect(await getSessionMaterial('session-b', first[0]!.materialId)).toBeNull();
// Rebinding the same owner upload into the same session is idempotent.
const rebound = await bindOwnerMaterialsToSession('session-a', 'owner-1', ['mat_owner']);
expect(rebound[0]!.materialId).toBe(first[0]!.materialId);
expect(await listSessionMaterials('session-a')).toHaveLength(1);
});
it('reuses and backfills a pre-upgrade legacy owner-material binding', async () => {
const { db: instance, bytes, puts, sessionStore } = await makeHost();
await sessionStore.createSession({ id: 'session-legacy', ownerId: 'owner-1', prompt: 'p' });
await sessionStore.createSession({ id: 'session-other', ownerId: 'owner-1', prompt: 'p' });
await seedOwnerMaterial(instance, 'mat_owner');
bytes.set('owner/mat_owner/raw', Buffer.from('PDF'));
// Exactly what the previous binder wrote: row id = owner upload id,
// owner_material_id NULL, copied bytes at the deterministic legacy key,
// and extraction already finished to prove the state must survive.
const legacyKey = legacyRawKey('session-legacy', 'mat_owner');
bytes.set(legacyKey, Buffer.from('PDF'));
await instance.query(
`INSERT INTO agent_session_materials
(id, session_id, kind, title, owner_material_id, raw_asset_id, text_chars,
extraction_status, extraction_attempts, extraction_stats, extractor_version, created_at)
VALUES ('mat_owner', 'session-legacy', 'source', 'textbook.pdf', NULL, $1, 0,
'done', 2, $2::jsonb, 'pdf@1', now())`,
[legacyKey, JSON.stringify({ chars: 1234, pages: 2, imageCount: 0 })],
);
const putsBefore = puts.length;
const rebound = await bindOwnerMaterialsToSession('session-legacy', 'owner-1', ['mat_owner']);
// The legacy row is reused, not duplicated or re-copied.
expect(rebound).toHaveLength(1);
expect(rebound[0]!.materialId).toBe('mat_owner');
expect(await listSessionMaterials('session-legacy')).toHaveLength(1);
expect(puts).toHaveLength(putsBefore);
expect(bytes.get(legacyKey)).toEqual(Buffer.from('PDF'));
const row = await getSessionMaterial('session-legacy', 'mat_owner');
expect(row).toMatchObject({
id: 'mat_owner',
ownerMaterialId: 'mat_owner',
title: 'textbook.pdf',
rawAssetId: legacyKey,
extraction: {
status: 'done',
attempts: 2,
extractorVersion: 'pdf@1',
stats: { chars: 1234, pages: 2, imageCount: 0 },
},
});
// The backfill lets the next bind take the fast path without copying.
const again = await bindOwnerMaterialsToSession('session-legacy', 'owner-1', ['mat_owner']);
expect(again[0]!.materialId).toBe('mat_owner');
expect(await listSessionMaterials('session-legacy')).toHaveLength(1);
expect(puts).toHaveLength(putsBefore);
// A different session is still a fresh row with its own byte copy.
const other = await bindOwnerMaterialsToSession('session-other', 'owner-1', ['mat_owner']);
expect(other[0]!.materialId).not.toBe('mat_owner');
expect(await listSessionMaterials('session-other')).toHaveLength(1);
expect(await getSessionMaterial('session-other', other[0]!.materialId)).toMatchObject({
ownerMaterialId: 'mat_owner',
rawAssetId: expect.any(String),
});
expect(puts).toHaveLength(putsBefore + 1);
});
it('removes the losing upload when a concurrent bind wins the same session', async () => {
const { db: instance, sessionStore } = await makeHost();
await sessionStore.createSession({ id: 'session-race', ownerId: 'owner-1', prompt: 'p' });
await seedOwnerMaterial(instance, 'mat_owner');
const ownerBytes = new Map<string, Buffer>([['owner/mat_owner/raw', Buffer.from('PDF')]]);
const sessionBytes = new Map<string, Buffer>();
let parkedLoser = true;
let winner: Awaited<ReturnType<typeof bindOwnerMaterialsToSession>> | undefined;
// The loser's upload of its own object is the pause point: the winner runs
// to completion there, so the loser is guaranteed to lose the unique index
// and must clean up the object it already stored.
setMaterialByteStoreForTests({
put: async (key, body) => {
sessionBytes.set(key, Buffer.from(body as Uint8Array));
if (parkedLoser) {
parkedLoser = false;
winner = await bindOwnerMaterialsToSession('session-race', 'owner-1', ['mat_owner']);
}
},
get: async (key) => {
const value = sessionBytes.get(key) ?? ownerBytes.get(key);
if (!value) throw new Error(`missing material bytes: ${key}`);
return value;
},
delete: async (key) => {
sessionBytes.delete(key);
ownerBytes.delete(key);
},
});
const loser = await bindOwnerMaterialsToSession('session-race', 'owner-1', ['mat_owner']);
expect(winner).toBeDefined();
expect(loser).toHaveLength(1);
expect(loser[0]!.materialId).toBe(winner![0]!.materialId);
const rows = await listSessionMaterials('session-race');
expect(rows).toHaveLength(1);
expect(rows[0]!.ownerMaterialId).toBe('mat_owner');
// Exactly one session byte object remains and it is the winner's.
expect([...sessionBytes.keys()]).toEqual([rows[0]!.rawAssetId]);
expect(ownerBytes.get('owner/mat_owner/raw')).toEqual(Buffer.from('PDF'));
});
it('POST /api/agent/sessions returns 202 when a second session reuses the upload', async () => {
const { db: instance, bytes } = await makeHost();
await seedOwnerMaterial(instance, 'mat_owner');
bytes.set('owner/mat_owner/raw', Buffer.from('PDF'));
const first = await post({ prompt: 'Build a course', materialIds: ['mat_owner'] });
expect(first.status).toBe(202);
// The regression: before the fix this second bind threw the primary-key
// violation, `withRequestOwnerId` swallowed it, and the response was 500.
const second = await post({ prompt: 'Build the sequel', materialIds: ['mat_owner'] });
expect(second.status).toBe(202);
});
});