From 6145abd12dc2db48f977dd5dee98a26a14043478 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:19:36 +0800 Subject: [PATCH 01/16] Archive v0.7: add catalog observer persistence --- .../migrations/0009_catalog_observer.sql | 72 +++++++++++++++++++ 1 file changed, 72 insertions(+) create mode 100644 packages/database/migrations/0009_catalog_observer.sql diff --git a/packages/database/migrations/0009_catalog_observer.sql b/packages/database/migrations/0009_catalog_observer.sql new file mode 100644 index 0000000..9570b69 --- /dev/null +++ b/packages/database/migrations/0009_catalog_observer.sql @@ -0,0 +1,72 @@ +BEGIN; + +SET search_path TO modelapse, public; + +CREATE TABLE catalog_observer_sources ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + provider_id uuid NOT NULL REFERENCES providers(id), + source_key text NOT NULL CHECK (source_key ~ '^[a-z0-9][a-z0-9-]*$'), + source_kind text NOT NULL CHECK (source_kind IN ('model_list', 'docs')), + url text NOT NULL CHECK (url ~ '^https://'), + title text NOT NULL, + parser text NOT NULL CHECK (parser IN ('openai_models', 'snapshot_only')), + credential_env text, + interval_seconds integer NOT NULL DEFAULT 21600 + CHECK (interval_seconds BETWEEN 300 AND 604800), + enabled boolean NOT NULL DEFAULT true, + next_run_at timestamptz NOT NULL DEFAULT now(), + last_attempted_at timestamptz, + last_succeeded_at timestamptz, + metadata jsonb NOT NULL DEFAULT '{}'::jsonb, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + UNIQUE (provider_id, source_key) +); + +CREATE TABLE catalog_collection_runs ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + observer_source_id uuid NOT NULL REFERENCES catalog_observer_sources(id), + status text NOT NULL + CHECK (status IN ('running', 'succeeded', 'partial', 'failed', 'skipped')), + started_at timestamptz NOT NULL, + completed_at timestamptz, + http_status integer, + item_count integer CHECK (item_count IS NULL OR item_count >= 0), + observations_emitted integer NOT NULL DEFAULT 0 + CHECK (observations_emitted >= 0), + collector_build text NOT NULL, + error_message text, + metadata jsonb NOT NULL DEFAULT '{}'::jsonb, + created_at timestamptz NOT NULL DEFAULT now(), + CHECK (completed_at IS NULL OR completed_at >= started_at) +); + +CREATE TABLE catalog_source_snapshots ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + observer_source_id uuid NOT NULL REFERENCES catalog_observer_sources(id), + collection_run_id uuid NOT NULL UNIQUE REFERENCES catalog_collection_runs(id), + source_record_id uuid NOT NULL REFERENCES source_records(id), + retrieved_at timestamptz NOT NULL, + content_sha256 char(64) NOT NULL + CHECK (content_sha256 ~ '^[0-9a-f]{64}$'), + content_type text, + etag text, + last_modified text, + response_body text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now() +); + +CREATE TRIGGER catalog_source_snapshots_append_only +BEFORE UPDATE OR DELETE ON catalog_source_snapshots +FOR EACH ROW EXECUTE FUNCTION prevent_append_only_mutation(); + +CREATE INDEX catalog_observer_sources_due_idx + ON catalog_observer_sources (enabled, next_run_at, provider_id); + +CREATE INDEX catalog_collection_runs_source_started_idx + ON catalog_collection_runs (observer_source_id, started_at DESC); + +CREATE INDEX catalog_source_snapshots_source_retrieved_idx + ON catalog_source_snapshots (observer_source_id, retrieved_at DESC); + +COMMIT; From 75ec3c425ba306fcddca6355ae4939a25194b8d2 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:21:05 +0800 Subject: [PATCH 02/16] Archive v0.7: implement scheduled catalog collector core --- .../catalog-admin/src/catalog-observer.ts | 780 ++++++++++++++++++ 1 file changed, 780 insertions(+) create mode 100644 packages/catalog-admin/src/catalog-observer.ts diff --git a/packages/catalog-admin/src/catalog-observer.ts b/packages/catalog-admin/src/catalog-observer.ts new file mode 100644 index 0000000..47d72f3 --- /dev/null +++ b/packages/catalog-admin/src/catalog-observer.ts @@ -0,0 +1,780 @@ +import { createHash } from "node:crypto"; +import { Pool, type PoolClient } from "pg"; +import { PgModelCatalogAdmin } from "./model-catalog.js"; + +export type CatalogObserverSourceKind = "model_list" | "docs"; +export type CatalogObserverParser = "openai_models" | "snapshot_only"; +export type CatalogCollectionStatus = + | "succeeded" + | "partial" + | "failed" + | "skipped"; + +export interface ObservedRemoteModel { + readonly id: string; + readonly providerSnapshotId: string | null; +} + +export interface CatalogCollectionResult { + readonly runId: string; + readonly sourceKey: string; + readonly providerSlug: string; + readonly status: CatalogCollectionStatus; + readonly httpStatus: number | null; + readonly contentSha256: string | null; + readonly itemCount: number | null; + readonly observationsEmitted: number; + readonly error: string | null; +} + +interface ObserverSource { + readonly id: string; + readonly providerId: string; + readonly providerSlug: string; + readonly sourceKey: string; + readonly sourceKind: CatalogObserverSourceKind; + readonly url: string; + readonly title: string; + readonly parser: CatalogObserverParser; + readonly credentialEnv: string | null; + readonly intervalSeconds: number; +} + +interface KnownBinding { + readonly canonicalSlug: string; + readonly apiModelId: string; + readonly providerSnapshotId: string | null; +} + +type FetchLike = typeof fetch; + +const DEFAULT_SOURCES = [ + { + providerSlug: "openai", + sourceKey: "models-api", + sourceKind: "model_list", + url: "https://api.openai.com/v1/models", + title: "OpenAI Models API", + parser: "openai_models", + credentialEnv: "OPENAI_API_KEY", + intervalSeconds: 21600, + }, + { + providerSlug: "openai", + sourceKey: "responses-docs", + sourceKind: "docs", + url: "https://platform.openai.com/docs/api-reference/responses", + title: "OpenAI Responses API reference", + parser: "snapshot_only", + credentialEnv: null, + intervalSeconds: 86400, + }, + { + providerSlug: "deepseek", + sourceKey: "models-api", + sourceKind: "model_list", + url: "https://api.deepseek.com/models", + title: "DeepSeek Models API", + parser: "openai_models", + credentialEnv: "DEEPSEEK_API_KEY", + intervalSeconds: 21600, + }, + { + providerSlug: "deepseek", + sourceKey: "responses-docs", + sourceKind: "docs", + url: "https://api-docs.deepseek.com/guides/responses_api/", + title: "DeepSeek Responses API guide", + parser: "snapshot_only", + credentialEnv: null, + intervalSeconds: 86400, + }, +] as const; + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} + +function optionalSnapshotId(value: Record): string | null { + for (const key of [ + "provider_snapshot_id", + "providerSnapshotId", + "snapshot", + "version", + "model_version", + ]) { + const raw = value[key]; + if (typeof raw === "string" && raw.trim()) return raw.trim(); + } + return null; +} + +export function parseOpenAICompatibleModelList( + input: unknown, +): readonly ObservedRemoteModel[] { + if (!isRecord(input) || !Array.isArray(input.data)) { + throw new Error("Model-list payload must contain a data array"); + } + + const models = new Map(); + for (const item of input.data) { + if (!isRecord(item) || typeof item.id !== "string" || !item.id.trim()) { + continue; + } + const id = item.id.trim(); + models.set(id, { + id, + providerSnapshotId: optionalSnapshotId(item), + }); + } + return [...models.values()].sort((a, b) => a.id.localeCompare(b.id)); +} + +function normalizedTimestamp(value: string | undefined): string { + if (!value) return new Date().toISOString(); + const parsed = new Date(value); + if (Number.isNaN(parsed.valueOf())) { + throw new Error("now must be an ISO-8601 timestamp"); + } + return parsed.toISOString(); +} + +function positiveInteger( + value: number | undefined, + fallback: number, + name: string, +): number { + const resolved = value ?? fallback; + if (!Number.isInteger(resolved) || resolved <= 0) { + throw new Error(name + " must be a positive integer"); + } + return resolved; +} + +function safeError(error: unknown): string { + if (error instanceof Error) return error.message.slice(0, 4000); + return String(error).slice(0, 4000); +} + +async function sourceRows( + client: PoolClient, + source: ObserverSource, + now: string, + limit: number, + force: boolean, +): Promise { + const result = await client.query<{ + id: string; + provider_id: string; + provider_slug: string; + source_key: string; + source_kind: CatalogObserverSourceKind; + url: string; + title: string; + parser: CatalogObserverParser; + credential_env: string | null; + interval_seconds: number; + }>( + `WITH due AS ( + SELECT candidate.id + FROM modelapse.catalog_observer_sources candidate + JOIN modelapse.providers provider + ON provider.id = candidate.provider_id + WHERE candidate.enabled = true + AND ($1::text IS NULL OR provider.slug = $1) + AND ($2::text IS NULL OR candidate.source_key = $2) + AND ($3::boolean OR candidate.next_run_at <= $4::timestamptz) + ORDER BY candidate.next_run_at, candidate.id + LIMIT $5 + FOR UPDATE OF candidate SKIP LOCKED + ), + claimed AS ( + UPDATE modelapse.catalog_observer_sources claimed_source + SET last_attempted_at = $4::timestamptz, + next_run_at = + $4::timestamptz + + make_interval(secs => claimed_source.interval_seconds), + updated_at = now() + FROM due + WHERE claimed_source.id = due.id + RETURNING claimed_source.* + ) + SELECT + claimed.id, + claimed.provider_id, + provider.slug AS provider_slug, + claimed.source_key, + claimed.source_kind, + claimed.url, + claimed.title, + claimed.parser, + claimed.credential_env, + claimed.interval_seconds + FROM claimed + JOIN modelapse.providers provider + ON provider.id = claimed.provider_id + ORDER BY claimed.next_run_at, claimed.id`, + [source.providerSlug, source.sourceKey, force, now, limit], + ); + + return result.rows.map((row) => ({ + id: row.id, + providerId: row.provider_id, + providerSlug: row.provider_slug, + sourceKey: row.source_key, + sourceKind: row.source_kind, + url: row.url, + title: row.title, + parser: row.parser, + credentialEnv: row.credential_env, + intervalSeconds: row.interval_seconds, + })); +} + +export class PgCatalogObserver { + private readonly fetchImpl: FetchLike; + private readonly credentialResolver: (name: string) => string | undefined; + private readonly timeoutMs: number; + private readonly maxResponseBytes: number; + private readonly modelAdmin: PgModelCatalogAdmin; + + constructor( + private readonly pool: Pool, + options: { + readonly fetchImpl?: FetchLike; + readonly credentialResolver?: (name: string) => string | undefined; + readonly timeoutMs?: number; + readonly maxResponseBytes?: number; + } = {}, + ) { + this.fetchImpl = options.fetchImpl ?? fetch; + this.credentialResolver = + options.credentialResolver ?? ((name) => process.env[name]); + this.timeoutMs = positiveInteger( + options.timeoutMs, + 30_000, + "Catalog observer timeout", + ); + this.maxResponseBytes = positiveInteger( + options.maxResponseBytes, + 2_000_000, + "Catalog observer maxResponseBytes", + ); + this.modelAdmin = new PgModelCatalogAdmin(pool); + } + + static connect( + connectionString: string, + options: { + readonly max?: number; + readonly fetchImpl?: FetchLike; + readonly credentialResolver?: (name: string) => string | undefined; + readonly timeoutMs?: number; + readonly maxResponseBytes?: number; + } = {}, + ): PgCatalogObserver { + const pool = new Pool({ + connectionString, + max: options.max ?? 3, + }); + return new PgCatalogObserver(pool, { + ...(options.fetchImpl ? { fetchImpl: options.fetchImpl } : {}), + ...(options.credentialResolver + ? { credentialResolver: options.credentialResolver } + : {}), + ...(options.timeoutMs ? { timeoutMs: options.timeoutMs } : {}), + ...(options.maxResponseBytes + ? { maxResponseBytes: options.maxResponseBytes } + : {}), + }); + } + + async close(): Promise { + await this.pool.end(); + } + + async ensureDefaultSources(): Promise { + for (const source of DEFAULT_SOURCES) { + await this.pool.query( + `INSERT INTO modelapse.catalog_observer_sources + ( + provider_id, + source_key, + source_kind, + url, + title, + parser, + credential_env, + interval_seconds + ) + SELECT + provider.id, + $2, + $3, + $4, + $5, + $6, + $7, + $8 + FROM modelapse.providers provider + WHERE provider.slug = $1 + ON CONFLICT (provider_id, source_key) DO NOTHING`, + [ + source.providerSlug, + source.sourceKey, + source.sourceKind, + source.url, + source.title, + source.parser, + source.credentialEnv, + source.intervalSeconds, + ], + ); + } + } + + async collectDue(input: { + readonly collectorBuild: string; + readonly now?: string; + readonly limit?: number; + readonly force?: boolean; + readonly providerSlug?: string; + readonly sourceKey?: string; + }): Promise { + const collectorBuild = input.collectorBuild.trim(); + if (!collectorBuild) throw new Error("collectorBuild is required"); + const now = normalizedTimestamp(input.now); + const limit = input.limit ?? 4; + if (!Number.isInteger(limit) || limit < 1 || limit > 20) { + throw new Error("Catalog observer limit must be an integer between 1 and 20"); + } + + await this.ensureDefaultSources(); + + const client = await this.pool.connect(); + let claimed: readonly ObserverSource[] = []; + try { + await client.query("BEGIN"); + claimed = await sourceRows( + client, + { + id: "", + providerId: "", + providerSlug: input.providerSlug?.trim() || "", + sourceKey: input.sourceKey?.trim() || "", + sourceKind: "docs", + url: "", + title: "", + parser: "snapshot_only", + credentialEnv: null, + intervalSeconds: 0, + }, + now, + limit, + input.force ?? false, + ); + await client.query("COMMIT"); + } catch (error) { + await client.query("ROLLBACK"); + throw error; + } finally { + client.release(); + } + + const results: CatalogCollectionResult[] = []; + for (const source of claimed) { + results.push(await this.collectSource(source, now, collectorBuild)); + } + return results; + } + + private async startRun( + sourceId: string, + startedAt: string, + collectorBuild: string, + ): Promise { + const result = await this.pool.query<{ id: string }>( + `INSERT INTO modelapse.catalog_collection_runs + (observer_source_id, status, started_at, collector_build) + VALUES ($1, 'running', $2, $3) + RETURNING id`, + [sourceId, startedAt, collectorBuild], + ); + const id = result.rows[0]?.id; + if (!id) throw new Error("Catalog collection run insert failed"); + return id; + } + + private async finishRun( + runId: string, + input: { + readonly status: CatalogCollectionStatus; + readonly completedAt: string; + readonly httpStatus?: number; + readonly itemCount?: number; + readonly observationsEmitted?: number; + readonly error?: string; + readonly metadata?: Record; + }, + ): Promise { + await this.pool.query( + `UPDATE modelapse.catalog_collection_runs + SET status = $2, + completed_at = $3, + http_status = $4, + item_count = $5, + observations_emitted = $6, + error_message = $7, + metadata = $8::jsonb + WHERE id = $1 + AND status = 'running'`, + [ + runId, + input.status, + input.completedAt, + input.httpStatus ?? null, + input.itemCount ?? null, + input.observationsEmitted ?? 0, + input.error ?? null, + JSON.stringify(input.metadata ?? {}), + ], + ); + } + + private async storeSnapshot(input: { + readonly source: ObserverSource; + readonly runId: string; + readonly retrievedAt: string; + readonly collectorBuild: string; + readonly body: string; + readonly contentSha256: string; + readonly contentType: string | null; + readonly etag: string | null; + readonly lastModified: string | null; + readonly httpStatus: number; + }): Promise { + const client = await this.pool.connect(); + try { + await client.query("BEGIN"); + const sourceType = + input.source.sourceKind === "model_list" + ? "provider_catalog" + : "provider_docs"; + const sourceRecord = await client.query<{ id: string }>( + `INSERT INTO modelapse.source_records + ( + source_type, + url, + title, + retrieved_at, + content_sha256, + metadata + ) + VALUES ($1, $2, $3, $4, $5, $6::jsonb) + RETURNING id`, + [ + sourceType, + input.source.url, + input.source.title, + input.retrievedAt, + input.contentSha256, + JSON.stringify({ + collector: "catalog-observer", + collectorBuild: input.collectorBuild, + observerSourceId: input.source.id, + collectionRunId: input.runId, + sourceKey: input.source.sourceKey, + parser: input.source.parser, + httpStatus: input.httpStatus, + contentType: input.contentType, + etag: input.etag, + lastModified: input.lastModified, + }), + ], + ); + const sourceRecordId = sourceRecord.rows[0]?.id; + if (!sourceRecordId) { + throw new Error("Catalog snapshot source record insert failed"); + } + + await client.query( + `INSERT INTO modelapse.catalog_source_snapshots + ( + observer_source_id, + collection_run_id, + source_record_id, + retrieved_at, + content_sha256, + content_type, + etag, + last_modified, + response_body + ) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)`, + [ + input.source.id, + input.runId, + sourceRecordId, + input.retrievedAt, + input.contentSha256, + input.contentType, + input.etag, + input.lastModified, + input.body, + ], + ); + await client.query("COMMIT"); + return sourceRecordId; + } catch (error) { + await client.query("ROLLBACK"); + throw error; + } finally { + client.release(); + } + } + + private async currentBindings( + providerId: string, + ): Promise { + const result = await this.pool.query<{ + canonical_slug: string; + api_model_id: string; + provider_snapshot_id: string | null; + }>( + `SELECT + model.canonical_slug, + binding.api_model_id, + snapshot.provider_snapshot_id + FROM modelapse.models model + JOIN modelapse.model_execution_bindings binding + ON binding.model_id = model.id + AND binding.valid_to IS NULL + JOIN modelapse.provider_endpoints endpoint + ON endpoint.id = binding.endpoint_id + AND endpoint.path = 'first_party_direct' + LEFT JOIN modelapse.model_snapshots snapshot + ON snapshot.id = binding.snapshot_id + WHERE model.provider_id = $1 + ORDER BY model.canonical_slug`, + [providerId], + ); + return result.rows.map((row) => ({ + canonicalSlug: row.canonical_slug, + apiModelId: row.api_model_id, + providerSnapshotId: row.provider_snapshot_id, + })); + } + + private async collectSource( + source: ObserverSource, + observedAt: string, + collectorBuild: string, + ): Promise { + const runId = await this.startRun(source.id, observedAt, collectorBuild); + const credential = source.credentialEnv + ? this.credentialResolver(source.credentialEnv)?.trim() + : undefined; + + if (source.credentialEnv && !credential) { + const error = "Missing collector credential: " + source.credentialEnv; + await this.finishRun(runId, { + status: "skipped", + completedAt: new Date().toISOString(), + error, + }); + return { + runId, + sourceKey: source.sourceKey, + providerSlug: source.providerSlug, + status: "skipped", + httpStatus: null, + contentSha256: null, + itemCount: null, + observationsEmitted: 0, + error, + }; + } + + let httpStatus: number | null = null; + try { + const headers: Record = { + accept: + source.sourceKind === "model_list" + ? "application/json" + : "text/html,application/xhtml+xml,text/plain;q=0.8,*/*;q=0.5", + "user-agent": "modelapse-catalog-observer/" + collectorBuild.slice(0, 64), + }; + if (credential) headers.authorization = "Bearer " + credential; + + const response = await this.fetchImpl(source.url, { + method: "GET", + headers, + redirect: "follow", + signal: AbortSignal.timeout(this.timeoutMs), + }); + httpStatus = response.status; + if (!response.ok) { + throw new Error("Catalog source returned HTTP " + response.status); + } + + const bytes = Buffer.from(await response.arrayBuffer()); + if (bytes.byteLength > this.maxResponseBytes) { + throw new Error( + "Catalog source exceeded max response bytes: " + bytes.byteLength, + ); + } + const body = bytes.toString("utf8"); + const contentSha256 = createHash("sha256").update(bytes).digest("hex"); + const contentType = response.headers.get("content-type"); + const etag = response.headers.get("etag"); + const lastModified = response.headers.get("last-modified"); + + const sourceRecordId = await this.storeSnapshot({ + source, + runId, + retrievedAt: observedAt, + collectorBuild, + body, + contentSha256, + contentType, + etag, + lastModified, + httpStatus: response.status, + }); + + if (source.parser === "snapshot_only") { + await this.pool.query( + `UPDATE modelapse.catalog_observer_sources + SET last_succeeded_at = $2, + updated_at = now() + WHERE id = $1`, + [source.id, observedAt], + ); + await this.finishRun(runId, { + status: "succeeded", + completedAt: new Date().toISOString(), + httpStatus: response.status, + observationsEmitted: 0, + metadata: { sourceRecordId, contentSha256 }, + }); + return { + runId, + sourceKey: source.sourceKey, + providerSlug: source.providerSlug, + status: "succeeded", + httpStatus: response.status, + contentSha256, + itemCount: null, + observationsEmitted: 0, + error: null, + }; + } + + const remoteModels = parseOpenAICompatibleModelList(JSON.parse(body)); + const remoteById = new Map(remoteModels.map((model) => [model.id, model])); + const knownBindings = await this.currentBindings(source.providerId); + const missingKnownApiModelIds: string[] = []; + const matchedApiModelIds = new Set(); + const observationErrors: string[] = []; + let observationsEmitted = 0; + + for (const binding of knownBindings) { + const remote = remoteById.get(binding.apiModelId); + if (!remote) { + missingKnownApiModelIds.push(binding.apiModelId); + continue; + } + matchedApiModelIds.add(binding.apiModelId); + try { + await this.modelAdmin.observeFirstPartyIdentity({ + providerSlug: source.providerSlug, + canonicalSlug: binding.canonicalSlug, + apiModelId: binding.apiModelId, + ...(remote.providerSnapshotId ?? binding.providerSnapshotId + ? { + providerSnapshotId: + remote.providerSnapshotId ?? binding.providerSnapshotId ?? undefined, + } + : {}), + sourceUrl: source.url, + sourceTitle: source.title, + sourceType: "provider_catalog", + sourceRecordId, + contentSha256, + observedAt, + collector: "catalog-observer", + }); + observationsEmitted += 1; + } catch (error) { + observationErrors.push( + binding.canonicalSlug + ": " + safeError(error), + ); + } + } + + const unmatchedRemoteModelIds = remoteModels + .map((model) => model.id) + .filter((id) => !matchedApiModelIds.has(id)) + .slice(0, 200); + const status: CatalogCollectionStatus = + observationErrors.length > 0 ? "partial" : "succeeded"; + const error = + observationErrors.length > 0 + ? observationErrors.slice(0, 20).join("; ") + : undefined; + + await this.pool.query( + `UPDATE modelapse.catalog_observer_sources + SET last_succeeded_at = $2, + updated_at = now() + WHERE id = $1`, + [source.id, observedAt], + ); + await this.finishRun(runId, { + status, + completedAt: new Date().toISOString(), + httpStatus: response.status, + itemCount: remoteModels.length, + observationsEmitted, + ...(error ? { error } : {}), + metadata: { + sourceRecordId, + contentSha256, + missingKnownApiModelIds, + unmatchedRemoteModelIds, + }, + }); + + return { + runId, + sourceKey: source.sourceKey, + providerSlug: source.providerSlug, + status, + httpStatus: response.status, + contentSha256, + itemCount: remoteModels.length, + observationsEmitted, + error: error ?? null, + }; + } catch (error) { + const message = safeError(error); + await this.finishRun(runId, { + status: "failed", + completedAt: new Date().toISOString(), + ...(httpStatus === null ? {} : { httpStatus }), + error: message, + }); + return { + runId, + sourceKey: source.sourceKey, + providerSlug: source.providerSlug, + status: "failed", + httpStatus, + contentSha256: null, + itemCount: null, + observationsEmitted: 0, + error: message, + }; + } + } +} From fcbbf3feaa4577289ddba22c29736d339618e7bd Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:21:31 +0800 Subject: [PATCH 03/16] Archive v0.7: let collector reuse fetched source records --- packages/catalog-admin/src/model-catalog.ts | 65 +++++++++++++++++---- 1 file changed, 55 insertions(+), 10 deletions(-) diff --git a/packages/catalog-admin/src/model-catalog.ts b/packages/catalog-admin/src/model-catalog.ts index 6e07ae8..51ea817 100644 --- a/packages/catalog-admin/src/model-catalog.ts +++ b/packages/catalog-admin/src/model-catalog.ts @@ -340,8 +340,11 @@ export class PgModelCatalogAdmin { readonly providerSnapshotId?: string; readonly sourceUrl: string; readonly sourceTitle: string; + readonly sourceType?: string; + readonly sourceRecordId?: string; readonly contentSha256?: string; readonly observedAt?: string; + readonly collector?: string; }): Promise { const providerSlug = input.providerSlug.trim(); const canonicalSlug = input.canonicalSlug.trim(); @@ -349,10 +352,21 @@ export class PgModelCatalogAdmin { const providerSnapshotId = input.providerSnapshotId?.trim() || null; const sourceUrl = input.sourceUrl.trim(); const sourceTitle = input.sourceTitle.trim(); + const sourceType = input.sourceType?.trim() || "provider_docs"; + const sourceRecordId = input.sourceRecordId?.trim() || null; const contentSha256 = input.contentSha256?.trim(); const observedAt = normalizedObservationTime(input.observedAt); + const collector = input.collector?.trim() || "catalog-admin-identity-observer"; - if (!providerSlug || !canonicalSlug || !apiModelId || !sourceUrl || !sourceTitle) { + if ( + !providerSlug || + !canonicalSlug || + !apiModelId || + !sourceUrl || + !sourceTitle || + !sourceType || + !collector + ) { throw new Error("Identity observation fields must be non-empty"); } if ( @@ -408,13 +422,43 @@ export class PgModelCatalogAdmin { ); } - const sourceId = await recordSourceObservation(client, { - sourceType: "provider_docs", - url: sourceUrl, - title: sourceTitle, - retrievedAt: observedAt, - ...(contentSha256 ? { contentSha256 } : {}), - }); + let sourceId = sourceRecordId; + if (sourceId) { + const existingSource = await client.query<{ + source_type: string; + url: string | null; + title: string | null; + retrieved_at: Date; + content_sha256: string | null; + }>( + `SELECT source_type, url, title, retrieved_at, content_sha256 + FROM modelapse.source_records + WHERE id = $1`, + [sourceId], + ); + const source = existingSource.rows[0]; + if ( + !source || + source.source_type !== sourceType || + source.url !== sourceUrl || + source.title !== sourceTitle || + source.retrieved_at.getTime() !== Date.parse(observedAt) || + (contentSha256 !== undefined && + source.content_sha256 !== contentSha256) + ) { + throw new Error( + "Existing source record does not match the identity observation", + ); + } + } else { + sourceId = await recordSourceObservation(client, { + sourceType, + url: sourceUrl, + title: sourceTitle, + retrievedAt: observedAt, + ...(contentSha256 ? { contentSha256 } : {}), + }); + } let snapshotId: string | null = null; if (providerSnapshotId) { @@ -567,19 +611,20 @@ export class PgModelCatalogAdmin { confidence, raw_observation ) - VALUES ($1, $2, $3, $4, 'provider_docs', $5, 1.0, $6::jsonb) + VALUES ($1, $2, $3, $4, $5, $6, 1.0, $7::jsonb) RETURNING id`, [ aliasId, modelId, snapshotId, observedAt, + sourceType, sourceId, JSON.stringify({ apiModelId, endpointId, providerSnapshotId, - collector: "catalog-admin-identity-observer", + collector, }), ], ); From b92de919c4c01356608196ab2e3c58cf799d6f70 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:21:55 +0800 Subject: [PATCH 04/16] Archive v0.7: fix observer filters and snapshot typing --- .../catalog-admin/src/catalog-observer.ts | 28 ++++++++----------- 1 file changed, 11 insertions(+), 17 deletions(-) diff --git a/packages/catalog-admin/src/catalog-observer.ts b/packages/catalog-admin/src/catalog-observer.ts index 47d72f3..10303ad 100644 --- a/packages/catalog-admin/src/catalog-observer.ts +++ b/packages/catalog-admin/src/catalog-observer.ts @@ -158,7 +158,10 @@ function safeError(error: unknown): string { async function sourceRows( client: PoolClient, - source: ObserverSource, + input: { + readonly providerSlug: string | null; + readonly sourceKey: string | null; + }, now: string, limit: number, force: boolean, @@ -214,7 +217,7 @@ async function sourceRows( JOIN modelapse.providers provider ON provider.id = claimed.provider_id ORDER BY claimed.next_run_at, claimed.id`, - [source.providerSlug, source.sourceKey, force, now, limit], + [input.providerSlug, input.sourceKey, force, now, limit], ); return result.rows.map((row) => ({ @@ -358,16 +361,8 @@ export class PgCatalogObserver { claimed = await sourceRows( client, { - id: "", - providerId: "", - providerSlug: input.providerSlug?.trim() || "", - sourceKey: input.sourceKey?.trim() || "", - sourceKind: "docs", - url: "", - title: "", - parser: "snapshot_only", - credentialEnv: null, - intervalSeconds: 0, + providerSlug: input.providerSlug?.trim() || null, + sourceKey: input.sourceKey?.trim() || null, }, now, limit, @@ -686,15 +681,14 @@ export class PgCatalogObserver { } matchedApiModelIds.add(binding.apiModelId); try { + const observedSnapshotId = + remote.providerSnapshotId ?? binding.providerSnapshotId; await this.modelAdmin.observeFirstPartyIdentity({ providerSlug: source.providerSlug, canonicalSlug: binding.canonicalSlug, apiModelId: binding.apiModelId, - ...(remote.providerSnapshotId ?? binding.providerSnapshotId - ? { - providerSnapshotId: - remote.providerSnapshotId ?? binding.providerSnapshotId ?? undefined, - } + ...(observedSnapshotId + ? { providerSnapshotId: observedSnapshotId } : {}), sourceUrl: source.url, sourceTitle: source.title, From a725250bf6a0cad336e970f907db6e4e539b7f80 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:22:26 +0800 Subject: [PATCH 05/16] Archive v0.7: export catalog observer --- packages/catalog-admin/src/index.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/packages/catalog-admin/src/index.ts b/packages/catalog-admin/src/index.ts index 1a5c0c6..b033c6b 100644 --- a/packages/catalog-admin/src/index.ts +++ b/packages/catalog-admin/src/index.ts @@ -2,3 +2,4 @@ export * from "./catalog.js"; export * from "./openai-smoke.js"; export * from "./direct-smoke.js"; export * from "./model-catalog.js"; +export * from "./catalog-observer.js"; From df148699727f421fe0e720659e632f1ca7b86755 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:22:29 +0800 Subject: [PATCH 06/16] Archive v0.7: add collector CLI and unit coverage --- packages/catalog-admin/package.json | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/packages/catalog-admin/package.json b/packages/catalog-admin/package.json index 7708ce3..e107c4c 100644 --- a/packages/catalog-admin/package.json +++ b/packages/catalog-admin/package.json @@ -10,12 +10,13 @@ "scripts": { "build": "tsc -p tsconfig.json", "typecheck": "tsc -p tsconfig.json --noEmit", - "test": "vitest run test/openai-smoke.test.ts", + "test": "vitest run test/openai-smoke.test.ts test/catalog-observer.test.ts", "test:integration": "vitest run test/bootstrap.integration.test.ts", "bootstrap:openai-smoke": "node dist/src/cli.js bootstrap-openai-smoke", "bootstrap:deepseek-smoke": "node dist/src/cli.js bootstrap-deepseek-smoke", "bootstrap:deepseek-flash-model": "node dist/src/cli.js bootstrap-deepseek-flash-model", - "observe:first-party-identity": "node dist/src/cli.js observe-first-party-identity" + "observe:first-party-identity": "node dist/src/cli.js observe-first-party-identity", + "collect:first-party-catalog": "node dist/src/cli.js collect-first-party-catalog" }, "dependencies": { "@modelapse/blob-store": "*", From 607c9430095be8c0ac85a6044ede8153936e657a Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:22:31 +0800 Subject: [PATCH 07/16] Archive v0.7: expose one-shot catalog collection --- packages/catalog-admin/src/cli.ts | 23 +++++++++++++++++++++-- 1 file changed, 21 insertions(+), 2 deletions(-) diff --git a/packages/catalog-admin/src/cli.ts b/packages/catalog-admin/src/cli.ts index 9813ff3..13516ef 100644 --- a/packages/catalog-admin/src/cli.ts +++ b/packages/catalog-admin/src/cli.ts @@ -1,5 +1,6 @@ import { FileSystemContentAddressedBlobStore } from "@modelapse/blob-store"; import { PgCatalogAdmin } from "./catalog.js"; +import { PgCatalogObserver } from "./catalog-observer.js"; import { PgModelCatalogAdmin } from "./model-catalog.js"; function requiredEnv(name: string): string { @@ -16,11 +17,29 @@ if ( command !== "observe-first-party-identity" ) { throw new Error( - "Usage: catalog-admin bootstrap-openai-smoke|bootstrap-deepseek-smoke|bootstrap-deepseek-flash-model|observe-first-party-identity", + "Usage: catalog-admin bootstrap-openai-smoke|bootstrap-deepseek-smoke|bootstrap-deepseek-flash-model|observe-first-party-identity|collect-first-party-catalog", ); } -if (command === "observe-first-party-identity") { +if (command === "collect-first-party-catalog") { + const observer = PgCatalogObserver.connect(requiredEnv("DATABASE_URL")); + try { + const providerSlug = + process.env.MODELAPSE_PROVIDER_SLUG?.trim() || undefined; + const sourceKey = + process.env.MODELAPSE_CATALOG_SOURCE_KEY?.trim() || undefined; + const result = await observer.collectDue({ + collectorBuild: + process.env.MODELAPSE_BUILD?.trim() || "catalog-admin", + force: process.env.MODELAPSE_COLLECT_FORCE === "true", + ...(providerSlug ? { providerSlug } : {}), + ...(sourceKey ? { sourceKey } : {}), + }); + process.stdout.write(JSON.stringify(result, null, 2) + "\n"); + } finally { + await observer.close(); + } +} else if (command === "observe-first-party-identity") { const models = PgModelCatalogAdmin.connect(requiredEnv("DATABASE_URL")); try { const providerSnapshotId = From f40c9ea4d6a16f0087e76b770df5237120b4ba07 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:22:43 +0800 Subject: [PATCH 08/16] Archive v0.7: test model-list normalization --- .../test/catalog-observer.test.ts | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) create mode 100644 packages/catalog-admin/test/catalog-observer.test.ts diff --git a/packages/catalog-admin/test/catalog-observer.test.ts b/packages/catalog-admin/test/catalog-observer.test.ts new file mode 100644 index 0000000..10b834d --- /dev/null +++ b/packages/catalog-admin/test/catalog-observer.test.ts @@ -0,0 +1,27 @@ +import { describe, expect, it } from "vitest"; +import { parseOpenAICompatibleModelList } from "../src/catalog-observer.js"; + +describe("catalog observer model-list parser", () => { + it("normalizes and deduplicates OpenAI-compatible model lists", () => { + expect( + parseOpenAICompatibleModelList({ + object: "list", + data: [ + { id: " model-b ", version: "2026-10-01" }, + { id: "model-a", provider_snapshot_id: "snap-a" }, + { id: "model-b", version: "2026-10-02" }, + { nope: true }, + ], + }), + ).toEqual([ + { id: "model-a", providerSnapshotId: "snap-a" }, + { id: "model-b", providerSnapshotId: "2026-10-02" }, + ]); + }); + + it("rejects payloads without a data array", () => { + expect(() => parseOpenAICompatibleModelList({ models: [] })).toThrow( + /data array/, + ); + }); +}); From f77821319fca239b0555719f232ecfda93342fe8 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:23:18 +0800 Subject: [PATCH 09/16] Archive v0.7: make catalog observer a runner dependency --- apps/runner/package.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/apps/runner/package.json b/apps/runner/package.json index 978b8ef..7eed66f 100644 --- a/apps/runner/package.json +++ b/apps/runner/package.json @@ -19,13 +19,13 @@ "@modelapse/provider-adapter": "*", "@modelapse/provider-openai": "*", "@modelapse/provider-deepseek": "*", - "@modelapse/evaluation": "*" + "@modelapse/evaluation": "*", + "@modelapse/catalog-admin": "*" }, "devDependencies": { "@types/pg": "8.23.1", "pg": "8.23.0", "vitest": "3.0.8", - "@modelapse/catalog-admin": "*", "@modelapse/database": "*" } } From 463b84591d7692901ed7ac5645264558d95da828 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:23:21 +0800 Subject: [PATCH 10/16] Archive v0.7: schedule catalog collection in queue runner --- apps/runner/src/index.ts | 95 ++++++++++++++++++++++++++++++++++++++++ 1 file changed, 95 insertions(+) diff --git a/apps/runner/src/index.ts b/apps/runner/src/index.ts index 449b24a..5138e7a 100644 --- a/apps/runner/src/index.ts +++ b/apps/runner/src/index.ts @@ -3,6 +3,10 @@ import { readFile } from "node:fs/promises"; import { hostname } from "node:os"; import { setTimeout as sleep } from "node:timers/promises"; import { FileSystemContentAddressedBlobStore } from "@modelapse/blob-store"; +import { + PgCatalogObserver, + type CatalogCollectionResult, +} from "@modelapse/catalog-admin"; import { parseDirectProviderRunRequest, PgRunJobQueue, @@ -57,6 +61,33 @@ function positiveIntegerEnv(name: string, fallback: number): number { return value; } +function booleanEnv(name: string, fallback: boolean): boolean { + const raw = process.env[name]?.trim().toLowerCase(); + if (!raw) return fallback; + if (["1", "true", "yes", "on"].includes(raw)) return true; + if (["0", "false", "no", "off"].includes(raw)) return false; + throw new Error(name + " must be a boolean"); +} + +function reportCatalogCollection(result: CatalogCollectionResult): void { + const line = JSON.stringify({ + catalogObserver: true, + provider: result.providerSlug, + source: result.sourceKey, + status: result.status, + httpStatus: result.httpStatus, + itemCount: result.itemCount, + observationsEmitted: result.observationsEmitted, + contentSha256: result.contentSha256, + error: result.error, + }); + if (result.status === "failed" || result.status === "partial") { + process.stderr.write(line + "\n"); + } else { + process.stdout.write(line + "\n"); + } +} + function summary(result: Awaited>) { return { runId: result.run.id, @@ -78,6 +109,22 @@ const providerTimeoutMs = positiveIntegerEnv( "MODELAPSE_PROVIDER_TIMEOUT_MS", 120_000, ); +const catalogObserverEnabled = booleanEnv( + "MODELAPSE_CATALOG_OBSERVER_ENABLED", + true, +); +const catalogObserverPollMs = positiveIntegerEnv( + "MODELAPSE_CATALOG_OBSERVER_POLL_MS", + 60_000, +); +const catalogObserverTimeoutMs = positiveIntegerEnv( + "MODELAPSE_CATALOG_OBSERVER_TIMEOUT_MS", + 30_000, +); +const catalogObserverMaxResponseBytes = positiveIntegerEnv( + "MODELAPSE_CATALOG_OBSERVER_MAX_RESPONSE_BYTES", + 2_000_000, +); const privateKey = await privateKeyPem(); const repository = PgRunRepository.connect(databaseUrl, { max: 2 }); const evaluations = PgEvaluationRepository.connect(databaseUrl, { max: 2 }); @@ -108,6 +155,13 @@ async function runStdinMode(): Promise { async function runQueueMode(): Promise { const queue = PgRunJobQueue.connect(databaseUrl, { max: 2 }); + const catalogObserver = catalogObserverEnabled + ? PgCatalogObserver.connect(databaseUrl, { + max: 3, + timeoutMs: catalogObserverTimeoutMs, + maxResponseBytes: catalogObserverMaxResponseBytes, + }) + : null; const workerId = process.env.MODELAPSE_WORKER_ID ?? hostname() + "-" + process.pid + "-" + randomUUID().slice(0, 8); @@ -123,6 +177,41 @@ async function runQueueMode(): Promise { } let stopping = false; + let catalogObserverTask: Promise | null = null; + let nextCatalogObserverAt = 0; + + const maybeCollectCatalog = () => { + if ( + !catalogObserver || + catalogObserverTask || + Date.now() < nextCatalogObserverAt + ) { + return; + } + + nextCatalogObserverAt = Date.now() + catalogObserverPollMs; + catalogObserverTask = catalogObserver + .collectDue({ + collectorBuild: runnerBuild, + limit: 4, + }) + .then((results) => { + for (const result of results) reportCatalogCollection(result); + }) + .catch((error: unknown) => { + process.stderr.write( + JSON.stringify({ + catalogObserver: true, + status: "loop_failed", + error: error instanceof Error ? error.message : String(error), + }) + "\n", + ); + }) + .finally(() => { + catalogObserverTask = null; + }); + }; + const stop = () => { stopping = true; }; @@ -131,6 +220,8 @@ async function runQueueMode(): Promise { try { while (!stopping) { + maybeCollectCatalog(); + const job = await processOneQueuedRunJob({ queue, repository, @@ -163,7 +254,11 @@ async function runQueueMode(): Promise { await sleep(pollMs); } } finally { + if (catalogObserverTask) { + await catalogObserverTask; + } await queue.close(); + await catalogObserver?.close(); } } From bbd76065c12ea311d770b307564602b3c777704d Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:24:09 +0800 Subject: [PATCH 11/16] Archive v0.7: verify scheduled collection feeds drift engine --- .../test/bootstrap.integration.test.ts | 180 +++++++++++++++++- 1 file changed, 179 insertions(+), 1 deletion(-) diff --git a/packages/catalog-admin/test/bootstrap.integration.test.ts b/packages/catalog-admin/test/bootstrap.integration.test.ts index 83b2fe8..2fcd494 100644 --- a/packages/catalog-admin/test/bootstrap.integration.test.ts +++ b/packages/catalog-admin/test/bootstrap.integration.test.ts @@ -8,7 +8,11 @@ import { migrateDatabase } from "@modelapse/database"; import { PgRunRepository } from "@modelapse/persistence"; import { Pool } from "pg"; import { afterAll, beforeAll, describe, expect, it } from "vitest"; -import { PgCatalogAdmin, PgModelCatalogAdmin } from "../src/index.js"; +import { + PgCatalogAdmin, + PgCatalogObserver, + PgModelCatalogAdmin, +} from "../src/index.js"; const ADMIN_DATABASE_URL = process.env.DATABASE_URL ?? @@ -200,6 +204,180 @@ describe("production catalog bootstrap", () => { } }); + it("periodically snapshots a first-party model list and feeds identity drift", async () => { + await catalog!.bootstrapDeepSeekSmoke({ + runnerBuild: "deepseek-observer-prerequisite", + }); + const registered = await modelCatalog!.bootstrapDeepSeekFlash(); + + let providerSnapshotId = "deepseek-flash-observer-a"; + const observer = PgCatalogObserver.connect(isolatedDatabaseUrl, { + fetchImpl: async () => + new Response( + JSON.stringify({ + object: "list", + data: [ + { + id: "deepseek-flash", + provider_snapshot_id: providerSnapshotId, + }, + { id: "unmapped-remote-model" }, + ], + }), + { + status: 200, + headers: { + "content-type": "application/json", + etag: '"observer-test"', + }, + }, + ), + credentialResolver: () => "observer-test-secret", + }); + + try { + const first = await observer.collectDue({ + collectorBuild: "observer-build-a", + providerSlug: "deepseek", + sourceKey: "models-api", + force: true, + now: "2099-02-01T00:00:00.000Z", + }); + providerSnapshotId = "deepseek-flash-observer-b"; + const second = await observer.collectDue({ + collectorBuild: "observer-build-b", + providerSlug: "deepseek", + sourceKey: "models-api", + force: true, + now: "2099-02-02T00:00:00.000Z", + }); + + expect(first).toEqual([ + expect.objectContaining({ + providerSlug: "deepseek", + sourceKey: "models-api", + status: "succeeded", + itemCount: 2, + observationsEmitted: 1, + }), + ]); + expect(second).toEqual([ + expect.objectContaining({ + providerSlug: "deepseek", + sourceKey: "models-api", + status: "succeeded", + itemCount: 2, + observationsEmitted: 1, + }), + ]); + + const verification = new Pool({ connectionString: isolatedDatabaseUrl }); + try { + const snapshots = await verification.query<{ + source_record_id: string; + content_sha256: string; + response_body: string; + }>( + `SELECT + snapshot.source_record_id, + snapshot.content_sha256, + snapshot.response_body + FROM modelapse.catalog_source_snapshots snapshot + JOIN modelapse.catalog_observer_sources source + ON source.id = snapshot.observer_source_id + JOIN modelapse.providers provider + ON provider.id = source.provider_id + WHERE provider.slug = 'deepseek' + AND source.source_key = 'models-api' + ORDER BY snapshot.retrieved_at`, + ); + + expect(snapshots.rows).toHaveLength(2); + expect( + new Set(snapshots.rows.map((row) => row.source_record_id)).size, + ).toBe(2); + expect( + snapshots.rows.map((row) => JSON.parse(row.response_body).data[0].provider_snapshot_id), + ).toEqual([ + "deepseek-flash-observer-a", + "deepseek-flash-observer-b", + ]); + + const observerDrift = await verification.query<{ + change_type: string; + changed_fields: string[]; + provider_snapshot_id: string | null; + }>( + `SELECT + drift.change_type, + drift.changed_fields, + current_snapshot.provider_snapshot_id + FROM modelapse.catalog_identity_drift_events drift + LEFT JOIN modelapse.model_snapshots current_snapshot + ON current_snapshot.id = drift.current_snapshot_id + WHERE drift.current_model_id = $1 + AND drift.current_source_id = ANY($2::uuid[]) + ORDER BY drift.occurred_at, drift.change_type`, + [ + registered.modelId, + snapshots.rows.map((row) => row.source_record_id), + ], + ); + + expect(observerDrift.rows).toHaveLength(4); + expect( + observerDrift.rows.filter( + (row) => row.change_type === "alias_target_changed", + ), + ).toHaveLength(2); + expect( + observerDrift.rows.filter( + (row) => row.change_type === "execution_binding_changed", + ), + ).toHaveLength(2); + expect( + observerDrift.rows.every((row) => + row.changed_fields.includes("snapshot"), + ), + ).toBe(true); + expect( + observerDrift.rows.some( + (row) => + row.provider_snapshot_id === "deepseek-flash-observer-b", + ), + ).toBe(true); + + const runMetadata = await verification.query<{ + status: string; + metadata: { + unmatchedRemoteModelIds?: string[]; + }; + }>( + `SELECT run.status, run.metadata + FROM modelapse.catalog_collection_runs run + JOIN modelapse.catalog_observer_sources source + ON source.id = run.observer_source_id + JOIN modelapse.providers provider + ON provider.id = source.provider_id + WHERE provider.slug = 'deepseek' + AND source.source_key = 'models-api' + ORDER BY run.started_at`, + ); + expect(runMetadata.rows).toHaveLength(2); + expect(runMetadata.rows[0]).toMatchObject({ + status: "succeeded", + metadata: { + unmatchedRemoteModelIds: ["unmapped-remote-model"], + }, + }); + } finally { + await verification.end(); + } + } finally { + await observer.close(); + } + }); + it("bootstraps DeepSeek idempotently and resolves its sourced direct target", async () => { const first = await catalog!.bootstrapDeepSeekSmoke({ runnerBuild: "deepseek-build-a", From bcb23332e76463179ee13e80538328d90e16957a Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:26:22 +0800 Subject: [PATCH 12/16] Archive v0.7: allow collector CLI command --- packages/catalog-admin/src/cli.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/packages/catalog-admin/src/cli.ts b/packages/catalog-admin/src/cli.ts index 13516ef..bc483ed 100644 --- a/packages/catalog-admin/src/cli.ts +++ b/packages/catalog-admin/src/cli.ts @@ -14,7 +14,8 @@ if ( command !== "bootstrap-openai-smoke" && command !== "bootstrap-deepseek-smoke" && command !== "bootstrap-deepseek-flash-model" && - command !== "observe-first-party-identity" + command !== "observe-first-party-identity" && + command !== "collect-first-party-catalog" ) { throw new Error( "Usage: catalog-admin bootstrap-openai-smoke|bootstrap-deepseek-smoke|bootstrap-deepseek-flash-model|observe-first-party-identity|collect-first-party-catalog", From 8feecbb894c48b478cca72456b19a9e3df23580e Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:26:29 +0800 Subject: [PATCH 13/16] Archive v0.7: keep collection completion timestamps valid --- packages/catalog-admin/src/catalog-observer.ts | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/packages/catalog-admin/src/catalog-observer.ts b/packages/catalog-admin/src/catalog-observer.ts index 10303ad..2007e98 100644 --- a/packages/catalog-admin/src/catalog-observer.ts +++ b/packages/catalog-admin/src/catalog-observer.ts @@ -156,6 +156,10 @@ function safeError(error: unknown): string { return String(error).slice(0, 4000); } +function completionTimestamp(startedAt: string): string { + return new Date(Math.max(Date.now(), Date.parse(startedAt))).toISOString(); +} + async function sourceRows( client: PoolClient, input: { @@ -574,7 +578,7 @@ export class PgCatalogObserver { const error = "Missing collector credential: " + source.credentialEnv; await this.finishRun(runId, { status: "skipped", - completedAt: new Date().toISOString(), + completedAt: completionTimestamp(observedAt), error, }); return { @@ -647,7 +651,7 @@ export class PgCatalogObserver { ); await this.finishRun(runId, { status: "succeeded", - completedAt: new Date().toISOString(), + completedAt: completionTimestamp(observedAt), httpStatus: response.status, observationsEmitted: 0, metadata: { sourceRecordId, contentSha256 }, @@ -726,7 +730,7 @@ export class PgCatalogObserver { ); await this.finishRun(runId, { status, - completedAt: new Date().toISOString(), + completedAt: completionTimestamp(observedAt), httpStatus: response.status, itemCount: remoteModels.length, observationsEmitted, @@ -754,7 +758,7 @@ export class PgCatalogObserver { const message = safeError(error); await this.finishRun(runId, { status: "failed", - completedAt: new Date().toISOString(), + completedAt: completionTimestamp(observedAt), ...(httpStatus === null ? {} : { httpStatus }), error: message, }); From 08ae8d153498c16fbaa2a85f638fa4e4fde5183e Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:26:38 +0800 Subject: [PATCH 14/16] Archive v0.7: avoid observer connection leak on invalid lease --- apps/runner/src/index.ts | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/apps/runner/src/index.ts b/apps/runner/src/index.ts index 5138e7a..632d75c 100644 --- a/apps/runner/src/index.ts +++ b/apps/runner/src/index.ts @@ -155,13 +155,6 @@ async function runStdinMode(): Promise { async function runQueueMode(): Promise { const queue = PgRunJobQueue.connect(databaseUrl, { max: 2 }); - const catalogObserver = catalogObserverEnabled - ? PgCatalogObserver.connect(databaseUrl, { - max: 3, - timeoutMs: catalogObserverTimeoutMs, - maxResponseBytes: catalogObserverMaxResponseBytes, - }) - : null; const workerId = process.env.MODELAPSE_WORKER_ID ?? hostname() + "-" + process.pid + "-" + randomUUID().slice(0, 8); @@ -176,6 +169,14 @@ async function runQueueMode(): Promise { ); } + const catalogObserver = catalogObserverEnabled + ? PgCatalogObserver.connect(databaseUrl, { + max: 3, + timeoutMs: catalogObserverTimeoutMs, + maxResponseBytes: catalogObserverMaxResponseBytes, + }) + : null; + let stopping = false; let catalogObserverTask: Promise | null = null; let nextCatalogObserverAt = 0; From ddce8b06f3f1bb9411ed5d2d1b80535dc73c431f Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:27:16 +0800 Subject: [PATCH 15/16] Archive v0.7: document scheduled catalog collection --- docs/catalog-observer.md | 119 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 119 insertions(+) create mode 100644 docs/catalog-observer.md diff --git a/docs/catalog-observer.md b/docs/catalog-observer.md new file mode 100644 index 0000000..1885362 --- /dev/null +++ b/docs/catalog-observer.md @@ -0,0 +1,119 @@ +# Catalog Observer / Scheduled Collection + +Archive v0.7 turns the v0.6 first-party identity observer into a periodic collection pipeline. + +The collector is intentionally **observe-first**. Fetching a first-party catalog does not, by itself, create a new canonical Model or remap an unknown model ID. Raw first-party snapshots are preserved first; only model IDs that already match a current sourced first-party execution binding are forwarded to the v0.6 identity observer. + +## Data flow + +```text +runner scheduler + | + v +catalog_observer_sources + | + | claim due rows (FOR UPDATE SKIP LOCKED) + v +first-party HTTP fetch + | + +--> source_records + +--> catalog_source_snapshots (append-only raw response) + +--> catalog_collection_runs + | + v +known current first-party bindings only + | + v +observeFirstPartyIdentity(...) + | + +--> model_snapshots + +--> alias_resolution_events + +--> model_execution_bindings + | + v +catalog_identity_drift_events (v0.6) +``` + +Collection and identity ingestion remain separate failure boundaries. A valid HTTP snapshot is retained even if parsing or one model observation fails. + +## Default sources + +Sources are registered lazily when their Provider already exists in the catalog. + +| Provider | Source | Kind | Default interval | Credential | +| --- | --- | --- | ---: | --- | +| OpenAI | `https://api.openai.com/v1/models` | model list | 6 h | `OPENAI_API_KEY` | +| OpenAI | Responses API reference | docs | 24 h | none | +| DeepSeek | `https://api.deepseek.com/models` | model list | 6 h | `DEEPSEEK_API_KEY` | +| DeepSeek | Responses API guide | docs | 24 h | none | + +The schema is source-driven: additional first-party sources can be added to `catalog_observer_sources` without adding another scheduler. + +## Identity policy + +For a `model_list` response, the collector compares returned IDs with current sourced `first_party_direct` execution bindings. + +- A matching ID is sent to `observeFirstPartyIdentity`. +- If the payload exposes an explicit snapshot/version field, it can advance the provider Snapshot and therefore feed the v0.6 drift engine. +- If no explicit snapshot/version is present, the current known Snapshot is preserved. Merely re-fetching a model list cannot erase a Snapshot. +- A remote ID that does not match an existing binding is retained in the collection-run metadata as `unmatchedRemoteModelIds`; it is **not** promoted to a canonical Model. +- A known bound ID missing from the remote list is recorded as `missingKnownApiModelIds`; absence alone does not retire or rewrite the Model. + +This keeps discovery evidence separate from canonical identity decisions. + +## Scheduling and concurrency + +The queue runner enables the Catalog Observer by default. It polls for due sources every 60 seconds, while each source controls its own collection interval through `interval_seconds` and `next_run_at`. + +Due work is claimed in a short PostgreSQL transaction with `FOR UPDATE SKIP LOCKED`. The claim advances `next_run_at` before network I/O. This prevents multiple runner processes from collecting the same source concurrently without holding database locks across HTTP requests. + +The collector itself runs beside the Run queue worker and does not block a provider test job while the HTTP collection is in flight. + +## Collection status + +Each attempt is recorded in `catalog_collection_runs`: + +- `succeeded`: snapshot stored and all eligible observations were emitted. +- `partial`: snapshot stored, but one or more eligible identity observations failed. +- `failed`: HTTP, size, JSON/parser, or persistence failure prevented a complete collection. +- `skipped`: a configured source requires a credential that is not available to the runner. + +Raw response bodies are stored in `catalog_source_snapshots` and are not exposed by the public Archive API. + +## Runner configuration + +Optional environment variables: + +```text +MODELAPSE_CATALOG_OBSERVER_ENABLED=true +MODELAPSE_CATALOG_OBSERVER_POLL_MS=60000 +MODELAPSE_CATALOG_OBSERVER_TIMEOUT_MS=30000 +MODELAPSE_CATALOG_OBSERVER_MAX_RESPONSE_BYTES=2000000 +``` + +Provider API keys already used by the runner are also used for authenticated model-list collection when the corresponding source requires them. + +## One-shot operator collection + +The same collector can be invoked without the long-running runner: + +```bash +npm run collect:first-party-catalog -w @modelapse/catalog-admin +``` + +Required: + +```text +DATABASE_URL=... +``` + +Optional narrowing / override: + +```text +MODELAPSE_PROVIDER_SLUG=deepseek +MODELAPSE_CATALOG_SOURCE_KEY=models-api +MODELAPSE_COLLECT_FORCE=true +MODELAPSE_BUILD= +``` + +`MODELAPSE_COLLECT_FORCE=true` bypasses `next_run_at` for matching enabled sources but still uses the same claim and collection path. From e9d59446b572c93bd37494259cf1364f35f532b7 Mon Sep 17 00:00:00 2001 From: WangEn Date: Thu, 1 Oct 2026 10:27:27 +0800 Subject: [PATCH 16/16] Archive v0.7: document runner observer deployment settings --- docs/deployment/coolify-ghcr.md | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/docs/deployment/coolify-ghcr.md b/docs/deployment/coolify-ghcr.md index bb6fd09..50ea748 100644 --- a/docs/deployment/coolify-ghcr.md +++ b/docs/deployment/coolify-ghcr.md @@ -137,8 +137,17 @@ MODELAPSE_JOB_LEASE_SECONDS=300 MODELAPSE_JOB_POLL_MS=1000 MODELAPSE_WORKER_ID=... MODELAPSE_EVIDENCE_COLLECTOR=... + +MODELAPSE_CATALOG_OBSERVER_ENABLED=true +MODELAPSE_CATALOG_OBSERVER_POLL_MS=60000 +MODELAPSE_CATALOG_OBSERVER_TIMEOUT_MS=30000 +MODELAPSE_CATALOG_OBSERVER_MAX_RESPONSE_BYTES=2000000 ``` +Archive v0.7 enables the Catalog Observer in queue mode by default. The runner only polls PostgreSQL for due observer sources at the configured poll interval; each source owns its longer collection cadence (currently 6 hours for model-list APIs and 24 hours for docs defaults). `FOR UPDATE SKIP LOCKED` coordinates claims across runner processes. Provider API keys configured for execution are reused server-side for authenticated model-list collection. Set `MODELAPSE_CATALOG_OBSERVER_ENABLED=false` to disable scheduled collection without disabling Run execution. + +See `docs/catalog-observer.md` for collection, snapshot, and identity-mapping semantics. + Attach persistent storage at: ```text