diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 53918f74..6c7e4039 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -72,6 +72,7 @@ jobs: message-passing/introduction message-passing/safe-message-handlers openai-agents + openrouter polling/infrequent ) for project in "${projects[@]}"; do diff --git a/.scripts/copy-shared-files.mjs b/.scripts/copy-shared-files.mjs index 36b71515..005ee967 100644 --- a/.scripts/copy-shared-files.mjs +++ b/.scripts/copy-shared-files.mjs @@ -64,6 +64,7 @@ const ESLINTIGNORE_EXCLUDE = [ const POST_CREATE_EXCLUDE = [ 'openai-agents', + 'openrouter', 'google-adk-agents', 'env-config', 'dsl-interpreter', diff --git a/.scripts/list-of-samples.json b/.scripts/list-of-samples.json index 4e635dba..ca313d05 100644 --- a/.scripts/list-of-samples.json +++ b/.scripts/list-of-samples.json @@ -37,6 +37,7 @@ "nexus-standalone-activity", "nexus-standalone-operations", "openai-agents", + "openrouter", "patching-api", "production", "protobufs", diff --git a/README.md b/README.md index adf0bf73..1e6e445f 100644 --- a/README.md +++ b/README.md @@ -190,6 +190,7 @@ and you'll be given the list of sample options. - [**Human in the Loop**](./google-adk-agents/src/human-in-the-loop): A `LongRunningFunctionTool` whose completion is gated by a Temporal Signal or Update. - [**Structured Output**](./google-adk-agents/src/structured-output): Schema-constrained agent output validated at the Workflow boundary. - [**Observability**](./google-adk-agents/src/observability): Token usage, latency, and call counts from the agent loop's OpenTelemetry spans, by composing `OpenTelemetryPlugin` onto the Worker alongside `GoogleAdkPlugin`. +- [**OpenRouter**](./openrouter): Call [OpenRouter](https://openrouter.ai/) from an Activity and fan a prompt batch out with bounded concurrency. Temporal owns the retries, `Retry-After` becomes the next retry delay, and OpenRouter's response cache makes a retried call free. ### Full-stack apps diff --git a/openrouter/.eslintignore b/openrouter/.eslintignore new file mode 100644 index 00000000..7bd99a41 --- /dev/null +++ b/openrouter/.eslintignore @@ -0,0 +1,3 @@ +node_modules +lib +.eslintrc.js \ No newline at end of file diff --git a/openrouter/.eslintrc.js b/openrouter/.eslintrc.js new file mode 100644 index 00000000..9f199cd9 --- /dev/null +++ b/openrouter/.eslintrc.js @@ -0,0 +1,48 @@ +const { builtinModules } = require('module'); + +const ALLOWED_NODE_BUILTINS = new Set(['assert']); + +module.exports = { + root: true, + parser: '@typescript-eslint/parser', + parserOptions: { + project: './tsconfig.json', + tsconfigRootDir: __dirname, + }, + plugins: ['@typescript-eslint', 'deprecation'], + extends: [ + 'eslint:recommended', + 'plugin:@typescript-eslint/eslint-recommended', + 'plugin:@typescript-eslint/recommended', + 'prettier', + ], + rules: { + // recommended for safety + '@typescript-eslint/no-floating-promises': 'error', // forgetting to await Activities and Workflow APIs is bad + 'deprecation/deprecation': 'warn', + + // code style preference + 'object-shorthand': ['error', 'always'], + + // relaxed rules, for convenience + '@typescript-eslint/no-unused-vars': [ + 'warn', + { + argsIgnorePattern: '^_', + varsIgnorePattern: '^_', + }, + ], + '@typescript-eslint/no-explicit-any': 'off', + }, + overrides: [ + { + files: ['src/**/workflows.ts', 'src/**/workflows-*.ts', 'src/**/workflows/*.ts'], + rules: { + 'no-restricted-imports': [ + 'error', + ...builtinModules.filter((m) => !ALLOWED_NODE_BUILTINS.has(m)).flatMap((m) => [m, `node:${m}`]), + ], + }, + }, + ], +}; diff --git a/openrouter/.gitignore b/openrouter/.gitignore new file mode 100644 index 00000000..a9f4ed54 --- /dev/null +++ b/openrouter/.gitignore @@ -0,0 +1,2 @@ +lib +node_modules \ No newline at end of file diff --git a/openrouter/.npmrc b/openrouter/.npmrc new file mode 100644 index 00000000..9cf94950 --- /dev/null +++ b/openrouter/.npmrc @@ -0,0 +1 @@ +package-lock=false \ No newline at end of file diff --git a/openrouter/.nvmrc b/openrouter/.nvmrc new file mode 100644 index 00000000..2bd5a0a9 --- /dev/null +++ b/openrouter/.nvmrc @@ -0,0 +1 @@ +22 diff --git a/openrouter/.post-create b/openrouter/.post-create new file mode 100644 index 00000000..20f776e8 --- /dev/null +++ b/openrouter/.post-create @@ -0,0 +1,22 @@ +To begin development, install the Temporal CLI: + +Mac: {cyan brew install temporal} +Other: Download and extract the latest release from https://github.com/temporalio/cli/releases/latest + +Start Temporal Server: + +{cyan temporal server start-dev} + +Use Node version 18+ (v22.x is recommended): + +Mac: {cyan brew install node@22} +Other: https://nodejs.org/en/download/ + +Set your OpenRouter API key (https://openrouter.ai/settings/keys) in the shell that runs the Worker: + +{cyan export OPENROUTER_API_KEY=sk-or-v1-...} + +Then, in the project directory, using two other shells, run these commands: + +{cyan npm run start.watch} +{cyan npm run workflow} diff --git a/openrouter/.prettierignore b/openrouter/.prettierignore new file mode 100644 index 00000000..7951405f --- /dev/null +++ b/openrouter/.prettierignore @@ -0,0 +1 @@ +lib \ No newline at end of file diff --git a/openrouter/.prettierrc b/openrouter/.prettierrc new file mode 100644 index 00000000..965d50bf --- /dev/null +++ b/openrouter/.prettierrc @@ -0,0 +1,2 @@ +printWidth: 120 +singleQuote: true diff --git a/openrouter/README.md b/openrouter/README.md new file mode 100644 index 00000000..ef21f235 --- /dev/null +++ b/openrouter/README.md @@ -0,0 +1,92 @@ +# OpenRouter + +Call [OpenRouter](https://openrouter.ai/) from a Temporal Activity and fan a prompt batch out, one Activity per prompt. OpenRouter serves hundreds of models from many providers behind one OpenAI-compatible API and one API key, and picks providers and models per request. Temporal handles everything around those calls: retries with backoff, fan-out with bounded concurrency, crash recovery, and a durable record of each prompt's result, cost, and retry history. + +This is the TypeScript port of the Python [`openrouter/prompt_batch`](https://github.com/temporalio/samples-python/tree/main/openrouter/prompt_batch) sample. The Python repo also has [`budget_gate`](https://github.com/temporalio/samples-python/tree/main/openrouter/budget_gate), a batch that pauses instead of failing when the budget or OpenRouter credits run out. + +## What this sample demonstrates + +- One Activity per prompt, run concurrently under a fixed number of runners, so a slow or failing prompt never blocks the others. +- OpenRouter's Auto Router (`openrouter/auto`) choosing a model per prompt, with the chosen model and OpenRouter's reported cost returned for each. +- Temporal-owned retries: the `openai` client is created with `maxRetries: 0`, so every attempt is one HTTP call driven by the Activity retry policy and Event History records the attempt count and last failure; 408, 429, 5xx, and OpenRouter's transient in-flight-budget 402 retry with backoff and honor `Retry-After`; other 4xx errors fail fast and the prompt is reported as skipped instead of failing the batch. Running out of money gets its own failure type, `OpenRouterOutOfCredits`, so a Workflow can pause on it: a 402 for the account or the API key (`error.metadata.limit_source` says which), or the 403 `Key limit exceeded` we have seen a per-key limit return in practice. OpenRouter can also return HTTP 200 with an `error` body and no `choices`, or with a partial answer and an `error` on the choice; the Activity checks for both. +- Retries served from OpenRouter's response cache at $0: the Activity sends `X-OpenRouter-Cache: true`, so if a Worker dies after OpenRouter answered but before Temporal recorded the result, the retried, byte-identical request is a cache hit. +- Heartbeats, so a dead Worker is detected after `heartbeatTimeout` (10s) rather than after the full `startToCloseTimeout`. + +## Running this sample + +1. `temporal server start-dev` to start [Temporal Server](https://github.com/temporalio/cli/#installation). +2. Set an [OpenRouter API key](https://openrouter.ai/settings/keys) in the Worker's environment. A few cents of credit is enough. + ```bash + export OPENROUTER_API_KEY="sk-or-v1-..." + ``` + Optional: `OPENROUTER_HTTP_REFERER` and `OPENROUTER_APP_TITLE` for [app attribution](https://openrouter.ai/docs/app-attribution). +3. `npm install` to install dependencies. +4. `npm run start.watch` to start the Worker. +5. In another shell, `npm run workflow -- "Explain retries in one sentence." "Write a haiku about databases."` to run the batch. + +``` +Starting openrouter-prompt-batch-... + +[deepseek/deepseek-v4-flash-0731] $0.000022 cache=MISS + Q: Explain retries in one sentence. + A: Retries are the automatic re-attempts of a failed operation ... + +[deepseek/deepseek-v4-flash-0731] $0.000525 cache=MISS + Q: Write a haiku about databases. + A: Columns and table, ... + +Reported cost: $0.000547 (what OpenRouter reported on each prompt's final attempt) +Inspect: temporal workflow show -w openrouter-prompt-batch-... +``` + +### See a retry that costs nothing + +`--fail-once` makes each Activity fail its first attempt _after_ OpenRouter has answered, which is what a Worker crash at the wrong moment looks like. The retry re-sends the identical request and OpenRouter serves it from cache: + +```bash +npm run workflow -- --fail-once "Explain idempotency in one sentence." +``` + +``` +[deepseek/deepseek-v4-flash-0731] $0.000000 cache=HIT + Q: Explain idempotency in one sentence. +``` + +`temporal workflow show -w ` shows the Activity completing on attempt 2 with the simulated failure as its last failure; the Worker log has one line per attempt with model, cost, and cache status. The cache is keyed on your API key and the exact request body, so nothing per-attempt goes in the body. OpenRouter writes the cache shortly after the response completes; a retry that arrives before that write lands is a `MISS` and is billed, which you may see occasionally with the one-second retry interval used here. + +### Other options + +- `--model `: any OpenRouter model instead of the Auto Router. +- `--max-concurrency `: how many prompts are in flight at once (default 5). + +## Using OpenRouter's SDKs instead + +This sample uses the `openai` package pointed at `https://openrouter.ai/api/v1`, which is the setup OpenRouter documents for OpenAI-compatible clients; OpenRouter-only fields such as `plugins` go in the request body. OpenRouter's own [`@openrouter/sdk`](https://www.npmjs.com/package/@openrouter/sdk) works too (it is ESM-only). If you use it, construct it with `retryConfig: { strategy: 'none' }`: by default it retries 5xx and connection errors for up to an hour, invisibly to Temporal. + +For agents built on the [Vercel AI SDK](../ai-sdk), [`@openrouter/ai-sdk-provider`](https://www.npmjs.com/package/@openrouter/ai-sdk-provider) is a drop-in `modelProvider` for `AiSdkPlugin`. For the [OpenAI Agents SDK](../openai-agents/src/model-providers), point the provider's `baseURL` at OpenRouter. + +## What Temporal does and does not guarantee + +Activities are at-least-once. If a Worker dies mid-call, the retry re-sends the request; within the cache TTL that retry costs nothing, but two identical requests in flight at the same time both miss the cache and both bill. Completed Activities are never re-run, so a Worker that restarts mid-batch picks up at the first unfinished prompt. + +The reported cost in the result is the sum of what OpenRouter reported on each prompt's final, successful attempt. An attempt that was billed but whose response never made it back to Temporal is not in that number (with `--fail-once`, the first attempt is billed and the result shows the $0 cache hit). For actual spend, use OpenRouter's dashboard or `GET /api/v1/key`. `unknownCostCount` is how many of the results had no cost at all. + +Each Activity adds a few events to the Workflow's Event History, and every answer is part of the Workflow result. The sample caps a batch at 100 prompts; for larger batches, use one Workflow per slice or continue-as-new. + +## Tests + +The tests replace OpenRouter with a fake `fetch` and the Activity with a fake, so they need no API key and make no network calls: + +```bash +npm test +``` + +## Files + +| File | Description | +| -------------------------------------- | --------------------------------------------------------------------------------------------- | +| [src/activities.ts](src/activities.ts) | `callOpenRouter`: one HTTP call per attempt, error classification, cache headers, heartbeats. | +| [src/workflows.ts](src/workflows.ts) | `promptBatch`: bounded fan-out, per-prompt failure handling, retry policy. | +| [src/worker.ts](src/worker.ts) | Builds the OpenRouter client once and runs the Worker. | +| [src/client.ts](src/client.ts) | Starts a batch and prints answer, model, cost, and cache status per prompt. | +| [src/shared.ts](src/shared.ts) | Types shared by client, Workflow, and Activity. | diff --git a/openrouter/package.json b/openrouter/package.json new file mode 100644 index 00000000..68315611 --- /dev/null +++ b/openrouter/package.json @@ -0,0 +1,51 @@ +{ + "name": "temporal-openrouter", + "version": "0.1.0", + "private": true, + "scripts": { + "build": "tsc --build", + "build.watch": "tsc --build --watch", + "format": "prettier --write .", + "format:check": "prettier --check .", + "lint": "eslint .", + "start": "ts-node src/worker.ts", + "start.watch": "nodemon src/worker.ts", + "workflow": "ts-node src/client.ts", + "test": "mocha --exit --require ts-node/register --require source-map-support/register src/mocha/*.test.ts" + }, + "nodemonConfig": { + "execMap": { + "ts": "ts-node" + }, + "ext": "ts", + "watch": [ + "src" + ] + }, + "dependencies": { + "@temporalio/activity": "^1.24.0", + "@temporalio/client": "^1.24.0", + "@temporalio/envconfig": "^1.24.0", + "@temporalio/worker": "^1.24.0", + "@temporalio/workflow": "^1.24.0", + "nanoid": "3.x", + "openai": "^6.0.0" + }, + "devDependencies": { + "@temporalio/testing": "^1.24.0", + "@tsconfig/node22": "^22.0.0", + "@types/mocha": "10.x", + "@types/node": "^22.9.1", + "@typescript-eslint/eslint-plugin": "^8.18.0", + "@typescript-eslint/parser": "^8.18.0", + "eslint": "^8.57.1", + "eslint-config-prettier": "^9.1.0", + "eslint-plugin-deprecation": "^3.0.0", + "mocha": "10.x", + "nodemon": "^3.1.7", + "prettier": "^3.4.2", + "source-map-support": "^0.5.21", + "ts-node": "^10.9.2", + "typescript": "^5.6.3" + } +} diff --git a/openrouter/src/activities.ts b/openrouter/src/activities.ts new file mode 100644 index 00000000..f1d3fea9 --- /dev/null +++ b/openrouter/src/activities.ts @@ -0,0 +1,229 @@ +import OpenAI, { APIError } from 'openai'; +import { ApplicationFailure, Context, log } from '@temporalio/activity'; +import { OPENROUTER_BASE_URL, OpenRouterRequest, OpenRouterResult } from './shared'; + +/** + * OpenAI SDK client pointed at OpenRouter. + * + * Client-side retries are disabled so that Temporal owns every retry: the + * attempt count and last failure land in Event History, and each attempt is + * logged below. (OpenRouter's official `@openrouter/sdk` + * retries 5xx and connection errors for up to an hour by default; if you use it + * instead, pass `retryConfig: { strategy: 'none' }`.) + */ +export function buildClient(apiKey = process.env.OPENROUTER_API_KEY): OpenAI { + if (!apiKey) { + throw new Error('OPENROUTER_API_KEY is required'); + } + const defaultHeaders: Record = {}; + // App attribution is optional. When set, OpenRouter lists your app in its + // public rankings; add X-OpenRouter-App-Visibility: hidden to opt out. + if (process.env.OPENROUTER_HTTP_REFERER) { + defaultHeaders['HTTP-Referer'] = process.env.OPENROUTER_HTTP_REFERER; + } + if (process.env.OPENROUTER_APP_TITLE) { + defaultHeaders['X-OpenRouter-Title'] = process.env.OPENROUTER_APP_TITLE; + } + return new OpenAI({ + baseURL: OPENROUTER_BASE_URL, + apiKey, + maxRetries: 0, + timeout: 60_000, + defaultHeaders, + }); +} + +/** Error type recorded in Event History for an OpenRouter HTTP status. */ +export function errorType(status: number): string { + return `OpenRouterHTTP${status}`; +} + +/** + * Thrown instead of an HTTP status type when the call failed for lack of + * money: a 402 (account or API key out of credits; `error.metadata.limit_source` + * says which) or, as observed in practice, a 403 "Key limit exceeded" for a + * per-key limit. A Workflow can pause on this and resume once someone tops up. + */ +export const OUT_OF_CREDITS = 'OpenRouterOutOfCredits'; + +/** Longest Retry-After the Activity passes through as the next retry delay. */ +const MAX_RETRY_AFTER_SECONDS = 300; + +/** Parse Retry-After in either its delta-seconds or HTTP-date form. */ +function retryAfter(headers: Headers | undefined): string | undefined { + const value = headers?.get('retry-after')?.trim(); + if (!value) return undefined; + let seconds = Number(value); + if (!Number.isFinite(seconds)) { + const delayMs = Date.parse(value) - Date.now(); + if (!Number.isFinite(delayMs)) return undefined; + seconds = Math.ceil(delayMs / 1000); + } + if (seconds <= 0) return undefined; + // Honor the server, within reason: nextRetryDelay overrides the retry + // policy's interval, so cap it rather than park a prompt for hours. + return `${Math.min(seconds, MAX_RETRY_AFTER_SECONDS)}s`; +} + +/** The `error` object OpenRouter returns, as far as this sample reads it. */ +interface OpenRouterErrorBody { + code?: number; + message?: string; + metadata?: { limit_source?: string }; +} + +function errorBody(body: unknown): OpenRouterErrorBody { + if (!body || typeof body !== 'object') return {}; + // openai's APIError.error is already the inner `error` object; a raw + // response body wraps it as `{ error: {...} }`. Accept both. + const inner = 'error' in body ? (body as { error: unknown }).error : body; + return inner && typeof inner === 'object' ? (inner as OpenRouterErrorBody) : {}; +} + +/** + * Turn an OpenRouter error into an ApplicationFailure with the right retry + * posture. Retryable: 408, 429 (honoring Retry-After), any 5xx, and the + * transient in-flight-budget 402. Non-retryable: other 4xx. 400 is a bad + * request, 401 a bad key, 403 a moderation or permission block. Out of money + * is its own type (OUT_OF_CREDITS). + */ +export function throwForStatus(status: number, error: OpenRouterErrorBody, headers?: Headers): never { + const message = error.message ?? ''; + // A 402 from the in-flight budget cap is transient: OpenRouter asks you to + // wait for Retry-After and try again. Every other 402, and the legacy 403 + // "Key limit exceeded", means someone has to add credits. + const transient402 = status === 402 && error.metadata?.limit_source === 'openrouter_in_flight_budget'; + if (!transient402 && (status === 402 || (status === 403 && message.toLowerCase().includes('limit exceeded')))) { + throw ApplicationFailure.create({ + message: `OpenRouter returned HTTP ${status}: ${message}`, + type: OUT_OF_CREDITS, + nonRetryable: true, + details: [{ status }], + }); + } + const retryable = transient402 || status === 408 || status === 429 || status >= 500; + throw ApplicationFailure.create({ + message: `OpenRouter returned HTTP ${status}: ${message}`, + type: errorType(status), + nonRetryable: !retryable, + nextRetryDelay: retryable ? retryAfter(headers) : undefined, + details: [{ status }], + }); +} + +function contentToText(content: unknown): string { + if (typeof content === 'string') return content; + if (!Array.isArray(content)) return ''; + return content + .flatMap((part) => (part && typeof part === 'object' && typeof part.text === 'string' ? [part.text] : [])) + .join('\n'); +} + +export function createActivities(client: OpenAI) { + return { + /** One chat completion. One HTTP call per attempt; Temporal retries. */ + async callOpenRouter(request: OpenRouterRequest): Promise { + const context = Context.current(); + // Heartbeat so a killed Worker is noticed after heartbeatTimeout rather + // than after the full startToCloseTimeout. + const heartbeatMs = context.info.heartbeatTimeoutMs; + const heartbeat = heartbeatMs + ? setInterval(() => context.heartbeat(context.info.attempt), heartbeatMs / 2) + : undefined; + try { + return await send(client, request, context); + } finally { + if (heartbeat) clearInterval(heartbeat); + } + }, + }; +} + +async function send(client: OpenAI, request: OpenRouterRequest, context: Context): Promise { + const attempt = context.info.attempt; + const params: OpenAI.Chat.ChatCompletionCreateParamsNonStreaming & { plugins?: unknown } = { + model: request.model, + messages: [{ role: 'user', content: request.prompt }], + }; + if (request.model === 'openrouter/auto') { + params.plugins = [{ id: 'auto-router', cost_tier: request.costTier }]; + } + + let data: OpenAI.Chat.ChatCompletion & { error?: OpenRouterErrorBody }; + let response: Response; + try { + ({ data, response } = await client.chat.completions + .create(params, { + // Abort the HTTP request if the Activity is cancelled. + signal: context.cancellationSignal, + headers: { + // Ask OpenRouter to cache the successful response. A retry of the + // byte-identical request within the TTL is served from cache and + // billed at $0. + 'X-OpenRouter-Cache': 'true', + 'X-OpenRouter-Cache-TTL': String(request.cacheTtlSeconds), + }, + }) + .withResponse()); + } catch (e) { + if (context.cancellationSignal.aborted) { + // Whatever the request did, the Activity was cancelled; surface that as + // a cancellation, not as a failed call. Rejects with CancelledFailure. + await context.cancelled; + } + if (e instanceof APIError && typeof e.status === 'number') { + const error = errorBody(e.error); + throwForStatus(e.status, { ...error, message: error.message ?? e.message }, e.headers); + } + // Connection errors and timeouts propagate as-is: Temporal retries them. + throw e; + } + + if (data.error) { + // OpenRouter can return HTTP 200 with an error body and no choices when + // the upstream provider failed after the request was accepted. + const error = errorBody(data); + throwForStatus(typeof error.code === 'number' ? error.code : 500, error, response.headers); + } + const choiceError = (data.choices?.[0] as { error?: OpenRouterErrorBody } | undefined)?.error; + if (choiceError) { + // Or a 200 with a partial answer and the provider's error on the choice + // itself; a partial answer is not an answer. + throwForStatus(typeof choiceError.code === 'number' ? choiceError.code : 500, choiceError, response.headers); + } + if (!data.choices?.length) { + // No error and no answer: treat like a server error and retry. + throwForStatus(500, { message: 'Response has no choices' }, response.headers); + } + + const usage = data.usage as (OpenAI.CompletionUsage & { cost?: number }) | undefined; + if (typeof usage?.cost !== 'number') { + // OpenRouter reports cost on every response; if it is ever missing, say + // so rather than pretending the call was free. + log.warn('OpenRouter response has no usage.cost'); + } + const result: OpenRouterResult = { + prompt: request.prompt, + model: data.model, + answer: contentToText(data.choices?.[0]?.message?.content), + costUsd: typeof usage?.cost === 'number' ? usage.cost : null, + generationId: data.id, + cacheStatus: response.headers.get('x-openrouter-cache-status') ?? '', + }; + log.info('OpenRouter call completed', { + attempt, + model: result.model, + costUsd: result.costUsd, + cacheStatus: result.cacheStatus, + generationId: result.generationId, + }); + if (request.failOnceAfterCall && attempt === 1) { + // Demo hook: the Worker "crashes" after the response arrived. The retry + // re-sends the identical request and gets a cache hit. + throw ApplicationFailure.create({ + message: 'Simulated failure after the response was received', + type: 'SimulatedFailure', + }); + } + return result; +} diff --git a/openrouter/src/client.ts b/openrouter/src/client.ts new file mode 100644 index 00000000..295bb12a --- /dev/null +++ b/openrouter/src/client.ts @@ -0,0 +1,66 @@ +import { Connection, Client } from '@temporalio/client'; +import { loadClientConnectConfig } from '@temporalio/envconfig'; +import { nanoid } from 'nanoid'; +import { promptBatch } from './workflows'; +import { DEFAULT_MODEL, TASK_QUEUE } from './shared'; + +const DEFAULT_PROMPTS = ['Explain retries in one sentence.', 'Write a haiku about databases.']; + +async function run() { + // Usage: npm run workflow -- [--fail-once] [--model ] [--max-concurrency ] [prompt ...] + const args = process.argv.slice(2); + let failOnceAfterCall = false; + let model = DEFAULT_MODEL; + let maxConcurrency = 5; + const prompts: string[] = []; + for (let i = 0; i < args.length; i++) { + const arg = args[i]; + if (arg === '--fail-once') failOnceAfterCall = true; + else if (arg === '--model' || arg === '--max-concurrency') { + const value = args[++i]; + if (value === undefined || value.startsWith('--')) throw new Error(`${arg} requires a value`); + if (arg === '--model') model = value; + else { + maxConcurrency = Number(value); + if (!Number.isInteger(maxConcurrency) || maxConcurrency < 1) { + throw new Error('--max-concurrency must be a positive integer'); + } + } + } else if (arg.startsWith('--')) throw new Error(`Unknown flag: ${arg}`); + else prompts.push(arg); + } + + const config = loadClientConnectConfig(); + const connection = await Connection.connect(config.connectionOptions); + const client = new Client({ connection, namespace: config.namespace ?? 'default' }); + + const workflowId = 'openrouter-prompt-batch-' + nanoid(); + console.log(`Starting ${workflowId}`); + const result = await client.workflow.execute(promptBatch, { + taskQueue: TASK_QUEUE, + workflowId, + args: [{ prompts: prompts.length ? prompts : DEFAULT_PROMPTS, model, maxConcurrency, failOnceAfterCall }], + }); + + for (const r of result.results) { + const cost = r.costUsd === null ? 'unknown' : `$${r.costUsd.toFixed(6)}`; + console.log(`\n[${r.model}] ${cost} cache=${r.cacheStatus || '-'}`); + console.log(` Q: ${r.prompt}`); + console.log(` A: ${r.answer.trim()}`); + } + for (const s of result.skipped) { + console.log(`\n[skipped: ${s.reason}] ${s.prompt}`); + } + console.log( + `\nReported cost: $${result.reportedCostUsd.toFixed(6)} (what OpenRouter reported on each prompt's final attempt)`, + ); + if (result.unknownCostCount > 0) { + console.log(` ${result.unknownCostCount} prompt(s) came back without a cost`); + } + console.log(`Inspect: temporal workflow show -w ${workflowId}`); +} + +run().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/openrouter/src/mocha/activities.test.ts b/openrouter/src/mocha/activities.test.ts new file mode 100644 index 00000000..9ce9a4d4 --- /dev/null +++ b/openrouter/src/mocha/activities.test.ts @@ -0,0 +1,269 @@ +import { MockActivityEnvironment } from '@temporalio/testing'; +import { ApplicationFailure, CancelledFailure } from '@temporalio/activity'; +import { describe, it } from 'mocha'; +import assert from 'assert'; +import OpenAI from 'openai'; +import { createActivities } from '../activities'; +import { OPENROUTER_BASE_URL, OpenRouterRequest, OpenRouterResult } from '../shared'; + +type FakeResponse = { status: number; body: unknown; headers?: Record }; + +/** Activities backed by a fake OpenRouter; no network, no API key. */ +function makeActivities(respond: (request: Request) => FakeResponse, seen: Request[] = []) { + const fetch = async (input: string | URL | Request, init?: RequestInit): Promise => { + const request = new Request(input, init); + seen.push(request); + const { status, body, headers } = respond(request); + return new Response(JSON.stringify(body), { + status, + headers: { 'content-type': 'application/json', ...headers }, + }); + }; + const client = new OpenAI({ baseURL: OPENROUTER_BASE_URL, apiKey: 'test-key', maxRetries: 0, fetch }); + return createActivities(client); +} + +const request: OpenRouterRequest = { + prompt: 'Explain retries in one sentence.', + model: 'openrouter/auto', + costTier: 'low', + cacheTtlSeconds: 600, + failOnceAfterCall: false, +}; + +function completion(cost: number | undefined = 0.000123, model = 'openai/gpt-4o-mini') { + return { + id: 'gen-123', + object: 'chat.completion', + created: 0, + model, + choices: [ + { index: 0, finish_reason: 'stop', message: { role: 'assistant', content: 'Retries repeat a failed call.' } }, + ], + usage: { prompt_tokens: 5, completion_tokens: 7, total_tokens: 12, cost }, + }; +} + +async function expectFailure(fn: () => Promise): Promise { + try { + await fn(); + } catch (e) { + assert.ok(e instanceof ApplicationFailure, `expected ApplicationFailure, got ${String(e)}`); + return e; + } + assert.fail('expected the activity to throw'); +} + +describe('callOpenRouter activity', () => { + it('returns model, cost, and cache status, with one HTTP call per attempt', async () => { + const seen: Request[] = []; + const activities = makeActivities( + () => ({ status: 200, body: completion(), headers: { 'X-OpenRouter-Cache-Status': 'MISS' } }), + seen, + ); + + const result = (await new MockActivityEnvironment().run(activities.callOpenRouter, request)) as OpenRouterResult; + + assert.deepStrictEqual(result, { + prompt: request.prompt, + model: 'openai/gpt-4o-mini', + answer: 'Retries repeat a failed call.', + costUsd: 0.000123, + generationId: 'gen-123', + cacheStatus: 'MISS', + }); + assert.strictEqual(seen.length, 1); + const body = (await seen[0].json()) as { model: string; plugins: unknown }; + assert.strictEqual(body.model, 'openrouter/auto'); + assert.deepStrictEqual(body.plugins, [{ id: 'auto-router', cost_tier: 'low' }]); + assert.strictEqual(seen[0].headers.get('x-openrouter-cache'), 'true'); + assert.strictEqual(seen[0].headers.get('x-openrouter-cache-ttl'), '600'); + }); + + it('treats 429 as retryable and honors Retry-After', async () => { + const activities = makeActivities(() => ({ + status: 429, + body: { error: { code: 429, message: 'Rate limited' } }, + headers: { 'Retry-After': '7' }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterHTTP429'); + assert.strictEqual(failure.nonRetryable, false); + assert.strictEqual(failure.nextRetryDelay, '7s'); + }); + + it('treats 402 insufficient credits as out of credits, non-retryable', async () => { + const activities = makeActivities(() => ({ + status: 402, + body: { error: { code: 402, message: 'Insufficient credits' } }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterOutOfCredits'); + assert.strictEqual(failure.nonRetryable, true); + assert.strictEqual(failure.message, 'OpenRouter returned HTTP 402: Insufficient credits'); + }); + + it('retries a transient in-flight-budget 402 after Retry-After', async () => { + const activities = makeActivities(() => ({ + status: 402, + body: { + error: { + code: 402, + message: 'In-flight budget exceeded', + metadata: { limit_source: 'openrouter_in_flight_budget' }, + }, + }, + headers: { 'Retry-After': '3' }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterHTTP402'); + assert.strictEqual(failure.nonRetryable, false); + assert.strictEqual(failure.nextRetryDelay, '3s'); + }); + + it('ignores an empty or non-positive Retry-After', async () => { + for (const value of ['', '-5', '0']) { + const activities = makeActivities(() => ({ + status: 429, + body: { error: { code: 429, message: 'Rate limited' } }, + headers: { 'Retry-After': value }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.nextRetryDelay, undefined, `Retry-After ${JSON.stringify(value)}`); + } + }); + + it('treats a plain 4xx as non-retryable with the server message', async () => { + const activities = makeActivities(() => ({ status: 400, body: { error: { code: 400, message: 'Bad prompt' } } })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterHTTP400'); + assert.strictEqual(failure.nonRetryable, true); + assert.strictEqual(failure.message, 'OpenRouter returned HTTP 400: Bad prompt'); + }); + + it('lets a connection error propagate unchanged so Temporal retries it', async () => { + const fetch = async (): Promise => { + throw new TypeError('fetch failed'); + }; + const client = new OpenAI({ baseURL: OPENROUTER_BASE_URL, apiKey: 'test-key', maxRetries: 0, fetch }); + const activities = createActivities(client); + await assert.rejects( + new MockActivityEnvironment().run(activities.callOpenRouter, request), + (e: unknown) => !(e instanceof ApplicationFailure) && e instanceof Error && /Connection error/.test(e.message), + ); + }); + + it('caps a huge Retry-After', async () => { + const activities = makeActivities(() => ({ + status: 429, + body: { error: { code: 429, message: 'Rate limited' } }, + headers: { 'Retry-After': '1000000000' }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.nextRetryDelay, '300s'); + }); + + it('retries a 200 with no choices and no error', async () => { + const body = { ...completion(), choices: [] }; + const activities = makeActivities(() => ({ status: 200, body })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterHTTP500'); + assert.strictEqual(failure.nonRetryable, false); + }); + + it('treats 403 key limit exceeded as out of credits too', async () => { + const activities = makeActivities(() => ({ + status: 403, + body: { error: { code: 403, message: 'Key limit exceeded (total limit)' } }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterOutOfCredits'); + assert.strictEqual(failure.nonRetryable, true); + }); + + it('classifies an error body inside a 200 by its code', async () => { + const activities = makeActivities(() => ({ + status: 200, + body: { error: { code: 403, message: 'Flagged by moderation' } }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterHTTP403'); + assert.strictEqual(failure.nonRetryable, true); + }); + + it('treats a provider error on the choice as an error, not an answer', async () => { + const body = completion() as ReturnType & { choices: Record[] }; + body.choices[0].finish_reason = 'error'; + body.choices[0].error = { code: 502, message: 'Provider died' }; + const activities = makeActivities(() => ({ status: 200, body })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterHTTP502'); + assert.strictEqual(failure.nonRetryable, false); + assert.match(failure.message, /Provider died/); + }); + + it('with failOnceAfterCall, fails the first attempt only', async () => { + const activities = makeActivities(() => ({ + status: 200, + body: completion(0), + headers: { 'X-OpenRouter-Cache-Status': 'HIT' }, + })); + const failOnce = { ...request, failOnceAfterCall: true }; + + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, failOnce)); + assert.strictEqual(failure.type, 'SimulatedFailure'); + assert.strictEqual(failure.nonRetryable, false); + + const result = (await new MockActivityEnvironment({ attempt: 2 }).run( + activities.callOpenRouter, + failOnce, + )) as OpenRouterResult; + assert.strictEqual(result.cacheStatus, 'HIT'); + assert.strictEqual(result.costUsd, 0); + }); + + it('aborts the HTTP request and surfaces cancellation when the Activity is cancelled', async () => { + let seenSignal: AbortSignal | undefined; + let requestStarted!: () => void; + const started = new Promise((resolve) => (requestStarted = resolve)); + const fetch = (input: string | URL | Request, init?: RequestInit): Promise => { + seenSignal = init?.signal ?? undefined; + requestStarted(); + return new Promise((_, reject) => { + init?.signal?.addEventListener('abort', () => reject(new DOMException('aborted', 'AbortError'))); + }); + }; + const client = new OpenAI({ baseURL: OPENROUTER_BASE_URL, apiKey: 'test-key', maxRetries: 0, fetch }); + const activities = createActivities(client); + const env = new MockActivityEnvironment({ heartbeatTimeoutMs: 1000 }); + + const run = env.run(activities.callOpenRouter, request); + await started; + env.cancel(); + + await assert.rejects(run, (e: unknown) => e instanceof CancelledFailure); + assert.ok(seenSignal?.aborted, 'the request signal should have been aborted'); + }); + + it('reports a missing cost as unknown', async () => { + const body = completion(); + delete (body.usage as { cost?: number }).cost; + const activities = makeActivities(() => ({ status: 200, body })); + const result = (await new MockActivityEnvironment().run(activities.callOpenRouter, request)) as OpenRouterResult; + assert.strictEqual(result.costUsd, null); + assert.strictEqual(result.cacheStatus, ''); + }); + + it('honors an HTTP-date Retry-After', async () => { + const when = new Date(Date.now() + 30_000).toUTCString(); + const activities = makeActivities(() => ({ + status: 503, + body: { error: { code: 503, message: 'No provider available' } }, + headers: { 'Retry-After': when }, + })); + const failure = await expectFailure(() => new MockActivityEnvironment().run(activities.callOpenRouter, request)); + assert.strictEqual(failure.type, 'OpenRouterHTTP503'); + const seconds = Number(String(failure.nextRetryDelay).replace('s', '')); + assert.ok(seconds > 25 && seconds <= 30, `unexpected delay ${failure.nextRetryDelay}`); + }); +}); diff --git a/openrouter/src/mocha/workflows.test.ts b/openrouter/src/mocha/workflows.test.ts new file mode 100644 index 00000000..b82170b5 --- /dev/null +++ b/openrouter/src/mocha/workflows.test.ts @@ -0,0 +1,150 @@ +import { TestWorkflowEnvironment } from '@temporalio/testing'; +import { after, before, describe, it } from 'mocha'; +import { Worker } from '@temporalio/worker'; +import { ApplicationFailure, CancelledFailure, Context } from '@temporalio/activity'; +import { WorkflowFailedError } from '@temporalio/client'; +import { nanoid } from 'nanoid'; +import assert from 'assert'; +import { promptBatch } from '../workflows'; +import { OpenRouterRequest, OpenRouterResult } from '../shared'; + +describe('promptBatch workflow', function () { + this.timeout(30_000); + + let testEnv: TestWorkflowEnvironment; + + before(async () => { + testEnv = await TestWorkflowEnvironment.createLocal(); + }); + + after(async () => { + await testEnv?.teardown(); + }); + + it('collects results and skips prompts that fail with a non-retryable error', async () => { + const taskQueue = 'test-openrouter-' + nanoid(); + const activities = { + async callOpenRouter(request: OpenRouterRequest): Promise { + if (request.prompt === 'bad') { + throw ApplicationFailure.create({ + message: 'OpenRouter returned HTTP 400: bad request', + type: 'OpenRouterHTTP400', + nonRetryable: true, + }); + } + return { + prompt: request.prompt, + model: 'openai/gpt-4o-mini', + answer: `Answer to: ${request.prompt}`, + costUsd: 0.001, + generationId: `gen-${request.prompt}`, + cacheStatus: 'MISS', + }; + }, + }; + + const worker = await Worker.create({ + connection: testEnv.nativeConnection, + taskQueue, + workflowsPath: require.resolve('../workflows'), + activities, + }); + + const result = await worker.runUntil( + testEnv.client.workflow.execute(promptBatch, { + args: [{ prompts: ['one', 'bad', 'two'], maxConcurrency: 2 }], + workflowId: 'test-openrouter-' + nanoid(), + taskQueue, + }), + ); + + assert.deepStrictEqual( + result.results.map((r) => r.prompt), + ['one', 'two'], + ); + assert.deepStrictEqual(result.skipped, [{ prompt: 'bad', reason: 'OpenRouterHTTP400' }]); + assert.strictEqual(result.reportedCostUsd, 0.002); + assert.strictEqual(result.unknownCostCount, 0); + }); + + it('counts prompts whose cost was unknown instead of treating them as free', async () => { + const taskQueue = 'test-openrouter-' + nanoid(); + const activities = { + async callOpenRouter(request: OpenRouterRequest): Promise { + return { + prompt: request.prompt, + model: 'm', + answer: 'ok', + costUsd: request.prompt === 'known' ? 0.001 : null, + generationId: `gen-${request.prompt}`, + cacheStatus: '', + }; + }, + }; + const worker = await Worker.create({ + connection: testEnv.nativeConnection, + taskQueue, + workflowsPath: require.resolve('../workflows'), + activities, + }); + const result = await worker.runUntil( + testEnv.client.workflow.execute(promptBatch, { + args: [{ prompts: ['known', 'unknown'] }], + workflowId: 'test-openrouter-' + nanoid(), + taskQueue, + }), + ); + assert.strictEqual(result.reportedCostUsd, 0.001); + assert.strictEqual(result.unknownCostCount, 1); + }); + + it('propagates Workflow cancellation instead of recording skipped prompts', async () => { + const taskQueue = 'test-openrouter-' + nanoid(); + const activities = { + async callOpenRouter(): Promise { + // Block until cancelled, then surface the cancellation. + await Context.current().cancelled; + throw new Error('unreachable'); + }, + }; + const worker = await Worker.create({ + connection: testEnv.nativeConnection, + taskQueue, + workflowsPath: require.resolve('../workflows'), + activities, + }); + await worker.runUntil(async () => { + const handle = await testEnv.client.workflow.start(promptBatch, { + args: [{ prompts: ['one', 'two'], maxConcurrency: 2 }], + workflowId: 'test-openrouter-' + nanoid(), + taskQueue, + }); + await new Promise((resolve) => setTimeout(resolve, 500)); + await handle.cancel(); + await assert.rejects( + handle.result(), + (err: unknown) => err instanceof WorkflowFailedError && err.cause instanceof CancelledFailure, + ); + }); + }); + + it('rejects a non-positive maxConcurrency', async () => { + const taskQueue = 'test-openrouter-' + nanoid(); + const worker = await Worker.create({ + connection: testEnv.nativeConnection, + taskQueue, + workflowsPath: require.resolve('../workflows'), + activities: { callOpenRouter: async () => assert.fail('should not run') }, + }); + await assert.rejects( + worker.runUntil( + testEnv.client.workflow.execute(promptBatch, { + args: [{ prompts: ['one'], maxConcurrency: 0 }], + workflowId: 'test-openrouter-' + nanoid(), + taskQueue, + }), + ), + (err: unknown) => /maxConcurrency/.test(String((err as { cause?: Error }).cause?.message)), + ); + }); +}); diff --git a/openrouter/src/shared.ts b/openrouter/src/shared.ts new file mode 100644 index 00000000..d6c03a40 --- /dev/null +++ b/openrouter/src/shared.ts @@ -0,0 +1,70 @@ +export const OPENROUTER_BASE_URL = 'https://openrouter.ai/api/v1'; + +// OpenRouter's Auto Router picks a concrete model per request. The response's +// `model` field reports which one it chose. +export const DEFAULT_MODEL = 'openrouter/auto'; + +export const TASK_QUEUE = 'openrouter-prompt-batch'; + +// Each Activity adds a few events to the Workflow's Event History and each +// answer is stored in the Workflow result payload. Keep batches small enough +// to stay well under the history and payload limits. +export const MAX_PROMPTS_PER_BATCH = 100; + +/** + * One chat completion request. Everything here ends up in the request body, + * so keep it free of per-attempt values: OpenRouter's response cache keys on + * the exact body, and a retried attempt should be byte-identical to the first. + */ +export interface OpenRouterRequest { + prompt: string; + model: string; + /** Auto Router cost tier: low, medium, high, xhigh, or max. */ + costTier: 'low' | 'medium' | 'high' | 'xhigh' | 'max'; + /** How long OpenRouter caches a successful response, in seconds. */ + cacheTtlSeconds: number; + /** + * Demo hook: fail the first attempt *after* the response arrives, so the + * retry shows a cache hit billed at $0 in Event History. + */ + failOnceAfterCall: boolean; +} + +export interface OpenRouterResult { + prompt: string; + model: string; + answer: string; + /** What OpenRouter reported for this attempt; null if the response had no usage.cost. */ + costUsd: number | null; + generationId: string; + /** "HIT" or "MISS" from X-OpenRouter-Cache-Status, or "" when absent. */ + cacheStatus: string; +} + +export interface SkippedPrompt { + prompt: string; + reason: string; +} + +export interface BatchInput { + prompts: string[]; + model?: string; + maxConcurrency?: number; + failOnceAfterCall?: boolean; +} + +export interface BatchResult { + results: OpenRouterResult[]; + skipped: SkippedPrompt[]; + /** + * Sum of the cost OpenRouter reported on each prompt's final, successful + * attempt. Attempts that were billed but whose result never reached Temporal + * are not included; OpenRouter's dashboard is the source of truth for spend. + */ + reportedCostUsd: number; + /** + * How many successful prompts came back without a cost. When this is not + * zero, reportedCostUsd is a subtotal of the known costs. + */ + unknownCostCount: number; +} diff --git a/openrouter/src/worker.ts b/openrouter/src/worker.ts new file mode 100644 index 00000000..661e0d40 --- /dev/null +++ b/openrouter/src/worker.ts @@ -0,0 +1,34 @@ +import { NativeConnection, Worker } from '@temporalio/worker'; +import { loadClientConnectConfig } from '@temporalio/envconfig'; +import { buildClient, createActivities } from './activities'; +import { TASK_QUEUE } from './shared'; + +async function run() { + // Same connection settings as the client, so a profile that points at a + // remote server moves both the starter and the Worker. + const config = loadClientConnectConfig(); + const connection = await NativeConnection.connect(config.connectionOptions); + try { + // One OpenRouter client for the Worker's lifetime, shared by every + // concurrent Activity. Reads OPENROUTER_API_KEY from the environment. + const activities = createActivities(buildClient()); + + const worker = await Worker.create({ + connection, + namespace: config.namespace ?? 'default', + taskQueue: TASK_QUEUE, + // Workflows are registered using a path as they run in a separate JS context. + workflowsPath: require.resolve('./workflows'), + activities, + }); + + await worker.run(); + } finally { + await connection.close(); + } +} + +run().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/openrouter/src/workflows.ts b/openrouter/src/workflows.ts new file mode 100644 index 00000000..1d9882d6 --- /dev/null +++ b/openrouter/src/workflows.ts @@ -0,0 +1,82 @@ +import { ActivityFailure, ApplicationFailure, isCancellation, log, proxyActivities } from '@temporalio/workflow'; +import type { createActivities } from './activities'; +import { + BatchInput, + BatchResult, + DEFAULT_MODEL, + MAX_PROMPTS_PER_BATCH, + OpenRouterResult, + SkippedPrompt, +} from './shared'; + +// Temporal owns retries: 1s, 2s, 4s, ... capped at 60s, five attempts. The +// Activity marks 4xx errors non-retryable and passes OpenRouter's Retry-After +// through as the next retry delay, so this policy only governs the rest. +const { callOpenRouter } = proxyActivities>({ + startToCloseTimeout: '90 seconds', + heartbeatTimeout: '10 seconds', + retry: { + initialInterval: '1 second', + backoffCoefficient: 2, + maximumInterval: '60 seconds', + maximumAttempts: 5, + }, +}); + +/** Fan one OpenRouter call out per prompt and collect the answers. */ +export async function promptBatch(batch: BatchInput): Promise { + if (batch.prompts.length > MAX_PROMPTS_PER_BATCH) { + throw ApplicationFailure.nonRetryable( + `Batch has ${batch.prompts.length} prompts; the limit is ${MAX_PROMPTS_PER_BATCH}. ` + + 'Split it, or see the README for the sliding-window pattern.', + ); + } + + const maxConcurrency = batch.maxConcurrency ?? 5; + if (!Number.isInteger(maxConcurrency) || maxConcurrency < 1) { + throw ApplicationFailure.nonRetryable('maxConcurrency must be a positive integer'); + } + + const outcomes: (OpenRouterResult | SkippedPrompt)[] = new Array(batch.prompts.length); + let next = 0; + // Bounded concurrency: N runners pull from the shared prompt list. + const runner = async () => { + while (next < batch.prompts.length) { + const index = next++; + outcomes[index] = await answer(batch.prompts[index], batch); + } + }; + await Promise.all(Array.from({ length: Math.min(maxConcurrency, batch.prompts.length) }, runner)); + + const results = outcomes.filter((o): o is OpenRouterResult => 'answer' in o); + const skipped = outcomes.filter((o): o is SkippedPrompt => 'reason' in o); + return { + results, + skipped, + reportedCostUsd: Number(results.reduce((sum, r) => sum + (r.costUsd ?? 0), 0).toFixed(6)), + unknownCostCount: results.filter((r) => r.costUsd === null).length, + }; +} + +async function answer(prompt: string, batch: BatchInput): Promise { + try { + return await callOpenRouter({ + prompt, + model: batch.model ?? DEFAULT_MODEL, + costTier: 'low', + cacheTtlSeconds: 600, + failOnceAfterCall: batch.failOnceAfterCall ?? false, + }); + } catch (e) { + // Workflow cancellation is not a per-prompt failure, and neither is + // anything other than the Activity itself failing. + if (isCancellation(e) || !(e instanceof ActivityFailure)) throw e; + // One bad prompt should not fail the batch. Record why and carry on; the + // caller decides what to do with skipped prompts. + const cause = e.cause; + const reason = + cause instanceof ApplicationFailure && cause.type ? cause.type : ((cause as Error)?.name ?? 'Unknown'); + log.warn('Skipping prompt', { prompt, reason }); + return { prompt, reason }; + } +} diff --git a/openrouter/tsconfig.json b/openrouter/tsconfig.json new file mode 100644 index 00000000..488f2c62 --- /dev/null +++ b/openrouter/tsconfig.json @@ -0,0 +1,13 @@ +{ + "extends": "@tsconfig/node22/tsconfig.json", + "version": "5.6.3", + "compilerOptions": { + "lib": ["es2021"], + "declaration": true, + "declarationMap": true, + "sourceMap": true, + "rootDir": "./src", + "outDir": "./lib" + }, + "include": ["src/**/*.ts"] +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index e0e52654..a21f7cb5 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -3114,6 +3114,76 @@ importers: specifier: ^5.6.3 version: 5.7.3 + openrouter: + dependencies: + '@temporalio/activity': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/client': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/envconfig': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/worker': + specifier: ^1.24.0 + version: 1.24.0(@swc/helpers@0.5.15) + '@temporalio/workflow': + specifier: ^1.24.0 + version: 1.24.0 + nanoid: + specifier: 3.x + version: 3.3.12 + openai: + specifier: ^6.0.0 + version: 6.49.0(@aws-sdk/credential-provider-node@3.972.71)(@smithy/signature-v4@5.6.9)(ws@8.18.0)(zod@4.4.3) + devDependencies: + '@temporalio/testing': + specifier: ^1.24.0 + version: 1.24.0(@swc/helpers@0.5.15) + '@tsconfig/node22': + specifier: ^22.0.0 + version: 22.0.5 + '@types/mocha': + specifier: 10.x + version: 10.0.10 + '@types/node': + specifier: ^22.9.1 + version: 22.12.0 + '@typescript-eslint/eslint-plugin': + specifier: ^8.18.0 + version: 8.22.0(@typescript-eslint/parser@8.22.0(eslint@8.57.1)(typescript@5.7.3))(eslint@8.57.1)(typescript@5.7.3) + '@typescript-eslint/parser': + specifier: ^8.18.0 + version: 8.22.0(eslint@8.57.1)(typescript@5.7.3) + eslint: + specifier: ^8.57.1 + version: 8.57.1 + eslint-config-prettier: + specifier: ^9.1.0 + version: 9.1.0(eslint@8.57.1) + eslint-plugin-deprecation: + specifier: ^3.0.0 + version: 3.0.0(eslint@8.57.1)(typescript@5.7.3) + mocha: + specifier: 10.x + version: 10.2.0(ts-node@10.9.2(@swc/core@1.10.11(@swc/helpers@0.5.15))(@types/node@22.12.0)(typescript@5.7.3)) + nodemon: + specifier: ^3.1.7 + version: 3.1.9 + prettier: + specifier: ^3.4.2 + version: 3.4.2 + source-map-support: + specifier: ^0.5.21 + version: 0.5.21 + ts-node: + specifier: ^10.9.2 + version: 10.9.2(@swc/core@1.10.11(@swc/helpers@0.5.15))(@types/node@22.12.0)(typescript@5.7.3) + typescript: + specifier: ^5.6.3 + version: 5.7.3 + patching-api: dependencies: '@temporalio/activity':