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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,15 @@ jobs:
- name: Checkout
uses: actions/checkout@v7

- name: Verify trained models are present
run: |
for f in F1Predictor.WebApi/models/podium-model.zip F1Predictor.WebApi/models/points-model.zip; do
if [ ! -s "$f" ]; then
echo "::error::Missing or empty $f — the Dockerfile bakes this into the production image; commit the trained model before merging." >&2
exit 1
fi
done

- name: Setup .NET 10
uses: actions/setup-dotnet@v6
with:
Expand Down
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -390,6 +390,7 @@ appsettings.*.local.json
# Local planning and prompt files
DOTNET_2026_STANDARDS_PROMPT.md
OTEL_ASPIRE_MIGRATION_PROMPT.md
/docs/

/graphify-out

Expand Down
10 changes: 5 additions & 5 deletions Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
</PropertyGroup>
<ItemGroup>
<PackageVersion Include="AspNetCore.HealthChecks.NpgSql" Version="9.0.0" />
<PackageVersion Include="Aspire.Hosting.JavaScript" Version="13.4.6" />
<PackageVersion Include="Aspire.Hosting.JavaScript" Version="13.5.3" />
<PackageVersion Include="FluentValidation" Version="12.1.1" />
<PackageVersion Include="MessagePack" Version="3.1.8" />
<PackageVersion Include="Microsoft.EntityFrameworkCore.Design" Version="10.0.9">
Expand All @@ -19,7 +19,7 @@
<PackageVersion Include="Microsoft.Extensions.Configuration.FileExtensions" Version="10.0.11" />
<PackageVersion Include="Microsoft.Extensions.Configuration.Json" Version="10.0.11" />
<PackageVersion Include="Microsoft.Extensions.TimeProvider.Testing" Version="10.5.0" />
<PackageVersion Include="Microsoft.Extensions.Logging.Debug" Version="10.0.9" />
<PackageVersion Include="Microsoft.Extensions.Logging.Debug" Version="10.0.11" />
<PackageVersion Include="Microsoft.Data.Sqlite.Core" Version="10.0.9" />
<PackageVersion Include="AspNetCore.HealthChecks.UI.Client" Version="9.0.0" />
<PackageVersion Include="Bogus" Version="35.6.5" />
Expand Down Expand Up @@ -50,10 +50,10 @@
<PackageVersion Include="SQLitePCLRaw.core" Version="2.1.12" />
<PackageVersion Include="Microsoft.EntityFrameworkCore.Tools" Version="10.0.9" />
<PackageVersion Include="Microsoft.Extensions.Azure" Version="1.13.1" />
<PackageVersion Include="Microsoft.Extensions.Configuration.UserSecrets" Version="10.0.9" />
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="10.0.9" />
<PackageVersion Include="Microsoft.Extensions.Configuration.UserSecrets" Version="10.0.11" />
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="10.0.11" />
<PackageVersion Include="Microsoft.Extensions.Diagnostics.HealthChecks.EntityFrameworkCore" Version="10.0.9" />
<PackageVersion Include="Microsoft.Extensions.Http" Version="10.0.9" />
<PackageVersion Include="Microsoft.Extensions.Http" Version="10.0.11" />
<PackageVersion Include="Microsoft.ML" Version="5.0.0" />
<PackageVersion Include="Microsoft.Extensions.ML" Version="5.0.0" />
<PackageVersion Include="Microsoft.Extensions.Logging.Abstractions" Version="10.0.11" />
Expand Down
2 changes: 1 addition & 1 deletion F1Predictor.AppHost/F1Predictor.AppHost.csproj
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
<Project Sdk="Aspire.AppHost.Sdk/13.4.6">
<Project Sdk="Aspire.AppHost.Sdk/13.5.3">

<PropertyGroup>
<OutputType>Exe</OutputType>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
using F1Predictor.Domain.RaceData.Entities;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.ChangeTracking;
using Microsoft.EntityFrameworkCore.Infrastructure;

namespace F1Predictor.Application.Abstractions.Data;

Expand All @@ -25,5 +26,21 @@ public interface IApplicationDbContext
/// </summary>
ChangeTracker ChangeTracker { get; }

/// <summary>
/// Provides access to database-level operations such as transactions. Used for advanced
/// scenarios like serializing concurrent writers with an advisory lock (see
/// <see cref="AcquireSessionAdvisoryLockAsync"/>).
/// </summary>
DatabaseFacade Database { get; }

/// <summary>
/// Takes a Postgres transaction-scoped advisory lock keyed on <paramref name="sessionKey"/>.
/// Must be called inside an active transaction — the lock releases automatically when that
/// transaction commits or rolls back. Serializes concurrent ingestion of the same session
/// across requests, scheduler ticks, and — once scaled beyond one replica — processes, none
/// of which otherwise coordinate with each other.
/// </summary>
Task AcquireSessionAdvisoryLockAsync(int sessionKey, CancellationToken cancellationToken);

Task<int> SaveChangesAsync(CancellationToken cancellationToken = default);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
namespace F1Predictor.Application.Features.Seasons.GetDataStatus;

/// <param name="IsStale">True if any of the three signals below means the pipeline needs a re-run.</param>
/// <param name="PendingRaceName">
/// Name of the earliest Grand Prix that has happened but has no classified result yet, or null if
/// none. Null cannot distinguish "nothing pending" from "OpenF1 never published this round's
/// result" — the same ambiguity <c>ChampionshipPoints</c> already accepts elsewhere.
/// </param>
/// <param name="PendingRaceDate">Start date of <see cref="PendingRaceName"/>, if any.</param>
/// <param name="FeaturesStale">True if a classified Grand Prix has no feature rows yet.</param>
/// <param name="ModelsAvailable">False if the models have never been trained.</param>
public sealed record DataStatusResponse(
bool IsStale,
string? PendingRaceName,
DateTimeOffset? PendingRaceDate,
bool FeaturesStale,
bool ModelsAvailable);
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
using F1Predictor.Application.Abstractions.Messaging;

namespace F1Predictor.Application.Features.Seasons.GetDataStatus;

/// <summary>
/// Whether a season's standings and predictions are current with what has actually raced.
/// </summary>
public sealed record GetDataStatusQuery(int Year) : IQuery<DataStatusResponse>;
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
using F1Predictor.Application.Abstractions.Data;
using F1Predictor.Application.Abstractions.MachineLearning;
using F1Predictor.Application.Abstractions.Messaging;
using Microsoft.EntityFrameworkCore;
using SharedKernel;

namespace F1Predictor.Application.Features.Seasons.GetDataStatus;

internal sealed class GetDataStatusQueryHandler(IApplicationDbContext context, IRacePredictor predictor)
: IQueryHandler<GetDataStatusQuery, DataStatusResponse>
{
public async Task<Result<DataStatusResponse>> Handle(
GetDataStatusQuery query,
CancellationToken cancellationToken)
{
var now = DateTimeOffset.UtcNow;

var pending = await (
from session in context.RaceSessions
join meeting in context.Meetings on session.MeetingKey equals meeting.MeetingKey
where meeting.Year == query.Year
&& !session.IsSprint
&& !session.IsClassified
&& session.DateStart < now
orderby session.DateStart
select new { meeting.MeetingName, session.DateStart })
.AsNoTracking()
.FirstOrDefaultAsync(cancellationToken);

var featuresStale = await context.RaceSessions
.Join(context.Meetings, session => session.MeetingKey, meeting => meeting.MeetingKey, (session, meeting) => new { session, meeting })
.Where(x => x.meeting.Year == query.Year && !x.session.IsSprint && x.session.IsClassified)
.AnyAsync(x => !context.DriverRaceFeatures.Any(f => f.SessionKey == x.session.SessionKey), cancellationToken);

var modelsAvailable = predictor.ModelsAvailable;

var isStale = pending is not null || featuresStale || !modelsAvailable;

return Result.Success(new DataStatusResponse(
isStale,
pending?.MeetingName,
pending?.DateStart,
featuresStale,
modelsAvailable));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -134,28 +134,31 @@ private static IngestOutcome Summarise(IReadOnlyList<IngestOutcome> outcomes)
var isSprint = string.Equals(session.SessionName, SprintSessionName, StringComparison.Ordinal);
var label = $"{meeting.MeetingName} ({session.SessionName})";

var existing = await context.RaceSessions
// Cheap early-out before bothering with a lock or the OpenF1 round trips below. The
// authoritative check — the one that actually decides skip vs. ingest — runs again
// after the advisory lock is acquired, since a concurrent writer may commit between
// this read and that point.
var precheck = await context.RaceSessions
.AsNoTracking()
.FirstOrDefaultAsync(s => s.SessionKey == session.SessionKey, cancellationToken);

// A session already stored with results is settled; anything else is re-checked so a
// scheduled race is picked up automatically once it runs. `force` re-fetches either way,
// which is the only route by which a provisional classification ever gets corrected.
if (existing is { IsClassified: true } && !force)
if (precheck is { IsClassified: true } && !force)
{
logger.LogInformation("{Label}: already ingested, skipping.", label);
return (IngestOutcome.AlreadyPresent, SessionTotals.Empty);
}

// No reason to hold a DB lock across network I/O, so fetch before taking it.
var payload = await FetchSessionAsync(session, allSessions, isSprint, cancellationToken);

if (existing is not null)
var written = await WriteSessionUnderLockAsync(session, isSprint, force, payload, cancellationToken);

if (!written)
{
await ClearSessionAsync(session.SessionKey, cancellationToken);
logger.LogInformation("{Label}: already ingested by a concurrent run, skipping.", label);
return (IngestOutcome.AlreadyPresent, SessionTotals.Empty);
}

await PersistSessionAsync(session, isSprint, payload, cancellationToken);

var totals = new SessionTotals(
payload.Results.Count,
payload.Grid.Count,
Expand All @@ -181,6 +184,48 @@ private static IngestOutcome Summarise(IReadOnlyList<IngestOutcome> outcomes)
return (IngestOutcome.Ingested, totals);
}

/// <summary>
/// Clears and re-persists one session's rows inside a transaction guarded by a Postgres
/// advisory lock keyed on <paramref name="session"/>'s key. Two triggers can otherwise race
/// to ingest the same session — a manual API call overlapping the Quartz coordinator's
/// startup tick, or, once this app scales beyond one replica, each replica's own scheduler
/// firing independently. The lock, held for the life of the transaction, serializes them at
/// the database so the loser re-checks post-commit state and skips instead of racing the
/// unique index on <c>DriverEntries</c>.
/// </summary>
/// <returns>False if a concurrent run already classified this session and <paramref name="force"/> is not set.</returns>
private async Task<bool> WriteSessionUnderLockAsync(
OpenF1Session session,
bool isSprint,
bool force,
SessionPayload payload,
CancellationToken cancellationToken)
{
await using var transaction = await context.Database.BeginTransactionAsync(cancellationToken);

await context.AcquireSessionAdvisoryLockAsync(session.SessionKey, cancellationToken);

var existing = await context.RaceSessions
.AsNoTracking()
.FirstOrDefaultAsync(s => s.SessionKey == session.SessionKey, cancellationToken);

if (existing is { IsClassified: true } && !force)
{
return false;
}

if (existing is not null)
{
await ClearSessionAsync(session.SessionKey, cancellationToken);
}

await PersistSessionAsync(session, isSprint, payload, cancellationToken);

await transaction.CommitAsync(cancellationToken);

return true;
}

/// <summary>
/// Pulls everything OpenF1 has for one session. How much that is depends on whether the
/// session has run — see the remarks on each call below.
Expand Down
4 changes: 4 additions & 0 deletions F1Predictor.Infrastructure/Database/ApplicationDbContext.cs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,10 @@ public class ApplicationDbContext(
public DbSet<DriverEntry> DriverEntries { get; set; } = null!;
public DbSet<DriverRaceFeature> DriverRaceFeatures { get; set; } = null!;

public Task AcquireSessionAdvisoryLockAsync(int sessionKey, CancellationToken cancellationToken) =>
Database.ExecuteSqlInterpolatedAsync(
$"SELECT pg_advisory_xact_lock({sessionKey})", cancellationToken);

public override async Task<int> SaveChangesAsync(CancellationToken cancellationToken = default)
{
int result = await base.SaveChangesAsync(cancellationToken);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ namespace F1Predictor.Infrastructure.Ingestion;
/// due — so this can run far more often than a season would ever need re-fetching without
/// wasting the free, rate-limited API on seasons that haven't changed.
/// </remarks>
[DisallowConcurrentExecution]
internal sealed class SeasonIngestionCoordinatorJob(
IApplicationDbContext dbContext,
ICommandHandler<IngestSeasonCommand, IngestSeasonResponse> ingestHandler,
Expand Down
Loading
Loading