Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/energy-identity-counter-continuity.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"ftw": patch
---

Preserve cumulative energy-counter continuity when a running device refines its identity from a MAC address to a confirmed hardware serial, so outages can still be recovered as counter gaps instead of starting a new baseline.
35 changes: 34 additions & 1 deletion go/cmd/ftw/energy_history.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package main

import (
"encoding/json"
"log/slog"
"math"

"github.com/srcfl/ftw/go/internal/state"
Expand All @@ -15,6 +16,32 @@ type energyIdentityLookup func(string) state.Device
// A known serial must be confirmed by the running driver before a weaker
// alias can write counters again. Another MAC or a newly reported serial is
// a replacement, not evidence that it owns the old device's counter history.
// energyCursorAlias returns a weaker identity that belongs to the same hardware
// and may safely seed ledger cursors when the live driver learns a serial.
// A prior different serial on the same MAC is a replacement and must never
// inherit the old counter state.
func energyCursorAlias(current state.Device, known []state.Device) string {
strongID := state.ResolveDeviceID(current.Make, current.Serial, "", "")
macID := state.ResolveDeviceID("", "", current.MAC, "")
if strongID == "" || macID == "" {
return ""
}
for _, previous := range known {
previousStrong := state.ResolveDeviceID(previous.Make, previous.Serial, "", "")
previousMAC := state.ResolveDeviceID("", "", previous.MAC, "")
if previousStrong != "" && previousStrong != strongID && previousMAC == macID {
return ""
}
}
for _, previous := range known {
if state.ResolveDeviceID(previous.Make, previous.Serial, "", "") == "" &&
state.ResolveDeviceID("", "", previous.MAC, "") == macID {
return macID
}
}
return ""
}

func confirmedEnergyDeviceID(current state.Device, known []state.Device) string {
id := state.ResolveDeviceID(current.Make, current.Serial, current.MAC, current.Endpoint)
if id == "" || state.ResolveDeviceID(current.Make, current.Serial, "", "") != "" {
Expand Down Expand Up @@ -62,7 +89,13 @@ func buildEnergyObservations(st *state.Store, tel *telemetry.Store, ctrl tickPer
asset := func(driver string, kind state.EnergyAssetKind) (string, string, bool) {
id, resolved := ids[driver]
if !resolved && identity != nil {
id = confirmedEnergyDeviceID(identity(driver), known)
current := identity(driver)
id = confirmedEnergyDeviceID(current, known)
if alias := energyCursorAlias(current, known); id != "" && alias != "" && alias != id {
if _, err := st.SeedEnergyCursorsFromDeviceAlias(alias, id); err != nil {
slog.Warn("energy cursor identity migration failed", "driver", driver, "from", alias, "to", id, "err", err)
}
}
ids[driver] = id
}
if id == "" {
Expand Down
56 changes: 56 additions & 0 deletions go/cmd/ftw/energy_identity_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -131,3 +131,59 @@ func TestStartupAliasCannotReplayCountersWhileRawHistoryKeepsWriting(t *testing.
}
}
}

func TestSerialRefinementSeedsMacEnergyCursorAndRecoversGap(t *testing.T) {
st, err := state.Open(filepath.Join(t.TempDir(), "state.db"))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = st.Close() })
mac := state.Device{DriverName: "meter", Make: "Pixii", MAC: "aa:bb:cc:dd:ee:ff", Endpoint: "modbus://site"}
if _, err := st.RegisterDevice(mac); err != nil {
t.Fatal(err)
}
oldAt := time.Now().Add(-24 * time.Hour).Truncate(time.Second)
oldCounter := 1000.0
macAsset := state.HardwareEnergyAssetID("mac:aabbccddeeff", state.AssetGridMeter)
if err := st.EnqueueTelemetryTick(nil, nil, []state.EnergyObservation{{
AssetID: macAsset, DeviceID: "mac:aabbccddeeff", AssetKind: state.AssetGridMeter,
Flow: state.FlowGridImport, AtMs: oldAt.UnixMilli(), CounterWh: &oldCounter,
}}); err != nil {
t.Fatal(err)
}
if err := st.FlushHistory(context.Background()); err != nil {
t.Fatal(err)
}

serial := state.Device{DriverName: "meter", Make: "Pixii", Serial: "235101100376", MAC: mac.MAC, Endpoint: mac.Endpoint}
if alias := energyCursorAlias(serial, st.CachedDevices()); alias != "mac:aabbccddeeff" {
t.Fatalf("alias = %q, want MAC identity", alias)
}
if _, err := st.SeedEnergyCursorsFromDeviceAlias("mac:aabbccddeeff", "pixii:235101100376"); err != nil {
t.Fatal(err)
}
newCounter := 1600.0
newAt := oldAt.Add(24 * time.Hour)
serialAsset := state.HardwareEnergyAssetID("pixii:235101100376", state.AssetGridMeter)
if err := st.EnqueueTelemetryTick(nil, nil, []state.EnergyObservation{{
AssetID: serialAsset, DeviceID: "pixii:235101100376", AssetKind: state.AssetGridMeter,
Flow: state.FlowGridImport, AtMs: newAt.UnixMilli(), CounterWh: &newCounter,
}}); err != nil {
t.Fatal(err)
}
if err := st.FlushHistory(context.Background()); err != nil {
t.Fatal(err)
}
start := oldAt.UnixMilli() / state.EnergyLedgerBucketMS * state.EnergyLedgerBucketMS
points, _, err := st.LoadEnergyHistory(state.EnergyHistoryQuery{AssetID: serialAsset, SinceMS: start, UntilMS: newAt.Add(time.Minute).UnixMilli(), BucketMS: state.EnergyLedgerBucketMS, Limit: 400})
if err != nil {
t.Fatal(err)
}
var total float64
for _, point := range points {
total += point.EnergyWh
}
if total < 599.999 || total > 600.001 {
t.Fatalf("recovered serial energy = %g Wh, want 600 Wh", total)
}
}
58 changes: 58 additions & 0 deletions go/internal/state/energy_ledger.go
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,64 @@ func HardwareEnergyAssetID(deviceID string, kind EnergyAssetKind) string {
return energyAssetID(deviceID, kind)
}

// SeedEnergyCursorsFromDeviceAlias carries cumulative-counter state from a
// weaker identity (for example mac:...) to a newly confirmed hardware identity
// (for example vendor:serial). Existing target cursors always win, making the
// operation idempotent and preventing a late alias from rewinding live state.
// Historical ledger entries stay on their original asset IDs; only the cursor
// continuity needed to recover the next counter gap is transferred.
func (s *Store) SeedEnergyCursorsFromDeviceAlias(fromDeviceID, toDeviceID string) (int64, error) {
fromDeviceID = strings.TrimSpace(fromDeviceID)
toDeviceID = strings.TrimSpace(toDeviceID)
if fromDeviceID == "" || toDeviceID == "" || fromDeviceID == toDeviceID {
return 0, nil
}
tx, err := s.history.Begin()
if err != nil {
return 0, err
}
defer tx.Rollback()
rows, err := tx.Query(`SELECT asset_id, kind FROM energy_assets WHERE device_id = ?`, fromDeviceID)
if err != nil {
return 0, err
}
type aliasAsset struct {
assetID string
kind EnergyAssetKind
}
var assets []aliasAsset
for rows.Next() {
var a aliasAsset
if err := rows.Scan(&a.assetID, &a.kind); err != nil {
rows.Close()
return 0, err
}
assets = append(assets, a)
}
if err := rows.Close(); err != nil {
return 0, err
}
var seeded int64
for _, a := range assets {
toAssetID := HardwareEnergyAssetID(toDeviceID, a.kind)
result, err := tx.Exec(`INSERT OR IGNORE INTO energy_ledger_cursors(asset_id, flow, cursor_kind, value, ts_ms)
SELECT ?, flow, cursor_kind, value, ts_ms
FROM energy_ledger_cursors WHERE asset_id = ?`, toAssetID, a.assetID)
if err != nil {
return 0, err
}
n, err := result.RowsAffected()
if err != nil {
return 0, err
}
seeded += n
}
if err := tx.Commit(); err != nil {
return 0, err
}
return seeded, nil
}

func validEnergyFlow(flow EnergyFlow) bool {
switch flow {
case FlowGridImport, FlowGridExport, FlowBatteryCharge, FlowBatteryDischarge,
Expand Down
Loading