Files

550 lines
25 KiB
TypeScript

import { createHash, randomUUID } from "node:crypto";
import { createRequire } from "node:module";
import type BetterSqlite3 from "better-sqlite3";
import {
auditRetentionMilliseconds,
ensureAdminOperationAuditSchema,
isSafeAuditRef,
isSafeAuditSummaryJson,
serializeAuditSummary,
} from "./audit-policy.js";
import { CreditError } from "./credit-errors.js";
export { CreditError } from "./credit-errors.js";
export type { CreditErrorCode } from "./credit-errors.js";
const require = createRequire(import.meta.url);
const Database = require("better-sqlite3") as typeof BetterSqlite3;
interface CreditAccountRow {
available_balance: number;
reserved_balance: number;
updated_at: number;
user_id: string;
}
interface LedgerRow {
amount: number;
available_after: number;
available_before: number;
created_at: number;
entry_status: "succeeded" | "frozen" | "committed" | "released";
entry_type: "registration_grant" | "generation_reserve" | "generation_commit" | "generation_release" | "admin_adjustment";
ledger_id: string;
model_id: string | null;
reason: string | null;
reference_id: string | null;
reference_type: "registration" | "generation" | "admin_adjustment" | null;
reserved_after: number;
reserved_before: number;
user_id: string;
}
interface ReservationRow {
amount: number;
generation_id: string;
model_id: string;
status: "reserved" | "committed" | "released";
user_id: string;
}
function iso(timestamp: number) {
return new Date(timestamp).toISOString();
}
function encodeCursor(row: Pick<LedgerRow, "created_at" | "ledger_id">) {
return Buffer.from(JSON.stringify([row.created_at, row.ledger_id]), "utf8").toString("base64url");
}
function decodeCursor(cursor: string | undefined) {
if (!cursor) return undefined;
try {
const value: unknown = JSON.parse(Buffer.from(cursor, "base64url").toString("utf8"));
if (!Array.isArray(value) || value.length !== 2 || !Number.isSafeInteger(value[0]) || typeof value[1] !== "string") {
throw new Error("cursor_invalid");
}
return { createdAt: value[0] as number, ledgerId: value[1] };
} catch {
throw new CreditError("credit_request_invalid");
}
}
export class CreditService {
readonly database: BetterSqlite3.Database;
private readonly clock: () => number;
constructor(input: { clock?: () => number; databasePath: string }) {
this.clock = input.clock ?? Date.now;
const nativeBinding = process.env.DADA_SQLITE_NATIVE_BINDING;
this.database = new Database(input.databasePath, nativeBinding ? { nativeBinding } : undefined);
this.database.pragma("journal_mode = WAL");
this.database.pragma("foreign_keys = ON");
this.database.pragma("synchronous = FULL");
this.database.pragma("busy_timeout = 5000");
this.database.function("dada_audit_ref_is_safe", { deterministic: true }, isSafeAuditRef);
this.database.function("dada_audit_summary_is_safe", { deterministic: true }, isSafeAuditSummaryJson);
this.database.function("dada_allow_retention_purge", { deterministic: false }, () => 0);
this.database.function("dada_retention_purge_now", { deterministic: false }, () => 0);
this.database.function("dada_allow_privacy_purge", { deterministic: false }, () => 0);
this.migrate();
}
close() {
this.database.close();
}
readAccount(userId: string) {
const row = this.database.prepare("SELECT * FROM credit_accounts WHERE user_id = ?").get(userId) as CreditAccountRow | undefined;
if (!row) throw new CreditError("credit_account_not_found");
return {
availableBalance: row.available_balance,
reservedBalance: row.reserved_balance,
updatedAt: iso(row.updated_at),
};
}
listLedger(input: {
cursor?: string;
eventType?: LedgerRow["entry_type"];
from?: string;
limit?: number;
to?: string;
userId: string;
}) {
const account = this.readAccount(input.userId);
const limit = input.limit ?? 20;
if (!Number.isSafeInteger(limit) || limit < 1 || limit > 100) throw new CreditError("credit_request_invalid");
const cursor = decodeCursor(input.cursor);
const conditions = ["user_id = ?"];
const values: Array<number | string> = [input.userId];
if (input.eventType) {
conditions.push("entry_type = ?");
values.push(input.eventType);
}
if (input.from) {
const from = Date.parse(input.from);
if (!Number.isFinite(from)) throw new CreditError("credit_request_invalid");
conditions.push("created_at >= ?");
values.push(from);
}
if (input.to) {
const to = Date.parse(input.to);
if (!Number.isFinite(to)) throw new CreditError("credit_request_invalid");
conditions.push("created_at <= ?");
values.push(to);
}
if (cursor) {
conditions.push("(created_at < ? OR (created_at = ? AND ledger_id < ?))");
values.push(cursor.createdAt, cursor.createdAt, cursor.ledgerId);
}
const rows = this.database.prepare(`
SELECT * FROM credit_ledger WHERE ${conditions.join(" AND ")}
ORDER BY created_at DESC, ledger_id DESC LIMIT ?
`).all(...values, limit + 1) as LedgerRow[];
const hasMore = rows.length > limit;
const page = rows.slice(0, limit);
return {
account,
entries: page.map((row) => ({
amount: row.amount,
availableAfter: row.available_after,
availableBefore: row.available_before,
createdAt: iso(row.created_at),
entryId: row.ledger_id,
entryType: row.entry_type,
modelId: row.model_id,
reason: row.reason,
referenceId: row.reference_id,
referenceType: row.reference_type,
reservedAfter: row.reserved_after,
reservedBefore: row.reserved_before,
status: row.entry_status,
})),
nextCursor: hasMore && page.length > 0 ? encodeCursor(page.at(-1)!) : null,
};
}
reserveGeneration(input: {
creditCost: number;
generationId: string;
modelId: string;
operationKey: string;
userId: string;
}) {
if (!Number.isSafeInteger(input.creditCost) || input.creditCost <= 0 || !input.modelId || !input.operationKey) {
throw new CreditError("credit_request_invalid");
}
return this.immediate(() => {
const replay = this.database.prepare("SELECT * FROM credit_ledger WHERE operation_key = ?")
.get(input.operationKey) as LedgerRow | undefined;
if (replay) {
if (replay.user_id !== input.userId || replay.reference_id !== input.generationId
|| replay.entry_type !== "generation_reserve" || replay.amount !== -input.creditCost
|| replay.model_id !== input.modelId) {
throw new CreditError("credit_operation_conflict");
}
return { availableBalance: replay.available_after, reservedBalance: replay.reserved_after, status: "reserved" as const };
}
const account = this.database.prepare("SELECT * FROM credit_accounts WHERE user_id = ?").get(input.userId) as CreditAccountRow | undefined;
if (!account) throw new CreditError("credit_account_not_found");
if (account.available_balance < input.creditCost) throw new CreditError("credit_insufficient");
if (!this.tableExists("generation_jobs")) throw new CreditError("credit_generation_not_found");
const generation = this.database.prepare(`
SELECT generation_id FROM generation_jobs
WHERE generation_id = ? AND owner_id = ? AND status IN ('queued', 'running')
`).get(input.generationId, input.userId);
if (!generation) throw new CreditError("credit_generation_not_found");
const existing = this.database.prepare("SELECT * FROM credit_reservations WHERE generation_id = ?")
.get(input.generationId) as ReservationRow | undefined;
if (existing) throw new CreditError("credit_operation_conflict");
const availableAfter = account.available_balance - input.creditCost;
const reservedAfter = account.reserved_balance + input.creditCost;
const now = this.clock();
const changed = this.database.prepare(`
UPDATE credit_accounts SET available_balance = ?, reserved_balance = ?, updated_at = ?
WHERE user_id = ? AND available_balance = ? AND reserved_balance = ?
`).run(availableAfter, reservedAfter, now, input.userId, account.available_balance, account.reserved_balance);
if (changed.changes !== 1) throw new CreditError("credit_invariant_failed");
this.database.prepare(`
INSERT INTO credit_reservations (generation_id, user_id, model_id, amount, status, created_at, finalized_at)
VALUES (?, ?, ?, ?, 'reserved', ?, NULL)
`).run(input.generationId, input.userId, input.modelId, input.creditCost, now);
this.database.prepare(`
UPDATE generation_jobs SET model_id = ?, confirmed_credit_cost = ?, reserved_credits = ?, final_credit_state = NULL
WHERE generation_id = ?
`).run(input.modelId, input.creditCost, input.creditCost, input.generationId);
this.insertLedger({
amount: -input.creditCost,
availableAfter,
availableBefore: account.available_balance,
createdAt: now,
entryStatus: "frozen",
entryType: "generation_reserve",
modelId: input.modelId,
operationKey: input.operationKey,
reason: null,
referenceId: input.generationId,
referenceType: "generation",
reservedAfter,
reservedBefore: account.reserved_balance,
userId: input.userId,
});
this.insertOutbox("generation_credit_reserved", input.generationId, input.operationKey, { amount: input.creditCost, user_id: input.userId }, now);
return { availableBalance: availableAfter, reservedBalance: reservedAfter, status: "reserved" as const };
});
}
finalizeGeneration(input: {
generationId: string;
operationKey: string;
outcome: "succeeded" | "failed" | "rejected";
}) {
if (!input.operationKey) throw new CreditError("credit_request_invalid");
return this.immediate(() => {
const reservation = this.database.prepare("SELECT * FROM credit_reservations WHERE generation_id = ?")
.get(input.generationId) as ReservationRow | undefined;
if (!reservation) throw new CreditError("credit_generation_not_found");
if (reservation.status !== "reserved") {
const replay = this.database.prepare(`
SELECT * FROM credit_ledger WHERE reference_id = ? AND entry_type IN ('generation_commit', 'generation_release')
ORDER BY created_at DESC, ledger_id DESC LIMIT 1
`).get(input.generationId) as LedgerRow | undefined;
if (!replay) throw new CreditError("credit_invariant_failed");
return {
availableBalance: replay.available_after,
creditState: reservation.status,
reservedBalance: replay.reserved_after,
status: "finalized" as const,
};
}
const operation = this.database.prepare("SELECT * FROM credit_ledger WHERE operation_key = ?").get(input.operationKey) as LedgerRow | undefined;
if (operation) throw new CreditError("credit_operation_conflict");
const account = this.database.prepare("SELECT * FROM credit_accounts WHERE user_id = ?").get(reservation.user_id) as CreditAccountRow | undefined;
if (!account || account.reserved_balance < reservation.amount) throw new CreditError("credit_invariant_failed");
const committed = input.outcome === "succeeded";
const creditState = committed ? "committed" as const : "released" as const;
const availableAfter = committed ? account.available_balance : account.available_balance + reservation.amount;
const reservedAfter = account.reserved_balance - reservation.amount;
if (!Number.isSafeInteger(availableAfter)) throw new CreditError("credit_invariant_failed");
const now = this.clock();
this.database.prepare(`
UPDATE credit_accounts SET available_balance = ?, reserved_balance = ?, updated_at = ? WHERE user_id = ?
`).run(availableAfter, reservedAfter, now, reservation.user_id);
this.database.prepare(`
UPDATE credit_reservations SET status = ?, finalized_at = ? WHERE generation_id = ? AND status = 'reserved'
`).run(creditState, now, input.generationId);
this.database.prepare("UPDATE generation_jobs SET final_credit_state = ? WHERE generation_id = ?")
.run(creditState, input.generationId);
this.insertLedger({
amount: committed ? -reservation.amount : reservation.amount,
availableAfter,
availableBefore: account.available_balance,
createdAt: now,
entryStatus: committed ? "committed" : "released",
entryType: committed ? "generation_commit" : "generation_release",
modelId: reservation.model_id,
operationKey: input.operationKey,
reason: null,
referenceId: input.generationId,
referenceType: "generation",
reservedAfter,
reservedBefore: account.reserved_balance,
userId: reservation.user_id,
});
this.insertOutbox(committed ? "generation_credit_committed" : "generation_credit_released", input.generationId, input.operationKey, { outcome: input.outcome }, now);
return { availableBalance: availableAfter, creditState, reservedBalance: reservedAfter, status: "finalized" as const };
});
}
adjustAvailable(input: {
adjustmentId: string;
adminId: string;
amount: number;
idempotencyKey: string;
reason: string;
userId: string;
}) {
const reason = input.reason.trim();
if (!Number.isSafeInteger(input.amount) || input.amount === 0 || !reason || reason.length > 500
|| input.idempotencyKey.length < 32 || input.idempotencyKey.length > 200
|| !/^[A-Za-z0-9_-]+$/.test(input.idempotencyKey)) {
throw new CreditError("credit_request_invalid");
}
const idempotencyKeyDigest = createHash("sha256").update(input.idempotencyKey, "utf8").digest("hex");
const requestHash = createHash("sha256")
.update(JSON.stringify([input.adjustmentId, input.adminId, input.userId, input.amount, reason]), "utf8")
.digest("hex");
return this.immediate(() => {
type Receipt = {
adjustment_id: string;
admin_id: string | null;
available_after: number;
idempotency_key_digest: string | null;
request_hash: string;
reserved_after: number;
user_id: string | null;
};
const receiptByKey = this.database.prepare(`
SELECT * FROM credit_adjustment_receipts
WHERE admin_id = ? AND user_id = ? AND idempotency_key_digest = ?
`).get(input.adminId, input.userId, idempotencyKeyDigest) as Receipt | undefined;
if (receiptByKey) {
if (receiptByKey.request_hash !== requestHash) throw new CreditError("credit_operation_conflict");
return {
adjustmentId: receiptByKey.adjustment_id,
availableBalance: receiptByKey.available_after,
reservedBalance: receiptByKey.reserved_after,
status: "adjusted" as const,
};
}
const receipt = this.database.prepare("SELECT * FROM credit_adjustment_receipts WHERE adjustment_id = ?")
.get(input.adjustmentId) as Receipt | undefined;
if (receipt) {
if (receipt.request_hash !== requestHash || receipt.admin_id !== input.adminId || receipt.user_id !== input.userId
|| receipt.idempotency_key_digest !== idempotencyKeyDigest) {
throw new CreditError("credit_operation_conflict");
}
return {
adjustmentId: input.adjustmentId,
availableBalance: receipt.available_after,
reservedBalance: receipt.reserved_after,
status: "adjusted" as const,
};
}
const admin = this.database.prepare(`
SELECT u.user_id FROM users u JOIN admin_access a ON a.user_id = u.user_id
WHERE u.user_id = ? AND u.role = 'super_admin' AND u.status = 'active' AND a.allowed = 1
`).get(input.adminId);
if (!admin) throw new CreditError("credit_account_not_found");
const target = this.database.prepare(`
SELECT c.* FROM credit_accounts c JOIN users u ON u.user_id = c.user_id
WHERE c.user_id = ? AND u.role = 'user' AND u.status <> 'deleted'
`).get(input.userId) as CreditAccountRow | undefined;
if (!target) throw new CreditError("credit_account_not_found");
const availableAfter = target.available_balance + input.amount;
if (!Number.isSafeInteger(availableAfter)) throw new CreditError("credit_request_invalid");
const now = this.clock();
this.database.prepare("UPDATE credit_accounts SET available_balance = ?, updated_at = ? WHERE user_id = ?")
.run(availableAfter, now, input.userId);
const ledgerId = this.insertLedger({
amount: input.amount,
availableAfter,
availableBefore: target.available_balance,
createdAt: now,
entryStatus: "succeeded",
entryType: "admin_adjustment",
modelId: null,
operationKey: `admin_adjustment:${createHash("sha256").update(`${input.adminId}\0${input.userId}\0${idempotencyKeyDigest}`, "utf8").digest("hex")}`,
reason,
referenceId: input.adjustmentId,
referenceType: "admin_adjustment",
reservedAfter: target.reserved_balance,
reservedBefore: target.reserved_balance,
userId: input.userId,
});
this.database.prepare(`
INSERT INTO admin_operation_logs (
log_id, actor_type, actor_ref, operation_type, target_type, target_ref,
result, before_summary, after_summary, occurred_at, expires_at
) VALUES (?, 'super_admin', ?, 'credit_adjustment', 'user_credit_account', ?, 'succeeded', ?, ?, ?, ?)
`).run(
randomUUID(), input.adminId, input.userId,
serializeAuditSummary({ available_balance: target.available_balance, reserved_balance: target.reserved_balance }),
serializeAuditSummary({ adjustment_amount: input.amount, available_balance: availableAfter, reserved_balance: target.reserved_balance }),
now, now + auditRetentionMilliseconds,
);
this.database.prepare(`
INSERT INTO credit_adjustment_receipts (
adjustment_id, admin_id, user_id, idempotency_key_digest, request_hash,
ledger_id, available_after, reserved_after, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
`).run(
input.adjustmentId, input.adminId, input.userId, idempotencyKeyDigest, requestHash,
ledgerId, availableAfter, target.reserved_balance, now,
);
return {
adjustmentId: input.adjustmentId,
availableBalance: availableAfter,
reservedBalance: target.reserved_balance,
status: "adjusted" as const,
};
});
}
private immediate<T>(action: () => T): T {
if (this.database.inTransaction) return action();
this.database.exec("BEGIN IMMEDIATE");
try {
const result = action();
this.database.exec("COMMIT");
return result;
} catch (error) {
if (this.database.inTransaction) this.database.exec("ROLLBACK");
throw error;
}
}
private insertLedger(input: {
amount: number;
availableAfter: number;
availableBefore: number;
createdAt: number;
entryStatus: LedgerRow["entry_status"];
entryType: LedgerRow["entry_type"];
modelId: string | null;
operationKey: string;
reason: string | null;
referenceId: string | null;
referenceType: LedgerRow["reference_type"];
reservedAfter: number;
reservedBefore: number;
userId: string;
}) {
const ledgerId = randomUUID();
this.database.prepare(`
INSERT INTO credit_ledger (
ledger_id, user_id, operation_key, entry_type, amount,
available_before, available_after, reserved_before, reserved_after, created_at,
reference_type, reference_id, model_id, reason, entry_status
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`).run(
ledgerId, input.userId, input.operationKey, input.entryType, input.amount,
input.availableBefore, input.availableAfter, input.reservedBefore, input.reservedAfter, input.createdAt,
input.referenceType, input.referenceId, input.modelId, input.reason, input.entryStatus,
);
return ledgerId;
}
private insertOutbox(topic: string, aggregateId: string, operationKey: string, payload: Record<string, unknown>, now: number) {
this.database.prepare(`
INSERT INTO outbox_events (
event_id, operation_key, topic, aggregate_type, aggregate_id, payload_json, status, created_at, published_at
) VALUES (?, ?, ?, 'generation', ?, ?, 'pending', ?, NULL)
`).run(randomUUID(), operationKey, topic, aggregateId, JSON.stringify(payload), now);
}
private tableExists(name: string) {
return Boolean(this.database.prepare("SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?").get(name));
}
private ensureColumn(table: string, column: string, definition: string) {
const columns = this.database.prepare(`PRAGMA table_info(${table})`).all() as Array<{ name: string }>;
if (!columns.some((value) => value.name === column)) this.database.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${definition}`);
}
private migrate() {
if (!this.tableExists("credit_accounts") || !this.tableExists("credit_ledger")) {
throw new Error("credit_schema_unavailable");
}
this.ensureColumn("credit_ledger", "reference_type", "TEXT");
this.ensureColumn("credit_ledger", "reference_id", "TEXT");
this.ensureColumn("credit_ledger", "model_id", "TEXT");
this.ensureColumn("credit_ledger", "reason", "TEXT");
this.ensureColumn("credit_ledger", "entry_status", "TEXT NOT NULL DEFAULT 'succeeded'");
if (this.tableExists("generation_jobs")) {
this.ensureColumn("generation_jobs", "model_id", "TEXT");
this.ensureColumn("generation_jobs", "model_config_version", "INTEGER");
this.ensureColumn("generation_jobs", "confirmed_credit_cost", "INTEGER");
this.ensureColumn("generation_jobs", "reserved_credits", "INTEGER NOT NULL DEFAULT 0");
this.ensureColumn("generation_jobs", "final_credit_state", "TEXT");
this.ensureColumn("generation_jobs", "finished_at", "INTEGER");
}
this.database.exec(`
CREATE TABLE IF NOT EXISTS credit_reservations (
generation_id TEXT PRIMARY KEY,
user_id TEXT NOT NULL REFERENCES users(user_id),
model_id TEXT NOT NULL,
amount INTEGER NOT NULL CHECK (amount > 0),
status TEXT NOT NULL CHECK (status IN ('reserved', 'committed', 'released')),
created_at INTEGER NOT NULL,
finalized_at INTEGER
);
CREATE INDEX IF NOT EXISTS credit_reservations_user_status ON credit_reservations(user_id, status, created_at);
CREATE TABLE IF NOT EXISTS credit_adjustment_receipts (
adjustment_id TEXT PRIMARY KEY,
admin_id TEXT REFERENCES users(user_id),
user_id TEXT REFERENCES users(user_id),
idempotency_key_digest TEXT CHECK (idempotency_key_digest IS NULL OR length(idempotency_key_digest) = 64),
request_hash TEXT NOT NULL CHECK (length(request_hash) = 64),
ledger_id TEXT NOT NULL UNIQUE REFERENCES credit_ledger(ledger_id),
available_after INTEGER NOT NULL,
reserved_after INTEGER NOT NULL CHECK (reserved_after >= 0),
created_at INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS outbox_events (
event_id TEXT PRIMARY KEY,
operation_key TEXT NOT NULL UNIQUE,
topic TEXT NOT NULL,
aggregate_type TEXT NOT NULL,
aggregate_id TEXT NOT NULL,
payload_json TEXT NOT NULL CHECK (json_valid(payload_json)),
status TEXT NOT NULL CHECK (status IN ('pending', 'published')),
created_at INTEGER NOT NULL,
published_at INTEGER
);
DROP TRIGGER IF EXISTS credit_ledger_shape_guard;
CREATE TRIGGER credit_ledger_shape_guard BEFORE INSERT ON credit_ledger
WHEN
NEW.entry_status NOT IN ('succeeded', 'frozen', 'committed', 'released')
OR (NEW.entry_type = 'admin_adjustment' AND (
NEW.reference_type <> 'admin_adjustment' OR NEW.reference_id IS NULL OR trim(COALESCE(NEW.reason, '')) = ''
))
OR (NEW.entry_type IN ('generation_reserve', 'generation_commit', 'generation_release') AND (
NEW.reference_type <> 'generation' OR NEW.reference_id IS NULL OR NEW.model_id IS NULL
))
BEGIN SELECT RAISE(ABORT, 'credit_ledger_shape_invalid'); END;
`);
this.ensureColumn("credit_adjustment_receipts", "admin_id", "TEXT REFERENCES users(user_id)");
this.ensureColumn("credit_adjustment_receipts", "user_id", "TEXT REFERENCES users(user_id)");
this.ensureColumn("credit_adjustment_receipts", "idempotency_key_digest", "TEXT");
this.database.exec(`
CREATE UNIQUE INDEX IF NOT EXISTS credit_adjustment_idempotency
ON credit_adjustment_receipts(admin_id, user_id, idempotency_key_digest)
WHERE admin_id IS NOT NULL AND user_id IS NOT NULL AND idempotency_key_digest IS NOT NULL;
`);
ensureAdminOperationAuditSchema(this.database, this.clock());
}
}