{"_id":"@aws/durable-execution-sdk-js-insight","name":"@aws/durable-execution-sdk-js-insight","dist-tags":{"beta":"0.1.0-alpha.0","latest":"0.1.0-alpha.0"},"versions":{"0.1.0-alpha.0":{"name":"@aws/durable-execution-sdk-js-insight","description":"Workflow Insight plugin for AWS Durable Execution SDK — curated execution observability","license":"Apache-2.0","version":"0.1.0-alpha.0","private":false,"repository":{"type":"git","url":"git+ssh://git@github.com/aws/aws-durable-execution-sdk-js.git","directory":"packages/aws-durable-execution-sdk-js-insight"},"homepage":"https://github.com/aws/aws-durable-execution-sdk-js/tree/main/packages/aws-durable-execution-sdk-js-insight","engines":{"node":">=22"},"main":"./dist-cjs/index.js","module":"./dist/index.mjs","types":"./dist-types/index.d.ts","exports":{".":{"types":"./dist-types/index.d.ts","require":"./dist-cjs/index.js","import":"./dist/index.mjs"}},"scripts":{"build":"concurrently npm:build:esm npm:build:cjs npm:build:types","build:esm":"rollup --config rollup.config.mjs --environment MODE:esm","build:cjs":"rollup --config rollup.config.mjs --environment MODE:cjs","build:types":"tsc --emitDeclarationOnly --declaration --declarationMap --outDir dist-types","lint":"eslint src --ext .ts","test":"jest","clean":"rm -rf dist dist-cjs dist-types"},"peerDependencies":{"@aws/durable-execution-sdk-js":">=2.1.0","@aws-sdk/client-s3":"^3.0.0","@aws-sdk/client-dynamodb":"^3.0.0","@aws-sdk/util-dynamodb":"^3.0.0","@aws-sdk/client-rds-data":"^3.0.0","@aws-sdk/client-cloudwatch-logs":"^3.0.0","@aws-sdk/client-firehose":"^3.0.0","@aws-sdk/client-eventbridge":"^3.0.0","@aws-sdk/client-redshift-data":"^3.0.0","@aws-sdk/credential-provider-node":"^3.0.0","@smithy/signature-v4":"^4.0.0","@aws-crypto/sha256-js":"^5.0.0","@aws-sdk/client-sqs":"^3.0.0"},"peerDependenciesMeta":{"@aws-sdk/client-s3":{"optional":true},"@aws-sdk/client-dynamodb":{"optional":true},"@aws-sdk/util-dynamodb":{"optional":true},"@aws-sdk/client-rds-data":{"optional":true},"@aws-sdk/client-cloudwatch-logs":{"optional":true},"@aws-sdk/client-firehose":{"optional":true},"@aws-sdk/client-eventbridge":{"optional":true},"@aws-sdk/client-redshift-data":{"optional":true},"@aws-sdk/credential-provider-node":{"optional":true},"@smithy/signature-v4":{"optional":true},"@aws-crypto/sha256-js":{"optional":true},"@aws-sdk/client-sqs":{"optional":true}},"devDependencies":{"@aws/durable-execution-sdk-js":"*","@typescript-eslint/eslint-plugin":"^8.25.0","@typescript-eslint/parser":"^8.25.0","eslint":"^9.21.0","eslint-plugin-tsdoc":"^0.5.2","jest":"^29.7.0","ts-jest":"^29.2.6","@aws-sdk/client-s3":"^3.658.0","@aws-sdk/client-dynamodb":"^3.658.0","@aws-sdk/util-dynamodb":"^3.658.0","@aws-sdk/client-rds-data":"^3.658.0","@aws-sdk/client-cloudwatch-logs":"^3.658.0","@aws-sdk/client-firehose":"^3.658.0","@aws-sdk/client-eventbridge":"^3.658.0","@aws-sdk/client-redshift-data":"^3.658.0","@aws-sdk/credential-provider-node":"^3.658.0","@smithy/signature-v4":"^4.0.0","@aws-crypto/sha256-js":"^5.0.0","@aws-sdk/client-sqs":"^3.658.0","concurrently":"^9.2.4","rollup":"^4.0.0","@rollup/plugin-typescript":"^12.0.0","@rollup/plugin-json":"^6.0.0","@rollup/plugin-esm-shim":"^0.1.0","@rollup/plugin-replace":"^6.0.0","typescript":"~5.8.0","tslib":"^2.6.0"},"_id":"@aws/durable-execution-sdk-js-insight@0.1.0-alpha.0","gitHead":"a479e2cab116a02e9b412a796d5c1a0066b38f26","bugs":{"url":"https://github.com/aws/aws-durable-execution-sdk-js/issues"},"_nodeVersion":"18.20.2","_npmVersion":"10.5.0","dist":{"integrity":"sha512-qp1IGadZ386pifluCTwrERyQYwqOUho2YadV44f0v+km4UoSoScRxg6iftzjdk4nBRqlXPjKwKWK8Fn2RISqug==","shasum":"fbf948287e60bc862b5ca2461a853dee94aaf748","tarball":"https://registry.npmjs.org/@aws/durable-execution-sdk-js-insight/-/durable-execution-sdk-js-insight-0.1.0-alpha.0.tgz","fileCount":70,"unpackedSize":327388,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIQD1njSPDf3vjRatGcmgpfHYuIbKacFhxvfETgBI/un0mQIgdu5Ct+uK3Gq0HudQs/VFc3FXRu6UYvIIBcM+lYtDtV8="}]},"_npmUser":{"name":"durable-execution-dev","email":"durable-execution-dev@amazon.com"},"directories":{},"maintainers":[{"name":"durable-execution-dev","email":"durable-execution-dev@amazon.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/durable-execution-sdk-js-insight_0.1.0-alpha.0_1784246639496_0.1881065901117709"},"_hasShrinkwrap":false}},"time":{"created":"2026-07-17T00:03:59.308Z","0.1.0-alpha.0":"2026-07-17T00:03:59.645Z","modified":"2026-07-17T00:03:59.922Z"},"maintainers":[{"name":"durable-execution-dev","email":"durable-execution-dev@amazon.com"}],"description":"Workflow Insight plugin for AWS Durable Execution SDK — curated execution observability","homepage":"https://github.com/aws/aws-durable-execution-sdk-js/tree/main/packages/aws-durable-execution-sdk-js-insight","repository":{"type":"git","url":"git+ssh://git@github.com/aws/aws-durable-execution-sdk-js.git","directory":"packages/aws-durable-execution-sdk-js-insight"},"bugs":{"url":"https://github.com/aws/aws-durable-execution-sdk-js/issues"},"license":"Apache-2.0","readme":"# @aws/durable-execution-sdk-js-insight\n\n**Workflow Insight** is an observability plugin for [AWS Lambda Durable Functions](https://docs.aws.amazon.com/lambda/latest/dg/durable-functions.html). It automatically captures execution state — status, timing, operations, input/output, and errors — and exports it to the destination(s) of your choice.\n\nOne line of configuration gives you full visibility into your durable workflows without building custom instrumentation.\n\n## What it does\n\nEvery time your durable function runs, Workflow Insight builds a **cumulative snapshot** of the execution (`WorkflowInsightRecord`) and sends it to one or more exporters. The record includes:\n\n- **Execution identity** — ARN, function name, region, account\n- **Status** — RUNNING, SUCCEEDED, FAILED\n- **Timing** — start/end time, total duration\n- **Input/Output** — the event and result of the execution\n- **Operations** — every step, wait, invoke, and callback with individual timing and status\n- **Errors** — error name and message when the execution fails\n\n## Record Schema (`WorkflowInsightRecord`)\n\nEvery emitted record has this shape:\n\n```typescript\ninterface WorkflowInsightRecord {\n  recordType: \"WorkflowInsight\"; // Fixed discriminator — use to filter insight records\n  schemaVersion: \"1.0\";\n  emittedAt: string; // ISO-8601 timestamp of emission\n\n  // Execution identity\n  executionArn: string; // Full ARN including execution/invocation IDs\n  executionName?: string; // Customer-provided name (--durable-execution-name)\n  functionName: string; // Lambda function name\n  functionQualifier: string; // Version or alias\n  region: string; // AWS region\n  accountId: string; // AWS account ID\n\n  // Execution state\n  status: \"RUNNING\" | \"SUCCEEDED\" | \"FAILED\"; // suspends (waits/timers) surface as RUNNING\n  startTime: string; // ISO-8601\n  endTime?: string; // ISO-8601 (absent while running)\n  durationMs?: number; // Total duration (absent while running)\n\n  // Payload\n  input?: unknown; // Execution input (the event)\n  output?: unknown; // Execution result (on success)\n  error?: {\n    // Error details (on failure)\n    name: string;\n    message: string;\n  };\n\n  // Operations (steps, waits, invokes, callbacks)\n  operations: OperationRecord[];\n\n  // Truncation markers — present only when the size limiter dropped data\n  truncated?: boolean; // true when data was dropped to fit maxRecordSizeBytes\n  droppedOperations?: number; // count of whole operations dropped\n  droppedInput?: boolean; // true if execution input was dropped (last resort)\n  droppedOutput?: boolean; // true if execution output was dropped (last resort)\n}\n\ninterface OperationRecord {\n  id: string; // Stable hash ID (same across replays)\n  name?: string; // Customer-provided name (from step(\"name\", fn))\n  type: string; // STEP | WAIT | CALLBACK | CHAINED_INVOKE | CONTEXT\n  subType?: string; // Additional categorization\n  parentId?: string; // Parent operation ID (for child contexts)\n  status: string; // STARTED | SUCCEEDED | FAILED | PENDING | CANCELLED\n  startTime?: string; // ISO-8601\n  endTime?: string; // ISO-8601\n  durationMs?: number; // Operation duration\n  attempt?: number; // Retry attempt count\n  error?: {\n    // Per-operation error\n    name: string;\n    message: string;\n  };\n  truncated?: boolean; // true if the size limiter dropped this operation's result\n}\n```\n\n### Example Record\n\n```json\n{\n  \"recordType\": \"WorkflowInsight\",\n  \"schemaVersion\": \"1.0\",\n  \"emittedAt\": \"2026-06-16T17:00:27.514Z\",\n  \"executionArn\": \"arn:aws:lambda:us-east-1:123456789012:function:order-processor:$LATEST/durable-execution/abc123/inv456\",\n  \"executionName\": \"abc123\",\n  \"functionName\": \"order-processor\",\n  \"functionQualifier\": \"$LATEST\",\n  \"region\": \"us-east-1\",\n  \"accountId\": \"123456789012\",\n  \"status\": \"SUCCEEDED\",\n  \"startTime\": \"2026-06-16T17:00:22.100Z\",\n  \"endTime\": \"2026-06-16T17:00:27.514Z\",\n  \"durationMs\": 5414,\n  \"input\": { \"orderId\": \"order-12345\", \"customerId\": \"cust-789\" },\n  \"output\": {\n    \"orderId\": \"order-12345\",\n    \"status\": \"completed\",\n    \"charged\": 99.99\n  },\n  \"operations\": [\n    {\n      \"id\": \"c4ca4238a0b92382\",\n      \"name\": \"validate-order\",\n      \"type\": \"STEP\",\n      \"subType\": \"Step\",\n      \"status\": \"SUCCEEDED\",\n      \"startTime\": \"2026-06-16T17:00:22.200Z\",\n      \"endTime\": \"2026-06-16T17:00:22.450Z\",\n      \"durationMs\": 250\n    },\n    {\n      \"id\": \"c81e728d9d4c2f63\",\n      \"name\": \"check-inventory\",\n      \"type\": \"STEP\",\n      \"subType\": \"Step\",\n      \"status\": \"SUCCEEDED\",\n      \"startTime\": \"2026-06-16T17:00:22.450Z\",\n      \"endTime\": \"2026-06-16T17:00:22.800Z\",\n      \"durationMs\": 350\n    },\n    {\n      \"id\": \"eccbc87e4b5ce2fe\",\n      \"name\": \"cool-down\",\n      \"type\": \"WAIT\",\n      \"subType\": \"Wait\",\n      \"status\": \"SUCCEEDED\",\n      \"startTime\": \"2026-06-16T17:00:22.800Z\",\n      \"endTime\": \"2026-06-16T17:00:27.800Z\",\n      \"durationMs\": 5000\n    },\n    {\n      \"id\": \"a87ff679a2f3e71d\",\n      \"name\": \"charge-payment\",\n      \"type\": \"STEP\",\n      \"subType\": \"Step\",\n      \"status\": \"SUCCEEDED\",\n      \"startTime\": \"2026-06-16T17:00:27.400Z\",\n      \"endTime\": \"2026-06-16T17:00:27.480Z\",\n      \"durationMs\": 80\n    }\n  ]\n}\n```\n\n## Installation\n\n```bash\nnpm install @aws/durable-execution-sdk-js-insight\n```\n\n**Requirements:**\n\n- Node.js ≥ 22\n- `@aws/durable-execution-sdk-js` ≥ 2.0.0-alpha.1 (peer dependency)\n- Lambda runtime: `nodejs22.x` or later\n\nExporter-specific AWS SDK packages (e.g., `@aws-sdk/client-s3`) are **optional peer dependencies** — they're already available in the Lambda runtime, so you don't need to install them. They're only needed if you bundle your own dependencies or run outside Lambda.\n\n## Quick Start\n\n```typescript\nimport { withDurableExecution } from \"@aws/durable-execution-sdk-js\";\nimport { workflowInsight } from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({});\n\nexport const handler = withDurableExecution(\n  async (event, context) => {\n    const result = await context.step(\"process\", async () => doWork(event));\n    return result;\n  },\n  { plugins: [insight] },\n);\n```\n\nThat's it. With zero configuration, insight records appear in your function's own CloudWatch log group as JSON (via the default `LambdaLogExporter`). No extra IAM permissions, no infrastructure to set up.\n\n## Configuration\n\n```typescript\nconst insight = workflowInsight({\n  // When to emit records (default: \"on-complete\")\n  emitMode: \"on-complete\",\n\n  // Sampling rate: 0.0–1.0 (default: 1.0 = all executions)\n  samplingRate: 1.0,\n\n  // Which operations to include (default: \"top-level\")\n  operationDetail: \"top-level\",\n\n  // Where to send records (default: [new LambdaLogExporter()])\n  exporters: [new LambdaLogExporter()],\n\n  // Control what data is included (default: include everything)\n  content: { ... },\n});\n```\n\n### `emitMode`\n\n| Mode                      | Behavior                                                    | Use case                                              |\n| ------------------------- | ----------------------------------------------------------- | ----------------------------------------------------- |\n| `\"on-complete\"` (default) | Emit one record when execution completes (SUCCEEDED/FAILED) | Low overhead; sufficient for post-hoc analysis        |\n| `\"on-change\"`             | Emit on every operation change + at end                     | Real-time monitoring; see executions as they progress |\n| `\"on-failure\"`            | Emit one record only when execution ends in FAILED          | Lowest overhead; error-focused alerting and triage    |\n\n### `samplingRate`\n\nA number between 0 and 1 (default `1.0` = every execution). When below 1.0, only a fraction of executions emit records; the rest are skipped entirely — no records and no exporter calls.\n\n```typescript\nsamplingRate: 0.1, // Only 10% of executions emit records\n```\n\nThe decision is **per-execution and all-or-nothing**: a sampled-in execution emits all of its records, a sampled-out execution emits none — you never get fragmented partial data. It is **deterministic across replays**: the decision is derived from a hash of the execution ARN, which is stable across replays, so a resumed execution always reaches the same decision. Values outside `[0, 1]` or non-numeric values are clamped/defaulted to `1.0` with a warning.\n\n### `operationDetail`\n\nControls which operations appear in each record's `operations` array.\n\n| Mode                    | Behavior                                                                             | Use case                                  |\n| ----------------------- | ------------------------------------------------------------------------------------ | ----------------------------------------- |\n| `\"top-level\"` (default) | Only top-level operations (anything with a `parentId` is dropped)                    | Consistent, resume-independent snapshots  |\n| `\"full-tree\"`           | Every operation, including children of contexts (parallel branches, map items, etc.) | Full nested detail (see the caveat below) |\n\n`\"top-level\"` is the default because it produces the **same set of operations regardless of when a record is emitted**, so an execution that suspends and resumes never yields a partially-populated tree.\n\n> **⚠️ `\"full-tree\"` and suspend/resume:** by default the backend prunes a\n> finished context's children from the state handed to later invocations (a\n> performance optimization). So for an execution that suspends/resumes, a\n> `\"full-tree\"` record only contains the children of contexts that were still\n> active in the invocation that emitted it — children of already-finished\n> contexts are missing. To keep the full tree across resume, also set\n> `childOperationsDepth` (below).\n\n### `childOperationsDepth` (on `withDurableExecution`)\n\nTo make `\"full-tree\"` complete across suspend/resume, tell the **core SDK** to\npreserve child operations, via `pluginsConfig.childOperationsDepth` on\n`withDurableExecution` (not on the plugin):\n\n```typescript\nexport const handler = withDurableExecution(myWorkflow, {\n  plugins: [workflowInsight({ operationDetail: \"full-tree\", exporters: [...] })],\n  pluginsConfig: {\n    // Children of top-level contexts = 1; their children = 2; whole tree = Infinity.\n    childOperationsDepth: 1,\n  },\n});\n```\n\n| Value         | Preserved across resume                                      |\n| ------------- | ------------------------------------------------------------ |\n| omitted / `0` | Nothing extra (top-level only survives resume) — **default** |\n| `1`           | Direct children of top-level contexts (map items, branches)  |\n| `2`           | Their children too (e.g. steps inside each map item)         |\n| `Infinity`    | The entire tree                                              |\n\n> **⚠️ Cost:** preservation forces the SDK's `ReplayChildren` mode on each\n> preserved context. This keeps its children in the execution state and, on\n> resume, rebuilds the context's result by **replaying** the already-checkpointed\n> children — the context's orchestration code re-runs, but the children are\n> **not** re-executed (no step bodies or side effects run again). The cost is the\n> extra replay pass over each preserved context plus carrying its children in the\n> state, and it grows with the depth you request. Enable only the depth you need.\n\n### `exporters`\n\nAn array of destinations. Records are sent to **all** exporters in parallel. If one fails, others still receive the record. Exporter errors never fail the execution.\n\n```typescript\nexporters: [\n  new S3Exporter({ bucket: \"my-bucket\" }),\n  new DynamoDBExporter({ tableName: \"insight\" }),\n  new OTelExporter({ endpoint: \"https://otlp.vendor.com/v1/logs\" }),\n],\n```\n\n### `content` (advanced)\n\nControl what data is included in records. By default, execution input/output are\nincluded as-is, per-operation errors are included, and operation results are\n**not** included. Use `content` to redact/reshape input/output, opt specific\noperation results in (with optional transforms), exclude operations, or drop\noperation errors:\n\n```typescript\ncontent: {\n  // Include/exclude/transform execution input\n  input: true,                          // include as-is (default)\n  input: false,                         // exclude\n  input: (i) => ({ id: i.orderId }),    // transform (redact sensitive fields)\n\n  // Same for output\n  output: true,\n\n  // Operation-level control\n  operations: {\n    includeErrors: true,                // include per-operation errors (default)\n    overrides: [\n      { operationName: \"charge-payment\", result: (r) => ({ amount: r.amount }) },\n      { operationName: \"internal-log\", exclude: true },\n    ],\n  },\n},\n```\n\n> **Note:** operation overrides are matched by `operationName`. If multiple\n> override entries (or multiple operations) share a name, the last matching\n> entry wins.\n\n> [!IMPORTANT]\n> **Operation `result` reflects the checkpointed, serialized value — not\n> necessarily your original return value.** The plugin passes your `result`\n> transform the operation's checkpointed result, JSON-parsed when it is valid\n> JSON and otherwise the raw string. It does **not** run your SDK\n> `Serdes.deserialize`. If the operation uses a custom `Serdes` — one that\n> serializes to a non-JSON format (e.g. XML), encrypts, or offloads large values\n> to external storage and checkpoints only a **pointer/filepath** (overflow\n> mode) — your transform receives that serialized form or pointer, not the\n> original deserialized object. Only enable operation results for operations\n> using the default JSON serialization, or whose serialized form your transform\n> can handle. Input/output transforms are not affected by this.\n\n## Querying by operation name (`operationsByName`)\n\n`operations` is a canonical **array** — lossless (it keeps every occurrence of a\nrepeated step name), and queryable in the analytical stores that support nested\ndata (Athena `UNNEST`, Postgres/Redshift JSON path, OpenSearch `nested`):\n\n```sql\n-- Postgres (JSONB): executions where \"convert_data\" ran under 5s\nWHERE record_json @? '$.operations[*] ? (@.name == \"convert_data\" && @.durationMs < 5000)';\n```\n\nPoint-access stores that can't filter \"the array element named X\" —\n**CloudWatch Logs** (`LambdaLogExporter`, `CloudWatchLogsExporter`) and\n**DynamoDB** (`DynamoDBExporter`) — emit an `operationsByName` map **instead of\nthe `operations` array**, so name-based queries become a simple dot-path (these\nstores trade the per-occurrence array detail for queryability):\n\n```\n# CloudWatch Logs Insights\nfields executionArn | filter operationsByName.convert_data.maxDurationMs < 5000\n```\n\nEach entry aggregates metrics across all occurrences of the name; `result`/`error`\nare kept only when the name ran exactly once:\n\n```json\n\"operationsByName\": {\n  \"insert_to_db\": {\n    \"type\": \"STEP\", \"subType\": \"Step\",\n    \"count\": 1, \"minDurationMs\": 6200, \"maxDurationMs\": 6200, \"totalDurationMs\": 6200,\n    \"failedCount\": 0, \"maxAttempt\": 1,\n    \"status\": \"SUCCEEDED\",\n    \"result\": { \"rows\": 1200 }\n  }\n}\n```\n\nNotes:\n\n- **Metrics** (`count`, `min`/`max`/`totalDurationMs`, `failedCount`, `maxAttempt`)\n  span all occurrences; `type`/`subType`/`status` reflect the most recently seen\n  occurrence.\n- **`result`/`error` are included only when the name occurs exactly once.** For a\n  repeated name (loops/retries/map) they're dropped — there's no single\n  representative value — but `failedCount` still flags failures.\n- Operations **without a name are excluded** (they can't be keyed or queried).\n- **Choosing store-safe operation names is your responsibility.** Names are used\n  verbatim as keys/identifiers — the library never sanitizes or escapes them. Any\n  character your target store treats specially (e.g. `.` in CloudWatch Logs\n  Insights / OpenSearch field paths, reserved or quoting characters in\n  DynamoDB attribute names and SQL identifiers, etc.) can make an operation hard\n  or impossible to query there. Stick to simple, portable names — letters,\n  digits, `-`, `_` — to stay safe across destinations.\n- The array remains the source of truth; array-native exporters (S3/Athena,\n  OpenSearch, Aurora, Redshift) emit only the array. See\n  [`docs/operations-shape.md`](./docs/operations-shape.md).\n\n## Record size & truncation\n\nDestinations cap payload size (CloudWatch Logs events at 256 KB, DynamoDB items\nat 400 KB, etc.). Each exporter carries a `maxRecordSizeBytes`; when a record's\nserialized JSON exceeds it, the plugin truncates a **per-exporter copy**\n(best-effort) before sending — the same record can go out full to one exporter\nand trimmed to another.\n\nDrop order, until the record fits:\n\n1. operation `result` fields, **oldest operation first**;\n2. whole operations, **oldest first**;\n3. as a last resort (once every operation is gone), execution `input`, then\n   `output`.\n\nIdentity/timeline fields are **never** dropped, and `input`/`output` are dropped\nonly after all operations are gone — so prefer `content.input` /\n`content.output` transforms to bound them earlier (those run before truncation).\nWhen anything is dropped, the emitted record carries `truncated: true`; each\noperation whose result was dropped is itself marked `truncated: true`, and\n`droppedOperations` (count), `droppedInput` / `droppedOutput` flags are set as\napplicable — so a trimmed record is always distinguishable from a complete one\n(\"cut, not missing\").\n\nThe size check measures the **exact shape each exporter emits**, not just the\ncanonical record. Exporters that reshape operations (the `operationsByName`\nexpansion used by CloudWatch Logs / DynamoDB / Lambda log, or the `\"both\"`\nformat) expose a `render` the limiter sizes against, so a record trimmed to its\nlimit reflects what is actually serialized. This bounds the serialized record\n_body_; it does not model the destination wire envelope (DynamoDB type\ndescriptors, CloudWatch Logs event framing, gzip, etc.) — which is why the\nfirst-party defaults below sit under each destination's hard limit to leave\nheadroom.\n\nPer-exporter defaults (override via each exporter's `maxRecordSizeBytes`):\n\n| Exporter(s)                                   | Default |\n| --------------------------------------------- | ------- |\n| Lambda log, CloudWatch Logs, SQS, EventBridge | 256 KB  |\n| DynamoDB                                      | 400 KB  |\n| Aurora, Redshift, Firehose, OTel              | 1 MB    |\n| S3                                            | 5 MB    |\n| OpenSearch                                    | 10 MB   |\n| HTTP, File                                    | none¹   |\n\n¹ No default — truncation is disabled unless you set `maxRecordSizeBytes`.\n\n## Exporters\n\n### Comparison Table\n\n| Exporter                 | Destination                     | Upsert        | Query Method     | Setup              | Best For                            |\n| ------------------------ | ------------------------------- | ------------- | ---------------- | ------------------ | ----------------------------------- |\n| `LambdaLogExporter`      | Function's CloudWatch log group | No            | Logs Insights    | None               | Getting started, zero-config        |\n| `CloudWatchLogsExporter` | Any CloudWatch log group        | No            | Logs Insights    | Log group + IAM    | Centralized logging, cross-function |\n| `S3Exporter`             | Amazon S3                       | Yes (by key)  | Athena SQL       | Bucket             | Analytics, long-term retention      |\n| `DynamoDBExporter`       | Amazon DynamoDB                 | Configurable  | GetItem / Query  | Table              | Fast lookups by execution           |\n| `AuroraExporter`         | Aurora MySQL/PostgreSQL         | Yes (UPSERT)  | SQL              | Cluster + Data API | Relational queries, joins           |\n| `RedshiftExporter`       | Amazon Redshift                 | Yes (MERGE)   | SQL              | Cluster/Serverless | Large-scale analytics               |\n| `OpenSearchExporter`     | Amazon OpenSearch               | Yes (by \\_id) | Full-text / DSL  | Domain             | Search, dashboards                  |\n| `FirehoseExporter`       | Kinesis Firehose → anywhere     | N/A           | Depends on dest. | Delivery stream    | Fan-out to S3/Redshift/Splunk       |\n| `EventBridgeExporter`    | Amazon EventBridge              | N/A           | Rules/patterns   | Event bus          | Event-driven reactions              |\n| `SQSExporter`            | Amazon SQS                      | N/A           | Consumer         | Queue              | Decoupled processing                |\n| `OTelExporter`           | Any OTLP backend                | N/A           | Varies           | Endpoint           | Third-party observability           |\n| `HttpExporter`           | Any HTTP endpoint               | N/A           | N/A              | URL                | Custom backends, webhooks           |\n| `FileExporter`           | Filesystem (EFS/mount)          | Configurable  | File read        | Directory          | Local dev, EFS persistence          |\n\n---\n\n### Destination Setup\n\nBefore using an exporter, you need to create the target resource and grant your Lambda function the required permissions. Below is what each exporter needs.\n\n#### LambdaLogExporter\n\n**No setup required.** Uses the function's own CloudWatch log group (automatically created by Lambda).\n\n#### CloudWatchLogsExporter\n\n**Resource:** A CloudWatch log group.\n\n```bash\naws logs create-log-group --log-group-name /custom/workflow-insight\naws logs put-retention-policy --log-group-name /custom/workflow-insight --retention-in-days 30\n```\n\n**IAM policy:**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": [\"logs:CreateLogStream\", \"logs:PutLogEvents\"],\n  \"Resource\": \"arn:aws:logs:*:*:log-group:/custom/workflow-insight:*\"\n}\n```\n\n#### S3Exporter\n\n**Resource:** An S3 bucket.\n\n```bash\naws s3 mb s3://my-insight-bucket --region us-east-1\n```\n\n**IAM policy:**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": \"s3:PutObject\",\n  \"Resource\": \"arn:aws:s3:::my-insight-bucket/workflow-insight/*\"\n}\n```\n\n**Optional (for Athena queries):** Create a Glue table or run a crawler over the prefix. The Hive-style partitioning (`year=YYYY/month=MM/day=DD/`) is auto-discovered by crawlers and `MSCK REPAIR TABLE`.\n\n#### DynamoDBExporter\n\n**Resource:** A DynamoDB table with the partition key (and optional sort key) matching your config.\n\n```bash\n# With sort key (full history per execution)\naws dynamodb create-table \\\n  --table-name workflow-insight \\\n  --attribute-definitions \\\n    AttributeName=pk,AttributeType=S \\\n    AttributeName=sk,AttributeType=S \\\n  --key-schema \\\n    AttributeName=pk,KeyType=HASH \\\n    AttributeName=sk,KeyType=RANGE \\\n  --billing-mode PAY_PER_REQUEST\n\n# Without sort key (upsert only)\naws dynamodb create-table \\\n  --table-name workflow-insight \\\n  --attribute-definitions AttributeName=pk,AttributeType=S \\\n  --key-schema AttributeName=pk,KeyType=HASH \\\n  --billing-mode PAY_PER_REQUEST\n```\n\n**IAM policy:**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": \"dynamodb:PutItem\",\n  \"Resource\": \"arn:aws:dynamodb:*:*:table/workflow-insight\"\n}\n```\n\n#### AuroraExporter\n\n**Resource:** An Aurora cluster with the **Data API enabled** and a Secrets Manager secret for credentials.\n\n**Table creation (PostgreSQL):**\n\n```sql\nCREATE TABLE workflow_insight (\n  execution_arn VARCHAR(512) PRIMARY KEY,\n  execution_name VARCHAR(256),\n  function_name VARCHAR(128),\n  status VARCHAR(20),\n  start_time VARCHAR(30),\n  end_time VARCHAR(30),\n  duration_ms BIGINT,\n  record_json TEXT,\n  emitted_at VARCHAR(30)\n);\n```\n\n**Table creation (MySQL):**\n\n```sql\nCREATE TABLE workflow_insight (\n  execution_arn VARCHAR(512) PRIMARY KEY,\n  execution_name VARCHAR(256),\n  function_name VARCHAR(128),\n  status VARCHAR(20),\n  start_time VARCHAR(30),\n  end_time VARCHAR(30),\n  duration_ms BIGINT,\n  record_json LONGTEXT,\n  emitted_at VARCHAR(30)\n);\n```\n\n**IAM policy:**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": [\"rds-data:ExecuteStatement\", \"secretsmanager:GetSecretValue\"],\n  \"Resource\": [\n    \"arn:aws:rds:*:*:cluster:my-cluster\",\n    \"arn:aws:secretsmanager:*:*:secret:my-db-creds-*\"\n  ]\n}\n```\n\n**Network:** No VPC required — the Data API is an HTTP endpoint. Your Lambda does NOT need to be in the Aurora VPC.\n\n#### RedshiftExporter\n\n**Resource:** A Redshift Serverless workgroup or provisioned cluster.\n\n**Table creation:**\n\n```sql\nCREATE TABLE public.workflow_insight (\n  execution_arn VARCHAR(512) PRIMARY KEY,\n  execution_name VARCHAR(256),\n  function_name VARCHAR(128),\n  status VARCHAR(20),\n  start_time VARCHAR(30),\n  end_time VARCHAR(30),\n  duration_ms BIGINT,\n  record_json VARCHAR(MAX),\n  emitted_at VARCHAR(30)\n);\n```\n\n**IAM policy (Serverless):**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": [\n    \"redshift-data:ExecuteStatement\",\n    \"redshift-serverless:GetCredentials\"\n  ],\n  \"Resource\": \"*\"\n}\n```\n\n**IAM policy (Provisioned with Secrets Manager):**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": [\"redshift-data:ExecuteStatement\", \"secretsmanager:GetSecretValue\"],\n  \"Resource\": [\n    \"arn:aws:redshift:*:*:cluster:my-cluster\",\n    \"arn:aws:secretsmanager:*:*:secret:redshift-creds-*\"\n  ]\n}\n```\n\n**Network:** No VPC required — the Redshift Data API is HTTP-based.\n\n#### OpenSearchExporter\n\n**Resource:** An Amazon OpenSearch Service domain (or self-managed cluster).\n\nThe index is auto-created on first write. No manual index creation needed (unless you want a custom mapping).\n\n**IAM policy (SigV4 auth):**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": \"es:ESHttpPut\",\n  \"Resource\": \"arn:aws:es:*:*:domain/my-domain/workflow-insight/*\"\n}\n```\n\n**Network:** If the OpenSearch domain is in a VPC, your Lambda must be in the same VPC (or a peered one) with a security group allowing HTTPS to the domain.\n\n#### FirehoseExporter\n\n**Resource:** A Kinesis Data Firehose delivery stream configured with your target destination (S3, Redshift, Splunk, HTTP endpoint, etc.).\n\n```bash\naws firehose create-delivery-stream \\\n  --delivery-stream-name workflow-insight-stream \\\n  --s3-destination-configuration \\\n    RoleARN=arn:aws:iam::123456789012:role/firehose-role,\\\n    BucketARN=arn:aws:s3:::my-bucket,\\\n    Prefix=workflow-insight/\n```\n\n**IAM policy:**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": \"firehose:PutRecord\",\n  \"Resource\": \"arn:aws:firehose:*:*:deliverystream/workflow-insight-stream\"\n}\n```\n\n#### EventBridgeExporter\n\n**Resource:** An EventBridge event bus (or use the `default` bus — no creation needed).\n\n```bash\n# Only if using a custom bus:\naws events create-event-bus --name workflow-insight-bus\n```\n\n**IAM policy:**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": \"events:PutEvents\",\n  \"Resource\": \"arn:aws:events:*:*:event-bus/default\"\n}\n```\n\n**Example rule** (trigger SNS on failure):\n\n```bash\naws events put-rule --name insight-failures \\\n  --event-pattern '{\"source\":[\"aws.durable-execution.insight\"],\"detail-type\":[\"FAILED\"]}'\naws events put-targets --rule insight-failures --targets Id=1,Arn=arn:aws:sns:...\n```\n\n#### SQSExporter\n\n**Resource:** An SQS queue (standard or FIFO).\n\n```bash\n# Standard queue\naws sqs create-queue --queue-name workflow-insight-queue\n\n# FIFO queue (content-based dedup enabled — exporter provides dedup IDs)\naws sqs create-queue --queue-name workflow-insight-queue.fifo \\\n  --attributes FifoQueue=true,ContentBasedDeduplication=false\n```\n\n**IAM policy:**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": \"sqs:SendMessage\",\n  \"Resource\": \"arn:aws:sqs:*:*:workflow-insight-queue*\"\n}\n```\n\n#### OTelExporter\n\n**Resource:** An OTLP-compatible endpoint (Datadog, Grafana Cloud, Splunk, New Relic, etc.). Setup varies by vendor.\n\n**No IAM policy needed** — uses HTTPS to an external endpoint. Authentication is via headers (API key, bearer token) configured in the exporter.\n\n**Network:** If the endpoint is external (internet), Lambda needs outbound internet access (default for non-VPC Lambdas, or NAT Gateway for VPC-attached Lambdas).\n\n#### HttpExporter\n\n**Resource:** Any HTTP endpoint that accepts JSON POSTs.\n\n**No IAM policy needed.** Authentication is via custom headers.\n\n**Network:** Same as OTelExporter — needs outbound access to the endpoint.\n\n#### FileExporter\n\n**Resource:** A writable directory. For Lambda:\n\n- **EFS mount:** Attach an EFS file system to your Lambda function via an access point.\n- **S3 File Gateway:** Mount an S3-backed NFS share.\n- **`/tmp`:** Ephemeral (512MB–10GB), lost between invocations. Only for testing.\n\n**EFS setup (summary):**\n\n1. Create an EFS file system and access point\n2. Add VPC config to your Lambda (same VPC as EFS)\n3. Configure the file system in the Lambda function:\n   ```bash\n   aws lambda update-function-configuration \\\n     --function-name my-fn \\\n     --file-system-configs Arn=arn:aws:elasticfilesystem:...:access-point/fsap-...,LocalMountPath=/mnt/efs\n   ```\n\n**IAM policy (EFS):**\n\n```json\n{\n  \"Effect\": \"Allow\",\n  \"Action\": [\"elasticfilesystem:ClientMount\", \"elasticfilesystem:ClientWrite\"],\n  \"Resource\": \"arn:aws:elasticfilesystem:*:*:access-point/fsap-*\"\n}\n```\n\n**Network:** Lambda must be in the same VPC as the EFS file system.\n\n---\n\n### CDK Infrastructure (Automated Setup)\n\n> **Note:** This CDK stack is designed for **getting started and testing**. It uses minimal configurations (single-AZ, no backups, `RemovalPolicy.DESTROY`) to keep costs low and teardown easy. For production, build your own infrastructure with proper security, redundancy, monitoring, and cost controls tailored to your workload.\n\nInstead of creating resources manually, you can use the included CDK stack to deploy all destination infrastructure, IAM permissions, and an example Lambda function with a single command.\n\n**Location:** `cdk/` directory within this package.\n\n#### Quick Start\n\n```bash\ncd packages/aws-durable-execution-sdk-js-insight/cdk\nnpm install\nnpx cdk deploy\n```\n\n> The CDK package depends on the `@aws/durable-execution-sdk-js` and\n> `@aws/durable-execution-sdk-js-insight` workspace packages (the example\n> Lambda imports them). The `deploy`, `synth`, `test`, and `typecheck` scripts\n> automatically build these dependencies first via a `build:deps` pre-step, so\n> no manual ordering is required. (If you prefer to build manually, run\n> `npm run build:deps`.)\n\n#### Configuration (`cdk/config.json`)\n\nEdit `config.json` to enable/disable destinations and configure settings:\n\n```json\n{\n  \"destinations\": {\n    \"cloudwatchLogs\": {\n      \"enabled\": true,\n      \"logGroupName\": \"/workflow-insight/demo\",\n      \"retentionDays\": 30\n    },\n    \"dynamodb\": { \"enabled\": true, \"tableName\": \"workflow-insight\" },\n    \"aurora\": {\n      \"enabled\": true,\n      \"tableName\": \"workflow_insight\",\n      \"databaseName\": \"postgres\",\n      \"minCapacity\": 0.5,\n      \"maxCapacity\": 1\n    },\n    \"s3\": { \"enabled\": false, \"bucketName\": \"workflow-insight-records\" },\n    \"redshift\": {\n      \"enabled\": false,\n      \"namespaceName\": \"insight-namespace\",\n      \"workgroupName\": \"insight-workgroup\",\n      \"databaseName\": \"dev\",\n      \"tableName\": \"workflow_insight\",\n      \"schema\": \"public\"\n    },\n    \"opensearch\": { \"enabled\": false, \"domainName\": \"workflow-insight\" },\n    \"firehose\": {\n      \"enabled\": false,\n      \"streamName\": \"workflow-insight\",\n      \"bufferIntervalSeconds\": 60,\n      \"bufferSizeMB\": 1\n    },\n    \"sqs\": { \"enabled\": false, \"queueName\": \"workflow-insight\", \"fifo\": false },\n    \"eventbridge\": { \"enabled\": false, \"eventBusName\": \"default\" }\n  },\n  \"lambda\": {\n    \"roleNames\": [],\n    \"discoverDurableFunctions\": false,\n    \"createExampleFunction\": true,\n    \"autoInvoke\": { \"enabled\": false, \"rateMinutes\": 5 }\n  }\n}\n```\n\n#### Lambda Settings\n\n| Setting                    | Description                                                                                                                                                                                                                                                                                                                                                                                                                                                                   |\n| -------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |\n| `roleNames`                | Array of **existing** IAM role names to attach the Insight permissions policy to. Default: `[]` (empty). The example function's role is added automatically when `createExampleFunction` is `true`, so you can leave this empty to start. Add your own durable function roles here to grant them access.                                                                                                                                                                      |\n| `discoverDurableFunctions` | When `true`, the CDK app lists all Lambda functions **at synth time** and identifies durable functions by the presence of `DurableConfig` (the native Lambda API field). The Insight permissions policy is attached directly to their exact execution-role ARNs — no runtime Lambda or `iam:PutRolePolicy` grant is used. Requires AWS credentials at synth time (same model as CDK context lookups), and `cdk synth`/`deploy` will reflect the account state at that moment. |\n| `createExampleFunction`    | Deploys an insurance claim processing workflow (with retry policies and random transient failures) pre-configured with the Insight plugin pointing to all enabled destinations                                                                                                                                                                                                                                                                                                |\n| `autoInvoke.enabled`       | Deploys a dispatcher Lambda + EventBridge rule that invokes the example function on a schedule with randomized input. **Default: `false`.** Enabling creates ongoing costs (Lambda invocations, Aurora ACU time, DynamoDB writes). Disable or destroy the stack when not actively testing.                                                                                                                                                                                    |\n| `autoInvoke.rateMinutes`   | How often the dispatcher triggers (default: 5 minutes)                                                                                                                                                                                                                                                                                                                                                                                                                        |\n\n#### What Gets Deployed\n\nWhen `createExampleFunction` and `autoInvoke` are both enabled, the stack creates:\n\n1. **Destination resources** — tables, clusters, queues, etc. for each enabled destination\n2. **IAM policy** — a single managed policy (`WorkflowInsightDestinations`) attached to all target roles\n3. **Example function** (`insight-example-workflow`) — a 6-step insurance claim workflow with:\n   - Retry policy (4 attempts, 1s initial backoff, coefficient 2, max 5s)\n   - ~30% random transient failure rate per step (generates retry attempt data)\n   - Three possible outcomes: APPROVED, REJECTED, or MORE_DOCUMENTS_REQUIRED\n4. **Dispatcher** (`insight-example-dispatcher`) — generates random claims with: `customerName`, `insuranceClaimNumber`, `claimAmount`, `claimType`\n5. **EventBridge rule** (`insight-example-dispatch`) — triggers the dispatcher every N minutes\n\nAfter deployment, data starts flowing automatically to all enabled destinations.\n\n#### Aurora Table Auto-Creation\n\nWhen Aurora is enabled, the stack deploys a custom resource that creates the `workflow_insight` table automatically via the RDS Data API. No manual SQL needed.\n\n#### Tear Down\n\n```bash\nnpx cdk destroy\n```\n\nThis removes all created resources (tables, clusters, functions, rules). Resources use `RemovalPolicy.DESTROY` by default for clean teardown.\n\n---\n\n### Exporter Usage & Examples\n\nDetailed configuration, code examples, and use-case guidance for each exporter.\n\n### LambdaLogExporter (default)\n\nWrites records to stdout via `console.log`. Since Lambda captures stdout to the function's CloudWatch log group, this requires **zero IAM permissions** and **zero setup**.\n\n```typescript\nimport {\n  workflowInsight,\n  LambdaLogExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [new LambdaLogExporter()],\n});\n```\n\n**When to use:** Getting started, quick debugging, or when you already query CloudWatch Logs. This is the default if you don't specify `exporters`.\n\n---\n\n### CloudWatchLogsExporter\n\nWrites records to a **specific** CloudWatch log group via `PutLogEvents`. Unlike `LambdaLogExporter`, you control the destination and can centralize records from multiple functions.\n\n```typescript\nimport {\n  workflowInsight,\n  CloudWatchLogsExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new CloudWatchLogsExporter({\n      logGroupName: \"/custom/workflow-insight\",\n      logStreamPrefix: \"workflow-insight/\", // optional (default)\n      region: \"us-east-1\", // optional\n    }),\n  ],\n});\n```\n\n**Log stream pattern:** `{prefix}{YYYY}/{MM}/{DD}` (one per day, auto-created).\n\n**IAM required:** `logs:CreateLogStream`, `logs:PutLogEvents` on the target log group.\n\n**When to use:** Centralizing insight data from multiple functions into one queryable log group.\n\n---\n\n### S3Exporter\n\nWrites records as JSON objects to Amazon S3 with Hive-style partitioning for Amazon Athena compatibility.\n\n```typescript\nimport {\n  workflowInsight,\n  S3Exporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new S3Exporter({\n      bucket: \"my-insight-bucket\",\n      prefix: \"workflow-insight/\", // optional (default)\n      partitioning: \"date\", // \"date\" | \"function-name\" | \"none\" (default: \"date\")\n      region: \"us-east-1\", // optional\n    }),\n  ],\n});\n```\n\n**S3 key pattern (date):** `workflow-insight/year=2026/month=06/day=16/{executionName}.json`\n\n**Upsert behavior:** Uses `executionName` as the object key — subsequent updates to the same execution overwrite the same file.\n\n**IAM required:** `s3:PutObject` on the bucket/prefix.\n\n**When to use:** Long-term retention, Athena queries, or feeding data into a data lake.\n\n---\n\n### DynamoDBExporter\n\nWrites records to DynamoDB via `PutItem`.\n\n```typescript\nimport {\n  workflowInsight,\n  DynamoDBExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new DynamoDBExporter({\n      tableName: \"workflow-insight\",\n      partitionKey: \"pk\", // optional (default: \"pk\"), value = executionArn\n      sortKey: \"sk\", // optional (default: \"sk\"), value = emittedAt\n      region: \"us-east-1\", // optional\n    }),\n  ],\n});\n```\n\n**Upsert behavior:**\n\n- With sort key (default): each emission creates a new item → full history per execution\n- Without sort key (`sortKey: undefined`): PutItem overwrites → only latest state kept\n\n**IAM required:** `dynamodb:PutItem` on the table.\n\n**When to use:** Fast lookups by execution ARN, or building real-time dashboards backed by DynamoDB.\n\n---\n\n### AuroraExporter\n\nWrites records to Aurora MySQL or PostgreSQL via the **RDS Data API** (no VPC needed).\n\n```typescript\nimport {\n  workflowInsight,\n  AuroraExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new AuroraExporter({\n      resourceArn: \"arn:aws:rds:us-east-1:123456789012:cluster:my-cluster\",\n      secretArn:\n        \"arn:aws:secretsmanager:us-east-1:123456789012:secret:db-creds\",\n      database: \"workflows\",\n      table: \"workflow_insight\", // optional (default)\n      engine: \"postgresql\", // \"postgresql\" | \"mysql\"\n      region: \"us-east-1\", // optional\n    }),\n  ],\n});\n```\n\n**Upsert behavior:** `INSERT ... ON CONFLICT DO UPDATE` (PostgreSQL) or `INSERT ... ON DUPLICATE KEY UPDATE` (MySQL) by `execution_arn`.\n\n**IAM required:** `rds-data:ExecuteStatement`, `secretsmanager:GetSecretValue`.\n\n**Table schema:**\n\n```sql\nCREATE TABLE workflow_insight (\n  execution_arn VARCHAR(512) PRIMARY KEY,\n  execution_name VARCHAR(256),\n  function_name VARCHAR(128),\n  status VARCHAR(20),\n  start_time VARCHAR(30),\n  end_time VARCHAR(30),\n  duration_ms BIGINT,\n  record_json TEXT,\n  emitted_at VARCHAR(30)\n);\n```\n\n**When to use:** Relational queries, joins with other business data, or teams already on Aurora.\n\n---\n\n### RedshiftExporter\n\nWrites records to Amazon Redshift via the **Redshift Data API**.\n\n```typescript\nimport {\n  workflowInsight,\n  RedshiftExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\n// Redshift Serverless\nconst insight = workflowInsight({\n  exporters: [\n    new RedshiftExporter({\n      workgroupName: \"my-workgroup\",\n      database: \"workflows\",\n      table: \"workflow_insight\", // optional (default)\n      schema: \"public\", // optional (default)\n      region: \"us-east-1\", // optional\n    }),\n  ],\n});\n\n// Provisioned cluster\nconst insight2 = workflowInsight({\n  exporters: [\n    new RedshiftExporter({\n      clusterIdentifier: \"my-cluster\",\n      database: \"workflows\",\n      secretArn: \"arn:aws:secretsmanager:...:secret:redshift-creds\",\n    }),\n  ],\n});\n```\n\n**Upsert behavior:** Uses `MERGE` to insert or update by `execution_arn`.\n\n**IAM required:** `redshift-data:ExecuteStatement` (+ `redshift-serverless:GetCredentials` or `redshift:GetClusterCredentialsWithIAM`).\n\n**When to use:** Large-scale cross-function analytics, BI dashboards, data warehouse integration.\n\n**Note:** The `record_json` column uses Redshift's SUPER type with `JSON_PARSE()` for navigable nested queries. SUPER values are limited to ~1 MB — executions with unusually large input/output payloads may fail the insert.\n\n---\n\n### OpenSearchExporter\n\nIndexes records to Amazon OpenSearch Service for full-text search and dashboards.\n\n```typescript\nimport {\n  workflowInsight,\n  OpenSearchExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\n// Amazon OpenSearch Service (IAM auth)\nconst insight = workflowInsight({\n  exporters: [\n    new OpenSearchExporter({\n      endpoint: \"https://my-domain.us-east-1.es.amazonaws.com\",\n      indexName: \"workflow-insight\", // optional (default)\n      region: \"us-east-1\",\n      auth: \"sigv4\", // optional (default)\n    }),\n  ],\n});\n\n// Basic auth (self-managed)\nconst insight2 = workflowInsight({\n  exporters: [\n    new OpenSearchExporter({\n      endpoint: \"https://opensearch.internal:9200\",\n      region: \"us-east-1\",\n      auth: \"basic\",\n      username: \"admin\",\n      password: process.env.OS_PASSWORD!,\n    }),\n  ],\n});\n```\n\n**Upsert behavior:** Document `_id` = `executionArn` — updates overwrite.\n\n**IAM required (SigV4):** `es:ESHttpPut` on the domain.\n\n**When to use:** Full-text search across executions, OpenSearch Dashboards/Kibana visualizations, complex filtering.\n\n---\n\n### FirehoseExporter\n\nSends records to Amazon Kinesis Data Firehose for delivery to S3, Redshift, Splunk, or any HTTP endpoint.\n\n```typescript\nimport {\n  workflowInsight,\n  FirehoseExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new FirehoseExporter({\n      deliveryStreamName: \"workflow-insight-stream\",\n      region: \"us-east-1\", // optional\n      operationsFormat: \"array\", // optional: \"array\" (default) | \"by-name\" | \"both\"\n    }),\n  ],\n});\n```\n\n**Format:** Newline-delimited JSON (NDJSON) — when Firehose batches records into S3 objects, they remain parseable.\n\n**IAM required:** `firehose:PutRecord` on the delivery stream.\n\n**When to use:** Fan-out to multiple destinations via one Firehose stream, buffered S3 delivery, or integration with Splunk/Datadog via Firehose HTTP endpoints.\n\n---\n\n### EventBridgeExporter\n\nPublishes records to Amazon EventBridge for event-driven reactions.\n\n```typescript\nimport {\n  workflowInsight,\n  EventBridgeExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new EventBridgeExporter({\n      eventBusName: \"default\", // optional (default)\n      source: \"aws.durable-execution.insight\", // optional (default)\n      region: \"us-east-1\", // optional\n      operationsFormat: \"array\", // optional: \"array\" (default) | \"by-name\" | \"both\"\n    }),\n  ],\n});\n```\n\n**Event structure:**\n\n- Source: `aws.durable-execution.insight`\n- DetailType: `SUCCEEDED` | `FAILED` | `RUNNING`\n- Detail: full record JSON\n\n**Example rule pattern:**\n\n```json\n{ \"source\": [\"aws.durable-execution.insight\"], \"detail-type\": [\"FAILED\"] }\n```\n\n**IAM required:** `events:PutEvents` on the event bus.\n\n**When to use:** Triggering notifications on failure, starting remediation workflows, fan-out to multiple consumers without coupling.\n\n---\n\n### SQSExporter\n\nSends records to Amazon SQS (standard or FIFO).\n\n```typescript\nimport {\n  workflowInsight,\n  SQSExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new SQSExporter({\n      queueUrl:\n        \"https://sqs.us-east-1.amazonaws.com/123456789012/insight-queue.fifo\",\n      messageGroupId: undefined, // optional (default: executionArn)\n      region: \"us-east-1\", // optional\n      operationsFormat: \"array\", // optional: \"array\" (default) | \"by-name\" | \"both\"\n    }),\n  ],\n});\n```\n\n**FIFO support:** Auto-detected from `.fifo` URL suffix. Dedup ID = `executionArn:emittedAt`. Message attributes include `status` and `functionName` for SQS message filtering.\n\n**IAM required:** `sqs:SendMessage` on the queue.\n\n**When to use:** Guaranteed single-consumer delivery, decoupling export processing from the Lambda invocation, ordered processing (FIFO).\n\n---\n\n### OTelExporter\n\nEmits records as OpenTelemetry log records via OTLP HTTP/JSON. Compatible with any OTLP backend (Datadog, Grafana, Splunk, New Relic, Honeycomb, etc.).\n\n```typescript\nimport {\n  workflowInsight,\n  OTelExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new OTelExporter({\n      endpoint: \"https://otlp.datadoghq.com/v1/logs\",\n      headers: { \"DD-API-KEY\": process.env.DD_API_KEY! },\n      protocol: \"http/json\", // optional (default)\n      operationsFormat: \"array\", // optional: \"array\" (default) | \"by-name\" | \"both\"\n    }),\n  ],\n});\n```\n\n**OTel mapping:**\n\n- Resource: `service.name`, `cloud.region`, `cloud.account.id`, `faas.name`, `faas.version`\n- Log attributes: `workflow.execution_arn`, `workflow.status`, `workflow.duration_ms`\n- Log body: full record JSON. `operationsFormat` controls how operations appear\n  in the body — the `operations` array (default), the `operationsByName` map, or\n  both. (Operations are only in the body, never attributes, so this never affects\n  attribute cardinality.)\n- Severity: ERROR for FAILED, INFO otherwise\n\n**No dependencies** — uses native `fetch`.\n\n**When to use:** Sending data to third-party observability platforms that support OTLP, unified observability across services.\n\n---\n\n### HttpExporter\n\nGeneric HTTP/Webhook exporter. POSTs (or PUTs) the full record as JSON to any URL.\n\n```typescript\nimport {\n  workflowInsight,\n  HttpExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nconst insight = workflowInsight({\n  exporters: [\n    new HttpExporter({\n      url: \"https://my-service.example.com/ingest\",\n      headers: { Authorization: \"Bearer \" + process.env.TOKEN },\n      method: \"POST\", // optional (default)\n      timeoutMs: 5000, // optional (default: 10000)\n      operationsFormat: \"array\", // optional: \"array\" (default) | \"by-name\" | \"both\"\n    }),\n  ],\n});\n```\n\n**No dependencies** — uses native `fetch` with configurable timeout.\n\n`operationsFormat` controls how operations are rendered in the posted body:\nthe canonical `operations` array (default), the name-keyed `operationsByName`\nmap, or both — pick what your endpoint consumes.\n\n**When to use:** Custom backends, internal microservices, SaaS integrations without dedicated exporters, prototyping.\n\n---\n\n### FileExporter\n\nWrites records to the filesystem — EFS mounts, S3 File Gateway, or any writable path.\n\n```typescript\nimport {\n  workflowInsight,\n  FileExporter,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\n// Append to daily NDJSON files on EFS\nconst insight = workflowInsight({\n  exporters: [\n    new FileExporter({\n      directory: \"/mnt/efs/workflow-insight\",\n      mode: \"ndjson\", // optional (default)\n      operationsFormat: \"array\", // optional: \"array\" (default) | \"by-name\" | \"both\"\n    }),\n  ],\n});\n\n// One JSON file per execution (upsert)\nconst insight2 = workflowInsight({\n  exporters: [\n    new FileExporter({\n      directory: \"/mnt/efs/workflow-insight\",\n      mode: \"json\",\n    }),\n  ],\n});\n```\n\n**Modes:**\n\n- `\"ndjson\"` (default): appends to `{YYYY-MM-DD}.ndjson` — good for bulk processing\n- `\"json\"`: one file per execution (`{executionName}.json`) — overwrites on update\n\n**No AWS SDK dependencies** — uses `node:fs/promises`.\n\n**When to use:** Lambda with EFS mount, local development/testing, S3 File Gateway integration.\n\n---\n\n## Multiple Exporters\n\nYou can use multiple exporters simultaneously. Records are sent to all of them in parallel:\n\n```typescript\nconst insight = workflowInsight({\n  exporters: [\n    new S3Exporter({ bucket: \"insight-archive\" }), // Long-term storage\n    new DynamoDBExporter({ tableName: \"insight-live\" }), // Fast lookups\n    new EventBridgeExporter({}), // Trigger alerts\n  ],\n  emitMode: \"on-change\",\n});\n```\n\n## Custom Exporters\n\nImplement the `InsightExporter` interface to build your own:\n\n```typescript\nimport {\n  InsightExporter,\n  WorkflowInsightRecord,\n} from \"@aws/durable-execution-sdk-js-insight\";\n\nclass MyExporter implements InsightExporter {\n  // Optional: opt into size-based truncation. When set, the plugin sends this\n  // exporter a copy trimmed to fit (drops operation results, then whole\n  // operations oldest-first, then execution input/output as a last resort) and\n  // sets `truncated: true` on it. Omit to receive full records.\n  readonly maxRecordSizeBytes = 256_000;\n\n  async export(record: WorkflowInsightRecord): Promise<void> {\n    // Send record wherever you want\n  }\n\n  async flush?(): Promise<void> {\n    // Optional: flush any buffered data before the invocation returns\n  }\n}\n```\n\n## Handling Backend-Initiated Events (STOPPED, TIMED_OUT)\n\nThe Workflow Insight plugin runs **inside your Lambda function**. This means it can only emit records when the function is invoked. Events that originate from the backend — such as `STOPPED` (manual stop) or `TIMED_OUT` (execution timeout exceeded) — happen **without a Lambda invocation**, so the plugin cannot capture them.\n\nTo get complete lifecycle coverage, subscribe to Lambda's durable execution lifecycle events via **Amazon EventBridge** and update your destination accordingly.\n\n### Step 1: Create a handler that updates your destination\n\n```typescript\n// lifecycle-handler.ts\nimport { DynamoDBClient, UpdateItemCommand } from \"@aws-sdk/client-dynamodb\";\n\nconst ddb = new DynamoDBClient({});\nconst TABLE = process.env.TABLE_NAME!;\n\nexport const handler = async (event: {\n  detail: {\n    executionArn: string;\n    status: string; // \"STOPPED\" | \"TIMED_OUT\"\n    timestamp: string;\n  };\n}) => {\n  const { executionArn, status, timestamp } = event.detail;\n\n  // Update the existing insight record in DynamoDB with the terminal status\n  await ddb.send(\n    new UpdateItemCommand({\n      TableName: TABLE,\n      Key: { pk: { S: executionArn } },\n      UpdateExpression: \"SET #s = :status, end_time = :ts\",\n      ExpressionAttributeNames: { \"#s\": \"status\" },\n      ExpressionAttributeValues: {\n        \":status\": { S: status },\n        \":ts\": { S: timestamp },\n      },\n    }),\n  );\n};\n```\n\n### Step 2: Create the EventBridge rule\n\nLambda emits durable execution lifecycle events to the **default** event bus:\n\n```bash\naws events put-rule \\\n  --name durable-execution-terminal-events \\\n  --event-pattern '{\n    \"source\": [\"aws.lambda\"],\n    \"detail-type\": [\"Lambda Durable Execution State Change\"],\n    \"detail\": {\n      \"status\": [\"STOPPED\", \"TIMED_OUT\"]\n    }\n  }'\n```\n\n### Step 3: Point the rule at your handler\n\n```bash\naws events put-targets \\\n  --rule durable-execution-terminal-events \\\n  --targets Id=1,Arn=arn:aws:lambda:us-east-1:123456789012:function:lifecycle-handler\n```\n\nGrant EventBridge permission to invoke it:\n\n```bash\naws lambda add-permission \\\n  --function-name lifecycle-handler \\\n  --statement-id eventbridge-lifecycle \\\n  --action lambda:InvokeFunction \\\n  --principal events.amazonaws.com \\\n  --source-arn arn:aws:events:us-east-1:123456789012:rule/durable-execution-terminal-events\n```\n\n### Step 4: For S3 / Aurora / Redshift destinations\n\nThe same pattern applies — the lifecycle handler reads the event and writes to your destination. For S3 you'd overwrite the execution's JSON file; for Aurora/Redshift you'd run an UPDATE on the `status` and `end_time` columns.\n\n### What each status means\n\n| Status    | Origin                                                              | Captured by plugin? | Captured by EventBridge? |\n| --------- | ------------------------------------------------------------------- | ------------------- | ------------------------ |\n| RUNNING   | Lambda invocation (incl. wait/suspend for timers & external events) | ✅ (on-change mode) | ❌                       |\n| SUCCEEDED | Lambda invocation                                                   | ✅                  | ✅                       |\n| FAILED    | Lambda invocation                                                   | ✅                  | ✅                       |\n| STOPPED   | Backend (manual stop API)                                           | ❌                  | ✅                       |\n| TIMED_OUT | Backend (ExecutionTimeout exceeded)                                 | ❌                  | ✅                       |\n\n> **Tip:** If you use the `EventBridgeExporter` alongside this pattern, your\n> plugin-emitted events (SUCCEEDED/FAILED) and the backend-emitted events\n> (STOPPED/TIMED_OUT) can share the same rule or feed the same consumer —\n> just adjust the event pattern to include both sources.\n\n## How It Works\n\nThe plugin hooks into the durable execution lifecycle:\n\n1. **`onInvocationStart`** — records execution start time\n2. **`onOperationChange`** — (on-change mode) schedules export of a RUNNING snapshot\n3. **`onInvocationEnd`** — schedules export of the snapshot, gated by `emitMode`: on terminal SUCCEEDED/FAILED (`on-complete`), FAILED only (`on-failure`), or every update including in-flight RUNNING snapshots (`on-change`)\n4. **`wrapInvocation`** — drains all pending exports before the Lambda returns\n\nExports are **coalesced**: if updates arrive faster than the exporter can handle, intermediate snapshots are dropped (each record is a complete snapshot, so the latest one supersedes all earlier ones). This prevents overlapping export calls and keeps overhead minimal.\n\nExporter errors **never fail the execution**. If an exporter throws, the error is swallowed and other exporters still receive the record.\n\n## License\n\nApache-2.0\n","readmeFilename":"README.md","_rev":"1-bfaa8b9327b225d4f21bff253bc12c7c"}