mirror of
https://github.com/melgarafael/DeskcommCRM.git
synced 2026-10-02 01:28:34 +08:00
fix(agent-engine): agendador avança o cron de organização parada sem enfileirar
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
9811bc613b
commit
fc094270dd
@@ -0,0 +1,67 @@
|
||||
/**
|
||||
* O AGENDADOR NÃO DISPARA FOLLOW-UP DE ORGANIZAÇÃO PARADA.
|
||||
*
|
||||
* O cron vencido da org suspensa/redigida/arquivada avança sem enfileirar job,
|
||||
* e o one-shot se encerra. Reativar não devolve o que venceu parado: a
|
||||
* reativação é sem rajada (spec §1.3).
|
||||
*/
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
import type pg from 'pg';
|
||||
|
||||
import { tickCron } from './scheduler';
|
||||
|
||||
const log = { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() } as never;
|
||||
const AGORA = Date.parse('2026-09-29T12:00:00Z');
|
||||
const HORA = 3_600_000;
|
||||
const cfg = { batchSize: 1, staggerWindowMs: 0, retryBaseMs: 1000, now: () => AGORA };
|
||||
|
||||
function cron(over: Record<string, unknown>) {
|
||||
return {
|
||||
id: 'cron-1', organization_id: 'org-1', contact_id: 'contato-1', kind: 'every',
|
||||
interval_ms: String(HORA), cron_expr: null, tz: 'UTC', job_kind: 'followup_turn', payload: {},
|
||||
next_run_at: new Date(AGORA - 1000), enabled: true, attempts: 0, max_attempts: 5,
|
||||
last_error: null, created_at: new Date(AGORA), updated_at: new Date(AGORA), operante: false,
|
||||
...over,
|
||||
};
|
||||
}
|
||||
|
||||
function poolCom(linha: Record<string, unknown>) {
|
||||
const sqls: Array<{ sql: string; params: unknown[] }> = [];
|
||||
const client = {
|
||||
query: vi.fn(async (sql: string, params: unknown[] = []) => {
|
||||
sqls.push({ sql, params });
|
||||
if (sql.includes('fn_org_operante')) return { rows: [linha] };
|
||||
if (sql.includes('insert into job_queue')) return { rows: [{ id: 'job-1' }] };
|
||||
return { rows: [] };
|
||||
}),
|
||||
release: vi.fn(),
|
||||
};
|
||||
return { pool: { connect: async () => client } as unknown as pg.Pool, sqls };
|
||||
}
|
||||
|
||||
describe('fireOneDue × organização parada', () => {
|
||||
it('recorrente: avança next_run_at, NÃO enfileira, e conta como skipped', async () => {
|
||||
const { pool, sqls } = poolCom(cron({}));
|
||||
const r = await tickCron(pool, cfg, log);
|
||||
expect(sqls.some((s) => s.sql.includes('insert into job_queue'))).toBe(false);
|
||||
const reagenda = sqls.find((s) => s.sql.startsWith('update cron_jobs set next_run_at'));
|
||||
expect(reagenda?.params).toEqual(['cron-1', new Date(AGORA - 1000 + HORA)]);
|
||||
expect(reagenda?.sql).toContain("last_error = 'org_nao_operante'");
|
||||
expect(sqls.some((s) => s.sql === 'commit')).toBe(true);
|
||||
expect(r).toEqual({ fired: 0, retried: 0, disabled: 0, skipped: 1 });
|
||||
});
|
||||
|
||||
it("one-shot ('at'): se encerra (enabled=false) sem enfileirar", async () => {
|
||||
const { pool, sqls } = poolCom(cron({ kind: 'at', interval_ms: null }));
|
||||
await tickCron(pool, cfg, log);
|
||||
expect(sqls.some((s) => s.sql.includes('insert into job_queue'))).toBe(false);
|
||||
expect(sqls.some((s) => s.sql.startsWith('update cron_jobs set enabled = false'))).toBe(true);
|
||||
});
|
||||
|
||||
it('org operante: enfileira como sempre (controle)', async () => {
|
||||
const { pool, sqls } = poolCom(cron({ operante: true }));
|
||||
const r = await tickCron(pool, cfg, log);
|
||||
expect(sqls.some((s) => s.sql.includes('insert into job_queue'))).toBe(true);
|
||||
expect(r.fired).toBe(1);
|
||||
});
|
||||
});
|
||||
@@ -150,6 +150,8 @@ export interface CronTickResult {
|
||||
fired: number;
|
||||
retried: number;
|
||||
disabled: number;
|
||||
/** Vencidos de organização não operante: avançados sem enfileirar. */
|
||||
skipped: number;
|
||||
}
|
||||
|
||||
type FailureOutcome = { outcome: 'retried' | 'disabled'; classification: string; attempts: number };
|
||||
@@ -198,13 +200,13 @@ async function fireOneDue(
|
||||
pool: pg.Pool,
|
||||
cfg: CronTickConfig,
|
||||
log: Logger,
|
||||
): Promise<'fired' | 'retried' | 'disabled' | 'empty'> {
|
||||
): Promise<'fired' | 'retried' | 'disabled' | 'skipped' | 'empty'> {
|
||||
const nowMs = (cfg.now ?? Date.now)();
|
||||
const client = await pool.connect();
|
||||
try {
|
||||
await client.query('begin');
|
||||
const { rows } = await client.query<CronJobRow>(
|
||||
`select * from cron_jobs
|
||||
const { rows } = await client.query<CronJobRow & { operante: boolean }>(
|
||||
`select *, public.fn_org_operante(organization_id) as operante from cron_jobs
|
||||
where enabled = true and next_run_at <= $1
|
||||
order by next_run_at
|
||||
limit 1
|
||||
@@ -217,6 +219,26 @@ async function fireOneDue(
|
||||
return 'empty';
|
||||
}
|
||||
const spec = specFromRow(cron);
|
||||
// Organização parada (suspensa, redigida, arquivada) não recebe follow-up:
|
||||
// o disparo avança sem enfileirar e o one-shot se encerra. Reativar não
|
||||
// devolve o que venceu parado — reativação é sem rajada (spec §1.3).
|
||||
if (!cron.operante) {
|
||||
const proximo = computeNextRunAt(spec, cron.next_run_at.getTime(), nowMs, cfg.staggerWindowMs, cron.contact_id);
|
||||
if (proximo === null) {
|
||||
await client.query(
|
||||
`update cron_jobs set enabled = false, last_error = 'org_nao_operante', updated_at = now() where id = $1`,
|
||||
[cron.id],
|
||||
);
|
||||
} else {
|
||||
await client.query(
|
||||
`update cron_jobs set next_run_at = $2, last_error = 'org_nao_operante', updated_at = now() where id = $1`,
|
||||
[cron.id, proximo],
|
||||
);
|
||||
}
|
||||
await client.query('commit');
|
||||
log.info('cron: organização não operante — disparo pulado sem enfileirar', { cron_job_id: cron.id });
|
||||
return 'skipped';
|
||||
}
|
||||
try {
|
||||
await client.query('savepoint fire');
|
||||
await enqueueJob(client, cron.organization_id, {
|
||||
@@ -274,7 +296,7 @@ async function fireOneDue(
|
||||
* um cron podre não derruba os demais). Para quando esgota os vencidos claimáveis.
|
||||
*/
|
||||
export async function tickCron(pool: pg.Pool, cfg: CronTickConfig, log: Logger): Promise<CronTickResult> {
|
||||
const result: CronTickResult = { fired: 0, retried: 0, disabled: 0 };
|
||||
const result: CronTickResult = { fired: 0, retried: 0, disabled: 0, skipped: 0 };
|
||||
for (let i = 0; i < cfg.batchSize; i += 1) {
|
||||
const outcome = await fireOneDue(pool, cfg, log);
|
||||
if (outcome === 'empty') break;
|
||||
@@ -302,7 +324,7 @@ export async function runCronLoop(
|
||||
while (!signal.aborted) {
|
||||
try {
|
||||
const tick = await tickCron(pool, cfg, log);
|
||||
if (tick.fired + tick.retried + tick.disabled > 0) log.info('cron: tick processado', { ...tick });
|
||||
if (tick.fired + tick.retried + tick.disabled + tick.skipped > 0) log.info('cron: tick processado', { ...tick });
|
||||
} catch (err) {
|
||||
log.error('cron: tick falhou — tenta no próximo intervalo', { error: errMsg(err) });
|
||||
}
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
import { afterAll, beforeAll, describe, expect, it } from "vitest";
|
||||
import pg from "pg";
|
||||
|
||||
import { tickCron } from "@/lib/agent-engine/cron/scheduler";
|
||||
import { createLogger } from "@/lib/agent-engine/obs/logger";
|
||||
|
||||
/**
|
||||
* O SQL REAL do agendador com a régua (migration 0492).
|
||||
*
|
||||
* `fireOneDue` chama `public.fn_org_operante(organization_id)` dentro do
|
||||
* `select … for update skip locked`, pelo pool `pg` do worker. O teste unitário
|
||||
* usa dublê; só aqui o Postgres executa o SQL com o papel do pool — um erro de
|
||||
* sintaxe ou de EXECUTE pararia TODOS os follow-ups da instalação.
|
||||
* Conexão copiada de `agent-watchdog.test.ts`.
|
||||
*/
|
||||
const container = process.env.TEST_DB_CONTAINER;
|
||||
if (!container) {
|
||||
throw new Error("TEST_DB_CONTAINER not set — rode via `pnpm test:db` (scripts/test-db.sh)");
|
||||
}
|
||||
|
||||
const PORT = Number(process.env.TEST_DB_PORT ?? 54329);
|
||||
const pool = new pg.Pool({
|
||||
connectionString: `postgresql://postgres:postgres@127.0.0.1:${PORT}/postgres`,
|
||||
max: 2,
|
||||
});
|
||||
const log = createLogger();
|
||||
|
||||
const ORG_PARADA = "c0de0492-8888-4000-8000-00000000000a";
|
||||
const ORG_ATIVA = "c0de0492-8888-4000-8000-00000000000b";
|
||||
const CONTATO = "c0de0492-8888-4000-8000-0000000000c1";
|
||||
const CRON_RECORRENTE = "c0de0492-8888-4000-8000-0000000000d1";
|
||||
const CRON_UNICO = "c0de0492-8888-4000-8000-0000000000d2";
|
||||
const HORA = 3_600_000;
|
||||
/**
|
||||
* `tickCron` reivindica QUALQUER cron vencido do banco compartilhado, na ordem
|
||||
* de `next_run_at` (`scheduler.ts:208-211`), e o `test:db` roda os arquivos em
|
||||
* ordem embaralhada (`--sequence.shuffle.files=true`). Com os NOSSOS dois crons
|
||||
* vencidos em 2000 — antes de qualquer cron que outro arquivo semeie — e
|
||||
* `batchSize: 2`, o tick pega exatamente os dois e não dispara cron alheio.
|
||||
* As asserções medem só as nossas linhas.
|
||||
*/
|
||||
const VENCIDO_EM = "2000-01-01T00:00:00Z";
|
||||
|
||||
beforeAll(async () => {
|
||||
await pool.query(`
|
||||
insert into public.organizations (id, slug, legal_name, display_name, status, suspended_kind, suspended_at) values
|
||||
('${ORG_PARADA}', 'cron-0492-parada', 'Parada', 'Parada', 'suspended', 'administrativa', now()),
|
||||
('${ORG_ATIVA}', 'cron-0492-ativa', 'Ativa', 'Ativa', 'active', null, null)
|
||||
on conflict (id) do nothing;
|
||||
insert into public.contacts (id, organization_id, display_name)
|
||||
values ('${CONTATO}', '${ORG_PARADA}', 'Contato cron 0492') on conflict (id) do nothing;
|
||||
insert into public.cron_jobs (id, organization_id, contact_id, kind, interval_ms, job_kind, next_run_at) values
|
||||
('${CRON_RECORRENTE}', '${ORG_PARADA}', '${CONTATO}', 'every', ${HORA}, 'followup_turn', '${VENCIDO_EM}'),
|
||||
('${CRON_UNICO}', '${ORG_PARADA}', '${CONTATO}', 'at', null, 'followup_turn', '${VENCIDO_EM}')
|
||||
on conflict (id) do nothing;
|
||||
`);
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
await pool.end();
|
||||
});
|
||||
|
||||
describe("tickCron × organização parada, no Postgres real", () => {
|
||||
it("controle: o papel do pool executa fn_org_operante e a régua responde", async () => {
|
||||
const { rows } = await pool.query<{ parada: boolean; ativa: boolean }>(
|
||||
"select public.fn_org_operante($1) as parada, public.fn_org_operante($2) as ativa",
|
||||
[ORG_PARADA, ORG_ATIVA],
|
||||
);
|
||||
expect(rows[0]).toEqual({ parada: false, ativa: true });
|
||||
});
|
||||
|
||||
it("os vencidos da org parada avançam/encerram sem job, com last_error org_nao_operante", async () => {
|
||||
await tickCron(pool, { batchSize: 2, staggerWindowMs: 0, retryBaseMs: 1000 }, log);
|
||||
// Só as NOSSAS linhas: nenhuma delas pode seguir vencida e habilitada.
|
||||
const { rows: pendentes } = await pool.query(
|
||||
"select count(*)::int as n from public.cron_jobs where id in ($1, $2) and enabled and next_run_at <= now()",
|
||||
[CRON_RECORRENTE, CRON_UNICO],
|
||||
);
|
||||
expect(pendentes[0].n, "o tick não reivindicou os dois crons da org parada").toBe(0);
|
||||
const { rows: jobs } = await pool.query("select count(*)::int as n from public.job_queue where organization_id = $1", [ORG_PARADA]);
|
||||
expect(jobs[0].n).toBe(0);
|
||||
const { rows: crons } = await pool.query<{ id: string; enabled: boolean; last_error: string; futuro: boolean }>(
|
||||
"select id, enabled, last_error, next_run_at > now() as futuro from public.cron_jobs where id in ($1, $2) order by id",
|
||||
[CRON_RECORRENTE, CRON_UNICO],
|
||||
);
|
||||
expect(crons).toEqual([
|
||||
{ id: CRON_RECORRENTE, enabled: true, last_error: "org_nao_operante", futuro: true },
|
||||
{ id: CRON_UNICO, enabled: false, last_error: "org_nao_operante", futuro: false },
|
||||
]);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user