2 Commits
29 changed files with 1287 additions and 0 deletions
+17
View File
@@ -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` - 当前任务:`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` - 本次变更摘要:确认 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=successB `mysql``redis-server``nginx``lisglosips@api``lisglosips@cdr-worker``lisglosips@recording-worker``lisglosips@config-publisher``heplify-server``lisglosips-prometheus``grafana-server` activepreflight okAPI 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 重新演练 - 总体状态: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=successB `mysql``redis-server``nginx``lisglosips@api``lisglosips@cdr-worker``lisglosips@recording-worker``lisglosips@config-publisher``heplify-server``lisglosips-prometheus``grafana-server` activepreflight okAPI 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 重新演练
+4
View File
@@ -12,6 +12,10 @@ V2 的目标是把当前纯前端原型转化为可部署、可联调、可上
本文作为下一阶段执行基线,后续数据库 DDL、OpenAPI、OpenSIPS 脚本和部署脚本应从本文拆分并版本化管理。 本文作为下一阶段执行基线,后续数据库 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. V2 范围
### 2.1 本期实施页面 ### 2.1 本期实施页面
+2
View File
@@ -1,4 +1,5 @@
import { Module } from '@nestjs/common'; import { Module } from '@nestjs/common';
import { CallerAnalyticsModule } from './caller-analytics/caller-analytics.module.js';
import { ConfigModule } from '@nestjs/config'; import { ConfigModule } from '@nestjs/config';
import crypto from 'node:crypto'; import crypto from 'node:crypto';
import { LoggerModule } from 'nestjs-pino'; import { LoggerModule } from 'nestjs-pino';
@@ -60,6 +61,7 @@ import { VendorsModule } from './vendors/vendors.module.js';
AuditModule, AuditModule,
AuthModule, AuthModule,
ActiveCallsModule, ActiveCallsModule,
CallerAnalyticsModule,
BusinessPrefixesModule, BusinessPrefixesModule,
CdrsModule, CdrsModule,
DashboardModule, DashboardModule,
@@ -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<string, unknown>, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); }
@Get('summary') summary(@Query() q: Record<string, unknown>, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); }
@Get('numbers') numbers(@Query() q: Record<string, unknown>, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); }
@Get('trends') trends(@Query() q: Record<string, unknown>, @CurrentUserParam() user: CurrentUser) { return this.service.overview(q, user); }
@Get('calls') calls(@Query() q: Record<string, unknown>, @CurrentUserParam() user: CurrentUser) { return this.service.detail(q, user); }
@Get('breakdown') breakdown(@Query() q: Record<string, unknown>, @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 {}
@@ -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();
});
});
@@ -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<string, unknown>;
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<string, string | number | null>;
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<string[] | null> {
if (user.permissions.includes('caller_analytics.view_all')) return null;
const access = await this.db.$queryRaw<Array<{ customer_id: string }>>`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<Row[]>(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<Row[]>(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<Array<{ payload: Record<string, unknown>; updated_at: Date }>>`SELECT payload, updated_at FROM caller_analysis_health WHERE id='worker'`,
this.db.$queryRaw<Array<{ customer_id: string; caller: string; payload: unknown; status: string }>>(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<Array<Record<string, unknown>>>(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<Row[]>(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<Array<{ id: string; name: string }>>(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 };
}
}
+3
View File
@@ -16,6 +16,7 @@ import {
normalizeVendorGateway, normalizeVendorGateway,
} from './utils/formatters.js'; } from './utils/formatters.js';
import { DashboardPage } from './pages/DashboardPage.jsx'; import { DashboardPage } from './pages/DashboardPage.jsx';
import { CallerAnalyticsPage } from './pages/CallerAnalyticsPage.jsx';
import { ActiveCallsPage } from './pages/ActiveCallsPage.jsx'; import { ActiveCallsPage } from './pages/ActiveCallsPage.jsx';
import { CustomersPage } from './pages/CustomersPage.jsx'; import { CustomersPage } from './pages/CustomersPage.jsx';
import { RechargeRecordsPage } from './pages/RechargeRecordsPage.jsx'; import { RechargeRecordsPage } from './pages/RechargeRecordsPage.jsx';
@@ -58,6 +59,7 @@ const navGroups = [
items: [ items: [
{ key: 'activeCalls', label: '当前通话', permissions: PAGE_PERMISSIONS.activeCalls }, { key: 'activeCalls', label: '当前通话', permissions: PAGE_PERMISSIONS.activeCalls },
{ key: 'cdr', label: '话单中心', permissions: PAGE_PERMISSIONS.cdr }, { key: 'cdr', label: '话单中心', permissions: PAGE_PERMISSIONS.cdr },
{ key: 'callerAnalytics', label: '主叫号码分析', permissions: PAGE_PERMISSIONS.callerAnalytics },
{ key: 'quality', label: '质检中心', permissions: PAGE_PERMISSIONS.quality }, { key: 'quality', label: '质检中心', permissions: PAGE_PERMISSIONS.quality },
{ key: 'sipops', label: 'SIP 运维', pending: true }, { key: 'sipops', label: 'SIP 运维', pending: true },
{ key: 'monitoring', label: '监控告警', pending: true }, { key: 'monitoring', label: '监控告警', pending: true },
@@ -76,6 +78,7 @@ const navGroups = [
const pages = { const pages = {
callerAnalytics: CallerAnalyticsPage,
dashboard: DashboardPage, dashboard: DashboardPage,
activeCalls: ActiveCallsPage, activeCalls: ActiveCallsPage,
customers: CustomersPage, customers: CustomersPage,
+3
View File
@@ -128,6 +128,9 @@ function queryString(params = {}) {
} }
export const api = { 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'), captcha: () => request('/auth/captcha'),
login: async (body) => { login: async (body) => {
const response = await request('/auth/login', { method: 'POST', body: jsonBody(body) }); const response = await request('/auth/login', { method: 'POST', body: jsonBody(body) });
+104
View File
@@ -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) => <Field label={label}><Select value={draft[key]} onChange={e => change(key, e.target.value)}><option value="">全部</option>{choices.map(c => <option key={c.id} value={c.id}>{c.name}</option>)}</Select></Field>;
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 <>
<PageTitle title="主叫号码分析" desc="接通:180 / 183 或后续阶段;应答:实际通话时长大于 0。" actions={<><Button variant="outline" onClick={() => setAuto(v => !v)}>{auto ? '自动刷新中' : '已暂停刷新'}</Button><Button disabled={loading} onClick={() => setTick(v => v + 1)}>{loading ? '更新中…' : '刷新'}</Button></>} />
{error ? <Alert tone="danger" title="更新失败,保留上次数据">{error}</Alert> : null}
<Panel title="筛选条件"><form onSubmit={e => { e.preventDefault(); setFilters({ ...draft, skip: 0 }); setTarget(null); }}>
<Toolbar>
<Field label="号码视角"><Select value={draft.view} onChange={e => change('view', e.target.value)}><option value="landing">落地主叫</option><option value="original">原始主叫</option></Select></Field>
<Field label="时间窗口"><Select value={draft.minutes} onChange={e => change('minutes', e.target.value)}>{[5,15,30,60].map(n => <option key={n} value={n}>最近 {n} 分钟</option>)}<option value="today">今日</option><option value="custom">自定义</option></Select></Field>
{draft.minutes === 'custom' ? <><Field label="开始时间"><Input type="datetime-local" required value={draft.from || ''} onChange={e => change('from', e.target.value)} /></Field><Field label="结束时间"><Input type="datetime-local" required value={draft.to || ''} onChange={e => change('to', e.target.value)} /></Field></> : null}
<Field label="主叫号码"><Input value={draft.caller} onChange={e => change('caller', e.target.value)} placeholder="精确匹配号码" /></Field>
{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]})))}
<Field label="被叫城市代码"><Input value={draft.city} onChange={e => change('city', e.target.value)} /></Field>
<Field label="最小呼叫数"><Input type="number" min="0" value={draft.minSamples} onChange={e => change('minSamples', e.target.value)} /></Field>
<Field label="排序"><Select value={draft.sort} onChange={e => change('sort', e.target.value)}>{[...countColumns, ...rateColumns].map(c => <option key={c.key} value={c.key}>{c.label}</option>)}</Select></Field>
<Field label="顺序"><Select value={draft.order} onChange={e => change('order', e.target.value)}><option value="desc">从高到低</option><option value="asc">从低到高</option></Select></Field>
<Field label="异常筛选"><Select value={String(draft.onlyAlerts)} onChange={e => change('onlyAlerts', e.target.value === 'true')}><option value="false">全部号码</option><option value="true">仅看告警</option></Select></Field>
</Toolbar>
<details className="cra-rate-filters"><summary>接通率 / 应答率区间</summary><Toolbar>{[['connection','接通率'],['overall','总体应答率'],['connected','已接通应答率']].map(([key,label]) => <Field key={key} label={`${label}%`}><span className="table-actions"><Input aria-label={`${label}下限`} type="number" min="0" max="100" value={draft[`${key}Min`] ?? ''} onChange={e => change(`${key}Min`, e.target.value)} placeholder="0" /><span></span><Input aria-label={`${label}上限`} type="number" min="0" max="100" value={draft[`${key}Max`] ?? ''} onChange={e => change(`${key}Max`, e.target.value)} placeholder="100" /></span></Field>)}</Toolbar></details>
<Button type="submit">查询</Button>
</form></Panel>
<section className="metric-grid cra-metrics">{[...countColumns, ...rateColumns].map(c => <div key={c.key} className="metric-card"><span>{c.label}</span><strong>{metrics ? c.render ? c.render(metrics) : metrics[c.key] : '—'}</strong></div>)}</section>
<p className="muted-text" role="status">{data ? `${formatDateTime(data.from)}${formatDateTime(data.to)} · ${data.qualityStatus === 'LIVE' ? '采集运行中' : '数据滞后或不完整'} · 待接通 ${metrics.pendingCalls} · 已结束未接通 ${metrics.failedCalls} · 未知 ${metrics.unknownCalls} · 更新于 ${formatDateTime(data.asOf)}` : '正在读取真实统计数据…'}</p>
{data ? <Alert title="口径与数据范围" tone={data.qualityStatus === 'LIVE' ? 'info' : 'warning'}>{data.historyNotice} {data.summaryScope}183 不代表实际振铃未接通数含待接通进行中指标为暂定值</Alert> : null}
<Panel title="分钟趋势" aside={<span className="muted-text">呼叫 / 接通 / 绿应答</span>}>
<div className="cra-chart" aria-label="分钟呼叫接通应答趋势">{timeline.map((r,i) => <div key={r.time ?? i} className="cra-chart-column" title={`${formatDateTime(Number(r.time))} 呼叫${r.totalCalls} 接通${r.connectedCalls} 应答${r.answeredCalls}`}><div className="cra-bars"><span style={{height:`${r.totalCalls/max*100}%`}} /><span style={{height:`${r.connectedCalls/max*100}%`}} /><span style={{height:`${r.answeredCalls/max*100}%`}} /></div><small>{new Date(Number(r.time)).toLocaleTimeString('zh-CN',{hour:'2-digit',minute:'2-digit'})}</small></div>)}</div>
{!timeline.length ? <p className="muted-text">当前窗口暂无已采集呼叫</p> : <details><summary>查看趋势数值与三个率</summary><SimpleTable rows={timeline.map((r,i) => ({ ...r, id:i, minute:formatDateTime(Number(r.time)) }))} columns={[{key:'minute',label:'分钟'},...countColumns,...rateColumns]} /></details>}
</Panel>
<Panel title="主叫号码列表" aside={<Badge>{data?.total ?? 0} 个号码</Badge>} className="wide-panel">
<SimpleTable rows={(data?.numbers ?? []).map(r => ({ ...r, id:`${r.customerId}:${r.caller}`, customerName:names[r.customerId] || r.customerId }))} loading={!data && loading} onRowClick={openDetail} columns={[
{key:'caller',label:'主叫号码',render:r => <Button size="sm" variant="ghost" onClick={() => openDetail(r)}>{r.caller}</Button>}, {key:'customerName',label:'客户'}, ...countColumns,...rateColumns,
{key:'pendingCalls',label:'待接通'}, {key:'unknownCalls',label:'未知'}, {key:'alert',label:'状态',render:r => <Badge tone={r.alert ? 'danger' : 'neutral'}>{r.alert ? r.alert.reasons?.join('、') || '告警' : r.unknownCalls ? '证据不足' : '观察'}</Badge>}
]} />
<Pagination total={data?.total ?? 0} take={filters.take} skip={filters.skip} count={data?.numbers.length ?? 0} loading={loading} onPageChange={skip => setFilters(f => ({...f,skip}))} onPageSizeChange={take => setFilters(f => ({...f,take,skip:0}))} />
</Panel>
{target ? <Drawer title={`号码详情 · ${target.caller}`} onClose={() => setTarget(null)}>
{detailError ? <Alert tone="danger" title="详情更新失败">{detailError}</Alert> : null}
<p>{names[target.customerId] || target.customerId} · {filters.view === 'landing' ? '落地主叫' : '原始主叫'}</p>
<Panel title="线路 / 地区 / 运营商分布"><SimpleTable rows={(detail?.breakdown ?? []).map((r,i) => ({...r,id:i,gateway:gatewayNames[r.vendorGatewayId] || r.vendorGatewayId || '业务呼叫汇总'}))} columns={[{key:'gateway',label:'落地网关'},{key:'city',label:'城市'},{key:'carrier',label:'运营商'},...countColumns,...rateColumns]} /></Panel>
<Panel title="关联呼叫与信令证据"><SimpleTable rows={(detail?.rows ?? []).map(r => {
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:'发起时间'}]} /></Panel>
<div className="table-actions"><Button variant="outline" disabled={!detailSkip} onClick={() => setDetailSkip(v => Math.max(0,v-25))}>上一页</Button><Button variant="outline" disabled={(detail?.rows.length ?? 0)<25} onClick={() => setDetailSkip(v => v+25)}>下一页</Button></div>
<p className="muted-text"> Call-ID 对应原始话单或信令旧固定时长话单仅作追溯不参与本统计校准</p>
</Drawer> : null}
</>;
}
+1
View File
@@ -1,4 +1,5 @@
export const PAGE_PERMISSIONS = { export const PAGE_PERMISSIONS = {
callerAnalytics: ['caller_analytics.view'],
dashboard: ['dashboard.view'], dashboard: ['dashboard.view'],
activeCalls: ['active_calls.view'], activeCalls: ['active_calls.view'],
customers: ['customers.view'], customers: ['customers.view'],
+9
View File
@@ -2893,3 +2893,12 @@ a {
animation: none; 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; }
+36
View File
@@ -0,0 +1,36 @@
import { PrismaClient, Prisma, analyticsHash, analyticsRates, COUNT_KEYS, type AnalyticsCounts } from '@lisglosips/database';
export async function evaluateCallerAlerts(db: PrismaClient) {
const end = new Date(Math.floor(Date.now() / 60000) * 60000); const start = new Date(+end - 15 * 60000);
const aggregates = Prisma.join(COUNT_KEYS.map(k => Prisma.sql`SUM(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,${`$.${k}`})) AS SIGNED)) AS ${Prisma.raw(k)}`));
const rows = await db.$queryRaw<Array<AnalyticsCounts & { customer_id: string; caller: string; view: string }>>(Prisma.sql`SELECT customer_id, caller, view, ${aggregates} FROM caller_analysis_minute_buckets WHERE started_at>=${start} AND started_at<${end} GROUP BY customer_id, caller, view HAVING SUM(CAST(JSON_UNQUOTE(JSON_EXTRACT(counts,'$.totalCalls')) AS SIGNED))>=50`);
for (const row of rows) {
const counts = Object.fromEntries(COUNT_KEYS.map(k => [k, Number(row[k])])) as unknown as AnalyticsCounts;
const rates = analyticsRates(counts);
if (rates.unknownCalls || rates.activeCalls / rates.totalCalls > 0.2) continue;
const reasons = [
rates.connectionRate !== null && rates.connectionRate < Number(process.env.ANALYTICS_LOW_CONNECTION_RATE ?? 20) ? '低接通率' : '',
rates.overallAnswerRate !== null && rates.overallAnswerRate < Number(process.env.ANALYTICS_LOW_OVERALL_RATE ?? 10) ? '低总体应答率' : '',
rates.connectedCalls >= 30 && rates.connectedAnswerRate !== null && rates.connectedAnswerRate < Number(process.env.ANALYTICS_LOW_CONNECTED_RATE ?? 20) ? '低已接通应答率' : ''
].filter(Boolean);
const previous = await db.$queryRaw<Array<AnalyticsCounts>>(Prisma.sql`SELECT ${aggregates} FROM caller_analysis_minute_buckets WHERE started_at>=${new Date(+start - 15 * 60000)} AND started_at<${start} AND customer_id=${row.customer_id} AND caller=${row.caller} AND view=${row.view}`);
const p = analyticsRates(Object.fromEntries(COUNT_KEYS.map(k => [k, Number(previous[0]?.[k] ?? 0)])) as unknown as AnalyticsCounts);
if (p.totalCalls >= 50 && !p.unknownCalls && p.activeCalls / p.totalCalls <= 0.2) {
for (const [key, label] of [['connectionRate','接通率骤降'],['overallAnswerRate','总体应答率骤降'],['connectedAnswerRate','已接通应答率骤降']] as const) {
if (key === 'connectedAnswerRate' && (p.connectedCalls < 30 || rates.connectedCalls < 30)) continue;
if (p[key] !== null && rates[key] !== null && p[key]! - rates[key]! > Number(process.env.ANALYTICS_RATE_DROP_POINTS ?? 15)) reasons.push(label);
}
}
const tail = await db.$queryRaw<Array<{ counts: AnalyticsCounts }>>`SELECT counts FROM caller_analysis_states WHERE customer_id=${row.customer_id} AND caller=${row.caller} AND view=${row.view} AND started_at>=${start} AND started_at<${end} ORDER BY started_at DESC, id DESC LIMIT 20`;
if (tail.length === 20 && tail.every(t => t.counts.failedCalls === 1)) reasons.push('连续20次已结束未接通');
const id = analyticsHash([row.customer_id, row.caller, row.view]);
const existing = await db.$queryRaw<Array<{ status: string; payload: { streak?: number; recovered?: number; since?: string } }>>`SELECT status,payload FROM caller_analysis_alerts WHERE id=${id}`;
const streak = reasons.length ? (existing[0]?.payload.streak ?? 0) + 1 : 0;
const recovered = reasons.length ? 0 : (existing[0]?.payload.recovered ?? 0) + 1;
const status = streak >= 3 ? 'ACTIVE' : existing[0]?.status === 'ACTIVE' && recovered < 3 ? 'ACTIVE' : 'OBSERVING';
const payload = JSON.stringify({ reasons, rates, streak, recovered, windowStart: start, windowEnd: end, since: existing[0]?.payload.since ?? new Date().toISOString() });
await db.$executeRaw`INSERT INTO caller_analysis_alerts(id,customer_id,view,caller,status,payload,updated_at) VALUES (${id},${row.customer_id},${row.view},${row.caller},${status},${payload},NOW(3)) ON DUPLICATE KEY UPDATE status=VALUES(status),payload=VALUES(payload),updated_at=VALUES(updated_at)`;
}
// Old low-sample windows are no longer active evidence, not permanent alarms.
await db.$executeRaw`UPDATE caller_analysis_alerts SET status='EXPIRED' WHERE updated_at<${new Date(Date.now() - 120000)} AND status='ACTIVE'`;
}
+102
View File
@@ -0,0 +1,102 @@
import { open, readFile, writeFile, rename, stat, appendFile } from 'node:fs/promises';
import { hostname } from 'node:os';
import { PrismaClient, Prisma, CallerAnalyticsStore, ANALYTICS_GROUP, ANALYTICS_STREAM, analyticsHash, parseAnalyticsEvent, type AnalyticsCall, type AnalyticsEvent } from '@lisglosips/database';
import { createRedisClient } from '@lisglosips/redis';
import { createLogger } from '@lisglosips/observability';
import { evaluateCallerAlerts } from './analytics-alerts.js';
const logger = createLogger('caller-analytics');
const db = new PrismaClient(); const redis = createRedisClient(process.env.REDIS_URL ?? '');
const store = new CallerAnalyticsStore(db);
const spool = process.env.ANALYTICS_SPOOL ?? '/var/log/lisglosips/caller-analytics.log';
const checkpoint = process.env.ANALYTICS_CHECKPOINT ?? '/var/lib/lisglosips/caller-analytics-offset.json';
let cursor = { offset: 0, ino: 0 }; let stopping = false; let lastAlerts = 0;
let dataThrough: string | null = null; let gap = false;
const consumer = `${hostname()}-${process.pid}`;
process.on('SIGTERM', () => { stopping = true; }); process.on('SIGINT', () => { stopping = true; });
function decodeLine(line: string): AnalyticsEvent | null {
const marker = line.indexOf('CRA1|'); if (marker < 0) return null;
const parts = line.slice(marker).trim().split('|');
if (parts.length !== 20) throw new Error('Malformed producer record');
const decoded = parts.slice(8).map(s => Buffer.from(s, 'base64').toString('utf8'));
const [uid, callId, customerId, customerGatewayId, vendorId, vendorGatewayId, caller, landingCaller, callee, city, carrier, attemptId] = decoded;
const at = Number(parts[2]) * 1000 + Math.floor(Number(parts[3]) / 1000);
const startedAt = Number(parts[4]) * 1000 + Math.floor(Number(parts[5]) / 1000);
const fields = { version: '1', eventId: analyticsHash(parts), callUid: analyticsHash(uid), callId, customerId, customerGatewayId, vendorId, vendorGatewayId, caller, landingCaller, callee, city, carrier, attemptId,
kind: parts[1], at: String(at), startedAt: String(startedAt), code: parts[6], source: parts[7], talkMs: '-1', sequence: '0' };
return parseAnalyticsEvent(Object.entries(fields).flat());
}
async function collect() {
const info = await stat(spool);
if (cursor.ino && (cursor.ino !== info.ino || info.size < cursor.offset)) { gap = true; cursor.offset = 0; }
cursor.ino = info.ino;
const file = await open(spool, 'r');
try {
const buffer = Buffer.alloc(Math.min(1024 * 1024, Math.max(0, info.size - cursor.offset)));
const { bytesRead } = await file.read(buffer, 0, buffer.length, cursor.offset);
const last = buffer.subarray(0, bytesRead).lastIndexOf(10); if (last < 0) return;
for (const line of buffer.subarray(0, last + 1).toString('utf8').split('\n').filter(Boolean)) {
let event: AnalyticsEvent | null;
try { event = decodeLine(line); } catch {
await appendFile(`${checkpoint}.deadletter`, `${JSON.stringify({ at:new Date(), line })}\n`, { mode:0o600 });
gap = true; continue;
}
if (event) await redis.xadd(ANALYTICS_STREAM, '*', 'event', JSON.stringify(event));
}
// Persist only after durable queue writes. Crash before checkpoint replays harmlessly.
cursor.offset += last + 1;
await writeFile(`${checkpoint}.tmp`, JSON.stringify(cursor), { mode: 0o600 }); await rename(`${checkpoint}.tmp`, checkpoint);
} finally { await file.close(); }
}
async function entries(entries: Array<[string, string[]]>) {
for (const [id, fields] of entries) {
const event = JSON.parse(fields[1]) as AnalyticsEvent;
await store.consume(event);
await redis.xack(ANALYTICS_STREAM, ANALYTICS_GROUP, id);
dataThrough = new Date(event.at).toISOString();
}
}
async function observeActive() {
const response = await fetch(process.env.ANALYTICS_MI_URL ?? 'http://127.0.0.1:8888/mi', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'dlg_list', params: [] }), signal: AbortSignal.timeout(2000) });
const body = await response.json() as { result?: { Dialogs?: Array<{ callid: string; state: number }> } };
if (!response.ok || !body.result?.Dialogs) throw new Error('Dialog observation unavailable');
const calls = await db.$queryRaw<Array<{ payload: AnalyticsCall }>>`SELECT c.payload FROM caller_analysis_calls c WHERE EXISTS (SELECT 1 FROM caller_analysis_states s WHERE s.call_key=c.id AND s.view='original' AND CAST(JSON_UNQUOTE(JSON_EXTRACT(s.counts,'$.activeCalls')) AS SIGNED)>0) LIMIT 2000`;
const live = new Map(body.result.Dialogs.map(d => [d.callid, Number(d.state)]));
const now = Date.now();
for (const { payload } of calls) for (const leg of Object.values(payload.legs)) {
if (leg.endedAt !== null || leg.acceptedAt === null || leg.talkMs > 0) continue;
// Only a confirmed live dialog supplies provisional positive duration. Never a browser timer.
if ((live.get(leg.event.callId) ?? 0) >= 3 && (live.get(leg.event.callId) ?? 0) <= 4 && now > leg.acceptedAt) {
await store.consume({ ...leg.event, kind: 'DURATION', at: now, talkMs: now - leg.acceptedAt, source: 'dialog-observation', eventId: analyticsHash([leg.event.callUid, leg.event.attemptId, 'positive-duration']) });
}
}
}
async function main() {
try { cursor = JSON.parse(await readFile(checkpoint, 'utf8')); } catch { /* First start reads the retained spool. */ }
await redis.connect(); await db.$connect();
try { await redis.xgroup('CREATE', ANALYTICS_STREAM, ANALYTICS_GROUP, '0', 'MKSTREAM'); } catch (e) { if (!(e instanceof Error) || !e.message.includes('BUSYGROUP')) throw e; }
let pendingCursor = '0-0';
while (!stopping) {
let degraded = gap; let reason = '';
try {
await collect();
const claimed = await redis.xautoclaim(ANALYTICS_STREAM, ANALYTICS_GROUP, consumer, 5000, pendingCursor, 'COUNT', 100) as [string, Array<[string, string[]]>];
pendingCursor = claimed[0]; await entries(claimed[1]);
const fresh = await redis.xreadgroup('GROUP', ANALYTICS_GROUP, consumer, 'COUNT', 200, 'BLOCK', 100, 'STREAMS', ANALYTICS_STREAM, '>') as Array<[string, Array<[string, string[]]>]> | null;
if (fresh) for (const [, batch] of fresh) await entries(batch);
await observeActive();
const groups = await redis.xinfo('GROUPS', ANALYTICS_STREAM) as Array<Array<string | number>>;
const group = groups.map(g => Object.fromEntries(Array.from({ length: g.length / 2 }, (_, i) => [g[i * 2], g[i * 2 + 1]]))).find(g => g.name === ANALYTICS_GROUP);
if (Number(group?.pending ?? 0) || Number(group?.lag ?? 0)) { degraded = true; reason = 'QUEUE_BACKLOG'; }
if (Date.now() - lastAlerts > 15000 && !degraded) { await evaluateCallerAlerts(db); lastAlerts = Date.now(); }
} catch (e) { degraded = true; reason = e instanceof Error ? e.message.slice(0, 200) : 'Worker error'; logger.error({ reason }, 'Analytics cycle failed'); }
try {
const payload = JSON.stringify({ degraded, reason, dataThrough, lagMs: degraded ? null : 0, captureSince: process.env.ANALYTICS_CAPTURE_SINCE ?? null, spoolOffset: cursor.offset, gap, updatedAt: new Date().toISOString() });
await db.$executeRaw(Prisma.sql`INSERT INTO caller_analysis_health(id,payload,updated_at) VALUES ('worker',${payload},NOW(3)) ON DUPLICATE KEY UPDATE payload=VALUES(payload),updated_at=VALUES(updated_at)`);
} catch { logger.error('Analytics heartbeat write failed'); }
await new Promise(resolve => setTimeout(resolve, 700));
}
redis.disconnect(); await db.$disconnect();
}
void main().catch(e => { logger.error({ message: e instanceof Error ? e.message : 'Startup failed' }, 'Analytics failed'); process.exitCode = 1; redis.disconnect(); void db.$disconnect(); });
+221
View File
@@ -0,0 +1,221 @@
# 主叫号码接通率与应答率实时分析设计方案
日期: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+PU 单列 |
| 应答数 | answeredCalls | A |
| 实时接通率 | connectionRate | C / T × 100% |
| 实时总体应答率 | overallAnswerRate | A / T × 100% |
| 实时已接通应答率 | connectedAnswerRate | A / C × 100% |
分母为 0 时接口返回 null、页面显示“—”。比率由整数计数计算,不平均各时间桶的百分比。数据不完整时,计数与比率应标注“已观测/暂定”,未知或采集缺口不可展示为完整结果。
旧方案的“已判定接通率”不作为主指标,本版以以上三个率为准;不使用未经解释的 ASR 简称混淆业务口径。
示例:T=100C=60(其中 A=30),F=25P=15U=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猜测匹配。
建议APIGET /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=1N=0A=1;三个率均100% |
| CRA-002 | INVITE→100→183→487,无正时长 | C=1N=0A=0;接通率100%,两个应答率0% |
| CRA-003 | INVITE→100→486,从未180/183/2xx | C=0N=F=1A=0;总体应答率0%,已接通应答率为null |
| CRA-004 | 未见100,直接180,随后取消 | C=1A=0,不因缺100漏计 |
| CRA-005 | 直接初始INVITE 200,正时长 | C=A=1DIRECT_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=0F=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。
- 告警通过独立进程环境变量配置三个低率阈值和前窗下降百分点;采用样本/未结束占比门槛、连续三次评估及三次恢复,站内展示,不执行停号换路。后续应按实际业务校准默认阈值。
+4
View File
@@ -1975,6 +1975,10 @@ corepack pnpm@10.33.0 exec vitest run apps/worker-recording/src/transfer.spec.ts
| 步骤 | 1. 导入器先干跑并确认未匹配为 0。2. 备份后导入。3. 核对四表数量。4. 四个页签执行首屏、翻页、页大小和筛选。 | | 步骤 | 1. 导入器先干跑并确认未匹配为 0。2. 备份后导入。3. 核对四表数量。4. 四个页签执行首屏、翻页、页大小和筛选。 |
| 预期结果 | 371 个行政城市、517,258 个有效七位号段、321 个区号、70 个前缀规则;首屏只读取 25 条;上一页/下一页与总数正确;页面加载不扫描或返回全表。 | | 预期结果 | 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. 缺陷分级 ## 9. 缺陷分级
| 级别 | 定义 | 示例 | | 级别 | 定义 | 示例 |
+87
View File
@@ -0,0 +1,87 @@
#!/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}`)')
# Previous delta releases stripped executable bits from bundled native engines.
find "$new/node_modules/.pnpm" -path '*/@prisma/engines/schema-engine-debian-openssl-3.0.x' -type f -exec chmod 0750 {} \;
/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"
@@ -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")
}
+2
View File
@@ -4,6 +4,8 @@ export const ACCESS_TOKEN_TYPE = 'Bearer';
export const PASSWORD_ALGO_ARGON2ID = 'argon2id'; export const PASSWORD_ALGO_ARGON2ID = 'argon2id';
export type PermissionKey = export type PermissionKey =
| 'caller_analytics.view'
| 'caller_analytics.view_all'
| 'dashboard.view' | 'dashboard.view'
| 'active_calls.view' | 'active_calls.view'
| 'active_calls.manage' | 'active_calls.manage'
@@ -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<string, AnalyticsLeg>; }
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<string, string> = {};
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) };
}
@@ -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 = <T>(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<Array<{ payload: AnalyticsCall | string | null }>>`SELECT payload FROM caller_analysis_calls WHERE id=${callKey} FOR UPDATE`;
const eventKey = analyticsHash(event.eventId);
const previous = await tx.$queryRaw<Array<{ id: string }>>`SELECT id FROM caller_analysis_events WHERE id=${eventKey}`;
if (previous.length) return 'duplicate';
const old = rows[0].payload ? json<AnalyticsCall>(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})`);
}
}
@@ -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> = {}): 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<AnalyticsCall | null>((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 });
});
});
+3
View File
@@ -1,3 +1,6 @@
export { Prisma, PrismaClient } from '@prisma/client'; export { Prisma, PrismaClient } from '@prisma/client';
export const DATABASE_PROVIDER = 'mysql'; export const DATABASE_PROVIDER = 'mysql';
export * from './caller-analytics-model.js';
export * from './caller-analytics-store.js';
@@ -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';
+85
View File
@@ -960,6 +960,91 @@ model OutboxEvent {
@@map("outbox_events") @@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 { model IdempotencyKey {
id String @id @db.VarChar(40) id String @id @db.VarChar(40)
key String @unique @db.VarChar(128) key String @unique @db.VarChar(128)
+2
View File
@@ -3,6 +3,8 @@ import { PrismaClient } from '@prisma/client';
const prisma = new PrismaClient(); const prisma = new PrismaClient();
const permissions = [ const permissions = [
['caller_analytics.view', 'caller_analytics', 'view', '查看授权客户的主叫分析'],
['caller_analytics.view_all', 'caller_analytics', 'view_all', '查看全部客户的主叫分析'],
['dashboard.view', 'dashboard', 'view', '查看 Dashboard 指标'], ['dashboard.view', 'dashboard', 'view', '查看 Dashboard 指标'],
['active_calls.view', 'active_calls', 'view', '查看当前通话'], ['active_calls.view', 'active_calls', 'view', '查看当前通话'],
['active_calls.manage', 'active_calls', 'manage', '强制挂断当前通话'], ['active_calls.manage', 'active_calls', 'manage', '强制挂断当前通话'],
+2
View File
@@ -74,6 +74,8 @@ const copyEntries = [
['infra/server-b/s56', 'infra/server-b/s56', true], ['infra/server-b/s56', 'infra/server-b/s56', true],
['infra/server-b/s57', 'infra/server-b/s57', true], ['infra/server-b/s57', 'infra/server-b/s57', true],
['infra/server-b/s58', 'infra/server-b/s58', 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/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], ['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], ['scripts/phase2-gateway-migration.mjs', 'scripts/phase2-gateway-migration.mjs', true],
+42
View File
@@ -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.b64encode})`).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.');
+58
View File
@@ -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(); }
+72
View File
@@ -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: <sip:uas@127.0.0.1:50630>', '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: <sip:99100001@127.0.0.1>;tag={tag}', f'To: <sip:{number}@127.0.0.1>', f'Call-ID: {callid}', f'Contact: <sip:99100001@127.0.0.1:{port}>', '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))