mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-07-26 09:52:11 +03:00
feat(db): add quotaConsumption sliding-window counter storage (B/F2)
Implements getBucket, incrementBucket (atomic UPSERT), getPair (curr+prev for sliding window formula), and gcOlderThan (stale bucket cleanup). Atomic increment uses INSERT ... ON CONFLICT DO UPDATE — no separate read-modify-write cycle needed.
This commit is contained in:
141
src/lib/db/quotaConsumption.ts
Normal file
141
src/lib/db/quotaConsumption.ts
Normal file
@@ -0,0 +1,141 @@
|
||||
/**
|
||||
* db/quotaConsumption.ts — Sliding Window Counter primitives for quota tracking.
|
||||
*
|
||||
* Implements low-level bucket read/write operations for the 2-bucket sliding
|
||||
* window counter algorithm. Each row is keyed on (api_key_id, dimension_key,
|
||||
* bucket_index) where dimension_key = "<poolId>:<unit>:<window>" and
|
||||
* bucket_index = floor(now_ms / window_ms).
|
||||
*
|
||||
* Atomicity: incrementBucket uses INSERT ... ON CONFLICT DO UPDATE (UPSERT)
|
||||
* which is a single atomic SQLite statement — no separate read-modify-write.
|
||||
*
|
||||
* Part of: Group B — Quota Sharing Engine (plan 22, frente F2).
|
||||
*/
|
||||
|
||||
import { getDbInstance } from "./core";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Internal helpers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
interface StatementLike<TRow = unknown> {
|
||||
all: (...params: unknown[]) => TRow[];
|
||||
get: (...params: unknown[]) => TRow | undefined;
|
||||
run: (...params: unknown[]) => { changes: number };
|
||||
}
|
||||
|
||||
interface DbLike {
|
||||
prepare: <TRow = unknown>(sql: string) => StatementLike<TRow>;
|
||||
}
|
||||
|
||||
function getDb(): DbLike {
|
||||
return getDbInstance() as unknown as DbLike;
|
||||
}
|
||||
|
||||
interface BucketRow {
|
||||
consumed: number;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Public API
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Read the consumed value for a single bucket. Returns 0 if no row exists.
|
||||
*/
|
||||
export function getBucket(
|
||||
apiKeyId: string,
|
||||
dimensionKey: string,
|
||||
bucketIndex: number
|
||||
): number {
|
||||
const row = getDb()
|
||||
.prepare<BucketRow>(
|
||||
`SELECT consumed FROM quota_consumption
|
||||
WHERE api_key_id = ? AND dimension_key = ? AND bucket_index = ?`
|
||||
)
|
||||
.get(apiKeyId, dimensionKey, bucketIndex);
|
||||
return row?.consumed ?? 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Atomically increment the consumed counter for a bucket.
|
||||
* Uses UPSERT: if the row does not exist it is created; if it exists the
|
||||
* delta is added to the existing consumed value and updated_at is refreshed.
|
||||
*
|
||||
* @param apiKeyId The API key being tracked.
|
||||
* @param dimensionKey "<poolId>:<unit>:<window>" string.
|
||||
* @param bucketIndex floor(nowMs / windowMs).
|
||||
* @param delta Amount to add (positive number).
|
||||
* @param nowMs Current epoch milliseconds (used for updated_at).
|
||||
*/
|
||||
export function incrementBucket(
|
||||
apiKeyId: string,
|
||||
dimensionKey: string,
|
||||
bucketIndex: number,
|
||||
delta: number,
|
||||
nowMs: number
|
||||
): void {
|
||||
getDb()
|
||||
.prepare(
|
||||
`INSERT INTO quota_consumption (api_key_id, dimension_key, bucket_index, consumed, updated_at)
|
||||
VALUES (?, ?, ?, ?, ?)
|
||||
ON CONFLICT(api_key_id, dimension_key, bucket_index)
|
||||
DO UPDATE SET
|
||||
consumed = consumed + excluded.consumed,
|
||||
updated_at = excluded.updated_at`
|
||||
)
|
||||
.run(apiKeyId, dimensionKey, bucketIndex, delta, nowMs);
|
||||
}
|
||||
|
||||
/**
|
||||
* Read the current and previous bucket values for the sliding window formula:
|
||||
* effective = prev × (1 − elapsed/window) + curr
|
||||
*
|
||||
* @param apiKeyId The API key being tracked.
|
||||
* @param dimensionKey "<poolId>:<unit>:<window>" string.
|
||||
* @param currentBucket The current bucket index (floor(nowMs / windowMs)).
|
||||
* @returns { curr, prev } — both default to 0 when row is absent.
|
||||
*/
|
||||
export function getPair(
|
||||
apiKeyId: string,
|
||||
dimensionKey: string,
|
||||
currentBucket: number
|
||||
): { curr: number; prev: number } {
|
||||
const prevBucket = currentBucket - 1;
|
||||
|
||||
const currRow = getDb()
|
||||
.prepare<BucketRow>(
|
||||
`SELECT consumed FROM quota_consumption
|
||||
WHERE api_key_id = ? AND dimension_key = ? AND bucket_index = ?`
|
||||
)
|
||||
.get(apiKeyId, dimensionKey, currentBucket);
|
||||
|
||||
const prevRow = getDb()
|
||||
.prepare<BucketRow>(
|
||||
`SELECT consumed FROM quota_consumption
|
||||
WHERE api_key_id = ? AND dimension_key = ? AND bucket_index = ?`
|
||||
)
|
||||
.get(apiKeyId, dimensionKey, prevBucket);
|
||||
|
||||
return {
|
||||
curr: currRow?.consumed ?? 0,
|
||||
prev: prevRow?.consumed ?? 0,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete rows whose updated_at is strictly less than maxUpdatedAtMs.
|
||||
* Used by GC background job to clean up stale bucket rows.
|
||||
*
|
||||
* Boundary semantics: rows with updated_at === maxUpdatedAtMs are KEPT.
|
||||
* Only rows with updated_at < maxUpdatedAtMs (strictly older) are deleted.
|
||||
*
|
||||
* @param maxUpdatedAtMs Epoch ms threshold (exclusive lower bound for kept rows).
|
||||
* @returns Number of rows deleted.
|
||||
*/
|
||||
export function gcOlderThan(maxUpdatedAtMs: number): number {
|
||||
const result = getDb()
|
||||
.prepare("DELETE FROM quota_consumption WHERE updated_at < ?")
|
||||
.run(maxUpdatedAtMs);
|
||||
return result.changes;
|
||||
}
|
||||
Reference in New Issue
Block a user