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
5 changes: 5 additions & 0 deletions .changeset/fork-cron-clear.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@moonshot-ai/kimi-code": patch
---

Cron tasks from the source session no longer fire inside a forked session.
4 changes: 2 additions & 2 deletions apps/vis/server/src/lib/agent-record-types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,12 +33,12 @@ import type {
ExportSessionManifest,
FileHistoryCheckpointed,
FileHistoryTracked,
Forked,
FullCompactionBegin,
FullCompactionCancel,
FullCompactionComplete,
GoalClear,
GoalCreate,
GoalForked,
GoalUpdate,
InteractionRequestEvent,
InteractionResolvedEvent,
Expand Down Expand Up @@ -174,7 +174,7 @@ export type AgentRecord =
| WireRecordOf<'cron.delete', CronDeletePayload>
| WireRecordOf<'file_history.checkpoint', FileHistoryCheckpointed>
| WireRecordOf<'file_history.tracked', FileHistoryTracked>
| WireRecordOf<'forked', GoalForked>
| WireRecordOf<'forked', Forked>
| WireRecordOf<'full_compaction.begin', FullCompactionBegin>
| WireRecordOf<'full_compaction.cancel', FullCompactionCancel>
| WireRecordOf<'full_compaction.complete', FullCompactionComplete>
Expand Down
4 changes: 2 additions & 2 deletions packages/agent-core-v2/docs/wire-manifest.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@
// cron.delete (none) src/features/cron/cronOps.ts
// file_history.checkpoint fileHistory src/features/fileHistory/fileHistoryOps.ts
// file_history.tracked fileHistory src/features/fileHistory/fileHistoryOps.ts
// forked (none) src/features/goal/goalOps.ts
// forked (none) src/session/agentLifecycle/forked.ts
// full_compaction.begin fullCompaction src/agent/fullCompaction/compactionOps.ts
// full_compaction.cancel fullCompaction src/agent/fullCompaction/compactionOps.ts
// full_compaction.complete fullCompaction src/agent/fullCompaction/compactionOps.ts
Expand Down Expand Up @@ -248,7 +248,7 @@ interface FileHistoryTrackedPayload {

/**
* states: (none)
* owner: src/features/goal/goalOps.ts
* owner: src/session/agentLifecycle/forked.ts
*/
interface ForkedPayload {
_name: 'forked';
Expand Down
5 changes: 4 additions & 1 deletion packages/agent-core-v2/src/features/cron/cronOps.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,10 @@ import type { CronJobOrigin } from '#/agent/contextMemory/types';
import type { CronTask } from '#/features/cron/cronTask';
import { Event2 } from '#/app/event/event2';

export type CronModelState = Map<string, CronTask>;
export interface CronModelState {
readonly tasks: Map<string, CronTask>;
readonly forkNotice: { reminderPending: boolean };
}

const cronTaskSchema = z.object({
id: z.string(),
Expand Down
78 changes: 58 additions & 20 deletions packages/agent-core-v2/src/features/cron/cronService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,14 +23,19 @@ import type { CronDeletedEvent, CronScheduledEvent } from '#/app/telemetry/event
import { ITelemetryService } from '#/app/telemetry/telemetry';
import { BugIndicatingError } from '#/errors';
import type { ContentPart } from '#human/llm/message';
import { ContextAppendMessage } from '#/agent/contextMemory/contextEvents';
import type { ContextMessage } from '#/agent/contextMemory/types';
import { IAgentReminderService } from '#/features/reminder/reminderService';
import { MAIN_AGENT_ID } from '#/session/agentLifecycle/agentLifecycle';
import { Forked } from '#/session/agentLifecycle/forked';
import { IEventDispatcher } from '#/state/eventDispatcher';

import { CronAdd, CronCursor, CronDelete, CronFired, type CronModelState } from './cronOps';

registerEvent2Class(CronAdd);
registerEvent2Class(CronDelete);
registerEvent2Class(CronCursor);
registerEvent2Class(Forked);

const STALE_THRESHOLD_MS = 7 * 24 * 60 * 60 * 1000;
const DEFAULT_POLL_INTERVAL_MS = 1_000;
Expand All @@ -43,14 +48,27 @@ export const CRON_FIRED = 'cron_fired' as const;
export const CRON_MISSED = 'cron_missed' as const;
export const CRON_DELETED = 'cron_deleted' as const;

const CRON_FORK_CLEARED_REMINDER = [
'This fork does not have any scheduled cron tasks.',
'Tasks from the source session continue to run in the source session.',
'Create new tasks here if needed.',
].join(' ');

const CRON_FORK_CLEARED_REMINDER_NAME = 'cron_fork_cleared';

function isCronForkClearedReminder(message: ContextMessage): boolean {
const origin = message.origin;
return origin?.kind === 'injection' && origin.variant === CRON_FORK_CLEARED_REMINDER_NAME;
}

interface CronActorContext {
readonly tasks: CronModelState;
readonly model: CronModelState;
readonly runtime: AgentActorContext<CronModelState>;
}

interface CronCommitEvent {
readonly type: 'cron.commit';
readonly tasks: CronModelState;
readonly model: CronModelState;
}

interface CronTickEvent {
Expand Down Expand Up @@ -154,7 +172,7 @@ function removeTasks(
runtime: AgentActorContext<CronModelState>,
ids: readonly string[],
): readonly string[] {
const removed = ids.filter((id) => runtime.getState().has(id));
const removed = ids.filter((id) => runtime.getState().tasks.has(id));
if (removed.length > 0) void runtime.dispatch(new CronDelete({ ids: removed }));
return removed;
}
Expand Down Expand Up @@ -254,7 +272,7 @@ async function processDue(
}
const advancedTo = lastDueMs ?? now;
state.lastSeenAt.set(task.id, advancedTo);
if (runtime.getState().has(task.id)) {
if (runtime.getState().tasks.has(task.id)) {
void runtime.dispatch(new CronCursor({ id: task.id, lastFiredAt: advancedTo }));
}
}
Expand All @@ -264,10 +282,10 @@ async function tickCron(
state: CronEffectState,
): Promise<void> {
await configOf(runtime).ready;
if (cronConfigOf(runtime).disabled || runtime.getState().size === 0) return;
if (cronConfigOf(runtime).disabled || runtime.getState().tasks.size === 0) return;
if (runtime.get(IAgentLoopService).snapshot().state === 'running') return;
const now = state.clocks.wallNow();
await Promise.all([...runtime.getState().values()].map((task) => processDue(runtime, state, task, now)));
await Promise.all([...runtime.getState().tasks.values()].map((task) => processDue(runtime, state, task, now)));
}

const cronEffects = fromCallback(({
Expand All @@ -283,6 +301,11 @@ const cronEffects = fromCallback(({
sendBack: (event: CronActorEvent) => void;
}) => {
if (input.runtime.agent.agentId !== MAIN_AGENT_ID) return;
if (input.runtime.getState().forkNotice.reminderPending) {
input.runtime.get(IAgentReminderService).notify(CRON_FORK_CLEARED_REMINDER, {
variant: CRON_FORK_CLEARED_REMINDER_NAME,
});
}
const timer = new IntervalTimer({ unref: true });
const state: CronEffectState = {
clocks: SYSTEM_CLOCKS,
Expand Down Expand Up @@ -356,7 +379,10 @@ const cronActorLogic = setup({
},
actors: { cronEffects },
}).createMachine({
context: ({ input }) => ({ tasks: new Map(), runtime: input }),
context: ({ input }) => ({
model: { tasks: new Map(), forkNotice: { reminderPending: false } },
runtime: input,
}),
initial: 'beforeRestore',
states: {
beforeRestore: {
Expand All @@ -383,7 +409,7 @@ const cronActorLogic = setup({
},
on: {
'cron.commit': {
actions: assign({ tasks: ({ event }) => event.tasks }),
actions: assign({ model: ({ event }) => event.model }),
},
},
});
Expand Down Expand Up @@ -429,24 +455,36 @@ export class AgentCronService extends AgentActorService<CronModelState> implemen
this.actor = this.attachActor(cronActorLogic, {
id: 'cron',
durable: {
events: [CronAdd, CronDelete, CronCursor],
events: [CronAdd, CronDelete, CronCursor, Forked, ContextAppendMessage],
undoable: false,
transition: (state, event) => {
if (event instanceof CronAdd) {
state.set(event.task.id, event.task);
state.tasks.set(event.task.id, event.task);
return;
}
if (event instanceof CronDelete) {
for (const id of event.ids) state.delete(id);
for (const id of event.ids) state.tasks.delete(id);
return;
}
if (event instanceof CronCursor) {
const task = state.get(event.id);
if (task !== undefined) state.set(event.id, { ...task, lastFiredAt: event.lastFiredAt });
const task = state.tasks.get(event.id);
if (task !== undefined) state.tasks.set(event.id, { ...task, lastFiredAt: event.lastFiredAt });
return;
}
if (event instanceof Forked) {
state.forkNotice.reminderPending =
state.tasks.size > 0 || state.forkNotice.reminderPending;
state.tasks.clear();
return;
}
if (event instanceof ContextAppendMessage) {
if (state.forkNotice.reminderPending && isCronForkClearedReminder(event.message)) {
state.forkNotice.reminderPending = false;
}
}
},
read: (snapshot) => (snapshot as CronActorSnapshot).context.tasks,
commit: (actor, tasks) => { actor.send({ type: 'cron.commit', tasks }); },
read: (snapshot) => (snapshot as CronActorSnapshot).context.model,
commit: (actor, model) => { actor.send({ type: 'cron.commit', model }); },
},
});
}
Expand All @@ -460,7 +498,7 @@ export class AgentCronService extends AgentActorService<CronModelState> implemen
}

addTask(init: CronTaskInit): CronTask {
const tasks = this.actor.getState();
const tasks = this.actor.getState().tasks;
let id: string | undefined;
for (let attempt = 0; attempt < MAX_ID_ATTEMPTS; attempt += 1) {
const candidate = ulid();
Expand All @@ -482,11 +520,11 @@ export class AgentCronService extends AgentActorService<CronModelState> implemen
}

getTask(id: string): CronTask | undefined {
return this.actor.getState().get(id);
return this.actor.getState().tasks.get(id);
}

list(): readonly CronTask[] {
return [...this.actor.getState().values()];
return [...this.actor.getState().tasks.values()];
}

isStale(task: CronTask): boolean {
Expand All @@ -495,15 +533,15 @@ export class AgentCronService extends AgentActorService<CronModelState> implemen

getNextFireTime(): number | null {
let min: number | null = null;
for (const task of this.actor.getState().values()) {
for (const task of this.actor.getState().tasks.values()) {
const next = nextFireFor(this.actor, task);
if (next !== null && (min === null || next < min)) min = next;
}
return min;
}

getNextFireForTask(taskId: string): number | null {
const task = this.actor.getState().get(taskId);
const task = this.actor.getState().tasks.get(taskId);
return task === undefined ? null : nextFireFor(this.actor, task);
}

Expand Down
11 changes: 0 additions & 11 deletions packages/agent-core-v2/src/features/goal/goalOps.ts
Original file line number Diff line number Diff line change
Expand Up @@ -111,17 +111,6 @@ export interface GoalClear {
readonly agentId: string;
}

const goalForkedSchema = z.object({ agentId: z.string() });

export class GoalForked extends AgentEvent2<z.infer<typeof goalForkedSchema>> {
static override readonly type = 'forked';
static override readonly durable = true;
static override readonly schema = goalForkedSchema;
}
export interface GoalForked {
readonly agentId: string;
}

export interface GoalUpdatedPayload {
readonly agentId: string;
snapshot: GoalSnapshot | null;
Expand Down
8 changes: 4 additions & 4 deletions packages/agent-core-v2/src/features/goal/goalService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ import {
type KimiErrorPayload,
} from '#/errors';
import { IAgentLifecycleService, MAIN_AGENT_ID } from '#/session/agentLifecycle/agentLifecycle';
import { Forked } from '#/session/agentLifecycle/forked';
import { ISessionUsageService } from '#/session/usage/sessionUsage';
import { IEventDispatcher } from '#/state/eventDispatcher';
import type { ExecutableToolResult } from '#/tool/toolContract';
Expand All @@ -58,7 +59,6 @@ import { IGoalDeadlineScheduler } from './goalDeadlineScheduler';
import {
GoalClear,
GoalCreate,
GoalForked,
GoalUpdate,
GoalUpdated,
type GoalModelState,
Expand All @@ -79,7 +79,7 @@ import type {
registerEvent2Class(GoalCreate);
registerEvent2Class(GoalUpdate);
registerEvent2Class(GoalClear);
registerEvent2Class(GoalForked);
registerEvent2Class(Forked);

const MAX_GOAL_OBJECTIVE_LENGTH = 4000;

Expand Down Expand Up @@ -1322,7 +1322,7 @@ export class AgentGoalService extends AgentActorService<GoalRuntimeState> implem
this.actor = this.attachActor(goalActorLogic, {
id: 'goal',
durable: {
events: [GoalCreate, GoalUpdate, GoalClear, GoalForked, ContextAppendMessage],
events: [GoalCreate, GoalUpdate, GoalClear, Forked, ContextAppendMessage],
undoable: false,
transition: (state, event) => {
if (event instanceof GoalCreate) {
Expand Down Expand Up @@ -1375,7 +1375,7 @@ export class AgentGoalService extends AgentActorService<GoalRuntimeState> implem
state.forkNotice.goalPresent = false;
return;
}
if (event instanceof GoalForked) {
if (event instanceof Forked) {
state.goal = null;
state.forkNotice.reminderPending =
state.forkNotice.goalPresent || state.forkNotice.reminderPending;
Expand Down
1 change: 1 addition & 0 deletions packages/agent-core-v2/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -463,6 +463,7 @@ export * from '#/features/cron/tools/cron-delete/cron-delete';
import '#/session/agentLifecycle/profile/profiles';
export * from '#/session/agentLifecycle/agentLifecycle';
export * from '#/session/agentLifecycle/agentLifecycleService';
export * from '#/session/agentLifecycle/forked';
export * from '#/session/agentLifecycle/mainAgent';
export * from '#/session/mcp/sessionMcpHandle';
import '#/app/mcpConfig/configSection';
Expand Down
15 changes: 15 additions & 0 deletions packages/agent-core-v2/src/session/agentLifecycle/forked.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
/* oxlint-disable typescript-eslint/no-unsafe-declaration-merging, eslint-plugin-import/namespace -- Event2 class+payload-interface declaration merging is the sanctioned event-declaration idiom. */
import { z } from 'zod';

import { AgentEvent2 } from '#/app/event/event2';

const forkedSchema = z.object({ agentId: z.string() });

export class Forked extends AgentEvent2<z.infer<typeof forkedSchema>> {
static override readonly type = 'forked';
static override readonly durable = true;
static override readonly schema = forkedSchema;
}
export interface Forked {
readonly agentId: string;
}
Loading
Loading