Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions docs/api-reference/openapi.json
Original file line number Diff line number Diff line change
Expand Up @@ -6864,6 +6864,21 @@
"description": "A persisted D1 record. Database timestamps are Unix milliseconds.",
"additionalProperties": true
},
"ProjectSourceLink": {
"type": "object",
"required": [
"item",
"added"
],
"properties": {
"item": {
"$ref": "#/components/schemas/ProjectItem"
},
"added": {
"type": "boolean"
}
}
},
"CreateProjectRequest": {
"type": "object",
"required": [
Expand Down
2 changes: 2 additions & 0 deletions docs/internals/endpoints.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ Called by the video2ctx web application or an explicit signed-in account action.
| `POST` | `/v1/notification-preferences/confirm-email` | `confirmNotificationEmail` | Confirm monitor email alerts from the signed-in dashboard | `sessionCookie` or `cliSession` or `bearerApiKey` or `apiKey` or `demoUser` | Enables email delivery only for the signed-in account after validating the confirmation token. |
| `DELETE` | `/v1/oauth/youtube` | `disconnectYouTube` | Disconnect the YouTube account | `sessionCookie` or `demoUser` | Mutates the account connection state; require an explicit user action. |
| `GET` | `/v1/oauth/youtube/connect` | `createYouTubeConnectUrl` | Create a YouTube OAuth URL | `sessionCookie` or `demoUser` | Returns a state-bound OAuth URL; start it only from a user-initiated connection flow. |
| `POST` | `/v1/projects/{id}/sources` | `linkProjectSource` | Add a saved Sources search or inspection to a project | `sessionCookie` or `demoUser` | Browser-session only. Verify that the project and saved source belong to the same signed-in user before linking shared asset references. |
| `GET` | `/v1/projects/{id}/sources/{itemId}` | `getProjectSource` | Restore a project source snapshot | `sessionCookie` or `demoUser` | Browser-session only. Verify project ownership and restore the item from that project without modifying saved references. |
| `POST` | `/v1/resolve` | `resolveInput` | Route universal UI input | `sessionCookie` or `cliSession` or `bearerApiKey` or `apiKey` or `demoUser` | First-party input router; its dispatch behavior is not a stable public API contract. |
| `POST` | `/v1/scale-inquiries` | `submitScaleInquiry` | Submit a Scale plan inquiry | Public or protocol-signed | Public lead form; validate Turnstile and rate limits, and never let the caller choose the notification recipient. |
| `GET` | `/v1/sources/recent` | `listRecentSources` | List recent Sources inputs | `sessionCookie` or `demoUser` | Browser-session only. Lists inputs belonging to the authenticated user without exposing shared asset references. |
Expand Down
127 changes: 122 additions & 5 deletions platform/src/durable-objects/user-account.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { AgentAdmissionQueue } from '../agents/runtime/admission-queue';
import type { AgentRequest, AgentAdmission } from '../agents/contracts';
import { DurableObject } from 'cloudflare:workers';
import { z } from 'zod';
import { RECENT_SOURCE_LIMIT, saveReferencedSourceSchema, sourceReferenceSchema, sourceIdentity, type RecentSource, type SaveReferencedSource, type SourceReference } from '../lib/source-history';
import { RECENT_SOURCE_LIMIT, saveReferencedSourceSchema, sourceReferenceSchema, sourceIdentity, sourceIdSchema, type RecentSource, type SaveReferencedSource, type SourceReference } from '../lib/source-history';

const MAX_SEARCH_TEXT_LENGTH = 32_000;
const MAX_TITLE_LENGTH = 80;
Expand Down Expand Up @@ -46,6 +46,26 @@ export interface UserSessionPage {
nextCursor: UserSessionCursor | null;
}

export interface ProjectSourceItem {
id: string;
source_id: string;
provider: 'youtube';
entity_type: 'search' | 'video' | 'playlist' | 'channel';
entity_id: string;
title: string;
created_at: number;
}

interface ProjectSourceRow extends Record<string, SqlStorageValue> {
id: string;
source_id: string;
input: string;
title: string;
kind: RecentSource['kind'];
snapshot: string;
created_at: number;
}

interface SessionRow extends Record<string, SqlStorageValue> {
conversation_id: string;
title: string;
Expand Down Expand Up @@ -102,6 +122,7 @@ export class UserAccountDO extends DurableObject<Env> {
this.ctx.storage.sql.exec('DELETE FROM user_sessions');
this.ctx.storage.sql.exec('DELETE FROM agent_conversations');
this.ctx.storage.sql.exec('DELETE FROM recent_sources');
this.ctx.storage.sql.exec('DELETE FROM project_sources');
// Keep only a tombstone so already-authenticated requests cannot recreate data.
}

Expand Down Expand Up @@ -153,10 +174,99 @@ export class UserAccountDO extends DurableObject<Env> {

getSource(id: string): { source: RecentSource; snapshot: SourceReference } | null {
this.assertActive();
const row = this.ctx.storage.sql.exec<{ snapshot: string }>('SELECT snapshot FROM recent_sources WHERE id = ?', z.string().uuid().parse(id)).toArray()[0];
if (!row) return null;
this.ctx.storage.sql.exec('UPDATE recent_sources SET updated_at = ? WHERE id = ?', this.nextSourceUpdate(), id);
return { source: this.listSources().find(source => source.id === id)!, snapshot: sourceReferenceSchema.parse(JSON.parse(row.snapshot)) };
const sourceId = sourceIdSchema.parse(id);
const row = this.ctx.storage.sql.exec<{ id: string; input: string; title: string; kind: RecentSource['kind']; updated_at: number; snapshot: string }>(
'SELECT id, input, title, kind, updated_at, snapshot FROM recent_sources WHERE id = ?', sourceId,
).toArray()[0];
if (row) {
const updatedAt = this.nextSourceUpdate();
this.ctx.storage.sql.exec('UPDATE recent_sources SET updated_at = ? WHERE id = ?', updatedAt, id);
return { source: { id: row.id, input: row.input, title: row.title, kind: row.kind, updatedAt },
snapshot: sourceReferenceSchema.parse(JSON.parse(row.snapshot)) };
}
return null;
}

saveSourceWithProject(value: SaveReferencedSource, projectId?: string) {
this.assertActive();
const project = projectId === undefined ? undefined : z.string().uuid().parse(projectId);
return this.ctx.storage.transactionSync(() => {
const source = this.saveSource(value);
const linked = project ? this.linkSourceToProject(project, source.id) : null;
if (project && !linked) throw new Error('Saved source could not be linked.');
return { source, linked };
});
}

getProjectSource(projectId: string, itemId: string): { source: RecentSource; snapshot: SourceReference } | null {
this.assertActive();
const saved = this.ctx.storage.sql.exec<ProjectSourceRow>(
'SELECT id, source_id, input, title, kind, snapshot, created_at FROM project_sources WHERE project_id = ? AND id = ?',
z.string().uuid().parse(projectId), sourceIdSchema.parse(itemId),
).toArray()[0];
if (!saved) return null;
return { source: { id: saved.source_id, input: saved.input, title: saved.title, kind: saved.kind, updatedAt: saved.created_at },
snapshot: sourceReferenceSchema.parse(JSON.parse(saved.snapshot)) };
}

linkSourceToProject(projectId: string, sourceId: string): { item: ProjectSourceItem; added: boolean } | null {
this.assertActive();
const project = z.string().uuid().parse(projectId);
const source = sourceIdSchema.parse(sourceId);
const recent = this.ctx.storage.sql.exec<{ source_key: string; input: string; title: string; kind: RecentSource['kind']; snapshot: string }>(
'SELECT source_key, input, title, kind, snapshot FROM recent_sources WHERE id = ?', source,
).toArray()[0];
if (!recent) {
const saved = this.ctx.storage.sql.exec<ProjectSourceRow>(
'SELECT id, source_id, input, title, kind, snapshot, created_at FROM project_sources WHERE project_id = ? AND source_id = ?', project, source,
).toArray()[0];
return saved ? { item: this.toProjectSourceItem(saved), added: false } : null;
}
const existing = this.ctx.storage.sql.exec<{ id: string }>(
'SELECT id FROM project_sources WHERE project_id = ? AND source_key = ?', project, recent.source_key,
).toArray()[0];
const id = existing?.id ?? crypto.randomUUID();
const createdAt = Date.now();
this.ctx.storage.sql.exec(`INSERT INTO project_sources
(id, project_id, source_key, source_id, input, title, kind, snapshot, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(project_id, source_key) DO UPDATE SET
source_id=excluded.source_id, input=excluded.input, title=excluded.title,
kind=excluded.kind, snapshot=excluded.snapshot`,
id, project, recent.source_key, source, recent.input, recent.title, recent.kind, recent.snapshot, createdAt);
const saved = this.ctx.storage.sql.exec<ProjectSourceRow>(
'SELECT id, source_id, input, title, kind, snapshot, created_at FROM project_sources WHERE id = ?', id,
).one();
return { item: this.toProjectSourceItem(saved), added: !existing };
}

listProjectSources(projectId: string): ProjectSourceItem[] {
this.assertActive();
const rows = this.ctx.storage.sql.exec<ProjectSourceRow>(
'SELECT id, source_id, input, title, kind, snapshot, created_at FROM project_sources WHERE project_id = ? ORDER BY created_at DESC',
z.string().uuid().parse(projectId),
).toArray();
return rows.map(row => this.toProjectSourceItem(row));
}

projectSourceCounts(): Array<{ projectId: string; count: number }> {
this.assertActive();
return this.ctx.storage.sql.exec<{ project_id: string; count: number }>(
'SELECT project_id, COUNT(*) AS count FROM project_sources GROUP BY project_id',
).toArray().map(row => ({ projectId: row.project_id, count: row.count }));
}

removeProjectSources(projectId: string): void {
this.assertActive();
this.ctx.storage.sql.exec('DELETE FROM project_sources WHERE project_id = ?', z.string().uuid().parse(projectId));
}

private toProjectSourceItem(row: ProjectSourceRow): ProjectSourceItem {
const snapshot = sourceReferenceSchema.parse(JSON.parse(row.snapshot));
return { id: row.id, source_id: row.source_id, provider: 'youtube',
entity_type: snapshot.kind === 'search' ? 'search' : snapshot.inspector.type,
entity_id: snapshot.kind === 'search' ? row.source_id : snapshot.inspector.id,
title: row.title, created_at: row.created_at };
}

private nextSourceUpdate(): number {
Expand Down Expand Up @@ -342,6 +452,13 @@ export class UserAccountDO extends DurableObject<Env> {
id TEXT PRIMARY KEY, source_key TEXT NOT NULL UNIQUE, input TEXT NOT NULL, title TEXT NOT NULL,
kind TEXT NOT NULL, updated_at INTEGER NOT NULL, snapshot TEXT NOT NULL
)`);
this.ctx.storage.sql.exec(`CREATE TABLE IF NOT EXISTS project_sources (
id TEXT PRIMARY KEY, project_id TEXT NOT NULL, source_key TEXT NOT NULL,
source_id TEXT NOT NULL, input TEXT NOT NULL, title TEXT NOT NULL,
kind TEXT NOT NULL, snapshot TEXT NOT NULL, created_at INTEGER NOT NULL,
UNIQUE(project_id, source_key)
)`);
this.ctx.storage.sql.exec('CREATE INDEX IF NOT EXISTS project_sources_project_idx ON project_sources (project_id, created_at DESC)');
this.ctx.storage.sql.exec('CREATE TABLE IF NOT EXISTS account_deletion (id INTEGER PRIMARY KEY)');
this.ctx.storage.sql.exec('CREATE TABLE IF NOT EXISTS agent_conversations (conversation_id TEXT PRIMARY KEY)');
this.ctx.storage.sql.exec(`
Expand Down
11 changes: 6 additions & 5 deletions platform/src/lib/exports.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { ApiError, now } from './http';
import { listProjectItems } from './project-items';

type Format = 'txt' | 'md' | 'json' | 'csv' | 'srt' | 'vtt';

Expand All @@ -17,11 +18,11 @@ export async function createProjectExport(env: Env, userId: string, projectId: s
const project = await env.DB.prepare('SELECT name,description FROM projects WHERE id=? AND user_id=?')
.bind(projectId, userId).first<{ name: string; description: string }>();
if (!project) throw new ApiError(404, 'PROJECT_NOT_FOUND', 'Project not found.');
const rows = await env.DB.prepare(
`SELECT provider,title,entity_type,entity_id,start_ms,end_ms,note,tags_json
FROM project_items WHERE project_id=? AND user_id=? ORDER BY created_at`
).bind(projectId, userId).all<ItemRow>();
const content = serialize(format, project, rows.results);
const items = await listProjectItems(env, userId, projectId);
const rows: ItemRow[] = items.reverse().map(({ provider, title, entity_type, entity_id, start_ms, end_ms, note, tags_json }) => ({
provider, title, entity_type, entity_id, start_ms, end_ms, note, tags_json,
}));
const content = serialize(format, project, rows);
const id = crypto.randomUUID();
const key = `private/${userId}/exports/${id}.${format}`;
await env.RESEARCH.put(key, content, { httpMetadata: { contentType: contentType(format) } });
Expand Down
18 changes: 18 additions & 0 deletions platform/src/lib/project-items.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
import { userAccountInstanceName } from '../agents/runtime/identity';

export interface ProjectItemRecord {
id: string; provider: string; entity_type: string; entity_id: string; title: string;
start_ms: number | null; end_ms: number | null; note: string; tags_json: string; created_at: number;
source_id?: string;
}

// Both dashboard detail and exports include legacy items and saved source references.
export async function listProjectItems(env: Env, userId: string, projectId: string): Promise<ProjectItemRecord[]> {
const legacy = await env.DB.prepare('SELECT * FROM project_items WHERE project_id=? AND user_id=? ORDER BY created_at DESC')
.bind(projectId, userId).all<ProjectItemRecord>();
const account = env.USER_ACCOUNT.getByName(await userAccountInstanceName(userId));
const sources = (await account.listProjectSources(projectId)).map(item => ({
...item, start_ms: null, end_ms: null, note: '', tags_json: '[]',
}));
return [...legacy.results, ...sources].sort((a, b) => b.created_at - a.created_at);
}
2 changes: 1 addition & 1 deletion platform/src/lib/source-history.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ export const sourceSnapshotSchema = z.discriminatedUnion('kind', [
})) }),
z.object({ kind: z.literal('inspection'), inspector }),
]);
export const saveSourceSchema = z.object({ input: z.string().trim().min(1).max(500), snapshot: z.discriminatedUnion('kind', [
export const saveSourceSchema = z.object({ projectId: z.string().uuid().optional(), input: z.string().trim().min(1).max(500), snapshot: z.discriminatedUnion('kind', [
z.object({ kind: z.literal('search'), selectedData: z.array(dataset).min(1) }),
z.object({ kind: z.literal('inspection'), inspector: inspector.pick({ provider: true, type: true, id: true, requestedData: true, dataErrors: true })
.extend({ loadedData: z.array(z.enum(['metadata', 'transcript', 'comments', 'channel'])) }) }),
Expand Down
4 changes: 4 additions & 0 deletions platform/src/openapi-audience.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ export const OPENAPI_OPERATION_AUDIENCE: Readonly<Record<string, OpenApiAudience
listRecentSources: 'first-party',
saveRecentSource: 'first-party',
getRecentSource: 'first-party',
linkProjectSource: 'first-party',
getProjectSource: 'first-party',
startAgentRun: 'consumer',
getAgentAccess: 'consumer',
listAgentSessions: 'consumer',
Expand Down Expand Up @@ -104,6 +106,8 @@ export const OPENAPI_INTERNAL_SAFETY: Readonly<Record<string, string>> = {
listRecentSources: 'Browser-session only. Lists inputs belonging to the authenticated user without exposing shared asset references.',
saveRecentSource: 'Browser-session only. Resolves existing provider assets server-side; never accepts client-supplied R2 keys or video payloads.',
getRecentSource: 'Browser-session only. Verifies user ownership before reading the saved shared references.',
getProjectSource: 'Browser-session only. Verify project ownership and restore the item from that project without modifying saved references.',
linkProjectSource: 'Browser-session only. Verify that the project and saved source belong to the same signed-in user before linking shared asset references.',
inspectLandingYouTubeVideo: 'Public, rate-limited demo route; do not use it as a credentialed bulk-data API.',
submitScaleInquiry: 'Public lead form; validate Turnstile and rate limits, and never let the caller choose the notification recipient.',
signInWithMagicLink: 'Sends account email; rate-limit callers and never disclose whether an address is registered.',
Expand Down
38 changes: 35 additions & 3 deletions platform/src/openapi.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1095,10 +1095,10 @@ export const openApiDocument = {
},
post: {
tags: ['Projects'], operationId: 'saveRecentSource', summary: 'Remember a Sources search or inspection', security: privateSecurity,
description: 'Stores the input and references in the user Durable Object. Provider data is read from existing shared storage. This does not fetch YouTube data.',
description: 'Stores the input and references in the user Durable Object. Provider data is read from existing shared storage. This does not fetch YouTube data. An optional projectId saves the project reference in the same user-storage transaction.',
requestBody: jsonBody(z.toJSONSchema(saveSourceSchema, { target: 'openapi-3.0' })),
responses: { '201': jsonResponse('Source remembered.', { type: 'object', properties: { source: schemaRef('RecentSource') } }),
'409': jsonResponse('Provider data is not yet present in shared storage.', schemaRef('Error')), ...standardErrors },
responses: { '201': jsonResponse('Source remembered.', { type: 'object', properties: { source: schemaRef('RecentSource'), linked: { anyOf: [schemaRef('ProjectSourceLink'), { type: 'null' }] } } }),
'409': jsonResponse('Provider data is not yet present in shared storage.', schemaRef('Error')), '404': responseRef('NotFound'), ...standardErrors },
},
},
'/v1/sources/recent/{id}': {
Expand Down Expand Up @@ -1181,6 +1181,35 @@ export const openApiDocument = {
},
},
},
'/v1/projects/{id}/sources/{itemId}': {
get: {
tags: ['Projects'], operationId: 'getProjectSource', summary: 'Restore a project source snapshot', security: privateSecurity,
description: 'Restores this project item without updating history or project references. Survives recent-source eviction.',
parameters: [idParameter, pathParameter('itemId', 'Project item UUID, not the recent source UUID.')],
responses: { '200': jsonResponse('Saved project source and displayed data.', { type: 'object', properties: {
source: schemaRef('RecentSource'), snapshot: z.toJSONSchema(sourceSnapshotSchema, { target: 'openapi-3.0' }),
} }), '404': responseRef('NotFound'), ...standardErrors },
},
},
'/v1/projects/{id}/sources': {
post: {
tags: ['Projects'],
operationId: 'linkProjectSource',
summary: 'Add a saved Sources search or inspection to a project',
description: 'Keeps a user-owned reference to the shared source assets even after the recent Sources list rotates.',
security: privateSecurity,
parameters: [idParameter],
requestBody: jsonBody({ type: 'object', required: ['sourceId'], properties: {
sourceId: { type: 'string', format: 'uuid' },
} }),
responses: {
'201': jsonResponse('Source linked to the project.', schemaRef('ProjectSourceLink')),
'200': jsonResponse('Existing project source refreshed.', schemaRef('ProjectSourceLink')),
...standardErrors,
'404': responseRef('NotFound'),
},
},
},
'/v1/imports': {
post: {
tags: ['Research'],
Expand Down Expand Up @@ -2142,6 +2171,9 @@ export const openApiDocument = {
Project: { allOf: [storedRecord, { properties: { id: { type: 'string' }, name: { type: 'string' }, description: { type: 'string' }, item_count: { type: 'integer' } } }] },
ProjectDetail: { allOf: [schemaRef('Project'), { type: 'object', required: ['items'], properties: { items: { type: 'array', items: schemaRef('ProjectItem') } } }] },
ProjectItem: storedRecord,
ProjectSourceLink: { type: 'object', required: ['item', 'added'], properties: {
item: schemaRef('ProjectItem'), added: { type: 'boolean' },
} },
CreateProjectRequest: {
type: 'object', required: ['name'], properties: {
name: { type: 'string', minLength: 1, maxLength: 120, example: 'Research inbox' },
Expand Down
Loading
Loading