From 61782740304a5e4b9e6e4ece19974dbfb3635d22 Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 16:04:46 +0530 Subject: [PATCH 1/8] feat(core): add count_pending_by_queue storage method --- .../src/storage/diesel_common/jobs.rs | 13 ++++++++++++ crates/taskito-core/src/storage/mod.rs | 9 ++++++++ .../src/storage/redis_backend/jobs/query.rs | 7 +++++++ .../taskito-core/src/storage/sqlite/tests.rs | 21 +++++++++++++++++++ crates/taskito-core/src/storage/traits.rs | 4 ++++ .../taskito-core/tests/rust/storage_tests.rs | 3 +++ 6 files changed, 57 insertions(+) diff --git a/crates/taskito-core/src/storage/diesel_common/jobs.rs b/crates/taskito-core/src/storage/diesel_common/jobs.rs index 91613d23..f3025a99 100644 --- a/crates/taskito-core/src/storage/diesel_common/jobs.rs +++ b/crates/taskito-core/src/storage/diesel_common/jobs.rs @@ -1976,6 +1976,19 @@ macro_rules! impl_diesel_job_ops { Ok(count) } + /// Count pending jobs on a queue (for the `max_pending` admission cap). + pub fn count_pending_by_queue(&self, queue_name: &str) -> Result { + let mut conn = self.conn()?; + + let count: i64 = jobs::table + .filter(jobs::queue.eq(queue_name)) + .filter(jobs::status.eq(JobStatus::Pending as i32)) + .count() + .get_result(&mut conn)?; + + Ok(count) + } + /// Purge job errors older than the given timestamp. /// /// Deletes in bounded batches, each its own txn — see diff --git a/crates/taskito-core/src/storage/mod.rs b/crates/taskito-core/src/storage/mod.rs index 92251195..bd27c8c7 100644 --- a/crates/taskito-core/src/storage/mod.rs +++ b/crates/taskito-core/src/storage/mod.rs @@ -835,6 +835,12 @@ macro_rules! impl_storage { ) -> $crate::error::Result { self.count_running_by_task(task_name) } + fn count_pending_by_queue( + &self, + queue_name: &str, + ) -> $crate::error::Result { + self.count_pending_by_queue(queue_name) + } fn stats_by_queue( &self, queue_name: &str, @@ -1444,6 +1450,9 @@ impl Storage for StorageBackend { fn count_running_by_task(&self, task_name: &str) -> Result { delegate!(self, count_running_by_task, task_name) } + fn count_pending_by_queue(&self, queue_name: &str) -> Result { + delegate!(self, count_pending_by_queue, queue_name) + } fn stats_by_queue(&self, queue_name: &str) -> Result { delegate!(self, stats_by_queue, queue_name) } diff --git a/crates/taskito-core/src/storage/redis_backend/jobs/query.rs b/crates/taskito-core/src/storage/redis_backend/jobs/query.rs index ffbceef1..57b3366d 100644 --- a/crates/taskito-core/src/storage/redis_backend/jobs/query.rs +++ b/crates/taskito-core/src/storage/redis_backend/jobs/query.rs @@ -254,6 +254,13 @@ impl RedisStorage { self.count_in_status(&mut conn, &by_task_key, JobStatus::Running) } + /// Count pending jobs on a queue (for the `max_pending` admission cap). + pub fn count_pending_by_queue(&self, queue_name: &str) -> Result { + let mut conn = self.conn()?; + let by_queue_key = self.key(&["jobs", "by_queue", queue_name]); + self.count_in_status(&mut conn, &by_queue_key, JobStatus::Pending) + } + pub fn stats_by_queue(&self, queue_name: &str) -> Result { let mut conn = self.conn()?; self.queue_stats(&mut conn, queue_name) diff --git a/crates/taskito-core/src/storage/sqlite/tests.rs b/crates/taskito-core/src/storage/sqlite/tests.rs index e0ceddb0..f6c2cf22 100644 --- a/crates/taskito-core/src/storage/sqlite/tests.rs +++ b/crates/taskito-core/src/storage/sqlite/tests.rs @@ -806,6 +806,27 @@ fn test_count_running_by_task() { assert_eq!(storage.count_running_by_task("no_such_task").unwrap(), 0); } +#[test] +fn test_count_pending_by_queue() { + let storage = test_storage(); + assert_eq!(storage.count_pending_by_queue("default").unwrap(), 0); + + storage.enqueue(make_job("task_a")).unwrap(); + storage.enqueue(make_job("task_a")).unwrap(); + let mut other = make_job("task_b"); + other.queue = "other".to_string(); + storage.enqueue(other).unwrap(); + + assert_eq!(storage.count_pending_by_queue("default").unwrap(), 2); + assert_eq!(storage.count_pending_by_queue("other").unwrap(), 1); + assert_eq!(storage.count_pending_by_queue("empty").unwrap(), 0); + + // Dequeue drops the job out of Pending → count decreases. + let now = now_millis() + 1000; + storage.dequeue("default", now, None).unwrap().unwrap(); + assert_eq!(storage.count_pending_by_queue("default").unwrap(), 1); +} + #[test] fn test_enqueue_rejects_missing_dependency() { let storage = test_storage(); diff --git a/crates/taskito-core/src/storage/traits.rs b/crates/taskito-core/src/storage/traits.rs index 91388a95..4236515f 100644 --- a/crates/taskito-core/src/storage/traits.rs +++ b/crates/taskito-core/src/storage/traits.rs @@ -329,6 +329,10 @@ pub trait Storage: Send + Sync + Clone { // ── Per-queue stats ────────────────────────────────────────── + /// Cheap count of pending jobs on a queue — the admission-cap primitive. + /// Single-status, unlike the full-breakdown `stats_by_queue`. + fn count_pending_by_queue(&self, queue_name: &str) -> Result; + fn stats_by_queue(&self, queue_name: &str) -> Result; fn stats_all_queues(&self) -> Result>; diff --git a/crates/taskito-core/tests/rust/storage_tests.rs b/crates/taskito-core/tests/rust/storage_tests.rs index de4730b5..b3144f8a 100644 --- a/crates/taskito-core/tests/rust/storage_tests.rs +++ b/crates/taskito-core/tests/rust/storage_tests.rs @@ -173,6 +173,8 @@ fn test_stats_by_queue_and_task(s: &impl Storage) { assert_eq!(st.pending, 3); assert_eq!(st.running, 0); assert_eq!(s.count_running_by_task(task).unwrap(), 0); + // Lean pending-count primitive agrees with the full breakdown. + assert_eq!(s.count_pending_by_queue(q).unwrap(), 3); // Run two of them. let d1 = s.dequeue(q, now_millis() + 1000, None).unwrap().unwrap(); @@ -181,6 +183,7 @@ fn test_stats_by_queue_and_task(s: &impl Storage) { let st = s.stats_by_queue(q).unwrap(); assert_eq!(st.running, 2); assert_eq!(st.pending, 1); + assert_eq!(s.count_pending_by_queue(q).unwrap(), 1); // Complete one — running drops, completed rises. s.complete(&d1.id, None).unwrap(); From 3f35332b6af02c96d918e31a60338afebc0bcdcc Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 16:05:25 +0530 Subject: [PATCH 2/8] feat(python): add max_pending admission cap --- crates/taskito-python/src/py_queue/mod.rs | 8 ++ sdks/python/taskito/__init__.py | 2 + sdks/python/taskito/_taskito.pyi | 1 + sdks/python/taskito/app.py | 37 ++++++- sdks/python/taskito/exceptions.py | 4 + sdks/python/taskito/mixins/runtime_config.py | 16 +++ sdks/python/tests/core/test_admission.py | 102 +++++++++++++++++++ 7 files changed, 169 insertions(+), 1 deletion(-) create mode 100644 sdks/python/tests/core/test_admission.py diff --git a/crates/taskito-python/src/py_queue/mod.rs b/crates/taskito-python/src/py_queue/mod.rs index 65332330..20b34825 100644 --- a/crates/taskito-python/src/py_queue/mod.rs +++ b/crates/taskito-python/src/py_queue/mod.rs @@ -589,6 +589,14 @@ impl PyQueue { }) } + /// Count pending jobs on a queue — the lean primitive behind the + /// `max_pending` admission cap (avoids the full `stats_by_queue` breakdown). + pub fn count_pending_by_queue(&self, queue_name: &str) -> PyResult { + self.storage + .count_pending_by_queue(queue_name) + .map_err(|e| pyo3::exceptions::PyRuntimeError::new_err(e.to_string())) + } + /// Get queue statistics broken down by queue name. pub fn stats_all_queues(&self) -> PyResult> { let all = self diff --git a/sdks/python/taskito/__init__.py b/sdks/python/taskito/__init__.py index 609b060d..d7ff5245 100644 --- a/sdks/python/taskito/__init__.py +++ b/sdks/python/taskito/__init__.py @@ -27,6 +27,7 @@ ProxyCleanupError, ProxyReconstructionError, QueueError, + QueueFullError, RateLimitExceededError, ResourceError, ResourceInitError, @@ -96,6 +97,7 @@ "ProxyReconstructionError", "Queue", "QueueError", + "QueueFullError", "RateLimitExceededError", "ResourceError", "ResourceInitError", diff --git a/sdks/python/taskito/_taskito.pyi b/sdks/python/taskito/_taskito.pyi index d1098a97..6ae724a2 100644 --- a/sdks/python/taskito/_taskito.pyi +++ b/sdks/python/taskito/_taskito.pyi @@ -163,6 +163,7 @@ class PyQueue: def purge_completed(self, older_than_seconds: int) -> int: ... def stats(self) -> dict[str, int]: ... def stats_by_queue(self, queue_name: str) -> dict[str, int]: ... + def count_pending_by_queue(self, queue_name: str) -> int: ... def stats_all_queues(self) -> dict[str, dict[str, int]]: ... def list_jobs_filtered( self, diff --git a/sdks/python/taskito/app.py b/sdks/python/taskito/app.py index 19fa5c51..30c43aee 100644 --- a/sdks/python/taskito/app.py +++ b/sdks/python/taskito/app.py @@ -30,7 +30,7 @@ from taskito.batching import BatchAccumulator, BatchConfig from taskito.codecs import CodecSerializer, PayloadCodec from taskito.events import EventBus, EventType -from taskito.exceptions import SerializationError +from taskito.exceptions import QueueFullError, SerializationError from taskito.interception import ArgumentInterceptor from taskito.interception.built_in import build_default_registry from taskito.interception.metrics import InterceptionMetrics @@ -152,6 +152,7 @@ def __init__( dlq_auto_retry_delay: int | None = None, dlq_auto_retry_max: int = 1, retention: Retention | None = None, + max_pending: dict[str, int] | None = None, ): """Initialize a new task queue. @@ -225,6 +226,14 @@ def __init__( is automatically retried. ``None`` disables auto-retry. dlq_auto_retry_max: Maximum number of DLQ auto-retries per entry before giving up. Defaults to 1. + max_pending: Opt-in per-queue admission cap, mapping queue name to a + maximum pending backlog. When set, ``enqueue``/``enqueue_many`` + raise :class:`QueueFullError` once a queue's pending count + reaches its cap. Queues absent from the map are uncapped (zero + overhead). The check is a non-atomic count-then-insert, so brief + overshoot is possible under concurrent producers — the same soft + guarantee as the rate limiter. Also settable at runtime via + ``set_queue_max_pending``. """ if backend == "sqlite": # Ensure parent directory exists for SQLite @@ -302,6 +311,8 @@ def __init__( self._init_predicate_state() self._drain_timeout = drain_timeout self._queue_configs: dict[str, dict[str, Any]] = {} + # Opt-in per-queue admission caps (queue -> max pending). Empty = uncapped. + self._max_pending: dict[str, int] = dict(max_pending or {}) self._event_bus = EventBus(max_workers=event_workers) self._webhook_manager = WebhookManager(queue_ref=self) @@ -509,6 +520,23 @@ def _serialize_result(self, task_name: str, result: Any) -> bytes: """ return self._serializer.dumps(result) + def _reject_if_queue_full(self, queue_name: str) -> None: + """Enforce the opt-in ``max_pending`` admission cap for a queue. + + Raises :class:`QueueFullError` when the queue's pending backlog has + reached its configured cap. No-op (and no query) for uncapped queues. + Non-atomic count-then-insert — brief overshoot is accepted, like the + rate limiter. + """ + cap = self._max_pending.get(queue_name) + if cap is None: + return + pending = self._inner.count_pending_by_queue(queue_name) + if pending >= cap: + raise QueueFullError( + f"queue '{queue_name}' is full: {pending} pending >= max_pending {cap}" + ) + def enqueue( self, task_name: str, @@ -664,6 +692,8 @@ def enqueue( if depends_on is not None: dep_ids = [depends_on] if isinstance(depends_on, str) else list(depends_on) + self._reject_if_queue_full(queue or "default") + py_job = self._inner.enqueue( task_name=task_name, payload=payload, @@ -894,6 +924,11 @@ def enqueue_many( delay=delays[i], ) + # Admission cap is all-or-nothing for the batch: if any distinct target + # queue is at its cap, reject before inserting any row. + for capped_queue in set(queues_list): + self._reject_if_queue_full(capped_queue) + py_jobs = self._inner.enqueue_batch( task_names=task_names, payloads=payloads, diff --git a/sdks/python/taskito/exceptions.py b/sdks/python/taskito/exceptions.py index c489f43d..0e228f58 100644 --- a/sdks/python/taskito/exceptions.py +++ b/sdks/python/taskito/exceptions.py @@ -77,6 +77,10 @@ class QueueError(TaskitoError): """Raised on queue-level operational errors.""" +class QueueFullError(QueueError): + """Raised when enqueuing would exceed a queue's ``max_pending`` admission cap.""" + + class ResourceError(TaskitoError): """Base exception for resource system errors.""" diff --git a/sdks/python/taskito/mixins/runtime_config.py b/sdks/python/taskito/mixins/runtime_config.py index 38d4dad0..217f2f38 100644 --- a/sdks/python/taskito/mixins/runtime_config.py +++ b/sdks/python/taskito/mixins/runtime_config.py @@ -25,6 +25,7 @@ class QueueRuntimeConfigMixin: """Rate limits, concurrency caps, and custom type registration.""" _queue_configs: dict[str, dict[str, Any]] + _max_pending: dict[str, int] _interceptor: ArgumentInterceptor | None def register_type( @@ -85,3 +86,18 @@ def set_queue_concurrency(self, queue_name: str, max_concurrent: int) -> None: from this queue. """ self._queue_configs.setdefault(queue_name, {})["max_concurrent"] = max_concurrent + + def set_queue_max_pending(self, queue_name: str, max_pending: int) -> None: + """Set an opt-in admission cap on a queue's pending backlog. + + Once the queue holds ``max_pending`` pending jobs, ``enqueue`` and + ``enqueue_many`` raise :class:`~taskito.exceptions.QueueFullError`. The + check is a non-atomic count-then-insert (brief overshoot is possible + under concurrent producers). Enforced producer-side, so it applies even + when no worker is running. + + Args: + queue_name: Queue name (e.g. ``"default"``). + max_pending: Maximum pending jobs allowed before enqueue is rejected. + """ + self._max_pending[queue_name] = max_pending diff --git a/sdks/python/tests/core/test_admission.py b/sdks/python/tests/core/test_admission.py new file mode 100644 index 00000000..c5bf49c2 --- /dev/null +++ b/sdks/python/tests/core/test_admission.py @@ -0,0 +1,102 @@ +"""S26 — opt-in ``max_pending`` admission cap. + +Jobs stay Pending without a running worker, so these tests exercise the cap +purely on the producer side. +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from taskito import Queue, QueueFullError + + +def _register(queue: Queue) -> None: + @queue.task(name="noop") + def noop() -> None: # pragma: no cover - never executed (no worker) + return None + + +def test_count_pending_by_queue_primitive(queue: Queue) -> None: + _register(queue) + assert queue._inner.count_pending_by_queue("default") == 0 + queue.enqueue("noop") + queue.enqueue("noop") + assert queue._inner.count_pending_by_queue("default") == 2 + + +def test_uncapped_queue_never_rejects(queue: Queue) -> None: + _register(queue) + for _ in range(50): + queue.enqueue("noop") + assert queue._inner.count_pending_by_queue("default") == 50 + + +def test_runtime_setter_rejects_at_cap(queue: Queue) -> None: + _register(queue) + queue.set_queue_max_pending("default", 2) + queue.enqueue("noop") + queue.enqueue("noop") + with pytest.raises(QueueFullError) as exc: + queue.enqueue("noop") + assert "max_pending 2" in str(exc.value) + # Rejected enqueue inserted nothing. + assert queue._inner.count_pending_by_queue("default") == 2 + + +def test_constructor_cap(tmp_path: Path) -> None: + q = Queue(db_path=str(tmp_path / "t.db"), workers=1, max_pending={"default": 1}) + _register(q) + q.enqueue("noop") + with pytest.raises(QueueFullError): + q.enqueue("noop") + + +def test_cap_is_per_queue(queue: Queue) -> None: + _register(queue) + queue.set_queue_max_pending("tight", 1) + queue.enqueue("noop", queue="tight") + with pytest.raises(QueueFullError): + queue.enqueue("noop", queue="tight") + # A different, uncapped queue is unaffected. + for _ in range(5): + queue.enqueue("noop", queue="wide") + + +def test_queue_full_is_queue_error_subclass() -> None: + from taskito.exceptions import QueueError, TaskitoError + + assert issubclass(QueueFullError, QueueError) + assert issubclass(QueueFullError, TaskitoError) + + +def test_enqueue_many_all_or_nothing(queue: Queue) -> None: + _register(queue) + queue.set_queue_max_pending("default", 3) + queue.enqueue("noop") + queue.enqueue("noop") + queue.enqueue("noop") # now at cap + with pytest.raises(QueueFullError): + queue.enqueue_many("noop", [(), (), ()]) + # None of the batch landed. + assert queue._inner.count_pending_by_queue("default") == 3 + + +def test_cap_frees_after_drain(queue: Queue, run_worker: object) -> None: + """Once pending drains below the cap, enqueue is admitted again.""" + from taskito import JobResult + + results: list[JobResult] = [] + + @queue.task(name="quick") + def quick() -> str: + return "ok" + + queue.set_queue_max_pending("default", 100) + # With a worker draining, pending stays well under the cap. + for _ in range(20): + results.append(queue.enqueue("quick")) + for r in results: + assert r.result(timeout=10) == "ok" From 442503b06953e3ebb58b9cbe9a44cca3a76257be Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 16:05:51 +0530 Subject: [PATCH 3/8] feat(node): add maxPending admission cap --- crates/taskito-node/src/queue/mod.rs | 10 ++++ sdks/node/src/errors.ts | 12 +++++ sdks/node/src/index.ts | 1 + sdks/node/src/queue.ts | 29 ++++++++++++ sdks/node/src/types.ts | 6 +++ sdks/node/test/core/admission.test.ts | 67 +++++++++++++++++++++++++++ 6 files changed, 125 insertions(+) create mode 100644 sdks/node/test/core/admission.test.ts diff --git a/crates/taskito-node/src/queue/mod.rs b/crates/taskito-node/src/queue/mod.rs index 0833eafe..8aa0066e 100644 --- a/crates/taskito-node/src/queue/mod.rs +++ b/crates/taskito-node/src/queue/mod.rs @@ -85,6 +85,16 @@ impl JsQueue { Ok(created.into_iter().map(|job| job.id).collect()) } + /// Count pending jobs on a queue — the lean primitive behind the + /// `maxPending` admission cap. Sync so the producer can gate a sync + /// `enqueue`/`enqueueMany` without a round trip to the event loop. + #[napi] + pub fn count_pending_by_queue(&self, queue: String) -> Result { + self.storage + .count_pending_by_queue(&queue) + .map_err(to_napi_err) + } + /// Fetch a job by id, or `null` if no such job exists. #[napi] pub fn get_job(&self, id: String) -> Result> { diff --git a/sdks/node/src/errors.ts b/sdks/node/src/errors.ts index db229355..35b347dc 100644 --- a/sdks/node/src/errors.ts +++ b/sdks/node/src/errors.ts @@ -79,6 +79,18 @@ export class QueueError extends TaskitoError { } } +/** Thrown by `enqueue`/`enqueueMany` when a queue's `maxPending` admission cap is reached. */ +export class QueueFullError extends QueueError { + constructor( + readonly queue: string, + readonly pending: number, + readonly cap: number, + ) { + super(`queue '${queue}' is full: ${pending} pending >= maxPending ${cap}`); + this.name = "QueueFullError"; + } +} + /** Thrown by {@link Queue.withLock} when the lock is held by another owner. */ export class LockNotAcquiredError extends TaskitoError { constructor(readonly lockName: string) { diff --git a/sdks/node/src/index.ts b/sdks/node/src/index.ts index 6dedd291..19bf22f4 100644 --- a/sdks/node/src/index.ts +++ b/sdks/node/src/index.ts @@ -25,6 +25,7 @@ export { PredicateRejectedError, ProxyError, QueueError, + QueueFullError, ResourceError, ResourceNotFoundError, ResourceScopeError, diff --git a/sdks/node/src/queue.ts b/sdks/node/src/queue.ts index a59abc3d..18319946 100644 --- a/sdks/node/src/queue.ts +++ b/sdks/node/src/queue.ts @@ -15,6 +15,7 @@ import { LockNotAcquiredError, PredicateRejectedError, QueueError, + QueueFullError, ResourceError, ResultTimeoutError, SerializationError, @@ -440,10 +441,26 @@ export class Queue { args?: Parameters, options?: EnqueueOptions, ): string { + this.rejectIfQueueFull(options?.queue ?? "default"); const { taskName, payload, options: nativeOpts } = this.prepareEnqueue(name, args, options); return this.native.enqueue(taskName, payload, nativeOpts); } + /** + * Enforce the opt-in `maxPending` admission cap for a queue (set via + * {@link Queue.configureQueue}). Throws {@link QueueFullError} once the + * queue's pending backlog reaches its cap; a no-op (and no query) for + * uncapped queues. Non-atomic count-then-insert, like the rate limiter. + */ + private rejectIfQueueFull(queue: string): void { + const cap = this.queueLimits.get(queue)?.maxPending; + if (cap === undefined) return; + const pending = this.native.countPendingByQueue(queue); + if (pending >= cap) { + throw new QueueFullError(queue, pending, cap); + } + } + /** * Enqueue many jobs of `name` in one storage round-trip. Each entry is its own * typed `args` + `options`. Returns the new job ids in input order. Unlike @@ -453,6 +470,10 @@ export class Queue { name: Name, jobs: ReadonlyArray<{ args?: Parameters; options?: EnqueueOptions }>, ): string[] { + // All-or-nothing: if any distinct target queue is at its cap, reject the batch. + for (const queue of new Set(jobs.map((job) => job.options?.queue ?? "default"))) { + this.rejectIfQueueFull(queue); + } const prepared = jobs.map((job) => { const { payload, options } = this.prepareEnqueue(name, job.args, job.options, { batch: true, @@ -831,6 +852,14 @@ export class Queue { return this.native.statsByQueue(queue); } + /** + * Count pending jobs on `queue` — the lean primitive behind the `maxPending` + * admission cap (avoids the full {@link Queue.statsByQueue} breakdown). + */ + countPendingByQueue(queue: string): number { + return this.native.countPendingByQueue(queue); + } + /** Job counts by status, keyed by queue name. */ statsAllQueues(): Promise> { return this.native.statsAllQueues(); diff --git a/sdks/node/src/types.ts b/sdks/node/src/types.ts index 01b05876..eb19697c 100644 --- a/sdks/node/src/types.ts +++ b/sdks/node/src/types.ts @@ -167,6 +167,12 @@ export interface PeriodicTask { export interface QueueLimits { maxConcurrent?: number; rateLimit?: RateLimit; + /** + * Opt-in admission cap on the queue's pending backlog. Once reached, + * `enqueue`/`enqueueMany` throw {@link QueueFullError}. Enforced producer-side + * (a non-atomic count-then-insert), so it applies even with no worker running. + */ + maxPending?: number; } /** A task handler plus its registration options. */ diff --git a/sdks/node/test/core/admission.test.ts b/sdks/node/test/core/admission.test.ts new file mode 100644 index 00000000..7a472c9e --- /dev/null +++ b/sdks/node/test/core/admission.test.ts @@ -0,0 +1,67 @@ +// S26 — opt-in `maxPending` admission cap. Jobs stay pending without a worker, +// so the cap is exercised purely producer-side. + +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { expect, it } from "vitest"; +import { Queue, QueueError, QueueFullError } from "../../src/index"; + +function newQueue(): Queue { + const dbPath = join(mkdtempSync(join(tmpdir(), "taskito-adm-")), "queue.db"); + return new Queue({ dbPath }); +} + +it("countPendingByQueue counts pending per queue", () => { + const queue = newQueue(); + queue.task("noop", () => undefined); + expect(queue.countPendingByQueue("default")).toBe(0); + queue.enqueue("noop"); + queue.enqueue("noop"); + expect(queue.countPendingByQueue("default")).toBe(2); +}); + +it("uncapped queue never rejects", () => { + const queue = newQueue(); + queue.task("noop", () => undefined); + for (let i = 0; i < 25; i++) queue.enqueue("noop"); + expect(queue.countPendingByQueue("default")).toBe(25); +}); + +it("rejects at the configured cap", () => { + const queue = newQueue(); + queue.task("noop", () => undefined); + queue.configureQueue("default", { maxPending: 2 }); + queue.enqueue("noop"); + queue.enqueue("noop"); + expect(() => queue.enqueue("noop")).toThrow(QueueFullError); + // Rejected enqueue inserted nothing. + expect(queue.countPendingByQueue("default")).toBe(2); +}); + +it("cap is per queue", () => { + const queue = newQueue(); + queue.task("noop", () => undefined); + queue.configureQueue("tight", { maxPending: 1 }); + queue.enqueue("noop", undefined, { queue: "tight" }); + expect(() => queue.enqueue("noop", undefined, { queue: "tight" })).toThrow(QueueFullError); + for (let i = 0; i < 5; i++) queue.enqueue("noop", undefined, { queue: "wide" }); +}); + +it("enqueueMany is all-or-nothing against the cap", () => { + const queue = newQueue(); + queue.task("noop", () => undefined); + queue.configureQueue("default", { maxPending: 3 }); + queue.enqueue("noop"); + queue.enqueue("noop"); + queue.enqueue("noop"); + expect(() => queue.enqueueMany("noop", [{}, {}, {}])).toThrow(QueueFullError); + expect(queue.countPendingByQueue("default")).toBe(3); +}); + +it("QueueFullError is a QueueError", () => { + const err = new QueueFullError("q", 5, 5); + expect(err).toBeInstanceOf(QueueError); + expect(err.queue).toBe("q"); + expect(err.cap).toBe(5); +}); From 295793ebc3e941e31b4528067dd45e555e071a47 Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 16:06:09 +0530 Subject: [PATCH 4/8] feat(java): add maxPending admission cap --- crates/taskito-java/src/queue/inspect.rs | 18 +++++ .../org/byteveda/taskito/DefaultTaskito.java | 34 ++++++++ .../java/org/byteveda/taskito/Taskito.java | 12 +++ .../taskito/errors/QueueFullException.java | 36 +++++++++ .../taskito/internal/JniQueueBackend.java | 5 ++ .../taskito/internal/NativeQueue.java | 2 + .../byteveda/taskito/spi/QueueBackend.java | 3 + .../byteveda/taskito/core/AdmissionTest.java | 79 +++++++++++++++++++ .../taskito/test/InMemoryQueueBackend.java | 5 ++ 9 files changed, 194 insertions(+) create mode 100644 sdks/java/src/main/java/org/byteveda/taskito/errors/QueueFullException.java create mode 100644 sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java diff --git a/crates/taskito-java/src/queue/inspect.rs b/crates/taskito-java/src/queue/inspect.rs index df73d3a5..0381baf9 100644 --- a/crates/taskito-java/src/queue/inspect.rs +++ b/crates/taskito-java/src/queue/inspect.rs @@ -64,6 +64,24 @@ pub extern "system" fn Java_org_byteveda_taskito_internal_NativeQueue_statsByQue }) } +/// `long countPendingByQueue(long handle, String queue)` — the lean primitive +/// behind the `maxPending` admission cap. +#[no_mangle] +pub extern "system" fn Java_org_byteveda_taskito_internal_NativeQueue_countPendingByQueue< + 'local, +>( + mut env: JNIEnv<'local>, + _class: JClass<'local>, + handle: jlong, + queue_name: JString<'local>, +) -> jlong { + guard(&mut env, 0, |env| { + let queue = unsafe { borrow_queue(handle) }; + let name = read_string(env, &queue_name)?; + Ok(queue.storage.count_pending_by_queue(&name)?) + }) +} + /// `String statsAllQueues(long handle)` — a JSON map of queue name to counts. #[no_mangle] pub extern "system" fn Java_org_byteveda_taskito_internal_NativeQueue_statsAllQueues<'local>( diff --git a/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java b/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java index 32bb1441..24d9738f 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java @@ -21,6 +21,7 @@ import org.byteveda.taskito.errors.EnqueueSkippedException; import org.byteveda.taskito.errors.InterceptionException; import org.byteveda.taskito.errors.PredicateRejectedException; +import org.byteveda.taskito.errors.QueueFullException; import org.byteveda.taskito.errors.SerializationException; import org.byteveda.taskito.errors.WorkflowException; import org.byteveda.taskito.interception.Interception; @@ -91,6 +92,8 @@ final class DefaultTaskito implements Taskito { private final Map> gates = new ConcurrentHashMap<>(); private final List interceptors = new CopyOnWriteArrayList<>(); private final List subscriptions = new CopyOnWriteArrayList<>(); + // Opt-in per-queue admission caps (queue -> max pending). Absent = uncapped. + private final Map maxPending = new ConcurrentHashMap<>(); DefaultTaskito(QueueBackend backend, Serializer serializer, Map codecs) { this.backend = backend; @@ -145,6 +148,30 @@ public Map resourceMetrics() { return resources.metrics(); } + @Override + public Taskito maxPending(String queue, int cap) { + maxPending.put(queue, cap); + return this; + } + + /** + * Enforce the opt-in {@code maxPending} admission cap for a queue. Throws + * {@link QueueFullException} once the queue's pending backlog reaches its + * cap; a no-op (and no query) for uncapped queues. Non-atomic + * count-then-insert, like the rate limiter. + */ + private void rejectIfQueueFull(String queueOrNull) { + String queue = queueOrNull == null ? "default" : queueOrNull; + Integer cap = maxPending.get(queue); + if (cap == null) { + return; + } + long pending = backend.countPendingByQueue(queue); + if (pending >= cap) { + throw new QueueFullException(queue, pending, cap); + } + } + @Override public Taskito predicate(String taskName, Predicate predicate) { return gate( @@ -267,6 +294,8 @@ private Optional dispatchEnqueue( .metadata(encode(context.metadata())) .build(); } + // Admission cap: reject before serializing/inserting if the target queue is full. + rejectIfQueueFull(finalOptions.queue()); // Serialize before codec-encoding so the idempotency key hashes the deterministic // pre-codec payload — a non-deterministic codec (e.g. an AES-GCM nonce) must not // change the dedup key. @@ -460,6 +489,11 @@ public QueueStats statsByQueue(String queue) { return decode(backend.statsByQueueJson(queue), QueueStats.class); } + @Override + public long countPendingByQueue(String queue) { + return backend.countPendingByQueue(queue); + } + @Override public Map statsAllQueues() { return decodeMap(backend.statsAllQueuesJson(), QueueStats.class); diff --git a/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java b/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java index 803ff455..e84c53cc 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java @@ -121,6 +121,15 @@ static Builder builder() { */ Taskito intercept(Interceptor interceptor); + /** + * Set an opt-in admission cap on {@code queue}'s pending backlog. Once the + * queue holds {@code cap} pending jobs, {@link #enqueue} throws + * {@link org.byteveda.taskito.errors.QueueFullException}. Enforced + * producer-side (a non-atomic count-then-insert), so it applies even with no + * worker running. Returns {@code this}. + */ + Taskito maxPending(String queue, int cap); + // ── Producer ──────────────────────────────────────────────────── /** Enqueue a typed payload using the task's default options; returns the job id. */ @@ -176,6 +185,9 @@ static Builder builder() { QueueStats statsByQueue(String queue); + /** Count pending jobs on {@code queue} — the primitive behind the {@code maxPending} cap. */ + long countPendingByQueue(String queue); + Map statsAllQueues(); List listJobs(JobFilter filter); diff --git a/sdks/java/src/main/java/org/byteveda/taskito/errors/QueueFullException.java b/sdks/java/src/main/java/org/byteveda/taskito/errors/QueueFullException.java new file mode 100644 index 00000000..509a3639 --- /dev/null +++ b/sdks/java/src/main/java/org/byteveda/taskito/errors/QueueFullException.java @@ -0,0 +1,36 @@ +package org.byteveda.taskito.errors; + +import org.byteveda.taskito.TaskitoException; + +/** + * An enqueue was rejected because the target queue reached its {@code maxPending} + * admission cap, so no job was created. Enforced producer-side (a non-atomic + * count-then-insert), so it applies even with no worker running. + */ +public class QueueFullException extends TaskitoException { + private final String queue; + private final long pending; + private final long cap; + + public QueueFullException(String queue, long pending, long cap) { + super("queue '" + queue + "' is full: " + pending + " pending >= maxPending " + cap); + this.queue = queue; + this.pending = pending; + this.cap = cap; + } + + /** The queue that rejected the enqueue. */ + public String queue() { + return queue; + } + + /** Pending count observed at rejection time. */ + public long pending() { + return pending; + } + + /** The configured cap. */ + public long cap() { + return cap; + } +} diff --git a/sdks/java/src/main/java/org/byteveda/taskito/internal/JniQueueBackend.java b/sdks/java/src/main/java/org/byteveda/taskito/internal/JniQueueBackend.java index a1172755..2c0d8acc 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/internal/JniQueueBackend.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/internal/JniQueueBackend.java @@ -97,6 +97,11 @@ public String statsByQueueJson(String queue) { return withOpenHandle(() -> NativeQueue.statsByQueue(handle, queue)); } + @Override + public long countPendingByQueue(String queue) { + return withOpenHandle(() -> NativeQueue.countPendingByQueue(handle, queue)); + } + @Override public String statsAllQueuesJson() { return withOpenHandle(() -> NativeQueue.statsAllQueues(handle)); diff --git a/sdks/java/src/main/java/org/byteveda/taskito/internal/NativeQueue.java b/sdks/java/src/main/java/org/byteveda/taskito/internal/NativeQueue.java index 46e3bb7f..0bf93b78 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/internal/NativeQueue.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/internal/NativeQueue.java @@ -43,6 +43,8 @@ private NativeQueue() {} public static native String statsByQueue(long handle, String queue); + public static native long countPendingByQueue(long handle, String queue); + public static native String statsAllQueues(long handle); public static native String listJobs(long handle, String filterJson); diff --git a/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java b/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java index 688569d4..9f973d34 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java @@ -42,6 +42,9 @@ public interface QueueBackend extends AutoCloseable { String statsByQueueJson(String queue); + /** Count pending jobs on {@code queue} — the primitive behind the {@code maxPending} cap. */ + long countPendingByQueue(String queue); + String statsAllQueuesJson(); String listJobsJson(String filterJson); diff --git a/sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java b/sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java new file mode 100644 index 00000000..8c6cc8aa --- /dev/null +++ b/sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java @@ -0,0 +1,79 @@ +package org.byteveda.taskito.core; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import java.nio.file.Path; +import org.byteveda.taskito.Taskito; +import org.byteveda.taskito.errors.QueueFullException; +import org.byteveda.taskito.task.EnqueueOptions; +import org.byteveda.taskito.task.Task; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.io.TempDir; + +/** S26 — opt-in {@code maxPending} admission cap. No worker runs, so jobs stay pending. */ +class AdmissionTest { + + private static Taskito sqlite(Path dir) { + return Taskito.builder() + .backend("sqlite") + .url(dir.resolve("t.db").toString()) + .open(); + } + + @Test + @Timeout(30) + void countsPendingPerQueue(@TempDir Path dir) { + Task noop = Task.of("noop", String.class); + try (Taskito queue = sqlite(dir)) { + assertEquals(0, queue.countPendingByQueue("default")); + queue.enqueue(noop, "a"); + queue.enqueue(noop, "b"); + assertEquals(2, queue.countPendingByQueue("default")); + } + } + + @Test + @Timeout(30) + void uncappedQueueNeverRejects(@TempDir Path dir) { + Task noop = Task.of("noop", String.class); + try (Taskito queue = sqlite(dir)) { + for (int i = 0; i < 25; i++) { + queue.enqueue(noop, String.valueOf(i)); + } + assertEquals(25, queue.countPendingByQueue("default")); + } + } + + @Test + @Timeout(30) + void rejectsAtCap(@TempDir Path dir) { + Task noop = Task.of("noop", String.class); + try (Taskito queue = sqlite(dir)) { + queue.maxPending("default", 2); + queue.enqueue(noop, "a"); + queue.enqueue(noop, "b"); + assertThrows(QueueFullException.class, () -> queue.enqueue(noop, "c")); + // Rejected enqueue inserted nothing. + assertEquals(2, queue.countPendingByQueue("default")); + } + } + + @Test + @Timeout(30) + void capIsPerQueue(@TempDir Path dir) { + Task noop = Task.of("noop", String.class); + try (Taskito queue = sqlite(dir)) { + queue.maxPending("tight", 1); + EnqueueOptions tight = EnqueueOptions.builder().queue("tight").build(); + queue.enqueue(noop, "a", tight); + assertThrows(QueueFullException.class, () -> queue.enqueue(noop, "b", tight)); + // A different, uncapped queue is unaffected. + EnqueueOptions wide = EnqueueOptions.builder().queue("wide").build(); + for (int i = 0; i < 5; i++) { + queue.enqueue(noop, String.valueOf(i), wide); + } + } + } +} diff --git a/sdks/java/test-support/src/main/java/org/byteveda/taskito/test/InMemoryQueueBackend.java b/sdks/java/test-support/src/main/java/org/byteveda/taskito/test/InMemoryQueueBackend.java index ed673c77..2d5539ee 100644 --- a/sdks/java/test-support/src/main/java/org/byteveda/taskito/test/InMemoryQueueBackend.java +++ b/sdks/java/test-support/src/main/java/org/byteveda/taskito/test/InMemoryQueueBackend.java @@ -182,6 +182,11 @@ public String statsByQueueJson(String queue) { return toJson(statsFor(queue)); } + @Override + public long countPendingByQueue(String queue) { + return count("pending", queue); + } + @Override public String statsAllQueuesJson() { Map out = new LinkedHashMap<>(); From 682eff3217940e1f95c89c61b80ea74f24d84126 Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 17:44:51 +0530 Subject: [PATCH 5/8] fix(python): size-aware max_pending, negative-cap guard, accumulator --- sdks/python/taskito/app.py | 37 +++++++--- sdks/python/taskito/mixins/runtime_config.py | 3 + sdks/python/tests/core/test_admission.py | 74 ++++++++++++++++---- 3 files changed, 90 insertions(+), 24 deletions(-) diff --git a/sdks/python/taskito/app.py b/sdks/python/taskito/app.py index 30c43aee..e552dbd0 100644 --- a/sdks/python/taskito/app.py +++ b/sdks/python/taskito/app.py @@ -21,6 +21,7 @@ import hashlib import logging import os +from collections import Counter from collections.abc import Callable, Sequence from concurrent.futures import ThreadPoolExecutor from typing import TYPE_CHECKING, Any @@ -313,6 +314,11 @@ def __init__( self._queue_configs: dict[str, dict[str, Any]] = {} # Opt-in per-queue admission caps (queue -> max pending). Empty = uncapped. self._max_pending: dict[str, int] = dict(max_pending or {}) + # A negative cap would make `pending + incoming > cap` always true and + # permanently reject every enqueue — reject it at construction instead. + for _queue, _cap in self._max_pending.items(): + if _cap < 0: + raise ValueError(f"max_pending for '{_queue}' must be non-negative") self._event_bus = EventBus(max_workers=event_workers) self._webhook_manager = WebhookManager(queue_ref=self) @@ -455,6 +461,9 @@ def _dispatch_batched_payload( del config # currently unused at dispatch time; reserved for telemetry payload = self._encode_payload(task_name, (items,), {}) self._check_payload_size(task_name, len(payload)) + # A flush is a single-job enqueue onto the default queue, so it honors the + # admission cap like any other producer path (no silent bypass). + self._reject_if_queue_full("default") py_job = self._inner.enqueue( task_name=task_name, payload=payload, @@ -520,21 +529,24 @@ def _serialize_result(self, task_name: str, result: Any) -> bytes: """ return self._serializer.dumps(result) - def _reject_if_queue_full(self, queue_name: str) -> None: + def _reject_if_queue_full(self, queue_name: str, incoming: int = 1) -> None: """Enforce the opt-in ``max_pending`` admission cap for a queue. - Raises :class:`QueueFullError` when the queue's pending backlog has - reached its configured cap. No-op (and no query) for uncapped queues. - Non-atomic count-then-insert — brief overshoot is accepted, like the - rate limiter. + Raises :class:`QueueFullError` when admitting ``incoming`` jobs would push + the queue's pending backlog past its configured cap. ``incoming`` is the + batch size, so a batch is rejected as a whole rather than overshooting the + cap by its full size. No-op (and no query) for uncapped queues. Non-atomic + count-then-insert — brief overshoot under concurrent producers is accepted, + like the rate limiter. """ cap = self._max_pending.get(queue_name) if cap is None: return pending = self._inner.count_pending_by_queue(queue_name) - if pending >= cap: + if pending + incoming > cap: raise QueueFullError( - f"queue '{queue_name}' is full: {pending} pending >= max_pending {cap}" + f"queue '{queue_name}' is full: {pending} pending + {incoming} " + f"would exceed max_pending {cap}" ) def enqueue( @@ -924,10 +936,13 @@ def enqueue_many( delay=delays[i], ) - # Admission cap is all-or-nothing for the batch: if any distinct target - # queue is at its cap, reject before inserting any row. - for capped_queue in set(queues_list): - self._reject_if_queue_full(capped_queue) + # Admission cap is all-or-nothing for the batch: reject before inserting + # any row if admitting this batch would push a target queue past its cap. + # Count the rows per queue so the check accounts for the whole batch + # rather than overshooting the cap by its size. + batch_by_queue = Counter(queues_list) + for capped_queue, incoming in batch_by_queue.items(): + self._reject_if_queue_full(capped_queue, incoming) py_jobs = self._inner.enqueue_batch( task_names=task_names, diff --git a/sdks/python/taskito/mixins/runtime_config.py b/sdks/python/taskito/mixins/runtime_config.py index 217f2f38..a17f040a 100644 --- a/sdks/python/taskito/mixins/runtime_config.py +++ b/sdks/python/taskito/mixins/runtime_config.py @@ -99,5 +99,8 @@ def set_queue_max_pending(self, queue_name: str, max_pending: int) -> None: Args: queue_name: Queue name (e.g. ``"default"``). max_pending: Maximum pending jobs allowed before enqueue is rejected. + Must be non-negative; ``0`` admits nothing. """ + if max_pending < 0: + raise ValueError("max_pending must be non-negative") self._max_pending[queue_name] = max_pending diff --git a/sdks/python/tests/core/test_admission.py b/sdks/python/tests/core/test_admission.py index c5bf49c2..372339e4 100644 --- a/sdks/python/tests/core/test_admission.py +++ b/sdks/python/tests/core/test_admission.py @@ -6,6 +6,7 @@ from __future__ import annotations +import time from pathlib import Path import pytest @@ -84,19 +85,66 @@ def test_enqueue_many_all_or_nothing(queue: Queue) -> None: assert queue._inner.count_pending_by_queue("default") == 3 -def test_cap_frees_after_drain(queue: Queue, run_worker: object) -> None: - """Once pending drains below the cap, enqueue is admitted again.""" - from taskito import JobResult +def test_enqueue_many_accounts_for_batch_size(queue: Queue) -> None: + _register(queue) + queue.set_queue_max_pending("default", 3) + # An empty queue but a batch bigger than the cap is rejected as a whole. + with pytest.raises(QueueFullError): + queue.enqueue_many("noop", [(), (), (), ()]) + assert queue._inner.count_pending_by_queue("default") == 0 + # A batch that exactly fits is admitted. + queue.enqueue_many("noop", [(), (), ()]) + assert queue._inner.count_pending_by_queue("default") == 3 + # Now full: even a single more is rejected. + with pytest.raises(QueueFullError): + queue.enqueue("noop") - results: list[JobResult] = [] - @queue.task(name="quick") - def quick() -> str: - return "ok" +def test_negative_cap_rejected(tmp_path: Path) -> None: + q = Queue(db_path=str(tmp_path / "n.db"), workers=1) + with pytest.raises(ValueError): + q.set_queue_max_pending("default", -1) + with pytest.raises(ValueError): + Queue(db_path=str(tmp_path / "n2.db"), workers=1, max_pending={"default": -1}) - queue.set_queue_max_pending("default", 100) - # With a worker draining, pending stays well under the cap. - for _ in range(20): - results.append(queue.enqueue("quick")) - for r in results: - assert r.result(timeout=10) == "ok" + +def test_cap_frees_after_drain(tmp_path: Path) -> None: + """Fill a small cap, observe rejection, then confirm draining re-admits. + + A gated task blocks the worker so the cap can actually be reached; releasing + the gate drains the backlog and a fresh enqueue must be admitted again. + """ + import threading + + q = Queue(db_path=str(tmp_path / "drain.db"), workers=1) + gate = threading.Event() + + @q.task(name="blocked") + def blocked() -> None: + gate.wait(timeout=10) + + q.set_queue_max_pending("default", 2) + # No worker yet: two enqueues fill the cap, the third is rejected. + q.enqueue("blocked") + q.enqueue("blocked") + with pytest.raises(QueueFullError): + q.enqueue("blocked") + + worker = threading.Thread(target=q.run_worker, daemon=True) + worker.start() + try: + # The worker claims a job (Running), so pending drops below the cap and a + # new enqueue is admitted where it was rejected a moment ago. + deadline = time.time() + 10 + admitted = False + while time.time() < deadline: + if q._inner.count_pending_by_queue("default") < 2: + q.enqueue("blocked") + admitted = True + break + time.sleep(0.05) + assert admitted, "draining below the cap must re-admit an enqueue" + finally: + gate.set() + q.shutdown() + worker.join(timeout=5) From 2d02998b44ec8291c3d2c89325f199786805e9a2 Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 17:45:21 +0530 Subject: [PATCH 6/8] fix(node): size-aware max_pending and negative-cap guard --- sdks/node/src/queue.ts | 30 ++++++++++++++++++++------- sdks/node/test/core/admission.test.ts | 19 +++++++++++++++++ 2 files changed, 41 insertions(+), 8 deletions(-) diff --git a/sdks/node/src/queue.ts b/sdks/node/src/queue.ts index 18319946..8c0963cf 100644 --- a/sdks/node/src/queue.ts +++ b/sdks/node/src/queue.ts @@ -387,6 +387,11 @@ export class Queue { /** Set per-queue concurrency / rate-limit applied when a worker runs. */ configureQueue(name: string, limits: QueueLimits): void { + // A negative cap would make `pending + incoming > cap` always true and + // permanently reject every enqueue — reject it at configuration time. + if (limits.maxPending !== undefined && limits.maxPending < 0) { + throw new RangeError("maxPending must be non-negative"); + } this.queueLimits.set(name, limits); } @@ -448,15 +453,17 @@ export class Queue { /** * Enforce the opt-in `maxPending` admission cap for a queue (set via - * {@link Queue.configureQueue}). Throws {@link QueueFullError} once the - * queue's pending backlog reaches its cap; a no-op (and no query) for - * uncapped queues. Non-atomic count-then-insert, like the rate limiter. + * {@link Queue.configureQueue}). Throws {@link QueueFullError} when admitting + * `incoming` jobs would push the queue's pending backlog past its cap; + * `incoming` is the batch size, so a batch is rejected as a whole rather than + * overshooting the cap by its full size. A no-op (and no query) for uncapped + * queues. Non-atomic count-then-insert, like the rate limiter. */ - private rejectIfQueueFull(queue: string): void { + private rejectIfQueueFull(queue: string, incoming = 1): void { const cap = this.queueLimits.get(queue)?.maxPending; if (cap === undefined) return; const pending = this.native.countPendingByQueue(queue); - if (pending >= cap) { + if (pending + incoming > cap) { throw new QueueFullError(queue, pending, cap); } } @@ -470,9 +477,16 @@ export class Queue { name: Name, jobs: ReadonlyArray<{ args?: Parameters; options?: EnqueueOptions }>, ): string[] { - // All-or-nothing: if any distinct target queue is at its cap, reject the batch. - for (const queue of new Set(jobs.map((job) => job.options?.queue ?? "default"))) { - this.rejectIfQueueFull(queue); + // All-or-nothing: reject the batch if admitting it would push a target queue + // past its cap. Count rows per queue so the check accounts for the whole + // batch rather than overshooting the cap by its size. + const perQueue = new Map(); + for (const job of jobs) { + const queue = job.options?.queue ?? "default"; + perQueue.set(queue, (perQueue.get(queue) ?? 0) + 1); + } + for (const [queue, incoming] of perQueue) { + this.rejectIfQueueFull(queue, incoming); } const prepared = jobs.map((job) => { const { payload, options } = this.prepareEnqueue(name, job.args, job.options, { diff --git a/sdks/node/test/core/admission.test.ts b/sdks/node/test/core/admission.test.ts index 7a472c9e..b9c5c372 100644 --- a/sdks/node/test/core/admission.test.ts +++ b/sdks/node/test/core/admission.test.ts @@ -59,6 +59,25 @@ it("enqueueMany is all-or-nothing against the cap", () => { expect(queue.countPendingByQueue("default")).toBe(3); }); +it("enqueueMany accounts for the batch size", () => { + const queue = newQueue(); + queue.task("noop", () => undefined); + queue.configureQueue("default", { maxPending: 3 }); + // Empty queue, but a batch bigger than the cap is rejected as a whole. + expect(() => queue.enqueueMany("noop", [{}, {}, {}, {}])).toThrow(QueueFullError); + expect(queue.countPendingByQueue("default")).toBe(0); + // A batch that exactly fits is admitted. + queue.enqueueMany("noop", [{}, {}, {}]); + expect(queue.countPendingByQueue("default")).toBe(3); + // Now full: one more is rejected. + expect(() => queue.enqueue("noop")).toThrow(QueueFullError); +}); + +it("configureQueue rejects a negative cap", () => { + const queue = newQueue(); + expect(() => queue.configureQueue("default", { maxPending: -1 })).toThrow(RangeError); +}); + it("QueueFullError is a QueueError", () => { const err = new QueueFullError("q", 5, 5); expect(err).toBeInstanceOf(QueueError); From 50a159b0622298e5f49745005448273334b31135 Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 17:45:39 +0530 Subject: [PATCH 7/8] fix(java): size-aware cap, negative-cap guard, optional SPI method --- .../org/byteveda/taskito/DefaultTaskito.java | 20 +++++++++++---- .../byteveda/taskito/spi/QueueBackend.java | 11 ++++++-- .../byteveda/taskito/core/AdmissionTest.java | 25 +++++++++++++++++++ 3 files changed, 49 insertions(+), 7 deletions(-) diff --git a/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java b/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java index 24d9738f..6550d137 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java @@ -150,24 +150,31 @@ public Map resourceMetrics() { @Override public Taskito maxPending(String queue, int cap) { + if (cap < 0) { + // A negative cap makes `pending + incoming > cap` always true and + // would permanently reject every enqueue for the queue. + throw new IllegalArgumentException("cap must be non-negative"); + } maxPending.put(queue, cap); return this; } /** * Enforce the opt-in {@code maxPending} admission cap for a queue. Throws - * {@link QueueFullException} once the queue's pending backlog reaches its - * cap; a no-op (and no query) for uncapped queues. Non-atomic + * {@link QueueFullException} when admitting {@code incoming} jobs would push + * the queue's pending backlog past its cap; {@code incoming} is the batch + * size, so a batch is rejected as a whole rather than overshooting the cap by + * its full size. A no-op (and no query) for uncapped queues. Non-atomic * count-then-insert, like the rate limiter. */ - private void rejectIfQueueFull(String queueOrNull) { + private void rejectIfQueueFull(String queueOrNull, int incoming) { String queue = queueOrNull == null ? "default" : queueOrNull; Integer cap = maxPending.get(queue); if (cap == null) { return; } long pending = backend.countPendingByQueue(queue); - if (pending >= cap) { + if (pending + incoming > cap) { throw new QueueFullException(queue, pending, cap); } } @@ -295,7 +302,7 @@ private Optional dispatchEnqueue( .build(); } // Admission cap: reject before serializing/inserting if the target queue is full. - rejectIfQueueFull(finalOptions.queue()); + rejectIfQueueFull(finalOptions.queue(), 1); // Serialize before codec-encoding so the idempotency key hashes the deterministic // pre-codec payload — a non-deterministic codec (e.g. an AES-GCM nonce) must not // change the dedup key. @@ -405,6 +412,9 @@ public List enqueueMany(Task task, List payloads, EnqueueOptio bytes[i] = encodeCodecs(payloadBytes, task.codecNames()); perJob.add(jobOptions); } + // The whole batch targets one queue (a single `options`); reject before + // inserting if admitting all of it would push that queue past its cap. + rejectIfQueueFull(options.queue(), payloads.size()); return Arrays.asList(backend.enqueueMany(task.name(), bytes, encode(perJob))); } diff --git a/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java b/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java index 9f973d34..47dd280e 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/spi/QueueBackend.java @@ -42,8 +42,15 @@ public interface QueueBackend extends AutoCloseable { String statsByQueueJson(String queue); - /** Count pending jobs on {@code queue} — the primitive behind the {@code maxPending} cap. */ - long countPendingByQueue(String queue); + /** + * Count pending jobs on {@code queue} — the primitive behind the + * {@code maxPending} cap. Defaults to unsupported so this optional capability + * doesn't break existing third-party {@code QueueBackend} implementations at + * compile time; it is only invoked when a queue actually has a cap set. + */ + default long countPendingByQueue(String queue) { + throw new UnsupportedOperationException("countPendingByQueue not supported by this backend"); + } String statsAllQueuesJson(); diff --git a/sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java b/sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java index 8c6cc8aa..7a1464c6 100644 --- a/sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java +++ b/sdks/java/src/test/java/org/byteveda/taskito/core/AdmissionTest.java @@ -60,6 +60,31 @@ void rejectsAtCap(@TempDir Path dir) { } } + @Test + @Timeout(30) + void enqueueManyAccountsForBatchSize(@TempDir Path dir) { + Task noop = Task.of("noop", String.class); + try (Taskito queue = sqlite(dir)) { + queue.maxPending("default", 3); + // Empty queue, but a batch bigger than the cap is rejected as a whole. + assertThrows( + QueueFullException.class, () -> queue.enqueueMany(noop, java.util.List.of("a", "b", "c", "d"))); + assertEquals(0, queue.countPendingByQueue("default")); + // A batch that exactly fits is admitted. + queue.enqueueMany(noop, java.util.List.of("a", "b", "c")); + assertEquals(3, queue.countPendingByQueue("default")); + // Now full: one more is rejected. + assertThrows(QueueFullException.class, () -> queue.enqueue(noop, "x")); + } + } + + @Test + void rejectsNegativeCap(@TempDir Path dir) { + try (Taskito queue = sqlite(dir)) { + assertThrows(IllegalArgumentException.class, () -> queue.maxPending("default", -1)); + } + } + @Test @Timeout(30) void capIsPerQueue(@TempDir Path dir) { From 3e0cebaf10f87b24c48f5d0e73f15a36421ba060 Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sat, 18 Jul 2026 18:00:17 +0530 Subject: [PATCH 8/8] test(python): assert worker stops in admission-cap cleanup --- sdks/python/tests/core/test_admission.py | 1 + 1 file changed, 1 insertion(+) diff --git a/sdks/python/tests/core/test_admission.py b/sdks/python/tests/core/test_admission.py index 372339e4..ae05fb07 100644 --- a/sdks/python/tests/core/test_admission.py +++ b/sdks/python/tests/core/test_admission.py @@ -148,3 +148,4 @@ def blocked() -> None: gate.set() q.shutdown() worker.join(timeout=5) + assert not worker.is_alive(), "worker did not stop during cleanup"