From 97bea09d2fbe0b900216732a5107899d0f296cb8 Mon Sep 17 00:00:00 2001 From: Dmitriy Vasilev Date: Mon, 21 Sep 2026 22:21:19 +0700 Subject: [PATCH] fix(queen): name where the server spins, and end a stalled loop in 3 min, not 12 Measured 2026-09-21 14:44: the agent server spun one core at 100% for twelve minutes (cpu 973 s -> 1512 s), answered nothing, and the swarm stood until the liveness loop's twelfth missed /health probe ended it. Nothing recorded where it spun: pino writes asynchronously, so the lines logged just before a spin sit in a buffer the blocked thread never flushes. lib/stall-watch.ts: every log line, every tool call and each compaction step is written into a SharedArrayBuffer ring (no I/O). A Worker thread watches the main thread's heartbeat; when it stops for 15 s the worker writes the last twelve activities straight to fd 2 ([stall] lines), and says how long the stall lasted when the loop turns again. docker-entrypoint.sh: the main thread touches /tmp/trios-loop-heartbeat from a timer, which only fires when the event loop turns. A file older than LIVENESS_STALL (180 s) ends the server at once. Busy is still not dead: a loop slowed by twenty bees keeps touching the file. The twelve-probe /health rule stays as the backstop, and is alone in charge when the file is missing. Verified locally with bun 1.3.11 under the extracted run_supervised: a 25 s sync spin is reported with its last activities; a spinning server is ended on the stale heartbeat (exit 143) after one missed probe; a server kept 80% busy for 30 s is not touched. Compaction tests: 125 pass. Co-Authored-By: Claude Opus 5 --- .../apps/server/src/agent/compaction.ts | 6 ++ .../apps/server/src/agent/compaction/utils.ts | 4 + trios/agent-server/apps/server/src/index.ts | 3 + .../apps/server/src/lib/logger.ts | 4 + .../apps/server/src/lib/stall-watch-worker.ts | 76 +++++++++++++++ .../apps/server/src/lib/stall-watch.ts | 96 +++++++++++++++++++ .../apps/server/src/tools/framework.ts | 11 +++ trios/agent-server/docker-entrypoint.sh | 21 +++- 8 files changed, 220 insertions(+), 1 deletion(-) create mode 100644 trios/agent-server/apps/server/src/lib/stall-watch-worker.ts create mode 100644 trios/agent-server/apps/server/src/lib/stall-watch.ts diff --git a/trios/agent-server/apps/server/src/agent/compaction.ts b/trios/agent-server/apps/server/src/agent/compaction.ts index cd6f154cf0..953d1a209a 100644 --- a/trios/agent-server/apps/server/src/agent/compaction.ts +++ b/trios/agent-server/apps/server/src/agent/compaction.ts @@ -6,6 +6,7 @@ import { streamText, } from 'ai' import { logger } from '../lib/logger' +import { markActivity } from '../lib/stall-watch' import { stripBinaryContent } from './compaction/content' import { buildSummarizationPrompt, @@ -150,6 +151,7 @@ async function compactMessages( config, ) + markActivity(`compaction: find split point (${messages.length} messages)`) const { splitIndex, turnStartIndex, isSplitTurn } = findSafeSplitPoint( messages, config.keepRecentTokens, @@ -353,6 +355,7 @@ export function createCompactionPrepareStep( return { messages, experimental_context: state } } + markActivity(`compaction: strip binary (${messages.length} messages)`) let current = stripBinaryContent(messages) currentTokens = estimateTokensForThreshold(current, config) if (currentTokens <= config.triggerThreshold) { @@ -360,6 +363,7 @@ export function createCompactionPrepareStep( } const keepRecent = AGENT_LIMITS.COMPACTION_PRUNE_KEEP_RECENT_MESSAGES + markActivity(`compaction: prune tool calls (${current.length} messages)`) const pruned = pruneMessages({ messages: current, toolCalls: `before-last-${keepRecent}-messages`, @@ -378,10 +382,12 @@ export function createCompactionPrepareStep( } } + markActivity(`compaction: reduce tool outputs (${current.length} messages)`) const reduced = reduceToolOutputs(current, { maxChars: config.toolOutputMaxChars, keepRecentCount: 2, }) + markActivity(`compaction: estimate tokens (${reduced.length} messages)`) currentTokens = estimateTokensForThreshold(reduced, config) if (currentTokens <= config.triggerThreshold) { return { messages: reduced, experimental_context: state } diff --git a/trios/agent-server/apps/server/src/agent/compaction/utils.ts b/trios/agent-server/apps/server/src/agent/compaction/utils.ts index b67ae9f160..9a4f977d08 100644 --- a/trios/agent-server/apps/server/src/agent/compaction/utils.ts +++ b/trios/agent-server/apps/server/src/agent/compaction/utils.ts @@ -6,6 +6,7 @@ import type { UserContent, } from 'ai' import { logger } from '../../lib/logger' +import { markActivity } from '../../lib/stall-watch' import { estimateToolResultOutput, stripToolResultOutput, @@ -459,6 +460,9 @@ export function slidingWindow( messages: ModelMessage[], maxTokens: number, ): ModelMessage[] { + markActivity( + `compaction: sliding window (${messages.length} messages, ${maxTokens} tokens)`, + ) let totalTokens = estimateTokens(messages) let startIndex = 0 diff --git a/trios/agent-server/apps/server/src/index.ts b/trios/agent-server/apps/server/src/index.ts index 7af0859106..30a7d118ab 100755 --- a/trios/agent-server/apps/server/src/index.ts +++ b/trios/agent-server/apps/server/src/index.ts @@ -22,6 +22,7 @@ import { CommanderError } from 'commander' import { loadServerConfig } from './config' import { isPortInUseError } from './lib/port-binding' import { Sentry } from './lib/sentry' +import { startStallWatch } from './lib/stall-watch' import { Application } from './main' const configResult = loadServerConfig() @@ -32,6 +33,8 @@ if (!configResult.ok) { process.exit(EXIT_CODES.GENERAL_ERROR) } +startStallWatch() + const app = new Application(configResult.value) try { diff --git a/trios/agent-server/apps/server/src/lib/logger.ts b/trios/agent-server/apps/server/src/lib/logger.ts index 92c7cefe39..a18abe9596 100644 --- a/trios/agent-server/apps/server/src/lib/logger.ts +++ b/trios/agent-server/apps/server/src/lib/logger.ts @@ -15,6 +15,7 @@ import path from 'node:path' import { CONTENT_LIMITS } from '@browseros/shared/constants/limits' import type { LoggerInterface, LogLevel } from '@browseros/shared/types/logger' import pino from 'pino' +import { markActivity } from './stall-watch' const isDev = process.env.NODE_ENV === 'development' const LOG_FILE_NAME = 'browseros-server.log' @@ -184,6 +185,9 @@ export class Logger implements LoggerInterface { message: string, meta?: Record, ): void { + // Before any I/O: the console write below is asynchronous, and a main + // thread that stalls right after it never flushes the line (stall-watch.ts). + markActivity(`${level} ${message}`) const logFn = this.consoleLogger[level].bind(this.consoleLogger) const fileLogFn = this.fileLogger?.[level].bind(this.fileLogger) diff --git a/trios/agent-server/apps/server/src/lib/stall-watch-worker.ts b/trios/agent-server/apps/server/src/lib/stall-watch-worker.ts new file mode 100644 index 0000000000..579a3fefa8 --- /dev/null +++ b/trios/agent-server/apps/server/src/lib/stall-watch-worker.ts @@ -0,0 +1,76 @@ +/** + * @license + * Copyright 2025 BrowserOS + * SPDX-License-Identifier: AGPL-3.0-or-later + * + * The watching half of lib/stall-watch.ts. Runs on its own thread, so it keeps + * running while the main thread is stuck, and writes to fd 2 directly: a + * console or pino call from here could route through the thread that is stuck. + */ + +import fs from 'node:fs' +import { HEADER_INTS, LABEL_BYTES, RING, SLOT_INTS } from './stall-watch' + +declare const self: Worker + +const REPORT_AFTER_S = Number(process.env.TRIOS_STALL_REPORT_SECONDS ?? 15) +const REPEAT_EVERY_S = 30 +const CHECK_EVERY_MS = 2000 + +const decoder = new TextDecoder() +const say = (line: string) => { + try { + fs.writeSync(2, `[stall] ${line}\n`) + } catch { + // Nowhere left to say it. + } +} + +self.onmessage = (event: MessageEvent<{ shared: SharedArrayBuffer }>) => { + const ints = new Int32Array(event.data.shared) + const bytes = new Uint8Array(event.data.shared) + + const recent = (now: number): string[] => { + const next = Atomics.load(ints, 1) + const lines: string[] = [] + for (let k = 1; k <= RING; k++) { + const slot = (((next - k) % RING) + RING) % RING + const base = HEADER_INTS + slot * SLOT_INTS + const at = Atomics.load(ints, base) + if (!at) continue + const len = Math.min(Atomics.load(ints, base + 1), LABEL_BYTES) + const label = decoder.decode( + bytes.slice((base + 2) * 4, (base + 2) * 4 + len), + ) + lines.push(` ${now - at}s ago: ${label}`) + } + return lines + } + + let stalledSince = 0 + let lastReport = 0 + let longest = 0 + + setInterval(() => { + const now = Math.floor(Date.now() / 1000) + const beat = Atomics.load(ints, 0) + const silent = now - beat + if (silent >= REPORT_AFTER_S) { + if (!stalledSince) stalledSince = beat + if (now - lastReport >= REPEAT_EVERY_S) { + lastReport = now + say( + `main thread has not turned its event loop for ${silent}s; the last things it did, newest first:\n${recent(now).join('\n')}`, + ) + } + } else if (stalledSince) { + const length = now - stalledSince + longest = Math.max(longest, length) + say( + `main thread turned again after ~${length}s (longest stall since boot: ${longest}s)`, + ) + stalledSince = 0 + lastReport = 0 + } + }, CHECK_EVERY_MS) +} diff --git a/trios/agent-server/apps/server/src/lib/stall-watch.ts b/trios/agent-server/apps/server/src/lib/stall-watch.ts new file mode 100644 index 0000000000..7f713977bd --- /dev/null +++ b/trios/agent-server/apps/server/src/lib/stall-watch.ts @@ -0,0 +1,96 @@ +/** + * @license + * Copyright 2025 BrowserOS + * SPDX-License-Identifier: AGPL-3.0-or-later + * + * STALL WATCH: say where the main thread stopped, and let the entrypoint tell a + * hung server from a busy one. + * + * Measured 2026-09-21 14:44: the server spun one core at 100% for twelve + * minutes (cpu 973 s -> 1512 s), answered nothing, and was ended by the + * entrypoint's liveness loop with no record of WHERE it spun. The log could not + * say: pino writes asynchronously, so the lines logged just before the spin sat + * in a buffer the blocked thread never flushed. + * + * So two things live outside the main thread's event loop: + * + * 1. Breadcrumbs. Every log line (and a few hand-placed marks in the hot paths) + * is written into a SharedArrayBuffer ring as it happens - no I/O, no + * event-loop turn needed. A Worker thread reads the ring, and when the main + * thread's heartbeat stops it writes the last activities straight to fd 2, + * which it can do while the main thread is stuck. + * + * 2. A heartbeat file. The main thread touches it every few seconds from a + * timer, which only fires when the event loop turns. A server busy with + * twenty bees still turns it, slowly; a spinning one does not turn it at + * all. docker-entrypoint.sh ends the server when the file goes stale for + * LIVENESS_STALL seconds - minutes sooner than the twelve missed /health + * probes it needed before, which could not tell busy from dead. + */ + +import fs from 'node:fs' + +export const HEARTBEAT_FILE = + process.env.TRIOS_LOOP_HEARTBEAT_FILE ?? '/tmp/trios-loop-heartbeat' + +// Layout of the shared buffer, in Int32 slots: +// [0] main-thread heartbeat, epoch seconds +// [1] index of the next ring slot to write +// then RING slots, each: [time, byteLength, ...LABEL_BYTES/4 ints] +export const RING = 12 +export const LABEL_BYTES = 240 +export const SLOT_INTS = 2 + LABEL_BYTES / 4 +export const HEADER_INTS = 2 +const TOTAL_INTS = HEADER_INTS + RING * SLOT_INTS + +let ints: Int32Array | null = null +let bytes: Uint8Array | null = null +const encoder = new TextEncoder() + +const nowSeconds = () => Math.floor(Date.now() / 1000) + +/** Record what the main thread is doing. Cheap: a few stores, no I/O. */ +export function markActivity(label: string): void { + if (!ints || !bytes) return + const slot = Atomics.add(ints, 1, 1) % RING + const base = HEADER_INTS + slot * SLOT_INTS + const view = bytes.subarray((base + 2) * 4, (base + 2) * 4 + LABEL_BYTES) + const { written } = encoder.encodeInto(label, view) + Atomics.store(ints, base + 1, written ?? 0) + Atomics.store(ints, base, nowSeconds()) +} + +/** Start the heartbeat and the watching thread. Safe to call once. */ +export function startStallWatch(): void { + if (ints || process.env.TRIOS_STALL_WATCH === '0') return + const shared = new SharedArrayBuffer(TOTAL_INTS * 4) + ints = new Int32Array(shared) + bytes = new Uint8Array(shared) + Atomics.store(ints, 0, nowSeconds()) + + try { + fs.writeFileSync(HEARTBEAT_FILE, `${process.pid}\n`) + } catch { + // No heartbeat file: the entrypoint falls back to /health probes alone. + } + let lastTouch = 0 + const beat = setInterval(() => { + const now = nowSeconds() + if (ints) Atomics.store(ints, 0, now) + if (now - lastTouch >= 5) { + lastTouch = now + try { + fs.utimesSync(HEARTBEAT_FILE, now, now) + } catch { + // Missing file only weakens the entrypoint's check; never fatal. + } + } + }, 1000) + beat.unref?.() + + const worker = new Worker(new URL('./stall-watch-worker.ts', import.meta.url)) + worker.postMessage({ shared, pid: process.pid }) + // Bun's Worker has unref(); the DOM typing this project compiles against + // does not declare it. Unref'd, the watcher never keeps the process alive. + ;(worker as Worker & { unref?: () => void }).unref?.() +} diff --git a/trios/agent-server/apps/server/src/tools/framework.ts b/trios/agent-server/apps/server/src/tools/framework.ts index b5f4219186..bed9455a6f 100644 --- a/trios/agent-server/apps/server/src/tools/framework.ts +++ b/trios/agent-server/apps/server/src/tools/framework.ts @@ -4,6 +4,7 @@ import type { ToolApprovalCategoryId } from '@browseros/shared/constants/tool-ap import type { AclRule } from '@browseros/shared/types/acl' import type { z } from 'zod' import type { Browser } from '../browser/browser' +import { markActivity } from '../lib/stall-watch' import { ToolResponse, type ToolResult } from './response' export interface ToolDefinition { @@ -128,6 +129,14 @@ export async function executeTool( } } + // Breadcrumbs for lib/stall-watch.ts: a tool that stalls the main thread + // (a regex over a huge output, a sync walk of a big tree) is named, with the + // start of its arguments, in the report the watcher writes. + let argText = '' + try { + argText = JSON.stringify(args).slice(0, 120) + } catch {} + markActivity(`tool ${tool.name} start ${argText}`) try { await tool.handler(args, ctx, response) } catch (err) { @@ -135,7 +144,9 @@ export async function executeTool( response.error(`Internal error in ${tool.name}: ${message}`) } + markActivity(`tool ${tool.name} build response`) const result = await response.build(ctx.browser) + markActivity(`tool ${tool.name} done`) const pageId = (args as Record).page if (typeof pageId === 'number') { diff --git a/trios/agent-server/docker-entrypoint.sh b/trios/agent-server/docker-entrypoint.sh index e28561af43..4e87054c6c 100755 --- a/trios/agent-server/docker-entrypoint.sh +++ b/trios/agent-server/docker-entrypoint.sh @@ -51,6 +51,18 @@ set -e # twelve misses in a row are needed: roughly twelve minutes of silence. A hang # of the kind this exists for - nine hours of `Application failed to respond` # - is caught all the same; a busy stretch is not. +# +# A STALLED LOOP IS DEAD, and sooner. Twelve minutes of silence was the price of +# telling busy from dead with /health alone. The server now tells them apart +# itself (apps/server/src/lib/stall-watch.ts): its main thread touches +# TRIOS_LOOP_HEARTBEAT_FILE from a timer, which fires only when the event loop +# turns. Twenty busy bees slow the loop; they do not stop it. Measured +# 2026-09-21 14:44 a server spun one core for twelve minutes with the loop +# stopped - so a file older than LIVENESS_STALL (180 s) ends the server without +# waiting for the twelfth missed probe. By then the stall watcher's own thread +# has written the last things the main thread did to the log ([stall] lines). +# No file (an older image, or a server that could not write it) leaves the +# /health rule alone in charge. # One line per reading of the hung server and every process under it, from # /proc alone (python3 is in the image; ps is not). liveness_forensics() { @@ -108,10 +120,17 @@ run_supervised() { interval="${LIVENESS_INTERVAL:-30}" fails_allowed="${LIVENESS_FAILS:-12}" probe_timeout="${LIVENESS_TIMEOUT:-30}" + stall_allowed="${LIVENESS_STALL:-180}" + heartbeat="${TRIOS_LOOP_HEARTBEAT_FILE:-/tmp/trios-loop-heartbeat}" sleep "${LIVENESS_GRACE:-240}" fails=0 while kill -0 "$server" 2>/dev/null; do - if python3 -c "import urllib.request; urllib.request.urlopen('http://127.0.0.1:$port/health', timeout=$probe_timeout)" >/dev/null 2>&1; then + stalled=$(python3 -c "import os,sys,time; print(int(time.time()-os.stat(sys.argv[1]).st_mtime))" "$heartbeat" 2>/dev/null || true) + if [ -n "$stalled" ] && [ "$stalled" -ge "$stall_allowed" ]; then + echo "[liveness] the server's event loop has not turned for ${stalled}s (limit ${stall_allowed}s)" + liveness_forensics "$server" + fails="$fails_allowed" + elif python3 -c "import urllib.request; urllib.request.urlopen('http://127.0.0.1:$port/health', timeout=$probe_timeout)" >/dev/null 2>&1; then fails=0 else fails=$((fails + 1))