Compare commits
2
Commits
813d96709e
...
c0bbaba99f
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c0bbaba99f | ||
|
|
9bd0d60cc8 |
@@ -214,20 +214,30 @@ else
|
|||||||
<MudText Typo="Typo.h6">@(_sourceEdit.Id == 0 ? "New source" : "Edit source")</MudText>
|
<MudText Typo="Typo.h6">@(_sourceEdit.Id == 0 ? "New source" : "Edit source")</MudText>
|
||||||
</TitleContent>
|
</TitleContent>
|
||||||
<DialogContent>
|
<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>())
|
@foreach (var type in Enum.GetValues<SourceType>())
|
||||||
{
|
{
|
||||||
<MudSelectItem T="SourceType" Value="type">@type</MudSelectItem>
|
<MudSelectItem T="SourceType" Value="type">@type</MudSelectItem>
|
||||||
}
|
}
|
||||||
</MudSelect>
|
</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">
|
if (ConnectorsFor(needed).Count == 0)
|
||||||
@foreach (var e in _endpoints)
|
{
|
||||||
{
|
<MudAlert Severity="Severity.Warning" Dense="true" Class="mb-2">
|
||||||
<MudSelectItem T="int?" Value="@((int?)e.Id)">@e.Name (@e.Type)</MudSelectItem>
|
No @needed connector yet — <MudLink Href="/admin/connectors">create one</MudLink>
|
||||||
}
|
(set it up once; every source then just picks it).
|
||||||
</MudSelect>
|
</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)
|
@if (_sourceEdit.SourceType == SourceType.HomeAssistant)
|
||||||
{
|
{
|
||||||
@@ -305,6 +315,7 @@ else
|
|||||||
if (source is null)
|
if (source is null)
|
||||||
{
|
{
|
||||||
_sourceEdit = new SourceEdit();
|
_sourceEdit = new SourceEdit();
|
||||||
|
OnSourceTypeChanged(_sourceEdit.SourceType);
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
{
|
{
|
||||||
@@ -332,6 +343,28 @@ else
|
|||||||
|
|
||||||
private async Task SaveSourceAsync()
|
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
|
var config = new SourceConfig
|
||||||
{
|
{
|
||||||
EntityId = Trim(_sourceEdit.EntityId),
|
EntityId = Trim(_sourceEdit.EntityId),
|
||||||
@@ -394,6 +427,44 @@ else
|
|||||||
|
|
||||||
private static string? Trim(string? value) => string.IsNullOrWhiteSpace(value) ? null : value.Trim();
|
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
|
private sealed class SourceEdit
|
||||||
{
|
{
|
||||||
public int Id { get; set; }
|
public int Id { get; set; }
|
||||||
|
|||||||
@@ -84,7 +84,7 @@ public sealed class MqttIngestionWorker(
|
|||||||
|
|
||||||
foreach (var endpoint in endpoints)
|
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);
|
var topics = await ResolveTopicsAsync(db, endpoint, cancellationToken).ConfigureAwait(false);
|
||||||
|
|
||||||
if (!client.IsConnected)
|
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();
|
var client = _factory.CreateMqttClient();
|
||||||
client.ApplicationMessageReceivedAsync += OnMessageAsync;
|
client.ApplicationMessageReceivedAsync += args => OnMessageAsync(endpointId, args);
|
||||||
return client;
|
return client;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -163,10 +166,12 @@ public sealed class MqttIngestionWorker(
|
|||||||
private static async Task<IReadOnlyList<string>> ResolveTopicsAsync(
|
private static async Task<IReadOnlyList<string>> ResolveTopicsAsync(
|
||||||
MeterVaultDbContext db, IngestionEndpoint endpoint, CancellationToken cancellationToken)
|
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
|
var sourceConfigs = await db.MeterSources
|
||||||
.Where(s => s.IsEnabled
|
.Where(s => s.IsEnabled
|
||||||
&& (s.SourceType == SourceType.Mqtt || s.SourceType == SourceType.Tasmota)
|
&& (s.SourceType == SourceType.Mqtt || s.SourceType == SourceType.Tasmota)
|
||||||
&& (s.EndpointId == endpoint.Id || s.EndpointId == null))
|
&& s.EndpointId == endpoint.Id)
|
||||||
.Select(s => s.Config)
|
.Select(s => s.Config)
|
||||||
.ToListAsync(cancellationToken).ConfigureAwait(false);
|
.ToListAsync(cancellationToken).ConfigureAwait(false);
|
||||||
|
|
||||||
@@ -188,7 +193,7 @@ public sealed class MqttIngestionWorker(
|
|||||||
return [.. topics];
|
return [.. topics];
|
||||||
}
|
}
|
||||||
|
|
||||||
private async Task OnMessageAsync(MqttApplicationMessageReceivedEventArgs args)
|
private async Task OnMessageAsync(int endpointId, MqttApplicationMessageReceivedEventArgs args)
|
||||||
{
|
{
|
||||||
var topic = args.ApplicationMessage.Topic;
|
var topic = args.ApplicationMessage.Topic;
|
||||||
var payload = args.ApplicationMessage.ConvertPayloadToString() ?? string.Empty;
|
var payload = args.ApplicationMessage.ConvertPayloadToString() ?? string.Empty;
|
||||||
@@ -197,7 +202,7 @@ public sealed class MqttIngestionWorker(
|
|||||||
{
|
{
|
||||||
await using var scope = _scopeFactory.CreateAsyncScope();
|
await using var scope = _scopeFactory.CreateAsyncScope();
|
||||||
var router = scope.ServiceProvider.GetRequiredService<MqttMessageRouter>();
|
var router = scope.ServiceProvider.GetRequiredService<MqttMessageRouter>();
|
||||||
await router.RouteAsync(topic, payload).ConfigureAwait(false);
|
await router.RouteAsync(endpointId, topic, payload).ConfigureAwait(false);
|
||||||
}
|
}
|
||||||
catch (Exception ex)
|
catch (Exception ex)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -6,10 +6,16 @@ using Microsoft.Extensions.Logging;
|
|||||||
namespace MeterVault.Infrastructure.Ingestion;
|
namespace MeterVault.Infrastructure.Ingestion;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// Routes an incoming MQTT message to every enabled MQTT/Tasmota source whose topic filter covers
|
/// Routes an incoming MQTT message to every enabled MQTT/Tasmota source that is bound to the
|
||||||
/// it, extracts the value (and payload timestamp), and ingests it (SDD §6.1). Decoupled from the
|
/// delivering broker <em>and</em> whose topic filter covers it, extracts the value (and payload
|
||||||
/// broker client so it can be exercised directly against the database in tests.
|
/// timestamp), and ingests it (SDD §6.1). Decoupled from the broker client so it can be exercised
|
||||||
|
/// directly against the database in tests.
|
||||||
/// </summary>
|
/// </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(
|
public sealed class MqttMessageRouter(
|
||||||
MeterVaultDbContext db, IngestionService ingestion, ILogger<MqttMessageRouter> logger)
|
MeterVaultDbContext db, IngestionService ingestion, ILogger<MqttMessageRouter> logger)
|
||||||
{
|
{
|
||||||
@@ -17,10 +23,13 @@ public sealed class MqttMessageRouter(
|
|||||||
private readonly IngestionService _ingestion = ingestion;
|
private readonly IngestionService _ingestion = ingestion;
|
||||||
private readonly ILogger<MqttMessageRouter> _logger = logger;
|
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
|
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);
|
.ToListAsync(cancellationToken).ConfigureAwait(false);
|
||||||
|
|
||||||
var routed = 0;
|
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()
|
public async Task Mqtt_router_ingests_a_tasmota_payload()
|
||||||
{
|
{
|
||||||
await using var db = fx.CreateContext();
|
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 router = new MqttMessageRouter(db, new IngestionService(db), NullLogger<MqttMessageRouter>.Instance);
|
||||||
|
|
||||||
var routed = await router.RouteAsync(
|
var routed = await router.RouteAsync(
|
||||||
|
brokerId,
|
||||||
"tele/plug7/SENSOR",
|
"tele/plug7/SENSOR",
|
||||||
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}""");
|
"""{"Time":"2024-03-01T10:00:00","ENERGY":{"Total":8421.0}}""");
|
||||||
|
|
||||||
@@ -115,9 +118,39 @@ public sealed class IngestionServiceTests(TimescaleFixture fx)
|
|||||||
await CleanupAsync(db, meterId);
|
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(
|
private static async Task<(int MeterId, int SourceId)> SetupAsync(
|
||||||
MeterVaultDbContext db, MeterMode mode, double scale = 1, double offset = 0,
|
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);
|
await DatabaseSeeder.SeedAsync(db);
|
||||||
var type = await db.EnergyTypes.FirstAsync(t => t.Key == "electricity");
|
var type = await db.EnergyTypes.FirstAsync(t => t.Key == "electricity");
|
||||||
@@ -136,6 +169,7 @@ public sealed class IngestionServiceTests(TimescaleFixture fx)
|
|||||||
{
|
{
|
||||||
MeterId = meter.Id,
|
MeterId = meter.Id,
|
||||||
SourceType = SourceType.Tasmota,
|
SourceType = SourceType.Tasmota,
|
||||||
|
EndpointId = endpointId ?? await CreateBrokerAsync(db),
|
||||||
ValueKind = SourceValueKind.Register,
|
ValueKind = SourceValueKind.Register,
|
||||||
Scale = scale,
|
Scale = scale,
|
||||||
Offset = offset,
|
Offset = offset,
|
||||||
@@ -147,10 +181,24 @@ public sealed class IngestionServiceTests(TimescaleFixture fx)
|
|||||||
return (meter.Id, source.Id);
|
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)
|
private static async Task CleanupAsync(MeterVaultDbContext db, int meterId)
|
||||||
{
|
{
|
||||||
await db.Readings.Where(r => r.MeterId == meterId).ExecuteDeleteAsync();
|
await db.Readings.Where(r => r.MeterId == meterId).ExecuteDeleteAsync();
|
||||||
await db.MeterEvents.Where(e => e.MeterId == meterId).ExecuteDeleteAsync();
|
await db.MeterEvents.Where(e => e.MeterId == meterId).ExecuteDeleteAsync();
|
||||||
await db.Meters.Where(m => m.Id == 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