MessageStorage
Stores Effect Cluster messages and replies behind a pluggable backend.
MessageStorage is the boundary between cluster runner logic and the storage system that keeps mailbox state recoverable. It saves requests, control envelopes, and replies; finds unprocessed messages for assigned shards; tracks duplicate requests; and manages reply handlers waiting for responses. This module also includes the encoded storage-driver contract and no-op or in-memory implementations for local use and tests.
Constructors
Signature
declare function make( storage: Omit< MessageStorage["Service"], "registerReplyHandler" | "unregisterReplyHandler" | "unregisterShardReplyHandlers" >,): Effect<{ readonly clearAddress: (address: EntityAddress) => Effect<void, PersistenceError>; readonly clearReplies: (requestId: Snowflake) => Effect<void, PersistenceError>; readonly registerReplyHandler: <R extends Any>( message: OutgoingRequest<R> | IncomingRequest<R>, ) => Effect<void, EntityNotAssignedToRunner>; readonly repliesFor: <R extends Any>( requests: Iterable<OutgoingRequest<R>>, ) => Effect<Array<Reply<R>>, MalformedMessage | PersistenceError>; readonly repliesForUnfiltered: ( requestIds: Iterable<Snowflake>, ) => Effect<Array<Encoded>, MalformedMessage | PersistenceError>; readonly requestIdForPrimaryKey: (options: { readonly address: EntityAddress; readonly id: string; readonly tag: string; }) => Effect<Option<Snowflake>, PersistenceError>; readonly resetAddress: (address: EntityAddress) => Effect<void, PersistenceError>; readonly resetShards: (shardIds: Iterable<ShardId>) => Effect<void, PersistenceError>; readonly saveEnvelope: ( envelope: OutgoingEnvelope, ) => Effect<void, MalformedMessage | PersistenceError>; readonly saveReply: <R extends Any>( reply: ReplyWithContext<R>, ) => Effect<void, MalformedMessage | PersistenceError>; readonly saveRequest: <R extends Any>( envelope: OutgoingRequest<R>, ) => Effect<SaveResult<R>, MalformedMessage | PersistenceError>; readonly unprocessedMessages: ( shardIds: Iterable<ShardId>, ) => Effect<Array<Incoming<any>>, PersistenceError>; readonly unprocessedMessagesById: <R extends Any>( messageIds: Iterable<Snowflake>, ) => Effect<Array<Incoming<R>>, PersistenceError>; readonly unregisterReplyHandler: (requestId: Snowflake) => Effect<void>; readonly unregisterShardReplyHandlers: (shardId: ShardId) => Effect<void>; readonly withTransaction: <A, E, R>(effect: Effect<A, E, R>) => Effect<A, E, R>;}>;makeEncoded
Builds a MessageStorage service from an encoded storage driver.
Details
The adapter handles envelope and reply encoding and decoding, primary-key generation, delayed delivery checks, duplicate decoding, and malformed-message defect replies.
Signature
declare const makeEncoded: ( encoded: Encoded,) => Effect.Effect<MessageStorage["Service"], never, Snowflake.Generator>;No-op MessageStorage service that does not persist messages or replies.
Signature
declare const noop: MessageStorage["Service"];SaveResult
Constructors and matchers for decoded save results.
Signature
declare const SaveResult: { readonly $is: <Tag extends "Success" | "Duplicate">( tag: Tag, ) => { <T extends SaveResult<any>>( u: T, ): u is T & { readonly _tag: Tag; }; (u: unknown): u is | Extract< Success, { readonly _tag: Tag; } > | Extract< Duplicate<never>, { readonly _tag: Tag; } >; }; readonly $match: { < A, B, C, D, Cases extends { Duplicate: (args: Duplicate<A extends Any ? A : never>) => any; Success: (args: Success) => any; }, >( cases: Cases, ): ( self: SaveResult<A extends Any ? A : never>, ) => Unify<ReturnType<Cases["Success" | "Duplicate"]>>; < A, B, C, D, Cases extends { Duplicate: (args: Duplicate<A extends Any ? A : never>) => any; Success: (args: Success) => any; }, >( self: SaveResult<A extends Any ? A : never>, cases: Cases, ): Unify<ReturnType<Cases["Success" | "Duplicate"]>>; }; Duplicate: <A>(args: { readonly lastReceivedReply: Option<Reply<A extends Any ? A : never>>; readonly originalId: Snowflake; }) => Duplicate<A extends Any ? A : never>; Success: <A>(args: void) => Success;};SaveResultEncoded
Constructors and matchers for encoded save results returned by storage drivers.
Signature
declare const SaveResultEncoded: { readonly $is: <Tag extends "Success" | "Duplicate">( tag: Tag, ) => (u: unknown) => u is | Extract< Success, { readonly _tag: Tag; } > | Extract< DuplicateEncoded, { readonly _tag: Tag; } >; readonly $match: { < Cases extends { Duplicate: (args: DuplicateEncoded) => any; Success: (args: Success) => any; }, >( cases: Cases, ): (value: Encoded) => Unify<ReturnType<Cases["Success" | "Duplicate"]>>; < Cases extends { Duplicate: (args: DuplicateEncoded) => any; Success: (args: Success) => any; }, >( value: Encoded, cases: Cases, ): Unify<ReturnType<Cases["Success" | "Duplicate"]>>; }; Duplicate: ConstructorFrom<DuplicateEncoded, "_tag">; Success: ConstructorFrom<Success, "_tag">;};Layers
layerMemory
Layer that provides in-memory message storage and its backing MemoryDriver.
Signature
declare const layerMemory: Layer.Layer<MessageStorage | MemoryDriver, never, ShardingConfig>;Layer that provides the no-op MessageStorage service.
Signature
declare const layerNoop: Layer.Layer<MessageStorage>;Models
EncodedRepliesOptions type
Cursor options for reading encoded replies across request sets.
Details
The fields distinguish existing requests from new requests and carry the driver-specific pagination cursor.
Signature
type EncodedRepliesOptions<A> = { readonly cursor: Option.Option<A>; readonly existingRequests: Array<string>; readonly newRequests: Array<string>;};EncodedUnprocessedOptions type
Cursor options for reading encoded unprocessed messages across shard sets.
Details
The fields distinguish existing shards from newly assigned shards and carry the driver-specific pagination cursor.
Signature
type EncodedUnprocessedOptions<A> = { readonly cursor: Option.Option<A>; readonly existingShards: Array<number>; readonly newShards: Array<number>;};MemoryEntry type
In-memory storage entry for a request envelope.
Details
It stores the encoded envelope, last acknowledged chunk, accumulated replies, and optional delivery time.
Signature
type MemoryEntry = { deliverAt: number | null; readonly envelope: Envelope.Encoded; lastReceivedChunk: Reply.ChunkEncoded | undefined; replies: Array<Reply.Encoded>;};SaveResult type
Result of saving a request or envelope into message storage.
Details
A duplicate result carries the original request ID and the last reply already received for the duplicated request.
Signature
type SaveResult<R extends Rpc.Any> = SaveResult.Success | SaveResult.Duplicate<R>;Other
SaveResult
Variants and helper types for SaveResult.
Services
Low-level storage-driver contract for encoded envelopes and replies.
Details
Implementations persist encoded messages, track primary keys and delayed delivery, read unprocessed messages, and provide transaction wrapping.
Signature
type Encoded = { readonly clearAddress: (address: EntityAddress) => Effect.Effect<void, PersistenceError>; readonly clearReplies: (requestId: Snowflake.Snowflake) => Effect.Effect<void, PersistenceError>; readonly repliesFor: ( requestIds: Arr.NonEmptyArray<string>, ) => Effect.Effect<Array<Reply.Encoded>, PersistenceError>; readonly repliesForUnfiltered: ( requestIds: Arr.NonEmptyArray<string>, ) => Effect.Effect<Array<Reply.Encoded>, PersistenceError>; readonly requestIdForPrimaryKey: ( primaryKey: string, ) => Effect.Effect<Option.Option<Snowflake.Snowflake>, PersistenceError>; readonly resetAddress: (address: EntityAddress) => Effect.Effect<void, PersistenceError>; readonly resetShards: ( shardIds: Arr.NonEmptyArray<string>, ) => Effect.Effect<void, PersistenceError>; readonly saveEnvelope: (options: { readonly deliverAt: number | null; readonly envelope: Envelope.Encoded; readonly primaryKey: string | null; }) => Effect.Effect<SaveResult.Encoded, PersistenceError>; readonly saveReply: (reply: Reply.Encoded) => Effect.Effect<void, PersistenceError>; readonly unprocessedMessages: ( shardIds: Arr.NonEmptyArray<string>, now: number, ) => Effect.Effect< Array<{ readonly envelope: Envelope.Encoded; readonly lastSentReply: Option.Option<Reply.Encoded>; }>, PersistenceError >; readonly unprocessedMessagesById: ( messageIds: Arr.NonEmptyArray<Snowflake.Snowflake>, now: number, ) => Effect.Effect< Array<{ readonly envelope: Envelope.Encoded; readonly lastSentReply: Option.Option<Reply.Encoded>; }>, PersistenceError >; readonly withTransaction: <A, E, R>(effect: Effect.Effect<A, E, R>) => Effect.Effect<A, E, R>;};MemoryDriver
Service that provides an in-memory message storage driver with inspectable backing state.
Details
It provides a MessageStorage service, the encoded driver implementation, and maps used to track requests, primary keys, unprocessed envelopes, reply IDs, and the journal.
Signature
declare class MemoryDriver extends Shape< "effect/cluster/MessageStorage/MemoryDriver", { cursors: WeakMap<{}, number>; encoded: Encoded; journal: Array<Encoded>; replyIds: Set<string>; requests: Map<string, MemoryEntry>; requestsByPrimaryKey: Map<string, MemoryEntry>; storage: { readonly clearAddress: (address: EntityAddress) => Effect<void, PersistenceError>; readonly clearReplies: (requestId: Snowflake) => Effect<void, PersistenceError>; readonly registerReplyHandler: <R extends Any>( message: OutgoingRequest<R> | IncomingRequest<R>, ) => Effect<void, EntityNotAssignedToRunner>; readonly repliesFor: <R extends Any>( requests: Iterable<OutgoingRequest<R>>, ) => Effect<Array<Reply<R>>, MalformedMessage | PersistenceError>; readonly repliesForUnfiltered: ( requestIds: Iterable<Snowflake>, ) => Effect<Array<Encoded>, MalformedMessage | PersistenceError>; readonly requestIdForPrimaryKey: (options: { readonly address: EntityAddress; readonly id: string; readonly tag: string; }) => Effect<Option<Snowflake>, PersistenceError>; readonly resetAddress: (address: EntityAddress) => Effect<void, PersistenceError>; readonly resetShards: (shardIds: Iterable<ShardId>) => Effect<void, PersistenceError>; readonly saveEnvelope: ( envelope: OutgoingEnvelope, ) => Effect<void, MalformedMessage | PersistenceError>; readonly saveReply: <R extends Any>( reply: ReplyWithContext<R>, ) => Effect<void, MalformedMessage | PersistenceError>; readonly saveRequest: <R extends Any>( envelope: OutgoingRequest<R>, ) => Effect<SaveResult<R>, MalformedMessage | PersistenceError>; readonly unprocessedMessages: ( shardIds: Iterable<ShardId>, ) => Effect<Array<Incoming<any>>, PersistenceError>; readonly unprocessedMessagesById: <R extends Any>( messageIds: Iterable<Snowflake>, ) => Effect<Array<Incoming<R>>, PersistenceError>; readonly unregisterReplyHandler: (requestId: Snowflake) => Effect<void>; readonly unregisterShardReplyHandlers: (shardId: ShardId) => Effect<void>; readonly withTransaction: <A, E, R>(effect: Effect<A, E, R>) => Effect<A, E, R>; }; unprocessed: Set<Encoded>; }, this> { constructor(_: never); static readonly layer: Layer<MemoryDriver>;}MemoryTransaction
Provides a context reference used in tests to simulate a transaction.
Signature
declare const MemoryTransaction: Reference<boolean>;MessageStorage
Service for cluster mailbox persistence and reply delivery.
Details
It stores outgoing requests, control envelopes, and replies; reads unprocessed messages; manages reply handlers; and provides transaction wrapping for storage operations.
Signature
declare class MessageStorage extends Shape< "effect/cluster/MessageStorage", { readonly clearAddress: (address: EntityAddress) => Effect<void, PersistenceError>; readonly clearReplies: (requestId: Snowflake) => Effect<void, PersistenceError>; readonly registerReplyHandler: <R extends Any>( message: OutgoingRequest<R> | IncomingRequest<R>, ) => Effect<void, EntityNotAssignedToRunner>; readonly repliesFor: <R extends Any>( requests: Iterable<OutgoingRequest<R>>, ) => Effect<Array<Reply<R>>, MalformedMessage | PersistenceError>; readonly repliesForUnfiltered: ( requestIds: Iterable<Snowflake>, ) => Effect<Array<Encoded>, MalformedMessage | PersistenceError>; readonly requestIdForPrimaryKey: (options: { readonly address: EntityAddress; readonly id: string; readonly tag: string; }) => Effect<Option<Snowflake>, PersistenceError>; readonly resetAddress: (address: EntityAddress) => Effect<void, PersistenceError>; readonly resetShards: (shardIds: Iterable<ShardId>) => Effect<void, PersistenceError>; readonly saveEnvelope: ( envelope: OutgoingEnvelope, ) => Effect<void, MalformedMessage | PersistenceError>; readonly saveReply: <R extends Any>( reply: ReplyWithContext<R>, ) => Effect<void, MalformedMessage | PersistenceError>; readonly saveRequest: <R extends Any>( envelope: OutgoingRequest<R>, ) => Effect<SaveResult<R>, MalformedMessage | PersistenceError>; readonly unprocessedMessages: ( shardIds: Iterable<ShardId>, ) => Effect<Array<Incoming<any>>, PersistenceError>; readonly unprocessedMessagesById: <R extends Any>( messageIds: Iterable<Snowflake>, ) => Effect<Array<Incoming<R>>, PersistenceError>; readonly unregisterReplyHandler: (requestId: Snowflake) => Effect<void>; readonly unregisterShardReplyHandlers: (shardId: ShardId) => Effect<void>; readonly withTransaction: <A, E, R>(effect: Effect<A, E, R>) => Effect<A, E, R>; }, this> { constructor(_: never);}
Wraps a concrete message storage implementation with reply-handler management.
Details
The returned service can register waiting reply handlers, notify them when replies are saved, and fail them when a request or shard is unregistered.