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
16 changes: 16 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,22 @@ Set `ShedduellerOptions.EnableJobLogCapture = true` to enable durable capture of

Use `NotBeforeUtc` for delayed jobs. Use `JobIdempotencyKind.MethodAndArguments` to reuse an existing queued job with the same target method and serialized arguments.

Submit dependency graphs as one atomic batch. `DependsOn` blocks a job until every referenced job in the same batch is terminal; completed, failed, and canceled prerequisites all satisfy the dependency. Graphs may contain arbitrary fan-out, fan-in, and depth.

```csharp
var fetchProfile = JobEnqueueItem.Create<ImportJobs>(
(jobs, ct) => jobs.FetchProfileAsync(managerId, ct));
var fetchRates = JobEnqueueItem.Create<ImportJobs>(
(jobs, ct) => jobs.FetchRatesAsync(managerId, ct));
var aggregate = JobEnqueueItem.Create<ImportJobs>(
(jobs, ct) => jobs.AggregateAsync(managerId, ct))
.DependsOn([fetchProfile, fetchRates]);

await enqueuer.EnqueueManyAsync([fetchProfile, fetchRates, aggregate], cancellationToken);
```

Every prerequisite must be present in the submitted batch. Cycles, self-dependencies, duplicate items, missing prerequisites, and idempotency within dependency graphs are rejected before enqueueing.

## Recurring Schedules

Recurring schedules are keyed definitions. Calling `CreateOrUpdateAsync` at startup is the intended reconciliation model.
Expand Down
26 changes: 26 additions & 0 deletions src/Sheddueller.Dashboard/Components/Pages/JobDetail.razor
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,32 @@
HrefFactory="GroupFilterHref"
LinkAriaLabelPrefix="Filter jobs by group" />
</div>

@if (_detail.PrerequisiteJobIds.Count > 0)
{
<div class="job-detail-chip-section">
<span class="job-detail-label">Prerequisites</span>
<div class="job-detail-chip-list">
@foreach (var prerequisiteJobId in _detail.PrerequisiteJobIds)
{
<JobLink CssClass="job-detail-chip" JobId="prerequisiteJobId" />
}
</div>
</div>
}

@if (_detail.DependentJobIds.Count > 0)
{
<div class="job-detail-chip-section">
<span class="job-detail-label">Dependents</span>
<div class="job-detail-chip-list">
@foreach (var dependentJobId in _detail.DependentJobIds)
{
<JobLink CssClass="job-detail-chip" JobId="dependentJobId" />
}
</div>
</div>
}
</section>

<div class="job-detail-operations">
Expand Down
2 changes: 2 additions & 0 deletions src/Sheddueller.Dashboard/Components/Pages/Jobs.razor
Original file line number Diff line number Diff line change
Expand Up @@ -1261,6 +1261,7 @@
{ Kind: JobQueuePositionKind.Claimable, Position: { } position } => string.Create(CultureInfo.InvariantCulture, $"#{position}"),
{ Kind: JobQueuePositionKind.Claimable } => "Ready",
{ Kind: JobQueuePositionKind.Claimed } => "Running",
{ Kind: JobQueuePositionKind.WaitingForDependencies } => "Dependencies",
{ Kind: JobQueuePositionKind.BlockedByConcurrency } => "Blocked",
{ Kind: JobQueuePositionKind.RetryWaiting } => "Retry",
{ Kind: JobQueuePositionKind.Delayed } => "Delayed",
Expand All @@ -1275,6 +1276,7 @@
{
JobQueuePositionKind.Claimable => "claimable",
JobQueuePositionKind.Claimed => "claimed",
JobQueuePositionKind.WaitingForDependencies => "waiting",
JobQueuePositionKind.BlockedByConcurrency => "blocked",
JobQueuePositionKind.RetryWaiting => "waiting",
JobQueuePositionKind.Delayed => "delayed",
Expand Down
1 change: 1 addition & 0 deletions src/Sheddueller.Dashboard/Internal/DashboardFormat.cs
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,7 @@ public static string QueueKind(
JobQueuePositionKind.Claimable => "claimable",
JobQueuePositionKind.Delayed => "delayed",
JobQueuePositionKind.RetryWaiting => "retry_waiting",
JobQueuePositionKind.WaitingForDependencies => "waiting_for_dependencies",
JobQueuePositionKind.BlockedByConcurrency => "blocked_by_concurrency",
JobQueuePositionKind.Claimed => "running_active",
JobQueuePositionKind.Terminal => "terminal",
Expand Down
129 changes: 129 additions & 0 deletions src/Sheddueller.Postgres/Internal/Operations/EnqueueJobOperation.cs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ public static async ValueTask<IReadOnlyList<EnqueueJobResult>> ExecuteManyAsync(
}

var requestSnapshot = requests.ToArray();
ValidateDependencyGraph(requestSnapshot);
foreach (var request in requestSnapshot)
{
ArgumentNullException.ThrowIfNull(request);
Expand All @@ -50,13 +51,15 @@ public static async ValueTask<IReadOnlyList<EnqueueJobResult>> ExecuteManyAsync(
await CreateStagingTablesAsync(connection, transaction, cancellationToken).ConfigureAwait(false);
await CopyJobsAsync(connection, requestSnapshot, cancellationToken).ConfigureAwait(false);
await CopyGroupsAsync(connection, requestSnapshot, cancellationToken).ConfigureAwait(false);
await CopyDependenciesAsync(connection, requestSnapshot, cancellationToken).ConfigureAwait(false);
await CopyTagsAsync(connection, requestSnapshot, cancellationToken).ConfigureAwait(false);
await CopyEventsAsync(connection, requestSnapshot, cancellationToken).ConfigureAwait(false);
await EnsureNoDuplicateJobIdsAsync(context, connection, transaction, cancellationToken).ConfigureAwait(false);
await LockIdempotencyKeysAsync(context, connection, transaction, cancellationToken).ConfigureAwait(false);

var results = await InsertStagedJobsAsync(context, connection, transaction, cancellationToken).ConfigureAwait(false);
await InsertStagedGroupsAsync(context, connection, transaction, cancellationToken).ConfigureAwait(false);
await InsertStagedDependenciesAsync(context, connection, transaction, cancellationToken).ConfigureAwait(false);
await InsertStagedTagsAsync(context, connection, transaction, cancellationToken).ConfigureAwait(false);
await InsertStagedEventsAsync(context, connection, transaction, cancellationToken).ConfigureAwait(false);
await PostgresMetricsRollups.RecordStagedQueuedJobsAsync(context, connection, transaction, cancellationToken)
Expand Down Expand Up @@ -106,6 +109,11 @@ create temp table sheddueller_enqueue_groups (
group_key text not null
) on commit drop;

create temp table sheddueller_enqueue_dependencies (
job_id uuid not null,
prerequisite_job_id uuid not null
) on commit drop;

create temp table sheddueller_enqueue_tags (
job_id uuid not null,
ordinal integer not null,
Expand Down Expand Up @@ -265,6 +273,34 @@ from stdin (format binary)
await importer.CompleteAsync(cancellationToken).ConfigureAwait(false);
}

private static async ValueTask CopyDependenciesAsync(
NpgsqlConnection connection,
EnqueueJobRequest[] requests,
CancellationToken cancellationToken)
{
await using var importer = await connection.BeginBinaryImportAsync(
"""
copy sheddueller_enqueue_dependencies (
job_id,
prerequisite_job_id)
from stdin (format binary)
""",
cancellationToken)
.ConfigureAwait(false);

foreach (var request in requests)
{
foreach (var prerequisiteJobId in request.PrerequisiteJobIds ?? [])
{
await importer.StartRowAsync(cancellationToken).ConfigureAwait(false);
await importer.WriteAsync(request.JobId, NpgsqlDbType.Uuid, cancellationToken).ConfigureAwait(false);
await importer.WriteAsync(prerequisiteJobId, NpgsqlDbType.Uuid, cancellationToken).ConfigureAwait(false);
}
}

await importer.CompleteAsync(cancellationToken).ConfigureAwait(false);
}

private static async ValueTask CopyEventsAsync(
NpgsqlConnection connection,
EnqueueJobRequest[] requests,
Expand Down Expand Up @@ -537,6 +573,28 @@ from sheddueller_enqueue_tags tag
cancellationToken)
.ConfigureAwait(false);

private static async ValueTask InsertStagedDependenciesAsync(
PostgresOperationContext context,
NpgsqlConnection connection,
NpgsqlTransaction transaction,
CancellationToken cancellationToken)
=> await PostgresOperationContext.ExecuteCountAsync(
connection,
transaction,
$"""
insert into {context.Names.JobDependencies} (job_id, prerequisite_job_id)
select dependency.job_id, dependency.prerequisite_job_id
from sheddueller_enqueue_dependencies dependency
join sheddueller_enqueue_results dependent_result on dependent_result.job_id = dependency.job_id
join sheddueller_enqueue_results prerequisite_result on prerequisite_result.job_id = dependency.prerequisite_job_id
where dependent_result.was_enqueued = true
and prerequisite_result.was_enqueued = true
on conflict (job_id, prerequisite_job_id) do nothing;
""",
static _ => { },
cancellationToken)
.ConfigureAwait(false);

private static async ValueTask InsertStagedEventsAsync(
PostgresOperationContext context,
NpgsqlConnection connection,
Expand Down Expand Up @@ -618,6 +676,77 @@ private static async ValueTask WriteNullableAsync<T>(
await importer.WriteAsync(value.Value, dbType, cancellationToken).ConfigureAwait(false);
}

private static void ValidateDependencyGraph(IReadOnlyList<EnqueueJobRequest> requests)
{
var requestsById = new Dictionary<Guid, EnqueueJobRequest>();
foreach (var request in requests)
{
ArgumentNullException.ThrowIfNull(request);
if (!requestsById.TryAdd(request.JobId, request))
{
throw new InvalidOperationException($"Job '{request.JobId}' appears more than once in the batch.");
}
}

var hasDependencies = requests.Any(static request => request.PrerequisiteJobIds is { Count: > 0 });
if (hasDependencies && requests.Any(static request => request.IdempotencyKey is not null))
{
throw new ArgumentException("Jobs in a dependency graph cannot use idempotency.", nameof(requests));
}

foreach (var request in requests)
{
var prerequisites = request.PrerequisiteJobIds ?? [];
if (prerequisites.Count != prerequisites.Distinct().Count())
{
throw new ArgumentException($"Job '{request.JobId}' contains duplicate prerequisites.", nameof(requests));
}

foreach (var prerequisiteJobId in prerequisites)
{
if (prerequisiteJobId == request.JobId)
{
throw new ArgumentException($"Job '{request.JobId}' cannot depend on itself.", nameof(requests));
}

if (!requestsById.ContainsKey(prerequisiteJobId))
{
throw new ArgumentException(
$"Prerequisite job '{prerequisiteJobId}' for job '{request.JobId}' is not included in the batch.",
nameof(requests));
}
}
}

var visiting = new HashSet<Guid>();
var visited = new HashSet<Guid>();
foreach (var request in requests)
{
visit(request.JobId);
}

void visit(Guid jobId)
{
if (visited.Contains(jobId))
{
return;
}

if (!visiting.Add(jobId))
{
throw new ArgumentException("Job dependency graphs cannot contain cycles.", nameof(requests));
}

foreach (var prerequisiteJobId in requestsById[jobId].PrerequisiteJobIds ?? [])
{
visit(prerequisiteJobId);
}

visiting.Remove(jobId);
visited.Add(jobId);
}
}

private static async ValueTask WriteNullableAsync(
NpgsqlBinaryImporter importer,
string? value,
Expand Down
Loading