Scheduled and loop agents get workspaces; schedules UI speaks the user's timezone (#1909)

* fix(app): schedules UI speaks the user's timezone

* refactor(server): require explicit agent placement and bind create-agent command

* fix(server): scheduled agents get a workspace and show in the sidebar

* fix(server): loop agents get a workspace and show in the sidebar

* fix(server): preserve scheduled and loop run semantics through the create command

* fix(server): serialize all schedule record writes through the per-schedule queue

* fix(server): tighten create-agent dependency contracts

* fix(server): handle schedule restamps and loop cancellations

* fix(server): preserve scheduled and loop create semantics

* refactor(server): schedule store owns atomic mutations

* fix(server): keep schedule stamp reuse deterministic

* fix(server): handle loop workspace stamp rejection

* refactor(server): schedule store owns typed identity upsert

* fix(server): keep schedule workspace stamps server-owned

* fix(server): preserve scheduled run semantics

* refactor: required nullable workspaceId replaces placement union; pure cadence policy

* fix(server): tighten scheduled and loop agent edges

* fix(server): fail canceled scheduled runs

* fix(server): close schedule and loop cwd edges

* fix(server): refresh schedule workspace stamps

* fix(server): preserve schedule timezone on replace

* fix(server): close loop stop and cadence reset edges
This commit is contained in:
Mohamed Boudra
2026-07-07 11:13:28 +02:00
committed by GitHub
parent 783f43c580
commit d3e1fe48f6
29 changed files with 4814 additions and 1067 deletions

View File

@@ -4,6 +4,10 @@ Paseo uses **file-based JSON persistence** instead of a traditional database. Al
All server-side stores live under `$PASEO_HOME` (defaults to `~/.paseo`).
## Store Surface Rules
Store APIs own persistence atomicity and should not make services coordinate raw reads and writes. A good store method maps cleanly to one SQL statement or one SQL transaction, even when the current implementation is JSON files. If a caller needs a queue, lock, read-merge-write loop, or uniqueness race workaround, that behavior belongs behind the store surface.
---
## Directory layout

View File

@@ -1,9 +1,19 @@
import { type ReactNode, useCallback, useMemo, useReducer, useRef, useState } from "react";
import {
type ReactNode,
useCallback,
useEffect,
useMemo,
useReducer,
useRef,
useState,
} from "react";
import { Pressable, Text, View } from "react-native";
import type { PressableStateCallbackType } from "react-native";
import { StyleSheet } from "react-native-unistyles";
import { AdaptiveTextInput } from "@/components/adaptive-modal-sheet";
import { SegmentedControl } from "@/components/ui/segmented-control";
import { getDeviceTimeZone } from "@/utils/device-timezone";
import { nextCronCadence } from "@/utils/schedule-cadence-policy";
import {
describeCron,
everyMsToParts,
@@ -20,8 +30,8 @@ interface CronPreset {
expression: string;
}
// 5-field expressions evaluated in UTC by the daemon. Each one round-trips
// through describeCron() so the chip and the live preview agree.
// 5-field expressions use the cadence timezone. Each one round-trips through
// describeCron() so the chip and the live preview agree.
const CRON_PRESETS: CronPreset[] = [
{ label: "Every hour", expression: "0 * * * *" },
{ label: "Daily 9:00", expression: "0 9 * * *" },
@@ -57,6 +67,25 @@ function describeInterval(value: number, unit: IntervalUnit): string {
return `Runs every ${value} ${noun}s`;
}
function getCronPreview(expression: string, timezone: string, error: string | null): string | null {
if (error) {
return null;
}
if (!expression) {
return null;
}
const described = describeCron({ type: "cron", expression, timezone });
if (described) {
return described;
}
return expression;
}
function intervalCadenceKey(cadence: Extract<ScheduleCadence, { type: "every" }>): string {
return `${cadence.type}:${cadence.everyMs}`;
}
export interface CadenceEditorProps {
value: ScheduleCadence;
onChange: (next: ScheduleCadence) => void;
@@ -65,6 +94,13 @@ export interface CadenceEditorProps {
export function CadenceEditor({ value, onChange, error }: CadenceEditorProps) {
const mode = value.type;
const deviceTimeZone = useMemo(getDeviceTimeZone, []);
const rememberedCronTimeZone = useRef(
value.type === "cron" ? (value.timezone ?? "UTC") : deviceTimeZone,
);
const emittedIntervalCadenceKey = useRef<string | null>(null);
const cronTimeZone =
value.type === "cron" ? (value.timezone ?? "UTC") : rememberedCronTimeZone.current;
// The numeric/text fields are native-owned (AdaptiveTextInput). We seed them
// once from the incoming cadence via lazy state initializers and bump
@@ -88,6 +124,24 @@ export function CadenceEditor({ value, onChange, error }: CadenceEditorProps) {
value.type === "cron" ? value.expression : DEFAULT_CRON_EXPRESSION,
);
useEffect(() => {
if (value.type === "cron") {
emittedIntervalCadenceKey.current = null;
rememberedCronTimeZone.current = value.timezone ?? "UTC";
lastCronExpression.current = value.expression;
return;
}
const cadenceKey = intervalCadenceKey(value);
if (emittedIntervalCadenceKey.current === cadenceKey) {
emittedIntervalCadenceKey.current = null;
return;
}
rememberedCronTimeZone.current = deviceTimeZone;
lastCronExpression.current = DEFAULT_CRON_EXPRESSION;
setCronText(DEFAULT_CRON_EXPRESSION);
bumpFieldResetKey();
}, [deviceTimeZone, value]);
const parsedIntervalValue = useMemo(() => {
const parsed = Number.parseInt(intervalValueText, 10);
return Number.isFinite(parsed) && parsed > 0 ? parsed : 1;
@@ -95,17 +149,24 @@ export function CadenceEditor({ value, onChange, error }: CadenceEditorProps) {
const emitInterval = useCallback(
(rawValue: number, unit: IntervalUnit) => {
onChange({ type: "every", everyMs: partsToEveryMs(rawValue, unit) });
if (value.type === "cron") {
rememberedCronTimeZone.current = value.timezone ?? "UTC";
}
const next = { type: "every" as const, everyMs: partsToEveryMs(rawValue, unit) };
emittedIntervalCadenceKey.current = intervalCadenceKey(next);
onChange(next);
},
[onChange],
[onChange, value],
);
const emitCron = useCallback(
(expression: string) => {
lastCronExpression.current = expression;
onChange({ type: "cron", expression });
const next = nextCronCadence(value, expression, rememberedCronTimeZone.current);
rememberedCronTimeZone.current = next.timezone ?? "UTC";
onChange(next);
},
[onChange],
[onChange, value],
);
const handleModeChange = useCallback(
@@ -161,7 +222,7 @@ export function CadenceEditor({ value, onChange, error }: CadenceEditorProps) {
const intervalPreview = describeInterval(parsedIntervalValue, intervalUnit);
const trimmedCron = cronText.trim();
const cronError = trimmedCron ? validateCron(trimmedCron) : null;
const cronPreview = cronError ? null : (describeCron(trimmedCron) ?? trimmedCron);
const cronPreview = getCronPreview(trimmedCron, cronTimeZone, cronError);
let cronFeedback: ReactNode = null;
if (cronError) {
@@ -231,7 +292,7 @@ export function CadenceEditor({ value, onChange, error }: CadenceEditorProps) {
style={styles.cronInput}
/>
{cronFeedback}
<Text style={styles.hint}>Times are in UTC</Text>
<Text style={styles.hint}>Times are in {cronTimeZone}</Text>
</View>
)}

View File

@@ -0,0 +1,3 @@
export function getDeviceTimeZone(): string {
return Intl.DateTimeFormat().resolvedOptions().timeZone;
}

View File

@@ -0,0 +1,73 @@
import { describe, expect, it } from "vitest";
import { nextCronCadence } from "./schedule-cadence-policy";
describe("nextCronCadence", () => {
it("preserves an existing cron cadence timezone when editing the expression", () => {
expect(
nextCronCadence(
{
type: "cron",
expression: "0 9 * * *",
timezone: "America/New_York",
},
"30 9 * * *",
"America/Los_Angeles",
),
).toEqual({
type: "cron",
expression: "30 9 * * *",
timezone: "America/New_York",
});
});
it("emits UTC for legacy cron cadences without a timezone", () => {
expect(
nextCronCadence(
{
type: "cron",
expression: "0 9 * * *",
},
"30 9 * * *",
"America/Los_Angeles",
),
).toEqual({
type: "cron",
expression: "30 9 * * *",
timezone: "UTC",
});
});
it("uses the device timezone when switching from interval to cron", () => {
expect(
nextCronCadence(
{
type: "every",
everyMs: 60 * 60_000,
},
"0 9 * * *",
"America/Los_Angeles",
),
).toEqual({
type: "cron",
expression: "0 9 * * *",
timezone: "America/Los_Angeles",
});
});
it("lets callers preserve a remembered cron timezone when toggling back from interval", () => {
expect(
nextCronCadence(
{
type: "every",
everyMs: 60 * 60_000,
},
"0 9 * * *",
"America/New_York",
),
).toEqual({
type: "cron",
expression: "0 9 * * *",
timezone: "America/New_York",
});
});
});

View File

@@ -0,0 +1,14 @@
import type { ScheduleCadence } from "@getpaseo/protocol/schedule/types";
type CronCadence = Extract<ScheduleCadence, { type: "cron" }>;
export function nextCronCadence(
current: ScheduleCadence,
expression: string,
deviceTimeZone: string,
): CronCadence {
if (current.type === "cron") {
return { type: "cron", expression, timezone: current.timezone ?? "UTC" };
}
return { type: "cron", expression, timezone: deviceTimeZone };
}

View File

@@ -85,17 +85,38 @@ describe("interval formatting", () => {
describe("describeCron", () => {
it("humanizes common fixed-time cron shapes", () => {
expect(describeCron("0 * * * *")).toBe("Every hour");
expect(describeCron("15 * * * *")).toBe("Every hour at :15");
expect(describeCron("0 9 * * *")).toBe("Daily at 09:00 UTC");
expect(describeCron("0 9 * * 1-5")).toBe("Weekdays at 09:00 UTC");
expect(describeCron("0 9 * * 0,6")).toBe("Weekends at 09:00 UTC");
expect(describeCron("0 9 * * 1")).toBe("Mondays at 09:00 UTC");
expect(describeCron({ type: "cron", expression: "0 * * * *" })).toBe("Every hour");
expect(describeCron({ type: "cron", expression: "15 * * * *" })).toBe("Every hour at :15");
expect(describeCron({ type: "cron", expression: "0 9 * * *" })).toBe("Daily at 09:00 UTC");
expect(describeCron({ type: "cron", expression: "0 9 * * 1-5" })).toBe("Weekdays at 09:00 UTC");
expect(describeCron({ type: "cron", expression: "0 9 * * 0,6" })).toBe("Weekends at 09:00 UTC");
expect(describeCron({ type: "cron", expression: "0 9 * * 1" })).toBe("Mondays at 09:00 UTC");
});
it("labels fixed-time cron cadences with their stored timezone", () => {
expect(
describeCron({
type: "cron",
expression: "0 9 * * *",
timezone: "America/New_York",
}),
).toBe("Daily at 09:00 America/New_York");
expect(
formatCadence({
type: "cron",
expression: "0 9 * * 1-5",
timezone: "Europe/Madrid",
}),
).toBe("Weekdays at 09:00 Europe/Madrid");
});
it("keeps timezone-less fixed-time cron cadences labeled as UTC", () => {
expect(formatCadence({ type: "cron", expression: "0 9 * * *" })).toBe("Daily at 09:00 UTC");
});
it("returns null for invalid or unrecognized valid cron expressions", () => {
expect(describeCron("not a cron")).toBeNull();
expect(describeCron("*/5 * * * *")).toBeNull();
expect(describeCron({ type: "cron", expression: "not a cron" })).toBeNull();
expect(describeCron({ type: "cron", expression: "*/5 * * * *" })).toBeNull();
expect(formatCadence({ type: "cron", expression: "0 9 * * *" })).toBe("Daily at 09:00 UTC");
});
});

View File

@@ -2,6 +2,7 @@ import type { ScheduleCadence, ScheduleSummary } from "@getpaseo/protocol/schedu
import { validateCronExpression } from "@getpaseo/protocol/schedule/cron-expression";
export type IntervalUnit = "minutes" | "hours" | "days";
type CronCadence = Extract<ScheduleCadence, { type: "cron" }>;
const MS_PER_MINUTE = 60_000;
const MS_PER_HOUR = MS_PER_MINUTE * 60;
@@ -82,7 +83,7 @@ export function formatCadence(cadence: ScheduleCadence): string {
if (cadence.type === "every") {
return formatEvery(cadence.everyMs);
}
return describeCron(cadence.expression) ?? cadence.expression;
return describeCron(cadence) ?? cadence.expression;
}
/**
@@ -90,8 +91,8 @@ export function formatCadence(cadence: ScheduleCadence): string {
* expression is valid but not one of the recognized patterns (callers fall
* back to showing the raw expression).
*/
export function describeCron(expr: string): string | null {
const trimmed = expr.trim();
export function describeCron(cadence: CronCadence): string | null {
const trimmed = cadence.expression.trim();
if (validateCron(trimmed) !== null) {
return null;
}
@@ -121,20 +122,21 @@ export function describeCron(expr: string): string | null {
return null;
}
const time = `${pad2(Number.parseInt(hour, 10))}:${pad2(minuteNum)}`;
const timezone = cadence.timezone ?? "UTC";
if (dayOfWeek === "*") {
return `Daily at ${time} UTC`;
return `Daily at ${time} ${timezone}`;
}
if (dayOfWeek === "1-5") {
return `Weekdays at ${time} UTC`;
return `Weekdays at ${time} ${timezone}`;
}
if (dayOfWeek === "0,6" || dayOfWeek === "6,0") {
return `Weekends at ${time} UTC`;
return `Weekends at ${time} ${timezone}`;
}
if (/^\d$/.test(dayOfWeek)) {
const day = DAY_NAMES[Number.parseInt(dayOfWeek, 10)];
if (day) {
return `${day}s at ${time} UTC`;
return `${day}s at ${time} ${timezone}`;
}
}
return null;

View File

@@ -0,0 +1,26 @@
import { describe, expect, it } from "vitest";
import { ScheduleCreateRequestSchema } from "./rpc-schemas.js";
describe("schedule RPC schemas", () => {
it("keeps new-agent workspace stamps out of create requests", () => {
const parsed = ScheduleCreateRequestSchema.parse({
type: "schedule/create",
requestId: "request-1",
prompt: "Run the task",
cadence: { type: "every", everyMs: 60_000 },
target: {
type: "new-agent",
config: {
provider: "claude",
cwd: "/tmp/project",
workspaceId: "client-owned-workspace",
},
},
});
expect(parsed.target.type).toBe("new-agent");
if (parsed.target.type === "new-agent") {
expect(parsed.target.config).not.toHaveProperty("workspaceId");
}
});
});

View File

@@ -7,6 +7,10 @@ import {
ScheduleTargetSchema,
} from "./types.js";
const ScheduleCreateNewAgentConfigSchema = ScheduleTargetSchema.options[1].shape.config.omit({
workspaceId: true,
});
const ScheduleCreateTargetSchema = z.discriminatedUnion("type", [
z.object({
type: z.literal("self"),
@@ -18,7 +22,7 @@ const ScheduleCreateTargetSchema = z.discriminatedUnion("type", [
}),
z.object({
type: z.literal("new-agent"),
config: ScheduleTargetSchema.options[1].shape.config,
config: ScheduleCreateNewAgentConfigSchema,
}),
]);

View File

@@ -27,6 +27,7 @@ export const ScheduleTargetSchema = z.discriminatedUnion("type", [
config: z.object({
provider: AgentProviderSchema,
cwd: z.string().trim().min(1),
workspaceId: z.string().optional(),
modeId: z.string().trim().min(1).optional(),
model: z.string().trim().min(1).optional(),
thinkingOptionId: z.string().trim().min(1).optional(),

View File

@@ -13,8 +13,17 @@ import {
const pendingAgentInitializations = new Map<string, Promise<ManagedAgent>>();
type AgentLoaderManager = Pick<
AgentManager,
| "createAgent"
| "getAgent"
| "getRegisteredProviderIds"
| "hydrateTimelineFromProvider"
| "resumeAgentFromPersistence"
>;
export interface EnsureAgentLoadedDeps {
agentManager: AgentManager;
agentManager: AgentLoaderManager;
agentStorage: AgentStorage;
validProviders?: Iterable<AgentProvider>;
logger: Logger;

View File

@@ -278,6 +278,7 @@ async function createManagedSession(
cwd: workdir,
},
agentId,
{ workspaceId: undefined },
);
return {
agentId,

File diff suppressed because it is too large Load Diff

View File

@@ -196,6 +196,16 @@ interface ProviderEnabledFlag {
type ProviderEnabledMap = Partial<Record<AgentProvider, ProviderEnabledFlag>>;
type ProviderClientMap = Partial<Record<AgentProvider, AgentClient>>;
export interface CreateAgentOptions {
labels?: Record<string, string>;
initialPrompt?: string;
env?: Record<string, string>;
persistSession?: boolean;
initialTitle?: string | null;
// undefined is an explicit decision: the agent never appears in the sidebar.
workspaceId: string | undefined;
}
export interface AgentManagerOptions {
clients?: ProviderClientMap;
providerDefinitions?: ProviderEnabledMap;
@@ -914,15 +924,8 @@ export class AgentManager {
async createAgent(
config: AgentSessionConfig,
agentId?: string,
options?: {
labels?: Record<string, string>;
initialPrompt?: string;
env?: Record<string, string>;
persistSession?: boolean;
initialTitle?: string | null;
workspaceId?: string;
},
agentId: string | undefined,
options: CreateAgentOptions,
): Promise<ManagedAgent> {
const resolvedAgentId = validateAgentId(agentId ?? this.idFactory(), "createAgent");
const { storedConfig, launchConfig } = await this.prepareSessionConfig(config, resolvedAgentId);
@@ -935,9 +938,9 @@ export class AgentManager {
const createOptions = this.buildCreateSessionOptions(options);
const session = await client.createSession(providerLaunchConfig, launchContext, createOptions);
return this.registerSession(session, storedConfig, resolvedAgentId, {
labels: options?.labels,
initialTitle: options?.initialTitle,
workspaceId: options?.workspaceId,
labels: options.labels,
initialTitle: options.initialTitle,
workspaceId: options.workspaceId,
});
}

View File

@@ -182,6 +182,15 @@ function buildBasePrompt(prompt: string, jsonSchema: JsonSchema): string {
].join("\n");
}
export function buildStructuredAgentResponsePrompt(options: {
prompt: string;
schema: z.ZodType | JsonSchema;
schemaName?: string;
}): string {
const validator = buildValidator(options.schema, options.schemaName ?? "Response");
return buildBasePrompt(options.prompt, validator.jsonSchema);
}
function buildRetryPrompt(basePrompt: string, errors: string[]): string {
const formattedErrors = errors.map((error) => `- ${error}`).join("\n");
return [
@@ -347,7 +356,10 @@ export async function generateStructuredAgentResponse<T>(
): Promise<T> {
const { manager, agentConfig, agentId, persistSession, prompt, schema, maxRetries, schemaName } =
options;
const agent = await manager.createAgent(agentConfig, agentId, { persistSession });
const agent = await manager.createAgent(agentConfig, agentId, {
persistSession,
workspaceId: undefined,
});
try {
const caller: AgentCaller = async (nextPrompt) => {
const result = await manager.runAgent(agent.id, nextPrompt);

View File

@@ -66,6 +66,7 @@ test("session create forwards clientMessageId to the initial prompt run options"
await createAgentCommand(dependencies, {
kind: "session",
config: { provider: "codex", cwd: "/tmp/paseo-create-test" },
workspaceId: "ws-create-test",
initialPrompt: "hello from create",
clientMessageId: "msg-create-1",
labels: {},
@@ -79,6 +80,52 @@ test("session create forwards clientMessageId to the initial prompt run options"
});
});
test("mcp create accepts provider-only internal input and leaves model undefined", async () => {
const snapshot = {
id: "agent-1",
provider: "claude",
cwd: "/tmp/paseo-create-test",
runtimeInfo: null,
} as ManagedAgent;
const createAgent = vi.fn(async () => snapshot);
const dependencies: Parameters<typeof createAgentCommand>[0] = {
agentManager: {
createAgent,
getAgent: vi.fn(() => snapshot),
} as unknown as Parameters<typeof createAgentCommand>[0]["agentManager"],
agentStorage: {} as Parameters<typeof createAgentCommand>[0]["agentStorage"],
logger: createTestLogger(),
providerSnapshotManager: {
resolveCreateConfig: vi.fn(async (input) => {
expect(input.provider).toBe("claude");
return {};
}),
} as Parameters<typeof createAgentCommand>[0]["providerSnapshotManager"],
};
await createAgentCommand(dependencies, {
kind: "mcp",
provider: "claude",
cwd: "/tmp/paseo-create-test",
workspaceId: "ws-create-test",
title: "provider default",
initialPrompt: "hello",
background: true,
notifyOnFinish: false,
});
expect(createAgent).toHaveBeenCalledWith(
expect.objectContaining({
provider: "claude",
model: undefined,
}),
undefined,
expect.objectContaining({
workspaceId: "ws-create-test",
}),
);
});
test("session create stamps the requested workspaceId when no worktree setup runs", async () => {
const workdir = mkdtempSync(join(tmpdir(), "create-agent-test-"));
const storage = new AgentStorage(join(workdir, "agents"), logger);
@@ -212,6 +259,7 @@ test("session create keeps the prompt title after the initial prompt settles", a
{
kind: "session",
config: { provider: "codex", cwd: workdir },
workspaceId: "ws-title-source",
initialPrompt: `${title}\n\ninclude tests`,
labels: {},
provisionalTitle: title,
@@ -249,6 +297,7 @@ test("session create keeps an explicit title after the initial prompt settles",
{
kind: "session",
config: { provider: "codex", cwd: workdir, title },
workspaceId: "ws-explicit-title-source",
initialPrompt: "Implement auth retries with backoff",
labels: {},
provisionalTitle: title,

View File

@@ -12,7 +12,7 @@ import type {
CreatePaseoWorktreeWorkflowResult,
} from "../../worktree-session.js";
import type { AgentAttachment, FirstAgentContext, GitSetupOptions } from "../../messages.js";
import type { AgentManager, ManagedAgent } from "../agent-manager.js";
import type { AgentManager, CreateAgentOptions, ManagedAgent } from "../agent-manager.js";
import type {
AgentPromptContentBlock,
AgentPromptInput,
@@ -22,8 +22,9 @@ import type {
import type { AgentStorage } from "../agent-storage.js";
import type { ProviderSnapshotManager } from "../provider-snapshot-manager.js";
import { setupFinishNotification, startCreatedAgentInitialPrompt } from "../agent-prompt.js";
import { resolveCreateAgentTitles } from "../create-agent-title.js";
import { normalizeClientMessageId, resolveClientMessageId } from "../../client-message-id.js";
import { resolveRequiredProviderModel } from "../mcp-shared.js";
import { resolveRequiredProviderModel, type ResolvedProviderModel } from "../mcp-shared.js";
import {
appendTimelineItemIfAgentKnown,
emitLiveTimelineItemIfAgentKnown,
@@ -37,7 +38,7 @@ export interface CreateAgentSessionWorktreeResult {
createdWorkspaceId?: string;
}
interface CreateAgentCommandDependencies {
export interface CreateAgentCommandDependencies {
agentManager: AgentManager;
agentStorage: AgentStorage;
logger: Logger;
@@ -47,16 +48,18 @@ interface CreateAgentCommandDependencies {
providerSnapshotManager: ProviderSnapshotManager;
createPaseoWorktree?: CreatePaseoWorktreeWorkflowFn;
// Mints a fresh directory workspace for a cwd and returns its id.
ensureWorkspaceForCreate?: (
cwd: string,
firstAgentContext?: FirstAgentContext,
) => Promise<string>;
ensureWorkspaceForCreate?: EnsureWorkspaceForCreate;
}
export type EnsureWorkspaceForCreate = (
cwd: string,
firstAgentContext?: FirstAgentContext,
) => Promise<string>;
export interface CreateAgentFromSessionInput {
kind: "session";
config: AgentSessionConfig;
workspaceId?: string;
workspaceId: string;
worktreeName?: string;
initialPrompt?: string;
clientMessageId?: string;
@@ -80,15 +83,19 @@ export interface CreateAgentFromMcpInput {
kind: "mcp";
provider: string;
title: string;
initialPrompt: string;
initialPrompt?: string;
config?: Partial<AgentSessionConfig>;
cwd?: string;
workspaceId?: string;
thinking?: string;
features?: Record<string, unknown>;
labels?: Record<string, string>;
mode?: string;
unattended?: boolean;
promptFailure?: CreateAgentPromptFailureMode;
background: boolean;
notifyOnFinish: boolean;
internal?: boolean;
detached?: boolean;
callerAgentId?: string;
callerContext?: {
@@ -107,33 +114,56 @@ export interface CreateAgentFromMcpInput {
}
export type CreateAgentCommandInput = CreateAgentFromSessionInput | CreateAgentFromMcpInput;
export type CreateAgentPromptFailureMode = "throw" | "log" | "return-error";
export interface CreateAgentCommandResult {
snapshot: ManagedAgent;
liveSnapshot: ManagedAgent;
background: boolean;
initialPromptStarted: boolean;
initialPromptError: unknown | null;
}
export type BoundCreateAgentCommand = (
input: CreateAgentCommandInput,
) => Promise<CreateAgentCommandResult>;
function requireResolvedWorkspaceId(workspaceId: string | undefined): string {
if (!workspaceId) {
throw new Error("createAgentCommand requires a resolved workspaceId");
}
return workspaceId;
}
export function formatProviderModel(provider: string, model: string | null | undefined): string {
if (!model || provider.includes("/")) {
return provider;
}
return `${provider}/${model}`;
}
function resolveProviderModel(providerValue: string): ResolvedProviderModel {
const providerInput = providerValue.trim();
if (providerInput.includes("/")) {
return resolveRequiredProviderModel(providerInput);
}
if (!providerInput) {
throw new Error("provider is required");
}
return { provider: providerInput, model: undefined };
}
interface ResolvedCreateAgent {
config: AgentSessionConfig;
createOptions?: AgentCreateOptions;
createOptions: CreateAgentOptions;
prompt?: AgentPromptInput;
runOptions?: AgentRunOptions;
setupContinuation?: AgentWorktreeSetupContinuation;
background: boolean;
promptFailure: "throw" | "log";
promptFailure: CreateAgentPromptFailureMode;
promptLogger?: Logger;
}
interface AgentCreateOptions {
labels?: Record<string, string>;
initialPrompt?: string;
env?: Record<string, string>;
initialTitle?: string | null;
workspaceId?: string;
}
export async function createAgentCommand(
dependencies: CreateAgentCommandDependencies,
input: CreateAgentCommandInput,
@@ -155,10 +185,12 @@ export async function createAgentCommand(
let liveSnapshot = snapshot;
let initialPromptStarted = false;
let initialPromptError: unknown | null = null;
if (resolved.prompt !== undefined) {
const sendResult = await sendInitialPrompt(dependencies, resolved, snapshot);
initialPromptStarted = sendResult.started;
liveSnapshot = sendResult.liveSnapshot;
initialPromptError = sendResult.error ?? null;
}
if (input.kind === "mcp" && input.notifyOnFinish && input.callerAgentId && initialPromptStarted) {
@@ -176,6 +208,7 @@ export async function createAgentCommand(
liveSnapshot,
background: resolved.background,
initialPromptStarted,
initialPromptError,
};
}
@@ -200,6 +233,7 @@ async function resolveSessionCreateAgent(
...(clientMessageId ? { messageId: clientMessageId } : {}),
}
: undefined;
const workspaceId = setupContinuation ? createdWorkspaceId : input.workspaceId;
return {
config: sessionConfig,
@@ -209,9 +243,9 @@ async function resolveSessionCreateAgent(
env: input.env,
initialTitle: input.provisionalTitle,
// A legacy git/worktreeName worktree creates a fresh workspace, so the
// agent belongs to that workspace, not the source one (mirrors the MCP
// path). createdWorkspaceId is the freshly created worktree's workspace.
workspaceId: setupContinuation ? createdWorkspaceId : input.workspaceId,
// agent belongs to that workspace, not the source one. createdWorkspaceId
// is the freshly created worktree's workspace.
workspaceId: requireResolvedWorkspaceId(workspaceId),
},
prompt: hasPromptContent ? prompt : undefined,
runOptions,
@@ -228,45 +262,34 @@ async function resolveMcpCreateAgent(
dependencies: CreateAgentCommandDependencies,
input: CreateAgentFromMcpInput,
): Promise<ResolvedCreateAgent> {
const resolvedProviderModel = resolveRequiredProviderModel(input.provider);
const resolvedProviderModel = resolveProviderModel(input.provider);
const provider = resolvedProviderModel.provider;
const parentAgent = input.callerAgentId
? requireParentAgent(dependencies.agentManager, input.callerAgentId)
: null;
const cwd = parentAgent
? resolveChildAgentCwd({
parentCwd: parentAgent.cwd,
requestedCwd: input.cwd,
lockedCwd: input.callerContext?.lockedCwd,
allowCustomCwd: input.callerContext?.allowCustomCwd ?? true,
})
: expandUserPath(input.cwd ?? process.cwd());
const cwd = resolveMcpInitialCwd(input, parentAgent);
const { resolvedCwd, setupContinuation, createdWorkspaceId } = await resolveMcpCwd({
dependencies,
cwd,
worktree: input.worktree,
initialPrompt: input.initialPrompt,
initialPrompt: input.initialPrompt ?? "",
});
// MCP callers resolve workspace ownership before this point. Worktree
// creation wins because the new agent lives in the fresh worktree workspace.
// Otherwise use the explicit workspace id, then the parent workspace for
// direct internal callers. Ownership is never resolved from cwd.
const workspaceId = setupContinuation
? createdWorkspaceId
: (input.workspaceId ??
parentAgent?.workspaceId ??
(await ensureWorkspaceForMcpCreate(dependencies, resolvedCwd, input.initialPrompt)));
const { modeId: resolvedMode, featureValues: resolvedFeatures } =
await dependencies.providerSnapshotManager.resolveCreateConfig({
cwd: resolvedCwd,
provider,
requestedMode: input.mode,
featureValues: input.features,
parent: parentAgent,
unattended: false,
});
const workspaceId = await resolveMcpWorkspaceId({
dependencies,
input,
parentAgent,
setupContinuation,
createdWorkspaceId,
resolvedCwd,
});
const resolvedCreateConfig = await resolveMcpProviderCreateConfig({
dependencies,
input,
provider,
resolvedCwd,
parentAgent,
});
const labels = mergeLabels({
callerAgentId: input.callerAgentId,
@@ -275,31 +298,122 @@ async function resolveMcpCreateAgent(
labels: input.labels,
});
const trimmedPrompt = input.initialPrompt.trim();
const trimmedPrompt = input.initialPrompt?.trim() ?? "";
return {
config: {
config: buildMcpSessionConfig({
input,
resolvedProviderModel,
provider,
cwd: resolvedCwd,
modeId: resolvedMode,
title: input.title.trim(),
model: resolvedProviderModel.model,
thinkingOptionId: input.thinking,
...(resolvedFeatures ? { featureValues: resolvedFeatures } : {}),
resolvedCwd,
trimmedPrompt,
resolvedMode: resolvedCreateConfig.modeId,
resolvedFeatures: resolvedCreateConfig.featureValues,
}),
createOptions: {
...(labels ? { labels } : {}),
workspaceId: requireResolvedWorkspaceId(workspaceId),
},
createOptions:
labels || workspaceId
? {
...(labels ? { labels } : {}),
...(workspaceId ? { workspaceId } : {}),
}
: undefined,
prompt: trimmedPrompt,
prompt: trimmedPrompt ? trimmedPrompt : undefined,
setupContinuation,
background: input.background,
promptFailure: "log",
promptFailure: input.promptFailure ?? "log",
};
}
function resolveMcpInitialCwd(
input: CreateAgentFromMcpInput,
parentAgent: ManagedAgent | null,
): string {
if (!parentAgent) {
return expandUserPath(input.cwd ?? process.cwd());
}
return resolveChildAgentCwd({
parentCwd: parentAgent.cwd,
requestedCwd: input.cwd,
lockedCwd: input.callerContext?.lockedCwd,
allowCustomCwd: input.callerContext?.allowCustomCwd ?? true,
});
}
async function resolveMcpWorkspaceId(params: {
dependencies: CreateAgentCommandDependencies;
input: CreateAgentFromMcpInput;
parentAgent: ManagedAgent | null;
setupContinuation?: AgentWorktreeSetupContinuation;
createdWorkspaceId?: string;
resolvedCwd: string;
}): Promise<string | undefined> {
// MCP callers resolve workspace ownership before this point. Worktree
// creation wins because the new agent lives in the fresh worktree workspace.
// Otherwise use the explicit workspace id, then the parent workspace for
// direct internal callers. Ownership is never resolved from cwd.
if (params.setupContinuation) {
return params.createdWorkspaceId;
}
if (params.input.workspaceId) {
return params.input.workspaceId;
}
if (params.parentAgent?.workspaceId) {
return params.parentAgent.workspaceId;
}
return ensureWorkspaceForMcpCreate(
params.dependencies,
params.resolvedCwd,
params.input.initialPrompt ?? "",
);
}
async function resolveMcpProviderCreateConfig(params: {
dependencies: CreateAgentCommandDependencies;
input: CreateAgentFromMcpInput;
provider: string;
resolvedCwd: string;
parentAgent: ManagedAgent | null;
}): Promise<{ modeId?: string; featureValues?: Record<string, unknown> }> {
const passthroughConfig = params.input.config;
return params.dependencies.providerSnapshotManager.resolveCreateConfig({
cwd: params.resolvedCwd,
provider: params.provider,
requestedMode: params.input.mode ?? passthroughConfig?.modeId,
featureValues: params.input.features ?? passthroughConfig?.featureValues,
parent: params.parentAgent,
unattended: params.input.unattended ?? false,
});
}
function buildMcpSessionConfig(params: {
input: CreateAgentFromMcpInput;
resolvedProviderModel: ResolvedProviderModel;
provider: string;
resolvedCwd: string;
trimmedPrompt: string;
resolvedMode?: string;
resolvedFeatures?: Record<string, unknown>;
}): AgentSessionConfig {
const passthroughConfig = params.input.config;
const { provisionalTitle } = resolveCreateAgentTitles({
configTitle: passthroughConfig?.title ?? params.input.title,
initialPrompt: params.trimmedPrompt,
});
const featureValues = params.resolvedFeatures ?? passthroughConfig?.featureValues;
const config: AgentSessionConfig = {
...passthroughConfig,
provider: params.provider,
cwd: params.resolvedCwd,
modeId: params.resolvedMode ?? passthroughConfig?.modeId,
model: params.resolvedProviderModel.model ?? passthroughConfig?.model,
thinkingOptionId: params.input.thinking ?? passthroughConfig?.thinkingOptionId,
internal: params.input.internal ?? passthroughConfig?.internal,
};
if (provisionalTitle) {
config.title = provisionalTitle;
}
if (featureValues) {
config.featureValues = featureValues;
}
return config;
}
async function ensureWorkspaceForMcpCreate(
dependencies: CreateAgentCommandDependencies,
cwd: string,
@@ -315,7 +429,7 @@ async function sendInitialPrompt(
dependencies: CreateAgentCommandDependencies,
resolved: ResolvedCreateAgent,
snapshot: ManagedAgent,
): Promise<{ started: boolean; liveSnapshot: ManagedAgent }> {
): Promise<{ started: boolean; liveSnapshot: ManagedAgent; error?: unknown }> {
try {
const prompt = resolved.prompt;
if (prompt === undefined) {
@@ -334,6 +448,9 @@ async function sendInitialPrompt(
if (resolved.promptFailure === "throw") {
throw error;
}
if (resolved.promptFailure === "return-error") {
return { started: false, liveSnapshot: snapshot, error };
}
dependencies.logger.error({ err: error, agentId: snapshot.id }, "Failed to run initial prompt");
return { started: false, liveSnapshot: snapshot };
}

View File

@@ -358,12 +358,16 @@ describe("Suite A: Core Fixes", () => {
expect(listenTarget?.type).toBe("tcp");
const cwd = await makeCwd("manager-direct-agent-cwd");
const snapshot = await daemonHandle.daemon.agentManager.createAgent({
provider: "claude",
cwd,
title: "Manager direct parity agent",
modeId: "bypassPermissions",
});
const snapshot = await daemonHandle.daemon.agentManager.createAgent(
{
provider: "claude",
cwd,
title: "Manager direct parity agent",
modeId: "bypassPermissions",
},
undefined,
{ workspaceId: undefined },
);
agentId = snapshot.id;
const expectedUrl = buildExpectedAgentMcpUrl({

View File

@@ -1600,7 +1600,10 @@ describe("create_agent MCP tool", () => {
thinkingOptionId: "think-hard",
}),
undefined,
{ labels: { source: "mcp" }, workspaceId: "workspace-created" },
{
labels: { source: "mcp" },
workspaceId: "workspace-created",
},
);
});
@@ -3091,7 +3094,9 @@ describe("create_agent MCP tool", () => {
});
expect(configArg.mcpServers).toBeUndefined();
expect(agentIdArg).toBeUndefined();
expect(optionsArg).toEqual({ workspaceId: "workspace-created" });
expect(optionsArg).toEqual({
workspaceId: "workspace-created",
});
});
it("rejects an explicit mode that is not valid for the target provider", async () => {

View File

@@ -250,6 +250,7 @@ describe("MockLoadTestAgentClient", () => {
model: "ten-second-stream",
},
"00000000-0000-4000-8000-000000000001",
{ workspaceId: undefined },
);
const resultPromise = manager.runAgent(

View File

@@ -56,10 +56,14 @@ async function createRewindHarness(options: { historyGate?: RewindHistoryGate }
logger: createTestLogger(),
idFactory: () => "00000000-0000-4000-8000-000000000901",
});
const agent = await manager.createAgent({
provider: "claude",
cwd: process.cwd(),
});
const agent = await manager.createAgent(
{
provider: "claude",
cwd: process.cwd(),
},
undefined,
{ workspaceId: undefined },
);
return { manager, session, agentId: agent.id };
}

View File

@@ -166,6 +166,10 @@ import { WorkspaceAutoName } from "./workspace-auto-name.js";
import { createGitMutationService } from "./session/git-mutation/git-mutation-service.js";
import { workspaceIdsOnCheckout } from "./workspace-directory.js";
import { resolveFirstAgentPromptTitle } from "./agent/create-agent-title.js";
import {
createAgentCommand,
type CreateAgentCommandDependencies,
} from "./agent/create-agent/create.js";
const MAX_MCP_DEBUG_BATCH_ITEMS = 10;
const REDACTED_LOG_VALUE = "[redacted]";
@@ -802,41 +806,6 @@ export async function createPaseoDaemon(
paseoHome: config.paseoHome,
workspaceGitService,
});
const loopService = new LoopService({
paseoHome: config.paseoHome,
logger,
agentManager,
providerSnapshotManager,
});
await loopService.initialize();
logger.info({ elapsed: elapsed() }, "Loop service initialized");
const scheduleService = new ScheduleService({
paseoHome: config.paseoHome,
logger,
agentManager,
agentStorage,
providerSnapshotManager,
});
await scheduleService.start();
agentManager.setAgentArchivedCallback(async (agentId) => {
try {
await scheduleService.completeForAgent(agentId);
} catch (error) {
logger.warn({ err: error, agentId }, "Failed to complete schedules for archived agent");
}
});
logger.info({ elapsed: elapsed() }, "Schedule service initialized");
logger.info({ elapsed: elapsed() }, "Loading persisted agent registry");
const persistedRecords = await agentStorage.list();
logger.info(
{ elapsed: elapsed() },
`Agent registry loaded (${persistedRecords.length} record${persistedRecords.length === 1 ? "" : "s"}); agents will initialize on demand`,
);
logger.info(
"Voice mode configured for agent-scoped resume flow (no dedicated voice assistant provider)",
);
logger.info({ elapsed: elapsed() }, "Preparing voice and MCP runtime");
const archiveWorkspaceRecordExternal = async (workspaceId: string) => {
const sessions = wsServer?.listActiveSessions() ?? [];
if (sessions.length > 0) {
@@ -903,6 +872,14 @@ export async function createPaseoDaemon(
),
);
};
const ensureWorkspaceForCreateAndBroadcastExternal = async (
cwd: string,
firstAgentContext?: FirstAgentContext,
): Promise<string> => {
const workspaceId = await ensureWorkspaceForCreateExternal(cwd, firstAgentContext);
await emitWorkspaceUpdatesExternal([workspaceId]);
return workspaceId;
};
const emitWorkspaceUpdateForCwdExternal = async (cwd: string) => {
const workspaceIds = workspaceIdsOnCheckout(await workspaceRegistry.list(), cwd);
await emitWorkspaceUpdatesExternal(workspaceIds);
@@ -996,6 +973,58 @@ export async function createPaseoDaemon(
);
};
const createAgentCommandDependencies: CreateAgentCommandDependencies = {
agentManager,
agentStorage,
logger,
paseoHome: config.paseoHome,
worktreesRoot: config.worktreesRoot,
terminalManager,
providerSnapshotManager,
createPaseoWorktree: createPaseoWorktreeForTools,
ensureWorkspaceForCreate: ensureWorkspaceForCreateAndBroadcastExternal,
};
const createAgent = (input: Parameters<typeof createAgentCommand>[1]) =>
createAgentCommand(createAgentCommandDependencies, input);
const loopService = new LoopService({
paseoHome: config.paseoHome,
logger,
agentManager,
createAgent,
ensureWorkspaceForCreate: ensureWorkspaceForCreateAndBroadcastExternal,
});
await loopService.initialize();
logger.info({ elapsed: elapsed() }, "Loop service initialized");
const scheduleService = new ScheduleService({
paseoHome: config.paseoHome,
logger,
agentManager,
agentStorage,
createAgent,
ensureWorkspaceForCreate: ensureWorkspaceForCreateAndBroadcastExternal,
workspaceRegistry,
});
await scheduleService.start();
agentManager.setAgentArchivedCallback(async (agentId) => {
try {
await scheduleService.completeForAgent(agentId);
} catch (error) {
logger.warn({ err: error, agentId }, "Failed to complete schedules for archived agent");
}
});
logger.info({ elapsed: elapsed() }, "Schedule service initialized");
logger.info({ elapsed: elapsed() }, "Loading persisted agent registry");
const persistedRecords = await agentStorage.list();
logger.info(
{ elapsed: elapsed() },
`Agent registry loaded (${persistedRecords.length} record${persistedRecords.length === 1 ? "" : "s"}); agents will initialize on demand`,
);
logger.info(
"Voice mode configured for agent-scoped resume flow (no dedicated voice assistant provider)",
);
logger.info({ elapsed: elapsed() }, "Preparing voice and MCP runtime");
const createAgentToolHostDependencies = (
runtime: PaseoToolRuntimeContext,
): PaseoToolHostDependencies => ({
@@ -1014,8 +1043,8 @@ export async function createPaseoDaemon(
workspaceRegistry,
markWorkspaceArchiving: markWorkspaceArchivingExternal,
clearWorkspaceArchiving: clearWorkspaceArchivingExternal,
ensureWorkspaceForCreate: ensureWorkspaceForCreateExternal,
createPaseoWorktree: createPaseoWorktreeForTools,
ensureWorkspaceForCreate: createAgentCommandDependencies.ensureWorkspaceForCreate,
createPaseoWorktree: createAgentCommandDependencies.createPaseoWorktree,
browserToolsEnabled: browserToolsPolicy.isEnabled(),
browserToolsBroker,
paseoHome: config.paseoHome,

View File

@@ -30,7 +30,11 @@ import type {
} from "./agent/agent-sdk-types.js";
import { AgentStorage } from "./agent/agent-storage.js";
import { AgentManager } from "./agent/agent-manager.js";
import { createAgentCommand } from "./agent/create-agent/create.js";
import type { ProviderSnapshotManager } from "./agent/provider-snapshot-manager.js";
import { createLocalCheckoutWorkspace } from "./paseo-worktree-service.js";
import { createNoopWorkspaceGitService } from "./test-utils/workspace-git-service-stub.js";
import { FileBackedProjectRegistry, FileBackedWorkspaceRegistry } from "./workspace-registry.js";
import { LoopService } from "./loop-service.js";
import { isPlatform } from "../test-utils/platform.js";
import { createTestLogger } from "../test-utils/test-logger.js";
@@ -46,11 +50,76 @@ const TEST_CAPABILITIES: AgentCapabilityFlags = {
const NO_UNATTENDED_LOOP_POLICY: Pick<ProviderSnapshotManager, "resolveCreateConfig"> = {
async resolveCreateConfig(input) {
expect(input).toMatchObject({ parent: null, unattended: true, requestedMode: undefined });
return { modeId: undefined, featureValues: input.featureValues };
expect(input).toMatchObject({ parent: null, unattended: true });
return {
modeId: input.unattended ? input.requestedMode : "interactive",
featureValues: input.featureValues,
};
},
};
interface TestLoopServiceOptions {
paseoHome: string;
agentManager: AgentManager;
agentStorage: AgentStorage;
logger: ReturnType<typeof createTestLogger>;
providerSnapshotManager?: Pick<ProviderSnapshotManager, "resolveCreateConfig">;
ensureWorkspaceForCreate?: (
cwd: string,
firstAgentContext?: { prompt: string },
) => Promise<string>;
}
function createLoopService(options: TestLoopServiceOptions): LoopService {
const providerSnapshotManager = options.providerSnapshotManager ?? NO_UNATTENDED_LOOP_POLICY;
const ensureWorkspaceForCreate =
options.ensureWorkspaceForCreate ?? (async () => "workspace-created-for-loop");
return new LoopService({
paseoHome: options.paseoHome,
agentManager: options.agentManager,
logger: options.logger,
ensureWorkspaceForCreate,
createAgent: (input) =>
createAgentCommand(
{
agentManager: options.agentManager,
agentStorage: options.agentStorage,
logger: options.logger,
providerSnapshotManager: providerSnapshotManager as ProviderSnapshotManager,
ensureWorkspaceForCreate,
},
input,
),
});
}
async function createRegistryBackedWorkspaceEnsure(rootDir: string): Promise<{
workspaceRegistry: FileBackedWorkspaceRegistry;
ensureWorkspaceForCreate: TestLoopServiceOptions["ensureWorkspaceForCreate"];
}> {
const workspaceRegistry = new FileBackedWorkspaceRegistry(
path.join(rootDir, "projects", "workspaces.json"),
createTestLogger(),
);
const projectRegistry = new FileBackedProjectRegistry(
path.join(rootDir, "projects", "projects.json"),
createTestLogger(),
);
await workspaceRegistry.initialize();
await projectRegistry.initialize();
const workspaceGitService = createNoopWorkspaceGitService();
return {
workspaceRegistry,
ensureWorkspaceForCreate: async (cwd, firstAgentContext) => {
const workspace = await createLocalCheckoutWorkspace(
{ cwd, title: firstAgentContext?.prompt ?? null },
{ projectRegistry, workspaceRegistry, workspaceGitService },
);
return workspace.workspaceId;
},
};
}
interface ScriptedAgentBehavior {
onRun(input: { config: AgentSessionConfig; prompt: string; turnId: string }): Promise<string>;
}
@@ -273,17 +342,18 @@ describe("LoopService", () => {
registry: storage,
logger,
});
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: NO_UNATTENDED_LOOP_POLICY,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Create done.txt when the task is actually fixed.",
cwd: workspaceDir,
model: "test-model",
verifyChecks: [
`${JSON.stringify(process.execPath)} ${JSON.stringify(path.basename(verifyScriptPath))}`,
],
@@ -329,11 +399,11 @@ describe("LoopService", () => {
registry: storage,
logger,
});
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: NO_UNATTENDED_LOOP_POLICY,
});
await service.initialize();
@@ -363,16 +433,150 @@ describe("LoopService", () => {
expect(workerConfigs[0]).toMatchObject({
provider: "codex",
model: "gpt-5.4",
internal: true,
});
expect(verifierConfigs).toHaveLength(1);
expect(verifierConfigs[0]).toMatchObject({
provider: "claude",
model: "sonnet",
internal: true,
});
});
test("loop worker and verifier agents share one registry workspace across iterations", async () => {
const { workspaceRegistry, ensureWorkspaceForCreate } =
await createRegistryBackedWorkspaceEnsure(tmpDir);
let verifierCount = 0;
const manager = new AgentManager({
clients: {
claude: new ScriptedAgentClient("claude", {
async onRun({ config }) {
if (config.title?.includes("worker")) {
return "worker finished";
}
verifierCount += 1;
return verifierCount >= 2
? '{"passed":true,"reason":"second verifier passed"}'
: '{"passed":false,"reason":"try again"}';
},
}),
},
registry: storage,
logger,
});
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
ensureWorkspaceForCreate,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Keep trying until the verifier passes.",
cwd: workspaceDir,
model: "test-model",
verifyPrompt: "Report whether the loop has passed.",
archive: true,
sleepMs: 1,
maxIterations: 2,
});
await waitForLoopCompletion(service, loop.id);
const finalLoop = await service.inspectLoop(loop.id);
expect(finalLoop.status).toBe("succeeded");
expect(finalLoop.iterations).toHaveLength(2);
const firstWorker = await storage.get(finalLoop.iterations[0]!.workerAgentId!);
const firstVerifier = await storage.get(finalLoop.iterations[0]!.verifierAgentId!);
const secondWorker = await storage.get(finalLoop.iterations[1]!.workerAgentId!);
const secondVerifier = await storage.get(finalLoop.iterations[1]!.verifierAgentId!);
const workspaceId = firstWorker?.workspaceId;
expect(workspaceId).toMatch(/^wks_/);
expect(firstVerifier?.workspaceId).toBe(workspaceId);
expect(secondWorker?.workspaceId).toBe(workspaceId);
expect(secondVerifier?.workspaceId).toBe(workspaceId);
expect(await workspaceRegistry.get(workspaceId!)).toMatchObject({
workspaceId,
cwd: workspaceDir,
});
expect(await workspaceRegistry.list()).toHaveLength(1);
});
test("rejects non-directory cwd before minting a loop workspace", async () => {
const filePath = path.join(tmpDir, "not-a-directory.txt");
writeFileSync(filePath, "not a directory");
let ensureCalls = 0;
const manager = new AgentManager({ registry: storage, logger });
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
ensureWorkspaceForCreate: async () => {
ensureCalls += 1;
return "workspace-created-for-file-cwd";
},
});
await service.initialize();
await expect(
service.runLoop({
prompt: "Use a file as cwd",
cwd: filePath,
verifyChecks: ["true"],
}),
).rejects.toThrow("is not a directory");
expect(ensureCalls).toBe(0);
});
test("model-less loop workers use provider defaults and keep fast worker logs", async () => {
const workerConfigs: AgentSessionConfig[] = [];
const manager = new AgentManager({
clients: {
claude: new ScriptedAgentClient("claude", {
async onRun({ config }) {
if (config.title?.includes("worker")) {
workerConfigs.push(config);
return "worker default model output";
}
return '{"passed":true,"reason":"ok"}';
},
}),
},
registry: storage,
logger,
});
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Use provider default model",
cwd: workspaceDir,
verifyChecks: ["true"],
maxIterations: 1,
});
await waitForLoopCompletion(service, loop.id);
const finalLoop = await service.inspectLoop(loop.id);
expect(finalLoop.status).toBe("succeeded");
expect(workerConfigs).toHaveLength(1);
expect(workerConfigs[0]).toMatchObject({
provider: "claude",
model: undefined,
});
expect(
finalLoop.logs.some(
(entry) => entry.source === "worker" && entry.text.includes("worker default model output"),
),
).toBe(true);
});
test("archives worker and verifier agents after each iteration when requested", async () => {
const archivedAgentIds: string[] = [];
const manager = new AgentManager({
@@ -395,17 +599,18 @@ describe("LoopService", () => {
archivedAgentIds.push(agentId);
await archiveAgent(agentId);
};
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: NO_UNATTENDED_LOOP_POLICY,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Create done.txt",
cwd: workspaceDir,
model: "test-model",
verifyPrompt: "Confirm that done.txt exists in the workspace.",
archive: true,
maxIterations: 1,
@@ -432,6 +637,118 @@ describe("LoopService", () => {
});
});
test("worker prompt-start failures fail the loop and archive the worker", async () => {
class StartFailureLoopSession implements AgentSession {
readonly provider = "claude";
readonly capabilities = TEST_CAPABILITIES;
readonly id = randomUUID();
async run(): Promise<AgentRunResult> {
return {
sessionId: this.id,
finalText: "",
timeline: [],
};
}
async startTurn(): Promise<{ turnId: string }> {
throw new Error("worker failed before starting");
}
subscribe(): () => void {
return () => {};
}
async *streamHistory(): AsyncGenerator<AgentStreamEvent> {}
async getRuntimeInfo(): Promise<AgentRuntimeInfo> {
return {
provider: this.provider,
sessionId: this.id,
model: null,
modeId: null,
};
}
async getAvailableModes(): Promise<AgentMode[]> {
return [];
}
async getCurrentMode(): Promise<string | null> {
return null;
}
async setMode(): Promise<void> {}
getPendingPermissions(): AgentPermissionRequest[] {
return [];
}
async respondToPermission(): Promise<void> {}
describePersistence(): AgentPersistenceHandle {
return {
provider: this.provider,
sessionId: this.id,
};
}
async interrupt(): Promise<void> {}
async close(): Promise<void> {}
async listCommands(): Promise<AgentSlashCommand[]> {
return [];
}
}
const manager = new AgentManager({
clients: {
claude: {
provider: "claude",
capabilities: TEST_CAPABILITIES,
createSession: async () => new StartFailureLoopSession(),
resumeSession: async () => new StartFailureLoopSession(),
fetchCatalog: async () => ({ models: [], modes: [] }),
isAvailable: async () => true,
},
},
registry: storage,
logger,
});
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Fail before starting",
cwd: workspaceDir,
model: "test-model",
verifyChecks: ["true"],
archive: true,
maxIterations: 1,
});
await waitForLoopCompletion(service, loop.id);
const finalLoop = await service.inspectLoop(loop.id);
expect(finalLoop.status).toBe("failed");
expect(finalLoop.iterations[0]).toMatchObject({
status: "failed",
workerOutcome: "failed",
failureReason: expect.stringContaining("worker failed before starting"),
});
const workerAgentId = finalLoop.iterations[0]?.workerAgentId;
expect(workerAgentId).toMatch(/^[0-9a-f-]{36}$/);
expect(await storage.get(workerAgentId!)).toMatchObject({
archivedAt: expect.any(String),
});
});
test("uses verifier prompt when provided", async () => {
const manager = new AgentManager({
clients: {
@@ -452,17 +769,18 @@ describe("LoopService", () => {
registry: storage,
logger,
});
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: NO_UNATTENDED_LOOP_POLICY,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Create done.txt",
cwd: workspaceDir,
model: "test-model",
verifyPrompt: "Confirm that done.txt exists in the workspace.",
maxIterations: 1,
});
@@ -499,14 +817,22 @@ describe("LoopService", () => {
registry: storage,
logger,
});
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: {
async resolveCreateConfig(input) {
expect(input).toMatchObject({ parent: null, unattended: true, requestedMode: undefined });
return { modeId: "bypassPermissions", featureValues: input.featureValues };
expect(input).toMatchObject({
parent: null,
unattended: true,
requestedMode: undefined,
});
return {
modeId: input.unattended ? "bypassPermissions" : "interactive",
featureValues: input.featureValues,
};
},
},
});
@@ -515,6 +841,7 @@ describe("LoopService", () => {
const loop = await service.runLoop({
prompt: "Create done.txt",
cwd: workspaceDir,
model: "test-model",
verifyPrompt: "Confirm that done.txt exists in the workspace.",
maxIterations: 1,
});
@@ -554,16 +881,23 @@ describe("LoopService", () => {
registry: storage,
logger,
});
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: {
async resolveCreateConfig(input) {
expect(input).toMatchObject({ parent: null, unattended: true, requestedMode: undefined });
expect(input).toMatchObject({
parent: null,
unattended: true,
requestedMode: undefined,
});
return {
modeId: "build",
featureValues: { ...input.featureValues, auto_accept: true },
modeId: input.unattended ? "build" : "interactive",
featureValues: input.unattended
? { ...input.featureValues, auto_accept: true }
: input.featureValues,
};
},
},
@@ -574,6 +908,7 @@ describe("LoopService", () => {
prompt: "Create done.txt",
cwd: workspaceDir,
provider: "opencode",
model: "test-model",
verifyPrompt: "Confirm that done.txt exists in the workspace.",
maxIterations: 1,
});
@@ -610,17 +945,18 @@ describe("LoopService", () => {
registry: storage,
logger,
});
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: NO_UNATTENDED_LOOP_POLICY,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Create done.txt",
cwd: workspaceDir,
model: "test-model",
modeId: "acceptEdits",
verifierModeId: "plan",
verifyPrompt: "Confirm that done.txt exists in the workspace.",
@@ -659,17 +995,18 @@ describe("LoopService", () => {
cancelledAgentIds.push(agentId);
return cancelAgentRun(agentId);
};
const service = new LoopService({
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
providerSnapshotManager: NO_UNATTENDED_LOOP_POLICY,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Wait forever",
cwd: workspaceDir,
model: "test-model",
verifyChecks: ["test -f never.txt"],
});
@@ -695,6 +1032,116 @@ describe("LoopService", () => {
expect(cancelledAgentIds).toEqual([workerAgentId]);
expect(finalLoop.logs.some((entry) => entry.text.includes("Stop requested"))).toBe(true);
});
test("stops while waiting for loop workspace provisioning without starting a worker", async () => {
let resolveWorkspace: ((workspaceId: string) => void) | null = null;
const workspaceProvisioned = new Promise<string>((resolve) => {
resolveWorkspace = resolve;
});
const manager = new AgentManager({
clients: {
claude: new ScriptedAgentClient("claude", {
async onRun() {
return "worker should not start";
},
}),
},
registry: storage,
logger,
});
const createAgent = manager.createAgent.bind(manager);
let createAgentCalls = 0;
manager.createAgent = async (...args) => {
createAgentCalls += 1;
return createAgent(...args);
};
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
ensureWorkspaceForCreate: async () => workspaceProvisioned,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Stop before workspace is ready",
cwd: workspaceDir,
model: "test-model",
verifyPrompt: "Should not run.",
maxIterations: 1,
});
await waitForLoopIteration(service, loop.id);
const stopPromise = service.stopLoop(loop.id);
await waitForStopRequested(service, loop.id);
resolveWorkspace?.("workspace-created-after-stop");
const stopped = await stopPromise;
expect(stopped.status).toBe("stopped");
expect(createAgentCalls).toBe(0);
const finalLoop = await service.inspectLoop(loop.id);
expect(finalLoop.status).toBe("stopped");
expect(finalLoop.activeWorkerAgentId).toBeNull();
expect(finalLoop.iterations[0]).toMatchObject({
workerAgentId: null,
status: "stopped",
failureReason: "Loop stopped",
});
});
test("treats externally canceled worker turns as failures", async () => {
let release: (() => void) | null = null;
const blocker = new Promise<void>((resolve) => {
release = resolve;
});
const manager = new AgentManager({
clients: {
claude: new ScriptedAgentClient("claude", {
async onRun({ config }) {
if (config.title?.includes("worker")) {
await blocker;
return "finished";
}
return '{"passed":true,"reason":"should not verify canceled worker"}';
},
}),
},
registry: storage,
logger,
});
const service = createLoopService({
paseoHome,
agentManager: manager,
agentStorage: storage,
logger,
});
await service.initialize();
const loop = await service.runLoop({
prompt: "Wait until canceled",
cwd: workspaceDir,
model: "test-model",
verifyPrompt: "Should not run.",
maxIterations: 1,
});
const workerAgentId = await waitForActiveWorkerRun(service, manager, loop.id);
await manager.cancelAgentRun(workerAgentId);
release?.();
await waitForLoopCompletion(service, loop.id);
const finalLoop = await service.inspectLoop(loop.id);
expect(finalLoop.status).toBe("failed");
expect(finalLoop.iterations).toHaveLength(1);
expect(finalLoop.iterations[0]).toMatchObject({
workerAgentId,
workerOutcome: "failed",
status: "failed",
});
expect(finalLoop.iterations[0]?.failureReason).toContain("was canceled");
expect(finalLoop.iterations[0]?.verifierAgentId).toBeNull();
});
});
async function fsMkdir(target: string): Promise<void> {
@@ -728,6 +1175,30 @@ async function waitForActiveWorkerRun(
throw new Error("Timed out waiting for loop worker run to start");
}
async function waitForLoopIteration(service: LoopService, loopId: string): Promise<void> {
const deadline = Date.now() + 2_000;
while (Date.now() < deadline) {
const loop = await service.inspectLoop(loopId);
if (loop.iterations.length > 0) {
return;
}
await new Promise((resolve) => setTimeout(resolve, 10));
}
throw new Error("Timed out waiting for loop iteration to start");
}
async function waitForStopRequested(service: LoopService, loopId: string): Promise<void> {
const deadline = Date.now() + 2_000;
while (Date.now() < deadline) {
const loop = await service.inspectLoop(loopId);
if (loop.stopRequestedAt) {
return;
}
await new Promise((resolve) => setTimeout(resolve, 10));
}
throw new Error("Timed out waiting for loop stop request");
}
async function waitForCancelledAgent(
cancelledAgentIds: readonly string[],
agentId: string,

View File

@@ -5,11 +5,18 @@ import { z } from "zod";
import type { Logger } from "pino";
import { writeJsonFileAtomic } from "./atomic-file.js";
import { curateAgentActivity } from "./agent/activity-curator.js";
import {
type BoundCreateAgentCommand,
type EnsureWorkspaceForCreate,
formatProviderModel,
} from "./agent/create-agent/create.js";
import type { AgentManager } from "./agent/agent-manager.js";
import { getStructuredAgentResponse } from "./agent/agent-response-loop.js";
import {
buildStructuredAgentResponsePrompt,
getStructuredAgentResponse,
} from "./agent/agent-response-loop.js";
import type {
AgentPromptInput,
AgentSessionConfig,
AgentStreamEvent,
AgentTimelineItem,
AgentProvider,
@@ -19,11 +26,6 @@ import {
createStringCommandShellEnvOverlay,
} from "../utils/string-command-shell.js";
import { execCommand } from "../utils/spawn.js";
import type {
ProviderSnapshotManager,
ResolvedProviderCreateConfig,
ResolveProviderCreateConfigOptions,
} from "./agent/provider-snapshot-manager.js";
const LOOP_ID_LENGTH = 8;
const DEFAULT_LOOP_PROVIDER: AgentProvider = "claude";
@@ -205,6 +207,21 @@ function ensureNonNegativeInteger(value: number | undefined, field: string): num
return value;
}
async function assertLoopCwdDirectory(cwd: string): Promise<void> {
let stats;
try {
stats = await fs.stat(cwd);
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") {
throw new Error(`Working directory ${cwd} no longer exists`, { cause: error });
}
throw error;
}
if (!stats.isDirectory()) {
throw new Error(`Working directory ${cwd} is not a directory`);
}
}
function buildWorkerTitle(loop: LoopRecord, iterationIndex: number): string {
const prefix = loop.name ?? loop.id;
return `${prefix} [loop ${iterationIndex} worker]`;
@@ -215,7 +232,14 @@ function buildVerifierTitle(loop: LoopRecord, iterationIndex: number): string {
return `${prefix} [loop ${iterationIndex} verifier]`;
}
type CreateConfigResolver = Pick<ProviderSnapshotManager, "resolveCreateConfig">;
type LoopAgentManager = Pick<
AgentManager,
"archiveAgent" | "cancelAgentRun" | "closeAgent" | "runAgent" | "subscribe" | "waitForAgentEvent"
>;
interface LoopExecutionContext {
workspaceId: Promise<string>;
}
function formatStreamLog(event: AgentStreamEvent): string | null {
switch (event.type) {
@@ -318,9 +342,10 @@ export class LoopService {
constructor(
private readonly options: {
paseoHome: string;
agentManager: AgentManager;
agentManager: LoopAgentManager;
logger: Logger;
providerSnapshotManager: CreateConfigResolver;
createAgent: BoundCreateAgentCommand;
ensureWorkspaceForCreate: EnsureWorkspaceForCreate;
},
) {
this.storePath = path.join(options.paseoHome, "loops", "loops.json");
@@ -389,13 +414,15 @@ export class LoopService {
if (!verifyPrompt && verifyChecks.length === 0) {
throw new Error("Loop requires --verify or at least one --verify-check");
}
const cwd = path.resolve(input.cwd);
await assertLoopCwdDirectory(cwd);
const createdAt = nowIso();
const record = LoopRecordSchema.parse({
id: createLoopId(),
name: normalizeName(input.name),
prompt,
cwd: path.resolve(input.cwd),
cwd,
provider: input.provider ?? DEFAULT_LOOP_PROVIDER,
model: normalizePrompt(input.model, "model"),
modeId: normalizePrompt(input.modeId, "modeId"),
@@ -489,10 +516,12 @@ export class LoopService {
level: "info",
text: "Stop requested.",
});
if (running) {
running.abortController.abort(new Error("Loop aborted"));
}
await this.persist();
if (running) {
running.abortController.abort(new Error("Loop aborted"));
if (loop.activeWorkerAgentId) {
await this.options.agentManager.cancelAgentRun(loop.activeWorkerAgentId).catch(() => {});
}
@@ -513,6 +542,11 @@ export class LoopService {
private async executeLoop(loopId: string, signal: AbortSignal): Promise<void> {
const loop = this.requireLoop(loopId);
const deadline = loop.maxTimeMs ? Date.now() + loop.maxTimeMs : null;
const workspaceId = this.options.ensureWorkspaceForCreate(loop.cwd, { prompt: loop.prompt });
workspaceId.catch(() => {});
const context: LoopExecutionContext = {
workspaceId,
};
try {
for (let index = 1; ; index += 1) {
@@ -551,14 +585,14 @@ export class LoopService {
});
await this.persist();
const workerPassed = await this.runWorkerIteration(loop, iteration, signal);
const workerPassed = await this.runWorkerIteration(loop, iteration, signal, context);
if (signal.aborted) {
throw new Error("Loop aborted");
}
if (!workerPassed) {
iteration.status = iteration.status === "stopped" ? "stopped" : "failed";
} else {
const verificationPassed = await this.runVerification(loop, iteration, signal);
const verificationPassed = await this.runVerification(loop, iteration, signal, context);
if (verificationPassed) {
iteration.status = "succeeded";
this.finishLoop(loop, "succeeded", `Iteration ${index} passed verification.`);
@@ -596,11 +630,11 @@ export class LoopService {
loopId: string,
error: unknown,
): Promise<void> {
const iteration = loop.activeIteration
? loop.iterations.find((candidate) => candidate.index === loop.activeIteration)
: null;
if (isAbortError(error)) {
this.finishLoop(loop, "stopped", "Loop stopped.");
const iteration = loop.activeIteration
? loop.iterations.find((candidate) => candidate.index === loop.activeIteration)
: null;
if (iteration && iteration.status === "running") {
iteration.status = "stopped";
iteration.failureReason = "Loop stopped";
@@ -613,9 +647,6 @@ export class LoopService {
const message = error instanceof Error ? error.message : String(error);
this.logger.error({ err: error, loopId }, "Loop execution failed");
this.finishLoop(loop, "failed", message);
const iteration = loop.activeIteration
? loop.iterations.find((candidate) => candidate.index === loop.activeIteration)
: null;
if (iteration && iteration.status === "running") {
iteration.status = "failed";
iteration.failureReason = message;
@@ -628,20 +659,40 @@ export class LoopService {
loop: LoopRecord,
iteration: LoopIterationRecord,
signal: AbortSignal,
context: LoopExecutionContext,
): Promise<boolean> {
const agent = await this.options.agentManager.createAgent(
await this.buildWorkerConfig(loop, iteration),
);
const workspaceId = await context.workspaceId;
if (signal.aborted) {
throw new Error("Loop aborted");
}
const created = await this.options.createAgent({
kind: "mcp",
provider: this.formatWorkerProviderModel(loop),
cwd: loop.cwd,
workspaceId,
title: buildWorkerTitle(loop, iteration.index),
mode: loop.modeId ?? undefined,
unattended: true,
promptFailure: "return-error",
background: true,
notifyOnFinish: false,
internal: true,
});
const agent = created.snapshot;
iteration.workerAgentId = agent.id;
loop.activeWorkerAgentId = agent.id;
loop.updatedAt = nowIso();
await this.persist();
let workerCanceledReason: string | null = null;
const unsubscribe = this.options.agentManager.subscribe(
(event) => {
if (event.type !== "agent_stream") {
return;
}
if (event.event.type === "turn_canceled") {
workerCanceledReason = event.event.reason;
}
const text = formatStreamLog(event.event);
if (!text) {
return;
@@ -658,14 +709,11 @@ export class LoopService {
);
try {
const prompt = this.toPrompt(loop.prompt);
const result = await this.options.agentManager.runAgent(agent.id, prompt);
const result = await this.options.agentManager.runAgent(agent.id, this.toPrompt(loop.prompt));
iteration.workerCompletedAt = nowIso();
iteration.workerOutcome = result.canceled ? "canceled" : "completed";
if (result.canceled) {
iteration.failureReason = "Worker run was canceled.";
iteration.status = "stopped";
return false;
iteration.workerOutcome = "completed";
if (result.canceled || workerCanceledReason) {
throw new Error(`Loop worker ${agent.id} was canceled: ${workerCanceledReason}`);
}
return true;
} catch (error) {
@@ -703,6 +751,7 @@ export class LoopService {
loop: LoopRecord,
iteration: LoopIterationRecord,
signal: AbortSignal,
context: LoopExecutionContext,
): Promise<boolean> {
for (const command of loop.verifyChecks) {
if (signal.aborted) {
@@ -736,9 +785,25 @@ export class LoopService {
}
const startedAt = nowIso();
const verifierAgent = await this.options.agentManager.createAgent(
await this.buildVerifierConfig(loop, iteration),
);
const initialVerifierPrompt = buildStructuredAgentResponsePrompt({
prompt: loop.verifyPrompt,
schema: LoopVerifyPromptSchema,
schemaName: "LoopVerifierResult",
});
const created = await this.options.createAgent({
kind: "mcp",
provider: this.formatVerifierProviderModel(loop),
cwd: loop.cwd,
workspaceId: await context.workspaceId,
title: buildVerifierTitle(loop, iteration.index),
mode: loop.verifierModeId ?? loop.modeId ?? undefined,
unattended: true,
promptFailure: "return-error",
background: true,
notifyOnFinish: false,
internal: true,
});
const verifierAgent = created.snapshot;
iteration.verifierAgentId = verifierAgent.id;
loop.activeVerifierAgentId = verifierAgent.id;
loop.updatedAt = nowIso();
@@ -765,8 +830,17 @@ export class LoopService {
);
try {
let waitingForInitialResponse = true;
const result = await getStructuredAgentResponse({
caller: async (nextPrompt) => {
if (waitingForInitialResponse) {
waitingForInitialResponse = false;
const run = await this.options.agentManager.runAgent(
verifierAgent.id,
initialVerifierPrompt,
);
return this.resolveFinalText(run.timeline, run.finalText);
}
const run = await this.options.agentManager.runAgent(
verifierAgent.id,
this.toPrompt(nextPrompt),
@@ -812,56 +886,18 @@ export class LoopService {
}
}
private async buildWorkerConfig(
loop: LoopRecord,
iteration: LoopIterationRecord,
): Promise<AgentSessionConfig> {
const provider = loop.workerProvider ?? loop.provider;
const resolvedUnattendedConfig = loop.modeId
? { modeId: loop.modeId, featureValues: undefined }
: await this.resolveProviderCreateConfig({ provider, cwd: loop.cwd });
return {
provider,
cwd: loop.cwd,
model: loop.workerModel ?? loop.model ?? undefined,
modeId: resolvedUnattendedConfig.modeId,
featureValues: resolvedUnattendedConfig.featureValues,
title: buildWorkerTitle(loop, iteration.index),
internal: true,
};
private formatWorkerProviderModel(loop: LoopRecord): string {
return formatProviderModel(
loop.workerProvider ?? loop.provider,
loop.workerModel ?? loop.model,
);
}
private async buildVerifierConfig(
loop: LoopRecord,
iteration: LoopIterationRecord,
): Promise<AgentSessionConfig> {
const provider = loop.verifierProvider ?? loop.provider;
const explicitModeId = loop.verifierModeId ?? loop.modeId;
const resolvedUnattendedConfig = explicitModeId
? { modeId: explicitModeId, featureValues: undefined }
: await this.resolveProviderCreateConfig({ provider, cwd: loop.cwd });
return {
provider,
cwd: loop.cwd,
model: loop.verifierModel ?? loop.model ?? undefined,
modeId: resolvedUnattendedConfig.modeId,
featureValues: resolvedUnattendedConfig.featureValues,
title: buildVerifierTitle(loop, iteration.index),
internal: true,
};
}
private resolveProviderCreateConfig(
input: Pick<ResolveProviderCreateConfigOptions, "provider" | "cwd">,
): Promise<ResolvedProviderCreateConfig> {
return this.options.providerSnapshotManager.resolveCreateConfig({
provider: input.provider,
cwd: input.cwd,
requestedMode: undefined,
featureValues: undefined,
parent: null,
unattended: true,
});
private formatVerifierProviderModel(loop: LoopRecord): string {
return formatProviderModel(
loop.verifierProvider ?? loop.provider,
loop.verifierModel ?? loop.model,
);
}
private resolveFinalText(timeline: AgentTimelineItem[], finalText: string): string {

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@@ -47,7 +47,7 @@ describe("ScheduleStore", () => {
expect(listed).toEqual([created]);
});
test("put round-trips an updated schedule to disk", async () => {
test("update round-trips an updated schedule to disk", async () => {
const created = await store.create({
name: "before",
prompt: "before",
@@ -79,7 +79,7 @@ describe("ScheduleStore", () => {
nextRunAt: "2026-01-01T09:00:00.000Z",
updatedAt: "2026-01-01T00:00:30.000Z",
};
await store.put(updated);
await store.update(created.id, () => updated);
const reloaded = await new ScheduleStore(tempDir).get(created.id);
expect(reloaded).toEqual(updated);
@@ -113,4 +113,195 @@ describe("ScheduleStore", () => {
expect(await store.get(created.id)).toBeNull();
expect(await store.list()).toEqual([]);
});
test("serializes concurrent updates on one schedule without losing writes", async () => {
const created = await store.create({
name: "before",
prompt: "before",
cadence: { type: "every", everyMs: 60_000 },
target: {
type: "new-agent",
config: { provider: "claude", cwd: tempDir },
},
status: "active",
createdAt: "2026-01-01T00:00:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
nextRunAt: "2026-01-01T00:01:00.000Z",
lastRunAt: null,
pausedAt: null,
expiresAt: null,
maxRuns: null,
runs: [],
});
let releaseFirstUpdate: (() => void) | null = null;
const firstUpdateBlocked = new Promise<void>((resolve) => {
releaseFirstUpdate = resolve;
});
let firstUpdaterEntered: (() => void) | null = null;
const firstUpdaterStarted = new Promise<void>((resolve) => {
firstUpdaterEntered = resolve;
});
let secondSawRunCount = -1;
const firstUpdate = store.update(created.id, async (schedule) => {
firstUpdaterEntered?.();
await firstUpdateBlocked;
return {
...schedule,
runs: [
...schedule.runs,
{
id: "run-1",
scheduledFor: "2026-01-01T00:01:00.000Z",
startedAt: "2026-01-01T00:01:00.000Z",
endedAt: null,
status: "running" as const,
agentId: null,
output: null,
error: null,
},
],
};
});
await firstUpdaterStarted;
const secondUpdate = store.update(created.id, (schedule) => {
secondSawRunCount = schedule.runs.length;
return {
...schedule,
prompt: "after",
};
});
releaseFirstUpdate?.();
const [, second] = await Promise.all([firstUpdate, secondUpdate]);
expect(secondSawRunCount).toBe(1);
expect(second).toMatchObject({
prompt: "after",
runs: [{ id: "run-1" }],
});
await expect(new ScheduleStore(tempDir).get(created.id)).resolves.toMatchObject({
prompt: "after",
runs: [{ id: "run-1" }],
});
});
test("revalidates a named target match after waiting for the schedule update queue", async () => {
class GatedListScheduleStore extends ScheduleStore {
private listGate: {
entered: () => void;
release: Promise<void>;
} | null = null;
gateNextList(gate: { entered: () => void; release: Promise<void> }): void {
this.listGate = gate;
}
override async list() {
const schedules = await super.list();
const gate = this.listGate;
if (gate) {
this.listGate = null;
gate.entered();
await gate.release;
}
return schedules;
}
}
const gatedStore = new GatedListScheduleStore(tempDir);
const target = {
type: "new-agent" as const,
config: { provider: "claude" as const, cwd: tempDir },
};
const created = await gatedStore.create({
name: "race",
prompt: "before",
cadence: { type: "every", everyMs: 60_000 },
target,
status: "active",
createdAt: "2026-01-01T00:00:00.000Z",
updatedAt: "2026-01-01T00:00:00.000Z",
nextRunAt: "2026-01-01T00:01:00.000Z",
lastRunAt: null,
pausedAt: null,
expiresAt: null,
maxRuns: null,
runs: [],
});
let releaseCompletion: (() => void) | null = null;
const completionBlocked = new Promise<void>((resolve) => {
releaseCompletion = resolve;
});
let completionEntered: (() => void) | null = null;
const completionStarted = new Promise<void>((resolve) => {
completionEntered = resolve;
});
const completeOriginal = gatedStore.update(created.id, async (schedule) => {
completionEntered?.();
await completionBlocked;
return {
...schedule,
status: "completed" as const,
nextRunAt: null,
updatedAt: "2026-01-01T00:00:30.000Z",
};
});
await completionStarted;
let releaseUpsertList: (() => void) | null = null;
const upsertListBlocked = new Promise<void>((resolve) => {
releaseUpsertList = resolve;
});
let upsertListEntered: (() => void) | null = null;
const upsertListed = new Promise<void>((resolve) => {
upsertListEntered = resolve;
});
gatedStore.gateNextList({
entered: () => upsertListEntered?.(),
release: upsertListBlocked,
});
const upsert = gatedStore.upsertByNameAndTarget("race", target, {
create: () => ({
name: "race",
prompt: "after",
cadence: { type: "every", everyMs: 60_000 },
target,
status: "active",
createdAt: "2026-01-01T00:01:00.000Z",
updatedAt: "2026-01-01T00:01:00.000Z",
nextRunAt: "2026-01-01T00:02:00.000Z",
lastRunAt: null,
pausedAt: null,
expiresAt: null,
maxRuns: null,
runs: [],
}),
update: () => {
throw new Error("stale identity match should not update");
},
});
await upsertListed;
releaseCompletion?.();
await completeOriginal;
releaseUpsertList?.();
const upserted = await upsert;
expect(upserted.id).not.toBe(created.id);
expect(upserted).toMatchObject({
name: "race",
prompt: "after",
status: "active",
});
await expect(gatedStore.get(created.id)).resolves.toMatchObject({
status: "completed",
prompt: "before",
});
expect(await gatedStore.list()).toHaveLength(2);
});
});

View File

@@ -1,14 +1,97 @@
import { randomBytes } from "node:crypto";
import { mkdir, readFile, readdir, rm } from "node:fs/promises";
import { join } from "node:path";
import { StoredScheduleSchema, type StoredSchedule } from "@getpaseo/protocol/schedule/types";
import {
StoredScheduleSchema,
type ScheduleTarget,
type StoredSchedule,
} from "@getpaseo/protocol/schedule/types";
import { writeJsonFileAtomic } from "../atomic-file.js";
function generateScheduleId(): string {
return randomBytes(4).toString("hex");
}
type ScheduleUpdater = (schedule: StoredSchedule) => StoredSchedule | Promise<StoredSchedule>;
interface ScheduleNameTargetUpsert {
create: () => Omit<StoredSchedule, "id"> | Promise<Omit<StoredSchedule, "id">>;
update: ScheduleUpdater;
}
function canonicalize(value: unknown): unknown {
if (Array.isArray(value)) {
return value.map(canonicalize);
}
if (value && typeof value === "object") {
const source = value as Record<string, unknown>;
return Object.fromEntries(
Object.keys(source)
.sort()
.map((key) => [key, canonicalize(source[key])]),
);
}
return value;
}
function normalizeScheduleName(name: string): string {
const trimmed = name.trim();
if (!trimmed) {
throw new Error("Schedule name is required");
}
return trimmed;
}
function normalizeOptionalScheduleName(name: string | null): string | null {
if (name === null) {
return null;
}
const trimmed = name.trim();
return trimmed ? trimmed : null;
}
function targetIdentity(target: ScheduleTarget): unknown {
if (target.type === "agent") {
return {
type: target.type,
agentId: target.agentId,
};
}
const { workspaceId: _workspaceId, ...config } = target.config;
return {
type: target.type,
config,
};
}
function nameTargetIdentityKey(name: string, target: ScheduleTarget): string {
return JSON.stringify(
canonicalize({
name: normalizeScheduleName(name),
target: targetIdentity(target),
}),
);
}
function matchesNameAndTarget(
schedule: StoredSchedule,
name: string,
target: ScheduleTarget,
): boolean {
const scheduleName = normalizeOptionalScheduleName(schedule.name);
return (
schedule.status !== "completed" &&
scheduleName !== null &&
scheduleName === normalizeScheduleName(name) &&
nameTargetIdentityKey(scheduleName, schedule.target) === nameTargetIdentityKey(name, target)
);
}
export class ScheduleStore {
private readonly scheduleMutations = new Map<string, Promise<unknown>>();
private readonly identityMutations = new Map<string, Promise<unknown>>();
constructor(private readonly dir: string) {}
private filePath(id: string): string {
@@ -47,18 +130,125 @@ export class ScheduleStore {
}
async create(schedule: Omit<StoredSchedule, "id">): Promise<StoredSchedule> {
const created = { ...schedule, id: generateScheduleId() };
await this.put(created);
const created = StoredScheduleSchema.parse({ ...schedule, id: generateScheduleId() });
await this.write(created);
return created;
}
async put(schedule: StoredSchedule): Promise<void> {
async update(id: string, updater: ScheduleUpdater): Promise<StoredSchedule | null> {
return this.serializeScheduleMutation(id, async () => {
const current = await this.get(id);
if (!current) {
return null;
}
const next = await updater(current);
if (next === current) {
return current;
}
if (next.id !== id) {
throw new Error(`Schedule update cannot change id: ${id}`);
}
const updated = StoredScheduleSchema.parse(next);
await this.write(updated);
return updated;
});
}
async upsertByNameAndTarget(
name: string,
target: ScheduleTarget,
options: ScheduleNameTargetUpsert,
): Promise<StoredSchedule> {
const identity = nameTargetIdentityKey(name, target);
return this.serializeIdentityMutation(identity, async () => {
while (true) {
const existing = (await this.list()).find((schedule) =>
matchesNameAndTarget(schedule, name, target),
);
if (!existing) {
const created = StoredScheduleSchema.parse({
...(await options.create()),
id: generateScheduleId(),
});
if (!matchesNameAndTarget(created, name, target)) {
throw new Error("Created schedule does not match requested identity");
}
await this.write(created);
return created;
}
const updated = await this.updateMatchedSchedule(existing.id, name, target, options.update);
if (updated) {
return updated;
}
}
});
}
private async write(schedule: StoredSchedule): Promise<void> {
await this.ensureDir();
await writeJsonFileAtomic(this.filePath(schedule.id), schedule);
}
async delete(id: string): Promise<void> {
await this.ensureDir();
await rm(this.filePath(id), { force: true });
await this.serializeScheduleMutation(id, async () => {
await this.ensureDir();
await rm(this.filePath(id), { force: true });
});
}
private async serializeScheduleMutation<T>(
scheduleId: string,
mutation: () => Promise<T>,
): Promise<T> {
return this.serializeMutation(this.scheduleMutations, scheduleId, mutation);
}
private async serializeIdentityMutation<T>(
identity: string,
mutation: () => Promise<T>,
): Promise<T> {
return this.serializeMutation(this.identityMutations, identity, mutation);
}
private async serializeMutation<T>(
promises: Map<string, Promise<unknown>>,
key: string,
mutation: () => Promise<T>,
): Promise<T> {
const previous = promises.get(key) ?? Promise.resolve();
const next = previous.catch(() => undefined).then(mutation);
promises.set(key, next);
try {
return await next;
} finally {
if (promises.get(key) === next) {
promises.delete(key);
}
}
}
private async updateMatchedSchedule(
id: string,
name: string,
target: ScheduleTarget,
updater: ScheduleUpdater,
): Promise<StoredSchedule | null> {
return this.serializeScheduleMutation(id, async () => {
const current = await this.get(id);
if (!current || !matchesNameAndTarget(current, name, target)) {
return null;
}
const next = await updater(current);
if (next.id !== id) {
throw new Error(`Schedule update cannot change id: ${id}`);
}
const updated = StoredScheduleSchema.parse(next);
if (!matchesNameAndTarget(updated, name, target)) {
throw new Error("Updated schedule does not match requested identity");
}
await this.write(updated);
return updated;
});
}
}

View File

@@ -29,11 +29,15 @@ describe("snapshot mutation ownership boundary", () => {
const cwd = mkdtempSync(path.join(os.tmpdir(), "snapshot-owner-live-"));
try {
const snapshot = await daemonHandle.daemon.agentManager.createAgent({
provider: "codex",
cwd,
model: "gpt-5.2-codex",
});
const snapshot = await daemonHandle.daemon.agentManager.createAgent(
{
provider: "codex",
cwd,
model: "gpt-5.2-codex",
},
undefined,
{ workspaceId: undefined },
);
await daemonHandle.daemon.agentManager.flush();
const applySnapshotSpy = vi.spyOn(daemonHandle.daemon.agentStorage, "applySnapshot");