Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
ALTER TABLE workspace_resource_summaries
ADD COLUMN memory_working_set_mean_bytes INTEGER;

ALTER TABLE workspace_resource_summaries
ADD COLUMN memory_working_set_peak_bytes INTEGER;

ALTER TABLE workspace_resource_summaries
ADD COLUMN memory_working_set_sample_count INTEGER NOT NULL DEFAULT 0;
3 changes: 3 additions & 0 deletions apps/api/src/db/schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1582,6 +1582,9 @@ export const workspaceResourceSummaries = sqliteTable(
memoryMeanBytes: integer('memory_mean_bytes'),
memoryPeakBytes: integer('memory_peak_bytes'),
memoryKernelPeakBytes: integer('memory_kernel_peak_bytes'),
memoryWorkingSetMeanBytes: integer('memory_working_set_mean_bytes'),
memoryWorkingSetPeakBytes: integer('memory_working_set_peak_bytes'),
memoryWorkingSetSampleCount: integer('memory_working_set_sample_count').notNull().default(0),
ioReadBytes: integer('io_read_bytes'),
ioWriteBytes: integer('io_write_bytes'),
oomCount: integer('oom_count').notNull().default(0),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ export const getResourceHistoryDef: AnthropicToolDef = {
description:
'Inspect bounded workspace resource history for a project session, task, or workspace. ' +
'Returns a cheap summary with server-resolved agentProfileId, skillId, and agentType plus a chunk index by default. Pass chunkId to load bounded downsampled samples and tool-span correlation for one chunk. ' +
'Working-set memory is the sizing figure; total memory includes reclaimable file cache. ' +
'Tool spans include the ACP kind and metadata-provided tool name when available. They are correlation windows, not causal per-process attribution, and stored payloads omit titles, prompts, commands, tool args/output, file paths, env, and secrets.',
input_schema: {
type: 'object',
Expand Down Expand Up @@ -84,6 +85,8 @@ export async function getResourceHistory(
...history,
notes: [
'Samples are workspace-level cgroup observations, not per-process attribution.',
'memoryWorkingSetMeanBytes and memoryWorkingSetPeakBytes estimate memory needed by excluding reclaimable inactive file cache; null means the VM agent did not report them.',
'memoryMeanBytes, memoryPeakBytes, and memoryKernelPeakBytes include cache and remain available for historical comparison.',
'Tool spans are timestamp correlation windows and may include an ACP kind and metadata-provided tool name; titles and inputs are never returned.',
'Chunk detail is returned only when chunkId is supplied; summary reads stay bounded.',
],
Expand Down
2 changes: 2 additions & 0 deletions apps/api/src/routes/mcp/session-tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,8 @@ export async function handleGetResourceHistory(
...history,
notes: [
'Samples are workspace-level cgroup observations, not per-process attribution.',
'memoryWorkingSetMeanBytes and memoryWorkingSetPeakBytes estimate memory needed by excluding reclaimable inactive file cache; null means the VM agent did not report them.',
'memoryMeanBytes, memoryPeakBytes, and memoryKernelPeakBytes include cache and remain available for historical comparison.',
'Tool spans are timestamp correlation windows and may include an ACP kind and metadata-provided tool name; titles and inputs are never returned.',
'Chunk detail is returned only when chunkId is supplied; summary reads stay bounded.',
],
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -171,7 +171,7 @@ export const PROJECT_AWARENESS_TOOLS = [
{
name: 'get_resource_history',
description:
'Inspect bounded workspace resource history for the current project. By default, MCP callers read their current session/task/workspace summary, including server-resolved agentProfileId, skillId, and agentType, plus the chunk index. Pass sessionId, taskId, or workspaceId to inspect a related scope. Pass chunkId to lazily load downsampled raw samples and tool-span correlation for that chunk, including ACP kind and metadata-provided tool name when available. This reports correlation, not causal per-process attribution, and never includes titles, prompts, commands, tool args/output, file paths, env, or secrets.',
'Inspect bounded workspace resource history for the current project. By default, MCP callers read their current session/task/workspace summary, including server-resolved agentProfileId, skillId, and agentType, plus the chunk index. Pass sessionId, taskId, or workspaceId to inspect a related scope. Pass chunkId to lazily load downsampled raw samples and tool-span correlation for that chunk, including ACP kind and metadata-provided tool name when available. Working-set memory is the sizing figure; total memory includes reclaimable file cache. This reports correlation, not causal per-process attribution, and never includes titles, prompts, commands, tool args/output, file paths, env, or secrets.',
inputSchema: {
type: 'object' as const,
properties: {
Expand Down
52 changes: 52 additions & 0 deletions apps/api/src/services/workspace-resource-history.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,9 @@
memoryMeanBytes?: number | null;
memoryPeakBytes?: number | null;
memoryKernelPeakBytes?: number | null;
memoryWorkingSetMeanBytes?: number | null;
memoryWorkingSetPeakBytes?: number | null;
memoryWorkingSetSampleCount?: number | null;
ioReadBytes?: number | null;
ioWriteBytes?: number | null;
oomCount?: number | null;
Expand All @@ -88,6 +91,7 @@
cpuMillis?: number;
memoryBytes?: number;
memoryPeakBytes?: number;
memoryWorkingSetBytes?: number;
ioReadBytes?: number;
ioWriteBytes?: number;
oom?: number;
Expand Down Expand Up @@ -144,6 +148,8 @@
memoryMeanBytes: number | null;
memoryPeakBytes: number | null;
memoryKernelPeakBytes: number | null;
memoryWorkingSetMeanBytes: number | null;
memoryWorkingSetPeakBytes: number | null;
ioReadBytes: number | null;
ioWriteBytes: number | null;
oomCount: number;
Expand Down Expand Up @@ -668,7 +674,7 @@
};
}

export async function storeWorkspaceResourceChunk(

Check failure on line 677 in apps/api/src/services/workspace-resource-history.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Refactor this function to reduce its Cognitive Complexity from 16 to the 15 allowed.

See more on https://sonarcloud.io/project/issues?id=raphaeltm_simple-agent-manager&issues=AaDsvPpmgwg8Ph2MtUHD&open=AaDsvPpmgwg8Ph2MtUHD&pullRequest=2184
env: Env,
projectId: string,
body: WorkspaceResourceUploadBody,
Expand Down Expand Up @@ -752,6 +758,25 @@
const skillId = workspace.skill_id;
const agentType = workspace.agent_type;
const runtime = normalizeNullable(body.runtime) ?? 'vm';
const memoryWorkingSetMeanBytes = assertFiniteMetric(
body.summary.memoryWorkingSetMeanBytes,
'summary.memoryWorkingSetMeanBytes'
);
const memoryWorkingSetPeakBytes = assertFiniteMetric(
body.summary.memoryWorkingSetPeakBytes,
'summary.memoryWorkingSetPeakBytes'
);
const memoryWorkingSetSampleCount =
memoryWorkingSetMeanBytes == null
? 0
: (assertFiniteMetric(
body.summary.memoryWorkingSetSampleCount,
'summary.memoryWorkingSetSampleCount'
) ?? body.sampleCount);
assertFiniteInteger(memoryWorkingSetSampleCount, 'summary.memoryWorkingSetSampleCount');
if (memoryWorkingSetSampleCount > body.sampleCount) {
throw errors.badRequest('summary.memoryWorkingSetSampleCount cannot exceed sampleCount');
}
const values = {
id: summaryId,
projectId,
Expand All @@ -776,6 +801,9 @@
body.summary.memoryKernelPeakBytes,
'summary.memoryKernelPeakBytes'
),
memoryWorkingSetMeanBytes,
memoryWorkingSetPeakBytes,
memoryWorkingSetSampleCount,
ioReadBytes: assertFiniteMetric(body.summary.ioReadBytes, 'summary.ioReadBytes'),
ioWriteBytes: assertFiniteMetric(body.summary.ioWriteBytes, 'summary.ioWriteBytes'),
oomCount: Math.trunc(assertFiniteMetric(body.summary.oomCount, 'summary.oomCount') ?? 0),
Expand Down Expand Up @@ -829,6 +857,27 @@
END`,
memoryPeakBytes: sql`MAX(COALESCE(${schema.workspaceResourceSummaries.memoryPeakBytes}, 0), ${values.memoryPeakBytes ?? 0})`,
memoryKernelPeakBytes: sql`MAX(COALESCE(${schema.workspaceResourceSummaries.memoryKernelPeakBytes}, 0), ${values.memoryKernelPeakBytes ?? 0})`,
memoryWorkingSetMeanBytes:
values.memoryWorkingSetMeanBytes == null || values.memoryWorkingSetSampleCount === 0
? schema.workspaceResourceSummaries.memoryWorkingSetMeanBytes
: sql`CASE
WHEN ${schema.workspaceResourceSummaries.memoryWorkingSetMeanBytes} IS NULL
OR ${schema.workspaceResourceSummaries.memoryWorkingSetSampleCount} = 0
THEN ${values.memoryWorkingSetMeanBytes}
ELSE CAST((
(${schema.workspaceResourceSummaries.memoryWorkingSetMeanBytes} * ${schema.workspaceResourceSummaries.memoryWorkingSetSampleCount}) +
(${values.memoryWorkingSetMeanBytes} * ${values.memoryWorkingSetSampleCount})
) / (${schema.workspaceResourceSummaries.memoryWorkingSetSampleCount} + ${values.memoryWorkingSetSampleCount}) AS INTEGER)
END`,
memoryWorkingSetPeakBytes:
values.memoryWorkingSetPeakBytes == null
? schema.workspaceResourceSummaries.memoryWorkingSetPeakBytes
: sql`CASE
WHEN ${schema.workspaceResourceSummaries.memoryWorkingSetPeakBytes} IS NULL
THEN ${values.memoryWorkingSetPeakBytes}
ELSE MAX(${schema.workspaceResourceSummaries.memoryWorkingSetPeakBytes}, ${values.memoryWorkingSetPeakBytes})
END`,
memoryWorkingSetSampleCount: sql`${schema.workspaceResourceSummaries.memoryWorkingSetSampleCount} + ${values.memoryWorkingSetSampleCount}`,
ioReadBytes: sql`COALESCE(${schema.workspaceResourceSummaries.ioReadBytes}, 0) + ${values.ioReadBytes ?? 0}`,
ioWriteBytes: sql`COALESCE(${schema.workspaceResourceSummaries.ioWriteBytes}, 0) + ${values.ioWriteBytes ?? 0}`,
oomCount: sql`${schema.workspaceResourceSummaries.oomCount} + ${values.oomCount}`,
Expand Down Expand Up @@ -908,6 +957,8 @@
memoryMeanBytes: row.memoryMeanBytes,
memoryPeakBytes: row.memoryPeakBytes,
memoryKernelPeakBytes: row.memoryKernelPeakBytes,
memoryWorkingSetMeanBytes: row.memoryWorkingSetMeanBytes,
memoryWorkingSetPeakBytes: row.memoryWorkingSetPeakBytes,
ioReadBytes: row.ioReadBytes,
ioWriteBytes: row.ioWriteBytes,
oomCount: row.oomCount,
Expand Down Expand Up @@ -947,6 +998,7 @@
Number(sample.cpuMillis ?? 0),
Number(sample.memoryBytes ?? 0) / (1024 * 1024),
Number(sample.memoryPeakBytes ?? 0) / (1024 * 1024),
Number(sample.memoryWorkingSetBytes ?? 0) / (1024 * 1024),
sample.gap ? Number.MAX_SAFE_INTEGER : 0
);
}
Expand Down
28 changes: 28 additions & 0 deletions apps/api/tests/helpers/resource-history.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
export function base64(bytes: Uint8Array): string {
let binary = '';
for (const byte of bytes) binary += String.fromCharCode(byte);

Check warning on line 3 in apps/api/tests/helpers/resource-history.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `String.fromCodePoint()` over `String.fromCharCode()`.

See more on https://sonarcloud.io/project/issues?id=raphaeltm_simple-agent-manager&issues=AaDs3r0bIl3dGzKsx5Hz&open=AaDs3r0bIl3dGzKsx5Hz&pullRequest=2184
return btoa(binary);
}

export async function sha256Hex(bytes: Uint8Array): Promise<string> {
const digest = await crypto.subtle.digest('SHA-256', bytes);
return [...new Uint8Array(digest)].map((byte) => byte.toString(16).padStart(2, '0')).join('');
}

export async function gzipText(value: string): Promise<Uint8Array> {
return gzipBytes(new TextEncoder().encode(value));
}

export async function gzipBytes(value: Uint8Array): Promise<Uint8Array> {
const stream = new Blob([value]).stream().pipeThrough(new CompressionStream('gzip'));
return new Uint8Array(await new Response(stream).arrayBuffer());
}

export function gzipJson(value: unknown): Promise<Uint8Array> {
return gzipText(JSON.stringify(value));
}

export async function gunzipJson(bytes: Uint8Array): Promise<unknown> {
const stream = new Blob([bytes]).stream().pipeThrough(new DecompressionStream('gzip'));
return JSON.parse(await new Response(stream).text()) as unknown;
}
192 changes: 192 additions & 0 deletions apps/api/tests/integration/workspace-resource-history-http.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,192 @@
import Database from 'better-sqlite3';
import { Hono } from 'hono';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';

import * as schema from '../../src/db/schema';
import type { Env } from '../../src/env';
import { handleAppError } from '../../src/middleware/app-error-handler';
import type { AuthContext } from '../../src/middleware/auth';
import { projectResourceHistoryRoutes } from '../../src/routes/projects/workspace-resource-history';
import { workspaceResourceHistoryCallbackRoute } from '../../src/routes/projects/workspace-resource-history-callback';
import { verifyCallbackToken } from '../../src/services/jwt';
import { base64, gzipJson, sha256Hex } from '../helpers/resource-history';
import { createSchemaTables, createSqliteD1 } from '../helpers/sqlite-d1';

vi.mock('../../src/services/jwt', async (importOriginal) => ({
...(await importOriginal<typeof import('../../src/services/jwt')>()),
verifyCallbackToken: vi.fn(),
}));

function authContext(): AuthContext {
return {
user: {
id: 'member-1',
email: 'member-1@example.test',
name: 'Resource history owner',
avatarUrl: null,
role: 'user',
status: 'active',
},
session: {
id: 'browser-session-1',
token: 'browser-token',
expiresAt: new Date('2027-01-01T00:00:00.000Z'),
},
};
}

function makeR2() {
const objects = new Map<string, Uint8Array>();
return {
objects,
binding: {
put: async (key: string, value: Uint8Array) => {
objects.set(key, value);
return null;
},
delete: async (key: string) => objects.delete(key),
get: async (key: string) => {
const value = objects.get(key);
if (!value) return null;
return {
arrayBuffer: async () =>
value.buffer.slice(value.byteOffset, value.byteOffset + value.byteLength),
body: new Blob([value]).stream(),
};
},
} as unknown as R2Bucket,
};
}

describe('workspace resource history HTTP vertical slice', () => {
let sqlite: Database.Database;
let env: Env;
let app: Hono<{ Bindings: Env }>;
let r2: ReturnType<typeof makeR2>;

beforeEach(() => {
vi.clearAllMocks();
vi.mocked(verifyCallbackToken).mockResolvedValue({
workspace: 'node-1',
type: 'callback',
scope: 'node',
});

sqlite = new Database(':memory:');
createSchemaTables(sqlite, [
schema.projects,
schema.projectMembers,
schema.workspaces,
schema.tasks,
schema.agentSessions,
schema.agentProfiles,
schema.skills,
schema.workspaceResourceSummaries,
schema.workspaceResourceChunks,
]);
sqlite
.prepare(`INSERT INTO projects (id, user_id, name) VALUES ('proj-1', 'member-1', 'Project')`)
.run();
sqlite
.prepare(
`INSERT INTO project_members (project_id, user_id, role, status)
VALUES ('proj-1', 'member-1', 'owner', 'active')`
)
.run();
sqlite
.prepare(
`INSERT INTO workspaces (id, project_id, node_id, chat_session_id)
VALUES ('ws-1', 'proj-1', 'node-1', 'session-1')`
)
.run();

r2 = makeR2();
env = {
DATABASE: createSqliteD1(sqlite),
PROJECT_DATA_ARCHIVE_R2: r2.binding,
} as unknown as Env;
app = new Hono<{ Bindings: Env }>();
app.use('*', async (c, next) => {
c.set('auth', authContext());
await next();
});
app.onError(handleAppError);
app.route('/api/projects', workspaceResourceHistoryCallbackRoute);
app.route('/api/projects', projectResourceHistoryRoutes);
});

afterEach(() => sqlite.close());

it('uploads and reads working-set summaries and raw samples through real routes', async () => {
const payload = {
samples: [
{ t: 1_000, memoryBytes: 1024, memoryWorkingSetBytes: 512 },
{ t: 1_500, memoryBytes: 4096, memoryWorkingSetBytes: 1024 },
{ t: 2_000, memoryBytes: 2048, memoryWorkingSetBytes: 768 },
],
toolSpans: [],
gaps: [],
};
const compressed = await gzipJson(payload);
const body = {
workspaceId: 'ws-1',
nodeId: 'node-1',
sessionId: 'session-1',
taskId: null,
sourceVersion: 1,
chunkSequence: 0,
startedAt: 1_000,
endedAt: 2_000,
sampleCount: 3,
gapCount: 0,
toolSpanCount: 0,
compressedBase64: base64(compressed),
compressedBytes: compressed.byteLength,
uncompressedBytes: JSON.stringify(payload).length,
sha256: await sha256Hex(compressed),
completeness: { status: 'complete' },
summary: {
memoryMeanBytes: 2389,
memoryPeakBytes: 4096,
memoryKernelPeakBytes: 8192,
memoryWorkingSetMeanBytes: 768,
memoryWorkingSetPeakBytes: 1024,
memoryWorkingSetSampleCount: 3,
},
};

const upload = await app.fetch(
new Request('https://api.test/api/projects/proj-1/workspace-resource-history', {
method: 'POST',
headers: { Authorization: 'Bearer callback-token', 'Content-Type': 'application/json' },
body: JSON.stringify(body),
}),
env
);
expect(upload.status).toBe(200);
const uploaded = (await upload.json()) as { chunkId: string };
expect(r2.objects.size).toBe(1);

const read = await app.fetch(
new Request(
`https://api.test/api/projects/proj-1/sessions/session-1/resource-history?chunkId=${encodeURIComponent(uploaded.chunkId)}`
),
env
);
expect(read.status).toBe(200);
const history = (await read.json()) as {
summary: Record<string, unknown>;
detail: { samples: Array<Record<string, unknown>> };
};
expect(history.summary).toMatchObject({
memoryPeakBytes: 4096,
memoryKernelPeakBytes: 8192,
memoryWorkingSetMeanBytes: 768,
memoryWorkingSetPeakBytes: 1024,
});
expect(history.detail.samples.map((sample) => sample.memoryBytes)).toEqual([1024, 4096, 2048]);
expect(history.detail.samples.map((sample) => sample.memoryWorkingSetBytes)).toEqual([
512, 1024, 768,
]);
});
});
Loading
Loading