增加订单实时状态同步
This commit is contained in:
@@ -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<number, RealtimeSubscriber>()
|
||||
const workerSubscribers = new Map<number, RealtimeSubscriber>()
|
||||
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<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 {
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user