From 2a3dcb0c68c52476d4ad38b9924607ce45acdcf5 Mon Sep 17 00:00:00 2001 From: hectorzhao Date: Mon, 31 Aug 2026 17:52:11 +0800 Subject: [PATCH] fix: reconcile caller aggregates and strengthen evidence handling --- .../caller-analytics.service.ts | 2 +- apps/web/src/pages/CallerAnalyticsPage.jsx | 5 +-- apps/worker-cdr/src/analytics-main.ts | 4 +-- .../database/src/caller-analytics-model.ts | 19 ++++++----- .../database/src/caller-analytics-store.ts | 2 +- .../database/src/caller-analytics.spec.ts | 13 ++++++++ scripts/rebuild-caller-analytics.mjs | 25 +++++++++++++++ tests/api/caller-analytics-real.mjs | 32 +++++++++++++++++-- 8 files changed, 85 insertions(+), 17 deletions(-) create mode 100644 scripts/rebuild-caller-analytics.mjs diff --git a/apps/api/src/modules/caller-analytics/caller-analytics.service.ts b/apps/api/src/modules/caller-analytics/caller-analytics.service.ts index b2d087a..4f86dcf 100644 --- a/apps/api/src/modules/caller-analytics/caller-analytics.service.ts +++ b/apps/api/src/modules/caller-analytics/caller-analytics.service.ts @@ -23,7 +23,7 @@ export function parseAnalyticsQuery(raw: Query, now = new Date()): AnalyticsQuer const from = raw.from ? new Date(str('from', 40)) : new Date(to.getTime() - number('minutes', 15, 1, 10080) * 60000); if (!Number.isFinite(+from) || !Number.isFinite(+to) || from >= to || +to - +from > 7 * 86400000 || +to > +now + 5000) throw new BadRequestException('Time range must be within 7 days and not in future'); const sort = str('sort', 32) || 'totalCalls'; - if (!['caller', ...COUNT_KEYS, 'connectionRate', 'overallAnswerRate', 'connectedAnswerRate'].includes(sort)) throw new BadRequestException('Invalid sort'); + if (!['caller', ...COUNT_KEYS, 'notConnectedCalls', 'connectionRate', 'overallAnswerRate', 'connectedAnswerRate'].includes(sort)) throw new BadRequestException('Invalid sort'); const order = str('order', 4) || 'desc'; if (!['asc', 'desc'].includes(order)) throw new BadRequestException('Invalid order'); const result: AnalyticsQuery = { view: view as AnalyticsQuery['view'], from, to, customerId: str('customerId', 32), caller: str('caller', 64), diff --git a/apps/web/src/pages/CallerAnalyticsPage.jsx b/apps/web/src/pages/CallerAnalyticsPage.jsx index 9edd593..a19df12 100644 --- a/apps/web/src/pages/CallerAnalyticsPage.jsx +++ b/apps/web/src/pages/CallerAnalyticsPage.jsx @@ -27,6 +27,7 @@ export function CallerAnalyticsPage() { setLoading(true); try { const q = { ...filters, onlyAlerts: String(filters.onlyAlerts) }; + if (q.minutes !== 'custom') { delete q.from; delete q.to; } if (q.minutes === 'today') { const now = new Date(); q.from = new Date(now.toLocaleDateString('en-CA', { timeZone: 'Asia/Shanghai' }) + 'T00:00:00+08:00').toISOString(); delete q.minutes; } if (q.minutes === 'custom') { q.from = new Date(q.from).toISOString(); q.to = new Date(q.to).toISOString(); delete q.minutes; } const result = await api.callerAnalytics(q, { signal: controller.signal }); @@ -43,7 +44,7 @@ export function CallerAnalyticsPage() { .then(r => { setDetail(r); setDetailError(''); }).catch(e => { if (!c.signal.aborted) setDetailError(explainApiError(e)); }); return () => c.abort(); }, [target, data, filters, detailSkip]); - const change = (key, value) => setDraft(d => ({ ...d, [key]: value, ...(key === 'view' && value === 'original' ? { vendorGatewayId: '' } : {}) })); + const change = (key, value) => setDraft(d => ({ ...d, [key]: value, ...(key === 'view' && value === 'original' ? { vendorGatewayId: '' } : {}), ...(key === 'customerId' ? {customerGatewayId:''} : {}) })); const names = Object.fromEntries(options.customers.map(c => [c.id, c.name])); const gatewayNames = Object.fromEntries(options.vendorGateways.map(c => [c.id, c.name])); const select = (label, key, choices) => ; @@ -75,7 +76,7 @@ export function CallerAnalyticsPage() {
{[...countColumns, ...rateColumns].map(c =>
{c.label}{metrics ? c.render ? c.render(metrics) : metrics[c.key] : '—'}
)}

{data ? `${formatDateTime(data.from)} — ${formatDateTime(data.to)} · ${data.qualityStatus === 'LIVE' ? '采集运行中' : '数据滞后或不完整'} · 待接通 ${metrics.pendingCalls} · 已结束未接通 ${metrics.failedCalls} · 未知 ${metrics.unknownCalls} · 更新于 ${formatDateTime(data.asOf)}` : '正在读取真实统计数据…'}

- {data ? {data.historyNotice} {data.summaryScope}。183 不代表实际振铃;未接通数含待接通,进行中指标为暂定值。 : null} + {data ? {data.historyNotice} {data.summaryScope}。183 不代表实际振铃;未接通数含待接通,进行中指标为暂定值。采集运行状态不代表窗口数据完整性已核验,端到端延迟尚未测定。 : null} 灰:呼叫 / 蓝:接通 / 绿:应答}>
{timeline.map((r,i) =>
{new Date(Number(r.time)).toLocaleTimeString('zh-CN',{hour:'2-digit',minute:'2-digit'})}
)}
{!timeline.length ?

当前窗口暂无已采集呼叫

:
查看趋势数值与三个率 ({ ...r, id:i, minute:formatDateTime(Number(r.time)) }))} columns={[{key:'minute',label:'分钟'},...countColumns,...rateColumns]} />
} diff --git a/apps/worker-cdr/src/analytics-main.ts b/apps/worker-cdr/src/analytics-main.ts index 1269070..fb03e15 100644 --- a/apps/worker-cdr/src/analytics-main.ts +++ b/apps/worker-cdr/src/analytics-main.ts @@ -54,7 +54,7 @@ async function entries(entries: Array<[string, string[]]>) { const event = JSON.parse(fields[1]) as AnalyticsEvent; await store.consume(event); await redis.xack(ANALYTICS_STREAM, ANALYTICS_GROUP, id); - dataThrough = new Date(event.at).toISOString(); + if (!dataThrough || event.at > Date.parse(dataThrough)) dataThrough = new Date(event.at).toISOString(); } } async function observeActive() { @@ -92,7 +92,7 @@ async function main() { if (Date.now() - lastAlerts > 15000 && !degraded) { await evaluateCallerAlerts(db); lastAlerts = Date.now(); } } catch (e) { degraded = true; reason = e instanceof Error ? e.message.slice(0, 200) : 'Worker error'; logger.error({ reason }, 'Analytics cycle failed'); } try { - const payload = JSON.stringify({ degraded, reason, dataThrough, lagMs: degraded ? null : 0, captureSince: process.env.ANALYTICS_CAPTURE_SINCE ?? null, spoolOffset: cursor.offset, gap, updatedAt: new Date().toISOString() }); + const payload = JSON.stringify({ degraded, reason, dataThrough, lagMs: null, coverageStatus: 'OBSERVATIONS_ONLY', captureSince: process.env.ANALYTICS_CAPTURE_SINCE ?? null, spoolOffset: cursor.offset, gap, updatedAt: new Date().toISOString() }); await db.$executeRaw(Prisma.sql`INSERT INTO caller_analysis_health(id,payload,updated_at) VALUES ('worker',${payload},NOW(3)) ON DUPLICATE KEY UPDATE payload=VALUES(payload),updated_at=VALUES(updated_at)`); } catch { logger.error('Analytics heartbeat write failed'); } await new Promise(resolve => setTimeout(resolve, 700)); diff --git a/packages/database/src/caller-analytics-model.ts b/packages/database/src/caller-analytics-model.ts index f14a9f2..504eebe 100644 --- a/packages/database/src/caller-analytics-model.ts +++ b/packages/database/src/caller-analytics-model.ts @@ -15,6 +15,7 @@ export interface AnalyticsLeg { event: AnalyticsEvent; first180At: number | null; first183At: number | null; acceptedAt: number | null; endedAt: number | null; finalCode: number; talkMs: number; durationAt: number; durationFinal: boolean; unknown: boolean; + coverageGap?: boolean; durationAuthority?: number; } export interface AnalyticsCall { start: AnalyticsEvent; legs: Record; } export interface AnalyticsCounts { @@ -67,26 +68,28 @@ export function foldAnalytics(call: AnalyticsCall | null, event: AnalyticsEvent) event, first180At: null, first183At: null, acceptedAt: null, endedAt: null, finalCode: 0, talkMs: 0, durationAt: 0, durationFinal: false, unknown: event.kind !== 'ATTEMPT' }; - if (event.kind === 'ATTEMPT') { leg.event = event; leg.unknown = false; } + if (event.kind === 'ATTEMPT') { leg.event = event; if (!leg.coverageGap) leg.unknown = false; } if (event.kind === 'PROGRESS' && event.code === 180) leg.first180At = earliest(leg.first180At, event.at); if (event.kind === 'PROGRESS' && event.code === 183) leg.first183At = earliest(leg.first183At, event.at); if (event.kind === 'ACCEPTED' && event.code >= 200 && event.code < 300) leg.acceptedAt = earliest(leg.acceptedAt, event.at); - if (event.kind === 'UNKNOWN') leg.unknown = true; + if (event.kind === 'UNKNOWN') { leg.unknown = true; leg.coverageGap = true; } if (event.kind === 'END' || event.kind === 'RECONCILE') { if (leg.endedAt === null || event.at >= leg.endedAt) { leg.endedAt = event.at; leg.finalCode = event.code; } // END carries the producer's accumulated phase flags through separate replayable events. - if (event.source === 'verified-dialog-final' || event.source === 'platform-final') leg.unknown = false; - if (event.source === 'verified-dialog-final' && leg.acceptedAt !== null && event.talkMs === null) { - leg.talkMs = Math.max(0, event.at - leg.acceptedAt); leg.durationAt = event.at; leg.durationFinal = true; + if (event.source === 'platform-final' && !event.attemptId && !leg.coverageGap) leg.unknown = false; + if (event.source === 'verified-dialog-final' && leg.acceptedAt !== null && event.talkMs === null && !leg.durationFinal) { + leg.talkMs = Math.max(0, event.at - leg.acceptedAt); leg.durationAt = event.at; leg.durationFinal = true; leg.durationAuthority = 1; } } if (event.kind === 'ACCEPTED' && leg.endedAt !== null && !leg.durationFinal) { - leg.talkMs = Math.max(0, leg.endedAt - event.at); leg.durationAt = leg.endedAt; leg.durationFinal = true; + leg.talkMs = Math.max(0, leg.endedAt - event.at); leg.durationAt = leg.endedAt; leg.durationFinal = true; leg.durationAuthority = 1; } if (event.talkMs !== null && ['DURATION', 'END', 'RECONCILE'].includes(event.kind)) { const final = event.kind !== 'DURATION'; - if ((!leg.durationFinal || final) && event.at >= leg.durationAt) { - leg.talkMs = event.talkMs; leg.durationAt = event.at; leg.durationFinal = final; + const authority = event.kind === 'RECONCILE' ? 3 : final ? 2 : 0; + const priorAuthority = leg.durationAuthority ?? (leg.durationFinal ? 1 : 0); + if (authority > priorAuthority || (authority === priorAuthority && event.at >= leg.durationAt)) { + leg.talkMs = event.talkMs; leg.durationAt = event.at; leg.durationFinal = final; leg.durationAuthority = authority; } } result.legs[key] = leg; diff --git a/packages/database/src/caller-analytics-store.ts b/packages/database/src/caller-analytics-store.ts index 9780a94..a27a15a 100644 --- a/packages/database/src/caller-analytics-store.ts +++ b/packages/database/src/caller-analytics-store.ts @@ -43,7 +43,7 @@ export class CallerAnalyticsStore { const minute = new Date(Math.floor(row.startedAt / 60000) * 60000); const id = analyticsHash([minute.toISOString(), row.view, row.customerId, row.customerGatewayId, row.vendorId, row.vendorGatewayId, row.caller, row.city, row.carrier]); const delta = JSON.stringify(Object.fromEntries(COUNT_KEYS.map(k => [k, sign * row[k]]))); - const additions = Prisma.join(COUNT_KEYS.flatMap(k => [Prisma.sql`${`$.${k}`}`, Prisma.sql`CAST(JSON_UNQUOTE(JSON_EXTRACT(counts, ${`$.${k}`})) AS SIGNED) + ${sign * row[k]}`])); + const additions = Prisma.join(COUNT_KEYS.flatMap(k => [Prisma.sql`${k}`, Prisma.sql`CAST(JSON_UNQUOTE(JSON_EXTRACT(counts, ${`$.${k}`})) AS SIGNED) + ${sign * row[k]}`])); await tx.$executeRaw(Prisma.sql`INSERT INTO caller_analysis_minute_buckets (id, started_at, view, customer_id, customer_gateway_id, vendor_id, vendor_gateway_id, caller, city, carrier, counts) VALUES (${id}, ${minute}, ${row.view}, ${row.customerId}, ${row.customerGatewayId}, ${row.vendorId}, ${row.vendorGatewayId}, ${row.caller}, ${row.city}, ${row.carrier}, ${delta}) diff --git a/packages/database/src/caller-analytics.spec.ts b/packages/database/src/caller-analytics.spec.ts index 507ffe6..264cdd8 100644 --- a/packages/database/src/caller-analytics.spec.ts +++ b/packages/database/src/caller-analytics.spec.ts @@ -46,4 +46,17 @@ describe('caller analytics business state', () => { it('calculates the three different business rates', () => { expect(analyticsRates({ totalCalls:100, connectedCalls:60, answeredCalls:30, failedCalls:25, pendingCalls:15, unknownCalls:0, activeCalls:15, talkMs:0 })).toMatchObject({ notConnectedCalls:40, connectionRate:60, overallAnswerRate:30, connectedAnswerRate:50 }); }); + it('does not erase missing progress evidence with an end event', () => { + expect(run([event('END',{code:487,source:'verified-dialog-final'})])[0]).toMatchObject({unknownCalls:1,failedCalls:0}); + for (const sequence of [[event('UNKNOWN'),event('END')],[event('END'),event('UNKNOWN')]]) { + expect(run([event('ATTEMPT'),...sequence,event('ATTEMPT')])[0]).toMatchObject({unknownCalls:1,failedCalls:0}); + } + }); + it('gives verified final zero priority over a newer provisional observation and delayed END', () => { + const provisional=event('DURATION',{at:start+4000,talkMs:3000}); + const final=event('RECONCILE',{at:start+3000,talkMs:0,source:'verified-dialog-final'}); + for (const sequence of [[provisional,final],[final,provisional]]) { + expect(run([event('ATTEMPT'),event('ACCEPTED',{code:200}),...sequence,event('END',{at:start+5000,source:'verified-dialog-final'})])[0]).toMatchObject({connectedCalls:1,answeredCalls:0,talkMs:0}); + } + }); }); diff --git a/scripts/rebuild-caller-analytics.mjs b/scripts/rebuild-caller-analytics.mjs new file mode 100644 index 0000000..cf8bcc3 --- /dev/null +++ b/scripts/rebuild-caller-analytics.mjs @@ -0,0 +1,25 @@ +// Maintenance only: stop caller-analytics worker first; does not touch calls or billing. +import { PrismaClient } from '@prisma/client'; +import { analyticsHash, COUNT_KEYS } from '../packages/database/dist/index.js'; +import assert from 'node:assert/strict'; +assert(process.argv.includes('--worker-stopped'), 'Stop the analytics worker before rebuilding derived buckets'); +const db = new PrismaClient(); +try { + await db.$transaction(async tx => { + const rows = await tx.$queryRaw`SELECT * FROM caller_analysis_states`; + assert(rows.length <= 100000, 'Use a batched maintenance job for larger data sets'); + const buckets = new Map(); + for (const row of rows) { + const minute = new Date(Math.floor(+row.started_at / 60000) * 60000); + const id = analyticsHash([minute.toISOString(), row.view, row.customer_id, row.customer_gateway_id, row.vendor_id, row.vendor_gateway_id, row.caller, row.city, row.carrier]); + const bucket = buckets.get(id) || {...row, id, started_at:minute, counts:Object.fromEntries(COUNT_KEYS.map(k=>[k,0]))}; + for (const key of COUNT_KEYS) bucket.counts[key] += Number(row.counts[key]); + buckets.set(id,bucket); + } + await tx.$executeRaw`DELETE FROM caller_analysis_minute_buckets`; + for (const r of buckets.values()) await tx.$executeRaw`INSERT INTO caller_analysis_minute_buckets + (id,started_at,view,customer_id,customer_gateway_id,vendor_id,vendor_gateway_id,caller,city,carrier,counts) + VALUES (${r.id},${r.started_at},${r.view},${r.customer_id},${r.customer_gateway_id},${r.vendor_id},${r.vendor_gateway_id},${r.caller},${r.city},${r.carrier},${JSON.stringify(r.counts)})`; + console.log(JSON.stringify({states:rows.length,buckets:buckets.size,rebuilt:true})); + },{timeout:60000}); +} finally {await db.$disconnect();} diff --git a/tests/api/caller-analytics-real.mjs b/tests/api/caller-analytics-real.mjs index c4709e2..7148d0c 100644 --- a/tests/api/caller-analytics-real.mjs +++ b/tests/api/caller-analytics-real.mjs @@ -4,13 +4,14 @@ import { signAccessToken } from '../../packages/auth/dist/index.js'; import { spawnSync } from 'node:child_process'; import { writeFile } from 'node:fs/promises'; import assert from 'node:assert/strict'; +import { CallerAnalyticsStore, COUNT_KEYS } from '../../packages/database/dist/index.js'; const db = new PrismaClient(); const fixture = 'cra_20260831'; const user = await db.user.findFirst({ where:{ status:'ENABLED', userRoles:{some:{roleId:'ROLE_SUPER_ADMIN'}} } }); assert(user, 'An existing test administrator is required'); const token = signAccessToken({sub:user.id,username:user.username,roles:['ROLE_SUPER_ADMIN'],typ:'access'}, {secret:process.env.AUTH_ACCESS_TOKEN_SECRET,issuer:process.env.AUTH_TOKEN_ISSUER||'lisglosips-api',audience:process.env.AUTH_TOKEN_AUDIENCE||'lisglosips-web',ttlSeconds:900}); -const get = async (path, authenticated = true) => { - const response = await fetch(`http://127.0.0.1:3000/api/v2${path}`, {headers:authenticated ? {Authorization:`Bearer ${token}`} : {}}); +const get = async (path, authenticated = true, authToken = token) => { + const response = await fetch(`http://127.0.0.1:3000/api/v2${path}`, {headers:authenticated ? {Authorization:`Bearer ${authToken}`} : {}}); return {status:response.status,body:await response.json()}; }; try { @@ -51,7 +52,32 @@ try { assert.equal(detail.status,200); assert.equal(detail.body.rows.length,6); const short = detail.body.rows.find(r=>r.callee==='99105'); assert(short.counts.talkMs>0 && short.counts.talkMs<1000,'Subsecond duration must count as answer'); - const evidence={date:new Date().toISOString(),from,api:result,detail:detail.body,sip:JSON.parse(sip.stdout),tests:['real SIP six scenarios','real API counts and three rates','real DB states','unauthenticated 401','invalid query 400','positive subsecond duration']}; + const auditFrom = new Date(Date.now()-86400000).toISOString(); + const auditPath = `/caller-analytics/overview?customerId=${fixture}&view=landing&from=${encodeURIComponent(auditFrom)}`; + const beforeReplay = await get(auditPath); assert.equal(beforeReplay.status,200); + const stateRows = await db.$queryRaw`SELECT counts FROM caller_analysis_states WHERE customer_id=${fixture} AND view='landing' AND started_at>=${new Date(auditFrom)}`; + for (const key of COUNT_KEYS) assert.equal(beforeReplay.body.summary[key],stateRows.reduce((sum,row)=>sum+row.counts[key],0),`Minute bucket reconciliation: ${key}`); + const events = await db.$queryRaw`SELECT payload FROM caller_analysis_events WHERE call_key=${detail.body.rows[0].callKey}`; + const store = new CallerAnalyticsStore(db); + for (const event of events) assert.equal(await store.consume(event.payload),'duplicate'); + assert.deepEqual((await get(auditPath)).body.summary,beforeReplay.body.summary,'Replay must not change counts'); + const testUserId = `cra_rbac_${Date.now()}`; + try { + await db.role.create({data:{id:testUserId,name:testUserId,permissions:{create:{permissionId:'caller_analytics.view'}}}}); + await db.user.create({data:{id:testUserId,username:testUserId,displayName:'临时主叫范围验收',userRoles:{create:{roleId:testUserId}}}}); + const scopedToken=signAccessToken({sub:testUserId,username:testUserId,roles:[],typ:'access'}, {secret:process.env.AUTH_ACCESS_TOKEN_SECRET,issuer:process.env.AUTH_TOKEN_ISSUER||'lisglosips-api',audience:process.env.AUTH_TOKEN_AUDIENCE||'lisglosips-web',ttlSeconds:60}); + assert.equal((await get(auditPath,true,scopedToken)).status,403,'No assigned customers must fail closed'); + await db.$executeRaw`INSERT INTO caller_analysis_access(user_id,customer_id) VALUES (${testUserId},${fixture})`; + const scoped = await get('/caller-analytics/overview?minutes=1440',true,scopedToken); + assert.equal(scoped.status,200); assert(scoped.body.numbers.every(r=>r.customerId===fixture)); + assert.equal((await get('/caller-analytics/overview?customerId=another_customer',true,scopedToken)).status,403); + const options=await get('/caller-analytics/options',true,scopedToken); + assert.equal(options.status,200); assert.deepEqual(options.body.customers.map(c=>c.id),[fixture]); + } finally { + await db.$executeRaw`DELETE FROM caller_analysis_access WHERE user_id=${testUserId}`; + await db.user.deleteMany({where:{id:testUserId}}); await db.role.deleteMany({where:{id:testUserId}}); + } + const evidence={date:new Date().toISOString(),from,api:result,detail:detail.body,sip:JSON.parse(sip.stdout),bucketAudit:beforeReplay.body.summary,replayed:events.length,tests:['real SIP six scenarios','real API counts and three rates','real DB states and minute buckets','database replay idempotency','customer scope and 403 isolation','scoped filter options','unauthenticated 401','invalid query 400','positive subsecond duration']}; await writeFile(process.env.CRA_REPORT || '/tmp/caller-analytics-real-evidence.json',JSON.stringify(evidence,null,2)); console.log(JSON.stringify({passed:true,summary:result.body.summary,shortMs:short.counts.talkMs,quality:result.body.qualityStatus})); }