Skip to content

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.

16 exports Added in v4.0.0 Source

Constructors

make

Added in v4.0.0 Source

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.

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

Added in v4.0.0 Source

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

noop

Added in v4.0.0 Source

No-op MessageStorage service that does not persist messages or replies.

Signature

declare const noop: MessageStorage["Service"];

SaveResult

Added in v4.0.0 Source

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

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

Added in v4.0.0 Source

Layer that provides in-memory message storage and its backing MemoryDriver.

Signature

declare const layerMemory: Layer.Layer<MessageStorage | MemoryDriver, never, ShardingConfig>;

layerNoop

Added in v4.0.0 Source

Layer that provides the no-op MessageStorage service.

Signature

declare const layerNoop: Layer.Layer<MessageStorage>;

Models

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

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

Added in v4.0.0 Source

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

Added in v4.0.0 Source

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

Added in v4.0.0 Source

Variants and helper types for SaveResult.

Services

Encoded type

Added in v4.0.0 Source

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

Added in v4.0.0 Source

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

Provides a context reference used in tests to simulate a transaction.

Signature

declare const MemoryTransaction: Reference<boolean>;

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