/** * @since 1.0.0 */ import type * as Rpc from "@effect/rpc/Rpc"; import * as Arr from "effect/Array"; import * as Context from "effect/Context"; import * as Data from "effect/Data"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import type { PersistenceError } from "./ClusterError.js"; import { EntityNotAssignedToRunner, MalformedMessage } from "./ClusterError.js"; import type { EntityAddress } from "./EntityAddress.js"; import * as Envelope from "./Envelope.js"; import * as Message from "./Message.js"; import * as Reply from "./Reply.js"; import type { ShardId } from "./ShardId.js"; import type { ShardingConfig } from "./ShardingConfig.js"; import * as Snowflake from "./Snowflake.js"; declare const MessageStorage_base: Context.TagClass(envelope: Message.OutgoingRequest) => Effect.Effect, PersistenceError | MalformedMessage>; /** * Save the provided message and its associated metadata. */ readonly saveEnvelope: (envelope: Message.OutgoingEnvelope) => Effect.Effect; /** * Save the provided `Reply` and its associated metadata. */ readonly saveReply: (reply: Reply.ReplyWithContext) => Effect.Effect; /** * Clear the `Reply`s for the given request id. */ readonly clearReplies: (requestId: Snowflake.Snowflake) => Effect.Effect; /** * Retrieves the replies for the specified requests. * * - Un-acknowledged chunk replies * - WithExit replies */ readonly repliesFor: (requests: Iterable>) => Effect.Effect>, PersistenceError | MalformedMessage>; /** * Retrieves the encoded replies for the specified request ids. */ readonly repliesForUnfiltered: (requestIds: Iterable) => Effect.Effect>, PersistenceError | MalformedMessage>; /** * Retrieves the request id for the specified primary key. */ readonly requestIdForPrimaryKey: (options: { readonly address: EntityAddress; readonly tag: string; readonly id: string; }) => Effect.Effect, PersistenceError>; /** * For locally sent messages, register a handler to process the replies. */ readonly registerReplyHandler: (message: Message.OutgoingRequest | Message.IncomingRequest) => Effect.Effect; /** * Unregister the reply handler for the specified message. */ readonly unregisterReplyHandler: (requestId: Snowflake.Snowflake) => Effect.Effect; /** * Unregister the reply handlers for the specified ShardId. */ readonly unregisterShardReplyHandlers: (shardId: ShardId) => Effect.Effect; /** * Retrieves the unprocessed messages for the specified shards. * * A message is unprocessed when: * * - Requests that have no WithExit replies * - Or they have no unacknowledged chunk replies * - The latest AckChunk envelope * - All Interrupt's for unprocessed requests */ readonly unprocessedMessages: (shardIds: Iterable) => Effect.Effect>, PersistenceError>; /** * Retrieves the unprocessed messages by id. */ readonly unprocessedMessagesById: (messageIds: Iterable) => Effect.Effect>, PersistenceError>; /** * Reset the mailbox state for the provided shards. */ readonly resetShards: (shardIds: Iterable) => Effect.Effect; /** * Reset the mailbox state for the provided address. */ readonly resetAddress: (address: EntityAddress) => Effect.Effect; /** * Clear all messages and replies for the provided address. */ readonly clearAddress: (address: EntityAddress) => Effect.Effect; }>; /** * @since 1.0.0 * @category context */ export declare class MessageStorage extends MessageStorage_base { } /** * @since 1.0.0 * @category SaveResult */ export type SaveResult = SaveResult.Success | SaveResult.Duplicate; /** * @since 1.0.0 * @category SaveResult */ export declare const SaveResult: { readonly Success: (args: void) => SaveResult.Success; readonly Duplicate: (args: { readonly originalId: Snowflake.Snowflake; readonly lastReceivedReply: Option.Option>; }) => SaveResult.Duplicate; readonly $is: (tag: Tag) => { >(u: T): u is T & { readonly _tag: Tag; }; (u: unknown): u is Extract | Extract, { readonly _tag: Tag; }>; }; readonly $match: { , const Cases extends { readonly [Tag in Self["_tag"]]: (args: Extract) => any; }>(self: Self, cases: Cases & { [K in Exclude]: never; }): import("effect/Unify").Unify>; any; readonly Duplicate: (args: SaveResult.Duplicate) => any; }>(cases: Cases & { [K in Exclude]: never; }): (self: SaveResult) => import("effect/Unify").Unify>; }; }; /** * @since 1.0.0 * @category SaveResult */ export declare const SaveResultEncoded: { readonly Success: Data.Case.Constructor; readonly Duplicate: Data.Case.Constructor; readonly $is: (tag: Tag) => (u: unknown) => u is Extract | Extract; readonly $match: { any; readonly Duplicate: (args: SaveResult.DuplicateEncoded) => any; }>(cases: Cases & { [K in Exclude]: never; }): (value: SaveResult.Encoded) => import("effect/Unify").Unify>; any; readonly Duplicate: (args: SaveResult.DuplicateEncoded) => any; }>(value: SaveResult.Encoded, cases: Cases & { [K in Exclude]: never; }): import("effect/Unify").Unify>; }; }; /** * @since 1.0.0 * @category SaveResult */ export declare namespace SaveResult { /** * @since 1.0.0 * @category SaveResult */ type Encoded = SaveResult.Success | SaveResult.DuplicateEncoded; /** * @since 1.0.0 * @category SaveResult */ interface Success { readonly _tag: "Success"; } /** * @since 1.0.0 * @category SaveResult */ interface Duplicate { readonly _tag: "Duplicate"; readonly originalId: Snowflake.Snowflake; readonly lastReceivedReply: Option.Option>; } /** * @since 1.0.0 * @category SaveResult */ interface DuplicateEncoded { readonly _tag: "Duplicate"; readonly originalId: Snowflake.Snowflake; readonly lastReceivedReply: Option.Option>; } /** * @since 1.0.0 * @category SaveResult */ interface Constructor extends Data.TaggedEnum.WithGenerics<1> { readonly taggedEnum: SaveResult; } } /** * @since 1.0.0 * @category Encoded */ export type Encoded = { /** * Save the provided message and its associated metadata. */ readonly saveEnvelope: (options: { readonly envelope: Envelope.Envelope.Encoded; readonly primaryKey: string | null; readonly deliverAt: number | null; }) => Effect.Effect; /** * Save the provided `Reply` and its associated metadata. */ readonly saveReply: (reply: Reply.ReplyEncoded) => Effect.Effect; /** * Remove the replies for the specified request. */ readonly clearReplies: (requestId: Snowflake.Snowflake) => Effect.Effect; /** * Retrieves the request id for the specified primary key. */ readonly requestIdForPrimaryKey: (primaryKey: string) => Effect.Effect, PersistenceError>; /** * Retrieves the replies for the specified requests. * * - Un-acknowledged chunk replies * - WithExit replies */ readonly repliesFor: (requestIds: Arr.NonEmptyArray) => Effect.Effect>, PersistenceError>; /** * Retrieves the replies for the specified request ids. */ readonly repliesForUnfiltered: (requestIds: Arr.NonEmptyArray) => Effect.Effect>, PersistenceError>; /** * Retrieves the unprocessed messages for the given shards. * * A message is unprocessed when: * * - Requests that have no WithExit replies * - Or they have no unacknowledged chunk replies * - The latest AckChunk envelope * - All Interrupt's for unprocessed requests */ readonly unprocessedMessages: (shardIds: Arr.NonEmptyArray, now: number) => Effect.Effect>; }>, PersistenceError>; /** * Retrieves the unprocessed messages by id. */ readonly unprocessedMessagesById: (messageIds: Arr.NonEmptyArray, now: number) => Effect.Effect>; }>, PersistenceError>; /** * Reset the mailbox state for the provided address. */ readonly resetAddress: (address: EntityAddress) => Effect.Effect; /** * Clear all messages and replies for the provided address. */ readonly clearAddress: (address: EntityAddress) => Effect.Effect; /** * Reset the mailbox state for the provided shards. */ readonly resetShards: (shardIds: Arr.NonEmptyArray) => Effect.Effect; }; /** * @since 1.0.0 * @category Encoded */ export type EncodedUnprocessedOptions = { readonly existingShards: Array; readonly newShards: Array; readonly cursor: Option.Option; }; /** * @since 1.0.0 * @category Encoded */ export type EncodedRepliesOptions = { readonly existingRequests: Array; readonly newRequests: Array; readonly cursor: Option.Option; }; /** * @since 1.0.0 * @category constructors */ export declare const make: (storage: Omit) => Effect.Effect; /** * @since 1.0.0 * @category constructors */ export declare const makeEncoded: (encoded: Encoded) => Effect.Effect; /** * @since 1.0.0 * @category Constructors */ export declare const noop: MessageStorage["Type"]; /** * @since 1.0.0 * @category Memory */ export type MemoryEntry = { readonly envelope: Envelope.Request.Encoded; lastReceivedChunk: Option.Option>; replies: Array>; deliverAt: number | null; }; declare const MemoryDriver_base: Effect.Service.Class]; readonly effect: Effect.Effect<{ readonly storage: { /** * Save the provided message and its associated metadata. */ readonly saveRequest: (envelope: Message.OutgoingRequest) => Effect.Effect, PersistenceError | MalformedMessage>; /** * Save the provided message and its associated metadata. */ readonly saveEnvelope: (envelope: Message.OutgoingEnvelope) => Effect.Effect; /** * Save the provided `Reply` and its associated metadata. */ readonly saveReply: (reply: Reply.ReplyWithContext) => Effect.Effect; /** * Clear the `Reply`s for the given request id. */ readonly clearReplies: (requestId: Snowflake.Snowflake) => Effect.Effect; /** * Retrieves the replies for the specified requests. * * - Un-acknowledged chunk replies * - WithExit replies */ readonly repliesFor: (requests: Iterable>) => Effect.Effect>, PersistenceError | MalformedMessage>; /** * Retrieves the encoded replies for the specified request ids. */ readonly repliesForUnfiltered: (requestIds: Iterable) => Effect.Effect>, PersistenceError | MalformedMessage>; /** * Retrieves the request id for the specified primary key. */ readonly requestIdForPrimaryKey: (options: { readonly address: EntityAddress; readonly tag: string; readonly id: string; }) => Effect.Effect, PersistenceError>; /** * For locally sent messages, register a handler to process the replies. */ readonly registerReplyHandler: (message: Message.OutgoingRequest | Message.IncomingRequest) => Effect.Effect; /** * Unregister the reply handler for the specified message. */ readonly unregisterReplyHandler: (requestId: Snowflake.Snowflake) => Effect.Effect; /** * Unregister the reply handlers for the specified ShardId. */ readonly unregisterShardReplyHandlers: (shardId: ShardId) => Effect.Effect; /** * Retrieves the unprocessed messages for the specified shards. * * A message is unprocessed when: * * - Requests that have no WithExit replies * - Or they have no unacknowledged chunk replies * - The latest AckChunk envelope * - All Interrupt's for unprocessed requests */ readonly unprocessedMessages: (shardIds: Iterable) => Effect.Effect>, PersistenceError>; /** * Retrieves the unprocessed messages by id. */ readonly unprocessedMessagesById: (messageIds: Iterable) => Effect.Effect>, PersistenceError>; /** * Reset the mailbox state for the provided shards. */ readonly resetShards: (shardIds: Iterable) => Effect.Effect; /** * Reset the mailbox state for the provided address. */ readonly resetAddress: (address: EntityAddress) => Effect.Effect; /** * Clear all messages and replies for the provided address. */ readonly clearAddress: (address: EntityAddress) => Effect.Effect; }; readonly encoded: Encoded; readonly requests: Map; readonly requestsByPrimaryKey: Map; readonly unprocessed: Set; readonly replyIds: Set; readonly journal: Envelope.Envelope.Encoded[]; readonly cursors: WeakMap<{}, number>; }, never, Snowflake.Generator>; }>; /** * @since 1.0.0 * @category Memory */ export declare class MemoryDriver extends MemoryDriver_base { } /** * @since 1.0.0 * @category layers */ export declare const layerNoop: Layer.Layer; /** * @since 1.0.0 * @category layers */ export declare const layerMemory: Layer.Layer; export {}; //# sourceMappingURL=MessageStorage.d.ts.map