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
8 changes: 8 additions & 0 deletions .claude/skills/env-reference/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,14 @@ See `apps/api/.env.example` for the full list. Key variables:
- `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`)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_ENABLED` — Enables automatic ProjectData cleanup that strips expandable `tool_metadata.content` payloads from old terminal-session tool messages under storage pressure (default: enabled)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TRIGGER_RATIO` — ProjectData storage usage ratio that starts automatic terminal-session tool payload cleanup (default: `0.8`)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TARGET_RATIO` — ProjectData storage usage ratio below which automatic tool payload cleanup stops (default: `0.75`)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_ROWS` — Maximum tool-message rows inspected by one automatic cleanup alarm batch (default: `500`)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_BYTES` — Maximum legacy `tool_metadata` bytes read into JS by one automatic cleanup alarm batch (default: `1048576`)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MIN_SESSION_AGE_DAYS` — Minimum terminal-session age before automatic cleanup may strip stored tool payload content (default: `7`)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_RECHECK_MS` — Delay before the next automatic cleanup alarm batch when more candidates remain (default: `60000`)
- `PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MAX_SESSIONS_PER_ALARM` — Maximum terminal sessions scanned by one automatic cleanup alarm batch (default: `25`)

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
8 changes: 8 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,14 @@ 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
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_ENABLED=true
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TRIGGER_RATIO=0.8
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TARGET_RATIO=0.75
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_ROWS=500
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_BYTES=1048576
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MIN_SESSION_AGE_DAYS=7
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_RECHECK_MS=60000
PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MAX_SESSIONS_PER_ALARM=25

# Development
NODE_ENV=development
Expand Down
9 changes: 9 additions & 0 deletions apps/api/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -801,6 +801,15 @@ INFOMANIAK_IP_POLL_INTERVAL_MS=3000
# PROJECT_COMMENT_LIST_MAX=300 # Max page size for the project-wide comment inbox
# PROJECT_COMMENT_LIST_MAX_BYTES=4000000 # Max estimated content bytes for the project-wide comment inbox
# DOCUMENT_CARD_RAW_OUTPUT_MAX_BYTES=16384
# PROJECT_DATA_TOOL_METADATA_MAX_BYTES=131072 # Max stored tool_metadata bytes for new messages before structured tool content is stripped
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_ENABLED=true # Auto-strip old terminal-session tool_metadata.content when ProjectData storage is under pressure
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TRIGGER_RATIO=0.8 # Start cleanup at this databaseSize / PROJECT_DATA_STORAGE_LIMIT_BYTES ratio
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TARGET_RATIO=0.75 # Continue cleanup until storage is below this ratio or candidates are exhausted
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_ROWS=500 # Max tool-message rows inspected per cleanup alarm batch
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_BYTES=1048576 # Max legacy tool_metadata bytes read into JS per cleanup alarm batch
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MIN_SESSION_AGE_DAYS=7 # Preserve newer terminal sessions from automated tool payload cleanup
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_RECHECK_MS=60000 # Delay before the next cleanup batch when more candidates remain
# PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MAX_SESSIONS_PER_ALARM=25 # Max terminal sessions scanned per cleanup alarm
# MESSAGE_SIZE_THRESHOLD=102400
# ACTIVITY_RETENTION_DAYS=90
# SESSION_IDLE_TIMEOUT_MINUTES=60
Expand Down
7 changes: 1 addition & 6 deletions apps/api/src/durable-objects/project-data/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1018,12 +1018,7 @@ export class ProjectData extends DurableObject<Env> {
if (await deferAlarmWhenDisabled(this.env, this.ctx.storage, 'ProjectData')) return;

try {
await storageSafety.measureAndPersistProjectDataStorage(
this.sql,
this.env,
this.getProjectId(),
'alarm'
);
await storageSafety.runProjectDataStorageSafetyAlarm(this.sql, this.env, this.getProjectId());
} catch (err) {
log.error('alarm.storage_safety_failed', {
error: err instanceof Error ? err.message : String(err),
Expand Down
37 changes: 37 additions & 0 deletions apps/api/src/durable-objects/project-data/storage-safety-meta.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
import { isJsonRecord } from '@simple-agent-manager/shared';

export const META_LAST_MEASURED_AT = 'storageSafetyLastMeasuredAt';
export const META_LAST_STATUS = 'storageSafetyLastStatus';
export const META_LAST_ERROR = 'storageSafetyLastError';

export function readStorageSafetyMeta(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;
const value = (row as Record<string, unknown>).value;
return typeof value === 'string' ? value : null;
}

export function readStorageSafetyMetaNumber(sql: SqlStorage, key: string): number | null {
const raw = readStorageSafetyMeta(sql, key);
if (!raw) return null;
const parsed = Number.parseInt(raw, 10);
return Number.isSafeInteger(parsed) ? parsed : null;
}

export function writeStorageSafetyMeta(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
);
}

export function deleteStorageSafetyMeta(sql: SqlStorage, key: string): void {
sql.exec('DELETE FROM do_meta WHERE key = ?', key);
}

export function truncateStorageSafetyMetaValue(value: string, maxLength: number): string {
return value.length <= maxLength ? value : value.slice(0, maxLength);
}
151 changes: 117 additions & 34 deletions apps/api/src/durable-objects/project-data/storage-safety.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,20 @@ import { isJsonRecord } from '@simple-agent-manager/shared';

import { createModuleLogger, serializeError } from '../../lib/logger';
import { persistError } from '../../services/observability';
import {
META_LAST_ERROR,
META_LAST_MEASURED_AT,
META_LAST_STATUS,
readStorageSafetyMeta as readMeta,
readStorageSafetyMetaNumber as readMetaNumber,
truncateStorageSafetyMetaValue as truncate,
writeStorageSafetyMeta as writeMeta,
} from './storage-safety-meta';
import {
type ProjectDataToolPayloadCleanupResult,
readProjectDataToolPayloadCleanupRecheckAt,
runProjectDataToolPayloadCleanup,
} from './tool-payload-cleanup';
import type { Env } from './types';

const log = createModuleLogger('project_data.storage_safety');
Expand All @@ -34,13 +48,16 @@ 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;
export const DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TRIGGER_RATIO = 0.8;
export const DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TARGET_RATIO = 0.75;
export const DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_ROWS = 500;
export const DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_BYTES = 1024 * 1024;
export const DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MIN_SESSION_AGE_DAYS = 7;
export const DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_RECHECK_MS = 60 * 1000;
export const DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MAX_SESSIONS_PER_ALARM = 25;

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;
Expand Down Expand Up @@ -76,7 +93,12 @@ export interface ProjectDataStorageEmergencyPurgeResult {
exhaustedCandidates: boolean;
}

interface StorageSafetyConfig {
export interface ProjectDataStorageAlarmResult {
measurement: ProjectDataStorageTelemetry | null;
cleanup: ProjectDataToolPayloadCleanupResult | null;
}

export interface StorageSafetyConfig {
enabled: boolean;
limitBytes: number;
measureIntervalMs: number;
Expand All @@ -88,6 +110,14 @@ interface StorageSafetyConfig {
emergencyTargetRatio: number;
emergencyBatchRows: number;
emergencyMaxBatches: number;
toolPayloadCleanupEnabled: boolean;
toolPayloadCleanupTriggerRatio: number;
toolPayloadCleanupTargetRatio: number;
toolPayloadCleanupBatchRows: number;
toolPayloadCleanupBatchBytes: number;
toolPayloadCleanupMinSessionAgeMs: number;
toolPayloadCleanupRecheckMs: number;
toolPayloadCleanupMaxSessionsPerAlarm: number;
}

function parsePositiveInteger(value: string | undefined, fallback: number): number {
Expand All @@ -96,6 +126,12 @@ function parsePositiveInteger(value: string | undefined, fallback: number): numb
return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : fallback;
}

function parseNonNegativeInteger(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);
Expand Down Expand Up @@ -127,6 +163,19 @@ export function resolveStorageSafetyConfig(env: Env): StorageSafetyConfig {

const thresholdsAreOrdered =
noticeRatio < warningRatio && warningRatio < criticalRatio && criticalRatio < degradedRatio;
const cleanupTriggerRatio = parseBoundedRatio(
env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TRIGGER_RATIO,
DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TRIGGER_RATIO
);
const cleanupTargetRatio = parseBoundedRatio(
env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TARGET_RATIO,
DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TARGET_RATIO
);
const cleanupRatiosAreOrdered = cleanupTargetRatio < cleanupTriggerRatio;
const cleanupMinSessionAgeDays = parseNonNegativeInteger(
env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MIN_SESSION_AGE_DAYS,
DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MIN_SESSION_AGE_DAYS
);

return {
enabled: envFlagEnabled(env.PROJECT_DATA_STORAGE_TELEMETRY_ENABLED),
Expand Down Expand Up @@ -162,6 +211,30 @@ export function resolveStorageSafetyConfig(env: Env): StorageSafetyConfig {
env.PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES,
DEFAULT_PROJECT_DATA_STORAGE_EMERGENCY_MAX_BATCHES
),
toolPayloadCleanupEnabled: envFlagEnabled(env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_ENABLED),
toolPayloadCleanupTriggerRatio: cleanupRatiosAreOrdered
? cleanupTriggerRatio
: DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TRIGGER_RATIO,
toolPayloadCleanupTargetRatio: cleanupRatiosAreOrdered
? cleanupTargetRatio
: DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_TARGET_RATIO,
toolPayloadCleanupBatchRows: parsePositiveInteger(
env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_ROWS,
DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_ROWS
),
toolPayloadCleanupBatchBytes: parsePositiveInteger(
env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_BYTES,
DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_BATCH_BYTES
),
toolPayloadCleanupMinSessionAgeMs: cleanupMinSessionAgeDays * 24 * 60 * 60 * 1000,
toolPayloadCleanupRecheckMs: parsePositiveInteger(
env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_RECHECK_MS,
DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_RECHECK_MS
),
toolPayloadCleanupMaxSessionsPerAlarm: parsePositiveInteger(
env.PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MAX_SESSIONS_PER_ALARM,
DEFAULT_PROJECT_DATA_TOOL_PAYLOAD_CLEANUP_MAX_SESSIONS_PER_ALARM
),
};
}

Expand All @@ -180,33 +253,6 @@ export function classifyStorageUsage(
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,
Expand Down Expand Up @@ -358,7 +404,24 @@ export function computeStorageSafetyAlarmTime(
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;
const measureAt = lastMeasuredAt === null ? now : lastMeasuredAt + config.measureIntervalMs;
const cleanupRecheckAt = config.toolPayloadCleanupEnabled
? readProjectDataToolPayloadCleanupRecheckAt(sql)
: null;
if (cleanupRecheckAt === null) return measureAt;
return Math.min(measureAt, cleanupRecheckAt);
}

export function shouldMeasureProjectDataStorage(
sql: SqlStorage,
env: Env,
now: number = Date.now()
): boolean {
const config = resolveStorageSafetyConfig(env);
if (!config.enabled) return false;
if (!readMeta(sql, 'projectId')) return false;
const lastMeasuredAt = readMetaNumber(sql, META_LAST_MEASURED_AT);
return lastMeasuredAt === null || now - lastMeasuredAt >= config.measureIntervalMs;
}

export async function measureAndPersistProjectDataStorage(
Expand Down Expand Up @@ -406,9 +469,29 @@ export async function measureAndPersistProjectDataStorage(
return telemetry;
}

export async function runProjectDataStorageSafetyAlarm(
sql: SqlStorage,
env: Env,
projectId: string | null
): Promise<ProjectDataStorageAlarmResult> {
const now = Date.now();
let measurement: ProjectDataStorageTelemetry | null = null;
if (shouldMeasureProjectDataStorage(sql, env, now)) {
measurement = await measureAndPersistProjectDataStorage(sql, env, projectId, 'alarm');
}
const config = resolveStorageSafetyConfig(env);
const cleanup = await runProjectDataToolPayloadCleanup(sql, env, projectId, config, {
allowStart: measurement !== null,
now,
classifyStatus: (databaseSizeBytes) => classifyStorageUsage(databaseSizeBytes, config),
recordTelemetry: (telemetry, fields) => upsertTelemetry(env, telemetry, fields),
});
return { measurement, cleanup };
}

function normalizeCount(row: unknown): number {
if (!isJsonRecord(row)) return 0;
const count = row.count;
const count = (row as Record<string, unknown>).count;
return typeof count === 'number' && Number.isFinite(count) ? count : 0;
}

Expand Down
66 changes: 66 additions & 0 deletions apps/api/src/durable-objects/project-data/tool-metadata-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,3 +139,69 @@ export function boundToolMetadataForStorage(
truncated: true,
};
}

export function stripToolMetadataPayloadForStorage(
toolMetadata: string | null,
env: Env
): {
value: string | null;
originalBytes: number;
storedBytes: number;
stripped: boolean;
failed: boolean;
} {
if (toolMetadata === null) {
return { value: null, originalBytes: 0, storedBytes: 0, stripped: false, failed: false };
}

const originalBytes = utf8Bytes(toolMetadata);
let parsed: unknown;
try {
parsed = JSON.parse(toolMetadata);
} catch {
return {
value: toolMetadata,
originalBytes,
storedBytes: originalBytes,
stripped: false,
failed: true,
};
}

let compact: unknown;
try {
compact = stripToolMetadataContent(parsed, resolveCompactMessageOptions(env));
} catch {
return {
value: toolMetadata,
originalBytes,
storedBytes: originalBytes,
stripped: false,
failed: true,
};
}
const compactJson = JSON.stringify(compact);
const compactBytes = utf8Bytes(compactJson);
if (compactBytes >= originalBytes) {
return {
value: toolMetadata,
originalBytes,
storedBytes: originalBytes,
stripped: false,
failed: false,
};
}

const maxBytes = resolveToolMetadataMaxBytes(env);
const value = compactBytes <= maxBytes
? compactJson
: serializeWithinLimit(buildMinimalToolMetadata(parsed, originalBytes), maxBytes);

return {
value,
originalBytes,
storedBytes: utf8Bytes(value),
stripped: true,
failed: false,
};
}
Loading
Loading