From 1958fc2c56c36eb58963e787a6f40289f1833ba0 Mon Sep 17 00:00:00 2001 From: Paul DeVito Date: Tue, 6 Oct 2026 00:57:00 +0300 Subject: [PATCH] feat(monitor): add a filtered job count Add BackWaveMonitor.GetJobCountAsync. It returns the number of jobs that match the scope filters of a JobQuery: state, queue, wire name, schedule id, and tag predicates. It ignores the paging fields (AfterSequence, SortDirection, MaxResults), and the monitor page size does not cap it. Add IJobStore.CountMatchingJobsAsync as a default interface method. The default pages through ListJobsAsync, so a custom store keeps compiling and gives the correct count. InMemory, SQLite, Postgres, SQL Server, and Oracle override it with one COUNT(*) that uses the same scope conditions as the listing. Every forwarding store wrapper in the tests forwards the new member. Add three conformance facts: the count equals the listing for each filter kind, the count ignores the paging fields, and an empty store gives 0. Add the read-only MCP tool count_jobs. It takes the same filters as search_jobs and uses the view gate. Add GET /monitor/jobs/count to the sample. --- samples/BackWave.Sample.Api/Program.cs | 12 ++ samples/BackWave.Sample.Api/README.md | 2 + src/BackWave.Conformance/ConformanceSuite.cs | 83 +++++++++ src/BackWave.Oracle/OracleJobStore.cs | 16 +- src/BackWave.Postgres/PostgresJobStore.cs | 16 +- src/BackWave.Pro.Mcp/ToolNames.cs | 1 + .../Tools/BackWaveProMcpJsonContext.cs | 1 + src/BackWave.Pro.Mcp/Tools/JobTools.cs | 78 ++++++-- src/BackWave.SqlServer/SqlServerJobStore.cs | 16 +- src/BackWave.Sqlite/SqliteJobStore.cs | 16 +- src/BackWave/Monitor/BackWaveMonitor.cs | 19 ++ src/BackWave/Storage/IJobStore.cs | 33 ++++ .../Storage/InMemory/InMemoryJobStore.cs | 13 +- .../AdaptivePollPacerTests.cs | 3 + .../BackWave.Hosting.Tests/FaultableStore.cs | 3 + .../PollCoalescingTests.cs | 3 + .../PreRelinquishStore.cs | 3 + tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs | 72 ++++++- .../BackWave.Pro.Mcp.Tests/ReadToolsTests.cs | 3 +- .../ToolListFilteringTests.cs | 12 +- tests/BackWave.Simulation/Simulator.cs | 6 + tests/BackWave.Tests/DispatchPolicyTests.cs | 3 + tests/BackWave.Tests/MonitorApiTests.cs | 176 ++++++++++++++++++ tests/BackWave.Tests/WorkflowsV2RetryTests.cs | 3 + 24 files changed, 562 insertions(+), 31 deletions(-) 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);