diff --git a/.agents/diary/how-to-build-modules.md b/.agents/diary/how-to-build-modules.md new file mode 100644 index 00000000..bf1edadd --- /dev/null +++ b/.agents/diary/how-to-build-modules.md @@ -0,0 +1 @@ +[unrelated bug] `docs/guides/getting-started.md` declares `const { quotes } = service.load();` twice in the gateway server example. This makes the copied example fail to compile and should be fixed in a separate documentation task. diff --git a/.agents/diary/queues-prototype.md b/.agents/diary/queues-prototype.md new file mode 100644 index 00000000..cb5b64ff --- /dev/null +++ b/.agents/diary/queues-prototype.md @@ -0,0 +1,34 @@ +# Queue prototype diary + +- **Unrelated bug.** `prisma-composer dev` printed React "Invalid hook call" + warnings twice before all local services became ready. The queue application + still ran correctly, so investigating the CLI rendering issue is outside this + prototype. + +- **Scope tangent.** Composer should validate that `PRISMA_WORKSPACE_ID` matches + the service token workspace before creating a project. A mismatch can make + project discovery miss an existing project while creation still succeeds in + the token workspace, producing an unwanted duplicate project. + +- **Scope tangent.** The deployed demo relies on the `--name queues-demo` + override while its root Module is named `queues-example`. A later deploy that + omits the override creates a second project; the example should make its + stable production name harder to omit. + +- **Unrelated bug.** A deploy that changed an artifact together with its input + document or dependency URL started the new artifact before the new binding was + active. The queue service entered a restart loop with its previous input, and + the dispatcher called previous queue and consumer routes. Reconciliation after + the bindings settled fixed both cases; deployment ordering should prevent this + state. + +- **Scope tangent.** Promoting a new Compute deployment leaves every previous + deployment running. This is unsafe for continuously polling drivers: six queue + dispatcher revisions competed for work, and older revisions used stale + bindings. Composer needs deployment retirement or a revision-fencing mechanism + before always-running drivers are production-safe. + +- **Unrelated bug.** The queue specification's relative links still use its old + `projects/queue-module/` depth after the file moved under `.drive/projects/`. + They currently resolve below `.drive/` instead of the repository `docs/` + directory and should be corrected separately. diff --git a/.drive/projects/queue-module/spec.md b/.drive/projects/queue-module/spec.md new file mode 100644 index 00000000..58011379 --- /dev/null +++ b/.drive/projects/queue-module/spec.md @@ -0,0 +1,702 @@ +# Summary + +Prisma Composer will provide a first-party queue capability for durable, +at-least-once work delivery on Prisma Cloud. It consists of a Queue Module and an +always-running dispatcher driver. Version one stores messages in Prisma Postgres +and pushes each queue's batches over authenticated HTTP to one statically wired +Composer consumer service. Any number of application services may produce to the +same queue. Later phases may add statically declared, authenticated external +producers and consumers. Publish and subscribe remains a separate capability and +will not change this queue's single-consumer delivery contract. + +## Prototype status + +The first walking skeleton is implemented. It deliberately proves the narrowest +durable path before the full version-one contract: + +- `defineQueues` creates the static typed queue catalogue. +- `queues` provisions one Postgres database and the queue service. +- `queueProducer` gives application services a typed `send` client. +- `queueConsumer` and `serveQueues` register typed handlers on an application + service. +- `queueDispatcher` runs separately, claims one message with a lease, pushes it + to one consumer, and completes or releases it. +- Each queue definition can set `maxAttempts` and a serializable fixed retry + delay. The queue service persists the exact next availability time and moves + exhausted messages into a terminal failed state. +- `examples/queues` proves one Compute service can produce and consume its own + messages through the durable queue, including a real fail-once retry. + +The prototype does not yet implement batches, exponential retry, failed-message +inspection and replay, pause and resume, or operational data. The detailed +requirements below remain the target design, not a claim about the current +implementation. + +# Context + +The primary consumers are application engineers building Prisma Apps who need +background work without operating a separate queue product. Consumer implementers +need a typed, Cloudflare-like handler model. On-call engineers need durable state +and clear attempt history when a message is delayed, retried, or delivered more +than once. + +This project starts with work queues, not publish and subscribe. Any number of +producer services may enqueue work, but each queue has exactly one logical +consumer service. A service may consume several queues, and the deployment target +may run several replicas of that service, but replicas are not separate Composer +consumer registrations. Different required outcomes belong on different queues +until publish and subscribe is designed as a separate capability. + +## At a glance + +Each producer receives a typed client through `service.load()` and commits a +message to the Queue Module. A dispatcher claims available messages from Postgres +and pushes a bounded batch to the queue's one consumer service. A successful +result completes the message; a failed or missing result makes it available for a +later attempt according to the queue's retry policy. + +```mermaid +flowchart LR + P1[Producer service A] -->|Typed enqueue client| Q[Queue service] + P2[Producer service B] -->|Typed enqueue client| Q + Q --> D[(Prisma Postgres)] + Q -->|Internal control port| X[Dispatcher] + X -->|One HTTP attempt| C[One logical consumer service] +``` + +Every attempt for a queue goes to the same logical consumer service. Runtime +replicas may share requests behind that service endpoint, but the queue topology +contains one queue-to-consumer edge. + +## Problem + +Composer has reusable modules for scheduled work, storage, streams, and email, but +no durable work queue. Applications can call a service directly, but this couples +the caller to the callee's availability and gives the application no durable +retry, backpressure, or attempt history. + +Prisma Compute is request-oriented and may scale to zero. Postgres can persist +messages and coordinate competing dispatchers, but it cannot by itself wake a +suspended Compute service when a new message or delayed retry becomes available. +The design therefore needs an explicit dispatcher lifecycle as well as durable +message state. + +The authoring experience should feel familiar to Cloudflare Queues while delivery +uses a Google Cloud Tasks style HTTP push. In-graph application code must not know +about database rows, leases, dispatcher instances, service URLs, or authentication +keys. + +## Approach + +The queue capability is composed from ordinary Composer primitives. The Queue +Module owns a Prisma Postgres database and a queue service. It exposes a typed +producer port and an internal dispatch-control port. A separate dispatcher driver +depends on the control port and the application-owned consumer port for each +queue. This remains target-specific functionality and does not add a queue +primitive to Composer core. + +A queue definition is static TypeScript data shared by the producer client, +consumer contract, and handler helper. It contains the payload schema and the +delivery policy that must be known at deployment. Consumer registrations are +also part of the static Composer graph, with exactly one consumer binding for each +queue. Only an explicit queue-to-consumer edge registers a consumer; exposing a +compatible but unwired consumer port does not. Any number of producer edges may +target the same queue. If topology wiring assigns a second distinct logical +consumer service, graph loading fails before execution or deployment and identifies +the queue and both conflicting services. Operational message state, attempts, +availability times, and leases are durable database state; consumer code and +deployment configuration are not stored in the database. + +Queue payloads must be JSON-compatible values. A payload may contain at most +128 KiB after JSON serialization and UTF-8 encoding. Schema validation runs +before this byte-size check, and either failure rejects the enqueue without +creating a message row. Queue envelope metadata does not count toward the +payload limit. Applications should store larger values elsewhere and enqueue a +reference to them. + +Each delivery contains at most the queue's configured batch size, which defaults +to 10 messages and accepts values from 1 through 100. The complete encoded HTTP +request, including envelope metadata, may not exceed 1 MiB. The dispatcher stops +adding messages when either bound is reached, never splits one message across +requests, and does not wait for a batch to fill before dispatching available +work. + +Each delivery is a bounded HTTP request. The dispatcher atomically claims +available messages with a lease, commits the claim, and only then calls the +consumer. It never holds a database transaction open while application code +runs. The consumer timeout defaults to 30 seconds and accepts values from 1 +through 50 seconds. The lease is not separately configurable; it lasts for the +consumer timeout plus 30 seconds. Lease expiry recovers work after a dispatcher +or consumer failure. Version one does not expose renewable leases or heartbeats +to consumer code. + +The consumer records per-message acknowledgement and retry decisions in memory +while it handles the batch. When the handler finishes, its helper returns every +decision in the response to the dispatcher's original HTTP request. The dispatcher +then commits the results in one short transaction. An ordinary thrown handler +error preserves decisions already recorded and retries undecided messages. A +timeout aborts the dispatcher's request but does not release the lease. A timeout, +invalid or lost response, or process crash applies no decisions, so the leased +batch becomes available again after lease expiry. + +One dispatcher driver remains running at all times. It serves every queue owned by +the Queue Module and processes due work without producer or web traffic. It drains +while work is available, then sleeps until the earliest known availability time +with a bounded fallback interval. Postgres remains the source of truth behind the +Queue Module's control port, so an enqueue wake signal may reduce latency but is +never required for correctness. ⚠️ **OF1** (always-running dispatcher support) +The Prisma Cloud target must support keeping one dispatcher instance running; +the driver must not emulate durability with `waitUntil` or a chain of +self-requests. + +Each queue has one queue-wide concurrency limit, which defaults to 10 and accepts +values from 1 through 100. One slot represents one active batch HTTP request, +regardless of how many messages the batch contains. Every slot targets the +queue's one logical consumer service. The dispatcher stops claiming when every +slot is in use and resumes when a request finishes or a lease expires. A timed-out +request releases its slot, but its messages remain leased until expiry. Version +one has no zero-value pause or per-replica capacity settings. + +Version one makes no message ordering guarantee. The dispatcher may prefer +currently available messages by availability time and creation order when +claiming work, but this is an implementation detail rather than a consumer +contract. Concurrent batches, retries, lease expiry, and replay can all cause a +later message to finish before an earlier one. Consumers must not depend on +first-in-first-out or per-key ordering. + +The initial delivery is attempt one. A queue defaults to five maximum attempts +and accepts values from 1 through 100; version one does not support unlimited +attempts. The dispatcher increments the attempt count in the same transaction +that claims the message. A claim therefore consumes an attempt even if the +dispatcher fails before the consumer receives the request. Producer validation +and enqueue failures do not consume attempts. + +Retry delays are static policy data built with `fixedBackoff` or +`exponentialBackoff`; arbitrary user callbacks are not accepted. +`exponentialBackoff` defaults to a five-second initial delay, a 15-minute +maximum, a factor of two, and equal jitter. For a calculated exponential limit +`d = min(max, initial × factor^(failed attempt - 1))`, no jitter waits exactly +`d`, full jitter selects from zero through `d`, and equal jitter selects from +`d / 2` through `d`. `fixedBackoff` uses the same delay after every failed +attempt. Delay values use one-second precision and must be from one second +through 24 hours. An exponential maximum cannot be shorter than its initial +delay, and its finite factor must be from 1.1 through 10. The helpers return +serializable descriptors, and invalid values fail while loading the graph. The +dispatcher calculates and persists each message's exact next availability time +once; restart does not recalculate it. + +A consumer may call `message.retry({ delay: "30s" })` to replace the queue +policy with one concrete delay for that message's next attempt. The delay must +be from one second through 24 hours. It does not reset the attempt count or +change the queue policy. A plain `message.retry()`, thrown error, timeout, +invalid response, or lost response uses the queue policy. The override takes +effect only when the dispatcher receives and commits the consumer's valid +response. + +Retention for an active message is measured from its original enqueue time and +does not reset on retry. It defaults to seven days and accepts values from one +hour through 30 days. When that period ends, an available, delayed, or retrying +message enters the terminal failed state with reason `expired` instead of being +silently deleted. Expiry during an active consumer request does not interrupt +the handler: a valid acknowledgement still completes the message, while a retry +or missing result produces the expired failure. + +Completed payloads are removed immediately. Their metadata and attempt history +remain for 24 hours by default, configurable from immediate deletion through +seven days. Failed records retain their payloads for replay for 30 days by +default, configurable from one through 90 days. Terminal retention starts when +the message enters the completed or failed state. A message and its attempt +history are deleted together after that period. Replay creates a new message +with a new active retention period. + +After the configured maximum attempts, a message enters an immutable terminal +failed state in its original queue and is no longer eligible for dispatch. Version +one does not route exhausted messages into a separate dead-letter queue. A +separate operational port can replay one retained failed message at a time. Replay +validates the original payload against the current queue schema and creates a new, +immediately available message with a new identifier, zero attempts, the current +retry policy, and a `replayedFrom` link. The original failed record remains +unchanged. The replay operation is durably idempotent. + +The same separately wired operational port provides paginated failed-message +inspection and graceful queue pause and resume controls. A pause transaction is +a durable barrier for new claims: a claim committed before the pause may still +be delivered, but no claim may commit after the pause does. Enqueue remains +available, active HTTP requests and leases continue, valid results are committed, +and retry and retention clocks keep advancing. Resume wakes the dispatcher to +drain accumulated work within the normal concurrency limit. Both operations are +idempotent. Runtime control state is durable, not deployment configuration. + +Failed-message listing returns metadata only: identifier, enqueue and failure +times, failure reason, attempt count, last consumer, retention expiry, and replay +status. Looking up one failed message returns its typed payload, full attempt +history, consumer and timing data for each attempt, safe error summaries, and +replay links. Payload access is therefore an explicit privileged operation and +never part of a list response. Listing uses an opaque cursor, defaults to 50 +records, permits at most 100, and orders by descending failure time followed by +message identifier. It can filter by failure time range, failure reason, replay +status, and last consumer. Version one does not provide offset pagination, +payload search, or arbitrary query filters. + +For version one, the Queue Compute service exposes current queue status and +recent activity through its separately wired operational RPC-over-HTTP port. A +Composer service wired to that port receives a typed client and the existing +per-binding service key; unwired callers are rejected. A later Prisma Console or +control-plane adapter can consume this contract, but that integration is outside +this phase. The Queue Module does not store general-purpose metric time series or +initiate outbound product analytics or metric export. + +Queue status returns the authored queue name, `running` or `paused` state, +optional pause time, counts for available, delayed, leased, and retained failed +messages, active and maximum batch request concurrency, optional age in +milliseconds of the oldest available message, optional last activity time, and +an `asOf` time. Optional values are absent when they do not apply. Timestamps use +RFC 3339, and counts are non-negative integers. + +The activity feed records normal enqueue and delivery work as batch summaries. +A delivery result reports acknowledged, retrying, and failed counts instead of +writing one successful event per message. Terminal message failure, expiry, +replay, pause, and resume are individual events. Per-message attempt history +remains available through failed-message detail rather than being duplicated in +the general feed. + +Activity is stored in a dedicated `queue_activity` table in the Queue Module's +Postgres database. Common indexed fields carry the event identifier, queue, +event type, occurrence and expiry times, and optional batch, message, and +consumer identifiers; type-specific safe metadata is stored separately from +those filter fields. Whenever an event describes a durable queue mutation, the +event and mutation commit in the same transaction. Activity is retained for +seven days by default, configurable from one through 30 days. Expired activity +is excluded from API results even before asynchronous physical cleanup. Activity +never contains message payloads. + +Every activity event has an opaque identifier, authored queue name, +discriminating type, and RFC 3339 occurrence time. The supported types and +additional fields are: + +- `enqueued`: message count. +- `delivery_started`: batch identifier, consumer identifier, and message count. +- `delivery_result`: batch and consumer identifiers, duration in milliseconds, + acknowledged, retrying, and failed counts, and one of `valid_response`, + `handler_error`, `timeout`, `network_error`, `http_error`, or + `invalid_response`. +- `message_failed`: message identifier, attempt count, and failure reason. +- `message_expired`: message identifier and enqueue time. +- `message_replayed`: original and new message identifiers. +- `queue_paused` and `queue_resumed`: no additional fields. + +Activity listing uses an opaque cursor, orders newest first, returns 50 events by +default and at most 100, and filters by occurrence time range, event type, +consumer identifier, and message identifier. It does not provide payload or +arbitrary text search. + +`queue_activity` is the Queue Module's only analytics record. It remains in the +user's Queue Module database and is available only through the explicitly wired +operational port. The queue service and dispatcher do not send product analytics, +internal error events, or queue activity to Prisma or any third-party telemetry +endpoint. + +Dedicated consumer Compute services are the documented default because they +isolate scaling, failures, and logs from the web application. An existing +Composer web service may instead be the queue's one logical consumer. The +application provisions each queue's consumer before the dispatcher and wires its +delivery port into the driver. A consumer may also depend on the Queue Module's +producer port, including for the same queue it consumes. Separating the queue +service from the dispatcher preserves this order and avoids a dependency cycle: +queue service, then consumer services, then dispatcher. + +Direct external enqueue is not part of version one. A later phase may allow a +queue to target one statically declared external HTTPS consumer instead of a +Composer consumer service. A separate external producer capability may expose a +queue-specific HTTPS endpoint to callers outside the Composer graph. Each external +producer's access is declared by name in the deployment topology, and Composer +provisions queue-specific credentials. External enqueue uses the same payload +schema, size limit, and durable-commit guarantee as an internal `send()` call. +Runtime self-registration, dynamic credential creation, and unauthenticated +enqueue remain outside the design. + +The selected authoring foundation separates pure queue definitions from Composer +nodes. `defineQueues` creates only a static typed catalog. Standalone +`queueProducer` and `queueConsumer` factories derive a producer dependency and a +consumer exposure contract from that catalog. The `queues` factory creates the +infrastructure Module, while `queueDispatcher` creates the separate delivery +driver. `serveQueues` is the consumer runtime adapter that turns typed handlers +into a Fetch handler; it provisions nothing. The fixed-backoff queue definition +syntax is implemented; batch handling and the remaining signatures stay open. + +> _Illustrative — names and exact syntax remain part of the authoring API design:_ +> +> ```ts +> const thumbnails = defineQueue({ +> message: type({ imageId: "string" }), +> retry: { +> maxAttempts: 5, +> delay: exponentialBackoff({ +> initial: "5s", +> max: "15m", +> factor: 2, +> jitter: "equal", +> }), +> }, +> }); +> +> const { thumbnails: queue } = service.load(); +> await queue.send({ imageId: "image-123" }); +> +> serveQueue(consumerService, thumbnails, async (batch) => { +> for (const message of batch.messages) { +> await createThumbnail(message.body); +> message.ack(); +> } +> }); +> ``` + +> _Illustrative — topology names and exact factory signatures remain part of the +> authoring API design:_ +> +> ```ts +> const queues = provision(queueModule({ queues: { thumbnails } })); +> const worker = provision(thumbnailWorker, { +> deps: { thumbnails: queues.thumbnails }, +> }); +> provision(queueDispatcher({ queues: { thumbnails } }), { +> deps: { +> control: queues.dispatch, +> consumers: { thumbnails: worker.queue }, +> }, +> }); +> ``` + +# Requirements + +## Functional Requirements + +- **FR1. Typed queue definition.** An application can define a queue from a + Standard Schema payload validator. The same definition types producer calls, + consumer handlers, and delivery validation. +- **FR2. Typed producer binding.** Any number of services wired to the Queue + Module may receive the same queue's typed client through `service.load()`. A + successful enqueue response means the message is durably committed. +- **FR3. Durable source of truth.** Prisma Postgres stores messages, availability, + attempt count, current lease, completion state, and enough failure information + to explain retry decisions. +- **FR4. Static internal consumer.** Each queue is wired at deployment to exactly + one logical Composer consumer service. One service may consume several queues. + Only an explicit consumer edge registers a service; exposing an unwired + compatible consumer port does not. Graph loading fails before execution or + deployment when a queue has no consumer, its consumer does not expose the + required contract, or topology wiring assigns a second distinct logical + consumer. A duplicate-consumer error identifies the queue and both conflicting + consumer services. Additional producer services remain allowed. +- **FR5. Single-consumer delivery.** Every attempt for a queue is sent to its one + logical consumer service. A successful attempt completes the message. Runtime + replicas behind that service endpoint are not separate consumer registrations. +- **FR6. Flexible consumer placement.** A consumer may be a dedicated Compute + service or an existing application Compute service. Documentation and examples + use a dedicated service by default. +- **FR7. HTTP push.** The dispatcher invokes the queue's consumer through a + bounded, authenticated HTTP request containing a batch of typed messages. The + response carries one acknowledgement or retry decision for every message. +- **FR8. Safe claiming.** Concurrent dispatcher work cannot claim the same + available message at the same time. Consumer application code runs outside the + claim transaction. +- **FR9. At-least-once recovery.** A failed request, retry result, lost response, + dispatcher failure, or expired lease can cause another attempt. The system does + not claim exactly-once execution. +- **FR10. Retry scheduling.** Retry policy determines when a failed message + becomes available and when automatic attempts stop. Retry state survives + Compute restarts. +- **FR11. Static policy, durable operations.** Queue and consumer configuration + are declared in TypeScript. Runtime message and attempt state is stored in + Postgres. +- **FR12. Internal authentication.** Composer provisions authentication for every + dispatcher-to-consumer edge. Application handlers do not receive or manage + service credentials. +- **FR13. Local operation.** The Queue Module and wired consumer services run under + `prisma-composer dev` with local durable storage and the same public producer and + consumer contracts used after deployment. +- **FR14. Failure visibility.** Runtime logs and stored attempt data identify the + queue, message, attempt number, consumer, outcome, and next retry time + without logging the message body by default. +- **FR15. Traffic-independent dispatch.** The queue capability includes one + dispatcher driver that remains running and processes new messages, delayed + retries, and expired leases without application traffic. +- **FR16. Queue-wide concurrency.** Each queue permits 10 active batch HTTP + requests by default and accepts a configured limit from 1 through 100. The + dispatcher does not claim more work while all slots are in use. Every slot + targets the queue's one logical consumer; version one does not configure + per-replica concurrency. +- **FR17. Consumer enqueue.** A consumer may receive the typed producer binding + for any queue, including the same queue it consumes, and enqueue follow-up work + without creating a dependency cycle. +- **FR18. Retry exhaustion.** A message that reaches its maximum attempts becomes + a terminal failed record in its original queue and is excluded from dispatch. +- **FR19. Failed-message replay.** A separately wired operational client can + replay one retained failed message. Replay is idempotent, preserves the failed + record, and returns a new message identifier linked to the original. +- **FR20. Unordered delivery.** Version one does not guarantee claim, delivery, + or completion order. Consumers cannot depend on first-in-first-out or per-key + ordering. +- **FR21. Bounded payload.** A queue accepts a JSON-compatible payload only when + it passes the queue schema and its serialized UTF-8 body is at most 128 KiB. + Rejection creates no message row. +- **FR22. Bounded batch.** A queue delivers 10 messages per batch by default and + accepts a configured limit from 1 through 100. The encoded request remains at + or below 1 MiB, may contain fewer available messages, and never splits a + message. +- **FR23. Bounded consumer request.** A consumer request times out after 30 + seconds by default and accepts a configured timeout from 1 through 50 seconds. + Its non-configurable lease lasts 30 seconds longer than that timeout, and a + timeout does not release the lease early. +- **FR24. Bounded attempts.** A queue allows five attempts by default and accepts + a configured maximum from 1 through 100. The initial delivery is attempt one, + and every durable claim consumes one attempt. +- **FR25. Retry policy builders.** A queue retry delay is a serializable + descriptor created by `fixedBackoff` or `exponentialBackoff`. + `exponentialBackoff` supports no, full, and equal jitter and defaults to equal + jitter. Delays range from one second through 24 hours, and exponential factors + range from 1.1 through 10. Invalid descriptors and arbitrary delay callbacks + fail graph loading. +- **FR26. Per-message retry delay.** A consumer can request a concrete delay from + one second through 24 hours for one message's next attempt. The override + changes neither its attempt count nor the queue policy and is applied only + from a valid consumer response. +- **FR27. Message retention.** Active message age is measured from enqueue and + does not reset on retry. Active retention defaults to seven days and ranges + from one hour through 30 days; expiry creates a replayable terminal failure. + Completed metadata defaults to 24 hours and ranges from immediate deletion + through seven days. Failed records default to 30 days and range from one + through 90 days. Completed payloads are removed immediately, and each terminal + record is deleted with its attempt history after retention. +- **FR28. Failed-message administration.** A separately wired operational client + can list retained failures without payloads, retrieve one failure with its + typed payload and attempt history, select replay candidates, and replay one + message at a time. Listing uses stable cursor pagination and filters for + failure time, reason, replay status, and last consumer. +- **FR29. Queue runtime control.** The operational client can pause and resume a + queue without changing its static topology or redeploying the application. + Pause prevents later claims but accepts enqueue and lets already claimed work + finish. Control state survives process restart, and both operations are + idempotent. +- **FR30. Operational status and activity.** The Queue Compute service exposes + current status and recent activity through a separately wired, typed + RPC-over-HTTP client using Composer service-key authentication. Normal traffic + is summarized by batch, while terminal failures, expiry, replay, pause, and + resume are individual events. Activity retention defaults to seven days and + ranges from one through 30 days. Durable activity and the queue mutation it + describes commit together in the Queue Module database. Status reports queue + state, message-state counts, concurrency use, oldest available message age, + and observation time. Activity uses a stable discriminated event contract, + cursor pagination, and filters for time, type, consumer, and message. + +## Non-Functional Requirements + +- **NFR1. Composer architecture.** The feature uses existing Module, service, + resource, dependency, and contract primitives. Composer core remains unaware of + Prisma Cloud queues. +- **NFR2. Deterministic topology.** Each queue's single consumer registration and + delivery edge are visible in the statically loaded graph. Version one performs + no runtime service discovery. +- **NFR3. Runtime neutrality.** Public authoring and handler types do not expose + Bun-only or Node-only types. Application code owns its build, and Composer does + not bundle or transform it. +- **NFR4. Durable acknowledgement boundary.** The system never reports enqueue + success before the message commit succeeds and never reports completion before + the acknowledgement state is committed. +- **NFR5. Failure isolation.** A consumer or dispatcher crash cannot leave a + message permanently unavailable. An expired lease eventually returns unfinished + work to the available state. +- **NFR6. Idempotent consumption.** Documentation and APIs state that consumers + must tolerate duplicate delivery. Message identifiers remain stable across + attempts so application code can implement durable deduplication where needed. +- **NFR7. Bounded work.** One delivery attempt has an HTTP deadline no longer + than 50 seconds and a lease exactly 30 seconds longer than that deadline. + Version one does not allow a handler to extend its lease. +- **NFR8. Durable clock.** Dispatcher availability is a Prisma Cloud target + property. Request-scoped `waitUntil`, self-invocation, and producer traffic are + not accepted as the clock that makes due work progress. +- **NFR9. No outbound product telemetry.** The queue service and dispatcher make + no outbound request for product analytics, OpenTelemetry export, or internal + error reporting. Operational analytics remain in `queue_activity` inside the + user's Queue Module database and are read only through the wired operational + port. + +## Non-goals + +- More than one logical consumer or subscription per queue, including publish and + subscribe delivery of one message to every subscription. A future topic and + subscription capability must not change the queue's single-consumer contract. +- Direct enqueue from services outside the Composer graph in version one. Until + external producer access is added, an external system must call an in-graph + application service that holds the typed producer binding. +- Statically declared external HTTPS consumers in version one; these are the next + planned phase. +- Dynamic consumer or producer registration, dynamic producer credential + creation, discovery, or heartbeats. +- Exactly-once execution. +- Consumer-controlled renewable leases or jobs that outlive one HTTP request. +- Automatic dead-letter routing and separate dead-letter queue consumers. +- Persisting consumer code or application deployment configuration in Postgres. +- Runtime editing of queue policy. +- Persisting each individual acknowledgement through a callback while a consumer + request is still running. +- Bulk replay of failed messages. +- Prisma Console or control-plane integration and a Queue Module-owned graphical + interface. +- Queue Module-owned outbound product analytics, Prometheus or OpenTelemetry + export, and durable metric time-series storage. Operational activity remains in + `queue_activity` inside the user's Queue Module database. +- Administrative operations beyond failed-message inspection, replay, pause and + resume, and queue observability. +- Strict first-in-first-out or per-key message ordering. + +# Acceptance Criteria + +- [ ] **AC1.** A producer sends a schema-valid message, receives success only + after commit, and the message remains available after all Queue Module Compute + processes restart. Covers FR1–FR3 and NFR4. +- [ ] **AC2.** A producer attempts to send a schema-invalid message and then a + schema-valid payload larger than 128 KiB. Each call returns the relevant + typed validation failure without creating a message row. Covers FR1, FR2, + and FR21. +- [ ] **AC3.** Two different producer services enqueue messages to the same queue. + Its one wired consumer service receives every message, and each message is + completed by one successful attempt. Covers FR2, FR4–FR7, and NFR2. +- [ ] **AC4.** Two dispatcher workers race for the same available messages and no + message has two active leases. Covers FR8. +- [ ] **AC5.** A consumer finishes its side effect but its HTTP response is lost. + The message is delivered again with the same message identifier and a higher + attempt number. Covers FR9, NFR5, and NFR6. +- [ ] **AC6.** A dispatcher stops after committing a claim. After lease expiry, + another dispatcher claims and delivers the message. Covers FR8, FR9, and NFR5. +- [ ] **AC7.** A retryable consumer result makes the message available at the time + calculated from queue policy, and restart does not change that time. Covers + FR10 and FR11. +- [ ] **AC8.** An unwired Compute service cannot invoke a queue consumer endpoint, + while the wired dispatcher can. Covers FR12. +- [ ] **AC9.** The same example application, including multiple producers wired to + one consumer, passes its conformance scenario under local development and a + deployed stage. Covers FR13 and NFR1–NFR3. +- [ ] **AC10.** Delivery logs and attempt records explain a failed attempt and its + next action without exposing the message body. Covers FR14. +- [ ] **AC11.** A consumer acknowledges some messages, requests retries for + others, and returns one valid response. The dispatcher commits each result + accordingly. If the handler throws an ordinary error before deciding every + message, recorded decisions are preserved and undecided messages are + retried. Covers FR7, FR9, and NFR4. +- [ ] **AC12.** With no producer or web traffic, a delayed retry becomes due and + the dispatcher delivers it. Restarting the dispatcher does not prevent due + work or expired leases from progressing. Covers FR15, NFR5, and NFR8. +- [ ] **AC13.** With a queue concurrency limit of ten, no more than ten batch HTTP + requests to its one logical consumer are active at once. Covers FR5 and + FR16. +- [ ] **AC14.** A consumer is wired to both a queue delivery contract and the same + queue's producer binding. The application graph loads without a dependency + cycle, and handling one message can durably enqueue follow-up work. Covers + FR4 and FR17. +- [ ] **AC15.** A message reaches its maximum attempts and becomes terminal + without further delivery. Replaying it twice with one logical replay request + returns the same new message, whose attempts start at zero and whose + `replayedFrom` points to the unchanged failed record. Covers FR18 and FR19. +- [ ] **AC16.** A delayed or retried message is overtaken by a later available + message, and both complete without the consumer relying on their order. + Covers FR20. +- [ ] **AC17.** With a configured batch size of 100 and enough available work, + delivery stops at 100 messages or before adding the first message that + would make the encoded request exceed 1 MiB. A smaller available set is + dispatched without waiting for more messages. Covers FR22 and NFR7. +- [ ] **AC18.** A consumer exceeds its configured timeout. The dispatcher aborts + the request, does not apply a result, and no dispatcher can claim the + messages before the original lease expires. Covers FR9, FR23, NFR5, and + NFR7. +- [ ] **AC19.** A queue configured with five maximum attempts repeatedly fails. + Each durable claim increments the attempt count, including a claim followed + by dispatcher failure before HTTP delivery. After attempt five fails, the + message becomes terminal and is not claimed again. Covers FR3, FR18, and + FR24. +- [ ] **AC20.** Fixed and exponential retry descriptors survive graph loading and + deployment without carrying executable user code. A failed attempt stores + one exact next availability time, and dispatcher restart does not change it. + Values outside the supported delay and factor ranges fail graph loading. + Covers FR10, FR11, and FR25. +- [ ] **AC21.** In one valid batch response, one message requests a 30-second + retry and another uses the queue policy. Their exact next availability + times reflect those separate decisions, while both attempt counts remain + unchanged by scheduling. Covers FR7, FR10, and FR26. +- [ ] **AC22.** An available or delayed message reaches its active retention + limit without exhausting its attempts. It becomes failed with reason + `expired`, remains inspectable and replayable during failed retention, and + is deleted with its attempt history when that period ends. Covers FR18, + FR19, and FR27. +- [ ] **AC23.** A consumer acknowledges a message after its active retention + deadline passes during the request. The message completes, its payload is + removed immediately, and its metadata and attempt history are deleted + after completed retention. A retry under the same timing instead produces + an expired failure retaining its payload. Covers FR7, FR18, and FR27. +- [ ] **AC24.** A pause races with a claim. A claim committed first may finish, + while no claim commits after the pause. Producers continue enqueueing + without consuming attempts. After process restart the queue remains paused; + resume wakes dispatch and both repeated controls return success. Covers + FR15, FR24, and FR29. +- [ ] **AC25.** Failed-message listing returns the documented metadata without + payloads. Looking up one listed identifier returns its typed payload, + attempts, consumers, safe error summaries, and replay links. An unwired + service can access neither operation. Covers FR12 and FR28. +- [ ] **AC26.** Failed-message listing returns 50 newest records by default and + never more than 100. Following its opaque cursor across equal failure times + produces each matching record once in stable order. Every documented + filter narrows the results without inspecting payloads. Covers FR28. +- [ ] **AC27.** Enqueue, delivery, retry, failure, replay, pause, and resume write + their documented operational activity only to `queue_activity`. Running + these operations makes no outbound product analytics, telemetry, or internal + error-reporting request. Covers FR30 and NFR9. +- [ ] **AC28.** Enqueue and delivery of one batch creates batch-level operational + activity with message and outcome counts rather than one success event per + message. A terminal failure, replay, pause, and resume each create an + individual event. Covers FR28–FR30. +- [ ] **AC29.** Operational activity survives Queue Compute restart, excludes + events after the configured retention deadline even before physical + cleanup, and never returns a message payload. Covers FR30. +- [ ] **AC30.** Status for a paused queue with available, delayed, leased, and + failed messages returns the authored queue name, pause time, exact + non-negative state counts, active and maximum concurrency, oldest available + age, last activity time, and `asOf`. Empty optional values are absent. + Covers FR29 and FR30. +- [ ] **AC31.** A transaction that commits a queue mutation also commits its + operational activity event; rolling back the mutation leaves neither visible. + Activity filtering uses indexed queue, time, type, consumer, batch, and + message fields without inspecting payloads. Covers FR30 and NFR4. +- [ ] **AC32.** Activity listing returns the 50 newest matching events by default + and never more than 100. Following its opaque cursor produces every event + once in stable order. Every event validates against its documented + discriminated shape and every supported filter narrows results. Covers + FR30. +- [ ] **AC33.** One consumer service is wired to a queue while a second service + exposes the same consumer contract without being wired. The graph loads and + messages go only to the wired consumer. Wiring the second service to the same + queue then makes graph loading fail before execution or deployment, and the + error identifies the queue and both conflicting consumer services. Adding + another producer does not fail graph loading. Covers FR2, FR4, FR5, and NFR2. + +# References + +- [Guiding principles](../../docs/design/01-principles/guiding-principles.md) +- [Architectural principles](../../docs/design/01-principles/architectural-principles.md) +- [ADR-0016: A Module has the same boundary as a service](../../docs/design/90-decisions/ADR-0016-a-module-has-the-same-boundary-as-a-service.md) +- [ADR-0020: Scheduled work is a driver, not a resource](../../docs/design/90-decisions/ADR-0020-scheduled-work-is-a-driver-not-a-resource.md) +- [ADR-0030: RPC callers use an auto-provisioned service key](../../docs/design/90-decisions/ADR-0030-rpc-callers-verified-with-an-auto-provisioned-service-key.md) +- [ADR-0037: Service RPC calls carry an idempotency key](../../docs/design/90-decisions/ADR-0037-service-rpc-calls-carry-an-idempotency-key.md) +- [Connection contracts](../../docs/design/10-domains/connection-contracts.md) +- [Google Cloud Tasks task model](https://cloud.google.com/tasks/docs/reference/rest/v2/projects.locations.queues.tasks) +- [Cloudflare Queues delivery guarantees](https://developers.cloudflare.com/queues/reference/delivery-guarantees/) +- [PostgreSQL `SELECT`](https://www.postgresql.org/docs/current/sql-select.html) + +# Deployment validation + +The walking skeleton is deployed as `queues-demo`. Its separate dispatcher +service remained active without direct requests, claimed five messages from the +queue service, and pushed all five to the application consumer. This validates +the target's always-running Compute pattern for the prototype. A deployed +fail-once message also moved from failed attempt one to successful attempt two +under the persisted five-second fixed retry policy. Restart recovery and the +remaining retry algorithms stay part of the full version-one work. diff --git a/.drive/projects/queue-module/trace.jsonl b/.drive/projects/queue-module/trace.jsonl new file mode 100644 index 00000000..84d1f7df --- /dev/null +++ b/.drive/projects/queue-module/trace.jsonl @@ -0,0 +1,4 @@ +{"event_id":"3624df39-5574-4c3c-b3e3-1a1f7e388b78","schema_version":"1","ts":"2026-07-31T15:50:52.276Z","project_run_id":"queue-module","orchestrator_agent_id":null,"event_type":"spec-amended","spec_path":".drive/projects/queue-module/spec.md","spec_kind":"project","byte_length":41966,"bytes_delta":378,"edge_cases_count":null,"open_questions_count":0,"dod_items_count":0,"reason":"scope-shift","sections_changed":["Summary","Prototype status","Context","At a glance","Approach","Functional Requirements","Non-Functional Requirements","Non-goals","Acceptance Criteria"]} +{"event_id":"87eae25d-6ad6-48ea-aac2-e315a9390bcb","schema_version":"1","ts":"2026-07-31T16:39:12.216Z","project_run_id":"queue-module","orchestrator_agent_id":null,"event_type":"spec-amended","spec_path":".drive/projects/queue-module/spec.md","spec_kind":"project","byte_length":43142,"bytes_delta":1176,"edge_cases_count":null,"open_questions_count":0,"dod_items_count":0,"reason":"operator-correction","sections_changed":["Approach","Functional Requirements","Acceptance Criteria"]} +{"event_id":"6e68fa34-5834-4857-8496-7c1ad1f805fb","schema_version":"1","ts":"2026-07-31T16:44:37.730Z","project_run_id":"queue-module","orchestrator_agent_id":null,"event_type":"spec-amended","spec_path":".drive/projects/queue-module/spec.md","spec_kind":"project","byte_length":44014,"bytes_delta":872,"edge_cases_count":null,"open_questions_count":0,"dod_items_count":0,"reason":"scope-shift","sections_changed":["Summary","Context","Approach","Non-goals"]} +{"event_id":"f020ea47-a748-419c-825d-7cacb0bcefed","schema_version":"1","ts":"2026-08-01T09:30:19.285Z","project_run_id":"queue-module","orchestrator_agent_id":null,"event_type":"spec-amended","spec_path":".drive/projects/queue-module/spec.md","spec_kind":"project","byte_length":42899,"bytes_delta":-1115,"edge_cases_count":null,"open_questions_count":0,"dod_items_count":0,"reason":"operator-correction","sections_changed":["Prototype status","Approach","Functional Requirements","Non-Functional Requirements","Non-goals","Acceptance Criteria"]} diff --git a/README.md b/README.md index 21277e23..ece5cede 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,8 @@ npx skills add prisma/composer --skill prisma-composer That's the whole setup. Your agent now knows the entire API and arrives prepped with the **building blocks** it can snap together — ready-made Modules for scheduled jobs, blob storage, and event streams, alongside the -ones you write. Ask it for what you want ("a Next.js storefront calling an +ones you write. Persistent work queues are available as an early prototype. +Ask it for what you want ("a Next.js storefront calling an orders API with its own Postgres, deployed to a staging stage") and let it compose. @@ -131,7 +132,7 @@ have. | Guide | Covers | | --- | --- | | [Getting started](docs/guides/getting-started.md) | Your first app end to end; porting an existing Node or Next.js app | -| [Building an app](docs/guides/building-an-app.md) | Contracts, databases (plain + Prisma Next-typed with migrations), reusable Modules, cron/storage/streams, config, secrets | +| [Building an app](docs/guides/building-an-app.md) | Contracts, databases (plain + Prisma Next-typed with migrations), reusable Modules, cron/storage/streams/queues, config, secrets | | [Testing](docs/guides/testing.md) | Unit tests with `mockService`, integration tests with `bootstrapService` | | [Deploying and operating](docs/guides/deploying.md) | Stages, destroy, CI, how apps behave in production | @@ -147,14 +148,16 @@ Complete, deployable apps under [`examples/`](examples/): | [cron](examples/cron/) | Scheduled jobs: `defineSchedule` + `serveSchedule` + the cron module | | [storage](examples/storage/) | The S3-backed blob store module | | [streams](examples/streams/) | Durable append-only event streams over storage | +| [queues](examples/queues/) | Persistent queue delivery from and back to one Compute service | ## Building blocks and extensions A **Module** is the unit of reuse: it owns its internals and exposes a typed port, so composing one is a couple of lines and never an integration you have -to invent. Three ship inside `@prisma/composer-prisma-cloud` today — `cron` -(scheduled jobs), `storage` (S3-backed blobs), and `streams` (durable event -streams) — alongside the Modules you write yourself. +to invent. Four ship inside `@prisma/composer-prisma-cloud` today — `cron` +(scheduled jobs), `storage` (S3-backed blobs), `streams` (durable event +streams), and the early `queues` prototype (Postgres-backed work delivery) — +alongside the Modules you write yourself. An **extension** is a package that brings its own Modules, resources, or deploy target. The convention is an npm package named `prisma-composer-*` — diff --git a/architecture.config.json b/architecture.config.json index 68eaf521..8ab760bf 100644 --- a/architecture.config.json +++ b/architecture.config.json @@ -672,6 +672,54 @@ "layer": "modules", "plane": "execution" }, + { + "glob": "packages/1-prisma-cloud/2-shared-modules/queues/src/*.ts", + "domain": "prisma-cloud", + "layer": "modules", + "plane": "shared" + }, + { + "glob": "packages/1-prisma-cloud/2-shared-modules/queues/src/execution/**", + "domain": "prisma-cloud", + "layer": "modules", + "plane": "execution" + }, + { + "glob": "packages/1-prisma-cloud/2-shared-modules/queues/src/exports/index.ts", + "domain": "prisma-cloud", + "layer": "modules", + "plane": "shared" + }, + { + "glob": "packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-service.ts", + "domain": "prisma-cloud", + "layer": "modules", + "plane": "shared" + }, + { + "glob": "packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-service.ts", + "domain": "prisma-cloud", + "layer": "modules", + "plane": "shared" + }, + { + "glob": "packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-entrypoint.ts", + "domain": "prisma-cloud", + "layer": "modules", + "plane": "execution" + }, + { + "glob": "packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-entrypoint.ts", + "domain": "prisma-cloud", + "layer": "modules", + "plane": "execution" + }, + { + "glob": "packages/9-public/composer-prisma-cloud/src/exports/queues.ts", + "domain": "public", + "layer": "public", + "plane": "shared" + }, { "glob": "packages/9-public/composer/src/exports/deploy.ts", "domain": "public", diff --git a/docs/guides/building-an-app.md b/docs/guides/building-an-app.md index 4354539c..401caf7b 100644 --- a/docs/guides/building-an-app.md +++ b/docs/guides/building-an-app.md @@ -3,7 +3,7 @@ This guide covers everything you reach for once [Getting started](getting-started.md) has shown you the shape: giving a service a database (plain or Prisma Next-typed), packaging pieces as reusable -Modules, the cron/storage/streams modules that ship with the framework, and +Modules, the cron/storage/streams/queues modules that ship with the framework, and the service input — configuration and secrets as one schema. ## How the pieces fit @@ -315,6 +315,7 @@ is a couple of lines: | `cron` from `@prisma/composer-prisma-cloud/cron` | A scheduler that fires your jobs at your service on an interval | nothing | | `storage` from `@prisma/composer-prisma-cloud/storage` | An S3-backed blob store, credentials included | `store` | | `streams` from `@prisma/composer-prisma-cloud/streams` | Durable append-only event streams, backed by a `store` | `streams` | +| `queues` from `@prisma/composer-prisma-cloud/queues` | Persistent work delivery backed by Postgres | `producer`, `dispatch` | Cron is the one most apps want first. You supply two things — a schedule and a runner service that exposes the `trigger` contract — and the module does @@ -357,6 +358,57 @@ provision(cron({ schedule, runner: promotionsService }), { [`examples/streams`](../../examples/streams/) show the other two, including the streams module's secret binding. +The queues prototype uses one static catalogue for producer and consumer +typing. The queue Module owns its Postgres database and queue service. A +separate always-running dispatcher claims messages and pushes them to the +consumer service: + +```ts +// queues.ts +import { defineQueues, fixedBackoff } from '@prisma/composer-prisma-cloud/queues'; +import { type } from 'arktype'; + +export const appQueues = defineQueues({ + thumbnails: { + message: type({ imageId: 'string' }), + retry: { + maxAttempts: 5, + delay: fixedBackoff({ delay: '5s' }), + }, + }, +}); +``` + +```ts +// worker/service.ts +import { queueConsumer, queueProducer } from '@prisma/composer-prisma-cloud/queues'; + +export default compute({ + name: 'worker', + deps: { queues: queueProducer(appQueues) }, + expose: { consumer: queueConsumer() }, + build: node({ module: import.meta.url, entry: '../../dist/worker/server.mjs' }), +}); +``` + +```ts +// module.ts +const queue = provision(queues({ definitions: appQueues })); +const worker = provision(workerService, { deps: { queues: queue.producer } }); + +provision(queueDispatcher(), { + deps: { queue: queue.dispatch, consumer: worker.consumer }, + input: { pollIntervalMs: 250, leaseSeconds: 30 }, +}); +``` + +Route `/rpc/*` in the worker to `serveQueues(workerService, appQueues, +handlers)`. A handler can enqueue follow-up work through the producer dependency, +including back to the same queue. This first prototype delivers one message per +request and supports per-queue fixed retry delays. Batches, exponential retry, +multiple competing consumers, replay, pause, and operational APIs remain future +work. See [`examples/queues`](../../examples/queues/) for the complete deployed app. + ### Where new blocks come from An **extension** is a package that brings its own Modules, resources, or diff --git a/examples/queues/README.md b/examples/queues/README.md new file mode 100644 index 00000000..95b49720 --- /dev/null +++ b/examples/queues/README.md @@ -0,0 +1,29 @@ +# Queues demo + +This deployed example proves one application Compute service can both produce +and consume persistent queue messages. Its browser sends messages through the +typed producer dependency. The separate dispatcher claims them from the queue +module's Postgres database and delivers them back to the same application +service's consumer port. + +Each queue definition owns its retry policy. This demo uses a fixed five-second +delay and stops after five delivery attempts: + +```ts +messages: { + message: type({ text: 'string' }), + retry: { + maxAttempts: 5, + delay: fixedBackoff({ delay: '5s' }), + }, +} +``` + +Enable **Fail the first delivery** before enqueuing to exercise the real retry +path. The consumer throws on attempt one. The dispatcher stores the next +availability time in Postgres and delivers the same message again on attempt +two. + +The consumed-message panel is a bounded in-memory demo feed. Queue durability +comes from Postgres; refreshing or restarting the application can clear the +display without losing queued messages. diff --git a/examples/queues/module.ts b/examples/queues/module.ts new file mode 100644 index 00000000..606a2f39 --- /dev/null +++ b/examples/queues/module.ts @@ -0,0 +1,13 @@ +import { module } from '@prisma/composer'; +import { queueDispatcher, queues } from '@prisma/composer-prisma-cloud/queues'; +import appService from './src/app/service.ts'; +import { demoQueues } from './src/queues.ts'; + +export default module('queues-example', ({ provision }) => { + const queue = provision(queues({ definitions: demoQueues })); + const app = provision(appService, { deps: { queues: queue.producer } }); + provision(queueDispatcher(), { + deps: { queue: queue.dispatch, consumer: app.consumer }, + input: { pollIntervalMs: 250, leaseSeconds: 30 }, + }); +}); diff --git a/examples/queues/package.json b/examples/queues/package.json new file mode 100644 index 00000000..833511b6 --- /dev/null +++ b/examples/queues/package.json @@ -0,0 +1,23 @@ +{ + "name": "@prisma/example-queues", + "version": "0.3.0", + "private": true, + "type": "module", + "scripts": { + "build": "tsdown", + "typecheck": "tsc --noEmit", + "deploy": "pnpm turbo run build --filter @prisma/example-queues... && ( set -a; . \"${PRISMA_DEPLOY_ENV:-../../.env}\"; set +a; bun node_modules/.bin/prisma-composer deploy module.ts ${QUEUES_STACK_NAME:+--name \"$QUEUES_STACK_NAME\"} )", + "destroy": "( set -a; . \"${PRISMA_DEPLOY_ENV:-../../.env}\"; set +a; bun node_modules/.bin/prisma-composer destroy module.ts --production ${QUEUES_STACK_NAME:+--name \"$QUEUES_STACK_NAME\"} )" + }, + "dependencies": { + "@prisma/composer": "workspace:0.3.0", + "@prisma/composer-prisma-cloud": "workspace:0.3.0", + "arktype": "^2.2.3" + }, + "devDependencies": { + "@types/bun": "^1.3.13", + "tsdown": "^0.22.7", + "unrun": "^0.3.1", + "typescript": "^6.0.3" + } +} diff --git a/examples/queues/prisma-composer.config.ts b/examples/queues/prisma-composer.config.ts new file mode 100644 index 00000000..2c085a38 --- /dev/null +++ b/examples/queues/prisma-composer.config.ts @@ -0,0 +1,8 @@ +import { defineConfig } from '@prisma/composer/config'; +import { nodeBuild } from '@prisma/composer/node/control'; +import { prismaCloud, prismaState } from '@prisma/composer-prisma-cloud/control'; + +export default defineConfig({ + extensions: [prismaCloud(), nodeBuild()], + state: prismaState(), +}); diff --git a/examples/queues/src/app/page.ts b/examples/queues/src/app/page.ts new file mode 100644 index 00000000..507e6b0b --- /dev/null +++ b/examples/queues/src/app/page.ts @@ -0,0 +1,519 @@ +import { DEMO_MAX_ATTEMPTS, DEMO_RETRY_DELAY } from '../queues.ts'; + +export const page = ` + + + + + Composer Queues + + + +
+

Composer Queues

+
queue-demo / messages
+
+
+
+
+

One application worker, two queue roles

+

The producer commits durable work. The independent dispatcher later pushes it back to the consumer port on the same Compute worker.

+
+
Ready for a message
+
+
+
+ Application Compute + queueDemo + Hosts the demo UI and application code. +
+ producer client + consumer port +
+
+ +
+ Queue Module + queues.service + Stores messages and coordinates leases. +
+ +
+ Prisma Postgres + queues.db + Durable messages and delivery state. +
+
+ Always-running Compute + queueDispatcher + Claims work and pushes it over HTTP. +
+ + +
queueDispatcher claims from queues.service, then pushes HTTP back to the queueDemo consumer port.
+
+
+
+
+
+

Produce

+ Persistent enqueue +
+
+
+ + +
+
+ + +
+
Retry policy: fixed ${DEMO_RETRY_DELAY} · ${DEMO_MAX_ATTEMPTS} attempts maximum
+
+ + +
+
+ + +
+
+
+
+
+

Consume

+ Waiting for messages +
+
Enqueue a message to see it consumed by this Compute app.
+
+
+ + +`; diff --git a/examples/queues/src/app/server.ts b/examples/queues/src/app/server.ts new file mode 100644 index 00000000..0bffc61f --- /dev/null +++ b/examples/queues/src/app/server.ts @@ -0,0 +1,111 @@ +import { type QueueConsumerMessage, serveQueues } from '@prisma/composer-prisma-cloud/queues'; +import { demoQueues } from '../queues.ts'; +import { page } from './page.ts'; +import service from './service.ts'; + +interface DeliveryEvent extends QueueConsumerMessage<{ text: string }> { + readonly status: 'failed' | 'consumed'; + readonly eventAt: string; +} + +const events: DeliveryEvent[] = []; +const MAX_FEED_MESSAGES = 200; +let failAttemptOne = false; +let failingMessageId: string | undefined; +let failureRecorded = false; + +function record(event: DeliveryEvent): void { + events.unshift(event); + if (events.length > MAX_FEED_MESSAGES) events.length = MAX_FEED_MESSAGES; +} + +const rpcHandler = serveQueues(service, demoQueues, { + messages: async (message) => { + if ( + failAttemptOne && + message.attempt === 1 && + (failingMessageId === undefined || failingMessageId === message.id) + ) { + failingMessageId = message.id; + if (!failureRecorded) { + failureRecorded = true; + record({ ...message, status: 'failed', eventAt: new Date().toISOString() }); + } + throw new Error('retry demo failed this consumer delivery'); + } + if (failingMessageId === message.id) { + failAttemptOne = false; + failingMessageId = undefined; + failureRecorded = false; + } + record({ ...message, status: 'consumed', eventAt: new Date().toISOString() }); + }, +}); + +function json(body: unknown, status = 200): Response { + return Response.json(body, { status, headers: { 'cache-control': 'no-store' } }); +} + +function enqueueInput( + value: unknown, +): { message: string; count: number; failFirstAttempt: boolean } | null { + if (typeof value !== 'object' || value === null) return null; + if (!('message' in value) || !('count' in value)) return null; + if ( + typeof value.message !== 'string' || + value.message.length < 1 || + value.message.length > 2000 + ) { + return null; + } + if (typeof value.count !== 'number' || !Number.isInteger(value.count)) return null; + if (value.count < 1 || value.count > 50) return null; + const failFirstAttempt = 'failFirstAttempt' in value ? value.failFirstAttempt : false; + if (failFirstAttempt !== undefined && typeof failFirstAttempt !== 'boolean') { + return null; + } + return { + message: value.message, + count: value.count, + failFirstAttempt: failFirstAttempt === true, + }; +} + +async function fetchHandler(request: Request): Promise { + const url = new URL(request.url); + if (url.pathname.startsWith('/rpc/')) return rpcHandler(request); + + if (request.method === 'GET' && url.pathname === '/') { + return new Response(page, { headers: { 'content-type': 'text/html; charset=utf-8' } }); + } + if (request.method === 'GET' && url.pathname === '/health') return json({ ok: true }); + if (request.method === 'GET' && url.pathname === '/api/consumed') { + return json({ messages: events }); + } + if (request.method === 'POST' && url.pathname === '/api/messages') { + let raw: unknown; + try { + raw = await request.json(); + } catch { + return json({ error: 'Request body must be JSON.' }, 400); + } + const input = enqueueInput(raw); + if (input === null) { + return json({ error: 'Message must be 1–2000 characters and copies must be 1–50.' }, 400); + } + + const { queues } = service.load(); + if (input.failFirstAttempt) { + failAttemptOne = true; + failingMessageId = undefined; + failureRecorded = false; + } + const ids = await Promise.all( + Array.from({ length: input.count }, () => queues.messages.send({ text: input.message })), + ); + return json({ count: ids.length, ids: ids.map(({ id }) => id) }, 202); + } + return new Response('Not found', { status: 404 }); +} + +Bun.serve({ port: service.port(), hostname: '0.0.0.0', fetch: fetchHandler }); diff --git a/examples/queues/src/app/service.ts b/examples/queues/src/app/service.ts new file mode 100644 index 00000000..340074c2 --- /dev/null +++ b/examples/queues/src/app/service.ts @@ -0,0 +1,11 @@ +import node from '@prisma/composer/node'; +import { compute } from '@prisma/composer-prisma-cloud'; +import { queueConsumer, queueProducer } from '@prisma/composer-prisma-cloud/queues'; +import { demoQueues } from '../queues.ts'; + +export default compute({ + name: 'queueDemo', + deps: { queues: queueProducer(demoQueues) }, + expose: { consumer: queueConsumer() }, + build: node({ module: import.meta.url, entry: '../../dist/app/server.mjs' }), +}); diff --git a/examples/queues/src/queues.ts b/examples/queues/src/queues.ts new file mode 100644 index 00000000..714846c4 --- /dev/null +++ b/examples/queues/src/queues.ts @@ -0,0 +1,15 @@ +import { defineQueues, fixedBackoff } from '@prisma/composer-prisma-cloud/queues'; +import { type } from 'arktype'; + +export const DEMO_RETRY_DELAY = '5s'; +export const DEMO_MAX_ATTEMPTS = 5; + +export const demoQueues = defineQueues({ + messages: { + message: type({ text: '1 <= string <= 2000' }), + retry: { + maxAttempts: DEMO_MAX_ATTEMPTS, + delay: fixedBackoff({ delay: DEMO_RETRY_DELAY }), + }, + }, +}); diff --git a/examples/queues/tsconfig.json b/examples/queues/tsconfig.json new file mode 100644 index 00000000..73a44bd1 --- /dev/null +++ b/examples/queues/tsconfig.json @@ -0,0 +1,7 @@ +{ + "extends": "../../tsconfig.base.json", + "compilerOptions": { + "types": ["bun"] + }, + "include": ["module.ts", "prisma-composer.config.ts", "src", "tests"] +} diff --git a/examples/queues/tsdown.config.ts b/examples/queues/tsdown.config.ts new file mode 100644 index 00000000..ebe2d4f3 --- /dev/null +++ b/examples/queues/tsdown.config.ts @@ -0,0 +1,14 @@ +import { defineConfig } from 'tsdown'; + +export default defineConfig({ + entry: { server: 'src/app/server.ts' }, + outDir: 'dist/app', + format: 'esm', + platform: 'node', + external: [/^bun$/, /^bun:/, /^node:/], + noExternal: [/.*/], + outputOptions: { inlineDynamicImports: true }, + dts: false, + sourcemap: false, + clean: true, +}); diff --git a/examples/queues/turbo.json b/examples/queues/turbo.json new file mode 100644 index 00000000..95960709 --- /dev/null +++ b/examples/queues/turbo.json @@ -0,0 +1,4 @@ +{ + "$schema": "https://turbo.build/schema.json", + "extends": ["//"] +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/package.json b/packages/1-prisma-cloud/2-shared-modules/queues/package.json new file mode 100644 index 00000000..b3c67ca2 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/package.json @@ -0,0 +1,52 @@ +{ + "name": "@internal/queues", + "version": "0.3.0", + "private": true, + "type": "module", + "description": "Persistent queues backed by Prisma Postgres and delivered by Prisma Compute.", + "exports": { + ".": { + "types": "./dist/index.d.mts", + "default": "./dist/index.mjs" + }, + "./queue-service": { + "types": "./dist/queue-service.d.mts", + "default": "./dist/queue-service.mjs" + }, + "./queue-entrypoint": { + "types": "./dist/queue-entrypoint.d.mts", + "default": "./dist/queue-entrypoint.mjs" + }, + "./dispatcher-service": { + "types": "./dist/dispatcher-service.d.mts", + "default": "./dist/dispatcher-service.mjs" + }, + "./dispatcher-entrypoint": { + "types": "./dist/dispatcher-entrypoint.d.mts", + "default": "./dist/dispatcher-entrypoint.mjs" + }, + "./package.json": "./package.json" + }, + "scripts": { + "typecheck": "tsc --noEmit", + "test": "bun test src", + "build": "tsdown", + "clean": "rm -rf dist" + }, + "dependencies": { + "@internal/core": "workspace:0.3.0", + "@internal/foundation": "workspace:0.3.0", + "@internal/node": "workspace:0.3.0", + "@internal/prisma-cloud": "workspace:0.3.0", + "@internal/service-rpc": "workspace:0.3.0", + "arktype": "^2.2.3" + }, + "devDependencies": { + "@internal/tsdown-config": "workspace:0.3.0", + "@types/bun": "^1.3.13", + "typescript": "^6.0.3", + "vitest": "^4.1.10", + "tsdown": "^0.22.7", + "unrun": "^0.3.1" + } +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/definitions.test.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/definitions.test.ts new file mode 100644 index 00000000..17b869a1 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/definitions.test.ts @@ -0,0 +1,35 @@ +import { describe, expect, test } from 'bun:test'; +import { type } from 'arktype'; +import { defineQueues, fixedBackoff } from '../exports/index.ts'; + +describe('queue retry policy', () => { + test('defines a fixed retry delay on one queue', () => { + const definitions = defineQueues({ + messages: { + message: type({ text: 'string' }), + retry: { maxAttempts: 5, delay: fixedBackoff({ delay: '5s' }) }, + }, + }); + + expect(definitions.messages.retry).toEqual({ + maxAttempts: 5, + delay: { kind: 'fixed', delaySeconds: 5 }, + }); + }); + + test('rejects retry delays outside one second through 24 hours', () => { + expect(() => fixedBackoff({ delay: '0s' })).toThrow('between 1 second and 24 hours'); + expect(() => fixedBackoff({ delay: '25h' })).toThrow('between 1 second and 24 hours'); + }); + + test('rejects an invalid maximum attempt count while defining queues', () => { + expect(() => + defineQueues({ + messages: { + message: type({ text: 'string' }), + retry: { maxAttempts: 0, delay: fixedBackoff({ delay: '5s' }) }, + }, + }), + ).toThrow('maxAttempts must be an integer from 1 through 100'); + }); +}); diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/dispatcher.test.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/dispatcher.test.ts new file mode 100644 index 00000000..bf1ab3c1 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/dispatcher.test.ts @@ -0,0 +1,64 @@ +import { describe, expect, test } from 'bun:test'; +import type { DispatcherClients } from '../execution/dispatcher.ts'; +import { dispatchOnce } from '../execution/dispatcher.ts'; + +function clients( + events: string[], + delivery?: { readonly error?: Error; readonly accepted?: boolean }, +): DispatcherClients { + return { + queue: { + claim: async () => ({ + message: { + id: 'message-1', + queue: 'messages', + body: { text: 'hello' }, + attempt: 1, + enqueuedAt: '2026-07-31T12:00:00.000Z', + }, + leaseToken: 'lease-1', + }), + complete: async () => { + events.push('complete'); + return { ok: true }; + }, + release: async () => { + events.push('release'); + return { outcome: 'retrying' }; + }, + }, + consumer: { + deliver: async () => { + events.push('deliver'); + if (delivery?.error !== undefined) throw delivery.error; + return { ok: delivery?.accepted ?? true }; + }, + }, + }; +} + +describe('dispatcher', () => { + test('completes a lease only after the consumer accepts the message', async () => { + const events: string[] = []; + expect(await dispatchOnce(clients(events), { leaseSeconds: 30 })).toBe(true); + expect(events).toEqual(['deliver', 'complete']); + }); + + test('releases the lease when delivery fails', async () => { + const events: string[] = []; + await expect( + dispatchOnce(clients(events, { error: new Error('consumer failed') }), { + leaseSeconds: 30, + }), + ).rejects.toThrow('consumer failed'); + expect(events).toEqual(['deliver', 'release']); + }); + + test('releases the lease when the consumer rejects the delivery', async () => { + const events: string[] = []; + expect(await dispatchOnce(clients(events, { accepted: false }), { leaseSeconds: 30 })).toBe( + true, + ); + expect(events).toEqual(['deliver', 'release']); + }); +}); diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/pg-harness.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/pg-harness.ts new file mode 100644 index 00000000..cbe35db2 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/pg-harness.ts @@ -0,0 +1,125 @@ +import { execFileSync, spawnSync } from 'node:child_process'; +import { randomUUID } from 'node:crypto'; +import * as fs from 'node:fs'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import { SQL } from 'bun'; + +export interface TestPostgres { + readonly url: string; + readonly stop: () => void; +} + +export interface TestDatabase { + readonly url: string; + readonly drop: () => Promise; +} + +const PG_ENV = { ...process.env, LC_ALL: 'C', LANG: 'C' }; + +const probe = (binary: string): boolean => + spawnSync(binary, ['--version'], { stdio: 'ignore', env: PG_ENV }).status === 0; + +function ubuntuCandidates(name: string): string[] { + try { + return fs + .readdirSync('/usr/lib/postgresql') + .map((version) => path.join('/usr/lib/postgresql', version, 'bin', name)); + } catch { + return []; + } +} + +function findBinary(name: string): string | undefined { + return [ + name, + `/opt/homebrew/opt/postgresql@15/bin/${name}`, + `/opt/homebrew/bin/${name}`, + `/usr/local/opt/postgresql@15/bin/${name}`, + `/usr/local/bin/${name}`, + ...ubuntuCandidates(name), + ].find(probe); +} + +function withDatabase(url: string, database: string): string { + const parsed = new URL(url); + parsed.pathname = `/${database}`; + return parsed.toString(); +} + +export async function createTestDatabase(baseUrl: string): Promise { + const name = `queues_test_${randomUUID().replaceAll('-', '')}`; + const admin = new SQL({ url: baseUrl, max: 1 }); + try { + await admin.unsafe(`create database "${name}"`); + } finally { + await admin.end(); + } + return { + url: withDatabase(baseUrl, name), + drop: async () => { + const sql = new SQL({ url: baseUrl, max: 1 }); + try { + await sql`select pg_terminate_backend(pid) from pg_stat_activity + where datname = ${name} and pid <> pg_backend_pid()`; + await sql.unsafe(`drop database if exists "${name}"`); + } finally { + await sql.end(); + } + }, + }; +} + +export function startTestPostgres(): TestPostgres | undefined { + const configured = process.env['STATE_TEST_DATABASE_URL']; + if (configured !== undefined) return { url: configured, stop: () => {} }; + + const initdb = findBinary('initdb'); + const pgCtl = findBinary('pg_ctl'); + if (initdb === undefined || pgCtl === undefined) { + if (process.env['CI'] !== undefined) { + throw new Error( + 'CI has no queue test Postgres: set STATE_TEST_DATABASE_URL or install initdb and pg_ctl.', + ); + } + return undefined; + } + + const baseDir = process.env['QUEUES_TEST_PG_TMPDIR'] ?? os.tmpdir(); + fs.mkdirSync(baseDir, { recursive: true }); + const dataDir = fs.mkdtempSync(path.join(baseDir, 'prisma-composer-queues-pg-')); + const logFile = path.join(dataDir, 'server.log'); + execFileSync( + initdb, + ['-D', dataDir, '-U', 'postgres', '--auth=trust', '-E', 'UTF8', '--locale=C'], + { stdio: 'pipe', env: PG_ENV }, + ); + + let lastError = 'unknown error'; + for (let attempt = 0; attempt < 5; attempt++) { + const port = 20000 + Math.floor(Math.random() * 20000); + const result = spawnSync( + pgCtl, + ['-D', dataDir, '-o', `-p ${port} -h 127.0.0.1`, '-w', '-l', logFile, 'start'], + { stdio: 'pipe', env: PG_ENV }, + ); + if (result.status === 0) { + return { + url: `postgres://postgres@127.0.0.1:${port}/postgres`, + stop: () => { + try { + execFileSync(pgCtl, ['-D', dataDir, '-m', 'fast', 'stop'], { + stdio: 'pipe', + env: PG_ENV, + }); + } finally { + fs.rmSync(dataDir, { recursive: true, force: true }); + } + }, + }; + } + lastError = result.stderr.toString(); + } + fs.rmSync(dataDir, { recursive: true, force: true }); + throw new Error(`queue test Postgres failed to start: ${lastError}`); +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/pg-queue-store.integration.test.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/pg-queue-store.integration.test.ts new file mode 100644 index 00000000..c83c9df5 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/__tests__/pg-queue-store.integration.test.ts @@ -0,0 +1,102 @@ +import { afterAll, beforeAll, describe, expect, test } from 'bun:test'; +import { createPgQueueStore } from '../execution/pg-queue-store.ts'; +import type { QueueStore } from '../queue-store.ts'; +import { createTestDatabase, startTestPostgres, type TestDatabase } from './pg-harness.ts'; + +const pg = startTestPostgres(); +const suite = pg ? describe : describe.skip; + +suite('persistent queue', () => { + let database: TestDatabase; + let store: QueueStore; + + beforeAll(async () => { + if (pg === undefined) throw new Error('no Postgres available'); + database = await createTestDatabase(pg.url); + store = await createPgQueueStore(database.url, [ + { name: 'messages', maxAttempts: 5, retryDelayMs: 1_000 }, + { name: 'delayed', maxAttempts: 5, retryDelayMs: 1_000 }, + { name: 'single-attempt', maxAttempts: 1, retryDelayMs: 1_000 }, + ]); + }); + + afterAll(async () => { + await database?.drop(); + pg?.stop(); + }); + + test('enqueue is idempotent and a completed message is not claimed again', async () => { + const enqueueKey = crypto.randomUUID(); + const first = await store.enqueue({ + queue: 'messages', + body: { text: 'hello' }, + enqueueKey, + }); + const duplicate = await store.enqueue({ + queue: 'messages', + body: { text: 'ignored duplicate' }, + enqueueKey, + }); + expect(duplicate.id).toBe(first.id); + + const leased = await store.claim(30); + expect(leased?.message).toEqual({ + id: first.id, + queue: 'messages', + body: { text: 'hello' }, + attempt: 1, + enqueuedAt: expect.any(String), + }); + if (leased === null) throw new Error('expected a leased message'); + expect(await store.complete(leased.message.id, leased.leaseToken)).toBe(true); + expect(await store.claim(30)).toBeNull(); + }); + + test('concurrent claims lease a message to only one dispatcher', async () => { + const enqueued = await store.enqueue({ + queue: 'messages', + body: { text: 'one owner' }, + enqueueKey: crypto.randomUUID(), + }); + const claims = await Promise.all([store.claim(30), store.claim(30)]); + const leased = claims.filter((claim) => claim !== null); + expect(leased).toHaveLength(1); + expect(leased[0]?.message.id).toBe(enqueued.id); + const winner = leased[0]; + if (winner === undefined) throw new Error('expected one winning claim'); + expect(await store.complete(winner.message.id, winner.leaseToken)).toBe(true); + }); + + test('persists the retry delay before making a failed delivery available', async () => { + await store.enqueue({ + queue: 'delayed', + body: { text: 'wait before retrying' }, + enqueueKey: crypto.randomUUID(), + }); + const first = await store.claim(30); + if (first === null) throw new Error('expected the message to lease'); + + expect(await store.release(first.message.id, first.leaseToken)).toBe('retrying'); + expect(await store.claim(30)).toBeNull(); + + await Bun.sleep(1_100); + const retried = await store.claim(30); + expect(retried?.message.id).toBe(first.message.id); + expect(retried?.message.attempt).toBe(2); + if (retried === null) throw new Error('expected the delayed retry'); + expect(await store.complete(retried.message.id, retried.leaseToken)).toBe(true); + }); + + test('marks a message failed after its configured maximum attempts', async () => { + await store.enqueue({ + queue: 'single-attempt', + body: { text: 'stop retrying' }, + enqueueKey: crypto.randomUUID(), + }); + const leased = await store.claim(30); + if (leased === null) throw new Error('expected the message to lease'); + + expect(await store.release(leased.message.id, leased.leaseToken)).toBe('failed'); + expect(await store.claim(30)).toBeNull(); + }); +}); diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/contracts.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/contracts.ts new file mode 100644 index 00000000..2599db7c --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/contracts.ts @@ -0,0 +1,102 @@ +import type { DependencyEnd } from '@internal/core'; +import { dependency, string } from '@internal/core'; +import { assertDefined } from '@internal/foundation/assertions'; +import { blindCast } from '@internal/foundation/casts'; +import { contract, makeClient, perBindingToken, rpc } from '@internal/service-rpc'; +import { type } from 'arktype'; +import type { QueueDefinition, QueueDefinitions } from './definitions.ts'; + +const unknownValue = type('unknown'); +const queueMessage = type({ + id: 'string', + queue: 'string', + body: unknownValue, + attempt: 'number.integer >= 1', + enqueuedAt: 'string', +}); +const queueClaim = type({ message: queueMessage, leaseToken: 'string' }).or({ message: 'null' }); + +export const queueProducerContract = contract({ + send: rpc({ + input: type({ queue: 'string', body: unknownValue }), + output: type({ id: 'string' }), + }), +}); + +export const queueControlContract = contract({ + claim: rpc({ + input: type({ leaseSeconds: '1 <= number.integer <= 300' }), + output: queueClaim, + }), + complete: rpc({ + input: type({ messageId: 'string', leaseToken: 'string' }), + output: type({ ok: 'boolean' }), + }), + release: rpc({ + input: type({ messageId: 'string', leaseToken: 'string' }), + output: type({ outcome: type("'retrying' | 'failed' | 'lost'") }), + }), +}); + +export const queueConsumerContract = contract({ + deliver: rpc({ input: queueMessage, output: type({ ok: 'boolean' }) }), +}); + +type MessageOf = Definition extends QueueDefinition ? Message : never; + +export interface QueueHandle { + send(message: Message): Promise<{ id: string }>; +} + +export type QueueProducer = { + readonly [Name in keyof Definitions]: QueueHandle>; +}; + +/** Declares a typed producer dependency with one queue handle per catalog entry. */ +export function queueProducer( + definitions: Definitions, +): DependencyEnd, typeof queueProducerContract> { + return dependency({ + type: 'rpc', + connection: { + params: { + url: string(), + serviceKey: string({ optional: true, provision: perBindingToken() }), + }, + hydrate: ({ url, serviceKey }) => { + const client = makeClient(queueProducerContract, url, { serviceKey }); + const producer: Record> = {}; + + for (const queueName of Object.keys(definitions)) { + const definition = definitions[queueName]; + assertDefined( + definition, + `queueProducer(): unreachable missing definition for "${queueName}".`, + ); + producer[queueName] = { + send: async (message) => { + const validated = definition.message(message); + if (validated instanceof type.errors) { + throw new Error( + `queue.${queueName}.send(): message does not match its schema: ${validated.summary}`, + ); + } + return client.send({ queue: queueName, body: validated }); + }, + }; + } + + return blindCast< + QueueProducer, + 'assembled from the literal queue catalog; every handle validates with its own schema before calling the shared producer contract' + >(producer); + }, + }, + required: queueProducerContract, + }); +} + +/** Declares the consumer port a dispatcher calls. Message typing is applied by serveQueues(). */ +export function queueConsumer(): typeof queueConsumerContract { + return queueConsumerContract; +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/definitions.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/definitions.ts new file mode 100644 index 00000000..35b6b1f2 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/definitions.ts @@ -0,0 +1,74 @@ +import type { Type } from 'arktype'; + +export type QueueDuration = `${number}${'s' | 'm' | 'h'}`; + +export interface FixedBackoff { + readonly kind: 'fixed'; + readonly delaySeconds: number; +} + +export interface QueueRetryPolicy { + readonly maxAttempts?: number; + readonly delay?: FixedBackoff; +} + +export interface QueueDefinition { + readonly message: Type; + readonly retry?: QueueRetryPolicy; +} + +// biome-ignore lint/suspicious/noExplicitAny: a queue catalog is heterogeneous; each key carries its own message type. +export type QueueDefinitions = Record>; + +/** Creates a serializable fixed delay used after every failed delivery attempt. */ +export function fixedBackoff(opts: { readonly delay: QueueDuration }): FixedBackoff { + const match = /^(\d+)(s|m|h)$/.exec(opts.delay); + if (match === null) { + throw new Error( + 'fixedBackoff(): delay must use whole seconds, minutes, or hours, for example "30s".', + ); + } + + const value = Number(match[1] ?? Number.NaN); + const unit = match[2]; + const multiplier = unit === 'h' ? 3600 : unit === 'm' ? 60 : 1; + const delaySeconds = value * multiplier; + if (delaySeconds < 1 || delaySeconds > 86_400) { + throw new Error('fixedBackoff(): delay must be between 1 second and 24 hours.'); + } + + return Object.freeze({ kind: 'fixed', delaySeconds }); +} + +/** Defines queue names and their message schemas without provisioning anything. */ +export function defineQueues< + // biome-ignore lint/suspicious/noExplicitAny: the self-referential bound preserves each queue's independent message type. + const Definitions extends { [Name in keyof Definitions]: QueueDefinition }, +>(definitions: Definitions): Definitions { + const catalog: QueueDefinitions = definitions; + for (const [queueName, definition] of Object.entries(catalog)) { + const maxAttempts = definition.retry?.maxAttempts; + if ( + maxAttempts !== undefined && + (!Number.isInteger(maxAttempts) || maxAttempts < 1 || maxAttempts > 100) + ) { + throw new Error( + `defineQueues(): queue "${queueName}" retry.maxAttempts must be an integer from 1 through 100.`, + ); + } + + const delay = definition.retry?.delay; + if ( + delay !== undefined && + (delay.kind !== 'fixed' || + !Number.isInteger(delay.delaySeconds) || + delay.delaySeconds < 1 || + delay.delaySeconds > 86_400) + ) { + throw new Error( + `defineQueues(): queue "${queueName}" retry.delay must come from fixedBackoff().`, + ); + } + } + return definitions; +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/dispatcher-service.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/dispatcher-service.ts new file mode 100644 index 00000000..eb6254aa --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/dispatcher-service.ts @@ -0,0 +1,28 @@ +import node from '@internal/node'; +import { compute } from '@internal/prisma-cloud'; +import { rpc } from '@internal/service-rpc'; +import { type } from 'arktype'; +import { queueConsumerContract, queueControlContract } from './contracts.ts'; + +const dispatcherInput = type({ + pollIntervalMs: '10 <= number.integer <= 60000', + leaseSeconds: '1 <= number.integer <= 300', +}); + +/** The always-running Compute service that moves messages from Postgres to a consumer. */ +export function queueDispatcher() { + return compute({ + name: 'queueDispatcher', + deps: { + queue: rpc(queueControlContract), + consumer: rpc(queueConsumerContract), + }, + input: dispatcherInput, + build: node({ + module: new URL('./dispatcher-service.mjs', import.meta.url).href, + entry: './dispatcher-entrypoint.mjs', + }), + }); +} + +export default queueDispatcher(); diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/dispatcher-entrypoint.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/dispatcher-entrypoint.ts new file mode 100644 index 00000000..c3650632 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/dispatcher-entrypoint.ts @@ -0,0 +1,9 @@ +import { queueDispatcher } from '../dispatcher-service.ts'; +import { runDispatcher } from './dispatcher.ts'; + +const service = queueDispatcher(); +const clients = service.load(); +const input = service.input(); + +console.info(`queue dispatcher ready with a ${input.pollIntervalMs}ms poll interval`); +runDispatcher(clients, input); diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/dispatcher.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/dispatcher.ts new file mode 100644 index 00000000..1964ff97 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/dispatcher.ts @@ -0,0 +1,70 @@ +import type { Client } from '@internal/service-rpc'; +import type { queueConsumerContract, queueControlContract } from '../contracts.ts'; + +export interface DispatcherClients { + readonly queue: Client; + readonly consumer: Client; +} + +/** Claims and delivers at most one message. Returns false when the queue is empty. */ +export async function dispatchOnce( + clients: DispatcherClients, + opts: { readonly leaseSeconds: number }, +): Promise { + const claimed = await clients.queue.claim({ leaseSeconds: opts.leaseSeconds }); + if (claimed.message === null) return false; + + try { + const result = await clients.consumer.deliver(claimed.message); + if (!result.ok) { + const released = await clients.queue.release({ + messageId: claimed.message.id, + leaseToken: claimed.leaseToken, + }); + if (released.outcome === 'lost') { + throw new Error(`queue dispatcher lost the lease for message "${claimed.message.id}"`); + } + return true; + } + } catch (error) { + const released = await clients.queue.release({ + messageId: claimed.message.id, + leaseToken: claimed.leaseToken, + }); + if (released.outcome === 'lost') { + throw new Error(`queue dispatcher lost the lease for message "${claimed.message.id}"`); + } + throw error; + } + + const completed = await clients.queue.complete({ + messageId: claimed.message.id, + leaseToken: claimed.leaseToken, + }); + if (!completed.ok) { + throw new Error(`queue dispatcher lost the lease for message "${claimed.message.id}"`); + } + return true; +} + +/** Starts the continuously running delivery loop. */ +export function runDispatcher( + clients: DispatcherClients, + opts: { + readonly leaseSeconds: number; + readonly pollIntervalMs: number; + }, +): void { + const tick = async (): Promise => { + let delay = 0; + try { + const delivered = await dispatchOnce(clients, opts); + delay = delivered ? 0 : opts.pollIntervalMs; + } catch (error) { + console.error('queue dispatcher delivery failed', error); + delay = opts.pollIntervalMs; + } + setTimeout(() => void tick(), delay); + }; + void tick(); +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/handlers.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/handlers.ts new file mode 100644 index 00000000..cbb7d55f --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/handlers.ts @@ -0,0 +1,38 @@ +import type { RpcHandlerContext } from '@internal/service-rpc'; +import type { QueueStore } from '../queue-store.ts'; + +export function createQueueHandlers(opts: { + readonly store: QueueStore; + readonly queueNames: readonly string[]; +}) { + const knownQueues = new Set(opts.queueNames); + + return { + send: async ( + input: { queue: string; body: unknown }, + _deps: unknown, + context: RpcHandlerContext, + ): Promise<{ id: string }> => { + if (!knownQueues.has(input.queue)) { + throw new Error(`unknown queue "${input.queue}"`); + } + return opts.store.enqueue({ + queue: input.queue, + body: input.body, + enqueueKey: context.idempotencyKey ?? crypto.randomUUID(), + }); + }, + claim: async (input: { leaseSeconds: number }) => { + const leased = await opts.store.claim(input.leaseSeconds); + return leased === null + ? { message: null } + : { message: leased.message, leaseToken: leased.leaseToken }; + }, + complete: async (input: { messageId: string; leaseToken: string }) => ({ + ok: await opts.store.complete(input.messageId, input.leaseToken), + }), + release: async (input: { messageId: string; leaseToken: string }) => ({ + outcome: await opts.store.release(input.messageId, input.leaseToken), + }), + }; +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/pg-queue-store.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/pg-queue-store.ts new file mode 100644 index 00000000..4a7f9530 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/pg-queue-store.ts @@ -0,0 +1,198 @@ +import { retryTransientConnect } from '@internal/prisma-cloud/connection'; +import { SQL } from 'bun'; +import type { + LeasedQueueMessage, + QueueReleaseOutcome, + QueueRuntimeConfiguration, + QueueStore, +} from '../queue-store.ts'; + +interface MessageRow { + readonly id: string; + readonly queue_name: string; + readonly body: unknown; + readonly attempts: number; + readonly created_at: Date | string; + readonly lease_token: string; +} + +function toIso(value: Date | string): string { + return value instanceof Date ? value.toISOString() : new Date(value).toISOString(); +} + +class PgQueueStore implements QueueStore { + constructor(private readonly sql: SQL) {} + + async enqueue(input: { + readonly queue: string; + readonly body: unknown; + readonly enqueueKey: string; + }): Promise<{ id: string }> { + const id = crypto.randomUUID(); + const inserted: Array<{ id: string }> = await this.sql` + insert into queue_messages (id, queue_name, body, enqueue_key) + values (${id}, ${input.queue}, ${input.body}, ${input.enqueueKey}) + on conflict (enqueue_key) do nothing + returning id`; + const insertedRow = inserted[0]; + if (insertedRow !== undefined) return { id: insertedRow.id }; + + const existing: Array<{ id: string }> = await this.sql` + select id from queue_messages where enqueue_key = ${input.enqueueKey}`; + const existingRow = existing[0]; + if (existingRow === undefined) { + throw new Error('queue enqueue conflict did not return its existing message'); + } + return { id: existingRow.id }; + } + + async claim(leaseSeconds: number): Promise { + const leaseToken = crypto.randomUUID(); + const rows: MessageRow[] = await this.sql` + with exhausted as ( + update queue_messages as message + set + state = 'failed', + failed_at = now(), + lease_token = null, + leased_until = null + from queue_configurations as config + where message.queue_name = config.queue_name + and message.state = 'leased' + and message.leased_until <= now() + and message.attempts >= config.max_attempts + ), candidate as ( + select message.id + from queue_messages as message + join queue_configurations as config on config.queue_name = message.queue_name + where + message.attempts < config.max_attempts + and ( + (message.state = 'ready' and message.available_at <= now()) + or (message.state = 'leased' and message.leased_until <= now()) + ) + order by message.available_at, message.created_at, message.id + for update skip locked + limit 1 + ) + update queue_messages as message + set + state = 'leased', + lease_token = ${leaseToken}, + leased_until = now() + (${leaseSeconds} * interval '1 second'), + attempts = attempts + 1 + from candidate + where message.id = candidate.id + returning message.id, message.queue_name, message.body, message.attempts, + message.created_at, message.lease_token`; + const row = rows[0]; + if (row === undefined) return null; + return { + leaseToken: row.lease_token, + message: { + id: row.id, + queue: row.queue_name, + body: row.body, + attempt: Number(row.attempts), + enqueuedAt: toIso(row.created_at), + }, + }; + } + + async complete(messageId: string, leaseToken: string): Promise { + const rows = await this.sql` + update queue_messages + set state = 'completed', completed_at = now(), lease_token = null, leased_until = null + where id = ${messageId} and state = 'leased' and lease_token = ${leaseToken} + returning id`; + return rows.length === 1; + } + + async release(messageId: string, leaseToken: string): Promise { + const rows: Array<{ state: string }> = await this.sql` + update queue_messages as message + set + state = case when message.attempts >= config.max_attempts then 'failed' else 'ready' end, + available_at = case + when message.attempts >= config.max_attempts then message.available_at + else now() + (config.retry_delay_ms * interval '1 millisecond') + end, + lease_token = null, + leased_until = null, + failed_at = case when message.attempts >= config.max_attempts then now() else null end + from queue_configurations as config + where message.queue_name = config.queue_name + and message.id = ${messageId} + and message.state = 'leased' + and message.lease_token = ${leaseToken} + returning message.state`; + const row = rows[0]; + if (row === undefined) return 'lost'; + return row.state === 'failed' ? 'failed' : 'retrying'; + } +} + +/** Connects to Postgres and applies the walking-skeleton queue schema idempotently. */ +export async function createPgQueueStore( + url: string, + queues: readonly QueueRuntimeConfiguration[], +): Promise { + const sql = new SQL({ url, max: 4, idleTimeout: 10 }); + await retryTransientConnect(async () => { + await sql` + create table if not exists queue_configurations ( + queue_name text primary key, + max_attempts integer not null check (max_attempts between 1 and 100), + retry_delay_ms integer not null check (retry_delay_ms between 1000 and 86400000), + updated_at timestamptz not null default now() + )`; + await sql` + create table if not exists queue_messages ( + id text primary key, + queue_name text not null, + body jsonb not null, + enqueue_key text not null unique, + state text not null default 'ready' + check (state in ('ready', 'leased', 'completed', 'failed')), + attempts integer not null default 0, + available_at timestamptz not null default now(), + lease_token text, + leased_until timestamptz, + created_at timestamptz not null default now(), + completed_at timestamptz, + failed_at timestamptz + )`; + await sql`alter table queue_messages add column if not exists failed_at timestamptz`; + await sql` + do $$ + declare current_definition text; + begin + select pg_get_constraintdef(oid) + into current_definition + from pg_constraint + where conrelid = 'queue_messages'::regclass + and conname = 'queue_messages_state_check'; + + if current_definition is null or current_definition not like '%failed%' then + alter table queue_messages drop constraint if exists queue_messages_state_check; + alter table queue_messages add constraint queue_messages_state_check + check (state in ('ready', 'leased', 'completed', 'failed')); + end if; + end $$`; + await sql` + create index if not exists queue_messages_claim_idx + on queue_messages (available_at, created_at, id) + where state in ('ready', 'leased')`; + + for (const queue of queues) { + await sql` + insert into queue_configurations (queue_name, max_attempts, retry_delay_ms) + values (${queue.name}, ${queue.maxAttempts}, ${queue.retryDelayMs}) + on conflict (queue_name) do update set + max_attempts = excluded.max_attempts, + retry_delay_ms = excluded.retry_delay_ms, + updated_at = now()`; + } + }); + return new PgQueueStore(sql); +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/queue-entrypoint.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/queue-entrypoint.ts new file mode 100644 index 00000000..b24086e4 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/execution/queue-entrypoint.ts @@ -0,0 +1,25 @@ +import { serve } from '@internal/service-rpc'; +import { queueService } from '../queue-service.ts'; +import { createQueueHandlers } from './handlers.ts'; +import { createPgQueueStore } from './pg-queue-store.ts'; + +const service = queueService(); +const { db } = service.load(); +const { queues } = service.input(); +const store = await createPgQueueStore(db.url, queues); +console.info(`queue service ready with ${queues.length} queue definition(s)`); +const handlers = createQueueHandlers({ + store, + queueNames: queues.map((queue) => queue.name), +}); + +const fetchHandler = serve(service, { + producer: { send: handlers.send }, + dispatch: { + claim: handlers.claim, + complete: handlers.complete, + release: handlers.release, + }, +}); + +Bun.serve({ port: service.port(), hostname: '0.0.0.0', fetch: fetchHandler }); diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-entrypoint.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-entrypoint.ts new file mode 100644 index 00000000..c59b3677 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-entrypoint.ts @@ -0,0 +1 @@ +import '../execution/dispatcher-entrypoint.ts'; diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-service.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-service.ts new file mode 100644 index 00000000..0fe764e7 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/dispatcher-service.ts @@ -0,0 +1 @@ +export { default } from '../dispatcher-service.ts'; diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/index.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/index.ts new file mode 100644 index 00000000..422b3114 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/index.ts @@ -0,0 +1,21 @@ +export type { QueueHandle, QueueProducer } from '../contracts.ts'; +export { + queueConsumer, + queueConsumerContract, + queueControlContract, + queueProducer, + queueProducerContract, +} from '../contracts.ts'; +export type { + FixedBackoff, + QueueDefinition, + QueueDefinitions, + QueueDuration, + QueueRetryPolicy, +} from '../definitions.ts'; +export { defineQueues, fixedBackoff } from '../definitions.ts'; +export { queueDispatcher } from '../dispatcher-service.ts'; +export { queues } from '../queues-module.ts'; +export type { QueueHandlers } from '../serve-queues.ts'; +export { serveQueues } from '../serve-queues.ts'; +export type { QueueConsumerMessage } from '../types.ts'; diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-entrypoint.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-entrypoint.ts new file mode 100644 index 00000000..8e73a74b --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-entrypoint.ts @@ -0,0 +1 @@ +import '../execution/queue-entrypoint.ts'; diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-service.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-service.ts new file mode 100644 index 00000000..1d74099d --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/exports/queue-service.ts @@ -0,0 +1 @@ +export { default } from '../queue-service.ts'; diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/queue-service.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/queue-service.ts new file mode 100644 index 00000000..a0fb3fc3 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/queue-service.ts @@ -0,0 +1,27 @@ +import node from '@internal/node'; +import { compute, postgres } from '@internal/prisma-cloud'; +import { type } from 'arktype'; +import { queueControlContract, queueProducerContract } from './contracts.ts'; + +const queueServiceInput = type({ + queues: type({ + name: 'string', + maxAttempts: '1 <= number.integer <= 100', + retryDelayMs: '1000 <= number.integer <= 86400000', + }).array(), +}); + +export function queueService() { + return compute({ + name: 'queues', + deps: { db: postgres() }, + input: queueServiceInput, + expose: { producer: queueProducerContract, dispatch: queueControlContract }, + build: node({ + module: new URL('./queue-service.mjs', import.meta.url).href, + entry: './queue-entrypoint.mjs', + }), + }); +} + +export default queueService(); diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/queue-store.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/queue-store.ts new file mode 100644 index 00000000..7d4b0965 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/queue-store.ts @@ -0,0 +1,31 @@ +export interface StoredQueueMessage { + readonly id: string; + readonly queue: string; + readonly body: unknown; + readonly attempt: number; + readonly enqueuedAt: string; +} + +export interface LeasedQueueMessage { + readonly message: StoredQueueMessage; + readonly leaseToken: string; +} + +export interface QueueRuntimeConfiguration { + readonly name: string; + readonly maxAttempts: number; + readonly retryDelayMs: number; +} + +export type QueueReleaseOutcome = 'retrying' | 'failed' | 'lost'; + +export interface QueueStore { + enqueue(input: { + readonly queue: string; + readonly body: unknown; + readonly enqueueKey: string; + }): Promise<{ id: string }>; + claim(leaseSeconds: number): Promise; + complete(messageId: string, leaseToken: string): Promise; + release(messageId: string, leaseToken: string): Promise; +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/queues-module.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/queues-module.ts new file mode 100644 index 00000000..2e71aa6c --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/queues-module.ts @@ -0,0 +1,36 @@ +import type { ModuleNode } from '@internal/core'; +import { module } from '@internal/core'; +import { postgres } from '@internal/prisma-cloud'; +import { queueControlContract, queueProducerContract } from './contracts.ts'; +import type { QueueDefinitions } from './definitions.ts'; +import { queueService } from './queue-service.ts'; + +/** Provisions the durable queue database and the service that owns it. */ +export function queues(opts: { + readonly definitions: QueueDefinitions; + readonly name?: string; +}): ModuleNode< + Record, + { producer: typeof queueProducerContract; dispatch: typeof queueControlContract }, + Record +> { + return module( + opts.name ?? 'queues', + { expose: { producer: queueProducerContract, dispatch: queueControlContract } }, + ({ provision }) => { + const db = provision(postgres({ name: 'db' }), { id: 'db' }); + const service = provision(queueService(), { + id: 'service', + deps: { db }, + input: { + queues: Object.entries(opts.definitions).map(([name, definition]) => ({ + name, + maxAttempts: definition.retry?.maxAttempts ?? 5, + retryDelayMs: (definition.retry?.delay?.delaySeconds ?? 5) * 1000, + })), + }, + }); + return { producer: service.producer, dispatch: service.dispatch }; + }, + ); +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/serve-queues.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/serve-queues.ts new file mode 100644 index 00000000..84fb7236 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/serve-queues.ts @@ -0,0 +1,67 @@ +import type { Deps, HydratedDeps, Params, RunnableServiceNode } from '@internal/core'; +import { assertDefined } from '@internal/foundation/assertions'; +import { blindCast } from '@internal/foundation/casts'; +import type { Handlers } from '@internal/service-rpc'; +import { serve } from '@internal/service-rpc'; +import { type } from 'arktype'; +import type { queueConsumerContract } from './contracts.ts'; +import type { QueueDefinition, QueueDefinitions } from './definitions.ts'; +import type { QueueConsumerMessage } from './types.ts'; + +type MessageOf = Definition extends QueueDefinition ? Message : never; + +export type QueueHandlers = { + readonly [Name in keyof Definitions]: ( + message: QueueConsumerMessage>, + deps: LoadedDeps, + ) => Promise; +}; + +type UntypedHandler = (message: QueueConsumerMessage, deps: unknown) => Promise; + +/** Validates a delivered body and routes it to the handler for its queue name. */ +export function serveQueues( + service: RunnableServiceNode, + definitions: Definitions, + handlers: QueueHandlers>, +): (request: Request) => Promise { + const byQueue = blindCast< + Record, + 'the caller supplies an exhaustive handler map keyed by the same queue catalog used for runtime dispatch' + >(handlers); + + const deliver = async ( + message: QueueConsumerMessage, + deps: HydratedDeps, + ): Promise<{ ok: boolean }> => { + const definition = definitions[message.queue]; + const handler = byQueue[message.queue]; + assertDefined(definition, `serveQueues(): unknown queue "${message.queue}".`); + assertDefined(handler, `serveQueues(): no handler for queue "${message.queue}".`); + + const validated = definition.message(message.body); + if (validated instanceof type.errors) { + throw new Error( + `serveQueues(): message for "${message.queue}" does not match its schema: ${validated.summary}`, + ); + } + try { + await handler({ ...message, body: validated }, deps); + return { ok: true }; + } catch (error) { + console.error( + `queue consumer handler failed for "${message.queue}" message "${message.id}" attempt ${message.attempt}`, + error, + ); + return { ok: false }; + } + }; + + return serve( + service, + blindCast< + Handlers, + 'deliver is typed from queueConsumerContract above; the cast only bridges the unresolved generic service projection' + >({ consumer: { deliver } }), + ); +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/src/types.ts b/packages/1-prisma-cloud/2-shared-modules/queues/src/types.ts new file mode 100644 index 00000000..59eaaaba --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/src/types.ts @@ -0,0 +1,7 @@ +export interface QueueConsumerMessage { + readonly id: string; + readonly queue: string; + readonly body: Body; + readonly attempt: number; + readonly enqueuedAt: string; +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/tsconfig.json b/packages/1-prisma-cloud/2-shared-modules/queues/tsconfig.json new file mode 100644 index 00000000..a16ed256 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/tsconfig.json @@ -0,0 +1,7 @@ +{ + "extends": "../../../../tsconfig.base.json", + "compilerOptions": { + "types": ["bun-types"] + }, + "include": ["src"] +} diff --git a/packages/1-prisma-cloud/2-shared-modules/queues/tsdown.config.ts b/packages/1-prisma-cloud/2-shared-modules/queues/tsdown.config.ts new file mode 100644 index 00000000..88c76e96 --- /dev/null +++ b/packages/1-prisma-cloud/2-shared-modules/queues/tsdown.config.ts @@ -0,0 +1,33 @@ +import { baseConfig } from '@internal/tsdown-config'; +import { defineConfig } from 'tsdown'; + +export default defineConfig([ + { + ...baseConfig, + entry: { + index: 'src/exports/index.ts', + 'queue-service': 'src/exports/queue-service.ts', + 'dispatcher-service': 'src/exports/dispatcher-service.ts', + }, + exports: false, + clean: true, + }, + { + ...baseConfig, + entry: { 'queue-entrypoint': 'src/exports/queue-entrypoint.ts' }, + exports: false, + clean: false, + skipNodeModulesBundle: false, + external: [/^bun$/, /^bun:/], + noExternal: [/^@internal\//, /^@prisma\//, /^arktype/, /^@standard-schema\//], + }, + { + ...baseConfig, + entry: { 'dispatcher-entrypoint': 'src/exports/dispatcher-entrypoint.ts' }, + exports: false, + clean: false, + skipNodeModulesBundle: false, + external: [/^bun$/, /^bun:/], + noExternal: [/^@internal\//, /^@prisma\//, /^arktype/, /^@standard-schema\//], + }, +]); diff --git a/packages/9-public/composer-prisma-cloud/package.json b/packages/9-public/composer-prisma-cloud/package.json index 381195d1..464d7f03 100644 --- a/packages/9-public/composer-prisma-cloud/package.json +++ b/packages/9-public/composer-prisma-cloud/package.json @@ -27,6 +27,9 @@ "./streams": "./dist/streams/index.mjs", "./streams/streams-entrypoint": "./dist/streams/streams-entrypoint.mjs", "./streams/testing": "./dist/streams/testing.mjs", + "./queues": "./dist/queues/index.mjs", + "./queues/queue-entrypoint": "./dist/queues/queue-entrypoint.mjs", + "./queues/dispatcher-entrypoint": "./dist/queues/dispatcher-entrypoint.mjs", "./package.json": "./package.json" }, "files": [ @@ -71,13 +74,15 @@ "@internal/nextjs": "workspace:0.3.0", "@internal/node": "workspace:0.3.0", "@internal/prisma-cloud": "workspace:0.3.0", + "@internal/queues": "workspace:0.3.0", "@internal/s3-protocol": "workspace:0.3.0", "@internal/service-rpc": "workspace:0.3.0", "@internal/storage": "workspace:0.3.0", "@internal/streams": "workspace:0.3.0", "@internal/tsdown-config": "workspace:0.3.0", "@types/node": "^26.0.1", - "typescript": "^6.0.3" + "typescript": "^6.0.3", + "unrun": "^0.3.1" }, "license": "Apache-2.0", "repository": { diff --git a/packages/9-public/composer-prisma-cloud/src/exports/queues.ts b/packages/9-public/composer-prisma-cloud/src/exports/queues.ts new file mode 100644 index 00000000..48d757cf --- /dev/null +++ b/packages/9-public/composer-prisma-cloud/src/exports/queues.ts @@ -0,0 +1 @@ +export * from '@internal/queues'; diff --git a/packages/9-public/composer-prisma-cloud/tsdown.config.ts b/packages/9-public/composer-prisma-cloud/tsdown.config.ts index 8ef835e0..b0153be7 100644 --- a/packages/9-public/composer-prisma-cloud/tsdown.config.ts +++ b/packages/9-public/composer-prisma-cloud/tsdown.config.ts @@ -58,6 +58,7 @@ const cronDist = '../../1-prisma-cloud/2-shared-modules/cron/dist'; const storageDist = '../../1-prisma-cloud/2-shared-modules/storage/dist'; const emailDist = '../../1-prisma-cloud/2-shared-modules/email/dist'; const streamsDist = '../../1-prisma-cloud/2-shared-modules/streams/dist'; +const queuesDist = '../../1-prisma-cloud/2-shared-modules/queues/dist'; const devEmulatorsDist = '../../1-prisma-cloud/0-lowering/dev-emulators/dist'; export default defineConfig([ { @@ -289,4 +290,34 @@ export default defineConfig([ noExternal: [/^@internal\//], plugins: [externalizeFramework], }, + { + ...baseConfig, + entry: { index: 'src/exports/queues.ts' }, + outDir: 'dist/queues', + exports: false, + clean: false, + skipNodeModulesBundle: false, + noExternal: [/^@internal\//], + plugins: [externalizeFramework], + }, + { + // Both reusable services resolve their entrypoints beside their own built + // service module. The internal package has already fully inlined each + // program; this pass places the four files in the published /queues tree. + ...baseConfig, + dts: false, + entry: { + 'queue-service': `${queuesDist}/queue-service.mjs`, + 'queue-entrypoint': `${queuesDist}/queue-entrypoint.mjs`, + 'dispatcher-service': `${queuesDist}/dispatcher-service.mjs`, + 'dispatcher-entrypoint': `${queuesDist}/dispatcher-entrypoint.mjs`, + }, + outDir: 'dist/queues', + exports: false, + clean: false, + skipNodeModulesBundle: false, + external: [/^bun$/, /^bun:/], + noExternal: [/^@internal\//], + plugins: [externalizeFramework], + }, ]); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index f5aaf3af..f4477b20 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -204,6 +204,31 @@ importers: specifier: ^6.0.3 version: 6.0.3 + examples/queues: + dependencies: + '@prisma/composer': + specifier: workspace:0.3.0 + version: link:../../packages/9-public/composer + '@prisma/composer-prisma-cloud': + specifier: workspace:0.3.0 + version: link:../../packages/9-public/composer-prisma-cloud + arktype: + specifier: ^2.2.3 + version: 2.2.3 + devDependencies: + '@types/bun': + specifier: ^1.3.13 + version: 1.3.14 + tsdown: + specifier: ^0.22.7 + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) + typescript: + specifier: ^6.0.3 + version: 6.0.3 + unrun: + specifier: ^0.3.1 + version: 0.3.1 + examples/storage: dependencies: '@aws-sdk/client-s3': @@ -500,7 +525,7 @@ importers: version: 1.3.14 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -516,7 +541,7 @@ importers: version: link:../tsdown-config tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -525,7 +550,7 @@ importers: devDependencies: tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -553,7 +578,7 @@ importers: version: 1.3.14 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -581,7 +606,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -606,7 +631,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -634,7 +659,7 @@ importers: version: 1.3.14 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -659,7 +684,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -696,7 +721,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -742,7 +767,7 @@ importers: version: 8.22.0 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -785,7 +810,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -822,7 +847,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -846,7 +871,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -928,7 +953,7 @@ importers: version: 8.20.0 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -1004,7 +1029,7 @@ importers: version: 8.20.0 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -1041,7 +1066,7 @@ importers: version: 1.3.14 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -1084,7 +1109,7 @@ importers: version: 8.0.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -1092,6 +1117,46 @@ importers: specifier: ^4.1.10 version: 4.1.10(@types/node@26.1.1)(vite@8.1.2(@types/node@26.1.1)(esbuild@0.28.1)(jiti@2.6.1)(yaml@2.9.0)) + packages/1-prisma-cloud/2-shared-modules/queues: + dependencies: + '@internal/core': + specifier: workspace:0.3.0 + version: link:../../../0-framework/1-core/core + '@internal/foundation': + specifier: workspace:0.3.0 + version: link:../../../0-framework/0-foundation/foundation + '@internal/node': + specifier: workspace:0.3.0 + version: link:../../../0-framework/2-authoring/node + '@internal/prisma-cloud': + specifier: workspace:0.3.0 + version: link:../../1-extensions/target + '@internal/service-rpc': + specifier: workspace:0.3.0 + version: link:../../../0-framework/2-authoring/service-rpc + arktype: + specifier: ^2.2.3 + version: 2.2.3 + devDependencies: + '@internal/tsdown-config': + specifier: workspace:0.3.0 + version: link:../../../0-framework/0-foundation/tsdown-config + '@types/bun': + specifier: ^1.3.13 + version: 1.3.14 + tsdown: + specifier: ^0.22.7 + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) + typescript: + specifier: ^6.0.3 + version: 6.0.3 + unrun: + specifier: ^0.3.1 + version: 0.3.1 + vitest: + specifier: ^4.1.10 + version: 4.1.10(@types/node@26.1.1)(vite@8.1.2(@types/node@26.1.1)(esbuild@0.28.1)(jiti@2.6.1)(yaml@2.9.0)) + packages/1-prisma-cloud/2-shared-modules/storage: dependencies: '@internal/core': @@ -1121,7 +1186,7 @@ importers: version: 1.3.14 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -1164,7 +1229,7 @@ importers: version: 1.3.14 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -1234,7 +1299,7 @@ importers: version: 26.1.1 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) typescript: specifier: ^6.0.3 version: 6.0.3 @@ -1297,7 +1362,7 @@ importers: version: 3.4.9 tsdown: specifier: ^0.22.7 - version: 0.22.12(typescript@6.0.3) + version: 0.22.12(typescript@6.0.3)(unrun@0.3.1) devDependencies: '@internal/auth': specifier: workspace:0.3.0 @@ -1332,6 +1397,9 @@ importers: '@internal/prisma-cloud': specifier: workspace:0.3.0 version: link:../../1-prisma-cloud/1-extensions/target + '@internal/queues': + specifier: workspace:0.3.0 + version: link:../../1-prisma-cloud/2-shared-modules/queues '@internal/s3-protocol': specifier: workspace:0.3.0 version: link:../../1-prisma-cloud/0-lowering/s3-protocol @@ -1353,6 +1421,9 @@ importers: typescript: specifier: ^6.0.3 version: 6.0.3 + unrun: + specifier: ^0.3.1 + version: 0.3.1 test/integration: devDependencies: @@ -5445,6 +5516,16 @@ packages: universal-user-agent@7.0.3: resolution: {integrity: sha512-TmnEAEAsBJVZM/AADELsK76llnwcf9vMKuPz8JflO1frO8Lchitr0fNaN9d+Ap0BjKtqWqd/J17qeDnXh8CL2A==} + unrun@0.3.1: + resolution: {integrity: sha512-onIck/oNnCaytwths1ZVp1LK2Gq2hPoyFhiHebObuUXqR3S0uHuLLaBK8K6mRRgV7Ptip8AnNvaUsgzwWwBZuA==} + engines: {node: ^22.13.0 || >=24.0.0} + hasBin: true + peerDependencies: + synckit: ^0.11.11 + peerDependenciesMeta: + synckit: + optional: true + util-deprecate@1.0.2: resolution: {integrity: sha512-EPD5q1uXyFxJpCrLnCc1nHnq3gOa6DZBocAIiI2TaSCA7VCJ1UJDMagCzIkXNsUYfD1daK//LTEQ8xiIbrHtcw==} @@ -9715,7 +9796,7 @@ snapshots: minimist: 1.2.8 strip-bom: 3.0.0 - tsdown@0.22.12(typescript@6.0.3): + tsdown@0.22.12(typescript@6.0.3)(unrun@0.3.1): dependencies: ansis: 4.3.1 cac: 7.0.0 @@ -9734,6 +9815,7 @@ snapshots: verkit: 0.1.2 optionalDependencies: typescript: 6.0.3 + unrun: 0.3.1 transitivePeerDependencies: - '@ts-macro/tsc' - '@typescript/native-preview' @@ -9805,6 +9887,10 @@ snapshots: universal-user-agent@7.0.3: {} + unrun@0.3.1: + dependencies: + rolldown: 1.2.0 + util-deprecate@1.0.2: {} uuid@14.0.1: {} diff --git a/skills/prisma-composer/SKILL.md b/skills/prisma-composer/SKILL.md index 7322ba3b..e1af78a6 100644 --- a/skills/prisma-composer/SKILL.md +++ b/skills/prisma-composer/SKILL.md @@ -53,7 +53,7 @@ Two packages, and only two, appear in your `package.json`: | Package | Provides | | --- | --- | | `@prisma/composer` | Core authoring: `module`, `secret`, `isSecretString`, `/arktype` (the `secretString()` schema leaf), `/rpc`, `/node`, `/nextjs`, `/config`, `/testing`, the `prisma-composer` CLI | -| `@prisma/composer-prisma-cloud` | The Prisma Cloud target: `compute`, `postgres`, `envSecret`, `envParam`, `/control`, `/testing`, and the shared `/cron`, `/storage`, `/streams`, `/prisma-next` modules | +| `@prisma/composer-prisma-cloud` | The Prisma Cloud target: `compute`, `postgres`, `envSecret`, `envParam`, `/control`, `/testing`, and the shared `/cron`, `/storage`, `/streams`, `/queues`, `/prisma-next` modules | ## Anatomy of a service @@ -378,6 +378,7 @@ growing: | `cron` from `/cron` | An always-on scheduler firing your schedule at your runner service | nothing | | `storage` from `/storage` | An S3-backed blob store (own Postgres + minted credentials) | `store` | | `streams` from `/streams` | Durable append-only event streams over a `store` | `streams` | +| `queues` from `/queues` | Persistent work delivery backed by Postgres | `producer`, `dispatch` | **Finding more.** A Composer extension — a package that brings its own Modules, resources, or deploy target — is published on npm under the name @@ -405,6 +406,43 @@ const handler = serveSchedule(service, schedule, { provision(cron({ schedule, runner: runnerService }), { deps: { worker: worker.rpc } }); ``` +Queues use one catalogue for typed producers and consumers. The Module owns +Postgres; a separate always-running dispatcher pushes claimed messages to a +consumer service: + +```ts +const definitions = defineQueues({ + thumbnails: { + message: type({ imageId: 'string' }), + retry: { + maxAttempts: 5, + delay: fixedBackoff({ delay: '5s' }), + }, + }, +}); + +// worker service +const workerService = compute({ + name: 'worker', + deps: { queues: queueProducer(definitions) }, + expose: { consumer: queueConsumer() }, + build: node({ module: import.meta.url, entry: '../dist/server.mjs' }), +}); + +// root module +const queue = provision(queues({ definitions })); +const worker = provision(workerService, { deps: { queues: queue.producer } }); +provision(queueDispatcher(), { + deps: { queue: queue.dispatch, consumer: worker.consumer }, + input: { pollIntervalMs: 250, leaseSeconds: 30 }, +}); +``` + +The worker routes `/rpc/*` to `serveQueues(workerService, definitions, +handlers)`. This early prototype delivers one message per request. Configurable +batches, exponential retry, competing consumers, replay, pause, and operational +APIs are not implemented yet. Use `examples/queues` as the complete reference. + ## Service input Choosing the channel is most of the decision: diff --git a/tsconfig.depcruise.json b/tsconfig.depcruise.json index 81ccd68d..35865536 100644 --- a/tsconfig.depcruise.json +++ b/tsconfig.depcruise.json @@ -117,6 +117,9 @@ "./packages/1-prisma-cloud/2-shared-modules/cron/src/exports/scheduler-entrypoint.ts" ], "@internal/cron": ["./packages/1-prisma-cloud/2-shared-modules/cron/src/exports/index.ts"], + "@internal/queues": [ + "./packages/1-prisma-cloud/2-shared-modules/queues/src/exports/index.ts" + ], "@prisma/composer/config": ["./packages/9-public/composer/src/exports/config.ts"], "@prisma/composer/deploy": ["./packages/9-public/composer/src/exports/deploy.ts"], "@prisma/composer/local-target": ["./packages/9-public/composer/src/exports/local-target.ts"], @@ -149,6 +152,9 @@ "@prisma/composer-prisma-cloud/cron": [ "./packages/9-public/composer-prisma-cloud/src/exports/cron.ts" ], + "@prisma/composer-prisma-cloud/queues": [ + "./packages/9-public/composer-prisma-cloud/src/exports/queues.ts" + ], "@prisma/composer-prisma-cloud": [ "./packages/9-public/composer-prisma-cloud/src/exports/index.ts" ],