Files
fusion/packages/core/src/reflection-store.ts
gsxdsm cca13737b6 FN-8603: reduce steady-state diagnostic log noise
Route routine core, engine, and dashboard diagnostics through debug-gated shared loggers.

- Demote steady-state diagnostic sites while preserving warnings and errors for actionable failures.
- Add cross-package severity contracts and manifest coverage for demoted log sites.
- Document logging severity guidance and add a patch changeset.

Files changed:
 .changeset/fn-8603-log-severity.md                 |  7 ++
 docs/diagnostics.md                                | 20 ++++--
 .../__tests__/log-severity-spam-contract.test.ts   | 71 ++++++++++++++++++
 packages/core/src/activity-analytics.ts            |  5 +-
 packages/core/src/ai-summarize.ts                  | 61 +++++++---------
 packages/core/src/async-mission-store.ts           |  5 +-
 packages/core/src/async-secrets-store.ts           |  7 +-
 packages/core/src/central-core.ts                  | 17 ++---
 packages/core/src/docker-provisioning.ts           | 13 ++--
 packages/core/src/index.ts                         |  1 +
 packages/core/src/master-key.ts                    |  9 ++-
 packages/core/src/memory-compaction.ts             | 29 ++++----
 packages/core/src/memory-insights.ts               |  7 +-
 packages/core/src/migration-orchestrator.ts        |  7 +-
 packages/core/src/mission-store.ts                 |  5 +-
 packages/core/src/node-discovery.ts                |  7 +-
 packages/core/src/notification/dispatcher.ts       |  9 ++-
 .../core/src/plugins/bundled-plugin-install.ts     | 11 +--
 packages/core/src/reflection-store.ts              |  5 +-
 packages/core/src/secrets-store.ts                 |  7 +-
 packages/core/src/task-store/agent-logs.ts         | 21 +++---
 packages/core/src/task-store/async-events.ts       |  5 +-
 packages/core/src/task-store/async-maintenance.ts  |  7 +-
 packages/core/src/task-store/comments-ops.ts       |  7 +-
 packages/core/src/task-store/task-mutation-ops.ts  | 11 +--
 packages/core/src/task-store/workflow-integrity.ts |  9 ++-
 packages/core/src/types/merge-policy.ts            |  5 +-
 packages/core/src/usage-events.ts                  |  5 +-
 .../__tests__/log-severity-spam-contract.test.ts   | 48 +++++++++++++
 packages/dashboard/src/ai-refine.ts                |  5 +-
 packages/dashboard/src/ai-session-diagnostics.ts   | 10 +--
 packages/dashboard/src/chat.ts                     |  8 ++-
 packages/dashboard/src/devserver-manager.ts        |  9 ++-
 packages/dashboard/src/file-service.ts             |  5 +-
 packages/dashboard/src/github-tracking-comments.ts |  7 +-
 .../dashboard/src/github-tracking-reconciler.ts    |  5 +-
 packages/dashboard/src/github-tracking-state.ts    |  5 +-
 packages/dashboard/src/gitlab-lifecycle.ts         |  5 +-
 packages/dashboard/src/insights-routes.ts          |  9 ++-
 packages/dashboard/src/issue-image-attachments.ts  |  5 +-
 packages/dashboard/src/knowledge-index.ts          |  5 +-
 packages/dashboard/src/plugin-routes.ts            |  7 +-
 packages/dashboard/src/routes/board-workflows.ts   |  5 +-
 packages/dashboard/src/routes/context.ts           |  5 +-
 .../dashboard/src/routes/register-auth-routes.ts   | 13 ++--
 .../routes/register-docker-provisioning-routes.ts  |  7 +-
 .../dashboard/src/routes/register-git-github.ts    | 21 +++---
 packages/dashboard/src/routes/register-gitlab.ts   |  7 +-
 .../src/routes/register-session-diff-routes.ts     |  9 ++-
 .../src/routes/register-settings-memory-routes.ts  |  7 +-
 .../src/routes/register-setup-activity-routes.ts   |  7 +-
 .../dashboard/src/routes/register-signal-routes.ts |  5 +-
 .../src/routes/register-task-workflow-routes.ts    | 11 +--
 packages/dashboard/src/runtime-logger.ts           | 11 +--
 packages/dashboard/src/server.ts                   |  7 +-
 packages/dashboard/src/sse.ts                      |  8 ++-
 packages/dashboard/src/terminal-service.ts         | 34 ++++-----
 packages/dashboard/src/view-chunk-manifest.ts      |  5 +-
 .../engine/src/__tests__/log-severity-manifest.ts  | 83 ++++++++++++++++++++++
 .../__tests__/log-severity-spam-contract.test.ts   | 40 ++++++++++-
 .../src/__tests__/logger-debug-gating.test.ts      |  7 +-
 packages/engine/src/goal-anchoring-audit.ts        |  5 +-
 packages/engine/src/plugin-runner.ts               | 44 ++++++------
 packages/engine/src/pty-native.ts                  |  9 ++-
 .../engine/src/runtimes/child-process-worker.ts    |  4 +-
 packages/engine/src/self-healing.ts                | 12 ++--
 packages/engine/src/worktree-hooks.ts              | 10 ++-
 67 files changed, 632 insertions(+), 250 deletions(-)

Fusion-Task-Id: FN-8603

Fusion-Task-Lineage: 53901db6-1af2-4bd7-b5ea-49507e048ef2

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
2026-07-26 12:01:19 -07:00

305 lines
9.4 KiB
TypeScript

import { createLogger } from "./logger.js";
const severityAuditLog = createLogger("core-reflection-store");
import { randomUUID } from "node:crypto";
import { EventEmitter } from "node:events";
import { existsSync } from "node:fs";
import { mkdir, readFile, unlink, writeFile } from "node:fs/promises";
import { join, resolve } from "node:path";
import type {
AgentPerformanceSummary,
AgentReflection,
ReflectionMetrics,
ReflectionTrigger,
} from "./types.js";
/** Events emitted by ReflectionStore. */
export interface ReflectionStoreEvents {
/** Emitted after a reflection is created and persisted. */
"reflection:created": (reflection: AgentReflection) => void;
/** Emitted when a performance summary is computed. */
"reflection:summary-computed": (summary: AgentPerformanceSummary) => void;
}
/** Constructor options for ReflectionStore. */
export interface ReflectionStoreOptions {
/** Root fn data directory (default: .fusion). */
rootDir?: string;
}
/** Input payload for creating a reflection. */
export interface CreateReflectionInput {
agentId: string;
trigger: ReflectionTrigger;
triggerDetail?: string;
taskId?: string;
metrics: ReflectionMetrics;
insights: string[];
suggestedImprovements: string[];
summary: string;
}
/** Options for computing a performance summary. */
export interface PerformanceSummaryOptions {
/** Time window in milliseconds to include reflections from. Defaults to 7 days. */
windowMs?: number;
}
interface AgentLock {
promise: Promise<unknown>;
}
const DEFAULT_REFLECTION_LIMIT = 50;
const DEFAULT_WINDOW_MS = 7 * 24 * 60 * 60 * 1000;
const SUMMARY_LIST_LIMIT = 10;
/**
* ReflectionStore persists agent self-reflection records in append-only JSONL files.
*
* Storage layout:
* - `.fusion/agents/{agentId}-reflections.jsonl`
*/
export class ReflectionStore extends EventEmitter {
private rootDir: string;
private agentsDir: string;
private locks: Map<string, AgentLock> = new Map();
constructor(options: ReflectionStoreOptions = {}) {
super();
if (!options.rootDir && process.env.VITEST === "true") {
throw new Error(
"ReflectionStore requires an explicit rootDir during test execution. Pass an absolute path to avoid writing to unintended locations.",
);
}
this.rootDir = options.rootDir ?? resolve(".fusion");
this.agentsDir = join(this.rootDir, "agents");
}
override on(event: "reflection:created", listener: ReflectionStoreEvents["reflection:created"]): this;
override on(
event: "reflection:summary-computed",
listener: ReflectionStoreEvents["reflection:summary-computed"],
): this;
override on(event: string | symbol, listener: (...args: any[]) => void): this {
return super.on(event, listener);
}
override emit(event: "reflection:created", reflection: AgentReflection): boolean;
override emit(event: "reflection:summary-computed", summary: AgentPerformanceSummary): boolean;
override emit(event: string | symbol, ...args: any[]): boolean {
return super.emit(event, ...args);
}
/** Ensure required directories exist. */
async init(): Promise<void> {
await mkdir(this.agentsDir, { recursive: true });
}
/** Create and append a reflection for an agent. */
async createReflection(input: CreateReflectionInput): Promise<AgentReflection> {
if (!input.agentId?.trim()) {
throw new Error("agentId is required");
}
return this.withLock(input.agentId, async () => {
const reflection: AgentReflection = {
id: `reflection-${randomUUID().slice(0, 8)}`,
agentId: input.agentId,
timestamp: new Date().toISOString(),
trigger: input.trigger,
triggerDetail: input.triggerDetail,
taskId: input.taskId,
metrics: input.metrics,
insights: input.insights,
suggestedImprovements: input.suggestedImprovements,
summary: input.summary,
};
const line = `${JSON.stringify(reflection)}\n`;
await writeFile(this.reflectionsPath(input.agentId), line, { flag: "a" });
this.emit("reflection:created", reflection);
return reflection;
});
}
/** Get recent reflections for an agent (newest first). */
async getReflections(agentId: string, limit = DEFAULT_REFLECTION_LIMIT): Promise<AgentReflection[]> {
if (!agentId?.trim()) {
return [];
}
const reflectionPath = this.reflectionsPath(agentId);
if (!existsSync(reflectionPath)) {
return [];
}
const reflections = await this.readReflectionsFromFile(agentId);
return reflections.slice(0, Math.max(0, limit));
}
/** Get the most recent reflection for an agent. */
async getLatestReflection(agentId: string): Promise<AgentReflection | null> {
const reflections = await this.getReflections(agentId, 1);
return reflections[0] ?? null;
}
/** Compute an aggregate performance summary from recent reflections. */
async getPerformanceSummary(
agentId: string,
options: PerformanceSummaryOptions = {},
): Promise<AgentPerformanceSummary> {
const windowMs = options.windowMs ?? DEFAULT_WINDOW_MS;
const cutoff = Date.now() - windowMs;
const allReflections = await this.getReflections(agentId, Number.MAX_SAFE_INTEGER);
const windowedReflections = allReflections.filter((reflection) => {
const timestamp = Date.parse(reflection.timestamp);
return Number.isFinite(timestamp) && timestamp >= cutoff;
});
let totalTasksCompleted = 0;
let totalTasksFailed = 0;
const durations: number[] = [];
const errorCounts = new Map<string, number>();
for (const reflection of windowedReflections) {
totalTasksCompleted += reflection.metrics.tasksCompleted ?? 0;
totalTasksFailed += reflection.metrics.tasksFailed ?? 0;
if (typeof reflection.metrics.avgDurationMs === "number") {
durations.push(reflection.metrics.avgDurationMs);
}
for (const error of reflection.metrics.commonErrors ?? []) {
const normalized = error.trim();
if (!normalized) continue;
errorCounts.set(normalized, (errorCounts.get(normalized) ?? 0) + 1);
}
}
const avgDurationMs = durations.length > 0
? durations.reduce((sum, value) => sum + value, 0) / durations.length
: 0;
const totalTasks = totalTasksCompleted + totalTasksFailed;
const successRate = totalTasks > 0 ? totalTasksCompleted / totalTasks : 0;
const commonErrors = Array.from(errorCounts.entries())
.sort((a, b) => {
if (b[1] !== a[1]) return b[1] - a[1];
return a[0].localeCompare(b[0]);
})
.slice(0, SUMMARY_LIST_LIMIT)
.map(([error]) => error);
const strengths = this.collectRecentUnique(
windowedReflections.flatMap((reflection) => reflection.insights),
SUMMARY_LIST_LIMIT,
);
const weaknesses = this.collectRecentUnique(
windowedReflections.flatMap((reflection) => reflection.suggestedImprovements),
SUMMARY_LIST_LIMIT,
);
const summary: AgentPerformanceSummary = {
agentId,
totalTasksCompleted,
totalTasksFailed,
avgDurationMs,
successRate,
commonErrors,
strengths,
weaknesses,
recentReflectionCount: windowedReflections.length,
computedAt: new Date().toISOString(),
};
this.emit("reflection:summary-computed", summary);
return summary;
}
/** Delete all persisted reflections for an agent. */
async deleteReflections(agentId: string): Promise<void> {
if (!agentId?.trim()) {
return;
}
await this.withLock(agentId, async () => {
try {
await unlink(this.reflectionsPath(agentId));
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "ENOENT") {
throw error;
}
}
});
}
private reflectionsPath(agentId: string): string {
return join(this.agentsDir, `${agentId}-reflections.jsonl`);
}
private async readReflectionsFromFile(agentId: string): Promise<AgentReflection[]> {
const reflectionPath = this.reflectionsPath(agentId);
try {
const content = await readFile(reflectionPath, "utf-8");
const lines = content.split("\n").filter(Boolean);
const reflections: AgentReflection[] = [];
for (const [index, line] of lines.entries()) {
try {
reflections.push(JSON.parse(line) as AgentReflection);
} catch (error) {
severityAuditLog.warn(
`[ReflectionStore] Skipping malformed reflection line ${index + 1} for ${agentId}`,
error,
);
}
}
return reflections.reverse();
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") {
return [];
}
throw error;
}
}
private collectRecentUnique(items: string[], maxItems: number): string[] {
const seen = new Set<string>();
const deduped: string[] = [];
for (const item of items) {
const normalized = item.trim();
if (!normalized || seen.has(normalized)) continue;
seen.add(normalized);
deduped.push(normalized);
if (deduped.length >= maxItems) {
break;
}
}
return deduped;
}
private async withLock<T>(agentId: string, fn: () => Promise<T>): Promise<T> {
let lock = this.locks.get(agentId);
if (!lock) {
lock = { promise: Promise.resolve() };
this.locks.set(agentId, lock);
}
const operation = lock.promise.then(fn, fn);
lock.promise = operation;
return operation as Promise<T>;
}
}