| @s.SourceType |
@SourceTarget(s) |
+
+ @{ var problem = ConnectorProblem(s); }
+ @if (problem is null)
+ {
+ @(_endpoints.FirstOrDefault(e => e.Id == s.EndpointId)?.Name ?? "—")
+ }
+ else
+ {
+
+ @problem
+
+ }
+ |
@(s.IsEnabled ? "yes" : "no") |
@(s.LastSeenAt?.ToString("yyyy-MM-dd HH:mm") ?? "—") |
@(s.LastValue is { } v ? Format.Number(v, 2) : "—") |
@@ -448,6 +470,12 @@ else
Snackbar.Add($"'{selected.Name}' is a {selected.Type} connector; a {_sourceEdit.SourceType} source needs {needed}.", Severity.Error);
return;
}
+
+ if (!selected.IsEnabled)
+ {
+ Snackbar.Add($"'{selected.Name}' is disabled, so this source would never ingest. Enable it first.", Severity.Error);
+ return;
+ }
}
else
{
@@ -527,8 +555,33 @@ else
_ => null,
};
+ // Only enabled connectors can ingest: both MQTT and HA workers filter on IsEnabled, so offering
+ // a disabled one would produce a source that saves cleanly and then never runs.
private List ConnectorsFor(EndpointType type) =>
- _endpoints.Where(e => e.Type == type).ToList();
+ _endpoints.Where(e => e.Type == type && e.IsEnabled).ToList();
+
+ ///
+ /// Why this source cannot ingest, or null if it can. Routing is endpoint-scoped, so an unbound
+ /// or mis-bound source is silently dead — and deleting a connector unlinks its sources, which
+ /// used to be harmless. Without this column such a source is indistinguishable from a healthy
+ /// one at "Enabled: yes".
+ ///
+ private string? ConnectorProblem(MeterSource source)
+ {
+ if (RequiredEndpointType(source.SourceType) is not { } needed)
+ {
+ return null;
+ }
+
+ var endpoint = _endpoints.FirstOrDefault(e => e.Id == source.EndpointId);
+ return endpoint switch
+ {
+ null => "no connector — never ingests",
+ { IsEnabled: false } => $"'{endpoint.Name}' is disabled",
+ _ when endpoint.Type != needed => $"'{endpoint.Name}' is {endpoint.Type}, needs {needed}",
+ _ => null,
+ };
+ }
// Changing the source type can invalidate the chosen connector (an HA connector cannot serve an
// MQTT source), so drop a selection that no longer fits rather than saving a mismatched pair.
diff --git a/src/Infrastructure/Dashboard/MeterPeriodService.cs b/src/Infrastructure/Dashboard/MeterPeriodService.cs
index be28b5f..529095f 100644
--- a/src/Infrastructure/Dashboard/MeterPeriodService.cs
+++ b/src/Infrastructure/Dashboard/MeterPeriodService.cs
@@ -43,6 +43,14 @@ public sealed class MeterPeriodService(
return null;
}
+ // Virtual meters are evaluated on read and only materialize into `consumption` when a cost
+ // category references them (SDD §14.1), so summing that table would report a confident zero
+ // for a meter that is working fine. Report nothing and let the page say why.
+ if (meter.Mode == MeterMode.Virtual)
+ {
+ return null;
+ }
+
var tz = ResolveTimeZone();
var today = DateOnly.FromDateTime(TimeZoneInfo.ConvertTime(DateTimeOffset.UtcNow, tz).Date);
diff --git a/src/Infrastructure/Import/ImportService.cs b/src/Infrastructure/Import/ImportService.cs
index 08a7a64..a6b0126 100644
--- a/src/Infrastructure/Import/ImportService.cs
+++ b/src/Infrastructure/Import/ImportService.cs
@@ -19,6 +19,7 @@ public sealed class ImportService(MeterVaultDbContext db, NormalizationService n
{
ArgumentNullException.ThrowIfNull(staged);
GuardDuplicateReadings(staged);
+ await GuardExistingReadingsAsync(staged, cancellationToken).ConfigureAwait(false);
await using var tx = await _db.Database.BeginTransactionAsync(cancellationToken).ConfigureAwait(false);
@@ -115,6 +116,50 @@ public sealed class ImportService(MeterVaultDbContext db, NormalizationService n
$"timestamp ({sample}). Check that no two mapped columns target the same meter.");
}
+ ///
+ /// Rejects readings that already exist in the database, which is what re-importing an overlapping
+ /// CSV produces — the commonest real duplicate, and one the in-batch guard cannot see because the
+ /// staged set is internally unique. Left to the database it surfaces as EF's
+ /// "An error occurred while saving the entity changes", naming neither meter nor date.
+ ///
+ ///
+ /// Checks one meter at a time, bounded by that meter's staged time range, so the query stays
+ /// proportional to the overlap rather than to the table.
+ ///
+ private async Task GuardExistingReadingsAsync(StagedImport staged, CancellationToken cancellationToken)
+ {
+ var clashes = new List();
+
+ foreach (var group in staged.Readings.GroupBy(r => r.MeterId))
+ {
+ var times = group.Select(r => r.Time).ToHashSet();
+ var from = times.Min();
+ var to = times.Max();
+
+ var existing = await _db.Readings.AsNoTracking()
+ .Where(r => r.MeterId == group.Key && r.Time >= from && r.Time <= to)
+ .Select(r => r.Time)
+ .ToListAsync(cancellationToken).ConfigureAwait(false);
+
+ foreach (var time in existing.Where(times.Contains).Take(3))
+ {
+ clashes.Add($"meter {group.Key} at {time:yyyy-MM-dd}");
+ }
+
+ if (clashes.Count >= 3)
+ {
+ break;
+ }
+ }
+
+ if (clashes.Count > 0)
+ {
+ throw new InvalidOperationException(
+ $"This import would overwrite readings that already exist ({string.Join("; ", clashes)}). "
+ + "Revert the earlier batch on the Import page first, or narrow the file's date range.");
+ }
+ }
+
private static IEnumerable AffectedMeters(StagedImport staged) =>
staged.Readings.Select(r => r.MeterId)
.Concat(staged.Events.Select(e => e.MeterId))
diff --git a/src/Infrastructure/Ingestion/IngestionService.cs b/src/Infrastructure/Ingestion/IngestionService.cs
index bac69c5..48b5955 100644
--- a/src/Infrastructure/Ingestion/IngestionService.cs
+++ b/src/Infrastructure/Ingestion/IngestionService.cs
@@ -100,11 +100,8 @@ public sealed class IngestionService(
/// Derives consumption for one meter after a batch of readings has been written. The public
/// counterpart to skipping renormalize on each individual ingest.
///
- public async Task RenormalizeMeterAsync(int meterId, CancellationToken cancellationToken = default)
- {
- await _normalization.RecomputeMeterAsync(meterId, batchId: null, cancellationToken).ConfigureAwait(false);
- await _db.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
- }
+ public Task RenormalizeMeterAsync(int meterId, CancellationToken cancellationToken = default) =>
+ RecomputeAtomicallyAsync(meterId, cancellationToken);
///
/// Derives consumption from the reading just written. Without this a live-ingested reading sits
@@ -130,8 +127,37 @@ public sealed class IngestionService(
// The reading must already be persisted: RecomputeMeterAsync re-reads the meter's readings
// from the database, so anything still pending in the change tracker would be missed.
+ await RecomputeAtomicallyAsync(meterId, cancellationToken).ConfigureAwait(false);
+ }
+
+ ///
+ /// Rebuilds a meter's consumption series as one atomic unit.
+ ///
+ ///
+ /// clears the series with
+ /// ExecuteDelete, which commits on its own when no transaction is ambient, and only then
+ /// adds the rebuilt rows. Without a transaction around both halves the meter has *no*
+ /// consumption in between: a dashboard read in that window reports zero, and a crash or a
+ /// cancelled request makes the loss permanent — for data the SDD treats as the long-term source
+ /// of truth (§5.5). Import and the events API already wrap their recomputes this way; live
+ /// ingestion was the path that did not.
+ ///
+ /// Respects an ambient transaction rather than nesting, so callers that already opened one keep
+ /// a single unit of work.
+ ///
+ private async Task RecomputeAtomicallyAsync(int meterId, CancellationToken cancellationToken)
+ {
+ if (_db.Database.CurrentTransaction is not null)
+ {
+ await _normalization.RecomputeMeterAsync(meterId, batchId: null, cancellationToken).ConfigureAwait(false);
+ await _db.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
+ return;
+ }
+
+ await using var tx = await _db.Database.BeginTransactionAsync(cancellationToken).ConfigureAwait(false);
await _normalization.RecomputeMeterAsync(meterId, batchId: null, cancellationToken).ConfigureAwait(false);
await _db.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
+ await tx.CommitAsync(cancellationToken).ConfigureAwait(false);
}
private async Task UpsertAsync(
diff --git a/src/Infrastructure/Persistence/Migrations/20260718174732_BindUnboundMqttSourcesToSoleEnabledBroker.Designer.cs b/src/Infrastructure/Persistence/Migrations/20260718174732_BindUnboundMqttSourcesToSoleEnabledBroker.Designer.cs
new file mode 100644
index 0000000..c479d5b
--- /dev/null
+++ b/src/Infrastructure/Persistence/Migrations/20260718174732_BindUnboundMqttSourcesToSoleEnabledBroker.Designer.cs
@@ -0,0 +1,931 @@
+//
+using System;
+using MeterVault.Infrastructure.Persistence;
+using Microsoft.EntityFrameworkCore;
+using Microsoft.EntityFrameworkCore.Infrastructure;
+using Microsoft.EntityFrameworkCore.Migrations;
+using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
+using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
+
+#nullable disable
+
+namespace MeterVault.Infrastructure.Persistence.Migrations
+{
+ [DbContext(typeof(MeterVaultDbContext))]
+ [Migration("20260718174732_BindUnboundMqttSourcesToSoleEnabledBroker")]
+ partial class BindUnboundMqttSourcesToSoleEnabledBroker
+ {
+ ///
+ protected override void BuildTargetModel(ModelBuilder modelBuilder)
+ {
+#pragma warning disable 612, 618
+ modelBuilder
+ .HasAnnotation("ProductVersion", "10.0.9")
+ .HasAnnotation("Relational:MaxIdentifierLength", 63);
+
+ NpgsqlModelBuilderExtensions.HasPostgresExtension(modelBuilder, "timescaledb");
+ NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
+
+ modelBuilder.Entity("MeterVault.Core.Domain.AppSetting", b =>
+ {
+ b.Property("Key")
+ .HasMaxLength(128)
+ .HasColumnType("character varying(128)")
+ .HasColumnName("key");
+
+ b.Property("Value")
+ .IsRequired()
+ .HasColumnType("jsonb")
+ .HasColumnName("value");
+
+ b.HasKey("Key")
+ .HasName("pk_app_setting");
+
+ b.ToTable("app_setting", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Consumption", b =>
+ {
+ b.Property("MeterId")
+ .HasColumnType("integer")
+ .HasColumnName("meter_id");
+
+ b.Property("Time")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("time");
+
+ b.Property("Kind")
+ .HasColumnType("smallint")
+ .HasColumnName("kind");
+
+ b.Property("Amount")
+ .HasColumnType("double precision")
+ .HasColumnName("amount");
+
+ b.Property("ImportBatchId")
+ .HasColumnType("integer")
+ .HasColumnName("import_batch_id");
+
+ b.Property("Quality")
+ .HasColumnType("smallint")
+ .HasColumnName("quality");
+
+ b.HasKey("MeterId", "Time", "Kind")
+ .HasName("pk_consumption");
+
+ b.HasIndex("ImportBatchId")
+ .HasDatabaseName("ix_consumption_import_batch_id");
+
+ b.ToTable("consumption", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.CostCategory", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("ColorHex")
+ .HasColumnType("text")
+ .HasColumnName("color_hex");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(128)
+ .HasColumnType("character varying(128)")
+ .HasColumnName("name");
+
+ b.Property("Sort")
+ .HasColumnType("integer")
+ .HasColumnName("sort");
+
+ b.HasKey("Id")
+ .HasName("pk_cost_category");
+
+ b.ToTable("cost_category", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.CostCategoryMember", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("CategoryId")
+ .HasColumnType("integer")
+ .HasColumnName("category_id");
+
+ b.Property("EnergyTypeId")
+ .HasColumnType("smallint")
+ .HasColumnName("energy_type_id");
+
+ b.Property("MeterId")
+ .HasColumnType("integer")
+ .HasColumnName("meter_id");
+
+ b.HasKey("Id")
+ .HasName("pk_cost_category_member");
+
+ b.HasIndex("CategoryId")
+ .HasDatabaseName("ix_cost_category_member_category_id");
+
+ b.HasIndex("EnergyTypeId")
+ .HasDatabaseName("ix_cost_category_member_energy_type_id");
+
+ b.HasIndex("MeterId")
+ .HasDatabaseName("ix_cost_category_member_meter_id");
+
+ b.ToTable("cost_category_member", null, t =>
+ {
+ t.HasCheckConstraint("ck_cost_category_member_target", "meter_id IS NOT NULL OR energy_type_id IS NOT NULL");
+ });
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.EnergyType", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("smallint")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("BaseUnit")
+ .IsRequired()
+ .HasMaxLength(16)
+ .HasColumnType("character varying(16)")
+ .HasColumnName("base_unit");
+
+ b.Property("ColorHex")
+ .HasColumnType("text")
+ .HasColumnName("color_hex");
+
+ b.Property("CreatedAt")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("created_at")
+ .HasDefaultValueSql("now()");
+
+ b.Property("DefaultMode")
+ .IsRequired()
+ .HasMaxLength(32)
+ .HasColumnType("character varying(32)")
+ .HasColumnName("default_mode");
+
+ b.Property("DisplayName")
+ .IsRequired()
+ .HasMaxLength(128)
+ .HasColumnType("character varying(128)")
+ .HasColumnName("display_name");
+
+ b.Property("Icon")
+ .HasColumnType("text")
+ .HasColumnName("icon");
+
+ b.Property("Key")
+ .IsRequired()
+ .HasMaxLength(64)
+ .HasColumnType("character varying(64)")
+ .HasColumnName("key");
+
+ b.HasKey("Id")
+ .HasName("pk_energy_type");
+
+ b.HasIndex("Key")
+ .IsUnique()
+ .HasDatabaseName("ix_energy_type_key");
+
+ b.ToTable("energy_type", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.ImportBatch", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("CreatedAt")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("created_at")
+ .HasDefaultValueSql("now()");
+
+ b.Property("Mapping")
+ .HasColumnType("jsonb")
+ .HasColumnName("mapping");
+
+ b.Property("RevertedAt")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("reverted_at");
+
+ b.Property("RowCount")
+ .HasColumnType("integer")
+ .HasColumnName("row_count");
+
+ b.Property("SourceName")
+ .HasColumnType("text")
+ .HasColumnName("source_name");
+
+ b.HasKey("Id")
+ .HasName("pk_import_batch");
+
+ b.ToTable("import_batch", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.IngestionEndpoint", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("Config")
+ .IsRequired()
+ .ValueGeneratedOnAdd()
+ .HasColumnType("jsonb")
+ .HasColumnName("config")
+ .HasDefaultValueSql("'{}'::jsonb");
+
+ b.Property("IsEnabled")
+ .HasColumnType("boolean")
+ .HasColumnName("is_enabled");
+
+ b.Property("LastSeenAt")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("last_seen_at");
+
+ b.Property("LastStatus")
+ .HasColumnType("text")
+ .HasColumnName("last_status");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasMaxLength(128)
+ .HasColumnType("character varying(128)")
+ .HasColumnName("name");
+
+ b.Property("Type")
+ .IsRequired()
+ .HasMaxLength(32)
+ .HasColumnType("character varying(32)")
+ .HasColumnName("type");
+
+ b.HasKey("Id")
+ .HasName("pk_ingestion_endpoint");
+
+ b.ToTable("ingestion_endpoint", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.ManualCost", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("Amount")
+ .HasColumnType("double precision")
+ .HasColumnName("amount");
+
+ b.Property("CategoryId")
+ .HasColumnType("integer")
+ .HasColumnName("category_id");
+
+ b.Property("Currency")
+ .IsRequired()
+ .HasMaxLength(8)
+ .HasColumnType("character varying(8)")
+ .HasColumnName("currency");
+
+ b.Property("ImportBatchId")
+ .HasColumnType("integer")
+ .HasColumnName("import_batch_id");
+
+ b.Property("MeterId")
+ .HasColumnType("integer")
+ .HasColumnName("meter_id");
+
+ b.Property("Notes")
+ .HasColumnType("text")
+ .HasColumnName("notes");
+
+ b.Property("PeriodEnd")
+ .HasColumnType("date")
+ .HasColumnName("period_end");
+
+ b.Property("PeriodStart")
+ .HasColumnType("date")
+ .HasColumnName("period_start");
+
+ b.HasKey("Id")
+ .HasName("pk_manual_cost");
+
+ b.HasIndex("CategoryId")
+ .HasDatabaseName("ix_manual_cost_category_id");
+
+ b.HasIndex("ImportBatchId")
+ .HasDatabaseName("ix_manual_cost_import_batch_id");
+
+ b.HasIndex("MeterId")
+ .HasDatabaseName("ix_manual_cost_meter_id");
+
+ b.ToTable("manual_cost", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Meter", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("CreatedAt")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("created_at")
+ .HasDefaultValueSql("now()");
+
+ b.Property("EnergyTypeId")
+ .HasColumnType("smallint")
+ .HasColumnName("energy_type_id");
+
+ b.Property("InitialBaseline")
+ .HasColumnType("double precision")
+ .HasColumnName("initial_baseline");
+
+ b.Property("InstalledAt")
+ .HasColumnType("date")
+ .HasColumnName("installed_at");
+
+ b.Property("IsActive")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("boolean")
+ .HasDefaultValue(true)
+ .HasColumnName("is_active");
+
+ b.Property("Location")
+ .HasColumnType("text")
+ .HasColumnName("location");
+
+ b.Property("Manufacturer")
+ .HasColumnType("text")
+ .HasColumnName("manufacturer");
+
+ b.Property("Meta")
+ .IsRequired()
+ .ValueGeneratedOnAdd()
+ .HasColumnType("jsonb")
+ .HasColumnName("meta")
+ .HasDefaultValueSql("'{}'::jsonb");
+
+ b.Property("Mode")
+ .IsRequired()
+ .HasMaxLength(32)
+ .HasColumnType("character varying(32)")
+ .HasColumnName("mode");
+
+ b.Property("Model")
+ .HasColumnType("text")
+ .HasColumnName("model");
+
+ b.Property("Name")
+ .IsRequired()
+ .HasColumnType("text")
+ .HasColumnName("name");
+
+ b.Property("RetiredAt")
+ .HasColumnType("date")
+ .HasColumnName("retired_at");
+
+ b.Property("SerialNumber")
+ .HasColumnType("text")
+ .HasColumnName("serial_number");
+
+ b.Property("Unit")
+ .IsRequired()
+ .HasColumnType("text")
+ .HasColumnName("unit");
+
+ b.Property("UpdatedAt")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("updated_at")
+ .HasDefaultValueSql("now()");
+
+ b.HasKey("Id")
+ .HasName("pk_meter");
+
+ b.HasIndex("EnergyTypeId", "IsActive")
+ .HasDatabaseName("ix_meter_energy_type_id_is_active");
+
+ b.ToTable("meter", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.MeterEvent", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("Amount")
+ .HasColumnType("double precision")
+ .HasColumnName("amount");
+
+ b.Property("EventType")
+ .IsRequired()
+ .HasMaxLength(32)
+ .HasColumnType("character varying(32)")
+ .HasColumnName("event_type");
+
+ b.Property("ImportBatchId")
+ .HasColumnType("integer")
+ .HasColumnName("import_batch_id");
+
+ b.Property("Meta")
+ .IsRequired()
+ .ValueGeneratedOnAdd()
+ .HasColumnType("jsonb")
+ .HasColumnName("meta")
+ .HasDefaultValueSql("'{}'::jsonb");
+
+ b.Property("MeterId")
+ .HasColumnType("integer")
+ .HasColumnName("meter_id");
+
+ b.Property("NewValue")
+ .HasColumnType("double precision")
+ .HasColumnName("new_value");
+
+ b.Property("Notes")
+ .HasColumnType("text")
+ .HasColumnName("notes");
+
+ b.Property("PrevValue")
+ .HasColumnType("double precision")
+ .HasColumnName("prev_value");
+
+ b.Property("Time")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("time");
+
+ b.Property("Unit")
+ .HasColumnType("text")
+ .HasColumnName("unit");
+
+ b.HasKey("Id")
+ .HasName("pk_meter_event");
+
+ b.HasIndex("ImportBatchId")
+ .HasDatabaseName("ix_meter_event_import_batch_id");
+
+ b.HasIndex("MeterId", "Time")
+ .HasDatabaseName("ix_meter_event_meter_id_time");
+
+ b.ToTable("meter_event", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.MeterLink", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("FromMeterId")
+ .HasColumnType("integer")
+ .HasColumnName("from_meter_id");
+
+ b.Property("ToMeterId")
+ .HasColumnType("integer")
+ .HasColumnName("to_meter_id");
+
+ b.HasKey("Id")
+ .HasName("pk_meter_link");
+
+ b.HasIndex("ToMeterId")
+ .HasDatabaseName("ix_meter_link_to_meter_id");
+
+ b.HasIndex("FromMeterId", "ToMeterId")
+ .IsUnique()
+ .HasDatabaseName("ix_meter_link_from_meter_id_to_meter_id");
+
+ b.ToTable("meter_link", null, t =>
+ {
+ t.HasCheckConstraint("ck_meter_link_distinct", "from_meter_id <> to_meter_id");
+ });
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.MeterSource", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("Config")
+ .IsRequired()
+ .ValueGeneratedOnAdd()
+ .HasColumnType("jsonb")
+ .HasColumnName("config")
+ .HasDefaultValueSql("'{}'::jsonb");
+
+ b.Property("EndpointId")
+ .HasColumnType("integer")
+ .HasColumnName("endpoint_id");
+
+ b.Property("IsEnabled")
+ .HasColumnType("boolean")
+ .HasColumnName("is_enabled");
+
+ b.Property("LastSeenAt")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("last_seen_at");
+
+ b.Property("LastStatus")
+ .HasColumnType("text")
+ .HasColumnName("last_status");
+
+ b.Property("LastValue")
+ .HasColumnType("double precision")
+ .HasColumnName("last_value");
+
+ b.Property("MeterId")
+ .HasColumnType("integer")
+ .HasColumnName("meter_id");
+
+ b.Property("Offset")
+ .HasColumnType("double precision")
+ .HasColumnName("offset");
+
+ b.Property("Priority")
+ .HasColumnType("integer")
+ .HasColumnName("priority");
+
+ b.Property("Scale")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("double precision")
+ .HasDefaultValue(1.0)
+ .HasColumnName("scale");
+
+ b.Property("SourceType")
+ .IsRequired()
+ .HasMaxLength(32)
+ .HasColumnType("character varying(32)")
+ .HasColumnName("source_type");
+
+ b.Property("ValueKind")
+ .IsRequired()
+ .HasMaxLength(16)
+ .HasColumnType("character varying(16)")
+ .HasColumnName("value_kind");
+
+ b.HasKey("Id")
+ .HasName("pk_meter_source");
+
+ b.HasIndex("EndpointId")
+ .HasDatabaseName("ix_meter_source_endpoint_id");
+
+ b.HasIndex("MeterId")
+ .HasDatabaseName("ix_meter_source_meter_id");
+
+ b.ToTable("meter_source", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Reading", b =>
+ {
+ b.Property("MeterId")
+ .HasColumnType("integer")
+ .HasColumnName("meter_id");
+
+ b.Property("Time")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("time");
+
+ b.Property("Flags")
+ .HasColumnType("integer")
+ .HasColumnName("flags");
+
+ b.Property("ImportBatchId")
+ .HasColumnType("integer")
+ .HasColumnName("import_batch_id");
+
+ b.Property("Quality")
+ .HasColumnType("smallint")
+ .HasColumnName("quality");
+
+ b.Property("SourceId")
+ .HasColumnType("integer")
+ .HasColumnName("source_id");
+
+ b.Property("Value")
+ .HasColumnType("double precision")
+ .HasColumnName("value");
+
+ b.HasKey("MeterId", "Time")
+ .HasName("pk_reading");
+
+ b.HasIndex("ImportBatchId")
+ .HasDatabaseName("ix_reading_import_batch_id");
+
+ b.ToTable("reading", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Tank", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("CachedAt")
+ .HasColumnType("timestamp with time zone")
+ .HasColumnName("cached_at");
+
+ b.Property("CachedBalance")
+ .HasColumnType("double precision")
+ .HasColumnName("cached_balance");
+
+ b.Property("Calibration")
+ .HasColumnType("jsonb")
+ .HasColumnName("calibration");
+
+ b.Property("Capacity")
+ .HasColumnType("double precision")
+ .HasColumnName("capacity");
+
+ b.Property("FixedRate")
+ .HasColumnType("double precision")
+ .HasColumnName("fixed_rate");
+
+ b.Property("LowThreshold")
+ .HasColumnType("double precision")
+ .HasColumnName("low_threshold");
+
+ b.Property("MeterId")
+ .HasColumnType("integer")
+ .HasColumnName("meter_id");
+
+ b.Property("RateMode")
+ .IsRequired()
+ .HasMaxLength(16)
+ .HasColumnType("character varying(16)")
+ .HasColumnName("rate_mode");
+
+ b.Property("ReorderThreshold")
+ .HasColumnType("double precision")
+ .HasColumnName("reorder_threshold");
+
+ b.Property("Unit")
+ .IsRequired()
+ .HasMaxLength(16)
+ .HasColumnType("character varying(16)")
+ .HasColumnName("unit");
+
+ b.HasKey("Id")
+ .HasName("pk_tank");
+
+ b.HasIndex("MeterId")
+ .IsUnique()
+ .HasDatabaseName("ix_tank_meter_id");
+
+ b.ToTable("tank", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Tariff", b =>
+ {
+ b.Property("Id")
+ .ValueGeneratedOnAdd()
+ .HasColumnType("integer")
+ .HasColumnName("id");
+
+ NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id"));
+
+ b.Property("Component")
+ .IsRequired()
+ .HasMaxLength(16)
+ .HasColumnType("character varying(16)")
+ .HasColumnName("component");
+
+ b.Property("Currency")
+ .IsRequired()
+ .HasMaxLength(8)
+ .HasColumnType("character varying(8)")
+ .HasColumnName("currency");
+
+ b.Property("Notes")
+ .HasColumnType("text")
+ .HasColumnName("notes");
+
+ b.Property("ScopeId")
+ .HasColumnType("integer")
+ .HasColumnName("scope_id");
+
+ b.Property("ScopeType")
+ .IsRequired()
+ .HasMaxLength(16)
+ .HasColumnType("character varying(16)")
+ .HasColumnName("scope_type");
+
+ b.Property("Unit")
+ .IsRequired()
+ .HasMaxLength(16)
+ .HasColumnType("character varying(16)")
+ .HasColumnName("unit");
+
+ b.Property("ValidFrom")
+ .HasColumnType("date")
+ .HasColumnName("valid_from");
+
+ b.Property("ValidTo")
+ .HasColumnType("date")
+ .HasColumnName("valid_to");
+
+ b.Property("Value")
+ .HasColumnType("double precision")
+ .HasColumnName("value");
+
+ b.HasKey("Id")
+ .HasName("pk_tariff");
+
+ b.HasIndex("ScopeType", "ScopeId", "Component", "ValidFrom")
+ .HasDatabaseName("ix_tariff_scope_type_scope_id_component_valid_from");
+
+ b.ToTable("tariff", (string)null);
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Consumption", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.Meter", null)
+ .WithMany()
+ .HasForeignKey("MeterId")
+ .OnDelete(DeleteBehavior.Restrict)
+ .IsRequired()
+ .HasConstraintName("fk_consumption_meter_meter_id");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.CostCategoryMember", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.CostCategory", "Category")
+ .WithMany("Members")
+ .HasForeignKey("CategoryId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired()
+ .HasConstraintName("fk_cost_category_member_cost_category_category_id");
+
+ b.HasOne("MeterVault.Core.Domain.EnergyType", null)
+ .WithMany()
+ .HasForeignKey("EnergyTypeId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .HasConstraintName("fk_cost_category_member_energy_type_energy_type_id");
+
+ b.HasOne("MeterVault.Core.Domain.Meter", null)
+ .WithMany()
+ .HasForeignKey("MeterId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .HasConstraintName("fk_cost_category_member_meter_meter_id");
+
+ b.Navigation("Category");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.ManualCost", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.CostCategory", null)
+ .WithMany()
+ .HasForeignKey("CategoryId")
+ .OnDelete(DeleteBehavior.SetNull)
+ .HasConstraintName("fk_manual_cost_cost_category_category_id");
+
+ b.HasOne("MeterVault.Core.Domain.Meter", null)
+ .WithMany()
+ .HasForeignKey("MeterId")
+ .OnDelete(DeleteBehavior.SetNull)
+ .HasConstraintName("fk_manual_cost_meter_meter_id");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Meter", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.EnergyType", "EnergyType")
+ .WithMany("Meters")
+ .HasForeignKey("EnergyTypeId")
+ .OnDelete(DeleteBehavior.Restrict)
+ .IsRequired()
+ .HasConstraintName("fk_meter_energy_type_energy_type_id");
+
+ b.Navigation("EnergyType");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.MeterEvent", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.Meter", null)
+ .WithMany()
+ .HasForeignKey("MeterId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired()
+ .HasConstraintName("fk_meter_event_meter_meter_id");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.MeterLink", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.Meter", "FromMeter")
+ .WithMany()
+ .HasForeignKey("FromMeterId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired()
+ .HasConstraintName("fk_meter_link_meter_from_meter_id");
+
+ b.HasOne("MeterVault.Core.Domain.Meter", "ToMeter")
+ .WithMany()
+ .HasForeignKey("ToMeterId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired()
+ .HasConstraintName("fk_meter_link_meter_to_meter_id");
+
+ b.Navigation("FromMeter");
+
+ b.Navigation("ToMeter");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.MeterSource", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.IngestionEndpoint", "Endpoint")
+ .WithMany()
+ .HasForeignKey("EndpointId")
+ .OnDelete(DeleteBehavior.SetNull)
+ .HasConstraintName("fk_meter_source_ingestion_endpoints_endpoint_id");
+
+ b.HasOne("MeterVault.Core.Domain.Meter", "Meter")
+ .WithMany("Sources")
+ .HasForeignKey("MeterId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired()
+ .HasConstraintName("fk_meter_source_meter_meter_id");
+
+ b.Navigation("Endpoint");
+
+ b.Navigation("Meter");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Reading", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.Meter", null)
+ .WithMany()
+ .HasForeignKey("MeterId")
+ .OnDelete(DeleteBehavior.Restrict)
+ .IsRequired()
+ .HasConstraintName("fk_reading_meter_meter_id");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Tank", b =>
+ {
+ b.HasOne("MeterVault.Core.Domain.Meter", "Meter")
+ .WithMany()
+ .HasForeignKey("MeterId")
+ .OnDelete(DeleteBehavior.Cascade)
+ .IsRequired()
+ .HasConstraintName("fk_tank_meter_meter_id");
+
+ b.Navigation("Meter");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.CostCategory", b =>
+ {
+ b.Navigation("Members");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.EnergyType", b =>
+ {
+ b.Navigation("Meters");
+ });
+
+ modelBuilder.Entity("MeterVault.Core.Domain.Meter", b =>
+ {
+ b.Navigation("Sources");
+ });
+#pragma warning restore 612, 618
+ }
+ }
+}
diff --git a/src/Infrastructure/Persistence/Migrations/20260718174732_BindUnboundMqttSourcesToSoleEnabledBroker.cs b/src/Infrastructure/Persistence/Migrations/20260718174732_BindUnboundMqttSourcesToSoleEnabledBroker.cs
new file mode 100644
index 0000000..d87294d
--- /dev/null
+++ b/src/Infrastructure/Persistence/Migrations/20260718174732_BindUnboundMqttSourcesToSoleEnabledBroker.cs
@@ -0,0 +1,42 @@
+using Microsoft.EntityFrameworkCore.Migrations;
+
+#nullable disable
+
+namespace MeterVault.Infrastructure.Persistence.Migrations
+{
+ ///
+ public partial class BindUnboundMqttSourcesToSoleEnabledBroker : Migration
+ {
+ ///
+ protected override void Up(MigrationBuilder migrationBuilder)
+ {
+ // Corrects BindUnboundMqttSourcesToSoleBroker, which counted brokers without regard to
+ // is_enabled. An instance with one live broker plus a disabled leftover counted two,
+ // declined to backfill on the grounds that the old routing was "ambiguous", and left its
+ // sources unbound — which under endpoint-scoped routing means permanently, silently dead.
+ //
+ // That reasoning was wrong for exactly this shape: MqttIngestionWorker only ever
+ // connected to enabled endpoints, so with a single enabled broker the mapping was never
+ // ambiguous. Re-run the backfill counting only enabled brokers.
+ //
+ // Idempotent and safe to follow the original: it touches only rows still NULL, so
+ // anything the first migration bound, or an operator has since bound by hand, is left
+ // alone. Instances that were already correct match nothing and are unaffected.
+ migrationBuilder.Sql("""
+ UPDATE meter_source AS s
+ SET endpoint_id = sole.id
+ FROM (SELECT id FROM ingestion_endpoint WHERE type = 'MqttBroker' AND is_enabled) AS sole
+ WHERE s.endpoint_id IS NULL
+ AND s.source_type IN ('Mqtt', 'Tasmota')
+ AND (SELECT count(*) FROM ingestion_endpoint WHERE type = 'MqttBroker' AND is_enabled) = 1;
+ """);
+ }
+
+ ///
+ protected override void Down(MigrationBuilder migrationBuilder)
+ {
+ // Intentionally empty, as for the migration this corrects: the rows it bound cannot be
+ // told apart from ones bound by hand, so clearing them would discard real configuration.
+ }
+ }
+}
diff --git a/tests/Integration.Tests/ImportRoundTripTests.cs b/tests/Integration.Tests/ImportRoundTripTests.cs
index f1f4171..254568f 100644
--- a/tests/Integration.Tests/ImportRoundTripTests.cs
+++ b/tests/Integration.Tests/ImportRoundTripTests.cs
@@ -85,4 +85,55 @@ public sealed class ImportRoundTripTests(TimescaleFixture fx)
// 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;
+ }
}
diff --git a/tests/Integration.Tests/MeterPeriodServiceTests.cs b/tests/Integration.Tests/MeterPeriodServiceTests.cs
index c36a301..34d8869 100644
--- a/tests/Integration.Tests/MeterPeriodServiceTests.cs
+++ b/tests/Integration.Tests/MeterPeriodServiceTests.cs
@@ -77,6 +77,19 @@ public sealed class MeterPeriodServiceTests(TimescaleFixture fx)
await CleanupAsync(db, meterId);
}
+ [Fact]
+ public async Task A_virtual_meter_reports_nothing_rather_than_a_confident_zero()
+ {
+ // Virtual meters evaluate on read and only materialize when a cost category references them
+ // (SDD §14.1). Summing `consumption` would render four zero tiles for a working meter.
+ await using var db = fx.CreateContext();
+ var meterId = await SetupAsync(db, MeterMode.Virtual);
+
+ Assert.Null(await NewService().GetAsync(meterId));
+
+ await CleanupAsync(db, meterId);
+ }
+
[Fact]
public async Task A_negative_previous_period_reports_no_basis_rather_than_an_inverted_percentage()
{