Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions samples/BackWave.Sample.Api/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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.");

Expand Down
2 changes: 2 additions & 0 deletions samples/BackWave.Sample.Api/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
83 changes: 83 additions & 0 deletions src/BackWave.Conformance/ConformanceSuite.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1795,6 +1795,89 @@ public async Task Clause_5_9_Reads_FilterAndCount_OverCommittedEffectsOnly()
Assert.Contains(new QueueStateCount("other", JobState.Scheduled, 1), counts);
}

/// <summary>
/// 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.
/// </summary>
[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));
}
}

/// <summary>
/// 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.
/// </summary>
[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 }));
}

/// <summary>
/// Certifies that the filtered count of an empty store is zero, with or without a filter.
/// </summary>
[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) ────────────────

/// <summary>
Expand Down
16 changes: 15 additions & 1 deletion src/BackWave.Oracle/OracleJobStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2588,7 +2588,21 @@ public async ValueTask<IReadOnlyList<QueueStateCount>> CountJobsAsync(Cancellati
return counts;
}

// Builds the §5.9 scope conditions shared by ListJobsAsync and FacetAsync - the scalar filters plus
/// <inheritdoc/>
public async ValueTask<long> 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<string>();
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)
Expand Down
16 changes: 15 additions & 1 deletion src/BackWave.Postgres/PostgresJobStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2176,7 +2176,21 @@ public async ValueTask<IReadOnlyList<QueueStateCount>> CountJobsAsync(Cancellati
return counts;
}

// Builds the spec §5.9 scope conditions shared by ListJobsAsync and FacetAsync — the scalar
/// <inheritdoc/>
public async ValueTask<long> 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<string>();
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
Expand Down
1 change: 1 addition & 0 deletions src/BackWave.Pro.Mcp/ToolNames.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
1 change: 1 addition & 0 deletions src/BackWave.Pro.Mcp/Tools/BackWaveProMcpJsonContext.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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))]
Expand Down
78 changes: 64 additions & 14 deletions src/BackWave.Pro.Mcp/Tools/JobTools.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<JobTools>() in AddMcp; never assembly scanning. Input parameter names are snake_case
// (the wire contract), which is why the C# parameters carry underscores.
Expand Down Expand Up @@ -53,17 +53,6 @@ public async Task<SearchJobsResult> SearchJobsAsync(
int? max_results = null,
CancellationToken cancellationToken = default)
{
JobState? parsedState = null;
if (state is not null)
{
if (!Enum.TryParse<JobState>(state, ignoreCase: true, out var s) || !Enum.IsDefined(s))
{
throw new McpException(
$"Unknown state '{state}'. Valid states: {string.Join(", ", Enum.GetNames<JobState>())}.");
}
parsedState = s;
}

var direction = sort switch
{
null => JobSortDirection.NewestFirst,
Expand Down Expand Up @@ -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,
Expand All @@ -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<CountJobsResult> 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",
Expand Down Expand Up @@ -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<JobState>(state, ignoreCase: true, out var parsed) || !Enum.IsDefined(parsed))
{
throw new McpException(
$"Unknown state '{state}'. Valid states: {string.Join(", ", Enum.GetNames<JobState>())}.");
}
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<JobTagPredicate> ParseTagPredicates(string[]? tags)
Expand Down Expand Up @@ -301,6 +343,14 @@ internal sealed record SearchJobsResult
public required bool HasMore { get; init; }
}

/// <summary>The structured result of <c>count_jobs</c>.</summary>
internal sealed record CountJobsResult
{
/// <summary>The number of jobs that match the filter.</summary>
[Description("The number of jobs that match the filter, across all matches; zero when nothing matches.")]
public required long Count { get; init; }
}

/// <summary>The structured result of <c>get_job</c>.</summary>
internal sealed record GetJobResult
{
Expand Down
16 changes: 15 additions & 1 deletion src/BackWave.SqlServer/SqlServerJobStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2354,7 +2354,21 @@ public async ValueTask<IReadOnlyList<QueueStateCount>> CountJobsAsync(Cancellati
return counts;
}

// Builds the §5.9 scope conditions shared by ListJobsAsync and FacetAsync — the scalar filters
/// <inheritdoc/>
public async ValueTask<long> 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<string>();
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
Expand Down
16 changes: 15 additions & 1 deletion src/BackWave.Sqlite/SqliteJobStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2042,7 +2042,21 @@ public async ValueTask<IReadOnlyList<QueueStateCount>> CountJobsAsync(Cancellati
return counts;
}

// Builds the §5.9 scope conditions shared by ListJobsAsync and FacetAsync — the scalar filters
/// <inheritdoc/>
public async ValueTask<long> 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<string>();
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.
Expand Down
19 changes: 19 additions & 0 deletions src/BackWave/Monitor/BackWaveMonitor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,25 @@ public async ValueTask<IReadOnlyList<JobSnapshot>> ListJobsAsync(
return [.. records.Select(ToSnapshot)];
}

/// <summary>
/// The number of jobs matching a filter, across the whole matching population. Use it to show a
/// total next to a filtered <see cref="ListJobsAsync"/> page (for example "1,204 quarantined jobs
/// on the <c>lab</c> queue"). The count uses the same filters as <see cref="ListJobsAsync"/>:
/// state, queue, wire name, schedule id, and tag predicates, each narrowing the count.
/// <para>
/// The query's paging fields do not apply: <see cref="JobQuery.AfterSequence"/>,
/// <see cref="JobQuery.SortDirection"/>, and <see cref="JobQuery.MaxResults"/> are ignored, and
/// <see cref="MaxMonitorPageSize"/> does not cap the count.
/// </para>
/// </summary>
/// <param name="query">
/// The filter. When <c>null</c>, an empty query is used, which counts all jobs.
/// </param>
/// <param name="cancellationToken">Cancels the read.</param>
/// <returns>The number of matching jobs; zero when nothing matches.</returns>
public ValueTask<long> GetJobCountAsync(JobQuery? query = null, CancellationToken cancellationToken = default)
=> store.CountMatchingJobsAsync(query ?? new JobQuery(), cancellationToken);

/// <summary>
/// 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
Expand Down
Loading
Loading