From fdb2135d0417d34f1c505e368ae95872b24ced2f Mon Sep 17 00:00:00 2001 From: elkaix Date: Wed, 30 Sep 2026 20:02:14 -0400 Subject: [PATCH 1/3] fix: restore filesystem watch defaults and align agent status and completion budget behavior Watchers default back on, NotifyUser nudges respect the host update panel, permission mode changes publish agent.status.updated, undo drops its turn's interruption reminder, forks clear inherited cron tasks with a notice, and the completion token cap is only sent when configured. --- apps/vis/server/src/lib/agent-record-types.ts | 88 ++++++++---------- docs/configuration/config-files.md | 4 +- docs/configuration/env-vars.md | 2 +- .../agent-core-v2/docs/wire-manifest.d.ts | 4 +- .../interruptionReminderService.ts | 5 +- .../agent/llmRequester/llmRequesterService.ts | 25 ++++-- .../permissionMode/permissionModeService.ts | 4 + .../src/agent/usage/usageEvents.ts | 1 + .../src/app/config/configService.ts | 2 +- .../src/features/cron/cronAgentRuntime.ts | 79 +++++++++++----- .../src/features/cron/cronOps.ts | 5 +- .../src/features/cron/cronService.ts | 78 +++++++++++----- .../src/features/goal/goalOps.ts | 11 --- .../src/features/goal/goalService.ts | 8 +- .../features/notify/notifyUserNudgeService.ts | 9 +- .../src/human/llm-pythinker/provider.ts | 3 +- .../src/human/llm-pythinker/trait.ts | 17 ++++ .../src/human/llm/protocol/format.ts | 2 +- .../llm/requester/bases/anthropic/profile.ts | 2 +- .../bases/openai-responses/requester.ts | 6 +- .../requester/bases/openai-responses/trait.ts | 6 ++ .../llm/requester/bases/openai/requester.ts | 6 +- .../human/llm/requester/bases/openai/trait.ts | 6 ++ .../src/human/test/llm/trait.test.ts | 18 +++- .../src/human/test/utils/watch.test.ts | 20 ++--- .../agent-core-v2/src/human/utils/watch.ts | 2 +- packages/agent-core-v2/src/index.ts | 1 + .../llm-adapter/model/completion-budget.ts | 33 ++----- .../src/llm-adapter/model/model.types.ts | 5 -- .../provider/provider-definition.ts | 2 + .../src/session/agentLifecycle/forked.ts | 15 ++++ .../fullCompaction/fullCompaction.test.ts | 4 +- .../test/agent/loop/loop.test.ts | 64 +++++++++---- .../permissionMode/permissionMode.test.ts | 23 ++++- .../workspaceAliasesService.test.ts | 3 - .../test/features/cron/sessionCron.test.ts | 90 ++++++++++++++----- .../notify/notifyUserNudgeService.test.ts | 25 +++++- .../test/features/plan/plan.test.ts | 10 ++- .../model/completionBudget.test.ts | 38 ++++---- .../protocol/protocolAdapterRegistry.test.ts | 4 +- .../agentProfileLoader.test.ts | 5 -- packages/agent-gateway/test/sessions.test.ts | 5 +- .../agent-gateway/test/workspaces.test.ts | 3 - 43 files changed, 477 insertions(+), 266 deletions(-) create mode 100644 packages/agent-core-v2/src/session/agentLifecycle/forked.ts diff --git a/apps/vis/server/src/lib/agent-record-types.ts b/apps/vis/server/src/lib/agent-record-types.ts index c2364c0e8..d01de9b87 100644 --- a/apps/vis/server/src/lib/agent-record-types.ts +++ b/apps/vis/server/src/lib/agent-record-types.ts @@ -1,7 +1,7 @@ // @ts-nocheck // apps/vis/server/src/lib/agent-record-types.ts // Single source of truth: engine shapes come from agent-core-v2 directly. -// Do NOT add local interfaces that duplicate engine shapes — the only +// Do NOT add local interfaces that duplicate upstream shapes — the only // exceptions are the legacy records below, which v2 never writes but old // (v1-written / pre-migration) wires still contain on disk. @@ -20,25 +20,26 @@ export { WIRE_PROTOCOL_VERSION } from '@pymodel/agent-core-v2/wire/migration/mig export type { AgentTaskInfo as BackgroundTaskInfo, AgentTaskStatus as BackgroundTaskStatus, -} from '@pymodel/agent-core-v2/agent/task/types'; +} from '@pymodel/agent-core-v2'; export type { SubagentTaskInfo as AgentBackgroundTaskInfo } from '@pymodel/agent-core-v2'; export type { ProcessTaskInfo as ProcessBackgroundTaskInfo } from '@pymodel/agent-core-v2/agent/tools/os/bash/process-task'; export type { QuestionTaskInfo as QuestionBackgroundTaskInfo } from '@pymodel/agent-core-v2/agent/tools/ask-user-question/question-background-task'; -import type { AgentTaskInfo as BackgroundTaskInfo } from '@pymodel/agent-core-v2/agent/task/types'; import type { + AgentTaskInfo as BackgroundTaskInfo, CronAddPayload, - CronTask, CronCursorPayload, CronDeletePayload, + CronTask, + ExportSessionManifest, + FileHistoryCheckpointed, + FileHistoryTracked, + Forked, FullCompactionBegin, FullCompactionCancel, FullCompactionComplete, - FileHistoryCheckpointed, - FileHistoryTracked, GoalClear, GoalCreate, - GoalForked, GoalUpdate, InteractionRequestEvent, InteractionResolvedEvent, @@ -78,7 +79,7 @@ import type { } from '@pymodel/agent-core-v2/agent/contextMemory/contextEvents'; import type { TurnCancel, TurnEnded, TurnPrompt, TurnSteer } from '@pymodel/agent-core-v2/agent/loop/turnOps'; import type { TurnStepInterrupted } from '@pymodel/agent-core-v2/agent/loop/turnEvents'; -import type { TurnStepRetrying } from '@pymodel/agent-core-v2/agent/stepRetry/stepRetryService'; +import type { TurnStepRetrying } from '@pymodel/agent-core-v2/agent/loop/turnEvents'; import type { UsageRecord } from '@pymodel/agent-core-v2/agent/usage/usageOps'; import type { ConfigUpdate, @@ -89,10 +90,7 @@ import type { import type { PermissionSetMode } from '@pymodel/agent-core-v2/agent/permissionMode/permissionModeOps'; import type { PermissionRecordApprovalResult } from '@pymodel/agent-core-v2/agent/permissionRules/permissionRulesOps'; import type { RuntimeSetBinding } from '@pymodel/agent-core-v2/agent/runtimeBinding/runtimeBindingOps'; -import type { - DynamicWorkflowModeEnter, - DynamicWorkflowModeExit, -} from '@pymodel/agent-core-v2/features/dynamic_workflow/dynamicWorkflowOps'; +import type { SwarmModeEnter, SwarmModeExit } from '@pymodel/agent-core-v2/features/swarm/swarmOps'; import type { TowerModeEnter, TowerModeExit } from '@pymodel/agent-core-v2/features/tower/towerOps'; import type { ToolsUpdateStore } from '@pymodel/agent-core-v2/features/todo/todoOps'; @@ -136,9 +134,19 @@ export interface StaleGuardClearedRecord { readonly time?: number; } +/** v2-dropped durable record: removed with the loop-side prompt admission + * facility, but old wires still contain it. */ +export interface PromptAcceptedRecord { + readonly type: 'prompt.accepted'; + readonly agentId: string; + readonly promptId: string; + readonly content?: unknown; + readonly time?: number; +} + /** The wire file header record. Declared locally (rather than via v2's * `WireMetadataRecord`) so the union member keeps concrete field types — - * the engine interface carries an index signature that would widen + * the upstream interface carries an index signature that would widen * `protocol_version` / `created_at` to `unknown`. */ export interface WireMetadataHeader { readonly type: 'metadata'; @@ -165,11 +173,9 @@ export type AgentRecord = | WireRecordOf<'cron.add', CronAddPayload> | WireRecordOf<'cron.cursor', CronCursorPayload> | WireRecordOf<'cron.delete', CronDeletePayload> - | WireRecordOf<'dynamic_workflow_mode.enter', DynamicWorkflowModeEnter> - | WireRecordOf<'dynamic_workflow_mode.exit', DynamicWorkflowModeExit> | 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> @@ -191,7 +197,7 @@ export type AgentRecord = | WireRecordOf<'plugin.session_start', PluginSessionStartEvent> | WireRecordOf<'profile.bind', ProfileBind> | WireRecordOf<'prompt.aborted', PromptAborted> - | WireRecordOf<'prompt.accepted', PromptAcceptedRecord> + | PromptAcceptedRecord | WireRecordOf<'prompt.completed', PromptCompleted> | WireRecordOf<'prompt.steered', PromptSteered> | WireRecordOf<'runtime.set_binding', RuntimeSetBinding> @@ -200,6 +206,8 @@ export type AgentRecord = | WireRecordOf<'subagent.failed', SubagentFailed> | WireRecordOf<'subagent.spawned', SubagentSpawned> | WireRecordOf<'subagent.started', SubagentStarted> + | WireRecordOf<'swarm_mode.enter', SwarmModeEnter> + | WireRecordOf<'swarm_mode.exit', SwarmModeExit> | WireRecordOf<'task.started', TaskStarted> | WireRecordOf<'task.terminated', TaskTerminated> | WireRecordOf<'task.waitDelivered', TaskWaitDelivered> @@ -234,35 +242,10 @@ export type AgentRecordOf = Extract< /** * `manifest.json` shape inside a `/export-debug-zip` bundle. Structural - * mirror of the engine's `ExportSessionManifest`, which is not re-exported - * from the package entry. All fields optional-tolerant because the manifest - * comes from another machine / pythinker-code version. + * current engine manifest with every field optional because the bundle may + * come from another machine or an older pythinker-code version. */ -export interface ImportManifest { - sessionId?: string; - exportedAt?: string; - pythinkerCodeVersion?: string; - wireProtocolVersion?: string; - os?: string; - nodejsVersion?: string; - sessionFirstActivity?: string; - sessionLastActivity?: string; - title?: string; - workspaceDir?: string; - sessionLogPath?: string; - globalLogPath?: string; - desktopLogPath?: string; - webLogPath?: string; - desktopVersion?: string; - installSource?: string; - shellEnv?: { - term?: string; - termProgram?: string; - termProgramVersion?: string; - multiplexer?: string; - shell?: string; - }; -} +export type ImportManifest = Partial; /** vis-side bookkeeping for one imported bundle, written to * `imported//import-meta.json`. */ @@ -325,13 +308,12 @@ export interface AgentInfo { wireExists: boolean; wireRecordCount: number; wireProtocolVersion: string | null; - /** Per-item dynamic_workflow work label persisted by the engine for - * dynamic-workflow-spawned sub-agents (`AgentMeta.dynamicWorkflowItem`, or - * `AgentMeta.labels.dynamicWorkflowItem` on v2-written sessions). `null` - * when the agent is not a dynamic_workflow item or when the value cannot - * be recovered (e.g. disk-only inventory of a session with a corrupt - * `state.json`). */ - dynamicWorkflowItem: string | null; + /** Per-item swarm work label persisted by the engine for swarm-spawned + * sub-agents (`AgentMeta.swarmItem`, or `AgentMeta.labels.swarmItem` on + * v2-written sessions). `null` when the agent is not a swarm item or when + * the value cannot be recovered (e.g. disk-only inventory of a session + * with a corrupt `state.json`). */ + swarmItem: string | null; } export interface SessionDetail { @@ -341,7 +323,7 @@ export interface SessionDetail { * which can drift after fork/rename. */ sessionDir: string; workDir: string; - state: unknown; // Preserve the source shape; the UI renders the actual state.json form. + state: unknown; // 原样透传,前端按 state.json 真实形状渲染 agents: AgentInfo[]; /** True for sessions imported from a debug zip. */ imported: boolean; diff --git a/docs/configuration/config-files.md b/docs/configuration/config-files.md index 91ebd573b..98163ae81 100644 --- a/docs/configuration/config-files.md +++ b/docs/configuration/config-files.md @@ -474,11 +474,11 @@ Both values must be positive integers. A call's `max_chars` overrides the defaul ## `watch` -`watch` controls filesystem watchers that reload local.toml, AGENTS.md, skills, MCP config, and `config.toml` itself. It defaults to off. Set `enabled` to `true` to attach watchers; with watchers off, changing the file later will not be picked up until restart. +`watch` controls filesystem watchers that reload local.toml, AGENTS.md, skills, MCP config, and `config.toml` itself. It defaults to on. Set `enabled` to `false` to start with no watchers; changing the file later will not be picked up until restart. | Field | Type | Default | Description | | --- | --- | --- | --- | -| `enabled` | `boolean` | `false` | Attach filesystem watchers; `false` disables every `watch()` for the process | +| `enabled` | `boolean` | `true` | Attach filesystem watchers; `false` disables every `watch()` for the process | `enabled` can be overridden by the `PYTHINKER_CODE_WATCH` environment variable, which takes higher priority than `config.toml`. diff --git a/docs/configuration/env-vars.md b/docs/configuration/env-vars.md index bebad2470..80fb8e139 100644 --- a/docs/configuration/env-vars.md +++ b/docs/configuration/env-vars.md @@ -146,7 +146,7 @@ Switches that control the behavior of subsystems such as telemetry, background t | `PYTHINKER_CODE_EXPERIMENTAL_SUBAGENT_FORK` | Enable the experimental `fork` parameter on the `Agent` and `AgentDynamicWorkflow` tools, letting the model start a subagent with a snapshot of the calling agent's conversation history instead of an empty context; the master `PYTHINKER_CODE_EXPERIMENTAL_FLAG=1` also enables it | Truthy: `1`/`true`/`yes`/`on`; falsy: `0`/`false`/`no`/`off` | | `PYTHINKER_CODE_EXPERIMENTAL_TOOL_SELECT` | Experimental on-demand tool loading: tools of MCP servers marked `deferred: true` stay out of the top-level tool list and are loaded via `select_tools`; also requires the model to declare the `dynamically_loaded_tools` capability — see [MCP](../customization/mcp.md#loading-tools-on-demand) | Truthy: `1`/`true`/`yes`/`on`; falsy: `0`/`false`/`no`/`off` | | `PYTHINKER_CODE_EXPERIMENTAL_TOWER` | Enable the experimental [`/tower`](../reference/slash-commands.md#modes--run-control) command for workspace-wide subagent coordination; the master `PYTHINKER_CODE_EXPERIMENTAL_FLAG=1` also enables it | Truthy: `1`/`true`/`yes`/`on`; falsy: `0`/`false`/`no`/`off` | -| `PYTHINKER_CODE_WATCH` | Attach filesystem watchers that reload config and workspace files; higher priority than `[watch] enabled` (default `false`) | Truthy: `1`/`true`/`yes`/`on`; falsy: `0`/`false`/`no`/`off` | +| `PYTHINKER_CODE_WATCH` | Attach filesystem watchers that reload config and workspace files; higher priority than `[watch] enabled` (default `true`) | Truthy: `1`/`true`/`yes`/`on`; falsy: `0`/`false`/`no`/`off` | | `PYTHINKER_CODE_SEARCH_WORKER` | Run the global search index in a dedicated worker thread; takes higher priority than `[database] search` in `config.toml` (default `true`) | Truthy: `1`/`true`/`yes`/`on`; falsy: `0`/`false`/`no`/`off` | | `PYTHINKER_CODE_PERSISTENCE_MINIDB_READMODEL` | Use the minidb-backed read model for session indexing; takes higher priority than `[database] base` in `config.toml` (default `true`) | Truthy: `1`/`true`/`yes`/`on`; falsy: `0`/`false`/`no`/`off` | | `PYTHINKER_MCP_STARTUP_TIMEOUT_MS` | Global default connection timeout (ms) for all MCP servers; takes higher priority than `[mcp] startup_timeout_ms` in `config.toml`, but a per-server `startupTimeoutMs` in `mcp.json` still wins (default `30000`) | Integer from `1` to `2147483647`; invalid values are ignored | diff --git a/packages/agent-core-v2/docs/wire-manifest.d.ts b/packages/agent-core-v2/docs/wire-manifest.d.ts index d690ec53d..0ff629a4b 100644 --- a/packages/agent-core-v2/docs/wire-manifest.d.ts +++ b/packages/agent-core-v2/docs/wire-manifest.d.ts @@ -38,7 +38,7 @@ // dynamic_workflow_mode.exit contextMemory, dynamic_workflow src/features/dynamic_workflow/dynamicWorkflowOps.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 @@ -268,7 +268,7 @@ interface FileHistoryTrackedPayload { /** * states: (none) - * owner: src/features/goal/goalOps.ts + * owner: src/session/agentLifecycle/forked.ts */ interface ForkedPayload { _name: 'forked'; diff --git a/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts index acea65627..9083d4b62 100644 --- a/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts +++ b/packages/agent-core-v2/src/agent/interruptionReminder/interruptionReminderService.ts @@ -2,6 +2,7 @@ import { Disposable } from '#/_base/di/lifecycle'; import { LifecycleScope } from '#/app/scopes'; import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory'; +import { isUndoAnchor } from '#/agent/contextMemory/conversationTime'; import type { ContextMessage } from '#/agent/contextMemory/types'; import { isVacuousContentPart } from '#/agent/contextMemory/vacuousContent'; import { TurnEnded } from '#/agent/loop/turnOps'; @@ -35,10 +36,12 @@ export class AgentInterruptionReminderService this._register( eventBus.subscribe(TurnEnded, (event) => { if (event.reason !== 'cancelled' || event.interruptReason !== 'user_cancelled') return; - const origin = lastComparableMessage(this.context.get())?.origin; + const history = this.context.get(); + const origin = lastComparableMessage(history)?.origin; if (origin?.kind === 'injection' && origin.variant === INTERRUPTION_REMINDER_VARIANT) return; this.reminder.notify(INTERRUPTION_REMINDER, { variant: INTERRUPTION_REMINDER_VARIANT, + ownerPromptId: history.findLast(isUndoAnchor)?.id, }); }), ); diff --git a/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts b/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts index fcef80b0c..178248543 100644 --- a/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts +++ b/packages/agent-core-v2/src/agent/llmRequester/llmRequesterService.ts @@ -660,26 +660,33 @@ export class AgentLLMRequesterService implements IAgentLLMRequesterService { const turnConfig = this.resolveTurnConfig(overrides.source); const resolved = turnConfig?.resolved ?? this.profile.resolveModelContext(); const baseParams = turnConfig?.params ?? this.profile.resolveRequestParams(); + const requester = this.modelCatalog.getRequester(resolved.modelAlias); + const maxCompletionTokensCap = + this.config.get('modelOverrides')?.maxCompletionTokens; + const usedContextTokens = + overrides.messages === undefined + ? this.tokenCounting.get(this.scopeContext.agentContext).size + : undefined; const budgetParams = completionBudgetParams({ budget: resolveCompletionBudget({ maxOutputSize: overrides.maxOutputSize ?? resolved.maxOutputSize, - reservedContextSize: resolved.reservedContextSize, - maxCompletionTokensCap: - this.config.get('modelOverrides')?.maxCompletionTokens, + maxCompletionTokensCap, }), capability: resolved.modelCapabilities, - usedContextTokens: - overrides.messages === undefined - ? this.tokenCounting.get(this.scopeContext.agentContext).measured - : undefined, + usedContextTokens, }); - const requester = this.modelCatalog.getRequester(resolved.modelAlias); + const optedOut = maxCompletionTokensCap !== undefined && maxCompletionTokensCap <= 0; const messages = overrides.messages ?? this.context.get(); return { requester, model: requester.model, - params: { ...baseParams, ...budgetParams }, + params: { + ...baseParams, + ...budgetParams, + ...(usedContextTokens === undefined ? {} : { usedContextTokens }), + ...(optedOut ? { maxCompletionTokens: maxCompletionTokensCap } : {}), + }, modelAlias: resolved.modelAlias, thinkingEffort: resolved.thinkingLevel, systemPrompt: overrides.systemPrompt ?? turnConfig?.systemPrompt ?? this.profile.getSystemPrompt(), diff --git a/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts b/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts index 5aaca3f17..a5f40575a 100644 --- a/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts +++ b/packages/agent-core-v2/src/agent/permissionMode/permissionModeService.ts @@ -14,6 +14,7 @@ import { MAIN_AGENT_ID, } from '#/session/agentLifecycle/agentLifecycle'; import { IAgentStateService } from '#/agent/state/agentState'; +import { AgentStatusUpdated } from '#/agent/usage/usageEvents'; import { IEventDispatcher } from '#/state/eventDispatcher'; import { IAgentPermissionModeService, type PermissionModeChangedContext } from './permissionMode'; import { @@ -58,6 +59,9 @@ export class AgentPermissionModeService extends Service implements IAgentPermiss void this.dispatcher.dispatch( new PermissionSetMode({ agentId: this.scopeContext.agentId, mode }), ); + void this.dispatcher.dispatch( + new AgentStatusUpdated({ agentId: this.scopeContext.agentId, permission: mode }), + ); if (changed) this._onDidChangeMode.fire({ mode, previousMode }); } diff --git a/packages/agent-core-v2/src/agent/usage/usageEvents.ts b/packages/agent-core-v2/src/agent/usage/usageEvents.ts index be7a22809..d1e1f2466 100644 --- a/packages/agent-core-v2/src/agent/usage/usageEvents.ts +++ b/packages/agent-core-v2/src/agent/usage/usageEvents.ts @@ -15,6 +15,7 @@ export interface AgentStatusUpdatedPayload { thinkingEffort?: string; maxContextTokens?: number; contextTokens?: number; + permission?: PermissionMode; } export class AgentStatusUpdated extends AgentEvent2 { diff --git a/packages/agent-core-v2/src/app/config/configService.ts b/packages/agent-core-v2/src/app/config/configService.ts index 4c3905a44..76400b0bf 100644 --- a/packages/agent-core-v2/src/app/config/configService.ts +++ b/packages/agent-core-v2/src/app/config/configService.ts @@ -669,7 +669,7 @@ export class ConfigService extends Disposable implements IConfigService { } private applyWatchEnabled(): void { - setWatchEnabled(this.get(WATCH_SECTION)?.enabled ?? false); + setWatchEnabled(this.get(WATCH_SECTION)?.enabled ?? true); } private deliveredValue(domain: string): unknown { diff --git a/packages/agent-core-v2/src/features/cron/cronAgentRuntime.ts b/packages/agent-core-v2/src/features/cron/cronAgentRuntime.ts index d527347b8..47387a369 100644 --- a/packages/agent-core-v2/src/features/cron/cronAgentRuntime.ts +++ b/packages/agent-core-v2/src/features/cron/cronAgentRuntime.ts @@ -3,6 +3,8 @@ import { assign, fromCallback, sendTo, setup, type Snapshot } from 'xstate'; import { IntervalTimer } from '#/_base/utils/timer'; import type { CronJobOrigin, CronMissedOrigin } from '#/agent/contextMemory/types'; +import { ContextAppendMessage } from '#/agent/contextMemory/contextEvents'; +import type { ContextMessage } from '#/agent/contextMemory/types'; import { IAgentLoopService, type Turn } from '#/agent/loop/loop'; import { defineAgentRuntimeContract, @@ -21,7 +23,9 @@ import type { CronDeletedEvent, CronScheduledEvent } from '#/app/telemetry/event import { ITelemetryService } from '#/app/telemetry/telemetry'; import { BugIndicatingError } from '#/errors'; import type { ContentPart } from '#/kosong/contract/message'; +import { IAgentReminderService } from '#/features/reminder/reminderService'; import { MAIN_AGENT_ID } from '#/session/agentLifecycle/agentLifecycle'; +import { Forked } from '#/session/agentLifecycle/forked'; import { CronAdd, CronCursor, CronDelete, CronFired, type CronModelState } from './cronOps'; @@ -36,14 +40,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: AgentRuntimeContext; } interface CronCommitEvent { readonly type: 'cron.commit'; - readonly tasks: CronModelState; + readonly model: CronModelState; } interface CronTickEvent { @@ -151,7 +168,7 @@ function removeTasks( runtime: AgentRuntimeContext, 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; } @@ -263,7 +280,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 })); } } @@ -276,11 +293,11 @@ async function tickCron( ): Promise { await config.ready; if (isDisposed()) return; - if (readCronConfig(config).disabled || runtime.getState().size === 0) return; + if (readCronConfig(config).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, isDisposed)), + [...runtime.getState().tasks.values()].map((task) => processDue(runtime, state, task, now, isDisposed)), ); } @@ -297,6 +314,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 config = configOf(input.runtime); const timer = new IntervalTimer({ unref: true }); const state: CronEffectState = { @@ -378,7 +400,7 @@ export class CronRuntime { } addTask(init: CronTaskInit): CronTask { - const tasks = this.runtime.getState(); + const tasks = this.runtime.getState().tasks; let id: string | undefined; for (let attempt = 0; attempt < MAX_ID_ATTEMPTS; attempt += 1) { const candidate = ulid(); @@ -400,11 +422,11 @@ export class CronRuntime { } getTask(id: string): CronTask | undefined { - return this.runtime.getState().get(id); + return this.runtime.getState().tasks.get(id); } list(): readonly CronTask[] { - return [...this.runtime.getState().values()]; + return [...this.runtime.getState().tasks.values()]; } isStale(task: CronTask): boolean { @@ -413,7 +435,7 @@ export class CronRuntime { getNextFireTime(): number | null { let min: number | null = null; - for (const task of this.runtime.getState().values()) { + for (const task of this.runtime.getState().tasks.values()) { const next = nextFireFor(this.runtime, task); if (next !== null && (min === null || next < min)) min = next; } @@ -421,7 +443,7 @@ export class CronRuntime { } getNextFireForTask(taskId: string): number | null { - const task = this.runtime.getState().get(taskId); + const task = this.runtime.getState().tasks.get(taskId); return task === undefined ? null : nextFireFor(this.runtime, task); } @@ -482,7 +504,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: { @@ -509,7 +534,7 @@ const cronActorLogic = setup({ }, on: { 'cron.commit': { - actions: assign({ tasks: ({ event }) => event.tasks }), + actions: assign({ model: ({ event }) => event.model }), }, }, }); @@ -521,28 +546,40 @@ export const cronAgentRuntimeProvider = defineAgentRuntimeProvider { 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 }); }, }, createApi: (context) => new CronRuntime(context), inspect: (snapshot) => - [...(snapshot as CronActorSnapshot).context.tasks.values()].map((task) => ({ + [...(snapshot as CronActorSnapshot).context.model.tasks.values()].map((task) => ({ id: task.id, cron: task.cron, recurring: task.recurring !== false, diff --git a/packages/agent-core-v2/src/features/cron/cronOps.ts b/packages/agent-core-v2/src/features/cron/cronOps.ts index 544aa817c..6c282e523 100644 --- a/packages/agent-core-v2/src/features/cron/cronOps.ts +++ b/packages/agent-core-v2/src/features/cron/cronOps.ts @@ -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; +export interface CronModelState { + readonly tasks: Map; + readonly forkNotice: { reminderPending: boolean }; +} const cronTaskSchema = z.object({ id: z.string(), diff --git a/packages/agent-core-v2/src/features/cron/cronService.ts b/packages/agent-core-v2/src/features/cron/cronService.ts index 1d8d0d7e8..30a305a21 100644 --- a/packages/agent-core-v2/src/features/cron/cronService.ts +++ b/packages/agent-core-v2/src/features/cron/cronService.ts @@ -23,7 +23,11 @@ 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'; @@ -31,6 +35,7 @@ import { CronAdd, CronCursor, CronDelete, CronFired, type CronModelState } from registerEvent2Class(CronAdd); registerEvent2Class(CronDelete); registerEvent2Class(CronCursor); +registerEvent2Class(Forked); const STALE_THRESHOLD_MS = 7 * 24 * 60 * 60 * 1000; const DEFAULT_POLL_INTERVAL_MS = 1_000; @@ -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; } interface CronCommitEvent { readonly type: 'cron.commit'; - readonly tasks: CronModelState; + readonly model: CronModelState; } interface CronTickEvent { @@ -154,7 +172,7 @@ function removeTasks( runtime: AgentActorContext, 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; } @@ -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 })); } } @@ -264,10 +282,10 @@ async function tickCron( state: CronEffectState, ): Promise { 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(({ @@ -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, @@ -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: { @@ -383,7 +409,7 @@ const cronActorLogic = setup({ }, on: { 'cron.commit': { - actions: assign({ tasks: ({ event }) => event.tasks }), + actions: assign({ model: ({ event }) => event.model }), }, }, }); @@ -429,24 +455,36 @@ export class AgentCronService extends AgentActorService 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 }); }, }, }); } @@ -460,7 +498,7 @@ export class AgentCronService extends AgentActorService 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(); @@ -482,11 +520,11 @@ export class AgentCronService extends AgentActorService 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 { @@ -495,7 +533,7 @@ export class AgentCronService extends AgentActorService 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; } @@ -503,7 +541,7 @@ export class AgentCronService extends AgentActorService implemen } 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); } diff --git a/packages/agent-core-v2/src/features/goal/goalOps.ts b/packages/agent-core-v2/src/features/goal/goalOps.ts index 6e2169b06..ffe5f6823 100644 --- a/packages/agent-core-v2/src/features/goal/goalOps.ts +++ b/packages/agent-core-v2/src/features/goal/goalOps.ts @@ -111,17 +111,6 @@ export interface GoalClear { readonly agentId: string; } -const goalForkedSchema = z.object({ agentId: z.string() }); - -export class GoalForked extends AgentEvent2> { - 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; diff --git a/packages/agent-core-v2/src/features/goal/goalService.ts b/packages/agent-core-v2/src/features/goal/goalService.ts index 713db4506..19ce2c3c5 100644 --- a/packages/agent-core-v2/src/features/goal/goalService.ts +++ b/packages/agent-core-v2/src/features/goal/goalService.ts @@ -49,6 +49,7 @@ import { type PythinkerErrorPayload, } 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'; @@ -58,7 +59,6 @@ import { IGoalDeadlineScheduler } from './goalDeadlineScheduler'; import { GoalClear, GoalCreate, - GoalForked, GoalUpdate, GoalUpdated, type GoalModelState, @@ -79,7 +79,7 @@ import type { registerEvent2Class(GoalCreate); registerEvent2Class(GoalUpdate); registerEvent2Class(GoalClear); -registerEvent2Class(GoalForked); +registerEvent2Class(Forked); const MAX_GOAL_OBJECTIVE_LENGTH = 4000; @@ -1322,7 +1322,7 @@ export class AgentGoalService extends AgentActorService 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) { @@ -1375,7 +1375,7 @@ export class AgentGoalService extends AgentActorService 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; diff --git a/packages/agent-core-v2/src/features/notify/notifyUserNudgeService.ts b/packages/agent-core-v2/src/features/notify/notifyUserNudgeService.ts index 2a0330ad4..1a35ff9c8 100644 --- a/packages/agent-core-v2/src/features/notify/notifyUserNudgeService.ts +++ b/packages/agent-core-v2/src/features/notify/notifyUserNudgeService.ts @@ -9,11 +9,12 @@ import { import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory'; import { IAgentScopeContext } from '#/agent/scopeContext/scopeContext'; import { IAgentToolRegistryService } from '#/agent/toolRegistry/toolRegistry'; +import { IBootstrapService } from '#/app/bootstrap/bootstrap'; import { IFlagService } from '#/app/flag/flag'; import { IAgentReminderService } from '#/features/reminder/reminderService'; import { IEventDispatcher } from '#/state/eventDispatcher'; -import { NOTIFY_USER_FLAG_ID } from './flag'; +import { notifyUserAvailable } from './notifyUserAvailability'; import { NOTIFY_USER_NUDGE_VARIANT, lastMidResponsePosition, @@ -38,11 +39,13 @@ const notifyUserNudgeReminders = fromCallback(({ }; }) => { const runtime = input.runtime; - if (!runtime.get(IFlagService).enabled(NOTIFY_USER_FLAG_ID)) return () => {}; + const available = (): boolean => + notifyUserAvailable(runtime.get(IFlagService), runtime.get(IBootstrapService)); + if (!available()) return () => {}; const registration = runtime.get(IAgentReminderService).register( NOTIFY_USER_NUDGE_VARIANT, ({ lastInjectedAt }): string | undefined => { - if (!runtime.get(IFlagService).enabled(NOTIFY_USER_FLAG_ID)) return undefined; + if (!available()) return undefined; if (runtime.get(IAgentToolRegistryService).resolve(NOTIFY_USER_TOOL_NAME) === undefined) { return undefined; } diff --git a/packages/agent-core-v2/src/human/llm-pythinker/provider.ts b/packages/agent-core-v2/src/human/llm-pythinker/provider.ts index 90bf8cc80..590a83194 100644 --- a/packages/agent-core-v2/src/human/llm-pythinker/provider.ts +++ b/packages/agent-core-v2/src/human/llm-pythinker/provider.ts @@ -3,7 +3,7 @@ import { anthropicBetaBase } from '#/llm/requester/bases/anthropic/requester'; import { openAIBase } from '#/llm/requester/bases/openai/requester'; import { openAIResponsesBase } from '#/llm/requester/bases/openai-responses/requester'; -import { pythinkerAnthropicTrait, pythinkerConnection, pythinkerOpenAITrait } from './trait'; +import { pythinkerAnthropicTrait, pythinkerConnection, pythinkerOpenAITrait, pythinkerResponsesTrait } from './trait'; import { classifyPythinkerQuotaError } from './errors'; import { pythinkerMediaContribution } from './media'; @@ -24,6 +24,7 @@ export const pythinkerProvider = createProvider({ }, openai_responses: { base: openAIResponsesBase, + trait: pythinkerResponsesTrait, connection: pythinkerConnection, classifyError: classifyPythinkerQuotaError, }, diff --git a/packages/agent-core-v2/src/human/llm-pythinker/trait.ts b/packages/agent-core-v2/src/human/llm-pythinker/trait.ts index d28ac17a0..cc6b1557b 100644 --- a/packages/agent-core-v2/src/human/llm-pythinker/trait.ts +++ b/packages/agent-core-v2/src/human/llm-pythinker/trait.ts @@ -1,3 +1,4 @@ +import type { LlmModel } from '#/llm/model'; import type { ProtocolEndpoint, ProviderConnection } from '#/llm/protocol/connection'; import type { ContentPart, ToolDescription } from '#/llm/message'; import { providerImagePolicy } from '#/llm/media/image-formats'; @@ -8,6 +9,7 @@ import type { OpenAIWireMessage, OpenAIWireToolCall, } from '#/llm/requester/bases/openai/contract'; +import type { OpenAIResponsesTrait } from '#/llm/requester/bases/openai-responses/trait'; import type { OpenAITrait } from '#/llm/requester/bases/openai/trait'; import { normalizePythinkerToolSchema } from './schema'; @@ -64,6 +66,19 @@ function convertPythinkerTool(tool: ToolDescription): Record { const pythinkerAcceptedImageMimes = (): ReadonlySet => providerImagePolicy('pythinker').acceptedMimes; +export function pythinkerUnsetCompletionTokens(input: { + readonly model: LlmModel; + readonly usedContextTokens?: number; +}): number | undefined { + const window = input.model.maxContextSize; + if (window === undefined || window <= 0 || input.usedContextTokens === undefined) return undefined; + return Math.max(1, window - input.usedContextTokens); +} + +export const pythinkerResponsesTrait: OpenAIResponsesTrait = { + completionTokensWhenUnset: pythinkerUnsetCompletionTokens, +}; + export const pythinkerOpenAITrait: OpenAITrait = { strictThinkingValidation: true, @@ -91,6 +106,8 @@ export const pythinkerOpenAITrait: OpenAITrait = { max_completion_tokens: maxCompletionTokens, }), + completionTokensWhenUnset: pythinkerUnsetCompletionTokens, + buildParams: (params) => { const { extra_body: extraBody, ...rest } = params; if (extraBody === undefined || extraBody === null) { diff --git a/packages/agent-core-v2/src/human/llm/protocol/format.ts b/packages/agent-core-v2/src/human/llm/protocol/format.ts index c25ed5e17..37fa76d42 100644 --- a/packages/agent-core-v2/src/human/llm/protocol/format.ts +++ b/packages/agent-core-v2/src/human/llm/protocol/format.ts @@ -12,7 +12,7 @@ export type FormatRequestInput = LlmRequestConfig & { export function resolveMaxCompletionCap(input: FormatRequestInput): number | undefined { const { maxCompletionTokens, usedContextTokens, maxContextTokens } = input; - if (maxCompletionTokens === undefined) { + if (maxCompletionTokens === undefined || maxCompletionTokens <= 0) { return undefined; } let cap = maxCompletionTokens; diff --git a/packages/agent-core-v2/src/human/llm/requester/bases/anthropic/profile.ts b/packages/agent-core-v2/src/human/llm/requester/bases/anthropic/profile.ts index 4959ee973..b351d787f 100644 --- a/packages/agent-core-v2/src/human/llm/requester/bases/anthropic/profile.ts +++ b/packages/agent-core-v2/src/human/llm/requester/bases/anthropic/profile.ts @@ -133,7 +133,7 @@ const CEILING_BY_FAMILY_VERSION: Readonly> = { 'haiku-3': 4096, }; -const FALLBACK_MAX_TOKENS = 128000; +const FALLBACK_MAX_TOKENS = 64000; function lookupClaudeCeiling(version: AnthropicModelVersion): number | undefined { const { family, major, minor } = version; diff --git a/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/requester.ts b/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/requester.ts index 357e54035..5aac0822a 100644 --- a/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/requester.ts +++ b/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/requester.ts @@ -81,7 +81,11 @@ export function prepareOpenAIResponsesRequest( encodeReasoningEffortFallback(t, ctx.model, trait?.strictThinkingValidation === true), ).kwargs; } - const cap = resolveMaxCompletionCap(input); + const requested = input.maxCompletionTokens; + const cap = + requested !== undefined && requested <= 0 + ? undefined + : (resolveMaxCompletionCap(input) ?? trait?.completionTokensWhenUnset?.(input)); if (cap !== undefined) { kwargs = { ...kwargs, diff --git a/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/trait.ts b/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/trait.ts index 304837481..d202bd39c 100644 --- a/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/trait.ts +++ b/packages/agent-core-v2/src/human/llm/requester/bases/openai-responses/trait.ts @@ -1,4 +1,5 @@ import type { ToolDescription } from '#/llm/message'; +import type { LlmModel } from '#/llm/model'; import type { TraitContext } from '#/llm/protocol/base'; import type { ThinkingStrategy } from '#/llm/protocol/thinking'; import type { ToolCallIdPolicy, ToolMessageConversion } from '#/llm/requester/requester'; @@ -19,6 +20,11 @@ export interface OpenAIResponsesTrait { ctx: TraitContext, ): Record | undefined; + completionTokensWhenUnset?(input: { + readonly model: LlmModel; + readonly usedContextTokens?: number; + }): number | undefined; + convertTool?(tool: ToolDescription, ctx: TraitContext): Record | undefined; mergeHistory?( diff --git a/packages/agent-core-v2/src/human/llm/requester/bases/openai/requester.ts b/packages/agent-core-v2/src/human/llm/requester/bases/openai/requester.ts index a19bfc1c6..70f631266 100644 --- a/packages/agent-core-v2/src/human/llm/requester/bases/openai/requester.ts +++ b/packages/agent-core-v2/src/human/llm/requester/bases/openai/requester.ts @@ -96,7 +96,11 @@ export function prepareOpenAIRequest( if (input.responseFormat !== undefined) { kwargs = { ...kwargs, response_format: responseFormatToOpenAI(input.responseFormat) }; } - const cap = resolveMaxCompletionCap(input); + const requested = input.maxCompletionTokens; + const cap = + requested !== undefined && requested <= 0 + ? undefined + : (resolveMaxCompletionCap(input) ?? trait?.completionTokensWhenUnset?.(input)); if (cap !== undefined) { kwargs = { ...kwargs, diff --git a/packages/agent-core-v2/src/human/llm/requester/bases/openai/trait.ts b/packages/agent-core-v2/src/human/llm/requester/bases/openai/trait.ts index d9377faa8..1a4784297 100644 --- a/packages/agent-core-v2/src/human/llm/requester/bases/openai/trait.ts +++ b/packages/agent-core-v2/src/human/llm/requester/bases/openai/trait.ts @@ -1,4 +1,5 @@ import type { Message, ToolDescription } from '#/llm/message'; +import type { LlmModel } from '#/llm/model'; import type { TraitContext } from '#/llm/protocol/base'; import type { ThinkingStrategy } from '#/llm/protocol/thinking'; import type { ToolCallIdPolicy, ToolMessageConversion } from '#/llm/requester/requester'; @@ -20,6 +21,11 @@ export interface OpenAITrait { ctx: TraitContext, ): Record | undefined; + completionTokensWhenUnset?(input: { + readonly model: LlmModel; + readonly usedContextTokens?: number; + }): number | undefined; + convertTool?(tool: ToolDescription, ctx: TraitContext): Record | undefined; convertMessage?( diff --git a/packages/agent-core-v2/src/human/test/llm/trait.test.ts b/packages/agent-core-v2/src/human/test/llm/trait.test.ts index 1f3964c2c..97f386493 100644 --- a/packages/agent-core-v2/src/human/test/llm/trait.test.ts +++ b/packages/agent-core-v2/src/human/test/llm/trait.test.ts @@ -739,6 +739,20 @@ describe('withMaxCompletionTokens', () => { ); expect(client.body()['max_completion_tokens']).toBe(1000); expect(client.body()['max_tokens']).toBeUndefined(); + + await requester.generate( + { model: { ...model, maxContextSize: 1000 } }, + { messages, usedContextTokens: 40 }, + { signal: new AbortController().signal }, + ); + expect(client.body()['max_completion_tokens']).toBe(960); + + await requester.generate( + { model: { ...model, maxContextSize: 1000 }, maxCompletionTokens: 0 }, + { messages, usedContextTokens: 40 }, + { signal: new AbortController().signal }, + ); + expect(client.body()['max_completion_tokens']).toBeUndefined(); }); it('uses max_completion_tokens for reasoning models without a trait', async () => { @@ -790,7 +804,7 @@ describe('withMaxCompletionTokens', () => { { messages }, { signal: new AbortController().signal }, ); - expect(client.body()['max_tokens']).toBe(128000); + expect(client.body()['max_tokens']).toBe(64000); const sonnet35 = { ...model, model: 'claude-3-5-sonnet-20241022' }; await requester.generate( @@ -1461,7 +1475,7 @@ describe('anthropic thinking kwargs', () => { expect(body['output_config']).toEqual({ effort: 'high' }); expect(body['betaFeatures']).toBeUndefined(); expect(body['betas']).toEqual(['context-management-2025-06-27']); - expect(body['max_tokens']).toBe(128000); + expect(body['max_tokens']).toBe(64000); expect(client.betaCalled()).toBe(true); let bodyMessages = body['messages'] as Record[]; expect(bodyMessages[0]?.['content']).toEqual([ diff --git a/packages/agent-core-v2/src/human/test/utils/watch.test.ts b/packages/agent-core-v2/src/human/test/utils/watch.test.ts index 67c518784..392cc5f67 100644 --- a/packages/agent-core-v2/src/human/test/utils/watch.test.ts +++ b/packages/agent-core-v2/src/human/test/utils/watch.test.ts @@ -20,14 +20,6 @@ const wait = (ms: number): Promise => new Promise((r) => setTimeout(r, ms) const longTempDir = (prefix: string): Promise => mkdtemp(join(realpathSync.native(tmpdir()), prefix)); -beforeEach(() => { - setWatchEnabled(true); -}); - -afterEach(() => { - setWatchEnabled(false); -}); - class TestNativeWatcher { private errorListener: ((error: NodeJS.ErrnoException) => void) | undefined; closed = false; @@ -477,10 +469,14 @@ describe('watch chokidar mode', () => { it('does not start a filesystem watch when watch is disabled', async () => { root = await mkdtemp(join(tmpdir(), 'watch-disabled-')); setWatchEnabled(false); - const events = await start(); - await writeFile(join(root, 'a.txt'), 'x'); - await wait(300); - expect(events).toHaveLength(0); + try { + const events = await start(); + await writeFile(join(root, 'a.txt'), 'x'); + await wait(300); + expect(events).toHaveLength(0); + } finally { + setWatchEnabled(true); + } }); it('stops firing after the handle is disposed', async () => { diff --git a/packages/agent-core-v2/src/human/utils/watch.ts b/packages/agent-core-v2/src/human/utils/watch.ts index e8fe23fbe..c1f4ca231 100644 --- a/packages/agent-core-v2/src/human/utils/watch.ts +++ b/packages/agent-core-v2/src/human/utils/watch.ts @@ -513,7 +513,7 @@ export const WATCH_ENV = 'PYTHINKER_CODE_WATCH'; const TRUE_WATCH_ENV = new Set(['1', 'true', 'yes', 'on']); const FALSE_WATCH_ENV = new Set(['0', 'false', 'no', 'off']); -let watchEnabledFromConfig = false; +let watchEnabledFromConfig = true; export function setWatchEnabled(enabled: boolean): void { watchEnabledFromConfig = enabled; diff --git a/packages/agent-core-v2/src/index.ts b/packages/agent-core-v2/src/index.ts index f83e5fbf2..f24ea0eba 100644 --- a/packages/agent-core-v2/src/index.ts +++ b/packages/agent-core-v2/src/index.ts @@ -465,6 +465,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'; diff --git a/packages/agent-core-v2/src/llm-adapter/model/completion-budget.ts b/packages/agent-core-v2/src/llm-adapter/model/completion-budget.ts index b4fe8c01e..0cbc60df4 100644 --- a/packages/agent-core-v2/src/llm-adapter/model/completion-budget.ts +++ b/packages/agent-core-v2/src/llm-adapter/model/completion-budget.ts @@ -1,50 +1,31 @@ import type { ModelCapability } from '../contract/capability'; -import type { CompletionBudgetConfig, CompletionBudgetParams } from './model.types'; +import type { CompletionBudgetParams } from './model.types'; const MIN_FLOOR = 1; -const DEFAULT_UNKNOWN_CONTEXT_FALLBACK = 32000; export function resolveCompletionBudget(args: { readonly maxOutputSize?: number; - readonly reservedContextSize?: number; readonly maxCompletionTokensCap?: number; -}): CompletionBudgetConfig | undefined { +}): number | undefined { if (args.maxCompletionTokensCap !== undefined) { if (args.maxCompletionTokensCap <= 0) return undefined; - return { hardCap: args.maxCompletionTokensCap }; + return args.maxCompletionTokensCap; } if (args.maxOutputSize !== undefined && args.maxOutputSize > 0) { - return { hardCap: args.maxOutputSize }; + return args.maxOutputSize; } - if (args.reservedContextSize !== undefined && args.reservedContextSize > 0) { - return { fallback: args.reservedContextSize }; - } - return { fallback: DEFAULT_UNKNOWN_CONTEXT_FALLBACK }; -} - -export function computeCompletionBudgetCap(args: { - readonly budget: CompletionBudgetConfig; - readonly capability: ModelCapability | undefined; -}): number { - const maxCtx = args.capability?.max_context_tokens ?? 0; - const cap = - args.budget.hardCap ?? - (maxCtx > 0 ? maxCtx : args.budget.fallback ?? DEFAULT_UNKNOWN_CONTEXT_FALLBACK); - return Math.max(MIN_FLOOR, cap); + return undefined; } export function completionBudgetParams(args: { - readonly budget: CompletionBudgetConfig | undefined; + readonly budget: number | undefined; readonly capability: ModelCapability | undefined; readonly usedContextTokens?: number; }): CompletionBudgetParams | undefined { if (args.budget === undefined) return undefined; return { - maxCompletionTokens: computeCompletionBudgetCap({ - budget: args.budget, - capability: args.capability, - }), + maxCompletionTokens: Math.max(MIN_FLOOR, args.budget), usedContextTokens: args.usedContextTokens, maxContextTokens: args.capability?.max_context_tokens, }; diff --git a/packages/agent-core-v2/src/llm-adapter/model/model.types.ts b/packages/agent-core-v2/src/llm-adapter/model/model.types.ts index 067924775..7b500d58b 100644 --- a/packages/agent-core-v2/src/llm-adapter/model/model.types.ts +++ b/packages/agent-core-v2/src/llm-adapter/model/model.types.ts @@ -8,11 +8,6 @@ export interface ModelOverrides { readonly maxCompletionTokens?: number; } -export interface CompletionBudgetConfig { - readonly hardCap?: number; - readonly fallback?: number; -} - export interface CompletionBudgetParams { readonly maxCompletionTokens: number; readonly usedContextTokens?: number; diff --git a/packages/agent-core-v2/src/llm-adapter/provider/provider-definition.ts b/packages/agent-core-v2/src/llm-adapter/provider/provider-definition.ts index f9809a471..cccfe8346 100644 --- a/packages/agent-core-v2/src/llm-adapter/provider/provider-definition.ts +++ b/packages/agent-core-v2/src/llm-adapter/provider/provider-definition.ts @@ -7,6 +7,7 @@ import { pythinkerAnthropicTrait, pythinkerConnection, pythinkerOpenAITrait, + pythinkerResponsesTrait, PYTHINKER_DEFAULT_BASE_URL, } from '#human/llm-pythinker/trait'; import { classifyPythinkerQuotaError } from '#human/llm-pythinker/errors'; @@ -246,6 +247,7 @@ registerProviderDefinition({ registerProviderDefinition({ id: 'pythinker', baseProtocol: 'openai_responses', + trait: pythinkerResponsesTrait, connection: pythinkerConnection, classifyError: classifyPythinkerQuotaError, endpoint: pythinkerEndpoint, diff --git a/packages/agent-core-v2/src/session/agentLifecycle/forked.ts b/packages/agent-core-v2/src/session/agentLifecycle/forked.ts new file mode 100644 index 000000000..30a910023 --- /dev/null +++ b/packages/agent-core-v2/src/session/agentLifecycle/forked.ts @@ -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> { + static override readonly type = 'forked'; + static override readonly durable = true; + static override readonly schema = forkedSchema; +} +export interface Forked { + readonly agentId: string; +} diff --git a/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts b/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts index 87c41d184..1817b1693 100644 --- a/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts +++ b/packages/agent-core-v2/test/agent/fullCompaction/fullCompaction.test.ts @@ -3017,7 +3017,7 @@ describe('FullCompaction', () => { const events = await ctx.untilTurnEnd(); expect(callCount).toBe(3); - expect(compactionMaxCompletionTokens).toEqual([32000]); + expect(compactionMaxCompletionTokens).toEqual([undefined]); expect(events).toContainEqual( expect.objectContaining({ event: 'compaction.started', @@ -3110,7 +3110,7 @@ describe('FullCompaction', () => { await ctx.untilTurnEnd(); expect(callCount).toBe(3); - expect(compactionMaxCompletionTokens).toEqual([undefined]); + expect(compactionMaxCompletionTokens).toEqual([Number(maxCompletionTokens)]); }, ); diff --git a/packages/agent-core-v2/test/agent/loop/loop.test.ts b/packages/agent-core-v2/test/agent/loop/loop.test.ts index 131fa8771..cdb485829 100644 --- a/packages/agent-core-v2/test/agent/loop/loop.test.ts +++ b/packages/agent-core-v2/test/agent/loop/loop.test.ts @@ -105,7 +105,7 @@ describe('Agent loop', () => { [wire] llm.tools_snapshot { "agentId": "main", "hash": "4f53cda18c2baa0c0354bb5f9a3ecbe5ed12ab4d8e11ba873c2f11161202b945", "tools": [], "time": "