From 11fb0351e4c307b7054d3aca08ab82bbcfdb22fd Mon Sep 17 00:00:00 2001 From: "jack.lewis" Date: Fri, 3 Jul 2026 16:54:36 +0100 Subject: [PATCH 1/4] Add optional correlation id for tracking job completion --- CLAUDE.md | 9 +- .../Data/BuilderJob.cs | 7 ++ .../Features/Jobs/CreateJob.cs | 1 + .../Features/Jobs/JobInstruction.cs | 7 ++ .../Features/Jobs/JobResponse.cs | 7 ++ .../Jobs/TextBuildJob.cs | 4 +- ...5_AddCorrelationIdToBuilderJob.Designer.cs | 112 ++++++++++++++++++ ...0703143635_AddCorrelationIdToBuilderJob.cs | 28 +++++ .../BuilderDbContextModelSnapshot.cs | 4 + .../JobCompletionNotification.cs | 3 +- .../TextServices.Search.Api.csproj | 1 + .../BuilderApi/JobResponseTests.cs | 22 ++++ .../BuilderApi/SnsJobNotifierTests.cs | 2 +- .../BuilderApi/TextBuildJobTests.cs | 29 ++++- 14 files changed, 229 insertions(+), 7 deletions(-) create mode 100644 src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.Designer.cs create mode 100644 src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.cs diff --git a/CLAUDE.md b/CLAUDE.md index 8d0f91a..18c7d85 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -91,9 +91,12 @@ Implementations: `FileSystemTextStore`, `S3TextStore`. ```json { "id": "2/books/my-book", - "sourceUri": "https://iiif.wellcomecollection.org/presentation/b21211024" + "sourceUri": "https://iiif.wellcomecollection.org/presentation/b21211024", + "correlationId": "caller-supplied-guid" } ``` +`correlationId` is optional and opaque — stored against the job and echoed back in the response and in the completion notification. + or with inline pages: ```json { @@ -104,7 +107,9 @@ or with inline pages: } ``` -**`BuilderJob` entity** (PostgreSQL, snake_case naming): `Id`, `SourceUri`, `SourceDataJson`, `Status`, `Created`, `Started`, `Finished`, `TotalPages`, `PagesCompleted`, `TotalWordCount`, `TotalImageCount`, `Errors`, `HangfireJobId`, `Services` (bitmask), `Title`, `CustomTypesJson`. +**`BuilderJob` entity** (PostgreSQL, snake_case naming): `Id`, `SourceUri`, `SourceDataJson`, `Status`, `Created`, `Started`, `Finished`, `TotalPages`, `PagesCompleted`, `TotalWordCount`, `TotalImageCount`, `Errors`, `HangfireJobId`, `Services` (bitmask), `Title`, `CustomTypesJson`, `CorrelationId`. + +**`CorrelationId`**: opaque string supplied by the caller in the job instruction (e.g. a pipelineJob GUID from iiif-presentation). Stored verbatim against the job and echoed back in `JobResponse` and in `JobCompletionNotification`, so callers can tie a completion notification back to the request that created the job. Not used for HTTP-level request tracing — that's handled separately by `CorrelationIdMiddleware` in `TextServices.Infrastructure`. **`TextBuildJob`** (Hangfire): fetches manifest/resources concurrently (bounded by `MaxConcurrentPageFetches`), feeds `TextBuilder`, persists all artefacts, records per-page warnings without aborting the job. diff --git a/src/TextServices.Builder.Api/Data/BuilderJob.cs b/src/TextServices.Builder.Api/Data/BuilderJob.cs index 11cfb2f..0540c72 100644 --- a/src/TextServices.Builder.Api/Data/BuilderJob.cs +++ b/src/TextServices.Builder.Api/Data/BuilderJob.cs @@ -58,4 +58,11 @@ public class BuilderJob /// A value of 0 means processing completed but no derivatives could be produced. /// public int? FulfilledServices { get; set; } + + /// + /// Opaque identifier supplied by the caller when the job was created (e.g. a + /// pipelineJob GUID). Echoed back in + /// so callers can tie a completion notification to the request that created the job. + /// + public string? CorrelationId { get; set; } } diff --git a/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs b/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs index 98076a8..301918b 100644 --- a/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs +++ b/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs @@ -43,6 +43,7 @@ public async Task Handle(CreateJobRequest request, Cancellation Created = DateTimeOffset.UtcNow, Services = (int)instruction.Services, Title = instruction.Title, + CorrelationId = instruction.CorrelationId, CustomTypesJson = instruction.CustomTypes != null ? System.Text.Json.JsonSerializer.Serialize(instruction.CustomTypes) : null, diff --git a/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs b/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs index 8768bb5..baa62be 100644 --- a/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs +++ b/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs @@ -44,6 +44,13 @@ public class JobInstruction : IValidatableObject /// public Dictionary? CustomTypes { get; set; } + /// + /// Opaque identifier supplied by the caller (e.g. a pipelineJob GUID). Stored against + /// the job and echoed back in and completion notifications so + /// callers can tie a job through to the request that created it. + /// + public string? CorrelationId { get; set; } + public IEnumerable Validate(ValidationContext validationContext) { if (Id.StartsWith('/')) diff --git a/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs b/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs index de27682..a9be356 100644 --- a/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs +++ b/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs @@ -25,6 +25,12 @@ public class JobResponse public int TotalImageCount { get; set; } public string? Errors { get; set; } + /// + /// Opaque identifier supplied by the caller when the job was created. Echoed back + /// unchanged so callers can tie this job to the request that created it. + /// + public string? CorrelationId { get; set; } + /// /// The services that were actually produced during the most recent run. /// Null for jobs processed before this field was introduced. @@ -119,6 +125,7 @@ public static JobResponse From(BuilderJob job, TextServicesOptions options) TotalWordCount = job.TotalWordCount, TotalImageCount = job.TotalImageCount, Errors = job.Errors, + CorrelationId = job.CorrelationId, SearchV1 = searchV1, AutocompleteV1 = autocompleteV1, SearchV2 = searchV2, diff --git a/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs b/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs index 95b4c2d..02168ca 100644 --- a/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs +++ b/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs @@ -83,7 +83,7 @@ public async Task ExecuteAsync(string jobId, IJobCancellationToken cancellationT await jobNotifier.Notify( new JobCompletionNotification(job.Id, job.Status, job.Finished, - job.TotalPages, job.TotalWordCount, job.Errors), + job.TotalPages, job.TotalWordCount, job.Errors, job.CorrelationId), cancellationToken.ShutdownToken); } catch (Exception ex) @@ -97,7 +97,7 @@ await jobNotifier.Notify( await jobNotifier.Notify( new JobCompletionNotification(job.Id, job.Status, job.Finished, - job.TotalPages, job.TotalWordCount, job.Errors), + job.TotalPages, job.TotalWordCount, job.Errors, job.CorrelationId), cancellationToken.ShutdownToken); } } diff --git a/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.Designer.cs b/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.Designer.cs new file mode 100644 index 0000000..2adb217 --- /dev/null +++ b/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.Designer.cs @@ -0,0 +1,112 @@ +// +using System; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; +using TextServices.Builder.Api.Data; + +#nullable disable + +namespace TextServices.Builder.Api.Migrations +{ + [DbContext(typeof(BuilderDbContext))] + [Migration("20260703143635_AddCorrelationIdToBuilderJob")] + partial class AddCorrelationIdToBuilderJob + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.7") + .HasAnnotation("Relational:MaxIdentifierLength", 63); + + NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + + modelBuilder.Entity("TextServices.Builder.Api.Data.BuilderJob", b => + { + b.Property("Id") + .HasMaxLength(500) + .HasColumnType("character varying(500)") + .HasColumnName("id"); + + b.Property("CorrelationId") + .HasColumnType("text") + .HasColumnName("correlation_id"); + + b.Property("Created") + .HasColumnType("timestamp with time zone") + .HasColumnName("created"); + + b.Property("CustomTypesJson") + .HasColumnType("text") + .HasColumnName("custom_types_json"); + + b.Property("Errors") + .HasColumnType("text") + .HasColumnName("errors"); + + b.Property("Finished") + .HasColumnType("timestamp with time zone") + .HasColumnName("finished"); + + b.Property("FulfilledServices") + .HasColumnType("integer") + .HasColumnName("fulfilled_services"); + + b.Property("HangfireJobId") + .HasColumnType("text") + .HasColumnName("hangfire_job_id"); + + b.Property("PagesCompleted") + .HasColumnType("integer") + .HasColumnName("pages_completed"); + + b.Property("Services") + .HasColumnType("integer") + .HasColumnName("services"); + + b.Property("SourceDataJson") + .HasColumnType("text") + .HasColumnName("source_data_json"); + + b.Property("SourceUri") + .HasColumnType("text") + .HasColumnName("source_uri"); + + b.Property("Started") + .HasColumnType("timestamp with time zone") + .HasColumnName("started"); + + b.Property("Status") + .IsRequired() + .HasColumnType("text") + .HasColumnName("status"); + + b.Property("Title") + .HasColumnType("text") + .HasColumnName("title"); + + b.Property("TotalImageCount") + .HasColumnType("integer") + .HasColumnName("total_image_count"); + + b.Property("TotalPages") + .HasColumnType("integer") + .HasColumnName("total_pages"); + + b.Property("TotalWordCount") + .HasColumnType("integer") + .HasColumnName("total_word_count"); + + b.HasKey("Id") + .HasName("pk_jobs"); + + b.ToTable("jobs", (string)null); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.cs b/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.cs new file mode 100644 index 0000000..1ad3bf9 --- /dev/null +++ b/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.cs @@ -0,0 +1,28 @@ +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace TextServices.Builder.Api.Migrations +{ + /// + public partial class AddCorrelationIdToBuilderJob : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.AddColumn( + name: "correlation_id", + table: "jobs", + type: "text", + nullable: true); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropColumn( + name: "correlation_id", + table: "jobs"); + } + } +} diff --git a/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs b/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs index 846902a..f80e576 100644 --- a/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs +++ b/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs @@ -29,6 +29,10 @@ protected override void BuildModel(ModelBuilder modelBuilder) .HasColumnType("character varying(500)") .HasColumnName("id"); + b.Property("CorrelationId") + .HasColumnType("text") + .HasColumnName("correlation_id"); + b.Property("Created") .HasColumnType("timestamp with time zone") .HasColumnName("created"); diff --git a/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs b/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs index 45638c7..57c673f 100644 --- a/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs +++ b/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs @@ -8,4 +8,5 @@ public record JobCompletionNotification( DateTimeOffset? Finished, int TotalPages, int TotalWordCount, - string? Errors); + string? Errors, + string? CorrelationId); diff --git a/src/TextServices.Search.Api/TextServices.Search.Api.csproj b/src/TextServices.Search.Api/TextServices.Search.Api.csproj index 2fa1628..518c738 100644 --- a/src/TextServices.Search.Api/TextServices.Search.Api.csproj +++ b/src/TextServices.Search.Api/TextServices.Search.Api.csproj @@ -12,6 +12,7 @@ + diff --git a/src/TextServices.Tests/BuilderApi/JobResponseTests.cs b/src/TextServices.Tests/BuilderApi/JobResponseTests.cs index a4f3ff4..3ad8b42 100644 --- a/src/TextServices.Tests/BuilderApi/JobResponseTests.cs +++ b/src/TextServices.Tests/BuilderApi/JobResponseTests.cs @@ -22,6 +22,28 @@ private static BuilderJob CompletedJob( Status = JobStatus.Completed, }; + // ------------------------------------------------------------------------- + // CorrelationId — echoed back unchanged + // ------------------------------------------------------------------------- + + [Fact] + public void From_JobWithCorrelationId_EchoesCorrelationId() + { + var job = CompletedJob(); + job.CorrelationId = "caller-supplied-guid"; + + var response = JobResponse.From(job, Options()); + + response.CorrelationId.ShouldBe("caller-supplied-guid"); + } + + [Fact] + public void From_JobWithoutCorrelationId_CorrelationIdIsNull() + { + var response = JobResponse.From(CompletedJob(), Options()); + response.CorrelationId.ShouldBeNull(); + } + // ------------------------------------------------------------------------- // FulfilledServices field // ------------------------------------------------------------------------- diff --git a/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs b/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs index 18ce0e6..b2d678f 100644 --- a/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs +++ b/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs @@ -15,7 +15,7 @@ public sealed class SnsJobNotifierTests private static JobCompletionNotification MakeNotification( string jobId = "my/job", JobStatus status = JobStatus.Completed) => - new(jobId, status, DateTimeOffset.UtcNow, 10, 500, null); + new(jobId, status, DateTimeOffset.UtcNow, 10, 500, null, null); // ------------------------------------------------------------------------- // No-op when TopicArn is absent diff --git a/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs b/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs index 6463f40..f355d7a 100644 --- a/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs +++ b/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs @@ -598,6 +598,32 @@ public async Task ExecuteAsync_OnSuccess_NotifiesWithCompletedStatus() notifier.Captured[0].Status.ShouldBe(JobStatus.Completed); } + [Fact] + public async Task ExecuteAsync_OnSuccess_NotificationEchoesCorrelationId() + { + var pages = new List + { + new() { Id = "https://example.org/c/1", Width = 1000, Height = 1500, + TextUri = "https://example.org/alto/1.xml" }, + }; + + var job = await CreateJob("test/notify-correlation-id", + sourceDataJson: JsonSerializer.Serialize(pages), + correlationId: "caller-supplied-guid"); + + var altoFetcher = FakeAlto(new Dictionary + { + ["https://example.org/alto/1.xml"] = SampleAlto("hello world"), + }); + + var notifier = new CapturingJobNotifier(); + var sut = MakeJob(altoFetcher: altoFetcher, jobNotifier: notifier); + await sut.ExecuteAsync(job.Id, FakeCancellationToken.Instance); + + notifier.Captured.ShouldHaveSingleItem(); + notifier.Captured[0].CorrelationId.ShouldBe("caller-supplied-guid"); + } + [Fact] public async Task ExecuteAsync_OnFailure_NotifiesWithFailedStatus() { @@ -622,7 +648,7 @@ public async Task ExecuteAsync_OnFailure_NotifiesWithFailedStatus() private async Task CreateJob(string id, string? sourceUri = null, string? sourceDataJson = null, - JobServices services = JobServices.All) + JobServices services = JobServices.All, string? correlationId = null) { var job = new BuilderJob { @@ -630,6 +656,7 @@ private async Task CreateJob(string id, SourceUri = sourceUri, SourceDataJson = sourceDataJson, Services = (int)services, + CorrelationId = correlationId, }; _db.Jobs.Add(job); await _db.SaveChangesAsync(); From a246c2d8ff8940f98676762072a50fd3b6c0ad48 Mon Sep 17 00:00:00 2001 From: "jack.lewis" Date: Mon, 6 Jul 2026 12:00:09 +0100 Subject: [PATCH 2/4] update to use invocation count --- CLAUDE.md | 9 +- docker-compose.yml | 9 +- .../Data/BuilderJob.cs | 8 +- .../Features/Jobs/CreateJob.cs | 1 - .../Features/Jobs/JobInstruction.cs | 7 - .../Features/Jobs/JobResponse.cs | 8 +- .../Features/Jobs/ReprocessJob.cs | 1 + .../Jobs/TextBuildJob.cs | 4 +- ...ddInvocationCountToBuilderJob.Designer.cs} | 224 +++++++++--------- ...6104802_AddInvocationCountToBuilderJob.cs} | 57 ++--- .../BuilderDbContextModelSnapshot.cs | 8 +- .../JobCompletionNotification.cs | 2 +- .../BuilderApi/JobResponseTests.cs | 22 +- .../BuilderApi/SnsJobNotifierTests.cs | 2 +- .../BuilderApi/TextBuildJobTests.cs | 12 +- 15 files changed, 186 insertions(+), 188 deletions(-) rename src/TextServices.Builder.Api/Migrations/{20260703143635_AddCorrelationIdToBuilderJob.Designer.cs => 20260706104802_AddInvocationCountToBuilderJob.Designer.cs} (91%) rename src/TextServices.Builder.Api/Migrations/{20260703143635_AddCorrelationIdToBuilderJob.cs => 20260706104802_AddInvocationCountToBuilderJob.cs} (62%) diff --git a/CLAUDE.md b/CLAUDE.md index 18c7d85..dbf6902 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -91,12 +91,9 @@ Implementations: `FileSystemTextStore`, `S3TextStore`. ```json { "id": "2/books/my-book", - "sourceUri": "https://iiif.wellcomecollection.org/presentation/b21211024", - "correlationId": "caller-supplied-guid" + "sourceUri": "https://iiif.wellcomecollection.org/presentation/b21211024" } ``` -`correlationId` is optional and opaque — stored against the job and echoed back in the response and in the completion notification. - or with inline pages: ```json { @@ -107,9 +104,9 @@ or with inline pages: } ``` -**`BuilderJob` entity** (PostgreSQL, snake_case naming): `Id`, `SourceUri`, `SourceDataJson`, `Status`, `Created`, `Started`, `Finished`, `TotalPages`, `PagesCompleted`, `TotalWordCount`, `TotalImageCount`, `Errors`, `HangfireJobId`, `Services` (bitmask), `Title`, `CustomTypesJson`, `CorrelationId`. +**`BuilderJob` entity** (PostgreSQL, snake_case naming): `Id`, `SourceUri`, `SourceDataJson`, `Status`, `Created`, `Started`, `Finished`, `TotalPages`, `PagesCompleted`, `TotalWordCount`, `TotalImageCount`, `Errors`, `HangfireJobId`, `Services` (bitmask), `Title`, `CustomTypesJson`, `InvocationCount`. -**`CorrelationId`**: opaque string supplied by the caller in the job instruction (e.g. a pipelineJob GUID from iiif-presentation). Stored verbatim against the job and echoed back in `JobResponse` and in `JobCompletionNotification`, so callers can tie a completion notification back to the request that created the job. Not used for HTTP-level request tracing — that's handled separately by `CorrelationIdMiddleware` in `TextServices.Infrastructure`. +**`InvocationCount`**: tracks how many times the job processor has run for this job. Set to `1` on initial `POST`, incremented by 1 on every `PUT` (reprocess). Existing rows default to `1`. Returned on the `POST`/`PUT`/`GET` `JobResponse` and included in `JobCompletionNotification`, so callers can tell which run a completion notification belongs to. Server-managed only — not settable by the caller. **`TextBuildJob`** (Hangfire): fetches manifest/resources concurrently (bounded by `MaxConcurrentPageFetches`), feeds `TextBuilder`, persists all artefacts, records per-page warnings without aborting the job. diff --git a/docker-compose.yml b/docker-compose.yml index 6eeee41..24dcd9a 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -8,10 +8,11 @@ services: image: postgres:14 hostname: postgres ports: - - "5452:5432" + - "5453:5432" volumes: - txt_postgres_data:/var/lib/postgresql/data - txt_postgres_data_backups:/backups + - $HOME\.aws:/root/.aws:ro environment: - POSTGRES_HOST=postgres - POSTGRES_PORT=5432 @@ -29,6 +30,7 @@ services: - postgres volumes: - txt_textservices_data:/data + - $HOME\.aws:/root/.aws:ro environment: - ASPNETCORE_ENVIRONMENT=Development - RunMigrations=true @@ -36,6 +38,8 @@ services: - TextServices__SearchApiBaseUrl=http://localhost:5294 - TextServices__Storage__FileSystem__RootPath=/data - CorsAllowedOrigins__0=http://localhost:5100 + env_file: + - .env search: build: @@ -45,11 +49,14 @@ services: - "5294:8080" volumes: - txt_textservices_data:/data + - $HOME\.aws:/root/.aws:ro environment: - ASPNETCORE_ENVIRONMENT=Development - TextServices__BaseUrl=http://localhost:5294 - TextServices__Storage__FileSystem__RootPath=/data - CorsAllowedOrigins__0=http://localhost:5100 + env_file: + - .env demo: build: diff --git a/src/TextServices.Builder.Api/Data/BuilderJob.cs b/src/TextServices.Builder.Api/Data/BuilderJob.cs index 0540c72..4a00787 100644 --- a/src/TextServices.Builder.Api/Data/BuilderJob.cs +++ b/src/TextServices.Builder.Api/Data/BuilderJob.cs @@ -60,9 +60,9 @@ public class BuilderJob public int? FulfilledServices { get; set; } /// - /// Opaque identifier supplied by the caller when the job was created (e.g. a - /// pipelineJob GUID). Echoed back in - /// so callers can tie a completion notification to the request that created the job. + /// Number of times the job processor has been invoked for this job. Set to 1 on + /// initial creation (POST) and incremented by 1 on every reprocess (PUT), + /// so callers can tell which run a completion notification belongs to. /// - public string? CorrelationId { get; set; } + public int InvocationCount { get; set; } = 1; } diff --git a/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs b/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs index 301918b..98076a8 100644 --- a/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs +++ b/src/TextServices.Builder.Api/Features/Jobs/CreateJob.cs @@ -43,7 +43,6 @@ public async Task Handle(CreateJobRequest request, Cancellation Created = DateTimeOffset.UtcNow, Services = (int)instruction.Services, Title = instruction.Title, - CorrelationId = instruction.CorrelationId, CustomTypesJson = instruction.CustomTypes != null ? System.Text.Json.JsonSerializer.Serialize(instruction.CustomTypes) : null, diff --git a/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs b/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs index baa62be..8768bb5 100644 --- a/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs +++ b/src/TextServices.Builder.Api/Features/Jobs/JobInstruction.cs @@ -44,13 +44,6 @@ public class JobInstruction : IValidatableObject /// public Dictionary? CustomTypes { get; set; } - /// - /// Opaque identifier supplied by the caller (e.g. a pipelineJob GUID). Stored against - /// the job and echoed back in and completion notifications so - /// callers can tie a job through to the request that created it. - /// - public string? CorrelationId { get; set; } - public IEnumerable Validate(ValidationContext validationContext) { if (Id.StartsWith('/')) diff --git a/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs b/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs index a9be356..abd56d6 100644 --- a/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs +++ b/src/TextServices.Builder.Api/Features/Jobs/JobResponse.cs @@ -26,10 +26,10 @@ public class JobResponse public string? Errors { get; set; } /// - /// Opaque identifier supplied by the caller when the job was created. Echoed back - /// unchanged so callers can tie this job to the request that created it. + /// Number of times the job processor has been invoked for this job. 1 on initial + /// creation, incremented on every reprocess (PUT). /// - public string? CorrelationId { get; set; } + public int InvocationCount { get; set; } = 1; /// /// The services that were actually produced during the most recent run. @@ -125,7 +125,7 @@ public static JobResponse From(BuilderJob job, TextServicesOptions options) TotalWordCount = job.TotalWordCount, TotalImageCount = job.TotalImageCount, Errors = job.Errors, - CorrelationId = job.CorrelationId, + InvocationCount = job.InvocationCount, SearchV1 = searchV1, AutocompleteV1 = autocompleteV1, SearchV2 = searchV2, diff --git a/src/TextServices.Builder.Api/Features/Jobs/ReprocessJob.cs b/src/TextServices.Builder.Api/Features/Jobs/ReprocessJob.cs index 52372cb..8c2f2ce 100644 --- a/src/TextServices.Builder.Api/Features/Jobs/ReprocessJob.cs +++ b/src/TextServices.Builder.Api/Features/Jobs/ReprocessJob.cs @@ -61,6 +61,7 @@ public async Task Handle(ReprocessJobRequest request, Cancel job.TotalImageCount = 0; job.Errors = null; job.HangfireJobId = null; + job.InvocationCount++; await db.SaveChangesAsync(ct); diff --git a/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs b/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs index 02168ca..8332848 100644 --- a/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs +++ b/src/TextServices.Builder.Api/Jobs/TextBuildJob.cs @@ -83,7 +83,7 @@ public async Task ExecuteAsync(string jobId, IJobCancellationToken cancellationT await jobNotifier.Notify( new JobCompletionNotification(job.Id, job.Status, job.Finished, - job.TotalPages, job.TotalWordCount, job.Errors, job.CorrelationId), + job.TotalPages, job.TotalWordCount, job.Errors, job.InvocationCount), cancellationToken.ShutdownToken); } catch (Exception ex) @@ -97,7 +97,7 @@ await jobNotifier.Notify( await jobNotifier.Notify( new JobCompletionNotification(job.Id, job.Status, job.Finished, - job.TotalPages, job.TotalWordCount, job.Errors, job.CorrelationId), + job.TotalPages, job.TotalWordCount, job.Errors, job.InvocationCount), cancellationToken.ShutdownToken); } } diff --git a/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.Designer.cs b/src/TextServices.Builder.Api/Migrations/20260706104802_AddInvocationCountToBuilderJob.Designer.cs similarity index 91% rename from src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.Designer.cs rename to src/TextServices.Builder.Api/Migrations/20260706104802_AddInvocationCountToBuilderJob.Designer.cs index 2adb217..10023ef 100644 --- a/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.Designer.cs +++ b/src/TextServices.Builder.Api/Migrations/20260706104802_AddInvocationCountToBuilderJob.Designer.cs @@ -1,112 +1,112 @@ -// -using System; -using Microsoft.EntityFrameworkCore; -using Microsoft.EntityFrameworkCore.Infrastructure; -using Microsoft.EntityFrameworkCore.Migrations; -using Microsoft.EntityFrameworkCore.Storage.ValueConversion; -using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; -using TextServices.Builder.Api.Data; - -#nullable disable - -namespace TextServices.Builder.Api.Migrations -{ - [DbContext(typeof(BuilderDbContext))] - [Migration("20260703143635_AddCorrelationIdToBuilderJob")] - partial class AddCorrelationIdToBuilderJob - { - /// - protected override void BuildTargetModel(ModelBuilder modelBuilder) - { -#pragma warning disable 612, 618 - modelBuilder - .HasAnnotation("ProductVersion", "10.0.7") - .HasAnnotation("Relational:MaxIdentifierLength", 63); - - NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); - - modelBuilder.Entity("TextServices.Builder.Api.Data.BuilderJob", b => - { - b.Property("Id") - .HasMaxLength(500) - .HasColumnType("character varying(500)") - .HasColumnName("id"); - - b.Property("CorrelationId") - .HasColumnType("text") - .HasColumnName("correlation_id"); - - b.Property("Created") - .HasColumnType("timestamp with time zone") - .HasColumnName("created"); - - b.Property("CustomTypesJson") - .HasColumnType("text") - .HasColumnName("custom_types_json"); - - b.Property("Errors") - .HasColumnType("text") - .HasColumnName("errors"); - - b.Property("Finished") - .HasColumnType("timestamp with time zone") - .HasColumnName("finished"); - - b.Property("FulfilledServices") - .HasColumnType("integer") - .HasColumnName("fulfilled_services"); - - b.Property("HangfireJobId") - .HasColumnType("text") - .HasColumnName("hangfire_job_id"); - - b.Property("PagesCompleted") - .HasColumnType("integer") - .HasColumnName("pages_completed"); - - b.Property("Services") - .HasColumnType("integer") - .HasColumnName("services"); - - b.Property("SourceDataJson") - .HasColumnType("text") - .HasColumnName("source_data_json"); - - b.Property("SourceUri") - .HasColumnType("text") - .HasColumnName("source_uri"); - - b.Property("Started") - .HasColumnType("timestamp with time zone") - .HasColumnName("started"); - - b.Property("Status") - .IsRequired() - .HasColumnType("text") - .HasColumnName("status"); - - b.Property("Title") - .HasColumnType("text") - .HasColumnName("title"); - - b.Property("TotalImageCount") - .HasColumnType("integer") - .HasColumnName("total_image_count"); - - b.Property("TotalPages") - .HasColumnType("integer") - .HasColumnName("total_pages"); - - b.Property("TotalWordCount") - .HasColumnType("integer") - .HasColumnName("total_word_count"); - - b.HasKey("Id") - .HasName("pk_jobs"); - - b.ToTable("jobs", (string)null); - }); -#pragma warning restore 612, 618 - } - } -} +// +using System; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; +using TextServices.Builder.Api.Data; + +#nullable disable + +namespace TextServices.Builder.Api.Migrations +{ + [DbContext(typeof(BuilderDbContext))] + [Migration("20260706104802_AddInvocationCountToBuilderJob")] + partial class AddInvocationCountToBuilderJob + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.7") + .HasAnnotation("Relational:MaxIdentifierLength", 63); + + NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + + modelBuilder.Entity("TextServices.Builder.Api.Data.BuilderJob", b => + { + b.Property("Id") + .HasMaxLength(500) + .HasColumnType("character varying(500)") + .HasColumnName("id"); + + b.Property("Created") + .HasColumnType("timestamp with time zone") + .HasColumnName("created"); + + b.Property("CustomTypesJson") + .HasColumnType("text") + .HasColumnName("custom_types_json"); + + b.Property("Errors") + .HasColumnType("text") + .HasColumnName("errors"); + + b.Property("Finished") + .HasColumnType("timestamp with time zone") + .HasColumnName("finished"); + + b.Property("FulfilledServices") + .HasColumnType("integer") + .HasColumnName("fulfilled_services"); + + b.Property("HangfireJobId") + .HasColumnType("text") + .HasColumnName("hangfire_job_id"); + + b.Property("InvocationCount") + .HasColumnType("integer") + .HasColumnName("invocation_count"); + + b.Property("PagesCompleted") + .HasColumnType("integer") + .HasColumnName("pages_completed"); + + b.Property("Services") + .HasColumnType("integer") + .HasColumnName("services"); + + b.Property("SourceDataJson") + .HasColumnType("text") + .HasColumnName("source_data_json"); + + b.Property("SourceUri") + .HasColumnType("text") + .HasColumnName("source_uri"); + + b.Property("Started") + .HasColumnType("timestamp with time zone") + .HasColumnName("started"); + + b.Property("Status") + .IsRequired() + .HasColumnType("text") + .HasColumnName("status"); + + b.Property("Title") + .HasColumnType("text") + .HasColumnName("title"); + + b.Property("TotalImageCount") + .HasColumnType("integer") + .HasColumnName("total_image_count"); + + b.Property("TotalPages") + .HasColumnType("integer") + .HasColumnName("total_pages"); + + b.Property("TotalWordCount") + .HasColumnType("integer") + .HasColumnName("total_word_count"); + + b.HasKey("Id") + .HasName("pk_jobs"); + + b.ToTable("jobs", (string)null); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.cs b/src/TextServices.Builder.Api/Migrations/20260706104802_AddInvocationCountToBuilderJob.cs similarity index 62% rename from src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.cs rename to src/TextServices.Builder.Api/Migrations/20260706104802_AddInvocationCountToBuilderJob.cs index 1ad3bf9..010c0d1 100644 --- a/src/TextServices.Builder.Api/Migrations/20260703143635_AddCorrelationIdToBuilderJob.cs +++ b/src/TextServices.Builder.Api/Migrations/20260706104802_AddInvocationCountToBuilderJob.cs @@ -1,28 +1,29 @@ -using Microsoft.EntityFrameworkCore.Migrations; - -#nullable disable - -namespace TextServices.Builder.Api.Migrations -{ - /// - public partial class AddCorrelationIdToBuilderJob : Migration - { - /// - protected override void Up(MigrationBuilder migrationBuilder) - { - migrationBuilder.AddColumn( - name: "correlation_id", - table: "jobs", - type: "text", - nullable: true); - } - - /// - protected override void Down(MigrationBuilder migrationBuilder) - { - migrationBuilder.DropColumn( - name: "correlation_id", - table: "jobs"); - } - } -} +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace TextServices.Builder.Api.Migrations +{ + /// + public partial class AddInvocationCountToBuilderJob : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.AddColumn( + name: "invocation_count", + table: "jobs", + type: "integer", + nullable: false, + defaultValue: 1); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropColumn( + name: "invocation_count", + table: "jobs"); + } + } +} diff --git a/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs b/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs index f80e576..d548f36 100644 --- a/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs +++ b/src/TextServices.Builder.Api/Migrations/BuilderDbContextModelSnapshot.cs @@ -29,10 +29,6 @@ protected override void BuildModel(ModelBuilder modelBuilder) .HasColumnType("character varying(500)") .HasColumnName("id"); - b.Property("CorrelationId") - .HasColumnType("text") - .HasColumnName("correlation_id"); - b.Property("Created") .HasColumnType("timestamp with time zone") .HasColumnName("created"); @@ -57,6 +53,10 @@ protected override void BuildModel(ModelBuilder modelBuilder) .HasColumnType("text") .HasColumnName("hangfire_job_id"); + b.Property("InvocationCount") + .HasColumnType("integer") + .HasColumnName("invocation_count"); + b.Property("PagesCompleted") .HasColumnType("integer") .HasColumnName("pages_completed"); diff --git a/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs b/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs index 57c673f..d4e59c9 100644 --- a/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs +++ b/src/TextServices.Builder.Api/Services/Notifications/JobCompletionNotification.cs @@ -9,4 +9,4 @@ public record JobCompletionNotification( int TotalPages, int TotalWordCount, string? Errors, - string? CorrelationId); + int InvocationCount); diff --git a/src/TextServices.Tests/BuilderApi/JobResponseTests.cs b/src/TextServices.Tests/BuilderApi/JobResponseTests.cs index 3ad8b42..f4e8031 100644 --- a/src/TextServices.Tests/BuilderApi/JobResponseTests.cs +++ b/src/TextServices.Tests/BuilderApi/JobResponseTests.cs @@ -23,25 +23,25 @@ private static BuilderJob CompletedJob( }; // ------------------------------------------------------------------------- - // CorrelationId — echoed back unchanged + // InvocationCount — reflects the job's current run count // ------------------------------------------------------------------------- [Fact] - public void From_JobWithCorrelationId_EchoesCorrelationId() + public void From_NewJob_InvocationCountIsOne() { - var job = CompletedJob(); - job.CorrelationId = "caller-supplied-guid"; - - var response = JobResponse.From(job, Options()); - - response.CorrelationId.ShouldBe("caller-supplied-guid"); + var response = JobResponse.From(CompletedJob(), Options()); + response.InvocationCount.ShouldBe(1); } [Fact] - public void From_JobWithoutCorrelationId_CorrelationIdIsNull() + public void From_ReprocessedJob_InvocationCountReflectsJobValue() { - var response = JobResponse.From(CompletedJob(), Options()); - response.CorrelationId.ShouldBeNull(); + var job = CompletedJob(); + job.InvocationCount = 3; + + var response = JobResponse.From(job, Options()); + + response.InvocationCount.ShouldBe(3); } // ------------------------------------------------------------------------- diff --git a/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs b/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs index b2d678f..f8e0818 100644 --- a/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs +++ b/src/TextServices.Tests/BuilderApi/SnsJobNotifierTests.cs @@ -15,7 +15,7 @@ public sealed class SnsJobNotifierTests private static JobCompletionNotification MakeNotification( string jobId = "my/job", JobStatus status = JobStatus.Completed) => - new(jobId, status, DateTimeOffset.UtcNow, 10, 500, null, null); + new(jobId, status, DateTimeOffset.UtcNow, 10, 500, null, 1); // ------------------------------------------------------------------------- // No-op when TopicArn is absent diff --git a/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs b/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs index f355d7a..adb2dc7 100644 --- a/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs +++ b/src/TextServices.Tests/BuilderApi/TextBuildJobTests.cs @@ -599,7 +599,7 @@ public async Task ExecuteAsync_OnSuccess_NotifiesWithCompletedStatus() } [Fact] - public async Task ExecuteAsync_OnSuccess_NotificationEchoesCorrelationId() + public async Task ExecuteAsync_OnSuccess_NotificationCarriesInvocationCount() { var pages = new List { @@ -607,9 +607,9 @@ public async Task ExecuteAsync_OnSuccess_NotificationEchoesCorrelationId() TextUri = "https://example.org/alto/1.xml" }, }; - var job = await CreateJob("test/notify-correlation-id", + var job = await CreateJob("test/notify-invocation-count", sourceDataJson: JsonSerializer.Serialize(pages), - correlationId: "caller-supplied-guid"); + invocationCount: 2); var altoFetcher = FakeAlto(new Dictionary { @@ -621,7 +621,7 @@ public async Task ExecuteAsync_OnSuccess_NotificationEchoesCorrelationId() await sut.ExecuteAsync(job.Id, FakeCancellationToken.Instance); notifier.Captured.ShouldHaveSingleItem(); - notifier.Captured[0].CorrelationId.ShouldBe("caller-supplied-guid"); + notifier.Captured[0].InvocationCount.ShouldBe(2); } [Fact] @@ -648,7 +648,7 @@ public async Task ExecuteAsync_OnFailure_NotifiesWithFailedStatus() private async Task CreateJob(string id, string? sourceUri = null, string? sourceDataJson = null, - JobServices services = JobServices.All, string? correlationId = null) + JobServices services = JobServices.All, int invocationCount = 1) { var job = new BuilderJob { @@ -656,7 +656,7 @@ private async Task CreateJob(string id, SourceUri = sourceUri, SourceDataJson = sourceDataJson, Services = (int)services, - CorrelationId = correlationId, + InvocationCount = invocationCount, }; _db.Jobs.Add(job); await _db.SaveChangesAsync(); From be66adc047aac6f447846314c0dead75d0e06693 Mon Sep 17 00:00:00 2001 From: "jack.lewis" Date: Mon, 6 Jul 2026 12:37:51 +0100 Subject: [PATCH 3/4] Add invocations to demo project --- docker-compose.yml | 8 +------- src/TextServices.Demo/wwwroot/builder.html | 1 + src/TextServices.Demo/wwwroot/js/builder.js | 1 + 3 files changed, 3 insertions(+), 7 deletions(-) diff --git a/docker-compose.yml b/docker-compose.yml index 24dcd9a..89d9b70 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -8,11 +8,10 @@ services: image: postgres:14 hostname: postgres ports: - - "5453:5432" + - "5452:5432" volumes: - txt_postgres_data:/var/lib/postgresql/data - txt_postgres_data_backups:/backups - - $HOME\.aws:/root/.aws:ro environment: - POSTGRES_HOST=postgres - POSTGRES_PORT=5432 @@ -30,7 +29,6 @@ services: - postgres volumes: - txt_textservices_data:/data - - $HOME\.aws:/root/.aws:ro environment: - ASPNETCORE_ENVIRONMENT=Development - RunMigrations=true @@ -38,8 +36,6 @@ services: - TextServices__SearchApiBaseUrl=http://localhost:5294 - TextServices__Storage__FileSystem__RootPath=/data - CorsAllowedOrigins__0=http://localhost:5100 - env_file: - - .env search: build: @@ -55,8 +51,6 @@ services: - TextServices__BaseUrl=http://localhost:5294 - TextServices__Storage__FileSystem__RootPath=/data - CorsAllowedOrigins__0=http://localhost:5100 - env_file: - - .env demo: build: diff --git a/src/TextServices.Demo/wwwroot/builder.html b/src/TextServices.Demo/wwwroot/builder.html index 33cf804..b73f002 100644 --- a/src/TextServices.Demo/wwwroot/builder.html +++ b/src/TextServices.Demo/wwwroot/builder.html @@ -60,6 +60,7 @@

Jobs

ID Source Status + Invocations Progress Words Actions diff --git a/src/TextServices.Demo/wwwroot/js/builder.js b/src/TextServices.Demo/wwwroot/js/builder.js index ab7c8bc..e182609 100644 --- a/src/TextServices.Demo/wwwroot/js/builder.js +++ b/src/TextServices.Demo/wwwroot/js/builder.js @@ -195,6 +195,7 @@ function renderTable(rows) { ${esc(id)} ${esc(truncate(sourceUri, 50))} ${esc(status)} + ${job?.invocationCount ?? '—'} ${renderProgress(job)} ${job?.totalWordCount != null ? job.totalWordCount.toLocaleString() : '—'} From e8c2704307d58908bc55138759138ea7a712de82 Mon Sep 17 00:00:00 2001 From: "jack.lewis" Date: Mon, 6 Jul 2026 12:43:56 +0100 Subject: [PATCH 4/4] Revert docker compose --- docker-compose.yml | 1 - 1 file changed, 1 deletion(-) diff --git a/docker-compose.yml b/docker-compose.yml index 89d9b70..6eeee41 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -45,7 +45,6 @@ services: - "5294:8080" volumes: - txt_textservices_data:/data - - $HOME\.aws:/root/.aws:ro environment: - ASPNETCORE_ENVIRONMENT=Development - TextServices__BaseUrl=http://localhost:5294