Skip to content

feat: restaura compare_against em poll_source_for_update - #1783

Merged
Winzen merged 39 commits into
mainfrom
feat/poll_compare_against
Aug 12, 2026
Merged

feat: restaura compare_against em poll_source_for_update#1783
Winzen merged 39 commits into
mainfrom
feat/poll_compare_against

Conversation

@Winzen

@Winzen Winzen commented Aug 10, 2026

Copy link
Copy Markdown
Collaborator

Contexto

Fecha a causa raiz sistêmica por trás da issue #1781 e dos bugs já corrigidos em br_me_caged (#1760) e br_ans_beneficiario (#1779).

O Prefect 0 tinha, na função check_if_data_is_outdated, um parâmetro date_type que deixava cada flow escolher o alvo da comparação de "há dados novos":

  • "data_max_date" (default) → Coverage.DateTimeRange (cobertura real)
  • "last_update_date"Table.Update.latest (timestamp de execução, bq.last_modified)

Auditoria (worktree pré-migração em 1dff7cfc) mostrou que 19 dos ~23 datasets usavam o default (data_max_date, correto). Esse parâmetro se perdeu na migração pro Prefect 3: poll_source_for_update_task não tem escolha nenhuma — sempre compara contra Table.Update.latest. Não foi uma decisão consciente, foi uma regressão silenciosa.

O que muda

  • MetadataClient.get_coverage_max_date(dataset_id, table_id) — novo método de leitura, lê Coverage.DateTimeRange e devolve o maior end entre todas as faixas (free + pro). Mesma fonte de dados que o antigo get_api_most_recent_date consultava.
  • poll_source_for_update(..., compare_against="coverage") — novo parâmetro. "coverage" (default) usa get_coverage_max_date; "table_update" usa get_table_update_latest (o comportamento antigo, para as fontes onde source_max_date é de fato um timestamp de publicação/execução, não uma competência).
  • poll_source_for_update_task(..., compare_against="coverage") — mesmo parâmetro exposto no wrapper Prefect.
  • Validação: compare_against fora de {"coverage", "table_update"} levanta ValueError antes de qualquer escrita (um typo não pode mais gravar o Poll e avaliar contra o alvo errado silenciosamente).
  • Logs: poll_source_for_update agora sempre loga a fonte e o alvo comparados ("Comparando fonte em X contra Y em Z"), nos dois desfechos — antes só dizia "há/não há atualizações" sem mostrar os valores.
  • compare_against aplicado explicitamente nos 32 flows que usam poll_source_for_update_task hoje (consolidado do fix: aplica compare_against em todos os flows com poll_source_for_update_task #1784 nesta branch) — 27 com "coverage", 5 com "table_update".
  • Testes novos/ajustados em test_client.py e test_register.py cobrindo o método novo e os dois modos de comparação, incluindo um teste que reproduz exatamente o cenário do bug (Table.Update "adiantado" bloqueando, Coverage real corretamente detectando a novidade) e um teste do ValueError de validação.

Sobre o default

Inicialmente o default era "table_update" (preservava o comportamento de todo mundo enquanto nada tinha sido auditado). Depois de aplicar compare_against explicitamente em todos os 32 flows, trocamos o default para "coverage", já que é o caminho correto na grande maioria dos casos e nenhum flow existente depende mais do valor implícito.

Testes realizados em produção

br_me_caged ainda não foi migrado para o poll.py (PR #1760 aberto) — usamos ele nesta branch para validar compare_against="coverage" ao vivo, sem esperar aquele PR: voltamos a Coverage.DateTimeRange de cada tabela em 1 mês (2026-062026-05) e rodamos com force_run=False, materialize_after_dump=False (só toca dev, sem risco de duplicar produção — o filtro is_incremental() do modelo dbt compara contra o máximo já materializado na própria tabela de destino, então mesmo reprocessar um mês já coberto não duplicaria a tabela final).

Confirma o log novo ("Comparando fonte em ... contra coverage em ...") funcionando nos dois desfechos:

Flow Run Resultado
microdados_movimentacao chi23-barisa-reach fonte=2026-06-01 > coverage=2026-05-01 → detectou atualização
microdados_movimentacao_excluida upsilon58-lothlorien-reach fonte=2026-06-01 > coverage=2026-05-01 → detectou atualização
microdados_movimentacao_fora_prazo mu-konishi-shift fonte=2026-06-01 > coverage=2026-05-01 → detectou atualização
microdados_movimentacao_fora_prazo theta5-enara-quadrant fonte=2026-06-01 = coverage=2026-06-01negativo correto, sem atualização depois de já ter caught up

Depois de cada teste, a Coverage.DateTimeRange foi restaurada para o valor real (2026-06), confirmado contra o BigQuery.

Pyrefly type check corrigido

O check estava vermelho por um problema sem relação com este PR: o PR #1793 (onboard do us_harvard_cbdb) quebrou o Pyrefly na mainmodels/us_harvard_cbdb/code/*.py usa imports bare (from schema_spec import ...), assumindo cwd=code/ na execução, que não resolve a partir da raiz do repo. Isso já falhava na própria main antes de qualquer coisa nossa, e entrou nesta branch só porque sincronizamos com a main via merge.

Corrigido adicionando models/us_harvard_cbdb/code em project-excludes no pyproject.toml — mesma política já aplicada a models/br_tse_eleicoes/code (pacote .py, não notebook, com o mesmo padrão de import relativo ao cwd). uv run pyrefly check confirma 0 diagnostics localmente, e o check já está verde no CI.

Optamos por aplicar só nesta branch em vez de abrir um PR separado contra a main — mais rápido pra destravar este PR, cientes de que a main continua com o problema até alguém do time do CBDB corrigir ou até isso ser levado pra lá separadamente.

Fix adicional: tabelas NonHistorical não têm baseline de coverage

Um comentário de code review em br_me_cnpj/flows.py:61 apontou que a tabela simples usa cobertura NonHistorical() — e register_table_materialization (register.py:282-287), para esse tier, só grava Table.Update.latest, sem nenhum Coverage.DateTimeRange. Com compare_against="coverage" (o novo default), get_coverage_max_date sempre retorna None pra essa tabela, e should_update_raw_source trata api_latest is None como "sempre há atualização" — ou seja, o poll nunca daria negativo, rodando o flow todo dia mesmo sem novidade na fonte.

Verificação contra o Prefect 0: lá, simples também não passava nenhum date_type explícito (mesmo default data_max_date/coverage), mas funcionava porque update_django_metadata com historical_database=False gravava uma Coverage.DateTimeRange "falsa" (tipo all_free, start=end) derivada de update_date_from_bq_metadata — ou seja, do próprio last_modified_time do BQ. No Prefect 3, Table.Update.latest é gravado com esse mesmo bq.last_modified(). Logo, compare_against="table_update" pra essa tabela compara exatamente o mesmo valor que o Prefect 0 comparava — só mudou onde ele fica guardado.

Achamos o mesmo padrão em br_rf_cnpj/flows.py, que também usa NonHistorical() para simples e dicionario, sem nenhum compare_against definido (caindo no default "coverage").

Corrigido nos dois: compare_against="table_update" só para as tabelas NonHistorical (simples em me_cnpj; simples/dicionario em rf_cnpj), mantendo "coverage" nas demais.

Fix adicional: commit_source_update_task gravava o Update da fonte tarde demais

Feedback de review: commit_source_update_task (grava RawDataSource.Update.latest) era chamado só no fim do flow, depois de register_table_materialization_task — ou seja, depois de baixar, subir e materializar tudo. Isso significa que se o flow falhar no meio (download, upload, dbt), o metadado da fonte original nunca avança, mesmo que a fonte realmente tenha publicado dado novo — quem olha a API não tem visibilidade de que havia novidade, só de que a run falhou.

Movido para logo depois do poll confirmar que há dado novo, antes de baixar/materializar. poll_source_for_update nunca lê RawDataSource.Update pra decidir nada (compara contra Coverage ou Table.Update, não contra ele mesmo), então adiantar essa escrita não tem risco de travar runs futuras — só ganha visibilidade: mesmo que o resto do flow falhe, o RawDataSource.Update já reflete que a fonte publicou.

commit_source_update_task ganhou dois parâmetros novos, update_metadata e materialize_after_dump (default True), e passou a decidir sozinha se grava (incluindo o source_max_date is not None) — isso evita repetir o mesmo if em cada um dos ~25 flows que chamam essa task; cada flow só passa os dois flags que já tem, sem precisar envolver a chamada num if externo. register_table_materialization_task continua no fim, sem mudança nenhuma.

Aplicado em todos os flows que usam esse par de funções, exceto onde a premissa não se aplica: br_me_siconfi e br_mf_divida_ativa (poll não-gating — a decisão de baixar vem de outra lógica, não do retorno do poll) e br_senado_dados_abertos (sem gate de poll nenhum).

Testado ao vivo em br_me_caged__microdados_movimentacao_excluida: apontamos temporariamente o deployment pra esta branch, voltamos Coverage e RawDataSource.Update em 1 mês, e confirmamos no log ("Data de atualização da fonte original modificada") que o Update da fonte é gravado logo após o poll, antes do download — nos dois runs que dispararam (inclusive um que teve uma falha de concorrência de IAM não relacionada ao nosso código, numa corrida entre duas execuções simultâneas do mesmo flow).

PRs relacionados

…1781)

O Prefect 0 tinha um parâmetro (date_type) em check_if_data_is_outdated
que deixava cada flow escolher se a comparação era contra a cobertura
real (Coverage.DateTimeRange, o default e o correto pra a maioria dos
~23 datasets auditados) ou contra Table.Update.latest (um timestamp de
execução). Esse parâmetro se perdeu na migração pro Prefect 3 —
poll_source_for_update sempre comparou contra Table.Update.latest,
causando o bug já corrigido pontualmente em br_me_caged (#1760) e
br_ans_beneficiario (#1779).

Adiciona compare_against="coverage"|"table_update" (default
"table_update", preserva o comportamento atual de todo mundo) em
poll_source_for_update e poll_source_for_update_task. "coverage" usa
o novo MetadataClient.get_coverage_max_date, que lê
Coverage.DateTimeRange (mesma fonte que o antigo get_api_most_recent_date
consultava) em vez de Table.Update.latest.

Próximo passo: setar compare_against="coverage" explicitamente nos
~17 flows que usavam data_max_date por padrão no Prefect 0 (auditoria
completa documentada no vault).
@coderabbitai

coderabbitai Bot commented Aug 10, 2026

Copy link
Copy Markdown

Review Change Stack

Note

Reviews paused

It looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

Adds MetadataClient.get_coverage_max_date and extends metadata polling with the compare_against option. Polling now defaults to coverage comparison and supports explicit comparison with Table.Update.latest. Dataset and crawler flows select the appropriate comparison mode.

Changes

Coverage-based metadata polling

Layer / File(s) Summary
Coverage date retrieval and validation
pipelines/utils/metadata/client.py, pipelines/utils/tests/metadata/conftest.py, pipelines/utils/tests/metadata/test_client.py
MetadataClient returns the latest valid coverage end date. Missing month and day values default to 1. Tests cover empty and multiple coverage ranges.
Configurable polling and flow wiring
pipelines/utils/metadata/register.py, pipelines/utils/metadata/tasks.py, pipelines/utils/tests/metadata/test_register.py, pipelines/datasets/*/flows.py, pipelines/crawler/*/flows.py
Polling validates compare_against, defaults to coverage, supports table_update, forwards the option through the task, and updates dataset and crawler flows. Tests cover both modes and invalid values.

Type-checking configuration

Layer / File(s) Summary
Pyrefly exclusion
pyproject.toml
Pyrefly excludes models/us_harvard_cbdb/code and documents its cwd-dependent imports.

Estimated code review effort: 3 (Moderate) | ~25 minutes

Sequence Diagram(s)

sequenceDiagram
  participant poll_source_for_update_task
  participant poll_source_for_update
  participant MetadataClient
  poll_source_for_update_task->>poll_source_for_update: pass compare_against
  alt compare_against is coverage
    poll_source_for_update->>MetadataClient: get_coverage_max_date
    MetadataClient-->>poll_source_for_update: maximum coverage date
  else compare_against is table_update
    poll_source_for_update->>MetadataClient: read Table.Update.latest
    MetadataClient-->>poll_source_for_update: table update date
  end
  poll_source_for_update-->>poll_source_for_update_task: return freshness boolean
Loading

Possibly related issues

Possibly related PRs

Suggested labels: check-metadata, bug

Suggested reviewers: aspeddro, laura-l-amaral

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 34.88% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
Title check ✅ Passed The title clearly identifies the restoration of compare_against in poll_source_for_update, which is the main change.
Description check ✅ Passed The description clearly explains the context, technical changes, tests, production validation, risks, and related PRs.
✨ Finishing Touches 💡 1
📝 Generate docstrings 💡
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feat/poll_compare_against

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 2

🧹 Nitpick comments (2)
pipelines/utils/tests/metadata/test_client.py (1)

291-315: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Test ranges without endYear.

MetadataClient.get_coverage_max_date skips a range when endYear is None. The new tests do not execute that branch. Add a mixed-range case with one missing endYear and one valid end date. Assert that the valid date is returned.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pipelines/utils/tests/metadata/test_client.py` around lines 291 - 315, Extend
test_coverage_max_date_defaults_missing_month_day_to_one with a mixed
datetimeRanges case containing one entry whose endYear is None and another with
a valid end date. Assert that client.get_coverage_max_date returns the valid
date, covering the skip behavior for ranges missing endYear while preserving
defaulting for missing month/day.
pipelines/utils/metadata/client.py (1)

244-255: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Add the required type hints and Google-style docstrings.

  • pipelines/utils/metadata/client.py#L244-L255: Add Args and Returns sections to get_coverage_max_date.
  • pipelines/utils/tests/metadata/conftest.py#L123-L136: Type the new coverage_max_date fixture parameter and document the constructor.
  • pipelines/utils/tests/metadata/conftest.py#L158-L159: Add parameter and return annotations plus a docstring.
  • pipelines/utils/tests/metadata/test_client.py#L248-L315: Add -> None annotations and docstrings to the new test functions.
  • pipelines/utils/metadata/register.py#L315-L358: Type client, dataset_id, and table_id.
  • pipelines/utils/tests/metadata/test_register.py#L110-L153: Add -> None annotations and docstrings to the new test functions.

As per coding guidelines, add Google-Style type hints and docstrings to functions.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pipelines/utils/metadata/client.py` around lines 244 - 255, Apply
Google-style type hints and docstrings across the coverage metadata additions:
in pipelines/utils/metadata/client.py:244-255, document dataset_id and table_id
under Args and the date-or-None result under Returns for get_coverage_max_date;
in pipelines/utils/tests/metadata/conftest.py:123-136, type coverage_max_date
and document the constructor, and at :158-159 add parameter/return annotations
and a docstring; in pipelines/utils/tests/metadata/test_client.py:248-315 and
pipelines/utils/tests/metadata/test_register.py:110-153, add -> None annotations
and docstrings to the new test functions; in
pipelines/utils/metadata/register.py:315-358, add types for client, dataset_id,
and table_id.

Source: Coding guidelines

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@pipelines/utils/metadata/register.py`:
- Around line 371-374: Validate compare_against before the poll-writing flow,
allowing only "table_update" and "coverage"; raise ValueError for any other
value before upsert_raw_source_poll or freshness evaluation. Preserve the
existing API selection for the two supported values, and add a test confirming
invalid input performs no write.

In `@pipelines/utils/tests/metadata/test_register.py`:
- Around line 110-122: Update test_poll_default_compares_against_table_update so
FakeMetadataClient.coverage_max_date is later than source_max_date, such as July
1, 2026, while keeping table_update_latest at May 1, 2026. Preserve the True
expectation and default comparison mode to verify the result uses table_update
rather than coverage.

---

Nitpick comments:
In `@pipelines/utils/metadata/client.py`:
- Around line 244-255: Apply Google-style type hints and docstrings across the
coverage metadata additions: in pipelines/utils/metadata/client.py:244-255,
document dataset_id and table_id under Args and the date-or-None result under
Returns for get_coverage_max_date; in
pipelines/utils/tests/metadata/conftest.py:123-136, type coverage_max_date and
document the constructor, and at :158-159 add parameter/return annotations and a
docstring; in pipelines/utils/tests/metadata/test_client.py:248-315 and
pipelines/utils/tests/metadata/test_register.py:110-153, add -> None annotations
and docstrings to the new test functions; in
pipelines/utils/metadata/register.py:315-358, add types for client, dataset_id,
and table_id.

In `@pipelines/utils/tests/metadata/test_client.py`:
- Around line 291-315: Extend
test_coverage_max_date_defaults_missing_month_day_to_one with a mixed
datetimeRanges case containing one entry whose endYear is None and another with
a valid end date. Assert that client.get_coverage_max_date returns the valid
date, covering the skip behavior for ranges missing endYear while preserving
defaulting for missing month/day.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro Plus

Run ID: 7dcdfd20-f7bd-4d69-81ee-fb4f105adc84

📥 Commits

Reviewing files that changed from the base of the PR and between 471c851 and 61ef8d0.

📒 Files selected for processing (6)
  • pipelines/utils/metadata/client.py
  • pipelines/utils/metadata/register.py
  • pipelines/utils/metadata/tasks.py
  • pipelines/utils/tests/metadata/conftest.py
  • pipelines/utils/tests/metadata/test_client.py
  • pipelines/utils/tests/metadata/test_register.py

Comment thread pipelines/utils/metadata/register.py
Comment thread pipelines/utils/tests/metadata/test_register.py Outdated
Aplica o compare_against restaurado em #1783 em todos os 18 flows que
ainda usam poll_source_for_update_task e têm equivalente no Prefect 0
(auditoria completa documentada no vault).

compare_against="coverage" (14 arquivos, 19 datasets — usavam
date_type="data_max_date" por padrão no Prefect 0, perdido na
migração):
- anatel (banda_larga_fixa, telefonia_movel)
- cgu (cartao_pagamento, servidores_publicos, licitacao_contrato,
  beneficios_cidadao)
- datasus (siasus/sihsus via _run_dbf_to_parquet, sinan)
- ibge_inflacao (inpc, ipca, ipca15)
- me_cnpj (empresas, socios, estabelecimentos, simples)
- rf (cno)
- br_anp_precos_combustiveis, br_bcb_agencia, br_bcb_estban,
  br_denatran_frota, br_inmet_bdmep, br_rf_cafir, br_stf_corte_aberta

compare_against="table_update" (4 arquivos — já usavam
date_type="last_update_date" de propósito no Prefect 0; valor
explícito por clareza, já era o default):
- cvm_fi, tse_eleicoes, br_cgu_emendas_parlamentares, br_ibge_pnadc

br_me_caged e br_ans_beneficiario ficam de fora — já corrigidos via
poll.py (#1760, #1779), não usam mais poll_source_for_update_task.
Winzen and others added 3 commits August 10, 2026 19:21
Completa os 9 flows (10 chamadas) que ainda dependiam do default
implícito de poll_source_for_update_task — os que ficaram de fora da
auditoria original por não terem equivalente no Prefect 0. Decisão
caso a caso pela natureza de cada fonte, não por convenção:

compare_against="coverage" (8 datasets — source_max_date representa
uma competência/data real dos dados, e na maioria dos casos o poll
efetivamente trava o flow):
- br_bcb_taxa_selic: max date da série diária buscada na API do BCB
- br_me_siconfi, br_mf_divida_ativa: max ano/trimestre disponível na
  fonte (poll não-gating nos dois, mas a comparação em si é correta)
- au_abs_cpi, au_abs_labour_force, us_bls_cpi, us_bls_qcew,
  world_cricsheet: max ano-mês/trimestre/data de partida publicado
  pela fonte, poll gating em todos

compare_against="table_update" (br_bndes_operacoes_contratadas, 2
flows no mesmo arquivo) — exceção deliberada, não bug: source_max_date
aqui é o last_modified do recurso no CKAN (timestamp de publicação,
granularidade dia/hora), enquanto a Coverage da tabela é anual
(coluna `ano`) — comparar contra Coverage misturaria um timestamp de
publicação com uma cobertura anual, duas granularidades desconexas. O
próprio código já documentava essa distinção antes desta mudança.
A auditoria (#1784) mostrou que 27 dos 32 flows que usam
poll_source_for_update_task precisam de compare_against="coverage" —
só 5 precisam de "table_update" de propósito. Com todos os callers já
explícitos (#1784), trocar o default não muda comportamento de
nenhum flow existente; só afeta flows novos que omitirem o parâmetro,
alinhando o caminho de menor esforço com a opção correta na maioria
dos casos.

Ajusta os testes que dependiam do default antigo:
- poll_source_for_update: 3 testes reescritos (default agora testa
  coverage; compare_against="table_update" explícito ganha teste
  próprio para os casos como br_bndes que precisam dele).
- register_source_poll (não exposto a compare_against, herda o
  default): 2 testes trocam o fixture de table_update_latest para
  coverage_max_date. Sem uso em nenhum flow hoje — só cobertura de
  teste.

Suite completa: 132 passando, 1 falha pré-existente e não relacionada
(test_policy.py::test_compute_part_bdpro_monthly_syncs_and_lags — data
"hoje" hardcoded que já cruzou de mês, falha igual sem esta mudança).
@rdahis

rdahis commented Aug 10, 2026

Copy link
Copy Markdown
Member

Closes #1687.

mergify Bot and others added 12 commits August 11, 2026 06:13
Qualquer valor diferente de "coverage" caía no branch else e avaliava
a novidade contra Table.Update.latest silenciosamente — um typo em
compare_against gravaria o Poll e compararia contra o alvo errado sem
avisar ninguém. Valida contra VALID_COMPARE_AGAINST = {"coverage",
"table_update"} e levanta ValueError antes de qualquer escrita.
Antes o log só dizia "há/não há atualizações", sem mostrar o que foi
comparado — dificultava diagnosticar exatamente o tipo de bug que
travou o CAGED e o ANS (Table.Update "adiantado" bloqueando). Agora
sempre loga fonte vs. alvo (coverage ou table_update) antes do
resultado, e repete os valores na mensagem final.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 2

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@pipelines/utils/metadata/register.py`:
- Line 326: Update the function containing compare_against in
pipelines/utils/metadata/register.py (326-326) with type hints for client,
dataset_id, and table_id. Add Google-style docstrings and -> None return
annotations to the four test functions in
pipelines/utils/tests/metadata/test_register.py at 114-114, 131-131, 139-139,
and 154-154, preserving the Python 3.10 and 79-character Ruff conventions.

In `@pipelines/utils/tests/metadata/test_register.py`:
- Around line 160-173: Update the FakeMetadataClient setup in the explicit-mode
test for poll_source_for_update so coverage_max_date is later than the
source_max_date, while table_update_latest remains earlier. Keep
compare_against="table_update" and the True assertion, ensuring the test
distinguishes the selected comparison target.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro Plus

Run ID: 7f5c8515-f1e4-4898-9462-ed3f7a81dea5

📥 Commits

Reviewing files that changed from the base of the PR and between 61ef8d0 and 4d71e30.

📒 Files selected for processing (3)
  • pipelines/utils/metadata/register.py
  • pipelines/utils/metadata/tasks.py
  • pipelines/utils/tests/metadata/test_register.py
🚧 Files skipped from review as they are similar to previous changes (1)
  • pipelines/utils/metadata/tasks.py

Comment thread pipelines/utils/metadata/register.py
Comment thread pipelines/utils/tests/metadata/test_register.py
… pro poll.py)

Teste da abordagem alternativa ao PR #1760 (migração completa pro
poll.py): br_me_caged ainda usa poll_source_for_update_task, mas agora
com compare_against="coverage" em vez do default antigo
(Table.Update.latest, o campo contaminado que causou o bug original).

register_table_materialization_task já atualiza Coverage.DateTimeRange
a partir do BigQuery corretamente — isso nunca foi o problema. Só a
comparação de leitura estava errada. Se compare_against="coverage"
sozinho resolver, não precisamos da migração inteira pro poll.py.

Afeta os 3 flows (microdados_movimentacao, _fora_prazo, _excluida) via
_run_me_caged compartilhado. Teste-piloto: microdados_movimentacao_fora_prazo.
@Winzen Winzen added the deploy-flow [PR] Dispara deploy dos flows alterados no work pool basedosdados-dev (Prefect 3 staging) label Aug 11, 2026
Winzen and others added 4 commits August 12, 2026 06:17
Trazido do PR #1784 pra testar o caminho table_update (só coverage
tinha sido validado ao vivo até agora, via br_me_caged). Prefect 0 já
usava date_type="last_update_date" de propósito nesse dataset
(datasets/br_ibge_pnadc/flows.py:57 no worktree pré-migração) — não é
um julgamento novo, é uma escolha já validada historicamente.
…ompare_against

Consolida o PR #1784 (compare_against aplicado em todos os flows
restantes) dentro do #1783, pra não depender de dois merges
sequenciais rápidos. br_me_caged (fora do #1784, migrado à parte pro
teste do compare_against) e br_ibge_pnadc (já trazido manualmente
antes) não aparecem no diff — já estavam corretos nesta branch.
@Winzen Winzen linked an issue Aug 12, 2026 that may be closed by this pull request
@Winzen
Winzen requested a review from a team August 12, 2026 09:49
mergify Bot and others added 5 commits August 12, 2026 09:50
Bug pré-existente na main (PR #1793, onboard CBDB): code/*.py usa
imports bare (from schema_spec import ...), assumindo cwd=code/ na
execução — não resolve a partir da raiz do repo. Mesmo padrão já
excluído para br_tse_eleicoes/code. Sem isso, o CI deste PR ficava
vermelho por um problema sem relação com compare_against.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@pipelines/crawler/me_cnpj/flows.py`:
- Line 61: The simples flow should not use coverage as the comparison baseline
because its NonHistorical configuration has no coverage end date. Update the
compare_against setting in the simples flow to use the compatible source
publication timestamp with table_update, or configure date-bearing coverage for
the table, and add a regression test verifying the scheduled no-op path
short-circuits.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Pro Plus

Run ID: b15a56de-4e98-4318-ad97-f02542bb3497

📥 Commits

Reviewing files that changed from the base of the PR and between 17878fd and 5ababf1.

📒 Files selected for processing (27)
  • pipelines/crawler/anatel/banda_larga_fixa/flows.py
  • pipelines/crawler/anatel/telefonia_movel/flows.py
  • pipelines/crawler/bndes/flows.py
  • pipelines/crawler/cgu/flows.py
  • pipelines/crawler/cvm/flows.py
  • pipelines/crawler/datasus/flows.py
  • pipelines/crawler/ibge_inflacao/flows.py
  • pipelines/crawler/me_cnpj/flows.py
  • pipelines/crawler/rf/flows.py
  • pipelines/crawler/tse_eleicoes/flows.py
  • pipelines/datasets/au_abs_cpi/flows.py
  • pipelines/datasets/au_abs_labour_force/flows.py
  • pipelines/datasets/br_anp_precos_combustiveis/flows.py
  • pipelines/datasets/br_bcb_agencia/flows.py
  • pipelines/datasets/br_bcb_estban/flows.py
  • pipelines/datasets/br_bcb_taxa_selic/flows.py
  • pipelines/datasets/br_cgu_emendas_parlamentares/flows.py
  • pipelines/datasets/br_denatran_frota/flows.py
  • pipelines/datasets/br_inmet_bdmep/flows.py
  • pipelines/datasets/br_me_siconfi/flows.py
  • pipelines/datasets/br_mf_divida_ativa/flows.py
  • pipelines/datasets/br_rf_cafir/flows.py
  • pipelines/datasets/br_stf_corte_aberta/flows.py
  • pipelines/datasets/us_bls_cpi/flows.py
  • pipelines/datasets/us_bls_qcew/flows.py
  • pipelines/datasets/world_cricsheet/flows.py
  • pyproject.toml

Comment thread pipelines/crawler/me_cnpj/flows.py Outdated
Winzen added 2 commits August 12, 2026 16:02
simples (me_cnpj) e simples/dicionario (rf_cnpj) usam coverage NonHistorical,
que não grava Coverage.DateTimeRange — só Table.Update.latest. Com
compare_against="coverage" (default), o poll nunca encontra baseline e
sempre reporta atualização. table_update compara contra o mesmo valor
(bq.last_modified) que o Prefect 0 usava para essas tabelas.
commit_source_update_task (grava RawDataSource.Update.latest) era chamado
só ao fim do flow, depois de register_table_materialization_task. Movido
para logo depois do poll confirmar dado novo, antes de baixar/materializar.

poll_source_for_update nunca lê RawDataSource.Update (compara contra
Coverage ou Table.Update), então não há risco de travar polls futuros ao
adiantar esse valor. O benefício: se o flow falhar no meio (download,
upload, dbt), o metadado da fonte já reflete que havia dado novo publicado,
mesmo que a tabela não tenha sido atualizada — dando visibilidade ao
usuário mesmo quando a materialização não termina.

commit_source_update_task ganhou os parâmetros update_metadata e
materialize_after_dump (default True) e passou a decidir sozinha se grava
(inclusive o "source_max_date is not None"). Isso evita repetir o mesmo
guard em ~25 call sites — cada flow só chama a task direto, passando os
flags que já tem; a task decide.

register_table_materialization_task continua no fim, sem mudança.

Aplicado em todos os flows que usam esse par de funções, exceto onde a
premissa não se aplica: br_me_siconfi e br_mf_divida_ativa (poll
não-gating — a decisão vem de outra lógica) e br_senado_dados_abertos
(sem gate de poll).

Docstrings de commit_source_update/commit_source_update_task corrigidas
para refletir a orientação nova.
@mergify

mergify Bot commented Aug 12, 2026

Copy link
Copy Markdown
Contributor

Tick the box to add this pull request to the merge queue (same as @mergifyio queue).

  • Queue this pull request

mergify Bot and others added 3 commits August 12, 2026 21:04
test_poll_explicit_table_update_still_works usava coverage_max_date antes
da fonte, igual ao table_update_latest — o teste passava mesmo se
compare_against="table_update" fosse ignorado e coverage lido por engano.
Agora coverage_max_date fica depois da fonte, então só passa se
table_update for de fato o campo lido.

Também adiciona type hints (dataset_id: str, table_id: str) em
poll_source_for_update, consistente com o resto do arquivo.
@Winzen
Winzen merged commit 4b71e36 into main Aug 12, 2026
10 checks passed
@Winzen
Winzen deleted the feat/poll_compare_against branch August 12, 2026 22:52
Winzen added a commit that referenced this pull request Aug 13, 2026
A migração pro poll.py que este PR propunha não é mais necessária: o
#1783 já corrigiu poll_source_for_update via compare_against, e
br_ans_beneficiario já usa esse default corrigido. Mantido só o que
continua valendo independente do mecanismo de poll:

- RAW_COLLUNS_TYPE: colunas de texto categóricas (baixa/média
  cardinalidade) trocadas de str para category — reduz bastante o
  footprint do DataFrame sem mudar o valor persistido no parquet.
- MODALIDADE_OPERADORA recasteada para category depois do remove_accents;
  del df + gc.collect() por estado em parquet_partition — sem isso a
  memória de cada estado se acumulava até o gc.collect() do loop de fora,
  já causou OOM num arquivo pequeno logo depois de um grande.
- source_format="parquet" nos dois upload_to_gcs: crawler_ans grava
  .parquet, e sem declarar o formato o dump_header procurava .csv e não
  achava nada.
- job_variables={"memory": "3Gi"}: pico medido em produção após a
  otimização foi ~1.78Gi; 3Gi dá ~1.7x de margem.
- compare_against="coverage" explícito no poll (já era o comportamento
  via default desde o #1783; deixado explícito por consistência com os
  outros 26 flows).

Também exclui models/world_aiddata_gcdf/code do Pyrefly (mesmo padrão do
br_tse_eleicoes/us_harvard_cbdb/us_cfpb_hmda: pacote .py com imports
relativos ao cwd, não notebook) — estava quebrando o type check na main,
sem relação com esta mudança.
Winzen added a commit that referenced this pull request Aug 13, 2026
O #1783 moveu commit_source_update_task pra logo após o poll, removendo
a chamada antiga do fim do flow em todos os arquivos afetados. O #1798,
mergeado depois mas cortado de uma base anterior a esse merge, ainda
tinha essa chamada antiga e só editou parâmetros dela (date_format).
O merge do git combinou as duas sem detectar a duplicação lógica —
commit_source_update_task passou a ser chamado duas vezes por run.

Sem impacto de dados (é idempotente), só uma escrita redundante.
Winzen added a commit that referenced this pull request Aug 14, 2026
…a em rf_cnpj (#1779)

A migração pro poll.py que este PR propunha originalmente não é mais
necessária: o #1783 já corrigiu poll_source_for_update via
compare_against, e br_ans_beneficiario já usa esse default corrigido.

Mantido só o que continua valendo independente do mecanismo de poll:

- br_ans_beneficiario: colunas de texto categóricas trocadas de str para
  category (reduz footprint de memória do DataFrame), del/gc.collect()
  por estado em parquet_partition (corrige OOM), source_format="parquet"
  nos upload_to_gcs (sem isso o dump_header procurava .csv e não achava
  nada), job_variables={"memory": "3Gi"} (pico medido ~1.78Gi depois da
  otimização), compare_against="coverage" explícito no poll.
- rf_cnpj: remove uma chamada duplicada de commit_source_update_task que
  sobrou de um merge entre o #1783 (moveu a chamada pra cima) e o #1798
  (mergeado depois, cortado de uma base anterior, ainda tinha a chamada
  antiga) — sem impacto de dados (idempotente), só uma escrita redundante.
- Exclui models/world_aiddata_gcdf/code do Pyrefly (mesmo padrão de
  br_tse_eleicoes/us_harvard_cbdb/us_cfpb_hmda) — destrava o CI deste PR.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

deploy-flow [PR] Dispara deploy dos flows alterados no work pool basedosdados-dev (Prefect 3 staging)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[bug] ajuste no sistema de atualização

3 participants