cedd60ab45
ci / build-test (push) Successful in 1m17s
Three defects found reviewing the last few commits. Deriving consumption on ingest made the batch reading endpoint quadratic. A recompute rewrites a meter's entire consumption series, and POST /api/v1/readings ran one per reading -- 500 readings for one meter meant 500 full rewrites. IngestByMeterAsync takes renormalize:false and the endpoint normalizes each touched meter once after the batch. Percentage change divided by a possibly negative baseline. A net-export meter going from -100 to -150 exported half again as much and would have been reported as "+50%", reading as more consumption. A non-positive baseline now reports no basis rather than a confident lie. The data-protection key ring had no persistent home outside Docker Compose. The LXC installer now creates /var/lib/metervault/keys at 0700 -- the app would otherwise create it under the default umask, leaving a key ring world-readable -- and the Unraid template maps it, since without that every UI-entered secret was lost whenever the container was recreated. README documents the variable and the trust boundary: keys on disk protect against leaked database content, not against an attacker who already has the host. Claude-Session: https://claude.ai/code/session_01V6joyergfvVLFEizH1hJLd
261 lines
10 KiB
C#
261 lines
10 KiB
C#
using System.Text.Json;
|
||
using MeterVault.Core.Domain;
|
||
using MeterVault.Infrastructure.Ingestion;
|
||
using MeterVault.Infrastructure.Persistence;
|
||
using Microsoft.EntityFrameworkCore;
|
||
using Microsoft.Extensions.Logging.Abstractions;
|
||
|
||
namespace MeterVault.Integration.Tests.Ingestion;
|
||
|
||
[Collection("Timescale")]
|
||
public sealed class IngestionServiceTests(TimescaleFixture fx)
|
||
{
|
||
private static readonly DateTimeOffset T0 = new(2024, 1, 1, 0, 0, 0, TimeSpan.Zero);
|
||
|
||
[Fact]
|
||
public async Task Writes_and_updates_idempotently_with_scale_and_offset()
|
||
{
|
||
await using var db = fx.CreateContext();
|
||
var (meterId, sourceId) = await SetupAsync(db, MeterMode.CumulativeCounter, scale: 0.001, offset: 0);
|
||
var service = NewIngestion(db);
|
||
|
||
// 1000 raw × 0.001 = 1.0.
|
||
Assert.Equal(IngestionOutcome.Written, await service.IngestAsync(sourceId, T0, 1000));
|
||
Assert.Equal(IngestionOutcome.Updated, await service.IngestAsync(sourceId, T0, 2000)); // same time → update
|
||
|
||
var reading = await db.Readings.SingleAsync(r => r.MeterId == meterId && r.Time == T0);
|
||
Assert.Equal(2.0, reading.Value, 6);
|
||
|
||
var source = await db.MeterSources.SingleAsync(s => s.Id == sourceId);
|
||
Assert.Equal("ok", source.LastStatus);
|
||
Assert.Equal(2.0, source.LastValue!.Value, 6);
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
[Fact]
|
||
public async Task Rejects_spurious_decrease_on_a_cumulative_register()
|
||
{
|
||
await using var db = fx.CreateContext();
|
||
var (meterId, sourceId) = await SetupAsync(db, MeterMode.CumulativeCounter);
|
||
var service = NewIngestion(db);
|
||
|
||
await service.IngestAsync(sourceId, T0, 500);
|
||
var outcome = await service.IngestAsync(sourceId, T0.AddHours(1), 400); // decrease, no event
|
||
|
||
Assert.Equal(IngestionOutcome.RejectedDecrease, outcome);
|
||
Assert.False(await db.Readings.AnyAsync(r => r.MeterId == meterId && r.Time == T0.AddHours(1)));
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
[Fact]
|
||
public async Task Old_reset_does_not_permanently_disable_the_decrease_guard()
|
||
{
|
||
await using var db = fx.CreateContext();
|
||
var (meterId, sourceId) = await SetupAsync(db, MeterMode.CumulativeCounter);
|
||
var service = NewIngestion(db);
|
||
|
||
// A reset early on explains an early decrease...
|
||
await service.IngestAsync(sourceId, T0, 100);
|
||
db.MeterEvents.Add(new MeterEvent { MeterId = meterId, Time = T0.AddMinutes(10), EventType = MeterEventType.CounterReset, NewValue = 0 });
|
||
await db.SaveChangesAsync();
|
||
Assert.Equal(IngestionOutcome.Written, await service.IngestAsync(sourceId, T0.AddMinutes(20), 30));
|
||
await service.IngestAsync(sourceId, T0.AddHours(1), 200);
|
||
|
||
// ...but a later spurious decrease with NO event in its window must still be rejected.
|
||
var outcome = await service.IngestAsync(sourceId, T0.AddHours(2), 150);
|
||
|
||
Assert.Equal(IngestionOutcome.RejectedDecrease, outcome);
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
[Fact]
|
||
public async Task Allows_decrease_when_a_swap_event_explains_it()
|
||
{
|
||
await using var db = fx.CreateContext();
|
||
var (meterId, sourceId) = await SetupAsync(db, MeterMode.CumulativeCounter);
|
||
var service = NewIngestion(db);
|
||
|
||
await service.IngestAsync(sourceId, T0, 500);
|
||
db.MeterEvents.Add(new MeterEvent
|
||
{
|
||
MeterId = meterId,
|
||
Time = T0.AddMinutes(30),
|
||
EventType = MeterEventType.MeterSwap,
|
||
PrevValue = 500,
|
||
NewValue = 0,
|
||
});
|
||
await db.SaveChangesAsync();
|
||
|
||
var outcome = await service.IngestAsync(sourceId, T0.AddHours(1), 20); // new meter reads low
|
||
|
||
Assert.Equal(IngestionOutcome.Written, outcome);
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
[Fact]
|
||
public async Task Mqtt_router_ingests_a_tasmota_payload()
|
||
{
|
||
await using var db = fx.CreateContext();
|
||
var brokerId = await CreateBrokerAsync(db);
|
||
var (meterId, _) = await SetupAsync(
|
||
db, MeterMode.CumulativeCounter, topic: "tele/plug7/SENSOR", endpointId: brokerId);
|
||
var router = new MqttMessageRouter(db, NewIngestion(db), NullLogger<MqttMessageRouter>.Instance);
|
||
|
||
var routed = await router.RouteAsync(
|
||
brokerId,
|
||
"tele/plug7/SENSOR",
|
||
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}""");
|
||
|
||
Assert.Equal(1, routed);
|
||
var reading = await db.Readings.SingleAsync(r => r.MeterId == meterId);
|
||
Assert.Equal(8421.0, reading.Value, 3);
|
||
Assert.Equal(new DateTimeOffset(2024, 3, 1, 10, 0, 0, TimeSpan.Zero), reading.Time);
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
[Fact]
|
||
public async Task Mqtt_router_ignores_a_source_bound_to_another_broker()
|
||
{
|
||
await using var db = fx.CreateContext();
|
||
var brokerA = await CreateBrokerAsync(db);
|
||
var brokerB = await CreateBrokerAsync(db);
|
||
|
||
// Topic filter that both brokers' traffic would match — the binding is the only thing
|
||
// separating them.
|
||
var (meterId, _) = await SetupAsync(
|
||
db, MeterMode.CumulativeCounter, topic: "tele/+/SENSOR", endpointId: brokerB);
|
||
var router = new MqttMessageRouter(db, NewIngestion(db), NullLogger<MqttMessageRouter>.Instance);
|
||
|
||
var routed = await router.RouteAsync(
|
||
brokerA,
|
||
"tele/plug7/SENSOR",
|
||
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}""");
|
||
|
||
Assert.Equal(0, routed);
|
||
Assert.False(await db.Readings.AnyAsync(r => r.MeterId == meterId));
|
||
|
||
// Same message on the broker it is actually bound to does land.
|
||
Assert.Equal(1, await router.RouteAsync(
|
||
brokerB,
|
||
"tele/plug7/SENSOR",
|
||
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}"""));
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
[Fact]
|
||
public async Task Ingesting_a_reading_derives_consumption_without_a_separate_recompute()
|
||
{
|
||
// Regression: live ingestion used to write only the raw reading, so consumption/generation
|
||
// stayed frozen at the last import until something else recomputed the meter.
|
||
await using var db = fx.CreateContext();
|
||
var (meterId, sourceId) = await SetupAsync(db, MeterMode.CumulativeCounter);
|
||
var service = NewIngestion(db);
|
||
|
||
await service.IngestAsync(sourceId, T0, 1000);
|
||
await service.IngestAsync(sourceId, T0.AddHours(1), 1250);
|
||
|
||
var consumption = await db.Consumption.AsNoTracking()
|
||
.Where(c => c.MeterId == meterId)
|
||
.OrderBy(c => c.Time)
|
||
.ToListAsync();
|
||
|
||
// The first reading is anchored against the meter's baseline (0), so it contributes 1000;
|
||
// what proves the fix is the second reading's 250 delta being there at all.
|
||
Assert.Equal(2, consumption.Count);
|
||
Assert.Equal(250d, consumption.Single(c => c.Time == T0.AddHours(1)).Amount, 3);
|
||
Assert.Equal(1250d, consumption.Sum(c => c.Amount), 3);
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
[Fact]
|
||
public async Task A_batch_can_defer_normalization_and_derive_the_same_series_once_at_the_end()
|
||
{
|
||
// Recomputing rewrites a meter's whole consumption series, so the batch endpoint skips it
|
||
// per reading and does it once. The result must be identical to normalizing as it goes.
|
||
await using var db = fx.CreateContext();
|
||
var (meterId, _) = await SetupAsync(db, MeterMode.CumulativeCounter);
|
||
var service = NewIngestion(db);
|
||
|
||
for (var hour = 0; hour < 5; hour++)
|
||
{
|
||
await service.IngestByMeterAsync(meterId, T0.AddHours(hour), 1000 + (hour * 10), renormalize: false);
|
||
}
|
||
|
||
Assert.False(await db.Consumption.AnyAsync(c => c.MeterId == meterId));
|
||
|
||
await service.RenormalizeMeterAsync(meterId);
|
||
|
||
var consumption = await db.Consumption.AsNoTracking().Where(c => c.MeterId == meterId).ToListAsync();
|
||
Assert.Equal(5, consumption.Count);
|
||
Assert.Equal(1040d, consumption.Sum(c => c.Amount), 3); // baseline 0 → 1000, then 4 × 10
|
||
|
||
await CleanupAsync(db, meterId);
|
||
}
|
||
|
||
private static IngestionService NewIngestion(MeterVaultDbContext db) =>
|
||
new(db, new MeterVault.Infrastructure.Normalization.NormalizationService(
|
||
db, MeterVault.Core.Normalization.NormalizationEngine.CreateDefault()));
|
||
|
||
private static async Task<(int MeterId, int SourceId)> SetupAsync(
|
||
MeterVaultDbContext db, MeterMode mode, double scale = 1, double offset = 0,
|
||
string topic = "tele/x/SENSOR", string? path = "ENERGY.Total", int? endpointId = null)
|
||
{
|
||
await DatabaseSeeder.SeedAsync(db);
|
||
var type = await db.EnergyTypes.FirstAsync(t => t.Key == "electricity");
|
||
|
||
var meter = new Meter
|
||
{
|
||
Name = $"ingest-{Guid.NewGuid():N}",
|
||
EnergyTypeId = type.Id,
|
||
Mode = mode,
|
||
Unit = "kWh",
|
||
};
|
||
db.Meters.Add(meter);
|
||
await db.SaveChangesAsync();
|
||
|
||
var source = new MeterSource
|
||
{
|
||
MeterId = meter.Id,
|
||
SourceType = SourceType.Tasmota,
|
||
EndpointId = endpointId ?? await CreateBrokerAsync(db),
|
||
ValueKind = SourceValueKind.Register,
|
||
Scale = scale,
|
||
Offset = offset,
|
||
Config = JsonSerializer.Serialize(new { topic, path }),
|
||
};
|
||
db.MeterSources.Add(source);
|
||
await db.SaveChangesAsync();
|
||
|
||
return (meter.Id, source.Id);
|
||
}
|
||
|
||
private static async Task<int> CreateBrokerAsync(MeterVaultDbContext db)
|
||
{
|
||
var endpoint = new IngestionEndpoint
|
||
{
|
||
Type = EndpointType.MqttBroker,
|
||
Name = $"broker-{Guid.NewGuid():N}",
|
||
Config = """{"host":"localhost","port":1883}""",
|
||
};
|
||
db.IngestionEndpoints.Add(endpoint);
|
||
await db.SaveChangesAsync();
|
||
return endpoint.Id;
|
||
}
|
||
|
||
private static async Task CleanupAsync(MeterVaultDbContext db, int meterId)
|
||
{
|
||
await db.Readings.Where(r => r.MeterId == meterId).ExecuteDeleteAsync();
|
||
await db.MeterEvents.Where(e => e.MeterId == meterId).ExecuteDeleteAsync();
|
||
await db.Consumption.Where(c => c.MeterId == meterId).ExecuteDeleteAsync();
|
||
await db.Meters.Where(m => m.Id == meterId).ExecuteDeleteAsync();
|
||
await db.IngestionEndpoints.Where(e => e.Name.StartsWith("broker-")).ExecuteDeleteAsync();
|
||
}
|
||
}
|