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..55e8305 100644 --- a/src/config.ts +++ b/src/config.ts @@ -15,6 +15,9 @@ export interface BridgeConfig { dataDir: string; siteTimezone: string; 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. */ @@ -102,6 +105,9 @@ 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), + 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 d49e243..1f40a04 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); + 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 workModeNotifier = new WorkModeNotifier( + { + apiUrl: config.cmsApiUrl, + bridgeToken: config.bridgeToken, + debouncePolls: config.notifyDebouncePolls, + timeoutMs: config.notifyTimeoutMs, + }, + store + ); 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..803a30e --- /dev/null +++ b/src/notify/workModeNotifier.test.ts @@ -0,0 +1,288 @@ +import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest'; +import type { ModbusRealtimeTelemetry } from '@checkmysolar/modbus-telemetry'; +import { + WorkModeNotifier, + type WorkModeNotifyStateStore, + 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', + }; +} + +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()); + }); + + 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, + timeoutMs: 5_000, + }, + createStateStore(0) + ); + + notifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); + expect(fetchMock).not.toHaveBeenCalled(); + + notifier.handleSample(sampleTelemetry(3)); + await flushNotifications(); + 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, + timeoutMs: 5_000, + }, + createStateStore(0) + ); + + notifier.handleSample(sampleTelemetry(3)); + notifier.handleSample(sampleTelemetry(4)); + await flushNotifications(); + expect(fetchMock).not.toHaveBeenCalled(); + + notifier.handleSample(sampleTelemetry(4)); + await flushNotifications(); + 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, + timeoutMs: 5_000, + }, + createStateStore(0) + ); + + 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); + }); +}); + +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', 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.objectContaining({ + signal: expect.any(AbortSignal), + }) + ); + }); +}); diff --git a/src/notify/workModeNotifier.ts b/src/notify/workModeNotifier.ts new file mode 100644 index 0000000..12f2b60 --- /dev/null +++ b/src/notify/workModeNotifier.ts @@ -0,0 +1,114 @@ +import type { ModbusRealtimeTelemetry } from '@checkmysolar/modbus-telemetry'; +import { formatError } from '../errors.js'; + +export interface WorkModeNotifierOptions { + apiUrl: string; + bridgeToken: string; + debouncePolls: number; + timeoutMs: number; +} + +export interface WorkModeNotifyStateStore { + getLastEmittedWorkMode(): number | undefined; + setLastEmittedWorkMode(workMode: number): void; +} + +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; + private inFlightWorkMode: number | undefined; + + constructor( + private readonly options: WorkModeNotifierOptions, + private readonly stateStore: WorkModeNotifyStateStore + ) { + this.lastEmittedWorkMode = stateStore.getLastEmittedWorkMode(); + } + + handleSample(telemetry: ModbusRealtimeTelemetry): void { + const workMode = telemetry.workMode; + if (workMode === undefined || workMode === null) { + return; + } + + if (workMode === this.lastEmittedWorkMode || workMode === this.inFlightWorkMode) { + 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; + } + + 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 { + await postWorkModeEvent(this.options, { + workMode, + sampledAt: telemetry.sampledAt, + previousWorkMode, + soc: telemetry.SoC, + }); + this.lastEmittedWorkMode = workMode; + 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, + 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), + signal: AbortSignal.timeout(options.timeoutMs), + }); + + 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}` : ''}`); + } +} 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')