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
4 changes: 2 additions & 2 deletions apps/runner/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": "*"
}
}
96 changes: 96 additions & 0 deletions apps/runner/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<ReturnType<typeof runDirectOpenAI>>) {
return {
runId: result.run.id,
Expand All @@ -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 });
Expand Down Expand Up @@ -122,7 +169,50 @@ async function runQueueMode(): Promise<void> {
);
}

const catalogObserver = catalogObserverEnabled
? PgCatalogObserver.connect(databaseUrl, {
max: 3,
timeoutMs: catalogObserverTimeoutMs,
maxResponseBytes: catalogObserverMaxResponseBytes,
})
: null;

let stopping = false;
let catalogObserverTask: Promise<void> | 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;
};
Expand All @@ -131,6 +221,8 @@ async function runQueueMode(): Promise<void> {

try {
while (!stopping) {
maybeCollectCatalog();

const job = await processOneQueuedRunJob({
queue,
repository,
Expand Down Expand Up @@ -163,7 +255,11 @@ async function runQueueMode(): Promise<void> {
await sleep(pollMs);
}
} finally {
if (catalogObserverTask) {
await catalogObserverTask;
}
await queue.close();
await catalogObserver?.close();
}
}

Expand Down
119 changes: 119 additions & 0 deletions docs/catalog-observer.md
Original file line number Diff line number Diff line change
@@ -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=<collector build id>
```

`MODELAPSE_COLLECT_FORCE=true` bypasses `next_run_at` for matching enabled sources but still uses the same claim and collection path.
9 changes: 9 additions & 0 deletions docs/deployment/coolify-ghcr.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 3 additions & 2 deletions packages/catalog-admin/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": "*",
Expand Down
Loading
Loading