Files
jobtrackingapp/JobTrackerApi/Services/CvProcessingQueue.cs
T
cesnimda c3c5af8329
CI and Deploy / test (pull_request) Failing after 1m39s
CI and Deploy / deploy (pull_request) Has been skipped
feat(cv)!: queue durable processing
CV upload now returns 202 with an owner-scoped operation instead of holding the request through parsing. Existing review approval remains required.

BREAKING CHANGE: profile-cv upload responses use the durable operation contract.
2026-08-09 15:11:03 +02:00

130 lines
5.4 KiB
C#

using JobTrackerApi.Controllers;
using JobTrackerApi.Data;
using JobTrackerApi.Models;
using Microsoft.EntityFrameworkCore;
namespace JobTrackerApi.Services;
public sealed record CvProcessingOutcome(
bool Succeeded,
string? FailureCategory = null,
string? FailureMessage = null,
bool Retryable = false,
string? Provider = null,
string? Model = null,
string? RouteReason = null);
public interface ICvProcessingQueue
{
Task<AiOperationAdmissionResult?> EnqueueAsync(int runId, CancellationToken cancellationToken = default);
}
/// <summary>
/// Compatibility name for the CV producer bridge. Durable scheduling and execution are owned by
/// the shared AI operation queue; this type does not keep an in-memory CV queue.
/// </summary>
public sealed class CvProcessingQueue(AiOperationAdmission admission) : ICvProcessingQueue
{
public const string TaskType = "cv.process";
public const string SubjectType = "cv_extraction_run";
public async Task<AiOperationAdmissionResult?> EnqueueAsync(int runId, CancellationToken cancellationToken = default)
=> await admission.EnqueueAsync(
TaskType,
$"run:{runId}",
SubjectType,
runId.ToString(System.Globalization.CultureInfo.InvariantCulture),
AiOperationPriorities.UserVisible,
cancellationToken);
}
public sealed class NoOpCvProcessingQueue : ICvProcessingQueue
{
public static readonly NoOpCvProcessingQueue Instance = new();
public Task<AiOperationAdmissionResult?> EnqueueAsync(int runId, CancellationToken cancellationToken = default)
=> Task.FromResult<AiOperationAdmissionResult?>(null);
}
public sealed class CvProcessingOperationHandler : IAiOperationHandler
{
public string TaskType => CvProcessingQueue.TaskType;
public async Task<AiOperationExecutionResult> ExecuteAsync(
AiOperationExecutionContext context,
IServiceProvider services,
CancellationToken cancellationToken)
{
if (!string.Equals(context.Lease.SubjectType, CvProcessingQueue.SubjectType, StringComparison.Ordinal) ||
!int.TryParse(context.Lease.SubjectId, out var runId) || runId <= 0)
{
throw new AiOperationFailure("invalid_cv_run", "The CV processing operation has an invalid run reference.", retryable: false);
}
CvProcessingOutcome? outcome;
try
{
outcome = await services.GetRequiredService<ProfileCvController>()
.ProcessQueuedRunAsync(runId, cancellationToken);
}
catch (OperationCanceledException)
{
var operation = await services.GetRequiredService<UserOperationStore>()
.GetAsync(context.Lease.OperationId, CancellationToken.None);
await SetRunStatusAsync(
services,
runId,
operation?.CancellationRequestedAtUtc is null ? "queued" : "cancelled",
operation?.CancellationRequestedAtUtc is null ? "CV processing timed out and may be retried." : "CV processing was cancelled.");
throw;
}
if (outcome is null)
throw new AiOperationFailure("cv_run_not_found", "The CV processing run is no longer available.", retryable: false);
if (!outcome.Succeeded)
{
if (outcome.Retryable)
{
var operation = await services.GetRequiredService<UserOperationStore>()
.GetAsync(context.Lease.OperationId, cancellationToken);
var canRetry = operation is not null && operation.AttemptCount < operation.MaxAttempts &&
(operation.DeadlineAtUtc is null || operation.DeadlineAtUtc > DateTime.UtcNow);
if (!canRetry)
await SetRunStatusAsync(services, runId, "failed", outcome.FailureMessage ?? "CV processing failed.");
}
if (outcome.Provider is not null || outcome.Model is not null || outcome.RouteReason is not null)
{
throw new AiGenerationException(
outcome.FailureCategory ?? "cv_processing_failed",
outcome.FailureMessage ?? "CV processing failed.",
outcome.Retryable,
outcome.Provider,
outcome.Model,
outcome.RouteReason);
}
throw new AiOperationFailure(
outcome.FailureCategory ?? "cv_processing_failed",
outcome.FailureMessage ?? "CV processing failed.",
outcome.Retryable);
}
return new AiOperationExecutionResult(
$"/api/profile-cv/runs/{runId}/diff",
outcome.Provider,
outcome.Model,
outcome.RouteReason ?? "cv_pipeline");
}
private static Task<int> SetRunStatusAsync(IServiceProvider services, int runId, string status, string message)
{
var completedAtUtc = status == "failed" || status == "cancelled" ? DateTimeOffset.UtcNow : (DateTimeOffset?)null;
return services.GetRequiredService<JobTrackerContext>().CvExtractionRuns
.Where(run => run.Id == runId)
.ExecuteUpdateAsync(setters => setters
.SetProperty(run => run.Status, status)
.SetProperty(run => run.ErrorMessage, message)
.SetProperty(run => run.CompletedAtUtc, completedAtUtc),
CancellationToken.None);
}
}