diff --git a/video-gen-admin/src/App.tsx b/video-gen-admin/src/App.tsx index 59222f02..e91dde54 100644 --- a/video-gen-admin/src/App.tsx +++ b/video-gen-admin/src/App.tsx @@ -45,6 +45,8 @@ import AdminApiKeys from './pages/AdminApiKeys'; import AdminApiModelPricings from './pages/AdminApiModelPricings'; import AdminApiUsage from './pages/AdminApiUsage'; import AdminInvoices from './pages/AdminInvoices'; +import AdminBankTransactions from './pages/AdminBankTransactions'; +import AdminScheduledTasks from './pages/AdminScheduledTasks'; import { useAdminStore } from './store'; @@ -112,6 +114,8 @@ const App = () => { } /> } /> } /> + } /> + } /> } /> } /> } /> diff --git a/video-gen-admin/src/api/index.ts b/video-gen-admin/src/api/index.ts index 269f89d7..0a764ade 100644 --- a/video-gen-admin/src/api/index.ts +++ b/video-gen-admin/src/api/index.ts @@ -18,6 +18,7 @@ import type { PrivatePortraitConfig, PrivatePortraitProjectListOut, PrivatePortraitAssetListOut, AdminUploadFileResult, AdminUploadResourceType, AdminUploadScene, VideoUpscaleConfigOut, VideoUpscaleConfigSavePayload, CreditProduct, + BankAccount, ScheduledTask, } from '../types'; import type { @@ -216,6 +217,102 @@ export async function toggleUserStatus(userId: string, isActive: boolean): Promi await api.put(`/admin/users/${userId}/status`, { is_active: isActive }); } +export async function queryBankTransactions(params: { + accountId: string; + startDate: string; + endDate: string; + dcFlag?: number; + page?: number; + pageSize?: number; +}): Promise<{ items: any[]; total: number }> { + const qs = new URLSearchParams(); + qs.set('account_id', params.accountId); + qs.set('start_date', params.startDate); + qs.set('end_date', params.endDate); + if (params.dcFlag !== undefined) qs.set('dc_flag', String(params.dcFlag)); + if (params.page) qs.set('page', String(params.page)); + if (params.pageSize) qs.set('page_size', String(params.pageSize)); + return api.get(`/admin/bank/transactions?${qs.toString()}`); +} + +// ── Bank Account Management ──────────────────────────────── + +export async function listBankAccounts(): Promise<{ items: BankAccount[] }> { + return api.get('/admin/bank/accounts'); +} + +export async function createBankAccount(payload: { + account_name: string; + bank_name: string; + account_no: string; + is_active?: boolean; + is_default?: boolean; + description?: string; +}): Promise<{ id: string; message: string }> { + return api.post('/admin/bank/accounts', payload); +} + +export async function updateBankAccount( + accountId: string, + payload: Partial<{ + account_name: string; + bank_name: string; + account_no: string; + is_active: boolean; + is_default: boolean; + description: string; + }>, +): Promise<{ message: string }> { + return api.put(`/admin/bank/accounts/${accountId}`, payload); +} + +export async function deleteBankAccount(accountId: string): Promise<{ message: string }> { + await api.delete(`/admin/bank/accounts/${accountId}`); + return { message: '删除成功' }; +} + +// ── Scheduled Tasks ──────────────────────────────────────── + +export async function listScheduledTasks(): Promise<{ items: ScheduledTask[] }> { + return api.get('/admin/scheduled-tasks'); +} + +export async function createScheduledTask(payload: { + name: string; + task_type: 'external_api' | 'internal_method'; + schedule: string; + config?: Record | string; + is_active?: boolean; +}): Promise<{ id: string; message: string }> { + return api.post('/admin/scheduled-tasks', payload); +} + +export async function updateScheduledTask( + taskId: string, + payload: Partial<{ + name: string; + task_type: 'external_api' | 'internal_method'; + schedule: string; + config: Record | string; + is_active: boolean; + }>, +): Promise<{ message: string }> { + return api.put(`/admin/scheduled-tasks/${taskId}`, payload); +} + +export async function deleteScheduledTask(taskId: string): Promise<{ message: string }> { + await api.delete(`/admin/scheduled-tasks/${taskId}`); + return { message: '删除成功' }; +} + +export async function runScheduledTask(taskId: string): Promise<{ message: string }> { + return api.post(`/admin/scheduled-tasks/${taskId}/run`, {}); +} + +export async function toggleScheduledTask(taskId: string): Promise<{ is_active: boolean; message: string }> { + return api.post(`/admin/scheduled-tasks/${taskId}/toggle`, {}); +} + export async function updateSingleDeviceLoginOverride(userId: string, override: boolean | null): Promise { await api.put(`/admin/users/${userId}/single-device-login-override`, { override }); } diff --git a/video-gen-admin/src/pages/AdminBankTransactions.tsx b/video-gen-admin/src/pages/AdminBankTransactions.tsx new file mode 100644 index 00000000..0fe0be0f --- /dev/null +++ b/video-gen-admin/src/pages/AdminBankTransactions.tsx @@ -0,0 +1,250 @@ +import React, { useState, useEffect, useCallback } from 'react'; +import { Table, Button, Space, Typography, message, Card, DatePicker, Select } from 'antd'; +import { BankOutlined, SearchOutlined } from '@ant-design/icons'; +import { queryBankTransactions, listBankAccounts } from '../api'; +import { formatDate } from '../utils/formatDate'; +import type { BankAccount } from '../types'; + +const { RangePicker } = DatePicker; + +interface TransactionRecord { + counterAcctNo?: string; + cnterName?: string; + cnterBankName?: string; + acctNo?: string; + transAmt?: number | string; + dcFlag?: number; + dcFlagLabel?: string; + transTimeStr?: string; + digestCode?: string; + remark?: string; + balance?: number | string; + tellerSeqno?: string; + transferId?: string; +} + +const AdminBankTransactions: React.FC = () => { + const [data, setData] = useState([]); + const [total, setTotal] = useState(0); + const [loading, setLoading] = useState(false); + const [page, setPage] = useState(1); + const [pageSize, setPageSize] = useState(20); + + // 银行账户 + const [accounts, setAccounts] = useState([]); + const [selectedAccountId, setSelectedAccountId] = useState(undefined); + + // 搜索条件 + const [dateRange, setDateRange] = useState<[string | null, string | null]>([null, null]); + const [dcFlag, setDcFlag] = useState(undefined); + + useEffect(() => { + loadAccounts(); + }, []); + + const loadAccounts = async () => { + try { + const res = await listBankAccounts(); + setAccounts(res.items || []); + } catch (err: any) { + message.error(err?.message || '加载银行账户失败'); + } + }; + + const fetchData = useCallback(async () => { + if (!selectedAccountId || !dateRange[0] || !dateRange[1]) return; + setLoading(true); + try { + const res = await queryBankTransactions({ + accountId: selectedAccountId, + startDate: dateRange[0], + endDate: dateRange[1], + dcFlag, + page, + pageSize, + }); + setData(res.items || []); + setTotal(res.total || 0); + } catch (err: any) { + message.error(err?.message || '查询失败'); + } finally { + setLoading(false); + } + }, [selectedAccountId, dateRange, dcFlag, page, pageSize]); + + useEffect(() => { + if (selectedAccountId && dateRange[0] && dateRange[1]) { + fetchData(); + } + }, [fetchData]); + + const handleSearch = () => { + if (!selectedAccountId) { + message.warning('请选择银行账户'); + return; + } + if (!dateRange[0] || !dateRange[1]) { + message.warning('请选择交易时间区间'); + return; + } + setPage(1); + fetchData(); + }; + + const selectedAccount = accounts.find(a => a.id === selectedAccountId); + + const columns = [ + { + title: '收款账号', + dataIndex: 'counterAcctNo', + width: 180, + render: (v: string) => v || '-', + }, + { + title: '收款户名', + dataIndex: 'cnterName', + width: 120, + render: (v: string) => v || '-', + }, + { + title: '收款开户行', + dataIndex: 'cnterBankName', + width: 150, + render: (v: string) => v || '-', + }, + { + title: '我方付款账号', + dataIndex: 'acctNo', + width: 180, + render: (v: string) => v || '-', + }, + { + title: '打款金额', + dataIndex: 'transAmt', + width: 120, + align: 'right' as const, + render: (v: number | string) => v != null ? Number(v).toFixed(2) : '-', + }, + { + title: '借贷方向', + dataIndex: 'dcFlagLabel', + width: 100, + render: (v: string, r: TransactionRecord) => v || (r.dcFlag === 0 ? '借/出金' : r.dcFlag === 1 ? '贷/入金' : '-'), + }, + { + title: '打款时间', + dataIndex: 'transTimeStr', + width: 160, + render: (v: string) => v ? formatDate(v) : '-', + }, + { + title: '摘要', + dataIndex: 'digestCode', + width: 100, + render: (v: string) => v || '-', + }, + { + title: '用途/备注', + dataIndex: 'remark', + width: 150, + render: (v: string) => v || '-', + }, + { + title: '账户余额', + dataIndex: 'balance', + width: 120, + align: 'right' as const, + render: (v: number | string) => v != null ? Number(v).toFixed(2) : '-', + }, + { + title: '流水号', + dataIndex: 'tellerSeqno', + width: 150, + render: (v: string) => v || '-', + }, + ]; + + return ( + + + + + + + + 银行交易查询 + 查询银行账户交易流水记录 + + + + + ({ + value: a.id, + label: `${a.bank_name} - ${a.account_no} (${a.account_name})`, + }))} + /> + { + if (Array.isArray(dateStrings)) { + setDateRange([dateStrings[0] || null, dateStrings[1] || null]); + } + }} + /> + + 借/出金 + 贷/入金 + + } onClick={handleSearch}> + 查询 + + + + {selectedAccount && ( + + 当前查询账户:{selectedAccount.bank_name} | {selectedAccount.account_no} | {selectedAccount.account_name} + + )} + + + + `${r.tellerSeqno || ''}-${r.transferId || ''}-${i}`} + columns={columns} + dataSource={data} + loading={loading} + scroll={{ x: 1400 }} + pagination={{ + current: page, + pageSize, + total, + showSizeChanger: true, + showTotal: (t) => `共 ${t} 条`, + onChange: (p, ps) => { + setPage(p); + setPageSize(ps); + }, + }} + /> + + + ); +}; + +export default AdminBankTransactions; diff --git a/video-gen-admin/src/pages/AdminScheduledTasks.tsx b/video-gen-admin/src/pages/AdminScheduledTasks.tsx new file mode 100644 index 00000000..ff39e1aa --- /dev/null +++ b/video-gen-admin/src/pages/AdminScheduledTasks.tsx @@ -0,0 +1,311 @@ +import React, { useState, useEffect, useCallback } from 'react'; +import { + Table, Button, Space, Typography, message, Card, Modal, Form, Input, Select, Switch, Tag, Popconfirm, Tabs, +} from 'antd'; +import { + ClockCircleOutlined, PlusOutlined, EditOutlined, DeleteOutlined, PlayCircleOutlined, StopOutlined, CheckCircleOutlined, CloseCircleOutlined, +} from '@ant-design/icons'; +import { + listScheduledTasks, + createScheduledTask, + updateScheduledTask, + deleteScheduledTask, + runScheduledTask, + toggleScheduledTask, +} from '../api'; +import type { ScheduledTask } from '../types'; + +const { TextArea } = Input; + +const SCHEDULE_TYPE_OPTIONS = [ + { value: 'external_api', label: '外部接口调用' }, + { value: 'internal_method', label: '内部方法执行' }, +]; + +const TASK_TYPE_LABELS: Record = { + external_api: '外部接口', + internal_method: '内部方法', +}; + +const STATUS_LABELS: Record = { + success: { label: '成功', color: 'green' }, + error: { label: '失败', color: 'red' }, +}; + +const AdminScheduledTasks: React.FC = () => { + const [data, setData] = useState([]); + const [loading, setLoading] = useState(false); + const [modal, setModal] = useState(false); + const [editing, setEditing] = useState(null); + const [form] = Form.useForm(); + const [saving, setSaving] = useState(false); + const [activeTab, setActiveTab] = useState<'basic' | 'config'>('basic'); + + const fetchData = useCallback(async () => { + setLoading(true); + try { + const res = await listScheduledTasks(); + setData(res.items || []); + } catch (err: any) { + message.error(err?.message || '加载失败'); + } finally { + setLoading(false); + } + }, []); + + useEffect(() => { + fetchData(); + }, [fetchData]); + + const openModal = (task?: ScheduledTask) => { + if (task) { + setEditing(task); + let configStr = ''; + if (task.config) { + try { + configStr = typeof task.config === 'string' ? JSON.stringify(JSON.parse(task.config), null, 2) : JSON.stringify(task.config, null, 2); + } catch { + configStr = task.config; + } + } + form.setFieldsValue({ + name: task.name, + task_type: task.task_type, + schedule: task.schedule, + config: configStr, + is_active: task.is_active, + }); + } else { + setEditing(null); + form.resetFields(); + form.setFieldsValue({ is_active: true, task_type: 'external_api' }); + } + setActiveTab('basic'); + setModal(true); + }; + + const handleSave = async () => { + try { + const values = await form.validateFields(); + let configVal = values.config; + if (configVal && typeof configVal === 'string') { + try { + configVal = JSON.stringify(JSON.parse(configVal)); + } catch { + message.warning('配置 JSON 格式不合法,将按原样保存'); + } + } + setSaving(true); + if (editing) { + await updateScheduledTask(editing.id, { ...values, config: configVal }); + message.success('更新成功'); + } else { + await createScheduledTask({ ...values, config: configVal }); + message.success('创建成功'); + } + setModal(false); + await fetchData(); + } catch (err: any) { + message.error(err?.message || '保存失败'); + } finally { + setSaving(false); + } + }; + + const handleDelete = async (taskId: string) => { + try { + await deleteScheduledTask(taskId); + message.success('删除成功'); + await fetchData(); + } catch (err: any) { + message.error(err?.message || '删除失败'); + } + }; + + const handleRun = async (taskId: string) => { + try { + await runScheduledTask(taskId); + message.success('任务已提交执行'); + setTimeout(fetchData, 1500); + } catch (err: any) { + message.error(err?.message || '执行失败'); + } + }; + + const handleToggle = async (task: ScheduledTask) => { + try { + const res = await toggleScheduledTask(task.id); + message.success(res.message); + await fetchData(); + } catch (err: any) { + message.error(err?.message || '操作失败'); + } + }; + + const taskType = Form.useWatch('task_type', form); + + const columns = [ + { title: '任务名称', dataIndex: 'name', width: 160, ellipsis: true }, + { + title: '类型', + dataIndex: 'task_type', + width: 100, + render: (v: string) => {TASK_TYPE_LABELS[v] || v}, + }, + { title: '调度表达式', dataIndex: 'schedule', width: 140, render: (v: string) => {v} }, + { + title: '状态', + dataIndex: 'is_active', + width: 80, + render: (v: boolean) => {v ? '启用' : '禁用'}, + }, + { + title: '最后执行', + dataIndex: 'last_run_at', + width: 160, + render: (v: string, r: ScheduledTask) => { + if (!v) return '-'; + const status = r.last_status ? STATUS_LABELS[r.last_status] : null; + return ( + + {v.includes('T') ? v.replace('T', ' ').slice(0, 19) : v} + {status && {status.label}} + + ); + }, + }, + { + title: '操作', + width: 220, + render: (_: any, record: ScheduledTask) => ( + + } onClick={() => handleRun(record.id)}>执行 + : } onClick={() => handleToggle(record)}> + {record.is_active ? '禁用' : '启用'} + + } onClick={() => openModal(record)}>编辑 + handleDelete(record.id)} okText="确定" cancelText="取消"> + }>删除 + + + ), + }, + ]; + + const scheduleHelp = ( + + • 纯数字:间隔秒数(如 60 = 每 60 秒执行一次) + • Cron 表达式(5 字段):分 时 日 月 周 + • 示例:* * * * * = 每分钟 | 0 * * * * = 每小时 + + ); + + return ( + + + + + + + + + 定时任务管理 + 自定义定时执行外部接口调用或内部方法 + + + } onClick={() => openModal()}> + 新增任务 + + + + + + `共 ${t} 条` }} + /> + + + {/* 新增/编辑弹窗 */} + setModal(false)} + confirmLoading={saving} + okText="保存" + cancelText="取消" + destroyOnClose + width={600} + > + + setActiveTab(k as 'basic' | 'config')} + items={[ + { + key: 'basic', + label: '基本设置', + children: ( + <> + + + + + + + + + + + + + > + ), + }, + { + key: 'config', + label: '任务配置', + children: ( + + + + ), + }, + ]} + /> + + + + ); +}; + +export default AdminScheduledTasks; diff --git a/video-gen-admin/src/pages/AdminSettings.tsx b/video-gen-admin/src/pages/AdminSettings.tsx index f494d4a4..d511e4e9 100644 --- a/video-gen-admin/src/pages/AdminSettings.tsx +++ b/video-gen-admin/src/pages/AdminSettings.tsx @@ -1,9 +1,9 @@ import React, { useEffect, useState } from 'react'; import { - Button, Card, Form, Input, InputNumber, message, Select, Space, Switch, Tabs, Typography, Upload, + Button, Card, Form, Input, InputNumber, message, Modal, Select, Space, Switch, Tabs, Typography, Upload, Table, Tag, Popconfirm, } from 'antd'; import { - SettingOutlined, SaveOutlined, UploadOutlined, FilePdfOutlined, EyeOutlined, DatabaseOutlined, VideoCameraOutlined, RobotOutlined, + SettingOutlined, SaveOutlined, UploadOutlined, FilePdfOutlined, EyeOutlined, DatabaseOutlined, VideoCameraOutlined, RobotOutlined, BankOutlined, PlusOutlined, EditOutlined, DeleteOutlined, } from '@ant-design/icons'; import { createSystemConfig, @@ -14,8 +14,12 @@ import { uploadLogo, uploadPdf, uploadLoginVideo, + listBankAccounts, + createBankAccount, + updateBankAccount, + deleteBankAccount, } from '../api'; -import type { ResourceCapacityUnit, SystemConfig } from '../types'; +import type { ResourceCapacityUnit, SystemConfig, BankAccount } from '../types'; const capacityUnitOptions: { value: ResourceCapacityUnit; label: string }[] = [ { value: 'MB', label: 'MB(1024 × 1024 字节)' }, @@ -30,10 +34,116 @@ const AdminSettings: React.FC = () => { const [uploading, setUploading] = useState(''); const [form] = Form.useForm(); + // 银行账户管理 + const [bankAccounts, setBankAccounts] = useState([]); + const [loadingAccounts, setLoadingAccounts] = useState(false); + const [accountModal, setAccountModal] = useState(false); + const [accountEditing, setAccountEditing] = useState(null); + const [accountForm] = Form.useForm(); + const [accountSaving, setAccountSaving] = useState(false); + useEffect(() => { load(); + loadBankAccounts(); }, []); + const loadBankAccounts = async () => { + setLoadingAccounts(true); + try { + const res = await listBankAccounts(); + setBankAccounts(res.items || []); + } catch (e: any) { + message.error(e?.message || '加载银行账户失败'); + } finally { + setLoadingAccounts(false); + } + }; + + const openAccountModal = (account?: BankAccount) => { + if (account) { + setAccountEditing(account); + accountForm.setFieldsValue({ + account_name: account.account_name, + bank_name: account.bank_name, + account_no: account.account_no, + is_active: account.is_active, + is_default: account.is_default, + description: account.description, + }); + } else { + setAccountEditing(null); + accountForm.resetFields(); + accountForm.setFieldsValue({ is_active: true, is_default: false }); + } + setAccountModal(true); + }; + + const handleAccountSave = async () => { + try { + const values = await accountForm.validateFields(); + setAccountSaving(true); + if (accountEditing) { + await updateBankAccount(accountEditing.id, values); + message.success('更新成功'); + } else { + await createBankAccount(values); + message.success('创建成功'); + } + setAccountModal(false); + await loadBankAccounts(); + } catch (e: any) { + message.error(e?.message || '保存失败'); + } finally { + setAccountSaving(false); + } + }; + + const handleAccountDelete = async (accountId: string) => { + try { + await deleteBankAccount(accountId); + message.success('删除成功'); + await loadBankAccounts(); + } catch (e: any) { + message.error(e?.message || '删除失败'); + } + }; + + const accountColumns = [ + { title: '账户名称', dataIndex: 'account_name', width: 150 }, + { title: '开户银行', dataIndex: 'bank_name', width: 150 }, + { + title: '银行账号', + dataIndex: 'account_no', + width: 180, + render: (v: string) => {v}, + }, + { + title: '状态', + dataIndex: 'is_active', + width: 80, + render: (v: boolean) => {v ? '启用' : '禁用'}, + }, + { + title: '默认', + dataIndex: 'is_default', + width: 80, + render: (v: boolean) => v && 默认, + }, + { title: '备注', dataIndex: 'description', ellipsis: true, render: (v: string) => v || '-' }, + { + title: '操作', + width: 140, + render: (_: any, record: BankAccount) => ( + + } onClick={() => openAccountModal(record)}>编辑 + handleAccountDelete(record.id)} okText="确定" cancelText="取消"> + }>删除 + + + ), + }, + ]; + const load = async () => { setLoading(true); try { @@ -49,6 +159,10 @@ const AdminSettings: React.FC = () => { if (!data.some(c => c.key === 'single_device_login_enabled')) { data.push({ id: 'cfg_single_device_login_enabled', key: 'single_device_login_enabled', value: 'false', description: '启用单设备登录(同端互斥):同一设备类型只允许一个登录会话' }); } + // 确保 OA 地址配置存在 + if (!data.some(c => c.key === 'oa_url')) { + data.push({ id: 'cfg_oa_url', key: 'oa_url', value: '', description: 'OA系统地址' }); + } setConfigs(data); const formValues: Record = {}; data.forEach(c => { formValues[c.key] = c.value; }); @@ -203,6 +317,7 @@ const AdminSettings: React.FC = () => { '协议配置': configs.filter(c => c.key === 'user_agreement_privacy_url'), 'SEO 设置': configs.filter(c => c.key.startsWith('seo_')), '用户积分配置': configs.filter(c => c.key.startsWith('user_') && c.key.includes('credits')), + '系统对接': configs.filter(c => c.key.startsWith('oa_')), '其他配置': configs.filter(c => c.key === 'operation_manual'), }; @@ -413,6 +528,52 @@ const AdminSettings: React.FC = () => { ), }, + { + key: 'integration', + label: '系统对接', + children: ( + + + + + 系统对接 + + {groupedConfigs['系统对接']?.map(config => ( + {config.description}} + extra={getFieldDescription(config)} + > + {getFieldComponent(config)} + + ))} + + + + {/* 银行账户管理 */} + + + + + 银行账户管理 + + } size="small" onClick={() => openAccountModal()}> + 新增账户 + + + + + + ), + }, { key: 'other', label: '其他配置', @@ -575,6 +736,41 @@ const AdminSettings: React.FC = () => { 保存配置 + + {/* 银行账户新增/编辑弹窗 */} + setAccountModal(false)} + confirmLoading={accountSaving} + okText="保存" + cancelText="取消" + destroyOnClose + > + + + + + + + + + + + + + + + + + + + + + + + ); }; diff --git a/video-gen-admin/src/types/index.ts b/video-gen-admin/src/types/index.ts index f5a3d074..ca7c36e0 100644 --- a/video-gen-admin/src/types/index.ts +++ b/video-gen-admin/src/types/index.ts @@ -1610,3 +1610,34 @@ export interface InvoiceDetail { updatedAt: string | null; orders: InvoiceOrder[]; } + +// ── Bank Account ─────────────────────────────────────── + +export interface BankAccount { + id: string; + account_name: string; + bank_name: string; + account_no: string; + is_active: boolean; + is_default: boolean; + description?: string | null; + created_at?: string | null; + updated_at?: string | null; +} + +// ── Scheduled Task ───────────────────────────────────── + +export interface ScheduledTask { + id: string; + name: string; + task_type: 'external_api' | 'internal_method'; + schedule: string; + config?: string | null; + is_active: boolean; + last_run_at?: string | null; + last_status?: string | null; + last_error?: string | null; + created_by?: string | null; + created_at?: string | null; + updated_at?: string | null; +} diff --git a/video-gen-api/alembic/versions/20260813_银行账户与定时任务.py b/video-gen-api/alembic/versions/20260813_银行账户与定时任务.py new file mode 100644 index 00000000..fc878974 --- /dev/null +++ b/video-gen-api/alembic/versions/20260813_银行账户与定时任务.py @@ -0,0 +1,53 @@ +"""银行账户与定时任务 + +Revision ID: 20260813_bank_account_scheduled_task +Revises: 20260813_split_device_type +Create Date: 2026-08-13 16:00:00.000000 +""" +from alembic import op +import sqlalchemy as sa + +# revision identifiers, used by Alembic. +revision = '20260813_bank_account_scheduled_task' +down_revision = '20260813_split_device_type' +branch_labels = None +depends_on = None + + +def upgrade(): + # 创建银行账户表 + op.create_table( + 'bank_accounts', + sa.Column('id', sa.String(32), primary_key=True), + sa.Column('account_name', sa.String(128), nullable=False, comment='账户名称'), + sa.Column('bank_name', sa.String(128), nullable=False, comment='开户银行'), + sa.Column('account_no', sa.String(64), unique=True, nullable=False, comment='银行账号'), + sa.Column('is_active', sa.Boolean, default=True, comment='是否启用'), + sa.Column('is_default', sa.Boolean, default=False, comment='是否默认账户'), + sa.Column('description', sa.String(256), nullable=True, comment='备注'), + sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now(), comment='创建时间'), + sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.func.now(), onupdate=sa.func.now(), comment='更新时间'), + ) + + # 创建定时任务表 + op.create_table( + 'scheduled_tasks', + sa.Column('id', sa.String(32), primary_key=True), + sa.Column('name', sa.String(128), nullable=False, comment='任务名称'), + sa.Column('task_type', sa.String(32), nullable=False, comment='类型: external_api / internal_method'), + sa.Column('schedule', sa.String(128), nullable=False, comment='Cron 表达式或间隔秒数'), + sa.Column('config', sa.Text, nullable=True, comment='任务配置 JSON'), + sa.Column('is_active', sa.Boolean, default=True, comment='是否启用'), + sa.Column('last_run_at', sa.String(64), nullable=True, comment='最后执行时间 ISO'), + sa.Column('last_status', sa.String(16), nullable=True, comment='最后执行状态'), + sa.Column('last_error', sa.Text, nullable=True, comment='最后执行错误信息'), + sa.Column('created_by', sa.String(32), nullable=True, comment='创建者管理员 ID'), + sa.Column('created_at', sa.DateTime(timezone=True), server_default=sa.func.now(), comment='创建时间'), + sa.Column('updated_at', sa.DateTime(timezone=True), server_default=sa.func.now(), onupdate=sa.func.now(), comment='更新时间'), + ) + + +def downgrade(): + # 删除新表 + op.drop_table('scheduled_tasks') + op.drop_table('bank_accounts') diff --git a/video-gen-api/app/admin_api/bank/__init__.py b/video-gen-api/app/admin_api/bank/__init__.py new file mode 100644 index 00000000..f4b8517a --- /dev/null +++ b/video-gen-api/app/admin_api/bank/__init__.py @@ -0,0 +1,3 @@ +from app.admin_api.bank.routes import router + +__all__ = ["router"] diff --git a/video-gen-api/app/admin_api/bank/routes.py b/video-gen-api/app/admin_api/bank/routes.py new file mode 100644 index 00000000..91668c3e --- /dev/null +++ b/video-gen-api/app/admin_api/bank/routes.py @@ -0,0 +1,173 @@ +"""银行账户管理与交易查询后台路由。""" + +from fastapi import APIRouter, Depends, HTTPException, Path, Query +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.dependencies import get_admin_user, get_db +from app.models.bank_account import BankAccount +from app.models.user import User +from app.services.bank.service import query_transactions_with_log +from app.utils.id_gen import generate_id + +router = APIRouter(prefix="/admin/bank", tags=["admin-bank"]) + + +# ============================================================ +# 银行账户 CRUD +# ============================================================ + + +@router.get("/accounts") +async def list_accounts( + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """列出所有银行账户。""" + result = await db.execute(select(BankAccount).order_by(BankAccount.is_default.desc(), BankAccount.created_at.desc())) + accounts = result.scalars().all() + return {"items": [ + { + "id": a.id, + "account_name": a.account_name, + "bank_name": a.bank_name, + "account_no": a.account_no, + "is_active": a.is_active, + "is_default": a.is_default, + "description": a.description, + "created_at": a.created_at.isoformat() if a.created_at else None, + "updated_at": a.updated_at.isoformat() if a.updated_at else None, + } + for a in accounts + ]} + + +@router.post("/accounts") +async def create_account( + body: dict, + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """新增银行账户。""" + account_name = (body.get("account_name") or "").strip() + bank_name = (body.get("bank_name") or "").strip() + account_no = (body.get("account_no") or "").strip() + if not account_name or not bank_name or not account_no: + raise HTTPException(status_code=400, detail="账户名称、开户银行、银行账号不能为空") + + # 检查账号唯一性 + existing = await db.execute(select(BankAccount).where(BankAccount.account_no == account_no)) + if existing.scalar_one_or_none(): + raise HTTPException(status_code=409, detail="该银行账号已存在") + + is_default = bool(body.get("is_default", False)) + # 如果设为默认,取消其他默认 + if is_default: + await db.execute( + BankAccount.__table__.update().where(BankAccount.is_default.is_(True)).values(is_default=False) + ) + + account = BankAccount( + id=generate_id(), + account_name=account_name, + bank_name=bank_name, + account_no=account_no, + is_active=bool(body.get("is_active", True)), + is_default=is_default, + description=body.get("description"), + ) + db.add(account) + await db.commit() + return {"id": account.id, "message": "创建成功"} + + +@router.put("/accounts/{account_id}") +async def update_account( + account_id: str = Path(..., description="账户 ID"), + body: dict = ..., + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """编辑银行账户。""" + result = await db.execute(select(BankAccount).where(BankAccount.id == account_id)) + account = result.scalar_one_or_none() + if account is None: + raise HTTPException(status_code=404, detail="账户不存在") + + if "account_name" in body: + account.account_name = str(body["account_name"]).strip() + if "bank_name" in body: + account.bank_name = str(body["bank_name"]).strip() + if "account_no" in body: + new_no = str(body["account_no"]).strip() + if new_no != account.account_no: + existing = await db.execute(select(BankAccount).where(BankAccount.account_no == new_no)) + if existing.scalar_one_or_none(): + raise HTTPException(status_code=409, detail="该银行账号已存在") + account.account_no = new_no + if "is_active" in body: + account.is_active = bool(body["is_active"]) + if "description" in body: + account.description = body.get("description") + + if body.get("is_default"): + await db.execute( + BankAccount.__table__.update() + .where(BankAccount.is_default.is_(True)) + .where(BankAccount.id != account_id) + .values(is_default=False) + ) + account.is_default = True + + await db.commit() + return {"message": "更新成功"} + + +@router.delete("/accounts/{account_id}") +async def delete_account( + account_id: str = Path(..., description="账户 ID"), + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """删除银行账户。""" + result = await db.execute(select(BankAccount).where(BankAccount.id == account_id)) + account = result.scalar_one_or_none() + if account is None: + raise HTTPException(status_code=404, detail="账户不存在") + await db.delete(account) + await db.commit() + return {"message": "删除成功"} + + +# ============================================================ +# 银行交易查询 +# ============================================================ + + +@router.get("/transactions") +async def list_transactions( + account_id: str = Query(..., description="银行账户 ID"), + start_date: str = Query(..., description="开始日期 (YYYY-MM-DD)"), + end_date: str = Query(..., description="结束日期 (YYYY-MM-DD)"), + dc_flag: int | None = Query(None, description="借贷方向: 0-借/出金, 1-贷/入金, 不传返回全部"), + page: int = Query(1, ge=1, description="页码"), + page_size: int = Query(20, ge=1, le=100, description="每页数量"), + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """查询银行交易流水(带接口请求记录到文件日志)。""" + result = await db.execute(select(BankAccount).where(BankAccount.id == account_id)) + account = result.scalar_one_or_none() + if account is None: + raise HTTPException(status_code=404, detail="银行账户不存在") + + data = await query_transactions_with_log( + admin.id, + acct_no=account.account_no, + start_date=start_date, + end_date=end_date, + dc_flag=dc_flag, + page=page, + page_size=page_size, + ) + return data diff --git a/video-gen-api/app/admin_api/scheduled_tasks/__init__.py b/video-gen-api/app/admin_api/scheduled_tasks/__init__.py new file mode 100644 index 00000000..571d3ff4 --- /dev/null +++ b/video-gen-api/app/admin_api/scheduled_tasks/__init__.py @@ -0,0 +1,3 @@ +from app.admin_api.scheduled_tasks.routes import router + +__all__ = ["router"] diff --git a/video-gen-api/app/admin_api/scheduled_tasks/routes.py b/video-gen-api/app/admin_api/scheduled_tasks/routes.py new file mode 100644 index 00000000..c432932e --- /dev/null +++ b/video-gen-api/app/admin_api/scheduled_tasks/routes.py @@ -0,0 +1,170 @@ +"""定时任务管理后台路由。""" + +import json + +from fastapi import APIRouter, Depends, HTTPException, Path +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.dependencies import get_admin_user, get_db +from app.models.scheduled_task import ScheduledTask +from app.models.user import User +from app.utils.id_gen import generate_id + +router = APIRouter(prefix="/admin/scheduled-tasks", tags=["admin-scheduled-tasks"]) + + +def _task_to_dict(task: ScheduledTask) -> dict: + return { + "id": task.id, + "name": task.name, + "task_type": task.task_type, + "schedule": task.schedule, + "config": task.config, + "is_active": task.is_active, + "last_run_at": task.last_run_at, + "last_status": task.last_status, + "last_error": task.last_error, + "created_by": task.created_by, + "created_at": task.created_at.isoformat() if task.created_at else None, + "updated_at": task.updated_at.isoformat() if task.updated_at else None, + } + + +@router.get("") +async def list_tasks( + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """列出所有定时任务。""" + result = await db.execute(select(ScheduledTask).order_by(ScheduledTask.created_at.desc())) + tasks = result.scalars().all() + return {"items": [_task_to_dict(t) for t in tasks]} + + +@router.post("") +async def create_task( + body: dict, + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """创建定时任务。""" + name = (body.get("name") or "").strip() + task_type = (body.get("task_type") or "").strip() + schedule = (body.get("schedule") or "").strip() + if not name or not task_type or not schedule: + raise HTTPException(status_code=400, detail="任务名称、类型、调度表达式不能为空") + if task_type not in ("external_api", "internal_method"): + raise HTTPException(status_code=400, detail="任务类型必须为 external_api 或 internal_method") + + config = body.get("config") + if isinstance(config, dict): + config = json.dumps(config, ensure_ascii=False) + elif isinstance(config, str): + # 验证 JSON 合法性 + try: + json.loads(config) + except json.JSONDecodeError: + raise HTTPException(status_code=400, detail="config 不是合法的 JSON") + + task = ScheduledTask( + id=generate_id(), + name=name, + task_type=task_type, + schedule=schedule, + config=config, + is_active=bool(body.get("is_active", True)), + created_by=admin.id, + ) + db.add(task) + await db.commit() + return {"id": task.id, "message": "创建成功"} + + +@router.put("/{task_id}") +async def update_task( + task_id: str = Path(..., description="任务 ID"), + body: dict = ..., + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """更新定时任务。""" + result = await db.execute(select(ScheduledTask).where(ScheduledTask.id == task_id)) + task = result.scalar_one_or_none() + if task is None: + raise HTTPException(status_code=404, detail="任务不存在") + + if "name" in body: + task.name = str(body["name"]).strip() + if "task_type" in body: + t = body["task_type"] + if t not in ("external_api", "internal_method"): + raise HTTPException(status_code=400, detail="任务类型必须为 external_api 或 internal_method") + task.task_type = t + if "schedule" in body: + task.schedule = str(body["schedule"]).strip() + if "is_active" in body: + task.is_active = bool(body["is_active"]) + if "config" in body: + config = body["config"] + if isinstance(config, dict): + config = json.dumps(config, ensure_ascii=False) + elif isinstance(config, str): + try: + json.loads(config) + except json.JSONDecodeError: + raise HTTPException(status_code=400, detail="config 不是合法的 JSON") + task.config = config + + await db.commit() + return {"message": "更新成功"} + + +@router.delete("/{task_id}") +async def delete_task( + task_id: str = Path(..., description="任务 ID"), + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """删除定时任务。""" + result = await db.execute(select(ScheduledTask).where(ScheduledTask.id == task_id)) + task = result.scalar_one_or_none() + if task is None: + raise HTTPException(status_code=404, detail="任务不存在") + await db.delete(task) + await db.commit() + return {"message": "删除成功"} + + +@router.post("/{task_id}/run") +async def run_task( + task_id: str = Path(..., description="任务 ID"), + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """手动执行一次定时任务。""" + result = await db.execute(select(ScheduledTask).where(ScheduledTask.id == task_id)) + task = result.scalar_one_or_none() + if task is None: + raise HTTPException(status_code=404, detail="任务不存在") + + from app.tasks.scheduled_tasks import execute_scheduled_task + + execute_scheduled_task.apply_async(args=[task_id]) + return {"message": "任务已提交执行"} + + +@router.post("/{task_id}/toggle") +async def toggle_task( + task_id: str = Path(..., description="任务 ID"), + admin: User = Depends(get_admin_user), + db: AsyncSession = Depends(get_db), +): + """启用/禁用定时任务。""" + result = await db.execute(select(ScheduledTask).where(ScheduledTask.id == task_id)) + task = result.scalar_one_or_none() + if task is None: + raise HTTPException(status_code=404, detail="任务不存在") + task.is_active = not task.is_active + await db.commit() + return {"is_active": task.is_active, "message": "已启用" if task.is_active else "已禁用"} diff --git a/video-gen-api/app/api/admin/__init__.py b/video-gen-api/app/api/admin/__init__.py index f8bb4294..3013af6c 100644 --- a/video-gen-api/app/api/admin/__init__.py +++ b/video-gen-api/app/api/admin/__init__.py @@ -15,6 +15,8 @@ from app.api.admin.contact import router as admin_contact_router from app.admin_api.api_keys import router as api_keys_admin_router from app.admin_api.api_model_pricings import router as api_model_pricings_admin_router from app.admin_api.vp_v3_quota import router as vp_v3_quota_admin_router +from app.admin_api.bank import router as bank_admin_router +from app.admin_api.scheduled_tasks import router as scheduled_tasks_admin_router router = APIRouter() router.include_router(video_prompt_schema_config_router) @@ -32,3 +34,5 @@ router.include_router(admin_contact_router) router.include_router(api_keys_admin_router) router.include_router(api_model_pricings_admin_router) router.include_router(vp_v3_quota_admin_router) +router.include_router(bank_admin_router) +router.include_router(scheduled_tasks_admin_router) diff --git a/video-gen-api/app/config.py b/video-gen-api/app/config.py index fcc401ec..1d17c06c 100644 --- a/video-gen-api/app/config.py +++ b/video-gen-api/app/config.py @@ -131,6 +131,10 @@ class Settings(BaseSettings): CAPTCHA_ENABLED: bool = True + # 银行交易查询接口配置 + BANK_API_BASE: str = "" + BANK_API_KEY: str = "" + BASE_URL: str = "" CORS_ORIGINS: list[str] = ["*"] diff --git a/video-gen-api/app/models/__init__.py b/video-gen-api/app/models/__init__.py index e5a95954..6a7d7342 100644 --- a/video-gen-api/app/models/__init__.py +++ b/video-gen-api/app/models/__init__.py @@ -43,6 +43,8 @@ from app.models.home_material import HomeMaterialAsset, HomeMaterialCategory, Ho from app.models.contact_request import ContactRequest from app.models.invoice import Invoice, InvoiceOrder from app.models.invoice_header import InvoiceHeader +from app.models.bank_account import BankAccount +from app.models.scheduled_task import ScheduledTask from app.models.private_portrait import PrivatePortraitProject, PrivatePortraitValidateSession, PrivatePortraitAssetGroup, PrivatePortraitAsset from app.models.api import ApiKey, ApiGenerationTask, ApiUsageLog, ApiKeyUpscaleConfig, ApiUpscaleLink @@ -68,4 +70,5 @@ __all__ = [ "ApiKey", "ApiGenerationTask", "ApiUsageLog", "ApiKeyUpscaleConfig", "ApiUpscaleLink", "ApiModelPricing", "Invoice", "InvoiceOrder", "InvoiceHeader", + "BankAccount", "ScheduledTask", ] diff --git a/video-gen-api/app/models/bank_account.py b/video-gen-api/app/models/bank_account.py new file mode 100644 index 00000000..acceb0dc --- /dev/null +++ b/video-gen-api/app/models/bank_account.py @@ -0,0 +1,18 @@ +"""银行账户信息模型。""" + +from sqlalchemy import Boolean, String +from sqlalchemy.orm import Mapped, mapped_column + +from app.models.base import Base, TimestampMixin + + +class BankAccount(Base, TimestampMixin): + __tablename__ = "bank_accounts" + + id: Mapped[str] = mapped_column(String(32), primary_key=True) + account_name: Mapped[str] = mapped_column(String(128), nullable=False, comment="账户名称") + bank_name: Mapped[str] = mapped_column(String(128), nullable=False, comment="开户银行") + account_no: Mapped[str] = mapped_column(String(64), unique=True, nullable=False, comment="银行账号") + is_active: Mapped[bool] = mapped_column(Boolean, default=True, comment="是否启用") + is_default: Mapped[bool] = mapped_column(Boolean, default=False, comment="是否默认账户") + description: Mapped[str | None] = mapped_column(String(256), nullable=True, comment="备注") diff --git a/video-gen-api/app/models/scheduled_task.py b/video-gen-api/app/models/scheduled_task.py new file mode 100644 index 00000000..e3d04943 --- /dev/null +++ b/video-gen-api/app/models/scheduled_task.py @@ -0,0 +1,23 @@ +"""定时任务配置模型。""" + +from sqlalchemy import Boolean, String, Text +from sqlalchemy.orm import Mapped, mapped_column + +from app.models.base import Base, TimestampMixin + + +class ScheduledTask(Base, TimestampMixin): + __tablename__ = "scheduled_tasks" + + id: Mapped[str] = mapped_column(String(32), primary_key=True) + name: Mapped[str] = mapped_column(String(128), nullable=False, comment="任务名称") + task_type: Mapped[str] = mapped_column(String(32), nullable=False, comment="类型: external_api / internal_method") + schedule: Mapped[str] = mapped_column(String(128), nullable=False, comment="Cron 表达式或间隔秒数") + config: Mapped[str | None] = mapped_column(Text, nullable=True, comment="任务配置 JSON") + # external_api config: {url, method, headers, payload} + # internal_method config: {module, function, args} + is_active: Mapped[bool] = mapped_column(Boolean, default=True, comment="是否启用") + last_run_at: Mapped[str | None] = mapped_column(String(64), nullable=True, comment="最后执行时间 ISO") + last_status: Mapped[str | None] = mapped_column(String(16), nullable=True, comment="最后执行状态: success / error") + last_error: Mapped[str | None] = mapped_column(Text, nullable=True, comment="最后执行错误信息") + created_by: Mapped[str | None] = mapped_column(String(32), nullable=True, comment="创建者管理员 ID") diff --git a/video-gen-api/app/services/bank/__init__.py b/video-gen-api/app/services/bank/__init__.py new file mode 100644 index 00000000..b8b4fc4f --- /dev/null +++ b/video-gen-api/app/services/bank/__init__.py @@ -0,0 +1,4 @@ +from app.services.bank.client import query_bank_transactions, BankApiError +from app.services.bank.service import query_transactions_with_log + +__all__ = ["query_bank_transactions", "BankApiError", "query_transactions_with_log"] diff --git a/video-gen-api/app/services/bank/client.py b/video-gen-api/app/services/bank/client.py new file mode 100644 index 00000000..203ec502 --- /dev/null +++ b/video-gen-api/app/services/bank/client.py @@ -0,0 +1,76 @@ +"""银行交易查询 API 客户端。""" + +import httpx + +from app.config import settings + + +class BankApiError(RuntimeError): + """银行接口错误。""" + + def __init__( + self, + message: str, + retryable: bool = False, + http_status: int | None = None, + ): + super().__init__(message) + self.retryable = retryable + self.http_status = http_status + + +async def query_bank_transactions( + acct_no: str, + start_date: str, + end_date: str, + dc_flag: int | None = None, + page: int = 1, + page_size: int = 20, +) -> dict: + """查询银行交易流水。 + + 调用外部银行接口,按账号和日期范围查询交易记录。 + + :param acct_no: 我方银行账号 + :param start_date: 开始日期 (YYYY-MM-DD) + :param end_date: 结束日期 (YYYY-MM-DD) + :param dc_flag: 借贷方向 (0-借/出金, 1-贷/入金, None-全部) + :param page: 页码 + :param page_size: 每页数量 (最大 100) + :return: 接口返回的原始数据 + """ + base = (settings.BANK_API_BASE or "").rstrip("/") + if not base: + raise BankApiError("银行接口地址未配置 (BANK_API_BASE)", retryable=False) + + url = f"{base}/api/v1/internal/get-blank-transfer-acct-time" + headers = { + "api-key": settings.BANK_API_KEY or "", + "Content-Type": "application/json", + } + payload: dict = { + "acct_no": acct_no, + "start_date": start_date, + "end_date": end_date, + "page": page, + "page_size": min(page_size, 100), + } + if dc_flag is not None: + payload["dc_flag"] = dc_flag + + timeout = httpx.Timeout(30.0) + try: + async with httpx.AsyncClient(timeout=timeout, follow_redirects=True) as client: + response = await client.post(url, headers=headers, json=payload) + response.raise_for_status() + return response.json() + except httpx.HTTPStatusError as e: + raise BankApiError( + f"银行接口返回错误: HTTP {e.response.status_code}", + retryable=e.response.status_code >= 500, + http_status=e.response.status_code, + ) from e + except httpx.TimeoutException as e: + raise BankApiError("银行接口请求超时", retryable=True) from e + except httpx.RequestError as e: + raise BankApiError(f"银行接口请求失败: {e}", retryable=True) from e diff --git a/video-gen-api/app/services/bank/file_logger.py b/video-gen-api/app/services/bank/file_logger.py new file mode 100644 index 00000000..6a0d9707 --- /dev/null +++ b/video-gen-api/app/services/bank/file_logger.py @@ -0,0 +1,50 @@ +"""银行接口请求文件日志记录器。 + +使用 get_logger 模式,日志输出到 log/bank_api/bank_api-YYYY-MM-DD.log。 +""" + +import json +from datetime import datetime, timezone + +from app.utils.logger import get_logger + +bank_api_logger = get_logger("bank_api", "bank_api") + + +def log_bank_api_request( + acct_no: str, + request_params: dict | None, + response_data: dict | None, + status_code: int | None, + is_success: bool, + error_msg: str | None = None, + admin_id: str | None = None, + duration_ms: int | None = None, +) -> None: + """记录银行接口请求到日志文件。 + + :param acct_no: 查询的银行账号 + :param request_params: 请求参数 + :param response_data: 响应数据 + :param status_code: HTTP 状态码 + :param is_success: 是否成功 + :param error_msg: 错误信息 + :param admin_id: 操作管理员 ID + :param duration_ms: 请求耗时(毫秒) + """ + log_data = { + "timestamp": datetime.now(timezone.utc).isoformat(), + "acct_no": acct_no, + "request_params": request_params, + "response_data": response_data, + "status_code": status_code, + "is_success": is_success, + "error_msg": error_msg, + "admin_id": admin_id, + "duration_ms": duration_ms, + } + line = json.dumps(log_data, ensure_ascii=False, default=str) + if is_success: + bank_api_logger.info(line) + else: + bank_api_logger.error(line) diff --git a/video-gen-api/app/services/bank/service.py b/video-gen-api/app/services/bank/service.py new file mode 100644 index 00000000..8ea6c03f --- /dev/null +++ b/video-gen-api/app/services/bank/service.py @@ -0,0 +1,47 @@ +"""银行交易查询服务 — 接口请求日志写入文件。""" + +import time + +from app.services.bank.client import query_bank_transactions, BankApiError +from app.services.bank.file_logger import log_bank_api_request + + +async def query_transactions_with_log( + admin_id: str, + **kwargs, +) -> dict: + """查询银行流水并记录接口日志到文件。 + + :param admin_id: 操作管理员 ID + :param kwargs: 透传给 query_bank_transactions 的参数 + :return: 接口返回的原始数据 + """ + acct_no = kwargs.get("acct_no", "") + request_params = {k: v for k, v in kwargs.items()} + start_time = time.monotonic() + try: + result = await query_bank_transactions(**kwargs) + duration_ms = int((time.monotonic() - start_time) * 1000) + log_bank_api_request( + acct_no=acct_no, + request_params=request_params, + response_data=result, + status_code=200, + is_success=True, + admin_id=admin_id, + duration_ms=duration_ms, + ) + return result + except BankApiError as e: + duration_ms = int((time.monotonic() - start_time) * 1000) + log_bank_api_request( + acct_no=acct_no, + request_params=request_params, + response_data=None, + status_code=e.http_status, + is_success=False, + error_msg=str(e), + admin_id=admin_id, + duration_ms=duration_ms, + ) + raise diff --git a/video-gen-api/app/tasks/celery_app.py b/video-gen-api/app/tasks/celery_app.py index f1d8f246..1b8efc87 100644 --- a/video-gen-api/app/tasks/celery_app.py +++ b/video-gen-api/app/tasks/celery_app.py @@ -40,6 +40,7 @@ CELERY_TASK_IMPORTS = ( "app.tasks.api_generation_tasks", "app.tasks.api_recovery_tasks", "app.tasks.api_upscale_tasks", + "app.tasks.scheduled_tasks", ) @@ -292,6 +293,85 @@ else: celery_app = None +def _parse_schedule_to_celery(schedule_str: str): + """将 schedule 字符串解析为 Celery 可识别的调度值。 + + - 纯数字:视为间隔秒数(返回 int) + - cron 表达式 (5 字段空格分隔):返回 crontab 对象 + """ + from celery.schedules import crontab + + s = (schedule_str or "").strip() + if not s: + return None + # 纯数字 → 间隔秒数 + if s.isdigit(): + return int(s) + # cron 表达式 (分 时 日 月 周) + parts = s.split() + if len(parts) == 5: + try: + return crontab( + minute=parts[0], + hour=parts[1], + day_of_month=parts[2], + month_of_year=parts[3], + day_of_week=parts[4], + ) + except Exception: + logger.exception("解析 cron 表达式失败: %s", s) + return None + logger.warning("无法解析 schedule 表达式: %s", s) + return None + + +@celery_app.on_after_configure.connect # type: ignore +def _setup_dynamic_beat_tasks(sender, **kwargs): + """从数据库加载活跃定时任务并注册到 Beat 调度。 + + 通过 @celery_app.on_after_configure.connect 在 Celery 配置完成后执行, + 适用于 Worker 和 Beat 启动场景。 + """ + if celery_app is None: + return + + async def _load(): + from sqlalchemy import select + + from app.models.base import async_session + from app.models.scheduled_task import ScheduledTask + + async with async_session() as db: + result = await db.execute( + select(ScheduledTask).where(ScheduledTask.is_active.is_(True)) + ) + tasks = result.scalars().all() + return tasks + + try: + active_tasks = run_async(_load()) + except Exception: + logger.exception("加载定时任务失败,跳过动态 Beat 注册") + return + + for task in active_tasks: + schedule_val = _parse_schedule_to_celery(task.schedule) + if schedule_val is None: + logger.warning("定时任务 %s schedule 无效,跳过注册: %s", task.id, task.schedule) + continue + beat_key = f"dynamic-scheduled-task-{task.id}" + sender.conf.beat_schedule[beat_key] = { + "task": "execute_scheduled_task", + "schedule": schedule_val, + "args": (task.id,), + "options": {"queue": RECOVERY_QUEUE}, + } + logger.info( + "动态注册定时任务到 Beat: %s (%s) schedule=%s", + task.name, task.id, task.schedule, + ) + + async def _try_acquire_startup_recovery_lock() -> bool: """任意 worker 启动时都可尝试抢恢复投递锁,避免依赖 hostname 命名。""" from app.services.redis_registry_service import redis_acquire_lock diff --git a/video-gen-api/app/tasks/scheduled_tasks.py b/video-gen-api/app/tasks/scheduled_tasks.py new file mode 100644 index 00000000..81b2dace --- /dev/null +++ b/video-gen-api/app/tasks/scheduled_tasks.py @@ -0,0 +1,138 @@ +"""定时任务执行器。 + +支持两种任务类型: +- external_api: 调用外部 HTTP 接口 +- internal_method: 动态导入并执行内部函数 + +任务在 Celery Worker 中执行,通过 run_async 桥接异步操作。 +""" + +import json +import logging +import time +from datetime import datetime, timezone + +import httpx +from sqlalchemy import select + +from app.models.base import async_session +from app.models.scheduled_task import ScheduledTask +from app.tasks.async_runner import run_async +from app.tasks.celery_app import celery_app + +logger = logging.getLogger("video_gen") + + +def _update_task_status(task_id: str, status: str, error_msg: str | None = None) -> None: + """更新任务最后执行状态。""" + + async def _do(): + async with async_session() as db: + result = await db.execute(select(ScheduledTask).where(ScheduledTask.id == task_id)) + task = result.scalar_one_or_none() + if task is None: + return + task.last_run_at = datetime.now(timezone.utc).isoformat() + task.last_status = status + task.last_error = error_msg + await db.commit() + + run_async(_do()) + + +def _execute_external_api(config: str | None) -> dict: + """执行外部 API 调用。""" + cfg = json.loads(config or "{}") + url = cfg.get("url", "").strip() + method = cfg.get("method", "GET").upper() + headers = cfg.get("headers") or {} + payload = cfg.get("payload") + timeout_val = float(cfg.get("timeout", 30)) + + if not url: + raise ValueError("外部接口地址 (url) 未配置") + + start = time.monotonic() + try: + with httpx.Client(timeout=httpx.Timeout(timeout_val, connect=10.0)) as client: + if method == "GET": + resp = client.get(url, headers=headers, params=payload) + elif method == "DELETE": + resp = client.delete(url, headers=headers, params=payload) + else: + resp = client.request(method, url, headers=headers, json=payload) + duration_ms = int((time.monotonic() - start) * 1000) + resp.raise_for_status() + return { + "status_code": resp.status_code, + "duration_ms": duration_ms, + "body": resp.text[:2000], + } + except httpx.HTTPError as e: + duration_ms = int((time.monotonic() - start) * 1000) + raise RuntimeError(f"外部接口请求失败: {e} (耗时 {duration_ms}ms)") from e + + +def _execute_internal_method(config: str | None) -> dict: + """执行内部方法调用。""" + cfg = json.loads(config or "{}") + module_path = cfg.get("module", "").strip() + function_name = cfg.get("function", "").strip() + args = cfg.get("args") or [] + + if not module_path or not function_name: + raise ValueError("内部方法需要指定 module 和 function") + + import importlib + + module = importlib.import_module(module_path) + func = getattr(module, function_name, None) + if func is None or not callable(func): + raise ValueError(f"模块 {module_path} 中不存在可调用函数 {function_name}") + + start = time.monotonic() + result = func(*args) + duration_ms = int((time.monotonic() - start) * 1000) + return { + "duration_ms": duration_ms, + "result": str(result)[:1000] if result is not None else None, + } + + +@celery_app.task(name="execute_scheduled_task", bind=True, ignore_result=True) # type: ignore[call-arg] +def execute_scheduled_task(self, task_id: str): + """执行定时任务(Celery 任务入口)。 + + 通过 run_async 桥接到异步上下文读取任务配置并执行。 + """ + + async def _run(): + async with async_session() as db: + result = await db.execute(select(ScheduledTask).where(ScheduledTask.id == task_id)) + task = result.scalar_one_or_none() + if task is None: + logger.warning("定时任务不存在: %s", task_id) + return + if not task.is_active: + logger.info("定时任务已禁用,跳过执行: %s", task_id) + return + + task_config = task.config + task_type = task.task_type + + try: + if task_type == "external_api": + exec_result = _execute_external_api(task_config) + elif task_type == "internal_method": + exec_result = _execute_internal_method(task_config) + else: + raise ValueError(f"未知的任务类型: {task_type}") + + _update_task_status(task_id, "success") + logger.info("定时任务执行成功: %s (%s) -> %s", task_id, task_type, exec_result) + except Exception as e: + error_msg = str(e) + _update_task_status(task_id, "error", error_msg) + logger.exception("定时任务执行失败: %s (%s)", task_id, task_type) + + run_async(_run())
{v}
* * * * *
0 * * * *