feat(recording): gravacao de chamadas + object storage abstraction

Fecha agente.md secao 90-94. A especificacao lista "Recording" e "Object
Storage" como dois passos separados na ordem de implementacao (secao
232), mas ficaram numa unica fase — sao acoplados o suficiente (Recording
precisa de um lugar pra guardar bytes) pra fazer sentido construir juntos.

## packages/storage — ObjectStorageProvider (secao 92)

Abstracao pequena: putObject/getObjectStream/deleteObject. Dois backends:
LocalObjectStorageProvider (filesystem, com checagem de path traversal
mesmo a key sendo sempre montada no servidor) e S3ObjectStorageProvider
(@aws-sdk/client-s3, preparado pra AWS S3 e MinIO via endpoint/
forcePathStyle customizaveis — nunca exercitado nesta sessao, sem
servidor S3 disponivel neste laboratorio). Escolhido por STORAGE_PROVIDER
env.

buildRecordingObjectKey (secao 93):
tenants/{tenant_id}/recordings/YYYY/MM/DD/{call_id}.wav, sempre montada
no servidor a partir de dados confiaveis.

## Bind mounts, nao volumes nomeados

/recordings e /data/object-storage usam bind mount pra um diretorio real
do host — apps/api roda no host, nao em Docker, e precisa enxergar os
mesmos arquivos que fs-events escreve. LOCAL_STORAGE_ROOT tem valores
diferentes por ambiente (mesmo padrao ja usado pra REDIS_URL).

## Quem grava: apps/predictive-dialer

So' chamadas originadas pelo discador com Campaign.recordingEnabled sao
gravadas nesta fase (unico caminho de originate que o sistema controla
hoje). RECORD_STEREO=true + execute_on_answer='record_session ...'
adicionados ao originate; origination_uuid pre-gerado (em vez de deixar o
provider sortear) porque o path de gravacao precisa dele antes do
comando de originate ser montado — o mesmo uuid vira Call.id no CDR.

## Quem sobe: apps/freeswitch-events/src/recording.ts

Em CALL_ENDED, encadeado depois do persistCallEvent terminar (nao em
paralelo) — uploadRecordingIfPresent le Call.talkTime/durationSeconds,
que e' exatamente o que persistCallEvent acabou de calcular no mesmo
evento (mesma classe de corrida ja corrigida uma vez na fase CDR, aqui
evitada por ordenacao). Sobe pro storage, cria Recording (retentionUntil
a partir de Plan.recordingRetentionDays), apaga o spool local.

## API + retencao

GET /recordings, GET /recordings/:id, GET /recordings/:id/audio (stream
autenticado, nunca URL direta pro storage). runRetentionSweep (secao 94)
no boot do apps/api + a cada hora — apaga o objeto, marca status=DELETED
(linha nunca apagada, fica como auditoria).

## Bug real achado testando esta fase

Recording.sizeBytes (BigInt) quebrava GET /recordings com 500 — Fastify
nao serializa BigInt nativamente (mesma classe de bug ja corrigida uma
vez no logger, fase Event Socket). Corrigido convertendo pra number na
resposta.

Verificado ponta a ponta: gravacao real criada (RIFF WAVE, PCM 16-bit,
ESTEREO 8000Hz — RECORD_STEREO confirmado), upload com path exato da
secao 93, download via API com md5 identico ao objeto original, varredura
de retencao apagando objeto + status DELETED + list/download bloqueados
depois. typecheck do workspace inteiro limpo.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01X1HxY46WGU4G1zmVDNKcWw
This commit is contained in:
2026-08-28 14:36:03 -03:00
parent 56499f4b99
commit c24a86776c
26 changed files with 1196 additions and 7 deletions

View File

@@ -15,6 +15,7 @@ COPY packages/types packages/types
COPY packages/shared packages/shared
COPY packages/telephony packages/telephony
COPY packages/database packages/database
COPY packages/storage packages/storage
COPY apps/freeswitch-events apps/freeswitch-events
RUN pnpm install --frozen-lockfile --filter @b2bcall/freeswitch-events...

View File

@@ -11,6 +11,7 @@
"dependencies": {
"@b2bcall/database": "workspace:*",
"@b2bcall/shared": "workspace:*",
"@b2bcall/storage": "workspace:*",
"@b2bcall/telephony": "workspace:*",
"esl": "11.2.1",
"ioredis": "^6.0.0"

View File

@@ -5,6 +5,7 @@ import { createLogger } from "@b2bcall/shared";
import { updateTrunkStatusFromGatewayEvent } from "./trunk-status";
import { resolveTenantIdForAgent, resolveTenantIdForQueue } from "./tenant-resolve";
import { persistCallEvent } from "./cdr";
import { uploadRecordingIfPresent } from "./recording";
const logger = createLogger("b2bcall-fs-events");
@@ -139,9 +140,23 @@ async function main() {
logger.error("falha ao publicar evento normalizado no Redis", { error: String(err) });
});
persistCallEvent(normalized).catch((err) => {
logger.error("falha ao persistir CDR", { error: String(err), type: normalized.type });
});
// .then() em vez de esperar aqui (await bloquearia o processamento do
// próximo evento ESL) — mas o upload da gravação só roda DEPOIS do
// persistCallEvent terminar de verdade, nunca em paralelo com ele:
// uploadRecordingIfPresent lê Call.talkTime/durationSeconds, que é
// exatamente o que persistCallEvent acabou de calcular no CALL_ENDED
// (mesma classe de corrida já corrigida uma vez neste arquivo, ver
// docs/CDR.md — aqui evitada por ordenação, não por retry).
persistCallEvent(normalized)
.then(() => {
if (normalized.type === "CALL_ENDED" && normalized.tenantId && normalized.callUuid) {
return uploadRecordingIfPresent(normalized.tenantId, normalized.callUuid);
}
return undefined;
})
.catch((err) => {
logger.error("falha ao persistir CDR/gravacao", { error: String(err), type: normalized.type });
});
logger.info(`evento: ${normalized.type}`, {
callUuid: normalized.callUuid,

View File

@@ -0,0 +1,103 @@
import { access, rm } from "node:fs/promises";
import { join } from "node:path";
import { getPrismaClient, withTenantContext } from "@b2bcall/database";
import { getObjectStorageProvider, buildRecordingObjectKey } from "@b2bcall/storage";
import { createLogger } from "@b2bcall/shared";
const logger = createLogger("b2bcall-fs-events");
const SPOOL_DIR = process.env.RECORDINGS_SPOOL_DIR ?? "/recordings";
const STORAGE_KIND = (process.env.STORAGE_PROVIDER ?? "local").toUpperCase() as "LOCAL" | "S3";
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
async function fileExists(path: string): Promise<boolean> {
try {
await access(path);
return true;
} catch {
return false;
}
}
/**
* Sobe pro object storage a gravação de uma chamada que acabou de
* terminar (agente.md secao 90-93) — chamado a partir de CALL_ENDED em
* main.ts, sempre depois do CDR já ter persistido `Call` (precisa de
* `Call.createdAt`/`talkTime` pra `recordedAt`/`durationSeconds`).
*
* `originateSimulatedAnswerLeg`/`originateRealPstnLeg`
* (apps/predictive-dialer) só setam `execute_on_answer='record_session
* ...'` quando `Campaign.recordingEnabled` é true — pra qualquer outra
* chamada, o arquivo simplesmente não existe, e essa função não faz nada
* (checagem por existência de arquivo, não por reconsultar a campanha).
*/
export async function uploadRecordingIfPresent(tenantId: string, callId: string): Promise<void> {
const spoolPath = join(SPOOL_DIR, `${callId}.wav`);
// record_session termina de gravar no CHANNEL_HANGUP_COMPLETE, mas pode
// levar um instante pra flush terminar — tenta algumas vezes antes de
// desistir, em vez de perder a gravação por uma corrida.
let exists = await fileExists(spoolPath);
for (let attempt = 0; !exists && attempt < 5; attempt++) {
await sleep(300);
exists = await fileExists(spoolPath);
}
if (!exists) return;
const prisma = getPrismaClient();
try {
const call = await withTenantContext(prisma, tenantId, (tx) => tx.call.findUniqueOrThrow({ where: { id: callId } }));
const tenant = await prisma.tenant.findUniqueOrThrow({ where: { id: tenantId }, include: { plan: true } });
const recordedAt = call.createdAt;
const objectKey = buildRecordingObjectKey(tenantId, callId, recordedAt);
const storage = getObjectStorageProvider();
const { sizeBytes, checksum } = await storage.putObject(objectKey, spoolPath);
const retentionDays = tenant.plan.recordingRetentionDays;
const retentionUntil = retentionDays
? new Date(recordedAt.getTime() + retentionDays * 24 * 60 * 60 * 1000)
: null;
await withTenantContext(prisma, tenantId, (tx) =>
tx.recording.create({
data: {
tenantId,
callId,
storageProvider: STORAGE_KIND,
objectKey,
format: "wav",
durationSeconds: call.talkTime ?? call.durationSeconds,
channels: 2,
sizeBytes,
checksum,
recordedAt,
retentionUntil,
},
}),
);
await rm(spoolPath, { force: true });
logger.info("gravacao enviada pro object storage", { callId, objectKey, sizeBytes });
} catch (err) {
logger.error("falha ao processar gravacao", { error: String(err), callId });
await withTenantContext(prisma, tenantId, (tx) =>
tx.recording
.create({
data: {
tenantId,
callId,
storageProvider: STORAGE_KIND,
objectKey: "",
recordedAt: new Date(),
status: "FAILED",
},
})
.catch(() => undefined),
);
}
}