import BackgroundJobsAdapter from "./adapter.js"; import Logger from "../logger.js"; export type PreparedBackgroundJob = { /** * - Serialized arguments. */ argsJson: string; /** * - Resolved concurrency. */ concurrency: { concurrencyKey: string; maxConcurrency: number; queueDerived: boolean; } | null; /** * - Creation timestamp. */ createdAtMs: number; /** * - Execution mode. */ executionMode: import("./types.js").BackgroundJobExecutionMode; /** * - New job id. */ jobId: string; /** * - Job name. */ jobName: string; /** * - Retry cap. */ maxRetries: number; /** * - Queue name. */ queue: string; /** * - Eligibility timestamp. */ scheduledAtMs: number; /** * - Per-job timeout override, or null when omitted. */ timeoutMs: number | null; }; export type BackgroundJobOrphanSelection = { /** * - Exact update fence. */ conditions: Record>; /** * - Selected active handoff. */ job: import("./types.js").BackgroundJobRow; }; export type BackgroundJobTransactionSerializationOptions = { /** * - Session lock held around the transaction. */ advisoryLock?: { failureMessage: string; name: string; }; }; export type BackgroundJobPruneCandidateMetadata = { /** * - Job ids selected for this batch. */ candidates: Readonly>; /** * - Terminal status the candidates were selected by. */ status: string; /** * - Terminal timestamp column compared against the cutoff. */ column: string; /** * - Cutoff timestamp the candidates were selected against. */ cutoff: number; }; export type BackgroundJobConcurrencyCountRow = { /** * - Persisted or aggregated active count. */ active_count: number | string; /** * - Durable cap identity. */ concurrency_key: string; }; export type BackgroundJobQueuedConcurrency = { /** * - Current concurrency key for queued work. */ concurrencyKey: string | null; /** * - Current concurrency cap for queued work. */ maxConcurrency: number | null; }; export declare const BACKGROUND_JOB_COUNTS_CHANNEL = "velocious-background-job-counts"; export declare const BACKGROUND_JOB_COUNT_BUCKETS: string[]; export default class BackgroundJobsStore extends BackgroundJobsAdapter { configuration: import("../configuration.js").default; databaseIdentifier: string | undefined; clock: { now: () => number; }; afterOwnedProducerValidation: ((producerProof: import("./types.js").BackgroundJobProducerProof) => void | Promise) | undefined; afterPruneCandidatesSelected: ((metadata: BackgroundJobPruneCandidateMetadata) => void | Promise) | undefined; logger: Logger; _readyPromise: Promise | null; _queueConcurrencyReconciled: boolean; /** * Runs constructor. * @param {object} args - Options. * @param {import("../configuration.js").default} args.configuration - Configuration. * @param {string} [args.databaseIdentifier] - Database identifier. * @param {{now: () => number}} [args.clock] - Injectable persistence clock. * @param {(producerProof: import("./types.js").BackgroundJobProducerProof) => void | Promise} [args.afterOwnedProducerValidation] - Exact owned-enqueue validation hook. * @param {(metadata: BackgroundJobPruneCandidateMetadata) => void | Promise} [args.afterPruneCandidatesSelected] - Optional barrier invoked after one retention batch's candidate discovery and before its serialized delete transaction. */ constructor({ configuration, databaseIdentifier, clock, afterOwnedProducerValidation, afterPruneCandidatesSelected }: { configuration: import("../configuration.js").default; databaseIdentifier?: string; clock?: { now: () => number; }; afterOwnedProducerValidation?: (producerProof: import("./types.js").BackgroundJobProducerProof) => void | Promise; afterPruneCandidatesSelected?: (metadata: BackgroundJobPruneCandidateMetadata) => void | Promise; }); /** * Runs get database identifier. * @returns {string} - Database identifier. */ getDatabaseIdentifier(): string; /** * Runs ensure ready. * @returns {Promise} - Resolves when ready. */ ensureReady(): Promise; /** * Ensures the background-jobs schema (tables + columns) exists on the configured * database, without initializing the runtime model. Lets `db:migrate` create the * framework's own schema deterministically alongside app migrations — and capture * it in the dumped structure SQL — instead of it only appearing once a store boots. * Idempotent: reuses the same `_ensureSchema` the runtime store uses, which skips * work already applied (tracked in `velocious_internal_migrations`). * @param {import("../database/drivers/base.js").default} [db] - Reuse an already * checked-out connection (e.g. the one `db:migrate` holds) rather than opening a * nested checkout that would deadlock a single-connection pool. * @returns {Promise} - Resolves when the schema is present. */ ensureSchema(db?: import("../database/drivers/base.js").default): Promise; /** * Reconciles queue-derived concurrency with the current configuration: the * explicit lifecycle path that adopts/releases persisted queued jobs onto * queue concurrency keys when `queues[name].maxConcurrent` is added, removed, * or changed. Called by the background-jobs main process on startup — the * deploy-time moment queue configuration changes take effect. Schema/tenant * checks and routine connection initialization deliberately never run this: * they stay read-only regarding queued job rows, because the broad * adoption/release UPDATEs deadlock against active job processes under * concurrent tenant initialization. Serialized across processes with a * database advisory lock so concurrently started mains cannot interleave the * UPDATEs; the per-instance memo only skips repeat work within this process. * @returns {Promise} - Resolves when reconciled. */ reconcileQueueConcurrency(): Promise; /** * Repairs durable active-count drift while a main process remains live. The * initial snapshot is read-only; only suspected mismatches take their * counter lock and re-count inside the serialized transaction path. * @returns {Promise} - Repair summary. */ reconcileActiveConcurrency(): Promise; /** * Runs enqueue. * @param {object} args - Options. * @param {string} args.jobName - Job name. * @param {Array>} args.args - Arguments. * @param {import("./types.js").BackgroundJobOptions} [args.options] - Options. * @returns {Promise} - Job id. */ enqueue({ jobName, args, options }: { jobName: string; args: Array>; options?: import("./types.js").BackgroundJobOptions; }): Promise; /** * Atomically validates an exact producing handoff and enqueues its follow-up. * Every exact request owns an internal durable replay identity, while queued * deduplication can point several distinct producer events at one covering row. * @param {object} args - Owned enqueue request. * @param {string} args.jobName - Job name. * @param {Array>} args.args - Arguments. * @param {import("./types.js").BackgroundJobOptions} [args.options] - Options. * @param {string} [args.producerInvocationId] - Stable identity for one owned enqueue invocation. * @param {import("./types.js").BackgroundJobProducerProof} args.producerProof - Exact producer lease. * @returns {Promise} - Durable follow-up id. */ enqueueFromOwnedHandoff({ jobName, args, options, producerInvocationId, producerProof }: { jobName: string; args: Array>; options?: import("./types.js").BackgroundJobOptions; producerInvocationId?: string; producerProof: import("./types.js").BackgroundJobProducerProof; }): Promise; /** * Finds the earliest queued job that covers this enqueue's identity and time. * @param {import("../database/drivers/base.js").default} db - Transaction connection. * @param {PreparedBackgroundJob} preparedJob - Normalized job. * @returns {Promise} - Covering job id. */ _deduplicatedQueuedJobId(db: import("../database/drivers/base.js").default, preparedJob: PreparedBackgroundJob): Promise; /** * Persists one internal exact-replay owner and its queued job in the caller's * producer-validation transaction. * @param {object} args - Transaction input. * @param {import("../database/drivers/base.js").default} args.db - Transaction connection. * @param {import("./types.js").BackgroundJobOptions} args.options - Enqueue options. * @param {PreparedBackgroundJob} args.preparedJob - Normalized job. * @param {string} args.producerInvocationId - Stable identity for one owned enqueue invocation. * @param {import("./types.js").BackgroundJobProducerProof} args.producerProof - Exact producer lease. * @returns {Promise} - Stable replay job id. */ _enqueueOwnedReplayInTransaction({ db, options, preparedJob, producerInvocationId, producerProof }: { db: import("../database/drivers/base.js").default; options: import("./types.js").BackgroundJobOptions; preparedJob: PreparedBackgroundJob; producerInvocationId: string; producerProof: import("./types.js").BackgroundJobProducerProof; }): Promise; /** * Atomically owns one durable idempotency scope and creates its job exactly once. * @param {object} args - Enqueue input. * @param {Array>} args.args - Job arguments. * @param {import("./types.js").BackgroundJobOptions} args.options - Job options. * @param {PreparedBackgroundJob} args.preparedJob - Normalized job. * @returns {Promise} - Stable original job id. */ _enqueueIdempotently({ args, options, preparedJob }: { args: Array>; options: import("./types.js").BackgroundJobOptions; preparedJob: PreparedBackgroundJob; }): Promise; /** * Owns or replays one public idempotency key inside the caller's transaction. * @param {object} args - Transaction input. * @param {Array>} args.args - Job arguments. * @param {boolean} [args.countRevisionLocked] - Whether the caller already owns count serialization. * @param {import("../database/drivers/base.js").default} args.db - Transaction connection. * @param {import("./types.js").BackgroundJobOptions} args.options - Job options. * @param {PreparedBackgroundJob} args.preparedJob - Normalized job. * @returns {Promise} - Stable original job id. */ _enqueueIdempotentlyInTransaction({ args, countRevisionLocked, db, options, preparedJob }: { args: Array>; countRevisionLocked?: boolean; db: import("../database/drivers/base.js").default; options: import("./types.js").BackgroundJobOptions; preparedJob: PreparedBackgroundJob; }): Promise; /** * Serializes one physical connection locally without taking ownership away * from the database uniqueness constraint shared by all processes. * @template T * @param {(db: import("../database/drivers/base.js").default) => Promise} callback - Transaction work. * @returns {Promise} - Callback result. */ _idempotentEnqueueTransaction(callback: (db: import("../database/drivers/base.js").default) => Promise): Promise; /** * Inserts an ownership row, resolving only a database uniqueness race. * @param {import("../database/drivers/base.js").default} db - Transaction connection. * @param {Record>} ownership - Ownership row. * @returns {Promise<{created: boolean, row: Record>}>} - Claim result. */ _claimIdempotencyOwnership(db: import("../database/drivers/base.js").default, ownership: Record>): Promise<{ created: boolean; row: Record>; }>; /** * Loads one durable enqueue owner. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} scopeDigest - Fixed-size scope digest. * @returns {Promise> | null>} - Row or null. */ _idempotencyOwnership(db: import("../database/drivers/base.js").default, scopeDigest: string): Promise> | null>; /** * Fails closed when a durable key is reused for a different canonical request. * @param {object} args - Validation input. * @param {Record>} args.existing - Stored owner. * @param {Record>} args.ownership - Requested owner. * @returns {void} */ _validateIdempotencyOwnership({ existing, ownership }: { existing: Record>; ownership: Record>; }): void; /** * Persists the built-in mail operation in the same first-enqueue transaction. * @param {import("../database/drivers/base.js").default} db - Transaction connection. * @param {object} args - Operation input. * @param {number} args.createdAtMs - Creation timestamp. * @param {string} args.jobId - Native job id. * @param {{operation: import("../mailer/index.js").MailerDeliveryOperation, payload: import("../mailer/index.js").MailerDeliveryPayload} | null} args.mailOperationInput - Mail operation. * @returns {Promise} - Resolves after persistence. */ _persistMailDeliveryOperation(db: import("../database/drivers/base.js").default, { createdAtMs, jobId, mailOperationInput }: { createdAtMs: number; jobId: string; mailOperationInput: { operation: import("../mailer/index.js").MailerDeliveryOperation; payload: import("../mailer/index.js").MailerDeliveryPayload; } | null; }): Promise; /** * Validates the durable mail row during an exact generic enqueue replay. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {object} args - Validation input. * @param {string} args.jobId - Owned job id. * @param {{operation: import("../mailer/index.js").MailerDeliveryOperation, payload: import("../mailer/index.js").MailerDeliveryPayload} | null} args.mailOperationInput - Mail operation. * @returns {Promise} - Resolves when exact. */ _validateMailDeliveryOperation(db: import("../database/drivers/base.js").default, { jobId, mailOperationInput }: { jobId: string; mailOperationInput: { operation: import("../mailer/index.js").MailerDeliveryOperation; payload: import("../mailer/index.js").MailerDeliveryPayload; } | null; }): Promise; /** * Loads a durable mail operation. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} operationKey - Fixed-size operation key. * @returns {Promise> | null>} - Row or null. */ _mailDeliveryOperation(db: import("../database/drivers/base.js").default, operationKey: string): Promise> | null>; /** * Compares provider-relevant durable mail operation fields. * @param {object} args - Validation input. * @param {Record>} args.existing - Stored row. * @param {Record>} args.requested - Requested row. * @returns {void} */ _validateMailDeliveryOperationRow({ existing, requested }: { existing: Record>; requested: Record>; }): void; /** * Canonical request digest excluding generated ids and immediate enqueue time. * @param {object} args - Digest input. * @param {Array>} args.args - Job arguments. * @param {import("./types.js").BackgroundJobOptions} args.options - Job options. * @param {PreparedBackgroundJob} args.preparedJob - Normalized job. * @returns {string} - SHA-256 digest. */ _idempotencyRequestDigest({ args, options, preparedJob }: { args: Array>; options: import("./types.js").BackgroundJobOptions; preparedJob: PreparedBackgroundJob; }): string; /** * Fixed-size globally indexed representation of the documented scope tuple. * @param {object} args - Scope input. * @param {string} args.idempotencyKey - Caller key. * @param {string} args.jobName - Job class name. * @param {string} args.queue - Queue name. * @returns {string} - SHA-256 scope digest. */ _idempotencyScopeDigest({ idempotencyKey, jobName, queue }: { idempotencyKey: string; jobName: string; queue: string; }): string; /** * Validates one caller key. * @param {string | undefined} idempotencyKey - Caller key. * @returns {string} - Valid key. */ _normalizeIdempotencyKey(idempotencyKey: string | undefined): string; /** * Canonical request identity for an internal owned-handoff replay. * Immediate enqueue wall time and generated job ids remain excluded. * @param {object} args - Digest input. * @param {import("./types.js").BackgroundJobOptions} args.options - Enqueue options. * @param {PreparedBackgroundJob} args.preparedJob - Normalized job. * @returns {string} - SHA-256 digest. */ _ownedEnqueueRequestDigest({ options, preparedJob }: { options: import("./types.js").BackgroundJobOptions; preparedJob: PreparedBackgroundJob; }): string; /** * Isolates internal producer replay ownership from caller idempotency scopes. * @param {object} args - Scope input. * @param {PreparedBackgroundJob} args.preparedJob - Normalized job. * @param {string} args.producerInvocationId - Stable identity for one owned enqueue invocation. * @param {import("./types.js").BackgroundJobProducerProof} args.producerProof - Exact producer lease. * @param {string} args.requestDigest - Canonical request digest. * @returns {string} - SHA-256 scope digest. */ _ownedEnqueueScopeDigest({ preparedJob, producerInvocationId, producerProof, requestDigest }: { preparedJob: PreparedBackgroundJob; producerInvocationId: string; producerProof: import("./types.js").BackgroundJobProducerProof; requestDigest: string; }): string; /** * Validates the untrusted identity of one producer-owned enqueue invocation. * @param {string | undefined} producerInvocationId - Producer invocation identity. * @returns {string} - Validated identity. */ _normalizeProducerInvocationId(producerInvocationId: string | undefined): string; /** * Validates the untrusted transport shape before transaction admission. * @param {import("./types.js").BackgroundJobProducerProof} producerProof - Producer proof. * @returns {import("./types.js").BackgroundJobProducerProof} - Normalized immutable proof. */ _normalizeProducerProof(producerProof: import("./types.js").BackgroundJobProducerProof): import("./types.js").BackgroundJobProducerProof; /** * Confirms exact active ownership while the enqueue transaction holds the * shared mutation fence used by terminal producer transitions. * @param {import("../database/drivers/base.js").default} db - Transaction connection. * @param {import("./types.js").BackgroundJobProducerProof} producerProof - Exact producer lease. * @returns {Promise} - Resolves while ownership remains exact. */ _validateOwnedProducerProof(db: import("../database/drivers/base.js").default, producerProof: import("./types.js").BackgroundJobProducerProof): Promise; /** * Replaces the queued owner of a stable schedule key with a new one-off job. * A handed-off owner is left running and reported truthfully. * @param {object} args - Options. * @param {string} args.scheduleKey - Stable logical schedule key. * @param {string} args.jobName - Job name. * @param {Array>} args.args - Arguments. * @param {import("./types.js").BackgroundJobOptions} [args.options] - Options. * @returns {Promise} - Replacement result. */ replaceScheduled({ scheduleKey, jobName, args, options }: { scheduleKey: string; jobName: string; args: Array>; options?: import("./types.js").BackgroundJobOptions; }): Promise; /** * Cancels the queued owner of a stable schedule key. A handed-off owner is * detached but not marked stopped because execution may already be running. * @param {string} scheduleKey - Stable logical schedule key. * @returns {Promise} - Cancellation result. */ cancelScheduled(scheduleKey: string): Promise; /** * Reads stable ownership and optional latest terminal history in one fenced transaction. * @param {string} scheduleKey - Stable logical schedule key. * @param {{includeLatestTerminal?: boolean}} [options] - Lookup options. * @returns {Promise} - Normalized public jobs. */ getScheduledJob(scheduleKey: string, { includeLatestTerminal }?: { includeLatestTerminal?: boolean; }): Promise; /** * Moves only a future queued stable owner to the current time. * @param {string} scheduleKey - Stable logical schedule key. * @returns {Promise} - Exact wake outcome. */ wakeScheduled(scheduleKey: string): Promise; /** * Runs next available job. * @param {object} [args] - Options. * @param {import("./types.js").BackgroundJobExecutionMode | import("./types.js").BackgroundJobExecutionMode[]} [args.executionMode] - Execution mode or modes to match. * @returns {Promise} - Next job. */ nextAvailableJob(args?: { executionMode?: import("./types.js").BackgroundJobExecutionMode | import("./types.js").BackgroundJobExecutionMode[]; }): Promise; /** * Returns the soonest future-scheduled queued job (one whose * `scheduled_at_ms` is in the future), or null when there are no * future-scheduled jobs. Used by the event-driven dispatcher to arm a * `setTimeout` for the exact moment the next scheduled job becomes * eligible, replacing the legacy 1-second polling loop. * @returns {Promise} - Soonest future-scheduled job, or null. */ nextScheduledJob(): Promise; /** * Runs next queued job. * @param {object} args - Options. * @param {import("../database/drivers/base.js").default} args.db - Database connection. * @param {"<=" | ">"} args.scheduledAtOperator - Scheduled timestamp operator. * @param {import("./types.js").BackgroundJobExecutionMode | import("./types.js").BackgroundJobExecutionMode[]} [args.executionMode] - Execution mode or modes to match. * @returns {Promise} - Next matching queued job. */ _nextQueuedJob({ db, scheduledAtOperator, executionMode }: { db: import("../database/drivers/base.js").default; scheduledAtOperator: "<=" | ">"; executionMode?: import("./types.js").BackgroundJobExecutionMode | import("./types.js").BackgroundJobExecutionMode[]; }): Promise; /** * Builds a raw SQL ORDER BY expression ranking queued jobs by their queue's * configured priority (`backgroundJobs.queues[queue].priority`, default `0`), * so the dispatcher picks higher-priority queues first regardless of enqueue * order. Only applied to the dispatch path (`scheduledAtOperator === "<="`); * the future-scheduled lookup must stay strictly time-ordered. Composes with * the concurrency EXISTS filter: a higher-priority queue already at its cap is * filtered out, so dispatch falls through to the next eligible lower-priority * job. Returns null when no queue configures a non-zero priority so the plain * FIFO ordering is left untouched (and no needless filesort is introduced). * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {string | null} - Raw SQL CASE expression, or null when no queue is prioritized. */ _queuePriorityOrderSql(db: import("../database/drivers/base.js").default): string | null; /** * Runs get job. * @param {string} jobId - Job id. * @returns {Promise} - Job row. */ getJob(jobId: string): Promise; /** * Counts jobs grouped by status. Used by the dashboard overview. * @returns {Promise>} - Counts keyed by status. */ countsByStatus(): Promise>; /** * Returns the authoritative dashboard count snapshot and its matching durable * revision. Locking the revision row before counting prevents a writer from * committing between the count query and revision read. * @returns {Promise<{counts: Record, revision: number, total: number}>} Snapshot. */ countSnapshot(): Promise<{ counts: Record; revision: number; total: number; }>; /** * Counts jobs matching the given filters. * @param {object} [args] - Options. * @param {string} [args.status] - Filter by status. * @param {string} [args.jobName] - Filter by job name. * @returns {Promise} - Matching job count. */ countJobs({ status, jobName }?: { status?: string; jobName?: string; }): Promise; /** * Lists jobs for the dashboard, filtered, sorted and paginated. * @param {object} [args] - Options. * @param {string} [args.status] - Filter by status. * @param {string} [args.jobName] - Filter by job name. * @param {number} [args.limit] - Maximum rows to return. * @param {number} [args.offset] - Rows to skip. * @param {string} [args.sortColumn] - Camel-cased column to sort by (see SORTABLE_COLUMNS). * @param {"ASC" | "DESC"} [args.sortDirection] - Sort direction. * @returns {Promise} - Normalized job rows. */ listJobs({ status, jobName, limit, offset, sortColumn, sortDirection }?: { status?: string; jobName?: string; limit?: number; offset?: number; sortColumn?: string; sortDirection?: "ASC" | "DESC"; }): Promise; /** * Runs mark handed off. * @param {object} args - Options. * @param {string} args.jobId - Job id. * @param {string} [args.handoffId] - Caller-selected exact lease id. Generated for legacy direct callers when omitted. * @param {string} [args.workerId] - Worker id. * @returns {Promise} - Claimed handoff lease, or null when no longer queued. */ markHandedOff({ jobId, handoffId, workerId }: { jobId: string; handoffId?: string; workerId?: string; }): Promise; /** * Runs mark completed. * @param {object} args - Options. * @param {string} args.jobId - Job id. * @param {string} [args.handoffId] - Handoff lease id. * @param {string} [args.workerId] - Worker id. * @param {number} [args.handedOffAtMs] - Handed off timestamp. * @returns {Promise} - Whether the fenced report was accepted. */ markCompleted({ jobId, handoffId, workerId, handedOffAtMs }: { jobId: string; handoffId?: string; workerId?: string; handedOffAtMs?: number; }): Promise; /** * Records pooled-child acceptance evidence for an active handoff: when the * executing runner child received and/or started the job, plus that child's * stable identity and pid. Only the fields supplied are written, so a * received-then-started observation lands as two fenced partial updates. The * update is fenced by the exact active handoff lease, so a report for a * reclaimed or re-handed-off job is dropped instead of stamping the wrong * attempt. * @param {object} args - Options. * @param {string} args.jobId - Job id. * @param {string} [args.handoffId] - Handoff lease id. * @param {string} [args.workerId] - Worker id. * @param {number} [args.handedOffAtMs] - Handed off timestamp. * @param {number} [args.receivedAtMs] - Epoch ms the runner child received the job. * @param {number} [args.startedAtMs] - Epoch ms the job's perform started in the child. * @param {string} [args.childInstanceId] - Stable pooled child identity. * @param {number} [args.childPid] - Pooled child OS pid. * @returns {Promise} - Whether the fenced report was accepted. */ markChildAccepted({ jobId, handoffId, workerId, handedOffAtMs, receivedAtMs, startedAtMs, childInstanceId, childPid }: { jobId: string; handoffId?: string; workerId?: string; handedOffAtMs?: number; receivedAtMs?: number; startedAtMs?: number; childInstanceId?: string; childPid?: number; }): Promise; /** * Returns the database data that clears pooled-child acceptance evidence. * @returns {Record>} - Cleared acceptance columns. */ _clearedChildAcceptanceData(): Record>; /** * Returns the row-shape counterpart of the cleared acceptance columns. * @returns {Pick} - Cleared acceptance fields. */ _clearedChildAcceptanceRow(): Pick; /** * Returns an active handoff to the queue at a caller-requested future time. * This is normal job control flow: it preserves failure attempts and metadata. * @param {object} args - Options. * @param {string} args.jobId - Job id. * @param {number} args.delayMs - Delay from persistence time in milliseconds. * @param {string} [args.handoffId] - Handoff lease id. * @param {string} [args.workerId] - Worker id. * @param {number} [args.handedOffAtMs] - Handed off timestamp. * @returns {Promise} - Whether the fenced report was accepted. */ markRescheduled({ jobId, delayMs, handoffId, workerId, handedOffAtMs }: { jobId: string; delayMs: number; handoffId?: string; workerId?: string; handedOffAtMs?: number; }): Promise; /** * Runs mark returned to queue. * @param {object} args - Options. * @param {string} args.jobId - Job id. * @param {string} args.handoffId - Handoff lease id. * @returns {Promise} - Resolves when updated. */ markReturnedToQueue({ jobId, handoffId }: { jobId: string; handoffId: string; }): Promise; /** * Returns the active `handed_off` jobs (jobId + handoffId) held under a worker * id. Used on worker reconnect: after a main restart a worker reconnects with * its stable id, and the fresh main adopts these leases so they are tracked — * and released if the reconnected worker later disconnects — instead of * sitting stuck until the age-based orphan sweep. This never reclaims, so a * gracefully-draining worker that keeps running its in-flight jobs is left * untouched. Rows with a null handoff id (legacy) are skipped; the orphan * sweep reclaims those via its `handed_off_at_ms` fence. * @param {object} args - Options. * @param {string} args.workerId - Worker id. * @returns {Promise>} - Active handoffs. */ handedOffJobsForWorker({ workerId }: { workerId: string; }): Promise>; /** * Snapshots exact, lease-aware active handoffs before a new main generation * starts accepting worker reconnects. Legacy rows without a complete worker, * lease, and timestamp identity stay owned by the age-based orphan sweep. * @returns {Promise} - Exact startup handoffs. */ snapshotHandedOffJobs(): Promise; /** * Reclaims only unchanged exact handoffs selected by a main-generation startup * snapshot. The ordinary orphan failure path owns retries, terminal status, * count transitions, schedule ownership, and concurrency release. * @param {object} args - Options. * @param {import("./types.js").BackgroundJobHandoffSnapshot[]} args.handoffs - Exact startup snapshots. * @param {ReturnType} args.error - Orphan reason. * @returns {Promise} - Accepted transitions. */ markOrphanedHandoffs({ handoffs, error }: { handoffs: import("./types.js").BackgroundJobHandoffSnapshot[]; error: ReturnType; }): Promise; /** * Runs mark failed. * @param {object} args - Options. * @param {string} args.jobId - Job id. * @param {ReturnType} args.error - Error. * @param {string} [args.handoffId] - Handoff lease id. * @param {string} [args.workerId] - Worker id. * @param {number} [args.handedOffAtMs] - Handed off timestamp. * @returns {Promise} - Updated job row when the report was accepted. */ markFailed({ jobId, error, handoffId, workerId, handedOffAtMs }: { jobId: string; error: ReturnType; handoffId?: string; workerId?: string; handedOffAtMs?: number; }): Promise; /** * Runs mark orphaned jobs. * @param {object} [args] - Options. * @param {number} [args.orphanedAfterMs] - Mark jobs orphaned after this duration. * @returns {Promise} - The jobs this sweep marked orphaned. */ markOrphanedJobs({ orphanedAfterMs }?: { orphanedAfterMs?: number; }): Promise; /** * Applies the common fenced orphan transition and records one aggregate count * delta for the accepted rows. * @param {object} args - Options. * @param {import("../database/drivers/base.js").default} args.db - Transaction connection. * @param {ReturnType} args.error - Orphan reason. * @param {BackgroundJobOrphanSelection[]} args.selections - Selected handoffs and exact fences. * @returns {Promise} - Accepted transitions. */ _markOrphanSelections({ db, error, selections }: { db: import("../database/drivers/base.js").default; error: ReturnType; selections: BackgroundJobOrphanSelection[]; }): Promise; /** * Deletes terminal job rows past their retention window so the jobs table * does not grow unbounded (completed rows in particular accumulate forever * otherwise). Batched by id — SELECT a page of ids, then * `DELETE ... WHERE id IN (...)` — rather than `DELETE ... LIMIT`, which not * every driver supports; each batch runs on its own connection so the sweep * yields between batches instead of holding one long transaction. * @param {object} [args] - Options. * @param {number | null} [args.completedTtlMs] - Delete `completed` jobs whose `completed_at_ms` is older than this many ms. Falsy or `<= 0` disables completed pruning. * @param {number | null} [args.failedTtlMs] - Delete terminal `failed`/`orphaned` jobs older than this many ms (by `failed_at_ms`/`orphaned_at_ms`). Falsy or `<= 0` disables. * @param {number} [args.batchSize] - Max rows deleted per batch. Default `1000`. * @returns {Promise} - Total rows deleted. */ pruneTerminalJobs({ completedTtlMs, failedTtlMs, batchSize }?: { completedTtlMs?: number | null; failedTtlMs?: number | null; batchSize?: number; }): Promise; /** * Deletes rows of one terminal status older than a cutoff, batch by batch, * until a candidate page returns fewer than `batchSize` rows. Candidate * discovery runs on a plain connection — never inside the serialized count * mutation — so a long scan cannot hold the count-revision lock and starve * enqueue acknowledgements; only the short delete transaction is serialized. * The delete revalidates status and cutoff for the selected ids, publishes * the delta from the actual affected-row count, and a page whose candidates * were already removed by a concurrent pruner still ends the pass only when * the page itself is short. * @param {object} args - Options. * @param {string} args.status - Terminal status to prune. * @param {string} args.column - Timestamp column compared against the cutoff. * @param {number} args.cutoff - Delete rows whose column value is `<= cutoff`. * @param {number} args.batchSize - Max rows per batch. * @returns {Promise} - Rows deleted for this status. */ _pruneStatusBatches({ status, column, cutoff, batchSize }: { status: string; column: string; cutoff: number; batchSize: number; }): Promise; /** * Runs clear all. * @returns {Promise} - Resolves when cleared. */ clearAll(): Promise; /** * Cancels a queued or handed-off job and releases any durable concurrency reservation. * @param {string} jobId - Job id. * @returns {Promise} - Whether the job was cancelled. */ cancel(jobId: string): Promise; /** * Runs get retry delay ms. * @param {number} retryCount - Retry attempt count (1-based). * @returns {number} - Delay in milliseconds. */ getRetryDelayMs(retryCount: number): number; /** * Normalizes one new job before entering its persistence transaction. * @param {object} args - Job input. * @param {Array>} args.args - Job arguments. * @param {string} args.jobName - Job name. * @param {import("./types.js").BackgroundJobOptions} [args.options] - Job options. * @returns {PreparedBackgroundJob} - Prepared job. */ _prepareJob({ args, jobName, options }: { args: Array>; jobName: string; options?: import("./types.js").BackgroundJobOptions; }): PreparedBackgroundJob; /** * Normalizes a per-job timeout while preserving omitted (worker fallback) * separately from explicitly disabled. * @param {import("./types.js").BackgroundJobOptions | undefined} options - Job options. * @returns {number | null} - Positive timeout, zero for disabled, or null when omitted. */ _normalizeJobTimeoutMs(options: import("./types.js").BackgroundJobOptions | undefined): number | null; /** * Inserts one prepared queued job, including its concurrency registration. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {object} args - Insert input. * @param {PreparedBackgroundJob} args.preparedJob - Prepared job. * @param {string | null} args.scheduleKey - Historical stable key. * @param {number | null} [args.scheduleOrder] - Monotonic stable ownership order. * @returns {Promise} - Resolves after insertion. */ _insertPreparedJob(db: import("../database/drivers/base.js").default, { preparedJob, scheduleKey, scheduleOrder }: { preparedJob: PreparedBackgroundJob; scheduleKey: string | null; scheduleOrder?: number | null; }): Promise; /** * Runs normalize max retries. * @param {number | null | undefined} maxRetries - Input. * @returns {number} - Normalized max retries. */ _normalizeMaxRetries(maxRetries: number | null | undefined): number; /** * Runs normalize scheduled at ms. * @param {number | undefined} scheduledAtMs - Requested dispatch timestamp. * @param {number} defaultScheduledAtMs - Default dispatch timestamp. * @returns {number} - Dispatch timestamp. */ _normalizeScheduledAtMs(scheduledAtMs: number | undefined, defaultScheduledAtMs: number): number; /** * Resolves a reschedule delay against persistence time. * @param {number} delayMs - Delay in milliseconds. * @returns {number} - Future eligibility timestamp. */ _rescheduledAtMs(delayMs: number): number; /** * Validates a public reschedule delay before persistence work begins. * @param {number} delayMs - Delay in milliseconds. * @returns {void} */ _validateRescheduleDelayMs(delayMs: number): void; /** * Validates a stable schedule key at the public storage boundary. * @param {string} scheduleKey - Stable logical schedule key. * @returns {string} - Validated key. */ _normalizeScheduleKey(scheduleKey: string): string; /** * Builds a bounded advisory-lock name for one stable schedule key. * @param {string} scheduleKey - Validated stable schedule key. * @returns {string} - Advisory-lock name. */ _scheduleKeyLockName(scheduleKey: string): string; /** * Ensures the background-jobs schema exists, reusing a caller-held connection when * one is given rather than checking out its own. * @param {import("../database/drivers/base.js").default} [existingDb] - Reuse an * already-checked-out connection (e.g. the one `db:migrate` holds) instead of * checking out a nested one — the nested checkout would deadlock a database * whose pool is capped at a single connection already held by the caller. * @returns {Promise} - Resolves when the schema is present. */ _ensureSchema(existingDb?: import("../database/drivers/base.js").default): Promise; /** * Serializes creation or upgrade of the background-jobs schema, checking out a * connection only after earlier schema work has completed when one is not supplied. * @param {import("../database/drivers/base.js").default} [existingDb] - Caller-owned * database connection. * @returns {Promise} - Resolves when the schema is present. */ _applySchema(existingDb?: import("../database/drivers/base.js").default): Promise; /** * Creates or upgrades the background-jobs tables, columns and concurrency rows on * the given connection. Serialized per process by {@link BackgroundJobsStore#_applySchema}. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when the schema is present. */ _applySchemaSteps(db: import("../database/drivers/base.js").default): Promise; /** * Runs ensure migrations table. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when complete. */ _ensureMigrationsTable(db: import("../database/drivers/base.js").default): Promise; /** * Runs has migration. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} [version] - Migration version. * @returns {Promise} - Whether migration exists. */ _hasMigration(db: import("../database/drivers/base.js").default, version?: string): Promise; /** * Runs apply migrations. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when complete. */ _applyMigrations(db: import("../database/drivers/base.js").default): Promise; /** * Runs ensure jobs table columns. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when complete. */ _ensureJobsTableColumns(db: import("../database/drivers/base.js").default): Promise; /** * Idempotently adds the pooled-child acceptance evidence columns to existing * job tables. They record when the executing runner child received and * started a job plus that child's identity, so a handed-off job can be told * apart from one whose runner never picked it up. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ensured. */ _ensureChildAcceptanceColumns(db: import("../database/drivers/base.js").default): Promise; /** * Repairs secondary indexes that older add-column upgrades declared but did * not create on every SQL driver. The migration ledger keeps routine store * readiness from repeatedly introspecting the full index set. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when all expected indexes exist. */ _ensureJobsTableIndexesOnce(db: import("../database/drivers/base.js").default): Promise; /** * Idempotently adds the per-job wall-clock timeout to existing job tables. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ensured. */ _ensureJobTimeoutColumn(db: import("../database/drivers/base.js").default): Promise; /** * Idempotently adds the historical stable schedule key to existing jobs. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ensured. */ _ensureScheduleKeyColumn(db: import("../database/drivers/base.js").default): Promise; /** * Idempotently adds monotonic schedule ownership history and its lookup index. * Existing rows remain null and use the documented legacy fallback ordering. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ensured. */ _ensureScheduleOrderColumn(db: import("../database/drivers/base.js").default): Promise; /** * Creates the retention-independent schedule-order high-water table and * initializes it from the greatest retained ordered row for every key. * Legacy rows whose order is null deliberately do not establish a watermark. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ensured and backfilled. */ _ensureScheduleOrderWatermarksTable(db: import("../database/drivers/base.js").default): Promise; /** * Backfills each key from its greatest retained non-legacy ownership order. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves after all retained keys are represented. */ _backfillScheduleOrderWatermarks(db: import("../database/drivers/base.js").default): Promise; /** * Idempotently adds the `queue` column to an existing jobs table. Existing * rows read back as the default queue (see {@link _normalizeJobRow}), so no * data backfill is required. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ensured. */ _ensureQueueColumn(db: import("../database/drivers/base.js").default): Promise; /** * Runs backfill execution modes once. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when complete. */ _backfillExecutionModesOnce(db: import("../database/drivers/base.js").default): Promise; /** * Rewrites pre-existing pooled rows (persisted as `execution_mode = "forked"` * plus a `velocious-pooled:*` handoff marker) to `execution_mode = "pooled"`, * clears the queued marker, then drops the now-redundant `forked` column so * `execution_mode` is the single source of truth. Runs once, guarded by the * migration ledger and a per-key advisory lock; a fresh table (created without * the column) short-circuits. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when complete. */ _dropForkedColumnOnce(db: import("../database/drivers/base.js").default): Promise; /** * Runs record migration. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} version - Migration version. * @returns {Promise} - Resolves when complete. */ _recordMigration(db: import("../database/drivers/base.js").default, version: string): Promise; _initializeModel(): Promise; /** * Runs get job row by id. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} jobId - Job id. * @returns {Promise} - Job row. */ _getJobRowById(db: import("../database/drivers/base.js").default, jobId: string): Promise; /** * Reads the job currently named by one stable owner row. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} scheduleKey - Validated stable schedule key. * @returns {Promise} - Normalized owner job. */ _scheduledOwnerJob(db: import("../database/drivers/base.js").default, scheduleKey: string): Promise; /** * Assigns the next ownership order while the caller holds the schedule-key * advisory lock and count-revision transaction fence. The independent * watermark survives both ownership release and terminal-history pruning. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} scheduleKey - Validated stable schedule key. * @returns {Promise} - Next monotonic ownership order. */ _nextScheduleOrder(db: import("../database/drivers/base.js").default, scheduleKey: string): Promise; /** * Finds the greatest retained non-legacy ownership order for migration and * rolling-upgrade compatibility. It is never the sole durability boundary. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} scheduleKey - Validated stable schedule key. * @returns {Promise} - Greatest retained order, or null. */ _greatestRetainedScheduleOrder(db: import("../database/drivers/base.js").default, scheduleKey: string): Promise; /** * Reads and validates one retention-independent schedule-order watermark. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} scheduleKey - Validated stable schedule key. * @returns {Promise} - Current watermark, or null before first ownership. */ _scheduleOrderWatermark(db: import("../database/drivers/base.js").default, scheduleKey: string): Promise; /** * Persists one schedule-order watermark without exposing it as a job row. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {object} args - Watermark identity and value. * @param {string} args.scheduleKey - Validated stable schedule key. * @param {number} args.scheduleOrder - Validated monotonic order. * @returns {Promise} - Resolves after persistence. */ _writeScheduleOrderWatermark(db: import("../database/drivers/base.js").default, { scheduleKey, scheduleOrder }: { scheduleKey: string; scheduleOrder: number; }): Promise; /** * Validates an ownership order loaded from durable storage. * @param {ReturnType} value - Stored order. * @returns {number} - Positive safe integer ownership order. */ _validatedScheduleOrder(value: ReturnType): number; /** * Builds a stable-schedule lookup exclusively from normalized job rows. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {object} args - Lookup options. * @param {boolean} args.includeLatestTerminal - Whether terminal history is requested. * @param {string} args.scheduleKey - Validated stable schedule key. * @returns {Promise} - Normalized public jobs. */ _scheduledJobLookup(db: import("../database/drivers/base.js").default, { includeLatestTerminal, scheduleKey }: { includeLatestTerminal: boolean; scheduleKey: string; }): Promise; /** * Releases ownership only when the key still points at the expected job. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {object} args - Ownership identity. * @param {string} args.jobId - Expected owner job id. * @param {string} args.scheduleKey - Stable schedule key. * @returns {Promise} - Resolves when deleted or already superseded. */ _releaseScheduleOwnership(db: import("../database/drivers/base.js").default, { jobId, scheduleKey }: { jobId: string; scheduleKey: string; }): Promise; /** * Releases a job's ownership when it has a historical schedule key. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {import("./types.js").BackgroundJobRow} job - Terminal job. * @returns {Promise} - Resolves when deleted or not applicable. */ _releaseScheduleOwnershipForJob(db: import("../database/drivers/base.js").default, job: import("./types.js").BackgroundJobRow): Promise; /** * Runs apply failure. * @param {object} args - Options. * @param {import("../database/drivers/base.js").default} args.db - Database connection. * @param {import("./types.js").BackgroundJobRow} args.job - Job row. * @param {ReturnType} args.error - Error. * @param {boolean} args.markOrphaned - Whether marking orphaned. * @param {Record>} [args.conditions] - Update fencing conditions. Defaults to the active-handoff lease match; the time-based orphan sweep overrides this with an id/status match so it can reclaim rows whose `handoff_id` is null (e.g. handed off by an older velocious before handoff-id fencing existed). * @returns {Promise} - Updated job row when the lease transition won. */ _applyFailure({ db, job, error, markOrphaned, conditions }: { db: import("../database/drivers/base.js").default; job: import("./types.js").BackgroundJobRow; error: ReturnType; markOrphaned: boolean; conditions?: Record>; }): Promise; /** * Runs failure update. * @param {object} args - Options. * @param {string} args.failureMessage - Last failure message. * @param {boolean} args.markOrphaned - Whether marking orphaned. * @param {number} args.nextAttempt - Next attempt count. * @param {number} args.now - Current timestamp. * @param {number | null} args.scheduledAt - Next scheduled timestamp. * @param {boolean} args.shouldRetry - Whether the job should retry. * @returns {Record>} - Database update data. */ _failureUpdate({ failureMessage, markOrphaned, nextAttempt, now, scheduledAt, shouldRetry }: { failureMessage: string; markOrphaned: boolean; nextAttempt: number; now: number; scheduledAt: number | null; shouldRetry: boolean; }): Record>; /** * Runs apply orphaned failure update. * @param {object} args - Options. * @param {boolean} args.markOrphaned - Whether marking orphaned. * @param {number} args.now - Current timestamp. * @param {Record>} args.update - Database update data. * @returns {void} */ _applyOrphanedFailureUpdate({ markOrphaned, now, update }: { markOrphaned: boolean; now: number; update: Record>; }): void; /** * Runs apply failure status update. * @param {object} args - Options. * @param {boolean} args.markOrphaned - Whether marking orphaned. * @param {number} args.now - Current timestamp. * @param {number | null} args.scheduledAt - Next scheduled timestamp. * @param {boolean} args.shouldRetry - Whether the job should retry. * @param {Record>} args.update - Database update data. * @returns {void} */ _applyFailureStatusUpdate({ markOrphaned, now, scheduledAt, shouldRetry, update }: { markOrphaned: boolean; now: number; scheduledAt: number | null; shouldRetry: boolean; update: Record>; }): void; /** * Runs normalize job row. * @param {Record>} row - Raw database row. * @returns {import("./types.js").BackgroundJobRow} - Normalized job row. */ _normalizeJobRow(row: Record>): import("./types.js").BackgroundJobRow; /** * Normalizes a job's queue name, defaulting to "default". * @param {import("./types.js").BackgroundJobOptions | undefined} options - Job options. * @returns {string} - Queue name. */ _normalizeQueue(options: import("./types.js").BackgroundJobOptions | undefined): string; /** * Resolves a job's durable concurrency. An explicit concurrencyKey/maxConcurrency * pair always wins. Otherwise, when the job's queue has a configured cap * (`backgroundJobs.queues[queue].maxConcurrent`), derive a queue-scoped * concurrency key so the queue cap is enforced cluster-wide through the * existing durable concurrency mechanism. * @param {import("./types.js").BackgroundJobOptions | undefined} options - Job options. * @param {string} queue - Normalized queue name. * @returns {{concurrencyKey: string, maxConcurrency: number, queueDerived: boolean} | null} - Resolved concurrency. */ _resolveConcurrency(options: import("./types.js").BackgroundJobOptions | undefined, queue: string): { concurrencyKey: string; maxConcurrency: number; queueDerived: boolean; } | null; /** * Applies the active generation's queue policy immediately before handoff. * Explicit concurrency remains owned by the enqueue request. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {import("./types.js").BackgroundJobRow} job - Queued job snapshot. * @returns {Promise} - Reconciled job, or null when its queued-state fence lost. */ _reconcileQueuedJobConcurrency(db: import("../database/drivers/base.js").default, job: import("./types.js").BackgroundJobRow): Promise; /** * Reads the configured max concurrency for a queue from the background-jobs config. * @param {string} queue - Queue name. * @returns {number | null} - Positive integer cap, or null when the queue has no configured cap. */ _queueMaxConcurrency(queue: string): number | null; /** * Like {@link _ensureConcurrencyKey}, but for queue-derived keys the configured * queue cap is the source of truth: if it changed, update the stored cap * instead of throwing on conflict (config-driven caps must be tunable). * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {{concurrencyKey: string, maxConcurrency: number}} concurrency - Concurrency configuration. * @returns {Promise} - Resolves when ensured. */ _ensureQueueConcurrencyKey(db: import("../database/drivers/base.js").default, { concurrencyKey, maxConcurrency }: { concurrencyKey: string; maxConcurrency: number; }): Promise; /** * Ensures the concurrency state table exists. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ready. */ _ensureConcurrencyTable(db: import("../database/drivers/base.js").default): Promise; /** * Ensures the stable schedule-key ownership table exists. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ready. */ _ensureScheduleKeysTable(db: import("../database/drivers/base.js").default): Promise; /** * Ensures durable generic enqueue ownership exists independently of job rows. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ready. */ _ensureIdempotencyKeysTable(db: import("../database/drivers/base.js").default): Promise; /** * Ensures durable provider-backed mail operation state exists independently of jobs. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when ready. */ _ensureMailDeliveryOperationsTable(db: import("../database/drivers/base.js").default): Promise; /** * Ensures the singleton durable count-revision row exists. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} Resolves when ready. */ _ensureCountRevisionTable(db: import("../database/drivers/base.js").default): Promise; /** * Records one logical count mutation atomically and broadcasts it after commit. * Zero entries are omitted; a wholly zero-net mutation does not consume a revision. * @param {import("../database/drivers/base.js").default} db - Transaction connection. * @param {Record} requestedDeltas - Signed bucket changes. * @returns {Promise} Resolves when recorded. */ _recordCountDelta(db: import("../database/drivers/base.js").default, requestedDeltas: Record): Promise; /** * Records a transition between persisted statuses. * @param {import("../database/drivers/base.js").default} db - Transaction connection. * @param {string} oldStatus - Previous status. * @param {string} newStatus - New status. * @returns {Promise} Resolves when recorded. */ _recordStatusTransition(db: import("../database/drivers/base.js").default, oldStatus: string, newStatus: string): Promise; /** * Reads the locked revision. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} Revision. */ _countRevision(db: import("../database/drivers/base.js").default): Promise; /** * Takes a portable write lock on the singleton revision row. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} Resolves when locked. */ _lockCountRevision(db: import("../database/drivers/base.js").default): Promise; /** * Builds zeroed canonical buckets. * @returns {Record} Zeroed canonical buckets. */ _emptyCountBuckets(): Record; /** * Counts normalized rows by canonical status. * @param {import("./types.js").BackgroundJobRow[]} jobs - Jobs. * @returns {Record} Counts. */ _statusCounts(jobs: import("./types.js").BackgroundJobRow[]): Record; /** * Reads a canonical snapshot after locking the revision row. * @param {import("../database/drivers/base.js").default} db - Transaction connection. * @returns {Promise<{counts: Record, revision: number, total: number}>} Snapshot. */ _countSnapshotOnLockedConnection(db: import("../database/drivers/base.js").default): Promise<{ counts: Record; revision: number; total: number; }>; /** * Registers or verifies a stable key configuration. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {object} concurrency - Concurrency configuration. * @param {string} concurrency.concurrencyKey - Concurrency key. * @param {number} concurrency.maxConcurrency - Stable cap. * @returns {Promise} - Resolves when verified. */ _ensureConcurrencyKey(db: import("../database/drivers/base.js").default, { concurrencyKey, maxConcurrency }: { concurrencyKey: string; maxConcurrency: number; }): Promise; /** * Locks the concurrency counter row so a job-release transaction acquires it *before* the job * row. {@link markHandedOff} reserves capacity (locking the counter row) before it updates the * job, so it locks concurrency-then-job; the release paths update the job before releasing * capacity, which is job-then-concurrency. Those opposite orders on the same shared counter row * are what deadlock (AB-BA) under a draining worker. Taking this lock first gives every * transaction a single concurrency-then-job order and removes the cycle. * * Uses a value-preserving `UPDATE` rather than `SELECT ... FOR UPDATE` so it stays portable * across drivers without row-level locking reads (e.g. SQLite); on row-locking engines the * matched row is write-locked for the rest of the transaction even though its value is unchanged. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string | null} concurrencyKey - Concurrency key. * @returns {Promise} - Resolves when the counter row is locked. */ _lockConcurrencyRow(db: import("../database/drivers/base.js").default, concurrencyKey: string | null): Promise; /** * Atomically reserves capacity for a key. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} concurrencyKey - Concurrency key. * @returns {Promise} - Whether capacity was reserved. */ _reserveConcurrency(db: import("../database/drivers/base.js").default, concurrencyKey: string): Promise; /** * Runs a portable update and returns its affected-row count. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {import("../database/drivers/base.js").UpdateSqlArgsType} args - Update options. * @returns {Promise} - Affected row count. */ _updateAffectedRows(db: import("../database/drivers/base.js").default, args: import("../database/drivers/base.js").UpdateSqlArgsType): Promise; /** * Releases capacity for a key. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string | null} concurrencyKey - Concurrency key. * @returns {Promise} - Resolves when released. */ _releaseConcurrency(db: import("../database/drivers/base.js").default, concurrencyKey: string | null): Promise; /** * Rebuilds durable counts from active handoffs. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {{insideTransaction?: boolean}} [options] - Reuse an enclosing transaction. * @returns {Promise} - Repair summary. */ _reconcileConcurrency(db: import("../database/drivers/base.js").default, { insideTransaction }?: { insideTransaction?: boolean; }): Promise; /** * Rebuilds one counter after locking it ahead of the job rows, matching the * lock order used by handoff and completion transitions. * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {string} concurrencyKey - Counter key. * @returns {Promise} - Applied repair. */ _reconcileConcurrencyKey(db: import("../database/drivers/base.js").default, concurrencyKey: string): Promise; /** * Validates a database count before it participates in reconciliation. * @param {number | string} value - Raw count. * @param {string} concurrencyKey - Counter key for diagnostics. * @returns {number} - Safe non-negative count. */ _validatedConcurrencyCount(value: number | string, concurrencyKey: string): number; /** * Reconciles queue-derived concurrency with the current configuration. Only * invoked through {@link reconcileQueueConcurrency} — the explicit lifecycle * path run at main-process startup under a cross-process advisory lock — * never from schema/tenant checks or routine connection initialization, * which stay read-only regarding queued job rows. The per-process memo is * latched by {@link reconcileQueueConcurrency} only after the following * count rebuild also succeeds, so a failed rebuild re-enters here on retry * (the adoption UPDATEs below are idempotent). Enqueue only consults config for new jobs, so a cap added, removed, or changed * while a backlog exists otherwise leaves persisted rows stale: pre-cap jobs * keep a null key and bypass the cap, post-removal jobs stay capped under a * now-unconfigured key, and a changed numeric cap stays stale until the next * enqueue. Bring queued durable state in line with config: sync each configured * queue's stored cap, adopt not-yet-keyed queued jobs onto their queue key, * and release queued jobs from queue keys whose queue is no longer capped. * Existing handoffs retain the policy and reservation they started with, so * reconciliation cannot race their completion/retry transitions. Runs before * {@link _reconcileConcurrency} so any pre-existing active counts are exact. * @param {import("../database/drivers/base.js").default} db - Database connection. * @returns {Promise} - Resolves when reconciled. */ _reconcileQueueConcurrency(db: import("../database/drivers/base.js").default): Promise; /** * Runs normalize number. * @param {ReturnType} value - Input value. * @returns {number | null} - Normalized number. */ _normalizeNumber(value: ReturnType): number | null; /** * Runs normalize execution mode. * @param {import("./types.js").BackgroundJobOptions} [options] - Job options. * @returns {import("./types.js").BackgroundJobExecutionMode} - Normalized execution mode. */ _normalizeExecutionMode(options?: import("./types.js").BackgroundJobOptions): import("./types.js").BackgroundJobExecutionMode; /** * Runs normalize execution mode name. * @param {string} executionMode - Execution mode name. * @returns {import("./types.js").BackgroundJobExecutionMode} - Normalized execution mode. */ _normalizeExecutionModeName(executionMode: string): import("./types.js").BackgroundJobExecutionMode; /** * Filters queued jobs by one or more execution modes against the * `execution_mode` column (the single source of truth). * @param {object} args - Options. * @param {import("../database/drivers/base.js").default} args.db - Database connection. * @param {import("./types.js").BackgroundJobExecutionMode | import("./types.js").BackgroundJobExecutionMode[]} args.executionMode - Runtime modes. * @param {import("../database/query/index.js").default} args.query - Query to filter. * @returns {import("../database/query/index.js").default} - Filtered query. */ _whereExecutionMode({ db, executionMode, query }: { db: import("../database/drivers/base.js").default; executionMode: import("./types.js").BackgroundJobExecutionMode | import("./types.js").BackgroundJobExecutionMode[]; query: import("../database/query/index.js").default; }): import("../database/query/index.js").default; /** * Runs parse args. * @param {ReturnType} value - Input value. * @returns {Array>} - Parsed args. */ _parseArgs(value: ReturnType): Array>; /** * Runs with db. * @template T * @param {(db: import("../database/drivers/base.js").default) => Promise} callback - Callback. * @returns {Promise} - Callback result. */ _withDb(callback: (db: import("../database/drivers/base.js").default) => Promise): Promise; /** * Runs a value-returning callback inside the driver's void-typed transaction API. * @template T * @param {import("../database/drivers/base.js").default} db - Database connection. * @param {() => Promise} callback - Transaction callback. * @returns {Promise} - Callback result. */ _transactionResult(db: import("../database/drivers/base.js").default, callback: () => Promise): Promise; /** * Serializes count-changing transactions before checking out their connection. * Database row locking still provides cross-process ordering; this guard * prevents concurrent callers on SQLite's shared connection from attempting * overlapping top-level transactions. * @template T * @param {(db: import("../database/drivers/base.js").default) => Promise} callback - Transaction callback. * @param {BackgroundJobTransactionSerializationOptions} [options] - Serialization options. * @returns {Promise} Callback result. */ _serializedCountMutation(callback: (db: import("../database/drivers/base.js").default) => Promise, options?: BackgroundJobTransactionSerializationOptions): Promise; /** * Runs a serialized callback inside one transaction. * @template T * @param {(db: import("../database/drivers/base.js").default) => Promise} callback - Transaction callback. * @param {BackgroundJobTransactionSerializationOptions} [options] - Serialization options. * @returns {Promise} Callback result. */ _serializedTransactionMutation(callback: (db: import("../database/drivers/base.js").default) => Promise, options?: BackgroundJobTransactionSerializationOptions): Promise; /** * Admits mutation callbacks to the process-local FIFO before they check out a * connection. Cross-process ordering remains the responsibility of durable * row/advisory locks and unique constraints acquired around the callback. * @template T * @param {(db: import("../database/drivers/base.js").default) => Promise} callback - Connection callback. * @param {BackgroundJobTransactionSerializationOptions} [options] - Serialization options. * @returns {Promise} Callback result. */ _serializedConnectionMutation(callback: (db: import("../database/drivers/base.js").default) => Promise, options?: BackgroundJobTransactionSerializationOptions): Promise; /** * Runs should accept report. * @param {object} args - Options. * @param {import("./types.js").BackgroundJobRow} args.job - Job row. * @param {string | null | undefined} args.handoffId - Handoff lease id from report. * @param {string | null | undefined} args.workerId - Worker id from report. * @param {number | null | undefined} args.handedOffAtMs - Handed off timestamp from report. * @returns {boolean} - Whether to accept the report. */ _shouldAcceptReport({ job, handoffId, workerId, handedOffAtMs }: { job: import("./types.js").BackgroundJobRow; handoffId: string | null | undefined; workerId: string | null | undefined; handedOffAtMs: number | null | undefined; }): boolean; /** * Runs active handoff conditions. * @param {import("./types.js").BackgroundJobRow} job - Job row. * @returns {Record} - Conditional transition fence. */ _activeHandoffConditions(job: import("./types.js").BackgroundJobRow): Record; /** * Runs handoff id report matches. * @param {object} args - Options. * @param {string | null | undefined} args.handoffId - Handoff lease id from report. * @param {import("./types.js").BackgroundJobRow} args.job - Job row. * @returns {boolean} - Whether the handoff lease matches. */ _handoffIdReportMatches({ handoffId, job }: { handoffId: string | null | undefined; job: import("./types.js").BackgroundJobRow; }): boolean; /** * Runs worker report matches. * @param {object} args - Options. * @param {import("./types.js").BackgroundJobRow} args.job - Job row. * @param {string | null | undefined} args.workerId - Worker id from report. * @returns {boolean} - Whether the worker report matches. */ _workerReportMatches({ job, workerId }: { job: import("./types.js").BackgroundJobRow; workerId: string | null | undefined; }): boolean; /** * Runs handoff report matches. * @param {object} args - Options. * @param {number | null | undefined} args.handedOffAtMs - Handed off timestamp from report. * @param {import("./types.js").BackgroundJobRow} args.job - Job row. * @returns {boolean} - Whether the handoff report matches. */ _handoffReportMatches({ handedOffAtMs, job }: { handedOffAtMs: number | null | undefined; job: import("./types.js").BackgroundJobRow; }): boolean; /** * Runs migration key. * @param {string} [version] - Migration version. * @returns {string} - Migration key. */ _migrationKey(version?: string): string; } //# sourceMappingURL=store.d.ts.map