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 @@ -221,6 +221,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 @@ -238,6 +245,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
17 changes: 12 additions & 5 deletions apps/server/src/engine/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -705,18 +705,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
66 changes: 66 additions & 0 deletions tests/model-worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,72 @@ test("the model worker keeps the text a model replies with when it calls no tool
}
});

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",
intelligenceApiKey: "test-project-key-never-sent",
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 });
}
});

test("browser reads keep observation identity distinct while reusing one session", async (t) => {
let currentUrl = "https://example.com/one";
const sessionIds = new Set<string>();
Expand Down
Loading