Skip to content

Commit db758d0

Browse files
authored
fix(iobroker): give the SDK sole reconnect ownership (#119)
StreamHandler ran its own five-attempt reconnect timer alongside the SDK's unbounded exponential backoff, and re-registered every listener on each scheduled reconnect (duplicate handlers over time). Listeners now register exactly once in connect(); stream_error and terminal auth_failure are handled so info.connection accurately reflects transient-down vs permanently-stopped state.
1 parent e82e608 commit db758d0

4 files changed

Lines changed: 146 additions & 103 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+
Give the SDK sole ownership of stream reconnection instead of running a competing five-attempt reconnect timer alongside it: `StreamHandler` now registers its stream listeners exactly once (fixing double-registration on repeated `connect()` calls) and handles `stream_error` and terminal `auth_failure`, so `info.connection` accurately reflects whether the stream is transiently down or has stopped for good.

packages/iobroker.teslemetry/lib/StreamHandler.ts

Lines changed: 68 additions & 98 deletions
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,12 @@ import { StateManager } from './StateManager.js';
33
import { EnergyHandler } from './EnergyHandler.js';
44

55
export class StreamHandler {
6-
private reconnectAttempts = 0;
7-
private maxReconnectAttempts = 5;
8-
private reconnectTimer?: NodeJS.Timeout;
9-
private isConnecting = false;
6+
// The SDK's TeslemetryStream owns reconnection (unbounded exponential
7+
// backoff, then a permanent stop after two consecutive auth failures) -
8+
// this class only listens and reflects that state, it never schedules
9+
// its own reconnect. Listeners are registered once; sse.on() replays
10+
// cached values into the handler and re-listening would double-dispatch.
11+
private listenersRegistered = false;
1012

1113
constructor(
1214
private adapter: ioBroker.Adapter,
@@ -16,67 +18,74 @@ export class StreamHandler {
1618
) {}
1719

1820
/**
19-
* Connect to SSE stream and set up event handlers
21+
* Register stream event handlers (once) and connect to the SSE stream.
22+
* Safe to call again after a connect: the SDK's connect() is a no-op
23+
* while already active.
2024
*/
2125
async connect(): Promise<void> {
22-
if (this.isConnecting) {
23-
this.adapter.log.debug('Stream connection already in progress');
24-
return;
26+
const sse = this.teslemetry.sse;
27+
28+
if (!this.listenersRegistered) {
29+
this.listenersRegistered = true;
30+
this.registerListeners(sse);
2531
}
2632

27-
this.isConnecting = true;
2833
this.adapter.log.info('Connecting to Teslemetry SSE stream...');
34+
await sse.connect();
35+
}
2936

30-
try {
31-
const sse = this.teslemetry.sse;
32-
33-
// Set up event handlers
34-
sse.on('connect', () => {
35-
this.adapter.log.info('SSE stream connected');
36-
this.reconnectAttempts = 0;
37-
this.isConnecting = false;
38-
this.adapter.setStateAsync('info.connection', true, true);
39-
});
40-
41-
sse.on('disconnect', () => {
42-
this.adapter.log.warn('SSE stream disconnected');
43-
this.isConnecting = false;
44-
this.adapter.setStateAsync('info.connection', false, true);
45-
this.scheduleReconnect();
46-
});
47-
48-
// Handle vehicle data updates
49-
sse.on('data', (event: any) => {
50-
this.handleDataEvent(event);
51-
});
52-
53-
// Handle vehicle state changes (online/asleep/offline)
54-
sse.on('state', (event: any) => {
55-
this.handleStateEvent(event);
56-
});
57-
58-
// Handle alerts
59-
sse.on('alerts', (event: any) => {
60-
this.handleAlertEvent(event);
61-
});
62-
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-
73-
// Connect to stream
74-
await sse.connect();
75-
} catch (error: any) {
76-
this.adapter.log.error(`Failed to connect to SSE stream: ${error.message}`);
77-
this.isConnecting = false;
78-
this.scheduleReconnect();
79-
}
37+
private registerListeners(sse: Teslemetry['sse']): void {
38+
sse.on('connect', () => {
39+
this.adapter.log.info('SSE stream connected');
40+
this.adapter.setStateAsync('info.connection', true, true);
41+
});
42+
43+
// The SDK retries transient disconnects on its own; this only reflects
44+
// current connectivity; it does not imply a terminal stop.
45+
sse.on('disconnect', () => {
46+
this.adapter.log.warn('SSE stream disconnected');
47+
this.adapter.setStateAsync('info.connection', false, true);
48+
});
49+
50+
sse.on('stream_error', ({ error, status, retries }) => {
51+
const message = error instanceof Error ? error.message : String(error);
52+
this.adapter.log.warn(
53+
`SSE stream error (status ${status ?? 'unknown'}, attempt ${retries}): ${message}`
54+
);
55+
});
56+
57+
sse.on('auth_failure', (error) => {
58+
this.adapter.log.error(
59+
`SSE stream authentication failed twice in a row and has stopped permanently: ${error.message}. ` +
60+
'Fix the access token and restart the adapter to resume streaming.'
61+
);
62+
this.adapter.setStateAsync('info.connection', false, true);
63+
});
64+
65+
// Handle vehicle data updates
66+
sse.on('data', (event: any) => {
67+
this.handleDataEvent(event);
68+
});
69+
70+
// Handle vehicle state changes (online/asleep/offline)
71+
sse.on('state', (event: any) => {
72+
this.handleStateEvent(event);
73+
});
74+
75+
// Handle alerts
76+
sse.on('alerts', (event: any) => {
77+
this.handleAlertEvent(event);
78+
});
79+
80+
// Handle energy site live power/battery/grid updates
81+
sse.on('live_status', (event: any) => {
82+
this.handleLiveStatusEvent(event);
83+
});
84+
85+
// Handle energy site settings updates (operation mode, reserves, ...)
86+
sse.on('site_info', (event: any) => {
87+
this.handleSiteInfoEvent(event);
88+
});
8089
}
8190

8291
/**
@@ -85,11 +94,6 @@ export class StreamHandler {
8594
disconnect(): void {
8695
this.adapter.log.info('Disconnecting from SSE stream...');
8796

88-
if (this.reconnectTimer) {
89-
clearTimeout(this.reconnectTimer);
90-
this.reconnectTimer = undefined;
91-
}
92-
9397
try {
9498
this.teslemetry.sse.disconnect();
9599
this.adapter.setStateAsync('info.connection', false, true);
@@ -202,38 +206,4 @@ export class StreamHandler {
202206
this.adapter.log.error(`Error handling site_info event: ${error.message}`);
203207
}
204208
}
205-
206-
/**
207-
* Schedule automatic reconnection with exponential backoff
208-
*/
209-
private scheduleReconnect(): void {
210-
if (this.reconnectAttempts >= this.maxReconnectAttempts) {
211-
this.adapter.log.error(
212-
`Max reconnection attempts (${this.maxReconnectAttempts}) reached. Stopping reconnection attempts.`
213-
);
214-
return;
215-
}
216-
217-
if (this.reconnectTimer) {
218-
clearTimeout(this.reconnectTimer);
219-
}
220-
221-
const delay = Math.min(1000 * Math.pow(2, this.reconnectAttempts), 30000);
222-
this.reconnectAttempts++;
223-
224-
this.adapter.log.info(
225-
`Scheduling reconnection attempt ${this.reconnectAttempts}/${this.maxReconnectAttempts} in ${delay / 1000}s`
226-
);
227-
228-
this.reconnectTimer = setTimeout(() => {
229-
this.connect();
230-
}, delay);
231-
}
232-
233-
/**
234-
* Reset reconnection attempts (call this after successful manual reconnection)
235-
*/
236-
resetReconnectAttempts(): void {
237-
this.reconnectAttempts = 0;
238-
}
239209
}

packages/iobroker.teslemetry/src/main.ts

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -153,13 +153,11 @@ class TeslemetryAdapter extends utils.Adapter {
153153
this.pollInterval = undefined;
154154
}
155155

156-
// Disconnect streaming
156+
// Disconnect streaming (StreamHandler.disconnect() also closes the
157+
// underlying SSE connection - no separate teslemetry.sse.disconnect() needed)
157158
if (this.streamHandler) {
158159
this.streamHandler.disconnect();
159-
}
160-
161-
// Disconnect SSE
162-
if (this.teslemetry) {
160+
} else if (this.teslemetry) {
163161
this.teslemetry.sse.disconnect();
164162
}
165163

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

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -127,3 +127,73 @@ test('live_status events for an unselected/unregistered energy site are ignored
127127

128128
assert.equal(states.get(`energy.${SITE_ID}.live.solar_power`), undefined);
129129
});
130+
131+
const STREAM_EVENTS = ['connect', 'disconnect', 'stream_error', 'auth_failure', 'data', 'state', 'alerts', 'live_status', 'site_info'] as const;
132+
133+
test('repeated disconnect/reconnect cycles keep listener counts constant (regression: connect() used to re-register a fresh set of listeners every time)', async () => {
134+
const teslemetry = new Teslemetry('fake-token');
135+
// The real SDK connect() no-ops while already active; simulate the SDK
136+
// marking itself active on each call, as it would after a real reconnect.
137+
(teslemetry.sse as any).connect = () => {
138+
(teslemetry.sse as any).active = true;
139+
return Promise.resolve();
140+
};
141+
142+
const { adapter } = createFakeAdapter();
143+
const stateManager = new StateManager(adapter);
144+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, new EnergyHandler(adapter, teslemetry, stateManager));
145+
146+
await streamHandler.connect();
147+
// data/state/alerts/live_status/site_info also carry the SDK's own
148+
// internal cache listener (TeslemetryStream's constructor), so their
149+
// baseline is 2 (SDK + ours) rather than 1 - what matters is that a
150+
// reconnect cycle never adds another one on top of that baseline.
151+
const countsAfterFirstConnect = STREAM_EVENTS.map((event) => teslemetry.sse.listenerCount(event));
152+
assert.ok(
153+
countsAfterFirstConnect.every((count) => count >= 1),
154+
`expected at least one listener per event after the first connect, got: ${JSON.stringify(countsAfterFirstConnect)}`,
155+
);
156+
157+
for (let i = 0; i < 3; i++) {
158+
teslemetry.sse.emit('disconnect');
159+
await streamHandler.connect();
160+
}
161+
162+
const countsAfterCycles = STREAM_EVENTS.map((event) => teslemetry.sse.listenerCount(event));
163+
assert.deepEqual(countsAfterCycles, countsAfterFirstConnect);
164+
});
165+
166+
test('stream_error is logged without scheduling a competing reconnect (regression: adapter used to run its own reconnect timer alongside the SDK)', async () => {
167+
const teslemetry = new Teslemetry('fake-token');
168+
(teslemetry.sse as any).connect = () => Promise.resolve();
169+
170+
const { adapter, logs } = createFakeAdapter();
171+
const stateManager = new StateManager(adapter);
172+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, new EnergyHandler(adapter, teslemetry, stateManager));
173+
await streamHandler.connect();
174+
175+
teslemetry.sse.emit('stream_error', { error: new Error('boom'), status: 503, retries: 2 } as any);
176+
await new Promise((resolve) => setTimeout(resolve, 0));
177+
178+
assert.ok(logs.some((l) => l.level === 'warn' && l.message.includes('boom') && l.message.includes('503')));
179+
});
180+
181+
test('terminal auth_failure marks info.connection false and logs an error (regression: adapter never listened for auth_failure at all)', async () => {
182+
const teslemetry = new Teslemetry('fake-token');
183+
(teslemetry.sse as any).connect = () => Promise.resolve();
184+
185+
const { adapter, states, logs } = createFakeAdapter();
186+
const stateManager = new StateManager(adapter);
187+
const streamHandler = new StreamHandler(adapter, teslemetry, stateManager, new EnergyHandler(adapter, teslemetry, stateManager));
188+
await streamHandler.connect();
189+
190+
teslemetry.sse.emit('connect');
191+
await new Promise((resolve) => setTimeout(resolve, 0));
192+
assert.equal(states.get('info.connection'), true);
193+
194+
teslemetry.sse.emit('auth_failure', new Error('token expired') as any);
195+
await new Promise((resolve) => setTimeout(resolve, 0));
196+
197+
assert.equal(states.get('info.connection'), false);
198+
assert.ok(logs.some((l) => l.level === 'error' && l.message.includes('stopped permanently')));
199+
});

0 commit comments

Comments
 (0)