Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions trios/agent-server/apps/server/src/agent/compaction.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -150,6 +151,7 @@ async function compactMessages(
config,
)

markActivity(`compaction: find split point (${messages.length} messages)`)
const { splitIndex, turnStartIndex, isSplitTurn } = findSafeSplitPoint(
messages,
config.keepRecentTokens,
Expand Down Expand Up @@ -353,13 +355,15 @@ 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) {
return { messages: current, experimental_context: state }
}

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`,
Expand All @@ -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 }
Expand Down
4 changes: 4 additions & 0 deletions trios/agent-server/apps/server/src/agent/compaction/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import type {
UserContent,
} from 'ai'
import { logger } from '../../lib/logger'
import { markActivity } from '../../lib/stall-watch'
import {
estimateToolResultOutput,
stripToolResultOutput,
Expand Down Expand Up @@ -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

Expand Down
3 changes: 3 additions & 0 deletions trios/agent-server/apps/server/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -32,6 +33,8 @@ if (!configResult.ok) {
process.exit(EXIT_CODES.GENERAL_ERROR)
}

startStallWatch()

const app = new Application(configResult.value)

try {
Expand Down
4 changes: 4 additions & 0 deletions trios/agent-server/apps/server/src/lib/logger.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -184,6 +185,9 @@ export class Logger implements LoggerInterface {
message: string,
meta?: Record<string, unknown>,
): 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)

Expand Down
76 changes: 76 additions & 0 deletions trios/agent-server/apps/server/src/lib/stall-watch-worker.ts
Original file line number Diff line number Diff line change
@@ -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)
}
96 changes: 96 additions & 0 deletions trios/agent-server/apps/server/src/lib/stall-watch.ts
Original file line number Diff line number Diff line change
@@ -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?.()
}
11 changes: 11 additions & 0 deletions trios/agent-server/apps/server/src/tools/framework.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -128,14 +129,24 @@ 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) {
const message = err instanceof Error ? err.message : String(err)
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<string, unknown>).page
if (typeof pageId === 'number') {
Expand Down
21 changes: 20 additions & 1 deletion trios/agent-server/docker-entrypoint.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down Expand Up @@ -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))
Expand Down
Loading