Skip to content

Commit b2a26ba

Browse files
authored
fix(iobroker): stream energy live_status/site_info in default streaming mode (#115)
* fix(iobroker): stream energy live_status/site_info in default streaming mode Energy data was fetched once at startup and then never updated when streaming was enabled (the default), since StreamHandler only listened to vehicle topics. Register live_status/site_info listeners alongside the existing vehicle ones and route them through the same flat energy parser used by the REST seed, so energy states keep updating from the stream instead of freezing after the initial fetch. * fix(iobroker): ignore energy stream events for unregistered sites live_status/site_info arrive on the account-level stream for every accessible site, not just the ones this adapter instance selected - without a check, an unselected site's events wrote to object paths that were never created. Filter both handlers against EnergyHandler.getRegisteredSites() before forwarding to the parser. * fix(iobroker): seed energy state before opening the SSE stream streamHandler.connect() only starts the SDK's background reconnect loop and returns immediately, so the previous ordering let a live_status/site_info event land before the initial energy REST fetch resolved - the fetch would then overwrite it with a stale snapshot. Fetch energy state first, then open the stream.
1 parent 89412c1 commit b2a26ba

6 files changed

Lines changed: 168 additions & 15 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"iobroker.teslemetry": patch
3+
---
4+
5+
Consume `live_status`/`site_info` SSE events per energy site so energy data keeps updating in the default streaming mode, instead of freezing after the initial startup fetch. The REST fetch still seeds a deterministic initial value; the stream now keeps it current afterward without polling.

packages/iobroker.teslemetry/README.md

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -24,10 +24,10 @@ ioBroker adapter for controlling Tesla vehicles and energy sites via the Tesleme
2424
- 🏠 **Off-Grid** - Configure backup reserves
2525

2626
### Real-Time Updates
27-
- 📡 **SSE Streaming** - Live vehicle updates via Server-Sent Events
27+
- 📡 **SSE Streaming** - Live vehicle and energy site updates via Server-Sent Events
2828
- 🔄 **Automatic Reconnection** - Robust connection handling
2929
- 💤 **Sleep Mode Aware** - Won't wake sleeping vehicles unnecessarily
30-
- 📊 **Polling** - Energy site data updates via polling; vehicles can also use polling instead of streaming
30+
- 📊 **Polling** - Optional fallback for vehicles and energy sites when streaming is disabled
3131

3232
## Prerequisites
3333

@@ -82,12 +82,11 @@ npm install iobroker.teslemetry
8282

8383
**Real-time Streaming (recommended)**
8484
- Enable **"Enable real-time streaming (SSE)"**
85-
- Provides instant updates when vehicle data changes
85+
- Provides instant updates when vehicle or energy site data changes
8686
- More efficient and responsive than polling
87-
- Energy site data is always updated via polling, regardless of this setting
8887

8988
**Polling Mode**
90-
- Disable streaming to use polling instead for vehicles too
89+
- Disable streaming to use polling instead for both vehicles and energy sites
9190
- Set **Poll Interval** (minimum 30 seconds)
9291
- Less efficient but works if streaming has issues
9392

packages/iobroker.teslemetry/lib/StateManager.ts

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -243,8 +243,11 @@ export class StateManager {
243243
}
244244

245245
/**
246-
* Update energy site data from the merged { ...siteInfo.response, ...liveStatus.response }
247-
* of getSiteInfo() and getLiveStatus() - both are flat, no `live_status` wrapper.
246+
* Update energy site data from a flat object matching the REST getSiteInfo()/getLiveStatus()
247+
* response shape - either their merged `{ ...siteInfo.response, ...liveStatus.response }`
248+
* (the REST seed), or a single SSE `live_status`/`site_info` event's payload (both are
249+
* flat, no wrapper, and use the same field names as their REST counterparts). Every field
250+
* is optional, so a partial single-topic payload only touches the states it carries.
248251
*/
249252
async updateEnergySiteData(siteId: number, data: any): Promise<void> {
250253
const base = `energy.${siteId}`;

packages/iobroker.teslemetry/lib/StreamHandler.ts

Lines changed: 59 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { Teslemetry } from '@teslemetry/api';
22
import { StateManager } from './StateManager.js';
3+
import { EnergyHandler } from './EnergyHandler.js';
34

45
export class StreamHandler {
56
private reconnectAttempts = 0;
@@ -10,7 +11,8 @@ export class StreamHandler {
1011
constructor(
1112
private adapter: ioBroker.Adapter,
1213
private teslemetry: Teslemetry,
13-
private stateManager: StateManager
14+
private stateManager: StateManager,
15+
private energyHandler: EnergyHandler
1416
) {}
1517

1618
/**
@@ -58,6 +60,16 @@ export class StreamHandler {
5860
this.handleAlertEvent(event);
5961
});
6062

63+
// Handle energy site live power/battery/grid updates
64+
sse.on('live_status', (event: any) => {
65+
this.handleLiveStatusEvent(event);
66+
});
67+
68+
// Handle energy site settings updates (operation mode, reserves, ...)
69+
sse.on('site_info', (event: any) => {
70+
this.handleSiteInfoEvent(event);
71+
});
72+
6173
// Connect to stream
6274
await sse.connect();
6375
} catch (error: any) {
@@ -145,6 +157,52 @@ export class StreamHandler {
145157
}
146158
}
147159

160+
/**
161+
* Handle energy site `live_status` events (solar/battery/grid/load power, SOC)
162+
*/
163+
private async handleLiveStatusEvent(event: any): Promise<void> {
164+
try {
165+
const { site_id, live_status } = event;
166+
167+
if (!site_id || !live_status) {
168+
return;
169+
}
170+
171+
const siteId = Number(site_id);
172+
if (!this.energyHandler.getRegisteredSites().includes(siteId)) {
173+
return;
174+
}
175+
176+
this.adapter.log.debug(`Received live_status event for energy site ${site_id}`);
177+
await this.stateManager.updateEnergySiteData(siteId, live_status);
178+
} catch (error: any) {
179+
this.adapter.log.error(`Error handling live_status event: ${error.message}`);
180+
}
181+
}
182+
183+
/**
184+
* Handle energy site `site_info` events (operation mode, reserves, tariff id, ...)
185+
*/
186+
private async handleSiteInfoEvent(event: any): Promise<void> {
187+
try {
188+
const { site_id, site_info } = event;
189+
190+
if (!site_id || !site_info) {
191+
return;
192+
}
193+
194+
const siteId = Number(site_id);
195+
if (!this.energyHandler.getRegisteredSites().includes(siteId)) {
196+
return;
197+
}
198+
199+
this.adapter.log.debug(`Received site_info event for energy site ${site_id}`);
200+
await this.stateManager.updateEnergySiteData(siteId, site_info);
201+
} catch (error: any) {
202+
this.adapter.log.error(`Error handling site_info event: ${error.message}`);
203+
}
204+
}
205+
148206
/**
149207
* Schedule automatic reconnection with exponential backoff
150208
*/

packages/iobroker.teslemetry/src/main.ts

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -112,23 +112,26 @@ class TeslemetryAdapter extends utils.Adapter {
112112
// Subscribe to all state changes
113113
this.subscribeStates('*');
114114

115+
// Seed energy state before opening the stream - connect() only starts the SDK's
116+
// background loop and returns immediately, so a stream event could otherwise land
117+
// before this REST fetch resolves and get overwritten by the stale snapshot.
118+
this.log.info('Fetching initial energy site data...');
119+
await this.energyHandler.fetchAllSiteData();
120+
115121
// Set up streaming or polling
116122
if (this.config.enableStreaming !== false) {
117123
this.log.info('Starting SSE streaming...');
118-
this.streamHandler = new StreamHandler(this, this.teslemetry, this.stateManager);
124+
this.streamHandler = new StreamHandler(this, this.teslemetry, this.stateManager, this.energyHandler);
119125
await this.streamHandler.connect();
120126
} else {
121127
this.log.info('SSE streaming disabled, using polling');
122128
this.startPolling();
123129
}
124130

125-
// Do initial data fetch
131+
// Do initial vehicle data fetch
126132
this.log.info('Fetching initial vehicle data...');
127133
await this.vehicleHandler.fetchAllVehicleData(false);
128134

129-
this.log.info('Fetching initial energy site data...');
130-
await this.energyHandler.fetchAllSiteData();
131-
132135
this.log.info('Teslemetry adapter started successfully');
133136
await this.setStateAsync('info.connection', true, true);
134137
} catch (error: any) {

packages/iobroker.teslemetry/test/StreamHandler.test.ts

Lines changed: 87 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,17 +3,20 @@ import assert from 'node:assert/strict';
33
import { Teslemetry } from '@teslemetry/api';
44
import { StreamHandler } from '../lib/StreamHandler.js';
55
import { StateManager } from '../lib/StateManager.js';
6+
import { EnergyHandler } from '../lib/EnergyHandler.js';
67
import { createFakeAdapter } from './fakeAdapter.js';
78

89
const VIN = '5YJSA1E14FF000000';
10+
const SITE_ID = 123;
911

1012
test('data events update states via the SSE flat-signal parser', async () => {
1113
const teslemetry = new Teslemetry('fake-token');
1214
// Avoid a real network connection; sse.connect() only needs to resolve.
1315
(teslemetry.sse as any).connect = () => Promise.resolve();
1416

1517
const { adapter, states } = createFakeAdapter();
16-
const streamHandler = new StreamHandler(adapter, teslemetry, new StateManager(adapter));
18+
const stateManager = new StateManager(adapter);
19+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, new EnergyHandler(adapter, teslemetry, stateManager));
1720
await streamHandler.connect();
1821

1922
teslemetry.sse.emit('data', { vin: VIN, data: { BatteryLevel: 71 } } as any);
@@ -28,7 +31,8 @@ test('alert events log each alert from the `alerts` array (regression: destructu
2831
(teslemetry.sse as any).connect = () => Promise.resolve();
2932

3033
const { adapter, logs } = createFakeAdapter();
31-
const streamHandler = new StreamHandler(adapter, teslemetry, new StateManager(adapter));
34+
const stateManager = new StateManager(adapter);
35+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, new EnergyHandler(adapter, teslemetry, stateManager));
3236
await streamHandler.connect();
3337

3438
teslemetry.sse.emit('alerts', {
@@ -42,3 +46,84 @@ test('alert events log each alert from the `alerts` array (regression: destructu
4246
`expected an alert log line, got: ${JSON.stringify(logs)}`,
4347
);
4448
});
49+
50+
test('live_status events update energy live-power states (regression: default streaming mode never wired up energy topics)', async () => {
51+
const teslemetry = new Teslemetry('fake-token');
52+
(teslemetry.sse as any).connect = () => Promise.resolve();
53+
54+
const { adapter, states } = createFakeAdapter();
55+
const stateManager = new StateManager(adapter);
56+
const energyHandler = new EnergyHandler(adapter, teslemetry, stateManager);
57+
energyHandler.registerSite(SITE_ID);
58+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, energyHandler);
59+
await streamHandler.connect();
60+
61+
teslemetry.sse.emit('live_status', { site_id: String(SITE_ID), live_status: { solar_power: 900, grid_status: 'Active' } } as any);
62+
await new Promise((resolve) => setTimeout(resolve, 0));
63+
64+
assert.equal(states.get(`energy.${SITE_ID}.live.solar_power`), 900);
65+
assert.equal(states.get(`energy.${SITE_ID}.live.grid_status`), 'Active');
66+
});
67+
68+
test('site_info events update energy operation states', async () => {
69+
const teslemetry = new Teslemetry('fake-token');
70+
(teslemetry.sse as any).connect = () => Promise.resolve();
71+
72+
const { adapter, states } = createFakeAdapter();
73+
const stateManager = new StateManager(adapter);
74+
const energyHandler = new EnergyHandler(adapter, teslemetry, stateManager);
75+
energyHandler.registerSite(SITE_ID);
76+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, energyHandler);
77+
await streamHandler.connect();
78+
79+
teslemetry.sse.emit('site_info', { site_id: String(SITE_ID), site_info: { default_real_mode: 'backup', backup_reserve_percent: 35 } } as any);
80+
await new Promise((resolve) => setTimeout(resolve, 0));
81+
82+
assert.equal(states.get(`energy.${SITE_ID}.operation.mode`), 'backup');
83+
assert.equal(states.get(`energy.${SITE_ID}.operation.backup_reserve_percent`), 35);
84+
});
85+
86+
test('energy live values change after the REST startup seed via live_status stream events, without polling (regression: default streaming mode left energy frozen after startup)', async () => {
87+
const teslemetry = new Teslemetry('fake-token');
88+
(teslemetry.sse as any).connect = () => Promise.resolve();
89+
90+
const site = teslemetry.api.getEnergySite(SITE_ID);
91+
(site as any).getLiveStatus = () => Promise.resolve({ response: { solar_power: 500 } });
92+
(site as any).getSiteInfo = () => Promise.resolve({ response: {} });
93+
94+
const { adapter, states } = createFakeAdapter();
95+
const stateManager = new StateManager(adapter);
96+
const energyHandler = new EnergyHandler(adapter, teslemetry, stateManager);
97+
energyHandler.registerSite(SITE_ID);
98+
99+
// Deterministic REST seed at startup
100+
await energyHandler.fetchSiteData(SITE_ID);
101+
assert.equal(states.get(`energy.${SITE_ID}.live.solar_power`), 500);
102+
103+
// The stream - not a poll interval - delivers the next value
104+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, energyHandler);
105+
await streamHandler.connect();
106+
107+
teslemetry.sse.emit('live_status', { site_id: String(SITE_ID), live_status: { solar_power: 1200 } } as any);
108+
await new Promise((resolve) => setTimeout(resolve, 0));
109+
110+
assert.equal(states.get(`energy.${SITE_ID}.live.solar_power`), 1200);
111+
});
112+
113+
test('live_status events for an unselected/unregistered energy site are ignored (regression: account-level stream delivers every accessible site, not just the ones this instance registered)', async () => {
114+
const teslemetry = new Teslemetry('fake-token');
115+
(teslemetry.sse as any).connect = () => Promise.resolve();
116+
117+
const { adapter, states } = createFakeAdapter();
118+
const stateManager = new StateManager(adapter);
119+
// Registers a different site; SITE_ID below was never selected for this instance.
120+
const energyHandler = new EnergyHandler(adapter, teslemetry, stateManager);
121+
energyHandler.registerSite(SITE_ID + 1);
122+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, energyHandler);
123+
await streamHandler.connect();
124+
125+
teslemetry.sse.emit('live_status', { site_id: String(SITE_ID), live_status: { solar_power: 900 } } as any);
126+
await new Promise((resolve) => setTimeout(resolve, 0));
127+
128+
assert.equal(states.get(`energy.${SITE_ID}.live.solar_power`), undefined);
129+
});

0 commit comments

Comments
 (0)