index.ts 133 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308
  1. /**
  2. * CPython subprocess code runtime: a fresh `python3` process runs each model program under an
  3. * asyncio event loop with top-level ``await``. Binding calls travel on fd 3 as JSON-lines,
  4. * leaving stdout/stderr free for the program's own output. This is containment, not a security
  5. * boundary: model code has bash-equivalent trust, contained by an empty environment, RLIMIT_CPU
  6. * + RLIMIT_AS, wall-clock timeout, and SIGTERM→grace→SIGKILL on the process group.
  7. *
  8. * The package owns the versionless fd-3 wire protocol between the Node host and
  9. * the CPython subprocess. The protocol's host-side codec and hostile-frame
  10. * validators are re-exported so every consumer of the wire shares one
  11. * vocabulary.
  12. * @module @deepseek-ai/dsh-experimental-code-runtime-python
  13. */
  14. import { spawn, type ChildProcessWithoutNullStreams } from 'node:child_process'
  15. import { accessSync, copyFileSync, constants as fsConstants, mkdtempSync, readFileSync, rmSync, statSync } from 'node:fs'
  16. import { tmpdir } from 'node:os'
  17. import { delimiter, dirname, isAbsolute, join, resolve } from 'node:path'
  18. import { fileURLToPath } from 'node:url'
  19. import type { Duplex } from 'node:stream'
  20. import { Context } from 'cordis'
  21. import z from 'schemastery'
  22. import { CodeRuntime, DUNDER_MEMBER, PORTABLE_RESERVED_WORDS, RESERVED_BINDING_GLOBALS, RESERVED_ERROR_MEMBERS } from '@deepseek-ai/dsh-code-runtime'
  23. import type { CodeBindingErrorClass, CodeBindingFunction, CodeJsonValue, CodeRunFailure, CodeRunRequest, CodeRunResult } from '@deepseek-ai/dsh-code-runtime'
  24. import { snapshotJsonValue } from '@deepseek-ai/dsh-session'
  25. import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
  26. import type { BootMessage, ChildToHost, ReplyMessage } from './protocol.ts'
  27. import { checkDoneValue, encodeJsonPlain, hasUnsafeIntegerToken, logTruncationMarker, validateChildFrame } from './protocol.ts'
  28. // Re-export the fd-3 wire vocabulary so the runtime and its tests share one
  29. // import surface; the protocol layer owns the definitions.
  30. export type { BootMessage, ChildToHost, ReplyMessage } from './protocol.ts'
  31. export {
  32. checkDoneValue,
  33. encodeJsonPlain,
  34. hasNonLosslessNumber,
  35. hasUnsafeIntegerToken,
  36. logTruncationMarker,
  37. validateChildFrame,
  38. } from './protocol.ts'
  39. /** Plugin config: every cap, changeable from `cordis.yml` (no hardcoded tunables). */
  40. export interface Config {
  41. /**
  42. * RLIMIT_CPU in whole seconds (a positive integer — `setrlimit` in the child
  43. * rejects a float). The child sets the soft limit to `cpuSeconds` and the
  44. * hard limit to `cpuSeconds + 1`: the kernel delivers SIGXCPU at the soft
  45. * limit, which the host classifies as a `timeout`; the +1s hard limit is a
  46. * SIGKILL backstop for a program that traps SIGXCPU. Granularity is seconds —
  47. * a coarser counterpart to the worker backend's millisecond `computeMs`.
  48. */
  49. cpuSeconds?: number
  50. /** Wall-clock ceiling in milliseconds; backstops CPU time for programs awaiting a promise nobody resolves. */
  51. maxWallMs?: number
  52. /**
  53. * RLIMIT_AS in mebibytes; caps address space so a runaway allocation fails
  54. * cleanly. Not applied on Darwin, where the dyld shared cache mapped into
  55. * every process at exec exceeds any practical cap and the kernel rejects
  56. * the call; `cpuSeconds` and `maxWallMs` still bound the run there. Bounds
  57. * `maxLogBytes`/`maxValueBytes` at load on EVERY platform (this static check
  58. * runs on Darwin too, where only the runtime `setrlimit` is skipped): each
  59. * budget times a worst-case Unicode expansion must fit this byte count minus a
  60. * fixed interpreter baseline, so a near-budget output cannot breach the address
  61. * space during the child's build-and-encode.
  62. */
  63. addressSpaceMb?: number
  64. /**
  65. * Shared byte budget for captured log text (host-side ledger). Bounded at load
  66. * against `addressSpaceMb`: the child builds and encodes a near-budget entry
  67. * under RLIMIT_AS with several copies live at once, so this cap times the
  68. * worst-case Unicode expansion must fit the address space left after the
  69. * interpreter baseline (see `addressSpaceMb`) — a load-time rejection, not a
  70. * runtime clamp.
  71. */
  72. maxLogBytes?: number
  73. /**
  74. * Byte cap for the completion value. Bounded at load against `addressSpaceMb`
  75. * the same way `maxLogBytes` is: the child builds and encodes a near-budget
  76. * value under RLIMIT_AS with several copies live at once, so this cap times the
  77. * worst-case Unicode expansion must fit the address space left after the
  78. * interpreter baseline.
  79. */
  80. maxValueBytes?: number
  81. /** SIGTERM→SIGKILL grace period on kill, matching bash-local's default. */
  82. graceMs?: number
  83. /**
  84. * Absolute path or basename of the CPython interpreter to spawn. Resolved
  85. * through `PATH` when a basename is given.
  86. */
  87. pythonBin?: string
  88. }
  89. /** {@link Config} with all defaults filled. */
  90. type ResolvedConfig = Required<Config>
  91. /**
  92. * The seam's language-portable identifier subset (see
  93. * `CodeBindingNamespace.global`) — identical to Python's identifier grammar,
  94. * so the shared contract needs no per-backend mapping here.
  95. */
  96. const IDENTIFIER = /^[A-Za-z_][A-Za-z0-9_]*$/
  97. /**
  98. * The seam's cross-language reserved-word union: the portable-identifier
  99. * contract promises a namespace list valid here is valid on every backend, so
  100. * a JS keyword like `typeof` is refused even though it is a legal Python name.
  101. */
  102. const RESERVED_NAMES = PORTABLE_RESERVED_WORDS
  103. /**
  104. * The seam's shared backend-owned globals (`console` is the worker's slot;
  105. * `__dsh_main__`/`__builtins__`/`__name__` are this bootstrap's wrapper and
  106. * seeded module globals). Shared so a namespace list valid on one backend is
  107. * valid on all — colliding with an owned slot would be silently overwritten
  108. * (or overwrite builtins), so the seam rejects them up front.
  109. */
  110. const RUNTIME_OWNED_GLOBALS = RESERVED_BINDING_GLOBALS
  111. /**
  112. * The seam's shared error-member exclusions (`RESERVED_ERROR_MEMBERS` +
  113. * dunder-form names) — enforced identically here and in the worker backend so
  114. * an errorClass valid on one backend is valid on all. Several dunders are
  115. * constrained CPython descriptors whose `setattr` raises while constructing
  116. * the very rejection it was meant to carry; the exact set is an interpreter
  117. * version detail, hence the dunder-wide rule at the seam.
  118. */
  119. const EXCEPTION_RESERVED_MEMBERS = RESERVED_ERROR_MEMBERS
  120. const DUNDER = DUNDER_MEMBER
  121. /**
  122. * The `py/` scripts the interpreter must be able to open: the entry script plus
  123. * every module it imports from its own directory. Kept beside the built JS so a
  124. * consumer package with `files: ['lib', 'py']` ships both.
  125. */
  126. const PY_SCRIPTS = ['bootstrap.py', 'protocol.py']
  127. /**
  128. * Copy the `py/` scripts to a real filesystem directory and return the entry
  129. * script's path there.
  130. *
  131. * The interpreter is an EXTERNAL process, so it can only open paths the OS
  132. * resolves. Inside the single-file Python-SDK executable, `import.meta.url`
  133. * resolves into pkg's virtual filesystem, which Node reads through its patched
  134. * `fs` but `python3` cannot see at all — the spawn fails with ENOENT on a path
  135. * that exists as far as the host is concerned. `bootstrap.py` additionally
  136. * inserts its own directory on `sys.path` to import the sibling `protocol.py`,
  137. * so both files must land in the SAME real directory.
  138. *
  139. * The copy is unconditional rather than gated on a bundled-runtime probe: the
  140. * read goes through Node's `fs` either way, and one code path means the
  141. * packaged deployment runs what the tests exercise. Placement is under
  142. * `os.tmpdir()` with `0o700` keeps the scripts off other users' reach, but NOT
  143. * the model's: the child runs as the same UID as the host, so a program can
  144. * rewrite the very files it was started from. Hence one copy per RUN, discarded
  145. * at settlement — a rewrite then damages only the run that performed it, which
  146. * is what fresh-subprocess-per-run already promises. Sharing one copy across
  147. * runs made an overwritten `bootstrap.py` break the next run.
  148. *
  149. * Deliberately SYNCHRONOUS. An `await` here would open an async boundary in
  150. * `execute` before the run is registered in `live` and before the abort
  151. * listener is installed, so a disposal or an abort landing in that window would
  152. * be missed: `teardown` would see no runs and return while the continuation
  153. * went on to spawn a subprocess, and an `addEventListener('abort')` installed
  154. * afterwards does not replay an event that already fired. Three small
  155. * filesystem operations per run are not worth that class of race, and `execute`
  156. * already runs synchronously up to `spawn`.
  157. *
  158. * A failed copy removes the directory here, so a partial attempt never outlives
  159. * the call that made it; a successful one is the caller's to remove, which it
  160. * derives from the returned path.
  161. *
  162. * @returns the absolute path of the materialized entry script.
  163. */
  164. function materializePyScripts(): string {
  165. const dir = mkdtempSync(join(tmpdir(), 'dsh-code-runtime-python-'))
  166. const source = fileURLToPath(new URL('../py/', import.meta.url))
  167. try {
  168. for (const name of PY_SCRIPTS) copyFileSync(join(source, name), join(dir, name))
  169. } catch (error: unknown) {
  170. try {
  171. rmSync(dir, { recursive: true, force: true })
  172. } catch {
  173. // Swallows only a failure to remove the partial staging directory. The
  174. // caller reports the copy failure that got us here, which is the
  175. // diagnosable one; nothing else can act on a temp dir we cannot unlink.
  176. }
  177. throw error
  178. }
  179. return join(dir, 'bootstrap.py')
  180. }
  181. /**
  182. * A frame's RAW length is capped before JSON.parse: the 64 MiB fd-3 frame
  183. * parse cap bounds the bytes, not the decoded structure, and a compact wide
  184. * frame near that ceiling (e.g. a huge array of tiny elements) could decode to
  185. * far more host memory than the wire admitted — an OOM inside the receive
  186. * path. 64 MiB raw admits every legal config (the widest in-tree completion
  187. * and binding frames are ~12 MB) while bounding decode amplification to a
  188. * roughly constant factor of the wire bytes. The unframed-buffer counter is
  189. * checked against this same cap BEFORE a `Buffer.concat` join, so an oversized
  190. * frame is dropped at one copy of its wire bytes. A hostile-peer invariant,
  191. * not a deployment choice.
  192. */
  193. const FRAME_PARSE_CAP_BYTES = 64 * 1024 * 1024
  194. /**
  195. * Fragments the unframed fd-3 buffer may hold before they are coalesced into
  196. * one Buffer, bounding retained per-chunk overhead that the byte cap cannot
  197. * see: the cap meters payload bytes, while each chunk is a distinct Buffer
  198. * with its own object and backing store. A
  199. * program writing single bytes without a newline produced one chunk per write.
  200. * 1024 keeps the overhead a small constant factor of the payload while leaving
  201. * normal pipe-sized reads (which arrive in far fewer, much larger chunks)
  202. * untouched. A framing invariant, not a deployment choice.
  203. */
  204. const MAX_PENDING_CHUNKS = 1024
  205. /**
  206. * Replies the host retains before fd 3 accepts them. The drain loop writes one
  207. * reply per iteration and waits for `drain` when the pipe is full; a child
  208. * that never reads its replies (hostile or wedged) leaves the pipe full, so
  209. * every call frame it keeps sending adds a reply the drain cannot write, and
  210. * the backlog would grow without bound until the wall clock. 1024 keeps
  211. * legitimate concurrent gathers (measured queue depths reach 11) far below
  212. * the ceiling while bounding the hostile backlog; the run settles as a
  213. * worker-exit past it, like the frame cap settles an oversized frame. A
  214. * framing invariant, not a deployment choice.
  215. */
  216. const MAX_PENDING_REPLIES = 1024
  217. /**
  218. * Bytes a frame spends on its own JSON structure around a capped payload, used
  219. * to bound `maxLogBytes`/`maxValueBytes` against {@link FRAME_PARSE_CAP_BYTES}
  220. * (the receive path rejects raw frames past that cap, settling the run as a
  221. * worker-exit).
  222. * The widest carrier is `{"type":"log","text":"","truncated":true}` at 41
  223. * bytes; 64 rounds that up so adding a field to either frame does not silently
  224. * invalidate the bound. A protocol constant, not a deployment choice.
  225. */
  226. const FRAME_ENVELOPE_BYTES = 64
  227. /**
  228. * Smallest `maxLogBytes` the backend can honor. The truncation marker alone
  229. * (`logTruncationMarker`) must serialize within the budget, or a marker-only
  230. * truncated run returns more than the configured cap: the marker text is
  231. * `[dsh-code-runtime-python] log capture truncated at <N> bytes` — 51 fixed
  232. * characters (the bracketed prefix `[dsh-code-runtime-python] log capture
  233. * truncated at ` counts both square brackets) plus the digits of N plus 6 —
  234. * and its serialized form adds 4 (two quotes, two array brackets), so the
  235. * smallest N that admits its own marker is 63 (51 + 2 + 6 + 4 = 63); 64 is the
  236. * floor with one byte of room. The marker itself remains envelope, not
  237. * payload, so a truncated run with admitted entries serializes to at most
  238. * `maxLogBytes + marker + envelope`.
  239. * `maxValueBytes` has no floor beyond the positive-integer requirement: a
  240. * completion can be as small as a single byte (`1`), and the done-frame
  241. * envelope is seam protocol cost, not the advertised completion budget.
  242. */
  243. const MIN_LOG_BYTES = 64
  244. /**
  245. * Extra time added to `graceMs` before the post-kill close-deadline force-settles
  246. * a run whose `close` never fires (a setsid-escaped orphan holds our inherited
  247. * stdio; see the `closeDeadline` arm in {@link PythonCodeRuntime.execute}). It
  248. * covers the OS reaping the killed child itself after SIGKILL — not a deployment
  249. * choice but a fixed safety margin, so it is a constant rather than a config knob.
  250. */
  251. const CLOSE_REAP_MARGIN_MS = 2_000
  252. /**
  253. * Worst-case peak child-process bytes a one-`maxLogBytes`/`maxValueBytes`-budget
  254. * output can transiently occupy while the child charges and frames it, expressed
  255. * as a multiple of the budget. The child's ledgers trigger on CHARACTER count
  256. * against a serialized-BYTE budget, and an astral character is one character but
  257. * four bytes of CPython `str` storage and four UTF-8 bytes — so a budget's worth
  258. * of astral characters is ~4x the budget in each string that holds it. The
  259. * heaviest path holds THREE such copies at once: a single
  260. * `sys.stdout.write(line + "\n")` keeps the caller's `text` argument (alive for
  261. * the whole `write` call, ~4x), the line slice `text[pos:newline]` handed to
  262. * `LogBuffer.push` (~4x), and the `text.encode("utf-8")` copy `_push_locked`
  263. * takes to charge and ship it (~4x). The settlement `flush_line` path holds only
  264. * two (its `"".join(...)` and that encode copy — it drops the pending chunks
  265. * before pushing), so the newline path is the binding worst case. Twelve covers
  266. * those three simultaneous ~4x copies. The interpreter baseline is NOT in this
  267. * multiple — it is reserved separately as {@link INTERPRETER_BASELINE_BYTES} —
  268. * because it is a fixed cost, not one that scales with the budget. Used to bound
  269. * `maxLogBytes`/`maxValueBytes` against `addressSpaceMb` at load, with a `>=` so
  270. * a budget whose worst-case peak exactly equals the room left after the baseline
  271. * is rejected (that peak plus the baseline is the whole address space, the
  272. * RLIMIT_AS edge), so a legitimate near-budget output truncates (log) or fails
  273. * as `output-limit` (value) rather than breaching `RLIMIT_AS` as `worker-exit`.
  274. * A fixed safety invariant tying the budgets to the address space, not a knob.
  275. */
  276. const OUTPUT_BUDGET_WORST_CASE_ADDRESS_SPACE_MULTIPLE = 12
  277. /**
  278. * Fixed address-space headroom reserved for the CPython interpreter itself
  279. * (loaded modules, the asyncio loop, import machinery) before the output-budget
  280. * multiple claims the rest. The budget check subtracts this from `addressSpaceMb`
  281. * so a budget sized right at `addressSpaceMb / MULTIPLE` — which the multiple
  282. * alone would admit — cannot leave the peak output allocation plus the
  283. * interpreter over the limit. Sized against ADDRESS SPACE, which is what
  284. * `RLIMIT_AS` bounds, not resident set: the bootstrap's own measurement is
  285. * 30.23 MiB of mappings for a `python3 -I` child (see `_make_cpu_enforcer`,
  286. * which also records the 64 MiB glibc per-thread arena reservation that pushes
  287. * it to 102.37 MiB when threads are used). 64 MiB is roughly twice the measured
  288. * baseline, leaving room for allocator arenas and import jitter. The value is a
  289. * fixed safety margin, not a deployment knob.
  290. */
  291. const INTERPRETER_BASELINE_BYTES = 64 * 1024 * 1024
  292. /**
  293. * Interval between process-group liveness probes while settlement waits for an
  294. * escalated SIGKILL to empty the group (see the `killing` branch in
  295. * {@link PythonCodeRuntime.execute}'s settle). A poll rather than an event
  296. * because the group members are the model's own descendants, which the host does
  297. * not `wait()` for and gets no exit signal from; the probe is a signal-0
  298. * `process.kill(-pid, 0)`, so the interval only bounds how promptly a now-empty
  299. * group is noticed, capped by `graceMs + CLOSE_REAP_MARGIN_MS`.
  300. */
  301. const GROUP_REAP_POLL_MS = 50
  302. /**
  303. * Extract a human message from an unknown thrown value.
  304. *
  305. * `String(error)` runs the value's own conversion, and a host binding may reject
  306. * with an object whose `Symbol.toPrimitive` or `toString` throws. One call site
  307. * is a detached async reply callback, where that throw escapes as an unhandled
  308. * rejection: the reply frame is never written, the program stays blocked on
  309. * `await`, and the run degrades to a `maxWallMs` timeout (a Node host without an
  310. * `unhandledRejection` listener exits outright). The conversion is therefore
  311. * wrapped, with a fixed literal as the fallback — the value already proved it
  312. * cannot be rendered, so nothing derived from it is safe to try.
  313. *
  314. * `Error.message` is typed `string` but is a plain writable property, so a
  315. * rejecting binding can hand back an `Error` carrying any value there. The
  316. * `Error` arm therefore goes through the same conversion rather than returning
  317. * `message` verbatim: the returned string crosses the wire under
  318. * `encodeJsonPlain`'s JSON-plain precondition, where a cyclic object grows the
  319. * encoder stack until the host exhausts memory and any other unsupported value
  320. * prevents the reply frame outright.
  321. *
  322. * The same conversion renders abort reasons, which reach an `AbortSignal`
  323. * listener: Node reports a throw from such a listener as an uncaught exception,
  324. * so an unwrapped conversion there can terminate the host with the run left
  325. * unsettled.
  326. *
  327. * @param error The thrown value, of unknown shape.
  328. * @returns The value's message or string form; a fixed placeholder when its own
  329. * conversion throws.
  330. */
  331. function messageOf(error: unknown): string {
  332. try {
  333. return String(error instanceof Error ? error.message : error)
  334. } catch {
  335. // Swallows only a throw from the value's own `message` getter or string
  336. // conversion. Nothing else runs inside the try, and the placeholder is a
  337. // literal, so this cannot throw again.
  338. return '<unrenderable rejection value>'
  339. }
  340. }
  341. /**
  342. * A process's start time, as the identity half of (pid, started).
  343. *
  344. * A pid is reusable the moment the kernel reaps it, so signalling one that a
  345. * later process inherited would terminate an unrelated process group. Start
  346. * time is what distinguishes the original from its replacement: `kill(pid, 0)`
  347. * answers "does this number exist", which is true for both.
  348. *
  349. * Linux reads field 22 of `/proc/<pid>/stat` (starttime in clock ticks); the
  350. * field is positional after the comm field's closing parenthesis, which is
  351. * parsed from the LAST such character because a process name may contain one.
  352. * Darwin has no `/proc`, so the caller gets `undefined` there and `killGroup`
  353. * signals the pgid without the identity re-check rather than paying a `ps`
  354. * fork on a teardown path. Any read failure is `undefined` for the same
  355. * reason: this
  356. * hardens a narrow race and must never be the thing that breaks teardown.
  357. * @param pid - the process to read.
  358. * @returns its start time, or undefined when unavailable.
  359. */
  360. export function readProcessStart(pid: number): string | undefined {
  361. /* v8 ignore next -- one arm per platform: the Linux coverage lane always takes the read path, and Darwin always this one. */
  362. if (process.platform !== 'linux') return undefined
  363. try {
  364. const stat = readFileSync(`/proc/${String(pid)}/stat`, 'utf8')
  365. const fields = stat.slice(stat.lastIndexOf(')') + 2).split(' ')
  366. // Field 22 overall; the slice above dropped pid and comm, so it is index 19.
  367. return fields[19]
  368. } catch {
  369. return undefined
  370. }
  371. }
  372. /**
  373. * Resolve `pythonBin` to an absolute path against the CURRENT process `PATH`,
  374. * BEFORE the child spawns with an empty environment. A basename (the default
  375. * `python3`) would otherwise fail: `env: {}` drops `PATH`, so Node's own lookup
  376. * falls back to the platform default (`/usr/bin:/bin`) and misses interpreters
  377. * that live only on the caller's `PATH` (Nix, pyenv, Homebrew, conda). An
  378. * absolute or explicitly relative path is validated directly: it must exist,
  379. * be executable, and be a regular file — a missing, non-executable, or
  380. * directory path is a self-contained configuration error that must fail at
  381. * load, not at the first run (the child spawns with an empty environment, so
  382. * execvp's platform default would otherwise silently mask the mistake). A
  383. * relative explicit path resolves against the host CWD, mirroring where
  384. * `spawn` would have looked for it. When no `PATH` entry holds an executable
  385. * match, `undefined` is returned and the LOAD check rejects the configuration:
  386. * falling back to the bare name would let spawn's `env: {}` execvp silently
  387. * start a system interpreter from the platform default PATH that the caller
  388. * never asked for.
  389. * @param bin - the configured interpreter (absolute or relative path, or bare command).
  390. * @returns an absolute path when resolvable, else `undefined`.
  391. */
  392. export function resolvePythonBin(bin: string): string | undefined {
  393. if (isAbsolute(bin) || bin.includes('/')) {
  394. // An explicit path is used as given (resolved against the host CWD when
  395. // relative), but only when it is a real executable regular file. The same
  396. // checks as the PATH branch below: `accessSync(X_OK)` admits directories,
  397. // so `isFile` narrows further, and a path that fails either is not a
  398. // usable interpreter.
  399. const candidate = resolve(bin)
  400. try {
  401. accessSync(candidate, fsConstants.X_OK)
  402. if (!statSync(candidate).isFile()) return undefined
  403. return candidate
  404. } catch {
  405. return undefined
  406. }
  407. }
  408. const path = process.env.PATH
  409. /* v8 ignore next -- PATH is set in every environment the runtime boots in; the guard is defensive. */
  410. if (path === undefined) return undefined
  411. for (const dir of path.split(delimiter)) {
  412. // An empty PATH segment (a `::`, implicitly CWD on POSIX) and a RELATIVE
  413. // segment (`bin` or `.`) are skipped: a basename must never resolve against
  414. // the working directory, and the returned candidate must be an absolute
  415. // path — spawn() resolves a relative pythonBin against the host CWD, which
  416. // is outside the seam contract.
  417. if (dir === '' || !isAbsolute(dir)) continue
  418. const candidate = join(dir, bin)
  419. try {
  420. accessSync(candidate, fsConstants.X_OK)
  421. // A directory passes X_OK too, so require a regular file: a PATH entry
  422. // named like the interpreter (e.g. a `python3` directory) must not be
  423. // chosen over a later real interpreter.
  424. if (!statSync(candidate).isFile()) continue
  425. return candidate
  426. } catch {
  427. // Not executable here; try the next PATH entry.
  428. }
  429. }
  430. return undefined
  431. }
  432. /** The marker appended when a diagnostic message is byte-capped host-side. */
  433. const TRUNCATION_MARKER = '… [truncated]'
  434. /**
  435. * The marker's own UTF-8 byte length, reserved out of the budget so a capped
  436. * message stays WITHIN `maxValueBytes` rather than exceeding it by the marker.
  437. * The ellipsis is 3 bytes, so this is 15, not the string's 13 code units.
  438. */
  439. const TRUNCATION_MARKER_BYTES = Buffer.byteLength(TRUNCATION_MARKER, 'utf8')
  440. // Fatal UTF-8 decoder for fd-3 frames: `toString('utf8')` replaces illegal
  441. // bytes with U+FFFD, which would silently corrupt a completion or binding
  442. // payload a forged frame smuggled in; a fatal decode throws instead and the
  443. // frame is dropped. Non-stream mode keeps it stateless across lines.
  444. const UTF8_FATAL = new TextDecoder('utf-8', { fatal: true })
  445. /**
  446. * Serialized JSON byte width of one character, given its code point and the
  447. * one-character string. Control characters below 0x20 escape to `\uXXXX` (6)
  448. * except the five with short forms `\b \t \n \f \r` (2); `"` and `\` escape to
  449. * 2; a LONE surrogate escapes to `\uXXXX` (6) under ES2019 well-formed
  450. * `JSON.stringify`; everything else rides at its raw UTF-8 width.
  451. * @param code - the character's code point.
  452. * @param character - the one-character (or one-code-point) string.
  453. * @returns the character's serialized JSON byte width.
  454. */
  455. function serializedCharCost(code: number, character: string): number {
  456. if (code < 0x20) return code === 0x08 || code === 0x09 || code === 0x0a || code === 0x0c || code === 0x0d ? 2 : 6
  457. if (code === 0x22 || code === 0x5c) return 2
  458. if (code >= 0xd800 && code <= 0xdfff) return 6
  459. return Buffer.byteLength(character, 'utf8')
  460. }
  461. /**
  462. * Serialized JSON-string cost of `text` (the two quotes plus each character's
  463. * escaped byte width), measured WITHOUT materializing the escaped copy, and
  464. * abandoned the instant it exceeds `maxBytes`. `JSON.stringify(text)` would
  465. * allocate the whole escaped form first — up to sixfold a control-char-dense
  466. * string — so a near-budget line under a large `maxLogBytes` could momentarily
  467. * allocate over a gigabyte just to measure it. This walks code point by code
  468. * point (a matched surrogate pair yields its combined code point ≥ 0x10000; a
  469. * lone surrogate yields a value in 0xD800–0xDFFF that {@link serializedCharCost}
  470. * charges the full six escaped bytes) and stops at the cap, allocating nothing.
  471. * @param text - the candidate string.
  472. * @param maxBytes - the largest serialized size the caller can admit.
  473. * @returns the exact serialized byte cost, or `undefined` once it exceeds `maxBytes`.
  474. */
  475. function jsonStringCostUpTo(text: string, maxBytes: number): number | undefined {
  476. if (maxBytes < 2) return undefined
  477. let bytes = 2 // the enclosing quotes
  478. for (const character of text) {
  479. bytes += serializedCharCost(character.codePointAt(0) as number, character)
  480. if (bytes > maxBytes) return undefined
  481. }
  482. return bytes
  483. }
  484. /**
  485. * Cross-chunk UTF-8 state for {@link accrueStrayCost}: `expected` continuation
  486. * bytes still needed to finish the in-progress sequence, its total `width`, and
  487. * `lowerFirst`/`upperFirst`, the valid range for the NEXT continuation byte
  488. * (only the first continuation of a 3- or 4-byte lead is range-restricted; once
  489. * consumed, later continuations accept the full 0x80–0xBF). All zero between
  490. * sequences. Carried on each {@link StrayBuffer} so a multibyte character split
  491. * across pipe `data` chunks is costed as one character.
  492. */
  493. interface Utf8CostState { expected: number; width: number; lowerFirst: number; upperFirst: number }
  494. /**
  495. * Accrue the serialized JSON cost of raw pipe bytes `buf`, decoding UTF-8 the way
  496. * `toString('utf8')` (WHATWG) would so a byte that renders as U+FFFD is charged
  497. * the three bytes that replacement character serializes to. A naive tally that
  498. * charged every byte 1 let a `b"\xff"` flood (every byte illegal → U+FFFD each)
  499. * grow the residual to a full budget's worth of raw bytes before flushing; near
  500. * a large `maxLogBytes` that retained ~256 MiB, then `flushStray`'s
  501. * `Buffer.concat` + `toString` expanded it to a ~1 GiB peak. Charging only the
  502. * structural width would leave the same gap for structurally-well-formed but
  503. * ILLEGAL sequences a flood produces just as cheaply — a CESU-8 surrogate
  504. * (`ED A0 80`) or an overlong (`E0 80 80`) decodes to THREE U+FFFD (cost 9), not
  505. * one width-3 character, so this validates each lead's first continuation range
  506. * (WHATWG: `E0`→A0-BF, `ED`→80-9F, `F0`→90-BF, `F4`→80-8F, others 80-BF) and
  507. * charges 3 per byte of any sequence that breaks. A control byte below 0x20
  508. * costs 6 (`\uXXXX`) or 2 (five short escapes); `"`/`\` cost 2; ASCII costs 1; a
  509. * fully valid multibyte sequence costs its byte width (2/3/4). `state` carries
  510. * the in-progress sequence across chunks; an unfinished tail at stream end is
  511. * decoded by the final `flushStray` and costed exactly there.
  512. * @param buf - raw bytes from a stdout/stderr pipe chunk.
  513. * @param state - the pipe's carried UTF-8 sequence state, mutated in place.
  514. * @returns the serialized cost accrued by the bytes that resolved in this call.
  515. */
  516. function accrueStrayCost(buf: Buffer, state: Utf8CostState): number {
  517. let cost = 0
  518. let index = 0
  519. while (index < buf.length) {
  520. const byte = buf[index] as number
  521. if (state.expected > 0) {
  522. // The valid range for THIS continuation: the lead-specific range applies
  523. // to the first continuation only, then reverts to the full 0x80–0xBF.
  524. const consumed = state.width - state.expected
  525. const lower = consumed === 1 ? state.lowerFirst : 0x80
  526. const upper = consumed === 1 ? state.upperFirst : 0xbf
  527. if (byte >= lower && byte <= upper) {
  528. state.expected -= 1
  529. if (state.expected === 0) {
  530. cost += state.width
  531. state.width = 0
  532. }
  533. index += 1
  534. continue
  535. }
  536. // The sequence broke. WHATWG's maximal-subpart rule folds the bytes
  537. // consumed so far into ONE U+FFFD (cost 3), then reprocesses this byte as
  538. // a fresh start (no index advance). Charging per consumed byte would
  539. // over-count, which is memory-safe but wrong; folding to one is exact.
  540. cost += 3
  541. state.expected = 0
  542. state.width = 0
  543. continue
  544. }
  545. if (byte < 0x20) {
  546. cost += byte === 0x08 || byte === 0x09 || byte === 0x0a || byte === 0x0c || byte === 0x0d ? 2 : 6
  547. } else if (byte === 0x22 || byte === 0x5c) {
  548. cost += 2
  549. } else if (byte < 0x80) {
  550. cost += 1
  551. } else if (byte >= 0xc2 && byte <= 0xdf) {
  552. state.expected = 1
  553. state.width = 2
  554. state.lowerFirst = 0x80
  555. state.upperFirst = 0xbf
  556. } else if (byte >= 0xe0 && byte <= 0xef) {
  557. state.expected = 2
  558. state.width = 3
  559. // Exclude the overlong (E0 80-9F) and CESU-8 surrogate (ED A0-BF) ranges.
  560. state.lowerFirst = byte === 0xe0 ? 0xa0 : 0x80
  561. state.upperFirst = byte === 0xed ? 0x9f : 0xbf
  562. } else if (byte >= 0xf0 && byte <= 0xf4) {
  563. state.expected = 3
  564. state.width = 4
  565. // Exclude the overlong (F0 80-8F) and out-of-range (F4 90-BF) leads.
  566. state.lowerFirst = byte === 0xf0 ? 0x90 : 0x80
  567. state.upperFirst = byte === 0xf4 ? 0x8f : 0xbf
  568. } else {
  569. // 0x80–0xc1 and 0xf5–0xff never begin a valid sequence: U+FFFD (3).
  570. cost += 3
  571. }
  572. index += 1
  573. }
  574. return cost
  575. }
  576. /**
  577. * Cap a done-frame `error.message` to `maxValueBytes` host-side: a forged done
  578. * frame can carry an arbitrarily long message, so truncate by RAW UTF-8 byte
  579. * length and append the shared marker on overflow. Completion VALUES are never
  580. * truncated — the seam forbids substitution, so an oversized value fails the run
  581. * as `output-limit` instead (see the done case in `execute`).
  582. *
  583. * This is the RECEIVE-side backstop, and it bills by raw bytes on purpose,
  584. * unlike the producing-side `_cap_message` in `py/bootstrap.py`, which bills by
  585. * SERIALIZED (JSON-escaped) cost. The split is deliberate: `_cap_message`'s
  586. * output has to cross fd 3 as a JSON string, so its escaped width is what the
  587. * frame ceiling bounds; this function's output goes straight into
  588. * `CodeRunResult.error.message` and never re-crosses a frame-bounded channel, so
  589. * the honest measure of what it retains is the raw length. An honest child has
  590. * already capped the diagnostic by serialized cost, and raw length ≤ serialized
  591. * cost, so a well-formed message passes through unchanged. A forged message with
  592. * control characters could serialize to roughly six times its raw length, but it
  593. * is not travelling any capped channel, so the raw-byte bound is the right one:
  594. * the value it protects is the model-visible size of `error.message`, not a wire
  595. * width.
  596. *
  597. * The marker's bytes are RESERVED from the budget, not added on top: the whole
  598. * returned string, marker included, is at most `maxValueBytes` bytes. Appending
  599. * the marker after retaining a full budget's worth of text would overrun the
  600. * very cap this function exists to enforce. The one exception is a configured
  601. * cap SMALLER than the marker itself, which leaves no room for message text at
  602. * all; the marker alone is returned there, so the bound is
  603. * `max(maxValueBytes, 15)`. Reporting the truncation is worth those 15 bytes,
  604. * and the default cap is 32 KiB.
  605. * @param message - the error message from an inbound (possibly forged) done frame.
  606. * @param maxValueBytes - the configured completion-value budget, reused here.
  607. * @returns the message unchanged, or its byte-capped form on overflow.
  608. */
  609. function capMessage(message: string, maxValueBytes: number): string {
  610. // Code-unit bounds BEFORE any encode, so a forged done frame carrying a
  611. // message anywhere below the 64 MiB fd-3 frame parse cap cannot force a
  612. // full-length UTF-8 copy under a 32 KiB cap. One UTF-16 code unit encodes to
  613. // at least one UTF-8 byte and at most three: three for a non-ASCII BMP
  614. // character, two apiece for the pair halves sharing an astral code point's
  615. // four bytes, and three for a LONE surrogate, which `Buffer.from` renders as
  616. // U+FFFD. So at most maxValueBytes/3 code units cannot overflow the cap and
  617. // need no encode at all...
  618. if (message.length * 3 <= maxValueBytes) return message
  619. // ...and nothing past the first maxValueBytes code units can fit inside it,
  620. // so only that prefix is ever encoded — at most 3 * maxValueBytes bytes.
  621. const keep = Math.min(message.length, maxValueBytes)
  622. const whole = keep === message.length
  623. const bytes = Buffer.from(whole ? message : message.slice(0, keep), 'utf8')
  624. // A message that fits is measured against the WHOLE cap: it gets no marker,
  625. // so reserving marker bytes here would truncate text that was within budget.
  626. if (whole && bytes.length <= maxValueBytes) return message
  627. // Past this point the message IS being truncated, so the marker WILL be
  628. // appended and its bytes come out of the cap instead of sitting on top of it.
  629. const budget = Math.max(0, maxValueBytes - TRUNCATION_MARKER_BYTES)
  630. // Trim back to the last complete UTF-8 sequence: a cut through a multibyte
  631. // character would decode as U+FFFD — corrupting the diagnostic AND
  632. // exceeding the byte cap, since the replacement character itself encodes
  633. // to three bytes. Continuation bytes are 0b10xxxxxx; at most three of them
  634. // precede a lead byte.
  635. //
  636. // This also covers a code-unit prefix ending on a HIGH SURROGATE whose low
  637. // half sits outside it, which `Buffer.from` encodes as U+FFFD: that orphan
  638. // occupies the last three bytes of `bytes`, and `bytes` is at least
  639. // `maxValueBytes + 2` long here (one byte per retained unit, three for the
  640. // orphan), so it starts past `budget` and is always cut. Reserving the
  641. // marker is what makes that hold; cutting at `maxValueBytes` itself did not,
  642. // and needed an explicit surrogate check.
  643. let end = Math.min(budget, bytes.length)
  644. while (end > 0 && ((bytes[end] as number) & 0b1100_0000) === 0b1000_0000) end--
  645. return `${bytes.subarray(0, end).toString('utf8')}${TRUNCATION_MARKER}`
  646. }
  647. /**
  648. * Copy an fd-3 line residual into a fresh, right-sized Buffer so it no longer
  649. * shares the joined-frame allocation it was sliced from.
  650. *
  651. * After the newline loop over a `Buffer.concat` of the pending chunks, the
  652. * leftover partial line is a `subarray` VIEW onto that concat's backing store.
  653. * A view keeps the ENTIRE backing allocation alive for as long as it is
  654. * retained, so carrying the view forward as the next pending chunk would pin a
  655. * whole large frame's worth of memory behind a tiny trailing fragment — and the
  656. * `pendingBytes` counter, set to the fragment's own length, would no longer
  657. * measure the memory actually held. `Buffer.from` allocates exactly
  658. * `residual.length` bytes and copies, letting the concat allocation be
  659. * collected; an empty residual carries nothing forward.
  660. * @param residual - the leftover slice after the last newline (a view).
  661. * @returns the pending-chunk list to carry forward: `[copy]`, or `[]` when empty.
  662. */
  663. export function detachResidual(residual: Buffer): Buffer[] {
  664. return residual.length > 0 ? [Buffer.from(residual)] : []
  665. }
  666. /** One namespace after seam validation: its callables plus the optional typed-rejection contract. */
  667. interface ValidatedNamespace {
  668. functions: Record<string, CodeBindingFunction>
  669. errorClass?: CodeBindingErrorClass
  670. }
  671. /**
  672. * One in-flight run's host-side state, tracked for disposal so teardown can
  673. * fail every live run as `abort` and AWAIT each child's exit.
  674. */
  675. interface LiveRun {
  676. kill(sig: NodeJS.Signals): void
  677. settle(failure: CodeRunFailure): void
  678. finished: Promise<void>
  679. }
  680. /**
  681. * The experimental {@link CodeRuntime} backend (private, not released) registering as `codeRuntime`. Every
  682. * cap is validated config; every long-running operation honors the request's
  683. * `AbortSignal`; every disposer awaits child-process exit.
  684. */
  685. export class PythonCodeRuntime extends CodeRuntime {
  686. static Config: z<Config> = z.object({
  687. cpuSeconds: z.number().default(60),
  688. maxWallMs: z.number().default(600_000),
  689. addressSpaceMb: z.number().default(512),
  690. maxLogBytes: z.number().default(65_536),
  691. maxValueBytes: z.number().default(32_768),
  692. graceMs: z.number().default(3_000),
  693. pythonBin: z.string().default('python3'),
  694. })
  695. readonly language = 'python'
  696. readonly isolation = 'process'
  697. private readonly config: ResolvedConfig
  698. private readonly live = new Set<LiveRun>()
  699. private disposed = false
  700. /* jscpd:ignore-start -- parallel to code-runtime-worker: sibling backends keep symmetric constructor/teardown/run shapes. */
  701. constructor(ctx: Context, config: Config) {
  702. super(ctx)
  703. // Reject at load on Windows: the bootstrap imports the POSIX-only `resource`
  704. // module for RLIMIT_CPU/RLIMIT_AS, spawns with a positional fd 3, and
  705. // terminates via negative-PID process-group signals — none of which exist
  706. // on Windows. Registering ctx.codeRuntime there would let assembly succeed
  707. // and defer the failure to the first run. The asymmetry with the worker
  708. // backend is intentional: that backend is cross-platform; this one is not.
  709. if (process.platform === 'win32') {
  710. throw new Error('dsh-code-runtime-python: this backend requires a Unix platform (POSIX rlimits, fd-3 stdio, process-group signals); it cannot run on Windows')
  711. }
  712. this.config = config as ResolvedConfig
  713. for (const [key, value] of Object.entries(this.config)) {
  714. if (typeof value === 'number' && !(Number.isFinite(value) && value > 0)) {
  715. throw new Error(`dsh-code-runtime-python: config.${key} must be a positive number, got ${String(value)}`)
  716. }
  717. }
  718. // cpuSeconds crosses to the child's setrlimit(RLIMIT_CPU) raw; a float
  719. // raises TypeError inside every child (a late per-run failure). Reject it
  720. // at load. maxLogBytes/maxValueBytes get their own integer gate below (the
  721. // child int()-truncates them, so a float would diverge from the host);
  722. // maxWallMs/graceMs/addressSpaceMb are consumed as numbers where a fraction
  723. // is harmless.
  724. if (!Number.isInteger(this.config.cpuSeconds)) {
  725. throw new Error(`dsh-code-runtime-python: config.cpuSeconds must be a positive integer, got ${String(this.config.cpuSeconds)}`)
  726. }
  727. // Finite is not the same as representable as an rlimit. `cpuSeconds` and its
  728. // `+ 1` hard limit both cross to `setrlimit` as integers, and `1e100` clears
  729. // `Number.isInteger` while being far past the safe range, so it cannot round
  730. // -trip: the child sees a different number than was configured. The `+ 1` is
  731. // what gets checked because that is the larger of the two values sent.
  732. if (!Number.isSafeInteger(this.config.cpuSeconds + 1)) {
  733. throw new Error(`dsh-code-runtime-python: config.cpuSeconds must be at most ${Number.MAX_SAFE_INTEGER - 1} (it and its +1 hard limit cross to setrlimit as exact integers), got ${String(this.config.cpuSeconds)}`)
  734. }
  735. // `addressSpaceMb` is multiplied by 1 MiB before it is framed, and a large
  736. // finite value overflows to `Infinity` there — which `encodeJsonPlain`
  737. // renders as `null`, so the child receives no limit at all and every run
  738. // ends in a bootstrap exception rather than a load-time configuration error.
  739. // Checking the DERIVED byte count is what catches it; the input itself looks
  740. // ordinary. Safe-integer, not merely finite, since the value must survive
  741. // the JSON round trip exactly.
  742. if (!Number.isSafeInteger(this.config.addressSpaceMb * 1024 * 1024)) {
  743. throw new Error(`dsh-code-runtime-python: config.addressSpaceMb must be at most ${Math.floor(Number.MAX_SAFE_INTEGER / (1024 * 1024))} (its byte count crosses the wire as an exact integer), got ${String(this.config.addressSpaceMb)}`)
  744. }
  745. // `pythonBin` reaches `spawn` as the executable path, where values the
  746. // string schema admits fail late and unhelpfully. An empty string makes
  747. // `spawn` throw `ERR_INVALID_ARG_VALUE` synchronously, and an embedded NUL
  748. // throws `ERR_INVALID_ARG_TYPE` — both from inside `run()`, so the method
  749. // REJECTS instead of resolving the `worker-exit` the seam promises for a
  750. // child that cannot start. A basename with no `PATH` match would silently
  751. // fall to execvp's platform default `PATH` under the empty spawn
  752. // environment (see the resolvePythonBin JSDoc), so it is rejected here
  753. // too. All three are self-contained configuration errors that fail at
  754. // load.
  755. if (this.config.pythonBin === '' || this.config.pythonBin.includes('\0')) {
  756. throw new Error(`dsh-code-runtime-python: config.pythonBin must be a non-empty path without NUL bytes, got ${JSON.stringify(this.config.pythonBin)}`)
  757. }
  758. // An explicit path that is not an executable regular file must fail at load
  759. // like any other self-contained configuration error (the empty/NUL cases
  760. // above); a basename that is not on PATH must fail at load, not silently
  761. // fall to execvp's platform default PATH (spawn runs with an EMPTY
  762. // environment, so execvp would resolve /usr/bin:/bin and could start a
  763. // system interpreter the caller never asked for). resolvePythonBin applies
  764. // the executable-regular-file check to both forms and returns undefined for
  765. // either failure; the message distinguishes the two so the fix is obvious.
  766. const resolvedBin = resolvePythonBin(this.config.pythonBin)
  767. if (resolvedBin === undefined) {
  768. const explicit = isAbsolute(this.config.pythonBin) || this.config.pythonBin.includes('/')
  769. throw new Error(`dsh-code-runtime-python: config.pythonBin ${JSON.stringify(this.config.pythonBin)} ${explicit ? 'is not an executable regular file' : 'does not resolve on PATH'}`)
  770. }
  771. // `maxWallMs` and `graceMs` are armed with setTimeout, which clamps any
  772. // delay past MAX_TIMER_DELAY_MS to 1 ms without a word — turning a
  773. // generous ceiling into an instant timeout and a generous grace period into
  774. // an instant SIGKILL. `graceMs` is checked against the margin the
  775. // close-deadline adds on top, since that sum is what gets armed.
  776. if (this.config.maxWallMs > MAX_TIMER_DELAY_MS) {
  777. throw new Error(`dsh-code-runtime-python: config.maxWallMs must not exceed ${MAX_TIMER_DELAY_MS} (setTimeout clamps a larger delay to 1ms), got ${String(this.config.maxWallMs)}`)
  778. }
  779. if (this.config.graceMs + CLOSE_REAP_MARGIN_MS > MAX_TIMER_DELAY_MS) {
  780. throw new Error(`dsh-code-runtime-python: config.graceMs must not exceed ${MAX_TIMER_DELAY_MS - CLOSE_REAP_MARGIN_MS} (its close deadline adds ${CLOSE_REAP_MARGIN_MS}ms, and setTimeout clamps a larger delay to 1ms), got ${String(this.config.graceMs)}`)
  781. }
  782. // The output caps are budgets for a payload that has to cross fd 3 inside
  783. // one frame, and the framing ceiling is fixed. A cap above what a frame can
  784. // carry is unsatisfiable: a completion or log entry that the cap admits
  785. // arrives as an over-ceiling frame and fails the run as `worker-exit`
  786. // instead of the `output-limit` the cap describes — a silent inversion, so
  787. // it fails at load. Both budgets are metered in SERIALIZED (JSON-escaped)
  788. // bytes — the host log ledger charges the serialized cost via
  789. // `jsonStringCostUpTo`, which walks to the cap without allocating the escaped
  790. // copy, `checkDoneValue` measures the escaped form, and the producing-side
  791. // `_cap_message` in the child also caps by serialized cost (which is why a
  792. // capped diagnostic still fits its frame) — so a payload admitted under the
  793. // cap occupies at most `cap + envelope` bytes on the wire; escaping is
  794. // already inside the charge and must not be multiplied in again. The
  795. // receive-side `capMessage` backstop is the one exception to this argument:
  796. // it bills a forged `done.error.message` by RAW bytes, but that output goes
  797. // into `CodeRunResult.error.message` and never re-crosses a frame-bounded
  798. // channel, so it is not part of the wire-width bound (see its JSDoc). The
  799. // admissible cap is therefore `parse-cap - envelope`: the receive path
  800. // rejects raw frames past FRAME_PARSE_CAP_BYTES before decoding (the run
  801. // settles as a worker-exit; a hostile compact-wide-frame OOM guard), so a
  802. // budget must not exceed what an honest child's frame can actually carry
  803. // through that parser.
  804. for (const key of ['maxLogBytes', 'maxValueBytes'] as const) {
  805. // Require an integer: the child reads these budgets through `int(...)`,
  806. // which silently floors a float, so `maxLogBytes: 3.5` would truncate at 3
  807. // bytes child-side while the host meters and marks at 3.5 — the two sides
  808. // enforcing different public config. Reject the float at load, as the
  809. // worker backend does for its byte budgets.
  810. if (!Number.isInteger(this.config[key])) {
  811. throw new Error(`dsh-code-runtime-python: config.${key} must be a positive integer (the child reads it as an int, so a float diverges from the host), got ${String(this.config[key])}`)
  812. }
  813. const limit = FRAME_PARSE_CAP_BYTES - FRAME_ENVELOPE_BYTES
  814. if (this.config[key] > limit) {
  815. throw new Error(`dsh-code-runtime-python: config.${key} must not exceed ${limit} (a payload that large cannot cross the fd-3 frame PARSER, which rejects raw frames past ${FRAME_PARSE_CAP_BYTES} bytes before decoding to bound host memory — a larger budget would admit a config whose honest child frames the host then rejects as a worker-exit), got ${String(this.config[key])}`)
  816. }
  817. // Reject a log budget too small to honor: the truncation marker alone
  818. // must serialize within the budget, or a marker-only truncated run
  819. // returns more than the configured cap. (With admitted entries the
  820. // marker is envelope, so the serialized logs run to
  821. // `maxLogBytes + marker + envelope`.)
  822. if (key === 'maxLogBytes' && this.config[key] < MIN_LOG_BYTES) {
  823. throw new Error(`dsh-code-runtime-python: config.maxLogBytes must be at least ${MIN_LOG_BYTES} (a smaller budget cannot serialize the truncation marker itself, so a marker-only truncated run would return more than the configured cap), got ${String(this.config[key])}`)
  824. }
  825. }
  826. // The child builds, charges, and frames a `maxLogBytes` log entry or a
  827. // `maxValueBytes` completion value under `RLIMIT_AS`, and both paths trigger
  828. // on CHARACTER count against a serialized-BYTE budget. An astral character is
  829. // one character but four bytes of `str` storage and four UTF-8 bytes, so a
  830. // budget's worth of them peaks at three simultaneous ~4x copies (the caller's
  831. // write argument, the line slice or joined pending handed to push, and the
  832. // encode push takes to charge and ship it). A budget approaching
  833. // `addressSpaceMb` therefore makes a LEGITIMATE near-budget output breach the
  834. // address space and die as `worker-exit` instead of truncating (log) or
  835. // failing as `output-limit` (value). Metering every child write against the
  836. // address space at runtime is the wrong fix — an exact serialized-cost check
  837. // is either a full encode (the allocation being avoided) or a per-character
  838. // Python loop that burns the CPU budget — so the incompatible pair is rejected
  839. // at load: each budget times the worst-case multiple must fit the address
  840. // space. Checked on every platform, not just where `RLIMIT_AS` is enforced:
  841. // the incompatibility is a property of the config values, and the child OOMs
  842. // on a Linux deployment regardless of the host that assembled the config, so a
  843. // uniform load-time rejection is the fail-loud contract (Darwin skips only the
  844. // runtime `setrlimit`).
  845. const addressSpaceBytes = this.config.addressSpaceMb * 1024 * 1024
  846. // Room left for the peak output allocation after the interpreter's own fixed
  847. // footprint. A budget must fit MULTIPLE times over into THIS, not the whole
  848. // address space, so a budget sized right at `addressSpaceMb / MULTIPLE` — which
  849. // the multiple alone would admit — cannot leave the peak plus the interpreter
  850. // over the limit.
  851. const budgetableBytes = addressSpaceBytes - INTERPRETER_BASELINE_BYTES
  852. // The largest budget that fits: the peak (budget * MULTIPLE) must leave room,
  853. // so a budget whose peak exactly equals `budgetableBytes` is rejected — that
  854. // peak plus the reserved baseline is the whole address space, the RLIMIT_AS
  855. // edge. `ceil(budgetableBytes / MULTIPLE) - 1` is the last integer strictly
  856. // under `budgetableBytes / MULTIPLE`.
  857. // Reject a too-small address space on its own terms FIRST. Once
  858. // `budgetableBytes` is zero or negative no budget can pass, and the loop
  859. // below would report "a limit of -1" (or -2796203 at addressSpaceMb 32) while
  860. // naming `maxLogBytes` -- pointing the operator at the knob that is not the
  861. // problem. The baseline is what `addressSpaceMb` must clear here.
  862. if (budgetableBytes <= 0) {
  863. throw new Error(`dsh-code-runtime-python: config.addressSpaceMb must exceed the ${INTERPRETER_BASELINE_BYTES}-byte interpreter baseline with room for the output budgets, so the child has address space left to build and encode them; got ${String(this.config.addressSpaceMb)} MiB (${addressSpaceBytes} bytes)`)
  864. }
  865. const admissibleBudget = Math.ceil(budgetableBytes / OUTPUT_BUDGET_WORST_CASE_ADDRESS_SPACE_MULTIPLE) - 1
  866. for (const key of ['maxLogBytes', 'maxValueBytes'] as const) {
  867. if (this.config[key] * OUTPUT_BUDGET_WORST_CASE_ADDRESS_SPACE_MULTIPLE >= budgetableBytes) {
  868. throw new Error(`dsh-code-runtime-python: config.${key} times the ${OUTPUT_BUDGET_WORST_CASE_ADDRESS_SPACE_MULTIPLE}x worst-case Unicode expansion must fit within the ${budgetableBytes} bytes left after the ${INTERPRETER_BASELINE_BYTES}-byte interpreter baseline within the ${addressSpaceBytes}-byte addressSpaceMb, so a near-budget output truncates rather than breaching RLIMIT_AS as worker-exit; got ${String(this.config[key])} against a limit of ${admissibleBudget}`)
  869. }
  870. }
  871. ctx.effect(() => () => this.teardown(), 'python code-runtime teardown')
  872. }
  873. /**
  874. * Dispose to quiescence: fail every in-flight run as aborted and AWAIT each
  875. * child's exit so no subprocess that stays in the child's process group
  876. * outlives the fiber. A descendant that escaped the group with `setsid()` /
  877. * `start_new_session=True` is unreachable by `kill(-pid)` and is the documented
  878. * exception (see the package README's Known Limitations); the process-group
  879. * teardown reaps everything that stays in the group.
  880. */
  881. private async teardown(): Promise<void> {
  882. this.disposed = true
  883. const runs = [...this.live]
  884. for (const run of runs) run.settle({ kind: 'abort', message: 'runtime disposed' })
  885. // Awaiting `finished` is also what clears staging: that promise resolves
  886. // inside the run's own `settle`, which removes its directory first. So there
  887. // is deliberately no sweep here — a second pass could only ever find an
  888. // empty set, and an unreachable cleanup path is worse than none, since it
  889. // reads as the real guarantee while never running.
  890. await Promise.all(runs.map(run => run.finished))
  891. }
  892. /**
  893. * Execute one program in a fresh Python subprocess. Success resolves with
  894. * `result.value` (and no `result.error`); failure — parse failure, thrown
  895. * exception, invalid completion, output overflow, budget expiry, abort, or
  896. * substrate death — resolves with `result.error` set (classified by
  897. * `CodeRunFailure.kind`). The method rejects only for seam misuse.
  898. */
  899. async run(request: CodeRunRequest): Promise<CodeRunResult> {
  900. if (this.disposed) throw new Error('dsh-code-runtime-python: run() after disposal')
  901. const bindings = this.validateBindings(request)
  902. if (request.signal?.aborted) {
  903. return { logs: [], error: { kind: 'abort', message: messageOf(request.signal.reason) } }
  904. }
  905. let bootstrapPath: string
  906. try {
  907. // The interpreter is an external process, so the entry script has to sit
  908. // on the real filesystem; see materializePyScripts. One copy PER RUN,
  909. // synchronously, so no async boundary opens before `execute` registers the
  910. // run and installs the abort listener.
  911. bootstrapPath = materializePyScripts()
  912. } catch (error: unknown) {
  913. // A full or read-only temp filesystem, or a packaged asset the deployment
  914. // failed to ship, is a SUBSTRATE failure — the same class as a child that
  915. // cannot start. The seam permits rejection only for misuse, so this
  916. // resolves as `worker-exit` rather than throwing out of `run()`.
  917. return { logs: [], error: { kind: 'worker-exit', message: `failed to stage the python bootstrap: ${messageOf(error)}` } }
  918. }
  919. return await this.execute(request, bindings, bootstrapPath)
  920. }
  921. /* jscpd:ignore-end */
  922. /**
  923. * Reject (seam misuse) malformed binding namespaces: non-identifier or
  924. * reserved globals/error classes, duplicates, and colliding or
  925. * runtime-owned injected globals.
  926. */
  927. private validateBindings(request: CodeRunRequest): Map<string, ValidatedNamespace> {
  928. const bindings = new Map<string, ValidatedNamespace>()
  929. // Every name the bootstrap injects into the program's one global namespace:
  930. // namespace globals plus error-class names. They must be a collision-free
  931. // set that avoids the runtime's own slots, or a later injection silently
  932. // overwrites an earlier one (or the completion/builtins slot) and the run
  933. // fails obscurely at execution time.
  934. const injectedGlobals = new Set<string>()
  935. const claimGlobal = (name: string, role: string): void => {
  936. if (RUNTIME_OWNED_GLOBALS.has(name)) {
  937. throw new Error(`dsh-code-runtime-python: ${role} ${JSON.stringify(name)} collides with a runtime-owned global`)
  938. }
  939. if (injectedGlobals.has(name)) {
  940. throw new Error(`dsh-code-runtime-python: ${role} ${JSON.stringify(name)} collides with another injected global`)
  941. }
  942. injectedGlobals.add(name)
  943. }
  944. for (const namespace of request.bindings) {
  945. // Snapshot the caller-supplied fields into plain values ONCE. The
  946. // namespace and errorClass objects may expose `global`/`name`/
  947. // `memberNameProperty` through getters: validation reads each several
  948. // times, and the ORIGINAL errorClass object would otherwise be retained
  949. // for the boot frame, whose JSON.stringify re-reads it after validation.
  950. // A getter that changes or throws on a later read would turn the
  951. // seam-misuse rejection into a worker-exit (or inject a different name
  952. // than validation approved); reading each field once here and keeping
  953. // the plain copy makes validation and the boot frame agree.
  954. const global = namespace.global
  955. if (!IDENTIFIER.test(global) || RESERVED_NAMES.has(global)) {
  956. throw new Error(`dsh-code-runtime-python: binding global ${JSON.stringify(global)} is not a usable Python identifier`)
  957. }
  958. if (bindings.has(global)) {
  959. throw new Error(`dsh-code-runtime-python: duplicate binding global ${JSON.stringify(global)}`)
  960. }
  961. claimGlobal(global, 'binding global')
  962. // The error class becomes a program global and its member property an
  963. // attribute name, so both face the Python identifier rules; the member
  964. // additionally must be assignable on a BaseException instance.
  965. const errorClass = namespace.errorClass
  966. let validatedErrorClass: CodeBindingErrorClass | undefined
  967. if (errorClass) {
  968. const name = errorClass.name
  969. const memberNameProperty = errorClass.memberNameProperty
  970. if (!IDENTIFIER.test(name) || RESERVED_NAMES.has(name)) {
  971. throw new Error(`dsh-code-runtime-python: errorClass.name ${JSON.stringify(name)} is not a usable Python identifier`)
  972. }
  973. // Any non-empty own attribute name is settable via setattr (the
  974. // program reads exotic names like `tool-name` with getattr), matching
  975. // the seam contract and the worker backend — only the seam-excluded
  976. // and protocol-reserved members below are refused.
  977. if (memberNameProperty.length === 0) {
  978. throw new Error('dsh-code-runtime-python: errorClass.memberNameProperty must be a non-empty attribute name')
  979. }
  980. if (EXCEPTION_RESERVED_MEMBERS.has(memberNameProperty) || DUNDER.test(memberNameProperty)) {
  981. throw new Error(`dsh-code-runtime-python: errorClass.memberNameProperty ${JSON.stringify(memberNameProperty)} is a reserved error member and cannot be assigned`)
  982. }
  983. claimGlobal(name, 'errorClass.name')
  984. validatedErrorClass = { name, memberNameProperty }
  985. }
  986. // Snapshot the callables into a plain own-property record before the
  987. // child can dispatch. `namespace.functions` is caller-supplied, so it may
  988. // expose members through getters or a Proxy; reading one of them inside
  989. // the fd-3 `data` callback would throw OUTSIDE the dispatcher's try and
  990. // terminate the host (defensive-patterns contain-callback-exceptions).
  991. // Reading every member here, in run()'s synchronous validation segment,
  992. // turns that throw into the seam-misuse rejection run() reserves for
  993. // malformed bindings. The snapshot is also the single key set the boot
  994. // frame advertises AND dispatch reads, so a getter whose keys differ
  995. // between reads cannot desynchronize the child's allowed names from what
  996. // the host will actually call. The record is null-prototype: the seam
  997. // contract treats member names like `__proto__` or `constructor` as
  998. // ordinary own properties, and a plain `{}` assignment of `__proto__`
  999. // would hit the prototype setter instead of creating the own property.
  1000. const functions = Object.create(null) as Record<string, CodeBindingFunction>
  1001. for (const name of Object.keys(namespace.functions)) {
  1002. // Only callables enter the snapshot: a getter exposing a non-function
  1003. // member would otherwise assign a value the dispatcher's `typeof fn
  1004. // !== 'function'` check rejects anyway, and keeping it out of the
  1005. // snapshot keeps the boot frame's name list and the dispatch key set
  1006. // one and the same.
  1007. const fn = namespace.functions[name]
  1008. if (typeof fn === 'function') functions[name] = fn
  1009. }
  1010. bindings.set(global, { functions, ...validatedErrorClass ? { errorClass: validatedErrorClass } : {} })
  1011. }
  1012. return bindings
  1013. }
  1014. /** Spawn the child for one validated run and drive it to settlement. */
  1015. private execute(
  1016. request: CodeRunRequest,
  1017. bindings: Map<string, ValidatedNamespace>,
  1018. bootstrapPath: string,
  1019. ): Promise<CodeRunResult> {
  1020. // This run's own staging directory, removed at settlement.
  1021. const bootstrapDir = dirname(bootstrapPath)
  1022. // Explicit pipe count of 4 puts the framed-JSON channel at fd 3 in the child.
  1023. // Resolve the interpreter against the current PATH first: the child's empty
  1024. // env would otherwise strip PATH and miss a basename python3 (see resolvePythonBin).
  1025. // `spawn` can throw SYNCHRONOUSLY — a descriptor-exhausted host (EMFILE) or a
  1026. // libuv-level failure surfaces here, before the Promise executor and its
  1027. // settlement path exist. Left uncaught it would REJECT run() (the seam
  1028. // permits rejection only for misuse) and strand this run's staging directory,
  1029. // which only settle() removes. Catch it, unlink the directory, and resolve a
  1030. // `worker-exit` — the same class as the async ENOENT `error` event below.
  1031. let child: ChildProcessWithoutNullStreams
  1032. let proto: Duplex | null
  1033. try {
  1034. // `-u` keeps the interpreter's own stdout/stderr UNBUFFERED: a program
  1035. // that writes through `sys.__stdout__`/`sys.__stderr__` (or C-stdio
  1036. // layered on the same fds) must have those bytes visible to the host's
  1037. // stray capture immediately — a block-buffered wrapper would otherwise
  1038. // hold them until an explicit flush, and the host SIGTERMs the child
  1039. // right after the done frame, before any finalization-time flush could
  1040. // run. The `_LogStream` replacement of `sys.stdout`/`sys.stderr` is
  1041. // unaffected (it is a Python object, not the C-level stdio buffer).
  1042. // Load validated that the configured interpreter resolves to an
  1043. // executable regular file (basename through PATH, explicit path
  1044. // directly). The type assertion is the load-time contract (see the
  1045. // pythonBin load checks); a PATH change between load and run would make
  1046. // this undefined and spawn throws synchronously, which the surrounding
  1047. // try settles as worker-exit like any other spawn failure.
  1048. const resolvedPythonBin = resolvePythonBin(this.config.pythonBin) as string
  1049. child = spawn(resolvedPythonBin, ['-u', '-I', bootstrapPath], {
  1050. env: {},
  1051. detached: true, // Own process group — kill(-pid, sig) reaches subprocesses the model program spawns.
  1052. stdio: ['pipe', 'pipe', 'pipe', 'pipe'],
  1053. })
  1054. // Fd 3 is a duplex pipe carrying protocol frames. Node types extra stdio
  1055. // entries as `Stream | null`; the runtime shape with `'pipe'` is a duplex,
  1056. // so we narrow at the boundary rather than smearing casts below. Stdout
  1057. // and stderr are guaranteed non-null under `'pipe'` and typed as such.
  1058. proto = child.stdio[3] as Duplex | null
  1059. /* v8 ignore next 3 -- `'pipe'` stdio always populates fd 3; guarding Node's `Stream | null` typing widening. */
  1060. if (proto === null) {
  1061. throw new Error('dsh-code-runtime-python: python subprocess spawned without a fd-3 pipe')
  1062. }
  1063. // Close the host's stdin write handle immediately: the program is an
  1064. // async body that reads nothing from fd 0, and a live pipe here would
  1065. // hold a host-side handle open past the run — a setsid-escaped descendant
  1066. // inheriting fd 0 would keep the host process from exiting even after the
  1067. // closeDeadline forced settlement. The child (and any descendant) reads
  1068. // EOF on fd 0 instead, and no host handle survives.
  1069. // oxlint-disable-next-line typescript/no-unnecessary-condition -- the boot-write-failure fake child has no stdin.
  1070. child.stdin?.destroy()
  1071. } catch (error: unknown) {
  1072. try {
  1073. rmSync(bootstrapDir, { recursive: true, force: true })
  1074. } catch {
  1075. // Same swallow as settle()'s removal: `force` already absorbs a missing
  1076. // directory, so only a filesystem-level refusal reaches here, and the
  1077. // staging copy holds nothing but two checked-in scripts.
  1078. }
  1079. return Promise.resolve({ logs: [], error: { kind: 'worker-exit' as const, message: `python spawn error: ${messageOf(error)}` } })
  1080. }
  1081. return new Promise<CodeRunResult>((resolve) => {
  1082. let settled = false
  1083. const logs: string[] = []
  1084. // An unterminated line flushed with the `open` flag: the next log frame
  1085. // appends to it (no fake newline between entries), and finish() pushes
  1086. // the residual if the run ends with it still open. Held as a fragment
  1087. // ARRAY, so k tiny open frames cost O(k) — re-joining and re-walking the
  1088. // whole held text per frame would be O(k * budget).
  1089. let openParts: string[] = []
  1090. // Past MAX_PENDING_CHUNKS, the held fragments are coalesced into sealed
  1091. // blocks (mirroring the fd-3 reader's `blocks` and the stray capture's
  1092. // seal): each fragment is a distinct array slot plus string object
  1093. // header — ~30x overhead the byte cap cannot see — so a budget-sized
  1094. // single-character open flood would otherwise accumulate thousands of
  1095. // slots. Sealing bounds the live fragment count exactly like the
  1096. // sibling paths; the merge reads sealed + current fragments. A block
  1097. // ARRAY (not one repeated string concat) matches the sibling shape and
  1098. // avoids depending on V8 ConsString amortization.
  1099. let openSealed: string[] = []
  1100. // Every truncation arm funnels here: the committed open prefix was
  1101. // ALREADY billed, so it is pushed BEFORE the marker — a flushed line is
  1102. // never lost (only the marker stays last), and no ledger re-charge
  1103. // happens. openParts is emptied here, so no later arm or finish() sees
  1104. // it.
  1105. const truncateLogs = (): void => {
  1106. logsTruncated = true
  1107. if (openSealed.length > 0 || openParts.length > 0) {
  1108. logs.push(openSealed.join('') + openParts.join(''))
  1109. openSealed = []
  1110. openParts = []
  1111. }
  1112. logs.push(logTruncationMarker(this.config.maxLogBytes))
  1113. clearStray(strayOut)
  1114. clearStray(strayErr)
  1115. }
  1116. // One host-side ledger covers normal frames, forged frames, and stray stdout bytes.
  1117. // The ledger starts one byte below maxLogBytes: each entry is charged its
  1118. // JSON-string cost plus one separator byte, and the serialized outer logs
  1119. // array adds one more byte of envelope (two brackets and n-1 commas over n
  1120. // entries' separators), so a result that exactly exhausts the ledger
  1121. // serializes to exactly maxLogBytes; WITHOUT the reserved byte it would
  1122. // serialize to maxLogBytes + 1. Reserving that byte keeps an admitted
  1123. // result within the configured cap; the truncation-marker entry is
  1124. // envelope, not payload, and rides uncharged.
  1125. let logBudget = this.config.maxLogBytes - 1
  1126. let logsTruncated = false
  1127. // Drop a pipe's buffered stray output wholesale: once the ledger has
  1128. // truncated, every byte of it would be no-op'd by admit(), so retaining
  1129. // it (and later Buffer.concat+decoding it in flushStray) would spend host
  1130. // memory on output that can never be admitted. Called from every arm that
  1131. // marks the ledger truncated — admit()'s two ceilings and the child-marker
  1132. // frame arm — so the end-path flushStray sees empty buffers and exits.
  1133. const clearStray = (stray: StrayBuffer): void => {
  1134. stray.chunks = []
  1135. stray.blocks = []
  1136. stray.cost = 0
  1137. stray.utf8 = { expected: 0, width: 0, lowerFirst: 0, upperFirst: 0 }
  1138. }
  1139. const admit = (text: string): void => {
  1140. // Post-truncation admits are no-ops: once the ledger has truncated, the
  1141. // marker is the last entry. Reachable within one `data` callback — a
  1142. // chunk carrying two newline-terminated lines where the first exhausts
  1143. // the budget hits this on the second — so it is a measured branch.
  1144. if (logsTruncated) return
  1145. // Each entry is charged its SERIALIZED cost — JSON.stringify's quotes
  1146. // and escapes plus one separator byte — because the seam bounds the
  1147. // serialized outer logs payload, and control characters expand
  1148. // several-fold under JSON escaping (a "\x00" flood would otherwise
  1149. // admit 6x its charge). The charge also puts a floor under an empty
  1150. // entry (its two quotes plus separator), so a `while True: print()`
  1151. // flood of zero-byte lines exhausts the ledger instead of growing the
  1152. // retained array without ever touching the budget. The one fixed
  1153. // truncation-marker entry is envelope, not payload, and rides
  1154. // uncharged.
  1155. //
  1156. // Cheap lower bound FIRST, before the escaped copy exists: every
  1157. // UTF-16 code unit costs at least one serialized byte (an ASCII
  1158. // character is one byte; a control character is six as `\uXXXX`; a
  1159. // non-ASCII BMP character is two or three; each half of a surrogate
  1160. // pair contributes two of the four bytes its code point encodes to),
  1161. // and the JSON form adds two quotes on top of the separator byte. So
  1162. // `text.length + 3` never exceeds the true cost, and a forged `log`
  1163. // frame carrying a control-heavy string anywhere below the 64 MiB
  1164. // frame parse cap truncates here instead of allocating a
  1165. // hundreds-of-megabytes escaped copy under a small maxLogBytes.
  1166. if (text.length + 3 > logBudget) {
  1167. // Release the buffered stray pipes: their bytes can never be
  1168. // admitted now (see clearStray).
  1169. truncateLogs()
  1170. return
  1171. }
  1172. // Past the lower bound, measure the exact serialized cost without
  1173. // allocating the escaped copy: `jsonStringCostUpTo` walks to the cap and
  1174. // stops, so even a near-budget control-char-dense line never materializes
  1175. // a sixfold-inflated `JSON.stringify` result. `+ 1` for the separator.
  1176. const measured = jsonStringCostUpTo(text, logBudget - 1)
  1177. if (measured === undefined) {
  1178. truncateLogs()
  1179. return
  1180. }
  1181. logBudget -= measured + 1
  1182. logs.push(text)
  1183. }
  1184. // Stray-byte capture: anything the child writes to its stdout/stderr
  1185. // (native prints, C-extension writes) still counts against the ledger.
  1186. //
  1187. // Output is admitted per LINE, not per transport chunk. `logs` entries
  1188. // are joined with `\n` downstream (Code Mode), so each entry must be one
  1189. // line: pushing a raw `data` chunk would turn every arbitrary pipe-read
  1190. // boundary into a model-visible newline, so a single 200 KiB native write
  1191. // split across pipe reads would read back with spurious line breaks. The
  1192. // child's own `log` frames are already line-granular; stray capture
  1193. // matches them by splitting on `\n`.
  1194. //
  1195. // Buffered as raw `Buffer` chunks with a running SERIALIZED-cost counter,
  1196. // exactly like the fd-3 reader below and for the same reasons: a string
  1197. // `+=` accumulator re-copies the whole residual on every pipe chunk
  1198. // (quadratic on a large newline-free write), and scanning it from index 0
  1199. // each chunk is a second quadratic. Appending a chunk is O(1); the split
  1200. // happens only when a `\n` actually arrived. A newline never appears inside
  1201. // a UTF-8 multibyte sequence (continuation bytes are 0x80–0xBF), so
  1202. // splitting on the raw 0x0a byte and decoding each complete line is safe
  1203. // without a streaming decoder — a line's bytes are whole by construction.
  1204. //
  1205. // `chunks` also seals into `blocks` past MAX_PENDING_CHUNKS, mirroring the
  1206. // fd-3 reader: without it a program pacing one-byte newline-free
  1207. // `os.write`s accumulates one Buffer object per write, and the object plus
  1208. // backing-store overhead — which no byte or cost count sees — exhausts the
  1209. // host heap far below the budget. Sealing bounds the live object count.
  1210. interface StrayBuffer { chunks: Buffer[]; blocks: Buffer[]; cost: number; utf8: Utf8CostState }
  1211. const strayOut: StrayBuffer = { chunks: [], blocks: [], cost: 0, utf8: { expected: 0, width: 0, lowerFirst: 0, upperFirst: 0 } }
  1212. const strayErr: StrayBuffer = { chunks: [], blocks: [], cost: 0, utf8: { expected: 0, width: 0, lowerFirst: 0, upperFirst: 0 } }
  1213. const captureStray = (stray: StrayBuffer, chunk: Buffer): void => {
  1214. // Once the ledger has truncated, stop buffering: admit() is a no-op past
  1215. // that point, so continuing to accumulate would retain host memory for
  1216. // output that can never be admitted.
  1217. if (logsTruncated) return
  1218. stray.chunks.push(chunk)
  1219. // Track SERIALIZED cost, not raw bytes: a control-char-dense residual
  1220. // (a NUL or illegal-UTF-8 flood) serializes several-fold, so a raw-byte
  1221. // threshold would let it grow to the full budget's worth of RAW bytes
  1222. // before flushing. `accrueStrayCost` decodes UTF-8 structurally across
  1223. // chunks (via `stray.utf8`) so a byte that renders as U+FFFD is charged
  1224. // its three serialized bytes, not one.
  1225. stray.cost += accrueStrayCost(chunk, stray.utf8)
  1226. // Bound the live fragment count (see the seal rationale above), before
  1227. // any concat so an over-count payload is never copied whole first.
  1228. if (stray.chunks.length >= MAX_PENDING_CHUNKS) {
  1229. stray.blocks.push(Buffer.concat(stray.chunks))
  1230. stray.chunks = []
  1231. }
  1232. if (chunk.includes(0x0a)) {
  1233. let buffered = Buffer.concat(stray.blocks.length > 0 ? [...stray.blocks, ...stray.chunks] : stray.chunks)
  1234. stray.blocks = []
  1235. let newline: number
  1236. while ((newline = buffered.indexOf(0x0a)) >= 0) {
  1237. admit(buffered.subarray(0, newline).toString('utf8'))
  1238. buffered = buffered.subarray(newline + 1)
  1239. }
  1240. // Carry the residual as a fresh right-sized copy, not the subarray view
  1241. // (which would pin the whole concat allocation). See detachResidual.
  1242. // The residual begins at a character boundary (a newline is never
  1243. // inside a multibyte sequence), so its cost and UTF-8 state recompute
  1244. // cleanly from a fresh walk.
  1245. // A line admitted inside the loop may have exhausted the ledger and
  1246. // cleared this pipe (see clearStray); the re-retain below must not
  1247. // resurrect the doomed residual.
  1248. // oxlint-disable-next-line typescript/no-unnecessary-condition -- admit() (a closure) sets it.
  1249. if (logsTruncated) return
  1250. stray.chunks = detachResidual(buffered)
  1251. stray.utf8 = { expected: 0, width: 0, lowerFirst: 0, upperFirst: 0 }
  1252. stray.cost = accrueStrayCost(buffered, stray.utf8)
  1253. }
  1254. // Newline-free residual is bounded by the ledger, not left to grow with
  1255. // the stream: an `os.write(1, b"A"*N)` flood carrying no newline would
  1256. // otherwise accumulate N bytes in host memory before `end`. The bound is
  1257. // on the COMBINED pending cost of both pipes, not each alone: stdout and
  1258. // stderr share one `logBudget`, so checking each against the full budget
  1259. // independently would let both retain nearly a budget's worth at once —
  1260. // ~2x peak, up to ~512 MiB near the ceiling — before either flushed.
  1261. // When the sum would cross the budget, flush both now. admit() charges
  1262. // the exact serialized cost, truncates, and marks the ledger, and the
  1263. // truncation short-circuit above stops buffering on the next chunk.
  1264. // `+ 3` covers the two quotes and one separator admit adds. The two
  1265. // pipes are independent OS streams whose `data` events already interleave
  1266. // nondeterministically with each other and with the child's own fd-3
  1267. // `log` frames, so `logs` carries no cross-pipe ordering guarantee to
  1268. // preserve here; a fixed drain order is as valid as any.
  1269. // Flushing is NOT a stream end: a multibyte UTF-8 character can be split
  1270. // across pipe `data` chunks, so the residual may end mid-sequence. A
  1271. // budget-triggered flush must decode only the complete prefix and carry
  1272. // the incomplete tail forward (≤3 bytes) on the same pipe's residual —
  1273. // decoding it here would render a legal character as U+FFFD in a released
  1274. // entry (see `flushStray`). This is unlike the `end`/closeDeadline paths
  1275. // below, where a trailing incomplete sequence is genuinely truncated input
  1276. // and U+FFFD is honest.
  1277. if (strayOut.cost + strayErr.cost + 3 > logBudget) {
  1278. flushStray(strayOut, true)
  1279. flushStray(strayErr, true)
  1280. }
  1281. }
  1282. // Flush a pipe's residual into `logs`. Called on the combined-budget
  1283. // threshold above, on the pipe's `end` (normal drain), and — for the
  1284. // setsid-escapee path where destroy() forces settlement without an `end` —
  1285. // explicitly in the closeDeadline handler. Idempotent: it clears what it
  1286. // admits, so a later flush is a no-op, and it returns early on an empty
  1287. // buffer so flushing the sibling that had nothing pending is a no-op. The
  1288. // `chunks`/`blocks` guard is the only emptiness check needed — `data` never
  1289. // emits a zero-length Buffer, so a non-empty fragment list always decodes
  1290. // to a non-empty tail.
  1291. //
  1292. // `retainPartialTail` is true only on the budget-triggered path: there the
  1293. // residual can end at an ARBITRARY pipe boundary, so if the incomplete
  1294. // trailing bytes of a UTF-8 lead sequence are pending (`stray.utf8.expected
  1295. // > 0`), they are withheld from the decode and re-carried on `chunks` for a
  1296. // later chunk to complete — decoding them here would render a LEGAL,
  1297. // un-finished character as U+FFFD in an admitted entry, and the next chunk's
  1298. // bytes would then each independently break into more U+FFFD. The withheld
  1299. // tail is `stray.utf8.width - stray.utf8.expected` bytes (the lead plus the
  1300. // continuations consumed so far), at most 3; `stray.utf8` is reset and the
  1301. // withheld tail re-accrued so the next chunk continues the walk correctly.
  1302. // The `end`/closeDeadline paths pass `false`: there a trailing incomplete
  1303. // sequence is real truncated input and the U+FFFD is the honest render.
  1304. function flushStray(stray: StrayBuffer, retainPartialTail?: boolean): void {
  1305. if (stray.chunks.length === 0 && stray.blocks.length === 0) return
  1306. // Concatenate the sealed blocks and the current-chunk residual together
  1307. // unconditionally (no `blocks.length > 0` ternary): a flush can run with
  1308. // either or both present, and a branch on their presence would need a
  1309. // test that flushes exactly at a seal boundary.
  1310. let full = Buffer.concat([...stray.blocks, ...stray.chunks])
  1311. // A budget flush landing exactly between a lead byte and its
  1312. // still-pending continuation requires the combined-cost threshold to trip
  1313. // on a specific mid-multibyte pipe boundary — not deterministically
  1314. // schedulable through the black-box seam, which observes only complete
  1315. // entries. So the retention arm is v8-ignored (exercised by review
  1316. // reasoning over the `stray.utf8` state, not by an in-tree test): it
  1317. // withholds the lead-plus-consumed-continuations tail (≤3 bytes, via
  1318. // `stray.utf8.width - stray.utf8.expected`) from the decode, re-carries it
  1319. // for a later chunk, and re-accrues the pipe's cost/UTF-8 state over it;
  1320. // decoding here would render a LEGAL, unfinished character as U+FFFD in an
  1321. // admitted entry. Every retainPartialTail=false call (the `end`/closeDeadline
  1322. // paths) and a budget flush with no partial tail in flight (`expected === 0`)
  1323. // falls through with `keep` unset: the FULL residual is decoded — there a
  1324. // trailing incomplete sequence is real truncated input and the U+FFFD is the
  1325. // honest render.
  1326. let keep: Buffer | undefined
  1327. /* v8 ignore next 18 -- mid-sequence budget-flush boundary is not schedulable from a test. */
  1328. if (retainPartialTail && stray.utf8.expected > 0) {
  1329. const drop = Math.min(stray.utf8.width - stray.utf8.expected, full.length)
  1330. keep = full.subarray(full.length - drop)
  1331. full = full.subarray(0, full.length - drop)
  1332. stray.chunks = detachResidual(keep)
  1333. // Re-accrue the withheld tail from a FRESH state: `stray.utf8` still
  1334. // holds the whole-pending state (`expected > 0`, i.e. the tail is
  1335. // mid-sequence), so metering `keep` against it would charge the carried
  1336. // LEAD byte as an illegal continuation. Reset, then walk `keep` so the
  1337. // resumed sequence re-claims its own lead.
  1338. stray.utf8 = { expected: 0, width: 0, lowerFirst: 0, upperFirst: 0 }
  1339. stray.cost = accrueStrayCost(keep, stray.utf8)
  1340. stray.blocks = []
  1341. // Do not admit an EMPTY entry: when the whole residual is a single
  1342. // unfinished multibyte sequence, `full` was drained into `keep` and no
  1343. // complete byte stream remains to admit. `admit('')` would push a
  1344. // model-visible bogus empty line (logs are joined with '\n' downstream).
  1345. if (full.length > 0) admit(full.toString('utf8'))
  1346. } else {
  1347. stray.chunks = []
  1348. stray.cost = 0
  1349. stray.utf8 = { expected: 0, width: 0, lowerFirst: 0, upperFirst: 0 }
  1350. stray.blocks = []
  1351. admit(full.toString('utf8'))
  1352. }
  1353. }
  1354. child.stdout.on('data', (chunk: Buffer) => { captureStray(strayOut, chunk) })
  1355. child.stderr.on('data', (chunk: Buffer) => { captureStray(strayErr, chunk) })
  1356. child.stdout.on('end', () => { flushStray(strayOut) })
  1357. child.stderr.on('end', () => { flushStray(strayErr) })
  1358. // Line-framed JSON reader over fd 3. The unframed buffer is bounded: a
  1359. // hostile program can loop `os.write(3, b"A"*4096)` with no newline to
  1360. // exhaust HOST memory, which the child's RLIMIT_AS does not cover. It is
  1361. // a memory-safety bound only: legitimate `call` frames may be large
  1362. // (binding traffic has no seam byte cap), so it never keys off
  1363. // maxValueBytes.
  1364. // Buffered as raw chunks with a running byte counter: appending is O(1)
  1365. // per chunk (a string `+=` accumulator would re-copy the whole prefix on
  1366. // every pipe chunk — quadratic on a large frame), joins happen only when
  1367. // a newline actually arrived, and the ceiling check reads the counter.
  1368. let pendingChunks: Buffer[] = []
  1369. // Fragments already merged into finished blocks. Kept separate from
  1370. // `pendingChunks` so sealing never re-copies what earlier seals produced;
  1371. // the two together are the unframed buffer, and `pendingBytes` counts both.
  1372. let sealedBlocks: Buffer[] = []
  1373. let pendingBytes = 0
  1374. proto.on('data', (chunk: Buffer) => {
  1375. // Once settled, stop accumulating: a hostile child that keeps flooding
  1376. // fd 3 between finish() and close must not regrow the host buffer.
  1377. /* v8 ignore next -- post-settlement data needs the child to outrace close after we decided. */
  1378. if (settled) return
  1379. pendingChunks.push(chunk)
  1380. pendingBytes += chunk.length
  1381. // Check the counter BEFORE the join, not the joined line afterwards:
  1382. // Buffer.concat allocates a second copy of everything held, so a line
  1383. // measured after the concat had already cost twice the ceiling — the
  1384. // ceiling this check exists to enforce. The counter is exact and free,
  1385. // and the retained chunks are released here so the rejected payload is
  1386. // not still held while the run settles.
  1387. //
  1388. // The counter charges the whole unframed buffer, which over-counts by at
  1389. // most the newline-bearing chunk's own length (one pipe read): the
  1390. // residual carried in is always a partial line, so nothing but the
  1391. // current line can be larger than that. That over-count is deliberate and
  1392. // load-bounded on the OTHER side: the config cap is `parse-cap - envelope`,
  1393. // and a legitimate near-cap frame plus a following chunk's leading bytes
  1394. // could in principle nudge the counter over the cap for one read window
  1395. // — but only when maxLogBytes/maxValueBytes is configured within one
  1396. // pipe read of the 64 MiB cap, orders of magnitude past the 32/64 KiB
  1397. // defaults.
  1398. //
  1399. // The cap is enforced ONLY when the held bytes are still a single
  1400. // unframed line (this chunk carries no newline, and earlier
  1401. // newline-bearing chunks were joined immediately): a frame past the cap
  1402. // would otherwise be fully `Buffer.concat`-ed (a second copy of its
  1403. // bytes) and only then dropped in the line loop — the peak-memory
  1404. // doubling this pre-concat check exists to prevent. Dropping the
  1405. // oversized unframed buffer before the join keeps the peak at one copy
  1406. // of the wire bytes. When this chunk DOES carry a newline the buffer
  1407. // holds several frames, so the FIRST-FRAME check below (not this
  1408. // counter, which charges them all) decides.
  1409. if (pendingBytes > FRAME_PARSE_CAP_BYTES && !chunk.includes(0x0a)) {
  1410. pendingChunks = []
  1411. sealedBlocks = []
  1412. pendingBytes = 0
  1413. finish({ error: { kind: 'worker-exit', message: `protocol frame exceeded ${FRAME_PARSE_CAP_BYTES} bytes on fd 3` } })
  1414. return
  1415. }
  1416. // Bound the FRAGMENT COUNT as well as the byte total, but only AFTER the
  1417. // ceiling check above: sealing first would `Buffer.concat` an already
  1418. // over-ceiling payload and allocate a second copy of it before the
  1419. // rejection ran, which is the peak-memory doubling that check exists to
  1420. // prevent.
  1421. //
  1422. // Fragment count needs its own bound because the ceiling meters payload
  1423. // bytes only, while each retained chunk is a separate Buffer with object
  1424. // and backing-store overhead no byte count sees: 5000 single-byte
  1425. // newline-free writes produced 5000 chunks holding 5031 bytes, so a
  1426. // program pacing such writes could accumulate millions of objects inside
  1427. // the wall budget and exhaust the host heap far below the ceiling.
  1428. //
  1429. // Sealing appends to a list of finished blocks instead of re-merging
  1430. // everything held. Concatenating the whole buffer at each threshold
  1431. // re-copied the entire accumulated prefix every time, so the cumulative
  1432. // copy volume was quadratic, not the amortized O(1) an earlier revision
  1433. // of this comment claimed: 10 MiB trickled a byte at a time copies
  1434. // 53.7 GB that way, and 64 MiB copies 2.2 TB. Here each byte is copied
  1435. // once into its block and never again, so the total stays linear, and the
  1436. // block list is itself bounded — every block holds at least
  1437. // `MAX_PENDING_CHUNKS - 1` bytes, so reaching the 64 MiB cap admits
  1438. // at most a few hundred thousand of them.
  1439. // Sealing runs ONLY on a newline-free chunk, and after the newline
  1440. // branch below: a chunk carrying a newline must reach the join (and its
  1441. // first-frame check) rather than being sealed into a block the check
  1442. // would then not scan for newlines. That keeps the invariant
  1443. // `sealedBlocks hold newline-free prefixes only` true, so the
  1444. // first-frame scan below can charge each sealed block's whole length
  1445. // toward the first frame without missing a newline inside it.
  1446. if (chunk.includes(0x0a)) {
  1447. // First-FRAME check before the join: measure the bytes up to the
  1448. // first newline across the held chunks. The byte counter cannot
  1449. // serve here — it charges the whole buffer, which legitimately
  1450. // holds several frames each within the cap. A first frame past the
  1451. // cap is dropped before the join (one copy of its wire bytes);
  1452. // later frames in the same buffer are handled line by line in the
  1453. // loop below.
  1454. let firstFrameLen = 0
  1455. let sawNewline = false
  1456. // Sealed blocks hold newline-free prefixes only (see the sealing
  1457. // gate below), so they are entirely part of the first frame.
  1458. for (const b of sealedBlocks) firstFrameLen += b.length
  1459. for (const c of pendingChunks) {
  1460. const nl = c.indexOf(0x0a)
  1461. if (nl >= 0) {
  1462. firstFrameLen += nl
  1463. sawNewline = true
  1464. break
  1465. }
  1466. firstFrameLen += c.length
  1467. }
  1468. if (sawNewline && firstFrameLen > FRAME_PARSE_CAP_BYTES) {
  1469. pendingChunks = []
  1470. sealedBlocks = []
  1471. pendingBytes = 0
  1472. finish({ error: { kind: 'worker-exit', message: `protocol frame exceeded ${FRAME_PARSE_CAP_BYTES} bytes on fd 3` } })
  1473. return
  1474. }
  1475. let buffered = Buffer.concat(sealedBlocks.length > 0 ? [...sealedBlocks, ...pendingChunks] : pendingChunks)
  1476. sealedBlocks = []
  1477. let newline: number
  1478. while ((newline = buffered.indexOf(0x0a)) >= 0) {
  1479. const line = buffered.subarray(0, newline)
  1480. buffered = buffered.subarray(newline + 1)
  1481. /* v8 ignore next -- an empty line comes only from a forged `\n\n` write. */
  1482. if (line.length === 0) continue
  1483. // No per-line cap check here: the pre-join counter (single unframed
  1484. // line) and the first-frame check (newline-bearing chunk) above
  1485. // reject any frame past FRAME_PARSE_CAP_BYTES before this join, so
  1486. // every line in this loop is within the cap by construction — a
  1487. // per-line check would be dead code.
  1488. // `toString('utf8')` would silently REPLACE illegal bytes with
  1489. // U+FFFD, corrupting a completion or binding payload a forged
  1490. // frame smuggled in (the honest child's lossless encoder never
  1491. // emits non-UTF-8, so such a frame is hostile traffic). The fatal
  1492. // decode throws on them and the frame is dropped — not accepted
  1493. // with a mangled value — the same treatment as the unsafe-integer
  1494. // check below.
  1495. let text: string
  1496. try {
  1497. text = UTF8_FATAL.decode(line)
  1498. } catch {
  1499. continue
  1500. }
  1501. // JSON.parse would silently ROUND an integer token outside the
  1502. // safe range before validation could see it, so a forged frame
  1503. // could smuggle a corrupted value into a dispatch or completion.
  1504. // An honest child never emits one (its validator rejects unsafe
  1505. // ints), so such a frame is hostile traffic: drop it like any
  1506. // other junk frame.
  1507. if (hasUnsafeIntegerToken(text)) continue
  1508. let parsed: unknown
  1509. try {
  1510. parsed = JSON.parse(text) as unknown
  1511. } catch {
  1512. continue // Junk frames drop silently (hostile-peer stance).
  1513. }
  1514. const message = validateChildFrame(parsed)
  1515. if (message) handleFrame(message)
  1516. }
  1517. // Carry the residual forward as a fresh, right-sized copy, NOT the
  1518. // `subarray` view: a view keeps the whole joined-frame allocation from
  1519. // the `Buffer.concat` above alive, so a large frame followed by a tiny
  1520. // trailing fragment would pin megabytes while `pendingBytes` reported
  1521. // only the fragment's length. See {@link detachResidual}.
  1522. pendingChunks = detachResidual(buffered)
  1523. pendingBytes = buffered.length
  1524. } else if (pendingChunks.length >= MAX_PENDING_CHUNKS) {
  1525. // A newline-free run past the fragment-count bound: seal the held
  1526. // chunks into one finished block (amortized O(1) per byte, see the
  1527. // comment above the count bound) and keep accumulating. The gate on
  1528. // `chunk.includes(0x0a)` is the ELSE half of the newline branch, so a
  1529. // newline-bearing chunk never lands in a sealed block.
  1530. sealedBlocks.push(Buffer.concat(pendingChunks))
  1531. pendingChunks = []
  1532. }
  1533. })
  1534. // Duplicate-call suppression against the honest child's id SEQUENCE, not
  1535. // a set of every id seen. `dispatch` sends consecutive ids from 0 with no
  1536. // gaps — it advances its counter only after the write succeeds, so a call
  1537. // rejected before reaching the wire consumes nothing — which makes the
  1538. // next legitimate id exactly `nextCallId`.
  1539. //
  1540. // Retaining a set instead let a program write an unbounded run of unique
  1541. // forged ids, each below the 64 MiB per-frame parse cap so nothing
  1542. // rejected them, and grow host memory for the whole run. Accepting any
  1543. // id above a high-water mark would have been just as wrong in the other
  1544. // direction: one forged `{"id": 9999}` would starve every honest call
  1545. // after it. The exact successor is the only test that both bounds the
  1546. // retained state to one number and cannot be poisoned by a forgery.
  1547. let nextCallId = 0
  1548. // Set by run() when the boot frame is written; the fd-3 handler calls it
  1549. // on boot-ack to send the run frame (see the seam's boot->boot-ack->run
  1550. // order). scoped per run. An object holder so the cross-closure
  1551. // assignment is a property write (eslint's prefer-const cannot see the
  1552. // reassignment through the closure).
  1553. const bootAckGate: { run?: () => void } = {}
  1554. const handleFrame = (message: ChildToHost): void => {
  1555. /* v8 ignore next -- late frame after settlement; defensive against forged post-settlement traffic. */
  1556. if (settled) return
  1557. switch (message.type) {
  1558. case 'boot-ack':
  1559. // The child accepted the boot frame (namespaces built); the run
  1560. // frame goes out now, not with the boot frame.
  1561. bootAckGate.run?.()
  1562. return
  1563. case 'log':
  1564. if (message.truncated === true) {
  1565. // The CHILD ledger hit its cap. Its marker is the last log text
  1566. // there will be, so record it and stop host capture at the same
  1567. // point: admitting it as ordinary text left the host budget open,
  1568. // so later direct `os.write(1, ...)` bytes were retained AFTER the
  1569. // marker and a host-side exhaustion could append a second one.
  1570. // Both ledgers are keyed to the same `maxLogBytes`, so one marker
  1571. // describes the run.
  1572. if (!logsTruncated) {
  1573. // The host's OWN marker, never the frame's text. `truncated` is
  1574. // attacker-reachable, so trusting the text let a program write
  1575. // `{"type":"log","truncated":true,"text":<1 MiB>}` and land all
  1576. // of it in `logs` under a 64-byte `maxLogBytes` — measured, the
  1577. // whole megabyte was retained, bypassing `admit` and its
  1578. // ceiling. Both ledgers key off the same `maxLogBytes`, so the
  1579. // marker the host generates says the same thing the child's
  1580. // would have.
  1581. truncateLogs()
  1582. }
  1583. return
  1584. }
  1585. if (message.open === true) {
  1586. // An explicit flush of an unterminated line: hold it so the next
  1587. // frame appends to the SAME entry (print('a', end='', flush=True)
  1588. // followed by print('b') reads back as one 'ab' entry, not a fake
  1589. // newline). Billed INCREMENTALLY so k tiny frames cost O(k), not
  1590. // O(k * budget) (re-walking the whole held text per frame): the
  1591. // first fragment is charged the full JSON-string cost plus the
  1592. // separator (quotes + content + newline), each continuation only
  1593. // its content (jsonStringCostUpTo includes the two quotes), and
  1594. // the closing frame only its own content — the merged entry's
  1595. // wire cost is billed exactly once, split across the fragments.
  1596. // Caps: the first fragment's exact-cost walk uses logBudget - 1
  1597. // (the ledger's reserved byte, matching admit), a continuation's
  1598. // logBudget + 2 (a continuation is billed WITHOUT quotes, so its
  1599. // billed cost cost - 2 fits exactly when the walk's cost is at
  1600. // most logBudget + 2).
  1601. if (!logsTruncated) {
  1602. // An EMPTY first open frame (openParts empty AND text '') bills
  1603. // cost + 1 = 3 but establishes no hold (the push is skipped),
  1604. // so the next frame is billed as a new first fragment. Not
  1605. // reachable from an honest child (_LogStream.write('') returns
  1606. // early; flush_line pushes only non-empty pending); for a
  1607. // forged frame it is a bounded over-charge in the safe
  1608. // direction (a flood exhausts the ledger into truncation).
  1609. const cap = openParts.length === 0 ? logBudget - 1 : logBudget + 2
  1610. const cost = jsonStringCostUpTo(message.text, cap)
  1611. if (cost === undefined) {
  1612. truncateLogs()
  1613. } else {
  1614. const bill = openParts.length === 0 ? cost + 1 : Math.max(cost - 2, 0)
  1615. logBudget -= bill
  1616. // A zero-content continuation (text '') bills 0; holding it
  1617. // would grow the fragment array without touching the ledger,
  1618. // so a forged empty-open flood could grow host memory — skip
  1619. // the push, the merge result is unchanged.
  1620. if (message.text !== '') {
  1621. if (openParts.length >= MAX_PENDING_CHUNKS) {
  1622. openSealed.push(openParts.join(''))
  1623. openParts = []
  1624. }
  1625. openParts.push(message.text)
  1626. }
  1627. }
  1628. }
  1629. return
  1630. }
  1631. if (openParts.length > 0) {
  1632. // Closing frame: the held fragments are already billed; bill only
  1633. // this frame's own content (the quotes and separator ride on the
  1634. // first fragment) and push the merged entry once. Cap is
  1635. // logBudget + 2 for the same reason as a continuation.
  1636. /* v8 ignore next -- logsTruncated is an invariant false here: an open
  1637. * frame that would trip the ledger resets openParts, so a non-empty
  1638. * hold implies the ledger never truncated. The guard is defensive. */
  1639. if (!logsTruncated) {
  1640. const cost = jsonStringCostUpTo(message.text, logBudget + 2)
  1641. if (cost === undefined) {
  1642. truncateLogs()
  1643. } else {
  1644. logBudget -= Math.max(cost - 2, 0)
  1645. logs.push(openSealed.join('') + openParts.join('') + message.text)
  1646. }
  1647. }
  1648. openSealed = []
  1649. openParts = []
  1650. return
  1651. }
  1652. admit(message.text)
  1653. return
  1654. case 'done': {
  1655. if (message.error) {
  1656. finish({ error: { kind: message.error.kind, message: capMessage(message.error.message, this.config.maxValueBytes) } })
  1657. return
  1658. }
  1659. if (message.value === undefined) {
  1660. finish({})
  1661. return
  1662. }
  1663. // Re-enforce the completion budget and number losslessness
  1664. // host-side: a forged done frame bypasses the Python-side
  1665. // _done_with_value check, and validateChildFrame no longer scans
  1666. // the value (an unbounded scan would push every member of a wide
  1667. // forgery before any cap ran). checkDoneValue folds both jobs into
  1668. // one bounded, iterative traversal — iterative because the seam's
  1669. // CodeJsonValue has no depth limit and an honest deep-but-small
  1670. // completion must cross intact rather than dying on stringify
  1671. // recursion; bounded because it stops at the cap without
  1672. // materializing the encoding, rejecting a forged value anywhere
  1673. // below the 64 MiB frame parse cap before it forces host-side copies.
  1674. // The seam forbids substituting a rendered/truncated value, so an
  1675. // oversized value fails the run as output-limit and a non-lossless
  1676. // number as invalid-output. The value is JSON-plain by construction
  1677. // (it came from JSON.parse of the frame), the traversal's precondition.
  1678. const check = checkDoneValue(message.value, this.config.maxValueBytes)
  1679. if (!check.ok) {
  1680. finish(check.reason === 'over-budget'
  1681. ? { error: { kind: 'output-limit', message: `completion value exceeded ${this.config.maxValueBytes} bytes` } }
  1682. : { error: { kind: 'invalid-output', message: 'completion value contained a non-lossless number' } })
  1683. return
  1684. }
  1685. finish({ value: message.value as CodeJsonValue })
  1686. return
  1687. }
  1688. case 'call': {
  1689. if (message.id !== nextCallId) return
  1690. nextCallId += 1
  1691. const record = bindings.get(message.global)?.functions
  1692. const fn = record && Object.hasOwn(record, message.name) ? record[message.name] : undefined
  1693. if (typeof fn !== 'function') {
  1694. // `call.global` and `call.name` are attacker-controlled strings
  1695. // with no byte cap of their own — only the 64 MiB fd-3 frame
  1696. // parse cap — so each is sliced to `maxValueBytes` CODE UNITS
  1697. // BEFORE it reaches the template. Interpolating them whole would
  1698. // copy them into the message, `JSON.stringify` would copy the
  1699. // escaped form, `encodeJsonPlain` the frame, and the pipe write
  1700. // again: four full-size host allocations off one below-ceiling
  1701. // forgery, past every hostile-peer bound the log and done-error
  1702. // paths apply. Nothing past the first `maxValueBytes` code units
  1703. // of either field can survive the byte cap anyway, so the slices
  1704. // lose only text `capMessage` would drop, and that final cap
  1705. // gives this reply the same budget and marker as a forged done
  1706. // error.
  1707. const cap = this.config.maxValueBytes
  1708. const target = `${message.global.slice(0, cap)}.${message.name.slice(0, cap)}`
  1709. // JSON.stringify on the WHOLE capped target would still allocate
  1710. // the escaped form — up to ~6x under control-heavy input, a
  1711. // multi-hundred-MB spike near the maxValueBytes ceiling that no
  1712. // hostile-peer bound would have admitted. The message only needs
  1713. // to identify the binding, so the escaped form is built from a
  1714. // 1 KiB prefix; capMessage then enforces the reply budget.
  1715. const preview = JSON.stringify(target.slice(0, 1024))
  1716. sendReply({ type: 'reply', id: message.id, ok: false, message: capMessage(`unknown binding ${preview}`, cap) })
  1717. return
  1718. }
  1719. // A binding that never settles (or resolves too slowly to keep up
  1720. // with the child's call rate) must not let the flood accumulate one
  1721. // async closure per frame until the wall clock: the reply cap only
  1722. // counts resolved calls, so it never trips for in-flight ones.
  1723. // Count the outstanding binding calls here, before dispatch, and
  1724. // release the slot in the body's finally — bounding in-flight
  1725. // closures to MAX_PENDING_REPLIES exactly like the reply backlog.
  1726. if (pendingCalls >= MAX_PENDING_REPLIES) {
  1727. finish({ error: { kind: 'worker-exit', message: `call backlog exceeded ${MAX_PENDING_REPLIES} in-flight binding calls (a binding never settled)` } })
  1728. return
  1729. }
  1730. pendingCalls += 1
  1731. void (async () => {
  1732. try {
  1733. const resolved = await fn(message.args)
  1734. // Drop a reply the run no longer needs BEFORE snapshotting it.
  1735. // `sendReply` also checks `settled`, but only after this value has
  1736. // been walked and copied: a binding that resolves a wide value
  1737. // after `maxWallMs`, an abort, or dispose already settled the run
  1738. // would spend host heap on a frame that is then discarded, and
  1739. // binding resolution carries no seam-level byte cap to bound it.
  1740. // oxlint-disable-next-line typescript/no-unnecessary-condition -- the run can settle while this binding is awaited.
  1741. if (settled) return
  1742. // The seam requires a lossy resolution to REJECT descriptively,
  1743. // not silently coerce: a raw JSON.stringify would turn NaN/
  1744. // Infinity into null and drop undefined fields. Snapshot through
  1745. // the same lossless-JSON boundary the worker backend uses (also
  1746. // iterative, so a deeply nested value cannot overflow the stack).
  1747. const value = snapshotJsonValue(resolved)
  1748. if (value === undefined) {
  1749. sendReply({ type: 'reply', id: message.id, ok: false, message: 'binding resolution must be lossless JSON' })
  1750. return
  1751. }
  1752. sendReply({ type: 'reply', id: message.id, ok: true, value })
  1753. } catch (error: unknown) {
  1754. // Check `settled` before formatting the error: a rejection that
  1755. // arrives after `maxWallMs`, an abort, or dispose has already
  1756. // settled the run, and `messageOf(error)` runs hostile getters
  1757. // before `sendReply` peeks at `settled`. Dropping the framed
  1758. // reply early spares the host heap and time for a run whose
  1759. // outcome is already fixed.
  1760. // (oxlint block-disable so both `v8 ignore next` and the rule
  1761. // suppression land on the `if`: `settled` flips true mid-wait,
  1762. // invisible to the type-aware lint, which narrows it to false.)
  1763. /* oxlint-disable typescript/no-unnecessary-condition */
  1764. /* v8 ignore next -- a rejection arriving after settlement is not schedulable from a test. */
  1765. if (settled) return
  1766. /* oxlint-enable typescript/no-unnecessary-condition */
  1767. sendReply({ type: 'reply', id: message.id, ok: false, message: messageOf(error) })
  1768. } finally {
  1769. // Release the in-flight slot on every exit — reply written,
  1770. // resolution rejected, or the run settling mid-wait (the
  1771. // `settled` early returns above). Without this, a binding that
  1772. // never resolves would leak its slot past the cap check and the
  1773. // flood bound would erode.
  1774. pendingCalls -= 1
  1775. }
  1776. })()
  1777. return
  1778. }
  1779. }
  1780. }
  1781. // Write one reply frame with the iterative encoder: a binding
  1782. // resolution has no seam-level depth or byte cap, so a deeply nested
  1783. // value must not die on JSON.stringify's recursion. The payload is
  1784. // JSON-plain by construction (snapshotJsonValue output, or literal
  1785. // strings/numbers), which is encodeJsonPlain's precondition. A closed
  1786. // pipe (child already gone) is swallowed since the close path settles
  1787. // the run.
  1788. //
  1789. // Replies are encoded and written ONE AT A TIME, waiting for `drain`
  1790. // whenever fd 3's buffer is full. Binding resolution carries no
  1791. // seam-level byte cap, so a program that resolves several large values in
  1792. // one `asyncio.gather` round would otherwise encode them all in the same
  1793. // turn and queue every frame in the writable stream's buffer -- measured
  1794. // to exhaust a 256 MiB Node heap, which kills the whole host process
  1795. // rather than failing this one run. Pacing changes no model-visible
  1796. // behavior: the child matches each reply to its `call` by id from a pump
  1797. // that reads fd 3 continuously, so arrival order was never observable,
  1798. // and the bindings themselves still run concurrently. Only the host's peak
  1799. // memory and the flush timing change.
  1800. const replyQueue: ReplyMessage[] = []
  1801. // Replies queued but not yet written, tracked separately from
  1802. // `replyQueue.length`: the drain loop clears consumed slots to `undefined`
  1803. // but does not shrink the array until it finishes, so `length` counts
  1804. // consumed frames too. The counter is what the cap in `sendReply` reads.
  1805. let pendingReplies = 0
  1806. // Binding calls dispatched but not yet settled (the async body below
  1807. // still awaits the binding's promise). The reply backlog cap only counts
  1808. // RESOLVED calls — `pendingReplies` grows after the await — so a child
  1809. // flooding calls against a binding that never settles would accumulate
  1810. // one async closure per frame until the wall clock without tripping it.
  1811. // Counted here before dispatch and released in the body's finally, the
  1812. // in-flight closures are bounded to the same MAX_PENDING_REPLIES.
  1813. let pendingCalls = 0
  1814. let draining = false
  1815. // Resolve when fd 3 can take another frame, OR when it is gone: a pipe
  1816. // destroyed under the drain (child exited, close-deadline teardown) never
  1817. // emits 'drain' again, so waiting on that event alone would hang the
  1818. // drain forever — `draining` stays true and the unconsumed queue is
  1819. // pinned with the closure. `once` plus the manual detach removes every
  1820. // listener whichever event wins, so a long backpressure wait leaves none
  1821. // behind.
  1822. const waitForDrain = (): Promise<void> => new Promise<void>((resolvePromise) => {
  1823. const finish = (): void => {
  1824. proto.off('drain', finish)
  1825. proto.off('close', finish)
  1826. proto.off('error', finish)
  1827. resolvePromise()
  1828. }
  1829. proto.once('drain', finish)
  1830. proto.once('close', finish)
  1831. proto.once('error', finish)
  1832. })
  1833. const drainReplies = async (): Promise<void> => {
  1834. if (draining) return
  1835. draining = true
  1836. let head = 0
  1837. try {
  1838. while (head < replyQueue.length) {
  1839. // Needs the run to settle between two queued frames. Measured queue
  1840. // depths reach 11 without the wall clock landing inside that window.
  1841. /* v8 ignore next -- see above; not schedulable from a test. */
  1842. if (settled) break
  1843. // A pipe destroyed under us (child exited, close deadline) will
  1844. // never emit 'drain' again; short-circuit before the write so the
  1845. // remaining frames are dropped by the `finally` below.
  1846. if (proto.destroyed) break
  1847. // Read by index, not `shift()`: a large `asyncio.gather` of wide
  1848. // bindings awaiting fd 3's `drain` can queue many frames, and each
  1849. // `shift()` re-slices the remaining array (O(n) per pop, O(n²) over
  1850. // the whole drain). A head cursor keeps the cost linear; the `finally`
  1851. // below discards everything consumed once the drain ends. The consumed
  1852. // slot is CLEARED here (not just advanced past) so a wide payload the
  1853. // pipe has already taken is released immediately: under sustained
  1854. // backpressure the drain loop can live across many `await drain`
  1855. // ticks, and leaving the slot set would pin the written value's bytes
  1856. // in `replyQueue` for the whole busy period, making host memory grow
  1857. // with cumulative processing rather than the current backlog.
  1858. const payload = replyQueue[head] as ReplyMessage
  1859. replyQueue[head] = undefined as unknown as ReplyMessage
  1860. head += 1
  1861. pendingReplies -= 1
  1862. // Compact the consumed prefix once it reaches the backlog bound:
  1863. // the array never shrinks until the drain finishes, and a child
  1864. // that reads replies just fast enough to keep the drain alive but
  1865. // never empty would otherwise grow the backing store linearly with
  1866. // cumulative throughput (consumed slots are undefined, but `length`
  1867. // keeps counting them). The splice is O(head) once per
  1868. // MAX_PENDING_REPLIES consumed frames — amortized O(1) per reply.
  1869. if (head >= MAX_PENDING_REPLIES) {
  1870. replyQueue.splice(0, head)
  1871. head = 0
  1872. }
  1873. // Encode inside the loop, not up front: a queued reply the run no
  1874. // longer needs is dropped by the `settled` check above without ever
  1875. // being serialized.
  1876. if (!proto.write(`${encodeJsonPlain(payload)}\n`)) {
  1877. await waitForDrain()
  1878. }
  1879. }
  1880. } catch {
  1881. // Pipe closed under us (child exited), or `drain` never arrives because
  1882. // the child died. The close path settles the run either way.
  1883. } finally {
  1884. draining = false
  1885. pendingReplies = 0
  1886. replyQueue.length = 0
  1887. }
  1888. }
  1889. const sendReply = (payload: ReplyMessage): void => {
  1890. /* v8 ignore next -- `settled` covers a race where the child exits between decision and write. */
  1891. if (settled) return
  1892. // A child that stops reading fd 3 leaves the drain loop blocked on
  1893. // `drain` forever while its call frames keep resolving into replies:
  1894. // the backlog would grow without bound until the wall clock, pinning
  1895. // every binding result the child provokes. Cap the retained backlog and
  1896. // settle the run as a worker-exit, the same hostile-peer bound the
  1897. // frame cap applies to inbound bytes.
  1898. if (pendingReplies >= MAX_PENDING_REPLIES) {
  1899. finish({ error: { kind: 'worker-exit', message: `reply queue exceeded ${MAX_PENDING_REPLIES} pending frames on fd 3 (the child stopped consuming its replies)` } })
  1900. return
  1901. }
  1902. pendingReplies += 1
  1903. replyQueue.push(payload)
  1904. void drainReplies()
  1905. }
  1906. // Escalate SIGTERM → grace → SIGKILL on the entire process group. Idempotent
  1907. // via `killing`.
  1908. let killing = false
  1909. let graceTimer: NodeJS.Timeout | undefined
  1910. // A backstop for the one case `close` cannot cover: model code that starts
  1911. // a descendant with `os.setsid()`/`start_new_session=True` moves it into a
  1912. // fresh process group, so the SIGTERM/SIGKILL aimed at the child's group
  1913. // (`kill(-pid)`) never reaches it. If that orphan inherited stdout/stderr/
  1914. // fd 3 and outlives the run, those pipes stay open and `close` never fires
  1915. // — leaving run() (and a teardown awaiting `finished`) hung indefinitely.
  1916. // finish() arms this deadline; when it fires we detach our stream handles
  1917. // and settle on the already-decided result regardless of the orphan.
  1918. let closeDeadline: NodeJS.Timeout | undefined
  1919. // The leader's start time, read once while it is certainly alive. `child.pid`
  1920. // keeps its numeric value after the leader is reaped (Node clears the
  1921. // internal handle, not the field), and `close` can trail `exit` by seconds
  1922. // while a pipe-holding descendant keeps the streams open. Signalling
  1923. // `-child.pid` in that window is a RAW syscall -- `child.kill()` would
  1924. // refuse, having dropped its handle, but `process.kill` has no such guard --
  1925. // so a recycled pgid would receive this run's SIGTERM and armed SIGKILL.
  1926. // `groupEmpty()` cannot cover it: it reports whether the group has members,
  1927. // not whether they are OURS, and it runs only after the first signal.
  1928. // The repository already takes this position in
  1929. // packages/subprocess/subprocess-local (`ProcessIdentity`, "preventing
  1930. // teardown escalation after PID reuse"); this is the same guard, kept local
  1931. // because a dependency on that package would be a new architectural edge.
  1932. const leaderStarted = child.pid === undefined ? undefined : readProcessStart(child.pid)
  1933. const killGroup = (sig: NodeJS.Signals): void => {
  1934. try {
  1935. /* v8 ignore next -- undefined pid means spawn never produced a process; finish() short-circuits before reaching kill(). */
  1936. if (child.pid === undefined) return
  1937. // A pid alone cannot answer this: `process.kill(pid, 0)` succeeds just
  1938. // as well for a REPLACEMENT process holding the recycled number. Only
  1939. // the start time distinguishes the two, so a reading that DISAGREES
  1940. // means the number now belongs to another process and must not be
  1941. // signalled.
  1942. //
  1943. // An ABSENT reading is the ordinary case, not a mismatch: once the
  1944. // leader is reaped its `/proc/<pid>/stat` is gone, while the group it
  1945. // led can still hold survivors that this teardown exists to reap. So
  1946. // only a present-and-different reading blocks the signal; undefined
  1947. // falls through, which is also the behavior on platforms with no
  1948. // `/proc` to read.
  1949. const nowStarted = readProcessStart(child.pid)
  1950. // The refusal arm needs a real pid recycled into a new group leader
  1951. // between spawn and teardown, which no test can schedule; the reader
  1952. // itself is covered directly by the process-identity test.
  1953. /* v8 ignore next -- unreachable without real pid reuse; see above. */
  1954. if (leaderStarted !== undefined && nowStarted !== undefined && nowStarted !== leaderStarted) return
  1955. process.kill(-child.pid, sig)
  1956. } catch {
  1957. // ESRCH — the process already died. Nothing to do.
  1958. }
  1959. }
  1960. const kill = (): void => {
  1961. /* v8 ignore next -- kill() is idempotent; tests do not double-invoke it. */
  1962. if (killing) return
  1963. killing = true
  1964. killGroup('SIGTERM')
  1965. // Escalate to SIGKILL after the grace window. The timer is `unref`'d so a
  1966. // pending SIGKILL never keeps the host process alive on its own; the
  1967. // guarantee that a same-group survivor is actually reaped before the fiber
  1968. // goes quiescent is enforced by settle() awaiting the group's death (see
  1969. // there), NOT by this timer firing during host lifetime. A setsid-escaped
  1970. // orphan in a FRESH group is the different case `closeDeadline` in finish()
  1971. // covers, since `close` never fires there.
  1972. graceTimer = setTimeout(() => { killGroup('SIGKILL') }, this.config.graceMs)
  1973. graceTimer.unref()
  1974. }
  1975. // True once the group has no members left: a signal-0 probe to the whole
  1976. // group (`kill(-pid, 0)`) throws ESRCH when empty (EPERM would still mean a
  1977. // member exists). Only meaningful once a spawn produced a pid.
  1978. const groupEmpty = (): boolean => {
  1979. /* v8 ignore next -- pid is always defined once escalation runs; the guard narrows the type. */
  1980. if (child.pid === undefined) return true
  1981. try {
  1982. process.kill(-child.pid, 0)
  1983. return false
  1984. } catch (error: unknown) {
  1985. return (error as NodeJS.ErrnoException).code === 'ESRCH'
  1986. }
  1987. }
  1988. let finishResolve!: () => void
  1989. const finished = new Promise<void>((done) => { finishResolve = done })
  1990. let resolved = false
  1991. // The decided terminal result for a live child, recorded by finish() and
  1992. // read by the `close` handler that settles it once the pipes have drained.
  1993. let decided: Omit<CodeRunResult, 'logs'>
  1994. // The single settlement point: resolve run() with the decided result and
  1995. // mark the fiber quiescent. Idempotent — the first call wins, so a later
  1996. // `close` after done/timeout/abort is absorbed as a no-op.
  1997. const settle = (result: Omit<CodeRunResult, 'logs'>): void => {
  1998. if (resolved) return
  1999. resolved = true
  2000. if (closeDeadline !== undefined) clearTimeout(closeDeadline)
  2001. // The child has exited by now (settle runs on `close`, or on a spawn
  2002. // that produced no pid), so its staging directory is no longer read and
  2003. // this run's copy goes away with it. Removed SYNCHRONOUSLY, before
  2004. // `resolve` below: a fire-and-forget removal left the directory on disk
  2005. // when `run()` resolved, so a caller could not observe the "gone by
  2006. // settlement" contract at all. Two files cost nothing to unlink here.
  2007. try {
  2008. rmSync(bootstrapDir, { recursive: true, force: true })
  2009. } catch {
  2010. // Swallows only a failure to remove this run's staging directory —
  2011. // `force` already absorbs a missing one, so what remains is a
  2012. // filesystem-level refusal. The run's own outcome is already decided
  2013. // and must still be delivered; the directory holds no secret, only a
  2014. // copy of two checked-in scripts. teardown deliberately does not
  2015. // sweep staging (its staging is cleared inside each run's settle), so
  2016. // a removal failure here is the one case the "gone by settlement"
  2017. // contract degrades on.
  2018. }
  2019. resolve({ ...result, logs })
  2020. // Mark the fiber quiescent for THIS run: drop it from `live` and resolve
  2021. // `finished` (what teardown awaits). Deferred until the process group is
  2022. // actually empty — dropping from `live` before then would let a
  2023. // `dispose()` that races a just-resolved run() snapshot an empty `live`
  2024. // and return while a same-group survivor is still alive, making teardown's
  2025. // "no SAME-GROUP subprocess outlives the fiber" guarantee false for that
  2026. // window (a setsid escapee is the documented exception — see teardown's
  2027. // JSDoc). Keeping the run in `live` until the group is reaped is exactly
  2028. // what makes a concurrent teardown await it.
  2029. const finalize = (): void => {
  2030. this.live.delete(live)
  2031. finishResolve()
  2032. }
  2033. // `finished` is what teardown awaits to honor "no same-group subprocess
  2034. // outlives the fiber". When no escalation ran (normal completion, no
  2035. // kill) or the group is already empty, cancel the pending SIGKILL and
  2036. // finalize now. Clearing it is what bounds the PID-reuse hazard: an armed
  2037. // `kill(-pid)` left to fire up to graceMs later could hit a RECYCLED pgid
  2038. // once the kernel reused the leader's pid, SIGKILLing an unrelated group.
  2039. // So the timer stays armed only while a real survivor exists — a
  2040. // same-group descendant that ignored SIGTERM but released the pipes,
  2041. // still alive here because its `close` is what got us to settle. In that
  2042. // case withhold finalize and poll the group on REF'd timers (a
  2043. // short-lived host would otherwise exit before the unref'd SIGKILL fired,
  2044. // reparenting the survivor to init), clearing the timer the moment the
  2045. // group empties. The wait is bounded by `graceMs + CLOSE_REAP_MARGIN_MS`
  2046. // in the normal case; if the host event loop was blocked past both timers
  2047. // the deadline branch below sends SIGKILL itself and grants ONE more reap
  2048. // margin, so the outer bound is `graceMs + 2 * CLOSE_REAP_MARGIN_MS`.
  2049. if (!killing || groupEmpty()) {
  2050. if (graceTimer !== undefined) clearTimeout(graceTimer)
  2051. finalize()
  2052. return
  2053. }
  2054. const deadline = Date.now() + this.config.graceMs + CLOSE_REAP_MARGIN_MS
  2055. // Once the deadline forces us to send SIGKILL ourselves, allow one more
  2056. // reap window for the kernel to tear the group down before giving up:
  2057. // SIGKILL is asynchronous, so the group is not gone the instant it is
  2058. // sent. `finalize` only runs on a confirmed-empty group, except at this
  2059. // final hard bound where nothing more can be done.
  2060. let hardDeadline = 0
  2061. const pollGroup = (): void => {
  2062. if (groupEmpty()) {
  2063. // The group is gone; the grace SIGKILL is moot. Cancel it (it may not
  2064. // have fired yet) and finalize. graceTimer is always defined here:
  2065. // pollGroup runs only when `killing` is set, and kill() armed it.
  2066. clearTimeout(graceTimer)
  2067. finalize()
  2068. return
  2069. }
  2070. if (hardDeadline === 0 && Date.now() >= deadline) {
  2071. // Deadline reached with the group still non-empty. This is reachable
  2072. // when the host event loop was blocked past both timers: Node runs
  2073. // this poll before the grace SIGKILL timer, so that SIGKILL may never
  2074. // have fired. Send it HERE (idempotent if the timer already ran) and
  2075. // keep polling for the group to actually empty — finalizing on mere
  2076. // signal delivery would declare quiescence while the group is still
  2077. // dying. Bound the extra wait by one more reap margin.
  2078. killGroup('SIGKILL')
  2079. clearTimeout(graceTimer)
  2080. hardDeadline = Date.now() + CLOSE_REAP_MARGIN_MS
  2081. }
  2082. // Hard bound: the self-sent SIGKILL delivered but `groupEmpty()` still
  2083. // reports the group non-empty for a full extra reap margin. This is
  2084. // reachable, not a kernel quirk: a SIGKILL'd same-group survivor
  2085. // lingers as a ZOMBIE until its parent `wait()`s it, and in a
  2086. // container whose PID 1 does not reap orphans the survivor is
  2087. // reparented to init and never waited, so the signal-0 probe keeps
  2088. // succeeding — the same environment dependence the Agent Note's
  2089. // rejected "assert the reap with process.kill(pid, 0)" alternative
  2090. // documents. The ignore stays because that container cannot be built
  2091. // deterministically across CI platforms, not because the branch is
  2092. // unreachable; finalizing here bounds the wait so such a deployment
  2093. // still goes quiescent within `graceMs + 2 * CLOSE_REAP_MARGIN_MS`.
  2094. /* v8 ignore next 4 -- reachable only in a PID-1-doesn't-reap container (zombie survivor); not deterministically buildable. */
  2095. if (hardDeadline !== 0 && Date.now() >= hardDeadline) {
  2096. finalize()
  2097. return
  2098. }
  2099. setTimeout(pollGroup, GROUP_REAP_POLL_MS)
  2100. }
  2101. pollGroup()
  2102. }
  2103. const finish = (result: Omit<CodeRunResult, 'logs'>): void => {
  2104. if (settled) return
  2105. settled = true
  2106. decided = result
  2107. clearTimeout(wallTimer)
  2108. request.signal?.removeEventListener('abort', onAbort)
  2109. // A spawn failure (ENOENT, EACCES) never produced a pid, so there is no
  2110. // process to kill: settle now. Its `close` still fires later and reaches
  2111. // the idempotent settle() again as a no-op.
  2112. // An unterminated flushed line never got a closing frame; it was
  2113. // billed incrementally, so push it directly (admit would re-bill).
  2114. // logsTruncated implies the hold is already empty (truncateLogs
  2115. // committed and cleared it), so this is reachable only when the run
  2116. // ends with the hold still open and untruncated.
  2117. if (openSealed.length > 0 || openParts.length > 0) {
  2118. logs.push(openSealed.join('') + openParts.join(''))
  2119. }
  2120. openSealed = []
  2121. openParts = []
  2122. if (child.pid === undefined) {
  2123. settle(result)
  2124. return
  2125. }
  2126. // Live child: SIGTERM→grace→SIGKILL, then let `close` (below) settle the
  2127. // run so any `done` frame buffered on fd 3 is handled first and the
  2128. // process is fully reaped before the fiber goes quiescent.
  2129. kill()
  2130. // `close` awaits every stdio stream draining, which a setsid-escaped
  2131. // orphan holding our inherited pipes can prevent forever. Bound that
  2132. // wait: after SIGKILL has had the grace window plus a margin to reap the
  2133. // child itself, force settlement on the decided result. Flush any
  2134. // newline-free stray residual FIRST — a leader that wrote a diagnostic
  2135. // with `os.write(1, ...)` and exited leaves it buffered, and destroying
  2136. // the stream below drops it before an `end`/`close` flush could run, so
  2137. // the diagnostic would be lost from `logs`. Detaching the stream handles
  2138. // then lets `close` land as a no-op if it ever arrives, and stops the
  2139. // orphan's stray output from being accounted against a run that already
  2140. // finished. `unref` so the deadline never keeps the host process alive.
  2141. closeDeadline = setTimeout(() => {
  2142. flushStray(strayOut)
  2143. flushStray(strayErr)
  2144. proto.destroy()
  2145. child.stdout.destroy()
  2146. child.stderr.destroy()
  2147. settle(result)
  2148. }, this.config.graceMs + CLOSE_REAP_MARGIN_MS)
  2149. closeDeadline.unref()
  2150. }
  2151. child.on('error', (error: Error) => {
  2152. finish({ error: { kind: 'worker-exit', message: `python spawn error: ${error.message}` } })
  2153. })
  2154. // `close` (not `exit`) is the settlement trigger: it fires only after the
  2155. // process exits AND every stdio stream — including the fd-3 protocol pipe —
  2156. // has drained, so a `done` frame the child wrote just before exiting is
  2157. // always handled before we settle. macOS can deliver `exit` before that
  2158. // final fd-3 data; keying off `close` makes the ordering irrelevant.
  2159. child.on('close', (code: number | null, signal: NodeJS.Signals | null) => {
  2160. // If done/timeout/abort already decided the result, finish() is a no-op
  2161. // and `decided` holds it — a SIGXCPU that arrives after a decision does
  2162. // not override it. Otherwise the child closed before completing: a
  2163. // SIGXCPU close is the kernel's own CPU meter firing — the RLIMIT_CPU
  2164. // soft limit, or the bootstrap's post-settlement getrusage check
  2165. // re-delivering SIGXCPU when a program trapped the soft limit and
  2166. // returned inside the soft-to-hard gap. That kernel-authoritative
  2167. // signal is the ONLY basis for the timeout classification: wall time
  2168. // is not evidence of CPU burn (a sleeping child SIGKILLed by a cgroup
  2169. // OOM killer, an operator, or itself consumed none), so every other
  2170. // signal or code — including an unsolicited SIGKILL, even the
  2171. // hard-limit one — reports as an opaque worker exit.
  2172. //
  2173. // The message names `cpuSeconds` as the CONFIGURED ceiling, not "the
  2174. // budget that fired": the child clamps RLIMIT_CPU to the stricter of
  2175. // `cpuSeconds` and any inherited soft limit, so under a tighter inherited
  2176. // cap SIGXCPU arrives before `cpuSeconds` — the host cannot see the
  2177. // effective value, so it states the ceiling it set rather than a second
  2178. // count it cannot guarantee.
  2179. finish(signal === 'SIGXCPU'
  2180. ? { error: { kind: 'timeout', message: `CPU time exhausted (limit at most the configured ${this.config.cpuSeconds}s; a stricter inherited RLIMIT_CPU can fire sooner)` } }
  2181. : { error: { kind: 'worker-exit', message: `python exited (code=${String(code)}, signal=${String(signal)}) before completing` } })
  2182. settle(decided)
  2183. })
  2184. // Fd-3 and the stdout/stderr pipes emit `error` on early child death
  2185. // (ECONNRESET/EPIPE); swallow them so they do not become uncaught. The
  2186. // authoritative failure signal is `child.on('close')` above.
  2187. const silenceStreamError = (): void => {}
  2188. proto.on('error', silenceStreamError)
  2189. child.stdout.on('error', silenceStreamError)
  2190. child.stderr.on('error', silenceStreamError)
  2191. /* jscpd:ignore-start -- wall-timer/abort/live-run wiring deliberately parallels code-runtime-worker; see the constructor note. */
  2192. const wallTimer = setTimeout(() => {
  2193. finish({ error: { kind: 'timeout', message: `wall-clock ceiling reached (${this.config.maxWallMs}ms)` } })
  2194. }, this.config.maxWallMs)
  2195. const onAbort = (): void => {
  2196. finish({ error: { kind: 'abort', message: messageOf(request.signal?.reason) } })
  2197. }
  2198. request.signal?.addEventListener('abort', onAbort, { once: true })
  2199. const live: LiveRun = {
  2200. kill,
  2201. finished,
  2202. settle: (failure: CodeRunFailure) => { finish({ error: failure }) },
  2203. }
  2204. this.live.add(live)
  2205. /* jscpd:ignore-end */
  2206. // Send the boot frame once fd 3 is writable. This runs LAST in run()'s
  2207. // synchronous setup: its failure path calls finish(), which reads
  2208. // wallTimer/onAbort and (through settle) live, so those bindings must
  2209. // already be initialized — issuing the write earlier hit their
  2210. // temporal dead zone and threw a ReferenceError that rejected run()
  2211. // instead of resolving the worker-exit it constructs here.
  2212. const boot: BootMessage = {
  2213. type: 'boot',
  2214. cpuSeconds: this.config.cpuSeconds,
  2215. addressSpaceBytes: this.config.addressSpaceMb * 1024 * 1024,
  2216. maxLogBytes: this.config.maxLogBytes,
  2217. maxValueBytes: this.config.maxValueBytes,
  2218. namespaces: [...bindings].map(([global, namespace]) => ({
  2219. global,
  2220. names: Object.keys(namespace.functions),
  2221. ...namespace.errorClass ? { errorClass: namespace.errorClass } : {},
  2222. })),
  2223. }
  2224. // The run frame is sent only after the child's boot-ack: the seam
  2225. // contract puts `run` after `boot-ack` (the ack confirms the namespaces
  2226. // were accepted), and sending it earlier would let a boot failure race
  2227. // the run frame. The ack handler below writes it.
  2228. let runSent = false
  2229. try {
  2230. proto.write(`${JSON.stringify(boot)}\n`)
  2231. } catch (error: unknown) {
  2232. finish({ error: { kind: 'worker-exit', message: `failed to boot python subprocess: ${messageOf(error)}` } })
  2233. return
  2234. }
  2235. // Register the ack gate with the frame handler before any data arrives.
  2236. bootAckGate.run = (): void => {
  2237. if (runSent) return
  2238. runSent = true
  2239. try {
  2240. proto.write(`${JSON.stringify({ type: 'run', program: request.program })}\n`)
  2241. } catch (error: unknown) {
  2242. /* v8 ignore next -- the child exited between its ack and this write; the run settles as worker-exit. */
  2243. finish({ error: { kind: 'worker-exit', message: `failed to boot python subprocess: ${messageOf(error)}` } })
  2244. }
  2245. }
  2246. })
  2247. }
  2248. }
  2249. export default PythonCodeRuntime