diff --git a/src/Tiger.Tests/BuildIngestionPriorityTests.cs b/src/Tiger.Tests/BuildIngestionPriorityTests.cs
new file mode 100644
index 0000000..52a9636
--- /dev/null
+++ b/src/Tiger.Tests/BuildIngestionPriorityTests.cs
@@ -0,0 +1,273 @@
+using System.Net;
+using Xunit;
+
+namespace Tiger.Tests;
+
+///
+/// Tests for the priority preemption behavior in .
+///
+public class BuildIngestionPriorityTests : IDisposable
+{
+ private readonly string _dbPath;
+ private readonly TigerDatabase _db;
+
+ public BuildIngestionPriorityTests()
+ {
+ _dbPath = Path.Combine(Path.GetTempPath(), $"tiger-test-{Guid.NewGuid()}.db");
+ _db = TigerDatabase.Open(_dbPath);
+ }
+
+ public void Dispose()
+ {
+ _db.Dispose();
+ if (File.Exists(_dbPath))
+ File.Delete(_dbPath);
+ }
+
+ ///
+ /// Verifies that PrioritizeBuild preempts in-flight 'tests' tasks.
+ /// The preempted tasks must be reset to 'pending' in the DB.
+ ///
+ [Fact]
+ public async Task PrioritizeBuild_PreemptsTestsTasks()
+ {
+ var normalFetchStarted = new SemaphoreSlim(0);
+ var priorityTaskStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ const int priorityBuildId = 100;
+
+ var handler = new DelegateHandler(async (request, ct) =>
+ {
+ var url = request.RequestUri?.ToString() ?? "";
+ if (url.Contains($"Build%2F{priorityBuildId}") || url.Contains($"Build/{priorityBuildId}"))
+ {
+ priorityTaskStarted.TrySetResult();
+ return CreateJsonResponse("""{"count":0,"value":[]}""");
+ }
+ normalFetchStarted.Release();
+ await Task.Delay(Timeout.Infinite, ct);
+ return null!;
+ });
+
+ var factory = new AzdoClientFactory((org, proj) => AzdoClient.Create(handler, org, proj));
+ var service = new BuildIngestionService(_db, factory, maxParallelism: 2);
+
+ InsertBuildWithTask("org", "proj", buildId: 1, taskType: "tests");
+ InsertBuildWithTask("org", "proj", buildId: 2, taskType: "tests");
+ InsertBuildWithTask("org", "proj", buildId: priorityBuildId, taskType: "tests");
+
+ service.Start();
+ await normalFetchStarted.WaitAsync();
+ await normalFetchStarted.WaitAsync();
+
+ service.PrioritizeBuild("org", priorityBuildId);
+ await priorityTaskStarted.Task;
+ await service.StopAsync();
+
+ var status1 = GetTaskStatus("org", 1, "tests");
+ var status2 = GetTaskStatus("org", 2, "tests");
+ Assert.Equal("pending", status1);
+ Assert.Equal("pending", status2);
+ }
+
+ ///
+ /// Verifies that PrioritizeBuild preempts in-flight 'timeline' tasks.
+ /// The preempted tasks must be reset to 'pending' in the DB.
+ ///
+ [Fact]
+ public async Task PrioritizeBuild_PreemptsTimelineTasks()
+ {
+ var normalFetchStarted = new SemaphoreSlim(0);
+ var priorityTaskStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ const int priorityBuildId = 100;
+
+ var handler = new DelegateHandler(async (request, ct) =>
+ {
+ var url = request.RequestUri?.ToString() ?? "";
+ if (url.Contains($"builds/{priorityBuildId}/timeline"))
+ {
+ priorityTaskStarted.TrySetResult();
+ return CreateJsonResponse("""{"records":[]}""");
+ }
+ normalFetchStarted.Release();
+ await Task.Delay(Timeout.Infinite, ct);
+ return null!;
+ });
+
+ var factory = new AzdoClientFactory((org, proj) => AzdoClient.Create(handler, org, proj));
+ var service = new BuildIngestionService(_db, factory, maxParallelism: 2);
+
+ InsertBuildWithTask("org", "proj", buildId: 1, taskType: "timeline");
+ InsertBuildWithTask("org", "proj", buildId: 2, taskType: "timeline");
+ InsertBuildWithTask("org", "proj", buildId: priorityBuildId, taskType: "timeline");
+
+ service.Start();
+ await normalFetchStarted.WaitAsync();
+ await normalFetchStarted.WaitAsync();
+
+ service.PrioritizeBuild("org", priorityBuildId);
+ await priorityTaskStarted.Task;
+ await service.StopAsync();
+
+ var status1 = GetTaskStatus("org", 1, "timeline");
+ var status2 = GetTaskStatus("org", 2, "timeline");
+ Assert.Equal("pending", status1);
+ Assert.Equal("pending", status2);
+ }
+
+ ///
+ /// Verifies that PrioritizeBuild preempts in-flight 'tests' tasks that are
+ /// blocked fetching Helix work items. The Helix HTTP calls receive the
+ /// cancellation token and are interrupted properly.
+ ///
+ [Fact]
+ public async Task PrioritizeBuild_PreemptsHelixFetch()
+ {
+ var normalHelixStarted = new SemaphoreSlim(0);
+ var priorityTaskStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ const int priorityBuildId = 100;
+
+ // AzDO handler: returns test results with helix info for normal builds,
+ // returns empty for priority build
+ var azdoHandler = new DelegateHandler((request, ct) =>
+ {
+ var url = request.RequestUri?.ToString() ?? "";
+
+ if (url.Contains($"Build%2F{priorityBuildId}") || url.Contains($"Build/{priorityBuildId}"))
+ {
+ priorityTaskStarted.TrySetResult();
+ return Task.FromResult(CreateJsonResponse("""{"count":0,"value":[]}"""));
+ }
+
+ // Return test runs with a helix-linked failure so the service proceeds to fetch helix data
+ if (url.Contains("test/runs") && url.Contains("includeRunDetails"))
+ {
+ return Task.FromResult(CreateJsonResponse(
+ """{"count":1,"value":[{"id":1,"name":"TestRun","totalTests":1,"passedTests":0,"notApplicableTests":0,"unanalyzedTests":0}]}"""));
+ }
+
+ if (url.Contains("test/Runs/") && url.Contains("results"))
+ {
+ return Task.FromResult(CreateJsonResponse("""
+ {"count":1,"value":[{
+ "id":1,
+ "testRun":{"id":"1","name":"TestRun"},
+ "testCaseTitle":"SomeTest",
+ "automatedTestName":"Namespace.SomeTest",
+ "outcome":"Failed",
+ "comment":"{\"HelixJobId\":\"test-job\",\"HelixWorkItemName\":\"test-workitem\"}"
+ }]}
+ """));
+ }
+
+ if (url.Contains("test/runs"))
+ {
+ return Task.FromResult(CreateJsonResponse(
+ """{"count":1,"value":[{"id":1,"name":"TestRun","totalTests":1,"passedTests":0,"notApplicableTests":0,"unanalyzedTests":0}]}"""));
+ }
+
+ return Task.FromResult(CreateJsonResponse("""{"count":0,"value":[]}"""));
+ });
+
+ // Helix handler: blocks on work item fetch, signals when started
+ var helixHandler = new DelegateHandler(async (request, ct) =>
+ {
+ normalHelixStarted.Release();
+ await Task.Delay(Timeout.Infinite, ct);
+ return null!;
+ });
+
+ var factory = new AzdoClientFactory((org, proj) => AzdoClient.Create(azdoHandler, org, proj));
+ Func helixFactory = () => HelixClient.Create(helixHandler);
+ var service = new BuildIngestionService(_db, factory, maxParallelism: 2, helixClientFactory: helixFactory);
+
+ InsertBuildWithTask("org", "proj", buildId: 1, taskType: "tests");
+ InsertBuildWithTask("org", "proj", buildId: 2, taskType: "tests");
+ InsertBuildWithTask("org", "proj", buildId: priorityBuildId, taskType: "tests");
+
+ service.Start();
+ await normalHelixStarted.WaitAsync();
+ await normalHelixStarted.WaitAsync();
+
+ service.PrioritizeBuild("org", priorityBuildId);
+ await priorityTaskStarted.Task;
+ await service.StopAsync();
+
+ var status1 = GetTaskStatus("org", 1, "tests");
+ var status2 = GetTaskStatus("org", 2, "tests");
+ Assert.Equal("pending", status1);
+ Assert.Equal("pending", status2);
+ }
+
+ // ── Helpers ─────────────────────────────────────────────────────
+
+ private void InsertBuildWithTask(string organization, string project, int buildId, string taskType)
+ {
+ _db.WithCommand(cmd =>
+ {
+ cmd.CommandText = """
+ INSERT OR IGNORE INTO builds
+ (organization, project, build_id, build_number, definition_name, definition_id,
+ status, result, source_branch)
+ VALUES
+ (@org, @proj, @buildId, @buildNumber, 'test-def', 1,
+ 'completed', 'failed', 'refs/heads/main')
+ """;
+ cmd.Parameters.AddWithValue("@org", organization);
+ cmd.Parameters.AddWithValue("@proj", project);
+ cmd.Parameters.AddWithValue("@buildId", buildId);
+ cmd.Parameters.AddWithValue("@buildNumber", $"20250101.{buildId}");
+ cmd.ExecuteNonQuery();
+ });
+
+ _db.WithCommand(cmd =>
+ {
+ cmd.CommandText = """
+ INSERT OR IGNORE INTO build_ingestion_tasks
+ (organization, build_id, task_type, status, is_complete, attempts)
+ VALUES
+ (@org, @buildId, @taskType, 'pending', 0, 0)
+ """;
+ cmd.Parameters.AddWithValue("@org", organization);
+ cmd.Parameters.AddWithValue("@buildId", buildId);
+ cmd.Parameters.AddWithValue("@taskType", taskType);
+ cmd.ExecuteNonQuery();
+ });
+ }
+
+ private string GetTaskStatus(string organization, int buildId, string taskType)
+ {
+ return _db.WithCommand(cmd =>
+ {
+ cmd.CommandText = """
+ SELECT status FROM build_ingestion_tasks
+ WHERE organization = @org AND build_id = @buildId AND task_type = @type
+ """;
+ cmd.Parameters.AddWithValue("@org", organization);
+ cmd.Parameters.AddWithValue("@buildId", buildId);
+ cmd.Parameters.AddWithValue("@type", taskType);
+ return cmd.ExecuteScalar() as string ?? throw new InvalidOperationException(
+ $"No task found for org={organization}, buildId={buildId}, taskType={taskType}");
+ });
+ }
+
+ // ── Shared Infrastructure ───────────────────────────────────────
+
+ ///
+ /// An HttpMessageHandler that delegates to a func, avoiding the need for
+ /// a separate subclass per test scenario.
+ ///
+ private sealed class DelegateHandler(
+ Func> handler) : HttpMessageHandler
+ {
+ protected override Task SendAsync(
+ HttpRequestMessage request, CancellationToken ct) => handler(request, ct);
+ }
+
+ private static HttpResponseMessage CreateJsonResponse(string json)
+ {
+ return new HttpResponseMessage(HttpStatusCode.OK)
+ {
+ Content = new StringContent(json, System.Text.Encoding.UTF8, "application/json"),
+ };
+ }
+}
diff --git a/src/Tiger/AzdoClient.cs b/src/Tiger/AzdoClient.cs
index d3f4b75..d006b2a 100644
--- a/src/Tiger/AzdoClient.cs
+++ b/src/Tiger/AzdoClient.cs
@@ -58,6 +58,23 @@ internal static AzdoClient Create(
return new AzdoClient(httpClient, organization, project, rateLimitState);
}
+ ///
+ /// Creates an with a custom .
+ /// Useful for testing — inject a handler that intercepts or blocks HTTP calls.
+ ///
+ public static AzdoClient Create(
+ HttpMessageHandler handler,
+ string organization = DefaultOrganization,
+ string project = DefaultProject)
+ {
+ var httpClient = new HttpClient(handler)
+ {
+ BaseAddress = new Uri($"https://dev.azure.com/{organization}/{project}/"),
+ };
+
+ return new AzdoClient(httpClient, organization, project, new AzdoRateLimitState());
+ }
+
///
public static Task CreateAsync(
TokenCredential tokenCredential,
@@ -96,7 +113,7 @@ private string GetBuildUri(int buildId) =>
return null;
}
- public async Task> GetRecentBuildsAsync(int? definitionId = null, int top = 10)
+ public async Task> GetRecentBuildsAsync(int? definitionId = null, int top = 10, CancellationToken ct = default)
{
var url = $"_apis/build/builds?api-version=7.1&$top={top}";
if (definitionId is not null)
@@ -104,10 +121,10 @@ public async Task> GetRecentBuildsAsync(int? definitionId = null
url += $"&definitions={definitionId}";
}
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var result = JsonSerializer.Deserialize>(json, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize builds response");
@@ -158,7 +175,7 @@ public async Task> GetCompletedBuildsSinceAsync(
return all;
}
- public async Task> GetBuildsForRepositoryAsync(string repository, int top = 10, string? reasonFilter = null)
+ public async Task> GetBuildsForRepositoryAsync(string repository, int top = 10, string? reasonFilter = null, CancellationToken ct = default)
{
var url = $"_apis/build/builds?api-version=7.1&$top={top}&repositoryId={Uri.EscapeDataString(repository)}&repositoryType=GitHub";
if (reasonFilter is not null)
@@ -166,25 +183,25 @@ public async Task> GetBuildsForRepositoryAsync(string repository
url += $"&reasonFilter={Uri.EscapeDataString(reasonFilter)}";
}
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var result = JsonSerializer.Deserialize>(json, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize builds response");
return result.Value.Select(MapBuild).ToList();
}
- public async Task> GetBuildsForPullRequestAsync(string repository, int prNumber, int top = 10)
+ public async Task> GetBuildsForPullRequestAsync(string repository, int prNumber, int top = 10, CancellationToken ct = default)
{
var branchName = $"refs/pull/{prNumber}/merge";
var url = $"_apis/build/builds?api-version=7.1&$top={top}&branchName={Uri.EscapeDataString(branchName)}&repositoryId={Uri.EscapeDataString(repository)}&repositoryType=GitHub";
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var result = JsonSerializer.Deserialize>(json, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize builds response");
@@ -197,15 +214,15 @@ public async Task> GetBuildsForPullRequestAsync(string repositor
/// The build ID.
/// Controls sub-result fetching for grouped results (e.g. xUnit theories):
/// null = do not fetch sub-results, positive = fetch up to that many, -1 = no limit.
- public async Task> GetTestFailuresAsync(int buildId, int? subResultCount = null)
+ public async Task> GetTestFailuresAsync(int buildId, int? subResultCount = null, CancellationToken ct = default)
{
var buildUri = $"vstfs:///Build/Build/{buildId}";
var runsUrl = $"_apis/test/runs?api-version=7.1&buildUri={Uri.EscapeDataString(buildUri)}";
- var runsResponse = await HttpClient.GetAsync(runsUrl);
+ var runsResponse = await HttpClient.GetAsync(runsUrl, ct);
runsResponse.EnsureSuccessStatusCode();
- var runsJson = await runsResponse.Content.ReadAsStringAsync();
+ var runsJson = await runsResponse.Content.ReadAsStringAsync(ct);
var runs = JsonSerializer.Deserialize>(runsJson, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize test runs response");
@@ -215,10 +232,10 @@ public async Task> GetTestFailuresAsync(int buildId, int? s
foreach (var run in runs.Value)
{
var resultsUrl = $"_apis/test/Runs/{run.Id}/results?api-version=7.1&outcomes=Failed";
- var resultsResponse = await HttpClient.GetAsync(resultsUrl);
+ var resultsResponse = await HttpClient.GetAsync(resultsUrl, ct);
resultsResponse.EnsureSuccessStatusCode();
- var resultsJson = await resultsResponse.Content.ReadAsStringAsync();
+ var resultsJson = await resultsResponse.Content.ReadAsStringAsync(ct);
var results = JsonSerializer.Deserialize>(resultsJson, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize test results response");
@@ -229,7 +246,7 @@ public async Task> GetTestFailuresAsync(int buildId, int? s
if (result.HasSubResults && canFetchSubs)
{
- var detailed = await GetTestResultWithSubResultsAsync(run.Id, result.Id);
+ var detailed = await GetTestResultWithSubResultsAsync(run.Id, result.Id, ct);
if (detailed is not null && detailed.SubResults.Count > 0)
{
var failedSubs = detailed.SubResults
@@ -256,38 +273,38 @@ public async Task> GetTestFailuresAsync(int buildId, int? s
///
/// Fetches a single test result with sub-results included.
///
- private async Task GetTestResultWithSubResultsAsync(int runId, int resultId)
+ private async Task GetTestResultWithSubResultsAsync(int runId, int resultId, CancellationToken ct = default)
{
var url = $"_apis/test/Runs/{runId}/results/{resultId}?api-version=7.1&detailsToInclude=SubResults";
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
return JsonSerializer.Deserialize(json, s_jsonOptions);
}
- public async Task> GetTestResultAttachmentsAsync(int runId, int testCaseResultId)
+ public async Task> GetTestResultAttachmentsAsync(int runId, int testCaseResultId, CancellationToken ct = default)
{
var url = $"_apis/test/Runs/{runId}/Results/{testCaseResultId}/attachments?api-version=7.2-preview.1";
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var result = JsonSerializer.Deserialize>(json, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize test result attachments response");
return result.Value;
}
- public async Task> GetTestSummaryByJobAsync(int buildId)
+ public async Task> GetTestSummaryByJobAsync(int buildId, CancellationToken ct = default)
{
var buildUri = $"vstfs:///Build/Build/{buildId}";
var runsUrl = $"_apis/test/runs?api-version=7.1&includeRunDetails=true&buildUri={Uri.EscapeDataString(buildUri)}";
- var response = await HttpClient.GetAsync(runsUrl);
+ var response = await HttpClient.GetAsync(runsUrl, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var runs = JsonSerializer.Deserialize>(json, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize test runs response");
@@ -305,13 +322,13 @@ public async Task> GetTestSummaryByJobAsync(int buildId
}).ToList();
}
- public async Task GetTimelineAsync(int buildId)
+ public async Task GetTimelineAsync(int buildId, CancellationToken ct = default)
{
var url = $"_apis/build/builds/{buildId}/timeline?api-version=7.1";
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var raw = JsonSerializer.Deserialize(json, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize timeline response");
@@ -342,13 +359,13 @@ public async Task GetTimelineAsync(int buildId)
};
}
- public async Task> GetArtifactsAsync(int buildId)
+ public async Task> GetArtifactsAsync(int buildId, CancellationToken ct = default)
{
var url = $"_apis/build/builds/{buildId}/artifacts?api-version=7.1";
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var raw = JsonSerializer.Deserialize>(json, s_jsonOptions)
?? throw new InvalidOperationException("Failed to deserialize artifacts response");
@@ -363,34 +380,34 @@ public async Task> GetArtifactsAsync(int buildId)
}).ToList();
}
- public async Task DownloadArtifactAsync(int buildId, string artifactName, string outputPath)
+ public async Task DownloadArtifactAsync(int buildId, string artifactName, string outputPath, CancellationToken ct = default)
{
- var artifacts = await GetArtifactsAsync(buildId);
+ var artifacts = await GetArtifactsAsync(buildId, ct);
var artifact = artifacts.FirstOrDefault(a => a.Name == artifactName)
?? throw new InvalidOperationException($"Artifact '{artifactName}' not found for build {buildId}");
var downloadUrl = artifact.DownloadUrl
?? throw new InvalidOperationException($"Artifact '{artifactName}' has no download URL");
- using var response = await HttpClient.GetAsync(downloadUrl, HttpCompletionOption.ResponseHeadersRead);
+ using var response = await HttpClient.GetAsync(downloadUrl, HttpCompletionOption.ResponseHeadersRead, ct);
response.EnsureSuccessStatusCode();
using var fileStream = File.Create(outputPath);
- await response.Content.CopyToAsync(fileStream);
+ await response.Content.CopyToAsync(fileStream, ct);
}
///
/// Lists files within an artifact. Supports both Container and PipelineArtifact types.
///
- public async Task> GetArtifactFilesAsync(int buildId, AzdoArtifact artifact)
+ public async Task> GetArtifactFilesAsync(int buildId, AzdoArtifact artifact, CancellationToken ct = default)
{
if (artifact.ResourceType == "Container")
{
- return await GetContainerFilesAsync(artifact);
+ return await GetContainerFilesAsync(artifact, ct);
}
else if (artifact.ResourceType == "PipelineArtifact")
{
- return await GetPipelineArtifactFilesAsync(buildId, artifact);
+ return await GetPipelineArtifactFilesAsync(buildId, artifact, ct);
}
else
{
@@ -401,7 +418,7 @@ public async Task> GetArtifactFilesAsync(int buildId, Az
///
/// Downloads a single file from an artifact to the specified output path.
///
- public async Task DownloadArtifactFileAsync(int buildId, AzdoArtifact artifact, ArtifactFileEntry file, string outputPath)
+ public async Task DownloadArtifactFileAsync(int buildId, AzdoArtifact artifact, ArtifactFileEntry file, string outputPath, CancellationToken ct = default)
{
string downloadUrl;
if (artifact.ResourceType == "Container")
@@ -415,7 +432,7 @@ public async Task DownloadArtifactFileAsync(int buildId, AzdoArtifact artifact,
downloadUrl = $"_apis/build/builds/{buildId}/artifacts?artifactName={Uri.EscapeDataString(artifact.Name)}&fileId={file.BlobId}&fileName={Uri.EscapeDataString(Path.GetFileName(file.Path))}&api-version=7.1";
}
- using var response = await HttpClient.GetAsync(downloadUrl, HttpCompletionOption.ResponseHeadersRead);
+ using var response = await HttpClient.GetAsync(downloadUrl, HttpCompletionOption.ResponseHeadersRead, ct);
response.EnsureSuccessStatusCode();
var dir = Path.GetDirectoryName(outputPath);
@@ -425,10 +442,10 @@ public async Task DownloadArtifactFileAsync(int buildId, AzdoArtifact artifact,
}
using var fileStream = File.Create(outputPath);
- await response.Content.CopyToAsync(fileStream);
+ await response.Content.CopyToAsync(fileStream, ct);
}
- private async Task> GetContainerFilesAsync(AzdoArtifact artifact)
+ private async Task> GetContainerFilesAsync(AzdoArtifact artifact, CancellationToken ct = default)
{
// resource.data is "#/{containerId}/{artifactName}"
var data = artifact.ResourceData
@@ -441,10 +458,10 @@ private async Task> GetContainerFilesAsync(AzdoArtifact
// The resource URL points to the container. Append $format=json if not present
var listUrl = url.Contains("?") ? $"{url}&$format=json" : $"{url}?$format=json";
- using var response = await HttpClient.GetAsync(listUrl);
+ using var response = await HttpClient.GetAsync(listUrl, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var container = JsonSerializer.Deserialize(json, s_jsonOptions);
if (container?.Value is null)
{
@@ -461,17 +478,17 @@ private async Task> GetContainerFilesAsync(AzdoArtifact
}).ToList();
}
- private async Task> GetPipelineArtifactFilesAsync(int buildId, AzdoArtifact artifact)
+ private async Task> GetPipelineArtifactFilesAsync(int buildId, AzdoArtifact artifact, CancellationToken ct = default)
{
var fileId = artifact.ResourceData
?? throw new InvalidOperationException($"Artifact '{artifact.Name}' has no resource data (file ID)");
var url = $"_apis/build/builds/{buildId}/artifacts?artifactName={Uri.EscapeDataString(artifact.Name)}&fileId={fileId}&fileName=manifest.json&api-version=7.1";
- using var response = await HttpClient.GetAsync(url);
+ using var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
var manifest = JsonSerializer.Deserialize(json, s_jsonOptions);
if (manifest?.Items is null)
{
diff --git a/src/Tiger/BuildIngestionService.cs b/src/Tiger/BuildIngestionService.cs
index 06034fa..9c45048 100644
--- a/src/Tiger/BuildIngestionService.cs
+++ b/src/Tiger/BuildIngestionService.cs
@@ -27,8 +27,10 @@ public sealed class BuildIngestionService : IDisposable
{
private readonly TigerDatabase _db;
private readonly AzdoClientFactory _clientFactory;
+ private readonly Func _helixClientFactory;
private readonly ServiceLog? _log;
private readonly ConcurrentStack _priorityTasks = new();
+ private readonly SemaphoreSlim _prioritySignal = new(0);
private CancellationTokenSource? _cts;
private Task? _workerTask;
@@ -36,7 +38,9 @@ public sealed class BuildIngestionService : IDisposable
private const int WorkerIntervalSeconds = 5;
private const int CircuitBreakerThreshold = 5;
private const int CircuitBreakerCooldownSeconds = 120;
- private const int MaxParallelism = 8;
+ private const int DefaultMaxParallelism = 8;
+
+ private readonly int _maxParallelism;
///
/// Backoff delays per attempt: 30s, 2min, 10min, 1hr, then abandon.
@@ -50,11 +54,13 @@ public sealed class BuildIngestionService : IDisposable
public bool IsRunning => _workerTask is not null && !_workerTask.IsCompleted;
- public BuildIngestionService(TigerDatabase db, AzdoClientFactory clientFactory, ServiceLog? log = null)
+ public BuildIngestionService(TigerDatabase db, AzdoClientFactory clientFactory, ServiceLog? log = null, int maxParallelism = DefaultMaxParallelism, Func? helixClientFactory = null)
{
_db = db;
_clientFactory = clientFactory;
_log = log;
+ _maxParallelism = maxParallelism;
+ _helixClientFactory = helixClientFactory ?? (() => HelixClient.Create());
}
///
@@ -252,11 +258,14 @@ INSERT OR IGNORE INTO test_results
cmd.ExecuteNonQuery();
}
- internal void InsertTimelineIssues(string organization, string project, int buildId, AzdoTimeline timeline)
+ internal void InsertTimelineIssues(string organization, string project, int buildId, AzdoTimeline timeline) =>
+ InsertTimelineIssues(_db, organization, project, buildId, timeline);
+
+ internal static void InsertTimelineIssues(TigerDatabase db, string organization, string project, int buildId, AzdoTimeline timeline)
{
var recordNames = timeline.Records.ToDictionary(r => r.Id, r => r.Name);
- _db.WithTransaction((conn, tx) =>
+ db.WithTransaction((conn, tx) =>
{
using var cmd = conn.CreateCommand();
cmd.Transaction = tx;
@@ -340,6 +349,7 @@ public async Task StopAsync()
///
/// Pushes all non-complete ingestion tasks for the specified build onto the
/// priority stack so they are picked up before the normal DB queue.
+ /// Signals the worker loop to preempt non-priority in-flight work if needed.
///
public void PrioritizeBuild(string organization, int buildId)
{
@@ -374,18 +384,38 @@ FROM build_ingestion_tasks t
{
_priorityTasks.Push(task);
}
+
+ // Wake the worker loop so it can preempt non-priority work
+ if (tasks.Count > 0)
+ {
+ _prioritySignal.Release();
+ }
}
///
- /// Maintains up to in-flight tasks at all times,
+ /// Tracks an in-flight task along with its per-task cancellation token, allowing
+ /// individual tasks to be cancelled for priority preemption without affecting
+ /// the overall worker loop.
+ ///
+ private sealed class InFlightTask
+ {
+ public required Task Task { get; init; }
+ public required CancellationTokenSource Cts { get; init; }
+ public required bool IsPriority { get; init; }
+ }
+
+ ///
+ /// Maintains up to in-flight tasks at all times,
/// adapting concurrency based on AzDO rate-limit headers.
/// When any task completes, its result is handled immediately and a new task
/// is claimed to fill the slot — no waiting for an entire batch to drain.
+ /// When priority work arrives, non-priority in-flight tasks are cancelled to
+ /// make room immediately.
///
private async Task WorkLoopAsync(CancellationToken ct)
{
var consecutiveFailures = 0;
- var inFlight = new List>(MaxParallelism);
+ var inFlight = new List(_maxParallelism);
while (!ct.IsCancellationRequested)
{
@@ -400,8 +430,10 @@ private async Task WorkLoopAsync(CancellationToken ct)
// Drain in-flight work before cooling down
while (inFlight.Count > 0)
{
- var done = await Task.WhenAny(inFlight);
- inFlight.Remove(done);
+ var done = await Task.WhenAny(inFlight.Select(f => f.Task));
+ var completed = inFlight.First(f => f.Task == done);
+ inFlight.Remove(completed);
+ completed.Cts.Dispose();
HandleCompletion(done);
}
@@ -410,29 +442,75 @@ private async Task WorkLoopAsync(CancellationToken ct)
continue;
}
+ // Check for priority preemption: if priority work is pending and all
+ // slots are full, cancel non-priority tasks to make room.
+ if (!_priorityTasks.IsEmpty && inFlight.Count >= GetEffectiveParallelism())
+ {
+ PreemptNonPriorityTasks(inFlight);
+
+ // Wait for cancelled tasks to complete and remove them from in-flight
+ var cancelled = inFlight.Where(f => f.Cts.IsCancellationRequested).ToList();
+ foreach (var flight in cancelled)
+ {
+ try
+ {
+ await flight.Task;
+ }
+ catch
+ {
+ }
+
+ inFlight.Remove(flight);
+ flight.Cts.Dispose();
+ // Don't count preempted tasks as failures
+ }
+ }
+
// Fill available slots up to effective parallelism
var effectiveParallelism = GetEffectiveParallelism();
while (inFlight.Count < effectiveParallelism)
{
- var task = GetNextReadyTask();
- if (task is null)
+ var next = GetNextReadyTask();
+ if (next is null)
{
break;
}
- inFlight.Add(RunTaskAsync(task, ct));
+ var (task, isPriority) = next.Value;
+ var taskCts = CancellationTokenSource.CreateLinkedTokenSource(ct);
+ var flight = new InFlightTask
+ {
+ Task = RunTaskAsync(task, taskCts.Token),
+ Cts = taskCts,
+ IsPriority = isPriority,
+ };
+ inFlight.Add(flight);
}
if (inFlight.Count == 0)
{
- await Task.Delay(TimeSpan.FromSeconds(WorkerIntervalSeconds), ct);
+ // Wait for either the poll interval or a priority signal
+ await WaitForWorkOrSignal(ct);
continue;
}
- // Wait for any one task to complete, then loop to refill the slot
- var completed = await Task.WhenAny(inFlight);
- inFlight.Remove(completed);
- HandleCompletion(completed);
+ // Wait for any one task to complete OR a priority signal to arrive
+ var signalTask = _prioritySignal.WaitAsync(ct);
+ var tasksToAwait = new List(inFlight.Count + 1);
+ tasksToAwait.AddRange(inFlight.Select(f => (Task)f.Task));
+ tasksToAwait.Add(signalTask);
+ var completedTask = await Task.WhenAny(tasksToAwait);
+
+ // If the priority signal fired, loop back to preempt
+ if (completedTask == signalTask)
+ {
+ continue;
+ }
+
+ var completedFlight = inFlight.First(f => f.Task == completedTask);
+ inFlight.Remove(completedFlight);
+ completedFlight.Cts.Dispose();
+ HandleCompletion(completedFlight.Task);
}
catch (OperationCanceledException)
{
@@ -446,20 +524,25 @@ private async Task WorkLoopAsync(CancellationToken ct)
}
// Drain remaining in-flight tasks on shutdown
- foreach (var remaining in inFlight)
+ foreach (var flight in inFlight)
{
try
{
- await remaining;
+ await flight.Cts.CancelAsync();
+ await flight.Task;
}
catch
{
}
+ finally
+ {
+ flight.Cts.Dispose();
+ }
}
void HandleCompletion(Task task)
{
- if (task.IsFaulted || task.Result is null)
+ if (task.IsFaulted || task.IsCanceled || task.Result is null)
{
consecutiveFailures++;
}
@@ -470,6 +553,34 @@ void HandleCompletion(Task task)
}
}
+ ///
+ /// Cancels non-priority in-flight tasks to make room for priority work.
+ /// Cancelled tasks are reset to 'pending' in the DB by .
+ ///
+ private void PreemptNonPriorityTasks(List inFlight)
+ {
+ var nonPriority = inFlight.Where(f => !f.IsPriority).ToList();
+ foreach (var flight in nonPriority)
+ {
+ _log?.Info("Worker", "Preempting non-priority task for priority work");
+ flight.Cts.Cancel();
+ }
+ }
+
+ ///
+ /// Waits for either the standard poll interval to elapse or a priority signal
+ /// to arrive, whichever comes first.
+ ///
+ private async Task WaitForWorkOrSignal(CancellationToken ct)
+ {
+ using var delayCts = CancellationTokenSource.CreateLinkedTokenSource(ct);
+ var delayTask = Task.Delay(TimeSpan.FromSeconds(WorkerIntervalSeconds), delayCts.Token);
+ var signalTask = _prioritySignal.WaitAsync(ct);
+
+ await Task.WhenAny(delayTask, signalTask);
+ await delayCts.CancelAsync();
+ }
+
///
/// Computes how many tasks to run concurrently based on AzDO rate-limit state.
/// Checks all known organizations and uses the most constrained one.
@@ -480,8 +591,8 @@ private int GetEffectiveParallelism()
return fraction switch
{
- > 0.5 => MaxParallelism,
- > 0.25 => MaxParallelism / 2,
+ > 0.5 => _maxParallelism,
+ > 0.25 => _maxParallelism / 2,
> 0.1 => 1,
_ => 0,
};
@@ -512,7 +623,8 @@ private double GetMinRemainingFraction()
///
/// Processes a single ingestion task and returns it for post-completion handling.
- /// Exceptions are captured so the caller can inspect them.
+ /// When cancelled (e.g. by priority preemption), resets the task to 'pending'
+ /// so it can be retried later. This does NOT increment the attempt counter.
///
private async Task RunTaskAsync(IngestionTask task, CancellationToken ct)
{
@@ -536,20 +648,25 @@ SELECT 1 FROM build_ingestion_tasks
try
{
- MarkRunning();
+ MarkRunning(_db, task);
await ProcessTaskAsync(task, ct);
- MarkComplete();
+ MarkComplete(_db, task);
}
catch (OperationCanceledException)
{
- throw;
+ // Preemption or shutdown — reset to pending so the task can be retried.
+ // No attempt increment, no backoff.
+ MarkPreempted(_db, task);
+ _log?.Info("Worker",
+ $"Task {task.TaskType} for build #{task.BuildId} preempted, reset to pending");
+ return null;
}
catch (Exception ex)
{
var newAttempts = task.Attempts + 1;
if (newAttempts >= MaxAttempts)
{
- MarkAbandoned(ex.Message);
+ MarkAbandoned(_db, task, ex.Message);
_log?.Error("Worker",
$"Task {task.TaskType} for build #{task.BuildId} abandoned after {newAttempts} attempts: {ex.Message}");
}
@@ -557,7 +674,7 @@ SELECT 1 FROM build_ingestion_tasks
{
var backoffIndex = Math.Min(newAttempts - 1, s_backoffSeconds.Length - 1);
var delaySecs = s_backoffSeconds[backoffIndex];
- MarkFailed(ex.Message, delaySecs);
+ MarkFailed(_db, task, ex.Message, delaySecs);
_log?.Warning("Worker",
$"Task {task.TaskType} for build #{task.BuildId} failed (attempt {newAttempts}), retry in {delaySecs}s: {ex.Message}");
}
@@ -571,9 +688,12 @@ SELECT 1 FROM build_ingestion_tasks
return task; // signals success
- void MarkRunning()
+ // All Mark* functions are static to guarantee they cannot capture the
+ // CancellationToken. DB writes must never be interrupted by cancellation.
+
+ static void MarkRunning(TigerDatabase db, IngestionTask task)
{
- _db.WithCommand(cmd =>
+ db.WithCommand(cmd =>
{
cmd.CommandText = """
UPDATE build_ingestion_tasks
@@ -587,9 +707,9 @@ UPDATE build_ingestion_tasks
});
}
- void MarkComplete()
+ static void MarkComplete(TigerDatabase db, IngestionTask task)
{
- _db.WithCommand(cmd =>
+ db.WithCommand(cmd =>
{
cmd.CommandText = """
UPDATE build_ingestion_tasks
@@ -603,9 +723,9 @@ UPDATE build_ingestion_tasks
});
}
- void MarkFailed(string error, int retryDelaySecs)
+ static void MarkFailed(TigerDatabase db, IngestionTask task, string error, int retryDelaySecs)
{
- _db.WithCommand(cmd =>
+ db.WithCommand(cmd =>
{
cmd.CommandText = $"""
UPDATE build_ingestion_tasks
@@ -624,9 +744,9 @@ UPDATE build_ingestion_tasks
});
}
- void MarkAbandoned(string error)
+ static void MarkAbandoned(TigerDatabase db, IngestionTask task, string error)
{
- _db.WithCommand(cmd =>
+ db.WithCommand(cmd =>
{
cmd.CommandText = """
UPDATE build_ingestion_tasks
@@ -644,6 +764,22 @@ UPDATE build_ingestion_tasks
cmd.ExecuteNonQuery();
});
}
+
+ static void MarkPreempted(TigerDatabase db, IngestionTask task)
+ {
+ db.WithCommand(cmd =>
+ {
+ cmd.CommandText = """
+ UPDATE build_ingestion_tasks
+ SET status = 'pending', next_retry_time = NULL
+ WHERE organization = @org AND build_id = @buildId AND task_type = @type
+ """;
+ cmd.Parameters.AddWithValue("@org", task.Organization);
+ cmd.Parameters.AddWithValue("@buildId", task.BuildId);
+ cmd.Parameters.AddWithValue("@type", task.TaskType);
+ cmd.ExecuteNonQuery();
+ });
+ }
}
private async Task ProcessTaskAsync(IngestionTask task, CancellationToken ct)
@@ -653,263 +789,262 @@ private async Task ProcessTaskAsync(IngestionTask task, CancellationToken ct)
switch (task.TaskType)
{
case "tests":
- await ProcessTestsAsync(client, task, ct);
+ await ProcessTests();
break;
case "timeline":
- await ProcessTimelineAsync(client, task, ct);
+ await ProcessTimeline();
break;
case "pr_info":
- await ProcessPrInfoAsync(task, ct);
+ await ProcessPrInfo();
break;
default:
_log?.Warning("Worker", $"Unknown task type: {task.TaskType}");
break;
}
- }
- private async Task ProcessTestsAsync(AzdoClient client, IngestionTask task, CancellationToken ct)
- {
- _log?.Info("Worker", $"Fetching tests for build #{task.BuildId}...");
+ return;
- var (testSummary, failures) = await FetchTestResults();
- var helixWorkItems = await FetchHelixWorkItems();
+ // ── Tests ───────────────────────────────────────────────────
- // Both fetches succeeded — persist everything in one transaction
- _db.WithTransaction((conn, tx) =>
+ async Task ProcessTests()
{
- InsertTestData(conn, tx);
- InsertHelixData(conn, tx);
- });
+ _log?.Info("Worker", $"Fetching tests for build #{task.BuildId}...");
+ var data = await FetchTestsDataAsync(client, _helixClientFactory, _db, _log, task, ct);
+ InsertTestsData(_db, task, data);
- if (failures.Count > 0)
- {
- _log?.Info("Worker",
- $" Build #{task.BuildId} — {failures.Count} test failure(s) across {failures.GroupBy(f => f.TestRunId).Count()} run(s)");
- }
- else
- {
- _log?.Info("Worker", $" Build #{task.BuildId} — tests complete (no failures)");
- }
+ if (data.Failures.Count > 0)
+ {
+ _log?.Info("Worker",
+ $" Build #{task.BuildId} — {data.Failures.Count} test failure(s) across {data.Failures.GroupBy(f => f.TestRunId).Count()} run(s)");
+ }
+ else
+ {
+ _log?.Info("Worker", $" Build #{task.BuildId} — tests complete (no failures)");
+ }
- if (helixWorkItems.Count > 0)
- {
- _log?.Info("Worker", $" Build #{task.BuildId} — helix complete ({helixWorkItems.Count} work item(s) fetched)");
+ if (data.HelixWorkItems.Count > 0)
+ {
+ _log?.Info("Worker", $" Build #{task.BuildId} — helix complete ({data.HelixWorkItems.Count} work item(s) fetched)");
+ }
}
- return;
-
- async Task<(List Summary, List Failures)> FetchTestResults()
+ static async Task FetchTestsDataAsync(
+ AzdoClient client, Func helixClientFactory, TigerDatabase db,
+ ServiceLog? log, IngestionTask task, CancellationToken ct)
{
- var summary = await client.GetTestSummaryByJobAsync(task.BuildId);
- var results = await client.GetTestFailuresAsync(task.BuildId, subResultCount: 50);
- return (summary, results);
- }
+ var summary = await client.GetTestSummaryByJobAsync(task.BuildId, ct);
+ var failures = await client.GetTestFailuresAsync(task.BuildId, subResultCount: 50, ct: ct);
- async Task> FetchHelixWorkItems()
- {
- // Extract helix job/work-item pairs directly from the fetched test results
+ // Fetch helix work items for any failures that reference them
var workItemKeys = failures
.Where(f => f.HelixJobName is not null && f.HelixWorkItemName is not null)
.Select(f => (f.HelixJobName!, f.HelixWorkItemName!))
.Distinct()
.ToList();
- if (workItemKeys.Count == 0)
+ var helixWorkItems = new List();
+ if (workItemKeys.Count > 0)
{
- return [];
- }
+ log?.Info("Worker", $"Fetching helix work items for build #{task.BuildId}...");
+ var helixClient = helixClientFactory();
- _log?.Info("Worker", $"Fetching helix work items for build #{task.BuildId}...");
-
- var helixClient = HelixClient.Create();
- var items = new List();
-
- foreach (var (jobName, workItemName) in workItemKeys)
- {
- if (ct.IsCancellationRequested)
+ foreach (var (jobName, workItemName) in workItemKeys)
{
- break;
- }
+ ct.ThrowIfCancellationRequested();
- // Skip if already fetched in a previous ingestion
- var exists = _db.WithCommand(cmd =>
- {
- cmd.CommandText = "SELECT 1 FROM helix_work_items WHERE job_name = @job AND work_item_name = @wi";
- cmd.Parameters.AddWithValue("@job", jobName);
- cmd.Parameters.AddWithValue("@wi", workItemName);
- return cmd.ExecuteScalar() is not null;
- });
+ var exists = db.WithCommand(cmd =>
+ {
+ cmd.CommandText = "SELECT 1 FROM helix_work_items WHERE job_name = @job AND work_item_name = @wi";
+ cmd.Parameters.AddWithValue("@job", jobName);
+ cmd.Parameters.AddWithValue("@wi", workItemName);
+ return cmd.ExecuteScalar() is not null;
+ });
- if (exists)
- {
- continue;
- }
+ if (exists)
+ {
+ continue;
+ }
- try
- {
- var workItem = await helixClient.GetWorkItemAsync(jobName, workItemName);
- items.Add(workItem);
- }
- catch (HttpRequestException ex)
- {
- _log?.Warning("Worker", $" Failed to fetch helix work item {jobName}/{workItemName}: {ex.Message}");
+ try
+ {
+ var workItem = await helixClient.GetWorkItemAsync(jobName, workItemName, ct);
+ helixWorkItems.Add(workItem);
+ }
+ catch (HttpRequestException ex)
+ {
+ log?.Warning("Worker", $" Failed to fetch helix work item {jobName}/{workItemName}: {ex.Message}");
+ }
}
}
- return items;
+ return new TestsData(summary, failures, helixWorkItems);
}
- void InsertTestData(SqliteConnection conn, SqliteTransaction tx)
+ static void InsertTestsData(TigerDatabase db, IngestionTask task, TestsData data)
{
- using var cmd = conn.CreateCommand();
- cmd.Transaction = tx;
-
- // Insert test runs for ALL runs from the summary
- foreach (var summary in testSummary)
- {
- InsertTestRun(cmd, task.Organization, task.Project, task.BuildId,
- summary.RunId, summary.JobName, summary.TotalCount, summary.PassedCount,
- summary.FailedCount, summary.SkippedCount, summary.Duration?.TotalSeconds);
- }
-
- // Insert individual failure results
- var runGroups = failures.GroupBy(f => f.TestRunId);
- foreach (var group in runGroups)
+ db.WithTransaction((conn, tx) =>
{
- var first = group.First();
+ using var cmd = conn.CreateCommand();
+ cmd.Transaction = tx;
- // If this run wasn't in the summary (unlikely but defensive), insert it now
- if (!testSummary.Any(s => s.RunId == group.Key))
+ foreach (var summary in data.Summary)
{
InsertTestRun(cmd, task.Organization, task.Project, task.BuildId,
- group.Key, first.TestRunName, group.Count(), 0, group.Count(), 0);
+ summary.RunId, summary.JobName, summary.TotalCount, summary.PassedCount,
+ summary.FailedCount, summary.SkippedCount, summary.Duration?.TotalSeconds);
}
- foreach (var r in group)
+ var runGroups = data.Failures.GroupBy(f => f.TestRunId);
+ foreach (var group in runGroups)
{
- InsertTestResult(cmd, task.Organization, task.Project, group.Key, r);
- }
- }
- }
+ var first = group.First();
+ if (!data.Summary.Any(s => s.RunId == group.Key))
+ {
+ InsertTestRun(cmd, task.Organization, task.Project, task.BuildId,
+ group.Key, first.TestRunName, group.Count(), 0, group.Count(), 0);
+ }
- void InsertHelixData(SqliteConnection conn, SqliteTransaction tx)
- {
- foreach (var workItem in helixWorkItems)
- {
- // Collect non-console-log files as JSON
- string? filesJson = null;
- if (workItem.Files is { Count: > 0 })
- {
- var filtered = workItem.Files
- .Where(f => !f.IsConsoleLog)
- .Select(f => new { fileName = f.FileName, uri = f.Uri })
- .ToList();
- if (filtered.Count > 0)
+ foreach (var r in group)
{
- filesJson = System.Text.Json.JsonSerializer.Serialize(filtered);
+ InsertTestResult(cmd, task.Organization, task.Project, group.Key, r);
}
}
- using (var cmd = conn.CreateCommand())
+ foreach (var workItem in data.HelixWorkItems)
{
- cmd.Transaction = tx;
- cmd.CommandText = """
- INSERT OR IGNORE INTO helix_work_items
- (job_name, work_item_name, state, exit_code, console_output_uri, files, is_deadletter)
- VALUES
- (@job, @wi, @state, @exitCode, @consoleUri, @files, @isDeadletter)
- """;
- cmd.Parameters.AddWithValue("@job", workItem.Job);
- cmd.Parameters.AddWithValue("@wi", workItem.Name);
- cmd.Parameters.AddWithValue("@state", workItem.State);
- cmd.Parameters.AddWithValue("@exitCode", workItem.ExitCode.HasValue ? workItem.ExitCode.Value : DBNull.Value);
- cmd.Parameters.AddWithValue("@consoleUri", (object?)workItem.ConsoleOutputUri ?? DBNull.Value);
- cmd.Parameters.AddWithValue("@files", (object?)filesJson ?? DBNull.Value);
- cmd.Parameters.AddWithValue("@isDeadletter", workItem.IsDeadLetter ? 1 : 0);
- cmd.ExecuteNonQuery();
- }
+ string? filesJson = null;
+ if (workItem.Files is { Count: > 0 })
+ {
+ var filtered = workItem.Files
+ .Where(f => !f.IsConsoleLog)
+ .Select(f => new { fileName = f.FileName, uri = f.Uri })
+ .ToList();
+ if (filtered.Count > 0)
+ {
+ filesJson = System.Text.Json.JsonSerializer.Serialize(filtered);
+ }
+ }
- if (workItem.IsDeadLetter)
- {
- using var cmd = conn.CreateCommand();
- cmd.Transaction = tx;
- cmd.CommandText = """
- UPDATE test_results
- SET error_message = 'Helix Work Item Dead Lettered. ' || COALESCE(error_message, '')
- WHERE is_helix_work_item = 1
- AND helix_job_name = @job
- AND helix_work_item_name = @wi
- AND organization = @org
- AND run_id IN (
- SELECT run_id FROM test_runs
- WHERE organization = @org AND build_id = @buildId
- )
- AND error_message NOT LIKE 'Helix Work Item Dead Lettered.%'
- """;
- cmd.Parameters.AddWithValue("@job", workItem.Job);
- cmd.Parameters.AddWithValue("@wi", workItem.Name);
- cmd.Parameters.AddWithValue("@org", task.Organization);
- cmd.Parameters.AddWithValue("@buildId", task.BuildId);
- cmd.ExecuteNonQuery();
- _log?.Warning("Worker", $" Helix work item {workItem.Name} is dead-lettered");
+ using (var helixCmd = conn.CreateCommand())
+ {
+ helixCmd.Transaction = tx;
+ helixCmd.CommandText = """
+ INSERT OR IGNORE INTO helix_work_items
+ (job_name, work_item_name, state, exit_code, console_output_uri, files, is_deadletter)
+ VALUES
+ (@job, @wi, @state, @exitCode, @consoleUri, @files, @isDeadletter)
+ """;
+ helixCmd.Parameters.AddWithValue("@job", workItem.Job);
+ helixCmd.Parameters.AddWithValue("@wi", workItem.Name);
+ helixCmd.Parameters.AddWithValue("@state", workItem.State);
+ helixCmd.Parameters.AddWithValue("@exitCode", workItem.ExitCode.HasValue ? workItem.ExitCode.Value : DBNull.Value);
+ helixCmd.Parameters.AddWithValue("@consoleUri", (object?)workItem.ConsoleOutputUri ?? DBNull.Value);
+ helixCmd.Parameters.AddWithValue("@files", (object?)filesJson ?? DBNull.Value);
+ helixCmd.Parameters.AddWithValue("@isDeadletter", workItem.IsDeadLetter ? 1 : 0);
+ helixCmd.ExecuteNonQuery();
+ }
+
+ if (workItem.IsDeadLetter)
+ {
+ using var dlCmd = conn.CreateCommand();
+ dlCmd.Transaction = tx;
+ dlCmd.CommandText = """
+ UPDATE test_results
+ SET error_message = 'Helix Work Item Dead Lettered. ' || COALESCE(error_message, '')
+ WHERE is_helix_work_item = 1
+ AND helix_job_name = @job
+ AND helix_work_item_name = @wi
+ AND organization = @org
+ AND run_id IN (
+ SELECT run_id FROM test_runs
+ WHERE organization = @org AND build_id = @buildId
+ )
+ AND error_message NOT LIKE 'Helix Work Item Dead Lettered.%'
+ """;
+ dlCmd.Parameters.AddWithValue("@job", workItem.Job);
+ dlCmd.Parameters.AddWithValue("@wi", workItem.Name);
+ dlCmd.Parameters.AddWithValue("@org", task.Organization);
+ dlCmd.Parameters.AddWithValue("@buildId", task.BuildId);
+ dlCmd.ExecuteNonQuery();
+ }
}
- }
+ });
}
- }
- private async Task ProcessTimelineAsync(AzdoClient client, IngestionTask task, CancellationToken ct)
- {
- _log?.Info("Worker", $"Fetching timeline for build #{task.BuildId}...");
- var timeline = await client.GetTimelineAsync(task.BuildId);
- InsertTimelineIssues(task.Organization, task.Project, task.BuildId, timeline);
+ // ── Timeline ────────────────────────────────────────────────
- var issueCount = timeline.Records.Sum(r => r.Issues.Count(i => i.Type is "error" or "warning"));
- _log?.Info("Worker", $" Build #{task.BuildId} — timeline complete ({issueCount} issues)");
- }
-
- private async Task ProcessPrInfoAsync(IngestionTask task, CancellationToken ct)
- {
- // Look up the build's PR number and repo
- var prInfo = _db.WithCommand(cmd =>
+ async Task ProcessTimeline()
{
- cmd.CommandText = "SELECT pr_number, repository_name FROM builds WHERE organization = @org AND build_id = @buildId";
- cmd.Parameters.AddWithValue("@org", task.Organization);
- cmd.Parameters.AddWithValue("@buildId", task.BuildId);
+ _log?.Info("Worker", $"Fetching timeline for build #{task.BuildId}...");
+ var timeline = await FetchTimelineDataAsync(client, task, ct);
+ InsertTimelineData(_db, task, timeline);
- using var reader = cmd.ExecuteReader();
- if (!reader.Read() || reader.IsDBNull(0) || reader.IsDBNull(1))
- {
- return ((int PrNumber, string Repository)?)null;
- }
+ var issueCount = timeline.Records.Sum(r => r.Issues.Count(i => i.Type is "error" or "warning"));
+ _log?.Info("Worker", $" Build #{task.BuildId} — timeline complete ({issueCount} issues)");
+ }
- return (reader.GetInt32(0), reader.GetString(1));
- });
+ static async Task FetchTimelineDataAsync(
+ AzdoClient client, IngestionTask task, CancellationToken ct)
+ {
+ return await client.GetTimelineAsync(task.BuildId, ct);
+ }
- if (prInfo is null)
+ static void InsertTimelineData(TigerDatabase db, IngestionTask task, AzdoTimeline timeline)
{
- _log?.Info("Worker", $" Build #{task.BuildId} — no PR info to fetch");
- return;
+ InsertTimelineIssues(db, task.Organization, task.Project, task.BuildId, timeline);
}
- var (prNumber, repository) = prInfo.Value;
+ // ── PR Info ─────────────────────────────────────────────────
- // Check if we already have this PR cached
- var exists = _db.WithCommand(cmd =>
- {
- cmd.CommandText = "SELECT 1 FROM pull_requests WHERE repository = @repo AND pr_number = @pr";
- cmd.Parameters.AddWithValue("@repo", repository);
- cmd.Parameters.AddWithValue("@pr", prNumber);
- return cmd.ExecuteScalar() is not null;
- });
- if (exists)
+ async Task ProcessPrInfo()
{
- _log?.Info("Worker", $" Build #{task.BuildId} — PR #{prNumber} already cached");
- return;
+ var prData = await FetchPrInfoDataAsync(_db, _log, task, ct);
+ if (prData is not null)
+ {
+ InsertPrInfoData(_db, prData.Value);
+ _log?.Info("Worker", $" Build #{task.BuildId} — PR #{prData.Value.PrNumber} info cached ({prData.Value.Author})");
+ }
}
- // Fetch PR info via gh CLI
- try
+ static async Task FetchPrInfoDataAsync(
+ TigerDatabase db, ServiceLog? log, IngestionTask task, CancellationToken ct)
{
+ var prInfo = db.WithCommand(cmd =>
+ {
+ cmd.CommandText = "SELECT pr_number, repository_name FROM builds WHERE organization = @org AND build_id = @buildId";
+ cmd.Parameters.AddWithValue("@org", task.Organization);
+ cmd.Parameters.AddWithValue("@buildId", task.BuildId);
+
+ using var reader = cmd.ExecuteReader();
+ if (!reader.Read() || reader.IsDBNull(0) || reader.IsDBNull(1))
+ {
+ return ((int PrNumber, string Repository)?)null;
+ }
+
+ return (reader.GetInt32(0), reader.GetString(1));
+ });
+
+ if (prInfo is null)
+ {
+ log?.Info("Worker", $" Build #{task.BuildId} — no PR info to fetch");
+ return null;
+ }
+
+ var (prNumber, repository) = prInfo.Value;
+
+ var exists = db.WithCommand(cmd =>
+ {
+ cmd.CommandText = "SELECT 1 FROM pull_requests WHERE repository = @repo AND pr_number = @pr";
+ cmd.Parameters.AddWithValue("@repo", repository);
+ cmd.Parameters.AddWithValue("@pr", prNumber);
+ return cmd.ExecuteScalar() is not null;
+ });
+ if (exists)
+ {
+ log?.Info("Worker", $" Build #{task.BuildId} — PR #{prNumber} already cached");
+ return null;
+ }
+
var psi = new System.Diagnostics.ProcessStartInfo("gh", $"pr view {prNumber} --repo {repository} --json title,author")
{
RedirectStandardOutput = true,
@@ -921,8 +1056,8 @@ private async Task ProcessPrInfoAsync(IngestionTask task, CancellationToken ct)
using var process = System.Diagnostics.Process.Start(psi);
if (process is null)
{
- _log?.Warning("Worker", $" Build #{task.BuildId} — failed to start gh process");
- return;
+ log?.Warning("Worker", $" Build #{task.BuildId} — failed to start gh process");
+ return null;
}
var output = await process.StandardOutput.ReadToEndAsync(ct);
@@ -930,47 +1065,54 @@ private async Task ProcessPrInfoAsync(IngestionTask task, CancellationToken ct)
if (process.ExitCode != 0)
{
- _log?.Warning("Worker", $" Build #{task.BuildId} — gh pr view failed (exit {process.ExitCode})");
- return;
+ log?.Warning("Worker", $" Build #{task.BuildId} — gh pr view failed (exit {process.ExitCode})");
+ return null;
}
- var prData = System.Text.Json.JsonDocument.Parse(output);
- var title = prData.RootElement.TryGetProperty("title", out var t) ? t.GetString() : null;
- var author = prData.RootElement.TryGetProperty("author", out var a) && a.TryGetProperty("login", out var login)
+ var prDoc = System.Text.Json.JsonDocument.Parse(output);
+ var title = prDoc.RootElement.TryGetProperty("title", out var t) ? t.GetString() : null;
+ var author = prDoc.RootElement.TryGetProperty("author", out var a) && a.TryGetProperty("login", out var login)
? login.GetString() : null;
- _db.WithCommand(cmd =>
+ return new PrInfoData(repository, prNumber, title, author);
+ }
+
+ static void InsertPrInfoData(TigerDatabase db, PrInfoData data)
+ {
+ db.WithCommand(cmd =>
{
cmd.CommandText = """
INSERT OR IGNORE INTO pull_requests (repository, pr_number, title, author)
VALUES (@repo, @pr, @title, @author)
""";
- cmd.Parameters.AddWithValue("@repo", repository);
- cmd.Parameters.AddWithValue("@pr", prNumber);
- cmd.Parameters.AddWithValue("@title", (object?)title ?? DBNull.Value);
- cmd.Parameters.AddWithValue("@author", (object?)author ?? DBNull.Value);
+ cmd.Parameters.AddWithValue("@repo", data.Repository);
+ cmd.Parameters.AddWithValue("@pr", data.PrNumber);
+ cmd.Parameters.AddWithValue("@title", (object?)data.Title ?? DBNull.Value);
+ cmd.Parameters.AddWithValue("@author", (object?)data.Author ?? DBNull.Value);
cmd.ExecuteNonQuery();
});
-
- _log?.Info("Worker", $" Build #{task.BuildId} — PR #{prNumber} info cached ({author})");
- }
- catch (Exception ex) when (ex is not OperationCanceledException)
- {
- _log?.Warning("Worker", $" Build #{task.BuildId} — PR info fetch failed: {ex.Message}");
}
}
+ private readonly record struct TestsData(
+ List Summary,
+ List Failures,
+ List HelixWorkItems);
+
+ private readonly record struct PrInfoData(
+ string Repository, int PrNumber, string? Title, string? Author);
+
// ── DB Helpers ──────────────────────────────────────────────────
- private IngestionTask? GetNextReadyTask()
+ private (IngestionTask Task, bool IsPriority)? GetNextReadyTask()
{
// Priority tasks take precedence over the normal DB query
if (_priorityTasks.TryPop(out var priority))
{
- return priority;
+ return (priority, true);
}
- return _db.WithCommand(cmd =>
+ var dbTask = _db.WithCommand(cmd =>
{
cmd.CommandText = """
SELECT t.organization, b.project, t.build_id, t.task_type, t.status, t.attempts
@@ -997,6 +1139,8 @@ LIMIT 1
return null;
});
+
+ return dbTask is not null ? (dbTask, false) : null;
}
///
@@ -1109,6 +1253,7 @@ public void Dispose()
{
_cts?.Cancel();
_cts?.Dispose();
+ _prioritySignal.Dispose();
}
private record IngestionTask(
diff --git a/src/Tiger/HelixClient.cs b/src/Tiger/HelixClient.cs
index ac379f9..7144a57 100644
--- a/src/Tiger/HelixClient.cs
+++ b/src/Tiger/HelixClient.cs
@@ -37,6 +37,16 @@ public static HelixClient Create(string? bearerToken = null)
return new HelixClient(httpClient);
}
+ ///
+ /// Creates a with a custom
+ /// for testing purposes.
+ ///
+ public static HelixClient Create(HttpMessageHandler handler)
+ {
+ var httpClient = new HttpClient(handler) { BaseAddress = new Uri(BaseUrl) };
+ return new HelixClient(httpClient);
+ }
+
///
public static Task CreateAsync(string? bearerToken = null) =>
Task.FromResult(Create(bearerToken));
@@ -50,39 +60,39 @@ public static string GetConsoleUrl(string jobName, string workItemName) =>
///
/// Get summary information about a single job.
///
- public async Task GetJobAsync(string jobName)
+ public async Task GetJobAsync(string jobName, CancellationToken ct = default)
{
var url = $"api/jobs/{Uri.EscapeDataString(jobName)}?api-version={ApiVersion}";
- return await GetAsync(url);
+ return await GetAsync(url, ct);
}
///
/// List all work items for a given job.
///
- public async Task> GetWorkItemsAsync(string jobName)
+ public async Task> GetWorkItemsAsync(string jobName, CancellationToken ct = default)
{
var url = $"api/jobs/{Uri.EscapeDataString(jobName)}/workitems?api-version={ApiVersion}";
- return await GetAsync>(url);
+ return await GetAsync>(url, ct);
}
///
/// Get detailed information about a single work item.
///
- public async Task GetWorkItemAsync(string jobName, string workItemName)
+ public async Task GetWorkItemAsync(string jobName, string workItemName, CancellationToken ct = default)
{
var url = $"api/jobs/{Uri.EscapeDataString(jobName)}/workitems/{Uri.EscapeDataString(workItemName)}?api-version={ApiVersion}";
- return await GetAsync(url);
+ return await GetAsync(url, ct);
}
///
/// Get console output for a specific work item.
///
- public async Task GetConsoleAsync(string jobName, string workItemName)
+ public async Task GetConsoleAsync(string jobName, string workItemName, CancellationToken ct = default)
{
var url = $"api/jobs/{Uri.EscapeDataString(jobName)}/workitems/{Uri.EscapeDataString(workItemName)}/console?api-version={ApiVersion}";
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var text = await response.Content.ReadAsStringAsync();
+ var text = await response.Content.ReadAsStringAsync(ct);
return new HelixWorkItemConsole
{
Job = jobName,
@@ -94,12 +104,12 @@ public async Task GetConsoleAsync(string jobName, string w
///
/// Get console output for multiple work items.
///
- public async Task> GetConsolesAsync(string jobName, List workItems)
+ public async Task> GetConsolesAsync(string jobName, List workItems, CancellationToken ct = default)
{
var list = new List();
foreach (var workItem in workItems)
{
- var console = await GetConsoleAsync(jobName, workItem.Name);
+ var console = await GetConsoleAsync(jobName, workItem.Name, ct);
list.Add(console);
}
return list;
@@ -108,35 +118,35 @@ public async Task> GetConsolesAsync(string jobName, L
///
/// List files uploaded from a specific work item.
///
- public async Task> GetFilesAsync(string jobName, string workItemName)
+ public async Task> GetFilesAsync(string jobName, string workItemName, CancellationToken ct = default)
{
var url = $"api/jobs/{Uri.EscapeDataString(jobName)}/workitems/{Uri.EscapeDataString(workItemName)}/files?api-version={ApiVersion}";
- return await GetAsync>(url);
+ return await GetAsync>(url, ct);
}
///
/// Download a specific file from a work item to a local path.
///
- public async Task DownloadFileAsync(string jobName, string workItemName, string fileName, string outputPath)
+ public async Task DownloadFileAsync(string jobName, string workItemName, string fileName, string outputPath, CancellationToken ct = default)
{
var url = $"api/jobs/{Uri.EscapeDataString(jobName)}/workitems/{Uri.EscapeDataString(workItemName)}/files/{Uri.EscapeDataString(fileName)}?api-version={ApiVersion}";
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var bytes = await response.Content.ReadAsByteArrayAsync();
+ var bytes = await response.Content.ReadAsByteArrayAsync(ct);
var dir = Path.GetDirectoryName(outputPath);
if (dir is not null)
{
Directory.CreateDirectory(dir);
}
- await File.WriteAllBytesAsync(outputPath, bytes);
+ await File.WriteAllBytesAsync(outputPath, bytes, ct);
}
///
/// Download all files from a work item to a directory.
///
- public async Task DownloadFilesAsync(string jobName, string workItemName, string outputDir)
+ public async Task DownloadFilesAsync(string jobName, string workItemName, string outputDir, CancellationToken ct = default)
{
- var files = await GetFilesAsync(jobName, workItemName);
+ var files = await GetFilesAsync(jobName, workItemName, ct);
using var httpClient = new HttpClient();
foreach (var file in files)
{
@@ -148,16 +158,16 @@ public async Task DownloadFilesAsync(string jobName, string workItemName, string
{
Directory.CreateDirectory(dir);
}
- var bytes = await httpClient.GetByteArrayAsync(file.Link);
- await File.WriteAllBytesAsync(filePath, bytes);
+ var bytes = await httpClient.GetByteArrayAsync(file.Link, ct);
+ await File.WriteAllBytesAsync(filePath, bytes, ct);
}
}
- private async Task GetAsync(string url)
+ private async Task GetAsync(string url, CancellationToken ct = default)
{
- var response = await HttpClient.GetAsync(url);
+ var response = await HttpClient.GetAsync(url, ct);
response.EnsureSuccessStatusCode();
- var json = await response.Content.ReadAsStringAsync();
+ var json = await response.Content.ReadAsStringAsync(ct);
return JsonSerializer.Deserialize(json, s_jsonOptions)
?? throw new InvalidOperationException($"Failed to deserialize response from {url}");
}