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.