From e945e8072a7e767eae1dae7aa41ebb3b23ba8e9b Mon Sep 17 00:00:00 2001 From: Jack Pearce <16779171+jkpe@users.noreply.github.com> Date: Thu, 23 Jul 2026 22:50:31 +0100 Subject: [PATCH 1/2] Introduces `WorkModeNotifier` to track and report changes in inverter work mode. Adds `NOTIFY_DEBOUNCE_POLLS` and `CMS_API_URL` configuration variables to control notification stability and the target API endpoint. --- README.md | 4 +- docker-compose.yml | 2 + src/config.ts | 4 + src/index.ts | 14 ++- src/notify/workModeNotifier.test.ts | 157 ++++++++++++++++++++++++++++ src/notify/workModeNotifier.ts | 92 ++++++++++++++++ 6 files changed, 271 insertions(+), 2 deletions(-) create mode 100644 src/notify/workModeNotifier.test.ts create mode 100644 src/notify/workModeNotifier.ts diff --git a/README.md b/README.md index 38cedf2..0964c62 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ Local Docker stack for **Fox ESS** inverters (H1, H3, KH, and OEM variants). Polls live data over Modbus TCP and sends snapshots to Check My Solar through a private tunnel. -The bridge **auto-detects** your inverter model from holding register 30000 on startup (same approach as [foxess_modbus](https://github.com/nathanmarlor/foxess_modbus)). +The bridge **auto-detects** your inverter model. **Full guide:** [checkmy.solar/docs/using-the-app/modbus-bridge/](https://checkmy.solar/docs/using-the-app/modbus-bridge/) @@ -30,6 +30,8 @@ export MODBUS_HOST='192.168.1.100' # Modbus adapter IP export BRIDGE_HOSTNAME='bridge-....modbus.internal' # from the app export TUNNEL_TOKEN='eyJ...' # from the app export SITE_TIMEZONE='Europe/London' # IANA timezone for hour buckets +# export CMS_API_URL='https://checkmy.solar' # override for dev/staging +# export NOTIFY_DEBOUNCE_POLLS=2 # stable polls before work mode push # export MODBUS_CONNECTION=aux # default; use lan for direct inverter LAN # export INVERTER_PROFILE=h3Modern # optional override # export BRIDGE_VERBOSE_LOG=true # log each Modbus poll and HTTP request diff --git a/docker-compose.yml b/docker-compose.yml index 66cb84e..e8c7789 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -18,6 +18,8 @@ services: MODBUS_CONNECTION: ${MODBUS_CONNECTION:-aux} BRIDGE_HOSTNAME: ${BRIDGE_HOSTNAME:?BRIDGE_HOSTNAME is required} SITE_TIMEZONE: ${SITE_TIMEZONE:?SITE_TIMEZONE is required} + CMS_API_URL: ${CMS_API_URL:-https://checkmy.solar} + NOTIFY_DEBOUNCE_POLLS: ${NOTIFY_DEBOUNCE_POLLS:-2} BRIDGE_VERBOSE_LOG: ${BRIDGE_VERBOSE_LOG:-false} networks: cms_net: diff --git a/src/config.ts b/src/config.ts index 6f49d7b..3fcf4a5 100644 --- a/src/config.ts +++ b/src/config.ts @@ -15,6 +15,8 @@ export interface BridgeConfig { dataDir: string; siteTimezone: string; bridgeHostname?: string; + cmsApiUrl: string; + notifyDebouncePolls: number; /** When true, log each Modbus poll and each HTTP request. */ verboseLogging: boolean; /** Force a register profile instead of auto-detecting from the inverter model. */ @@ -102,6 +104,8 @@ export function loadConfig(): BridgeConfig { dataDir: process.env.BRIDGE_DATA_DIR?.trim() || '/data', siteTimezone: readTimezone('SITE_TIMEZONE'), bridgeHostname: readOptional('BRIDGE_HOSTNAME'), + cmsApiUrl: readOptional('CMS_API_URL') ?? 'https://checkmy.solar', + notifyDebouncePolls: readInt('NOTIFY_DEBOUNCE_POLLS', 2), verboseLogging: readBoolean('BRIDGE_VERBOSE_LOG', false), inverterProfile: readProfileId('INVERTER_PROFILE'), modbusConnection: readConnectionType('MODBUS_CONNECTION', 'aux'), diff --git a/src/index.ts b/src/index.ts index d49e243..74322e3 100644 --- a/src/index.ts +++ b/src/index.ts @@ -4,6 +4,7 @@ import { HourlyAggregator } from './aggregation/hourlyAggregator.js'; import { startBridgeHttpServer } from './http/server.js'; import { FoxModbusClient } from './modbus/client.js'; import { mapH1G2TodayTotalsSnapshotToFoxShape } from './modbus/h1g2TodayTotals.js'; +import { WorkModeNotifier } from './notify/workModeNotifier.js'; import { RealtimeStore } from './storage/sqlite.js'; import { formatStoredTelemetryLog } from './telemetryLog.js'; @@ -16,6 +17,7 @@ function sleep(ms: number): Promise { async function runPollCycle( modbus: FoxModbusClient, store: RealtimeStore, + workModeNotifier: WorkModeNotifier, aggregator: HourlyAggregator, verboseLogging: boolean ): Promise { @@ -26,6 +28,7 @@ async function runPollCycle( ]); store.upsert(telemetry, telemetry.sampledAt); + await workModeNotifier.handleSample(telemetry); const todayTotals = mapH1G2TodayTotalsSnapshotToFoxShape(todayTotalsSnapshot); if (todayTotals) { store.upsertTodayTotals(todayTotals, todayTotalsSnapshot.sampledAt); @@ -42,6 +45,15 @@ async function main(): Promise { const config = loadConfig(); const store = new RealtimeStore(config.dataDir, { verboseLogging: config.verboseLogging }); const aggregator = new HourlyAggregator(store, config.siteTimezone); + const latestSnapshot = store.getLatest(); + const workModeNotifier = new WorkModeNotifier( + { + apiUrl: config.cmsApiUrl, + bridgeToken: config.bridgeToken, + debouncePolls: config.notifyDebouncePolls, + }, + latestSnapshot?.telemetry.workMode + ); let detectedInverter: ReturnType = null; @@ -96,7 +108,7 @@ async function main(): Promise { backoffMs = config.pollIntervalMs; while (true) { - await runPollCycle(modbus, store, aggregator, config.verboseLogging); + await runPollCycle(modbus, store, workModeNotifier, aggregator, config.verboseLogging); await sleep(config.pollIntervalMs); } } catch (error) { diff --git a/src/notify/workModeNotifier.test.ts b/src/notify/workModeNotifier.test.ts new file mode 100644 index 0000000..eb45bee --- /dev/null +++ b/src/notify/workModeNotifier.test.ts @@ -0,0 +1,157 @@ +import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest'; +import type { ModbusRealtimeTelemetry } from '@checkmysolar/modbus-telemetry'; +import { WorkModeNotifier, postWorkModeEvent } from './workModeNotifier.js'; + +function sampleTelemetry(workMode: number): ModbusRealtimeTelemetry { + return { + loadsPower: 2.88, + pvPower: 3.5, + pv1Power: 2, + pv2Power: 1.5, + pvStringCount: 2, + pvStringPowers: { pv1Power: 2, pv2Power: 1.5 }, + feedinPower: 1.2, + gridConsumptionPower: 0, + batChargePower: 0, + batDischargePower: 0.5, + SoC: 85, + ResidualEnergy: 10.5, + batVoltage: 51.2, + batCurrent: -1, + batTemperature: 28, + gridVoltage: 230, + gridCurrent: 5, + gridFrequency: 50, + meterPower2: 0.1, + ambientTemperature: 25, + deviceTemperature: 45, + runningState: 163, + isOffGrid: false, + epsPower: 0, + epsPowerR: 0, + epsVoltR: 240, + epsCurrentR: 0, + workMode, + sampledAt: '2026-07-09T11:59:30.000Z', + }; +} + +describe('WorkModeNotifier', () => { + beforeEach(() => { + vi.stubGlobal('fetch', vi.fn()); + }); + + afterEach(() => { + vi.unstubAllGlobals(); + }); + + it('requires stable polls before posting a work mode change', async () => { + const fetchMock = vi.mocked(fetch); + fetchMock.mockResolvedValue(new Response(null, { status: 204 })); + + const notifier = new WorkModeNotifier( + { + apiUrl: 'https://checkmy.solar', + bridgeToken: 'cms_bridge_test', + debouncePolls: 2, + }, + 0 + ); + + await notifier.handleSample(sampleTelemetry(3)); + expect(fetchMock).not.toHaveBeenCalled(); + + await notifier.handleSample(sampleTelemetry(3)); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(fetchMock).toHaveBeenCalledWith( + 'https://checkmy.solar/api/bridge/events/work-mode', + expect.objectContaining({ + method: 'POST', + headers: expect.objectContaining({ + Authorization: 'Bearer cms_bridge_test', + }), + body: JSON.stringify({ + workMode: 3, + sampledAt: '2026-07-09T11:59:30.000Z', + previousWorkMode: 0, + soc: 85, + }), + }) + ); + }); + + it('resets debounce when the pending mode changes before stability', async () => { + const fetchMock = vi.mocked(fetch); + fetchMock.mockResolvedValue(new Response(null, { status: 204 })); + + const notifier = new WorkModeNotifier( + { + apiUrl: 'https://checkmy.solar', + bridgeToken: 'cms_bridge_test', + debouncePolls: 2, + }, + 0 + ); + + await notifier.handleSample(sampleTelemetry(3)); + await notifier.handleSample(sampleTelemetry(4)); + expect(fetchMock).not.toHaveBeenCalled(); + + await notifier.handleSample(sampleTelemetry(4)); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(fetchMock.mock.calls[0]?.[1]).toEqual( + expect.objectContaining({ + body: JSON.stringify({ + workMode: 4, + sampledAt: '2026-07-09T11:59:30.000Z', + previousWorkMode: 0, + soc: 85, + }), + }) + ); + }); + + it('does not fail the poll loop when the API rejects the event', async () => { + const fetchMock = vi.mocked(fetch); + fetchMock.mockResolvedValue(new Response('bad request', { status: 400 })); + + const notifier = new WorkModeNotifier( + { + apiUrl: 'https://checkmy.solar', + bridgeToken: 'cms_bridge_test', + debouncePolls: 1, + }, + 0 + ); + + await expect(notifier.handleSample(sampleTelemetry(1))).resolves.toBeUndefined(); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); +}); + +describe('postWorkModeEvent', () => { + beforeEach(() => { + vi.stubGlobal('fetch', vi.fn()); + }); + + afterEach(() => { + vi.unstubAllGlobals(); + }); + + it('accepts 204 responses', async () => { + const fetchMock = vi.mocked(fetch); + fetchMock.mockResolvedValue(new Response(null, { status: 204 })); + + await expect( + postWorkModeEvent( + { apiUrl: 'https://checkmy.solar/', bridgeToken: 'cms_bridge_test' }, + { workMode: 2, sampledAt: '2026-07-09T11:59:30.000Z' } + ) + ).resolves.toBeUndefined(); + + expect(fetchMock).toHaveBeenCalledWith( + 'https://checkmy.solar/api/bridge/events/work-mode', + expect.any(Object) + ); + }); +}); diff --git a/src/notify/workModeNotifier.ts b/src/notify/workModeNotifier.ts new file mode 100644 index 0000000..7a51d86 --- /dev/null +++ b/src/notify/workModeNotifier.ts @@ -0,0 +1,92 @@ +import type { ModbusRealtimeTelemetry } from '@checkmysolar/modbus-telemetry'; +import { formatError } from '../errors.js'; + +export interface WorkModeNotifierOptions { + apiUrl: string; + bridgeToken: string; + debouncePolls: number; +} + +export interface WorkModeEventPayload { + workMode: number; + sampledAt: string; + previousWorkMode?: number; + soc?: number; +} + +export class WorkModeNotifier { + private pendingWorkMode: number | undefined; + private pendingCount = 0; + private lastEmittedWorkMode: number | undefined; + + constructor( + private readonly options: WorkModeNotifierOptions, + initialWorkMode?: number + ) { + this.lastEmittedWorkMode = initialWorkMode; + } + + async handleSample(telemetry: ModbusRealtimeTelemetry): Promise { + const workMode = telemetry.workMode; + if (workMode === undefined || workMode === null) { + return; + } + + if (workMode === this.lastEmittedWorkMode) { + this.pendingWorkMode = undefined; + this.pendingCount = 0; + return; + } + + if (workMode === this.pendingWorkMode) { + this.pendingCount += 1; + } else { + this.pendingWorkMode = workMode; + this.pendingCount = 1; + } + + if (this.pendingCount < this.options.debouncePolls) { + return; + } + + const previousWorkMode = this.lastEmittedWorkMode; + + try { + await postWorkModeEvent(this.options, { + workMode, + sampledAt: telemetry.sampledAt, + previousWorkMode, + soc: telemetry.SoC, + }); + this.lastEmittedWorkMode = workMode; + this.pendingWorkMode = undefined; + this.pendingCount = 0; + } catch (error) { + console.error(`Work mode notification failed: ${formatError(error)}`); + } + } +} + +export async function postWorkModeEvent( + options: Pick, + payload: WorkModeEventPayload +): Promise { + const url = `${options.apiUrl.replace(/\/$/, '')}/api/bridge/events/work-mode`; + const response = await fetch(url, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + Authorization: `Bearer ${options.bridgeToken}`, + }, + body: JSON.stringify(payload), + }); + + if (response.status === 401) { + throw new Error('Bridge token rejected'); + } + + if (!response.ok && response.status !== 204) { + const text = await response.text().catch(() => ''); + throw new Error(`Work mode event rejected: ${response.status}${text ? ` ${text}` : ''}`); + } +} From bdb63e13ada5a0e410cb7456cd5df9cfd926e1d8 Mon Sep 17 00:00:00 2001 From: Jack Pearce <16779171+jkpe@users.noreply.github.com> Date: Thu, 23 Jul 2026 23:03:50 +0100 Subject: [PATCH 2/2] Fix work mode notify reliability after restart and during polling. Persist the last successfully notified work mode, post CMS events in the background with a timeout, and avoid blocking the Modbus poll loop. Co-authored-by: Cursor --- src/config.ts | 2 + src/index.ts | 6 +- src/notify/workModeNotifier.test.ts | 155 +++++++++++++++++++++++++--- src/notify/workModeNotifier.ts | 36 +++++-- src/storage/sqlite.ts | 32 ++++++ 5 files changed, 209 insertions(+), 22 deletions(-) diff --git a/src/config.ts b/src/config.ts index 3fcf4a5..55e8305 100644 --- a/src/config.ts +++ b/src/config.ts @@ -17,6 +17,7 @@ export interface BridgeConfig { bridgeHostname?: string; cmsApiUrl: string; notifyDebouncePolls: number; + notifyTimeoutMs: number; /** When true, log each Modbus poll and each HTTP request. */ verboseLogging: boolean; /** Force a register profile instead of auto-detecting from the inverter model. */ @@ -106,6 +107,7 @@ export function loadConfig(): BridgeConfig { bridgeHostname: readOptional('BRIDGE_HOSTNAME'), cmsApiUrl: readOptional('CMS_API_URL') ?? 'https://checkmy.solar', notifyDebouncePolls: readInt('NOTIFY_DEBOUNCE_POLLS', 2), + notifyTimeoutMs: readInt('NOTIFY_TIMEOUT_MS', 5_000), verboseLogging: readBoolean('BRIDGE_VERBOSE_LOG', false), inverterProfile: readProfileId('INVERTER_PROFILE'), modbusConnection: readConnectionType('MODBUS_CONNECTION', 'aux'), diff --git a/src/index.ts b/src/index.ts index 74322e3..1f40a04 100644 --- a/src/index.ts +++ b/src/index.ts @@ -28,7 +28,7 @@ async function runPollCycle( ]); store.upsert(telemetry, telemetry.sampledAt); - await workModeNotifier.handleSample(telemetry); + workModeNotifier.handleSample(telemetry); const todayTotals = mapH1G2TodayTotalsSnapshotToFoxShape(todayTotalsSnapshot); if (todayTotals) { store.upsertTodayTotals(todayTotals, todayTotalsSnapshot.sampledAt); @@ -45,14 +45,14 @@ async function main(): Promise { const config = loadConfig(); const store = new RealtimeStore(config.dataDir, { verboseLogging: config.verboseLogging }); const aggregator = new HourlyAggregator(store, config.siteTimezone); - const latestSnapshot = store.getLatest(); const workModeNotifier = new WorkModeNotifier( { apiUrl: config.cmsApiUrl, bridgeToken: config.bridgeToken, debouncePolls: config.notifyDebouncePolls, + timeoutMs: config.notifyTimeoutMs, }, - latestSnapshot?.telemetry.workMode + store ); let detectedInverter: ReturnType = null; diff --git a/src/notify/workModeNotifier.test.ts b/src/notify/workModeNotifier.test.ts index eb45bee..803a30e 100644 --- a/src/notify/workModeNotifier.test.ts +++ b/src/notify/workModeNotifier.test.ts @@ -1,6 +1,10 @@ import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest'; import type { ModbusRealtimeTelemetry } from '@checkmysolar/modbus-telemetry'; -import { WorkModeNotifier, postWorkModeEvent } from './workModeNotifier.js'; +import { + WorkModeNotifier, + type WorkModeNotifyStateStore, + postWorkModeEvent, +} from './workModeNotifier.js'; function sampleTelemetry(workMode: number): ModbusRealtimeTelemetry { return { @@ -36,6 +40,21 @@ function sampleTelemetry(workMode: number): ModbusRealtimeTelemetry { }; } +function createStateStore(initialWorkMode?: number): WorkModeNotifyStateStore { + let lastEmittedWorkMode = initialWorkMode; + return { + getLastEmittedWorkMode: () => lastEmittedWorkMode, + setLastEmittedWorkMode: (workMode) => { + lastEmittedWorkMode = workMode; + }, + }; +} + +async function flushNotifications(): Promise { + await Promise.resolve(); + await Promise.resolve(); +} + describe('WorkModeNotifier', () => { beforeEach(() => { vi.stubGlobal('fetch', vi.fn()); @@ -54,14 +73,17 @@ describe('WorkModeNotifier', () => { apiUrl: 'https://checkmy.solar', bridgeToken: 'cms_bridge_test', debouncePolls: 2, + timeoutMs: 5_000, }, - 0 + createStateStore(0) ); - await notifier.handleSample(sampleTelemetry(3)); + notifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); expect(fetchMock).not.toHaveBeenCalled(); - await notifier.handleSample(sampleTelemetry(3)); + notifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); expect(fetchMock).toHaveBeenCalledTimes(1); expect(fetchMock).toHaveBeenCalledWith( 'https://checkmy.solar/api/bridge/events/work-mode', @@ -89,15 +111,18 @@ describe('WorkModeNotifier', () => { apiUrl: 'https://checkmy.solar', bridgeToken: 'cms_bridge_test', debouncePolls: 2, + timeoutMs: 5_000, }, - 0 + createStateStore(0) ); - await notifier.handleSample(sampleTelemetry(3)); - await notifier.handleSample(sampleTelemetry(4)); + notifier.handleSample(sampleTelemetry(3)); + notifier.handleSample(sampleTelemetry(4)); + await flushNotifications(); expect(fetchMock).not.toHaveBeenCalled(); - await notifier.handleSample(sampleTelemetry(4)); + notifier.handleSample(sampleTelemetry(4)); + await flushNotifications(); expect(fetchMock).toHaveBeenCalledTimes(1); expect(fetchMock.mock.calls[0]?.[1]).toEqual( expect.objectContaining({ @@ -120,11 +145,115 @@ describe('WorkModeNotifier', () => { apiUrl: 'https://checkmy.solar', bridgeToken: 'cms_bridge_test', debouncePolls: 1, + timeoutMs: 5_000, }, - 0 + createStateStore(0) ); - await expect(notifier.handleSample(sampleTelemetry(1))).resolves.toBeUndefined(); + notifier.handleSample(sampleTelemetry(1)); + await flushNotifications(); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); + + it('retries after a failed post when the mode is still unchanged', async () => { + const fetchMock = vi.mocked(fetch); + let resolveFirst: ((response: Response) => void) | undefined; + fetchMock + .mockImplementationOnce( + () => + new Promise((resolve) => { + resolveFirst = resolve; + }) + ) + .mockResolvedValueOnce(new Response(null, { status: 204 })); + + const stateStore = createStateStore(0); + const notifier = new WorkModeNotifier( + { + apiUrl: 'https://checkmy.solar', + bridgeToken: 'cms_bridge_test', + debouncePolls: 1, + timeoutMs: 5_000, + }, + stateStore + ); + + notifier.handleSample(sampleTelemetry(3)); + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(1)); + const firstRequest = fetchMock.mock.results[0]?.value as Promise; + resolveFirst!(new Response('bad request', { status: 400 })); + await firstRequest; + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(stateStore.getLastEmittedWorkMode()).toBe(0); + + notifier.handleSample(sampleTelemetry(3)); + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(2)); + await vi.waitFor(() => expect(stateStore.getLastEmittedWorkMode()).toBe(3)); + }); + + it('does not post duplicate notifications while a request is in flight', async () => { + const fetchMock = vi.mocked(fetch); + let resolveFetch: (() => void) | undefined; + fetchMock.mockImplementation( + () => + new Promise((resolve) => { + resolveFetch = () => resolve(new Response(null, { status: 204 })); + }) + ); + + const notifier = new WorkModeNotifier( + { + apiUrl: 'https://checkmy.solar', + bridgeToken: 'cms_bridge_test', + debouncePolls: 1, + timeoutMs: 5_000, + }, + createStateStore(0) + ); + + notifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); + expect(fetchMock).toHaveBeenCalledTimes(1); + + notifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); + expect(fetchMock).toHaveBeenCalledTimes(1); + + resolveFetch?.(); + await flushNotifications(); + }); + + it('restores last emitted work mode from persistent state on startup', async () => { + const fetchMock = vi.mocked(fetch); + fetchMock.mockResolvedValue(new Response(null, { status: 204 })); + + const stateStore = createStateStore(0); + const notifier = new WorkModeNotifier( + { + apiUrl: 'https://checkmy.solar', + bridgeToken: 'cms_bridge_test', + debouncePolls: 1, + timeoutMs: 5_000, + }, + stateStore + ); + + notifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); + expect(stateStore.getLastEmittedWorkMode()).toBe(3); + + const restartedNotifier = new WorkModeNotifier( + { + apiUrl: 'https://checkmy.solar', + bridgeToken: 'cms_bridge_test', + debouncePolls: 1, + timeoutMs: 5_000, + }, + stateStore + ); + + restartedNotifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); expect(fetchMock).toHaveBeenCalledTimes(1); }); }); @@ -144,14 +273,16 @@ describe('postWorkModeEvent', () => { await expect( postWorkModeEvent( - { apiUrl: 'https://checkmy.solar/', bridgeToken: 'cms_bridge_test' }, + { apiUrl: 'https://checkmy.solar/', bridgeToken: 'cms_bridge_test', timeoutMs: 5_000 }, { workMode: 2, sampledAt: '2026-07-09T11:59:30.000Z' } ) ).resolves.toBeUndefined(); expect(fetchMock).toHaveBeenCalledWith( 'https://checkmy.solar/api/bridge/events/work-mode', - expect.any(Object) + expect.objectContaining({ + signal: expect.any(AbortSignal), + }) ); }); }); diff --git a/src/notify/workModeNotifier.ts b/src/notify/workModeNotifier.ts index 7a51d86..12f2b60 100644 --- a/src/notify/workModeNotifier.ts +++ b/src/notify/workModeNotifier.ts @@ -5,6 +5,12 @@ export interface WorkModeNotifierOptions { apiUrl: string; bridgeToken: string; debouncePolls: number; + timeoutMs: number; +} + +export interface WorkModeNotifyStateStore { + getLastEmittedWorkMode(): number | undefined; + setLastEmittedWorkMode(workMode: number): void; } export interface WorkModeEventPayload { @@ -18,21 +24,22 @@ export class WorkModeNotifier { private pendingWorkMode: number | undefined; private pendingCount = 0; private lastEmittedWorkMode: number | undefined; + private inFlightWorkMode: number | undefined; constructor( private readonly options: WorkModeNotifierOptions, - initialWorkMode?: number + private readonly stateStore: WorkModeNotifyStateStore ) { - this.lastEmittedWorkMode = initialWorkMode; + this.lastEmittedWorkMode = stateStore.getLastEmittedWorkMode(); } - async handleSample(telemetry: ModbusRealtimeTelemetry): Promise { + handleSample(telemetry: ModbusRealtimeTelemetry): void { const workMode = telemetry.workMode; if (workMode === undefined || workMode === null) { return; } - if (workMode === this.lastEmittedWorkMode) { + if (workMode === this.lastEmittedWorkMode || workMode === this.inFlightWorkMode) { this.pendingWorkMode = undefined; this.pendingCount = 0; return; @@ -49,6 +56,17 @@ export class WorkModeNotifier { return; } + this.inFlightWorkMode = workMode; + this.pendingWorkMode = undefined; + this.pendingCount = 0; + + void this.emitWorkMode(workMode, telemetry); + } + + private async emitWorkMode( + workMode: number, + telemetry: ModbusRealtimeTelemetry + ): Promise { const previousWorkMode = this.lastEmittedWorkMode; try { @@ -59,16 +77,19 @@ export class WorkModeNotifier { soc: telemetry.SoC, }); this.lastEmittedWorkMode = workMode; - this.pendingWorkMode = undefined; - this.pendingCount = 0; + this.stateStore.setLastEmittedWorkMode(workMode); } catch (error) { console.error(`Work mode notification failed: ${formatError(error)}`); + } finally { + if (this.inFlightWorkMode === workMode) { + this.inFlightWorkMode = undefined; + } } } } export async function postWorkModeEvent( - options: Pick, + options: Pick, payload: WorkModeEventPayload ): Promise { const url = `${options.apiUrl.replace(/\/$/, '')}/api/bridge/events/work-mode`; @@ -79,6 +100,7 @@ export async function postWorkModeEvent( Authorization: `Bearer ${options.bridgeToken}`, }, body: JSON.stringify(payload), + signal: AbortSignal.timeout(options.timeoutMs), }); if (response.status === 401) { diff --git a/src/storage/sqlite.ts b/src/storage/sqlite.ts index fc78b27..c56c2ed 100644 --- a/src/storage/sqlite.ts +++ b/src/storage/sqlite.ts @@ -93,6 +93,12 @@ CREATE TABLE IF NOT EXISTS samples ( ts TEXT PRIMARY KEY, telemetry_json TEXT NOT NULL ); + +CREATE TABLE IF NOT EXISTS work_mode_notify_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + last_emitted_work_mode INTEGER, + updated_at TEXT NOT NULL +); `; function mapHourRow(row: HourRowDb): DayHourlyRow { @@ -164,6 +170,32 @@ export class RealtimeStore { this.logWrite('upsert realtime_snapshot', `sampledAt=${sampledAt}`); } + getLastEmittedWorkMode(): number | undefined { + const row = this.db + .prepare('SELECT last_emitted_work_mode FROM work_mode_notify_state WHERE id = 1') + .get() as { last_emitted_work_mode: number | null } | undefined; + + if (!row || row.last_emitted_work_mode === null) { + return undefined; + } + + return row.last_emitted_work_mode; + } + + setLastEmittedWorkMode(workMode: number): void { + const updatedAt = new Date().toISOString(); + this.db + .prepare( + `INSERT INTO work_mode_notify_state (id, last_emitted_work_mode, updated_at) + VALUES (1, ?, ?) + ON CONFLICT(id) DO UPDATE SET + last_emitted_work_mode = excluded.last_emitted_work_mode, + updated_at = excluded.updated_at` + ) + .run(workMode, updatedAt); + this.logWrite('setLastEmittedWorkMode', `workMode=${workMode}`); + } + getLatest(): RealtimeSnapshot | null { const row = this.db .prepare('SELECT telemetry_json, sampled_at, updated_at FROM realtime_snapshot WHERE id = 1')