Skip to content

Commit b6b3af9

Browse files
authored
fix(emulator): prevent race condition when stop() called during TCP initialization (#292)
Fixes race condition when stop() is called while start() Promise is still pending, plus three critical issues identified during code review. Root Cause: The stop() method checked if (!this.started || !this.server) and returned early if initialization hadn't completed. However, this.server is assigned immediately when start() creates the server, while this.started is only set to true after initialization completes. If stop() is called during this window, it returns without canceling initialization or cleaning up. Primary Fix (Issue #275): - Store the reject callback from start() Promise in initReject property - When stop() is called during initialization, reject the Promise - Clean up event listeners regardless of started state - Only attempt to close the server if it fully initialized Critical Fixes from Code Review: - Issue #293: Clear initialization timeout to prevent memory leak - Issue #294: Move listener cleanup before Promise rejection to prevent race condition - Issue #295: Add comprehensive restart test to verify transport reusability - TypeScript: Use delete operator for optional properties (exactOptionalPropertyTypes) Test Coverage: - 37/37 tests passing (100%) - Coverage: 96.22% statements, 82% branches, 100% functions - Added test for stop during initialization - Added test for restart after stop during initialization Closes #275, #293, #294, #295 Related: #296, #297, #298, #299
1 parent 46f887f commit b6b3af9

2 files changed

Lines changed: 199 additions & 13 deletions

File tree

packages/emulator/src/transports/tcp.test.ts

Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -180,6 +180,154 @@ describe('TcpTransport', () => {
180180
})
181181

182182
describe('error handling', () => {
183+
it('should handle stop() called during initialization', async () => {
184+
// Test for issue #275: Race condition when stop() is called during initialization
185+
// When start() Promise is pending and stop() is called, the start() Promise should
186+
// reject and event listeners should be cleaned up properly
187+
const { ServerTCP } = await import('modbus-serial')
188+
let localMockInstance: any
189+
;(ServerTCP as any).mockImplementationOnce((vector: any, _options: any) => {
190+
capturedServiceVector = vector
191+
initializedListener = undefined
192+
errorListener = undefined
193+
eventListeners = new Map()
194+
195+
// Create mock that extends EventEmitter
196+
localMockInstance = Object.create(EventEmitter.prototype)
197+
Object.assign(localMockInstance, {
198+
close: jest.fn((cb: (err: Error | null) => void) => cb(null)),
199+
socks: new Map(),
200+
})
201+
202+
// Initialize EventEmitter state
203+
EventEmitter.call(localMockInstance)
204+
205+
// Override on() to track listeners
206+
const originalOn = localMockInstance.on.bind(localMockInstance)
207+
localMockInstance.on = jest.fn((event: string, listener: any) => {
208+
if (event === 'initialized') {
209+
initializedListener = listener
210+
} else if (event === 'error') {
211+
errorListener = listener
212+
}
213+
if (!eventListeners.has(event)) {
214+
eventListeners.set(event, new Set())
215+
}
216+
eventListeners.get(event)!.add(listener)
217+
return originalOn(event, listener)
218+
})
219+
220+
// Don't emit initialized automatically - let test control timing
221+
222+
return localMockInstance
223+
})
224+
225+
const testTransport = new TcpTransport({ host: 'localhost', port: 8502 })
226+
const startPromise = testTransport.start()
227+
228+
// Call stop() while initialization is still pending
229+
await testTransport.stop()
230+
231+
// start() Promise should reject with appropriate error
232+
await expect(startPromise).rejects.toThrow('Transport stopped during initialization')
233+
234+
// Verify event listeners are cleaned up
235+
expect(localMockInstance.listenerCount('initialized')).toBe(0)
236+
expect(localMockInstance.listenerCount('error')).toBe(0)
237+
238+
// Server should not be closed (it never fully initialized), just cleaned up
239+
expect(localMockInstance.close).not.toHaveBeenCalled()
240+
}, 15000)
241+
242+
it('should allow restart after stop() during initialization', async () => {
243+
// Test for issue #275 acceptance criteria: rapid start() → stop() → start() sequence
244+
// Verifies transport is reusable after stop() interrupts initialization
245+
const { ServerTCP } = await import('modbus-serial')
246+
let mockInstance1: any
247+
let mockInstance2: any
248+
let callCount = 0
249+
250+
const originalMock = (ServerTCP as any).getMockImplementation()
251+
252+
;(ServerTCP as any).mockImplementation((vector: any, _options: any) => {
253+
capturedServiceVector = vector
254+
initializedListener = undefined
255+
errorListener = undefined
256+
eventListeners = new Map()
257+
callCount++
258+
259+
// Create mock that extends EventEmitter
260+
const localMockInstance = Object.create(EventEmitter.prototype)
261+
Object.assign(localMockInstance, {
262+
close: jest.fn((cb: (err: Error | null) => void) => cb(null)),
263+
socks: new Map(),
264+
})
265+
266+
// Initialize EventEmitter state
267+
EventEmitter.call(localMockInstance)
268+
269+
// Override on() to track listeners
270+
const originalOn = localMockInstance.on.bind(localMockInstance)
271+
localMockInstance.on = jest.fn((event: string, listener: any) => {
272+
if (event === 'initialized') {
273+
initializedListener = listener
274+
} else if (event === 'error') {
275+
errorListener = listener
276+
}
277+
if (!eventListeners.has(event)) {
278+
eventListeners.set(event, new Set())
279+
}
280+
eventListeners.get(event)!.add(listener)
281+
return originalOn(event, listener)
282+
})
283+
284+
// Store reference based on call count
285+
if (callCount === 1) {
286+
mockInstance1 = localMockInstance
287+
// Don't emit initialized for first attempt
288+
} else if (callCount === 2) {
289+
mockInstance2 = localMockInstance
290+
// Emit initialized for second attempt after a tick
291+
setImmediate(() => {
292+
if (initializedListener) {
293+
initializedListener()
294+
}
295+
})
296+
}
297+
298+
return localMockInstance
299+
})
300+
301+
const testTransport = new TcpTransport({ host: 'localhost', port: 8502 })
302+
303+
// First attempt - stop during init
304+
const startPromise1 = testTransport.start()
305+
await testTransport.stop()
306+
await expect(startPromise1).rejects.toThrow('Transport stopped during initialization')
307+
308+
// Verify first server was cleaned up
309+
expect(mockInstance1.listenerCount('initialized')).toBe(0)
310+
expect(mockInstance1.listenerCount('error')).toBe(0)
311+
312+
// Second attempt - should succeed
313+
const startPromise2 = testTransport.start()
314+
await expect(startPromise2).resolves.not.toThrow()
315+
316+
// Verify second server initialized properly
317+
expect(mockInstance2).toBeDefined()
318+
319+
// Clean up
320+
await testTransport.stop()
321+
322+
// Restore original mock implementation for subsequent tests
323+
if (originalMock) {
324+
;(ServerTCP as any).mockImplementation(originalMock)
325+
} else {
326+
// Restore default behavior
327+
jest.clearAllMocks()
328+
}
329+
}, 15000)
330+
183331
it('should not leak event listeners when start() fails due to port binding error', async () => {
184332
// Test for issue #274: Event listeners should be cleaned up on initialization failure
185333
// Before the fix, when start() fails with an error event, the 'initialized' and 'error'

packages/emulator/src/transports/tcp.ts

Lines changed: 51 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,8 @@ export class TcpTransport extends BaseTransport {
3030
private requestHandler?: (slaveId: number, request: Buffer) => Promise<Buffer>
3131
private server?: ServerTCP
3232
private started = false
33+
private initReject?: (reason: Error) => void
34+
private initTimeout?: NodeJS.Timeout
3335

3436
constructor(config: TcpTransportConfig) {
3537
super()
@@ -86,30 +88,46 @@ export class TcpTransport extends BaseTransport {
8688
// Wait for initialized event with timeout
8789
return new Promise<void>((resolve, reject) => {
8890
let isInitializing = true
91+
this.initReject = reject
8992

90-
const timeout = setTimeout(() => {
93+
const cleanupOnError = (): void => {
9194
isInitializing = false
95+
delete this.initReject
96+
if (this.initTimeout) {
97+
clearTimeout(this.initTimeout)
98+
delete this.initTimeout
99+
}
92100
EventEmitter.prototype.removeAllListeners.call(server, 'initialized')
93101
EventEmitter.prototype.removeAllListeners.call(server, 'error')
94102
delete this.server
103+
}
104+
105+
const cleanupOnSuccess = (): void => {
106+
isInitializing = false
107+
delete this.initReject
108+
if (this.initTimeout) {
109+
clearTimeout(this.initTimeout)
110+
delete this.initTimeout
111+
}
112+
EventEmitter.prototype.removeAllListeners.call(server, 'initialized')
113+
EventEmitter.prototype.removeAllListeners.call(server, 'error')
114+
}
115+
116+
this.initTimeout = setTimeout(() => {
117+
cleanupOnError()
95118
reject(new Error('TCP server initialization timeout after 10s'))
96119
}, 10000)
97120

98121
server.on('initialized', () => {
99-
clearTimeout(timeout)
100-
isInitializing = false
122+
cleanupOnSuccess()
101123
this.started = true
102124
resolve()
103125
})
104126

105127
// Handle errors - cast to EventEmitter as modbus-serial types don't include 'error' event
106128
;(server as unknown as EventEmitter).on('error', (err: Error) => {
107129
if (isInitializing) {
108-
clearTimeout(timeout)
109-
isInitializing = false
110-
EventEmitter.prototype.removeAllListeners.call(server, 'initialized')
111-
EventEmitter.prototype.removeAllListeners.call(server, 'error')
112-
delete this.server
130+
cleanupOnError()
113131
reject(err)
114132
} else {
115133
// Log errors that occur after initialization
@@ -120,14 +138,34 @@ export class TcpTransport extends BaseTransport {
120138
}
121139

122140
async stop(): Promise<void> {
123-
if (!this.started || !this.server) {
141+
// Clear initialization timeout if pending
142+
if (this.initTimeout) {
143+
clearTimeout(this.initTimeout)
144+
delete this.initTimeout
145+
}
146+
147+
// Clean up event listeners FIRST to prevent race condition (issue #294)
148+
// This ensures 'initialized' can't fire after we reject the Promise
149+
if (this.server) {
150+
EventEmitter.prototype.removeAllListeners.call(this.server, 'initialized')
151+
EventEmitter.prototype.removeAllListeners.call(this.server, 'error')
152+
}
153+
154+
// If initialization is pending, reject it
155+
if (this.initReject) {
156+
this.initReject(new Error('Transport stopped during initialization'))
157+
delete this.initReject
158+
}
159+
160+
if (!this.server) {
124161
return
125162
}
126163

127-
// Clean up event listeners to prevent memory leaks (issue #253)
128-
// Use EventEmitter.prototype since modbus-serial doesn't expose these methods in types
129-
EventEmitter.prototype.removeAllListeners.call(this.server, 'initialized')
130-
EventEmitter.prototype.removeAllListeners.call(this.server, 'error')
164+
// Only close the server if it was started, otherwise just clean up
165+
if (!this.started) {
166+
delete this.server
167+
return
168+
}
131169

132170
return new Promise<void>((resolve, reject) => {
133171
if (!this.server) {

0 commit comments

Comments
 (0)