From ac0a71652b58ff8309d12180b32b80531b65896c Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 14:37:56 -0400 Subject: [PATCH 1/8] feat(hosted): add execution lifecycle primitives --- drizzle/0001_execution_keys.sql | 3 + drizzle/meta/0001_snapshot.json | 1058 +++++++++++++++++++++++ drizzle/meta/_journal.json | 7 + packages/filerouter/src/client.ts | 10 +- packages/filerouter/src/documents.ts | 15 + packages/filerouter/src/hosted.ts | 5 + packages/filerouter/src/index.ts | 4 +- packages/filerouter/src/jobs.ts | 46 +- packages/filerouter/test/client.test.ts | 30 +- packages/filerouter/test/jobs.test.ts | 117 ++- src/api/app.ts | 32 +- src/api/contracts.ts | 50 +- src/db/schema.ts | 2 + src/lib/document-deletion.server.ts | 83 +- src/lib/document-jobs.server.ts | 39 +- src/lib/document-source.server.ts | 63 +- src/lib/hmac.server.ts | 64 ++ src/lib/provider-completion.server.ts | 139 +++ src/workflows/document-workflow.ts | 99 ++- test/worker/api.test.ts | 130 ++- test/worker/provider-completion.test.ts | 146 ++++ test/worker/retention.test.ts | 1 + 22 files changed, 2001 insertions(+), 142 deletions(-) create mode 100644 drizzle/0001_execution_keys.sql create mode 100644 drizzle/meta/0001_snapshot.json create mode 100644 src/lib/hmac.server.ts create mode 100644 src/lib/provider-completion.server.ts create mode 100644 test/worker/provider-completion.test.ts diff --git a/drizzle/0001_execution_keys.sql b/drizzle/0001_execution_keys.sql new file mode 100644 index 0000000..c6bbdbc --- /dev/null +++ b/drizzle/0001_execution_keys.sql @@ -0,0 +1,3 @@ +ALTER TABLE `document_execution` ADD `key` text NOT NULL DEFAULT '';--> statement-breakpoint +UPDATE `document_execution` SET `key` = `provider`;--> statement-breakpoint +CREATE UNIQUE INDEX `document_execution_job_id_key_idx` ON `document_execution` (`job_id`,`key`); diff --git a/drizzle/meta/0001_snapshot.json b/drizzle/meta/0001_snapshot.json new file mode 100644 index 0000000..eef7094 --- /dev/null +++ b/drizzle/meta/0001_snapshot.json @@ -0,0 +1,1058 @@ +{ + "version": "6", + "dialect": "sqlite", + "id": "5cbe51f4-3afe-4440-a6a2-27d75e76a3c3", + "prevId": "16c52000-39d4-457f-9a19-b49bad0e9fc7", + "tables": { + "account": { + "name": "account", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "account_id": { + "name": "account_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "provider_id": { + "name": "provider_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "user_id": { + "name": "user_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "access_token": { + "name": "access_token", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "refresh_token": { + "name": "refresh_token", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "id_token": { + "name": "id_token", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "access_token_expires_at": { + "name": "access_token_expires_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "refresh_token_expires_at": { + "name": "refresh_token_expires_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "scope": { + "name": "scope", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "password": { + "name": "password", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "account_user_id_idx": { + "name": "account_user_id_idx", + "columns": ["user_id"], + "isUnique": false + } + }, + "foreignKeys": { + "account_user_id_user_id_fk": { + "name": "account_user_id_user_id_fk", + "tableFrom": "account", + "tableTo": "user", + "columnsFrom": ["user_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "apikey": { + "name": "apikey", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "config_id": { + "name": "config_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'default'" + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "start": { + "name": "start", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "reference_id": { + "name": "reference_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "prefix": { + "name": "prefix", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "key": { + "name": "key", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "refill_interval": { + "name": "refill_interval", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "refill_amount": { + "name": "refill_amount", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "last_refill_at": { + "name": "last_refill_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "enabled": { + "name": "enabled", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false, + "default": true + }, + "rate_limit_enabled": { + "name": "rate_limit_enabled", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false, + "default": true + }, + "rate_limit_time_window": { + "name": "rate_limit_time_window", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "rate_limit_max": { + "name": "rate_limit_max", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "request_count": { + "name": "request_count", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false, + "default": 0 + }, + "remaining": { + "name": "remaining", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "last_request": { + "name": "last_request", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "expires_at": { + "name": "expires_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "permissions": { + "name": "permissions", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "metadata": { + "name": "metadata", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": { + "apikey_key_unique": { + "name": "apikey_key_unique", + "columns": ["key"], + "isUnique": true + }, + "apikey_config_id_idx": { + "name": "apikey_config_id_idx", + "columns": ["config_id"], + "isUnique": false + }, + "apikey_reference_id_idx": { + "name": "apikey_reference_id_idx", + "columns": ["reference_id"], + "isUnique": false + } + }, + "foreignKeys": { + "apikey_reference_id_user_id_fk": { + "name": "apikey_reference_id_user_id_fk", + "tableFrom": "apikey", + "tableTo": "user", + "columnsFrom": ["reference_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "device_code": { + "name": "device_code", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "device_code": { + "name": "device_code", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "user_code": { + "name": "user_code", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "user_id": { + "name": "user_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "client_id": { + "name": "client_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "scope": { + "name": "scope", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'pending'" + }, + "expires_at": { + "name": "expires_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "last_polled_at": { + "name": "last_polled_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "polling_interval": { + "name": "polling_interval", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": { + "device_code_device_code_unique": { + "name": "device_code_device_code_unique", + "columns": ["device_code"], + "isUnique": true + }, + "device_code_user_code_unique": { + "name": "device_code_user_code_unique", + "columns": ["user_code"], + "isUnique": true + }, + "device_code_user_code_idx": { + "name": "device_code_user_code_idx", + "columns": ["user_code"], + "isUnique": false + } + }, + "foreignKeys": { + "device_code_user_id_user_id_fk": { + "name": "device_code_user_id_user_id_fk", + "tableFrom": "device_code", + "tableTo": "user", + "columnsFrom": ["user_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "document": { + "name": "document", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "user_id": { + "name": "user_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'ready'" + }, + "file_name": { + "name": "file_name", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "content_type": { + "name": "content_type", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "size": { + "name": "size", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "etag": { + "name": "etag", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "object_key": { + "name": "object_key", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "idempotency_key_hash": { + "name": "idempotency_key_hash", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "request_hash": { + "name": "request_hash", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "expires_at": { + "name": "expires_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "document_expires_at_idx": { + "name": "document_expires_at_idx", + "columns": ["expires_at"], + "isUnique": false + }, + "document_user_id_idempotency_key_idx": { + "name": "document_user_id_idempotency_key_idx", + "columns": ["user_id", "idempotency_key_hash"], + "isUnique": true + } + }, + "foreignKeys": { + "document_user_id_user_id_fk": { + "name": "document_user_id_user_id_fk", + "tableFrom": "document", + "tableTo": "user", + "columnsFrom": ["user_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "document_execution": { + "name": "document_execution", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "job_id": { + "name": "job_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "key": { + "name": "key", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "provider": { + "name": "provider", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "position": { + "name": "position", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'queued'" + }, + "outputs": { + "name": "outputs", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "result_key": { + "name": "result_key", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "result_expires_at": { + "name": "result_expires_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "page_count": { + "name": "page_count", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "duration_ms": { + "name": "duration_ms", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "usage": { + "name": "usage", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "error_code": { + "name": "error_code", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "error_message": { + "name": "error_message", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "completed_at": { + "name": "completed_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": { + "document_execution_job_id_idx": { + "name": "document_execution_job_id_idx", + "columns": ["job_id"], + "isUnique": false + }, + "document_execution_result_expires_at_idx": { + "name": "document_execution_result_expires_at_idx", + "columns": ["result_expires_at"], + "isUnique": false + }, + "document_execution_job_id_key_idx": { + "name": "document_execution_job_id_key_idx", + "columns": ["job_id", "key"], + "isUnique": true + } + }, + "foreignKeys": { + "document_execution_job_id_document_job_id_fk": { + "name": "document_execution_job_id_document_job_id_fk", + "tableFrom": "document_execution", + "tableTo": "document_job", + "columnsFrom": ["job_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "document_job": { + "name": "document_job", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "user_id": { + "name": "user_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "document_id": { + "name": "document_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'queued'" + }, + "idempotency_key_hash": { + "name": "idempotency_key_hash", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "request_hash": { + "name": "request_hash", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "metadata": { + "name": "metadata", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "metering_status": { + "name": "metering_status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'pending'" + }, + "error": { + "name": "error", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "document_job_document_id_idx": { + "name": "document_job_document_id_idx", + "columns": ["document_id"], + "isUnique": false + }, + "document_job_created_at_idx": { + "name": "document_job_created_at_idx", + "columns": ["created_at"], + "isUnique": false + }, + "document_job_user_id_idempotency_key_idx": { + "name": "document_job_user_id_idempotency_key_idx", + "columns": ["user_id", "idempotency_key_hash"], + "isUnique": true + } + }, + "foreignKeys": { + "document_job_user_id_user_id_fk": { + "name": "document_job_user_id_user_id_fk", + "tableFrom": "document_job", + "tableTo": "user", + "columnsFrom": ["user_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + }, + "document_job_document_id_document_id_fk": { + "name": "document_job_document_id_document_id_fk", + "tableFrom": "document_job", + "tableTo": "document", + "columnsFrom": ["document_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "session": { + "name": "session", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "expires_at": { + "name": "expires_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "token": { + "name": "token", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "ip_address": { + "name": "ip_address", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "user_agent": { + "name": "user_agent", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "user_id": { + "name": "user_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "session_token_unique": { + "name": "session_token_unique", + "columns": ["token"], + "isUnique": true + }, + "session_user_id_idx": { + "name": "session_user_id_idx", + "columns": ["user_id"], + "isUnique": false + } + }, + "foreignKeys": { + "session_user_id_user_id_fk": { + "name": "session_user_id_user_id_fk", + "tableFrom": "session", + "tableTo": "user", + "columnsFrom": ["user_id"], + "columnsTo": ["id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "user": { + "name": "user", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "email": { + "name": "email", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "email_verified": { + "name": "email_verified", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": false + }, + "image": { + "name": "image", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "is_anonymous": { + "name": "is_anonymous", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false, + "default": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "user_email_unique": { + "name": "user_email_unique", + "columns": ["email"], + "isUnique": true + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "verification": { + "name": "verification", + "columns": { + "id": { + "name": "id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "identifier": { + "name": "identifier", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "value": { + "name": "value", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "expires_at": { + "name": "expires_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "verification_identifier_idx": { + "name": "verification_identifier_idx", + "columns": ["identifier"], + "isUnique": false + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + } + }, + "views": {}, + "enums": {}, + "_meta": { + "schemas": {}, + "tables": {}, + "columns": {} + }, + "internal": { + "indexes": {} + } +} diff --git a/drizzle/meta/_journal.json b/drizzle/meta/_journal.json index adfe67a..15b48d3 100644 --- a/drizzle/meta/_journal.json +++ b/drizzle/meta/_journal.json @@ -8,6 +8,13 @@ "when": 1785020813925, "tag": "0000_baseline", "breakpoints": true + }, + { + "idx": 1, + "version": "6", + "when": 1785252291672, + "tag": "0001_execution_keys", + "breakpoints": true } ] } diff --git a/packages/filerouter/src/client.ts b/packages/filerouter/src/client.ts index ee9a15b..f8a526c 100644 --- a/packages/filerouter/src/client.ts +++ b/packages/filerouter/src/client.ts @@ -6,6 +6,7 @@ import type { HostedDocumentCreateOptions, HostedDocumentDeleteOptions, HostedDocumentGetOptions, + HostedDocumentReleaseOptions, } from "./documents" import { FileRouterError } from "./errors" import { HostedExecutions } from "./executions" @@ -16,10 +17,10 @@ import type { import { FILEROUTER_DEFAULT_API_URL } from "./hosted" import type { HostedCompareResult, + HostedProviderTarget, HostedJob, HostedParseResult, HostedProviderOptions, - HostedProviderTarget, } from "./hosted" import { describeInput } from "./internal/input" import { HostedTransport } from "./internal/hosted-transport" @@ -124,7 +125,7 @@ export class FileRouter { options, async ({ documentId, job, signal }) => { const execution = job.executions.find( - (candidate) => candidate.provider === provider + (candidate) => candidate.key === provider ) if (!execution || execution.status !== "complete") { throw executionError(execution?.error?.message ?? job.error) @@ -164,8 +165,9 @@ export class FileRouter { providers: results, resources: { documentId, - executions: job.executions.map(({ id, provider }) => ({ + executions: job.executions.map(({ id, key, provider }) => ({ id, + key, provider, })), jobId: job.id, @@ -313,6 +315,7 @@ function providerTarget( ...(options.pageFields && { pageFields: options.pageFields }), ...(options.pages && { pages: options.pages }), ...(value !== undefined && { providerOptions: value }), + key: provider, provider, } as HostedProviderTarget } @@ -349,6 +352,7 @@ export type { HostedDocumentCreateOptions, HostedDocumentDeleteOptions, HostedDocumentGetOptions, + HostedDocumentReleaseOptions, HostedExecutionResultOptions, HostedExecutionWaitOptions, HostedJobCreateInput, diff --git a/packages/filerouter/src/documents.ts b/packages/filerouter/src/documents.ts index fbc42af..feede2e 100644 --- a/packages/filerouter/src/documents.ts +++ b/packages/filerouter/src/documents.ts @@ -23,6 +23,7 @@ export interface HostedDocumentGetOptions { } export type HostedDocumentDeleteOptions = HostedDocumentGetOptions +export type HostedDocumentReleaseOptions = HostedDocumentGetOptions export interface FileRouterDocuments { create( @@ -31,6 +32,7 @@ export interface FileRouterDocuments { ): Promise delete(id: string, options?: HostedDocumentDeleteOptions): Promise get(id: string, options?: HostedDocumentGetOptions): Promise + release(id: string, options?: HostedDocumentReleaseOptions): Promise } export class HostedDocuments implements FileRouterDocuments { @@ -99,6 +101,19 @@ export class HostedDocuments implements FileRouterDocuments { } ) } + + async release( + id: string, + options: HostedDocumentReleaseOptions = {} + ): Promise { + await this.#transport.request( + `${HOSTED_DOCUMENTS_PATH}/${encodeURIComponent(id)}/release`, + { + method: "POST", + ...(options.signal && { signal: options.signal }), + } + ) + } } async function resolveUpload( diff --git a/packages/filerouter/src/hosted.ts b/packages/filerouter/src/hosted.ts index b103637..773222b 100644 --- a/packages/filerouter/src/hosted.ts +++ b/packages/filerouter/src/hosted.ts @@ -15,6 +15,7 @@ export const FILEROUTER_CLI_SCOPE = "jobs:create jobs:read" export const FILEROUTER_DEFAULT_API_URL = "https://filerouter.dev" export const MAX_HOSTED_JOB_REQUEST_BYTES = 64 * 1024 export const MAX_HOSTED_METADATA_ENTRIES = 50 +export const MAX_HOSTED_JOB_EXECUTIONS = 25 export const HOSTED_DOCUMENTS_PATH = "/api/v1/documents" export const HOSTED_EXECUTIONS_PATH = "/api/v1/executions" @@ -53,6 +54,7 @@ export interface HostedDocument { interface HostedProviderTargetBase { includeRaw?: boolean + key: string outputs?: Array pageFields?: Array pages?: Array @@ -76,6 +78,7 @@ export interface HostedExecution { error?: { code?: string; message: string } id: string jobId: string + key: string outputs: Array pageCount?: number provider: ProviderId @@ -97,6 +100,7 @@ export interface HostedJob { } export interface HostedJobAccepted { + executions: Array id: string status: HostedJobStatus } @@ -109,6 +113,7 @@ export interface HostedProvider { export interface HostedExecutionReference { id: string + key: string provider: ProviderId } diff --git a/packages/filerouter/src/index.ts b/packages/filerouter/src/index.ts index 9f5f036..6c447d9 100644 --- a/packages/filerouter/src/index.ts +++ b/packages/filerouter/src/index.ts @@ -2,7 +2,6 @@ export { DirectFileRouter, assertProviderOutputs, assertProviderPageFields, - createDirectFileRouter, serializeProviderError, type DirectFileRouterOptions, } from "./router" @@ -42,6 +41,7 @@ export { HOSTED_JOBS_PATH, HOSTED_PROVIDERS_PATH, MAX_HOSTED_JOB_REQUEST_BYTES, + MAX_HOSTED_JOB_EXECUTIONS, MAX_HOSTED_METADATA_ENTRIES, hostedDocumentStatuses, hostedExecutionStatuses, @@ -53,6 +53,7 @@ export { type HostedExecution, type HostedExecutionReference, type HostedExecutionStatus, + type HostedProviderTarget, type HostedJob, type HostedJobAccepted, type HostedJobStatus, @@ -60,7 +61,6 @@ export { type HostedParseResult, type HostedProvider, type HostedProviderOptions, - type HostedProviderTarget, } from "./hosted" export { normalizeDocumentFileName, diff --git a/packages/filerouter/src/jobs.ts b/packages/filerouter/src/jobs.ts index 40c79e8..6d60eca 100644 --- a/packages/filerouter/src/jobs.ts +++ b/packages/filerouter/src/jobs.ts @@ -2,15 +2,16 @@ import { FileRouterError } from "./errors" import { HOSTED_JOBS_PATH, MAX_HOSTED_JOB_REQUEST_BYTES, + MAX_HOSTED_JOB_EXECUTIONS, MAX_HOSTED_METADATA_ENTRIES, } from "./hosted" import type { HostedExecution, + HostedExecutionReference, + HostedProviderTarget, HostedJob, HostedJobAccepted, - HostedProviderTarget, } from "./hosted" -import type { ProviderId } from "./catalog" import type { HostedTransport } from "./internal/hosted-transport" import { assertPageFields } from "./internal/provider-options" import { abortableSleep } from "./internal/sleep" @@ -58,7 +59,7 @@ export interface FileRouterJobs { ): Promise waitForExecution( job: HostedJobAccepted | string, - provider: ProviderId, + execution: HostedExecutionReference | string, options?: HostedExecutionWaitOptions ): Promise } @@ -113,14 +114,19 @@ export class HostedJobs implements FileRouterJobs { waitForExecution( job: HostedJobAccepted | string, - provider: ProviderId, + execution: HostedExecutionReference | string, options: HostedExecutionWaitOptions = {} ): Promise { return withTimeout( options.timeoutMs ?? DEFAULT_HOSTED_JOB_TIMEOUT_MS, options.signal, (signal) => - this.#waitForExecution(jobId(job), provider, signal, options.onStatus) + this.#waitForExecution( + jobId(job), + executionId(execution), + signal, + options.onStatus + ) ) } @@ -146,19 +152,19 @@ export class HostedJobs implements FileRouterJobs { async #waitForExecution( id: string, - provider: ProviderId, + executionId: string, signal: AbortSignal, onStatus?: (execution: HostedExecution) => void ): Promise { let previousStatus: HostedExecution["status"] | undefined for await (const job of this.#pollJobs(id, signal)) { const execution = job.executions.find( - (candidate) => candidate.provider === provider + (candidate) => candidate.id === executionId ) if (!execution) { if (job.status === "complete" || job.status === "failed") { throw new FileRouterError( - `Hosted job ${id} does not include provider ${provider}.`, + `Hosted job ${id} does not include execution ${executionId}.`, { code: "InvalidInput" } ) } @@ -173,8 +179,8 @@ export class HostedJobs implements FileRouterJobs { } if (job.status === "complete" || job.status === "failed") { throw new FileRouterError( - `Hosted job ${id} ended before ${provider} reached a terminal status.`, - { code: "ParseFailed", providerId: provider } + `Hosted job ${id} ended before execution ${executionId} reached a terminal status.`, + { code: "ParseFailed" } ) } } @@ -195,6 +201,10 @@ function jobId(job: HostedJobAccepted | string): string { return typeof job === "string" ? job : job.id } +function executionId(execution: HostedExecutionReference | string): string { + return typeof execution === "string" ? execution : execution.id +} + export function assertHostedJobDraft(input: HostedJobCreateDraft): void { serializeJobInput({ ...input, @@ -208,11 +218,17 @@ function serializeJobInput(input: HostedJobCreateInput): string { code: "InvalidInput", }) } + if (input.providers.length > MAX_HOSTED_JOB_EXECUTIONS) { + throw new FileRouterError( + `Hosted jobs accept at most ${MAX_HOSTED_JOB_EXECUTIONS} provider executions.`, + { code: "InvalidInput" } + ) + } if ( - new Set(input.providers.map((target) => target.provider)).size !== + new Set(input.providers.map((provider) => provider.key)).size !== input.providers.length ) { - throw new FileRouterError("Each provider may appear only once.", { + throw new FileRouterError("Each provider key must be unique.", { code: "InvalidInput", }) } @@ -226,6 +242,12 @@ function serializeJobInput(input: HostedJobCreateInput): string { ) } for (const target of input.providers) { + if (target.key.trim().length === 0 || target.key.length > 64) { + throw new FileRouterError( + "Provider keys must be non-blank and at most 64 characters.", + { code: "InvalidInput" } + ) + } assertPageFields( target.outputs ?? [DEFAULT_PARSE_OUTPUT], target.pageFields diff --git a/packages/filerouter/test/client.test.ts b/packages/filerouter/test/client.test.ts index 48758d4..c874cfe 100644 --- a/packages/filerouter/test/client.test.ts +++ b/packages/filerouter/test/client.test.ts @@ -37,7 +37,13 @@ describe("FileRouter", () => { }) expect(readJsonBody(fetchMock, 1)).toEqual({ documentId: "document-1", - providers: [{ outputs: ["markdown"], provider: "llamaparse" }], + providers: [ + { + key: "llamaparse", + outputs: ["markdown"], + provider: "llamaparse", + }, + ], }) }) @@ -115,7 +121,7 @@ describe("FileRouter", () => { }) }) - test("maps provider options onto explicit execution targets", async () => { + test("maps provider options onto explicit provider entries", async () => { const fetchMock = vi .fn() .mockResolvedValueOnce(Response.json(document("document-3"))) @@ -142,6 +148,7 @@ describe("FileRouter", () => { expect(readJsonBody(fetchMock, 1)).toMatchObject({ providers: [ { + key: "llamaparse", pageFields: ["markdown"], pages: [1, 3], provider: "llamaparse", @@ -179,8 +186,8 @@ describe("FileRouter", () => { expect(comparison.resources).toEqual({ documentId: "document-4", executions: [ - { id: "execution-4a", provider: "llamaparse" }, - { id: "execution-4b", provider: "liteparse" }, + { id: "execution-4a", key: "llamaparse", provider: "llamaparse" }, + { id: "execution-4b", key: "liteparse", provider: "liteparse" }, ], jobId: "job-4", }) @@ -272,6 +279,20 @@ describe("FileRouter", () => { expect.objectContaining({ method: "DELETE" }) ) }) + + test("releases stored artifacts without deleting job history", async () => { + const fetchMock = vi + .fn() + .mockResolvedValueOnce(new Response(null, { status: 204 })) + + await expect( + createClient(fetchMock).documents.release("document-6") + ).resolves.toBeUndefined() + expect(fetchMock).toHaveBeenCalledWith( + "https://example.com/api/v1/documents/document-6/release", + expect.objectContaining({ method: "POST" }) + ) + }) }) function createClient(fetchMock: typeof fetch): FileRouter { @@ -305,6 +326,7 @@ function execution( durationMs: 10, id, jobId: "job", + key: provider, outputs: ["markdown"], provider, resultAvailable: true, diff --git a/packages/filerouter/test/jobs.test.ts b/packages/filerouter/test/jobs.test.ts index 9b79f5c..0d1c3d2 100644 --- a/packages/filerouter/test/jobs.test.ts +++ b/packages/filerouter/test/jobs.test.ts @@ -1,31 +1,58 @@ import { describe, expect, test, vi } from "vite-plus/test" import { FileRouter } from "../src/client" -import { MAX_HOSTED_JOB_REQUEST_BYTES } from "../src/hosted" +import { + MAX_HOSTED_JOB_REQUEST_BYTES, + MAX_HOSTED_JOB_EXECUTIONS, +} from "../src/hosted" describe("hosted resources", () => { test("creates recoverable jobs from stored documents", async () => { - const fetchMock = vi - .fn() - .mockResolvedValue(Response.json({ id: "job-1", status: "queued" })) + const fetchMock = vi.fn().mockResolvedValue( + Response.json({ + executions: [ + { + id: "execution-1", + key: "primary", + provider: "llamaparse", + }, + ], + id: "job-1", + status: "queued", + }) + ) const client = createClient(fetchMock) await expect( client.jobs.create( { documentId: "document-1", - providers: [{ outputs: ["markdown"], provider: "llamaparse" }], + providers: [ + { + key: "primary", + outputs: ["markdown"], + provider: "llamaparse", + }, + ], }, { idempotencyKey: "job-create-1" } ) - ).resolves.toEqual({ id: "job-1", status: "queued" }) + ).resolves.toEqual({ + executions: [ + { id: "execution-1", key: "primary", provider: "llamaparse" }, + ], + id: "job-1", + status: "queued", + }) const headers = new Headers(fetchMock.mock.calls[0]?.[1]?.headers) expect(headers.get("idempotency-key")).toBe("job-create-1") await expect( new Request("https://example.com", fetchMock.mock.calls[0]?.[1]).json() ).resolves.toEqual({ documentId: "document-1", - providers: [{ outputs: ["markdown"], provider: "llamaparse" }], + providers: [ + { key: "primary", outputs: ["markdown"], provider: "llamaparse" }, + ], }) }) @@ -46,7 +73,7 @@ describe("hosted resources", () => { expect(statuses).toEqual(["queued", "running", "complete"]) }) - test("waits for one provider without waiting for the whole job", async () => { + test("waits for one execution without waiting for the whole job", async () => { const fetchMock = vi .fn() .mockResolvedValueOnce( @@ -68,9 +95,13 @@ describe("hosted resources", () => { const statuses: Array = [] await expect( - createClient(fetchMock).jobs.waitForExecution("job-2", "liteparse", { - onStatus: (value) => statuses.push(value.status), - }) + createClient(fetchMock).jobs.waitForExecution( + "job-2", + "execution-liteparse", + { + onStatus: (value) => statuses.push(value.status), + } + ) ).resolves.toMatchObject({ provider: "liteparse", status: "complete" }) expect(statuses).toEqual(["queued", "complete"]) expect(fetchMock).toHaveBeenCalledTimes(2) @@ -110,7 +141,34 @@ describe("hosted resources", () => { expect(fetchMock).toHaveBeenCalledOnce() }) - test("rejects ambiguous or oversized jobs before sending them", async () => { + test("accepts repeated providers with distinct execution keys", async () => { + const fetchMock = vi.fn().mockResolvedValue( + Response.json({ + executions: [ + { + id: "execution-1", + key: "accurate OCR", + provider: "llamaparse", + }, + { id: "execution-2", key: "second", provider: "llamaparse" }, + ], + id: "job-1", + status: "queued", + }) + ) + + await createClient(fetchMock).jobs.create({ + documentId: "document-1", + providers: [ + { key: "accurate OCR", provider: "llamaparse" }, + { key: "second", provider: "llamaparse" }, + ], + }) + + expect(fetchMock).toHaveBeenCalledOnce() + }) + + test("rejects invalid or oversized jobs before sending them", async () => { const fetchMock = vi.fn() const jobs = createClient(fetchMock).jobs const base = { @@ -120,7 +178,10 @@ describe("hosted resources", () => { await expect( jobs.create({ ...base, - providers: [{ provider: "llamaparse" }, { provider: "llamaparse" }], + providers: [ + { key: "duplicate", provider: "llamaparse" }, + { key: "duplicate", provider: "liteparse" }, + ], }) ).rejects.toMatchObject({ code: "InvalidInput" }) await expect(jobs.create({ ...base, providers: [] })).rejects.toMatchObject( @@ -131,7 +192,19 @@ describe("hosted resources", () => { await expect( jobs.create({ ...base, - providers: [{ pageFields: ["markdown"], provider: "llamaparse" }], + providers: [{ key: " ", provider: "llamaparse" }], + }) + ).rejects.toMatchObject({ code: "InvalidInput" }) + await expect( + jobs.create({ + ...base, + providers: [ + { + key: "primary", + pageFields: ["markdown"], + provider: "llamaparse", + }, + ], }) ).rejects.toMatchObject({ code: "InvalidInput" }) await expect( @@ -140,7 +213,7 @@ describe("hosted resources", () => { metadata: Object.fromEntries( Array.from({ length: 51 }, (_, index) => [`key-${index}`, "value"]) ), - providers: [{ provider: "llamaparse" }], + providers: [{ key: "primary", provider: "llamaparse" }], }) ).rejects.toMatchObject({ code: "InvalidInput" }) await expect( @@ -148,6 +221,7 @@ describe("hosted resources", () => { ...base, providers: [ { + key: "primary", providerOptions: { agentic_options: { custom_prompt: "x".repeat(MAX_HOSTED_JOB_REQUEST_BYTES), @@ -158,6 +232,18 @@ describe("hosted resources", () => { ], }) ).rejects.toMatchObject({ code: "InvalidInput" }) + await expect( + jobs.create({ + ...base, + providers: Array.from( + { length: MAX_HOSTED_JOB_EXECUTIONS + 1 }, + (_, index) => ({ + key: `target-${index}`, + provider: "liteparse" as const, + }) + ), + }) + ).rejects.toMatchObject({ code: "InvalidInput" }) expect(fetchMock).not.toHaveBeenCalled() }) }) @@ -193,6 +279,7 @@ function execution( createdAt: "2026-07-18T00:00:00.000Z", id: `execution-${provider}`, jobId: "job", + key: provider, outputs: ["markdown"], provider, resultAvailable: status === "complete", diff --git a/src/api/app.ts b/src/api/app.ts index 7d170b1..b010e9f 100644 --- a/src/api/app.ts +++ b/src/api/app.ts @@ -24,6 +24,7 @@ import { getJobRoute, IdempotencyKeySchema, listProvidersRoute, + releaseDocumentRoute, } from "@/api/contracts" import { problemResponse } from "@/api/problem" import { @@ -40,7 +41,10 @@ import { } from "@/lib/document-jobs.server" import { getProviderSourceResponse } from "@/lib/document-source.server" import { createDocument, getDocument } from "@/lib/documents.server" -import { deleteDocument } from "@/lib/document-deletion.server" +import { + deleteDocument, + releaseDocumentArtifacts, +} from "@/lib/document-deletion.server" import { HttpError } from "@/lib/http.server" import { hostedProviderCatalog } from "@/lib/hosted-providers.server" import { @@ -48,6 +52,7 @@ import { serializeError, type WideEvent, } from "@/observability/log" +import { receiveProviderCompletion } from "@/lib/provider-completion.server" type ApiRequestEvent = Partial & { credential_id?: string @@ -210,6 +215,19 @@ api.notFound((context) => { ) }) +api.post( + "/api/v1/provider-completions/:jobId/:executionId", + async (context) => { + const { executionId, jobId } = context.req.param() + return receiveProviderCompletion( + context.req.raw, + context.env, + jobId, + executionId + ) + } +) + const limitDocumentJsonBody = bodyLimit({ maxSize: MAX_HOSTED_JOB_REQUEST_BYTES, onError: () => { @@ -233,6 +251,7 @@ api.use( `${HOSTED_DOCUMENTS_PATH}/:documentId`, requireApiKey((method) => (method === "DELETE" ? "create" : "read")) ) +api.use(`${HOSTED_DOCUMENTS_PATH}/:documentId/release`, requireApiKey("create")) api.use( HOSTED_JOBS_PATH, requireApiKey("create"), @@ -338,6 +357,17 @@ api.openapi(deleteDocumentRoute, async (context) => { return context.body(null, 204) }) +api.openapi(releaseDocumentRoute, async (context) => { + const { documentId } = context.req.valid("param") + context.get("requestEvent").document_id = documentId + await releaseDocumentArtifacts( + documentId, + context.get("principal").userId, + context.env + ) + return context.body(null, 204) +}) + api.openapi(createJobRoute, async (context) => { const input = context.req.valid("json") const idempotencyKey = context.req.valid("header")["idempotency-key"] diff --git a/src/api/contracts.ts b/src/api/contracts.ts index bd79478..c490e2d 100644 --- a/src/api/contracts.ts +++ b/src/api/contracts.ts @@ -9,6 +9,7 @@ import { hostedExecutionStatuses, hostedJobStatuses, MAX_HOSTED_METADATA_ENTRIES, + MAX_HOSTED_JOB_EXECUTIONS, } from "@file_router/sdk/hosted" import { providerIds } from "@file_router/sdk/catalog" @@ -65,6 +66,12 @@ const DocumentSchema = z const ProviderTargetSchema = z .object({ includeRaw: z.boolean().optional(), + key: z + .string() + .max(64) + .refine((value) => value.trim().length > 0, { + message: "Provider key cannot be blank.", + }), outputs: z.array(ParseOutputSchema).min(1).optional(), pageFields: z.array(ParsePageFieldSchema).min(1).optional(), pages: z.array(z.number().int().positive()).min(1).optional(), @@ -86,15 +93,18 @@ export const CreateJobRequestSchema = z } ) .optional(), - providers: z.array(ProviderTargetSchema).min(1).max(providerIds.length), + providers: z + .array(ProviderTargetSchema) + .min(1) + .max(MAX_HOSTED_JOB_EXECUTIONS), }) .strict() .superRefine((value, context) => { - const providers = value.providers.map((target) => target.provider) - if (new Set(providers).size !== providers.length) { + const keys = value.providers.map((provider) => provider.key) + if (new Set(keys).size !== keys.length) { context.addIssue({ code: "custom", - message: "Each provider may appear only once per job.", + message: "Each provider key must be unique per job.", path: ["providers"], }) } @@ -114,6 +124,7 @@ const ExecutionSchema = z error: ExecutionErrorSchema.optional(), id: ExecutionIdSchema, jobId: JobIdSchema, + key: z.string(), outputs: z.array(ParseOutputSchema), pageCount: z.number().int().nonnegative().optional(), provider: ProviderIdSchema, @@ -144,7 +155,17 @@ const JobSchema = z .openapi("Job") const JobAcceptedSchema = z - .object({ id: JobIdSchema, status: z.enum(hostedJobStatuses) }) + .object({ + executions: z.array( + z.object({ + id: ExecutionIdSchema, + key: z.string(), + provider: ProviderIdSchema, + }) + ), + id: JobIdSchema, + status: z.enum(hostedJobStatuses), + }) .openapi("JobAccepted") const ProblemSchema = z @@ -240,6 +261,25 @@ export const deleteDocumentRoute = createRoute({ tags: ["Documents"], }) +export const releaseDocumentRoute = createRoute({ + description: + "Deletes retained source and result artifacts after document jobs finish while preserving job and execution records.", + method: "post", + path: `${HOSTED_DOCUMENTS_PATH}/{documentId}/release`, + request: { params: z.object({ documentId: DocumentIdSchema }) }, + responses: { + 204: { description: "Document artifacts released" }, + 400: problem, + 401: problem, + 409: problem, + 429: problem, + 500: problem, + }, + security: [{ BearerAuth: [] }], + summary: "Release document artifacts", + tags: ["Documents"], +}) + export const createJobRoute = createRoute({ description: "Runs one or more provider executions against a stored document.", diff --git a/src/db/schema.ts b/src/db/schema.ts index 10ec92b..ba6f415 100644 --- a/src/db/schema.ts +++ b/src/db/schema.ts @@ -248,6 +248,7 @@ export const documentExecution = sqliteTable( jobId: text("job_id") .notNull() .references(() => documentJob.id, { onDelete: "cascade" }), + key: text("key").notNull(), provider: text("provider").$type().notNull(), position: integer("position").notNull(), status: text("status", { enum: hostedExecutionStatuses }) @@ -279,5 +280,6 @@ export const documentExecution = sqliteTable( (table) => [ index("document_execution_job_id_idx").on(table.jobId), index("document_execution_result_expires_at_idx").on(table.resultExpiresAt), + uniqueIndex("document_execution_job_id_key_idx").on(table.jobId, table.key), ] ) diff --git a/src/lib/document-deletion.server.ts b/src/lib/document-deletion.server.ts index 3fbe898..f68586a 100644 --- a/src/lib/document-deletion.server.ts +++ b/src/lib/document-deletion.server.ts @@ -1,8 +1,89 @@ -import { and, eq, inArray } from "drizzle-orm" +import { and, eq, inArray, notExists } from "drizzle-orm" import { document, documentExecution, documentJob } from "@/db/schema" import { createDb } from "@/db/server" import { deleteR2Objects } from "@/lib/r2-objects.server" +import { HttpError } from "@/lib/http.server" + +const activeJobStatuses = ["queued", "running"] as const + +export async function releaseDocumentArtifacts( + id: string, + userId: string, + env: Cloudflare.Env +): Promise { + const db = createDb(env.DB) + const now = new Date() + const activeJobs = db + .select({ id: documentJob.id }) + .from(documentJob) + .where( + and( + eq(documentJob.documentId, document.id), + inArray(documentJob.status, activeJobStatuses) + ) + ) + const releasable = await db + .update(document) + .set({ status: "expired", updatedAt: now }) + .where( + and( + eq(document.id, id), + eq(document.userId, userId), + notExists(activeJobs) + ) + ) + .returning({ objectKey: document.objectKey }) + .get() + if (!releasable) { + const stored = await db + .select({ id: document.id }) + .from(document) + .where(and(eq(document.id, id), eq(document.userId, userId))) + .get() + if (!stored) { + return + } + throw new HttpError(409, "Document still has active jobs.", { + code: "document_active", + }) + } + + const executions = await db + .select({ + id: documentExecution.id, + resultKey: documentExecution.resultKey, + }) + .from(documentExecution) + .innerJoin(documentJob, eq(documentExecution.jobId, documentJob.id)) + .where(eq(documentJob.documentId, id)) + .all() + await deleteR2Objects(env.FILEROUTER_FILES, [ + releasable.objectKey, + ...executions.map((execution) => execution.resultKey), + ]) + + const releaseDocument = db + .update(document) + .set({ objectKey: null, updatedAt: now }) + .where(and(eq(document.id, id), eq(document.userId, userId))) + if (executions.length === 0) { + await releaseDocument + return + } + await db.batch([ + releaseDocument, + db + .update(documentExecution) + .set({ resultExpiresAt: now, resultKey: null, updatedAt: now }) + .where( + inArray( + documentExecution.id, + executions.map((execution) => execution.id) + ) + ), + ]) +} export async function deleteDocument( id: string, diff --git a/src/lib/document-jobs.server.ts b/src/lib/document-jobs.server.ts index e64af86..0903470 100644 --- a/src/lib/document-jobs.server.ts +++ b/src/lib/document-jobs.server.ts @@ -36,10 +36,11 @@ export interface CreateJobResult { } type NormalizedTarget = DocumentWorkflowTarget & { + key: string position: number } -type UnvalidatedProviderTarget = Omit< +type UnvalidatedExecutionTarget = Omit< HostedProviderTarget, "providerOptions" > & { @@ -47,7 +48,7 @@ type UnvalidatedProviderTarget = Omit< } type CreateDocumentJobInput = Omit & { - providers: Array + providers: Array } export async function createDocumentJob( @@ -133,6 +134,7 @@ export async function createDocumentJob( createdAt: now, id: target.executionId, jobId, + key: target.key, outputs: target.outputs, position: target.position, provider: target.provider, @@ -172,7 +174,7 @@ export async function createDocumentJob( }, jobId, requestId, - targets: targets.map(({ position: _, ...target }) => target), + targets: targets.map(({ key: _, position: __, ...target }) => target), userId, } try { @@ -182,7 +184,18 @@ export async function createDocumentJob( throw error } - return { job: { id: jobId, status: "queued" }, replayed: false } + return { + job: { + executions: targets.map(({ executionId, key, provider }) => ({ + id: executionId, + key, + provider, + })), + id: jobId, + status: "queued", + }, + replayed: false, + } } function documentIsAvailable( @@ -296,6 +309,7 @@ function normalizeTargets( pageFields: [...new Set(target.pageFields)], }), ...(target.pages && { pages: [...new Set(target.pages)] }), + key: target.key, position, provider: target.provider, ...(target.providerOptions && { @@ -357,6 +371,7 @@ function serializeExecution( }), id: execution.id, jobId: execution.jobId, + key: execution.key, outputs: execution.outputs, ...(execution.pageCount !== null && { pageCount: execution.pageCount }), provider: execution.provider, @@ -419,8 +434,22 @@ async function replayJob( { code: "idempotency_conflict" } ) } + const executions = await db + .select({ + id: documentExecution.id, + key: documentExecution.key, + provider: documentExecution.provider, + }) + .from(documentExecution) + .where(eq(documentExecution.jobId, existing.id)) + .orderBy(documentExecution.position) + .all() return { - job: { id: existing.id, status: existing.status }, + job: { + executions, + id: existing.id, + status: existing.status, + }, replayed: true, } } diff --git a/src/lib/document-source.server.ts b/src/lib/document-source.server.ts index 3f4f065..bb133e0 100644 --- a/src/lib/document-source.server.ts +++ b/src/lib/document-source.server.ts @@ -1,9 +1,10 @@ import { normalizeDocumentFileName } from "@file_router/sdk" +import { signHmac, verifyHmac } from "@/lib/hmac.server" import { HttpError } from "@/lib/http.server" const SOURCE_URL_TTL_SECONDS = 30 * 60 -const encoder = new TextEncoder() +const SOURCE_TOKEN_PURPOSE = "filerouter-source-v2" export async function createProviderSourceUrl( env: Cloudflare.Env, @@ -159,13 +160,11 @@ async function signSourceToken( fileName: string, expires: number ): Promise { - const key = await sourceSigningKey(secret, ["sign"]) - const signature = await crypto.subtle.sign( - "HMAC", - key, - sourceTokenPayload(documentId, fileName, expires).buffer as ArrayBuffer + return signHmac( + secret, + SOURCE_TOKEN_PURPOSE, + sourceTokenPayload(documentId, fileName, expires) ) - return toBase64Url(new Uint8Array(signature)) } async function verifySourceToken( @@ -175,29 +174,11 @@ async function verifySourceToken( expires: number, token: string ): Promise { - try { - const key = await sourceSigningKey(secret, ["verify"]) - return crypto.subtle.verify( - "HMAC", - key, - fromBase64Url(token).buffer as ArrayBuffer, - sourceTokenPayload(documentId, fileName, expires).buffer as ArrayBuffer - ) - } catch { - return false - } -} - -function sourceSigningKey( - secret: string, - usages: KeyUsage[] -): Promise { - return crypto.subtle.importKey( - "raw", - encoder.encode(`filerouter-source-v2:${secret}`), - { hash: "SHA-256", name: "HMAC" }, - false, - usages + return verifyHmac( + secret, + SOURCE_TOKEN_PURPOSE, + sourceTokenPayload(documentId, fileName, expires), + token ) } @@ -205,8 +186,8 @@ function sourceTokenPayload( documentId: string, fileName: string, expires: number -): Uint8Array { - return encoder.encode(`${documentId}\n${fileName}\n${expires}`) +): string { + return `${documentId}\n${fileName}\n${expires}` } function parseExpiration(value: string | undefined): number | undefined { @@ -218,21 +199,3 @@ function parseExpiration(value: string | undefined): number | undefined { ? expires : undefined } - -function toBase64Url(value: Uint8Array): string { - return btoa(String.fromCharCode(...value)) - .replaceAll("+", "-") - .replaceAll("/", "_") - .replace(/=+$/, "") -} - -function fromBase64Url(value: string): Uint8Array { - if (!/^[A-Za-z0-9_-]+$/.test(value)) { - throw new Error("Invalid source token.") - } - const padded = value - .replaceAll("-", "+") - .replaceAll("_", "/") - .padEnd(Math.ceil(value.length / 4) * 4, "=") - return Uint8Array.from(atob(padded), (character) => character.charCodeAt(0)) -} diff --git a/src/lib/hmac.server.ts b/src/lib/hmac.server.ts new file mode 100644 index 0000000..4c248eb --- /dev/null +++ b/src/lib/hmac.server.ts @@ -0,0 +1,64 @@ +const encoder = new TextEncoder() + +export async function signHmac( + secret: string, + purpose: string, + payload: string +): Promise { + const signature = await crypto.subtle.sign( + "HMAC", + await signingKey(secret, purpose, ["sign"]), + encoder.encode(payload) + ) + return toBase64Url(new Uint8Array(signature)) +} + +export async function verifyHmac( + secret: string, + purpose: string, + payload: string, + signature: string +): Promise { + try { + return crypto.subtle.verify( + "HMAC", + await signingKey(secret, purpose, ["verify"]), + fromBase64Url(signature).buffer as ArrayBuffer, + encoder.encode(payload) + ) + } catch { + return false + } +} + +function signingKey( + secret: string, + purpose: string, + usages: KeyUsage[] +): Promise { + return crypto.subtle.importKey( + "raw", + encoder.encode(`${purpose}:${secret}`), + { hash: "SHA-256", name: "HMAC" }, + false, + usages + ) +} + +function toBase64Url(value: Uint8Array): string { + return btoa(String.fromCharCode(...value)) + .replaceAll("+", "-") + .replaceAll("/", "_") + .replace(/=+$/, "") +} + +function fromBase64Url(value: string): Uint8Array { + if (!/^[A-Za-z0-9_-]+$/.test(value)) { + throw new Error("Invalid HMAC signature.") + } + const padded = value + .replaceAll("-", "+") + .replaceAll("_", "/") + .padEnd(Math.ceil(value.length / 4) * 4, "=") + return Uint8Array.from(atob(padded), (character) => character.charCodeAt(0)) +} diff --git a/src/lib/provider-completion.server.ts b/src/lib/provider-completion.server.ts new file mode 100644 index 0000000..5571e28 --- /dev/null +++ b/src/lib/provider-completion.server.ts @@ -0,0 +1,139 @@ +import type { ProviderId } from "@file_router/sdk/catalog" + +import { signHmac, verifyHmac } from "@/lib/hmac.server" +import { HttpError } from "@/lib/http.server" + +const COMPLETION_TOKEN_PURPOSE = "filerouter-provider-completion-v1" +const COMPLETION_URL_TTL_SECONDS = 30 * 60 +const TERMINAL_WORKFLOW_STATUSES = new Set([ + "complete", + "errored", + "terminated", + "unknown", +]) + +type CompletionProvider = "datalab" | "llamaparse" + +export function providerCompletionEventType( + provider: ProviderId, + executionId: string +): string | undefined { + return isCompletionProvider(provider) + ? executionCompletionEventType(executionId) + : undefined +} + +export async function withProviderCompletion( + env: Cloudflare.Env, + jobId: string, + executionId: string, + provider: ProviderId, + providerOptions: Record | undefined +): Promise | undefined> { + if (!isCompletionProvider(provider)) { + return providerOptions + } + const callbackUrl = await createCallbackUrl(env, jobId, executionId) + return provider === "datalab" + ? { ...providerOptions, webhook_url: callbackUrl } + : { + ...providerOptions, + webhook_configurations: [ + { + webhook_events: [ + "parse.success", + "parse.error", + "parse.partial_success", + "parse.cancelled", + ], + webhook_output_format: "json", + webhook_url: callbackUrl, + }, + ], + } +} + +export async function receiveProviderCompletion( + request: Request, + env: Cloudflare.Env, + jobId: string, + executionId: string +): Promise { + const url = new URL(request.url) + const expires = readExpiration(url.searchParams.get("expires")) + const token = url.searchParams.get("token") + if ( + !expires || + !token || + !(await verifyHmac( + env.BETTER_AUTH_SECRET, + COMPLETION_TOKEN_PURPOSE, + tokenPayload(jobId, executionId, expires), + token + )) + ) { + throw new HttpError(404, "Provider completion not found.", { + code: "provider_completion_not_found", + }) + } + + const instance = await env.DOCUMENT_WORKFLOW.get(jobId) + try { + await instance.sendEvent({ + payload: {}, + type: executionCompletionEventType(executionId), + }) + } catch (error) { + const status = await instance.status().catch(() => undefined) + if (!status || !TERMINAL_WORKFLOW_STATUSES.has(status.status)) { + throw error + } + } + return new Response(null, { status: 204 }) +} + +function isCompletionProvider(value: string): value is CompletionProvider { + return value === "datalab" || value === "llamaparse" +} + +function executionCompletionEventType(executionId: string): string { + return `provider-${executionId}` +} + +async function createCallbackUrl( + env: Cloudflare.Env, + jobId: string, + executionId: string +): Promise { + const expires = Math.floor(Date.now() / 1000) + COMPLETION_URL_TTL_SECONDS + const token = await signHmac( + env.BETTER_AUTH_SECRET, + COMPLETION_TOKEN_PURPOSE, + tokenPayload(jobId, executionId, expires) + ) + const url = new URL( + `/api/v1/provider-completions/${encodeURIComponent(jobId)}/${encodeURIComponent(executionId)}`, + new URL(env.BETTER_AUTH_URL).origin + ) + url.searchParams.set("expires", String(expires)) + url.searchParams.set("token", token) + return url.toString() +} + +function tokenPayload( + jobId: string, + executionId: string, + expires: number +): string { + return `${jobId}\n${executionId}\n${expires}` +} + +function readExpiration(value: string | null): number | undefined { + if (!value || !/^\d{10}$/.test(value)) { + return undefined + } + const expires = Number(value) + return Number.isSafeInteger(expires) && expires >= Date.now() / 1000 + ? expires + : undefined +} diff --git a/src/workflows/document-workflow.ts b/src/workflows/document-workflow.ts index 6e77239..13dd4c9 100644 --- a/src/workflows/document-workflow.ts +++ b/src/workflows/document-workflow.ts @@ -26,6 +26,10 @@ import { captureServerTelemetry } from "@/integrations/posthog/server" import { createProviderSourceUrl } from "@/lib/document-source.server" import { resultExpiresAt } from "@/lib/document-retention" import { createHostedProviders } from "@/lib/hosted-providers.server" +import { + providerCompletionEventType, + withProviderCompletion, +} from "@/lib/provider-completion.server" import { emitWideEvent, serializeError } from "@/observability/log" import { storeProviderResult, @@ -60,13 +64,14 @@ const PROVIDER_STATUS_STEP = { timeout: "1 minute", } as const satisfies WorkflowStepConfig +const PROVIDER_COMPLETION_TIMEOUT = "10 seconds" +const PROVIDER_POLL_INTERVAL = "2 seconds" + const METERING_STEP = { retries: { backoff: "exponential", delay: "5 seconds", limit: 5 }, timeout: "1 minute", } as const satisfies WorkflowStepConfig -const SYNC_PROVIDER_ATTEMPTS = 4 - type ProviderStepStatus = | { status: "pending" | "running" } | { error: string; status: "failed" } @@ -113,6 +118,7 @@ export class DocumentWorkflow extends WorkflowEntrypoint< configured[target.provider], target, input, + params.jobId, this.env ) ) @@ -188,6 +194,7 @@ async function processExecution( provider: FileRouterProvider, target: DocumentWorkflowTarget, input: ProviderInput, + jobId: string, env: Cloudflare.Env ): Promise { const startedAt = Date.now() @@ -195,7 +202,7 @@ async function processExecution( try { outcome = provider.jobs - ? await processAsyncProvider(step, provider, target, input, env) + ? await processAsyncProvider(step, provider, target, input, jobId, env) : await processSyncProvider(step, provider, target, input, env) } catch (error) { outcome = { @@ -218,36 +225,19 @@ async function processSyncProvider( input: ProviderInput, env: Cloudflare.Env ): Promise> { - for (let attempt = 1; attempt <= SYNC_PROVIDER_ATTEMPTS; attempt += 1) { - try { - return await step.do( - `process ${target.executionId} ${attempt}`, - PROVIDER_EXECUTION_STEP, - async () => - storeProviderResult( - env.FILEROUTER_FILES, - target.executionId, - selectPageFields( - await provider.parse(input, parseOptions(target)), - target.pageFields - ) - ) - ) - } catch (error) { - if ( - attempt === SYNC_PROVIDER_ATTEMPTS || - !FileRouterError.isInstance(error) || - !error.retryable - ) { - throw error - } - await step.sleep( - `wait to retry ${target.executionId} ${attempt}`, - `${attempt} second${attempt === 1 ? "" : "s"}` + return step.do( + `process ${target.executionId}`, + PROVIDER_EXECUTION_STEP, + async () => + storeProviderResult( + env.FILEROUTER_FILES, + target.executionId, + selectPageFields( + await provider.parse(input, parseOptions(target)), + target.pageFields + ) ) - } - } - throw new Error("Sync provider retry loop exhausted.") + ) } async function processAsyncProvider( @@ -255,16 +245,31 @@ async function processAsyncProvider( provider: FileRouterProvider, target: DocumentWorkflowTarget, input: ProviderInput, + jobId: string, env: Cloudflare.Env ): Promise> { const jobs = provider.jobs if (!jobs) { throw new Error(`Provider ${provider.id} does not support durable jobs.`) } + const completionEvent = providerCompletionEventType( + target.provider, + target.executionId + ) + const options = parseOptions(target) const job = await step.do( `submit ${target.executionId}`, PROVIDER_EXECUTION_STEP, - () => jobs.submit(input, parseOptions(target)) + async () => { + const providerOptions = await withProviderCompletion( + env, + jobId, + target.executionId, + target.provider, + target.providerOptions + ) + return jobs.submit(input, parseOptions(target, providerOptions)) + } ) const deadline = new Date(job.submittedAt).getTime() + 14 * 60 * 1000 let attempt = 0 @@ -275,7 +280,7 @@ async function processAsyncProvider( `check ${target.executionId} ${attempt}`, PROVIDER_STATUS_STEP, async () => { - const current = await jobs.get(job, parseOptions(target)) + const current = await jobs.get(job, options) if (current.status !== "complete") { return current } @@ -295,7 +300,22 @@ async function processAsyncProvider( providerId: provider.id, }) } - await step.sleep(`wait for ${target.executionId} ${attempt}`, "10 seconds") + if (completionEvent) { + try { + await step.waitForEvent(`wait for ${target.executionId} ${attempt}`, { + timeout: PROVIDER_COMPLETION_TIMEOUT, + type: completionEvent, + }) + } catch { + // A timeout is the polling watchdog when a provider webhook is delayed + // or lost. The next provider status request remains authoritative. + } + } else { + await step.sleep( + `wait for ${target.executionId} ${attempt}`, + PROVIDER_POLL_INTERVAL + ) + } } throw new FileRouterError(`${provider.id} job timed out.`, { @@ -442,16 +462,19 @@ async function markWorkflowFailure( }) } -function parseOptions(target: DocumentWorkflowTarget) { +function parseOptions( + target: DocumentWorkflowTarget, + providerOptions = target.providerOptions +) { return { includeRaw: target.includeRaw, outputs: target.outputs, ...(target.pageFields && { pageFields: target.pageFields }), ...(target.pages && { pages: target.pages }), provider: target.provider, - ...(target.providerOptions && { + ...(providerOptions && { providerOptions: { - [target.provider]: target.providerOptions, + [target.provider]: providerOptions, } satisfies ProviderParseOptions, }), } diff --git a/test/worker/api.test.ts b/test/worker/api.test.ts index c89fad6..7d3aeb9 100644 --- a/test/worker/api.test.ts +++ b/test/worker/api.test.ts @@ -22,6 +22,9 @@ describe("FileRouter Worker", () => { expect(response.status).toBe(200) expect(specification.openapi).toBe("3.1.0") expect(specification.paths).toHaveProperty("/api/v1/documents") + expect(specification.paths).toHaveProperty( + "/api/v1/documents/{documentId}/release" + ) expect(specification.paths).toHaveProperty("/api/v1/jobs") expect(specification.paths).toHaveProperty( "/api/v1/executions/{executionId}/result" @@ -96,8 +99,13 @@ describe("FileRouter Worker", () => { body: JSON.stringify({ documentId: document.id, providers: [ - { outputs: ["markdown"], provider: "llamaparse" }, { + key: "primary", + outputs: ["markdown"], + provider: "llamaparse", + }, + { + key: "fallback", outputs: ["pages"], pageFields: ["markdown"], provider: "liteparse", @@ -115,7 +123,14 @@ describe("FileRouter Worker", () => { testEnv ) expect(jobResponse.status).toBe(202) - const accepted = await jobResponse.json<{ id: string }>() + const accepted = await jobResponse.json<{ + executions: Array<{ id: string; key: string; provider: string }> + id: string + }>() + expect(accepted.executions).toEqual([ + expect.objectContaining({ key: "primary", provider: "llamaparse" }), + expect.objectContaining({ key: "fallback", provider: "liteparse" }), + ]) expect(create).toHaveBeenCalledWith( expect.objectContaining({ id: accepted.id, @@ -143,15 +158,15 @@ describe("FileRouter Worker", () => { const storedJob = await job.json<{ createdAt: string documentId: string - executions: Array<{ provider: string; status: string }> + executions: Array<{ key: string; provider: string; status: string }> id: string status: string }>() expect(storedJob).toMatchObject({ documentId: document.id, executions: [ - { provider: "llamaparse", status: "queued" }, - { provider: "liteparse", status: "queued" }, + { key: "primary", provider: "llamaparse", status: "queued" }, + { key: "fallback", provider: "liteparse", status: "queued" }, ], id: accepted.id, status: "queued", @@ -176,6 +191,18 @@ describe("FileRouter Worker", () => { ) expect(deniedDelete.status).toBe(401) + const activeRelease = await api.fetch( + new Request( + `https://filerouter.test/api/v1/documents/${document.id}/release`, + { + headers: { Authorization: `Bearer ${apiKey}` }, + method: "POST", + } + ), + testEnv + ) + expect(activeRelease.status).toBe(409) + const deleted = await api.fetch( new Request(`https://filerouter.test/api/v1/documents/${document.id}`, { headers: { Authorization: `Bearer ${apiKey}` }, @@ -268,7 +295,7 @@ describe("FileRouter Worker", () => { new Request("https://filerouter.test/api/v1/jobs", { body: JSON.stringify({ documentId: stored.id, - providers: [{ provider: "llamaparse" }], + providers: [{ key: "primary", provider: "llamaparse" }], }), headers: { Authorization: `Bearer ${apiKey}`, @@ -294,6 +321,97 @@ describe("FileRouter Worker", () => { await cleanupUser(userId) }) + test("releases artifacts without deleting completed job history", async () => { + const userId = "user-release-artifacts" + const apiKey = await createApiKey(userId) + const testEnv = envWithWorkflow(vi.fn().mockResolvedValue({})) + const documentResponse = await api.fetch( + new Request("https://filerouter.test/api/v1/documents", { + body: "%PDF-release", + headers: { + Authorization: `Bearer ${apiKey}`, + "Content-Type": "application/octet-stream", + "Idempotency-Key": "document-release-1", + }, + method: "POST", + }), + testEnv + ) + const storedDocument = await documentResponse.json<{ id: string }>() + const jobResponse = await api.fetch( + new Request("https://filerouter.test/api/v1/jobs", { + body: JSON.stringify({ + documentId: storedDocument.id, + providers: [{ key: "primary", provider: "liteparse" }], + }), + headers: { + Authorization: `Bearer ${apiKey}`, + "Content-Type": "application/json", + "Idempotency-Key": "job-release-1", + }, + method: "POST", + }), + testEnv + ) + const accepted = await jobResponse.json<{ + executions: Array<{ id: string }> + id: string + }>() + const executionId = accepted.executions[0]?.id + if (!executionId) { + throw new Error("Expected an execution reference.") + } + const resultKey = `executions/${executionId}/result.json` + await env.FILEROUTER_FILES.put(resultKey, "{}") + const now = Math.floor(Date.now() / 1_000) + await env.DB.batch([ + env.DB.prepare( + "UPDATE document_job SET status = 'complete' WHERE id = ?" + ).bind(accepted.id), + env.DB.prepare( + "UPDATE document_execution SET status = 'complete', duration_ms = 10, page_count = 1, result_key = ?, result_expires_at = ?, completed_at = ? WHERE id = ?" + ).bind(resultKey, now + 3_600, now, executionId), + ]) + + const released = await api.fetch( + new Request( + `https://filerouter.test/api/v1/documents/${storedDocument.id}/release`, + { + headers: { Authorization: `Bearer ${apiKey}` }, + method: "POST", + } + ), + testEnv + ) + + expect(released.status).toBe(204) + expect( + await env.FILEROUTER_FILES.head(`documents/${storedDocument.id}/source`) + ).toBeNull() + expect(await env.FILEROUTER_FILES.head(resultKey)).toBeNull() + const preservedJob = await api.fetch( + new Request(`https://filerouter.test/api/v1/jobs/${accepted.id}`, { + headers: { Authorization: `Bearer ${apiKey}` }, + }), + testEnv + ) + expect(preservedJob.status).toBe(200) + await expect(preservedJob.json()).resolves.toMatchObject({ + executions: [ + { + id: executionId, + key: "primary", + resultAvailable: false, + status: "complete", + }, + ], + id: accepted.id, + status: "complete", + }) + + await cleanupUser(userId) + }) + test("rejects oversized job requests at the HTTP boundary", async () => { const userId = "user-job-size" const apiKey = await createApiKey(userId) diff --git a/test/worker/provider-completion.test.ts b/test/worker/provider-completion.test.ts new file mode 100644 index 0000000..5d45659 --- /dev/null +++ b/test/worker/provider-completion.test.ts @@ -0,0 +1,146 @@ +import { env } from "cloudflare:workers" +import { describe, expect, test, vi } from "vite-plus/test" + +import { api } from "@/api/app" +import { + providerCompletionEventType, + withProviderCompletion, +} from "@/lib/provider-completion.server" + +describe("provider completion", () => { + test("injects only provider-owned completion options", async () => { + const datalab = await withProviderCompletion( + env, + "job-1", + "execution-1", + "datalab", + { mode: "balanced" } + ) + expect(datalab).toMatchObject({ + mode: "balanced", + webhook_url: expect.stringContaining( + "/api/v1/provider-completions/job-1/execution-1" + ), + }) + expect(providerCompletionEventType("datalab", "execution-1")).toBe( + "provider-execution-1" + ) + + const llamaparse = await withProviderCompletion( + env, + "job-1", + "execution-2", + "llamaparse", + { tier: "fast" } + ) + expect(llamaparse).toMatchObject({ + tier: "fast", + webhook_configurations: [ + { + webhook_events: expect.arrayContaining([ + "parse.success", + "parse.error", + ]), + webhook_output_format: "json", + webhook_url: expect.stringContaining( + "/api/v1/provider-completions/job-1/execution-2" + ), + }, + ], + }) + expect(providerCompletionEventType("llamaparse", "execution-2")).toBe( + "provider-execution-2" + ) + + await expect( + withProviderCompletion( + env, + "job-1", + "execution-3", + "mistral-ocr", + undefined + ) + ).resolves.toBeUndefined() + expect( + providerCompletionEventType("mistral-ocr", "execution-3") + ).toBeUndefined() + }) + + test("authenticates callbacks before waking the matching execution", async () => { + const sendEvent = vi.fn().mockResolvedValue(undefined) + const get = vi.fn().mockResolvedValue({ + sendEvent, + status: vi.fn().mockResolvedValue({ status: "waiting" }), + }) + const testEnv = { + ...env, + DOCUMENT_WORKFLOW: { + create: vi.fn(), + createBatch: vi.fn(), + get, + } as Cloudflare.Env["DOCUMENT_WORKFLOW"], + } + const providerOptions = await withProviderCompletion( + testEnv, + "job-2", + "execution-4", + "datalab", + undefined + ) + const callbackUrl = providerOptions?.webhook_url + if (typeof callbackUrl !== "string") { + throw new Error("Expected a Datalab completion URL.") + } + + const accepted = await api.fetch( + new Request(callbackUrl, { method: "POST" }), + testEnv + ) + expect(accepted.status).toBe(204) + expect(get).toHaveBeenCalledWith("job-2") + expect(sendEvent).toHaveBeenCalledWith({ + payload: {}, + type: "provider-execution-4", + }) + + const invalidUrl = new URL(callbackUrl) + invalidUrl.searchParams.set("token", "invalid") + const rejected = await api.fetch( + new Request(invalidUrl, { method: "POST" }), + testEnv + ) + expect(rejected.status).toBe(404) + expect(sendEvent).toHaveBeenCalledOnce() + }) + + test("accepts late callbacks after the workflow is terminal", async () => { + const testEnv = { + ...env, + DOCUMENT_WORKFLOW: { + create: vi.fn(), + createBatch: vi.fn(), + get: vi.fn().mockResolvedValue({ + sendEvent: vi.fn().mockRejectedValue(new Error("already complete")), + status: vi.fn().mockResolvedValue({ status: "complete" }), + }), + } as Cloudflare.Env["DOCUMENT_WORKFLOW"], + } + const providerOptions = await withProviderCompletion( + testEnv, + "job-3", + "execution-5", + "datalab", + undefined + ) + const callbackUrl = providerOptions?.webhook_url + if (typeof callbackUrl !== "string") { + throw new Error("Expected a Datalab completion URL.") + } + + const response = await api.fetch( + new Request(callbackUrl, { method: "POST" }), + testEnv + ) + expect(response.status).toBe(204) + }) +}) diff --git a/test/worker/retention.test.ts b/test/worker/retention.test.ts index bcdbbd1..2d28c19 100644 --- a/test/worker/retention.test.ts +++ b/test/worker/retention.test.ts @@ -221,6 +221,7 @@ function storedExecution(input: { return { id: input.id, jobId: input.jobId, + key: input.id, outputs: ["markdown" as const], position: 0, provider: "llamaparse" as const, From c091aa2ffc24a40a7e6f2fef97656e8295ac8850 Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 14:38:14 -0400 Subject: [PATCH 2/8] fix(sdk): preserve timeouts and page output semantics --- packages/filerouter/src/datalab.ts | 90 ++++++++++++++++--- packages/filerouter/src/llamaparse.ts | 3 +- packages/filerouter/src/mistral.ts | 19 +++- packages/filerouter/src/router.ts | 78 +++++++++------- packages/filerouter/src/testing.ts | 1 - packages/filerouter/src/types.ts | 3 +- packages/filerouter/test/datalab.test.ts | 86 ++++++++++++++++++ packages/filerouter/test/mistral.test.ts | 50 +++++++++++ packages/filerouter/test/router.test.ts | 13 ++- .../engines/pdf-inspector/server.ts | 3 +- src/lib/native-parser.server.ts | 5 +- 11 files changed, 289 insertions(+), 62 deletions(-) diff --git a/packages/filerouter/src/datalab.ts b/packages/filerouter/src/datalab.ts index 873ecd7..a55e19e 100644 --- a/packages/filerouter/src/datalab.ts +++ b/packages/filerouter/src/datalab.ts @@ -28,6 +28,7 @@ const OUTPUTS = [ "json", "markdown", "metadata", + "pages", ] satisfies Array export interface DatalabProviderOptions { @@ -61,6 +62,8 @@ export interface DatalabParseOptions { output_format?: Array | string page_range?: string paginate?: boolean + processing_location?: string + /** @deprecated Use `processing_location`. */ processing_region?: string raw?: Record save_checkpoint?: boolean @@ -81,20 +84,20 @@ interface DatalabSubmitResponse { export function datalab( options: DatalabProviderOptions = {} -): FileRouterProvider { +): FileRouterProvider { const jobs = datalabJobs(options) return { capabilities: { execution: "async", features: ["page-selection"], outputs: OUTPUTS, + pageFields: ["html", "json", "markdown", "metadata"], }, id: PROVIDER_ID, jobs, name: "Datalab", parse: (input, parseOptions) => parseDatalab(input, parseOptions, jobs, options.pollingIntervalMs), - raw: options, } } @@ -235,7 +238,13 @@ async function createFormData( ...raw, ...native, })) { - if (key !== "file" && key !== "file_url" && key !== "output_format") { + if ( + value !== undefined && + value !== null && + key !== "file" && + key !== "file_url" && + key !== "output_format" + ) { body.set(key, formValue(value)) } } @@ -254,6 +263,13 @@ async function createFormData( "output_format", nativeOutputFormats(nativeOptions.output_format, outputs).join(",") ) + if ( + outputs.includes("pages") && + (parseOptions.pageFields === undefined || + parseOptions.pageFields.includes("markdown")) + ) { + body.set("include_markdown_in_chunks", "true") + } if (parseOptions.pages) { body.set("page_range", parseOptions.pages.map((page) => page - 1).join(",")) } @@ -268,7 +284,7 @@ function normalizeDatalab( includeRaw: boolean, startedAt: Date ): ParseResult { - const pages = readRecords(raw.pages).map(normalizePage) + const pages = normalizePages(raw.json) const markdown = readString(raw.markdown) const html = readString(raw.html) const json = raw.json @@ -302,6 +318,7 @@ function normalizeDatalab( json, markdown, metadata, + pages, }), pageCount, provider: PROVIDER_ID, @@ -324,26 +341,71 @@ function normalizeDatalab( } } +function normalizePages(value: unknown): Array { + if (!isRecord(value)) { + return [] + } + return readRecords(value.children) + .filter((block) => block.block_type === "Page") + .map(normalizePage) +} + function normalizePage(raw: Record, index: number): ParsePage { - const html = readString(raw.html) - const markdown = readString(raw.markdown) - const text = readString(raw.text) + const children = readRecords(raw.children) + const html = readString(raw.html) ?? joinBlockField(children, "html") + const markdown = + readString(raw.markdown) ?? joinBlockField(children, "markdown") + const metadata = { + ...(raw.bbox !== undefined && { bbox: raw.bbox }), + ...(typeof raw.id === "string" && { blockId: raw.id }), + ...(raw.polygon !== undefined && { polygon: raw.polygon }), + ...(raw.section_hierarchy !== undefined && { + sectionHierarchy: raw.section_hierarchy, + }), + } return { ...(html && { html }), - ...(raw.json !== undefined && { json: raw.json }), + json: raw, ...(markdown && { markdown }), - ...(isRecord(raw.metadata) && { metadata: raw.metadata }), - pageNumber: readNumber(raw.page_number) ?? index + 1, - ...(text && { text }), + ...(Object.keys(metadata).length > 0 && { metadata }), + pageNumber: datalabPageNumber(raw, index), warnings: [], } } +function datalabPageNumber( + page: Record, + index: number +): number { + const explicit = readNumber(page.page_number) + if ( + explicit !== undefined && + Number.isSafeInteger(explicit) && + explicit > 0 + ) { + return explicit + } + const id = readString(page.id) + const zeroBased = id?.match(/(?:^|\/)page\/(\d+)(?:\/|$)/i)?.[1] + return zeroBased === undefined ? index + 1 : Number(zeroBased) + 1 +} + +function joinBlockField( + blocks: Array>, + field: "html" | "markdown" +): string | undefined { + const values = blocks.flatMap((block) => { + const value = readString(block[field]) + return value ? [value] : [] + }) + return values.length > 0 ? values.join("\n\n") : undefined +} + function datalabOutputs( outputs: Array ): Array { - const supported = new Set() + const supported = new Set() for (const output of outputs) { if ( output === "chunks" || @@ -353,14 +415,14 @@ function datalabOutputs( ) { supported.add(output) } - if (["images", "metadata"].includes(output)) { + if (output === "images" || output === "metadata" || output === "pages") { supported.add("json") } } if (supported.size === 0) { supported.add(DEFAULT_PARSE_OUTPUT) } - return [...supported] as Array + return [...supported] } function nativeOutputFormats( diff --git a/packages/filerouter/src/llamaparse.ts b/packages/filerouter/src/llamaparse.ts index 44a3237..50acf34 100644 --- a/packages/filerouter/src/llamaparse.ts +++ b/packages/filerouter/src/llamaparse.ts @@ -71,7 +71,7 @@ const LLAMAPARSE_OUTPUTS: Array = [ export const llamaparse = ( opts: LlamaParseProviderOptions = {} -): FileRouterProvider => { +): FileRouterProvider => { const jobs = llamaParseJobs(opts) return { capabilities: { @@ -83,7 +83,6 @@ export const llamaparse = ( id: "llamaparse", jobs, name: "LlamaParse", - raw: opts.client, parse: async (input, options) => { const job = await jobs.submit(input, options) return waitForProviderJob( diff --git a/packages/filerouter/src/mistral.ts b/packages/filerouter/src/mistral.ts index 2d01066..1b2cc9e 100644 --- a/packages/filerouter/src/mistral.ts +++ b/packages/filerouter/src/mistral.ts @@ -21,6 +21,7 @@ import type { ParseOptions, ParseOutput, ParsePage, + ParsePageField, ParseResult, ProviderInput, } from "./types" @@ -59,7 +60,7 @@ export type MistralOcrParseOptions = Partial> export function mistralOcr( options: MistralOcrProviderOptions = {} -): FileRouterProvider { +): FileRouterProvider { return { capabilities: { execution: "sync", @@ -80,7 +81,6 @@ export function mistralOcr( id: PROVIDER_ID, name: "Mistral OCR", parse: (input, parseOptions) => parseMistral(input, parseOptions, options), - raw: options.client, } } @@ -96,6 +96,8 @@ async function parseMistral( PROVIDER_ID ) const outputs = parseOptions.outputs ?? [DEFAULT_PARSE_OUTPUT] + const pageFieldRequested = (field: ParsePageField): boolean => + parseOptions.pageFields?.includes(field) === true const requestOptions = mistralRequestOptions(parseOptions) let uploadedFileId: string | undefined @@ -115,14 +117,23 @@ async function parseMistral( ...native, document, includeImageBase64: - outputs.includes("images") || native.includeImageBase64 === true, + outputs.includes("images") || + pageFieldRequested("images") || + native.includeImageBase64 === true, model: native.model ?? options.model ?? "mistral-ocr-4-0", ...(parseOptions.pages && { pages: parseOptions.pages.map((page) => page - 1), }), - ...(outputs.includes("tables") && { + ...((outputs.includes("tables") || pageFieldRequested("tables")) && { tableFormat: native.tableFormat ?? "markdown", }), + ...(pageFieldRequested("blocks") && { includeBlocks: true }), + ...(pageFieldRequested("confidence") && { + confidenceScoresGranularity: + native.confidenceScoresGranularity ?? "page", + }), + ...(pageFieldRequested("footer") && { extractFooter: true }), + ...(pageFieldRequested("header") && { extractHeader: true }), } const response = await client.ocr.process(request, requestOptions) diff --git a/packages/filerouter/src/router.ts b/packages/filerouter/src/router.ts index 460dff2..2840ce0 100644 --- a/packages/filerouter/src/router.ts +++ b/packages/filerouter/src/router.ts @@ -6,6 +6,7 @@ import { assertPages, assertTimeoutMs, } from "./internal/provider-options" +import { withTimeout } from "./internal/timeout" import { DEFAULT_PARSE_OUTPUT } from "./types" import type { CompareOptions, @@ -52,19 +53,28 @@ export class DirectFileRouter { assertPageFields(outputs, options.pageFields) assertProviderOutputs(provider, outputs) assertProviderPageFields(provider, options.pageFields) - const normalizedInput = await resolveParseInput(input) - try { - return selectPageFields( - await provider.parse(normalizedInput, { ...options, outputs }), - options.pageFields - ) - } catch (error) { - throw toFileRouterError(error, { - code: "ParseFailed", - providerId: provider.id, - }) + const run = async (signal: AbortSignal | undefined) => { + const normalizedInput = await resolveParseInput(input, signal) + try { + return selectPageFields( + await provider.parse(normalizedInput, { + ...options, + outputs, + ...(signal && { signal }), + }), + options.pageFields + ) + } catch (error) { + throw toFileRouterError(error, { + code: "ParseFailed", + providerId: provider.id, + }) + } } + return options.timeoutMs === undefined + ? run(options.signal) + : withTimeout(options.timeoutMs, options.signal, run) } async compare( @@ -77,27 +87,33 @@ export class DirectFileRouter { const outputs = options.outputs ?? [DEFAULT_PARSE_OUTPUT] assertPageFields(outputs, options.pageFields) const providerIds = options.providers ?? Object.keys(this.#providers) - const normalizedInput = await resolveParseInput(input) - const providers = await Promise.all( - providerIds.map((providerId) => - this.#compareProvider(providerId, normalizedInput, { - ...options, - outputs, - }) + const run = async (signal: AbortSignal | undefined) => { + const normalizedInput = await resolveParseInput(input, signal) + const providers = await Promise.all( + providerIds.map((providerId) => + this.#compareProvider(providerId, normalizedInput, { + ...options, + outputs, + ...(signal && { signal }), + }) + ) ) - ) - const completedAt = new Date() + const completedAt = new Date() - return { - input: describeInput(input), - outputs, - providers, - timing: { - completedAt: completedAt.toISOString(), - durationMs: completedAt.getTime() - startedAt.getTime(), - startedAt: startedAt.toISOString(), - }, + return { + input: describeInput(input), + outputs, + providers, + timing: { + completedAt: completedAt.toISOString(), + durationMs: completedAt.getTime() - startedAt.getTime(), + startedAt: startedAt.toISOString(), + }, + } } + return options.timeoutMs === undefined + ? run(options.signal) + : withTimeout(options.timeoutMs, options.signal, run) } async #compareProvider( @@ -187,10 +203,6 @@ export class DirectFileRouter { } } -export const createDirectFileRouter = ( - opts: DirectFileRouterOptions -): DirectFileRouter => new DirectFileRouter(opts) - export const assertProviderOutputs = ( provider: FileRouterProvider, outputs: Array diff --git a/packages/filerouter/src/testing.ts b/packages/filerouter/src/testing.ts index 676696d..cd3b3c7 100644 --- a/packages/filerouter/src/testing.ts +++ b/packages/filerouter/src/testing.ts @@ -34,7 +34,6 @@ export const fakeProvider = ( }, id, name: opts.name ?? "Fake Provider", - raw: opts, parse: ( _input: ProviderInput, options: ParseOptions diff --git a/packages/filerouter/src/types.ts b/packages/filerouter/src/types.ts index 47459ef..d0a204c 100644 --- a/packages/filerouter/src/types.ts +++ b/packages/filerouter/src/types.ts @@ -265,12 +265,11 @@ export interface ProviderJobs { ) => Promise } -export interface FileRouterProvider { +export interface FileRouterProvider { readonly capabilities: ProviderCapabilities readonly id: string readonly jobs?: ProviderJobs readonly name: string - readonly raw?: Raw parse: (input: ProviderInput, options: ParseOptions) => Promise } diff --git a/packages/filerouter/test/datalab.test.ts b/packages/filerouter/test/datalab.test.ts index b31b959..dcbb544 100644 --- a/packages/filerouter/test/datalab.test.ts +++ b/packages/filerouter/test/datalab.test.ts @@ -136,6 +136,92 @@ describe("Datalab provider", () => { }) }) + test("normalizes page-level JSON into portable pages", async () => { + const fetchMock = vi + .fn() + .mockResolvedValueOnce( + Response.json({ + request_check_url: "https://www.datalab.to/api/v1/convert/request-1", + request_id: "request-1", + success: true, + }) + ) + .mockResolvedValueOnce( + Response.json({ + json: { + children: [ + { + bbox: [0, 0, 100, 100], + block_type: "Page", + children: [ + { + html: "

First

", + markdown: "# First", + }, + ], + id: "/page/0/Page/0", + }, + { + block_type: "Page", + children: [ + { + html: "

Third

", + markdown: "Third", + }, + ], + id: "/page/2/Page/2", + }, + ], + }, + page_count: 2, + status: "complete", + success: true, + }) + ) + const router = new DirectFileRouter({ + providers: { + datalab: datalab({ + apiKey: "test-key", + fetch: fetchMock, + pollingIntervalMs: 0, + }), + }, + }) + + const result = await router.parse("https://example.com/report.pdf", { + outputs: ["pages"], + pageFields: ["html", "json", "markdown", "metadata"], + pages: [1, 3], + }) + + expect(result.outputs.pages).toEqual([ + { + html: "

First

", + json: expect.objectContaining({ id: "/page/0/Page/0" }), + markdown: "# First", + metadata: { + bbox: [0, 0, 100, 100], + blockId: "/page/0/Page/0", + }, + pageNumber: 1, + warnings: [], + }, + { + html: "

Third

", + json: expect.objectContaining({ id: "/page/2/Page/2" }), + markdown: "Third", + metadata: { blockId: "/page/2/Page/2" }, + pageNumber: 3, + warnings: [], + }, + ]) + const body = fetchMock.mock.calls[0]?.[1]?.body + expect(body).toBeInstanceOf(FormData) + expect((body as FormData).get("output_format")).toBe("json") + expect((body as FormData).get("include_markdown_in_chunks")).toBe("true") + expect((body as FormData).get("page_range")).toBe("0,2") + }) + test("rejects untrusted polling URLs", async () => { const fetchMock = vi.fn().mockResolvedValue( Response.json({ diff --git a/packages/filerouter/test/mistral.test.ts b/packages/filerouter/test/mistral.test.ts index bfc0fb9..1ac5ad1 100644 --- a/packages/filerouter/test/mistral.test.ts +++ b/packages/filerouter/test/mistral.test.ts @@ -161,4 +161,54 @@ describe("Mistral OCR provider", () => { expect.any(Object) ) }) + + test("requests optional provider data for selected page fields", async () => { + const process = vi.fn().mockResolvedValue({ + model: "mistral-ocr-latest", + pages: [ + { + dimensions: null, + images: [], + index: 0, + markdown: "# Page", + tables: [], + }, + ], + usageInfo: { docSizeBytes: 10, pagesProcessed: 1 }, + }) + const router = new DirectFileRouter({ + providers: { + mistral: mistralOcr({ + client: { + files: { delete: vi.fn(), upload: vi.fn() }, + ocr: { process }, + }, + }), + }, + }) + + await router.parse("https://example.com/report.pdf", { + outputs: ["pages"], + pageFields: [ + "blocks", + "confidence", + "footer", + "header", + "images", + "tables", + ], + }) + + expect(process).toHaveBeenCalledWith( + expect.objectContaining({ + confidenceScoresGranularity: "page", + extractFooter: true, + extractHeader: true, + includeBlocks: true, + includeImageBase64: true, + tableFormat: "markdown", + }), + expect.any(Object) + ) + }) }) diff --git a/packages/filerouter/test/router.test.ts b/packages/filerouter/test/router.test.ts index 63a5484..56fe913 100644 --- a/packages/filerouter/test/router.test.ts +++ b/packages/filerouter/test/router.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, test } from "vite-plus/test" +import { describe, expect, test, vi } from "vite-plus/test" import { DirectFileRouter, @@ -138,6 +138,17 @@ describe("DirectFileRouter", () => { ) }) + test("applies a direct timeout before starting provider work", async () => { + const provider = fakeProvider() + const parse = vi.spyOn(provider, "parse") + const router = new DirectFileRouter({ providers: { fake: provider } }) + + await expect( + router.parse(new Blob(["document"]), { timeoutMs: 0 }) + ).rejects.toMatchObject({ code: "Timeout" }) + expect(parse).not.toHaveBeenCalled() + }) + test("serializes DirectFileRouter errors from another package copy", () => { const error = Object.assign(new Error("rate limited"), { code: "RateLimit", diff --git a/services/native-parsers/engines/pdf-inspector/server.ts b/services/native-parsers/engines/pdf-inspector/server.ts index 6af8c03..2b747aa 100644 --- a/services/native-parsers/engines/pdf-inspector/server.ts +++ b/services/native-parsers/engines/pdf-inspector/server.ts @@ -79,8 +79,9 @@ function parsePdf(bytes: Buffer, value: unknown): NativeParserResult { pagesWithColumns: extraction.pagesWithColumns, pagesWithTables: extraction.pagesWithTables, pdfType: classification.pdfType, + sourcePageCount: classification.pageCount, }, - pageCount: classification.pageCount, + pageCount: pages.length, pages, warnings, }, diff --git a/src/lib/native-parser.server.ts b/src/lib/native-parser.server.ts index 2e16b33..66f9fe2 100644 --- a/src/lib/native-parser.server.ts +++ b/src/lib/native-parser.server.ts @@ -94,9 +94,6 @@ async function parseNative( } const startedAt = new Date() - const supportedPageFields = new Set(config.capabilities.pageFields) - const selectedNativePageFields = - options.pageFields?.filter((field) => supportedPageFields.has(field)) ?? [] const response = await config.fetch( new Request(`https://native-parsers.internal/v1/${config.id}/parse`, { headers: { @@ -106,7 +103,7 @@ async function parseNative( options.outputs && { outputs: options.outputs }), ...(!options.includeRaw && options.pageFields && { - pageFields: selectedNativePageFields, + pageFields: options.pageFields, }), ...(options.pages && { pages: options.pages }), ...(options.providerOptions?.[config.id] && { From cc5f0c9a70075fb6a714da691d10a9054af609ea Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 14:38:29 -0400 Subject: [PATCH 3/8] docs: document execution lifecycle primitives --- docs/api/jobs.mdx | 29 ++++++++++++++++++----------- docs/concepts/how-it-works.mdx | 13 +++++++++---- docs/concepts/results.mdx | 2 ++ docs/guides/cheap-path.mdx | 21 +++++++++++++++++---- docs/guides/durable-jobs.mdx | 18 ++++++++++++------ docs/guides/fast-then-strong.mdx | 15 +++++++++++---- docs/guides/overview.mdx | 17 +++++++++-------- 7 files changed, 78 insertions(+), 37 deletions(-) diff --git a/docs/api/jobs.mdx b/docs/api/jobs.mdx index b532c3e..3b173f0 100644 --- a/docs/api/jobs.mdx +++ b/docs/api/jobs.mdx @@ -33,19 +33,19 @@ curl https://filerouter.dev/api/v1/documents \ The response contains an immutable document ID that can be reused across jobs until its `expiresAt` time, seven days after upload. -Delete the document, its jobs, and retained results before then when you no -longer need them: +Release the source and retained results while preserving job history when you +no longer need the artifacts: ```bash -curl https://filerouter.dev/api/v1/documents/DOCUMENT_ID \ - --request DELETE \ +curl https://filerouter.dev/api/v1/documents/DOCUMENT_ID/release \ + --request POST \ --header "Authorization: Bearer $FILEROUTER_API_KEY" ``` ## Create a job Each provider entry is an independent execution. Provider-specific options stay -on that target instead of leaking into the rest of the job. +on that entry instead of leaking into the rest of the job. ```bash curl https://filerouter.dev/api/v1/jobs \ @@ -56,8 +56,13 @@ curl https://filerouter.dev/api/v1/jobs \ --data '{ "documentId": "550e8400-e29b-41d4-a716-446655440000", "providers": [ - { "provider": "llamaparse", "outputs": ["markdown", "tables"] }, { + "key": "primary", + "provider": "llamaparse", + "outputs": ["markdown", "tables"] + }, + { + "key": "inspection", "provider": "liteparse", "outputs": ["pages"], "pageFields": ["markdown"], @@ -67,7 +72,9 @@ curl https://filerouter.dev/api/v1/jobs \ }' ``` -If a target omits `outputs`, FileRouter requests `markdown`. `pageFields` +`key` is a caller-defined identifier unique within the job. It lets one +provider appear more than once with different options. If an entry omits +`outputs`, FileRouter requests `markdown`. `pageFields` requires the `pages` output and limits the fields retained on each page. `pageNumber` and `warnings` are always retained. @@ -83,7 +90,7 @@ curl https://filerouter.dev/api/v1/jobs/550e8400-e29b-41d4-a716-446655440000 \ ``` Jobs move through `queued`, `running`, and then `complete` or `failed`. The job -contains one execution per provider, including its status, duration, usage, and +contains one execution per provider entry, including its key, provider, status, duration, usage, and result availability. Results stay separate so one slow or failed provider does not hide successful work. @@ -95,7 +102,7 @@ curl https://filerouter.dev/api/v1/executions/EXECUTION_ID/result \ The TypeScript SDK's `parse` and `compare` methods perform these steps for you. Use `router.documents`, `router.jobs`, and `router.executions` when your app needs to persist IDs or compose its own workflow. Use -`router.jobs.waitForExecution(job, provider)` to consume one provider without +`router.jobs.waitForExecution(job, execution)` to consume one execution without waiting for the complete job. Call -`router.documents.delete(documentId)` to remove a document and its retained job -data immediately. +`router.documents.release(documentId)` to remove stored artifacts or +`router.documents.delete(documentId)` to remove the document and its job history. diff --git a/docs/concepts/how-it-works.mdx b/docs/concepts/how-it-works.mdx index 5defaf8..253f4e9 100644 --- a/docs/concepts/how-it-works.mdx +++ b/docs/concepts/how-it-works.mdx @@ -50,12 +50,14 @@ try { documentId: document.id, providers: [ { + key: "initial", provider: "liteparse", outputs: ["markdown", "pages"], pageFields: ["markdown"], providerOptions: { ocr: "auto" }, }, { + key: "secondary", provider: "llamaparse", outputs: ["markdown", "tables"], providerOptions: { tier: "agentic" }, @@ -63,20 +65,23 @@ try { ], }) - const execution = await router.jobs.waitForExecution(job, "liteparse") + const executionRef = job.executions.find((item) => item.key === "initial") + if (!executionRef) throw new Error("Initial execution was not created.") + const execution = await router.jobs.waitForExecution(job, executionRef) if (execution.status === "complete" && execution.resultAvailable) { const result = await router.executions.result(execution.id) console.log(result.outputs.markdown) } } finally { - await router.documents.delete(document.id) + await router.documents.release(document.id) } ``` For step-by-step recipes, see [Guides](/guides/overview). -`waitForExecution` lets an application consume one provider as soon as it -finishes without waiting for the other executions in the job. +`waitForExecution` lets an application consume one execution as soon as it +finishes without waiting for the others. Provider keys are caller-defined labels; +FileRouter does not assign roles or decide which result wins. ## Built-in providers diff --git a/docs/concepts/results.mdx b/docs/concepts/results.mdx index dbd1a61..0f31a21 100644 --- a/docs/concepts/results.mdx +++ b/docs/concepts/results.mdx @@ -43,6 +43,8 @@ interface ParseResult { ``` Only requested and supported values appear under `outputs`. +`pageCount` is the number of pages processed in this result; when a provider +exposes the source document's total separately, it appears in `metadata`. ## Page fields diff --git a/docs/guides/cheap-path.mdx b/docs/guides/cheap-path.mdx index a35e08b..d8a6309 100644 --- a/docs/guides/cheap-path.mdx +++ b/docs/guides/cheap-path.mdx @@ -63,6 +63,7 @@ async function parseClaim(claimId: string, filePath: string) { documentId: document.id, providers: [ { + key: "initial", outputs: ["markdown", "pages"], provider: "liteparse", pageFields: ["markdown"], @@ -73,9 +74,13 @@ async function parseClaim(claimId: string, filePath: string) { { idempotencyKey: `job:${claimId}:liteparse` } ) + const initialRef = cheapAccepted.executions.find( + (item) => item.key === "initial" + ) + if (!initialRef) throw new Error("Initial execution was not created.") const cheapExec = await router.jobs.waitForExecution( cheapAccepted, - "liteparse" + initialRef ) if (cheapExec.status === "complete" && cheapExec.resultAvailable) { const cheap = await router.executions.result(cheapExec.id) @@ -88,22 +93,30 @@ async function parseClaim(claimId: string, filePath: string) { { documentId: document.id, providers: [ - { outputs: ["markdown", "tables"], provider: "llamaparse" }, + { + key: "secondary", + outputs: ["markdown", "tables"], + provider: "llamaparse", + }, ], }, { idempotencyKey: `job:${claimId}:llamaparse` } ) + const secondaryRef = strongAccepted.executions.find( + (item) => item.key === "secondary" + ) + if (!secondaryRef) throw new Error("Secondary execution was not created.") const strongExec = await router.jobs.waitForExecution( strongAccepted, - "llamaparse" + secondaryRef ) if (strongExec.status !== "complete" || !strongExec.resultAvailable) { throw new Error("Strong parse did not complete.") } return await router.executions.result(strongExec.id) } finally { - await router.documents.delete(document.id) + await router.documents.release(document.id) } } ``` diff --git a/docs/guides/durable-jobs.mdx b/docs/guides/durable-jobs.mdx index f47c7d4..3b9b811 100644 --- a/docs/guides/durable-jobs.mdx +++ b/docs/guides/durable-jobs.mdx @@ -25,7 +25,9 @@ const document = await router.documents.create("./claim.pdf", { const accepted = await router.jobs.create( { documentId: document.id, - providers: [{ outputs: ["markdown"], provider: "liteparse" }], + providers: [ + { key: "primary", outputs: ["markdown"], provider: "liteparse" }, + ], }, { idempotencyKey: `job:${claimId}:liteparse` } ) @@ -65,13 +67,15 @@ try { const accepted = await router.jobs.create( { documentId: document.id, - providers: [{ outputs: ["markdown"], provider: "liteparse" }], + providers: [ + { key: "primary", outputs: ["markdown"], provider: "liteparse" }, + ], }, { idempotencyKey: `job:${claimId}:liteparse` } ) const job = await router.jobs.wait(accepted) - const execution = job.executions.find((item) => item.provider === "liteparse") + const execution = job.executions.find((item) => item.key === "primary") if (execution?.status === "complete" && execution.resultAvailable) { const result = await router.executions.result(execution.id) @@ -79,11 +83,13 @@ try { console.log(result.pageCount, result.timing.durationMs) } } finally { - await router.documents.delete(document.id) + await router.documents.release(document.id) } ``` -`documents.delete` removes the source, related jobs, and retained results. Prefer it at the end of ingest handlers and one-shot tools. +`documents.release` removes the source and retained results while preserving job and execution records. It refuses while a job is active. + +`documents.delete` removes the document, related jobs, and retained results. Use it when you do not need the history. For hosted `parse` / `compare`, delete via the returned resources when you no longer need retained outputs: @@ -92,7 +98,7 @@ const parsed = await router.parse("./claim.pdf", { provider: "liteparse", outputs: ["markdown"], }) -await router.documents.delete(parsed.resources.documentId) +await router.documents.release(parsed.resources.documentId) ``` ## Partial failure diff --git a/docs/guides/fast-then-strong.mdx b/docs/guides/fast-then-strong.mdx index ab34949..30c8a8d 100644 --- a/docs/guides/fast-then-strong.mdx +++ b/docs/guides/fast-then-strong.mdx @@ -30,33 +30,39 @@ try { documentId: document.id, providers: [ { + key: "preview", outputs: ["markdown", "pages"], provider: "liteparse", pageFields: ["markdown"], providerOptions: { ocr: "auto" }, }, { + key: "full", outputs: ["markdown", "pages"], provider: "llamaparse", }, ], }) - const preview = await router.jobs.waitForExecution(job, "liteparse") + const previewRef = job.executions.find((item) => item.key === "preview") + if (!previewRef) throw new Error("Preview execution was not created.") + const preview = await router.jobs.waitForExecution(job, previewRef) if (preview.status === "complete" && preview.resultAvailable) { const fast = await router.executions.result(preview.id) // Render draft UI from fast.outputs.markdown console.log("preview", fast.outputs.markdown?.slice(0, 200)) } - const strong = await router.jobs.waitForExecution(job, "llamaparse") + const fullRef = job.executions.find((item) => item.key === "full") + if (!fullRef) throw new Error("Full execution was not created.") + const strong = await router.jobs.waitForExecution(job, fullRef) if (strong.status === "complete" && strong.resultAvailable) { const full = await router.executions.result(strong.id) // Replace draft UI with full.outputs.markdown console.log("final", full.timing.durationMs, full.pageCount) } } finally { - await router.documents.delete(document.id) + await router.documents.release(document.id) } ``` @@ -64,7 +70,8 @@ try { - One stored document. Both engines see the same bytes. - Executions are independent. A slow or failed LlamaParse run does not erase a finished LiteParse result. -- `waitForExecution` returns when that provider hits `complete` or `failed`. Sibling work can keep running. +- `waitForExecution` returns when that execution hits `complete` or `failed`. Sibling work can keep running. +- `preview` and `full` are caller-defined keys for this application; FileRouter does not assign those roles. ## Output rules diff --git a/docs/guides/overview.mdx b/docs/guides/overview.mdx index fdbe2d7..eeb3a80 100644 --- a/docs/guides/overview.mdx +++ b/docs/guides/overview.mdx @@ -37,13 +37,14 @@ They are application-owned patterns. FileRouter does not auto-pick engines or fa ## Primitives cheat sheet -| Piece | Role | -| ----------------------------- | -------------------------------------------- | -| `documents.create` / `delete` | Store once; remove source + retained results | -| `jobs.create` | Start one or more provider executions | -| `jobs.waitForExecution` | Unblock when a single provider finishes | -| `jobs.wait` | Wait until the whole job ends | -| `executions.result` | Fetch a normalized `ParseResult` | -| `parse` / `compare` | Convenience wrappers over the same resources | +| Piece | Role | +| ------------------------------ | --------------------------------------------- | +| `documents.create` / `release` | Store once; release source + retained results | +| `documents.delete` | Remove the document and its job history | +| `jobs.create` | Start one or more provider executions | +| `jobs.waitForExecution` | Unblock when a single execution finishes | +| `jobs.wait` | Wait until the whole job ends | +| `executions.result` | Fetch a normalized `ParseResult` | +| `parse` / `compare` | Convenience wrappers over the same resources | See [How FileRouter works](/concepts/how-it-works) for the contract, and [Documents and jobs](/api/jobs) for the HTTP shape. From f7357cdc9cd534363fb00d13df8d0dd67067e0a6 Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 15:30:03 -0400 Subject: [PATCH 4/8] fix(hosted): harden release and callback handling --- src/api/app.ts | 4 ++ src/lib/document-deletion.server.ts | 63 ++++++++++++++--------------- src/lib/hmac.server.ts | 2 +- src/lib/r2-objects.server.ts | 4 +- test/worker/api.test.ts | 39 ++++++++++++++++++ 5 files changed, 77 insertions(+), 35 deletions(-) diff --git a/src/api/app.ts b/src/api/app.ts index b010e9f..6056255 100644 --- a/src/api/app.ts +++ b/src/api/app.ts @@ -219,6 +219,10 @@ api.post( "/api/v1/provider-completions/:jobId/:executionId", async (context) => { const { executionId, jobId } = context.req.param() + Object.assign(context.get("requestEvent"), { + execution_id: executionId, + job_id: jobId, + }) return receiveProviderCompletion( context.req.raw, context.env, diff --git a/src/lib/document-deletion.server.ts b/src/lib/document-deletion.server.ts index f68586a..c32da50 100644 --- a/src/lib/document-deletion.server.ts +++ b/src/lib/document-deletion.server.ts @@ -1,4 +1,4 @@ -import { and, eq, inArray, notExists } from "drizzle-orm" +import { and, eq, inArray, isNotNull, notExists } from "drizzle-orm" import { document, documentExecution, documentJob } from "@/db/schema" import { createDb } from "@/db/server" @@ -49,37 +49,32 @@ export async function releaseDocumentArtifacts( }) } - const executions = await db - .select({ - id: documentExecution.id, - resultKey: documentExecution.resultKey, - }) - .from(documentExecution) - .innerJoin(documentJob, eq(documentExecution.jobId, documentJob.id)) + const documentJobIds = db + .select({ id: documentJob.id }) + .from(documentJob) .where(eq(documentJob.documentId, id)) + const resultKeys = await db + .select({ key: documentExecution.resultKey }) + .from(documentExecution) + .where(inArray(documentExecution.jobId, documentJobIds)) .all() await deleteR2Objects(env.FILEROUTER_FILES, [ releasable.objectKey, - ...executions.map((execution) => execution.resultKey), + ...resultKeys.map((result) => result.key), ]) - const releaseDocument = db - .update(document) - .set({ objectKey: null, updatedAt: now }) - .where(and(eq(document.id, id), eq(document.userId, userId))) - if (executions.length === 0) { - await releaseDocument - return - } await db.batch([ - releaseDocument, + db + .update(document) + .set({ objectKey: null, updatedAt: now }) + .where(and(eq(document.id, id), eq(document.userId, userId))), db .update(documentExecution) .set({ resultExpiresAt: now, resultKey: null, updatedAt: now }) .where( - inArray( - documentExecution.id, - executions.map((execution) => execution.id) + and( + isNotNull(documentExecution.resultKey), + inArray(documentExecution.jobId, documentJobIds) ) ), ]) @@ -113,17 +108,21 @@ export async function deleteDocument( .filter((job) => job.status === "queued" || job.status === "running") .map((job) => terminateWorkflow(env.DOCUMENT_WORKFLOW, job.id)) ) - const jobIds = jobs.map((job) => job.id) - const resultKeys = - jobIds.length === 0 - ? [] - : ( - await db - .select({ key: documentExecution.resultKey }) - .from(documentExecution) - .where(inArray(documentExecution.jobId, jobIds)) - .all() - ).map((execution) => execution.key) + const resultKeys = ( + await db + .select({ key: documentExecution.resultKey }) + .from(documentExecution) + .where( + inArray( + documentExecution.jobId, + db + .select({ id: documentJob.id }) + .from(documentJob) + .where(eq(documentJob.documentId, id)) + ) + ) + .all() + ).map((execution) => execution.key) await deleteR2Objects(env.FILEROUTER_FILES, [stored.objectKey, ...resultKeys]) await db .delete(document) diff --git a/src/lib/hmac.server.ts b/src/lib/hmac.server.ts index 4c248eb..88f39d0 100644 --- a/src/lib/hmac.server.ts +++ b/src/lib/hmac.server.ts @@ -20,7 +20,7 @@ export async function verifyHmac( signature: string ): Promise { try { - return crypto.subtle.verify( + return await crypto.subtle.verify( "HMAC", await signingKey(secret, purpose, ["verify"]), fromBase64Url(signature).buffer as ArrayBuffer, diff --git a/src/lib/r2-objects.server.ts b/src/lib/r2-objects.server.ts index afcd86a..149901c 100644 --- a/src/lib/r2-objects.server.ts +++ b/src/lib/r2-objects.server.ts @@ -3,7 +3,7 @@ export async function deleteR2Objects( keys: Array ): Promise { const uniqueKeys = [...new Set(keys.filter((key): key is string => !!key))] - if (uniqueKeys.length > 0) { - await bucket.delete(uniqueKeys) + for (let index = 0; index < uniqueKeys.length; index += 1000) { + await bucket.delete(uniqueKeys.slice(index, index + 1000)) } } diff --git a/test/worker/api.test.ts b/test/worker/api.test.ts index 7d3aeb9..614f73a 100644 --- a/test/worker/api.test.ts +++ b/test/worker/api.test.ts @@ -202,6 +202,12 @@ describe("FileRouter Worker", () => { testEnv ) expect(activeRelease.status).toBe(409) + await expect(activeRelease.json()).resolves.toMatchObject({ + code: "document_active", + }) + expect( + await testEnv.FILEROUTER_FILES.head(`documents/${document.id}/source`) + ).not.toBeNull() const deleted = await api.fetch( new Request(`https://filerouter.test/api/v1/documents/${document.id}`, { @@ -364,6 +370,7 @@ describe("FileRouter Worker", () => { const resultKey = `executions/${executionId}/result.json` await env.FILEROUTER_FILES.put(resultKey, "{}") const now = Math.floor(Date.now() / 1_000) + const historicalJobId = crypto.randomUUID() await env.DB.batch([ env.DB.prepare( "UPDATE document_job SET status = 'complete' WHERE id = ?" @@ -371,6 +378,31 @@ describe("FileRouter Worker", () => { env.DB.prepare( "UPDATE document_execution SET status = 'complete', duration_ms = 10, page_count = 1, result_key = ?, result_expires_at = ?, completed_at = ? WHERE id = ?" ).bind(resultKey, now + 3_600, now, executionId), + env.DB.prepare( + "INSERT INTO document_job (id, user_id, document_id, status, idempotency_key_hash, request_hash, metering_status, created_at, updated_at) VALUES (?, ?, ?, 'complete', ?, ?, 'skipped', ?, ?)" + ).bind( + historicalJobId, + userId, + storedDocument.id, + `historical-${historicalJobId}`, + historicalJobId, + now, + now + ), + ...Array.from({ length: 101 }, (_, index) => + env.DB.prepare( + "INSERT INTO document_execution (id, job_id, key, provider, position, status, outputs, error_message, created_at, updated_at, completed_at) VALUES (?, ?, ?, 'liteparse', ?, 'failed', ?, 'historical failure', ?, ?, ?)" + ).bind( + `historical-execution-${index}`, + historicalJobId, + `historical-${index}`, + index, + JSON.stringify(["markdown"]), + now, + now, + now + ) + ), ]) const released = await api.fetch( @@ -389,6 +421,13 @@ describe("FileRouter Worker", () => { await env.FILEROUTER_FILES.head(`documents/${storedDocument.id}/source`) ).toBeNull() expect(await env.FILEROUTER_FILES.head(resultKey)).toBeNull() + await expect( + env.DB.prepare( + "SELECT result_expires_at FROM document_execution WHERE id = ?" + ) + .bind("historical-execution-0") + .first<{ result_expires_at: number | null }>() + ).resolves.toEqual({ result_expires_at: null }) const preservedJob = await api.fetch( new Request(`https://filerouter.test/api/v1/jobs/${accepted.id}`, { headers: { Authorization: `Bearer ${apiKey}` }, From 410a08f520da7238a2cf65261128658170c73b17 Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 15:30:18 -0400 Subject: [PATCH 5/8] fix(sdk): tighten Datalab pages and job validation --- packages/filerouter/src/datalab.ts | 87 ++++++++++++++++++------ packages/filerouter/src/jobs.ts | 22 +++--- packages/filerouter/src/router.ts | 2 + packages/filerouter/test/datalab.test.ts | 59 +++++++++++++++- packages/filerouter/test/jobs.test.ts | 11 +++ 5 files changed, 149 insertions(+), 32 deletions(-) diff --git a/packages/filerouter/src/datalab.ts b/packages/filerouter/src/datalab.ts index a55e19e..7f7ad90 100644 --- a/packages/filerouter/src/datalab.ts +++ b/packages/filerouter/src/datalab.ts @@ -232,12 +232,13 @@ async function createFormData( PROVIDER_ID ) const { raw, ...native } = nativeOptions - - for (const [key, value] of Object.entries({ + const nativeFields = { ...options.raw, ...raw, ...native, - })) { + } + + for (const [key, value] of Object.entries(nativeFields)) { if ( value !== undefined && value !== null && @@ -266,7 +267,8 @@ async function createFormData( if ( outputs.includes("pages") && (parseOptions.pageFields === undefined || - parseOptions.pageFields.includes("markdown")) + parseOptions.pageFields.includes("markdown")) && + nativeFields.include_markdown_in_chunks === undefined ) { body.set("include_markdown_in_chunks", "true") } @@ -284,13 +286,16 @@ function normalizeDatalab( includeRaw: boolean, startedAt: Date ): ParseResult { - const pages = normalizePages(raw.json) + const pageBlocks = datalabPageBlocks(raw.json) + const pages = requestedOutputs.includes("pages") + ? pageBlocks.map(normalizePage) + : [] const markdown = readString(raw.markdown) const html = readString(raw.html) const json = raw.json const chunks = raw.chunks const images = normalizeImages(raw.images) - const pageCount = readNumber(raw.page_count) ?? pages.length + const pageCount = readNumber(raw.page_count) ?? pageBlocks.length const qualityScore = readNumber(raw.parse_quality_score) const costBreakdown = isRecord(raw.cost_breakdown) ? raw.cost_breakdown @@ -341,20 +346,19 @@ function normalizeDatalab( } } -function normalizePages(value: unknown): Array { +function datalabPageBlocks(value: unknown): Array> { if (!isRecord(value)) { return [] } - return readRecords(value.children) - .filter((block) => block.block_type === "Page") - .map(normalizePage) + return readRecords(value.children).filter( + (block) => block.block_type === "Page" + ) } function normalizePage(raw: Record, index: number): ParsePage { - const children = readRecords(raw.children) - const html = readString(raw.html) ?? joinBlockField(children, "html") - const markdown = - readString(raw.markdown) ?? joinBlockField(children, "markdown") + const blocks = indexDatalabBlocks(raw) + const html = resolveBlockField(raw, "html", blocks) + const markdown = resolveBlockField(raw, "markdown", blocks) const metadata = { ...(raw.bbox !== undefined && { bbox: raw.bbox }), ...(typeof raw.id === "string" && { blockId: raw.id }), @@ -391,15 +395,56 @@ function datalabPageNumber( return zeroBased === undefined ? index + 1 : Number(zeroBased) + 1 } -function joinBlockField( - blocks: Array>, - field: "html" | "markdown" +function indexDatalabBlocks( + root: Record +): Map> { + const blocks = new Map>() + const pending = [root] + while (pending.length > 0) { + const block = pending.pop() + if (!block) { + continue + } + const id = readString(block.id) + if (id) { + blocks.set(id, block) + } + pending.push(...readRecords(block.children)) + } + return blocks +} + +function resolveBlockField( + root: Record, + field: "html" | "markdown", + blocks: Map> ): string | undefined { - const values = blocks.flatMap((block) => { + const rootId = readString(root.id) + const visiting = new Set(rootId ? [rootId] : []) + const resolve = (block: Record): string | undefined => { const value = readString(block[field]) - return value ? [value] : [] - }) - return values.length > 0 ? values.join("\n\n") : undefined + if (!value) { + const children = readRecords(block.children).flatMap((child) => { + const childValue = resolve(child) + return childValue ? [childValue] : [] + }) + return children.length > 0 ? children.join("\n\n") : undefined + } + return value.replace( + /]*\bsrc=(["'])([^"']+)\1[^>]*>\s*<\/content-ref>/gi, + (reference, _quote: string, childId: string) => { + const child = blocks.get(childId) + if (!child || visiting.has(childId)) { + return reference + } + visiting.add(childId) + const resolved = resolve(child) + visiting.delete(childId) + return resolved ?? reference + } + ) + } + return resolve(root) } function datalabOutputs( diff --git a/packages/filerouter/src/jobs.ts b/packages/filerouter/src/jobs.ts index 6d60eca..03e574c 100644 --- a/packages/filerouter/src/jobs.ts +++ b/packages/filerouter/src/jobs.ts @@ -224,14 +224,6 @@ function serializeJobInput(input: HostedJobCreateInput): string { { code: "InvalidInput" } ) } - if ( - new Set(input.providers.map((provider) => provider.key)).size !== - input.providers.length - ) { - throw new FileRouterError("Each provider key must be unique.", { - code: "InvalidInput", - }) - } if ( input.metadata && Object.keys(input.metadata).length > MAX_HOSTED_METADATA_ENTRIES @@ -242,7 +234,11 @@ function serializeJobInput(input: HostedJobCreateInput): string { ) } for (const target of input.providers) { - if (target.key.trim().length === 0 || target.key.length > 64) { + if ( + typeof target.key !== "string" || + target.key.trim().length === 0 || + target.key.length > 64 + ) { throw new FileRouterError( "Provider keys must be non-blank and at most 64 characters.", { code: "InvalidInput" } @@ -253,6 +249,14 @@ function serializeJobInput(input: HostedJobCreateInput): string { target.pageFields ) } + if ( + new Set(input.providers.map((provider) => provider.key)).size !== + input.providers.length + ) { + throw new FileRouterError("Each provider key must be unique.", { + code: "InvalidInput", + }) + } return stringifyJson(input) } diff --git a/packages/filerouter/src/router.ts b/packages/filerouter/src/router.ts index 2840ce0..3a820e5 100644 --- a/packages/filerouter/src/router.ts +++ b/packages/filerouter/src/router.ts @@ -56,6 +56,7 @@ export class DirectFileRouter { const run = async (signal: AbortSignal | undefined) => { const normalizedInput = await resolveParseInput(input, signal) + signal?.throwIfAborted() try { return selectPageFields( await provider.parse(normalizedInput, { @@ -89,6 +90,7 @@ export class DirectFileRouter { const providerIds = options.providers ?? Object.keys(this.#providers) const run = async (signal: AbortSignal | undefined) => { const normalizedInput = await resolveParseInput(input, signal) + signal?.throwIfAborted() const providers = await Promise.all( providerIds.map((providerId) => this.#compareProvider(providerId, normalizedInput, { diff --git a/packages/filerouter/test/datalab.test.ts b/packages/filerouter/test/datalab.test.ts index dcbb544..1e00d19 100644 --- a/packages/filerouter/test/datalab.test.ts +++ b/packages/filerouter/test/datalab.test.ts @@ -155,11 +155,25 @@ describe("Datalab provider", () => { block_type: "Page", children: [ { - html: "

First

", - markdown: "# First", + block_type: "SectionHeader", + children: [ + { + block_type: "Text", + html: "First", + id: "/page/0/Text/2", + markdown: "First", + }, + ], + html: "

", + id: "/page/0/SectionHeader/1", + markdown: + "# ", }, ], + html: "", id: "/page/0/Page/0", + markdown: + "", }, { block_type: "Page", @@ -222,6 +236,47 @@ describe("Datalab provider", () => { expect((body as FormData).get("page_range")).toBe("0,2") }) + test("preserves an explicit markdown-in-JSON opt-out", async () => { + const fetchMock = vi + .fn() + .mockResolvedValueOnce( + Response.json({ + request_check_url: "https://www.datalab.to/api/v1/convert/request-2", + request_id: "request-2", + success: true, + }) + ) + .mockResolvedValueOnce( + Response.json({ + json: { children: [] }, + page_count: 0, + status: "complete", + success: true, + }) + ) + const router = new DirectFileRouter({ + providers: { + datalab: datalab({ + apiKey: "test-key", + fetch: fetchMock, + pollingIntervalMs: 0, + }), + }, + }) + + await router.parse("https://example.com/report.pdf", { + outputs: ["pages"], + pageFields: ["markdown"], + providerOptions: { + datalab: { include_markdown_in_chunks: false }, + }, + }) + + const body = fetchMock.mock.calls[0]?.[1]?.body + expect(body).toBeInstanceOf(FormData) + expect((body as FormData).get("include_markdown_in_chunks")).toBe("false") + }) + test("rejects untrusted polling URLs", async () => { const fetchMock = vi.fn().mockResolvedValue( Response.json({ diff --git a/packages/filerouter/test/jobs.test.ts b/packages/filerouter/test/jobs.test.ts index 0d1c3d2..8fc96ec 100644 --- a/packages/filerouter/test/jobs.test.ts +++ b/packages/filerouter/test/jobs.test.ts @@ -195,6 +195,17 @@ describe("hosted resources", () => { providers: [{ key: " ", provider: "llamaparse" }], }) ).rejects.toMatchObject({ code: "InvalidInput" }) + await expect( + jobs.create({ + ...base, + providers: [ + { + key: undefined as unknown as string, + provider: "llamaparse", + }, + ], + }) + ).rejects.toMatchObject({ code: "InvalidInput" }) await expect( jobs.create({ ...base, From dd9e6bb448419b31d2def8f5bbce2722b24538c9 Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 15:30:31 -0400 Subject: [PATCH 6/8] docs: clarify artifact release preconditions --- docs/api/jobs.mdx | 9 +++++---- docs/concepts/how-it-works.mdx | 3 +++ docs/guides/cheap-path.mdx | 2 +- docs/guides/durable-jobs.mdx | 6 +++--- docs/guides/fast-then-strong.mdx | 2 +- 5 files changed, 13 insertions(+), 9 deletions(-) diff --git a/docs/api/jobs.mdx b/docs/api/jobs.mdx index 3b173f0..273d6d2 100644 --- a/docs/api/jobs.mdx +++ b/docs/api/jobs.mdx @@ -34,7 +34,8 @@ The response contains an immutable document ID that can be reused across jobs until its `expiresAt` time, seven days after upload. Release the source and retained results while preserving job history when you -no longer need the artifacts: +no longer need the artifacts. Every job for the document must be terminal; +release returns `409 document_active` while any job is queued or running. ```bash curl https://filerouter.dev/api/v1/documents/DOCUMENT_ID/release \ @@ -78,9 +79,9 @@ provider appear more than once with different options. If an entry omits requires the `pages` output and limits the fields retained on each page. `pageNumber` and `warnings` are always retained. -A new job returns `202 Accepted`. Reusing the key with the same request returns -the existing job and `Idempotent-Replayed: true`; changing the request returns a -conflict. +A new job returns `202 Accepted`. Reusing the same `Idempotency-Key` header with +the same request returns the existing job and `Idempotent-Replayed: true`; +changing the request with that header returns a conflict. ## Poll and retrieve results diff --git a/docs/concepts/how-it-works.mdx b/docs/concepts/how-it-works.mdx index 253f4e9..8670ccc 100644 --- a/docs/concepts/how-it-works.mdx +++ b/docs/concepts/how-it-works.mdx @@ -45,6 +45,7 @@ const document = await router.documents.create(stream, { mimeType: "application/pdf", }) +let jobId: string | undefined try { const job = await router.jobs.create({ documentId: document.id, @@ -64,6 +65,7 @@ try { }, ], }) + jobId = job.id const executionRef = job.executions.find((item) => item.key === "initial") if (!executionRef) throw new Error("Initial execution was not created.") @@ -73,6 +75,7 @@ try { console.log(result.outputs.markdown) } } finally { + if (jobId) await router.jobs.wait(jobId) await router.documents.release(document.id) } ``` diff --git a/docs/guides/cheap-path.mdx b/docs/guides/cheap-path.mdx index d8a6309..328f46d 100644 --- a/docs/guides/cheap-path.mdx +++ b/docs/guides/cheap-path.mdx @@ -128,7 +128,7 @@ async function parseClaim(claimId: string, filePath: string) { - Triage with `liteparse` or `pdf-inspector` - Escalate with a second job on the same document, not a second blind `parse()` that re-uploads - Request only outputs each engine supports -- Delete documents when the pipeline finishes +- Release stored artifacts when the pipeline finishes - Use [Bake off engines](/guides/bake-off) offline to learn which engine wins on your corpus ## Related diff --git a/docs/guides/durable-jobs.mdx b/docs/guides/durable-jobs.mdx index 3b9b811..59e1bac 100644 --- a/docs/guides/durable-jobs.mdx +++ b/docs/guides/durable-jobs.mdx @@ -1,6 +1,6 @@ --- title: "Idempotency and cleanup" -description: "Safe retries, retention windows, and deleting hosted documents." +description: "Safe retries, retention windows, and cleaning up hosted documents." --- Hosted FileRouter is a job system: store a document, run executions, retain results for a while, then clean up. Treat document IDs, job IDs, and idempotency keys as part of your application contract. @@ -47,7 +47,7 @@ From the hosted product rules: | Normalized results | Available while retained; tied to the document cleanup window | | Job records | Scheduled for deletion after 30 days | -Storage in these windows does not consume credits. Do not rely on expiry for sensitive data. Delete when the workflow finishes. +Storage in these windows does not consume credits. Do not rely on expiry for sensitive data. Release stored artifacts when the workflow finishes. ## Cleanup after a job @@ -91,7 +91,7 @@ try { `documents.delete` removes the document, related jobs, and retained results. Use it when you do not need the history. -For hosted `parse` / `compare`, delete via the returned resources when you no longer need retained outputs: +For hosted `parse` / `compare`, release via the returned resources when you no longer need retained outputs: ```ts const parsed = await router.parse("./claim.pdf", { diff --git a/docs/guides/fast-then-strong.mdx b/docs/guides/fast-then-strong.mdx index 30c8a8d..7a54f5d 100644 --- a/docs/guides/fast-then-strong.mdx +++ b/docs/guides/fast-then-strong.mdx @@ -13,7 +13,7 @@ If you want to **avoid** starting the expensive engine until the cheap one looks 2. Create one job with a fast engine and a stronger engine. 3. `waitForExecution` on the fast provider and render. 4. `waitForExecution` on the strong provider and replace or merge. -5. Delete the document when you are done. +5. Release the stored artifacts when you are done. ```ts import { FileRouter } from "@file_router/sdk" From 14bffb9b69d1c851f48fef6d598631dbeca5e9bf Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 15:59:10 -0400 Subject: [PATCH 7/8] refactor(sdk): remove speculative provider fallbacks Keep Datalab options aligned with the current Convert API. Let FileRouter own portable output and page selection. Reject invalid statuses, page IDs, and block references instead of guessing. --- packages/filerouter/src/datalab.ts | 152 +++++++++++------------ packages/filerouter/test/datalab.test.ts | 3 +- src/lib/provider-completion.server.ts | 1 - 3 files changed, 73 insertions(+), 83 deletions(-) diff --git a/packages/filerouter/src/datalab.ts b/packages/filerouter/src/datalab.ts index 7f7ad90..44ba84c 100644 --- a/packages/filerouter/src/datalab.ts +++ b/packages/filerouter/src/datalab.ts @@ -12,6 +12,7 @@ import type { ParseOptions, ParseOutput, ParsePage, + ParsePageField, ParseResult, ProviderJobReference, ProviderJobs, @@ -41,35 +42,27 @@ export interface DatalabProviderOptions { } export type DatalabMode = "accurate" | "balanced" | "fast" -export type DatalabOutputFormat = "chunks" | "html" | "json" | "markdown" +type DatalabOutputFormat = "chunks" | "html" | "json" | "markdown" /** Per-request options named after Datalab's Convert API fields. */ export interface DatalabParseOptions { add_block_ids?: boolean additional_config?: string - checkpoint_id?: string disable_image_captions?: boolean disable_image_extraction?: boolean eval_rubric_id?: number extras?: string fence_synthetic_captions?: boolean - format_lines?: boolean - image_resolution?: number include_markdown_in_chunks?: boolean max_pages?: number mode?: DatalabMode model_override_settings?: string - output_format?: Array | string - page_range?: string paginate?: boolean processing_location?: string - /** @deprecated Use `processing_location`. */ - processing_region?: string raw?: Record save_checkpoint?: boolean skip_cache?: boolean token_efficient_markdown?: boolean - use_llm?: boolean webhook_url?: string word_bboxes?: boolean workflowstepdata_id?: number @@ -198,26 +191,30 @@ async function getDatalabJob( }) assertSuccessful(raw) - const status = readString(raw.status)?.toLowerCase() - if (status === "complete" || status === "completed") { + const status = readString(raw.status) + if (status === "complete") { return { result: normalizeDatalab( raw, job.id, parseOptions.outputs ?? [DEFAULT_PARSE_OUTPUT], + parseOptions.pageFields, parseOptions.includeRaw === true, new Date(job.submittedAt) ), status: "complete", } } - if (["cancelled", "error", "failed"].includes(status ?? "")) { + if (status === "failed") { return { - error: readString(raw.error) ?? `Datalab job ${status}.`, + error: readString(raw.error) ?? "Datalab job failed.", status: "failed", } } - return { status: status === "running" ? "running" : "pending" } + if (status === "processing") { + return { status: "running" } + } + throw datalabParseError("Datalab returned an invalid job status.") } async function createFormData( @@ -244,7 +241,8 @@ async function createFormData( value !== null && key !== "file" && key !== "file_url" && - key !== "output_format" + key !== "output_format" && + key !== "page_range" ) { body.set(key, formValue(value)) } @@ -260,10 +258,7 @@ async function createFormData( if (mode) { body.set("mode", mode) } - body.set( - "output_format", - nativeOutputFormats(nativeOptions.output_format, outputs).join(",") - ) + body.set("output_format", datalabOutputs(outputs).join(",")) if ( outputs.includes("pages") && (parseOptions.pageFields === undefined || @@ -283,12 +278,14 @@ function normalizeDatalab( raw: Record, id: string, requestedOutputs: Array, + requestedPageFields: Array | undefined, includeRaw: boolean, startedAt: Date ): ParseResult { - const pageBlocks = datalabPageBlocks(raw.json) - const pages = requestedOutputs.includes("pages") - ? pageBlocks.map(normalizePage) + const pagesRequested = requestedOutputs.includes("pages") + const pageBlocks = datalabPageBlocks(raw.json, pagesRequested) + const pages = pagesRequested + ? pageBlocks.map((page) => normalizePage(page, requestedPageFields)) : [] const markdown = readString(raw.markdown) const html = readString(raw.html) @@ -300,8 +297,7 @@ function normalizeDatalab( const costBreakdown = isRecord(raw.cost_breakdown) ? raw.cost_breakdown : undefined - const totalCost = - readNumber(costBreakdown?.total) ?? readNumber(raw.total_cost) + const totalCost = readNumber(costBreakdown?.total) const metadata = { ...(isRecord(raw.metadata) ? raw.metadata : {}), ...(typeof raw.checkpoint_id === "string" && { @@ -346,8 +342,14 @@ function normalizeDatalab( } } -function datalabPageBlocks(value: unknown): Array> { +function datalabPageBlocks( + value: unknown, + required: boolean +): Array> { if (!isRecord(value)) { + if (required) { + throw datalabParseError("Datalab did not return JSON page output.") + } return [] } return readRecords(value.children).filter( @@ -355,10 +357,20 @@ function datalabPageBlocks(value: unknown): Array> { ) } -function normalizePage(raw: Record, index: number): ParsePage { +function normalizePage( + raw: Record, + requestedFields: Array | undefined +): ParsePage { const blocks = indexDatalabBlocks(raw) - const html = resolveBlockField(raw, "html", blocks) - const markdown = resolveBlockField(raw, "markdown", blocks) + const allFields = requestedFields === undefined + const html = + allFields || requestedFields.includes("html") + ? resolveBlockField(raw, "html", blocks) + : undefined + const markdown = + allFields || requestedFields.includes("markdown") + ? resolveBlockField(raw, "markdown", blocks) + : undefined const metadata = { ...(raw.bbox !== undefined && { bbox: raw.bbox }), ...(typeof raw.id === "string" && { blockId: raw.id }), @@ -373,26 +385,19 @@ function normalizePage(raw: Record, index: number): ParsePage { json: raw, ...(markdown && { markdown }), ...(Object.keys(metadata).length > 0 && { metadata }), - pageNumber: datalabPageNumber(raw, index), + pageNumber: datalabPageNumber(raw), warnings: [], } } -function datalabPageNumber( - page: Record, - index: number -): number { - const explicit = readNumber(page.page_number) - if ( - explicit !== undefined && - Number.isSafeInteger(explicit) && - explicit > 0 - ) { - return explicit - } +function datalabPageNumber(page: Record): number { const id = readString(page.id) const zeroBased = id?.match(/(?:^|\/)page\/(\d+)(?:\/|$)/i)?.[1] - return zeroBased === undefined ? index + 1 : Number(zeroBased) + 1 + const pageNumber = zeroBased === undefined ? NaN : Number(zeroBased) + 1 + if (!Number.isSafeInteger(pageNumber) || pageNumber <= 0) { + throw datalabParseError("Datalab returned a page with an invalid ID.") + } + return pageNumber } function indexDatalabBlocks( @@ -424,29 +429,44 @@ function resolveBlockField( const resolve = (block: Record): string | undefined => { const value = readString(block[field]) if (!value) { - const children = readRecords(block.children).flatMap((child) => { - const childValue = resolve(child) - return childValue ? [childValue] : [] - }) - return children.length > 0 ? children.join("\n\n") : undefined + return undefined } return value.replace( /]*\bsrc=(["'])([^"']+)\1[^>]*>\s*<\/content-ref>/gi, - (reference, _quote: string, childId: string) => { + (_reference, _quote: string, childId: string) => { const child = blocks.get(childId) - if (!child || visiting.has(childId)) { - return reference + if (!child) { + throw datalabParseError( + `Datalab page references a missing block: ${childId}.` + ) + } + if (visiting.has(childId)) { + throw datalabParseError( + `Datalab page contains a cyclic block reference: ${childId}.` + ) } visiting.add(childId) const resolved = resolve(child) visiting.delete(childId) - return resolved ?? reference + if (resolved === undefined) { + throw datalabParseError( + `Datalab block is missing ${field}: ${childId}.` + ) + } + return resolved } ) } return resolve(root) } +function datalabParseError(message: string): FileRouterError { + return new FileRouterError(message, { + code: "ParseFailed", + providerId: PROVIDER_ID, + }) +} + function datalabOutputs( outputs: Array ): Array { @@ -470,36 +490,6 @@ function datalabOutputs( return [...supported] } -function nativeOutputFormats( - value: DatalabParseOptions["output_format"], - outputs: Array -): Array { - const native = Array.isArray(value) - ? value - : typeof value === "string" - ? value.split(",").map((format) => format.trim()) - : [] - const supported = new Set(datalabOutputs(outputs)) - for (const format of native) { - if ( - format !== "chunks" && - format !== "html" && - format !== "json" && - format !== "markdown" - ) { - throw new FileRouterError( - `Unsupported Datalab output format: ${format}`, - { - code: "ParseFailed", - providerId: PROVIDER_ID, - } - ) - } - supported.add(format) - } - return [...supported] -} - function normalizeImages(value: unknown): Array { if (!isRecord(value)) { return [] diff --git a/packages/filerouter/test/datalab.test.ts b/packages/filerouter/test/datalab.test.ts index 1e00d19..2dd0d2c 100644 --- a/packages/filerouter/test/datalab.test.ts +++ b/packages/filerouter/test/datalab.test.ts @@ -44,7 +44,6 @@ describe("Datalab provider", () => { datalab: { add_block_ids: true, mode: "accurate", - output_format: "json", }, llamaparse: { tier: "fast" }, }, @@ -183,7 +182,9 @@ describe("Datalab provider", () => { markdown: "Third", }, ], + html: "

Third

", id: "/page/2/Page/2", + markdown: "Third", }, ], }, diff --git a/src/lib/provider-completion.server.ts b/src/lib/provider-completion.server.ts index 5571e28..3b55d99 100644 --- a/src/lib/provider-completion.server.ts +++ b/src/lib/provider-completion.server.ts @@ -9,7 +9,6 @@ const TERMINAL_WORKFLOW_STATUSES = new Set([ "complete", "errored", "terminated", - "unknown", ]) type CompletionProvider = "datalab" | "llamaparse" From 51e0ef1cb87b8db59992957bae3e5d4b7f316f88 Mon Sep 17 00:00:00 2001 From: Urjit Chakraborty <135136842+urjitc@users.noreply.github.com> Date: Tue, 28 Jul 2026 22:13:36 -0400 Subject: [PATCH 8/8] fix(sdk): honor Datalab page field selection --- packages/filerouter/src/datalab.ts | 24 ++++++++++++++---------- packages/filerouter/test/datalab.test.ts | 23 ++++++++++++++++++++--- 2 files changed, 34 insertions(+), 13 deletions(-) diff --git a/packages/filerouter/src/datalab.ts b/packages/filerouter/src/datalab.ts index 44ba84c..cd43fa7 100644 --- a/packages/filerouter/src/datalab.ts +++ b/packages/filerouter/src/datalab.ts @@ -371,20 +371,24 @@ function normalizePage( allFields || requestedFields.includes("markdown") ? resolveBlockField(raw, "markdown", blocks) : undefined - const metadata = { - ...(raw.bbox !== undefined && { bbox: raw.bbox }), - ...(typeof raw.id === "string" && { blockId: raw.id }), - ...(raw.polygon !== undefined && { polygon: raw.polygon }), - ...(raw.section_hierarchy !== undefined && { - sectionHierarchy: raw.section_hierarchy, - }), - } + const json = allFields || requestedFields.includes("json") ? raw : undefined + const metadata = + allFields || requestedFields.includes("metadata") + ? { + ...(raw.bbox !== undefined && { bbox: raw.bbox }), + ...(typeof raw.id === "string" && { blockId: raw.id }), + ...(raw.polygon !== undefined && { polygon: raw.polygon }), + ...(raw.section_hierarchy !== undefined && { + sectionHierarchy: raw.section_hierarchy, + }), + } + : undefined return { ...(html && { html }), - json: raw, + ...(json && { json }), ...(markdown && { markdown }), - ...(Object.keys(metadata).length > 0 && { metadata }), + ...(metadata && Object.keys(metadata).length > 0 && { metadata }), pageNumber: datalabPageNumber(raw), warnings: [], } diff --git a/packages/filerouter/test/datalab.test.ts b/packages/filerouter/test/datalab.test.ts index 2dd0d2c..1f500ed 100644 --- a/packages/filerouter/test/datalab.test.ts +++ b/packages/filerouter/test/datalab.test.ts @@ -249,8 +249,18 @@ describe("Datalab provider", () => { ) .mockResolvedValueOnce( Response.json({ - json: { children: [] }, - page_count: 0, + json: { + children: [ + { + bbox: [0, 0, 100, 100], + block_type: "Page", + html: "

First

", + id: "/page/0/Page/0", + markdown: "First", + }, + ], + }, + page_count: 1, status: "complete", success: true, }) @@ -265,7 +275,7 @@ describe("Datalab provider", () => { }, }) - await router.parse("https://example.com/report.pdf", { + const result = await router.parse("https://example.com/report.pdf", { outputs: ["pages"], pageFields: ["markdown"], providerOptions: { @@ -276,6 +286,13 @@ describe("Datalab provider", () => { const body = fetchMock.mock.calls[0]?.[1]?.body expect(body).toBeInstanceOf(FormData) expect((body as FormData).get("include_markdown_in_chunks")).toBe("false") + expect(result.outputs.pages).toEqual([ + { + markdown: "First", + pageNumber: 1, + warnings: [], + }, + ]) }) test("rejects untrusted polling URLs", async () => {