feat: add caller connection and answer rate analytics
This commit is contained in:
@@ -0,0 +1,36 @@
|
||||
import { PrismaClient, Prisma, analyticsHash, analyticsRates, COUNT_KEYS, type AnalyticsCounts } from '@lisglosips/database';
|
||||
|
||||
export async function evaluateCallerAlerts(db: PrismaClient) {
|
||||
const end = new Date(Math.floor(Date.now() / 60000) * 60000); const start = new Date(+end - 15 * 60000);
|
||||
const aggregates = Prisma.join(COUNT_KEYS.map(k => Prisma.sql`SUM(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,${`$.${k}`})) AS SIGNED)) AS ${Prisma.raw(k)}`));
|
||||
const rows = await db.$queryRaw<Array<AnalyticsCounts & { customer_id: string; caller: string; view: string }>>(Prisma.sql`SELECT customer_id, caller, view, ${aggregates} FROM caller_analysis_minute_buckets WHERE started_at>=${start} AND started_at<${end} GROUP BY customer_id, caller, view HAVING SUM(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,'$.totalCalls')) AS SIGNED))>=50`);
|
||||
for (const row of rows) {
|
||||
const counts = Object.fromEntries(COUNT_KEYS.map(k => [k, Number(row[k])])) as unknown as AnalyticsCounts;
|
||||
const rates = analyticsRates(counts);
|
||||
if (rates.unknownCalls || rates.activeCalls / rates.totalCalls > 0.2) continue;
|
||||
const reasons = [
|
||||
rates.connectionRate !== null && rates.connectionRate < Number(process.env.ANALYTICS_LOW_CONNECTION_RATE ?? 20) ? '低接通率' : '',
|
||||
rates.overallAnswerRate !== null && rates.overallAnswerRate < Number(process.env.ANALYTICS_LOW_OVERALL_RATE ?? 10) ? '低总体应答率' : '',
|
||||
rates.connectedCalls >= 30 && rates.connectedAnswerRate !== null && rates.connectedAnswerRate < Number(process.env.ANALYTICS_LOW_CONNECTED_RATE ?? 20) ? '低已接通应答率' : ''
|
||||
].filter(Boolean);
|
||||
const previous = await db.$queryRaw<Array<AnalyticsCounts>>(Prisma.sql`SELECT ${aggregates} FROM caller_analysis_minute_buckets WHERE started_at>=${new Date(+start - 15 * 60000)} AND started_at<${start} AND customer_id=${row.customer_id} AND caller=${row.caller} AND view=${row.view}`);
|
||||
const p = analyticsRates(Object.fromEntries(COUNT_KEYS.map(k => [k, Number(previous[0]?.[k] ?? 0)])) as unknown as AnalyticsCounts);
|
||||
if (p.totalCalls >= 50 && !p.unknownCalls && p.activeCalls / p.totalCalls <= 0.2) {
|
||||
for (const [key, label] of [['connectionRate','接通率骤降'],['overallAnswerRate','总体应答率骤降'],['connectedAnswerRate','已接通应答率骤降']] as const) {
|
||||
if (key === 'connectedAnswerRate' && (p.connectedCalls < 30 || rates.connectedCalls < 30)) continue;
|
||||
if (p[key] !== null && rates[key] !== null && p[key]! - rates[key]! > Number(process.env.ANALYTICS_RATE_DROP_POINTS ?? 15)) reasons.push(label);
|
||||
}
|
||||
}
|
||||
const tail = await db.$queryRaw<Array<{ counts: AnalyticsCounts }>>`SELECT counts FROM caller_analysis_states WHERE customer_id=${row.customer_id} AND caller=${row.caller} AND view=${row.view} AND started_at>=${start} AND started_at<${end} ORDER BY started_at DESC, id DESC LIMIT 20`;
|
||||
if (tail.length === 20 && tail.every(t => t.counts.failedCalls === 1)) reasons.push('连续20次已结束未接通');
|
||||
const id = analyticsHash([row.customer_id, row.caller, row.view]);
|
||||
const existing = await db.$queryRaw<Array<{ status: string; payload: { streak?: number; recovered?: number; since?: string } }>>`SELECT status,payload FROM caller_analysis_alerts WHERE id=${id}`;
|
||||
const streak = reasons.length ? (existing[0]?.payload.streak ?? 0) + 1 : 0;
|
||||
const recovered = reasons.length ? 0 : (existing[0]?.payload.recovered ?? 0) + 1;
|
||||
const status = streak >= 3 ? 'ACTIVE' : existing[0]?.status === 'ACTIVE' && recovered < 3 ? 'ACTIVE' : 'OBSERVING';
|
||||
const payload = JSON.stringify({ reasons, rates, streak, recovered, windowStart: start, windowEnd: end, since: existing[0]?.payload.since ?? new Date().toISOString() });
|
||||
await db.$executeRaw`INSERT INTO caller_analysis_alerts(id,customer_id,view,caller,status,payload,updated_at) VALUES (${id},${row.customer_id},${row.view},${row.caller},${status},${payload},NOW(3)) ON DUPLICATE KEY UPDATE status=VALUES(status),payload=VALUES(payload),updated_at=VALUES(updated_at)`;
|
||||
}
|
||||
// Old low-sample windows are no longer active evidence, not permanent alarms.
|
||||
await db.$executeRaw`UPDATE caller_analysis_alerts SET status='EXPIRED' WHERE updated_at<${new Date(Date.now() - 120000)} AND status='ACTIVE'`;
|
||||
}
|
||||
@@ -0,0 +1,102 @@
|
||||
import { open, readFile, writeFile, rename, stat, appendFile } from 'node:fs/promises';
|
||||
import { hostname } from 'node:os';
|
||||
import { PrismaClient, Prisma, CallerAnalyticsStore, ANALYTICS_GROUP, ANALYTICS_STREAM, analyticsHash, parseAnalyticsEvent, type AnalyticsCall, type AnalyticsEvent } from '@lisglosips/database';
|
||||
import { createRedisClient } from '@lisglosips/redis';
|
||||
import { createLogger } from '@lisglosips/observability';
|
||||
import { evaluateCallerAlerts } from './analytics-alerts.js';
|
||||
|
||||
const logger = createLogger('caller-analytics');
|
||||
const db = new PrismaClient(); const redis = createRedisClient(process.env.REDIS_URL ?? '');
|
||||
const store = new CallerAnalyticsStore(db);
|
||||
const spool = process.env.ANALYTICS_SPOOL ?? '/var/log/lisglosips/caller-analytics.log';
|
||||
const checkpoint = process.env.ANALYTICS_CHECKPOINT ?? '/var/lib/lisglosips/caller-analytics-offset.json';
|
||||
let cursor = { offset: 0, ino: 0 }; let stopping = false; let lastAlerts = 0;
|
||||
let dataThrough: string | null = null; let gap = false;
|
||||
const consumer = `${hostname()}-${process.pid}`;
|
||||
process.on('SIGTERM', () => { stopping = true; }); process.on('SIGINT', () => { stopping = true; });
|
||||
|
||||
function decodeLine(line: string): AnalyticsEvent | null {
|
||||
const marker = line.indexOf('CRA1|'); if (marker < 0) return null;
|
||||
const parts = line.slice(marker).trim().split('|');
|
||||
if (parts.length !== 20) throw new Error('Malformed producer record');
|
||||
const decoded = parts.slice(8).map(s => Buffer.from(s, 'base64').toString('utf8'));
|
||||
const [uid, callId, customerId, customerGatewayId, vendorId, vendorGatewayId, caller, landingCaller, callee, city, carrier, attemptId] = decoded;
|
||||
const at = Number(parts[2]) * 1000 + Math.floor(Number(parts[3]) / 1000);
|
||||
const startedAt = Number(parts[4]) * 1000 + Math.floor(Number(parts[5]) / 1000);
|
||||
const fields = { version: '1', eventId: analyticsHash(parts), callUid: analyticsHash(uid), callId, customerId, customerGatewayId, vendorId, vendorGatewayId, caller, landingCaller, callee, city, carrier, attemptId,
|
||||
kind: parts[1], at: String(at), startedAt: String(startedAt), code: parts[6], source: parts[7], talkMs: '-1', sequence: '0' };
|
||||
return parseAnalyticsEvent(Object.entries(fields).flat());
|
||||
}
|
||||
async function collect() {
|
||||
const info = await stat(spool);
|
||||
if (cursor.ino && (cursor.ino !== info.ino || info.size < cursor.offset)) { gap = true; cursor.offset = 0; }
|
||||
cursor.ino = info.ino;
|
||||
const file = await open(spool, 'r');
|
||||
try {
|
||||
const buffer = Buffer.alloc(Math.min(1024 * 1024, Math.max(0, info.size - cursor.offset)));
|
||||
const { bytesRead } = await file.read(buffer, 0, buffer.length, cursor.offset);
|
||||
const last = buffer.subarray(0, bytesRead).lastIndexOf(10); if (last < 0) return;
|
||||
for (const line of buffer.subarray(0, last + 1).toString('utf8').split('\n').filter(Boolean)) {
|
||||
let event: AnalyticsEvent | null;
|
||||
try { event = decodeLine(line); } catch {
|
||||
await appendFile(`${checkpoint}.deadletter`, `${JSON.stringify({ at:new Date(), line })}\n`, { mode:0o600 });
|
||||
gap = true; continue;
|
||||
}
|
||||
if (event) await redis.xadd(ANALYTICS_STREAM, '*', 'event', JSON.stringify(event));
|
||||
}
|
||||
// Persist only after durable queue writes. Crash before checkpoint replays harmlessly.
|
||||
cursor.offset += last + 1;
|
||||
await writeFile(`${checkpoint}.tmp`, JSON.stringify(cursor), { mode: 0o600 }); await rename(`${checkpoint}.tmp`, checkpoint);
|
||||
} finally { await file.close(); }
|
||||
}
|
||||
async function entries(entries: Array<[string, string[]]>) {
|
||||
for (const [id, fields] of entries) {
|
||||
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();
|
||||
}
|
||||
}
|
||||
async function observeActive() {
|
||||
const response = await fetch(process.env.ANALYTICS_MI_URL ?? 'http://127.0.0.1:8888/mi', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'dlg_list', params: [] }), signal: AbortSignal.timeout(2000) });
|
||||
const body = await response.json() as { result?: { Dialogs?: Array<{ callid: string; state: number }> } };
|
||||
if (!response.ok || !body.result?.Dialogs) throw new Error('Dialog observation unavailable');
|
||||
const calls = await db.$queryRaw<Array<{ payload: AnalyticsCall }>>`SELECT c.payload FROM caller_analysis_calls c WHERE EXISTS (SELECT 1 FROM caller_analysis_states s WHERE s.call_key=c.id AND s.view='original' AND CAST(JSON_UNQUOTE(JSON_EXTRACT(s.counts,'$.activeCalls')) AS SIGNED)>0) LIMIT 2000`;
|
||||
const live = new Map(body.result.Dialogs.map(d => [d.callid, Number(d.state)]));
|
||||
const now = Date.now();
|
||||
for (const { payload } of calls) for (const leg of Object.values(payload.legs)) {
|
||||
if (leg.endedAt !== null || leg.acceptedAt === null || leg.talkMs > 0) continue;
|
||||
// Only a confirmed live dialog supplies provisional positive duration. Never a browser timer.
|
||||
if ((live.get(leg.event.callId) ?? 0) >= 3 && (live.get(leg.event.callId) ?? 0) <= 4 && now > leg.acceptedAt) {
|
||||
await store.consume({ ...leg.event, kind: 'DURATION', at: now, talkMs: now - leg.acceptedAt, source: 'dialog-observation', eventId: analyticsHash([leg.event.callUid, leg.event.attemptId, 'positive-duration']) });
|
||||
}
|
||||
}
|
||||
}
|
||||
async function main() {
|
||||
try { cursor = JSON.parse(await readFile(checkpoint, 'utf8')); } catch { /* First start reads the retained spool. */ }
|
||||
await redis.connect(); await db.$connect();
|
||||
try { await redis.xgroup('CREATE', ANALYTICS_STREAM, ANALYTICS_GROUP, '0', 'MKSTREAM'); } catch (e) { if (!(e instanceof Error) || !e.message.includes('BUSYGROUP')) throw e; }
|
||||
let pendingCursor = '0-0';
|
||||
while (!stopping) {
|
||||
let degraded = gap; let reason = '';
|
||||
try {
|
||||
await collect();
|
||||
const claimed = await redis.xautoclaim(ANALYTICS_STREAM, ANALYTICS_GROUP, consumer, 5000, pendingCursor, 'COUNT', 100) as [string, Array<[string, string[]]>];
|
||||
pendingCursor = claimed[0]; await entries(claimed[1]);
|
||||
const fresh = await redis.xreadgroup('GROUP', ANALYTICS_GROUP, consumer, 'COUNT', 200, 'BLOCK', 100, 'STREAMS', ANALYTICS_STREAM, '>') as Array<[string, Array<[string, string[]]>]> | null;
|
||||
if (fresh) for (const [, batch] of fresh) await entries(batch);
|
||||
await observeActive();
|
||||
const groups = await redis.xinfo('GROUPS', ANALYTICS_STREAM) as Array<Array<string | number>>;
|
||||
const group = groups.map(g => Object.fromEntries(Array.from({ length: g.length / 2 }, (_, i) => [g[i * 2], g[i * 2 + 1]]))).find(g => g.name === ANALYTICS_GROUP);
|
||||
if (Number(group?.pending ?? 0) || Number(group?.lag ?? 0)) { degraded = true; reason = 'QUEUE_BACKLOG'; }
|
||||
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() });
|
||||
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));
|
||||
}
|
||||
redis.disconnect(); await db.$disconnect();
|
||||
}
|
||||
void main().catch(e => { logger.error({ message: e instanceof Error ? e.message : 'Startup failed' }, 'Analytics failed'); process.exitCode = 1; redis.disconnect(); void db.$disconnect(); });
|
||||
Reference in New Issue
Block a user