Files
order_site/apps/backend/src/services/realtime/realtime-event-service.ts
T

521 lines
16 KiB
TypeScript

import type { Response } from 'express'
import {
getWorkerFinanceRequestById,
getWorkerUserById,
getWorkOrderById,
listWorkerWorkOrderNotes,
listWorkerWorkOrderViews,
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 {
isWorkOrderSharingFilled,
mapFinanceRequest,
mapWorkOrderAdmin,
mapWorkOrderForWorker,
mapWorkerUser,
resolveWorkerPermissions,
} from '../worker-platform/mappers.js'
export type RealtimeEventType =
| 'admin_notification.changed'
| 'order.changed'
| 'work_order.changed'
| 'worker_finance_request.changed'
| 'worker_wallet.changed'
export type RealtimeEvent = {
eventId: number
version: number
type: RealtimeEventType
entityId: number
scopes: string[]
operation: 'upsert' | 'remove' | 'refresh'
data?: unknown
occurredAt: string
}
type RealtimeEventInput = {
type: RealtimeEventType
entityId?: number
scopes?: string[]
operation?: RealtimeEvent['operation']
data?: unknown
admin?: boolean
workerIds?: number[]
broadcastWorkers?: boolean
}
type RealtimeSubscriber = {
id: number
response: Response
workerId: number | null
}
type StoredRealtimeEvent = RealtimeEvent & {
adminAudience: boolean
workerIds: number[]
broadcastWorkers: boolean
}
const adminSubscribers = new Map<number, RealtimeSubscriber>()
const workerSubscribers = new Map<number, RealtimeSubscriber>()
let nextSubscriberId = 1
let heartbeatTimer: NodeJS.Timeout | null = null
let publishQueue: Promise<void> = Promise.resolve()
/** 建立后台实时事件流,事件快照用于前端定点更新缓存。 */
export function openAdminRealtimeStream(response: Response, lastEventId = 0) {
return openRealtimeStream(response, null, lastEventId)
}
/** 建立打手实时事件流,只向该打手推送其可见的订单和资金变化。 */
export function openWorkerRealtimeStream(response: Response, workerId: number, lastEventId = 0) {
return openRealtimeStream(response, Number(workerId), lastEventId)
}
/** 发布带版本号和数据快照的实时事件;数据库仍是最终权威数据源。 */
export function publishRealtimeEvent(input: RealtimeEventInput) {
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)
}
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)
}
if (storedEvent.broadcastWorkers) {
broadcast(workerSubscribers.values(), storedEvent)
return
}
if (storedEvent.workerIds.length > 0) {
broadcast(
[...workerSubscribers.values()].filter((subscriber) =>
storedEvent.workerIds.includes(Number(subscriber.workerId)),
),
storedEvent,
)
}
}
/** 工单变更统一通知后台;大厅与打手本人按可见范围分别接收。 */
export function publishWorkOrderRealtimeChange(input: {
workOrderId: number
workerIds?: number[]
hallChanged?: boolean
}) {
const workOrderId = Math.max(0, Number(input.workOrderId) || 0)
void publishWorkOrderSnapshots({ ...input, workOrderId }).catch(() => {
publishRealtimeEvent({
type: 'work_order.changed',
entityId: workOrderId,
scopes: ['admin_work_orders'],
operation: 'refresh',
admin: true,
})
})
}
/** 资金申请变更会同步后台资金页与申请所属打手的钱包页面。 */
export function publishWorkerFinanceRealtimeChange(input: {
requestId: number
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, shares),
pendingUnfreezeAmount: pendingUnfreezeByOrderId.get(workOrderId) || 0,
sharingProgress: {
joinedQuantity,
pendingSubmissionCount,
remainingQuantity: Math.max(
0,
Number(workOrder.sharing_total_quantity || 1) - joinedQuantity,
),
filled: isWorkOrderSharingFilled(
workOrder.sharing_enabled === true,
Number(workOrder.sharing_total_quantity || 1),
joinedQuantity,
),
},
}
publishRealtimeEvent({
type: 'work_order.changed',
entityId: workOrderId,
scopes: ['admin_work_orders'],
operation: 'upsert',
data: adminSnapshot,
admin: true,
})
if (input.hallChanged) {
const sharingFilled = isWorkOrderSharingFilled(
workOrder.sharing_enabled === true,
Number(workOrder.sharing_total_quantity || 1),
joinedQuantity,
)
publishRealtimeEvent({
type: 'work_order.changed',
entityId: workOrderId,
scopes: ['worker_hall'],
operation: workOrder.status === 'open' && !sharingFilled ? '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, views] = await Promise.all([
listWorkerWorkOrderNotes(workerId, [workOrderId]),
listWorkerWorkOrderViews(workerId, [workOrderId]),
])
const note = notes.find((item) => Number(item.work_order_id) === workOrderId)?.note || ''
const viewedAt =
views.find((item) => Number(item.work_order_id) === workOrderId)?.viewed_at || null
const share = shares.find((item) => Number(item.worker_id) === workerId)
const isAssignedWorker = Number(workOrder.assigned_worker_id || 0) === workerId
const workerOrderVisible = isAssignedWorker || Boolean(share && share.status !== 'cancelled')
publishRealtimeEvent({
type: 'work_order.changed',
entityId: workOrderId,
scopes: ['worker_my_orders'],
operation: workerOrderVisible ? 'upsert' : 'remove',
data: mapWorkOrderForWorker(
workOrder,
resolveWorkerPermissions(worker),
share,
note,
viewedAt,
),
workerIds: [workerId],
})
}
}
function mapWorkerHallWorkOrder(
workOrder: NonNullable<Awaited<ReturnType<typeof getWorkOrderById>>>,
shares: Awaited<ReturnType<typeof listWorkOrderShares>>,
) {
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 && 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, lastEventId: number) {
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(),
lastEventId: Math.max(0, Number(lastEventId) || 0),
})
void replayEvents(subscriber, Math.max(0, Number(lastEventId) || 0))
ensureHeartbeat()
response.on('close', () => {
subscribers.delete(subscriber.id)
stopHeartbeatWhenIdle()
})
}
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<RealtimeSubscriber>, 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 {
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)
}
}
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
}