Skip to content
Open
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 30 additions & 2 deletions src/collector/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,13 @@ export interface CollectorCoreHandle {

const msgRetryCounterCache = new NodeCache() as CacheStore

// Teto da espera do stop() pelo connect() em voo. O connect() contém uma chamada HTTP
// (fetchLatestBaileysVersion) que numa rede ruim fica pendurada — e o stop() não pode ficar preso
// atrás dela, porque é depois dele que o outbox é fechado (checkpoint do WAL do SQLite). Sem este
// teto, um restart durante reconexão cairia no forceMs do shutdown (saída com código 1) sem nunca
// fechar o outbox. Em teste o connect() resolve em milissegundos: quem ganha a corrida é a promise.
const STOP_CONNECT_WAIT_MS = 2_000

// converte a metadata crua do baileys no shape persistido em cache (group-metadata.json).
function toGroupInfo(meta: import('@whiskeysockets/baileys').GroupMetadata): GroupInfo {
return {
Expand Down Expand Up @@ -139,6 +146,12 @@ export function startCollectorCore(deps: CollectorCoreDeps): CollectorCoreHandle

let stopped = false
let sock: ReturnType<typeof makeWASocketReal> | null = null
// aponta pro connect() em voo mais recente (inicial ou reconexão via 'close'). stop() espera essa
// promise antes de terminar — sem isso, um connect() já passado dos awaits de auth/versão no
// momento do stop() ainda abriria socket e registraria handlers depois que o chamador já
// considerava o coletor parado (produção: reconecta ao WhatsApp mesmo "parado"; testes: o
// connect() de um caso que já terminou seguia mexendo no mesmo authDir do caso seguinte).
let connectPromise: Promise<void> = Promise.resolve()

async function fetchGroupsMetadata(activeSock: ReturnType<typeof makeWASocketReal>) {
const cache = loadGroupCache()
Expand Down Expand Up @@ -176,6 +189,11 @@ export function startCollectorCore(deps: CollectorCoreDeps): CollectorCoreHandle

const { version } = await fetchLatestBaileysVersion()

// stop() pode ter sido chamado enquanto ainda esperávamos o auth state / a versão mais
// recente — sem este guard o connect() abandonado abriria socket e registraria ev.process
// mesmo assim, "reconectando" depois que o coletor já era pra estar parado.
if (stopped) return

const activeSock = makeSocket({
version,
logger: deps.baileysLogger,
Expand Down Expand Up @@ -324,7 +342,7 @@ export function startCollectorCore(deps: CollectorCoreDeps): CollectorCoreHandle
}
}
deps.onStatus?.('connecting')
connect()
connectPromise = connect()
}
if (qr) {
deps.onQr?.(qr)
Expand All @@ -333,14 +351,24 @@ export function startCollectorCore(deps: CollectorCoreDeps): CollectorCoreHandle
})
}

connect()
connectPromise = connect()

return {
async stop() {
stopped = true
stopSender()
stopHeartbeat()
stopRetention()
// espera o connect() em voo terminar (ou abortar pelo guard acima) antes de fechar outbox
// e encerrar o socket — sem isso, stop() "concluía" enquanto uma conexão/reconexão ainda
// em andamento seguia criando authDir/socket por trás. Com teto: ver STOP_CONNECT_WAIT_MS.
await Promise.race([
connectPromise.catch(() => {}),
new Promise<void>((resolve) => {
const timer = setTimeout(resolve, STOP_CONNECT_WAIT_MS)
timer.unref?.()
}),
])
outbox?.close()
sock?.end(undefined)
},
Expand Down
66 changes: 46 additions & 20 deletions tests/collector/core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -78,8 +78,9 @@ function makeFakeSocket(opts: {
}
// emit: dispara o loop de ev.process com o lote de eventos informado.
const emit = async (events: Record<string, unknown>) => {
// connect() é async (await useMultiFileAuthState + fetchLatestBaileysVersion); damos uns ticks.
for (let i = 0; i < 100 && !processCb; i++) await new Promise((r) => setImmediate(r))
// connect() é async (await useMultiFileAuthState + fetchLatestBaileysVersion) — espera por
// tempo real, não por ticks (ver waitUntil).
await waitUntil(() => processCb !== null)
assert.ok(processCb, 'ev.process não foi registrado')
await processCb!(events)
}
Expand All @@ -90,11 +91,35 @@ const makeFakeMakeSocket = (sock: any) => (() => sock) as any

const flush = async (n = 4) => { for (let i = 0; i < n; i++) await new Promise((r) => setImmediate(r)) }

// espera até `cond()` virar true, por até `timeoutMs` de tempo REAL (não um número fixo de voltas
// do event loop). connect() faz I/O de disco de verdade (useMultiFileAuthState: stat/mkdir/readFile
// via libuv); sob carga — a suíte inteira rodando junto, ~20+ processos disputando disco/CPU — essa
// I/O pode demorar mais do que 100 `setImmediate` levam para passar (às vezes só ~2ms de relógio
// quando o event loop está ocioso). Contar tempo real em vez de tick corrige a flakiness raiz deste
// arquivo (evidência no relatório da correção: 100 ticks/2ms decorridos com useMultiFileAuthState
// ainda em voo).
async function waitUntil(cond: () => boolean, timeoutMs = 5000): Promise<boolean> {
const start = Date.now()
while (!cond()) {
if (Date.now() - start > timeoutMs) return false
await new Promise((r) => setTimeout(r, 1))
}
return true
}

// cada teste ganha sua PRÓPRIA subpasta de authDir — nunca a mesma de outro teste. Motivo: todos
// os testes usavam o mesmo 'baileys_auth_info' relativo (cwd fixo no arquivo todo), então uma
// reconexão que sobrevive ao teste que a originou (ver correção em src/collector/core.ts) ficava
// batendo no MESMO diretório/mutex de arquivo (mutex é module-level no baileys) que o teste seguinte
// — ruído de I/O concorrente que também alimentava a flakiness.
let authDirSeq = 0
const nextAuthDir = () => `auth-${++authDirSeq}/baileys_auth_info`

test('repassa cada evento p/ saveEvent (trilha NDJSON), router e onEvent', async () => {
const fake = makeFakeSocket()
const seen: Array<[string, unknown]> = []
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -128,7 +153,7 @@ test('connection.update: open dispara onStatus(connected,{me}) e busca metadata
const statuses: Array<{ s: string; me?: string }> = []
const events: string[] = []
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -158,7 +183,7 @@ test('connection.update: connecting e qr propagam via onStatus/onQr', async () =
const statuses: string[] = []
let qr: string | null = null
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand All @@ -179,10 +204,11 @@ test('connection.update: connecting e qr propagam via onStatus/onQr', async () =

test('loggedOut limpa o authDir e reconecta', async () => {
const fake = makeFakeSocket()
const authDir = path.join(tmpDir, 'baileys_auth_info')
const authRel = nextAuthDir()
const authDir = path.join(tmpDir, authRel)
const statuses: string[] = []
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: authRel,
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand All @@ -191,8 +217,8 @@ test('loggedOut limpa o authDir e reconecta', async () => {
makeSocket: makeFakeMakeSocket(fake.sock),
})

// useMultiFileAuthState cria o dir na conexão inicial — espera-o existir.
for (let i = 0; i < 100 && !fs.existsSync(authDir); i++) await new Promise((r) => setImmediate(r))
// useMultiFileAuthState cria o dir na conexão inicial — espera-o existir (tempo real, não ticks).
await waitUntil(() => fs.existsSync(authDir))
assert.ok(fs.existsSync(authDir), 'authDir deveria ter sido criado pelo connect inicial')

// DisconnectReason.loggedOut === 401
Expand All @@ -209,7 +235,7 @@ test('webhook != null liga o coletor: cria o outbox.db', async () => {
const fake = makeFakeSocket()
const outboxPath = path.join(tmpDir, 'collector-on.db')
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'collector-on.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -240,7 +266,7 @@ test('webhook == null mantém o coletor desligado (sem outbox.db)', async () =>
const fake = makeFakeSocket()
const outboxPath = path.join(tmpDir, 'collector-off.db')
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'collector-off.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand All @@ -258,7 +284,7 @@ test('webhook == null mantém o coletor desligado (sem outbox.db)', async () =>
test('stop() encerra o socket e sendMessage delega ao socket', async () => {
const fake = makeFakeSocket()
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -287,7 +313,7 @@ test('messages.upsert com /ban de admin aciona a remoção (fiação core → ba
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -323,7 +349,7 @@ test('messages.upsert com /admin on de admin aciona groupSettingUpdate (fiação
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -361,7 +387,7 @@ test('messages.upsert com /ban por telefone remove da COMUNIDADE (fiação core
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -395,7 +421,7 @@ test('group-participants.update invalida o diretório (próximo comando refaz o
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -440,7 +466,7 @@ test('messages.upsert com /kick remove só do grupo (fiação core → kick-comm
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -478,7 +504,7 @@ test('group-participants.update add de quem foi banido dispara a remoção (fia
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -529,7 +555,7 @@ test('/unban tira da denylist e a reentrada volta a ser permitida (fiação core
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down Expand Up @@ -586,7 +612,7 @@ test('comando em grupo comum é apagado e reportado no grupo de log (fiação co
},
})
const handle = startCollectorCore({
authDir: 'baileys_auth_info',
authDir: nextAuthDir(),
outboxPath: 'outbox.db',
logger: silentLogger,
baileysLogger: silentLogger,
Expand Down