Queues
A queue takes a stream of items and drains them through a worker effect — one item at a time, or many in parallel, with priority, de-duplication, retries, and back-pressure. In this toolkit a queue is a resource: you declare it once, and everywhere you yield* MyQueue you get a single handle that does everything — enqueue work, watch it drain, and steer it — through the same value.
That handle is the whole surface. There is no separate producer and admin API: the code that enqueues an email can also pause the queue, read how many are pending, and subscribe to every completion. And because its a resource, the handle reads identically whether the queue runs in this process or across the network — only the layer that provides it changes.
This guide starts at the start: the smallest queue that works, then each piece it was built from.
Your first queue
A queue has two halves: a Tag (what the queue is — its item type and name) and a layer (how it runs — the worker). Here is the whole thing:
import { import QueueHyperlinkQueueHyperlink } from "hyperlink-ts"
import { import EffectEffect, import SchemaSchema } from "effect"
// The item: a plain schema. This is the queue's payload type.
const const EmailJob: Schema.Struct<{
readonly to: Schema.String
readonly subject: Schema.String
}>
const EmailJob: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly to: Schema.String; readonly subject: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type…;
rebuild: (ast: Objects) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Typ…;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String …;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailJob = import SchemaSchema.function Struct<{
readonly to: Schema.String;
readonly subject: Schema.String;
}>(fields: {
readonly to: Schema.String;
readonly subject: Schema.String;
}): Schema.Struct<{
readonly to: Schema.String;
readonly subject: Schema.String;
}>
Defines a struct schema from a map of field schemas.
Details
Each field value is a schema. Use
optionalKey
or
optional
to
mark fields as optional, and
mutableKey
to mark them as mutable.
The resulting schema's Type is a readonly object type with the fields'
decoded types. The Encoded form mirrors the field schemas' encoded types.
Example (Defining a basic struct)
import { Schema } from "effect"
const Person = Schema.Struct({
name: Schema.String,
age: Schema.Number,
email: Schema.optionalKey(Schema.String)
})
// { readonly name: string; readonly age: number; readonly email?: string }
type Person = typeof Person.Type
const alice = Schema.decodeUnknownSync(Person)({ name: "Alice", age: 30 })
console.log(alice)
// { name: 'Alice', age: 30 }
Struct({
to: Schema.String(property) to: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
to: import SchemaSchema.const String: Schema.Stringconst String: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Type-level representation of
String
.
Schema for string values. Validates that the input is typeof "string".
String,
subject: Schema.String(property) subject: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
subject: import SchemaSchema.const String: Schema.Stringconst String: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Type-level representation of
String
.
Schema for string values. Validates that the input is typeof "string".
String,
})
// The tag: the contract. `Self` is the class itself (Effect's two-stage form).
class class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
}
Emails extends import QueueHyperlinkQueueHyperlink.Tag<class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
}
Emails>()("app/Emails", {
payload: Schema.Struct<{
readonly to: Schema.String
readonly subject: Schema.String
}>
(property) payload: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly to: Schema.String; readonly subject: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type…;
rebuild: (ast: Objects) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Typ…;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String …;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
payload: const EmailJob: Schema.Struct<{
readonly to: Schema.String
readonly subject: Schema.String
}>
const EmailJob: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly to: Schema.String; readonly subject: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type…;
rebuild: (ast: Objects) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Typ…;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String …;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailJob,
}) {}
// The layer: the worker. `effect` runs once per item.
const const EmailsLive: anyconst EmailsLive: {
build: (memoMap: MemoMap, scope: Scope) => Effect.Effect<Context<Storage | Emails | Local<Emails>>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailsLive = import QueueHyperlinkQueueHyperlink.layer(class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails, {
effect: (
job: any
) => Effect.Effect<void, never, never>
effect: (job: any(parameter) job: {
to: string;
subject: string;
}
job) => import EffectEffect.const log: (
...message: ReadonlyArray<any>
) => Effect.Effect<void>
Logs one or more messages using the default log level.
Example (Logging at the default level)
import { Effect } from "effect"
const program = Effect.gen(function*() {
yield* Effect.log("Starting computation")
const result = 2 + 2
yield* Effect.log("Result:", result)
yield* Effect.log("Multiple", "values", "can", "be", "logged")
return result
})
Effect.runPromise(program).then(console.log)
// Output:
// timestamp=2023-... level=INFO message="Starting computation"
// timestamp=2023-... level=INFO message="Result: 4"
// timestamp=2023-... level=INFO message="Multiple values can be logged"
// 4
log(`sending "${job: any(parameter) job: {
to: string;
subject: string;
}
job.subject}" to ${job: any(parameter) job: {
to: string;
subject: string;
}
job.to}`),
concurrency: numberconcurrency: 4,
})Thats a complete, running queue. To use it, yield* Emails for the handle and add an item — anywhere the layer is provided:
const const program: Effect.Effect<
void,
unknown,
unknown
>
const program: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
program = import EffectEffect.const gen: <any, void>(f: () => Generator<any, void, never>) => Effect.Effect<void, unknown, unknown> (+1 overload)Provides a way to write effectful code using generator functions, simplifying
control flow and error handling.
When to use
Use when you want to write effectful code that looks and behaves like
synchronous code, while still handling asynchronous tasks, errors, and complex
control flow such as loops and conditions.
Generator functions work similarly to async/await but keep errors,
requirements, and interruption in the Effect type. You can yield* values
from effects and return the final result at the end.
Example (Sequencing effects with generators)
import { Data, Effect } from "effect"
class DiscountRateError extends Data.TaggedError("DiscountRateError")<{}> {}
const addServiceCharge = (amount: number) => amount + 1
const applyDiscount = (
total: number,
discountRate: number
): Effect.Effect<number, DiscountRateError> =>
discountRate === 0
? Effect.fail(new DiscountRateError())
: Effect.succeed(total - (total * discountRate) / 100)
const fetchTransactionAmount = Effect.promise(() => Promise.resolve(100))
const fetchDiscountRate = Effect.promise(() => Promise.resolve(5))
export const program = Effect.gen(function*() {
const transactionAmount = yield* fetchTransactionAmount
const discountRate = yield* fetchDiscountRate
const discountedAmount = yield* applyDiscount(
transactionAmount,
discountRate
)
const finalAmount = addServiceCharge(discountedAmount)
return `Final amount to charge: ${finalAmount}`
})
gen(function* () {
const const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails = yield* class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails
yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.add({ to: stringto: "[email protected]", subject: stringsubject: "Welcome" })
})Provide EmailsLive to program and the item drains through the worker. Nothing else is wired: the worker pool, the retry machinery, and the observability store all come with the layer.
The handle
Hover emails and youll see its type — the named handle:
const const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails = yield* class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails
})QueueHyperlink<{ to: string; subject: string }, void, never, never> reads as QueueHyperlink<Payload, Success, Error, Requirements>:
Payload — the decoded item type. What
addaccepts.Success — the workers return value (here
void; see Success values).Error — the workers typed failure channel (here
never— this queues worker cant fail in a typed way; see When work fails).Requirements — what a local
yield*needs (never); for a remote client its the transport.
The handle groups its members by what theyre for:
Enqueue —
add,prioritize,defer,enqueue.Observe —
size,isEmpty,status,events,metrics.Control —
start,pause,resume,shutdown,clear.Route —
release,deadLetter,drop.
The rest of this guide walks those groups.
The item, and its schema
The payload schema is the single source of truth for the item type. It is a real Effect Schema, not just a type: it decodes the item on the way in and — when the queue is served over RPC — validates it on the wire, so a bad item is rejected before it ever reaches a worker. The decoded type flows everywhere: add(item), the workers argument, and every event that carries the item.
Use whatever Schema shape fits — structs, unions, branded strings, nested data. The only rule is that the payload is a single schema (a Schema.Struct is the common case), so the wire contract is unambiguous.
The worker
The worker lives on the layer. effect is the only required field; the rest tune how it drains:
const const EmailsLive: anyconst EmailsLive: {
build: (memoMap: MemoMap, scope: Scope) => Effect.Effect<Context<Storage | Emails | Local<Emails>>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailsLive = import QueueHyperlinkQueueHyperlink.layer(class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails, {
// runs once per item; the second arg is per-attempt context
effect: (
job: any,
ctx: any
) => Effect.Effect<void, never, never>
effect: (job: any(parameter) job: {
to: string;
subject: string;
}
job, ctx: any(parameter) ctx: {
add: QueueEnqueue<T, EEnqueue, R>;
prioritize: QueueEnqueue<T, EEnqueue, R>;
defer: QueueEnqueue<T, EEnqueue, R>;
attempts: number;
enqueuedAt: number;
priority: Priority;
}
ctx) =>
import EffectEffect.const log: (
...message: ReadonlyArray<any>
) => Effect.Effect<void>
Logs one or more messages using the default log level.
Example (Logging at the default level)
import { Effect } from "effect"
const program = Effect.gen(function*() {
yield* Effect.log("Starting computation")
const result = 2 + 2
yield* Effect.log("Result:", result)
yield* Effect.log("Multiple", "values", "can", "be", "logged")
return result
})
Effect.runPromise(program).then(console.log)
// Output:
// timestamp=2023-... level=INFO message="Starting computation"
// timestamp=2023-... level=INFO message="Result: 4"
// timestamp=2023-... level=INFO message="Multiple values can be logged"
// 4
log(`send ${job: any(parameter) job: {
to: string;
subject: string;
}
job.to} (attempt ${ctx: any(parameter) ctx: {
add: QueueEnqueue<T, EEnqueue, R>;
prioritize: QueueEnqueue<T, EEnqueue, R>;
defer: QueueEnqueue<T, EEnqueue, R>;
attempts: number;
enqueuedAt: number;
priority: Priority;
}
ctx.attempts}, ${ctx: any(parameter) ctx: {
add: QueueEnqueue<T, EEnqueue, R>;
prioritize: QueueEnqueue<T, EEnqueue, R>;
defer: QueueEnqueue<T, EEnqueue, R>;
attempts: number;
enqueuedAt: number;
priority: Priority;
}
ctx.priority})`),
concurrency: numberconcurrency: 4, // worker pool size — up to 4 items in flight
attempts: numberattempts: 3, // 1 try + 2 retries, then dead-lettered
key: (job: any) => anykey: (job: any(parameter) job: {
to: string;
subject: string;
}
job) => job: any(parameter) job: {
to: string;
subject: string;
}
job.to, // de-dup: same key is skipped while one is in flight
})concurrency— how many items drain at once. Default is 1 (strictly sequential); raise it for I/O-bound work.attempts— total tries per item. On the last failure the item is dead-lettered rather than retried.key— a de-duplication key. While an item with a given key is in flight, another with the same key is skipped — handy for refresh user 7 style work.
The worker effect may require services (R); the layer captures that context at build time and provides it to every run, so the resulting service needs nothing beyond what the layer itself requires.
When work fails
A queues worker either succeeds, or it fails in a way the queue declares. The tags error schema is that declaration — and its enforced: if you dont declare an error, the workers typed error channel is never, so a worker that Effect.fails wont even compile. Its failures must become defects (orDie), which the queue still catches — they just arent a typed part of the contract.
Declare an error schema and the worker may fail with it; that typed failure then rides the Failed events cause:
const const EmailJob: Schema.Struct<{
readonly to: Schema.String
readonly subject: Schema.String
}>
const EmailJob: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly to: Schema.String; readonly subject: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type…;
rebuild: (ast: Objects) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Typ…;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String …;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailJob = import SchemaSchema.function Struct<{
readonly to: Schema.String;
readonly subject: Schema.String;
}>(fields: {
readonly to: Schema.String;
readonly subject: Schema.String;
}): Schema.Struct<{
readonly to: Schema.String;
readonly subject: Schema.String;
}>
Defines a struct schema from a map of field schemas.
Details
Each field value is a schema. Use
optionalKey
or
optional
to
mark fields as optional, and
mutableKey
to mark them as mutable.
The resulting schema's Type is a readonly object type with the fields'
decoded types. The Encoded form mirrors the field schemas' encoded types.
Example (Defining a basic struct)
import { Schema } from "effect"
const Person = Schema.Struct({
name: Schema.String,
age: Schema.Number,
email: Schema.optionalKey(Schema.String)
})
// { readonly name: string; readonly age: number; readonly email?: string }
type Person = typeof Person.Type
const alice = Schema.decodeUnknownSync(Person)({ name: "Alice", age: 30 })
console.log(alice)
// { name: 'Alice', age: 30 }
Struct({ to: Schema.String(property) to: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
to: import SchemaSchema.const String: Schema.Stringconst String: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Type-level representation of
String
.
Schema for string values. Validates that the input is typeof "string".
String, subject: Schema.String(property) subject: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
subject: import SchemaSchema.const String: Schema.Stringconst String: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Type-level representation of
String
.
Schema for string values. Validates that the input is typeof "string".
String })
// declare the failure type on the tag…
class class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
}
Emails extends import QueueHyperlinkQueueHyperlink.Tag<class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
}
Emails>()("app/Emails", {
payload: Schema.Struct<{
readonly to: Schema.String
readonly subject: Schema.String
}>
(property) payload: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly to: Schema.String; readonly subject: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type…;
rebuild: (ast: Objects) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Typ…;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String …;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
payload: const EmailJob: Schema.Struct<{
readonly to: Schema.String
readonly subject: Schema.String
}>
const EmailJob: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly to: Schema.String; readonly subject: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type…;
rebuild: (ast: Objects) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Typ…;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String …;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailJob,
error: Schema.String(property) error: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
error: import SchemaSchema.const String: Schema.Stringconst String: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Type-level representation of
String
.
Schema for string values. Validates that the input is typeof "string".
String, // the worker may fail with a string
}) {}
// …and now the worker is allowed to fail with it
const const EmailsLive: anyconst EmailsLive: {
build: (memoMap: MemoMap, scope: Scope) => Effect.Effect<Context<Emails | Storage | Local<Emails>>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailsLive = import QueueHyperlinkQueueHyperlink.layer(class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails, {
effect: (
job: any
) =>
| Effect.Effect<void, never, never>
| Effect.Effect<never, string, never>
effect: (job: any(parameter) job: {
to: string;
subject: string;
}
job) =>
job: any(parameter) job: {
to: string;
subject: string;
}
job.to.includes("@")
? import EffectEffect.const void: Effect.Effect<void, never, never>
export void
(alias) const void: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Returns an effect that succeeds with void.
void
: import EffectEffect.const fail: <string>(
error: string
) => Effect.Effect<never, string, never>
Creates an Effect that represents a recoverable error.
When to use
Use to explicitly signal a recoverable error in an Effect.
Details
The error keeps propagating unless it is handled. You can handle tagged
errors with functions like
catchTag
or
catchTags
.
Example (Creating a failed effect)
import { Data, Effect } from "effect"
class OperationFailedError extends Data.TaggedError("OperationFailedError")<{}> {}
// ┌─── Effect<never, OperationFailedError, never>
// ▼
const failure = Effect.fail(
new OperationFailedError()
)
fail(`invalid address: ${job: any(parameter) job: {
to: string;
subject: string;
}
job.to}`),
attempts: numberattempts: 1,
})The handle now types as QueueHyperlink<…, void, string, never> — the string error is visible to anyone watching events. The rule is deliberate: the tag is the error contract, and workers conform to it. A queues declared failures are part of its public shape, not an implementation detail.
Enqueueing
Four verbs put work in. Three are priority lanes:
yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.add({ to: stringto: "[email protected]", subject: stringsubject: "Welcome" }) // normal
yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.prioritize({ to: stringto: "[email protected]", subject: stringsubject: "Reset code" }) // jumps the line
yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.defer({ to: stringto: "[email protected]", subject: stringsubject: "Newsletter" }) // sinks to the back
})add, prioritize, and defer each also accept an array — one call enqueues a batch, which matters over RPC (one round trip, not N). The fourth verb, enqueue, re-injects existing entries (from a release, below) with their attempt counts preserved.
Observing
Three kinds of read, for three questions.
How much is waiting, right now? — size, isEmpty, and status are reactive Subscribables. .get reads the current value once; .changes is a live stream you can render:
const const pending: anypending = yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.size.get // a number, right now
const const empty: anyempty = yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.isEmpty.get // boolean
yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.size.changes.pipe(import StreamStream.const runForEach: <number, void, never, never>(f: (a: number) => Effect.Effect<void, never, never>) => <E, R>(self: Stream.Stream<number, E, R>) => Effect.Effect<void, E, R> (+1 overload)Runs the provided effectful callback for each element of the stream.
Example (Running an effect for each value)
import { Console, Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3)
const program = Effect.gen(function*() {
yield* Stream.runForEach(stream, (n) => Console.log(`Processing: ${n}`))
})
Effect.runPromise(program)
// Processing: 1
// Processing: 2
// Processing: 3
runForEach(const onDepth: (
n: number
) => Effect.Effect<void>
onDepth)) // live, every change
})What just happened? — events is a stream of discrete facts: Enqueued, Started, Completed, Failed, RetryScheduled, RetryExhausted, and more. Subscribe once, off-fiber, and dispatch by tag:
yield* import EffectEffect.const forkScoped: <any>(
effectOrOptions?: any,
options?:
| {
readonly startImmediately?:
| boolean
| undefined
readonly uninterruptible?:
| boolean
| "inherit"
| undefined
}
| undefined
) => Effect.Effect<
Fiber<unknown, unknown>,
never,
unknown
>
Forks the fiber in a Scope, interrupting it when the scope is closed.
Example (Forking into the current scope)
import { Effect } from "effect"
const backgroundTask = Effect.gen(function*() {
yield* Effect.sleep("5 seconds")
yield* Effect.log("Background task completed")
return "result"
})
const program = Effect.scoped(
Effect.gen(function*() {
const fiber = yield* backgroundTask.pipe(Effect.forkScoped)
// or fork a fiber that starts immediately:
yield* backgroundTask.pipe(Effect.forkScoped({ startImmediately: true }))
yield* Effect.log("Task forked in scope")
yield* Effect.sleep("1 second")
// Fiber will be interrupted when scope closes
return "scope completed"
})
)
forkScoped(
const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.events.pipe(
import HyperlinkHyperlink.runForEachTag({
type Completed: (
e: any
) => Effect.Effect<void, never, never>
Completed: (e: any(parameter) e: {
_tag: "Completed";
entry: QueueEntry<T>;
success: A;
elapsed: Duration.Duration;
}
e) => import EffectEffect.const log: (
...message: ReadonlyArray<any>
) => Effect.Effect<void>
Logs one or more messages using the default log level.
Example (Logging at the default level)
import { Effect } from "effect"
const program = Effect.gen(function*() {
yield* Effect.log("Starting computation")
const result = 2 + 2
yield* Effect.log("Result:", result)
yield* Effect.log("Multiple", "values", "can", "be", "logged")
return result
})
Effect.runPromise(program).then(console.log)
// Output:
// timestamp=2023-... level=INFO message="Starting computation"
// timestamp=2023-... level=INFO message="Result: 4"
// timestamp=2023-... level=INFO message="Multiple values can be logged"
// 4
log(`sent → ${e: any(parameter) e: {
_tag: "Completed";
entry: QueueEntry<T>;
success: A;
elapsed: Duration.Duration;
}
e.entry.item.to}`),
type RetryExhausted: (
e: any
) => Effect.Effect<void, never, never>
RetryExhausted: (e: any(parameter) e: {
_tag: "RetryExhausted";
entry: QueueEntry<T>;
cause: Cause.Cause<E>;
}
e) =>
import EffectEffect.const logError: (
...message: ReadonlyArray<any>
) => Effect.Effect<void>
Logs one or more messages at the ERROR level.
Example (Logging errors)
import { Effect } from "effect"
const program = Effect.gen(function*() {
yield* Effect.logError("Database connection failed")
yield* Effect.logError(
"Error code:",
500,
"Message:",
"Internal server error"
)
// Can be used with error objects
const error = new Error("Something went wrong")
yield* Effect.logError("Caught error:", error.message)
})
Effect.runPromise(program)
// Output:
// timestamp=2023-... level=ERROR message="Database connection failed"
// timestamp=2023-... level=ERROR message="Error code: 500 Message: Internal server error"
// timestamp=2023-... level=ERROR message="Caught error: Something went wrong"
logError(`dead-letter ${e: any(parameter) e: {
_tag: "RetryExhausted";
entry: QueueEntry<T>;
cause: Cause.Cause<E>;
}
e.entry.item.to}: ${import CauseCause.const pretty: <unknown>(
cause: Cause.Cause<unknown>
) => string
Formats a Cause as a human-readable string for logging or debugging.
When to use
Use to render a whole cause as one human-readable string for logs or
diagnostics.
Details
Delegates to
prettyErrors
to convert each reason to an Error,
then joins their stack traces with newlines. Nested Error.cause chains
are rendered inline with indentation:
ErrorName: message
at ...
at ... {
[cause]: NestedError: message
at ...
}
Span annotations are appended to the relevant stack frames when available.
Gotchas
Rendering an empty cause produces an empty string because there are no
errors to render.
Example (Rendering a cause)
import { Cause } from "effect"
const rendered = Cause.pretty(Cause.fail("something went wrong"))
console.log(rendered.includes("something went wrong")) // true
pretty(e: any(parameter) e: {
_tag: "RetryExhausted";
entry: QueueEntry<T>;
cause: Cause.Cause<E>;
}
e.cause)}`),
}),
),
)
})How is it trending? — metrics.stream emits a windowed aggregate (throughput, average wait, per-window counts) once per window; metrics.query reads historical windows back from the store.
Effect queues cant enumerate their pending items — theres no list. You read counts (size, per-priority sizes on status) and facts (events), and you target what you already know by key (drop, deadLetter). Thats a feature: it keeps the queue O(1) to observe at any depth.
Controlling
The same handle steers the queue. pause stops draining (items still enqueue and accumulate); resume starts again; shutdown drains gracefully and stops; clear empties the pending items and returns how many it cleared; start forks the worker pool (idempotent — layers do this for you).
Three verbs route work out of the queue: release exports pending entries and removes them (hand them to another runtime, then enqueue them there); deadLetter removes entries matching a selector and records them as dead-lettered; drop removes them without a trace. You target these by what you know — an entry id or a matching item — never by listing.
Success values
If the worker returns a value, declare a success schema and that value flows onto the Completed event and the stores analytics:
const const Job: Schema.Struct<{
readonly id: Schema.String
}>
const Job: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly id: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly id: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly id: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly id: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>>) => Schema.Struct<{ readonly id: Schema.String }>;
rebuild: (ast: Objects) => Schema.Struct<{ readonly id: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Job = import SchemaSchema.function Struct<{
readonly id: Schema.String;
}>(fields: {
readonly id: Schema.String;
}): Schema.Struct<{
readonly id: Schema.String;
}>
Defines a struct schema from a map of field schemas.
Details
Each field value is a schema. Use
optionalKey
or
optional
to
mark fields as optional, and
mutableKey
to mark them as mutable.
The resulting schema's Type is a readonly object type with the fields'
decoded types. The Encoded form mirrors the field schemas' encoded types.
Example (Defining a basic struct)
import { Schema } from "effect"
const Person = Schema.Struct({
name: Schema.String,
age: Schema.Number,
email: Schema.optionalKey(Schema.String)
})
// { readonly name: string; readonly age: number; readonly email?: string }
type Person = typeof Person.Type
const alice = Schema.decodeUnknownSync(Person)({ name: "Alice", age: 30 })
console.log(alice)
// { name: 'Alice', age: 30 }
Struct({ id: Schema.String(property) id: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
id: import SchemaSchema.const String: Schema.Stringconst String: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<string, readonly []>) => Schema.String;
annotateKey: (annotations: Schema.Annotations.Key<string>) => Schema.String;
check: (checks_0: Check<string>, ...checks: Array<Check<string>>) => Schema.String;
rebuild: (ast: String) => Schema.String;
make: (input: string, options?: Schema.MakeOptions) => string;
makeOption: (input: string, options?: Schema.MakeOptions) => Option<string>;
makeEffect: (input: string, options?: Schema.MakeOptions) => Effect.Effect<string, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Type-level representation of
String
.
Schema for string values. Validates that the input is typeof "string".
String })
class class Doublerclass Doubler {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
}
Doubler extends import QueueHyperlinkQueueHyperlink.Tag<class Doublerclass Doubler {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
}
Doubler>()("app/Doubler", {
payload: Schema.Struct<{
readonly id: Schema.String
}>
(property) payload: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly id: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly id: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly id: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly id: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>>) => Schema.Struct<{ readonly id: Schema.String }>;
rebuild: (ast: Objects) => Schema.Struct<{ readonly id: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
payload: const Job: Schema.Struct<{
readonly id: Schema.String
}>
const Job: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly id: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly id: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly id: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly id: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>>) => Schema.Struct<{ readonly id: Schema.String }>;
rebuild: (ast: Objects) => Schema.Struct<{ readonly id: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>>;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly id: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly id: Schema.String }, 'Type'>, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Job,
success: Schema.Number(property) success: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<number, readonly []>) => Schema.Number;
annotateKey: (annotations: Schema.Annotations.Key<number>) => Schema.Number;
check: (checks_0: Check<number>, ...checks: Array<Check<number>>) => Schema.Number;
rebuild: (ast: Number) => Schema.Number;
make: (input: number, options?: Schema.MakeOptions) => number;
makeOption: (input: number, options?: Schema.MakeOptions) => Option<number>;
makeEffect: (input: number, options?: Schema.MakeOptions) => Effect.Effect<number, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
success: import SchemaSchema.const Number: Schema.Numberconst Number: {
Rebuild: Rebuild;
Iso: Iso;
ast: Ast;
Type: T;
Encoded: E;
DecodingServices: RD;
EncodingServices: RE;
annotate: (annotations: Schema.Annotations.Bottom<number, readonly []>) => Schema.Number;
annotateKey: (annotations: Schema.Annotations.Key<number>) => Schema.Number;
check: (checks_0: Check<number>, ...checks: Array<Check<number>>) => Schema.Number;
rebuild: (ast: Number) => Schema.Number;
make: (input: number, options?: Schema.MakeOptions) => number;
makeOption: (input: number, options?: Schema.MakeOptions) => Option<number>;
makeEffect: (input: number, options?: Schema.MakeOptions) => Effect.Effect<number, Schema.SchemaError, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Type-level representation of
Number
.
Schema for number values, including NaN, Infinity, and -Infinity.
Details
Default JSON serializer:
- Finite numbers are serialized as numbers.
- Non-finite values are serialized as strings (
"NaN", "Infinity", "-Infinity").
Number, // the worker returns a number
}) {}
const const DoublerLive: anyconst DoublerLive: {
build: (memoMap: MemoMap, scope: Scope) => Effect.Effect<Context<Doubler | Storage | Local<Doubler>>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
DoublerLive = import QueueHyperlinkQueueHyperlink.layer(class Doublerclass Doubler {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ id: string }, number, never, never>) => QueueHyperlink.QueueHyperlink<{ id: string }, number, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ id: string }, number, never, never>) => Context<Doubler>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ id: string }, number, never, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Doubler | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ id: string }, number, never, never>) => A) => Effect.Effect<A, never, Doubler>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Doubler, {
effect: (
job: any
) => Effect.Effect<number, never, never>
effect: (job: any(parameter) job: {
id: string;
}
job) => import EffectEffect.const succeed: <number>(
value: number
) => Effect.Effect<number, never, never>
Creates an Effect that always succeeds with a given value.
When to use
Use when an effect should complete successfully with a specific value without any errors
or external dependencies.
Example (Creating a successful effect)
import { Effect } from "effect"
// Creating an effect that represents a successful scenario
//
// ┌─── Effect<number, never, never>
// ▼
const success = Effect.succeed(42)
succeed(job: any(parameter) job: {
id: string;
}
job.id.length * 2),
})The handle types as QueueHyperlink<{ id: string }, number, never, never>, and Completed.success carries the number.
The .Service shorthand
Tag + layer keeps the contract and the worker separate — which is what makes a queue location-transparent (the same tag, a different layer, and it runs remotely). When you dont need that split, QueueHyperlink.Service fuses both into one class — a self-contained Service:
class class Emailsclass Emails {
Service: Service;
key: Identifier;
}
Emails extends import QueueHyperlinkQueueHyperlink.Service<class Emailsclass Emails {
Service: Service;
key: Identifier;
}
Emails, typeof const EmailJob: Schema.Struct<{
readonly to: Schema.String
readonly subject: Schema.String
}>
const EmailJob: {
Type: Struct.Type<Fields>;
Encoded: Struct.Encoded<Fields>;
DecodingServices: Struct.DecodingServices<Fields>;
EncodingServices: Struct.EncodingServices<Fields>;
Iso: Struct.Iso<Fields>;
fields: Fields;
mapFields: (f: (fields: { readonly to: Schema.String; readonly subject: Schema.String }) => To, options?: { readonly unsafePreserveChecks?: boolean | undefined } | undefined) => Schema.Struct<{ [K in keyof Readonly<To>]: Readonly<To>[K]; }>;
Rebuild: Rebuild;
ast: Ast;
annotate: (annotations: Schema.Annotations.Bottom<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String; }, 'Type'>, readonly []>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String; }>;
annotateKey: (annotations: Schema.Annotations.Key<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
check: (checks_0: Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>>, ...checks: Array<Check<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type…;
rebuild: (ast: Objects) => Schema.Struct<{ readonly to: Schema.String; readonly subject: Schema.String }>;
make: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Type'>;
makeOption: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Option<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String }, 'Typ…;
makeEffect: (input: Schema.Struct.ReadonlyMakeIn<{ readonly to: Schema.String; readonly subject: Schema.String }>, options?: Schema.MakeOptions) => Effect.Effect<Schema.Struct.ReadonlySide<{ readonly to: Schema.String; readonly subject: Schema.String …;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailJob.Struct<{ readonly to: String; readonly subject: String; }>["Type"]: Schema.Struct.ReadonlySide<{
readonly to: Schema.String;
readonly subject: Schema.String;
}, "Type">
(property) Struct<{ readonly to: String; readonly subject: String; }>["Type"]: {
to: string;
subject: string;
}
Type, never>()(
"app/Emails",
{
concurrency: numberconcurrency: 4,
effect: (
job: any
) => Effect.Effect<void, never, never>
effect: (job: any(parameter) job: {
to: string;
subject: string;
}
job) => import EffectEffect.const log: (
...message: ReadonlyArray<any>
) => Effect.Effect<void>
Logs one or more messages using the default log level.
Example (Logging at the default level)
import { Effect } from "effect"
const program = Effect.gen(function*() {
yield* Effect.log("Starting computation")
const result = 2 + 2
yield* Effect.log("Result:", result)
yield* Effect.log("Multiple", "values", "can", "be", "logged")
return result
})
Effect.runPromise(program).then(console.log)
// Output:
// timestamp=2023-... level=INFO message="Starting computation"
// timestamp=2023-... level=INFO message="Result: 4"
// timestamp=2023-... level=INFO message="Multiple values can be logged"
// 4
log(`send ${job: any(parameter) job: {
to: string;
subject: string;
}
job.to}`),
},
) {}yield* Emails yields the exact same handle type. Reach for Service for a self-contained local queue; reach for Tag + layer when the queue might move.
Everything so far is the basic queue. The rest of this guide is the operating surface — the controls you reach for once a queue is real.
Handling failure
attempts is the blunt instrument: try N times, then dead-letter. Real failure handling is per-error, and thats what onFailure is for. It runs when an attempt fails, receives the entry and the Cause, and decides what happens next — retry, dead-letter, or drop:
const const EmailsLive: anyconst EmailsLive: {
build: (memoMap: MemoMap, scope: Scope) => Effect.Effect<Context<Storage | Emails | Local<Emails>>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailsLive = import QueueHyperlinkQueueHyperlink.layer(class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, string, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails, {
effect: (
job: any
) => Effect.Effect<void, string, never>
effect: (job: any(parameter) job: {
to: string;
subject: string;
}
job) => const send: (job: {
to: string
}) => Effect.Effect<void, string>
send(job: any(parameter) job: {
to: string;
subject: string;
}
job),
attempts: numberattempts: 5,
onFailure: (
entry: any,
cause: any
) =>
| Effect.Effect<"retry", never, never>
| Effect.Effect<"dead-letter", never, never>
onFailure: (entry: any(parameter) entry: {
item: T;
entryId: string;
key: string;
priority: Priority;
level: number;
attempts: number;
timestamps: QueueEntryTimestamps;
batchId: string;
releaseId: string;
sourceHyperlinkId: string;
attributes: { readonly [key: string]: unknown };
}
entry, cause: any(parameter) cause: {
reasons: ReadonlyArray<Reason<E>>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
cause) =>
const isTransient: (
cause: Cause.Cause<string>
) => boolean
isTransient(cause: any(parameter) cause: {
reasons: ReadonlyArray<Reason<E>>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
cause)
? import EffectEffect.const succeed: <"retry">(value: "retry") => Effect.Effect<"retry", never, never>Creates an Effect that always succeeds with a given value.
When to use
Use when an effect should complete successfully with a specific value without any errors
or external dependencies.
Example (Creating a successful effect)
import { Effect } from "effect"
// Creating an effect that represents a successful scenario
//
// ┌─── Effect<number, never, never>
// ▼
const success = Effect.succeed(42)
succeed("retry" as type const = "retry"const) // a blip — spend an attempt
: import EffectEffect.const succeed: <"dead-letter">(value: "dead-letter") => Effect.Effect<"dead-letter", never, never>Creates an Effect that always succeeds with a given value.
When to use
Use when an effect should complete successfully with a specific value without any errors
or external dependencies.
Example (Creating a successful effect)
import { Effect } from "effect"
// Creating an effect that represents a successful scenario
//
// ┌─── Effect<number, never, never>
// ▼
const success = Effect.succeed(42)
succeed("dead-letter" as type const = "dead-letter"const), // a bad address — set it aside
})Three dispositions: "retry" re-enqueues (until attempts runs out), "dead-letter" sets the entry aside as failed (a DeadLettered event), and "drop" discards it silently. Without onFailure, the default is retry until attempts, then dead-letter. For retrying the effect itself (backoff, jitter), put Effect.retry on your worker effect — thats a different layer of the onion: onFailure decides the entrys fate after the effect has given up.
Rate limiting the drain
A queue that hammers a downstream API needs a ceiling. rateLimit caps how many items start per window; excess wait, and a RateLimitExceeded event fires when the ceiling bites:
const const EmailsLive: anyconst EmailsLive: {
build: (memoMap: MemoMap, scope: Scope) => Effect.Effect<Context<Storage | Emails | Local<Emails>>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailsLive = import QueueHyperlinkQueueHyperlink.layer(class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails, {
effect: (
job: any
) => Effect.Effect<void, never, never>
effect: (job: any(parameter) job: {
to: string;
subject: string;
}
job) => import EffectEffect.const log: (
...message: ReadonlyArray<any>
) => Effect.Effect<void>
Logs one or more messages using the default log level.
Example (Logging at the default level)
import { Effect } from "effect"
const program = Effect.gen(function*() {
yield* Effect.log("Starting computation")
const result = 2 + 2
yield* Effect.log("Result:", result)
yield* Effect.log("Multiple", "values", "can", "be", "logged")
return result
})
Effect.runPromise(program).then(console.log)
// Output:
// timestamp=2023-... level=INFO message="Starting computation"
// timestamp=2023-... level=INFO message="Result: 4"
// timestamp=2023-... level=INFO message="Multiple values can be logged"
// 4
log(`send ${job: any(parameter) job: {
to: string;
subject: string;
}
job.to}`),
concurrency: numberconcurrency: 8,
rateLimit: {
limit: number
window: Duration.Duration
}
rateLimit: { limit: numberlimit: 100, window: Duration.Duration(property) window: {
value: DurationValue;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
window: import DurationDuration.const seconds: (
seconds: number
) => Duration.Duration
Creates a Duration from seconds.
Example (Creating durations from seconds)
import { Duration } from "effect"
const duration = Duration.seconds(30)
console.log(Duration.toMillis(duration)) // 30000
seconds(1) }, // ≤ 100 starts/sec
})concurrency and rateLimit are orthogonal: concurrency bounds in-flight work, rate limit bounds start rate. Use both — a pool of 8 workers that collectively start no faster than 100/sec.
Bootstrapping: start paused
Sometimes you want to load a queue before it drains — seed a backlog, wire up a subscriber, then let it rip. Start it paused and resume when ready:
yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.add({ to: stringto: "[email protected]", subject: stringsubject: "queued while paused" })
yield* const emails: anyconst emails: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
}
emails.resume // now it drains
})Pulling work in
The queues so far are push — something calls add. A queue can also pull, with refill: a loader that the engine calls to fetch work. onStart seeds it once on boot; onDrained re-polls the source every time the queue empties — turning a queue into a durable poller over an external source (a table, a topic, an inbox):
const const EmailsLive: anyconst EmailsLive: {
build: (memoMap: MemoMap, scope: Scope) => Effect.Effect<Context<Storage | Emails | Local<Emails>>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailsLive = import QueueHyperlinkQueueHyperlink.layer(class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Effect.Effect<A, E, R>) => Effect.Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => A) => Effect.Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails, {
effect: (
job: any
) => Effect.Effect<void, never, never>
effect: (job: any(parameter) job: {
to: string;
subject: string;
}
job) => import EffectEffect.const log: (
...message: ReadonlyArray<any>
) => Effect.Effect<void>
Logs one or more messages using the default log level.
Example (Logging at the default level)
import { Effect } from "effect"
const program = Effect.gen(function*() {
yield* Effect.log("Starting computation")
const result = 2 + 2
yield* Effect.log("Result:", result)
yield* Effect.log("Multiple", "values", "can", "be", "logged")
return result
})
Effect.runPromise(program).then(console.log)
// Output:
// timestamp=2023-... level=INFO message="Starting computation"
// timestamp=2023-... level=INFO message="Result: 4"
// timestamp=2023-... level=INFO message="Multiple values can be logged"
// 4
log(job: any(parameter) job: {
to: string;
subject: string;
}
job.to),
refill: {
onStart: boolean
onDrained: boolean
load: (
queue: any
) => Effect.Effect<unknown, unknown, unknown>
}
refill: {
onStart: booleanonStart: true, // seed on boot
onDrained: booleanonDrained: true, // re-poll when empty
load: (
queue: any
) => Effect.Effect<unknown, unknown, unknown>
load: (queue: anyqueue) => import EffectEffect.const flatMap: <readonly {
to: string;
subject: string;
}[], never, never, unknown, unknown, unknown>(self: Effect.Effect<readonly {
to: string;
subject: string;
}[], never, never>, f: (a: readonly {
to: string;
subject: string;
}[]) => Effect.Effect<unknown, unknown, unknown>) => Effect.Effect<unknown, unknown, unknown> (+1 overload)
Chains effects to produce new Effect instances, useful for combining
operations that depend on previous results.
When to use
Use when you need to chain multiple effects, ensuring that each
step produces a new Effect while flattening any nested effects that may
occur.
Details
flatMap lets you sequence effects so that the result of one effect can be
used in the next step. It is similar to flatMap used with arrays but works
specifically with Effect instances, allowing you to avoid deeply nested
effect structures.
Since effects are immutable, flatMap always returns a new effect instead of
changing the original one.
Example (Choosing flatMap syntax variants)
import { Effect, pipe } from "effect"
const myEffect = Effect.succeed(1)
const transformation = (n: number) => Effect.succeed(n + 1)
const flatMappedWithPipe = pipe(myEffect, Effect.flatMap(transformation))
const flatMappedWithDataFirst = Effect.flatMap(myEffect, transformation)
const flatMappedWithMethod = myEffect.pipe(Effect.flatMap(transformation))
Example (Sequencing dependent effects)
import { Data, Effect, pipe } from "effect"
class DiscountRateError extends Data.TaggedError("DiscountRateError")<{}> {}
// Function to apply a discount safely to a transaction amount
const applyDiscount = (
total: number,
discountRate: number
): Effect.Effect<number, DiscountRateError> =>
discountRate === 0
? Effect.fail(new DiscountRateError())
: Effect.succeed(total - (total * discountRate) / 100)
// Simulated asynchronous task to fetch a transaction amount from database
const fetchTransactionAmount = Effect.promise(() => Promise.resolve(100))
// Chaining the fetch and discount application using `flatMap`
const finalAmount = pipe(
fetchTransactionAmount,
Effect.flatMap((amount) => applyDiscount(amount, 5))
)
Effect.runPromise(finalAmount).then(console.log)
// Output: 95
flatMap(const nextBatch: Effect.Effect<
readonly {
to: string
subject: string
}[],
never,
never
>
const nextBatch: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
nextBatch, queue: anyqueue.add),
},
})The loader gets the queue handle, so it enqueues with the same verbs you do.
Operating a live queue
Three streams tell you what a running queue is doing, from three angles.
events is the fact log. Every discrete thing the queue does is a tagged event, and they fall into four families:
lifecycle —
Enqueued,Started,Completedfailure —
Failed,RetryScheduled,RetryExhaustedrouting —
Released,DeadLettered,Dropped,Clearedqueue-level —
Start,RateLimitExceeded,ShutdownRequested,ShutdownComplete,Drained
You never handle all of them — pick the tags you care about with Hyperlink.runForEachTag and ignore the rest.
status is the current-state snapshot (a Subscribable): per-priority pending sizes, how many are inFlight, the running completed count, whether its paused, and its phase — running, then draining after a shutdown request, then off. Its the one value a dashboard renders.
metrics is the aggregate view: metrics.stream emits one windowed summary per window (throughput, average wait and execution time, per-window counts); metrics.query reads past windows back from the store for charts and trends.
status answers what is true now, events answers what happened, metrics answers how is it trending. Reach for the one that matches the question — they dont overlap.
Persistence and analytics
Every queue comes with an observability store already wired in. By default its in-memory: lifecycle events and metric windows are recorded, and metrics.query reads them back — for the life of the process. Provide a durable store instead (the toolkit ships a SQLite-backed one) and that history survives restarts: a dashboard reconnecting after a redeploy still sees yesterdays throughput.
The store is also an analytics surface in its own right — beyond metrics.query it answers questions like the slowest completions and how many have completed, computed over the recorded history rather than the live queue. You reach it with QueueHyperlink.store(tag).
Running it across the network
This is the payoff of the tag/layer split. The tag is the contract; the layer decides where the work runs — and nothing else in your code changes.
Provide QueueHyperlink.layer and the queue is local. Provide QueueHyperlink.serve instead and the worker runs behind an RPC server, its handlers mounted for callers. A different process then provides Hyperlink.client(Tag) (or Hyperlink.connect(Tag, Hyperlink.protocolHttp(port)) over HTTP), and the same yield* Tag code drives the remote queue — add, size, events, pause, all of it — as if it were in-process. The handles Requirements param is the only tell: never locally, the transport for a client.
For moving pending work between runtimes, release exports entries decoded and releaseEncoded exports them in wire form (no item schema needed on the receiver); the other side enqueues them, attempt budgets intact.
Reconfiguring at runtime
A queues concurrency, rateLimit, or paused state isnt frozen at definition. QueueHyperlink.configure(Tag, patch) is a layer that overlays a config patch on top of the base — Layer.provideMerge it, and the queue drains under the merged config:
const const Tuned: Layer.Layer<any, any, any>const Tuned: {
build: (memoMap: Layer.MemoMap, scope: Scope) => Effect<Context<Emails>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
Tuned = const EmailsLive: Layer.Layer<
Emails,
never,
never
>
const EmailsLive: {
build: (memoMap: Layer.MemoMap, scope: Scope) => Effect<Context<Emails>, never, never>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
EmailsLive.Pipeable.pipe<Layer.Layer<Emails, never, never>, Layer.Layer<any, any, any>>(this: Layer.Layer<Emails, never, never>, ab: (_: Layer.Layer<Emails, never, never>) => Layer.Layer<any, any, any>): Layer.Layer<any, any, any> (+21 overloads)pipe(
import LayerLayer.const provideMerge: <any>(that: any) => <A, E, R>(self: Layer.Layer<A, E, R>) => Layer.Layer<any, any, any> (+3 overloads)Feeds the output services of the dependency layer into the requirements of
this layer, returning a layer that provides both sets of services.
When to use
Use when you need to compose Layers while keeping both the constructed
service and the dependency used to build it available.
Details
Prefer
provide
when the dependency should stay private.
Example (Providing dependencies while retaining services)
import { Context, Effect, Layer } from "effect"
class Database extends Context.Service<Database, {
readonly query: (sql: string) => Effect.Effect<string>
}>()("Database") {}
class Logger extends Context.Service<Logger, {
readonly log: (msg: string) => Effect.Effect<void>
}>()("Logger") {}
class UserService extends Context.Service<UserService, {
readonly getUser: (id: string) => Effect.Effect<{
id: string
name: string
}>
}>()("UserService") {}
// Create dependency layers
const databaseLayer = Layer.succeed(Database, {
query: Effect.fn("Database.query")((sql: string) => Effect.succeed(`DB: ${sql}`))
})
const loggerLayer = Layer.succeed(Logger, {
log: Effect.fn("Logger.log")((msg: string) => Effect.sync(() => console.log(`[LOG] ${msg}`)))
})
// UserService depends on Database and Logger
const userServiceLayer = Layer.effect(UserService, Effect.gen(function*() {
const database = yield* Database
const logger = yield* Logger
return {
getUser: Effect.fn("UserService.getUser")(function*(id: string) {
yield* logger.log(`Looking up user ${id}`)
const result = yield* database.query(
`SELECT * FROM users WHERE id = ${id}`
)
return { id, name: result }
})
}
}))
// Provide dependencies and merge all services together
const allServicesLayer = userServiceLayer.pipe(
Layer.provideMerge(Layer.mergeAll(databaseLayer, loggerLayer))
)
// Now the resulting layer provides UserService, Database, AND Logger
const program = Effect.gen(function*() {
const userService = yield* UserService
const logger = yield* Logger // Still available!
const database = yield* Database // Still available!
const user = yield* userService.getUser("123")
yield* logger.log(`Found user: ${user.name}`)
return user
}).pipe(
Effect.provide(allServicesLayer)
)
provideMerge(import QueueHyperlinkQueueHyperlink.configure(class Emailsclass Emails {
key: Identifier;
Service: {
status: Hyperlink.Subscribable<QueueStatus>;
size: Hyperlink.Subscribable<number>;
isEmpty: Hyperlink.Subscribable<boolean>;
start: Effect.Effect<void, never, Requirements>;
pause: Effect.Effect<void>;
resume: Effect.Effect<void>;
shutdown: Effect.Effect<void>;
clear: Effect.Effect<number, never, Requirements>;
metrics: { readonly stream: Stream.Stream<QueueMetrics>; readonly query: (input: { readonly limit?: number; readonly since?: DateTime.Utc; readonly until?: DateTime.Utc }) => Effect.Effect<ReadonlyArray<QueueMetrics>, never, Requirements> };
add: QueueEnqueue<Payload, never, Requirements>;
prioritize: QueueEnqueue<Payload, never, Requirements>;
defer: QueueEnqueue<Payload, never, Requirements>;
enqueue: (entries: ReadonlyArray<QueueEntry<Payload>>) => Effect.Effect<void, never, Requirements>;
release: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
releaseEncoded: (input: { readonly options?: QueueReleaseOptions }) => Effect.Effect<ReadonlyArray<Hyperlink.Decoded<typeof queueEncodedEntry>>, QueueReleaseEncodingError, Requirements>;
deadLetter: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
drop: (input: { readonly selector: QueueEntrySelector<Payload> | QueueEntry<Payload>; readonly options: QueueRouteOptions }) => Effect.Effect<ReadonlyArray<QueueEntry<Payload>>, never, Requirements>;
events: Stream.Stream<QueueEvent<Payload, Error, Success>>;
};
groupId: string;
description: string | undefined;
of: (this: void, self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>;
context: (self: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Context<Emails>;
use: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => Effect<A, E, R>) => Effect<A, E, Emails | R>;
useSync: (f: (service: QueueHyperlink.QueueHyperlink<{ to: string; subject: string }, void, never, never>) => A) => Effect<A, never, Emails>;
Identifier: Identifier;
stack: string | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Emails, { concurrency: numberconcurrency: 16 })),
)Because its just a layer, the patch can come from anywhere a layer can — an env flag, or a live DynamicConfig swap that re-tunes the queue while it runs.
Custom priority lanes
high / normal / defer covers most needs, but some domains have their own ordering — tiers, SLAs, numbered levels. CustomQueueHyperlink is the same queue with arbitrary lanes: you define the levels, and add targets one by name. The handle reads the same; only the priority axis is yours to shape.
A live control panel
Because the handle is the whole surface, a UI is just another consumer of it. These docs render one inline: a ``` queue block mounts a live control panel for a declared queue — buttons that call add / prioritize / pause / resume / clear, and stats read straight off the status stream. Enqueue an item and watch pending climb, then drain to completed:
The same handle that runs the queue drives the panel — no separate admin API, no extra wiring.
The raw engine
Under the resource wrapper is a plain queue engine. QueueHyperlink.make(config) returns a handle directly — the workers, retries, and events, without the Tag, Layer, or RPC machinery. layer and Service are built on it; reach for make only when you want to embed a queue inside something else and manage its scope yourself. For everything else, the tag is the queue.