Skip to content

Commit d262b1c

Browse files
committed
fix(logging): split queue diagnostic runtime
1 parent 39f22ef commit d262b1c

4 files changed

Lines changed: 56 additions & 36 deletions

File tree

src/logging/diagnostic-runtime.ts

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,40 @@
1+
import { emitDiagnosticEvent } from "../infra/diagnostic-events.js";
2+
import { createSubsystemLogger } from "./subsystem.js";
3+
4+
const diag = createSubsystemLogger("diagnostic");
5+
let lastActivityAt = 0;
6+
7+
export const diagnosticLogger = diag;
8+
9+
export function markDiagnosticActivity(): void {
10+
lastActivityAt = Date.now();
11+
}
12+
13+
export function getLastDiagnosticActivityAt(): number {
14+
return lastActivityAt;
15+
}
16+
17+
export function resetDiagnosticActivityForTest(): void {
18+
lastActivityAt = 0;
19+
}
20+
21+
export function logLaneEnqueue(lane: string, queueSize: number): void {
22+
diag.debug(`lane enqueue: lane=${lane} queueSize=${queueSize}`);
23+
emitDiagnosticEvent({
24+
type: "queue.lane.enqueue",
25+
lane,
26+
queueSize,
27+
});
28+
markDiagnosticActivity();
29+
}
30+
31+
export function logLaneDequeue(lane: string, waitMs: number, queueSize: number): void {
32+
diag.debug(`lane dequeue: lane=${lane} waitMs=${waitMs} queueSize=${queueSize}`);
33+
emitDiagnosticEvent({
34+
type: "queue.lane.dequeue",
35+
lane,
36+
queueSize,
37+
waitMs,
38+
});
39+
markDiagnosticActivity();
40+
}

src/logging/diagnostic.ts

Lines changed: 10 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,12 @@
11
import { getRuntimeConfig } from "../config/config.js";
22
import type { OpenClawConfig } from "../config/types.openclaw.js";
33
import { emitDiagnosticEvent } from "../infra/diagnostic-events.js";
4+
import {
5+
diagnosticLogger as diag,
6+
getLastDiagnosticActivityAt,
7+
markDiagnosticActivity as markActivity,
8+
resetDiagnosticActivityForTest,
9+
} from "./diagnostic-runtime.js";
410
import {
511
diagnosticSessionStates,
612
getDiagnosticSessionState,
@@ -10,9 +16,7 @@ import {
1016
type SessionRef,
1117
type SessionStateValue,
1218
} from "./diagnostic-session-state.js";
13-
import { createSubsystemLogger } from "./subsystem.js";
14-
15-
const diag = createSubsystemLogger("diagnostic");
19+
export { diagnosticLogger, logLaneDequeue, logLaneEnqueue } from "./diagnostic-runtime.js";
1620

1721
const webhookStats = {
1822
received: 0,
@@ -21,7 +25,6 @@ const webhookStats = {
2125
lastReceived: 0,
2226
};
2327

24-
let lastActivityAt = 0;
2528
const DEFAULT_STUCK_SESSION_WARN_MS = 120_000;
2629
const MIN_STUCK_SESSION_WARN_MS = 1_000;
2730
const MAX_STUCK_SESSION_WARN_MS = 24 * 60 * 60 * 1000;
@@ -34,10 +37,6 @@ function loadCommandPollBackoffRuntime() {
3437
return commandPollBackoffRuntimePromise;
3538
}
3639

37-
function markActivity() {
38-
lastActivityAt = Date.now();
39-
}
40-
4140
export function resolveStuckSessionWarnMs(config?: OpenClawConfig): number {
4241
const raw = config?.diagnostics?.stuckSessionWarnMs;
4342
if (typeof raw !== "number" || !Number.isFinite(raw)) {
@@ -244,27 +243,6 @@ export function logSessionStuck(params: SessionRef & { state: SessionStateValue;
244243
markActivity();
245244
}
246245

247-
export function logLaneEnqueue(lane: string, queueSize: number) {
248-
diag.debug(`lane enqueue: lane=${lane} queueSize=${queueSize}`);
249-
emitDiagnosticEvent({
250-
type: "queue.lane.enqueue",
251-
lane,
252-
queueSize,
253-
});
254-
markActivity();
255-
}
256-
257-
export function logLaneDequeue(lane: string, waitMs: number, queueSize: number) {
258-
diag.debug(`lane dequeue: lane=${lane} waitMs=${waitMs} queueSize=${queueSize}`);
259-
emitDiagnosticEvent({
260-
type: "queue.lane.dequeue",
261-
lane,
262-
queueSize,
263-
waitMs,
264-
});
265-
markActivity();
266-
}
267-
268246
export function logRunAttempt(params: SessionRef & { runId: string; attempt: number }) {
269247
diag.debug(
270248
`run attempt: sessionId=${params.sessionId ?? "unknown"} sessionKey=${
@@ -360,15 +338,15 @@ export function startDiagnosticHeartbeat(
360338
0,
361339
);
362340
const hasActivity =
363-
lastActivityAt > 0 ||
341+
getLastDiagnosticActivityAt() > 0 ||
364342
webhookStats.received > 0 ||
365343
activeCount > 0 ||
366344
waitingCount > 0 ||
367345
totalQueued > 0;
368346
if (!hasActivity) {
369347
return;
370348
}
371-
if (now - lastActivityAt > 120_000 && activeCount === 0 && waitingCount === 0) {
349+
if (now - getLastDiagnosticActivityAt() > 120_000 && activeCount === 0 && waitingCount === 0) {
372350
return;
373351
}
374352

@@ -425,12 +403,10 @@ export function getDiagnosticSessionStateCountForTest(): number {
425403

426404
export function resetDiagnosticStateForTest(): void {
427405
resetDiagnosticSessionStateForTest();
406+
resetDiagnosticActivityForTest();
428407
webhookStats.received = 0;
429408
webhookStats.processed = 0;
430409
webhookStats.errors = 0;
431410
webhookStats.lastReceived = 0;
432-
lastActivityAt = 0;
433411
stopDiagnosticHeartbeat();
434412
}
435-
436-
export { diag as diagnosticLogger };

src/process/command-queue.test.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ const diagnosticMocks = vi.hoisted(() => ({
1212
},
1313
}));
1414

15-
vi.mock("../logging/diagnostic.js", () => ({
15+
vi.mock("../logging/diagnostic-runtime.js", () => ({
1616
logLaneEnqueue: diagnosticMocks.logLaneEnqueue,
1717
logLaneDequeue: diagnosticMocks.logLaneDequeue,
1818
diagnosticLogger: diagnosticMocks.diag,

src/process/command-queue.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,8 @@
1-
import { diagnosticLogger as diag, logLaneDequeue, logLaneEnqueue } from "../logging/diagnostic.js";
1+
import {
2+
diagnosticLogger as diag,
3+
logLaneDequeue,
4+
logLaneEnqueue,
5+
} from "../logging/diagnostic-runtime.js";
26
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
37
import { CommandLane } from "./lanes.js";
48
/**

0 commit comments

Comments
 (0)