using MeterVault.Core.Domain;
using MeterVault.Core.Normalization;
using MeterVault.Infrastructure.Import;
using MeterVault.Infrastructure.Normalization;
using MeterVault.Infrastructure.Persistence;
using Microsoft.EntityFrameworkCore;
using static MeterVault.Integration.Tests.Reconciliation.ReconciliationSupport;
namespace MeterVault.Integration.Tests;
///
/// End-to-end persistence: stage the water CSV → commit as a batch → normalized consumption lands
/// in the hypertable → revert removes everything and rebases consumption (SDD §6.3, FR-6).
///
[Collection("Timescale")]
public sealed class ImportRoundTripTests(TimescaleFixture fx)
{
[Fact]
public async Task Commit_persists_consumption_and_revert_removes_it()
{
await using var db = fx.CreateContext();
await DatabaseSeeder.SeedAsync(db);
var waterType = await db.EnergyTypes.FirstAsync(t => t.Key == "water");
var meter = new Meter
{
Name = "Zähler Wasser (round-trip test)",
EnergyTypeId = waterType.Id,
Mode = MeterMode.CumulativeCounter,
Unit = "m3",
InitialBaseline = 820,
};
db.Meters.Add(meter);
await db.SaveChangesAsync();
var profile = new MappingProfile
{
Name = "test-water",
DateColumn = 0,
DateKind = DateKind.MonthName,
FirstDataRowIndex = 1,
DetectCumulativeSwaps = true,
Columns =
[
new ColumnMapping
{
Index = 1,
Role = MappingRole.Reading,
MeterId = meter.Id,
Unit = "m3",
SwapConsumptionColumn = 2,
},
],
};
StagedImport staged;
using (var reader = new StreamReader(FixturePath(Water)))
{
staged = new CsvImporter().Stage(profile, reader);
}
var normalization = new NormalizationService(db, NormalizationEngine.CreateDefault());
var service = new ImportService(db, normalization);
// Commit.
var batchId = await service.CommitAsync(staged, sourceName: "Wasser.csv", mappingJson: null);
Assert.True(await db.Readings.AnyAsync(r => r.MeterId == meter.Id));
Assert.True(await db.MeterEvents.AnyAsync(e => e.MeterId == meter.Id && e.EventType == MeterEventType.MeterSwap));
var maerz = await db.Consumption.SingleAsync(c =>
c.MeterId == meter.Id && c.Time == new DateTimeOffset(2023, 3, 1, 0, 0, 0, TimeSpan.Zero));
Assert.Equal(12d, maerz.Amount, 3);
// Revert.
await service.RevertAsync(batchId);
Assert.False(await db.Readings.AnyAsync(r => r.MeterId == meter.Id));
Assert.False(await db.MeterEvents.AnyAsync(e => e.MeterId == meter.Id));
Assert.False(await db.Consumption.AnyAsync(c => c.MeterId == meter.Id));
var batch = await db.ImportBatches.SingleAsync(b => b.Id == batchId);
Assert.NotNull(batch.RevertedAt);
// Cleanup so the shared container stays tidy for other tests. ExecuteDelete bypasses the
// change tracker (which holds stale entries after the revert's set-based deletes).
await db.Meters.Where(m => m.Id == meter.Id).ExecuteDeleteAsync();
}
[Fact]
public async Task Re_importing_the_same_file_is_refused_by_meter_and_date()
{
// The commonest real duplicate: import a file, then import an overlapping one. The in-batch
// guard cannot see it (that set is internally unique), so left to the database it surfaced as
// EF's "An error occurred while saving the entity changes", naming neither meter nor date.
await using var db = fx.CreateContext();
await DatabaseSeeder.SeedAsync(db);
var type = await db.EnergyTypes.FirstAsync(t => t.Key == "electricity");
var meter = new Meter
{
Name = $"reimport-{Guid.NewGuid():N}",
EnergyTypeId = type.Id,
Mode = MeterMode.CumulativeCounter,
Unit = "kWh",
};
db.Meters.Add(meter);
await db.SaveChangesAsync();
var service = new ImportService(db, new NormalizationService(db, NormalizationEngine.CreateDefault()));
StagedImport Stage() => Staged(meter.Id, new DateTimeOffset(2024, 5, 1, 0, 0, 0, TimeSpan.Zero), 1200);
await service.CommitAsync(Stage(), "first.csv", mappingJson: null);
var error = await Assert.ThrowsAsync(
() => service.CommitAsync(Stage(), "again.csv", mappingJson: null));
Assert.Contains("already exist", error.Message, StringComparison.OrdinalIgnoreCase);
Assert.Contains("2024-05-01", error.Message, StringComparison.Ordinal);
Assert.Contains($"meter {meter.Id}", error.Message, StringComparison.Ordinal);
await db.Consumption.Where(c => c.MeterId == meter.Id).ExecuteDeleteAsync();
await db.Readings.Where(r => r.MeterId == meter.Id).ExecuteDeleteAsync();
await db.Meters.Where(m => m.Id == meter.Id).ExecuteDeleteAsync();
}
private static StagedImport Staged(int meterId, DateTimeOffset time, double value)
{
var staged = new StagedImport();
staged.Readings.Add(new Reading
{
MeterId = meterId,
Time = time,
Value = value,
Quality = ReadingQuality.Imported,
});
return staged;
}
}