mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
refactor: extract helper methods to reduce duplication and improve readability
DRY up repeated patterns across 7 core files: - state-store: extract serializeState() and split assembleStateJson() into 3 focused methods - session: extract _resetBuffers(), _clearAllTimers(), _handleJsonMessage() - ralph-tracker: extract completeAllTodos() (was 4x duplicated), emitValidationWarning(), similarity constants - subagent-watcher: extract markSubagentAsCompleted(), extractFirstTextContent(), emitToolResult(), findOldestInactiveAgent() - respawn-controller: extract recoveryResetToWatching(), canAutoAccept(), formatRemainingSeconds(), validatePositiveTimeout() - tmux-manager: replace 15 path.includes() checks with single UNSAFE_PATH_CHARS regex - session-auto-ops: extract executeWhenIdle() shared retry helper for checkAutoCompact/checkAutoClear Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
+58
-66
@@ -100,6 +100,18 @@ const TODO_CLEANUP_INTERVAL_MS = INACTIVITY_TIMEOUT_MS;
|
||||
*/
|
||||
const TODO_SIMILARITY_THRESHOLD = 0.85;
|
||||
|
||||
/**
|
||||
* Similarity threshold for short todo content (<30 chars).
|
||||
* Higher threshold reduces false positive deduplication of short strings.
|
||||
*/
|
||||
const SIMILARITY_THRESHOLD_SHORT = 0.95;
|
||||
|
||||
/**
|
||||
* Similarity threshold for medium-length todo content (30-60 chars).
|
||||
* Slightly relaxed compared to short strings.
|
||||
*/
|
||||
const SIMILARITY_THRESHOLD_MEDIUM = 0.9;
|
||||
|
||||
/**
|
||||
* Debounce interval for event emissions (milliseconds).
|
||||
* Prevents UI jitter from rapid consecutive updates.
|
||||
@@ -1299,6 +1311,24 @@ export class RalphTracker extends EventEmitter {
|
||||
this.detectTodoItems(trimmed);
|
||||
}
|
||||
|
||||
/**
|
||||
* Mark all tracked todos as completed and emit todoUpdate if any changed.
|
||||
* @returns true if any todo was updated
|
||||
*/
|
||||
private completeAllTodos(): boolean {
|
||||
let updated = false;
|
||||
for (const todo of this._todos.values()) {
|
||||
if (todo.status !== 'completed') {
|
||||
todo.status = 'completed';
|
||||
updated = true;
|
||||
}
|
||||
}
|
||||
if (updated) {
|
||||
this.emit('todoUpdate', this.todos);
|
||||
}
|
||||
return updated;
|
||||
}
|
||||
|
||||
/**
|
||||
* Detect "all tasks complete" messages.
|
||||
*/
|
||||
@@ -1318,16 +1348,7 @@ export class RalphTracker extends EventEmitter {
|
||||
return;
|
||||
}
|
||||
|
||||
let updated = false;
|
||||
for (const todo of this._todos.values()) {
|
||||
if (todo.status !== 'completed') {
|
||||
todo.status = 'completed';
|
||||
updated = true;
|
||||
}
|
||||
}
|
||||
if (updated) {
|
||||
this.emit('todoUpdate', this.todos);
|
||||
}
|
||||
this.completeAllTodos();
|
||||
|
||||
if (this._loopState.completionPhrase) {
|
||||
this._loopState.active = false;
|
||||
@@ -1425,16 +1446,7 @@ export class RalphTracker extends EventEmitter {
|
||||
|
||||
if (bareCount > 1) return;
|
||||
|
||||
let updated = false;
|
||||
for (const todo of this._todos.values()) {
|
||||
if (todo.status !== 'completed') {
|
||||
todo.status = 'completed';
|
||||
updated = true;
|
||||
}
|
||||
}
|
||||
if (updated) {
|
||||
this.emit('todoUpdate', this.todos);
|
||||
}
|
||||
this.completeAllTodos();
|
||||
|
||||
this._loopState.active = false;
|
||||
this._loopState.lastActivity = Date.now();
|
||||
@@ -1480,16 +1492,7 @@ export class RalphTracker extends EventEmitter {
|
||||
if (canonicalCount >= 2 || this._loopState.active) {
|
||||
this._loopState.active = false;
|
||||
this._loopState.lastActivity = Date.now();
|
||||
let updated = false;
|
||||
for (const todo of this._todos.values()) {
|
||||
if (todo.status !== 'completed') {
|
||||
todo.status = 'completed';
|
||||
updated = true;
|
||||
}
|
||||
}
|
||||
if (updated) {
|
||||
this.emit('todoUpdate', this.todos);
|
||||
}
|
||||
this.completeAllTodos();
|
||||
this.emit('completionDetected', matchedPhrase);
|
||||
this.emit('loopUpdate', this.loopState);
|
||||
return;
|
||||
@@ -1497,16 +1500,7 @@ export class RalphTracker extends EventEmitter {
|
||||
}
|
||||
|
||||
if (this._loopState.active || count >= 2) {
|
||||
let updated = false;
|
||||
for (const todo of this._todos.values()) {
|
||||
if (todo.status !== 'completed') {
|
||||
todo.status = 'completed';
|
||||
updated = true;
|
||||
}
|
||||
}
|
||||
if (updated) {
|
||||
this.emit('todoUpdate', this.todos);
|
||||
}
|
||||
this.completeAllTodos();
|
||||
|
||||
this._loopState.active = false;
|
||||
this._loopState.lastActivity = Date.now();
|
||||
@@ -1532,41 +1526,39 @@ export class RalphTracker extends EventEmitter {
|
||||
const suggestedPhrase = `${phrase}_${uniqueSuffix}`;
|
||||
|
||||
if (COMMON_COMPLETION_PHRASES.has(normalized)) {
|
||||
console.warn(
|
||||
`[RalphTracker] Warning: Completion phrase "${phrase}" is very common and may cause false positives. Consider using: "${suggestedPhrase}"`
|
||||
);
|
||||
this.emit('phraseValidationWarning', {
|
||||
phrase,
|
||||
reason: 'common',
|
||||
suggestedPhrase,
|
||||
});
|
||||
this.emitValidationWarning(phrase, 'common', suggestedPhrase);
|
||||
return;
|
||||
}
|
||||
|
||||
if (normalized.length < MIN_RECOMMENDED_PHRASE_LENGTH) {
|
||||
console.warn(
|
||||
`[RalphTracker] Warning: Completion phrase "${phrase}" is too short (${normalized.length} chars). Consider using: "${suggestedPhrase}"`
|
||||
);
|
||||
this.emit('phraseValidationWarning', {
|
||||
phrase,
|
||||
reason: 'short',
|
||||
suggestedPhrase,
|
||||
});
|
||||
this.emitValidationWarning(phrase, 'short', suggestedPhrase);
|
||||
return;
|
||||
}
|
||||
|
||||
if (/^\d+$/.test(normalized)) {
|
||||
console.warn(
|
||||
`[RalphTracker] Warning: Completion phrase "${phrase}" is numeric-only and may cause false positives. Consider using: "${suggestedPhrase}"`
|
||||
);
|
||||
this.emit('phraseValidationWarning', {
|
||||
phrase,
|
||||
reason: 'numeric',
|
||||
suggestedPhrase,
|
||||
});
|
||||
this.emitValidationWarning(phrase, 'numeric', suggestedPhrase);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Emit a phrase validation warning with a console message and event.
|
||||
*/
|
||||
private emitValidationWarning(phrase: string, reason: 'common' | 'short' | 'numeric', suggestedPhrase: string): void {
|
||||
const descriptions: Record<'common' | 'short' | 'numeric', string> = {
|
||||
common: 'is very common and may cause false positives',
|
||||
short: `is too short (${phrase.toUpperCase().replace(/[\s_\-.]+/g, '').length} chars)`,
|
||||
numeric: 'is numeric-only and may cause false positives',
|
||||
};
|
||||
console.warn(
|
||||
`[RalphTracker] Warning: Completion phrase "${phrase}" ${descriptions[reason]}. Consider using: "${suggestedPhrase}"`
|
||||
);
|
||||
this.emit('phraseValidationWarning', {
|
||||
phrase,
|
||||
reason,
|
||||
suggestedPhrase,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Activate the loop if not already active.
|
||||
*/
|
||||
@@ -1977,9 +1969,9 @@ export class RalphTracker extends EventEmitter {
|
||||
|
||||
let threshold: number;
|
||||
if (normalized.length < 30) {
|
||||
threshold = 0.95;
|
||||
threshold = SIMILARITY_THRESHOLD_SHORT;
|
||||
} else if (normalized.length < 60) {
|
||||
threshold = 0.9;
|
||||
threshold = SIMILARITY_THRESHOLD_MEDIUM;
|
||||
} else {
|
||||
threshold = TODO_SIMILARITY_THRESHOLD;
|
||||
}
|
||||
|
||||
+89
-51
@@ -551,6 +551,14 @@ export interface RespawnEvents {
|
||||
respawnBlocked: (data: { reason: string; details: string }) => void;
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert milliseconds to a non-negative whole number of seconds for countdown display.
|
||||
* Rounds up so that e.g. 1200 ms shows as 2 s (never under-reports remaining time).
|
||||
*/
|
||||
function formatRemainingSeconds(ms: number): number {
|
||||
return Math.max(0, Math.ceil(ms / 1000));
|
||||
}
|
||||
|
||||
/** Default configuration values */
|
||||
const DEFAULT_CONFIG: RespawnConfig = {
|
||||
idleTimeoutMs: 10000, // 10 seconds of no activity after prompt (legacy, still used as fallback)
|
||||
@@ -839,12 +847,25 @@ export class RespawnController extends EventEmitter {
|
||||
private validateConfig(): void {
|
||||
const c = this.config;
|
||||
|
||||
/**
|
||||
* Validate that a timeout value is positive (or non-negative when allowZero is true).
|
||||
* Falls back to the DEFAULT_CONFIG value if invalid.
|
||||
*/
|
||||
const validatePositiveTimeout = (field: keyof RespawnConfig, allowZero = false): void => {
|
||||
const value = c[field] as number;
|
||||
const invalid = allowZero ? value < 0 : value <= 0;
|
||||
if (invalid) {
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
(c as any)[field] = DEFAULT_CONFIG[field];
|
||||
}
|
||||
};
|
||||
|
||||
// Ensure timeouts are positive
|
||||
if (c.idleTimeoutMs <= 0) c.idleTimeoutMs = DEFAULT_CONFIG.idleTimeoutMs;
|
||||
if (c.completionConfirmMs <= 0) c.completionConfirmMs = DEFAULT_CONFIG.completionConfirmMs;
|
||||
if (c.noOutputTimeoutMs <= 0) c.noOutputTimeoutMs = DEFAULT_CONFIG.noOutputTimeoutMs;
|
||||
if (c.autoAcceptDelayMs < 0) c.autoAcceptDelayMs = DEFAULT_CONFIG.autoAcceptDelayMs;
|
||||
if (c.interStepDelayMs <= 0) c.interStepDelayMs = DEFAULT_CONFIG.interStepDelayMs;
|
||||
validatePositiveTimeout('idleTimeoutMs');
|
||||
validatePositiveTimeout('completionConfirmMs');
|
||||
validatePositiveTimeout('noOutputTimeoutMs');
|
||||
validatePositiveTimeout('autoAcceptDelayMs', true);
|
||||
validatePositiveTimeout('interStepDelayMs');
|
||||
|
||||
// Ensure completion confirm doesn't exceed no-output timeout
|
||||
if (c.completionConfirmMs > c.noOutputTimeoutMs) {
|
||||
@@ -852,14 +873,14 @@ export class RespawnController extends EventEmitter {
|
||||
}
|
||||
|
||||
// Ensure AI check timeouts are positive
|
||||
if (c.aiIdleCheckTimeoutMs <= 0) c.aiIdleCheckTimeoutMs = DEFAULT_CONFIG.aiIdleCheckTimeoutMs;
|
||||
if (c.aiIdleCheckCooldownMs < 0) c.aiIdleCheckCooldownMs = DEFAULT_CONFIG.aiIdleCheckCooldownMs;
|
||||
if (c.aiIdleCheckMaxContext <= 0) c.aiIdleCheckMaxContext = DEFAULT_CONFIG.aiIdleCheckMaxContext;
|
||||
validatePositiveTimeout('aiIdleCheckTimeoutMs');
|
||||
validatePositiveTimeout('aiIdleCheckCooldownMs', true);
|
||||
validatePositiveTimeout('aiIdleCheckMaxContext');
|
||||
|
||||
// Ensure plan check timeouts are positive
|
||||
if (c.aiPlanCheckTimeoutMs <= 0) c.aiPlanCheckTimeoutMs = DEFAULT_CONFIG.aiPlanCheckTimeoutMs;
|
||||
if (c.aiPlanCheckCooldownMs < 0) c.aiPlanCheckCooldownMs = DEFAULT_CONFIG.aiPlanCheckCooldownMs;
|
||||
if (c.aiPlanCheckMaxContext <= 0) c.aiPlanCheckMaxContext = DEFAULT_CONFIG.aiPlanCheckMaxContext;
|
||||
validatePositiveTimeout('aiPlanCheckTimeoutMs');
|
||||
validatePositiveTimeout('aiPlanCheckCooldownMs', true);
|
||||
validatePositiveTimeout('aiPlanCheckMaxContext');
|
||||
}
|
||||
|
||||
/** Wire up AI checker events to controller events (removes existing listeners first to prevent duplicates) */
|
||||
@@ -987,11 +1008,11 @@ export class RespawnController extends EventEmitter {
|
||||
waitingFor = 'AI verdict (IDLE or WORKING)';
|
||||
} else if (this._state === 'confirming_idle') {
|
||||
statusText = `Confirming idle (${confidence}% confidence)`;
|
||||
waitingFor = `${Math.max(0, Math.ceil((this.config.completionConfirmMs - msSinceLastOutput) / 1000))}s more silence`;
|
||||
waitingFor = `${formatRemainingSeconds(this.config.completionConfirmMs - msSinceLastOutput)}s more silence`;
|
||||
} else if (this._state === 'watching') {
|
||||
const aiState = this.aiChecker.getState();
|
||||
if (aiState.status === 'cooldown') {
|
||||
const remaining = Math.ceil(this.aiChecker.getCooldownRemainingMs() / 1000);
|
||||
const remaining = formatRemainingSeconds(this.aiChecker.getCooldownRemainingMs());
|
||||
statusText = `AI Check: WORKING (cooldown ${remaining}s)`;
|
||||
waitingFor = 'Cooldown to expire';
|
||||
} else if (completionMessageDetected) {
|
||||
@@ -1798,24 +1819,28 @@ export class RespawnController extends EventEmitter {
|
||||
case 'sending_init':
|
||||
case 'sending_kickstart':
|
||||
// For sending states, retry the send
|
||||
this.log('Recovery: returning to watching state');
|
||||
this.setState('watching');
|
||||
this.startNoOutputTimer();
|
||||
this.startPreFilterTimer();
|
||||
if (this.config.autoAcceptPrompts) {
|
||||
this.startAutoAcceptTimer();
|
||||
}
|
||||
this.recoveryResetToWatching('returning to watching state');
|
||||
break;
|
||||
|
||||
default:
|
||||
// Fallback: reset to watching
|
||||
this.log('Recovery: fallback to watching state');
|
||||
this.setState('watching');
|
||||
this.startNoOutputTimer();
|
||||
this.startPreFilterTimer();
|
||||
if (this.config.autoAcceptPrompts) {
|
||||
this.startAutoAcceptTimer();
|
||||
}
|
||||
this.recoveryResetToWatching('fallback to watching state');
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reset the controller to watching state during stuck-state recovery.
|
||||
* Sets state to watching and restarts all detection timers.
|
||||
*
|
||||
* @param reason - Human-readable reason for the reset (logged)
|
||||
*/
|
||||
private recoveryResetToWatching(reason: string): void {
|
||||
this.log(`Recovery: ${reason}`);
|
||||
this.setState('watching');
|
||||
this.startNoOutputTimer();
|
||||
this.startPreFilterTimer();
|
||||
if (this.config.autoAcceptPrompts) {
|
||||
this.startAutoAcceptTimer();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2051,7 +2076,7 @@ export class RespawnController extends EventEmitter {
|
||||
// If on cooldown, don't start check - wait for cooldown to expire
|
||||
if (this.aiChecker.isOnCooldown()) {
|
||||
this.log(
|
||||
`AI check on cooldown (${Math.ceil(this.aiChecker.getCooldownRemainingMs() / 1000)}s remaining), waiting...`
|
||||
`AI check on cooldown (${formatRemainingSeconds(this.aiChecker.getCooldownRemainingMs())}s remaining), waiting...`
|
||||
);
|
||||
return;
|
||||
}
|
||||
@@ -2198,36 +2223,15 @@ export class RespawnController extends EventEmitter {
|
||||
* @fires planCheckStarted
|
||||
*/
|
||||
private tryAutoAccept(): void {
|
||||
// Only auto-accept in watching state (not during a respawn cycle)
|
||||
if (this._state !== 'watching') return;
|
||||
if (!this.canAutoAccept()) return;
|
||||
|
||||
// Don't auto-accept if a completion message was detected (normal idle handles it)
|
||||
if (this.completionMessageTime !== null) return;
|
||||
|
||||
// Don't auto-accept if disabled
|
||||
if (!this.config.autoAcceptPrompts) return;
|
||||
|
||||
// Don't auto-accept if we haven't received any output yet (prevents spurious Enter on fresh start)
|
||||
if (!this.hasReceivedOutput) return;
|
||||
|
||||
// Don't auto-accept if an elicitation dialog (AskUserQuestion) was detected
|
||||
if (this.elicitationDetected) {
|
||||
this.log('Skipping auto-accept: elicitation dialog detected (AskUserQuestion)');
|
||||
return;
|
||||
}
|
||||
|
||||
// Stage 1: Pre-filter — check if buffer looks like plan mode
|
||||
const buffer = this.terminalBuffer.value;
|
||||
if (!this.isPlanModePreFilterMatch(buffer)) {
|
||||
this.log('Skipping auto-accept: pre-filter did not match plan mode patterns');
|
||||
return;
|
||||
}
|
||||
|
||||
// Stage 2: AI confirmation (if enabled and available)
|
||||
if (this.config.aiPlanCheckEnabled && this.planChecker.status !== 'disabled') {
|
||||
if (this.planChecker.isOnCooldown()) {
|
||||
this.log(
|
||||
`Skipping auto-accept: plan checker on cooldown (${Math.ceil(this.planChecker.getCooldownRemainingMs() / 1000)}s remaining)`
|
||||
`Skipping auto-accept: plan checker on cooldown (${formatRemainingSeconds(this.planChecker.getCooldownRemainingMs())}s remaining)`
|
||||
);
|
||||
return;
|
||||
}
|
||||
@@ -2244,6 +2248,40 @@ export class RespawnController extends EventEmitter {
|
||||
this.sendAutoAcceptEnter();
|
||||
}
|
||||
|
||||
/**
|
||||
* Check whether all preconditions for auto-accept are met.
|
||||
* Validates state, config, and pre-filter conditions before attempting auto-accept.
|
||||
*
|
||||
* @returns True if auto-accept should proceed to the AI confirmation stage
|
||||
*/
|
||||
private canAutoAccept(): boolean {
|
||||
// Only auto-accept in watching state (not during a respawn cycle)
|
||||
if (this._state !== 'watching') return false;
|
||||
|
||||
// Don't auto-accept if a completion message was detected (normal idle handles it)
|
||||
if (this.completionMessageTime !== null) return false;
|
||||
|
||||
// Don't auto-accept if disabled
|
||||
if (!this.config.autoAcceptPrompts) return false;
|
||||
|
||||
// Don't auto-accept if we haven't received any output yet (prevents spurious Enter on fresh start)
|
||||
if (!this.hasReceivedOutput) return false;
|
||||
|
||||
// Don't auto-accept if an elicitation dialog (AskUserQuestion) was detected
|
||||
if (this.elicitationDetected) {
|
||||
this.log('Skipping auto-accept: elicitation dialog detected (AskUserQuestion)');
|
||||
return false;
|
||||
}
|
||||
|
||||
// Stage 1: Pre-filter — check if buffer looks like plan mode
|
||||
if (!this.isPlanModePreFilterMatch(this.terminalBuffer.value)) {
|
||||
this.log('Skipping auto-accept: pre-filter did not match plan mode patterns');
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if the terminal buffer matches plan mode pre-filter patterns.
|
||||
* Only checks the last 2000 chars (plan mode UI appears at the bottom).
|
||||
|
||||
+98
-51
@@ -27,6 +27,57 @@ const COMPACT_COOLDOWN_MS = 10000;
|
||||
/** Cooldown after clear completes before re-enabling (5 seconds) */
|
||||
const CLEAR_COOLDOWN_MS = 5000;
|
||||
|
||||
/**
|
||||
* Executes an action when the session becomes idle, retrying if currently working.
|
||||
*
|
||||
* @param action - The async action to execute once idle
|
||||
* @param isActive - Returns whether this operation is still active (not cancelled)
|
||||
* @param isWorking - Returns whether the session is currently working
|
||||
* @param isStopped - Returns whether the session has been stopped
|
||||
* @param retryMs - Delay between retry attempts when working
|
||||
* @param cooldownMs - Delay after action completes before calling onCooldownDone
|
||||
* @param setTimer - Stores the timer reference for cleanup
|
||||
* @param onCooldownDone - Called after cooldown to reset state
|
||||
*/
|
||||
async function executeWhenIdle(
|
||||
action: () => Promise<void>,
|
||||
isActive: () => boolean,
|
||||
isWorking: () => boolean,
|
||||
isStopped: () => boolean,
|
||||
retryMs: number,
|
||||
cooldownMs: number,
|
||||
setTimer: (timer: NodeJS.Timeout | null) => void,
|
||||
onCooldownDone: () => void
|
||||
): Promise<void> {
|
||||
if (isStopped()) return;
|
||||
if (!isActive()) return;
|
||||
|
||||
if (!isWorking()) {
|
||||
if (isStopped()) return;
|
||||
|
||||
await action();
|
||||
|
||||
if (!isStopped()) {
|
||||
setTimer(
|
||||
setTimeout(() => {
|
||||
if (isStopped()) return;
|
||||
setTimer(null);
|
||||
onCooldownDone();
|
||||
}, cooldownMs)
|
||||
);
|
||||
}
|
||||
} else {
|
||||
if (!isStopped()) {
|
||||
setTimer(
|
||||
setTimeout(
|
||||
() => executeWhenIdle(action, isActive, isWorking, isStopped, retryMs, cooldownMs, setTimer, onCooldownDone),
|
||||
retryMs
|
||||
)
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Minimum valid threshold for auto-clear/compact (1000 tokens) */
|
||||
const MIN_AUTO_THRESHOLD = 1000;
|
||||
|
||||
@@ -181,37 +232,35 @@ export class SessionAutoOps extends EventEmitter {
|
||||
`[SessionAutoOps] Auto-compact triggered: ${totalTokens} tokens >= ${this._autoCompactThreshold} threshold`
|
||||
);
|
||||
|
||||
const checkAndCompact = async () => {
|
||||
if (this.callbacks.isStopped()) return;
|
||||
if (!this._isCompacting) return;
|
||||
|
||||
if (!this.callbacks.isWorking()) {
|
||||
if (this.callbacks.isStopped()) return;
|
||||
|
||||
const compactCmd = this._autoCompactPrompt ? `/compact ${this._autoCompactPrompt}\r` : '/compact\r';
|
||||
await this.callbacks.writeCommand(compactCmd);
|
||||
this.emit('autoCompact', {
|
||||
tokens: totalTokens,
|
||||
threshold: this._autoCompactThreshold,
|
||||
prompt: this._autoCompactPrompt || undefined,
|
||||
});
|
||||
|
||||
if (!this.callbacks.isStopped()) {
|
||||
this._autoCompactTimer = setTimeout(() => {
|
||||
if (this.callbacks.isStopped()) return;
|
||||
this._autoCompactTimer = null;
|
||||
this._isCompacting = false;
|
||||
}, COMPACT_COOLDOWN_MS);
|
||||
}
|
||||
} else {
|
||||
if (!this.callbacks.isStopped()) {
|
||||
this._autoCompactTimer = setTimeout(checkAndCompact, AUTO_RETRY_DELAY_MS);
|
||||
}
|
||||
}
|
||||
const action = async () => {
|
||||
const compactCmd = this._autoCompactPrompt ? `/compact ${this._autoCompactPrompt}\r` : '/compact\r';
|
||||
await this.callbacks.writeCommand(compactCmd);
|
||||
this.emit('autoCompact', {
|
||||
tokens: totalTokens,
|
||||
threshold: this._autoCompactThreshold,
|
||||
prompt: this._autoCompactPrompt || undefined,
|
||||
});
|
||||
};
|
||||
|
||||
if (!this.callbacks.isStopped()) {
|
||||
this._autoCompactTimer = setTimeout(checkAndCompact, AUTO_INITIAL_DELAY_MS);
|
||||
this._autoCompactTimer = setTimeout(
|
||||
() =>
|
||||
executeWhenIdle(
|
||||
action,
|
||||
() => this._isCompacting,
|
||||
() => this.callbacks.isWorking(),
|
||||
() => this.callbacks.isStopped(),
|
||||
AUTO_RETRY_DELAY_MS,
|
||||
COMPACT_COOLDOWN_MS,
|
||||
(timer) => {
|
||||
this._autoCompactTimer = timer;
|
||||
},
|
||||
() => {
|
||||
this._isCompacting = false;
|
||||
}
|
||||
),
|
||||
AUTO_INITIAL_DELAY_MS
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -231,32 +280,30 @@ export class SessionAutoOps extends EventEmitter {
|
||||
`[SessionAutoOps] Auto-clear triggered: ${totalTokens} tokens >= ${this._autoClearThreshold} threshold`
|
||||
);
|
||||
|
||||
const checkAndClear = async () => {
|
||||
if (this.callbacks.isStopped()) return;
|
||||
if (!this._isClearing) return;
|
||||
|
||||
if (!this.callbacks.isWorking()) {
|
||||
if (this.callbacks.isStopped()) return;
|
||||
|
||||
await this.callbacks.writeCommand('/clear\r');
|
||||
this.emit('autoClear', { tokens: totalTokens, threshold: this._autoClearThreshold });
|
||||
|
||||
if (!this.callbacks.isStopped()) {
|
||||
this._autoClearTimer = setTimeout(() => {
|
||||
if (this.callbacks.isStopped()) return;
|
||||
this._autoClearTimer = null;
|
||||
this._isClearing = false;
|
||||
}, CLEAR_COOLDOWN_MS);
|
||||
}
|
||||
} else {
|
||||
if (!this.callbacks.isStopped()) {
|
||||
this._autoClearTimer = setTimeout(checkAndClear, AUTO_RETRY_DELAY_MS);
|
||||
}
|
||||
}
|
||||
const action = async () => {
|
||||
await this.callbacks.writeCommand('/clear\r');
|
||||
this.emit('autoClear', { tokens: totalTokens, threshold: this._autoClearThreshold });
|
||||
};
|
||||
|
||||
if (!this.callbacks.isStopped()) {
|
||||
this._autoClearTimer = setTimeout(checkAndClear, AUTO_INITIAL_DELAY_MS);
|
||||
this._autoClearTimer = setTimeout(
|
||||
() =>
|
||||
executeWhenIdle(
|
||||
action,
|
||||
() => this._isClearing,
|
||||
() => this.callbacks.isWorking(),
|
||||
() => this.callbacks.isStopped(),
|
||||
AUTO_RETRY_DELAY_MS,
|
||||
CLEAR_COOLDOWN_MS,
|
||||
(timer) => {
|
||||
this._autoClearTimer = timer;
|
||||
},
|
||||
() => {
|
||||
this._isClearing = false;
|
||||
}
|
||||
),
|
||||
AUTO_INITIAL_DELAY_MS
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+116
-116
@@ -876,13 +876,7 @@ export class Session extends EventEmitter {
|
||||
throw new Error('Session already has a running process');
|
||||
}
|
||||
|
||||
this._status = 'busy';
|
||||
this._terminalBuffer.clear();
|
||||
this._textOutput.clear();
|
||||
this._errorBuffer = '';
|
||||
this._messages = [];
|
||||
this._lineBuffer = '';
|
||||
this._lastActivityAt = Date.now();
|
||||
this._resetBuffers();
|
||||
|
||||
const modeLabel = this.mode === 'opencode' ? 'OpenCode' : 'Claude';
|
||||
console.log(
|
||||
@@ -1257,13 +1251,7 @@ export class Session extends EventEmitter {
|
||||
throw new Error('Session already has a running process');
|
||||
}
|
||||
|
||||
this._status = 'busy';
|
||||
this._terminalBuffer.clear();
|
||||
this._textOutput.clear();
|
||||
this._errorBuffer = '';
|
||||
this._messages = [];
|
||||
this._lineBuffer = '';
|
||||
this._lastActivityAt = Date.now();
|
||||
this._resetBuffers();
|
||||
|
||||
// Use user's default shell or bash
|
||||
const shell = process.env.SHELL || '/bin/bash';
|
||||
@@ -1448,13 +1436,7 @@ export class Session extends EventEmitter {
|
||||
return;
|
||||
}
|
||||
|
||||
this._status = 'busy';
|
||||
this._terminalBuffer.clear();
|
||||
this._textOutput.clear();
|
||||
this._errorBuffer = '';
|
||||
this._messages = [];
|
||||
this._lineBuffer = '';
|
||||
this._lastActivityAt = Date.now();
|
||||
this._resetBuffers();
|
||||
this._promptResolved = false; // Reset race condition guard
|
||||
|
||||
this.resolvePromise = resolve;
|
||||
@@ -1565,6 +1547,117 @@ export class Session extends EventEmitter {
|
||||
});
|
||||
}
|
||||
|
||||
private _resetBuffers(): void {
|
||||
this._status = 'busy';
|
||||
this._terminalBuffer.clear();
|
||||
this._textOutput.clear();
|
||||
this._errorBuffer = '';
|
||||
this._messages = [];
|
||||
this._lineBuffer = '';
|
||||
this._lastActivityAt = Date.now();
|
||||
}
|
||||
|
||||
private _clearAllTimers(): void {
|
||||
// Clear activity timeout to prevent memory leak
|
||||
if (this.activityTimeout) {
|
||||
clearTimeout(this.activityTimeout);
|
||||
this.activityTimeout = null;
|
||||
}
|
||||
|
||||
// Clear line buffer flush timer
|
||||
if (this._lineBufferFlushTimer) {
|
||||
clearTimeout(this._lineBufferFlushTimer);
|
||||
this._lineBufferFlushTimer = null;
|
||||
}
|
||||
|
||||
// Destroy auto-compact/auto-clear automation (clears its timers)
|
||||
this._autoOps.destroy();
|
||||
|
||||
// Clear prompt check timers
|
||||
if (this._promptCheckInterval) {
|
||||
clearInterval(this._promptCheckInterval);
|
||||
this._promptCheckInterval = null;
|
||||
}
|
||||
if (this._promptCheckTimeout) {
|
||||
clearTimeout(this._promptCheckTimeout);
|
||||
this._promptCheckTimeout = null;
|
||||
}
|
||||
|
||||
// Clear shell idle timer
|
||||
if (this._shellIdleTimer) {
|
||||
clearTimeout(this._shellIdleTimer);
|
||||
this._shellIdleTimer = null;
|
||||
}
|
||||
|
||||
// Clear expensive processing timer
|
||||
if (this._expensiveProcessTimer) {
|
||||
clearTimeout(this._expensiveProcessTimer);
|
||||
this._expensiveProcessTimer = null;
|
||||
}
|
||||
this._pendingCleanData = '';
|
||||
}
|
||||
|
||||
private _handleJsonMessage(cleanLine: string, rawLine: string): void {
|
||||
try {
|
||||
const msg = JSON.parse(cleanLine) as ClaudeMessage;
|
||||
this._messages.push(msg);
|
||||
this.emit('message', msg);
|
||||
|
||||
// Trim messages array for long-running sessions
|
||||
if (this._messages.length > MAX_MESSAGES) {
|
||||
this._messages = this._messages.slice(-Math.floor(MAX_MESSAGES * 0.8));
|
||||
}
|
||||
|
||||
// Extract Claude session ID from messages (can be in any message type)
|
||||
// Support both sessionId (camelCase) and session_id (snake_case)
|
||||
const msgSessionId =
|
||||
((msg as unknown as Record<string, unknown>).sessionId as string | undefined) ?? msg.session_id;
|
||||
if (msgSessionId && !this._claudeSessionId) {
|
||||
this._claudeSessionId = msgSessionId;
|
||||
}
|
||||
|
||||
// Process message for task tracking
|
||||
this._taskTracker.processMessage(msg);
|
||||
|
||||
if (msg.type === 'assistant' && msg.message?.content) {
|
||||
for (const block of msg.message.content) {
|
||||
if (block.type === 'text' && block.text) {
|
||||
this._textOutput.append(block.text);
|
||||
}
|
||||
}
|
||||
// Track tokens from usage (with validation)
|
||||
if (msg.message.usage) {
|
||||
const inputDelta = msg.message.usage.input_tokens || 0;
|
||||
const outputDelta = msg.message.usage.output_tokens || 0;
|
||||
|
||||
// Sanity check: max 100k tokens per message (generous limit)
|
||||
const MAX_TOKENS_PER_MESSAGE = 100_000;
|
||||
if (inputDelta > 0 && inputDelta <= MAX_TOKENS_PER_MESSAGE) {
|
||||
this._totalInputTokens += inputDelta;
|
||||
}
|
||||
if (outputDelta > 0 && outputDelta <= MAX_TOKENS_PER_MESSAGE) {
|
||||
this._totalOutputTokens += outputDelta;
|
||||
}
|
||||
|
||||
// Check if we should auto-compact or auto-clear
|
||||
this._autoOps.checkAutoCompact();
|
||||
this._autoOps.checkAutoClear();
|
||||
}
|
||||
}
|
||||
|
||||
if (msg.type === 'result' && msg.total_cost_usd) {
|
||||
this._totalCost = msg.total_cost_usd;
|
||||
}
|
||||
} catch (parseErr) {
|
||||
// Not JSON, just regular output - this is expected for non-JSON lines
|
||||
console.debug(
|
||||
'[Session] Line not JSON (expected for text output):',
|
||||
parseErr instanceof Error ? parseErr.message : parseErr
|
||||
);
|
||||
this._textOutput.append(rawLine + '\n');
|
||||
}
|
||||
}
|
||||
|
||||
private processOutput(data: string): void {
|
||||
// Early return if session is stopped to prevent any processing or timer creation
|
||||
if (this._isStopped) return;
|
||||
@@ -1606,64 +1699,7 @@ export class Session extends EventEmitter {
|
||||
const cleanLine = trimmed.replace(ANSI_ESCAPE_PATTERN_FULL, '');
|
||||
|
||||
if (cleanLine.startsWith('{') && cleanLine.endsWith('}')) {
|
||||
try {
|
||||
const msg = JSON.parse(cleanLine) as ClaudeMessage;
|
||||
this._messages.push(msg);
|
||||
this.emit('message', msg);
|
||||
|
||||
// Trim messages array for long-running sessions
|
||||
if (this._messages.length > MAX_MESSAGES) {
|
||||
this._messages = this._messages.slice(-Math.floor(MAX_MESSAGES * 0.8));
|
||||
}
|
||||
|
||||
// Extract Claude session ID from messages (can be in any message type)
|
||||
// Support both sessionId (camelCase) and session_id (snake_case)
|
||||
const msgSessionId =
|
||||
((msg as unknown as Record<string, unknown>).sessionId as string | undefined) ?? msg.session_id;
|
||||
if (msgSessionId && !this._claudeSessionId) {
|
||||
this._claudeSessionId = msgSessionId;
|
||||
}
|
||||
|
||||
// Process message for task tracking
|
||||
this._taskTracker.processMessage(msg);
|
||||
|
||||
if (msg.type === 'assistant' && msg.message?.content) {
|
||||
for (const block of msg.message.content) {
|
||||
if (block.type === 'text' && block.text) {
|
||||
this._textOutput.append(block.text);
|
||||
}
|
||||
}
|
||||
// Track tokens from usage (with validation)
|
||||
if (msg.message.usage) {
|
||||
const inputDelta = msg.message.usage.input_tokens || 0;
|
||||
const outputDelta = msg.message.usage.output_tokens || 0;
|
||||
|
||||
// Sanity check: max 100k tokens per message (generous limit)
|
||||
const MAX_TOKENS_PER_MESSAGE = 100_000;
|
||||
if (inputDelta > 0 && inputDelta <= MAX_TOKENS_PER_MESSAGE) {
|
||||
this._totalInputTokens += inputDelta;
|
||||
}
|
||||
if (outputDelta > 0 && outputDelta <= MAX_TOKENS_PER_MESSAGE) {
|
||||
this._totalOutputTokens += outputDelta;
|
||||
}
|
||||
|
||||
// Check if we should auto-compact or auto-clear
|
||||
this._autoOps.checkAutoCompact();
|
||||
this._autoOps.checkAutoClear();
|
||||
}
|
||||
}
|
||||
|
||||
if (msg.type === 'result' && msg.total_cost_usd) {
|
||||
this._totalCost = msg.total_cost_usd;
|
||||
}
|
||||
} catch (parseErr) {
|
||||
// Not JSON, just regular output - this is expected for non-JSON lines
|
||||
console.debug(
|
||||
'[Session] Line not JSON (expected for text output):',
|
||||
parseErr instanceof Error ? parseErr.message : parseErr
|
||||
);
|
||||
this._textOutput.append(line + '\n');
|
||||
}
|
||||
this._handleJsonMessage(cleanLine, line);
|
||||
} else if (trimmed) {
|
||||
this._textOutput.append(line + '\n');
|
||||
}
|
||||
@@ -2030,43 +2066,7 @@ export class Session extends EventEmitter {
|
||||
// Set stopped flag first to prevent new timers from being created
|
||||
this._isStopped = true;
|
||||
|
||||
// Clear activity timeout to prevent memory leak
|
||||
if (this.activityTimeout) {
|
||||
clearTimeout(this.activityTimeout);
|
||||
this.activityTimeout = null;
|
||||
}
|
||||
|
||||
// Clear line buffer flush timer
|
||||
if (this._lineBufferFlushTimer) {
|
||||
clearTimeout(this._lineBufferFlushTimer);
|
||||
this._lineBufferFlushTimer = null;
|
||||
}
|
||||
|
||||
// Destroy auto-compact/auto-clear automation (clears its timers)
|
||||
this._autoOps.destroy();
|
||||
|
||||
// Clear prompt check timers
|
||||
if (this._promptCheckInterval) {
|
||||
clearInterval(this._promptCheckInterval);
|
||||
this._promptCheckInterval = null;
|
||||
}
|
||||
if (this._promptCheckTimeout) {
|
||||
clearTimeout(this._promptCheckTimeout);
|
||||
this._promptCheckTimeout = null;
|
||||
}
|
||||
|
||||
// Clear shell idle timer
|
||||
if (this._shellIdleTimer) {
|
||||
clearTimeout(this._shellIdleTimer);
|
||||
this._shellIdleTimer = null;
|
||||
}
|
||||
|
||||
// Clear expensive processing timer
|
||||
if (this._expensiveProcessTimer) {
|
||||
clearTimeout(this._expensiveProcessTimer);
|
||||
this._expensiveProcessTimer = null;
|
||||
}
|
||||
this._pendingCleanData = '';
|
||||
this._clearAllTimers();
|
||||
|
||||
// Immediately cleanup Promise callbacks to prevent orphaned references
|
||||
// during the rest of stop() processing (e.g., if mux kill times out)
|
||||
|
||||
+48
-50
@@ -195,16 +195,7 @@ export class StateStore {
|
||||
* Only dirty sessions are re-serialized; clean sessions use cached JSON fragments.
|
||||
*/
|
||||
private assembleStateJson(): string {
|
||||
// Re-serialize dirty sessions and update cache
|
||||
for (const id of this.dirtySessions) {
|
||||
const session = this.state.sessions[id];
|
||||
if (session) {
|
||||
this.cachedSessionJsons.set(id, JSON.stringify(session));
|
||||
} else {
|
||||
this.cachedSessionJsons.delete(id);
|
||||
}
|
||||
}
|
||||
this.dirtySessions.clear();
|
||||
this.updateDirtySessionCache();
|
||||
|
||||
// Build sessions object from cached fragments
|
||||
const sessionParts: string[] = [];
|
||||
@@ -218,6 +209,25 @@ export class StateStore {
|
||||
sessionParts.push(`${JSON.stringify(id)}:${json}`);
|
||||
}
|
||||
|
||||
this.pruneStaleCacheEntries();
|
||||
|
||||
return this.buildPartialJson(sessionParts);
|
||||
}
|
||||
|
||||
private updateDirtySessionCache(): void {
|
||||
// Re-serialize dirty sessions and update cache
|
||||
for (const id of this.dirtySessions) {
|
||||
const session = this.state.sessions[id];
|
||||
if (session) {
|
||||
this.cachedSessionJsons.set(id, JSON.stringify(session));
|
||||
} else {
|
||||
this.cachedSessionJsons.delete(id);
|
||||
}
|
||||
}
|
||||
this.dirtySessions.clear();
|
||||
}
|
||||
|
||||
private pruneStaleCacheEntries(): void {
|
||||
// Prune stale cache entries (sessions removed via direct state mutation)
|
||||
if (this.cachedSessionJsons.size > Object.keys(this.state.sessions).length) {
|
||||
for (const cachedId of this.cachedSessionJsons.keys()) {
|
||||
@@ -226,7 +236,9 @@ export class StateStore {
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private buildPartialJson(sessionParts: string[]): string {
|
||||
// Build final JSON: sessions from cache, everything else re-serialized (tiny)
|
||||
const sessionsJson = `{${sessionParts.join(',')}}`;
|
||||
|
||||
@@ -249,6 +261,28 @@ export class StateStore {
|
||||
return `{${parts.join(',')}}`;
|
||||
}
|
||||
|
||||
private serializeState(): string | null {
|
||||
try {
|
||||
return this.assembleStateJson();
|
||||
} catch (assembleErr) {
|
||||
// Fallback to full serialization if incremental assembly fails
|
||||
console.warn('[StateStore] assembleStateJson failed, falling back to full serialize:', assembleErr);
|
||||
this.cachedSessionJsons.clear();
|
||||
this.dirtySessions.clear();
|
||||
try {
|
||||
return JSON.stringify(this.state);
|
||||
} catch (err) {
|
||||
console.error('[StateStore] Failed to serialize state (circular reference or invalid data):', err);
|
||||
this.consecutiveSaveFailures++;
|
||||
if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) {
|
||||
console.error('[StateStore] Circuit breaker OPEN - serialization failing repeatedly');
|
||||
this.circuitBreakerOpen = true;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async _doSaveAsync(): Promise<void> {
|
||||
this.saveDeb.cancel();
|
||||
if (!this.dirty) {
|
||||
@@ -265,28 +299,10 @@ export class StateStore {
|
||||
|
||||
const tempPath = this.filePath + '.tmp';
|
||||
const backupPath = this.filePath + '.bak';
|
||||
let json: string;
|
||||
|
||||
// Step 1: Serialize state (validates it's JSON-safe)
|
||||
try {
|
||||
json = this.assembleStateJson();
|
||||
} catch (assembleErr) {
|
||||
// Fallback to full serialization if incremental assembly fails
|
||||
console.warn('[StateStore] assembleStateJson failed, falling back to full serialize:', assembleErr);
|
||||
this.cachedSessionJsons.clear();
|
||||
this.dirtySessions.clear();
|
||||
try {
|
||||
json = JSON.stringify(this.state);
|
||||
} catch (err) {
|
||||
console.error('[StateStore] Failed to serialize state (circular reference or invalid data):', err);
|
||||
this.consecutiveSaveFailures++;
|
||||
if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) {
|
||||
console.error('[StateStore] Circuit breaker OPEN - serialization failing repeatedly');
|
||||
this.circuitBreakerOpen = true;
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
const json = this.serializeState();
|
||||
if (json === null) return;
|
||||
|
||||
// Clear dirty flag BEFORE async I/O so mutations during write re-set it.
|
||||
// The state snapshot is already captured in `json` above.
|
||||
@@ -351,27 +367,9 @@ export class StateStore {
|
||||
|
||||
const tempPath = this.filePath + '.tmp';
|
||||
const backupPath = this.filePath + '.bak';
|
||||
let json: string;
|
||||
|
||||
try {
|
||||
json = this.assembleStateJson();
|
||||
} catch (assembleErr) {
|
||||
// Fallback to full serialization if incremental assembly fails
|
||||
console.warn('[StateStore] assembleStateJson failed, falling back to full serialize:', assembleErr);
|
||||
this.cachedSessionJsons.clear();
|
||||
this.dirtySessions.clear();
|
||||
try {
|
||||
json = JSON.stringify(this.state);
|
||||
} catch (err) {
|
||||
console.error('[StateStore] Failed to serialize state (circular reference or invalid data):', err);
|
||||
this.consecutiveSaveFailures++;
|
||||
if (this.consecutiveSaveFailures >= MAX_CONSECUTIVE_FAILURES) {
|
||||
console.error('[StateStore] Circuit breaker OPEN - serialization failing repeatedly');
|
||||
this.circuitBreakerOpen = true;
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
const json = this.serializeState();
|
||||
if (json === null) return;
|
||||
|
||||
// Backup via atomic copy (avoids reading entire file into memory)
|
||||
try {
|
||||
|
||||
+94
-83
@@ -222,6 +222,83 @@ export class SubagentWatcher extends EventEmitter {
|
||||
return INTERNAL_AGENT_PATTERNS.some((pattern) => pattern.test(description));
|
||||
}
|
||||
|
||||
/**
|
||||
* Mark a subagent as completed: clear PID, set status, clean up pending tool calls, emit event.
|
||||
*/
|
||||
private markSubagentAsCompleted(info: SubagentInfo): void {
|
||||
info.pid = undefined;
|
||||
info.status = 'completed';
|
||||
this.pendingToolCalls.delete(info.agentId);
|
||||
this.emit('subagent:completed', info);
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract text from message content, handling both string and array formats.
|
||||
* For array content, returns the text from the first 'text' block.
|
||||
*/
|
||||
private extractFirstTextContent(
|
||||
content: string | Array<{ type: string; text?: string }> | undefined
|
||||
): string | undefined {
|
||||
if (!content) return undefined;
|
||||
if (typeof content === 'string') {
|
||||
const trimmed = content.trim();
|
||||
return trimmed.length > 0 ? trimmed : undefined;
|
||||
}
|
||||
if (Array.isArray(content)) {
|
||||
const firstContent = content[0];
|
||||
if (firstContent?.type === 'text' && firstContent.text) {
|
||||
const trimmed = firstContent.text.trim();
|
||||
return trimmed.length > 0 ? trimmed : undefined;
|
||||
}
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Process a tool_result content block: look up pending tool call, emit tool_result event.
|
||||
*/
|
||||
private emitToolResult(
|
||||
content: { tool_use_id: string; content?: string | Array<{ type: string; text?: string }>; is_error?: boolean },
|
||||
agentId: string,
|
||||
sessionId: string,
|
||||
timestamp: string
|
||||
): void {
|
||||
const resultContent = this.extractToolResultContent(content.content);
|
||||
const agentPendingCalls = this.pendingToolCalls.get(agentId);
|
||||
const pendingCall = agentPendingCalls?.get(content.tool_use_id);
|
||||
const toolName = pendingCall?.toolName;
|
||||
// Delete after lookup to prevent memory leak
|
||||
agentPendingCalls?.delete(content.tool_use_id);
|
||||
|
||||
const toolResult: SubagentToolResult = {
|
||||
agentId,
|
||||
sessionId,
|
||||
timestamp,
|
||||
toolUseId: content.tool_use_id,
|
||||
tool: toolName,
|
||||
preview: resultContent.substring(0, MESSAGE_TEXT_LIMIT),
|
||||
contentLength: resultContent.length,
|
||||
isError: content.is_error || false,
|
||||
};
|
||||
this.emit('subagent:tool_result', toolResult);
|
||||
}
|
||||
|
||||
/**
|
||||
* Find the oldest inactive (non-active) agent for LRU eviction.
|
||||
* Returns the agent ID of the oldest inactive agent, or null if all are active.
|
||||
*/
|
||||
private findOldestInactiveAgent(): string | null {
|
||||
let oldestId: string | null = null;
|
||||
let oldestTime = Infinity;
|
||||
for (const [id, existing] of this.agentInfo) {
|
||||
if (existing.status !== 'active' && existing.lastActivityAt < oldestTime) {
|
||||
oldestTime = existing.lastActivityAt;
|
||||
oldestId = id;
|
||||
}
|
||||
}
|
||||
return oldestId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract short model identifier from full model name
|
||||
*/
|
||||
@@ -307,10 +384,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
|
||||
const alive = this.checkSubagentAliveFromPidMap(info, pidMap);
|
||||
if (!alive) {
|
||||
info.pid = undefined;
|
||||
info.status = 'completed';
|
||||
this.pendingToolCalls.delete(info.agentId);
|
||||
this.emit('subagent:completed', info);
|
||||
this.markSubagentAsCompleted(info);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -677,10 +751,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
const pid = await this.findSubagentProcess(info.sessionId);
|
||||
if (pid) {
|
||||
process.kill(pid, 'SIGTERM');
|
||||
info.pid = undefined;
|
||||
info.status = 'completed';
|
||||
this.pendingToolCalls.delete(info.agentId);
|
||||
this.emit('subagent:completed', info);
|
||||
this.markSubagentAsCompleted(info);
|
||||
return true;
|
||||
}
|
||||
} catch {
|
||||
@@ -688,10 +759,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
}
|
||||
|
||||
// Mark as completed even if we couldn't find the process
|
||||
info.pid = undefined;
|
||||
info.status = 'completed';
|
||||
this.pendingToolCalls.delete(info.agentId);
|
||||
this.emit('subagent:completed', info);
|
||||
this.markSubagentAsCompleted(info);
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -843,19 +911,9 @@ export class SubagentWatcher extends EventEmitter {
|
||||
}
|
||||
} else if (entry.type === 'user' && entry.message?.content) {
|
||||
// Handle both string and array content formats
|
||||
if (typeof entry.message.content === 'string') {
|
||||
const text = entry.message.content.trim();
|
||||
if (text.length < 100 && !text.includes('{')) {
|
||||
lines.push(`${this.formatTime(entry.timestamp)} 📥 User: ${text.substring(0, USER_TEXT_PREVIEW_LENGTH)}`);
|
||||
}
|
||||
} else {
|
||||
const firstContent = entry.message.content[0];
|
||||
if (firstContent?.type === 'text' && firstContent.text) {
|
||||
const text = firstContent.text.trim();
|
||||
if (text.length < 100 && !text.includes('{')) {
|
||||
lines.push(`${this.formatTime(entry.timestamp)} 📥 User: ${text.substring(0, USER_TEXT_PREVIEW_LENGTH)}`);
|
||||
}
|
||||
}
|
||||
const text = this.extractFirstTextContent(entry.message.content);
|
||||
if (text && text.length < 100 && !text.includes('{')) {
|
||||
lines.push(`${this.formatTime(entry.timestamp)} 📥 User: ${text.substring(0, USER_TEXT_PREVIEW_LENGTH)}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1027,15 +1085,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
try {
|
||||
const entry = JSON.parse(line);
|
||||
if (entry.type === 'user' && entry.message?.content) {
|
||||
let text: string | undefined;
|
||||
if (typeof entry.message.content === 'string') {
|
||||
text = entry.message.content.trim();
|
||||
} else if (Array.isArray(entry.message.content)) {
|
||||
const firstContent = entry.message.content[0];
|
||||
if (firstContent?.type === 'text' && firstContent.text) {
|
||||
text = firstContent.text.trim();
|
||||
}
|
||||
}
|
||||
const text = this.extractFirstTextContent(entry.message.content);
|
||||
if (text) {
|
||||
resolved = true;
|
||||
rl.close();
|
||||
@@ -1286,14 +1336,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
|
||||
// Enforce MAX_TRACKED_AGENTS during insertion — evict oldest inactive agent
|
||||
if (this.agentInfo.size >= MAX_TRACKED_AGENTS) {
|
||||
let oldestId: string | null = null;
|
||||
let oldestTime = Infinity;
|
||||
for (const [id, existing] of this.agentInfo) {
|
||||
if (existing.status !== 'active' && existing.lastActivityAt < oldestTime) {
|
||||
oldestTime = existing.lastActivityAt;
|
||||
oldestId = id;
|
||||
}
|
||||
}
|
||||
const oldestId = this.findOldestInactiveAgent();
|
||||
if (oldestId) {
|
||||
this.removeAgent(oldestId);
|
||||
}
|
||||
@@ -1386,15 +1429,7 @@ export class SubagentWatcher extends EventEmitter {
|
||||
let description = await this.extractDescriptionFromParentTranscript(info.projectHash, info.sessionId, agentId);
|
||||
// Fallback: extract smart title from the prompt content
|
||||
if (!description) {
|
||||
let text: string | undefined;
|
||||
if (typeof entry.message.content === 'string') {
|
||||
text = entry.message.content.trim();
|
||||
} else if (Array.isArray(entry.message.content)) {
|
||||
const firstContent = entry.message.content[0];
|
||||
if (firstContent?.type === 'text' && firstContent.text) {
|
||||
text = firstContent.text.trim();
|
||||
}
|
||||
}
|
||||
const text = this.extractFirstTextContent(entry.message.content);
|
||||
if (text) {
|
||||
description = this.extractSmartTitle(text);
|
||||
}
|
||||
@@ -1480,24 +1515,12 @@ export class SubagentWatcher extends EventEmitter {
|
||||
}
|
||||
} else if (content.type === 'tool_result' && content.tool_use_id) {
|
||||
// Extract tool result
|
||||
const resultContent = this.extractToolResultContent(content.content);
|
||||
const agentPendingCalls = this.pendingToolCalls.get(agentId);
|
||||
const pendingCall = agentPendingCalls?.get(content.tool_use_id);
|
||||
const toolName = pendingCall?.toolName;
|
||||
// Delete after lookup to prevent memory leak
|
||||
agentPendingCalls?.delete(content.tool_use_id);
|
||||
|
||||
const toolResult: SubagentToolResult = {
|
||||
this.emitToolResult(
|
||||
{ tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error },
|
||||
agentId,
|
||||
sessionId,
|
||||
timestamp: entry.timestamp,
|
||||
toolUseId: content.tool_use_id,
|
||||
tool: toolName,
|
||||
preview: resultContent.substring(0, MESSAGE_TEXT_LIMIT),
|
||||
contentLength: resultContent.length,
|
||||
isError: content.is_error || false,
|
||||
};
|
||||
this.emit('subagent:tool_result', toolResult);
|
||||
entry.timestamp
|
||||
);
|
||||
} else if (content.type === 'text' && content.text) {
|
||||
const text = content.text.trim();
|
||||
if (text.length > 0) {
|
||||
@@ -1531,24 +1554,12 @@ export class SubagentWatcher extends EventEmitter {
|
||||
// Check for tool_result blocks in user messages (common pattern)
|
||||
for (const content of entry.message.content) {
|
||||
if (content.type === 'tool_result' && content.tool_use_id) {
|
||||
const resultContent = this.extractToolResultContent(content.content);
|
||||
const agentPendingCalls = this.pendingToolCalls.get(agentId);
|
||||
const pendingCall = agentPendingCalls?.get(content.tool_use_id);
|
||||
const toolName = pendingCall?.toolName;
|
||||
// Delete after lookup to prevent memory leak
|
||||
agentPendingCalls?.delete(content.tool_use_id);
|
||||
|
||||
const toolResult: SubagentToolResult = {
|
||||
this.emitToolResult(
|
||||
{ tool_use_id: content.tool_use_id, content: content.content, is_error: content.is_error },
|
||||
agentId,
|
||||
sessionId,
|
||||
timestamp: entry.timestamp,
|
||||
toolUseId: content.tool_use_id,
|
||||
tool: toolName,
|
||||
preview: resultContent.substring(0, MESSAGE_TEXT_LIMIT),
|
||||
contentLength: resultContent.length,
|
||||
isError: content.is_error || false,
|
||||
};
|
||||
this.emit('subagent:tool_result', toolResult);
|
||||
entry.timestamp
|
||||
);
|
||||
} else if (content.type === 'text' && content.text) {
|
||||
const userText = content.text.trim();
|
||||
if (userText.length > 0 && userText.length < 500) {
|
||||
|
||||
+4
-17
@@ -98,6 +98,9 @@ const LEGACY_MUX_NAME_PATTERN = /^claudeman-[a-f0-9-]+$/;
|
||||
/** Regex to validate tmux pane targets (e.g., "%0", "%1", "0", "1") */
|
||||
const SAFE_PANE_TARGET_PATTERN = /^(%\d+|\d+)$/;
|
||||
|
||||
/** Characters unsafe in paths — shell metacharacters, quotes, and control chars */
|
||||
const UNSAFE_PATH_CHARS = /[;&|$`(){}<>'"\n\r]/;
|
||||
|
||||
/**
|
||||
* Validates that a session name contains only safe characters.
|
||||
* Prevents command injection via malformed session IDs.
|
||||
@@ -111,23 +114,7 @@ function isValidMuxName(name: string): boolean {
|
||||
* Prevents command injection via malformed paths.
|
||||
*/
|
||||
function isValidPath(path: string): boolean {
|
||||
if (
|
||||
path.includes(';') ||
|
||||
path.includes('&') ||
|
||||
path.includes('|') ||
|
||||
path.includes('$') ||
|
||||
path.includes('`') ||
|
||||
path.includes('(') ||
|
||||
path.includes(')') ||
|
||||
path.includes('{') ||
|
||||
path.includes('}') ||
|
||||
path.includes('<') ||
|
||||
path.includes('>') ||
|
||||
path.includes("'") ||
|
||||
path.includes('"') ||
|
||||
path.includes('\n') ||
|
||||
path.includes('\r')
|
||||
) {
|
||||
if (UNSAFE_PATH_CHARS.test(path)) {
|
||||
return false;
|
||||
}
|
||||
if (path.includes('..')) {
|
||||
|
||||
Reference in New Issue
Block a user