feat: complete TASK-WP2-04 credit ledger
This commit is contained in:
@@ -10,6 +10,7 @@ import {
|
||||
AccountProfileUpdateResponseSchema,
|
||||
AccountSettingsResponseSchema,
|
||||
AdminAuthenticatedUserSchema,
|
||||
AdminCreditParamsSchema,
|
||||
AdminLoginCompleteRequestSchema,
|
||||
AdminLoginCompleteResponseSchema,
|
||||
AdminLoginSendRequestSchema,
|
||||
@@ -21,6 +22,16 @@ import {
|
||||
CorrelationIdSchema,
|
||||
AuthenticatedUserSchema,
|
||||
CreditSummarySchema,
|
||||
CreditAdjustmentHeadersSchema,
|
||||
CreditAdjustmentRequestSchema,
|
||||
CreditAdjustmentResponseSchema,
|
||||
CreditBalanceResponseSchema,
|
||||
CreditEntryStatusSchema,
|
||||
CreditEntryTypeSchema,
|
||||
CreditLedgerEntrySchema,
|
||||
CreditLedgerQuerySchema,
|
||||
CreditLedgerResponseSchema,
|
||||
CreditReferenceTypeSchema,
|
||||
CsrfHeadersSchema,
|
||||
ErrorDetailsSchema,
|
||||
ErrorEnvelopeSchema,
|
||||
@@ -70,6 +81,10 @@ import {
|
||||
type AdminLoginSendRequest,
|
||||
type AccountDeletionCompleteRequest,
|
||||
type AccountProfileUpdateRequest,
|
||||
type AdminCreditParams,
|
||||
type CreditAdjustmentHeaders,
|
||||
type CreditAdjustmentRequest,
|
||||
type CreditLedgerQuery,
|
||||
type LoginCompleteRequest,
|
||||
type LoginSendRequest,
|
||||
type FailedEmptyTrashRequest,
|
||||
@@ -98,6 +113,8 @@ import {
|
||||
type BrowserUnsupportedReason,
|
||||
} from "./browser-support.js";
|
||||
import { EventHub } from "./event-hub.js";
|
||||
import { CreditError } from "./credit-errors.js";
|
||||
import type { CreditService } from "./credits.js";
|
||||
import type { PublicAssetResolver } from "./local-data-root.js";
|
||||
import { isAllowedNetworkRequest, type NetworkBoundaryOptions } from "./network-boundary.js";
|
||||
import { ProjectError } from "./project-errors.js";
|
||||
@@ -125,6 +142,7 @@ export interface CreateAppOptions {
|
||||
browserGate?: boolean;
|
||||
browserSupportRelease?: BrowserSupportRelease;
|
||||
browserSupportSecret?: Buffer;
|
||||
credits?: CreditService;
|
||||
eventHub?: EventHub;
|
||||
networkBoundary?: NetworkBoundaryOptions;
|
||||
publicAssets?: PublicAssetResolver;
|
||||
@@ -215,6 +233,21 @@ function projectFailure(reply: FastifyReply, correlationId: string, error: unkno
|
||||
return reply.code(mapping[error.code]).send(null);
|
||||
}
|
||||
|
||||
function creditFailure(reply: FastifyReply, correlationId: string, error: unknown) {
|
||||
if (!(error instanceof CreditError)) {
|
||||
return reply.code(503).send(createErrorEnvelope({ code: "AUTH_SERVICE_UNAVAILABLE", correlationId }));
|
||||
}
|
||||
if (error.code === "credit_operation_conflict") {
|
||||
return reply.code(409).send(createErrorEnvelope({ code: "IDEMPOTENCY_KEY_CONFLICT", correlationId }));
|
||||
}
|
||||
const status = error.code === "credit_request_invalid"
|
||||
? 400
|
||||
: error.code === "credit_insufficient" || error.code === "credit_invariant_failed"
|
||||
? 409
|
||||
: 404;
|
||||
return reply.code(status).send(null);
|
||||
}
|
||||
|
||||
type ProjectSummaryView = ReturnType<ProjectService["listProjects"]>[number];
|
||||
type ProjectDetailView = ReturnType<ProjectService["getProject"]>;
|
||||
|
||||
@@ -331,6 +364,17 @@ export async function createApp(options: CreateAppOptions = {}) {
|
||||
AdminLoginCompleteResponseSchema,
|
||||
AdminSessionResponseSchema,
|
||||
CreditSummarySchema,
|
||||
CreditEntryTypeSchema,
|
||||
CreditEntryStatusSchema,
|
||||
CreditReferenceTypeSchema,
|
||||
CreditBalanceResponseSchema,
|
||||
CreditLedgerEntrySchema,
|
||||
CreditLedgerQuerySchema,
|
||||
CreditLedgerResponseSchema,
|
||||
AdminCreditParamsSchema,
|
||||
CreditAdjustmentHeadersSchema,
|
||||
CreditAdjustmentRequestSchema,
|
||||
CreditAdjustmentResponseSchema,
|
||||
CsrfHeadersSchema,
|
||||
AccountSettingsResponseSchema,
|
||||
AccountProfileUpdateRequestSchema,
|
||||
@@ -952,6 +996,199 @@ export async function createApp(options: CreateAppOptions = {}) {
|
||||
},
|
||||
);
|
||||
|
||||
app.get(
|
||||
"/api/v1/me/credits",
|
||||
{
|
||||
schema: {
|
||||
operationId: "getMyCredits",
|
||||
response: {
|
||||
200: Type.Ref(CreditBalanceResponseSchema),
|
||||
401: Type.Ref(ErrorEnvelopeSchema),
|
||||
404: Type.Null(),
|
||||
503: Type.Ref(ErrorEnvelopeSchema),
|
||||
},
|
||||
tags: ["Credits"],
|
||||
},
|
||||
},
|
||||
async (request, reply) => {
|
||||
if (!options.registration || !options.credits) {
|
||||
return reply.code(503).send(createErrorEnvelope({ code: "AUTH_SERVICE_UNAVAILABLE", correlationId: request.id }));
|
||||
}
|
||||
const token = cookieValue(headerValue(request.headers.cookie), userSessionCookieName);
|
||||
const session = token ? options.registration.readUserSession(token) : undefined;
|
||||
if (!session) return reply.code(401).send(createErrorEnvelope({ code: "AUTH_SESSION_INVALID", correlationId: request.id }));
|
||||
try {
|
||||
const account = options.credits.readAccount(session.userId);
|
||||
return {
|
||||
available_balance: account.availableBalance,
|
||||
reserved_balance: account.reservedBalance,
|
||||
updated_at: account.updatedAt,
|
||||
};
|
||||
} catch (error) {
|
||||
return creditFailure(reply, request.id, error);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
app.get(
|
||||
"/api/v1/me/credit-ledger",
|
||||
{
|
||||
attachValidation: true,
|
||||
schema: {
|
||||
operationId: "getMyCreditLedger",
|
||||
querystring: Type.Ref(CreditLedgerQuerySchema),
|
||||
response: {
|
||||
200: Type.Ref(CreditLedgerResponseSchema),
|
||||
400: Type.Null(),
|
||||
401: Type.Ref(ErrorEnvelopeSchema),
|
||||
404: Type.Null(),
|
||||
503: Type.Ref(ErrorEnvelopeSchema),
|
||||
},
|
||||
tags: ["Credits"],
|
||||
},
|
||||
},
|
||||
async (request, reply) => {
|
||||
if (request.validationError) return reply.code(400).send(null);
|
||||
if (!options.registration || !options.credits) {
|
||||
return reply.code(503).send(createErrorEnvelope({ code: "AUTH_SERVICE_UNAVAILABLE", correlationId: request.id }));
|
||||
}
|
||||
const token = cookieValue(headerValue(request.headers.cookie), userSessionCookieName);
|
||||
const session = token ? options.registration.readUserSession(token) : undefined;
|
||||
if (!session) return reply.code(401).send(createErrorEnvelope({ code: "AUTH_SESSION_INVALID", correlationId: request.id }));
|
||||
try {
|
||||
const query = request.query as CreditLedgerQuery;
|
||||
const ledger = options.credits.listLedger({
|
||||
...(query.cursor ? { cursor: query.cursor } : {}),
|
||||
...(query.event_type ? { eventType: query.event_type } : {}),
|
||||
...(query.from ? { from: query.from } : {}),
|
||||
...(query.limit ? { limit: query.limit } : {}),
|
||||
...(query.to ? { to: query.to } : {}),
|
||||
userId: session.userId,
|
||||
});
|
||||
return {
|
||||
credits: {
|
||||
available_balance: ledger.account.availableBalance,
|
||||
reserved_balance: ledger.account.reservedBalance,
|
||||
},
|
||||
entries: ledger.entries.map((entry) => ({
|
||||
amount: entry.amount,
|
||||
available_after: entry.availableAfter,
|
||||
available_before: entry.availableBefore,
|
||||
created_at: entry.createdAt,
|
||||
entry_id: entry.entryId,
|
||||
entry_type: entry.entryType,
|
||||
model_id: entry.modelId,
|
||||
reason: entry.reason,
|
||||
reference_id: entry.referenceId,
|
||||
reference_type: entry.referenceType,
|
||||
reserved_after: entry.reservedAfter,
|
||||
reserved_before: entry.reservedBefore,
|
||||
status: entry.status,
|
||||
})),
|
||||
next_cursor: ledger.nextCursor,
|
||||
updated_at: ledger.account.updatedAt,
|
||||
};
|
||||
} catch (error) {
|
||||
return creditFailure(reply, request.id, error);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
app.get(
|
||||
"/api/v1/admin/users/:userId/credits",
|
||||
{
|
||||
attachValidation: true,
|
||||
schema: {
|
||||
operationId: "getAdminUserCredits",
|
||||
params: Type.Ref(AdminCreditParamsSchema),
|
||||
response: {
|
||||
200: Type.Ref(CreditBalanceResponseSchema),
|
||||
400: Type.Null(),
|
||||
401: Type.Ref(ErrorEnvelopeSchema),
|
||||
404: Type.Null(),
|
||||
503: Type.Ref(ErrorEnvelopeSchema),
|
||||
},
|
||||
tags: ["Admin Credits"],
|
||||
},
|
||||
},
|
||||
async (request, reply) => {
|
||||
if (request.validationError) return reply.code(400).send(null);
|
||||
if (!options.registration || !options.credits) {
|
||||
return reply.code(503).send(createErrorEnvelope({ code: "AUTH_SERVICE_UNAVAILABLE", correlationId: request.id }));
|
||||
}
|
||||
const token = cookieValue(headerValue(request.headers.cookie), adminSessionCookieName);
|
||||
const session = token ? options.registration.readAdminSession(token) : undefined;
|
||||
if (!session) return reply.code(401).send(createErrorEnvelope({ code: "AUTH_SESSION_INVALID", correlationId: request.id }));
|
||||
try {
|
||||
const account = options.credits.readAccount((request.params as AdminCreditParams).userId);
|
||||
return {
|
||||
available_balance: account.availableBalance,
|
||||
reserved_balance: account.reservedBalance,
|
||||
updated_at: account.updatedAt,
|
||||
};
|
||||
} catch (error) {
|
||||
return creditFailure(reply, request.id, error);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
app.post(
|
||||
"/api/v1/admin/users/:userId/credit-adjustments",
|
||||
{
|
||||
attachValidation: true,
|
||||
schema: {
|
||||
body: Type.Ref(CreditAdjustmentRequestSchema),
|
||||
headers: Type.Ref(CreditAdjustmentHeadersSchema),
|
||||
operationId: "adjustAdminUserCredits",
|
||||
params: Type.Ref(AdminCreditParamsSchema),
|
||||
response: {
|
||||
200: Type.Ref(CreditAdjustmentResponseSchema),
|
||||
400: Type.Null(),
|
||||
401: Type.Ref(ErrorEnvelopeSchema),
|
||||
403: Type.Ref(ErrorEnvelopeSchema),
|
||||
404: Type.Null(),
|
||||
409: Type.Union([Type.Ref(ErrorEnvelopeSchema), Type.Null()]),
|
||||
503: Type.Ref(ErrorEnvelopeSchema),
|
||||
},
|
||||
tags: ["Admin Credits"],
|
||||
},
|
||||
},
|
||||
async (request, reply) => {
|
||||
if (request.validationError) return reply.code(400).send(null);
|
||||
if (!options.registration || !options.credits) {
|
||||
return reply.code(503).send(createErrorEnvelope({ code: "AUTH_SERVICE_UNAVAILABLE", correlationId: request.id }));
|
||||
}
|
||||
const token = cookieValue(headerValue(request.headers.cookie), adminSessionCookieName);
|
||||
const headers = request.headers as CreditAdjustmentHeaders;
|
||||
if (!token) return reply.code(401).send(createErrorEnvelope({ code: "AUTH_SESSION_INVALID", correlationId: request.id }));
|
||||
try {
|
||||
const admin = options.registration.authorizeAdminMutation({
|
||||
csrfToken: headers["x-csrf-token"],
|
||||
sessionToken: token,
|
||||
});
|
||||
const body = request.body as CreditAdjustmentRequest;
|
||||
const adjusted = options.credits.adjustAvailable({
|
||||
adjustmentId: body.adjustment_id,
|
||||
adminId: admin.userId,
|
||||
amount: body.amount,
|
||||
idempotencyKey: headers["idempotency-key"],
|
||||
reason: body.reason,
|
||||
userId: (request.params as AdminCreditParams).userId,
|
||||
});
|
||||
return {
|
||||
adjustment_id: adjusted.adjustmentId,
|
||||
available_balance: adjusted.availableBalance,
|
||||
reserved_balance: adjusted.reservedBalance,
|
||||
status: adjusted.status,
|
||||
};
|
||||
} catch (error) {
|
||||
return error instanceof RegistrationError
|
||||
? registrationFailure(reply, request.id, error)
|
||||
: creditFailure(reply, request.id, error);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
app.post(
|
||||
"/api/v1/account/deletion/send",
|
||||
{
|
||||
|
||||
@@ -0,0 +1,16 @@
|
||||
export type CreditErrorCode =
|
||||
| "credit_account_not_found"
|
||||
| "credit_generation_not_found"
|
||||
| "credit_insufficient"
|
||||
| "credit_invariant_failed"
|
||||
| "credit_operation_conflict"
|
||||
| "credit_request_invalid";
|
||||
|
||||
export class CreditError extends Error {
|
||||
readonly code: CreditErrorCode;
|
||||
|
||||
constructor(code: CreditErrorCode) {
|
||||
super(code);
|
||||
this.code = code;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,548 @@
|
||||
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 {
|
||||
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());
|
||||
}
|
||||
}
|
||||
@@ -7,6 +7,7 @@ import { createApp } from "./app.js";
|
||||
import { readBrowserSupportRelease } from "./browser-support.js";
|
||||
import { defaultInstanceConfigPath, readConfiguredLocalDataRoot } from "./local-data-root.js";
|
||||
import { ManagedStorage } from "./managed-storage.js";
|
||||
import { CreditService } from "./credits.js";
|
||||
import { ProjectService } from "./projects.js";
|
||||
import { RegistrationService } from "./registration.js";
|
||||
import { MockResendAdapter } from "./resend-adapter.js";
|
||||
@@ -17,6 +18,7 @@ import { attachApiSupervisorControl, initializeApiCredentialClients, receiveApiC
|
||||
const credentialChannelEnabled = process.argv.includes("--dada-credential-stdin");
|
||||
let registration: RegistrationService | undefined;
|
||||
let projects: ProjectService | undefined;
|
||||
let credits: CreditService | undefined;
|
||||
const instanceConfigPath = process.env.DADA_INSTANCE_CONFIG_PATH ?? defaultInstanceConfigPath();
|
||||
if (credentialChannelEnabled) {
|
||||
const clients = initializeApiCredentialClients(await receiveApiCredentials());
|
||||
@@ -36,8 +38,11 @@ if (credentialChannelEnabled) {
|
||||
sessionPepper: derivePepper("session-pepper"),
|
||||
});
|
||||
projects = new ProjectService({ databasePath });
|
||||
credits = new CreditService({ databasePath });
|
||||
registration.applySecureConfig(readSecureConfigCandidate(instanceConfigPath));
|
||||
} catch (error) {
|
||||
credits?.close();
|
||||
credits = undefined;
|
||||
projects?.close();
|
||||
projects = undefined;
|
||||
registration?.close();
|
||||
@@ -51,6 +56,7 @@ if (credentialChannelEnabled) {
|
||||
const browserSupportRelease = readBrowserSupportRelease(resolve("RELEASE.json"));
|
||||
const app = await createApp({
|
||||
...(browserSupportRelease ? { browserSupportRelease } : {}),
|
||||
...(credits ? { credits } : {}),
|
||||
...(projects ? { projects } : {}),
|
||||
...(registration ? { registration } : {}),
|
||||
});
|
||||
@@ -67,6 +73,7 @@ if (controlPipeIndex >= 0) {
|
||||
let storage: ManagedStorage | undefined;
|
||||
const control = attachApiSupervisorControl(controlPipe, async () => {
|
||||
await app.close();
|
||||
credits?.close();
|
||||
projects?.close();
|
||||
registration?.close();
|
||||
storage?.close();
|
||||
|
||||
@@ -783,6 +783,12 @@ export class ProjectService {
|
||||
prompt TEXT NOT NULL CHECK (length(prompt) BETWEEN 1 AND 4000),
|
||||
ratio TEXT NOT NULL CHECK (ratio IN ('3:4', '1:1', '4:3', '9:16')),
|
||||
status TEXT NOT NULL CHECK (status IN ('queued', 'running', 'succeeded', 'failed', 'rejected')),
|
||||
model_id TEXT,
|
||||
model_config_version INTEGER,
|
||||
confirmed_credit_cost INTEGER CHECK (confirmed_credit_cost IS NULL OR confirmed_credit_cost > 0),
|
||||
reserved_credits INTEGER NOT NULL DEFAULT 0 CHECK (reserved_credits >= 0),
|
||||
final_credit_state TEXT CHECK (final_credit_state IS NULL OR final_credit_state IN ('committed', 'released')),
|
||||
finished_at INTEGER,
|
||||
error_category TEXT CHECK (error_category IS NULL OR error_category IN (
|
||||
'upstream_timeout', 'upstream_failed', 'safety_rejected', 'model_disabled',
|
||||
'gateway_balance_insufficient', 'gateway_contract_invalid', 'reference_invalid',
|
||||
|
||||
@@ -990,6 +990,11 @@ export class RegistrationService {
|
||||
return { userId: session.user_id };
|
||||
}
|
||||
|
||||
authorizeAdminMutation(input: { csrfToken: string; sessionToken: string }) {
|
||||
const session = this.authenticatedAdminMutationSession(input.sessionToken, input.csrfToken, this.options.clock());
|
||||
return { userId: session.user_id };
|
||||
}
|
||||
|
||||
logoutUser(input: { csrfToken: string; sessionToken: string }) {
|
||||
const now = this.options.clock();
|
||||
this.runImmediate("session_revoke", () => {
|
||||
@@ -1531,6 +1536,26 @@ export class RegistrationService {
|
||||
return session;
|
||||
}
|
||||
|
||||
private authenticatedAdminMutationSession(sessionToken: string, csrfToken: string, now: number) {
|
||||
const session = this.database.prepare(`
|
||||
SELECT s.session_id, s.user_id, s.csrf_token_digest
|
||||
FROM sessions s
|
||||
JOIN users u ON u.user_id = s.user_id
|
||||
JOIN admin_access a ON a.user_id = u.user_id
|
||||
WHERE s.token_digest = ? AND s.audience = 'admin' AND s.revoked_at IS NULL
|
||||
AND s.expires_at > ? AND u.role = 'super_admin' AND u.status = 'active' AND a.allowed = 1
|
||||
`).get(digest(sessionToken), now) as {
|
||||
csrf_token_digest: string | null;
|
||||
session_id: string;
|
||||
user_id: string;
|
||||
} | undefined;
|
||||
if (!session) throw new RegistrationError("AUTH_SESSION_INVALID", "session_invalid");
|
||||
if (!session.csrf_token_digest || !constantTimeTextEqual(session.csrf_token_digest, digest(csrfToken))) {
|
||||
throw new RegistrationError("AUTH_CSRF_INVALID", "csrf_invalid");
|
||||
}
|
||||
return session;
|
||||
}
|
||||
|
||||
private queueOwnedManagedFiles(userId: string, now: number) {
|
||||
const table = this.database.prepare(`
|
||||
SELECT 1 AS present FROM sqlite_master WHERE type = 'table' AND name = 'managed_files'
|
||||
|
||||
Reference in New Issue
Block a user