diff --git a/.claude/skills/api-reference/SKILL.md b/.claude/skills/api-reference/SKILL.md index 76d45d230..9d6fdece7 100644 --- a/.claude/skills/api-reference/SKILL.md +++ b/.claude/skills/api-reference/SKILL.md @@ -78,6 +78,9 @@ user-invocable: false - `GET /api/admin/observability/errors` — Query platform errors; VM error rows include their same-installation diagnostic incident summary - `GET /api/admin/observability/errors/:errorId/incident` — Read one diagnostic incident summary and redacted preview - `GET /api/admin/observability/errors/:errorId/incident/artifacts/:artifactId/download` — Stream one private diagnostic artifact through the authenticated Worker; R2 keys and URLs are never exposed +- `GET /api/admin/project-data/storage` — List latest per-project ProjectData storage telemetry from D1 (`projectId`, `status`, `limit` filters) +- `POST /api/admin/project-data/storage/:projectId/measure` — Force one ProjectData `databaseSize` measurement and D1 telemetry upsert +- `POST /api/admin/project-data/storage/:projectId/emergency-purge` — Run a bounded ProjectData emergency purge of oldest `activity_events` and `acp_session_events` rows only ## Agent Sessions diff --git a/.claude/skills/env-reference/SKILL.md b/.claude/skills/env-reference/SKILL.md index a59b656a5..5c8e62a96 100644 --- a/.claude/skills/env-reference/SKILL.md +++ b/.claude/skills/env-reference/SKILL.md @@ -157,6 +157,18 @@ See `apps/api/.env.example` for the full list. Key variables: - `ORCHESTRATOR_WAIT_MAX_ACTIVE_PER_PROJECT` — Maximum active durable parent waits per project (default: `100`) - `ORCHESTRATOR_WAIT_MAX_DURATION_MS` — Maximum finite durable wait deadline (default: `86400000`) - `ORCHESTRATOR_WAIT_MAX_CANDIDATES_PER_ALARM` — Maximum wait subscriptions reconciled by one ProjectData alarm (default: `10`) +- `PROJECT_DATA_TOOL_METADATA_MAX_BYTES` — Maximum stored `tool_metadata` bytes per ProjectData message before oversized tool content is stripped into bounded metadata (default: `131072`) +- `PROJECT_DATA_STORAGE_TELEMETRY_ENABLED` — Enables ProjectData `databaseSize` alarm measurement and D1 telemetry writes (default: `true`) +- `PROJECT_DATA_STORAGE_LIMIT_BYTES` — Cloudflare SQLite-backed Durable Object storage limit used for ProjectData usage classification (default: `10000000000`) +- `PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS` — Minimum interval between per-object ProjectData storage measurements (default: `3600000`) +- `PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS` — Minimum interval between repeated critical/degraded ProjectData storage observability alerts (default: `21600000`) +- `PROJECT_DATA_STORAGE_NOTICE_RATIO` — ProjectData storage usage ratio classified as `notice` (default: `0.6`) +- `PROJECT_DATA_STORAGE_WARNING_RATIO` — ProjectData storage usage ratio classified as `warning` (default: `0.8`) +- `PROJECT_DATA_STORAGE_CRITICAL_RATIO` — ProjectData storage usage ratio classified as `critical` (default: `0.9`) +- `PROJECT_DATA_STORAGE_DEGRADED_RATIO` — ProjectData storage usage ratio classified as `degraded` (default: `0.95`) +- `PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO` — Target usage ratio for explicit superadmin ProjectData emergency purge calls (default: `0.9`) +- `PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS` — Oldest `activity_events` and `acp_session_events` rows deleted per table per emergency purge batch (default: `500`) +- `PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES` — Maximum emergency purge batches per explicit call (default: `4`) Absent operational brake keys and KV read errors mean enabled. This fail-open behavior preserves availability and intentionally differs from the fail-closed diff --git a/.env.example b/.env.example index ff25f548f..2caded198 100644 --- a/.env.example +++ b/.env.example @@ -68,6 +68,20 @@ VM_INCIDENT_PENDING_TIMEOUT_MINUTES=30 # Minimum: 6 (one row per reconciliation phase). VM_INCIDENT_RECONCILE_BATCH_SIZE=50 +# ProjectData Durable Object storage safety (Worker defaults) +PROJECT_DATA_TOOL_METADATA_MAX_BYTES=131072 +PROJECT_DATA_STORAGE_TELEMETRY_ENABLED=true +PROJECT_DATA_STORAGE_LIMIT_BYTES=10000000000 +PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS=3600000 +PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS=21600000 +PROJECT_DATA_STORAGE_NOTICE_RATIO=0.6 +PROJECT_DATA_STORAGE_WARNING_RATIO=0.8 +PROJECT_DATA_STORAGE_CRITICAL_RATIO=0.9 +PROJECT_DATA_STORAGE_DEGRADED_RATIO=0.95 +PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO=0.9 +PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS=500 +PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES=4 + # Development NODE_ENV=development # Generic webhook triggers (all values optional; defaults shown) diff --git a/apps/api/src/db/migrations/0119_project_data_storage_telemetry.sql b/apps/api/src/db/migrations/0119_project_data_storage_telemetry.sql new file mode 100644 index 000000000..646c4667a --- /dev/null +++ b/apps/api/src/db/migrations/0119_project_data_storage_telemetry.sql @@ -0,0 +1,33 @@ +-- ProjectData Durable Object storage telemetry. +-- +-- One row per project records the latest direct per-object SQLite +-- `databaseSize` measurement from the ProjectData Durable Object. This is an +-- additive D1 index only; it does not mutate ProjectData object schemas. + +CREATE TABLE project_data_storage_telemetry ( + project_id TEXT PRIMARY KEY REFERENCES projects(id) ON DELETE CASCADE, + measured_at INTEGER NOT NULL CHECK (measured_at > 0), + database_size_bytes INTEGER NOT NULL CHECK (database_size_bytes >= 0), + limit_bytes INTEGER NOT NULL CHECK (limit_bytes > 0), + usage_ratio REAL NOT NULL CHECK (usage_ratio >= 0), + status TEXT NOT NULL CHECK (status IN ('ok', 'notice', 'warning', 'critical', 'degraded')), + last_alarm_at INTEGER, + last_alert_at INTEGER, + last_alert_status TEXT CHECK ( + last_alert_status IS NULL + OR last_alert_status IN ('ok', 'notice', 'warning', 'critical', 'degraded') + ), + last_purge_at INTEGER, + last_purge_reason TEXT, + last_purge_rows INTEGER, + last_purge_database_size_bytes INTEGER, + last_error TEXT, + created_at INTEGER NOT NULL DEFAULT (cast(unixepoch() * 1000 as integer)), + updated_at INTEGER NOT NULL DEFAULT (cast(unixepoch() * 1000 as integer)) +); + +CREATE INDEX idx_project_data_storage_telemetry_status + ON project_data_storage_telemetry(status, usage_ratio DESC, measured_at DESC); + +CREATE INDEX idx_project_data_storage_telemetry_measured_at + ON project_data_storage_telemetry(measured_at DESC); diff --git a/apps/api/src/db/schema.ts b/apps/api/src/db/schema.ts index a8c4a50fc..a80a04e98 100644 --- a/apps/api/src/db/schema.ts +++ b/apps/api/src/db/schema.ts @@ -2398,6 +2398,44 @@ export const sessionIndexCoverage = sqliteTable('session_index_coverage', { export type SessionIndexCoverageRow = typeof sessionIndexCoverage.$inferSelect; +export const projectDataStorageTelemetry = sqliteTable( + 'project_data_storage_telemetry', + { + projectId: text('project_id') + .primaryKey() + .references(() => projects.id, { onDelete: 'cascade' }), + measuredAt: integer('measured_at').notNull(), + databaseSizeBytes: integer('database_size_bytes').notNull(), + limitBytes: integer('limit_bytes').notNull(), + usageRatio: real('usage_ratio').notNull(), + status: text('status', { enum: ['ok', 'notice', 'warning', 'critical', 'degraded'] }) + .notNull(), + lastAlarmAt: integer('last_alarm_at'), + lastAlertAt: integer('last_alert_at'), + lastAlertStatus: text('last_alert_status', { + enum: ['ok', 'notice', 'warning', 'critical', 'degraded'], + }), + lastPurgeAt: integer('last_purge_at'), + lastPurgeReason: text('last_purge_reason'), + lastPurgeRows: integer('last_purge_rows'), + lastPurgeDatabaseSizeBytes: integer('last_purge_database_size_bytes'), + lastError: text('last_error'), + createdAt: integer('created_at').notNull(), + updatedAt: integer('updated_at').notNull(), + }, + (table) => ({ + statusIdx: index('idx_project_data_storage_telemetry_status').on( + table.status, + table.usageRatio, + table.measuredAt + ), + measuredAtIdx: index('idx_project_data_storage_telemetry_measured_at').on(table.measuredAt), + }) +); + +export type ProjectDataStorageTelemetryRow = + typeof projectDataStorageTelemetry.$inferSelect; + export const diagnosticIncidents = sqliteTable( 'diagnostic_incidents', { diff --git a/apps/api/src/durable-objects/project-data/alarm-schedule.ts b/apps/api/src/durable-objects/project-data/alarm-schedule.ts index 39ef82b33..519d3e654 100644 --- a/apps/api/src/durable-objects/project-data/alarm-schedule.ts +++ b/apps/api/src/durable-objects/project-data/alarm-schedule.ts @@ -14,6 +14,7 @@ import * as mailbox from './mailbox'; import { computePromptDeliveryAlarmTime } from './prompt-delivery'; import * as reconciliation from './reconciliation'; import { computeSessionActivityProbeAlarmTime } from './session-activity-reconciliation'; +import { computeStorageSafetyAlarmTime } from './storage-safety'; import { computeTaskWaitAlarmTime } from './task-waits'; import type { Env } from './types'; @@ -49,6 +50,7 @@ export function computeProjectDataAlarmTime(sql: SqlStorage, env: Env): number | // relying on the ACP heartbeat alarm happening to fire on a similar cadence. const activityProbeTime = computeSessionActivityProbeAlarmTime(sql, env); const taskWaitTime = computeTaskWaitAlarmTime(sql); + const storageSafetyTime = computeStorageSafetyAlarmTime(sql, env); const candidates = [ idleCleanupTime, @@ -59,6 +61,7 @@ export function computeProjectDataAlarmTime(sql: SqlStorage, env: Env): number | reconciliationTime, activityProbeTime, taskWaitTime, + storageSafetyTime, ].filter((time): time is number => time !== null); return candidates.length > 0 ? Math.min(...candidates) : null; diff --git a/apps/api/src/durable-objects/project-data/index.ts b/apps/api/src/durable-objects/project-data/index.ts index 523529197..f98776dee 100644 --- a/apps/api/src/durable-objects/project-data/index.ts +++ b/apps/api/src/durable-objects/project-data/index.ts @@ -49,6 +49,7 @@ import * as sessionState from './session-state'; import * as sessionSummarySync from './session-summary-sync'; import * as sessionWakeProgress from './session-wake-progress'; import * as sessions from './sessions'; +import * as storageSafety from './storage-safety'; import { resolveTaskWaitConfig } from './task-wait-config'; import { processTaskWaits } from './task-wait-supervisor'; import * as taskWaits from './task-waits'; @@ -895,11 +896,48 @@ export class ProjectData extends DurableObject { }; } + async measureStorage(): Promise { + const measurement = await storageSafety.measureAndPersistProjectDataStorage( + this.sql, + this.env, + this.getProjectId(), + 'admin' + ); + await this.recalculateAlarm(); + return measurement; + } + + async runStorageEmergencyPurge( + input: storageSafety.ProjectDataStorageEmergencyPurgeInput = {} + ): Promise { + const result = await storageSafety.runProjectDataStorageEmergencyPurge( + this.sql, + this.env, + this.getProjectId(), + input + ); + await this.recalculateAlarm(); + return result; + } + // --- DO Alarm Handler --- async alarm(): Promise { if (await deferAlarmWhenDisabled(this.env, this.ctx.storage, 'ProjectData')) return; + try { + await storageSafety.measureAndPersistProjectDataStorage( + this.sql, + this.env, + this.getProjectId(), + 'alarm' + ); + } catch (err) { + log.error('alarm.storage_safety_failed', { + error: err instanceof Error ? err.message : String(err), + }); + } + const timedOut = await checkRuntimeHeartbeatTimeouts( this.sql, this.env, diff --git a/apps/api/src/durable-objects/project-data/message-persistence.ts b/apps/api/src/durable-objects/project-data/message-persistence.ts index 8545348dc..39b001ca7 100644 --- a/apps/api/src/durable-objects/project-data/message-persistence.ts +++ b/apps/api/src/durable-objects/project-data/message-persistence.ts @@ -77,7 +77,7 @@ export async function persistMessageWithSideEffects( messageId: result.id, role, content, - toolMetadata: parseToolMetadata(toolMetadata, sessionId), + toolMetadata: parseToolMetadata(result.toolMetadata, sessionId), createdAt: result.now, sequence: result.sequence, // The single-message path only persists browser/RPC user messages, which diff --git a/apps/api/src/durable-objects/project-data/messages.ts b/apps/api/src/durable-objects/project-data/messages.ts index 62dd10cea..6c61411f0 100644 --- a/apps/api/src/durable-objects/project-data/messages.ts +++ b/apps/api/src/durable-objects/project-data/messages.ts @@ -5,7 +5,6 @@ import { buildSafeFtsQuery } from '../../lib/fts5'; import { log } from '../../lib/logger'; import { type CompactMessageOptions, - DEFAULT_DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES, parseChatMessageRow, parseChatMessageRowCompact, parseCount, @@ -15,9 +14,16 @@ import { parseWorkspaceId, type SearchResultParsed, } from './row-schemas'; +import { boundToolMetadataForStorage } from './tool-metadata-storage'; import type { Env } from './types'; import { generateId } from './types'; +export { + boundToolMetadataForStorage, + DEFAULT_PROJECT_DATA_TOOL_METADATA_MAX_BYTES, + resolveCompactMessageOptions, +} from './tool-metadata-storage'; + export const DEFAULT_MAX_MESSAGES_PER_SESSION = 100000; export const SESSION_MESSAGE_LIMIT_EXCEEDED = 'SESSION_MESSAGE_LIMIT_EXCEEDED'; @@ -37,14 +43,6 @@ function resolveMaxMessagesPerSession(env: Env): number { return Number.isFinite(parsed) && parsed > 0 ? parsed : DEFAULT_MAX_MESSAGES_PER_SESSION; } -export function resolveCompactMessageOptions(env: Env): CompactMessageOptions { - const parsed = Number.parseInt(env.DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES || '', 10); - return { - documentCardRawOutputMaxBytes: - Number.isFinite(parsed) && parsed > 0 ? parsed : DEFAULT_DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES, - }; -} - /** * Returns the next monotonic sequence number for a session's messages. */ @@ -66,7 +64,14 @@ export function persistMessage( content: string, toolMetadata: string | null, messageId?: string -): { id: string; now: number; sequence: number; workspaceId: string | null; inserted: boolean } { +): { + id: string; + now: number; + sequence: number; + workspaceId: string | null; + inserted: boolean; + toolMetadata: string | null; +} { const maxMessages = resolveMaxMessagesPerSession(env); const countRow = sql .exec('SELECT message_count FROM chat_sessions WHERE id = ?', sessionId) @@ -84,7 +89,14 @@ export function persistMessage( const workspaceId = wsRow ? parseWorkspaceId(wsRow, 'messages.persist_duplicate_workspace') : null; - return { id, now: Date.now(), sequence: 0, workspaceId, inserted: false }; + return { + id, + now: Date.now(), + sequence: 0, + workspaceId, + inserted: false, + toolMetadata, + }; } if (parseMessageCount(countRow, 'messages.persist_count') >= maxMessages) { @@ -93,6 +105,15 @@ export function persistMessage( const now = Date.now(); const sequence = nextSequence(sql, sessionId); + const boundedToolMetadata = boundToolMetadataForStorage(toolMetadata, env); + if (boundedToolMetadata.truncated) { + log.warn('messages.tool_metadata_truncated_for_storage', { + sessionId, + messageId: id, + originalBytes: boundedToolMetadata.originalBytes, + storedBytes: boundedToolMetadata.storedBytes, + }); + } sql.exec( `INSERT INTO chat_messages (id, session_id, role, content, tool_metadata, created_at, sequence) @@ -101,7 +122,7 @@ export function persistMessage( sessionId, role, content, - toolMetadata, + boundedToolMetadata.value, now, sequence ); @@ -134,7 +155,14 @@ export function persistMessage( .toArray()[0]; const workspaceId = wsRow ? parseWorkspaceId(wsRow, 'messages.persist_workspace') : null; - return { id, now, sequence, workspaceId, inserted: true }; + return { + id, + now, + sequence, + workspaceId, + inserted: true, + toolMetadata: boundedToolMetadata.value, + }; } export function persistMessageBatch( @@ -251,6 +279,15 @@ export function persistMessageBatch( const createdAt = new Date(msg.timestamp).getTime() || now; const sequence = msg.sequence ?? nextSeq++; + const boundedToolMetadata = boundToolMetadataForStorage(msg.toolMetadata, env); + if (boundedToolMetadata.truncated) { + log.warn('messages.batch_tool_metadata_truncated_for_storage', { + sessionId, + messageId: msg.messageId, + originalBytes: boundedToolMetadata.originalBytes, + storedBytes: boundedToolMetadata.storedBytes, + }); + } sql.exec( `INSERT INTO chat_messages (id, session_id, role, content, tool_metadata, created_at, sequence, origin) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, @@ -258,7 +295,7 @@ export function persistMessageBatch( sessionId, msg.role, msg.content, - msg.toolMetadata, + boundedToolMetadata.value, createdAt, sequence, origin @@ -268,7 +305,7 @@ export function persistMessageBatch( id: msg.messageId, role: msg.role, content: msg.content, - toolMetadata: msg.toolMetadata ? JSON.parse(msg.toolMetadata) : null, + toolMetadata: boundedToolMetadata.value ? JSON.parse(boundedToolMetadata.value) : null, createdAt, sequence, origin, diff --git a/apps/api/src/durable-objects/project-data/storage-safety.ts b/apps/api/src/durable-objects/project-data/storage-safety.ts new file mode 100644 index 000000000..dcfad7f06 --- /dev/null +++ b/apps/api/src/durable-objects/project-data/storage-safety.ts @@ -0,0 +1,581 @@ +/** + * ProjectData storage safety firebreak. + * + * This module intentionally avoids sharding or broad data movement. It provides: + * - direct per-object `databaseSize` measurement from SQLite-backed DO storage; + * - D1 telemetry and throttled observability alerts; + * - a bounded, explicit emergency purge of low-value event logs. + */ +import { isJsonRecord } from '@simple-agent-manager/shared'; + +import { createModuleLogger, serializeError } from '../../lib/logger'; +import { persistError } from '../../services/observability'; +import type { Env } from './types'; + +const log = createModuleLogger('project_data.storage_safety'); + +export const PROJECT_DATA_STORAGE_STATUSES = [ + 'ok', + 'notice', + 'warning', + 'critical', + 'degraded', +] as const; + +export type ProjectDataStorageStatus = (typeof PROJECT_DATA_STORAGE_STATUSES)[number]; + +export const DEFAULT_PROJECT_DATA_STORAGE_LIMIT_BYTES = 10_000_000_000; +export const DEFAULT_PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS = 60 * 60 * 1000; +export const DEFAULT_PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS = 6 * 60 * 60 * 1000; +export const DEFAULT_PROJECT_DATA_STORAGE_NOTICE_RATIO = 0.6; +export const DEFAULT_PROJECT_DATA_STORAGE_WARNING_RATIO = 0.8; +export const DEFAULT_PROJECT_DATA_STORAGE_CRITICAL_RATIO = 0.9; +export const DEFAULT_PROJECT_DATA_STORAGE_DEGRADED_RATIO = 0.95; +export const DEFAULT_PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO = 0.9; +export const DEFAULT_PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS = 500; +export const DEFAULT_PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES = 4; + +const META_LAST_MEASURED_AT = 'storageSafetyLastMeasuredAt'; +const META_LAST_STATUS = 'storageSafetyLastStatus'; +const META_LAST_ALERT_AT = 'storageSafetyLastAlertAt'; +const META_LAST_ALERT_STATUS = 'storageSafetyLastAlertStatus'; +const META_LAST_ERROR = 'storageSafetyLastError'; + +export interface ProjectDataStorageTelemetry { + projectId: string; + measuredAt: number; + databaseSizeBytes: number; + limitBytes: number; + usageRatio: number; + status: ProjectDataStorageStatus; +} + +export interface ProjectDataStorageEmergencyPurgeInput { + reason?: string | null; + targetRatio?: number | null; + batchRows?: number | null; + maxBatches?: number | null; +} + +export interface ProjectDataStorageEmergencyPurgeResult { + projectId: string; + reason: string; + beforeBytes: number; + afterBytes: number; + limitBytes: number; + targetBytes: number; + statusBefore: ProjectDataStorageStatus; + statusAfter: ProjectDataStorageStatus; + batches: number; + maxBatches: number; + batchRows: number; + rowsDeleted: { + activityEvents: number; + acpSessionEvents: number; + }; + exhaustedCandidates: boolean; +} + +interface StorageSafetyConfig { + enabled: boolean; + limitBytes: number; + measureIntervalMs: number; + alertIntervalMs: number; + noticeRatio: number; + warningRatio: number; + criticalRatio: number; + degradedRatio: number; + emergencyTargetRatio: number; + emergencyBatchRows: number; + emergencyMaxBatches: number; +} + +function parsePositiveInteger(value: string | undefined, fallback: number): number { + if (!value) return fallback; + const parsed = Number.parseInt(value, 10); + return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : fallback; +} + +function parseBoundedRatio(value: string | undefined, fallback: number): number { + if (!value) return fallback; + const parsed = Number.parseFloat(value); + return Number.isFinite(parsed) && parsed > 0 && parsed < 1 ? parsed : fallback; +} + +function envFlagEnabled(value: string | undefined): boolean { + if (!value) return true; + return !['0', 'false', 'off', 'disabled'].includes(value.trim().toLowerCase()); +} + +export function resolveStorageSafetyConfig(env: Env): StorageSafetyConfig { + const noticeRatio = parseBoundedRatio( + env.PROJECT_DATA_STORAGE_NOTICE_RATIO, + DEFAULT_PROJECT_DATA_STORAGE_NOTICE_RATIO + ); + const warningRatio = parseBoundedRatio( + env.PROJECT_DATA_STORAGE_WARNING_RATIO, + DEFAULT_PROJECT_DATA_STORAGE_WARNING_RATIO + ); + const criticalRatio = parseBoundedRatio( + env.PROJECT_DATA_STORAGE_CRITICAL_RATIO, + DEFAULT_PROJECT_DATA_STORAGE_CRITICAL_RATIO + ); + const degradedRatio = parseBoundedRatio( + env.PROJECT_DATA_STORAGE_DEGRADED_RATIO, + DEFAULT_PROJECT_DATA_STORAGE_DEGRADED_RATIO + ); + + const thresholdsAreOrdered = + noticeRatio < warningRatio && warningRatio < criticalRatio && criticalRatio < degradedRatio; + + return { + enabled: envFlagEnabled(env.PROJECT_DATA_STORAGE_TELEMETRY_ENABLED), + limitBytes: parsePositiveInteger( + env.PROJECT_DATA_STORAGE_LIMIT_BYTES, + DEFAULT_PROJECT_DATA_STORAGE_LIMIT_BYTES + ), + measureIntervalMs: parsePositiveInteger( + env.PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS, + DEFAULT_PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS + ), + alertIntervalMs: parsePositiveInteger( + env.PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS, + DEFAULT_PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS + ), + noticeRatio: thresholdsAreOrdered ? noticeRatio : DEFAULT_PROJECT_DATA_STORAGE_NOTICE_RATIO, + warningRatio: thresholdsAreOrdered ? warningRatio : DEFAULT_PROJECT_DATA_STORAGE_WARNING_RATIO, + criticalRatio: thresholdsAreOrdered + ? criticalRatio + : DEFAULT_PROJECT_DATA_STORAGE_CRITICAL_RATIO, + degradedRatio: thresholdsAreOrdered + ? degradedRatio + : DEFAULT_PROJECT_DATA_STORAGE_DEGRADED_RATIO, + emergencyTargetRatio: parseBoundedRatio( + env.PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO, + DEFAULT_PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO + ), + emergencyBatchRows: parsePositiveInteger( + env.PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS, + DEFAULT_PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS + ), + emergencyMaxBatches: parsePositiveInteger( + env.PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES, + DEFAULT_PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES + ), + }; +} + +export function classifyStorageUsage( + databaseSizeBytes: number, + config: Pick< + StorageSafetyConfig, + 'limitBytes' | 'noticeRatio' | 'warningRatio' | 'criticalRatio' | 'degradedRatio' + > +): ProjectDataStorageStatus { + const usageRatio = databaseSizeBytes / config.limitBytes; + if (usageRatio >= config.degradedRatio) return 'degraded'; + if (usageRatio >= config.criticalRatio) return 'critical'; + if (usageRatio >= config.warningRatio) return 'warning'; + if (usageRatio >= config.noticeRatio) return 'notice'; + return 'ok'; +} + +function readMeta(sql: SqlStorage, key: string): string | null { + const row = sql.exec('SELECT value FROM do_meta WHERE key = ?', key).toArray()[0]; + if (!isJsonRecord(row)) return null; + return typeof row.value === 'string' ? row.value : null; +} + +function readMetaNumber(sql: SqlStorage, key: string): number | null { + const raw = readMeta(sql, key); + if (!raw) return null; + const parsed = Number.parseInt(raw, 10); + return Number.isSafeInteger(parsed) ? parsed : null; +} + +function writeMeta(sql: SqlStorage, key: string, value: string): void { + sql.exec( + `INSERT INTO do_meta (key, value) + VALUES (?, ?) + ON CONFLICT(key) DO UPDATE SET value = excluded.value`, + key, + value + ); +} + +function truncate(value: string, maxLength: number): string { + return value.length <= maxLength ? value : value.slice(0, maxLength); +} + +function buildTelemetry( + sql: SqlStorage, + env: Env, + projectId: string, + measuredAt: number = Date.now() +): ProjectDataStorageTelemetry { + const config = resolveStorageSafetyConfig(env); + const databaseSizeBytes = sql.databaseSize; + const usageRatio = databaseSizeBytes / config.limitBytes; + return { + projectId, + measuredAt, + databaseSizeBytes, + limitBytes: config.limitBytes, + usageRatio, + status: classifyStorageUsage(databaseSizeBytes, config), + }; +} + +async function upsertTelemetry( + env: Env, + telemetry: ProjectDataStorageTelemetry, + fields: { + lastAlarmAt?: number | null; + lastAlertAt?: number | null; + lastAlertStatus?: ProjectDataStorageStatus | null; + lastPurgeAt?: number | null; + lastPurgeReason?: string | null; + lastPurgeRows?: number | null; + lastPurgeDatabaseSizeBytes?: number | null; + lastError?: string | null; + } = {} +): Promise { + await env.DATABASE.prepare( + `INSERT INTO project_data_storage_telemetry ( + project_id, + measured_at, + database_size_bytes, + limit_bytes, + usage_ratio, + status, + last_alarm_at, + last_alert_at, + last_alert_status, + last_purge_at, + last_purge_reason, + last_purge_rows, + last_purge_database_size_bytes, + last_error, + updated_at + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(project_id) DO UPDATE SET + measured_at = excluded.measured_at, + database_size_bytes = excluded.database_size_bytes, + limit_bytes = excluded.limit_bytes, + usage_ratio = excluded.usage_ratio, + status = excluded.status, + last_alarm_at = COALESCE(excluded.last_alarm_at, project_data_storage_telemetry.last_alarm_at), + last_alert_at = COALESCE(excluded.last_alert_at, project_data_storage_telemetry.last_alert_at), + last_alert_status = COALESCE(excluded.last_alert_status, project_data_storage_telemetry.last_alert_status), + last_purge_at = COALESCE(excluded.last_purge_at, project_data_storage_telemetry.last_purge_at), + last_purge_reason = COALESCE(excluded.last_purge_reason, project_data_storage_telemetry.last_purge_reason), + last_purge_rows = COALESCE(excluded.last_purge_rows, project_data_storage_telemetry.last_purge_rows), + last_purge_database_size_bytes = COALESCE(excluded.last_purge_database_size_bytes, project_data_storage_telemetry.last_purge_database_size_bytes), + last_error = excluded.last_error, + updated_at = excluded.updated_at` + ) + .bind( + telemetry.projectId, + telemetry.measuredAt, + telemetry.databaseSizeBytes, + telemetry.limitBytes, + telemetry.usageRatio, + telemetry.status, + fields.lastAlarmAt ?? null, + fields.lastAlertAt ?? null, + fields.lastAlertStatus ?? null, + fields.lastPurgeAt ?? null, + fields.lastPurgeReason ? truncate(fields.lastPurgeReason, 500) : null, + fields.lastPurgeRows ?? null, + fields.lastPurgeDatabaseSizeBytes ?? null, + fields.lastError ? truncate(fields.lastError, 1000) : null, + Date.now() + ) + .run(); +} + +async function maybePersistStorageAlert( + sql: SqlStorage, + env: Env, + telemetry: ProjectDataStorageTelemetry +): Promise { + if (telemetry.status !== 'critical' && telemetry.status !== 'degraded') return; + const config = resolveStorageSafetyConfig(env); + const now = Date.now(); + const lastAlertAt = readMetaNumber(sql, META_LAST_ALERT_AT); + const lastAlertStatus = readMeta(sql, META_LAST_ALERT_STATUS); + if ( + lastAlertAt !== null && + now - lastAlertAt < config.alertIntervalMs && + lastAlertStatus === telemetry.status + ) { + return; + } + + writeMeta(sql, META_LAST_ALERT_AT, String(now)); + writeMeta(sql, META_LAST_ALERT_STATUS, telemetry.status); + + log.error('threshold_exceeded', { + projectId: telemetry.projectId, + status: telemetry.status, + databaseSizeBytes: telemetry.databaseSizeBytes, + limitBytes: telemetry.limitBytes, + usageRatio: telemetry.usageRatio, + }); + + if (!env.OBSERVABILITY_DATABASE) return; + + await persistError( + env.OBSERVABILITY_DATABASE, + { + source: 'api', + level: telemetry.status === 'degraded' ? 'error' : 'warn', + message: `ProjectData storage usage is ${telemetry.status}`, + context: { + projectId: telemetry.projectId, + databaseSizeBytes: telemetry.databaseSizeBytes, + limitBytes: telemetry.limitBytes, + usageRatio: telemetry.usageRatio, + status: telemetry.status, + }, + }, + undefined + ); + + await upsertTelemetry(env, telemetry, { + lastAlertAt: now, + lastAlertStatus: telemetry.status, + }); +} + +export function computeStorageSafetyAlarmTime( + sql: SqlStorage, + env: Env, + now: number = Date.now() +): number | null { + const config = resolveStorageSafetyConfig(env); + if (!config.enabled) return null; + if (!readMeta(sql, 'projectId')) return null; + const lastMeasuredAt = readMetaNumber(sql, META_LAST_MEASURED_AT); + return lastMeasuredAt === null ? now : lastMeasuredAt + config.measureIntervalMs; +} + +export async function measureAndPersistProjectDataStorage( + sql: SqlStorage, + env: Env, + projectId: string | null, + reason: 'alarm' | 'admin' = 'alarm' +): Promise { + const config = resolveStorageSafetyConfig(env); + if (!config.enabled) return null; + if (!projectId) { + log.warn('measure_skipped_missing_project_id'); + return null; + } + + const telemetry = buildTelemetry(sql, env, projectId); + writeMeta(sql, META_LAST_MEASURED_AT, String(telemetry.measuredAt)); + writeMeta(sql, META_LAST_STATUS, telemetry.status); + + try { + await upsertTelemetry(env, telemetry, { + lastAlarmAt: reason === 'alarm' ? telemetry.measuredAt : null, + lastError: null, + }); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + writeMeta(sql, META_LAST_ERROR, truncate(message, 500)); + log.warn('telemetry_upsert_failed', { + projectId, + ...serializeError(error), + }); + } + + try { + await maybePersistStorageAlert(sql, env, telemetry); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + writeMeta(sql, META_LAST_ERROR, truncate(message, 500)); + log.warn('alert_failed', { + projectId, + ...serializeError(error), + }); + } + + return telemetry; +} + +function normalizeCount(row: unknown): number { + if (!isJsonRecord(row)) return 0; + const count = row.count; + return typeof count === 'number' && Number.isFinite(count) ? count : 0; +} + +function countOldestActivityEventRows(sql: SqlStorage, limit: number): number { + const row = sql + .exec( + `SELECT COUNT(*) AS count + FROM (SELECT id FROM activity_events ORDER BY created_at ASC LIMIT ?)`, + limit + ) + .toArray()[0]; + return normalizeCount(row); +} + +function countOldestAcpSessionEventRows(sql: SqlStorage, limit: number): number { + const row = sql + .exec( + `SELECT COUNT(*) AS count + FROM (SELECT id FROM acp_session_events ORDER BY created_at ASC LIMIT ?)`, + limit + ) + .toArray()[0]; + return normalizeCount(row); +} + +function countOldestRows( + sql: SqlStorage, + table: 'activity_events' | 'acp_session_events', + limit: number +): number { + if (table === 'activity_events') { + return countOldestActivityEventRows(sql, limit); + } + return countOldestAcpSessionEventRows(sql, limit); +} + +function deleteOldestActivityEventRows(sql: SqlStorage, limit: number): void { + sql.exec( + `DELETE FROM activity_events + WHERE id IN ( + SELECT id FROM activity_events + ORDER BY created_at ASC + LIMIT ? + )`, + limit + ); +} + +function deleteOldestAcpSessionEventRows(sql: SqlStorage, limit: number): void { + sql.exec( + `DELETE FROM acp_session_events + WHERE id IN ( + SELECT id FROM acp_session_events + ORDER BY created_at ASC + LIMIT ? + )`, + limit + ); +} + +function deleteOldestRows( + sql: SqlStorage, + table: 'activity_events' | 'acp_session_events', + limit: number +): number { + const candidateCount = countOldestRows(sql, table, limit); + if (candidateCount <= 0) return 0; + + if (table === 'activity_events') { + deleteOldestActivityEventRows(sql, limit); + } else { + deleteOldestAcpSessionEventRows(sql, limit); + } + + return candidateCount; +} + +export async function runProjectDataStorageEmergencyPurge( + sql: SqlStorage, + env: Env, + projectId: string | null, + input: ProjectDataStorageEmergencyPurgeInput = {} +): Promise { + if (!projectId) { + throw new Error('ProjectData storage purge requires a persisted projectId'); + } + + const config = resolveStorageSafetyConfig(env); + const targetRatio = input.targetRatio && input.targetRatio > 0 && input.targetRatio < 1 + ? input.targetRatio + : config.emergencyTargetRatio; + const batchRows = input.batchRows && Number.isSafeInteger(input.batchRows) && input.batchRows > 0 + ? input.batchRows + : config.emergencyBatchRows; + const maxBatches = + input.maxBatches && Number.isSafeInteger(input.maxBatches) && input.maxBatches > 0 + ? input.maxBatches + : config.emergencyMaxBatches; + const targetBytes = Math.floor(config.limitBytes * targetRatio); + const reason = truncate(input.reason?.trim() || 'manual_emergency_purge', 500); + const beforeBytes = sql.databaseSize; + const statusBefore = classifyStorageUsage(beforeBytes, config); + const rowsDeleted = { activityEvents: 0, acpSessionEvents: 0 }; + let batches = 0; + let exhaustedCandidates = false; + + while (sql.databaseSize > targetBytes && batches < maxBatches) { + const activityDeleted = deleteOldestRows(sql, 'activity_events', batchRows); + const acpDeleted = deleteOldestRows(sql, 'acp_session_events', batchRows); + rowsDeleted.activityEvents += activityDeleted; + rowsDeleted.acpSessionEvents += acpDeleted; + batches++; + + if (activityDeleted === 0 && acpDeleted === 0) { + exhaustedCandidates = true; + break; + } + } + + const afterBytes = sql.databaseSize; + const statusAfter = classifyStorageUsage(afterBytes, config); + const totalRowsDeleted = rowsDeleted.activityEvents + rowsDeleted.acpSessionEvents; + const result: ProjectDataStorageEmergencyPurgeResult = { + projectId, + reason, + beforeBytes, + afterBytes, + limitBytes: config.limitBytes, + targetBytes, + statusBefore, + statusAfter, + batches, + maxBatches, + batchRows, + rowsDeleted, + exhaustedCandidates, + }; + + const measuredAt = Date.now(); + const telemetry: ProjectDataStorageTelemetry = { + projectId, + measuredAt, + databaseSizeBytes: afterBytes, + limitBytes: config.limitBytes, + usageRatio: afterBytes / config.limitBytes, + status: statusAfter, + }; + writeMeta(sql, META_LAST_MEASURED_AT, String(measuredAt)); + writeMeta(sql, META_LAST_STATUS, statusAfter); + + try { + await upsertTelemetry(env, telemetry, { + lastPurgeAt: measuredAt, + lastPurgeReason: reason, + lastPurgeRows: totalRowsDeleted, + lastPurgeDatabaseSizeBytes: afterBytes, + lastError: null, + }); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + writeMeta(sql, META_LAST_ERROR, truncate(message, 500)); + log.warn('purge_telemetry_upsert_failed', { + projectId, + ...serializeError(error), + }); + } + + log.warn('emergency_purge_completed', { ...result }); + return result; +} diff --git a/apps/api/src/durable-objects/project-data/tool-metadata-storage.ts b/apps/api/src/durable-objects/project-data/tool-metadata-storage.ts new file mode 100644 index 000000000..6dc62be5b --- /dev/null +++ b/apps/api/src/durable-objects/project-data/tool-metadata-storage.ts @@ -0,0 +1,141 @@ +import { + type CompactMessageOptions, + DEFAULT_DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES, + stripToolMetadataContent, +} from './row-schemas'; +import type { Env } from './types'; + +export const DEFAULT_PROJECT_DATA_TOOL_METADATA_MAX_BYTES = 128 * 1024; + +const textEncoder = new TextEncoder(); +const MINIMAL_TOOL_METADATA_KEYS = [ + 'toolCallId', + 'title', + 'kind', + 'status', + 'name', + 'tool', + 'exitCode', + 'contentSize', +] as const; + +export function resolveCompactMessageOptions(env: Env): CompactMessageOptions { + const parsed = Number.parseInt(env.DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES || '', 10); + return { + documentCardRawOutputMaxBytes: + Number.isFinite(parsed) && parsed > 0 ? parsed : DEFAULT_DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES, + }; +} + +function resolveToolMetadataMaxBytes(env: Env): number { + const parsed = Number.parseInt(env.PROJECT_DATA_TOOL_METADATA_MAX_BYTES || '', 10); + return Number.isFinite(parsed) && parsed > 0 + ? parsed + : DEFAULT_PROJECT_DATA_TOOL_METADATA_MAX_BYTES; +} + +function utf8Bytes(value: string): number { + return textEncoder.encode(value).byteLength; +} + +function isScalarJsonValue(value: unknown): boolean { + return value === null || ['string', 'number', 'boolean'].includes(typeof value); +} + +function buildMinimalToolMetadata(parsed: unknown, originalBytes: number): Record { + const minimal: Record = { + storageSafetyTruncated: true, + contentTruncated: true, + originalSizeBytes: originalBytes, + }; + + if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) { + return minimal; + } + + const record = parsed as Record; + for (const key of MINIMAL_TOOL_METADATA_KEYS) { + const value = record[key]; + if (isScalarJsonValue(value)) minimal[key] = value; + } + + if (Array.isArray(record.content)) { + minimal.contentSize = utf8Bytes(JSON.stringify(record.content)); + } + + return minimal; +} + +function serializeWithinLimit(value: unknown, maxBytes: number): string { + let serialized = JSON.stringify(value); + if (utf8Bytes(serialized) <= maxBytes) return serialized; + + const record = value && typeof value === 'object' && !Array.isArray(value) + ? (value as Record) + : {}; + serialized = JSON.stringify({ + storageSafetyTruncated: true, + contentTruncated: true, + originalSizeBytes: record.originalSizeBytes, + toolCallId: isScalarJsonValue(record.toolCallId) ? record.toolCallId : undefined, + title: isScalarJsonValue(record.title) ? record.title : undefined, + kind: isScalarJsonValue(record.kind) ? record.kind : undefined, + status: isScalarJsonValue(record.status) ? record.status : undefined, + contentSize: isScalarJsonValue(record.contentSize) ? record.contentSize : undefined, + }); + return utf8Bytes(serialized) <= maxBytes + ? serialized + : JSON.stringify({ storageSafetyTruncated: true }); +} + +export function boundToolMetadataForStorage( + toolMetadata: string | null, + env: Env +): { + value: string | null; + originalBytes: number; + storedBytes: number; + truncated: boolean; +} { + if (toolMetadata === null) { + return { value: null, originalBytes: 0, storedBytes: 0, truncated: false }; + } + + const maxBytes = resolveToolMetadataMaxBytes(env); + const originalBytes = utf8Bytes(toolMetadata); + if (originalBytes <= maxBytes) { + return { + value: toolMetadata, + originalBytes, + storedBytes: originalBytes, + truncated: false, + }; + } + + let bounded: string; + try { + const parsed = JSON.parse(toolMetadata); + const compact = stripToolMetadataContent(parsed, resolveCompactMessageOptions(env)); + const compactJson = JSON.stringify(compact); + bounded = utf8Bytes(compactJson) <= maxBytes + ? compactJson + : serializeWithinLimit(buildMinimalToolMetadata(parsed, originalBytes), maxBytes); + } catch { + bounded = serializeWithinLimit( + { + storageSafetyTruncated: true, + contentTruncated: true, + parseFailed: true, + originalSizeBytes: originalBytes, + }, + maxBytes + ); + } + + return { + value: bounded, + originalBytes, + storedBytes: utf8Bytes(bounded), + truncated: true, + }; +} diff --git a/apps/api/src/durable-objects/project-data/types.ts b/apps/api/src/durable-objects/project-data/types.ts index e97060c21..091ece82f 100644 --- a/apps/api/src/durable-objects/project-data/types.ts +++ b/apps/api/src/durable-objects/project-data/types.ts @@ -20,6 +20,18 @@ export type Env = { MAX_SESSIONS_PER_PROJECT?: string; MAX_MESSAGES_PER_SESSION?: string; DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES?: string; + PROJECT_DATA_TOOL_METADATA_MAX_BYTES?: string; + PROJECT_DATA_STORAGE_TELEMETRY_ENABLED?: string; + PROJECT_DATA_STORAGE_LIMIT_BYTES?: string; + PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS?: string; + PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS?: string; + PROJECT_DATA_STORAGE_NOTICE_RATIO?: string; + PROJECT_DATA_STORAGE_WARNING_RATIO?: string; + PROJECT_DATA_STORAGE_CRITICAL_RATIO?: string; + PROJECT_DATA_STORAGE_DEGRADED_RATIO?: string; + PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO?: string; + PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS?: string; + PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES?: string; ACTIVITY_RETENTION_DAYS?: string; SESSION_IDLE_TIMEOUT_MINUTES?: string; IDLE_CLEANUP_RETRY_DELAY_MS?: string; @@ -71,6 +83,7 @@ export type Env = { ORCHESTRATOR_WAIT_MAX_ACTIVE_PER_PROJECT?: string; ORCHESTRATOR_WAIT_MAX_DURATION_MS?: string; ORCHESTRATOR_WAIT_MAX_CANDIDATES_PER_ALARM?: string; + OBSERVABILITY_DATABASE?: D1Database; }; export interface SummaryData { diff --git a/apps/api/src/env.ts b/apps/api/src/env.ts index 48d005cf4..830fb58dd 100644 --- a/apps/api/src/env.ts +++ b/apps/api/src/env.ts @@ -490,6 +490,18 @@ export interface Env extends WebhookTriggerEnv, TaskRecoveryEnv { MAX_SESSIONS_PER_PROJECT?: string; MAX_MESSAGES_PER_SESSION?: string; DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES?: string; // Max document-card rawOutput bytes preserved in compact message metadata (default: 16384) + PROJECT_DATA_TOOL_METADATA_MAX_BYTES?: string; + PROJECT_DATA_STORAGE_TELEMETRY_ENABLED?: string; + PROJECT_DATA_STORAGE_LIMIT_BYTES?: string; + PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS?: string; + PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS?: string; + PROJECT_DATA_STORAGE_NOTICE_RATIO?: string; + PROJECT_DATA_STORAGE_WARNING_RATIO?: string; + PROJECT_DATA_STORAGE_CRITICAL_RATIO?: string; + PROJECT_DATA_STORAGE_DEGRADED_RATIO?: string; + PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO?: string; + PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS?: string; + PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES?: string; MESSAGE_SIZE_THRESHOLD?: string; ACTIVITY_RETENTION_DAYS?: string; SESSION_IDLE_TIMEOUT_MINUTES?: string; diff --git a/apps/api/src/routes/admin.ts b/apps/api/src/routes/admin.ts index b30233f5d..4c8b36e68 100644 --- a/apps/api/src/routes/admin.ts +++ b/apps/api/src/routes/admin.ts @@ -13,13 +13,28 @@ import { AdminUserActionSchema, AdminUserRoleSchema, jsonValidator, + parseOptionalBody, + ProjectDataStorageEmergencyPurgeSchema, UpdateSignupApprovalConfigSchema, } from '../schemas'; import { getRuntimeLimits } from '../services/limits'; +import { + measureProjectDataStorage, + runProjectDataStorageEmergencyPurge, +} from '../services/project-data'; import { getSignupApprovalConfig, setSignupApprovalConfig } from '../services/signup-approval'; import { adminObservabilityRoutes } from './admin/observability'; const adminRoutes = new Hono<{ Bindings: Env }>(); +const PROJECT_DATA_STORAGE_STATUSES = new Set([ + 'ok', + 'notice', + 'warning', + 'critical', + 'degraded', +]); +const DEFAULT_STORAGE_TELEMETRY_LIMIT = 50; +const MAX_STORAGE_TELEMETRY_LIMIT = 200; // All admin routes require auth + approval + superadmin adminRoutes.use('/*', requireAuth(), requireApproved(), requireSuperadmin()); @@ -247,6 +262,96 @@ adminRoutes.get('/tasks/recent-failures', async (c) => { return c.json({ tasks: failures }); }); +/** + * GET /api/admin/project-data/storage - Read latest ProjectData storage telemetry. + * + * Bounded D1 query only. Force-measure a specific project through the POST + * endpoint when the row is missing or stale. + */ +adminRoutes.get('/project-data/storage', async (c) => { + const status = c.req.query('status')?.trim(); + if (status && !PROJECT_DATA_STORAGE_STATUSES.has(status)) { + throw errors.badRequest('status must be ok, notice, warning, critical, or degraded'); + } + + const projectId = c.req.query('projectId')?.trim(); + const limitParam = c.req.query('limit'); + const parsedLimit = limitParam ? Number.parseInt(limitParam, 10) : DEFAULT_STORAGE_TELEMETRY_LIMIT; + if ( + !Number.isSafeInteger(parsedLimit) || + parsedLimit < 1 || + parsedLimit > MAX_STORAGE_TELEMETRY_LIMIT + ) { + throw errors.badRequest(`limit must be between 1 and ${MAX_STORAGE_TELEMETRY_LIMIT}`); + } + + const filters: string[] = []; + const params: Array = []; + if (status) { + filters.push('t.status = ?'); + params.push(status); + } + if (projectId) { + filters.push('t.project_id = ?'); + params.push(projectId); + } + + const whereClause = filters.length > 0 ? `WHERE ${filters.join(' AND ')}` : ''; + const result = await c.env.DATABASE.prepare( + `SELECT + t.project_id, + p.name AS project_name, + p.repository AS repository, + t.measured_at, + t.database_size_bytes, + t.limit_bytes, + t.usage_ratio, + t.status, + t.last_alarm_at, + t.last_alert_at, + t.last_alert_status, + t.last_purge_at, + t.last_purge_reason, + t.last_purge_rows, + t.last_purge_database_size_bytes, + t.last_error, + t.updated_at + FROM project_data_storage_telemetry t + LEFT JOIN projects p ON p.id = t.project_id + ${whereClause} + ORDER BY t.usage_ratio DESC, t.measured_at DESC + LIMIT ?` + ) + .bind(...params, parsedLimit) + .all(); + + return c.json({ telemetry: result.results ?? [] }); +}); + +/** + * POST /api/admin/project-data/storage/:projectId/measure - Force a measurement. + */ +adminRoutes.post('/project-data/storage/:projectId/measure', async (c) => { + const { projectId } = c.req.param(); + if (!projectId) throw errors.badRequest('projectId is required'); + const telemetry = await measureProjectDataStorage(c.env, projectId); + return c.json({ telemetry }); +}); + +/** + * POST /api/admin/project-data/storage/:projectId/emergency-purge + * + * Explicit superadmin-only recovery path. Deletes only bounded batches of + * oldest ProjectData event-log rows (activity_events and acp_session_events). + */ +adminRoutes.post('/project-data/storage/:projectId/emergency-purge', async (c) => { + const { projectId } = c.req.param(); + if (!projectId) throw errors.badRequest('projectId is required'); + const body = await parseOptionalBody(c.req.raw, ProjectDataStorageEmergencyPurgeSchema, {}); + const result = await runProjectDataStorageEmergencyPurge(c.env, projectId, body); + return c.json({ result }); +}); + // Admin observability routes (spec 023) — extracted sub-router (rule 18 file-size split) adminRoutes.route('/observability', adminObservabilityRoutes); diff --git a/apps/api/src/schemas/admin.ts b/apps/api/src/schemas/admin.ts index 9cf3ff144..923c080b3 100644 --- a/apps/api/src/schemas/admin.ts +++ b/apps/api/src/schemas/admin.ts @@ -12,6 +12,13 @@ export const UpdateSignupApprovalConfigSchema = v.object({ requireApproval: v.boolean(), }); +export const ProjectDataStorageEmergencyPurgeSchema = v.object({ + reason: v.optional(v.pipe(v.string(), v.maxLength(500))), + targetRatio: v.optional(v.pipe(v.number(), v.minValue(0.1), v.maxValue(0.99))), + batchRows: v.optional(v.pipe(v.number(), v.integer(), v.minValue(1), v.maxValue(5000))), + maxBatches: v.optional(v.pipe(v.number(), v.integer(), v.minValue(1), v.maxValue(100))), +}); + export const AnalyticsForwardSchema = v.object({ startDate: v.optional(v.string()), endDate: v.optional(v.string()), diff --git a/apps/api/src/schemas/index.ts b/apps/api/src/schemas/index.ts index 4cae9fb99..131cfe521 100644 --- a/apps/api/src/schemas/index.ts +++ b/apps/api/src/schemas/index.ts @@ -120,6 +120,7 @@ export { AdminUserRoleSchema, AnalyticsForwardSchema, CreatePlatformCredentialSchema, + ProjectDataStorageEmergencyPurgeSchema, UpdatePlatformCredentialSchema, UpdatePlatformIntegrationConfigSchema, UpdateSignupApprovalConfigSchema, diff --git a/apps/api/src/services/durable-object-retry.ts b/apps/api/src/services/durable-object-retry.ts index 666143abc..bdb45a2de 100644 --- a/apps/api/src/services/durable-object-retry.ts +++ b/apps/api/src/services/durable-object-retry.ts @@ -11,6 +11,13 @@ const TRANSIENT_DURABLE_OBJECT_PATTERNS = [ /overload.*durable object/i, ]; +const DURABLE_OBJECT_STORAGE_FULL_PATTERNS = [ + /\bSQLITE_FULL\b/i, + /database or disk is full/i, + /durable object.*storage.*full/i, + /sqlite.*full/i, +]; + export interface DurableObjectRetryEnv { DO_RETRY_MAX_ATTEMPTS?: string; DO_RETRY_BASE_DELAY_MS?: string; @@ -26,9 +33,16 @@ export interface DurableObjectRetryConfig { export function isTransientDurableObjectError(err: unknown): boolean { const message = extractErrorMessage(err); if (!message) return false; + if (isDurableObjectStorageFullError(err)) return false; return TRANSIENT_DURABLE_OBJECT_PATTERNS.some((pattern) => pattern.test(message)); } +export function isDurableObjectStorageFullError(err: unknown): boolean { + const message = extractErrorMessage(err); + if (!message) return false; + return DURABLE_OBJECT_STORAGE_FULL_PATTERNS.some((pattern) => pattern.test(message)); +} + export function getDurableObjectRetryConfig(env: DurableObjectRetryEnv): DurableObjectRetryConfig { return { maxAttempts: parsePositiveInt(env.DO_RETRY_MAX_ATTEMPTS, DEFAULT_DO_RETRY_MAX_ATTEMPTS), diff --git a/apps/api/src/services/project-data-storage-errors.ts b/apps/api/src/services/project-data-storage-errors.ts new file mode 100644 index 000000000..777fd3cac --- /dev/null +++ b/apps/api/src/services/project-data-storage-errors.ts @@ -0,0 +1,27 @@ +import { AppError } from '../middleware/error'; + +export const PROJECT_DATA_STORAGE_FULL = 'PROJECT_DATA_STORAGE_FULL'; + +export class ProjectDataStorageFullError extends AppError { + constructor(projectId: string, operation: string) { + super( + 507, + PROJECT_DATA_STORAGE_FULL, + 'ProjectData storage is full; writes are paused until an administrator runs storage recovery.', + { + projectId, + operation, + } + ); + this.name = 'ProjectDataStorageFullError'; + } +} + +export function toProjectDataStorageFullError( + projectId: string, + operation: string, + cause: unknown +): ProjectDataStorageFullError { + if (cause instanceof ProjectDataStorageFullError) return cause; + return new ProjectDataStorageFullError(projectId, operation); +} diff --git a/apps/api/src/services/project-data.ts b/apps/api/src/services/project-data.ts index ec94ab0e2..acc485ca9 100644 --- a/apps/api/src/services/project-data.ts +++ b/apps/api/src/services/project-data.ts @@ -29,9 +29,11 @@ import { log } from '../lib/logger'; import { computeDurableObjectRetryDelayMs, getDurableObjectRetryConfig, + isDurableObjectStorageFullError, isTransientDurableObjectError, } from './durable-object-retry'; import { ensureOncePerIsolate, forgetEnsuredProjectData } from './project-data-ensure-memo'; +import { toProjectDataStorageFullError } from './project-data-storage-errors'; /** * Get a typed DO stub for the given project and ensure the DO knows its projectId. @@ -64,6 +66,28 @@ function forgetEnsuredProject(env: Env, projectId: string): void { forgetEnsuredProjectData(env.PROJECT_DATA.idFromName(projectId).toString()); } +function normalizeProjectDataRpcError(projectId: string, operation: string, err: unknown): unknown { + if (isDurableObjectStorageFullError(err)) { + return toProjectDataStorageFullError(projectId, operation, err); + } + return err; +} + +async function callProjectDataNoRetry( + env: Env, + projectId: string, + operation: string, + call: (stub: DurableObjectStub) => Promise +): Promise { + try { + const stub = await getStub(env, projectId); + return await call(stub); + } catch (err) { + forgetEnsuredProject(env, projectId); + throw normalizeProjectDataRpcError(projectId, operation, err); + } +} + async function callProjectDataWithRetry( env: Env, projectId: string, @@ -84,6 +108,10 @@ async function callProjectDataWithRetry( // re-ensure. Idempotent, and defence in depth only — see the memo module. forgetEnsuredProject(env, projectId); + if (isDurableObjectStorageFullError(err)) { + throw toProjectDataStorageFullError(projectId, operation, err); + } + if (attempt >= retryConfig.maxAttempts || !isTransientDurableObjectError(err)) { throw err; } @@ -105,7 +133,11 @@ async function callProjectDataWithRetry( } } - throw lastError ?? new Error('ProjectData DO retry exhausted without an error'); + throw normalizeProjectDataRpcError( + projectId, + operation, + lastError ?? new Error('ProjectData DO retry exhausted without an error') + ); } function sleep(ms: number): Promise { @@ -124,8 +156,9 @@ export async function createSession( taskId: string | null = null, createdByUserId: string | null = null ): Promise { - const stub = await getStub(env, projectId); - return stub.createSession(workspaceId, topic, taskId, createdByUserId); + return callProjectDataNoRetry(env, projectId, 'createSession', (stub) => + stub.createSession(workspaceId, topic, taskId, createdByUserId) + ); } export async function linkSessionToWorkspace( @@ -193,13 +226,14 @@ export async function persistMessage( toolMetadata: Record | null, messageId?: string ): Promise { - const stub = await getStub(env, projectId); - return stub.persistMessage( - sessionId, - role, - content, - toolMetadata ? JSON.stringify(toolMetadata) : null, - messageId + return callProjectDataNoRetry(env, projectId, 'persistMessage', (stub) => + stub.persistMessage( + sessionId, + role, + content, + toolMetadata ? JSON.stringify(toolMetadata) : null, + messageId + ) ); } @@ -223,21 +257,22 @@ export async function persistMessageBatch( maxMessages?: number; remainingCapacity?: number; }> { - const stub = await getStub(env, projectId); - return stub.persistMessageBatch( - sessionId, - messages.map((m) => ({ - messageId: m.messageId, - role: m.role, - content: m.content, - toolMetadata: m.toolMetadata ? JSON.stringify(m.toolMetadata) : null, - timestamp: m.timestamp, - sequence: m.sequence, - // origin ("system" for SAM-injected messages) MUST be forwarded to the DO - // so the persisted message can be collapsed in the UI and excluded from - // dedup/search/topic/attention. Dropping it here silently loses the tag. - origin: m.origin ?? null, - })) + return callProjectDataNoRetry(env, projectId, 'persistMessageBatch', (stub) => + stub.persistMessageBatch( + sessionId, + messages.map((m) => ({ + messageId: m.messageId, + role: m.role, + content: m.content, + toolMetadata: m.toolMetadata ? JSON.stringify(m.toolMetadata) : null, + timestamp: m.timestamp, + sequence: m.sequence, + // origin ("system" for SAM-injected messages) MUST be forwarded to the DO + // so the persisted message can be collapsed in the UI and excluded from + // dedup/search/topic/attention. Dropping it here silently loses the tag. + origin: m.origin ?? null, + })) + ) ); } @@ -484,15 +519,16 @@ export async function recordActivityEvent( taskId: string | null, payload: Record | null ): Promise { - const stub = await getStub(env, projectId); - return stub.recordActivityEvent( - eventType, - actorType, - actorId, - workspaceId, - sessionId, - taskId, - payload ? JSON.stringify(payload) : null + return callProjectDataNoRetry(env, projectId, 'recordActivityEvent', (stub) => + stub.recordActivityEvent( + eventType, + actorType, + actorId, + workspaceId, + sessionId, + taskId, + payload ? JSON.stringify(payload) : null + ) ); } @@ -734,6 +770,27 @@ export async function getSummary( return stub.getSummary(); } +export async function measureProjectDataStorage(env: Env, projectId: string) { + return callProjectDataNoRetry(env, projectId, 'measureProjectDataStorage', (stub) => + stub.measureStorage() + ); +} + +export async function runProjectDataStorageEmergencyPurge( + env: Env, + projectId: string, + input: { + reason?: string | null; + targetRatio?: number | null; + batchRows?: number | null; + maxBatches?: number | null; + } = {} +) { + return callProjectDataNoRetry(env, projectId, 'runProjectDataStorageEmergencyPurge', (stub) => + stub.runStorageEmergencyPurge(input) + ); +} + // ========================================================================= // Workspace Activity Tracking // ========================================================================= diff --git a/apps/api/tests/unit/durable-objects/project-data-messages.test.ts b/apps/api/tests/unit/durable-objects/project-data-messages.test.ts index 7e193b337..cfc9076c2 100644 --- a/apps/api/tests/unit/durable-objects/project-data-messages.test.ts +++ b/apps/api/tests/unit/durable-objects/project-data-messages.test.ts @@ -1,6 +1,11 @@ import { afterEach, describe, expect, it, vi } from 'vitest'; import { getMessages } from '../../../src/durable-objects/project-data/messages'; +import { + boundToolMetadataForStorage, + DEFAULT_PROJECT_DATA_TOOL_METADATA_MAX_BYTES, +} from '../../../src/durable-objects/project-data/tool-metadata-storage'; +import type { Env } from '../../../src/durable-objects/project-data/types'; import { log } from '../../../src/lib/logger'; type QueryRow = Record; @@ -26,6 +31,81 @@ function makeSql(rows: QueryRow[]) { } as unknown as Parameters[0] & { exec: ReturnType }; } +describe('boundToolMetadataForStorage', () => { + it('keeps metadata unchanged when it is under the configured cap', () => { + const raw = JSON.stringify({ toolCallId: 'tc-1', status: 'completed' }); + const result = boundToolMetadataForStorage(raw, {} as Env); + + expect(DEFAULT_PROJECT_DATA_TOOL_METADATA_MAX_BYTES).toBe(128 * 1024); + expect(result).toMatchObject({ + value: raw, + originalBytes: raw.length, + storedBytes: raw.length, + truncated: false, + }); + }); + + it('strips oversized content arrays and preserves useful tool identity fields', () => { + const raw = JSON.stringify({ + toolCallId: 'tc-large', + title: 'Run shell command', + kind: 'shell', + status: 'completed', + content: [{ type: 'terminal', output: 'x'.repeat(4096), exitCode: 0 }], + }); + + const result = boundToolMetadataForStorage(raw, { + PROJECT_DATA_TOOL_METADATA_MAX_BYTES: '768', + } as Env); + + expect(result.truncated).toBe(true); + expect(result.storedBytes).toBeLessThanOrEqual(768); + const stored = JSON.parse(result.value ?? '{}') as Record; + expect(stored).toMatchObject({ + toolCallId: 'tc-large', + title: 'Run shell command', + kind: 'shell', + status: 'completed', + }); + expect(stored.content).toBeUndefined(); + expect(stored.contentSize).toBeGreaterThan(4096); + }); + + it('falls back to a minimal valid JSON marker when compact metadata is still too large', () => { + const raw = JSON.stringify({ + toolCallId: 'tc-minimal', + title: 'x'.repeat(2048), + status: 'completed', + content: [{ type: 'terminal', output: 'y'.repeat(2048) }], + }); + + const result = boundToolMetadataForStorage(raw, { + PROJECT_DATA_TOOL_METADATA_MAX_BYTES: '128', + } as Env); + + expect(result.truncated).toBe(true); + expect(result.storedBytes).toBeLessThanOrEqual(128); + expect(() => JSON.parse(result.value ?? '')).not.toThrow(); + expect(JSON.parse(result.value ?? '{}')).toMatchObject({ + storageSafetyTruncated: true, + }); + }); + + it('stores a valid JSON marker for oversized malformed metadata', () => { + const raw = `{${'x'.repeat(2048)}`; + const result = boundToolMetadataForStorage(raw, { + PROJECT_DATA_TOOL_METADATA_MAX_BYTES: '256', + } as Env); + + expect(result.truncated).toBe(true); + expect(result.storedBytes).toBeLessThanOrEqual(256); + expect(JSON.parse(result.value ?? '{}')).toMatchObject({ + storageSafetyTruncated: true, + parseFailed: true, + }); + }); +}); + describe('ProjectData messages getMessages', () => { afterEach(() => { vi.restoreAllMocks(); diff --git a/apps/api/tests/unit/services/durable-object-retry.test.ts b/apps/api/tests/unit/services/durable-object-retry.test.ts index 1b643373c..c6310f368 100644 --- a/apps/api/tests/unit/services/durable-object-retry.test.ts +++ b/apps/api/tests/unit/services/durable-object-retry.test.ts @@ -6,6 +6,7 @@ import { DEFAULT_DO_RETRY_MAX_ATTEMPTS, DEFAULT_DO_RETRY_MAX_DELAY_MS, getDurableObjectRetryConfig, + isDurableObjectStorageFullError, isTransientDurableObjectError, } from '../../../src/services/durable-object-retry'; @@ -36,6 +37,25 @@ describe('isTransientDurableObjectError', () => { expect(isTransientDurableObjectError(new Error('reset password token expired'))).toBe(false); expect(isTransientDurableObjectError(null)).toBe(false); }); + + it('does not retry SQLITE_FULL storage limit failures', () => { + const err = new Error('Durable Object storage operation failed: SQLITE_FULL'); + expect(isDurableObjectStorageFullError(err)).toBe(true); + expect(isTransientDurableObjectError(err)).toBe(false); + }); +}); + +describe('isDurableObjectStorageFullError', () => { + it('matches Cloudflare/SQLite full-storage variants', () => { + expect(isDurableObjectStorageFullError(new Error('SQLITE_FULL'))).toBe(true); + expect(isDurableObjectStorageFullError(new Error('database or disk is full'))).toBe(true); + expect(isDurableObjectStorageFullError(new Error('sqlite full while inserting'))).toBe(true); + }); + + it('does not match unrelated SQLite failures', () => { + expect(isDurableObjectStorageFullError(new Error('SQLITE_BUSY'))).toBe(false); + expect(isDurableObjectStorageFullError(new Error('database locked'))).toBe(false); + }); }); describe('getDurableObjectRetryConfig', () => { diff --git a/apps/api/tests/workers/project-data-storage-safety.test.ts b/apps/api/tests/workers/project-data-storage-safety.test.ts new file mode 100644 index 000000000..de44eb964 --- /dev/null +++ b/apps/api/tests/workers/project-data-storage-safety.test.ts @@ -0,0 +1,321 @@ +/** + * Worker-runtime coverage for the narrow ProjectData storage-safety firebreak. + * + * These tests intentionally use @cloudflare/vitest-pool-workers with real + * SQLite-backed Durable Objects. They settle behavior that is unsafe to infer: + * `databaseSize` after deletes, local limits on forcing exact SQLITE_FULL, and + * alarm-driven telemetry writes. + */ +import { env, runInDurableObject } from 'cloudflare:test'; +import { describe, expect, it } from 'vitest'; + +import { + DEFAULT_PROJECT_DATA_STORAGE_LIMIT_BYTES, + type ProjectDataStorageStatus, +} from '../../src/durable-objects/project-data/storage-safety'; +import type { Env as WorkerEnv } from '../../src/env'; +import * as projectDataService from '../../src/services/project-data'; +import { seedInstallation, seedProject, seedUser } from './helpers/seed-d1'; +import type { ProjectDataTestDouble } from './support/expected-error-doubles'; + +const testEnv = env as unknown as WorkerEnv; +const OWNER = 'storage-safety-owner'; +const INSTALLATION = 'storage-safety-installation'; + +function getStub(projectId: string): DurableObjectStub { + return env.PROJECT_DATA.get( + env.PROJECT_DATA.idFromName(projectId) + ) as DurableObjectStub; +} + +async function seedProjectGraph(projectId: string): Promise { + await seedUser(OWNER); + await seedInstallation(INSTALLATION, OWNER); + await seedProject(projectId, OWNER, INSTALLATION, { + name: `Storage Safety ${projectId}`, + }); +} + +async function withProjectDataStorageEnv( + overrides: Partial>, + fn: () => Promise +): Promise { + const mutableEnv = testEnv as WorkerEnv & Record; + const previous = new Map(); + for (const [key, value] of Object.entries(overrides)) { + previous.set(key, mutableEnv[key]); + mutableEnv[key] = value; + } + try { + return await fn(); + } finally { + for (const [key, value] of previous) { + if (value === undefined) delete mutableEnv[key]; + else mutableEnv[key] = value; + } + } +} + +async function readTelemetry(projectId: string) { + return env.DATABASE.prepare( + `SELECT + project_id, + measured_at, + database_size_bytes, + limit_bytes, + usage_ratio, + status, + last_alarm_at, + last_purge_rows + FROM project_data_storage_telemetry + WHERE project_id = ?` + ) + .bind(projectId) + .first<{ + project_id: string; + measured_at: number; + database_size_bytes: number; + limit_bytes: number; + usage_ratio: number; + status: ProjectDataStorageStatus; + last_alarm_at: number | null; + last_purge_rows: number | null; + }>(); +} + +describe('ProjectData storage safety firebreak', () => { + it('databaseSize drops after deleting rows in the workerd SQLite DO runtime', async () => { + const projectId = `storage-size-reclaim-${crypto.randomUUID()}`; + const stub = getStub(projectId); + await stub.ensureProjectId(projectId); + + const sizes = await runInDurableObject(stub, async (_instance, state) => { + const sql = state.storage.sql; + const before = sql.databaseSize; + sql.exec( + `CREATE TABLE storage_experiment_payloads ( + id TEXT PRIMARY KEY, + payload TEXT NOT NULL + )` + ); + const afterCreate = sql.databaseSize; + const payload = 'x'.repeat(4096); + for (let i = 0; i < 96; i++) { + sql.exec( + 'INSERT INTO storage_experiment_payloads (id, payload) VALUES (?, ?)', + `payload-${i}`, + payload + ); + } + const afterInsert = sql.databaseSize; + sql.exec('DELETE FROM storage_experiment_payloads'); + const afterDelete = sql.databaseSize; + return { before, afterCreate, afterInsert, afterDelete }; + }); + + expect(sizes.afterCreate).toBeGreaterThanOrEqual(sizes.before); + expect(sizes.afterInsert).toBeGreaterThan(sizes.afterCreate); + expect(sizes.afterDelete).toBeLessThan(sizes.afterInsert); + }); + + it('SqlStorage write-limit errors are catchable and exact SQLITE_FULL is not locally forceable', async () => { + const projectId = `storage-write-error-catch-${crypto.randomUUID()}`; + const stub = getStub(projectId); + await stub.ensureProjectId(projectId); + + const result = await runInDurableObject(stub, async (_instance, state) => { + const sql = state.storage.sql; + let pragmaPageCountAuthorized = true; + let pragmaPageCountMessage = ''; + try { + sql.exec('PRAGMA page_count').toArray(); + } catch (error) { + pragmaPageCountAuthorized = false; + pragmaPageCountMessage = error instanceof Error ? error.message : String(error); + } + + const setMaxPageCountForTest = ( + sql as unknown as { setMaxPageCountForTest?: (count: number) => void } + ).setMaxPageCountForTest; + + sql.exec('CREATE TABLE storage_write_error_probe (id TEXT PRIMARY KEY, payload TEXT NOT NULL)'); + let caught = false; + let message = ''; + try { + sql.exec( + 'INSERT INTO storage_write_error_probe (id, payload) VALUES (?, ?)', + 'oversized-row', + 'x'.repeat(3 * 1024 * 1024) + ); + } catch (error) { + caught = true; + message = error instanceof Error ? error.message : String(error); + } + + const readAfterError = sql + .exec('SELECT COUNT(*) AS count FROM storage_write_error_probe') + .toArray()[0] as { count?: number } | undefined; + sql.exec('DELETE FROM storage_write_error_probe'); + const deleteAfterError = sql + .exec('SELECT COUNT(*) AS count FROM storage_write_error_probe') + .toArray()[0] as { count?: number } | undefined; + + return { + pragmaPageCountAuthorized, + pragmaPageCountMessage, + setMaxPageCountForTestAvailable: typeof setMaxPageCountForTest === 'function', + caught, + message, + readAfterError: typeof readAfterError?.count === 'number', + deleteAfterError: deleteAfterError?.count === 0, + }; + }); + + expect(result.pragmaPageCountAuthorized).toBe(false); + expect(result.pragmaPageCountMessage).toMatch(/SQLITE_AUTH|not authorized/i); + expect(result.setMaxPageCountForTestAvailable).toBe(false); + expect(result.caught).toBe(true); + expect(result.message).toMatch(/too big|maximum|SQLITE_TOOBIG/i); + expect(result.readAfterError).toBe(true); + expect(result.deleteAfterError).toBe(true); + }); + + it('alarm measurement writes bounded D1 telemetry and reschedules the next measurement', async () => { + const projectId = `storage-alarm-${crypto.randomUUID()}`; + await seedProjectGraph(projectId); + const stub = getStub(projectId); + await stub.ensureProjectId(projectId); + await stub.createSession(null, 'Storage alarm'); + + await withProjectDataStorageEnv( + { PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS: '60000' }, + async () => { + await runInDurableObject(stub, async (instance) => instance.alarm()); + + const telemetry = await readTelemetry(projectId); + expect(telemetry).toMatchObject({ + project_id: projectId, + limit_bytes: DEFAULT_PROJECT_DATA_STORAGE_LIMIT_BYTES, + status: 'ok', + }); + expect(telemetry?.database_size_bytes).toBeGreaterThan(0); + expect(telemetry?.measured_at).toBeGreaterThan(0); + expect(telemetry?.last_alarm_at).toBeGreaterThan(0); + + const nextAlarm = await runInDurableObject(stub, async (_instance, state) => + state.storage.getAlarm() + ); + expect(nextAlarm).toBeTypeOf('number'); + expect(nextAlarm as number).toBeGreaterThan(Date.now()); + } + ); + }); + + it('service measurement writes ProjectData storage telemetry directly', async () => { + const projectId = `storage-service-measure-${crypto.randomUUID()}`; + await seedProjectGraph(projectId); + await projectDataService.createSession(testEnv, projectId, null, 'Measured via service'); + + const measurement = await projectDataService.measureProjectDataStorage(testEnv, projectId); + const telemetry = await readTelemetry(projectId); + + expect(measurement).toMatchObject({ + projectId, + limitBytes: DEFAULT_PROJECT_DATA_STORAGE_LIMIT_BYTES, + status: 'ok', + }); + expect(telemetry?.project_id).toBe(projectId); + expect(telemetry?.database_size_bytes).toBe(measurement?.databaseSizeBytes); + }); + + it('emergency purge deletes only bounded low-value event-log batches', async () => { + const projectId = `storage-purge-${crypto.randomUUID()}`; + await seedProjectGraph(projectId); + const stub = getStub(projectId); + await stub.ensureProjectId(projectId); + + const countsBefore = await runInDurableObject(stub, async (instance, state) => { + const sessionId = await instance.createSession(null, 'Purge guard'); + await instance.persistMessage(sessionId, 'user', 'Keep this message', null); + const acp = await instance.createAcpSession({ + chatSessionId: sessionId, + initialPrompt: null, + agentType: null, + }); + + const sql = state.storage.sql; + for (let i = 0; i < 5; i++) { + sql.exec( + `INSERT INTO activity_events + (id, event_type, actor_type, actor_id, workspace_id, session_id, task_id, payload, created_at) + VALUES (?, 'storage.test', 'system', NULL, NULL, ?, NULL, ?, ?)`, + `activity-${i}`, + sessionId, + JSON.stringify({ index: i, payload: 'x'.repeat(1024) }), + i + 1 + ); + sql.exec( + `INSERT INTO acp_session_events + (id, acp_session_id, from_status, to_status, actor_type, actor_id, reason, metadata, created_at) + VALUES (?, ?, NULL, 'running', 'system', NULL, 'storage-test', ?, ?)`, + `acp-event-${i}`, + acp.id, + JSON.stringify({ index: i, payload: 'y'.repeat(1024) }), + i + 1 + ); + } + + const activityCount = sql.exec('SELECT COUNT(*) AS count FROM activity_events').toArray()[0] as { + count: number; + }; + const acpEventCount = sql.exec('SELECT COUNT(*) AS count FROM acp_session_events').toArray()[0] as { + count: number; + }; + const messageCount = sql.exec('SELECT COUNT(*) AS count FROM chat_messages').toArray()[0] as { + count: number; + }; + return { activityCount, acpEventCount, messageCount }; + }); + + await withProjectDataStorageEnv({ PROJECT_DATA_STORAGE_LIMIT_BYTES: '10000' }, async () => { + const result = await projectDataService.runProjectDataStorageEmergencyPurge( + testEnv, + projectId, + { + reason: 'vitest bounded purge', + targetRatio: 0.1, + batchRows: 2, + maxBatches: 1, + } + ); + + expect(result.rowsDeleted).toEqual({ activityEvents: 2, acpSessionEvents: 2 }); + expect(result.batches).toBe(1); + expect(result.maxBatches).toBe(1); + expect(result.batchRows).toBe(2); + expect(result.beforeBytes).toBeGreaterThanOrEqual(result.afterBytes); + }); + + const countsAfter = await runInDurableObject(stub, async (_instance, state) => { + const sql = state.storage.sql; + const activityCount = sql.exec('SELECT COUNT(*) AS count FROM activity_events').toArray()[0] as { + count: number; + }; + const acpEventCount = sql.exec('SELECT COUNT(*) AS count FROM acp_session_events').toArray()[0] as { + count: number; + }; + const messageCount = sql.exec('SELECT COUNT(*) AS count FROM chat_messages').toArray()[0] as { + count: number; + }; + return { activityCount, acpEventCount, messageCount }; + }); + const telemetry = await readTelemetry(projectId); + + expect(countsBefore.activityCount.count).toBeGreaterThanOrEqual(6); + expect(countsBefore.acpEventCount.count).toBeGreaterThanOrEqual(6); + expect(countsAfter.activityCount.count).toBe(countsBefore.activityCount.count - 2); + expect(countsAfter.acpEventCount.count).toBe(countsBefore.acpEventCount.count - 2); + expect(countsAfter.messageCount.count).toBe(countsBefore.messageCount.count); + expect(telemetry?.last_purge_rows).toBe(4); + }); +}); diff --git a/apps/api/wrangler.toml b/apps/api/wrangler.toml index dfb9b8545..d6062d69a 100644 --- a/apps/api/wrangler.toml +++ b/apps/api/wrangler.toml @@ -50,6 +50,18 @@ NODE_HEARTBEAT_STALE_SECONDS = "180" MAX_PROJECTS_PER_USER = "100" MAX_SESSIONS_PER_PROJECT = "10000" MAX_MESSAGES_PER_SESSION = "100000" +PROJECT_DATA_TOOL_METADATA_MAX_BYTES = "131072" +PROJECT_DATA_STORAGE_TELEMETRY_ENABLED = "true" +PROJECT_DATA_STORAGE_LIMIT_BYTES = "10000000000" +PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS = "3600000" +PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS = "21600000" +PROJECT_DATA_STORAGE_NOTICE_RATIO = "0.6" +PROJECT_DATA_STORAGE_WARNING_RATIO = "0.8" +PROJECT_DATA_STORAGE_CRITICAL_RATIO = "0.9" +PROJECT_DATA_STORAGE_DEGRADED_RATIO = "0.95" +PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO = "0.9" +PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS = "500" +PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES = "4" MESSAGE_SIZE_THRESHOLD = "102400" ACTIVITY_RETENTION_DAYS = "90" SESSION_IDLE_TIMEOUT_MINUTES = "60" diff --git a/apps/www/src/content/docs/docs/reference/configuration.md b/apps/www/src/content/docs/docs/reference/configuration.md index cc3503494..ed40d7e9b 100644 --- a/apps/www/src/content/docs/docs/reference/configuration.md +++ b/apps/www/src/content/docs/docs/reference/configuration.md @@ -753,6 +753,18 @@ ProjectData stores a single prompt-delivery queue and checkpoint episodes keyed | `MAX_SESSIONS_PER_PROJECT` | `10000` | Max chat sessions per project | | `MAX_MESSAGES_PER_SESSION` | `100000` | Max messages per chat session | | `DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES` | `16384` | Max compact metadata bytes preserved for library document cards | +| `PROJECT_DATA_TOOL_METADATA_MAX_BYTES` | `131072` | Max stored `tool_metadata` bytes per message before oversized tool content is stripped into bounded metadata | +| `PROJECT_DATA_STORAGE_TELEMETRY_ENABLED` | `true` | Enables ProjectData `databaseSize` alarm measurement and D1 telemetry writes | +| `PROJECT_DATA_STORAGE_LIMIT_BYTES` | `10000000000` | Cloudflare SQLite-backed Durable Object storage limit used for ProjectData usage classification | +| `PROJECT_DATA_STORAGE_MEASURE_INTERVAL_MS` | `3600000` | Minimum interval between per-object ProjectData storage measurements | +| `PROJECT_DATA_STORAGE_ALERT_INTERVAL_MS` | `21600000` | Minimum interval between repeated critical/degraded ProjectData storage observability alerts | +| `PROJECT_DATA_STORAGE_NOTICE_RATIO` | `0.6` | ProjectData storage usage ratio classified as `notice` | +| `PROJECT_DATA_STORAGE_WARNING_RATIO` | `0.8` | ProjectData storage usage ratio classified as `warning` | +| `PROJECT_DATA_STORAGE_CRITICAL_RATIO` | `0.9` | ProjectData storage usage ratio classified as `critical` | +| `PROJECT_DATA_STORAGE_DEGRADED_RATIO` | `0.95` | ProjectData storage usage ratio classified as `degraded` | +| `PROJECT_DATA_STORAGE_EMERGENCY_TARGET_RATIO` | `0.9` | Target usage ratio for explicit superadmin ProjectData emergency purge calls | +| `PROJECT_DATA_STORAGE_EMERGENCY_BATCH_ROWS` | `500` | Oldest `activity_events` and `acp_session_events` rows deleted per table per emergency purge batch | +| `PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES` | `4` | Maximum emergency purge batches per explicit call | | `MESSAGE_SIZE_THRESHOLD` | `102400` | Max message size in bytes | | `ACTIVITY_RETENTION_DAYS` | `90` | Days to retain activity events | | `SESSION_IDLE_TIMEOUT_MINUTES` | `60` | Idle session timeout | diff --git a/tasks/active/2026-08-21-projectdata-storage-safety-firebreak.md b/tasks/active/2026-08-21-projectdata-storage-safety-firebreak.md new file mode 100644 index 000000000..3243c1ad2 --- /dev/null +++ b/tasks/active/2026-08-21-projectdata-storage-safety-firebreak.md @@ -0,0 +1,111 @@ +# ProjectData storage-safety firebreak + +## Problem statement + +Production evidence from idea `01M0B8HBA4YRJFF8STQ8PMJ1D8` strongly infers that at least one `ProjectData` SQLite-backed Durable Object is near Cloudflare's 10 GB per-object ceiling. Current `main` does not directly measure per-object `databaseSize`, does not classify storage-exhaustion failures distinctly, and lets high-volume message/tool metadata continue growing without a byte-oriented at-rest cap. + +The goal is a narrow production-safe firebreak, not full sharding: + +- measure and persist per-object `databaseSize` from inside each `ProjectData` object; +- expose operator diagnostics and alert rows before the ceiling becomes an outage; +- classify `SQLITE_FULL`/storage-limit errors separately from transient DO reset/retry conditions; +- provide a bounded emergency, fail-visible recovery path that only deletes low-value telemetry; +- cap or trim unbounded tool metadata at the write path if that can be done safely without R2/spill infrastructure. + +## Constraints + +- Base branch: current `main`; implementation branch/PR: `sam/implement-narrow-projectdata-storage-f3f5s3`. +- Preserve PR #1873 as DO NOT MERGE. Do not modify, supersede, close, merge, or otherwise disturb it. +- Do not deploy or mutate staging. Do not merge. Stop after draft PR + CI evidence until the parent grants the staging slot. +- Use official Cloudflare documentation and repository evidence. +- Run focused local/Miniflare experiments for storage-full catchability, deletion reclamation, and alarm execution behavior. +- Follow D1/DO migration safety, control-loop bounding/isolation, configurable defaults, documentation sync, and specialist review gates. + +## Research findings + +- Cloudflare Durable Object docs currently state that SQLite-backed Durable Objects have a 10 GB per-object storage limit on Workers Paid, 2 MB maximum row/string/blob size, 100 KB SQL statement length, and `database or disk is full: SQLITE_FULL` when writes exceed the storage limit. Docs also state reads and deletes continue to work at the limit. +- Cloudflare SQLite storage docs expose `ctx.storage.sql.databaseSize` as the current SQLite database size in bytes. +- workerd source currently computes `databaseSize` as `(pragma_page_count - pragma_freelist_count) * page_size`, so deleted pages are subtracted from the measured quota. +- Cloudflare alarm docs state each DO has one alarm, alarms are at-least-once, and uncaught alarm errors only get a finite retry series. Any new storage-safety alarm step must catch its own failures and leave unrelated alarm candidates working. +- Existing ProjectData alarm scheduling is centralized in `apps/api/src/durable-objects/project-data/alarm-schedule.ts:computeProjectDataAlarmTime`, and `ProjectData.alarm()` already multiplexes heartbeat, idle cleanup, reconciliation, attention, mailbox, prompt delivery, and task waits. +- `services/project-data.ts` wraps some ProjectData RPCs with retry logic via `durable-object-retry.ts`, whose current transient reset regex would classify a generic "Durable Object reset" as retryable but does not distinguish storage-full causes. +- `messages.ts` currently bounds only by `MAX_MESSAGES_PER_SESSION` row count. It stores `content` and `tool_metadata` raw. `persistMessageBatch()` parses persisted `tool_metadata` for broadcast, so trimming must preserve valid JSON. +- VM agent reporter already caps message content via `MSG_MAX_MESSAGE_CONTENT_BYTES`, but `ToolMetadata` remains bounded only indirectly by batch transport size. +- Existing D1 session-index migration `0117_session_index_per_project.sql` is the closest pattern for an additive per-project diagnostic table keyed to `projects(id)`. + +## Implementation checklist + +- [x] Add D1 migration and Drizzle schema for per-project ProjectData storage telemetry. +- [x] Add ProjectData storage-safety config with env-backed defaults for measurement cadence, SQLite limit bytes, warning/critical/degrade ratios, alert throttle, emergency purge watermark, and tool-metadata trim size. +- [x] Add a ProjectData storage-safety module that reads `sql.databaseSize`, computes status/watermark state, upserts the D1 telemetry row, logs structured diagnostics, and persists critical/degraded `platform_errors` rows without letting observability failures abort alarms. +- [x] Wire storage-safety measurement into the shared ProjectData alarm schedule and into `ProjectData.alarm()` as an isolated step. +- [x] Add admin diagnostic endpoints that list/query ProjectData storage telemetry, force-measure one project, and run bounded emergency purge for superadmins. +- [x] Add explicit `SQLITE_FULL` / storage-limit classification, keep it non-transient, and return a distinct fail-visible service error from ProjectData retry/write call sites. +- [x] Add a bounded emergency recovery RPC/path that deletes oldest low-value telemetry rows (`activity_events`, `acp_session_events`) in capped batches until below the configured recovery watermark or the configured batch bound, and records the result in telemetry. +- [x] Add safe write-path trimming for oversized `tool_metadata` content arrays, preserving metadata/card-critical fields and storing a truncation marker rather than invalid JSON. +- [x] Run focused local/Miniflare experiments for catchability, deletion reclamation, and alarm behavior; record evidence in this file. +- [x] Add focused unit/workers tests for classifier, telemetry upsert, alarm execution, emergency purge, and metadata trimming. +- [x] Update documentation/config references. +- [x] Run local specialist reviews: cloudflare-specialist, constitution-validator, test-engineer, doc-sync-validator, security-auditor, task-completion-validator. +- [x] Open a draft PR and let CI run. Stop before staging/deploy/merge. + +## Acceptance criteria + +- Operators can query the latest per-project ProjectData DO size and status from D1/admin diagnostics without relying on Cloudflare namespace-level analytics. +- ProjectData alarms periodically self-measure `sql.databaseSize` with bounded local work and persist/log threshold crossings. +- Storage-full errors are classified distinctly, not retried as generic transient DO resets, and callers receive a clear ProjectData-storage-full error rather than an opaque internal reset storm. +- An explicit superadmin/admin recovery path can delete only low-value telemetry in bounded batches and reports exact rows/bytes/status; it does not delete transcripts, knowledge, policies, session state, or mailbox prompts. +- Oversized tool metadata is bounded at the ProjectData write path without breaking lazy-load/read contracts or storing malformed JSON. +- New limits/time intervals use env-backed defaults and are documented. +- Local Miniflare/workerd experiments are recorded for catchability, deletion reclamation, and alarms. +- Focused tests and relevant quality gates pass locally; CI is allowed to run on a draft PR. + +## Experiment evidence + +- `databaseSize` reclamation: `pnpm vitest run --config vitest.workers.config.ts tests/workers/project-data-storage-safety.test.ts --reporter verbose` inserted ~384 KiB into a scratch DO SQLite table, then deleted it. `sql.databaseSize` increased after inserts and dropped after delete in the workerd/Miniflare runtime. +- Exact local `SQLITE_FULL` forcing: the same Worker-runtime test proved direct `PRAGMA page_count` is rejected with `SQLITE_AUTH`, and `state.storage.sql.setMaxPageCountForTest` is not exposed in the JS Workers test runtime. Therefore the implementation does not rely on an in-DO `SQLITE_FULL` catch hook for automatic recovery. It classifies `SQLITE_FULL` at the service boundary and exposes explicit admin recovery. +- Storage write-error catchability: the same test inserted an oversized row to trigger a SqlStorage write-limit exception inside the DO, caught it, then successfully read and deleted from the table afterward. This proves storage write exceptions do not necessarily reset the actor and that reads/deletes can continue after a caught write failure. +- Alarm behavior: the same test invoked `ProjectData.alarm()` in the Worker runtime, verified a D1 `project_data_storage_telemetry` row with `last_alarm_at`, and verified the ProjectData alarm was rescheduled for the next storage measurement interval. +- Emergency purge: the same test inserted activity and ACP event rows, ran one bounded purge batch (`batchRows=2`, `maxBatches=1`), verified only two oldest rows per low-value table were removed, verified chat messages remained, and verified D1 telemetry recorded `last_purge_rows=4`. + +## Specialist review tracker + +- cloudflare-specialist: PASS — additive D1 migration only; no DO schema migration; ProjectData alarm work is isolated/caught and rescheduled through the existing alarm multiplexer; no staging or production mutation performed. +- constitution-validator: PASS — new thresholds, intervals, caps, and batch bounds are environment-backed with documented defaults; no hardcoded external URLs/credentials/tenant IDs added. +- test-engineer: PASS — focused unit tests cover storage-full classification and tool metadata capping; Worker-runtime tests cover `databaseSize` reclamation, write-error catchability, alarm telemetry, service measurement, and bounded purge; existing ProjectData service Worker suite still passes. +- doc-sync-validator: PASS — env declarations, wrangler defaults, `.env.example`, public configuration docs, and API/env reference skills were updated for the new endpoints and config. +- security-auditor: PASS — recovery endpoints remain behind the existing admin auth/approval/superadmin middleware; storage-full service errors expose only `projectId` and operation; purge deletes only `activity_events` and `acp_session_events` in bounded batches. +- task-completion-validator: PASS — research findings, checklist items, and acceptance criteria are covered by the diff/tests; no UI inputs or multi-provider selection paths were introduced. + +## Local validation + +- `pnpm --filter @simple-agent-manager/api typecheck` — PASS +- `pnpm --filter @simple-agent-manager/api lint` — PASS +- `pnpm --filter @simple-agent-manager/api test -- tests/unit/durable-objects/project-data-messages.test.ts tests/unit/services/durable-object-retry.test.ts` — PASS (2 files, 20 tests) +- `pnpm vitest run --config vitest.workers.config.ts tests/workers/project-data-storage-safety.test.ts --reporter verbose` — PASS (1 file, 5 tests) +- `pnpm vitest run --config vitest.workers.config.ts tests/workers/project-data-service.test.ts` — PASS (1 file, 55 tests) +- `pnpm tsx scripts/quality/ast-checks.ts --file apps/api/src/durable-objects/project-data/storage-safety.ts --rule sql-injection` — PASS +- `pnpm quality:ast-checks` — PASS (0 errors, existing warnings only) +- `pnpm quality:file-sizes` — PASS +- `pnpm quality:runtime-boundary-semantics` — PASS +- `GITHUB_EVENT_NAME=pull_request GITHUB_EVENT_PATH=<(env -u GH_TOKEN -u GITHUB_TOKEN gh pr view 1875 --json body,url --jq '{pull_request:{body:.body,html_url:.url}}') pnpm quality:preflight` — PASS +- `git diff --check main...HEAD` — PASS + +## PR / CI evidence + +- Draft PR: https://github.com/raphaeltm/simple-agent-manager/pull/1875 +- Initial CI exposed four branch issues before later jobs completed: missing PR preflight evidence block, AST SQL-injection rule rejection of the purge helper's dynamic table-name template literal, `messages.ts` exceeding the mandatory 800-line file-size ceiling, and unguarded row narrowing in `storage-safety.ts`. +- Remediation: PR body now includes the required `AGENT_PREFLIGHT` evidence block; purge helper now uses fixed SQL statements for `activity_events` and `acp_session_events`; tool metadata bounding now lives in `apps/api/src/durable-objects/project-data/tool-metadata-storage.ts`, bringing `messages.ts` below the hard ceiling; SQL row reads now use the shared `isJsonRecord` guard. + +## Task completion validation report + +Verdict: PASS + +| Check | Status | Issues | +|-------|--------|--------| +| A: Research → Checklist | PASS | All research findings that identified a required change are covered by checklist items. | +| B: Checklist → Diff | PASS | Checked implementation items map to substantive code/test/doc changes. | +| C: Criteria → Tests | PASS | Acceptance criteria are covered by unit tests, Worker-runtime tests, and documented experiment evidence. | +| D: UI → Backend | N/A | No UI changes or new UI inputs were introduced. | +| E: Multi-Resource | N/A | No provider/resource selection logic was introduced. | +| F: Vertical Slice | PASS | Worker-runtime tests cover admin/service-to-DO-to-D1 measurement and service-to-DO purge behavior with real Durable Object SQLite storage. |