test(storage): isolate PostgreSQL lifecycle coverage

This commit is contained in:
joshpatra committed 2026-07-27 08:43:58 -04:00
1 parent 43f32d93ba
commit 96b556b446
6 files changed
+168 -117

No files matched your search

@@ -64,6 +64,24 @@ public sealed class DurableStorageRuntimeProbeTests : IAsyncLifetime
Assert.Equal(2, factory.CreateCount);
}
[Fact]
public async Task ProbePreservesCallerCancellation()
{
var options = RuntimeOptions();
var state = new DurableStorageState(options);
using var probe = new DurableStorageRuntimeProbe(
new CountingDbContextFactory(_database.Options),
options,
state,
_clock,
NullLogger<DurableStorageRuntimeProbe>.Instance);
using var cancellation = new CancellationTokenSource();
cancellation.Cancel();
await Assert.ThrowsAnyAsync<OperationCanceledException>(
() => probe.CheckAsync(cancellation.Token));
}
[Fact]
public async Task ProbeMarksRuntimeSchemaDriftUnready()
{
+56
View File
@@ -145,6 +145,54 @@ public sealed class DurableStorageTests : IAsyncLifetime
Assert.True(await ColumnExists(context, "backups", "RestoreStatus"));
}
[Fact]
public async Task ProbeCacheSnapshot_AppliesToFreshPostgres()
{
await using var context = await Factory().CreateDbContextAsync();
await context.Database.MigrateAsync();
Assert.Contains(
"20260724012448_ProbeCacheSnapshot",
await context.Database.GetAppliedMigrationsAsync());
Assert.Equal("boolean", await ColumnType(context, "playlist_links", "Enabled"));
}
[Fact]
public async Task ProbeCacheSnapshot_UpgradesLegacyIntegerEnabledColumn()
{
await using var context = await Factory().CreateDbContextAsync();
var migrator = context.GetService<IMigrator>();
await migrator.MigrateAsync("20260723233918_AddDownloadArtifactMediaFacts");
await context.Database.ExecuteSqlRawAsync("""
ALTER TABLE playlist_links
ALTER COLUMN "Enabled" DROP DEFAULT,
ALTER COLUMN "Enabled" TYPE integer USING (CASE WHEN "Enabled" THEN 1 ELSE 0 END),
ALTER COLUMN "Enabled" SET DEFAULT 1
""");
await migrator.MigrateAsync();
Assert.Equal("boolean", await ColumnType(context, "playlist_links", "Enabled"));
}
[Fact]
public async Task SchemaCompatibility_RejectsCaseDivergentMigrationId()
{
await using var context = await Factory().CreateDbContextAsync();
await context.Database.MigrateAsync();
const string migration = "20260724012448_ProbeCacheSnapshot";
var divergent = migration.ToLowerInvariant();
await context.Database.ExecuteSqlInterpolatedAsync(
$"UPDATE \"__EFMigrationsHistory\" SET \"MigrationId\" = {divergent} WHERE \"MigrationId\" = {migration}");
var compatibility = await DurableSchemaCompatibility.InspectAsync(context);
Assert.Equal(DurableSchemaCompatibilityStatus.UnsupportedVersion, compatibility.Status);
Assert.Contains(migration, compatibility.MissingMigrations);
Assert.Contains(divergent, compatibility.UnknownMigrations);
}
[Fact]
public async Task OnboardingMigration_BackfillsExistingBackendIdentity()
{
@@ -262,6 +310,14 @@ public sealed class DurableStorageTests : IAsyncLifetime
$"SELECT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema = 'public' AND table_name = {table} AND column_name = {column}) AS \"Value\"")
.SingleAsync();
private static async Task<string> ColumnType(
AllstarrDbContext context,
string table,
string column) =>
await context.Database.SqlQuery<string>(
$"SELECT data_type AS \"Value\" FROM information_schema.columns WHERE table_schema = 'public' AND table_name = {table} AND column_name = {column}")
.SingleAsync();
public async Task DisposeAsync()
{
await _database.DisposeAsync();
+20 -113
View File
@@ -20,7 +20,6 @@ using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.FileProviders;
using Npgsql;
namespace allstarr.Tests;
@@ -30,14 +29,8 @@ public sealed class PostgresStorageIntegrationTests
[Trait("Category", "Postgres")]
public async Task NativePostgresLineageConstraints_RejectCrossTenantFavoriteJob()
{
var connectionString = Environment.GetEnvironmentVariable("ALLSTARR_TEST_POSTGRES");
if (string.IsNullOrWhiteSpace(connectionString)) return;
var options = new DbContextOptionsBuilder<AllstarrDbContext>()
.UseNpgsql(connectionString)
.Options;
await using var db = new AllstarrDbContext(options);
await db.Database.ExecuteSqlRawAsync("DROP SCHEMA IF EXISTS public CASCADE; CREATE SCHEMA public");
await using var database = await PostgresTestDatabase.CreateAsync();
await using var db = new AllstarrDbContext(database.Options);
await db.Database.MigrateAsync();
var now = DateTimeOffset.UtcNow;
@@ -90,14 +83,11 @@ public sealed class PostgresStorageIntegrationTests
[Trait("Category", "Postgres")]
public async Task NativePostgresHostOptions_SupportIdentityJobAndOutboxTransactions()
{
var connectionString = Environment.GetEnvironmentVariable("ALLSTARR_TEST_POSTGRES");
if (string.IsNullOrWhiteSpace(connectionString)) return;
await using var services = BuildHostStorageServices(connectionString);
await using var database = await PostgresTestDatabase.CreateAsync();
await using var services = BuildHostStorageServices(database.ConnectionString);
var factory = services.GetRequiredService<IDbContextFactory<AllstarrDbContext>>();
await using (var db = await factory.CreateDbContextAsync())
{
await db.Database.ExecuteSqlRawAsync("DROP SCHEMA IF EXISTS public CASCADE; CREATE SCHEMA public");
await db.Database.MigrateAsync();
}
@@ -290,17 +280,12 @@ public sealed class PostgresStorageIntegrationTests
[Trait("Category", "Postgres")]
public async Task NativePostgresLegacyEnvMigration_AtomicallyAppliesAndDecryptsSharedAccount()
{
var connectionString = Environment.GetEnvironmentVariable("ALLSTARR_TEST_POSTGRES");
if (string.IsNullOrWhiteSpace(connectionString))
{
return;
}
await using var database = await PostgresTestDatabase.CreateAsync();
var root = Path.Combine(Path.GetTempPath(), "allstarr-tests", Guid.NewGuid().ToString("N"));
Directory.CreateDirectory(root);
try
{
await using var hostServices = BuildHostStorageServices(connectionString);
await using var hostServices = BuildHostStorageServices(database.ConnectionString);
var factory = hostServices.GetRequiredService<IDbContextFactory<AllstarrDbContext>>();
await using (var strategyContext = await factory.CreateDbContextAsync())
{
@@ -310,8 +295,6 @@ public sealed class PostgresStorageIntegrationTests
var userId = Guid.CreateVersion7();
await using (var db = await factory.CreateDbContextAsync())
{
await db.Database.ExecuteSqlRawAsync(
"DROP SCHEMA IF EXISTS public CASCADE; CREATE SCHEMA public");
await db.Database.MigrateAsync();
db.Tenants.Add(new TenantRecord
{
@@ -465,18 +448,8 @@ public sealed class PostgresStorageIntegrationTests
[Trait("Category", "Postgres")]
public async Task NativePostgresAdditiveMigrations_CanRollBackToFoundationAndReapply()
{
var connectionString = Environment.GetEnvironmentVariable("ALLSTARR_TEST_POSTGRES");
if (string.IsNullOrWhiteSpace(connectionString))
{
return;
}
var dbOptions = new DbContextOptionsBuilder<AllstarrDbContext>()
.UseNpgsql(connectionString)
.Options;
await using var context = new AllstarrDbContext(dbOptions);
await context.Database.ExecuteSqlRawAsync(
"DROP SCHEMA IF EXISTS public CASCADE; CREATE SCHEMA public");
await using var database = await PostgresTestDatabase.CreateAsync();
await using var context = new AllstarrDbContext(database.Options);
var migrator = context.GetService<IMigrator>();
await migrator.MigrateAsync();
@@ -511,16 +484,8 @@ public sealed class PostgresStorageIntegrationTests
[Trait("Category", "Postgres")]
public async Task NativePostgresMigrationLockAndDurableQueue_WorkAgainstSelectedDatabase()
{
var connectionString = Environment.GetEnvironmentVariable("ALLSTARR_TEST_POSTGRES");
if (string.IsNullOrWhiteSpace(connectionString))
{
return;
}
var dbOptions = new DbContextOptionsBuilder<AllstarrDbContext>()
.UseNpgsql(connectionString)
.Options;
var factory = new TestDbContextFactory(dbOptions);
await using var database = await PostgresTestDatabase.CreateAsync();
var factory = new TestDbContextFactory(database.Options);
await using (var reset = await factory.CreateDbContextAsync())
{
await reset.Database.ExecuteSqlRawAsync(
@@ -539,7 +504,7 @@ public sealed class PostgresStorageIntegrationTests
var options = new DurableStorageOptions
{
Provider = "Postgres",
ConnectionString = connectionString,
ConnectionString = database.ConnectionString,
AutoMigrate = true,
ConnectionRetryCount = 0,
BackupDirectory = Path.Combine(Path.GetTempPath(), "allstarr-postgres-backups")
@@ -634,16 +599,8 @@ public sealed class PostgresStorageIntegrationTests
[Trait("Category", "Postgres")]
public async Task NativePostgresCacheLoss_PreservesDurableWorkAndProgressAcrossCacheRestart()
{
var connectionString = Environment.GetEnvironmentVariable("ALLSTARR_TEST_POSTGRES");
if (string.IsNullOrWhiteSpace(connectionString))
{
return;
}
var dbOptions = new DbContextOptionsBuilder<AllstarrDbContext>()
.UseNpgsql(connectionString)
.Options;
var factory = new TestDbContextFactory(dbOptions);
await using var database = await PostgresTestDatabase.CreateAsync();
var factory = new TestDbContextFactory(database.Options);
await using (var reset = await factory.CreateDbContextAsync())
{
await reset.Database.ExecuteSqlRawAsync(
@@ -860,14 +817,8 @@ public sealed class PostgresStorageIntegrationTests
[Trait("Category", "Postgres")]
public async Task NativePostgresBackup_VerifiesAndRestoresIntoIsolatedDatabase()
{
var connectionString = Environment.GetEnvironmentVariable("ALLSTARR_TEST_POSTGRES");
if (string.IsNullOrWhiteSpace(connectionString))
{
return;
}
var sourceBuilder = new NpgsqlConnectionStringBuilder(connectionString);
var targetDatabase = $"allstarr_restore_{Guid.NewGuid():N}";
await using var sourceDatabase = await PostgresTestDatabase.CreateAsync();
await using var targetDatabase = await PostgresTestDatabase.CreateAsync();
var backupRoot = Path.Combine(
Path.GetTempPath(),
"allstarr-tests",
@@ -875,14 +826,9 @@ public sealed class PostgresStorageIntegrationTests
Directory.CreateDirectory(backupRoot);
try
{
var dbOptions = new DbContextOptionsBuilder<AllstarrDbContext>()
.UseNpgsql(connectionString)
.Options;
var factory = new TestDbContextFactory(dbOptions);
var factory = new TestDbContextFactory(sourceDatabase.Options);
await using (var reset = await factory.CreateDbContextAsync())
{
await reset.Database.ExecuteSqlRawAsync(
"DROP SCHEMA IF EXISTS public CASCADE; CREATE SCHEMA public");
await reset.Database.MigrateAsync();
reset.Jobs.Add(new DurableJobRecord
{
@@ -903,7 +849,7 @@ public sealed class PostgresStorageIntegrationTests
var options = new DurableStorageOptions
{
Provider = "Postgres",
ConnectionString = connectionString,
ConnectionString = sourceDatabase.ConnectionString,
BackupDirectory = backupRoot
};
var state = new DurableStorageState(options);
@@ -917,20 +863,15 @@ public sealed class PostgresStorageIntegrationTests
var artifact = await service.CreateAsync();
Assert.True(File.Exists(artifact.ArtifactPath));
Assert.True(File.Exists(artifact.ManifestPath));
await CreateDatabase(sourceBuilder, targetDatabase);
var targetBuilder = new NpgsqlConnectionStringBuilder(connectionString)
{
Database = targetDatabase
};
await service.RestorePostgresAsync(
artifact,
targetBuilder.ConnectionString,
targetDatabase.ConnectionString,
destructiveRestoreConfirmed: true,
isolatedTargetDatabaseConfirmation: targetDatabase);
isolatedTargetDatabaseConfirmation: targetDatabase.DatabaseName);
var restoredOptions = new DbContextOptionsBuilder<AllstarrDbContext>()
.UseNpgsql(targetBuilder.ConnectionString)
.UseNpgsql(targetDatabase.ConnectionString)
.Options;
await using var restored = new AllstarrDbContext(restoredOptions);
var restoredJob = await restored.Jobs.AsNoTracking().SingleAsync();
@@ -939,7 +880,6 @@ public sealed class PostgresStorageIntegrationTests
}
finally
{
await DropDatabase(sourceBuilder, targetDatabase);
if (Directory.Exists(backupRoot))
{
Directory.Delete(backupRoot, recursive: true);
@@ -1023,39 +963,6 @@ public sealed class PostgresStorageIntegrationTests
return Convert.ToInt32(await command.ExecuteScalarAsync()) == 1;
}
private static async Task CreateDatabase(
NpgsqlConnectionStringBuilder source,
string database)
{
var admin = new NpgsqlConnectionStringBuilder(source.ConnectionString)
{
Database = "postgres"
};
await using var connection = new NpgsqlConnection(admin.ConnectionString);
await connection.OpenAsync();
await using var command = connection.CreateCommand();
command.CommandText = $"CREATE DATABASE {QuoteIdentifier(database)}";
await command.ExecuteNonQueryAsync();
}
private static async Task DropDatabase(
NpgsqlConnectionStringBuilder source,
string database)
{
var admin = new NpgsqlConnectionStringBuilder(source.ConnectionString)
{
Database = "postgres"
};
await using var connection = new NpgsqlConnection(admin.ConnectionString);
await connection.OpenAsync();
await using var command = connection.CreateCommand();
command.CommandText = $"DROP DATABASE IF EXISTS {QuoteIdentifier(database)} WITH (FORCE)";
await command.ExecuteNonQueryAsync();
}
private static string QuoteIdentifier(string identifier) =>
'"' + identifier.Replace("\"", "\"\"", StringComparison.Ordinal) + '"';
private sealed class TestDbContextFactory(DbContextOptions<AllstarrDbContext> options)
: IDbContextFactory<AllstarrDbContext>
{
+68 -2
View File
@@ -371,6 +371,60 @@ public sealed class SidecarReadinessTests : IAsyncLifetime
Assert.DoesNotContain(_root, System.Text.Json.JsonSerializer.Serialize(snapshot), StringComparison.Ordinal);
}
[Fact]
public async Task ReadinessPreservesCallerCancellationFromSecretProbe()
{
var keyRingPath = Path.Combine(_root, "keyring.json");
await File.WriteAllTextAsync(keyRingPath, JsonSerializer.Serialize(new
{
activeKeyId = "active",
keys = new Dictionary<string, string>
{
["active"] = Convert.ToBase64String(RandomNumberGenerator.GetBytes(32))
}
}));
if (!OperatingSystem.IsWindows())
{
File.SetUnixFileMode(keyRingPath, UnixFileMode.UserRead | UnixFileMode.UserWrite);
}
using var cancellation = new CancellationTokenSource();
cancellation.Cancel();
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => Readiness(
new SidecarStatusCatalog(new SidecarHealthOptions()),
new ReadinessOptions { RequireSecretKeyRing = true },
keyRingPath: keyRingPath).CheckAsync(cancellation.Token));
}
[Fact]
public async Task SidecarProbePreservesCallerCancellation()
{
var options = new SidecarHealthOptions
{
Targets =
[
new SidecarProbeTarget
{
Id = "blocking",
ProviderId = "fixture",
BaseUrl = "http://sidecar.test/"
}
]
};
var monitor = new SidecarHealthMonitor(
new HandlerFactory(new BlockingHandler()),
options,
new SidecarStatusCatalog(options),
_health,
NullLogger<SidecarHealthMonitor>.Instance,
_clock);
using var cancellation = new CancellationTokenSource();
cancellation.Cancel();
await Assert.ThrowsAnyAsync<OperationCanceledException>(
() => monitor.ProbeAllOnceAsync(cancellation.Token));
}
[Fact]
public async Task RequiredDirectory_MustExistAndAcceptAWriteProbeWithoutExposingItsPath()
{
@@ -404,11 +458,12 @@ public sealed class SidecarReadinessTests : IAsyncLifetime
private PlatformReadinessService Readiness(
SidecarStatusCatalog catalog,
ReadinessOptions options,
IDurableStorageRuntimeProbe? storageProbe = null)
IDurableStorageRuntimeProbe? storageProbe = null,
string? keyRingPath = null)
{
var secretOptions = new SecretStoreOptions
{
KeyRingPath = Path.Combine(_root, "missing-keyring.json")
KeyRingPath = keyRingPath ?? Path.Combine(_root, "missing-keyring.json")
};
return new PlatformReadinessService(
_storageState,
@@ -455,6 +510,17 @@ public sealed class SidecarReadinessTests : IAsyncLifetime
public HttpClient CreateClient(string name) => new(handler, disposeHandler: false);
}
private sealed class BlockingHandler : HttpMessageHandler
{
protected override async Task<HttpResponseMessage> SendAsync(
HttpRequestMessage request,
CancellationToken cancellationToken)
{
await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken);
return new HttpResponseMessage(HttpStatusCode.OK);
}
}
private sealed class FakeClock(DateTimeOffset now) : IPlatformClock
{
public DateTimeOffset UtcNow { get; private set; } = now;
@@ -81,6 +81,10 @@ public sealed class PlatformReadinessService
components.Add(new ReadinessComponent("secret-key-ring", "ready", true));
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
throw;
}
catch
{
components.Add(new ReadinessComponent(
@@ -35,8 +35,8 @@ public static class DurableSchemaCompatibility
}
var applied = (await context.Database.GetAppliedMigrationsAsync(cancellationToken)).ToArray();
var knownSet = known.ToHashSet(StringComparer.OrdinalIgnoreCase);
var appliedSet = applied.ToHashSet(StringComparer.OrdinalIgnoreCase);
var knownSet = known.ToHashSet(StringComparer.Ordinal);
var appliedSet = applied.ToHashSet(StringComparer.Ordinal);
var unknown = applied.Where(migration => !knownSet.Contains(migration)).ToArray();
var missing = known.Where(migration => !appliedSet.Contains(migration)).ToArray();
var status = unknown.Length > 0