Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions packages/coding-agent/.changes/factory-dag-executor.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- Added the state-machine factory executor: `rlm.factory.run/status/stop/resume` canonicalizes the stored factory (dag sugar compiles to machine form), enters entry states, and drives the machine to quiescence — every settle evaluates its guarded transitions (fan-out legal, max_entries blocking recorded as transition_blocked, re-entry re-binds inputs), with bounded re-entry, per-entry foreach/budgets/retries, fail_fast/continue/escalate policies, cancellation cascades, and one quiet notice per milestone (finished/failed/paused/budget_exceeded/max_transitions_exceeded). `status()` reports per-state entries_used/max_entries and a transitions_fired usage count over a 200-event window. Wait states stay specified but gated: the watch host handlers (`rlm.watch.*`) arrive with the communication series, so the validator rejects wait blocks until then.
4 changes: 3 additions & 1 deletion packages/coding-agent/src/core/agent-messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type { HostRequestHandler } from "./kernel/index.js";
import type { CustomMessage } from "./messages.js";
import {
ASYNC_BASH_COMPLETION_CUSTOM_TYPE,
FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE,
HEARTBEAT_PROMPT_CUSTOM_TYPE,
sanitizeMessageHeaderValue,
} from "./messages.js";
Expand Down Expand Up @@ -447,7 +448,8 @@ export function startsAgentRun(message: AgentMessage): boolean {
isAgentSessionMessage(message) ||
(message.role === "custom" &&
(message.customType === HEARTBEAT_PROMPT_CUSTOM_TYPE ||
message.customType === ASYNC_BASH_COMPLETION_CUSTOM_TYPE))
message.customType === ASYNC_BASH_COMPLETION_CUSTOM_TYPE ||
message.customType === FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE))
);
}

Expand Down
31 changes: 31 additions & 0 deletions packages/coding-agent/src/core/agent-session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,7 @@ import {
convertToLlm,
createAsyncBashCompletionMessage,
createCompactionOutcomeMessage,
createFactoryProgressMessage,
createHarnessDigestMessage,
createHeartbeatPromptMessage,
createRefinementNoticeMessage,
Expand All @@ -204,6 +205,8 @@ import {
createRlmChildTerminalNoticeMessage,
createSessionSlashCommandMessage,
createSessionSlashCommandResultMessage,
FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE,
FACTORY_PROGRESS_PREVIEW_LABEL,
HARNESS_DIGEST_CUSTOM_TYPE,
type HarnessDigestDetails,
HEARTBEAT_PROMPT_CUSTOM_TYPE,
Expand Down Expand Up @@ -268,6 +271,7 @@ import {
createAsyncBashCompletionHostHandler,
createAsyncBashConsumedHostHandler,
createDefaultRlmSubagentSessionName,
createFactoryProgressHostHandler,
createRlmCollectHostHandler,
createRlmCreateSessionHostHandler,
createRlmDeleteSubagentHostHandler,
Expand Down Expand Up @@ -956,6 +960,8 @@ function injectedMessagePreviewLabel(message: CustomMessage): string | undefined
return HEARTBEAT_PROMPT_PREVIEW_LABEL;
case ASYNC_BASH_COMPLETION_CUSTOM_TYPE:
return ASYNC_BASH_COMPLETION_PREVIEW_LABEL;
case FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE:
return FACTORY_PROGRESS_PREVIEW_LABEL;
case GOAL_CONTEXT_CUSTOM_TYPE:
return GOAL_CONTEXT_PREVIEW_LABEL;
default:
Expand Down Expand Up @@ -10981,6 +10987,31 @@ export class AgentSession {
),
"rlm.progress.note": createRlmProgressNoteHostHandler((message) => this.noteRlmProgress(message)),
"rlm.delete_subagent": createRlmDeleteSubagentHostHandler((target) => this.deleteRlmSubagent(target)),
"factory.progress": createFactoryProgressHostHandler(async (details) => {
const message = createFactoryProgressMessage(details);
const disposeSignal = this._sessionActionCommitDisposeAbortController.signal;
while (true) {
let admissionCommitted = false;
try {
await this._promptInjectedMessage(message.content, message, {
streamingBehavior: "steer",
queueIfBusy: true,
resumeIfIdle: true,
returnAfterAccepted: true,
suppressAutonomousContinuation: true,
admissionCommitted: () => {
admissionCommitted = true;
},
});
return;
} catch (error) {
if (admissionCommitted || !(error instanceof SessionInputAdmissionPausedError)) throw error;
while (this._sessionInputAdmissionPauses.size > 0 && !disposeSignal.aborted) {
await this._waitForSessionActivityChange(disposeSignal);
}
}
}
}),
"model.info": async () => ({
id: this.model?.id ?? null,
provider: this.model?.provider ?? null,
Expand Down
30 changes: 30 additions & 0 deletions packages/coding-agent/src/core/messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ export const RLM_CHILD_FAILURE_CUSTOM_TYPE = "rlm_child_failure";
export const RLM_CHILD_TERMINAL_NOTICE_CUSTOM_TYPE = "rlm_child_terminal_notice";
export const ASYNC_BASH_COMPLETION_CUSTOM_TYPE = "async_bash_completion";
export const ASYNC_BASH_COMPLETION_PREVIEW_LABEL = "Background command finished";
export const FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE = "factory_progress_notice";
export const FACTORY_PROGRESS_PREVIEW_LABEL = "Factory progress";

/**
* Names and other metadata interpolated into a `[<kind> ...]` header line must not
Expand Down Expand Up @@ -267,6 +269,34 @@ Command: ${JSON.stringify(details.command)}`,
};
}

export interface FactoryProgressDetails {
runId: string;
kind: "finished" | "failed" | "paused" | "budget_exceeded" | "max_transitions_exceeded";
node?: string;
detail: string;
}

interface FactoryProgressMessage extends CustomMessage<FactoryProgressDetails> {
customType: typeof FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE;
content: string;
}

export function createFactoryProgressMessage(
details: FactoryProgressDetails,
timestamp = Date.now(),
): FactoryProgressMessage {
// The wire kind budget_exceeded renders as the friendlier budget-exceeded label.
const kind = details.kind === "budget_exceeded" ? "budget-exceeded" : details.kind;
return {
role: "custom",
customType: FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE,
content: `[factory-progress run:${sanitizeMessageHeaderValue(details.runId)}] ${kind}: ${details.detail}`,
display: true,
details,
timestamp,
};
}

export function createRlmChildFailureMessage(
details: RlmChildFailureDetails,
timestamp = Date.now(),
Expand Down
2 changes: 1 addition & 1 deletion packages/coding-agent/src/core/prompts/rlm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ const REPL_CONTROL_PROMPT = [
"",
"Continual harness state is available as `rlm.harness` and `rlm.get_harness_state()`. CRUD calls are local to this Prime Agent session by default: `rlm.harness.create_memory(...)`, `rlm.harness.update_memory(...)`, `rlm.harness.delete_memory(...)`, `rlm.harness.create_skill(...)`, `rlm.harness.update_skill(...)`, `rlm.harness.delete_skill(...)`, `rlm.harness.create_subagent(...)`, `rlm.harness.update_subagent(...)`, `rlm.harness.delete_subagent(...)`, `rlm.harness.create_factory(...)`, `rlm.harness.update_factory(...)`, `rlm.harness.delete_factory(...)`, `rlm.harness.create_prompt_note(...)`, `rlm.harness.update_prompt_note(...)`, `rlm.harness.delete_prompt_note(...)`, plus `rlm.harness.record_refinement(...)` and `rlm.harness.overview()`. Use `global_=True` only for stable cross-session lessons; Python reserves `global`, so literal `global=True` is invalid syntax.",
"",
"Factory entries declare validated state-machine workflows of subagent states in arguments['machine'] (DAG specs in arguments['dag'] compile to machine form); execution lands in a follow-up PR.",
"Factory entries declare validated state-machine workflows of subagent states in arguments['machine'] (DAG specs in arguments['dag'] compile to machine form): run a stored factory with await rlm.factory.run('<id>'), watch it with rlm.factory.status(run_id), stop it with rlm.factory.stop(run_id), and resume an escalate-paused run with rlm.factory.resume(run_id).",
"",
"Terminology: continual harness names the persisted prompt, memory, skill, and subagent layer; RLM names the runtime, Python REPL kernel, and native call interface exposed to the model.",
"",
Expand Down
13 changes: 9 additions & 4 deletions packages/coding-agent/src/core/refinement/refinement.ts
Original file line number Diff line number Diff line change
Expand Up @@ -718,7 +718,7 @@ export function formatHarnessStateForPrompt(
);
} else if (kind === "factory" && entries.length > 0 && includeIpythonExamples) {
lines.push(
`${kind}: ${entries.length} (state-machine workflow specs; run one with \`await rlm.factory.run('<id>')\`; execution lands in a follow-up PR)`,
`${kind}: ${entries.length} (state-machine workflow specs; run one with \`await rlm.factory.run('<id>')\`; watch with \`rlm.factory.status(run_id)\`, stop with \`rlm.factory.stop(run_id)\`)`,
);
} else {
lines.push(`${kind}: ${entries.length}`);
Expand Down Expand Up @@ -1067,10 +1067,15 @@ function validateEdit(edit: RefinementEdit, computedId?: string): string | undef
}
if (edit.action !== "delete" && edit.kind === "factory") {
// Structural check only: the kernel validator (rlm.factory) enforces the full
// DAG semantics at write time; do not reimplement it here.
// machine semantics at write time; do not reimplement it here.
const dag = edit.arguments?.dag;
if (typeof dag !== "object" || dag === null || Array.isArray(dag)) {
return "factory entry requires a dag object in arguments";
const machine = edit.arguments?.machine;
if (dag !== undefined && machine !== undefined) {
return "pass either dag or machine form, not both";
}
const spec = machine ?? dag;
if (typeof spec !== "object" || spec === null || Array.isArray(spec)) {
return "factory entry requires a dag or machine object in arguments";
}
}
return undefined;
Expand Down
57 changes: 57 additions & 0 deletions packages/coding-agent/src/core/rlm-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,30 @@ interface AsyncBashConsumedRequest {
}

type AsyncBashConsumedHandler = (request: AsyncBashConsumedRequest) => void | Promise<void>;

export type FactoryProgressKind = "finished" | "failed" | "paused" | "budget_exceeded" | "max_transitions_exceeded";

export interface FactoryProgressRequest {
runId: string;
kind: FactoryProgressKind;
node?: string;
detail: string;
}

export type FactoryProgressHandler = (request: FactoryProgressRequest) => void | Promise<void>;

const FACTORY_PROGRESS_KINDS: readonly FactoryProgressKind[] = [
"finished",
"failed",
"paused",
"budget_exceeded",
"max_transitions_exceeded",
];

function isFactoryProgressKind(value: unknown): value is FactoryProgressKind {
return typeof value === "string" && (FACTORY_PROGRESS_KINDS as readonly string[]).includes(value);
}

export type RlmListSubagentsHandler = () => RlmListSubagentsResult | Promise<RlmListSubagentsResult>;
export type RlmDeleteSubagentHandler = (target: string) => Promise<RlmDeleteSubagentResult>;
export type RlmFindModelsHandler = (query: string, limit: number) => RlmFindModelsResult | Promise<RlmFindModelsResult>;
Expand Down Expand Up @@ -331,6 +355,39 @@ export function createAsyncBashCompletionHostHandler(handler: AsyncBashCompletio
};
}

/** Adapt a factory executor milestone into a validated `factory.progress` host notification. */
export function createFactoryProgressHostHandler(handler: FactoryProgressHandler): HostRequestHandler {
return async (payload) => {
const runId = payload.run_id;
if (typeof runId !== "string" || !runId.trim()) {
throw new Error("factory.progress run_id must be a non-empty string");
}
const kind = payload.kind;
if (!isFactoryProgressKind(kind)) {
throw new Error(
`factory.progress kind must be one of ${FACTORY_PROGRESS_KINDS.join(", ")}, got ${JSON.stringify(kind)}`,
);
}
const detail = payload.detail;
if (typeof detail !== "string" || !detail.trim()) {
throw new Error("factory.progress detail must be a non-empty string");
}
const node = payload.node;
if (node !== undefined && typeof node !== "string") {
throw new Error("factory.progress node must be a string when provided");
}
if (node !== undefined && !node.trim()) {
throw new Error("factory.progress node must be a non-empty string when provided");
}
const request: FactoryProgressRequest = { runId: runId.trim(), kind, detail };
if (typeof node === "string" && node.trim()) {
request.node = node.trim();
}
await handler(request);
return {};
};
}

/** The kernel read a finished command's result, so its completion notice is stale. */
export function createAsyncBashConsumedHostHandler(handler: AsyncBashConsumedHandler): HostRequestHandler {
return async (payload) => {
Expand Down
93 changes: 93 additions & 0 deletions packages/coding-agent/test/factory-executor.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
import { describe, expect, it, vi } from "vitest";
import { startsAgentRun } from "../src/core/agent-messages.js";
import {
convertToLlm,
createFactoryProgressMessage,
FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE,
} from "../src/core/messages.js";
import { createFactoryProgressHostHandler } from "../src/core/rlm-runtime.js";

describe("factory progress", () => {
it("creates a bracket-grammar milestone notice and keeps it model-visible", () => {
const message = createFactoryProgressMessage({
runId: "run-1",
kind: "budget_exceeded",
detail: "run budget_ms 600000 exceeded after 612000ms",
});

expect(message.customType).toBe(FACTORY_PROGRESS_NOTICE_CUSTOM_TYPE);
expect(message.content).toBe(
"[factory-progress run:run-1] budget-exceeded: run budget_ms 600000 exceeded after 612000ms",
);
expect(convertToLlm([message])).toEqual([
{
role: "user",
content: [{ type: "text", text: message.content }],
timestamp: message.timestamp,
},
]);
});

it("renders the max_transitions pause kind verbatim (distinct from budget_exceeded)", () => {
const message = createFactoryProgressMessage({
runId: "run-2",
kind: "max_transitions_exceeded",
detail: "max_transitions 1 exceeded; no new entries",
});
expect(message.content).toBe(
"[factory-progress run:run-2] max_transitions_exceeded: max_transitions 1 exceeded; no new entries",
);
});

it("renders the finished, failed, and paused milestone kinds verbatim", () => {
expect(createFactoryProgressMessage({ runId: "r", kind: "finished", detail: "all 3 nodes done" }).content).toBe(
"[factory-progress run:r] finished: all 3 nodes done",
);
expect(createFactoryProgressMessage({ runId: "r", kind: "failed", detail: "node review failed" }).content).toBe(
"[factory-progress run:r] failed: node review failed",
);
expect(
createFactoryProgressMessage({ runId: "r", kind: "paused", detail: "escalate: awaiting resume" }).content,
).toBe("[factory-progress run:r] paused: escalate: awaiting resume");
});

it("starts a new agent run for a factory milestone follow-up", () => {
const message = createFactoryProgressMessage({ runId: "r", kind: "finished", detail: "done" });
expect(startsAgentRun(message)).toBe(true);
});

it("validates and forwards kernel milestone payloads", async () => {
const milestone = vi.fn();
const handler = createFactoryProgressHostHandler(milestone);
const payload = { run_id: " run-1 ", kind: "paused", node: "review", detail: "node review failed" };

await expect(handler(payload)).resolves.toEqual({});
expect(milestone).toHaveBeenCalledWith({
runId: "run-1",
kind: "paused",
node: "review",
detail: "node review failed",
});
});

it("omits the node field when the kernel does not provide one", async () => {
const milestone = vi.fn();
const handler = createFactoryProgressHostHandler(milestone);

await handler({ run_id: "run-1", kind: "finished", detail: "all nodes done" });
expect(milestone).toHaveBeenCalledWith({ runId: "run-1", kind: "finished", detail: "all nodes done" });
});

it.each([
[{ kind: "finished", detail: "done" }, "run_id must be a non-empty string"],
[{ run_id: "", kind: "finished", detail: "done" }, "run_id must be a non-empty string"],
[{ run_id: "r", kind: "weird", detail: "done" }, "kind must be one of"],
[{ run_id: "r", kind: "finished" }, "detail must be a non-empty string"],
[{ run_id: "r", kind: "finished", detail: " " }, "detail must be a non-empty string"],
[{ run_id: "r", kind: "finished", detail: "done", node: 5 }, "node must be a string when provided"],
[{ run_id: "r", kind: "finished", detail: "done", node: "" }, "node must be a non-empty string"],
])("rejects an invalid payload %#", async (payload, error) => {
const handler = createFactoryProgressHostHandler(() => undefined);
await expect(handler(payload)).rejects.toThrow(error);
});
});
Loading
Loading