From 7100e22890026c37b92a1d511393c2c8c0a6754e Mon Sep 17 00:00:00 2001 From: michelhelsdingen Date: Thu, 3 Sep 2026 09:39:43 +0200 Subject: [PATCH] fix(collab): de feed-loop per team stopte nooit en stapelde zich op als wees collab-launch.sh startte per team een inline `while true` die messages.jsonl naar feed.txt kopieert, zonder enige stopconditie. disbandTeam schreef wel de .finished-marker maar liet die loop met rust, en collab-cleanup.sh haalt alleen mappen weg. Op 01-09-2026 draaiden er zestien, de oudste elf dagen, voor teams die de service niet meer kende; ze zijn toen met de hand gekild. De loop staat nu in scripts/collab-poller.sh en houdt zelf op zodra zijn team voorbij is: bij de .finished-marker (na een laatste flush), bij een verdwenen runtime-map, wanneer de service 404 geeft of het team disbanded meldt, en na COLLAB_POLLER_MAX_API_FAILURES mislukte checks op rij als de service weg is. Hij schrijft zijn eigen poller.pid en ruimt die bij vertrek op. disbandTeam stuurt hem daarnaast meteen een SIGTERM, zodat een disband nooit een loop achterlaat, ook niet tot de volgende tick. Co-Authored-By: Claude Fable 5.1 --- docs/architecture.md | 1 + scripts/collab-launch.sh | 18 +- scripts/collab-poller.sh | 79 +++++++++ services/ensemble-service.ts | 13 +- tests/collab-poller.test.ts | 307 +++++++++++++++++++++++++++++++++++ 5 files changed, 401 insertions(+), 17 deletions(-) create mode 100755 scripts/collab-poller.sh create mode 100644 tests/collab-poller.test.ts diff --git a/docs/architecture.md b/docs/architecture.md index f7f60f7..cd9d3c0 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -61,6 +61,7 @@ ensemble/ ├── scripts/ │ ├── collab-launch.sh # All-in-one team launcher │ ├── collab-poll.sh # Single-shot message poller +│ ├── collab-poller.sh # Per-team feed loop, stops with its team │ ├── collab-livefeed.sh # Continuous live feed │ ├── collab-status.sh # Multi-team dashboard │ ├── collab-replay.sh # Session replay diff --git a/scripts/collab-launch.sh b/scripts/collab-launch.sh index 845306a..71deea2 100755 --- a/scripts/collab-launch.sh +++ b/scripts/collab-launch.sh @@ -118,7 +118,6 @@ RUNTIME_DIR="$(collab_runtime_dir "$TEAM_ID")" MESSAGES_FILE="$(collab_messages_file "$TEAM_ID")" BRIDGE_PID_FILE="$(collab_bridge_pid "$TEAM_ID")" BRIDGE_LOG_FILE="$(collab_bridge_log "$TEAM_ID")" -POLLER_PID_FILE="$(collab_poller_pid "$TEAM_ID")" FEED_FILE="$(collab_feed_file "$TEAM_ID")" TEAM_ID_FILE="$(collab_team_id_file "$TEAM_ID")" @@ -231,21 +230,8 @@ else MONITOR_MODE="session" fi -# ─── 5. Background poller ─── -nohup bash -c ' -TID="'"$TEAM_ID"'" -MESSAGES_FILE="'"$MESSAGES_FILE"'" -FEED_FILE="'"$FEED_FILE"'" -S=0 -while true; do - M=$(wc -l < "$MESSAGES_FILE" 2>/dev/null | tr -d " "); [ -z "$M" ] && M=0 - if [ "$M" -gt "$S" ]; then - tail -n +"$((S+1))" "$MESSAGES_FILE" >> "$FEED_FILE" 2>/dev/null - S=$M - fi - sleep 5 -done' > /dev/null 2>&1 & -printf '%s\n' "$!" > "$POLLER_PID_FILE" +# ─── 5. Background poller (writes its own PID file, stops when the team is over) ─── +nohup "$SCRIPT_DIR/collab-poller.sh" "$TEAM_ID" "$API" > /dev/null 2>&1 & # ─── 6. Wait for agents ─── echo -ne " ${SPIN} Agents spawning..." diff --git a/scripts/collab-poller.sh b/scripts/collab-poller.sh new file mode 100755 index 0000000..daffff7 --- /dev/null +++ b/scripts/collab-poller.sh @@ -0,0 +1,79 @@ +#!/usr/bin/env bash +# collab-poller.sh — Copies new lines from messages.jsonl into feed.txt for one team. +# Usage: collab-poller.sh [api-url] +# +# Started in the background by collab-launch.sh. Stops on its own when the team +# is finished or gone, so it can never outlive its team: +# * the .finished marker appears (written by disbandTeam), +# * the runtime directory is removed (collab-cleanup.sh), +# * the service answers 404 or reports the team as disbanded, +# * the service stays unreachable for COLLAB_POLLER_MAX_API_FAILURES checks. +# +# Env: COLLAB_POLL_SECS (default 5), COLLAB_POLLER_CHECK_EVERY (default 12, so +# the service is asked once a minute), COLLAB_POLLER_MAX_API_FAILURES (default 10). +set -uo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +# shellcheck source=./collab-paths.sh +source "$SCRIPT_DIR/collab-paths.sh" + +TEAM_ID="${1:?Usage: collab-poller.sh [api-url]}" +API="${2:-${ENSEMBLE_URL:-http://localhost:23000}}" +RUNTIME_DIR="$(collab_runtime_dir "$TEAM_ID")" +MESSAGES_FILE="$(collab_messages_file "$TEAM_ID")" +FEED_FILE="$(collab_feed_file "$TEAM_ID")" +PID_FILE="$(collab_poller_pid "$TEAM_ID")" +FINISHED_FILE="$(collab_finished_marker "$TEAM_ID")" + +POLL_SECS="${COLLAB_POLL_SECS:-5}" +CHECK_EVERY="${COLLAB_POLLER_CHECK_EVERY:-12}" +MAX_API_FAILURES="${COLLAB_POLLER_MAX_API_FAILURES:-10}" + +[ -d "$RUNTIME_DIR" ] || exit 0 +printf '%s\n' "$$" > "$PID_FILE" +trap 'rm -f "$PID_FILE"' EXIT INT TERM + +flush() { + local m + m=$(wc -l < "$MESSAGES_FILE" 2>/dev/null | tr -d ' '); [ -z "$m" ] && m=0 + if [ "$m" -gt "$SEEN" ]; then + tail -n +"$((SEEN + 1))" "$MESSAGES_FILE" >> "$FEED_FILE" 2>/dev/null + SEEN=$m + fi +} + +# Prints "gone" when the service says the team is over, "down" when the service +# cannot be reached, and "alive" otherwise. +team_state() { + local code body + body=$(curl -s -o - -w '\n%{http_code}' --max-time 5 "$API/api/ensemble/teams/$TEAM_ID" 2>/dev/null) || { echo down; return; } + code="${body##*$'\n'}" + case "$code" in + 404) echo gone ;; + 200) if printf '%s' "$body" | grep -q '"status":[[:space:]]*"disbanded"'; then echo gone; else echo alive; fi ;; + *) echo down ;; + esac +} + +SEEN=0 +TICK=0 +API_FAILURES=0 +while true; do + flush + [ -d "$RUNTIME_DIR" ] || exit 0 + [ -f "$FINISHED_FILE" ] && exit 0 + + TICK=$((TICK + 1)) + if [ "$TICK" -ge "$CHECK_EVERY" ]; then + TICK=0 + case "$(team_state)" in + gone) exit 0 ;; + down) + API_FAILURES=$((API_FAILURES + 1)) + [ "$API_FAILURES" -ge "$MAX_API_FAILURES" ] && exit 0 + ;; + alive) API_FAILURES=0 ;; + esac + fi + sleep "$POLL_SECS" +done diff --git a/services/ensemble-service.ts b/services/ensemble-service.ts index 7892a79..13d854b 100644 --- a/services/ensemble-service.ts +++ b/services/ensemble-service.ts @@ -23,7 +23,7 @@ import { exportObservation, checkMemoryEndpoint } from '../lib/memory-export' import { AgentWatchdog } from '../lib/agent-watchdog' import { collabPromptFile, collabDeliveryFile, collabSummaryFile, collabMessagesFile, - collabRuntimeDir, collabFinishedMarker, collabBridgePosted, + collabRuntimeDir, collabFinishedMarker, collabBridgePosted, collabPollerPid, collabBridgeResult, ensureCollabDirs, } from '../lib/collab-paths' import fs from 'fs' @@ -1141,6 +1141,17 @@ export async function disbandTeam(teamId: string): Promise 0) { + try { process.kill(pollerPid, 'SIGTERM') } catch { /* already gone */ } + } + fs.unlinkSync(pollerPidFile) + } const deliveryDir = path.join(collabRuntimeDir(teamId), 'delivery') if (fs.existsSync(deliveryDir)) fs.rmSync(deliveryDir, { recursive: true, force: true }) for (const f of [collabBridgeResult(teamId), collabBridgePosted(teamId)]) { diff --git a/tests/collab-poller.test.ts b/tests/collab-poller.test.ts new file mode 100644 index 0000000..3715634 --- /dev/null +++ b/tests/collab-poller.test.ts @@ -0,0 +1,307 @@ +/** + * Regression test for the feed poller that never stopped. + * + * collab-launch.sh starts one background loop per team that copies new lines + * from messages.jsonl into feed.txt. Until 2026-09 that loop was an inline + * `while true` with no exit condition, and nothing killed it on disband. On + * 2026-09-01 sixteen of them were still running, the oldest eleven days old, + * for teams the service no longer knew about. + * + * The poller now lives in scripts/collab-poller.sh and stops on its own when + * its team is finished or gone, and disbandTeam() kills it outright. + */ +import fs from 'fs' +import http from 'http' +import os from 'os' +import path from 'path' +import { spawn, type ChildProcess } from 'child_process' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { EnsembleTeam } from '../types/ensemble' + +const REPO = process.cwd() +const POLLER = path.join(REPO, 'scripts/collab-poller.sh') +const LAUNCH = path.join(REPO, 'scripts/collab-launch.sh') +const RUNTIME_ROOT = '/tmp/ensemble' + +function newTeamId(): string { + return `test-poller-${process.pid}-${Math.random().toString(36).slice(2, 8)}` +} + +function runtimeDir(teamId: string): string { + return path.join(RUNTIME_ROOT, teamId) +} + +async function waitFor(check: () => boolean, timeoutMs: number, label: string): Promise { + const deadline = Date.now() + timeoutMs + while (Date.now() < deadline) { + if (check()) return + await new Promise(r => setTimeout(r, 50)) + } + throw new Error(`timed out after ${timeoutMs}ms waiting for: ${label}`) +} + +function isAlive(pid: number): boolean { + try { + process.kill(pid, 0) + return true + } catch { + return false + } +} + +/** A stand-in for the ensemble service: the test decides what /teams/:id answers. */ +function fakeService(handler: (res: http.ServerResponse) => void): Promise<{ url: string; close: () => void }> { + return new Promise(resolve => { + const server = http.createServer((req, res) => { + if (req.url === '/api/v1/health') { + res.writeHead(200, { 'Content-Type': 'application/json' }) + res.end('{"status":"healthy"}') + return + } + handler(res) + }) + server.listen(0, '127.0.0.1', () => { + const { port } = server.address() as { port: number } + resolve({ url: `http://127.0.0.1:${port}`, close: () => server.close() }) + }) + }) +} + +describe('collab-poller.sh', () => { + let teamId: string + let child: ChildProcess | undefined + let exited = false + let exitCode: number | null = null + let closeService: (() => void) | undefined + + beforeEach(() => { + teamId = newTeamId() + fs.mkdirSync(runtimeDir(teamId), { recursive: true }) + fs.writeFileSync(path.join(runtimeDir(teamId), 'messages.jsonl'), '') + exited = false + exitCode = null + }) + + afterEach(() => { + if (child) { + // Stop reporting for this child, or its SIGKILL exit lands in the next test's variables. + child.removeAllListeners('exit') + if (child.pid && isAlive(child.pid)) child.kill('SIGKILL') + } + child = undefined + closeService?.() + closeService = undefined + fs.rmSync(runtimeDir(teamId), { recursive: true, force: true }) + }) + + function start(apiUrl: string, extraEnv: Record = {}): ChildProcess { + child = spawn('bash', [POLLER, teamId, apiUrl], { + stdio: 'ignore', + env: { + ...process.env, + COLLAB_POLL_SECS: '0.2', + COLLAB_POLLER_CHECK_EVERY: '1', + COLLAB_POLLER_MAX_API_FAILURES: '50', + ...extraEnv, + }, + }) + child.on('exit', code => { exited = true; exitCode = code }) + return child + } + + it('copies new messages into feed.txt', async () => { + const svc = await fakeService(res => { res.writeHead(200); res.end('{"team":{"status":"active"}}') }) + closeService = svc.close + start(svc.url) + const messages = path.join(runtimeDir(teamId), 'messages.jsonl') + const feed = path.join(runtimeDir(teamId), 'feed.txt') + + fs.appendFileSync(messages, '{"from":"codex","content":"hoi"}\n') + await waitFor(() => fs.existsSync(feed) && fs.readFileSync(feed, 'utf8').includes('hoi'), 3000, 'first line in feed') + fs.appendFileSync(messages, '{"from":"claude","content":"dag"}\n') + await waitFor(() => fs.readFileSync(feed, 'utf8').includes('dag'), 3000, 'second line in feed') + + expect(fs.readFileSync(feed, 'utf8')).toBe( + '{"from":"codex","content":"hoi"}\n{"from":"claude","content":"dag"}\n', + ) + expect(exited).toBe(false) + }) + + it('writes its own pid file and removes it on exit', async () => { + const svc = await fakeService(res => { res.writeHead(200); res.end('{"team":{"status":"active"}}') }) + closeService = svc.close + const proc = start(svc.url) + const pidFile = path.join(runtimeDir(teamId), 'poller.pid') + + await waitFor(() => fs.existsSync(pidFile), 3000, 'poller.pid written') + expect(Number(fs.readFileSync(pidFile, 'utf8').trim())).toBe(proc.pid) + + fs.writeFileSync(path.join(runtimeDir(teamId), '.finished'), new Date().toISOString()) + await waitFor(() => exited, 3000, 'poller exit after .finished') + expect(exitCode).toBe(0) + expect(fs.existsSync(pidFile)).toBe(false) + }) + + it('stops once the team is disbanded (.finished marker), after a last flush', async () => { + const svc = await fakeService(res => { res.writeHead(200); res.end('{"team":{"status":"active"}}') }) + closeService = svc.close + start(svc.url) + const feed = path.join(runtimeDir(teamId), 'feed.txt') + await waitFor(() => fs.existsSync(path.join(runtimeDir(teamId), 'poller.pid')), 3000, 'poller up') + + // Disband writes the last "X has left" messages and the marker right after each other. + fs.appendFileSync(path.join(runtimeDir(teamId), 'messages.jsonl'), '{"from":"ensemble","content":"codex has left"}\n') + fs.writeFileSync(path.join(runtimeDir(teamId), '.finished'), new Date().toISOString()) + + await waitFor(() => exited, 3000, 'poller exit after .finished') + expect(fs.readFileSync(feed, 'utf8')).toContain('codex has left') + }) + + it('stops when the runtime directory is gone', async () => { + const svc = await fakeService(res => { res.writeHead(200); res.end('{"team":{"status":"active"}}') }) + closeService = svc.close + start(svc.url) + await waitFor(() => fs.existsSync(path.join(runtimeDir(teamId), 'poller.pid')), 3000, 'poller up') + + fs.rmSync(runtimeDir(teamId), { recursive: true, force: true }) + await waitFor(() => exited, 3000, 'poller exit after rm -rf') + expect(exitCode).toBe(0) + }) + + it('stops when the service no longer knows the team', async () => { + const svc = await fakeService(res => { res.writeHead(404); res.end('{"error":"Team not found"}') }) + closeService = svc.close + start(svc.url) + await waitFor(() => exited, 3000, 'poller exit on 404') + expect(exitCode).toBe(0) + }) + + it('stops when the service reports the team as disbanded', async () => { + const svc = await fakeService(res => { res.writeHead(200); res.end('{"team":{"status":"disbanded"}}') }) + closeService = svc.close + start(svc.url) + await waitFor(() => exited, 3000, 'poller exit on disbanded') + expect(exitCode).toBe(0) + }) + + it('survives a service that is briefly unreachable', async () => { + // Port with nothing listening: connection refused on every check. + start('http://127.0.0.1:1', { COLLAB_POLLER_MAX_API_FAILURES: '50' }) + await new Promise(r => setTimeout(r, 1000)) + expect(exited).toBe(false) + }) + + it('stops when the service stays unreachable', async () => { + start('http://127.0.0.1:1', { COLLAB_POLLER_MAX_API_FAILURES: '2' }) + await waitFor(() => exited, 4000, 'poller exit after repeated failures') + expect(exitCode).toBe(0) + }) +}) + +describe('collab-launch.sh wiring', () => { + it('starts the poller script instead of an inline loop', () => { + const src = fs.readFileSync(LAUNCH, 'utf8') + expect(src).toContain('collab-poller.sh') + expect(src).not.toMatch(/while true; do\s*\n\s*M=\$\(wc -l/) + }) +}) + +describe('disbandTeam() kills the poller', () => { + const originalDataDir = process.env.ENSEMBLE_DATA_DIR + let tempRoot: string + let teamId: string + let sleeper: ChildProcess | undefined + + beforeEach(() => { + tempRoot = fs.mkdtempSync(path.join(os.tmpdir(), 'ensemble-poller-')) + process.env.ENSEMBLE_DATA_DIR = tempRoot + teamId = newTeamId() + fs.mkdirSync(runtimeDir(teamId), { recursive: true }) + }) + + afterEach(() => { + if (sleeper && sleeper.pid && isAlive(sleeper.pid)) sleeper.kill('SIGKILL') + sleeper = undefined + vi.restoreAllMocks() + vi.resetModules() + vi.doUnmock('../lib/ensemble-registry') + vi.doUnmock('../lib/agent-spawner') + vi.doUnmock('../lib/hosts-config') + vi.doUnmock('../lib/agent-runtime') + vi.doUnmock('../lib/agent-config') + vi.doUnmock('../lib/memory-export') + if (originalDataDir === undefined) { + delete process.env.ENSEMBLE_DATA_DIR + } else { + process.env.ENSEMBLE_DATA_DIR = originalDataDir + } + fs.rmSync(runtimeDir(teamId), { recursive: true, force: true }) + fs.rmSync(tempRoot, { recursive: true, force: true }) + }) + + it('the process in poller.pid is gone after disband', async () => { + // Stand-in for the poller: any long-lived process whose pid is in poller.pid. + sleeper = spawn('sleep', ['300'], { stdio: 'ignore' }) + const pid = sleeper.pid! + fs.writeFileSync(path.join(runtimeDir(teamId), 'poller.pid'), `${pid}\n`) + + const team: EnsembleTeam = { + id: teamId, + name: 'poller-team', + description: 'test', + status: 'active', + agents: [], + createdBy: 'test', + createdAt: '2026-09-03T08:00:00.000Z', + feedMode: 'live', + } + + vi.doMock('../lib/ensemble-registry', () => ({ + getMessages: vi.fn(() => []), + loadTeams: vi.fn(() => [team]), + appendMessage: vi.fn(), + updateTeam: vi.fn((_id: string, updates: Partial) => ({ ...team, ...updates })), + createTeam: vi.fn(), + getTeam: vi.fn(() => team), + saveTeams: vi.fn(), + })) + vi.doMock('../lib/agent-spawner', () => ({ + spawnLocalAgent: vi.fn(), + killLocalAgent: vi.fn(), + spawnRemoteAgent: vi.fn(), + killRemoteAgent: vi.fn(), + postRemoteSessionCommand: vi.fn(), + isRemoteSessionReady: vi.fn(), + getAgentTokenUsage: vi.fn(async () => 'unknown'), + })) + vi.doMock('../lib/hosts-config', () => ({ + isSelf: vi.fn(() => true), + getHostById: vi.fn(), + getSelfHostId: vi.fn(() => 'local'), + })) + vi.doMock('../lib/agent-runtime', () => ({ + getRuntime: vi.fn(() => ({ capturePane: vi.fn(), sendKeys: vi.fn(), pasteFromFile: vi.fn() })), + })) + vi.doMock('../lib/agent-config', () => ({ + resolveAgentProgram: vi.fn(() => ({ readyMarker: '>', inputMethod: 'sendKeys' })), + resolveAgentProgramDetailed: vi.fn((program: string) => ({ + agent: { command: program, readyMarker: '>', inputMethod: 'sendKeys' }, + how: 'exact', + requested: program, + })), + availableAgentKeys: vi.fn(() => ['claude', 'codex']), + })) + vi.doMock('../lib/memory-export', () => ({ + checkMemoryEndpoint: vi.fn(async () => ({ ok: true, endpoint: 'mock' })), + exportObservation: vi.fn(async () => ({ ok: true, endpoint: 'mock' })), + pendingExportFile: vi.fn((dir: string) => path.join(dir, 'pending-export.json')), + resolveMemoryEndpoint: vi.fn(() => 'mock'), + })) + + const mod = await import('../services/ensemble-service') + expect(isAlive(pid)).toBe(true) + await mod.disbandTeam(teamId) + + await waitFor(() => !isAlive(pid), 3000, 'poller process killed by disband') + }) +})