Skip to content

feat: make the notification queue durable with claim leases - #879

Open
Evaristus023 wants to merge 1 commit into
Core-Foundry:mainfrom
Evaristus023:feat/persistent-notification-queue-780
Open

Evaristus023 wants to merge 1 commit into
Core-Foundry:mainfrom
Evaristus023:feat/persistent-notification-queue-780

Conversation

@Evaristus023

Copy link
Copy Markdown

Overview

The listener already stores scheduled notifications in scheduled_notifications and dequeues them with an atomic UPDATE ... WHERE id IN (SELECT ...) claim, but the claim had no heartbeat: lock_expires_at was written once at dequeue time and never extended, while both schedulers call recoverStaleLocks() at the top of every poll cycle. Any delivery that outlived SCHEDULER_LOCK_TIMEOUT_MS (60s by default) was therefore reset to PENDING by the next poll and handed to another worker while the first worker was still sending it — the same notification delivered twice, and a job that was still running counted as a failed attempt.

This change turns the claim into a durable, renewable lease:

  • the worker that owns a claimed job renews its lease at a third of the lock window (owner-checked, so a worker that genuinely lost the claim cannot take it back);
  • the dequeue path now honours the persisted retry backoff, so a job inside its next_retry_at window is no longer claimed immediately by the main scheduler while the retry scheduler owns it;
  • failed jobs keep their recorded processing state instead of having it overwritten by a concurrent claimant.

Related Issue

Closes #780

Changes

Persistent claim + heartbeat

  • [ADD] listener/src/services/notification-claim-lease.ts

    • NotificationClaimLease renews a claim every lockTimeoutMs / 3 via calculateLeaseRenewIntervalMs and is strictly dependency-free (timers and the renewal function are injectable), so the renewal policy is unit-testable without a database.
    • Renewals are serialised: a slow heartbeat can never overlap the next one.
    • A renewal that resolves false means the claim is gone: the lease is marked lost, the owner is notified once, and it never renews again.
    • A transient renewal error keeps the claim (the existing lease is still valid) and retries on the next tick; the serialisation chain is kept alive so a throwing callback cannot silently kill the heartbeat.
    • stop() clears the timer and drains the in-flight renewal, so the database can be closed immediately afterwards.
    • startClaimLease() returns null for repositories that have no renewLock, so existing test doubles keep working unchanged.
  • [MODIFY] listener/src/services/scheduled-notification-repository.ts

    • fetchAndLockPendingNotifications (dequeue) now skips rows whose next_retry_at is still in the future, so backoff state is respected and only one dequeue path owns a retrying job.
    • New renewLock(id, processorId, lockTimeoutMs) heartbeat: extends lock_expires_at only while the row is still PROCESSING and owned by that processor, and reports whether the lease is still held.
  • [MODIFY] listener/src/services/notification-scheduler.ts

    • Starts a lease heartbeat for each claimed notification and stops it in the finally block next to workerManager.completeJob, logging lost leases and failed renewals.
  • [MODIFY] listener/src/services/retry-scheduler.ts

    • The same heartbeat around processRetry.

Tests

  • [ADD] listener/src/services/notification-claim-lease.test.ts

    • Renewal interval bounds, renewal argument identity, overlapping-heartbeat serialisation, lease-lost, transient-error recovery, stop() draining, idempotent start/stop, and both startClaimLease branches.
  • [ADD] listener/src/services/scheduled-notification-queue.test.ts

    • Acceptance coverage: enqueue → close the database → reopen → the job is still claimable and its payload intact; a job whose worker died mid-flight is requeued and claimable again; three due jobs land with exactly one of two workers; a renewed lease is not recovered while a dead worker's is; a job inside its backoff window keeps retry_count/last_error and is not returned by either dequeue path; a permanently failed job keeps retry_count, last_error, error_details, processing_started_at, processing_completed_at and is neither claimable nor renewable.

Verification Results

Implemented via GitHub REST / Git Data API (no local clone, no dependency install).

1. tsc --noEmit (repo tsconfig)      -> 6 errors, all pre-existing syntax errors in files this PR does not touch
                                        (src/index.ts, src/middleware/security-headers.ts, src/services/discord-notification.ts)
                                        ... 0 errors in the changed files.
2. tsc --noEmit (focused program)    -> 0 errors in notification-claim-lease.ts and scheduled-notification-repository.ts
                                        (residual TS7016/TS7006/TS7031 errors come from local .d.ts shims standing in
                                        for the uninstalled sqlite3 / winston / @types/node-cache).
3. Executable harness (real TS, real SQLite via Node's built-in node:sqlite with a local sqlite3 shim):
   lease + repository acceptance checks  -> 15/15 passed
   scheduler integration checks          -> 3/3 passed
     * 400 ms delivery against a 120 ms lease: recoverStaleLocks() -> 0 with the heartbeat
     * control (same run, heartbeat disabled): recoverStaleLocks() -> 1, job stolen mid-flight
     * retry scheduler: claim held for the whole delivery

The jest suites added here were NOT executed: listener/node_modules is not installed and
main does not parse (src/utils/request-id.ts is missing the `}` + `/**` that close
generateCorrelationId, alongside the three syntax errors above), so no listener test can run
until those pre-existing breakages are fixed. The behaviour they assert was verified by the
executable harness above.
Acceptance Criteria Status
Queued notifications survive restarts ✅ rows stay in scheduled_notifications; verified by dequeue after close → reopen, and by requeueing the job of a worker that died mid-flight
A job cannot be processed by multiple workers simultaneously ✅ atomic claim + owner-checked lease heartbeat; verified live (0 stolen) and by the no-heartbeat control (1 stolen)
Failed jobs retain their processing state ✅ retry_count, last_error, error_details, processing_started_at, processing_completed_at preserved; terminal jobs are not claimable or renewable and their backoff window is honoured

Closes #780

@drips-wave

drips-wave Bot commented Sep 28, 2026

Copy link
Copy Markdown

@Evaristus023 Great news! 🎉 Based on an automated assessment of this PR, the linked Wave issue(s) no longer count against your application limits.

You can now already apply to more issues while waiting for a review of this PR. Keep up the great work! 🚀

Learn more about application limits

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Implement Persistent Notification Queue

1 participant