fix(tsconnect-worker): preserve in-flight results across ownership transfer
Results of reads/accepts that were in flight when a resource was detached were delivered to the old owner's dead port, and accepted conns were registered on (and leaked with) the old client entry. Each resource now carries a generation-tagged channel; stale completions are deposited into a backlog the next owner's bridge drains first, serialized so stream order is preserved. TTL expiry and cancel close backlogged conns. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -228,13 +228,12 @@ if (!WASM_BUILT) {
|
||||
)
|
||||
const listener = await handle.listen("tcp", "0.0.0.0:0")
|
||||
const addr = listener.addr
|
||||
const pendingAccept = listener
|
||||
.accept()
|
||||
.then(() => "")
|
||||
.catch((e: Error) => e.message)
|
||||
await handle.transfer("popout", "move", listener)
|
||||
let deadAccept = ""
|
||||
try {
|
||||
await listener.accept()
|
||||
} catch (e) {
|
||||
deadAccept = (e as Error).message
|
||||
}
|
||||
const deadAccept = await pendingAccept
|
||||
let deadTransfer = ""
|
||||
try {
|
||||
await listener.transferState()
|
||||
|
||||
@@ -1,6 +1,12 @@
|
||||
import { test, suite } from "node:test"
|
||||
import assert from "node:assert/strict"
|
||||
import { createTransferRegistry, type AnyResource } from "./transfers.js"
|
||||
import {
|
||||
chanServe,
|
||||
createTransferRegistry,
|
||||
newChan,
|
||||
type AnyResource,
|
||||
type Chan,
|
||||
} from "./transfers.js"
|
||||
import type { TransferToken } from "./protocol.js"
|
||||
|
||||
// ── Helpers ───────────────────────────────────────────────────────────────────
|
||||
@@ -30,7 +36,7 @@ function fakeConn(): {
|
||||
function connResource(): { res: AnyResource; port1: MessagePort; obj: { closed: boolean } } {
|
||||
const { obj } = fakeConn()
|
||||
const { port1 } = new MessageChannel()
|
||||
const res = { kind: "conn", obj, port1 } as unknown as AnyResource
|
||||
const res = { kind: "conn", obj, port1, chan: newChan() } as unknown as AnyResource
|
||||
return { res, port1, obj }
|
||||
}
|
||||
|
||||
@@ -172,3 +178,150 @@ suite("onEmpty", () => {
|
||||
await waitFor(() => empties === 1)
|
||||
})
|
||||
})
|
||||
|
||||
// ── chanServe ─────────────────────────────────────────────────────────────────
|
||||
|
||||
suite("chanServe", () => {
|
||||
test("delivers the produced value while the generation is current", async () => {
|
||||
const chan = newChan<string>()
|
||||
const got: string[] = []
|
||||
chanServe(
|
||||
chan,
|
||||
async () => "a",
|
||||
(v) => got.push(v),
|
||||
() => assert.fail("unexpected error"),
|
||||
)
|
||||
await waitFor(() => got.length === 1)
|
||||
assert.deepEqual(got, ["a"])
|
||||
assert.equal(chan.backlog.length, 0)
|
||||
})
|
||||
|
||||
test("deposits an in-flight result into the backlog when detached mid-produce", async () => {
|
||||
const chan = newChan<string>()
|
||||
let resolve: ((v: string) => void) | undefined
|
||||
const got: string[] = []
|
||||
chanServe(
|
||||
chan,
|
||||
() => new Promise<string>((r) => (resolve = r)),
|
||||
(v) => got.push(v),
|
||||
() => assert.fail("unexpected error"),
|
||||
)
|
||||
await waitFor(() => resolve !== undefined)
|
||||
chan.gen++
|
||||
resolve!("late")
|
||||
await waitFor(() => chan.backlog.length === 1)
|
||||
assert.deepEqual(got, [])
|
||||
assert.deepEqual(chan.backlog, ["late"])
|
||||
})
|
||||
|
||||
test("serves the backlog to the next owner before producing again", async () => {
|
||||
const chan = newChan<string>()
|
||||
let resolve: ((v: string) => void) | undefined
|
||||
chanServe(
|
||||
chan,
|
||||
() => new Promise<string>((r) => (resolve = r)),
|
||||
() => assert.fail("stale delivery"),
|
||||
() => assert.fail("stale error"),
|
||||
)
|
||||
await waitFor(() => resolve !== undefined)
|
||||
chan.gen++
|
||||
|
||||
const got: string[] = []
|
||||
let produced = 0
|
||||
chanServe(
|
||||
chan,
|
||||
async () => {
|
||||
produced++
|
||||
return "fresh"
|
||||
},
|
||||
(v) => got.push(v),
|
||||
() => assert.fail("unexpected error"),
|
||||
)
|
||||
resolve!("crossed")
|
||||
await waitFor(() => got.length === 1)
|
||||
assert.deepEqual(got, ["crossed"])
|
||||
assert.equal(produced, 0)
|
||||
})
|
||||
|
||||
test("drops errors that land after detach", async () => {
|
||||
const chan = newChan<string>()
|
||||
let reject: ((e: Error) => void) | undefined
|
||||
let failed = false
|
||||
let settled = false
|
||||
chanServe(
|
||||
chan,
|
||||
() => new Promise<string>((_, r) => (reject = r)),
|
||||
() => assert.fail("unexpected delivery"),
|
||||
() => (failed = true),
|
||||
)
|
||||
await waitFor(() => reject !== undefined)
|
||||
chan.gen++
|
||||
reject!(new Error("late failure"))
|
||||
chan.lock = chan.lock.then(() => {
|
||||
settled = true
|
||||
})
|
||||
await waitFor(() => settled)
|
||||
assert.equal(failed, false)
|
||||
})
|
||||
})
|
||||
|
||||
// ── backlog lifecycle ─────────────────────────────────────────────────────────
|
||||
|
||||
suite("backlog lifecycle", () => {
|
||||
function listenerResource(): {
|
||||
res: AnyResource
|
||||
obj: { closed: boolean }
|
||||
chan: Chan<{ closed: boolean; close(): void }>
|
||||
} {
|
||||
const obj = {
|
||||
addr: "127.0.0.1:80",
|
||||
closed: false,
|
||||
close() {
|
||||
this.closed = true
|
||||
},
|
||||
}
|
||||
const chan = newChan<{ closed: boolean; close(): void }>()
|
||||
const { port1 } = new MessageChannel()
|
||||
const res = { kind: "listener", obj, port1, chan } as unknown as AnyResource
|
||||
return { res, obj, chan }
|
||||
}
|
||||
|
||||
test("detach bumps the channel generation", () => {
|
||||
const registry = createTransferRegistry()
|
||||
const { res, chan } = listenerResource()
|
||||
registry.detach(entryWith(res), 0)
|
||||
assert.equal(chan.gen, 1)
|
||||
})
|
||||
|
||||
test("TTL expiry closes backlogged accepted conns", async () => {
|
||||
const registry = createTransferRegistry({ ttl: 10 })
|
||||
const { res, obj, chan } = listenerResource()
|
||||
const stray = {
|
||||
closed: false,
|
||||
close() {
|
||||
this.closed = true
|
||||
},
|
||||
}
|
||||
registry.detach(entryWith(res), 0)
|
||||
chan.backlog.push(stray)
|
||||
await waitFor(() => obj.closed)
|
||||
assert.equal(stray.closed, true)
|
||||
})
|
||||
|
||||
test("claim hands back the resource with its backlog intact", () => {
|
||||
const registry = createTransferRegistry()
|
||||
const { res, chan } = listenerResource()
|
||||
const stray = {
|
||||
closed: false,
|
||||
close() {
|
||||
this.closed = true
|
||||
},
|
||||
}
|
||||
const { transferId } = registry.detach(entryWith(res), 0)
|
||||
chan.backlog.push(stray)
|
||||
const claimed = registry.claim(transferId)
|
||||
assert.equal(claimed.chan, chan)
|
||||
assert.equal(claimed.chan.backlog.length, 1)
|
||||
assert.equal(stray.closed, false)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,10 +1,49 @@
|
||||
import type { Conn, TCPListener, PacketConn } from "@webnet/tsconnect"
|
||||
import type { ResourceMeta, TransferToken, W2C } from "./protocol.js"
|
||||
|
||||
// Persists across bridge re-creation on transfer: results of operations that
|
||||
// were in flight when the resource was detached (gen bumped) are deposited in
|
||||
// the backlog and served to the next owner's bridge before touching the
|
||||
// underlying resource again. The lock serializes produce calls across bridge
|
||||
// generations so stream ordering is preserved.
|
||||
export type Chan<T> = { gen: number; backlog: T[]; lock: Promise<void> }
|
||||
|
||||
export function newChan<T>(): Chan<T> {
|
||||
return { gen: 0, backlog: [], lock: Promise.resolve() }
|
||||
}
|
||||
|
||||
export function chanServe<T>(
|
||||
chan: Chan<T>,
|
||||
produce: () => Promise<T>,
|
||||
deliver: (value: T) => void,
|
||||
fail: (err: unknown) => void,
|
||||
): void {
|
||||
const gen = chan.gen
|
||||
chan.lock = chan.lock.then(async () => {
|
||||
const buffered = chan.backlog.shift()
|
||||
if (buffered !== undefined) {
|
||||
deliver(buffered)
|
||||
return
|
||||
}
|
||||
try {
|
||||
const value = await produce()
|
||||
if (chan.gen !== gen) chan.backlog.push(value)
|
||||
else deliver(value)
|
||||
} catch (err) {
|
||||
if (chan.gen === gen) fail(err)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
export type AnyResource =
|
||||
| { kind: "conn"; obj: Conn; port1: MessagePort }
|
||||
| { kind: "listener"; obj: TCPListener; port1: MessagePort }
|
||||
| { kind: "packetconn"; obj: PacketConn; port1: MessagePort }
|
||||
| { kind: "conn"; obj: Conn; port1: MessagePort; chan: Chan<Uint8Array> }
|
||||
| { kind: "listener"; obj: TCPListener; port1: MessagePort; chan: Chan<Conn> }
|
||||
| {
|
||||
kind: "packetconn"
|
||||
obj: PacketConn
|
||||
port1: MessagePort
|
||||
chan: Chan<{ data: Uint8Array; addr: string }>
|
||||
}
|
||||
|
||||
type TransferMessage = Extract<W2C, { type: "transfer" }>
|
||||
|
||||
@@ -29,6 +68,24 @@ export type Flush = (msg: TransferMessage, transferables: Transferable[]) => voi
|
||||
|
||||
export type TransferRegistry = ReturnType<typeof createTransferRegistry>
|
||||
|
||||
function closeResource(res: AnyResource): void {
|
||||
try {
|
||||
res.obj.close()
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
if (res.kind === "listener") {
|
||||
for (const conn of res.chan.backlog) {
|
||||
try {
|
||||
conn.close()
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
}
|
||||
res.chan.backlog.length = 0
|
||||
}
|
||||
}
|
||||
|
||||
function metaOf(res: AnyResource): ResourceMeta {
|
||||
if (res.kind === "conn")
|
||||
return { kind: "conn", localAddr: res.obj.localAddr, remoteAddr: res.obj.remoteAddr }
|
||||
@@ -51,11 +108,7 @@ export function createTransferRegistry(opts: { ttl?: number; onEmpty?: () => voi
|
||||
if (!pending) return
|
||||
clearTimeout(pending.timer)
|
||||
pendingTransfers.delete(transferId)
|
||||
try {
|
||||
pending.resource.obj.close()
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
closeResource(pending.resource)
|
||||
}
|
||||
|
||||
function detach(
|
||||
@@ -67,15 +120,12 @@ export function createTransferRegistry(opts: { ttl?: number; onEmpty?: () => voi
|
||||
entry.resources.delete(resourceId)
|
||||
res.port1.onmessage = null
|
||||
res.port1.close()
|
||||
res.chan.gen++
|
||||
|
||||
const transferId = crypto.randomUUID()
|
||||
const timer = setTimeout(() => {
|
||||
pendingTransfers.delete(transferId)
|
||||
try {
|
||||
res.obj.close()
|
||||
} catch {
|
||||
/* already closed */
|
||||
}
|
||||
closeResource(res)
|
||||
checkEmpty()
|
||||
}, ttl)
|
||||
pendingTransfers.set(transferId, { resource: res, timer })
|
||||
|
||||
@@ -25,8 +25,8 @@ import type {
|
||||
} from "./protocol.js"
|
||||
import { pumpStreamToPort, portToReadableStream } from "./protocol.js"
|
||||
import type { TransferToken } from "./protocol.js"
|
||||
import { createTransferRegistry } from "./transfers.js"
|
||||
import type { AnyResource } from "./transfers.js"
|
||||
import { chanServe, createTransferRegistry, newChan } from "./transfers.js"
|
||||
import type { AnyResource, Chan } from "./transfers.js"
|
||||
|
||||
const sw = self as unknown as SharedWorkerGlobalScope
|
||||
|
||||
@@ -84,7 +84,11 @@ function wrapDispatch(s: ReturnType<typeof buildIpnStore>): void {
|
||||
|
||||
// ── Worker-side Conn bridge ──────────────────────────────────────────────────
|
||||
|
||||
function bridgeConn(conn: Conn, onClose?: () => void): { port1: MessagePort; port2: MessagePort } {
|
||||
function bridgeConn(
|
||||
conn: Conn,
|
||||
chan: Chan<Uint8Array>,
|
||||
onClose?: () => void,
|
||||
): { port1: MessagePort; port2: MessagePort } {
|
||||
const ch = new MessageChannel()
|
||||
const { port1, port2 } = ch
|
||||
port1.onmessage = async (e: MessageEvent) => {
|
||||
@@ -97,13 +101,17 @@ function bridgeConn(conn: Conn, onClose?: () => void): { port1: MessagePort; por
|
||||
return
|
||||
}
|
||||
if (msg.type === "read") {
|
||||
try {
|
||||
const data = await conn.read()
|
||||
const reply: ConnW2C = { type: "data", id: msg.id, data }
|
||||
port1.postMessage(reply, [data.buffer])
|
||||
} catch (err) {
|
||||
port1.postMessage({ type: "readError", id: msg.id, error: errMsg(err) } satisfies ConnW2C)
|
||||
}
|
||||
chanServe(
|
||||
chan,
|
||||
() => conn.read(),
|
||||
(data) => {
|
||||
const reply: ConnW2C = { type: "data", id: msg.id, data }
|
||||
port1.postMessage(reply, [data.buffer])
|
||||
},
|
||||
(err) => {
|
||||
port1.postMessage({ type: "readError", id: msg.id, error: errMsg(err) } satisfies ConnW2C)
|
||||
},
|
||||
)
|
||||
return
|
||||
}
|
||||
if (msg.type === "write") {
|
||||
@@ -124,6 +132,7 @@ function bridgeConn(conn: Conn, onClose?: () => void): { port1: MessagePort; por
|
||||
function bridgeListener(
|
||||
listener: TCPListener,
|
||||
entry: ClientEntry,
|
||||
chan: Chan<Conn>,
|
||||
onClose?: () => void,
|
||||
): { port1: MessagePort; port2: MessagePort } {
|
||||
const ch = new MessageChannel()
|
||||
@@ -138,25 +147,29 @@ function bridgeListener(
|
||||
return
|
||||
}
|
||||
if (msg.type === "accept") {
|
||||
try {
|
||||
const conn = await listener.accept()
|
||||
const { port2: connPort, resourceId } = registerConnResource(entry, conn)
|
||||
const reply: ListenerW2C = {
|
||||
type: "accepted",
|
||||
id: msg.id,
|
||||
localAddr: conn.localAddr,
|
||||
remoteAddr: conn.remoteAddr,
|
||||
port: connPort,
|
||||
resourceId,
|
||||
}
|
||||
port1.postMessage(reply, [connPort])
|
||||
} catch (err) {
|
||||
port1.postMessage({
|
||||
type: "acceptError",
|
||||
id: msg.id,
|
||||
error: errMsg(err),
|
||||
} satisfies ListenerW2C)
|
||||
}
|
||||
chanServe(
|
||||
chan,
|
||||
() => listener.accept(),
|
||||
(conn) => {
|
||||
const { port2: connPort, resourceId } = registerConnResource(entry, conn)
|
||||
const reply: ListenerW2C = {
|
||||
type: "accepted",
|
||||
id: msg.id,
|
||||
localAddr: conn.localAddr,
|
||||
remoteAddr: conn.remoteAddr,
|
||||
port: connPort,
|
||||
resourceId,
|
||||
}
|
||||
port1.postMessage(reply, [connPort])
|
||||
},
|
||||
(err) => {
|
||||
port1.postMessage({
|
||||
type: "acceptError",
|
||||
id: msg.id,
|
||||
error: errMsg(err),
|
||||
} satisfies ListenerW2C)
|
||||
},
|
||||
)
|
||||
}
|
||||
}
|
||||
port1.start()
|
||||
@@ -167,6 +180,7 @@ function bridgeListener(
|
||||
|
||||
function bridgePacketConn(
|
||||
pc: PacketConn,
|
||||
chan: Chan<{ data: Uint8Array; addr: string }>,
|
||||
onClose?: () => void,
|
||||
): { port1: MessagePort; port2: MessagePort } {
|
||||
const ch = new MessageChannel()
|
||||
@@ -181,17 +195,21 @@ function bridgePacketConn(
|
||||
return
|
||||
}
|
||||
if (msg.type === "readFrom") {
|
||||
try {
|
||||
const { data, addr } = await pc.readFrom()
|
||||
const reply: PacketConnW2C = { type: "packet", id: msg.id, data, addr }
|
||||
port1.postMessage(reply, [data.buffer])
|
||||
} catch (err) {
|
||||
port1.postMessage({
|
||||
type: "readError",
|
||||
id: msg.id,
|
||||
error: errMsg(err),
|
||||
} satisfies PacketConnW2C)
|
||||
}
|
||||
chanServe(
|
||||
chan,
|
||||
() => pc.readFrom(),
|
||||
({ data, addr }) => {
|
||||
const reply: PacketConnW2C = { type: "packet", id: msg.id, data, addr }
|
||||
port1.postMessage(reply, [data.buffer])
|
||||
},
|
||||
(err) => {
|
||||
port1.postMessage({
|
||||
type: "readError",
|
||||
id: msg.id,
|
||||
error: errMsg(err),
|
||||
} satisfies PacketConnW2C)
|
||||
},
|
||||
)
|
||||
return
|
||||
}
|
||||
if (msg.type === "writeTo") {
|
||||
@@ -216,30 +234,35 @@ function bridgePacketConn(
|
||||
function registerConnResource(
|
||||
entry: ClientEntry,
|
||||
conn: Conn,
|
||||
chan: Chan<Uint8Array> = newChan(),
|
||||
): { port2: MessagePort; resourceId: number } {
|
||||
const resourceId = nextResourceId++
|
||||
const { port1, port2 } = bridgeConn(conn, () => entry.resources.delete(resourceId))
|
||||
entry.resources.set(resourceId, { kind: "conn", obj: conn, port1 })
|
||||
const { port1, port2 } = bridgeConn(conn, chan, () => entry.resources.delete(resourceId))
|
||||
entry.resources.set(resourceId, { kind: "conn", obj: conn, port1, chan })
|
||||
return { port2, resourceId }
|
||||
}
|
||||
|
||||
function registerListenerResource(
|
||||
entry: ClientEntry,
|
||||
listener: TCPListener,
|
||||
chan: Chan<Conn> = newChan(),
|
||||
): { port2: MessagePort; resourceId: number } {
|
||||
const resourceId = nextResourceId++
|
||||
const { port1, port2 } = bridgeListener(listener, entry, () => entry.resources.delete(resourceId))
|
||||
entry.resources.set(resourceId, { kind: "listener", obj: listener, port1 })
|
||||
const { port1, port2 } = bridgeListener(listener, entry, chan, () =>
|
||||
entry.resources.delete(resourceId),
|
||||
)
|
||||
entry.resources.set(resourceId, { kind: "listener", obj: listener, port1, chan })
|
||||
return { port2, resourceId }
|
||||
}
|
||||
|
||||
function registerPacketConnResource(
|
||||
entry: ClientEntry,
|
||||
pc: PacketConn,
|
||||
chan: Chan<{ data: Uint8Array; addr: string }> = newChan(),
|
||||
): { port2: MessagePort; resourceId: number } {
|
||||
const resourceId = nextResourceId++
|
||||
const { port1, port2 } = bridgePacketConn(pc, () => entry.resources.delete(resourceId))
|
||||
entry.resources.set(resourceId, { kind: "packetconn", obj: pc, port1 })
|
||||
const { port1, port2 } = bridgePacketConn(pc, chan, () => entry.resources.delete(resourceId))
|
||||
entry.resources.set(resourceId, { kind: "packetconn", obj: pc, port1, chan })
|
||||
return { port2, resourceId }
|
||||
}
|
||||
|
||||
@@ -602,7 +625,7 @@ async function handleCall(
|
||||
case "claimResource": {
|
||||
const resource = registry.claim(args[0] as string)
|
||||
if (resource.kind === "conn") {
|
||||
const { port2, resourceId } = registerConnResource(entry, resource.obj)
|
||||
const { port2, resourceId } = registerConnResource(entry, resource.obj, resource.chan)
|
||||
const msg: W2C = {
|
||||
type: "conn",
|
||||
id,
|
||||
@@ -613,7 +636,7 @@ async function handleCall(
|
||||
}
|
||||
send(port, msg, [port2])
|
||||
} else if (resource.kind === "listener") {
|
||||
const { port2, resourceId } = registerListenerResource(entry, resource.obj)
|
||||
const { port2, resourceId } = registerListenerResource(entry, resource.obj, resource.chan)
|
||||
const msg: W2C = {
|
||||
type: "listener",
|
||||
id,
|
||||
@@ -623,7 +646,11 @@ async function handleCall(
|
||||
}
|
||||
send(port, msg, [port2])
|
||||
} else {
|
||||
const { port2, resourceId } = registerPacketConnResource(entry, resource.obj)
|
||||
const { port2, resourceId } = registerPacketConnResource(
|
||||
entry,
|
||||
resource.obj,
|
||||
resource.chan,
|
||||
)
|
||||
const msg: W2C = {
|
||||
type: "packetconn",
|
||||
id,
|
||||
|
||||
Reference in New Issue
Block a user