diff --git a/CLAUDE.md b/CLAUDE.md index 422401e..53ae339 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -6,7 +6,7 @@ This file provides guidance to Claude Code (claude.ai/code) when working with co MeterVault is a self-hosted, local-first energy & utility metering platform: it ingests meter data from Home Assistant, Tasmota and MQTT on a schedule, stores every reading timestamped and immutable, normalizes it into consumption, and turns it into cost dashboards. Energy types (electricity, water, heating oil, gas, …) and meters are **user-defined, never hardcoded**. -**Status: implemented (M0–M7) + SDD §8 panels.** The full solution is built and green — five projects, ~108 tests, working Docker deploy. `docs/SDD.md` remains the design reference; the milestone map (§12) matches the git history (M0…M7 commits). The dedicated **PV/Solar** (`/solar`), **Oil/consumable** (`/consumables`) and **meter-detail** (`/meters/{id}`) views (SDD §8.4–§8.6) are implemented as read models in `Infrastructure/Dashboard` (`SolarService`, `ConsumableService`, `MeterDetailService`) — PV meters are found by `Mode == GenerationCounter` and grid/load meters by a `role` tag in `Meter.Meta` (`MeterRoles`/`MeterMeta`), so nothing is hardcoded by name. **Admin write-CRUD** (SDD §8.7) is implemented as MudBlazor inline-dialog pages: energy types, meters (+ recompute on mode/baseline change), a meter's ingest sources (meter-detail Sources tab), tariffs, cost categories + members, and connectors (`ingestion_endpoint`, secrets by env-var reference only). `/admin/settings` is a read-only effective-config view (settings are env-driven and reproducible, not DB-stored). **Home Assistant reading** is configured here: an HA connector (`BaseUrl` + `TokenEnv`) + an HA source (entity id) drives `HomeAssistantWorker`'s REST poll; `HaConnectionTester` powers the connector "Test connection" button. Remaining refinements (HA WebSocket *push* — REST poll works today; commit-arbitrary-CSV-from-UI needs a meter-mapping wizard; full de-DE UI-string localization) are noted at their commits. Set `MeterVault__SeedReferenceData=true` (compose: `METERVAULT_SEED=true`) for a one-command populated demo. +**Status: implemented (M0–M7) + SDD §8 panels.** The full solution is built and green — five projects, ~108 tests, working Docker deploy. `docs/SDD.md` remains the design reference; the milestone map (§12) matches the git history (M0…M7 commits). The dedicated **PV/Solar** (`/solar`), **Oil/consumable** (`/consumables`) and **meter-detail** (`/meters/{id}`) views (SDD §8.4–§8.6) are implemented as read models in `Infrastructure/Dashboard` (`SolarService`, `ConsumableService`, `MeterDetailService`) — PV meters are found by `Mode == GenerationCounter` and grid/load meters by a `role` tag in `Meter.Meta` (`MeterRoles`/`MeterMeta`), so nothing is hardcoded by name. **Admin write-CRUD** (SDD §8.7) is implemented as MudBlazor inline-dialog pages: energy types, meters (+ recompute on mode/baseline change), a meter's ingest sources (meter-detail Sources tab), tariffs, cost categories + members, and connectors (`ingestion_endpoint`, secrets by env-var reference only). `/admin/settings` is a read-only effective-config view (settings are env-driven and reproducible, not DB-stored). **Home Assistant reading** is configured here: an HA connector (`BaseUrl` + `TokenEnv`) + an HA source (entity id) drives `HomeAssistantWorker`'s REST poll; `HaConnectionTester` powers the connector "Test connection" button. **Meter topology & flow**: `MeterLink` (a directed `from→to` edge; a downstream meter is a *subsection* of an upstream one, multi-parent allowed) drives a per-energy-type page `/energy/{id}` with a hand-rolled SVG **Sankey** (`SankeyChart.razor`, since ApexCharts has no Sankey type) computed by `FlowService` (link value = downstream consumption, split proportionally across multiple parents; unaccounted remainder → an "Other" node). Upstream meters are wired cycle-safely in the meter editor; the nav lists a link per energy type. Remaining refinements (HA WebSocket *push* — REST poll works today; commit-arbitrary-CSV-from-UI needs a meter-mapping wizard; full de-DE UI-string localization) are noted at their commits. Set `MeterVault__SeedReferenceData=true` (compose: `METERVAULT_SEED=true`) for a one-command populated demo. ## Source of truth diff --git a/README.md b/README.md index 14fc2ef..7cb737d 100644 --- a/README.md +++ b/README.md @@ -26,9 +26,14 @@ full design. savings), an **oil/consumable panel** (tank gauge, deliveries, burner runtime, effective L/h, forecast-to-empty) and a **per-meter detail view** (raw readings, consumption, sources, tariff timeline, events), one-click reference-data load, CSV dry-run. +- **Per-energy-type flow pages** (Electricity, Water, …): a **Sankey diagram** of the meter chain — + a downstream meter is a *subsection* of an upstream one (main → car, pool, garden, …), arrow + thickness ∝ amount, with an auto-computed "Other/unmetered" remainder. Meters can have several + upstreams (a merge, e.g. grid + solar → house). - **Admin UI**: full create/edit/delete for energy types, meters (with consumption recompute on - mode/baseline change), ingest sources, tariffs, cost categories, and MQTT/Home-Assistant - connectors; a "Test connection" for Home Assistant; effective-settings view. + mode/baseline change, and cycle-safe upstream-meter wiring), ingest sources, tariffs, cost + categories, and MQTT/Home-Assistant connectors; a "Test connection" for Home Assistant; + effective-settings view. - **REST API + OpenAPI/Swagger**, API-key auth, reverse-proxy trust (Authelia/Traefik). - **JSON config export/import** for portability; Docker Compose + multi-arch image. diff --git a/src/App/Components/Layout/NavMenu.razor b/src/App/Components/Layout/NavMenu.razor index ead1a9b..2b5012f 100644 --- a/src/App/Components/Layout/NavMenu.razor +++ b/src/App/Components/Layout/NavMenu.razor @@ -1,6 +1,16 @@ +@inject Microsoft.EntityFrameworkCore.IDbContextFactory DbFactory +@using Microsoft.EntityFrameworkCore +@using MeterVault.Core.Domain + Overview Trends + + @foreach (var type in _energyTypes) + { + @type.DisplayName + } + Solar / PV Oil / consumables Meters @@ -14,3 +24,32 @@ Settings + +@code { + private List _energyTypes = []; + + protected override async Task OnInitializedAsync() + { + try + { + await using var db = await DbFactory.CreateDbContextAsync(); + _energyTypes = await db.EnergyTypes.AsNoTracking().OrderBy(t => t.Id).ToListAsync(); + } + catch (Exception) + { + // Nav must never break the layout — a DB hiccup just hides the per-type links. + _energyTypes = []; + } + } + + // Map the energy type's stored icon name to a Material icon; fall back to a generic gauge. + private static string TypeIcon(string? icon) => icon switch + { + "bolt" => Icons.Material.Filled.Bolt, + "water_drop" => Icons.Material.Filled.WaterDrop, + "local_gas_station" => Icons.Material.Filled.LocalGasStation, + "gas_meter" => Icons.Material.Filled.GasMeter, + "thermostat" => Icons.Material.Filled.Thermostat, + _ => Icons.Material.Filled.Bolt, + }; +} diff --git a/src/App/Components/Pages/EnergyView.razor b/src/App/Components/Pages/EnergyView.razor new file mode 100644 index 0000000..4f9e561 --- /dev/null +++ b/src/App/Components/Pages/EnergyView.razor @@ -0,0 +1,173 @@ +@page "/energy/{Id:int}" +@rendermode InteractiveServer +@inject FlowService Flow +@inject MeterVault.Infrastructure.Costing.CostService Costs +@inject Microsoft.EntityFrameworkCore.IDbContextFactory DbFactory +@using Microsoft.EntityFrameworkCore +@using MudBlazor + +MeterVault — @(_graph?.EnergyType ?? "Energy") + +
+ @(_graph?.EnergyType ?? "Energy") flow + + Last 12 months + Last 24 months + Last 5 years + All time + +
+ +@if (_graph is null) +{ + +} +else if (!_graph.HasData) +{ + + No meters for this energy type yet. Add meters in Meters, or load the + reference data from Import. + +} +else +{ + + + + Top-level consumption + @Format.Number(_graph.Total, 0) @_graph.Unit + + + + + Cost (range) + @Format.Euro(_cost) + + + + + Meters + @_meters.Count + + + + + + Flow + @if (_graph.HasChain) + { + + Where the top-level flow goes. Arrow thickness ∝ amount; "Other" is the unmetered remainder. + + + } + else + { + + No meter chain configured yet. In Meters → edit a sub-meter and set its + upstream meter(s) to show where the main meter's flow divides (e.g. main → car, pool, other). + + @if (_graph.Nodes.Count > 0) + { + + MeterConsumption + + @foreach (var node in _graph.Nodes.OrderByDescending(n => n.Value)) + { + + @node.Label + @Format.Number(node.Value, 0) @_graph.Unit + + } + + + } + } + + + + Meters + + NameModeUpstream ofConsumption + + @foreach (var meter in _meters) + { + + @meter.Name + @meter.Mode + @UpstreamLabel(meter.Id) + @Format.Number(NodeValue(meter.Id), 0) @_graph.Unit + + } + + + +} + +@code { + [Parameter] + public int Id { get; set; } + + private int _months = 60; + private bool _loading; + private FlowGraph? _graph; + private double _cost; + private List _meters = []; + private Dictionary> _downstream = []; + + protected override Task OnParametersSetAsync() => LoadAsync(); + + private async Task OnRangeChanged(int months) + { + _months = months; + await LoadAsync(); + } + + private async Task LoadAsync() + { + if (_loading) + { + return; + } + + _loading = true; + _graph = null; + try + { + var asOf = DateOnly.FromDateTime(DateTime.UtcNow); + var from = new DateOnly(asOf.AddMonths(-_months).Year, asOf.AddMonths(-_months).Month, 1); + var to = asOf.AddMonths(1); + var typeId = (short)Id; + + _graph = await Flow.GetFlowAsync(typeId, from, to); + + await using var db = await DbFactory.CreateDbContextAsync(); + _meters = await db.Meters.AsNoTracking().Where(m => m.EnergyTypeId == typeId).OrderBy(m => m.Name).ToListAsync(); + var links = await db.MeterLinks.AsNoTracking() + .Where(l => _meters.Select(m => m.Id).Contains(l.FromMeterId)) + .ToListAsync(); + var names = _meters.ToDictionary(m => m.Id, m => m.Name); + _downstream = links + .GroupBy(l => l.FromMeterId) + .ToDictionary(g => g.Key, g => g.Select(l => names.GetValueOrDefault(l.ToMeterId, $"#{l.ToMeterId}")).ToList()); + + var fromUtc = new DateTimeOffset(from.Year, from.Month, from.Day, 0, 0, 0, TimeSpan.Zero); + var toUtc = new DateTimeOffset(to.Year, to.Month, to.Day, 0, 0, 0, TimeSpan.Zero); + double cost = 0; + foreach (var meter in _meters) + { + cost += (await Costs.GetMeterCostsAsync(meter.Id, fromUtc, toUtc, CostBucket.Month)).Sum(c => c.Cost); + } + _cost = cost; + } + finally + { + _loading = false; + } + } + + private double NodeValue(int meterId) => _graph?.Nodes.FirstOrDefault(n => n.MeterId == meterId)?.Value ?? 0; + + private string UpstreamLabel(int meterId) => + _downstream.TryGetValue(meterId, out var children) && children.Count > 0 ? string.Join(", ", children) : "—"; +} diff --git a/src/App/Components/Pages/Meters.razor b/src/App/Components/Pages/Meters.razor index 97071f2..d105937 100644 --- a/src/App/Components/Pages/Meters.razor +++ b/src/App/Components/Pages/Meters.razor @@ -87,6 +87,15 @@ else grid_import grid_export + + @foreach (var m in AvailableUpstream()) + { + @m.Name + } +
@@ -108,6 +117,7 @@ else @code { private List? _meters; private List _energyTypes = []; + private List _allLinks = []; private bool _editOpen; private EditModel _working = new(); private readonly DialogOptions _dialogOptions = new() { MaxWidth = MaxWidth.Small, FullWidth = true }; @@ -118,6 +128,7 @@ else { await using var db = await DbFactory.CreateDbContextAsync(); _energyTypes = await db.EnergyTypes.AsNoTracking().OrderBy(t => t.DisplayName).ToListAsync(); + _allLinks = await db.MeterLinks.AsNoTracking().ToListAsync(); _meters = await db.Meters .AsNoTracking() .Include(m => m.EnergyType) @@ -126,6 +137,49 @@ else .ToListAsync(); } + // Upstream candidates: same energy type, not self, and not a descendant (would create a cycle). + private IEnumerable AvailableUpstream() + { + if (_meters is null) + { + return []; + } + + var descendants = Descendants(_working.Id); + return _meters.Where(m => m.EnergyTypeId == _working.EnergyTypeId && m.Id != _working.Id && !descendants.Contains(m.Id)); + } + + private HashSet Descendants(int meterId) + { + var result = new HashSet(); + if (meterId == 0) + { + return result; + } + + var queue = new Queue(); + queue.Enqueue(meterId); + while (queue.Count > 0) + { + var current = queue.Dequeue(); + foreach (var link in _allLinks.Where(l => l.FromMeterId == current)) + { + if (result.Add(link.ToMeterId)) + { + queue.Enqueue(link.ToMeterId); + } + } + } + + return result; + } + + private string UpstreamText(IReadOnlyList ids) + { + var names = ids.Select(idText => int.TryParse(idText, out var id) ? _meters?.FirstOrDefault(m => m.Id == id)?.Name ?? idText : idText); + return string.Join(", ", names); + } + private void OpenEdit(Meter? meter) { if (meter is null) @@ -150,6 +204,7 @@ else Manufacturer = meter.Manufacturer, Model = meter.Model, IsActive = meter.IsActive, + Upstream = _allLinks.Where(l => l.ToMeterId == meter.Id).Select(l => l.FromMeterId).ToHashSet(), }; } _editOpen = true; @@ -164,9 +219,10 @@ else } await using var db = await DbFactory.CreateDbContextAsync(); + int meterId; if (_working.Id == 0) { - db.Meters.Add(new Meter + var meter = new Meter { Name = _working.Name.Trim(), EnergyTypeId = _working.EnergyTypeId, @@ -179,8 +235,10 @@ else Manufacturer = Trim(_working.Manufacturer), Model = Trim(_working.Model), IsActive = _working.IsActive, - }); + }; + db.Meters.Add(meter); await db.SaveChangesAsync(); + meterId = meter.Id; } else { @@ -208,13 +266,35 @@ else } await tx.CommitAsync(); + meterId = existing.Id; } + await SyncUpstreamAsync(db, meterId, _working.Upstream); + _editOpen = false; Snackbar.Add("Saved.", Severity.Success); await LoadAsync(); } + /// Reconciles the meter's incoming flow links to the selected upstream meters. + private static async Task SyncUpstreamAsync(MeterVault.Infrastructure.Persistence.MeterVaultDbContext db, int meterId, IEnumerable desiredUpstream) + { + var desired = desiredUpstream.Where(id => id != meterId).ToHashSet(); + var existing = await db.MeterLinks.Where(l => l.ToMeterId == meterId).ToListAsync(); + + foreach (var link in existing.Where(l => !desired.Contains(l.FromMeterId))) + { + db.MeterLinks.Remove(link); + } + + foreach (var fromId in desired.Where(id => existing.All(l => l.FromMeterId != id))) + { + db.MeterLinks.Add(new MeterLink { FromMeterId = fromId, ToMeterId = meterId }); + } + + await db.SaveChangesAsync(); + } + private async Task DeleteAsync(Meter meter) { await using var db = await DbFactory.CreateDbContextAsync(); @@ -258,6 +338,7 @@ else public string? Manufacturer { get; set; } public string? Model { get; set; } public bool IsActive { get; set; } = true; + public IReadOnlyCollection Upstream { get; set; } = new HashSet(); public bool RecomputeNeeded => Mode != OriginalMode || Math.Abs(InitialBaseline - OriginalBaseline) > 1e-9; } diff --git a/src/App/Components/Shared/SankeyChart.razor b/src/App/Components/Shared/SankeyChart.razor new file mode 100644 index 0000000..10b86b7 --- /dev/null +++ b/src/App/Components/Shared/SankeyChart.razor @@ -0,0 +1,165 @@ +@using System.Globalization +@using System.Text +@using System.Net +@using MeterVault.Infrastructure.Dashboard + +@if (string.IsNullOrEmpty(_svg)) +{ + No flow to show for this period. +} +else +{ +
+ @((MarkupString)_svg) +
+} + +@code { + private const double W = 1000; + private const double NodeWidth = 16; + private const double NodeGap = 12; + private const double LeftPad = 8; + private const double RightPad = 8; + + [Parameter, EditorRequired] + public IReadOnlyList Nodes { get; set; } = []; + + [Parameter, EditorRequired] + public IReadOnlyList Links { get; set; } = []; + + [Parameter] + public string Unit { get; set; } = ""; + + private string _svg = ""; + + protected override void OnParametersSet() => _svg = BuildSvg(); + + private string BuildSvg() + { + if (Nodes.Count == 0) + { + return ""; + } + + var maxDepth = Nodes.Max(n => n.Depth); + var columns = Nodes.GroupBy(n => n.Depth).ToDictionary(g => g.Key, g => g.OrderByDescending(n => n.Value).ToList()); + var maxCount = columns.Values.Max(c => c.Count); + var height = Math.Max(320, (maxCount * 46) + 40); + + // One value→pixel scale so every column fits (flow is conserved → column totals are ~equal; + // the densest column with the most gaps constrains the scale). + var scale = double.MaxValue; + foreach (var col in columns.Values) + { + var sum = col.Sum(n => n.Value); + if (sum > 0) + { + scale = Math.Min(scale, (height - ((col.Count - 1) * NodeGap) - 20) / sum); + } + } + + if (double.IsInfinity(scale) || scale <= 0) + { + scale = 1; + } + + var colStep = maxDepth == 0 ? 0 : (W - LeftPad - RightPad - NodeWidth) / maxDepth; + var geo = new Dictionary(); + foreach (var (depth, col) in columns) + { + var heights = col.Select(n => Math.Max(3, n.Value * scale)).ToList(); + var colHeight = heights.Sum() + ((col.Count - 1) * NodeGap); + var y = (height - colHeight) / 2; + var x = LeftPad + (depth * colStep); + for (var i = 0; i < col.Count; i++) + { + geo[col[i].Id] = new NodeGeo(x, y, heights[i]); + y += heights[i] + NodeGap; + } + } + + // Ribbon band offsets: order each source's out-links by target y, each target's in-links by source y. + var srcOffset = new Dictionary(); + var dstOffset = new Dictionary(); + var srcBand = new Dictionary(); + var dstBand = new Dictionary(); + foreach (var group in Links.GroupBy(l => l.From)) + { + foreach (var link in group.OrderBy(l => geo.TryGetValue(l.To, out var g) ? g.Y : 0)) + { + srcBand[link] = srcOffset.GetValueOrDefault(group.Key); + srcOffset[group.Key] = srcBand[link] + (link.Value * scale); + } + } + + foreach (var group in Links.GroupBy(l => l.To)) + { + foreach (var link in group.OrderBy(l => geo.TryGetValue(l.From, out var g) ? g.Y : 0)) + { + dstBand[link] = dstOffset.GetValueOrDefault(group.Key); + dstOffset[group.Key] = dstBand[link] + (link.Value * scale); + } + } + + var color = Nodes.ToDictionary(n => n.Id, n => n.ColorHex ?? "#607D8B"); + var label = Nodes.ToDictionary(n => n.Id, n => n.Label); + + var sb = new StringBuilder(); + sb.Append(CultureInfo.InvariantCulture, + $""); + + // Ribbons first (under nodes). + foreach (var link in Links) + { + if (!geo.TryGetValue(link.From, out var s) || !geo.TryGetValue(link.To, out var t)) + { + continue; + } + + var band = link.Value * scale; + var sy0 = s.Y + srcBand[link]; + var ty0 = t.Y + dstBand[link]; + var sx = s.X + NodeWidth; + var tx = t.X; + var midX = (sx + tx) / 2; + var path = + $"M{F(sx)},{F(sy0)} C{F(midX)},{F(sy0)} {F(midX)},{F(ty0)} {F(tx)},{F(ty0)} " + + $"L{F(tx)},{F(ty0 + band)} C{F(midX)},{F(ty0 + band)} {F(midX)},{F(sy0 + band)} {F(sx)},{F(sy0 + band)} Z"; + var tip = Enc($"{label.GetValueOrDefault(link.From)} → {label.GetValueOrDefault(link.To)}: {Fmt(link.Value)}"); + sb.Append(CultureInfo.InvariantCulture, + $"{tip}"); + } + + // Nodes + labels. + foreach (var node in Nodes) + { + if (!geo.TryGetValue(node.Id, out var g)) + { + continue; + } + + var rightmost = node.Depth == maxDepth; + var labelX = rightmost ? g.X - 6 : g.X + NodeWidth + 6; + var anchor = rightmost ? "end" : "start"; + var fill = Enc(node.ColorHex ?? "#607D8B"); + var name = Enc(node.Label); + var val = Enc(Fmt(node.Value)); + sb.Append(CultureInfo.InvariantCulture, + $"{name}: {val}"); + sb.Append(CultureInfo.InvariantCulture, + $"" + + $"{name}{val}"); + } + + sb.Append(""); + return sb.ToString(); + } + + private string Fmt(double value) => $"{value.ToString("N0", CultureInfo.GetCultureInfo("de-DE"))} {Unit}".Trim(); + + private static string F(double value) => value.ToString("0.##", CultureInfo.InvariantCulture); + + private static string Enc(string value) => WebUtility.HtmlEncode(value); + + private readonly record struct NodeGeo(double X, double Y, double H); +} diff --git a/src/Core/Domain/MeterLink.cs b/src/Core/Domain/MeterLink.cs new file mode 100644 index 0000000..6a49a1d --- /dev/null +++ b/src/Core/Domain/MeterLink.cs @@ -0,0 +1,23 @@ +namespace MeterVault.Core.Domain; + +/// +/// A directed flow edge in the meter topology: energy/water measured by is a +/// subsection of the flow through (upstream → downstream). Not an +/// addition — a sub-meter shows where an upstream meter's flow goes. Multiple edges into one meter +/// model a merge (e.g. house load fed by grid + solar); multiple edges out model a split +/// (main → car, pool, …). Drives the per-energy-type flow (Sankey) view. +/// +public sealed class MeterLink +{ + public int Id { get; set; } + + /// Upstream meter (the larger flow this edge draws from). + public int FromMeterId { get; set; } + + public Meter? FromMeter { get; set; } + + /// Downstream meter (the subsection). + public int ToMeterId { get; set; } + + public Meter? ToMeter { get; set; } +} diff --git a/src/Infrastructure/Dashboard/FlowModels.cs b/src/Infrastructure/Dashboard/FlowModels.cs new file mode 100644 index 0000000..ddb4469 --- /dev/null +++ b/src/Infrastructure/Dashboard/FlowModels.cs @@ -0,0 +1,26 @@ +namespace MeterVault.Infrastructure.Dashboard; + +/// A node in the flow graph: a meter, or a synthetic "Other/unmetered" remainder. +public sealed record FlowNode(string Id, string Label, double Value, int Depth, string? ColorHex, bool IsOther, int? MeterId); + +/// A directed flow edge with the quantity that flows along it, in the energy type's base unit. +public sealed record FlowLink(string From, string To, double Value); + +/// +/// The per-energy-type flow graph (SDD-style topology view): meters as nodes sized by consumption, +/// directed edges sized by the flow along each configured link, plus "Other" remainders where an +/// upstream meter's flow isn't fully accounted for by its sub-meters. Rendered as a Sankey diagram. +/// +public sealed record FlowGraph( + short EnergyTypeId, + string EnergyType, + string Unit, + double Total, + IReadOnlyList Nodes, + IReadOnlyList Links) +{ + public bool HasData => Nodes.Count > 0; + + /// True when meters are actually chained (not just a flat, unlinked list). + public bool HasChain => Links.Count > 0; +} diff --git a/src/Infrastructure/Dashboard/FlowService.cs b/src/Infrastructure/Dashboard/FlowService.cs new file mode 100644 index 0000000..fdda31b --- /dev/null +++ b/src/Infrastructure/Dashboard/FlowService.cs @@ -0,0 +1,142 @@ +using MeterVault.Core.Domain; +using MeterVault.Infrastructure.Persistence; +using Microsoft.EntityFrameworkCore; + +namespace MeterVault.Infrastructure.Dashboard; + +/// +/// Builds the per-energy-type flow graph (Sankey) from the meter topology () +/// and consumption over a period. Each meter is a node sized by its consumption; each configured +/// edge carries the downstream meter's consumption (split proportionally when a meter has several +/// upstreams); the unaccounted remainder under a meter becomes a synthetic "Other" node. Nothing is +/// hardcoded per energy type — it works for electricity, water, gas, … alike. DbContext factory +/// keeps it Blazor-circuit safe. +/// +public sealed class FlowService(IDbContextFactory contextFactory) +{ + private const double Epsilon = 0.01; + + private readonly IDbContextFactory _contextFactory = contextFactory; + + public async Task GetFlowAsync(short energyTypeId, DateOnly from, DateOnly to, CancellationToken cancellationToken = default) + { + await using var db = await _contextFactory.CreateDbContextAsync(cancellationToken).ConfigureAwait(false); + + var energyType = await db.EnergyTypes.AsNoTracking().FirstOrDefaultAsync(t => t.Id == energyTypeId, cancellationToken).ConfigureAwait(false); + var meters = await db.Meters.AsNoTracking().Where(m => m.EnergyTypeId == energyTypeId).ToListAsync(cancellationToken).ConfigureAwait(false); + if (energyType is null || meters.Count == 0) + { + return new FlowGraph(energyTypeId, energyType?.DisplayName ?? "", energyType?.BaseUnit ?? "", 0, [], []); + } + + var meterIds = meters.Select(m => m.Id).ToHashSet(); + var fromUtc = ToUtc(from); + var toUtc = ToUtc(to); + + var sums = await db.Consumption.AsNoTracking() + .Where(c => c.Time >= fromUtc && c.Time < toUtc && c.Kind == ConsumptionKind.Consumption) + .GroupBy(c => c.MeterId) + .Select(g => new { MeterId = g.Key, Total = g.Sum(x => x.Amount) }) + .ToListAsync(cancellationToken).ConfigureAwait(false); + var value = sums.Where(s => meterIds.Contains(s.MeterId)).ToDictionary(s => s.MeterId, s => s.Total); + double V(int id) => value.GetValueOrDefault(id); + + var links = await db.MeterLinks.AsNoTracking() + .Where(l => meterIds.Contains(l.FromMeterId) && meterIds.Contains(l.ToMeterId)) + .ToListAsync(cancellationToken).ConfigureAwait(false); + + var parents = meters.ToDictionary(m => m.Id, _ => new List()); + var children = meters.ToDictionary(m => m.Id, _ => new List()); + foreach (var link in links) + { + children[link.FromMeterId].Add(link.ToMeterId); + parents[link.ToMeterId].Add(link.FromMeterId); + } + + var depth = ComputeDepths(meters.Select(m => m.Id).ToList(), parents, children); + + // Link value: a child's consumption flows in from its parent(s); with several parents it is + // split proportionally to the parents' own consumption (equal split if those are all zero). + var flowLinks = new List(); + var outgoingByParent = meters.ToDictionary(m => m.Id, _ => 0d); + foreach (var (childId, parentIds) in parents) + { + if (parentIds.Count == 0) + { + continue; + } + + var parentTotal = parentIds.Sum(V); + foreach (var parentId in parentIds) + { + var share = parentIds.Count == 1 ? 1d + : parentTotal > Epsilon ? V(parentId) / parentTotal + : 1d / parentIds.Count; + var linkValue = V(childId) * share; + if (linkValue > Epsilon) + { + flowLinks.Add(new FlowLink(NodeId(parentId), NodeId(childId), linkValue)); + outgoingByParent[parentId] += linkValue; + } + } + } + + var nodes = new List(); + foreach (var meter in meters) + { + // Keep a meter node if it carries flow or participates in the topology. + if (V(meter.Id) <= Epsilon && children[meter.Id].Count == 0 && parents[meter.Id].Count == 0) + { + continue; + } + + nodes.Add(new FlowNode(NodeId(meter.Id), meter.Name, V(meter.Id), depth.GetValueOrDefault(meter.Id), energyType.ColorHex, false, meter.Id)); + + // Unaccounted remainder under a meter with sub-meters → "Other". + if (children[meter.Id].Count > 0) + { + var remainder = V(meter.Id) - outgoingByParent[meter.Id]; + if (remainder > Epsilon) + { + var otherId = $"other{meter.Id}"; + nodes.Add(new FlowNode(otherId, $"Other ({meter.Name})", remainder, depth.GetValueOrDefault(meter.Id) + 1, "#78909C", true, null)); + flowLinks.Add(new FlowLink(NodeId(meter.Id), otherId, remainder)); + } + } + } + + var total = meters.Where(m => parents[m.Id].Count == 0).Sum(m => V(m.Id)); + return new FlowGraph(energyTypeId, energyType.DisplayName, energyType.BaseUnit, total, nodes, flowLinks); + } + + /// Longest-path depth from the roots (Kahn topological relaxation); robust to stray cycles. + private static Dictionary ComputeDepths( + List ids, Dictionary> parents, Dictionary> children) + { + var depth = ids.ToDictionary(id => id, _ => 0); + var indegree = ids.ToDictionary(id => id, id => parents[id].Count); + var queue = new Queue(ids.Where(id => indegree[id] == 0)); + var processed = 0; + + while (queue.Count > 0) + { + var node = queue.Dequeue(); + processed++; + foreach (var child in children[node]) + { + depth[child] = Math.Max(depth[child], depth[node] + 1); + if (--indegree[child] == 0) + { + queue.Enqueue(child); + } + } + } + + // Any nodes left (a cycle) keep depth 0 — the admin prevents cycles, this is just a guard. + return depth; + } + + private static string NodeId(int meterId) => $"m{meterId}"; + + private static DateTimeOffset ToUtc(DateOnly date) => new(date.Year, date.Month, date.Day, 0, 0, 0, TimeSpan.Zero); +} diff --git a/src/Infrastructure/DependencyInjection.cs b/src/Infrastructure/DependencyInjection.cs index 315d7d2..41bf6ad 100644 --- a/src/Infrastructure/DependencyInjection.cs +++ b/src/Infrastructure/DependencyInjection.cs @@ -42,6 +42,7 @@ public static class DependencyInjection services.AddScoped(); services.AddScoped(); services.AddScoped(); + services.AddScoped(); services.AddScoped(); return services; diff --git a/src/Infrastructure/Import/ReferenceDataImporter.cs b/src/Infrastructure/Import/ReferenceDataImporter.cs index e0d0557..8b5e365 100644 --- a/src/Infrastructure/Import/ReferenceDataImporter.cs +++ b/src/Infrastructure/Import/ReferenceDataImporter.cs @@ -68,6 +68,10 @@ public sealed class ReferenceDataImporter(MeterVaultDbContext db, ImportService Calibration = $"{{\"volumePerUnit\":{ReferenceProfiles.OilLitresPerCm.ToString(System.Globalization.CultureInfo.InvariantCulture)}}}", }); + // Demo flow chain: car charging is a subsection of total house load (Haus → Auto), so the + // electricity flow view shows Haus dividing into Auto + an "Other" remainder. + _db.MeterLinks.Add(new MeterLink { FromMeterId = haus.Id, ToMeterId = auto.Id }); + AddElectricityTariffs(electricity); AddTariff(TariffScope.EnergyType, water, 5.00, "EUR/m3", new DateOnly(2022, 11, 1)); // Strom and Wasser costs are computed from meters + tariffs; only Heizung comes from the diff --git a/src/Infrastructure/Persistence/MeterVaultDbContext.cs b/src/Infrastructure/Persistence/MeterVaultDbContext.cs index 80509ac..abf62af 100644 --- a/src/Infrastructure/Persistence/MeterVaultDbContext.cs +++ b/src/Infrastructure/Persistence/MeterVaultDbContext.cs @@ -16,6 +16,7 @@ public sealed class MeterVaultDbContext(DbContextOptions op public DbSet EnergyTypes => Set(); public DbSet Meters => Set(); public DbSet MeterSources => Set(); + public DbSet MeterLinks => Set(); public DbSet Readings => Set(); public DbSet Consumption => Set(); public DbSet MeterEvents => Set(); @@ -74,6 +75,16 @@ public sealed class MeterVaultDbContext(DbContextOptions op e.HasIndex(x => x.MeterId); }); + b.Entity(e => + { + e.ToTable("meter_link"); + e.HasKey(x => x.Id); + e.HasOne(x => x.FromMeter).WithMany().HasForeignKey(x => x.FromMeterId).OnDelete(DeleteBehavior.Cascade); + e.HasOne(x => x.ToMeter).WithMany().HasForeignKey(x => x.ToMeterId).OnDelete(DeleteBehavior.Cascade); + e.HasIndex(x => new { x.FromMeterId, x.ToMeterId }).IsUnique(); + e.ToTable(t => t.HasCheckConstraint("ck_meter_link_distinct", "from_meter_id <> to_meter_id")); + }); + // Hypertable — the (meter_id, time) PK contains the partition column (time), // which Timescale requires. Converted to a hypertable in a raw-SQL migration. b.Entity(e => diff --git a/src/Infrastructure/Persistence/Migrations/20260714114901_AddMeterLinks.Designer.cs b/src/Infrastructure/Persistence/Migrations/20260714114901_AddMeterLinks.Designer.cs new file mode 100644 index 0000000..09af149 --- /dev/null +++ b/src/Infrastructure/Persistence/Migrations/20260714114901_AddMeterLinks.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("20260714114901_AddMeterLinks")] + partial class AddMeterLinks + { + /// + 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/20260714114901_AddMeterLinks.cs b/src/Infrastructure/Persistence/Migrations/20260714114901_AddMeterLinks.cs new file mode 100644 index 0000000..40a2b62 --- /dev/null +++ b/src/Infrastructure/Persistence/Migrations/20260714114901_AddMeterLinks.cs @@ -0,0 +1,60 @@ +using Microsoft.EntityFrameworkCore.Migrations; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; + +#nullable disable + +namespace MeterVault.Infrastructure.Persistence.Migrations +{ + /// + public partial class AddMeterLinks : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "meter_link", + columns: table => new + { + id = table.Column(type: "integer", nullable: false) + .Annotation("Npgsql:ValueGenerationStrategy", NpgsqlValueGenerationStrategy.IdentityByDefaultColumn), + from_meter_id = table.Column(type: "integer", nullable: false), + to_meter_id = table.Column(type: "integer", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("pk_meter_link", x => x.id); + table.CheckConstraint("ck_meter_link_distinct", "from_meter_id <> to_meter_id"); + table.ForeignKey( + name: "fk_meter_link_meter_from_meter_id", + column: x => x.from_meter_id, + principalTable: "meter", + principalColumn: "id", + onDelete: ReferentialAction.Cascade); + table.ForeignKey( + name: "fk_meter_link_meter_to_meter_id", + column: x => x.to_meter_id, + principalTable: "meter", + principalColumn: "id", + onDelete: ReferentialAction.Cascade); + }); + + migrationBuilder.CreateIndex( + name: "ix_meter_link_from_meter_id_to_meter_id", + table: "meter_link", + columns: new[] { "from_meter_id", "to_meter_id" }, + unique: true); + + migrationBuilder.CreateIndex( + name: "ix_meter_link_to_meter_id", + table: "meter_link", + column: "to_meter_id"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "meter_link"); + } + } +} diff --git a/src/Infrastructure/Persistence/Migrations/MeterVaultDbContextModelSnapshot.cs b/src/Infrastructure/Persistence/Migrations/MeterVaultDbContextModelSnapshot.cs index 81c37a9..679098a 100644 --- a/src/Infrastructure/Persistence/Migrations/MeterVaultDbContextModelSnapshot.cs +++ b/src/Infrastructure/Persistence/Migrations/MeterVaultDbContextModelSnapshot.cs @@ -499,6 +499,39 @@ namespace MeterVault.Infrastructure.Persistence.Migrations 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") @@ -812,6 +845,27 @@ namespace MeterVault.Infrastructure.Persistence.Migrations .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") diff --git a/tests/Integration.Tests/DashboardRenderTests.cs b/tests/Integration.Tests/DashboardRenderTests.cs index 4a78c59..af9dc3f 100644 --- a/tests/Integration.Tests/DashboardRenderTests.cs +++ b/tests/Integration.Tests/DashboardRenderTests.cs @@ -46,6 +46,7 @@ public sealed class DashboardRenderTests(TimescaleFixture fx) // Panel read models compute real figures from the reference data (SDD §8.4–§8.6). int hausId; + short electricityTypeId; using (var scope = factory.Services.CreateScope()) { var services = scope.ServiceProvider; @@ -72,6 +73,13 @@ public sealed class DashboardRenderTests(TimescaleFixture fx) Assert.NotNull(detail); Assert.True(detail!.ReadingCount > 0); Assert.True(detail.TotalConsumption > 0); + + // Flow graph: the demo Haus → Auto chain yields a link + an "Other (Haus)" remainder. + electricityTypeId = await db.EnergyTypes.Where(t => t.Key == "electricity").Select(t => t.Id).FirstAsync(); + var flow = await services.GetRequiredService() + .GetFlowAsync(electricityTypeId, new DateOnly(1997, 1, 1), new DateOnly(2027, 1, 1)); + Assert.True(flow.HasChain); + Assert.Contains(flow.Nodes, n => n.IsOther); } using var client = factory.CreateClient(); @@ -91,6 +99,7 @@ public sealed class DashboardRenderTests(TimescaleFixture fx) "/meters", "/trends", "/solar", "/consumables", "/import", "/admin/tariffs", "/admin/energy-types", "/admin/categories", "/admin/connectors", "/admin/settings", $"/meters/{hausId}", + $"/energy/{electricityTypeId}", }) { var response = await client.GetAsync(new Uri(path, UriKind.Relative)); @@ -106,6 +115,7 @@ public sealed class DashboardRenderTests(TimescaleFixture fx) private static async Task ClearDataAsync(MeterVaultDbContext db) { + await db.MeterLinks.ExecuteDeleteAsync(); await db.Consumption.ExecuteDeleteAsync(); await db.Readings.ExecuteDeleteAsync(); await db.MeterEvents.ExecuteDeleteAsync(); diff --git a/tests/Integration.Tests/FlowServiceTests.cs b/tests/Integration.Tests/FlowServiceTests.cs new file mode 100644 index 0000000..8a146c9 --- /dev/null +++ b/tests/Integration.Tests/FlowServiceTests.cs @@ -0,0 +1,111 @@ +using MeterVault.Core.Domain; +using MeterVault.Infrastructure.Dashboard; +using MeterVault.Infrastructure.Persistence; +using Microsoft.EntityFrameworkCore; + +namespace MeterVault.Integration.Tests; + +/// +/// The per-energy-type flow graph (Sankey): a single-parent chain attributes the child's full +/// consumption to its parent and shows the remainder as "Other"; a two-parent merge splits the +/// child's consumption proportionally to the parents' own consumption. +/// +[Collection("Timescale")] +public sealed class FlowServiceTests(TimescaleFixture fx) +{ + [Fact] + public async Task Single_parent_chain_makes_other_remainder() + { + await using var db = fx.CreateContext(); + try + { + var type = await SeedTypeAsync(db, "flow_elec_a"); + var main = await AddMeterAsync(db, "Main", type); + var car = await AddMeterAsync(db, "Car", type); + db.MeterLinks.Add(new MeterLink { FromMeterId = main.Id, ToMeterId = car.Id }); + await db.SaveChangesAsync(); + + await AddConsumptionAsync(db, main.Id, 100); + await AddConsumptionAsync(db, car.Id, 30); + + var graph = await new FlowService(fx).GetFlowAsync(type, new DateOnly(2024, 1, 1), new DateOnly(2024, 12, 31)); + + Assert.Equal(100, graph.Total, 1); + var link = Assert.Single(graph.Links, l => l.To == $"m{car.Id}"); + Assert.Equal(30, link.Value, 1); // full child consumption flows from its single parent + var other = Assert.Single(graph.Nodes, n => n.IsOther); + Assert.Equal(70, other.Value, 1); // 100 − 30 + } + finally + { + await ClearAsync(db); + } + } + + [Fact] + public async Task Two_parents_split_child_proportionally() + { + await using var db = fx.CreateContext(); + try + { + var type = await SeedTypeAsync(db, "flow_elec_b"); + var grid = await AddMeterAsync(db, "Grid", type); + var solar = await AddMeterAsync(db, "Solar draw", type); + var house = await AddMeterAsync(db, "House", type); + db.MeterLinks.Add(new MeterLink { FromMeterId = grid.Id, ToMeterId = house.Id }); + db.MeterLinks.Add(new MeterLink { FromMeterId = solar.Id, ToMeterId = house.Id }); + await db.SaveChangesAsync(); + + await AddConsumptionAsync(db, grid.Id, 75); + await AddConsumptionAsync(db, solar.Id, 25); + await AddConsumptionAsync(db, house.Id, 40); + + var graph = await new FlowService(fx).GetFlowAsync(type, new DateOnly(2024, 1, 1), new DateOnly(2024, 12, 31)); + + // House (40) splits 75:25 → 30 from grid, 10 from solar. + Assert.Equal(30, graph.Links.Single(l => l.From == $"m{grid.Id}" && l.To == $"m{house.Id}").Value, 1); + Assert.Equal(10, graph.Links.Single(l => l.From == $"m{solar.Id}" && l.To == $"m{house.Id}").Value, 1); + } + finally + { + await ClearAsync(db); + } + } + + private static async Task SeedTypeAsync(MeterVaultDbContext db, string key) + { + var type = new EnergyType { Key = key, DisplayName = key, BaseUnit = "kWh", DefaultMode = MeterMode.CumulativeCounter }; + db.EnergyTypes.Add(type); + await db.SaveChangesAsync(); + return type.Id; + } + + private static async Task AddMeterAsync(MeterVaultDbContext db, string name, short type) + { + var meter = new Meter { Name = name, EnergyTypeId = type, Mode = MeterMode.DirectDelta, Unit = "kWh" }; + db.Meters.Add(meter); + await db.SaveChangesAsync(); + return meter; + } + + private static async Task AddConsumptionAsync(MeterVaultDbContext db, int meterId, double amount) + { + db.Consumption.Add(new Consumption + { + MeterId = meterId, + Time = new DateTimeOffset(2024, 6, 15, 0, 0, 0, TimeSpan.Zero), + Amount = amount, + Kind = ConsumptionKind.Consumption, + Quality = ReadingQuality.Manual, + }); + await db.SaveChangesAsync(); + } + + private static async Task ClearAsync(MeterVaultDbContext db) + { + await db.MeterLinks.ExecuteDeleteAsync(); + await db.Consumption.ExecuteDeleteAsync(); + await db.Meters.ExecuteDeleteAsync(); + await db.EnergyTypes.Where(t => t.Key.StartsWith("flow_elec_")).ExecuteDeleteAsync(); + } +}