|
| 1 | +using System.Text.Json; |
| 2 | +using CodeSpace.Core.Persistence.Db; |
| 3 | +using CodeSpace.Core.Persistence.Entities; |
| 4 | +using CodeSpace.Core.Services.Workflows.Artifacts.Profiles; |
| 5 | +using CodeSpace.Core.Services.Workflows.Artifacts.Providers.Local; |
| 6 | +using CodeSpace.Core.Settings; |
| 7 | +using CodeSpace.Messages.Constants; |
| 8 | +using Microsoft.EntityFrameworkCore; |
| 9 | +using Npgsql; |
| 10 | + |
| 11 | +namespace CodeSpace.Core.Services.Agents.AgentRunLogging; |
| 12 | + |
| 13 | +/// <summary> |
| 14 | +/// Gives an unconfigured team one real Settings-visible local route only after the deployment explicitly qualifies its |
| 15 | +/// local-rwx namespace as shared. This is a missing-only bootstrap, not a routing fallback: an existing route in any |
| 16 | +/// lifecycle state or any collision on the reserved profile name remains authoritative, and later Settings revisions can |
| 17 | +/// move the data class to any registered cloud provider without changing Agent Run capture. |
| 18 | +/// </summary> |
| 19 | +public sealed class AgentRunLogStorageReadiness : IAgentRunLogStorageReadiness |
| 20 | +{ |
| 21 | + internal const string DefaultProfileStableName = "codespace-agent-run-log-default"; |
| 22 | + private const int AdvisoryLockNamespace = 117; |
| 23 | + private readonly CodeSpaceDbContext _db; |
| 24 | + private readonly TimeProvider _clock; |
| 25 | + |
| 26 | + public AgentRunLogStorageReadiness(CodeSpaceDbContext db, TimeProvider clock) |
| 27 | + { |
| 28 | + _db = db; |
| 29 | + _clock = clock; |
| 30 | + } |
| 31 | + |
| 32 | + public async Task EnsureDefaultRouteAsync(Guid teamId, CancellationToken cancellationToken) |
| 33 | + { |
| 34 | + if (teamId == Guid.Empty || !RuntimeSettings.Current.ArtifactLocalRwxShared) return; |
| 35 | + |
| 36 | + await using var transaction = await _db.Database.BeginTransactionAsync(cancellationToken).ConfigureAwait(false); |
| 37 | + try |
| 38 | + { |
| 39 | + await _db.Database.ExecuteSqlInterpolatedAsync($"SELECT pg_advisory_xact_lock(hashtextextended({teamId.ToString()}, {AdvisoryLockNamespace}))", cancellationToken).ConfigureAwait(false); |
| 40 | + if (!await _db.Team.AsNoTracking().AnyAsync(team => team.Id == teamId, cancellationToken).ConfigureAwait(false) |
| 41 | + || await _db.StorageRoute.AsNoTracking().AnyAsync(route => route.TeamId == teamId && route.DataClassTypeKey == AgentRunLogStorageResolver.DataClassTypeKey, cancellationToken).ConfigureAwait(false)) |
| 42 | + { |
| 43 | + await transaction.CommitAsync(cancellationToken).ConfigureAwait(false); |
| 44 | + return; |
| 45 | + } |
| 46 | + |
| 47 | + var rootPath = DurableRoots.ArtifactStore(RuntimeSettings.Current.ArtifactStoreDirectory); |
| 48 | + var canonicalConfig = CanonicalConfig(rootPath); |
| 49 | + if (await _db.StorageProfile.AsNoTracking().AnyAsync( |
| 50 | + value => value.TeamId == teamId && value.StableName == DefaultProfileStableName, cancellationToken).ConfigureAwait(false)) |
| 51 | + { |
| 52 | + await transaction.CommitAsync(cancellationToken).ConfigureAwait(false); |
| 53 | + return; |
| 54 | + } |
| 55 | + |
| 56 | + var profile = BuildProfile(teamId, canonicalConfig, _clock.GetUtcNow()); |
| 57 | + _db.StorageProfile.Add(profile); |
| 58 | + var route = BuildRoute(teamId, profile.Id, _clock.GetUtcNow()); |
| 59 | + _db.StorageRoute.Add(route); |
| 60 | + await _db.SaveChangesAsync(cancellationToken).ConfigureAwait(false); |
| 61 | + route.State = StorageRouteState.Active; |
| 62 | + route.LastModifiedDate = _clock.GetUtcNow(); |
| 63 | + route.LastModifiedBy = SystemUsers.SeederId; |
| 64 | + await _db.SaveChangesAsync(cancellationToken).ConfigureAwait(false); |
| 65 | + await transaction.CommitAsync(cancellationToken).ConfigureAwait(false); |
| 66 | + } |
| 67 | + catch (Exception exception) when (IsUniqueViolation(exception)) |
| 68 | + { |
| 69 | + await transaction.RollbackAsync(CancellationToken.None).ConfigureAwait(false); |
| 70 | + _db.ChangeTracker.Clear(); |
| 71 | + } |
| 72 | + } |
| 73 | + |
| 74 | + private static StorageProfile BuildProfile(Guid teamId, string canonicalConfig, DateTimeOffset now) |
| 75 | + { |
| 76 | + var profile = new StorageProfile |
| 77 | + { |
| 78 | + Id = Guid.NewGuid(), TeamId = teamId, StableName = DefaultProfileStableName, CurrentRevision = 1, |
| 79 | + State = StorageProfileState.Active, CreatedDate = now, CreatedBy = SystemUsers.SeederId, |
| 80 | + LastModifiedDate = now, LastModifiedBy = SystemUsers.SeederId, |
| 81 | + }; |
| 82 | + using var document = JsonDocument.Parse(canonicalConfig); |
| 83 | + profile.Revisions.Add(new StorageProfileRevision |
| 84 | + { |
| 85 | + Id = Guid.NewGuid(), TeamId = teamId, StorageProfileId = profile.Id, Revision = 1, |
| 86 | + ProviderTypeKey = LocalRwxArtifactStorageDriverFactory.TypeKey, NonSecretConfigJson = canonicalConfig, |
| 87 | + CredentialRef = null, NamespaceFingerprint = StorageProfileRules.NamespaceFingerprint(LocalRwxArtifactStorageDriverFactory.TypeKey, document.RootElement), |
| 88 | + CreatedDate = now, CreatedBy = SystemUsers.SeederId, |
| 89 | + }); |
| 90 | + return profile; |
| 91 | + } |
| 92 | + |
| 93 | + private static StorageRoute BuildRoute(Guid teamId, Guid profileId, DateTimeOffset now) |
| 94 | + { |
| 95 | + var route = new StorageRoute |
| 96 | + { |
| 97 | + Id = Guid.NewGuid(), TeamId = teamId, DataClassTypeKey = AgentRunLogStorageResolver.DataClassTypeKey, |
| 98 | + CurrentRevision = 1, State = StorageRouteState.Draft, CreatedDate = now, CreatedBy = SystemUsers.SeederId, |
| 99 | + LastModifiedDate = now, LastModifiedBy = SystemUsers.SeederId, |
| 100 | + }; |
| 101 | + route.Revisions.Add(new StorageRouteRevision |
| 102 | + { |
| 103 | + Id = Guid.NewGuid(), TeamId = teamId, StorageRouteId = route.Id, Revision = 1, StorageProfileId = profileId, |
| 104 | + ProfileRevisionMode = StorageProfileRevisionMode.CurrentAtWrite, PinnedProfileRevision = null, |
| 105 | + CreatedDate = now, CreatedBy = SystemUsers.SeederId, |
| 106 | + }); |
| 107 | + return route; |
| 108 | + } |
| 109 | + |
| 110 | + private static string CanonicalConfig(string rootPath) |
| 111 | + { |
| 112 | + using var document = JsonDocument.Parse(JsonSerializer.Serialize(new { rootPath })); |
| 113 | + return StorageProfileRules.CanonicalJson(document.RootElement); |
| 114 | + } |
| 115 | + |
| 116 | + private static bool IsUniqueViolation(Exception exception) |
| 117 | + { |
| 118 | + for (Exception? current = exception; current != null; current = current.InnerException) |
| 119 | + if (current is PostgresException { SqlState: PostgresErrorCodes.UniqueViolation }) return true; |
| 120 | + return false; |
| 121 | + } |
| 122 | +} |
0 commit comments