diff --git a/.changeset/tesla-telemetry-recovery.md b/.changeset/tesla-telemetry-recovery.md new file mode 100644 index 000000000..38b6d47fa --- /dev/null +++ b/.changeset/tesla-telemetry-recovery.md @@ -0,0 +1,5 @@ +--- +"ftw": patch +--- + +Vehicle telemetry refresh no longer needs a fresh SoC to find its recipient. Core retains vehicle source time and saves a wake budget across restarts. The release includes Tesla BLE driver 0.2.5 for bounded recovery; charging commands keep their existing freshness checks. It also carries the merged Heishamon 0.8.0 and MyUplink 1.2.3 driver updates. diff --git a/docs/writing-a-driver.md b/docs/writing-a-driver.md index bd85a08d9..1635f68ff 100644 --- a/docs/writing-a-driver.md +++ b/docs/writing-a-driver.md @@ -208,3 +208,27 @@ For live work, start with telemetry only and a physically supervised device. Compare FTW, vendor UI and the site meter before sending a non-zero command. Test charge, discharge, zero, offline/default mode and reconnect. Record device-specific safety knowledge in the driver next to the code it constrains. + +### Vehicle telemetry recovery + +A vehicle driver may declare `DRIVER.telemetry_wake = true` for a `wake_up` +command that only requests telemetry. Core selects it from configuration, +without requiring a fresh SoC. Without an explicit vehicle-to-loadpoint +binding, refresh requires one enabled vehicle and one connected loadpoint. +Read-only and observe-only drivers cannot receive it. Planning and +`charge_start` still require fresh vehicle evidence. + +`host.unix_ms()` returns Unix milliseconds. `host.millis()` and `host.now_ms()` +remain process uptime. A vehicle emit may include `soc_observed_at_ms` and +`soc_fresh`; Core retains the source time and rejects invalid, future and +repeated timestamps as new SoC observations. Old drivers that omit source +time keep their receipt-time behavior. + +Call `host.reserve_vehicle_wake()` before a telemetry wake. It returns +`allowed, retry_ms, err` and saves the attempt before the network call, keyed +by the make and serial set by the driver. The limit is three attempts per +30-minute window, at least 90 seconds apart. Reload, rename, failed HTTP and +Core restart do not reset it. Missing identity, storage failure and hosts +without the function must prevent recovery wakes. This API reserves a budget; +it does not grant a network or command capability. Proxy cache hits do not +prove a new observation, even when HTTP succeeds. diff --git a/drivers/BUNDLED_SOURCE.json b/drivers/BUNDLED_SOURCE.json index 612e290c8..32ba7bbc1 100644 --- a/drivers/BUNDLED_SOURCE.json +++ b/drivers/BUNDLED_SOURCE.json @@ -16,7 +16,7 @@ "from the signed channel. Run scripts/sync-bundled-drivers.sh to update." ], "repository": "srcfl/device-drivers", - "commit": "7c3440c9d10faedfd554f1c795cdbf7b197f3339", + "commit": "c9de25106b07071fde55c9f8f5e95b69c58e6fc9", "source_dir": "drivers/lua", "drivers": [ "ambibox_v2x", "ctek", "ctek_hybrid", "ctek_v2", "deye", "easee_cloud", diff --git a/go/cmd/ftw/driver_registry.go b/go/cmd/ftw/driver_registry.go index f96194d5e..5a208d231 100644 --- a/go/cmd/ftw/driver_registry.go +++ b/go/cmd/ftw/driver_registry.go @@ -1,6 +1,8 @@ package main import ( + "time" + "github.com/srcfl/ftw/go/internal/config" "github.com/srcfl/ftw/go/internal/drivers" "github.com/srcfl/ftw/go/internal/state" @@ -21,5 +23,8 @@ func newDriverRegistry(tel *telemetry.Store, st *state.Store) *drivers.Registry reg.SecretOverride = func(owner, key string) (string, bool) { return st.LoadConfig(driverSecretKey(owner, key)) } + reg.VehicleWakeReservation = func(makeName, serial string) (bool, time.Duration, error) { + return st.ReserveVehicleWake(state.ResolveDeviceID(makeName, serial, "", ""), time.Now()) + } return reg } diff --git a/go/cmd/ftw/main.go b/go/cmd/ftw/main.go index 4ae5b19fb..fd329b2e3 100644 --- a/go/cmd/ftw/main.go +++ b/go/cmd/ftw/main.go @@ -2119,6 +2119,22 @@ func main() { return pick.Driver, pick.ChargingState, true }) + lpController.SetVehicleRefreshTarget(func(lpID string) (string, error) { + lp, ok := lpMgr.State(lpID) + if !ok || !lp.PluggedIn { + return "", nil + } + connected := 0 + for _, other := range lpMgr.States() { + if other.PluggedIn { + connected++ + } + } + cfgMu.RLock() + defer cfgMu.RUnlock() + return configuredVehicleRefreshDriver(cfg.Drivers, connected) + }) + lpController.SetVehicleChargeState(func(lpID string) (loadpoint.VehicleChargeState, bool) { pick := telemetry.PickVehicleForCompletion(tel, time.Now()) if pick.Driver == "" || pick.Stale || !lpMgr.VehicleObservationApplies(lpID, pick.UpdatedAt) { @@ -3154,18 +3170,10 @@ func main() { // plan once with measured-truth instead of the pluginSoC // estimate the startup replan used. if mpcSvc != nil && !vehicleReplanFired { - for _, vr := range tel.ReadingsByType(telemetry.DerVehicle) { - if vr.SoC == nil { - continue - } - if h := tel.DriverHealth(vr.Driver); h == nil || !h.IsOnline() { - continue - } + if pick := telemetry.PickBestVehicle(tel, time.Now()); pick.Driver != "" { vehicleReplanFired = true go mpcSvc.Replan(ctx) - slog.Info("first vehicle SoC seen → MPC replan triggered", - "driver", vr.Driver, "soc", *vr.SoC) - break + slog.Info("fresh vehicle SoC received; MPC replan requested", "driver", pick.Driver, "soc", pick.SoC) } } diff --git a/go/cmd/ftw/vehicle_refresh.go b/go/cmd/ftw/vehicle_refresh.go new file mode 100644 index 000000000..9a746c08c --- /dev/null +++ b/go/cmd/ftw/vehicle_refresh.go @@ -0,0 +1,44 @@ +package main + +import ( + "fmt" + "slices" + + "github.com/srcfl/ftw/go/internal/config" + "github.com/srcfl/ftw/go/internal/drivers" +) + +// Until a loadpoint has an explicit vehicle binding, one configured vehicle +// and one connected charger are the only unambiguous telemetry recipient. +// Freshness remains a requirement for planning and charge_start, not for wake. +func configuredVehicleRefreshDriver(configured []config.Driver, connected int) (string, error) { + if connected != 1 { + return "", fmt.Errorf("vehicle refresh requires one connected loadpoint") + } + var candidate *config.Driver + var entry drivers.CatalogEntry + for i := range configured { + d := &configured[i] + if d.Disabled { + continue + } + e, err := drivers.ParseCatalogFile(d.Lua) + if err != nil { + return "", fmt.Errorf("vehicle refresh: cannot read driver metadata") + } + if !slices.Contains(e.Capabilities, "vehicle") { + continue + } + if candidate != nil { + return "", fmt.Errorf("vehicle refresh is ambiguous with multiple configured vehicles") + } + candidate, entry = d, e + } + if candidate == nil { + return "", nil + } + if candidate.ObserveOnly || entry.ReadOnly || !entry.TelemetryWake { + return "", nil + } + return candidate.Name, nil +} diff --git a/go/cmd/ftw/vehicle_refresh_test.go b/go/cmd/ftw/vehicle_refresh_test.go new file mode 100644 index 000000000..98e44b002 --- /dev/null +++ b/go/cmd/ftw/vehicle_refresh_test.go @@ -0,0 +1,122 @@ +package main + +import ( + "context" + "github.com/srcfl/ftw/go/internal/config" + "github.com/srcfl/ftw/go/internal/state" + "github.com/srcfl/ftw/go/internal/telemetry" + "os" + "path/filepath" + "testing" + "time" +) + +func TestVehicleRefreshSelectsOnlyUnambiguousAuthorizedConfiguration(t *testing.T) { + dir := t.TempDir() + vehicle := func(name string, readOnly, wake bool) config.Driver { + path := filepath.Join(dir, name+".lua") + ro, wk := "false", "false" + if readOnly { + ro = "true" + } + if wake { + wk = "true" + } + source := `DRIVER = { + id = "test", + capabilities = { "vehicle" }, + read_only = ` + ro + `, + telemetry_wake = ` + wk + `, +}` + if err := os.WriteFile(path, []byte(source), 0600); err != nil { + t.Fatal(err) + } + return config.Driver{Name: name, Lua: path} + } + tesla := vehicle("tesla", false, true) + other := vehicle("other", true, false) + cases := []struct { + name string + drivers []config.Driver + connected int + want string + wantErr bool + }{ + {"missing or stale soc", []config.Driver{tesla}, 1, "tesla", false}, + {"two vehicles", []config.Driver{tesla, other}, 1, "", true}, + {"two chargers", []config.Driver{tesla}, 2, "", true}, + {"read only", []config.Driver{other}, 1, "", false}, + {"no wake declaration", []config.Driver{vehicle("legacy", false, false)}, 1, "", false}, + {"disabled other", []config.Driver{tesla, {Name: other.Name, Lua: other.Lua, Disabled: true}}, 1, "tesla", false}, + {"observe only", []config.Driver{{Name: tesla.Name, Lua: tesla.Lua, ObserveOnly: true}}, 1, "", false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + got, err := configuredVehicleRefreshDriver(tc.drivers, tc.connected) + if got != tc.want || (err != nil) != tc.wantErr { + t.Fatalf("got=%q err=%v", got, err) + } + }) + } +} + +func TestVehicleRefreshPollsPromptlyAndKeepsBudgetAcrossRename(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "vehicle.lua") + source := `DRIVER = { + id = "vehicle-test", + capabilities = { "vehicle" }, + telemetry_wake = true, +} +function driver_init(config) + host.set_make("Tesla") + host.set_sn("VIN-A") + host.set_poll_interval(3600000) +end +function driver_poll() + host.emit("vehicle", {soc=100, soc_fresh=true, soc_observed_at_ms=host.unix_ms(), charging_state="Complete"}) + host.set_poll_interval(3600000) +end +function driver_command(action) + assert(action == "wake_up") + local allowed = host.reserve_vehicle_wake() + host.emit_metric("reserved", allowed and 1 or 0) + host.set_poll_interval(10) + return true +end +function driver_default_mode() end` + if err := os.WriteFile(path, []byte(source), 0600); err != nil { + t.Fatal(err) + } + st, err := state.Open(filepath.Join(dir, "state.db")) + if err != nil { + t.Fatal(err) + } + defer st.Close() + tel := telemetry.NewStore() + reg := newDriverRegistry(tel, st) + defer reg.ShutdownAll() + for i, name := range []string{"original", "renamed"} { + if err := reg.Add(context.Background(), config.Driver{Name: name, Lua: path}); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + err := reg.Send(ctx, name, []byte(`{"action":"wake_up"}`)) + cancel() + if err != nil { + t.Fatal(err) + } + expected := float64(1 - i) + if value, _, ok := tel.LatestMetric(name, "reserved"); !ok || value != expected { + t.Fatalf("reservation after rename: %v", value) + } + deadline := time.Now().Add(time.Second) + for tel.Get(name, telemetry.DerVehicle) == nil && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + if reading := tel.Get(name, telemetry.DerVehicle); reading == nil || reading.SoC == nil || *reading.SoC != 1 { + t.Fatal("wake interval did not reach the active poll timer") + } + reg.Remove(name) + } +} diff --git a/go/internal/drivers/catalog.go b/go/internal/drivers/catalog.go index c1277767f..0a27882bb 100644 --- a/go/internal/drivers/catalog.go +++ b/go/internal/drivers/catalog.go @@ -94,6 +94,8 @@ type CatalogEntry struct { // write path, rather than FTW matching on a filename or vendor name. // Read-only remains the default for every driver in the catalog. WriteCapabilities []string `json:"write_capabilities,omitempty"` + // TelemetryWake declares a wake-only command, separate from charging. + TelemetryWake bool `json:"telemetry_wake,omitempty"` // Replaces names catalog driver ids this driver takes over, such as // esphome-dsmr folded into esphome_dsmr. A device running one of them is // moved to this driver when the release ships it. @@ -227,6 +229,7 @@ func parseCatalogEntry(path string) (CatalogEntry, error) { e.TestedModels = pickList(block, "tested_models") e.ConfigSecrets = pickList(block, "config_secrets") e.WriteCapabilities = pickList(block, "write_capabilities") + e.TelemetryWake = pickBool(block, "telemetry_wake") e.Replaces = pickList(block, "replaces") e.AuthPostPath = pickString(block, "auth_post_path") e.AuthPostPaths = pickList(block, "auth_post_paths") diff --git a/go/internal/drivers/heishamon_test.go b/go/internal/drivers/heishamon_test.go index e30a27578..87affde80 100644 --- a/go/internal/drivers/heishamon_test.go +++ b/go/internal/drivers/heishamon_test.go @@ -37,6 +37,7 @@ func TestHeishamonEmitsMetricsWithoutFakeBattery(t *testing.T) { } mqtt.Push("panasonic_heat_pump/main/Outside_Temp", "-4.5") + mqtt.Push("panasonic_heat_pump/main/Heat_Power_Consumption", "1250") mqtt.Push("panasonic_heat_pump/main/Main_Inlet_Temp", "31.2") mqtt.Push("panasonic_heat_pump/main/Main_Outlet_Temp", "35.7") mqtt.Push("panasonic_heat_pump/main/Main_Target_Temp", "36") @@ -46,7 +47,8 @@ func TestHeishamonEmitsMetricsWithoutFakeBattery(t *testing.T) { } wants := map[string]float64{ - "hp_outside_temp_c": -4.5, + "hp_outdoor_temp_c": -4.5, + "hp_power_w": 1250, "hp_inlet_temp_c": 31.2, "hp_outlet_temp_c": 35.7, "hp_target_temp_c": 36, diff --git a/go/internal/drivers/host.go b/go/internal/drivers/host.go index 778bd68bb..0af09602d 100644 --- a/go/internal/drivers/host.go +++ b/go/internal/drivers/host.go @@ -179,6 +179,10 @@ type HostEnv struct { // to 1 MiB and checks signed secret-key grants for managed drivers. PersistSecret func(key, value string) error + // ReserveVehicleWake checkpoints a telemetry wake before the driver sends it. + // The host binds this to hardware identity, independent of the driver name. + ReserveVehicleWake func(makeName, serial string) (bool, time.Duration, error) + // Poll-scoped Modbus evidence prevents a Modbus driver from turning failed // reads into fresh zero-valued telemetry. The Lua runtime holds emissions // until driver_poll ends, then commits them only when every read succeeded. diff --git a/go/internal/drivers/lua.go b/go/internal/drivers/lua.go index c39769291..77316cb30 100644 --- a/go/internal/drivers/lua.go +++ b/go/internal/drivers/lua.go @@ -643,6 +643,28 @@ func registerHost(L *lua.LState, env *HostEnv) { return 1 })) + // Wall-clock time is separate from millis/now_ms, which measure uptime. + host.RawSetString("unix_ms", L.NewFunction(func(L *lua.LState) int { + L.Push(lua.LNumber(time.Now().UnixMilli())) + return 1 + })) + host.RawSetString("reserve_vehicle_wake", L.NewFunction(func(L *lua.LState) int { + allowed, retry, err := false, 30*time.Minute, ErrNoCapability + if !env.ProbeReadOnly && !driverDeclaresReadOnly(L) && + (env.RuntimePolicy == nil || !env.RuntimePolicy.ReadOnly) && env.ReserveVehicleWake != nil { + makeName, serial := env.Identity() + allowed, retry, err = env.ReserveVehicleWake(makeName, serial) + } + L.Push(lua.LBool(allowed && err == nil)) + L.Push(lua.LNumber(retry.Milliseconds())) + if err != nil { + L.Push(lua.LString("vehicle wake reservation failed")) + } else { + L.Push(lua.LNil) + } + return 3 + })) + // host.sleep(ms) — block the driver goroutine for ms milliseconds. // Used for vendor-required inter-write pacing (Solis 100ms, Deye 50ms); // safe because each driver has its own goroutine and VM lock. diff --git a/go/internal/drivers/lua_vehicle_wake_test.go b/go/internal/drivers/lua_vehicle_wake_test.go new file mode 100644 index 000000000..572046883 --- /dev/null +++ b/go/internal/drivers/lua_vehicle_wake_test.go @@ -0,0 +1,62 @@ +package drivers + +import ( + "context" + "os" + "path/filepath" + "testing" + "time" + + "github.com/srcfl/ftw/go/internal/telemetry" +) + +func TestLuaVehicleWakeHostBindingAndReadOnlyBoundary(t *testing.T) { + for _, readonly := range []bool{false, true} { + t.Run(map[bool]string{false: "vehicle", true: "read_only"}[readonly], func(t *testing.T) { + declaration := "false" + if readonly { + declaration = "true" + } + source := `DRIVER = { read_only = ` + declaration + ` } +function driver_init(config) host.set_make("Tesla"); host.set_sn("VIN-A") end +function driver_poll() + assert(host.unix_ms() > 1700000000000) + local ok, retry, err = host.reserve_vehicle_wake() + assert(retry > 0) + if DRIVER.read_only then assert(not ok and err) else assert(ok and not err) end +end +function driver_default_mode() end` + path := filepath.Join(t.TempDir(), "vehicle.lua") + if err := os.WriteFile(path, []byte(source), 0600); err != nil { + t.Fatal(err) + } + env := NewHostEnv("renamed-driver", telemetry.NewStore()) + calls := 0 + env.ReserveVehicleWake = func(makeName, serial string) (bool, time.Duration, error) { + if makeName != "Tesla" || serial != "VIN-A" { + t.Fatal("reservation lost hardware identity") + } + calls++ + return true, 90 * time.Second, nil + } + d, err := NewLuaDriver(path, env) + if err != nil { + t.Fatal(err) + } + defer d.Cleanup() + if err := d.Init(context.Background(), nil); err != nil { + t.Fatal(err) + } + if _, err := d.Poll(context.Background()); err != nil { + t.Fatal(err) + } + want := 1 + if readonly { + want = 0 + } + if calls != want { + t.Fatalf("reservation calls=%d want=%d", calls, want) + } + }) + } +} diff --git a/go/internal/drivers/registry.go b/go/internal/drivers/registry.go index c1e3236e7..504e2a5e1 100644 --- a/go/internal/drivers/registry.go +++ b/go/internal/drivers/registry.go @@ -101,7 +101,8 @@ type Registry struct { // driver (the counterpart to SecretPersister). Applied over the // config.yaml value at driver_init so a rotated token survives a // restart. Returns ("", false) when no override exists. - SecretOverride func(owner, key string) (string, bool) + SecretOverride func(owner, key string) (string, bool) + VehicleWakeReservation func(makeName, serial string) (bool, time.Duration, error) // RuntimePolicyResolver returns the verified signed policy of a managed // read-only artifact. Nil means bundled, local and control-capable // signed drivers, which run without one. @@ -540,6 +541,12 @@ func (r *Registry) add(ctx context.Context, cfg config.Driver, startupDefault bo // Wire secret write-back (rotated OAuth tokens). The host must install // SecretPersister and SecretOverride before Add: init may persist a // secret, and the poll loop starts before Add returns. + env.ReserveVehicleWake = func(makeName, serial string) (bool, time.Duration, error) { + if cfg.ObserveOnly || r.VehicleWakeReservation == nil { + return false, 30 * time.Minute, ErrNoCapability + } + return r.VehicleWakeReservation(makeName, serial) + } secretOwner := cfg.SecretOwner() env.PersistSecret = func(key, value string) error { if r.SecretPersister == nil { @@ -1019,6 +1026,11 @@ func (r *Registry) runLoop(rd *runningDriver) { } else { rd.markCommandApplied() } + if action == "wake_up" || action == "ev_wake" { + // A telemetry refresh may request an early read even when its + // wake was throttled. Apply that interval to the active timer. + timer.Reset(rd.env.PollInterval()) + } finishCommand() if err == nil && cyclePause { rd.evPausePending = true diff --git a/go/internal/loadpoint/controller.go b/go/internal/loadpoint/controller.go index 517365cc8..d4d5f49a7 100644 --- a/go/internal/loadpoint/controller.go +++ b/go/internal/loadpoint/controller.go @@ -108,8 +108,9 @@ type Controller struct { // vehicleStatus is the bound vehicle driver and its charging state. // nil disables auto-wake. - vehicleStatus func(loadpointID string) (driver, chargingState string, ok bool) - vehicleChargeState func(loadpointID string) (VehicleChargeState, bool) + vehicleStatus func(loadpointID string) (driver, chargingState string, ok bool) + vehicleRefreshTarget func(loadpointID string) (string, error) + vehicleChargeState func(loadpointID string) (VehicleChargeState, bool) // peakRemainingSurplusW is the best PV-minus-load surplus left today. // Below the 3Φ minimum, surplus_only locks 1Φ for the day. nil keeps 3Φ. @@ -555,6 +556,14 @@ func (c *Controller) SetVehicleStatus(f func(loadpointID string) (driver, chargi c.vehicleStatus = f } +// SetVehicleRefreshTarget selects an authorized telemetry recipient from +// configuration. It must not require a fresh SoC or permit charging. +func (c *Controller) SetVehicleRefreshTarget(f func(string) (string, error)) { + if c != nil { + c.vehicleRefreshTarget = f + } +} + // SetPeakRemainingSurplusW wires the forecast-based "best surplus // we'll see for the rest of the day" reader used by surplus_only's // 1Φ-lock decision. Typical implementation in main.go iterates the @@ -879,48 +888,23 @@ func (c *Controller) AnyLoadpointSurplusActive() bool { return false } -// RefreshVehicle sends a one-off wake command to the vehicle driver -// bound to the given loadpoint, bypassing the auto-wake cooldown. -// Used by the API when the operator edits the schedule — wakes -// Tesla / BMW / whichever vehicle driver is bound so the next poll -// surfaces any vehicle-side limit / SoC / connection changes -// immediately rather than waiting up to the next natural wake -// window. -// -// Sends the generic `wake_up` action (cross-driver protocol — any -// vehicle driver implements it against its own back-end). Distinct -// from `charge_start` (still used by the auto-wake loop when it's -// trying to convince a detached car to actually start drawing -// current) — `wake_up` is purely a telemetry refresh, no charge -// side effects. Returns nil if no vehicle driver is bound or the -// controller isn't fully wired (no-op). Errors from the send hop -// are returned for the caller to surface to the operator. +// RefreshVehicle requests telemetry only. The driver retains its durable +// wake limit; refreshing a goal never resets the charge_start cooldown. func (c *Controller) RefreshVehicle(ctx context.Context, lpID string) error { - if c == nil || c.vehicleStatus == nil || c.send == nil { + if c == nil || c.vehicleRefreshTarget == nil || c.send == nil { return nil } - driver, _, ok := c.vehicleStatus(lpID) - if !ok || driver == "" { + driver, err := c.vehicleRefreshTarget(lpID) + if err != nil { + return err + } + if driver == "" { return nil } payload, err := json.Marshal(map[string]any{"action": "wake_up"}) if err != nil { return err } - // Reset the auto-wake throttle so a manual refresh doesn't leave - // the LP in a long backoff afterwards — the operator just told us - // they want a fresh read, no reason to apply 90s cooldown to the - // next legitimate auto-wake. - c.wakeMu.Lock() - if c.wakeLast == nil { - c.wakeLast = map[string]time.Time{} - c.wakeKickUntil = map[string]time.Time{} - c.wakeAttempts = map[string]int{} - } - delete(c.wakeAttempts, lpID) - c.wakeLast[lpID] = time.Now() - c.wakeMu.Unlock() - slog.Info("loadpoint manual wake (schedule edit)", "lp", lpID, "vehicle_driver", driver) return c.sendVehicle(ctx, driver, payload) } diff --git a/go/internal/loadpoint/controller_dispatch_outcome_test.go b/go/internal/loadpoint/controller_dispatch_outcome_test.go index 493734d4d..bd97a78ba 100644 --- a/go/internal/loadpoint/controller_dispatch_outcome_test.go +++ b/go/internal/loadpoint/controller_dispatch_outcome_test.go @@ -411,7 +411,7 @@ func TestScheduleRefreshKeepsVehicleCallerDeadline(t *testing.T) { } }) c := NewController(NewManager(), nil, nil, send) - c.SetVehicleStatus(func(string) (string, string, bool) { return "tesla", "Stopped", true }) + c.SetVehicleRefreshTarget(func(string) (string, error) { return "tesla", nil }) c.SetCommandTimeout(10 * time.Millisecond) ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() diff --git a/go/internal/loadpoint/vehicle_refresh_test.go b/go/internal/loadpoint/vehicle_refresh_test.go new file mode 100644 index 000000000..df869bab8 --- /dev/null +++ b/go/internal/loadpoint/vehicle_refresh_test.go @@ -0,0 +1,41 @@ +package loadpoint + +import ( + "context" + "encoding/json" + "errors" + "testing" + "time" +) + +func TestRefreshVehicleIndependentOfFreshSoCAndChargeStart(t *testing.T) { + commands := 0 + c := NewController(NewManager(), nil, nil, SenderFunc(func(_ context.Context, driver string, payload []byte) error { + var action struct { + Action string `json:"action"` + } + if err := json.Unmarshal(payload, &action); err != nil { + t.Fatal(err) + } + if driver != "tesla" || action.Action != "wake_up" { + t.Fatalf("driver=%q payload=%s", driver, payload) + } + commands++ + return nil + })) + c.SetVehicleStatus(func(string) (string, string, bool) { return "", "", false }) + c.SetVehicleRefreshTarget(func(string) (string, error) { return "tesla", nil }) + previous := time.Now().Add(-time.Hour) + c.wakeLast = map[string]time.Time{"garage": previous} + c.wakeAttempts = map[string]int{"garage": 5} + if err := c.RefreshVehicle(context.Background(), "garage"); err != nil { + t.Fatal(err) + } + if commands != 1 || c.wakeAttempts["garage"] != 5 || !c.wakeLast["garage"].Equal(previous) { + t.Fatal("refresh changed charge-start state") + } + c.SetVehicleRefreshTarget(func(string) (string, error) { return "", errors.New("ambiguous vehicle") }) + if err := c.RefreshVehicle(context.Background(), "garage"); err == nil || commands != 1 { + t.Fatal("ambiguous vehicle refreshed") + } +} diff --git a/go/internal/state/vehicle_wake.go b/go/internal/state/vehicle_wake.go new file mode 100644 index 000000000..7ae451b9f --- /dev/null +++ b/go/internal/state/vehicle_wake.go @@ -0,0 +1,77 @@ +package state + +import ( + "database/sql" + "encoding/json" + "errors" + "fmt" + "time" +) + +const vehicleWakeGap = 90 * time.Second +const vehicleWakeWindow = 30 * time.Minute +const vehicleWakeBudget = 3 + +type vehicleWakeCheckpoint struct { + WindowStartMs int64 `json:"window_start_ms"` + LastAttemptMs int64 `json:"last_attempt_ms"` + Attempts int `json:"attempts"` +} + +// ReserveVehicleWake commits the attempt before any network request. A crash, +// driver rename or failed HTTP call cannot restore the wake budget. This uses +// existing config rows and does not change the state schema. +func (s *Store) ReserveVehicleWake(deviceID string, now time.Time) (bool, time.Duration, error) { + if deviceID == "" || now.IsZero() { + return false, vehicleWakeWindow, fmt.Errorf("vehicle wake requires hardware identity and time") + } + allowed := false + retry := vehicleWakeWindow + err := s.durableConfigWrite(func(tx *sql.Tx) error { + key := "vehicle_wake:" + deviceID + var raw string + c := vehicleWakeCheckpoint{} + err := tx.QueryRow(`SELECT value FROM config WHERE key = ?`, key).Scan(&raw) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return err + } + if err == nil { + if err := json.Unmarshal([]byte(raw), &c); err != nil { + return err + } + if c.Attempts < 1 || c.Attempts > vehicleWakeBudget || c.WindowStartMs <= 0 || c.LastAttemptMs < c.WindowStartMs { + return fmt.Errorf("invalid vehicle wake checkpoint") + } + } + nowMs := now.UnixMilli() + windowEnd := time.UnixMilli(c.WindowStartMs).Add(vehicleWakeWindow) + // Clock rollback retains the budget instead of treating the record as absent. + if c.Attempts > 0 && now.Before(windowEnd) { + next := time.UnixMilli(c.LastAttemptMs).Add(vehicleWakeGap) + if c.Attempts >= vehicleWakeBudget { + next = windowEnd + } + if now.Before(next) { + retry = next.Sub(now) + return nil + } + } else { + c = vehicleWakeCheckpoint{WindowStartMs: nowMs} + } + c.Attempts++ + c.LastAttemptMs = nowMs + encoded, err := json.Marshal(c) + if err != nil { + return err + } + if err := saveConfigValues(tx, map[string]string{key: string(encoded)}); err != nil { + return err + } + allowed, retry = true, vehicleWakeGap + return nil + }) + if err != nil { + return false, vehicleWakeWindow, err + } + return allowed, retry, nil +} diff --git a/go/internal/state/vehicle_wake_test.go b/go/internal/state/vehicle_wake_test.go new file mode 100644 index 000000000..257c91e83 --- /dev/null +++ b/go/internal/state/vehicle_wake_test.go @@ -0,0 +1,90 @@ +package state + +import ( + "path/filepath" + "sync" + "testing" + "time" +) + +func TestVehicleWakeBudgetSurvivesRestartAndIsPerHardware(t *testing.T) { + path := filepath.Join(t.TempDir(), "state.db") + s, err := Open(path) + if err != nil { + t.Fatal(err) + } + now := time.Unix(1760000000, 0) + for attempt := 0; attempt < 3; attempt++ { + allowed, _, err := s.ReserveVehicleWake("tesla:VIN-A", now) + if err != nil || !allowed { + t.Fatalf("attempt %d: %v %v", attempt, allowed, err) + } + if err := s.Close(); err != nil { + t.Fatal(err) + } + s, err = Open(path) + if err != nil { + t.Fatal(err) + } + allowed, retry, err := s.ReserveVehicleWake("tesla:VIN-A", now.Add(time.Second)) + if err != nil || allowed || retry <= 0 { + t.Fatalf("restart renewed budget: %v %v %v", allowed, retry, err) + } + now = now.Add(vehicleWakeGap) + } + defer s.Close() + if allowed, _, err := s.ReserveVehicleWake("tesla:VIN-A", now); allowed || err != nil { + t.Fatalf("fourth wake: %v %v", allowed, err) + } + if allowed, _, err := s.ReserveVehicleWake("tesla:VIN-B", now); !allowed || err != nil { + t.Fatalf("other VIN: %v %v", allowed, err) + } + if allowed, _, err := s.ReserveVehicleWake("tesla:VIN-A", time.Unix(1760000000, 0).Add(vehicleWakeWindow)); !allowed || err != nil { + t.Fatalf("new window: %v %v", allowed, err) + } +} + +func TestVehicleWakeFailsClosed(t *testing.T) { + s := freshStore(t) + now := time.Now() + if allowed, _, err := s.ReserveVehicleWake("", now); allowed || err == nil { + t.Fatal("missing identity admitted") + } + if err := s.SaveConfig("vehicle_wake:tesla:bad", "broken"); err != nil { + t.Fatal(err) + } + if allowed, _, err := s.ReserveVehicleWake("tesla:bad", now); allowed || err == nil { + t.Fatal("broken checkpoint admitted") + } + if _, err := s.db.Exec(`CREATE TRIGGER fail_wake BEFORE INSERT ON config WHEN NEW.key LIKE 'vehicle_wake:%' BEGIN SELECT RAISE(ABORT, 'disk failure'); END`); err != nil { + t.Fatal(err) + } + if allowed, _, err := s.ReserveVehicleWake("tesla:unwritten", now); allowed || err == nil { + t.Fatal("uncommitted attempt admitted") + } +} + +func TestVehicleWakeConcurrentReservationsAndClockRollback(t *testing.T) { + s := freshStore(t) + now := time.Now() + var wg sync.WaitGroup + var mu sync.Mutex + accepted := 0 + for i := 0; i < 8; i++ { + wg.Go(func() { + allowed, _, _ := s.ReserveVehicleWake("tesla:one", now) + if allowed { + mu.Lock() + accepted++ + mu.Unlock() + } + }) + } + wg.Wait() + if accepted != 1 { + t.Fatalf("concurrent wakes accepted=%d", accepted) + } + if allowed, retry, err := s.ReserveVehicleWake("tesla:one", now.Add(-time.Minute)); allowed || retry < time.Minute || err != nil { + t.Fatalf("clock rollback: %v %v %v", allowed, retry, err) + } +} diff --git a/go/internal/telemetry/store.go b/go/internal/telemetry/store.go index 2f24abc0a..3e65f4eeb 100644 --- a/go/internal/telemetry/store.go +++ b/go/internal/telemetry/store.go @@ -90,7 +90,7 @@ type DerReading struct { RawW float64 SmoothedW float64 SoC *float64 // optional; 0..1 fraction for every DER including vehicles - SoCUpdatedAt time.Time // last fresh SoC receipt time; preserved across cached updates + SoCUpdatedAt time.Time // last SoC source observation (receipt time for legacy drivers) Data json.RawMessage UpdatedAt time.Time } @@ -426,6 +426,21 @@ func (s *Store) Update(driver string, t DerType, rawW float64, soc *float64, dat var socUpdatedAt time.Time if socFresh { socUpdatedAt = now + if t == DerVehicle { + var fields map[string]json.RawMessage + _ = json.Unmarshal(data, &fields) + if raw, present := fields["soc_observed_at_ms"]; present { + var ms int64 + if json.Unmarshal(raw, &ms) != nil || ms <= 0 || ms > now.UnixMilli() { + socFresh = false + } else { + socUpdatedAt = time.UnixMilli(ms) + if previous := s.readings[k]; previous != nil && !socUpdatedAt.After(previous.SoCUpdatedAt) { + socFresh = false + } + } + } + } } // Preserve last-known SoC when the new emit doesn't include one or marks a // non-nil value as a cached replay. @@ -435,6 +450,7 @@ func (s *Store) Update(driver string, t DerType, rawW float64, soc *float64, dat // persisted metrics can still distinguish it from a fresh number. // A battery SoC older than BatterySoCMaxAge is not carried: it is unknown. if !socFresh { + socUpdatedAt = time.Time{} prev, ok := s.readings[k] switch { case !ok || prev.SoC == nil: @@ -470,7 +486,7 @@ func (s *Store) Update(driver string, t DerType, rawW float64, soc *float64, dat ) if socFresh { s.pending = append(s.pending, - MetricSample{Driver: driver, Metric: t.String() + "_soc", TsMs: tsMs, Value: *soc}, + MetricSample{Driver: driver, Metric: t.String() + "_soc", TsMs: socUpdatedAt.UnixMilli(), Value: *soc}, ) } s.pendingMu.Unlock() diff --git a/go/internal/telemetry/vehicle_source_time_test.go b/go/internal/telemetry/vehicle_source_time_test.go new file mode 100644 index 000000000..371186ccc --- /dev/null +++ b/go/internal/telemetry/vehicle_source_time_test.go @@ -0,0 +1,54 @@ +package telemetry + +import ( + "encoding/json" + "fmt" + "testing" + "time" +) + +func TestVehicleSourceTimeControlsFreshnessAndFullCarRecovery(t *testing.T) { + s := NewStore() + s.DriverHealthMut("tesla").RecordSuccess() + now := time.Now() + emit := func(at time.Time, value float64) { + s.Update("tesla", DerVehicle, 0, &value, json.RawMessage(fmt.Sprintf(`{"soc":%g,"soc_fresh":true,"soc_observed_at_ms":%d,"charging_state":"Complete"}`, value*100, at.UnixMilli()))) + } + old := now.Add(-10 * time.Minute) + emit(old, 0.5) + if got := s.Get("tesla", DerVehicle); !got.SoCUpdatedAt.Equal(time.UnixMilli(old.UnixMilli())) { + t.Fatalf("source time replaced: %v", got.SoCUpdatedAt) + } + if pick := PickBestVehicle(s, now); pick.Driver != "" { + t.Fatal("old first response became planning truth") + } + emit(old, 0.5) + if pick := PickBestVehicle(s, now); pick.Driver != "" { + t.Fatal("cached response renewed observation") + } + s.FlushSamples() + fresh := now.Add(-time.Second) + emit(fresh, 1) + pick := PickBestVehicle(s, now) + if pick.Driver != "tesla" || pick.SoC != 1 { + t.Fatalf("full car did not replace inferred demand: %+v", pick) + } + if samples := s.FlushSamples(); len(samples) != 2 { + t.Fatalf("fresh source not stored: %+v", samples) + } + emit(fresh, 0.1) + if got := s.Get("tesla", DerVehicle); *got.SoC != 1 || !got.SoCUpdatedAt.Equal(pick.UpdatedAt) { + t.Fatal("identical source timestamp replaced full-car reading") + } +} + +func TestVehicleInvalidSourceTimeFailsClosed(t *testing.T) { + for _, stamp := range []string{`null`, `"yesterday"`, `0`, `-1`, fmt.Sprint(time.Now().Add(time.Hour).UnixMilli())} { + s := NewStore() + soc := 0.8 + s.Update("tesla", DerVehicle, 0, &soc, json.RawMessage(`{"soc_fresh":true,"soc_observed_at_ms":`+stamp+`}`)) + if got := s.Get("tesla", DerVehicle); got.SoC != nil || !got.SoCUpdatedAt.IsZero() { + t.Fatalf("invalid source %s admitted", stamp) + } + } +}