Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c0bbaba99f | ||
|
|
9bd0d60cc8 |
@@ -214,21 +214,31 @@ else
|
||||
<MudText Typo="Typo.h6">@(_sourceEdit.Id == 0 ? "New source" : "Edit source")</MudText>
|
||||
</TitleContent>
|
||||
<DialogContent>
|
||||
<MudSelect T="SourceType" @bind-Value="_sourceEdit.SourceType" Label="Source type" Class="mb-2">
|
||||
<MudSelect T="SourceType" Value="_sourceEdit.SourceType" ValueChanged="OnSourceTypeChanged" Label="Source type" Class="mb-2">
|
||||
@foreach (var type in Enum.GetValues<SourceType>())
|
||||
{
|
||||
<MudSelectItem T="SourceType" Value="type">@type</MudSelectItem>
|
||||
}
|
||||
</MudSelect>
|
||||
@if (_sourceEdit.SourceType is SourceType.HomeAssistant or SourceType.Mqtt or SourceType.Tasmota)
|
||||
@if (RequiredEndpointType(_sourceEdit.SourceType) is { } needed)
|
||||
{
|
||||
<MudSelect T="int?" @bind-Value="_sourceEdit.EndpointId" Label="Connector" Clearable="true" Class="mb-2">
|
||||
@foreach (var e in _endpoints)
|
||||
if (ConnectorsFor(needed).Count == 0)
|
||||
{
|
||||
<MudSelectItem T="int?" Value="@((int?)e.Id)">@e.Name (@e.Type)</MudSelectItem>
|
||||
<MudAlert Severity="Severity.Warning" Dense="true" Class="mb-2">
|
||||
No @needed connector yet — <MudLink Href="/admin/connectors">create one</MudLink>
|
||||
(set it up once; every source then just picks it).
|
||||
</MudAlert>
|
||||
}
|
||||
else
|
||||
{
|
||||
<MudSelect T="int?" @bind-Value="_sourceEdit.EndpointId" Label="Connector" Required="true" Class="mb-2">
|
||||
@foreach (var e in ConnectorsFor(needed))
|
||||
{
|
||||
<MudSelectItem T="int?" Value="@((int?)e.Id)">@e.Name</MudSelectItem>
|
||||
}
|
||||
</MudSelect>
|
||||
}
|
||||
}
|
||||
@if (_sourceEdit.SourceType == SourceType.HomeAssistant)
|
||||
{
|
||||
<MudTextField @bind-Value="_sourceEdit.EntityId" Label="Entity id (e.g. sensor.house_power)" Class="mb-2" />
|
||||
@@ -305,6 +315,7 @@ else
|
||||
if (source is null)
|
||||
{
|
||||
_sourceEdit = new SourceEdit();
|
||||
OnSourceTypeChanged(_sourceEdit.SourceType);
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -332,6 +343,28 @@ else
|
||||
|
||||
private async Task SaveSourceAsync()
|
||||
{
|
||||
// A live source without a matching connector has no connection details and would silently
|
||||
// never ingest, so refuse it here rather than letting it look configured.
|
||||
if (RequiredEndpointType(_sourceEdit.SourceType) is { } needed)
|
||||
{
|
||||
var selected = _endpoints.FirstOrDefault(e => e.Id == _sourceEdit.EndpointId);
|
||||
if (selected is null)
|
||||
{
|
||||
Snackbar.Add($"Pick a {needed} connector for this {_sourceEdit.SourceType} source.", Severity.Error);
|
||||
return;
|
||||
}
|
||||
|
||||
if (selected.Type != needed)
|
||||
{
|
||||
Snackbar.Add($"'{selected.Name}' is a {selected.Type} connector; a {_sourceEdit.SourceType} source needs {needed}.", Severity.Error);
|
||||
return;
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
_sourceEdit.EndpointId = null;
|
||||
}
|
||||
|
||||
var config = new SourceConfig
|
||||
{
|
||||
EntityId = Trim(_sourceEdit.EntityId),
|
||||
@@ -394,6 +427,44 @@ else
|
||||
|
||||
private static string? Trim(string? value) => string.IsNullOrWhiteSpace(value) ? null : value.Trim();
|
||||
|
||||
/// <summary>
|
||||
/// Which connector kind a source type needs, or null if it needs none (manual/import/virtual).
|
||||
/// Tasmota has no endpoint kind of its own — it is served by an MQTT broker connector.
|
||||
/// </summary>
|
||||
private static EndpointType? RequiredEndpointType(SourceType sourceType) => sourceType switch
|
||||
{
|
||||
SourceType.HomeAssistant => EndpointType.HomeAssistant,
|
||||
SourceType.Mqtt or SourceType.Tasmota => EndpointType.MqttBroker,
|
||||
_ => null,
|
||||
};
|
||||
|
||||
private List<IngestionEndpoint> ConnectorsFor(EndpointType type) =>
|
||||
_endpoints.Where(e => e.Type == type).ToList();
|
||||
|
||||
// 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.
|
||||
private void OnSourceTypeChanged(SourceType sourceType)
|
||||
{
|
||||
_sourceEdit.SourceType = sourceType;
|
||||
|
||||
var needed = RequiredEndpointType(sourceType);
|
||||
var selected = _endpoints.FirstOrDefault(e => e.Id == _sourceEdit.EndpointId);
|
||||
if (needed is null || (selected is not null && selected.Type != needed))
|
||||
{
|
||||
_sourceEdit.EndpointId = null;
|
||||
}
|
||||
|
||||
// Sole candidate: preselect it, so the common single-broker / single-HA setup is one click.
|
||||
if (needed is not null && _sourceEdit.EndpointId is null)
|
||||
{
|
||||
var candidates = _endpoints.Where(e => e.Type == needed).ToList();
|
||||
if (candidates.Count == 1)
|
||||
{
|
||||
_sourceEdit.EndpointId = candidates[0].Id;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private sealed class SourceEdit
|
||||
{
|
||||
public int Id { get; set; }
|
||||
|
||||
@@ -84,7 +84,7 @@ public sealed class MqttIngestionWorker(
|
||||
|
||||
foreach (var endpoint in endpoints)
|
||||
{
|
||||
var client = _clients.GetOrAdd(endpoint.Id, _ => CreateClient());
|
||||
var client = _clients.GetOrAdd(endpoint.Id, id => CreateClient(id));
|
||||
var topics = await ResolveTopicsAsync(db, endpoint, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
if (!client.IsConnected)
|
||||
@@ -114,10 +114,13 @@ public sealed class MqttIngestionWorker(
|
||||
}
|
||||
}
|
||||
|
||||
private IMqttClient CreateClient()
|
||||
// One client per endpoint, with the endpoint id captured in the handler: MQTTnet's event args
|
||||
// carry the topic but not which connection delivered it, and the router needs that to keep
|
||||
// sources bound to one broker from ingesting another's traffic.
|
||||
private IMqttClient CreateClient(int endpointId)
|
||||
{
|
||||
var client = _factory.CreateMqttClient();
|
||||
client.ApplicationMessageReceivedAsync += OnMessageAsync;
|
||||
client.ApplicationMessageReceivedAsync += args => OnMessageAsync(endpointId, args);
|
||||
return client;
|
||||
}
|
||||
|
||||
@@ -163,10 +166,12 @@ public sealed class MqttIngestionWorker(
|
||||
private static async Task<IReadOnlyList<string>> ResolveTopicsAsync(
|
||||
MeterVaultDbContext db, IngestionEndpoint endpoint, CancellationToken cancellationToken)
|
||||
{
|
||||
// Bound sources only, matching MqttMessageRouter: an unbound source is not routed, so
|
||||
// subscribing its topic on every broker would only invite traffic nothing consumes.
|
||||
var sourceConfigs = await db.MeterSources
|
||||
.Where(s => s.IsEnabled
|
||||
&& (s.SourceType == SourceType.Mqtt || s.SourceType == SourceType.Tasmota)
|
||||
&& (s.EndpointId == endpoint.Id || s.EndpointId == null))
|
||||
&& s.EndpointId == endpoint.Id)
|
||||
.Select(s => s.Config)
|
||||
.ToListAsync(cancellationToken).ConfigureAwait(false);
|
||||
|
||||
@@ -188,7 +193,7 @@ public sealed class MqttIngestionWorker(
|
||||
return [.. topics];
|
||||
}
|
||||
|
||||
private async Task OnMessageAsync(MqttApplicationMessageReceivedEventArgs args)
|
||||
private async Task OnMessageAsync(int endpointId, MqttApplicationMessageReceivedEventArgs args)
|
||||
{
|
||||
var topic = args.ApplicationMessage.Topic;
|
||||
var payload = args.ApplicationMessage.ConvertPayloadToString() ?? string.Empty;
|
||||
@@ -197,7 +202,7 @@ public sealed class MqttIngestionWorker(
|
||||
{
|
||||
await using var scope = _scopeFactory.CreateAsyncScope();
|
||||
var router = scope.ServiceProvider.GetRequiredService<MqttMessageRouter>();
|
||||
await router.RouteAsync(topic, payload).ConfigureAwait(false);
|
||||
await router.RouteAsync(endpointId, topic, payload).ConfigureAwait(false);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
|
||||
@@ -6,10 +6,16 @@ using Microsoft.Extensions.Logging;
|
||||
namespace MeterVault.Infrastructure.Ingestion;
|
||||
|
||||
/// <summary>
|
||||
/// Routes an incoming MQTT message to every enabled MQTT/Tasmota source whose topic filter covers
|
||||
/// it, extracts the value (and payload timestamp), and ingests it (SDD §6.1). Decoupled from the
|
||||
/// broker client so it can be exercised directly against the database in tests.
|
||||
/// Routes an incoming MQTT message to every enabled MQTT/Tasmota source that is bound to the
|
||||
/// delivering broker <em>and</em> whose topic filter covers it, extracts the value (and payload
|
||||
/// timestamp), and ingests it (SDD §6.1). Decoupled from the broker client so it can be exercised
|
||||
/// directly against the database in tests.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// The endpoint predicate is load-bearing, not defensive: topic filters routinely overlap between
|
||||
/// brokers (every Tasmota install publishes <c>tele/+/SENSOR</c>), so matching on topic alone would
|
||||
/// let a message from one broker be ingested by a source bound to another.
|
||||
/// </remarks>
|
||||
public sealed class MqttMessageRouter(
|
||||
MeterVaultDbContext db, IngestionService ingestion, ILogger<MqttMessageRouter> logger)
|
||||
{
|
||||
@@ -17,10 +23,13 @@ public sealed class MqttMessageRouter(
|
||||
private readonly IngestionService _ingestion = ingestion;
|
||||
private readonly ILogger<MqttMessageRouter> _logger = logger;
|
||||
|
||||
public async Task<int> RouteAsync(string topic, string payload, CancellationToken cancellationToken = default)
|
||||
public async Task<int> RouteAsync(
|
||||
int endpointId, string topic, string payload, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var sources = await _db.MeterSources
|
||||
.Where(s => s.IsEnabled && (s.SourceType == SourceType.Mqtt || s.SourceType == SourceType.Tasmota))
|
||||
.Where(s => s.IsEnabled
|
||||
&& (s.SourceType == SourceType.Mqtt || s.SourceType == SourceType.Tasmota)
|
||||
&& s.EndpointId == endpointId)
|
||||
.ToListAsync(cancellationToken).ConfigureAwait(false);
|
||||
|
||||
var routed = 0;
|
||||
|
||||
+931
@@ -0,0 +1,931 @@
|
||||
// <auto-generated />
|
||||
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("20260718091623_BindUnboundMqttSourcesToSoleBroker")]
|
||||
partial class BindUnboundMqttSourcesToSoleBroker
|
||||
{
|
||||
/// <inheritdoc />
|
||||
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<string>("Key")
|
||||
.HasMaxLength(128)
|
||||
.HasColumnType("character varying(128)")
|
||||
.HasColumnName("key");
|
||||
|
||||
b.Property<string>("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<int>("MeterId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("meter_id");
|
||||
|
||||
b.Property<DateTimeOffset>("Time")
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("time");
|
||||
|
||||
b.Property<short>("Kind")
|
||||
.HasColumnType("smallint")
|
||||
.HasColumnName("kind");
|
||||
|
||||
b.Property<double>("Amount")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("amount");
|
||||
|
||||
b.Property<int?>("ImportBatchId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("import_batch_id");
|
||||
|
||||
b.Property<short>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<string>("ColorHex")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("color_hex");
|
||||
|
||||
b.Property<string>("Name")
|
||||
.IsRequired()
|
||||
.HasMaxLength(128)
|
||||
.HasColumnType("character varying(128)")
|
||||
.HasColumnName("name");
|
||||
|
||||
b.Property<int>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<int>("CategoryId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("category_id");
|
||||
|
||||
b.Property<short?>("EnergyTypeId")
|
||||
.HasColumnType("smallint")
|
||||
.HasColumnName("energy_type_id");
|
||||
|
||||
b.Property<int?>("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<short>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("smallint")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<short>("Id"));
|
||||
|
||||
b.Property<string>("BaseUnit")
|
||||
.IsRequired()
|
||||
.HasMaxLength(16)
|
||||
.HasColumnType("character varying(16)")
|
||||
.HasColumnName("base_unit");
|
||||
|
||||
b.Property<string>("ColorHex")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("color_hex");
|
||||
|
||||
b.Property<DateTimeOffset>("CreatedAt")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("created_at")
|
||||
.HasDefaultValueSql("now()");
|
||||
|
||||
b.Property<string>("DefaultMode")
|
||||
.IsRequired()
|
||||
.HasMaxLength(32)
|
||||
.HasColumnType("character varying(32)")
|
||||
.HasColumnName("default_mode");
|
||||
|
||||
b.Property<string>("DisplayName")
|
||||
.IsRequired()
|
||||
.HasMaxLength(128)
|
||||
.HasColumnType("character varying(128)")
|
||||
.HasColumnName("display_name");
|
||||
|
||||
b.Property<string>("Icon")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("icon");
|
||||
|
||||
b.Property<string>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<DateTimeOffset>("CreatedAt")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("created_at")
|
||||
.HasDefaultValueSql("now()");
|
||||
|
||||
b.Property<string>("Mapping")
|
||||
.HasColumnType("jsonb")
|
||||
.HasColumnName("mapping");
|
||||
|
||||
b.Property<DateTimeOffset?>("RevertedAt")
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("reverted_at");
|
||||
|
||||
b.Property<int>("RowCount")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("row_count");
|
||||
|
||||
b.Property<string>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<string>("Config")
|
||||
.IsRequired()
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("jsonb")
|
||||
.HasColumnName("config")
|
||||
.HasDefaultValueSql("'{}'::jsonb");
|
||||
|
||||
b.Property<bool>("IsEnabled")
|
||||
.HasColumnType("boolean")
|
||||
.HasColumnName("is_enabled");
|
||||
|
||||
b.Property<DateTimeOffset?>("LastSeenAt")
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("last_seen_at");
|
||||
|
||||
b.Property<string>("LastStatus")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("last_status");
|
||||
|
||||
b.Property<string>("Name")
|
||||
.IsRequired()
|
||||
.HasMaxLength(128)
|
||||
.HasColumnType("character varying(128)")
|
||||
.HasColumnName("name");
|
||||
|
||||
b.Property<string>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<double>("Amount")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("amount");
|
||||
|
||||
b.Property<int?>("CategoryId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("category_id");
|
||||
|
||||
b.Property<string>("Currency")
|
||||
.IsRequired()
|
||||
.HasMaxLength(8)
|
||||
.HasColumnType("character varying(8)")
|
||||
.HasColumnName("currency");
|
||||
|
||||
b.Property<int?>("ImportBatchId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("import_batch_id");
|
||||
|
||||
b.Property<int?>("MeterId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("meter_id");
|
||||
|
||||
b.Property<string>("Notes")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("notes");
|
||||
|
||||
b.Property<DateOnly>("PeriodEnd")
|
||||
.HasColumnType("date")
|
||||
.HasColumnName("period_end");
|
||||
|
||||
b.Property<DateOnly>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<DateTimeOffset>("CreatedAt")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("created_at")
|
||||
.HasDefaultValueSql("now()");
|
||||
|
||||
b.Property<short>("EnergyTypeId")
|
||||
.HasColumnType("smallint")
|
||||
.HasColumnName("energy_type_id");
|
||||
|
||||
b.Property<double>("InitialBaseline")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("initial_baseline");
|
||||
|
||||
b.Property<DateOnly?>("InstalledAt")
|
||||
.HasColumnType("date")
|
||||
.HasColumnName("installed_at");
|
||||
|
||||
b.Property<bool>("IsActive")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("boolean")
|
||||
.HasDefaultValue(true)
|
||||
.HasColumnName("is_active");
|
||||
|
||||
b.Property<string>("Location")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("location");
|
||||
|
||||
b.Property<string>("Manufacturer")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("manufacturer");
|
||||
|
||||
b.Property<string>("Meta")
|
||||
.IsRequired()
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("jsonb")
|
||||
.HasColumnName("meta")
|
||||
.HasDefaultValueSql("'{}'::jsonb");
|
||||
|
||||
b.Property<string>("Mode")
|
||||
.IsRequired()
|
||||
.HasMaxLength(32)
|
||||
.HasColumnType("character varying(32)")
|
||||
.HasColumnName("mode");
|
||||
|
||||
b.Property<string>("Model")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("model");
|
||||
|
||||
b.Property<string>("Name")
|
||||
.IsRequired()
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("name");
|
||||
|
||||
b.Property<DateOnly?>("RetiredAt")
|
||||
.HasColumnType("date")
|
||||
.HasColumnName("retired_at");
|
||||
|
||||
b.Property<string>("SerialNumber")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("serial_number");
|
||||
|
||||
b.Property<string>("Unit")
|
||||
.IsRequired()
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("unit");
|
||||
|
||||
b.Property<DateTimeOffset>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<double?>("Amount")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("amount");
|
||||
|
||||
b.Property<string>("EventType")
|
||||
.IsRequired()
|
||||
.HasMaxLength(32)
|
||||
.HasColumnType("character varying(32)")
|
||||
.HasColumnName("event_type");
|
||||
|
||||
b.Property<int?>("ImportBatchId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("import_batch_id");
|
||||
|
||||
b.Property<string>("Meta")
|
||||
.IsRequired()
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("jsonb")
|
||||
.HasColumnName("meta")
|
||||
.HasDefaultValueSql("'{}'::jsonb");
|
||||
|
||||
b.Property<int>("MeterId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("meter_id");
|
||||
|
||||
b.Property<double?>("NewValue")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("new_value");
|
||||
|
||||
b.Property<string>("Notes")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("notes");
|
||||
|
||||
b.Property<double?>("PrevValue")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("prev_value");
|
||||
|
||||
b.Property<DateTimeOffset>("Time")
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("time");
|
||||
|
||||
b.Property<string>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<int>("FromMeterId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("from_meter_id");
|
||||
|
||||
b.Property<int>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<string>("Config")
|
||||
.IsRequired()
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("jsonb")
|
||||
.HasColumnName("config")
|
||||
.HasDefaultValueSql("'{}'::jsonb");
|
||||
|
||||
b.Property<int?>("EndpointId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("endpoint_id");
|
||||
|
||||
b.Property<bool>("IsEnabled")
|
||||
.HasColumnType("boolean")
|
||||
.HasColumnName("is_enabled");
|
||||
|
||||
b.Property<DateTimeOffset?>("LastSeenAt")
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("last_seen_at");
|
||||
|
||||
b.Property<string>("LastStatus")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("last_status");
|
||||
|
||||
b.Property<double?>("LastValue")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("last_value");
|
||||
|
||||
b.Property<int>("MeterId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("meter_id");
|
||||
|
||||
b.Property<double>("Offset")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("offset");
|
||||
|
||||
b.Property<int>("Priority")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("priority");
|
||||
|
||||
b.Property<double>("Scale")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("double precision")
|
||||
.HasDefaultValue(1.0)
|
||||
.HasColumnName("scale");
|
||||
|
||||
b.Property<string>("SourceType")
|
||||
.IsRequired()
|
||||
.HasMaxLength(32)
|
||||
.HasColumnType("character varying(32)")
|
||||
.HasColumnName("source_type");
|
||||
|
||||
b.Property<string>("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<int>("MeterId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("meter_id");
|
||||
|
||||
b.Property<DateTimeOffset>("Time")
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("time");
|
||||
|
||||
b.Property<int>("Flags")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("flags");
|
||||
|
||||
b.Property<int?>("ImportBatchId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("import_batch_id");
|
||||
|
||||
b.Property<short>("Quality")
|
||||
.HasColumnType("smallint")
|
||||
.HasColumnName("quality");
|
||||
|
||||
b.Property<int?>("SourceId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("source_id");
|
||||
|
||||
b.Property<double>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<DateTimeOffset?>("CachedAt")
|
||||
.HasColumnType("timestamp with time zone")
|
||||
.HasColumnName("cached_at");
|
||||
|
||||
b.Property<double?>("CachedBalance")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("cached_balance");
|
||||
|
||||
b.Property<string>("Calibration")
|
||||
.HasColumnType("jsonb")
|
||||
.HasColumnName("calibration");
|
||||
|
||||
b.Property<double>("Capacity")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("capacity");
|
||||
|
||||
b.Property<double?>("FixedRate")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("fixed_rate");
|
||||
|
||||
b.Property<double?>("LowThreshold")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("low_threshold");
|
||||
|
||||
b.Property<int>("MeterId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("meter_id");
|
||||
|
||||
b.Property<string>("RateMode")
|
||||
.IsRequired()
|
||||
.HasMaxLength(16)
|
||||
.HasColumnType("character varying(16)")
|
||||
.HasColumnName("rate_mode");
|
||||
|
||||
b.Property<double?>("ReorderThreshold")
|
||||
.HasColumnType("double precision")
|
||||
.HasColumnName("reorder_threshold");
|
||||
|
||||
b.Property<string>("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<int>("Id")
|
||||
.ValueGeneratedOnAdd()
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("id");
|
||||
|
||||
NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property<int>("Id"));
|
||||
|
||||
b.Property<string>("Component")
|
||||
.IsRequired()
|
||||
.HasMaxLength(16)
|
||||
.HasColumnType("character varying(16)")
|
||||
.HasColumnName("component");
|
||||
|
||||
b.Property<string>("Currency")
|
||||
.IsRequired()
|
||||
.HasMaxLength(8)
|
||||
.HasColumnType("character varying(8)")
|
||||
.HasColumnName("currency");
|
||||
|
||||
b.Property<string>("Notes")
|
||||
.HasColumnType("text")
|
||||
.HasColumnName("notes");
|
||||
|
||||
b.Property<int?>("ScopeId")
|
||||
.HasColumnType("integer")
|
||||
.HasColumnName("scope_id");
|
||||
|
||||
b.Property<string>("ScopeType")
|
||||
.IsRequired()
|
||||
.HasMaxLength(16)
|
||||
.HasColumnType("character varying(16)")
|
||||
.HasColumnName("scope_type");
|
||||
|
||||
b.Property<string>("Unit")
|
||||
.IsRequired()
|
||||
.HasMaxLength(16)
|
||||
.HasColumnType("character varying(16)")
|
||||
.HasColumnName("unit");
|
||||
|
||||
b.Property<DateOnly>("ValidFrom")
|
||||
.HasColumnType("date")
|
||||
.HasColumnName("valid_from");
|
||||
|
||||
b.Property<DateOnly?>("ValidTo")
|
||||
.HasColumnType("date")
|
||||
.HasColumnName("valid_to");
|
||||
|
||||
b.Property<double>("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
|
||||
}
|
||||
}
|
||||
}
|
||||
+47
@@ -0,0 +1,47 @@
|
||||
using Microsoft.EntityFrameworkCore.Migrations;
|
||||
|
||||
#nullable disable
|
||||
|
||||
namespace MeterVault.Infrastructure.Persistence.Migrations
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public partial class BindUnboundMqttSourcesToSoleBroker : Migration
|
||||
{
|
||||
/// <inheritdoc />
|
||||
protected override void Up(MigrationBuilder migrationBuilder)
|
||||
{
|
||||
// MQTT routing now honours meter_source.endpoint_id (SDD §6.1): a source is served only
|
||||
// by the broker it is bound to. Previously an unbound source was subscribed on every
|
||||
// broker and matched on topic alone, so unbound sources that work today would silently
|
||||
// go quiet after this deploy.
|
||||
//
|
||||
// Backfill them onto the single broker only when exactly one exists — then the old
|
||||
// "any broker" behaviour and the new "its broker" behaviour are the same thing, so the
|
||||
// rewrite is provably lossless. With zero brokers there is nothing to bind to; with two
|
||||
// or more the old behaviour was already ambiguous and a guess could route a meter's
|
||||
// data to the wrong broker, so those are left for the operator to resolve in the UI.
|
||||
//
|
||||
// Enums persist as their C# names (HasConversion<string>), hence 'Mqtt'/'MqttBroker'.
|
||||
// HomeAssistant sources are deliberately excluded: the HA workers have always required
|
||||
// endpoint_id, so an unbound HA source is already inert and binding it here would
|
||||
// activate ingestion the operator never had running.
|
||||
migrationBuilder.Sql("""
|
||||
UPDATE meter_source AS s
|
||||
SET endpoint_id = sole.id
|
||||
FROM (SELECT id FROM ingestion_endpoint WHERE type = 'MqttBroker') AS sole
|
||||
WHERE s.endpoint_id IS NULL
|
||||
AND s.source_type IN ('Mqtt', 'Tasmota')
|
||||
AND (SELECT count(*) FROM ingestion_endpoint WHERE type = 'MqttBroker') = 1;
|
||||
""");
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
protected override void Down(MigrationBuilder migrationBuilder)
|
||||
{
|
||||
// Intentionally empty. The rows this bound are indistinguishable from ones the operator
|
||||
// bound by hand, so clearing endpoint_id on the way down would discard real
|
||||
// configuration. Leaving the binding in place is harmless under the old routing, which
|
||||
// ignored endpoint_id entirely.
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -100,10 +100,13 @@ public sealed class IngestionServiceTests(TimescaleFixture fx)
|
||||
public async Task Mqtt_router_ingests_a_tasmota_payload()
|
||||
{
|
||||
await using var db = fx.CreateContext();
|
||||
var (meterId, _) = await SetupAsync(db, MeterMode.CumulativeCounter, topic: "tele/plug7/SENSOR");
|
||||
var brokerId = await CreateBrokerAsync(db);
|
||||
var (meterId, _) = await SetupAsync(
|
||||
db, MeterMode.CumulativeCounter, topic: "tele/plug7/SENSOR", endpointId: brokerId);
|
||||
var router = new MqttMessageRouter(db, new IngestionService(db), NullLogger<MqttMessageRouter>.Instance);
|
||||
|
||||
var routed = await router.RouteAsync(
|
||||
brokerId,
|
||||
"tele/plug7/SENSOR",
|
||||
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}""");
|
||||
|
||||
@@ -115,9 +118,39 @@ public sealed class IngestionServiceTests(TimescaleFixture fx)
|
||||
await CleanupAsync(db, meterId);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task Mqtt_router_ignores_a_source_bound_to_another_broker()
|
||||
{
|
||||
await using var db = fx.CreateContext();
|
||||
var brokerA = await CreateBrokerAsync(db);
|
||||
var brokerB = await CreateBrokerAsync(db);
|
||||
|
||||
// Topic filter that both brokers' traffic would match — the binding is the only thing
|
||||
// separating them.
|
||||
var (meterId, _) = await SetupAsync(
|
||||
db, MeterMode.CumulativeCounter, topic: "tele/+/SENSOR", endpointId: brokerB);
|
||||
var router = new MqttMessageRouter(db, new IngestionService(db), NullLogger<MqttMessageRouter>.Instance);
|
||||
|
||||
var routed = await router.RouteAsync(
|
||||
brokerA,
|
||||
"tele/plug7/SENSOR",
|
||||
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}""");
|
||||
|
||||
Assert.Equal(0, routed);
|
||||
Assert.False(await db.Readings.AnyAsync(r => r.MeterId == meterId));
|
||||
|
||||
// Same message on the broker it is actually bound to does land.
|
||||
Assert.Equal(1, await router.RouteAsync(
|
||||
brokerB,
|
||||
"tele/plug7/SENSOR",
|
||||
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}"""));
|
||||
|
||||
await CleanupAsync(db, meterId);
|
||||
}
|
||||
|
||||
private static async Task<(int MeterId, int SourceId)> SetupAsync(
|
||||
MeterVaultDbContext db, MeterMode mode, double scale = 1, double offset = 0,
|
||||
string topic = "tele/x/SENSOR", string? path = "ENERGY.Total")
|
||||
string topic = "tele/x/SENSOR", string? path = "ENERGY.Total", int? endpointId = null)
|
||||
{
|
||||
await DatabaseSeeder.SeedAsync(db);
|
||||
var type = await db.EnergyTypes.FirstAsync(t => t.Key == "electricity");
|
||||
@@ -136,6 +169,7 @@ public sealed class IngestionServiceTests(TimescaleFixture fx)
|
||||
{
|
||||
MeterId = meter.Id,
|
||||
SourceType = SourceType.Tasmota,
|
||||
EndpointId = endpointId ?? await CreateBrokerAsync(db),
|
||||
ValueKind = SourceValueKind.Register,
|
||||
Scale = scale,
|
||||
Offset = offset,
|
||||
@@ -147,10 +181,24 @@ public sealed class IngestionServiceTests(TimescaleFixture fx)
|
||||
return (meter.Id, source.Id);
|
||||
}
|
||||
|
||||
private static async Task<int> CreateBrokerAsync(MeterVaultDbContext db)
|
||||
{
|
||||
var endpoint = new IngestionEndpoint
|
||||
{
|
||||
Type = EndpointType.MqttBroker,
|
||||
Name = $"broker-{Guid.NewGuid():N}",
|
||||
Config = """{"host":"localhost","port":1883}""",
|
||||
};
|
||||
db.IngestionEndpoints.Add(endpoint);
|
||||
await db.SaveChangesAsync();
|
||||
return endpoint.Id;
|
||||
}
|
||||
|
||||
private static async Task CleanupAsync(MeterVaultDbContext db, int meterId)
|
||||
{
|
||||
await db.Readings.Where(r => r.MeterId == meterId).ExecuteDeleteAsync();
|
||||
await db.MeterEvents.Where(e => e.MeterId == meterId).ExecuteDeleteAsync();
|
||||
await db.Meters.Where(m => m.Id == meterId).ExecuteDeleteAsync();
|
||||
await db.IngestionEndpoints.Where(e => e.Name.StartsWith("broker-")).ExecuteDeleteAsync();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user