diff --git a/IMPLEMENTATION_STATUS.md b/IMPLEMENTATION_STATUS.md index ca61d4b..c7570ae 100644 --- a/IMPLEMENTATION_STATUS.md +++ b/IMPLEMENTATION_STATUS.md @@ -4,6 +4,23 @@ ## 当前状态 +### 2026-08-31 主叫号码实时分析实施(测试发布验证中) + +- 新增独立分析进程、MySQL事件/状态/分钟汇总/权限范围/告警表、真实API和主叫分析页面,计费Worker不变。 +- 新采集使用OpenSIPS进展及结束观察钩子、本机rsyslog持久队列/日志、Redis Stream,再由独立进程事务去重及聚合。 +- 本地校验:TypeScript、ESLint、Web构建通过;全量Vitest 29文件/114测试通过;Prisma schema通过;B候选OpenSIPS配置校验通过。测试部署/真实呼叫结果见后续发布记录,不能据此认定已上线。 +- 现场发现旧成功CDR固定duration_sec=6,不能作新应答统计权威来源。保留旧计费不变,新功能仅信任独立实际时间证据。T节点SSH暂不可达,真实SIP测试改为B本机独立零费率入口/落地,不拨外部号码。 +- Git中代码审查报告属于并行任务,不纳入本次提交。 + +### 2026-08-31 文档增补:主叫号码接通率与应答率实时分析 + +- 本轮任务:方案落文档并按用户要求修订接通/应答业务口径;文档已完成,功能实现待开始。 +- 产物:[专项设计V1.1](docs/CALLER_REALTIME_ANALYTICS_DESIGN.md),主设计文档索引及测试计划CRA-001至CRA-024验收索引。 +- 口径:180/183或后续阶段计接通;正实际通话时长计应答;总体应答率=A/T,已接通应答率=A/C。待接通与未知单列说明,不等同最终失败。 +- 验证:本地源码及RFC 3261只读核对;文档公式、场景、链接和Git差异检查。24项功能用例仅为验收设计,均未执行。 +- 变更边界:仅Markdown文档;无代码、迁移、生产变更、真实外呼或Git提交/推送。下方原有发布状态保留为历史记录,未在本轮重新验证。 +- 遗留/下一步:实施前核验生产呼叫过程、实际通话时长来源与精度、重试机制及容量,按专项方案实施并真实联调;历史缺少180/183的数据不得伪造完整接通统计。 + - 当前任务:`S56 - 全部当前服务收敛到 Server B` - 本次变更摘要:确认 B 已同时运行 OpenSIPS、RTPEngine、录音守护进程、Web/API/Worker、MySQL、Redis 与监控;将 Recording Worker 从 B 到 A 的 SSH 拉取改为 B 本机 ready 目录读取,安装本机 finalizer timer,并通过 systemd BindPaths/附加组暴露最小目录权限。当前 release 为 `s56-b-consolidation-20260827T095812Z`,提交为 `309bffc7ffa70a760bce21ae59cb6c457c9c57fe`。 - 总体状态:S30 已完成;本地 KVM 开发环境 V2 闭环原冻结 release 为 `s28-v2-20260621220924`;B 二期冻结基线为 `s42-phase2-business-prefix-20260624140000`,当前 release 为 `s53-brand-logo-sip-text-20260630105500`;稳定版名称 `20260630release`,Git tag `20260630release` 和分支 `stable/20260630release` 均指向稳定代码提交 `9920575bbae7fbef4929843632f94d687ffbe289`。2026-06-30 开机后已恢复 A/B/T 服务:A `opensips`、`rtpengine-daemon`、`rtpengine-recording-daemon`、`lisglosips-redis-auth-proxy`、`lisglosips-node-exporter`、`lisglosips-recording-finalize.timer` active,`lisglosips-redis-hotpath-load.service` Result=success;B `mysql`、`redis-server`、`nginx`、`lisglosips@api`、`lisglosips@cdr-worker`、`lisglosips@recording-worker`、`lisglosips@config-publisher`、`heplify-server`、`lisglosips-prometheus`、`grafana-server` active,preflight ok,API ready ok,封版 MySQL 备份 `/data/backups/mysql/20260630T061723Z`,Redis 备份 `/data/backups/redis/20260630T061723Z`;T `opensips`、`rtpengine-daemon`、`rtpengine-recording-daemon`、`lisglosips-s28-uas`、`apache2`、`mariadb/mysql` active。B preflight、远程 smoke、8.2/8.7/8.8 定向复测均通过;业务前缀、号码库、客户网关编辑与前缀选择 Web-only 修复后,s53 当前首页引用 Web JS `index-CrjpbrU-.js`、CSS `index-CrJGI5_R.css`,并包含 `/favicon.ico`、`/brand/logo1.png`、`/brand/logo2.png`;封版远程 smoke `REMOTE_SMOKE_20260630T061636Z.md` 通过;完整按钮级远程浏览器巡检 `REMOTE_WEB_UI_SMOKE_20260629T104040Z.md` 通过,覆盖 15 个核心菜单及安全主操作,客户网关管理编辑报告包含 `business-prefix-checkbox-toggle-ok: false->true->false`;封版报告为 `docs/20260630_RELEASE_BASELINE_REPORT.md`;当前仍未导入真实 80 万手机号段;阿里云迁移前需按 S30 Runbook 重新演练 diff --git a/SOFTSWITCH_PLATFORM_DESIGN_V2.md b/SOFTSWITCH_PLATFORM_DESIGN_V2.md index 696001d..8e1142d 100644 --- a/SOFTSWITCH_PLATFORM_DESIGN_V2.md +++ b/SOFTSWITCH_PLATFORM_DESIGN_V2.md @@ -12,6 +12,10 @@ V2 的目标是把当前纯前端原型转化为可部署、可联调、可上 本文作为下一阶段执行基线,后续数据库 DDL、OpenAPI、OpenSIPS 脚本和部署脚本应从本文拆分并版本化管理。 +### 1.1 2026-08-31 新增设计:主叫号码实时分析 + +详见 [主叫号码接通率与应答率实时分析设计方案](docs/CALLER_REALTIME_ANALYTICS_DESIGN.md)。该功能处于设计完成、实现待开始阶段。新功能“接通数”按180 Ringing/183 Session Progress或后续阶段统计;“应答数”按实际通话时长大于0统计;实时总体应答率=应答数/总呼叫数,实时已接通应答率=应答数/接通数。直接成功、缺少100、未接通/待接通/未知数据及历史回算边界以专项文档V1.1为准。既有首页历史指标与计费规则不在本次文档任务中变更。 + ## 2. V2 范围 ### 2.1 本期实施页面 diff --git a/apps/api/src/modules/app.module.ts b/apps/api/src/modules/app.module.ts index 49d9b1f..9deaa54 100644 --- a/apps/api/src/modules/app.module.ts +++ b/apps/api/src/modules/app.module.ts @@ -1,4 +1,5 @@ import { Module } from '@nestjs/common'; +import { CallerAnalyticsModule } from './caller-analytics/caller-analytics.module.js'; import { ConfigModule } from '@nestjs/config'; import crypto from 'node:crypto'; import { LoggerModule } from 'nestjs-pino'; @@ -60,6 +61,7 @@ import { VendorsModule } from './vendors/vendors.module.js'; AuditModule, AuthModule, ActiveCallsModule, + CallerAnalyticsModule, BusinessPrefixesModule, CdrsModule, DashboardModule, diff --git a/apps/api/src/modules/caller-analytics/caller-analytics.module.ts b/apps/api/src/modules/caller-analytics/caller-analytics.module.ts new file mode 100644 index 0000000..9fd93f3 --- /dev/null +++ b/apps/api/src/modules/caller-analytics/caller-analytics.module.ts @@ -0,0 +1,18 @@ +import { Controller, Get, Inject, Module, Query } from '@nestjs/common'; +import { CurrentUserParam, RequirePermissions, type CurrentUser } from '../security/security.metadata.js'; +import { CallerAnalyticsService } from './caller-analytics.service.js'; + +@Controller('caller-analytics') +@RequirePermissions('caller_analytics.view') +class CallerAnalyticsController { + constructor(@Inject(CallerAnalyticsService) private readonly service: CallerAnalyticsService) {} + @Get('overview') overview(@Query() q: Record, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); } + @Get('summary') summary(@Query() q: Record, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); } + @Get('numbers') numbers(@Query() q: Record, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); } + @Get('trends') trends(@Query() q: Record, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); } + @Get('calls') calls(@Query() q: Record, @CurrentUserParam() user: CurrentUser) { return this.service.detail(q, user); } + @Get('breakdown') breakdown(@Query() q: Record, @CurrentUserParam() user: CurrentUser) { return this.service.detail(q, user); } + @Get('options') options(@CurrentUserParam() user: CurrentUser) { return this.service.options(user); } +} +@Module({ controllers: [CallerAnalyticsController], providers: [CallerAnalyticsService] }) +export class CallerAnalyticsModule {} diff --git a/apps/api/src/modules/caller-analytics/caller-analytics.service.spec.ts b/apps/api/src/modules/caller-analytics/caller-analytics.service.spec.ts new file mode 100644 index 0000000..6da0de6 --- /dev/null +++ b/apps/api/src/modules/caller-analytics/caller-analytics.service.spec.ts @@ -0,0 +1,15 @@ +import { describe, expect, it } from 'vitest'; +import { parseAnalyticsQuery } from './caller-analytics.service.js'; +describe('analytics query boundary', () => { + it('rejects inverted, excessive and invalid date windows', () => { + for (const q of [{from:'bad'}, {from:'2026-01-01',to:'2026-02-01'}, {from:'2026-08-02',to:'2026-08-01'}]) expect(() => parseAnalyticsQuery(q)).toThrow(); + }); + it('rejects sort injection and mixed-grain vendor filters', () => { + expect(() => parseAnalyticsQuery({sort:'caller; DROP TABLE x'})).toThrow(); + expect(() => parseAnalyticsQuery({view:'original',vendorGatewayId:'g'})).toThrow(); + }); + it('rejects invalid percentage and pagination values', () => { + expect(() => parseAnalyticsQuery({connectionMin:'101'})).toThrow(); + expect(() => parseAnalyticsQuery({take:'1.5'})).toThrow(); + }); +}); diff --git a/apps/api/src/modules/caller-analytics/caller-analytics.service.ts b/apps/api/src/modules/caller-analytics/caller-analytics.service.ts new file mode 100644 index 0000000..7736ca9 --- /dev/null +++ b/apps/api/src/modules/caller-analytics/caller-analytics.service.ts @@ -0,0 +1,113 @@ +import { BadRequestException, ForbiddenException, Inject, Injectable } from '@nestjs/common'; +import { ANALYTICS_VERSION, COUNT_KEYS, Prisma, analyticsRates, emptyCounts, type AnalyticsCounts } from '@lisglosips/database'; +import { PrismaService } from '../database/prisma.service.js'; +import type { CurrentUser } from '../security/security.metadata.js'; + +type Query = Record; +export interface AnalyticsQuery { + view: 'original' | 'landing'; from: Date; to: Date; customerId: string; caller: string; + customerGatewayId: string; vendorId: string; vendorGatewayId: string; city: string; carrier: string; + minSamples: number; skip: number; take: number; sort: string; order: string; onlyAlerts: boolean; + connectionMin: number; connectionMax: number; overallMin: number; overallMax: number; connectedMin: number; connectedMax: number; +} +export function parseAnalyticsQuery(raw: Query, now = new Date()): AnalyticsQuery { + const str = (key: string, max: number) => { + const v = raw[key] ?? ''; if (typeof v !== 'string' || v.length > max) throw new BadRequestException(`Invalid ${key}`); return v.trim(); + }; + const number = (key: string, fallback: number, min: number, max: number) => { + const v = raw[key] === undefined || raw[key] === '' ? fallback : Number(raw[key]); + if (!Number.isFinite(v) || v < min || v > max) throw new BadRequestException(`Invalid ${key}`); return v; + }; + const view = str('view', 16) || 'landing'; if (!['original', 'landing'].includes(view)) throw new BadRequestException('Invalid view'); + const to = raw.to ? new Date(str('to', 40)) : now; + 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'); + 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), + customerGatewayId: str('customerGatewayId', 32), vendorId: str('vendorId', 32), vendorGatewayId: str('vendorGatewayId', 32), city: str('city', 12), carrier: str('carrier', 20), + minSamples: number('minSamples', 0, 0, 1e9), skip: number('skip', 0, 0, 1e6), take: number('take', 25, 1, 100), sort, order, + onlyAlerts: raw.onlyAlerts === 'true', connectionMin: number('connectionMin', 0, 0, 100), connectionMax: number('connectionMax', 100, 0, 100), + overallMin: number('overallMin', 0, 0, 100), overallMax: number('overallMax', 100, 0, 100), connectedMin: number('connectedMin', 0, 0, 100), connectedMax: number('connectedMax', 100, 0, 100) + }; + if (!Number.isInteger(result.skip) || !Number.isInteger(result.take) || result.connectionMin > result.connectionMax || result.overallMin > result.overallMax || result.connectedMin > result.connectedMax) throw new BadRequestException('Invalid range'); + if (view === 'original' && (result.vendorId || result.vendorGatewayId)) throw new BadRequestException('线路筛选请切换落地主叫,避免混用业务呼叫与尝试分母'); + return result; +} +type Row = Record; +const countSql = Prisma.join(COUNT_KEYS.map(k => Prisma.sql`COALESCE(SUM(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts, ${`$.${k}`})) AS SIGNED)),0) AS ${Prisma.raw(k)}`)); +const normalize = (row: Row) => ({ ...row, caller: String(row.caller ?? ''), customerId: String(row.customerId ?? ''), ...analyticsRates(Object.fromEntries(COUNT_KEYS.map(k => [k, Number(row[k] ?? 0)])) as unknown as AnalyticsCounts) }); + +@Injectable() +export class CallerAnalyticsService { + constructor(@Inject(PrismaService) private readonly db: PrismaService) {} + + private async scope(user: CurrentUser, customerId: string): Promise { + if (user.permissions.includes('caller_analytics.view_all')) return null; + const access = await this.db.$queryRaw>`SELECT customer_id FROM caller_analysis_access WHERE user_id=${user.id}`; + const ids = access.map(a => a.customer_id); + if (!ids.length || (customerId && !ids.includes(customerId))) throw new ForbiddenException('没有该客户的主叫分析权限'); + return ids; + } + private where(q: AnalyticsQuery, scope: string[] | null): Prisma.Sql { + const parts = [Prisma.sql`view=${q.view}`]; + const fields = { customerId: 'customer_id', customerGatewayId: 'customer_gateway_id', vendorId: 'vendor_id', vendorGatewayId: 'vendor_gateway_id', caller: 'caller', city: 'city', carrier: 'carrier' } as const; + for (const [key, column] of Object.entries(fields)) { const v = q[key as keyof typeof fields]; if (v) parts.push(Prisma.sql`${Prisma.raw(column)}=${v}`); } + if (scope) parts.push(Prisma.sql`customer_id IN (${Prisma.join(scope)})`); + return Prisma.join(parts, ' AND '); + } + private source(q: AnalyticsQuery, scope: string[] | null) { + const start = new Date(Math.ceil(+q.from / 60000) * 60000); const end = new Date(Math.floor(+q.to / 60000) * 60000); + const where = this.where(q, scope); + if (start >= end) return Prisma.sql`SELECT * FROM caller_analysis_states WHERE ${where} AND started_at>=${q.from} AND started_at<${q.to}`; + const columns = Prisma.raw('started_at, customer_id, customer_gateway_id, vendor_id, vendor_gateway_id, caller, city, carrier, counts'); + return Prisma.sql`SELECT ${columns} FROM caller_analysis_minute_buckets WHERE ${where} AND started_at>=${start} AND started_at<${end} + UNION ALL SELECT ${columns} FROM caller_analysis_states WHERE ${where} AND started_at>=${q.from} AND started_at<${q.to} AND (started_at<${start} OR started_at>=${end})`; + } + async overview(raw: Query, user: CurrentUser) { + const q = parseAnalyticsQuery(raw); const scope = await this.scope(user, q.customerId); const source = this.source(q, scope); + const [numbers, trends, health, alerts] = await this.db.$transaction([ + this.db.$queryRaw(Prisma.sql`SELECT customer_id AS customerId, caller, ${countSql} FROM (${source}) a GROUP BY customer_id, caller HAVING SUM(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,'$.totalCalls')) AS SIGNED))>0 LIMIT 10001`), + this.db.$queryRaw(Prisma.sql`SELECT FLOOR(UNIX_TIMESTAMP(started_at)/60)*60000 AS time, ${countSql} FROM (${source}) a GROUP BY time ORDER BY time`), + this.db.$queryRaw; updated_at: Date }>>`SELECT payload, updated_at FROM caller_analysis_health WHERE id='worker'`, + this.db.$queryRaw>(Prisma.sql`SELECT customer_id, caller, payload, status FROM caller_analysis_alerts WHERE view=${q.view} ${scope ? Prisma.sql`AND customer_id IN (${Prisma.join(scope)})` : Prisma.empty} AND status='ACTIVE'`) + ]); + if (numbers.length > 10000) throw new BadRequestException('号码超过10000个,请缩短时间或筛选客户'); + const annotated = numbers.map(r => ({ ...normalize(r), alert: alerts.find(a => a.customer_id === r.customerId && a.caller === r.caller)?.payload ?? null })); + const fits = (v: number | null, min: number, max: number) => v === null ? min === 0 && max === 100 : v >= min && v <= max; + const filtered = annotated.filter(r => r.totalCalls >= q.minSamples && (!q.onlyAlerts || r.alert) && fits(r.connectionRate, q.connectionMin, q.connectionMax) && fits(r.overallAnswerRate, q.overallMin, q.overallMax) && fits(r.connectedAnswerRate, q.connectedMin, q.connectedMax)); + filtered.sort((a, b) => { + const av = a[q.sort as keyof typeof a]; const bv = b[q.sort as keyof typeof b]; + if (av === null) return bv === null ? 0 : 1; if (bv === null) return -1; + const cmp = typeof av === 'string' && typeof bv === 'string' ? av.localeCompare(bv) : Number(av) - Number(bv); + return (q.order === 'asc' ? cmp : -cmp) || String(a.caller).localeCompare(String(b.caller)); + }); + const sum = annotated.reduce((a, r) => { for (const k of COUNT_KEYS) a[k] += r[k]; return a; }, emptyCounts()); + const h = health[0]; const quality = !h || Date.now() - +h.updated_at > 10000 || h.payload.degraded ? 'DEGRADED' : 'LIVE'; + return { metricVersion: ANALYTICS_VERSION, view: q.view, from: q.from, to: q.to, asOf: new Date(), qualityStatus: quality, + dataThrough: h?.payload.dataThrough ?? null, lagMs: h?.payload.lagMs ?? null, durationPrecision: 'milliseconds', + health: h?.payload ?? null, summary: analyticsRates(sum), summaryScope: '时间及业务筛选;样本/比例/异常筛选仅影响号码列表', + numbers: filtered.slice(q.skip, q.skip + q.take), total: filtered.length, skip: q.skip, take: q.take, + trends: trends.map(normalize), historyNotice: '仅统计新采集链路的数据;旧固定6秒话单不参与应答统计或校准。' }; + } + async detail(raw: Query, user: CurrentUser) { + const q = parseAnalyticsQuery(raw); if (!q.caller || !q.customerId) throw new BadRequestException('请选择客户和主叫号码'); + const scope = await this.scope(user, q.customerId); const where = this.where(q, scope); + const rows = await this.db.$queryRaw>>(Prisma.sql`SELECT id, call_key AS callKey, started_at AS startedAt, caller, callee, vendor_gateway_id AS vendorGatewayId, counts, payload FROM caller_analysis_states WHERE ${where} AND started_at>=${q.from} AND started_at<${q.to} ORDER BY started_at DESC, id DESC LIMIT ${q.take} OFFSET ${q.skip}`); + const source = this.source(q, scope); + const breakdown = await this.db.$queryRaw(Prisma.sql`SELECT vendor_gateway_id AS vendorGatewayId, city, carrier, ${countSql} FROM (${source}) a GROUP BY vendor_gateway_id, city, carrier HAVING SUM(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,'$.totalCalls')) AS SIGNED))>0 LIMIT 500`); + return { rows, breakdown: breakdown.map(normalize), skip: q.skip, take: q.take }; + } + async options(user: CurrentUser) { + const scope = await this.scope(user, ''); + const [customers, customerGateways, vendorGateways] = await Promise.all([ + this.db.customer.findMany({ where: { ...(scope ? { id: { in: scope } } : {}), deletedAt: null }, select: { id: true, name: true }, take: 2000 }), + this.db.customerGateway.findMany({ where: { ...(scope ? { customerId: { in: scope } } : {}), deletedAt: null }, select: { id: true, name: true, customerId: true }, take: 2000 }), + // Limited users only receive gateway IDs they actually have traffic on, not the global vendor configuration. + scope ? this.db.$queryRaw>(Prisma.sql`SELECT DISTINCT vendor_gateway_id AS id, vendor_gateway_id AS name FROM caller_analysis_states WHERE customer_id IN (${Prisma.join(scope)}) AND vendor_gateway_id<>'' LIMIT 2000`) : this.db.vendorGateway.findMany({ where: { deletedAt: null }, select: { id: true, name: true }, take: 2000 }) + ]); + return { customers, customerGateways, vendorGateways }; + } +} diff --git a/apps/web/src/App.jsx b/apps/web/src/App.jsx index 495898b..4557d32 100644 --- a/apps/web/src/App.jsx +++ b/apps/web/src/App.jsx @@ -16,6 +16,7 @@ import { normalizeVendorGateway, } from './utils/formatters.js'; import { DashboardPage } from './pages/DashboardPage.jsx'; +import { CallerAnalyticsPage } from './pages/CallerAnalyticsPage.jsx'; import { ActiveCallsPage } from './pages/ActiveCallsPage.jsx'; import { CustomersPage } from './pages/CustomersPage.jsx'; import { RechargeRecordsPage } from './pages/RechargeRecordsPage.jsx'; @@ -58,6 +59,7 @@ const navGroups = [ items: [ { key: 'activeCalls', label: '当前通话', permissions: PAGE_PERMISSIONS.activeCalls }, { key: 'cdr', label: '话单中心', permissions: PAGE_PERMISSIONS.cdr }, + { key: 'callerAnalytics', label: '主叫号码分析', permissions: PAGE_PERMISSIONS.callerAnalytics }, { key: 'quality', label: '质检中心', permissions: PAGE_PERMISSIONS.quality }, { key: 'sipops', label: 'SIP 运维', pending: true }, { key: 'monitoring', label: '监控告警', pending: true }, @@ -76,6 +78,7 @@ const navGroups = [ const pages = { + callerAnalytics: CallerAnalyticsPage, dashboard: DashboardPage, activeCalls: ActiveCallsPage, customers: CustomersPage, diff --git a/apps/web/src/api.js b/apps/web/src/api.js index 9254bcf..f62c7e8 100644 --- a/apps/web/src/api.js +++ b/apps/web/src/api.js @@ -128,6 +128,9 @@ function queryString(params = {}) { } export const api = { + callerAnalytics: (params, options) => request(`/caller-analytics/overview${queryString(params)}`, options), + callerAnalyticsCalls: (params, options) => request(`/caller-analytics/calls${queryString(params)}`, options), + callerAnalyticsOptions: options => request('/caller-analytics/options', options), captcha: () => request('/auth/captcha'), login: async (body) => { const response = await request('/auth/login', { method: 'POST', body: jsonBody(body) }); diff --git a/apps/web/src/pages/CallerAnalyticsPage.jsx b/apps/web/src/pages/CallerAnalyticsPage.jsx new file mode 100644 index 0000000..9edd593 --- /dev/null +++ b/apps/web/src/pages/CallerAnalyticsPage.jsx @@ -0,0 +1,104 @@ +import { useEffect, useRef, useState } from 'react'; +import { Alert, Badge, Button, Field, Input, Select } from '../components/ui.jsx'; +import { PageTitle, Panel, Toolbar, SimpleTable, Drawer, Pagination } from '../components/layout.jsx'; +import { api, explainApiError } from '../api.js'; +import { formatDateTime } from '../utils/formatters.js'; + +const percent = v => v === null || v === undefined ? '—' : `${Number(v).toFixed(2)}%`; +const rateColumns = [ + { key: 'connectionRate', label: '实时接通率', render: r => percent(r.connectionRate) }, + { key: 'overallAnswerRate', label: '实时总体应答率', render: r => percent(r.overallAnswerRate) }, + { key: 'connectedAnswerRate', label: '实时已接通应答率', render: r => percent(r.connectedAnswerRate) } +]; +const countColumns = [{ key: 'totalCalls', label: '总呼叫数' }, { key: 'connectedCalls', label: '接通数' }, { key: 'notConnectedCalls', label: '未接通数' }, { key: 'answeredCalls', label: '应答数' }]; +const initial = { view: 'landing', minutes: '15', caller: '', customerId: '', customerGatewayId: '', vendorGatewayId: '', city: '', carrier: '', minSamples: '0', sort: 'totalCalls', order: 'desc', skip: 0, take: 25, onlyAlerts: false }; +export function CallerAnalyticsPage() { + const [draft, setDraft] = useState(initial); const [filters, setFilters] = useState(initial); + const [data, setData] = useState(null); const [error, setError] = useState(''); const [loading, setLoading] = useState(false); + const [options, setOptions] = useState({ customers: [], customerGateways: [], vendorGateways: [] }); + const [auto, setAuto] = useState(true); const [tick, setTick] = useState(0); + const [target, setTarget] = useState(null); const [detail, setDetail] = useState(null); const [detailError, setDetailError] = useState(''); + const [detailSkip, setDetailSkip] = useState(0); const sequence = useRef(0); + useEffect(() => { const c = new AbortController(); api.callerAnalyticsOptions({ signal: c.signal }).then(setOptions).catch(e => { if (!c.signal.aborted) setError(explainApiError(e)); }); return () => c.abort(); }, []); + useEffect(() => { + const controller = new AbortController(); const id = ++sequence.current; + let timer; + const refresh = async () => { + setLoading(true); + try { + const q = { ...filters, onlyAlerts: String(filters.onlyAlerts) }; + 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 }); + if (id === sequence.current) { setData(result); setError(''); } + } catch (e) { if (!controller.signal.aborted && id === sequence.current) setError(explainApiError(e)); } + finally { if (!controller.signal.aborted && id === sequence.current) { setLoading(false); if (auto) timer = setTimeout(refresh, 4000); } } + }; + void refresh(); return () => { controller.abort(); clearTimeout(timer); }; + }, [filters, auto, tick]); + useEffect(() => { + if (!target || !data) return undefined; + const c = new AbortController(); + api.callerAnalyticsCalls({ ...filters, from: data.from, to: data.to, caller: target.caller, customerId: target.customerId, skip: detailSkip, take: 25 }, { signal: c.signal }) + .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 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) => ; + const openDetail = row => { setTarget(row); setDetail(null); setDetailSkip(0); setDetailError(''); }; + const metrics = data?.summary; + const timeline = data?.trends ?? []; + const max = Math.max(1, ...timeline.map(r => r.totalCalls)); + return <> + } /> + {error ? {error} : null} +
{ e.preventDefault(); setFilters({ ...draft, skip: 0 }); setTarget(null); }}> + + + + {draft.minutes === 'custom' ? <> change('from', e.target.value)} /> change('to', e.target.value)} /> : null} + change('caller', e.target.value)} placeholder="精确匹配号码" /> + {select('客户', 'customerId', options.customers)} + {select('客户网关', 'customerGatewayId', options.customerGateways.filter(g => !draft.customerId || g.customerId === draft.customerId))} + {draft.view === 'landing' ? select('落地网关', 'vendorGatewayId', options.vendorGateways) : null} + {select('被叫运营商', 'carrier', ['MOBILE','UNICOM','TELECOM','BROADCAST','MVNO','UNKNOWN'].map((id,i) => ({id,name:['移动','联通','电信','广电','虚拟运营商','未知'][i]})))} + change('city', e.target.value)} /> + change('minSamples', e.target.value)} /> + + + + +
接通率 / 应答率区间{[['connection','接通率'],['overall','总体应答率'],['connected','已接通应答率']].map(([key,label]) => change(`${key}Min`, e.target.value)} placeholder="0" /> change(`${key}Max`, e.target.value)} placeholder="100" />)}
+ +
+
{[...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} + 灰:呼叫 / 蓝:接通 / 绿:应答}> +
{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]} />
} +
+ {data?.total ?? 0} 个号码} className="wide-panel"> + ({ ...r, id:`${r.customerId}:${r.caller}`, customerName:names[r.customerId] || r.customerId }))} loading={!data && loading} onRowClick={openDetail} columns={[ + {key:'caller',label:'主叫号码',render:r => }, {key:'customerName',label:'客户'}, ...countColumns,...rateColumns, + {key:'pendingCalls',label:'待接通'}, {key:'unknownCalls',label:'未知'}, {key:'alert',label:'状态',render:r => {r.alert ? r.alert.reasons?.join('、') || '告警' : r.unknownCalls ? '证据不足' : '观察'}} + ]} /> + setFilters(f => ({...f,skip}))} onPageSizeChange={take => setFilters(f => ({...f,take,skip:0}))} /> + + {target ? setTarget(null)}> + {detailError ? {detailError} : null} +

{names[target.customerId] || target.customerId} · {filters.view === 'landing' ? '落地主叫' : '原始主叫'}

+ ({...r,id:i,gateway:gatewayNames[r.vendorGatewayId] || r.vendorGatewayId || '业务呼叫汇总'}))} columns={[{key:'gateway',label:'落地网关'},{key:'city',label:'城市'},{key:'carrier',label:'运营商'},...countColumns,...rateColumns]} /> + { + const p=r.payload; const legs=p.legs ? Object.values(p.legs) : [p]; + return {...r, callId:p.start?.callId || p.event?.callId, original:p.start?.caller || p.event?.caller, + phase:legs.map(l => l.first180At ? '180 Ringing' : l.first183At ? '183 Session Progress' : l.acceptedAt ? '直接成功' : '未达接通').join(' / '), + final:legs.map(l => l.finalCode || '进行中').join(' / '), duration:`${(r.counts.talkMs/1000).toFixed(3)} 秒${r.counts.activeCalls ? '(暂定)' : ''}`, time:formatDateTime(r.startedAt)}; + })} columns={[{key:'callId',label:'Call-ID'},{key:'original',label:'原始主叫'},{key:'callee',label:'被叫'},{key:'phase',label:'接通证据'},{key:'final',label:'最终结果'},{key:'duration',label:'通话时长'},{key:'time',label:'发起时间'}]} /> +
+

按 Call-ID 对应原始话单或信令;旧固定时长话单仅作追溯,不参与本统计校准。

+
: null} + ; +} diff --git a/apps/web/src/permissions.js b/apps/web/src/permissions.js index 74d5a0a..d24bf2d 100644 --- a/apps/web/src/permissions.js +++ b/apps/web/src/permissions.js @@ -1,4 +1,5 @@ export const PAGE_PERMISSIONS = { + callerAnalytics: ['caller_analytics.view'], dashboard: ['dashboard.view'], activeCalls: ['active_calls.view'], customers: ['customers.view'], diff --git a/apps/web/src/styles.css b/apps/web/src/styles.css index 3250b2d..dcb933e 100644 --- a/apps/web/src/styles.css +++ b/apps/web/src/styles.css @@ -2893,3 +2893,12 @@ a { animation: none; } } +.cra-metrics { grid-template-columns: repeat(auto-fit, minmax(150px, 1fr)); margin: 20px 0; } +.cra-rate-filters { margin: 12px 0; } +.cra-chart { display: flex; gap: 8px; overflow-x: auto; min-height: 160px; padding: 12px 0; } +.cra-chart-column { flex: 1 0 38px; max-width: 70px; text-align: center; } +.cra-bars { display: flex; align-items: flex-end; justify-content: center; height: 120px; gap: 2px; } +.cra-bars span { width: 9px; min-height: 1px; background: #94a3b8; border-radius: 2px 2px 0 0; } +.cra-bars span:nth-child(2) { background: #2563eb; } +.cra-bars span:nth-child(3) { background: #16a34a; } +.cra-chart small { color: #64748b; font-size: 10px; } diff --git a/apps/worker-cdr/src/analytics-alerts.ts b/apps/worker-cdr/src/analytics-alerts.ts new file mode 100644 index 0000000..17a356b --- /dev/null +++ b/apps/worker-cdr/src/analytics-alerts.ts @@ -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>(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>(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>`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>`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'`; +} diff --git a/apps/worker-cdr/src/analytics-main.ts b/apps/worker-cdr/src/analytics-main.ts new file mode 100644 index 0000000..9d4d4e8 --- /dev/null +++ b/apps/worker-cdr/src/analytics-main.ts @@ -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>`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>; + 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(); }); diff --git a/docs/CALLER_REALTIME_ANALYTICS_DESIGN.md b/docs/CALLER_REALTIME_ANALYTICS_DESIGN.md new file mode 100644 index 0000000..13c3336 --- /dev/null +++ b/docs/CALLER_REALTIME_ANALYTICS_DESIGN.md @@ -0,0 +1,220 @@ +# 主叫号码接通率与应答率实时分析设计方案 + +日期:2026-08-31 +版本:V1.1(按用户修订后的业务口径) +状态:设计文档已完成;代码、数据库迁移、真实联调与上线均待实施。 + +实施更新(2026-08-31):独立分析链路、API、页面和测试脚本已编写;部署与真实验收结果以实施状态及发布报告为准,上述为方案编制时状态。 + +## 1. 目标与范围 + +新增“主叫号码分析”页面,按号码实时查看呼叫是否到达振铃/会话进展阶段,以及是否产生正通话时长。支持异常发现、线路对比、失败原因分析和话单追溯。 + +本方案中的“接通”是平台自定义的业务阶段指标,不等于 SIP 最终成功应答,也不等于真人接听。此前方案中将接通数定义为 INVITE 成功应答的口径由本文替代。所有新增页面、接口、告警和测试应使用同一个口径版本 `caller-analytics-v1.1`。 + +第一期只做统计和站内异常提示,不自动停号、换号、切换线路或修改计费。外部通知渠道另行配置授权,不在本轮发送消息。 + +## 2. 当前源码基础与实施差距 + +本地源码核对依据,不代表生产状态已核实: + +- `prisma/schema.prisma` 的 `RawCdr` 已有 caller、landingCaller、客户/网关/线路维度、startedAt、answeredAt、endedAt、durationSec、sipCode 等字段。 +- `packages/redis/src/cdr-stream.ts` 已有 CDR 事件流契约及消费机制;`apps/worker-cdr/src/rating.ts` 负责话单落库及计费。 +- `apps/api/src/modules/dashboard/` 现有指标按话单 SIP 2xx 计算,不能直接复用为本文的“接通数”或“应答数”。 +- 现有 RawCdr 没有独立的 first180At、first183At 等过程字段,最终响应码不能证明之前是否发生过振铃/进展。 +- 当前通话页面的 durationSec 存在从 startedAt 计算的逻辑,不能未经核验直接作为本文的通话时长。 + +需新增呼叫过程采集、分析状态及预聚合。既有计费规则不随本次业务术语改变;首页旧指标保留旧版本语义,不可只改标签冒充新口径。后续统一首页时,应单独替换数据来源并展示口径版本。 + +## 3. 业务口径 + +### 3.1 总呼叫数与统计视角 + +默认“落地主叫”视角,可切换“原始主叫”: + +| 视角 | 号码与总呼叫数 T | 用途 | +| --- | --- | --- | +| 原始主叫 | 客户送入号码;每个已识别客户的独立业务呼叫计一次 | 分析客户号码整体效果,包括平台内业务拒绝 | +| 落地主叫 | 平台向线路发送的号码;每个实际发出的线路尝试计一次 | 分析具体出局号码及线路表现 | + +未认证扫描、OPTIONS、REGISTER 等不计入。客户一次呼叫的鉴权续发、SIP 重传不新增业务呼叫;线路重试用 attemptId 区分。原始主叫视角任一尝试接通则该业务呼叫接通,任一有效通话产生正时长则该业务呼叫应答,均最多计一次。不同线路尝试可分别计数,但不得相加后冒充客户呼叫数。 + +落地主叫不等于已经验证的被叫终端显示号码。保存号码原值和标准化值,不得无规则删除前缀;同号不同客户不能混淆归属。 + +### 3.2 接通数 C:达到 180 / 183 或更后阶段 + +接通数指初始 INVITE 已进入有效呼叫处理,并首次达到以下任一条件的去重呼叫数: + +1. 收到与该初始 INVITE/线路尝试匹配的 `180 Ringing`; +2. 收到匹配的 `183 Session Progress`,不要求必须带 SDP; +3. 没有观测到 180/183,但直接收到初始 INVITE 的有效成功 2xx(通常为 `200 OK`):作为已进入后续阶段计入接通,原因记录为 `DIRECT_FINAL_SUCCESS`; +4. 最终权威话单证明该呼叫产生正通话时长,但过程事件缺失:补计接通与应答,标注 `CDR_POSITIVE_DURATION`,不能伪造 180/183 时间。 + +“INVITE、Trying 均成功”按业务意图理解为 INVITE 已受理并继续推进,不硬性要求抓到 100 Trying。100 是临时处理响应;180 表示尝试提醒用户;183 表示会话进展,不证明实际振铃。状态代理不会向上游转发收到的 100,因此不能以缺少 100 排除后续有效进展。上述协议含义参考 [RFC 3261 §21.1](https://www.rfc-editor.org/rfc/rfc3261.html#section-21.1)。本文选择 180/183 作为业务接通门槛,是产品口径。 + +仅 INVITE 发出、100、181 或 182 不构成接通。匹配必须核实请求方法、事务与呼叫标识,不能使用 BYE、CANCEL、PRACK、OPTIONS 的 200,也不能把 re-INVITE 当新呼叫。响应与方法相关的依据见 [RFC 3261 §21.2.1](https://www.rfc-editor.org/rfc/rfc3261.html#section-21.2.1)。 + +线路视角只采信该出局尝试的实际下游响应,不能用本机合成振铃冒充线路进展。接通一旦成立,后续 486、487、超时等不会撤销;只能在证据校正或误关联修复时重算。 + +### 3.3 未接通数 N:未达到接通门槛 + +业务定义:未达到上述振铃/会话进展或后续阶段的呼叫数。实时展示必须区分“尚未达到”和“已经结束且未达到”: + +- F:已结束且从未接通的呼叫,包括明确拒绝、取消或确认的呼叫超时。 +- P:仍在处理中、尚未接通的呼叫;页面显示“待接通”,不作为最终失败。 +- U:证据不足而无法判定是否接通的呼叫;页面显示“状态未知”,不强行归类。 + +已确定的实时未接通数 `N = F + P`,并标注“含待接通 P 次”。在无未知数据时 `T = C + N`;有未知时 `T = C + N + U`,未接通率只能展示已确认范围或暂不展示完整率,不能把 U 隐式算作失败。全部呼叫结束且完成对账后 P=U=0,此时 `N=F=T-C`。 + +接通后无人接听、接通后取消、接通后最终失败,均归入“已接通、未应答”,不归入未接通数。另设“最终结果码分布”,避免把协议最终失败数等同未接通数。 + +### 3.4 应答数 A:通话时长大于 0 + +应答数指通话时长严格大于 0 的去重呼叫数。不以单独收到 200、存在 answeredAt 或计费时长大于 0 代替。 + +- 通话时长采用业务通话段的实际时长,不含呼叫建立、振铃、183 早期媒体时长,不用计费向上取整时长。 +- 新采集建议保留 `talkDurationMs` 毫秒值,以 `talkDurationMs > 0` 判断,再格式化显示秒数。正的亚秒通话也应计入,不能因显示向下取整为 0 秒而漏计。 +- 原始时长只有整数秒的历史话单,`durationSec > 0` 可计入应答;`durationSec = 0` 不得凭空推定存在亚秒通话,展示历史精度限制。 +- 进行中的呼叫由权威会话状态及通话计时起点计算已发生的正时长,实时标注暂定。不能由前端收到 200 后自行运行计时器,也不能用 INVITE 发起时间作通话起点。 +- 正时长证据不要求有 RTP 音频分析或真人检测;ACK/会话状态用于核实计时来源,不额外改变“实际通话时长 > 0”的最终业务条件。 +- 最终以核验过来源和计时口径的 CDR/会话时长校准,允许撤销错误的暂定应答,并在同一事务调整统计。字段矛盾、负时长或缺少可靠计时证据进入数据质量核对。 + +应答数是接通数的子集,满足 `0 ≤ A ≤ C ≤ T`。其中 C-A 包含仍在等待应答、已结束但没有正时长,以及应答时长待核实的呼叫,不能全部称为“应答失败”。 + +### 3.5 指标公式 + +所有分子分母必须使用相同客户权限、号码视角、线路维度、发起时间窗口和快照,不混用全量总数与筛选后数量。 + +| 中文指标 | 接口建议字段 | 公式 | +| --- | --- | --- | +| 总呼叫数 | totalCalls | T | +| 接通数 | connectedCalls | C | +| 未接通数(含待接通) | notConnectedCalls | N=F+P;U 单列 | +| 应答数 | answeredCalls | A | +| 实时接通率 | connectionRate | C / T × 100% | +| 实时总体应答率 | overallAnswerRate | A / T × 100% | +| 实时已接通应答率 | connectedAnswerRate | A / C × 100% | + +分母为 0 时接口返回 null、页面显示“—”。比率由整数计数计算,不平均各时间桶的百分比。数据不完整时,计数与比率应标注“已观测/暂定”,未知或采集缺口不可展示为完整结果。 + +旧方案的“已判定接通率”不作为主指标,本版以以上三个率为准;不使用未经解释的 ASR 简称混淆业务口径。 + +示例:T=100,C=60(其中 A=30),F=25,P=15,U=0。未接通数 N=40,其中待接通15;实时接通率60%,实时总体应答率30%,实时已接通应答率50%。C 中其余30次可能仍在振铃,也可能已经结束但没有正时长。 + +## 4. 时间窗口与状态更新 + +- 最近5/15/30/60分钟、今日、自定义;以发起时间归桶,后续事件回写原桶。原始视角用业务呼叫发起时间,线路视角用尝试发起时间。 +- 聚合粒度默认1分钟。分钟对齐窗口必须在界面展示实际起止;自定义非整分钟边界从状态/事实数据补算边缘,禁止悄悄扩大范围。 +- 数据库使用UTC,今日按Asia/Shanghai自然日计算。分子、分母与对比窗口使用统一服务端asOf。 +- 记录独立事实标志:是否接通、是否正时长、是否结束、是否数据完整,而非只依赖最终 SIP 码。 +- 正常过程:初始计T与P;达到180/183计C并减P;正通话时长计A;结束只更新状态/时长,不重复增加C或A。 +- 未接通直接结束:P减1、F加1,N不变。已接通后失败:C不变、A按时长判定,F不增加。 +- 业务呼叫未完成全部重试前,不能因一个尝试失败就判定客户呼叫最终未接通。 +- 乱序、重复、跨分钟、跨日和迟到事件按稳定标识重算贡献差值;接通数不因183转180重复增加。 +- 采集失联不能直接判业务超时;超时失败必须有事务/会话结束依据,否则进入U和数据异常提示。 + +## 5. 页面与交互 + +新增“主叫号码分析”,默认落地主叫,支持原始主叫切换。筛选包含时间、号码精确搜索、客户、客户网关、供应商、落地网关、被叫运营商/地区、最小样本、三个率的区间和仅看异常。 + +列表核心列:主叫号码、客户、总呼叫数、接通数、未接通数、应答数、实时接通率、实时总体应答率、实时已接通应答率、异常状态。待接通、未知、已结束未接通及前窗百分点变化可作为补充列或展开内容;宽表沿用横向滚动与固定标识列,不把业务数字藏到提示框。 + +号码详情展示分钟趋势(T/C/A及三个率)、线路对比、被叫运营商/地区拆分、180/183/直接成功分布、未接通原因、接通未应答的最终结果、主叫改写映射和关联话单。页面文案解释183不保证实际振铃,接通不等于应答。 + +每5秒刷新,保留筛选、排序、分页和展开详情。展示数据截至时间、采集延迟、统计口径版本和完整性。刷新失败保留旧数据并提示,不能清空成零。号码访问与导出按客户范围/权限控制,查询接口自行校验权限,不能只依赖前端隐藏。 + +## 6. 事件链路与数据设计 + +采用“OpenSIPS呼叫事件 → 独立Redis Stream → 分析Worker → MySQL状态/分钟汇总 → API及缓存 → 页面”。复用现有技术栈,不改变计费Worker消费组或计费规则。 + +事件类型建议为 CALL_STARTED、ATTEMPT_STARTED、TRYING、RINGING、SESSION_PROGRESS、INVITE_ACCEPTED、TALK_DURATION_OBSERVED、CALL_ENDED、CDR_RECONCILED。均需schemaVersion、eventId、callUid、attemptId(适用时)、nodeId/实例启动标识、occurredAt、receivedAt、初始事务标识、客户/线路/号码快照及证据来源。保留SIP Call-ID用于追溯,但不单独用它承担业务全局唯一性。 + +过程字段至少包括 first180At、first183At、inviteAcceptedAt、talkStartedAt、endedAt、talkDurationMs、connectedEvidence、finalSipCode、dataQuality。直接成功和CDR补计不得伪造first180At/first183At。通话计时起点须通过当前OpenSIPS实际事件与CDR生成代码核验;字段名不作为真实性证明。 + +| 建议数据对象 | 关键内容 | +| --- | --- | +| caller_analysis_events | eventId唯一约束、呼叫/尝试、事件时间、类型、来源、重放与去重证据 | +| caller_analysis_states | 以视角粒度保存T/C/N/A贡献、F/P/U状态、号码/维度快照、版本和时长来源 | +| caller_analysis_minute_buckets | 分钟、客户、号码类型/标准化号码、网关/线路、被叫分类、口径版本及各项计数/时长 | +| caller_analysis_alerts | 规则、触发快照、样本量、持续与恢复状态、冷却时间 | + +原始呼叫总览和线路尝试汇总分开维护,或通过显式grain字段区分,避免多线路重复累计。分钟聚合的唯一键应使用稳定非空维度键,缺失维度用受控unknown键,不能依赖MySQL可空复合唯一键实现幂等。索引与维度组合在获知号码基数、呼叫量、保留期后验证,不宣称当前机器已具备容量。 + +分析Worker的去重记录、状态更新和统计差值必须同一数据库事务提交;提交后ACK。Redis缓存不是唯一事实来源;数据库提交后缓存更新失败可失效重建。消费者独立,使用独立去重命名空间,不复用计费锁以免互相吞事件。消费组重领/重放场景须验证;机制参考 [Redis Streams](https://redis.io/docs/latest/develop/data-types/streams/)。 + +采集失败使用有界补送/持久化方案,设计重试期限、积压容量和告警;不得让统计故障长时间阻塞SIP热路径。超过留存或补送能力产生的缺口必须可见,不能承诺无条件零丢失。最终CDR通过同一callUid/attemptId校准状态,不作为新呼叫重复累计。旧CDR关联不足时先核对,不能仅凭Call-ID猜测匹配。 + +建议API:GET /api/v2/caller-analytics/summary、/numbers、/trends、/breakdown、/calls;共享视角、时间与维度筛选契约。响应包含counts、rates、asOf、dataThrough、lagMs、metricVersion、qualityStatus、unknownCalls及durationPrecision。完整性无法判定时lagMs可为null,不用“最后一个业务事件的年龄”冒充采集延迟。 + +## 7. 历史回算与兼容 + +1. 有可信正时长的历史CDR可回算A并证明对应C;仅有2xx而无时长需核验CDR字段语义才能证明C,不能证明A。 +2. 最终486/487/超时的话单可能先有180/183,仅凭最终码无法判定未接通。 +3. 若历史信令留存完整且能准确关联,可以补回过程;否则C是已知下界,N和A/C不能输出为完整历史指标。展示“历史过程数据不足”,完整接通率/已接通应答率返回null,允许另列清楚标注的已观测计数。 +4. 历史时长精度与新采集不同,应带durationPrecision,不把历史整数秒与新毫秒精度差异解释成业务变化。 +5. 新实时数据生效时间、口径版本、回算批次明确记录。新旧版本不能无提示混画同一趋势或用于同比告警。 + +## 8. 异常规则 + +分别支持“低接通率”“低总体应答率”“低已接通应答率”“较前窗下降”和“连续已结束未接通”。接通率低侧重排查建立链路;接通率正常但已接通应答率低,侧重排查接通后的结果。仅作排查方向,不凭比例或单个SIP码认定封号。 + +建议初始模板:最近15分钟、总样本≥50;按已接通应答率触发时还需C≥30;与前一等长窗口比较时两窗均满足样本门槛;下降超过15个百分点且持续多个评估周期再提示。三个率的绝对低值阈值分别配置,原方案20%的示例不能无校准套用到新接通口径。 + +设置未结束占比上限、数据完整性门槛、冷却、合并和恢复迟滞;仍大量振铃、持续通话缺少可信时长、采集滞后或未知数据过多时暂缓业务告警,转为数据质量提示。连续未接通按发起顺序和已确定结果评估,不能按乱序到达的事件简单累加。 + +## 9. 实施步骤、风险与回滚 + +1. 核验当前生产OpenSIPS/CDR事件、通话计时精度、线路重试、号码基数、峰值CPS、保留期与机器资源;只读核验与受控呼叫分开安排。 +2. 实施话单分析基础版:视角、列表、正时长应答、历史完整性提示、话单钻取。缺少过程数据时不宣称已支持完整接通统计。 +3. 实施过程事件、独立Worker、预聚合、真实后端接口、进行中正时长更新、最终对账与站内异常,完成实时版。 +4. 在受控真实呼叫与故障恢复验收后灰度启用;先观察数据质量和统计/计费隔离,再扩大覆盖。 + +新增影响包括OpenSIPS事件采集开销、Redis队列容量、MySQL写放大、查询和缓存负载。统计新表/新服务独立部署并使用功能开关;回滚可关闭采集钩子和分析Worker/入口,保留证据及新表,不删除既有话单、不调整计费余额和路由。生产修改及受控外呼须按届时任务授权执行。 + +## 10. 验收用例(设计,尚未执行) + +以下单呼叫终态用例默认T=1;应答时长均指可信实际通话时长。未说明缺口的用例假设完整采集。 + +| ID | 场景 | 预期 | +| --- | --- | --- | +| CRA-001 | INVITE→100→180→200,正通话时长 | C=1,N=0,A=1;三个率均100% | +| CRA-002 | INVITE→100→183→487,无正时长 | C=1,N=0,A=0;接通率100%,两个应答率0% | +| CRA-003 | INVITE→100→486,从未180/183/2xx | C=0,N=F=1,A=0;总体应答率0%,已接通应答率为null | +| CRA-004 | 未见100,直接180,随后取消 | C=1,A=0,不因缺100漏计 | +| CRA-005 | 直接初始INVITE 200,正时长 | C=A=1;DIRECT_FINAL_SUCCESS,无伪造振铃时间 | +| CRA-006 | 180或直接200,最终实际时长严格为0 | C=1,A=0,不能仅凭200计应答 | +| CRA-007 | 183不带SDP,或带早期媒体但未建立通话 | C=1,A=0,早期媒体不计通话时长 | +| CRA-008 | 仅100且仍在处理 | C=A=0,P=N=1,显示待接通,不标最终失败 | +| CRA-009 | 仅181/182后失败 | C=A=0,F=N=1 | +| CRA-010 | 重复INVITE、180、183、200、CDR及消费重放 | T/C/A不重复,N不为负,不重复扣费 | +| CRA-011 | BYE/PRACK/CANCEL的200、re-INVITE | 不新增T/C/A | +| CRA-012 | 先180后486或确认超时 | C=1,A=0;最终失败分布增加,未接通数不增加 | +| CRA-013 | 可信通话时长500ms,页面秒值取整为0 | 新精度A=1;旧整数秒0保持历史精度限制 | +| CRA-014 | 通话进行中已产生正时长,最终CDR校正 | 实时A可见且标暂定;结束不重复,错误暂定值可事务回退 | +| CRA-015 | 一次客户呼叫经两条线路尝试,第二条应答 | 原始视角T=1;线路视角共2次尝试,归属正确 | +| CRA-016 | 同号不同客户、主叫改写及客户越权查询 | 数据不串号/不越权,两个视角的分母保持一致 | +| CRA-017 | 乱序、跨窗口/跨日应答、迟到与重启重领 | 回写发起桶,重放前后计数一致,计费链路不受影响 | +| CRA-018 | 零呼叫或零接通 | 分母0返回null;不显示NaN或无依据的0% | +| CRA-019 | 历史487,缺少过程事件 | 不断言N=1;标未知,完整接通率及A/C不可用 | +| CRA-020 | Redis/MySQL/Worker故障、补送溢出 | 显示滞后/缺口,停止相关业务误告警,恢复重放可对账 | +| CRA-021 | T=100/C=60/A=30/F=25/P=15 | N=40;三个率依次60%/30%/50% | +| CRA-022 | 本机合成180,下游直接失败且无进展 | 线路视角不计C,不能用合成回复美化线路结果 | +| CRA-023 | 只有可信最终正时长CDR,过程丢失 | 补C/A并标CDR证据,不伪造180/183,不重复T | +| CRA-024 | 采集失联导致一个呼叫阶段无法判定 | U单列,T=C+N+U,不自动判失败,率标暂定/不完整 | + +功能验收须形成真实SIP信令、API返回、数据库状态/聚合与最终话单的对应证据;mock、静态页面或localStorage不构成功能验收。单元测试用于公式/状态转换,真实联调用于端到端与故障恢复。正常约定负载下事件或可信正时长可用至页面可见的P95目标≤5秒;这是待测目标,不是当前容量承诺。需分别记录采集、消费、API和页面延迟及原有呼叫/计费回归结果。 + +## 11. 本次交付边界 + +本次仅完成方案、业务口径、文档索引和验收用例编写与一致性检查。没有修改应用代码、执行数据库迁移、提交/推送Git、连接生产实施或发起真实呼叫。后续开发应以本文V1.1为准,不沿用“2xx=接通/应答”的旧统计实现。 + +## 12. 实施补充(2026-08-31) + +- 第11节为方案编制轮次的交付边界。后续用户已授权实施、提交推送与部署测试环境。 +- B现场核验发现旧CDR在收到200后直接填入固定6秒,新链路禁止用该字段校准。新链路使用同一OpenSIPS实例的成功应答和结束观察时间差,毫秒精度;进行中以MI确认会话存在后标记暂定正时长。原计费规则与旧CDR生成保持不变,既有固定时长问题需独立整改。 +- 采集先由OpenSIPS写结构化本机日志,经rsyslog队列保留,再写独立Redis Stream;本机日志游标在XADD成功后落盘,重放由数据库事务去重。轮转/截断及异常记录进入缺口提示,不承诺无限留存或无条件零丢失。 +- 新权限为caller_analytics.view及caller_analytics.view_all;迁移只默认授予超级管理员。其他账号需配置caller_analysis_access客户范围,否则拒绝访问;普通view不隐含全客户权限。 +- 卡片和趋势按时间及业务筛选展示整体,最小样本/比例区间/仅告警只筛选号码列表,界面明确提示该范围。单次最多10000个聚合号码,超出要求缩小时间或筛选客户;详细分布最多500个组合,查询窗口最多7天。 +- 实际当前OpenSIPS路由只选择一次落地,没有切换重试流程;采集标记单attempt,状态层支持多attempt且有测试,未来新增真实重试必须同步采集独立attemptId,不能直接复用单attempt钩子。 +- 历史无过程证据的话单不自动导入完整统计。页面显示新采集生效说明,旧数据不可用不等于呼叫数为0;历史补录仅接受已核验来源的RECONCILE,不接受旧固定时长CDR。 +- 告警通过独立进程环境变量配置三个低率阈值和前窗下降百分点;采用样本/未结束占比门槛、连续三次评估及三次恢复,站内展示,不执行停号换路。后续应按实际业务校准默认阈值。 diff --git a/docs/TEST_PLAN_AND_CASES.md b/docs/TEST_PLAN_AND_CASES.md index 232e9e3..223a83b 100644 --- a/docs/TEST_PLAN_AND_CASES.md +++ b/docs/TEST_PLAN_AND_CASES.md @@ -1975,6 +1975,10 @@ corepack pnpm@10.33.0 exec vitest run apps/worker-recording/src/transfer.spec.ts | 步骤 | 1. 导入器先干跑并确认未匹配为 0。2. 备份后导入。3. 核对四表数量。4. 四个页签执行首屏、翻页、页大小和筛选。 | | 预期结果 | 371 个行政城市、517,258 个有效七位号段、321 个区号、70 个前缀规则;首屏只读取 25 条;上一页/下一页与总数正确;页面加载不扫描或返回全表。 | +## 8A. 主叫号码实时分析验收设计索引(2026-08-31) + +[主叫号码实时分析V1.1](CALLER_REALTIME_ANALYTICS_DESIGN.md)第10节定义CRA-001至CRA-024,覆盖180/183业务接通、正时长应答、总体应答率A/T、已接通应答率A/C、直接200/缺少100/零时长/亚秒时长、待接通与未知、线路重试、重复乱序、历史缺口、权限隔离和故障恢复。该组用例目前全部待执行,不计入既有通过数量;最终验收必须有真实SIP、API、数据库及CDR对账证据。 + ## 9. 缺陷分级 | 级别 | 定义 | 示例 | diff --git a/infra/server-b/caller-analytics/deploy.sh b/infra/server-b/caller-analytics/deploy.sh new file mode 100644 index 0000000..26ec422 --- /dev/null +++ b/infra/server-b/caller-analytics/deploy.sh @@ -0,0 +1,85 @@ +#!/usr/bin/env bash +set -Eeuo pipefail +commit=${1:?commit required} +archive=${2:?archive required} +expected=${3:?sha256 required} +[[ "$commit" =~ ^[0-9a-f]{40}$ ]] +test "$(sha256sum "$archive" | awk '{print $1}')" = "$expected" +stamp=$(date -u +%Y%m%dT%H%M%SZ) +old=$(readlink -f /opt/lisglosips/current) +new=/opt/lisglosips/releases/caller-analytics-$stamp +backup=/var/backups/lisglosips-caller-analytics/$stamp +test -d "$old" +test ! -e "$new" +# Refuse to interrupt calls, including unrelated test calls. +curl -fsS -H 'Content-Type: application/json' --data '{"jsonrpc":"2.0","id":1,"method":"dlg_list","params":[]}' http://127.0.0.1:8888/mi | python3 -c 'import sys,json; assert json.load(sys.stdin)["result"]["Dialogs"]==[], "Active calls: retry after calls end"' +systemctl start lisglosips-backup.service +install -d -m 0700 "$backup" +cp -a /etc/opensips/opensips.cfg "$backup/opensips.cfg" +printf '%s\n' "$old" > "$backup/previous-release" +for file in /etc/lisglosips/caller-analytics.env /etc/rsyslog.d/35-caller-analytics.conf; do + if test -f "$file"; then cp -a "$file" "$backup/$(basename "$file")"; fi +done +changed=0 +rollback() { + rc=$? + if test "$changed" = 1; then + systemctl stop lisglosips@caller-analytics || true + cp -a "$backup/opensips.cfg" /etc/opensips/opensips.cfg + ln -sfn "$old" /opt/lisglosips/current + systemctl restart opensips lisglosips@api || true + fi + echo "FAILED rc=$rc backup=$backup (additive analytics tables retained)" + exit "$rc" +} +trap rollback ERR +cp -a --reflink=auto "$old" "$new" +tar -xzf "$archive" -C "$new" +printf '%s\n' "$commit" > "$new/.deployed-commit" +printf '%s\n' "$old" > "$new/.delta-base-release" +source_cfg="$backup/opensips.cfg" +if grep -q 'CRA1|' "$source_cfg"; then + # A code-only redeploy keeps previously reviewed instrumentation unchanged. + cp "$source_cfg" "$new/opensips-candidate.cfg" +else + /usr/bin/node "$new/scripts/instrument-caller-analytics.mjs" "$source_cfg" "$new/opensips-candidate.cfg" +fi +opensips -C -f "$new/opensips-candidate.cfg" +set -a +source /etc/lisglosips/secrets/mysql-migrate.env +set +a +export DATABASE_URL +DATABASE_URL=$(/usr/bin/node -e 'console.log(`mysql://${encodeURIComponent(process.env.MYSQL_USER)}:${encodeURIComponent(process.env.MYSQL_PASSWORD)}@${process.env.MYSQL_HOST}:${process.env.MYSQL_PORT}/${process.env.MYSQL_DATABASE}`)') +/usr/bin/node "$new/node_modules/prisma/build/index.js" migrate deploy --schema "$new/prisma/schema.prisma" +unset DATABASE_URL MYSQL_PASSWORD +install -d -o syslog -g lisglosips -m 0750 /var/log/lisglosips +touch /var/log/lisglosips/caller-analytics.log +chown syslog:lisglosips /var/log/lisglosips/caller-analytics.log +chmod 0640 /var/log/lisglosips/caller-analytics.log +install -m 0644 "$new/infra/server-b/caller-analytics/rsyslog.conf" /etc/rsyslog.d/35-caller-analytics.conf +rsyslogd -N1 +if ! test -f /etc/lisglosips/caller-analytics.env; then + grep -E '^(DATABASE_URL|REDIS_URL)=' /etc/lisglosips/cdr-worker.env > /etc/lisglosips/caller-analytics.env + printf '%s\n' 'LISGLOSIPS_ENTRYPOINT=apps/worker-cdr/dist/analytics-main.js' 'LISGLOSIPS_SERVICE_NAME=caller-analytics' "ANALYTICS_CAPTURE_SINCE=$(date -u +%Y-%m-%dT%H:%M:%SZ)" >> /etc/lisglosips/caller-analytics.env + chown root:lisglosips /etc/lisglosips/caller-analytics.env + chmod 0640 /etc/lisglosips/caller-analytics.env +fi +chown -R root:lisglosips "$new/apps/api/dist" "$new/apps/worker-cdr/dist" "$new/packages/database/dist" "$new/packages/auth/dist" "$new/public" "$new/scripts" "$new/infra/server-b/caller-analytics" +chmod -R u=rwX,g=rX,o= "$new/apps/api/dist" "$new/apps/worker-cdr/dist" "$new/packages/database/dist" "$new/packages/auth/dist" "$new/public" "$new/scripts" "$new/infra/server-b/caller-analytics" +changed=1 +systemctl restart rsyslog +install -o root -g root -m 0644 "$new/opensips-candidate.cfg" /etc/opensips/opensips.cfg +ln -sfn "$new" /opt/lisglosips/current +systemctl restart opensips lisglosips@api +systemctl enable --now lisglosips@caller-analytics +systemctl restart lisglosips@caller-analytics +for i in $(seq 1 15); do + if curl -fsS http://127.0.0.1:3000/api/v2/health/ready >/dev/null; then break; fi + sleep 1 +done +curl -fsS http://127.0.0.1:3000/api/v2/health/ready +systemctl is-active opensips lisglosips@api lisglosips@cdr-worker lisglosips@caller-analytics +test "$(cat /opt/lisglosips/current/.deployed-commit)" = "$commit" +changed=0 +trap - ERR +printf '\nRELEASE=%s\nBACKUP=%s\nCOMMIT=%s\n' "$new" "$backup" "$commit" diff --git a/infra/server-b/caller-analytics/rsyslog.conf b/infra/server-b/caller-analytics/rsyslog.conf new file mode 100644 index 0000000..9b722bf --- /dev/null +++ b/infra/server-b/caller-analytics/rsyslog.conf @@ -0,0 +1,8 @@ +# OpenSIPS emits only analytics observations here. No network destination. +if ($programname == 'opensips' and $msg contains 'CRA1|') then { + action(type="omfile" file="/var/log/lisglosips/caller-analytics.log" + fileOwner="syslog" fileGroup="lisglosips" fileCreateMode="0640" + queue.type="LinkedList" queue.filename="caller-analytics" + queue.maxDiskSpace="256m" queue.saveOnShutdown="on" + action.resumeRetryCount="-1") +} diff --git a/packages/auth/src/index.ts b/packages/auth/src/index.ts index d2f2b35..7c88d49 100644 --- a/packages/auth/src/index.ts +++ b/packages/auth/src/index.ts @@ -4,6 +4,8 @@ export const ACCESS_TOKEN_TYPE = 'Bearer'; export const PASSWORD_ALGO_ARGON2ID = 'argon2id'; export type PermissionKey = + | 'caller_analytics.view' + | 'caller_analytics.view_all' | 'dashboard.view' | 'active_calls.view' | 'active_calls.manage' diff --git a/packages/database/src/caller-analytics-model.ts b/packages/database/src/caller-analytics-model.ts new file mode 100644 index 0000000..f14a9f2 --- /dev/null +++ b/packages/database/src/caller-analytics-model.ts @@ -0,0 +1,130 @@ +import { createHash } from 'node:crypto'; + +export const ANALYTICS_VERSION = 'caller-analytics-v1.1'; +export const ANALYTICS_STREAM = 'stream:caller_analytics'; +export const ANALYTICS_GROUP = 'caller-analytics-workers'; +export type AnalyticsKind = 'START' | 'ATTEMPT' | 'PROGRESS' | 'ACCEPTED' | 'DURATION' | 'END' | 'RECONCILE' | 'UNKNOWN'; +export interface AnalyticsEvent { + eventId: string; callUid: string; callId: string; attemptId: string; + kind: AnalyticsKind; at: number; startedAt: number; + customerId: string; customerGatewayId: string; vendorId: string; vendorGatewayId: string; + caller: string; landingCaller: string; callee: string; city: string; carrier: string; + code: number; talkMs: number | null; source: string; sequence: number; +} +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; +} +export interface AnalyticsCall { start: AnalyticsEvent; legs: Record; } +export interface AnalyticsCounts { + totalCalls: number; connectedCalls: number; answeredCalls: number; + failedCalls: number; pendingCalls: number; unknownCalls: number; activeCalls: number; talkMs: number; +} +export const COUNT_KEYS = ['totalCalls', 'connectedCalls', 'answeredCalls', 'failedCalls', 'pendingCalls', 'unknownCalls', 'activeCalls', 'talkMs'] as const; +export const emptyCounts = (): AnalyticsCounts => ({ totalCalls: 0, connectedCalls: 0, answeredCalls: 0, failedCalls: 0, pendingCalls: 0, unknownCalls: 0, activeCalls: 0, talkMs: 0 }); +export function analyticsHash(value: unknown): string { return createHash('sha256').update(JSON.stringify(value)).digest('hex'); } + +export function parseAnalyticsEvent(fields: string[]): AnalyticsEvent { + const f: Record = {}; + for (let i = 0; i < fields.length; i += 2) f[fields[i]] = fields[i + 1]; + const text = (key: string, max = 255): string => { + const value = f[key] ?? ''; + if (value.length > max || [...value].some(c => c.charCodeAt(0) < 32)) throw new Error(`Invalid ${key}`); + return value === 'none' ? '' : value; + }; + const num = (key: string, fallback = 0): number => { + const value = f[key] === undefined ? fallback : Number(f[key]); + if (!Number.isSafeInteger(value) || value < 0) throw new Error(`Invalid ${key}`); + return value; + }; + const kind = text('kind') as AnalyticsKind; + if (f.version !== '1' || !['START', 'ATTEMPT', 'PROGRESS', 'ACCEPTED', 'DURATION', 'END', 'RECONCILE', 'UNKNOWN'].includes(kind)) throw new Error('Invalid analytics schema'); + const at = num('at'); const startedAt = num('startedAt'); + if (at < 1_700_000_000_000 || at > Date.now() + 60_000 || startedAt > at || startedAt < 1_700_000_000_000) throw new Error('Invalid event time'); + const event: AnalyticsEvent = { + eventId: text('eventId', 255), callUid: text('callUid', 255), callId: text('callId'), attemptId: text('attemptId', 80), + kind, at, startedAt, customerId: text('customerId', 32), customerGatewayId: text('customerGatewayId', 32), + vendorId: text('vendorId', 32), vendorGatewayId: text('vendorGatewayId', 32), + caller: text('caller', 64), landingCaller: text('landingCaller', 64), callee: text('callee', 64), + city: text('city', 12), carrier: text('carrier', 20), code: num('code'), + talkMs: f.talkMs === undefined || f.talkMs === '-1' ? null : num('talkMs'), source: text('source', 64), sequence: num('sequence') + }; + if (!event.eventId || !event.callUid || !event.callId || !event.customerId || !event.caller) throw new Error('Missing call identity'); + if (event.code > 699 || (kind === 'PROGRESS' && ![100, 180, 181, 182, 183].includes(event.code))) throw new Error('Invalid response code'); + if (kind === 'RECONCILE' && event.source !== 'verified-dialog-final') throw new Error('Unverified CDR cannot reconcile'); + return event; +} + +const earliest = (a: number | null, b: number): number => a === null ? b : Math.min(a, b); +export function foldAnalytics(call: AnalyticsCall | null, event: AnalyticsEvent): AnalyticsCall { + const result: AnalyticsCall = call ? structuredClone(call) : { start: event, legs: {} }; + if (result.start.customerId !== event.customerId || result.start.caller !== event.caller) throw new Error('Conflicting call identity'); + if (event.startedAt < result.start.startedAt || (event.kind === 'START' && event.startedAt === result.start.startedAt)) result.start = event; + if (event.kind === 'START') return result; + const key = event.attemptId || 'platform'; + const leg: AnalyticsLeg = result.legs[key] ?? { + 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 === '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 === '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.kind === 'ACCEPTED' && leg.endedAt !== null && !leg.durationFinal) { + leg.talkMs = Math.max(0, leg.endedAt - event.at); leg.durationAt = leg.endedAt; leg.durationFinal = true; + } + 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; + } + } + result.legs[key] = leg; + return result; +} + +export interface AnalyticsProjection extends AnalyticsCounts { + id: string; callKey: string; view: 'original' | 'landing'; startedAt: number; + customerId: string; customerGatewayId: string; vendorId: string; vendorGatewayId: string; + caller: string; callee: string; city: string; carrier: string; payload: AnalyticsCall | AnalyticsLeg; +} +function legCounts(leg: AnalyticsLeg): AnalyticsCounts { + const connected = leg.first180At !== null || leg.first183At !== null || leg.acceptedAt !== null || leg.talkMs > 0; + const unknown = !connected && leg.unknown; + return { totalCalls: 1, connectedCalls: +connected, answeredCalls: +(leg.talkMs > 0), failedCalls: +(!connected && !unknown && leg.endedAt !== null), pendingCalls: +(!connected && !unknown && leg.endedAt === null), unknownCalls: +unknown, activeCalls: +(leg.endedAt === null), talkMs: leg.talkMs }; +} +export function projectAnalytics(call: AnalyticsCall): AnalyticsProjection[] { + const e = call.start; const key = analyticsHash(e.callUid); + const legs = Object.values(call.legs); const counts = legs.map(legCounts); + const connected = counts.some(c => c.connectedCalls > 0); const active = !legs.length || counts.some(c => c.activeCalls > 0); + const unknown = !connected && counts.some(c => c.unknownCalls > 0); + const base = { customerId: e.customerId, customerGatewayId: e.customerGatewayId, callee: e.callee, city: e.city, carrier: e.carrier, callKey: key }; + return [{ ...base, id: analyticsHash([key, 'original']), view: 'original', startedAt: e.startedAt, caller: e.caller, + vendorId: '', vendorGatewayId: '', payload: call, totalCalls: 1, connectedCalls: +connected, + answeredCalls: +(counts.some(c => c.answeredCalls > 0)), failedCalls: +(!connected && !unknown && !active), + pendingCalls: +(!connected && !unknown && active), unknownCalls: +unknown, activeCalls: +active, + talkMs: counts.reduce((sum, c) => sum + c.talkMs, 0) + }, ...legs.filter(l => l.event.attemptId && l.event.vendorGatewayId).map(leg => ({ + ...base, ...legCounts(leg), id: analyticsHash([key, leg.event.attemptId]), view: 'landing' as const, + startedAt: leg.event.startedAt, caller: leg.event.landingCaller || leg.event.caller, + vendorId: leg.event.vendorId, vendorGatewayId: leg.event.vendorGatewayId, payload: leg + }))]; +} + +export function analyticsRates(counts: AnalyticsCounts) { + const ratio = (n: number, d: number) => d ? Number((n / d * 100).toFixed(4)) : null; + return { ...counts, notConnectedCalls: counts.failedCalls + counts.pendingCalls, + connectionRate: counts.unknownCalls ? null : ratio(counts.connectedCalls, counts.totalCalls), + overallAnswerRate: ratio(counts.answeredCalls, counts.totalCalls), + connectedAnswerRate: counts.unknownCalls ? null : ratio(counts.answeredCalls, counts.connectedCalls) }; +} diff --git a/packages/database/src/caller-analytics-store.ts b/packages/database/src/caller-analytics-store.ts new file mode 100644 index 0000000..9780a94 --- /dev/null +++ b/packages/database/src/caller-analytics-store.ts @@ -0,0 +1,52 @@ +import { Prisma, PrismaClient } from '@prisma/client'; +import { analyticsHash, foldAnalytics, projectAnalytics, COUNT_KEYS, type AnalyticsCall, type AnalyticsEvent, type AnalyticsProjection } from './caller-analytics-model.js'; + +type Tx = Prisma.TransactionClient; +const json = (value: T | string): T => typeof value === 'string' ? JSON.parse(value) as T : value; +export class CallerAnalyticsStore { + constructor(private readonly db: PrismaClient) {} + + async consume(event: AnalyticsEvent): Promise<'processed' | 'duplicate'> { + // The legacy router puts diagnostic markers in customer_id for unauthenticated traffic. + // Such traffic is outside the business population and must not create analysis calls. + if (!await this.db.customer.findUnique({ where: { id: event.customerId }, select: { id: true } })) return 'duplicate'; + return this.db.$transaction(async tx => { + const callKey = analyticsHash(event.callUid); + await tx.$executeRaw`INSERT IGNORE INTO caller_analysis_calls (id, payload, updated_at) VALUES (${callKey}, NULL, NOW(3))`; + const rows = await tx.$queryRaw>`SELECT payload FROM caller_analysis_calls WHERE id=${callKey} FOR UPDATE`; + const eventKey = analyticsHash(event.eventId); + const previous = await tx.$queryRaw>`SELECT id FROM caller_analysis_events WHERE id=${eventKey}`; + if (previous.length) return 'duplicate'; + const old = rows[0].payload ? json(rows[0].payload) : null; + const next = foldAnalytics(old, event); + const before = new Map((old ? projectAnalytics(old) : []).map(row => [row.id, row])); + for (const row of projectAnalytics(next)) { + const prior = before.get(row.id); + if (prior) await this.bucket(tx, prior, -1); + await this.bucket(tx, row, 1); + await this.state(tx, row); + } + await tx.$executeRaw`UPDATE caller_analysis_calls SET payload=${JSON.stringify(next)}, updated_at=NOW(3) WHERE id=${callKey}`; + await tx.$executeRaw`INSERT INTO caller_analysis_events (id, call_key, occurred_at, payload, created_at) VALUES (${eventKey}, ${callKey}, ${new Date(event.at)}, ${JSON.stringify(event)}, NOW(3))`; + return 'processed'; + }, { timeout: 15000 }); + } + + private async state(tx: Tx, row: AnalyticsProjection) { + await tx.$executeRaw`INSERT INTO caller_analysis_states + (id, call_key, view, started_at, customer_id, customer_gateway_id, vendor_id, vendor_gateway_id, caller, callee, city, carrier, counts, payload, updated_at) + VALUES (${row.id}, ${row.callKey}, ${row.view}, ${new Date(row.startedAt)}, ${row.customerId}, ${row.customerGatewayId}, ${row.vendorId}, ${row.vendorGatewayId}, ${row.caller}, ${row.callee}, ${row.city}, ${row.carrier}, ${JSON.stringify(Object.fromEntries(COUNT_KEYS.map(k => [k, row[k]])))}, ${JSON.stringify(row.payload)}, NOW(3)) + ON DUPLICATE KEY UPDATE started_at=VALUES(started_at), customer_gateway_id=VALUES(customer_gateway_id), vendor_id=VALUES(vendor_id), vendor_gateway_id=VALUES(vendor_gateway_id), caller=VALUES(caller), callee=VALUES(callee), city=VALUES(city), carrier=VALUES(carrier), counts=VALUES(counts), payload=VALUES(payload), updated_at=NOW(3)`; + } + + private async bucket(tx: Tx, row: AnalyticsProjection, sign: number) { + 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]}`])); + 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}) + ON DUPLICATE KEY UPDATE counts=JSON_OBJECT(${additions})`); + } +} diff --git a/packages/database/src/caller-analytics.spec.ts b/packages/database/src/caller-analytics.spec.ts new file mode 100644 index 0000000..507ffe6 --- /dev/null +++ b/packages/database/src/caller-analytics.spec.ts @@ -0,0 +1,49 @@ +import { describe, expect, it } from 'vitest'; +import { analyticsRates, foldAnalytics, projectAnalytics, type AnalyticsEvent, type AnalyticsCall } from './caller-analytics-model.js'; +const start = Date.UTC(2026, 7, 31, 0); +const event = (kind: AnalyticsEvent['kind'], extra: Partial = {}): AnalyticsEvent => ({ eventId: kind, callUid: 'u1', callId: 'c1', attemptId: '1', kind, at: start + 1000, startedAt: start, customerId: 'c', customerGatewayId: 'cg', vendorId: 'v', vendorGatewayId: 'vg', caller: '123', landingCaller: '456', callee: '789', city: 'city', carrier: 'UNKNOWN', code: 0, talkMs: null, source: 'opensips', sequence: 0, ...extra }); +function run(events: AnalyticsEvent[]) { return projectAnalytics(events.reduce((c, e) => foldAnalytics(c, e), null)!); } +describe('caller analytics business state', () => { + it.each([180,183])('counts %i as connected even if final call fails', code => { + const [r] = run([event('ATTEMPT'), event('PROGRESS', { code }), event('END', { code: 487, source: 'verified-dialog-final' })]); + expect(analyticsRates(r)).toMatchObject({ totalCalls: 1, connectedCalls: 1, answeredCalls: 0, failedCalls: 0, notConnectedCalls: 0, connectionRate: 100, overallAnswerRate: 0, connectedAnswerRate: 0 }); + }); + it.each([100,181,182])('does not treat %i as connected', code => { + const [r] = run([event('ATTEMPT'), event('PROGRESS', { code }), event('END', { code: 486, source: 'verified-dialog-final' })]); + expect(analyticsRates(r)).toMatchObject({ connectedCalls: 0, failedCalls: 1, connectedAnswerRate: null }); + }); + it('counts positive sub-second duration, not just 200', () => { + const [zero] = run([event('ATTEMPT'), event('ACCEPTED', { code: 200 })]); + expect(zero.answeredCalls).toBe(0); + const [positive] = run([event('ATTEMPT'), event('ACCEPTED', { code: 200 }), event('END', { at: start + 1500, source: 'verified-dialog-final' })]); + expect(positive).toMatchObject({ connectedCalls: 1, answeredCalls: 1, talkMs: 500 }); + }); + it('preserves zero actual duration', () => { + const [r] = run([event('ATTEMPT'), event('ACCEPTED', { code: 200 }), event('END', { talkMs: 0, source: 'verified-dialog-final' })]); + expect(r.answeredCalls).toBe(0); + }); + it('is insensitive to duplicate progress and delayed attempt', () => { + const e = event('PROGRESS', { code: 183 }); + const [r] = run([e, e, event('END', { source: 'verified-dialog-final', code: 487 }), event('ATTEMPT')]); + expect(r).toMatchObject({ totalCalls: 1, connectedCalls: 1, answeredCalls: 0, activeCalls: 0 }); + }); + it('corrects provisional duration and rejects late provisional overwrite', () => { + const [r] = run([event('ATTEMPT'), event('ACCEPTED', { code: 200 }), event('DURATION', { at: start + 2000, talkMs: 1000 }), event('RECONCILE', { at: start + 3000, talkMs: 0, source: 'verified-dialog-final' }), event('DURATION', { at: start + 4000, talkMs: 3000 })]); + expect(r).toMatchObject({ connectedCalls: 1, answeredCalls: 0, talkMs: 0 }); + }); + it('keeps original call separate from multiple landing attempts', () => { + const rows = run([event('ATTEMPT'), event('END', { code: 486 }), event('ATTEMPT', { attemptId: '2' }), event('PROGRESS', { attemptId: '2', code: 180 })]); + expect(rows).toHaveLength(3); expect(rows[0]).toMatchObject({ totalCalls: 1, connectedCalls: 1 }); + expect(rows.slice(1).reduce((n,r) => n+r.totalCalls,0)).toBe(2); + }); + it('distinguishes pending and unknown', () => { + expect(run([event('ATTEMPT')])[0]).toMatchObject({ pendingCalls: 1, failedCalls: 0 }); + expect(run([event('UNKNOWN')])[0]).toMatchObject({ unknownCalls: 1, failedCalls: 0, pendingCalls: 0 }); + }); + it('does not merge different customer identities', () => { + expect(() => run([event('ATTEMPT'), event('PROGRESS', { customerId: 'other', code: 180 })])).toThrow('Conflicting'); + }); + 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 }); + }); +}); diff --git a/packages/database/src/index.ts b/packages/database/src/index.ts index a03832b..be302f5 100644 --- a/packages/database/src/index.ts +++ b/packages/database/src/index.ts @@ -1,3 +1,6 @@ export { Prisma, PrismaClient } from '@prisma/client'; export const DATABASE_PROVIDER = 'mysql'; + +export * from './caller-analytics-model.js'; +export * from './caller-analytics-store.js'; diff --git a/prisma/migrations/20260831090000_caller_analytics/migration.sql b/prisma/migrations/20260831090000_caller_analytics/migration.sql new file mode 100644 index 0000000..8e4417d --- /dev/null +++ b/prisma/migrations/20260831090000_caller_analytics/migration.sql @@ -0,0 +1,43 @@ +CREATE TABLE caller_analysis_calls ( + id CHAR(64) NOT NULL PRIMARY KEY, payload JSON NULL, updated_at DATETIME(3) NOT NULL +); +CREATE TABLE caller_analysis_events ( + id CHAR(64) NOT NULL PRIMARY KEY, call_key CHAR(64) NOT NULL, occurred_at DATETIME(3) NOT NULL, + payload JSON NOT NULL, created_at DATETIME(3) NOT NULL, + INDEX caller_events_call (call_key, occurred_at), INDEX caller_events_retention (created_at) +); +CREATE TABLE caller_analysis_states ( + id CHAR(64) NOT NULL PRIMARY KEY, call_key CHAR(64) NOT NULL, view VARCHAR(16) NOT NULL, + started_at DATETIME(3) NOT NULL, customer_id VARCHAR(32) NOT NULL, customer_gateway_id VARCHAR(32) NOT NULL, + vendor_id VARCHAR(32) NOT NULL, vendor_gateway_id VARCHAR(32) NOT NULL, caller VARCHAR(64) NOT NULL, + callee VARCHAR(64) NOT NULL, city VARCHAR(12) NOT NULL, carrier VARCHAR(20) NOT NULL, + counts JSON NOT NULL, payload JSON NOT NULL, updated_at DATETIME(3) NOT NULL, + INDEX caller_states_window (view, started_at), INDEX caller_states_customer (customer_id, view, started_at), + INDEX caller_states_number (view, caller, started_at), INDEX caller_states_call (call_key) +); +CREATE TABLE caller_analysis_minute_buckets ( + id CHAR(64) NOT NULL PRIMARY KEY, started_at DATETIME(3) NOT NULL, view VARCHAR(16) NOT NULL, + customer_id VARCHAR(32) NOT NULL, customer_gateway_id VARCHAR(32) NOT NULL, + vendor_id VARCHAR(32) NOT NULL, vendor_gateway_id VARCHAR(32) NOT NULL, caller VARCHAR(64) NOT NULL, + city VARCHAR(12) NOT NULL, carrier VARCHAR(20) NOT NULL, counts JSON NOT NULL, + INDEX caller_buckets_window (view, started_at), INDEX caller_buckets_customer (customer_id, view, started_at), + INDEX caller_buckets_number (view, caller, started_at) +); +CREATE TABLE caller_analysis_access ( + user_id VARCHAR(32) NOT NULL, customer_id VARCHAR(32) NOT NULL, PRIMARY KEY(user_id, customer_id) +); +CREATE TABLE caller_analysis_health ( + id VARCHAR(64) NOT NULL PRIMARY KEY, payload JSON NOT NULL, updated_at DATETIME(3) NOT NULL +); +CREATE TABLE caller_analysis_alerts ( + id CHAR(64) NOT NULL PRIMARY KEY, customer_id VARCHAR(32) NOT NULL, view VARCHAR(16) NOT NULL, + caller VARCHAR(64) NOT NULL, status VARCHAR(20) NOT NULL, payload JSON NOT NULL, + updated_at DATETIME(3) NOT NULL, INDEX caller_alerts_customer (customer_id, status, updated_at) +); +INSERT INTO permissions (id, module, action, description, created_at, updated_at) +VALUES ('caller_analytics.view', 'caller_analytics', 'view', '查看授权客户的主叫分析', NOW(3), NOW(3)), + ('caller_analytics.view_all', 'caller_analytics', 'view_all', '查看全部客户的主叫分析', NOW(3), NOW(3)); +INSERT INTO role_permissions (role_id, permission_id, created_at) +SELECT id, 'caller_analytics.view', NOW(3) FROM roles WHERE id='ROLE_SUPER_ADMIN'; +INSERT INTO role_permissions (role_id, permission_id, created_at) +SELECT id, 'caller_analytics.view_all', NOW(3) FROM roles WHERE id='ROLE_SUPER_ADMIN'; diff --git a/prisma/schema.prisma b/prisma/schema.prisma index be7f9e6..f0caf9c 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -960,6 +960,91 @@ model OutboxEvent { @@map("outbox_events") } +model CallerAnalysisCall { + id String @id @db.Char(64) + payload Json? + updatedAt DateTime @map("updated_at") @db.DateTime(3) + @@map("caller_analysis_calls") +} + +model CallerAnalysisEvent { + id String @id @db.Char(64) + callKey String @map("call_key") @db.Char(64) + occurredAt DateTime @map("occurred_at") @db.DateTime(3) + payload Json + createdAt DateTime @map("created_at") @db.DateTime(3) + @@index([callKey, occurredAt], map: "caller_events_call") + @@index([createdAt], map: "caller_events_retention") + @@map("caller_analysis_events") +} + +model CallerAnalysisState { + id String @id @db.Char(64) + callKey String @map("call_key") @db.Char(64) + view String @db.VarChar(16) + startedAt DateTime @map("started_at") @db.DateTime(3) + customerId String @map("customer_id") @db.VarChar(32) + customerGatewayId String @map("customer_gateway_id") @db.VarChar(32) + vendorId String @map("vendor_id") @db.VarChar(32) + vendorGatewayId String @map("vendor_gateway_id") @db.VarChar(32) + caller String @db.VarChar(64) + callee String @db.VarChar(64) + city String @db.VarChar(12) + carrier String @db.VarChar(20) + counts Json + payload Json + updatedAt DateTime @map("updated_at") @db.DateTime(3) + @@index([view, startedAt], map: "caller_states_window") + @@index([customerId, view, startedAt], map: "caller_states_customer") + @@index([view, caller, startedAt], map: "caller_states_number") + @@index([callKey], map: "caller_states_call") + @@map("caller_analysis_states") +} + +model CallerAnalysisMinuteBucket { + id String @id @db.Char(64) + startedAt DateTime @map("started_at") @db.DateTime(3) + view String @db.VarChar(16) + customerId String @map("customer_id") @db.VarChar(32) + customerGatewayId String @map("customer_gateway_id") @db.VarChar(32) + vendorId String @map("vendor_id") @db.VarChar(32) + vendorGatewayId String @map("vendor_gateway_id") @db.VarChar(32) + caller String @db.VarChar(64) + city String @db.VarChar(12) + carrier String @db.VarChar(20) + counts Json + @@index([view, startedAt], map: "caller_buckets_window") + @@index([customerId, view, startedAt], map: "caller_buckets_customer") + @@index([view, caller, startedAt], map: "caller_buckets_number") + @@map("caller_analysis_minute_buckets") +} + +model CallerAnalysisAccess { + userId String @map("user_id") @db.VarChar(32) + customerId String @map("customer_id") @db.VarChar(32) + @@id([userId, customerId]) + @@map("caller_analysis_access") +} + +model CallerAnalysisHealth { + id String @id @db.VarChar(64) + payload Json + updatedAt DateTime @map("updated_at") @db.DateTime(3) + @@map("caller_analysis_health") +} + +model CallerAnalysisAlert { + id String @id @db.Char(64) + customerId String @map("customer_id") @db.VarChar(32) + view String @db.VarChar(16) + caller String @db.VarChar(64) + status String @db.VarChar(20) + payload Json + updatedAt DateTime @map("updated_at") @db.DateTime(3) + @@index([customerId, status, updatedAt], map: "caller_alerts_customer") + @@map("caller_analysis_alerts") +} + model IdempotencyKey { id String @id @db.VarChar(40) key String @unique @db.VarChar(128) diff --git a/prisma/seed.ts b/prisma/seed.ts index ae0f659..42c8067 100644 --- a/prisma/seed.ts +++ b/prisma/seed.ts @@ -3,6 +3,8 @@ import { PrismaClient } from '@prisma/client'; const prisma = new PrismaClient(); const permissions = [ + ['caller_analytics.view', 'caller_analytics', 'view', '查看授权客户的主叫分析'], + ['caller_analytics.view_all', 'caller_analytics', 'view_all', '查看全部客户的主叫分析'], ['dashboard.view', 'dashboard', 'view', '查看 Dashboard 指标'], ['active_calls.view', 'active_calls', 'view', '查看当前通话'], ['active_calls.manage', 'active_calls', 'manage', '强制挂断当前通话'], diff --git a/scripts/build-release-artifact.mjs b/scripts/build-release-artifact.mjs index bac3a5b..520d1ac 100644 --- a/scripts/build-release-artifact.mjs +++ b/scripts/build-release-artifact.mjs @@ -74,6 +74,8 @@ const copyEntries = [ ['infra/server-b/s56', 'infra/server-b/s56', true], ['infra/server-b/s57', 'infra/server-b/s57', true], ['infra/server-b/s58', 'infra/server-b/s58', true], + ['infra/server-b/caller-analytics', 'infra/server-b/caller-analytics', true], + ['scripts/instrument-caller-analytics.mjs', 'scripts/instrument-caller-analytics.mjs', true], ['infra/ops/startup/lisglosips-start-b.sh', 'infra/ops/startup/lisglosips-start-b.sh', true], ['infra/server-a/s28/lisglosips_hotpath.lua', 'infra/server-a/s28/lisglosips_hotpath.lua', true], ['scripts/phase2-gateway-migration.mjs', 'scripts/phase2-gateway-migration.mjs', true], diff --git a/scripts/instrument-caller-analytics.mjs b/scripts/instrument-caller-analytics.mjs new file mode 100644 index 0000000..c802963 --- /dev/null +++ b/scripts/instrument-caller-analytics.mjs @@ -0,0 +1,42 @@ +// Add observational hooks to the verified S28 routing layout, without replacing routing or billing. +import { readFileSync, writeFileSync } from 'node:fs'; +const [source, output] = process.argv.slice(2); +if (!source || !output) throw new Error('Usage: node instrument-caller-analytics.mjs source.cfg output.cfg'); +let cfg = readFileSync(source, 'utf8').replaceAll('\r\n', '\n').replace(/^\uFEFF/, ''); +if (cfg.includes('CRA1|')) throw new Error('Already instrumented; use the recorded pre-analytics backup'); +function insert(marker, extra, after = false) { + if (cfg.split(marker).length !== 2) throw new Error(`Unsupported routing layout: ${marker}`); + cfg = cfg.replace(marker, after ? marker + extra : extra + marker); +} +const metadata = { + uid: '$ci + "~" + $ft + "~" + $TS', callid: '$ci', customer: '$var(s28_customer_id)', + gateway: '$var(s28_gateway_id)', vendor: '$var(s28_vendor_id)', vendorGateway: '$var(s28_vendor_gateway_id)', + caller: '$fU', landing: '$var(s28_landing_caller)', callee: '$var(s28_real_callee)', + city: '$var(s28_callee_city_code)', carrier: '$var(s28_callee_operator)', startsec: '$Ts', startmicro: '$Tsm', attempt: '"none"' +}; +const init = Object.entries(metadata).map(([k,v]) => ` $avp(cra_${k}) = ${v};`).join('\n'); +insert(' if ($var(s28_decision) != "allow") {', `${init}\n $var(cra_kind)="START"; $var(cra_code)=0; $var(cra_source)="opensips"; route(CRA_EMIT_AVP);\n\n`); +const bind = Object.keys(metadata).map(k => ` $dlg_val(cra_${k}) = $avp(cra_${k});`).join('\n'); +insert(' $dlg_val(s28_customer_id) = $var(s28_customer_id);', `${bind}\n dlg_on_hangup("CRA_HANGUP");\n dlg_on_timeout("CRA_TIMEOUT");\n`); +insert(' t_on_reply("S28_REPLY");', ' $dlg_val(cra_attempt)="1"; $avp(cra_attempt)="1";\n $var(cra_kind)="ATTEMPT"; $var(cra_code)=0; $var(cra_source)="opensips"; route(CRA_EMIT_DLG);\n t_on_failure("CRA_FAILURE");\n'); +insert('onreply_route[S28_REPLY] {', '\n if ($rs==180 || $rs==183 || $rs==100 || $rs==181 || $rs==182) {\n $var(cra_kind)="PROGRESS"; $var(cra_code)=$rs; $var(cra_source)="opensips"; route(CRA_EMIT_DLG);\n }\n if ($rs>=200 && $rs<300) {\n $var(cra_kind)="ACCEPTED"; $var(cra_code)=$rs; $var(cra_source)="opensips"; route(CRA_EMIT_DLG);\n }\n', true); +insert('route[S28_CDR_FAILURE] {', '\n $var(cra_kind)="END"; $var(cra_code)=$var(s28_reply_code); $var(cra_source)="platform-final"; route(CRA_EMIT_AVP);\n', true); +const fields = ['uid', 'callid', 'customer', 'gateway', 'vendor', 'vendorGateway', 'caller', 'landing', 'callee', 'city', 'carrier', 'attempt']; +function emitter(type, prefix) { + const ref = k => `$${prefix}(cra_${k})`; + const encoded = fields.map(k => `$(${prefix}(cra_${k}){s.encode.base64})`).join('|'); + return `\nroute[CRA_EMIT_${type}] {\n if (${ref('customer')}==NULL || ${ref('customer')}=="none") return;\n xlog("L_NOTICE", "CRA1|$var(cra_kind)|$Ts|$Tsm|${ref('startsec')}|${ref('startmicro')}|$var(cra_code)|$var(cra_source)|${encoded}\\n");\n}\n`; +} +cfg += emitter('AVP', 'avp') + emitter('DLG', 'dlg_val') + ` +route[CRA_HANGUP] { + $var(cra_kind)="END"; $var(cra_code)=200; $var(cra_source)="verified-dialog-final"; route(CRA_EMIT_DLG); +} +route[CRA_TIMEOUT] { + $var(cra_kind)="END"; $var(cra_code)=408; $var(cra_source)="verified-dialog-final"; route(CRA_EMIT_DLG); +} +failure_route[CRA_FAILURE] { + $var(cra_kind)="END"; $var(cra_code)=$T_reply_code; $var(cra_source)="verified-dialog-final"; route(CRA_EMIT_DLG); +} +`; +writeFileSync(output, cfg, 'utf8'); +console.log('Analytics observation hooks generated; OpenSIPS config validation is required before activation.'); diff --git a/tests/api/caller-analytics-real.mjs b/tests/api/caller-analytics-real.mjs new file mode 100644 index 0000000..4f7b9e9 --- /dev/null +++ b/tests/api/caller-analytics-real.mjs @@ -0,0 +1,58 @@ +// Run on test B from a staged release, with api.env loaded. No public telephone destinations. +import { PrismaClient } from '@prisma/client'; +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'; +const db = new PrismaClient(); +const fixture = 'cra_20260831'; +const user = await db.user.findFirst({ where:{ status:'ENABLED', roles:{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}`} : {}}); + return {status:response.status,body:await response.json()}; +}; +try { + if (process.argv.includes('--setup')) { + assert.equal(await db.customerGateway.count({where:{sourceIp:'127.0.0.1',id:{not:fixture}}}),0,'Loopback gateway already in use'); + await db.$transaction(async tx => { + await tx.customer.upsert({where:{id:fixture},create:{id:fixture,name:'主叫分析受控测试客户',balance:100,creditLimit:100},update:{status:'ENABLED'}}); + await tx.vendor.upsert({where:{id:fixture},create:{id:fixture,name:'主叫分析本机模拟供应商'},update:{status:'ENABLED'}}); + await tx.vendorGateway.upsert({where:{id:fixture},create:{id:fixture,vendorId:fixture,name:'主叫分析本机模拟落地',authMode:'IP',host:'127.0.0.1',port:50630,cycleRate:0},update:{status:'ENABLED'}}); + await tx.landingLineGroup.upsert({where:{id:fixture},create:{id:fixture,name:'主叫分析本机测试线路组'},update:{status:'ENABLED'}}); + await tx.landingLineGroupItem.upsert({where:{id:fixture},create:{id:fixture,lineGroupId:fixture,vendorGatewayId:fixture,priority:1},update:{status:'ENABLED'}}); + await tx.customerGateway.upsert({where:{id:fixture},create:{id:fixture,customerId:fixture,name:'主叫分析本机测试入口',authMode:'IP',sourceIp:'127.0.0.1',lineGroupId:fixture,cycleRate:0},update:{status:'ENABLED'}}); + await tx.outboxEvent.create({data:{id:`cra_${Date.now()}`,aggregateType:'customer_gateway_config',aggregateId:fixture,eventType:'CONFIG_CHANGED',payload:{reason:'isolated caller analytics test fixture'}}}); + }); + console.log('Loopback-only zero-rate fixture prepared; waiting for normal config publisher.'); + } else if (process.argv.includes('--disable')) { + await db.customerGateway.update({where:{id:fixture},data:{status:'DISABLED'}}); + await db.vendorGateway.update({where:{id:fixture},data:{status:'DISABLED'}}); + await db.outboxEvent.create({data:{id:`cra_${Date.now()}`,aggregateType:'customer_gateway_config',aggregateId:fixture,eventType:'CONFIG_CHANGED',payload:{reason:'disable isolated test fixture'}}}); + console.log('Isolated test ingress and egress disabled; evidence retained.'); + } else { + const from = new Date().toISOString(); + const sip = spawnSync('python3',['tests/api/caller-analytics-sip.py'],{encoding:'utf8',timeout:90000}); + assert.equal(sip.status,0,sip.stderr || sip.stdout); + console.log(sip.stdout); + let result; + for (let i=0;i<12;i++) { + result = await get(`/caller-analytics/overview?customerId=${fixture}&view=landing&from=${encodeURIComponent(from)}`); + if (result.status===200 && result.body.summary.totalCalls===6 && result.body.summary.activeCalls===0) break; + await new Promise(r=>setTimeout(r,1000)); + } + assert.equal(result.status,200,JSON.stringify(result.body)); + assert.deepEqual(Object.fromEntries(['totalCalls','connectedCalls','answeredCalls','failedCalls'].map(k=>[k,result.body.summary[k]])),{totalCalls:6,connectedCalls:4,answeredCalls:2,failedCalls:2}); + assert.equal((await get('/caller-analytics/overview',false)).status,401); + assert.equal((await get('/caller-analytics/overview?view=invalid')).status,400); + assert.equal(result.body.summary.connectedAnswerRate,50); + const detail = await get(`/caller-analytics/calls?customerId=${fixture}&caller=99100001&view=landing&from=${encodeURIComponent(from)}`); + 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']}; + 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})); + } +} finally { await db.$disconnect(); } diff --git a/tests/api/caller-analytics-sip.py b/tests/api/caller-analytics-sip.py new file mode 100644 index 0000000..53a226b --- /dev/null +++ b/tests/api/caller-analytics-sip.py @@ -0,0 +1,72 @@ +#!/usr/bin/env python3 +"""Isolated UDP SIP UAC/UAS on test B. No PSTN routing or remote targets.""" +import json +import re +import socket +import threading +import time +import uuid + +events = [] +stop = threading.Event() + +def header(message, name): + match = re.search(r'^' + re.escape(name) + r':\s*(.+)$', message, re.I | re.M) + return match.group(1).strip() if match else '' + +def reply(message, code, reason): + to = header(message, 'To') + if ';tag=' not in to: + to += ';tag=cra-uas' + vias = re.findall(r'^Via:\s*(.+)$', message, re.I | re.M) + routes = re.findall(r'^Record-Route:\s*(.+)$', message, re.I | re.M) + lines = [f'SIP/2.0 {code} {reason}'] + [f'Via: {v.strip()}' for v in vias] + lines += [f'From: {header(message,"From")}', f'To: {to}', f'Call-ID: {header(message,"Call-ID")}', f'CSeq: {header(message,"CSeq")}'] + lines += [f'Record-Route: {r.strip()}' for r in routes] + lines += ['Contact: ', 'Content-Length: 0', '', ''] + return '\r\n'.join(lines).encode() + +def uas(): + with socket.socket(socket.AF_INET, socket.SOCK_DGRAM) as sock: + sock.bind(('127.0.0.1',50630)); sock.settimeout(.2) + while not stop.is_set(): + try: data, addr = sock.recvfrom(65535) + except socket.timeout: continue + msg = data.decode(); method = msg.split()[0] + if method == 'INVITE': + number = re.search(r'INVITE sip:([^@]+)',msg).group(1) + cases = {'99101':[(180,'Ringing'),(486,'Busy Here')], '99102':[(183,'Session Progress'),(487,'Request Terminated')], '99103':[(486,'Busy Here')], '99104':[(200,'OK')], '99105':[(180,'Ringing'),(200,'OK')], '99106':[(100,'Trying'),(486,'Busy Here')]} + for code, reason in cases[number]: + sock.sendto(reply(msg,code,reason),addr) + events.append({'callId':header(msg,'Call-ID'),'direction':'UAS','code':code,'time':time.time()}) + time.sleep(.15 if number!='99106' else 1.5) + elif method == 'BYE': sock.sendto(reply(msg,200,'OK'),addr) + +thread = threading.Thread(target=uas,daemon=True); thread.start(); time.sleep(.3) +try: + for number in ['99101','99102','99103','99104','99105','99106']: + with socket.socket(socket.AF_INET,socket.SOCK_DGRAM) as sock: + sock.bind(('127.0.0.1',0)); sock.settimeout(8); port=sock.getsockname()[1] + callid=f'cra-{number}-{uuid.uuid4().hex}@loopback'; tag=uuid.uuid4().hex[:8] + branch='z9hG4bK'+uuid.uuid4().hex + base=[f'Via: SIP/2.0/UDP 127.0.0.1:{port};branch={branch};rport', f'From: ;tag={tag}', f'To: ', f'Call-ID: {callid}', f'Contact: ', 'Max-Forwards: 70'] + invite='\r\n'.join([f'INVITE sip:{number}@127.0.0.1:15060 SIP/2.0',*base,'CSeq: 1 INVITE','Content-Length: 0','','']) + sock.sendto(invite.encode(),('127.0.0.1',15060)) + while True: + msg=sock.recv(65535).decode(); code=int(msg.split()[1]); events.append({'callId':callid,'direction':'UAC','code':code,'time':time.time()}) + if code<200: continue + if code>=300: + ack=invite.replace('INVITE sip:','ACK sip:',1).replace('CSeq: 1 INVITE','CSeq: 1 ACK').replace(base[2],f'To: {header(msg,"To")}') + sock.sendto(ack.encode(),('127.0.0.1',15060)); break + routes = list(reversed(re.findall(r'^Record-Route:\s*(.+)$',msg,re.I|re.M))) + for method,cseq in [('ACK',1),('BYE',2)]: + if method=='BYE': time.sleep(.3 if number=='99105' else 1.5) + seq=[f'{method} sip:uas@127.0.0.1:50630 SIP/2.0',f'Via: SIP/2.0/UDP 127.0.0.1:{port};branch=z9hG4bK{uuid.uuid4().hex};rport',base[1],f'To: {header(msg,"To")}',base[3],*['Route: '+r.strip() for r in routes],f'CSeq: {cseq} {method}','Max-Forwards: 70','Content-Length: 0','',''] + sock.sendto('\r\n'.join(seq).encode(),('127.0.0.1',15060)) + if method=='BYE': + final=sock.recv(65535).decode(); assert final.startswith('SIP/2.0 200'), final + break + time.sleep(.3) +finally: + stop.set(); thread.join(timeout=2) +print(json.dumps(events))