From 2aa8769661915728799901cb169179787fff30f1 Mon Sep 17 00:00:00 2001 From: yml2213 Date: Mon, 17 Aug 2026 17:32:02 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=B9=E9=80=A0=E5=AE=9E=E6=97=B6=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6=E7=BC=93=E5=AD=98=E5=90=8C=E6=AD=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/db/migrations/030_realtime_events.sql | 19 + .../src/repositories/realtime-event-repo.ts | 71 ++++ apps/backend/src/routes/admin/realtime.ts | 11 +- apps/backend/src/routes/worker.ts | 6 +- .../admin/admin-notification-service.ts | 8 +- .../src/services/order/order-service.ts | 9 +- .../kuaishou-industry/send-code-service.ts | 9 +- .../realtime/realtime-event-service.ts | 367 +++++++++++++++--- .../services/worker-platform/admin-service.ts | 37 +- .../worker-platform/worker-service.ts | 37 +- apps/frontend/src/hooks/useRealtimeEvents.ts | 26 +- apps/frontend/src/layouts/AdminLayout.tsx | 320 +++++++++------ apps/frontend/src/layouts/WorkerLayout.tsx | 194 +++++++-- apps/frontend/src/main.tsx | 3 + .../src/pages/worker/WorkerHallPage.tsx | 38 -- apps/frontend/src/styles/worker.css | 33 -- 16 files changed, 837 insertions(+), 351 deletions(-) create mode 100644 apps/backend/src/db/migrations/030_realtime_events.sql create mode 100644 apps/backend/src/repositories/realtime-event-repo.ts diff --git a/apps/backend/src/db/migrations/030_realtime_events.sql b/apps/backend/src/db/migrations/030_realtime_events.sql new file mode 100644 index 00000000..20566903 --- /dev/null +++ b/apps/backend/src/db/migrations/030_realtime_events.sql @@ -0,0 +1,19 @@ +-- 030_realtime_events.sql —— SSE 事件持久化与断线补发。 + +CREATE TABLE IF NOT EXISTS realtime_events ( + id BIGSERIAL PRIMARY KEY, + event_type TEXT NOT NULL, + entity_id BIGINT NOT NULL DEFAULT 0, + scopes JSONB NOT NULL DEFAULT '[]'::jsonb, + operation TEXT NOT NULL DEFAULT 'refresh', + data JSONB, + admin_audience BOOLEAN NOT NULL DEFAULT FALSE, + worker_ids JSONB NOT NULL DEFAULT '[]'::jsonb, + broadcast_workers BOOLEAN NOT NULL DEFAULT FALSE, + occurred_at TIMESTAMPTZ NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_realtime_events_occurred_at + ON realtime_events(occurred_at DESC); + +COMMENT ON TABLE realtime_events IS 'SSE 实时事件日志,保留最近事件以支持断线后的增量补发'; diff --git a/apps/backend/src/repositories/realtime-event-repo.ts b/apps/backend/src/repositories/realtime-event-repo.ts new file mode 100644 index 00000000..c2c75ac3 --- /dev/null +++ b/apps/backend/src/repositories/realtime-event-repo.ts @@ -0,0 +1,71 @@ +import { query } from '../db/client.js' + +export type RealtimeEventRow = { + id: number + event_type: string + entity_id: number + scopes: unknown + operation: string + data: unknown + admin_audience: boolean + worker_ids: unknown + broadcast_workers: boolean + occurred_at: string +} + +export async function createRealtimeEvent(input: { + eventType: string + entityId: number + scopes: string[] + operation: string + data?: unknown + adminAudience: boolean + workerIds: number[] + broadcastWorkers: boolean + occurredAt: string +}): Promise { + const result = await query( + ` + INSERT INTO realtime_events ( + event_type, entity_id, scopes, operation, data, + admin_audience, worker_ids, broadcast_workers, occurred_at + ) VALUES ($1, $2, $3::jsonb, $4, $5::jsonb, $6, $7::jsonb, $8, $9) + RETURNING * + `, + [ + input.eventType, + input.entityId, + JSON.stringify(input.scopes), + input.operation, + input.data === undefined ? null : JSON.stringify(input.data), + input.adminAudience, + JSON.stringify(input.workerIds), + input.broadcastWorkers, + input.occurredAt, + ], + ) + const row = result.rows[0] + if (!row) throw new Error('realtime_event_create_failed') + + // 事件仅用于短期断线补发,控制表体积和重连回放成本。 + await query( + ` + DELETE FROM realtime_events + WHERE id < GREATEST(0, (SELECT COALESCE(MAX(id), 0) - 500 FROM realtime_events)) + `, + ) + return row +} + +export async function listRealtimeEventsAfter(eventId: number, limit = 500) { + const result = await query( + ` + SELECT * FROM realtime_events + WHERE id > $1 + ORDER BY id ASC + LIMIT $2 + `, + [Math.max(0, Number(eventId) || 0), Math.min(500, Math.max(1, Number(limit) || 500))], + ) + return result.rows +} diff --git a/apps/backend/src/routes/admin/realtime.ts b/apps/backend/src/routes/admin/realtime.ts index 26e2f5a9..4d579505 100644 --- a/apps/backend/src/routes/admin/realtime.ts +++ b/apps/backend/src/routes/admin/realtime.ts @@ -5,8 +5,15 @@ import { requireAdminRoles } from './session.js' const router = Router() -router.get('/realtime', requireAdminRoles(['admin', 'operator', 'support']), (_req, res) => { - openAdminRealtimeStream(res) +router.get('/realtime', requireAdminRoles(['admin', 'operator', 'support']), (req, res) => { + openAdminRealtimeStream(res, resolveLastEventId(req)) }) +function resolveLastEventId(request: { + get(name: string): string | undefined + query: Record +}) { + return Number(request.get('Last-Event-ID') || request.query.lastEventId || 0) || 0 +} + export default router diff --git a/apps/backend/src/routes/worker.ts b/apps/backend/src/routes/worker.ts index 3d9f0a05..fb4cc13e 100644 --- a/apps/backend/src/routes/worker.ts +++ b/apps/backend/src/routes/worker.ts @@ -180,7 +180,11 @@ router.delete( ) router.get('/realtime', (req, res) => { - openWorkerRealtimeStream(res, getRequiredWorkerSession(req).workerId) + openWorkerRealtimeStream( + res, + getRequiredWorkerSession(req).workerId, + Number(req.get('Last-Event-ID') || req.query.lastEventId || 0) || 0, + ) }) router.get( diff --git a/apps/backend/src/services/admin/admin-notification-service.ts b/apps/backend/src/services/admin/admin-notification-service.ts index 7b3dc3f3..b55b3199 100644 --- a/apps/backend/src/services/admin/admin-notification-service.ts +++ b/apps/backend/src/services/admin/admin-notification-service.ts @@ -122,6 +122,8 @@ export async function remindWorkerAcceptanceAdminNotification(workOrderId: numbe type: 'admin_notification.changed', entityId: Number(notification.id), scopes: ['admin_notifications'], + operation: 'upsert', + data: mapAdminNotification(notification), admin: true, }) } @@ -147,6 +149,8 @@ export async function resolveAdminNotificationEntity( type: 'admin_notification.changed', entityId, scopes: ['admin_notifications'], + operation: 'remove', + data: { notificationId: entityId }, admin: true, }) } @@ -163,7 +167,7 @@ export async function listAdminPendingNotifications(limit = 20) { } } -function mapAdminNotification(row: AdminNotificationRow) { +export function mapAdminNotification(row: AdminNotificationRow) { return { notificationId: Number(row.id), notificationType: row.notification_type, @@ -215,6 +219,8 @@ async function createConfiguredWorkerPlatformNotification( type: 'admin_notification.changed', entityId: Number(notification.id), scopes: ['admin_notifications'], + operation: 'upsert', + data: mapAdminNotification(notification), admin: true, }) } diff --git a/apps/backend/src/services/order/order-service.ts b/apps/backend/src/services/order/order-service.ts index d1be405e..41bfb229 100644 --- a/apps/backend/src/services/order/order-service.ts +++ b/apps/backend/src/services/order/order-service.ts @@ -12,7 +12,7 @@ import { } from '../fulfillment/order-fulfillment-readiness-service.js' import { syncWorkerOrdersForSourceOrder } from '../worker-platform/index.js' import { - publishRealtimeEvent, + publishAdminOrderRealtimeChange, publishWorkOrderRealtimeChange, } from '../realtime/realtime-event-service.js' import { nowIso } from '../../utils/time.js' @@ -246,12 +246,7 @@ function publishSourceOrderRealtimeChanges( orderId: number, workOrders: Array<{ workOrderId: number }> = [], ) { - publishRealtimeEvent({ - type: 'order.changed', - entityId: Number(orderId), - scopes: ['admin_orders'], - admin: true, - }) + publishAdminOrderRealtimeChange(Number(orderId)) for (const workOrder of workOrders) { publishWorkOrderRealtimeChange({ workOrderId: workOrder.workOrderId }) } 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 43eef94c..b88b2fff 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 @@ -32,7 +32,7 @@ import { import { bindKuaishouIndustryVouchersToOrderTasks } from './voucher-binding-service.js' import { parseAmountToFen } from '../../../utils/money.js' import { - publishRealtimeEvent, + publishAdminOrderRealtimeChange, publishWorkOrderRealtimeChange, } from '../../realtime/realtime-event-service.js' import type { KuaishouIndustryVoucherRow, OrderRow } from '../../../types/repository/rows.js' @@ -136,12 +136,7 @@ export async function handleSendCode(rawBody: JsonObject = {}) { now, }) if (syncResult.orderId) { - publishRealtimeEvent({ - type: 'order.changed', - entityId: syncResult.orderId, - scopes: ['admin_orders'], - admin: true, - }) + publishAdminOrderRealtimeChange(syncResult.orderId) } for (const workOrderId of syncResult.workOrderIds) { publishWorkOrderRealtimeChange({ workOrderId }) diff --git a/apps/backend/src/services/realtime/realtime-event-service.ts b/apps/backend/src/services/realtime/realtime-event-service.ts index d6aaf917..070d0b02 100644 --- a/apps/backend/src/services/realtime/realtime-event-service.ts +++ b/apps/backend/src/services/realtime/realtime-event-service.ts @@ -1,5 +1,31 @@ import type { Response } from 'express' +import { + getWorkerFinanceRequestById, + getWorkerUserById, + getWorkOrderById, + listWorkerWorkOrderNotes, + listWorkOrderShares, + sumPendingUnfreezeByOrderIds, +} from '../../repositories/worker-platform/index.js' +import { getOrderById } from '../../repositories/order-repo.js' +import { listTasksByOrderId } from '../../repositories/task-repo.js' +import { + createRealtimeEvent, + listRealtimeEventsAfter, + type RealtimeEventRow, +} from '../../repositories/realtime-event-repo.js' +import { getAdminDashboardSummary } from '../admin/admin-dashboard-service.js' +import { mapAdminOrderListItem } from '../admin/admin-order-read-helpers.js' +import { mapAdminTaskListItem } from '../admin/admin-task-read-helpers.js' +import { + mapFinanceRequest, + mapWorkOrderAdmin, + mapWorkOrderForWorker, + mapWorkerUser, + resolveWorkerPermissions, +} from '../worker-platform/mappers.js' + export type RealtimeEventType = | 'admin_notification.changed' | 'order.changed' @@ -9,9 +35,12 @@ export type RealtimeEventType = export type RealtimeEvent = { eventId: number + version: number type: RealtimeEventType entityId: number scopes: string[] + operation: 'upsert' | 'remove' | 'refresh' + data?: unknown occurredAt: string } @@ -19,6 +48,8 @@ type RealtimeEventInput = { type: RealtimeEventType entityId?: number scopes?: string[] + operation?: RealtimeEvent['operation'] + data?: unknown admin?: boolean workerIds?: number[] broadcastWorkers?: boolean @@ -30,57 +61,92 @@ type RealtimeSubscriber = { workerId: number | null } +type StoredRealtimeEvent = RealtimeEvent & { + adminAudience: boolean + workerIds: number[] + broadcastWorkers: boolean +} + const adminSubscribers = new Map() const workerSubscribers = new Map() let nextSubscriberId = 1 -let nextEventId = 1 let heartbeatTimer: NodeJS.Timeout | null = null +let publishQueue: Promise = Promise.resolve() -/** 建立后台实时事件流,事件仅用于提示前端重新读取权威数据。 */ -export function openAdminRealtimeStream(response: Response) { - return openRealtimeStream(response, null) +/** 建立后台实时事件流,事件快照用于前端定点更新缓存。 */ +export function openAdminRealtimeStream(response: Response, lastEventId = 0) { + return openRealtimeStream(response, null, lastEventId) } /** 建立打手实时事件流,只向该打手推送其可见的订单和资金变化。 */ -export function openWorkerRealtimeStream(response: Response, workerId: number) { - return openRealtimeStream(response, Number(workerId)) +export function openWorkerRealtimeStream(response: Response, workerId: number, lastEventId = 0) { + return openRealtimeStream(response, Number(workerId), lastEventId) } -/** 发布轻量实时事件;数据库仍是唯一权威数据源。 */ +/** 发布带版本号和数据快照的实时事件;数据库仍是最终权威数据源。 */ export function publishRealtimeEvent(input: RealtimeEventInput) { - const event: RealtimeEvent = { - eventId: nextEventId++, + const workerIds = [ + ...new Set( + (input.workerIds || []) + .map((item) => Number(item)) + .filter((item) => Number.isInteger(item) && item > 0), + ), + ] + const event = { type: input.type, entityId: Math.max(0, Number(input.entityId) || 0), scopes: [ ...new Set((input.scopes || []).map((item) => String(item || '').trim()).filter(Boolean)), ], + operation: input.operation || 'refresh', + data: input.data, + adminAudience: input.admin === true, + workerIds, + broadcastWorkers: input.broadcastWorkers === true, occurredAt: new Date().toISOString(), } + // 按提交顺序串行写入,保证 Last-Event-ID 与推送顺序一致。 + publishQueue = publishQueue.then(() => persistAndBroadcast(event)).catch(() => undefined) +} - if (input.admin) { - broadcast(adminSubscribers.values(), event) +async function persistAndBroadcast(input: { + type: RealtimeEventType + entityId: number + scopes: string[] + operation: RealtimeEvent['operation'] + data?: unknown + adminAudience: boolean + workerIds: number[] + broadcastWorkers: boolean + occurredAt: string +}) { + const row = await createRealtimeEvent({ + eventType: input.type, + entityId: input.entityId, + scopes: input.scopes, + operation: input.operation, + data: input.data, + adminAudience: input.adminAudience, + workerIds: input.workerIds, + broadcastWorkers: input.broadcastWorkers, + occurredAt: input.occurredAt, + }) + const storedEvent = mapStoredEvent(row) + if (storedEvent.adminAudience) { + broadcast(adminSubscribers.values(), storedEvent) } - - 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 (storedEvent.broadcastWorkers) { + broadcast(workerSubscribers.values(), storedEvent) + return } - - if (workerIds.size > 0) { + if (storedEvent.workerIds.length > 0) { broadcast( [...workerSubscribers.values()].filter((subscriber) => - workerIds.has(Number(subscriber.workerId)), + storedEvent.workerIds.includes(Number(subscriber.workerId)), ), - event, + storedEvent, ) } - return event } /** 工单变更统一通知后台;大厅与打手本人按可见范围分别接收。 */ @@ -90,28 +156,15 @@ export function publishWorkOrderRealtimeChange(input: { hallChanged?: boolean }) { const workOrderId = Math.max(0, Number(input.workOrderId) || 0) - publishRealtimeEvent({ - type: 'work_order.changed', - entityId: workOrderId, - scopes: ['admin_work_orders'], - admin: true, + void publishWorkOrderSnapshots({ ...input, workOrderId }).catch(() => { + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['admin_work_orders'], + operation: 'refresh', + 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, - }) - } } /** 资金申请变更会同步后台资金页与申请所属打手的钱包页面。 */ @@ -120,29 +173,188 @@ export function publishWorkerFinanceRealtimeChange(input: { workerId: number walletChanged?: boolean }) { + void publishWorkerFinanceSnapshots(input).catch(() => { + publishRealtimeEvent({ + type: 'worker_finance_request.changed', + entityId: input.requestId, + scopes: ['admin_worker_finance'], + operation: 'refresh', + admin: true, + }) + }) +} + +/** 钱包事件携带当前钱包快照,避免客户端为余额变化重新请求个人资料。 */ +export function publishWorkerWalletRealtimeChange(workerId: number) { + const normalizedWorkerId = Math.max(0, Number(workerId) || 0) + void getWorkerUserById(normalizedWorkerId) + .then((worker) => { + publishRealtimeEvent({ + type: 'worker_wallet.changed', + entityId: normalizedWorkerId, + scopes: ['worker_wallet'], + operation: worker ? 'upsert' : 'refresh', + ...(worker ? { data: { wallet: mapWorkerUser(worker).wallet } } : {}), + workerIds: [normalizedWorkerId], + }) + }) + .catch(() => undefined) +} + +/** 订单变更推送后台列表行快照,前端无需重新读取整页订单。 */ +export function publishAdminOrderRealtimeChange(orderId: number) { + const normalizedOrderId = Math.max(0, Number(orderId) || 0) + void getOrderById(normalizedOrderId) + .then(async (order) => { + const [dashboard, tasks] = await Promise.all([ + getAdminDashboardSummary(), + order ? listTasksByOrderId(normalizedOrderId) : [], + ]) + const data = { + order: order ? await mapAdminOrderListItem(order) : { orderId: normalizedOrderId }, + dashboard, + tasks: tasks.map((task) => mapAdminTaskListItem(task)), + } + publishRealtimeEvent({ + type: 'order.changed', + entityId: normalizedOrderId, + scopes: ['admin_orders'], + operation: order ? 'upsert' : 'remove', + data, + admin: true, + }) + }) + .catch(() => undefined) +} + +async function publishWorkOrderSnapshots(input: { + workOrderId: number + workerIds?: number[] + hallChanged?: boolean +}) { + const workOrderId = Math.max(0, Number(input.workOrderId) || 0) + const workOrder = await getWorkOrderById(workOrderId) + if (!workOrder) { + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['admin_work_orders'], + operation: 'remove', + data: { workOrderId }, + admin: true, + }) + return + } + + const pendingUnfreezeByOrderId = await sumPendingUnfreezeByOrderIds([workOrderId]) + const shares = await listWorkOrderShares(workOrderId) + const joinedQuantity = shares + .filter((share) => share.status !== 'cancelled') + .reduce((sum, share) => sum + Number(share.quantity || 0), 0) + const pendingSubmissionCount = shares.filter((share) => share.status === 'joined').length + const adminSnapshot = { + ...mapWorkOrderAdmin(workOrder), + pendingUnfreezeAmount: pendingUnfreezeByOrderId.get(workOrderId) || 0, + sharingProgress: { joinedQuantity, pendingSubmissionCount }, + } + + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['admin_work_orders'], + operation: 'upsert', + data: adminSnapshot, + admin: true, + }) + + if (input.hallChanged) { + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['worker_hall'], + operation: workOrder.status === 'open' ? 'upsert' : 'remove', + data: mapWorkerHallWorkOrder(workOrder, shares), + broadcastWorkers: true, + }) + } + + for (const workerId of new Set( + (input.workerIds || []).map((item) => Number(item)).filter(Boolean), + )) { + const worker = await getWorkerUserById(workerId) + if (!worker) continue + const notes = await listWorkerWorkOrderNotes(workerId, [workOrderId]) + const note = notes.find((item) => Number(item.work_order_id) === workOrderId)?.note || '' + const share = shares.find((item) => Number(item.worker_id) === workerId) + publishRealtimeEvent({ + type: 'work_order.changed', + entityId: workOrderId, + scopes: ['worker_my_orders'], + operation: 'upsert', + data: mapWorkOrderForWorker(workOrder, resolveWorkerPermissions(worker), share, note), + workerIds: [workerId], + }) + } +} + +function mapWorkerHallWorkOrder( + workOrder: NonNullable>>, + shares: Awaited>, +) { + const mapped = mapWorkOrderAdmin(workOrder) + const { worker: _worker, orderId: _orderId, orderItemId: _orderItemId, ...visibleOrder } = mapped + const joinedQuantity = shares + .filter((share) => share.status !== 'cancelled') + .reduce((sum, share) => sum + Number(share.quantity || 0), 0) + const pendingSubmissionCount = shares.filter((share) => share.status === 'joined').length + return { + ...visibleOrder, + // 大厅接收端根据自身等级重新计算实际冻结金额,避免不同等级看到同一笔错误押金。 + requiredDepositAmount: mapped.requiredDepositAmount, + sharingProgress: { joinedQuantity, pendingSubmissionCount }, + } +} + +async function publishWorkerFinanceSnapshots(input: { + requestId: number + workerId: number + walletChanged?: boolean +}) { + const [request, worker] = await Promise.all([ + getWorkerFinanceRequestById(input.requestId), + getWorkerUserById(input.workerId), + ]) + const data = request ? mapFinanceRequest(request) : { requestId: Number(input.requestId) } + const operation = request ? 'upsert' : 'remove' publishRealtimeEvent({ type: 'worker_finance_request.changed', entityId: input.requestId, scopes: ['admin_worker_finance'], + operation, + data, admin: true, }) publishRealtimeEvent({ type: 'worker_finance_request.changed', entityId: input.requestId, scopes: ['worker_finance_requests'], + operation, + data, workerIds: [input.workerId], }) - if (input.walletChanged) { + if (input.walletChanged && worker) { publishRealtimeEvent({ type: 'worker_wallet.changed', entityId: input.workerId, scopes: ['worker_wallet'], + operation: 'upsert', + data: { wallet: mapWorkerUser(worker).wallet }, workerIds: [input.workerId], }) } } -function openRealtimeStream(response: Response, workerId: number | null) { +function openRealtimeStream(response: Response, workerId: number | null, lastEventId: number) { const subscriber: RealtimeSubscriber = { id: nextSubscriberId++, response, @@ -160,7 +372,9 @@ function openRealtimeStream(response: Response, workerId: number | null) { subscribers.set(subscriber.id, subscriber) writeEvent(subscriber, 'ready', { connectedAt: new Date().toISOString(), + lastEventId: Math.max(0, Number(lastEventId) || 0), }) + void replayEvents(subscriber, Math.max(0, Number(lastEventId) || 0)) ensureHeartbeat() response.on('close', () => { @@ -169,6 +383,56 @@ function openRealtimeStream(response: Response, workerId: number | null) { }) } +async function replayEvents(subscriber: RealtimeSubscriber, lastEventId: number) { + const events = await listRealtimeEventsAfter(lastEventId) + for (const row of events) { + const event = mapStoredEvent(row) + if (event.eventId <= lastEventId || !isEventVisibleToSubscriber(event, subscriber)) continue + writeEvent(subscriber, 'realtime', event) + } +} + +function mapStoredEvent(row: RealtimeEventRow): StoredRealtimeEvent { + const scopes = normalizeStringArray(row.scopes) + const workerIds = normalizeStringArray(row.worker_ids) + .map((item) => Number(item)) + .filter((item) => Number.isInteger(item) && item > 0) + const operation = + row.operation === 'upsert' || row.operation === 'remove' ? row.operation : ('refresh' as const) + return { + eventId: Number(row.id), + version: Number(row.id), + type: row.event_type as RealtimeEventType, + entityId: Number(row.entity_id || 0), + scopes, + operation, + ...(row.data === null || row.data === undefined ? {} : { data: row.data }), + occurredAt: row.occurred_at, + adminAudience: row.admin_audience === true, + workerIds, + broadcastWorkers: row.broadcast_workers === true, + } +} + +function normalizeStringArray(value: unknown): string[] { + const values = + typeof value === 'string' + ? (() => { + try { + return JSON.parse(value) + } catch { + return [] + } + })() + : value + return Array.isArray(values) ? values.map((item) => String(item || '')).filter(Boolean) : [] +} + +function isEventVisibleToSubscriber(event: StoredRealtimeEvent, subscriber: RealtimeSubscriber) { + if (subscriber.workerId === null) return event.adminAudience + return event.broadcastWorkers || event.workerIds.includes(Number(subscriber.workerId)) +} + function broadcast(subscribers: Iterable, event: RealtimeEvent) { for (const subscriber of subscribers) { writeEvent(subscriber, 'realtime', event) @@ -182,7 +446,12 @@ function writeEvent(subscriber: RealtimeSubscriber, name: string, payload: unkno } try { - subscriber.response.write(`event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`) + const eventId = + name === 'realtime' && payload && typeof payload === 'object' && 'eventId' in payload + ? Number((payload as { eventId?: unknown }).eventId || 0) + : 0 + const idLine = eventId > 0 ? `id: ${eventId}\n` : '' + subscriber.response.write(`${idLine}event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`) } catch { removeSubscriber(subscriber) } diff --git a/apps/backend/src/services/worker-platform/admin-service.ts b/apps/backend/src/services/worker-platform/admin-service.ts index 31de18a2..55e2bc41 100644 --- a/apps/backend/src/services/worker-platform/admin-service.ts +++ b/apps/backend/src/services/worker-platform/admin-service.ts @@ -70,8 +70,8 @@ import { resolveAdminNotificationEntity, } from '../admin/admin-notification-service.js' import { - publishRealtimeEvent, publishWorkerFinanceRealtimeChange, + publishWorkerWalletRealtimeChange, publishWorkOrderRealtimeChange, } from '../realtime/realtime-event-service.js' import { @@ -544,12 +544,7 @@ export async function assignAdminWorkOrderToWorker( 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)], - }) + publishWorkerWalletRealtimeChange(Number(worker.id)) return { order: mapWorkOrderAdmin(updated), voucherConsume } } @@ -578,12 +573,7 @@ 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)], - }) + publishWorkerWalletRealtimeChange(Number(workerId)) return { wallet: mapWallet(wallet) } } @@ -1515,12 +1505,7 @@ export async function acceptAdminWorkOrder(workOrderId: number | string, actorNa workerIds, }) for (const workerId of workerIds) { - publishRealtimeEvent({ - type: 'worker_wallet.changed', - entityId: workerId, - scopes: ['worker_wallet'], - workerIds: [workerId], - }) + publishWorkerWalletRealtimeChange(workerId) } return { order: mapWorkOrderAdmin(updated || workOrder) } } @@ -1603,12 +1588,7 @@ export async function unassignAdminWorkOrder(workOrderId: number | string, actor 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)], - }) + publishWorkerWalletRealtimeChange(Number(workOrder.assigned_worker_id)) return { order: mapWorkOrderAdmin(updated) } } @@ -1673,12 +1653,7 @@ export async function cancelAdminWorkOrder( hallChanged: action === 'return_to_hall', }) if (workerId) { - publishRealtimeEvent({ - type: 'worker_wallet.changed', - entityId: workerId, - scopes: ['worker_wallet'], - workerIds: [workerId], - }) + publishWorkerWalletRealtimeChange(workerId) } return { order: mapWorkOrderAdmin(updated), action } } diff --git a/apps/backend/src/services/worker-platform/worker-service.ts b/apps/backend/src/services/worker-platform/worker-service.ts index 090ded50..4104bc25 100644 --- a/apps/backend/src/services/worker-platform/worker-service.ts +++ b/apps/backend/src/services/worker-platform/worker-service.ts @@ -83,8 +83,8 @@ import { remindWorkerAcceptanceAdminNotification, } from '../admin/admin-notification-service.js' import { - publishRealtimeEvent, publishWorkerFinanceRealtimeChange, + publishWorkerWalletRealtimeChange, publishWorkOrderRealtimeChange, } from '../realtime/realtime-event-service.js' @@ -1223,12 +1223,7 @@ export async function grabWorkerHallOrder(workOrderId: number | string, session: workerIds: [Number(worker.id)], hallChanged: true, }) - publishRealtimeEvent({ - type: 'worker_wallet.changed', - entityId: Number(worker.id), - scopes: ['worker_wallet'], - workerIds: [Number(worker.id)], - }) + publishWorkerWalletRealtimeChange(Number(worker.id)) return { order: mapWorkOrderForWorker(grabbed.order, permissions) } } @@ -1266,12 +1261,7 @@ export async function settleDueDepositUnfreezes( }) if (released) { processedCount += 1 - publishRealtimeEvent({ - type: 'worker_wallet.changed', - entityId: Number(released.worker_id), - scopes: ['worker_wallet'], - workerIds: [Number(released.worker_id)], - }) + publishWorkerWalletRealtimeChange(Number(released.worker_id)) } } return { @@ -1313,12 +1303,7 @@ export async function settleOverdueWorkOrders(options: { limit?: number; workerI 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)], - }) + publishWorkerWalletRealtimeChange(Number(workOrder.assigned_worker_id)) } } else { skippedCount += 1 @@ -1652,12 +1637,7 @@ export async function submitWorkerOrderAcceptance( workOrderId: Number(workOrder.id), workerIds: [Number(worker.id)], }) - publishRealtimeEvent({ - type: 'worker_wallet.changed', - entityId: Number(worker.id), - scopes: ['worker_wallet'], - workerIds: [Number(worker.id)], - }) + publishWorkerWalletRealtimeChange(Number(worker.id)) return { order: mapWorkOrderForWorker(acceptedOrder, resolveWorkerPermissions(worker)), } @@ -2271,12 +2251,7 @@ async function publishWorkOrderParticipantChange( ...(options.hallChanged ? { hallChanged: true } : {}), }) for (const workerId of options.walletChangedWorkerIds || []) { - publishRealtimeEvent({ - type: 'worker_wallet.changed', - entityId: workerId, - scopes: ['worker_wallet'], - workerIds: [workerId], - }) + publishWorkerWalletRealtimeChange(workerId) } } diff --git a/apps/frontend/src/hooks/useRealtimeEvents.ts b/apps/frontend/src/hooks/useRealtimeEvents.ts index a692e263..fc6f6b40 100644 --- a/apps/frontend/src/hooks/useRealtimeEvents.ts +++ b/apps/frontend/src/hooks/useRealtimeEvents.ts @@ -2,6 +2,7 @@ import { useEffect, useRef } from 'react' export type RealtimeEvent = { eventId: number + version: number type: | 'admin_notification.changed' | 'order.changed' @@ -10,6 +11,8 @@ export type RealtimeEvent = { | 'worker_wallet.changed' entityId: number scopes: string[] + operation: 'upsert' | 'remove' | 'refresh' + data?: unknown occurredAt: string } @@ -35,18 +38,21 @@ export function useRealtimeEvents({ url, token, onEvent, onConnected }: UseRealt let stopped = false let abortController: AbortController | null = null let reconnectTimer = 0 + let lastEventId = 0 const connect = () => { if (stopped || document.visibilityState !== 'visible') return abortController?.abort() abortController = new AbortController() + const headers: Record = { + Accept: 'text/event-stream', + Authorization: `Bearer ${token}`, + } + if (lastEventId > 0) headers['Last-Event-ID'] = String(lastEventId) void fetch(url, { method: 'GET', - headers: { - Accept: 'text/event-stream', - Authorization: `Bearer ${token}`, - }, + headers, cache: 'no-store', signal: abortController.signal, }) @@ -58,7 +64,11 @@ export function useRealtimeEvents({ url, token, onEvent, onConnected }: UseRealt await readSseStream(response.body, (eventName, data) => { if (eventName !== 'realtime') return const event = normalizeRealtimeEvent(data) - if (event) onEventRef.current(event) + if (event) { + if (event.eventId <= lastEventId) return + lastEventId = Math.max(lastEventId, event.eventId) + onEventRef.current(event) + } }) }) .catch(() => undefined) @@ -153,11 +163,17 @@ function normalizeRealtimeEvent(value: string): RealtimeEvent | null { } return { eventId: Number(parsed.eventId || 0), + version: Number(parsed.version || 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) : [], + operation: + parsed.operation === 'upsert' || parsed.operation === 'remove' + ? parsed.operation + : 'refresh', + data: parsed.data, occurredAt: String(parsed.occurredAt || ''), } } catch { diff --git a/apps/frontend/src/layouts/AdminLayout.tsx b/apps/frontend/src/layouts/AdminLayout.tsx index b6ac9826..8f5a0092 100644 --- a/apps/frontend/src/layouts/AdminLayout.tsx +++ b/apps/frontend/src/layouts/AdminLayout.tsx @@ -32,7 +32,7 @@ import { } from 'antd' import type { MenuProps } from 'antd' import { useQuery, useQueryClient } from '@tanstack/react-query' -import { useEffect, useMemo, useRef, useState } from 'react' +import { useMemo, useState } from 'react' import { Outlet, useLocation, useNavigate } from 'react-router' import { @@ -53,7 +53,6 @@ import { useRealtimeEvents, type RealtimeEvent } from '@/hooks/useRealtimeEvents const SIDEBAR_COLLAPSED_KEY = 'react-admin-sidebar-collapsed' const NOTIFICATION_SOUND_KEY = 'admin-notification-sound-enabled' -const NOTIFICATION_SOUND_COOLDOWN_MS = 30_000 export default function AdminLayout() { const navigate = useNavigate() @@ -69,22 +68,16 @@ export default function AdminLayout() { const [soundEnabled, setSoundEnabled] = useState( () => window.localStorage.getItem(NOTIFICATION_SOUND_KEY) === '1', ) - const seenNotificationVersions = useRef | null>(null) - const lastSoundAt = useRef(0) const devMockStatusQuery = useQuery({ queryKey: ['admin-dev-mock-status'], queryFn: () => fetchDevMockStatus(), - staleTime: 60_000, retry: false, }) const devMockEnabled = Boolean(devMockStatusQuery.data?.data?.enabled) const notificationsQuery = useQuery({ queryKey: ['admin-notifications'], queryFn: fetchAdminNotifications, - refetchInterval: 10_000, - refetchIntervalInBackground: false, - refetchOnWindowFocus: true, retry: false, }) const notifications = notificationsQuery.data?.data.items || [] @@ -93,39 +86,9 @@ export default function AdminLayout() { useRealtimeEvents({ url: '/api/v1/admin/realtime', token: getAdminToken(), - onEvent: (event) => handleAdminRealtimeEvent(event, queryClient), - onConnected: () => { - void invalidateAdminRealtimeQueries(queryClient) - }, + onEvent: (event) => applyAdminRealtimeEvent(event, queryClient), }) - useEffect(() => { - const currentVersions = new Map( - notifications.map((notification) => [notification.notificationId, notification.updatedAt]), - ) - if (seenNotificationVersions.current === null) { - seenNotificationVersions.current = currentVersions - return - } - const changedNotifications = notifications.filter( - (notification) => - seenNotificationVersions.current?.get(notification.notificationId) !== - notification.updatedAt, - ) - seenNotificationVersions.current = currentVersions - if (changedNotifications.length === 0) return - - void Promise.all([ - queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-finance-requests'] }), - queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-orders'] }), - queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-summary'] }), - ]) - if (soundEnabled && Date.now() - lastSoundAt.current >= NOTIFICATION_SOUND_COOLDOWN_MS) { - lastSoundAt.current = Date.now() - announceNotification(selectNotificationToAnnounce(changedNotifications)) - } - }, [notifications, queryClient, soundEnabled]) - const menuItems = useMemo(() => { const operationItems: MenuProps['items'] = [ { key: '/admin/dashboard', icon: , label: '概览' }, @@ -354,38 +317,215 @@ export default function AdminLayout() { ) } -function handleAdminRealtimeEvent( +function applyAdminRealtimeEvent( event: RealtimeEvent, queryClient: ReturnType, ) { - if (event.type === 'admin_notification.changed') { - void queryClient.invalidateQueries({ queryKey: ['admin-notifications'] }) + if (event.type === 'admin_notification.changed' && isRecord(event.data)) { + queryClient.setQueryData(['admin-notifications'], (current: unknown) => { + if (!isRecord(current) || !isRecord(current.data)) return current + const currentData = current.data as { items?: unknown[]; pendingCount?: number } + const items = Array.isArray(currentData.items) ? currentData.items : [] + const notification = event.data as Record + const notificationId = Number(notification.notificationId || event.entityId) + const nextItems = + event.operation === 'remove' + ? items.filter( + (item) => + Number((item as { notificationId?: unknown })?.notificationId) !== notificationId, + ) + : upsertById(items, notification, 'notificationId', true) + return { + ...current, + data: { + ...currentData, + items: nextItems, + pendingCount: + event.operation === 'remove' + ? Math.max(0, Number(currentData.pendingCount || 0) - 1) + : currentData.pendingCount, + }, + } + }) } - 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' && + event.scopes.includes('admin_work_orders') && + isRecord(event.data) + ) { + const data = event.data + patchAdminWorkOrderCaches(queryClient, { ...event, data }) } - 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' && + event.scopes.includes('admin_worker_finance') && + isRecord(event.data) + ) { + const data = event.data + queryClient.setQueriesData( + { queryKey: ['admin-worker-platform-finance-requests'] }, + (current: unknown) => patchFinanceRequestListEnvelope(current, { ...event, data }), + ) } - if (event.type === 'worker_finance_request.changed') { - void queryClient.invalidateQueries({ queryKey: ['admin-worker-platform-finance-requests'] }) + + if ( + event.type === 'order.changed' && + event.scopes.includes('admin_orders') && + isRecord(event.data) + ) { + const payload = event.data + if (isRecord(payload.order)) { + patchAdminOrderCaches(queryClient, { ...event, data: payload.order }) + } + if (isRecord(payload.dashboard)) { + queryClient.setQueryData(['admin-dashboard-summary'], (current: unknown) => { + if (!isRecord(current)) return current + return { ...current, data: payload.dashboard } + }) + } + if (Array.isArray(payload.tasks)) { + patchAdminTaskCaches(queryClient, payload.tasks) + } } } -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'] }), - ]) +function patchAdminWorkOrderCaches( + queryClient: ReturnType, + event: RealtimeEvent, +) { + const item = event.data as { + workOrderId?: unknown + status?: unknown + productName?: unknown + platformOrderId?: unknown + worker?: { workerId?: unknown } | null + } + for (const [queryKey, current] of queryClient.getQueriesData({ + queryKey: ['admin-worker-platform-orders'], + })) { + const statuses = String(queryKey[1] || '') + .split(',') + .filter(Boolean) + const keyword = String(queryKey[2] || '') + const workerId = Number(queryKey[3] || 0) + const targetWorkOrderId = Number(queryKey[4] || 0) + const page = Number(queryKey[5] || 1) + const matches = + (!statuses.length || statuses.includes(String(item.status || ''))) && + matchesWorkOrderKeyword(item, keyword) && + (!workerId || workerId === Number(item.worker?.workerId || 0)) && + (!targetWorkOrderId || targetWorkOrderId === Number(item.workOrderId || 0)) + const shouldInsert = event.operation === 'upsert' && page === 1 && matches + const operation = event.operation === 'remove' || !matches ? 'remove' : 'upsert' + queryClient.setQueryData( + queryKey, + patchWorkOrderListEnvelope(current, { ...event, operation }, shouldInsert), + ) + } +} + +function patchAdminTaskCaches(queryClient: ReturnType, tasks: unknown[]) { + queryClient.setQueriesData({ queryKey: ['admin-tasks'] }, (current: unknown) => { + if (!isRecord(current) || !isRecord(current.data) || !Array.isArray(current.data.items)) { + return current + } + const items = current.data.items as unknown[] + const nextItems = items.map((item) => { + const taskId = Number((item as { taskId?: unknown }).taskId || 0) + return ( + tasks.find((task) => Number((task as { taskId?: unknown }).taskId || 0) === taskId) || item + ) + }) + return { ...current, data: { ...current.data, items: nextItems } } + }) +} + +function patchAdminOrderCaches( + queryClient: ReturnType, + event: RealtimeEvent, +) { + const item = event.data as { orderId?: unknown; platformOrderId?: unknown; payStatus?: unknown } + for (const [queryKey, current] of queryClient.getQueriesData({ queryKey: ['admin-orders'] })) { + const params = isRecord(queryKey[1]) ? queryKey[1] : {} + const page = Number(params.page || 1) + const matches = + (!String(params.platformOrderId || '').trim() || + String(item.platformOrderId || '') === String(params.platformOrderId || '').trim()) && + (!String(params.payStatus || '').trim() || + String(item.payStatus || '') === String(params.payStatus || '').trim()) + const shouldInsert = event.operation === 'upsert' && page === 1 && matches && !params.skuCode + const operation = event.operation === 'remove' || !matches ? 'remove' : 'upsert' + queryClient.setQueryData( + queryKey, + patchAdminOrderListEnvelope(current, { ...event, operation }, shouldInsert), + ) + } +} + +function patchAdminOrderListEnvelope(current: unknown, event: RealtimeEvent, insert = false) { + if (!isRecord(current) || !isRecord(current.data) || !Array.isArray(current.data.items)) { + return current + } + const id = Number((event.data as { orderId?: unknown }).orderId || event.entityId) + const items = current.data.items as unknown[] + const nextItems = + event.operation === 'remove' + ? items.filter((entry) => Number((entry as { orderId?: unknown })?.orderId) !== id) + : upsertById(items, event.data, 'orderId', insert) + return { ...current, data: { ...current.data, items: nextItems } } +} + +function patchWorkOrderListEnvelope(current: unknown, event: RealtimeEvent, insert = false) { + if (!isRecord(current) || !isRecord(current.data) || !Array.isArray(current.data.items)) { + return current + } + const item = event.data + const id = Number((item as { workOrderId?: unknown }).workOrderId || event.entityId) + const items = current.data.items as unknown[] + const nextItems = + event.operation === 'remove' + ? items.filter((entry) => Number((entry as { workOrderId?: unknown })?.workOrderId) !== id) + : upsertById(items, item, 'workOrderId', insert) + return { ...current, data: { ...current.data, items: nextItems } } +} + +function patchFinanceRequestListEnvelope(current: unknown, event: RealtimeEvent) { + if (!isRecord(current) || !isRecord(current.data) || !Array.isArray(current.data.items)) { + return current + } + const id = Number((event.data as { requestId?: unknown }).requestId || event.entityId) + const items = current.data.items as unknown[] + const nextItems = + event.operation === 'remove' + ? items.filter((entry) => Number((entry as { requestId?: unknown })?.requestId) !== id) + : upsertById(items, event.data, 'requestId') + return { ...current, data: { ...current.data, items: nextItems } } +} + +function upsertById(items: unknown[], nextItem: unknown, idKey: string, insert = false) { + const id = Number((nextItem as Record)[idKey] || 0) + const index = items.findIndex( + (item) => Number((item as Record)?.[idKey] || 0) === id, + ) + if (index < 0) return insert ? [nextItem, ...items] : items + return items.map((item, itemIndex) => (itemIndex === index ? nextItem : item)) +} + +function matchesWorkOrderKeyword( + item: { productName?: unknown; platformOrderId?: unknown }, + keyword: string, +) { + const normalizedKeyword = keyword.trim().toLowerCase() + if (!normalizedKeyword) return true + return [item.productName, item.platformOrderId] + .map((value) => String(value || '').toLowerCase()) + .some((value) => value.includes(normalizedKeyword)) +} + +function isRecord(value: unknown): value is Record { + return Boolean(value && typeof value === 'object' && !Array.isArray(value)) } type NotificationAlertKind = @@ -396,60 +536,6 @@ type NotificationAlertKind = | 'reminder' | 'default' -function selectNotificationToAnnounce(notifications: AdminNotification[]) { - return [...notifications].sort((left, right) => { - const priorityDiff = notificationAlertRank(right) - notificationAlertRank(left) - if (priorityDiff !== 0) return priorityDiff - return right.notificationId - left.notificationId - })[0] -} - -function notificationAlertRank(notification: AdminNotification) { - if (notification.notificationType === 'worker_acceptance_reminded') return 3 - if (notification.notificationType === 'worker_withdraw_requested') return 2 - return 1 -} - -function announceNotification(notification?: AdminNotification) { - if (!notification?.soundEnabled) return - const { text, kind } = resolveNotificationAnnouncement(notification) - playNotificationSound(kind) - if (!('speechSynthesis' in window) || typeof SpeechSynthesisUtterance === 'undefined') return - - try { - window.speechSynthesis.cancel() - const utterance = new SpeechSynthesisUtterance(text) - utterance.lang = 'zh-CN' - utterance.rate = 1.05 - utterance.volume = 0.9 - window.speechSynthesis.speak(utterance) - } catch { - // 语音播报不可用时,前面的区分音调仍会提示新待办。 - } -} - -function resolveNotificationAnnouncement(notification?: AdminNotification): { - text: string - kind: NotificationAlertKind -} { - if (notification?.notificationType === 'worker_acceptance_reminded') { - return { text: '你有新的催验收工单', kind: 'reminder' } - } - if (notification?.notificationType === 'worker_withdraw_requested') { - return { text: '你有新的提现申请', kind: 'withdraw' } - } - if (notification?.notificationType === 'worker_recharge_requested') { - return { text: '你有新的充值申请', kind: 'recharge' } - } - if (notification?.notificationType === 'worker_acceptance_submitted') { - return { text: '你有新的待验收工单', kind: 'acceptance' } - } - if (notification?.notificationType === 'work_order_material_required') { - return { text: '你有工单待补资料', kind: 'material' } - } - return { text: '你有新的待处理事项', kind: 'default' } -} - function playNotificationSound(kind: NotificationAlertKind) { try { const audio = new AudioContext() @@ -477,7 +563,7 @@ function playNotificationSound(kind: NotificationAlertKind) { oscillator.stop(audio.currentTime + sequence.length * duration) oscillator.addEventListener('ended', () => void audio.close()) } catch { - // 部分浏览器禁止音频或不支持 Web Audio,不影响待办轮询。 + // 部分浏览器禁止音频或不支持 Web Audio,不影响待办展示。 } } diff --git a/apps/frontend/src/layouts/WorkerLayout.tsx b/apps/frontend/src/layouts/WorkerLayout.tsx index f54a1fa7..b736df75 100644 --- a/apps/frontend/src/layouts/WorkerLayout.tsx +++ b/apps/frontend/src/layouts/WorkerLayout.tsx @@ -32,18 +32,15 @@ export default function WorkerLayout() { const location = useLocation() const isMobile = useIsMobile() const { message, modal } = App.useApp() + const queryClient = useQueryClient() 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) - }, + onEvent: (event) => applyWorkerRealtimeEvent(event, queryClient), }) async function submitLogout() { @@ -156,40 +153,179 @@ export default function WorkerLayout() { ) } -function handleWorkerRealtimeEvent( +function applyWorkerRealtimeEvent( 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.type === 'work_order.changed' && isRecord(event.data)) { + const data = event.data if (event.scopes.includes('worker_my_orders')) { - void queryClient.invalidateQueries({ queryKey: ['worker-my-orders'] }) - void queryClient.invalidateQueries({ queryKey: ['worker-my-orders-overview'] }) + patchWorkerMyOrderCaches(queryClient, { ...event, data }) + } + if (event.scopes.includes('worker_hall')) { + patchWorkerHallCaches(queryClient, { ...event, data }) } } - 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' && isRecord(event.data)) { + const walletData = event.data + queryClient.setQueryData(['worker-profile'], (current: unknown) => { + if (!isRecord(current) || !isRecord(current.data) || !isRecord(current.data.worker)) + return current + return { + ...current, + data: { + ...current.data, + worker: { ...current.data.worker, wallet: walletData.wallet || walletData }, + }, + } + }) } - if (event.type === 'worker_wallet.changed') { - void queryClient.invalidateQueries({ queryKey: ['worker-wallet-ledgers'] }) - void queryClient.invalidateQueries({ queryKey: ['worker-profile'] }) + + if (event.type === 'worker_finance_request.changed' && isRecord(event.data)) { + const data = event.data + queryClient.setQueriesData({ queryKey: ['worker-finance-requests'] }, (current: unknown) => + patchFinanceRequestListEnvelope(current, { ...event, data }), + ) } } -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 patchWorkerMyOrderCaches( + queryClient: ReturnType, + event: RealtimeEvent, +) { + for (const [queryKey, current] of queryClient.getQueriesData({ + queryKey: ['worker-my-orders'], + })) { + const status = String(queryKey[1] || '') + const keyword = String(queryKey[2] || '') + const page = Number(queryKey[3] || 1) + const item = event.data as { + status?: unknown + productName?: unknown + platformOrderId?: unknown + } + const matchesStatus = !status || status === String(item.status || '') + const matchesKeyword = matchesWorkOrderKeyword(item, keyword) + const matches = matchesStatus && matchesKeyword + const shouldInsert = event.operation === 'upsert' && page === 1 && matches + const operation = event.operation === 'remove' || !matches ? 'remove' : 'upsert' + queryClient.setQueryData( + queryKey, + patchWorkOrderListEnvelope(current, { ...event, operation }, shouldInsert), + ) + } +} + +function patchWorkerHallCaches( + queryClient: ReturnType, + event: RealtimeEvent, +) { + const workerSnapshot = applyWorkerDepositRule(queryClient, event.data) + const item = workerSnapshot as { + status?: unknown + categoryId?: unknown + productName?: unknown + platformOrderId?: unknown + } + for (const [queryKey, current] of queryClient.getQueriesData({ + queryKey: ['worker-hall-orders'], + })) { + const keyword = String(queryKey[1] || '') + const categoryId = Number(queryKey[2] || 0) + const matchesCategory = !categoryId || categoryId === Number(item.categoryId || 0) + const matches = + String(item.status || '') === 'open' && + matchesCategory && + matchesWorkOrderKeyword(item, keyword) + const shouldInsert = event.operation === 'upsert' && matches + const operation = event.operation === 'remove' || !matches ? 'remove' : 'upsert' + queryClient.setQueryData( + queryKey, + patchInfiniteWorkOrderEnvelope( + current, + { ...event, data: workerSnapshot, operation }, + shouldInsert, + ), + ) + } +} + +function applyWorkerDepositRule(queryClient: ReturnType, data: unknown) { + if (!isRecord(data)) return data + const profile = queryClient.getQueryData(['worker-profile']) + const depositFreeAmount = + isRecord(profile) && isRecord(profile.data) && isRecord(profile.data.worker) + ? Number( + (profile.data.worker.level as { permissions?: { depositFreeAmount?: unknown } } | null) + ?.permissions?.depositFreeAmount || 0, + ) + : 0 + const requiredDepositAmount = Number(data.requiredDepositAmount || 0) + return { + ...data, + freezeDepositAmount: Math.max(0, requiredDepositAmount - depositFreeAmount), + } +} + +function matchesWorkOrderKeyword( + item: { productName?: unknown; platformOrderId?: unknown }, + keyword: string, +) { + const normalizedKeyword = keyword.trim().toLowerCase() + if (!normalizedKeyword) return true + return [item.productName, item.platformOrderId] + .map((value) => String(value || '').toLowerCase()) + .some((value) => value.includes(normalizedKeyword)) +} + +function patchWorkOrderListEnvelope(current: unknown, event: RealtimeEvent, insert = false) { + if (!isRecord(current) || !isRecord(current.data) || !Array.isArray(current.data.items)) { + return current + } + const id = Number((event.data as { workOrderId?: unknown }).workOrderId || event.entityId) + const items = current.data.items as unknown[] + const nextItems = + event.operation === 'remove' + ? items.filter((entry) => Number((entry as { workOrderId?: unknown })?.workOrderId) !== id) + : upsertById(items, event.data, 'workOrderId', insert) + return { ...current, data: { ...current.data, items: nextItems } } +} + +function patchInfiniteWorkOrderEnvelope(current: unknown, event: RealtimeEvent, insert = false) { + if (!isRecord(current) || !Array.isArray(current.pages)) return current + return { + ...current, + pages: current.pages.map((page: unknown, index) => + patchWorkOrderListEnvelope(page, event, insert && index === 0), + ), + } +} + +function patchFinanceRequestListEnvelope(current: unknown, event: RealtimeEvent) { + if (!isRecord(current) || !isRecord(current.data) || !Array.isArray(current.data.items)) { + return current + } + const id = Number((event.data as { requestId?: unknown }).requestId || event.entityId) + const items = current.data.items as unknown[] + const nextItems = + event.operation === 'remove' + ? items.filter((entry) => Number((entry as { requestId?: unknown })?.requestId) !== id) + : upsertById(items, event.data, 'requestId') + return { ...current, data: { ...current.data, items: nextItems } } +} + +function upsertById(items: unknown[], nextItem: unknown, idKey: string, insert = false) { + const id = Number((nextItem as Record)[idKey] || 0) + const index = items.findIndex( + (item) => Number((item as Record)?.[idKey] || 0) === id, + ) + if (index < 0) return insert ? [nextItem, ...items] : items + return items.map((item, itemIndex) => (itemIndex === index ? nextItem : item)) +} + +function isRecord(value: unknown): value is Record { + return Boolean(value && typeof value === 'object' && !Array.isArray(value)) } function resolveSelectedKey(pathname: string) { diff --git a/apps/frontend/src/main.tsx b/apps/frontend/src/main.tsx index ac7860b8..503b1c56 100644 --- a/apps/frontend/src/main.tsx +++ b/apps/frontend/src/main.tsx @@ -14,7 +14,10 @@ import './styles/main.css' const queryClient = new QueryClient({ defaultOptions: { queries: { + // 页面切换和网络恢复不主动重新请求,数据由首次加载、用户操作或手动刷新触发更新。 + staleTime: Infinity, refetchOnWindowFocus: false, + refetchOnReconnect: false, retry: 1, }, }, diff --git a/apps/frontend/src/pages/worker/WorkerHallPage.tsx b/apps/frontend/src/pages/worker/WorkerHallPage.tsx index a2a641d2..308741ea 100644 --- a/apps/frontend/src/pages/worker/WorkerHallPage.tsx +++ b/apps/frontend/src/pages/worker/WorkerHallPage.tsx @@ -19,7 +19,6 @@ import { Modal, Space, Spin, - Switch, Tabs, Tag, Tooltip, @@ -57,22 +56,8 @@ export default function WorkerHallPage() { const [leaderboardOpen, setLeaderboardOpen] = useState(false) const [sharingQuantity, setSharingQuantity] = useState(1) const [submittingSharing, setSubmittingSharing] = useState(false) - const [autoRefresh, setAutoRefresh] = useState(loadAutoRefreshPreference) const loadMoreRef = useRef(null) - useEffect(() => { - if (!autoRefresh) return - const timer = window.setInterval(() => { - void refreshAll() - }, 5000) - return () => window.clearInterval(timer) - }, [autoRefresh]) - - function toggleAutoRefresh(enabled: boolean) { - setAutoRefresh(enabled) - localStorage.setItem(AUTO_REFRESH_STORAGE_KEY, enabled ? '1' : '0') - } - const ordersQuery = useInfiniteQuery({ queryKey: ['worker-hall-orders', keyword, categoryId], initialPageParam: 1, @@ -129,7 +114,6 @@ export default function WorkerHallPage() { queryKey: ['worker-hall-leaderboard'], queryFn: () => fetchWorkerHallLeaderboard(), enabled: leaderboardOpen, - staleTime: 30_000, retry: false, }) @@ -299,10 +283,6 @@ export default function WorkerHallPage() { onClick={() => setLeaderboardOpen(true)} /> -
- 自动刷新 - -