feat(followup): refactor text application for follow-ups and enhance tests

- Renamed `aplicarTextoNosFollowups` to `aplicarTextoAosEnrollmentsEmEspera` for clarity and updated its parameters to improve functionality.
- Introduced a new test suite for `aplicarTextoNosFollowups`, ensuring proper application of text to follow-ups and validating the expected behavior.
- Enhanced existing tests to verify the integration of the new text application logic, improving overall test coverage and reliability.
This commit is contained in:
Ian Couto
2026-08-25 19:23:27 -03:00
parent 249eedae08
commit f88aa227bf
4 changed files with 83 additions and 11 deletions
+11
View File
@@ -1,3 +1,6 @@
import { readFileSync } from "node:fs";
import { join } from "node:path";
import { describe, expect, it } from "vitest";
import { inboundEhDestaPergunta, textoDoPayloadInbound } from "./aplicar-inbound";
@@ -25,3 +28,11 @@ describe("inboundEhDestaPergunta", () => {
expect(inboundEhDestaPergunta("2026-08-25T19:25:00.000Z", "2026-08-25T19:19:00.000Z")).toBe(true);
});
});
describe("aplicarTextoNosFollowups — segunda passada", () => {
it("reaplica waiting_reply depois de avançar active (corrida cap_nome)", () => {
const fonte = readFileSync(join(process.cwd(), "lib/followup/aplicar-inbound.ts"), "utf8");
const chamadas = fonte.match(/await aplicarTextoAosEnrollmentsEmEspera\(/g) ?? [];
expect(chamadas.length).toBe(2);
});
});
+24 -11
View File
@@ -72,6 +72,25 @@ async function ultimoInboundDoContato(
return typeof data?.body === "string" ? data.body.trim() : "";
}
async function aplicarTextoAosEnrollmentsEmEspera(
admin: SupabaseClient,
orgId: string,
contactIds: string[],
texto: string,
deps: TickDeps,
): Promise<void> {
const { data, error } = await admin
.from("followup_enrollments")
.select("*")
.eq("organization_id", orgId)
.in("contact_id", contactIds)
.eq("status", "waiting_reply");
if (error) throw new Error(error.message);
for (const row of data ?? []) {
await aplicarRespostaInbound(deps, row as EnrollmentRow, texto);
}
}
export async function aplicarTextoNosFollowups(
admin: SupabaseClient,
sinal: SinalDeInboundFollowup,
@@ -79,18 +98,9 @@ export async function aplicarTextoNosFollowups(
const contactIds = await idsDoContatoEGemeos(admin, sinal.organizationId, sinal.contactId);
const texto = (sinal.texto?.trim() || (await ultimoInboundDoContato(admin, sinal.organizationId, contactIds))).trim();
if (!texto) return;
const { data, error } = await admin
.from("followup_enrollments")
.select("*")
.eq("organization_id", sinal.organizationId)
.in("contact_id", contactIds)
.in("status", ["waiting_reply", "active"]);
if (error) throw new Error(error.message);
const deps = tickDepsDe(admin);
for (const row of data ?? []) {
if (row.status !== "waiting_reply") continue;
await aplicarRespostaInbound(deps, row as EnrollmentRow, texto);
}
// 1ª passada: resposta chegou com o enrollment já em waiting_reply.
await aplicarTextoAosEnrollmentsEmEspera(admin, sinal.organizationId, contactIds, texto, deps);
for (let i = 0; i < 6; i++) {
const agora = new Date().toISOString();
const { data: vivos, error: vivosErr } = await admin
@@ -108,4 +118,7 @@ export async function aplicarTextoNosFollowups(
const enviados = await enviarTextoFixoPendente(admin, contactIds);
if (!(vivos?.length) && !enviados) break;
}
// 2ª passada: o SIM pode ter chegado enquanto o nó ainda era `active` (ex.:
// cap_nome enfileirando a pergunta de confirmação neste mesmo request).
await aplicarTextoAosEnrollmentsEmEspera(admin, sinal.organizationId, contactIds, texto, deps);
}
+40
View File
@@ -14,6 +14,7 @@ import { createHmac, timingSafeEqual } from "node:crypto";
import { audit } from "@/lib/audit";
import { sincronizarSaudeDaConexao } from "@/lib/channels/health";
import { aplicarEfeitosPosEntrada } from "@/lib/channels/pos-entrada";
import { acelerarPipelineDeEventos } from "@/lib/dev/kick-local-pipeline";
import { canonicalPhoneBR } from "@/lib/channels/phone-variants";
import { estamparAtribuicaoDoContato } from "@/lib/leads/atribuicao-de-anuncio";
import { extrairAtribuicaoWaha } from "@/lib/waha/atribuicao-de-anuncio";
@@ -459,6 +460,25 @@ async function markConversation(
/**
* Mensagem recebida (fromMe=false). Contato = remetente (`from`).
*/
async function mensagemIngeridaPorExternalId(
admin: Admin,
orgId: string,
externalId: string,
): Promise<{ id: string; contact_id: string; body: string | null } | null> {
const { data, error } = await admin
.from("messages")
.select("id, contact_id, body")
.eq("organization_id", orgId)
.eq("external_id", externalId)
.eq("direction", "inbound")
.maybeSingle();
if (error) {
logger.warn("waha.ingest: dedup sem ler mensagem existente", { detail: error.message });
return null;
}
return data ?? null;
}
async function handleInbound(
admin: Admin,
session: Session,
@@ -548,6 +568,26 @@ async function handleInbound(
external_id: p.id,
direcao: "inbound",
});
// A 1ª entrega pode ter gravado a mensagem e estourado o tempo ANTES de
// `aplicarEfeitosPosEntrada` — a reentrega cai aqui. Reacelerar só o
// pipeline (sem re-despachar o agente) destrava o match_reply.
const existente = await mensagemIngeridaPorExternalId(admin, session.organization_id, p.id);
if (existente) {
try {
await acelerarPipelineDeEventos(admin, {
organizationId: session.organization_id,
contactId: existente.contact_id,
messageId: existente.id,
texto: existente.body,
});
} catch (err) {
logger.warn("waha.ingest: dedup nao reacelerou pipeline", {
organization_id: session.organization_id,
external_id: p.id,
detail: err instanceof Error ? err.message : String(err),
});
}
}
return;
}
@@ -103,4 +103,12 @@ describe("dedup do ingest deixa rastro", () => {
);
}
});
it("dedup inbound reacelera o pipeline de follow-up (âncora no fonte)", () => {
const fonte = readFileSync(join(process.cwd(), "lib/waha/ingest.ts"), "utf8");
const inbound = fonte.split('if (insertErr?.code === "23505")')[1] ?? "";
expect(inbound.slice(0, 2000), "reentrega WAHA tem que destravar match_reply").toMatch(
/acelerarPipelineDeEventos/,
);
});
});