diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 53918f74..895d6e2b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -60,6 +60,7 @@ jobs: eager-workflow-start early-return empty + external-storage google-adk-agents hello-world langsmith diff --git a/.scripts/copy-shared-files.mjs b/.scripts/copy-shared-files.mjs index 36b71515..47226efb 100644 --- a/.scripts/copy-shared-files.mjs +++ b/.scripts/copy-shared-files.mjs @@ -31,6 +31,7 @@ const TSCONFIG_EXCLUDE = [ 'langsmith', ]; const GITIGNORE_EXCLUDE = [ + 'external-storage', 'nextjs-ecommerce-oneclick', 'monorepo-folders', 'production', diff --git a/.scripts/list-of-samples.json b/.scripts/list-of-samples.json index 4e635dba..ef981eff 100644 --- a/.scripts/list-of-samples.json +++ b/.scripts/list-of-samples.json @@ -17,6 +17,7 @@ "encryption", "env-config", "expense", + "external-storage", "fetch-esm", "food-delivery", "google-adk-agents", diff --git a/README.md b/README.md index adf0bf73..5e4947ea 100644 --- a/README.md +++ b/README.md @@ -145,6 +145,7 @@ and you'll be given the list of sample options. - [**Worker Versioning**](./worker-versioning): Version Workers with Build IDs in order to deploy incompatible changes to Workflow code. - [**Protobufs**](./protobufs): Use [Protobufs](https://docs.temporal.io/security/#default-data-converter). - [**Custom Payload Converter**](./ejson): Customize data serialization by creating a `PayloadConverter` that uses EJSON to convert Dates, binary, and regexes. +- [**External Storage**](./external-storage): Keep large payloads out of Workflow History by writing a custom `StorageDriver` that offloads them to disk, where other Workers can read them back. - **Monorepos**: - [`/monorepo-folders`](./monorepo-folders): yarn workspace with packages for a web frontend, API server, Worker, and Workflows/Activities. - [`psigen/temporal-ts-example`](https://github.com/psigen/temporal-ts-example): yarn workspace containerized with [tilt](https://tilt.dev/). Includes `temporalite`, `parcel`, and different packages for Workflows and Activities. diff --git a/external-storage/.eslintignore b/external-storage/.eslintignore new file mode 100644 index 00000000..7bd99a41 --- /dev/null +++ b/external-storage/.eslintignore @@ -0,0 +1,3 @@ +node_modules +lib +.eslintrc.js \ No newline at end of file diff --git a/external-storage/.eslintrc.js b/external-storage/.eslintrc.js new file mode 100644 index 00000000..9f199cd9 --- /dev/null +++ b/external-storage/.eslintrc.js @@ -0,0 +1,48 @@ +const { builtinModules } = require('module'); + +const ALLOWED_NODE_BUILTINS = new Set(['assert']); + +module.exports = { + root: true, + parser: '@typescript-eslint/parser', + parserOptions: { + project: './tsconfig.json', + tsconfigRootDir: __dirname, + }, + plugins: ['@typescript-eslint', 'deprecation'], + extends: [ + 'eslint:recommended', + 'plugin:@typescript-eslint/eslint-recommended', + 'plugin:@typescript-eslint/recommended', + 'prettier', + ], + rules: { + // recommended for safety + '@typescript-eslint/no-floating-promises': 'error', // forgetting to await Activities and Workflow APIs is bad + 'deprecation/deprecation': 'warn', + + // code style preference + 'object-shorthand': ['error', 'always'], + + // relaxed rules, for convenience + '@typescript-eslint/no-unused-vars': [ + 'warn', + { + argsIgnorePattern: '^_', + varsIgnorePattern: '^_', + }, + ], + '@typescript-eslint/no-explicit-any': 'off', + }, + overrides: [ + { + files: ['src/**/workflows.ts', 'src/**/workflows-*.ts', 'src/**/workflows/*.ts'], + rules: { + 'no-restricted-imports': [ + 'error', + ...builtinModules.filter((m) => !ALLOWED_NODE_BUILTINS.has(m)).flatMap((m) => [m, `node:${m}`]), + ], + }, + }, + ], +}; diff --git a/external-storage/.gitignore b/external-storage/.gitignore new file mode 100644 index 00000000..1aec77ff --- /dev/null +++ b/external-storage/.gitignore @@ -0,0 +1,3 @@ +lib +node_modules +storage diff --git a/external-storage/.npmrc b/external-storage/.npmrc new file mode 100644 index 00000000..9cf94950 --- /dev/null +++ b/external-storage/.npmrc @@ -0,0 +1 @@ +package-lock=false \ No newline at end of file diff --git a/external-storage/.nvmrc b/external-storage/.nvmrc new file mode 100644 index 00000000..2bd5a0a9 --- /dev/null +++ b/external-storage/.nvmrc @@ -0,0 +1 @@ +22 diff --git a/external-storage/.post-create b/external-storage/.post-create new file mode 100644 index 00000000..055c11e9 --- /dev/null +++ b/external-storage/.post-create @@ -0,0 +1,18 @@ +To begin development, install the Temporal CLI: + +Mac: {cyan brew install temporal} +Other: Download and extract the latest release from https://github.com/temporalio/cli/releases/latest + +Start Temporal Server: + +{cyan temporal server start-dev} + +Use Node version 18+ (v22.x is recommended): + +Mac: {cyan brew install node@22} +Other: https://nodejs.org/en/download/ + +Then, in the project directory, using two other shells, run these commands: + +{cyan npm run start.watch} +{cyan npm run workflow} diff --git a/external-storage/.prettierignore b/external-storage/.prettierignore new file mode 100644 index 00000000..7951405f --- /dev/null +++ b/external-storage/.prettierignore @@ -0,0 +1 @@ +lib \ No newline at end of file diff --git a/external-storage/.prettierrc b/external-storage/.prettierrc new file mode 100644 index 00000000..965d50bf --- /dev/null +++ b/external-storage/.prettierrc @@ -0,0 +1,2 @@ +printWidth: 120 +singleQuote: true diff --git a/external-storage/README.md b/external-storage/README.md new file mode 100644 index 00000000..e7d1819d --- /dev/null +++ b/external-storage/README.md @@ -0,0 +1,156 @@ +# External Storage + +> **Experimental.** External storage shipped in SDK 1.21.0 and is marked experimental; +> its API may change. + +Temporal stores every Workflow argument, Activity result, Signal, and heartbeat detail in +Workflow History, and enforces a size limit on each one. External storage moves the large +ones out: payloads over a size threshold are written to storage you control and replaced +on the wire by a small reference. The retrieving side resolves the reference before your +code ever sees it, so Workflow and Activity code stays unchanged. + +This sample writes a **custom driver** end to end. It keeps payloads on the filesystem, so +a Client in one process can hand a 1 MiB argument to a Worker in another process without +either of them putting it in the Temporal database. + +## How it works + +`ExternalStorage` is configured on the `DataConverter`, on both the Client and the Worker: + +```ts +new Client({ + connection, + dataConverter: { + externalStorage: new ExternalStorage({ + drivers: [new FileSystemStorageDriver({ rootDir })], + payloadSizeThreshold: 32 * 1024, + }), + }, +}); +``` + +A driver is four members: + +```ts +interface StorageDriver { + readonly name: string; // routing key written into the reference; must match across processes + readonly type: string; // stable implementation ID, reported via Worker heartbeat + store(context, payloads): Promise; + retrieve(context, claims): Promise; +} +``` + +The SDK handles the rest: it measures each payload, batches the over-threshold ones per +driver, calls `store`, and swaps in a reference carrying the driver name and the claim you +returned. Ordering is your contract to keep, one claim per payload. A claim is an opaque +`Record`, and it lands in Workflow History, so keep it small and free of +secrets. + +`store` receives a `target` describing the Workflow or Activity that produced the payloads +(namespace, ID, run ID, type), which this driver uses to lay out keys. If a driver throws, +the enclosing Workflow or Activity Task fails **retryably**, so transient I/O errors +recover on their own. + +## Code + +- [`filesystem-storage-driver.ts`](./src/filesystem-storage-driver.ts) — the custom driver. + Content-addresses each payload by SHA-256, writes it atomically, verifies the hash on + read, and refuses keys that escape the storage root. The comments cover the decisions + and the alternatives at each one. +- [`data-converter.ts`](./src/data-converter.ts) — wires the driver into an + `ExternalStorage`, with the size threshold and the multi-driver `driverSelector` option. +- [`workflows.ts`](./src/workflows.ts) — an ordinary Workflow. Documents which of its four + payloads get offloaded, by whom, and when. +- [`activities.ts`](./src/activities.ts) — one Activity that takes and returns a large + payload, one that takes a large payload and returns a small one. +- [`worker.ts`](./src/worker.ts) / [`client.ts`](./src/client.ts) — both sides configured. +- [`inspect.ts`](./src/inspect.ts) — prints the references Temporal Server actually holds + alongside the blobs on disk they point at. + +## Running this sample + +1. `temporal server start-dev` to start [Temporal Server](https://github.com/temporalio/cli/#installation). +1. `npm install` to install dependencies. +1. `npm run start.watch` to start the Worker. +1. In another shell, `npm run workflow` to run the Workflow Client. +1. `npm run inspect ` to see what was stored where. + +The Client passes a 1 MiB document and gets ~1 MiB of extracted text back. Inline, that +same payload would cross the wire five times, each crossing over the SDK's default 512 KiB +outbound size warning and pushing toward the server's 2 MiB per-payload limit: + +``` +Starting workflow with a 1048590 byte document +Payloads of 32768 bytes or more are offloaded to .../external-storage/storage +Started workflow document-V1StGXR8_Z5jdHi6B +Summary: 17404 lines, 165338 words, 1048590 characters +Received 1048590 bytes of extracted text + +To see what the server actually stored, run: + npm run inspect document-V1StGXR8_Z5jdHi6B +``` + +`npm run inspect` then shows the same execution from the server's side. Every large +payload is a reference of a few hundred bytes; the small one was left inline (keys +abbreviated here): + +``` +History payloads for document-V1StGXR8_Z5jdHi6B: + + #1 WORKFLOW_EXECUTION_STARTED: 300 bytes on the wire (reference to 1066065 bytes in 'sample.filesystemdriver', key=v1/wf/default/processDocument/document-.../null/sha256/1558f39e...) + #5 ACTIVITY_TASK_SCHEDULED: 332 bytes on the wire (reference to 1066065 bytes in 'sample.filesystemdriver', key=v1/wf/.../d0dd96bc-.../sha256/1558f39e...) + #7 ACTIVITY_TASK_COMPLETED: 332 bytes on the wire (reference to 1066023 bytes in 'sample.filesystemdriver', key=v1/wf/.../d0dd96bc-.../sha256/dfb0fab1...) + #11 ACTIVITY_TASK_SCHEDULED: 332 bytes on the wire (reference to 1066023 bytes in 'sample.filesystemdriver', key=v1/wf/.../d0dd96bc-.../sha256/dfb0fab1...) + #13 ACTIVITY_TASK_COMPLETED: 47 bytes on the wire (inline) + #17 WORKFLOW_EXECUTION_COMPLETED: 332 bytes on the wire (reference to 1066137 bytes in 'sample.filesystemdriver', key=v1/wf/.../d0dd96bc-.../sha256/0a4efcbf...) +``` + +Four things this output shows: + +- **The Client and the Worker are separate processes.** The blob behind event #1 was + written by the Client and read by the Worker. Nothing coordinates them beyond both + drivers resolving the same directory. +- **The threshold is per payload, not per Activity.** `summarize` returned 47 bytes at + #13, under the 32 KiB threshold, so it stayed inline. Its 1 MiB _argument_ at #11 did + not. +- **Identical content deduplicates.** #7 and #11 are the same key: the extracted text was + stored once when the Activity completed, and the reference was reused when it became the + next Activity's argument. +- **Deduplication stops at the key prefix.** #1 and #5 have the same hash under different + prefixes, so the document is on disk twice. The Client stored it before a run ID + existed, hence the `null` segment; the Worker stored it again under the real run ID. + Four blobs, ~4 MiB, for one 1 MiB document. `buildKeyPrefix` in the driver explains the + tradeoff and how to trade it the other way. + +## Testing + +`npm test` runs both suites: + +- [`filesystem-storage-driver.test.ts`](./src/test/filesystem-storage-driver.test.ts) — + the driver on its own: byte-for-byte round trips, deduplication, a second driver + instance reading what the first wrote, and the failure paths (corrupted blob, hostile + Workflow ID, claim pointing outside the storage root, oversized payload). +- [`workflows.test.ts`](./src/test/workflows.test.ts) — the Workflow against a test + server, asserting the large payloads are references in History and that below-threshold + payloads are not offloaded at all. + +## Before using this in production + +- **Prefer a vended driver.** The SDK ships + [`@temporalio/external-storage-s3`](https://github.com/temporalio/sdk-typescript/tree/main/contrib/external-storage-s3) + and + [`@temporalio/external-storage-gcs`](https://github.com/temporalio/sdk-typescript/tree/main/contrib/external-storage-gcs). + Write your own only if neither backend fits. This driver exists to show what the + interface asks of you. +- **Shared storage is required.** A local directory only works here because everything + runs on one machine. Real Workers need storage all of them can reach. +- **Nothing deletes blobs.** Retention is on you: an object-store lifecycle policy, or a + reaper keyed on Workflow completion. Blobs must outlive every History that references + them, including retention on completed Executions, replay, and Workflow Reset. +- **Payloads leave Temporal's trust boundary.** Whatever you offload is now protected by + your storage's access control and encryption at rest, not Temporal's. Combine with a + `PayloadCodec` if you need the bytes encrypted before they land there (see the + [encryption](../encryption) sample). +- **The Web UI shows references, not values.** Offloaded payloads render as + `ExternalStorageReference` in History. A [Codec Server](https://docs.temporal.io/production-deployment/data-encryption) + can resolve them for viewing. diff --git a/external-storage/package.json b/external-storage/package.json new file mode 100644 index 00000000..e6da4610 --- /dev/null +++ b/external-storage/package.json @@ -0,0 +1,54 @@ +{ + "name": "temporal-external-storage", + "version": "0.1.0", + "private": true, + "scripts": { + "build": "tsc --build", + "build.watch": "tsc --build --watch", + "format": "prettier --write .", + "format:check": "prettier --check .", + "lint": "eslint .", + "test": "mocha --require ts-node/register --require source-map-support/register src/test/*.test.ts", + "test.watch": "mocha --require ts-node/register --require source-map-support/register src/test/*.test.ts -w --watch-files src", + "start": "ts-node src/worker.ts", + "start.watch": "nodemon src/worker.ts", + "workflow": "ts-node src/client.ts", + "inspect": "ts-node src/inspect.ts" + }, + "nodemonConfig": { + "execMap": { + "ts": "ts-node" + }, + "ext": "ts", + "watch": [ + "src" + ] + }, + "dependencies": { + "@temporalio/activity": "^1.24.0", + "@temporalio/client": "^1.24.0", + "@temporalio/common": "^1.24.0", + "@temporalio/envconfig": "^1.24.0", + "@temporalio/proto": "^1.24.0", + "@temporalio/worker": "^1.24.0", + "@temporalio/workflow": "^1.24.0", + "nanoid": "^3.3.8" + }, + "devDependencies": { + "@temporalio/testing": "^1.24.0", + "@tsconfig/node22": "^22.0.0", + "@types/mocha": "^9.1.1", + "@types/node": "^22.9.1", + "@typescript-eslint/eslint-plugin": "^8.18.0", + "@typescript-eslint/parser": "^8.18.0", + "eslint": "^8.57.1", + "eslint-config-prettier": "^9.1.0", + "eslint-plugin-deprecation": "^3.0.0", + "mocha": "^10.0.0", + "nodemon": "^3.1.7", + "prettier": "^3.4.2", + "source-map-support": "^0.5.21", + "ts-node": "^10.9.2", + "typescript": "^5.6.3" + } +} diff --git a/external-storage/src/activities.ts b/external-storage/src/activities.ts new file mode 100644 index 00000000..dd67da8b --- /dev/null +++ b/external-storage/src/activities.ts @@ -0,0 +1,35 @@ +import { log } from '@temporalio/activity'; +import type { Document } from './shared'; + +/** + * Stands in for OCR or text extraction: takes a large document in and returns a large + * string out. Both the argument and the return value are offloaded, so neither ever + * reaches Temporal Server. + * + * Nothing here is aware of external storage. By the time the Activity runs, the Worker + * has already retrieved the argument through the driver; the returned string is stored + * on the way back out. + */ +export async function extractText(document: Document): Promise { + log.info('extracting text', { document: document.name, contentLength: document.content.length }); + + const extractedText = document.content + .split('\n') + .map((line) => line.trim()) + .filter((line) => line.length > 0) + .join('\n'); + + log.info('extracted text', { extractedLength: extractedText.length }); + return extractedText; +} + +/** + * Takes a large argument and returns a small result, so its argument is offloaded and + * its return value stays inline. Mixing both in one Workflow shows the threshold at + * work: offloading is decided per payload, by size, not per Activity. + */ +export async function summarize(text: string): Promise { + const lines = text.split('\n'); + const words = text.split(/\s+/).filter((word) => word.length > 0); + return `${lines.length} lines, ${words.length} words, ${text.length} characters`; +} diff --git a/external-storage/src/client.ts b/external-storage/src/client.ts new file mode 100644 index 00000000..96babe06 --- /dev/null +++ b/external-storage/src/client.ts @@ -0,0 +1,43 @@ +import { Client, Connection } from '@temporalio/client'; +import { loadClientConnectConfig } from '@temporalio/envconfig'; +import { nanoid } from 'nanoid'; +import { PAYLOAD_SIZE_THRESHOLD, STORAGE_ROOT, createDataConverter } from './data-converter'; +import { TASK_QUEUE, makeDocument } from './shared'; +import { processDocument } from './workflows'; + +const DOCUMENT_SIZE_BYTES = 1024 * 1024; + +async function run() { + const config = loadClientConnectConfig(); + const connection = await Connection.connect(config.connectionOptions); + + // The Client offloads too. Without external storage configured here, the 1 MiB + // argument below would be sent inline and rejected for exceeding Temporal's + // per-payload limit, and an offloaded result would come back as an unreadable + // reference. + const client = new Client({ connection, dataConverter: createDataConverter() }); + + const document = makeDocument('quarterly-report.txt', DOCUMENT_SIZE_BYTES); + console.log(`Starting workflow with a ${document.content.length} byte document`); + console.log(`Payloads of ${PAYLOAD_SIZE_THRESHOLD} bytes or more are offloaded to ${STORAGE_ROOT}`); + + const workflowId = `document-${nanoid()}`; + const handle = await client.workflow.start(processDocument, { + args: [document], + taskQueue: TASK_QUEUE, + workflowId, + }); + console.log(`Started workflow ${handle.workflowId}`); + + const result = await handle.result(); + console.log(`Summary: ${result.summary}`); + console.log(`Received ${result.extractedText.length} bytes of extracted text`); + console.log(`\nTo see what the server actually stored, run:\n npm run inspect ${workflowId}`); + + await connection.close(); +} + +run().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/external-storage/src/data-converter.ts b/external-storage/src/data-converter.ts new file mode 100644 index 00000000..78bdae27 --- /dev/null +++ b/external-storage/src/data-converter.ts @@ -0,0 +1,62 @@ +import * as path from 'node:path'; +import { ExternalStorage } from '@temporalio/common'; +import type { DataConverter } from '@temporalio/common'; +import { FileSystemStorageDriver } from './filesystem-storage-driver'; + +/** + * Where the blobs live. Every process in this sample resolves the same directory, which + * is what lets the Worker read a payload the Client wrote, in a different process. + * + * A real deployment points this at storage all Workers can reach: a shared volume, or + * a driver backed by S3 or GCS. The SDK ships drivers for both, so if object storage + * is where you are headed, prefer `@temporalio/external-storage-s3` or + * `@temporalio/external-storage-gcs` over writing your own. This sample writes a + * driver from scratch to show what the interface asks of you. + */ +export const STORAGE_ROOT = process.env.EXTERNAL_STORAGE_DIR ?? path.resolve(__dirname, '..', 'storage'); + +/** + * Payloads at or above this size are offloaded; smaller ones are sent inline. Set low + * here so the sample offloads without needing huge test data. The SDK default is + * 256 KiB, which is a reasonable production starting point: it comfortably clears + * Temporal's 2 MiB per-payload limit while leaving small payloads (the vast majority) + * on the fast path, with no storage round trip. + * + * Setting this to `0` offloads every payload regardless of size. That is occasionally + * useful for compliance ("no business data in the Temporal database"), but it puts a + * storage round trip in front of every Signal, Query, and Activity argument. + */ +export const PAYLOAD_SIZE_THRESHOLD = 32 * 1024; + +/** + * Builds the DataConverter shared by the Worker and the Client. + * + * Both sides need it, and their driver `name`s must match: the reference payload in + * History names the driver that wrote it, and the retrieving side looks it up by that + * name. A Client without external storage configured cannot read an offloaded result; + * it sees the raw reference and fails to deserialize it. + * + * Unlike `payloadConverterPath`, which is a *path* because the Workflow sandbox has to + * load it separately, `externalStorage` is passed as a live object. Storing and + * retrieving happen in Worker and Client code, outside the sandbox, so the driver is + * free to do I/O and hold connections. + */ +// @@@SNIPSTART typescript-custom-driver-data-converter +export function createDataConverter(rootDir: string = STORAGE_ROOT): DataConverter { + return { + externalStorage: new ExternalStorage({ + drivers: [new FileSystemStorageDriver({ rootDir })], + payloadSizeThreshold: PAYLOAD_SIZE_THRESHOLD, + + // With one driver registered, every offloaded payload goes to it. Register more + // than one and a `driverSelector` becomes required, letting you route per + // payload: a cheap archive tier for a known-bulky Workflow type, a driver per + // tenant, or a per-region bucket. Returning `null` from the selector keeps that + // payload inline, which is how you exempt specific payloads from offloading. + // + // driverSelector: (context, _payload) => + // context.target?.type === 'processDocument' ? coldDriver : hotDriver, + }), + }; +} +// @@@SNIPEND diff --git a/external-storage/src/filesystem-storage-driver.ts b/external-storage/src/filesystem-storage-driver.ts new file mode 100644 index 00000000..65c6c78e --- /dev/null +++ b/external-storage/src/filesystem-storage-driver.ts @@ -0,0 +1,329 @@ +import { createHash, randomUUID } from 'node:crypto'; +import { access, mkdir, readFile, rename, unlink, writeFile } from 'node:fs/promises'; +import * as path from 'node:path'; +import { StorageDriverClaim } from '@temporalio/common'; +import type { + Payload, + StorageDriver, + StorageDriverRetrieveContext, + StorageDriverStoreContext, + StorageDriverTargetInfo, +} from '@temporalio/common'; +import { temporal } from '@temporalio/proto'; + +const PayloadProto = temporal.api.common.v1.Payload; + +/** + * Prefix on every key. Keys are stored verbatim in the claim, so old claims keep + * resolving after you change the layout below: bump this and new blobs land under + * the new scheme while existing references still point at the old one. + */ +const KEY_LAYOUT_VERSION = 'v1'; + +/** Guardrail against a single runaway payload filling the disk. */ +const DEFAULT_MAX_PAYLOAD_SIZE = 50 * 1024 * 1024; + +const HASH_ALGORITHM = 'sha256'; + +/** Key segment written when a target field is absent (e.g. a run ID we don't know yet). */ +const NULL_SEGMENT = 'null'; + +/** Characters allowed in a path segment. Everything else is escaped. */ +const UNSAFE_SEGMENT_CHARS = /[^A-Za-z0-9._-]/g; + +/** Cap on a single path segment, to stay well inside filesystem limits (255 bytes on ext4/APFS). */ +const MAX_SEGMENT_LENGTH = 120; + +export interface FileSystemStorageDriverOptions { + /** + * Directory that holds the blobs. Every process that needs to read a payload back + * must be able to reach this same directory, so in a real deployment this is a + * shared mount (NFS, EFS, a Kubernetes RWX volume), not a local temp dir. + */ + rootDir: string; + + /** + * Routing name written into the reference payload on the wire. The retrieving side + * looks the driver up by this name, so it must match across every Client and Worker + * that touches the workflow, and it must stay stable for as long as any history + * still references blobs written by this driver. + */ + driverName?: string; + + /** Maximum serialized size of a single payload, in bytes. Defaults to 50 MiB. */ + maxPayloadSize?: number; +} + +/** + * A custom {@link StorageDriver} that offloads large payloads to a filesystem + * directory instead of sending them to Temporal Server. + * + * The SDK calls {@link store} on the way out and replaces each stored payload with a + * small reference containing this driver's `name` and the returned claim. On the way + * in, it reads the reference, finds the driver by name, and calls {@link retrieve}. + * Workflow and Activity code never sees any of this: it gets the original value. + * + * Blobs are content-addressed by a SHA-256 hash of the serialized payload: + * + * - Writing the same bytes twice is a no-op, which matters because Temporal retries. + * An Activity that heartbeats the same large details repeatedly, or a Workflow Task + * that replays after a failure, does not accumulate duplicate blobs. + * - The hash doubles as an integrity check on read (see {@link retrieve}). + * + * The alternative, keying by a random UUID, makes cleanup easier to reason about + * (one blob has exactly one referrer) at the cost of both properties above. + * + * Note on process boundaries: an in-memory `Map` driver is tempting and about ten + * lines, but it only works while the storing and retrieving code share a process. + * The moment a second Worker picks up the Activity Task, or the Client tries to read + * a Workflow result, retrieval fails. That cross-process handoff is the whole point of + * external storage, so this driver uses the filesystem. + */ +export class FileSystemStorageDriver implements StorageDriver { + readonly name: string; + + /** + * Stable identifier for this *implementation*, shared by every instance and reported + * to the server via Worker heartbeat for observability. Distinct from `name`, which + * identifies this particular configured instance. The SDK's own drivers use values + * like `aws.s3driver` and `gcp.gcsdriver`. + */ + readonly type = 'sample.filesystemdriver'; + + private readonly rootDir: string; + private readonly maxPayloadSize: number; + + constructor({ rootDir, driverName = 'sample.filesystemdriver', maxPayloadSize }: FileSystemStorageDriverOptions) { + this.rootDir = path.resolve(rootDir); + this.name = driverName; + this.maxPayloadSize = maxPayloadSize ?? DEFAULT_MAX_PAYLOAD_SIZE; + } + + /** + * Called with every payload the SDK decided to offload, batched per driver. Must + * return one claim per payload, in the same order. + * + * Throwing here fails the enclosing Workflow or Activity Task *retryably*, so a + * transient I/O error is retried rather than killing the Execution. That makes it + * safe to let errors propagate instead of, say, silently falling back to inline + * payloads, which would defeat the point of the size threshold. + */ + // @@@SNIPSTART typescript-custom-storage-driver + async store(context: StorageDriverStoreContext, payloads: Payload[]): Promise { + const keyPrefix = buildKeyPrefix(context.target); + return runAllAbortingOnFirstError(context.abortSignal, (signal) => + payloads.map((payload) => this.storePayload(payload, keyPrefix, signal)), + ); + } + + /** Inverse of {@link store}: one payload per claim, in the same order. */ + async retrieve(context: StorageDriverRetrieveContext, claims: StorageDriverClaim[]): Promise { + return runAllAbortingOnFirstError(context.abortSignal, (signal) => + claims.map((claim) => this.retrievePayload(claim, signal)), + ); + } + // @@@SNIPEND + + private async storePayload( + payload: Payload, + keyPrefix: string, + abortSignal: AbortSignal, + ): Promise { + // Store the encoded Payload proto, not just `payload.data`. A Payload also carries + // metadata (the `encoding` key, protobuf message names, anything a custom + // PayloadConverter added), and metadata values are arbitrary bytes that would not + // survive a round trip through the string-valued claim map. Encoding the whole + // message keeps the payload byte-for-byte identical end to end. + const payloadBytes = PayloadProto.encode(payload).finish(); + if (payloadBytes.length > this.maxPayloadSize) { + throw new Error( + `Payload of ${payloadBytes.length} bytes exceeds the configured maxPayloadSize of ${this.maxPayloadSize} bytes`, + ); + } + + const hashValue = createHash(HASH_ALGORITHM).update(payloadBytes).digest('hex'); + const key = `${keyPrefix}/${HASH_ALGORITHM}/${hashValue}`; + + try { + await this.writeIfAbsent(key, payloadBytes, abortSignal); + } catch (err) { + // Wrapping adds the context that makes a retry loop diagnosable: a bare ENOENT + // says nothing about which key or which storage root. On ES2022 and above, prefer + // `new Error(message, { cause: err })` to keep the original stack attached; these + // samples target ES2021, so the message is folded in instead. + throw new Error( + `FileSystemStorageDriver failed to store [rootDir=${this.rootDir}, key=${key}]: ${describe(err)}`, + ); + } + + // The claim is the only thing that reaches Temporal Server, embedded in the + // reference payload that replaces the real one. Keep it small and keep it free of + // anything sensitive: it is visible in Workflow History and in the Web UI. + // + // `rootDir` is deliberately absent. It is deployment configuration, and a second + // Worker may well mount the same storage at a different path; putting it in the + // claim would pin every historical reference to one machine's filesystem layout. + return new StorageDriverClaim({ key, hashAlgorithm: HASH_ALGORITHM, hashValue }); + } + + /** + * Writes the blob unless it is already there. The write goes to a temporary file and + * is then renamed, because `rename` is atomic: a reader can only ever observe the + * complete blob, never a half-written one. Without that, a crash mid-write would + * leave a file whose name promises content it does not contain, and content + * addressing would hand it to a reader as valid. + */ + private async writeIfAbsent(key: string, payloadBytes: Uint8Array, abortSignal: AbortSignal): Promise { + const filePath = this.resolveKey(key); + if (await exists(filePath)) return; + + await mkdir(path.dirname(filePath), { recursive: true }); + const tempPath = `${filePath}.${randomUUID()}.tmp`; + try { + await writeFile(tempPath, payloadBytes, { signal: abortSignal }); + await rename(tempPath, filePath); + } catch (err) { + await unlink(tempPath).catch(() => undefined); + // Two Workers can race to store identical bytes. On POSIX the loser's rename + // silently replaces an identical file; on Windows it fails. Either way the blob + // is present and correct, so treat that as success. + if (await exists(filePath)) return; + throw err; + } + } + + private async retrievePayload(claim: StorageDriverClaim, abortSignal: AbortSignal): Promise { + const { key, hashAlgorithm, hashValue: expectedHash } = claim.claimData; + // Claims come off the wire, so validate rather than assume. A missing field means a + // claim written by a different driver, or by an older version of this one. + if (!key) { + throw new Error("FileSystemStorageDriver claim is missing required field 'key'"); + } + if (hashAlgorithm !== HASH_ALGORITHM || !expectedHash) { + throw new Error( + `FileSystemStorageDriver claim [key=${key}] must carry hashAlgorithm='${HASH_ALGORITHM}' and a hashValue, ` + + `got hashAlgorithm='${hashAlgorithm ?? ''}'`, + ); + } + + const filePath = this.resolveKey(key); + let payloadBytes: Uint8Array; + try { + payloadBytes = await readFile(filePath, { signal: abortSignal }); + } catch (err) { + throw new Error( + `FileSystemStorageDriver failed to retrieve [rootDir=${this.rootDir}, key=${key}]: ${describe(err)}`, + ); + } + + // Verifying is cheap next to the read and catches truncation, corruption, and a + // claim pointing at the wrong blob. Failing loudly here is much better than + // handing malformed bytes to the PayloadConverter, where the error would surface + // as a confusing deserialization failure far from its cause. + const actualHash = createHash(HASH_ALGORITHM).update(payloadBytes).digest('hex'); + if (actualHash !== expectedHash) { + throw new Error( + `FileSystemStorageDriver integrity check failed [key=${key}]: ` + + `expected ${HASH_ALGORITHM}:${expectedHash}, got ${HASH_ALGORITHM}:${actualHash}`, + ); + } + + return PayloadProto.decode(payloadBytes); + } + + /** + * Maps a key to a path under `rootDir`, refusing anything that escapes it. Keys + * arrive from Workflow History, which we should not treat as trusted input: a claim + * containing `../../etc/passwd` must not turn into a read outside the blob store. + */ + private resolveKey(key: string): string { + const filePath = path.resolve(this.rootDir, key); + if (filePath !== this.rootDir && !filePath.startsWith(this.rootDir + path.sep)) { + throw new Error(`FileSystemStorageDriver refused a key that resolves outside rootDir [key=${key}]`); + } + return filePath; + } +} + +/** + * Builds the directory prefix for a blob from the Workflow or Activity that produced + * it. Nothing functionally depends on this (the claim carries the full key), but it + * makes the store browsable and gives cleanup something to work with: "delete blobs + * for this run" becomes a directory removal. + * + * The tradeoff is that deduplication only reaches within a prefix, so identical bytes + * stored under two different prefixes are written twice. This is visible in the sample: + * a Client stores a Workflow argument before the run ID exists, so its blob lands under + * a `null` run ID, and the Worker writes the same bytes again under the real run ID once + * the Workflow schedules an Activity with them. + * + * Dropping the prefix for a flat `sha256/` namespace buys global deduplication and + * gives up per-run locality, which is what cleanup keys off. The vended S3 and GCS + * drivers make the same choice this one does; whether it is right for you depends on + * whether your payloads repeat across Executions and how you plan to expire them. + */ +function buildKeyPrefix(target: StorageDriverTargetInfo | undefined): string { + if (target === undefined) return KEY_LAYOUT_VERSION; + const segments = + target.kind === 'workflow' + ? ['wf', target.namespace, target.type, target.id, target.runId] + : ['act', target.namespace, target.type, target.id, target.runId]; + return [KEY_LAYOUT_VERSION, ...segments.map(toSafeSegment)].join('/'); +} + +/** + * Escapes a value for use as a single path segment. Workflow IDs are caller-supplied + * and may contain `/`, `..`, or characters the filesystem reserves, so allow a known + * set and escape everything else rather than blocklisting known-bad input. + */ +function toSafeSegment(value: string | undefined): string { + if (!value) return NULL_SEGMENT; + let escaped = value.replace(UNSAFE_SEGMENT_CHARS, (char) => `%${char.charCodeAt(0).toString(16).padStart(2, '0')}`); + + // `.` and `..` are made of allowed characters but mean something to the filesystem, so + // they need escaping too: a Workflow ID of `..` would otherwise become a path + // component that walks up a level. `resolveKey` would still refuse to read or write + // outside `rootDir`, but the blob would land somewhere surprising on the way there. + if (escaped === '.' || escaped === '..') escaped = escaped.replace(/\./g, '%2e'); + + // Truncation can make two very long IDs share a prefix directory. Harmless, since the + // blob's own name is its content hash, but worth knowing when browsing the store. + return escaped.length <= MAX_SEGMENT_LENGTH ? escaped : escaped.slice(0, MAX_SEGMENT_LENGTH); +} + +function describe(err: unknown): string { + return err instanceof Error ? err.message : String(err); +} + +async function exists(filePath: string): Promise { + try { + await access(filePath); + return true; + } catch { + return false; + } +} + +/** + * Runs the per-payload operations concurrently and, on the first failure, aborts the + * rest before propagating the error. The SDK hands us an `abortSignal` and expects + * siblings to be cancelled on first error: once one payload in a batch fails, the + * enclosing Task is going to fail anyway, so finishing the other writes is wasted work. + */ +async function runAllAbortingOnFirstError( + externalSignal: AbortSignal | undefined, + makeTasks: (signal: AbortSignal) => Promise[], +): Promise { + const controller = new AbortController(); + const signal = externalSignal ? AbortSignal.any([externalSignal, controller.signal]) : controller.signal; + const tasks = makeTasks(signal); + try { + return await Promise.all(tasks); + } catch (err) { + controller.abort(); + // Wait for the aborted siblings to settle so no write is still in flight (and no + // temp file un-cleaned) by the time the Task retries. + await Promise.allSettled(tasks); + throw err; + } +} diff --git a/external-storage/src/inspect.ts b/external-storage/src/inspect.ts new file mode 100644 index 00000000..d99eeede --- /dev/null +++ b/external-storage/src/inspect.ts @@ -0,0 +1,134 @@ +import { readdir, stat } from 'node:fs/promises'; +import * as path from 'node:path'; +import { Client, Connection } from '@temporalio/client'; +import { loadClientConnectConfig } from '@temporalio/envconfig'; +import { temporal } from '@temporalio/proto'; +import { STORAGE_ROOT } from './data-converter'; + +/** Metadata written by the SDK on the reference payload that replaces an offloaded one. */ +const REFERENCE_ENCODING = 'json/protobuf'; +const REFERENCE_MESSAGE_TYPE = 'temporal.api.sdk.v1.ExternalStorageReference'; + +/** + * Shows both halves of the round trip for one Workflow Execution: the small references + * Temporal Server holds, and the blobs on disk they point at. + * + * Usage: `npm run inspect ` + */ +async function run() { + const workflowId = process.argv[2]; + if (!workflowId) { + console.error('Usage: npm run inspect '); + process.exit(1); + } + + const config = loadClientConnectConfig(); + const connection = await Connection.connect(config.connectionOptions); + + // Deliberately *without* external storage. A Client configured with it retrieves + // references transparently, which is exactly what we want to look behind here. + const client = new Client({ connection }); + + const { events } = await client.workflow.getHandle(workflowId).fetchHistory(); + + console.log(`History payloads for ${workflowId}:\n`); + for (const event of events ?? []) { + const eventType = (temporal.api.enums.v1.EventType[event.eventType ?? 0] ?? 'UNKNOWN').replace('EVENT_TYPE_', ''); + for (const payload of collectPayloads(event)) { + const wireSize = payload.data?.length ?? 0; + const reference = describeReference(payload); + const detail = reference ?? 'inline'; + console.log(` #${event.eventId} ${eventType}: ${wireSize} bytes on the wire (${detail})`); + } + } + + const blobs = await listBlobs(STORAGE_ROOT); + console.log(`\nBlobs under ${STORAGE_ROOT}:\n`); + if (blobs.length === 0) { + console.log(' (none: no payload has crossed the size threshold yet)'); + } + let total = 0; + for (const blob of blobs) { + total += blob.size; + console.log(` ${blob.size} bytes ${blob.relativePath}`); + } + console.log(`\n${blobs.length} blob(s), ${total} bytes total`); + + await connection.close(); +} + +/** + * Finds every Payload nested anywhere in a history event. + * + * Payloads hang off dozens of differently shaped event attributes (`input`, `result`, + * `details`, `lastHeartbeatDetails`, memo fields, ...), so this walks the decoded event + * instead of enumerating them. Fine for an inspection script; application code should + * reach for the specific field it cares about. + */ +function collectPayloads( + node: unknown, + found: temporal.api.common.v1.IPayload[] = [], +): temporal.api.common.v1.IPayload[] { + if (node === null || typeof node !== 'object' || node instanceof Uint8Array) return found; + if (isPayload(node)) { + found.push(node); + return found; + } + for (const value of Object.values(node)) { + collectPayloads(value, found); + } + return found; +} + +function isPayload(node: object): node is temporal.api.common.v1.IPayload { + const candidate = node as { metadata?: unknown; data?: unknown }; + return candidate.data instanceof Uint8Array && typeof candidate.metadata === 'object' && candidate.metadata !== null; +} + +/** Describes a payload if it is an external storage reference, else `null`. */ +function describeReference(payload: temporal.api.common.v1.IPayload): string | null { + if ( + readMetadata(payload, 'encoding') !== REFERENCE_ENCODING || + readMetadata(payload, 'messageType') !== REFERENCE_MESSAGE_TYPE + ) { + return null; + } + + // The reference is a protobuf-JSON encoded ExternalStorageReference: the driver name + // plus the claim that driver handed back. The original size travels alongside it in + // `externalPayloads`, so tooling can report the real payload size without a fetch. + const { driverName, claimData } = JSON.parse(Buffer.from(payload.data ?? []).toString()) as { + driverName?: string; + claimData?: Record; + }; + const originalSize = payload.externalPayloads?.[0]?.sizeBytes; + return `reference to ${originalSize ?? '?'} bytes in '${driverName}', key=${claimData?.key}`; +} + +function readMetadata(payload: temporal.api.common.v1.IPayload, key: string): string | undefined { + const raw = payload.metadata?.[key]; + return raw ? Buffer.from(raw).toString() : undefined; +} + +async function listBlobs(rootDir: string): Promise<{ relativePath: string; size: number }[]> { + let entries: string[]; + try { + entries = await readdir(rootDir, { recursive: true }); + } catch { + return []; + } + + const blobs = []; + for (const entry of entries) { + const stats = await stat(path.join(rootDir, entry)); + if (stats.isFile()) { + blobs.push({ relativePath: entry, size: stats.size }); + } + } + return blobs.sort((a, b) => a.relativePath.localeCompare(b.relativePath)); +} + +run().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/external-storage/src/shared.ts b/external-storage/src/shared.ts new file mode 100644 index 00000000..d3194a18 --- /dev/null +++ b/external-storage/src/shared.ts @@ -0,0 +1,34 @@ +export const TASK_QUEUE = 'external-storage'; + +export interface Document { + name: string; + /** Raw text of a scanned document. Large enough to be worth keeping out of Workflow History. */ + content: string; +} + +export interface ProcessingResult { + documentName: string; + /** Small enough to stay inline. */ + summary: string; + /** Large, and offloaded on its way back to the Client. */ + extractedText: string; +} + +const LOREM = [ + 'the quick brown fox jumps over the lazy dog', + 'invoice total due on receipt net thirty terms apply', + 'shipment manifest reviewed and countersigned by the depot', + 'all measurements recorded in metric units unless noted', +]; + +/** Builds a deterministic document of roughly `sizeBytes` characters. */ +export function makeDocument(name: string, sizeBytes: number): Document { + const lines: string[] = []; + let length = 0; + for (let i = 0; length < sizeBytes; i++) { + const line = `${String(i).padStart(6, '0')} ${LOREM[i % LOREM.length]}`; + lines.push(line); + length += line.length + 1; + } + return { name, content: lines.join('\n') }; +} diff --git a/external-storage/src/test/filesystem-storage-driver.test.ts b/external-storage/src/test/filesystem-storage-driver.test.ts new file mode 100644 index 00000000..eb09c9a0 --- /dev/null +++ b/external-storage/src/test/filesystem-storage-driver.test.ts @@ -0,0 +1,166 @@ +import assert from 'assert'; +import { mkdtemp, readdir, stat, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import * as path from 'node:path'; +import { before, describe, it } from 'mocha'; +import { StorageDriverClaim } from '@temporalio/common'; +import type { Payload, StorageDriverTargetInfo } from '@temporalio/common'; +import { defaultPayloadConverter } from '@temporalio/common'; +import { FileSystemStorageDriver } from '../filesystem-storage-driver'; + +const workflowTarget: StorageDriverTargetInfo = { + kind: 'workflow', + namespace: 'default', + id: 'doc-1', + runId: 'run-1', + type: 'processDocument', +}; + +function toPayload(value: unknown): Payload { + return defaultPayloadConverter.toPayload(value)!; +} + +function toBytes(data: Uint8Array | null | undefined): Buffer { + return Buffer.from(data ?? new Uint8Array()); +} + +function readEncoding(payload: Payload): string { + return Buffer.from(payload.metadata?.encoding ?? new Uint8Array()).toString(); +} + +async function listFiles(rootDir: string): Promise { + const files: string[] = []; + for (const entry of await readdir(rootDir, { recursive: true })) { + const filePath = path.join(rootDir, entry); + if ((await stat(filePath)).isFile()) files.push(filePath); + } + return files; +} + +describe('FileSystemStorageDriver', function () { + let rootDir: string; + let driver: FileSystemStorageDriver; + + before(async () => { + rootDir = await mkdtemp(path.join(tmpdir(), 'external-storage-test-')); + driver = new FileSystemStorageDriver({ rootDir }); + }); + + it('round-trips a payload byte-for-byte', async () => { + const original = toPayload({ text: 'x'.repeat(100_000), nested: { flag: true } }); + + const [claim] = await driver.store({ target: workflowTarget }, [original]); + const [retrieved] = await driver.retrieve({}, [claim]); + + // Compare bytes rather than the objects: protobuf decoding yields `Buffer` where the + // PayloadConverter produced `Uint8Array`. Buffer extends Uint8Array, so the SDK + // treats them alike, but `deepStrictEqual` compares prototypes and would fail. + assert.strictEqual(Buffer.compare(toBytes(retrieved.data), toBytes(original.data)), 0, 'payload data should match'); + assert.deepStrictEqual(Object.keys(retrieved.metadata ?? {}), Object.keys(original.metadata ?? {})); + assert.strictEqual(readEncoding(retrieved), readEncoding(original), 'metadata should survive the round trip'); + assert.deepStrictEqual( + defaultPayloadConverter.fromPayload(retrieved), + defaultPayloadConverter.fromPayload(original), + ); + }); + + it('returns one claim per payload, in order', async () => { + const payloads = [toPayload('first'), toPayload('second'), toPayload('third')]; + + const claims = await driver.store({ target: workflowTarget }, payloads); + const retrieved = await driver.retrieve({}, claims); + + assert.strictEqual(claims.length, payloads.length); + assert.deepStrictEqual( + retrieved.map((payload) => defaultPayloadConverter.fromPayload(payload)), + ['first', 'second', 'third'], + ); + }); + + it('writes identical payloads once', async () => { + const localRoot = await mkdtemp(path.join(tmpdir(), 'external-storage-dedupe-')); + const localDriver = new FileSystemStorageDriver({ rootDir: localRoot }); + const payload = toPayload('the same bytes twice'); + + const [first] = await localDriver.store({ target: workflowTarget }, [payload]); + const [second] = await localDriver.store({ target: workflowTarget }, [payload]); + + assert.deepStrictEqual(first.claimData, second.claimData, 'identical content should produce identical claims'); + assert.strictEqual((await listFiles(localRoot)).length, 1); + }); + + it('keys blobs under the storing workflow', async () => { + const localRoot = await mkdtemp(path.join(tmpdir(), 'external-storage-key-')); + const localDriver = new FileSystemStorageDriver({ rootDir: localRoot }); + + const [claim] = await localDriver.store({ target: workflowTarget }, [toPayload('keyed')]); + + assert.match(claim.claimData.key, /^v1\/wf\/default\/processDocument\/doc-1\/run-1\/sha256\/[0-9a-f]{64}$/); + }); + + // Workflow IDs are caller-supplied and become part of the key, so they must not be + // able to steer a write out of the blob store. + for (const hostileWorkflowId of ['../../escape attempt', '..', '.', 'a/b']) { + it(`neutralizes a workflow ID of '${hostileWorkflowId}'`, async () => { + const localRoot = await mkdtemp(path.join(tmpdir(), 'external-storage-escape-')); + const localDriver = new FileSystemStorageDriver({ rootDir: localRoot }); + + const [claim] = await localDriver.store({ target: { ...workflowTarget, id: hostileWorkflowId } }, [ + toPayload('escaped'), + ]); + + const components = claim.claimData.key.split('/'); + assert.ok(!components.includes('..') && !components.includes('.'), `traversal in key: ${claim.claimData.key}`); + + const [file] = await listFiles(localRoot); + assert.ok(file?.startsWith(localRoot + path.sep), `blob written outside rootDir: ${file}`); + + // And it still reads back, which is the part a hostile ID must not break. + const [retrieved] = await localDriver.retrieve({}, [claim]); + assert.strictEqual(defaultPayloadConverter.fromPayload(retrieved), 'escaped'); + }); + } + + it('a second driver instance reads what the first one wrote', async () => { + // The point of external storage: the process that stores a payload is usually not + // the process that reads it back. + const [claim] = await driver.store({ target: workflowTarget }, [toPayload('written by another process')]); + + const otherWorkerDriver = new FileSystemStorageDriver({ rootDir }); + const [retrieved] = await otherWorkerDriver.retrieve({}, [claim]); + + assert.strictEqual(defaultPayloadConverter.fromPayload(retrieved), 'written by another process'); + }); + + it('rejects a corrupted blob', async () => { + const [claim] = await driver.store({ target: workflowTarget }, [toPayload('will be corrupted')]); + await writeFile(path.join(rootDir, claim.claimData.key), 'not the promised bytes'); + + await assert.rejects(driver.retrieve({}, [claim]), /integrity check failed/); + }); + + it('rejects a claim whose key escapes rootDir', async () => { + const claim = new StorageDriverClaim({ + key: '../../../etc/passwd', + hashAlgorithm: 'sha256', + hashValue: 'f'.repeat(64), + }); + + await assert.rejects(driver.retrieve({}, [claim]), /resolves outside rootDir/); + }); + + it('rejects a claim without integrity information', async () => { + const claim = new StorageDriverClaim({ key: 'v1/sha256/deadbeef' }); + + await assert.rejects(driver.retrieve({}, [claim]), /hashAlgorithm/); + }); + + it('refuses a payload larger than maxPayloadSize', async () => { + const smallLimitDriver = new FileSystemStorageDriver({ rootDir, maxPayloadSize: 1024 }); + + await assert.rejects( + smallLimitDriver.store({ target: workflowTarget }, [toPayload('x'.repeat(2048))]), + /exceeds the configured maxPayloadSize/, + ); + }); +}); diff --git a/external-storage/src/test/workflows.test.ts b/external-storage/src/test/workflows.test.ts new file mode 100644 index 00000000..8e8093e1 --- /dev/null +++ b/external-storage/src/test/workflows.test.ts @@ -0,0 +1,136 @@ +import assert from 'assert'; +import { mkdtemp, readdir, stat } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import * as path from 'node:path'; +import { after, before, describe, it } from 'mocha'; +import { nanoid } from 'nanoid'; +import { Client } from '@temporalio/client'; +import { TestWorkflowEnvironment } from '@temporalio/testing'; +import { DefaultLogger, Runtime, Worker } from '@temporalio/worker'; +import * as activities from '../activities'; +import { PAYLOAD_SIZE_THRESHOLD, createDataConverter } from '../data-converter'; +import { makeDocument } from '../shared'; +import { processDocument } from '../workflows'; + +const REFERENCE_MESSAGE_TYPE = 'temporal.api.sdk.v1.ExternalStorageReference'; + +describe('processDocument with external storage', function () { + let env: TestWorkflowEnvironment; + let storageRoot: string; + + this.slow(10_000); + this.timeout(60_000); + + before(async function () { + Runtime.install({ logger: new DefaultLogger('WARN') }); + env = await TestWorkflowEnvironment.createTimeSkipping(); + storageRoot = await mkdtemp(path.join(tmpdir(), 'external-storage-e2e-')); + }); + + after(async () => { + await env?.teardown(); + }); + + /** + * Runs the Workflow with external storage configured on both the Client and the + * Worker, as a deployment would. Returns the Workflow ID so a test can go back and + * look at what the server actually recorded. + */ + async function runWorkflow(documentSizeBytes: number): Promise<{ workflowId: string; extractedLength: number }> { + const taskQueue = `test-${nanoid()}`; + const workflowId = `test-${nanoid()}`; + const dataConverter = createDataConverter(storageRoot); + + const worker = await Worker.create({ + connection: env.nativeConnection, + taskQueue, + workflowsPath: require.resolve('../workflows'), + activities, + dataConverter, + }); + const client = new Client({ connection: env.connection, dataConverter }); + + const result = await worker.runUntil( + client.workflow.execute(processDocument, { + args: [makeDocument('report.txt', documentSizeBytes)], + taskQueue, + workflowId, + }), + ); + + return { workflowId, extractedLength: result.extractedText.length }; + } + + it('round-trips a document larger than the threshold', async () => { + const documentSizeBytes = PAYLOAD_SIZE_THRESHOLD * 8; + const { extractedLength } = await runWorkflow(documentSizeBytes); + + // Workflow and Activity code see the whole document; the offloading is invisible to + // them. + assert.ok( + extractedLength >= documentSizeBytes, + `expected at least ${documentSizeBytes} characters of extracted text, got ${extractedLength}`, + ); + assert.ok((await listBlobs(storageRoot)).length > 0, 'expected blobs to be written to external storage'); + }); + + it('keeps the large payloads out of Workflow History', async () => { + const { workflowId } = await runWorkflow(PAYLOAD_SIZE_THRESHOLD * 8); + + // A Client *without* external storage sees what the server holds: references. + // A Client *with* it configured, as in the test above, gets the real values back. + const { events } = await env.client.workflow.getHandle(workflowId).fetchHistory(); + + const startInput = events?.find((event) => event.workflowExecutionStartedEventAttributes) + ?.workflowExecutionStartedEventAttributes?.input?.payloads?.[0]; + assert.ok(startInput, 'expected a start input payload'); + assert.strictEqual(readMetadata(startInput.metadata, 'messageType'), REFERENCE_MESSAGE_TYPE); + assert.ok( + (startInput.data?.length ?? 0) < PAYLOAD_SIZE_THRESHOLD, + 'the reference that replaced the document should be small', + ); + + const activityInput = events?.find((event) => event.activityTaskScheduledEventAttributes) + ?.activityTaskScheduledEventAttributes?.input?.payloads?.[0]; + assert.ok(activityInput, 'expected an activity input payload'); + assert.strictEqual(readMetadata(activityInput.metadata, 'messageType'), REFERENCE_MESSAGE_TYPE); + }); + + it('leaves payloads below the threshold inline', async () => { + const emptyRoot = await mkdtemp(path.join(tmpdir(), 'external-storage-inline-')); + const taskQueue = `test-${nanoid()}`; + const dataConverter = createDataConverter(emptyRoot); + + const worker = await Worker.create({ + connection: env.nativeConnection, + taskQueue, + workflowsPath: require.resolve('../workflows'), + activities, + dataConverter, + }); + const client = new Client({ connection: env.connection, dataConverter }); + + await worker.runUntil( + client.workflow.execute(processDocument, { + args: [makeDocument('memo.txt', 256)], + taskQueue, + workflowId: `test-${nanoid()}`, + }), + ); + + assert.deepStrictEqual(await listBlobs(emptyRoot), [], 'nothing should be offloaded below the threshold'); + }); +}); + +function readMetadata(metadata: Record | null | undefined, key: string): string | undefined { + const raw = metadata?.[key]; + return raw ? Buffer.from(raw).toString() : undefined; +} + +async function listBlobs(rootDir: string): Promise { + const blobs: string[] = []; + for (const entry of await readdir(rootDir, { recursive: true })) { + if ((await stat(path.join(rootDir, entry))).isFile()) blobs.push(entry); + } + return blobs; +} diff --git a/external-storage/src/worker.ts b/external-storage/src/worker.ts new file mode 100644 index 00000000..ae8943a9 --- /dev/null +++ b/external-storage/src/worker.ts @@ -0,0 +1,27 @@ +import { Worker } from '@temporalio/worker'; +import * as activities from './activities'; +import { PAYLOAD_SIZE_THRESHOLD, STORAGE_ROOT, createDataConverter } from './data-converter'; +import { TASK_QUEUE } from './shared'; + +async function run() { + console.log(`Offloading payloads of ${PAYLOAD_SIZE_THRESHOLD} bytes or more to ${STORAGE_ROOT}`); + + const worker = await Worker.create({ + workflowsPath: require.resolve('./workflows'), + activities, + taskQueue: TASK_QUEUE, + + // The Worker stores payloads on the way out (Activity arguments, Activity and + // Workflow results, heartbeat details) and retrieves them on the way in. Run a + // second copy of this Worker and it will read blobs the first one wrote, without + // any coordination beyond both drivers pointing at the same storage. + dataConverter: createDataConverter(), + }); + + await worker.run(); +} + +run().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/external-storage/src/workflows.ts b/external-storage/src/workflows.ts new file mode 100644 index 00000000..31ca1f6f --- /dev/null +++ b/external-storage/src/workflows.ts @@ -0,0 +1,44 @@ +import { proxyActivities } from '@temporalio/workflow'; +import type * as activities from './activities'; +import type { Document, ProcessingResult } from './shared'; + +const { extractText, summarize } = proxyActivities({ + startToCloseTimeout: '1 minute', +}); + +/** + * Ordinary Workflow code. Nothing in it refers to external storage, and that is the + * point: offloading is a DataConverter concern, so turning it on does not change how a + * Workflow is written. + * + * Four payloads cross a process boundary here, and each is offloaded or inlined purely + * on its own size: + * + * 1. `document` — stored by the *Client* before it calls StartWorkflowExecution, and + * retrieved by the Worker before this function is first invoked. + * 2. the `extractText` argument — stored by the Worker when it completes the Workflow + * Task carrying the ScheduleActivityTask command, and retrieved by whichever Worker + * picks up that Activity Task. On a multi-Worker Task Queue that is usually a + * different process, which is why the blobs have to be somewhere shared. + * 3. the `extractText` result — stored when the Activity completes, retrieved when the + * result is delivered into this Workflow's next activation. + * 4. the return value below — stored when the Workflow completes, retrieved by the + * Client in `handle.result()`. + * + * The `summarize` result is small, so it stays inline and never touches the driver. + * + * Note that no driver code runs inside the Workflow sandbox: the sandbox has no I/O, + * and store/retrieve happen in Worker code on either side of the activation. + */ +export async function processDocument(document: Document): Promise { + const extractedText = await extractText(document); + const summary = await summarize(extractedText); + + // Returning a large result is only practical *because* of external storage; inline it + // would run into Temporal's per-payload size limit. Small results are still the + // better default, since every Client that reads this one pays a storage round trip. + // The alternative is to return a claim of your own (an object key, a row ID) and let + // callers fetch the bulk themselves, at the cost of doing by hand exactly what the + // driver already does. + return { documentName: document.name, summary, extractedText }; +} diff --git a/external-storage/tsconfig.json b/external-storage/tsconfig.json new file mode 100644 index 00000000..488f2c62 --- /dev/null +++ b/external-storage/tsconfig.json @@ -0,0 +1,13 @@ +{ + "extends": "@tsconfig/node22/tsconfig.json", + "version": "5.6.3", + "compilerOptions": { + "lib": ["es2021"], + "declaration": true, + "declarationMap": true, + "sourceMap": true, + "rootDir": "./src", + "outDir": "./lib" + }, + "include": ["src/**/*.ts"] +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index e0e52654..77c351da 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -1120,6 +1120,79 @@ importers: specifier: ^5.6.3 version: 5.7.3 + external-storage: + dependencies: + '@temporalio/activity': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/client': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/common': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/envconfig': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/proto': + specifier: ^1.24.0 + version: 1.24.0 + '@temporalio/worker': + specifier: ^1.24.0 + version: 1.24.0(@swc/helpers@0.5.15) + '@temporalio/workflow': + specifier: ^1.24.0 + version: 1.24.0 + nanoid: + specifier: ^3.3.8 + version: 3.3.12 + devDependencies: + '@temporalio/testing': + specifier: ^1.24.0 + version: 1.24.0(@swc/helpers@0.5.15) + '@tsconfig/node22': + specifier: ^22.0.0 + version: 22.0.5 + '@types/mocha': + specifier: ^9.1.1 + version: 9.1.1 + '@types/node': + specifier: ^22.9.1 + version: 22.12.0 + '@typescript-eslint/eslint-plugin': + specifier: ^8.18.0 + version: 8.22.0(@typescript-eslint/parser@8.22.0(eslint@8.57.1)(typescript@5.7.3))(eslint@8.57.1)(typescript@5.7.3) + '@typescript-eslint/parser': + specifier: ^8.18.0 + version: 8.22.0(eslint@8.57.1)(typescript@5.7.3) + eslint: + specifier: ^8.57.1 + version: 8.57.1 + eslint-config-prettier: + specifier: ^9.1.0 + version: 9.1.0(eslint@8.57.1) + eslint-plugin-deprecation: + specifier: ^3.0.0 + version: 3.0.0(eslint@8.57.1)(typescript@5.7.3) + mocha: + specifier: ^10.0.0 + version: 10.2.0(ts-node@10.9.2(@swc/core@1.10.11(@swc/helpers@0.5.15))(@types/node@22.12.0)(typescript@5.7.3)) + nodemon: + specifier: ^3.1.7 + version: 3.1.9 + prettier: + specifier: ^3.4.2 + version: 3.4.2 + source-map-support: + specifier: ^0.5.21 + version: 0.5.21 + ts-node: + specifier: ^10.9.2 + version: 10.9.2(@swc/core@1.10.11(@swc/helpers@0.5.15))(@types/node@22.12.0)(typescript@5.7.3) + typescript: + specifier: ^5.6.3 + version: 5.7.3 + fetch-esm: dependencies: '@temporalio/activity': @@ -6633,183 +6706,155 @@ packages: resolution: {integrity: sha512-9B+taZ8DlyyqzZQnoeIvDVR/2F4EbMepXMc/NdVbkzsJbzkUjhXv/70GQJ7tdLA4YJgNP25zukcxpX2/SueNrA==} cpu: [arm64] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-arm64@1.2.4': resolution: {integrity: sha512-excjX8DfsIcJ10x1Kzr4RcWe1edC9PquDRRPx3YVCvQv+U5p7Yin2s32ftzikXojb1PIFc/9Mt28/y+iRklkrw==} cpu: [arm64] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-arm@1.0.5': resolution: {integrity: sha512-gvcC4ACAOPRNATg/ov8/MnbxFDJqf/pDePbBnuBDcjsI8PssmjoKMAz4LtLaVi+OnSb5FK/yIOamqDwGmXW32g==} cpu: [arm] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-arm@1.2.4': resolution: {integrity: sha512-bFI7xcKFELdiNCVov8e44Ia4u2byA+l3XtsAj+Q8tfCwO6BQ8iDojYdvoPMqsKDkuoOo+X6HZA0s0q11ANMQ8A==} cpu: [arm] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-ppc64@1.2.4': resolution: {integrity: sha512-FMuvGijLDYG6lW+b/UvyilUWu5Ayu+3r2d1S8notiGCIyYU/76eig1UfMmkZ7vwgOrzKzlQbFSuQfgm7GYUPpA==} cpu: [ppc64] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-riscv64@1.2.4': resolution: {integrity: sha512-oVDbcR4zUC0ce82teubSm+x6ETixtKZBh/qbREIOcI3cULzDyb18Sr/Wcyx7NRQeQzOiHTNbZFF1UwPS2scyGA==} cpu: [riscv64] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-s390x@1.0.4': resolution: {integrity: sha512-u7Wz6ntiSSgGSGcjZ55im6uvTrOxSIS8/dgoVMoiGE9I6JAfU50yH5BoDlYA1tcuGS7g/QNtetJnxA6QEsCVTA==} cpu: [s390x] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-s390x@1.2.4': resolution: {integrity: sha512-qmp9VrzgPgMoGZyPvrQHqk02uyjA0/QrTO26Tqk6l4ZV0MPWIW6LTkqOIov+J1yEu7MbFQaDpwdwJKhbJvuRxQ==} cpu: [s390x] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-x64@1.0.4': resolution: {integrity: sha512-MmWmQ3iPFZr0Iev+BAgVMb3ZyC4KeFc3jFxnNbEPas60e1cIfevbtuyf9nDGIzOaW9PdnDciJm+wFFaTlj5xYw==} cpu: [x64] os: [linux] - libc: [glibc] '@img/sharp-libvips-linux-x64@1.2.4': resolution: {integrity: sha512-tJxiiLsmHc9Ax1bz3oaOYBURTXGIRDODBqhveVHonrHJ9/+k89qbLl0bcJns+e4t4rvaNBxaEZsFtSfAdquPrw==} cpu: [x64] os: [linux] - libc: [glibc] '@img/sharp-libvips-linuxmusl-arm64@1.0.4': resolution: {integrity: sha512-9Ti+BbTYDcsbp4wfYib8Ctm1ilkugkA/uscUn6UXK1ldpC1JjiXbLfFZtRlBhjPZ5o1NCLiDbg8fhUPKStHoTA==} cpu: [arm64] os: [linux] - libc: [musl] '@img/sharp-libvips-linuxmusl-arm64@1.2.4': resolution: {integrity: sha512-FVQHuwx1IIuNow9QAbYUzJ+En8KcVm9Lk5+uGUQJHaZmMECZmOlix9HnH7n1TRkXMS0pGxIJokIVB9SuqZGGXw==} cpu: [arm64] os: [linux] - libc: [musl] '@img/sharp-libvips-linuxmusl-x64@1.0.4': resolution: {integrity: sha512-viYN1KX9m+/hGkJtvYYp+CCLgnJXwiQB39damAO7WMdKWlIhmYTfHjwSbQeUK/20vY154mwezd9HflVFM1wVSw==} cpu: [x64] os: [linux] - libc: [musl] '@img/sharp-libvips-linuxmusl-x64@1.2.4': resolution: {integrity: sha512-+LpyBk7L44ZIXwz/VYfglaX/okxezESc6UxDSoyo2Ks6Jxc4Y7sGjpgU9s4PMgqgjj1gZCylTieNamqA1MF7Dg==} cpu: [x64] os: [linux] - libc: [musl] '@img/sharp-linux-arm64@0.33.5': resolution: {integrity: sha512-JMVv+AMRyGOHtO1RFBiJy/MBsgz0x4AWrT6QoEVVTyh1E39TrCUpTRI7mx9VksGX4awWASxqCYLCV4wBZHAYxA==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm64] os: [linux] - libc: [glibc] '@img/sharp-linux-arm64@0.34.5': resolution: {integrity: sha512-bKQzaJRY/bkPOXyKx5EVup7qkaojECG6NLYswgktOZjaXecSAeCWiZwwiFf3/Y+O1HrauiE3FVsGxFg8c24rZg==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm64] os: [linux] - libc: [glibc] '@img/sharp-linux-arm@0.33.5': resolution: {integrity: sha512-JTS1eldqZbJxjvKaAkxhZmBqPRGmxgu+qFKSInv8moZ2AmT5Yib3EQ1c6gp493HvrvV8QgdOXdyaIBrhvFhBMQ==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm] os: [linux] - libc: [glibc] '@img/sharp-linux-arm@0.34.5': resolution: {integrity: sha512-9dLqsvwtg1uuXBGZKsxem9595+ujv0sJ6Vi8wcTANSFpwV/GONat5eCkzQo/1O6zRIkh0m/8+5BjrRr7jDUSZw==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm] os: [linux] - libc: [glibc] '@img/sharp-linux-ppc64@0.34.5': resolution: {integrity: sha512-7zznwNaqW6YtsfrGGDA6BRkISKAAE1Jo0QdpNYXNMHu2+0dTrPflTLNkpc8l7MUP5M16ZJcUvysVWWrMefZquA==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [ppc64] os: [linux] - libc: [glibc] '@img/sharp-linux-riscv64@0.34.5': resolution: {integrity: sha512-51gJuLPTKa7piYPaVs8GmByo7/U7/7TZOq+cnXJIHZKavIRHAP77e3N2HEl3dgiqdD/w0yUfiJnII77PuDDFdw==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [riscv64] os: [linux] - libc: [glibc] '@img/sharp-linux-s390x@0.33.5': resolution: {integrity: sha512-y/5PCd+mP4CA/sPDKl2961b+C9d+vPAveS33s6Z3zfASk2j5upL6fXVPZi7ztePZ5CuH+1kW8JtvxgbuXHRa4Q==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [s390x] os: [linux] - libc: [glibc] '@img/sharp-linux-s390x@0.34.5': resolution: {integrity: sha512-nQtCk0PdKfho3eC5MrbQoigJ2gd1CgddUMkabUj+rBevs8tZ2cULOx46E7oyX+04WGfABgIwmMC0VqieTiR4jg==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [s390x] os: [linux] - libc: [glibc] '@img/sharp-linux-x64@0.33.5': resolution: {integrity: sha512-opC+Ok5pRNAzuvq1AG0ar+1owsu842/Ab+4qvU879ippJBHvyY5n2mxF1izXqkPYlGuP/M556uh53jRLJmzTWA==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [x64] os: [linux] - libc: [glibc] '@img/sharp-linux-x64@0.34.5': resolution: {integrity: sha512-MEzd8HPKxVxVenwAa+JRPwEC7QFjoPWuS5NZnBt6B3pu7EG2Ge0id1oLHZpPJdn3OQK+BQDiw9zStiHBTJQQQQ==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [x64] os: [linux] - libc: [glibc] '@img/sharp-linuxmusl-arm64@0.33.5': resolution: {integrity: sha512-XrHMZwGQGvJg2V/oRSUfSAfjfPxO+4DkiRh6p2AFjLQztWUuY/o8Mq0eMQVIY7HJ1CDQUJlxGGZRw1a5bqmd1g==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm64] os: [linux] - libc: [musl] '@img/sharp-linuxmusl-arm64@0.34.5': resolution: {integrity: sha512-fprJR6GtRsMt6Kyfq44IsChVZeGN97gTD331weR1ex1c1rypDEABN6Tm2xa1wE6lYb5DdEnk03NZPqA7Id21yg==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [arm64] os: [linux] - libc: [musl] '@img/sharp-linuxmusl-x64@0.33.5': resolution: {integrity: sha512-WT+d/cgqKkkKySYmqoZ8y3pxx7lx9vVejxW/W4DOFMYVSkErR+w7mf2u8m/y4+xHe7yY9DAXQMWQhpnMuFfScw==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [x64] os: [linux] - libc: [musl] '@img/sharp-linuxmusl-x64@0.34.5': resolution: {integrity: sha512-Jg8wNT1MUzIvhBFxViqrEhWDGzqymo3sV7z7ZsaWbZNDLXRJZoRGrjulp60YYtV4wfY8VIKcWidjojlLcWrd8Q==} engines: {node: ^18.17.0 || ^20.3.0 || >=21.0.0} cpu: [x64] os: [linux] - libc: [musl] '@img/sharp-wasm32@0.33.5': resolution: {integrity: sha512-ykUW4LVGaMcU9lu9thv85CbRMAwfeadCJHRsg2GmeRa/cJxsVY9Rbd57JcMxBkKHag5U/x7TSBpScF4U8ElVzg==} @@ -7337,56 +7382,48 @@ packages: engines: {node: '>= 10'} cpu: [arm64] os: [linux] - libc: [glibc] '@next/swc-linux-arm64-gnu@16.2.9': resolution: {integrity: sha512-hBD75iWpUtkL9SmQmcRhmLomn9jgkPzCEkbOcLgHymPEKzv+6ONy13RRiIEz/iEObjkS2Jlb5gYS2XGoS3X4rw==} engines: {node: '>= 10'} cpu: [arm64] os: [linux] - libc: [glibc] '@next/swc-linux-arm64-musl@15.1.6': resolution: {integrity: sha512-+n3u//bfsrIaZch4cgOJ3tXCTbSxz0s6brJtU3SzLOvkJlPQMJ+eHVRi6qM2kKKKLuMY+tcau8XD9CJ1OjeSQQ==} engines: {node: '>= 10'} cpu: [arm64] os: [linux] - libc: [musl] '@next/swc-linux-arm64-musl@16.2.9': resolution: {integrity: sha512-qZTI3pf9SGc/obr8NkQAekBxmp1QK+kVm+VAf3BALLfFAj+1kUhkTxmrWpVos9R/UYIA8AWX2p6cGI5WdwzVUA==} engines: {node: '>= 10'} cpu: [arm64] os: [linux] - libc: [musl] '@next/swc-linux-x64-gnu@15.1.6': resolution: {integrity: sha512-SpuDEXixM3PycniL4iVCLyUyvcl6Lt0mtv3am08sucskpG0tYkW1KlRhTgj4LI5ehyxriVVcfdoxuuP8csi3kQ==} engines: {node: '>= 10'} cpu: [x64] os: [linux] - libc: [glibc] '@next/swc-linux-x64-gnu@16.2.9': resolution: {integrity: sha512-xm0HfRNX+UkH4R3c18ynswjj5o5uEj/7iI9p9omdtTSIsRCzQqkGMA+10nzJ4EHnYC3as65IMhbbl5fWRUWHYg==} engines: {node: '>= 10'} cpu: [x64] os: [linux] - libc: [glibc] '@next/swc-linux-x64-musl@15.1.6': resolution: {integrity: sha512-L4druWmdFSZIIRhF+G60API5sFB7suTbDRhYWSjiw0RbE+15igQvE2g2+S973pMGvwN3guw7cJUjA/TmbPWTHQ==} engines: {node: '>= 10'} cpu: [x64] os: [linux] - libc: [musl] '@next/swc-linux-x64-musl@16.2.9': resolution: {integrity: sha512-QumimHkGEG6vM3PfEDWKyKen03NcqLOkeKB1EfcPe7VxzmEiCa4jNnMyBn/US5zcd/VE1CI+O8Ovb3lfjVHfGw==} engines: {node: '>= 10'} cpu: [x64] os: [linux] - libc: [musl] '@next/swc-win32-arm64-msvc@15.1.6': resolution: {integrity: sha512-s8w6EeqNmi6gdvM19tqKKWbCyOBvXFbndkGHl+c9YrzsLARRdCHsD9S1fMj8gsXm9v8vhC8s3N8rjuC/XrtkEg==} @@ -8244,79 +8281,66 @@ packages: resolution: {integrity: sha512-Q8CBCCQtDFrYtXoeUXSrnFXKOnyUhx6bz+SkL6A0E7V8kAiCJ5pamq1WtbfpVGhR5TSpXY6ak3avmDc5fHTyJA==} cpu: [arm] os: [linux] - libc: [glibc] '@rollup/rollup-linux-arm-musleabihf@4.61.1': resolution: {integrity: sha512-nwnhk1581l0FBVellGcVCAT0Oi06onEA3WB53sf01VO3I0UPBkMH9sXONYME2K0ovXcNayJfNtHfm6mpJElatQ==} cpu: [arm] os: [linux] - libc: [musl] '@rollup/rollup-linux-arm64-gnu@4.61.1': resolution: {integrity: sha512-x5Xr49hwt3hdW75UOZm3395YwwzPyauktslv29KpWL/T+vVAzoT3azLcTWv0eMciBNrx+DYjH4paehHoLpPvpg==} cpu: [arm64] os: [linux] - libc: [glibc] '@rollup/rollup-linux-arm64-musl@4.61.1': resolution: {integrity: sha512-unMS3H73DpaoPyyEVPjGKleM/s0mkmsauTENpw4INQY8y4+IuLNjkueQ5QCtC0D3N38Y38yhAU8OoZ20S2Tm6w==} cpu: [arm64] os: [linux] - libc: [musl] '@rollup/rollup-linux-loong64-gnu@4.61.1': resolution: {integrity: sha512-zNZzGRnAhwjFEYmvphJRV5XaQGjs62cCmeYYHUT//NbvEnHauw+I85nGG+SiVg5ld4GX8D1IbKIX+ozITQnhMQ==} cpu: [loong64] os: [linux] - libc: [glibc] '@rollup/rollup-linux-loong64-musl@4.61.1': resolution: {integrity: sha512-LdpWGL8X209B2SIvWjqlc8VZgM6PKfontSerGepuldQmHYrAOtnMCXeJkxXGbC+PPZVOuu5czJo7fNV6aeW8rQ==} cpu: [loong64] os: [linux] - libc: [musl] '@rollup/rollup-linux-ppc64-gnu@4.61.1': resolution: {integrity: sha512-EC5kTtNaNGOmbMGqar8dvJy6y/hg99GAwjfBz++pxZhQATXGcRjd6c5en5wcbru0vkRmiMGsQKdMJOOf6sza4g==} cpu: [ppc64] os: [linux] - libc: [glibc] '@rollup/rollup-linux-ppc64-musl@4.61.1': resolution: {integrity: sha512-8hiwp6D4acEcNK78I4rP0/XtS1sknWIAMJBPdR4l6zUtyTm5KiTDr5bXmWt4foY7nAN7AThDHgkLIEZOWKbzWw==} cpu: [ppc64] os: [linux] - libc: [musl] '@rollup/rollup-linux-riscv64-gnu@4.61.1': resolution: {integrity: sha512-10dh/h/BqA7DuMPWSxkR8uks18FRwnwOEqr5zOTEl+NOwP/OMzKX8OFR/Of9xxDA7D5qef1Nzar5WDD2kCCr1g==} cpu: [riscv64] os: [linux] - libc: [glibc] '@rollup/rollup-linux-riscv64-musl@4.61.1': resolution: {integrity: sha512-YKJ5lg35DP17gcAOggnihe+APw9HLyj1Xn7gsmGumBJAUDa6NGXNixJzmkWLhcK9TOuuyQjdamzvJefkO7qHZQ==} cpu: [riscv64] os: [linux] - libc: [musl] '@rollup/rollup-linux-s390x-gnu@4.61.1': resolution: {integrity: sha512-Mlil5G2Jj6a7B3LWGctg+XPL9vdXYuzCtNXfxOQ0nPjc2m6ueUktocPGH9bnAM0bNRKb/bAWTujUU7IJQdQA+g==} cpu: [s390x] os: [linux] - libc: [glibc] '@rollup/rollup-linux-x64-gnu@4.61.1': resolution: {integrity: sha512-bVWIOIk6pV01p4CdUbPP7CJ/434z+OooYjDuFcR+44N35YvKUC66G8MGnvcWx5mWKW3g61J+t74l3Kj15Kwn2Q==} cpu: [x64] os: [linux] - libc: [glibc] '@rollup/rollup-linux-x64-musl@4.61.1': resolution: {integrity: sha512-qy5pBvZbqNFheBz61R1rzsezjm0J7O2oNGoWtGoY89SZYLUfxAJTBAqDChqAIdB4rCiIbi9nF7yZ83GnNiLwSw==} cpu: [x64] os: [linux] - libc: [musl] '@rollup/rollup-openbsd-x64@4.61.1': resolution: {integrity: sha512-E83TXjI4zm0+5f2qO+UOudaCYIhYwpJ5jq6YCZNIZ+6CbfhKrkAGezeiASBL9ElxAxFsRS9ZhESv8mfnj6TKeg==} @@ -8575,28 +8599,24 @@ packages: engines: {node: '>=10'} cpu: [arm64] os: [linux] - libc: [glibc] '@swc/core-linux-arm64-musl@1.10.11': resolution: {integrity: sha512-2mMscXe/ivq8c4tO3eQSbQDFBvagMJGlalXCspn0DgDImLYTEnt/8KHMUMGVfh0gMJTZ9q4FlGLo7mlnbx99MQ==} engines: {node: '>=10'} cpu: [arm64] os: [linux] - libc: [musl] '@swc/core-linux-x64-gnu@1.10.11': resolution: {integrity: sha512-eu2apgDbC4xwsigpl6LS+iyw6a3mL6kB4I+6PZMbFF2nIb1Dh7RGnu70Ai6mMn1o80fTmRSKsCT3CKMfVdeNFg==} engines: {node: '>=10'} cpu: [x64] os: [linux] - libc: [glibc] '@swc/core-linux-x64-musl@1.10.11': resolution: {integrity: sha512-0n+wPWpDigwqRay4IL2JIvAqSKCXv6nKxPig9M7+epAlEQlqX+8Oq/Ap3yHtuhjNPb7HmnqNJLCXT1Wx+BZo0w==} engines: {node: '>=10'} cpu: [x64] os: [linux] - libc: [musl] '@swc/core-win32-arm64-msvc@1.10.11': resolution: {integrity: sha512-7+bMSIoqcbXKosIVd314YjckDRPneA4OpG1cb3/GrkQTEDXmWT3pFBBlJf82hzJfw7b6lfv6rDVEFBX7/PJoLA==} @@ -9568,11 +9588,6 @@ packages: engines: {node: '>=0.4.0'} hasBin: true - acorn@8.15.0: - resolution: {integrity: sha512-NZyJarBfL7nWwIq+FDL6Zp/yHEhePMNnnJ0y3qfieCrmNvYct8uvtiV41UvlSe6apAfk0fY1FbWx+NwfmpvtTg==} - engines: {node: '>=0.4.0'} - hasBin: true - acorn@8.16.0: resolution: {integrity: sha512-UVJyE9MttOsBQIDKw1skb9nAwQuR5wuGD3+82K6JgJlm/Y+KI92oNsMNGZCYdDsVtRHSak0pcV5Dno5+4jh9sw==} engines: {node: '>=0.4.0'} @@ -18522,7 +18537,7 @@ snapshots: dependencies: '@babel/compat-data': 7.26.5 '@babel/helper-validator-option': 7.25.9 - browserslist: 4.24.4 + browserslist: 4.28.2 lru-cache: 5.1.1 semver: 6.3.1 @@ -21587,7 +21602,7 @@ snapshots: '@mikro-orm/mariadb@6.6.16(@mikro-orm/core@6.6.16)(pg@8.20.0)': dependencies: '@mikro-orm/core': 6.6.16 - '@mikro-orm/knex': 6.6.16(@mikro-orm/core@6.6.16)(mariadb@3.4.5)(mysql2@3.20.0(@types/node@22.12.0))(pg@8.20.0) + '@mikro-orm/knex': 6.6.16(@mikro-orm/core@6.6.16)(mariadb@3.4.5)(pg@8.20.0)(sqlite3@5.1.7) mariadb: 3.4.5 transitivePeerDependencies: - better-sqlite3 @@ -21640,7 +21655,7 @@ snapshots: '@mikro-orm/postgresql@6.6.16(@mikro-orm/core@6.6.16)(mariadb@3.4.5)': dependencies: '@mikro-orm/core': 6.6.16 - '@mikro-orm/knex': 6.6.16(@mikro-orm/core@6.6.16)(mariadb@3.4.5)(mysql2@3.20.0(@types/node@22.12.0))(pg@8.20.0) + '@mikro-orm/knex': 6.6.16(@mikro-orm/core@6.6.16)(mariadb@3.4.5)(pg@8.20.0)(sqlite3@5.1.7) pg: 8.20.0 postgres-array: 3.0.4 postgres-date: 2.1.0 @@ -24465,7 +24480,7 @@ snapshots: dependencies: '@mapbox/node-pre-gyp': 1.0.11(encoding@0.1.13) '@rollup/pluginutils': 4.2.1 - acorn: 8.15.0 + acorn: 8.16.0 async-sema: 3.1.1 bindings: 1.5.0 estree-walker: 2.0.2 @@ -24768,8 +24783,6 @@ snapshots: acorn@7.4.1: {} - acorn@8.15.0: {} - acorn@8.16.0: {} address@1.2.2: {}