diff --git a/.claude/skills/api-reference/SKILL.md b/.claude/skills/api-reference/SKILL.md index 2faa998a6..87bb3b426 100644 --- a/.claude/skills/api-reference/SKILL.md +++ b/.claude/skills/api-reference/SKILL.md @@ -58,6 +58,8 @@ user-invocable: false - `GET /api/projects/:projectId/sessions/:sessionId/state` — Get lightweight ACP activity state for a chat session - `GET /api/projects/:projectId/sessions/:sessionId/messages` — List persisted session messages (supports `roles`, `before`, `after`, `limit`, `compact`, `order=asc|desc`). `before`/`after` accept an exact `[createdAt,sequence,id]` cursor taken from a page's edge row — required to resume inside a group of rows sharing one timestamp — or a legacy millisecond timestamp that excludes every row at that time. Without `order`, an `after`-only request reads ascending; anything else reads newest-first. Pages are always returned oldest-first. - `GET /api/projects/:projectId/sessions/:sessionId/messages/:messageId/tool-content` — Lazy-load stored tool content for compact messages, falling back to the private R2 archive when inline payloads have been stripped +- `GET /api/projects/:projectId/sessions/:sessionId/resource-timeline` — Whole-session resource timeline index: every retained chunk (ascending) with its summary and per-minute `rollup`, `runs` per workspace with `reservation`, `omittedChunkCount` past `WORKSPACE_RESOURCE_TIMELINE_MAX_CHUNKS`, and `collection` (`collected`/`pending`/`unsupported`/`expired`) explaining an empty timeline +- `GET /api/projects/:projectId/sessions/:sessionId/resource-timeline/chunks/:chunkId` — One chunk's 5-second samples and tool spans, scoped to the path's project and session (foreign chunk → 404) - `GET /api/projects/:projectId/sessions/:sessionId/comments` — List message-anchored comment threads (supports `messageId`, `status=open|sent|resolved`, `afterSequence`, `limit`) - `POST /api/projects/:projectId/sessions/:sessionId/comments` — Create a message-anchored comment thread (`{ messageId, body, quote?, clientMutationId? }`) - `POST /api/projects/:projectId/sessions/:sessionId/comments/:threadId/replies` — Append a comment reply (`{ body, clientMutationId? }`) diff --git a/.claude/skills/env-reference/SKILL.md b/.claude/skills/env-reference/SKILL.md index 1d499e015..24d83a3c3 100644 --- a/.claude/skills/env-reference/SKILL.md +++ b/.claude/skills/env-reference/SKILL.md @@ -192,6 +192,8 @@ Activity coalescing and binding caches are per Worker isolate. Delayed flushes c - `NODE_STOPPED_HANDOFF_REQUEST_TIMEOUT_MS` — Per-candidate provider/DNS deadline during stopped-node handoff, capped by remaining sweep time; provider failures enter cleanup backoff (default: `5000`) - `IDLE_CLEANUP_MAX_RESIDENCE_MS` — Maximum ProjectData idle-cleanup schedule residence before preserved/error outcomes stop re-arming and surface attention (default: `7200000`) - `IDLE_CLEANUP_MAX_CANDIDATES_PER_SWEEP` — Maximum idle-cleanup schedules or workspace idle checks one ProjectData alarm pass takes, and maximum reporter-scoped task candidates inspected per check (default: `5`) +- `WORKSPACE_RESOURCE_TIMELINE_MAX_CHUNKS` — Max chunks the session resource timeline lists; older ones are disclosed as omitted (default: `1000`) +- `WORKSPACE_RESOURCE_ROLLUP_BUCKET_MS` / `WORKSPACE_RESOURCE_ROLLUP_MAX_BUCKETS` — Per-chunk rollup bucket width computed on upload, and the bucket cap that widens it (defaults: `60000` / `60`) - `WORKSPACE_IDLE_TIMEOUT_MS` — Installation default for how long an active chat session's workspace can go without messages or terminal activity before ProjectData retires it, once its runtime is conclusively dead; the project Workspace Idle Timeout setting overrides it (default: `7200000`) - `WORKSPACE_IDLE_BACKOFF_BASE_MS` — First retry delay after a ProjectData workspace-idle check finds an idle workspace it cannot retire yet: inconclusive task candidates, a live or unprovable runtime, a missing project identity, or a failed check (default: `600000`; `apps/api/src/durable-objects/project-data/workspace-idle-timeouts.ts`) - `WORKSPACE_IDLE_BACKOFF_MAX_MS` — Maximum workspace-idle retry delay; the delay doubles from the base and resets on new activity or when the session wakes (default: `21600000`) diff --git a/apps/api/.env.example b/apps/api/.env.example index adfed65e9..c52125b4c 100644 --- a/apps/api/.env.example +++ b/apps/api/.env.example @@ -1026,6 +1026,9 @@ INFOMANIAK_IP_POLL_INTERVAL_MS=3000 # WORKSPACE_RESOURCE_LIST_LIMIT=24 # WORKSPACE_RESOURCE_CLEANUP_BATCH_SIZE=50 # WORKSPACE_RESOURCE_OBJECT_CLEANUP_LIMIT=5000 +# WORKSPACE_RESOURCE_TIMELINE_MAX_CHUNKS=1000 +# WORKSPACE_RESOURCE_ROLLUP_BUCKET_MS=60000 +# WORKSPACE_RESOURCE_ROLLUP_MAX_BUCKETS=60 # PROJECT_DATA_ARCHIVE_GLOBAL_SWEEP_ENABLED=false # Separate fail-closed switch for the unscoped scheduled archive-sharding sweep # PROJECT_DATA_ARCHIVE_GLOBAL_SWEEP_INTERVAL_MS=86400000 # Persisted cadence gate for unscoped scheduled archive-sharding sweeps (code fallback daily; wrangler.toml ships 1080000: one claim per 20 min on the 5-min cron) # PROJECT_DATA_ARCHIVE_SHARD_COUNT=128 # Deterministic ProjectData archive shard fanout diff --git a/apps/api/src/db/migrations/0177_workspace_resource_chunk_rollups.sql b/apps/api/src/db/migrations/0177_workspace_resource_chunk_rollups.sql new file mode 100644 index 000000000..17c57b607 --- /dev/null +++ b/apps/api/src/db/migrations/0177_workspace_resource_chunk_rollups.sql @@ -0,0 +1,8 @@ +-- Per-minute resource rollup for each stored chunk, computed on upload from the +-- decoded chunk payload. It lets the session resource timeline draw a whole +-- multi-day session from D1 alone, without downloading every R2 chunk. +-- +-- Existing rows stay NULL. They remain fully readable: the timeline falls back +-- to the chunk's own summary_json (one aggregate spanning the chunk), and the +-- full-resolution detail is still read from R2 on zoom. +ALTER TABLE workspace_resource_chunks ADD COLUMN rollup_json TEXT; diff --git a/apps/api/src/db/schema.ts b/apps/api/src/db/schema.ts index d470b26ba..604bf9076 100644 --- a/apps/api/src/db/schema.ts +++ b/apps/api/src/db/schema.ts @@ -1647,6 +1647,8 @@ export const workspaceResourceChunks = sqliteTable( toolSpanCount: integer('tool_span_count').notNull().default(0), completenessJson: text('completeness_json').notNull(), summaryJson: text('summary_json').notNull(), + /** Per-minute rollup computed on upload; NULL for chunks uploaded before migration 0177. */ + rollupJson: text('rollup_json'), createdAt: integer('created_at').notNull(), expiresAt: integer('expires_at').notNull(), uploadedByNodeId: text('uploaded_by_node_id'), diff --git a/apps/api/src/env.ts b/apps/api/src/env.ts index 5121d6e5f..1176d69a9 100644 --- a/apps/api/src/env.ts +++ b/apps/api/src/env.ts @@ -643,6 +643,9 @@ export interface Env extends WebhookTriggerEnv, TaskRecoveryEnv { WORKSPACE_RESOURCE_LIST_LIMIT?: string; // Max chunk indexes returned by detail list (default: 24) WORKSPACE_RESOURCE_CLEANUP_BATCH_SIZE?: string; // Max expired chunks/summaries cleaned per sweep (default: 50) WORKSPACE_RESOURCE_OBJECT_CLEANUP_LIMIT?: string; // Max R2 objects deleted for one project/workspace resource-history prefix cleanup (default: 5000) + WORKSPACE_RESOURCE_TIMELINE_MAX_CHUNKS?: string; // Max chunks listed by the whole-session resource timeline index; older ones are disclosed as omitted (default: 1000) + WORKSPACE_RESOURCE_ROLLUP_BUCKET_MS?: string; // Width of each per-chunk rollup bucket computed on upload (default: 60000) + WORKSPACE_RESOURCE_ROLLUP_MAX_BUCKETS?: string; // Max rollup buckets stored per chunk; the bucket width widens to fit (default: 60) PROJECT_DATA_STORAGE_TELEMETRY_LIST_LIMIT_DEFAULT?: string; PROJECT_DATA_STORAGE_TELEMETRY_LIST_LIMIT_MAX?: string; PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_ENABLED?: string; diff --git a/apps/api/src/routes/projects/workspace-resource-history.ts b/apps/api/src/routes/projects/workspace-resource-history.ts index 38e4dff61..87c11903b 100644 --- a/apps/api/src/routes/projects/workspace-resource-history.ts +++ b/apps/api/src/routes/projects/workspace-resource-history.ts @@ -7,6 +7,10 @@ import { getUserId } from '../../middleware/auth'; import { errors } from '../../middleware/error'; import { requireProjectAccess } from '../../middleware/project-auth'; import { getWorkspaceResourceHistory } from '../../services/workspace-resource-history'; +import { + getSessionResourceTimeline, + getSessionResourceTimelineChunk, +} from '../../services/workspace-resource-timeline'; const projectResourceHistoryRoutes = new Hono<{ Bindings: Env }>(); @@ -38,6 +42,25 @@ projectResourceHistoryRoutes.get('/:id/sessions/:sessionId/resource-history', as ); }); +projectResourceHistoryRoutes.get('/:id/sessions/:sessionId/resource-timeline', async (c) => { + const projectId = c.req.param('id'); + const sessionId = c.req.param('sessionId'); + await requireAccess(c.env, projectId, getUserId(c)); + return c.json(await getSessionResourceTimeline(c.env, { projectId, sessionId })); +}); + +projectResourceHistoryRoutes.get( + '/:id/sessions/:sessionId/resource-timeline/chunks/:chunkId', + async (c) => { + const projectId = c.req.param('id'); + const sessionId = c.req.param('sessionId'); + await requireAccess(c.env, projectId, getUserId(c)); + const chunkId = optionalDetailChunkId(c.req.param('chunkId')); + if (!chunkId) throw errors.badRequest('chunkId is required'); + return c.json(await getSessionResourceTimelineChunk(c.env, { projectId, sessionId, chunkId })); + } +); + projectResourceHistoryRoutes.get('/:id/tasks/:taskId/resource-history', async (c) => { const projectId = c.req.param('id'); const taskId = c.req.param('taskId'); diff --git a/apps/api/src/services/workspace-resource-history.ts b/apps/api/src/services/workspace-resource-history.ts index e48a6cb8d..73c16a880 100644 --- a/apps/api/src/services/workspace-resource-history.ts +++ b/apps/api/src/services/workspace-resource-history.ts @@ -8,6 +8,7 @@ import { log } from '../lib/logger'; import { parsePositiveInt } from '../lib/route-helpers'; import { parseJsonRecord } from '../lib/runtime-validation'; import { AppError, errors } from '../middleware/error'; +import { buildStoredRollupJson } from './workspace-resource-rollup'; export const WORKSPACE_RESOURCE_STORAGE_FORMAT = 'resource-history-gzip-json-v1'; const R2_PREFIX = 'resource-history/v1'; @@ -629,6 +630,7 @@ async function validateWorkspaceResourceChunkBytes( actualSha: string; compressedBytes: number; uncompressedBytes: number; + payload: WorkspaceResourceChunkPayload; }> { const bytes = decodeBase64(body.compressedBase64); if (bytes.byteLength !== body.compressedBytes) { @@ -671,6 +673,7 @@ async function validateWorkspaceResourceChunkBytes( actualSha: await sha256Hex(sanitized.bytes), compressedBytes: sanitized.bytes.byteLength, uncompressedBytes: sanitized.uncompressedBytes, + payload: decodedChunk.value, }; } @@ -694,7 +697,7 @@ export async function storeWorkspaceResourceChunk( ); const nodeId = validateWorkspaceUploadIdentity(workspace, body, uploadedByNodeId); assertWorkspaceResourceUploadMetadata(body); - const { bytes, actualSha, compressedBytes, uncompressedBytes } = + const { bytes, actualSha, compressedBytes, uncompressedBytes, payload } = await validateWorkspaceResourceChunkBytes(env, body); const db = drizzle(env.DATABASE, { schema }); @@ -911,6 +914,11 @@ export async function storeWorkspaceResourceChunk( toolSpanCount: body.toolSpanCount ?? 0, completenessJson, summaryJson, + rollupJson: buildStoredRollupJson(env, payload, body, { + projectId, + workspaceId, + chunkSequence: body.chunkSequence, + }), createdAt: now, expiresAt, uploadedByNodeId, @@ -1031,7 +1039,7 @@ export function downsamplePreservingSpikes( return { samples: selected, downsampled: true }; } -async function readChunkPayload(env: Env, chunk: schema.WorkspaceResourceChunkRow) { +export async function readChunkPayload(env: Env, chunk: schema.WorkspaceResourceChunkRow) { const object = await env.PROJECT_DATA_ARCHIVE_R2.get(chunk.r2Key); if (!object) throw errors.notFound('Resource history chunk'); const compressedBytes = new Uint8Array(await object.arrayBuffer()); diff --git a/apps/api/src/services/workspace-resource-rollup.ts b/apps/api/src/services/workspace-resource-rollup.ts new file mode 100644 index 000000000..f0354ee92 --- /dev/null +++ b/apps/api/src/services/workspace-resource-rollup.ts @@ -0,0 +1,273 @@ +/** + * Per-minute rollup of one resource-history chunk. + * + * Computed once on upload from the already-decoded chunk payload and stored in + * `workspace_resource_chunks.rollup_json`, so the session timeline can draw a + * whole multi-day session from D1 without downloading every R2 chunk. The + * layout is columnar (one array per metric) to keep a 15-minute chunk's rollup + * around a kilobyte. + */ +import type { Env } from '../env'; +import { log } from '../lib/logger'; +import { parsePositiveInt } from '../lib/route-helpers'; +import type { ResourceSamplePoint, ResourceToolSpan } from './workspace-resource-history'; + +export const WORKSPACE_RESOURCE_ROLLUP_VERSION = 1; +const DEFAULT_ROLLUP_BUCKET_MS = 60_000; +const DEFAULT_ROLLUP_MAX_BUCKETS = 60; + +/** Columnar rollup. Every array has one entry per bucket, ascending by `start`. */ +export interface WorkspaceResourceRollup { + v: typeof WORKSPACE_RESOURCE_ROLLUP_VERSION; + bucketMs: number; + start: number[]; + end: number[]; + /** Measured samples in the bucket (gaps and unsupported samples excluded). */ + samples: number[]; + /** CPU in cores (1 = one core fully busy). */ + cpuMeanCores: Array; + cpuMaxCores: Array; + memoryMeanBytes: Array; + memoryMaxBytes: Array; + workingSetMeanBytes: Array; + workingSetMaxBytes: Array; + ioReadBytes: Array; + ioWriteBytes: Array; + oomKills: number[]; + toolCallStarts: number[]; +} + +export interface RollupConfig { + bucketMs: number; + maxBuckets: number; +} + +export function getRollupConfig(env: Env): RollupConfig { + return { + bucketMs: parsePositiveInt(env.WORKSPACE_RESOURCE_ROLLUP_BUCKET_MS, DEFAULT_ROLLUP_BUCKET_MS), + maxBuckets: parsePositiveInt( + env.WORKSPACE_RESOURCE_ROLLUP_MAX_BUCKETS, + DEFAULT_ROLLUP_MAX_BUCKETS + ), + }; +} + +interface Accumulator { + start: number; + samples: number; + cpuSum: number; + cpuCount: number; + cpuMax: number | null; + memorySum: number; + memoryCount: number; + memoryMax: number | null; + workingSetSum: number; + workingSetCount: number; + workingSetMax: number | null; + ioRead: number | null; + ioWrite: number | null; + oomKills: number; + toolCallStarts: number; +} + +function finiteNonNegative(value: unknown): number | null { + return typeof value === 'number' && Number.isFinite(value) && value >= 0 ? value : null; +} + +function emptyAccumulator(start: number): Accumulator { + return { + start, + samples: 0, + cpuSum: 0, + cpuCount: 0, + cpuMax: null, + memorySum: 0, + memoryCount: 0, + memoryMax: null, + workingSetSum: 0, + workingSetCount: 0, + workingSetMax: null, + ioRead: null, + ioWrite: null, + oomKills: 0, + toolCallStarts: 0, + }; +} + +function max(current: number | null, value: number): number { + return current == null ? value : Math.max(current, value); +} + +/** + * Bucket width: the configured width, widened in whole multiples until the + * chunk's span fits in `maxBuckets`. A normal 15-minute chunk keeps one-minute + * buckets; a malformed or unusually long one cannot produce an unbounded rollup. + */ +function bucketWidth(spanMs: number, config: RollupConfig): number { + const needed = Math.ceil(Math.max(spanMs, 1) / config.maxBuckets); + return Math.max(1, Math.ceil(needed / config.bucketMs)) * config.bucketMs; +} + +export function computeWorkspaceResourceRollup( + payload: { samples?: ResourceSamplePoint[]; toolSpans?: ResourceToolSpan[] }, + window: { startedAt: number; endedAt: number }, + config: RollupConfig +): WorkspaceResourceRollup { + const width = bucketWidth(window.endedAt - window.startedAt, config); + const buckets = new Map(); + const bucketFor = (t: number): Accumulator => { + // Clamp into the chunk window so a stray timestamp cannot add buckets beyond the bound. + const clamped = Math.min(Math.max(t, window.startedAt), window.endedAt); + const key = Math.floor(clamped / width) * width; + let bucket = buckets.get(key); + if (!bucket) { + bucket = emptyAccumulator(key); + buckets.set(key, bucket); + } + return bucket; + }; + + for (const sample of payload.samples ?? []) { + if (sample.gap || sample.unsupported) continue; + // A sample's timestamp closes the interval it measures, (t - interval, t]; a + // sample stamped exactly on a minute boundary belongs to the minute it ends. + const bucket = bucketFor(sample.t - 1); + bucket.samples += 1; + + const intervalMs = finiteNonNegative(sample.intervalMillis); + const cpuMillis = finiteNonNegative(sample.cpuMillis); + if (intervalMs && cpuMillis != null) { + const cores = cpuMillis / intervalMs; + bucket.cpuSum += cores; + bucket.cpuCount += 1; + bucket.cpuMax = max(bucket.cpuMax, cores); + } + + const memory = finiteNonNegative(sample.memoryBytes); + if (memory != null) { + bucket.memorySum += memory; + bucket.memoryCount += 1; + bucket.memoryMax = max(bucket.memoryMax, memory); + } + // The kernel's memory.peak catches spikes between two samples. + const memoryPeak = finiteNonNegative(sample.memoryPeakBytes); + if (memoryPeak != null) bucket.memoryMax = max(bucket.memoryMax, memoryPeak); + + const workingSet = finiteNonNegative(sample.memoryWorkingSetBytes); + if (workingSet != null) { + bucket.workingSetSum += workingSet; + bucket.workingSetCount += 1; + bucket.workingSetMax = max(bucket.workingSetMax, workingSet); + } + + // I/O counters are per-sample deltas, so the bucket total is their sum. + const ioRead = finiteNonNegative(sample.ioReadBytes); + if (ioRead != null) bucket.ioRead = (bucket.ioRead ?? 0) + ioRead; + const ioWrite = finiteNonNegative(sample.ioWriteBytes); + if (ioWrite != null) bucket.ioWrite = (bucket.ioWrite ?? 0) + ioWrite; + + // Same accounting as the collector's chunk summary: oom + oom_kill deltas. + bucket.oomKills += + (finiteNonNegative(sample.oom) ?? 0) + (finiteNonNegative(sample.oomKill) ?? 0); + } + + for (const span of payload.toolSpans ?? []) { + bucketFor(span.startedAt).toolCallStarts += 1; + } + + const ordered = [...buckets.values()].sort((a, b) => a.start - b.start); + const mean = (sum: number, count: number) => (count > 0 ? sum / count : null); + return { + v: WORKSPACE_RESOURCE_ROLLUP_VERSION, + bucketMs: width, + start: ordered.map((b) => Math.max(b.start, window.startedAt)), + end: ordered.map((b) => Math.min(b.start + width, window.endedAt)), + samples: ordered.map((b) => b.samples), + cpuMeanCores: ordered.map((b) => mean(b.cpuSum, b.cpuCount)), + cpuMaxCores: ordered.map((b) => b.cpuMax), + memoryMeanBytes: ordered.map((b) => { + const value = mean(b.memorySum, b.memoryCount); + return value == null ? null : Math.round(value); + }), + memoryMaxBytes: ordered.map((b) => b.memoryMax), + workingSetMeanBytes: ordered.map((b) => { + const value = mean(b.workingSetSum, b.workingSetCount); + return value == null ? null : Math.round(value); + }), + workingSetMaxBytes: ordered.map((b) => b.workingSetMax), + ioReadBytes: ordered.map((b) => b.ioRead), + ioWriteBytes: ordered.map((b) => b.ioWrite), + oomKills: ordered.map((b) => b.oomKills), + toolCallStarts: ordered.map((b) => b.toolCallStarts), + }; +} + +const ROLLUP_COLUMNS = [ + 'start', + 'end', + 'samples', + 'cpuMeanCores', + 'cpuMaxCores', + 'memoryMeanBytes', + 'memoryMaxBytes', + 'workingSetMeanBytes', + 'workingSetMaxBytes', + 'ioReadBytes', + 'ioWriteBytes', + 'oomKills', + 'toolCallStarts', +] as const; + +/** + * Parses a stored rollup, or returns null when it is absent or malformed. A + * malformed rollup must never break the timeline: the caller falls back to the + * chunk summary. + */ +export function parseWorkspaceResourceRollup(raw: string | null): WorkspaceResourceRollup | null { + if (!raw) return null; + let value: unknown; + try { + value = JSON.parse(raw); + } catch { + return null; + } + if (typeof value !== 'object' || value === null) return null; + const record = value as Record; + if (record.v !== WORKSPACE_RESOURCE_ROLLUP_VERSION) return null; + if (typeof record.bucketMs !== 'number' || !(record.bucketMs > 0)) return null; + const starts = record.start; + if (!Array.isArray(starts)) return null; + for (const column of ROLLUP_COLUMNS) { + const entries = record[column]; + if (!Array.isArray(entries) || entries.length !== starts.length) return null; + if ( + !entries.every( + (entry) => entry === null || (typeof entry === 'number' && Number.isFinite(entry)) + ) + ) { + return null; + } + } + return value as WorkspaceResourceRollup; +} + +/** + * The stored rollup is an optimisation for the timeline: failing to build one + * must never reject the chunk upload. The timeline falls back to the summary. + */ +export function buildStoredRollupJson( + env: Env, + payload: { samples?: ResourceSamplePoint[]; toolSpans?: ResourceToolSpan[] }, + window: { startedAt: number; endedAt: number }, + context: { projectId: string; workspaceId: string; chunkSequence: number } +): string | null { + try { + return JSON.stringify(computeWorkspaceResourceRollup(payload, window, getRollupConfig(env))); + } catch (error) { + log.warn('workspace_resource_history.rollup_failed', { + ...context, + error: error instanceof Error ? error.message : String(error), + }); + return null; + } +} diff --git a/apps/api/src/services/workspace-resource-timeline.ts b/apps/api/src/services/workspace-resource-timeline.ts new file mode 100644 index 000000000..79d99564c --- /dev/null +++ b/apps/api/src/services/workspace-resource-timeline.ts @@ -0,0 +1,322 @@ +/** + * Whole-session resource timeline. + * + * A session's resource history is stored as immutable ~15-minute chunks, one + * series per workspace lifetime (every wake provisions a fresh workspace). The + * index returns every chunk of the session with its per-minute rollup, so the + * client can draw the whole session from D1 alone and fetch full-resolution + * chunks only for the window the user zooms into. + */ +import { and, count, desc, eq } from 'drizzle-orm'; +import { drizzle } from 'drizzle-orm/d1'; + +import * as schema from '../db/schema'; +import type { Env } from '../env'; +import { log } from '../lib/logger'; +import { parsePositiveInt } from '../lib/route-helpers'; +import { errors } from '../middleware/error'; +import { parseStoredResolvedReservationJson } from './resource-requirements-input'; +import { readChunkPayload } from './workspace-resource-history'; +import { + parseWorkspaceResourceRollup, + type WorkspaceResourceRollup, +} from './workspace-resource-rollup'; + +const DEFAULT_TIMELINE_MAX_CHUNKS = 1000; +/** `nodes.runtime` for Instant sessions, whose container runtime records no resource history yet. */ +const UNSUPPORTED_RUNTIMES: ReadonlySet = new Set(['cf-container']); + +export interface ResourceTimelineChunk { + id: string; + workspaceId: string; + startedAt: number; + endedAt: number; + sampleCount: number; + gapCount: number; + toolSpanCount: number; + summary: unknown; + completeness: unknown; + /** Per-minute rollup; null for chunks uploaded before rollups existed. */ + rollup: WorkspaceResourceRollup | null; +} + +export interface ResourceTimelineRun { + workspaceId: string; + nodeId: string | null; + runtime: string | null; + startedAt: number; + endedAt: number; + /** What the workspace reserved; null when unknown or unparseable. */ + reservation: { cpuMillis: number; memoryMb: number } | null; +} + +export interface ResourceTimelineIndexResponse { + sessionId: string; + runs: ResourceTimelineRun[]; + /** Ascending by `startedAt`. */ + chunks: ResourceTimelineChunk[]; + totalChunkCount: number; + /** Oldest chunks left out because the session exceeds `maxChunks`. */ + omittedChunkCount: number; + maxChunks: number; + /** + * Why an empty timeline is empty: `unsupported` when the session's runtime + * collects no resource history, `expired` when samples were collected but have + * passed their retention (the session summary outlives them), and `pending` + * when nothing has been uploaded yet. + */ + collection: 'collected' | 'pending' | 'unsupported' | 'expired'; + runtime: string | null; +} + +export function getTimelineMaxChunks(env: Env): number { + return parsePositiveInt(env.WORKSPACE_RESOURCE_TIMELINE_MAX_CHUNKS, DEFAULT_TIMELINE_MAX_CHUNKS); +} + +function jsonOrNull(value: string): unknown { + try { + return JSON.parse(value); + } catch { + return null; + } +} + +/** Only what the index needs: skips storage bookkeeping (R2 key, checksum, sizes) on up to maxChunks rows. */ +const TIMELINE_CHUNK_COLUMNS = { + id: schema.workspaceResourceChunks.id, + workspaceId: schema.workspaceResourceChunks.workspaceId, + startedAt: schema.workspaceResourceChunks.startedAt, + endedAt: schema.workspaceResourceChunks.endedAt, + sampleCount: schema.workspaceResourceChunks.sampleCount, + gapCount: schema.workspaceResourceChunks.gapCount, + toolSpanCount: schema.workspaceResourceChunks.toolSpanCount, + summaryJson: schema.workspaceResourceChunks.summaryJson, + completenessJson: schema.workspaceResourceChunks.completenessJson, + rollupJson: schema.workspaceResourceChunks.rollupJson, +}; + +type TimelineChunkRow = { + [K in keyof typeof TIMELINE_CHUNK_COLUMNS]: schema.WorkspaceResourceChunkRow[K]; +}; + +function timelineChunk(row: TimelineChunkRow): ResourceTimelineChunk { + return { + id: row.id, + workspaceId: row.workspaceId, + startedAt: row.startedAt, + endedAt: row.endedAt, + sampleCount: row.sampleCount, + gapCount: row.gapCount, + toolSpanCount: row.toolSpanCount, + summary: jsonOrNull(row.summaryJson), + completeness: jsonOrNull(row.completenessJson), + rollup: parseWorkspaceResourceRollup(row.rollupJson), + }; +} + +interface RunWorkspaceRow { + id: string; + node_id: string | null; + resolved_reservation_json: string | null; + runtime: string | null; +} + +function reservationOf(row: RunWorkspaceRow): ResourceTimelineRun['reservation'] { + try { + const reservation = parseStoredResolvedReservationJson(row.resolved_reservation_json); + return reservation + ? { cpuMillis: reservation.cpuMillis, memoryMb: reservation.memoryMb } + : null; + } catch (error) { + // A malformed reservation hides the reservation line for that run; it must not fail the index. + log.warn('workspace_resource_timeline.reservation_skipped', { + workspaceId: row.id, + error: error instanceof Error ? error.message : String(error), + }); + return null; + } +} + +/** Workspaces behind the session's chunks, scoped to the project on both sides of the join. */ +async function loadRunWorkspaces( + env: Env, + projectId: string, + sessionId: string +): Promise> { + const { results } = await env.DATABASE.prepare( + `SELECT w.id, w.node_id, w.resolved_reservation_json, n.runtime + FROM workspaces w + LEFT JOIN nodes n ON n.id = w.node_id + WHERE w.project_id = ? + AND w.id IN ( + SELECT DISTINCT workspace_id + FROM workspace_resource_chunks + WHERE project_id = ? AND session_id = ? + )` + ) + .bind(projectId, projectId, sessionId) + .all(); + return new Map(results.map((row) => [row.id, row])); +} + +/** Runtime of the session's current workspace, used only to explain an empty timeline. */ +async function loadSessionRuntime( + env: Env, + projectId: string, + sessionId: string +): Promise { + const row = await env.DATABASE.prepare( + `SELECT n.runtime + FROM workspaces w + JOIN nodes n ON n.id = w.node_id + WHERE w.project_id = ? AND w.chat_session_id = ? + ORDER BY w.created_at DESC + LIMIT 1` + ) + .bind(projectId, sessionId) + .first<{ runtime: string | null }>(); + return row?.runtime ?? null; +} + +/** A retained summary with no chunks means the samples existed and expired. */ +async function hasSessionSummary(env: Env, projectId: string, sessionId: string): Promise { + const row = await env.DATABASE.prepare( + `SELECT 1 AS present + FROM workspace_resource_summaries + WHERE project_id = ? AND session_id = ? + LIMIT 1` + ) + .bind(projectId, sessionId) + .first<{ present: number }>(); + return row != null; +} + +function groupRuns( + chunks: readonly ResourceTimelineChunk[], + workspaces: ReadonlyMap +): ResourceTimelineRun[] { + const runs = new Map(); + for (const chunk of chunks) { + const existing = runs.get(chunk.workspaceId); + if (existing) { + existing.startedAt = Math.min(existing.startedAt, chunk.startedAt); + existing.endedAt = Math.max(existing.endedAt, chunk.endedAt); + continue; + } + const workspace = workspaces.get(chunk.workspaceId); + runs.set(chunk.workspaceId, { + workspaceId: chunk.workspaceId, + nodeId: workspace?.node_id ?? null, + runtime: workspace?.runtime ?? null, + startedAt: chunk.startedAt, + endedAt: chunk.endedAt, + reservation: workspace ? reservationOf(workspace) : null, + }); + } + return [...runs.values()].sort((a, b) => a.startedAt - b.startedAt); +} + +export async function getSessionResourceTimeline( + env: Env, + input: { projectId: string; sessionId: string } +): Promise { + const { projectId, sessionId } = input; + const db = drizzle(env.DATABASE, { schema }); + const maxChunks = getTimelineMaxChunks(env); + const scope = and( + eq(schema.workspaceResourceChunks.projectId, projectId), + eq(schema.workspaceResourceChunks.sessionId, sessionId) + ); + + // Newest first so that, past the cap, it is the oldest chunks that are left out (and disclosed). + // Ordered by `started_at` alone so idx_workspace_resource_chunks_project_session serves the + // ORDER BY and LIMIT stops the scan; an `id` tie-break adds a sort over the session's whole + // history. Chunks within a run are sequential, so equal start times do not occur in practice. + const rows = await db + .select(TIMELINE_CHUNK_COLUMNS) + .from(schema.workspaceResourceChunks) + .where(scope) + .orderBy(desc(schema.workspaceResourceChunks.startedAt)) + .limit(maxChunks); + + let totalChunkCount = rows.length; + if (rows.length >= maxChunks) { + const [total] = await db + .select({ value: count() }) + .from(schema.workspaceResourceChunks) + .where(scope); + totalChunkCount = total?.value ?? rows.length; + } + + const chunks: ResourceTimelineChunk[] = []; + for (const row of rows) { + try { + chunks.push(timelineChunk(row)); + } catch (error) { + log.warn('workspace_resource_timeline.chunk_skipped', { + projectId, + sessionId, + chunkId: row.id, + error: error instanceof Error ? error.message : String(error), + }); + } + } + chunks.reverse(); + + if (chunks.length === 0) { + const [runtime, hasSummary] = await Promise.all([ + loadSessionRuntime(env, projectId, sessionId), + hasSessionSummary(env, projectId, sessionId), + ]); + return { + sessionId, + runs: [], + chunks, + totalChunkCount, + omittedChunkCount: 0, + maxChunks, + collection: + runtime && UNSUPPORTED_RUNTIMES.has(runtime) + ? 'unsupported' + : hasSummary + ? 'expired' + : 'pending', + runtime, + }; + } + + const runs = groupRuns(chunks, await loadRunWorkspaces(env, projectId, sessionId)); + return { + sessionId, + runs, + chunks, + totalChunkCount, + omittedChunkCount: Math.max(0, totalChunkCount - rows.length), + maxChunks, + collection: 'collected', + runtime: runs.at(-1)?.runtime ?? null, + }; +} + +export async function getSessionResourceTimelineChunk( + env: Env, + input: { projectId: string; sessionId: string; chunkId: string } +) { + const db = drizzle(env.DATABASE, { schema }); + const chunk = await db + .select() + .from(schema.workspaceResourceChunks) + .where( + and( + eq(schema.workspaceResourceChunks.id, input.chunkId), + eq(schema.workspaceResourceChunks.projectId, input.projectId), + eq(schema.workspaceResourceChunks.sessionId, input.sessionId) + ) + ) + .get(); + // Defence in depth: the predicate already scopes, but a missing or foreign chunk must look identical. + if (!chunk || chunk.projectId !== input.projectId || chunk.sessionId !== input.sessionId) { + throw errors.notFound('Resource history chunk'); + } + return readChunkPayload(env, chunk); +} diff --git a/apps/api/tests/unit/routes/workspace-resource-history-routes.test.ts b/apps/api/tests/unit/routes/workspace-resource-history-routes.test.ts index a4a0614cc..7cee0056a 100644 --- a/apps/api/tests/unit/routes/workspace-resource-history-routes.test.ts +++ b/apps/api/tests/unit/routes/workspace-resource-history-routes.test.ts @@ -10,6 +10,8 @@ const mocks = vi.hoisted(() => ({ db: { id: 'mock-db' }, requireProjectAccess: vi.fn(), getWorkspaceResourceHistory: vi.fn(), + getSessionResourceTimeline: vi.fn(), + getSessionResourceTimelineChunk: vi.fn(), })); vi.mock('drizzle-orm/d1', () => ({ @@ -28,6 +30,11 @@ vi.mock('../../../src/services/workspace-resource-history', () => ({ getWorkspaceResourceHistory: mocks.getWorkspaceResourceHistory, })); +vi.mock('../../../src/services/workspace-resource-timeline', () => ({ + getSessionResourceTimeline: mocks.getSessionResourceTimeline, + getSessionResourceTimelineChunk: mocks.getSessionResourceTimelineChunk, +})); + describe('workspace resource history project routes', () => { let app: Hono<{ Bindings: Env }>; const env = { DATABASE: {} as Env['DATABASE'] } as Env; @@ -108,4 +115,64 @@ describe('workspace resource history project routes', () => { expect(res.status).toBe(404); expect(mocks.getWorkspaceResourceHistory).not.toHaveBeenCalled(); }); + + it('serves the whole-session timeline index after project authorization', async () => { + mocks.getSessionResourceTimeline.mockResolvedValueOnce({ sessionId: 'sess-1', chunks: [] }); + + const res = await app.request( + '/api/projects/proj-1/sessions/sess-1/resource-timeline', + {}, + env + ); + + expect(res.status).toBe(200); + expect(mocks.requireProjectAccess).toHaveBeenCalledWith(mocks.db, 'proj-1', 'member-1'); + expect(mocks.getSessionResourceTimeline).toHaveBeenCalledWith(env, { + projectId: 'proj-1', + sessionId: 'sess-1', + }); + expect(await res.json()).toEqual({ sessionId: 'sess-1', chunks: [] }); + }); + + it('serves one timeline chunk scoped to the project and session in the path', async () => { + mocks.getSessionResourceTimelineChunk.mockResolvedValueOnce({ + chunkId: 'wrchunk:1', + samples: [], + }); + + const res = await app.request( + '/api/projects/proj-1/sessions/sess-1/resource-timeline/chunks/wrchunk%3A1', + {}, + env + ); + + expect(res.status).toBe(200); + expect(mocks.getSessionResourceTimelineChunk).toHaveBeenCalledWith(env, { + projectId: 'proj-1', + sessionId: 'sess-1', + chunkId: 'wrchunk:1', + }); + }); + + it('reads no timeline data when project access is denied', async () => { + mocks.requireProjectAccess.mockRejectedValue( + Object.assign(new Error('Project not found'), { statusCode: 404, error: 'NOT_FOUND' }) + ); + + const index = await app.request( + '/api/projects/proj-1/sessions/sess-1/resource-timeline', + {}, + env + ); + const chunk = await app.request( + '/api/projects/proj-1/sessions/sess-1/resource-timeline/chunks/wrchunk%3A1', + {}, + env + ); + + expect(index.status).toBe(404); + expect(chunk.status).toBe(404); + expect(mocks.getSessionResourceTimeline).not.toHaveBeenCalled(); + expect(mocks.getSessionResourceTimelineChunk).not.toHaveBeenCalled(); + }); }); diff --git a/apps/api/tests/unit/workspace-resource-timeline.test.ts b/apps/api/tests/unit/workspace-resource-timeline.test.ts new file mode 100644 index 000000000..6ab2ef931 --- /dev/null +++ b/apps/api/tests/unit/workspace-resource-timeline.test.ts @@ -0,0 +1,523 @@ +import Database from 'better-sqlite3'; +import { describe, expect, it } from 'vitest'; + +import * as schema from '../../src/db/schema'; +import type { Env } from '../../src/env'; +import { + storeWorkspaceResourceChunk, + type WorkspaceResourceUploadBody, +} from '../../src/services/workspace-resource-history'; +import { + buildStoredRollupJson, + computeWorkspaceResourceRollup, + parseWorkspaceResourceRollup, +} from '../../src/services/workspace-resource-rollup'; +import { + getSessionResourceTimeline, + getSessionResourceTimelineChunk, +} from '../../src/services/workspace-resource-timeline'; +import { createSchemaTables, createSqliteD1 } from '../helpers/sqlite-d1'; + +const MINUTE = 60_000; +const T0 = 1_790_000_000_000 - (1_790_000_000_000 % (15 * MINUTE)); +const CONFIG = { bucketMs: MINUTE, maxBuckets: 60 }; +const RESERVATION = { + cpuMillis: 2_000, + memoryMb: 4_096, + diskMb: 40_960, + exclusiveNode: false, + source: 'platform', + sourceId: 'platform', + version: 1, +}; + +function base64(bytes: Uint8Array): string { + let binary = ''; + for (const byte of bytes) binary += String.fromCharCode(byte); + return btoa(binary); +} + +async function sha256Hex(bytes: Uint8Array): Promise { + const digest = await crypto.subtle.digest('SHA-256', bytes); + return [...new Uint8Array(digest)].map((byte) => byte.toString(16).padStart(2, '0')).join(''); +} + +async function gzipJson(value: unknown): Promise { + const stream = new Blob([JSON.stringify(value)]) + .stream() + .pipeThrough(new CompressionStream('gzip')); + return new Uint8Array(await new Response(stream).arrayBuffer()); +} + +function makeR2(): R2Bucket { + const objects = new Map(); + return { + 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), + }; + }, + } as unknown as R2Bucket; +} + +/** Twelve 5-second samples per minute across a 15-minute chunk. */ +function chunkPayload(start: number, cpuMillis = 2_500) { + const samples = Array.from({ length: 180 }, (_, index) => ({ + t: start + (index + 1) * 5_000, + intervalMillis: 5_000, + cpuMillis, + memoryBytes: 1_000_000 + index, + ioReadBytes: 10, + ioWriteBytes: 20, + })); + return { + samples, + toolSpans: [{ id: 'span-1', kind: 'execute', toolName: 'Bash', startedAt: start + 90_000 }], + }; +} + +async function upload( + env: Env, + input: { + projectId: string; + workspaceId: string; + sessionId: string; + sequence: number; + start: number; + } +) { + const payload = chunkPayload(input.start); + const json = JSON.stringify(payload); + const compressed = await gzipJson(payload); + const body: WorkspaceResourceUploadBody = { + workspaceId: input.workspaceId, + sessionId: input.sessionId, + sourceVersion: 1, + chunkSequence: input.sequence, + startedAt: input.start, + endedAt: input.start + 15 * MINUTE, + sampleCount: payload.samples.length, + toolSpanCount: 1, + compressedBase64: base64(compressed), + compressedBytes: compressed.byteLength, + uncompressedBytes: new TextEncoder().encode(json).byteLength, + sha256: await sha256Hex(compressed), + completeness: { status: 'complete' }, + summary: { cpuMeanMillis: 2_500, cpuPeakMillis: 2_500, sampleIntervalMillis: 5_000 }, + }; + return storeWorkspaceResourceChunk(env, input.projectId, body, null); +} + +function setup(envOverrides: Record = {}) { + const sqlite = new Database(':memory:'); + createSchemaTables(sqlite, [ + schema.nodes, + schema.workspaces, + schema.tasks, + schema.agentSessions, + schema.agentProfiles, + schema.skills, + schema.workspaceResourceSummaries, + schema.workspaceResourceChunks, + ]); + const env = { + DATABASE: createSqliteD1(sqlite), + PROJECT_DATA_ARCHIVE_R2: makeR2(), + ...envOverrides, + } as unknown as Env; + const addNode = (id: string, runtime: string) => + sqlite.prepare(`INSERT INTO nodes (id, runtime) VALUES (?, ?)`).run(id, runtime); + const addWorkspace = ( + id: string, + projectId: string, + nodeId: string, + sessionId: string | null, + reservation: string | null = null, + createdAt = '2026-09-29T00:00:00Z' + ) => + sqlite + .prepare( + `INSERT INTO workspaces (id, project_id, node_id, chat_session_id, resolved_reservation_json, created_at) + VALUES (?, ?, ?, ?, ?, ?)` + ) + .run(id, projectId, nodeId, sessionId, reservation, createdAt); + return { sqlite, env, addNode, addWorkspace }; +} + +describe('computeWorkspaceResourceRollup', () => { + it('buckets samples per minute with CPU in cores, memory mean/max and summed I/O', () => { + const rollup = computeWorkspaceResourceRollup( + chunkPayload(T0), + { startedAt: T0, endedAt: T0 + 15 * MINUTE }, + CONFIG + ); + + expect(rollup.bucketMs).toBe(MINUTE); + expect(rollup.start).toHaveLength(15); + expect(rollup.start[0]).toBe(T0); + expect(rollup.end.at(-1)).toBe(T0 + 15 * MINUTE); + // 2 500 ms of CPU in a 5 000 ms interval is half a core. + expect(rollup.cpuMeanCores[0]).toBeCloseTo(0.5); + expect(rollup.cpuMaxCores[0]).toBeCloseTo(0.5); + expect(rollup.memoryMaxBytes[0]).toBeGreaterThanOrEqual(rollup.memoryMeanBytes[0] ?? Infinity); + // A sample stamped T0+60s closes the first minute, so every minute holds exactly twelve. + expect(rollup.samples).toEqual(Array.from({ length: 15 }, () => 12)); + expect(rollup.ioReadBytes.reduce((a, b) => a + (b ?? 0), 0)).toBe(1_800); + expect(rollup.toolCallStarts[1]).toBe(1); + expect(rollup.workingSetMeanBytes.every((value) => value === null)).toBe(true); + }); + + it('skips gap and unsupported samples, reads working set and memory.peak, and counts OOM kills', () => { + const rollup = computeWorkspaceResourceRollup( + { + samples: [ + { + t: T0 + 5_000, + intervalMillis: 5_000, + cpuMillis: 5_000, + memoryBytes: 100, + memoryWorkingSetBytes: 60, + }, + { t: T0 + 10_000, gap: true, cpuMillis: 99_999, memoryBytes: 99_999 }, + { t: T0 + 15_000, unsupported: 'no cgroup', memoryBytes: 99_999 }, + { + t: T0 + 20_000, + intervalMillis: 5_000, + cpuMillis: 0, + memoryBytes: 300, + memoryPeakBytes: 900, + memoryWorkingSetBytes: 80, + oom: 1, + oomKill: 1, + }, + ], + }, + { startedAt: T0, endedAt: T0 + MINUTE }, + CONFIG + ); + + expect(rollup.samples).toEqual([2]); + expect(rollup.cpuMeanCores[0]).toBeCloseTo(0.5); + expect(rollup.cpuMaxCores[0]).toBeCloseTo(1); + expect(rollup.memoryMeanBytes).toEqual([200]); + expect(rollup.memoryMaxBytes).toEqual([900]); + expect(rollup.workingSetMeanBytes).toEqual([70]); + expect(rollup.workingSetMaxBytes).toEqual([80]); + expect(rollup.oomKills).toEqual([2]); + }); + + it('widens the bucket so a long or malformed window stays within the configured bucket count', () => { + const day = 24 * 60 * MINUTE; + const samples = Array.from({ length: 2_000 }, (_, index) => ({ + t: T0 + index * 43_000, + intervalMillis: 5_000, + cpuMillis: 100, + })); + const rollup = computeWorkspaceResourceRollup( + { samples: [...samples, { t: T0 + 50 * day, intervalMillis: 5_000, cpuMillis: 1 }] }, + { startedAt: T0, endedAt: T0 + day }, + { bucketMs: MINUTE, maxBuckets: 10 } + ); + + expect(rollup.start.length).toBeLessThanOrEqual(11); + expect(rollup.bucketMs % MINUTE).toBe(0); + expect(rollup.end.at(-1)).toBeLessThanOrEqual(T0 + day); + }); + + it('degrades to no rollup instead of failing the upload when building one throws', () => { + const env = {} as Env; + const window = { startedAt: T0, endedAt: T0 + 15 * MINUTE }; + const context = { projectId: 'proj-1', workspaceId: 'ws-1', chunkSequence: 0 }; + const hostile = { + get samples(): never { + throw new Error('decoder bug'); + }, + }; + + expect(buildStoredRollupJson(env, hostile, window, context)).toBeNull(); + // Control: a well-formed payload still produces a stored rollup. + expect( + parseWorkspaceResourceRollup(buildStoredRollupJson(env, chunkPayload(T0), window, context)) + ).not.toBeNull(); + }); + + it('round-trips through JSON and rejects malformed stored rollups', () => { + const rollup = computeWorkspaceResourceRollup( + chunkPayload(T0), + { startedAt: T0, endedAt: T0 + 15 * MINUTE }, + CONFIG + ); + expect(parseWorkspaceResourceRollup(JSON.stringify(rollup))).toEqual(rollup); + expect(parseWorkspaceResourceRollup(null)).toBeNull(); + expect(parseWorkspaceResourceRollup('{not json')).toBeNull(); + expect(parseWorkspaceResourceRollup(JSON.stringify({ ...rollup, v: 99 }))).toBeNull(); + expect( + parseWorkspaceResourceRollup(JSON.stringify({ ...rollup, cpuMeanCores: [1] })) + ).toBeNull(); + expect( + parseWorkspaceResourceRollup( + JSON.stringify({ ...rollup, samples: rollup.samples.map(() => 'x') }) + ) + ).toBeNull(); + }); +}); + +describe('session resource timeline', () => { + it('stores a rollup on upload and serves every chunk of a multi-wake session in order', async () => { + const { env, addNode, addWorkspace } = setup(); + addNode('node-a', 'vm'); + addNode('node-b', 'vm'); + addWorkspace('ws-wake-1', 'proj-1', 'node-a', null, JSON.stringify(RESERVATION)); + addWorkspace('ws-wake-2', 'proj-1', 'node-b', 'sess-1', null); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-wake-1', + sessionId: 'sess-1', + sequence: 0, + start: T0, + }); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-wake-1', + sessionId: 'sess-1', + sequence: 1, + start: T0 + 15 * MINUTE, + }); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-wake-2', + sessionId: 'sess-1', + sequence: 0, + start: T0 + 10 * 60 * MINUTE, + }); + + const timeline = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-1', + }); + + expect(timeline.collection).toBe('collected'); + expect(timeline.chunks.map((chunk) => chunk.startedAt)).toEqual([ + T0, + T0 + 15 * MINUTE, + T0 + 10 * 60 * MINUTE, + ]); + expect(timeline.totalChunkCount).toBe(3); + expect(timeline.omittedChunkCount).toBe(0); + expect(timeline.chunks.every((chunk) => chunk.rollup?.start.length === 15)).toBe(true); + expect(timeline.runs.map((run) => [run.workspaceId, run.startedAt, run.endedAt])).toEqual([ + ['ws-wake-1', T0, T0 + 30 * MINUTE], + ['ws-wake-2', T0 + 10 * 60 * MINUTE, T0 + 10 * 60 * MINUTE + 15 * MINUTE], + ]); + expect(timeline.runs[0]?.reservation).toEqual({ cpuMillis: 2_000, memoryMb: 4_096 }); + expect(timeline.runs[1]?.reservation).toBeNull(); + expect(timeline.runtime).toBe('vm'); + }); + + it('excludes another project and another session (attack) while serving the owner (control)', async () => { + const { sqlite, env, addNode, addWorkspace } = setup(); + addNode('node-a', 'vm'); + addWorkspace('ws-owner', 'proj-1', 'node-a', 'sess-1'); + addWorkspace('ws-foreign', 'proj-2', 'node-a', 'sess-1'); + addWorkspace('ws-other-session', 'proj-1', 'node-a', 'sess-2'); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-owner', + sessionId: 'sess-1', + sequence: 0, + start: T0, + }); + await upload(env, { + projectId: 'proj-2', + workspaceId: 'ws-foreign', + sessionId: 'sess-1', + sequence: 0, + start: T0, + }); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-other-session', + sessionId: 'sess-2', + sequence: 0, + start: T0, + }); + expect(sqlite.prepare('SELECT COUNT(*) AS n FROM workspace_resource_chunks').get()).toEqual({ + n: 3, + }); + + const timeline = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-1', + }); + + expect(timeline.chunks.map((chunk) => chunk.workspaceId)).toEqual(['ws-owner']); + expect(timeline.runs.map((run) => run.workspaceId)).toEqual(['ws-owner']); + }); + + it('keeps the newest chunks past the cap and discloses how many older ones were left out', async () => { + const { env, addNode, addWorkspace } = setup({ WORKSPACE_RESOURCE_TIMELINE_MAX_CHUNKS: '2' }); + addNode('node-a', 'vm'); + addWorkspace('ws-1', 'proj-1', 'node-a', 'sess-1'); + for (let sequence = 0; sequence < 3; sequence += 1) { + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-1', + sessionId: 'sess-1', + sequence, + start: T0 + sequence * 15 * MINUTE, + }); + } + + const timeline = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-1', + }); + + expect(timeline.maxChunks).toBe(2); + expect(timeline.totalChunkCount).toBe(3); + expect(timeline.omittedChunkCount).toBe(1); + expect(timeline.chunks.map((chunk) => chunk.startedAt)).toEqual([ + T0 + 15 * MINUTE, + T0 + 30 * MINUTE, + ]); + }); + + it('tolerates a malformed rollup and a malformed reservation without failing the index', async () => { + const { sqlite, env, addNode, addWorkspace } = setup(); + addNode('node-a', 'vm'); + addWorkspace('ws-1', 'proj-1', 'node-a', 'sess-1', '{"source":"nope"}'); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-1', + sessionId: 'sess-1', + sequence: 0, + start: T0, + }); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-1', + sessionId: 'sess-1', + sequence: 1, + start: T0 + 15 * MINUTE, + }); + sqlite + .prepare( + `UPDATE workspace_resource_chunks SET rollup_json = '{broken' WHERE chunk_sequence = 0` + ) + .run(); + + const timeline = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-1', + }); + + expect(timeline.chunks).toHaveLength(2); + expect(timeline.chunks[0]?.rollup).toBeNull(); + expect(timeline.chunks[0]?.summary).toMatchObject({ cpuMeanMillis: 2_500 }); + expect(timeline.chunks[1]?.rollup).not.toBeNull(); + expect(timeline.runs[0]?.reservation).toBeNull(); + }); + + it('explains an empty timeline: Instant sessions are unsupported, VM sessions are pending', async () => { + const { env, addNode, addWorkspace } = setup(); + addNode('node-container', 'cf-container'); + addNode('node-vm', 'vm'); + addWorkspace('ws-instant', 'proj-1', 'node-container', 'sess-instant'); + addWorkspace('ws-vm', 'proj-1', 'node-vm', 'sess-vm'); + // Same session id in another project must not decide this project's answer. + addWorkspace('ws-foreign', 'proj-2', 'node-container', 'sess-vm', null, '2026-09-30T00:00:00Z'); + + const instant = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-instant', + }); + const vm = await getSessionResourceTimeline(env, { projectId: 'proj-1', sessionId: 'sess-vm' }); + const unknown = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-none', + }); + + expect(instant).toMatchObject({ + collection: 'unsupported', + runtime: 'cf-container', + chunks: [], + }); + expect(vm).toMatchObject({ collection: 'pending', runtime: 'vm', chunks: [] }); + expect(unknown).toMatchObject({ collection: 'pending', runtime: null }); + }); + + it('reports expired history when samples aged out but the session summary remains', async () => { + const { sqlite, env, addNode, addWorkspace } = setup(); + addNode('node-a', 'vm'); + addWorkspace('ws-1', 'proj-1', 'node-a', 'sess-old'); + addWorkspace('ws-2', 'proj-2', 'node-a', 'sess-foreign'); + await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-1', + sessionId: 'sess-old', + sequence: 0, + start: T0, + }); + await upload(env, { + projectId: 'proj-2', + workspaceId: 'ws-2', + sessionId: 'sess-foreign', + sequence: 0, + start: T0, + }); + // Retention cleanup deletes the chunks and keeps the longer-lived summary. + sqlite.prepare('DELETE FROM workspace_resource_chunks').run(); + + const expired = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-old', + }); + // Another project's summary for a session id must not make this project's answer "expired". + const foreign = await getSessionResourceTimeline(env, { + projectId: 'proj-1', + sessionId: 'sess-foreign', + }); + + expect(expired).toMatchObject({ collection: 'expired', chunks: [] }); + expect(foreign).toMatchObject({ collection: 'pending', chunks: [] }); + }); + + it('reads a chunk only through its own project and session', async () => { + const { env, addNode, addWorkspace } = setup(); + addNode('node-a', 'vm'); + addWorkspace('ws-1', 'proj-1', 'node-a', 'sess-1'); + const { chunkId } = await upload(env, { + projectId: 'proj-1', + workspaceId: 'ws-1', + sessionId: 'sess-1', + sequence: 0, + start: T0, + }); + + const owner = await getSessionResourceTimelineChunk(env, { + projectId: 'proj-1', + sessionId: 'sess-1', + chunkId, + }); + expect(owner.chunkId).toBe(chunkId); + expect(owner.samples).toHaveLength(180); + expect(owner.toolSpans[0]).toMatchObject({ toolName: 'Bash' }); + + await expect( + getSessionResourceTimelineChunk(env, { projectId: 'proj-1', sessionId: 'sess-2', chunkId }) + ).rejects.toMatchObject({ statusCode: 404 }); + await expect( + getSessionResourceTimelineChunk(env, { projectId: 'proj-2', sessionId: 'sess-1', chunkId }) + ).rejects.toMatchObject({ statusCode: 404 }); + }); +}); diff --git a/apps/web/package.json b/apps/web/package.json index 19ccc8383..c79439893 100644 --- a/apps/web/package.json +++ b/apps/web/package.json @@ -46,6 +46,7 @@ "recharts": "3.10.1", "remark-gfm": "catalog:", "tailwindcss": "4.3.3", + "uplot": "1.6.32", "valibot": "catalog:" }, "devDependencies": { diff --git a/apps/web/src/components/chat/SessionResourceHistoryDrawer.tsx b/apps/web/src/components/chat/SessionResourceHistoryDrawer.tsx index 7fc29fe13..59b087a63 100644 --- a/apps/web/src/components/chat/SessionResourceHistoryDrawer.tsx +++ b/apps/web/src/components/chat/SessionResourceHistoryDrawer.tsx @@ -1,559 +1,27 @@ import { Spinner } from '@simple-agent-manager/ui'; -import { keepPreviousData, useQuery } from '@tanstack/react-query'; -import { - Activity, - AlertTriangle, - ChevronRight, - Cpu, - Database, - HardDrive, - MemoryStick, - X, -} from 'lucide-react'; -import { useEffect, useMemo, useRef, useState } from 'react'; +import { Activity, X } from 'lucide-react'; +import { lazy, Suspense, useRef } from 'react'; import { createPortal } from 'react-dom'; -import { - getSessionResourceHistory, - type WorkspaceResourceChunk, - type WorkspaceResourceHistoryResponse, - type WorkspaceResourceSample, - type WorkspaceResourceSummary, - type WorkspaceResourceToolSpan, -} from '../../lib/api'; +import { importWithRetry } from '../../lib/lazy-with-retry'; +import type { ResourceHistorySource } from './resource-timeline/resource-source'; import { useDialogFocusTrap } from './useDialogFocusTrap'; +/** Loaded on first open so the chat route never downloads the charting code. */ +const ResourceTimeline = lazy(() => + importWithRetry(() => import('./resource-timeline/ResourceTimeline')).then((module) => ({ + default: module.ResourceTimeline, + })) +); + interface SessionResourceHistoryDrawerProps { - projectId: string; - sessionId: string; + source: ResourceHistorySource; onClose: () => void; } -function formatBytes(value: number | null | undefined): string { - if (value == null || !Number.isFinite(value)) return '—'; - if (value <= 0) return '0 B'; - const units = ['B', 'KB', 'MB', 'GB', 'TB']; - let next = value; - let unit = 0; - while (next >= 1024 && unit < units.length - 1) { - next /= 1024; - unit += 1; - } - return `${next >= 10 ? next.toFixed(0) : next.toFixed(1)} ${units[unit]}`; -} - -function formatDuration(startedAt: number, endedAt: number): string { - const seconds = Math.max(0, Math.round((endedAt - startedAt) / 1000)); - if (seconds < 90) return `${seconds}s`; - const minutes = Math.round(seconds / 60); - if (minutes < 90) return `${minutes}m`; - return `${(minutes / 60).toFixed(1)}h`; -} - -function formatTime(ts: number): string { - return new Date(ts).toLocaleTimeString([], { hour: '2-digit', minute: '2-digit' }); -} - -function toolSpanLabel(span: WorkspaceResourceToolSpan): string { - const toolName = span.toolName?.trim(); - if (toolName) return toolName; - const kind = span.kind?.trim(); - if (!kind || kind === 'acp_tool_call') return 'tool'; - return kind.replaceAll('_', ' '); -} - -function sampleMemoryMiB(sample: WorkspaceResourceSample): number { - return Number(sample.memoryBytes ?? 0) / (1024 * 1024); -} - -function sampleWorkingSetMiB(sample: WorkspaceResourceSample): number | null { - if (sample.memoryWorkingSetBytes == null) return null; - return Number(sample.memoryWorkingSetBytes) / (1024 * 1024); -} - -function sampleCpuMillis(sample: WorkspaceResourceSample): number { - return Number(sample.cpuMillis ?? 0); -} - -function clampPercent(value: number): number { - return Math.max(0, Math.min(100, value)); -} - -function sampleX(samples: WorkspaceResourceSample[], sample: WorkspaceResourceSample): number { - if (samples.length <= 1) return 0; - const start = samples[0]?.t; - const end = samples.at(-1)?.t; - if (typeof start !== 'number' || typeof end !== 'number') return 0; - if (!Number.isFinite(start) || !Number.isFinite(end) || end <= start) return 0; - return clampPercent(((sample.t - start) / (end - start)) * 100); -} - -function seriesPoints( - samples: WorkspaceResourceSample[], - valueForSample: (sample: WorkspaceResourceSample) => number | null, - scaleMax?: number -): string { - const values = samples.map(valueForSample).filter((value): value is number => value != null); - const max = Math.max(1, scaleMax ?? 0, ...values); - return samples - .filter((sample) => valueForSample(sample) != null) - .map((sample) => { - const x = sampleX(samples, sample); - const value = valueForSample(sample) ?? 0; - const y = 100 - (value / max) * 84 - 8; - return `${x.toFixed(2)},${y.toFixed(2)}`; - }) - .join(' '); -} - -function seriesSegments( - samples: WorkspaceResourceSample[], - valueForSample: (sample: WorkspaceResourceSample) => number | null, - scaleMax?: number -): string[] { - const values = samples.map(valueForSample).filter((value): value is number => value != null); - const max = Math.max(1, scaleMax ?? 0, ...values); - const segments: string[][] = []; - let current: string[] = []; - - for (const sample of samples) { - const value = valueForSample(sample); - if (value == null) { - if (current.length > 0) segments.push(current); - current = []; - continue; - } - const x = sampleX(samples, sample); - const y = 100 - (value / max) * 84 - 8; - current.push(`${x.toFixed(2)},${y.toFixed(2)}`); - } - if (current.length > 0) segments.push(current); - return segments.map((segment) => segment.join(' ')); -} - -export function ResourceSparkline({ - samples, - toolSpans, -}: Readonly<{ samples: WorkspaceResourceSample[]; toolSpans: WorkspaceResourceToolSpan[] }>) { - const cpuPoints = useMemo(() => seriesPoints(samples, sampleCpuMillis), [samples]); - const memoryScaleMax = useMemo(() => Math.max(1, ...samples.map(sampleMemoryMiB)), [samples]); - const totalMemoryPoints = useMemo( - () => seriesPoints(samples, sampleMemoryMiB, memoryScaleMax), - [memoryScaleMax, samples] - ); - const workingSetSegments = useMemo( - () => seriesSegments(samples, sampleWorkingSetMiB, memoryScaleMax), - [memoryScaleMax, samples] - ); - const eventMarkers = useMemo( - () => - samples - .map((sample) => ({ sample, x: sampleX(samples, sample) })) - .filter(({ sample }) => sample.gap || sample.oom || sample.oomKill || sample.counterReset), - [samples] - ); - const chunkIoRead = samples.reduce((total, sample) => total + Number(sample.ioReadBytes ?? 0), 0); - const chunkIoWrite = samples.reduce( - (total, sample) => total + Number(sample.ioWriteBytes ?? 0), - 0 - ); - - if (samples.length === 0) { - return ( -
- No samples in this chunk. -
- ); - } - - const start = samples[0]?.t ?? 0; - const end = samples.at(-1)?.t ?? start; - const spanWidth = Math.max(1, end - start); - - return ( -
- {/* - `preserveAspectRatio="none"`: without it the default `xMidYMid meet` - letterboxes the 1:1 viewBox into a centred SQUARE — measured 224px of - drawn width inside a 323px card (144px at this `h-36`), so a third of - the chart area was empty and the timeline was compressed to a third of - its width. Every polyline already carries - `vectorEffect="non-scaling-stroke"`, which only matters under - non-uniform scaling: stretching was always the intent. - */} - - {toolSpans.map((span) => { - const x = ((span.startedAt - start) / spanWidth) * 100; - const width = - ((Math.max(span.endedAt ?? span.startedAt, span.startedAt + 1) - span.startedAt) / - spanWidth) * - 100; - return ( - - ); - })} - - - {workingSetSegments.map((points, index) => ( - - ))} - {eventMarkers.map(({ sample, x }) => ( - - - - - ))} - -
- CPU: green solid line, normalized to the CPU peak for this chunk. - Memory needed: purple solid line (working set, when reported). - Total RAM: faint purple dashed line, including reclaimable file cache. - - Blue bands: concurrent tool windows. Dashed markers: gaps, counter resets, or OOM samples. - - - Chunk I/O deltas: {formatBytes(chunkIoRead)} read · {formatBytes(chunkIoWrite)} write. - -
- {toolSpans.length > 0 && ( -
-
- Tool windows -
-
    - {toolSpans.map((span) => { - const label = toolSpanLabel(span); - return ( -
  • - - {label} · {formatTime(span.startedAt)} - {span.approximate ? ' · approximate end' : ''} - - - {formatDuration(span.startedAt, span.endedAt ?? span.startedAt)} - {span.concurrency ? ` · ${span.concurrency} concurrent` : ''} - -
  • - ); - })} -
-
- )} -
- ); -} - -export function StatCard({ - icon: Icon, - label, - value, -}: Readonly<{ icon: typeof Cpu; label: string; value: string }>) { - return ( -
-
- - {label} -
-
{value}
-
- ); -} - -function formatCpuPeak(summary: WorkspaceResourceSummary): string { - if (summary.cpuPeakMillis == null) return '—'; - return `${Math.round(summary.cpuPeakMillis)} ms/sample`; -} - -function detailPointLabel(detail: WorkspaceResourceHistoryResponse['detail']): string { - if (!detail) return ''; - if (detail.downsampled) return `${detail.samples.length}/${detail.originalSampleCount} points`; - return `${detail.samples.length} points`; -} - -export function ResourceHistoryContent({ - isLoading, - isError, - isFetching, - history, - summary, - effectiveChunkId, - detail, - onSelectChunk, -}: Readonly<{ - isLoading: boolean; - isError: boolean; - isFetching: boolean; - history: WorkspaceResourceHistoryResponse | undefined; - summary: WorkspaceResourceSummary | null | undefined; - effectiveChunkId: string | null; - detail: WorkspaceResourceHistoryResponse['detail']; - onSelectChunk: (chunkId: string) => void; -}>) { - if (isLoading && !history) { - return ( -
- -
- ); - } - - if (isError) { - return ( -
- Resource history could not be loaded. -
- ); - } - - const chunks = history?.chunks ?? []; - if (!summary && chunks.length === 0 && !detail) { - return ( -
- No retained resource history is available for this session yet. -
- ); - } - - return ( - <> - {/* 1. Stat cards */} - {summary && ( -
- - - - - - -
- )} - - {/* 2. OOM banner */} - {summary && summary.oomCount > 0 && ( -
- - - {summary.oomCount} OOM event{summary.oomCount === 1 ? '' : 's'} observed in retained - samples. - -
- )} - - {/* 3. Chart (auto-loaded for newest chunk) */} - {detail && ( -
-
-

- Detail timeline -

- {detailPointLabel(detail)} -
- -
- )} - - {effectiveChunkId && !detail && isFetching && ( -
- -
- )} - - {/* 4. Correlation disclaimer (contextual, after chart) */} - {detail && ( -
- Correlation is based on concurrent tool windows and background resource usage. It is not - per-process causal attribution. Disk space is not sampled on the hot loop. -
- )} - - {/* 5. Chunks disclosure (collapsed by default) */} - {chunks.length > 0 && ( - - )} - - ); -} - -export function ChunkButton({ - chunk, - selected, - onSelect, -}: Readonly<{ chunk: WorkspaceResourceChunk; selected: boolean; onSelect: () => void }>) { - return ( - - ); -} - -function ChunksDisclosure({ - chunks, - effectiveChunkId, - onSelectChunk, -}: Readonly<{ - chunks: WorkspaceResourceChunk[]; - effectiveChunkId: string | null; - onSelectChunk: (chunkId: string) => void; -}>) { - const [open, setOpen] = useState(false); - return ( -
- - {open && ( -
- {chunks.map((chunk) => ( - onSelectChunk(chunk.id)} - /> - ))} -
- )} -
- ); -} - -export function SessionResourceHistoryDrawer({ - projectId, - sessionId, - onClose, -}: SessionResourceHistoryDrawerProps) { +export function SessionResourceHistoryDrawer({ source, onClose }: SessionResourceHistoryDrawerProps) { const panelRef = useRef(null); useDialogFocusTrap(panelRef, onClose); - const [selectedChunkId, setSelectedChunkId] = useState(null); - - const query = useQuery({ - queryKey: ['session-resource-history', projectId, sessionId, selectedChunkId], - queryFn: () => getSessionResourceHistory(projectId, sessionId, { chunkId: selectedChunkId }), - staleTime: 15_000, - placeholderData: keepPreviousData, - }); - - const history = query.data; - const summary = history?.summary; - const latestChunk = history?.chunks[0] ?? null; - const effectiveChunkId = selectedChunkId ?? latestChunk?.id ?? null; - const detail = query.isPlaceholderData ? undefined : history?.detail; - - useEffect(() => { - setSelectedChunkId(null); - }, [projectId, sessionId]); - - useEffect(() => { - if (selectedChunkId === null && latestChunk?.id && !query.isPlaceholderData) { - setSelectedChunkId(latestChunk.id); - } - }, [selectedChunkId, latestChunk?.id, query.isPlaceholderData]); return createPortal( <> @@ -586,17 +54,16 @@ export function SessionResourceHistoryDrawer({ -
- +
+ + +
+ } + > + +
, diff --git a/apps/web/src/components/chat/resource-timeline/ResourceTimeline.tsx b/apps/web/src/components/chat/resource-timeline/ResourceTimeline.tsx new file mode 100644 index 000000000..6ca4a6ccb --- /dev/null +++ b/apps/web/src/components/chat/resource-timeline/ResourceTimeline.tsx @@ -0,0 +1,272 @@ +import { Button, Spinner } from '@simple-agent-manager/ui'; +import { AlertTriangle } from 'lucide-react'; +import { type ReactNode, useCallback, useMemo, useRef, useState } from 'react'; + +import { formatBytes, formatDayTime, formatElapsed, formatMinutes } from './format'; +import { type PanelState, PLOT_LEFT_GUTTER_PX, PLOT_RIGHT_PADDING_PX } from './panels'; +import { readoutAtCursor, readoutForRange } from './readout'; +import type { ResourceHistorySource } from './resource-source'; +import { buildSeries, findPeaks, type UsagePeak } from './series'; +import { sleeps, type TimeAxisMode, toAxis, toReal } from './time-axis'; +import { useTimelineView } from './timeline-view'; +import { PEAK_CONTEXT_MS, PeakList, RangeControls, ReadoutBar } from './TimelineControls'; +import { TimelineNavigator } from './TimelineNavigator'; +import { TimelinePanels } from './TimelinePanels'; +import type { ResourceTimelineIndex } from './types'; +import { useElementWidth } from './useElementWidth'; +import { + useResourceTimelineData, + useResourceTimelineIndex, + useTimelineAxis, +} from './useResourceTimelineData'; + +/** Peaks closer together than this count as one moment. */ +const PEAK_SEPARATION_MS = 20 * 60_000; +const PEAKS_PER_METRIC = 3; +/** "Latest" shows this much of the newest data. */ +const LATEST_SPAN_MS = 15 * 60_000; + +/** + * The whole-session resource timeline. The drawer loads this module lazily, so + * the project chat route never downloads uPlot unless someone opens Resources. + */ +export function ResourceTimeline({ source }: Readonly<{ source: ResourceHistorySource }>) { + const indexQuery = useResourceTimelineIndex(source); + + if (indexQuery.isPending) { + return ( +
+ +
+ ); + } + if (!indexQuery.data) { + return ( +
+

Resource history could not be loaded.

+ +
+ ); + } + if (indexQuery.data.chunks.length === 0) { + switch (indexQuery.data.collection) { + case 'unsupported': + return ( + + Instant sessions run in a lightweight container that does not record CPU, memory or disk + usage yet. Sessions on a VM workspace record the full timeline. + + ); + case 'expired': + return ( + + This session's CPU, memory and disk samples are older than the retention period and + have been deleted. + + ); + default: + return ( + + The workspace samples CPU, memory and disk every few seconds and uploads them every{' '} + {formatMinutes(indexQuery.data.uploadIntervalMs)}, so the first data appears about{' '} + {formatMinutes(indexQuery.data.uploadIntervalMs)} after the session starts. + + ); + } + } + return ; +} + +function EmptyNotice({ title, children }: Readonly<{ title: string; children: ReactNode }>) { + return ( +
+

{title}

+

{children}

+
+ ); +} + +function SessionSummary({ + index, + activeMs, +}: Readonly<{ index: ResourceTimelineIndex; activeMs: number }>) { + const first = index.runs[0]?.startedAt ?? index.chunks[0]?.startedAt ?? 0; + const last = index.chunks.at(-1)?.endedAt ?? first; + const nodes = new Set(index.runs.map((run) => run.nodeId).filter(Boolean)).size; + const highWater = Math.max(0, ...index.chunks.map((chunk) => chunk.memoryHighWaterBytes ?? 0)); + const wakes = index.runs.length; + return ( +
+

+ {formatElapsed(activeMs)} active + {wakes > 1 && ` over ${formatElapsed(last - first)} · ${wakes} wake cycles`} + {nodes > 1 && ` on ${nodes} nodes`} +

+

+ Data until {formatDayTime(last)} · uploaded every {formatMinutes(index.uploadIntervalMs)} + {highWater > 0 && ` · kernel memory peak ${formatBytes(highWater)}`} +

+ {index.completeness.kind === 'truncated' && ( +

+

+ )} +
+ ); +} + +function TimelineBody({ + source, + index, +}: Readonly<{ source: ResourceHistorySource; index: ResourceTimelineIndex }>) { + const [axisMode, setAxisMode] = useState('active'); + const axis = useTimelineAxis(index, axisMode); + const { view, intent, setRange, zoom, showSpan, showLatest, showAll } = useTimelineView(axis); + /** The instant under the finger or pointer, as wall-clock time so it survives axis changes. */ + const [cursorT, setCursorT] = useState(null); + + const measureRef = useRef(null); + const width = useElementWidth(measureRef); + const plotWidth = Math.max(1, width - PLOT_LEFT_GUTTER_PX - PLOT_RIGHT_PADDING_PX); + + const data = useResourceTimelineData(source, index, axis, { + min: view.min, + max: view.max, + widthPx: plotWidth, + }); + const series = useMemo( + () => buildSeries(data.aggregates, axis, view.min, view.max, plotWidth, index.sampleIntervalMs), + [data.aggregates, axis, view.min, view.max, plotWidth, index.sampleIntervalMs] + ); + const hasWorkingSet = useMemo( + () => data.overview.some((aggregate) => aggregate.workingSetMaxBytes != null), + [data.overview] + ); + const panelState = useMemo( + () => ({ + axis, + viewMin: view.min, + viewMax: view.max, + series, + runs: index.runs, + toolSpans: data.toolSpans, + hasWorkingSet, + }), + [axis, view.min, view.max, series, index.runs, data.toolSpans, hasWorkingSet] + ); + + const cursorX = cursorT == null ? null : toAxis(axis, cursorT); + const cursorInView = cursorX != null && cursorX >= view.min && cursorX <= view.max; + const readout = cursorInView + ? readoutAtCursor(cursorX, series, axis, index.runs, data.toolSpans, index.sampleIntervalMs) + : readoutForRange( + view.min, + view.max, + axis, + data.aggregates, + index.completeness.kind === 'truncated' + ); + + const peaks = useMemo( + () => + [ + ...findPeaks(data.overview, 'memory', PEAKS_PER_METRIC, PEAK_SEPARATION_MS), + ...findPeaks(data.overview, 'cpu', PEAKS_PER_METRIC, PEAK_SEPARATION_MS), + ].sort((a, b) => a.at - b.at), + [data.overview] + ); + + const onCursor = useCallback( + (x: number | null) => setCursorT(x == null ? null : toReal(axis, x)), + [axis] + ); + const onAxisMode = useCallback((mode: TimeAxisMode) => setAxisMode(mode), []); + const centre = cursorInView ? cursorX : (view.min + view.max) / 2; + const runLabel = useCallback( + (t: number) => { + const position = index.runs.findIndex((run) => t >= run.startedAt && t <= run.endedAt); + return position === -1 ? 'between runs' : `run ${position + 1} of ${index.runs.length}`; + }, + [index.runs] + ); + const onPeak = useCallback( + (peak: UsagePeak) => { + setCursorT(peak.at); + showSpan(PEAK_CONTEXT_MS, toAxis(axis, peak.at)); + }, + [axis, showSpan] + ); + + return ( +
+ + setCursorT(null)} + /> +
+ + +
+ 0} + onAll={showAll} + onSpan={(spanMs) => showSpan(spanMs, centre)} + onLatest={() => showLatest(LATEST_SPAN_MS)} + onAxisMode={onAxisMode} + /> + +
+ About this data +
    +
  • + Sampled every {formatElapsed(index.sampleIntervalMs)} from the workspace + container's cgroup. It covers everything in the container, not one process. +
  • +
  • CPU is in cores: 1.0 means one core fully busy.
  • +
  • + Memory “used” is the working set the kernel cannot reclaim; “incl. + cache” adds page cache it can drop under pressure. Size workspaces by + “used”. +
  • +
  • + Tool calls line up with usage by time only. A spike during a tool call is a correlation, + not proof that the tool caused it. +
  • +
  • + Zoomed out, each point is an average with its peak shaded; zoom in for individual + samples. +
  • +
+
+
+ ); +} diff --git a/apps/web/src/components/chat/resource-timeline/TimelineControls.tsx b/apps/web/src/components/chat/resource-timeline/TimelineControls.tsx new file mode 100644 index 000000000..387640835 --- /dev/null +++ b/apps/web/src/components/chat/resource-timeline/TimelineControls.tsx @@ -0,0 +1,195 @@ +import { AlertTriangle, ChevronsRight, Loader2, X } from 'lucide-react'; + +import { formatBytes, formatCores, formatDayTime, formatElapsed } from './format'; +import type { Readout } from './readout'; +import type { UsagePeak } from './series'; +import type { TimeAxisMode } from './time-axis'; + +const PRESET_SPANS_MS = [ + { label: '5m', ms: 5 * 60_000 }, + { label: '15m', ms: 15 * 60_000 }, + { label: '1h', ms: 60 * 60_000 }, + { label: '3h', ms: 3 * 60 * 60_000 }, +]; + +/** Tapping a peak shows this much time around it. */ +export const PEAK_CONTEXT_MS = 10 * 60_000; + +const chip = + 'rounded-full border px-2.5 py-1 text-xs transition-colors focus-visible:outline-2 focus-visible:outline-focus-ring'; +const chipIdle = `${chip} border-border-default text-fg-muted hover:bg-surface-hover hover:text-fg-primary`; +const chipActive = `${chip} border-accent bg-accent/15 text-fg-primary`; + +export function ReadoutBar({ + readout, + pendingChunks, + failedChunks, + onClear, +}: Readonly<{ + readout: Readout; + pendingChunks: number; + failedChunks: number; + onClear: () => void; +}>) { + return ( +
+
+
{readout.time}
+
{readout.context}
+ {readout.oomKills > 0 && ( +
+
+ )} +
+
+ {pendingChunks > 0 && ( + + + )} + {failedChunks > 0 && pendingChunks === 0 && ( + + detail unavailable + + )} + {readout.mode === 'cursor' && ( + + )} +
+
+ ); +} + +export function RangeControls({ + fullSpanMs, + viewSpanMs, + followingLatest, + axisMode, + hasSleeps, + onAll, + onSpan, + onLatest, + onAxisMode, +}: Readonly<{ + fullSpanMs: number; + viewSpanMs: number; + followingLatest: boolean; + axisMode: TimeAxisMode; + hasSleeps: boolean; + onAll: () => void; + onSpan: (spanMs: number) => void; + onLatest: () => void; + onAxisMode: (mode: TimeAxisMode) => void; +}>) { + const showingAll = viewSpanMs >= fullSpanMs * 0.999; + const presets = PRESET_SPANS_MS.filter((preset) => preset.ms < fullSpanMs * 0.9); + return ( +
+ + {presets.map((preset) => { + const active = !showingAll && Math.abs(viewSpanMs - preset.ms) < preset.ms * 0.02; + return ( + + ); + })} + + {hasSleeps && ( +
+ {(['active', 'clock'] as const).map((mode) => ( + + ))} +
+ )} +
+ ); +} + +export function PeakList({ + peaks, + runLabel, + onSelect, +}: Readonly<{ + peaks: UsagePeak[]; + runLabel: (t: number) => string; + onSelect: (peak: UsagePeak) => void; +}>) { + if (peaks.length === 0) return null; + return ( +
+

+ Busiest moments +

+
    + {peaks.map((peak) => ( +
  • + +
  • + ))} +
+
+ ); +} diff --git a/apps/web/src/components/chat/resource-timeline/TimelineNavigator.tsx b/apps/web/src/components/chat/resource-timeline/TimelineNavigator.tsx new file mode 100644 index 000000000..73123bf81 --- /dev/null +++ b/apps/web/src/components/chat/resource-timeline/TimelineNavigator.tsx @@ -0,0 +1,245 @@ +import { type PointerEvent, useEffect, useLayoutEffect, useMemo, useRef } from 'react'; +import uPlot from 'uplot'; + +import { type ChartTheme, useChartTheme, withAlpha } from './chart-theme'; +import { formatCompactDuration, formatDayTime } from './format'; +import { createFrameCoalescer } from './frame-coalescer'; +import { PLOT_LEFT_GUTTER_PX, PLOT_RIGHT_PADDING_PX } from './panels'; +import { buildSeries } from './series'; +import { sleeps, type TimeAxis } from './time-axis'; +import type { ViewRange } from './timeline-view'; +import type { ResourceAggregate } from './types'; +import { useElementWidth } from './useElementWidth'; +import { useUplot } from './useUplot'; + +const HEIGHT = 40; +/** Sleeps at least this long get their length written on the strip. */ +const LABELLED_SLEEP_MS = 60 * 60_000; +/** Grips are thin to look at but wide to touch. */ +const GRIP_HIT_PX = 44; + +type Drag = + | { kind: 'move'; startX: number; view: ViewRange } + | { kind: 'min' | 'max'; view: ViewRange } + | { kind: 'draw'; anchor: number }; + +interface TimelineNavigatorProps { + axis: TimeAxis; + overview: readonly ResourceAggregate[]; + view: ViewRange; + sessionStart: number | null; + sessionEnd: number | null; + onRange: (min: number, max: number) => void; +} + +function navigatorOptions( + theme: ChartTheme, + axisRef: { current: TimeAxis } +): Omit { + return { + height: HEIGHT, + legend: { show: false }, + cursor: { show: false }, + select: { show: false, left: 0, top: 0, width: 0, height: 0 }, + padding: [2, PLOT_RIGHT_PADDING_PX, 2, 0], + scales: { + x: { time: false, range: () => [axisRef.current.min, axisRef.current.max] }, + cpu: { range: (_u, _min, max) => [0, Math.max(0.5, max ?? 0) * 1.1] }, + memory: { range: (_u, _min, max) => [0, Math.max(1, max ?? 0) * 1.1] }, + }, + axes: [ + { show: false }, + { + scale: 'cpu', + size: PLOT_LEFT_GUTTER_PX, + show: true, + values: () => [], + grid: { show: false }, + ticks: { show: false }, + }, + ], + series: [ + {}, + { + scale: 'cpu', + stroke: theme.cpu, + width: 1, + fill: withAlpha(theme.cpu, 0.22), + points: { show: false }, + }, + { scale: 'memory', stroke: withAlpha(theme.memory, 0.8), width: 1, points: { show: false } }, + ], + hooks: { + drawClear: [ + (u) => { + const { ctx, bbox } = u; + const px = uPlot.pxRatio; + ctx.save(); + ctx.font = `${10 * px}px system-ui, sans-serif`; + ctx.textAlign = 'center'; + ctx.textBaseline = 'top'; + for (const sleep of sleeps(axisRef.current)) { + const x0 = u.valToPos(sleep.axisStart, 'x', true); + const x1 = u.valToPos(sleep.axisEnd, 'x', true); + ctx.fillStyle = theme.sleep; + ctx.fillRect(x0, bbox.top, Math.max(px, x1 - x0), bbox.height); + const label = formatCompactDuration(sleep.realEnd - sleep.realStart); + if ( + sleep.realEnd - sleep.realStart >= LABELLED_SLEEP_MS && + x1 - x0 >= ctx.measureText(label).width + 4 * px + ) { + ctx.fillStyle = theme.mutedText; + ctx.fillText(label, (x0 + x1) / 2, bbox.top + 2 * px); + } + } + ctx.restore(); + }, + ], + }, + }; +} + +export function TimelineNavigator({ + axis, + overview, + view, + sessionStart, + sessionEnd, + onRange, +}: Readonly) { + const theme = useChartTheme(); + const containerRef = useRef(null); + const trackRef = useRef(null); + const axisRef = useRef(axis); + const trackWidth = useElementWidth(trackRef); + const drag = useRef(null); + // Drags report every pointer event; apply at most one view change per frame. + const onRangeRef = useRef(onRange); + const rangeFrame = useMemo( + () => createFrameCoalescer((min: number, max: number) => onRangeRef.current(min, max)), + [] + ); + useLayoutEffect(() => { + onRangeRef.current = onRange; + }, [onRange]); + useEffect(() => () => rangeFrame.cancel(), [rangeFrame]); + + useLayoutEffect(() => { + axisRef.current = axis; + }, [axis]); + + const options = useMemo(() => navigatorOptions(theme, axisRef), [theme]); + const data = useMemo((): uPlot.AlignedData => { + const series = buildSeries(overview, axis, axis.min, axis.max, Math.max(1, trackWidth), 0); + return [series.x, series.cpuMax, series.memoryMax]; + }, [overview, axis, trackWidth]); + useUplot(containerRef, options, data); + + const span = Math.max(1, axis.max - axis.min); + const toPx = (x: number) => ((x - axis.min) / span) * trackWidth; + const toAxisX = (clientX: number) => { + const rect = trackRef.current?.getBoundingClientRect(); + if (!rect) return axis.min; + return axis.min + ((clientX - rect.left) / Math.max(1, rect.width)) * span; + }; + const left = toPx(view.min); + const width = Math.max(2, toPx(view.max) - left); + + const onPointerDown = (event: PointerEvent) => { + const x = event.clientX - (trackRef.current?.getBoundingClientRect().left ?? 0); + const nearMin = Math.abs(x - left) <= GRIP_HIT_PX / 2; + const nearMax = Math.abs(x - (left + width)) <= GRIP_HIT_PX / 2; + if (nearMin || nearMax) { + // Prefer the grip the pointer is closer to when the window is narrow. + drag.current = { + kind: Math.abs(x - left) < Math.abs(x - left - width) ? 'min' : 'max', + view, + }; + } else if (x > left && x < left + width) { + drag.current = { kind: 'move', startX: event.clientX, view }; + } else { + drag.current = { kind: 'draw', anchor: toAxisX(event.clientX) }; + } + event.currentTarget.setPointerCapture(event.pointerId); + }; + + const onPointerMove = (event: PointerEvent) => { + const current = drag.current; + if (!current) return; + const at = toAxisX(event.clientX); + switch (current.kind) { + case 'move': { + const shift = ((event.clientX - current.startX) / Math.max(1, trackWidth)) * span; + rangeFrame.schedule(current.view.min + shift, current.view.max + shift); + break; + } + case 'min': + rangeFrame.schedule(Math.min(at, current.view.max - 1), current.view.max); + break; + case 'max': + rangeFrame.schedule(current.view.min, Math.max(at, current.view.min + 1)); + break; + case 'draw': + if (Math.abs(toPx(at) - toPx(current.anchor)) > 6) + rangeFrame.schedule(Math.min(at, current.anchor), Math.max(at, current.anchor)); + break; + } + }; + + const onPointerUp = (event: PointerEvent) => { + const current = drag.current; + drag.current = null; + // A tap outside the window recentres it there. + if ( + current?.kind === 'draw' && + Math.abs(toPx(toAxisX(event.clientX)) - toPx(current.anchor)) <= 6 + ) { + const half = (view.max - view.min) / 2; + onRange(current.anchor - half, current.anchor + half); + } + }; + + return ( +
+
+