Capacity unit consolidation. Three coherent themes, small commits inside. ## Census before/after (`node scripts/lifecycle-column-census.mjs`) | | before | after | |---|---:|---:| | triage column guards (the bar) | 10 | **10** | | `--strict` on main | ❌ **RED** | ✅ green | | baseline staleness | 14 files stale | **0** | This branch does **not** move the triage bar — its remaining 10 are moves.ts (dies with the flag), the dashboard cluster, and one deliberate site. It fixes the instrument that measures the bar, plus a live defect the comparison count cannot see. --- ## 1. `--strict` was RED on clean `origin/main`, and it was my fault ``` packages/dashboard/src/routes/register-task-workflow-routes.ts: 22 -> 23 ``` My merged #2621 added a v1-IR pre-WIP fallback answering a greptile P1 and shipped no marker or baseline update, so the program's measuring instrument has been failing on main since it landed. Fixed **at the site** with a `DELIBERATE-LITERAL` marker, not by bumping the baseline. That branch runs only when the IR declares no columns and no nodes, so there is no role to resolve — `resolveLifecycleColumns` returns nothing and the legacy pre-implementation ids are the only pre-WIP signal that exists there. It is *unconvertible*, not unfinished; the sibling `else` two lines down is the trait path for every IR that can answer. A rise that is genuinely correct belongs where a reader will see it. ## 2. The baseline was stale for 14 files — a hole, not cosmetics A stale allowance lets converted guards return while the check stays green. Measured gaps: ``` self-healing.ts allows 126, tree has 111 executor.ts allows 112, tree has 104 moves.ts allows 44, tree has 39 default-workflow-hooks allows 25, tree has 7 mission-feature-sync allows 5, tree has 0 MissionControlPanel allows 4, tree has 0 (+8 more) ``` **Only two of the fourteen are mine.** The other twelve are already-merged conversions by other workers where nobody re-recorded. Re-recorded all fourteen here rather than waiting for twelve PRs, because until it happens the ratchet is not holding the 779 it exists to hold. Flagging it plainly: those drops are other people's work being locked in, not mine being claimed. ## 3. Routines created tasks into the column U11 deleted The routine editor's "Target Column" defaulted to `triage`. That value is submitted as the create step's `taskColumn`, and an **explicit** column bypasses the workflow entry-column resolution added for column-less creates (#2589) — so every routine saved with the untouched default seeded its tasks into a column the board does not declare. Defaulting to `todo` would be the same mistake one column over: a custom workflow declaring no `todo` is seeded into an undeclared column just as surely, because an explicit column overrides entry resolution whatever its value. So the default sends **nothing** and each workflow's own intake resolution decides. The `triage` **option** is removed too, not merely un-defaulted — fixing the initializer alone left the operator able to pick the deleted column one click away, and it was the option labelled "Planning", the name the merged `todo` column now displays. Removing it retires that label inversion as well. Found by scanning **membership** forms rather than comparisons: the comparison census cannot see a `?? "triage"` default, so no count showed this and nobody was looking. Revert-proof — restoring the default fails with *"the default must not name a column at all"*. ## 4. "Worktrees off is INERT" had one unaudited reader The constraint was that `maxWorktrees` become genuinely inert, "not set very high and not skipped by convention". `resolveWorktreeCapacityLimit` returns `null` for that, and its unit tests can only prove the **resolver** is right — they cannot see a second reader, which is the only way the constraint breaks. Audited every `maxWorktrees` read that bounds anything. **Exactly two:** `scheduler.ts` (the admission gate, via the resolver, single call site, optional gate snapshot) and `self-healing.ts`'s `enforceWorktreeCap` — `(settings.maxWorktrees ?? 4) * 2`, a **raw** read. The second is **not a bug** and is left alone: it bounds worktree *directories on disk* and only removes *idle* ones. Worktrees still exist in OFF mode, so that bound must keep applying or idle directories accumulate unbounded. Recorded consequence: in OFF mode the number still governs disk retention while gating no admission — an edge you scoped out. The note says explicitly **not** to unify the two readers: routing hygiene through the resolver returns `null` in OFF mode and silently removes the disk bound, which is a leak dressed as a simplification. New ratchet requires every file bounding on `maxWorktrees` to be named with a reason, and rejects a **stale** allowlist entry. Proven by injecting `active >= (settings.maxWorktrees ?? 4)` into `hybrid-executor.ts`. --- ## Deliberately NOT included - **My own census script.** #2633 landed the canonical one, and it is better than mine — an AST classifier *plus* an independent text classifier with `--compare`, and a baseline that fails on unrecorded **drops** as well as rises. Mine only caught rises. I deleted mine rather than ship a second measuring instrument; three copies of "strip comments" is the drift shape this program keeps paying for, so the worktree ratchet now imports #2633's `stripComments`. - **My TaskContextMenu fix.** Superseded, and by a better answer: main's `isPureIntakeColumn` (intake *without* hold) keeps the merged Planning column shown and suppresses only a bare Ideas capture, which resolves the exact hold-lane objection coderabbit raised against my version. I briefly clobbered that merged work by checking my old file out wholesale, caught it in the diff, and reverted. ## Verification `pnpm lint` clean · core + dashboard `tsc` clean · census suite 23/23 · worktree ratchet 8/8 · RoutineEditor 49/49 · `routes-task-retry-planning-column` 16/16 · `lifecycle-column-census --strict` exits 0. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --- ## Added after review (all four greptile threads were real, and two of them mattered) **The routine fix was half a fix.** `routine-runner.ts:515` *and* `cron-runner.ts:982` both did `column: (step.taskColumn as Column) || "triage"` **after** the step is read, so every routine — including ones saved through the fixed editor — still created tasks into the deleted column. Both now omit it. **The advanced steps editor MANUFACTURED the defect.** `ScheduleStepsEditor.tsx` had three `triage` defaults: the new-step template (`:64`), the per-step initializer (`:95`), and the select still offering it (`:344`). So the path I had *not* fixed produced the bug by default, on fresh data. Template names no column; initializer coerces a persisted `triage`; `triage` removed from the options; empty submits `undefined`. **Four pre-existing tests pinned the defect** and are rewritten to the corrected invariant rather than appeased: | test | asserted | |---|---| | `cron-runner`: "defaults column to triage when taskColumn is not set" | `column: "triage"` | | `ScheduleStepsEditor`: "adds a create-task step..." | `taskColumn` toBe `"triage"` | | `ScheduleStepsEditor`: "allows saving create-task step..." | the legacy column is **resubmitted** | | plus the explicit-column case added beside each, so the fix cannot swallow a deliberate choice | **The allowlist hole was the worst finding.** `AUDITED_BOUNDS` was keyed by FILE, so every bounding expression in an allowlisted file was exempt — a second raw bound in `scheduler.ts` stayed green, the one case that ratchet exists for. Per-expression now, and making it so **immediately surfaced a real second bound the file-level version was hiding** (`maxWorktreesGate.used >= maxWorktreesGate.limit`, safe by construction since the snapshot is `undefined` in OFF mode). Proven by injection. ## Found while re-reading my own deletion, not reported A **rendered tooltip** still named a deleted cap. The "Queued to plan" badge read *"planning starts when a concurrency slot frees up (maxConcurrent / globalMaxConcurrent)"*. The cross-project cap is gone — capacity is two numbers per project — so it told operators their planning waited on a limiter they can no longer find a setting for. Names the surviving dimension only now. ## Coding (Ideas): enforcing #2651 rather than repeating it I took the unowned coding-ideas IR merge, concluded it must not be done, then found **#2651 had already implemented, reverted and documented exactly that** — with better grounding than my own argument. It added no test, so nothing stops the next person reaching the same dead end. So this ships their reasoning as a ratchet, not a second opinion: triage discovery keys on the column's `autoTriage`, so a merged column is either never scanned (cards sit on a bootstrap stub until the **capacity hold** releases them, sending **unplanned** work into in-progress — worse than stalling) or scanning wins and the manual gate is gone. Their scope caveat is kept: `autoTriage` is a general trait field, so only *this preset's* collapse is dead, not manual intake as a concept. The registry does not reject the merged shape, which is why prose was not enough. ## Verification (re-run) `pnpm lint` clean · core + engine + dashboard-app `tsc` clean · `lifecycle-column-census --strict` exits 0 ("every file matches its baseline exactly") · routine-runner 24/24 · cron-runner 156/156 · ScheduleStepsEditor 41/41 · RoutineEditor 49/49 · worktree + coding-ideas 12/12. TaskCard has 2 failures **pre-existing on main** — confirmed identical with my changes stashed. --- ## Bears directly on the closing bar: this PR already removes the 67-guard ratchet slack Measured on current `origin/main` with the census itself: ``` tree total: 787 baseline total: 854 SLACK: 67 FILES ABOVE BASELINE (1): +1 packages/dashboard/src/routes/register-task-workflow-routes.ts (22 -> 23) FILES BELOW BASELINE: 13, totalling 68 unrecorded conversions -18 core/default-workflow-hooks.ts (25->7) -15 engine/self-healing.ts (126->111) -8 engine/executor.ts (112->104) -5 core/task-store/moves.ts (44->39) -5 engine/mission-feature-sync.ts (5->0) -4 core/live-agent-count.ts (10->6) ``` **The slack is not regression — it is 13 files of merged conversions nobody re-recorded**, against exactly **one** rise. This PR re-records the baseline **854 → 782 across 140 files**, which closes it. **And the "+3 that slipped in" is +1, and it is mine.** `register-task-workflow-routes.ts 22 → 23` is the v1-IR pre-WIP fallback my #2621 added; it is justified (that branch runs only when the IR declares no columns or nodes, so there is no role to resolve) but it shipped with no marker and no baseline update — which is why `--strict` has been **red on main since it merged**. Fixed here at the site with a `DELIBERATE-LITERAL` marker rather than by bumping the baseline, because a rise that is genuinely correct belongs where a reader will see it. Sequencing note for the auto-lowering change: if this lands first, that work is purely the mechanism (auto-lower, or fail with tighten instructions) rather than a cleanup, and the two re-records will not collide in the same file. Also worth carrying into that mechanism, from building the same guard here: **`--update` must refuse to RAISE.** An earlier version of mine wrote current counts verbatim, so a developer who added a literal and ran the documented update command locked the regression in as the new ceiling — the mirror of the high-water problem. Lowering can be unattended; raising should be a hand edit with the reason recorded. ## Third piece of residue from my own deletion `updateGlobalConcurrency` in the dashboard API client PUT to `/api/global-concurrency`, a route removed when the machine-wide cap went. Zero callers; the only reference was the `legacy.ts` barrel re-export. Deleted both. `fetchGlobalConcurrency` **survives on purpose** — the GET route remains and serves live utilization telemetry to the footer and Command Center; nothing gates on it. That is the third: after the second raw `maxWorktrees` reader and the "Queued to plan" tooltip. A deletion is not finished when the enforcement goes — the client, the label and the tooltip outlive it. --- ## Re-greened the dashboard API tests: 117 failures on main, ONE root cause These would have polluted the closing verification pass, and nobody owned them. `api()` builds headers via `new Headers(...)` and returns `Object.fromEntries(headers.entries())` — and `Headers.entries()` **lowercases every key**, so the object reaching `fetch` is `content-type`, not `Content-Type`. `ab87d0d80` then added `x-fusion-client: dashboard-ui` for run-audit attribution. Both changes are correct; neither is visible at a call site, so **114 assertions across 7 files** kept asserting the old shape and went red together. Fixed by naming the shape **once** in `app/test/apiRequestHeaders.ts` rather than patching 114 literals — restating a shared fact 114 times is what made a two-line client change look like 117 failures. Deliberately not a loose `objectContaining`: these tests are the only thing pinning that the attribution header is sent *at all*. **117 → 4.** The remaining 4 are unrelated pre-existing CSS failures (`task-detail-modal-tablet-width` ×3, `space-token-defined` ×1) — confirmed identical on clean main with my changes stashed. ### A gap this surfaced, recorded not papered over Three routes failed in the *opposite* direction — they send the old shape because they call `fetch()` **directly**, bypassing `api()`, so they never get the attribution header. `client.ts` claims the opposite: > "Applied once here rather than per-call so no future mutation route has to remember it." That does not hold for a route that bypasses the helper it is applied in. **Measured in `app/api/`: 8 files make direct `fetch()` calls and 7 include mutations (POST/DELETE)** — among them `ai-sessions.ts`'s DELETE, which is the same class as the four-delete incident the header was added for. So the attribution fix has a hole in exactly its motivating case. Not fixed here: routing those onto `api()` is a behaviour change across the API layer and belongs to its owner, not to a test re-green. Those assertions use a separate `API_JSON_HEADERS_NO_ATTRIBUTION` constant so the gap stays **visible** — if a route is later moved onto `api()`, its test fails and points at the note explaining why. --- ## This branch takes the triage bar 10 → 5, and makes `--strict` green `node scripts/lifecycle-column-census.mjs` on this branch reports **triage 5**, against **10** on `origin/main`. The five removed are the ScheduleStepsEditor template/initializer/option and the RoutineEditor default/option — the automation paths that were creating tasks into the deleted column. **`--strict` was also RED on clean main, twice over, and both causes were the same mistake:** a thorough written rationale the tool cannot read, because the marker was not where the census looks. The census reads a comparison node's **leading comments**; a `DELIBERATE-LITERAL` in the JSDoc above the enclosing function or declaration does not reach the comparison inside it. | site | why it is legitimate | why the tool could not see it | |---|---|---| | `columnRoles.ts:80` `isHoldColumnRole` | degrades to `columnId === "todo"` only when a column has **no resolved traits** — identical in kind to `LEGACY_PRE_IMPLEMENTATION_COLUMN_IDS` directly above, which escapes counting only because a Set is a membership form | rationale written, **no marker token** | | `MissionControlPanel.tsx` ×3 | the SDLC funnel **alias table** — maps `to-do`/`ready`/`review`/`shipped` onto one display stage with an explicit `other` bucket, and nothing branches on it | marker in the JSDoc; the comparisons are arrow bodies **inside the array literal**, which it does not reach | The second only surfaced because converting the `triage` stage to a Set removed its count and exposed the siblings — red gate, justification sitting three lines above, unreachable. Both are markers, no behaviour change. Neither is a conversion candidate: resolving the funnel table to traits would **drop the non-column aliases it exists to accept**. **For the auto-lowering work:** the marker-placement rule is now the recurring trap — three instances, three different authors, including me. A marker that does not register is indistinguishable from no marker, and the failure mode is a red gate with a written explanation nobody can act on. If the census accepted a marker anywhere in the enclosing declaration's comments, none of the three would have happened. Baseline re-recorded per the tool's own instruction ("Re-record the baseline in the SAME PR that lowered the count"). <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit - **Bug Fixes** - Task “Actions” menus no longer appear on bare cards in the Planning column. - Routines, scheduled tasks, and create-task steps now respect each board’s configured workflow intake column instead of using a retired default. - Legacy tasks saved with the retired intake column are migrated to automatic workflow resolution. - Target-column selection now offers only “Automatic (workflow intake)” and “Planning,” removing the obsolete option. - Capacity/planning messaging and related UI tooltip text were clarified; concurrency cap updates are managed per project. - **Tests** - Added/updated coverage for workflow intake resolution, create-task target column behavior (including legacy coercion), capacity safeguards, and API request consistency. <!-- end of auto-generated comment: release notes by coderabbit.ai --> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
706 lines
26 KiB
TypeScript
706 lines
26 KiB
TypeScript
/**
|
|
* RoutineRunner — orchestrates routine execution via the heartbeat system.
|
|
*
|
|
* - Validates routine state before execution (enabled, has assigned agent)
|
|
* - Enforces concurrency policies (parallel/skip/queue/replace)
|
|
* - Handles catch-up for missed runs
|
|
* - Triggers heartbeat execution for routines
|
|
*/
|
|
|
|
import { CronExpressionParser } from "cron-parser";
|
|
import { isInProcessBackupCommand, isInProcessMemoryBackupCommand } from "./cron-runner.js";
|
|
import type {
|
|
RoutineStore,
|
|
Routine,
|
|
RoutineExecutionResult,
|
|
AutomationRunResult,
|
|
AutomationStep,
|
|
AutomationStepResult,
|
|
Column,
|
|
TaskCreateInput,
|
|
TaskStore,
|
|
} from "@fusion/core";
|
|
import type { HeartbeatMonitor } from "./agent-heartbeat.js";
|
|
import type { AiPromptExecutor, AiPromptLiveCallbacks } from "./cron-runner.js";
|
|
import { createLogger } from "./logger.js";
|
|
import { defaultShell } from "./shell-utils.js";
|
|
import { resolveSandboxBackend } from "./sandbox/index.js";
|
|
import { createRunAuditor, generateSyntheticRunId } from "./run-audit.js";
|
|
import type { EngineRunContext, RunAuditor } from "./run-audit.js";
|
|
import type { SandboxBackend } from "./sandbox/types.js";
|
|
|
|
const log = createLogger("routine-runner");
|
|
const DEFAULT_TIMEOUT_MS = 5 * 60 * 1000;
|
|
const MAX_BUFFER = 1024 * 1024;
|
|
const MAX_OUTPUT_LENGTH = 10 * 1024;
|
|
|
|
|
|
/** Options for RoutineRunner constructor */
|
|
/*
|
|
FNXC:AutomationLiveOutput 2026-06-26-00:00:
|
|
Routine manual triggers share the automation live-output contract. Thread optional callbacks through the runner so routes can stream step boundaries, AI text/tool events, and final output without changing scheduled/background execution behavior.
|
|
*/
|
|
export type RoutineLiveRunCallbacks = AiPromptLiveCallbacks & {
|
|
onStep?: (data: Record<string, unknown>) => void;
|
|
};
|
|
|
|
export interface RoutineRunnerOptions {
|
|
/** RoutineStore for querying and updating routines */
|
|
routineStore: RoutineStore;
|
|
/** HeartbeatMonitor for triggering agent execution */
|
|
heartbeatMonitor: HeartbeatMonitor;
|
|
/** Project root directory */
|
|
rootDir: string;
|
|
/** Optional task store used when routines execute schedule-style actions. */
|
|
taskStore?: TaskStore;
|
|
/** Optional AI prompt executor for ai-prompt action steps. */
|
|
aiPromptExecutor?: AiPromptExecutor;
|
|
}
|
|
|
|
/**
|
|
* Maximum number of catch-up executions to prevent runaway loops.
|
|
*/
|
|
const MAX_CATCH_UP_INTERVALS = 10;
|
|
|
|
/**
|
|
* RoutineRunner orchestrates routine execution via the heartbeat system.
|
|
*
|
|
* Key behaviors:
|
|
* - Enforces concurrency policies before starting executions
|
|
* - Handles catch-up for missed runs based on catch-up policy
|
|
* - Triggers heartbeats with routine context in the trigger detail
|
|
*/
|
|
export class RoutineRunner {
|
|
private options: RoutineRunnerOptions;
|
|
/** Tracks currently-running executions by routine ID */
|
|
private inFlightExecutions: Map<string, Promise<RoutineExecutionResult>> = new Map();
|
|
|
|
constructor(options: RoutineRunnerOptions) {
|
|
this.options = options;
|
|
}
|
|
|
|
/**
|
|
* Execute a routine by ID with a given trigger type.
|
|
*
|
|
* @param routineId - ID of the routine to execute
|
|
* @param triggerType - What triggered this execution: "cron", "webhook", or "api"
|
|
* @param context - Additional context passed to the heartbeat execution
|
|
* @returns The execution result
|
|
* @throws Error if routine not found or disabled
|
|
*/
|
|
async executeRoutine(
|
|
routineId: string,
|
|
triggerType: "cron" | "webhook" | "api",
|
|
context?: Record<string, unknown>,
|
|
liveCallbacks?: RoutineLiveRunCallbacks,
|
|
): Promise<RoutineExecutionResult> {
|
|
// 1. Load routine
|
|
let routine: Routine;
|
|
try {
|
|
routine = await this.options.routineStore.getRoutine(routineId);
|
|
} catch {
|
|
throw new Error(`Routine '${routineId}' not found`);
|
|
}
|
|
|
|
// 2. Validate routine state
|
|
if (!routine.enabled) {
|
|
throw new Error(`Routine '${routineId}' is disabled`);
|
|
}
|
|
|
|
if (!this.hasRoutineAction(routine) && !routine.agentId) {
|
|
throw new Error(`Routine '${routineId}' has no assigned agent`);
|
|
}
|
|
|
|
// 3. Enforce concurrency policy
|
|
const concurrency = routine.executionPolicy ?? "queue";
|
|
|
|
if (concurrency === "reject" && this.inFlightExecutions.has(routineId)) {
|
|
log.log(`Routine ${routineId} rejected — already running`);
|
|
// Return a failed result without creating an execution record
|
|
return {
|
|
routineId,
|
|
success: false,
|
|
output: "Routine rejected — already running",
|
|
error: "Routine rejected — already running",
|
|
startedAt: new Date().toISOString(),
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
|
|
// If queue, wait for existing execution
|
|
if (concurrency === "queue" && this.inFlightExecutions.has(routineId)) {
|
|
log.log(`Routine ${routineId} queued — waiting for existing execution`);
|
|
const existingResult = await this.inFlightExecutions.get(routineId);
|
|
if (existingResult) {
|
|
await existingResult;
|
|
}
|
|
}
|
|
|
|
// 4. Record execution start
|
|
const startedAt = new Date().toISOString();
|
|
|
|
// Set in-flight BEFORE starting execution to prevent race conditions
|
|
const executionPromise = this.runExecution(routine, triggerType, context, startedAt, liveCallbacks, true);
|
|
this.inFlightExecutions.set(routineId, executionPromise);
|
|
|
|
try {
|
|
await this.options.routineStore.startRoutineExecution(routineId, {
|
|
triggeredAt: startedAt,
|
|
invocationSource: "routine",
|
|
});
|
|
|
|
const result = await executionPromise;
|
|
return result;
|
|
} finally {
|
|
this.inFlightExecutions.delete(routineId);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Execute an already-claimed central global routine.
|
|
* The scheduler owns central run-state persistence because project RoutineStore
|
|
* cannot address central.global_routines.
|
|
*/
|
|
async executeGlobalRoutine(
|
|
routine: Routine,
|
|
triggerType: "cron" | "webhook" | "api",
|
|
context?: Record<string, unknown>,
|
|
): Promise<RoutineExecutionResult> {
|
|
const startedAt = new Date().toISOString();
|
|
const executionPromise = this.runExecution(routine, triggerType, context, startedAt, undefined, false);
|
|
this.inFlightExecutions.set(routine.id, executionPromise);
|
|
try {
|
|
return await executionPromise;
|
|
} finally {
|
|
this.inFlightExecutions.delete(routine.id);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Internal execution logic for a routine.
|
|
*/
|
|
private async runExecution(
|
|
routine: Routine,
|
|
triggerType: string,
|
|
context: Record<string, unknown> | undefined,
|
|
startedAt: string,
|
|
liveCallbacks: RoutineLiveRunCallbacks | undefined,
|
|
persistProjectRunState: boolean,
|
|
): Promise<RoutineExecutionResult> {
|
|
const routineId = routine.id;
|
|
|
|
try {
|
|
const actionResult = this.hasRoutineAction(routine)
|
|
? await this.executeRoutineAction(routine, startedAt, liveCallbacks)
|
|
: await this.executeAgentRoutine(routine, triggerType, context);
|
|
|
|
if (persistProjectRunState) {
|
|
await this.options.routineStore.completeRoutineExecution(routineId, {
|
|
completedAt: actionResult.completedAt,
|
|
success: actionResult.success,
|
|
resultJson: actionResult.success ? { output: actionResult.output } : undefined,
|
|
output: actionResult.output,
|
|
error: actionResult.error,
|
|
triggerType: triggerType as RoutineExecutionResult["triggerType"],
|
|
stepResults: actionResult.stepResults,
|
|
});
|
|
}
|
|
|
|
return {
|
|
routineId,
|
|
success: actionResult.success,
|
|
output: actionResult.output,
|
|
startedAt,
|
|
completedAt: actionResult.completedAt,
|
|
error: actionResult.error,
|
|
triggerType: triggerType as RoutineExecutionResult["triggerType"],
|
|
stepResults: actionResult.stepResults,
|
|
};
|
|
} catch (err) {
|
|
const errorMessage = err instanceof Error ? err.message : String(err);
|
|
log.error(`Routine ${routineId} execution failed: ${errorMessage}`);
|
|
|
|
// Record failure
|
|
if (persistProjectRunState) {
|
|
try {
|
|
await this.options.routineStore.completeRoutineExecution(routineId, {
|
|
completedAt: new Date().toISOString(),
|
|
success: false,
|
|
error: errorMessage,
|
|
});
|
|
} catch (persistError) {
|
|
log.error(`[${routineId}] Failed to persist error state: ${persistError}`);
|
|
}
|
|
}
|
|
|
|
return {
|
|
routineId,
|
|
success: false,
|
|
output: errorMessage,
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
error: errorMessage,
|
|
};
|
|
}
|
|
}
|
|
|
|
private hasRoutineAction(routine: Routine): boolean {
|
|
return Boolean((routine.steps && routine.steps.length > 0) || routine.command?.trim());
|
|
}
|
|
|
|
private async executeAgentRoutine(
|
|
routine: Routine,
|
|
triggerType: string,
|
|
context: Record<string, unknown> | undefined,
|
|
): Promise<AutomationRunResult> {
|
|
const run = await this.options.heartbeatMonitor.executeHeartbeat({
|
|
agentId: routine.agentId,
|
|
source: "routine",
|
|
triggerDetail: `routine:${routine.id}:${triggerType}`,
|
|
contextSnapshot: {
|
|
routineId: routine.id,
|
|
routineName: routine.name,
|
|
triggerType,
|
|
...context,
|
|
},
|
|
});
|
|
|
|
if (run.status === "failed" || run.status === "terminated") {
|
|
const error = run.stderrExcerpt || `Run ${run.status}`;
|
|
return {
|
|
success: false,
|
|
output: error,
|
|
error,
|
|
startedAt: new Date().toISOString(),
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
|
|
return {
|
|
success: true,
|
|
output: run.resultJson ? JSON.stringify(run.resultJson) : "Routine completed successfully",
|
|
startedAt: new Date().toISOString(),
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
|
|
private async executeRoutineAction(
|
|
routine: Routine,
|
|
startedAt: string,
|
|
liveCallbacks?: RoutineLiveRunCallbacks,
|
|
): Promise<AutomationRunResult> {
|
|
if (routine.steps && routine.steps.length > 0) {
|
|
return this.executeSteps(routine, startedAt, liveCallbacks);
|
|
}
|
|
liveCallbacks?.onStep?.({ stepIndex: 0, stepId: "command", stepName: routine.name, stepType: "command", status: "started" });
|
|
const result = await this.executeCommand(routine, routine.command ?? "", routine.timeoutMs, startedAt);
|
|
liveCallbacks?.onStep?.({ stepIndex: 0, stepId: "command", stepName: routine.name, stepType: "command", status: "completed", success: result.success, error: result.error });
|
|
if (result.output) liveCallbacks?.onText?.(result.output);
|
|
return result;
|
|
}
|
|
|
|
private getRoutineCommandAuditor(routine: Routine): RunAuditor | undefined {
|
|
if (!this.options.taskStore) {
|
|
return undefined;
|
|
}
|
|
|
|
// FN-4689: close FN-4640 follow-up by wiring routine command sandbox execution through RunAuditor.
|
|
const engineRunContext: EngineRunContext = {
|
|
runId: generateSyntheticRunId("routine", routine.id),
|
|
agentId: routine.agentId ?? "routine-runner",
|
|
phase: "routine-execute",
|
|
source: "routine",
|
|
};
|
|
|
|
return createRunAuditor(this.options.taskStore, engineRunContext);
|
|
}
|
|
|
|
private async executeCommand(
|
|
routine: Routine,
|
|
command: string,
|
|
timeoutMs: number | undefined,
|
|
startedAt: string,
|
|
): Promise<AutomationRunResult> {
|
|
// Intercept the auto-backup command so it runs in-process via the engine's
|
|
// existing TaskStore instead of shelling out to a globally-installed
|
|
// fusion binary (which may be older than the running engine and re-create
|
|
// the nested `.fusion/.fusion/` directory). Mirrors cron-runner.ts.
|
|
if (isInProcessBackupCommand(command) && this.options.taskStore) {
|
|
try {
|
|
const { runBackupCommand, resolveGlobalBackupRoot } = await import("@fusion/core");
|
|
const fusionDir = this.options.taskStore.getFusionDir();
|
|
const settings = await this.options.taskStore.getSettings();
|
|
const result = await runBackupCommand(resolveGlobalBackupRoot(this.options.taskStore), settings);
|
|
const output = truncateOutput(result.output ?? "", "");
|
|
return {
|
|
success: result.success,
|
|
output,
|
|
error: result.success ? undefined : formatInProcessBackupError(output, fusionDir),
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
} catch (err) {
|
|
const message = formatInProcessBackupError(err, this.options.taskStore.getFusionDir());
|
|
return {
|
|
success: false,
|
|
output: "",
|
|
error: message,
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
}
|
|
|
|
if (isInProcessMemoryBackupCommand(command) && this.options.taskStore) {
|
|
try {
|
|
const { runMemoryBackupCommand } = await import("@fusion/core");
|
|
const fusionDir = this.options.taskStore.getFusionDir();
|
|
const settings = await this.options.taskStore.getSettings();
|
|
const result = await runMemoryBackupCommand(fusionDir, settings);
|
|
return {
|
|
success: result.success,
|
|
output: truncateOutput(result.output ?? "", ""),
|
|
error: result.success ? undefined : result.output,
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
} catch (err) {
|
|
const message = err instanceof Error ? err.message : String(err);
|
|
return {
|
|
success: false,
|
|
output: "",
|
|
error: message,
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
}
|
|
|
|
const auditor = this.getRoutineCommandAuditor(routine);
|
|
const backend: SandboxBackend = auditor
|
|
? resolveSandboxBackend({ auditor })
|
|
: resolveSandboxBackend();
|
|
await backend.prepare({ allowNetwork: true });
|
|
const result = await backend.run(command, {
|
|
cwd: this.options.rootDir,
|
|
timeoutMs: timeoutMs ?? DEFAULT_TIMEOUT_MS,
|
|
maxBuffer: MAX_BUFFER,
|
|
shell: defaultShell,
|
|
});
|
|
|
|
if (result.exitCode === 0 && !result.signal && !result.timedOut && !result.bufferExceeded && !result.spawnError) {
|
|
return {
|
|
success: true,
|
|
output: truncateOutput(result.stdout, result.stderr),
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
|
|
const error = result.timedOut
|
|
? `Command timed out after ${(timeoutMs ?? DEFAULT_TIMEOUT_MS) / 1000}s`
|
|
: result.spawnError?.message ?? "Command failed";
|
|
return {
|
|
success: false,
|
|
output: truncateOutput(result.stdout, result.stderr),
|
|
error,
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
|
|
private async executeSteps(routine: Routine, startedAt: string, liveCallbacks?: RoutineLiveRunCallbacks): Promise<AutomationRunResult> {
|
|
const steps = routine.steps ?? [];
|
|
const stepResults: AutomationStepResult[] = [];
|
|
let overallSuccess = true;
|
|
let stoppedEarly = false;
|
|
|
|
for (let i = 0; i < steps.length; i++) {
|
|
const step = steps[i];
|
|
liveCallbacks?.onStep?.({ stepIndex: i, stepId: step.id, stepName: step.name, stepType: step.type, status: "started" });
|
|
const result = await this.executeStep(routine, step, i, liveCallbacks);
|
|
liveCallbacks?.onStep?.({ stepIndex: i, stepId: step.id, stepName: step.name, stepType: step.type, status: "completed", success: result.success, error: result.error });
|
|
if (step.type !== "ai-prompt" && result.output) liveCallbacks?.onText?.(result.output);
|
|
stepResults.push(result);
|
|
|
|
if (!result.success) {
|
|
overallSuccess = false;
|
|
if (!step.continueOnFailure) {
|
|
stoppedEarly = true;
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
const outputParts: string[] = [];
|
|
for (const sr of stepResults) {
|
|
outputParts.push(`=== Step ${sr.stepIndex + 1}: ${sr.stepName} (${sr.success ? "success" : "FAILED"}) ===`);
|
|
if (sr.output) outputParts.push(sr.output);
|
|
if (sr.error) outputParts.push(`Error: ${sr.error}`);
|
|
}
|
|
const failedSteps = stepResults.filter((sr) => !sr.success);
|
|
|
|
return {
|
|
success: overallSuccess,
|
|
output: truncateOutput(outputParts.join("\n"), ""),
|
|
error: failedSteps.length > 0
|
|
? `${failedSteps.length} step(s) failed: ${failedSteps.map((s) => s.stepName).join(", ")}${stoppedEarly ? " (execution stopped)" : ""}`
|
|
: undefined,
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
stepResults,
|
|
};
|
|
}
|
|
|
|
private async executeStep(
|
|
routine: Routine,
|
|
step: AutomationStep,
|
|
stepIndex: number,
|
|
liveCallbacks?: RoutineLiveRunCallbacks,
|
|
): Promise<AutomationStepResult> {
|
|
const startedAt = new Date().toISOString();
|
|
const timeoutMs = step.timeoutMs ?? routine.timeoutMs ?? DEFAULT_TIMEOUT_MS;
|
|
|
|
if (step.type === "command") {
|
|
const result = await this.executeCommand(routine, step.command ?? "", timeoutMs, startedAt);
|
|
return {
|
|
stepId: step.id,
|
|
stepName: step.name,
|
|
stepIndex,
|
|
success: result.success,
|
|
output: result.output,
|
|
error: result.error,
|
|
startedAt,
|
|
completedAt: result.completedAt,
|
|
};
|
|
}
|
|
|
|
if (step.type === "ai-prompt") {
|
|
if (!step.prompt?.trim()) {
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: false, output: "", error: "AI prompt step has no prompt specified", startedAt, completedAt: new Date().toISOString() };
|
|
}
|
|
if (!this.options.aiPromptExecutor) {
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: false, output: "", error: "AI execution is not configured", startedAt, completedAt: new Date().toISOString() };
|
|
}
|
|
try {
|
|
/*
|
|
FNXC:Automations 2026-07-12-20:30:
|
|
Routine AI-prompt steps share the CronRunner AiPromptExecutor seam. Pass the persisted per-step thinking level before live callbacks so explicit reasoning effort applies and omitted/blank values inherit defaults.
|
|
*/
|
|
const output = await Promise.race([
|
|
this.options.aiPromptExecutor(step.prompt, step.modelProvider, step.modelId, step.allowedTools, step.thinkingLevel?.trim() || undefined, liveCallbacks),
|
|
new Promise<never>((_resolve, reject) => setTimeout(() => reject(new Error(`AI prompt step timed out after ${timeoutMs / 1000}s`)), timeoutMs)),
|
|
]);
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: true, output: truncateOutput(output, ""), startedAt, completedAt: new Date().toISOString() };
|
|
} catch (err) {
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: false, output: "", error: err instanceof Error ? err.message : String(err), startedAt, completedAt: new Date().toISOString() };
|
|
}
|
|
}
|
|
|
|
if (step.type === "create-task") {
|
|
if (!this.options.taskStore) {
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: false, output: "", error: "Task creation is not configured", startedAt, completedAt: new Date().toISOString() };
|
|
}
|
|
if (!step.taskDescription?.trim()) {
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: false, output: "", error: "Create-task step has no task description specified", startedAt, completedAt: new Date().toISOString() };
|
|
}
|
|
|
|
/*
|
|
FNXC:Automations 2026-07-12-20:30:
|
|
Routine create-task steps map their persisted thinking level onto the spawned task. Blank values remain undefined to keep the task's normal thinking-level inheritance.
|
|
*/
|
|
const taskInput: TaskCreateInput = {
|
|
title: step.taskTitle?.trim() || undefined,
|
|
description: step.taskDescription.trim(),
|
|
/*
|
|
FNXC:Automations 2026-07-30-16:40 (greptile #2652 — the UI fix was only half of it):
|
|
Was `(step.taskColumn as Column) || "triage"`. U11 deletes `triage` from the default workflow, so
|
|
a step with no explicit column created its task into a column the board does not declare — and it
|
|
did so for EVERY routine, including ones saved through the fixed editor, because this substitution
|
|
happens after the step is read. Fixing the form's default alone changed nothing at runtime.
|
|
|
|
Omitted instead of defaulted: `createTask` resolves the workflow's own intake column when no
|
|
column is given (#2589), which is the only answer correct for every board, custom workflows
|
|
included. An explicit column on the step is still honoured.
|
|
*/
|
|
column: step.taskColumn ? (step.taskColumn as Column) : undefined,
|
|
modelProvider: step.modelProvider?.trim() || undefined,
|
|
modelId: step.modelId?.trim() || undefined,
|
|
thinkingLevel: (step.thinkingLevel?.trim() || undefined) as TaskCreateInput["thinkingLevel"],
|
|
source: {
|
|
sourceType: "automation",
|
|
sourceMetadata: { routineId: routine.id, stepId: step.id },
|
|
},
|
|
};
|
|
try {
|
|
const task = await this.options.taskStore.createTask(taskInput);
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: true, output: `Created task ${task.id}: ${task.title || task.description.slice(0, 80)}`, startedAt, completedAt: new Date().toISOString() };
|
|
} catch (err) {
|
|
return { stepId: step.id, stepName: step.name, stepIndex, success: false, output: "", error: err instanceof Error ? err.message : String(err), startedAt, completedAt: new Date().toISOString() };
|
|
}
|
|
}
|
|
|
|
return {
|
|
stepId: step.id,
|
|
stepName: step.name,
|
|
stepIndex,
|
|
success: false,
|
|
output: "",
|
|
error: `Unknown step type: "${String((step as unknown as Record<string, unknown>).type)}"`,
|
|
startedAt,
|
|
completedAt: new Date().toISOString(),
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Handle catch-up for missed routine executions based on the catch-up policy.
|
|
*
|
|
* @param routine - The routine to check for catch-up
|
|
*/
|
|
async handleCatchUp(routine: Routine): Promise<void> {
|
|
const catchUpPolicy = routine.catchUpPolicy ?? "skip";
|
|
|
|
if (catchUpPolicy === "skip") {
|
|
return;
|
|
}
|
|
|
|
// "run_one" or "run" policy - need to catch up
|
|
if (!routine.lastRunAt) {
|
|
// Never run before — nothing to catch up
|
|
return;
|
|
}
|
|
|
|
// Calculate missed intervals
|
|
if (!routine.cronExpression) {
|
|
return;
|
|
}
|
|
|
|
try {
|
|
const cronExpr = CronExpressionParser.parse(routine.cronExpression, {
|
|
currentDate: new Date(routine.lastRunAt ?? Date.now()),
|
|
});
|
|
const lastRun = new Date(routine.lastRunAt ?? Date.now());
|
|
const now = new Date();
|
|
|
|
const missedIntervals: Date[] = [];
|
|
|
|
// Get next interval after lastRun, then iterate
|
|
let intervalDate = new Date(cronExpr.next().toISOString() ?? Date.now());
|
|
|
|
while (intervalDate.getTime() <= now.getTime() && missedIntervals.length < MAX_CATCH_UP_INTERVALS) {
|
|
if (intervalDate.getTime() > lastRun.getTime()) {
|
|
missedIntervals.push(new Date(intervalDate));
|
|
}
|
|
const nextIso = cronExpr.next().toISOString();
|
|
if (!nextIso) break;
|
|
intervalDate = new Date(nextIso);
|
|
}
|
|
|
|
if (missedIntervals.length === 0) {
|
|
return;
|
|
}
|
|
|
|
log.log(`[${routine.id}] Running ${missedIntervals.length} catch-up executions`);
|
|
|
|
// Execute each missed interval
|
|
for (const missedInterval of missedIntervals) {
|
|
try {
|
|
await this.executeRoutine(routine.id, "cron", {
|
|
catchUp: true,
|
|
missedInterval: missedInterval.toISOString(),
|
|
});
|
|
} catch (err) {
|
|
log.error(`[${routine.id}] Catch-up execution failed: ${err}`);
|
|
}
|
|
}
|
|
} catch (err) {
|
|
log.error(`[${routine.id}] Error calculating catch-up intervals: ${err}`);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Trigger a routine manually (via API).
|
|
*
|
|
* @param routineId - The ID of the routine to trigger
|
|
* @returns The execution result
|
|
* @throws Error if routine not found or disabled
|
|
*/
|
|
async triggerManual(routineId: string, liveCallbacks?: RoutineLiveRunCallbacks): Promise<RoutineExecutionResult> {
|
|
const routine = await this.options.routineStore.getRoutine(routineId);
|
|
|
|
if (!routine.enabled) {
|
|
throw new Error(`Routine '${routineId}' is disabled`);
|
|
}
|
|
|
|
return this.executeRoutine(routineId, "api", undefined, liveCallbacks);
|
|
}
|
|
|
|
/**
|
|
* Trigger a routine via webhook.
|
|
*
|
|
* @param routineId - The ID of the routine to trigger
|
|
* @param payload - The webhook payload
|
|
* @param _signature - The webhook signature (verified by RoutineScheduler)
|
|
* @returns The execution result
|
|
* @throws Error if routine not found, not a webhook trigger, or disabled
|
|
*/
|
|
async triggerWebhook(
|
|
routineId: string,
|
|
payload: Record<string, unknown>,
|
|
_signature?: string
|
|
): Promise<RoutineExecutionResult> {
|
|
const routine = await this.options.routineStore.getRoutine(routineId);
|
|
|
|
if (routine.trigger.type !== "webhook") {
|
|
throw new Error(
|
|
`Routine '${routineId}' does not have webhook trigger type`
|
|
);
|
|
}
|
|
|
|
if (!routine.enabled) {
|
|
throw new Error(`Routine '${routineId}' is disabled`);
|
|
}
|
|
|
|
return this.executeRoutine(routineId, "webhook", { webhookPayload: payload });
|
|
}
|
|
|
|
/**
|
|
* Get the number of currently-running executions.
|
|
*/
|
|
getInFlightCount(): number {
|
|
return this.inFlightExecutions.size;
|
|
}
|
|
|
|
/**
|
|
* Check if a routine is currently being executed.
|
|
*/
|
|
isRoutineRunning(routineId: string): boolean {
|
|
return this.inFlightExecutions.has(routineId);
|
|
}
|
|
}
|
|
|
|
/*
|
|
FNXC:DatabaseBackup 2026-06-26-12:00:
|
|
Routine-runner in-process backups persist AutomationRunResult.error directly to lastRunResult. Normalize empty or opaque failures here so Database Backup cards always show a DB-qualified cause.
|
|
*/
|
|
function formatInProcessBackupError(err: unknown, fusionDir: string): string {
|
|
const message = err instanceof Error ? err.message.trim() : String(err ?? "").trim();
|
|
const cause = message || "unknown error";
|
|
if (cause.includes("project DB") || cause.includes("central DB")) {
|
|
return cause;
|
|
}
|
|
return `project PostgreSQL run backup command failed; project state: ${fusionDir}; cause: ${cause}`;
|
|
}
|
|
|
|
function truncateOutput(stdout: string, stderr: string): string {
|
|
let output = stdout;
|
|
if (stderr) {
|
|
output += stdout ? "\n--- stderr ---\n" : "";
|
|
output += stderr;
|
|
}
|
|
if (output.length > MAX_OUTPUT_LENGTH) {
|
|
return `${output.slice(0, MAX_OUTPUT_LENGTH)}\n[output truncated]`;
|
|
}
|
|
return output;
|
|
}
|