feat: add telephony layer and asterisk-events worker

- packages/telephony: cliente AMI proprio sobre TCP puro (sem dependencia
  de terceiros pouco mantida) + interface TelephonyProvider +
  AsteriskTelephonyProvider (Originate, Hangup, QueuePause/Add/Remove,
  QueueStatus, ExtensionState, DeviceState, PJSIPShowEndpoints/Contacts,
  Reload, runCommand, stream de eventos). Testado contra o Asterisk real —
  o formato de resposta do Command mudou entre versoes do Asterisk
  (headers 'Output:' repetidos em vez de 'Response: Follows'/'--END
  COMMAND--'), corrigido apos inspecionar os bytes crus do protocolo
- apps/asterisk-events: worker dedicado a manter a conexao AMI viva,
  normalizar eventos (Newchannel, DialBegin/End, Hangup, DeviceStateChange,
  ContactStatus, eventos de fila/agente), persistir ExtensionState no
  Postgres e publicar em Redis pub/sub para consumo em tempo real.
  Containerizado, alcanca o Asterisk (host network) via
  host.docker.internal a partir da rede bridge. Heartbeat no Redis para
  health check
- packages/database: novos modelos Trunk, Extension, ExtensionState
  (migration aplicada)
- packages/shared: secret-crypto.ts (AES-256-GCM para credenciais de trunk
  e senha SIP em repouso, master key externa ao banco)

Testado ponta a ponta: chamada real originada -> eventos normalizados
recebidos via Redis SUBSCRIBE, heartbeat renovando no TTL correto.
This commit is contained in:
2026-08-27 12:41:41 -03:00
parent a2898fa566
commit 6f3d731581
20 changed files with 1059 additions and 4 deletions

View File

@@ -0,0 +1,22 @@
{
"name": "@b2bcall/asterisk-events",
"version": "0.1.0",
"private": true,
"main": "dist/main.js",
"scripts": {
"build": "tsc",
"start": "node dist/main.js",
"start:dev": "ts-node src/main.ts"
},
"dependencies": {
"@b2bcall/database": "workspace:*",
"@b2bcall/telephony": "workspace:*",
"ioredis": "^5.4.2",
"pino": "^9.6.0"
},
"devDependencies": {
"@types/node": "^24.0.0",
"ts-node": "^10.9.2",
"typescript": "^5.7.3"
}
}

View File

@@ -0,0 +1,7 @@
import pino from 'pino';
// Logs estruturados JSON (agente.md seção 60) — nunca loga segredos.
export const logger = pino({
level: process.env.LOG_LEVEL ?? 'info',
base: { service: 'b2bcall-asterisk-events' },
});

View File

@@ -0,0 +1,83 @@
import { AsteriskTelephonyProvider, type AmiMessage } from '@b2bcall/telephony';
import { PrismaClient } from '@b2bcall/database';
import Redis from 'ioredis';
import { logger } from './logger';
import { extractExtensionStatePatch, RELEVANT_EVENTS } from './normalize';
const REDIS_CHANNEL_EXTENSIONS = 'b2bcall:events:extensions';
const REDIS_CHANNEL_ASTERISK = 'b2bcall:events:asterisk';
const HEARTBEAT_KEY = 'b2bcall:asterisk-events:heartbeat';
const HEARTBEAT_TTL_SECONDS = 15;
async function main() {
const prisma = new PrismaClient();
const redis = new Redis(process.env.REDIS_URL!);
const provider = new AsteriskTelephonyProvider({
host: process.env.ASTERISK_HOST!,
amiPort: Number(process.env.AMI_PORT ?? 5038),
amiUsername: process.env.AMI_USERNAME!,
amiSecret: process.env.AMI_SECRET!,
reconnect: true,
});
provider.onEvent((event: AmiMessage) => {
void handleEvent(event).catch((err) => logger.error({ err, event }, 'Falha ao processar evento AMI'));
});
async function handleEvent(event: AmiMessage) {
if (!event.Event || !RELEVANT_EVENTS.has(event.Event)) return;
await redis.publish(REDIS_CHANNEL_ASTERISK, JSON.stringify(event));
const patch = extractExtensionStatePatch(event);
if (!patch) return;
const updated = await prisma.extensionState.upsert({
where: { extension: patch.extension },
create: {
extension: patch.extension,
deviceState: patch.deviceState,
contactStatus: patch.contactStatus,
contactUri: patch.contactUri,
},
update: {
...(patch.deviceState !== undefined ? { deviceState: patch.deviceState } : {}),
...(patch.contactStatus !== undefined ? { contactStatus: patch.contactStatus } : {}),
...(patch.contactUri !== undefined ? { contactUri: patch.contactUri } : {}),
},
});
await redis.publish(REDIS_CHANNEL_EXTENSIONS, JSON.stringify(updated));
}
async function heartbeat() {
try {
await redis.set(HEARTBEAT_KEY, new Date().toISOString(), 'EX', HEARTBEAT_TTL_SECONDS);
} catch (err) {
logger.error({ err }, 'Falha ao gravar heartbeat no Redis');
}
}
logger.info('Conectando ao AMI...');
await provider.connect();
logger.info('Conectado ao AMI. Escutando eventos.');
await heartbeat();
setInterval(() => void heartbeat(), (HEARTBEAT_TTL_SECONDS * 1000) / 2);
const shutdown = async () => {
logger.info('Encerrando apps/asterisk-events...');
provider.disconnect();
await redis.quit();
await prisma.$disconnect();
process.exit(0);
};
process.on('SIGTERM', () => void shutdown());
process.on('SIGINT', () => void shutdown());
}
main().catch((err) => {
logger.error({ err }, 'Erro fatal ao iniciar apps/asterisk-events');
process.exit(1);
});

View File

@@ -0,0 +1,61 @@
import type { AmiMessage } from '@b2bcall/telephony';
export interface ExtensionStatePatch {
extension: string;
deviceState?: string;
contactStatus?: string;
contactUri?: string;
}
/**
* Extrai uma atualização de estado de ramal a partir de um evento AMI bruto,
* ou `null` se o evento não for relevante para o painel de ramais. Os
* ramais PJSIP são nomeados pelo próprio número (ver ExtensionsService), por
* isso "PJSIP/1001" -> "1001" e EndpointName já É o número do ramal.
*/
export function extractExtensionStatePatch(event: AmiMessage): ExtensionStatePatch | null {
switch (event.Event) {
case 'DeviceStateChange': {
const device = event.Device ?? '';
const [tech, exten] = device.split('/');
if (tech !== 'PJSIP' || !exten) return null;
return { extension: exten, deviceState: event.State };
}
case 'ContactStatus': {
const endpoint = event.EndpointName;
if (!endpoint) return null;
return {
extension: endpoint,
contactStatus: event.ContactStatus,
contactUri: event.URI,
};
}
default:
return null;
}
}
// Eventos administrativos/de chamada que a spec pede para consumir (agente.md
// seção 8), além dos dois acima já usados para estado de ramal. Persistência
// rica (call_events/queue_events) chega nas Fases 6/7 quando as tabelas de
// campanha/fila existirem — por ora, apenas repassados ao Redis para quem
// quiser consumir em tempo real (ex.: futura Fase 5 de monitoramento de filas).
export const RELEVANT_EVENTS = new Set([
'Newchannel',
'DialBegin',
'DialEnd',
'BridgeEnter',
'BridgeLeave',
'Hangup',
'Newstate',
'DeviceStateChange',
'QueueMemberStatus',
'QueueMemberPause',
'AgentConnect',
'AgentComplete',
'QueueCallerJoin',
'QueueCallerLeave',
'QueueCallerAbandon',
'ContactStatus',
'PeerStatus',
]);

View File

@@ -0,0 +1,14 @@
{
"compilerOptions": {
"module": "commonjs",
"moduleResolution": "node",
"target": "ES2022",
"outDir": "dist",
"rootDir": "src",
"declaration": false,
"esModuleInterop": true,
"skipLibCheck": true,
"strict": false
},
"include": ["src"]
}