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
1 change: 1 addition & 0 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
18 changes: 2 additions & 16 deletions scripts/collab-launch.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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")"

Expand Down Expand Up @@ -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..."
Expand Down
79 changes: 79 additions & 0 deletions scripts/collab-poller.sh
Original file line number Diff line number Diff line change
@@ -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 <team-id> [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 <team-id> [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
13 changes: 12 additions & 1 deletion services/ensemble-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -1141,6 +1141,17 @@ export async function disbandTeam(teamId: string): Promise<ServiceResult<{ team:

// Soft cleanup: remove ephemeral files, keep messages/summary/log, write .finished marker
try {
// The feed poller started by collab-launch.sh also stops on the .finished
// marker, but only on its next tick; kill it here so a disband never leaves
// a loop behind (on 2026-09-01 sixteen of them had outlived their teams).
const pollerPidFile = collabPollerPid(teamId)
if (fs.existsSync(pollerPidFile)) {
const pollerPid = parseInt(fs.readFileSync(pollerPidFile, 'utf8').trim(), 10)
if (pollerPid > 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)]) {
Expand Down
Loading
Loading