Skip to content

Commit f9216c5

Browse files
authored
refactor: abort signal (#46)
1 parent 6fee25a commit f9216c5

20 files changed

Lines changed: 252 additions & 187 deletions

File tree

src/adapters/broadcast-channel/index.ts

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,10 @@ export function createContext(channel: BroadcastChannel, options?: BroadcastChan
7171

7272
return {
7373
context: ctx,
74-
dispose: () => {
74+
dispose: (reason?: unknown) => {
75+
// Cascade-cancel any in-flight `defineInvoke(...)` so callers don't hang
76+
// on a torn-down channel (especially when `closeOnDispose: true`).
77+
ctx.abort(reason ?? new Error('eventa: invoke cancelled, BroadcastChannel disposed'))
7578
cleanupRemoval.forEach(removal => removal.remove())
7679
if (closeOnDispose) {
7780
channel.close?.()

src/adapters/electron/main.ts

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -103,7 +103,10 @@ export function createContext(ipcMain: IpcMain, window?: BrowserWindow, options?
103103

104104
return {
105105
context: ctx,
106-
dispose: () => {
106+
dispose: (reason?: unknown) => {
107+
// Cascade-cancel any in-flight `defineInvoke(...)` so renderer-bound
108+
// RPCs don't hang after the main-side adapter is torn down.
109+
ctx.abort(reason ?? new Error('eventa: invoke cancelled, electron main ipc disposed'))
107110
cleanupRemoval.forEach(removal => removal.remove())
108111
},
109112
}

src/adapters/electron/renderer.ts

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,10 @@ export function createContext(ipcRenderer: IpcRenderer, options?: {
6565

6666
return {
6767
context: ctx,
68-
dispose: () => {
68+
dispose: (reason?: unknown) => {
69+
// Cascade-cancel any in-flight `defineInvoke(...)` so main-bound
70+
// RPCs don't hang after the renderer-side adapter is torn down.
71+
ctx.abort(reason ?? new Error('eventa: invoke cancelled, electron renderer ipc disposed'))
6972
cleanupRemoval.forEach(removal => removal.remove())
7073
},
7174
}

src/adapters/event-emitter/index.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -64,7 +64,8 @@ export function createContext(eventTarget: NodeJS.EventEmitter, options?: {
6464

6565
return {
6666
context: ctx,
67-
dispose: () => {
67+
dispose: (reason?: unknown) => {
68+
ctx.abort(reason ?? new Error('eventa: invoke cancelled, EventEmitter adapter disposed'))
6869
cleanupRemoval.forEach(removal => removal.remove())
6970
},
7071
}

src/adapters/event-target/index.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,8 @@ export function createContext(eventTarget: EventTarget, options?: {
7171

7272
return {
7373
context: ctx,
74-
dispose: () => {
74+
dispose: (reason?: unknown) => {
75+
ctx.abort(reason ?? new Error('eventa: invoke cancelled, EventTarget adapter disposed'))
7576
cleanupRemoval.forEach(removal => removal.remove())
7677
},
7778
}

src/adapters/websocket/h3/peer.test.ts

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -154,4 +154,60 @@ describe('h3 websocket adapter', { timeout: 2000 }, async () => {
154154
expect(res.output).toEqual('500')
155155
}
156156
})
157+
158+
// ROOT CAUSE:
159+
//
160+
// Before this fix, createPeerContext only wired the `message` hook, leaving
161+
// `close` and `error` to consumers. A server-side defineInvoke() back to a
162+
// client would hang forever when that peer disconnected, because nothing
163+
// aborted the peer's ctx.signal. This mismatched the native ws adapter's
164+
// self-aborting contract.
165+
//
166+
// We fixed this by wiring the `close` / `error` hooks directly in
167+
// createPeerContext to call ctx.abort(...) in addition to emitting the
168+
// business event — same shape as native ws onclose/onerror.
169+
it('issue: rejects pending server-side invoke when peer disconnects', async (testCtx) => {
170+
const port = randomBetween(40000, 50000)
171+
const app = new H3()
172+
173+
const { untilLeastOneConnected, hooks } = createPeerHooks()
174+
app.get('/ws', defineWebSocketHandler(hooks))
175+
176+
{
177+
const server = serve(app, {
178+
port,
179+
plugins: [ws({
180+
resolve: async (req) => {
181+
const response = (await app.fetch(req)) as Response & { crossws: Partial<Hooks> }
182+
return response.crossws
183+
},
184+
})],
185+
})
186+
187+
testCtx.onTestFinished(() => {
188+
server.close()
189+
})
190+
}
191+
192+
const opened = createUntil<void>()
193+
const wsConn = new WebSocket(`ws://localhost:${port}/ws`)
194+
wsConn.onopen = () => opened.handler()
195+
await opened.promise
196+
197+
// Intentionally do NOT register a handler on the client side — the server
198+
// will issue an invoke that has no responder, so the only path out of
199+
// the pending promise is the disconnect-driven abort.
200+
createContext(wsConn)
201+
const { context: serverPeerContext } = await untilLeastOneConnected
202+
203+
const events = defineInvokeEventa<string, string>('test:peer-abort-cascade')
204+
const invoke = defineInvoke(serverPeerContext, events)
205+
const pending = invoke('hello')
206+
207+
// Drop the client side; crossws will fire the `close` hook on the server
208+
// peer; createPeerContext.hooks.close calls ctx.abort(...).
209+
wsConn.close()
210+
211+
await expect(pending).rejects.toThrowError(/peer disconnected/i)
212+
})
157213
})

src/adapters/websocket/h3/peer.ts

Lines changed: 36 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,6 @@ import type { EventContext } from '../../../context'
44
import type { DirectionalEventa, Eventa } from '../../../eventa'
55

66
import { createContext as createBaseContext } from '../../../context'
7-
import { registerInvokeAbortEventListeners } from '../../../context-extension-invoke-internal'
87
import { and, defineEventa, defineInboundEventa, defineOutboundEventa, EventaFlowDirection, matchBy } from '../../../eventa'
98
import { generateWebsocketPayload, parseWebsocketPayload } from '../internal'
109

@@ -13,28 +12,12 @@ export const wsDisconnectedEvent = defineEventa<{ id: string }>('eventa:adapters
1312
export const wsErrorEvent = defineEventa<{ error: unknown }>('eventa:adapters:websocket-peer:error')
1413

1514
export function createPeerContext(peer: Peer): {
16-
hooks: Pick<Hooks, 'message'>
15+
hooks: Pick<Hooks, 'message' | 'close' | 'error'>
1716
context: EventContext<any, { raw: { message: Message } }>
1817
} {
1918
const peerId = peer.id
2019
const ctx = createBaseContext<any, { raw: { message: Message } }>()
2120

22-
// Reject any in-flight `defineInvoke(...)` promises if this peer's transport
23-
// dies. Mirrors the native ws adapter so server-side code that issues an
24-
// invoke back to a client (push-style RPC) doesn't hang on close.
25-
registerInvokeAbortEventListeners(ctx, wsDisconnectedEvent, (payload) => {
26-
if (payload.id === wsDisconnectedEvent.id) {
27-
const id = (payload as Eventa<{ id?: string }>).body?.id
28-
return new Error(`eventa: invoke cancelled, peer disconnected${id ? ` (${id})` : ''}`)
29-
}
30-
if (payload.id === wsErrorEvent.id) {
31-
const err = (payload as Eventa<{ error?: unknown }>).body?.error
32-
return err instanceof Error ? err : new Error('eventa: invoke cancelled, peer error')
33-
}
34-
return undefined
35-
})
36-
registerInvokeAbortEventListeners(ctx, wsErrorEvent)
37-
3821
ctx.on(and(
3922
matchBy((e: DirectionalEventa<any>) => e._flowDirection === EventaFlowDirection.Outbound || !e._flowDirection),
4023
matchBy('*'),
@@ -52,11 +35,30 @@ export function createPeerContext(peer: Peer): {
5235
ctx.emit(defineInboundEventa(type), payload.body, { raw: { message } })
5336
}
5437
catch (error) {
38+
// Per-message parse failure — recoverable, do NOT abort lifetime.
5539
console.error('Failed to parse WebSocket message:', error)
5640
ctx.emit(wsErrorEvent, { error }, { raw: { message } })
5741
}
5842
}
5943
},
44+
close(peer, details) {
45+
// crossws fires close for ANY peer; filter to our own.
46+
if (peer.id !== peerId) {
47+
return
48+
}
49+
const reasonText = details.reason ? ` (${details.reason})` : ''
50+
// Cascade-cancel any in-flight `defineInvoke(...)` so server-side code
51+
// that issued an invoke back to this peer doesn't hang on close.
52+
ctx.abort(new Error(`eventa: invoke cancelled, peer disconnected${reasonText}`))
53+
ctx.emit(wsDisconnectedEvent, { id: peerId })
54+
},
55+
error(peer, error) {
56+
if (peer.id !== peerId) {
57+
return
58+
}
59+
ctx.abort(error instanceof Error ? error : new Error('eventa: invoke cancelled, peer error'))
60+
ctx.emit(wsErrorEvent, { error })
61+
},
6062
},
6163
context: ctx,
6264
}
@@ -70,18 +72,30 @@ export function createPeerHooks(): { hooks: Partial<Hooks>, untilLeastOneConnect
7072
resolve = r
7173
})
7274

75+
// NOTICE: single-peer model — these closure-scoped hook refs get overwritten
76+
// when a second peer connects, so `createPeerHooks` only correctly serves the
77+
// most-recently-opened peer. Multi-peer support requires a peerId-keyed Map
78+
// and a different "untilLeastOneConnected" semantic; out of scope here.
7379
let message: Hooks['message'] | undefined
80+
let close: Hooks['close'] | undefined
81+
let error: Hooks['error'] | undefined
7482

75-
const hooks: Pick<Hooks, 'open' | 'message'> = {
83+
const hooks: Pick<Hooks, 'open' | 'message' | 'close' | 'error'> = {
7684
open: (peer) => {
7785
const { context, hooks } = createPeerContext(peer)
7886
message = hooks.message
87+
close = hooks.close
88+
error = hooks.error
7989
resolve({ peer, context })
8090
},
8191
message: (peer, msg) => {
82-
if (message != null) {
83-
message(peer, msg)
84-
}
92+
message?.(peer, msg)
93+
},
94+
close: (peer, details) => {
95+
close?.(peer, details)
96+
},
97+
error: (peer, err) => {
98+
error?.(peer, err)
8599
},
86100
}
87101

src/adapters/websocket/native/index.spec.ts

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -164,10 +164,11 @@ describe('browser websocket adapter', () => {
164164
// had to maintain its own pending-RPC tracker and reject manually on
165165
// disconnect.
166166
//
167-
// We fixed this by registering wsDisconnectedEvent and wsErrorEvent as abort
168-
// events on the context inside `createContext`, with a mapAbortError that
169-
// produces real Error instances. defineInvoke's existing abortOnEvents
170-
// machinery then rejects every in-flight invoke when either event fires.
167+
// We fixed this by giving every EventContext a lifetime AbortSignal
168+
// (`ctx.signal`) and an `abort(reason)` method. The native ws adapter calls
169+
// `ctx.abort(error)` from `onclose` / `onerror`; defineInvoke hooks
170+
// `ctx.signal` so transport death cascades into a synchronous reject of
171+
// every in-flight invoke. Modeled after Go's context.Context.
171172
it('issue: rejects pending invoke when socket closes mid-flight', async (testCtx) => {
172173
const port = randomBetween(40000, 50000)
173174
const app = new H3()
@@ -212,8 +213,8 @@ describe('browser websocket adapter', () => {
212213
const invocation = invoke('hello')
213214

214215
// Drop the underlying socket without ever delivering a response. The
215-
// adapter emits wsDisconnectedEvent, defineInvoke sees it via abortOnEvents,
216-
// and rejects with the mapAbortError-produced Error.
216+
// adapter calls ctx.abort(error) from onclose; defineInvoke is hooked to
217+
// ctx.signal and rejects with the adapter-supplied Error.
217218
wsConn.close()
218219

219220
await expect(invocation).rejects.toThrowError(/websocket disconnected/i)

src/adapters/websocket/native/index.ts

Lines changed: 5 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@ import type { EventContext } from '../../../context'
22
import type { DirectionalEventa, Eventa } from '../../../eventa'
33

44
import { createContext as createBaseContext } from '../../../context'
5-
import { registerInvokeAbortEventListeners } from '../../../context-extension-invoke-internal'
65
import { and, defineEventa, defineInboundEventa, defineOutboundEventa, EventaFlowDirection, matchBy } from '../../../eventa'
76
import { generateWebsocketPayload, parseWebsocketPayload } from '../internal'
87

@@ -13,22 +12,6 @@ export const wsErrorEvent = defineEventa<{ error: unknown }>()
1312
export function createContext(wsConn: WebSocket) {
1413
const ctx = createBaseContext() as EventContext<any, { raw: { message?: any, open?: Event, error?: Event, close?: CloseEvent } }>
1514

16-
// Reject any in-flight `defineInvoke(...)` promises when the socket dies.
17-
// Without this, callers wait forever for a response that the closed transport
18-
// can never deliver
19-
registerInvokeAbortEventListeners(ctx, wsDisconnectedEvent, (payload) => {
20-
if (payload.id === wsDisconnectedEvent.id) {
21-
const url = (payload as Eventa<{ url?: string }>).body?.url
22-
return new Error(`eventa: invoke cancelled, websocket disconnected${url ? ` (${url})` : ''}`)
23-
}
24-
if (payload.id === wsErrorEvent.id) {
25-
const err = (payload as Eventa<{ error?: unknown }>).body?.error
26-
return err instanceof Error ? err : new Error('eventa: invoke cancelled, websocket error')
27-
}
28-
return undefined
29-
})
30-
registerInvokeAbortEventListeners(ctx, wsErrorEvent)
31-
3215
ctx.on(and(
3316
matchBy((e: DirectionalEventa<any>) => e._flowDirection === EventaFlowDirection.Outbound || !e._flowDirection),
3417
matchBy('*'),
@@ -53,10 +36,15 @@ export function createContext(wsConn: WebSocket) {
5336
}
5437

5538
wsConn.onerror = (error) => {
39+
// Socket-level error (not a per-message parse failure — those stay
40+
// recoverable in `onmessage` above). Abort lifetime so any in-flight
41+
// invoke rejects; emit the business event for non-invoke listeners.
42+
ctx.abort(new Error('eventa: invoke cancelled, websocket error'))
5643
ctx.emit(wsErrorEvent, { error }, { raw: { error } })
5744
}
5845

5946
wsConn.onclose = (close) => {
47+
ctx.abort(new Error(`eventa: invoke cancelled, websocket disconnected (${wsConn.url})`))
6048
ctx.emit(wsDisconnectedEvent, { url: wsConn.url }, { raw: { close } })
6149
}
6250

src/adapters/webworkers/index.ts

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@ import type { EventContext } from '../../context'
22
import type { DirectionalEventa, Eventa } from '../../eventa'
33

44
import { createContext as createBaseContext } from '../../context'
5-
import { registerInvokeAbortEventListeners } from '../../context-extension-invoke-internal'
65
import { and, defineInboundEventa, defineOutboundEventa, EventaFlowDirection, matchBy } from '../../eventa'
76
import { generateWorkerPayload, parseWorkerPayload } from './internal'
87
import { isWorkerEventa, normalizeOnListenerParameters, workerErrorEvent } from './shared'
@@ -15,8 +14,6 @@ export function createContext(worker: Worker) {
1514
},
1615
{ raw: { message?: MessageEvent, error?: ErrorEvent, messageError?: MessageEvent }, transfer?: Transferable[] }
1716
>
18-
// Configure invoke to fail fast on fatal worker errors (load/syntax/runtime).
19-
registerInvokeAbortEventListeners(ctx, workerErrorEvent)
2017

2118
ctx.on(and(
2219
matchBy((e: DirectionalEventa<any>) => e._flowDirection === EventaFlowDirection.Outbound || !e._flowDirection),
@@ -49,10 +46,14 @@ export function createContext(worker: Worker) {
4946
}
5047

5148
worker.onerror = (error) => {
49+
// Fatal worker error (load / syntax / runtime). Abort lifetime so any
50+
// in-flight invoke rejects; emit the business event for non-invoke listeners.
51+
ctx.abort(error instanceof Error ? error : new Error('eventa: invoke cancelled, webworker error'))
5252
ctx.emit(workerErrorEvent, { error }, { raw: { error } })
5353
}
5454

5555
worker.onmessageerror = (error) => {
56+
ctx.abort(new Error('eventa: invoke cancelled, webworker messageerror'))
5657
ctx.emit(workerErrorEvent, { error }, { raw: { messageError: error } })
5758
}
5859

0 commit comments

Comments
 (0)