Skip to content
Merged
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
14 changes: 14 additions & 0 deletions apps/server/src/engine/model.ts
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,13 @@ export async function executeModelTask(
async (data) => {
const key = createHash("sha256").update(JSON.stringify(data)).digest("hex");
const action = await service.prepare(owner, task, { kind: "email.send", data }, key, ctx);
if (action.status === "succeeded") {
task = await ctx.checkpoint({
state: { ...task.state, approvalResult: action.result },
actionId: null,
});
return { status: "succeeded", actionId: action.id, result: action.result };
}
outcome = { status: "waiting_approval", actionId: action.id };
return { status: "waiting_approval", actionId: action.id };
},
Expand All @@ -245,6 +252,13 @@ export async function executeModelTask(
key,
ctx,
);
if (action.status === "succeeded") {
task = await ctx.checkpoint({
state: { ...task.state, approvalResult: action.result },
actionId: null,
});
return { status: "succeeded", actionId: action.id, result: action.result };
}
outcome = { status: "waiting_approval", actionId: action.id };
return { status: "waiting_approval", actionId: action.id };
},
Expand Down
142 changes: 119 additions & 23 deletions apps/server/src/engine/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -272,7 +272,7 @@ export class AgentService {
? "Paused. Resume when you're ready."
: "",
...(task.kind === "monitor" && action === "resume"
? { state: { ...task.state, failures: 0, notice: null } }
? { state: { ...task.state, failures: 0, notice: null, resumingMonitor: false } }
: {}),
},
);
Expand All @@ -287,6 +287,8 @@ export class AgentService {
{
status: action === "cancel" ? "stopped" : action === "pause" ? "paused" : "active",
nextCheckAt: date(),
// Clearing the error fences out a failure reconcile that read the task before this.
...(action === "resume" || action === "retry" ? { error: null } : {}),
},
);
if (action === "cancel" && task.actionId) {
Expand Down Expand Up @@ -396,10 +398,26 @@ export class AgentService {
return monitor;
}
private async activateMonitor(owner: string, monitor: Monitor) {
if (monitor.status !== "active") return;
if (monitor.status !== "active") return null;
const task = await this.getTask(owner, monitor.taskId);
if (task.status !== "paused" || !task.state.initializingMonitor) return;
await this.db.compareAndSwap(
if (task.status !== "paused") return null;
if (task.state.resumingMonitor)
return this.db.compareAndSwap(
owner,
"tasks",
task.id,
{ status: "paused", state: { resumingMonitor: true } },
{
status: "queued",
nextRunAt: date(),
leaseId: null,
leaseUntil: null,
error: null,
state: { ...task.state, resumingMonitor: false, failures: 0, notice: null },
},
);
if (!task.state.initializingMonitor) return null;
return this.db.compareAndSwap(
owner,
"tasks",
task.id,
Expand All @@ -415,18 +433,79 @@ export class AgentService {
if (!monitor) throw new AppError("Monitor not found", 404);
if (monitor.status === "stopped" && action !== "stop")
throw new AppError("Create a new watch to restart this stopped monitor", 409);
const status = action === "pause" ? "paused" : action === "stop" ? "stopped" : "active";
const saved = await this.db.put(owner, "monitors", { ...monitor, status, nextCheckAt: date() });
const task = await this.getTask(owner, monitor.taskId);
if (action === "pause" || action === "stop")
if (action === "pause" || action === "stop") {
const status = action === "pause" ? "paused" : "stopped";
const saved = await this.db.put(owner, "monitors", {
...monitor,
status,
nextCheckAt: date(),
});
const task = await this.getTask(owner, monitor.taskId);
await this.control(owner, task.id, action === "pause" ? "pause" : "cancel");
else {
return saved;
}
let monitorStatus = monitor.status;
for (let attempt = 0; attempt < 2; attempt++) {
const task = await this.getTask(owner, monitor.taskId);
if (task.status === "cancelled") break;
if (task.status === "paused") {
// Mark the paused task before activating the monitor so no worker can claim it in between.
const marked = await this.db.compareAndSwap(
owner,
"tasks",
task.id,
{ status: "paused", leaseId: task.leaseId ?? null },
{ state: { ...task.state, resumingMonitor: true } },
);
if (!marked) break;
const activated = await this.db.compareAndSwap<Monitor>(
owner,
"monitors",
id,
{ status: monitor.status },
{ status: "active", nextCheckAt: date(), error: null },
);
const saved = activated ?? (await this.db.get<Monitor>(owner, "monitors", id));
if (saved?.status === "active" && (await this.activateMonitor(owner, saved))) return saved;
const unmarked = await this.db.compareAndSwap(
owner,
"tasks",
task.id,
{ status: "paused", state: { resumingMonitor: true } },
{ state: { ...task.state, resumingMonitor: false } },
);
// Another request or maintenance may have finished this resume first.
if (!unmarked && saved?.status === "active") {
const latest = await this.getTask(owner, task.id);
if (["queued", "running", "scheduled"].includes(latest.status)) return saved;
}
if (activated)
await this.db.compareAndSwap(
owner,
"monitors",
id,
{ status: "active" },
{ status: monitor.status, error: monitor.error ?? null },
);
break;
}
// Only activate the monitor we read, so a concurrent stop is never undone.
const saved = await this.db.compareAndSwap<Monitor>(
owner,
"monitors",
id,
{ status: monitorStatus },
{ status: "active", nextCheckAt: date() },
);
if (!saved) break;
monitorStatus = "active";
this.worker.abort(task.id);
await this.db.compareAndSwap(
const queued = await this.db.compareAndSwap(
owner,
"tasks",
task.id,
{ status: task.status, leaseId: task.leaseId ?? null },
// updatedAt fences out a whole run finishing in between, which would move the baseline.
{ status: task.status, leaseId: task.leaseId ?? null, updatedAt: task.updatedAt },
{
status: "queued",
nextRunAt: date(),
Expand All @@ -436,8 +515,9 @@ export class AgentService {
state: { ...task.state, failures: 0, notice: null },
},
);
if (queued) return saved;
}
return saved;
throw new AppError("The watch changed while updating. Try again.", 409);
}
async refreshIdeas(owner: string) {
const w = await this.workspace.snapshot(owner);
Expand Down Expand Up @@ -626,18 +706,25 @@ export class AgentService {
409,
);
const proposal = await this.actions.propose(owner, input, `${task.id}:${key}`, task.id);
if (proposal.status === "succeeded") return proposal;
if (proposal.status !== "awaiting_review" && proposal.status !== "executing")
throw new AppError(
`Reviewed action ${proposal.status}: ${proposal.error ?? "No further action was taken"}`,
409,
);
try {
await context.checkpoint({ actionId: proposal.id });
} catch (error) {
if (proposal.status === "awaiting_review")
await this.actions.decide(owner, proposal.id, proposal.hash, "deny");
throw error;
}
await context.event(
"approval",
proposal.title,
`Review prepared for ${proposal.account ?? "the connected account"}`,
);
if (proposal.status === "awaiting_review")
await context.event(
"approval",
proposal.title,
`Review prepared for ${proposal.account ?? "the connected account"}`,
);
return proposal;
}
private async execute(
Expand Down Expand Up @@ -684,11 +771,7 @@ export class AgentService {
"monitors",
String(task.input.monitorId),
{ status: "active" },
{
error: detail,
nextCheckAt,
...(failures >= 5 ? { status: "paused" } : {}),
},
{ error: detail, nextCheckAt },
);
await context.event(
"error",
Expand All @@ -702,6 +785,7 @@ export class AgentService {
state: {
...task.state,
failures,
resumingMonitor: false,
notice: {
title: "Watch needs attention",
body: detail,
Expand Down Expand Up @@ -805,6 +889,15 @@ export class AgentService {
.safeParse(task.state.notice);
if ((task.status === "scheduled" || (task.status === "paused" && task.error)) && notice.success)
await this.notify(owner, notice.data.title, notice.data.body, task.id, notice.data.key);
// A watch pauses after repeated failures only once that task outcome has committed.
if (task.kind === "monitor" && task.status === "paused" && task.error)
await this.db.compareAndSwap(
owner,
"monitors",
String(task.input.monitorId),
{ status: "active", error: task.error },
{ status: "paused" },
);
}
private async document(
owner: string,
Expand Down Expand Up @@ -926,7 +1019,8 @@ export class AgentService {
}
const text = observation.text.replace(/\s+/g, " ").trim();
const currentHash = hash(text);
const previousHash = monitor.lastHash;
const previousHash =
typeof task.state.lastHash === "string" ? task.state.lastHash : monitor.lastHash;
const matched =
monitor.condition === "change"
? Boolean(previousHash && previousHash !== currentHash)
Expand Down Expand Up @@ -971,6 +1065,8 @@ export class AgentService {
state: {
...task.state,
sessionId: observation.sessionId,
lastHash: currentHash,
resumingMonitor: false,
matched,
failures: 0,
notice: shouldNotify
Expand Down
65 changes: 65 additions & 0 deletions tests/model-worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -114,3 +114,68 @@ test("CopilotKit model worker executes server tools and persists the confirmed o
await rm(directory, { recursive: true, force: true });
}
});

test("replaying a completed prepared action returns its receipt without reopening approval", async (t) => {
const directory = await mkdtemp(join(tmpdir(), "openmuse-model-replay-"));
const db = await createStore();
const draft = {
title: "Sample walk",
start: "2026-10-10T10:00:00-07:00",
end: "2026-10-10T11:00:00-07:00",
};
let calls: ({ name: string; arguments: object } | undefined)[] = [
{ name: "prepare_event", arguments: draft },
];
const { requests } = await modelFixture(t, (index) => calls[index]);
const server = await createApp(db, {
mode: "sample",
port: 8787,
host: "127.0.0.1",
publicUrl: "http://localhost:8787",
dataDir: directory,
agentBackend: "model",
model: "openai/fixture",
googleRedirectUri: "http://localhost:8787/api/google/callback",
allowedOrigins: [],
});
try {
const task = await server.agent.createTask("replay-owner", {
prompt: "Put a sample walk on my calendar",
});
await server.agent.worker.tick();
const pending = await server.agent.getTask("replay-owner", task.id);
assert.equal(pending.status, "waiting_approval");
assert.ok(pending.actionId);
const proposal = await db.get<ActionProposal>("replay-owner", "actions", pending.actionId);
assert.ok(proposal);
const completed = await server.actions.decide(
"replay-owner",
proposal.id,
proposal.hash,
"approve",
);
assert.equal(completed.status, "succeeded");

requests.length = 0;
calls = [
{ name: "prepare_event", arguments: draft },
{ name: "finish_task", arguments: { summary: "The reviewed event is already complete." } },
];
await server.agent.worker.tick();

const finished = await server.agent.getTask("replay-owner", task.id);
assert.equal(finished.status, "succeeded", finished.error ?? finished.question);
assert.equal(finished.actionId, null);
assert.equal(finished.state.approvalResult, completed.result);
const actions = (await db.list<ActionProposal>("replay-owner", "actions")).filter(
(action) => action.taskId === task.id,
);
assert.equal(actions.length, 1);
assert.equal(actions[0].status, "succeeded");
assert.ok(requests.some((request) => request.body.includes(String(completed.result))));
} finally {
await server.agent.stop();
await db.close();
await rm(directory, { recursive: true, force: true });
}
});
Loading