v0.3.0-rc.3: event-driven awareness + skill pre-emption #26

Closed
halofourteen wants to merge 2 commits from v0.3.0-rc.3 into v0.3.0-rc.2
8 changed files with 490 additions and 5 deletions
+19
View File
@@ -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=<paste your key here>
# 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
+45 -1
View File
@@ -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
+2 -2
View File
@@ -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",
+143
View File
@@ -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,
};
+145
View File
@@ -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);
});
+30
View File
@@ -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) => {
+60
View File
@@ -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();
}
});
+46 -2
View File
@@ -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,
};