mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 00:02:10 +00:00
test(ssh,relay): pin the real bounds of the closed-generation list and the port-scan exit
Corrects the generation-leak claim: the ranges list is bounded by the number of concurrently live generations, not flattened unconditionally, and a generation that is never closed leaves a permanent gap. Makes the structure log-time so that degraded shape is never worse than the Set it replaced. Extends the port-scan fixture to several listeners so the exit condition is distinguishable from "result.size > 0", and pins the shared-inode attribution change the early exit introduces.
This commit is contained in:
@@ -0,0 +1,65 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { SshPtyClosedGenerationRanges } from './ssh-pty-closed-generation-ranges'
|
||||
|
||||
// Deterministic so a membership divergence reproduces from the failure output alone.
|
||||
function lcg(seed: number): () => number {
|
||||
let state = seed >>> 0
|
||||
return () => {
|
||||
state = (state * 1_664_525 + 1_013_904_223) >>> 0
|
||||
return state / 0x1_0000_0000
|
||||
}
|
||||
}
|
||||
|
||||
describe('SshPtyClosedGenerationRanges', () => {
|
||||
it('answers membership identically to a plain set under arbitrary insertion order', () => {
|
||||
const random = lcg(0xc0ffee)
|
||||
const ranges = new SshPtyClosedGenerationRanges()
|
||||
const oracle = new Set<number>()
|
||||
for (let step = 0; step < 4_000; step += 1) {
|
||||
const generation = 1 + Math.floor(random() * 400)
|
||||
ranges.add(generation)
|
||||
oracle.add(generation)
|
||||
}
|
||||
for (let generation = 0; generation <= 402; generation += 1) {
|
||||
expect([generation, ranges.has(generation)]).toEqual([generation, oracle.has(generation)])
|
||||
}
|
||||
})
|
||||
|
||||
it('merges an inserted generation that bridges two ranges', () => {
|
||||
const ranges = new SshPtyClosedGenerationRanges()
|
||||
ranges.add(1)
|
||||
ranges.add(3)
|
||||
expect(ranges.size).toBe(2)
|
||||
|
||||
ranges.add(2)
|
||||
|
||||
expect(ranges.size).toBe(1)
|
||||
expect([1, 2, 3].map((generation) => ranges.has(generation))).toEqual([true, true, true])
|
||||
expect(ranges.has(4)).toBe(false)
|
||||
})
|
||||
|
||||
it('treats a repeated close as a no-op', () => {
|
||||
const ranges = new SshPtyClosedGenerationRanges()
|
||||
ranges.add(7)
|
||||
ranges.add(7)
|
||||
expect(ranges.size).toBe(1)
|
||||
expect(ranges.activeGaps).toBe(6)
|
||||
})
|
||||
|
||||
it('keeps insertion sublinear so a fragmented list cannot become quadratic', () => {
|
||||
// The scan this replaced took ~220ms to insert 20k non-adjacent generations and ~1300ms for
|
||||
// 100k membership probes against them; both are microseconds once the lookup is log-time.
|
||||
const ranges = new SshPtyClosedGenerationRanges()
|
||||
const startedAt = performance.now()
|
||||
for (let generation = 2; generation <= 80_000; generation += 2) {
|
||||
ranges.add(generation)
|
||||
}
|
||||
for (let probe = 0; probe < 100_000; probe += 1) {
|
||||
ranges.has(80_001)
|
||||
}
|
||||
const elapsed = performance.now() - startedAt
|
||||
|
||||
expect(ranges.size).toBe(40_000)
|
||||
expect(elapsed).toBeLessThan(1_000)
|
||||
})
|
||||
})
|
||||
@@ -7,10 +7,7 @@ export class SshPtyClosedGenerationRanges {
|
||||
private readonly ranges: ClosedGenerationRange[] = []
|
||||
|
||||
add(generation: number): void {
|
||||
let index = 0
|
||||
while (index < this.ranges.length && this.ranges[index]!.end + 1 < generation) {
|
||||
index++
|
||||
}
|
||||
const index = this.firstRangeReachableFrom(generation)
|
||||
const current = this.ranges[index]
|
||||
if (!current || generation + 1 < current.start) {
|
||||
this.ranges.splice(index, 0, { start: generation, end: generation })
|
||||
@@ -26,15 +23,27 @@ export class SshPtyClosedGenerationRanges {
|
||||
}
|
||||
|
||||
has(generation: number): boolean {
|
||||
for (const range of this.ranges) {
|
||||
if (generation < range.start) {
|
||||
return false
|
||||
}
|
||||
if (generation <= range.end) {
|
||||
return true
|
||||
const range = this.ranges[this.firstRangeReachableFrom(generation)]
|
||||
return range !== undefined && generation >= range.start && generation <= range.end
|
||||
}
|
||||
|
||||
// Why binary search rather than a scan: ranges only collapse to one while every allocated
|
||||
// generation is eventually closed. A generation that is allocated and never closed leaves a
|
||||
// permanent gap, and `has` is on the per-output-chunk admission path, so a linear scan turns
|
||||
// fragmentation into a hot-path cost (measured ~1800x a plain Set at 20k ranges) and makes `add`
|
||||
// quadratic. Log-time keeps the degraded shape no worse than the Set this replaced.
|
||||
private firstRangeReachableFrom(generation: number): number {
|
||||
let low = 0
|
||||
let high = this.ranges.length
|
||||
while (low < high) {
|
||||
const mid = (low + high) >> 1
|
||||
if (this.ranges[mid]!.end + 1 < generation) {
|
||||
low = mid + 1
|
||||
} else {
|
||||
high = mid
|
||||
}
|
||||
}
|
||||
return false
|
||||
return low
|
||||
}
|
||||
|
||||
get size(): number {
|
||||
|
||||
@@ -1,6 +1,17 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { SshPtyClosedGenerationRanges } from './ssh-pty-closed-generation-ranges'
|
||||
|
||||
const HOSTS = 8
|
||||
const RECONNECTS = 20_000
|
||||
|
||||
function lcg(seed: number): () => number {
|
||||
let state = seed >>> 0
|
||||
return () => {
|
||||
state = (state * 1_664_525 + 1_013_904_223) >>> 0
|
||||
return state / 0x1_0000_0000
|
||||
}
|
||||
}
|
||||
|
||||
// Provider generations come from one process-global counter shared by every SSH target
|
||||
// (ssh-pty-output-intake-registry.ts), so lower-numbered generations are routinely still live on a
|
||||
// different host. The closed set must answer exactly; a high-water approximation would reject a
|
||||
@@ -16,15 +27,64 @@ describe('closed provider generations across concurrent SSH targets', () => {
|
||||
expect(closed.has(1)).toBe(false)
|
||||
})
|
||||
|
||||
it('stays flat across a long monotonic run of reconnects', () => {
|
||||
// The leak this replaces: one retained entry per relay reconnect, forever.
|
||||
it('collapses to one range when two targets take turns reconnecting', () => {
|
||||
// Each reconnect closes the generation it is replacing, so alternating hosts still produce a
|
||||
// contiguous closed run -- interleaving alone does not fragment the list.
|
||||
const closed = new SshPtyClosedGenerationRanges()
|
||||
for (let generation = 1; generation <= 100_000; generation += 1) {
|
||||
closed.add(generation)
|
||||
let nextGeneration = 1
|
||||
const live = [nextGeneration++, nextGeneration++]
|
||||
for (let reconnect = 0; reconnect < RECONNECTS; reconnect += 1) {
|
||||
const host = reconnect % live.length
|
||||
closed.add(live[host]!)
|
||||
live[host] = nextGeneration++
|
||||
}
|
||||
|
||||
expect(closed.size).toBe(1)
|
||||
expect(closed.has(100_000)).toBe(true)
|
||||
expect(closed.has(100_001)).toBe(false)
|
||||
expect(closed.has(live[0]!)).toBe(false)
|
||||
expect(closed.has(live[1]!)).toBe(false)
|
||||
})
|
||||
|
||||
it('bounds ranges by the live generation count, not the reconnect count', () => {
|
||||
// Eight hosts reconnecting in scrambled order, with closes landing out of order the way a
|
||||
// deferred model migration settles them. Every gap is a generation that is still live, so the
|
||||
// list can never hold more than one range per live generation plus one.
|
||||
const random = lcg(0x5eed)
|
||||
const closed = new SshPtyClosedGenerationRanges()
|
||||
let nextGeneration = 1
|
||||
const live = Array.from({ length: HOSTS }, () => nextGeneration++)
|
||||
const settling: number[] = []
|
||||
let peakRanges = 0
|
||||
for (let reconnect = 0; reconnect < RECONNECTS; reconnect += 1) {
|
||||
const host = Math.floor(random() * HOSTS)
|
||||
settling.push(live[host]!)
|
||||
live[host] = nextGeneration++
|
||||
if (settling.length > 4) {
|
||||
closed.add(settling.splice(Math.floor(random() * settling.length), 1)[0]!)
|
||||
}
|
||||
peakRanges = Math.max(peakRanges, closed.size)
|
||||
}
|
||||
for (const generation of settling) {
|
||||
closed.add(generation)
|
||||
}
|
||||
|
||||
expect(peakRanges).toBeLessThanOrEqual(HOSTS + 4 + 1)
|
||||
expect(closed.size).toBeLessThanOrEqual(HOSTS + 1)
|
||||
for (const generation of live) {
|
||||
expect(closed.has(generation)).toBe(false)
|
||||
}
|
||||
})
|
||||
|
||||
it('retains one range for every generation that is allocated and never closed', () => {
|
||||
// The honest limitation, pinned rather than assumed away: ranges compact only against closes.
|
||||
// A generation that is allocated and abandoned without a close leaves a permanent gap, so this
|
||||
// structure is bounded by unclosed generations -- not by anything the reconnect loop does.
|
||||
const closed = new SshPtyClosedGenerationRanges()
|
||||
let nextGeneration = 1
|
||||
for (let reconnect = 0; reconnect < 5_000; reconnect += 1) {
|
||||
nextGeneration += 1
|
||||
closed.add(nextGeneration++)
|
||||
}
|
||||
|
||||
expect(closed.size).toBe(5_000)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -30,9 +30,12 @@ export class SshPtyModelAdmission {
|
||||
private readonly idleWaiters = new Map<string, Set<() => void>>()
|
||||
// Why ranges rather than a Set: provider generations are a process-global monotonic counter
|
||||
// shared by every SSH target, so this grew one entry per relay reconnect for the life of the
|
||||
// process. Contiguous closed runs collapse into a single range, and membership stays exact --
|
||||
// process. Because a reconnect closes the generation it replaces, the closed run stays
|
||||
// contiguous apart from the generations that are still live, which bounds the list at roughly
|
||||
// one range per concurrently connected target however hosts interleave. Membership stays exact:
|
||||
// generations below the high-water mark can still be live on another host, so a high-water
|
||||
// approximation would reject a healthy target's output.
|
||||
// approximation would reject a healthy target's output. Do not scope this per target -- with a
|
||||
// global counter each target's own closes are non-adjacent, which is the growth this avoids.
|
||||
private readonly closingGenerations = new SshPtyClosedGenerationRanges()
|
||||
private readonly migratingPtys = new Set<string>()
|
||||
private globalSourceUnits = 0
|
||||
|
||||
@@ -70,28 +70,48 @@ function createDeferred<T>(): { promise: Promise<T>; resolve: (value: T) => void
|
||||
return { promise, resolve }
|
||||
}
|
||||
|
||||
type FixtureListener = { port: number; inode: number }
|
||||
|
||||
const DEFAULT_LISTENER: FixtureListener = { port: 3000, inode: 11_111 }
|
||||
|
||||
function tcpRow(index: number, { port, inode }: FixtureListener): string {
|
||||
const hexPort = port.toString(16).toUpperCase().padStart(4, '0')
|
||||
return `${index}: 0100007F:${hexPort} 00000000:0000 0A 00000000:00000000 00:00000000 00000000 1000 0 ${inode}`
|
||||
}
|
||||
|
||||
function mockLinuxProcScan({
|
||||
pidCount,
|
||||
fdCount,
|
||||
firstReadlink
|
||||
firstReadlink,
|
||||
listeners = [DEFAULT_LISTENER],
|
||||
inodesByPidOffset,
|
||||
cmdlineByPidOffset
|
||||
}: {
|
||||
pidCount: number
|
||||
fdCount: number
|
||||
firstReadlink?: Promise<string>
|
||||
// Listening rows in /proc/net/tcp. Several rows are what makes an early exit keyed on
|
||||
// `result.size === inodes.size` distinguishable from one keyed on `result.size > 0`.
|
||||
listeners?: readonly FixtureListener[]
|
||||
// Socket inode per fd index, keyed by the pid's offset from PID_BASE. Offsets left out hold no
|
||||
// listening socket at all; fd indexes past the end of a list link to a non-socket path.
|
||||
inodesByPidOffset?: ReadonlyMap<number, readonly number[]>
|
||||
cmdlineByPidOffset?: ReadonlyMap<number, string>
|
||||
}): void {
|
||||
const tcpHeader =
|
||||
'sl local_address rem_address st tx_queue rx_queue tr tm->when retrnsmt uid timeout inode'
|
||||
const tcpRow =
|
||||
'0: 0100007F:0BB8 00000000:0000 0A 00000000:00000000 00:00000000 00000000 1000 0 11111'
|
||||
readFileMock.mockImplementation(async (path: string) => {
|
||||
if (path === '/proc/net/tcp') {
|
||||
return `${tcpHeader}\n${tcpRow}\n`
|
||||
return `${tcpHeader}\n${listeners.map((listener, index) => tcpRow(index, listener)).join('\n')}\n`
|
||||
}
|
||||
if (path === '/proc/net/tcp6') {
|
||||
return `${tcpHeader}\n`
|
||||
}
|
||||
if (path.endsWith('/cmdline')) {
|
||||
return '/usr/bin/node\0server.js'
|
||||
const cmdlineMatch = path.match(/^\/proc\/(\d+)\/cmdline$/)
|
||||
if (cmdlineMatch) {
|
||||
return (
|
||||
cmdlineByPidOffset?.get(Number(cmdlineMatch[1]) - PID_BASE) ?? '/usr/bin/node\0server.js'
|
||||
)
|
||||
}
|
||||
throw new Error(`unexpected readFile: ${path}`)
|
||||
})
|
||||
@@ -109,13 +129,21 @@ function mockLinuxProcScan({
|
||||
})
|
||||
|
||||
let first = true
|
||||
readlinkMock.mockImplementation(() => {
|
||||
readlinkMock.mockImplementation((path: string) => {
|
||||
if (first && firstReadlink) {
|
||||
first = false
|
||||
return firstReadlink
|
||||
}
|
||||
first = false
|
||||
return Promise.resolve('socket:[11111]')
|
||||
if (!inodesByPidOffset) {
|
||||
return Promise.resolve(`socket:[${DEFAULT_LISTENER.inode}]`)
|
||||
}
|
||||
const match = path.match(/^\/proc\/(\d+)\/fd\/(\d+)$/)
|
||||
if (!match) {
|
||||
throw new Error(`unexpected readlink: ${path}`)
|
||||
}
|
||||
const inode = inodesByPidOffset.get(Number(match[1]) - PID_BASE)?.[Number(match[2])]
|
||||
return Promise.resolve(inode === undefined ? '/dev/null' : `socket:[${inode}]`)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -142,6 +170,93 @@ describe('PortScanHandler Linux walk bounds', () => {
|
||||
|
||||
expect(readlinkMock).toHaveBeenCalledTimes(6)
|
||||
})
|
||||
|
||||
it('resolves an owner for every listening socket before it exits', async () => {
|
||||
// The exit condition has to be "every inode is attributed", not "some inode is". With one
|
||||
// fixture row those are the same assertion, which is how an exit-after-the-first-listener bug
|
||||
// would slip through: ports 3001 and 3002 would come back ownerless.
|
||||
mockLinuxProcScan({
|
||||
pidCount: 500,
|
||||
fdCount: 1,
|
||||
listeners: [
|
||||
{ port: 3000, inode: 11_111 },
|
||||
{ port: 3001, inode: 22_222 },
|
||||
{ port: 3002, inode: 33_333 }
|
||||
],
|
||||
inodesByPidOffset: new Map([
|
||||
[0, [11_111]],
|
||||
[1, [22_222]],
|
||||
[2, [33_333]]
|
||||
])
|
||||
})
|
||||
|
||||
await expect(capturePortDetectHandler()({}, requestContext())).resolves.toEqual({
|
||||
platform: 'linux',
|
||||
ports: [
|
||||
{ host: '127.0.0.1', port: 3000, pid: PID_BASE, processName: 'node' },
|
||||
{ host: '127.0.0.1', port: 3001, pid: PID_BASE + 1, processName: 'node' },
|
||||
{ host: '127.0.0.1', port: 3002, pid: PID_BASE + 2, processName: 'node' }
|
||||
]
|
||||
})
|
||||
// Three owners found means three readlinks: it stops at the third pid, not the five hundredth.
|
||||
expect(readlinkMock).toHaveBeenCalledTimes(3)
|
||||
})
|
||||
|
||||
it('keeps walking past an attributed socket to reach a later owner', async () => {
|
||||
// A partially attributed table is the sharpest case: the first pid resolves one inode, and the
|
||||
// second inode is only reachable at the end of the walk.
|
||||
mockLinuxProcScan({
|
||||
pidCount: 4,
|
||||
fdCount: 1,
|
||||
listeners: [
|
||||
{ port: 3000, inode: 11_111 },
|
||||
{ port: 3001, inode: 22_222 }
|
||||
],
|
||||
inodesByPidOffset: new Map([
|
||||
[0, [11_111]],
|
||||
[3, [22_222]]
|
||||
])
|
||||
})
|
||||
|
||||
await expect(capturePortDetectHandler()({}, requestContext())).resolves.toEqual({
|
||||
platform: 'linux',
|
||||
ports: [
|
||||
{ host: '127.0.0.1', port: 3000, pid: PID_BASE, processName: 'node' },
|
||||
{ host: '127.0.0.1', port: 3001, pid: PID_BASE + 3, processName: 'node' }
|
||||
]
|
||||
})
|
||||
expect(readlinkMock).toHaveBeenCalledTimes(4)
|
||||
})
|
||||
})
|
||||
|
||||
describe('PortScanHandler shared listening inode attribution', () => {
|
||||
it('attributes a shared inode to the first holder the walk reaches', async () => {
|
||||
// An nginx master and its workers (or a Node cluster) share one listening inode. The walk used
|
||||
// to overwrite the entry for every later holder, so the last pid in readdir order won; exiting
|
||||
// as soon as the inode is attributed makes the first one win instead. That is the better
|
||||
// answer -- the master owns the socket -- but it is a visible change to the name in the ports
|
||||
// UI, so pin it here rather than let it drift.
|
||||
mockLinuxProcScan({
|
||||
pidCount: 3,
|
||||
fdCount: 1,
|
||||
inodesByPidOffset: new Map([
|
||||
[0, [DEFAULT_LISTENER.inode]],
|
||||
[1, [DEFAULT_LISTENER.inode]],
|
||||
[2, [DEFAULT_LISTENER.inode]]
|
||||
]),
|
||||
cmdlineByPidOffset: new Map([
|
||||
[0, '/usr/sbin/nginx\0master process'],
|
||||
[1, '/usr/sbin/nginx-worker\0worker process'],
|
||||
[2, '/usr/sbin/nginx-worker\0worker process']
|
||||
])
|
||||
})
|
||||
|
||||
await expect(capturePortDetectHandler()({}, requestContext())).resolves.toEqual({
|
||||
platform: 'linux',
|
||||
ports: [{ host: '127.0.0.1', port: 3000, pid: PID_BASE, processName: 'nginx' }]
|
||||
})
|
||||
expect(readlinkMock).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
})
|
||||
|
||||
describe('PortScanHandler Linux cancellation', () => {
|
||||
|
||||
Reference in New Issue
Block a user