fix: complete caller analytics reliability remediation

This commit is contained in:
hectorzhao
2026-09-01 10:03:29 +08:00
parent a116926867
commit 7f8825cf3a
25 changed files with 850 additions and 202 deletions
+19 -35
View File
@@ -1,36 +1,20 @@
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'`;
import { COUNT_KEYS, Prisma, PrismaClient, analyticsHash, analyticsRates, type AnalyticsCounts } from '@lisglosips/database';
type Row={customer_id:string;caller:string;view:string}&AnalyticsCounts;
type Prior={status:string;payload:{streak?:number;recovered?:number;since?:string;evaluationKey?:string}};
const sums=()=>Prisma.join(COUNT_KEYS.map(k=>Prisma.sql`SUM(COALESCE(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,${`$.${k}`})) AS SIGNED),0)) AS ${Prisma.raw(k)}`));
const numeric=(row:Partial<Row>):AnalyticsCounts=>Object.fromEntries(COUNT_KEYS.map(k=>[k,Number(row[k]??0)]))as unknown as AnalyticsCounts;
export async function suspendCallerAlerts(db:PrismaClient,reason:string){await db.$executeRaw`UPDATE caller_analysis_alerts SET status='SUSPENDED',payload=JSON_SET(payload,'$.suspendedReason',${reason},'$.suspendedAt',${new Date().toISOString()},'$.streak',0,'$.recovered',0),updated_at=NOW(3) WHERE status IN ('ACTIVE','OBSERVING')`;}
export async function evaluateCallerAlerts(db:PrismaClient){
const end=new Date(Math.floor(Date.now()/60000)*60000),start=new Date(+end-15*60000),evaluationKey=end.toISOString(),aggregates=sums();
const rows=await db.$queryRaw<Row[]>(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(COALESCE(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,'$.totalCalls')) AS SIGNED),0))>=50`);
for(const row of rows){const counts=numeric(row),rates=analyticsRates(counts);if(rates.unknownCalls||rates.durationUnknownCalls||rates.activeCalls/rates.totalCalls>0.2)continue;const reasons=new Map<string,string>();
if(rates.connectionRate!==null&&rates.connectionRate<Number(process.env.ANALYTICS_LOW_CONNECTION_RATE??20))reasons.set('LOW_CONNECTION','低接通率');
if(rates.overallAnswerRate!==null&&rates.overallAnswerRate<Number(process.env.ANALYTICS_LOW_OVERALL_RATE??10))reasons.set('LOW_OVERALL','低总体应答率');
if(rates.connectedCalls>=30&&rates.connectedAnswerRate!==null&&rates.connectedAnswerRate<Number(process.env.ANALYTICS_LOW_CONNECTED_RATE??20))reasons.set('LOW_CONNECTED','低已接通应答率');
const previous=await db.$queryRaw<Row[]>(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}`),p=analyticsRates(numeric(previous[0]??{}));
if(p.totalCalls>=50&&!p.unknownCalls&&!p.durationUnknownCalls&&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.set(`DROP_${key}`,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.set('CONSECUTIVE_FAILURE','连续20次已结束未接通');
for(const ruleId of['LOW_CONNECTION','LOW_OVERALL','LOW_CONNECTED','DROP_connectionRate','DROP_overallAnswerRate','DROP_connectedAnswerRate','CONSECUTIVE_FAILURE']){const id=analyticsHash([row.customer_id,row.caller,row.view,ruleId,'v2']),existing=await db.$queryRaw<Prior[]>`SELECT status,payload FROM caller_analysis_alerts WHERE id=${id}`,old=existing[0];if(old?.payload.evaluationKey===evaluationKey)continue;const hit=reasons.has(ruleId),streak=hit?(old?.payload.streak??0)+1:0,recovered=hit?0:(old?.payload.recovered??0)+1,status=streak>=3?'ACTIVE':old?.status==='ACTIVE'&&recovered<3?'ACTIVE':'OBSERVING',payload=JSON.stringify({ruleId,reason:reasons.get(ruleId)??null,rates,streak,recovered,evaluationKey,windowStart:start,windowEnd:end,scope:{customerId:row.customer_id,caller:row.caller,view:row.view,minutes:15},since:old?.payload.since??new Date().toISOString(),cooldownSeconds:60});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)`;}
}
await db.$executeRaw`UPDATE caller_analysis_alerts SET status='STALE' WHERE updated_at<${new Date(Date.now()-120000)} AND status IN ('ACTIVE','OBSERVING','SUSPENDED')`;
}
+29 -92
View File
@@ -1,102 +1,39 @@
import { open, readFile, writeFile, rename, stat, appendFile } from 'node:fs/promises';
import { appendFile, open, readFile, rename, stat, writeFile } 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 { ANALYTICS_GROUP, ANALYTICS_STREAM, CallerAnalyticsStore, Prisma, PrismaClient, analyticsHash, parseAnalyticsEvent, validateAnalyticsEvent, type AnalyticsCall, type AnalyticsEvent } from '@lisglosips/database';
import { createRedisClient } from '@lisglosips/redis';
import { createLogger } from '@lisglosips/observability';
import { evaluateCallerAlerts } from './analytics-alerts.js';
import { evaluateCallerAlerts, suspendCallerAlerts } 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 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; });
interface CollectorCursor { offset:number; ino:number; gap?:boolean; eventsSinceProducerStart?:number; lastProducerCount?:number }
let cursor:CollectorCursor={offset:0,ino:0,eventsSinceProducerStart:0}; let stopping=false; let lastAlerts=0; let lastTrim=0; let lastAlertSuspend=0;
let dataThrough:string|null=null; let ingestLagMs:number|null=null; let gap=false; let bootId='unknown';
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' };
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 [uid,callId,customerId,customerGatewayId,vendorId,vendorGatewayId,caller,landingCaller,callee,city,carrier,attemptId]=parts.slice(8).map(s=>Buffer.from(s,'base64').toString('utf8'));
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 receivedAt=Date.now();
const fields={version:'2',eventId:analyticsHash(parts),callUid:analyticsHash([uid,customerId,hostname()]),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:String(at*1000),nodeId:hostname(),bootId,receivedAt:String(receivedAt),method:'INVITE',evidenceVersion:'2'};
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, gap }), { mode: 0o600 }); await rename(`${checkpoint}.tmp`, checkpoint);
} finally { await file.close(); }
async function recordGap(kind:string,detail:Record<string,unknown>){const id=analyticsHash([kind,detail,Date.now()]);await db.$executeRaw`INSERT INTO caller_analysis_gaps(id,node_id,kind,status,started_at,payload,created_at,updated_at) VALUES (${id},${hostname()},${kind},'OPEN',NOW(3),${JSON.stringify(detail)},NOW(3),NOW(3))`;}
async function persistCursor(){await writeFile(`${checkpoint}.tmp`,JSON.stringify({...cursor,gap}),{mode:0o600});await rename(`${checkpoint}.tmp`,checkpoint);}
async function collect(){
const info=await stat(spool); if(cursor.ino&&(cursor.ino!==info.ino||info.size<cursor.offset)){gap=true;await recordGap('SPOOL_ROTATED_OR_TRUNCATED',{previousIno:cursor.ino,currentIno:info.ino,previousOffset:cursor.offset,currentSize:info.size});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)){try{const event=decodeLine(line);if(event){await redis.xadd(ANALYTICS_STREAM,'*','event',JSON.stringify(event));cursor.eventsSinceProducerStart=(cursor.eventsSinceProducerStart??0)+1;}}catch(error){const message=error instanceof Error?error.message:'decode error';await appendFile(`${checkpoint}.deadletter`,`${JSON.stringify({at:new Date(),line,error:message})}\n`,{mode:0o600});gap=true;await recordGap('MALFORMED_SPOOL_RECORD',{offset:cursor.offset,error:message});}}
cursor.offset+=last+1;await persistCursor();}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);
if (!dataThrough || event.at > Date.parse(dataThrough)) 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 { const saved = JSON.parse(await readFile(checkpoint, 'utf8')); cursor = saved; gap = saved.gap === true; } 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: 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));
}
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(); });
async function consumeEntries(items:Array<[string,string[]]>){for(const[id,fields]of items){try{const i=fields.findIndex(v=>v==='event');if(i<0||!fields[i+1])throw new Error('Missing event field');const event=validateAnalyticsEvent(JSON.parse(fields[i+1]));await store.consume(event);ingestLagMs=Math.max(0,(event.receivedAt??Date.now())-event.at);if(!dataThrough||event.at>Date.parse(dataThrough))dataThrough=new Date(event.at).toISOString();}catch(error){const message=error instanceof Error?error.message.slice(0,500):'Invalid queued event';const deadId=analyticsHash([ANALYTICS_STREAM,id]);await db.$executeRaw`INSERT IGNORE INTO caller_analysis_dead_letters(id,stream_id,payload,error,status,created_at,updated_at) VALUES (${deadId},${id},${JSON.stringify(fields)},${message},'OPEN',NOW(3),NOW(3))`;gap=true;}await redis.xack(ANALYTICS_STREAM,ANALYTICS_GROUP,id);}}
async function producerCount():Promise<number>{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:'get_statistics',params:{statistics:['cra_event_total']}}),signal:AbortSignal.timeout(2000)});const body=await response.json()as{result?:Record<string,number>};const value=body.result?.['script:cra_event_total'];if(!response.ok||typeof value!=='number'||!Number.isSafeInteger(value))throw new Error('Producer counter unavailable');return value;}
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 live=new Map(body.result.Dialogs.map(d=>[d.callid,Number(d.state)]));let lastId='';for(;;){const calls=await db.$queryRaw<Array<{id:string;payload:AnalyticsCall}>>`SELECT c.id,c.payload FROM caller_analysis_calls c WHERE c.id>${lastId} AND 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) ORDER BY c.id LIMIT 500`;if(!calls.length)break;const now=Date.now();for(const{payload}of calls)for(const leg of Object.values(payload.legs)){if(leg.endedAt!==null||leg.talkStartedAt==null||leg.talkMs>0)continue;if((live.get(leg.event.callId)??0)===4&&now>leg.talkStartedAt)await store.consume({...leg.event,kind:'DURATION',at:now,talkMs:now-leg.talkStartedAt,source:'dialog-acked-observation',eventId:analyticsHash([leg.event.callUid,leg.event.attemptId,'positive-duration'])});}lastId=calls[calls.length-1].id;}}
async function main(){try{bootId=(await readFile('/proc/sys/kernel/random/boot_id','utf8')).trim();}catch{bootId=`${hostname()}-unknown`;}try{const saved=JSON.parse(await readFile(checkpoint,'utf8'))as CollectorCursor;cursor=saved;gap=saved.gap===true;}catch{/* first start */}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 produced=await producerCount();if(cursor.lastProducerCount!==undefined&&produced<cursor.lastProducerCount){gap=degraded=true;reason='PRODUCER_RESTART';await recordGap(reason,{previous:cursor.lastProducerCount,current:produced});cursor.eventsSinceProducerStart=0;}cursor.lastProducerCount=produced;await persistCursor();const claimed=await redis.xautoclaim(ANALYTICS_STREAM,ANALYTICS_GROUP,consumer,5000,pendingCursor,'COUNT',100)as[string,Array<[string,string[]]>];pendingCursor=claimed[0];await consumeEntries(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 consumeEntries(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(produced!==(cursor.eventsSinceProducerStart??0)){degraded=true;reason='PRODUCER_COLLECTOR_COUNT_MISMATCH';}if(!Number(group?.pending??0)&&Date.now()-lastTrim>60000){await redis.xtrim(ANALYTICS_STREAM,'MAXLEN','~',Number(process.env.ANALYTICS_STREAM_MAXLEN??1000000));lastTrim=Date.now();}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');}if(degraded&&Date.now()-lastAlertSuspend>15000){try{await suspendCallerAlerts(db,reason||'ANALYTICS_DEGRADED');lastAlertSuspend=Date.now();}catch{logger.error('Analytics alert suspension failed');}}try{const spoolSize=await stat(spool).then(s=>s.size).catch(()=>null);const coverageStatus=!degraded&&!gap&&cursor.lastProducerCount===cursor.eventsSinceProducerStart?'CONTINUOUS':'DEGRADED';const payload=JSON.stringify({degraded,reason,dataThrough,lagMs:ingestLagMs,coverageStatus,captureSince:process.env.ANALYTICS_CAPTURE_SINCE??null,spoolOffset:cursor.offset,spoolUnreadBytes:spoolSize===null?null:Math.max(0,spoolSize-cursor.offset),gap,producerEventCount:cursor.lastProducerCount??null,collectedEventCount:cursor.eventsSinceProducerStart??null,bootId,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();});