import type { OperationDispatcher, OperationRecord, OperationStore } from '../types/operations.js'; /** * The Redis commands the operation store needs. * * Structural rather than tied to one client, so `ioredis`, `node-redis`, or a cluster proxy all * satisfy it without this package depending on any of them. */ export interface RedisOperationLikeClient { /** Reads a hash field. */ hget(key: string, field: string): Promise | string | null; /** Sets a hash field. */ hset(key: string, field: string, value: string): Promise | unknown; /** Deletes a hash field. */ hdel(key: string, field: string): Promise | unknown; /** Reads every hash value. */ hvals(key: string): Promise | string[]; /** Optional CAS primitive. When absent the store falls back to a read-compare-write. */ eval?(script: string, numKeys: number, ...args: string[]): Promise | unknown; } /** Options for the Redis operation store. */ export interface RedisOperationStoreOptions { /** Key prefix. Defaults to `nexus-ai-pro:operations:`. */ prefix?: string; /** Disables the Lua compare-and-set even when the client exposes `eval`. */ useEval?: boolean; } /** * Redis-backed operation storage, so a submitted operation survives a restart. * * Updates are a compare-and-set on `sequence`. When the client exposes `eval` the comparison and * the write happen in one Lua call, which is genuinely atomic; otherwise the store falls back to a * read-compare-write that narrows the race but cannot close it. Prefer a client with `eval` * whenever more than one worker can touch the same operation. */ export declare class RedisOperationStore implements OperationStore { private readonly client; private readonly prefix; private readonly useEval; constructor(client: RedisOperationLikeClient, options?: RedisOperationStoreOptions); /** Stores a new record and indexes its idempotency key. Refuses records carrying raw bytes. */ create(record: OperationRecord): Promise; /** Reads a record. */ read(id: string): Promise | undefined>; /** * Writes a record when its stored sequence still equals `expectedSequence`. Resolves false when * another worker got there first. */ update(record: OperationRecord, expectedSequence: number): Promise; /** Deletes a record and its idempotency index. Resolves true when it existed. */ delete(id: string): Promise; /** * Records whose lease has expired, or running records without one, up to `limit`, for another * worker to take over. */ claimExpired(now: string, limit: number): Promise>>; /** Finds the record that claimed an idempotency key. */ findByIdempotencyKey(key: string): Promise | undefined>; /** Every record. */ list(): Promise>>; private recordsKey; private idempotencyKey; } /** The part of a BullMQ `Queue` the dispatcher needs. */ export interface BullMQLikeOperationQueue { /** Adds a job. */ add(name: string, data: unknown, options?: Record): Promise<{ id?: string | number; }> | { id?: string | number; }; } /** Options for the BullMQ operation dispatcher. */ export interface BullMQOperationDispatcherOptions { /** Job name used for every dispatched operation. Defaults to `nexus-operation`. */ jobName?: string; /** Merged into the BullMQ job options, for priority, delay, or removal policy. */ jobOptions?: Record; } /** * Hands accepted operations to a BullMQ queue for a worker process to execute. * * Only the operation id and its routing metadata are queued — never the record's result or any * payload — so the job stays small and the store remains the single source of truth. The worker * reads the record by id, which is also what makes a redelivered job safe. */ export declare class BullMQOperationDispatcher implements OperationDispatcher { private readonly queue; private readonly options; constructor(queue: BullMQLikeOperationQueue, options?: BullMQOperationDispatcherOptions); /** * Queues an operation's id and routing metadata, using the operation id as the job id so a * duplicate dispatch is ignored. */ dispatch(record: OperationRecord): Promise; }