From 146ca0d99e6ce668146b06f71117dd2cba368380 Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:09:35 +0300 Subject: [PATCH 01/12] docs(integration-toolkit): document poll delivery JSONata transform Add a Payload Mapping section to Pollable Outbound: enqueue-time evaluation, bindings, object output, mapping_version, failure handling per poison_policy, redrive with reapply_mapping, the mapping preview endpoint and a multi-organization example. Document the org_id and mapping_version envelope fields, late-arrival re-keying with MSG_LATE_ARRIVAL, and removal of internal keys from poll payloads. --- .../integration-toolkit/pollable-outbound.md | 253 +++++++++++++++++- 1 file changed, 241 insertions(+), 12 deletions(-) diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index f1752d2..142190d 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -30,7 +30,7 @@ sequenceDiagram participant ERP as ERP / Middleware EC->>Q: Standardized business event (e.g. contract.updated) - Note over Q: Item stored with raw event payload
(idempotent — redelivery creates no duplicate) + Note over Q: Optional JSONata transform applied,
item stored (idempotent — redelivery creates no duplicate) ERP->>Q: POST /outbound/messages/poll { limit } Q-->>ERP: Leased batch + lease_token per message Note over ERP: Durably persist the items @@ -44,7 +44,7 @@ Key properties: - **Lease + ack/delete (at-least-once).** A poll leases a batch under a visibility timeout, hiding it from concurrent polls. Items you do not acknowledge in time reappear on a later poll — a consumer crash never loses data, but you must handle occasional redelivery (deduplicate by `id` or `event_id`). - **One polling loop per integration.** A single poll returns the merged feed across **all** of the integration's poll-mode use cases. Each message carries `use_case_id` and `event_name` for routing on your side. - **FIFO ordering, promised per entity.** Updates to the same entity are never delivered out of order — even across lease timeouts and retries. See [Ordering Guarantees](#ordering-guarantees). -- **Raw standardized events.** Poll messages carry the [Core Event](/docs/integrations/core-events) payload **as-is** — no JSONata mapping is applied in poll mode (see [Payload Contract](#payload-contract)). +- **Raw or mapped payloads.** By default poll messages carry the [Core Event](/docs/integrations/core-events) payload **as-is**. An optional JSONata transform reshapes it at enqueue time, so a consumer can receive one consistent shape (see [Payload Mapping](#payload-mapping)). - **Long, configurable retention.** Undelivered items are kept for `retention_days` (default 30, max 90) — designed for consumers that are legitimately offline for days. ## Configuration @@ -81,14 +81,14 @@ curl -X POST 'https://integration-toolkit.sls.epilot.io/v1/integrations/{integra | Property | Type | Required | Default | Bounds | Description | |----------|------|----------|---------|--------|-------------| | `retention_days` | integer | No | `30` | min 1, max 90 | How long undelivered queue items are retained before expiry | -| `poison_policy` | string | No | `"dead_letter"` | `dead_letter` \| `block` | What happens when an item exhausts `max_delivery_attempts` — see [Poison Messages](#poison-messages-dead-letter-vs-block) | +| `poison_policy` | string | No | `"dead_letter"` | `dead_letter` \| `block` | What happens when an item exhausts `max_delivery_attempts` — see [Poison Messages](#poison-messages-dead_letter-vs-block) | | `max_delivery_attempts` | integer | No | `5` | min 1, max 100 | Delivery (lease) attempts before the `poison_policy` is applied | Validation rules: - At most **one** poll mapping per use case (regardless of `enabled`). Webhook mappings may coexist alongside it — push and poll for the same event is allowed. - A `poll` delivery must not carry webhook fields (`webhook_id`, `webhook_name`), and a `webhook` delivery must not carry poll fields (`retention_days`, `poison_policy`, `max_delivery_attempts`). -- `jsonata_expression` is required for `webhook` mappings only; for `poll` mappings it is permitted but ignored. +- `jsonata_expression` is required for `webhook` mappings only. For `poll` mappings it is optional: when set, it transforms each payload at enqueue time (see [Payload Mapping](#payload-mapping)). It is validated on save — JSONata syntax, the 10,000-character limit, and only the bindings `$env`, `$mapValue` and `$mapKey`. When a poll use case is enabled, the event-catalog event is enabled as usual, but **no webhook configuration is created** — events route to the queue instead. @@ -98,6 +98,7 @@ When a poll use case is enabled, the event-catalog event is enabled as usual, bu |-----------|----------------| | `poll`, `ack` | `integration:consume`, resource-scoped to the integration | | `dlq`, `dlq/redrive`, `unblock` (operator) | `integration:manage` | +| `outbound/mapping-simulation` ([preview](#preview-a-transform)) | `integration:view`, resource-scoped to the integration | The [`integration:consume` grant](../../auth/grant-actions.md) enables least-privilege middleware tokens: a role granting only `integration:consume` on one integration can poll and acknowledge that feed and nothing else — no configuration reads, no other integrations, no entity access. Consumer tokens with only `integration:consume` receive `403` on the operator endpoints. @@ -133,11 +134,13 @@ Response — a leased batch spanning every enabled poll use case of the integrat { "id": "msg_9f3c8a1b…", "lease_token": "lt_a1b2…", + "org_id": "123456", "use_case_id": "uc_contract_sync", "event_name": "contract.updated", "event_id": "evt_77…", "group": "0", - "payload": { "...": "the standardized event-catalog event, as-is" }, + "payload": { "...": "the mapped output, or the standardized event as-is" }, + "mapping_version": "3f9a1c0b7d2e4a65", "enqueued_at": "2026-06-09T08:00:00Z" } ], @@ -150,11 +153,13 @@ Response — a leased batch spanning every enabled poll use case of the integrat |-------|-------------| | `id` | Opaque message id (`msg_…`) — stable per message across leases; use it for deduplication | | `lease_token` | Opaque lease token (`lt_…`) — must be echoed back on ack; changes when a lapsed message is re-leased | +| `org_id` | The epilot organization the message belongs to — always populated. Useful for middleware that polls several organizations into one pipeline | | `use_case_id` | The poll-mode use case that produced this message — route on this when consuming multiple use cases | | `event_name` | Standardized event name (e.g. `contract.updated`) | | `event_id` | Unique id of the originating event | | `group` | Ordering group — messages sharing a group are strictly ordered, distinct groups are independent. Constant `"0"` in v1 | -| `payload` | The raw standardized event — always inlined, regardless of size | +| `payload` | The mapped output when the use case has a `jsonata_expression`, otherwise the raw standardized event — always inlined, regardless of size | +| `mapping_version` | Version of the transform that produced `payload` (see [Mapping version](#mapping-version)). Absent for raw payloads | | `visibility_timeout_seconds` | Effective visibility timeout for this lease (per-integration server-side setting, default `300`) | | `has_more` | Whether more messages are waiting beyond this batch | @@ -226,18 +231,187 @@ Consequences of FIFO with a single stream: - **One in-flight batch per stream.** While a batch is leased, concurrent polls return an empty batch. Consumer-side parallelism does not increase throughput within a stream. - **The stream spans all of an integration's poll use cases.** A blocking message from one use case also holds back the others' messages (relevant for `poison_policy: "block"`, below). +### Late arrivals + +The stream is ordered by event time. Occasionally an event reaches the queue after the consumer has already moved past the position where it would sort — for example, when the event itself was delayed on the way in. Such an event is never inserted behind the consumer's position, where it would be skipped. Instead it is **re-keyed to the tail** of the stream and delivered after everything already enqueued, and the `MSG_LATE_ARRIVAL` monitoring code is emitted with the original and the re-keyed sequence time in its detail. + +This keeps the ordering promise unchanged: order is promised **per entity** only. A late arrival can be delivered after a newer event for a different entity. If your consumer needs the time the business event actually happened, read the event's `_event_time` field rather than relying on delivery order. + ## Payload Contract -Poll messages carry the **raw standardized event-catalog payload** — no JSONata transform is applied in poll mode. The mapping's `jsonata_expression` is ignored for poll mappings (it is not even syntax-checked). Condition filtering (`filterConditions`) is likewise not available for poll mode: a poll use case enqueues **every** event of its configured type, and the consumer filters on its side if needed. +Without a transform, poll messages carry the **raw standardized event-catalog payload**, exactly as the event catalog emitted it. With a `jsonata_expression` on the poll mapping, they carry the **mapped output** instead — see [Payload Mapping](#payload-mapping). The `mapping_version` envelope field tells the two apart: it is present on mapped payloads and absent on raw ones. -This is an intentional difference between the two delivery modes: JSONata mapping and condition filtering are push-pipeline features executed inside the webhook delivery infrastructure, which a poll use case bypasses entirely. +The internal keys `_downgrades` and `_automation_chain` are removed from poll payloads. They are bookkeeping for epilot's own event pipeline and carry no business data. + +Webhook condition filtering (`filterConditions`) is not available for poll mode. To limit which events a poll use case enqueues, use the use case's `event_filter` — a JSONata predicate evaluated against the same standardized event before anything is enqueued. Events the filter rejects never reach the queue. :::caution Switching delivery types changes the payload shape -A mapping switched from `webhook` to `poll` (or vice versa) changes what the consumer receives: webhook consumers get the **mapped** (JSONata-transformed) payload, poll consumers get the **raw** standardized event. +A mapping switched from `webhook` to `poll` (or vice versa) changes what the consumer receives, even with the same `jsonata_expression`. See [Webhook and poll mode are not interchangeable](#webhook-and-poll-mode-are-not-interchangeable). ::: Payloads are copied onto the queue item at enqueue time, so an item stays consumable for its full retention window even after the source event has aged out of the event catalog. Payloads are **always inlined** in the poll response regardless of size — the batch simply includes as many messages as fit within `limit` and the ~5.5 MB response cap. +## Payload Mapping + +A poll mapping can carry an optional `jsonata_expression` that reshapes each event before it is enqueued. The typical reason is a middleware that consumes many epilot organizations: each organization's configuration differs a little (status codes, reason labels, which attributes are filled), and a per-organization mapping evens those differences out so the middleware receives **one consistent shape**. + +```json +{ + "id": "b8f1c9a0-58dd-4f7a-9a3e-000000000001", + "name": "Meter readings to middleware", + "enabled": true, + "jsonata_expression": "{ \"meter_number\": meter_number, \"reason\": $mapValue($env.reading_reason, reason, \"OTHER\") }", + "delivery": { + "type": "poll", + "poison_policy": "dead_letter" + } +} +``` + +### How the transform is evaluated + +- **When:** once, at **enqueue time**, right after the use case's `event_filter` has accepted the event. The mapped output is stored on the queue item, so every poll of that item returns the same payload — a lease lapse or a redelivery never re-evaluates the expression. +- **Input:** the standardized event-catalog event, hydrated in full — the same root that `event_filter` sees. Field paths start at the top level of the [Core Event](/docs/integrations/core-events) (`_event_id`, `_event_time`, `meter_number`, …). +- **Bindings:** `$env` (the organization's non-secret environment variables, including [Key/Value Maps](./key-value-maps.md)), `$mapValue` and `$mapKey`. No other bindings are available — in particular there is no `$now`, because the output must depend only on the event and the configuration. An expression that references any other `$`-binding is rejected on save. +- **Output:** must be a **JSON object**. An array, a scalar, `null`, or an undefined result is a mapping failure (`invalid_output`). +- **Empty means raw.** An absent, empty, or whitespace-only `jsonata_expression` applies no transform, and the raw standardized event is delivered — the behavior of every poll use case that has no expression. + +The expression is validated on save: JSONata syntax, a maximum of 10,000 characters, and no bindings outside `$env`, `$mapValue` and `$mapKey`. + +### Mapping version + +Every mapped message carries a `mapping_version`: the first 16 hexadecimal characters of the SHA-256 hash of the trimmed `jsonata_expression`. Two messages with the same `mapping_version` were produced by the same expression. The field is absent on raw payloads. + +Because the transform runs at enqueue time, **changing the expression affects new items only**. Items already on the queue keep the payload their original expression produced. While the queue drains after a change, a consumer can therefore see old and new shapes side by side (or raw and mapped ones, when an expression was added or removed). Use `mapping_version` to tell them apart, and roll out shape changes that your consumer can read in both versions until the old items are gone. + +### Mapping failures + +A failure to evaluate — a runtime error, a timeout, or an output that is not an object — does not drop the event. The item is still enqueued at its normal position in the stream, marked as failed, with the raw payload kept alongside it. At enqueue time epilot emits the `MAPPING_EXPRESSION_FAILED` monitoring error, with the use case, event id, event name, `mapping_version` and the error message in its context. + +A failed item is **never delivered to the consumer**. When it reaches the head of the stream it is handled immediately — without waiting for `max_delivery_attempts`, because evaluating the same expression against the same event would fail the same way every time. What happens next follows the use case's `poison_policy`: + +| Policy | What happens to a failed item at the head | +|--------|-------------------------------------------| +| `dead_letter` (default) | It moves straight to the [dead-letter queue](#dead-letter-queue-and-operator-actions) with `reason: "mapping_failed"`, and the stream moves on. `MSG_DEAD_LETTERED` is emitted with `reason`, `mapping_error` and `mapping_version` in its detail | +| `block` | The stream halts on it. `MSG_HEAD_BLOCKED` is emitted with `reason: "mapping_failed"` in its detail. Release it with [`unblock`](#unblock--skip-a-blocked-head), as for any other blocked head | + +To recover, fix the expression (or the missing environment variable), then [redrive](#redrive--re-enqueue-dead-lettered-messages) the affected DLQ entries. A redrive always re-applies the **current** mapping to entries that failed mapping. + +:::tip A missing map is a mapping failure +`$mapValue` and `$mapKey` fail when their first argument is not an object — for example when `$env.reading_reason` has not been created yet in an organization. Create the environment variables before enabling the use case, and use the [preview endpoint](#preview-a-transform) with a real `event_id` to confirm. +::: + +### Preview a transform + +``` +POST /v1/integrations/{integrationId}/outbound/mapping-simulation +``` + +Evaluates an expression exactly as enqueue would, without enqueuing anything. Requires the `integration:view` grant on the integration. + +```json +{ + "jsonata_expression": "{ \"meter_number\": meter_number, \"reason\": $mapValue($env.reading_reason, reason, \"OTHER\") }", + "event_id": "01J9Z…", + "event_catalog_event": "MeterReadingAdded" +} +``` + +| Field | Type | Required | Description | +|-------|------|----------|-------------| +| `jsonata_expression` | string | Yes | The expression to evaluate | +| `payload` | object | One of `payload` / `event_id` | An event to evaluate against, supplied inline | +| `event_id` | string | One of `payload` / `event_id` | A historical event. epilot loads it from the event catalog and hydrates it exactly as enqueue does, so the preview matches what the queue would store | +| `event_catalog_event` | string | No | The event-catalog event the `event_id` belongs to | + +Supply exactly one of `payload` or `event_id`. `$env`, `$mapValue` and `$mapKey` resolve on the server from the organization's non-secret environment variables and key/value maps, the same way as at enqueue. + +```json +{ + "valid": true, + "output": { "meter_number": "A-1002", "reason": "PERIODIC" }, + "mapping_version": "3f9a1c0b7d2e4a65", + "input": { "_event_id": "01J9Z…", "_event_name": "MeterReadingAdded", "...": "the hydrated event" } +} +``` + +| Field | Description | +|-------|-------------| +| `valid` | Whether the expression evaluated to a JSON object | +| `output` | The mapped output — present when `valid` is `true` | +| `error` | Present when `valid` is `false`: a `code`, a `message`, and for syntax errors the `position` in the expression | +| `mapping_version` | The version this expression would stamp on messages | +| `input` | The hydrated event the expression ran against — returned when the request used `event_id` | + +| Error `code` | Meaning | +|--------------|---------| +| `syntax_error` | The expression is not valid JSONata | +| `unknown_binding` | The expression uses a `$`-binding other than `$env`, `$mapValue` or `$mapKey` | +| `evaluation_error` | The expression failed while running (for example, a missing key/value map) | +| `timeout` | The evaluation took too long | +| `invalid_output` | The result is not a JSON object | +| `expression_too_long` | The expression exceeds 10,000 characters | + +A mapping error is a normal `200` response with `valid: false`. A `4xx` status means the request itself was malformed — for example, both or neither of `payload` and `event_id`. + +### Webhook and poll mode are not interchangeable + +The same `jsonata_expression` is **not guaranteed to produce the same output** in webhook and poll mode. The webhook service builds its own input for the expression and enriches the event differently; in poll mode the input is the standardized event-catalog event described above. Write and test a poll expression against poll input — the [preview endpoint](#preview-a-transform) with a real `event_id` is the reliable check — rather than copying a webhook expression across unchanged. + +### Example: one shape for a multi-organization middleware + +A middleware collects meter readings from several utilities, each its own epilot organization. The organizations label reading reasons differently — one uses English labels, another German ones — but the middleware wants a single set of reason codes. + +Each organization declares the same [Key/Value Map](./key-value-maps.md) under the same environment key, `reading_reason`, with its own entries: + +```json title="Organization A — reading_reason" +{ "Periodic reading": "PERIODIC", "Move-out": "MOVE_OUT", "Meter change": "METER_CHANGE" } +``` + +```json title="Organization B — reading_reason" +{ "Turnusablesung": "PERIODIC", "Auszug": "MOVE_OUT", "Zählerwechsel": "METER_CHANGE" } +``` + +Both organizations then use the **same** expression on their `MeterReadingAdded` poll use case. Because it references only `$env`, it contains nothing organization-specific and can ship unchanged in a blueprint: + +```jsonata title="Poll mapping — MeterReadingAdded" +{ + "event_id": _event_id, + "occurred_at": _event_time, + "meter_number": meter_number, + "contract_number": contract_number, + "reason": $mapValue($env.reading_reason, reason, "OTHER"), + "readings": [meter_readings.{ + "obis_number": obis_number, + "value": value, + "direction": direction, + "read_at": reading_timestamp + }] +} +``` + +Whichever organization a message comes from, the middleware receives the same object: + +```json +{ + "event_id": "01J9Z…", + "occurred_at": "2026-06-09T07:58:12.000Z", + "meter_number": "A-1002", + "contract_number": "C-20931", + "reason": "PERIODIC", + "readings": [ + { "obis_number": "1-0:1.8.0", "value": 18234, "direction": "feed-out", "read_at": "2026-06-09T07:55:00.000Z" } + ] +} +``` + +A few details make this robust: + +- The `[ … ]` around `meter_readings.{ … }` keeps `readings` an array even when an event carries a single reading. +- The `"OTHER"` default keeps an unmapped label from producing a missing field; the middleware sees an explicit code it can route for review. +- The message envelope's `org_id` tells the middleware which organization a message came from, so the payload does not need to repeat it. +- `occurred_at` carries the business event time, which stays correct even for [late arrivals](#late-arrivals). + ## Retention and Expiry - Undelivered items expire after the mapping's `retention_days` (default 30, max 90), counted from enqueue time. @@ -259,6 +433,7 @@ Notes: - Delivery attempts only increment when a message is actually **leased** — an offline consumer never triggers poison handling by mere absence. - With `block`, remember the blast radius: the stream spans all of the integration's poll use cases, so a blocking message from one use case also holds back the others. That is the deliberate trade-off of `block`. - A blocked stream raises the `MSG_HEAD_BLOCKED` monitoring error and a `stream_blocked` conflict in [`outbound-status`](#queue-health-in-outbound-status). +- An item whose [payload mapping failed](#mapping-failures) is poisoned from the start: the policy applies as soon as it reaches the head, without waiting for `max_delivery_attempts`. Because it is never delivered, a consumer acknowledgement cannot release it — under `block`, use [`unblock`](#unblock--skip-a-blocked-head). ## Dead-Letter Queue and Operator Actions @@ -285,12 +460,31 @@ Returns dead-lettered messages oldest first, paginated via an opaque `next_token "delivery_attempts": 5, "reason": "max_delivery_attempts exhausted", "expires_at": "2026-07-09T03:00:00Z" + }, + { + "id": "msg_7a41…", + "use_case_id": "uc_meter_readings", + "event_name": "MeterReadingAdded", + "event_id": "evt_93…", + "enqueued_at": "2026-06-09T02:40:00Z", + "dead_lettered_at": "2026-06-09T02:41:00Z", + "delivery_attempts": 0, + "reason": "mapping_failed", + "mapping_error": "$mapValue: first argument must be an object", + "mapping_version": "3f9a1c0b7d2e4a65", + "expires_at": "2026-07-09T02:41:00Z" } ], "next_token": "…" } ``` +| Field | Description | +|-------|-------------| +| `reason` | Why the message was dead-lettered: the policy (`max_delivery_attempts` exhausted), the operator's `unblock` reason, or `mapping_failed` for a [payload mapping failure](#mapping-failures) | +| `mapping_error` | The mapping error message (truncated to 1,024 characters) — present when the payload mapping failed | +| `mapping_version` | Version of the expression that produced the stored payload, or that failed on it | + ### Redrive — re-enqueue dead-lettered messages ``` @@ -299,11 +493,41 @@ POST /v1/integrations/{integrationId}/outbound/messages/dlq/redrive ```json { - "ids": ["msg_5d2e…"] + "ids": ["msg_5d2e…", "msg_7a41…"], + "reapply_mapping": true } ``` -Takes 1–100 message ids and reports a per-id outcome — `redriven`, or `not_found` for unknown ids and entries concurrently redriven or expired. The redriven copy is re-enqueued with **zero delivery attempts and a fresh retention window**; the original DLQ entry is removed. +| Field | Type | Required | Default | Description | +|-------|------|----------|---------|-------------| +| `ids` | string[] | Yes | — | 1–100 message ids to redrive | +| `reapply_mapping` | boolean | No | `false` | Re-apply the use case's **current** `jsonata_expression` to the raw payload instead of re-sending the stored payload | + +The response reports a per-id outcome: + +```json +{ + "results": [ + { "id": "msg_5d2e…", "status": "redriven" }, + { "id": "msg_7a41…", "status": "mapping_failed", "mapping_error": "$mapValue: first argument must be an object" } + ] +} +``` + +| `status` | Meaning | +|----------|---------| +| `redriven` | Re-enqueued at the tail | +| `not_found` | Unknown id, or the entry was concurrently redriven or expired | +| `mapping_failed` | Re-applying the mapping failed again. The entry stays in the DLQ with its `mapping_error` and `mapping_version` updated; `mapping_error` is also returned in the result | + +The redriven copy is re-enqueued with **zero delivery attempts and a fresh retention window**; the original DLQ entry is removed. + +How the payload of the redriven copy is chosen: + +- By default (`reapply_mapping: false`) the stored payload is re-sent unchanged — the behavior before payload mapping existed. +- With `reapply_mapping: true` the current expression is evaluated again against the kept raw payload. Use it after fixing an expression, so already dead-lettered messages go out in the corrected shape. +- An entry that was dead-lettered **because its mapping failed** (`reason: "mapping_failed"`) is always re-mapped with the current configuration, whatever `reapply_mapping` says — it has no usable mapped payload to re-send. +- When re-mapping and the use case no longer has an expression, the raw payload is delivered. :::caution Redrive ordering A redriven message is re-enqueued at the **tail** with a new id and sequence — it is delivered out of its original per-entity order, because the stream has moved on. This is inherent to redrive and matches SQS DLQ semantics. If your consumer is order-sensitive, reconcile redriven messages explicitly (e.g. compare against current entity state). @@ -338,6 +562,10 @@ Poll-queue message lifecycle events flow through the standard monitoring pipelin | `MSG_EXPIRED_UNPOLLED` | error | An item's retention window elapsed without it ever being consumed — the offline-consumer loss signal | | `MSG_DEAD_LETTERED` | error | A message exhausted `max_delivery_attempts` under the `dead_letter` policy, or an operator skipped a blocked head (includes `delivery_attempts` in the event detail) | | `MSG_HEAD_BLOCKED` | error | The stream halted on a poisoned head under the `block` policy — emitted **once per blocked episode**, not on every poll (includes `delivery_attempts` in the event detail) | +| `MSG_LATE_ARRIVAL` | warning | An event arrived after the consumer had moved past its position and was re-keyed to the tail of the stream (includes `original_sequence_time` and `rekeyed_sequence_time` in the event detail) — see [Late arrivals](#late-arrivals) | +| `MAPPING_EXPRESSION_FAILED` | error | The poll mapping's `jsonata_expression` failed at enqueue time (includes `use_case_id`, `event_id`, `event_name`, `mapping_version` and the error message) — see [Mapping failures](#mapping-failures) | + +When a failed mapping reaches the head, `MSG_DEAD_LETTERED` or `MSG_HEAD_BLOCKED` carries `reason: "mapping_failed"` in its detail, so the mapping failure and its consequence for the stream can be told apart from ordinary poison messages. Lease lapses (a message reappearing after a visibility timeout) are deliberately **not** a per-occurrence signal — they are normal at-least-once behavior and would be noisy. The attempt count is reported on `MSG_DEAD_LETTERED` / `MSG_HEAD_BLOCKED`, which are the actionable events. @@ -377,6 +605,7 @@ The poll and ack endpoints are part of the ERP Integration API's OpenAPI spec an ## Current Limitations -- **No JSONata payload mapping** and **no condition filtering** for poll mappings — poll delivers raw standardized events; the consumer transforms and filters on its side (see [Payload Contract](#payload-contract)). +- **No webhook condition filtering** (`filterConditions`) for poll mappings — use the use case's `event_filter` instead (see [Payload Contract](#payload-contract)). +- **Payload mapping bindings are limited** to `$env`, `$mapValue` and `$mapKey`, and the output must be a JSON object (see [Payload Mapping](#payload-mapping)). - **At most one poll mapping per use case.** - **One in-flight batch per integration stream** — consumer-side parallelism does not increase throughput. The contract is shard-ready (per-entity ordering promise, per-message `group` field), but sharding is a server-side setting that is not yet enabled. From 20ee45b1b3c00ecc5a4e58f9617949a58f6d065e Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:09:35 +0300 Subject: [PATCH 02/12] docs(integration-toolkit): cross-link poll payload mapping from related pages Update configuration, key/value maps, overview and use cases, which stated that poll delivery applies no JSONata transform. --- docs/integrations/integration-toolkit/configuration.md | 6 +++--- docs/integrations/integration-toolkit/key-value-maps.md | 3 ++- docs/integrations/integration-toolkit/overview.md | 1 + docs/integrations/integration-toolkit/use-cases.md | 2 ++ 4 files changed, 8 insertions(+), 4 deletions(-) diff --git a/docs/integrations/integration-toolkit/configuration.md b/docs/integrations/integration-toolkit/configuration.md index 7f623e6..49afcc3 100644 --- a/docs/integrations/integration-toolkit/configuration.md +++ b/docs/integrations/integration-toolkit/configuration.md @@ -192,7 +192,7 @@ curl -X POST 'https://integration-toolkit.sls.epilot.io/v1/integrations/{integra | `id` | string (UUID) | No | Unique identifier for the mapping; generated when omitted | | `name` | string | Yes | Display name for the mapping | | `enabled` | boolean | Yes | Whether this mapping is active | -| `jsonata_expression` | string | For `webhook` delivery | JSONata expression to transform the event payload. Required for `webhook`, ignored for `poll`, and rejected for `file_proxy` delivery | +| `jsonata_expression` | string | For `webhook` delivery | JSONata expression to transform the event payload. Required for `webhook` (evaluated by the webhook service). Optional for `poll`: evaluated at enqueue time against the standardized event-catalog event with `$env` / `$mapValue` / `$mapKey`, must return a JSON object, and an empty value delivers the raw event — see [Payload Mapping](./pollable-outbound.md#payload-mapping). Rejected for `file_proxy` delivery | | `delivery` | object | Yes | How the event is delivered — discriminated on `type`: `webhook`, `poll`, or `file_proxy` | #### Delivery Types @@ -211,7 +211,7 @@ curl -X POST 'https://integration-toolkit.sls.epilot.io/v1/integrations/{integra | `webhook_id` | string | Yes | Reference to the webhook configuration in epilot Webhooks | | `webhook_name` | string | No | Cached webhook name for display purposes | -**Poll delivery (pull):** for ERPs that cannot expose an inbound HTTP endpoint (firewalled, on-prem, batch systems). Items are placed on a pull-based queue that your system fetches and acknowledges. Poll items carry the **raw standardized event payload** — no JSONata transform is applied. See [Pollable Outbound](./pollable-outbound.md) for the full feature documentation (polling API, ordering guarantees, dead-letter handling, monitoring): +**Poll delivery (pull):** for ERPs that cannot expose an inbound HTTP endpoint (firewalled, on-prem, batch systems). Items are placed on a pull-based queue that your system fetches and acknowledges. Poll items carry the **raw standardized event payload**, unless the mapping sets a `jsonata_expression` — then they carry its output, evaluated once at enqueue time (see [Payload Mapping](./pollable-outbound.md#payload-mapping)). See [Pollable Outbound](./pollable-outbound.md) for the full feature documentation (polling API, payload mapping, ordering guarantees, dead-letter handling, monitoring): ```jsonc // DeliveryConfig — poll variant @@ -234,7 +234,7 @@ curl -X POST 'https://integration-toolkit.sls.epilot.io/v1/integrations/{integra - A `poll` delivery must not carry webhook fields (`webhook_id`, `webhook_name`), and a `webhook` delivery must not carry poll fields (`retention_days`, `poison_policy`, `max_delivery_attempts`). ::: -Everything beyond the configuration contract — the polling and acknowledgement API, lease and ordering semantics, retention and expiry behavior, the dead-letter queue and operator actions, and poll-mode monitoring — is documented on the dedicated [Pollable Outbound](./pollable-outbound.md) page. +Everything beyond the configuration contract — the polling and acknowledgement API, payload mapping and its preview endpoint, lease and ordering semantics, retention and expiry behavior, the dead-letter queue and operator actions, and poll-mode monitoring — is documented on the dedicated [Pollable Outbound](./pollable-outbound.md) page. **File proxy delivery (push):** points to an upload-direction `file_proxy` use case in the same integration. The referenced recipe owns fan-out, payload mapping, authentication, and HTTP steps. `jsonata_expression` is rejected on this mapping type, and the use case's `event_catalog_event` must declare `event_attachments`. See [Outbound File Delivery](./outbound-file-delivery.md) for the complete setup and runtime behavior. diff --git a/docs/integrations/integration-toolkit/key-value-maps.md b/docs/integrations/integration-toolkit/key-value-maps.md index 752676f..91699ed 100644 --- a/docs/integrations/integration-toolkit/key-value-maps.md +++ b/docs/integrations/integration-toolkit/key-value-maps.md @@ -51,13 +51,14 @@ Rules worth knowing: - Lookup keys are compared as strings — `$mapValue($env.salutation, 1)` and `$mapValue($env.salutation, "1")` are the same lookup. - `$mapKey` compares values with strict equality: a map value `"1"` does not match the number `1`. Coerce with `$string()` when the source field is numeric. - When no `default` is given and nothing matches, the result is `undefined` and the mapped attribute is simply omitted — the same behaviour as any other undefined JSONata result. -- If the first argument is not an object (for example the environment variable does not exist yet), the expression fails with `$mapValue: first argument must be an object` / `$mapKey: …`. In inbound use cases this surfaces as a mapping error in monitoring; in webhooks the delivery fails. +- If the first argument is not an object (for example the environment variable does not exist yet), the expression fails with `$mapValue: first argument must be an object` / `$mapKey: …`. In inbound use cases this surfaces as a mapping error in monitoring; in webhooks the delivery fails; in pollable outbound the queue item is marked as a [mapping failure](./pollable-outbound.md#mapping-failures). - Values are read through the environments cache, so a change to a map becomes visible to running integrations within about 60 seconds. ### Where `$env`, `$mapValue` and `$mapKey` are available - Inbound use cases — every `jsonataExpression` field mapping and entity-level `jsonata` expression, including the mapping simulation endpoint. - Outbound webhooks — the payload transformation and multipart form-field expressions. +- Pollable outbound — the poll mapping's [payload transform](./pollable-outbound.md#payload-mapping), evaluated at enqueue time, including its preview endpoint. See the [multi-organization example](./pollable-outbound.md#example-one-shape-for-a-multi-organization-middleware) for maps that give many organizations one payload shape. - Outbound file proxy — request body templates and delivery expressions. ## Recommended shape: one map, both directions diff --git a/docs/integrations/integration-toolkit/overview.md b/docs/integrations/integration-toolkit/overview.md index 5c24daf..6d8f83d 100644 --- a/docs/integrations/integration-toolkit/overview.md +++ b/docs/integrations/integration-toolkit/overview.md @@ -120,6 +120,7 @@ See the [Configuration Guide](./configuration.md#secure-proxy-use-cases) for set - Inbound event processing (ERP to epilot entity mapping) - Outbound webhook payloads (epilot event to ERP format) +- Pollable outbound payloads, transformed at enqueue time ([Payload Mapping](./pollable-outbound.md#payload-mapping)) - The Map Data flow building block ### Monitoring and Alerting diff --git a/docs/integrations/integration-toolkit/use-cases.md b/docs/integrations/integration-toolkit/use-cases.md index a2ed6dc..21d761b 100644 --- a/docs/integrations/integration-toolkit/use-cases.md +++ b/docs/integrations/integration-toolkit/use-cases.md @@ -344,6 +344,8 @@ This use case is typically combined with [Keep Customer In Sync](#keep-customer- Outbound use cases push epilot events to your ERP via [Webhooks](/docs/integrations/webhooks). When a user performs a self-service action in a portal, journey, or epilot 360, an automation triggers a [Core Event](/docs/integrations/core-events) which is delivered to your middle layer webhook endpoint. Your middle layer then processes the event and calls the appropriate ERP API. +The same events can also be delivered through [Pollable Outbound](./pollable-outbound.md) when the middle layer cannot receive webhooks. A poll mapping can reshape each event with an optional JSONata transform — see [Payload Mapping](./pollable-outbound.md#payload-mapping). + ### Submit Meter Reading A portal user or service agent submits a new meter reading. From b8569ddf90373ee1dc3cf224c695ca4cab423b5d Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:12:09 +0300 Subject: [PATCH 03/12] docs(integration-toolkit): document the outbound event_filter Describe event_filter as a general outbound use case setting: input, bindings, result, error handling, save-time validation, and which deliveries it applies to. Link it from the poll payload contract. --- .../integration-toolkit/configuration.md | 20 +++++++++++++++++++ .../integration-toolkit/pollable-outbound.md | 2 +- 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/docs/integrations/integration-toolkit/configuration.md b/docs/integrations/integration-toolkit/configuration.md index 49afcc3..68b089f 100644 --- a/docs/integrations/integration-toolkit/configuration.md +++ b/docs/integrations/integration-toolkit/configuration.md @@ -185,6 +185,26 @@ curl -X POST 'https://integration-toolkit.sls.epilot.io/v1/integrations/{integra }' ``` +#### Event Filter + +`event_filter` is an optional JSONata predicate on the use case configuration, next to `event_catalog_event`. It narrows which events of that name the use case handles — for example, only tickets with a certain purpose, or only certain contract types: + +```json +{ + "event_catalog_event": "CustomerRequestSubmitted", + "event_filter": "$count(ticket._purpose[$ = $env.move_request_purpose]) > 0", + "mappings": [ … ] +} +``` + +- **Input:** the full hydrated event-catalog event, so relation nodes such as `ticket` and `contact` are populated. +- **Bindings:** `$env` (the organization's non-secret environment variables, including [Key/Value Maps](./key-value-maps.md)), `$mapValue` and `$mapKey`. Referencing `$env` instead of hard-coding organization-specific values such as entity IDs keeps the filter portable between organizations, for example in a blueprint. +- **Result:** the use case handles the event only when the filter evaluates truthy. When `event_filter` is absent, every event of the configured name is handled. +- **Errors:** a filter that throws while evaluating is treated as **no match** and logged, so one malformed filter cannot stop the other use cases subscribed to the same event. +- **Validation on save:** the filter must be a non-empty string, valid JSONata, and use no bindings other than `$env`, `$mapValue` and `$mapKey` (names the expression binds itself with `:=` are allowed). + +An event the filter rejects is not processed by the use case at all: no [Pollable Outbound](./pollable-outbound.md) queue item and no [file delivery](./outbound-file-delivery.md). The Integration Toolkit evaluates the filter for poll and file proxy deliveries. Webhook deliveries are sent by epilot Webhooks, which does not evaluate `event_filter`. + #### Mapping Properties | Property | Type | Required | Description | diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index 142190d..24791c6 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -243,7 +243,7 @@ Without a transform, poll messages carry the **raw standardized event-catalog pa The internal keys `_downgrades` and `_automation_chain` are removed from poll payloads. They are bookkeeping for epilot's own event pipeline and carry no business data. -Webhook condition filtering (`filterConditions`) is not available for poll mode. To limit which events a poll use case enqueues, use the use case's `event_filter` — a JSONata predicate evaluated against the same standardized event before anything is enqueued. Events the filter rejects never reach the queue. +Webhook condition filtering (`filterConditions`) is not available for poll mode. To limit which events a poll use case enqueues, use the use case's [`event_filter`](./configuration.md#event-filter) — a JSONata predicate evaluated against the same standardized event before anything is enqueued. Events the filter rejects never reach the queue. :::caution Switching delivery types changes the payload shape A mapping switched from `webhook` to `poll` (or vice versa) changes what the consumer receives, even with the same `jsonata_expression`. See [Webhook and poll mode are not interchangeable](#webhook-and-poll-mode-are-not-interchangeable). From 58e92443e21dd405041a8371c283aa08fcfc0cf6 Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:13:25 +0300 Subject: [PATCH 04/12] docs(integration-toolkit): word event_filter delivery scope neutrally --- docs/integrations/integration-toolkit/configuration.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/integrations/integration-toolkit/configuration.md b/docs/integrations/integration-toolkit/configuration.md index 68b089f..1d129c8 100644 --- a/docs/integrations/integration-toolkit/configuration.md +++ b/docs/integrations/integration-toolkit/configuration.md @@ -203,7 +203,7 @@ curl -X POST 'https://integration-toolkit.sls.epilot.io/v1/integrations/{integra - **Errors:** a filter that throws while evaluating is treated as **no match** and logged, so one malformed filter cannot stop the other use cases subscribed to the same event. - **Validation on save:** the filter must be a non-empty string, valid JSONata, and use no bindings other than `$env`, `$mapValue` and `$mapKey` (names the expression binds itself with `:=` are allowed). -An event the filter rejects is not processed by the use case at all: no [Pollable Outbound](./pollable-outbound.md) queue item and no [file delivery](./outbound-file-delivery.md). The Integration Toolkit evaluates the filter for poll and file proxy deliveries. Webhook deliveries are sent by epilot Webhooks, which does not evaluate `event_filter`. +An event the filter rejects is not processed by the use case at all: no [Pollable Outbound](./pollable-outbound.md) queue item and no [file delivery](./outbound-file-delivery.md). `event_filter` currently applies to poll and file proxy deliveries; it is not evaluated for webhook deliveries. #### Mapping Properties From ec0a876b5389d6d81173b4705361b451312d2d32 Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:33:52 +0300 Subject: [PATCH 05/12] docs(integration-toolkit): align late arrivals and deduplication with backend behaviour Late arrivals are measured against everything already leased, carry no envelope marker and are signalled only by MSG_LATE_ARRIVAL. A rare lease race can deliver one event under two message ids, so consumers should deduplicate on use_case_id plus event_id. --- .../integration-toolkit/pollable-outbound.md | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index 24791c6..535b705 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -41,9 +41,9 @@ sequenceDiagram Key properties: -- **Lease + ack/delete (at-least-once).** A poll leases a batch under a visibility timeout, hiding it from concurrent polls. Items you do not acknowledge in time reappear on a later poll — a consumer crash never loses data, but you must handle occasional redelivery (deduplicate by `id` or `event_id`). +- **Lease + ack/delete (at-least-once).** A poll leases a batch under a visibility timeout, hiding it from concurrent polls. Items you do not acknowledge in time reappear on a later poll — a consumer crash never loses data, but you must handle occasional redelivery (deduplicate on `use_case_id` + `event_id` — see [A typical polling loop](#a-typical-polling-loop)). - **One polling loop per integration.** A single poll returns the merged feed across **all** of the integration's poll-mode use cases. Each message carries `use_case_id` and `event_name` for routing on your side. -- **FIFO ordering, promised per entity.** Updates to the same entity are never delivered out of order — even across lease timeouts and retries. See [Ordering Guarantees](#ordering-guarantees). +- **FIFO ordering, promised per entity.** Updates to the same entity are never delivered out of order — even across lease timeouts and retries. The one exception is an event that reaches the queue late, which goes to the tail and is flagged in monitoring (see [Late arrivals](#late-arrivals)). See [Ordering Guarantees](#ordering-guarantees). - **Raw or mapped payloads.** By default poll messages carry the [Core Event](/docs/integrations/core-events) payload **as-is**. An optional JSONata transform reshapes it at enqueue time, so a consumer can receive one consistent shape (see [Payload Mapping](#payload-mapping)). - **Long, configurable retention.** Undelivered items are kept for `retention_days` (default 30, max 90) — designed for consumers that are legitimately offline for days. @@ -207,7 +207,7 @@ loop (every N seconds / on schedule): batch = POST …/outbound/messages/poll { limit: 100 } if batch.messages is empty: sleep / wait for next run for message in batch.messages (in order): - persist message durably (dedupe on message.id) + persist message durably (dedupe on use_case_id + event_id) POST …/outbound/messages/ack { acks: all (id, lease_token) pairs } if batch.has_more: poll again immediately ``` @@ -216,7 +216,7 @@ Practical guidance: - **Finish well inside the visibility timeout.** If processing a batch can exceed `visibility_timeout_seconds`, lower your `limit` — a lapsed lease means the whole batch is re-delivered and your acks come back `stale_lease`. - **Ack in stream order**, ideally the whole batch at once. Partial acks are fine as long as they are contiguous from the head of the batch. -- **Deduplicate.** At-least-once delivery means a message can arrive twice (with a fresh `lease_token`). The `id` is stable across redeliveries. +- **Deduplicate on `use_case_id` + `event_id`.** At-least-once delivery means a message can arrive twice. A lease redelivery keeps the same `id` (with a fresh `lease_token`), but in a rare lease race the same event can also be delivered twice under **different** message ids — so `id` alone is not a sufficient deduplication key. The same `event_id` legitimately appears once per poll use case it matches, which is why the key includes `use_case_id`. - **Do not parallelize polls of one integration.** Only one batch can be in flight per stream; concurrent polls receive empty batches (this is by design, to preserve ordering). ## Ordering Guarantees @@ -233,9 +233,11 @@ Consequences of FIFO with a single stream: ### Late arrivals -The stream is ordered by event time. Occasionally an event reaches the queue after the consumer has already moved past the position where it would sort — for example, when the event itself was delayed on the way in. Such an event is never inserted behind the consumer's position, where it would be skipped. Instead it is **re-keyed to the tail** of the stream and delivered after everything already enqueued, and the `MSG_LATE_ARRIVAL` monitoring code is emitted with the original and the re-keyed sequence time in its detail. +The stream is ordered by event time. Occasionally an event reaches the queue after the stream has already moved past the position where it would sort — for example, when the event itself was delayed on the way in. "Moved past" covers everything already handed out in a lease, acknowledged or not. Such an event is never inserted behind that position, where it would be skipped. Instead it is **re-keyed to the tail** of the stream and delivered after everything already enqueued, and the `MSG_LATE_ARRIVAL` monitoring warning is emitted with `message_id`, `event_name`, `original_sequence_time` and `rekeyed_sequence_time` in its detail. -This keeps the ordering promise unchanged: order is promised **per entity** only. A late arrival can be delivered after a newer event for a different entity. If your consumer needs the time the business event actually happened, read the event's `_event_time` field rather than relying on delivery order. +Only the queue position changes. The payload is untouched — its `_event_time` still carries the original event time — and the poll message envelope has no late-arrival marker: `MSG_LATE_ARRIVAL` in monitoring is the only signal. + +Order is promised **per entity** only. Because a late arrival goes to the tail, it is delivered after messages that were already in the stream — including, in the rare case, a newer event for the same entity. If your consumer applies state changes, compare the event's `_event_time` with the last one you applied for that entity rather than relying on delivery order alone. ## Payload Contract From 13874bbcc4a64f60d8616df93801b68ea19d02c4 Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:36:30 +0300 Subject: [PATCH 06/12] docs(integration-toolkit): reconcile poll docs with final hardening behaviour Deduplicate by event_id, count expired leases toward the late-arrival watermark, and state that internal keys are stripped before storage and before the payload mapping runs. --- .../integration-toolkit/pollable-outbound.md | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index 535b705..01b01c3 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -41,7 +41,7 @@ sequenceDiagram Key properties: -- **Lease + ack/delete (at-least-once).** A poll leases a batch under a visibility timeout, hiding it from concurrent polls. Items you do not acknowledge in time reappear on a later poll — a consumer crash never loses data, but you must handle occasional redelivery (deduplicate on `use_case_id` + `event_id` — see [A typical polling loop](#a-typical-polling-loop)). +- **Lease + ack/delete (at-least-once).** A poll leases a batch under a visibility timeout, hiding it from concurrent polls. Items you do not acknowledge in time reappear on a later poll — a consumer crash never loses data, but you must handle occasional redelivery (deduplicate by `event_id` — see [A typical polling loop](#a-typical-polling-loop)). - **One polling loop per integration.** A single poll returns the merged feed across **all** of the integration's poll-mode use cases. Each message carries `use_case_id` and `event_name` for routing on your side. - **FIFO ordering, promised per entity.** Updates to the same entity are never delivered out of order — even across lease timeouts and retries. The one exception is an event that reaches the queue late, which goes to the tail and is flagged in monitoring (see [Late arrivals](#late-arrivals)). See [Ordering Guarantees](#ordering-guarantees). - **Raw or mapped payloads.** By default poll messages carry the [Core Event](/docs/integrations/core-events) payload **as-is**. An optional JSONata transform reshapes it at enqueue time, so a consumer can receive one consistent shape (see [Payload Mapping](#payload-mapping)). @@ -207,7 +207,7 @@ loop (every N seconds / on schedule): batch = POST …/outbound/messages/poll { limit: 100 } if batch.messages is empty: sleep / wait for next run for message in batch.messages (in order): - persist message durably (dedupe on use_case_id + event_id) + persist message durably (dedupe on message.event_id) POST …/outbound/messages/ack { acks: all (id, lease_token) pairs } if batch.has_more: poll again immediately ``` @@ -216,7 +216,7 @@ Practical guidance: - **Finish well inside the visibility timeout.** If processing a batch can exceed `visibility_timeout_seconds`, lower your `limit` — a lapsed lease means the whole batch is re-delivered and your acks come back `stale_lease`. - **Ack in stream order**, ideally the whole batch at once. Partial acks are fine as long as they are contiguous from the head of the batch. -- **Deduplicate on `use_case_id` + `event_id`.** At-least-once delivery means a message can arrive twice. A lease redelivery keeps the same `id` (with a fresh `lease_token`), but in a rare lease race the same event can also be delivered twice under **different** message ids — so `id` alone is not a sufficient deduplication key. The same `event_id` legitimately appears once per poll use case it matches, which is why the key includes `use_case_id`. +- **Deduplicate by `event_id`.** At-least-once delivery means an event can arrive twice. A lease redelivery keeps the same `id` (with a fresh `lease_token`), but in a rare race the same event can also be delivered twice under **different** message ids — so `id` alone is not a sufficient deduplication key. If several of your poll use cases subscribe to the same event, each produces its own message for it; scope the `event_id` check per `use_case_id` in that case. - **Do not parallelize polls of one integration.** Only one batch can be in flight per stream; concurrent polls receive empty batches (this is by design, to preserve ordering). ## Ordering Guarantees @@ -233,7 +233,7 @@ Consequences of FIFO with a single stream: ### Late arrivals -The stream is ordered by event time. Occasionally an event reaches the queue after the stream has already moved past the position where it would sort — for example, when the event itself was delayed on the way in. "Moved past" covers everything already handed out in a lease, acknowledged or not. Such an event is never inserted behind that position, where it would be skipped. Instead it is **re-keyed to the tail** of the stream and delivered after everything already enqueued, and the `MSG_LATE_ARRIVAL` monitoring warning is emitted with `message_id`, `event_name`, `original_sequence_time` and `rekeyed_sequence_time` in its detail. +The stream is ordered by event time. Occasionally an event reaches the queue after the stream has already moved past the position where it would sort — for example, when the event itself was delayed on the way in. "Moved past" covers everything already acknowledged and everything ever handed out in a lease — including leases that expired without an acknowledgement. Such an event is never inserted behind that position, where it would be skipped. Instead it is **re-keyed to the tail** of the stream and delivered after everything already enqueued, and the `MSG_LATE_ARRIVAL` monitoring warning is emitted with `message_id`, `event_name`, `original_sequence_time` and `rekeyed_sequence_time` in its detail. Only the queue position changes. The payload is untouched — its `_event_time` still carries the original event time — and the poll message envelope has no late-arrival marker: `MSG_LATE_ARRIVAL` in monitoring is the only signal. @@ -243,7 +243,7 @@ Order is promised **per entity** only. Because a late arrival goes to the tail, Without a transform, poll messages carry the **raw standardized event-catalog payload**, exactly as the event catalog emitted it. With a `jsonata_expression` on the poll mapping, they carry the **mapped output** instead — see [Payload Mapping](#payload-mapping). The `mapping_version` envelope field tells the two apart: it is present on mapped payloads and absent on raw ones. -The internal keys `_downgrades` and `_automation_chain` are removed from poll payloads. They are bookkeeping for epilot's own event pipeline and carry no business data. +The internal keys `_downgrades` and `_automation_chain` are removed from poll payloads before the item is stored, and again on delivery. They are bookkeeping for epilot's own event pipeline and carry no business data. A [payload mapping](#payload-mapping) is evaluated against the stripped event, so expressions never see these keys either. Webhook condition filtering (`filterConditions`) is not available for poll mode. To limit which events a poll use case enqueues, use the use case's [`event_filter`](./configuration.md#event-filter) — a JSONata predicate evaluated against the same standardized event before anything is enqueued. Events the filter rejects never reach the queue. @@ -273,7 +273,7 @@ A poll mapping can carry an optional `jsonata_expression` that reshapes each eve ### How the transform is evaluated - **When:** once, at **enqueue time**, right after the use case's `event_filter` has accepted the event. The mapped output is stored on the queue item, so every poll of that item returns the same payload — a lease lapse or a redelivery never re-evaluates the expression. -- **Input:** the standardized event-catalog event, hydrated in full — the same root that `event_filter` sees. Field paths start at the top level of the [Core Event](/docs/integrations/core-events) (`_event_id`, `_event_time`, `meter_number`, …). +- **Input:** the standardized event-catalog event, hydrated in full — the same root that `event_filter` sees, minus the internal `_downgrades` and `_automation_chain` keys. Field paths start at the top level of the [Core Event](/docs/integrations/core-events) (`_event_id`, `_event_time`, `meter_number`, …). - **Bindings:** `$env` (the organization's non-secret environment variables, including [Key/Value Maps](./key-value-maps.md)), `$mapValue` and `$mapKey`. No other bindings are available — in particular there is no `$now`, because the output must depend only on the event and the configuration. An expression that references any other `$`-binding is rejected on save. - **Output:** must be a **JSON object**. An array, a scalar, `null`, or an undefined result is a mapping failure (`invalid_output`). - **Empty means raw.** An absent, empty, or whitespace-only `jsonata_expression` applies no transform, and the raw standardized event is delivered — the behavior of every poll use case that has no expression. From 26001d1e7f8e75d5d8e9eaf2a8678a18590956fa Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:37:44 +0300 Subject: [PATCH 07/12] docs(integration-toolkit): deduplicate poll messages on use_case_id plus event_id --- docs/integrations/integration-toolkit/pollable-outbound.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index 01b01c3..ef9bb8c 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -41,7 +41,7 @@ sequenceDiagram Key properties: -- **Lease + ack/delete (at-least-once).** A poll leases a batch under a visibility timeout, hiding it from concurrent polls. Items you do not acknowledge in time reappear on a later poll — a consumer crash never loses data, but you must handle occasional redelivery (deduplicate by `event_id` — see [A typical polling loop](#a-typical-polling-loop)). +- **Lease + ack/delete (at-least-once).** A poll leases a batch under a visibility timeout, hiding it from concurrent polls. Items you do not acknowledge in time reappear on a later poll — a consumer crash never loses data, but you must handle occasional redelivery (deduplicate on `use_case_id` + `event_id` — see [A typical polling loop](#a-typical-polling-loop)). - **One polling loop per integration.** A single poll returns the merged feed across **all** of the integration's poll-mode use cases. Each message carries `use_case_id` and `event_name` for routing on your side. - **FIFO ordering, promised per entity.** Updates to the same entity are never delivered out of order — even across lease timeouts and retries. The one exception is an event that reaches the queue late, which goes to the tail and is flagged in monitoring (see [Late arrivals](#late-arrivals)). See [Ordering Guarantees](#ordering-guarantees). - **Raw or mapped payloads.** By default poll messages carry the [Core Event](/docs/integrations/core-events) payload **as-is**. An optional JSONata transform reshapes it at enqueue time, so a consumer can receive one consistent shape (see [Payload Mapping](#payload-mapping)). @@ -207,7 +207,7 @@ loop (every N seconds / on schedule): batch = POST …/outbound/messages/poll { limit: 100 } if batch.messages is empty: sleep / wait for next run for message in batch.messages (in order): - persist message durably (dedupe on message.event_id) + persist message durably (dedupe on use_case_id + event_id) POST …/outbound/messages/ack { acks: all (id, lease_token) pairs } if batch.has_more: poll again immediately ``` @@ -216,7 +216,7 @@ Practical guidance: - **Finish well inside the visibility timeout.** If processing a batch can exceed `visibility_timeout_seconds`, lower your `limit` — a lapsed lease means the whole batch is re-delivered and your acks come back `stale_lease`. - **Ack in stream order**, ideally the whole batch at once. Partial acks are fine as long as they are contiguous from the head of the batch. -- **Deduplicate by `event_id`.** At-least-once delivery means an event can arrive twice. A lease redelivery keeps the same `id` (with a fresh `lease_token`), but in a rare race the same event can also be delivered twice under **different** message ids — so `id` alone is not a sufficient deduplication key. If several of your poll use cases subscribe to the same event, each produces its own message for it; scope the `event_id` check per `use_case_id` in that case. +- **Deduplicate on `use_case_id` + `event_id`.** At-least-once delivery means an event can arrive twice. A lease redelivery keeps the same `id` (with a fresh `lease_token`), but in a rare race the same event can also be delivered twice under **different** message ids — so `id` alone is not a sufficient deduplication key. The same `event_id` legitimately appears once per poll use case it matches, which is why the key includes `use_case_id`. - **Do not parallelize polls of one integration.** Only one batch can be in flight per stream; concurrent polls receive empty batches (this is by design, to preserve ordering). ## Ordering Guarantees From 5331792341de912f361bd281809d9a7e9c8fd751 Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 13:45:01 +0300 Subject: [PATCH 08/12] docs(integration-toolkit): reconcile poll mapping docs with final feature behaviour Document the preview endpoint's 400/404 rules, the required event_catalog_event with event_id, the 500 ms and depth guardrails, enqueue-time re-validation, the mapping_error format, zero delivery attempts on failed-mapping dead letters, and the fix, unblock, redrive recovery workflow. --- .../integration-toolkit/pollable-outbound.md | 43 +++++++++++++------ 1 file changed, 29 insertions(+), 14 deletions(-) diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index ef9bb8c..353e99b 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -278,7 +278,9 @@ A poll mapping can carry an optional `jsonata_expression` that reshapes each eve - **Output:** must be a **JSON object**. An array, a scalar, `null`, or an undefined result is a mapping failure (`invalid_output`). - **Empty means raw.** An absent, empty, or whitespace-only `jsonata_expression` applies no transform, and the raw standardized event is delivered — the behavior of every poll use case that has no expression. -The expression is validated on save: JSONata syntax, a maximum of 10,000 characters, and no bindings outside `$env`, `$mapValue` and `$mapKey`. +The expression is validated on save: JSONata syntax, a maximum of 10,000 characters, and no bindings outside `$env`, `$mapValue` and `$mapKey`. The same checks run again at enqueue time, so an expression stored before these checks existed that does not pass them produces [failed items](#mapping-failures) rather than being silently skipped. + +Each evaluation is guarded by a 500 ms time limit and an evaluation depth limit. Exceeding either is a mapping failure with the code `timeout`. ### Mapping version @@ -290,14 +292,20 @@ Because the transform runs at enqueue time, **changing the expression affects ne A failure to evaluate — a runtime error, a timeout, or an output that is not an object — does not drop the event. The item is still enqueued at its normal position in the stream, marked as failed, with the raw payload kept alongside it. At enqueue time epilot emits the `MAPPING_EXPRESSION_FAILED` monitoring error, with the use case, event id, event name, `mapping_version` and the error message in its context. +The error is recorded as `mapping_error` in the form `: ` — for example `evaluation_error: $mapValue: first argument must be an object` — using the same codes as the [preview endpoint](#preview-a-transform), truncated to 1,024 characters. + A failed item is **never delivered to the consumer**. When it reaches the head of the stream it is handled immediately — without waiting for `max_delivery_attempts`, because evaluating the same expression against the same event would fail the same way every time. What happens next follows the use case's `poison_policy`: | Policy | What happens to a failed item at the head | |--------|-------------------------------------------| -| `dead_letter` (default) | It moves straight to the [dead-letter queue](#dead-letter-queue-and-operator-actions) with `reason: "mapping_failed"`, and the stream moves on. `MSG_DEAD_LETTERED` is emitted with `reason`, `mapping_error` and `mapping_version` in its detail | -| `block` | The stream halts on it. `MSG_HEAD_BLOCKED` is emitted with `reason: "mapping_failed"` in its detail. Release it with [`unblock`](#unblock--skip-a-blocked-head), as for any other blocked head | +| `dead_letter` (default) | It moves straight to the [dead-letter queue](#dead-letter-queue-and-operator-actions) with `reason: "mapping_failed"` and `delivery_attempts: 0` (it was never leased), and the stream moves on. `MSG_DEAD_LETTERED` is emitted with `reason`, `mapping_error` and `mapping_version` in its detail | +| `block` | The stream halts on it. `MSG_HEAD_BLOCKED` is emitted with `reason: "mapping_failed"` in its detail. Release it with [`unblock`](#unblock--skip-a-blocked-head), which dead-letters it with `reason: "mapping_failed"`. A consumer acknowledgement cannot release it, because it is never delivered | + +To recover: -To recover, fix the expression (or the missing environment variable), then [redrive](#redrive--re-enqueue-dead-lettered-messages) the affected DLQ entries. A redrive always re-applies the **current** mapping to entries that failed mapping. +1. Fix the expression, or create the missing environment variable. Check the fix with the [preview endpoint](#preview-a-transform) against the failed `event_id`. +2. Under `block`, [unblock](#unblock--skip-a-blocked-head) the stream. The failed head moves to the DLQ. +3. [Redrive](#redrive--re-enqueue-dead-lettered-messages) the affected DLQ entries. A redrive always re-applies the **current** mapping to entries that failed mapping, so `reapply_mapping` is not needed for them. :::tip A missing map is a mapping failure `$mapValue` and `$mapKey` fail when their first argument is not an object — for example when `$env.reading_reason` has not been created yet in an organization. Create the environment variables before enabling the use case, and use the [preview endpoint](#preview-a-transform) with a real `event_id` to confirm. @@ -309,7 +317,7 @@ To recover, fix the expression (or the missing environment variable), then [redr POST /v1/integrations/{integrationId}/outbound/mapping-simulation ``` -Evaluates an expression exactly as enqueue would, without enqueuing anything. Requires the `integration:view` grant on the integration. +Evaluates an expression exactly as enqueue would, without enqueuing anything. Requires the `integration:view` grant, scoped to the integration. An unknown integration, or one belonging to another organization, returns `404`. ```json { @@ -321,10 +329,10 @@ Evaluates an expression exactly as enqueue would, without enqueuing anything. Re | Field | Type | Required | Description | |-------|------|----------|-------------| -| `jsonata_expression` | string | Yes | The expression to evaluate | +| `jsonata_expression` | string | Yes | The expression to evaluate. A blank expression returns `400` | | `payload` | object | One of `payload` / `event_id` | An event to evaluate against, supplied inline | | `event_id` | string | One of `payload` / `event_id` | A historical event. epilot loads it from the event catalog and hydrates it exactly as enqueue does, so the preview matches what the queue would store | -| `event_catalog_event` | string | No | The event-catalog event the `event_id` belongs to | +| `event_catalog_event` | string | With `event_id` | The event-catalog event type used to load and hydrate the historical event, for example `MeterReadingAdded`. Required when `event_id` is given | Supply exactly one of `payload` or `event_id`. `$env`, `$mapValue` and `$mapKey` resolve on the server from the organization's non-secret environment variables and key/value maps, the same way as at enqueue. @@ -343,18 +351,23 @@ Supply exactly one of `payload` or `event_id`. `$env`, `$mapValue` and `$mapKey` | `output` | The mapped output — present when `valid` is `true` | | `error` | Present when `valid` is `false`: a `code`, a `message`, and for syntax errors the `position` in the expression | | `mapping_version` | The version this expression would stamp on messages | -| `input` | The hydrated event the expression ran against — returned when the request used `event_id` | +| `input` | The hydrated event the expression ran against — returned when the request used `event_id`, on failures as well as successes | | Error `code` | Meaning | |--------------|---------| | `syntax_error` | The expression is not valid JSONata | | `unknown_binding` | The expression uses a `$`-binding other than `$env`, `$mapValue` or `$mapKey` | | `evaluation_error` | The expression failed while running (for example, a missing key/value map) | -| `timeout` | The evaluation took too long | +| `timeout` | The evaluation exceeded the 500 ms time limit or the evaluation depth limit | | `invalid_output` | The result is not a JSON object | | `expression_too_long` | The expression exceeds 10,000 characters | -A mapping error is a normal `200` response with `valid: false`. A `4xx` status means the request itself was malformed — for example, both or neither of `payload` and `event_id`. +A mapping error is a normal `200` response with `valid: false`. A `4xx` status means the request itself could not be served: + +| Status | When | +|--------|------| +| `400` | The expression is blank, both or neither of `payload` and `event_id` are given, or `event_id` is given without `event_catalog_event` | +| `404` | The integration is unknown or belongs to another organization, or the event catalog does not know the requested event | ### Webhook and poll mode are not interchangeable @@ -472,7 +485,7 @@ Returns dead-lettered messages oldest first, paginated via an opaque `next_token "dead_lettered_at": "2026-06-09T02:41:00Z", "delivery_attempts": 0, "reason": "mapping_failed", - "mapping_error": "$mapValue: first argument must be an object", + "mapping_error": "evaluation_error: $mapValue: first argument must be an object", "mapping_version": "3f9a1c0b7d2e4a65", "expires_at": "2026-07-09T02:41:00Z" } @@ -483,8 +496,8 @@ Returns dead-lettered messages oldest first, paginated via an opaque `next_token | Field | Description | |-------|-------------| -| `reason` | Why the message was dead-lettered: the policy (`max_delivery_attempts` exhausted), the operator's `unblock` reason, or `mapping_failed` for a [payload mapping failure](#mapping-failures) | -| `mapping_error` | The mapping error message (truncated to 1,024 characters) — present when the payload mapping failed | +| `reason` | Why the message was dead-lettered: the policy (`max_delivery_attempts` exhausted), the operator's `unblock` reason, or `mapping_failed` for a [payload mapping failure](#mapping-failures) — whether it was dead-lettered by policy or by `unblock` | +| `mapping_error` | The mapping error as `: ` (truncated to 1,024 characters) — present when the payload mapping failed | | `mapping_version` | Version of the expression that produced the stored payload, or that failed on it | ### Redrive — re-enqueue dead-lettered messages @@ -511,7 +524,7 @@ The response reports a per-id outcome: { "results": [ { "id": "msg_5d2e…", "status": "redriven" }, - { "id": "msg_7a41…", "status": "mapping_failed", "mapping_error": "$mapValue: first argument must be an object" } + { "id": "msg_7a41…", "status": "mapping_failed", "mapping_error": "evaluation_error: $mapValue: first argument must be an object" } ] } ``` @@ -549,6 +562,8 @@ POST /v1/integrations/{integrationId}/outbound/messages/unblock For streams halted under `poison_policy: "block"`. **Skip equals dead-letter:** unblocking dead-letters the blocked head (recording the optional `reason`, max 500 characters) and emits `MSG_DEAD_LETTERED` — the message then becomes redrivable from the DLQ like any other dead-lettered item. The next message becomes the head and the stream resumes. +When the blocked head is a [failed-mapping item](#mapping-failures), it is dead-lettered with `reason: "mapping_failed"` instead, so a later redrive re-applies the current mapping to it. + The response reports `unblocked: true` with the `dead_lettered_id` of the skipped head, or `unblocked: false` as a safe no-op when the stream is not currently blocked. A late acknowledgement from the consumer also unblocks the stream naturally — no operator action needed. ## Monitoring From e847f58bd51d83e63cb9b1f11d3a3f321951c638 Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 25 Sep 2026 14:23:02 +0300 Subject: [PATCH 09/12] docs(integration-toolkit): document the 5 MiB cap on poll mapping output --- docs/integrations/integration-toolkit/pollable-outbound.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index 353e99b..9246de8 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -275,7 +275,7 @@ A poll mapping can carry an optional `jsonata_expression` that reshapes each eve - **When:** once, at **enqueue time**, right after the use case's `event_filter` has accepted the event. The mapped output is stored on the queue item, so every poll of that item returns the same payload — a lease lapse or a redelivery never re-evaluates the expression. - **Input:** the standardized event-catalog event, hydrated in full — the same root that `event_filter` sees, minus the internal `_downgrades` and `_automation_chain` keys. Field paths start at the top level of the [Core Event](/docs/integrations/core-events) (`_event_id`, `_event_time`, `meter_number`, …). - **Bindings:** `$env` (the organization's non-secret environment variables, including [Key/Value Maps](./key-value-maps.md)), `$mapValue` and `$mapKey`. No other bindings are available — in particular there is no `$now`, because the output must depend only on the event and the configuration. An expression that references any other `$`-binding is rejected on save. -- **Output:** must be a **JSON object**. An array, a scalar, `null`, or an undefined result is a mapping failure (`invalid_output`). +- **Output:** must be a **JSON object**. An array, a scalar, `null`, or an undefined result is a mapping failure (`invalid_output`). The mapped output is also capped at **5 MiB**; a larger result is a mapping failure with the same code (`mapped output exceeds 5 MiB`). - **Empty means raw.** An absent, empty, or whitespace-only `jsonata_expression` applies no transform, and the raw standardized event is delivered — the behavior of every poll use case that has no expression. The expression is validated on save: JSONata syntax, a maximum of 10,000 characters, and no bindings outside `$env`, `$mapValue` and `$mapKey`. The same checks run again at enqueue time, so an expression stored before these checks existed that does not pass them produces [failed items](#mapping-failures) rather than being silently skipped. @@ -359,7 +359,7 @@ Supply exactly one of `payload` or `event_id`. `$env`, `$mapValue` and `$mapKey` | `unknown_binding` | The expression uses a `$`-binding other than `$env`, `$mapValue` or `$mapKey` | | `evaluation_error` | The expression failed while running (for example, a missing key/value map) | | `timeout` | The evaluation exceeded the 500 ms time limit or the evaluation depth limit | -| `invalid_output` | The result is not a JSON object | +| `invalid_output` | The result is not a JSON object, or the mapped output exceeds 5 MiB | | `expression_too_long` | The expression exceeds 10,000 characters | A mapping error is a normal `200` response with `valid: false`. A `4xx` status means the request itself could not be served: From 02f6c444b1d12170b23c01c2e9be1c12b1df878a Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 2 Oct 2026 15:26:19 +0300 Subject: [PATCH 10/12] docs(monitoring): regenerate monitoring codes from the merged backend snapshot Adds MSG_LATE_ARRIVAL, moves MSG_ACKED and ACK_CONFIRMED to success and MSG_DEAD_LETTERED, MSG_EXPIRED_UNPOLLED and MSG_HEAD_BLOCKED to error, and picks up the other codes added on the backend since the last sync. --- .../integration-toolkit/monitoring/codes.md | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/docs/integrations/integration-toolkit/monitoring/codes.md b/docs/integrations/integration-toolkit/monitoring/codes.md index 57ae52e..42817d2 100644 --- a/docs/integrations/integration-toolkit/monitoring/codes.md +++ b/docs/integrations/integration-toolkit/monitoring/codes.md @@ -15,7 +15,7 @@ Every monitoring event carries a **code** and a **level**. The code says what happened; the level says how much you should care. Both are filterable in the Integration Hub's [Monitoring tab](./overview.md) and through the events API. -There are 80 codes. You are most likely here because you saw one in a failed +There are 85 codes. You are most likely here because you saw one in a failed event — find it below. :::tip @@ -31,6 +31,7 @@ Something failed and the event did not do what it was meant to do. These are wha |---|---| | `ATTACHMENT_NOT_FOUND` | The file no longer exists — it was removed between the event and the delivery | | `ATTRIBUTE_TYPE_MISMATCH` | An attribute value did not match the type declared in the entity schema | +| `CONDITIONAL_VARIANT_WRITE_FAILED` | Conditional pricing refused one variant write. The item is dropped and not retried; the rest of the batch and the run continue. details.code names the reason where pricing gave one — UNKNOWN_ERROR means it did not | | `DEPRECATED_ENDPOINT` | This endpoint version is deprecated | | `DIRECT_ENTITY_NOT_ALLOWED` | The entity is not permitted by the use case entity allowlist | | `DIRECT_PAYLOAD_INVALID` | The direct mode payload failed validation against the versioned payload schema | @@ -53,11 +54,14 @@ Something failed and the event did not do what it was meant to do. These are wha | `METER_READING_GROUP_RETRYING` | A batch write of meter readings failed and will be retried automatically — one event per attempt covering the whole group (reading_count and external_ids in details) | | `MISSING_REQUIRED_PARAM` | A required parameter is missing from the request | | `MISSING_UNIQUE_IDENTIFIERS` | The event is missing the unique identifier field(s) required to match an entity | +| `MSG_DEAD_LETTERED` | Outbound message moved to the dead-letter queue after exhausting delivery attempts, or via an operator skip | +| `MSG_EXPIRED_UNPOLLED` | Outbound message expired before being consumed — retention elapsed without a successful poll | +| `MSG_HEAD_BLOCKED` | Outbound stream halted by a poison head message (block policy) — requires operator unblock or consumer acknowledgement | | `OAUTH2_TOKEN_FAILURE` | Failed to obtain an OAuth2 access token | | `PAYLOAD_TOO_LARGE` | The payload exceeded the maximum size accepted by the receiving system | | `PRUNE_SCOPE_PARTIAL_FAILURE` | Scope pruning completed with some failures | | `RECURSION_DEPTH_EXCEEDED` | Maximum recursion depth was exceeded during processing | -| `RELATION_REF_ITEM_NOT_FOUND` | The relation_ref target entity exists but the referenced item/value could not be matched — skipped as non-retryable. Check the mapping configuration and the entity data. | +| `RELATION_REF_ITEM_NOT_FOUND` | The relation_ref value could not be matched on the target entity (invalid mapped value, or still no match after writing it to the target) — skipped as non-retryable. Check the mapping configuration and the entity data. | | `RELATION_REF_VALUE_UNDEFINED` | A relation_ref mapping value resolved to undefined — check the mapping expression | | `REQUIRED_PARAM_MISSING` | A param the use case marks as required resolved to nothing, so the delivery was stopped before anything was sent — see param_name | | `SECURE_PROXY_DISABLED` | The secure proxy use case is disabled | @@ -75,6 +79,7 @@ Something failed and the event did not do what it was meant to do. These are wha | `SIGNATURE_VERIFICATION_UNAVAILABLE` | The file service could not be reached to verify the request signature | | `STEP_DISABLED` | A request step's "run this step when" expression returned false, so this step and every step after it were skipped. This is the configuration working as written, not a fault. | | `TIMEOUT` | The operation timed out | +| `UNIQUE_ID_LOOKUP_UNRESOLVABLE` | A related entity created for this event still could not be found by its unique ID, so the relation was skipped instead of creating a duplicate. Check that the unique ID value type matches the entity schema. | | `UNIQUE_ID_MULTIPLE_MATCHES` | Multiple entities matched the unique ID | | `UNIQUE_ID_NOT_IN_SCHEMA` | The unique ID attribute is not defined in the entity schema | | `UNKNOWN_ERROR` | An unexpected error occurred during processing | @@ -90,9 +95,11 @@ Processing continued, but something needs a human eye — often a retry in fligh | Code | What it means | |---|---| | `ACK_TIMEOUT` | Acknowledgement timed out waiting for the ERP system | +| `CONDITIONAL_VARIANT_WRITE_WARNING` | One variant write succeeded with a warning. The variant is stored and the rest of the batch and the run continue. details.code names the warning | | `EXTERNAL_WARNING` | A warning span pushed by an external system via the external monitoring events endpoint. | | `FILE_PROXY_UPLOAD_RETRYING` | A file upload failed with a retryable error and will be retried automatically — one event per attempt | | `LOOKUP_UNMAPPED` | A value was not listed in a lookup table and its fallback was used — see lookup_name and lookup_key for the gap | +| `MSG_LATE_ARRIVAL` | An event arrived after the poll consumer had already received later events, so it was placed at the end of the stream instead of at its event time — it is delivered, but out of event-time order | | `SOFT_DELETED_ENTITY_MATCHED` | A soft-deleted entity matched the unique ID — it will be resurrected on upsert, or referenced as-is by a relation. Investigate why the ERP source is sending events for a deleted entity. | ## Success @@ -101,6 +108,8 @@ The event did what it was meant to do. Useful for confirming a sync actually lan | Code | What it means | |---|---| +| `ACK_CONFIRMED` | Acknowledgement was confirmed by the ERP system | +| `CONDITIONAL_VARIANTS_WRITTEN` | Conditional price variants were written for one imported entity: emitted once per chunk, with a per-outcome count and the variant ids in details | | `ENTITY_CREATED` | A new entity was created in epilot | | `ENTITY_DELETED` | An entity was deleted from epilot | | `ENTITY_NO_OP` | No changes were needed for the entity | @@ -110,6 +119,7 @@ The event did what it was meant to do. Useful for confirming a sync actually lan | `FILE_PROXY_UPLOADED` | The external system accepted the file | | `METER_READING_DELETED` | One or more meter readings were deleted — emitted once per batch, not per reading (reading_count and external_ids in details) | | `METER_READING_UPSERTED` | One or more meter readings were created or updated — emitted once per batch, not per reading (reading_count and external_ids in details) | +| `MSG_ACKED` | Outbound message delivered: the polling consumer acknowledged it and it was removed from the queue | | `PRUNE_SCOPE_COMPLETED` | Scope pruning completed successfully | | `WEBHOOK_DELIVERED` | Webhook was delivered successfully | @@ -119,17 +129,12 @@ Lifecycle markers rather than outcomes: a message was queued, a duplicate was ig | Code | What it means | |---|---| -| `ACK_CONFIRMED` | Acknowledgement was confirmed by the ERP system | | `ACK_PENDING` | Acknowledgement is pending from the ERP system | | `DUPLICATE_EVENT` | This event was already processed (duplicate) | | `EXTERNAL_INFO` | An informational span pushed by an external system via the external monitoring events endpoint. | | `FAN_OUT_EMPTY` | The split expression returned an empty list, so nothing was sent — expected for events that carry no relevant items | | `FILE_PROXY_UPLOAD_ENQUEUED` | A per-file upload was accepted for delivery during fan-out | -| `MSG_ACKED` | Outbound message acknowledged by the polling consumer and removed from the queue | -| `MSG_DEAD_LETTERED` | Outbound message moved to the dead-letter queue after exhausting delivery attempts, or via an operator skip | | `MSG_ENQUEUED` | Outbound message enqueued to the poll queue, awaiting consumption by the ERP | -| `MSG_EXPIRED_UNPOLLED` | Outbound message expired before being consumed — retention elapsed without a successful poll | -| `MSG_HEAD_BLOCKED` | Outbound stream halted by a poison head message (block policy) — requires operator unblock or consumer acknowledgement | ## Status-code families Some codes are generated from an upstream response rather than drawn from the fixed list above. From dddc86ae144c79a486f3e70bf955b8091c6a4f5d Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 2 Oct 2026 15:26:19 +0300 Subject: [PATCH 11/12] docs(monitoring): classify ACK_CONFIRMED as success --- docs/integrations/integration-toolkit/monitoring/acks.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/docs/integrations/integration-toolkit/monitoring/acks.md b/docs/integrations/integration-toolkit/monitoring/acks.md index 470b8e5..566f50f 100644 --- a/docs/integrations/integration-toolkit/monitoring/acks.md +++ b/docs/integrations/integration-toolkit/monitoring/acks.md @@ -81,11 +81,12 @@ filterable in the Monitoring tab: | Code | Level | When | |---|---|---| | `ACK_PENDING` | info | The event was delivered and epilot is waiting for the acknowledgement | -| `ACK_CONFIRMED` | info | Your acknowledgement arrived | +| `ACK_CONFIRMED` | success | Your acknowledgement arrived | | `ACK_TIMEOUT` | warning | No acknowledgement within the timeout window | -`ACK_PENDING` and `ACK_CONFIRMED` are **info**-level: they are lifecycle markers, not -outcomes, so they are counted in total events but deliberately excluded from the +`ACK_PENDING` is **info**-level: a lifecycle marker, not an outcome, so it is counted +in total events but deliberately excluded from the success rate. `ACK_CONFIRMED` is a +**success** — your system confirmed it processed the event, so it counts towards the success rate. `ACK_TIMEOUT` is a **warning** — the delivery itself worked, so it is not an error on epilot's side, but something on yours needs attention. From 3c76545bd71adfcd961acaeadf5229cadb861564 Mon Sep 17 00:00:00 2001 From: Karl Rankla Date: Fri, 2 Oct 2026 15:26:19 +0300 Subject: [PATCH 12/12] docs(integration-toolkit): align poll delivery docs with the shipped backend Correct the DLQ reason values, the 5 MiB error message, redrive and unblock retention, which poll mapping a redrive re-applies, and the failure detail fields. Document the depth limit, retries on environment load failures, batches ending before a failed item, the preview 403, and the v1 and v2 outbound validation differences. --- .../integration-toolkit/configuration.md | 2 ++ .../integration-toolkit/pollable-outbound.md | 35 ++++++++++--------- 2 files changed, 21 insertions(+), 16 deletions(-) diff --git a/docs/integrations/integration-toolkit/configuration.md b/docs/integrations/integration-toolkit/configuration.md index 1d129c8..31f1050 100644 --- a/docs/integrations/integration-toolkit/configuration.md +++ b/docs/integrations/integration-toolkit/configuration.md @@ -215,6 +215,8 @@ An event the filter rejects is not processed by the use case at all: no [Pollabl | `jsonata_expression` | string | For `webhook` delivery | JSONata expression to transform the event payload. Required for `webhook` (evaluated by the webhook service). Optional for `poll`: evaluated at enqueue time against the standardized event-catalog event with `$env` / `$mapValue` / `$mapKey`, must return a JSON object, and an empty value delivers the raw event — see [Payload Mapping](./pollable-outbound.md#payload-mapping). Rejected for `file_proxy` delivery | | `delivery` | object | Yes | How the event is delivered — discriminated on `type`: `webhook`, `poll`, or `file_proxy` | +Outbound configurations are validated on save. The v1 use case endpoints require a `jsonata_expression` on every `webhook` mapping. The v2 integration upsert (`POST` / `PUT /v2/integrations`) validates an outbound use case only when it is new or its configuration changed, so configurations stored before a rule existed can be re-sent unchanged. It also accepts a `webhook` mapping with **no** `jsonata_expression` at all, for configurations that predate that requirement — such a mapping never enables or updates its webhook. An empty expression and invalid JSONata are rejected on both versions. + #### Delivery Types **Webhook delivery (push):** the event payload is transformed with the mapping's `jsonata_expression` and pushed to a pre-configured webhook (epilot Webhooks): diff --git a/docs/integrations/integration-toolkit/pollable-outbound.md b/docs/integrations/integration-toolkit/pollable-outbound.md index 9246de8..1607492 100644 --- a/docs/integrations/integration-toolkit/pollable-outbound.md +++ b/docs/integrations/integration-toolkit/pollable-outbound.md @@ -275,12 +275,14 @@ A poll mapping can carry an optional `jsonata_expression` that reshapes each eve - **When:** once, at **enqueue time**, right after the use case's `event_filter` has accepted the event. The mapped output is stored on the queue item, so every poll of that item returns the same payload — a lease lapse or a redelivery never re-evaluates the expression. - **Input:** the standardized event-catalog event, hydrated in full — the same root that `event_filter` sees, minus the internal `_downgrades` and `_automation_chain` keys. Field paths start at the top level of the [Core Event](/docs/integrations/core-events) (`_event_id`, `_event_time`, `meter_number`, …). - **Bindings:** `$env` (the organization's non-secret environment variables, including [Key/Value Maps](./key-value-maps.md)), `$mapValue` and `$mapKey`. No other bindings are available — in particular there is no `$now`, because the output must depend only on the event and the configuration. An expression that references any other `$`-binding is rejected on save. -- **Output:** must be a **JSON object**. An array, a scalar, `null`, or an undefined result is a mapping failure (`invalid_output`). The mapped output is also capped at **5 MiB**; a larger result is a mapping failure with the same code (`mapped output exceeds 5 MiB`). +- **Output:** must be a **JSON object**. An array, a scalar, `null`, or an undefined result is a mapping failure (`invalid_output`). The mapped output is also capped at **5 MiB**; a larger result is a mapping failure with the same code (`The mapped output exceeds 5 MiB`). An output that cannot be serialized as JSON — for example one containing a function — is an `invalid_output` failure too. - **Empty means raw.** An absent, empty, or whitespace-only `jsonata_expression` applies no transform, and the raw standardized event is delivered — the behavior of every poll use case that has no expression. The expression is validated on save: JSONata syntax, a maximum of 10,000 characters, and no bindings outside `$env`, `$mapValue` and `$mapKey`. The same checks run again at enqueue time, so an expression stored before these checks existed that does not pass them produces [failed items](#mapping-failures) rather than being silently skipped. -Each evaluation is guarded by a 500 ms time limit and an evaluation depth limit. Exceeding either is a mapping failure with the code `timeout`. +Each evaluation is guarded by a 500 ms time limit and an evaluation depth limit of 500. Exceeding either is a mapping failure with the code `timeout`. + +Only failures the expression itself causes are mapping failures. If the organization's environment variables cannot be loaded, nothing is recorded on an item: the event is retried and mapped once the environment is reachable. ### Mapping version @@ -290,16 +292,16 @@ Because the transform runs at enqueue time, **changing the expression affects ne ### Mapping failures -A failure to evaluate — a runtime error, a timeout, or an output that is not an object — does not drop the event. The item is still enqueued at its normal position in the stream, marked as failed, with the raw payload kept alongside it. At enqueue time epilot emits the `MAPPING_EXPRESSION_FAILED` monitoring error, with the use case, event id, event name, `mapping_version` and the error message in its context. +A failure to evaluate — a runtime error, a timeout, or an output that is not an object — does not drop the event. The item is still enqueued at its normal position in the stream, marked as failed, with the raw payload kept alongside it. At enqueue time epilot emits the `MAPPING_EXPRESSION_FAILED` monitoring error once, with `use_case_id`, `event_id`, `event_name`, `message_id`, `mapping_version`, `error_code` and the error message in its detail. The error is recorded as `mapping_error` in the form `: ` — for example `evaluation_error: $mapValue: first argument must be an object` — using the same codes as the [preview endpoint](#preview-a-transform), truncated to 1,024 characters. -A failed item is **never delivered to the consumer**. When it reaches the head of the stream it is handled immediately — without waiting for `max_delivery_attempts`, because evaluating the same expression against the same event would fail the same way every time. What happens next follows the use case's `poison_policy`: +A failed item is **never delivered to the consumer**. When it reaches the head of the stream it is handled immediately — without waiting for `max_delivery_attempts`, because evaluating the same expression against the same event would fail the same way every time. A poll batch never reaches past a failed item: when one sits behind deliverable messages, the batch ends before it (with `has_more: true`), and it is handled once it is the head. What happens next follows the use case's `poison_policy`: | Policy | What happens to a failed item at the head | |--------|-------------------------------------------| | `dead_letter` (default) | It moves straight to the [dead-letter queue](#dead-letter-queue-and-operator-actions) with `reason: "mapping_failed"` and `delivery_attempts: 0` (it was never leased), and the stream moves on. `MSG_DEAD_LETTERED` is emitted with `reason`, `mapping_error` and `mapping_version` in its detail | -| `block` | The stream halts on it. `MSG_HEAD_BLOCKED` is emitted with `reason: "mapping_failed"` in its detail. Release it with [`unblock`](#unblock--skip-a-blocked-head), which dead-letters it with `reason: "mapping_failed"`. A consumer acknowledgement cannot release it, because it is never delivered | +| `block` | The stream halts on it. `MSG_HEAD_BLOCKED` is emitted with `reason: "mapping_failed"` in its detail. Release it with [`unblock`](#unblock--skip-a-blocked-head), which dead-letters it (with `reason: "mapping_failed"` unless you give your own). A consumer acknowledgement cannot release it, because it is never delivered | To recover: @@ -351,7 +353,7 @@ Supply exactly one of `payload` or `event_id`. `$env`, `$mapValue` and `$mapKey` | `output` | The mapped output — present when `valid` is `true` | | `error` | Present when `valid` is `false`: a `code`, a `message`, and for syntax errors the `position` in the expression | | `mapping_version` | The version this expression would stamp on messages | -| `input` | The hydrated event the expression ran against — returned when the request used `event_id`, on failures as well as successes | +| `input` | The hydrated event the expression ran against, without the internal `_downgrades` and `_automation_chain` keys — returned when the request used `event_id`, on failures as well as successes | | Error `code` | Meaning | |--------------|---------| @@ -367,6 +369,7 @@ A mapping error is a normal `200` response with `valid: false`. A `4xx` status m | Status | When | |--------|------| | `400` | The expression is blank, both or neither of `payload` and `event_id` are given, or `event_id` is given without `event_catalog_event` | +| `403` | The token lacks `integration:view` on the integration | | `404` | The integration is unknown or belongs to another organization, or the event catalog does not know the requested event | ### Webhook and poll mode are not interchangeable @@ -432,7 +435,7 @@ A few details make this robust: - Undelivered items expire after the mapping's `retention_days` (default 30, max 90), counted from enqueue time. - Changing `retention_days` affects **new items only** — already-enqueued items keep the TTL computed at enqueue time. - An expired item is never delivered: the poll API filters expired items even before the storage layer reaps them. Each expiry of an item that was **never consumed** emits an `MSG_EXPIRED_UNPOLLED` monitoring error, so silent data loss is always visible. -- Dead-lettered items get a **re-armed retention window** at dead-letter time — a full `retention_days` from that moment — giving operators the whole window to redrive instead of whatever sliver remained. +- Dead-lettered items get a **re-armed retention window** at dead-letter time — a full `retention_days` from that moment when the poison policy dead-letters them, and 30 days when an operator [unblock](#unblock--skip-a-blocked-head) does — giving operators the whole window to redrive instead of whatever sliver remained. ## Poison Messages: `dead_letter` vs `block` @@ -473,7 +476,7 @@ Returns dead-lettered messages oldest first, paginated via an opaque `next_token "enqueued_at": "2026-06-08T22:10:00Z", "dead_lettered_at": "2026-06-09T03:00:00Z", "delivery_attempts": 5, - "reason": "max_delivery_attempts exhausted", + "reason": "max_delivery_attempts_exhausted", "expires_at": "2026-07-09T03:00:00Z" }, { @@ -496,9 +499,9 @@ Returns dead-lettered messages oldest first, paginated via an opaque `next_token | Field | Description | |-------|-------------| -| `reason` | Why the message was dead-lettered: the policy (`max_delivery_attempts` exhausted), the operator's `unblock` reason, or `mapping_failed` for a [payload mapping failure](#mapping-failures) — whether it was dead-lettered by policy or by `unblock` | +| `reason` | Why the message was dead-lettered: `max_delivery_attempts_exhausted`, `mapping_failed` for a [payload mapping failure](#mapping-failures), or for an [`unblock`](#unblock--skip-a-blocked-head) the operator's reason — `operator_skip` (or `mapping_failed` for a failed-mapping head) when none was given | | `mapping_error` | The mapping error as `: ` (truncated to 1,024 characters) — present when the payload mapping failed | -| `mapping_version` | Version of the expression that produced the stored payload, or that failed on it | +| `mapping_version` | Version of the expression that produced the stored payload, or that failed on it. Absent for raw payloads | ### Redrive — re-enqueue dead-lettered messages @@ -533,16 +536,16 @@ The response reports a per-id outcome: |----------|---------| | `redriven` | Re-enqueued at the tail | | `not_found` | Unknown id, or the entry was concurrently redriven or expired | -| `mapping_failed` | Re-applying the mapping failed again. The entry stays in the DLQ with its `mapping_error` and `mapping_version` updated; `mapping_error` is also returned in the result | +| `mapping_failed` | Re-applying the mapping failed again. The entry stays in the DLQ, now with `reason: "mapping_failed"`, the new `mapping_error` and `mapping_version`, a new `dead_lettered_at`, and a fresh 30-day retention window; `mapping_error` is also returned in the result | -The redriven copy is re-enqueued with **zero delivery attempts and a fresh retention window**; the original DLQ entry is removed. +The redriven copy is re-enqueued with **zero delivery attempts and a fresh 30-day retention window** (the default, whatever the use case's `retention_days`); the original DLQ entry is removed. How the payload of the redriven copy is chosen: - By default (`reapply_mapping: false`) the stored payload is re-sent unchanged — the behavior before payload mapping existed. - With `reapply_mapping: true` the current expression is evaluated again against the kept raw payload. Use it after fixing an expression, so already dead-lettered messages go out in the corrected shape. -- An entry that was dead-lettered **because its mapping failed** (`reason: "mapping_failed"`) is always re-mapped with the current configuration, whatever `reapply_mapping` says — it has no usable mapped payload to re-send. -- When re-mapping and the use case no longer has an expression, the raw payload is delivered. +- An entry **whose mapping failed** (it carries a `mapping_error`) is always re-mapped with the current configuration, whatever `reapply_mapping` says — it has no usable mapped payload to re-send. +- "Current configuration" means the **enabled** poll mapping of the use case, and only while the use case itself is enabled. When re-mapping and there is no such mapping, or it has no expression, the raw payload is delivered. :::caution Redrive ordering A redriven message is re-enqueued at the **tail** with a new id and sequence — it is delivered out of its original per-entity order, because the stream has moved on. This is inherent to redrive and matches SQS DLQ semantics. If your consumer is order-sensitive, reconcile redriven messages explicitly (e.g. compare against current entity state). @@ -562,7 +565,7 @@ POST /v1/integrations/{integrationId}/outbound/messages/unblock For streams halted under `poison_policy: "block"`. **Skip equals dead-letter:** unblocking dead-letters the blocked head (recording the optional `reason`, max 500 characters) and emits `MSG_DEAD_LETTERED` — the message then becomes redrivable from the DLQ like any other dead-lettered item. The next message becomes the head and the stream resumes. -When the blocked head is a [failed-mapping item](#mapping-failures), it is dead-lettered with `reason: "mapping_failed"` instead, so a later redrive re-applies the current mapping to it. +Without a `reason`, the entry records `operator_skip`, or `mapping_failed` when the blocked head is a [failed-mapping item](#mapping-failures). Either way a failed-mapping entry is re-mapped by a later redrive. The DLQ entry from an unblock gets a 30-day retention window. The response reports `unblocked: true` with the `dead_lettered_id` of the skipped head, or `unblocked: false` as a safe no-op when the stream is not currently blocked. A late acknowledgement from the consumer also unblocks the stream naturally — no operator action needed. @@ -580,7 +583,7 @@ Poll-queue message lifecycle events flow through the standard monitoring pipelin | `MSG_DEAD_LETTERED` | error | A message exhausted `max_delivery_attempts` under the `dead_letter` policy, or an operator skipped a blocked head (includes `delivery_attempts` in the event detail) | | `MSG_HEAD_BLOCKED` | error | The stream halted on a poisoned head under the `block` policy — emitted **once per blocked episode**, not on every poll (includes `delivery_attempts` in the event detail) | | `MSG_LATE_ARRIVAL` | warning | An event arrived after the consumer had moved past its position and was re-keyed to the tail of the stream (includes `original_sequence_time` and `rekeyed_sequence_time` in the event detail) — see [Late arrivals](#late-arrivals) | -| `MAPPING_EXPRESSION_FAILED` | error | The poll mapping's `jsonata_expression` failed at enqueue time (includes `use_case_id`, `event_id`, `event_name`, `mapping_version` and the error message) — see [Mapping failures](#mapping-failures) | +| `MAPPING_EXPRESSION_FAILED` | error | The poll mapping's `jsonata_expression` failed at enqueue time (includes `use_case_id`, `event_id`, `event_name`, `message_id`, `mapping_version`, `error_code` and the error message) — see [Mapping failures](#mapping-failures) | When a failed mapping reaches the head, `MSG_DEAD_LETTERED` or `MSG_HEAD_BLOCKED` carries `reason: "mapping_failed"` in its detail, so the mapping failure and its consequence for the stream can be told apart from ordinary poison messages.