Skip to content

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.ConversationId
declare 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


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>
}

Source

Since v0.0.0

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>
}

Source

Since v0.0.0

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
: never

Source

Since v0.0.0

Build a delivery (an identity helper for inference).

Signature

declare const delivery: <E = never, R = never>(delivery: Delivery<E, R>) => Delivery<E, R>

Source

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>

Source

Since v0.0.0

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>

Source

Since v0.0.0