fix: reconcile caller aggregates and strengthen evidence handling

This commit is contained in:
hectorzhao
2026-08-31 17:52:11 +08:00
parent d78616561e
commit 2a3dcb0c68
8 changed files with 85 additions and 17 deletions
@@ -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),
+3 -2
View File
@@ -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) => <Field label={label}><Select value={draft[key]} onChange={e => change(key, e.target.value)}><option value="">全部</option>{choices.map(c => <option key={c.id} value={c.id}>{c.name}</option>)}</Select></Field>;
@@ -75,7 +76,7 @@ export function CallerAnalyticsPage() {
</form></Panel>
<section className="metric-grid cra-metrics">{[...countColumns, ...rateColumns].map(c => <div key={c.key} className="metric-card"><span>{c.label}</span><strong>{metrics ? c.render ? c.render(metrics) : metrics[c.key] : '—'}</strong></div>)}</section>
<p className="muted-text" role="status">{data ? `${formatDateTime(data.from)}${formatDateTime(data.to)} · ${data.qualityStatus === 'LIVE' ? '采集运行中' : '数据滞后或不完整'} · 待接通 ${metrics.pendingCalls} · 已结束未接通 ${metrics.failedCalls} · 未知 ${metrics.unknownCalls} · 更新于 ${formatDateTime(data.asOf)}` : '正在读取真实统计数据…'}</p>
{data ? <Alert title="口径与数据范围" tone={data.qualityStatus === 'LIVE' ? 'info' : 'warning'}>{data.historyNotice} {data.summaryScope}183 不代表实际振铃未接通数含待接通进行中指标为暂定值</Alert> : null}
{data ? <Alert title="口径与数据范围" tone={data.qualityStatus === 'LIVE' ? 'info' : 'warning'}>{data.historyNotice} {data.summaryScope}183 不代表实际振铃未接通数含待接通进行中指标为暂定值采集运行状态不代表窗口数据完整性已核验端到端延迟尚未测定</Alert> : null}
<Panel title="分钟趋势" aside={<span className="muted-text">呼叫 / 接通 / 绿应答</span>}>
<div className="cra-chart" aria-label="分钟呼叫接通应答趋势">{timeline.map((r,i) => <div key={r.time ?? i} className="cra-chart-column" title={`${formatDateTime(Number(r.time))} 呼叫${r.totalCalls} 接通${r.connectedCalls} 应答${r.answeredCalls}`}><div className="cra-bars"><span style={{height:`${r.totalCalls/max*100}%`}} /><span style={{height:`${r.connectedCalls/max*100}%`}} /><span style={{height:`${r.answeredCalls/max*100}%`}} /></div><small>{new Date(Number(r.time)).toLocaleTimeString('zh-CN',{hour:'2-digit',minute:'2-digit'})}</small></div>)}</div>
{!timeline.length ? <p className="muted-text">当前窗口暂无已采集呼叫</p> : <details><summary>查看趋势数值与三个率</summary><SimpleTable rows={timeline.map((r,i) => ({ ...r, id:i, minute:formatDateTime(Number(r.time)) }))} columns={[{key:'minute',label:'分钟'},...countColumns,...rateColumns]} /></details>}
+2 -2
View File
@@ -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));
@@ -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<string, AnalyticsLeg>; }
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;
@@ -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})
@@ -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});
}
});
});
+25
View File
@@ -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();}
+29 -3
View File
@@ -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}));
}