diff --git a/.env.example b/.env.example index 037f70a..bbdec81 100644 --- a/.env.example +++ b/.env.example @@ -29,6 +29,25 @@ GOOGLE_API_KEY= PI_DEFAULT_PROVIDER=openai PI_DEFAULT_MODEL=gpt-5-mini +# --- Fast LLM tier (v0.3.0+) -------------------------------------------------- +# Pi (CLI) handles slow deep analytics — post-mortem and 30-min reflection. +# A parallel "fast advisor" tier handles tactical "what to do RIGHT NOW" +# decisions when the reflex is stuck / preempted. Any OpenAI-compatible +# endpoint works: TimeWeb, OpenAI direct, Groq, OpenRouter, local Ollama. +# +# Leave blank to disable — the bot then runs without the fast tier +# (rc.1 scaffold is a safe no-op when API_KEY is unset). +# +# Example for TimeWeb (replace with your own endpoint + key + model): +# PEPA_FAST_LLM_BASE_URL=https://api.gen-api.ru/api/v1 +# PEPA_FAST_LLM_API_KEY= +# PEPA_FAST_LLM_MODEL=gpt-5-mini +PEPA_FAST_LLM_BASE_URL= +PEPA_FAST_LLM_API_KEY= +PEPA_FAST_LLM_MODEL= +# Optional: request timeout in ms (default 8000) +PEPA_FAST_LLM_TIMEOUT_MS= + # --- Bot behaviour ------------------------------------------------------------ # How often (seconds) the autonomous tick prompt fires. Set to 0 to disable. TICK_INTERVAL_SECONDS=60 diff --git a/dev/v0.3.0/STATUS.md b/dev/v0.3.0/STATUS.md index 2c60554..5834587 100644 --- a/dev/v0.3.0/STATUS.md +++ b/dev/v0.3.0/STATUS.md @@ -108,7 +108,51 @@ Tests: 315 green (was 279 on rc.1, +36 new): - `runtime/reflex.test.js` — 2 new integration tests (manifesto-on overrides curriculum; well-fed bot pursues tools_stone) -### rc.3 — (pending) Event-driven awareness + skill pre-emption +### rc.3 — Event-driven awareness + skill pre-emption +**Root problem solved**: in v0.2.x the reflex was purely polling. The +loop took a snapshot every DISPATCH_INTERVAL_MS (~2s) and decided what +to do, but anything that happened **between** ticks was invisible. +Concretely: when the operator dug a path that let the bot fall to a +new area, the bot continued executing its prior `explore.far` against +stale assumptions until the next tick. By then it had wandered further +off course, and the cycle never broke. Same problem for hostile spawns +and HP plunges — the reflex saw them only after the current skill ran +its 30-90s timeout. + +This rc gives the reflex an event-driven layer that **preempts** the +in-flight skill within ~100ms of an environmental shock. + +- [`runtime/awareness/events.js`](../../runtime/awareness/events.js) — + wires direct `bot.on(...)` listeners and surfaces them as flags + an + optional preempt callback: + - `bot.on("move")` — single-tick position jump ≥ 5 blocks (teleport, + fall, pathfinder snap, operator pushed us) → `forced_move` + - `bot.on("health")` — HP drop ≥ 2 in one tick → `health_plunge` + - `bot.on("entitySpawn")` — hostile mob spawns within 12 blocks → + `hostile_added` + - `bot.on("blockUpdate")` — block change within manhattan 4 → + `env_changed` (informational only, NOT preempting; throttled 800ms) +- [`runtime/skills/index.js`](../../runtime/skills/index.js): + - `RUNNER_CODES.PREEMPTED` — new stable failure code + - `runSkill()` now races `execute()` with `ctx.abortSignal`. If the + signal fires mid-await, the skill returns `{ ok: false, code: + "preempted" }` within one microtask — no skill code change needed. + Long-running skills (`gather.logs`, `explore.far`, + `recovery.tunnel-out`, `survive.pillar-up`) get this for free. +- [`runtime/bot.js`](../../runtime/bot.js): + - `dispatchAction` creates a fresh `AbortController` per dispatch + and stores it on `reflexCtx.currentAbort` + `reflexCtx.abortSignal` + - `bot.once("spawn")` calls `attachAwareness(bot, {onPreempt})` + where `onPreempt` aborts the current dispatch + - `reflexCtx.lastPreempt` records the most recent shock for + snapshot/telemetry consumers + +Tests: 332 green (was 315 on rc.2, +17 new): +- `runtime/awareness/events.test.js` — 12 tests (each event type, + thresholds, throttling, hostile filter) +- `runtime/skills/contract.test.js` — 3 new preempt tests (mid-flight + abort, pre-armed signal, clean signal doesn't interfere) +- 2 extra contract sanity checks shaken out by signal plumbing ## Next session quick start diff --git a/package.json b/package.json index 8934550..a656f5b 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "pepa-pi-bot", - "version": "0.3.0-rc.2", + "version": "0.3.0-rc.3", "private": true, "description": "An autonomous, self-extending Minecraft player powered by Pi and Mineflayer.", "license": "MIT", @@ -16,7 +16,7 @@ "tui": "tsx tui/tui.tsx", "propose:apply": "node scripts/propose-apply.js", "stop": "bash scripts/stop.sh", - "test": "node --test runtime/skills/contract.test.js runtime/skills/groups.test.js runtime/skills/compat.test.js runtime/skills/recovery-tunnel-out.test.js runtime/skills/pillar-up.test.js runtime/curriculum.test.js runtime/social/social.test.js runtime/social/conversation.test.js runtime/social/chat-history.test.js runtime/social/reply-pi.test.js runtime/stuck-incident.test.js runtime/compat.test.js runtime/reflex.test.js runtime/base-site.test.js runtime/locations.test.js runtime/watch-filter.test.js runtime/world-journal.test.js runtime/scenario-memory.test.js runtime/critic.test.js runtime/skill-library.test.js runtime/skill-registry.test.js runtime/modes.test.js runtime/pathfinder-watchdog.test.js runtime/manifesto/needs.test.js runtime/manifesto/state.test.js runtime/knowledge/knowledge.test.js runtime/llm/provider.test.js runtime/coach/postmortem.test.js runtime/coach/advice.test.js runtime/coach/reflect.test.js runtime/coach/fast-advisor.test.js runtime/persona/chatter.test.js scripts/edit-scope.test.js scripts/lint-patch.test.js" + "test": "node --test runtime/skills/contract.test.js runtime/skills/groups.test.js runtime/skills/compat.test.js runtime/skills/recovery-tunnel-out.test.js runtime/skills/pillar-up.test.js runtime/curriculum.test.js runtime/social/social.test.js runtime/social/conversation.test.js runtime/social/chat-history.test.js runtime/social/reply-pi.test.js runtime/stuck-incident.test.js runtime/compat.test.js runtime/reflex.test.js runtime/base-site.test.js runtime/locations.test.js runtime/watch-filter.test.js runtime/world-journal.test.js runtime/scenario-memory.test.js runtime/critic.test.js runtime/skill-library.test.js runtime/skill-registry.test.js runtime/modes.test.js runtime/pathfinder-watchdog.test.js runtime/manifesto/needs.test.js runtime/manifesto/state.test.js runtime/awareness/events.test.js runtime/knowledge/knowledge.test.js runtime/llm/provider.test.js runtime/coach/postmortem.test.js runtime/coach/advice.test.js runtime/coach/reflect.test.js runtime/coach/fast-advisor.test.js runtime/persona/chatter.test.js scripts/edit-scope.test.js scripts/lint-patch.test.js" }, "dependencies": { "better-sqlite3": "^11.10.0", diff --git a/runtime/awareness/events.js b/runtime/awareness/events.js new file mode 100644 index 0000000..04589f1 --- /dev/null +++ b/runtime/awareness/events.js @@ -0,0 +1,143 @@ +// Event-driven awareness. The reflex used to be polling-only: every +// DISPATCH_INTERVAL_MS the loop took a snapshot and decided what to do. +// That means anything happening *between* ticks — a creeper spawning, +// the bot taking damage, the bot being teleported by a falling block — +// was invisible until the next tick, and any active skill kept running +// against stale assumptions. +// +// This module wires direct mineflayer listeners that update a small +// flags object the reflex can consume each tick AND that triggers +// "preempt" callbacks (registered by the dispatcher) when something +// significant happens. The skill currently in flight can react by +// observing ctx.abortSignal.aborted between awaits. + +import { info } from "../log.js"; + +const HOSTILE_NAMES = new Set([ + "zombie", "skeleton", "creeper", "spider", "cave_spider", "witch", + "husk", "stray", "drowned", "phantom", "blaze", "ghast", "magma_cube", + "pillager", "vindicator", "vex", "wither_skeleton", "wither", "ravager", + "enderman", "endermite", "guardian", "elder_guardian", "evoker", "silverfish", + "hoglin", "zoglin", "piglin", "piglin_brute", "shulker", "warden", +]); + +// Heuristic thresholds — tunable later. +const FORCED_MOVE_BLOCKS = 5; // single tick movement > this = forced (teleport/fall/push) +const HEALTH_PLUNGE_DELTA = 2; // HP dropped by ≥ this in one tick = take note +const HOSTILE_CLOSE_BLOCKS = 12; // entity spawning within = preempt +const BLOCK_UPDATE_RADIUS = 4; // blockUpdate within manhattan = env-changed +const ENV_CHANGE_THROTTLE_MS = 800; + +export function attachAwareness(bot, { onPreempt = null } = {}) { + if (!bot || typeof bot.on !== "function") { + throw new Error("attachAwareness: bot.on missing"); + } + const state = createAwarenessState(); + let lastPos = bot.entity?.position ? cloneVec(bot.entity.position) : null; + let lastHealth = typeof bot.health === "number" ? bot.health : null; + let lastEnvChangeAt = 0; + + function preempt(reason, payload) { + try { onPreempt?.({ reason, payload, at: Date.now() }); } catch (e) { + info("awareness", `preempt callback threw: ${e?.message ?? e}`); + } + } + + bot.on("move", () => { + const pos = bot.entity?.position; + if (!pos) return; + const cur = cloneVec(pos); + if (lastPos) { + const dist = Math.hypot(cur.x - lastPos.x, cur.y - lastPos.y, cur.z - lastPos.z); + if (dist >= FORCED_MOVE_BLOCKS) { + state.flags.forcedMove = { at: Date.now(), from: lastPos, to: cur, distance: Math.round(dist * 10) / 10 }; + info("awareness", `forced move: ${state.flags.forcedMove.distance}b from (${Math.round(lastPos.x)}, ${Math.round(lastPos.y)}, ${Math.round(lastPos.z)}) to (${Math.round(cur.x)}, ${Math.round(cur.y)}, ${Math.round(cur.z)})`); + preempt("forced_move", state.flags.forcedMove); + } + } + lastPos = cur; + }); + + bot.on("health", () => { + const hp = bot.health; + if (typeof hp !== "number") return; + if (lastHealth !== null && hp + HEALTH_PLUNGE_DELTA <= lastHealth) { + state.flags.healthPlunge = { at: Date.now(), from: lastHealth, to: hp, delta: lastHealth - hp }; + info("awareness", `hp plunge: ${lastHealth} → ${hp}`); + preempt("health_plunge", state.flags.healthPlunge); + } + lastHealth = hp; + }); + + bot.on("entitySpawn", (entity) => { + if (!entity) return; + const name = (entity.name ?? "").toLowerCase(); + if (!HOSTILE_NAMES.has(name)) return; + const me = bot.entity?.position; + if (!me || !entity.position) return; + const dist = me.distanceTo(entity.position); + if (dist > HOSTILE_CLOSE_BLOCKS) return; + state.flags.hostileAdded = { at: Date.now(), name, distance: Math.round(dist * 10) / 10 }; + info("awareness", `hostile near: ${name}@${state.flags.hostileAdded.distance}m`); + preempt("hostile_added", state.flags.hostileAdded); + }); + + bot.on("blockUpdate", (oldBlock, newBlock) => { + const me = bot.entity?.position; + if (!me) return; + const block = newBlock ?? oldBlock; + const at = block?.position; + if (!at) return; + const manhattan = Math.abs(at.x - me.x) + Math.abs(at.y - me.y) + Math.abs(at.z - me.z); + if (manhattan > BLOCK_UPDATE_RADIUS) return; + const now = Date.now(); + if (now - lastEnvChangeAt < ENV_CHANGE_THROTTLE_MS) return; + lastEnvChangeAt = now; + state.flags.envChanged = { at: now, blockName: block?.name ?? "?", distance: manhattan }; + // envChanged is informational only — does NOT trigger preempt by + // default (block updates are too frequent during gather skills). + }); + + state._teardown = () => { + // node:events doesn't expose direct unbind without storing refs. + // In tests we just drop the bot. Real reflex never detaches. + }; + + info("awareness", "attached (forced_move + health_plunge + hostile_added + env_changed)"); + return state; +} + +export function createAwarenessState() { + return { + flags: { + forcedMove: null, + healthPlunge: null, + hostileAdded: null, + envChanged: null, + }, + consume() { + const out = { ...this.flags }; + this.flags = { + forcedMove: null, + healthPlunge: null, + hostileAdded: null, + envChanged: null, + }; + return out; + }, + hasPreempting() { + const f = this.flags; + return !!(f.forcedMove || f.healthPlunge || f.hostileAdded); + }, + }; +} + +function cloneVec(v) { + return { x: v.x, y: v.y, z: v.z }; +} + +// Test exports +export const __testing = { + HOSTILE_NAMES, FORCED_MOVE_BLOCKS, HEALTH_PLUNGE_DELTA, + HOSTILE_CLOSE_BLOCKS, BLOCK_UPDATE_RADIUS, ENV_CHANGE_THROTTLE_MS, +}; diff --git a/runtime/awareness/events.test.js b/runtime/awareness/events.test.js new file mode 100644 index 0000000..33e2c77 --- /dev/null +++ b/runtime/awareness/events.test.js @@ -0,0 +1,145 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { EventEmitter } from "node:events"; + +import { attachAwareness, createAwarenessState, __testing } from "./events.js"; + +function vec(x, y, z) { + return { + x, y, z, + distanceTo(other) { + return Math.hypot(this.x - other.x, this.y - other.y, this.z - other.z); + }, + }; +} + +function makeBot(pos = vec(0, 64, 0), hp = 20) { + const bot = new EventEmitter(); + bot.entity = { position: pos }; + bot.health = hp; + return bot; +} + +test("attachAwareness: throws when bot has no on()", () => { + assert.throws(() => attachAwareness({}), /bot\.on missing/); +}); + +test("createAwarenessState: starts with null flags, consume resets", () => { + const s = createAwarenessState(); + assert.equal(s.flags.forcedMove, null); + s.flags.forcedMove = { at: 1, from: {}, to: {}, distance: 7 }; + assert.equal(s.hasPreempting(), true); + const out = s.consume(); + assert.equal(out.forcedMove.distance, 7); + assert.equal(s.flags.forcedMove, null); +}); + +test("forcedMove: jump > threshold flags + preempts", () => { + const calls = []; + const bot = makeBot(vec(0, 64, 0)); + const state = attachAwareness(bot, { onPreempt: (e) => calls.push(e) }); + // move within threshold — no flag + bot.entity.position = vec(1, 64, 0); + bot.emit("move"); + assert.equal(state.flags.forcedMove, null); + assert.equal(calls.length, 0); + // teleport / fall — far jump + bot.entity.position = vec(20, 64, 0); + bot.emit("move"); + assert.ok(state.flags.forcedMove, "forcedMove flag set"); + assert.ok(state.flags.forcedMove.distance >= 18); + assert.equal(calls.length, 1); + assert.equal(calls[0].reason, "forced_move"); +}); + +test("healthPlunge: HP drop ≥ delta flags + preempts", () => { + const calls = []; + const bot = makeBot(vec(0, 64, 0), 20); + const state = attachAwareness(bot, { onPreempt: (e) => calls.push(e) }); + // trivial HP change does NOT flag + bot.health = 19; + bot.emit("health"); + assert.equal(state.flags.healthPlunge, null); + // big drop + bot.health = 12; + bot.emit("health"); + assert.ok(state.flags.healthPlunge); + assert.equal(state.flags.healthPlunge.from, 19); + assert.equal(state.flags.healthPlunge.to, 12); + assert.equal(calls.length, 1); + assert.equal(calls[0].reason, "health_plunge"); +}); + +test("hostileAdded: zombie nearby triggers preempt", () => { + const calls = []; + const bot = makeBot(); + const state = attachAwareness(bot, { onPreempt: (e) => calls.push(e) }); + const zombie = { name: "zombie", position: vec(2, 64, 0) }; + bot.emit("entitySpawn", zombie); + assert.ok(state.flags.hostileAdded); + assert.equal(state.flags.hostileAdded.name, "zombie"); + assert.equal(state.flags.hostileAdded.distance, 2); + assert.equal(calls.length, 1); + assert.equal(calls[0].reason, "hostile_added"); +}); + +test("hostileAdded: far hostile ignored", () => { + const bot = makeBot(); + const state = attachAwareness(bot); + const far = { name: "creeper", position: vec(50, 64, 0) }; + bot.emit("entitySpawn", far); + assert.equal(state.flags.hostileAdded, null); +}); + +test("hostileAdded: passive mob ignored", () => { + const bot = makeBot(); + const state = attachAwareness(bot); + const cow = { name: "cow", position: vec(2, 64, 0) }; + bot.emit("entitySpawn", cow); + assert.equal(state.flags.hostileAdded, null); +}); + +test("envChanged: nearby blockUpdate flags but does NOT preempt", () => { + const calls = []; + const bot = makeBot(); + const state = attachAwareness(bot, { onPreempt: (e) => calls.push(e) }); + const newBlock = { name: "cobblestone", position: vec(1, 64, 0) }; + bot.emit("blockUpdate", null, newBlock); + assert.ok(state.flags.envChanged); + assert.equal(state.flags.envChanged.blockName, "cobblestone"); + assert.equal(calls.length, 0, "env changes are observational, not preempting"); +}); + +test("envChanged: throttled", () => { + const bot = makeBot(); + const state = attachAwareness(bot); + const near = { name: "stone", position: vec(2, 64, 0) }; + bot.emit("blockUpdate", null, near); + const firstAt = state.flags.envChanged.at; + bot.emit("blockUpdate", null, near); + // second one within throttle window keeps the first timestamp + assert.equal(state.flags.envChanged.at, firstAt); +}); + +test("envChanged: far blockUpdate ignored", () => { + const bot = makeBot(); + const state = attachAwareness(bot); + const far = { name: "stone", position: vec(20, 64, 0) }; + bot.emit("blockUpdate", null, far); + assert.equal(state.flags.envChanged, null); +}); + +test("hasPreempting: true only for forcedMove/healthPlunge/hostileAdded", () => { + const s = createAwarenessState(); + assert.equal(s.hasPreempting(), false); + s.flags.envChanged = { at: 1, blockName: "stone", distance: 2 }; + assert.equal(s.hasPreempting(), false, "envChanged alone does not preempt"); + s.flags.hostileAdded = { at: 1, name: "creeper", distance: 5 }; + assert.equal(s.hasPreempting(), true); +}); + +test("thresholds: constants are sane", () => { + assert.ok(__testing.FORCED_MOVE_BLOCKS >= 3 && __testing.FORCED_MOVE_BLOCKS <= 10); + assert.ok(__testing.HEALTH_PLUNGE_DELTA >= 1 && __testing.HEALTH_PLUNGE_DELTA <= 5); + assert.ok(__testing.HOSTILE_CLOSE_BLOCKS >= 8); +}); diff --git a/runtime/bot.js b/runtime/bot.js index 977ff25..dcc1a22 100644 --- a/runtime/bot.js +++ b/runtime/bot.js @@ -61,6 +61,7 @@ import { initKnowledge } from "./knowledge/index.js"; import { attach as attachCoach } from "./coach/postmortem.js"; import { attach as attachReflect } from "./coach/reflect.js"; import { attach as attachChatter } from "./persona/chatter.js"; +import { attachAwareness } from "./awareness/events.js"; fs.mkdirSync(stateDir, { recursive: true }); const JOINED_FLAG = path.join(stateDir, "joined-before.flag"); @@ -78,6 +79,7 @@ const ESCALATION_COOLDOWN_MS = 10 * 60 * 1000; let bot = null; let pathWatchdog = null; +let awarenessState = null; let reflexPaused = false; let tickTimer = null; let reconnectTimer = null; @@ -237,6 +239,14 @@ function dispatchAction(fn, label, opts = {}) { } reflexCtx.busy = true; reflexCtx.currentActionLabel = label; + // v0.3.0-rc.3 — pre-emption: each dispatch gets a fresh AbortController. + // awareness/events.js#onPreempt fires controller.abort() when the env + // shocks (forced move, HP plunge, hostile spawn) the current skill + // shouldn't run against. runSkill races execute() with the signal and + // returns code: "preempted" within one microtask. + const dispatchAbort = new AbortController(); + reflexCtx.currentAbort = dispatchAbort; + reflexCtx.abortSignal = dispatchAbort.signal; const startedAt = Date.now(); // Capture the situation hash BEFORE the action runs so a failure is // attributable to the state at dispatch time, not the state after the @@ -312,6 +322,10 @@ function dispatchAction(fn, label, opts = {}) { .finally(() => { reflexCtx.busy = false; reflexCtx.currentActionLabel = null; + if (reflexCtx.currentAbort === dispatchAbort) { + reflexCtx.currentAbort = null; + reflexCtx.abortSignal = null; + } }); } @@ -673,6 +687,22 @@ function connect() { try { attachCoach(bot, { stateDir, askPi }); } catch (e) { warn("coach", `attach: ${e?.message ?? e}`); } try { attachReflect({ bot, stateDir, askPi, getSnapshot: () => lastSnapshot }); } catch (e) { warn("reflect", `attach: ${e?.message ?? e}`); } try { attachChatter(bot, { getSnapshot: () => lastSnapshot }); } catch (e) { warn("persona", `attach: ${e?.message ?? e}`); } + // v0.3.0-rc.3 — awareness layer: listens to bot.on('move'/'health'/ + // 'entitySpawn'/'blockUpdate') and aborts the current dispatch via + // reflexCtx.currentAbort when something disrupts the in-flight skill. + try { + awarenessState = attachAwareness(bot, { + onPreempt: ({ reason, payload }) => { + const abort = reflexCtx.currentAbort; + if (abort && !abort.signal.aborted) { + info("preempt", `aborting ${reflexCtx.currentActionLabel ?? "?"} due to ${reason}`); + abort.abort(); + } + reflexCtx.lastPreempt = { reason, payload, at: Date.now() }; + }, + }); + reflexCtx.awareness = awarenessState; + } catch (e) { warn("awareness", `attach: ${e?.message ?? e}`); } }); bot.on("messagestr", (text) => { diff --git a/runtime/skills/contract.test.js b/runtime/skills/contract.test.js index 566e010..27667ad 100644 --- a/runtime/skills/contract.test.js +++ b/runtime/skills/contract.test.js @@ -166,3 +166,63 @@ test("result missing code defaults to runner DONE on success", async () => { teardown(); } }); + +test("abortSignal: mid-execute abort surfaces code: preempted", async () => { + const teardown = _registerForTest({ + id: "test.preempt-midflight", + timeoutMs: 5000, + preconditions: () => ({ ok: true }), + execute: async () => { + await new Promise((r) => setTimeout(r, 1500)); + return { ok: true }; + }, + }); + const controller = new AbortController(); + const runP = runSkill("test.preempt-midflight", { abortSignal: controller.signal }); + setTimeout(() => controller.abort(), 30); + try { + const res = await runP; + assert.equal(res.ok, false); + assert.equal(res.code, RUNNER_CODES.PREEMPTED); + } finally { + teardown(); + } +}); + +test("abortSignal: pre-aborted signal short-circuits to preempted", async () => { + const teardown = _registerForTest({ + id: "test.preempt-prearm", + timeoutMs: 5000, + preconditions: () => ({ ok: true }), + execute: async () => { + await new Promise((r) => setTimeout(r, 200)); + return { ok: true }; + }, + }); + const controller = new AbortController(); + controller.abort(); + try { + const res = await runSkill("test.preempt-prearm", { abortSignal: controller.signal }); + assert.equal(res.ok, false); + assert.equal(res.code, RUNNER_CODES.PREEMPTED); + } finally { + teardown(); + } +}); + +test("abortSignal: not aborted → skill completes normally", async () => { + const teardown = _registerForTest({ + id: "test.preempt-clear", + timeoutMs: 5000, + preconditions: () => ({ ok: true }), + execute: async () => ({ ok: true, code: "done" }), + }); + const controller = new AbortController(); + try { + const res = await runSkill("test.preempt-clear", { abortSignal: controller.signal }); + assert.equal(res.ok, true); + assert.equal(res.code, "done"); + } finally { + teardown(); + } +}); diff --git a/runtime/skills/index.js b/runtime/skills/index.js index 0fcc65d..5440028 100644 --- a/runtime/skills/index.js +++ b/runtime/skills/index.js @@ -118,6 +118,7 @@ export const RUNNER_CODES = Object.freeze({ TIMEOUT: "timeout", THREW: "threw", VALIDATION_FAILED: "validation_failed", + PREEMPTED: "preempted", DONE: "done", }); @@ -139,6 +140,42 @@ function withTimeout(promise, ms, label) { return Promise.race([promise, timeout]).finally(() => clearTimeout(timer)); } +// v0.3.0-rc.3 — wrap execute() so that if ctx.abortSignal fires we +// stop awaiting (and surface code: "preempted"). The skill itself +// doesn't need to read the signal — the race below ensures runSkill +// returns control to the reflex within one microtask of abort(). The +// skill's own async work may continue in the background harmlessly, +// because the next dispatch will overwrite any shared state. +function raceWithAbort(promise, signal) { + if (!signal) return promise; + if (signal.aborted) { + return Promise.reject(Object.assign(new Error("preempted"), { _preempted: true })); + } + return new Promise((resolve, reject) => { + let settled = false; + const onAbort = () => { + if (settled) return; + settled = true; + reject(Object.assign(new Error("preempted"), { _preempted: true })); + }; + signal.addEventListener("abort", onAbort, { once: true }); + promise.then( + (v) => { + if (settled) return; + settled = true; + signal.removeEventListener?.("abort", onAbort); + resolve(v); + }, + (e) => { + if (settled) return; + settled = true; + signal.removeEventListener?.("abort", onAbort); + reject(e); + }, + ); + }); +} + // Drive one skill through its full lifecycle. The caller (typically reflex.js // or, eventually, a higher-level scheduler) decides when to invoke; runSkill // only owns the contract enforcement. @@ -172,12 +209,19 @@ export async function runSkill(id, ctx, args = {}) { const timeoutMs = skill.timeoutMs ?? 30_000; let raw; try { - raw = await withTimeout(skill.execute(ctx, args), timeoutMs, `skill(${id})`); + raw = await withTimeout( + raceWithAbort(skill.execute(ctx, args), ctx?.abortSignal), + timeoutMs, + `skill(${id})`, + ); } catch (e) { const isTimeout = /timed out after/.test(e.message); + const isPreempted = e?._preempted === true; const result = { ok: false, - code: isTimeout ? RUNNER_CODES.TIMEOUT : RUNNER_CODES.THREW, + code: isPreempted + ? RUNNER_CODES.PREEMPTED + : isTimeout ? RUNNER_CODES.TIMEOUT : RUNNER_CODES.THREW, detail: e.message, worldDelta: null, };