diff --git a/apps/backend/src/routes/admin.ts b/apps/backend/src/routes/admin.ts index 22b96d89..d26bde61 100644 --- a/apps/backend/src/routes/admin.ts +++ b/apps/backend/src/routes/admin.ts @@ -9,6 +9,7 @@ import filesRouter from './admin/files.js' import kuaishouIndustryRouter from './admin/kuaishou-industry.js' import ordersRouter from './admin/orders.js' import notificationsRouter from './admin/notifications.js' +import realtimeRouter from './admin/realtime.js' import platformConfigRouter from './admin/platform-config.js' import { requireAdminSession } from './admin/session.js' import tasksRouter from './admin/tasks.js' @@ -28,6 +29,7 @@ router.use(platformConfigRouter) router.use(kuaishouIndustryRouter) router.use(ordersRouter) router.use(notificationsRouter) +router.use(realtimeRouter) router.use(tasksRouter) router.use(cloudtentaclesRecordsRouter) router.use(workerPlatformRouter) diff --git a/apps/backend/src/routes/admin/realtime.ts b/apps/backend/src/routes/admin/realtime.ts new file mode 100644 index 00000000..26e2f5a9 --- /dev/null +++ b/apps/backend/src/routes/admin/realtime.ts @@ -0,0 +1,12 @@ +import { Router } from 'express' + +import { openAdminRealtimeStream } from '../../services/realtime/realtime-event-service.js' +import { requireAdminRoles } from './session.js' + +const router = Router() + +router.get('/realtime', requireAdminRoles(['admin', 'operator', 'support']), (_req, res) => { + openAdminRealtimeStream(res) +}) + +export default router diff --git a/apps/backend/src/routes/worker.ts b/apps/backend/src/routes/worker.ts index 89925cfa..46192485 100644 --- a/apps/backend/src/routes/worker.ts +++ b/apps/backend/src/routes/worker.ts @@ -37,6 +37,7 @@ import { requireActiveWorker, requireWorkerSession, } from './worker/session.js' +import { openWorkerRealtimeStream } from '../services/realtime/realtime-event-service.js' const router = Router() @@ -105,6 +106,10 @@ router.post( router.use(requireWorkerSession) +router.get('/realtime', (req, res) => { + openWorkerRealtimeStream(res, getRequiredWorkerSession(req).workerId) +}) + router.get( '/profile', createRouteHandler((req) => getWorkerProfile(getRequiredWorkerSession(req)), { diff --git a/apps/backend/src/services/admin/admin-notification-service.ts b/apps/backend/src/services/admin/admin-notification-service.ts index 0cf21752..7b3dc3f3 100644 --- a/apps/backend/src/services/admin/admin-notification-service.ts +++ b/apps/backend/src/services/admin/admin-notification-service.ts @@ -12,6 +12,7 @@ import { getWorkerPlatformNotificationConfig, type WorkerPlatformNotificationEventKey, } from '../worker-platform/worker-platform-notification-config-service.js' +import { publishRealtimeEvent } from '../realtime/realtime-event-service.js' const ADMIN_WORKER_PLATFORM_PATH = '/admin/worker-platform' @@ -116,6 +117,14 @@ export async function remindWorkerAcceptanceAdminNotification(workOrderId: numbe cooldownMs: 30 * 60 * 1000, }) } + if (notification?.visible) { + publishRealtimeEvent({ + type: 'admin_notification.changed', + entityId: Number(notification.id), + scopes: ['admin_notifications'], + admin: true, + }) + } return notification } @@ -134,6 +143,12 @@ export async function resolveAdminNotificationEntity( reason, now: new Date().toISOString(), }) + publishRealtimeEvent({ + type: 'admin_notification.changed', + entityId, + scopes: ['admin_notifications'], + admin: true, + }) } export async function listAdminPendingNotifications(limit = 20) { @@ -195,6 +210,14 @@ async function createConfiguredWorkerPlatformNotification( cooldownKey: notification.dedupe_key, }) } + if (notification?.visible) { + publishRealtimeEvent({ + type: 'admin_notification.changed', + entityId: Number(notification.id), + scopes: ['admin_notifications'], + admin: true, + }) + } return notification } diff --git a/apps/backend/src/services/order/order-service.ts b/apps/backend/src/services/order/order-service.ts index aa9a8064..d1be405e 100644 --- a/apps/backend/src/services/order/order-service.ts +++ b/apps/backend/src/services/order/order-service.ts @@ -11,6 +11,10 @@ import { syncOrderFulfillmentAttachments, } from '../fulfillment/order-fulfillment-readiness-service.js' import { syncWorkerOrdersForSourceOrder } from '../worker-platform/index.js' +import { + publishRealtimeEvent, + publishWorkOrderRealtimeChange, +} from '../realtime/realtime-event-service.js' import { nowIso } from '../../utils/time.js' import { logIntegration } from '../../utils/logger.js' import { createHttpError } from '../../utils/http.js' @@ -192,10 +196,11 @@ export async function upsertOrderFromSource( { level: 'warn' }, ) - await syncWorkerOrdersForSourceOrder(order, orderItems, { + const workOrderSync = await syncWorkerOrdersForSourceOrder(order, orderItems, { source: `${sourceLabel}_order_upsert`, autoOnly: true, }) + publishSourceOrderRealtimeChanges(order.id, workOrderSync.created) return { ignoreReason: readiness.blockReason || 'fulfillment_not_ready', order, @@ -212,10 +217,11 @@ export async function upsertOrderFromSource( source: `${sourceLabel}_order_upsert`, now, }) - await syncWorkerOrdersForSourceOrder(order, orderItems, { + const workOrderSync = await syncWorkerOrdersForSourceOrder(order, orderItems, { source: `${sourceLabel}_order_upsert`, autoOnly: true, }) + publishSourceOrderRealtimeChanges(order.id, workOrderSync.created) logIntegration('[order-service]', `${sourceLabel} 订单 upsert 完成`, { orderId: order.id, @@ -236,6 +242,21 @@ export async function upsertOrderFromSource( } } +function publishSourceOrderRealtimeChanges( + orderId: number, + workOrders: Array<{ workOrderId: number }> = [], +) { + publishRealtimeEvent({ + type: 'order.changed', + entityId: Number(orderId), + scopes: ['admin_orders'], + admin: true, + }) + for (const workOrder of workOrders) { + publishWorkOrderRealtimeChange({ workOrderId: workOrder.workOrderId }) + } +} + const ORDER_STATUS_PRIORITY = { created: 0, paid: 1, diff --git a/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts b/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts index 412dfb4d..43eef94c 100644 --- a/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts +++ b/apps/backend/src/services/platforms/kuaishou-industry/send-code-service.ts @@ -31,6 +31,10 @@ import { } from './voucher-service.js' import { bindKuaishouIndustryVouchersToOrderTasks } from './voucher-binding-service.js' import { parseAmountToFen } from '../../../utils/money.js' +import { + publishRealtimeEvent, + publishWorkOrderRealtimeChange, +} from '../../realtime/realtime-event-service.js' import type { KuaishouIndustryVoucherRow, OrderRow } from '../../../types/repository/rows.js' type SendCodeCallbackParams = { @@ -118,23 +122,36 @@ export async function handleSendCode(rawBody: JsonObject = {}) { }) } - await syncWorkOrdersFromKuaishouSendCode({ - oid: normalizedOid, - sellerId: params.sellerId, - itemId: params.itemId, - itemTitle: params.itemTitle, - skuId: params.skuId, - num: params.num, - paymentFen: resolveSendCodePaymentFen(params.ext), - rawParams: params, - vouchers, - now, - }).catch((error) => { + try { + const syncResult = await syncWorkOrdersFromKuaishouSendCode({ + oid: normalizedOid, + sellerId: params.sellerId, + itemId: params.itemId, + itemTitle: params.itemTitle, + skuId: params.skuId, + num: params.num, + paymentFen: resolveSendCodePaymentFen(params.ext), + rawParams: params, + vouchers, + now, + }) + if (syncResult.orderId) { + publishRealtimeEvent({ + type: 'order.changed', + entityId: syncResult.orderId, + scopes: ['admin_orders'], + admin: true, + }) + } + for (const workOrderId of syncResult.workOrderIds) { + publishWorkOrderRealtimeChange({ workOrderId }) + } + } catch (error) { logWarn('[kuaishou-industry/send-code]', '发单平台工单同步失败,不影响发码主流程', { oid: normalizedOid, error: error instanceof Error ? error.message : String(error), }) - }) + } scheduleKuaishouIndustrySendCallbackRetry(normalizedOid, SEND_CALLBACK_INITIAL_DELAY_MS, { source: 'send_code_accepted', diff --git a/apps/backend/src/services/realtime/realtime-event-service.ts b/apps/backend/src/services/realtime/realtime-event-service.ts new file mode 100644 index 00000000..d6aaf917 --- /dev/null +++ b/apps/backend/src/services/realtime/realtime-event-service.ts @@ -0,0 +1,219 @@ +import type { Response } from 'express' + +export type RealtimeEventType = + | 'admin_notification.changed' + | 'order.changed' + | 'work_order.changed' + | 'worker_finance_request.changed' + | 'worker_wallet.changed' + +export type RealtimeEvent = { + eventId: number + type: RealtimeEventType + entityId: number + scopes: string[] + occurredAt: string +} + +type RealtimeEventInput = { + type: RealtimeEventType + entityId?: number + scopes?: string[] + admin?: boolean + workerIds?: number[] + broadcastWorkers?: boolean +} + +type RealtimeSubscriber = { + id: number + response: Response + workerId: number | null +} + +const adminSubscribers = new Map() +const workerSubscribers = new Map() +let nextSubscriberId = 1 +let nextEventId = 1 +let heartbeatTimer: NodeJS.Timeout | null = null + +/** 建立后台实时事件流,事件仅用于提示前端重新读取权威数据。 */ +export function openAdminRealtimeStream(response: Response) { + return openRealtimeStream(response, null) +} + +/** 建立打手实时事件流,只向该打手推送其可见的订单和资金变化。 */ +export function openWorkerRealtimeStream(response: Response, workerId: number) { + return openRealtimeStream(response, Number(workerId)) +} + +/** 发布轻量实时事件;数据库仍是唯一权威数据源。 */ +export function publishRealtimeEvent(input: RealtimeEventInput) { + const event: RealtimeEvent = { + eventId: nextEventId++, + type: input.type, + entityId: Math.max(0, Number(input.entityId) || 0), + scopes: [ + ...new Set((input.scopes || []).map((item) => String(item || '').trim()).filter(Boolean)), + ], + occurredAt: new Date().toISOString(), + } + + if (input.admin) { + broadcast(adminSubscribers.values(), event) + } + + const workerIds = new Set( + (input.workerIds || []) + .map((item) => Number(item)) + .filter((item) => Number.isInteger(item) && item > 0), + ) + if (input.broadcastWorkers) { + broadcast(workerSubscribers.values(), event) + return event + } + + if (workerIds.size > 0) { + broadcast( + [...workerSubscribers.values()].filter((subscriber) => + workerIds.has(Number(subscriber.workerId)), + ), + event, + ) + } + return event +} + +/** 工单变更统一通知后台;大厅与打手本人按可见范围分别接收。 */ +export function publishWorkOrderRealtimeChange(input: { + workOrderId: number + workerIds?: number[] + hallChanged?: boolean +}) { + const workOrderId = Math.max(0, Number(input.workOrderId) || 0) + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['admin_work_orders'], + admin: true, + }) + if (input.hallChanged) { + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['worker_hall'], + broadcastWorkers: true, + }) + } + if (input.workerIds?.length) { + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['worker_my_orders'], + workerIds: input.workerIds, + }) + } +} + +/** 资金申请变更会同步后台资金页与申请所属打手的钱包页面。 */ +export function publishWorkerFinanceRealtimeChange(input: { + requestId: number + workerId: number + walletChanged?: boolean +}) { + publishRealtimeEvent({ + type: 'worker_finance_request.changed', + entityId: input.requestId, + scopes: ['admin_worker_finance'], + admin: true, + }) + publishRealtimeEvent({ + type: 'worker_finance_request.changed', + entityId: input.requestId, + scopes: ['worker_finance_requests'], + workerIds: [input.workerId], + }) + if (input.walletChanged) { + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: input.workerId, + scopes: ['worker_wallet'], + workerIds: [input.workerId], + }) + } +} + +function openRealtimeStream(response: Response, workerId: number | null) { + const subscriber: RealtimeSubscriber = { + id: nextSubscriberId++, + response, + workerId, + } + const subscribers = workerId === null ? adminSubscribers : workerSubscribers + + response.status(200) + response.setHeader('Content-Type', 'text/event-stream; charset=utf-8') + response.setHeader('Cache-Control', 'no-cache, no-transform') + response.setHeader('Connection', 'keep-alive') + response.setHeader('X-Accel-Buffering', 'no') + response.flushHeaders() + response.write('retry: 3000\n\n') + subscribers.set(subscriber.id, subscriber) + writeEvent(subscriber, 'ready', { + connectedAt: new Date().toISOString(), + }) + ensureHeartbeat() + + response.on('close', () => { + subscribers.delete(subscriber.id) + stopHeartbeatWhenIdle() + }) +} + +function broadcast(subscribers: Iterable, event: RealtimeEvent) { + for (const subscriber of subscribers) { + writeEvent(subscriber, 'realtime', event) + } +} + +function writeEvent(subscriber: RealtimeSubscriber, name: string, payload: unknown) { + if (subscriber.response.writableEnded || subscriber.response.destroyed) { + removeSubscriber(subscriber) + return + } + + try { + subscriber.response.write(`event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`) + } catch { + removeSubscriber(subscriber) + } +} + +function removeSubscriber(subscriber: RealtimeSubscriber) { + const subscribers = subscriber.workerId === null ? adminSubscribers : workerSubscribers + subscribers.delete(subscriber.id) + stopHeartbeatWhenIdle() +} + +function ensureHeartbeat() { + if (heartbeatTimer) return + heartbeatTimer = setInterval(() => { + for (const subscriber of [...adminSubscribers.values(), ...workerSubscribers.values()]) { + if (subscriber.response.writableEnded || subscriber.response.destroyed) { + removeSubscriber(subscriber) + continue + } + try { + subscriber.response.write(': ping\n\n') + } catch { + removeSubscriber(subscriber) + } + } + }, 25_000) + heartbeatTimer.unref() +} + +function stopHeartbeatWhenIdle() { + if (adminSubscribers.size > 0 || workerSubscribers.size > 0 || !heartbeatTimer) return + clearInterval(heartbeatTimer) + heartbeatTimer = null +} diff --git a/apps/backend/src/services/worker-platform/admin-service.ts b/apps/backend/src/services/worker-platform/admin-service.ts index 227bc848..c8e2d8a2 100644 --- a/apps/backend/src/services/worker-platform/admin-service.ts +++ b/apps/backend/src/services/worker-platform/admin-service.ts @@ -65,6 +65,11 @@ import { createWorkOrderMaterialAdminNotification, resolveAdminNotificationEntity, } from '../admin/admin-notification-service.js' +import { + publishRealtimeEvent, + publishWorkerFinanceRealtimeChange, + publishWorkOrderRealtimeChange, +} from '../realtime/realtime-event-service.js' import { consumeIndustryVouchersBeforeWorkOrderAssign, consumeIndustryVouchersBeforeWorkOrderPublish, @@ -512,6 +517,17 @@ export async function assignAdminWorkOrderToWorker( }, ) } + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: [Number(worker.id)], + hallChanged: workOrder.status === WORK_ORDER_STATUS.OPEN, + }) + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: Number(worker.id), + scopes: ['worker_wallet'], + workerIds: [Number(worker.id)], + }) return { order: mapWorkOrderAdmin(updated), voucherConsume } } @@ -540,6 +556,12 @@ export async function creditAdminWorkerWallet(workerId: number | string, payload payloadJson: JSON.stringify({ source: 'admin_manual_credit' }), now: nowIso(), }) + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: Number(workerId), + scopes: ['worker_wallet'], + workerIds: [Number(workerId)], + }) return { wallet: mapWallet(wallet) } } @@ -693,6 +715,11 @@ export async function reviewAdminWorkerFinanceRequest( Number(reviewed.request.id), 'finance_request_reviewed', ) + publishWorkerFinanceRealtimeChange({ + requestId: Number(reviewed.request.id), + workerId: Number(reviewed.request.worker_id), + walletChanged: reviewed.request.request_type === 'withdraw' && status === 'approved', + }) return { request: mapFinanceRequest(reviewed.request), @@ -817,6 +844,10 @@ export async function deductAdminWorkOrderPendingDeposit( payloadJson: JSON.stringify({ amount, note }), now: nowIso(), }) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: await resolveWorkOrderRealtimeWorkerIds(workOrder), + }) return { deductedAmount: result.deductedAmount } } @@ -891,6 +922,10 @@ export async function updateAdminWorkOrderSharing( }), now, }) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + hallChanged: workOrder.status === WORK_ORDER_STATUS.OPEN, + }) return { order: mapWorkOrderAdmin(updated || workOrder) } } @@ -983,6 +1018,10 @@ export async function submitAdminWorkOrderMaterial( if (nextStatus !== WORK_ORDER_STATUS.PENDING_MATERIAL) { await resolveAdminNotificationEntity('work_order', Number(workOrder.id), 'material_completed') } + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: await resolveWorkOrderRealtimeWorkerIds(workOrder), + }) return { order: mapWorkOrderAdmin(updated || workOrder), complete } } @@ -1152,6 +1191,10 @@ export async function updateAdminWorkOrder( }), now, }) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + hallChanged: workOrder.status === WORK_ORDER_STATUS.OPEN, + }) return { order: mapWorkOrderAdmin(updated || workOrder) } } @@ -1190,6 +1233,10 @@ export async function deleteAdminWorkOrder(workOrderId: number | string) { errorCode: 'work_order_delete_failed', }) } + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + hallChanged: workOrder.status === WORK_ORDER_STATUS.OPEN, + }) return { deleted: true } } @@ -1256,6 +1303,7 @@ export async function publishAdminWorkOrder(workOrderId: number | string, actorN payloadJson: JSON.stringify({ voucherConsume }), now, }) + publishWorkOrderRealtimeChange({ workOrderId: Number(workOrder.id), hallChanged: true }) return { order: mapWorkOrderAdmin(updated), voucherConsume } } @@ -1282,6 +1330,7 @@ export async function unpublishAdminWorkOrder(workOrderId: number | string, acto toStatus: WORK_ORDER_STATUS.UNASSIGNED, now, }) + publishWorkOrderRealtimeChange({ workOrderId: Number(workOrder.id), hallChanged: true }) return { order: mapWorkOrderAdmin(updated || workOrder) } } @@ -1312,6 +1361,7 @@ export async function pinAdminWorkOrder( toStatus: workOrder.status, now, }) + publishWorkOrderRealtimeChange({ workOrderId: Number(workOrder.id), hallChanged: true }) return { order: mapWorkOrderAdmin(updated || workOrder) } } @@ -1351,6 +1401,10 @@ export async function markAdminWorkOrderProblem( if (workOrder.status === WORK_ORDER_STATUS.PENDING_ACCEPTANCE) { await resolveAdminNotificationEntity('work_order', Number(workOrder.id), 'marked_problem') } + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: await resolveWorkOrderRealtimeWorkerIds(workOrder), + }) return { order: mapWorkOrderAdmin(updated || workOrder) } } @@ -1389,6 +1443,10 @@ export async function resolveAdminProblemWorkOrder( errorCode: 'work_order_problem_resolution_conflict', }) } + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: await resolveWorkOrderRealtimeWorkerIds(workOrder), + }) return { order: mapWorkOrderAdmin(updated) } } @@ -1421,6 +1479,19 @@ export async function acceptAdminWorkOrder(workOrderId: number | string, actorNa }) } await resolveAdminNotificationEntity('work_order', Number(workOrder.id), 'accepted') + const workerIds = await resolveWorkOrderRealtimeWorkerIds(workOrder) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds, + }) + for (const workerId of workerIds) { + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: workerId, + scopes: ['worker_wallet'], + workerIds: [workerId], + }) + } return { order: mapWorkOrderAdmin(updated || workOrder) } } @@ -1498,6 +1569,16 @@ export async function unassignAdminWorkOrder(workOrderId: number | string, actor errorCode: 'work_order_unassign_conflict', }) } + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: [Number(workOrder.assigned_worker_id)], + }) + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: Number(workOrder.assigned_worker_id), + scopes: ['worker_wallet'], + workerIds: [Number(workOrder.assigned_worker_id)], + }) return { order: mapWorkOrderAdmin(updated) } } @@ -1555,9 +1636,30 @@ export async function cancelAdminWorkOrder( errorCode: 'work_order_cancel_conflict', }) } + const workerId = Number(workOrder.assigned_worker_id) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: workerId ? [workerId] : [], + hallChanged: action === 'return_to_hall', + }) + if (workerId) { + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: workerId, + scopes: ['worker_wallet'], + workerIds: [workerId], + }) + } return { order: mapWorkOrderAdmin(updated), action } } +async function resolveWorkOrderRealtimeWorkerIds(workOrder: WorkOrderRow) { + const shares = await listWorkOrderShares(workOrder.id) + return [Number(workOrder.assigned_worker_id), ...shares.map((share) => Number(share.worker_id))] + .filter((workerId) => Number.isInteger(workerId) && workerId > 0) + .filter((workerId, index, values) => values.indexOf(workerId) === index) +} + export async function getAdminWorkerPlatformSummary() { const [pendingWorkers, pendingMaterial, openOrders, inProgressOrders] = await Promise.all([ listWorkerUsers({ page: 1, pageSize: 1, status: 'pending_review' }), diff --git a/apps/backend/src/services/worker-platform/sync-work-orders-from-send-code.ts b/apps/backend/src/services/worker-platform/sync-work-orders-from-send-code.ts index 28b15df9..3df90339 100644 --- a/apps/backend/src/services/worker-platform/sync-work-orders-from-send-code.ts +++ b/apps/backend/src/services/worker-platform/sync-work-orders-from-send-code.ts @@ -129,6 +129,9 @@ export async function syncWorkOrdersFromKuaishouSendCode( }, }) createdWorkOrderCount = syncResult.createdCount + for (const workOrder of syncResult.created) { + workOrderIds.push(Number(workOrder.workOrderId)) + } for (const item of syncResult.skipped) { skipped.push(`work_order:${item.reason}`) } diff --git a/apps/backend/src/services/worker-platform/worker-service.ts b/apps/backend/src/services/worker-platform/worker-service.ts index aba1aeee..6aa7dc10 100644 --- a/apps/backend/src/services/worker-platform/worker-service.ts +++ b/apps/backend/src/services/worker-platform/worker-service.ts @@ -42,6 +42,7 @@ import { listWorkCategories, listWorkOrders, listWorkOrderEventsByOrderId, + listWorkOrderShares, listWorkerLevels, listDueDepositUnfreezes, listOverdueWorkOrders, @@ -75,6 +76,11 @@ import { getWorkerAcceptanceReminder, remindWorkerAcceptanceAdminNotification, } from '../admin/admin-notification-service.js' +import { + publishRealtimeEvent, + publishWorkerFinanceRealtimeChange, + publishWorkOrderRealtimeChange, +} from '../realtime/realtime-event-service.js' import { DEFAULT_CATEGORY_KEY, @@ -701,6 +707,10 @@ export async function createWorkerRechargeRequest( workerName: worker.display_name || worker.username, amountFen: Number(created.amount || amount), }) + publishWorkerFinanceRealtimeChange({ + requestId: Number(created.id), + workerId: Number(worker.id), + }) return { request: mapFinanceRequest(created), } @@ -797,6 +807,10 @@ export async function createWorkerWithdrawRequest( amountFen: Number(created.amount || amount), channel: created.account_channel, }) + publishWorkerFinanceRealtimeChange({ + requestId: Number(created.id), + workerId: Number(worker.id), + }) return { request: mapFinanceRequest(created), availableForWithdraw, @@ -1049,6 +1063,17 @@ export async function grabWorkerHallOrder(workOrderId: number | string, session: }), now: nowIso(), }) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: [Number(worker.id)], + hallChanged: true, + }) + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: Number(worker.id), + scopes: ['worker_wallet'], + workerIds: [Number(worker.id)], + }) return { order: mapWorkOrderForWorker(grabbed.order, permissions) } } @@ -1084,7 +1109,15 @@ export async function settleDueDepositUnfreezes( unfreezeId: unfreeze.id, now: nowIso(), }) - if (released) processedCount += 1 + if (released) { + processedCount += 1 + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: Number(released.worker_id), + scopes: ['worker_wallet'], + workerIds: [Number(released.worker_id)], + }) + } } return { checkedCount: due.length, @@ -1117,6 +1150,21 @@ export async function settleOverdueWorkOrders(options: { limit?: number; workerI }) if (result.order) { processedCount += 1 + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: Number(workOrder.assigned_worker_id) + ? [Number(workOrder.assigned_worker_id)] + : [], + hallChanged: result.order.status === WORK_ORDER_STATUS.OPEN, + }) + if (Number(workOrder.assigned_worker_id)) { + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: Number(workOrder.assigned_worker_id), + scopes: ['worker_wallet'], + workerIds: [Number(workOrder.assigned_worker_id)], + }) + } } else { skippedCount += 1 } @@ -1288,6 +1336,11 @@ export async function joinWorkerSharingOrder( if (!result.share) { throw resolveWorkOrderShareJoinError(result.failureReason) } + await publishWorkOrderParticipantChange(Number(workOrderId), { + workerIds: [Number(worker.id)], + hallChanged: true, + walletChangedWorkerIds: [Number(worker.id)], + }) return { share: mapWorkOrderShare(result.share) } } @@ -1440,6 +1493,16 @@ export async function submitWorkerOrderAcceptance( }), now, }) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: [Number(worker.id)], + }) + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: Number(worker.id), + scopes: ['worker_wallet'], + workerIds: [Number(worker.id)], + }) return { order: mapWorkOrderForWorker(acceptedOrder, resolveWorkerPermissions(worker)), } @@ -1451,6 +1514,10 @@ export async function submitWorkerOrderAcceptance( workerName: worker.display_name || worker.username, imageCount: imageUrls.length, }) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + workerIds: [Number(worker.id)], + }) return { order: mapWorkOrderForWorker(updated || workOrder, resolveWorkerPermissions(worker)), } @@ -1645,6 +1712,7 @@ export async function supplementWorkerAcceptedOrderEvidence( }), now, }) + publishWorkOrderRealtimeChange({ workOrderId: Number(workOrder.id) }) return { order: mapWorkOrderForWorker(updated || workOrder, resolveWorkerPermissions(worker)), } @@ -1697,6 +1765,9 @@ async function submitWorkerSharingAcceptance( errorCode: 'work_order_sharing_share_not_found', }) } + await publishWorkOrderParticipantChange(Number(workOrder.id), { + workerIds: [Number(session.workerId)], + }) return { share: mapWorkOrderShare(result.share), sharingPendingCount: await countWorkOrderPendingSharingSubmissions(Number(workOrder.id)), @@ -1789,6 +1860,9 @@ export async function saveWorkerOrderAcceptanceDraft( payloadJson: JSON.stringify({ imageCount: imageUrls.length }), now, }) + await publishWorkOrderParticipantChange(Number(workOrder.id), { + workerIds: [Number(session.workerId)], + }) return { order: mapWorkOrderForWorker(updated || workOrder, resolveWorkerPermissions(worker)), } @@ -1823,6 +1897,9 @@ export async function saveWorkerOrderAcceptanceDraft( payloadJson: JSON.stringify({ imageCount: imageUrls.length }), now, }) + await publishWorkOrderParticipantChange(Number(workOrder.id), { + workerIds: [Number(session.workerId)], + }) return { order: mapWorkOrderForWorker(updated || workOrder, resolveWorkerPermissions(worker)), } @@ -1881,6 +1958,9 @@ async function saveWorkerSharingAcceptanceDraft( }), now, }) + await publishWorkOrderParticipantChange(Number(workOrder.id), { + workerIds: [Number(session.workerId)], + }) return { share: share ? mapWorkOrderShare(share) : mapWorkOrderShare(sharingShare) } } @@ -1915,6 +1995,9 @@ async function saveWorkerSharingAcceptanceDraft( }), now, }) + await publishWorkOrderParticipantChange(Number(workOrder.id), { + workerIds: [Number(session.workerId)], + }) return { share: share ? mapWorkOrderShare(share) : mapWorkOrderShare(sharingShare) } } @@ -2001,12 +2084,47 @@ export async function collectSubmitWorkOrder(payload: JsonObject = {}) { payloadJson: JSON.stringify({ fields: submittedFields, complete }), now, }) + publishWorkOrderRealtimeChange({ + workOrderId: Number(workOrder.id), + }) return { order: mapWorkOrderPublic(updated || workOrder), complete, } } +/** 拼单工单同步所有参与打手,避免仅主打手页面更新。 */ +async function publishWorkOrderParticipantChange( + workOrderId: number, + options: { + workerIds?: number[] + hallChanged?: boolean + walletChangedWorkerIds?: number[] + } = {}, +) { + const shares = await listWorkOrderShares(workOrderId) + const workerIds = [ + ...(options.workerIds || []), + ...shares.map((share) => Number(share.worker_id)), + ].filter( + (workerId, index, values) => + Number.isInteger(workerId) && workerId > 0 && values.indexOf(workerId) === index, + ) + publishWorkOrderRealtimeChange({ + workOrderId, + workerIds, + ...(options.hallChanged ? { hallChanged: true } : {}), + }) + for (const workerId of options.walletChangedWorkerIds || []) { + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: workerId, + scopes: ['worker_wallet'], + workerIds: [workerId], + }) + } +} + export function resolveCollectSubmitTargetWorkOrder( orderNo: string, workOrders: WorkOrderRow[], diff --git a/apps/frontend/src/hooks/useRealtimeEvents.ts b/apps/frontend/src/hooks/useRealtimeEvents.ts new file mode 100644 index 00000000..a692e263 --- /dev/null +++ b/apps/frontend/src/hooks/useRealtimeEvents.ts @@ -0,0 +1,166 @@ +import { useEffect, useRef } from 'react' + +export type RealtimeEvent = { + eventId: number + type: + | 'admin_notification.changed' + | 'order.changed' + | 'work_order.changed' + | 'worker_finance_request.changed' + | 'worker_wallet.changed' + entityId: number + scopes: string[] + occurredAt: string +} + +type UseRealtimeEventsOptions = { + url: string + token: string + onEvent: (event: RealtimeEvent) => void + onConnected?: () => void +} + +const RECONNECT_DELAY_MS = 3_000 + +/** 通过 fetch 保持带 Bearer 鉴权的 SSE 连接,避免把登录 token 暴露到 URL。 */ +export function useRealtimeEvents({ url, token, onEvent, onConnected }: UseRealtimeEventsOptions) { + const onEventRef = useRef(onEvent) + const onConnectedRef = useRef(onConnected) + onEventRef.current = onEvent + onConnectedRef.current = onConnected + + useEffect(() => { + if (!token) return + + let stopped = false + let abortController: AbortController | null = null + let reconnectTimer = 0 + + const connect = () => { + if (stopped || document.visibilityState !== 'visible') return + abortController?.abort() + abortController = new AbortController() + + void fetch(url, { + method: 'GET', + headers: { + Accept: 'text/event-stream', + Authorization: `Bearer ${token}`, + }, + cache: 'no-store', + signal: abortController.signal, + }) + .then(async (response) => { + if (!response.ok || !response.body) { + throw new Error(`realtime_stream_${response.status}`) + } + onConnectedRef.current?.() + await readSseStream(response.body, (eventName, data) => { + if (eventName !== 'realtime') return + const event = normalizeRealtimeEvent(data) + if (event) onEventRef.current(event) + }) + }) + .catch(() => undefined) + .finally(() => { + if (stopped || document.visibilityState !== 'visible') return + reconnectTimer = window.setTimeout(connect, RECONNECT_DELAY_MS) + }) + } + + const handleVisibilityChange = () => { + if (document.visibilityState === 'visible') { + connect() + } else { + abortController?.abort() + window.clearTimeout(reconnectTimer) + } + } + + document.addEventListener('visibilitychange', handleVisibilityChange) + connect() + return () => { + stopped = true + document.removeEventListener('visibilitychange', handleVisibilityChange) + abortController?.abort() + window.clearTimeout(reconnectTimer) + } + }, [token, url]) +} + +async function readSseStream( + stream: ReadableStream, + onMessage: (eventName: string, data: string) => void, +) { + const reader = stream.getReader() + const decoder = new TextDecoder() + let buffer = '' + let eventName = 'message' + let dataLines: string[] = [] + + function dispatch() { + if (dataLines.length > 0) { + onMessage(eventName, dataLines.join('\n')) + } + eventName = 'message' + dataLines = [] + } + + function parseLine(line: string) { + if (!line) { + dispatch() + return + } + if (line.startsWith(':')) return + const separatorIndex = line.indexOf(':') + const field = separatorIndex >= 0 ? line.slice(0, separatorIndex) : line + const value = separatorIndex >= 0 ? line.slice(separatorIndex + 1).replace(/^ /, '') : '' + if (field === 'event') eventName = value || 'message' + if (field === 'data') dataLines.push(value) + } + + try { + while (true) { + const { done, value } = await reader.read() + if (done) break + buffer += decoder.decode(value, { stream: true }) + const lines = buffer.split(/\r?\n/) + buffer = lines.pop() || '' + lines.forEach(parseLine) + } + buffer += decoder.decode() + if (buffer) parseLine(buffer) + dispatch() + } finally { + reader.releaseLock() + } +} + +function normalizeRealtimeEvent(value: string): RealtimeEvent | null { + try { + const parsed = JSON.parse(value) as Partial + const type = String(parsed.type || '') + if ( + ![ + 'admin_notification.changed', + 'order.changed', + 'work_order.changed', + 'worker_finance_request.changed', + 'worker_wallet.changed', + ].includes(type) + ) { + return null + } + return { + eventId: Number(parsed.eventId || 0), + type: type as RealtimeEvent['type'], + entityId: Number(parsed.entityId || 0), + scopes: Array.isArray(parsed.scopes) + ? parsed.scopes.map((item) => String(item || '')).filter(Boolean) + : [], + occurredAt: String(parsed.occurredAt || ''), + } + } catch { + return null + } +} diff --git a/apps/frontend/src/layouts/AdminLayout.tsx b/apps/frontend/src/layouts/AdminLayout.tsx index b6adacd5..b6ac9826 100644 --- a/apps/frontend/src/layouts/AdminLayout.tsx +++ b/apps/frontend/src/layouts/AdminLayout.tsx @@ -44,10 +44,12 @@ import { import { clearAdminSession, getAdminRole, + getAdminToken, getAdminTokenExpiresAt, getAdminUsername, } from '@/utils/admin-auth' import { formatAdminDateTime } from '@/utils/admin-time' +import { useRealtimeEvents, type RealtimeEvent } from '@/hooks/useRealtimeEvents' const SIDEBAR_COLLAPSED_KEY = 'react-admin-sidebar-collapsed' const NOTIFICATION_SOUND_KEY = 'admin-notification-sound-enabled' @@ -88,6 +90,15 @@ export default function AdminLayout() { const notifications = notificationsQuery.data?.data.items || [] const pendingNotificationCount = Number(notificationsQuery.data?.data.pendingCount || 0) + useRealtimeEvents({ + url: '/api/v1/admin/realtime', + token: getAdminToken(), + onEvent: (event) => handleAdminRealtimeEvent(event, queryClient), + onConnected: () => { + void invalidateAdminRealtimeQueries(queryClient) + }, + }) + useEffect(() => { const currentVersions = new Map( notifications.map((notification) => [notification.notificationId, notification.updatedAt]), @@ -343,6 +354,40 @@ export default function AdminLayout() { ) } +function handleAdminRealtimeEvent( + event: RealtimeEvent, + queryClient: ReturnType, +) { + if (event.type === 'admin_notification.changed') { + void queryClient.invalidateQueries({ queryKey: ['admin-notifications'] }) + } + if (event.type === 'order.changed') { + void queryClient.invalidateQueries({ queryKey: ['admin-orders'] }) + void queryClient.invalidateQueries({ queryKey: ['admin-order-detail', String(event.entityId)] }) + void queryClient.invalidateQueries({ queryKey: ['admin-tasks'] }) + void queryClient.invalidateQueries({ queryKey: ['admin-dashboard-summary'] }) + } + if (event.type === 'work_order.changed') { + void queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-orders'] }) + void queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-summary'] }) + } + if (event.type === 'worker_finance_request.changed') { + void queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-finance-requests'] }) + } +} + +function invalidateAdminRealtimeQueries(queryClient: ReturnType) { + return Promise.all([ + queryClient.invalidateQueries({ queryKey: ['admin-notifications'] }), + queryClient.invalidateQueries({ queryKey: ['admin-orders'] }), + queryClient.invalidateQueries({ queryKey: ['admin-tasks'] }), + queryClient.invalidateQueries({ queryKey: ['admin-dashboard-summary'] }), + queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-orders'] }), + queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-finance-requests'] }), + queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-summary'] }), + ]) +} + type NotificationAlertKind = | 'withdraw' | 'recharge' diff --git a/apps/frontend/src/layouts/WorkerLayout.tsx b/apps/frontend/src/layouts/WorkerLayout.tsx index 9c417635..f54a1fa7 100644 --- a/apps/frontend/src/layouts/WorkerLayout.tsx +++ b/apps/frontend/src/layouts/WorkerLayout.tsx @@ -5,12 +5,15 @@ import { UserOutlined, } from '@ant-design/icons' import { App, Button, Layout, Menu, Typography } from 'antd' +import { useQueryClient } from '@tanstack/react-query' import type { MenuProps } from 'antd' import { Outlet, useLocation, useNavigate } from 'react-router' import { logoutWorker } from '@/services/worker' import { useIsMobile } from '@/utils/use-is-mobile' import { clearWorkerSession, getWorkerStatus, getWorkerUsername } from '@/utils/worker-auth' +import { getWorkerToken } from '@/utils/worker-auth' +import { useRealtimeEvents, type RealtimeEvent } from '@/hooks/useRealtimeEvents' const menuItems: MenuProps['items'] = [ { key: '/hall', icon: , label: '抢单大厅' }, @@ -31,8 +34,18 @@ export default function WorkerLayout() { const { message, modal } = App.useApp() const username = getWorkerUsername() || '打手' const status = getWorkerStatus() + const queryClient = useQueryClient() const currentKey = resolveSelectedKey(location.pathname) + useRealtimeEvents({ + url: '/api/v1/worker/realtime', + token: getWorkerToken(), + onEvent: (event) => handleWorkerRealtimeEvent(event, queryClient), + onConnected: () => { + void invalidateWorkerRealtimeQueries(queryClient) + }, + }) + async function submitLogout() { try { await logoutWorker() @@ -143,6 +156,42 @@ export default function WorkerLayout() { ) } +function handleWorkerRealtimeEvent( + event: RealtimeEvent, + queryClient: ReturnType, +) { + if (event.type === 'work_order.changed') { + if (event.scopes.includes('worker_hall')) { + void queryClient.invalidateQueries({ queryKey: ['worker-hall-orders'] }) + void queryClient.invalidateQueries({ queryKey: ['worker-hall-summary'] }) + } + if (event.scopes.includes('worker_my_orders')) { + void queryClient.invalidateQueries({ queryKey: ['worker-my-orders'] }) + void queryClient.invalidateQueries({ queryKey: ['worker-my-orders-overview'] }) + } + } + if (event.type === 'worker_finance_request.changed') { + void queryClient.invalidateQueries({ queryKey: ['worker-finance-requests'] }) + void queryClient.invalidateQueries({ queryKey: ['worker-profile'] }) + } + if (event.type === 'worker_wallet.changed') { + void queryClient.invalidateQueries({ queryKey: ['worker-wallet-ledgers'] }) + void queryClient.invalidateQueries({ queryKey: ['worker-profile'] }) + } +} + +function invalidateWorkerRealtimeQueries(queryClient: ReturnType) { + return Promise.all([ + queryClient.invalidateQueries({ queryKey: ['worker-hall-orders'] }), + queryClient.invalidateQueries({ queryKey: ['worker-hall-summary'] }), + queryClient.invalidateQueries({ queryKey: ['worker-my-orders'] }), + queryClient.invalidateQueries({ queryKey: ['worker-my-orders-overview'] }), + queryClient.invalidateQueries({ queryKey: ['worker-profile'] }), + queryClient.invalidateQueries({ queryKey: ['worker-wallet-ledgers'] }), + queryClient.invalidateQueries({ queryKey: ['worker-finance-requests'] }), + ]) +} + function resolveSelectedKey(pathname: string) { if (pathname.startsWith('/orders')) return '/orders' if (pathname.startsWith('/profile')) return '/profile'