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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .claude/skills/api-reference/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
12 changes: 12 additions & 0 deletions .claude/skills/env-reference/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
14 changes: 14 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
33 changes: 33 additions & 0 deletions apps/api/src/db/migrations/0119_project_data_storage_telemetry.sql
Original file line number Diff line number Diff line change
@@ -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);
38 changes: 38 additions & 0 deletions apps/api/src/db/schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
{
Expand Down
3 changes: 3 additions & 0 deletions apps/api/src/durable-objects/project-data/alarm-schedule.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand Down Expand Up @@ -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,
Expand All @@ -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;
Expand Down
38 changes: 38 additions & 0 deletions apps/api/src/durable-objects/project-data/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -895,11 +896,48 @@ export class ProjectData extends DurableObject<Env> {
};
}

async measureStorage(): Promise<storageSafety.ProjectDataStorageTelemetry | null> {
const measurement = await storageSafety.measureAndPersistProjectDataStorage(
this.sql,
this.env,
this.getProjectId(),
'admin'
);
await this.recalculateAlarm();
return measurement;
}

async runStorageEmergencyPurge(
input: storageSafety.ProjectDataStorageEmergencyPurgeInput = {}
): Promise<storageSafety.ProjectDataStorageEmergencyPurgeResult> {
const result = await storageSafety.runProjectDataStorageEmergencyPurge(
this.sql,
this.env,
this.getProjectId(),
input
);
await this.recalculateAlarm();
return result;
}

// --- DO Alarm Handler ---

async alarm(): Promise<void> {
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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
67 changes: 52 additions & 15 deletions apps/api/src/durable-objects/project-data/messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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';

Expand All @@ -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.
*/
Expand All @@ -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)
Expand All @@ -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) {
Expand All @@ -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)
Expand All @@ -101,7 +122,7 @@ export function persistMessage(
sessionId,
role,
content,
toolMetadata,
boundedToolMetadata.value,
now,
sequence
);
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -251,14 +279,23 @@ 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 (?, ?, ?, ?, ?, ?, ?, ?)`,
msg.messageId,
sessionId,
msg.role,
msg.content,
msg.toolMetadata,
boundedToolMetadata.value,
createdAt,
sequence,
origin
Expand All @@ -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,
Expand Down
Loading
Loading