Ingress
Ingress adapters: turn an at-least-once source (a message queue, a webhook log, a Queue) into idempotent
submissions. Each delivery is submitted with requestId = the message’s ID, so Conversation.submit returns the
existing submission for a redelivered message instead of admitting it twice; the delivery is acknowledged only after
the admitting commit succeeded. A crash between admission and acknowledgement leads to a redelivery, which the
requestId absorbs.
import { type Ids, Ingress } from "@effective-harness/core"import * as Effect from "effect/Effect"import * as Stream from "effect/Stream"
declare const conversationId: Ids.ConversationIddeclare const broker: Stream.Stream<{ readonly id: string; readonly text: string; readonly ack: Effect.Effect<void> }>
const program = Ingress.run( Stream.map(broker, (message) => Ingress.delivery({ id: message.id, conversationId, submission: { type: "input", content: message.text }, ack: message.ack }) ))Since v0.0.0
Delivery (interface)
Section titled “Delivery (interface)”One delivery of an at-least-once source.
Signature
export interface Delivery<E = never, R = never> { /** The message's stable ID; redeliveries carry the same one. Used as the submission's `requestId`. */ readonly id: string readonly conversationId: ConversationId readonly submission: SubmissionInput /** Acknowledge the message to its source; runs after the durable admission. */ readonly ack: Effect.Effect<void, E, R>}Since v0.0.0
RunOptions (interface)
Section titled “RunOptions (interface)”Options of run.
Signature
export interface RunOptions<E = never, R = never> { /** * Deliveries submitted at once. Default 1, which admits messages in source order; deliveries of different * conversations may use more. */ readonly concurrency?: number /** * A delivery refused with `ConversationBusy` (an input with `whenBusy: "reject"`). Default: leave it unacknowledged * for the source to redeliver, and go on. */ readonly onBusy?: (delivery: Delivery<any, any>, error: ConversationBusy) => Effect.Effect<void, E, R>}Since v0.0.0
SubmissionInput (type alias)
Section titled “SubmissionInput (type alias)”A submission draft without its requestId: the delivery’s ID takes that place.
Signature
type SubmissionInput = SubmissionDraft extends infer D ? D extends SubmissionDraft ? Omit<D, "requestId"> : never : neverSince v0.0.0
delivery
Section titled “delivery”Build a delivery (an identity helper for inference).
Signature
declare const delivery: <E = never, R = never>(delivery: Delivery<E, R>) => Delivery<E, R>Since v0.0.0
Consume a source until it ends: each delivery is submitted and acknowledged (see submit). A commit failure or a
failed acknowledgement ends the run with that failure; the deliveries it left unacknowledged are redelivered by
the source and deduplicated on the next run.
Signature
declare const run: <EA, RA, E, R, EB = never, RB = never>( source: Stream.Stream<Delivery<EA, RA>, E, R>, options?: RunOptions<EB, RB>) => Effect.Effect<void, E | EA | EB | CommitError, Harness | R | RA | RB>Since v0.0.0
submit
Section titled “submit”Submit one delivery with requestId = its ID, then acknowledge it. A delivery already admitted (a redelivery)
returns the existing submission and is acknowledged again. A failed admission acknowledges nothing, so the source
redelivers.
Signature
declare const submit: <E, R>( delivery: Delivery<E, R>) => Effect.Effect<SubmissionId, CommitError | ConversationBusy | E, Harness | R>Since v0.0.0