Files
pepa-pi-bot/runtime/bot.js
T
mayatnikovandClaude Opus 4.7 2308596164 feat(runtime): observability + no-progress detector (Phase 1)
Phase 1 of plans/autonomous-survival-bot-prd.md. The bot must always be
able to answer "what am I doing and why am I not doing more?" without
parsing the log stream.

New modules:
- runtime/state.js: pure FSM classifier emitting emergency / working /
  recovering / planning / social / idle from snapshot + reflex context.
- runtime/no-progress.js: sliding-window detector that watches position
  and inventory; when both are unchanged for 60 s+, emits one stable
  reason code from REASONS (waiting_for_day, night_hostile_nearby,
  no_food_source, inventory_full, no_reachable_target, planner_empty,
  awaiting_action_cooldown).
- runtime/viewer.js: optional prismarine-viewer launcher behind
  VIEWER_PORT. Lazy import so the dep is not required by default.

Wiring:
- runtime/bot.js: tick() now computes runtimeState + noProgressReason
  every tick and stamps them on the snapshot along with activeSkill,
  currentMilestone (read from plan.md, cached 30 s), lastResult,
  failuresByCode and lastEscalation.
- runtime/bot.js: dispatchAction records lastResult and lastFailureAt
  for the recovering-state classifier.
- runtime/planner.js: exports isPlannerBusy(), readNextMilestone()
  and planExists() so the runtime can show planning state + current
  milestone without spawning extra Pi calls.
- runtime/config.js: adds VIEWER_PORT support.

TUI:
- tui/tui.tsx: StatusBar gains a state badge, current-skill row,
  milestone row, no-progress reason warning, last-result line with
  ok/fail color, failures-by-class summary and last-escalation age.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-25 22:10:41 +03:00

722 lines
24 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Long-running Mineflayer process. Owns:
// - the MC TCP connection + reconnect policy
// - the reflex tick loop (no LLM in hot path)
// - the IPC server for TUI clients
// - on-demand Pi-headless escalation (manual or automatic)
//
// MC chat is dialog-only (see plans/autonomous-survival-bot-prd.md, FR1).
// Player/operator chat may produce a social reply but never dispatches a
// movement/build/mining task. TUI is the only local control plane.
//
// Lifecycle: started by `npm run bot`. Connects to MC, spawns IPC server,
// ticks every TICK_INTERVAL_SECONDS, broadcasts STATUS to clients each tick.
// SIGINT / SIGTERM: graceful disconnect + socket cleanup + exit.
import fs from "node:fs";
import path from "node:path";
import mineflayer from "mineflayer";
import { config, stateDir, redactedConfig } from "./config.js";
import { info, warn, error } from "./log.js";
import { snapshot as buildSnapshot } from "./perceive.js";
import { runTick } from "./reflex.js";
import { createIpcServer } from "./ipc-server.js";
import { askPi } from "./pi-bridge.js";
import { COMMAND_TYPES, EVENT_TYPES } from "./ipc-protocol.js";
import {
readCurrentTask,
writeCurrentTask,
clearCurrentTask,
appendDiary,
writeProposal,
listProposals,
readProposal,
approveProposal,
} from "./state-store.js";
import { startAutoImprover } from "./auto-improve.js";
import { startPlanner, isPlannerBusy, readNextMilestone, planExists } from "./planner.js";
import { computeState, STATES } from "./state.js";
import { createNoProgressDetector } from "./no-progress.js";
import { maybeStartViewer } from "./viewer.js";
fs.mkdirSync(stateDir, { recursive: true });
const JOINED_FLAG = path.join(stateDir, "joined-before.flag");
// Auto-escalation tunables. With tick=3s, 20 noops ≈ 1 minute idle before we
// even consider asking Pi. Cooldown prevents spamming the LLM when the bot
// is permanently stuck on the same situation.
const ESCALATE_AFTER_NOOPS = 20;
const ESCALATION_COOLDOWN_MS = 10 * 60 * 1000;
let bot = null;
let reflexPaused = false;
let tickTimer = null;
let reconnectTimer = null;
let shuttingDown = false;
let lastSnapshot = { connected: false };
let consecutiveNoops = 0;
let lastEscalationAt = 0;
// Observability state — surfaced in every STATUS snapshot so the TUI (and
// future Telegram/diary surfaces) can answer "what is the bot doing and why
// isn't it doing more?" without parsing the log stream.
const noProgress = createNoProgressDetector();
let lastResult = null; // { label, ok, code, detail, ts }
let lastFailureAt = 0;
let lastPlanReadAt = 0;
let cachedMilestone = null;
let cachedPlanExists = false;
const MILESTONE_CACHE_MS = 30_000;
// Reflex context — passed into reflex.js every tick. Mutable across ticks.
const reflexCtx = {
bot: null,
snapshot: lastSnapshot,
busy: false,
currentActionLabel: null,
idleCounter: 0,
lastEatAt: 0,
lastSleepAttemptAt: 0,
// Tracks repeated failure of the same labelled action — triggers a proposal.
recentFailures: [], // [{label, detail, ts}], capped at 10
dispatch: dispatchAction,
};
let chatTimestamps = [];
const CHAT_WINDOW_MS = 60_000;
let ipc;
// ---- chat rate limit -------------------------------------------------------
function chatRateAllowed() {
const now = Date.now();
chatTimestamps = chatTimestamps.filter((t) => now - t < CHAT_WINDOW_MS);
if (chatTimestamps.length >= config.chatRateLimitPerMin) return false;
chatTimestamps.push(now);
return true;
}
function botChat(text) {
if (!bot) return;
if (!chatRateAllowed()) {
warn("chat", `dropped chat (rate-limited): ${text.slice(0, 60)}`);
return;
}
bot.chat(text);
}
// ---- auth ------------------------------------------------------------------
function hasJoinedBefore() {
return fs.existsSync(JOINED_FLAG);
}
function markJoinedBefore() {
try {
fs.writeFileSync(JOINED_FLAG, new Date().toISOString());
} catch (e) {
warn("auth", `could not write joined flag: ${e.message}`);
}
}
function maybeHandleAuthPrompt(text) {
if (!bot || !config.authmePassword) return;
const lower = text.toLowerCase();
const sawRegister = lower.includes("/register");
const sawLogin = lower.includes("/login");
if (!sawRegister && !sawLogin) return;
const cmd = sawLogin ? "login" : hasJoinedBefore() ? "login" : "register";
if (cmd === "register") {
bot.chat(`/register ${config.authmePassword} ${config.authmePassword}`);
info("auth", "sent /register (password redacted)");
} else {
bot.chat(`/login ${config.authmePassword}`);
info("auth", "sent /login (password redacted)");
}
markJoinedBefore();
}
// ---- action dispatch -------------------------------------------------------
// Reflexes call this to fire an async action without blocking the tick.
// Sets busy=true, runs fn, clears busy when done; optional onComplete callback
// receives the action's { ok, detail } result.
function dispatchAction(fn, label, opts = {}) {
if (reflexCtx.busy) {
warn("dispatch", `tried to dispatch ${label} while busy with ${reflexCtx.currentActionLabel}`);
return;
}
reflexCtx.busy = true;
reflexCtx.currentActionLabel = label;
// current-task is a resume anchor — keep it small. Embedding the full
// perception snapshot blows the file up to ~3 KB per write × every action.
writeCurrentTask({ label, status: "in_progress", position: lastSnapshot.position });
info("dispatch", `→ ${label}`);
Promise.resolve()
.then(() => fn())
.then((res) => {
const ok = !!res?.ok;
info(
"dispatch",
`← ${label} ${ok ? "ok" : "fail"}${res?.detail ? ` (${JSON.stringify(res.detail).slice(0, 80)})` : ""}`,
);
writeCurrentTask({ label, status: ok ? "completed" : "failed", detail: res?.detail });
lastResult = {
label,
ok,
code: res?.code ?? (ok ? "done" : classifyFailure(res?.detail)),
detail: res?.detail,
ts: Date.now(),
};
if (!ok) {
lastFailureAt = lastResult.ts;
recordFailure(label, res?.detail);
} else {
clearRecentFailures(label);
}
if (opts.onComplete) {
try {
opts.onComplete(res ?? { ok: false, detail: "no result" });
} catch (e) {
warn("dispatch", `onComplete threw for ${label}: ${e.message}`);
}
}
})
.catch((e) => {
warn("dispatch", `${label} threw: ${e?.message ?? e}`);
writeCurrentTask({ label, status: "threw", detail: String(e?.message ?? e) });
lastResult = {
label,
ok: false,
code: "threw",
detail: String(e?.message ?? e),
ts: Date.now(),
};
lastFailureAt = lastResult.ts;
recordFailure(label, String(e?.message ?? e));
})
.finally(() => {
reflexCtx.busy = false;
reflexCtx.currentActionLabel = null;
});
}
// ---- failure tracking + proposal detection --------------------------------
//
// A proposal is a request to the LLM to patch the codebase. They cost tokens
// and may produce risky patches that need rolling back. We file them ONLY for
// failures that genuinely look like bugs the reflex layer can't handle on its
// own. Everything else is a feature gap the script should solve via reflex
// chain reordering, cooldowns, or new primitives.
const PROPOSAL_THRESHOLD = 5; // raised from 3 to dampen spam
let lastProposalAt = 0;
const PROPOSAL_COOLDOWN_MS = 30 * 60 * 1000;
// Detail substrings that mean "this is a known feature gap, the bot handles
// it via reflex routing already". Don't file a proposal — the bot will switch
// strategies on its own. If something here is wrong, fix the routing.
const NORMAL_FAILURE_SUBSTRINGS = [
"no reachable log",
"no log within",
"no bed in range",
"no food in inventory",
"no target in reach",
"rate-limited",
"can't see you nearby",
"returned false",
"no result",
];
// Detail substrings that look like a real bug — patch-worthy.
const BUG_FAILURE_SUBSTRINGS = [
"TypeError",
"ReferenceError",
"Cannot read properties",
"is not a function",
"is not iterable",
"is not defined",
"unknown block",
"unknown item",
];
function classifyFailure(detail) {
const s = String(detail ?? "");
if (BUG_FAILURE_SUBSTRINGS.some((sub) => s.includes(sub))) return "bug";
if (NORMAL_FAILURE_SUBSTRINGS.some((sub) => s.includes(sub))) return "feature-gap";
if (s.includes("timed out")) return "timeout";
return "other";
}
function recordFailure(label, detail) {
const kind = classifyFailure(detail);
reflexCtx.recentFailures.push({ ts: Date.now(), label, detail, kind });
if (reflexCtx.recentFailures.length > 20) reflexCtx.recentFailures.shift();
maybeFileProposal(label);
}
function clearRecentFailures(label) {
reflexCtx.recentFailures = reflexCtx.recentFailures.filter((f) => f.label !== label);
}
function maybeFileProposal(label) {
// Same-label trailing run.
const trailing = [];
for (let i = reflexCtx.recentFailures.length - 1; i >= 0; i--) {
const f = reflexCtx.recentFailures[i];
if (f.label === label) trailing.push(f);
else break;
}
if (trailing.length < PROPOSAL_THRESHOLD) return;
if (Date.now() - lastProposalAt < PROPOSAL_COOLDOWN_MS) return;
// Only file when the run is dominated by bug-class failures (any single
// bug counts) OR persistent timeouts on the same operation. Feature gaps
// are skipped — the reflex layer should re-route, not the LLM.
const anyBug = trailing.some((f) => f.kind === "bug");
const allTimeout = trailing.every((f) => f.kind === "timeout");
if (!anyBug && !allTimeout) return;
lastProposalAt = Date.now();
const summary = `${label} failed ${trailing.length}× in a row (${anyBug ? "bug" : "persistent timeout"})`;
const slimSnapshot = lastSnapshot && {
position: lastSnapshot.position,
health: lastSnapshot.health,
food: lastSnapshot.food,
inventory: lastSnapshot.inventory,
isDay: lastSnapshot.isDay,
closestHostile: lastSnapshot.closestHostile,
dimension: lastSnapshot.dimension,
};
const body = [
`# Repeated failure: ${label}`,
"",
`Class: **${anyBug ? "bug" : "persistent timeout"}**.`,
"",
"## What happened",
"",
`The reflex layer dispatched \`${label}\` ${trailing.length} times in succession without a single success.`,
"",
"## Most recent failures",
"",
...trailing
.slice(0, 5)
.map(
(f, i) =>
`${i + 1}. \`${new Date(f.ts).toISOString()}\` [${f.kind}] ${JSON.stringify(f.detail).slice(0, 200)}`,
),
"",
"## Slim snapshot",
"",
"```json",
JSON.stringify(slimSnapshot, null, 2),
"```",
"",
"## Constraints for the patch",
"",
"- Touch only files under `runtime/`. Don't touch `extensions/`, `tui/`, or any docs.",
"- Don't introduce new npm dependencies.",
"- Don't change `.env` or anything in `state/`.",
"- Don't push, don't open a PR. Commit on the current branch only.",
"- Prefer the smallest viable fix. A 3-line guard is better than a 30-line refactor.",
"- If the failure is genuinely irrecoverable (server-side, not code), document it in a code comment and exit 1.",
].join("\n");
const { filename } = writeProposal({ kind: `repeated-fail-${label}`, summary, body });
warn("proposal", `filed ${filename}: ${summary}`);
appendDiary(`proposal filed: ${filename} (${summary})`);
}
// ---- chat replies (dialog-only) --------------------------------------------
//
// As of Phase 0 of the survival-bot PRD, MC chat is dialog-only for everyone,
// including OPERATOR_USERNAMES. The bot may reply socially or answer status
// questions, but it does NOT dispatch movement/build/mining tasks from chat.
// If a player addresses the bot with a command-like verb (come, follow, build,
// pause, stop, give), the bot acknowledges that chat is dialog-only and
// records the ignored command in the diary.
let lastChatReplyAt = 0;
const CHAT_REPLY_COOLDOWN_MS = 30_000;
const GREETING_RE = /\b(hi|hello|hey|yo|sup|hola|привет|здаров|здарова|здорова|здравствуй|здравствуйте|салам)\b/i;
const STATUS_RE = /\b(status|how are you|what are you doing|whats up|what[']?s up|чё делаешь|что делаешь|как ты|как дела|статус)\b/i;
const COMMAND_LIKE_RE = /\b(come|follow|build|pause|resume|stop|go to|goto|tp|teleport|give|drop|attack|kill|dig|mine|chop|farm|harvest|sleep here|иди сюда|подойди|следуй|остановись|стоп|пауза|строй|копай|дай)\b/i;
function isOperator(username) {
if (!username) return false;
return config.operators.includes(username.toLowerCase());
}
function buildStatusReply() {
const s = lastSnapshot;
const parts = [];
if (reflexCtx.busy) parts.push(`busy=${reflexCtx.currentActionLabel}`);
if (s.health !== undefined) parts.push(`hp=${s.health}/20`);
if (s.food !== undefined) parts.push(`food=${s.food}/20`);
if (s.position) parts.push(`pos=${s.position.x},${s.position.y},${s.position.z}`);
if (s.hostileCount) parts.push(`hostiles=${s.hostileCount}`);
return parts.join(" ") || "alive";
}
function handleChat(username, text) {
if (!bot) return;
const trimmed = text.trim();
const lower = trimmed.toLowerCase();
const botname = bot.username.toLowerCase();
const addressed = lower.includes(botname);
const looksLikeCommand = addressed && COMMAND_LIKE_RE.test(lower);
// Command-like chat (from anyone, including operators) is recorded but not
// dispatched. We tell the speaker once per cooldown so they aren't left
// wondering why nothing happened.
if (looksLikeCommand) {
const op = isOperator(username) ? "operator" : "player";
info("chat", `ignored command-like chat from ${op} ${username}: ${trimmed.slice(0, 80)}`);
appendDiary(`ignored command-like chat from ${username}: ${trimmed.slice(0, 120)}`);
const since = Date.now() - lastChatReplyAt;
if (since >= CHAT_REPLY_COOLDOWN_MS) {
lastChatReplyAt = Date.now();
botChat(`${username}: MC chat is dialog-only — operator uses the TUI to drive me.`);
}
return;
}
// Dialog: greeting, status question, or addressed banter.
const wantsStatus = addressed && STATUS_RE.test(lower);
const isGreeting = GREETING_RE.test(lower);
if (!addressed && !isGreeting) return;
const since = Date.now() - lastChatReplyAt;
if (since < CHAT_REPLY_COOLDOWN_MS) return;
lastChatReplyAt = Date.now();
if (wantsStatus) {
botChat(`${username}: ${buildStatusReply()}`);
return;
}
const greetings = ["yo", "hey", "hi", "привет"];
const reply = greetings[Math.floor(Math.random() * greetings.length)];
botChat(`${username}: ${reply}`);
}
// ---- connect ---------------------------------------------------------------
function connect() {
if (bot) return;
info("mc", `connecting as ${config.username}${config.host}:${config.port} (v${config.version})`);
bot = mineflayer.createBot({
host: config.host,
port: config.port,
username: config.username,
auth: config.authMode === "microsoft" ? "microsoft" : "offline",
version: config.mineflayerVersion,
hideErrors: false,
});
reflexCtx.bot = bot;
bot.once("spawn", () => {
info("mc", `spawned at ${JSON.stringify(bot.entity.position)}`);
appendDiary(`spawned at ${bot.entity.position.x.toFixed(0)},${bot.entity.position.y.toFixed(0)},${bot.entity.position.z.toFixed(0)}`);
ipc?.broadcast(EVENT_TYPES.STATUS, buildSnapshot(bot));
maybeStartViewer(bot).catch((e) => warn("viewer", `start threw: ${e?.message ?? e}`));
});
bot.on("messagestr", (text) => {
ipc?.broadcast(EVENT_TYPES.CHAT, { from: "server", text, kind: "system" });
maybeHandleAuthPrompt(text);
});
bot.on("chat", (username, message) => {
if (username === bot.username) return;
ipc?.broadcast(EVENT_TYPES.CHAT, { from: username, text: message, kind: "player" });
try {
handleChat(username, message);
} catch (e) {
warn("chat", `chat handler threw: ${e.message}`);
}
});
bot.on("death", () => {
const pos = bot.entity?.position;
warn("mc", `died at ${JSON.stringify(pos)}`);
appendDiary(`died at ${pos?.x.toFixed(0)},${pos?.y.toFixed(0)},${pos?.z.toFixed(0)}`);
ipc?.broadcast(EVENT_TYPES.DEATH, { reason: "unknown", position: pos });
clearCurrentTask();
});
bot.on("kicked", (reason) => {
warn("mc", `kicked: ${reason}`);
});
bot.on("error", (err) => {
error("mc", `bot error: ${err?.message ?? err}`);
});
bot.on("end", (reason) => {
warn("mc", `connection ended: ${reason}`);
bot = null;
reflexCtx.bot = null;
lastSnapshot = { connected: false };
if (!shuttingDown) scheduleReconnect();
});
}
function scheduleReconnect() {
if (reconnectTimer || shuttingDown) return;
const delay = 5000;
info("mc", `reconnecting in ${delay}ms`);
reconnectTimer = setTimeout(() => {
reconnectTimer = null;
connect();
}, delay);
}
// ---- escalation ------------------------------------------------------------
function maybeAutoEscalate() {
if (reflexCtx.busy) return;
if (consecutiveNoops < ESCALATE_AFTER_NOOPS) return;
const since = Date.now() - lastEscalationAt;
if (since < ESCALATION_COOLDOWN_MS) return;
lastEscalationAt = Date.now();
consecutiveNoops = 0;
const promptCtx = JSON.stringify(lastSnapshot, null, 2);
const prompt = [
`You are the escalation cortex for an autonomous Minecraft bot. The bot's`,
`script-driven reflex loop has produced no useful action for ${ESCALATE_AFTER_NOOPS} consecutive ticks`,
`(~${Math.round((ESCALATE_AFTER_NOOPS * config.tickIntervalMs) / 1000)}s). The reflex chain is:`,
` defend > eat > sleep > tech-tree > autonomous > idle`,
`Snapshot:`,
"```json",
promptCtx,
"```",
`Decide ONE next thing the bot should attempt. Output as plain text — what it should do and why, in 3 lines max.`,
`If "do nothing" is the right answer, say so.`,
`Do NOT propose code changes. Do NOT propose anything beyond what mineflayer + the existing actions (attack, flee, eat, sleep, goTo) can do.`,
].join("\n");
info("escalation", `firing askPi after ${ESCALATE_AFTER_NOOPS} noops`);
askPi({
prompt,
onChunk: (chunk) => ipc?.broadcast(EVENT_TYPES.ASK_PI_CHUNK, chunk),
onDone: (result) => {
info("escalation", `pi done code=${result.code} dur=${result.durationMs}ms`);
ipc?.broadcast(EVENT_TYPES.ASK_PI_DONE, result);
},
});
}
// ---- tick ------------------------------------------------------------------
function failuresByCode() {
const counts = {};
for (const f of reflexCtx.recentFailures) {
const k = f.kind || "other";
counts[k] = (counts[k] ?? 0) + 1;
}
return counts;
}
function refreshMilestoneCache(now) {
if (now - lastPlanReadAt < MILESTONE_CACHE_MS) return;
lastPlanReadAt = now;
try {
cachedMilestone = readNextMilestone();
cachedPlanExists = planExists();
} catch (e) {
warn("planner", `milestone read failed: ${e.message}`);
}
}
function tick() {
if (shuttingDown) return;
const now = Date.now();
if (bot && bot.entity) {
lastSnapshot = buildSnapshot(bot);
lastSnapshot.pendingProposals = listProposals().length;
lastSnapshot.lastReflex = reflexCtx.lastReflex ?? null;
lastSnapshot.busy = reflexCtx.busy
? { label: reflexCtx.currentActionLabel ?? "?" }
: null;
reflexCtx.snapshot = lastSnapshot;
if (!reflexPaused) {
const result = runTick(reflexCtx);
if (!result || result.action === "noop") {
consecutiveNoops++;
maybeAutoEscalate();
} else if (result.action === "skipped") {
// busy — neither productive nor stuck; don't increment noops
} else {
consecutiveNoops = 0;
}
}
// Observability: compute runtime state + no-progress reason and
// stamp them on the snapshot so the TUI / future surfaces can show
// one concrete answer to "why is the bot idle?".
refreshMilestoneCache(now);
const plannerInFlight = isPlannerBusy();
const runtimeState = computeState({
snapshot: lastSnapshot,
ctx: reflexCtx,
plannerInFlight,
lastChatReplyAt,
lastFailureAt,
now,
});
const noProgressReason = noProgress.detect({
snapshot: lastSnapshot,
ctx: reflexCtx,
planExists: cachedPlanExists,
now,
});
lastSnapshot.runtimeState = runtimeState;
lastSnapshot.activeSkill = reflexCtx.busy
? reflexCtx.currentActionLabel
: reflexCtx.lastReflex?.label ?? null;
lastSnapshot.currentMilestone = cachedMilestone;
lastSnapshot.lastResult = lastResult;
lastSnapshot.noProgressReason = noProgressReason;
lastSnapshot.failuresByCode = failuresByCode();
lastSnapshot.lastEscalation = lastEscalationAt
? { ts: lastEscalationAt, ageMs: now - lastEscalationAt }
: null;
lastSnapshot.reflexPaused = reflexPaused;
ipc?.broadcast(EVENT_TYPES.STATUS, lastSnapshot);
} else {
lastSnapshot = { connected: false };
ipc?.broadcast(EVENT_TYPES.STATUS, lastSnapshot);
}
}
function startTickLoop() {
if (tickTimer) clearInterval(tickTimer);
tickTimer = setInterval(tick, config.tickIntervalMs);
}
// ---- IPC commands ----------------------------------------------------------
function handleCommand(msg, send) {
switch (msg.type) {
case COMMAND_TYPES.PAUSE:
reflexPaused = true;
info("ipc", "reflex paused by client");
break;
case COMMAND_TYPES.RESUME:
reflexPaused = false;
info("ipc", "reflex resumed by client");
break;
case COMMAND_TYPES.STOP:
info("ipc", "stop requested by client");
gracefulExit(0);
break;
case COMMAND_TYPES.CHAT: {
const text = (msg.payload?.text || "").trim();
if (!text || !bot) return;
if (!chatRateAllowed()) {
send(EVENT_TYPES.ERROR, { source: "chat", text: "rate-limited" });
return;
}
bot.chat(text);
break;
}
case COMMAND_TYPES.ASK_PI: {
const prompt = msg.payload?.prompt;
if (!prompt) return;
askPi({
prompt,
onChunk: (chunk) => ipc?.broadcast(EVENT_TYPES.ASK_PI_CHUNK, chunk),
onDone: (result) => ipc?.broadcast(EVENT_TYPES.ASK_PI_DONE, result),
});
break;
}
case COMMAND_TYPES.SNAPSHOT:
send(EVENT_TYPES.STATUS, lastSnapshot);
break;
case COMMAND_TYPES.PROPOSAL_LATEST: {
const all = listProposals();
if (all.length === 0) {
send(EVENT_TYPES.PROPOSAL, { filename: null, body: null, total: 0 });
return;
}
const filename = all[all.length - 1];
send(EVENT_TYPES.PROPOSAL, {
filename,
body: readProposal(filename),
total: all.length,
});
break;
}
case COMMAND_TYPES.PROPOSAL_APPROVE: {
const filename = msg.payload?.filename;
if (!filename) return;
try {
const dst = approveProposal(filename);
info("proposal", `approved ${filename}${dst}`);
appendDiary(`proposal approved: ${filename}`);
} catch (e) {
send(EVENT_TYPES.ERROR, { source: "proposal", text: e.message });
}
break;
}
default:
warn("ipc", `unknown command type: ${msg.type}`);
}
}
// ---- shutdown --------------------------------------------------------------
function gracefulExit(code) {
if (shuttingDown) return;
shuttingDown = true;
info("runtime", "shutting down");
if (tickTimer) clearInterval(tickTimer);
if (reconnectTimer) clearTimeout(reconnectTimer);
try {
bot?.quit("shutdown");
} catch {}
ipc?.close();
setTimeout(() => process.exit(code), 500);
}
process.on("SIGINT", () => gracefulExit(0));
process.on("SIGTERM", () => gracefulExit(0));
info("runtime", `pepa runtime starting; cfg=${JSON.stringify(redactedConfig())}`);
// Resume info: surface stale state across restarts. We don't auto-resume any
// action — but we tell the operator if the bot died mid-task last time, and
// the count of pending proposals.
const lastTask = readCurrentTask();
if (lastTask && lastTask.label) {
info("resume", `previous task: ${lastTask.label} (${lastTask.status ?? "?"}, ${lastTask.ts ?? "?"})`);
if (lastTask.status === "in_progress") {
warn("resume", `last shutdown happened mid-action — operator should review state/<host>/current-task.json`);
}
}
const pendingProposals = listProposals();
if (pendingProposals.length > 0) {
warn("resume", `${pendingProposals.length} pending proposal(s) — see state/<host>/proposals/`);
}
ipc = createIpcServer({
getStatusSnapshot: () => ({
...lastSnapshot,
pendingProposals: listProposals().length,
}),
onCommand: handleCommand,
});
connect();
startTickLoop();
startAutoImprover();
startPlanner(() => lastSnapshot);