From 2b92dfe5a67c55dab1557707c9339810d020e0ae Mon Sep 17 00:00:00 2001 From: Ruhan Freitas Date: Fri, 21 Aug 2026 16:08:33 -0300 Subject: [PATCH] =?UTF-8?q?feat(scraper):=20implementar=20lock=20distribu?= =?UTF-8?q?=C3=ADdo=20para=20impedir=20execu=C3=A7=C3=B5es=20simult=C3=A2n?= =?UTF-8?q?eas=20(PAV-119)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 2 + BACKEND.md | 5 + SCRAPER.md | 40 +- .../modules/admin/scrapers/scraperClient.ts | 20 +- .../admin/scrapers/scrapers.controller.ts | 39 +- .../admin/scrapers/scrapers.service.ts | 11 +- .../modules/admin/scrapers/scrapers.types.ts | 11 + .../tests/unit/modules/admin/scrapers.test.ts | 145 ++++++-- docker-compose.yml | 6 +- scraper-go/cmd/server/admin_handlers.go | 76 +++- scraper-go/cmd/server/admin_handlers_test.go | 144 ++++++++ scraper-go/cmd/server/handlers.go | 24 +- scraper-go/cmd/server/main.go | 2 + scraper-go/cmd/server/server.go | 25 +- scraper-go/internal/config/config.go | 82 ++++- scraper-go/internal/config/config_test.go | 71 +++- scraper-go/internal/cronjob/cronjob.go | 159 ++++++-- scraper-go/internal/cronjob/cronjob_test.go | 267 +++++++++++++- scraper-go/internal/pipeline/pipeline.go | 26 +- scraper-go/internal/pipeline/search.go | 35 +- .../internal/pipeline/search_lock_test.go | 223 +++++++++++ .../internal/pipeline/source_schedule_test.go | 43 +++ scraper-go/internal/runlock/runlock.go | 276 ++++++++++++++ scraper-go/internal/runlock/runlock_test.go | 346 ++++++++++++++++++ scraper-go/internal/runlock/valkey.go | 137 +++++++ 25 files changed, 2108 insertions(+), 107 deletions(-) create mode 100644 scraper-go/cmd/server/admin_handlers_test.go create mode 100644 scraper-go/internal/pipeline/search_lock_test.go create mode 100644 scraper-go/internal/runlock/runlock.go create mode 100644 scraper-go/internal/runlock/runlock_test.go create mode 100644 scraper-go/internal/runlock/valkey.go diff --git a/.env.example b/.env.example index 1223345..a060c46 100644 --- a/.env.example +++ b/.env.example @@ -31,6 +31,8 @@ SEARCH_KEYWORDS=UX Designer,UI Designer,Product Manager,Product Owner # Scraping behavior SCRAPER_MAX_CONCURRENCY=12 +SCRAPER_RUN_LOCK_TTL=120s +SCRAPER_RUN_LOCK_RENEW_INTERVAL=30s GOMAXPROCS=2 GOMEMLIMIT=1500MiB WAIT_BETWEEN_SEARCHES_MS=5000 diff --git a/BACKEND.md b/BACKEND.md index 126b0de..9ab5a1a 100644 --- a/BACKEND.md +++ b/BACKEND.md @@ -201,6 +201,10 @@ Base: `/` - `PATCH /admin/users/:id/unblock` — desbloqueia usuário. - `POST /admin/users/:id/reset` — reseta credenciais/senha conforme regra do serviço. - `POST /admin/scrapers/run` — dispara execução dos scrapers. + - Sucesso: `202` com `{ ok: true, message }`. + - Execução concorrente: `409` com `{ ok: false, code: "SCRAPER_ALREADY_RUNNING", message }`. + - Lock/Valkey indisponível: `503` com `{ ok: false, code: "SCRAPER_RUN_LOCK_UNAVAILABLE", message }`. + - `POST /admin/scrapers/:id/run` — aplica o mesmo contrato ao scraper nomeado (`go-scraper`). - `GET /admin/observability/metrics` — visão de métricas administrativas. - `GET /admin/observability/dashboards` — lista dashboards de observabilidade. - `GET /admin/audit` — consulta logs de auditoria. @@ -250,6 +254,7 @@ Definidas/consumidas em `src/config.ts` e outros módulos: - `goScraper.ts` faz POST em `${GO_SCRAPER_URL}/scrape` com `ScrapeParams` e valida `ScrapeResponse`. - `goKeywords.ts` consulta e publica keywords via endpoints do serviço Go (`/api/keywords`). - O backend lê os índices criados pelo scraper no Valkey, incluindo `scraper:jobs:keyword:*`, `scraper:jobs:family:*`, `scraper:jobs:technology:*` e `scraper:jobs:seniority:*`. +- Disparos administrativos usam `scraperClient` e preservam os códigos operacionais do serviço Go. O código `SCRAPER_ALREADY_RUNNING` é um conflito esperado; `SCRAPER_RUN_LOCK_UNAVAILABLE` indica política fail-closed e não inicia coleta. ## Banco de dados diff --git a/SCRAPER.md b/SCRAPER.md index 2378208..089f04c 100644 --- a/SCRAPER.md +++ b/SCRAPER.md @@ -193,7 +193,7 @@ go run ./cmd/server Docker: há um `Dockerfile` em `scraper-go/`. No Docker Compose, configure `VALKEY_URL=redis://valkey:6379/0` no `.env` da raiz para que o scraper acesse o Valkey pelo nome do serviço na rede Docker. -No Compose da raiz, o serviço escuta em . +No Compose da raiz, a porta `8081` fica exposta apenas na rede interna `vagas-net`; ela não é publicada no host. O backend acessa o serviço por `http://scraper-go:8081`. Para testes locais diretos, execute o binário fora do Compose ou use um override de desenvolvimento que publique a porta somente em interface confiável. ### Limites globais de execução @@ -206,7 +206,26 @@ O scraper possui um orçamento global de concorrência por execução controlado - `POST /scrape` preserva o contrato atual: quando `maxConcurrency` não é informado, ou vem como `0`/negativo, usa o limite global; quando vem positivo abaixo do teto, usa o valor solicitado; quando vem acima do teto, usa o teto global. - A concorrência efetiva é calculada antes da chave de cache e é o mesmo valor usado pelo pipeline, logs e semáforo. -O semáforo global atual é criado uma vez por chamada do pipeline. Portanto, o limite é por execução: duas execuções simultâneas ainda podem possuir dois semáforos independentes com a mesma capacidade. Lock entre cron/manual, prevenção de simultaneidade e limites por provider pertencem às próximas sub-issues. +O semáforo global atual é criado uma vez por chamada do pipeline. O lock distribuído abaixo impede que duas execuções mantenham semáforos independentes ao mesmo tempo; limites específicos por provider permanecem para uma sub-issue posterior. + +### Lock distribuído de execução + +Todas as origens que iniciam adapters compartilham o lock `scraper:run:lock` no Valkey: + +- cron (`source=cron`); +- disparo administrativo (`source=admin_manual`); +- cache miss de `POST /scrape` (`source=public_endpoint`). + +Cache hits de `POST /scrape` não executam adapters e, por isso, não adquirem o lock. A aquisição usa `SET ... NX PX` com token aleatório por execução. Renovação e liberação usam scripts Lua que comparam o token; não existe `DEL` incondicional. O estado informativo fica no hash `scraper:run:state`, com `runId`, `source`, `startedAt` e `lockExpiresAt`, e possui o mesmo TTL do lock. + +Configuração: + +- `SCRAPER_RUN_LOCK_TTL`: padrão `120s`; +- `SCRAPER_RUN_LOCK_RENEW_INTERVAL`: padrão `30s`, obrigatoriamente positivo e menor que o TTL. + +O mecanismo é fail-closed: se o Valkey não confirmar a aquisição, nenhum adapter é iniciado. Erros temporários de renovação são tolerados até a margem segura; perda confirmada do token ou ausência de confirmação antes dessa margem cancela o contexto da execução. A liberação ocorre no encerramento e só remove chaves pertencentes ao token atual; em crash abrupto, o TTL é a proteção final. No graceful shutdown, o scheduler deixa de aceitar novos disparos e aguarda as execuções ativas liberarem o lock antes do processo encerrar. + +O token proprietário nunca é gravado no estado operacional nem nos logs. Um `runId` independente identifica a execução para observabilidade sem expor a credencial usada pelos scripts de renovação e liberação. No Docker Compose de produção, o serviço `scraper-go` também define: @@ -224,7 +243,7 @@ O serviço expõe endpoints HTTP (implementação em `cmd/server` e arquivos ass - GET `/metrics` — métricas Prometheus. - GET `/api/keywords` — retorna as keywords atualmente carregadas. - POST `/api/keywords` — atualiza/persiste as keywords (aceita `keywords: string[]`). -- POST `/admin/scrape` — dispara uma execução manual em background; retorna 409 se já houver execução em andamento. +- POST `/admin/scrape` — dispara uma execução manual em background; retorna `409` com `SCRAPER_ALREADY_RUNNING` se já houver execução e `503` com `SCRAPER_RUN_LOCK_UNAVAILABLE` se o Valkey não confirmar a aquisição. - GET `/admin/scrape/status` — informa se existe uma execução em andamento. - GET `/admin/jobs/count` — retorna a quantidade de vagas persistidas no Valkey. - GET `/admin/jobs` — lista uma amostra das vagas persistidas no Valkey; aceita `limit`. @@ -356,6 +375,8 @@ Confirme no serviço `scraper-go` os equivalentes de `SCRAPER_MAX_CONCURRENCY=12 - `VALKEY_URL` — conexão Redis/Valkey. Em Docker Compose, use `redis://valkey:6379/0`; em execução local fora do Docker, use uma URL acessível pelo host, por exemplo `redis://localhost:6379/0`. - `SCRAPER_MAX_CONCURRENCY` — teto global de concorrência por execução. Padrão: `12`. Configuração explícita inválida impede a inicialização. +- `SCRAPER_RUN_LOCK_TTL` — duração do lock distribuído. Padrão: `120s`. +- `SCRAPER_RUN_LOCK_RENEW_INTERVAL` — intervalo de renovação. Padrão: `30s`; deve ser menor que `SCRAPER_RUN_LOCK_TTL`. - `GOMAXPROCS` — limite efetivo de threads executando código Go simultaneamente. Valor inicial no Compose: `2`. - `GOMEMLIMIT` — meta de memória do runtime/GC. Valor inicial no Compose: `1500MiB`; não substitui `mem_limit` do container. - `JOOBLE_API_KEY` — Jooble integration. @@ -376,3 +397,16 @@ Confirme no serviço `scraper-go` os equivalentes de `SCRAPER_MAX_CONCURRENCY=12 - Projetado para rodar frequentemente; use caching e indexação para reduzir chamadas repetidas. - Monitorar erros 429 e ajustar `WaitBetweenSearchesMs` / semáforos por adaptador. - Verifique logs estruturados (slog JSON) e `/metrics` para métricas de sucesso/falhas por adaptador. +- Logs `scraper run lock acquired`, `scraper execution skipped`, `scraper run lock lost` e `scraper run lock released` identificam `source` e `run_id`. + +### Verificação operacional do lock + +Durante uma execução controlada: + +```bash +docker exec vagas-valkey valkey-cli GET scraper:run:lock +docker exec vagas-valkey valkey-cli PTTL scraper:run:lock +docker exec vagas-valkey valkey-cli HGETALL scraper:run:state +``` + +Uma segunda execução manual deve retornar conflito sem iniciar adapters. Após o término, as duas chaves devem desaparecer. Nunca remova a chave manualmente apenas porque ela existe: primeiro confirme que não há processo correspondente ativo e registre valor e TTL. diff --git a/backend/src/modules/admin/scrapers/scraperClient.ts b/backend/src/modules/admin/scrapers/scraperClient.ts index 69caf52..b118c9e 100644 --- a/backend/src/modules/admin/scrapers/scraperClient.ts +++ b/backend/src/modules/admin/scrapers/scraperClient.ts @@ -11,12 +11,25 @@ const TIMEOUT_MS = 5000; /** Lançado quando o Go scraper responde 409 (já em execução) */ export class ScraperAlreadyRunningError extends Error { - constructor(message = "scraper já está em execução") { + readonly code = "SCRAPER_ALREADY_RUNNING"; + + constructor(message = "Já existe uma execução do scraper em andamento.") { super(message); this.name = "ScraperAlreadyRunningError"; } } +export class ScraperRunLockUnavailableError extends Error { + readonly code = "SCRAPER_RUN_LOCK_UNAVAILABLE"; + + constructor( + message = "Não foi possível confirmar a disponibilidade do scraper.", + ) { + super(message); + this.name = "ScraperRunLockUnavailableError"; + } +} + async function request(path: string, init?: RequestInit): Promise { const response = await fetch(`${config.scraperUrl}${path}`, { ...init, @@ -28,6 +41,11 @@ async function request(path: string, init?: RequestInit): Promise { throw new ScraperAlreadyRunningError(body?.message); } + if (response.status === 503) { + const body = await response.json().catch(() => null); + throw new ScraperRunLockUnavailableError(body?.message); + } + if (!response.ok) { throw new Error(`scraper respondeu HTTP ${response.status}`); } diff --git a/backend/src/modules/admin/scrapers/scrapers.controller.ts b/backend/src/modules/admin/scrapers/scrapers.controller.ts index a37af29..b2ecf78 100644 --- a/backend/src/modules/admin/scrapers/scrapers.controller.ts +++ b/backend/src/modules/admin/scrapers/scrapers.controller.ts @@ -1,6 +1,9 @@ import type { Request, Response } from "express"; import type { AuditService } from "../audit/audit.service"; -import { ScraperAlreadyRunningError } from "./scraperClient"; +import { + ScraperAlreadyRunningError, + ScraperRunLockUnavailableError, +} from "./scraperClient"; import { ScrapersService } from "./scrapers.service"; export class ScrapersController { @@ -18,7 +21,19 @@ export class ScrapersController { res.status(202).json(result); } catch (error) { if (error instanceof ScraperAlreadyRunningError) { - res.status(409).json({ ok: false, message: error.message }); + res.status(409).json({ + ok: false, + code: error.code, + message: error.message, + }); + return; + } + if (error instanceof ScraperRunLockUnavailableError) { + res.status(503).json({ + ok: false, + code: error.code, + message: error.message, + }); return; } res.status(500).json({ ok: false, message: "erro ao iniciar scraper" }); @@ -38,12 +53,25 @@ export class ScrapersController { res.status(202).json({ ...result, scraper: scraperName }); } catch (error) { if (error instanceof ScraperAlreadyRunningError) { - res.status(409).json({ ok: false, message: error.message }); + res.status(409).json({ + ok: false, + code: error.code, + message: error.message, + }); + return; + } + if (error instanceof ScraperRunLockUnavailableError) { + res.status(503).json({ + ok: false, + code: error.code, + message: error.message, + }); return; } const message = - error instanceof Error && error.message.startsWith("scraper desconhecido") + error instanceof Error && + error.message.startsWith("scraper desconhecido") ? error.message : "erro ao iniciar scraper"; @@ -87,7 +115,8 @@ export class ScrapersController { async listJobs(req: Request, res: Response): Promise { try { const rawLimit = Number(req.query?.limit); - const limit = Number.isFinite(rawLimit) && rawLimit > 0 ? rawLimit : undefined; + const limit = + Number.isFinite(rawLimit) && rawLimit > 0 ? rawLimit : undefined; const result = await this.scrapersService.getJobs(limit); this.auditService.fromRequest(req, "scrapers.read", { diff --git a/backend/src/modules/admin/scrapers/scrapers.service.ts b/backend/src/modules/admin/scrapers/scrapers.service.ts index 657eb02..7c8c262 100644 --- a/backend/src/modules/admin/scrapers/scrapers.service.ts +++ b/backend/src/modules/admin/scrapers/scrapers.service.ts @@ -1,4 +1,8 @@ -import { ScraperAlreadyRunningError, scraperClient } from "./scraperClient"; +import { + ScraperAlreadyRunningError, + ScraperRunLockUnavailableError, + scraperClient, +} from "./scraperClient"; import type { AdminScraper, GetJobsResult, @@ -13,7 +17,10 @@ export class ScrapersService { try { return await scraperClient.triggerScrape(); } catch (error) { - if (error instanceof ScraperAlreadyRunningError) { + if ( + error instanceof ScraperAlreadyRunningError || + error instanceof ScraperRunLockUnavailableError + ) { // repropaga para o controller decidir o status HTTP (409) throw error; } diff --git a/backend/src/modules/admin/scrapers/scrapers.types.ts b/backend/src/modules/admin/scrapers/scrapers.types.ts index 81eadec..c9b2cfc 100644 --- a/backend/src/modules/admin/scrapers/scrapers.types.ts +++ b/backend/src/modules/admin/scrapers/scrapers.types.ts @@ -25,6 +25,17 @@ export const TriggerScrapeResultSchema = z.object({ }); export type TriggerScrapeResult = z.infer; +export const TriggerScrapeErrorSchema = z.object({ + ok: z.literal(false), + code: z.enum([ + "SCRAPER_ALREADY_RUNNING", + "SCRAPER_RUN_LOCK_UNAVAILABLE", + "SCRAPER_RUN_LOCK_LOST", + ]), + message: z.string(), +}); +export type TriggerScrapeError = z.infer; + // --- ScraperStatus --- export const ScraperStatusSchema = z.object({ name: z.string().optional(), diff --git a/backend/tests/unit/modules/admin/scrapers.test.ts b/backend/tests/unit/modules/admin/scrapers.test.ts index 9f3bc0d..a957746 100644 --- a/backend/tests/unit/modules/admin/scrapers.test.ts +++ b/backend/tests/unit/modules/admin/scrapers.test.ts @@ -2,6 +2,7 @@ import type { Request, Response } from "express"; import { beforeEach, describe, expect, it, vi } from "vitest"; import { ScraperAlreadyRunningError, + ScraperRunLockUnavailableError, scraperClient, } from "../../../../src/modules/admin/scrapers/scraperClient"; import { ScrapersController } from "../../../../src/modules/admin/scrapers/scrapers.controller"; @@ -54,7 +55,11 @@ describe("scraperClient", () => { vi.fn().mockResolvedValue({ ok: false, status: 409, - json: () => Promise.resolve({ message: "busy" }), + json: () => + Promise.resolve({ + code: "SCRAPER_ALREADY_RUNNING", + message: "busy", + }), }), ); @@ -63,6 +68,27 @@ describe("scraperClient", () => { ); }); + it("throws ScraperRunLockUnavailableError for 503 responses", async () => { + vi.stubGlobal( + "fetch", + vi.fn().mockResolvedValue({ + ok: false, + status: 503, + json: () => + Promise.resolve({ + code: "SCRAPER_RUN_LOCK_UNAVAILABLE", + message: "lock unavailable", + }), + }), + ); + + await expect(scraperClient.triggerScrape()).rejects.toMatchObject({ + name: "ScraperRunLockUnavailableError", + code: "SCRAPER_RUN_LOCK_UNAVAILABLE", + message: "lock unavailable", + }); + }); + it("throws generic errors for non-ok responses", async () => { vi.stubGlobal( "fetch", @@ -131,7 +157,10 @@ describe("ScrapersService", () => { lastRunAt: "2026-07-02T10:00:00.000Z", }); vi.spyOn(scraperClient, "getJobsCount").mockResolvedValue({ total: 9 }); - vi.spyOn(scraperClient, "getJobs").mockResolvedValue({ jobs: [], total: 0 }); + vi.spyOn(scraperClient, "getJobs").mockResolvedValue({ + jobs: [], + total: 0, + }); const service = new ScrapersService(); @@ -159,6 +188,13 @@ describe("ScrapersService", () => { ScraperAlreadyRunningError, ); + vi.spyOn(scraperClient, "triggerScrape").mockRejectedValueOnce( + new ScraperRunLockUnavailableError("down"), + ); + await expect(service.triggerScrape()).rejects.toBeInstanceOf( + ScraperRunLockUnavailableError, + ); + vi.spyOn(scraperClient, "triggerScrape").mockRejectedValueOnce( new Error("network"), ); @@ -186,7 +222,9 @@ describe("ScrapersService", () => { it("returns null jobsCollected when count fails", async () => { vi.spyOn(scraperClient, "getStatus").mockResolvedValue({ running: true }); - vi.spyOn(scraperClient, "getJobsCount").mockRejectedValue(new Error("down")); + vi.spyOn(scraperClient, "getJobsCount").mockRejectedValue( + new Error("down"), + ); await expect(new ScrapersService().listScrapers()).resolves.toEqual([ { @@ -223,14 +261,19 @@ describe("ScrapersController", () => { reprocessJobs: vi.fn(), }; const auditService = { fromRequest: vi.fn() }; - const controller = new ScrapersController(service as any, auditService as any); + const controller = new ScrapersController( + service as any, + auditService as any, + ); beforeEach(() => { vi.clearAllMocks(); service.triggerScrape.mockResolvedValue({ ok: true, message: "started" }); service.triggerScraper.mockResolvedValue({ ok: true, message: "started" }); service.getStatus.mockResolvedValue({ running: false }); - service.listScrapers.mockResolvedValue([{ name: "go-scraper", running: false }]); + service.listScrapers.mockResolvedValue([ + { name: "go-scraper", running: false }, + ]); service.getJobs.mockResolvedValue({ jobs: [], total: 0 }); service.getJobsCount.mockResolvedValue({ total: 0 }); service.reprocessJobs.mockResolvedValue({ ok: true, message: "queued" }); @@ -241,7 +284,10 @@ describe("ScrapersController", () => { await controller.trigger(req, res); expect(res.status).toHaveBeenCalledWith(202); - expect(auditService.fromRequest).toHaveBeenCalledWith(req, "scrapers.trigger"); + expect(auditService.fromRequest).toHaveBeenCalledWith( + req, + "scrapers.trigger", + ); service.triggerScrape.mockRejectedValueOnce( new ScraperAlreadyRunningError("busy"), @@ -249,6 +295,25 @@ describe("ScrapersController", () => { const conflictRes = response(); await controller.trigger(req, conflictRes); expect(conflictRes.status).toHaveBeenCalledWith(409); + expect(conflictRes.json).toHaveBeenCalledWith({ + ok: false, + code: "SCRAPER_ALREADY_RUNNING", + message: "busy", + }); + expect(auditService.fromRequest).toHaveBeenCalledTimes(1); + + service.triggerScrape.mockRejectedValueOnce( + new ScraperRunLockUnavailableError("lock unavailable"), + ); + const unavailableRes = response(); + await controller.trigger(req, unavailableRes); + expect(unavailableRes.status).toHaveBeenCalledWith(503); + expect(unavailableRes.json).toHaveBeenCalledWith({ + ok: false, + code: "SCRAPER_RUN_LOCK_UNAVAILABLE", + message: "lock unavailable", + }); + expect(auditService.fromRequest).toHaveBeenCalledTimes(1); }); it("triggers one scraper with 202 and maps failures", async () => { @@ -275,6 +340,26 @@ describe("ScrapersController", () => { conflictRes, ); expect(conflictRes.status).toHaveBeenCalledWith(409); + expect(conflictRes.json).toHaveBeenCalledWith({ + ok: false, + code: "SCRAPER_ALREADY_RUNNING", + message: "busy", + }); + + service.triggerScraper.mockRejectedValueOnce( + new ScraperRunLockUnavailableError("lock unavailable"), + ); + const unavailableRes = response(); + await controller.triggerOne( + { ...req, params: { id: "go-scraper" } } as unknown as Request, + unavailableRes, + ); + expect(unavailableRes.status).toHaveBeenCalledWith(503); + expect(unavailableRes.json).toHaveBeenCalledWith({ + ok: false, + code: "SCRAPER_RUN_LOCK_UNAVAILABLE", + message: "lock unavailable", + }); service.triggerScraper.mockRejectedValueOnce( new Error("scraper desconhecido: lever"), @@ -293,22 +378,38 @@ describe("ScrapersController", () => { await controller.listJobs(req, response()); await controller.jobsCount(req, response()); - expect(auditService.fromRequest).toHaveBeenCalledWith(req, "scrapers.read", { - type: "scrapers", - id: "status", - }); - expect(auditService.fromRequest).toHaveBeenCalledWith(req, "scrapers.read", { - type: "scrapers", - id: "list", - }); - expect(auditService.fromRequest).toHaveBeenCalledWith(req, "scrapers.read", { - type: "scrapers", - id: "jobs", - }); - expect(auditService.fromRequest).toHaveBeenCalledWith(req, "scrapers.read", { - type: "scrapers", - id: "jobs-count", - }); + expect(auditService.fromRequest).toHaveBeenCalledWith( + req, + "scrapers.read", + { + type: "scrapers", + id: "status", + }, + ); + expect(auditService.fromRequest).toHaveBeenCalledWith( + req, + "scrapers.read", + { + type: "scrapers", + id: "list", + }, + ); + expect(auditService.fromRequest).toHaveBeenCalledWith( + req, + "scrapers.read", + { + type: "scrapers", + id: "jobs", + }, + ); + expect(auditService.fromRequest).toHaveBeenCalledWith( + req, + "scrapers.read", + { + type: "scrapers", + id: "jobs-count", + }, + ); }); it("reprocesses jobs and maps failures to 500", async () => { diff --git a/docker-compose.yml b/docker-compose.yml index f658b5e..5538859 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -14,6 +14,8 @@ services: - GO_SCRAPER_ADDR=:8081 - VALKEY_URL=redis://valkey:6379/0 - SCRAPER_MAX_CONCURRENCY=${SCRAPER_MAX_CONCURRENCY-12} + - SCRAPER_RUN_LOCK_TTL=${SCRAPER_RUN_LOCK_TTL:-120s} + - SCRAPER_RUN_LOCK_RENEW_INTERVAL=${SCRAPER_RUN_LOCK_RENEW_INTERVAL:-30s} - GOMAXPROCS=${GOMAXPROCS:-2} - GOMEMLIMIT=${GOMEMLIMIT:-1500MiB} - GUPY_ENABLED=${GUPY_ENABLED:-true} @@ -31,8 +33,8 @@ services: - LEVER_ENABLED=${LEVER_ENABLED:-false} - LEVER_COMPANIES_FILE=/app/internal/interfaces/leverCompanies.json - LEVER_INCLUDE_ALL_JOBS=${LEVER_INCLUDE_ALL_JOBS:-true} - ports: - - "8081:8081" + expose: + - "8081" healthcheck: test: ["CMD", "wget", "--spider", "http://localhost:8081/health"] interval: 5s diff --git a/scraper-go/cmd/server/admin_handlers.go b/scraper-go/cmd/server/admin_handlers.go index be67d0b..ff6b9c7 100644 --- a/scraper-go/cmd/server/admin_handlers.go +++ b/scraper-go/cmd/server/admin_handlers.go @@ -3,30 +3,47 @@ package main import ( "context" "encoding/json" + "errors" "net/http" "strconv" "time" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/cronjob" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/jobstore" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" ) // handleTriggerScrape dispara o scraper manualmente via POST /admin/scrape // retorna 409 se já houver uma execução em andamento. -func handleTriggerScrape(scheduler *cronjob.Scheduler) http.HandlerFunc { +func handleTriggerScrape(scheduler *cronjob.Scheduler, executionCtx context.Context) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { // A execução manual precisa sobreviver ao fim da request HTTP; se usarmos // r.Context(), o scraper nasce com contexto cancelado assim que respondemos. - if err := scheduler.RunNow(context.Background()); err != nil { - if err == cronjob.ErrAlreadyRunning { - w.WriteHeader(http.StatusConflict) - json.NewEncoder(w).Encode(map[string]any{ - "ok": false, - "message": "scraper já está em execução", - }) + if err := scheduler.RunNow(executionCtx); err != nil { + if errors.Is(err, cronjob.ErrAlreadyRunning) { + writeScraperError( + w, + http.StatusConflict, + "SCRAPER_ALREADY_RUNNING", + "Já existe uma execução do scraper em andamento.", + ) return } - http.Error(w, "erro ao iniciar scraper", http.StatusInternalServerError) + if errors.Is(err, cronjob.ErrLockUnavailable) { + writeScraperError( + w, + http.StatusServiceUnavailable, + "SCRAPER_RUN_LOCK_UNAVAILABLE", + "Não foi possível confirmar a disponibilidade do scraper.", + ) + return + } + writeScraperError( + w, + http.StatusInternalServerError, + "SCRAPER_START_FAILED", + "Erro ao iniciar scraper.", + ) return } @@ -38,6 +55,47 @@ func handleTriggerScrape(scheduler *cronjob.Scheduler) http.HandlerFunc { } } +func writeScraperError(w http.ResponseWriter, status int, code, message string) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(map[string]any{ + "ok": false, + "code": code, + "message": message, + }) +} + +func mapRunLockError(w http.ResponseWriter, err error) bool { + switch { + case errors.Is(err, runlock.ErrAlreadyHeld): + writeScraperError( + w, + http.StatusConflict, + "SCRAPER_ALREADY_RUNNING", + "Já existe uma execução do scraper em andamento.", + ) + return true + case errors.Is(err, runlock.ErrUnavailable): + writeScraperError( + w, + http.StatusServiceUnavailable, + "SCRAPER_RUN_LOCK_UNAVAILABLE", + "Não foi possível confirmar a disponibilidade do scraper.", + ) + return true + case errors.Is(err, runlock.ErrLost): + writeScraperError( + w, + http.StatusServiceUnavailable, + "SCRAPER_RUN_LOCK_LOST", + "A execução perdeu o lock distribuído e foi cancelada.", + ) + return true + default: + return false + } +} + // handleScraperStatus retorna o estado atual do scheduler via GET /admin/scrape/status func handleScraperStatus(scheduler *cronjob.Scheduler) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { diff --git a/scraper-go/cmd/server/admin_handlers_test.go b/scraper-go/cmd/server/admin_handlers_test.go new file mode 100644 index 0000000..88aa46a --- /dev/null +++ b/scraper-go/cmd/server/admin_handlers_test.go @@ -0,0 +1,144 @@ +package main + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/cache" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/cronjob" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/domain" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/jobstore" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/keywords" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/ports" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" + "github.com/alicebob/miniredis/v2" + "github.com/redis/go-redis/v9" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type immediateAdapter struct{} + +func (immediateAdapter) SourceName() string { + return "immediate" +} + +func (immediateAdapter) Search(context.Context, string, domain.ScrapeRequest) ([]domain.Job, error) { + return nil, nil +} + +func TestManualScrapeIsAcceptedWhenLockIsFree(t *testing.T) { + scheduler, _, client, _ := newHandlerTestScheduler(t) + recorder := httptest.NewRecorder() + + handleTriggerScrape(scheduler, context.Background()).ServeHTTP( + recorder, + httptest.NewRequest(http.MethodPost, "/admin/scrape", nil), + ) + + assert.Equal(t, http.StatusOK, recorder.Code) + var body map[string]any + require.NoError(t, json.Unmarshal(recorder.Body.Bytes(), &body)) + assert.Equal(t, true, body["ok"]) + require.Eventually(t, func() bool { + return client.Exists(context.Background(), runlock.LockKey).Val() == 0 + }, time.Second, 10*time.Millisecond) +} + +func TestManualScrapeReturnsConflictWithoutStartingWhenLockIsHeld(t *testing.T) { + scheduler, manager, _, _ := newHandlerTestScheduler(t) + held, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + t.Cleanup(func() { _ = held.Release(context.Background()) }) + recorder := httptest.NewRecorder() + + handleTriggerScrape(scheduler, context.Background()).ServeHTTP( + recorder, + httptest.NewRequest(http.MethodPost, "/admin/scrape", nil), + ) + + assert.Equal(t, http.StatusConflict, recorder.Code) + assertScraperError(t, recorder, "SCRAPER_ALREADY_RUNNING") + assert.False(t, scheduler.IsRunning()) +} + +func TestManualScrapeFailsClosedWhenValkeyIsUnavailable(t *testing.T) { + scheduler, _, _, server := newHandlerTestScheduler(t) + server.Close() + recorder := httptest.NewRecorder() + + handleTriggerScrape(scheduler, context.Background()).ServeHTTP( + recorder, + httptest.NewRequest(http.MethodPost, "/admin/scrape", nil), + ) + + assert.Equal(t, http.StatusServiceUnavailable, recorder.Code) + assertScraperError(t, recorder, "SCRAPER_RUN_LOCK_UNAVAILABLE") + assert.False(t, scheduler.IsRunning()) +} + +func TestRunLockErrorsUseStableHTTPContract(t *testing.T) { + cases := []struct { + err error + status int + code string + }{ + {runlock.ErrAlreadyHeld, http.StatusConflict, "SCRAPER_ALREADY_RUNNING"}, + {runlock.ErrUnavailable, http.StatusServiceUnavailable, "SCRAPER_RUN_LOCK_UNAVAILABLE"}, + {runlock.ErrLost, http.StatusServiceUnavailable, "SCRAPER_RUN_LOCK_LOST"}, + } + + for _, tc := range cases { + t.Run(tc.code, func(t *testing.T) { + recorder := httptest.NewRecorder() + require.True(t, mapRunLockError(recorder, tc.err)) + assert.Equal(t, tc.status, recorder.Code) + assertScraperError(t, recorder, tc.code) + }) + } +} + +func newHandlerTestScheduler( + t *testing.T, +) (*cronjob.Scheduler, *runlock.Manager, *redis.Client, *miniredis.Miniredis) { + t.Helper() + + server := miniredis.RunT(t) + client := redis.NewClient(&redis.Options{Addr: server.Addr()}) + t.Cleanup(func() { _ = client.Close() }) + + manager, err := runlock.New(runlock.NewValkeyStore(client), runlock.Config{ + TTL: 2 * time.Second, + RenewInterval: 500 * time.Millisecond, + }) + require.NoError(t, err) + + keywordStore := keywords.NewStore(cache.NewRedisCache(client)) + require.NoError(t, keywordStore.Save(context.Background(), []string{"go"})) + + cfg := cronjob.DefaultConfig() + cfg.MaxConcurrency = 1 + cfg.ScrapeTimeout = time.Second + scheduler := cronjob.New( + cfg, + keywordStore, + jobstore.New(client), + []ports.JobSource{immediateAdapter{}}, + client, + manager, + ) + return scheduler, manager, client, server +} + +func assertScraperError(t *testing.T, recorder *httptest.ResponseRecorder, code string) { + t.Helper() + var body map[string]any + require.NoError(t, json.Unmarshal(recorder.Body.Bytes(), &body)) + assert.Equal(t, false, body["ok"]) + assert.Equal(t, code, body["code"]) + assert.NotEmpty(t, body["message"]) +} diff --git a/scraper-go/cmd/server/handlers.go b/scraper-go/cmd/server/handlers.go index 32f86f7..78b0d4e 100644 --- a/scraper-go/cmd/server/handlers.go +++ b/scraper-go/cmd/server/handlers.go @@ -14,6 +14,7 @@ import ( "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/keywords" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/pipeline" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/ports" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" ) const ( @@ -21,7 +22,14 @@ const ( scrapeTimeout = 15 * time.Minute ) -func handleScrape(adapterList []ports.JobSource, kwStore *keywords.Store, c cache.Cache, rdb *redis.Client, runtimeCfg config.RuntimeConfig) http.HandlerFunc { +func handleScrape( + adapterList []ports.JobSource, + kwStore *keywords.Store, + c cache.Cache, + rdb *redis.Client, + runLock *runlock.Manager, + runtimeCfg config.RuntimeConfig, +) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req domain.ScrapeRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { @@ -46,8 +54,20 @@ func handleScrape(adapterList []ports.JobSource, kwStore *keywords.Store, c cach start := time.Now() - result, err := pipeline.SearchJobs(ctx, c, searchConfig, adapterList, scrapeTTL, rdb) + result, err := pipeline.SearchJobs( + ctx, + c, + searchConfig, + adapterList, + scrapeTTL, + rdb, + runLock, + "public_endpoint", + ) if err != nil { + if mapRunLockError(w, err) { + return + } http.Error(w, "Erro ao buscar vagas.", http.StatusInternalServerError) return } diff --git a/scraper-go/cmd/server/main.go b/scraper-go/cmd/server/main.go index eeefc5c..c3f1256 100644 --- a/scraper-go/cmd/server/main.go +++ b/scraper-go/cmd/server/main.go @@ -49,6 +49,8 @@ func logRuntimeConfig(cfg config.RuntimeConfig) { slog.Info("scraper runtime configurado", "max_concurrency", cfg.MaxConcurrency, "max_concurrency_source", cfg.MaxConcurrencySource, + "run_lock_ttl", cfg.RunLockTTL, + "run_lock_renew_interval", cfg.RunLockRenewInterval, "gomaxprocs_effective", runtime.GOMAXPROCS(0), "gomaxprocs_source", envSource(gomaxprocsSet), "gomemlimit_effective_bytes", memLimit, diff --git a/scraper-go/cmd/server/server.go b/scraper-go/cmd/server/server.go index 9fe15a8..959f65c 100644 --- a/scraper-go/cmd/server/server.go +++ b/scraper-go/cmd/server/server.go @@ -16,6 +16,7 @@ import ( "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/jobstore" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/keywords" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/ports" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" "github.com/prometheus/client_golang/prometheus/promhttp" "github.com/redis/go-redis/v9" ) @@ -42,28 +43,39 @@ func run(adapterList []ports.JobSource, runtimeCfg config.RuntimeConfig) { // ── Módulos ── kwStore := keywords.NewStore(c) jobStore := jobstore.New(rdb) + runLock, err := runlock.New(runlock.NewValkeyStore(rdb), runlock.Config{ + TTL: runtimeCfg.RunLockTTL, + RenewInterval: runtimeCfg.RunLockRenewInterval, + }) + if err != nil { + slog.Error("configuração inválida do lock distribuído", "error", err) + os.Exit(1) + } // ── Scheduler (cronjob) ── schedulerCfg := cronjob.DefaultConfig() schedulerCfg.MaxConcurrency = runtimeCfg.MaxConcurrency - scheduler := cronjob.New(schedulerCfg, kwStore, jobStore, adapterList, rdb) + scheduler := cronjob.New(schedulerCfg, kwStore, jobStore, adapterList, rdb, runLock) scheduler.OnComplete = func(kws []string, scraped, saved int, duration time.Duration) { printSummary(len(adapterList), kws, scraped, duration) } + bgCtx, bgCancel := context.WithCancel(context.Background()) + defer bgCancel() + // ── Rotas ── mux := http.NewServeMux() // Públicas - mux.Handle("POST /scrape", handleScrape(adapterList, kwStore, c, rdb, runtimeCfg)) + mux.Handle("POST /scrape", handleScrape(adapterList, kwStore, c, rdb, runLock, runtimeCfg)) mux.Handle("GET /health", handleHealth(c)) mux.Handle("GET /metrics", promhttp.Handler()) mux.Handle("GET /api/keywords", handleGetKeywords(kwStore)) mux.Handle("POST /api/keywords", handleSaveKeywords(kwStore)) // Administrativas - mux.Handle("POST /admin/scrape", handleTriggerScrape(scheduler)) + mux.Handle("POST /admin/scrape", handleTriggerScrape(scheduler, bgCtx)) mux.Handle("GET /admin/scrape/status", handleScraperStatus(scheduler)) mux.Handle("GET /admin/jobs", handleGetJobs(jobStore)) mux.Handle("GET /admin/jobs/count", handleJobsCount(jobStore)) @@ -77,8 +89,6 @@ func run(adapterList []ports.JobSource, runtimeCfg config.RuntimeConfig) { } // ── Background: inicia o scheduler ── - bgCtx, bgCancel := context.WithCancel(context.Background()) - defer bgCancel() scheduler.Start(bgCtx) // ── Graceful shutdown ── @@ -99,7 +109,10 @@ func run(adapterList []ports.JobSource, runtimeCfg config.RuntimeConfig) { shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() - scheduler.Stop() + bgCancel() + if err := scheduler.Shutdown(shutdownCtx); err != nil { + slog.Error("erro ao aguardar liberação do lock no shutdown", "error", err) + } if err := srv.Shutdown(shutdownCtx); err != nil { slog.Error("erro durante o shutdown", "error", err) diff --git a/scraper-go/internal/config/config.go b/scraper-go/internal/config/config.go index f9d2b26..767826b 100644 --- a/scraper-go/internal/config/config.go +++ b/scraper-go/internal/config/config.go @@ -5,20 +5,27 @@ import ( "os" "strconv" "strings" + "time" ) const ( - DefaultMaxConcurrency = 12 + DefaultMaxConcurrency = 12 + DefaultRunLockTTL = 120 * time.Second + DefaultRunLockRenewInterval = 30 * time.Second SourceEnvironment = "environment" SourceInternalDefault = "internal_default" - ScraperMaxConcurrencyEnv = "SCRAPER_MAX_CONCURRENCY" + ScraperMaxConcurrencyEnv = "SCRAPER_MAX_CONCURRENCY" + ScraperRunLockTTLEnv = "SCRAPER_RUN_LOCK_TTL" + ScraperRunLockRenewIntervalEnv = "SCRAPER_RUN_LOCK_RENEW_INTERVAL" ) type RuntimeConfig struct { MaxConcurrency int MaxConcurrencySource string + RunLockTTL time.Duration + RunLockRenewInterval time.Duration } func LoadRuntimeConfig() (RuntimeConfig, error) { @@ -26,31 +33,70 @@ func LoadRuntimeConfig() (RuntimeConfig, error) { } func LoadRuntimeConfigFromLookup(lookup func(string) (string, bool)) (RuntimeConfig, error) { + cfg := RuntimeConfig{ + MaxConcurrency: DefaultMaxConcurrency, + MaxConcurrencySource: SourceInternalDefault, + RunLockTTL: DefaultRunLockTTL, + RunLockRenewInterval: DefaultRunLockRenewInterval, + } + value, ok := lookup(ScraperMaxConcurrencyEnv) - if !ok { - return RuntimeConfig{ - MaxConcurrency: DefaultMaxConcurrency, - MaxConcurrencySource: SourceInternalDefault, - }, nil + if ok { + trimmed := strings.TrimSpace(value) + if trimmed == "" { + return RuntimeConfig{}, fmt.Errorf("%s must be a positive integer", ScraperMaxConcurrencyEnv) + } + + parsed, err := strconv.Atoi(trimmed) + if err != nil { + return RuntimeConfig{}, fmt.Errorf("%s must be a positive integer: %w", ScraperMaxConcurrencyEnv, err) + } + if parsed <= 0 { + return RuntimeConfig{}, fmt.Errorf("%s must be greater than zero", ScraperMaxConcurrencyEnv) + } + + cfg.MaxConcurrency = parsed + cfg.MaxConcurrencySource = SourceEnvironment } - trimmed := strings.TrimSpace(value) - if trimmed == "" { - return RuntimeConfig{}, fmt.Errorf("%s must be a positive integer", ScraperMaxConcurrencyEnv) + var err error + cfg.RunLockTTL, err = durationFromLookup(lookup, ScraperRunLockTTLEnv, DefaultRunLockTTL) + if err != nil { + return RuntimeConfig{}, err + } + cfg.RunLockRenewInterval, err = durationFromLookup(lookup, ScraperRunLockRenewIntervalEnv, DefaultRunLockRenewInterval) + if err != nil { + return RuntimeConfig{}, err } + if cfg.RunLockRenewInterval >= cfg.RunLockTTL { + return RuntimeConfig{}, fmt.Errorf( + "%s must be shorter than %s", + ScraperRunLockRenewIntervalEnv, + ScraperRunLockTTLEnv, + ) + } + + return cfg, nil +} - parsed, err := strconv.Atoi(trimmed) +func durationFromLookup( + lookup func(string) (string, bool), + key string, + fallback time.Duration, +) (time.Duration, error) { + value, ok := lookup(key) + if !ok { + return fallback, nil + } + + parsed, err := time.ParseDuration(strings.TrimSpace(value)) if err != nil { - return RuntimeConfig{}, fmt.Errorf("%s must be a positive integer: %w", ScraperMaxConcurrencyEnv, err) + return 0, fmt.Errorf("%s must be a valid positive duration: %w", key, err) } if parsed <= 0 { - return RuntimeConfig{}, fmt.Errorf("%s must be greater than zero", ScraperMaxConcurrencyEnv) + return 0, fmt.Errorf("%s must be greater than zero", key) } - - return RuntimeConfig{ - MaxConcurrency: parsed, - MaxConcurrencySource: SourceEnvironment, - }, nil + return parsed, nil } func ResolveEffectiveConcurrency(requested, globalMax int) int { diff --git a/scraper-go/internal/config/config_test.go b/scraper-go/internal/config/config_test.go index 86bf76e..52b4d6d 100644 --- a/scraper-go/internal/config/config_test.go +++ b/scraper-go/internal/config/config_test.go @@ -2,6 +2,7 @@ package config import ( "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -15,11 +16,16 @@ func TestLoadRuntimeConfigUsesDefaultWhenEnvMissing(t *testing.T) { require.NoError(t, err) assert.Equal(t, 12, cfg.MaxConcurrency) assert.Equal(t, SourceInternalDefault, cfg.MaxConcurrencySource) + assert.Equal(t, 120*time.Second, cfg.RunLockTTL) + assert.Equal(t, 30*time.Second, cfg.RunLockRenewInterval) } func TestLoadRuntimeConfigUsesEnvValue(t *testing.T) { - cfg, err := LoadRuntimeConfigFromLookup(func(string) (string, bool) { - return "8", true + cfg, err := LoadRuntimeConfigFromLookup(func(key string) (string, bool) { + if key == ScraperMaxConcurrencyEnv { + return "8", true + } + return "", false }) require.NoError(t, err) @@ -56,3 +62,64 @@ func TestResolveEffectiveConcurrency(t *testing.T) { assert.Equal(t, 8, ResolveEffectiveConcurrency(8, 12)) assert.Equal(t, 12, ResolveEffectiveConcurrency(40, 12)) } + +func TestLoadRuntimeConfigUsesRunLockDurationsFromEnvironment(t *testing.T) { + values := map[string]string{ + ScraperRunLockTTLEnv: "3m", + ScraperRunLockRenewIntervalEnv: "45s", + } + + cfg, err := LoadRuntimeConfigFromLookup(func(key string) (string, bool) { + value, ok := values[key] + return value, ok + }) + + require.NoError(t, err) + assert.Equal(t, 3*time.Minute, cfg.RunLockTTL) + assert.Equal(t, 45*time.Second, cfg.RunLockRenewInterval) +} + +func TestLoadRuntimeConfigRejectsInvalidRunLockDurations(t *testing.T) { + cases := []struct { + name string + values map[string]string + }{ + { + name: "invalid ttl", + values: map[string]string{ScraperRunLockTTLEnv: "invalid"}, + }, + { + name: "zero ttl", + values: map[string]string{ScraperRunLockTTLEnv: "0s"}, + }, + { + name: "negative renewal interval", + values: map[string]string{ScraperRunLockRenewIntervalEnv: "-1s"}, + }, + { + name: "renewal equals ttl", + values: map[string]string{ + ScraperRunLockTTLEnv: "30s", + ScraperRunLockRenewIntervalEnv: "30s", + }, + }, + { + name: "renewal exceeds ttl", + values: map[string]string{ + ScraperRunLockTTLEnv: "30s", + ScraperRunLockRenewIntervalEnv: "31s", + }, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + _, err := LoadRuntimeConfigFromLookup(func(key string) (string, bool) { + value, ok := tc.values[key] + return value, ok + }) + + require.Error(t, err) + }) + } +} diff --git a/scraper-go/internal/cronjob/cronjob.go b/scraper-go/internal/cronjob/cronjob.go index dc04ef6..23ad327 100644 --- a/scraper-go/internal/cronjob/cronjob.go +++ b/scraper-go/internal/cronjob/cronjob.go @@ -2,6 +2,8 @@ package cronjob import ( "context" + "errors" + "fmt" "log/slog" "sync" "time" @@ -13,6 +15,7 @@ import ( "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/keywords" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/pipeline" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/ports" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" ) type Config struct { @@ -43,21 +46,33 @@ type Scheduler struct { jobStore *jobstore.Store adapterList []ports.JobSource rdb *redis.Client + runLock *runlock.Manager OnComplete func(keywords []string, scraped, saved int, duration time.Duration) mu sync.Mutex running bool + stopped bool lastRunAt time.Time lastJobs int stop chan struct{} + stopOnce sync.Once + active sync.WaitGroup } -func New(cfg Config, kwStore *keywords.Store, jobStore *jobstore.Store, adapterList []ports.JobSource, rdb *redis.Client) *Scheduler { +func New( + cfg Config, + kwStore *keywords.Store, + jobStore *jobstore.Store, + adapterList []ports.JobSource, + rdb *redis.Client, + runLock *runlock.Manager, +) *Scheduler { return &Scheduler{ cfg: cfg, kwStore: kwStore, jobStore: jobStore, adapterList: adapterList, rdb: rdb, + runLock: runLock, OnComplete: nil, stop: make(chan struct{}), } @@ -67,7 +82,7 @@ func (s *Scheduler) Start(ctx context.Context) { slog.Info("cronjob: scheduler iniciado", "interval", s.cfg.Interval) go func() { - s.run(ctx, "cron") + s.runCron(ctx) ticker := time.NewTicker(s.cfg.Interval) defer ticker.Stop() @@ -75,7 +90,7 @@ func (s *Scheduler) Start(ctx context.Context) { for { select { case <-ticker.C: - s.run(ctx, "cron") + s.runCron(ctx) case <-s.stop: slog.Info("cronjob: scheduler encerrado") return @@ -88,18 +103,50 @@ func (s *Scheduler) Start(ctx context.Context) { } func (s *Scheduler) Stop() { - close(s.stop) + s.stopOnce.Do(func() { + s.mu.Lock() + s.stopped = true + s.mu.Unlock() + close(s.stop) + }) +} + +// Shutdown impede novos disparos e aguarda as execuções ativas liberarem o lock. +func (s *Scheduler) Shutdown(ctx context.Context) error { + s.Stop() + + done := make(chan struct{}) + go func() { + s.active.Wait() + close(done) + }() + + select { + case <-done: + return nil + case <-ctx.Done(): + return fmt.Errorf("cronjob: shutdown: %w", ctx.Err()) + } } func (s *Scheduler) RunNow(ctx context.Context) error { - s.mu.Lock() - if s.running { - s.mu.Unlock() - return ErrAlreadyRunning + lease, err := s.acquire(ctx, "admin_manual") + if err != nil { + return err } - s.mu.Unlock() - go s.run(ctx, "admin_manual") + // Contabiliza a execução antes do yield da goroutine para o Shutdown não + // retornar enquanto o lease ainda estiver ativo. + s.active.Add(1) + go func() { + defer s.active.Done() + if runErr := s.runWithLease(lease); runErr != nil { + slog.Error("cronjob: execução manual falhou", + "source", "admin_manual", + "error", runErr, + ) + } + }() return nil } @@ -115,13 +162,47 @@ func (s *Scheduler) Snapshot() (running bool, lastRunAt time.Time, jobsCollected return s.running, s.lastRunAt, s.lastJobs } -func (s *Scheduler) run(ctx context.Context, origin string) { - s.mu.Lock() - if s.running { - s.mu.Unlock() - slog.Warn("cronjob: execução ignorada, já há uma em andamento") +func (s *Scheduler) runCron(ctx context.Context) { + lease, err := s.acquire(ctx, "cron") + if errors.Is(err, runlock.ErrAlreadyHeld) { + slog.Info("scraper execution skipped", + "source", "cron", + "reason", "run_lock_already_held", + "attempted_at", time.Now().UTC(), + ) return } + if err != nil { + slog.Error("cronjob: falha ao adquirir lock distribuído", + "source", "cron", + "attempted_at", time.Now().UTC(), + "error", err, + ) + return + } + + s.active.Add(1) + defer s.active.Done() + if err := s.runWithLease(lease); err != nil { + slog.Error("cronjob: execução falhou", "source", "cron", "error", err) + } +} + +func (s *Scheduler) acquire(ctx context.Context, source string) (*runlock.Lease, error) { + s.mu.Lock() + stopped := s.stopped + s.mu.Unlock() + if stopped { + return nil, fmt.Errorf("%w: scheduler is shutting down", runlock.ErrUnavailable) + } + if s.runLock == nil { + return nil, fmt.Errorf("%w: lock service is not configured", runlock.ErrUnavailable) + } + return s.runLock.Acquire(ctx, source) +} + +func (s *Scheduler) runWithLease(lease *runlock.Lease) (runErr error) { + s.mu.Lock() s.running = true s.mu.Unlock() @@ -131,19 +212,33 @@ func (s *Scheduler) run(ctx context.Context, origin string) { s.mu.Unlock() }() + defer func() { + releaseCtx, cancelRelease := context.WithTimeout(context.Background(), 5*time.Second) + defer cancelRelease() + if err := lease.Release(releaseCtx); err != nil && runErr == nil { + runErr = err + } + }() + start := time.Now() - scrapeCtx, cancel := context.WithTimeout(ctx, s.cfg.ScrapeTimeout) + scrapeCtx, cancel := context.WithTimeout(lease.Context(), s.cfg.ScrapeTimeout) defer cancel() kws, err := s.kwStore.Load(scrapeCtx) - if err != nil || len(kws) == 0 { + if err != nil { slog.Error("cronjob: falha ao carregar keywords", "error", err) - return + return fmt.Errorf("cronjob: load keywords: %w", err) + } + if len(kws) == 0 { + err = errors.New("no keywords configured") + slog.Error("cronjob: falha ao carregar keywords", "error", err) + return fmt.Errorf("cronjob: load keywords: %w", err) } config := s.searchConfig(kws) slog.Info("scraper execução iniciada", - "origin", origin, + "source", lease.State().Source, + "run_id", lease.RunID(), "max_concurrency_configured", s.cfg.MaxConcurrency, "max_concurrency_effective", config.MaxConcurrency, "keywords", len(kws), @@ -153,14 +248,20 @@ func (s *Scheduler) run(ctx context.Context, origin string) { jobs, err := pipeline.ScrapeAllSources(scrapeCtx, config, s.adapterList, s.rdb) if err != nil { slog.Error("cronjob: scrape falhou", "error", err) - return + return fmt.Errorf("cronjob: scrape: %w", err) + } + if err := executionError(scrapeCtx); err != nil { + return err } // Salva apenas vagas novas saved, err := s.jobStore.SaveBatch(scrapeCtx, jobs) if err != nil { slog.Error("cronjob: erro ao salvar vagas", "error", err) - return + return fmt.Errorf("cronjob: save jobs: %w", err) + } + if err := executionError(scrapeCtx); err != nil { + return err } // ✅ Constrói o índice invertido para buscas por keyword @@ -182,6 +283,7 @@ func (s *Scheduler) run(ctx context.Context, origin string) { if s.OnComplete != nil { s.OnComplete(kws, len(jobs), saved, time.Since(start)) } + return nil } func (s *Scheduler) searchConfig(kws []string) pipeline.SearchConfig { @@ -195,10 +297,15 @@ func (s *Scheduler) searchConfig(kws []string) pipeline.SearchConfig { } } -type alreadyRunningError struct{} - -func (e alreadyRunningError) Error() string { - return "cronjob: scraper já está em execução" +func executionError(ctx context.Context) error { + if cause := context.Cause(ctx); cause != nil { + return cause + } + return ctx.Err() } -var ErrAlreadyRunning error = alreadyRunningError{} +var ( + ErrAlreadyRunning = runlock.ErrAlreadyHeld + ErrLockUnavailable = runlock.ErrUnavailable + ErrLockLost = runlock.ErrLost +) diff --git a/scraper-go/internal/cronjob/cronjob_test.go b/scraper-go/internal/cronjob/cronjob_test.go index be0701d..16b03be 100644 --- a/scraper-go/internal/cronjob/cronjob_test.go +++ b/scraper-go/internal/cronjob/cronjob_test.go @@ -1,11 +1,79 @@ package cronjob import ( + "context" + "sync/atomic" "testing" + "time" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/cache" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/domain" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/jobstore" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/keywords" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/pipeline" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/ports" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" + "github.com/alicebob/miniredis/v2" + "github.com/redis/go-redis/v9" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) +type blockingAdapter struct { + started chan struct{} + release chan struct{} + calls atomic.Int32 +} + +type shutdownBlockingAdapter struct { + started chan struct{} + canceled chan struct{} + finish chan struct{} +} + +func (a *blockingAdapter) SourceName() string { + return "blocking" +} + +func (a *blockingAdapter) Search(ctx context.Context, _ string, _ domain.ScrapeRequest) ([]domain.Job, error) { + if a.calls.Add(1) == 1 { + close(a.started) + } + select { + case <-a.release: + return []domain.Job{schedulerTestJob()}, nil + case <-ctx.Done(): + return []domain.Job{schedulerTestJob()}, nil + } +} + +func (a *shutdownBlockingAdapter) SourceName() string { + return "shutdown-blocking" +} + +func (a *shutdownBlockingAdapter) Search( + ctx context.Context, + _ string, + _ domain.ScrapeRequest, +) ([]domain.Job, error) { + close(a.started) + <-ctx.Done() + close(a.canceled) + <-a.finish + return nil, context.Cause(ctx) +} + +func schedulerTestJob() domain.Job { + return domain.Job{ + Title: "Backend Engineer", + Company: "Candidate", + Location: "Remote", + URL: "https://example.com/backend-engineer", + Source: "test", + Description: "Backend engineer building Go APIs and distributed systems.", + } +} + func TestDefaultConfigUsesSafeMaxConcurrency(t *testing.T) { cfg := DefaultConfig() @@ -16,7 +84,7 @@ func TestSchedulerSearchConfigReceivesGlobalMaxConcurrency(t *testing.T) { cfg := DefaultConfig() cfg.MaxConcurrency = 9 - scheduler := New(cfg, nil, nil, nil, nil) + scheduler := New(cfg, nil, nil, nil, nil, nil) searchConfig := scheduler.searchConfig([]string{"go"}) assert.Equal(t, 9, searchConfig.MaxConcurrency) @@ -27,7 +95,202 @@ func TestAdminManualSharesSchedulerConfig(t *testing.T) { cfg := DefaultConfig() cfg.MaxConcurrency = 7 - scheduler := New(cfg, nil, nil, nil, nil) + scheduler := New(cfg, nil, nil, nil, nil, nil) assert.Equal(t, 7, scheduler.cfg.MaxConcurrency) } + +func TestTwoManualExecutionsStartOnlyOneAdapter(t *testing.T) { + scheduler, adapter, _ := newConcurrentTestScheduler(t) + + require.NoError(t, scheduler.RunNow(context.Background())) + waitForAdapterStart(t, adapter) + + err := scheduler.RunNow(context.Background()) + + require.ErrorIs(t, err, ErrAlreadyRunning) + assert.Equal(t, int32(1), adapter.calls.Load()) + close(adapter.release) + require.Eventually(t, func() bool { return !scheduler.IsRunning() }, time.Second, 10*time.Millisecond) +} + +func TestCronSkipsWhileManualExecutionHoldsLock(t *testing.T) { + scheduler, adapter, _ := newConcurrentTestScheduler(t) + + require.NoError(t, scheduler.RunNow(context.Background())) + waitForAdapterStart(t, adapter) + scheduler.runCron(context.Background()) + + assert.Equal(t, int32(1), adapter.calls.Load()) + close(adapter.release) + require.Eventually(t, func() bool { return !scheduler.IsRunning() }, time.Second, 10*time.Millisecond) +} + +func TestManualIsRejectedWhileCronExecutionHoldsLock(t *testing.T) { + scheduler, adapter, _ := newConcurrentTestScheduler(t) + + go scheduler.runCron(context.Background()) + waitForAdapterStart(t, adapter) + + err := scheduler.RunNow(context.Background()) + + require.ErrorIs(t, err, ErrAlreadyRunning) + assert.Equal(t, int32(1), adapter.calls.Load()) + close(adapter.release) + require.Eventually(t, func() bool { return !scheduler.IsRunning() }, time.Second, 10*time.Millisecond) +} + +func TestSecondCronExecutionDoesNotStartAdapters(t *testing.T) { + scheduler, adapter, _ := newConcurrentTestScheduler(t) + + go scheduler.runCron(context.Background()) + waitForAdapterStart(t, adapter) + scheduler.runCron(context.Background()) + + assert.Equal(t, int32(1), adapter.calls.Load()) + close(adapter.release) + require.Eventually(t, func() bool { return !scheduler.IsRunning() }, time.Second, 10*time.Millisecond) +} + +func TestManualExecutionReleasesLockWhenApplicationContextIsCanceled(t *testing.T) { + scheduler, adapter, client := newConcurrentTestScheduler(t) + ctx, cancel := context.WithCancel(context.Background()) + + require.NoError(t, scheduler.RunNow(ctx)) + waitForAdapterStart(t, adapter) + cancel() + + require.Eventually(t, func() bool { return !scheduler.IsRunning() }, time.Second, 10*time.Millisecond) + assert.Equal(t, int64(0), client.Exists(context.Background(), runlock.LockKey).Val()) +} + +func TestShutdownWaitsForActiveExecutionToReleaseLock(t *testing.T) { + scheduler, _, client := newConcurrentTestScheduler(t) + adapter := &shutdownBlockingAdapter{ + started: make(chan struct{}), + canceled: make(chan struct{}), + finish: make(chan struct{}), + } + scheduler.adapterList = []ports.JobSource{adapter} + runCtx, cancelRun := context.WithCancel(context.Background()) + + require.NoError(t, scheduler.RunNow(runCtx)) + select { + case <-adapter.started: + case <-time.After(time.Second): + t.Fatal("adapter did not start") + } + cancelRun() + select { + case <-adapter.canceled: + case <-time.After(time.Second): + t.Fatal("adapter context was not canceled") + } + + shutdownCtx, cancelShutdown := context.WithTimeout(context.Background(), time.Second) + defer cancelShutdown() + shutdownDone := make(chan error, 1) + go func() { + shutdownDone <- scheduler.Shutdown(shutdownCtx) + }() + + select { + case err := <-shutdownDone: + t.Fatalf("shutdown returned before active execution finished: %v", err) + case <-time.After(50 * time.Millisecond): + } + + close(adapter.finish) + require.NoError(t, <-shutdownDone) + assert.False(t, scheduler.IsRunning()) + assert.Equal(t, int64(0), client.Exists(context.Background(), runlock.LockKey).Val()) + require.ErrorIs(t, scheduler.RunNow(context.Background()), ErrLockUnavailable) +} + +func TestSchedulerAndPublicSearchShareLockWithoutStartingOrPersistingSecondRun(t *testing.T) { + scheduler, adapter, client := newConcurrentTestScheduler(t) + + require.NoError(t, scheduler.RunNow(context.Background())) + waitForAdapterStart(t, adapter) + + _, err := pipeline.SearchJobs( + context.Background(), + cache.NewMemoryCache(), + pipeline.SearchConfig{Keywords: []string{"go"}, MaxConcurrency: 1}, + scheduler.adapterList, + time.Minute, + client, + scheduler.runLock, + "public_endpoint", + ) + + require.ErrorIs(t, err, runlock.ErrAlreadyHeld) + assert.Equal(t, int32(1), adapter.calls.Load()) + assert.Equal(t, int64(0), client.Exists(context.Background(), "scraper:jobs:index").Val()) + close(adapter.release) + require.Eventually(t, func() bool { return !scheduler.IsRunning() }, time.Second, 10*time.Millisecond) +} + +func TestConfirmedLockLossPreventsSchedulerPersistenceAndIndexing(t *testing.T) { + scheduler, adapter, client := newConcurrentTestScheduler(t) + + require.NoError(t, scheduler.RunNow(context.Background())) + waitForAdapterStart(t, adapter) + require.NoError(t, client.Set( + context.Background(), + runlock.LockKey, + "new-owner", + time.Second, + ).Err()) + + require.Eventually(t, func() bool { return !scheduler.IsRunning() }, 2*time.Second, 10*time.Millisecond) + + job := schedulerTestJob() + jobID := jobstore.StableID(&job) + assert.Equal(t, int64(0), client.Exists(context.Background(), "scraper:job:"+jobID).Val()) + assert.Equal(t, int64(0), client.Exists(context.Background(), "scraper:jobs:index").Val()) + assert.Equal(t, "new-owner", client.Get(context.Background(), runlock.LockKey).Val()) +} + +func newConcurrentTestScheduler(t *testing.T) (*Scheduler, *blockingAdapter, *redis.Client) { + t.Helper() + + server := miniredis.RunT(t) + client := redis.NewClient(&redis.Options{Addr: server.Addr()}) + t.Cleanup(func() { _ = client.Close() }) + + lockManager, err := runlock.New(runlock.NewValkeyStore(client), runlock.Config{ + TTL: 2 * time.Second, + RenewInterval: 500 * time.Millisecond, + }) + require.NoError(t, err) + + keywordStore := keywords.NewStore(cache.NewRedisCache(client)) + require.NoError(t, keywordStore.Save(context.Background(), []string{"go"})) + + adapter := &blockingAdapter{ + started: make(chan struct{}), + release: make(chan struct{}), + } + cfg := DefaultConfig() + cfg.ScrapeTimeout = 2 * time.Second + cfg.MaxConcurrency = 1 + + return New( + cfg, + keywordStore, + jobstore.New(client), + []ports.JobSource{adapter}, + client, + lockManager, + ), adapter, client +} + +func waitForAdapterStart(t *testing.T, adapter *blockingAdapter) { + t.Helper() + select { + case <-adapter.started: + case <-time.After(time.Second): + t.Fatal("adapter did not start") + } +} diff --git a/scraper-go/internal/pipeline/pipeline.go b/scraper-go/internal/pipeline/pipeline.go index 690c64d..2ea0e2f 100644 --- a/scraper-go/internal/pipeline/pipeline.go +++ b/scraper-go/internal/pipeline/pipeline.go @@ -58,9 +58,22 @@ func Run(ctx context.Context, adapterList []ports.JobSource, req domain.ScrapeRe results := make(chan result, len(tasks)) var wg sync.WaitGroup +schedule: for _, t := range tasks { + if ctx.Err() != nil { + break schedule + } + select { + case sem <- struct{}{}: + case <-ctx.Done(): + break schedule + } + if ctx.Err() != nil { + <-sem + break schedule + } + wg.Add(1) - sem <- struct{}{} go func(t adapterTask) { defer wg.Done() @@ -109,6 +122,9 @@ func Run(ctx context.Context, adapterList []ports.JobSource, req domain.ScrapeRe allJobs = append(allJobs, r.jobs...) } } + if cause := context.Cause(ctx); cause != nil { + return nil, cause + } deduped := dedup.DedupeJobs(allJobs) classified := classifier.ClassifyJobs(deduped) @@ -134,7 +150,7 @@ func runAdapterTask(ctx context.Context, t adapterTask, req domain.ScrapeRequest } func IndexJobsInValkey(ctx context.Context, rdb *redis.Client, jobs []domain.Job, keywords []string) { - if rdb == nil || len(jobs) == 0 { + if rdb == nil || len(jobs) == 0 || ctx.Err() != nil { return } @@ -157,6 +173,9 @@ func IndexJobsInValkey(ctx context.Context, rdb *redis.Client, jobs []domain.Job kwIndex := make(map[string][]string) // finalKey → []id for _, job := range jobs { + if ctx.Err() != nil { + return + } id := jobstore.StableID(&job) if id == "" { continue @@ -203,6 +222,9 @@ func IndexJobsInValkey(ctx context.Context, rdb *redis.Client, jobs []domain.Job // Publica os índices de keyword com RENAME atômico // Fluxo: escreve em :next → RENAME :next → final → Expire no final for finalKey, ids := range kwIndex { + if ctx.Err() != nil { + return + } tempKey := finalKey + ":next" pipe := rdb.Pipeline() diff --git a/scraper-go/internal/pipeline/search.go b/scraper-go/internal/pipeline/search.go index 2304930..8b5eff8 100644 --- a/scraper-go/internal/pipeline/search.go +++ b/scraper-go/internal/pipeline/search.go @@ -12,6 +12,7 @@ import ( "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/domain" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/inflight" "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/ports" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" ) type SearchResult struct { @@ -28,6 +29,8 @@ func SearchJobs( adapterList []ports.JobSource, ttl time.Duration, rdb *redis.Client, + runLock *runlock.Manager, + source string, ) (SearchResult, error) { config = normalizeSearchConfig(config) cacheKey := BuildCacheKey(config) @@ -39,7 +42,7 @@ func SearchJobs( return result, nil } - return inflight.Do(cacheKey, func() (SearchResult, error) { + return inflight.Do(cacheKey, func() (result SearchResult, resultErr error) { if result, found, err := cache.GetAs[SearchResult](c, ctx, cacheKey); err != nil { return SearchResult{}, fmt.Errorf("pipeline.SearchJobs: cache re-check: %w", err) } else if found { @@ -47,21 +50,43 @@ func SearchJobs( return result, nil } - jobs, err := ScrapeAllSources(ctx, config, adapterList, rdb) + if runLock == nil { + return SearchResult{}, fmt.Errorf("%w: lock service is not configured", runlock.ErrUnavailable) + } + lease, err := runLock.Acquire(ctx, source) + if err != nil { + return SearchResult{}, err + } + defer func() { + releaseCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := lease.Release(releaseCtx); err != nil && resultErr == nil { + resultErr = err + } + }() + + runCtx := lease.Context() + jobs, err := ScrapeAllSources(runCtx, config, adapterList, rdb) if err != nil { return SearchResult{}, fmt.Errorf("pipeline.SearchJobs: scrape: %w", err) } + if err := context.Cause(runCtx); err != nil { + return SearchResult{}, err + } - IndexJobsInValkey(ctx, rdb, jobs, config.Keywords) + IndexJobsInValkey(runCtx, rdb, jobs, config.Keywords) + if err := context.Cause(runCtx); err != nil { + return SearchResult{}, err + } - result := SearchResult{ + result = SearchResult{ Jobs: jobs, Total: len(jobs), CachedAt: time.Now(), FromCache: false, } - if err := c.Set(ctx, cacheKey, result, ttl); err != nil { + if err := c.Set(runCtx, cacheKey, result, ttl); err != nil { slog.Error("pipeline.SearchJobs: cache write failed", "key", cacheKey, "error", err, diff --git a/scraper-go/internal/pipeline/search_lock_test.go b/scraper-go/internal/pipeline/search_lock_test.go new file mode 100644 index 0000000..7f18f77 --- /dev/null +++ b/scraper-go/internal/pipeline/search_lock_test.go @@ -0,0 +1,223 @@ +package pipeline + +import ( + "context" + "errors" + "sync/atomic" + "testing" + "time" + + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/cache" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/domain" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/ports" + "github.com/Benevanio/Jobs_Scraper_Global/scraper-go/internal/runlock" + "github.com/alicebob/miniredis/v2" + "github.com/redis/go-redis/v9" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type searchLockAdapter struct { + calls atomic.Int32 +} + +type blockingSearchLockAdapter struct { + started chan struct{} +} + +func (a *blockingSearchLockAdapter) SourceName() string { + return "blocking-search-lock" +} + +func (a *blockingSearchLockAdapter) Search(ctx context.Context, _ string, _ domain.ScrapeRequest) ([]domain.Job, error) { + select { + case <-a.started: + default: + close(a.started) + } + <-ctx.Done() + return nil, context.Cause(ctx) +} + +func (a *searchLockAdapter) SourceName() string { + return "search-lock" +} + +func (a *searchLockAdapter) Search(context.Context, string, domain.ScrapeRequest) ([]domain.Job, error) { + panic("batch path expected") +} + +func (a *searchLockAdapter) SearchBatch(context.Context, []string, domain.ScrapeRequest) ([]domain.Job, error) { + a.calls.Add(1) + return []domain.Job{{ + Title: "Go Developer", + Company: "Candidate", + Location: "Remote", + URL: "https://example.com/jobs/1", + Source: "test", + }}, nil +} + +func TestSearchJobsDoesNotStartAdapterWhenAnotherOriginOwnsLock(t *testing.T) { + manager, client, _ := newSearchLockManager(t) + held, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + t.Cleanup(func() { _ = held.Release(context.Background()) }) + adapter := &searchLockAdapter{} + + _, err = SearchJobs( + context.Background(), + cache.NewMemoryCache(), + SearchConfig{Keywords: []string{"go"}, MaxConcurrency: 1}, + []ports.JobSource{adapter}, + time.Minute, + client, + manager, + "public_endpoint", + ) + + require.ErrorIs(t, err, runlock.ErrAlreadyHeld) + assert.Equal(t, int32(0), adapter.calls.Load()) + assert.Equal(t, int64(0), client.Exists(context.Background(), "scraper:jobs:index").Val()) +} + +func TestSearchJobsFailsClosedWhenValkeyIsUnavailable(t *testing.T) { + manager, client, server := newSearchLockManager(t) + server.Close() + adapter := &searchLockAdapter{} + + _, err := SearchJobs( + context.Background(), + cache.NewMemoryCache(), + SearchConfig{Keywords: []string{"go"}, MaxConcurrency: 1}, + []ports.JobSource{adapter}, + time.Minute, + client, + manager, + "public_endpoint", + ) + + require.ErrorIs(t, err, runlock.ErrUnavailable) + assert.Equal(t, int32(0), adapter.calls.Load()) +} + +func TestSearchJobsRunsAndReleasesLockWhenItIsFree(t *testing.T) { + manager, client, _ := newSearchLockManager(t) + adapter := &searchLockAdapter{} + + result, err := SearchJobs( + context.Background(), + cache.NewMemoryCache(), + SearchConfig{Keywords: []string{"go"}, MaxConcurrency: 1}, + []ports.JobSource{adapter}, + time.Minute, + client, + manager, + "public_endpoint", + ) + + require.NoError(t, err) + assert.Equal(t, 0, result.Total) + assert.Equal(t, int32(1), adapter.calls.Load()) + assert.Equal(t, int64(0), client.Exists(context.Background(), runlock.LockKey).Val()) +} + +func TestSearchJobsCancelsWithoutIndexingAfterConfirmedLockLoss(t *testing.T) { + server := miniredis.RunT(t) + client := redis.NewClient(&redis.Options{Addr: server.Addr()}) + t.Cleanup(func() { _ = client.Close() }) + manager, err := runlock.New(runlock.NewValkeyStore(client), runlock.Config{ + TTL: 200 * time.Millisecond, + RenewInterval: 25 * time.Millisecond, + }) + require.NoError(t, err) + adapter := &blockingSearchLockAdapter{started: make(chan struct{})} + result := make(chan error, 1) + + go func() { + _, searchErr := SearchJobs( + context.Background(), + cache.NewMemoryCache(), + SearchConfig{Keywords: []string{"go"}, MaxConcurrency: 1}, + []ports.JobSource{adapter}, + time.Minute, + client, + manager, + "public_endpoint", + ) + result <- searchErr + }() + + select { + case <-adapter.started: + case <-time.After(time.Second): + t.Fatal("adapter did not start") + } + require.NoError(t, client.Set(context.Background(), runlock.LockKey, "new-owner", time.Second).Err()) + + select { + case err := <-result: + require.ErrorIs(t, err, runlock.ErrLost) + case <-time.After(time.Second): + t.Fatal("search did not stop after lock loss") + } + assert.Equal(t, int64(0), client.Exists(context.Background(), "scraper:jobs:index").Val()) + assert.Equal(t, "new-owner", client.Get(context.Background(), runlock.LockKey).Val()) +} + +func TestPipelineErrorIsNotMaskedByReleaseFailure(t *testing.T) { + manager, err := runlock.New(releaseFailingStore{}, runlock.Config{ + TTL: time.Second, + RenewInterval: 250 * time.Millisecond, + }) + require.NoError(t, err) + + _, err = SearchJobs( + context.Background(), + cache.NewMemoryCache(), + SearchConfig{Keywords: []string{"go"}, MaxConcurrency: 0}, + []ports.JobSource{&searchLockAdapter{}}, + time.Minute, + nil, + manager, + "public_endpoint", + ) + + require.ErrorContains(t, err, "max concurrency") + assert.NotContains(t, err.Error(), "release failed") +} + +func newSearchLockManager( + t *testing.T, +) (*runlock.Manager, *redis.Client, *miniredis.Miniredis) { + t.Helper() + + server := miniredis.RunT(t) + client := redis.NewClient(&redis.Options{Addr: server.Addr()}) + t.Cleanup(func() { _ = client.Close() }) + + manager, err := runlock.New(runlock.NewValkeyStore(client), runlock.Config{ + TTL: 2 * time.Second, + RenewInterval: 500 * time.Millisecond, + }) + require.NoError(t, err) + return manager, client, server +} + +type releaseFailingStore struct{} + +func (releaseFailingStore) TryAcquire(context.Context, string, time.Duration) (bool, error) { + return true, nil +} + +func (releaseFailingStore) WriteState(context.Context, string, runlock.State, time.Duration) (bool, error) { + return true, nil +} + +func (releaseFailingStore) Renew(context.Context, string, runlock.State, time.Duration) (bool, error) { + return true, nil +} + +func (releaseFailingStore) Release(context.Context, string, string) (bool, error) { + return false, errors.New("release failed") +} diff --git a/scraper-go/internal/pipeline/source_schedule_test.go b/scraper-go/internal/pipeline/source_schedule_test.go index 79623e7..b55c89c 100644 --- a/scraper-go/internal/pipeline/source_schedule_test.go +++ b/scraper-go/internal/pipeline/source_schedule_test.go @@ -2,6 +2,7 @@ package pipeline import ( "context" + "sync/atomic" "testing" "time" @@ -31,6 +32,23 @@ type batchRunTestAdapter struct { keywords []string } +type cancelAwareAdapter struct { + started chan struct{} + calls atomic.Int32 +} + +func (a *cancelAwareAdapter) SourceName() string { + return "Cancel Test" +} + +func (a *cancelAwareAdapter) Search(ctx context.Context, _ string, _ domain.ScrapeRequest) ([]domain.Job, error) { + if a.calls.Add(1) == 1 { + close(a.started) + } + <-ctx.Done() + return nil, context.Cause(ctx) +} + func (a *batchRunTestAdapter) SourceName() string { return "Batch Test" } @@ -121,3 +139,28 @@ func TestRunExecutesWithValidMaxConcurrency(t *testing.T) { assert.Empty(t, jobs) assert.Equal(t, 1, adapter.batchCalls) } + +func TestRunStopsSchedulingAdaptersAfterContextCancellation(t *testing.T) { + adapter := &cancelAwareAdapter{started: make(chan struct{})} + ctx, cancel := context.WithCancelCause(context.Background()) + result := make(chan error, 1) + + go func() { + _, err := Run(ctx, []adapters.Adapter{adapter}, domain.ScrapeRequest{ + Keywords: []string{"go", "java", "python"}, + MaxConcurrency: 1, + }) + result <- err + }() + + select { + case <-adapter.started: + case <-time.After(time.Second): + t.Fatal("first adapter did not start") + } + lostErr := assert.AnError + cancel(lostErr) + + require.ErrorIs(t, <-result, lostErr) + assert.Equal(t, int32(1), adapter.calls.Load()) +} diff --git a/scraper-go/internal/runlock/runlock.go b/scraper-go/internal/runlock/runlock.go new file mode 100644 index 0000000..0777884 --- /dev/null +++ b/scraper-go/internal/runlock/runlock.go @@ -0,0 +1,276 @@ +package runlock + +import ( + "context" + "crypto/rand" + "encoding/hex" + "errors" + "fmt" + "log/slog" + "sync" + "time" +) + +var ( + ErrAlreadyHeld = errors.New("scraper run lock already held") + ErrUnavailable = errors.New("scraper run lock unavailable") + ErrLost = errors.New("scraper run lock lost") +) + +type Config struct { + TTL time.Duration + RenewInterval time.Duration +} + +func (c Config) Validate() error { + if c.TTL <= 0 { + return errors.New("runlock: ttl must be greater than zero") + } + if c.RenewInterval <= 0 { + return errors.New("runlock: renew interval must be greater than zero") + } + if c.RenewInterval >= c.TTL { + return errors.New("runlock: renew interval must be shorter than ttl") + } + return nil +} + +type State struct { + RunID string + Source string + StartedAt time.Time + LockExpiresAt time.Time +} + +type Manager struct { + store Store + cfg Config + now func() time.Time +} + +func New(store Store, cfg Config) (*Manager, error) { + if store == nil { + return nil, errors.New("runlock: store is required") + } + if err := cfg.Validate(); err != nil { + return nil, err + } + return &Manager{ + store: store, + cfg: cfg, + now: time.Now, + }, nil +} + +func (m *Manager) Acquire(ctx context.Context, source string) (*Lease, error) { + token, err := newToken() + if err != nil { + return nil, fmt.Errorf("runlock: generate token: %w", err) + } + runID, err := newToken() + if err != nil { + return nil, fmt.Errorf("runlock: generate run id: %w", err) + } + + acquired, err := m.store.TryAcquire(ctx, token, m.cfg.TTL) + if err != nil { + return nil, fmt.Errorf("%w: acquire: %v", ErrUnavailable, err) + } + if !acquired { + return nil, ErrAlreadyHeld + } + + startedAt := m.now().UTC() + state := State{ + RunID: runID, + Source: source, + StartedAt: startedAt, + LockExpiresAt: startedAt.Add(m.cfg.TTL), + } + if written, stateErr := m.store.WriteState(ctx, token, state, m.cfg.TTL); stateErr != nil || !written { + slog.Warn("scraper run state unavailable", + "source", source, + "run_id", runID, + "error", stateErr, + ) + } + + leaseCtx, cancel := context.WithCancelCause(ctx) + lease := &Lease{ + manager: m, + token: token, + runID: runID, + source: source, + startedAt: startedAt, + ctx: leaseCtx, + cancel: cancel, + stopRenewal: make(chan struct{}), + renewalDone: make(chan struct{}), + } + + slog.Info("scraper run lock acquired", + "source", source, + "run_id", runID, + "lock_expires_at", state.LockExpiresAt, + ) + go lease.renew() + return lease, nil +} + +type Lease struct { + manager *Manager + token string + runID string + source string + startedAt time.Time + + ctx context.Context + cancel context.CancelCauseFunc + + stopRenewal chan struct{} + renewalDone chan struct{} + stopOnce sync.Once + releaseOnce sync.Once + releaseErr error +} + +func (l *Lease) Context() context.Context { + return l.ctx +} + +func (l *Lease) Token() string { + return l.token +} + +func (l *Lease) RunID() string { + return l.runID +} + +func (l *Lease) State() State { + return State{ + RunID: l.runID, + Source: l.source, + StartedAt: l.startedAt, + LockExpiresAt: l.manager.now().UTC().Add(l.manager.cfg.TTL), + } +} + +func (l *Lease) Release(ctx context.Context) error { + l.releaseOnce.Do(func() { + l.cancel(nil) + l.stopOnce.Do(func() { + close(l.stopRenewal) + }) + <-l.renewalDone + + released, err := l.manager.store.Release(ctx, l.token, l.runID) + switch { + case err != nil: + l.releaseErr = fmt.Errorf("%w: release: %v", ErrUnavailable, err) + slog.Error("scraper run lock release failed", + "source", l.source, + "run_id", l.runID, + "error", err, + ) + case released: + slog.Info("scraper run lock released", + "source", l.source, + "run_id", l.runID, + ) + default: + l.releaseErr = fmt.Errorf("%w: ownership changed during release", ErrLost) + slog.Warn("scraper run lock not released because ownership changed", + "source", l.source, + "run_id", l.runID, + ) + } + }) + return l.releaseErr +} + +func (l *Lease) renew() { + defer close(l.renewalDone) + + ticker := time.NewTicker(l.manager.cfg.RenewInterval) + defer ticker.Stop() + + safetyTimer := time.NewTimer(l.manager.cfg.TTL - l.manager.cfg.RenewInterval) + defer safetyTimer.Stop() + + for { + select { + case <-l.stopRenewal: + return + case <-l.ctx.Done(): + return + case <-safetyTimer.C: + cause := fmt.Errorf("%w: renewal could not be confirmed before safety deadline", ErrLost) + slog.Error("scraper run lock lost", + "source", l.source, + "run_id", l.runID, + "reason", "renewal_safety_deadline", + ) + l.cancel(cause) + return + case <-ticker.C: + state := l.State() + renewCtx, cancelRenew := context.WithTimeout(l.ctx, l.renewTimeout()) + renewed, err := l.manager.store.Renew( + renewCtx, + l.token, + state, + l.manager.cfg.TTL, + ) + cancelRenew() + if err != nil { + if l.ctx.Err() != nil { + return + } + slog.Warn("scraper run lock renewal failed temporarily", + "source", l.source, + "run_id", l.runID, + "error", err, + ) + continue + } + if !renewed { + cause := fmt.Errorf("%w: token no longer owns lock", ErrLost) + slog.Error("scraper run lock lost", + "source", l.source, + "run_id", l.runID, + "reason", "ownership_changed", + ) + l.cancel(cause) + return + } + + resetTimer(safetyTimer, l.manager.cfg.TTL-l.manager.cfg.RenewInterval) + } + } +} + +func (l *Lease) renewTimeout() time.Duration { + safetyMargin := l.manager.cfg.TTL - l.manager.cfg.RenewInterval + if safetyMargin < l.manager.cfg.RenewInterval { + return safetyMargin + } + return l.manager.cfg.RenewInterval +} + +func resetTimer(timer *time.Timer, duration time.Duration) { + if !timer.Stop() { + select { + case <-timer.C: + default: + } + } + timer.Reset(duration) +} + +func newToken() (string, error) { + var value [32]byte + if _, err := rand.Read(value[:]); err != nil { + return "", err + } + return hex.EncodeToString(value[:]), nil +} diff --git a/scraper-go/internal/runlock/runlock_test.go b/scraper-go/internal/runlock/runlock_test.go new file mode 100644 index 0000000..3da0e32 --- /dev/null +++ b/scraper-go/internal/runlock/runlock_test.go @@ -0,0 +1,346 @@ +package runlock + +import ( + "context" + "errors" + "fmt" + "sync/atomic" + "testing" + "time" + + "github.com/alicebob/miniredis/v2" + "github.com/redis/go-redis/v9" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestConfigValidation(t *testing.T) { + require.NoError(t, (Config{TTL: 2 * time.Second, RenewInterval: time.Second}).Validate()) + + invalid := []Config{ + {TTL: 0, RenewInterval: time.Second}, + {TTL: time.Second, RenewInterval: 0}, + {TTL: time.Second, RenewInterval: time.Second}, + {TTL: time.Second, RenewInterval: 2 * time.Second}, + } + for _, cfg := range invalid { + require.Error(t, cfg.Validate()) + } +} + +func TestAcquireCommandUsesAtomicSetNXPX(t *testing.T) { + assert.Equal(t, []any{ + "SET", + LockKey, + "owner-token", + "NX", + "PX", + int64(120000), + }, acquireCommand("owner-token", 120*time.Second)) +} + +func TestAcquireRejectsExistingLockAndCreatesOperationalState(t *testing.T) { + manager, client, _ := newTestManager(t, 2*time.Second, 500*time.Millisecond) + ctx := context.Background() + + lease, err := manager.Acquire(ctx, "cron") + require.NoError(t, err) + t.Cleanup(func() { _ = lease.Release(context.Background()) }) + + assert.Equal(t, lease.Token(), client.Get(ctx, LockKey).Val()) + assert.Greater(t, client.PTTL(ctx, LockKey).Val(), time.Duration(0)) + + state, err := client.HGetAll(ctx, StateKey).Result() + require.NoError(t, err) + assert.Equal(t, lease.RunID(), state["runId"]) + assert.NotEqual(t, lease.Token(), state["runId"]) + assert.Equal(t, "cron", state["source"]) + assert.NotEmpty(t, state["startedAt"]) + assert.NotEmpty(t, state["lockExpiresAt"]) + + _, err = manager.Acquire(ctx, "admin_manual") + require.ErrorIs(t, err, ErrAlreadyHeld) +} + +func TestTokensAreUniqueAcrossExecutions(t *testing.T) { + manager, _, _ := newTestManager(t, 2*time.Second, 500*time.Millisecond) + + first, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + firstToken := first.Token() + firstRunID := first.RunID() + require.NoError(t, first.Release(context.Background())) + + second, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + defer second.Release(context.Background()) + + assert.NotEqual(t, firstToken, second.Token()) + assert.NotEqual(t, firstRunID, second.RunID()) +} + +func TestOnlyOwnerCanRenewOrRelease(t *testing.T) { + manager, client, server := newTestManager(t, 2*time.Second, 500*time.Millisecond) + ctx := context.Background() + + lease, err := manager.Acquire(ctx, "cron") + require.NoError(t, err) + + wrongState := lease.State() + wrongState.RunID = "different-run" + renewed, err := manager.store.Renew(ctx, "different-token", wrongState, 2*time.Second) + require.NoError(t, err) + assert.False(t, renewed) + + released, err := manager.store.Release(ctx, "different-token", "different-run") + require.NoError(t, err) + assert.False(t, released) + assert.Equal(t, lease.Token(), client.Get(ctx, LockKey).Val()) + + server.FastForward(time.Second) + renewed, err = manager.store.Renew(ctx, lease.Token(), lease.State(), 2*time.Second) + require.NoError(t, err) + assert.True(t, renewed) + assert.Greater(t, client.PTTL(ctx, LockKey).Val(), 1500*time.Millisecond) + assert.Greater(t, client.PTTL(ctx, StateKey).Val(), 1500*time.Millisecond) + + require.NoError(t, lease.Release(ctx)) + assert.Equal(t, redis.Nil, client.Get(ctx, LockKey).Err()) + assert.Empty(t, client.HGetAll(ctx, StateKey).Val()) +} + +func TestOldExecutionCannotRemoveNewLock(t *testing.T) { + manager, client, _ := newTestManager(t, 2*time.Second, 500*time.Millisecond) + ctx := context.Background() + + oldLease, err := manager.Acquire(ctx, "cron") + require.NoError(t, err) + + require.NoError(t, client.Set(ctx, LockKey, "new-owner", 2*time.Second).Err()) + require.NoError(t, client.HSet(ctx, StateKey, "runId", "new-run").Err()) + + require.ErrorIs(t, oldLease.Release(ctx), ErrLost) + assert.Equal(t, "new-owner", client.Get(ctx, LockKey).Val()) + assert.Equal(t, "new-run", client.HGet(ctx, StateKey, "runId").Val()) +} + +func TestLockExpiresAfterOwnerCrashes(t *testing.T) { + manager, client, server := newTestManager(t, 2*time.Second, 500*time.Millisecond) + ctx, cancel := context.WithCancel(context.Background()) + + lease, err := manager.Acquire(ctx, "cron") + require.NoError(t, err) + cancel() + + select { + case <-lease.renewalDone: + case <-time.After(time.Second): + t.Fatal("renewal routine did not stop after context cancellation") + } + + server.FastForward(3 * time.Second) + assert.Equal(t, redis.Nil, client.Get(context.Background(), LockKey).Err()) + assert.Empty(t, client.HGetAll(context.Background(), StateKey).Val()) +} + +func TestAcquireFailsClosedWhenValkeyIsUnavailable(t *testing.T) { + manager, _, server := newTestManager(t, 2*time.Second, 500*time.Millisecond) + server.Close() + + _, err := manager.Acquire(context.Background(), "cron") + + require.ErrorIs(t, err, ErrUnavailable) +} + +func TestReleaseStopsRenewalRoutine(t *testing.T) { + manager, _, _ := newTestManager(t, 2*time.Second, 500*time.Millisecond) + lease, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + + require.NoError(t, lease.Release(context.Background())) + + select { + case <-lease.renewalDone: + case <-time.After(time.Second): + t.Fatal("renewal routine was left running") + } +} + +func TestConfirmedOwnershipLossCancelsLease(t *testing.T) { + manager, client, _ := newTestManager(t, 200*time.Millisecond, 25*time.Millisecond) + lease, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + + require.NoError(t, client.Set(context.Background(), LockKey, "new-owner", time.Second).Err()) + + select { + case <-lease.Context().Done(): + case <-time.After(time.Second): + t.Fatal("lease context was not canceled after ownership loss") + } + require.ErrorIs(t, context.Cause(lease.Context()), ErrLost) + + require.ErrorIs(t, lease.Release(context.Background()), ErrLost) + assert.Equal(t, "new-owner", client.Get(context.Background(), LockKey).Val()) +} + +func TestTemporaryRenewalFailureDoesNotCancelLease(t *testing.T) { + var renewCalls atomic.Int32 + store := &fakeStore{ + renew: func(context.Context, State, time.Duration) (bool, error) { + if renewCalls.Add(1) == 1 { + return false, errors.New("temporary network error") + } + return true, nil + }, + } + manager, err := New(store, Config{ + TTL: 200 * time.Millisecond, + RenewInterval: 50 * time.Millisecond, + }) + require.NoError(t, err) + + lease, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + time.Sleep(125 * time.Millisecond) + + assert.NoError(t, context.Cause(lease.Context())) + assert.GreaterOrEqual(t, renewCalls.Load(), int32(2)) + require.NoError(t, lease.Release(context.Background())) +} + +func TestContinuousRenewalFailureCancelsBeforeTTL(t *testing.T) { + store := &fakeStore{ + renew: func(context.Context, State, time.Duration) (bool, error) { + return false, errors.New("valkey unavailable") + }, + } + manager, err := New(store, Config{ + TTL: 200 * time.Millisecond, + RenewInterval: 50 * time.Millisecond, + }) + require.NoError(t, err) + + lease, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + + select { + case <-lease.Context().Done(): + case <-time.After(190 * time.Millisecond): + t.Fatal("lease was not canceled before ttl") + } + require.ErrorIs(t, context.Cause(lease.Context()), ErrLost) + require.NoError(t, lease.Release(context.Background())) +} + +func TestReleaseReturnsStoreFailure(t *testing.T) { + store := &fakeStore{releaseErr: errors.New("connection reset")} + manager, err := New(store, Config{ + TTL: time.Second, + RenewInterval: 250 * time.Millisecond, + }) + require.NoError(t, err) + lease, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + + err = lease.Release(context.Background()) + + require.ErrorIs(t, err, ErrUnavailable) + require.ErrorContains(t, err, "connection reset") +} + +func TestReleaseReportsOwnershipChangeAsLockLost(t *testing.T) { + manager, client, _ := newTestManager(t, 2*time.Second, 500*time.Millisecond) + lease, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + require.NoError(t, client.Set(context.Background(), LockKey, "new-owner", time.Second).Err()) + + err = lease.Release(context.Background()) + + require.ErrorIs(t, err, ErrLost) + assert.Equal(t, "new-owner", client.Get(context.Background(), LockKey).Val()) +} + +func TestReleaseCancelsInFlightRenewalBeforeWaiting(t *testing.T) { + renewalStarted := make(chan struct{}) + store := &fakeStore{ + renew: func(ctx context.Context, _ State, _ time.Duration) (bool, error) { + close(renewalStarted) + <-ctx.Done() + return false, context.Cause(ctx) + }, + } + manager, err := New(store, Config{ + TTL: time.Second, + RenewInterval: 20 * time.Millisecond, + }) + require.NoError(t, err) + lease, err := manager.Acquire(context.Background(), "cron") + require.NoError(t, err) + + select { + case <-renewalStarted: + case <-time.After(time.Second): + t.Fatal("renewal did not start") + } + + released := make(chan error, 1) + go func() { + released <- lease.Release(context.Background()) + }() + + select { + case err := <-released: + require.NoError(t, err) + case <-time.After(200 * time.Millisecond): + t.Fatal("release waited indefinitely for in-flight renewal") + } +} + +func newTestManager( + t *testing.T, + ttl time.Duration, + renewInterval time.Duration, +) (*Manager, *redis.Client, *miniredis.Miniredis) { + t.Helper() + + server := miniredis.RunT(t) + client := redis.NewClient(&redis.Options{Addr: server.Addr()}) + t.Cleanup(func() { + _ = client.Close() + }) + + manager, err := New(NewValkeyStore(client), Config{ + TTL: ttl, + RenewInterval: renewInterval, + }) + require.NoError(t, err) + return manager, client, server +} + +type fakeStore struct { + renew func(context.Context, State, time.Duration) (bool, error) + releaseErr error +} + +func (s *fakeStore) TryAcquire(context.Context, string, time.Duration) (bool, error) { + return true, nil +} + +func (s *fakeStore) WriteState(context.Context, string, State, time.Duration) (bool, error) { + return true, nil +} + +func (s *fakeStore) Renew(ctx context.Context, _ string, state State, ttl time.Duration) (bool, error) { + if s.renew == nil { + return true, nil + } + return s.renew(ctx, state, ttl) +} + +func (s *fakeStore) Release(context.Context, string, string) (bool, error) { + if s.releaseErr != nil { + return false, fmt.Errorf("fake release: %w", s.releaseErr) + } + return true, nil +} diff --git a/scraper-go/internal/runlock/valkey.go b/scraper-go/internal/runlock/valkey.go new file mode 100644 index 0000000..a143946 --- /dev/null +++ b/scraper-go/internal/runlock/valkey.go @@ -0,0 +1,137 @@ +package runlock + +import ( + "context" + "strconv" + "time" + + "github.com/redis/go-redis/v9" +) + +const ( + LockKey = "scraper:run:lock" + StateKey = "scraper:run:state" +) + +var writeStateScript = redis.NewScript(` +if redis.call("GET", KEYS[1]) ~= ARGV[1] then + return 0 +end +redis.call("DEL", KEYS[2]) +redis.call("HSET", KEYS[2], + "runId", ARGV[2], + "source", ARGV[3], + "startedAt", ARGV[4], + "lockExpiresAt", ARGV[5]) +redis.call("PEXPIRE", KEYS[2], ARGV[6]) +return 1 +`) + +var renewScript = redis.NewScript(` +if redis.call("GET", KEYS[1]) ~= ARGV[1] then + return 0 +end +redis.call("PEXPIRE", KEYS[1], ARGV[6]) +if redis.call("HGET", KEYS[2], "runId") ~= ARGV[2] then + redis.call("DEL", KEYS[2]) +end +redis.call("HSET", KEYS[2], + "runId", ARGV[2], + "source", ARGV[3], + "startedAt", ARGV[4], + "lockExpiresAt", ARGV[5]) +redis.call("PEXPIRE", KEYS[2], ARGV[6]) +return 1 +`) + +var releaseScript = redis.NewScript(` +if redis.call("GET", KEYS[1]) ~= ARGV[1] then + return 0 +end +redis.call("DEL", KEYS[1]) +if redis.call("HGET", KEYS[2], "runId") == ARGV[2] then + redis.call("DEL", KEYS[2]) +end +return 1 +`) + +type Store interface { + TryAcquire(ctx context.Context, token string, ttl time.Duration) (bool, error) + WriteState(ctx context.Context, token string, state State, ttl time.Duration) (bool, error) + Renew(ctx context.Context, token string, state State, ttl time.Duration) (bool, error) + Release(ctx context.Context, token, runID string) (bool, error) +} + +type ValkeyStore struct { + client *redis.Client +} + +func NewValkeyStore(client *redis.Client) *ValkeyStore { + return &ValkeyStore{client: client} +} + +func (s *ValkeyStore) TryAcquire(ctx context.Context, token string, ttl time.Duration) (bool, error) { + result, err := s.client.Do(ctx, acquireCommand(token, ttl)...).Text() + if err == redis.Nil { + return false, nil + } + if err != nil { + return false, err + } + return result == "OK", nil +} + +func acquireCommand(token string, ttl time.Duration) []any { + return []any{"SET", LockKey, token, "NX", "PX", ttl.Milliseconds()} +} + +func (s *ValkeyStore) WriteState( + ctx context.Context, + token string, + state State, + ttl time.Duration, +) (bool, error) { + result, err := writeStateScript.Run( + ctx, + s.client, + []string{LockKey, StateKey}, + token, + state.RunID, + state.Source, + state.StartedAt.UTC().Format(time.RFC3339Nano), + state.LockExpiresAt.UTC().Format(time.RFC3339Nano), + strconv.FormatInt(ttl.Milliseconds(), 10), + ).Int64() + return result == 1, err +} + +func (s *ValkeyStore) Renew( + ctx context.Context, + token string, + state State, + ttl time.Duration, +) (bool, error) { + result, err := renewScript.Run( + ctx, + s.client, + []string{LockKey, StateKey}, + token, + state.RunID, + state.Source, + state.StartedAt.UTC().Format(time.RFC3339Nano), + state.LockExpiresAt.UTC().Format(time.RFC3339Nano), + strconv.FormatInt(ttl.Milliseconds(), 10), + ).Int64() + return result == 1, err +} + +func (s *ValkeyStore) Release(ctx context.Context, token, runID string) (bool, error) { + result, err := releaseScript.Run( + ctx, + s.client, + []string{LockKey, StateKey}, + token, + runID, + ).Int64() + return result == 1, err +}