duraq 3.0.0
duraq: ^3.0.0 copied to clipboard
A durable queuing system implemented in Dart, with a SQLite backend
Changelog #
All notable changes to DuraQ will be documented in this file.
The format is based on Keep a Changelog, and this project adheres to Semantic Versioning.
3.0.0 - 2026-09-14 #
One semantic change: an entry id now identifies an entry within its queue, rather than across the whole database. Nothing else in the API moves.
This is the change the 2.0.0 notes said should be revisited. It was deferred because rebuilding the SQLite table needed a migration path that did not exist yet; 2.0.0 built one, and this release uses it.
Existing databases are migrated in place on first open, and cannot be opened by 2.x afterwards. Read Upgrading from 2.0.x below before deploying.
Upgrading from 2.0.x #
For most callers this is a no-op recompile: every method that takes an entry id already took a queue name beside it, so no call site changes shape.
1. Check whether you relied on ids being unique across the database. You did
if you used DuplicateEntryException to find out whether an id was in use
anywhere, or if you treat an id as addressing an entry without saying which
queue it is in. Both now need the queue name to be meaningful.
// Before: the second store threw, whatever queue it named.
await storage.store('orders', entry('order-42'));
await storage.store('shipping', entry('order-42')); // DuplicateEntryException
// Now: two queues, two entries, neither affecting the other.
await storage.store('orders', entry('order-42'));
await storage.store('shipping', entry('order-42')); // fine
A repeated id within one queue is still a conflict, and StoreConflict.replace
and .ignore still do exactly what they did.
2. Back up the database file if you may need to roll back. The first open
migrates it to schema version 2, and 2.x will then refuse it with a
SchemaVersionException naming both versions. That refusal is deliberate: an
older release would read the file as if ids were still global and could delete a
second queue's entry under StoreConflict.replace. Refusing to open is the
better failure.
The migration itself is safe to interrupt — it runs in one transaction that rolls back whole — and cannot lose a row, because the version 1 key was strictly stricter than the version 2 key.
3. If you use the Isar backend, move to duraq_isar 2.0.0 at the same time.
The two release together.
4. If you implement StorageInterface yourself, nothing forces a change —
but your backend should now treat (queueName, entryId) as the identity of an
entry. The shared conformance suite in test/support/ covers this, and running
your backend against it is the quickest way to find out where you stand.
Breaking #
-
An entry id is now unique within its queue, not across the whole database. Storing
order-42in one queue no longer stops another queue from holding an entry with that id: they are different entries, and neither affects the other. Retrieval, replacement, removal and status changes all act on the entry in the queue named, as they already did — the public API does not change shape, because every method that takes an entry id already took a queue name beside it.Queues are namespaces. Two independent producers writing
order-42into queues of their own is reasonable, and until now the second one failed with aDuplicateEntryExceptionfor a reason that had nothing to do with either of them. Storing the same id twice in one queue is still a conflict, andStoreConflict.replaceand.ignorestill do what they did.You are affected if you relied on global uniqueness — using
DuplicateEntryExceptionto find out whether an id was in use anywhere, or treating an id as addressing an entry without saying which queue it is in. -
The SQLite schema is at version 2. A version 1 database is migrated in place the first time this release opens it: the
queue_entriestable is rebuilt withPRIMARY KEY (queue_name, id)and its indexes are recreated, in one transaction that rolls back whole if anything in it fails. No row can be lost to the change, because the old key was strictly stricter than the new one.Once migrated, the database can no longer be opened by duraq 2.x, which would read it as if ids were still global. That is deliberate: a
SchemaVersionExceptionnaming both versions is better than an older release quietly deleting a second queue's entry underStoreConflict.replace. Back the file up before upgrading if you may need to roll back.
Changed #
DuplicateEntryException's message now names the queue, matching what it has always meant:Queue "orders" already holds an entry with id "order-42".
2.0.0 - 2026-09-14 #
A durability release. Six defects that could lose or duplicate a job are fixed, every finding from the durability audit of 12 September 2026 is closed, and the package is now tested by CI for the first time.
If you are upgrading from 1.0.x, read the Breaking section below: dequeue
now means something different, and the Isar backend has moved to its own
package.
Fixed #
- The health check API is now exported from
package:duraq/duraq.dart.HealthCheck,StorageHealthCheck,MetricsHealthCheck,QueueHealthCheck,HealthCheckAggregator,HealthStatusandHealthCheckResultwere documented in the README but never exported, so following the README gave "Method not found". The existing tests passed only because they imported thesrc/path directly. QueueHealthChecknow reports on the queues that exist. It asked the metrics collector for the size of a queue calleddefault, a name nothing in the system used, so the figure was zero unless a caller happened to record one there. It now reads every queue from the storage — or the ones named in the newqueueNames— and reports what is waiting and what is ready per queue.MetricsHealthCheckno longer probes a queue calledhealth-check, which nothing used either. It reads the figures the collector holds for the system as a whole and reports them.ExponentialBackoff.getRetryDelay()no longer overflows. The delay wasbaseDelay.inMilliseconds * pow(2, attempts)in integer arithmetic, which wraps a 64-bit int at attempt 63.maxDelaycould not catch the wrapped value, so attempt 58 scheduled a retry roughly 60,000 years out, and from attempt 64 every delay came out zero — backoff became a tight retry loop. The delay is now computed in floating point and capped before jitter, so it is never negative, never zero, and never abovemaxDelayat any attempt count.- Jitter now applies after the cap rather than before it. Previously every
attempt at or past the ceiling returned exactly
maxDelaywith no spread at all, which is the point in a backoff where spreading retries matters most. Delays at the ceiling now land between 75% and 100% ofmaxDelay. updateEntryStatus()no longer clears fields the caller did not mention. Every update wroteerror_messageandnext_retry_at, using null when they were not supplied, so completing an entry erased the error that explained its last failure and any status change dropped a pending retry time. Fields left out now keep their stored values, with one deliberate exception: moving an entry topendingwithout a retry time clears the backoff, because that is what making an entry available again means.- Dead lettered and expired entries now release their claim. The lock was
released for
completed,failedandpendingonly, so an entry kept its lease for the rest of its duration after the work was over. The visible symptom wasretryDeadLetter()appearing to do nothing: the entry went back to pending still locked by the claim it died under, and the scan skipped it until the lease ran out. - Lock ids are now unique per acquisition. They were built from the queue name, the entry id and the clock in milliseconds, so two claims of the same entry inside one millisecond produced the same id, and a consumer holding the older one was accepted as the current holder — defeating the ownership check that exists to stop exactly that. Found by a test that flaked once in three runs.
- Taking a lock no longer reports every failure as contention. Any exception during the insert answered "the entry is already locked", so a broken schema, a failing disk, or a database that went away moved the scan on to the next candidate and made a failing storage look like an empty queue. Only a real conflict — a primary key collision on SQLite, a unique index violation on Isar — now means locked; everything else is raised.
dequeue()now removes the entry it returns. It previously claimed the entry and handed back the payload with no way to acknowledge it, leaving the row inprocessing; once the lease expired the entry was handed out again, so every dequeued item came back. The claim and the delete now happen in one transaction. That makesdequeue()at-most-once: useprocessNext(), which is unchanged, when an item must not be lost.QueueManager.queue<T>()no longer throws a cast error when a queue is asked for under a second element type. The cache was keyed by name alone, so the first caller's type won for the life of the process and a queue first touched untyped could never be fetched typed. Queues are now cached per name and type, so a worker readingInvoiceand an admin tool readingdynamiccan share one queue.- Removing an entry now releases its lock. A deleted id stayed marked as claimed until its lease ran out, so the same id could not be enqueued and picked up again in that window.
- Concurrent calls to
SQLiteStorage.transaction()no longer nest inside one another. Previously two overlapping transactions shared a single depth counter, so one caller's rollback discarded another caller's committed rows while that caller still reported success. - Concurrent calls to
SQLiteStorage.retrieve()no longer return null while entries are still pending. Ten parallel consumers against five entries now receive all five, each exactly once. - Retry delays are now honoured. Both backends exclude an entry from retrieval
until its
nextRetryAthas passed, soExponentialBackoffand any otherRetryPolicytake effect instead of a failed entry being re-delivered immediately. Measured cost of the added predicate is 0.44 microseconds. SQLiteStorage.store()now persistsnextRetryAt. The column was missing from the insert, so a pre-built entry lost its retry time silently.- The Isar backend no longer stores one row per
storecall.entryIdcarried a non-unique index next to an auto-incrementing key, so storing the same entry twice created a second row, counted it twice, and handed the same payload to a processor twice.storenow upserts through an entry identity index. IsarStorage.transaction()now provides atomicity. The body runs in a single Isar write transaction and operations called inside it join that transaction, so a body that throws leaves nothing behind. Previously each write committed on its own and a failure halfway left the earlier writes in place.- Two processes sharing one SQLite file no longer fail each other's writes. The busy timeout was left at zero, so a second writer got an immediate "database is locked" instead of waiting its turn. Two processes doing 300 enqueues and 50 claims each now complete with no failures, where 47 of 300 enqueues failed before.
- Switching a new database file to write-ahead logging no longer crashes when two processes open it at the same moment. That switch needs an exclusive lock and does not go through the busy handler, so it is retried briefly and then accepts the mode the file is in.
dispose()no longer throws when another connection holds the write lock. It released its locks through a future whose failure nothing handled, so a contended shutdown surfaced as an unhandled exception.- Entries are no longer stranded when a consumer dies. An entry left in
processingwith no live lock is returned to the queue on the next retrieval, so a crash, a kill, or adequeue()that is never acknowledged no longer loses the job. Both backends. - A consumer whose lease expired can no longer finish an entry that has since
been given to someone else. Claims carry a lease id now:
retrievereturns it on the entry,updateEntryStatustakes it, and a status change made against a lease that is no longer the live one is discarded rather than applied over the work of whoever holds the entry.Queue.processNextpasses it for you. The lock id had been generated and returned since the beginning and then kept by nobody, so release deleted whatever lock was on the entry. dispose()no longer releases locks held by other consumers. The release was an unfiltered delete over the lock table, so one process shutting down freed every in-flight entry in the database, including entries other processes were still working on. Each lock manager now tracks and releases only its own.
Upgrading from 1.0.x #
Four things need attention. Everything else is source-compatible.
1. If you use the Isar backend, add the package and one import. Existing databases open unchanged.
dependencies:
duraq: ^2.0.0
duraq_isar: ^1.0.0 # new
import 'package:duraq/duraq.dart';
import 'package:duraq_isar/duraq_isar.dart'; // new
Pass ...IsarStorage.requiredSchemas to Isar.open rather than listing the
collections yourself — the set gained a schema-version collection, and a
hand-written list will fail to open with an error saying so.
2. If you call dequeue(), decide whether you meant it. It now removes the
entry as it hands it over, which is what it always claimed to do. That makes it
at most once: work in flight when the process dies is gone. Previously the
entry was left claimed and redelivered when its lease expired, so a dequeue()
loop was silently redelivering everything it processed.
final item = await queue.dequeue(); // at most once; lost if you crash
await queue.processNext(handleItem); // at least once; retried and reclaimed
Use processNext for anything that must not be lost. If you were relying on the
old redelivery, you were relying on a bug, but the behaviour you want is
processNext.
3. If you implement StorageInterface yourself, add close() and
countReady(). Dart's implements copies signatures only, so the interface's
default bodies do not reach you — the compiler will say so. See
test/custom_backend_test.dart in this package for the smallest version that
satisfies it.
4. If you call updateEntryStatus() directly, it now throws
EntryNotFoundException when no entry has that id, rather than reporting
success. A change discarded because your lease expired still returns quietly.
Also worth knowing: the SDK floor is now 3.2.0. That corrects a false claim
rather than dropping support — sqlite3 has required 3.2.0 for some time, so
1.0.x could never actually resolve on 3.0 or 3.1.
Breaking #
- The declared SDK floor moves from
>=3.0.0to>=3.2.0. This corrects a claim rather than dropping support:sqlite32.2.0, the oldest this package allows, itself requires 3.2.0, soduraqcould never have resolved on 3.0 or 3.1. Found by the CI job that builds on the advertised floor. - The Isar backend moved to its own package,
duraq_isar. A project using only SQLite no longer resolvesisar, its generated code, or its version constraint — which was the point:isaris pinned to 3.1.0+1, whose generated code raises analyzer warnings on current Dart. Projects using Isar addduraq_isarto their pubspec and one import;IsarStoragebehaves as before and existing databases open unchanged. See that package's changelog. StorageInterfacegainedclose(). Custom backends must declare it; both built-in backends delegate to their existingdispose().StorageInterfacegainedcountReady(). Custom backends must declare it; the interface's default body filtersretrieveAll, which is correct but reads every entry, so a real backend should answer with a query.IsarStorage.requiredSchemasnow includesQueueMetaCollectionSchema. Callers already passing...IsarStorage.requiredSchemastoIsar.openneed no change and their databases upgrade in place. A caller that listed DuraQ's collections by hand must add it; doing so raises aDuraQExceptionnaming the fix rather than Isar's ownMissing TypeSchema, which names neither the collection nor what to do.updateEntryStatus()now throwsEntryNotFoundExceptionwhen no entry with that id is in the queue. It previously reported success, so a typo, a stale id, or an entry already removed by a retention pass all looked like work completing normally. A change discarded because the caller's lease expired still returns quietly: the entry exists and someone else holds the claim.StorageInterface.store()takes a newonConflictparameter,updateEntryStatus()takes a newleaseIdparameter, andStorageInterfacegainedrunMaintenance(). Custom backends must add all three. Dart'simplementscopies only signatures, so the defaultrunMaintenancebody does not reach a class that implements the interface rather than extending it.- Storing an entry whose id is already in the storage now throws
DuplicateEntryExceptionon every backend. SQLite previously threw a rawSqliteExceptioncarrying the failing statement and its parameters, which put the entry payload into the error text; Isar silently added a second row. PassStoreConflict.replaceorStoreConflict.ignorewhere a repeated store is expected. IsarStorage.beginTransaction(),commitTransaction()androllbackTransaction()now throwUnsupportedError. Isar write transactions take their work as a callback, so a transaction cannot be opened in one call and closed in another. Previously these moved a counter and guaranteed nothing: work done between them was already committed, and a rollback discarded nothing. Usetransaction(), which now runs the whole body in one Isar write transaction.- The Isar schema changed: indexes that no queries used were removed, and indexes matching the queries this package runs were added. Existing databases open and migrate on their own; no action is needed beyond the duplicate cleanup below.
Performance #
- Dequeue no longer slows down as the backlog grows. The per-queue expiry sweep
on the retrieval path could not use the partial expiry index and fell back to
scanning every pending entry in the queue. A
(queue_name, expires_at)partial index turns it into a seek: 6.29 microseconds against 1,086 at a backlog of 20,000. End to end, a dequeue and acknowledge at that depth went from 1,187 microseconds to 294, and the cost is now flat across depths. - Retrieval reads candidates in batches instead of one row at a time with a growing offset, which re-read the head of the queue on every locked candidate. Both backends.
- The Isar backend now uses indexes. Every query ran as a full collection scan followed by an in-memory sort, while the model declared ten indexes no query referenced. Queries now go through an entry identity index, a composite index that returns entries already in retrieval order, and a per-queue expiry index; the unused ones are gone. At a backlog of 20,000 a dequeue and acknowledge went from 3,599 microseconds to 542, enqueue from 2,563 per second to 5,960, and neither now degrades as the backlog grows.
Added #
runMaintenance({policy, queueName})onStorageInterfaceand both backends: one periodic pass that returns entries whose consumer died, marks entries that outlived their deadline, and deletes finished entries older than the retention policy allows. It reports what it did as aMaintenanceReport. Nothing calls it on a schedule; run it from your own timer or at startup.RetentionPolicy, controlling how long completed, failed, dead lettered and expired entries are kept. Defaults to 7 days for completed and failed, 30 for dead letters, 1 for expired.RetentionPolicy.keepEverything()deletes nothing.StoreConflictand aDuraQExceptionhierarchy withDuplicateEntryException, so a duplicate id can be handled without catching driver-specific errors. The entry payload is deliberately kept out of the error text.onConflictonQueue.enqueueEntry, for producers that retry an enqueue they got no answer for.busyTimeoutonSQLiteStorage: how long to keep trying to start a write when another process or isolate holds the write lock. Defaults to 5 seconds. The wait is spent in short slices with the isolate free in between, rather than one long block inside the driver.StorageBusyException, thrown when that budget runs out. It replaces the rawSqliteExceptiona caller would otherwise have to recognise, and the work is untouched, so retrying the call is safe.leaseDurationon both storage backends: how long a retrieved entry stays claimed before another consumer may take it. Defaults to five minutes, which is the duration that was previously hardcoded.maxDeliveryAttemptson both storage backends: how many times an entry may be delivered before an expiring lease sends it to the dead letter queue instead of back to the queue. Defaults to 5, so a job that crashes its consumer cannot cycle forever.reclaimStaleEntries({String? queueName})on both storage backends, for recovering entries at startup that a previous run left claimed. Returns the number of entries returned to pending.IsarStorage.removeDuplicateEntries(), a one-off cleanup for databases written by earlier versions. Those versions could store several rows for one entry; this collapses them, keeping the most recently updated row, and returns how many rows it removed. Run it once after upgrading.close()onStorageInterface, implemented by both backends. Neitherdisposenorcloseappeared on the interface before, and the two backends disagreed on the shape — SQLite's was synchronous and returned void, Isar's was asynchronous — so shutdown could not be written without knowing which backend was underneath.dispose()stays on both concrete classes for callers already using it; the Isar backend still leaves the Isar instance open, since the caller owns it.synchronousonSQLiteStorage, withSqliteSynchronousand theactiveSynchronousgetter that reads the setting back from the live connection. Write-ahead logging runs atNORMAL, which is unchanged and still the default: a DuraQ process that dies loses nothing, but the machine going down can lose the most recent commits — for a queue, jobs that were accepted.SqliteSynchronous.fullcloses that at a measured cost of about half the single-enqueue throughput (13,496/s to 6,913/s, median of five interleaved rounds).countReady()onStorageInterface, both backends andQueue.readyLength: the number of entries that can be handed out now.count()includes entries scheduled for later and entries waiting out a retry backoff, so a queue could report a length of two and hand out nothing — an autoscaler reading it scales up for work that is not due.count()is unchanged and still reports the backlog; its documentation now says which is which. The interface carries a working default that filtersretrieveAll, but Dart'simplementscopies signatures only, so a custom backend must declare it.- A
metricsparameter onQueueandQueueManager. Nothing in the library recorded a metric, so every rate aQueueMetricscould report read zero however busy the system was. A queue given a collector now records enqueues, dequeues, completions, failures, latencies and processing times, each labelled with the queue's name. Throughput counts attempts rather than successes, sogetErrorRatestays a proportion. maxReadyBacklogonQueueHealthCheck, which reports degraded when too much work is ready to run. Deliberately compared against ready rather than waiting: a queue full of entries scheduled for next week is not falling behind. Running the check also samples each queue's size into the metrics collector, sogetCurrentQueueSizestops reading zero forever.- Schema versioning on both backends. SQLite records its version in the
user_versionpragma; Isar records it in a newQueueMetaCollectionrow, which is what Isar has no equivalent of. A database written by a newer DuraQ is now refused on open withSchemaVersionExceptionrather than being read as if nothing had changed, and both backends have a migration runner so the next change of shape or meaning has a way to reach databases already in the field. Before this, tables were created if absent and never versioned. SQLiteStorage.schemaVersionandIsarStorage.schemaVersion, the version this release writes, alongsideSQLiteStorage.storedSchemaVersionandIsarStorage.storedSchemaVersion()for what a given database is at.SchemaVersionException, raised when a storage is at a version this release does not understand, and when an entry carries a status string this release has no name for. The latter previously surfaced asArgumentError: Invalid argument (name), naming neither the entry nor why.QueueCodec<T>and acodecparameter onQueue,DeadLetterQueueandQueueManager.queue. Payloads are stored as JSON, which limited a queue to whatjsonEncodeaccepts no matter what its type argument said. A codec makes that boundary explicit and lifts it, soQueue<Invoice>can hold anInvoice.QueueCodec.from(encode:, decode:)builds one from a pair of functions.PayloadCodecException, raised when a payload cannot cross that boundary: an unencodable payload with no codec, a codec that threw, or a queue read through an element type its entries were not written with. Each replaces an error that named only the failing conversion —JsonUnsupportedObjectError, or a bareTypeErrorabout two unrelated types — with one naming the queue, the type, and the way out.QueueEntry.withData<R>(), which copies an entry around a payload of a different type.copyWithcannot change the payload type, and both encoding and decoding do.tool/verify.sh, the project's gate:dart analyze --fatal-infos --fatal-warningsfollowed by the test suite.tool/verify.sh --flake Nruns the suite N times instead, to surface timing flakes.tool/install-hooks.shpoints git at.githooks/, whosepre-pushhook runs the gate before anything leaves the machine.analysis_options.yaml. Thelintsdev dependency has been declared since 1.0.0 but was never applied, because nothing told the analyzer to use it. Generated Isar code is excluded;test/analysis_options.yamlrelaxes two rules that only make sense for shipped code.- A GitHub Actions workflow running the same script. The repository had no workflows before, so nothing had ever been checked automatically. A second, non-blocking job reports whether the advertised Dart 3.0 floor still builds and re-runs the suite to watch for flakes.
Changed #
- Operations on a
SQLiteStorageinstance are serialized, and calls made from inside atransaction()body remain re-entrant. Measured cost is under 0.3 microseconds per operation. beginTransaction()now holds exclusive access to the storage until the transaction is committed or rolled back. A manual transaction that is never closed will block later operations; prefertransaction().- Raw generic types are written out.
QueueEntryin a signature now readsQueueEntry<dynamic>, which is the same type spelled honestly: the storage layer is untyped by design. No behaviour changes. sqlite3now requires 2.2.0 or later, up from 2.1.0, which is the version that replaced the deprecatedgetUpdatedRows()withupdatedRows. The package already resolved well above this floor in practice.testnow requires 1.25.0 or later, for the reportertool/verify.shuses. Dev dependency only; consumers are unaffected.
1.0.0 - 2026-03-22 #
Changed #
- BREAKING CHANGE:
IsarStoragenow requires an external Isar instance instead of creating its own IsarStorage.create()factory method has been removedIsarStorageconstructor now accepts anIsarinstance parameterIsarStorage.dispose()no longer closes the Isar instance (caller responsibility)- Added
IsarStorage.requiredSchemasstatic getter to help users configure Isar with required schemas
Added #
- Support for shared Isar instances across multiple components
- Better integration with external systems that manage Isar lifecycle
- Comprehensive documentation for new Isar usage patterns
Migration Guide #
Before:
final storage = await IsarStorage.create(dbPath: 'path/to/queue.isar');
// ... use storage
await storage.dispose(); // Closed Isar automatically
After:
await Isar.initializeIsarCore(download: true);
final isar = await Isar.open([
...IsarStorage.requiredSchemas,
// Add your other schemas here
], directory: 'path/to/db');
final storage = IsarStorage(isar);
// ... use storage
await storage.dispose(); // Releases locks only
await isar.close(); // Caller manages Isar lifecycle