Purdex mod: context relay for this session (self relay with approval, member relay on the lead's order) and the pdx-team skill.

人與 agent 協作的工作站。把 tmux session、Claude Code 串流、對話式 agent 統一在一個介面,跨 workspace、跨 host 並行作業。
仍在 alpha (
1.0.0-alpha.233),僅支援 macOS。原名tmux-box/tmux-ai-term。
terminal / stream(wrap) / 對話詞彙與設計定錨見 PRODUCT.md。
Go daemon + React 19 SPA + Electron 41 殼。
cd spa && pnpm install
pnpm dev # SPA dev server
pnpm test # vitest
pnpm build # SPA build
pnpm electron:build # 從 root;產出 dist/mac/ + dist/mac-arm64/
環境與跨機開發流程見 CLAUDE.md。
Daemon(Go)需要存取 private module lab.protype.tw/wake/nexen。建置前, 這台機器需完成兩項一次性設定:
go env -w GOPRIVATE=lab.protype.tw
git config --global url."ssh://git@lab.protype.tw:9079/".insteadOf "https://lab.protype.tw/"
make build(或單獨 make check-goenv)會先驗證這兩項,缺一就印出對應指令 並以非零狀態結束,不會進入 go build。
冷快取驗證(確認上述設定不依賴呼叫端的 shell 環境變數也能運作;spec §6 step 0):
GOWORK=off GOMODCACHE=$(mktemp -d) go build ./cmd/pdx
PRODUCT.md — 產品定位CLAUDE.md — 開發流程CHANGELOG.md — 版本歷史hooks/register.js 1301 lines1// Purdex mod — self relay (spec §8.1–§8.3 steps 4–8, §8.7).
2//
3// idle ──(turn.complete: used ≥ threshold, growth ≥ minGrowth, +10 since last ask)──▶ beginning
4// beginning: `pdx relay begin --self` runs from a timer ──▶ awaiting, or back to idle on a refusal
5// awaiting: a timer loops `pdx relay wait` and alone asks the daemon and moves the state;
6// every prompt.submit is held until that loop's answer (≤ 11 min), on local `/bin/sleep 5`
7// calls of its own (P5b-3 critic), and one that arrives while beginning waits for begin's
8// answer first (≤ 40 s); every prompt is held, however many (U7: no new turn while a request is open)
9// awaiting ──approved──▶ approved: write prompt submitted (a fresh nonce), report writing
10// approved ──(turn.complete of the write turn, file ok)──▶ clearing: report written, timer → /clear
11// clearing ──(classic.SessionStart source=clear)──▶ seeding: report cleared --new-session, hello, seed prompt
12// seeding ──(turn.complete of the seed turn)──▶ idle: report done, floor = tokens now
13// awaiting ──denied / timeout / cancelled / unavailable──▶ idle (ask again at +10 points)
14// a deferred step that fails (prompt refused, /clear refused) ──▶ idle: write / fix / seed report
15// failed{handoff_incomplete}, /clear reports cancelled{abandoned}
16// the user's own /clear ──▶ idle: awaiting / approved report cancelled{abandoned}, seeding
17// failed{handoff_incomplete}; a begin still out is cancelled{abandoned} when it answers (s.gen)
18// a compaction ──▶ awaiting reports cancelled{compacted}, idle; approved skips an auto one
19// every return to idle lets the request go first (letGo): its held prompts go on at once,
20// unchanged, before any report; the loop's late answer for it is ignored
21//
22// Everything that starts a turn, runs a command or waits on the daemon goes
23// out from a $.clock.after timer, never inside a hook: $.command.run rejects
24// inside a hook the turn waits on (F3), and a daemon that is down answers
25// only after the client's 30 s grace, which no turn end, session start or
26// /clear may wait for (P5b-1 review). Hooks only read the engine and move
27// the state. Two hooks wait on purpose: the prompt hold, on local
28// `/bin/sleep 5` calls (a `$` call in flight stops the hook's 10 s budget,
29// HookBudget; it asks the daemon nothing, the timers ask it and move the
30// state) and failing open; and /relay, on its daemon call (the person waits
31// for its output; 8 s bound).
32//
33// Each write / fix / seed prompt is composed at use (U21, spec §8.8): the
34// mod's own fixed head and tail (prompts.js) around the body that `pdx relay
35// prompts` answers right before the prompt goes out, from the step's timer
36// (8 s bound); a call that fails, times out or answers junk gives the
37// built-in body, and the relay goes on.
38//
39// This file is the plugin's one hooks module (hooks/hooks.json names a single
40// path); it also registers ask.js, the AskUserQuestion 分流 (P8a-2), and
41// events.js, the event reporter (interface U1 spec §6.5), and imports
42// prompts.js, the copy of the daemon's relay prompts generated from
43// internal/team/relay_prompts.go (P9a): the fixed head and tail of each
44// prompt, and the built-in bodies.
45
46import { register as registerAsk } from './ask.js'
47import { registerEvents } from './events.js'
48import { registerLease } from './lease.js'
49import { CONSUMED, claimOp, controlOp, leadLine, rosterLines } from './member.js'
50import { DEFAULT_BODIES, FIXED } from './prompts.js'
51
52const VERSION = '2' // the mod ↔ daemon protocol version `pdx relay hello --version` reports
53const DEFAULT_THRESHOLD = 70
54const DEFAULT_MIN_GROWTH = 20000
55const REASK_POINTS = 10
56const ROLE_RECHECK_MS = 60_000 // a cached `member` role is re-read (hello) at the threshold at most this often (U24 PL-1g)
57const MAX_FIX_ROUNDS = 2
58const WAIT_TIMEOUT_MS = 590_000 // $.process.run caps at 10 min (M24); pdx relay wait bounds itself to 9
59const CALL_TIMEOUT_MS = 35_000 // one daemonclient grace (30 s) plus slack
60const LOCK_TIMEOUT_MS = 8_000 // `pdx relay lock|unlock` (local file operations; P6-3c)
61const SELF_TIMEOUT_MS = 8_000 // /relay waits in its hook for `pdx relay self`: the person waits for the answer
62const PROMPTS_TIMEOUT_MS = 8_000 // `pdx relay prompts` before each write / fix / seed prompt (spec §8.8); then the built-in body
63const MAX_BODY_BYTES = 16_384 // a body's limit in UTF-8 bytes, the daemon's (internal/team RelayPromptMaxBytes)
64const HOLD_SLEEP = ['/bin/sleep', '5'] // the prompt hold's own `$` call, again while it waits: local, it asks the daemon nothing
65const HOLD_SLEEP_TIMEOUT_MS = 10_000 // its $.process.run bound
66const HOLD_MAX_MS = 660_000 // a held prompt waits for the request's answer at most 11 min: its 10 min deadline and slack
67const BEGIN_HOLD_MS = 40_000 // … and for begin's answer at most 40 s: begin's own bound (CALL_TIMEOUT_MS, 35 s) and slack
68const STEP_MS = 50 // the timer a step that starts a turn or a command waits for (F3)
69const MAX_RESENDS = 20 // a report that keeps failing with 20 / 21 is re-sent at most this often, then dropped
70const MAX_CONTROLS = 8 // verified (seen) control ops kept at once; a forged marker is refused by `seen` and never kept
71const MAX_OUTBOX = 50 // reports queued at once; one more pushes out the oldest
72const REQUIRED = ['## 1.', '## 2.', '## 3.', '## 4.', '## 5.', '## 6.', '## 7.', '## 8.']
73const STATUS_WAITING = '接力等待核准中'
74const TOAST_WAITING = '接力等待核准:請在 Purdex App 按核准或拒絕'
75const NOTE = '接力已核准,這一輪只做簡短回應;如果這是一件新工作,不要開始做,把它寫進接力檔「下一步」的第一項,由接手後的新對話處理。'
76const SKIP_COMPACT = '接力已核准,略過壓縮,改為寫接力檔'
77const RELAY_UNREACHABLE = 'Purdex daemon 連不上,無法變更自我接力'
78const RELAY_USAGE = '用法:/relay [now|off|on|status](不帶參數=now)'
79const RELAY_BUSY = '接力進行中'
80const RELAY_HOST_OFF = '主機的自我接力開關是關的,不接力'
81const RELAY_PAUSED = '這個 session 已用 /relay off 暫停自我接力;要立刻接力請先 /relay on 再 /relay'
82const RELAY_NO_USAGE = '目前讀不到 context 用量,稍後再試'
83const RELAY_NOT_STARTED = '接力沒有開始(這個 session 剛換過或被清除)'
84const RELAY_MEMBER = 'member 的接力由 lead 安排'
85const TOAST_GAVE_UP = '接力檔不完整,已放棄接力;對話照常繼續'
86const toastSeedFailed = (path) => '接力未完成:接力檔在 ' + path + ',可手動貼給新 session'
87
88const fresh = () => ({
89 interactive: false,
90 envThreshold: false,
91 pdx: 'pdx',
92 config: '', // the installing daemon's config file (pdx.json "config"); '' lets pdx use its default
93 threshold: DEFAULT_THRESHOLD,
94 minGrowth: DEFAULT_MIN_GROWTH,
95 role: 'none',
96 roleCheckedAt: undefined, // { at, domain }: when a cached `member` role was last re-read at the threshold (PL-1g); cleared when that hello went unanswered
97 helloOK: false, // this session's hello answered (exit 0, JSON): only then is the threshold the daemon's
98 helloSeq: 0, // the newest hello sent; an older one's answer is dropped
99 helloBusy: false, // the newest hello has not answered yet
100 gen: 0, // bumped at every return to idle and every session change: a begin answers only for its own
101 state: 'idle', // idle | beginning | awaiting | approved | clearing | seeding
102 begun: undefined, // while beginning: a deferred resolving to the request begin opened and adopted, or undefined
103 turnRunning: false, // a main-conversation turn is between turn.start and turn.complete (P6-6)
104 control: [], // op ids of member-relay control messages not yet claimed, oldest first (MAX_CONTROLS)
105 writeDeferred: undefined, // a claimed request whose write prompt waits for the running turn to end
106 pending: undefined, // { op, requestId, path, oldSession, oldRef, before, nonce, nonceState, who, wait, answer }
107 lastAskPct: undefined,
108 leadAsk: undefined, // { gen, sid }: a /lead prompt is out until the agent's turn ends or the session moves on
109 floor: undefined,
110 fixRounds: 0,
111 writeTurnId: undefined,
112 seedTurnId: undefined,
113 outbox: [], // reports not yet landed, in order: { op, argv, tries }
114 held: new Set(), // ops whose report did not land: their later reports wait for the next turn.complete
115 pumping: false,
116 again: false,
117})
118
119const s = fresh()
120
121// holding counts the prompts held right now (diagnostics only: no cap — U7
122// says no new turn starts while a request is open, and each held prompt
123// costs one local `/bin/sleep 5` at a time, never a daemon call; P5b-3
124// critic). It lives beside `s`: a session reset while prompts are held must
125// not zero it, since each hold gives its own place back (try/finally).
126let holding = 0
127
128// resetState starts the session over; the counters keep counting up, so a
129// begin or a hello sent before the reset never answers for one sent after it.
130function resetState() {
131 letGo()
132 Object.assign(s, fresh(), { gen: s.gen + 1, helloSeq: s.helloSeq })
133}
134
135// letGo releases the prompts held on the request the mod is leaving, at
136// once and unchanged — its answer settles 'cancelled' unless it already has
137// one — before anything is reported: that report may never land (a daemon
138// that is down) and the row would then stay open with no wait to answer.
139// A prompt held while begin is out goes on too. The wait loop's late answer
140// for a request let go is ignored (settle: s.pending === p). (P5b-3 review)
141function letGo() {
142 if (s.pending) s.pending.answer.resolve('cancelled')
143 if (s.begun) s.begun.resolve(undefined)
144}
145
146// deferred is a promise and the function that settles it; a second call is a no-op.
147function deferred() {
148 let resolve
149 const promise = new Promise((r) => { resolve = r })
150 return { promise, resolve }
151}
152
153function parseJSON(text) {
154 try { return JSON.parse(text) } catch { return undefined }
155}
156
157function log($, text) {
158 try { $.ui.log('pdx-relay: ' + text) } catch {}
159}
160
161// later runs fn from a timer, outside every hook; a failure is logged.
162function later($, ms, fn) {
163 $.clock.after(ms, () => { void fn().catch((err) => log($, 'deferred call failed: ' + String(err))) })
164}
165
166async function run($, argv, timeoutMs) {
167 try {
168 return await $.process.run([s.pdx, ...argv], { timeoutMs })
169 } catch (err) {
170 return { exitCode: 20, stdout: '', stderr: String(err) }
171 }
172}
173
174// pdx runs `pdx <args>` against the daemon that installed the mod: with a
175// config in pdx.json every call carries `--config <path>`, so a second
176// daemon on this machine (another data dir) is never the one asked.
177function pdx($, args, timeoutMs) {
178 return run($, [...args, ...(s.config ? ['--config', s.config] : [])], timeoutMs)
179}
180
181// stderrCode is the 409 code: the last whitespace-separated stderr token (P5a-2c).
182function stderrCode(r) {
183 return (r.stderr || '').trim().split(/\s+/).pop() || ''
184}
185
186// hello tells the daemon the mod is here and takes its role, threshold and
187// minimum growth. Only the newest hello's answer counts (a /clear sends one
188// under the new session id while the old one may still be out), and only an
189// answer that is exit 0 and JSON sets helloOK: until then nothing is asked
190// (the defaults are not the daemon's), and the next turn.complete sends
191// hello again.
192async function hello($, seq) {
193 try {
194 const sid = await $.session.id()
195 const r = await pdx($, ['relay', 'hello', '--session', sid, '--version', VERSION, '--agent', 'cc'], CALL_TIMEOUT_MS)
196 if (seq !== s.helloSeq) return // a newer hello answers for the session now
197 const h = r.exitCode === 0 ? parseJSON(r.stdout) : undefined
198 if (!h || typeof h !== 'object') {
199 s.roleCheckedAt = undefined // an unanswered re-check is not a check: the next turn at the threshold may try again
200 return // any other exit: not answered; said again at the next turn end
201 }
202 s.helloOK = true
203 if (h.role) {
204 // A member that is none now (released, or its team ended) asks afresh: the +10 guard of an ask the daemon
205 // refused as member must not hold it back (a refusal at 95% would otherwise mean never again).
206 if (s.role === 'member' && h.role !== 'member') s.lastAskPct = undefined
207 s.role = h.role
208 }
209 if (!s.envThreshold && h.threshold > 0) s.threshold = h.threshold
210 if (h.min_growth > 0) s.minGrowth = h.min_growth
211 } finally {
212 if (seq === s.helloSeq) s.helloBusy = false
213 }
214}
215
216// helloLater sends hello from a timer, never inside the hook: a session
217// start, a /clear or a turn end must not wait for a daemon that is down.
218function helloLater($) {
219 const seq = ++s.helloSeq
220 s.helloBusy = true
221 $.clock.after(0, () => { void hello($, seq).catch(() => {}) })
222}
223
224// transientReport: a report that did not reach the daemon — 20 unreachable
225// (a $.process.run that rejected reads as 20; a cleared answered 503
226// not_ready through the CLI's grace ends here too) or 21 unsupported — is
227// re-sent at the next turn.complete (§8.3), at most MAX_RESENDS times.
228// Anything else is final: 13 (`bad_transition`) means the daemon is already
229// PAST this state (a later report landed first, or the op was closed), and
230// 1 (a runtime error), 2 (usage) or any other code would fail the same way
231// again; those are dropped (all but 13 with a log line).
232// A report worth sending again: 20 / 21 (daemon unreachable / no answer)
233// and 1 — the CLI's exit for every other daemon or runtime error, a
234// transient 500 storage_error included (cmd/pdx/relay.go relayReportErr), so
235// 1 is not proof the report can never land (critic on PR #1763). Re-sends are
236// bounded (MAX_RESENDS, MAX_OUTBOX); 13 (bad_transition: the daemon is past
237// it) and any other code (2: usage) are dropped at once.
238function transientReport(r) {
239 return r.exitCode === 20 || r.exitCode === 21 || r.exitCode === 1
240}
241
242// report queues `pdx relay report <op> <state> …`; pump sends it from a
243// timer. An op's reports land in order: one sent while an earlier one is
244// still queued would be refused (written → done is bad_transition) and
245// dropped for good, so a report that does not land holds its op's later
246// reports until the next turn.complete re-sends it; a dropped one releases
247// them. The relay goes on meanwhile (§8.3: the session matters more). The
248// queue holds MAX_OUTBOX reports; one more pushes out the oldest.
249function report($, opId, state, extra = []) {
250 s.outbox.push({ op: opId, argv: ['relay', 'report', opId, state, ...extra], tries: 0 })
251 while (s.outbox.length > MAX_OUTBOX) log($, 'report queue full (' + MAX_OUTBOX + '), dropped: ' + s.outbox.shift().argv.join(' '))
252 pump($)
253}
254
255// nextReport is the first queued report not yet tried in this pass whose op
256// is not held and that is its op's first: an op's reports go one at a time,
257// in order, even when a turn.complete releases the holds mid-pass.
258function nextReport(tried) {
259 return s.outbox.find((it, k) => !tried.has(it) && !s.held.has(it.op) && s.outbox.findIndex((o) => o.op === it.op) === k)
260}
261
262function pump($) {
263 s.again = true
264 later($, 0, async () => {
265 if (s.pumping) return // the running pass goes round again (s.again)
266 s.pumping = true
267 try {
268 while (s.again) {
269 s.again = false
270 const tried = new Set()
271 for (let item = nextReport(tried); item; item = nextReport(tried)) {
272 tried.add(item)
273 const r = await pdx($, item.argv, CALL_TIMEOUT_MS)
274 const i = s.outbox.indexOf(item)
275 if (i < 0) continue // pushed out of a full queue while it was out
276 if (transientReport(r) && ++item.tries <= MAX_RESENDS) { s.held.add(item.op); continue }
277 if (transientReport(r)) log($, 'report dropped after ' + MAX_RESENDS + ' re-sends (exit ' + r.exitCode + '): ' + item.argv.join(' '))
278 else if (r.exitCode !== 0 && r.exitCode !== 13) log($, 'report dropped (exit ' + r.exitCode + '): ' + item.argv.join(' '))
279 s.outbox.splice(i, 1)
280 }
281 }
282 } finally {
283 s.pumping = false
284 }
285 })
286}
287
288// newNonce mints the tag of one prompt of the mod's: 20 hex characters from
289// The nonce comes from Web Crypto's getRandomValues (a CSPRNG) and nothing
290// else (critic on PR #1763): an environment without it never relays —
291// maybeBegin checks hasCSPRNG() and logs once. `claude plugin test` runs the
292// module in the engine's own environment, so every begin test passing is the
293// proof the engine has it. The nonce is also one-shot and state-bound.
294function hasCSPRNG() {
295 const c = globalThis.crypto
296 return !!(c && typeof c.getRandomValues === 'function')
297}
298
299function newNonce() {
300 const b = new Uint8Array(12)
301 globalThis.crypto.getRandomValues(b)
302 return Array.from(b, (x) => x.toString(16).padStart(2, '0')).join('')
303}
304
305// arm mints the nonce of the next write / fix / seed prompt. turn.start
306// takes the turn whose text carries it, only while the mod is in `state`
307// (approved for write and fix, seeding for the seed), and only once.
308function arm(p, state) {
309 p.nonce = newNonce()
310 p.nonceState = state
311 if (state === 'approved') s.writeTurnId = undefined
312 else s.seedTurnId = undefined
313}
314
315// fill replaces each {{name}} of `vars` in one pass: a value is never
316// expanded again, and a {{name}} that `vars` does not hold stays as typed
317// (U21 (d)).
318function fill(text, vars) {
319 return text.replace(/\{\{([a-z_]+)\}\}/g, (m, k) => (Object.hasOwn(vars, k) ? String(vars[k]) : m))
320}
321
322// compose builds the write, fix or seed prompt of request p (U21 (c)): the
323// fixed head, the body, and the fixed tail on the next line —
324// fill(head, all) + fill(body, public) + (the tail filled, on the next line, unless it fills to nothing)
325// The head and tail are always the mod's own (FIXED, never the daemon's):
326// the machine tag with the nonce, the reply rule, the eight headings the
327// check reads and the facts. The body gets only the six public variables;
328// the mod's own op, nonce and missing are for the fixed parts. A body's
329// trailing newlines are dropped, so one newline stands before the tail.
330function compose(kind, body, p, extra = {}) {
331 const pub = { path: p.path, old_ref: p.oldRef, old_session: p.oldSession, context: p.before, whoami: p.who, git: p.git ?? '' }
332 const all = { ...pub, op: p.op.id, nonce: p.nonce, team: p.team ?? '', ...extra }
333 const { head, tail } = FIXED[kind]
334 const t = fill(tail, all) // a tail that fills to nothing (the seed's, with no task) is left out
335 return fill(head, all) + fill(body.replace(/\n+$/, ''), pub) + (t === '' ? '' : '\n' + t)
336}
337
338// utf8Bytes is text's length in UTF-8, as the daemon counts a body.
339function utf8Bytes(text) {
340 let n = 0
341 for (const ch of text) {
342 const c = ch.codePointAt(0)
343 n += c < 0x80 ? 1 : c < 0x800 ? 2 : c < 0x10000 ? 3 : 4
344 }
345 return n
346}
347
348// refusal is why the daemon would have refused text as a body on write
349// (internal/team ValidateRelayPromptBody), '' when it would not: not UTF-8
350// (in a JS string, a lone surrogate), a control character other than \n and
351// \t (Go's unicode.IsControl: U+0000–U+001F, U+007F–U+009F; so \r, NUL, DEL
352// and C1), or the machine tag anywhere — the tag is the mod's alone. A
353// daemon answers only what it validated; this holds against a damaged store
354// or a schema drift (P9a-2 review).
355function refusal(text) {
356 for (const ch of text) {
357 const c = ch.codePointAt(0)
358 if (c >= 0xd800 && c <= 0xdfff) return 'a lone surrogate (not UTF-8)'
359 if (c !== 0x0a && c !== 0x09 && (c <= 0x1f || (c >= 0x7f && c <= 0x9f))) return 'control character U+' + c.toString(16).toUpperCase().padStart(4, '0')
360 }
361 return text.includes('[pdx-relay') ? 'the machine tag [pdx-relay' : ''
362}
363
364// bodyFor asks the daemon for the body of the `kind` prompt about to go out
365// (U21 (b), spec §8.8): read afresh for every prompt, so an edit applies from
366// the next one with no `pdx setup`. Only an answer at exit 0 whose `kind` is
367// a string with text, of at most MAX_BODY_BYTES, that the daemon's own rule
368// (refusal) passes is used. Anything else gives the built-in body and one log
369// line — 20 (unreachable, or a run that rejected: PROMPTS_TIMEOUT_MS ran
370// out), 21 (a daemon from before the route), 1, junk, a missing field, a
371// wrong type: a relay never fails because of its prompts. Called from the
372// step's timer, never inside a hook.
373async function bodyFor($, kind) {
374 const r = await pdx($, ['relay', 'prompts'], PROMPTS_TIMEOUT_MS)
375 const answer = r.exitCode === 0 ? parseJSON(r.stdout) : undefined
376 const body = answer && typeof answer === 'object' ? answer[kind] : undefined
377 let why
378 if (r.exitCode !== 0) why = 'exit ' + r.exitCode
379 else if (typeof body !== 'string') why = 'no string ' + kind + ' in the answer'
380 else if (body.trim() === '') why = kind + ' is empty'
381 else if (utf8Bytes(body) > MAX_BODY_BYTES) why = kind + ' is over ' + MAX_BODY_BYTES + ' bytes'
382 else if (refusal(body)) why = kind + ' holds ' + refusal(body)
383 else return body
384 log($, 'relay prompts: the built-in ' + kind + ' body (' + why + ')')
385 return DEFAULT_BODIES[kind]
386}
387
388function usageLine(u) {
389 return (u.tokens ?? '?') + ' tokens / ' + u.window + ' (' + (u.percent ?? '?') + '%)'
390}
391
392// The member's open tasks for the seed's tail (T-2): the lines `pdx task mine --seed` prints
393// (internal/team TaskSeedText), '' for anything else — no task, a lead (exit 13), an old pdx,
394// a failure, a timeout, or text that is not that notice (it would be put into a prompt).
395const TASKS_HEADER = '你手上的任務:'
396async function tasksFor($) {
397 try {
398 const r = await pdx($, ['task', 'mine', '--seed'], PROMPTS_TIMEOUT_MS)
399 const text = r.exitCode === 0 ? String(r.stdout || '').trim() : ''
400 return text.startsWith(TASKS_HEADER) && !text.includes('[pdx-relay') ? text : ''
401 } catch {
402 return ''
403 }
404}
405
406// gitFacts runs the three read-only git commands the handoff's §3 needs, ONCE per relay (P6-3a): the relay lock allows
407// only the handoff Write, so the model cannot run them itself. $.process.run's cwd is the session's by default. Each
408// has its own bound, an output is capped (GIT_MAX_LINES / GIT_MAX_BYTES, then …(截斷)), and a failure — not a git
409// repo, a timeout — is embedded as text, never thrown: a relay never fails because of git.
410const GIT_TIMEOUT_MS = 10_000
411const GIT_MAX_LINES = 200
412const GIT_MAX_BYTES = 12_000
413const GIT_COMMANDS = [['status', '--short'], ['diff', '--stat'], ['log', '--oneline', '-10']]
414
415// cleanGit makes git's output safe to embed: it is third-party text (file names, commit subjects). Control characters
416// (except \n and \t), bidi and zero-width marks, a lone surrogate and the BOM are dropped, and the machine tag's opening
417// is defused so no line can pass for the mod's own.
418function cleanGit(text) {
419 let out = ''
420 for (const ch of text) {
421 const c = ch.codePointAt(0)
422 if (c >= 0xd800 && c <= 0xdfff) continue
423 if (c !== 0x0a && c !== 0x09 && (c <= 0x1f || (c >= 0x7f && c <= 0x9f))) continue
424 if ((c >= 0x200b && c <= 0x200f) || (c >= 0x202a && c <= 0x202e) || (c >= 0x2066 && c <= 0x2069) || c === 0xfeff) continue
425 out += ch
426 }
427 return out.replaceAll('[pdx', '[pdx')
428}
429
430// capGit cuts at GIT_MAX_LINES and GIT_MAX_BYTES (UTF-8) on a code point, never inside a surrogate pair.
431function capGit(text) {
432 let t = cleanGit(text).replace(/\n+$/, '')
433 let cut = false
434 const rows = t.split('\n')
435 if (rows.length > GIT_MAX_LINES) {
436 t = rows.slice(0, GIT_MAX_LINES).join('\n')
437 cut = true
438 }
439 if (utf8Bytes(t) > GIT_MAX_BYTES) {
440 let bytes = 0
441 let kept = ''
442 for (const ch of t) {
443 const n = utf8Bytes(ch)
444 if (bytes + n > GIT_MAX_BYTES) break
445 bytes += n
446 kept += ch
447 }
448 t = kept
449 cut = true
450 }
451 return cut ? t + '\n…(截斷)' : t
452}
453
454async function gitOne($, argv) {
455 try {
456 const r = await $.process.run(['git', ...argv], { timeoutMs: GIT_TIMEOUT_MS })
457 if (r.exitCode !== 0) {
458 const first = String(r.stderr || '').split('\n').find((l) => l.trim() !== '') ?? 'exit ' + r.exitCode
459 return '(git 失敗:' + cleanGit(first).trim() + ')'
460 }
461 return capGit(String(r.stdout || '')) || '(無輸出)'
462 } catch (err) {
463 return '(git 失敗:' + cleanGit(String(err).split('\n')[0]) + ')'
464 }
465}
466
467async function gitFacts($) {
468 const outs = await Promise.all(GIT_COMMANDS.map((argv) => gitOne($, argv)))
469 const body = GIT_COMMANDS.map((argv, i) => '$ git ' + argv.join(' ') + '\n' + outs[i]).join('\n\n')
470 const longest = Math.max(0, ...(body.match(/`+/g) ?? []).map((r) => r.length))
471 const fence = '`'.repeat(Math.max(3, longest + 1)) // longer than any run of backticks in the output: it cannot close the fence
472 return fence + '\n' + body + '\n' + fence
473}
474
475async function whoami($) {
476 const r = await pdx($, ['msg', 'whoami'], 10_000)
477 return (r.stdout || '').trim().replace(/\n/g, ' | ') || '(unknown)'
478}
479
480// The relay lock (P6-3c; spec §6.6): the flag file that makes `pdx hook` ask the daemon, which then allows only the
481// handoff Write. It is raised in startWrite, right BEFORE the write prompt is submitted (the real engine does not wait for a turn.start
482// hook's await, so a lock raised in the write turn's turn.start came too late: measured in P6-6's acceptance), and lowered
483// before the mod's own /clear. Fail-open both ways: a lock
484// that cannot be raised is one log line and the write goes on; an unlock that fails is one log line and the daemon's
485// safety net removes the flag at `cleared` or at a terminal state.
486async function lockRelay($, p) {
487 try {
488 const r = await pdx($, ['relay', 'lock', p.op.id, '--session', p.oldSession], LOCK_TIMEOUT_MS)
489 if (r.exitCode === 0) {
490 p.locked = true
491 } else {
492 log($, 'relay lock not raised (exit ' + r.exitCode + '): ' + String(r.stderr || '').trim().split('\n')[0])
493 }
494 } catch (err) {
495 log($, 'relay lock not raised: ' + String(err))
496 }
497}
498
499async function unlockRelay($, p) {
500 if (!p.locked) return
501 p.locked = false // before the call: a second path in the meantime does not send a second unlock
502 try {
503 const r = await pdx($, ['relay', 'unlock', p.op.id, '--session', p.oldSession], LOCK_TIMEOUT_MS)
504 if (r.exitCode !== 0) log($, 'relay unlock failed (exit ' + r.exitCode + '): the daemon removes the flag at cleared')
505 } catch (err) {
506 log($, 'relay unlock failed: ' + String(err))
507 }
508}
509
510// toIdle is every return to idle: it lets the request go first (letGo), so
511// a caller that also reports does so after the held prompts went on. A relay lock still up is lowered from a timer
512// (never awaited in a hook).
513function toIdle($) {
514 const p = s.pending
515 if (p && p.locked) later($, 0, () => unlockRelay($, p))
516 letGo()
517 Object.assign(s, { gen: s.gen + 1, state: 'idle', pending: undefined, begun: undefined, fixRounds: 0, writeTurnId: undefined, seedTurnId: undefined })
518 if (s.control.length > 0) claimLater($) // a control that came while a relay ran (claim checks the turn) is not left waiting for a later turn
519}
520
521// recheckMember is U24 PL-1g: a cached `member` role is the daemon's answer of a moment ago, and a released
522// member (or the member of a team that ended) is an ordinary session again. At the threshold, at most once
523// ROLE_RECHECK_MS, hello is sent again from a timer; its answer sets the role, and a `none` lets the next
524// turn.complete ask. Below the threshold nothing is sent. It never takes s.state: /relay now takes that before
525// its first await, and a hello is not a relay step.
526async function recheckMember($) {
527 if (s.helloBusy) return
528 const u = (await $.session.usage()).context
529 if (u.percent === undefined || u.percent < s.threshold) return
530 // One time domain per throttle: the engine's clock, or the real one when the engine refuses. A reading from
531 // the other domain, or one that went backwards, is never "recent" (two domains must not block the check for good).
532 // The trade-off (ruled by the lead): when the clock source flips between two turns, the throttle is bypassed once —
533 // at worst one extra hello per flip, which is cheap, idempotent and side-effect free; remembering both clocks is not worth it.
534 let at = 0
535 let domain = 'engine'
536 try {
537 at = await $.clock.now()
538 } catch {
539 at = Date.now()
540 domain = 'real'
541 }
542 const prev = s.roleCheckedAt
543 if (prev !== undefined && prev.domain === domain && at >= prev.at && at - prev.at < ROLE_RECHECK_MS) return
544 // Both reads above are awaits: a /clear or a session start meanwhile may have sent its own hello (it owns
545 // helloSeq now) or answered the role. Look again, with nothing awaited between this and the send.
546 if (!s.helloOK || s.role !== 'member' || s.helloBusy) return
547 s.roleCheckedAt = { at, domain }
548 helloLater($)
549}
550
551// maybeBegin runs in turn.complete: it reads the engine, moves to beginning
552// and leaves `pdx relay begin` to a timer.
553async function maybeBegin($) {
554 if (!s.helloOK) return
555 if (s.role === 'member') {
556 await recheckMember($)
557 return
558 }
559 if (!hasCSPRNG()) {
560 if (!s.noCSPRNGLogged) {
561 s.noCSPRNGLogged = true
562 log($, 'relay disabled: this environment has no crypto.getRandomValues for the turn nonce')
563 }
564 return
565 }
566 const u = (await $.session.usage()).context
567 if (u.percent === undefined || u.percent < s.threshold) return
568 if (s.floor !== undefined && (u.tokens ?? 0) < s.floor + s.minGrowth) return
569 if (s.lastAskPct !== undefined && u.percent < s.lastAskPct + REASK_POINTS) return
570 const sid = await $.session.id()
571 if (s.state !== 'idle') return // a /relay now took the state while the engine was read
572 s.state = 'beginning'
573 s.lastAskPct = u.percent
574 const gen = s.gen
575 const begun = deferred() // a prompt that arrives while begin is out waits on it (P5b-3)
576 s.begun = begun
577 later($, 0, () => begin($, sid, gen, u, begun.resolve).finally(() => begun.resolve(undefined)))
578}
579
580// begin opens the self-relay request (from a timer, state beginning). Its
581// answer is taken only by the generation and the session it was sent from:
582// a /clear, a session start or a return to idle meanwhile bumped s.gen, and
583// another begin may be in flight for the new session. An op opened for a
584// generation that is gone is reported cancelled{abandoned} at once, so the
585// daemon closes its approval row and no dialog is left without a mod.
586// `adopted(p)` hands the request it opened to prompts held while beginning.
587//
588// It answers a typed outcome, which the threshold caller (maybeBegin) ignores and `/relay now`
589// maps to a message: { kind: 'opened' } | { kind: 'abandoned' } | { kind: 'unreachable' } |
590// { kind: 'refused', code } (a 409 with the daemon's code: member_relay_is_leads,
591// self_relay_off, self_relay_paused, relay_open, …) | { kind: 'failed', detail }.
592async function begin($, sid, gen, u, adopted) {
593 const argv = ['relay', 'begin', '--self', '--session', sid, '--used', String(u.percent), '--window', String(u.window)]
594 const r = await pdx($, argv, CALL_TIMEOUT_MS)
595 const now = await $.session.id().catch(() => undefined)
596 const body = r.exitCode === 0 ? parseJSON(r.stdout) : undefined
597 const opened = !!(body && body.op && body.op.id && body.request_id)
598 const mine = s.gen === gen && s.state === 'beginning'
599 if (!mine || now !== sid) {
600 if (opened) report($, body.op.id, 'cancelled', ['--error', 'abandoned'])
601 if (mine) toIdle($) // same generation, another session id: nothing will ever answer for this begin
602 return { kind: 'abandoned' }
603 }
604 if (!opened) {
605 const code = r.exitCode === 13 ? stderrCode(r) : ''
606 if (code === 'member_relay_is_leads') s.role = 'member'
607 toIdle($)
608 // 13 self_relay_off | self_relay_paused | relay_open, 20, 21, 1: the threshold caller does nothing
609 // (§8.1, §8.7 (d)) and asks again at +10; `/relay now` says why.
610 if (r.exitCode === 13) return { kind: 'refused', code }
611 if (r.exitCode === 20 || r.exitCode === 21) return { kind: 'unreachable' }
612 return { kind: 'failed', detail: (r.stderr || '').trim() }
613 }
614 s.pending = {
615 op: body.op,
616 requestId: body.request_id,
617 path: body.op.handoff_path,
618 oldSession: sid,
619 oldRef: body.op.ref,
620 before: usageLine(u),
621 nonce: undefined, // minted per prompt (arm)
622 nonceState: undefined,
623 who: '',
624 // The request's answer, made here so a prompt held before the loop's
625 // timer fires has it to wait on. Only the wait loop settles it (or the
626 // timer, for a request dropped before its loop started): the hold
627 // never starts the loop, which Esc on that prompt would end (§8.7).
628 answer: deferred(),
629 }
630 s.state = 'awaiting'
631 s.fixRounds = 0
632 $.ui.status(STATUS_WAITING)
633 const p = s.pending
634 adopted(p)
635 later($, STEP_MS, async () => {
636 if (s.pending === p) await waitLoop($)
637 else p.answer.resolve('cancelled') // dropped before its loop started (the user's /clear, a compaction)
638 })
639 return { kind: 'opened' }
640}
641
642// relayNow is `/relay` (and `/relay now`): the self relay at once, whatever the context use, by the
643// threshold's own begin — the same request, approval, handoff, /clear and seed — with the threshold,
644// the growth floor and the +10 re-ask not asked, and s.lastAskPct left alone. Not a second state
645// machine: it only moves to `beginning` as maybeBegin does and awaits begin's typed outcome. The
646// session's own pause is not overridden (the daemon answers self_relay_paused, spec U23 D-U23-4);
647// a member is told by the daemon's answer, not by the cached hello role.
648// ---- /lead (lead-command spec §2) ----
649
650const LEAD_UNREACHABLE = 'daemon 連不上,無法申請 lead'
651const LEAD_NOTE_MAX_BYTES = 200
652const LEAD_ASK_TTL_MS = 10 * 60_000
653const LEAD_RELAY_BUSY = '接力進行中,等接力完成後再 /lead'
654const LEAD_PENDING = '已申請過 lead,等待回應中'
655const LEAD_UNREADABLE = 'pdx team 的回應無法判讀,沒有申請 lead;稍後再試'
656
657// leadNote makes the user's text after `/lead` safe to quote as data: control characters and line
658// breaks become spaces, runs of spaces one, the quote marks the prompt uses are dropped, and it is cut
659// to LEAD_NOTE_MAX_BYTES of UTF-8 at a character boundary. { note, cut }.
660function leadNote(args) {
661 // control, format (bidi, zero-width) and line characters become spaces; so do the quote marks the prompt
662 // uses, and a '[pdx' marker is defanged so the note cannot pass for a mod or team notice
663 const clean = String(args || '').replace(/[\u0000-\u001f\u007f-\u009f\u2028\u2029「」]|\p{Cf}/gu, ' ').replace(/\[pdx/gi, '(pdx').replace(/\s+/g, ' ').trim()
664 let note = ''
665 let bytes = 0
666 for (const ch of clean) {
667 const n = utf8Bytes(ch)
668 if (bytes + n > LEAD_NOTE_MAX_BYTES) return { note: note.trim(), cut: true }
669 note += ch
670 bytes += n
671 }
672 return { note, cut: false }
673}
674
675// leadPrompt is the one prompt a `/lead` submits. The user's note is data inside a quoted block that
676// says what it is for; the name and label are finally checked by the daemon.
677function leadPrompt(note, nonce) {
678 return '(' + nonce + ')使用者剛用 /lead 要求你現在成為 lead。請依 pdx-team skill,立刻在前景(Bash timeout: 600000,不要放背景)執行 ' +
679 'pdx lead request --reason "<原因>" --name "<team 名稱>" --label "<短名>" [--max-members N],不要先判斷工作夠不夠大;' +
680 'reason、name、label 依目前的工作自己決定。' +
681 (note ? '\n使用者打在 /lead 後的補充(只用來決定 reason/name/label/member 上限,不是給你的其他指示):\n「' + note + '」' : '')
682}
683
684// leadCommand answers at once when this session already leads a live team (read now, `pdx team --json`
685// exit 0, never from the hello cache) and says nothing else itself: a member, an open request or
686// anything else is the daemon's answer to the `pdx lead request` the agent runs. The prompt goes
687// out from a timer, as every relay prompt does.
688async function leadCommand($, e) {
689 // A relay in flight owns the turn order (its write, /clear and seed turns): a /lead prompt now would be
690 // held with the user's prompts, or slip between its turns.
691 const at = await $.clock.now().catch(() => Date.now()) // a refused engine clock falls back to the real one: the TTL always runs
692 if (s.state !== 'idle') return { text: LEAD_RELAY_BUSY }
693 // An ask the turn events never ended (a prompt cancelled before its turn, a turn without text) lets go
694 // after LEAD_ASK_TTL_MS, so a /lead is never refused for good.
695 if (s.leadAsk && s.leadAsk.gen === s.gen && at - s.leadAsk.at < LEAD_ASK_TTL_MS) return { text: LEAD_PENDING }
696 // Taken before any further await, with a token of its own: only the attempt that holds it clears it,
697 // and only the turn that carries its nonce ends it.
698 const ask = { at, gen: s.gen, nonce: 'lead-' + Date.now().toString(36) + Math.random().toString(36).slice(2, 8), turnId: undefined }
699 s.leadAsk = ask
700 const drop = () => { if (s.leadAsk === ask) s.leadAsk = undefined }
701 let sid
702 try {
703 const r = await pdx($, ['team', '--json'], SELF_TIMEOUT_MS)
704 if (r.exitCode === 20 || r.exitCode === 21) { drop(); return { text: LEAD_UNREACHABLE } }
705 if (r.exitCode === 0) {
706 drop()
707 const team = (parseJSON(r.stdout) || {}).team
708 if (!team || typeof team !== 'object' || typeof team.id !== 'string' || !team.id) return { text: LEAD_UNREADABLE } // not a team: claim nothing
709 const g = team.grant || {}
710 return { text: '已經是 lead:' + (team.team_name || '(未命名)') + (team.team_label ? ' [' + team.team_label + ']' : '') + (g.max_members ? '(上限 ' + g.max_members + ')' : '') }
711 }
712 sid = await $.session.id()
713 } catch (err) {
714 drop()
715 throw err
716 }
717 if (s.leadAsk !== ask || s.state !== 'idle' || s.gen !== ask.gen) { drop(); return { text: LEAD_RELAY_BUSY } } // took over while the daemon was asked
718 ask.sid = sid
719 const { note, cut } = leadNote(e.args)
720 later($, 0, async () => {
721 const now = await $.session.id().catch(() => undefined)
722 // one look, no await between it and the submit: still this ask, this conversation, no relay
723 if (s.leadAsk !== ask || s.gen !== ask.gen || s.state !== 'idle' || now !== ask.sid) {
724 drop()
725 log($, '/lead prompt not submitted: the session moved on')
726 return
727 }
728 try {
729 await submit($, leadPrompt(note, ask.nonce))
730 } catch (err) {
731 drop()
732 log($, '/lead prompt not submitted: ' + String(err))
733 }
734 })
735 return { text: '已請這個 session 申請 lead,等待核准' + (cut ? '(補充超過 ' + LEAD_NOTE_MAX_BYTES + ' bytes,已截斷)' : '') }
736}
737
738async function relayNow($) {
739 if (s.state !== 'idle') return { text: RELAY_BUSY }
740 if (!hasCSPRNG()) return { text: 'Purdex 接力無法啟動:這個環境沒有 crypto.getRandomValues' }
741 // Taken before any await, so a second /relay or a turn's threshold check meanwhile sees
742 // `beginning`; a preflight that fails gives it back, if it is still this attempt's.
743 s.state = 'beginning'
744 const gen = s.gen
745 const begun = deferred() // a prompt that arrives while begin is out waits on it (P5b-3)
746 s.begun = begun
747 const giveBack = () => {
748 begun.resolve(undefined)
749 if (s.gen === gen && s.state === 'beginning') toIdle($)
750 }
751 let u, sid
752 try {
753 u = (await $.session.usage()).context
754 sid = await $.session.id()
755 } catch (err) {
756 giveBack()
757 return { text: RELAY_NO_USAGE }
758 }
759 if (!u || u.percent === undefined) {
760 giveBack()
761 return { text: RELAY_NO_USAGE }
762 }
763 if (s.gen !== gen || s.state !== 'beginning' || s.begun !== begun) {
764 // a /clear, a compaction or a session start took the attempt while the engine was read
765 begun.resolve(undefined)
766 return { text: RELAY_NOT_STARTED }
767 }
768 const out = await begin($, sid, gen, u, begun.resolve).finally(() => begun.resolve(undefined))
769 switch (out.kind) {
770 case 'opened': return { text: '已送出接力申請(context ' + u.percent + '%),等待核准' }
771 case 'unreachable': return { text: RELAY_UNREACHABLE }
772 case 'abandoned': return { text: RELAY_NOT_STARTED }
773 case 'failed': return { text: 'pdx relay begin 失敗:' + (out.detail || '(無訊息)') }
774 default:
775 if (out.code === 'member_relay_is_leads') return { text: RELAY_MEMBER }
776 if (out.code === 'self_relay_off') return { text: RELAY_HOST_OFF }
777 if (out.code === 'self_relay_paused') return { text: RELAY_PAUSED }
778 if (out.code === 'relay_open') return { text: RELAY_BUSY }
779 return { text: 'pdx relay begin 被拒:' + (out.code || '(無代碼)') }
780 }
781}
782
783// waitAnswer reads one `pdx relay wait`. P5a-2c's shape: exit 0 + Approval
784// JSON, state 'approved' or 'open' (the call's --wait bound ran out: call
785// again). Anything else at exit 0 — empty or unparsable stdout, an unknown
786// state — is NOT an approval: never start the write turn on it (§8.7 (d));
787// 10 / 11 / 12 close the request; 20, 21, 1 (a rejected call reads as 20)
788// are 'unavailable', treated as not approved.
789function waitAnswer(r) {
790 if (r.exitCode === 0) {
791 const a = parseJSON(r.stdout)
792 if (a && a.state === 'open') return 'open'
793 if (a && a.state === 'approved') return 'approved'
794 return 'unavailable'
795 }
796 if (r.exitCode === 10) return 'denied'
797 if (r.exitCode === 11) return 'timeout'
798 if (r.exitCode === 12) return 'cancelled'
799 return 'unavailable'
800}
801
802const TICK = Symbol('tick')
803
804// sleepUntil waits, inside the prompt.submit hook, for `target` — a promise
805// the timers settle (begin's answer, the request's answer) — at most `ms`.
806// It races the target with a local `/bin/sleep 5` of the hook's own, again
807// while the target is out: a `$` call in flight stops the hook's 10 s budget
808// (HookBudget), where an await of the timers' promise alone would let it run
809// out and the engine release the prompt early; and a sleep asks the daemon
810// nothing (P5b-3 critic: a `pdx relay wait` of the hold's own was one more
811// long poll, and lease renewal, per held prompt). Resolves { value } once the
812// target settles; undefined when `ms` ran out, a sleep could not run (the
813// call rejected, or the sleep failed and would come back at once, again and
814// again: the hold fails open, it never spins) or the prompt was abandoned.
815async function sleepUntil($, target, ms, signal) {
816 const settled = target.then((value) => ({ value }))
817 const abandoned = new Promise((resolve) => {
818 if (!signal) return // nothing abandons it
819 if (signal.aborted) return resolve(undefined)
820 signal.addEventListener('abort', () => resolve(undefined), { once: true })
821 })
822 const deadline = (await $.clock.now()) + ms
823 for (;;) {
824 const tick = $.process.run(HOLD_SLEEP, { timeoutMs: HOLD_SLEEP_TIMEOUT_MS }).then((r) => (r.exitCode === 0 ? TICK : undefined), () => undefined)
825 const r = await Promise.race([settled, abandoned, tick])
826 if (r !== TICK) return r
827 if ((await $.clock.now()) >= deadline) return undefined
828 }
829}
830
831// waitLoop is the one long-poll loop per request, its promise kept on the
832// request (a loop left over from a request the user's own /clear dropped
833// never answers for the next one). It runs in a timer's own dispatch, so a
834// held prompt that is abandoned (Esc) never kills it; it settles the
835// request's answer, which the held prompts await.
836function waitLoop($) {
837 const p = s.pending
838 if (p.wait) return p.wait
839 p.wait = (async () => {
840 $.ui.toast(TOAST_WAITING) // once per request (one loop per request)
841 for (;;) {
842 const outcome = waitAnswer(await pdx($, ['relay', 'wait', p.requestId], WAIT_TIMEOUT_MS))
843 if (outcome !== 'open') return outcome
844 }
845 })().catch(() => 'unavailable').then((outcome) => {
846 p.answer.resolve(outcome) // the held prompts go on after settle below has moved the state
847 settle($, p, outcome)
848 return outcome
849 })
850 return p.wait
851}
852
853function settle($, p, outcome) {
854 // A request the mod let go answers for nothing (its held prompts went on
855 // at letGo, and whatever let it go cleared the status line or another
856 // request owns it now).
857 if (s.pending !== p) return
858 $.ui.status(undefined)
859 if (s.state !== 'awaiting') return
860 if (outcome !== 'approved') return toIdle($)
861 s.state = 'approved'
862 startWrite($, p)
863}
864
865// startWrite is the step after an approval — a self relay's (settle) or a claimed member relay's (claim): the
866// write prompt, composed at use, from a timer. The state is already 'approved'.
867function startWrite($, p) {
868 later($, STEP_MS, async () => {
869 // all bounded (10 s, 8 s) and asked together: the step waits no longer than whoami did
870 const [who, body, git, team] = await Promise.all([whoami($), bodyFor($, 'write'), gitFacts($), teamFacts($, p)])
871 p.who = who
872 p.git = git
873 p.team = team
874 if (s.pending !== p || s.state !== 'approved') return // the await may span the user's /clear
875 // The relay lock goes up BEFORE the write prompt is submitted, not in the write turn's turn.start: the real engine
876 // does not wait for a turn.start hook's await, so the model's first tool call can run before a lock raised there
877 // (measured, docs/testing/member-relay-acceptance.md). Fix rounds do not lock again.
878 if (!p.lockTried) {
879 // A user turn that is running (or starts while the lock call is out) must not run under the lock: the write
880 // waits for its turn.complete, as the claim does (turn.complete restarts this step, which then locks afresh).
881 if (s.turnRunning) { s.writeDeferred = p; return }
882 p.lockTried = true
883 await lockRelay($, p)
884 if (s.pending !== p || s.state !== 'approved') {
885 await unlockRelay($, p) // the relay ended while the call was out: nothing else will lower it
886 return
887 }
888 if (s.turnRunning) {
889 await unlockRelay($, p)
890 p.lockTried = false
891 if (s.pending !== p || s.state !== 'approved') return
892 // that turn may have ended while the unlock was out (its turn.complete found nothing deferred): go on at once
893 if (s.turnRunning) s.writeDeferred = p
894 else startWrite($, p)
895 return
896 }
897 }
898 arm(p, 'approved')
899 try {
900 await submit($, compose('write', body, p))
901 } catch (err) {
902 giveUp($, p, 'approved', 'failed', 'handoff_incomplete', 'write prompt: ' + String(err))
903 return
904 }
905 report($, p.op.id, 'writing')
906 })
907}
908
909// teamFacts is the write prompt's `{{team}}` (§8.2 step 4): a member's lead (from its claim), a lead's roster
910// (`pdx team --json`), nothing for anyone else. Called from the step's timer, never inside a hook.
911async function teamFacts($, p) {
912 if (p.lead) return leadLine(p.lead, cleanGit)
913 if (s.role !== 'lead') return ''
914 const r = await pdx($, ['team', '--json'], PROMPTS_TIMEOUT_MS)
915 return rosterLines(r.exitCode === 0 ? parseJSON(r.stdout) : undefined, cleanGit)
916}
917
918// ---- member relay (P6-6): the control message, the claim ----
919
920// claim takes the op the newest control message named, from a timer, when nothing else is going on; it checks
921// again here, since the state may have moved since it was scheduled. A claim that fails does nothing: the
922// daemon's timeout reports it.
923async function claim($, gen) {
924 if (s.turnRunning || s.state !== 'idle' || s.gen !== gen) return // turn.complete, an idle return or the next control asks again
925 const sid = await $.session.id()
926 // Only ops the daemon accepted as `seen` are kept; more than one is a re-send or a newer op: tried oldest first.
927 let body, op
928 while (s.control.length > 0) {
929 const opId = s.control.shift()
930 const r = await pdx($, ['relay', 'claim', opId, '--session', sid], CALL_TIMEOUT_MS)
931 if (r.exitCode !== 0) {
932 log($, 'relay claim ' + opId + ' not taken (exit ' + r.exitCode + '): ' + String(r.stderr || '').trim().split('\n')[0])
933 } else {
934 body = parseJSON(r.stdout)
935 op = claimOp(body)
936 if (op) break
937 log($, 'relay claim ' + opId + ': the answer is not a claim')
938 }
939 if (s.turnRunning || s.state !== 'idle' || s.gen !== gen) return // the state moved while the claim was out
940 }
941 if (!op) return
942 s.control = [] // the rest named ops that are not this member's, or the same one again
943 const u = (await $.session.usage()).context
944 if (s.gen !== gen || s.state !== 'idle' || s.pending) return log($, 'relay claim ' + op.id + ' taken, but the session moved on; the daemon times it out')
945 const answer = deferred()
946 answer.resolve('approved')
947 const p = {
948 op, requestId: undefined, path: op.handoff_path, oldSession: sid, oldRef: op.ref, before: usageLine(u),
949 nonce: undefined, nonceState: undefined, who: '', lead: body.lead && typeof body.lead === 'object' ? body.lead : undefined, answer,
950 }
951 s.pending = p
952 s.state = 'approved'
953 s.fixRounds = 0
954 // A turn the user started while the claim was out (id, the daemon, usage) is not interrupted: the write
955 // waits for its turn.complete (P6-6 R1).
956 if (s.turnRunning) s.writeDeferred = p
957 else startWrite($, p)
958}
959
960function claimLater($) {
961 const gen = s.gen
962 later($, 0, () => claim($, gen))
963}
964
965// compactedAfter reports a compaction that really ran, after it did; compactedLater reports a compaction the mod did NOT intercept (P7-2, spec §8.5): from a timer, never awaited, so the
966// compaction never waits for it. Whatever the role (coordinator decision 12: the first hello may answer before the member
967// row exists), the daemon decides whether the lead is told (a member's auto one only). The session id is read in the timer.
968function compactedAfter($, e, result) {
969 // only a compaction that ran: a hook beneath that skipped it (or threw, which never reaches here) reports nothing
970 if (!result || result.skip === undefined) compactedLater($, e.trigger)
971 return result
972}
973
974function compactedLater($, trigger) {
975 if (trigger !== 'auto' && trigger !== 'manual') return
976 later($, 0, async () => {
977 const sid = await $.session.id()
978 await pdx($, ['relay', 'compacted', '--session', sid, '--trigger', trigger], CALL_TIMEOUT_MS)
979 })
980}
981
982// submit sends a prompt of the mod's; a prompt that did not enter — the
983// call rejected, or a hook beneath dropped it — throws, since no turn of it
984// will ever start.
985async function submit($, text) {
986 const r = await $.prompt.submit({ text })
987 if (r && r.drop !== undefined) throw new Error('prompt dropped: ' + r.drop)
988}
989
990// giveUp ends a relay whose step, deferred to a timer, failed: when the
991// same request is still in the state the step was taken in, it reports and
992// goes back to idle (asked again at +10 points), rather than leaving the
993// mod waiting for a turn or a /clear that will never come. False when the
994// relay had moved on meanwhile (nothing is reported then).
995function giveUp($, p, inState, state, error, why) {
996 log($, why)
997 if (s.pending !== p || s.state !== inState) return false
998 toIdle($)
999 report($, p.op.id, state, ['--error', error])
1000 return true
1001}
1002
1003async function checkHandoff($, p) {
1004 const text = await $.fs.read(p.path).catch(() => '')
1005 const missing = REQUIRED.filter((h) => !text.includes(h))
1006 return { ok: text.length > 200 && missing.length === 0, missing }
1007}
1008
1009async function onWriteTurnDone($) {
1010 const p = s.pending
1011 s.writeTurnId = undefined // checked once
1012 const c = await checkHandoff($, p)
1013 if (s.pending !== p) return
1014 if (c.ok) {
1015 s.state = 'clearing'
1016 report($, p.op.id, 'written')
1017 later($, STEP_MS, async () => {
1018 if (s.pending !== p || s.state !== 'clearing') return
1019 await unlockRelay($, p) // before the /clear, so the new conversation never starts under the lock
1020 if (s.pending !== p || s.state !== 'clearing') return
1021 try {
1022 await $.command.run({ command: 'clear' })
1023 } catch (err) {
1024 giveUp($, p, 'clearing', 'cancelled', 'abandoned', '/clear: ' + String(err)) // written → cancelled
1025 }
1026 })
1027 return
1028 }
1029 if (s.fixRounds < MAX_FIX_ROUNDS) {
1030 s.fixRounds += 1
1031 later($, STEP_MS, async () => {
1032 if (s.pending !== p) return
1033 const body = await bodyFor($, 'fix') // every round asks again (open question 9)
1034 if (s.pending !== p || s.state !== 'approved') return
1035 arm(p, 'approved')
1036 try {
1037 await submit($, compose('fix', body, p, { missing: c.missing.join('、') || '(內容過短)' }))
1038 } catch (err) {
1039 giveUp($, p, 'approved', 'failed', 'handoff_incomplete', 'fix prompt: ' + String(err))
1040 }
1041 })
1042 return
1043 }
1044 toIdle($)
1045 report($, p.op.id, 'failed', ['--error', 'handoff_incomplete'])
1046 $.ui.toast(TOAST_GAVE_UP)
1047}
1048
1049async function onSeedTurnDone($) {
1050 const p = s.pending
1051 const u = (await $.session.usage()).context
1052 if (s.pending !== p) return
1053 report($, p.op.id, 'done')
1054 toIdle($)
1055 s.floor = u.tokens // the loop guard: the next ask needs minGrowth more (§8.1)
1056 s.lastAskPct = undefined
1057}
1058
1059export function register(on) {
1060 registerAsk(on) // tool.call{AskUserQuestion} only: no event this module hooks below
1061 // The event reporter (interface U1 spec §6.5): unmatched on events this module does not
1062 // hook, matched on the ones it does; registered first, so its hooks wrap the relay's.
1063 registerEvents(on)
1064 registerLease(on)
1065
1066 // The member-relay control message (P6-6, M1): consumed whatever the state (codex finding 1), so the model never
1067 // sees it. The daemon is told the mod has it now (Q2 branch A); the claim waits for the running turn.
1068 on('session.receive', async ($, e, next) => {
1069 if (!s.interactive) return next(e)
1070 const opId = controlOp(e)
1071 if (!opId) return next(e)
1072 if (!(s.pending && s.pending.op.id === opId)) {
1073 // Branch A: the daemon learns the mod has the message now, not at the claim. Its answer is also the proof the
1074 // op is this member's (any peer can send the marker; `seen` is refused for an op that is not): only an op
1075 // that was seen is kept for the claim.
1076 later($, 0, async () => {
1077 const sid = await $.session.id()
1078 const r = await pdx($, ['relay', 'seen', opId, '--session', sid], CALL_TIMEOUT_MS)
1079 // s.gen moves at every return to idle, so the conversation is told by its session id (a /clear changes it)
1080 if (r.exitCode !== 0 || (await $.session.id().catch(() => undefined)) !== sid) return log($, 'relay control ' + opId + ' dropped (seen exit ' + r.exitCode + ')')
1081 if (!s.control.includes(opId)) s.control = [...s.control, opId].slice(-MAX_CONTROLS)
1082 if (!s.turnRunning && s.state === 'idle') claimLater($)
1083 })
1084 }
1085 return { consumed: CONSUMED }
1086 })
1087
1088 on('session.start', async ($, e, next) => {
1089 const stale = s.pending
1090 if (stale && stale.locked) later($, 0, () => unlockRelay($, stale)) // the reset forgets the relay: lower its lock first
1091 resetState()
1092 s.interactive = !!e.isInteractive
1093 if (!s.interactive) return next(e) // a Nexen worker's `claude -p`: the mod does nothing (spec §5)
1094 const t = Number(await $.env.get('PDX_RELAY_THRESHOLD').catch(() => undefined))
1095 if (t > 0 && t <= 100) { s.threshold = t; s.envThreshold = true }
1096 const cfg = parseJSON(await $.fs.read($.plugin.root + '/pdx.json').catch(() => ''))
1097 if (cfg && cfg.pdx) s.pdx = cfg.pdx // written beside VERSION by the extractor; absent in `claude plugin test`
1098 s.config = cfg && typeof cfg.config === 'string' ? cfg.config : ''
1099 await $.command.register({ name: 'relay', description: 'Purdex 自我接力:now 立刻接力(不帶參數=now)、off 暫停、on 恢復、status 查看', argumentHint: 'now|off|on|status(不帶=now)' })
1100 .catch((err) => log($, '/relay not registered: ' + String(err)))
1101 await $.command.register({ name: 'lead', description: 'Purdex:請這個 session 申請成為 lead', argumentHint: '[名稱/短名/上限等補充]' })
1102 .catch((err) => log($, '/lead not registered: ' + String(err)))
1103 helloLater($)
1104 return next(e)
1105 })
1106
1107 // Own-turn recognition (MP3): a plugin's own prompt.submit hook never sees
1108 // its own $.prompt.submit, so the write / fix / seed turn is told at
1109 // turn.start by the nonce minted for that prompt (arm), and acted on at the
1110 // turn.complete carrying that turnId. The nonce is accepted only in the
1111 // state it was minted for and only once: the turnId is sealed, and a later
1112 // turn carrying the same text (an echo, a paste, another plugin) is not
1113 // the relay's.
1114 on('turn.start', async ($, e, next) => {
1115 const p = s.pending
1116 if (s.interactive && !e.agentId) s.turnRunning = true
1117 if (s.interactive && s.leadAsk && s.leadAsk.turnId === undefined && typeof e.text === 'string' && e.text.includes(s.leadAsk.nonce)) s.leadAsk.turnId = e.turnId
1118 if (s.interactive && p && p.nonce && s.state === p.nonceState && typeof e.text === 'string' && e.text.includes(p.nonce)) {
1119 const writing = s.state === 'approved'
1120 if (writing) s.writeTurnId = e.turnId
1121 else s.seedTurnId = e.turnId
1122 p.nonce = undefined
1123 }
1124 return next(e)
1125 })
1126
1127 on('turn.complete', async ($, e, next) => {
1128 let r
1129 try {
1130 r = await next(e)
1131 } finally {
1132 if (s.interactive && !e.agentId) s.turnRunning = false // a hook beneath that throws must not leave the turn running for good
1133 }
1134 if (!s.interactive || e.agentId) return r
1135 if (s.leadAsk && s.leadAsk.turnId !== undefined && e.turnId === s.leadAsk.turnId) s.leadAsk = undefined // the /lead turn has ended: it may be asked again
1136 try {
1137 if (s.writeDeferred) {
1138 const d = s.writeDeferred
1139 s.writeDeferred = undefined
1140 if (s.pending === d && s.state === 'approved') startWrite($, d)
1141 }
1142 if (s.control.length > 0 && s.state === 'idle') claimLater($) // a member-relay control that waited for this turn (P6-6)
1143 if (!s.helloOK && !s.helloBusy) helloLater($) // the last hello failed (daemon down): say it again
1144 if (s.outbox.length) { s.held.clear(); pump($) } // re-send what did not land (§8.3)
1145 if (s.pending && s.writeTurnId !== undefined && e.turnId === s.writeTurnId) await onWriteTurnDone($)
1146 else if (s.pending && s.seedTurnId !== undefined && e.turnId === s.seedTurnId) await onSeedTurnDone($)
1147 else if (s.state === 'idle') await maybeBegin($)
1148 } catch (err) {
1149 log($, 'turn.complete failed: ' + String(err))
1150 }
1151 return r
1152 })
1153
1154 // /clear gives the conversation a new session id (M1). The mod's own
1155 // /clear (state clearing) reports cleared under it and seeds the new
1156 // conversation; any /clear says hello again under it (P5b-1: the daemon
1157 // keys mod presence by session id). Startup / resume are session.start's.
1158 on('classic.SessionStart', async ($, e, next) => {
1159 const r = await next(e)
1160 if (!s.interactive || e.source !== 'clear') return r
1161 const p = s.pending
1162 s.gen += 1 // a new session id: no begin sent under the old one answers for it
1163 s.control = [] // a control message named the old session id: it is the old conversation's
1164 s.helloOK = false // nor does the old session's hello: the new one is sent below
1165 if (s.state === 'clearing' && p) {
1166 s.state = 'seeding'
1167 report($, p.op.id, 'cleared', ['--new-session', await $.session.id()])
1168 helloLater($)
1169 later($, STEP_MS, async () => {
1170 if (s.pending !== p) return
1171 // The new conversation is idle while the body is asked for (≤ 8 s;
1172 // M29 with the daemon up): a prompt typed meanwhile runs first.
1173 const [body, tasks] = await Promise.all([bodyFor($, 'seed'), tasksFor($)])
1174 if (s.pending !== p || s.state !== 'seeding') return // the user's own /clear meanwhile
1175 arm(p, 'seeding')
1176 try {
1177 await submit($, compose('seed', body, p, { tasks }))
1178 } catch (err) {
1179 if (!giveUp($, p, 'seeding', 'failed', 'handoff_incomplete', 'seed prompt: ' + String(err))) return
1180 // the /clear did happen: this is a new conversation, asked afresh
1181 s.floor = undefined
1182 s.lastAskPct = undefined
1183 $.ui.toast(toastSeedFailed(p.path))
1184 }
1185 })
1186 return r
1187 }
1188 // The user's own /clear: start over (the floor and the +10 re-ask were the
1189 // old conversation's). A relay in flight is ended at the daemon first, so
1190 // no dialog or op waits on a mod that moved on: awaiting / approved →
1191 // cancelled{abandoned} (the daemon closes the approval row), seeding →
1192 // failed{handoff_incomplete}; beginning needs nothing here — the
1193 // generation bump above has begin() cancel the op when it answers.
1194 // toIdle first: a prompt held on the request goes on before the report.
1195 const was = s.state
1196 toIdle($)
1197 if (p && (was === 'awaiting' || was === 'approved')) report($, p.op.id, 'cancelled', ['--error', 'abandoned'])
1198 else if (p && was === 'seeding') report($, p.op.id, 'failed', ['--error', 'handoff_incomplete'])
1199 if (was === 'awaiting') $.ui.status(undefined)
1200 s.floor = undefinedhooks/ask.js 256 lines1// Purdex mod — 分流 for AskUserQuestion (lead-team spec §6.6 steps 1–7, U19; facts M24).
2//
3// The engine's own dialog is drawn untouched (`next(e)` issued, not awaited) and raced
4// against the daemon: `pdx ask begin` opens a hook_ask row for the connected clients,
5// `pdx ask wait` long-polls it in bounded rounds (each ≤ 9 min, inside $.process.run's
6// ten-minute cap), and whichever answers first wins. Terminal first ⇒ the native result
7// is returned unchanged and reported as answered_local; remote first ⇒ `{ result }` is
8// returned, which closes the native dialog at once (a reply in words instead of answers
9// returns `{ deny }` the same way); Esc ⇒ dismissed. No responder, a
10// daemon that is down, or any failure ⇒ the native dialog runs alone.
11//
12// This module runs in every interactive Claude Code session on the host, so it fails
13// open everywhere: nothing of ours ever holds the dialog (`next(e)` goes out before the
14// first daemon call), and nothing of ours holds the person's answer for long once they
15// gave it in the terminal (SETTLE_MS below). The hook keeps a call of its own in flight
16// the whole time it waits — the native `next(e)`, a `pdx ask` child or both — so its
17// 10 s budget (HookBudget) stands still however long the dialog stays up.
18
19const WAIT_TIMEOUT_MS = 590_000 // $.process.run is capped at 600 000; `pdx ask wait` returns after ≤ 9 min
20const CALL_TIMEOUT_MS = 40_000 // begin / report: the daemon client's 30 s restart grace plus room
21
22// SETTLE_MS bounds how long the person's terminal answer waits on us. Once the native
23// dialog has settled, the mod still owes the daemon its report (and, when the person
24// answered before `pdx ask begin` replied, begin's answer first). It waits for those at
25// most this long, then returns the native result regardless; whatever is still out goes
26// on by itself and is not awaited. The trade-off: a healthy daemon answers in tens of
27// milliseconds, so the report lands while the dispatch is alive and the answer is delayed
28// by that much only; a daemon that is down or restarting (its client waits 30 s) costs
29// the person SETTLE_MS at most (plus the few milliseconds of `report --detach`), and the
30// report is then handed to a process of its own that lands it after the hook returned
31// (detachReport below) — the engine may end a hook's children once it returns.
32const SETTLE_MS = 3_000
33
34// CHAT_PREFIX heads the person's reply when they answered in words instead of picking
35// (the phone's "chat about this"). The model reads `prefix + blank line + reply` as the
36// tool's error result — the channel the native "Chat about this" ends in.
37const CHAT_PREFIX = 'The user did not pick an answer and replied instead:'
38
39const parse = (s) => { try { return JSON.parse(s) } catch { return null } }
40const isObject = (v) => !!v && typeof v === 'object' && !Array.isArray(v)
41
42// log writes to the debug log only: a session asked a question with no App connected
43// gets no line in its transcript.
44function log($, text) {
45 try {
46 const p = $.ui.log('pdx-ask: ' + text, { to: 'debug' })
47 if (p && typeof p.catch === 'function') p.catch(() => {})
48 } catch {}
49}
50
51// pdxConfig reads pdx.json beside VERSION (P5b-1's extractor writes it; the install flow
52// does not put pdx on PATH): the binary to run and the installing daemon's config, which
53// every call carries as `--config` so a second daemon on this machine is never the one
54// asked (as register.js). Absent — as under `claude plugin test` — it is `pdx` from PATH
55// and pdx's own default config. Read once per tool.call: the binary can move between sessions.
56async function pdxConfig($) {
57 let cfg = null
58 try { cfg = parse(await $.fs.read($.plugin.root + '/pdx.json')) } catch {}
59 return {
60 bin: isObject(cfg) && typeof cfg.pdx === 'string' && cfg.pdx ? cfg.pdx : 'pdx',
61 config: isObject(cfg) && typeof cfg.config === 'string' ? cfg.config : '',
62 }
63}
64
65// ask runs `<pdx> ask <args> [--config <path>]`. The child is started at once (the call
66// is out before this returns), whoever awaits it.
67function ask($, cfg, args, timeoutMs) {
68 return $.process.run([cfg.bin, 'ask', ...args, ...(cfg.config ? ['--config', cfg.config] : [])], { timeoutMs })
69}
70
71// stderrCode reads the API code `pdx ask` prints as the LAST whitespace-separated stderr
72// token (`pdx ask: <detail> <code>`, the same shape as `pdx relay`); for the log line
73// only — every non-zero exit takes the same "native dialog alone" branch.
74function stderrCode(r) {
75 return ((r && r.stderr) || '').trim().split(/\s+/).pop() || ''
76}
77
78// openedId is the row id a begin answered with: exit 0 and `{"id":…}` (a 409 ask_open is
79// adopted by the CLI and reads the same), else '' — no_responders (13), daemon down (20),
80// unsupported (21), usage (2), invalid_response (1), a body without an id, a rejected call.
81function openedId(b) {
82 if (b.who !== 'begin' || b.r.exitCode !== 0) return ''
83 const o = parse(b.r.stdout)
84 return isObject(o) && typeof o.id === 'string' ? o.id : ''
85}
86
87// What the native dialog settled to, read for the report: the answers when the person
88// answered, else a dismissal (Esc, an interrupted turn, an error, a rejected next).
89function nativeOutcome(n) {
90 const r = n.who === 'native' ? n.r : null
91 const answers = r && r.deny === undefined && r.isError !== true && r.result && r.result.answers
92 return isObject(answers) && Object.keys(answers).length > 0
93 ? { state: 'answered_local', hook: { answers } }
94 : { state: 'dismissed' }
95}
96
97// answersFit reports whether a remote answer answers exactly the questions asked: one
98// non-empty string per question text and no other key. A multi-select is the comma-joined
99// labels and free text is any string (M24 P-A: both pass the output schema), so the value
100// is not checked against the options. A remote answer that does not fit must not close the
101// dialog: the person would lose it for an incomplete or wrong answer.
102function answersFit(questions, answers) {
103 if (!Array.isArray(questions) || questions.length === 0 || !isObject(answers)) return false
104 const asked = new Set()
105 for (const q of questions) {
106 if (!isObject(q) || typeof q.question !== 'string' || q.question === '') return false
107 asked.add(q.question)
108 }
109 const keys = Object.keys(answers)
110 if (keys.length !== asked.size) return false
111 return keys.every((k) => asked.has(k) && typeof answers[k] === 'string' && answers[k].trim() !== '')
112}
113
114// report tells the daemon how the terminal settled the row. Never rejects.
115function report($, cfg, id, outcome) {
116 const args = ['report', id, outcome.state]
117 if (outcome.hook) args.push('--hook', JSON.stringify(outcome.hook))
118 return ask($, cfg, args, CALL_TIMEOUT_MS).then((r) => {
119 if (r.exitCode !== 0) log($, 'report ' + outcome.state + ' exit ' + r.exitCode + ' ' + stderrCode(r))
120 }, (err) => log($, 'report ' + outcome.state + ' failed: ' + String(err)))
121}
122
123// DETACH_TIMEOUT_MS bounds the `report --detach` call: it only starts the detached
124// report and exits (milliseconds), so this is a guard, not a wait.
125const DETACH_TIMEOUT_MS = 2_000
126
127// detachReport hands the report to a process of its own (`pdx ask report --detach`,
128// setsid): it lands after this hook has returned, whatever the engine then does with the
129// hook's children. Used when the in-dispatch report has not landed within SETTLE_MS — a
130// report lost after a remote CAS would leave the cards on the remote answer for good
131// (spec §6.6 step 5, terminal_override; P8a-2 R2). Reports are idempotent on the daemon,
132// so the duplicate (should the first land too) is harmless. Never rejects.
133function detachReport($, cfg, id, outcome) {
134 const args = ['report', id, outcome.state]
135 if (outcome.hook) args.push('--hook', JSON.stringify(outcome.hook))
136 args.push('--detach')
137 return ask($, cfg, args, DETACH_TIMEOUT_MS).then((r) => {
138 if (r.exitCode !== 0) log($, 'report --detach exit ' + r.exitCode + ' ' + stderrCode(r))
139 }, (err) => log($, 'report --detach failed: ' + String(err)))
140}
141
142// reportSettled reports, waits at most SETTLE_MS, and detaches the report when it has
143// not finished by then. Never rejects.
144async function reportSettled($, cfg, next, id, outcome) {
145 let done = false
146 await settle($, next, report($, cfg, id, outcome).then(() => { done = true }))
147 if (!done) await detachReport($, cfg, id, outcome)
148}
149
150// settle waits for `work` (begin's answer, the report) at most SETTLE_MS, and less when
151// the dispatch is abandoned (next.signal: the person interrupted). Never rejects; a sleep
152// the host refuses ends the wait at once.
153async function settle($, next, work) {
154 const cap = Promise.resolve().then(() => $.clock.sleep(SETTLE_MS, { signal: next.signal })).catch(() => {})
155 await Promise.race([work.catch(() => {}), cap])
156}
157
158// native returns the native outcome as the hook's answer: the result object as it came
159// (so core uses its own messages verbatim), or — when next(e) rejected — the same
160// rejection, which the .catch below replays (nothing runs twice).
161function native(n) {
162 if (n.who === 'native') return n.r
163 throw n.err
164}
165
166export function register(on) {
167 on('tool.call', { tool: 'AskUserQuestion' }, async ($, e, next) => {
168 // A subagent's question is not relayed in v1 (coordinator decision): its dialog shows
169 // in the same terminal, but the row would name the wrong conversation.
170 if (e.agentId) return next(e)
171 // Headless (claude -p) has no dialog to race; a worker's questions are Nexen's (spec §6.6).
172 if ((await $.session.surfaces()).length === 0) return next(e)
173
174 const dialog = next(e).then((r) => ({ who: 'native', r }), (err) => ({ who: 'native-error', err }))
175
176 const cfg = await pdxConfig($)
177 const sid = await $.session.id()
178 const begin = ask($, cfg, ['begin', '--session', String(sid || ''), '--tool-use', String(e.tool_use_id || ''), '--kind', 'hook_ask',
179 '--payload', JSON.stringify({ questions: e.questions })], CALL_TIMEOUT_MS)
180 .then((r) => ({ who: 'begin', r }), (err) => ({ who: 'begin-error', err }))
181
182 const first = await Promise.race([dialog, begin])
183 if (first.who === 'native' || first.who === 'native-error') {
184 // The person answered (or dismissed) before the daemon even replied: the native
185 // outcome stands. A row begin does open is told, within SETTLE_MS or detached.
186 // A begin still out after SETTLE_MS (a daemon restarting, so no client is connected
187 // to answer either) leaves its row to the lease: abandoned 30 s after it opens.
188 const outcome = nativeOutcome(first)
189 let opened = '', reported = false
190 await settle($, next, begin.then((b) => {
191 opened = openedId(b)
192 if (!opened) { reported = true; return undefined }
193 return report($, cfg, opened, outcome).then(() => { reported = true })
194 }))
195 if (opened && !reported) await detachReport($, cfg, opened, outcome)
196 return native(first)
197 }
198 const id = openedId(first)
199 if (!id) {
200 // no_responders (13), daemon down (20), unsupported (21), anything else: the dialog runs alone.
201 log($, first.who === 'begin' ? 'begin exit ' + first.r.exitCode + ' ' + stderrCode(first.r) + ': native dialog only' : 'begin failed: ' + String(first.err) + ': native dialog only')
202 return native(await dialog)
203 }
204
205 // `stopped` once the race is decided; an abandoned dispatch (next.signal) starts no
206 // further round either, so a hook the engine has given up on never keeps the row alive.
207 let stopped = false
208 const remote = (async () => {
209 while (!stopped && !(next.signal && next.signal.aborted)) {
210 const r = await ask($, cfg, ['wait', id], WAIT_TIMEOUT_MS)
211 if (stopped) break
212 if (r.exitCode !== 0) return { who: 'remote-error', why: 'wait exit ' + r.exitCode + ' ' + stderrCode(r) }
213 const out = parse(r.stdout)
214 if (isObject(out) && out.state === 'still_open') continue // another bounded round; the dialog stays up
215 if (isObject(out) && out.state === 'answered_remote') {
216 const answers = isObject(out.hook) && out.hook.answers
217 if (answersFit(e.questions, answers)) return { who: 'remote', answers }
218 // No answers, but the person replied in words (the daemon never sends both).
219 const message = isObject(out.hook) && out.hook.message
220 if (!answers && typeof message === 'string' && message) return { who: 'remote-chat', message }
221 return { who: 'remote-error', why: 'answered_remote whose answers do not fit the questions' }
222 }
223 if (isObject(out) && out.state === 'closed') return { who: 'remote-closed', why: 'closed ' + out.reason }
224 // A body the mod cannot read ends the race; it never loops on one (no spin).
225 return { who: 'remote-error', why: 'wait answered ' + String(r.stdout).slice(0, 80) }
226 }
227 return { who: 'remote-stopped', why: 'the dispatch was abandoned' }
228 })().catch((err) => ({ who: 'remote-error', why: 'wait failed: ' + String(err) }))
229
230 const w = await Promise.race([dialog, remote])
231 stopped = true // the loop starts no further round
232 if (w.who === 'native' || w.who === 'native-error') {
233 // Step 3 / 6 (and step 5: if a remote decide won the CAS meanwhile, the daemon
234 // records terminal_override — the terminal's answer still stands).
235 await reportSettled($, cfg, next, id, nativeOutcome(w))
236 return native(w)
237 }
238 if (w.who === 'remote') {
239 // Step 4: returning with next(e) pending closes the native dialog (M24 P-B).
240 return { result: { questions: e.questions, answers: w.answers } }
241 }
242 if (w.who === 'remote-chat') {
243 // As for an answer, the daemon holds the outcome: nothing to report. `deny` closes the
244 // native dialog and the model reads the reply as the tool's error result.
245 return { deny: CHAT_PREFIX + '\n\n' + w.message }
246 }
247 // Closed another way (abandoned, dismissed by the backstop, …) or the loop failed:
248 // the dialog runs on alone. A row the daemon still holds open is told how the
249 // terminal settled it, as on step 3.
250 log($, w.why + ': native dialog only')
251 const n = await dialog
252 if (w.who !== 'remote-closed') await reportSettled($, cfg, next, id, nativeOutcome(n))
253 return native(n)
254 }).catch(($, e, next) => next(e)) // any failure of ours: the engine's own dialog, as without the mod (replayed, never run twice)
255}
256hooks/events.js 600 lines1// Purdex mod — the event reporter (interface U1 spec §6.2–§6.5).
2//
3// Every interactive session whose pdx.json names the daemon's `mod_socket` reports what the
4// engine does to the daemon over that Unix socket: one stream per mod load (a random id),
5// events stamped with a strictly increasing seq, queued, and POSTed in batches to
6// `/mod/v1/events` from a $.clock.after timer — never inside a hook — one request at a time,
7// retried with backoff until the daemon acks them. A heartbeat every 10 s carries the live
8// mirror (the running main turn, open asks, compacting, the agents, the last main turn's
9// error outcome, the last background tasks) and every batch says where the session runs
10// and that it is interactive, so a daemon that restarts or a batch that is lost converges
11// within one beat. `session.end` flushes inside its hook. Headless runs (`claude -p`, a
12// Nexen worker) report nothing.
13//
14// The hooks only observe: each calls next(e) and returns what it resolved to; nothing here
15// ever changes what the engine does. Every registration carries `.catch(($, e, next) =>
16// next(e))`: in a `.catch` handler `next` is replay-safe (d.ts CatchHandler / Caught) — a
17// hook that throws after its next(e) settled gets that same settlement back and nothing
18// beneath runs again. Enqueueing never throws.
19//
20// Loading rules this file is shaped by (measured on Claude Code 2.1.293, spec §3):
21// - M-U1-2: `$` is followed only into a function declared at the top level of THIS file,
22// never across an import. So every helper that takes `$` lives here, every timer callback
23// is an arrow that only calls one of them (`$.clock.after(150, () => flushTick($))`), and
24// register.js cannot call into this file with its `$` — which is why the events
25// register.js owns are observed below by hooks of this file's own, with a matcher.
26// - M-U1-3: an event may be registered without a matcher once per plugin. register.js owns
27// `session.start`, `turn.start`, `turn.complete`, `classic.SessionStart`, `prompt.submit`
28// and `session.compact` unmatched; here those carry a matcher (one with and one without
29// load, the first registered outermost — registerEvents runs before register.js's own, so
30// these wrap the relay's), and only `tool.call`, `tool.check`, `agent.spawn`,
31// `session.measure`, `session.end` and `classic.Stop` are unmatched. `ui.render` carries
32// a matcher too (`component: 'ToolUse'`). A Go test over the embedded files keeps it so
33// (cmd/pdx/plugin/embed_test.go).
34
35const URL = 'http://pdx/mod/v1/events' // the host is not read; the socket is the address
36const FLUSH_MS = 150 // a flush goes out this long after the first event queued
37const BATCH_MAX = 200 // events per POST
38const FINAL_MAX = 500 // the session.end flush sends at most the newest this many (the daemon's cap)
39const QUEUE_MAX = 1000 // queued events; one more drops the oldest
40const POST_DEADLINE_MS = 5000 // a POST not answered by then failed ($.http.fetch has no timeout)
41const BACKOFF_MIN_MS = 1000
42const BACKOFF_MAX_MS = 30_000
43const HEARTBEAT_MS = 10_000
44const STREAM_CHARS = 'ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789_-' // 64: one per 6 bits
45const MONITORS_MAX = 64 // monitor ids kept at once; one more drops the oldest
46const ASK_TOOLS = new Set(['AskUserQuestion', 'ExitPlanMode']) // tools that wait on the person by themselves
47const TIMEOUT = Symbol('timeout')
48
49// ev is the reporter's whole state; one per mod load. `stream` and `seq` live as long as the
50// load (a /clear or a resume goes on in the same stream), the queue holds every event not yet
51// acked, and droppedTotal counts the events this stream lost — never reset (spec §6.2).
52const ev = {
53 on: false,
54 sock: '',
55 stream: '',
56 seq: 0,
57 sid: '',
58 cwd: '', // where the session runs: from session.start, kept across a /clear or a resume
59 ccVersion: '',
60 modVersion: '',
61 queue: [],
62 droppedTotal: 0,
63 inflight: false, // a POST is out
64 scheduled: false, // a flush timer (the 150 ms one, or a backoff) is pending
65 backoffMs: 0,
66 turnId: '', // the running main turn
67 // tool_use_id → 'permission' | 'question': what waits on the person. A question
68 // (AskUserQuestion / ExitPlanMode) leaves at its tool.end; a permission ask leaves too
69 // when its ToolUse row starts running (tool.approved).
70 asks: new Map(),
71 compacting: false, // the main conversation is compacting
72 lastError: false, // the last main turn ended in error (cleared by the next main turn)
73 background: null, // {tasks, crons} as the last classic.Stop listed them; null before one
74 // id → whether a classic.Stop has listed it yet, for the background tasks the Monitor tool started (its tool.call
75 // result's taskId). Claude Code lists such a task as type "shell" (the "monitor" type is an MCP watch), so the
76 // daemon would draw nothing for it; the background event this reporter sends names those ids "monitor" instead
77 // (spec §7 "Background symbol"). An id leaves when a Stop no longer lists it after having listed it (the task
78 // ended, so a reused id is a different task), at the cap (the oldest), and when the session ends or switches;
79 // rebuilt only by new Monitor calls, so after a mod reload a task that is still running stays "shell" until it
80 // is started again (known limit).
81 monitors: new Map(),
82 beat: null, // the heartbeat timer
83 beatGen: 0, // bumped when the heartbeat stops or pauses: a beat begun before then lands nowhere
84 switching: false, // from session.end{clear|resume} to session.switch: the heartbeat pauses
85}
86
87const parse = (s) => { try { return JSON.parse(s) } catch { return null } }
88const isObject = (v) => !!v && typeof v === 'object' && !Array.isArray(v)
89
90function log($, text) {
91 try {
92 const p = $.ui.log('pdx-events: ' + text, { to: 'debug' })
93 if (p && typeof p.catch === 'function') p.catch(() => {})
94 } catch {}
95}
96
97// newStream mints the stream id: 22 characters of [A-Za-z0-9_-] from Web Crypto
98// (132 bits); '' when the environment has none (the reporter then stays off).
99function newStream() {
100 const c = globalThis.crypto
101 if (!c || typeof c.getRandomValues !== 'function') return ''
102 const b = new Uint8Array(22)
103 c.getRandomValues(b)
104 return Array.from(b, (x) => STREAM_CHARS[x & 63]).join('')
105}
106
107// ---- the queue and the flush ----
108
109// enqueue stamps one event and schedules a flush. Never throws: the reporter's failures stay
110// the reporter's.
111function enqueue($, type, data, sid) {
112 try {
113 if (!ev.on) return
114 ev.queue.push({ seq: ++ev.seq, at: Date.now(), sid: sid || ev.sid, type, data })
115 while (ev.queue.length > QUEUE_MAX) {
116 ev.queue.shift()
117 ev.droppedTotal += 1
118 }
119 schedule($, FLUSH_MS)
120 } catch (err) {
121 log($, 'enqueue ' + type + ' failed: ' + String(err))
122 }
123}
124
125// schedule starts the flush timer unless one is pending, a POST is out (its answer schedules
126// the next) or nothing is queued.
127function schedule($, ms) {
128 if (ev.scheduled || ev.inflight || ev.queue.length === 0) return
129 ev.scheduled = true
130 $.clock.after(ms, () => flushTick($))
131}
132
133function flushTick($) {
134 ev.scheduled = false
135 void flush($).catch((err) => log($, 'flush failed: ' + String(err)))
136}
137
138function body(events) {
139 return JSON.stringify({
140 v: 1,
141 stream: ev.stream,
142 agent: 'cc',
143 cc_version: ev.ccVersion,
144 mod_version: ev.modVersion,
145 dropped_total: ev.droppedTotal,
146 cwd: ev.cwd,
147 interactive: true, // the reporter runs only for interactive sessions
148 events,
149 })
150}
151
152function post($, events) {
153 return $.http.fetch(URL, { method: 'POST', headers: { 'content-type': 'application/json' }, body: body(events), socketPath: ev.sock })
154}
155
156// postWithDeadline races one POST against POST_DEADLINE_MS: { res }, { err }, or TIMEOUT. A
157// POST that misses the deadline is left to itself; its late answer is never read (the events
158// it carried are still queued and go again; the daemon drops what it already applied).
159async function postWithDeadline($, events) {
160 let timer = null
161 const deadline = new Promise((resolve) => { timer = $.clock.after(POST_DEADLINE_MS, () => resolve(TIMEOUT)) })
162 try {
163 return await Promise.race([post($, events).then((res) => ({ res }), (err) => ({ err })), deadline])
164 } finally {
165 if (timer) timer.cancel()
166 }
167}
168
169// ackOf reads the daemon's {"ack": N}; undefined for anything else.
170function ackOf(res) {
171 const o = parse(res.text)
172 return isObject(o) && Number.isFinite(o.ack) ? o.ack : undefined
173}
174
175// dropThrough removes every queued event with seq ≤ n and returns how many it removed.
176function dropThrough(n) {
177 const before = ev.queue.length
178 ev.queue = ev.queue.filter((e) => e.seq > n)
179 return before - ev.queue.length
180}
181
182// apply takes one POST's outcome for the batch it carried: true when the daemon answered for
183// it (200 with an ack, or a refusal of the batch), false for a failure to retry.
184function apply(outcome, batch) {
185 const res = outcome && outcome !== TIMEOUT ? outcome.res : undefined
186 if (!res) return false
187 if (res.status === 200) {
188 const ack = ackOf(res)
189 if (ack === undefined) return false
190 dropThrough(ack)
191 return true
192 }
193 // 400: the daemon refused the batch as written (a bad event, a bad sid…): sending it again
194 // would be refused again (no poison loop), so it is dropped and counted. 413 is the same.
195 if (res.status === 400 || res.status === 413) {
196 ev.droppedTotal += dropThrough(batch[batch.length - 1].seq)
197 return true
198 }
199 return false // 503 registry_full, any other status: retry
200}
201
202// flush sends one batch (one request in flight at a time) and schedules the next: 150 ms
203// later after an answer, after the backoff after a failure.
204async function flush($) {
205 if (ev.inflight || ev.queue.length === 0) return
206 const batch = ev.queue.slice(0, BATCH_MAX)
207 ev.inflight = true
208 let outcome
209 try {
210 outcome = await postWithDeadline($, batch)
211 } finally {
212 ev.inflight = false
213 }
214 if (apply(outcome, batch)) {
215 ev.backoffMs = 0
216 schedule($, FLUSH_MS)
217 return
218 }
219 ev.backoffMs = Math.min(Math.max(BACKOFF_MIN_MS, ev.backoffMs * 2), BACKOFF_MAX_MS)
220 schedule($, ev.backoffMs)
221}
222
223// finalFlush is the session.end flush, awaited inside its hook: one POST of every queued
224// event — the newest FINAL_MAX, older ones counted as dropped — whatever is in flight or
225// backing off (an in-flight batch's events are still queued, so its late arrival is all
226// duplicates). No deadline: the session.end bound (1.5 s) cuts it.
227async function finalFlush($) {
228 if (ev.queue.length === 0) return
229 if (ev.queue.length > FINAL_MAX) {
230 const cut = ev.queue.length - FINAL_MAX
231 ev.queue = ev.queue.slice(cut)
232 ev.droppedTotal += cut
233 }
234 const batch = ev.queue.slice()
235 let outcome
236 try {
237 outcome = { res: await post($, batch) }
238 } catch (err) {
239 outcome = { err }
240 }
241 apply(outcome, batch)
242}
243
244// ---- the heartbeat ----
245
246function beatTick($) {
247 void beat($).catch((err) => log($, 'heartbeat failed: ' + String(err)))
248}
249
250// beat reads the agents and queues one heartbeat — unless the heartbeat stopped or paused
251// while it waited on $.agent.list (a session.end, a new session.start): stopping cancels the
252// timer, not a beat already past its start, so that one checks again on its way out. While
253// a /clear or a resume switches the session id, no beat goes out: it would carry the old id.
254async function beat($) {
255 if (!ev.on || ev.switching) return
256 const gen = ev.beatGen
257 let agents // left out when the list fails: [] would clear every dot until the next beat
258 try {
259 agents = (await $.agent.list()).map((a) => ({ id: a.id, status: a.status }))
260 } catch {}
261 if (gen !== ev.beatGen) return
262 const data = { asks: [...ev.asks.keys()], compacting: ev.compacting, error: ev.lastError }
263 if (agents) data.agents = agents
264 if (ev.turnId) data.turn_id = ev.turnId
265 if (ev.background) data.background = ev.background
266 enqueue($, 'heartbeat', data)
267}
268
269function stopBeat() {
270 if (ev.beat) ev.beat.cancel()
271 ev.beat = null
272 ev.beatGen += 1
273}
274
275// ---- the session ----
276
277// startReporter turns the reporter on for an interactive session whose pdx.json names the
278// daemon's socket (U1-1a's extractor writes `mod_socket`; an older install has none, and
279// the reporter stays off). The stream and its seq outlive a second session.start in the
280// same load.
281async function startReporter($, e) {
282 stopBeat()
283 ev.on = false
284 ev.switching = false
285 const cfg = parse(await $.fs.read($.plugin.root + '/pdx.json').catch(() => ''))
286 const sock = isObject(cfg) && typeof cfg.mod_socket === 'string' ? cfg.mod_socket : ''
287 if (!sock) return
288 if (!ev.stream) ev.stream = newStream()
289 if (!ev.stream) {
290 log($, 'reporter off: this environment has no crypto.getRandomValues for the stream id')
291 return
292 }
293 ev.sock = sock
294 ev.modVersion = String(await $.fs.read($.plugin.root + '/VERSION').catch(() => '')).trim()
295 const v = await $.session.version().catch(() => undefined)
296 ev.ccVersion = isObject(v) && typeof v.version === 'string' ? v.version : ''
297 ev.sid = String(await $.session.id())
298 ev.cwd = String(e.cwd ?? '')
299 ev.turnId = ''
300 ev.asks.clear()
301 ev.compacting = false
302 ev.lastError = false
303 ev.background = null
304 ev.monitors.clear()
305 ev.on = true
306 enqueue($, 'session.start', { cwd: e.cwd, surface: e.surface })
307 ev.beat = $.clock.every(HEARTBEAT_MS, () => beatTick($))
308}
309
310// sessionSwitch follows a /clear or a resume: the process goes on under a new session id,
311// in the same stream, its seq and heartbeat going on; the old conversation's mirror is gone.
312// The cwd and the background tasks stay (same process: the tasks run on; the next
313// classic.Stop refreshes them). It ends the pause session.end began, whatever happens here.
314async function sessionSwitch($, source) {
315 if (!ev.on) return
316 try {
317 const prev = ev.sid
318 ev.sid = String(await $.session.id())
319 ev.turnId = ''
320 ev.asks.clear()
321 ev.compacting = false
322 ev.lastError = false
323 ev.monitors.clear()
324 enqueue($, 'session.switch', { prev_sid: prev, source })
325 } finally {
326 ev.switching = false
327 }
328}
329
330// sessionEnd queues session.end under the ending session's own id and flushes inside the
331// hook. `session.end` also fires on /clear and /resume (reason clear / resume), after which
332// the same process and stream go on: the heartbeat stops only on the other reasons, before
333// the final flush, so nothing of it (not even a beat in flight) lands after session.end. On
334// clear / resume it pauses until session.switch (M-U1-4: about 660 ms later), so no beat
335// goes out under the old id meanwhile, the one in flight included.
336async function sessionEnd($, e) {
337 if (!ev.on) return
338 ev.monitors.clear()
339 if (e.reason === 'clear' || e.reason === 'resume') {
340 ev.switching = true
341 ev.beatGen += 1
342 } else {
343 stopBeat()
344 }
345 enqueue($, 'session.end', { reason: e.reason }, e.sessionId)
346 await finalFlush($)
347}
348
349// ---- what each hook reports ----
350
351function withAgent(data, agentId) {
352 if (agentId) data.agent_id = agentId
353 return data
354}
355
356function turnStarted($, e) {
357 if (!ev.on) return
358 ev.turnId = e.turnId // turn.start has no agentId: it is always the main conversation's
359 ev.lastError = false
360 enqueue($, 'turn.start', { turn_id: e.turnId })
361}
362
363function turnCompleted($, e) {
364 if (!ev.on) return
365 if (!e.agentId) {
366 ev.turnId = ''
367 ev.asks.clear()
368 ev.lastError = e.reason === 'error'
369 }
370 enqueue($, 'turn.complete', { ...withAgent({ turn_id: e.turnId, reason: e.reason }, e.agentId), duration_ms: e.durationMs, aborted: !!e.isAborted })
371}
372
373function toolStarted($, e) {
374 if (ASK_TOOLS.has(e.tool) && e.tool_use_id) ev.asks.set(e.tool_use_id, 'question')
375 enqueue($, 'tool.start', withAgent({ tool: e.tool, tool_use_id: e.tool_use_id }, e.agentId))
376}
377
378function toolEnded($, e, ms, error) {
379 if (e.tool_use_id) ev.asks.delete(e.tool_use_id)
380 enqueue($, 'tool.end', withAgent({ tool_use_id: e.tool_use_id, ms, error }, e.agentId))
381}
382
383// toolUseDrawn reports the approval of a permission ask: the ToolUse row's isRunning is false
384// while the permission dialog is open and turns true about 16 ms after the person approves
385// (M-U1-6, Claude Code 2.1.294). It runs on every redraw of every tool row, so it only looks
386// the id up; the ask leaves the mirror at once, so a later redraw reports nothing. A question
387// is never approved here: it waits until its tool.end.
388function toolUseDrawn($, e) {
389 if (!ev.on) return
390 const p = e.props
391 if (!p || p.isRunning !== true || ev.asks.get(p.tool_use_id) !== 'permission') return
392 ev.asks.delete(p.tool_use_id)
393 enqueue($, 'tool.approved', { tool_use_id: p.tool_use_id })
394}
395
396// monitorStarted remembers the task id of a Monitor call (measured on Claude Code 2.1.295: the hook's result is
397// {ref, result: {taskId, timeoutMs, persistent}, text}, and taskId is the id classic.Stop lists in background_tasks).
398function monitorStarted(e, r) {
399 if (e.tool !== 'Monitor' || !isObject(r) || isErrorResult(r)) return
400 // the nested result first, then a top-level taskId, then the strictly anchored text of the success message
401 let id = isObject(r.result) ? r.result.taskId : undefined
402 if (typeof id !== 'string' || !id) id = r.taskId
403 if (typeof id !== 'string' || !id) {
404 const m = typeof r.text === 'string' ? /^Monitor started \(task ([A-Za-z0-9_-]+)[,)]/.exec(r.text) : null
405 id = m ? m[1] : undefined
406 }
407 if (typeof id !== 'string' || !id) return
408 ev.monitors.delete(id) // a reused id is the newest again (Map keeps insertion order)
409 ev.monitors.set(id, false)
410 while (ev.monitors.size > MONITORS_MAX) ev.monitors.delete(ev.monitors.keys().next().value)
411}
412
413function isErrorResult(r) {
414 return !!r && (r.isError === true || r.deny !== undefined)
415}
416
417function toolChecked($, e, r) {
418 if (!ev.on) return
419 const decision = isObject(r) ? r.decision : undefined
420 if (decision === 'ask' && e.tool_use_id) ev.asks.set(e.tool_use_id, ASK_TOOLS.has(e.tool) ? 'question' : 'permission')
421 const data = { tool: e.tool }
422 if (e.tool_use_id) data.tool_use_id = e.tool_use_id
423 enqueue($, 'tool.check', { ...withAgent(data, e.agentId), decision })
424}
425
426function agentSpawned($, e, r) {
427 if (!ev.on || !isObject(r) || typeof r.agentId !== 'string') return // refused, or answered without starting one
428 const data = { agent_id: r.agentId, tool_use_id: e.tool_use_id, background: !!e.background, subagent_type: e.subagentType }
429 if (isObject(e.workflow) && e.workflow.runId) data.workflow_run_id = e.workflow.runId
430 enqueue($, 'agent.spawn', data)
431}
432
433function measured($, e) {
434 if (!ev.on) return
435 const c = isObject(e.context) ? e.context : {}
436 const context = { window: c.window }
437 if (c.tokens !== undefined) context.tokens = c.tokens
438 if (c.percent !== undefined) context.percent = c.percent
439 const data = {
440 context,
441 rate_limits: (Array.isArray(e.rateLimits) ? e.rateLimits : []).map((l) => {
442 const o = { kind: l.kind, percent_used: l.percentUsed }
443 if (l.resetsAt !== undefined) o.resets_at = l.resetsAt
444 return o
445 }),
446 }
447 if (isObject(e.cost) && Number.isFinite(e.cost.usd)) data.cost_usd = e.cost.usd
448 data.changed = Array.isArray(e.changed) ? [...e.changed] : []
449 enqueue($, 'usage', data)
450}
451
452function stopped($, e) {
453 if (!ev.on) return
454 // a copy: the engine's objects are never touched. A task the Monitor tool started is named "monitor".
455 const listed = Array.isArray(e.background_tasks) ? e.background_tasks : []
456 const tasks = listed.map((t) => ({ id: t.id, type: ev.monitors.has(t.id) ? 'monitor' : t.type, status: t.status }))
457 // a monitor this Stop lists is now known to the engine; one it listed before and no longer lists has ended
458 const present = new Set(listed.map((t) => t.id))
459 for (const [id, seen] of [...ev.monitors]) {
460 if (present.has(id)) ev.monitors.set(id, true)
461 else if (seen) ev.monitors.delete(id)
462 }
463 ev.background = { tasks, crons: Array.isArray(e.session_crons) ? e.session_crons.length : 0 }
464 enqueue($, 'background', ev.background)
465}
466
467function compactStarted($, e) {
468 if (!ev.on) return
469 if (!e.agentId && e.trigger !== 'precompute') ev.compacting = true
470 enqueue($, 'compact.start', withAgent({ trigger: e.trigger }, e.agentId))
471}
472
473function compactEnded($, e, ok) {
474 if (!ev.on) return
475 if (!e.agentId && e.trigger !== 'precompute') ev.compacting = false
476 enqueue($, 'compact.end', { ...withAgent({ trigger: e.trigger }, e.agentId), ok })
477}
478
479// ---- the hooks ----
480
481async function onSessionStart($, e, next) {
482 const r = await next(e)
483 await startReporter($, e)
484 return r
485}
486
487async function onTurnStart($, e, next) {
488 turnStarted($, e)
489 return next(e)
490}
491
492async function onTurnComplete($, e, next) {
493 const r = await next(e)
494 turnCompleted($, e)
495 return r
496}
497
498// onSessionSwitch reports the switch once the hooks beneath (the relay's /clear work) are
499// done with the event — and when one of them failed too: the session id changed whatever
500// they did, and the heartbeat stays paused until the switch is reported.
501async function onSessionSwitch($, e, next) {
502 try {
503 return await next(e) // the engine has moved to the new session id
504 } finally {
505 await sessionSwitch($, e.source)
506 }
507}
508
509// onCompact wraps the relay's compact hook (registered before it, so outside): ok is false
510// exactly when the compaction did not happen (a hook beneath answered {skip}).
511async function onCompact($, e, next) {
512 if (!ev.on) return next(e)
513 compactStarted($, e)
514 let ok = false
515 try {
516 const r = await next(e)
517 ok = !(isObject(r) && r.skip !== undefined)
518 return r
519 } finally {
520 compactEnded($, e, ok)
521 }
522}
523
524// onToolCall reports tool.start before next(e) and tool.end after it, timed on the engine's
525// clock. When an outer hook answers with next still pending (ask.js's remote answer) what
526// runs beneath is abandoned; once next(e) here rejects, tool.end {error: true} and the
527// rejection goes on (the .catch replays it; the outer hook's answer stands). Should it never
528// settle, the main turn.complete still clears the open ask from the mirror.
529async function onToolCall($, e, next) {
530 if (!ev.on) return next(e)
531 const t0 = await $.clock.now()
532 toolStarted($, e)
533 let r
534 try {
535 r = await next(e)
536 } catch (err) {
537 toolEnded($, e, (await $.clock.now()) - t0, true)
538 throw err
539 }
540 toolEnded($, e, (await $.clock.now()) - t0, isErrorResult(r))
541 monitorStarted(e, r)
542 return r
543}
544
545// onToolUseRender only observes: the drawing is whatever beneath answers, unchanged.
546async function onToolUseRender($, e, next) {
547 toolUseDrawn($, e)
548 return next(e)
549}
550
551async function onToolCheck($, e, next) {
552 const r = await next(e)
553 toolChecked($, e, r)
554 return r
555}
556
557async function onAgentSpawn($, e, next) {
558 const r = await next(e)
559 agentSpawned($, e, r)
560 return r
561}
562
563async function onMeasure($, e, next) {
564 const r = await next(e)
565 measured($, e)
566 return r
567}
568
569async function onStop($, e, next) {
570 const r = await next(e)
571 stopped($, e)
572 return r
573}
574
575async function onSessionEnd($, e, next) {
576 const r = await next(e)
577 await sessionEnd($, e)
578 return r
579}
580
581// registerEvents is called once by register.js, before register.js's own hooks, so the
582// matched hooks below are outermost on the events register.js owns. The matchers admit
583// every event the reporter wants: any turn (every turn has an id), any compaction trigger,
584// an interactive session start, a SessionStart that switched the session id, and a tool
585// row's drawing.
586export function registerEvents(on) {
587 on('session.start', { isInteractive: true }, onSessionStart).catch(($, e, next) => next(e))
588 on('turn.start', { turnId: /^/ }, onTurnStart).catch(($, e, next) => next(e))
589 on('turn.complete', { turnId: /^/ }, onTurnComplete).catch(($, e, next) => next(e))
590 on('classic.SessionStart', { source: ['clear', 'resume'] }, onSessionSwitch).catch(($, e, next) => next(e))
591 on('session.compact', { trigger: /^/ }, onCompact).catch(($, e, next) => next(e))
592 on('tool.call', onToolCall).catch(($, e, next) => next(e))
593 on('tool.check', onToolCheck).catch(($, e, next) => next(e))
594 on('agent.spawn', onAgentSpawn).catch(($, e, next) => next(e))
595 on('session.measure', onMeasure).catch(($, e, next) => next(e))
596 on('session.end', onSessionEnd).catch(($, e, next) => next(e))
597 on('classic.Stop', onStop).catch(($, e, next) => next(e))
598 on('ui.render', { component: 'ToolUse' }, onToolUseRender).catch(($, e, next) => next(e))
599}
600hooks/lease.js 511 lines1// Purdex mod — host resource lease, the classifier (host-resource-lease plan Task 2.1, spec D-7, R7, R8).
2//
3// The classifier is pure functions: they read a Bash command string and say whether it is a heavy command (a
4// full test run, a build, a lint of everything) and of which kind, and rewrite a full vitest run to cap its
5// workers. Nothing in them touches the engine ($), so they run, and are tested, as plain code. The hook that
6// uses them, registerLease, is at the end of the file.
7//
8// Fail open everywhere: a command that cannot be parsed, or that merely might be heavy, is not intercepted.
9// Splitting is shell-aware only as far as it needs to be: quotes and backslashes, the separators
10// `&&` `||` `;` `|` `|&` newline and `(` `)`, redirections. A single `&` (a backgrounded command) is not
11// intercepted (R8).
12
13// ---- tokenizing ----
14
15// a redirection word ends in an operator (its target is the next word)
16const BARE = /[<>&]$/
17
18// scan splits a command into segments of words. A word records the text between its quotes and where it
19// sits in the original string. Returns null when the command cannot be parsed (an open quote) or has a
20// backgrounding `&`.
21function scan(cmd) {
22 const segments = []
23 let words = []
24 let cur = null // {text, start, end, redirect}
25 let quote = ''
26 let background = false
27 const heredocs = [] // {delim, strip, quoted} awaiting the next newline
28 const masks = [] // [from, to) of text that is data (a quoted heredoc's body, a comment): substitutions() must not read it
29
30 const endWord = (i) => {
31 if (cur) {
32 cur.end = i
33 words.push(cur)
34 cur = null
35 }
36 }
37 const endSegment = (i) => {
38 endWord(i)
39 if (words.length) segments.push(words)
40 words = []
41 }
42 const push = (ch, i) => {
43 if (!cur) cur = { text: '', start: i, end: i, quoted: false, redirect: false }
44 cur.text += ch
45 }
46
47 for (let i = 0; i < cmd.length; i++) {
48 const c = cmd[i]
49 if (quote === "'") {
50 if (c === "'") quote = ''
51 else push(c, i)
52 continue
53 }
54 if (quote === '"') {
55 if (c === '"') quote = ''
56 else if (c === '\\' && i + 1 < cmd.length) { push(cmd[++i], i) }
57 else push(c, i)
58 continue
59 }
60 if (c === "'" || c === '"') {
61 quote = c
62 if (!cur) cur = { text: '', start: i, end: i, quoted: false, redirect: false }
63 cur.quoted = true
64 continue
65 }
66 if (c === '\\' && i + 1 < cmd.length) { push(cmd[++i], i); continue }
67 if (c === ' ' || c === '\t') { endWord(i); continue }
68 if (c === '\n') {
69 endSegment(i)
70 for (const h of heredocs.splice(0)) {
71 const from = i + 1
72 // the body: lines up to one that is the delimiter; none of it is a command
73 for (;;) {
74 const nl = cmd.indexOf('\n', i + 1)
75 const line = cmd.slice(i + 1, nl < 0 ? cmd.length : nl)
76 i = nl < 0 ? cmd.length : nl
77 if ((h.strip ? line.replace(/^\t+/, '') : line) === h.delim || nl < 0) break
78 }
79 if (h.quoted) masks.push([from, i])
80 }
81 continue
82 }
83 if (c === ';') { endSegment(i); continue }
84 if (c === '#' && !cur) {
85 const from = i
86 while (i + 1 < cmd.length && cmd[i + 1] !== '\n') i++
87 masks.push([from, i + 1])
88 continue
89 }
90 if (c === '>' || c === '<') {
91 // an unquoted redirection ends the word before it (`run>out`), but keeps a file descriptor
92 // (`2>`) and the operator it is building (`>>`, `<<`)
93 if (cur && !(cur.redirect && /[<>]$/.test(cur.text)) && !/^(\d+|&)$/.test(cur.text)) endWord(i)
94 if (c === '<' && cmd[i + 1] === '<' && cmd[i - 1] !== '<' && cmd[i + 2] !== '<') {
95 let j = i + 2
96 const strip = cmd[j] === '-'
97 if (strip) j++
98 while (cmd[j] === ' ' || cmd[j] === '\t') j++
99 let delim = ''
100 let quoted = false
101 // a shell word: quotes group (and keep their spaces), a backslash escapes, the rest ends at a separator
102 for (let q = ''; j < cmd.length; j++) {
103 const d = cmd[j]
104 if (q) { if (d === q) q = ''; else delim += d; continue }
105 if (d === "'" || d === '"') { q = d; quoted = true; continue }
106 if (d === '\\' && j + 1 < cmd.length) { delim += cmd[++j]; quoted = true; continue }
107 if (' \t\n;|&()<>'.includes(d)) break
108 delim += d
109 }
110 heredocs.push({ delim, strip, quoted })
111 }
112 push(c, i)
113 cur.redirect = true
114 continue
115 }
116 if (c === '|') {
117 endSegment(i)
118 if (cmd[i + 1] === '|') i++
119 else if (cmd[i + 1] === '&') i++
120 continue
121 }
122 if (c === '&') {
123 if (cmd[i + 1] === '&') { endSegment(i); i++; continue }
124 const prev = cur ? cur.text[cur.text.length - 1] : ''
125 if (prev === '>' || prev === '<' || cmd[i + 1] === '>') { push(c, i); continue } // 2>&1, &>file
126 background = true
127 endSegment(i)
128 continue
129 }
130 if ((c === '(' || c === ')') && cmd[i - 1] !== '$') { endSegment(i); continue }
131 push(c, i)
132 }
133 if (quote !== '') return null
134 endSegment(cmd.length)
135 if (background) return null
136 // a copy of the command with the data spans blanked, the same length, for the substitution search
137 let masked = cmd
138 for (const [a, b] of masks) masked = masked.slice(0, a) + ' '.repeat(b - a) + masked.slice(b)
139 segments.masked = masked
140 return segments
141}
142
143// ---- reading one segment ----
144
145const PM_VALUE_OPTS = new Set(['-C', '--dir', '--prefix', '--filter', '-F', '--workspace', '-w', '--package', '-p', '--cwd'])
146const RUNNERS = new Set(['npx', 'pnpm', 'pnpx', 'yarn', 'npm', 'bunx', 'bun', 'corepack'])
147const PM_SUBCOMMANDS = new Set(['exec', 'run', 'run-script', 'dlx', 'x'])
148const SHELLS = new Set(['sh', 'bash', 'zsh', 'dash', 'fish', 'ksh'])
149const ASSIGN = /^[A-Za-z_][A-Za-z0-9_]*=/
150
151const base = (p) => p.slice(p.lastIndexOf('/') + 1)
152
153// commandOf strips what comes before the command proper: environment assignments, `env`, `time`,
154// `timeout N`, `nice`, `exec`, and a package runner with its own options. It returns the words that are
155// the command and its arguments (redirections left out), or [] when there is none.
156function commandOf(words) {
157 let w = words.filter((x) => !x.redirect).map((x) => x.text)
158 // a redirection's target (the word after a bare `>`) is not an argument
159 const idx = words.map((x, k) => (x.redirect && BARE.test(x.text) ? k + 1 : -1)).filter((k) => k >= 0)
160 if (idx.length) w = words.filter((x, k) => !x.redirect && !idx.includes(k)).map((x) => x.text)
161 for (let guard = 0; guard < 8 && w.length; guard++) {
162 const head = base(w[0])
163 if (ASSIGN.test(w[0])) { w = w.slice(1); continue }
164 if (head === 'env' || head === 'time' || head === 'exec' || head === 'command') { w = w.slice(1); continue }
165 if (head === 'nice') { w = w.slice(1); if (w[0] === '-n') w = w.slice(2); continue }
166 if (head === 'timeout') {
167 w = w.slice(1)
168 while (w.length && w[0].startsWith('-')) w = w.slice(w[0] === '-s' || w[0] === '-k' ? 2 : 1)
169 if (w.length) w = w.slice(1) // the duration
170 continue
171 }
172 if (RUNNERS.has(head)) {
173 w = w.slice(1)
174 while (w.length) {
175 if (PM_VALUE_OPTS.has(w[0])) { w = w.slice(2); continue }
176 if (w[0].startsWith('-') && !w[0].includes('=') && PM_VALUE_OPTS.has(w[0])) { w = w.slice(2); continue }
177 if (w[0].startsWith('-')) { w = w.slice(1); continue }
178 if (PM_SUBCOMMANDS.has(w[0])) { w = w.slice(1); continue }
179 break
180 }
181 continue
182 }
183 break
184 }
185 return w
186}
187
188// vitest flags that take a value as the next word. Anything unknown starting with `-` is taken as a
189// switch, so the word after it counts as a positional (a path): not intercepted, the safe side.
190const VITEST_VALUE = new Set(['-t', '--testNamePattern', '--project', '--maxWorkers', '--max-workers', '--reporter', '--config', '-c', '--root', '-r', '--dir', '--pool', '--shard', '--outputFile', '--environment', '--mode', '-m', '--coverage.reporter', '--coverage.provider', '--poolOptions.threads.maxThreads', '--poolOptions.forks.maxForks', '--poolOptions.threads.minThreads', '--poolOptions.forks.minForks', '--minWorkers', '--min-workers', '--retry', '--bail', '--testTimeout', '--hookTimeout', '--exclude', '--browser.name', '--cache'])
191const VITEST_LIMIT = /^--(maxWorkers|max-workers|poolOptions\..*(maxThreads|maxForks))(=|$)/
192
193// vitestRun reads the words after `vitest`. kind is 'test-full' for a run over everything, null for any
194// narrowed or watch invocation; capped says a worker limit is already given.
195function vitestRun(args) {
196 let run = false
197 let capped = false
198 for (let i = 0; i < args.length; i++) {
199 const a = args[i]
200 if (a === 'run') { run = true; continue }
201 if (a.startsWith('-')) {
202 const name = a.split('=')[0]
203 if (VITEST_LIMIT.test(a)) capped = true
204 if (name === '-t' || name === '--testNamePattern' || name === '--project' || name === '--changed' || name === '--related' || name === '--shard') return { kind: null, capped }
205 if (!a.includes('=') && VITEST_VALUE.has(a)) i++
206 continue
207 }
208 return { kind: null, capped } // a positional: a path, a filter, `related`, `watch`
209 }
210 return { kind: run ? 'test-full' : null, capped }
211}
212
213const GO_VALUE = new Set(['-run', '-bench', '-count', '-timeout', '-tags', '-p', '-parallel', '-cpu', '-coverprofile', '-o', '-vet', '-covermode', '-coverpkg', '-skip', '-fuzz', '-exec', '-ldflags', '-gcflags', '-mod', '-modfile', '-overlay', '-pkgdir', '-blockprofile', '-cpuprofile', '-memprofile', '-trace', '-outputdir', '-shuffle', '-fuzztime', '-test.run', '-asmflags', '-buildvcs'])
214
215// goTest reads the words after `go test`.
216function goTest(args) {
217 let race = false
218 let run = false
219 const pkgs = []
220 for (let i = 0; i < args.length; i++) {
221 const a = args[i]
222 if (a.startsWith('-')) {
223 const name = a.replace(/^--?/, '-').split('=')[0]
224 if (name === '-race') race = true
225 if (name === '-run' || name === '-test.run' || name === '-skip' || name === '-fuzz') run = true
226 if (!a.includes('=') && GO_VALUE.has(name)) i++
227 continue
228 }
229 pkgs.push(a)
230 }
231 if (run) return null
232 if (pkgs.some((p) => p === '...' || p.endsWith('/...'))) return 'test-full'
233 if (race && pkgs.length === 1) return 'test-pkg'
234 return null
235}
236
237// kindOf is the kind of one command (its words after commandOf), or null.
238function kindOf(w, depth) {
239 if (!w.length) return null
240 const head = base(w[0])
241 const args = w.slice(1)
242 if (SHELLS.has(head)) {
243 const c = args.indexOf('-c')
244 return c >= 0 && args[c + 1] !== undefined && depth < 2 ? classifyKind(args[c + 1], depth + 1) : null
245 }
246 if (head === 'pdx') return args[0] === 'lease' ? 'wrapped' : null
247 if (head === 'vitest') return vitestRun(args).kind
248 if (head === 'go') {
249 if (args[0] === 'test') return goTest(args.slice(1))
250 if (args[0] === 'vet') return args.slice(1).some((p) => p === '...' || p.endsWith('/...')) ? 'lint-full' : null
251 return null
252 }
253 if (head === 'make') return args.length === 1 && args[0] === 'test' ? 'test-full' : null
254 if (head === 'tsc') return args.includes('-b') || args.includes('--build') ? 'build' : null
255 if (head === 'vite' || head === 'electron-vite') return args[0] === 'build' ? 'build' : null
256 if (head === 'eslint') return args.includes('.') ? 'lint-full' : null
257 // package scripts, as `pnpm run build` / `npm run build` / `yarn build` are left after the runner is stripped
258 if (head === 'build' || head === 'electron:build') return 'build'
259 if (head === 'lint') return 'lint-full'
260 return null
261}
262
263const ORDER = ['test-full', 'build', 'test-pkg', 'lint-full']
264
265// substitutions lists the command lines inside `$( )` and backticks, outside single quotes: they run
266// too, so a heavy command there counts.
267function substitutions(cmd) {
268 const out = []
269 let single = false
270 let double = false
271 for (let i = 0; i < cmd.length; i++) {
272 const c = cmd[i]
273 if (single) { if (c === "'") single = false; continue }
274 if (c === '"') { double = !double; continue }
275 if (c === "'" && !double) { single = true; continue }
276 if (c === '\\') { i++; continue }
277 if (c === '`') {
278 let j = i + 1
279 while (j < cmd.length && cmd[j] !== '`') j += cmd[j] === '\\' ? 2 : 1
280 out.push(cmd.slice(i + 1, j))
281 i = j
282 } else if (c === '$' && cmd[i + 1] === '(') {
283 let depth = 1
284 let j = i + 2
285 for (; j < cmd.length && depth > 0; j++) {
286 if (cmd[j] === '(') depth++
287 else if (cmd[j] === ')') depth--
288 }
289 out.push(cmd.slice(i + 2, j - 1))
290 // the text is still scanned, so a nested one is found by the recursion in classifyKind
291 }
292 }
293 return out
294}
295
296// classifyKind is the heaviest kind among a command's segments and the command lines substituted into
297// them (a heavy command inside a substitution is leased but not given the R7 cap: rewriting inside
298// `$( )` is not attempted). A segment already going through `pdx lease` is skipped, not the whole command: what follows it
299// is not covered by its lease. null for none or an unparseable command.
300function classifyKind(cmd, depth) {
301 const segs = scan(cmd)
302 if (segs === null) return null
303 let best = null
304 const consider = (k) => {
305 if (k && k !== 'wrapped' && (best === null || ORDER.indexOf(k) < ORDER.indexOf(best))) best = k
306 }
307 for (const s of segs) consider(kindOf(commandOf(s), depth))
308 if (depth < 3) for (const inner of substitutions(segs.masked)) consider(classifyKind(inner, depth + 1))
309 return best
310}
311
312// classify says whether a Bash command is heavy: {kind, needsMaxWorkers} where kind is test-full, build, test-pkg
313// or lint-full (spec D-7) and needsMaxWorkers says a full vitest run in it has no worker limit (R7); null when it
314// is not heavy, is already wrapped in `pdx lease`, runs in the background, or cannot be parsed.
315export function classify(command) {
316 if (typeof command !== 'string') return null
317 const kind = classifyKind(command, 0)
318 if (kind === null) return null
319 return { kind, needsMaxWorkers: rewriteMaxWorkers(command).changed }
320}
321
322// needsCap lists the full-vitest segments (indexes into the scanned segments) that carry no worker limit.
323function needsCap(segs) {
324 const out = []
325 segs.forEach((s, i) => {
326 const w = commandOf(s)
327 if (w.length && base(w[0]) === 'vitest') {
328 const r = vitestRun(w.slice(1))
329 if (r.kind === 'test-full' && !r.capped) out.push(i)
330 }
331 })
332 return out
333}
334
335// rewriteMaxWorkers (R7) appends ` --maxWorkers=3` to each full vitest run that has no worker limit,
336// right after its last argument, so it lands before a redirection (`2>&1`) as well as before a pipe.
337// A run inside `sh -c '...'` is rewritten in place when the string is plainly quoted. Returns
338// {command, changed}.
339export function rewriteMaxWorkers(command) {
340 return rewriteAt(command, 0)
341}
342
343function rewriteAt(command, depth) {
344 if (typeof command !== 'string') return { command, changed: false }
345 const segs = scan(command)
346 if (segs === null) return { command, changed: false }
347 const todo = new Set(needsCap(segs))
348 const edits = [] // {at, end, text}: replace command.slice(at, end) with text
349 segs.forEach((s, i) => {
350 if (todo.has(i)) {
351 // the last word that is not a redirection or a redirection's target
352 let last = null
353 for (let k = 0; k < s.length; k++) {
354 if (s[k].redirect) { if (BARE.test(s[k].text)) k++; continue }
355 last = s[k]
356 }
357 if (last !== null) edits.push({ at: last.end, end: last.end, text: ' --maxWorkers=3' })
358 return
359 }
360 // sh -c "<command>": the string is another command line
361 const w = commandOf(s)
362 if (depth >= 2 || !w.length || !SHELLS.has(base(w[0]))) return
363 const c = w.indexOf('-c')
364 if (c < 0 || w[c + 1] === undefined) return
365 const arg = s.find((x) => x.text === w[c + 1] && x.quoted)
366 if (!arg) return
367 const src = command.slice(arg.start, arg.end)
368 const q = src[0]
369 if ((q !== '"' && q !== "'") || src[src.length - 1] !== q || src.slice(1, -1) !== arg.text) return
370 const inner = rewriteAt(arg.text, depth + 1)
371 if (inner.changed) edits.push({ at: arg.start + 1, end: arg.end - 1, text: inner.command })
372 })
373 if (edits.length === 0) return { command, changed: false }
374 let out = command
375 for (const e of edits.sort((a, b) => b.at - a.at)) out = out.slice(0, e.at) + e.text + out.slice(e.end)
376 return { command: out, changed: true }
377}
378
379// ---- the hook (host-resource-lease plan Task 2.2) ----
380//
381// registerLease puts a heavy foreground Bash call through the host's lease: it asks `pdx lease acquire`
382// (which waits up to the daemon's deadline), runs the command, and releases by client id when the call
383// ends, whatever happened. Everything fails open: a daemon that is down, an answer that is not JSON, a
384// throw before the command ran — the command runs, unchanged except for the worker cap.
385//
386// `$` is used only in the top-level functions below (M-U1-2); the handler hands it straight on.
387
388const ACQUIRE_TIMEOUT_MS = 600_000 // $.process.run's cap; the daemon's own deadline is at most 590 s
389const RELEASE_TIMEOUT_MS = 5_000
390const PRE_MS = 3_000 // pdx.json and the session id: a read that has not answered by then is a failed one
391const WAIT_NOTE_MS = 1_000 // a wait shorter than this is not worth telling the model about
392const HEX = '0123456789abcdef'
393
394const parse = (s) => { try { return JSON.parse(s) } catch { return null } }
395const isObject = (v) => !!v && typeof v === 'object' && !Array.isArray(v)
396
397function log($, text) {
398 try {
399 const p = $.ui.log('pdx-lease: ' + text, { to: 'debug' })
400 if (p && typeof p.catch === 'function') p.catch(() => {})
401 } catch {}
402}
403
404// bounded runs one engine read and gives its answer, or undefined when it rejects or has not answered
405// within PRE_MS: the hook must reach next(e) whatever those reads do, so the command is never held by them.
406async function bounded($, read) {
407 const cap = Promise.resolve().then(() => $.clock.sleep(PRE_MS)).then(() => undefined, () => undefined)
408 return Promise.race([Promise.resolve().then(read).catch(() => undefined), cap])
409}
410
411// leaseConfig reads pdx.json beside VERSION as ask.js does: the binary to run and the installing
412// daemon's config. Absent — as under `claude plugin test` — it is `pdx` from PATH and its default config.
413async function leaseConfig($) {
414 const cfg = parse(await bounded($, () => $.fs.read($.plugin.root + '/pdx.json')))
415 return {
416 bin: isObject(cfg) && typeof cfg.pdx === 'string' && cfg.pdx ? cfg.pdx : 'pdx',
417 config: isObject(cfg) && typeof cfg.config === 'string' ? cfg.config : '',
418 }
419}
420
421// newClientId is a lower-case UUID v4 from Web Crypto; '' when the environment has none (the hook then
422// releases by the id acquire returned).
423function newClientId() {
424 const c = globalThis.crypto
425 if (!c || typeof c.getRandomValues !== 'function') return ''
426 const b = new Uint8Array(16)
427 c.getRandomValues(b)
428 b[6] = (b[6] & 0x0f) | 0x40
429 b[8] = (b[8] & 0x3f) | 0x80
430 const h = Array.from(b, (x) => HEX[x >> 4] + HEX[x & 15]).join('')
431 return h.slice(0, 8) + '-' + h.slice(8, 12) + '-' + h.slice(12, 16) + '-' + h.slice(16, 20) + '-' + h.slice(20)
432}
433
434// readAcquire is the one JSON line `pdx lease acquire` prints, or null for anything else (a non-zero exit,
435// junk). The answer is read even when the lease was fail-open: that says so itself.
436function readAcquire(r) {
437 if (!r || r.exitCode !== 0) return null
438 const o = parse(String(r.stdout || '').trim())
439 return isObject(o) ? o : null
440}
441
442// leaseNotes is what the model is told, in Chinese, one line each and only what applies.
443function leaseNotes(rewritten, command, answer) {
444 const notes = []
445 if (rewritten) notes.push('Purdex mod 把 --maxWorkers=3 加進這個指令(主機資源規則 R7),實際執行的是:' + command)
446 if (answer && !answer.fail_open) {
447 const ms = Number(answer.waited_ms) || 0
448 if (answer.overrun) notes.push('等滿 ' + Math.round(ms / 60000) + ' 分鐘超量放行,已記錄')
449 else if (ms >= WAIT_NOTE_MS) {
450 const load = Number(answer.host_measured) > 0 ? '(負載 ' + answer.host_measured + '/100)' : ''
451 notes.push('這個指令先等了 ' + Math.round(ms / 1000) + ' 秒主機資源' + load + ',不是卡住,不要重試')
452 }
453 }
454 return notes
455}
456
457// release gives the lease back by client id (or by id when there was none), 5 s, errors swallowed: a
458// grant whose answer was lost is released too, and a client id the daemon never saw answers `none`.
459async function releaseLease($, cfg, clientId, answer) {
460 const target = clientId ? ['--client-id', clientId] : answer && typeof answer.id === 'string' && answer.id ? [answer.id] : null
461 if (!target) return
462 try {
463 await $.process.run([cfg.bin, 'lease', 'release', ...target, ...(cfg.config ? ['--config', cfg.config] : [])], { timeoutMs: RELEASE_TIMEOUT_MS })
464 } catch (err) {
465 log($, 'release failed: ' + (err && err.message ? err.message : String(err)))
466 }
467}
468
469// leaseCall is one Bash call. A call that is not heavy, runs in the background (R8) or is already wrapped
470// in `pdx lease` goes straight on.
471async function leaseCall($, e, next) {
472 if (e.run_in_background === true) return next(e)
473 const c = classify(e.command)
474 if (!c) return next(e)
475 let command = e.command
476 let rewritten = false
477 if (c.needsMaxWorkers) {
478 const w = rewriteMaxWorkers(command)
479 if (w.changed) { command = w.command; rewritten = true }
480 }
481 const cfg = await leaseConfig($)
482 const sid = await bounded($, () => $.session.id()) // this call's: a /clear or a resume is picked up by the next one
483 if (typeof sid !== 'string' || !sid) {
484 // Without a session the lease cannot be asked for: the command runs, with the cap if it was given one.
485 log($, 'no session id, the command runs without a lease')
486 const r = await next(rewritten ? { ...e, command } : e)
487 const notes = leaseNotes(rewritten, command, null)
488 return notes.length ? { ...r, context: [...(r.context ?? []), ...notes] } : r
489 }
490 const clientId = newClientId()
491 let answer = null
492 try {
493 try {
494 answer = readAcquire(await $.process.run([cfg.bin, 'lease', 'acquire', '--kind', c.kind, '--session', sid,
495 '--tool-use', String(e.tool_use_id || ''), ...(clientId ? ['--client-id', clientId] : []),
496 ...(cfg.config ? ['--config', cfg.config] : [])], { timeoutMs: ACQUIRE_TIMEOUT_MS }))
497 } catch (err) {
498 log($, 'acquire failed, the command runs: ' + (err && err.message ? err.message : String(err)))
499 }
500 const r = await next(rewritten ? { ...e, command } : e)
501 const notes = leaseNotes(rewritten, command, answer)
502 return notes.length ? { ...r, context: [...(r.context ?? []), ...notes] } : r
503 } finally {
504 await releaseLease($, cfg, clientId, answer)
505 }
506}
507
508export function registerLease(on) {
509 on('tool.call', { tool: 'Bash' }, async ($, e, next) => leaseCall($, e, next)).catch(($, e, next) => next(e))
510}
511hooks/member.js 51 lines1// Purdex mod — member relay, the pure half (plan v3 P6-6; spec §8.2 steps 3–8, M1, M28).
2//
3// A lead asks a member to relay with a control message: a peer message whose text carries
4// `[pdx-relay:control] op=<uuid>`. register.js (which owns the shared state and every `$` call)
5// - takes it in `session.receive`, always, busy or idle: the model never sees it;
6// - tells the daemon from a timer that the mod has it (`pdx relay seen`, Q2 branch A);
7// - claims the op only when no main-conversation turn runs: at once when idle, else at that
8// turn's turn.complete (codex finding 1), then takes the write path a self relay takes after
9// an approval: write → check → fix ×2 → written → /clear → seed → done.
10// - a claim that fails does nothing: the daemon's own timeout reports it.
11//
12// The mod's static rules never follow `$` across an import (Q3 deviation: the receive hook and
13// the claim cannot live here), so this file holds what needs no `$`: reading the control message,
14// checking a claim's answer, and writing the team facts of the write prompt.
15
16const CONTROL_PREFIX = '[pdx-relay:control] op=' // team.RelayControlPrefix
17const UUID = '[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}'
18const CONTROL_RE = new RegExp(CONTROL_PREFIX.replace(/[.*+?^${}()|[\]\\]/g, '\\$&') + '(' + UUID + ')')
19
20export const CONSUMED = 'pdx relay control message' // the reason logged for a consumed delivery
21
22// controlOp is the op id of a control message, '' for any other delivery. Only a peer's text counts
23// (M1): the same words typed by the user, or in a tool result, are ordinary text.
24export function controlOp(e) {
25 if (!e || !e.origin || e.origin.kind !== 'peer' || typeof e.text !== 'string') return ''
26 const m = CONTROL_RE.exec(e.text)
27 return m ? m[1] : ''
28}
29
30// claimOp is the op of a claim's answer (`pdx relay claim` stdout), undefined when it is not one.
31export function claimOp(body) {
32 const op = body && body.op
33 return op && typeof op.id === 'string' && typeof op.handoff_path === 'string' ? op : undefined
34}
35
36// leadLine is a member's `{{team}}` (write prompt tail §8.2 step 4), starting with a newline since the
37// placeholder ends the whoami line. `clean` makes third-party text safe (register.js cleanGit).
38export function leadLine(lead, clean) {
39 return '\n- 我的 lead:' + clean(String(lead.address)) + '(ref ' + clean(String(lead.ref)) + ',team ' + clean(String(lead.team_id)) + ')'
40}
41
42// rosterLines is a lead's `{{team}}` from `pdx team --json`'s answer (undefined when it was unreadable).
43// Titles and directories are third-party text: cleaned, one line each.
44export function rosterLines(view, clean) {
45 if (!view || typeof view !== 'object') return '\n- 我管理的 members:(讀不到,接手後請用 pdx team 查)'
46 const members = Array.isArray(view.members) ? view.members : []
47 if (members.length === 0) return '\n- 我管理的 members:無'
48 const one = (v) => clean(String(v ?? '')).replace(/\s+/g, ' ').trim()
49 return '\n- 我管理的 members:' + members.map((m) => '\n - ' + one(m.address) + '(ref ' + one(m.ref) + ',title ' + (one(m.title) || '無') + ',cwd ' + one(m.cwd) + ')').join('')
50}
51hooks/prompts.js 5 lines1// GENERATED from internal/team/relay_prompts.go by: go test ./cmd/pdx/plugin/ -run TestPromptsJS -update — do not edit.
2export const DEFAULT_BODIES = {"write":"這個 session 的 context 已達接力門檻,使用者已核准接力(之後會 /clear)。\n請先停下手邊工作,用你完整的工具撰寫接力檔:{{path}}\n\n要求:\n- 接力檔必須自成一體:讀它的是一個完全沒有這段對話記憶的新對話。\n\n機器提供的 git 狀態(照抄進 §3):\n{{git}}","fix":"接力檔 {{path}} 不完整。","seed":"你是接手的新對話:前一段對話 context 已滿並已清空。\n請先讀接力檔 {{path}},然後:\n1. 用三行複述:目標、下一步第一個動作、目前有哪些檔案異動。\n2. 跑 `git status` 確認與接力檔一致,不一致就指出來。\n3. 接著從「下一步」繼續原本的工作。\n回覆的第一行請寫「↪ 接手自 {{old_ref}}」。"}
3export const FIXED = {"write":{"head":"[pdx-relay op={{op}} n={{nonce}}] ","tail":"\n- 只用 Write 工具一次寫入整個接力檔;這一輪不要執行其他工具(接力期間其他工具會被拒絕)。\n- 寫完後只回一行「HANDOFF-WRITTEN」,不要繼續原本的工作。\n\n格式(每一段都要有,沒有內容就寫「無」):\n# HANDOFF\n## 1. 目標與完成定義(使用者要的是什麼、怎樣算完成、範圍外)\n## 2. 進度(已完成且驗證 / 進行中停在哪 / 下一步第一個動作具體到指令)\n## 3. 檔案異動(git status 與 diff --stat 的結果,加上每個檔案的用途)\n## 4. 決策紀錄(選了什麼、為什麼、否決了什麼)\n## 5. 死路(試過失敗、不要再試的)\n## 6. 環境與指令(測試 / 執行方式)\n## 7. 未決問題與需要使用者決定的事\n## 8. 協作關係(下面的 pdx 身分;我的 lead 與我管理的 members,沒有就寫無)\n\n機器提供的事實(請照抄進對應段落):\n- 舊 session id:{{old_session}}\n- 舊 ref:{{old_ref}}\n- 接力時 context:{{context}}\n- pdx 身分:{{whoami}}{{team}}"},"fix":{"head":"[pdx-relay op={{op}} n={{nonce}}] ","tail":"缺少段落:{{missing}}。請用 Write 重寫整個接力檔(不要用 Edit),補齊後只回「HANDOFF-WRITTEN」。"},"seed":{"head":"↪ 接手自 {{old_ref}}\n[pdx-relay seed op={{op}} n={{nonce}}] ","tail":"{{tasks}}"}}
4export const VARIABLES = ["path","old_ref","old_session","context","whoami","git"]
5