diff --git a/samples/BackWave.Sample.Api/Program.cs b/samples/BackWave.Sample.Api/Program.cs
index 8c97c17..5f862f8 100644
--- a/samples/BackWave.Sample.Api/Program.cs
+++ b/samples/BackWave.Sample.Api/Program.cs
@@ -641,6 +641,18 @@ await m.GetJobAsync(id) is { } snapshot ? Results.Ok(snapshot) : Results.NotFoun
Results.Ok(await m.ListJobsAsync(new JobQuery { State = state, Queue = queue })))
.WithSummary("List jobs, optionally filtered by state and/or queue.");
+monitor.MapGet("/jobs/count", async (BackWaveMonitor m, JobState? state, string? queue, string? tenant) =>
+ Results.Ok(new
+ {
+ count = await m.GetJobCountAsync(new JobQuery
+ {
+ State = state,
+ Queue = queue,
+ TagPredicates = tenant is null ? [] : [JobTagPredicate.HasKeyValue("tenant", tenant)],
+ }),
+ }))
+ .WithSummary("Count the jobs that match the filters: ?state=, ?queue=, ?tenant=. Not capped by the page size.");
+
monitor.MapGet("/queues", async (BackWaveMonitor m) => Results.Ok(await m.GetQueueDepthsAsync()))
.WithSummary("Queue depths: job counts by queue and state.");
diff --git a/samples/BackWave.Sample.Api/README.md b/samples/BackWave.Sample.Api/README.md
index 0535f86..fa2adec 100644
--- a/samples/BackWave.Sample.Api/README.md
+++ b/samples/BackWave.Sample.Api/README.md
@@ -131,6 +131,8 @@ enqueued atomically (ADR 0023). Each renders as a graph at `/backwave/workflows/
**Monitor** (`/monitor`)
- `GET /monitor/jobs/{id}` · `GET /monitor/jobs?state=&queue=` · `GET /monitor/queues`
+- `GET /monitor/jobs/count?state=&queue=&tenant=` - the number of jobs that match the filters.
+ The page size does not cap it.
- `GET /monitor/workflows` — every Workflow with its derived status + member count.
- `GET /monitor/workflows/{id}` — one Workflow's graph: members, structural Dependency edges, status.
- `GET /monitor/tagged?tenant=&label=` — filter jobs by Tags; predicates **AND** (a Keyed
diff --git a/src/BackWave.Conformance/ConformanceSuite.cs b/src/BackWave.Conformance/ConformanceSuite.cs
index 0750340..daaa74f 100644
--- a/src/BackWave.Conformance/ConformanceSuite.cs
+++ b/src/BackWave.Conformance/ConformanceSuite.cs
@@ -1795,6 +1795,89 @@ public async Task Clause_5_9_Reads_FilterAndCount_OverCommittedEffectsOnly()
Assert.Contains(new QueueStateCount("other", JobState.Scheduled, 1), counts);
}
+ ///
+ /// Certifies that the filtered count equals the size of the filtered listing for every filter kind -
+ /// state, queue, wire name, schedule id, each tag predicate shape, and a combination - and is zero
+ /// for a filter that matches nothing.
+ ///
+ [Fact]
+ public async Task Clause_5_9_CountMatchingJobs_EqualsTheListing_ForEveryFilterKind()
+ {
+ var store = await CreateStoreAsync();
+ await store.EnqueueAsync(Job() with { Tags = JobTags.Empty.WithTag("tenant", "acme").WithLabel("urgent") }, T0);
+ await store.EnqueueAsync(Job() with { Tags = JobTags.Empty.WithTag("tenant", "globex") }, T0);
+ await store.EnqueueAsync(Job(), T0);
+ await store.EnqueueAsync(Job(queue: "other", wireName: "other-job") with { Tags = JobTags.Empty.WithTag("tenant", "acme") }, T0);
+ var claimed = Assert.Single(await ClaimAsync(store, T0, maxJobs: 1));
+ await store.ReportOutcomeAsync(claimed.JobId, "w1", claimed.Attempt, new JobOutcome.Success(), T0);
+ await store.UpsertScheduleAsync(Schedule("nightly", cursor: T0));
+ var tick = T0.AddDays(1).AddHours(3);
+ Assert.Equal(1, await store.MintDueAsync(
+ [new MintDecision("nightly", ExpectedCursor: T0, NewCursor: tick, Ticks: [tick], SkippedTicks: [])]));
+
+ (JobQuery Query, long Expected)[] cases =
+ [
+ (new JobQuery(), 5),
+ (new JobQuery { State = JobState.Succeeded }, 1),
+ (new JobQuery { State = JobState.Scheduled }, 4),
+ (new JobQuery { Queue = "other" }, 1),
+ (new JobQuery { WireName = "other-job" }, 1),
+ (new JobQuery { ScheduleId = "nightly" }, 1),
+ (new JobQuery { TagPredicates = [JobTagPredicate.HasKeyValue("tenant", "acme")] }, 2),
+ (new JobQuery { TagPredicates = [JobTagPredicate.HasKey("tenant")] }, 3),
+ (new JobQuery { TagPredicates = [JobTagPredicate.HasLabel("urgent")] }, 1),
+ (new JobQuery
+ {
+ Queue = "default",
+ TagPredicates = [JobTagPredicate.HasKeyValue("tenant", "acme"), JobTagPredicate.HasLabel("urgent")],
+ }, 1),
+ (new JobQuery { Queue = "nowhere" }, 0),
+ ];
+ foreach (var (query, expected) in cases)
+ {
+ Assert.Equal(expected, await store.CountMatchingJobsAsync(query));
+ Assert.Equal((await store.ListJobsAsync(query)).Count, await store.CountMatchingJobsAsync(query));
+ }
+ }
+
+ ///
+ /// Certifies that the filtered count ignores the query's paging fields - the after-cursor, the sort
+ /// direction, and the page size - and is never capped by the monitor page bound.
+ ///
+ [Fact]
+ public async Task Clause_5_9_CountMatchingJobs_IgnoresPaging_AndIsNotCappedByThePageBound()
+ {
+ var store = await CreateStoreAsync();
+ var total = store.Bounds.MaxMonitorPageSize + 1;
+ for (var i = 0; i < total; i++)
+ {
+ await store.EnqueueAsync(Job(), now: T0);
+ }
+ await store.EnqueueAsync(Job(queue: "other"), now: T0);
+ var first = (await store.ListJobsAsync(new JobQuery { Queue = "default", MaxResults = 1 }))[0];
+
+ Assert.Equal(total, await store.CountMatchingJobsAsync(new JobQuery { Queue = "default" }));
+ Assert.Equal(total, await store.CountMatchingJobsAsync(new JobQuery { Queue = "default", MaxResults = 1 }));
+ Assert.Equal(total, await store.CountMatchingJobsAsync(
+ new JobQuery { Queue = "default", AfterSequence = first.Sequence }));
+ Assert.Equal(total, await store.CountMatchingJobsAsync(
+ new JobQuery { Queue = "default", SortDirection = JobSortDirection.NewestFirst, AfterSequence = first.Sequence }));
+ Assert.Equal(total + 1, await store.CountMatchingJobsAsync(new JobQuery { MaxResults = 1 }));
+ }
+
+ ///
+ /// Certifies that the filtered count of an empty store is zero, with or without a filter.
+ ///
+ [Fact]
+ public async Task Clause_5_9_CountMatchingJobs_EmptyStore_IsZero()
+ {
+ var store = await CreateStoreAsync();
+
+ Assert.Equal(0, await store.CountMatchingJobsAsync(new JobQuery()));
+ Assert.Equal(0, await store.CountMatchingJobsAsync(
+ new JobQuery { State = JobState.Scheduled, TagPredicates = [JobTagPredicate.HasLabel("urgent")] }));
+ }
+
// ── Mutation teeth: boundary/misc contract facts (issue 0235) ────────────────
///
diff --git a/src/BackWave.Oracle/OracleJobStore.cs b/src/BackWave.Oracle/OracleJobStore.cs
index 995c299..c76afc4 100644
--- a/src/BackWave.Oracle/OracleJobStore.cs
+++ b/src/BackWave.Oracle/OracleJobStore.cs
@@ -2588,7 +2588,21 @@ public async ValueTask> CountJobsAsync(Cancellati
return counts;
}
- // Builds the §5.9 scope conditions shared by ListJobsAsync and FacetAsync - the scalar filters plus
+ ///
+ public async ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ {
+ await EnsureReadyAsync(cancellationToken).ConfigureAwait(false);
+
+ await using var connection = await OpenAsync(cancellationToken).ConfigureAwait(false);
+ await using var command = new OracleCommand { Connection = connection, BindByName = true };
+ var conditions = new List();
+ AppendScopeConditions(query, conditions, command);
+ var where = conditions.Count > 0 ? $"WHERE {string.Join(" AND ", conditions)}" : string.Empty;
+ command.CommandText = _schema.Rewrite($"SELECT count(*) FROM backwave.jobs {where}");
+ return Convert.ToInt64(await command.ExecuteScalarCountedAsync(cancellationToken).ConfigureAwait(false));
+ }
+
+ // Builds the §5.9 scope conditions shared by the job list, count, and facet reads - the scalar filters plus
// the AND-ed tag predicates, each an EXISTS over job_tags correlated to the job row (has-key-any-value
// omits the value condition). Everything is parameterized onto `command`. Pagination is NOT a scope
// condition - the caller adds it. Empty key/value encode to the CHR(1) sentinel so a Label (empty key)
diff --git a/src/BackWave.Postgres/PostgresJobStore.cs b/src/BackWave.Postgres/PostgresJobStore.cs
index 46ac03e..474e30a 100644
--- a/src/BackWave.Postgres/PostgresJobStore.cs
+++ b/src/BackWave.Postgres/PostgresJobStore.cs
@@ -2176,7 +2176,21 @@ public async ValueTask> CountJobsAsync(Cancellati
return counts;
}
- // Builds the spec §5.9 scope conditions shared by ListJobsAsync and FacetAsync — the scalar
+ ///
+ public async ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ {
+ await EnsureReadyAsync(cancellationToken).ConfigureAwait(false);
+
+ await using var connection = await _dataSource.OpenConnectionAsync(cancellationToken).ConfigureAwait(false);
+ await using var command = new NpgsqlCommand { Connection = connection };
+ var conditions = new List();
+ AppendScopeConditions(query, conditions, command);
+ var where = conditions.Count > 0 ? $"WHERE {string.Join(" AND ", conditions)}" : string.Empty;
+ command.CommandText = _schema.Rewrite($"SELECT count(*) FROM backwave.jobs {where}");
+ return (long)(await command.ExecuteScalarCountedAsync(cancellationToken).ConfigureAwait(false))!;
+ }
+
+ // Builds the spec §5.9 scope conditions shared by the job list, count, and facet reads — the scalar
// filters plus the AND-ed Job Tag predicates (ADR 0022), each an EXISTS over job_tags correlated
// to the job row (has-key-any-value omits the value condition). Everything is parameterized onto
// `command`. Pagination is NOT a scope condition — the caller adds it. The empty-string key
diff --git a/src/BackWave.Pro.Mcp/ToolNames.cs b/src/BackWave.Pro.Mcp/ToolNames.cs
index 7a6471d..9efdcfc 100644
--- a/src/BackWave.Pro.Mcp/ToolNames.cs
+++ b/src/BackWave.Pro.Mcp/ToolNames.cs
@@ -12,6 +12,7 @@ internal static class ToolNames
{
// JobTools
public const string SearchJobs = "search_jobs";
+ public const string CountJobs = "count_jobs";
public const string GetJob = "get_job";
public const string GetJobHistory = "get_job_history";
public const string GetJobDependencies = "get_job_dependencies";
diff --git a/src/BackWave.Pro.Mcp/Tools/BackWaveProMcpJsonContext.cs b/src/BackWave.Pro.Mcp/Tools/BackWaveProMcpJsonContext.cs
index e767a38..8fb8254 100644
--- a/src/BackWave.Pro.Mcp/Tools/BackWaveProMcpJsonContext.cs
+++ b/src/BackWave.Pro.Mcp/Tools/BackWaveProMcpJsonContext.cs
@@ -15,6 +15,7 @@ namespace BackWave.Pro.Mcp.Tools;
[JsonSerializable(typeof(AuditRecordsResult))]
[JsonSerializable(typeof(CancelJobResult))]
[JsonSerializable(typeof(CancelWorkflowToolResult))]
+[JsonSerializable(typeof(CountJobsResult))]
[JsonSerializable(typeof(GetJobDependenciesResult))]
[JsonSerializable(typeof(GetJobHistoryResult))]
[JsonSerializable(typeof(GetJobResult))]
diff --git a/src/BackWave.Pro.Mcp/Tools/JobTools.cs b/src/BackWave.Pro.Mcp/Tools/JobTools.cs
index dceea9f..22ef782 100644
--- a/src/BackWave.Pro.Mcp/Tools/JobTools.cs
+++ b/src/BackWave.Pro.Mcp/Tools/JobTools.cs
@@ -7,8 +7,8 @@
namespace BackWave.Pro.Mcp.Tools;
-// The job read tools (mcp-0003 inventory, issue 0225): search_jobs, get_job, get_job_history,
-// get_job_dependencies — all behind the view gate, all wrapping BackWaveMonitor reads 1:1.
+// The job read tools (mcp-0003 inventory, issue 0225): search_jobs, count_jobs, get_job,
+// get_job_history, get_job_dependencies — all behind the view gate, all wrapping BackWaveMonitor reads 1:1.
// Internal: the tool surface is wire-level (MCP), never a C# API. Registered explicitly via
// WithTools() in AddMcp; never assembly scanning. Input parameter names are snake_case
// (the wire contract), which is why the C# parameters carry underscores.
@@ -53,17 +53,6 @@ public async Task SearchJobsAsync(
int? max_results = null,
CancellationToken cancellationToken = default)
{
- JobState? parsedState = null;
- if (state is not null)
- {
- if (!Enum.TryParse(state, ignoreCase: true, out var s) || !Enum.IsDefined(s))
- {
- throw new McpException(
- $"Unknown state '{state}'. Valid states: {string.Join(", ", Enum.GetNames())}.");
- }
- parsedState = s;
- }
-
var direction = sort switch
{
null => JobSortDirection.NewestFirst,
@@ -92,7 +81,7 @@ _ when sort.Equals("oldest_first", StringComparison.OrdinalIgnoreCase) => JobSor
var query = new JobQuery
{
- State = parsedState,
+ State = ParseState(state),
Queue = queue,
WireName = wire_name,
ScheduleId = schedule_id,
@@ -115,6 +104,45 @@ _ when sort.Equals("oldest_first", StringComparison.OrdinalIgnoreCase) => JobSor
};
}
+ [McpServerTool(
+ Name = ToolNames.CountJobs,
+ Title = "Count jobs",
+ UseStructuredContent = true,
+ ReadOnly = true,
+ Idempotent = true,
+ OpenWorld = false)]
+ [Description(
+ "Count the jobs that match a filter, across all matches rather than one page. Takes the " +
+ "same filters as search_jobs: every filter is optional and omitted filters match " +
+ "everything; supplied filters are AND-ed together. Use it to answer \"how many\" " +
+ "questions without paging through search_jobs.")]
+ public async Task CountJobsAsync(
+ [Description("Only jobs in this state: Scheduled, AwaitingParent, Leased, Succeeded, Cancelled, DeadLettered, or Quarantined. Omit to match any state.")]
+ string? state = null,
+ [Description("Only jobs on this queue. Omit to match any queue.")]
+ string? queue = null,
+ [Description("Only jobs of this wire name (the job type's stable string identity; list_wire_names enumerates them). Omit to match any type.")]
+ string? wire_name = null,
+ [Description("Only jobs minted by this recurring schedule. Omit to match jobs from any source.")]
+ string? schedule_id = null,
+ [Description("Tag predicates, AND-ed together. Each is one of three forms: \"key=value\" (the job carries that keyed tag), \"key=*\" (the job carries any value under that key), or a bare \"value\" (the job carries that label). A label containing '=' cannot be expressed here.")]
+ string[]? tags = null,
+ CancellationToken cancellationToken = default)
+ {
+ var query = new JobQuery
+ {
+ State = ParseState(state),
+ Queue = queue,
+ WireName = wire_name,
+ ScheduleId = schedule_id,
+ TagPredicates = ParseTagPredicates(tags),
+ };
+ return new CountJobsResult
+ {
+ Count = await monitor.GetJobCountAsync(query, cancellationToken).ConfigureAwait(false),
+ };
+ }
+
[McpServerTool(
Name = ToolNames.GetJob,
Title = "Get job",
@@ -242,6 +270,20 @@ private static Guid ParseJobId(string jobId)
: throw new McpException(
$"Invalid job_id '{jobId}': expected a GUID, e.g. \"7f4df6f2-8c3a-4a0e-9d1a-2f6b8c1d5e3f\".");
+ private static JobState? ParseState(string? state)
+ {
+ if (state is null)
+ {
+ return null;
+ }
+ if (!Enum.TryParse(state, ignoreCase: true, out var parsed) || !Enum.IsDefined(parsed))
+ {
+ throw new McpException(
+ $"Unknown state '{state}'. Valid states: {string.Join(", ", Enum.GetNames())}.");
+ }
+ return parsed;
+ }
+
// The compact tag-predicate grammar (mcp-0003): "key=value" keyed, "key=*" any-value-under-key,
// bare "value" a label. Split on the FIRST '=' only, so a value containing '=' stays intact.
private static IReadOnlyList ParseTagPredicates(string[]? tags)
@@ -301,6 +343,14 @@ internal sealed record SearchJobsResult
public required bool HasMore { get; init; }
}
+/// The structured result of count_jobs.
+internal sealed record CountJobsResult
+{
+ /// The number of jobs that match the filter.
+ [Description("The number of jobs that match the filter, across all matches; zero when nothing matches.")]
+ public required long Count { get; init; }
+}
+
/// The structured result of get_job.
internal sealed record GetJobResult
{
diff --git a/src/BackWave.SqlServer/SqlServerJobStore.cs b/src/BackWave.SqlServer/SqlServerJobStore.cs
index 0fd9d4b..9ad73fa 100644
--- a/src/BackWave.SqlServer/SqlServerJobStore.cs
+++ b/src/BackWave.SqlServer/SqlServerJobStore.cs
@@ -2354,7 +2354,21 @@ public async ValueTask> CountJobsAsync(Cancellati
return counts;
}
- // Builds the §5.9 scope conditions shared by ListJobsAsync and FacetAsync — the scalar filters
+ ///
+ public async ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ {
+ await EnsureReadyAsync(cancellationToken).ConfigureAwait(false);
+
+ await using var connection = await OpenAsync(cancellationToken).ConfigureAwait(false);
+ await using var command = new SqlCommand { Connection = connection };
+ var conditions = new List();
+ AppendScopeConditions(query, conditions, command);
+ var where = conditions.Count > 0 ? $"WHERE {string.Join(" AND ", conditions)}" : string.Empty;
+ command.CommandText = _schema.Rewrite($"SELECT COUNT_BIG(*) FROM backwave.jobs {where}");
+ return (long)(await command.ExecuteScalarCountedAsync(cancellationToken).ConfigureAwait(false))!;
+ }
+
+ // Builds the §5.9 scope conditions shared by the job list, count, and facet reads — the scalar filters
// plus the AND-ed tag predicates (ADR 0022), each an EXISTS over job_tags correlated to the job
// row (has-key-any-value omits the value condition). Everything is parameterized onto `command`.
// Pagination is NOT a scope condition — the caller adds it. The empty-string key sentinel carries
diff --git a/src/BackWave.Sqlite/SqliteJobStore.cs b/src/BackWave.Sqlite/SqliteJobStore.cs
index c8a207e..6f32ba4 100644
--- a/src/BackWave.Sqlite/SqliteJobStore.cs
+++ b/src/BackWave.Sqlite/SqliteJobStore.cs
@@ -2042,7 +2042,21 @@ public async ValueTask> CountJobsAsync(Cancellati
return counts;
}
- // Builds the §5.9 scope conditions shared by ListJobsAsync and FacetAsync — the scalar filters
+ ///
+ public async ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ {
+ await EnsureReadyAsync(cancellationToken).ConfigureAwait(false);
+
+ await using var connection = await OpenAsync(cancellationToken).ConfigureAwait(false);
+ await using var command = connection.CreateCommand();
+ var conditions = new List();
+ AppendScopeConditions(query, conditions, command);
+ var where = conditions.Count > 0 ? $"WHERE {string.Join(" AND ", conditions)}" : string.Empty;
+ command.CommandText = _schema.Rewrite($"SELECT count(*) FROM backwave_jobs {where}");
+ return (long)(await command.ExecuteScalarCountedAsync(cancellationToken).ConfigureAwait(false))!;
+ }
+
+ // Builds the §5.9 scope conditions shared by the job list, count, and facet reads — the scalar filters
// plus the AND-ed tag predicates (ADR 0022), each an EXISTS over backwave_job_tags correlated to
// the job row. Pagination is NOT a scope condition — the caller adds it. The empty-string key
// sentinel carries Labels.
diff --git a/src/BackWave/Monitor/BackWaveMonitor.cs b/src/BackWave/Monitor/BackWaveMonitor.cs
index 98fa868..a6c079d 100644
--- a/src/BackWave/Monitor/BackWaveMonitor.cs
+++ b/src/BackWave/Monitor/BackWaveMonitor.cs
@@ -127,6 +127,25 @@ public async ValueTask> ListJobsAsync(
return [.. records.Select(ToSnapshot)];
}
+ ///
+ /// The number of jobs matching a filter, across the whole matching population. Use it to show a
+ /// total next to a filtered page (for example "1,204 quarantined jobs
+ /// on the lab queue"). The count uses the same filters as :
+ /// state, queue, wire name, schedule id, and tag predicates, each narrowing the count.
+ ///
+ /// The query's paging fields do not apply: ,
+ /// , and are ignored, and
+ /// does not cap the count.
+ ///
+ ///
+ ///
+ /// The filter. When null, an empty query is used, which counts all jobs.
+ ///
+ /// Cancels the read.
+ /// The number of matching jobs; zero when nothing matches.
+ public ValueTask GetJobCountAsync(JobQuery? query = null, CancellationToken cancellationToken = default)
+ => store.CountMatchingJobsAsync(query ?? new JobQuery(), cancellationToken);
+
///
/// One job's payload bytes, rendered best-effort for display. The payload is opaque to BackWave —
/// serialized by your own serializer and never parsed here — so it is decoded as UTF-8 with a hex
diff --git a/src/BackWave/Storage/IJobStore.cs b/src/BackWave/Storage/IJobStore.cs
index 0edebfa..54227ff 100644
--- a/src/BackWave/Storage/IJobStore.cs
+++ b/src/BackWave/Storage/IJobStore.cs
@@ -464,6 +464,39 @@ ValueTask SetConcurrencyLimitAsync(
/// The matching jobs in the requested sort order, capped at the store's maximum page size.
ValueTask> ListJobsAsync(JobQuery query, CancellationToken cancellationToken = default);
+ ///
+ /// Counts the jobs matching a filter, for monitoring. The count uses exactly the filter predicates
+ /// applies (state, queue, wire name, schedule id, and AND-ed tag
+ /// predicates). The query's pagination and sort fields (,
+ /// , ) do not apply: the count
+ /// always spans the whole matching population, never a single page, and the monitor page cap does
+ /// not bound it.
+ ///
+ /// The default implementation pages through oldest-first and counts the
+ /// rows, so an adapter that does not override it still returns the correct number, at the cost of
+ /// one read per page. An overriding adapter counts in the store (for example a single
+ /// SELECT COUNT(*) over the same predicates).
+ ///
+ ///
+ /// The filter to count. Its pagination and sort fields are ignored.
+ /// Cancels the operation.
+ /// The number of jobs that match the filter; zero when nothing matches.
+ async ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ {
+ var scope = query with { AfterSequence = null, SortDirection = JobSortDirection.OldestFirst, MaxResults = int.MaxValue };
+ long count = 0;
+ while (true)
+ {
+ var page = await ListJobsAsync(scope, cancellationToken).ConfigureAwait(false);
+ if (page.Count == 0)
+ {
+ return count;
+ }
+ count += page.Count;
+ scope = scope with { AfterSequence = page[^1].Sequence };
+ }
+ }
+
///
/// Reads job counts grouped by Queue and state — the queue depths — for monitoring.
///
diff --git a/src/BackWave/Storage/InMemory/InMemoryJobStore.cs b/src/BackWave/Storage/InMemory/InMemoryJobStore.cs
index d70ef88..05c898b 100644
--- a/src/BackWave/Storage/InMemory/InMemoryJobStore.cs
+++ b/src/BackWave/Storage/InMemory/InMemoryJobStore.cs
@@ -1712,9 +1712,18 @@ public ValueTask> ListJobsAsync(
}
}
+ ///
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ {
+ lock (_gate)
+ {
+ return ValueTask.FromResult((long)_jobs.Values.Count(j => MatchesScope(j, query)));
+ }
+ }
+
///
- /// The scope predicate shared by and
- /// : the scalar filters AND-ed with the tag predicates. Pagination
+ /// The scope predicate shared by , ,
+ /// and : the scalar filters AND-ed with the tag predicates. Pagination
/// (cursor/sort/take) is NOT part of the scope — facets count the whole matching population.
///
private static bool MatchesScope(JobRecord j, JobQuery query)
diff --git a/tests/BackWave.Hosting.Tests/AdaptivePollPacerTests.cs b/tests/BackWave.Hosting.Tests/AdaptivePollPacerTests.cs
index fe7907a..68ed898 100644
--- a/tests/BackWave.Hosting.Tests/AdaptivePollPacerTests.cs
+++ b/tests/BackWave.Hosting.Tests/AdaptivePollPacerTests.cs
@@ -277,6 +277,9 @@ public ValueTask> ListJobsAsync(JobQuery query, Cancell
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
=> inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);
diff --git a/tests/BackWave.Hosting.Tests/FaultableStore.cs b/tests/BackWave.Hosting.Tests/FaultableStore.cs
index 273a8ee..ffc8714 100644
--- a/tests/BackWave.Hosting.Tests/FaultableStore.cs
+++ b/tests/BackWave.Hosting.Tests/FaultableStore.cs
@@ -414,6 +414,9 @@ public ValueTask> ListJobsAsync(JobQuery query, Cancell
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
=> inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);
diff --git a/tests/BackWave.Hosting.Tests/PollCoalescingTests.cs b/tests/BackWave.Hosting.Tests/PollCoalescingTests.cs
index 71834d0..4e37ccc 100644
--- a/tests/BackWave.Hosting.Tests/PollCoalescingTests.cs
+++ b/tests/BackWave.Hosting.Tests/PollCoalescingTests.cs
@@ -190,6 +190,9 @@ public ValueTask> ListJobsAsync(JobQuery query, Cancell
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
=> inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);
diff --git a/tests/BackWave.Hosting.Tests/PreRelinquishStore.cs b/tests/BackWave.Hosting.Tests/PreRelinquishStore.cs
index dc3360a..a392bef 100644
--- a/tests/BackWave.Hosting.Tests/PreRelinquishStore.cs
+++ b/tests/BackWave.Hosting.Tests/PreRelinquishStore.cs
@@ -103,6 +103,9 @@ public ValueTask> ListJobsAsync(
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue,
CancellationToken cancellationToken = default)
diff --git a/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs b/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs
index 1b34afa..0e80f74 100644
--- a/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs
+++ b/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs
@@ -12,16 +12,17 @@ namespace BackWave.Pro.Mcp.Tests;
///
/// The job read tools end-to-end through the mounted MCP endpoint (issue 0225): search_jobs with
-/// its tag-predicate grammar and cursor paging, get_job's found-not-error contract, the
-/// self-explaining get_job_history response, and get_job_dependencies edges.
+/// its tag-predicate grammar and cursor paging, count_jobs over the same filters, get_job's
+/// found-not-error contract, the self-explaining get_job_history response, and get_job_dependencies
+/// edges.
///
public sealed class JobToolsTests
{
private static readonly string[] JobToolNames =
- ["search_jobs", "get_job", "get_job_history", "get_job_dependencies"];
+ ["search_jobs", "count_jobs", "get_job", "get_job_history", "get_job_dependencies"];
[Fact]
- public async Task ToolsList_ShowsAllFourJobTools_WithOutputSchemas()
+ public async Task ToolsList_ShowsAllFiveJobTools_WithOutputSchemas()
{
await using var server = await McpTestServer.StartAsync();
@@ -40,6 +41,7 @@ JsonElement OutputProperties(string name) =>
Assert.True(OutputProperties("search_jobs").TryGetProperty("jobs", out _));
Assert.True(OutputProperties("search_jobs").TryGetProperty("nextCursor", out _));
Assert.True(OutputProperties("search_jobs").TryGetProperty("hasMore", out _));
+ Assert.True(OutputProperties("count_jobs").TryGetProperty("count", out _));
Assert.True(OutputProperties("get_job").TryGetProperty("found", out _));
Assert.True(OutputProperties("get_job_history").TryGetProperty("transitions", out _));
Assert.True(OutputProperties("get_job_history").TryGetProperty("historyPolicy", out _));
@@ -53,6 +55,12 @@ JsonElement OutputProperties(string name) =>
{
Assert.True(searchInputs.TryGetProperty(parameter, out _), $"search_jobs is missing input '{parameter}'");
}
+
+ // count_jobs takes the search filters but none of the paging inputs.
+ var countInputs = tools.Single(t => t.Name == "count_jobs").InputSchema!.Value.GetProperty("properties");
+ Assert.Equal(
+ ["queue", "schedule_id", "state", "tags", "wire_name"],
+ countInputs.EnumerateObject().Select(p => p.Name).Order());
}
[Fact]
@@ -406,6 +414,52 @@ public async Task SearchJobs_OldestFirst_ReversesTheOrder()
Assert.Equal(second, jobs[1].GetProperty("jobId").GetGuid());
}
+ [Fact]
+ public async Task CountJobs_CountsEveryMatch_PastTheStoreCap_WithTheSearchFilters()
+ {
+ // The count spans every match, not one page: seed past the store's page cap and the total
+ // still comes back whole. Each filter narrows the count exactly as it narrows search_jobs.
+ await using var server = await McpTestServer.StartAsync(
+ bounds: StoreBounds.Default with { MaxMonitorPageSize = 10 });
+ for (var i = 0; i < 12; i++)
+ {
+ await server.SeedJobAsync("critical", "send-email");
+ }
+ await server.SeedJobAsync("bulk", "resize-image");
+ await SeedTaggedJobAsync(server.Store, JobTags.Empty.WithLabel("urgent").WithTag("tenant", "acme"));
+ await SeedTaggedJobAsync(server.Store, JobTags.Empty.WithTag("tenant", "globex"));
+
+ Assert.Equal(15, await CountAsync(server, []));
+ Assert.Equal(14, await CountAsync(server, new() { ["queue"] = "critical" }));
+ Assert.Equal(12, await CountAsync(server, new() { ["wire_name"] = "send-email" }));
+ Assert.Equal(15, await CountAsync(server, new() { ["state"] = "scheduled" }));
+ Assert.Equal(0, await CountAsync(server, new() { ["state"] = "Succeeded" }));
+ Assert.Equal(0, await CountAsync(server, new() { ["schedule_id"] = "nightly" }));
+ Assert.Equal(2, await CountAsync(server, new() { ["tags"] = new[] { "tenant=*" } }));
+ Assert.Equal(1, await CountAsync(server, new() { ["tags"] = new[] { "tenant=acme", "urgent" } }));
+ Assert.Equal(0, await CountAsync(server, new() { ["queue"] = "bulk", ["tags"] = new[] { "urgent" } }));
+ }
+
+ [Fact]
+ public async Task CountJobs_InvalidStateAndTag_AreInvalidInputErrors()
+ {
+ await using var server = await McpTestServer.StartAsync();
+
+ var badState = await server.Client.CallToolAsync("count_jobs", new Dictionary
+ {
+ ["state"] = "Exploded",
+ });
+ Assert.True(badState.IsError);
+ Assert.Contains("Scheduled", badState.Text);
+
+ var badTag = await server.Client.CallToolAsync("count_jobs", new Dictionary
+ {
+ ["tags"] = new[] { "=acme" },
+ });
+ Assert.True(badTag.IsError);
+ Assert.Contains("key=value", badTag.Text);
+ }
+
[Fact]
public async Task GetJob_ReturnsTheSnapshot()
{
@@ -560,6 +614,9 @@ public async Task DeniedViewGate_CoversTheJobTools()
var result = await server.Client.CallToolAsync("search_jobs");
Assert.True(result.IsError);
Assert.Contains("Permission denied", result.Text);
+ var count = await server.Client.CallToolAsync("count_jobs");
+ Assert.True(count.IsError);
+ Assert.Contains("Permission denied", count.Text);
}
private static async Task SeedTaggedJobAsync(InMemoryJobStore store, JobTags tags)
@@ -572,6 +629,13 @@ private static async Task SeedTaggedJobAsync(InMemoryJobStore store, JobTa
return id;
}
+ private static async Task CountAsync(McpTestServer server, Dictionary args)
+ {
+ var result = await server.Client.CallToolAsync("count_jobs", args);
+ Assert.False(result.IsError);
+ return result.StructuredContent!.Value.GetProperty("count").GetInt64();
+ }
+
private static async Task SearchByTagsAsync(McpTestServer server, string[] tags)
{
var result = await server.Client.CallToolAsync("search_jobs", new Dictionary
diff --git a/tests/BackWave.Pro.Mcp.Tests/ReadToolsTests.cs b/tests/BackWave.Pro.Mcp.Tests/ReadToolsTests.cs
index d25e72a..25f55b6 100644
--- a/tests/BackWave.Pro.Mcp.Tests/ReadToolsTests.cs
+++ b/tests/BackWave.Pro.Mcp.Tests/ReadToolsTests.cs
@@ -23,11 +23,12 @@ public async Task UnconfiguredHost_ListsTheWholeNonSensitiveReadSurface_AndNothi
var tools = await server.Client.ListToolsAsync();
- // The complete 14-tool non-sensitive read surface (the 12 plain reads plus the two
+ // The complete 15-tool non-sensitive read surface (the 13 plain reads plus the two
// workflow reads of 0228); this assertion is exact so it fails loudly if a sensitive or
// write tool ever leaks into the unconfigured default.
string[] expected =
[
+ "count_jobs",
"get_job",
"get_job_dependencies",
"get_job_history",
diff --git a/tests/BackWave.Pro.Mcp.Tests/ToolListFilteringTests.cs b/tests/BackWave.Pro.Mcp.Tests/ToolListFilteringTests.cs
index 6b00308..7bd173b 100644
--- a/tests/BackWave.Pro.Mcp.Tests/ToolListFilteringTests.cs
+++ b/tests/BackWave.Pro.Mcp.Tests/ToolListFilteringTests.cs
@@ -20,11 +20,11 @@ public sealed class ToolListFilteringTests
// literals, deliberately NOT the ToolNames constants, so a rename of a tool's wire value breaks
// this test instead of silently sliding through. It catches the rename-ships-ungated trap two
// ways: a write tool whose attribute Name drifts from its gate key would surface in the default
- // (ungranted) list and fail DefaultServer_PresentsExactlyTheFourteenReadTools; a read tool
+ // (ungranted) list and fail DefaultServer_PresentsExactlyTheFifteenReadTools; a read tool
// dropped or renamed would fail the same set-equality. The all-granted set pins the full surface.
private static readonly string[] DefaultReadTools =
[
- "search_jobs", "get_job", "get_job_history", "get_job_dependencies",
+ "search_jobs", "count_jobs", "get_job", "get_job_history", "get_job_dependencies",
"get_observer_lag", "list_observer_dead_letters",
"list_workflows", "get_workflow",
"get_queue_settings", "get_tag_facet", "list_wire_names", "list_schedules", "list_audit_records",
@@ -33,11 +33,11 @@ public sealed class ToolListFilteringTests
private static readonly string[] SensitiveDataTools = ["get_job_payload", "get_job_output"];
- private static readonly string[] AllTwentyThreeTools =
+ private static readonly string[] AllTwentyFourTools =
[.. DefaultReadTools, .. AllWriteTools, "cancel_workflow", .. SensitiveDataTools];
[Fact]
- public async Task DefaultServer_PresentsExactlyTheFourteenReadTools()
+ public async Task DefaultServer_PresentsExactlyTheFifteenReadTools()
{
// No write grants and sensitive data not authorized: tools/list is exactly the read surface.
// Set-equality (not Contains) is the point — an accidentally-ungated write tool would appear,
@@ -50,7 +50,7 @@ public async Task DefaultServer_PresentsExactlyTheFourteenReadTools()
}
[Fact]
- public async Task EveryGateGranted_PresentsExactlyTheFullTwentyThreeTools()
+ public async Task EveryGateGranted_PresentsExactlyTheFullTwentyFourTools()
{
// Every write callback granted AND sensitive data exposed: tools/list is the entire surface,
// no more and no less. This pins the full membership so an added or removed tool is caught.
@@ -69,7 +69,7 @@ public async Task EveryGateGranted_PresentsExactlyTheFullTwentyThreeTools()
var tools = (await server.Client.ListToolsAsync()).Select(t => t.Name).ToHashSet();
- Assert.Equal(AllTwentyThreeTools.ToHashSet(), tools);
+ Assert.Equal(AllTwentyFourTools.ToHashSet(), tools);
}
[Fact]
diff --git a/tests/BackWave.Simulation/Simulator.cs b/tests/BackWave.Simulation/Simulator.cs
index e81c7eb..04bded0 100644
--- a/tests/BackWave.Simulation/Simulator.cs
+++ b/tests/BackWave.Simulation/Simulator.cs
@@ -3422,6 +3422,9 @@ public ValueTask> ListJobsAsync(JobQuery query, Cancell
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
=> inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);
@@ -3611,6 +3614,9 @@ public ValueTask> ListJobsAsync(JobQuery query, Cancell
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
=> inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);
diff --git a/tests/BackWave.Tests/DispatchPolicyTests.cs b/tests/BackWave.Tests/DispatchPolicyTests.cs
index a91f70c..8cba48c 100644
--- a/tests/BackWave.Tests/DispatchPolicyTests.cs
+++ b/tests/BackWave.Tests/DispatchPolicyTests.cs
@@ -480,6 +480,9 @@ public ValueTask> ListJobsAsync(JobQuery query, Cancell
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
=> inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);
diff --git a/tests/BackWave.Tests/MonitorApiTests.cs b/tests/BackWave.Tests/MonitorApiTests.cs
index 20efad2..5d0de7d 100644
--- a/tests/BackWave.Tests/MonitorApiTests.cs
+++ b/tests/BackWave.Tests/MonitorApiTests.cs
@@ -156,4 +156,180 @@ public async Task ListJobs_PageIsBoundedByMaxMonitorPageSize()
Assert.Equal(2, (await fixture.Monitor.ListJobsAsync()).Count);
Assert.Single(await fixture.Monitor.ListJobsAsync(new JobQuery { MaxResults = 1 }));
}
+
+
+ [Fact]
+ public async Task GetJobCount_CountsTheFilteredPopulation_PastThePageBound()
+ {
+ var fixture = CreateFixture(new StoreBounds { MaxMonitorPageSize = 2 });
+
+ await fixture.Client.EnqueueAsync(new InventorySync("a"), dueTime: T0);
+ await fixture.Client.EnqueueAsync(new InventorySync("b"), dueTime: T0);
+ await fixture.Client.EnqueueAsync(new InventorySync("c"), dueTime: T0);
+ await fixture.Client.EnqueueAsync(new InventorySync("d"), dueTime: T0.AddDays(1), queue: "reports");
+
+ await fixture.Pump.PumpAsync(T0); // "a", "b", "c": Succeeded; "d": still Scheduled
+
+ Assert.Equal(4, await fixture.Monitor.GetJobCountAsync());
+ Assert.Equal(3, await fixture.Monitor.GetJobCountAsync(new JobQuery { State = JobState.Succeeded }));
+ Assert.Equal(1, await fixture.Monitor.GetJobCountAsync(new JobQuery { Queue = "reports" }));
+ Assert.Equal(3, await fixture.Monitor.GetJobCountAsync(new JobQuery { Queue = "default", MaxResults = 1 }));
+ Assert.Equal(0, await fixture.Monitor.GetJobCountAsync(new JobQuery { WireName = "no-such-wire-name" }));
+ }
+
+ [Fact]
+ public async Task GetJobCount_OnAStoreWithoutACountOverride_PagesThroughTheListing()
+ {
+ var inner = new InMemoryJobStore(new StoreBounds { MaxMonitorPageSize = 2 });
+ var store = new ListingOnlyStore(inner);
+ for (var i = 0; i < 5; i++)
+ {
+ await inner.EnqueueAsync(new NewJob(Guid.NewGuid(), "inventory-sync", "{}"u8.ToArray(), "default", T0), T0);
+ }
+ await inner.EnqueueAsync(new NewJob(Guid.NewGuid(), "inventory-sync", "{}"u8.ToArray(), "reports", T0), T0);
+ var monitor = new BackWaveMonitor(store);
+
+ // Paging fields on the query are ignored: the default resets the cursor and walks every page.
+ Assert.Equal(5, await monitor.GetJobCountAsync(
+ new JobQuery { Queue = "default", AfterSequence = 3, SortDirection = JobSortDirection.NewestFirst, MaxResults = 1 }));
+ Assert.Equal(4, store.ListCalls); // pages of 2, 2, 1, then the empty page that ends the walk
+
+ Assert.Equal(6, await monitor.GetJobCountAsync());
+ Assert.Equal(0, await monitor.GetJobCountAsync(new JobQuery { Queue = "nowhere" }));
+ }
+
+ ///
+ /// Forwards every member to the In-Memory Store except the filtered count, so the count falls back
+ /// to the store contract's default, which pages through the listing. Counts the listing reads.
+ ///
+ private sealed class ListingOnlyStore(IJobStore inner) : IJobStore
+ {
+ public int ListCalls { get; private set; }
+
+ public ValueTask> ClaimAsync(ClaimRequest request, CancellationToken cancellationToken = default)
+ => inner.ClaimAsync(request, cancellationToken);
+
+ public bool SupportsTransactionalEnqueue => inner.SupportsTransactionalEnqueue;
+
+ public ValueTask EnqueueAsync(
+ NewJob job, DateTimeOffset now, System.Data.Common.DbTransaction? transaction = null, CancellationToken cancellationToken = default)
+ => inner.EnqueueAsync(job, now, transaction, cancellationToken);
+
+ public ValueTask ReportOutcomeAsync(
+ Guid jobId, string workerId, int attempt, JobOutcome outcome, DateTimeOffset now,
+ string? failureDetail = null, JobTags? addedTags = null, ReadOnlyMemory? output = null,
+ CancellationToken cancellationToken = default)
+ => inner.ReportOutcomeAsync(jobId, workerId, attempt, outcome, now, failureDetail, addedTags, output, cancellationToken);
+
+ public ValueTask?> GetJobOutputAsync(Guid jobId, CancellationToken cancellationToken = default)
+ => inner.GetJobOutputAsync(jobId, cancellationToken);
+
+ public ValueTask> HeartbeatAsync(
+ string workerId, IReadOnlyList jobIds, TimeSpan leaseDuration, DateTimeOffset now, CancellationToken cancellationToken = default)
+ => inner.HeartbeatAsync(workerId, jobIds, leaseDuration, now, cancellationToken);
+
+ public ValueTask ExpireLeasesAsync(
+ DateTimeOffset now, int maxJobs, IReadOnlyList queues, RetryDisposition disposition, CancellationToken cancellationToken = default)
+ => inner.ExpireLeasesAsync(now, maxJobs, queues, disposition, cancellationToken);
+
+ public ValueTask CancelJobAsync(Guid jobId, string actor, DateTimeOffset now, CancellationToken cancellationToken = default)
+ => inner.CancelJobAsync(jobId, actor, now, cancellationToken);
+
+ public ValueTask RequeueAsync(Guid jobId, string actor, DateTimeOffset now, CancellationToken cancellationToken = default)
+ => inner.RequeueAsync(jobId, actor, now, cancellationToken);
+
+ public ValueTask PauseQueueAsync(string queue, string actor, DateTimeOffset now, CancellationToken cancellationToken = default)
+ => inner.PauseQueueAsync(queue, actor, now, cancellationToken);
+
+ public ValueTask ResumeQueueAsync(string queue, string actor, DateTimeOffset now, CancellationToken cancellationToken = default)
+ => inner.ResumeQueueAsync(queue, actor, now, cancellationToken);
+
+ public ValueTask TriggerScheduleNowAsync(string scheduleId, string actor, DateTimeOffset now, CancellationToken cancellationToken = default)
+ => inner.TriggerScheduleNowAsync(scheduleId, actor, now, cancellationToken);
+
+ public ValueTask> ListAuditRecordsAsync(string target, CancellationToken cancellationToken = default)
+ => inner.ListAuditRecordsAsync(target, cancellationToken);
+
+ public ValueTask UpsertScheduleAsync(ScheduleRecord schedule, CancellationToken cancellationToken = default)
+ => inner.UpsertScheduleAsync(schedule, cancellationToken);
+
+ public ValueTask RemoveScheduleAsync(string scheduleId, CancellationToken cancellationToken = default)
+ => inner.RemoveScheduleAsync(scheduleId, cancellationToken);
+
+ public ValueTask> ListSchedulesAsync(CancellationToken cancellationToken = default)
+ => inner.ListSchedulesAsync(cancellationToken);
+
+ public ValueTask MintDueAsync(IReadOnlyList decisions, CancellationToken cancellationToken = default)
+ => inner.MintDueAsync(decisions, cancellationToken);
+
+ public ValueTask SetConcurrencyLimitAsync(string queue, int? limit, string actor, DateTimeOffset now, CancellationToken cancellationToken = default)
+ => inner.SetConcurrencyLimitAsync(queue, limit, actor, now, cancellationToken);
+
+ public ValueTask GetJobAsync(Guid jobId, CancellationToken cancellationToken = default)
+ => inner.GetJobAsync(jobId, cancellationToken);
+
+ public ValueTask> GetJobHistoryAsync(Guid jobId, CancellationToken cancellationToken = default)
+ => inner.GetJobHistoryAsync(jobId, cancellationToken);
+
+ public ValueTask> ListJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ {
+ ListCalls++;
+ return inner.ListJobsAsync(query, cancellationToken);
+ }
+
+ public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
+ => inner.CountJobsAsync(cancellationToken);
+
+ public ValueTask> FacetAsync(
+ string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
+ => inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);
+
+ public ValueTask> SuggestTagsAsync(TagSuggestQuery query, CancellationToken cancellationToken = default)
+ => inner.SuggestTagsAsync(query, cancellationToken);
+
+ public ValueTask> ListQueueSettingsAsync(CancellationToken cancellationToken = default)
+ => inner.ListQueueSettingsAsync(cancellationToken);
+
+ public ValueTask GetDependencyEdgesAsync(Guid jobId, CancellationToken cancellationToken = default)
+ => inner.GetDependencyEdgesAsync(jobId, cancellationToken);
+
+ public ValueTask EnqueueWorkflowAsync(
+ WorkflowDefinition workflow, DateTimeOffset now, System.Data.Common.DbTransaction? transaction = null,
+ CancellationToken cancellationToken = default)
+ => inner.EnqueueWorkflowAsync(workflow, now, transaction, cancellationToken);
+
+ public ValueTask> ListWorkflowsAsync(CancellationToken cancellationToken = default)
+ => inner.ListWorkflowsAsync(cancellationToken);
+
+ public ValueTask GetWorkflowAsync(Guid workflowId, CancellationToken cancellationToken = default)
+ => inner.GetWorkflowAsync(workflowId, cancellationToken);
+
+ public ValueTask PurgeTerminalAsync(
+ TerminalStateClass stateClass, DateTimeOffset terminalBefore, int maxJobs, CancellationToken cancellationToken = default)
+ => inner.PurgeTerminalAsync(stateClass, terminalBefore, maxJobs, cancellationToken);
+
+ public ValueTask ClaimObserverDeliveriesAsync(
+ ObserverClaimRequest request, CancellationToken cancellationToken = default)
+ => inner.ClaimObserverDeliveriesAsync(request, cancellationToken);
+
+ public ValueTask ReportObserverDeliveriesAsync(
+ ObserverDeliveryReport report, CancellationToken cancellationToken = default)
+ => inner.ReportObserverDeliveriesAsync(report, cancellationToken);
+
+ // Forwarded as well as the void twin: without this override the interface default answers
+ // Unreported and flattens the inner store's fence verdict before the caller ever sees it.
+ public ValueTask TryReportObserverDeliveriesAsync(
+ ObserverDeliveryReport report, CancellationToken cancellationToken = default)
+ => inner.TryReportObserverDeliveriesAsync(report, cancellationToken);
+
+ public ValueTask GetObserverCursorAsync(string observerId, CancellationToken cancellationToken = default)
+ => inner.GetObserverCursorAsync(observerId, cancellationToken);
+
+ public ValueTask GetObserverLagAsync(ObserverLagRequest request, CancellationToken cancellationToken = default)
+ => inner.GetObserverLagAsync(request, cancellationToken);
+
+ public ValueTask> ListObserverDeadLettersAsync(
+ string observerId, CancellationToken cancellationToken = default)
+ => inner.ListObserverDeadLettersAsync(observerId, cancellationToken);
+ }
}
diff --git a/tests/BackWave.Tests/WorkflowsV2RetryTests.cs b/tests/BackWave.Tests/WorkflowsV2RetryTests.cs
index f1b0c64..2570a98 100644
--- a/tests/BackWave.Tests/WorkflowsV2RetryTests.cs
+++ b/tests/BackWave.Tests/WorkflowsV2RetryTests.cs
@@ -383,6 +383,9 @@ public ValueTask> ListJobsAsync(JobQuery query, Cancell
public ValueTask> CountJobsAsync(CancellationToken cancellationToken = default)
=> inner.CountJobsAsync(cancellationToken);
+ public ValueTask CountMatchingJobsAsync(JobQuery query, CancellationToken cancellationToken = default)
+ => inner.CountMatchingJobsAsync(query, cancellationToken);
+
public ValueTask> FacetAsync(
string key, JobQuery? baseQuery = null, int maxResults = int.MaxValue, CancellationToken cancellationToken = default)
=> inner.FacetAsync(key, baseQuery, maxResults, cancellationToken);