diff --git a/pipelines/utils/metadata/bq.py b/pipelines/utils/metadata/bq.py index 7a7c38195..7cbb6614d 100644 --- a/pipelines/utils/metadata/bq.py +++ b/pipelines/utils/metadata/bq.py @@ -22,6 +22,7 @@ ) from pipelines.utils.metadata.utils import ( able_to_query_bigquery_metadata, + extract_first_date_from_bq, extract_last_date_from_bq, update_date_from_bq_metadata, update_row_access_policy, @@ -71,6 +72,29 @@ def read_max_date( coverage.date_format.value, ).date() + def read_min_date( + self, dataset_id: str, table_id: str, coverage: CoverageSpec + ) -> datetime.date: + """Início real da série — MIN da coluna de cobertura. Simétrica a + `read_max_date`; usada só para `part_bdpro`, para emitir o range free + completo e impedir que ele inverta (ver `compute_coverage_ranges`).""" + first_date = extract_first_date_from_bq( + dataset_id, + table_id, + # pyrefly: ignore [missing-attribute] + coverage.date_format.value, + # pyrefly: ignore [missing-attribute] + date_column_to_legacy_dict(coverage.date_column), + self.billing_project_id, + self.bq_project, + ) + return datetime.datetime.strptime( + # pyrefly: ignore [missing-attribute] + first_date, + # pyrefly: ignore [missing-attribute] + coverage.date_format.value, + ).date() + def last_modified( self, dataset_id: str, table_id: str ) -> datetime.datetime: diff --git a/pipelines/utils/metadata/domain.py b/pipelines/utils/metadata/domain.py index d65aca187..4d3a1b0aa 100644 --- a/pipelines/utils/metadata/domain.py +++ b/pipelines/utils/metadata/domain.py @@ -86,6 +86,12 @@ def as_relativedelta(self) -> relativedelta: class _CoverageBase(BaseModel): date_column: DateColumn date_format: DateFormat + # Passo da série, em unidades da granularidade de `date_format` — o "(N)" da + # notação de cobertura da BD ("2004(1)2022"). Quase sempre 1 (série contínua + # anual/mensal/diária), mas nem sempre: eleições brasileiras são bienais, com + # `interval=2`. Cada `DateTimeRange` gravado carrega este valor, então uma + # série bienal precisa declará-lo aqui — senão a escrita o fixaria em 1. + interval: int = Field(default=1, ge=1) @model_validator(mode="after") def _column_matches_format(self) -> _CoverageBase: diff --git a/pipelines/utils/metadata/dto.py b/pipelines/utils/metadata/dto.py index 29cb44388..9c5124567 100644 --- a/pipelines/utils/metadata/dto.py +++ b/pipelines/utils/metadata/dto.py @@ -22,6 +22,7 @@ Year = Annotated[int, Field(ge=1900, le=2100)] Month = Annotated[int, Field(ge=1, le=12)] Day = Annotated[int, Field(ge=1, le=31)] +Interval = Annotated[int, Field(ge=1)] def _to_iso8601(value: object) -> str: @@ -46,7 +47,17 @@ def _to_iso8601(value: object) -> str: class DateTimeRangeInput(BaseModel): - """Payload de `CreateUpdateDateTimeRange` (Coverage.DateTimeRange).""" + """Payload de `CreateUpdateDateTimeRange` (Coverage.DateTimeRange). + + `interval` é sempre enviado. O backend o **exige** ao criar um range com + início e fim ("Interval must exist in ranges with start and end dates"); sem + ele, `upsert_coverage_datetime_range` só conseguia UPDATE de um range já + existente, nunca CREATE. O valor é o passo da série (o "(N)" da notação de + cobertura): 1 para a maioria (anual/mensal/diária contínua), mas nem sempre — + eleições brasileiras são bienais (`interval=2`). No caminho das pipelines ele + vem de `CoverageSpec.interval`; o default 1 aqui só serve à construção direta + do DTO. + """ coverage: UUIDStr startYear: Year | None = None @@ -55,6 +66,7 @@ class DateTimeRangeInput(BaseModel): endYear: Year | None = None endMonth: Month | None = None endDay: Day | None = None + interval: Interval = 1 @model_validator(mode="after") def _shape_consistent(self) -> DateTimeRangeInput: diff --git a/pipelines/utils/metadata/policy.py b/pipelines/utils/metadata/policy.py index 84dc39836..e331d551b 100644 --- a/pipelines/utils/metadata/policy.py +++ b/pipelines/utils/metadata/policy.py @@ -125,7 +125,10 @@ def _next_period(d: date, fmt: DateFormat) -> date: def compute_coverage_ranges( - spec: CoverageSpec, source_end: date, coverage_ids: CoverageIds + spec: CoverageSpec, + source_end: date, + coverage_ids: CoverageIds, + source_start: date | None = None, ) -> CoverageRanges: """Calcula os DateTimeRange free e/ou pro. @@ -133,17 +136,28 @@ def compute_coverage_ranges( - all_bdpro → range pro terminando em `source_end`. - part_bdpro→ pro termina em `source_end`; free termina em `source_end - free_lag`; pro começa onde a free termina. + + `source_start` (o início real da série, lido do BigQuery) é usado só no ramo + part_bdpro: quando informado, o range free sai **completo** (início e fim), + de modo que a escrita seja auto-consistente e não dependa do início já + gravado. Isso corrige a inversão do range free numa tabela cujo histórico + inteiro cabe dentro da janela do paywall (`source_start > free_end` — um + snapshot sem história), em que o update só-do-fim deixava início > fim e o + backend recusava. Quando `source_start` é `None` (chamadas legadas/testes), + o free sai só-com-fim, exatamente como antes. """ if isinstance(spec, NonHistorical): raise ValueError("NonHistorical não usa compute_coverage_ranges") fmt = spec.date_format + interval = spec.interval # passo da série (o "(N)" da notação); ver domain if isinstance(spec, AllFree): return CoverageRanges( free=DateTimeRangeInput( # pyrefly: ignore [bad-argument-type] coverage=coverage_ids.free, + interval=interval, **_components(source_end, fmt, "end"), ), free_end=source_end, @@ -154,20 +168,33 @@ def compute_coverage_ranges( pro=DateTimeRangeInput( # pyrefly: ignore [bad-argument-type] coverage=coverage_ids.pro, + interval=interval, **_components(source_end, fmt, "end"), ) ) # part_bdpro free_end = source_end - spec.free_lag.as_relativedelta() + free_fields = _components(free_end, fmt, "end") + if source_start is not None: + # Emite o range free completo (início + fim) e trava o início em + # `free_end`: numa tabela cujo histórico inteiro está dentro da janela + # paga (`source_start > free_end`), `min` colapsa o range para + # `[free_end, free_end]` — válido e não-invertido — em vez de início > + # fim. Assim que a história ultrapassa `free_end`, `min` escolhe + # `source_start` e o free passa a cobrir a série toda normalmente. + free_start = min(source_start, free_end) + free_fields = {**_components(free_start, fmt, "start"), **free_fields} free = DateTimeRangeInput( # pyrefly: ignore [bad-argument-type] coverage=coverage_ids.free, - **_components(free_end, fmt, "end"), + interval=interval, + **free_fields, ) pro = DateTimeRangeInput( # pyrefly: ignore [bad-argument-type] coverage=coverage_ids.pro, + interval=interval, **_components(source_end, fmt, "end"), # free termina em free_end inclusive, então pro começa no período # seguinte: as coberturas são mutuamente exclusivas. diff --git a/pipelines/utils/metadata/poll.py b/pipelines/utils/metadata/poll.py index 0284093b0..91c2be708 100644 --- a/pipelines/utils/metadata/poll.py +++ b/pipelines/utils/metadata/poll.py @@ -20,7 +20,11 @@ from pipelines.utils.metadata import policy from pipelines.utils.metadata.client import MetadataClient -from pipelines.utils.metadata.domain import CoverageSpec, NonHistorical +from pipelines.utils.metadata.domain import ( + CoverageSpec, + NonHistorical, + PartBdpro, +) from pipelines.utils.metadata.register import BQReader from pipelines.utils.utils import log @@ -170,8 +174,18 @@ def sync_table_coverage( policy.assert_coverage_topology(coverage, coverage_ids) + # Só part_bdpro lê o início da série (para o range free sair completo e não + # inverter — ver compute_coverage_ranges); os demais tiers evitam o scan. + source_start = ( + bq.read_min_date( + dataset_id=dataset_id, table_id=table_id, coverage=coverage + ) + if isinstance(coverage, PartBdpro) + else None + ) + ranges = policy.compute_coverage_ranges( - coverage, source_coverage, coverage_ids + coverage, source_coverage, coverage_ids, source_start=source_start ) for coverage_range in ranges.to_list(): diff --git a/pipelines/utils/metadata/register.py b/pipelines/utils/metadata/register.py index 054186cd7..369b440ea 100644 --- a/pipelines/utils/metadata/register.py +++ b/pipelines/utils/metadata/register.py @@ -48,6 +48,9 @@ class BQReader(Protocol): def read_max_date( self, dataset_id: str, table_id: str, coverage: CoverageSpec ) -> datetime.date: ... + def read_min_date( + self, dataset_id: str, table_id: str, coverage: CoverageSpec + ) -> datetime.date: ... def last_modified( self, dataset_id: str, table_id: str ) -> datetime.datetime: ... @@ -292,8 +295,18 @@ def register_table_materialization( coverage_ids = client.get_coverage_ids(dataset_id, table_id) policy.assert_coverage_topology(coverage, coverage_ids) + # Só part_bdpro precisa do início da série (para o range free sair completo + # e não inverter); evita um scan extra nos demais tiers. + source_start = ( + bq.read_min_date(dataset_id, table_id, coverage) + if isinstance(coverage, PartBdpro) + else None + ) + # Cálculo puro dos ranges de cobertura. - ranges = policy.compute_coverage_ranges(coverage, source_end, coverage_ids) + ranges = policy.compute_coverage_ranges( + coverage, source_end, coverage_ids, source_start=source_start + ) for dtr in ranges.to_list(): client.upsert_coverage_datetime_range(dtr) diff --git a/pipelines/utils/metadata/utils.py b/pipelines/utils/metadata/utils.py index f1037e789..bbcbb4983 100644 --- a/pipelines/utils/metadata/utils.py +++ b/pipelines/utils/metadata/utils.py @@ -167,6 +167,59 @@ def extract_last_date_from_bq( raise +def extract_first_date_from_bq( + dataset_id: str, + table_id: str, + date_format: str, + date_column: dict, + billing_project_id: str, + project_id: str = "basedosdados", +) -> str: + """Extrai o início real da série (MIN da coluna de cobertura). + + Simétrica a `extract_last_date_from_bq`, mas para o começo da cobertura. É + usada só no cálculo do range free de `part_bdpro` (sempre histórico), então + não tem o ramo non-historical. + + Filtra valores anteriores a 1900 em qualquer granularidade — um ano < 1900 + não é sequer representável em `DateTimeRangeInput` (`Year >= 1900`) e + derrubaria a escrita — e, só para colunas de data (`{'date'}`), também + valores futuros, pelo mesmo motivo do MAX (um typo distorceria a cobertura). + Para colunas ano/ano-mês/ano-tri um rótulo futuro pode ser legítimo (ano + orçamentário, safra), então lá não se filtra o teto. + + Returns: + str: a primeira data no formato "%Y", "%Y-%m" ou "%Y-%m-%d". + """ + query_date_column = format_date_column(date_column) + + date_filter = f"\n WHERE {query_date_column} >= DATE '1900-01-01'" + if date_column.keys() == {"date"}: + date_filter += f"\n AND {query_date_column} <= CURRENT_DATE()" + + try: + query_bd = f""" + SELECT + MIN({query_date_column}) as min_date + FROM + `{project_id}.{dataset_id}.{table_id}`{date_filter} + """ + log(query_bd) + t = bd.read_sql( + query=query_bd, + billing_project_id=billing_project_id, + from_file=True, + ) + + first_date = t["min_date"][0].strftime(date_format) + log(f"Primeira data: {first_date}") + + return first_date + except Exception as e: + log(f"An error occurred while extracting the first date: {e!s}") + raise + + def format_date_column(date_column: dict) -> str: if date_column.keys() == {"date"}: query_date_column = date_column["date"] diff --git a/pipelines/utils/tests/metadata/conftest.py b/pipelines/utils/tests/metadata/conftest.py index cf9b3fd64..198443a4c 100644 --- a/pipelines/utils/tests/metadata/conftest.py +++ b/pipelines/utils/tests/metadata/conftest.py @@ -19,6 +19,8 @@ from __future__ import annotations +import datetime + import pytest @@ -180,8 +182,13 @@ def __init__( max_date=None, last_modified=None, can_read=True, + min_date=None, ): self._max_date = max_date + # Sem min_date explícito, usa max_date: assim um teste part_bdpro que só + # passa max_date ainda exercita o caminho de produção (source_start não + # nulo), e não o ramo legado source_start=None de compute_coverage_ranges. + self._min_date = max_date if min_date is None else min_date self._last_modified = last_modified self._can_read = can_read self.rap_calls: list[tuple] = [] @@ -189,6 +196,22 @@ def __init__( def read_max_date(self, dataset_id, table_id, coverage): return self._max_date + def read_min_date( + self, dataset_id: str, table_id: str, coverage + ) -> datetime.date | None: + """Devolve o início da série configurado (espelha `read_min_date` real). + + Args: + dataset_id: ID do dataset (ignorado pelo fake). + table_id: ID da tabela (ignorado pelo fake). + coverage: `CoverageSpec` da tabela (ignorado pelo fake). + + Returns: + A data mínima configurada — `min_date`, ou `max_date` quando aquele + não foi informado. + """ + return self._min_date + def last_modified(self, dataset_id, table_id): return self._last_modified diff --git a/pipelines/utils/tests/metadata/test_bq.py b/pipelines/utils/tests/metadata/test_bq.py index 0c266d82c..482193e00 100644 --- a/pipelines/utils/tests/metadata/test_bq.py +++ b/pipelines/utils/tests/metadata/test_bq.py @@ -45,6 +45,22 @@ def test_read_max_date_translates_and_parses(mock_extract): assert args[4] == "proj" # billing_project_id +@patch("pipelines.utils.metadata.bq.extract_first_date_from_bq") +def test_read_min_date_translates_and_parses(mock_extract) -> None: + """read_min_date traduz o DateColumn e faz o parse simétrico a read_max_date.""" + mock_extract.return_value = "2000-10-19" + bq = BigQueryReader(billing_project_id="proj", bq_project="basedosdados") + + out = bq.read_min_date("br_x", "tab", _daily()) + + assert out == datetime.date(2000, 10, 19) + # tradução: DateOnly → {"date":"data"} e formato "%Y-%m-%d" + args = mock_extract.call_args.args + assert args[2] == "%Y-%m-%d" + assert args[3] == {"date": "data"} + assert args[4] == "proj" # billing_project_id + + @patch("pipelines.utils.metadata.bq.update_date_from_bq_metadata") def test_last_modified_delegates(mock_lm): mock_lm.return_value = datetime.datetime(2026, 6, 2, 10, 0) diff --git a/pipelines/utils/tests/metadata/test_client.py b/pipelines/utils/tests/metadata/test_client.py index c9ebd8bee..ffd4a788a 100644 --- a/pipelines/utils/tests/metadata/test_client.py +++ b/pipelines/utils/tests/metadata/test_client.py @@ -128,6 +128,11 @@ def test_upsert_coverage_datetime_range_passes_dto_fields(client, backend): v = _input(backend.mutation_for("CreateUpdateDateTimeRange")) assert v["coverage"] == UUID assert v["endYear"] == 2026 and v["endMonth"] == 6 and v["endDay"] == 1 + # allDatetimerange devolveu None (range inexistente) → é um CREATE, e o + # backend exige `interval` para criar. Antes ele não saía no payload e o + # CREATE era impossível; agora sai sempre (default 1). + assert v["interval"] == 1 + assert "id" not in v # sem id ⇒ create, não update def test_write_carries_auth_header(client, backend): diff --git a/pipelines/utils/tests/metadata/test_dto.py b/pipelines/utils/tests/metadata/test_dto.py index 721dd7712..dfa511efa 100644 --- a/pipelines/utils/tests/metadata/test_dto.py +++ b/pipelines/utils/tests/metadata/test_dto.py @@ -58,6 +58,33 @@ def test_datetimerange_bad_coverage_uuid(): DateTimeRangeInput(coverage=BAD_UUID, endYear=2026) +def test_datetimerange_interval_defaults_to_one_and_is_emitted() -> None: + """O default 1 sai no payload — condição para o CREATE. + + O backend exige `interval` ao CRIAR um range com início e fim; sem ele o + upsert só conseguia UPDATE. Como não é None, sobrevive ao + `model_dump(exclude_none=True)`. + """ + dto = DateTimeRangeInput( + coverage=UUID, startYear=2026, startMonth=1, endYear=2026, endMonth=6 + ) + assert dto.interval == 1 + assert dto.model_dump(exclude_none=True)["interval"] == 1 + + +def test_datetimerange_interval_respects_explicit_value() -> None: + """Um interval explícito (ex.: 2, série bienal) é preservado, não fixado.""" + dto = DateTimeRangeInput(coverage=UUID, endYear=2026, interval=3) + assert dto.interval == 3 + + +@pytest.mark.parametrize("bad", [0, -1]) +def test_datetimerange_interval_must_be_positive(bad: int) -> None: + """interval é `ge=1`: zero ou negativo é rejeitado.""" + with pytest.raises(ValidationError): + DateTimeRangeInput(coverage=UUID, endYear=2026, interval=bad) + + # --- PollInput: FK, ISO 8601, defaults ---------------------------------------- def test_poll_input_valid(): dto = PollInput(rawDataSource=UUID, latest="2026-06-01", entity=UUID) diff --git a/pipelines/utils/tests/metadata/test_policy.py b/pipelines/utils/tests/metadata/test_policy.py index 33d4c2a87..7081b25da 100644 --- a/pipelines/utils/tests/metadata/test_policy.py +++ b/pipelines/utils/tests/metadata/test_policy.py @@ -187,6 +187,121 @@ def test_compute_part_bdpro_ranges_never_overlap(): assert (r.pro.startYear, r.pro.startMonth, r.pro.startDay) == (2026, 2, 1) +def test_compute_part_bdpro_free_end_only_without_source_start() -> None: + """Sem source_start (chamada legada), o free continua saindo só-com-fim. + + O comportamento antigo é preservado: o início gravado no backend é mantido + pelo update, que só toca o fim. + """ + spec = PartBdpro( + date_column=DateOnly(col="data"), + date_format=DateFormat.YEAR_MD, + free_lag=FreeLag(unit="months", value=6), + ) + r = compute_coverage_ranges(spec, date(2026, 8, 25), IDS_BOTH) + # pyrefly: ignore [missing-attribute] + assert (r.free.startYear, r.free.startMonth, r.free.startDay) == ( + None, + None, + None, + ) + # pyrefly: ignore [missing-attribute] + assert (r.free.endYear, r.free.endMonth, r.free.endDay) == (2026, 2, 25) + + +def test_compute_part_bdpro_free_full_range_when_history_exists() -> None: + """Com source_start bem antes de free_end, o free sai completo e válido. + + Início = início real da série, fim = source_end - lag; início <= fim, então + o range não inverte. + """ + spec = PartBdpro( + date_column=DateOnly(col="data"), + date_format=DateFormat.YEAR_MD, + free_lag=FreeLag(unit="months", value=6), + ) + r = compute_coverage_ranges( + spec, + date(2026, 8, 25), + IDS_BOTH, + source_start=date(2000, 10, 19), + ) + # pyrefly: ignore [missing-attribute] + free_start = (r.free.startYear, r.free.startMonth, r.free.startDay) + # pyrefly: ignore [missing-attribute] + free_end = (r.free.endYear, r.free.endMonth, r.free.endDay) + assert free_start == (2000, 10, 19) + assert free_end == (2026, 2, 25) + assert free_start <= free_end # início <= fim: não invertido + + +def test_compute_part_bdpro_free_collapses_when_history_all_paywalled() -> ( + None +): + """Snapshot sem história: o free colapsa em vez de inverter. + + Quando todo o dado é mais novo que free_end (source_start > free_end), o + free vira [free_end, free_end] — válido e degenerado — em vez de início > + fim, que o backend recusa. + """ + spec = PartBdpro( + date_column=DateOnly(col="data_extracao"), + date_format=DateFormat.YEAR_MD, + free_lag=FreeLag(unit="months", value=6), + ) + r = compute_coverage_ranges( + spec, + date(2026, 8, 25), + IDS_BOTH, + source_start=date(2026, 8, 25), # 1 único snapshot, hoje + ) + free_end_expected = (2026, 2, 25) # 2026-08-25 menos 6 meses + # pyrefly: ignore [missing-attribute] + free_start = (r.free.startYear, r.free.startMonth, r.free.startDay) + # pyrefly: ignore [missing-attribute] + free_end = (r.free.endYear, r.free.endMonth, r.free.endDay) + assert free_start == free_end_expected + assert free_end == free_end_expected + # start == end: range degenerado porém VÁLIDO (não invertido) + assert free_start <= free_end + # pro segue cobrindo até source_end, sem inverter + # pyrefly: ignore [missing-attribute] + assert (r.pro.endYear, r.pro.endMonth, r.pro.endDay) == (2026, 8, 25) + + +def test_compute_ranges_carry_spec_interval() -> None: + """O interval da spec vai para os ranges — não é fixado em 1. + + Eleições brasileiras são bienais (interval=2): tanto o range free (all_free) + quanto ambos os ranges de um part_bdpro precisam carregar esse passo. + """ + biennial_free = AllFree( + date_column=YearOnly(col="ano"), + date_format=DateFormat.YEAR, + interval=2, + ) + r = compute_coverage_ranges(biennial_free, date(2022, 1, 1), IDS_BOTH) + # pyrefly: ignore [missing-attribute] + assert r.free.interval == 2 + + biennial_part = PartBdpro( + date_column=YearOnly(col="ano"), + date_format=DateFormat.YEAR, + free_lag=FreeLag(unit="years", value=4), + interval=2, + ) + r = compute_coverage_ranges( + biennial_part, + date(2022, 1, 1), + IDS_BOTH, + source_start=date(1998, 1, 1), + ) + # pyrefly: ignore [missing-attribute] + assert r.free.interval == 2 + # pyrefly: ignore [missing-attribute] + assert r.pro.interval == 2 + + def test_compute_all_bdpro_annual(): spec = AllBdpro( date_column=YearOnly(col="ano"), date_format=DateFormat.YEAR diff --git a/pipelines/utils/tests/metadata/test_register.py b/pipelines/utils/tests/metadata/test_register.py index 8579ccbf9..5d1e6adfa 100644 --- a/pipelines/utils/tests/metadata/test_register.py +++ b/pipelines/utils/tests/metadata/test_register.py @@ -222,6 +222,41 @@ def test_part_bdpro_writes_coverages_table_update_and_rap(): assert len(bq.rap_calls) == 1 +def test_part_bdpro_history_less_writes_non_inverting_free_range() -> None: + """Regressão: um snapshot sem história (min == max) não inverte o free. + + Antes o range free invertia (início gravado > free_end) e o backend recusava + a escrita. Agora o orquestrador lê o início da série (read_min_date) e emite + um free completo e travado em free_end, com start <= end. + """ + client = FakeMetadataClient( + coverage_ids=CoverageIds( + free="ffffffff-ffff-4fff-8fff-ffffffffffff", + pro="aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa", + ) + ) + spec = PartBdpro( + date_column=DateOnly(col="data_extracao"), + date_format=DateFormat.YEAR_MD, + ) + bq = FakeBQ( + max_date=datetime.date(2026, 8, 25), + min_date=datetime.date(2026, 8, 25), # 1 único snapshot + last_modified=datetime.datetime(2026, 8, 26), + can_read=True, + ) + register_table_materialization(client, bq, "br_x", "tab", spec) + + # primeiro write "coverage" é o range free (ranges.to_list() = [free, pro]) + free_dto = next( + args[0] for entity, args, _ in client.writes if entity == "coverage" + ) + start = (free_dto.startYear, free_dto.startMonth, free_dto.startDay) + end = (free_dto.endYear, free_dto.endMonth, free_dto.endDay) + assert start <= end # não invertido + assert start == end == (2026, 2, 25) # colapsado em free_end + + def test_all_free_writes_only_free_coverage_and_table_update_no_rap(): client = FakeMetadataClient( coverage_ids=CoverageIds(