atendimentosService.js 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  1. import { Model } from "objection";
  2. import { config } from "../config/index.js";
  3. import { Atendimento } from "../models/Atendimento.model.js";
  4. import { AtendimentoCliente } from "../models/AtendimentoCliente.model.js";
  5. import { AtendimentoMensagem } from "../models/AtendimentoMensagem.model.js";
  6. import {
  7. atendimentoPayloadSchema,
  8. protocolosAbertosSchema,
  9. protocoloMensagensSchema
  10. } from "../middleware/schemas/Atendimento.Schema.js";
  11. import { assertPublicHttpsUrl } from "./ingestService.js";
  12. import { BadRequestError } from "../shared/errors/index.js";
  13. const MAX_REDIRECTS = 5;
  14. const MENSAGENS_BATCH = 500;
  15. const SYNC_CONCURRENCY = 4;
  16. function ifbotAuthHeaders() {
  17. const { token, authHeader } = config.ifbot;
  18. if (!token) return {};
  19. const value = authHeader === "Authorization" && !token.includes(" ") ? `Bearer ${token}` : token;
  20. return { [authHeader]: value };
  21. }
  22. async function fetchJson(urlStr, extraHeaders = {}) {
  23. let currentUrl = urlStr;
  24. let res;
  25. for (let hop = 0; ; hop += 1) {
  26. assertPublicHttpsUrl(currentUrl);
  27. res = await fetch(currentUrl, {
  28. headers: {
  29. "User-Agent": "Mozilla/5.0 star-oraculo/1.0",
  30. Accept: "application/json",
  31. ...extraHeaders
  32. },
  33. redirect: "manual",
  34. signal: AbortSignal.timeout(15_000)
  35. });
  36. if (res.status < 300 || res.status >= 400) break;
  37. const location = res.headers.get("location");
  38. if (!location || hop >= MAX_REDIRECTS) {
  39. const err = new Error(`url_fetch_error:${res.status}`);
  40. err.statusCode = 502;
  41. throw err;
  42. }
  43. currentUrl = new URL(location, currentUrl).toString();
  44. }
  45. if (!res.ok) {
  46. const err = new Error(`url_fetch_error:${res.status}`);
  47. err.statusCode = 502;
  48. throw err;
  49. }
  50. try {
  51. return await res.json();
  52. } catch {
  53. const err = new Error("url_invalid_json");
  54. err.statusCode = 502;
  55. throw err;
  56. }
  57. }
  58. export async function fetchAtendimentoJson(urlStr) {
  59. const json = await fetchJson(urlStr);
  60. const parsed = atendimentoPayloadSchema.safeParse(json);
  61. if (!parsed.success) {
  62. throw new BadRequestError("invalid_atendimento_payload");
  63. }
  64. return parsed.data;
  65. }
  66. // prefixo do código do protocolo identifica o setor (ex.: SUP0000014690/2026 → SUP)
  67. export function setorFromCodigo(codigo) {
  68. const m = /^([A-Z]{3})\d/.exec(String(codigo ?? ""));
  69. return m ? m[1] : null;
  70. }
  71. function toDate(value) {
  72. if (!value) return null;
  73. const d = new Date(value);
  74. return Number.isNaN(d.getTime()) ? null : d;
  75. }
  76. export async function salvarAtendimento({ payload, sourceUrl, abertura, telefone }) {
  77. const now = new Date();
  78. const trx = await Model.startTransaction();
  79. try {
  80. const c = payload.Cliente;
  81. const clienteMerge = ["Nome", "Foto", "Wid", "Ultima2", "Ultimacliente2", "Deadmensage2", "UpdatedAt"];
  82. const clienteTelefone = telefone ?? c.Telefone ?? null;
  83. if (clienteTelefone) clienteMerge.push("Telefone");
  84. await AtendimentoCliente.query(trx)
  85. .insert({
  86. Id: c.Id,
  87. Nome: c.Nome,
  88. Foto: c.Foto ?? null,
  89. Wid: c.Wid ?? null,
  90. Telefone: clienteTelefone,
  91. Ultima2: c.Ultima2 ?? null,
  92. Ultimacliente2: c.Ultimacliente2 ?? null,
  93. Deadmensage2: c.Deadmensage2 ?? null,
  94. UpdatedAt: now
  95. })
  96. .onConflict("Id")
  97. .merge(clienteMerge);
  98. const p = payload.Protocolo;
  99. const setor = setorFromCodigo(p.Codigo);
  100. const atendimentoMerge = ["Codigo", "Status", "ClienteId", "Setor", "IngestedAt", "UpdatedAt"];
  101. if (sourceUrl) atendimentoMerge.push("SourceUrl");
  102. if (abertura) atendimentoMerge.push("Abertura");
  103. await Atendimento.query(trx)
  104. .insert({
  105. Id: p.Id,
  106. Codigo: p.Codigo,
  107. Status: p.Status,
  108. ClienteId: c.Id,
  109. Setor: setor,
  110. SourceUrl: sourceUrl ?? null,
  111. Abertura: toDate(abertura),
  112. IngestedAt: now,
  113. UpdatedAt: now
  114. })
  115. .onConflict("Id")
  116. .merge(atendimentoMerge);
  117. const mensagens = (payload.Mensagens ?? []).map((m) => ({
  118. Id: m.Id,
  119. AtendimentoId: p.Id,
  120. Body: m.Body ?? null,
  121. Resposta: m.Resposta ?? null,
  122. Timestamp: toDate(m.Timestamp),
  123. Tipodemidia: m.Tipodemidia ?? null,
  124. Midia: m.Midia ?? null,
  125. Autor: m.Autor ?? null,
  126. Citacao: m.Citacao ?? null,
  127. Transcricao: m.Transcricao ?? null
  128. }));
  129. // batch insert via knex: o Objection não suporta insert em lote no MySQL
  130. for (let i = 0; i < mensagens.length; i += MENSAGENS_BATCH) {
  131. await trx(AtendimentoMensagem.tableName)
  132. .insert(mensagens.slice(i, i + MENSAGENS_BATCH))
  133. .onConflict("Id")
  134. .merge(["AtendimentoId", "Body", "Resposta", "Timestamp", "Tipodemidia", "Midia", "Autor", "Citacao", "Transcricao"]);
  135. }
  136. await trx.commit();
  137. return {
  138. protocoloId: p.Id,
  139. codigo: p.Codigo,
  140. clienteId: c.Id,
  141. mensagens: mensagens.length
  142. };
  143. } catch (err) {
  144. await trx.rollback();
  145. throw err;
  146. }
  147. }
  148. async function importarProtocolo(item, baseUrl) {
  149. const mensagensUrl = `${baseUrl}/protocolos/${item.Id}/mensagens`;
  150. const json = await fetchJson(mensagensUrl, ifbotAuthHeaders());
  151. // resposta completa (Protocolo + Cliente + Mensagens) ou só a lista de mensagens
  152. let payload;
  153. const full = atendimentoPayloadSchema.safeParse(json);
  154. if (full.success) {
  155. payload = full.data;
  156. } else {
  157. const soMensagens = protocoloMensagensSchema.safeParse(json);
  158. if (!soMensagens.success) throw new BadRequestError("invalid_mensagens_payload");
  159. payload = {
  160. Protocolo: { Id: item.Id, Codigo: item.Codigo, Status: 1 },
  161. Cliente: item.Cliente,
  162. Mensagens: soMensagens.data.Mensagens
  163. };
  164. }
  165. return salvarAtendimento({
  166. payload,
  167. sourceUrl: mensagensUrl,
  168. abertura: item.Abertura ?? null,
  169. telefone: item.Cliente?.Telefone ?? null
  170. });
  171. }
  172. export async function sincronizarProtocolosAbertos({ pagina = 1, limite = 50, inicio, fim } = {}) {
  173. const baseUrl = config.ifbot.baseUrl.replace(/\/+$/, "");
  174. const params = new URLSearchParams({ pagina: String(pagina), limite: String(limite) });
  175. if (inicio) params.set("inicio", inicio);
  176. if (fim) params.set("fim", fim);
  177. const listaJson = await fetchJson(`${baseUrl}/protocolos/abertos?${params.toString()}`, ifbotAuthHeaders());
  178. const lista = protocolosAbertosSchema.safeParse(listaJson);
  179. if (!lista.success) throw new BadRequestError("invalid_protocolos_payload");
  180. const protocolos = lista.data.Protocolos;
  181. const erros = [];
  182. let importados = 0;
  183. let mensagens = 0;
  184. for (let i = 0; i < protocolos.length; i += SYNC_CONCURRENCY) {
  185. const chunk = protocolos.slice(i, i + SYNC_CONCURRENCY);
  186. const resultados = await Promise.allSettled(chunk.map((item) => importarProtocolo(item, baseUrl)));
  187. resultados.forEach((r, idx) => {
  188. if (r.status === "fulfilled") {
  189. importados += 1;
  190. mensagens += r.value.mensagens;
  191. } else {
  192. erros.push({
  193. protocoloId: chunk[idx].Id,
  194. codigo: chunk[idx].Codigo,
  195. erro: r.reason?.message ?? "erro_desconhecido"
  196. });
  197. }
  198. });
  199. }
  200. return {
  201. total: lista.data.Total ?? protocolos.length,
  202. pagina,
  203. limite,
  204. processados: protocolos.length,
  205. importados,
  206. mensagens,
  207. erros
  208. };
  209. }
  210. let fullSyncEmAndamento = false;
  211. export async function sincronizarTodosProtocolosAbertos({ limite = 200, maxPaginas = 50 } = {}) {
  212. if (fullSyncEmAndamento) {
  213. const err = new Error("sync_em_andamento");
  214. err.statusCode = 409;
  215. throw err;
  216. }
  217. fullSyncEmAndamento = true;
  218. try {
  219. const agregado = { total: 0, paginas: 0, processados: 0, importados: 0, mensagens: 0, erros: [] };
  220. for (let pagina = 1; pagina <= maxPaginas; pagina += 1) {
  221. const r = await sincronizarProtocolosAbertos({ pagina, limite });
  222. agregado.total = r.total;
  223. agregado.paginas = pagina;
  224. agregado.processados += r.processados;
  225. agregado.importados += r.importados;
  226. agregado.mensagens += r.mensagens;
  227. agregado.erros.push(...r.erros);
  228. if (r.processados < limite || agregado.processados >= agregado.total) break;
  229. }
  230. return agregado;
  231. } finally {
  232. fullSyncEmAndamento = false;
  233. }
  234. }
  235. export async function listarAtendimentos({ page = 1, pageSize = 20, setor, status } = {}) {
  236. const query = Atendimento.query()
  237. .withGraphFetched("cliente")
  238. .orderBy("IngestedAt", "desc")
  239. .page(page - 1, pageSize);
  240. if (setor) query.where("Setor", setor);
  241. if (status !== undefined) query.where("Status", status);
  242. const { results, total } = await query;
  243. return { results, total, page, pageSize };
  244. }
  245. export async function obterAtendimento(id) {
  246. return Atendimento.query()
  247. .findById(id)
  248. .withGraphFetched("[cliente, mensagens, avaliacao]")
  249. .modifyGraph("mensagens", (q) => q.orderBy("Timestamp", "asc").orderBy("Id", "asc"));
  250. }