diff --git a/.changeset/energy-identity-counter-continuity.md b/.changeset/energy-identity-counter-continuity.md new file mode 100644 index 000000000..fddf794db --- /dev/null +++ b/.changeset/energy-identity-counter-continuity.md @@ -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. diff --git a/go/cmd/ftw/energy_history.go b/go/cmd/ftw/energy_history.go index ad5453aac..0dc51dcf2 100644 --- a/go/cmd/ftw/energy_history.go +++ b/go/cmd/ftw/energy_history.go @@ -2,6 +2,7 @@ package main import ( "encoding/json" + "log/slog" "math" "github.com/srcfl/ftw/go/internal/state" @@ -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, "", "") != "" { @@ -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 == "" { diff --git a/go/cmd/ftw/energy_identity_test.go b/go/cmd/ftw/energy_identity_test.go index 30baf1033..920429e24 100644 --- a/go/cmd/ftw/energy_identity_test.go +++ b/go/cmd/ftw/energy_identity_test.go @@ -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) + } +} diff --git a/go/internal/state/energy_ledger.go b/go/internal/state/energy_ledger.go index 7aa1427ef..926f6630d 100644 --- a/go/internal/state/energy_ledger.go +++ b/go/internal/state/energy_ledger.go @@ -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,