'use strict'; const crypto = require('crypto'); const { getResponseDataUsage } = require('../utils/fastgptResponseData'); const DOC_KIND = { BID: 'bid', PLAN: 'plan', }; const EVENT_TYPE = { CREATE: 'create', REWRITE: 'rewrite', GENERATE: 'generate', }; const BILLING_STATUS = { FREE: 'free', PENDING: 'pending', CONFIRMED: 'confirmed', ROLLED_BACK: 'rolled_back', SKIP: 'skip', }; const QUOTA_ERROR_CODE = { CREATE_LIMIT_REACHED: 'CREATE_LIMIT_REACHED', REWRITE_LIMIT_REACHED: 'REWRITE_LIMIT_REACHED', RECHARGE_REQUIRED: 'RECHARGE_REQUIRED', CHARGE_CONFIRM_REQUIRED: 'CHARGE_CONFIRM_REQUIRED', CHARGE_CONFIRM_INVALID: 'CHARGE_CONFIRM_INVALID', BILLING_SERVICE_UNAVAILABLE: 'BILLING_SERVICE_UNAVAILABLE', }; const QUOTA_ENFORCED_REGISTER_SOURCES = new Set(['OFFICIALWEBSITE', 'PHONE']); const CHARGEABLE_REGISTER_SOURCES = new Set(['OFFICIALWEBSITE']); const toSafePositiveNumber = (value, fallback) => { const n = Number(value); return Number.isFinite(n) && n > 0 ? n : fallback; }; const buildQuotaError = (code, message) => { const error = new Error(message); error.code = code; return error; }; const normalizeUserId = (value) => { const text = String(value ?? '').trim(); if (!text) return null; if (!/^\d+$/.test(text)) return null; const n = Number(text); if (!Number.isSafeInteger(n) || n <= 0) return null; return n; }; const normalizeBoolean = (value) => { if (value === true || value === false) return value; const text = String(value ?? '').trim().toLowerCase(); return text === '1' || text === 'true' || text === 'yes' || text === 'on'; }; const toUsageFromResponse = (responseBody = {}) => { // FastGPT top-level usage may be a placeholder. Only detail=true responseData is trusted. return getResponseDataUsage(responseBody); }; const buildService = (app) => { const getOrm = () => { const orm = app?.fs?.dc?.orm; if (!orm) { throw buildQuotaError( QUOTA_ERROR_CODE.BILLING_SERVICE_UNAVAILABLE, '计费服务未初始化:数据库连接未就绪' ); } return orm; }; // 延迟获取 orm,避免服务初始化时机早于 sequelize 注入导致启动报错。 const orm = new Proxy({}, { get(_target, prop) { const ormInstance = getOrm(); const raw = ormInstance[prop]; return typeof raw === 'function' ? raw.bind(ormInstance) : raw; }, }); const QueryTypes = new Proxy({}, { get(_target, prop) { return getOrm()?.QueryTypes?.[prop]; }, }); const billingConfig = app.fs.config.billing || {}; const tokenPerYuan = toSafePositiveNumber(billingConfig.tokenPerYuan, 10000); const createChargeYuan = toSafePositiveNumber( process.env.BILLING_CREATE_CHARGE_YUAN ?? billingConfig.createChargeYuan, 5 ); const rewriteLimitPerHour = Math.max( 1, Number(process.env.BILLING_REWRITE_LIMIT_PER_HOUR || 20) || 20 ); const rewriteWindowMinutes = Math.max( 1, Number(process.env.BILLING_REWRITE_WINDOW_MINUTES || 60) || 60 ); const confirmTokenTtlMs = Math.max( 30 * 1000, Number(billingConfig.confirmTokenTtlMs || billingConfig.confirmTokenTtl || 5 * 60 * 1000) || (5 * 60 * 1000) ); const chargeConfirmStore = new Map(); const calculateAmountYuan = (totalTokens) => { // 按 TOKEN_PER_YUAN 换算金额,并固定保留 6 位小数作为账务精度。 const tokens = Number(totalTokens || 0); if (!Number.isFinite(tokens) || tokens <= 0) return 0; return Number((tokens / tokenPerYuan).toFixed(6)); }; const extractFirstRow = (queryResult) => { if (!Array.isArray(queryResult)) return null; if (Array.isArray(queryResult[0])) return queryResult[0][0] || null; if (queryResult[0] && typeof queryResult[0] === 'object') return queryResult[0]; return null; }; const lockByKey = async (transaction, lockKey) => { await orm.query( 'SELECT pg_advisory_xact_lock(hashtext(:lockKey));', { replacements: { lockKey: String(lockKey || '') }, type: QueryTypes.SELECT, transaction, } ); }; const pruneExpiredConfirmToken = () => { const now = Date.now(); for (const [token, item] of chargeConfirmStore.entries()) { if (!item?.expireAt || item.expireAt <= now) { chargeConfirmStore.delete(token); } } }; const createChargeConfirmToken = ({ userId, docKind, action, aideductUserId }) => { pruneExpiredConfirmToken(); const token = crypto.randomBytes(24).toString('hex'); const expireAt = Date.now() + confirmTokenTtlMs; chargeConfirmStore.set(token, { userId: Number(userId), docKind: String(docKind || ''), action: String(action || ''), aideductUserId: String(aideductUserId || ''), expireAt, }); return { token, expireAt, }; }; const consumeChargeConfirmToken = ({ token, userId, docKind, action, aideductUserId, }) => { pruneExpiredConfirmToken(); const key = String(token || '').trim(); if (!key) return false; const payload = chargeConfirmStore.get(key); if (!payload) return false; const valid = ( Number(payload.userId) === Number(userId) && String(payload.docKind || '') === String(docKind || '') && String(payload.action || '') === String(action || '') && String(payload.aideductUserId || '') === String(aideductUserId || '') ); if (!valid) return false; chargeConfirmStore.delete(key); return true; }; const resolveAideductUserId = async ({ transaction, userId, externalUserId = '' }) => { const preferredExternalUserId = String(externalUserId || '').trim(); if (preferredExternalUserId && preferredExternalUserId !== String(userId)) { return preferredExternalUserId; } const [userRow] = await orm.query( 'SELECT external_user_id FROM tender_user WHERE id = :userId LIMIT 1;', { replacements: { userId }, type: QueryTypes.SELECT, transaction, } ); const external = String(userRow?.external_user_id || '').trim(); if (external) return external; return String(userId); }; const resolveUserBillingMode = async ({ transaction, userId }) => { const normalizedUserId = normalizeUserId(userId); if (!normalizedUserId) { return { registerSource: '', quotaEnabled: false, chargeable: false, }; } const replacements = { userId: normalizedUserId }; const queryOptions = transaction ? { replacements, type: QueryTypes.SELECT, transaction } : { replacements, type: QueryTypes.SELECT }; const [userRow] = await orm.query( 'SELECT register_source FROM tender_user WHERE id = :userId LIMIT 1;', queryOptions ); const registerSource = String(userRow?.register_source || '').trim().toUpperCase(); return { registerSource, quotaEnabled: QUOTA_ENFORCED_REGISTER_SOURCES.has(registerSource), chargeable: CHARGEABLE_REGISTER_SOURCES.has(registerSource), }; }; const checkWalletChargeDecision = async ({ transaction, userId, docKind, action, chargeAmountYuan = 0, confirmCharge = false, confirmToken = '', externalUserId = '', }) => { const billingMode = await resolveUserBillingMode({ transaction, userId }); if (!billingMode.chargeable) { const isRewriteAction = String(action || '') === EVENT_TYPE.REWRITE; return { passed: false, code: isRewriteAction ? QUOTA_ERROR_CODE.REWRITE_LIMIT_REACHED : QUOTA_ERROR_CODE.CREATE_LIMIT_REACHED, message: isRewriteAction ? `重编试用次数已达上限,请${rewriteWindowMinutes}分钟后再试` : '今日试用次数已用完,请明日再试', }; } const aideductService = app.fs?.aideductService; if (!aideductService?.getBalance) { return { passed: false, code: QUOTA_ERROR_CODE.BILLING_SERVICE_UNAVAILABLE, message: '扣费服务暂不可用,请稍后重试', }; } const aideductUserId = await resolveAideductUserId({ transaction, userId, externalUserId, }); let balance = 0; try { const balanceRes = await aideductService.getBalance({ userId: aideductUserId }); balance = Number(balanceRes?.data?.balance || 0); } catch (error) { app?.logger?.log?.(error); return { passed: false, code: QUOTA_ERROR_CODE.BILLING_SERVICE_UNAVAILABLE, message: '扣费服务暂不可用,请稍后重试', }; } if (!Number.isFinite(balance) || balance <= 0) { return { passed: false, code: QUOTA_ERROR_CODE.RECHARGE_REQUIRED, message: '免费次数已用完,请充值后继续使用', }; } const isConfirmed = normalizeBoolean(confirmCharge); if (!isConfirmed) { const confirmInfo = createChargeConfirmToken({ userId, docKind, action, aideductUserId, }); return { passed: false, code: QUOTA_ERROR_CODE.CHARGE_CONFIRM_REQUIRED, message: chargeAmountYuan > 0 ? `免费次数已用完,继续使用将扣费${Number(chargeAmountYuan)}元` : '免费次数已用完,继续使用将进行扣费', confirmToken: confirmInfo.token, confirmExpireAt: new Date(confirmInfo.expireAt).toISOString(), }; } const validConfirm = consumeChargeConfirmToken({ token: confirmToken, userId, docKind, action, aideductUserId, }); if (!validConfirm) { return { passed: false, code: QUOTA_ERROR_CODE.CHARGE_CONFIRM_INVALID, message: '扣费确认已失效,请重新操作', }; } return { passed: true, chargeRequired: chargeAmountYuan > 0, chargeAmountYuan: chargeAmountYuan > 0 ? Number(chargeAmountYuan) : 0, }; }; const ensureUserProfileAndPolicy = async (transaction, userId) => { const [userRow] = await orm.query( 'SELECT pep_id FROM tender_user WHERE id = :userId LIMIT 1;', { replacements: { userId }, type: QueryTypes.SELECT, transaction, } ); const policyCode = userRow?.pep_id ? 'default' : 'non_project_enterprise'; await orm.query( `INSERT INTO user_billing_profile (user_id, policy_code, activated_at, created_at, updated_at) VALUES (:userId, :policyCode, now(), now(), now()) ON CONFLICT (user_id) DO UPDATE SET policy_code = EXCLUDED.policy_code, updated_at = now();`, { replacements: { userId, policyCode }, type: QueryTypes.INSERT, transaction, } ); const [policy] = await orm.query( `SELECT bp.* FROM user_billing_profile ubp JOIN billing_policy bp ON bp.code = ubp.policy_code WHERE ubp.user_id = :userId LIMIT 1;`, { replacements: { userId }, type: QueryTypes.SELECT, transaction, } ); if (!policy) { throw new Error('用户计费策略不存在'); } return policy; }; const getUserQuotaSnapshot = async ({ userId }) => { const normalizedUserId = normalizeUserId(userId); if (!normalizedUserId) return null; const billingMode = await resolveUserBillingMode({ userId: normalizedUserId }); if (!billingMode.quotaEnabled) return null; const transaction = await orm.transaction(); try { const policy = await ensureUserProfileAndPolicy(transaction, normalizedUserId); const createRows = await orm.query( `SELECT doc_kind, COUNT(*)::BIGINT AS used_count FROM quota_event WHERE user_id = :userId AND event_type = 'create' GROUP BY doc_kind;`, { replacements: { userId: normalizedUserId }, type: QueryTypes.SELECT, transaction, } ); const [rewriteRow] = await orm.query( `SELECT COUNT(*)::BIGINT AS rewrite_count FROM quota_event WHERE user_id = :userId AND event_type = 'rewrite' AND created_at >= now() - make_interval(mins => :windowMinutes);`, { replacements: { userId: normalizedUserId, windowMinutes: rewriteWindowMinutes, }, type: QueryTypes.SELECT, transaction, } ); await transaction.commit(); const createCountMap = new Map( (Array.isArray(createRows) ? createRows : []).map((row) => [ String(row?.doc_kind || '').trim().toLowerCase(), Number(row?.used_count || 0), ]) ); const bidLimit = Number(policy?.seed_bid_limit || 0); const planLimit = Number(policy?.seed_plan_limit || 0); const bidUsed = Number(createCountMap.get(DOC_KIND.BID) || 0); const planUsed = Number(createCountMap.get(DOC_KIND.PLAN) || 0); const rewriteUsed = Number(rewriteRow?.rewrite_count || 0); return { policyCode: String(policy?.code || ''), policyName: String(policy?.name || ''), create: { bid: { limit: bidLimit, used: bidUsed, remaining: Math.max(0, bidLimit - bidUsed), }, plan: { limit: planLimit, used: planUsed, remaining: Math.max(0, planLimit - planUsed), }, }, rewrite: { limitPerHour: rewriteLimitPerHour, windowMinutes: rewriteWindowMinutes, used: rewriteUsed, remaining: Math.max(0, rewriteLimitPerHour - rewriteUsed), }, generatedAt: new Date().toISOString(), }; } catch (error) { await transaction.rollback(); throw error; } }; const checkCreateQuota = async ({ transaction, userId, docKind, confirmCharge = false, confirmToken = '', externalUserId = '', }) => { if (!userId || !docKind) { return { passed: true }; } const billingMode = await resolveUserBillingMode({ transaction, userId }); if (!billingMode.quotaEnabled) { return { passed: true, chargeRequired: false, chargeAmountYuan: 0, }; } // 创建配额并发敏感,按 user+docKind 加事务级 advisory lock。 await lockByKey(transaction, `create:${userId}:${docKind}`); const policy = await ensureUserProfileAndPolicy(transaction, userId); const [counter] = await orm.query( `SELECT COUNT(*)::BIGINT AS total_created, COUNT(*) FILTER ( WHERE created_at >= now() - make_interval(hours => :rollingHours) )::BIGINT AS rolling_created, COUNT(*) FILTER ( WHERE created_at >= date_trunc('day', now()) )::BIGINT AS today_created FROM quota_event WHERE user_id = :userId AND event_type = 'create' AND doc_kind = :docKind;`, { replacements: { userId, docKind, rollingHours: Number(policy.rolling_create_hours || 24), }, type: QueryTypes.SELECT, transaction, } ); const totalCreated = Number(counter?.total_created || 0); const todayCreated = Number(counter?.today_created || 0); const seedLimitRaw = docKind === DOC_KIND.BID ? policy.seed_bid_limit : policy.seed_plan_limit; const seedLimit = Number.isFinite(Number(seedLimitRaw)) ? Number(seedLimitRaw) : 3; // 先消耗永久种子额度。 if (totalCreated < seedLimit) return { passed: true }; // 超过免费总额度后,每日首次创建继续免费。 if (todayCreated === 0) return { passed: true }; // 手机号用户仅走试用次数逻辑,超限后需等待,不进入扣费链路。 if (!billingMode.chargeable) { return { passed: false, code: QUOTA_ERROR_CODE.CREATE_LIMIT_REACHED, message: '今日试用次数已用完,请明日再试', }; } return checkWalletChargeDecision({ transaction, userId, docKind, action: EVENT_TYPE.CREATE, chargeAmountYuan: createChargeYuan, confirmCharge, confirmToken, externalUserId, }); }; const recordCreateEvent = async ({ transaction, userId, docKind, docId, source = 'manual', billingStatus = BILLING_STATUS.SKIP, amountYuan = 0, requestId = null, }) => { if (!userId || !docKind || !docId) return null; const billingMode = await resolveUserBillingMode({ transaction, userId }); if (!billingMode.quotaEnabled) return null; // 删除不回退额度,因此用事件流水统计,且同文档创建事件幂等写入。 const result = await orm.query( `INSERT INTO quota_event ( user_id, doc_id, doc_kind, event_type, source, amount_yuan, billing_status, request_id, created_at ) VALUES ( :userId, :docId, :docKind, 'create', :source, :amountYuan, :billingStatus, :requestId, now() ) ON CONFLICT DO NOTHING RETURNING id;`, { replacements: { userId, docId, docKind, source, amountYuan, billingStatus, requestId, }, transaction, } ); return extractFirstRow(result)?.id || null; }; const checkAndRecordRewrite = async ({ transaction, userId, docId, docKind, chapterKey = null, action = 'rewrite', consumeCount = 1, confirmCharge = false, confirmToken = '', externalUserId = '', }) => { if (!userId || !docId) { return { passed: true }; } const billingMode = await resolveUserBillingMode({ transaction, userId }); if (!billingMode.quotaEnabled) { return { passed: true, overLimit: false, chargeRequired: false, chargeAmountYuan: 0, }; } const normalizedConsumeCount = Math.max(1, Math.floor(Number(consumeCount || 1))); // 重写限流按 user 粒度,全文档共享每小时次数。 await lockByKey(transaction, `rewrite:${userId}`); const [counter] = await orm.query( `SELECT COUNT(*)::BIGINT AS rewrite_count FROM quota_event WHERE user_id = :userId AND event_type = 'rewrite' AND created_at >= now() - make_interval(mins => :windowMinutes);`, { replacements: { userId, windowMinutes: rewriteWindowMinutes, }, type: QueryTypes.SELECT, transaction, } ); const rewriteCount = Number(counter?.rewrite_count || 0); const overLimit = rewriteCount + normalizedConsumeCount > rewriteLimitPerHour; if (overLimit) { if (!billingMode.chargeable) { return { passed: false, code: QUOTA_ERROR_CODE.REWRITE_LIMIT_REACHED, message: `重编试用次数已达上限,请${rewriteWindowMinutes}分钟后再试`, }; } const decision = await checkWalletChargeDecision({ transaction, userId, docKind, action: EVENT_TYPE.REWRITE, confirmCharge, confirmToken, externalUserId, }); if (decision?.passed === false) return decision; return { ...decision, passed: true, overLimit: true, }; } await orm.query( `INSERT INTO quota_event ( user_id, doc_id, doc_kind, chapter_key, event_type, action, billing_status, created_at ) SELECT :userId, :docId, :docKind, :chapterKey, 'rewrite', :action, 'skip', now() FROM generate_series(1, :consumeCount);`, { replacements: { userId, docId, docKind: docKind || null, chapterKey, action, consumeCount: normalizedConsumeCount, }, type: QueryTypes.INSERT, transaction, } ); return { passed: true, overLimit: false, }; }; const markChapterFirstFree = async ({ userId, docId, chapterKey }) => { if (!userId || !docId || !chapterKey) return false; // 首次章节免费:唯一键冲突即表示“已免过”,无需额外读查询。 const result = await orm.query( `INSERT INTO quota_chapter_first_free ( user_id, doc_id, chapter_key, first_generated_at ) VALUES ( :userId, :docId, :chapterKey, now() ) ON CONFLICT DO NOTHING RETURNING user_id;`, { replacements: { userId, docId, chapterKey }, } ); return !!extractFirstRow(result); }; const isDocGenerateBillingExemptByCreate = async ({ userId, docId, docKind }) => { if (!userId || !docId || !docKind) return false; const billingMode = await resolveUserBillingMode({ userId }); if (!billingMode.quotaEnabled || !billingMode.chargeable) return true; const [row] = await orm.query( `SELECT billing_status FROM quota_event WHERE user_id = :userId AND doc_id = :docId AND doc_kind = :docKind AND event_type = 'create' ORDER BY id DESC LIMIT 1;`, { replacements: { userId, docId, docKind }, type: QueryTypes.SELECT, } ); const status = String(row?.billing_status || '').trim().toLowerCase(); // 免费创建(skip/free)文档:目录与正文生成统一免扣费。 return status === BILLING_STATUS.SKIP || status === BILLING_STATUS.FREE; }; const recordGenerateEvent = async ({ userId, docId, docKind, chapterKey = null, usage = {}, billingStatus = BILLING_STATUS.SKIP, tokenPerYuanSnapshot = tokenPerYuan, amountYuan = 0, requestId = null, source = 'manual', }) => { if (!userId || !docId) return null; // 生成事件统一沉淀 token 与金额快照,后续审计不受配置变更影响。 const inputTokens = Number(usage.inputTokens || 0); const outputTokens = Number(usage.outputTokens || 0); const embeddingTokens = Number(usage.embeddingTokens || 0); const totalTokens = Number(usage.totalTokens || 0); const result = await orm.query( `INSERT INTO quota_event ( user_id, doc_id, doc_kind, chapter_key, event_type, source, input_tokens, output_tokens, embedding_tokens, total_tokens, token_per_yuan, amount_yuan, billing_status, request_id, created_at ) VALUES ( :userId, :docId, :docKind, :chapterKey, 'generate', :source, :inputTokens, :outputTokens, :embeddingTokens, :totalTokens, :tokenPerYuanSnapshot, :amountYuan, :billingStatus, :requestId, now() ) RETURNING id;`, { replacements: { userId, docId, docKind: docKind || null, chapterKey, source, inputTokens, outputTokens, embeddingTokens, totalTokens, tokenPerYuanSnapshot, amountYuan, billingStatus, requestId, }, } ); return extractFirstRow(result)?.id || null; }; const updateGenerateBillingStatus = async ({ eventId, billingStatus }) => { if (!eventId) return; await orm.query( `UPDATE quota_event SET billing_status = :billingStatus WHERE id = :eventId;`, { replacements: { eventId, billingStatus }, type: QueryTypes.UPDATE, } ); }; const createRequestId = ({ userId, docKind, docId, chapterKey = '' }) => { return [ 'gen', userId, docKind || 'unknown', docId, chapterKey || 'none', Date.now(), crypto.randomBytes(4).toString('hex'), ].join('_'); }; const createCreateChargeRequestId = ({ userId, docKind, docId }) => { return [ 'create', userId, docKind || 'unknown', docId || 'none', Date.now(), crypto.randomBytes(4).toString('hex'), ].join('_'); }; const resolveChargeUserId = async ({ transaction, userId, externalUserId = '' }) => { return resolveAideductUserId({ transaction, userId, externalUserId, }); }; const settleGenerateBilling = async ({ eventId, aideductUserId, amountYuan, requestId, logger, }) => { const service = app.fs.aideductService; if (!service?.preDeduct || !service?.confirmDeduct || !service?.rollbackDeduct) { // 没有外部钱包服务时,保留 usage 流水但标记为 skip。 await updateGenerateBillingStatus({ eventId, billingStatus: BILLING_STATUS.SKIP }); return BILLING_STATUS.SKIP; } try { // 两阶段扣费:先预扣再确认,确保外部账务一致性。 await service.preDeduct({ userId: aideductUserId, amount: amountYuan, requestId }); await service.confirmDeduct({ userId: aideductUserId, amount: amountYuan, requestId }); await updateGenerateBillingStatus({ eventId, billingStatus: BILLING_STATUS.CONFIRMED }); return BILLING_STATUS.CONFIRMED; } catch (error) { try { // 任何失败都尝试回退预扣,避免冻结金额悬挂。 await service.rollbackDeduct({ userId: aideductUserId, amount: amountYuan, requestId }); } catch (rollbackError) { logger?.log?.(rollbackError); } await updateGenerateBillingStatus({ eventId, billingStatus: BILLING_STATUS.ROLLED_BACK }); logger?.log?.(error); return BILLING_STATUS.ROLLED_BACK; } }; return { DOC_KIND, EVENT_TYPE, BILLING_STATUS, QUOTA_ERROR_CODE, tokenPerYuan, createChargeYuan, normalizeUserId, toUsageFromResponse, calculateAmountYuan, getUserQuotaSnapshot, checkCreateQuota, recordCreateEvent, checkAndRecordRewrite, markChapterFirstFree, isDocGenerateBillingExemptByCreate, recordGenerateEvent, updateGenerateBillingStatus, createRequestId, createCreateChargeRequestId, resolveChargeUserId, settleGenerateBilling, }; }; module.exports = async (app) => { app.fs = app.fs || {}; app.fs.quotaBillingService = buildService(app); };