/**
 * @beignet/core/outbox
 *
 * Durable outbox primitives for transactionally recording events and jobs that
 * should be delivered after the owning database transaction commits.
 */

import {
  type EventPayloadDef,
  type EventPublishOptions,
  type InferEventPayload,
  parseEventPayload,
} from "../events/index.js";
import {
  getJobRetryDelayMs,
  getJobRetryMaxAttempts,
  type InferJobPayload,
  type JobDef,
  parseJobPayload,
  SINGLE_ATTEMPT_DISPATCH,
  type SingleAttemptJobDispatch,
  shouldRetryJob,
} from "../jobs/index.js";
import type { JobDispatcherPort } from "../ports/events.js";
import type { DomainEventRecorderPort } from "../ports/unit-of-work.js";
import {
  type BaseProviderInstrumentationEvent,
  createProviderInstrumentation,
  type ProviderInstrumentationTarget,
} from "../providers/index.js";
import {
  captureTraceCarrier,
  parseTraceCarrier,
  resolveTracingPort,
  runWithTracing,
  type TraceCarrier,
  type TracingPort,
} from "../tracing/index.js";

/**
 * Value or promise of that value.
 */
export type MaybePromise<T> = T | Promise<T>;

/**
 * Default lease duration for claimed outbox messages.
 */
export const DEFAULT_OUTBOX_LEASE_MS = 30_000;
/**
 * Default maximum delivery attempts before a message is dead-lettered.
 */
export const DEFAULT_OUTBOX_MAX_ATTEMPTS = 3;

/**
 * Message kinds supported by the Beignet outbox.
 */
export type OutboxMessageKind = "event" | "job";

/**
 * Delivery status for an outbox message.
 */
export type OutboxMessageStatus =
  | "pending"
  | "claimed"
  | "delivered"
  | "deadLettered";

/**
 * JSON-serializable value accepted by the outbox.
 */
export type OutboxJsonValue =
  | null
  | string
  | number
  | boolean
  | readonly OutboxJsonValue[]
  | { readonly [key: string]: OutboxJsonValue };

/**
 * JSON object accepted by outbox payload helpers.
 */
export type OutboxJsonObject = { readonly [key: string]: OutboxJsonValue };

/**
 * Serialized delivery error stored on failed messages.
 */
export interface OutboxErrorInfo {
  /**
   * Error name when available.
   */
  name?: string;
  /**
   * Error message.
   */
  message: string;
  /**
   * Error stack when available.
   */
  stack?: string;
}

/**
 * Input for enqueueing a raw outbox message.
 */
export interface OutboxEnqueueInput {
  /**
   * Optional caller-provided message ID.
   */
  id?: string;
  /**
   * Message kind.
   */
  kind: OutboxMessageKind;
  /**
   * Event or job name.
   */
  name: string;
  /**
   * JSON-serializable payload.
   */
  payload: OutboxJsonValue;
  /** Versioned trace context captured when the message was recorded. */
  trace?: TraceCarrier;
  /**
   * Earliest time the message may be claimed.
   */
  availableAt?: Date;
  /**
   * Maximum delivery attempts before dead-lettering.
   */
  maxAttempts?: number;
}

/**
 * Durable outbox message record.
 */
export interface OutboxMessage {
  /**
   * Stable message ID.
   */
  id: string;
  /**
   * Message kind.
   */
  kind: OutboxMessageKind;
  /**
   * Event or job name.
   */
  name: string;
  /**
   * JSON-serializable payload.
   */
  payload: OutboxJsonValue;
  /** Versioned trace context captured when the message was recorded. */
  trace?: TraceCarrier;
  /**
   * Current delivery status.
   */
  status: OutboxMessageStatus;
  /**
   * Number of claim attempts.
   */
  attempts: number;
  /**
   * Maximum delivery attempts before dead-lettering.
   */
  maxAttempts: number;
  /**
   * Earliest time the message may be claimed.
   */
  availableAt: Date;
  /**
   * Last claim timestamp.
   */
  claimedAt: Date | null;
  /**
   * Lease expiration timestamp for claimed messages.
   */
  lockedUntil: Date | null;
  /**
   * Token required to mark a claimed message delivered or failed.
   */
  claimToken: string | null;
  /**
   * Delivery timestamp.
   */
  deliveredAt: Date | null;
  /**
   * Last delivery error.
   */
  lastError: OutboxErrorInfo | null;
  /**
   * Creation timestamp.
   */
  createdAt: Date;
  /**
   * Last update timestamp.
   */
  updatedAt: Date;
}

/**
 * Message returned from a successful outbox claim.
 */
export interface ClaimedOutboxMessage
  extends Omit<
    OutboxMessage,
    "claimToken" | "claimedAt" | "lockedUntil" | "status"
  > {
  status: "claimed";
  claimToken: string;
  claimedAt: Date;
  lockedUntil: Date;
}

/**
 * Options for leasing a batch of pending outbox messages for delivery.
 */
export interface OutboxClaimBatchOptions {
  /**
   * Maximum messages to claim in one batch.
   */
  limit: number;
  /**
   * Claim timestamp.
   */
  now?: Date;
  /**
   * Lease duration in milliseconds.
   */
  leaseMs?: number;
}

/**
 * Input for marking a claimed message delivered.
 */
export interface OutboxMarkDeliveredInput {
  /**
   * Claimed message ID.
   */
  id: string;
  /**
   * Claim token returned by `claimBatch(...)`.
   */
  claimToken: string;
  /**
   * Delivery timestamp.
   */
  now?: Date;
}

/**
 * Input for marking a claimed message failed.
 */
export interface OutboxMarkFailedInput {
  /**
   * Claimed message ID.
   */
  id: string;
  /**
   * Claim token returned by `claimBatch(...)`.
   */
  claimToken: string;
  /**
   * Delivery error.
   */
  error?: unknown;
  /**
   * Next time the message may be claimed.
   */
  retryAt?: Date;
  /**
   * Whether this failure should dead-letter the message.
   */
  deadLetter?: boolean;
  /**
   * Failure timestamp.
   */
  now?: Date;
}

/**
 * Default maximum messages returned from outbox admin list calls.
 */
export const DEFAULT_OUTBOX_ADMIN_LIST_LIMIT = 50;

/**
 * Shared filters for outbox admin read operations.
 */
export interface OutboxMessageQuery {
  /**
   * Status or statuses to include.
   */
  status?: OutboxMessageStatus | readonly OutboxMessageStatus[];
  /**
   * Message kind to include.
   */
  kind?: OutboxMessageKind;
  /**
   * Event or job name to include.
   */
  name?: string;
  /**
   * Include messages last updated before this timestamp.
   */
  updatedBefore?: Date;
  /**
   * Include delivered messages delivered before this timestamp.
   */
  deliveredBefore?: Date;
}

/**
 * Options for listing outbox messages through the admin port.
 */
export interface OutboxListMessagesOptions extends OutboxMessageQuery {
  /**
   * Maximum messages to return. Defaults to
   * `DEFAULT_OUTBOX_ADMIN_LIST_LIMIT`.
   */
  limit?: number;
}

/**
 * Options for counting outbox messages through the admin port.
 */
export interface OutboxCountMessagesOptions extends OutboxMessageQuery {}

/**
 * Input for returning a dead-lettered message to the pending queue.
 */
export interface OutboxRequeueMessageInput {
  /**
   * Dead-lettered message ID.
   */
  id: string;
  /**
   * Earliest time the message may be claimed again. Defaults to `now`.
   */
  availableAt?: Date;
  /**
   * Reset attempts to zero before requeueing. Defaults to preserving the
   * attempt count so operators can decide whether to grant a fresh retry
   * budget.
   */
  resetAttempts?: boolean;
  /**
   * Requeue timestamp.
   */
  now?: Date;
}

/**
 * Input for purging dead-lettered messages.
 */
export interface OutboxPurgeDeadLetteredInput {
  /**
   * Only purge messages last updated before this timestamp. Omit only when the
   * caller intentionally wants to purge all dead-lettered messages.
   */
  before?: Date;
  /**
   * Maximum messages to purge.
   */
  limit?: number;
}

/**
 * Input for pruning delivered messages.
 */
export interface OutboxPruneDeliveredInput {
  /**
   * Prune delivered messages delivered before this timestamp.
   */
  before: Date;
  /**
   * Maximum messages to prune.
   */
  limit?: number;
}

/**
 * Result for outbox admin delete operations.
 */
export interface OutboxDeleteResult {
  /**
   * Number of rows deleted.
   */
  deleted: number;
}

/**
 * App-facing outbox storage port.
 *
 * Durable adapters should claim messages atomically and require `claimToken`
 * for delivery/failure updates.
 */
export interface OutboxPort {
  /**
   * Enqueue a new pending message.
   */
  enqueue(input: OutboxEnqueueInput): Promise<OutboxMessage>;
  /**
   * Atomically claim eligible messages for one worker.
   */
  claimBatch(options: OutboxClaimBatchOptions): Promise<ClaimedOutboxMessage[]>;
  /**
   * Mark a claimed message delivered.
   */
  markDelivered(input: OutboxMarkDeliveredInput): Promise<void>;
  /**
   * Mark a claimed message failed, retryable, or dead-lettered.
   */
  markFailed(input: OutboxMarkFailedInput): Promise<void>;
}

/**
 * Operational outbox admin port.
 *
 * Keep this separate from `OutboxPort` so request and drain contexts can expose
 * only the hot-path delivery operations. Wire this port into maintenance
 * contexts for CLI commands, runbooks, and devtools.
 */
export interface OutboxAdminPort {
  /**
   * List messages ordered by newest update first.
   */
  listMessages(
    options?: OutboxListMessagesOptions,
  ): Promise<readonly OutboxMessage[]>;
  /**
   * Count messages matching admin filters.
   */
  countMessages(options?: OutboxCountMessagesOptions): Promise<number>;
  /**
   * Fetch one message by ID.
   */
  getMessage(id: string): Promise<OutboxMessage | null>;
  /**
   * Return a dead-lettered message to pending state.
   */
  requeueMessage(input: OutboxRequeueMessageInput): Promise<OutboxMessage>;
  /**
   * Delete dead-lettered messages.
   */
  purgeDeadLettered(
    input?: OutboxPurgeDeadLetteredInput,
  ): Promise<OutboxDeleteResult>;
  /**
   * Delete delivered messages older than a retention cutoff.
   */
  pruneDelivered(input: OutboxPruneDeliveredInput): Promise<OutboxDeleteResult>;
}

/**
 * In-memory outbox for tests and local examples.
 */
export interface MemoryOutboxPort extends OutboxPort, OutboxAdminPort {
  /**
   * Current message snapshots.
   */
  readonly messages: readonly OutboxMessage[];
  /**
   * Remove all messages.
   */
  clear(): void;
}

/**
 * Options for `createOutboxMessage(...)`.
 */
export interface CreateOutboxMessageOptions {
  /**
   * Generated message ID override. Wins over `input.id`.
   */
  id?: string;
  /**
   * Fallback ID factory used when neither `id` nor `input.id` is provided.
   */
  createId?: () => string;
  /**
   * Timestamp used for created/updated/available dates.
   */
  now?: Date;
}

/**
 * Options for typed event/job enqueue helpers.
 */
export interface EnqueueTypedOutboxOptions {
  /**
   * Optional caller-provided message ID.
   */
  id?: string;
  /**
   * Earliest time the message may be claimed.
   */
  availableAt?: Date;
  /**
   * Maximum delivery attempts before dead-lettering.
   */
  maxAttempts?: number;
  /** Explicit trace context to persist with this message. */
  trace?: TraceCarrier;
  /** Tracing port used to capture the active context at enqueue time. */
  tracing?: TracingPort;
}

/**
 * Registry of definitions that `drainOutbox(...)` can deliver.
 */
export interface OutboxRegistry {
  /**
   * Event definitions keyed by event name.
   */
  readonly events: ReadonlyMap<string, EventPayloadDef>;
  /**
   * Job definitions keyed by job name.
   */
  readonly jobs: ReadonlyMap<string, JobDef>;
}

/**
 * Input for defining an outbox registry.
 */
export interface DefineOutboxRegistryInput {
  /**
   * Events that may be delivered from the outbox.
   */
  events?: readonly EventPayloadDef[];
  /**
   * Jobs that may be delivered from the outbox.
   */
  jobs?: readonly JobDef[];
}

/**
 * Correlation fields attached to outbox instrumentation events.
 */
export type OutboxInstrumentationContext = Pick<
  BaseProviderInstrumentationEvent,
  "requestId" | "traceId" | "spanId" | "parentSpanId" | "traceparent"
>;

/**
 * Options for draining one outbox batch.
 */
export interface DrainOutboxOptions {
  /**
   * Outbox storage port.
   */
  outbox: OutboxPort;
  /**
   * Registry used to resolve message names to event/job definitions.
   */
  registry: OutboxRegistry;
  /**
   * Event bus used for event messages.
   */
  eventBus?: {
    publish<E extends EventPayloadDef>(
      event: E,
      payload: InferEventPayload<E>,
      options?: EventPublishOptions,
    ): MaybePromise<void>;
  };
  /**
   * Job dispatcher used for job messages.
   */
  jobs?: JobDispatcherPort;
  /**
   * Maximum messages to claim in one drain pass.
   */
  batchSize?: number;
  /**
   * Timestamp used for claiming and state updates.
   */
  now?: Date;
  /**
   * Claim lease duration in milliseconds.
   */
  leaseMs?: number;
  /**
   * Retry delay in milliseconds or function for per-message delay.
   */
  retryDelayMs?:
    | number
    | ((args: {
        message: ClaimedOutboxMessage;
        error: unknown;
        now: Date;
      }) => number);
  /**
   * Optional instrumentation target for delivery, retry, and dead-letter
   * visibility.
   */
  instrumentation?: ProviderInstrumentationTarget;
  /**
   * Optional correlation fields attached to outbox instrumentation events.
   */
  instrumentationContext?: OutboxInstrumentationContext;
  /**
   * Observer called when delivery fails. Observer failures are ignored so the
   * original delivery failure still controls retry/dead-letter behavior.
   */
  onError?: (
    error: unknown,
    message: ClaimedOutboxMessage,
  ) => MaybePromise<void>;
  /**
   * Observer called after a failed delivery is successfully moved to the dead
   * letter state. Observer failures are ignored.
   */
  onDeadLetter?: (
    error: unknown,
    message: ClaimedOutboxMessage,
  ) => MaybePromise<void>;
  /**
   * Observer called when a failed delivery cannot be settled as retryable or
   * dead-lettered. Observer failures are ignored.
   */
  onSettlementError?: (
    settlementError: unknown,
    message: ClaimedOutboxMessage,
    deliveryError: unknown,
  ) => MaybePromise<void>;
}

/**
 * Summary returned from one `drainOutbox(...)` pass.
 */
export interface DrainOutboxResult {
  /**
   * Messages claimed in this batch.
   */
  claimed: number;
  /**
   * Messages delivered successfully.
   */
  delivered: number;
  /**
   * Messages scheduled for retry.
   */
  retried: number;
  /**
   * Messages moved to dead letter state.
   */
  deadLettered: number;
}

/**
 * Error thrown when an outbox payload is not JSON serializable.
 */
export class OutboxSerializationError extends Error {
  constructor(message: string) {
    super(message);
    this.name = "OutboxSerializationError";
  }
}

/**
 * Error thrown when an outbox message cannot be resolved through the registry.
 */
export class OutboxRegistryError extends Error {
  constructor(message: string) {
    super(message);
    this.name = "OutboxRegistryError";
  }
}

/**
 * Error thrown when a claimed message cannot be updated with the supplied token.
 */
export class OutboxClaimError extends Error {
  /**
   * Message ID involved in the claim error.
   */
  readonly id: string;

  constructor(args: { id: string; message: string }) {
    super(args.message);
    this.name = "OutboxClaimError";
    this.id = args.id;
  }
}

/**
 * Error thrown when an outbox admin operation cannot be completed safely.
 */
export class OutboxAdminError extends Error {
  /**
   * Message ID involved in the admin error, when applicable.
   */
  readonly id?: string;

  constructor(args: { message: string; id?: string }) {
    super(args.message);
    this.name = "OutboxAdminError";
    this.id = args.id;
  }
}

function assertNonEmptyString(name: string, value: string): void {
  if (typeof value !== "string" || value.trim().length === 0) {
    throw new Error(`${name} must be a non-empty string`);
  }
}

function assertPositiveInteger(name: string, value: number): void {
  if (!Number.isInteger(value) || value <= 0) {
    throw new Error(`${name} must be a positive integer`);
  }
}

function assertValidDate(name: string, value: Date): void {
  if (!(value instanceof Date) || Number.isNaN(value.getTime())) {
    throw new Error(`${name} must be a valid Date`);
  }
}

function cloneDate(value: Date): Date {
  return new Date(value.getTime());
}

function createId(): string {
  if (!globalThis.crypto?.randomUUID) {
    throw new Error("crypto.randomUUID is required to create outbox IDs.");
  }

  return globalThis.crypto.randomUUID();
}

function assertJsonValue(
  value: unknown,
  path: readonly string[] = [],
  seen: WeakSet<object> = new WeakSet(),
): OutboxJsonValue {
  const label = path.length > 0 ? path.join(".") : "payload";

  if (value === null) return null;

  if (
    typeof value === "string" ||
    typeof value === "boolean" ||
    typeof value === "number"
  ) {
    if (typeof value === "number" && !Number.isFinite(value)) {
      throw new OutboxSerializationError(
        `Outbox ${label} must be a finite number.`,
      );
    }
    return value;
  }

  if (Array.isArray(value)) {
    if (seen.has(value)) {
      throw new OutboxSerializationError(
        `Outbox ${label} must be JSON serializable. Circular references are not supported.`,
      );
    }
    seen.add(value);
    try {
      return value.map((item, index) =>
        assertJsonValue(item, [...path, String(index)], seen),
      );
    } finally {
      seen.delete(value);
    }
  }

  if (typeof value === "object") {
    if (value instanceof Date) {
      throw new OutboxSerializationError(
        `Outbox ${label} must be JSON serializable. Convert Date values to strings before enqueueing.`,
      );
    }
    if (seen.has(value)) {
      throw new OutboxSerializationError(
        `Outbox ${label} must be JSON serializable. Circular references are not supported.`,
      );
    }

    seen.add(value);
    try {
      const record = value as Record<string, unknown>;
      const output: Record<string, OutboxJsonValue> = {};
      for (const key of Object.keys(record)) {
        const child = record[key];
        if (child === undefined) {
          throw new OutboxSerializationError(
            `Outbox ${[...path, key].join(".")} cannot be undefined.`,
          );
        }
        output[key] = assertJsonValue(child, [...path, key], seen);
      }
      return output;
    } finally {
      seen.delete(value);
    }
  }

  throw new OutboxSerializationError(
    `Outbox ${label} must be JSON serializable. Received ${typeof value}.`,
  );
}

/**
 * Convert an unknown value to an outbox-safe JSON value.
 *
 * Dates, undefined values, functions, non-finite numbers, symbols, and circular
 * references are rejected so durable adapters can store the payload safely.
 */
export function toOutboxJsonValue(value: unknown): OutboxJsonValue {
  const jsonValue = assertJsonValue(value);
  return JSON.parse(JSON.stringify(jsonValue)) as OutboxJsonValue;
}

/**
 * Serialize an unknown delivery error into outbox error metadata.
 */
export function serializeOutboxError(error: unknown): OutboxErrorInfo {
  if (error instanceof Error) {
    return {
      name: error.name,
      message: error.message,
      stack: error.stack,
    };
  }

  if (typeof error === "string") {
    return { message: error };
  }

  return { message: "Unknown outbox delivery error" };
}

function copyMessage(message: OutboxMessage): OutboxMessage {
  return {
    ...message,
    availableAt: cloneDate(message.availableAt),
    claimedAt: message.claimedAt ? cloneDate(message.claimedAt) : null,
    lockedUntil: message.lockedUntil ? cloneDate(message.lockedUntil) : null,
    deliveredAt: message.deliveredAt ? cloneDate(message.deliveredAt) : null,
    createdAt: cloneDate(message.createdAt),
    updatedAt: cloneDate(message.updatedAt),
    lastError: message.lastError ? { ...message.lastError } : null,
  };
}

function normalizeStatusFilter(
  status: OutboxMessageQuery["status"],
): ReadonlySet<OutboxMessageStatus> | undefined {
  if (status === undefined) return undefined;
  return new Set(Array.isArray(status) ? status : [status]);
}

function messageMatchesQuery(
  message: OutboxMessage,
  query: OutboxMessageQuery = {},
): boolean {
  const statuses = normalizeStatusFilter(query.status);
  if (statuses && !statuses.has(message.status)) return false;
  if (query.kind !== undefined && message.kind !== query.kind) return false;
  if (query.name !== undefined && message.name !== query.name) return false;
  if (
    query.updatedBefore !== undefined &&
    message.updatedAt.getTime() >= query.updatedBefore.getTime()
  ) {
    return false;
  }
  if (
    query.deliveredBefore !== undefined &&
    (message.deliveredAt === null ||
      message.deliveredAt.getTime() >= query.deliveredBefore.getTime())
  ) {
    return false;
  }
  return true;
}

function compareMessagesForAdminList(
  left: OutboxMessage,
  right: OutboxMessage,
): number {
  const updated = right.updatedAt.getTime() - left.updatedAt.getTime();
  if (updated !== 0) return updated;
  const created = right.createdAt.getTime() - left.createdAt.getTime();
  if (created !== 0) return created;
  return right.id.localeCompare(left.id);
}

function compareMessagesForAdminCleanup(
  left: OutboxMessage,
  right: OutboxMessage,
): number {
  const updated = left.updatedAt.getTime() - right.updatedAt.getTime();
  if (updated !== 0) return updated;
  const created = left.createdAt.getTime() - right.createdAt.getTime();
  if (created !== 0) return created;
  return left.id.localeCompare(right.id);
}

function compareDeliveredMessagesForPrune(
  left: OutboxMessage,
  right: OutboxMessage,
): number {
  const delivered =
    (left.deliveredAt?.getTime() ?? 0) - (right.deliveredAt?.getTime() ?? 0);
  if (delivered !== 0) return delivered;
  return compareMessagesForAdminCleanup(left, right);
}

function resolveAdminListLimit(limit: number | undefined): number {
  const resolved = limit ?? DEFAULT_OUTBOX_ADMIN_LIST_LIMIT;
  assertPositiveInteger("limit", resolved);
  return resolved;
}

function toClaimedMessage(message: OutboxMessage): ClaimedOutboxMessage {
  if (
    message.status !== "claimed" ||
    !message.claimToken ||
    !message.claimedAt ||
    !message.lockedUntil
  ) {
    throw new OutboxClaimError({
      id: message.id,
      message: `Outbox message "${message.id}" is not claimed.`,
    });
  }

  return {
    ...copyMessage(message),
    status: "claimed",
    claimToken: message.claimToken,
    claimedAt: cloneDate(message.claimedAt),
    lockedUntil: cloneDate(message.lockedUntil),
  };
}

/**
 * Create a validated pending outbox message.
 */
export function createOutboxMessage(
  input: OutboxEnqueueInput,
  options: CreateOutboxMessageOptions = {},
): OutboxMessage {
  assertNonEmptyString("kind", input.kind);
  assertNonEmptyString("name", input.name);
  if (input.id !== undefined) assertNonEmptyString("id", input.id);
  if (input.maxAttempts !== undefined) {
    assertPositiveInteger("maxAttempts", input.maxAttempts);
  }

  const now = options.now ?? new Date();
  const trace = parseTraceCarrier(input.trace);
  return {
    id: options.id ?? input.id ?? options.createId?.() ?? createId(),
    kind: input.kind,
    name: input.name,
    payload: toOutboxJsonValue(input.payload),
    ...(trace ? { trace } : {}),
    status: "pending",
    attempts: 0,
    maxAttempts: input.maxAttempts ?? DEFAULT_OUTBOX_MAX_ATTEMPTS,
    availableAt: input.availableAt
      ? cloneDate(input.availableAt)
      : cloneDate(now),
    claimedAt: null,
    lockedUntil: null,
    claimToken: null,
    deliveredAt: null,
    lastError: null,
    createdAt: cloneDate(now),
    updatedAt: cloneDate(now),
  };
}

function isEligible(message: OutboxMessage, now: Date): boolean {
  if (message.status === "pending") {
    return message.availableAt.getTime() <= now.getTime();
  }

  return (
    message.status === "claimed" &&
    message.lockedUntil !== null &&
    message.lockedUntil.getTime() <= now.getTime()
  );
}

/**
 * Options for `createMemoryOutbox(...)`.
 */
export interface MemoryOutboxOptions {
  /**
   * Message and claim-token ID factory. Defaults to `crypto.randomUUID()`.
   */
  id?: () => string;
  /**
   * Clock used for enqueue, claim, and completion timestamps when a call does
   * not supply its own `now`. Defaults to the system clock.
   */
  now?: () => Date;
}

/**
 * Create an in-memory outbox for tests and local examples.
 *
 * The memory outbox is process-local and not durable.
 */
export function createMemoryOutbox(
  storeOptions: MemoryOutboxOptions = {},
): MemoryOutboxPort {
  const createStoreId = storeOptions.id ?? createId;
  const storeNow = storeOptions.now ?? (() => new Date());
  const messages = new Map<string, OutboxMessage>();

  function getClaimedOrThrow(id: string, claimToken: string): OutboxMessage {
    const message = messages.get(id);
    if (!message) {
      throw new OutboxClaimError({
        id,
        message: `Outbox message "${id}" does not exist.`,
      });
    }
    if (message.status !== "claimed" || message.claimToken !== claimToken) {
      throw new OutboxClaimError({
        id,
        message: `Outbox message "${id}" is not claimed by this worker.`,
      });
    }
    return message;
  }

  return {
    get messages() {
      return [...messages.values()].map(copyMessage);
    },

    async listMessages(options = {}) {
      const limit = resolveAdminListLimit(options.limit);
      if (options.updatedBefore) {
        assertValidDate("updatedBefore", options.updatedBefore);
      }
      if (options.deliveredBefore) {
        assertValidDate("deliveredBefore", options.deliveredBefore);
      }

      return [...messages.values()]
        .filter((message) => messageMatchesQuery(message, options))
        .sort(compareMessagesForAdminList)
        .slice(0, limit)
        .map(copyMessage);
    },

    async countMessages(options = {}) {
      if (options.updatedBefore) {
        assertValidDate("updatedBefore", options.updatedBefore);
      }
      if (options.deliveredBefore) {
        assertValidDate("deliveredBefore", options.deliveredBefore);
      }

      return [...messages.values()].filter((message) =>
        messageMatchesQuery(message, options),
      ).length;
    },

    async getMessage(id) {
      assertNonEmptyString("id", id);
      const message = messages.get(id);
      return message ? copyMessage(message) : null;
    },

    async requeueMessage(input) {
      assertNonEmptyString("id", input.id);
      const message = messages.get(input.id);
      if (!message) {
        throw new OutboxAdminError({
          id: input.id,
          message: `Outbox message "${input.id}" does not exist.`,
        });
      }
      if (message.status !== "deadLettered") {
        throw new OutboxAdminError({
          id: input.id,
          message: `Outbox message "${input.id}" is not dead-lettered.`,
        });
      }

      const now = input.now ?? storeNow();
      const availableAt = input.availableAt ?? now;
      assertValidDate("now", now);
      assertValidDate("availableAt", availableAt);

      message.status = "pending";
      message.availableAt = cloneDate(availableAt);
      message.claimToken = null;
      message.claimedAt = null;
      message.lockedUntil = null;
      if (input.resetAttempts) message.attempts = 0;
      message.updatedAt = cloneDate(now);

      return copyMessage(message);
    },

    async purgeDeadLettered(input = {}) {
      const limit =
        input.limit === undefined
          ? undefined
          : resolveAdminListLimit(input.limit);
      if (input.before) assertValidDate("before", input.before);

      const candidates = [...messages.values()]
        .filter(
          (message) =>
            message.status === "deadLettered" &&
            (input.before === undefined ||
              message.updatedAt.getTime() < input.before.getTime()),
        )
        .sort(compareMessagesForAdminCleanup)
        .slice(0, limit);

      for (const message of candidates) {
        messages.delete(message.id);
      }

      return { deleted: candidates.length };
    },

    async pruneDelivered(input) {
      assertValidDate("before", input.before);
      const limit =
        input.limit === undefined
          ? undefined
          : resolveAdminListLimit(input.limit);
      const candidates = [...messages.values()]
        .filter(
          (message) =>
            message.status === "delivered" &&
            message.deliveredAt !== null &&
            message.deliveredAt.getTime() < input.before.getTime(),
        )
        .sort(compareDeliveredMessagesForPrune)
        .slice(0, limit);

      for (const message of candidates) {
        messages.delete(message.id);
      }

      return { deleted: candidates.length };
    },

    async enqueue(input) {
      const message = createOutboxMessage(input, {
        createId: createStoreId,
        now: storeNow(),
      });
      if (messages.has(message.id)) {
        throw new Error(`Outbox message "${message.id}" already exists.`);
      }
      messages.set(message.id, message);
      return copyMessage(message);
    },

    async claimBatch(options) {
      assertPositiveInteger("limit", options.limit);
      const now = options.now ?? storeNow();
      const leaseMs = options.leaseMs ?? DEFAULT_OUTBOX_LEASE_MS;
      assertPositiveInteger("leaseMs", leaseMs);
      const lockedUntil = new Date(now.getTime() + leaseMs);
      const claimed: ClaimedOutboxMessage[] = [];

      const eligible = [...messages.values()]
        .filter((message) => isEligible(message, now))
        .sort((a, b) => {
          const available = a.availableAt.getTime() - b.availableAt.getTime();
          if (available !== 0) return available;
          return a.createdAt.getTime() - b.createdAt.getTime();
        })
        .slice(0, options.limit);

      for (const message of eligible) {
        message.status = "claimed";
        message.attempts += 1;
        message.claimToken = createStoreId();
        message.claimedAt = cloneDate(now);
        message.lockedUntil = cloneDate(lockedUntil);
        message.updatedAt = cloneDate(now);
        claimed.push(toClaimedMessage(message));
      }

      return claimed;
    },

    async markDelivered(input) {
      assertNonEmptyString("id", input.id);
      assertNonEmptyString("claimToken", input.claimToken);
      const message = getClaimedOrThrow(input.id, input.claimToken);
      const now = input.now ?? storeNow();

      message.status = "delivered";
      message.deliveredAt = cloneDate(now);
      message.claimToken = null;
      message.claimedAt = null;
      message.lockedUntil = null;
      message.updatedAt = cloneDate(now);
    },

    async markFailed(input) {
      assertNonEmptyString("id", input.id);
      assertNonEmptyString("claimToken", input.claimToken);
      const message = getClaimedOrThrow(input.id, input.claimToken);
      const now = input.now ?? storeNow();

      message.status = input.deadLetter ? "deadLettered" : "pending";
      message.lastError = serializeOutboxError(input.error);
      message.availableAt = input.retryAt
        ? cloneDate(input.retryAt)
        : cloneDate(now);
      message.claimToken = null;
      message.claimedAt = null;
      message.lockedUntil = null;
      message.updatedAt = cloneDate(now);
    },

    clear() {
      messages.clear();
    },
  };
}

function mapDefinitions<T extends { name: string }>(
  kind: string,
  defs: readonly T[],
): ReadonlyMap<string, T> {
  const map = new Map<string, T>();
  for (const def of defs) {
    if (map.has(def.name)) {
      throw new OutboxRegistryError(
        `Duplicate ${kind} definition "${def.name}" in outbox registry.`,
      );
    }
    map.set(def.name, def);
  }
  return map;
}

/**
 * Define the events and jobs that an outbox drain worker can deliver.
 *
 * Duplicate names throw because message delivery resolves by name.
 */
export function defineOutboxRegistry(
  input: DefineOutboxRegistryInput,
): OutboxRegistry {
  return {
    events: mapDefinitions("event", input.events ?? []),
    jobs: mapDefinitions("job", input.jobs ?? []),
  };
}

/**
 * Validate an event payload and enqueue it as an outbox message.
 */
export async function enqueueEvent<E extends EventPayloadDef>(
  outbox: OutboxPort,
  event: E,
  payload: InferEventPayload<E>,
  options: EnqueueTypedOutboxOptions = {},
): Promise<OutboxMessage> {
  await parseEventPayload(event, payload);
  const trace =
    parseTraceCarrier(options.trace) ?? captureTraceCarrier(options.tracing);
  return outbox.enqueue({
    id: options.id,
    kind: "event",
    name: event.name,
    payload: toOutboxJsonValue(payload),
    trace,
    availableAt: options.availableAt,
    maxAttempts: options.maxAttempts,
  });
}

/**
 * Validate a job payload and enqueue it as an outbox message.
 */
export async function enqueueJob<J extends JobDef>(
  outbox: OutboxPort,
  job: J,
  payload: InferJobPayload<J>,
  options: EnqueueTypedOutboxOptions = {},
): Promise<OutboxMessage> {
  await parseJobPayload(job, payload);
  const trace =
    parseTraceCarrier(options.trace) ?? captureTraceCarrier(options.tracing);
  return outbox.enqueue({
    id: options.id,
    kind: "job",
    name: job.name,
    payload: toOutboxJsonValue(payload),
    trace,
    availableAt: options.availableAt,
    maxAttempts: options.maxAttempts ?? getJobRetryMaxAttempts(job.retry),
  });
}

/**
 * Create a domain event recorder that writes events to the outbox.
 */
export function createOutboxEventRecorder(
  outbox: OutboxPort,
  options: EnqueueTypedOutboxOptions = {},
): DomainEventRecorderPort {
  return {
    async record(event, payload, publishOptions) {
      await enqueueEvent(outbox, event, payload, {
        ...options,
        trace: publishOptions?.trace ?? options.trace,
      });
    },
  };
}

/**
 * Create a job dispatcher that writes jobs to the outbox.
 */
export function createOutboxJobDispatcher(
  outbox: OutboxPort,
  options: EnqueueTypedOutboxOptions = {},
): JobDispatcherPort {
  return {
    async dispatch(job, payload, dispatchOptions) {
      await enqueueJob(outbox, job, payload, {
        ...options,
        trace: dispatchOptions?.trace ?? options.trace,
      });
    },
  };
}

function resolveRetryDelayMs(
  options: DrainOutboxOptions,
  message: ClaimedOutboxMessage,
  error: unknown,
  now: Date,
): number {
  if (typeof options.retryDelayMs === "function") {
    const delay = options.retryDelayMs({ message, error, now });
    assertPositiveInteger("retryDelayMs", delay);
    return delay;
  }

  if (options.retryDelayMs !== undefined) {
    assertPositiveInteger("retryDelayMs", options.retryDelayMs);
    return options.retryDelayMs;
  }

  if (message.kind === "job") {
    const job = options.registry.jobs.get(message.name);
    if (job?.retry) {
      return getJobRetryDelayMs(job.retry, {
        attempt: message.attempts,
        error,
        jobName: message.name,
      });
    }
  }

  return Math.min(60_000, 1000 * 2 ** Math.max(0, message.attempts - 1));
}

function shouldRetryOutboxMessage(
  options: DrainOutboxOptions,
  message: ClaimedOutboxMessage,
  error: unknown,
): boolean {
  if (message.kind !== "job") {
    return message.attempts < message.maxAttempts;
  }

  const job = options.registry.jobs.get(message.name);
  return shouldRetryJob(job?.retry, {
    attempt: message.attempts,
    error,
    jobName: message.name,
    maxAttempts: message.maxAttempts,
  });
}

function outboxInstrumentationDetails(
  message: ClaimedOutboxMessage,
  details?: Record<string, unknown>,
): Record<string, unknown> {
  return {
    attempt: message.attempts,
    maxAttempts: message.maxAttempts,
    messageId: message.id,
    messageKind: message.kind,
    messageName: message.name,
    ...details,
  };
}

async function deliverOutboxMessage(
  options: DrainOutboxOptions,
  message: ClaimedOutboxMessage,
  trace?: TraceCarrier,
): Promise<void> {
  if (message.kind === "event") {
    if (!options.eventBus) {
      throw new OutboxRegistryError(
        `Cannot deliver event "${message.name}" without an event bus.`,
      );
    }

    const event = options.registry.events.get(message.name);
    if (!event) {
      throw new OutboxRegistryError(
        `Outbox registry does not include event "${message.name}".`,
      );
    }

    await parseEventPayload(event, message.payload);
    await options.eventBus.publish(
      event,
      message.payload as never,
      trace ? { trace } : undefined,
    );
    return;
  }

  if (!options.jobs) {
    throw new OutboxRegistryError(
      `Cannot deliver job "${message.name}" without a job dispatcher.`,
    );
  }

  const job = options.registry.jobs.get(message.name);
  if (!job) {
    throw new OutboxRegistryError(
      `Outbox registry does not include job "${message.name}".`,
    );
  }

  await parseJobPayload(job, message.payload);

  // The drain owns execution retries: failed deliveries are rescheduled with
  // the job's own policy via markFailed/retryAt. When the dispatcher exposes
  // a single-attempt dispatch (the inline dispatcher does), use it so the
  // retry policy runs in exactly one layer. Durable providers do not expose
  // it — for them dispatch is an enqueue and the queue owns execution.
  const singleAttempt = (
    options.jobs as {
      [SINGLE_ATTEMPT_DISPATCH]?: SingleAttemptJobDispatch;
    }
  )[SINGLE_ATTEMPT_DISPATCH];
  if (singleAttempt) {
    await singleAttempt(job, message.payload as never, {
      attempt: message.attempts,
      maxAttempts: message.maxAttempts,
      trace,
    });
    return;
  }

  await options.jobs.dispatch(
    job,
    message.payload as never,
    trace ? { trace } : undefined,
  );
}

/**
 * Claim and deliver one batch of outbox messages.
 *
 * This does not loop forever; production workers should call it on their own
 * polling cadence. Event and job messages require matching registry entries.
 * Failed messages are retried with backoff until `maxAttempts`, then
 * dead-lettered.
 */
export async function drainOutbox(
  options: DrainOutboxOptions,
): Promise<DrainOutboxResult> {
  const batchSize = options.batchSize ?? 100;
  assertPositiveInteger("batchSize", batchSize);
  const instrumentation = createProviderInstrumentation(
    options.instrumentation,
    {
      providerName: "outbox",
      watcher: "outbox",
    },
  );
  const jobInstrumentation = createProviderInstrumentation(
    options.instrumentation,
    {
      providerName: "outbox",
      watcher: "jobs",
    },
  );
  const tracing = resolveTracingPort(options.instrumentation);

  const now = options.now ?? new Date();
  const messages = await options.outbox.claimBatch({
    limit: batchSize,
    now,
    leaseMs: options.leaseMs,
  });
  const result: DrainOutboxResult = {
    claimed: messages.length,
    delivered: 0,
    retried: 0,
    deadLettered: 0,
  };

  for (const message of messages) {
    try {
      const parentTrace = parseTraceCarrier(message.trace);
      const traceAttributes = {
        "beignet.outbox.message_kind": message.kind,
        "beignet.outbox.message_name": message.name,
      } as const;
      await runWithTracing(
        tracing,
        {
          name: `beignet.outbox deliver ${message.name}`,
          type: "outbox",
          kind: "consumer",
          parent: parentTrace,
          attributes: traceAttributes,
          metricAttributes: traceAttributes,
        },
        (span) =>
          deliverOutboxMessage(
            options,
            message,
            captureTraceCarrier(span?.context ?? parentTrace),
          ),
      );
      await options.outbox.markDelivered({
        id: message.id,
        claimToken: message.claimToken,
        now,
      });
      instrumentation.record({
        type: "outbox",
        ...options.instrumentationContext,
        messageId: message.id,
        messageKind: message.kind,
        messageName: message.name,
        status: "delivered",
        details: outboxInstrumentationDetails(message),
      });
      result.delivered += 1;
    } catch (error) {
      try {
        await options.onError?.(error, message);
      } catch {
        // Preserve the delivery failure path so the message is retried or
        // dead-lettered even if the observer fails.
      }
      const shouldRetry = shouldRetryOutboxMessage(options, message, error);
      const deadLetter = !shouldRetry;
      const retryDelayMs = deadLetter
        ? 0
        : resolveRetryDelayMs(options, message, error, now);
      try {
        await options.outbox.markFailed({
          id: message.id,
          claimToken: message.claimToken,
          error,
          deadLetter,
          now,
          retryAt: deadLetter
            ? undefined
            : new Date(now.getTime() + retryDelayMs),
        });
      } catch (settlementError) {
        try {
          await options.onSettlementError?.(settlementError, message, error);
        } catch {
          // Preserve the settlement failure when its observer also fails.
        }
        instrumentation.custom({
          name: "outbox.settlement.failed",
          label: "Outbox settlement failed",
          summary: `Could not settle failed ${message.kind} "${message.name}"`,
          details: outboxInstrumentationDetails(message, {
            deliveryError: serializeOutboxError(error),
            settlementError: serializeOutboxError(settlementError),
          }),
        });
        continue;
      }

      if (deadLetter) {
        try {
          await options.onDeadLetter?.(error, message);
        } catch {
          // Dead-letter observers must not change the settled message state.
        }
        instrumentation.record({
          type: "outbox",
          ...options.instrumentationContext,
          messageId: message.id,
          messageKind: message.kind,
          messageName: message.name,
          status: "deadLettered",
          details: outboxInstrumentationDetails(message, {
            error: serializeOutboxError(error),
          }),
        });
        if (message.kind === "job") {
          jobInstrumentation.record({
            type: "job",
            ...options.instrumentationContext,
            jobName: message.name,
            status: "deadLettered",
            details: outboxInstrumentationDetails(message, {
              error: serializeOutboxError(error),
            }),
          });
        }
        result.deadLettered += 1;
      } else {
        const retryAt = new Date(now.getTime() + retryDelayMs).toISOString();
        instrumentation.record({
          type: "outbox",
          ...options.instrumentationContext,
          messageId: message.id,
          messageKind: message.kind,
          messageName: message.name,
          status: "retryScheduled",
          details: outboxInstrumentationDetails(message, {
            retryDelayMs,
            retryAt,
            error: serializeOutboxError(error),
          }),
        });
        if (message.kind === "job") {
          jobInstrumentation.record({
            type: "job",
            ...options.instrumentationContext,
            jobName: message.name,
            status: "retryScheduled",
            details: outboxInstrumentationDetails(message, {
              retryDelayMs,
              retryAt,
              error: serializeOutboxError(error),
            }),
          });
        }
        result.retried += 1;
      }
    }
  }

  return result;
}

/**
 * Domain event recorder port re-exported for outbox integrations.
 */
export type { DomainEventRecorderPort };
