Files
OmniRoute/src/lib/usage/aggregateHistory.ts
Diego Rodrigues de Sa e Souza 191009dd23 Release v3.8.7 (#2919)
* feat(plugins): WordPress-style plugin system backend

* fix(plugins): address code review feedback

- Path traversal guard: validate entryPoint stays within plugin dir
- install() now handles direct plugin directories (not just parent dirs)
- Non-null assertion replaced with explicit null check
- require efficiency: allowedModules map moved outside function
- Source wrapper: add newlines to prevent trailing comment issues
- Config validation: validate values against configSchema on save
- Dynamic import comment: clarify Node.js caching behavior

Co-Authored-By: OpenClaude (mimo-v2.5-pro) <openclaude@gitlawb.com>

* fix(plugins): replace vm with child_process, add auth to all routes

Addresses all remaining code review feedback:

1. **Loader rewrite**: Replaced Node.js vm module with child_process.fork()
   for proper process-level isolation. Complies with Rule 3 (no eval).
   Each plugin runs in a separate Node.js process with IPC communication.

2. **Auth on all routes**: Added requireManagementAuth to all 6 plugin
   API route files (list, install, scan, details, activate, deactivate, config).

3. **Env filtering**: Only safe env vars passed to plugin processes unless
   "env" permission is granted.

Co-Authored-By: OpenClaude (mimo-v2.5-pro) <openclaude@gitlawb.com>

* fix(plugins): security + ESM fixes for loader and manager

loader.ts:
- Fix IPC: use process.send()/process.on("message") instead of worker_threads.parentPort
- Fix ESM: write host script as .mjs (not .js) to force ESM execution
- Add timeout: 10s default on callHook() with Promise.race
- Add SIGKILL escalation: SIGTERM first, then SIGKILL after 3s grace
- Fix env filtering: use allowlist (safeKeys) instead of passing all env vars
- Clear timeout on successful IPC response (no timer leak)

manager.ts:
- Fix path traversal: use fs.realpath() instead of startsWith()
- Fix imports: use registerHook/unregisterHooks from hooks.ts
- Register hooks individually via registerHook(event, name, handler)

hooks.ts:
- Copied from feat/plugin-custom-hooks (canonical registry)

* feat(discovery): add discovery tool stub service

Phase 1 scaffold for automated provider discovery:
- DiscoveryConfig, DiscoveryResult types
- probeEndpoint() for URL availability checking
- scanProvider() stub (Phase 2 will implement real scanning)
- getDiscoveryResults() stub
- Default config: disabled (opt-in)

* chore(plugins): slop cleanup — pino logger, remove redundant sorts

- index.ts: replace console.log/error with pino structured logging
- hooks.ts: remove redundant .sort() in emitHookBlocking/runOnResponse (already sorted on registration)
- manager.ts: add readFile import

* test(plugins): add scanner, loader, manager unit tests

- scanner: 9 tests (discovery, hidden dirs, validation, entry point, multiple)
- loader: 5 tests (type contracts, Plugin/PluginContext/PluginResult interfaces)
- manager: 6 tests (singleton, lifecycle methods, error on unknown)
- Total: 20 tests, all passing

* fix(settings): add missing home page pin keys to updateSettingsSchema

* feat(plugins): add i18n keys to all 42 locales

* fix(settings): add missing security keys to updateSettingsSchema and add tests

* fix(usage): analytics route reads combo_name/requested_model from call_logs only

The 3.8.6 variant of #2904 added SELECTs of combo_name/requested_model
against usage_history, but those columns only exist in call_logs (no
migration adds them to usage_history). This returned HTTP 500 on
/api/usage/analytics. Restore the working query shape from the 3.8.7
variant. Fixes 18 failing usage-analytics-route tests.

* fix(types,test): resolve noImplicitAny in progressiveAging + align semaphore test to #2903 gate pruning

- progressiveAging: type compression results so messages[0].content is
  indexable (was TS7053 against {}); restores typecheck:noimplicit:core gate.
- services-branch-hardening: #2903 (perf-ram) prunes idle rate-limit gates
  on zero; assert no-running/empty-queue without assuming the entry persists.

* fix(analytics): address merged review regressions

* fix(executor): normalize max effort for openai shape providers

* Make zero-latency combo optimizations opt-in

* Address zero-latency combo review feedback

* chore(release): sync v3.8.7 touchpoints + credit contributors

- llm.txt → 3.8.7 (Current version + Key Features header)
- CHANGELOG: add Dmitry Kuznetsov & Nikolay Alafuzov to 3.8.6 Hall of Contributors
- version already 3.8.7 across package.json/open-sse/electron/openapi (from #2909)

* fix(cleanup): restore usage history cutoff boundary

* docs(changelog): rank 3.8.6 contributors in a commits table with their PRs

* fix(dashboard): theme ReactFlow Controls +/- buttons for dark mode

* fix(settings): add missing home page pin keys to updateSettingsSchema

* fix(settings): add missing security keys to updateSettingsSchema and add tests

* fix(executor): normalize max effort for openai shape providers

* Make zero-latency combo optimizations opt-in

* Address zero-latency combo review feedback

* fix(analytics): address merged review regressions

* fix(cleanup): restore usage history cutoff boundary

* feat(plugins): WordPress-style plugin system backend

* fix(plugins): address code review feedback

- Path traversal guard: validate entryPoint stays within plugin dir
- install() now handles direct plugin directories (not just parent dirs)
- Non-null assertion replaced with explicit null check
- require efficiency: allowedModules map moved outside function
- Source wrapper: add newlines to prevent trailing comment issues
- Config validation: validate values against configSchema on save
- Dynamic import comment: clarify Node.js caching behavior

Co-Authored-By: OpenClaude (mimo-v2.5-pro) <openclaude@gitlawb.com>

* fix(plugins): replace vm with child_process, add auth to all routes

Addresses all remaining code review feedback:

1. **Loader rewrite**: Replaced Node.js vm module with child_process.fork()
   for proper process-level isolation. Complies with Rule 3 (no eval).
   Each plugin runs in a separate Node.js process with IPC communication.

2. **Auth on all routes**: Added requireManagementAuth to all 6 plugin
   API route files (list, install, scan, details, activate, deactivate, config).

3. **Env filtering**: Only safe env vars passed to plugin processes unless
   "env" permission is granted.

Co-Authored-By: OpenClaude (mimo-v2.5-pro) <openclaude@gitlawb.com>

* fix(plugins): security + ESM fixes for loader and manager

loader.ts:
- Fix IPC: use process.send()/process.on("message") instead of worker_threads.parentPort
- Fix ESM: write host script as .mjs (not .js) to force ESM execution
- Add timeout: 10s default on callHook() with Promise.race
- Add SIGKILL escalation: SIGTERM first, then SIGKILL after 3s grace
- Fix env filtering: use allowlist (safeKeys) instead of passing all env vars
- Clear timeout on successful IPC response (no timer leak)

manager.ts:
- Fix path traversal: use fs.realpath() instead of startsWith()
- Fix imports: use registerHook/unregisterHooks from hooks.ts
- Register hooks individually via registerHook(event, name, handler)

hooks.ts:
- Copied from feat/plugin-custom-hooks (canonical registry)

* feat(discovery): add discovery tool stub service

Phase 1 scaffold for automated provider discovery:
- DiscoveryConfig, DiscoveryResult types
- probeEndpoint() for URL availability checking
- scanProvider() stub (Phase 2 will implement real scanning)
- getDiscoveryResults() stub
- Default config: disabled (opt-in)

* chore(plugins): slop cleanup — pino logger, remove redundant sorts

- index.ts: replace console.log/error with pino structured logging
- hooks.ts: remove redundant .sort() in emitHookBlocking/runOnResponse (already sorted on registration)
- manager.ts: add readFile import

* test(plugins): add scanner, loader, manager unit tests

- scanner: 9 tests (discovery, hidden dirs, validation, entry point, multiple)
- loader: 5 tests (type contracts, Plugin/PluginContext/PluginResult interfaces)
- manager: 6 tests (singleton, lifecycle methods, error on unknown)
- Total: 20 tests, all passing

* feat(plugins): add i18n keys to all 42 locales

* chore(plugins): remove duplicate migration 059_create_plugins.sql

* chore(plugins): remove duplicate migration 059_create_plugins.sql (post-merge)

* fix(sse): guard non-string error.code in proxyFetch + harden model parsing (#2463) (#2923)

Integrated into release/v3.8.7

* fix(docker): add runner-web stage with Playwright Chromium (#2832) (#2846)

Integrated into release/v3.8.7

* docs(changelog): document NVIDIA NIM and error code type-crash fix (#2463)

* test: ignore NVIDIA_BASE_URL and NVIDIA_MODEL in env contract check

---------

Co-authored-by: oyi77 <oyi77@users.noreply.github.com>
Co-authored-by: OpenClaude (mimo-v2.5-pro) <openclaude@gitlawb.com>
Co-authored-by: Apostol Apostolov <theapoapostolov@gmail.com>
Co-authored-by: Halil Tezcan KARABULUT <info@hlltzcnkb.com>
Co-authored-by: R.D. <rogerproself@gmail.com>
2026-05-29 19:54:00 -03:00

213 lines
7.3 KiB
TypeScript

/**
* Aggregation utility functions for usage data summarization.
* Rolls up usage_history (and quota_snapshots) into daily summary tables.
*
* @module lib/usage/aggregateHistory
*/
import { getDbInstance } from "../db/core";
import { getUserDatabaseSettings } from "../db/databaseSettings";
interface AggregationResult {
processed: number;
inserted: number;
errors: number;
}
/**
* Roll up quota_snapshots into daily_usage_summary table.
* Aggregates by provider, model, and date.
*
* @param fromDate - Start date (YYYY-MM-DD format)
* @param toDate - End date (YYYY-MM-DD format)
* @returns Aggregation result with counts
*/
export async function rollupDailyUsage(
fromDate: string,
toDate: string
): Promise<AggregationResult> {
const db = getDbInstance();
const result: AggregationResult = {
processed: 0,
inserted: 0,
errors: 0,
};
try {
// Aggregate quota_snapshots by provider, model, and date
const aggregateQuery = `
INSERT INTO daily_usage_summary (provider, model, date, total_requests, total_input_tokens, total_output_tokens, total_cost)
SELECT
provider,
COALESCE(json_extract(raw_data, '$.model'), 'unknown') as model,
DATE(created_at) as date,
COUNT(*) as total_requests,
COALESCE(SUM(CAST(json_extract(raw_data, '$.input_tokens') AS INTEGER)), 0) as total_input_tokens,
COALESCE(SUM(CAST(json_extract(raw_data, '$.output_tokens') AS INTEGER)), 0) as total_output_tokens,
COALESCE(SUM(CAST(json_extract(raw_data, '$.cost') AS REAL)), 0.0) as total_cost
FROM quota_snapshots
WHERE DATE(created_at) >= ? AND DATE(created_at) <= ?
GROUP BY provider, model, DATE(created_at)
ON CONFLICT(provider, model, date) DO UPDATE SET
total_requests = excluded.total_requests,
total_input_tokens = excluded.total_input_tokens,
total_output_tokens = excluded.total_output_tokens,
total_cost = excluded.total_cost
`;
const stmt = db.prepare(aggregateQuery);
const runResult = stmt.run(fromDate, toDate);
result.processed = runResult.changes;
result.inserted = runResult.changes;
console.log(`[Aggregation] Daily rollup: ${result.inserted} rows for ${fromDate} to ${toDate}`);
} catch (err: any) {
console.error("[Aggregation] Daily rollup error:", err);
result.errors++;
}
return result;
}
/**
* Roll up quota_snapshots into hourly_usage_summary table.
* Aggregates by provider, model, and hour.
*
* @param fromDate - Start datetime (YYYY-MM-DD HH:MM:SS format)
* @param toDate - End datetime (YYYY-MM-DD HH:MM:SS format)
* @returns Aggregation result with counts
*/
export async function rollupHourlyQuota(
fromDate: string,
toDate: string
): Promise<AggregationResult> {
const db = getDbInstance();
const result: AggregationResult = {
processed: 0,
inserted: 0,
errors: 0,
};
try {
// Aggregate quota_snapshots by provider, model, and hour
const aggregateQuery = `
INSERT INTO hourly_usage_summary (provider, model, date_hour, total_requests, total_input_tokens, total_output_tokens, total_cost)
SELECT
provider,
COALESCE(json_extract(raw_data, '$.model'), 'unknown') as model,
datetime(strftime('%Y-%m-%d %H:00:00', created_at)) as date_hour,
COUNT(*) as total_requests,
COALESCE(SUM(CAST(json_extract(raw_data, '$.input_tokens') AS INTEGER)), 0) as total_input_tokens,
COALESCE(SUM(CAST(json_extract(raw_data, '$.output_tokens') AS INTEGER)), 0) as total_output_tokens,
COALESCE(SUM(CAST(json_extract(raw_data, '$.cost') AS REAL)), 0.0) as total_cost
FROM quota_snapshots
WHERE created_at >= ? AND created_at <= ?
GROUP BY provider, model, datetime(strftime('%Y-%m-%d %H:00:00', created_at))
ON CONFLICT(provider, model, date_hour) DO UPDATE SET
total_requests = excluded.total_requests,
total_input_tokens = excluded.total_input_tokens,
total_output_tokens = excluded.total_output_tokens,
total_cost = excluded.total_cost
`;
const stmt = db.prepare(aggregateQuery);
const runResult = stmt.run(fromDate, toDate);
result.processed = runResult.changes;
result.inserted = runResult.changes;
console.log(
`[Aggregation] Hourly rollup: ${result.inserted} rows for ${fromDate} to ${toDate}`
);
} catch (err: any) {
console.error("[Aggregation] Hourly rollup error:", err);
result.errors++;
}
return result;
}
/**
* Roll up usage_history into daily_usage_summary before raw rows are deleted.
* This is the authoritative rollup — sourced from actual per-request token data,
* not from quota_snapshots. Should be called before cleanupUsageHistory() deletes rows.
*
* The ON CONFLICT clause uses SUM so re-running is additive-safe: if a date already
* has a partial rollup (e.g. from a previous partial cleanup), new rows accumulate.
*
* @param beforeDate - ISO timestamp/date boundary. Rows strictly before this value are rolled up.
* @returns Aggregation result with counts
*/
export async function rollupUsageHistoryBeforeDate(beforeDate: string): Promise<AggregationResult> {
const db = getDbInstance();
const result: AggregationResult = {
processed: 0,
inserted: 0,
errors: 0,
};
try {
const aggregateQuery = `
INSERT INTO daily_usage_summary (provider, model, date, total_requests, total_input_tokens, total_output_tokens, total_cost)
SELECT
LOWER(provider) as provider,
LOWER(model) as model,
DATE(timestamp) as date,
COUNT(*) as total_requests,
COALESCE(SUM(tokens_input), 0) as total_input_tokens,
COALESCE(SUM(tokens_output), 0) as total_output_tokens,
0.0 as total_cost
FROM usage_history
WHERE timestamp < ?
AND provider IS NOT NULL AND provider != ''
AND model IS NOT NULL AND model != ''
GROUP BY LOWER(provider), LOWER(model), DATE(timestamp)
ON CONFLICT(provider, model, date) DO UPDATE SET
total_requests = daily_usage_summary.total_requests + excluded.total_requests,
total_input_tokens = daily_usage_summary.total_input_tokens + excluded.total_input_tokens,
total_output_tokens = daily_usage_summary.total_output_tokens + excluded.total_output_tokens
`;
const stmt = db.prepare(aggregateQuery);
const runResult = stmt.run(beforeDate);
result.processed = runResult.changes;
result.inserted = runResult.changes;
console.log(
`[Aggregation] usage_history rollup: ${result.inserted} rows for dates before ${beforeDate}`
);
} catch (err: any) {
console.error("[Aggregation] usage_history rollup error:", err);
result.errors++;
}
return result;
}
/**
* Get the cutoff date for raw data based on retention settings.
* Data older than this should be aggregated and cleaned up.
*
* @returns ISO date string (YYYY-MM-DD)
*/
export async function getRawDataCutoffDate(): Promise<string> {
const rawDataRetentionDays = getUserDatabaseSettings().aggregation.rawDataRetentionDays;
const cutoffDate = new Date();
cutoffDate.setDate(cutoffDate.getDate() - rawDataRetentionDays);
return cutoffDate.toISOString().split("T")[0];
}
/**
* Check if aggregation is enabled in settings.
*/
export async function isAggregationEnabled(): Promise<boolean> {
return getUserDatabaseSettings().aggregation.enabled;
}