v4.3.13: handleTaskControlAction → gateway/tasks/handle-task-control.ts 분리 (5,021→4,668줄)
3자 순환 의존성(executeTool→handleTaskControlAction→handleChat→executeTool) 발견: - inferTaskChannelFromSession/findBlockedTaskForSession/parseTaskIdFromText는 실제로는 handleChat에 의존하지 않는 순수 쿼리 헬퍼임을 확인하고 팩토리 밖으로 빼서 일반 export로 전환 → handle-chat.ts가 deps 주입 없이 직접 import하도록 변경 (HandleChatDeps 2개 필드 제거) - handleTaskControlAction 자체만 진짜 순환: launchBackgroundTaskRunner가 handleChat을 필요로 하는데 handleChat이 executeTool을, executeTool이 handleTaskControlAction을 필요로 함 → lazy binding으로 해결: executeTool 인스턴스화 이전에 hoisted wrapper 함수 (let _handleTaskControlActionImpl + function handleTaskControlAction)를 선언해 dep로 넘기고, handleChat 생성 후 실제 구현을 대입 실제 백그라운드 태스크 생성→조회→삭제 왕복으로 순환 경로 전체 검증 완료.
This commit is contained in:
+23
-376
@@ -70,6 +70,12 @@ import { createBuildTools } from './chat/build-tools';
|
||||
import { createExecuteTool, type ToolResult } from './chat/execute-tool';
|
||||
import { createHandleChat } from './chat/handle-chat';
|
||||
import { createPersonalityContext } from './chat/personality-context';
|
||||
import {
|
||||
createHandleTaskControl,
|
||||
inferTaskChannelFromSession,
|
||||
findBlockedTaskForSession,
|
||||
parseTaskIdFromText,
|
||||
} from './tasks/handle-task-control';
|
||||
import { registerOpenAIAuthRoutes } from './routes/routes-auth-openai';
|
||||
import { registerUserAppSessionRoutes } from './routes/routes-user-app-sessions';
|
||||
import { registerChannelMappingRoutes } from './routes/routes-channel-mappings';
|
||||
@@ -1506,17 +1512,6 @@ async function imageSearch(query: string, count: number = 3): Promise<string> {
|
||||
.join('\n\n');
|
||||
}
|
||||
|
||||
interface TaskControlResponse {
|
||||
success: boolean;
|
||||
action: string;
|
||||
code?: string;
|
||||
message?: string;
|
||||
scope?: string;
|
||||
task?: Record<string, any> | null;
|
||||
tasks?: Array<Record<string, any>>;
|
||||
candidates?: Array<Record<string, any>>;
|
||||
}
|
||||
|
||||
function normalizeToolArgs(rawArgs: any): any {
|
||||
if (rawArgs == null) return {};
|
||||
if (typeof rawArgs === 'string') {
|
||||
@@ -1581,6 +1576,17 @@ function repairJson(input: string): string {
|
||||
return s;
|
||||
}
|
||||
|
||||
// handleTaskControlAction needs handleChat (via launchBackgroundTaskRunner), and
|
||||
// handleChat needs executeTool — but executeTool must exist before handleChat can
|
||||
// be created. Lazy-bind: executeTool gets this hoisted wrapper function as its dep;
|
||||
// the real implementation is plugged into _handleTaskControlActionImpl once handleChat
|
||||
// exists below. Safe because the wrapper is only ever invoked at request time, well
|
||||
// after module init has finished assigning the real implementation.
|
||||
let _handleTaskControlActionImpl: (sessionId: string, args: any) => Promise<any>;
|
||||
function handleTaskControlAction(sessionId: string, args: any): Promise<any> {
|
||||
return _handleTaskControlActionImpl(sessionId, args);
|
||||
}
|
||||
|
||||
const executeTool = createExecuteTool({
|
||||
isPathInsideDir,
|
||||
cronScheduler,
|
||||
@@ -2061,7 +2067,6 @@ const { handleChat, handleCodeChat } = createHandleChat({
|
||||
recordOrchestrationEvent,
|
||||
telegramChannel,
|
||||
makeBroadcastForTask,
|
||||
inferTaskChannelFromSession,
|
||||
resolveToolImageContent,
|
||||
isSkillEnabledForUser,
|
||||
executeTool,
|
||||
@@ -2082,7 +2087,6 @@ const { handleChat, handleCodeChat } = createHandleChat({
|
||||
getPreemptSessionCount,
|
||||
getOrchestrationSessionStats,
|
||||
getModelProfileExtraPrompt,
|
||||
findBlockedTaskForSession,
|
||||
extractLikelyUrl,
|
||||
evaluateBrowserSnapshotQuality,
|
||||
buildPersonalityContext,
|
||||
@@ -2091,373 +2095,16 @@ const { handleChat, handleCodeChat } = createHandleChat({
|
||||
CONFIG_DIR_PATH,
|
||||
});
|
||||
|
||||
_handleTaskControlActionImpl = createHandleTaskControl({
|
||||
telegramChannel,
|
||||
makeBroadcastForTask,
|
||||
handleChat,
|
||||
}).handleTaskControlAction;
|
||||
|
||||
function createSSESender(res: express.Response): (event: string, data: any) => void {
|
||||
return (type: string, data: any) => { try { res.write(`data: ${JSON.stringify({ type, ...data })}\n\n`); } catch {} };
|
||||
}
|
||||
|
||||
const ACTIVE_TASK_STATUSES: TaskStatus[] = [
|
||||
'queued',
|
||||
'running',
|
||||
'paused',
|
||||
'stalled',
|
||||
'needs_assistance',
|
||||
'failed',
|
||||
'waiting_subagent',
|
||||
];
|
||||
|
||||
function inferTaskChannelFromSession(sessionId: string): 'web' | 'telegram' {
|
||||
return String(sessionId || '').startsWith('telegram_') ? 'telegram' : 'web';
|
||||
}
|
||||
|
||||
function latestTaskForSession(sessionId: string, statuses: TaskStatus[]): TaskRecord | null {
|
||||
const tasks = listTasks({ status: statuses })
|
||||
.filter(t => t.sessionId === sessionId)
|
||||
.sort((a, b) => b.lastProgressAt - a.lastProgressAt);
|
||||
return tasks[0] || null;
|
||||
}
|
||||
|
||||
function findBlockedTaskForSession(sessionId: string): TaskRecord | null {
|
||||
const blocked = listTasks({ status: ['needs_assistance', 'stalled', 'paused', 'failed'] })
|
||||
.filter(t => t.sessionId === sessionId)
|
||||
.filter(t =>
|
||||
t.status === 'needs_assistance'
|
||||
|| t.status === 'stalled'
|
||||
|| t.status === 'failed'
|
||||
|| (t.status === 'paused' && t.pauseReason !== 'user_pause'),
|
||||
)
|
||||
.sort((a, b) => b.lastProgressAt - a.lastProgressAt);
|
||||
return blocked[0] || null;
|
||||
}
|
||||
|
||||
function getLatestPauseContext(task: TaskRecord): { reason: string; detail: string } {
|
||||
const latestPause = [...(task.journal || [])].reverse().find((j) => j.type === 'pause');
|
||||
if (latestPause) {
|
||||
return {
|
||||
reason: String(latestPause.content || '').replace(/^Task paused for assistance:\s*/i, '').slice(0, 220),
|
||||
detail: String(latestPause.detail || '').slice(0, 420),
|
||||
};
|
||||
}
|
||||
return { reason: task.pauseReason || 'paused', detail: '' };
|
||||
}
|
||||
|
||||
function summarizeTaskRecord(task: TaskRecord): Record<string, any> {
|
||||
const total = Array.isArray(task.plan) ? task.plan.length : 0;
|
||||
const step = Math.min((task.currentStepIndex || 0) + 1, Math.max(1, total));
|
||||
const done = (task.plan || []).filter((s) => s.status === 'done' || s.status === 'skipped').length;
|
||||
const latestPause = getLatestPauseContext(task);
|
||||
return {
|
||||
task_id: task.id,
|
||||
title: task.title,
|
||||
status: task.status,
|
||||
pause_reason: task.pauseReason || null,
|
||||
step,
|
||||
total_steps: Math.max(1, total),
|
||||
completed_steps: done,
|
||||
last_issue: latestPause.reason || null,
|
||||
last_issue_detail: latestPause.detail || null,
|
||||
channel: task.channel,
|
||||
session_id: task.sessionId,
|
||||
last_progress_at: task.lastProgressAt,
|
||||
last_progress_iso: new Date(task.lastProgressAt).toISOString(),
|
||||
started_at: task.startedAt,
|
||||
started_at_iso: new Date(task.startedAt).toISOString(),
|
||||
completed_at: task.completedAt || null,
|
||||
completed_at_iso: task.completedAt ? new Date(task.completedAt).toISOString() : null,
|
||||
};
|
||||
}
|
||||
|
||||
function buildBlockedTaskStatusMessage(task: TaskRecord): string {
|
||||
const summary = summarizeTaskRecord(task);
|
||||
const lines = [
|
||||
`Task status: ${summary.title}`,
|
||||
`Status: ${summary.status}`,
|
||||
`Step: ${summary.step}/${summary.total_steps} (${summary.completed_steps} completed)`,
|
||||
summary.last_issue ? `Last issue: ${summary.last_issue}` : '',
|
||||
summary.last_issue_detail ? `Details: ${summary.last_issue_detail}` : '',
|
||||
`Task ID: ${summary.task_id}`,
|
||||
`You can say: "resume task ${summary.task_id}" or "rerun task ${summary.task_id}".`,
|
||||
];
|
||||
return lines.filter(Boolean).join('\n');
|
||||
}
|
||||
|
||||
function parseTaskStatusFilter(raw: any): TaskStatus[] | undefined {
|
||||
if (raw === undefined || raw === null || raw === '') return undefined;
|
||||
const valid = new Set<TaskStatus>([
|
||||
'queued',
|
||||
'running',
|
||||
'paused',
|
||||
'stalled',
|
||||
'needs_assistance',
|
||||
'failed',
|
||||
'complete',
|
||||
'waiting_subagent',
|
||||
]);
|
||||
const values = String(raw)
|
||||
.split(/[,\s]+/)
|
||||
.map(v => v.trim())
|
||||
.filter(Boolean) as TaskStatus[];
|
||||
const filtered = values.filter(v => valid.has(v));
|
||||
return filtered.length > 0 ? filtered : undefined;
|
||||
}
|
||||
|
||||
function getTaskScopeBuckets(sessionId: string, statuses?: TaskStatus[]) {
|
||||
const all = listTasks(statuses ? { status: statuses } : undefined).sort((a, b) => b.lastProgressAt - a.lastProgressAt);
|
||||
const sessionTasks = all.filter(t => t.sessionId === sessionId);
|
||||
const channel = inferTaskChannelFromSession(sessionId);
|
||||
const channelTasks = all.filter(t => t.channel === channel && t.sessionId !== sessionId);
|
||||
return { all, sessionTasks, channelTasks, channel };
|
||||
}
|
||||
|
||||
function parseTaskIdFromText(text: string): string | null {
|
||||
const m = String(text || '').match(/\b([a-f0-9]{8}-[a-f0-9-]{27,})\b/i);
|
||||
return m ? m[1] : null;
|
||||
}
|
||||
|
||||
function launchBackgroundTaskRunner(taskId: string): void {
|
||||
const runner = new BackgroundTaskRunner(taskId, handleChat, makeBroadcastForTask(taskId), telegramChannel);
|
||||
runner.start().catch(err => console.error(`[BackgroundTaskRunner] task_control start ${taskId} error:`, err.message));
|
||||
}
|
||||
|
||||
async function handleTaskControlAction(sessionId: string, args: any): Promise<TaskControlResponse> {
|
||||
const action = String(args?.action || '').trim().toLowerCase();
|
||||
const taskId = String(args?.task_id || args?.id || '').trim();
|
||||
const includeAllSessions = args?.include_all_sessions === true;
|
||||
const note = String(args?.note || '').trim();
|
||||
const limitRaw = Number(args?.limit);
|
||||
const limit = Number.isFinite(limitRaw) && limitRaw > 0 ? Math.min(100, Math.floor(limitRaw)) : 20;
|
||||
const statusFilter = parseTaskStatusFilter(args?.status);
|
||||
|
||||
if (!action) {
|
||||
return { success: false, action: 'unknown', code: 'invalid_action', message: 'task_control requires action.' };
|
||||
}
|
||||
|
||||
if (action === 'list' || action === 'latest') {
|
||||
const statuses = statusFilter || (action === 'list' ? ACTIVE_TASK_STATUSES : undefined);
|
||||
const scope = getTaskScopeBuckets(sessionId, statuses);
|
||||
const tasks = includeAllSessions
|
||||
? scope.all
|
||||
: [...scope.sessionTasks, ...scope.channelTasks];
|
||||
if (action === 'latest') {
|
||||
const latest = tasks[0] || null;
|
||||
return {
|
||||
success: true,
|
||||
action,
|
||||
scope: includeAllSessions ? 'all_sessions' : `session+${scope.channel}`,
|
||||
task: latest ? summarizeTaskRecord(latest) : null,
|
||||
message: latest ? `Latest task is "${latest.title}" (${latest.status}).` : 'No tasks found.',
|
||||
};
|
||||
}
|
||||
const summarized = tasks.slice(0, limit).map(summarizeTaskRecord);
|
||||
return {
|
||||
success: true,
|
||||
action,
|
||||
scope: includeAllSessions ? 'all_sessions' : `session+${scope.channel}`,
|
||||
tasks: summarized,
|
||||
message: summarized.length > 0 ? `Found ${summarized.length} task(s).` : 'No tasks found.',
|
||||
};
|
||||
}
|
||||
|
||||
if (action === 'get') {
|
||||
if (!taskId) return { success: false, action, code: 'missing_task_id', message: 'task_control(get) requires task_id.' };
|
||||
const task = loadTask(taskId);
|
||||
if (!task) return { success: false, action, code: 'not_found', message: `Task not found: ${taskId}` };
|
||||
return { success: true, action, task: summarizeTaskRecord(task), message: `Loaded task "${task.title}".` };
|
||||
}
|
||||
|
||||
const resolveCandidateForAction = (candidateAction: 'resume' | 'rerun' | 'pause' | 'cancel' | 'delete') => {
|
||||
if (taskId) {
|
||||
const exact = loadTask(taskId);
|
||||
if (!exact) return { task: null as TaskRecord | null, err: `Task not found: ${taskId}` };
|
||||
return { task: exact, err: '' };
|
||||
}
|
||||
|
||||
const preferredStatuses: TaskStatus[] =
|
||||
candidateAction === 'rerun'
|
||||
? ['needs_assistance', 'stalled', 'paused', 'failed', 'complete']
|
||||
: candidateAction === 'delete'
|
||||
? ['needs_assistance', 'stalled', 'paused', 'failed', 'queued', 'complete', 'waiting_subagent']
|
||||
// 'running' included so tasks stuck in running state (dead runner) can be resumed
|
||||
: ['needs_assistance', 'stalled', 'paused', 'failed', 'queued', 'running'];
|
||||
const scope = getTaskScopeBuckets(sessionId, preferredStatuses);
|
||||
let preferred = [...scope.sessionTasks, ...scope.channelTasks];
|
||||
if (preferred.length === 0) {
|
||||
preferred = scope.all;
|
||||
}
|
||||
if (preferred.length === 0) {
|
||||
return { task: null as TaskRecord | null, err: 'No matching task found in current scope.' };
|
||||
}
|
||||
if (preferred.length === 1) {
|
||||
return { task: preferred[0], err: '' };
|
||||
}
|
||||
return { task: null as TaskRecord | null, err: 'AMBIGUOUS', candidates: preferred.slice(0, 3) };
|
||||
};
|
||||
|
||||
if (action === 'resume' || action === 'rerun') {
|
||||
const resolved = resolveCandidateForAction(action);
|
||||
if (!resolved.task) {
|
||||
if (resolved.err === 'AMBIGUOUS') {
|
||||
return {
|
||||
success: false,
|
||||
action,
|
||||
code: 'ambiguous',
|
||||
message: 'Multiple tasks match. Provide task_id.',
|
||||
candidates: (resolved.candidates || []).map(summarizeTaskRecord),
|
||||
};
|
||||
}
|
||||
return { success: false, action, code: 'no_candidate', message: resolved.err };
|
||||
}
|
||||
const task = loadTask(resolved.task.id);
|
||||
if (!task) return { success: false, action, code: 'not_found', message: `Task not found: ${resolved.task.id}` };
|
||||
if (BackgroundTaskRunner.isRunning(task.id)) {
|
||||
return {
|
||||
success: true,
|
||||
action,
|
||||
task: summarizeTaskRecord(task),
|
||||
message: `Task "${task.title}" is already actively running (runner is live).`,
|
||||
};
|
||||
}
|
||||
|
||||
// Status is 'running' but no active runner found — runner died without cleanup.
|
||||
// Auto-correct the stale status so the resume proceeds normally.
|
||||
if (task.status === 'running') {
|
||||
appendJournal(task.id, { type: 'status_push', content: 'Stale running status detected (no active runner). Auto-correcting to paused for resume.' });
|
||||
BackgroundTaskRunner.forceRelease(task.id); // clear any ghost activeRunners entry
|
||||
updateTaskStatus(task.id, 'paused', { pauseReason: 'error' });
|
||||
task.status = 'paused'; // keep local ref in sync
|
||||
}
|
||||
|
||||
if (action === 'resume') {
|
||||
if (task.status === 'complete') {
|
||||
return { success: false, action, code: 'already_complete', message: `Task "${task.title}" is complete. Use rerun to restart.` };
|
||||
}
|
||||
updateTaskStatus(task.id, 'queued');
|
||||
appendJournal(task.id, { type: 'resume', content: `task_control resume${note ? `: ${note.slice(0, 220)}` : ''}` });
|
||||
// Reset self-heal counter so the user's manual intervention gives a fresh start
|
||||
task.selfHealAttempts = 0;
|
||||
task.resynthAttempts = 0;
|
||||
saveTask(task);
|
||||
if (note) {
|
||||
const resumeMessages = Array.isArray(task.resumeContext?.messages) ? task.resumeContext.messages : [];
|
||||
updateResumeContext(task.id, {
|
||||
messages: [
|
||||
...resumeMessages,
|
||||
{ role: 'user', content: `[TASK USER FOLLOW-UP]\n${note}`, timestamp: Date.now() },
|
||||
].slice(-80),
|
||||
});
|
||||
}
|
||||
launchBackgroundTaskRunner(task.id);
|
||||
const refreshed = loadTask(task.id) || task;
|
||||
return {
|
||||
success: true,
|
||||
action,
|
||||
task: summarizeTaskRecord(refreshed),
|
||||
message: `Resumed task "${refreshed.title}" at step ${refreshed.currentStepIndex + 1}/${Math.max(1, refreshed.plan.length)}.`,
|
||||
};
|
||||
}
|
||||
|
||||
// rerun
|
||||
task.status = 'queued';
|
||||
task.pauseReason = undefined;
|
||||
task.currentStepIndex = 0;
|
||||
task.completedAt = undefined;
|
||||
task.finalSummary = undefined;
|
||||
task.lastToolCall = undefined;
|
||||
task.lastToolCallAt = undefined;
|
||||
task.lastProgressAt = Date.now();
|
||||
task.plan = (task.plan || []).map((step, idx) => ({
|
||||
...step,
|
||||
index: idx,
|
||||
status: 'pending',
|
||||
completedAt: undefined,
|
||||
notes: undefined,
|
||||
}));
|
||||
task.resumeContext = {
|
||||
...(task.resumeContext || {
|
||||
messages: [],
|
||||
browserSessionActive: false,
|
||||
round: 0,
|
||||
orchestrationLog: [],
|
||||
}),
|
||||
messages: [],
|
||||
browserSessionActive: false,
|
||||
browserUrl: undefined,
|
||||
round: 0,
|
||||
orchestrationLog: [],
|
||||
fileOpState: undefined,
|
||||
};
|
||||
saveTask(task);
|
||||
appendJournal(task.id, { type: 'status_push', content: `task_control rerun${note ? `: ${note.slice(0, 220)}` : ''}` });
|
||||
launchBackgroundTaskRunner(task.id);
|
||||
const refreshed = loadTask(task.id) || task;
|
||||
return {
|
||||
success: true,
|
||||
action,
|
||||
task: summarizeTaskRecord(refreshed),
|
||||
message: `Rerunning task "${refreshed.title}" from step 1/${Math.max(1, refreshed.plan.length)}.`,
|
||||
};
|
||||
}
|
||||
|
||||
if (action === 'pause' || action === 'cancel') {
|
||||
const resolved = resolveCandidateForAction(action as any);
|
||||
if (!resolved.task) {
|
||||
if (resolved.err === 'AMBIGUOUS') {
|
||||
return {
|
||||
success: false,
|
||||
action,
|
||||
code: 'ambiguous',
|
||||
message: 'Multiple tasks match. Provide task_id.',
|
||||
candidates: (resolved.candidates || []).map(summarizeTaskRecord),
|
||||
};
|
||||
}
|
||||
return { success: false, action, code: 'no_candidate', message: resolved.err };
|
||||
}
|
||||
if (action === 'cancel' && args?.confirm !== true) {
|
||||
return { success: false, action, code: 'needs_confirmation', message: 'cancel requires confirm=true.' };
|
||||
}
|
||||
const task = loadTask(resolved.task.id);
|
||||
if (!task) return { success: false, action, code: 'not_found', message: `Task not found: ${resolved.task.id}` };
|
||||
if (BackgroundTaskRunner.isRunning(task.id)) {
|
||||
BackgroundTaskRunner.requestPause(task.id);
|
||||
}
|
||||
updateTaskStatus(task.id, 'paused', { pauseReason: 'user_pause' });
|
||||
appendJournal(task.id, { type: 'pause', content: `task_control ${action}${note ? `: ${note.slice(0, 220)}` : ''}` });
|
||||
const refreshed = loadTask(task.id) || task;
|
||||
return {
|
||||
success: true,
|
||||
action,
|
||||
task: summarizeTaskRecord(refreshed),
|
||||
message: `${action === 'cancel' ? 'Cancelled' : 'Paused'} task "${refreshed.title}".`,
|
||||
};
|
||||
}
|
||||
|
||||
if (action === 'delete') {
|
||||
if (args?.confirm !== true) {
|
||||
return { success: false, action, code: 'needs_confirmation', message: 'delete requires confirm=true.' };
|
||||
}
|
||||
const resolved = resolveCandidateForAction('delete');
|
||||
if (!resolved.task) {
|
||||
if (resolved.err === 'AMBIGUOUS') {
|
||||
return {
|
||||
success: false,
|
||||
action,
|
||||
code: 'ambiguous',
|
||||
message: 'Multiple tasks match. Provide task_id.',
|
||||
candidates: (resolved.candidates || []).map(summarizeTaskRecord),
|
||||
};
|
||||
}
|
||||
return { success: false, action, code: 'no_candidate', message: resolved.err };
|
||||
}
|
||||
if (BackgroundTaskRunner.isRunning(resolved.task.id)) {
|
||||
return { success: false, action, code: 'running', message: `Task "${resolved.task.title}" is running. Pause it before delete.` };
|
||||
}
|
||||
const ok = deleteTask(resolved.task.id);
|
||||
if (!ok) return { success: false, action, code: 'not_found', message: `Task not found: ${resolved.task.id}` };
|
||||
return { success: true, action, message: `Deleted task "${resolved.task.title}" (${resolved.task.id}).` };
|
||||
}
|
||||
|
||||
return { success: false, action, code: 'invalid_action', message: `Unsupported task_control action: ${action}` };
|
||||
}
|
||||
|
||||
function renderTaskCandidatesForHuman(candidates: Array<Record<string, any>>): string {
|
||||
if (!Array.isArray(candidates) || candidates.length === 0) return 'No candidates found.';
|
||||
return candidates
|
||||
|
||||
Reference in New Issue
Block a user