改造实时事件缓存同步
This commit is contained in:
@@ -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 实时事件日志,保留最近事件以支持断线后的增量补发';
|
||||
@@ -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<RealtimeEventRow> {
|
||||
const result = await query<RealtimeEventRow>(
|
||||
`
|
||||
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<RealtimeEventRow>(
|
||||
`
|
||||
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
|
||||
}
|
||||
@@ -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<string, unknown>
|
||||
}) {
|
||||
return Number(request.get('Last-Event-ID') || request.query.lastEventId || 0) || 0
|
||||
}
|
||||
|
||||
export default router
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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 })
|
||||
}
|
||||
|
||||
@@ -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 })
|
||||
|
||||
@@ -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<number, RealtimeSubscriber>()
|
||||
const workerSubscribers = new Map<number, RealtimeSubscriber>()
|
||||
let nextSubscriberId = 1
|
||||
let nextEventId = 1
|
||||
let heartbeatTimer: NodeJS.Timeout | null = null
|
||||
let publishQueue: Promise<void> = 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<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) {
|
||||
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<RealtimeSubscriber>, 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)
|
||||
}
|
||||
|
||||
@@ -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 }
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user