From 8e4c4469ee0981f0b04abfbe8991fceb4c52092c Mon Sep 17 00:00:00 2001 From: stromanni Date: Thu, 23 Jul 2026 20:48:48 +0530 Subject: [PATCH 1/4] feat(node): expose per-subscription topic stats Binds the core topic_backlog_stats aggregate, mirroring Python's topic_stats(). --- crates/taskito-node/src/convert/mod.rs | 4 +- crates/taskito-node/src/convert/pubsub.rs | 38 +++++++++++++ crates/taskito-node/src/queue/pubsub.rs | 20 +++++-- sdks/node/src/native.ts | 1 + sdks/node/src/queue.ts | 12 +++++ sdks/node/src/types.ts | 1 + sdks/node/test/core/pubsub.test.ts | 66 +++++++++++++++++++++++ 7 files changed, 137 insertions(+), 5 deletions(-) diff --git a/crates/taskito-node/src/convert/mod.rs b/crates/taskito-node/src/convert/mod.rs index 2c9eec59..d09a9332 100644 --- a/crates/taskito-node/src/convert/mod.rs +++ b/crates/taskito-node/src/convert/mod.rs @@ -23,8 +23,8 @@ pub use ops::{ pub use outcome::{outcome_to_js, JsOutcome}; pub use periodic::{periodic_to_js, JsPeriodicTask}; pub use pubsub::{ - subscription_to_js, topic_log_stat_to_js, topic_message_to_js, topic_to_js, JsSubscription, - JsTopic, JsTopicLogStat, JsTopicMessage, + subscription_to_js, topic_log_stat_to_js, topic_message_to_js, topic_stat_to_js, topic_to_js, + JsSubscription, JsTopic, JsTopicLogStat, JsTopicMessage, JsTopicStat, }; pub use stats::{ dead_job_to_js, job_error_to_js, metric_to_js, stats_to_js, status_code, worker_to_js, diff --git a/crates/taskito-node/src/convert/pubsub.rs b/crates/taskito-node/src/convert/pubsub.rs index 1d1ae068..fd42722f 100644 --- a/crates/taskito-node/src/convert/pubsub.rs +++ b/crates/taskito-node/src/convert/pubsub.rs @@ -3,6 +3,7 @@ use napi::bindgen_prelude::Buffer; use napi_derive::napi; use taskito_core::storage::records::{Subscription, Topic, TopicLogStats, TopicMessage}; +use taskito_core::storage::SubscriptionBacklogStats; /// A topic subscription: routes messages published to `topic` to `taskName` /// jobs on `queue`, one delivery per active subscription. @@ -28,6 +29,43 @@ pub fn subscription_to_js(row: Subscription) -> JsSubscription { } } +/// Backlog snapshot for one topic subscription. One entry per *registered* +/// subscription — durable or ephemeral, active or paused — even at zero +/// backlog, so a caller renders the full subscriber list from a single call. +#[napi(object)] +pub struct JsTopicStat { + pub topic: String, + pub subscription: String, + pub task_name: String, + pub queue: String, + pub active: bool, + pub durable: bool, + /// Deliveries waiting to run. + pub pending: i64, + /// Deliveries currently executing. + pub running: i64, + /// Deliveries in the dead-letter queue. + pub dead: i64, + /// Age (ms) of the oldest still-pending delivery, absent at zero backlog. + pub oldest_pending_age_ms: Option, +} + +/// Convert a core [`SubscriptionBacklogStats`] into its JS-facing shape. +pub fn topic_stat_to_js(stat: SubscriptionBacklogStats) -> JsTopicStat { + JsTopicStat { + topic: stat.topic, + subscription: stat.subscription_name, + task_name: stat.task_name, + queue: stat.queue, + active: stat.active, + durable: stat.durable, + pending: stat.pending, + running: stat.running, + dead: stat.dead, + oldest_pending_age_ms: stat.oldest_pending_age_ms, + } +} + /// A message pulled from a log topic. `id` is the cursor token to pass to /// `ackTopic`; `payload` is the opaque published bytes. #[napi(object)] diff --git a/crates/taskito-node/src/queue/pubsub.rs b/crates/taskito-node/src/queue/pubsub.rs index 3b5fc242..ff741070 100644 --- a/crates/taskito-node/src/queue/pubsub.rs +++ b/crates/taskito-node/src/queue/pubsub.rs @@ -12,9 +12,9 @@ use taskito_core::Storage; use super::JsQueue; use crate::config::PublishOptions; use crate::convert::{ - job_to_js, subscription_to_js, topic_log_stat_to_js, topic_message_to_js, topic_to_js, JsJob, - JsSubscription, JsTopic, JsTopicLogStat, JsTopicMessage, DEFAULT_MAX_RETRIES, DEFAULT_PRIORITY, - DEFAULT_TIMEOUT_MS, + job_to_js, subscription_to_js, topic_log_stat_to_js, topic_message_to_js, topic_stat_to_js, + topic_to_js, JsJob, JsSubscription, JsTopic, JsTopicLogStat, JsTopicMessage, JsTopicStat, + DEFAULT_MAX_RETRIES, DEFAULT_PRIORITY, DEFAULT_TIMEOUT_MS, }; use crate::error::{invalid_arg, join_to_napi_err, non_negative, to_napi_err}; @@ -274,6 +274,20 @@ impl JsQueue { .map_err(join_to_napi_err)? } + /// Backlog snapshot per registered subscription — one entry each, even at + /// zero backlog. Counts are aggregated live off the delivery-attribution + /// indexes, so they can never drift the way a maintained counter would. + #[napi] + pub async fn topic_backlog_stats(&self) -> Result> { + let storage = self.storage.clone(); + spawn_blocking(move || { + let stats = storage.topic_backlog_stats().map_err(to_napi_err)?; + Ok(stats.into_iter().map(topic_stat_to_js).collect()) + }) + .await + .map_err(join_to_napi_err)? + } + /// Drop ephemeral subscriptions whose owning worker is gone. Runs on the /// heartbeat cadence. Returns the number of subscriptions removed. /// diff --git a/sdks/node/src/native.ts b/sdks/node/src/native.ts index 342dcb75..09d6bf55 100644 --- a/sdks/node/src/native.ts +++ b/sdks/node/src/native.ts @@ -38,6 +38,7 @@ export type { JsTopic, JsTopicLogStat, JsTopicMessage, + JsTopicStat, JsWorkerRow, JsWorkflowAdvance, JsWorkflowNode, diff --git a/sdks/node/src/queue.ts b/sdks/node/src/queue.ts index 90066842..518d33f1 100644 --- a/sdks/node/src/queue.ts +++ b/sdks/node/src/queue.ts @@ -86,6 +86,7 @@ import type { TaskOptions, TopicLogStat, TopicMessage, + TopicStat, WorkerInfo, WorkerRunOptions, } from "./types"; @@ -635,6 +636,17 @@ export class Queue { return this.native.listSubscriptions(topic); } + /** + * Backlog snapshot per subscription, optionally filtered to one `topic`. Every + * registered subscription appears — paused or ephemeral ones included — even + * with nothing queued, so the full subscriber list comes from one call. + * Counts are computed live off indexed columns, so this is safe to poll. + */ + async topicStats(topic?: string): Promise { + const stats = await this.native.topicBacklogStats(); + return topic === undefined ? stats : stats.filter((stat) => stat.topic === topic); + } + /** * Register a durable **log** subscription: a named cursor over `topic`. Unlike * `subscriber`, it has no handler — the topic's publishes are stored once each diff --git a/sdks/node/src/types.ts b/sdks/node/src/types.ts index 87ef9090..504141e6 100644 --- a/sdks/node/src/types.ts +++ b/sdks/node/src/types.ts @@ -24,6 +24,7 @@ export type { JsTaskLog as TaskLog, JsTopic as DeclaredTopic, JsTopicLogStat as TopicLogStat, + JsTopicStat as TopicStat, JsWorkerRow as WorkerInfo, MeshWorkerConfig, } from "./native"; diff --git a/sdks/node/test/core/pubsub.test.ts b/sdks/node/test/core/pubsub.test.ts index 56e4ad7c..c59d4b5c 100644 --- a/sdks/node/test/core/pubsub.test.ts +++ b/sdks/node/test/core/pubsub.test.ts @@ -272,3 +272,69 @@ it("reapEphemeralSubscriptions spares fresh rows even with a dead owner", async new Set(["durable_task", "ephemeral_task"]), ); }); + +it("topicStats reports one backlog row per subscription", async () => { + const queue = newQueue(); + queue.subscriber("orders", "send_email", () => undefined, { subscriptionName: "email" }); + queue.subscriber("orders", "track_order", () => undefined, { subscriptionName: "analytics" }); + await queue.declareSubscriptions(); + await queue.publish("orders", [1]); + await queue.publish("orders", [2]); + + const stats = new Map((await queue.topicStats("orders")).map((s) => [s.subscription, s])); + expect(new Set(stats.keys())).toEqual(new Set(["email", "analytics"])); + const email = stats.get("email"); + expect(email?.taskName).toBe("send_email"); + expect(email?.queue).toBe("default"); + expect(email?.active).toBe(true); + expect(email?.durable).toBe(true); + expect(email?.pending).toBe(2); + expect(email?.running).toBe(0); + expect(email?.dead).toBe(0); + expect(email?.oldestPendingAgeMs).toBeGreaterThanOrEqual(0); + expect(stats.get("analytics")?.pending).toBe(2); +}); + +it("topicStats reports zeros for an idle subscription", async () => { + const queue = newQueue(); + queue.subscriber("orders", "send_email", () => undefined, { subscriptionName: "email" }); + await queue.declareSubscriptions(); + + const [stat] = await queue.topicStats("orders"); + expect(stat?.pending).toBe(0); + expect(stat?.running).toBe(0); + expect(stat?.dead).toBe(0); + expect(stat?.oldestPendingAgeMs).toBeUndefined(); +}); + +it("topicStats counts a failed delivery as dead", async () => { + const queue = newQueue(); + queue.subscriber( + "orders", + "flaky", + () => { + throw new Error("boom"); + }, + { subscriptionName: "flaky", maxRetries: 0 }, + ); + await queue.declareSubscriptions(); + await queue.publish("orders", [1]); + + worker = queue.runWorker(); + const deadCounted = await waitFor(async () => { + const [stat] = await queue.topicStats("orders"); + return stat?.dead === 1 && stat.pending === 0; + }); + expect(deadCounted).toBe(true); +}); + +it("topicStats filters by topic and ignores non-pubsub jobs", async () => { + const queue = newQueue(); + queue.task("plain", () => undefined); + queue.subscriber("orders", "send_email", () => undefined, { subscriptionName: "email" }); + await queue.declareSubscriptions(); + queue.enqueue("plain"); + + expect((await queue.topicStats()).map((s) => s.subscription)).toEqual(["email"]); + expect(await queue.topicStats("other-topic")).toEqual([]); +}); From e910c98d2d1a299afa8c8858a97ce3a591634cbc Mon Sep 17 00:00:00 2001 From: stromanni Date: Thu, 23 Jul 2026 20:48:53 +0530 Subject: [PATCH 2/4] feat(java): expose per-subscription topic stats Binds the core topic_backlog_stats aggregate, mirroring Python's topic_stats(). --- crates/taskito-java/src/convert.rs | 37 ++++++- crates/taskito-java/src/queue/pubsub.rs | 20 +++- .../org/byteveda/taskito/DefaultTaskito.java | 17 +++ .../java/org/byteveda/taskito/Taskito.java | 12 +++ .../taskito/internal/JniQueueBackend.java | 5 + .../taskito/internal/NativeQueue.java | 3 + .../org/byteveda/taskito/model/TopicStat.java | 62 +++++++++++ .../byteveda/taskito/spi/QueueBackend.java | 5 + .../org/byteveda/taskito/core/PubSubTest.java | 100 ++++++++++++++++++ 9 files changed, 259 insertions(+), 2 deletions(-) create mode 100644 sdks/java/src/main/java/org/byteveda/taskito/model/TopicStat.java diff --git a/crates/taskito-java/src/convert.rs b/crates/taskito-java/src/convert.rs index 1af3e88a..ac6f4356 100644 --- a/crates/taskito-java/src/convert.rs +++ b/crates/taskito-java/src/convert.rs @@ -12,7 +12,7 @@ use taskito_core::storage::records::{ CircuitBreakerState, JobError, LockInfo, PeriodicTask, ReplayEntry, Subscription, TaskLogEntry, TaskMetric, Topic, TopicLogStats, TopicMessage, WorkerInfo, }; -use taskito_core::storage::{DeadJob, QueueStats}; +use taskito_core::storage::{DeadJob, QueueStats, SubscriptionBacklogStats}; use crate::error::BindingError; @@ -294,6 +294,41 @@ impl From<&TopicLogStats> for TopicLogStatsView { } } +/// Java-facing backlog snapshot for one topic subscription. One entry per +/// *registered* subscription — paused and ephemeral ones included — even at zero +/// backlog. `oldest_pending_age_ms` is null when nothing is pending. +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct TopicStatsView<'a> { + pub topic: &'a str, + pub subscription: &'a str, + pub task_name: &'a str, + pub queue: &'a str, + pub active: bool, + pub durable: bool, + pub pending: i64, + pub running: i64, + pub dead: i64, + pub oldest_pending_age_ms: Option, +} + +impl<'a> From<&'a SubscriptionBacklogStats> for TopicStatsView<'a> { + fn from(s: &'a SubscriptionBacklogStats) -> Self { + Self { + topic: &s.topic, + subscription: &s.subscription_name, + task_name: &s.task_name, + queue: &s.queue, + active: s.active, + durable: s.durable, + pending: s.pending, + running: s.running, + dead: s.dead, + oldest_pending_age_ms: s.oldest_pending_age_ms, + } + } +} + /// Java-facing view of a declared topic. `retention_ms` is null when the backlog /// is kept until consumed; `created_at` is Unix milliseconds. #[derive(Serialize)] diff --git a/crates/taskito-java/src/queue/pubsub.rs b/crates/taskito-java/src/queue/pubsub.rs index 40df802a..06973638 100644 --- a/crates/taskito-java/src/queue/pubsub.rs +++ b/crates/taskito-java/src/queue/pubsub.rs @@ -15,7 +15,7 @@ use taskito_core::Storage; use crate::backend; use crate::convert::{ build_publish_request, parse_json, to_json, JobView, PublishOptions, SubscriptionView, - TopicLogStatsView, TopicMessageView, TopicView, + TopicLogStatsView, TopicMessageView, TopicStatsView, TopicView, }; use crate::ffi::{guard, new_string, read_bytes, read_optional_string, read_string}; @@ -154,6 +154,24 @@ pub extern "system" fn Java_org_byteveda_taskito_internal_NativeQueue_setSubscri }) } +/// `String topicBacklogStats(long handle)` — a JSON array of `TopicStatsView`, +/// one backlog snapshot per registered subscription (even at zero backlog). +/// Counts are aggregated live off the delivery-attribution indexes, so they can +/// never drift the way a maintained counter would. +#[no_mangle] +pub extern "system" fn Java_org_byteveda_taskito_internal_NativeQueue_topicBacklogStats<'local>( + mut env: JNIEnv<'local>, + _class: JClass<'local>, + handle: jlong, +) -> jstring { + guard(&mut env, std::ptr::null_mut(), |env| { + let queue = unsafe { borrow_queue(handle) }; + let stats = queue.storage.topic_backlog_stats()?; + let views: Vec = stats.iter().map(TopicStatsView::from).collect(); + new_string(env, to_json(&views)?) + }) +} + /// `long reapEphemeralSubscriptions(long handle)` — drop ephemeral subscriptions /// whose owning worker is gone; returns the count removed. #[no_mangle] 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 ef5f7414..8584461e 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java @@ -52,6 +52,7 @@ import org.byteveda.taskito.model.Topic; import org.byteveda.taskito.model.TopicLogStat; import org.byteveda.taskito.model.TopicMessage; +import org.byteveda.taskito.model.TopicStat; import org.byteveda.taskito.model.WorkerInfo; import org.byteveda.taskito.model.WorkflowRunInfo; import org.byteveda.taskito.predicates.EnqueueDecision; @@ -924,6 +925,22 @@ public List listSubscriptions(String topic) { return decodeList(backend.listSubscriptionsJson(topic), Subscription.class); } + @Override + public List topicStats() { + return decodeList(backend.topicBacklogStatsJson(), TopicStat.class); + } + + @Override + public List topicStats(String topic) { + List filtered = new ArrayList<>(); + for (TopicStat stat : topicStats()) { + if (stat.topic.equals(topic)) { + filtered.add(stat); + } + } + return filtered; + } + @Override public List listTopics() { Set topics = new LinkedHashSet<>(); 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 d949b168..a4ef7037 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java @@ -41,6 +41,7 @@ import org.byteveda.taskito.model.Topic; import org.byteveda.taskito.model.TopicLogStat; import org.byteveda.taskito.model.TopicMessage; +import org.byteveda.taskito.model.TopicStat; import org.byteveda.taskito.model.WorkerInfo; import org.byteveda.taskito.model.WorkflowRunInfo; import org.byteveda.taskito.predicates.EnqueueGate; @@ -459,6 +460,17 @@ static Builder builder() { /** One topic's active subscriptions. */ List listSubscriptions(String topic); + /** + * Backlog snapshot per subscription, across all topics. Every registered + * subscription appears — paused and ephemeral ones included — even with nothing + * queued, so the full subscriber list comes from one call. Counts are computed + * live off indexed columns, so this is safe to poll. + */ + List topicStats(); + + /** As {@link #topicStats()}, filtered to one topic. */ + List topicStats(String topic); + /** Distinct topics that currently have at least one subscription. */ List listTopics(); 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 213df5d6..4678c9f4 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 @@ -393,6 +393,11 @@ public boolean setSubscriptionActive(String topic, String subscriptionName, bool return withOpenHandle(() -> NativeQueue.setSubscriptionActive(handle, topic, subscriptionName, active)); } + @Override + public String topicBacklogStatsJson() { + return withOpenHandle(() -> NativeQueue.topicBacklogStats(handle)); + } + @Override public long reapEphemeralSubscriptions() { return withOpenHandle(() -> NativeQueue.reapEphemeralSubscriptions(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 6e9dcb83..1430f025 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 @@ -178,6 +178,9 @@ public static native void registerSubscription( public static native boolean setSubscriptionActive( long handle, String topic, String subscriptionName, boolean active); + /** A JSON array of per-subscription backlog snapshots, one per registered subscription. */ + public static native String topicBacklogStats(long handle); + /** Drop ephemeral subscriptions whose owning worker is gone; returns the count removed. */ public static native long reapEphemeralSubscriptions(long handle); diff --git a/sdks/java/src/main/java/org/byteveda/taskito/model/TopicStat.java b/sdks/java/src/main/java/org/byteveda/taskito/model/TopicStat.java new file mode 100644 index 00000000..b413ece0 --- /dev/null +++ b/sdks/java/src/main/java/org/byteveda/taskito/model/TopicStat.java @@ -0,0 +1,62 @@ +package org.byteveda.taskito.model; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonProperty; + +/** Backlog snapshot for one topic subscription: how much of its fan-out is outstanding. */ +@JsonIgnoreProperties(ignoreUnknown = true) +public final class TopicStat { + public final String topic; + + /** Stable subscription identity; unique per topic. */ + public final String subscription; + + /** Task enqueued for each published message. */ + public final String taskName; + + /** Queue deliveries are enqueued into. */ + public final String queue; + + /** Whether the subscription currently receives deliveries (false = paused). */ + public final boolean active; + + /** Whether the registration persists across restarts (false = ephemeral). */ + public final boolean durable; + + /** Deliveries waiting to run. */ + public final long pending; + + /** Deliveries currently executing. */ + public final long running; + + /** Deliveries in the dead-letter queue. */ + public final long dead; + + /** Age (ms) of the oldest still-pending delivery, or {@code null} at zero backlog. */ + public final Long oldestPendingAgeMs; + + @JsonCreator + public TopicStat( + @JsonProperty("topic") String topic, + @JsonProperty("subscription") String subscription, + @JsonProperty("taskName") String taskName, + @JsonProperty("queue") String queue, + @JsonProperty("active") boolean active, + @JsonProperty("durable") boolean durable, + @JsonProperty("pending") long pending, + @JsonProperty("running") long running, + @JsonProperty("dead") long dead, + @JsonProperty("oldestPendingAgeMs") Long oldestPendingAgeMs) { + this.topic = topic; + this.subscription = subscription; + this.taskName = taskName; + this.queue = queue; + this.active = active; + this.durable = durable; + this.pending = pending; + this.running = running; + this.dead = dead; + this.oldestPendingAgeMs = oldestPendingAgeMs; + } +} 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 b12f1447..1cb9381f 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 @@ -294,6 +294,11 @@ default boolean setSubscriptionActive(String topic, String subscriptionName, boo throw new UnsupportedOperationException(PUBSUB_UNSUPPORTED); } + /** A JSON array of per-subscription backlog snapshots, one per registered subscription. */ + default String topicBacklogStatsJson() { + throw new UnsupportedOperationException(PUBSUB_UNSUPPORTED); + } + /** Drop ephemeral subscriptions whose owning worker is gone; returns the count removed. */ default long reapEphemeralSubscriptions() { throw new UnsupportedOperationException(PUBSUB_UNSUPPORTED); diff --git a/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java b/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java index 9fc28ece..f7c0b2fd 100644 --- a/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java +++ b/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java @@ -2,6 +2,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -10,14 +12,20 @@ import java.time.Duration; import java.util.ArrayList; import java.util.Collections; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import org.byteveda.taskito.Taskito; import org.byteveda.taskito.TaskitoException; +import org.byteveda.taskito.events.EventName; import org.byteveda.taskito.internal.JniQueueBackend; import org.byteveda.taskito.model.Job; import org.byteveda.taskito.model.JobStatus; import org.byteveda.taskito.model.Subscription; +import org.byteveda.taskito.model.TopicStat; import org.byteveda.taskito.pubsub.PublishOptions; import org.byteveda.taskito.pubsub.SubscriptionOptions; import org.byteveda.taskito.task.Task; @@ -303,6 +311,98 @@ void unsubscribeDropsTheLocalDeclarationSoWorkerStartDoesNotResurrectIt(@TempDir } } + @Test + void topicStatsReportsOneBacklogRowPerSubscription(@TempDir Path dir) { + try (Taskito queue = open(dir)) { + queue.subscribe( + "orders", + SEND_EMAIL, + SubscriptionOptions.builder().name("email").build()); + queue.subscribe( + "orders", + TRACK_ORDER, + SubscriptionOptions.builder().name("analytics").build()); + queue.publish("orders", "o-1"); + queue.publish("orders", "o-2"); + + Map stats = bySubscription(queue.topicStats("orders")); + assertEquals(Set.of("email", "analytics"), stats.keySet()); + TopicStat email = stats.get("email"); + assertEquals("orders", email.topic); + assertEquals(SEND_EMAIL.name(), email.taskName); + assertEquals("default", email.queue); + assertTrue(email.active); + assertTrue(email.durable); + assertEquals(2, email.pending); + assertEquals(0, email.running); + assertEquals(0, email.dead); + assertNotNull(email.oldestPendingAgeMs); + assertTrue(email.oldestPendingAgeMs >= 0); + assertEquals(2, stats.get("analytics").pending); + } + } + + @Test + void topicStatsReportsZerosForAnIdleSubscription(@TempDir Path dir) { + try (Taskito queue = open(dir)) { + queue.subscribe("orders", SEND_EMAIL); + + List stats = queue.topicStats("orders"); + assertEquals(1, stats.size()); + TopicStat idle = stats.get(0); + assertEquals(0, idle.pending); + assertEquals(0, idle.running); + assertEquals(0, idle.dead); + assertNull(idle.oldestPendingAgeMs); + } + } + + @Test + void topicStatsFiltersByTopicAndIgnoresNonPubSubJobs(@TempDir Path dir) { + try (Taskito queue = open(dir)) { + queue.subscribe( + "orders", + SEND_EMAIL, + SubscriptionOptions.builder().name("email").build()); + queue.enqueue(Task.of("pubsub.plain", String.class), "p-1"); + + assertEquals(1, queue.topicStats().size()); + assertEquals("email", queue.topicStats().get(0).subscription); + assertTrue(queue.topicStats("other-topic").isEmpty()); + } + } + + @Test + @Timeout(30) + void topicStatsCountsAFailedDeliveryAsDead(@TempDir Path dir) throws Exception { + Task flaky = Task.of("pubsub.flaky", String.class).maxRetries(0); + try (Taskito queue = open(dir)) { + queue.subscribe( + "orders", flaky, SubscriptionOptions.builder().name("flaky").build()); + + CountDownLatch dead = new CountDownLatch(1); + try (Worker worker = queue.worker() + .handle(flaky, (String payload) -> { + throw new IllegalStateException("boom"); + }) + .on(EventName.DEAD, event -> dead.countDown()) + .start()) { + queue.publish("orders", "o-1"); + assertTrue(dead.await(20, TimeUnit.SECONDS), "the delivery should dead-letter"); + } + + TopicStat stat = queue.topicStats("orders").get(0); + assertEquals(1, stat.dead); + assertEquals(0, stat.pending); + } + } + + private static Map bySubscription(List stats) { + Map bySubscription = new LinkedHashMap<>(); + stats.forEach(stat -> bySubscription.put(stat.subscription, stat)); + return bySubscription; + } + private static List ids(List jobs) { List ids = new ArrayList<>(); jobs.forEach(job -> ids.add(job.id)); From 5fa2356a6c23a19c11f4c0b83b9729a107fccbdb Mon Sep 17 00:00:00 2001 From: stromanni Date: Thu, 23 Jul 2026 20:48:53 +0530 Subject: [PATCH 3/4] docs: document topicStats in the Node and Java SDKs --- CHANGELOG.md | 4 ++++ docs/content/docs/java/api-reference/queue/pubsub.mdx | 1 + docs/content/docs/node/api-reference/queue/pubsub.mdx | 1 + 3 files changed, 6 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index da20da37..e464bf39 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -27,6 +27,10 @@ underlying Rust crates are released together, in lock-step. - **Node `request` resource scope.** A fresh instance on every `useResource()` call, each disposed when the task ends — matching the Java scope of the same name. +- **Pub/sub backlog stats in the Node and Java SDKs.** `queue.topicStats(topic?)` (Node) and + `topicStats()` / `topicStats(topic)` (Java) return the per-subscription snapshot Python has + as `topic_stats()`: `pending`, `running`, `dead`, and `oldestPendingAgeMs`, with a row for + every registered subscription even at zero backlog. - Job outcome events report how long the task ran: `durationMs()` on Java's `OutcomeEvent`, `durationMs` on Node's, and `duration_ms` on Python's job event payloads. Java also gains `NodeSnapshot.durationMs()` / `compensationDurationMs()` and `TaskContext.elapsedMs()`. diff --git a/docs/content/docs/java/api-reference/queue/pubsub.mdx b/docs/content/docs/java/api-reference/queue/pubsub.mdx index c2652181..cd187183 100644 --- a/docs/content/docs/java/api-reference/queue/pubsub.mdx +++ b/docs/content/docs/java/api-reference/queue/pubsub.mdx @@ -14,6 +14,7 @@ semantics, lifecycle, cross-SDK topics). | `pauseSubscription(topic, name)` / `resumeSubscription(topic, name)` | Stop/resume deliveries without unregistering; `false` if none matched. | | `listSubscriptions()` / `listSubscriptions(topic)` | Every subscription, or one topic's active ones. | | `listTopics()` | Distinct topics with at least one subscription. | +| `topicStats()` / `topicStats(topic)` → `List` | Backlog snapshot per subscription, across all topics or filtered to one: `topic`, `subscription`, `taskName`, `queue`, `active`, `durable`, `pending`, `running`, `dead`, `oldestPendingAgeMs`. Every registered subscription appears — paused and ephemeral ones included — even at zero backlog. Computed live off indexed columns, so it is safe to poll. | | `subscribeLog(topic, name)` | Register a durable log subscription — a named cursor with no handler. Writes immediately, so register it before the publishes it should see. | | `logConsumer(topic, name, payloadType, handler)` / `logConsumer(topic, name, payloadType, handler, options)` | Register a **managed consumer**: the durable log subscription plus, once a worker runs, a daemon thread that pulls messages, decodes each into `payloadType`, invokes `handler`, and advances the cursor. `LogConsumerOptions` sets `pollIntervalMs` (default 1000), `batchSize` (default 100), and `onError` (`"retry"` leaves a failed message un-acked to re-read, default; `"skip"` acks past it). | | `declareTopic(name)` / `declareTopic(name, retention)` | Declare a log topic so its publishes are retained even with no subscriber (removing the late-join boundary). `retention` (a `Duration`) bounds a sub-less backlog. Idempotent. | diff --git a/docs/content/docs/node/api-reference/queue/pubsub.mdx b/docs/content/docs/node/api-reference/queue/pubsub.mdx index 7b89f410..364ec812 100644 --- a/docs/content/docs/node/api-reference/queue/pubsub.mdx +++ b/docs/content/docs/node/api-reference/queue/pubsub.mdx @@ -15,6 +15,7 @@ semantics, lifecycle, cross-SDK topics). | `pauseSubscription(topic, name)` / `resumeSubscription(topic, name)` → `Promise` | Stop/resume deliveries without unregistering. | | `listSubscriptions(topic?)` → `Promise` | All subscriptions, or one topic's active ones. | | `listTopics()` → `Promise` | Distinct topics with at least one subscription. | +| `topicStats(topic?)` → `Promise` | Backlog snapshot per subscription, optionally filtered to one topic: `topic`, `subscription`, `taskName`, `queue`, `active`, `durable`, `pending`, `running`, `dead`, `oldestPendingAgeMs`. Every registered subscription appears — paused and ephemeral ones included — even at zero backlog. Computed live off indexed columns, so it is safe to poll. | | `reapEphemeralSubscriptions()` → `Promise` | Drop ephemeral subscriptions whose owning worker is gone. Runs automatically on the worker heartbeat cadence; exposed for operational tooling. | | `subscribeLog(topic, name)` → `Promise` | Register a durable log subscription — a named cursor with no handler. Writes immediately, so register it before the publishes it should see. | | `logConsumer(topic, name, handler, opts?)` → `this` | Register a **managed consumer**: the durable log subscription plus, once a worker runs, a poll loop that pulls messages, calls `handler(...args)` per message (the handler may return a `Promise` — it is awaited), and advances the cursor. `opts`: `pollIntervalMs` (default 1000), `batchSize` (default 100), `onError` (`"retry"` leaves a failed message un-acked to re-read, default; `"skip"` acks past it). | From 8f54100f81df8fdb5ababa6b1845d57f89cd5a15 Mon Sep 17 00:00:00 2001 From: stromanni Date: Thu, 23 Jul 2026 21:28:07 +0530 Subject: [PATCH 4/4] fix(java): treat a null topicStats filter as no filter --- .../src/main/java/org/byteveda/taskito/DefaultTaskito.java | 4 ++++ sdks/java/src/main/java/org/byteveda/taskito/Taskito.java | 2 +- .../src/test/java/org/byteveda/taskito/core/PubSubTest.java | 2 ++ 3 files changed, 7 insertions(+), 1 deletion(-) 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 8584461e..eeb4e295 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/DefaultTaskito.java @@ -932,6 +932,10 @@ public List topicStats() { @Override public List topicStats(String topic) { + // A null topic is "no filter", matching listSubscriptions(null). + if (topic == null) { + return topicStats(); + } List filtered = new ArrayList<>(); for (TopicStat stat : topicStats()) { if (stat.topic.equals(topic)) { 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 a4ef7037..c1d59b01 100644 --- a/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java +++ b/sdks/java/src/main/java/org/byteveda/taskito/Taskito.java @@ -468,7 +468,7 @@ static Builder builder() { */ List topicStats(); - /** As {@link #topicStats()}, filtered to one topic. */ + /** As {@link #topicStats()}, filtered to one topic; a {@code null} topic means no filter. */ List topicStats(String topic); /** Distinct topics that currently have at least one subscription. */ diff --git a/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java b/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java index f7c0b2fd..d62a6282 100644 --- a/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java +++ b/sdks/java/src/test/java/org/byteveda/taskito/core/PubSubTest.java @@ -369,6 +369,8 @@ void topicStatsFiltersByTopicAndIgnoresNonPubSubJobs(@TempDir Path dir) { assertEquals(1, queue.topicStats().size()); assertEquals("email", queue.topicStats().get(0).subscription); assertTrue(queue.topicStats("other-topic").isEmpty()); + // A null topic is "no filter", not "a topic nothing matches". + assertEquals(1, queue.topicStats(null).size()); } }