Pular para conteúdo

API reference

Every signature, parameter table and accepted value below is rendered from the docstrings of the version named in the header, so it changes with the code. In Python, help(sus.load) shows the same text for the version you have installed. Which years, ufs and months each dataset takes is in Bases e argumentos. To install, follow the installation guide.

For researchers

Four calls cover an analysis from download to citation:

import omnisus as sus

dados = sus.load("sim_obitos", years=[2023], ufs=["RR"])        # rows, codes as published
dados = sus.label("sim_obitos", dados, columns=["sexo", "racacor"])  # + sexo_rotulo, racacor_rotulo
sus.check_columns("sim_obitos", dados)                           # can I trust each column?
with sus.LakeReader() as lake:
    print(sus.cite(lake, dataset="sim_obitos").text)             # files, SHA-256, snapshot
  • load downloads what DATASUS publishes into the lake (where it lives) and returns the rows. A second call downloads nothing. When the files are validated sources it also adds harmonised categories (idade_anos_completos, sexo_categoria, <date>_data); otherwise it warns which were left out (see ADR 0003).
  • label puts the dictionary's label next to each code; a code the dictionary does not know gets no label.
  • check_columns reports, per column, empties, unlabelled codes and date ranges.

Download what DATASUS publishes for these scopes and return the rows.

The one-call path for a script or notebook cell. It imports into the lake (a second call downloads nothing new) and returns the requested scopes' rows with codes as DATASUS published them. When every file comes from a validated source (see describe_dataset(dataset)["analytics"]["validated_sources"]) it also adds the harmonised categories: idade_anos_completos, idade_status, sexo_categoria, sexo_status and <date>_data columns. Otherwise it leaves them out and warns once. Readable labels come from :func:label.

Parameters:

Name Type Description Default
dataset str | Dataset

DATASUS FTP dataset name, e.g. "sim_obitos", "sih_aih_reduzida"; see :func:datasets. IBGE population, SIGTAP and CNES names have their own functions (:func:import_ibge_populacao, :func:import_sigtap, :func:import_cnes_master).

required
years Iterable[int]

Years, e.g. [2023] or range(2020, 2024).

required
ufs Sequence[str] | None

UF abbreviations, e.g. ["RR"]. Required for datasets published per UF; pass :data:ALL_UFS for all of Brazil. National datasets (SINAN, SIM subsets) accept only None.

None
months Iterable[int] | None

Months 1-12 of monthly datasets (SIH, SIA, CNES), e.g. [1, 2]; None (default) means all twelve.

None
target str | None

DuckLake target; None (default) is data/raw/omnisus.ducklake under the working directory, or under $OMNISUS_DATA_DIR.

None
policy ImportPolicy

"skip_same" (default) reuses scopes already imported from the same file; "replace" rewrites scopes imported from another file, omnisus release or dictionary version; "append" and "error_if_exists" as in :func:import_dataset.

'skip_same'

Returns:

Type Description
DataFrame

A polars DataFrame with one row per record of the requested scopes.

Raises:

Type Description
ValueError

dataset is unknown or not an FTP dataset, ufs is missing for a UF dataset, or given for a national one.

RuntimeError

a scope failed to import; the message names each one.

LookupError

DATASUS publishes none of the requested scopes.

Warns:

Type Description
UserWarning

harmonised categories were left out, and why.

Examples:

>>> import omnisus as sus
>>> dados = sus.load("sim_obitos", years=[2023], ufs=["RR"])
>>> dados.height
3311
>>> dados = sus.label("sim_obitos", dados, columns=["sexo"])
Source code in src/omnisus/__init__.py
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
def load(
    dataset: str | Dataset,
    *,
    years: Iterable[int],
    ufs: Sequence[str] | None = None,
    months: Iterable[int] | None = None,
    target: str | None = None,
    policy: ImportPolicy = "skip_same",
) -> pl.DataFrame:
    """Download what DATASUS publishes for these scopes and return the rows.

    The one-call path for a script or notebook cell. It imports into the lake (a
    second call downloads nothing new) and returns the requested scopes' rows with
    codes as DATASUS published them. When every file comes from a validated source
    (see ``describe_dataset(dataset)["analytics"]["validated_sources"]``) it also adds
    the harmonised categories: ``idade_anos_completos``, ``idade_status``,
    ``sexo_categoria``, ``sexo_status`` and ``<date>_data`` columns. Otherwise it
    leaves them out and warns once. Readable labels come from :func:`label`.

    Args:
        dataset: DATASUS FTP dataset name, e.g. ``"sim_obitos"``,
            ``"sih_aih_reduzida"``; see :func:`datasets`. IBGE population, SIGTAP and
            CNES names have their own functions (:func:`import_ibge_populacao`,
            :func:`import_sigtap`, :func:`import_cnes_master`).
        years: Years, e.g. ``[2023]`` or ``range(2020, 2024)``.
        ufs: UF abbreviations, e.g. ``["RR"]``. Required for datasets published per
            UF; pass :data:`ALL_UFS` for all of Brazil. National datasets (SINAN,
            SIM subsets) accept only ``None``.
        months: Months 1-12 of monthly datasets (SIH, SIA, CNES), e.g. ``[1, 2]``;
            ``None`` (default) means all twelve.
        target: DuckLake target; ``None`` (default) is ``data/raw/omnisus.ducklake``
            under the working directory, or under ``$OMNISUS_DATA_DIR``.
        policy: ``"skip_same"`` (default) reuses scopes already imported from the
            same file; ``"replace"`` rewrites scopes imported from another file, omnisus
            release or dictionary version; ``"append"`` and ``"error_if_exists"`` as in
            :func:`import_dataset`.

    Returns:
        A polars DataFrame with one row per record of the requested scopes.

    Raises:
        ValueError: ``dataset`` is unknown or not an FTP dataset, ``ufs`` is missing
            for a UF dataset, or given for a national one.
        RuntimeError: a scope failed to import; the message names each one.
        LookupError: DATASUS publishes none of the requested scopes.

    Warns:
        UserWarning: harmonised categories were left out, and why.

    Examples:
        >>> import omnisus as sus
        >>> dados = sus.load("sim_obitos", years=[2023], ufs=["RR"])  # doctest: +SKIP
        >>> dados.height  # doctest: +SKIP
        3311
        >>> dados = sus.label("sim_obitos", dados, columns=["sexo"])  # doctest: +SKIP
    """
    d = _loadable(dataset)
    if ufs is None and d.geography != "national":
        raise ValueError(
            f"{d.name} is published per UF: pass ufs=['RR', ...], or ufs=sus.ALL_UFS "
            "for all of Brazil"
        )
    scopes = scopes_for(d, years=years, ufs=ufs, months=months)
    report = import_dataset(d, scopes=scopes, target=target, policy=policy)
    if report.failed:
        reasons = "; ".join(f"{o.scope}: {o.reason}" for o in report.failed)
        raise RuntimeError(
            f"{d.name} import failed ({reasons}). To rewrite a scope already in the "
            "lake, run this load again with policy='replace'."
        )
    rows = _read_with_harmonised(d, scopes, target)
    if rows.is_empty():
        raise LookupError(f"{d.name}: DATASUS publishes none of these scopes: {scopes}")
    return rows

Add the dictionary's label next to each coded column, keeping the code.

Each coded column c gains c_rotulo right after it. A code the dictionary does not know gets a null label, and so do a null and a blank the map does not declare: nothing is guessed, and the code stays in c.

Parameters:

Name Type Description Default
dataset str

Dataset name, e.g. "sim_obitos"; see sus.datasets().

required
df DataFrame

Rows of that dataset, e.g. from :func:omnisus.load.

required
columns Sequence[str] | None

Columns to label, e.g. ["sexo", "racacor"]. None (the default) labels every column of df that has a code map.

None

Returns:

Type Description
DataFrame

df with one <column>_rotulo (String) column after each labelled column.

Raises:

Type Description
ValueError

a name in columns is not in df or has no code map.

FileNotFoundError

dataset has no packaged dictionary.

Examples:

>>> import polars as pl
>>> import omnisus as sus
>>> dados = pl.DataFrame({"sexo": ["1", "2", " "], "idade": ["435", "450", "401"]})
>>> sus.label("sim_obitos", dados)
shape: (3, 3)
┌──────┬─────────────┬───────┐
│ sexo ┆ sexo_rotulo ┆ idade │
│ ---  ┆ ---         ┆ ---   │
│ str  ┆ str         ┆ str   │
╞══════╪═════════════╪═══════╡
│ 1    ┆ Masculino   ┆ 435   │
│ 2    ┆ Feminino    ┆ 450   │
│      ┆ null        ┆ 401   │
└──────┴─────────────┴───────┘
Source code in src/omnisus/transforms/dictionaries.py
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
def label(dataset: str, df: pl.DataFrame, *, columns: Sequence[str] | None = None) -> pl.DataFrame:
    """Add the dictionary's label next to each coded column, keeping the code.

    Each coded column ``c`` gains ``c_rotulo`` right after it. A code the
    dictionary does not know gets a null label, and so do a null and a blank the
    map does not declare: nothing is guessed, and the code stays in ``c``.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``; see ``sus.datasets()``.
        df: Rows of that dataset, e.g. from :func:`omnisus.load`.
        columns: Columns to label, e.g. ``["sexo", "racacor"]``. ``None`` (the
            default) labels every column of ``df`` that has a code map.

    Returns:
        ``df`` with one ``<column>_rotulo`` (``String``) column after each labelled column.

    Raises:
        ValueError: a name in ``columns`` is not in ``df`` or has no code map.
        FileNotFoundError: ``dataset`` has no packaged dictionary.

    Examples:
        >>> import polars as pl
        >>> import omnisus as sus
        >>> dados = pl.DataFrame({"sexo": ["1", "2", " "], "idade": ["435", "450", "401"]})
        >>> sus.label("sim_obitos", dados)
        shape: (3, 3)
        ┌──────┬─────────────┬───────┐
        │ sexo ┆ sexo_rotulo ┆ idade │
        │ ---  ┆ ---         ┆ ---   │
        │ str  ┆ str         ┆ str   │
        ╞══════╪═════════════╪═══════╡
        │ 1    ┆ Masculino   ┆ 435   │
        │ 2    ┆ Feminino    ┆ 450   │
        │      ┆ null        ┆ 401   │
        └──────┴─────────────┴───────┘
    """
    dicionario = load_dicionario(dataset)
    maps = {
        column: definition["x-decode"]
        for column in df.columns
        if (definition := dicionario.field_def(column)) and definition.get("x-decode")
    }
    if columns is not None:
        missing = [column for column in columns if column not in maps]
        if missing:
            raise ValueError(
                f"{dataset}: no code map for {', '.join(missing)} in this DataFrame; "
                f"columns with one: {', '.join(maps) or 'none'}"
            )
        maps = {column: maps[column] for column in columns}
    selection: list[pl.Expr] = []
    for column in df.columns:
        selection.append(pl.col(column))
        if column in maps:
            codes = df[column].drop_nulls().unique().to_list()
            labels = {code: _lookup_decode(maps[column], code) for code in codes}
            selection.append(
                pl.col(column)
                .replace_strict(labels, default=None, return_dtype=pl.String)
                .alias(f"{column}_rotulo")
            )
    return df.select(selection)

Report, per column, the dictionary's description and what the rows show.

Answers "can I trust this column?" before an analysis: how much is empty, which published codes have no label, and whether dates parse and fall in a plausible range (impossible dates such as 1899-12-30 parse, so look at date_min). A value counts as empty when it is null or an empty string.

Parameters:

Name Type Description Default
dataset str

Dataset name, e.g. "sim_obitos"; see sus.datasets().

required
df DataFrame

Rows of that dataset, e.g. from :func:omnisus.load.

required

Returns:

Type Description
DataFrame

One row per column of df, in order, with columns:

DataFrame

column; label (the dictionary's, None when undeclared);

DataFrame

rule ("code map (n codes)", "date ddMMyyyy", ``"not in

DataFrame

dictionary"orNone);pct_empty;distinct`` (filled values);

DataFrame

example_code and example_label (the first filled value and its

DataFrame

label); unlabelled_codes (list of filled codes the code map does not

DataFrame

label) and unlabelled_rows; pct_invalid_dates, date_min and

DataFrame

date_max (Date) for date fields, None otherwise.

Raises:

Type Description
FileNotFoundError

dataset has no packaged dictionary.

Examples:

SIH publishes homonimo 2, which no DATASUS table labels:

>>> import polars as pl
>>> import omnisus as sus
>>> dados = pl.DataFrame({"homonimo": ["0", "2", ""], "nasc": ["19850320", "18991230", ""]})
>>> relatorio = sus.check_columns("sih_aih_reduzida", dados)
>>> relatorio.select("column", "pct_empty", "unlabelled_codes", "date_min").rows()
[('homonimo', 33.3, ['2'], None), ('nasc', 33.3, [], datetime.date(1899, 12, 30))]
Source code in src/omnisus/transforms/columns.py
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
def check_columns(dataset: str, df: pl.DataFrame) -> pl.DataFrame:
    """Report, per column, the dictionary's description and what the rows show.

    Answers "can I trust this column?" before an analysis: how much is empty,
    which published codes have no label, and whether dates parse and fall in a
    plausible range (impossible dates such as ``1899-12-30`` parse, so look at
    ``date_min``). A value counts as empty when it is null or an empty string.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``; see ``sus.datasets()``.
        df: Rows of that dataset, e.g. from :func:`omnisus.load`.

    Returns:
        One row per column of ``df``, in order, with columns:
        ``column``; ``label`` (the dictionary's, ``None`` when undeclared);
        ``rule`` (``"code map (n codes)"``, ``"date ddMMyyyy"``, ``"not in
        dictionary"`` or ``None``); ``pct_empty``; ``distinct`` (filled values);
        ``example_code`` and ``example_label`` (the first filled value and its
        label); ``unlabelled_codes`` (list of filled codes the code map does not
        label) and ``unlabelled_rows``; ``pct_invalid_dates``, ``date_min`` and
        ``date_max`` (``Date``) for date fields, ``None`` otherwise.

    Raises:
        FileNotFoundError: ``dataset`` has no packaged dictionary.

    Examples:
        SIH publishes ``homonimo`` ``2``, which no DATASUS table labels:

        >>> import polars as pl
        >>> import omnisus as sus
        >>> dados = pl.DataFrame({"homonimo": ["0", "2", ""], "nasc": ["19850320", "18991230", ""]})
        >>> relatorio = sus.check_columns("sih_aih_reduzida", dados)
        >>> relatorio.select("column", "pct_empty", "unlabelled_codes", "date_min").rows()
        [('homonimo', 33.3, ['2'], None), ('nasc', 33.3, [], datetime.date(1899, 12, 30))]
    """
    dicionario = load_dicionario(dataset)
    rows = []
    for column in df.columns:
        field = dicionario.field_def(column)
        text = df[column].cast(pl.String)
        filled = text.filter(text.is_not_null() & (text != ""))
        example = filled[0] if filled.len() else None
        decode_map = (field or {}).get("x-decode")
        unlabelled: list[str] = []
        unlabelled_rows = 0
        if decode_map:
            for code, count in filled.value_counts(sort=True).rows():
                if _lookup_decode(decode_map, code) is None:
                    unlabelled.append(code)
                    unlabelled_rows += count
        invalid = date_min = date_max = None
        pattern = DATE_FORMATS.get(field.get("x-format", "")) if field else None
        if field and field.get("type") == "date" and pattern:
            dates = filled.str.strptime(pl.Date, pattern, strict=False)
            invalid = _percent(dates.null_count(), filled.len())
            date_min, date_max = dates.min(), dates.max()
        rows.append(
            {
                "column": column,
                "label": field.get("label") if field else None,
                "rule": _rule(field),
                "pct_empty": _percent(df.height - filled.len(), df.height),
                "distinct": filled.n_unique(),
                "example_code": example,
                "example_label": _lookup_decode(decode_map, example) if decode_map else None,
                "unlabelled_codes": unlabelled,
                "unlabelled_rows": unlabelled_rows,
                "pct_invalid_dates": invalid,
                "date_min": date_min,
                "date_max": date_max,
            }
        )
    return pl.DataFrame(rows, schema=_SCHEMA)

The Portuguese citation of what a lake holds: files, hashes, snapshot and run.

Parameters:

Name Type Description Default
lake Lake | LakeReader

An open :class:~omnisus.Lake or :class:~omnisus.LakeReader.

required
dataset str | None

Cite only this dataset, e.g. "sim_obitos"; None (default) cites every active publication.

None
snapshot_id int | None

The snapshot your analysis read; None (default) is the newest.

None
run_id str | None

Cite only this run's publications.

None
accessed date | None

Access date to print; None (default) is today.

None

Returns:

Name Type Description
A Citation

class:~omnisus.Citation; paste citation.text into the methods.

Raises:

Type Description
LookupError

snapshot_id is None and the lake has no snapshots.

Examples:

>>> import omnisus as sus
>>> with sus.LakeReader() as lake:
...     print(sus.cite(lake, dataset="sim_obitos").text)
Source code in src/omnisus/research.py
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
def cite(
    lake: Lake | LakeReader,
    *,
    dataset: str | None = None,
    snapshot_id: int | None = None,
    run_id: str | None = None,
    accessed: date | None = None,
) -> Citation:
    """The Portuguese citation of what a lake holds: files, hashes, snapshot and run.

    Args:
        lake: An open :class:`~omnisus.Lake` or :class:`~omnisus.LakeReader`.
        dataset: Cite only this dataset, e.g. ``"sim_obitos"``; ``None`` (default)
            cites every active publication.
        snapshot_id: The snapshot your analysis read; ``None`` (default) is the
            newest.
        run_id: Cite only this run's publications.
        accessed: Access date to print; ``None`` (default) is today.

    Returns:
        A :class:`~omnisus.Citation`; paste ``citation.text`` into the methods.

    Raises:
        LookupError: ``snapshot_id`` is ``None`` and the lake has no snapshots.

    Examples:
        >>> import omnisus as sus
        >>> with sus.LakeReader() as lake:  # doctest: +SKIP
        ...     print(sus.cite(lake, dataset="sim_obitos").text)
    """
    pinned = latest_snapshot_id(lake) if snapshot_id is None else snapshot_id
    if dataset == "ibge_populacao":
        rows = _ibge_manifest(lake)
    else:
        rows = [
            row
            for row in lake.publications(run_id=run_id)
            if row.get("active", True) and (dataset is None or row.get("dataset") == dataset)
        ]
    return citation_from_publications(
        rows,
        snapshot_id=pinned,
        dataset=dataset,
        run_id=run_id,
        accessed=accessed,
    )

Discovery

Ask the server what exists before deciding what to import.

The scopes DATASUS publishes now: what an import can actually get.

Pass the result to :func:~omnisus.import_dataset; unlike :func:~omnisus.scopes_for, it plans only files the server lists. Files of other datasets sharing a directory are skipped.

Parameters:

Name Type Description Default
dataset str | Dataset

Dataset name, e.g. "sim_obitos", or a :class:Dataset.

required
years Iterable[int] | None

Keep these years, e.g. [2023, 2024]; None (default) keeps all.

None
ufs Sequence[str] | None

Keep these UFs, e.g. ["RR"]; None keeps all. National datasets accept only None.

None
months Iterable[int] | None

Keep these months 1-12; None keeps all. National datasets accept only None.

None
refresh bool

True ignores the local listing cache and asks the server again.

False

Returns:

Type Description
list[ScopeKey]

The scopes, ordered by year, UF and month.

Raises:

Type Description
ValueError

as :func:available_releases.

FtpUnavailable

the server did not answer after retries.

Examples:

>>> import omnisus as sus
>>> sus.available("sim_obitos", years=[2023], ufs=["RR"])
[ScopeKey(uf='RR', ano=2023, mes=None)]
Source code in src/omnisus/sources/datasus_ftp/inventory.py
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
def available(
    dataset: str | Dataset,
    *,
    years: Iterable[int] | None = None,
    ufs: Sequence[str] | None = None,
    months: Iterable[int] | None = None,
    refresh: bool = False,
) -> list[ScopeKey]:
    """The scopes DATASUS publishes now: what an import can actually get.

    Pass the result to :func:`~omnisus.import_dataset`; unlike
    :func:`~omnisus.scopes_for`, it plans only files the server lists. Files of
    other datasets sharing a directory are skipped.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``, or a :class:`Dataset`.
        years: Keep these years, e.g. ``[2023, 2024]``; ``None`` (default) keeps all.
        ufs: Keep these UFs, e.g. ``["RR"]``; ``None`` keeps all. National datasets
            accept only ``None``.
        months: Keep these months 1-12; ``None`` keeps all. National datasets accept
            only ``None``.
        refresh: ``True`` ignores the local listing cache and asks the server again.

    Returns:
        The scopes, ordered by year, UF and month.

    Raises:
        ValueError: as :func:`available_releases`.
        FtpUnavailable: the server did not answer after retries.

    Examples:
        >>> import omnisus as sus
        >>> sus.available("sim_obitos", years=[2023], ufs=["RR"])  # doctest: +SKIP
        [ScopeKey(uf='RR', ano=2023, mes=None)]
    """
    return list(available_releases(dataset, years=years, ufs=ufs, months=months, refresh=refresh))

The scopes DATASUS publishes now, and whether each is final or preliminary.

One cached listing per directory. A scope found in two directories is a server inconsistency and raises; it is never resolved by preference.

Parameters:

Name Type Description Default
dataset str | Dataset

Dataset name, e.g. "sim_obitos", or a :class:Dataset.

required
years Iterable[int] | None

Keep these years, e.g. [2023, 2024]; None (default) keeps all.

None
ufs Sequence[str] | None

Keep these UFs, e.g. ["RR"]; None keeps all. National datasets accept only None.

None
months Iterable[int] | None

Keep these months 1-12; None keeps all. National datasets accept only None.

None
refresh bool

True ignores the local listing cache and asks the server again.

False

Returns:

Type Description
dict[ScopeKey, Release]

{scope: "final" | "prelim"}, ordered by year, UF and month.

Raises:

Type Description
ValueError

dataset is unknown, a national dataset got ufs/months, or a scope is listed in two directories.

FtpUnavailable

the server did not answer after retries.

Examples:

>>> import omnisus as sus
>>> sus.available_releases("sim_obitos", years=[2024], ufs=["RR"])
{ScopeKey(uf='RR', ano=2024, mes=None): 'prelim'}
Source code in src/omnisus/sources/datasus_ftp/inventory.py
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
def available_releases(
    dataset: str | Dataset,
    *,
    years: Iterable[int] | None = None,
    ufs: Sequence[str] | None = None,
    months: Iterable[int] | None = None,
    refresh: bool = False,
) -> dict[ScopeKey, Release]:
    """The scopes DATASUS publishes now, and whether each is final or preliminary.

    One cached listing per directory. A scope found in two directories is a server
    inconsistency and raises; it is never resolved by preference.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``, or a :class:`Dataset`.
        years: Keep these years, e.g. ``[2023, 2024]``; ``None`` (default) keeps all.
        ufs: Keep these UFs, e.g. ``["RR"]``; ``None`` keeps all. National datasets
            accept only ``None``.
        months: Keep these months 1-12; ``None`` keeps all. National datasets accept
            only ``None``.
        refresh: ``True`` ignores the local listing cache and asks the server again.

    Returns:
        ``{scope: "final" | "prelim"}``, ordered by year, UF and month.

    Raises:
        ValueError: ``dataset`` is unknown, a national dataset got ``ufs``/``months``,
            or a scope is listed in two directories.
        FtpUnavailable: the server did not answer after retries.

    Examples:
        >>> import omnisus as sus
        >>> sus.available_releases("sim_obitos", years=[2024], ufs=["RR"])  # doctest: +SKIP
        {ScopeKey(uf='RR', ano=2024, mes=None): 'prelim'}
    """
    d = resolve(dataset)
    if d.geography == "national" and (ufs is not None or months is not None):
        raise ValueError(f"{d.name} is national: it has no UF or month to filter by")
    wanted = set(years) if years is not None else None
    wanted_ufs = {u.upper() for u in ufs} if ufs is not None else None
    wanted_months = set(months) if months is not None else None
    found = {
        scope: source.release
        for scope, source in list_sources(d, refresh=refresh).items()
        if (wanted is None or scope.ano in wanted)
        and (wanted_ufs is None or scope.uf in wanted_ufs)
        and (wanted_months is None or scope.mes in wanted_months)
    }
    logger.info("inventory.available", dataset=d.name, scopes=len(found))
    return found

List any DATASUS FTP path, including datasets the package does not curate.

The open-world counterpart to :func:available: it reaches other SINAN agravos, CIHA or PCE. Eager, so a huge tree is listed whole; call omnisus.sources.datasus_ftp.inventory.crawl for a lazy walk.

Parameters:

Name Type Description Default
path str

Absolute FTP path, e.g. "/dissemin/publicos/SINAN/DADOS/FINAIS".

required
depth int

How many directory levels to list; 1 (default) lists path only.

1
refresh bool

True ignores the local listing cache and asks the server again.

False

Returns:

Name Type Description
One list[FtpEntry]

class:FtpEntry per file or directory, with size and server time.

Raises:

Type Description
FtpPathNotFound

the server has no such path.

FtpUnavailable

the server did not answer after retries.

Examples:

>>> import omnisus as sus
>>> entries = sus.browse("/dissemin/publicos/SINAN/DADOS/FINAIS")
>>> entries[0].name
'ACBIBR06.dbc'
Source code in src/omnisus/__init__.py
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
def browse(path: str, *, depth: int = 1, refresh: bool = False) -> list[FtpEntry]:
    """List any DATASUS FTP path, including datasets the package does not curate.

    The open-world counterpart to :func:`available`: it reaches other SINAN agravos,
    CIHA or PCE. Eager, so a huge tree is listed whole; call
    ``omnisus.sources.datasus_ftp.inventory.crawl`` for a lazy walk.

    Args:
        path: Absolute FTP path, e.g. ``"/dissemin/publicos/SINAN/DADOS/FINAIS"``.
        depth: How many directory levels to list; ``1`` (default) lists ``path`` only.
        refresh: ``True`` ignores the local listing cache and asks the server again.

    Returns:
        One :class:`FtpEntry` per file or directory, with size and server time.

    Raises:
        FtpPathNotFound: the server has no such path.
        FtpUnavailable: the server did not answer after retries.

    Examples:
        >>> import omnisus as sus
        >>> entries = sus.browse("/dissemin/publicos/SINAN/DADOS/FINAIS")  # doctest: +SKIP
        >>> entries[0].name  # doctest: +SKIP
        'ACBIBR06.dbc'
    """
    return list(_crawl(path, depth=depth, refresh=refresh))

One line of a DATASUS FTP directory listing.

Source code in src/omnisus/sources/datasus_ftp/inventory.py
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
@dataclass(frozen=True)
class FtpEntry:
    """One line of a DATASUS FTP directory listing."""

    name: str
    """File or directory name, spaces preserved."""

    path: str
    """Absolute remote path."""

    parent: str
    """Absolute path of the containing directory."""

    is_dir: bool

    size_bytes: int
    """0 for directories. Exceeds 32 bits in the wild."""

    modified: datetime
    """Server-reported mtime. Detects DATASUS republishing a file we ingested."""

modified instance-attribute

Server-reported mtime. Detects DATASUS republishing a file we ingested.

name instance-attribute

File or directory name, spaces preserved.

parent instance-attribute

Absolute path of the containing directory.

path instance-attribute

Absolute remote path.

size_bytes instance-attribute

0 for directories. Exceeds 32 bits in the wild.

Releases

Some datasets publish final and preliminary files under the same names in two directories (a row's prelim_dir). available_releases reports which directory each scope came from; outdated compares the files a lake has published against what the server lists today — by path, and by size or server mtime when recorded — and returns the scopes that differ (a moved directory, or a same-name republish), to be re-imported with import_dataset(..., policy="replace"). See reprocessing and maintenance.

Scopes whose imported files differ from what the server lists now.

Different means another path (a preliminary year that became final, or a moved directory), or the same path with another size or server time (DATASUS republished it). Publications without a recorded source are always reported. Scopes the server no longer lists are a withdrawal, not returned here. Read-only.

Parameters:

Name Type Description Default
dataset str | Dataset

Dataset name, e.g. "sim_obitos", or a :class:Dataset.

required
lake Lake | LakeReader

An open :class:Lake or :class:LakeReader.

required
refresh bool

True ignores the local listing cache and asks the server again.

False

Returns:

Type Description
list[ScopeKey]

The scopes to re-import with :func:import_dataset and

list[ScopeKey]

policy="replace", ordered by year, UF and month.

Raises:

Type Description
ValueError

dataset is unknown.

FtpUnavailable

the server did not answer after retries.

Examples:

>>> import omnisus as sus
>>> with sus.LakeReader() as lake:
...     stale = sus.outdated("sim_obitos", lake=lake)
>>> sus.import_dataset(
...     "sim_obitos", scopes=stale, policy="replace", run_id="refresh-2026-09"
... )
Source code in src/omnisus/__init__.py
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
def outdated(
    dataset: str | Dataset,
    *,
    lake: Lake | LakeReader,
    refresh: bool = False,
) -> list[ScopeKey]:
    """Scopes whose imported files differ from what the server lists now.

    Different means another path (a preliminary year that became final, or a moved
    directory), or the same path with another size or server time (DATASUS
    republished it). Publications without a recorded source are always reported.
    Scopes the server no longer lists are a withdrawal, not returned here.
    Read-only.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``, or a :class:`Dataset`.
        lake: An open :class:`Lake` or :class:`LakeReader`.
        refresh: ``True`` ignores the local listing cache and asks the server again.

    Returns:
        The scopes to re-import with :func:`import_dataset` and
        ``policy="replace"``, ordered by year, UF and month.

    Raises:
        ValueError: ``dataset`` is unknown.
        FtpUnavailable: the server did not answer after retries.

    Examples:
        >>> import omnisus as sus
        >>> with sus.LakeReader() as lake:  # doctest: +SKIP
        ...     stale = sus.outdated("sim_obitos", lake=lake)
        >>> sus.import_dataset(  # doctest: +SKIP
        ...     "sim_obitos", scopes=stale, policy="replace", run_id="refresh-2026-09"
        ... )
    """
    d = resolve(dataset)
    current = list_sources(d, refresh=refresh)
    changed = {
        row["scope"]
        for row in lake.publications()
        if row["dataset"] == d.name
        and row["active"]
        and row["scope"] in current
        and _files_differ(row["sources"], current[row["scope"]].files)
    }
    return sorted(changed, key=lambda s: (s.ano, s.uf or "", s.mes or 0))

Planning

Planning is composition: build a list of scopes any way you like and hand it to import_dataset. There is no planner flag on the Python API.

Every scope (UF x year [x month]) of a dataset, without asking the server.

Planning is composition: pass the result, or any other list of scopes, to :func:import_dataset. Scopes DATASUS does not publish are later reported as skipped; :func:available plans only what the server lists.

Parameters:

Name Type Description Default
dataset str | Dataset

Dataset name, e.g. "sim_obitos", or a :class:Dataset; see :func:datasets.

required
years Iterable[int]

Years, e.g. [2022, 2023] or range(2020, 2024).

required
ufs Sequence[str] | None

UF abbreviations, e.g. ["RR", "AC"]. None (default) means all 27 (:data:ALL_UFS). National datasets (SINAN, SIM subsets) accept only None.

None
months Iterable[int] | None

Months 1-12 for monthly datasets (SIH, SIA, CNES); None (default) means all twelve. Ignored for yearly UF datasets; national datasets accept only None.

None

Returns:

Type Description
list[ScopeKey]

The scopes, ordered by year, then UF, then month.

Raises:

Type Description
ValueError

dataset is unknown, or a national dataset got ufs or months.

Examples:

>>> import omnisus as sus
>>> sus.scopes_for("sih_aih_reduzida", years=[2024], ufs=["RR"], months=[1, 2])
[ScopeKey(uf='RR', ano=2024, mes=1), ScopeKey(uf='RR', ano=2024, mes=2)]
>>> sus.scopes_for("sinan_chagas", years=[2023])
[ScopeKey(uf=None, ano=2023, mes=None)]
Source code in src/omnisus/__init__.py
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
def scopes_for(
    dataset: str | Dataset,
    *,
    years: Iterable[int],
    ufs: Sequence[str] | None = None,
    months: Iterable[int] | None = None,
) -> list[ScopeKey]:
    """Every scope (UF x year [x month]) of a dataset, without asking the server.

    Planning is composition: pass the result, or any other list of scopes, to
    :func:`import_dataset`. Scopes DATASUS does not publish are later reported as
    ``skipped``; :func:`available` plans only what the server lists.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``, or a :class:`Dataset`; see
            :func:`datasets`.
        years: Years, e.g. ``[2022, 2023]`` or ``range(2020, 2024)``.
        ufs: UF abbreviations, e.g. ``["RR", "AC"]``. ``None`` (default) means all 27
            (:data:`ALL_UFS`). National datasets (SINAN, SIM subsets) accept only ``None``.
        months: Months 1-12 for monthly datasets (SIH, SIA, CNES); ``None`` (default)
            means all twelve. Ignored for yearly UF datasets; national datasets accept
            only ``None``.

    Returns:
        The scopes, ordered by year, then UF, then month.

    Raises:
        ValueError: ``dataset`` is unknown, or a national dataset got ``ufs`` or
            ``months``.

    Examples:
        >>> import omnisus as sus
        >>> sus.scopes_for("sih_aih_reduzida", years=[2024], ufs=["RR"], months=[1, 2])
        [ScopeKey(uf='RR', ano=2024, mes=1), ScopeKey(uf='RR', ano=2024, mes=2)]
        >>> sus.scopes_for("sinan_chagas", years=[2023])
        [ScopeKey(uf=None, ano=2023, mes=None)]
    """
    d = resolve(dataset)
    if d.geography == "national":
        if ufs is not None:
            raise ValueError(
                "national datasets do not accept UF filters; filter records after import"
            )
        if months is not None:
            raise ValueError("national yearly datasets do not accept month filters")
        return [ScopeKey(uf=None, ano=year) for year in years]
    uf_list = tuple(ufs) if ufs is not None else ALL_UFS
    month_list = tuple(months) if months is not None else tuple(range(1, 13))
    scopes: list[ScopeKey] = []
    for year in years:
        for uf in uf_list:
            if d.monthly:
                scopes.extend(ScopeKey(uf=uf, ano=year, mes=m) for m in month_list)
            else:
                scopes.append(ScopeKey(uf=uf, ano=year))
    return scopes

The 27 federative units. A two-letter code outside this set is not a state.

Importing

The FTP importers return ImportReport; inspect report.failed, report.skipped and report.ok. import_ibge_populacao returns list[ImportResult], while import_cnes_master returns the number of records written.

import_research is the researcher door: it requires run_id and defaults to policy="skip_same". It refuses append. import_dataset remains the operator API and still appends unless a policy is set.

ImportAbortedError interrupts an FTP run when it cannot safely continue. Inspect its report for determined outcomes and unresolved for (input_index, ScopeKey) pairs before retrying. Imports append data unless an explicit replay policy is selected. FTP imports accept append (default), skip_same, error_if_exists and replace. See reprocessing and maintenance for legacy scope restrictions, run IDs and byte budgets.

Every import function also runs from inside a notebook cell, where an event loop is already running. Importing cnes_estabelecimentos by any function refreshes the aux_cnes view.

Import the given scopes of a DATASUS FTP dataset into the lake (operator door).

The import iterates scopes and never asks where they came from: :func:scopes_for plans blindly and :func:available asks the server first. Every run lists the server once and imports what it lists; callers never choose a release or file name. Importing cnes_estabelecimentos also refreshes the aux_cnes view. Works inside a notebook's running event loop.

Parameters:

Name Type Description Default
dataset str | Dataset

Dataset name, e.g. "sia_bpa_individualizado", or a :class:Dataset; see :func:datasets.

required
scopes Sequence[ScopeKey]

The scopes to import, e.g. from :func:scopes_for or :func:available.

required
target str | None

DuckLake target such as "ducklake:/path/omnisus.ducklake" or "ducklake:postgresql://…?storage=s3://…". None (default) is data/raw/omnisus.ducklake under the working directory, or under $OMNISUS_DATA_DIR.

None
concurrency int

Downloads in flight. DATASUS FTP is a shared public server; raise it only with reason.

DEFAULT_CONCURRENCY
batch_size int

Scopes committed in one DuckLake transaction, hence one snapshot.

DEFAULT_BATCH_SIZE
policy ImportPolicy

What to do with a scope already in the lake: "append" (default) adds rows again; "skip_same" skips it when file and parser are unchanged, without downloading when the listing shows the same path, size and server time; "error_if_exists" fails that scope; "replace" rewrites it. :func:load and :func:import_research default to "skip_same".

'append'
run_id str | None

Your identifier for this run, chosen before it starts, to reconcile an interrupted commit through Lake.publications(run_id=...). None generates one.

None
max_payload_bytes int

Largest compressed download accepted per file (512 MiB).

DEFAULT_MAX_PAYLOAD_BYTES
max_inflight_bytes int

Compressed bytes held at once before parsing (1 GiB). Neither bounds total memory: a decompressed DBF is held whole.

DEFAULT_MAX_INFLIGHT_BYTES

Returns:

Name Type Description
One ImportReport

class:ScopeOutcome per scope: ok, skipped (e.g. the server

ImportReport

does not list it) or failed, with code and reason. Inspect

ImportReport

report.failed; later batches continue after a failed scope.

Raises:

Type Description
ImportAbortedError

the lake transaction failed mid-run; its report and unresolved say what was committed. Inspect before retrying.

ValueError

dataset is unknown or policy is not one of the four values.

Examples:

>>> import omnisus as sus
>>> scopes = sus.available("sim_obitos", years=[2023], ufs=["RR"])
>>> report = sus.import_dataset("sim_obitos", scopes=scopes)
>>> [(str(o.scope), o.status) for o in report.outcomes]
[('RR_2023', 'ok')]
Source code in src/omnisus/__init__.py
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
def import_dataset(
    dataset: str | Dataset,
    *,
    scopes: Sequence[ScopeKey],
    target: str | None = None,
    concurrency: int = DEFAULT_CONCURRENCY,
    batch_size: int = DEFAULT_BATCH_SIZE,
    policy: ImportPolicy = "append",
    run_id: str | None = None,
    max_payload_bytes: int = DEFAULT_MAX_PAYLOAD_BYTES,
    max_inflight_bytes: int = DEFAULT_MAX_INFLIGHT_BYTES,
) -> ImportReport:
    """Import the given scopes of a DATASUS FTP dataset into the lake (operator door).

    The import iterates ``scopes`` and never asks where they came from:
    :func:`scopes_for` plans blindly and :func:`available` asks the server first.
    Every run lists the server once and imports what it lists; callers never choose
    a release or file name. Importing ``cnes_estabelecimentos`` also refreshes the
    ``aux_cnes`` view. Works inside a notebook's running event loop.

    Args:
        dataset: Dataset name, e.g. ``"sia_bpa_individualizado"``, or a
            :class:`Dataset`; see :func:`datasets`.
        scopes: The scopes to import, e.g. from :func:`scopes_for` or :func:`available`.
        target: DuckLake target such as ``"ducklake:/path/omnisus.ducklake"`` or
            ``"ducklake:postgresql://…?storage=s3://…"``. ``None`` (default) is
            ``data/raw/omnisus.ducklake`` under the working directory, or under
            ``$OMNISUS_DATA_DIR``.
        concurrency: Downloads in flight. DATASUS FTP is a shared public server; raise
            it only with reason.
        batch_size: Scopes committed in one DuckLake transaction, hence one snapshot.
        policy: What to do with a scope already in the lake: ``"append"`` (default)
            adds rows again; ``"skip_same"`` skips it when file and parser are
            unchanged, without downloading when the listing shows the same path,
            size and server time; ``"error_if_exists"`` fails that scope;
            ``"replace"`` rewrites it. :func:`load` and :func:`import_research` default to ``"skip_same"``.
        run_id: Your identifier for this run, chosen before it starts, to reconcile an
            interrupted commit through ``Lake.publications(run_id=...)``. ``None``
            generates one.
        max_payload_bytes: Largest compressed download accepted per file (512 MiB).
        max_inflight_bytes: Compressed bytes held at once before parsing (1 GiB).
            Neither bounds total memory: a decompressed DBF is held whole.

    Returns:
        One :class:`ScopeOutcome` per scope: ``ok``, ``skipped`` (e.g. the server
        does not list it) or ``failed``, with ``code`` and ``reason``. Inspect
        ``report.failed``; later batches continue after a failed scope.

    Raises:
        ImportAbortedError: the lake transaction failed mid-run; its ``report`` and
            ``unresolved`` say what was committed. Inspect before retrying.
        ValueError: ``dataset`` is unknown or ``policy`` is not one of the four values.

    Examples:
        >>> import omnisus as sus
        >>> scopes = sus.available("sim_obitos", years=[2023], ufs=["RR"])  # doctest: +SKIP
        >>> report = sus.import_dataset("sim_obitos", scopes=scopes)  # doctest: +SKIP
        >>> [(str(o.scope), o.status) for o in report.outcomes]  # doctest: +SKIP
        [('RR_2023', 'ok')]
    """
    d = resolve(dataset)

    async def run() -> ImportReport:
        with Lake.local(target) as lake:
            report = await _run_scopes_ftp(
                d,
                scopes=scopes,
                lake=lake,
                concurrency=concurrency,
                batch_size=batch_size,
                policy=policy,
                run_id=run_id,
                max_payload_bytes=max_payload_bytes,
                max_inflight_bytes=max_inflight_bytes,
            )
            if after_import := _AFTER_IMPORT.get(d.name):
                after_import(lake)
            return report

    return run_sync(run)

Import for a citable run: you choose the run_id before anything downloads.

Same as :func:~omnisus.import_dataset, but run_id is required, so an interrupted run is reconciled through Lake.publications(run_id=...), and the default policy never duplicates rows on a retry.

Parameters:

Name Type Description Default
dataset str | Dataset

Dataset name, e.g. "sim_obitos", or a :class:~omnisus.Dataset.

required
scopes Sequence[ScopeKey]

The scopes to import, e.g. from :func:~omnisus.scopes_for.

required
run_id str

Your non-empty identifier for this run, e.g. "tese-cap2-2026-09-23".

required
target str | None

DuckLake target; None (default) is data/raw/omnisus.ducklake under the working directory, or under $OMNISUS_DATA_DIR.

None
concurrency int

Downloads in flight; see :func:~omnisus.import_dataset.

DEFAULT_CONCURRENCY
batch_size int

Scopes committed per DuckLake transaction.

DEFAULT_BATCH_SIZE
policy ImportPolicy

"skip_same" (default), "error_if_exists" or "replace". "append" is refused: it duplicates rows on retry.

'skip_same'
max_payload_bytes int

Largest compressed download accepted per file (512 MiB).

DEFAULT_MAX_PAYLOAD_BYTES
max_inflight_bytes int

Compressed bytes held at once before parsing (1 GiB).

DEFAULT_MAX_INFLIGHT_BYTES

Returns:

Name Type Description
The ImportReport

class:~omnisus.ImportReport, as :func:~omnisus.import_dataset.

Raises:

Type Description
ValueError

run_id is empty or policy is "append".

ImportAbortedError

as :func:~omnisus.import_dataset.

Examples:

>>> import omnisus as sus
>>> scopes = sus.scopes_for("sim_obitos", years=[2023], ufs=["RR"])
>>> sus.import_research("sim_obitos", scopes=scopes, run_id="cap2")
Source code in src/omnisus/research.py
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
def import_research(
    dataset: str | Dataset,
    *,
    scopes: Sequence[ScopeKey],
    run_id: str,
    target: str | None = None,
    concurrency: int = DEFAULT_CONCURRENCY,
    batch_size: int = DEFAULT_BATCH_SIZE,
    policy: ImportPolicy = "skip_same",
    max_payload_bytes: int = DEFAULT_MAX_PAYLOAD_BYTES,
    max_inflight_bytes: int = DEFAULT_MAX_INFLIGHT_BYTES,
) -> ImportReport:
    """Import for a citable run: you choose the ``run_id`` before anything downloads.

    Same as :func:`~omnisus.import_dataset`, but ``run_id`` is required, so an
    interrupted run is reconciled through ``Lake.publications(run_id=...)``, and the
    default policy never duplicates rows on a retry.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``, or a :class:`~omnisus.Dataset`.
        scopes: The scopes to import, e.g. from :func:`~omnisus.scopes_for`.
        run_id: Your non-empty identifier for this run, e.g. ``"tese-cap2-2026-09-23"``.
        target: DuckLake target; ``None`` (default) is ``data/raw/omnisus.ducklake``
            under the working directory, or under ``$OMNISUS_DATA_DIR``.
        concurrency: Downloads in flight; see :func:`~omnisus.import_dataset`.
        batch_size: Scopes committed per DuckLake transaction.
        policy: ``"skip_same"`` (default), ``"error_if_exists"`` or ``"replace"``.
            ``"append"`` is refused: it duplicates rows on retry.
        max_payload_bytes: Largest compressed download accepted per file (512 MiB).
        max_inflight_bytes: Compressed bytes held at once before parsing (1 GiB).

    Returns:
        The :class:`~omnisus.ImportReport`, as :func:`~omnisus.import_dataset`.

    Raises:
        ValueError: ``run_id`` is empty or ``policy`` is ``"append"``.
        ImportAbortedError: as :func:`~omnisus.import_dataset`.

    Examples:
        >>> import omnisus as sus
        >>> scopes = sus.scopes_for("sim_obitos", years=[2023], ufs=["RR"])
        >>> sus.import_research("sim_obitos", scopes=scopes, run_id="cap2")  # doctest: +SKIP
    """
    from omnisus import import_dataset

    if not str(run_id).strip():
        raise ValueError("import_research requires a nonempty run_id")
    if policy == "append":
        raise ValueError(
            "import_research refuses policy='append'; use import_dataset or policy='replace'"
        )
    return import_dataset(
        dataset,
        scopes=scopes,
        target=target,
        concurrency=concurrency,
        batch_size=batch_size,
        policy=policy,
        run_id=run_id,
        max_payload_bytes=max_payload_bytes,
        max_inflight_bytes=max_inflight_bytes,
    )

Import IBGE population by municipality: a census or the year's latest estimate.

Each population edition is its own publication, identified by the returned publication_id in ibge_population_manifest; it never appears in Lake.publications(). Reconcile an interrupted run through that manifest.

Parameters:

Name Type Description Default
years Iterable[int]

Years, e.g. [2022]. Censuses exist for 2010 and 2022; an estimate is importable only as the aggregate's latest edition, and years without a verified territorial universe are refused.

required
census bool

Required, with no default, so nobody gets an estimate thinking it is the census. True imports the census; False the latest estimate.

required
target str | None

DuckLake target; None (default) is data/raw/omnisus.ducklake under the working directory, or under $OMNISUS_DATA_DIR.

None

Returns:

Name Type Description
One list[ImportResult]

class:ImportResult per year, with rows and publication_id.

Raises:

Type Description
ValueError

the year has no such edition.

TypeError

census was not given.

Examples:

>>> import omnisus as sus
>>> (censo,) = sus.import_ibge_populacao(years=[2022], census=True)
>>> censo.rows
5570
Source code in src/omnisus/__init__.py
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
def import_ibge_populacao(
    *,
    years: Iterable[int],
    census: bool,
    target: str | None = None,
) -> list[ImportResult]:
    """Import IBGE population by municipality: a census or the year's latest estimate.

    Each population edition is its own publication, identified by the returned
    ``publication_id`` in ``ibge_population_manifest``; it never appears in
    ``Lake.publications()``. Reconcile an interrupted run through that manifest.

    Args:
        years: Years, e.g. ``[2022]``. Censuses exist for 2010 and 2022; an estimate
            is importable only as the aggregate's latest edition, and years without a
            verified territorial universe are refused.
        census: Required, with no default, so nobody gets an estimate thinking it is
            the census. ``True`` imports the census; ``False`` the latest estimate.
        target: DuckLake target; ``None`` (default) is ``data/raw/omnisus.ducklake``
            under the working directory, or under ``$OMNISUS_DATA_DIR``.

    Returns:
        One :class:`ImportResult` per year, with ``rows`` and ``publication_id``.

    Raises:
        ValueError: the year has no such edition.
        TypeError: ``census`` was not given.

    Examples:
        >>> import omnisus as sus
        >>> (censo,) = sus.import_ibge_populacao(years=[2022], census=True)  # doctest: +SKIP
        >>> censo.rows  # doctest: +SKIP
        5570
    """
    from omnisus.sources.ibge.importers.pop import import_pop_year

    edition = "census" if census else "estimate"

    async def run() -> list[ImportResult]:
        results: list[ImportResult] = []
        with Lake.local(target) as lake:
            for y in years:
                results.append(await import_pop_year(year=y, lake=lake, product=edition))
        return results

    return run_sync(run)

Fetch CNES establishment names from the public API into cnes_master.

The CNES-ST DBF has no establishment names; this pulls them from apidadosabertos.saude.gov.br and joins them into the aux_cnes view. It is an idempotent upsert, not a publication: no run_id, no policy, no manifest row. To finish an interrupted run, run it again.

Parameters:

Name Type Description Default
codes Sequence[str] | None

7-digit CNES codes, e.g. ["2789590"]. None (default) takes every code in cnes_estabelecimentos.

None
target str | None

DuckLake target; None (default) is data/raw/omnisus.ducklake under the working directory, or under $OMNISUS_DATA_DIR.

None
concurrency int

API requests in flight (default 5).

5
only_missing bool

True (default) skips codes already in cnes_master, so reruns are incremental; False fetches every code again.

True
progress Callable[[int, int], None] | None

Optional progress(done, total) callback after each fetch.

None

Returns:

Type Description
int

The number of records written in this run.

Raises:

Type Description
ValueError

a code is not seven ASCII digits, or the API returned conflicting records for one code.

Examples:

>>> import omnisus as sus
>>> sus.import_cnes_master(codes=["2789590"])
1
Source code in src/omnisus/__init__.py
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
def import_cnes_master(
    *,
    codes: Sequence[str] | None = None,
    target: str | None = None,
    concurrency: int = 5,
    only_missing: bool = True,
    progress: Callable[[int, int], None] | None = None,
) -> int:
    """Fetch CNES establishment names from the public API into ``cnes_master``.

    The CNES-ST DBF has no establishment names; this pulls them from
    ``apidadosabertos.saude.gov.br`` and joins them into the ``aux_cnes`` view. It
    is an idempotent upsert, not a publication: no ``run_id``, no ``policy``, no
    manifest row. To finish an interrupted run, run it again.

    Args:
        codes: 7-digit CNES codes, e.g. ``["2789590"]``. ``None`` (default) takes
            every code in ``cnes_estabelecimentos``.
        target: DuckLake target; ``None`` (default) is ``data/raw/omnisus.ducklake``
            under the working directory, or under ``$OMNISUS_DATA_DIR``.
        concurrency: API requests in flight (default 5).
        only_missing: ``True`` (default) skips codes already in ``cnes_master``, so
            reruns are incremental; ``False`` fetches every code again.
        progress: Optional ``progress(done, total)`` callback after each fetch.

    Returns:
        The number of records written in this run.

    Raises:
        ValueError: a code is not seven ASCII digits, or the API returned conflicting
            records for one code.

    Examples:
        >>> import omnisus as sus
        >>> sus.import_cnes_master(codes=["2789590"])  # doctest: +SKIP
        1
    """
    from omnisus.sources.cnes.importers.master import (
        import_cnes_master as _impl,
    )

    return _impl(
        codes=codes,
        target=target,
        concurrency=concurrency,
        only_missing=only_missing,
        progress=progress,
    )

Import SIGTAP procedures, one competência per year x month.

Each competência is one national monthly publication from ftp2.datasus.gov.br/public/sistemas/tup/downloads into aux_sigtap_procedimentos. Nothing joins at import: join on co_procedimento and the record's competência yourself.

Parameters:

Name Type Description Default
years Iterable[int]

Years, e.g. [2024].

required
months Iterable[int] | None

Months 1-12, e.g. [1, 2]; None (default) means all twelve.

None
target str | None

DuckLake target; None (default) is data/raw/omnisus.ducklake under the working directory, or under $OMNISUS_DATA_DIR.

None
policy ImportPolicy

"skip_same" (default) reports an unchanged zip as skipped/unchanged; "append", "error_if_exists" and "replace" as in :func:import_dataset.

'skip_same'

Returns:

Name Type Description
One ImportReport

class:ScopeOutcome per competência; one the server does not list is

ImportReport

skipped/not_listed. :func:cite names each zip.

Examples:

>>> import omnisus as sus
>>> report = sus.import_sigtap(years=[2024], months=[1])
>>> [(str(o.scope), o.status) for o in report.outcomes]
[('national_2024_01', 'ok')]
Source code in src/omnisus/__init__.py
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
def import_sigtap(
    *,
    years: Iterable[int],
    months: Iterable[int] | None = None,
    target: str | None = None,
    policy: ImportPolicy = "skip_same",
) -> ImportReport:
    """Import SIGTAP procedures, one competência per year x month.

    Each competência is one national monthly publication from
    ``ftp2.datasus.gov.br/public/sistemas/tup/downloads`` into
    ``aux_sigtap_procedimentos``. Nothing joins at import: join on
    ``co_procedimento`` and the record's competência yourself.

    Args:
        years: Years, e.g. ``[2024]``.
        months: Months 1-12, e.g. ``[1, 2]``; ``None`` (default) means all twelve.
        target: DuckLake target; ``None`` (default) is ``data/raw/omnisus.ducklake``
            under the working directory, or under ``$OMNISUS_DATA_DIR``.
        policy: ``"skip_same"`` (default) reports an unchanged zip as
            ``skipped``/``unchanged``; ``"append"``, ``"error_if_exists"`` and
            ``"replace"`` as in :func:`import_dataset`.

    Returns:
        One :class:`ScopeOutcome` per competência; one the server does not list is
        ``skipped``/``not_listed``. :func:`cite` names each zip.

    Examples:
        >>> import omnisus as sus
        >>> report = sus.import_sigtap(years=[2024], months=[1])  # doctest: +SKIP
        >>> [(str(o.scope), o.status) for o in report.outcomes]  # doctest: +SKIP
        [('national_2024_01', 'ok')]
    """
    from omnisus.sources.sigtap import import_sigtap as _impl

    month_list = tuple(months) if months is not None else tuple(range(1, 13))
    wanted = [(year, month) for year in years for month in month_list]
    with Lake.local(target) as lake:
        return _impl(wanted, lake=lake, policy=policy)

Every SIGTAP competência the server lists now, oldest first.

Returns:

Type Description
list[tuple[int, int]]

(year, month) pairs, e.g. [(2008, 1), (2008, 2), ...].

Raises:

Type Description
FtpUnavailable

the server did not answer after retries.

Examples:

>>> import omnisus as sus
>>> sus.available_sigtap()[:2]
[(2008, 1), (2008, 2)]
Source code in src/omnisus/__init__.py
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
def available_sigtap() -> list[tuple[int, int]]:
    """Every SIGTAP competência the server lists now, oldest first.

    Returns:
        ``(year, month)`` pairs, e.g. ``[(2008, 1), (2008, 2), ...]``.

    Raises:
        FtpUnavailable: the server did not answer after retries.

    Examples:
        >>> import omnisus as sus
        >>> sus.available_sigtap()[:2]  # doctest: +SKIP
        [(2008, 1), (2008, 2)]
    """
    from omnisus.sources.sigtap import available_sigtap as _impl

    return _impl()

Vocabularies

The bootstrap tables (aux_uf, aux_municipios, aux_cid10, aux_ocupacoes, aux_paises) come from DATASUS's SIM/CID10/TABELAS; the CBO 2002 codes in aux_ocupacoes come from CBO2002.CNV in SIM/CID10/TAB/OBITOS_CID10_TAB.zip. Each file is hashed in the source registry. Nothing joins at import: reference_join_sql writes the LEFT JOIN a field's dictionary declares. See vocabularies.

The LEFT JOIN from a coded field to the vocabulary its dictionary declares.

The declaration is the field's foreignKeys entry; an x-join-rule adds a condition. Nothing joins at import: paste the result after FROM lake.<dataset> AS <alias>.

Parameters:

Name Type Description Default
dataset str

Dataset name, e.g. "sim_obitos".

required
field str

A field that declares a reference, e.g. "causabas" (CID-10) or "ocup" (CBO).

required
alias str

Alias of the dataset table in your query (default "d").

'd'
lake_alias str

Alias the lake is attached as (default "lake").

'lake'

Returns:

Type Description
str

SQL text; the vocabulary is aliased ref_<field>.

Raises:

Type Description
ValueError

the field declares no reference, or more than one.

FileNotFoundError

dataset has no packaged dictionary.

Examples:

>>> import omnisus as sus
>>> sus.reference_join_sql("sim_obitos", "causabas")
'LEFT JOIN "lake"."aux_cid10" AS "ref_causabas" ON d."causabas" = "ref_causabas"."codigo"'
Source code in src/omnisus/research.py
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
def reference_join_sql(
    dataset: str, field: str, *, alias: str = "d", lake_alias: str = "lake"
) -> str:
    """The ``LEFT JOIN`` from a coded field to the vocabulary its dictionary declares.

    The declaration is the field's ``foreignKeys`` entry; an ``x-join-rule`` adds a
    condition. Nothing joins at import: paste the result after
    ``FROM lake.<dataset> AS <alias>``.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``.
        field: A field that declares a reference, e.g. ``"causabas"`` (CID-10) or
            ``"ocup"`` (CBO).
        alias: Alias of the dataset table in your query (default ``"d"``).
        lake_alias: Alias the lake is attached as (default ``"lake"``).

    Returns:
        SQL text; the vocabulary is aliased ``ref_<field>``.

    Raises:
        ValueError: the field declares no reference, or more than one.
        FileNotFoundError: ``dataset`` has no packaged dictionary.

    Examples:
        >>> import omnisus as sus
        >>> sus.reference_join_sql("sim_obitos", "causabas")
        'LEFT JOIN "lake"."aux_cid10" AS "ref_causabas" ON d."causabas" = "ref_causabas"."codigo"'
    """
    definition = load_dicionario(dataset).field_def(field) or {}
    keys = definition.get("foreignKeys", [])
    if len(keys) != 1:
        raise ValueError(f"{dataset}.{field} declares no reference (or more than one)")
    reference = keys[0]["reference"]
    ref = quote_identifier(f"ref_{field}")
    sql = (
        f"LEFT JOIN {qualified(lake_alias, reference['resource'])} AS {ref} "
        f"ON {alias}.{quote_identifier(field)} = {ref}.{quote_identifier(reference['fields'])}"
    )
    rule = keys[0].get("x-join-rule")
    if rule:
        sql += " AND (" + rule.replace("{d}", alias).replace("{r}", ref) + ")"
    return sql

Dictionaries, labels and harmonised categories

Each dataset ships a dictionary: fields, labels, code maps and audited analytical rules. display_row is the per-row door to the same lookup as label. analytical_projection returns the SQL behind load's harmonised categories, for queries you write yourself; take the schema and the sources from the snapshot you query. See consumption.

A dataset's packaged metadata: fields, sources, analytical rules and a digest.

Works offline. schema.fields keeps the authored definitions; fields holds one self-contained document per column with its review state. Neither describes a live lake: use DESCRIBE at the snapshot you query.

Parameters:

Name Type Description Default
dataset str

Dataset name, e.g. "sim_obitos"; see :func:~omnisus.datasets.

required

Returns:

Type Description
dict[str, Any]

A new JSON-compatible dict on every call, with metadata_hash to pin it and

dict[str, Any]

analytics (version, validated_sources) for the harmonised

dict[str, Any]

categories.

Raises:

Type Description
FileNotFoundError

dataset has no packaged dictionary.

Examples:

>>> import omnisus as sus
>>> metadata = sus.describe_dataset("sim_obitos")
>>> sorted(metadata)[:3]
['analytics', 'dataset', 'dictionary_version']
>>> len(metadata["metadata_hash"])
64
Source code in src/omnisus/metadata.py
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
def describe_dataset(dataset: str) -> dict[str, Any]:
    """A dataset's packaged metadata: fields, sources, analytical rules and a digest.

    Works offline. ``schema.fields`` keeps the authored definitions;
    ``fields`` holds one self-contained document per column with its review state.
    Neither describes a live lake: use ``DESCRIBE`` at the snapshot you query.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``; see :func:`~omnisus.datasets`.

    Returns:
        A new JSON-compatible dict on every call, with ``metadata_hash`` to pin it and
        ``analytics`` (``version``, ``validated_sources``) for the harmonised
        categories.

    Raises:
        FileNotFoundError: ``dataset`` has no packaged dictionary.

    Examples:
        >>> import omnisus as sus
        >>> metadata = sus.describe_dataset("sim_obitos")
        >>> sorted(metadata)[:3]
        ['analytics', 'dataset', 'dictionary_version']
        >>> len(metadata["metadata_hash"])
        64
    """
    dictionary = load_dicionario(dataset)
    raw = deepcopy(dictionary.raw)
    registry = sources_registry()
    sources = {source["id"]: source for source in registry["sources"]}
    analytics = raw.get("x-analytics")
    if analytics:
        for name in ("age", "sex"):
            rule = analytics.get(name)
            if rule is None:
                continue
            if not rule.get("evidence"):
                raise ValueError(f"Analytical rule requires evidence: {name}")
            for evidence in rule["evidence"]:
                source_id = evidence["source_id"]
                if source_id not in sources:
                    raise ValueError(f"Unresolved analytical source: {source_id}")
                if sources[source_id]["authority"] != "official":
                    raise ValueError(f"Analytical rule requires official evidence: {name}")
    result = {
        "schema_version": SCHEMA_VERSION,
        "dataset": dictionary.name,
        "title": dictionary.title,
        "dictionary_version": dictionary.version,
        "schema": raw["schema"],
        "fields": [
            _column(dictionary.name, dictionary.version, field, sources)
            for field in raw["schema"]["fields"]
        ],
        "sources": registry["sources"],
        "analytics": raw.get("x-analytics"),
        "storage": {
            "policy": "source_physical_types",
            "dictionary_types": "descriptive",
            "provenance": "LakeReader.publications()",
            "observed_schema": None,
        },
    }
    used_sources = {source["id"] for column in result["fields"] for source in column["sources"]}
    if analytics:
        for name in ("age", "sex"):
            used_sources.update(
                evidence["source_id"] for evidence in analytics.get(name, {}).get("evidence", [])
            )
    result["sources"] = [sources[key] for key in sorted(used_sources)]
    result["metadata_hash"] = hashlib.sha256(canonical_json(result)).hexdigest()
    return result

Replace each code in one row by its dictionary label, without mutating the input.

The per-row door to the same lookup as :func:label: a coded field gets its label, or None when the dictionary does not know the code; every other field is returned as published.

Parameters:

Name Type Description Default
dataset str

Dataset name, e.g. "sim_obitos"; see sus.datasets().

required
row Mapping[str, Any]

One record as {column: value}. Keys match case-insensitively.

required

Returns:

Type Description
dict[str, Any]

A new dict with the same keys.

Raises:

Type Description
FileNotFoundError

dataset has no packaged dictionary.

Examples:

>>> import omnisus as sus
>>> sus.display_row("sim_obitos", {"sexo": "2", "dtobito": "01012023"})
{'sexo': 'Feminino', 'dtobito': '01012023'}
Source code in src/omnisus/transforms/dictionaries.py
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
def display_row(dataset: str, row: Mapping[str, Any]) -> dict[str, Any]:
    """Replace each code in one row by its dictionary label, without mutating the input.

    The per-row door to the same lookup as :func:`label`: a coded field gets its
    label, or ``None`` when the dictionary does not know the code; every other
    field is returned as published.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``; see ``sus.datasets()``.
        row: One record as ``{column: value}``. Keys match case-insensitively.

    Returns:
        A new dict with the same keys.

    Raises:
        FileNotFoundError: ``dataset`` has no packaged dictionary.

    Examples:
        >>> import omnisus as sus
        >>> sus.display_row("sim_obitos", {"sexo": "2", "dtobito": "01012023"})
        {'sexo': 'Feminino', 'dtobito': '01012023'}
    """
    return load_dicionario(dataset).decode_row(dict(row))

SQL for the harmonised categories of a dataset, enabled only for validated sources.

Opens no lake and changes no stored record: it returns expressions to put in a SELECT. :func:~omnisus.load already does this for you. Each category is enabled only when every source in scopes is listed, by scope, release and SHA-256, in the rule's validated_sources (ADR 0003); otherwise it is reported in unavailable with the reason. Take the schema and the sources from the same snapshot as the rows you select.

Parameters:

Name Type Description Default
dataset str

Dataset name, e.g. "sim_obitos".

required
observed_schema Mapping[str, str]

{column: DuckDB type} of the relation you will select from, e.g. from DESCRIBE. Names match case-insensitively.

required
scopes Sequence[SourceContext]

The :class:SourceContext of every source in that relation, e.g. SourceContext.from_publication(row) for each active publication.

required
rule_version str | None

Require this rules edition, e.g. "1.1.0"; None (default) takes the packaged one.

None

Returns:

Type Description
AnalyticalProjection

columns (name, DuckDB type, expression), unavailable (field, reason:

AnalyticalProjection

missing_source_field, unconfirmed_source_scope,

AnalyticalProjection

unsupported_dataset, ...), the rule version and the metadata hash.

Raises:

Type Description
ValueError

rule_version is not the packaged one, two observed columns differ only by case, or a derived name collides with a source column.

Examples:

>>> import omnisus as sus
>>> dorr2023 = sus.SourceContext(
...     sus.ScopeKey("RR", 2023),
...     "final",
...     "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
... )
>>> projection = sus.analytical_projection(
...     "sim_obitos", observed_schema={"sexo": "VARCHAR"}, scopes=[dorr2023]
... )
>>> [column.name for column in projection.columns]
['sexo_categoria', 'sexo_status']
Source code in src/omnisus/transforms/analytics.py
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
def analytical_projection(
    dataset: str,
    *,
    observed_schema: Mapping[str, str],
    scopes: Sequence[SourceContext],
    rule_version: str | None = None,
) -> AnalyticalProjection:
    """SQL for the harmonised categories of a dataset, enabled only for validated sources.

    Opens no lake and changes no stored record: it returns expressions to put in a
    ``SELECT``. :func:`~omnisus.load` already does this for you. Each category is
    enabled only when every source in ``scopes`` is listed, by scope, release and
    SHA-256, in the rule's ``validated_sources`` (ADR 0003); otherwise it is
    reported in ``unavailable`` with the reason. Take the schema and the sources from
    the same snapshot as the rows you select.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``.
        observed_schema: ``{column: DuckDB type}`` of the relation you will select
            from, e.g. from ``DESCRIBE``. Names match case-insensitively.
        scopes: The :class:`SourceContext` of every source in that relation, e.g.
            ``SourceContext.from_publication(row)`` for each active publication.
        rule_version: Require this rules edition, e.g. ``"1.1.0"``; ``None``
            (default) takes the packaged one.

    Returns:
        ``columns`` (name, DuckDB type, expression), ``unavailable`` (field, reason:
        ``missing_source_field``, ``unconfirmed_source_scope``,
        ``unsupported_dataset``, ...), the rule version and the metadata hash.

    Raises:
        ValueError: ``rule_version`` is not the packaged one, two observed columns
            differ only by case, or a derived name collides with a source column.

    Examples:
        >>> import omnisus as sus
        >>> dorr2023 = sus.SourceContext(
        ...     sus.ScopeKey("RR", 2023),
        ...     "final",
        ...     "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
        ... )
        >>> projection = sus.analytical_projection(
        ...     "sim_obitos", observed_schema={"sexo": "VARCHAR"}, scopes=[dorr2023]
        ... )
        >>> [column.name for column in projection.columns]
        ['sexo_categoria', 'sexo_status']
    """
    metadata = describe_dataset(dataset)
    rules = metadata["analytics"]
    version = rules["version"] if rules else None
    if rule_version is not None and rule_version != version:
        raise ValueError(f"Unavailable analytical rule version: {rule_version}")
    columns: list[DerivedColumn] = []
    unavailable: list[UnavailableField] = []
    names = {name.lower(): name for name in observed_schema}
    if len(names) != len(observed_schema):
        raise ValueError("Ambiguous case-insensitive source columns")
    definitions = {field["name"]: field for field in metadata["schema"]["fields"]}

    def available(output: str, inputs: Sequence[str], rule: Mapping[str, Any]) -> bool:
        if any(name not in names for name in inputs):
            unavailable.append(UnavailableField(output, "missing_source_field"))
            return False
        if not _confirmed(
            scopes, rule.get("validated_sources", rules.get("validated_sources", []))
        ):
            unavailable.append(UnavailableField(output, "unconfirmed_source_scope"))
            return False
        return True

    def add(name: str, dtype: str, expression: str) -> None:
        if name.lower() in names or any(column.name == name for column in columns):
            raise ValueError(f"Analytical column collision: {name}")
        columns.append(DerivedColumn(name, dtype, expression))

    if not rules:
        return AnalyticalProjection(
            (),
            (
                UnavailableField("idade_anos_completos", "unsupported_dataset"),
                UnavailableField("sexo_categoria", "unsupported_dataset"),
            ),
            None,
            metadata["metadata_hash"],
        )
    age = rules.get("age")
    if age is not None:
        inputs = [age["field"]] + ([age["unit_field"]] if age.get("unit_field") else [])
        if available("idade_anos_completos", inputs, age):
            sql = age_expressions(
                age,
                quote_identifier(names[inputs[0]]),
                quote_identifier(names[inputs[1]]) if len(inputs) > 1 else None,
            )
            add("idade_anos_completos", "INTEGER", sql.years)
            add("idade_status", "VARCHAR", sql.status)
            add("idade_quantidade", "INTEGER", sql.quantity)
            add("idade_unidade", "VARCHAR", sql.unit)
    else:
        unavailable.append(UnavailableField("idade_anos_completos", "unsupported_field"))
    sex = rules.get("sex")
    if sex is not None and available("sexo_categoria", [sex["field"]], sex):
        value, status = _sex_expressions(sex, quote_identifier(names[sex["field"]]))
        add("sexo_categoria", "VARCHAR", value)
        add("sexo_status", "VARCHAR", status)
    for name in rules.get("dates", []):
        output = f"{name}_data"
        if not available(output, [name], rules):
            continue
        definition = definitions[name]
        format_ = definition.get("x-format")
        if format_ not in DATE_FORMATS:
            unavailable.append(UnavailableField(output, "unsupported_date_format"))
            continue
        source_name = names[name]
        value, status = _date_expressions(
            quote_identifier(source_name), observed_schema[source_name], format_
        )
        add(output, "DATE", value)
        add(f"{output}_status", "VARCHAR", status)
    return AnalyticalProjection(
        tuple(columns), tuple(unavailable), version, metadata["metadata_hash"]
    )

Source identity, read from the same snapshot as the queried records.

Source code in src/omnisus/transforms/analytics.py
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
@dataclass(frozen=True)
class SourceContext:
    """Source identity, read from the same snapshot as the queried records."""

    scope: ScopeKey
    release: str | None = None
    source_sha256: str | None = None

    @classmethod
    def from_publication(cls, row: Mapping[str, Any]) -> SourceContext:
        """The source identity recorded by one ``Lake.publications()`` row.

        Args:
            row: A publication row with ``dataset``, ``scope_json``, ``source_uri``
                and ``source_sha256``.

        Returns:
            Scope, release (from the file's directory) and aggregate SHA-256.

        Raises:
            ValueError: the row lacks one of those keys or has no supported scope.

        Examples:
            >>> import omnisus as sus
            >>> sus.SourceContext.from_publication({
            ...     "dataset": "sim_obitos",
            ...     "scope_json": '{"ano": 2023, "uf": "RR"}',
            ...     "source_uri": "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/DORR2023.dbc",
            ...     "source_sha256": "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
            ... }).release
            'final'
        """
        try:
            fields = json.loads(row["scope_json"])
            if not isinstance(fields, dict):
                raise ValueError("Publication scope must be an object")
            scope = scope_from_fields(fields)
            if scope is None:
                raise ValueError("Publication has no supported scope")
            return cls(
                scope,
                release_from_uri(row["dataset"], row.get("source_uri")),
                row.get("source_sha256"),
            )
        except (KeyError, TypeError, json.JSONDecodeError) as exc:
            raise ValueError("Incomplete publication identity") from exc

from_publication(row) classmethod

The source identity recorded by one Lake.publications() row.

Parameters:

Name Type Description Default
row Mapping[str, Any]

A publication row with dataset, scope_json, source_uri and source_sha256.

required

Returns:

Type Description
SourceContext

Scope, release (from the file's directory) and aggregate SHA-256.

Raises:

Type Description
ValueError

the row lacks one of those keys or has no supported scope.

Examples:

>>> import omnisus as sus
>>> sus.SourceContext.from_publication({
...     "dataset": "sim_obitos",
...     "scope_json": '{"ano": 2023, "uf": "RR"}',
...     "source_uri": "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/DORR2023.dbc",
...     "source_sha256": "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
... }).release
'final'
Source code in src/omnisus/transforms/analytics.py
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
@classmethod
def from_publication(cls, row: Mapping[str, Any]) -> SourceContext:
    """The source identity recorded by one ``Lake.publications()`` row.

    Args:
        row: A publication row with ``dataset``, ``scope_json``, ``source_uri``
            and ``source_sha256``.

    Returns:
        Scope, release (from the file's directory) and aggregate SHA-256.

    Raises:
        ValueError: the row lacks one of those keys or has no supported scope.

    Examples:
        >>> import omnisus as sus
        >>> sus.SourceContext.from_publication({
        ...     "dataset": "sim_obitos",
        ...     "scope_json": '{"ano": 2023, "uf": "RR"}',
        ...     "source_uri": "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/DORR2023.dbc",
        ...     "source_sha256": "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
        ... }).release
        'final'
    """
    try:
        fields = json.loads(row["scope_json"])
        if not isinstance(fields, dict):
            raise ValueError("Publication scope must be an object")
        scope = scope_from_fields(fields)
        if scope is None:
            raise ValueError("Publication has no supported scope")
        return cls(
            scope,
            release_from_uri(row["dataset"], row.get("source_uri")),
            row.get("source_sha256"),
        )
    except (KeyError, TypeError, json.JSONDecodeError) as exc:
        raise ValueError("Incomplete publication identity") from exc

Results

ImportResult.bytes_written measures the temporary staging Parquet file, not final lake storage growth. snapshot_id=None can mean an uncommitted result or an unavailable snapshot ID; it does not alone establish whether a write committed. Managed FTP results also carry run_id, batch_id and publication_id; IBGE results carry their canonical publication_id.

Per-scope outcomes of one import run.

Normal FTP completion reports each requested input position. An ImportAbortedError carries a partial report of determined outcomes and separately identifies unresolved positions; inspect before retrying.

Inspect :attr:failed — never the report's truthiness. An empty run and a run where everything failed are different facts, and no falsy sentinel stands in for either.

Source code in src/omnisus/sources/_base.py
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
@dataclass(frozen=True)
class ImportReport:
    """Per-scope outcomes of one import run.

    Normal FTP completion reports each requested input position. An
    ImportAbortedError carries a partial report of determined outcomes and
    separately identifies unresolved positions; inspect before retrying.

    Inspect :attr:`failed` — never the report's truthiness. An empty run and a
    run where everything failed are different facts, and no falsy sentinel
    stands in for either.
    """

    outcomes: tuple[ScopeOutcome, ...]
    run_id: str | None = None

    @property
    def rows(self) -> int:
        """Total rows ingested across successful scopes."""
        return sum(o.result.rows for o in self.outcomes if o.result is not None)

    @property
    def ok(self) -> tuple[ScopeOutcome, ...]:
        return tuple(o for o in self.outcomes if o.status == "ok")

    @property
    def skipped(self) -> tuple[ScopeOutcome, ...]:
        return tuple(o for o in self.outcomes if o.status == "skipped")

    @property
    def failed(self) -> tuple[ScopeOutcome, ...]:
        return tuple(o for o in self.outcomes if o.status == "failed")

rows property

Total rows ingested across successful scopes.

What happened to one scope in an import run.

skipped and failed are different facts and must not be collapsed. skipped means the listing does not have the scope (not_listed), it is outside the row's coverage (outside_coverage), or the same source is already published (unchanged). failed is worth retrying.

Source code in src/omnisus/sources/_base.py
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
@dataclass(frozen=True)
class ScopeOutcome:
    """What happened to one scope in an import run.

    ``skipped`` and ``failed`` are different facts and must not be collapsed.
    ``skipped`` means the listing does not have the scope (``not_listed``), it
    is outside the row's coverage (``outside_coverage``), or the same source is
    already published (``unchanged``). ``failed`` is worth retrying.
    """

    scope: ScopeKey
    status: ScopeStatus
    result: ImportResult | None = None
    """Set when ``status == "ok"``, ``None`` otherwise."""

    reason: str | None = None
    """Set when ``status`` is ``skipped`` or ``failed``, ``None`` otherwise."""

    code: ScopeCode | None = None
    """Why the scope did not import; ``None`` only for ``ok``."""

code = None class-attribute instance-attribute

Why the scope did not import; None only for ok.

reason = None class-attribute instance-attribute

Set when status is skipped or failed, None otherwise.

result = None class-attribute instance-attribute

Set when status == "ok", None otherwise.

Outcome of one import_scope call.

Source code in src/omnisus/sources/_base.py
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
@dataclass
class ImportResult:
    """Outcome of one import_scope call."""

    rows: int
    bytes_written: int
    duration_seconds: float
    snapshot_id: int | None = None
    run_id: str | None = None
    batch_id: str | None = None
    publication_id: str | None = None

Identifies one slice of a dataset (e.g., SP/2024 or MG/2024/01).

mes is None for yearly datasets, set for monthly (SIH).

Source code in src/omnisus/sources/_base.py
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
@dataclass(frozen=True)
class ScopeKey:
    """Identifies one slice of a dataset (e.g., SP/2024 or MG/2024/01).

    `mes` is None for yearly datasets, set for monthly (SIH).
    """

    uf: str | None
    """None denotes a national source, never an artificial record-level UF."""
    ano: int
    mes: int | None = None

    def __str__(self) -> str:
        if self.mes is None:
            return f"{self.uf or 'national'}_{self.ano}"
        return f"{self.uf or 'national'}_{self.ano}_{self.mes:02d}"

uf instance-attribute

None denotes a national source, never an artificial record-level UF.

What :func:delete_scope removed: data rows, and the manifest rows retired with them.

Source code in src/omnisus/lake/publication.py
50
51
52
53
54
55
@dataclass(frozen=True)
class DeletionResult:
    """What :func:`delete_scope` removed: data rows, and the manifest rows retired with them."""

    rows_deleted: int
    publications_retired: int

The lake

Lake.local() and LakeReader() without a target open the default lake, in the folder set by set_lake_dir (where the lake lives). Pass a target such as ducklake:/path/omnisus.ducklake for another lake. Local handles enforce a cooperative writer lock for their lifetime. Use Lake.transaction() for managed writes; raw SQL transaction control is outside this contract. Managed transactions cannot nest.

To read, open a LakeReader on the same target: it attaches the catalog read-only, takes no lock, creates nothing and sets no option, so it runs alongside an import. Pass snapshot_id to pin the session; without it every statement reads the latest committed snapshot.

import omnisus as sus

with sus.LakeReader() as reader:
    latest = reader.snapshots()[-1]["snapshot_id"]
    rows = reader.connect().execute("SELECT count(*) FROM lake.sim_obitos").fetchone()

with sus.LakeReader(snapshot_id=latest) as reader:
    ...  # every statement here sees exactly that snapshot
import omnisus as sus
import polars as pl

with sus.Lake.local() as lake:
    with lake.transaction() as receipt:
        result = lake.ingest("example", pl.DataFrame({"id": [1]}).lazy())
        assert result.snapshot_id is None
    assert receipt.committed
    print(receipt.snapshot_id, result.rows)

The receipt's snapshot_id may remain None after a successful commit when no new snapshot was created or the lookup was unavailable. See Architecture for rollback and recovery boundaries.

Lake.publications(run_id=...) reads the durable source-publication manifest; each row carries scope, the ScopeKey the package wrote (None for a shape this version does not write), alongside the raw scope_json. Lake.attempts(run_id=...) reads separately recorded known failures. Lake.ingest_parquet appends a staging file directly. Lake.publish_scope adds scope validation, source identity and replay policy to that write. Lake.delete_scope(table, scope) removes one source scope and retires every publication within it in the same transaction; a yearly scope on a monthly table covers all its months. import_ibge_populacao and import_cnes_master do not take part in this manifest — see their docstrings for how each reconciles.

Lake.optimize(table) merges adjacent files. Lake.expire_snapshots and Lake.cleanup_files take older_than as a timezone-aware datetime and default to dry_run=True. Each returns a list of result dictionaries.

Keep the default lake in path for the rest of this session.

Every call without target (load, LakeReader, import_*) then reads and writes the lake there. On Colab, point it at the mounted Google Drive so the lake survives the runtime. It sets $OMNISUS_DATA_DIR, which scripts can set directly instead. The folder is created by the first import.

Parameters:

Name Type Description Default
path str | PathLike[str]

The lake's folder; ~ and relative paths are resolved now.

required

Returns:

Type Description
Path

The absolute folder.

Examples:

>>> import omnisus as sus
>>> sus.set_lake_dir("/content/drive/MyDrive/omnisus")
PosixPath('/content/drive/MyDrive/omnisus')
Source code in src/omnisus/lake/catalog.py
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
def set_lake_dir(path: str | os.PathLike[str]) -> Path:
    """Keep the default lake in ``path`` for the rest of this session.

    Every call without ``target`` (``load``, ``LakeReader``, ``import_*``) then reads
    and writes the lake there. On Colab, point it at the mounted Google Drive so the
    lake survives the runtime. It sets ``$OMNISUS_DATA_DIR``, which scripts can set
    directly instead. The folder is created by the first import.

    Args:
        path: The lake's folder; ``~`` and relative paths are resolved now.

    Returns:
        The absolute folder.

    Examples:
        >>> import omnisus as sus
        >>> sus.set_lake_dir("/content/drive/MyDrive/omnisus")  # doctest: +SKIP
        PosixPath('/content/drive/MyDrive/omnisus')
    """
    folder = Path(path).expanduser().resolve()
    os.environ[DATA_DIR_ENV] = str(folder)
    return folder

Bases: Session

Writer handle on a DuckLake.

Use :meth:local with a ducklake: target or :meth:cloud for a Postgres-backed catalog. Managed ingestion requires one writer per lake, including when the catalog is hosted in Postgres. To only read, open a :class:~omnisus.lake.session.LakeReader instead: it takes no lock.

Source code in src/omnisus/lake/operations.py
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
class Lake(Session):
    """Writer handle on a DuckLake.

    Use :meth:`local` with a ``ducklake:`` target or :meth:`cloud` for a
    Postgres-backed catalog. Managed ingestion requires one writer per lake,
    including when the catalog is hosted in Postgres. To only read, open a
    :class:`~omnisus.lake.session.LakeReader` instead: it takes no lock.
    """

    def __init__(self, *, target: CatalogURI, alias: str = "lake") -> None:
        self._target = target
        self._in_transaction = False
        self._pending_results: list[ImportResult] = []
        self._ensured: set[str] = set()
        """Tables this handle has already created and partitioned. Ensuring a
        table costs a schema read of the staging file; once per run is enough."""
        self._columns: dict[str, set[str]] = {}
        """Known column names per table, so schema reconciliation does not
        re-query the catalog for every scope."""
        from omnisus.lake.locking import WriterLock

        self._writer_lock = WriterLock(target.catalog_uri)
        try:
            con = make_connection(
                catalog_uri=target.catalog_uri,
                storage_root=target.storage_root,
                alias=alias,
            )
        except BaseException:
            self._writer_lock.close()
            raise
        super().__init__(con, alias)

    @classmethod
    def local(cls, target: str | None = None) -> Self:
        """Open, or create, a lake on local disk for writing.

        One writer per lake: a second handle opening it for writing, in this process
        or another, fails at once with ``WriterBusyError``. Use
        :class:`~omnisus.LakeReader` to read while another process writes.

        Args:
            target: ``"ducklake:<dir>/omnisus.ducklake"``: data files go in that
                directory and the SQLite catalog next to it. ``None`` (default) is
                ``data/raw/omnisus.ducklake`` under the working directory, or under
                ``$OMNISUS_DATA_DIR``.

        Returns:
            An open :class:`Lake`; use it as a context manager so it closes.

        Raises:
            ValueError: ``target`` does not start with ``ducklake:``.
            WriterBusyError: another handle already has this lake open for writing
                (``omnisus.lake.locking.WriterBusyError``).
            CatalogAttachError: DuckDB could not attach the catalog.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.bootstrap_auxiliares()
        """
        return cls(target=parse_target(resolve_target(target)))

    @classmethod
    def cloud(cls, *, catalog: str, storage: str) -> Self:
        """Open a lake whose catalog is in PostgreSQL and data in object storage.

        Args:
            catalog: ``"postgresql://user:password@host/db"``.
            storage: Data path, e.g. ``"s3://bucket/lake"``.

        Returns:
            An open :class:`Lake`.

        Raises:
            ValueError: ``catalog`` is not a ``postgresql://`` URI.
            CatalogAttachError: DuckDB could not attach the catalog.

        Examples:
            >>> import omnisus as sus
            >>> lake = sus.Lake.cloud(  # doctest: +SKIP
            ...     catalog="postgresql://u:p@db/lake", storage="s3://bucket/lake"
            ... )
        """
        if not catalog.startswith(("postgresql://", "postgres://")):
            raise ValueError("cloud catalog must be postgresql://")
        return cls(target=CatalogURI(catalog_uri=catalog, storage_root=storage))

    @property
    def in_transaction(self) -> bool:
        return self._in_transaction

    def _read_snapshot(self) -> int | None:
        row = self._con.execute(
            "SELECT id FROM ducklake_last_committed_snapshot(?)", [self._alias]
        ).fetchone()
        return None if row is None or row[0] is None else int(row[0])

    @contextlib.contextmanager
    def transaction(self) -> Iterator[TransactionReceipt]:
        """Group writes into one DuckLake transaction, hence one snapshot.

        Yields:
            A receipt whose ``snapshot_id`` is set once the context commits.

        Raises:
            RuntimeError: already inside a transaction on this handle.
            TransactionStateError: the transaction could not begin or end cleanly;
                the handle becomes unusable, so close it and inspect the catalog.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake, lake.transaction() as receipt:  # doctest: +SKIP
            ...     lake.ensure_aux_cnes_view()
        """
        self.connect()  # Reject closed or invalidated handles.
        if self._in_transaction:
            raise RuntimeError("nested Lake.transaction is not supported")
        self._columns.clear()
        self._ensured.clear()
        try:
            before = self._read_snapshot()
            self._con.execute("BEGIN TRANSACTION")
        except BaseException as exc:
            self._unusable = True
            if not isinstance(exc, Exception):
                raise
            raise TransactionStateError("could not begin managed transaction") from exc

        receipt = TransactionReceipt()
        self._in_transaction = True
        self._pending_results = []
        try:
            try:
                yield receipt
            except BaseException as original:
                try:
                    self._con.execute("ROLLBACK")
                except BaseException as rollback_error:
                    self._unusable = True
                    original.add_note(f"rollback also failed: {rollback_error}")
                    if isinstance(original, Exception):
                        if not isinstance(rollback_error, Exception):
                            raise rollback_error from original
                        raise TransactionStateError(
                            "rollback failed; handle unusable"
                        ) from original
                raise
            else:
                try:
                    self._con.execute("COMMIT")
                except BaseException as original:
                    self._unusable = True
                    try:
                        self._con.execute("ROLLBACK")
                    except BaseException as rollback_error:
                        original.add_note(f"rollback cleanup also failed: {rollback_error}")
                        if isinstance(original, Exception) and not isinstance(
                            rollback_error, Exception
                        ):
                            raise rollback_error from original
                    if not isinstance(original, Exception):
                        raise
                    raise CommitOutcomeUnknown(
                        "commit outcome unknown; inspect before retry"
                    ) from original

                receipt.committed = True
                try:
                    after = self._read_snapshot()
                except Exception:
                    logger.warning("lake.snapshot_unavailable_after_commit")
                else:
                    receipt.snapshot_id = after if after != before else None
                for result in self._pending_results:
                    result.snapshot_id = receipt.snapshot_id
        finally:
            self._in_transaction = False
            self._pending_results = []
            self._columns.clear()
            self._ensured.clear()

    def optimize(self, table: str) -> list[dict]:
        """Merge a table's small data files, keeping every historical snapshot.

        Args:
            table: Table name, e.g. ``"sim_obitos"``.

        Returns:
            One dict per maintenance step DuckLake reports.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.optimize("sim_obitos")
        """
        from omnisus.lake.maintenance import run_maintenance

        return run_maintenance(self.connect(), alias=self._alias, operation="compact", table=table)

    def expire_snapshots(self, *, older_than: datetime, dry_run: bool = True) -> list[dict]:
        """Forget snapshots older than a cutoff; only simulates unless told otherwise.

        Expired snapshots can no longer be read or cited.

        Args:
            older_than: Timezone-aware ``datetime`` cutoff.
            dry_run: ``True`` (default) lists what would expire; ``False`` expires it.

        Returns:
            The snapshots expired, or that would be.

        Raises:
            ValueError: ``older_than`` has no timezone.

        Examples:
            >>> from datetime import UTC, datetime
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.expire_snapshots(older_than=datetime(2026, 1, 1, tzinfo=UTC))
        """
        from omnisus.lake.maintenance import run_maintenance

        return run_maintenance(
            self.connect(),
            alias=self._alias,
            operation="expire",
            older_than=older_than,
            dry_run=dry_run,
        )

    def cleanup_files(self, *, older_than: datetime, dry_run: bool = True) -> list[dict]:
        """Delete data files no snapshot needs, older than a cutoff; simulates by default.

        Args:
            older_than: Timezone-aware ``datetime`` cutoff.
            dry_run: ``True`` (default) lists the files; ``False`` deletes them.

        Returns:
            The files deleted, or that would be.

        Raises:
            ValueError: ``older_than`` has no timezone.

        Examples:
            >>> from datetime import UTC, datetime
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.cleanup_files(older_than=datetime(2026, 1, 1, tzinfo=UTC))
        """
        from omnisus.lake.maintenance import run_maintenance

        return run_maintenance(
            self.connect(),
            alias=self._alias,
            operation="cleanup",
            older_than=older_than,
            dry_run=dry_run,
        )

    def bootstrap_auxiliares(self) -> None:
        """Load the ``aux_*`` vocabularies (UF, municipalities, CID-10, CBO, countries).

        They come from the package's bootstrap zip, each file hashed in the source
        registry. Re-running replaces the tables' contents.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.bootstrap_auxiliares()
        """
        self.connect()

        import io
        import os
        import tempfile
        import zipfile
        from importlib.resources import files

        zip_bytes = (files("omnisus.data") / "auxiliares-bootstrap.zip").read_bytes()
        with zipfile.ZipFile(io.BytesIO(zip_bytes)) as zf:
            for name in zf.namelist():
                if not name.endswith(".parquet"):
                    continue
                table = name.removesuffix(".parquet")
                with tempfile.NamedTemporaryFile(suffix=".parquet", delete=False) as tmp:
                    tmp.write(zf.read(name))
                    tmp_path = tmp.name
                try:
                    self._con.execute(
                        f"CREATE OR REPLACE TABLE {qualified(self._alias, table)} AS "
                        f"SELECT * FROM read_parquet({quote_literal(tmp_path)})"
                    )
                finally:
                    os.unlink(tmp_path)

    def ensure_aux_cnes_view(self) -> bool:
        """Create or refresh ``aux_cnes``: one row per CNES establishment.

        Every import of ``cnes_estabelecimentos`` already calls this. The view joins
        the latest competência of ``cnes_estabelecimentos`` with ``cnes_master``
        names: ``cnes`` (7 digits), ``nome`` (NULL until
        :func:`~omnisus.import_cnes_master` runs), ``tp_unid``, ``codufmun`` and
        ``yyyymm_max``.

        Returns:
            ``False`` when ``cnes_estabelecimentos`` does not exist yet, else ``True``.

        Raises:
            RuntimeError: the handle is closed or unusable.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.ensure_aux_cnes_view()
            True
        """
        self.connect()
        tables = set(self.tables())
        if "cnes_estabelecimentos" not in tables:
            return False

        operational = qualified(self._alias, "cnes_estabelecimentos")
        view = qualified(self._alias, "aux_cnes")
        if "cnes_master" in tables:
            nome_select = "m.nome"
            join_clause = f"LEFT JOIN {qualified(self._alias, 'cnes_master')} m USING (cnes)"
        else:
            nome_select = "CAST(NULL AS VARCHAR) AS nome"
            join_clause = ""

        # Select all rows at the latest competence. Identical repeated imports
        # can collapse in this view; conflicting rows must never be picked by
        # incidental load order. The guard remains in the view for later inserts.
        latest_cte = f"""
            WITH ranked AS (
                SELECT cnes, tp_unid, codufmun, CAST(ano AS INTEGER) * 100 + CAST(mes AS INTEGER) AS yyyymm_max,
                       DENSE_RANK() OVER (PARTITION BY cnes ORDER BY ano DESC NULLS LAST, mes DESC NULLS LAST) AS rank
                FROM {operational} WHERE cnes IS NOT NULL
            ), latest AS (
                SELECT DISTINCT cnes, tp_unid, codufmun, yyyymm_max FROM ranked WHERE rank=1
            ), checked AS (
                SELECT *, COUNT(*) OVER (PARTITION BY cnes) AS variants FROM latest
            )
        """
        conflict = self._con.execute(
            latest_cte
            + "SELECT cnes FROM checked WHERE variants > 1 OR yyyymm_max IS NULL LIMIT 1"
        ).fetchone()
        if conflict:
            raise ValueError(
                "ambiguous or conflicting latest CNES rows; reconcile source versions before refreshing"
            )
        self._con.execute(f"""
            CREATE OR REPLACE VIEW {view} AS
            {latest_cte}
            SELECT CASE WHEN variants != 1 OR yyyymm_max IS NULL
                        THEN error('conflicting latest CNES rows') ELSE checked.cnes END AS cnes,
                   {nome_select}, checked.tp_unid, checked.codufmun, checked.yyyymm_max
            FROM checked {join_clause}
            WHERE CASE WHEN variants != 1 OR yyyymm_max IS NULL THEN error('conflicting latest CNES rows') ELSE true END
        """)
        return True

    def _staging_columns(self, staging: str) -> list[tuple[str, str]]:
        """(name, type) of the staging file, read from the Parquet footer."""
        return [
            (str(name), str(dtype))
            for name, dtype, *_ in self._con.execute(
                f"DESCRIBE SELECT * FROM read_parquet({quote_literal(str(staging))})"
            ).fetchall()
        ]

    def _table_columns(self, table: str) -> set[str]:
        if table not in self._columns:
            rows = self._con.execute(
                """
                SELECT column_name FROM information_schema.columns
                WHERE table_catalog = ? AND table_schema = 'main' AND table_name = ?
                """,
                [self._alias, table],
            ).fetchall()
            self._columns[table] = {str(name) for (name,) in rows}
        return self._columns[table]

    def _ensure_table(self, table: str, staging: str, partition_by: tuple[str, ...]) -> None:
        """Create the table and set its partitioning, then keep its schema wide
        enough for what is being inserted.

        DATASUS changes its layouts between eras: SIM-DO goes 42 columns in
        1996, 45 in 2005, 61 in 2010, 90 in 2015 and 89 in 2020, adding *and*
        removing columns. A table created from whichever scope landed first
        therefore rejects most of the others, which is why a multi-year import
        could not build a lake at all. New columns are added as they appear;
        columns a given era lacks are left NULL by ``INSERT ... BY NAME``.
        """
        if table not in set(self.tables()):
            self._con.execute(
                f"CREATE TABLE {qualified(self._alias, table)} AS "
                f"SELECT * FROM read_parquet({quote_literal(str(staging))}) WHERE 1=0"
            )
            if partition_by:
                cols = ", ".join(quote_identifier(c) for c in partition_by)
                self._con.execute(
                    f"ALTER TABLE {qualified(self._alias, table)} SET PARTITIONED BY ({cols})"
                )
            self._columns.pop(table, None)
            self._ensured.add(table)
            return

        import pyarrow as pa
        import pyarrow.parquet as pq

        from omnisus.lake.schema import compatible_type

        known = self._table_columns(table)
        existing_types = dict(
            self._con.execute(
                "SELECT column_name, data_type FROM information_schema.columns WHERE table_catalog = ? AND table_schema = 'main' AND table_name = ?",
                [self._alias, table],
            ).fetchall()
        )
        incoming = self._staging_columns(staging)
        null_fields = {f.name for f in pq.read_schema(staging) if pa.types.is_null(f.type)}
        promotions = []
        # Validate ALL shared columns before altering any of them.
        for name, dtype in incoming:
            if name in existing_types and name not in null_fields:
                target_type = compatible_type(existing_types[name], dtype)
                if target_type != existing_types[name]:
                    promotions.append((name, target_type))
        for name, dtype in promotions:
            self._con.execute(
                f"ALTER TABLE {qualified(self._alias, table)} ALTER COLUMN {quote_identifier(name)} SET TYPE {dtype}"
            )
        missing = [(n, t) for n, t in incoming if n not in known]
        for name, dtype in missing:
            self._con.execute(
                f"ALTER TABLE {qualified(self._alias, table)} ADD COLUMN {quote_identifier(name)} {dtype}"
            )
            known.add(name)
        if missing:
            logger.info(
                "lake.schema_widened",
                table=table,
                added=[n for n, _ in missing],
            )
        self._ensured.add(table)

    def ingest(
        self,
        table: str,
        lazyframe: pl.LazyFrame,
        *,
        partition_by: tuple[str, ...] = (),
    ) -> ImportResult:
        """Append a polars LazyFrame to a table, without a publication record.

        Outside :meth:`transaction` it commits before returning; inside, the result's
        ``snapshot_id`` stays ``None`` until the context commits.

        Args:
            table: Destination table, e.g. ``"minha_tabela"``.
            lazyframe: The rows to materialize and insert.
            partition_by: Partition columns, applied once when the table is created.

        Returns:
            Rows inserted, staging bytes and duration.

        Raises:
            RuntimeError: the handle is closed or unusable.

        Examples:
            >>> import polars as pl
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.ingest("minha_tabela", pl.LazyFrame({"x": [1, 2]}))
        """
        import tempfile
        import time
        from pathlib import Path

        self.connect()
        if not self.in_transaction:
            with self.transaction():
                return self.ingest(table, lazyframe, partition_by=partition_by)
        start = time.monotonic()
        with tempfile.TemporaryDirectory(prefix="omnisus-staging-") as tmp:
            staging = Path(tmp) / "data.parquet"
            lazyframe.sink_parquet(staging, row_group_size=1_000_000, compression="zstd")
            result = self.ingest_parquet(table, staging, partition_by=partition_by)
        result.duration_seconds = time.monotonic() - start
        return result

    def ingest_parquet(
        self, table: str, staging: str | Path, *, partition_by: tuple[str, ...] = ()
    ) -> ImportResult:
        """Append a validated Parquet staging file to a table, without rewriting it.

        Args:
            table: Destination table, e.g. ``"minha_tabela"``.
            staging: Path of the Parquet file.
            partition_by: Partition columns, applied once when the table is created.

        Returns:
            Rows inserted, the staging file's size and duration.

        Raises:
            RuntimeError: the handle is closed or unusable.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.ingest_parquet("minha_tabela", "staging.parquet")
        """
        import time
        from pathlib import Path

        from omnisus.sources._base import ImportResult

        self.connect()
        staging = Path(staging)
        if not self.in_transaction:
            with self.transaction():
                return self.ingest_parquet(table, staging, partition_by=partition_by)
        start = time.monotonic()
        self._ensure_table(table, str(staging), partition_by)
        inserted = self._con.execute(
            f"INSERT INTO {qualified(self._alias, table)} BY NAME SELECT * FROM read_parquet({quote_literal(str(staging))})"
        ).fetchone()
        result = ImportResult(
            rows=0 if inserted is None else int(inserted[0]),
            bytes_written=staging.stat().st_size,
            duration_seconds=time.monotonic() - start,
        )
        self._pending_results.append(result)
        return result

    def publish_scope(self, table: str, staging: Path, **kwargs: Any) -> ImportResult | None:
        """Publish one scope from a staging file, recording its source and policy.

        The low-level step :func:`~omnisus.import_dataset` runs per scope.

        Args:
            table: Dataset table, e.g. ``"sim_obitos"``.
            staging: Path of the validated Parquet staging file.
            **kwargs: ``scope``, ``source_sha256``, ``parser_version``, ``policy``,
                ``run_id``, ``source_uri`` and ``source_files``, as recorded in the
                publication.

        Returns:
            The result, or ``None`` when ``policy="skip_same"`` found the same source.

        Raises:
            ValueError: the staging file does not match the scope or the table.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.publish_scope("sim_obitos", "staging.parquet", scope=scope, ...)
        """
        from omnisus.lake.publication import publish_scope

        return publish_scope(self, table, staging, **kwargs)

    def delete_scope(self, table: str, scope: ScopeKey) -> DeletionResult:
        """Delete one scope's rows and retire its publications, in one transaction.

        Args:
            table: Dataset table, e.g. ``"sim_obitos"``.
            scope: The scope to delete, e.g. ``sus.ScopeKey("RR", 2023)``.

        Returns:
            Rows deleted and publications retired.

        Raises:
            ValueError: the table's national/state shape does not match ``scope``.

        Examples:
            >>> import omnisus as sus
            >>> with sus.Lake.local() as lake:  # doctest: +SKIP
            ...     lake.delete_scope("sim_obitos", sus.ScopeKey("RR", 2023))
        """
        from omnisus.lake.publication import delete_scope

        return delete_scope(self, table, scope)

    def close(self) -> None:
        """Close the connection and release the writer lock.

        Examples:
            >>> import omnisus as sus
            >>> lake = sus.Lake.local()  # doctest: +SKIP
            >>> lake.close()  # doctest: +SKIP
        """
        try:
            super().close()
        finally:
            self._writer_lock.close()
            self._columns.clear()
            self._ensured.clear()

bootstrap_auxiliares()

Load the aux_* vocabularies (UF, municipalities, CID-10, CBO, countries).

They come from the package's bootstrap zip, each file hashed in the source registry. Re-running replaces the tables' contents.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.bootstrap_auxiliares()
Source code in src/omnisus/lake/operations.py
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
def bootstrap_auxiliares(self) -> None:
    """Load the ``aux_*`` vocabularies (UF, municipalities, CID-10, CBO, countries).

    They come from the package's bootstrap zip, each file hashed in the source
    registry. Re-running replaces the tables' contents.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.bootstrap_auxiliares()
    """
    self.connect()

    import io
    import os
    import tempfile
    import zipfile
    from importlib.resources import files

    zip_bytes = (files("omnisus.data") / "auxiliares-bootstrap.zip").read_bytes()
    with zipfile.ZipFile(io.BytesIO(zip_bytes)) as zf:
        for name in zf.namelist():
            if not name.endswith(".parquet"):
                continue
            table = name.removesuffix(".parquet")
            with tempfile.NamedTemporaryFile(suffix=".parquet", delete=False) as tmp:
                tmp.write(zf.read(name))
                tmp_path = tmp.name
            try:
                self._con.execute(
                    f"CREATE OR REPLACE TABLE {qualified(self._alias, table)} AS "
                    f"SELECT * FROM read_parquet({quote_literal(tmp_path)})"
                )
            finally:
                os.unlink(tmp_path)

cleanup_files(*, older_than, dry_run=True)

Delete data files no snapshot needs, older than a cutoff; simulates by default.

Parameters:

Name Type Description Default
older_than datetime

Timezone-aware datetime cutoff.

required
dry_run bool

True (default) lists the files; False deletes them.

True

Returns:

Type Description
list[dict]

The files deleted, or that would be.

Raises:

Type Description
ValueError

older_than has no timezone.

Examples:

>>> from datetime import UTC, datetime
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.cleanup_files(older_than=datetime(2026, 1, 1, tzinfo=UTC))
Source code in src/omnisus/lake/operations.py
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
def cleanup_files(self, *, older_than: datetime, dry_run: bool = True) -> list[dict]:
    """Delete data files no snapshot needs, older than a cutoff; simulates by default.

    Args:
        older_than: Timezone-aware ``datetime`` cutoff.
        dry_run: ``True`` (default) lists the files; ``False`` deletes them.

    Returns:
        The files deleted, or that would be.

    Raises:
        ValueError: ``older_than`` has no timezone.

    Examples:
        >>> from datetime import UTC, datetime
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.cleanup_files(older_than=datetime(2026, 1, 1, tzinfo=UTC))
    """
    from omnisus.lake.maintenance import run_maintenance

    return run_maintenance(
        self.connect(),
        alias=self._alias,
        operation="cleanup",
        older_than=older_than,
        dry_run=dry_run,
    )

close()

Close the connection and release the writer lock.

Examples:

>>> import omnisus as sus
>>> lake = sus.Lake.local()
>>> lake.close()
Source code in src/omnisus/lake/operations.py
620
621
622
623
624
625
626
627
628
629
630
631
632
633
def close(self) -> None:
    """Close the connection and release the writer lock.

    Examples:
        >>> import omnisus as sus
        >>> lake = sus.Lake.local()  # doctest: +SKIP
        >>> lake.close()  # doctest: +SKIP
    """
    try:
        super().close()
    finally:
        self._writer_lock.close()
        self._columns.clear()
        self._ensured.clear()

cloud(*, catalog, storage) classmethod

Open a lake whose catalog is in PostgreSQL and data in object storage.

Parameters:

Name Type Description Default
catalog str

"postgresql://user:password@host/db".

required
storage str

Data path, e.g. "s3://bucket/lake".

required

Returns:

Type Description
Self

An open :class:Lake.

Raises:

Type Description
ValueError

catalog is not a postgresql:// URI.

CatalogAttachError

DuckDB could not attach the catalog.

Examples:

>>> import omnisus as sus
>>> lake = sus.Lake.cloud(
...     catalog="postgresql://u:p@db/lake", storage="s3://bucket/lake"
... )
Source code in src/omnisus/lake/operations.py
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
@classmethod
def cloud(cls, *, catalog: str, storage: str) -> Self:
    """Open a lake whose catalog is in PostgreSQL and data in object storage.

    Args:
        catalog: ``"postgresql://user:password@host/db"``.
        storage: Data path, e.g. ``"s3://bucket/lake"``.

    Returns:
        An open :class:`Lake`.

    Raises:
        ValueError: ``catalog`` is not a ``postgresql://`` URI.
        CatalogAttachError: DuckDB could not attach the catalog.

    Examples:
        >>> import omnisus as sus
        >>> lake = sus.Lake.cloud(  # doctest: +SKIP
        ...     catalog="postgresql://u:p@db/lake", storage="s3://bucket/lake"
        ... )
    """
    if not catalog.startswith(("postgresql://", "postgres://")):
        raise ValueError("cloud catalog must be postgresql://")
    return cls(target=CatalogURI(catalog_uri=catalog, storage_root=storage))

delete_scope(table, scope)

Delete one scope's rows and retire its publications, in one transaction.

Parameters:

Name Type Description Default
table str

Dataset table, e.g. "sim_obitos".

required
scope ScopeKey

The scope to delete, e.g. sus.ScopeKey("RR", 2023).

required

Returns:

Type Description
DeletionResult

Rows deleted and publications retired.

Raises:

Type Description
ValueError

the table's national/state shape does not match scope.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.delete_scope("sim_obitos", sus.ScopeKey("RR", 2023))
Source code in src/omnisus/lake/operations.py
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
def delete_scope(self, table: str, scope: ScopeKey) -> DeletionResult:
    """Delete one scope's rows and retire its publications, in one transaction.

    Args:
        table: Dataset table, e.g. ``"sim_obitos"``.
        scope: The scope to delete, e.g. ``sus.ScopeKey("RR", 2023)``.

    Returns:
        Rows deleted and publications retired.

    Raises:
        ValueError: the table's national/state shape does not match ``scope``.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.delete_scope("sim_obitos", sus.ScopeKey("RR", 2023))
    """
    from omnisus.lake.publication import delete_scope

    return delete_scope(self, table, scope)

ensure_aux_cnes_view()

Create or refresh aux_cnes: one row per CNES establishment.

Every import of cnes_estabelecimentos already calls this. The view joins the latest competência of cnes_estabelecimentos with cnes_master names: cnes (7 digits), nome (NULL until :func:~omnisus.import_cnes_master runs), tp_unid, codufmun and yyyymm_max.

Returns:

Type Description
bool

False when cnes_estabelecimentos does not exist yet, else True.

Raises:

Type Description
RuntimeError

the handle is closed or unusable.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.ensure_aux_cnes_view()
True
Source code in src/omnisus/lake/operations.py
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
def ensure_aux_cnes_view(self) -> bool:
    """Create or refresh ``aux_cnes``: one row per CNES establishment.

    Every import of ``cnes_estabelecimentos`` already calls this. The view joins
    the latest competência of ``cnes_estabelecimentos`` with ``cnes_master``
    names: ``cnes`` (7 digits), ``nome`` (NULL until
    :func:`~omnisus.import_cnes_master` runs), ``tp_unid``, ``codufmun`` and
    ``yyyymm_max``.

    Returns:
        ``False`` when ``cnes_estabelecimentos`` does not exist yet, else ``True``.

    Raises:
        RuntimeError: the handle is closed or unusable.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.ensure_aux_cnes_view()
        True
    """
    self.connect()
    tables = set(self.tables())
    if "cnes_estabelecimentos" not in tables:
        return False

    operational = qualified(self._alias, "cnes_estabelecimentos")
    view = qualified(self._alias, "aux_cnes")
    if "cnes_master" in tables:
        nome_select = "m.nome"
        join_clause = f"LEFT JOIN {qualified(self._alias, 'cnes_master')} m USING (cnes)"
    else:
        nome_select = "CAST(NULL AS VARCHAR) AS nome"
        join_clause = ""

    # Select all rows at the latest competence. Identical repeated imports
    # can collapse in this view; conflicting rows must never be picked by
    # incidental load order. The guard remains in the view for later inserts.
    latest_cte = f"""
        WITH ranked AS (
            SELECT cnes, tp_unid, codufmun, CAST(ano AS INTEGER) * 100 + CAST(mes AS INTEGER) AS yyyymm_max,
                   DENSE_RANK() OVER (PARTITION BY cnes ORDER BY ano DESC NULLS LAST, mes DESC NULLS LAST) AS rank
            FROM {operational} WHERE cnes IS NOT NULL
        ), latest AS (
            SELECT DISTINCT cnes, tp_unid, codufmun, yyyymm_max FROM ranked WHERE rank=1
        ), checked AS (
            SELECT *, COUNT(*) OVER (PARTITION BY cnes) AS variants FROM latest
        )
    """
    conflict = self._con.execute(
        latest_cte
        + "SELECT cnes FROM checked WHERE variants > 1 OR yyyymm_max IS NULL LIMIT 1"
    ).fetchone()
    if conflict:
        raise ValueError(
            "ambiguous or conflicting latest CNES rows; reconcile source versions before refreshing"
        )
    self._con.execute(f"""
        CREATE OR REPLACE VIEW {view} AS
        {latest_cte}
        SELECT CASE WHEN variants != 1 OR yyyymm_max IS NULL
                    THEN error('conflicting latest CNES rows') ELSE checked.cnes END AS cnes,
               {nome_select}, checked.tp_unid, checked.codufmun, checked.yyyymm_max
        FROM checked {join_clause}
        WHERE CASE WHEN variants != 1 OR yyyymm_max IS NULL THEN error('conflicting latest CNES rows') ELSE true END
    """)
    return True

expire_snapshots(*, older_than, dry_run=True)

Forget snapshots older than a cutoff; only simulates unless told otherwise.

Expired snapshots can no longer be read or cited.

Parameters:

Name Type Description Default
older_than datetime

Timezone-aware datetime cutoff.

required
dry_run bool

True (default) lists what would expire; False expires it.

True

Returns:

Type Description
list[dict]

The snapshots expired, or that would be.

Raises:

Type Description
ValueError

older_than has no timezone.

Examples:

>>> from datetime import UTC, datetime
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.expire_snapshots(older_than=datetime(2026, 1, 1, tzinfo=UTC))
Source code in src/omnisus/lake/operations.py
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
def expire_snapshots(self, *, older_than: datetime, dry_run: bool = True) -> list[dict]:
    """Forget snapshots older than a cutoff; only simulates unless told otherwise.

    Expired snapshots can no longer be read or cited.

    Args:
        older_than: Timezone-aware ``datetime`` cutoff.
        dry_run: ``True`` (default) lists what would expire; ``False`` expires it.

    Returns:
        The snapshots expired, or that would be.

    Raises:
        ValueError: ``older_than`` has no timezone.

    Examples:
        >>> from datetime import UTC, datetime
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.expire_snapshots(older_than=datetime(2026, 1, 1, tzinfo=UTC))
    """
    from omnisus.lake.maintenance import run_maintenance

    return run_maintenance(
        self.connect(),
        alias=self._alias,
        operation="expire",
        older_than=older_than,
        dry_run=dry_run,
    )

ingest(table, lazyframe, *, partition_by=())

Append a polars LazyFrame to a table, without a publication record.

Outside :meth:transaction it commits before returning; inside, the result's snapshot_id stays None until the context commits.

Parameters:

Name Type Description Default
table str

Destination table, e.g. "minha_tabela".

required
lazyframe LazyFrame

The rows to materialize and insert.

required
partition_by tuple[str, ...]

Partition columns, applied once when the table is created.

()

Returns:

Type Description
ImportResult

Rows inserted, staging bytes and duration.

Raises:

Type Description
RuntimeError

the handle is closed or unusable.

Examples:

>>> import polars as pl
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.ingest("minha_tabela", pl.LazyFrame({"x": [1, 2]}))
Source code in src/omnisus/lake/operations.py
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
def ingest(
    self,
    table: str,
    lazyframe: pl.LazyFrame,
    *,
    partition_by: tuple[str, ...] = (),
) -> ImportResult:
    """Append a polars LazyFrame to a table, without a publication record.

    Outside :meth:`transaction` it commits before returning; inside, the result's
    ``snapshot_id`` stays ``None`` until the context commits.

    Args:
        table: Destination table, e.g. ``"minha_tabela"``.
        lazyframe: The rows to materialize and insert.
        partition_by: Partition columns, applied once when the table is created.

    Returns:
        Rows inserted, staging bytes and duration.

    Raises:
        RuntimeError: the handle is closed or unusable.

    Examples:
        >>> import polars as pl
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.ingest("minha_tabela", pl.LazyFrame({"x": [1, 2]}))
    """
    import tempfile
    import time
    from pathlib import Path

    self.connect()
    if not self.in_transaction:
        with self.transaction():
            return self.ingest(table, lazyframe, partition_by=partition_by)
    start = time.monotonic()
    with tempfile.TemporaryDirectory(prefix="omnisus-staging-") as tmp:
        staging = Path(tmp) / "data.parquet"
        lazyframe.sink_parquet(staging, row_group_size=1_000_000, compression="zstd")
        result = self.ingest_parquet(table, staging, partition_by=partition_by)
    result.duration_seconds = time.monotonic() - start
    return result

ingest_parquet(table, staging, *, partition_by=())

Append a validated Parquet staging file to a table, without rewriting it.

Parameters:

Name Type Description Default
table str

Destination table, e.g. "minha_tabela".

required
staging str | Path

Path of the Parquet file.

required
partition_by tuple[str, ...]

Partition columns, applied once when the table is created.

()

Returns:

Type Description
ImportResult

Rows inserted, the staging file's size and duration.

Raises:

Type Description
RuntimeError

the handle is closed or unusable.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.ingest_parquet("minha_tabela", "staging.parquet")
Source code in src/omnisus/lake/operations.py
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
def ingest_parquet(
    self, table: str, staging: str | Path, *, partition_by: tuple[str, ...] = ()
) -> ImportResult:
    """Append a validated Parquet staging file to a table, without rewriting it.

    Args:
        table: Destination table, e.g. ``"minha_tabela"``.
        staging: Path of the Parquet file.
        partition_by: Partition columns, applied once when the table is created.

    Returns:
        Rows inserted, the staging file's size and duration.

    Raises:
        RuntimeError: the handle is closed or unusable.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.ingest_parquet("minha_tabela", "staging.parquet")
    """
    import time
    from pathlib import Path

    from omnisus.sources._base import ImportResult

    self.connect()
    staging = Path(staging)
    if not self.in_transaction:
        with self.transaction():
            return self.ingest_parquet(table, staging, partition_by=partition_by)
    start = time.monotonic()
    self._ensure_table(table, str(staging), partition_by)
    inserted = self._con.execute(
        f"INSERT INTO {qualified(self._alias, table)} BY NAME SELECT * FROM read_parquet({quote_literal(str(staging))})"
    ).fetchone()
    result = ImportResult(
        rows=0 if inserted is None else int(inserted[0]),
        bytes_written=staging.stat().st_size,
        duration_seconds=time.monotonic() - start,
    )
    self._pending_results.append(result)
    return result

local(target=None) classmethod

Open, or create, a lake on local disk for writing.

One writer per lake: a second handle opening it for writing, in this process or another, fails at once with WriterBusyError. Use :class:~omnisus.LakeReader to read while another process writes.

Parameters:

Name Type Description Default
target str | None

"ducklake:<dir>/omnisus.ducklake": data files go in that directory and the SQLite catalog next to it. None (default) is data/raw/omnisus.ducklake under the working directory, or under $OMNISUS_DATA_DIR.

None

Returns:

Type Description
Self

An open :class:Lake; use it as a context manager so it closes.

Raises:

Type Description
ValueError

target does not start with ducklake:.

WriterBusyError

another handle already has this lake open for writing (omnisus.lake.locking.WriterBusyError).

CatalogAttachError

DuckDB could not attach the catalog.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.bootstrap_auxiliares()
Source code in src/omnisus/lake/operations.py
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
@classmethod
def local(cls, target: str | None = None) -> Self:
    """Open, or create, a lake on local disk for writing.

    One writer per lake: a second handle opening it for writing, in this process
    or another, fails at once with ``WriterBusyError``. Use
    :class:`~omnisus.LakeReader` to read while another process writes.

    Args:
        target: ``"ducklake:<dir>/omnisus.ducklake"``: data files go in that
            directory and the SQLite catalog next to it. ``None`` (default) is
            ``data/raw/omnisus.ducklake`` under the working directory, or under
            ``$OMNISUS_DATA_DIR``.

    Returns:
        An open :class:`Lake`; use it as a context manager so it closes.

    Raises:
        ValueError: ``target`` does not start with ``ducklake:``.
        WriterBusyError: another handle already has this lake open for writing
            (``omnisus.lake.locking.WriterBusyError``).
        CatalogAttachError: DuckDB could not attach the catalog.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.bootstrap_auxiliares()
    """
    return cls(target=parse_target(resolve_target(target)))

optimize(table)

Merge a table's small data files, keeping every historical snapshot.

Parameters:

Name Type Description Default
table str

Table name, e.g. "sim_obitos".

required

Returns:

Type Description
list[dict]

One dict per maintenance step DuckLake reports.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.optimize("sim_obitos")
Source code in src/omnisus/lake/operations.py
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
def optimize(self, table: str) -> list[dict]:
    """Merge a table's small data files, keeping every historical snapshot.

    Args:
        table: Table name, e.g. ``"sim_obitos"``.

    Returns:
        One dict per maintenance step DuckLake reports.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.optimize("sim_obitos")
    """
    from omnisus.lake.maintenance import run_maintenance

    return run_maintenance(self.connect(), alias=self._alias, operation="compact", table=table)

publish_scope(table, staging, **kwargs)

Publish one scope from a staging file, recording its source and policy.

The low-level step :func:~omnisus.import_dataset runs per scope.

Parameters:

Name Type Description Default
table str

Dataset table, e.g. "sim_obitos".

required
staging Path

Path of the validated Parquet staging file.

required
**kwargs Any

scope, source_sha256, parser_version, policy, run_id, source_uri and source_files, as recorded in the publication.

{}

Returns:

Type Description
ImportResult | None

The result, or None when policy="skip_same" found the same source.

Raises:

Type Description
ValueError

the staging file does not match the scope or the table.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
...     lake.publish_scope("sim_obitos", "staging.parquet", scope=scope, ...)
Source code in src/omnisus/lake/operations.py
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
def publish_scope(self, table: str, staging: Path, **kwargs: Any) -> ImportResult | None:
    """Publish one scope from a staging file, recording its source and policy.

    The low-level step :func:`~omnisus.import_dataset` runs per scope.

    Args:
        table: Dataset table, e.g. ``"sim_obitos"``.
        staging: Path of the validated Parquet staging file.
        **kwargs: ``scope``, ``source_sha256``, ``parser_version``, ``policy``,
            ``run_id``, ``source_uri`` and ``source_files``, as recorded in the
            publication.

    Returns:
        The result, or ``None`` when ``policy="skip_same"`` found the same source.

    Raises:
        ValueError: the staging file does not match the scope or the table.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake:  # doctest: +SKIP
        ...     lake.publish_scope("sim_obitos", "staging.parquet", scope=scope, ...)
    """
    from omnisus.lake.publication import publish_scope

    return publish_scope(self, table, staging, **kwargs)

transaction()

Group writes into one DuckLake transaction, hence one snapshot.

Yields:

Type Description
TransactionReceipt

A receipt whose snapshot_id is set once the context commits.

Raises:

Type Description
RuntimeError

already inside a transaction on this handle.

TransactionStateError

the transaction could not begin or end cleanly; the handle becomes unusable, so close it and inspect the catalog.

Examples:

>>> import omnisus as sus
>>> with sus.Lake.local() as lake, lake.transaction() as receipt:
...     lake.ensure_aux_cnes_view()
Source code in src/omnisus/lake/operations.py
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
@contextlib.contextmanager
def transaction(self) -> Iterator[TransactionReceipt]:
    """Group writes into one DuckLake transaction, hence one snapshot.

    Yields:
        A receipt whose ``snapshot_id`` is set once the context commits.

    Raises:
        RuntimeError: already inside a transaction on this handle.
        TransactionStateError: the transaction could not begin or end cleanly;
            the handle becomes unusable, so close it and inspect the catalog.

    Examples:
        >>> import omnisus as sus
        >>> with sus.Lake.local() as lake, lake.transaction() as receipt:  # doctest: +SKIP
        ...     lake.ensure_aux_cnes_view()
    """
    self.connect()  # Reject closed or invalidated handles.
    if self._in_transaction:
        raise RuntimeError("nested Lake.transaction is not supported")
    self._columns.clear()
    self._ensured.clear()
    try:
        before = self._read_snapshot()
        self._con.execute("BEGIN TRANSACTION")
    except BaseException as exc:
        self._unusable = True
        if not isinstance(exc, Exception):
            raise
        raise TransactionStateError("could not begin managed transaction") from exc

    receipt = TransactionReceipt()
    self._in_transaction = True
    self._pending_results = []
    try:
        try:
            yield receipt
        except BaseException as original:
            try:
                self._con.execute("ROLLBACK")
            except BaseException as rollback_error:
                self._unusable = True
                original.add_note(f"rollback also failed: {rollback_error}")
                if isinstance(original, Exception):
                    if not isinstance(rollback_error, Exception):
                        raise rollback_error from original
                    raise TransactionStateError(
                        "rollback failed; handle unusable"
                    ) from original
            raise
        else:
            try:
                self._con.execute("COMMIT")
            except BaseException as original:
                self._unusable = True
                try:
                    self._con.execute("ROLLBACK")
                except BaseException as rollback_error:
                    original.add_note(f"rollback cleanup also failed: {rollback_error}")
                    if isinstance(original, Exception) and not isinstance(
                        rollback_error, Exception
                    ):
                        raise rollback_error from original
                if not isinstance(original, Exception):
                    raise
                raise CommitOutcomeUnknown(
                    "commit outcome unknown; inspect before retry"
                ) from original

            receipt.committed = True
            try:
                after = self._read_snapshot()
            except Exception:
                logger.warning("lake.snapshot_unavailable_after_commit")
            else:
                receipt.snapshot_id = after if after != before else None
            for result in self._pending_results:
                result.snapshot_id = receipt.snapshot_id
    finally:
        self._in_transaction = False
        self._pending_results = []
        self._columns.clear()
        self._ensured.clear()

Bases: Session

Read-only session on an existing lake.

target uses the same grammar as :meth:Lake.local — ducklake:./omnisus.ducklake or ducklake:postgresql://…?storage=… — and the storage part is ignored: the catalog records its own data path. Nothing is created and no option is set, so a reader never takes the writer lock and never blocks a writer; DuckLake refuses writes on the connection itself. snapshot_id pins the session to one snapshot (snapshots()[-1]["snapshot_id"] is the latest); without it every statement reads the latest committed snapshot.

Raises :class:~omnisus.lake.CatalogAttachError when the catalog does not exist or the pinned snapshot is unknown.

Source code in src/omnisus/lake/session.py
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
class LakeReader(Session):
    """Read-only session on an existing lake.

    ``target`` uses the same grammar as :meth:`Lake.local` —
    ``ducklake:./omnisus.ducklake`` or ``ducklake:postgresql://…?storage=…`` —
    and the storage part is ignored: the catalog records its own data path.
    Nothing is created and no option is set, so a reader never takes the
    writer lock and never blocks a writer; DuckLake refuses writes on the
    connection itself. ``snapshot_id`` pins the session to one snapshot
    (``snapshots()[-1]["snapshot_id"]`` is the latest); without it every
    statement reads the latest committed snapshot.

    Raises :class:`~omnisus.lake.CatalogAttachError` when the catalog does
    not exist or the pinned snapshot is unknown.
    """

    snapshot_id: int | None
    """The pinned snapshot, or ``None`` for a session that reads the latest one."""

    def __init__(
        self, target: str | None = None, *, snapshot_id: int | None = None, alias: str = "lake"
    ) -> None:
        if snapshot_id is not None and (type(snapshot_id) is not int or snapshot_id < 0):
            raise ValueError("snapshot_id must be a non-negative integer")
        self.snapshot_id = snapshot_id
        catalog = parse_target(resolve_target(target))
        con = make_reader_connection(
            catalog_uri=catalog.catalog_uri, alias=alias, snapshot_id=snapshot_id
        )
        super().__init__(con, alias)

snapshot_id = snapshot_id instance-attribute

The pinned snapshot, or None for a session that reads the latest one.

Research citations and joins

import_research and cite are above. cite reads the publication manifest (or ibge_population_manifest) and returns Portuguese text matching the reproducibility guide. latest_snapshot_id is the newest catalog snapshot. Municipality helpers take the leftmost 6 or 7 digits; they do not pad and they do not rewrite stored columns.

Structured citation plus the Portuguese paragraph the guide uses.

Source code in src/omnisus/research.py
28
29
30
31
32
33
34
35
36
37
@dataclass(frozen=True)
class Citation:
    """Structured citation plus the Portuguese paragraph the guide uses."""

    text: str
    snapshot_id: int
    omnisus: str
    dataset: str | None
    run_id: str | None
    publications: tuple[dict[str, Any], ...]

Format publication rows you already have, e.g. kept in a JSON file.

Parameters:

Name Type Description Default
publications Sequence[Mapping[str, Any]]

Rows of Lake.publications(), or IBGE manifest rows.

required
snapshot_id int

The lake snapshot the analysis read.

required
dataset str | None

Dataset name, used when a row lacks it; "ibge_populacao" selects the IBGE wording.

None
run_id str | None

The run to name; None uses each row's own run_id.

None
accessed date | None

Access date to print; None (default) is today.

None

Returns:

Name Type Description
A Citation

class:~omnisus.Citation whose text is the Portuguese paragraph.

Examples:

>>> from datetime import date
>>> import omnisus as sus
>>> row = {
...     "dataset": "sim_obitos",
...     "source_uri": "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/DORR2023.dbc",
...     "source_sha256": "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
...     "run_id": "cap2",
... }
>>> print(sus.citation_from_publications([row], snapshot_id=7, accessed=date(2026, 9, 23)).text)
SIM — Declarações de Óbito (sim_obitos), arquivo DORR2023.dbc (...), SHA-256 15b5...baf0, acessado em 2026-09-23 pelo DATASUS. Importado com omnisus ..., lake snapshot 7, execução cap2.
Source code in src/omnisus/research.py
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
def citation_from_publications(
    publications: Sequence[Mapping[str, Any]],
    *,
    snapshot_id: int,
    dataset: str | None = None,
    run_id: str | None = None,
    accessed: date | None = None,
) -> Citation:
    """Format publication rows you already have, e.g. kept in a JSON file.

    Args:
        publications: Rows of ``Lake.publications()``, or IBGE manifest rows.
        snapshot_id: The lake snapshot the analysis read.
        dataset: Dataset name, used when a row lacks it; ``"ibge_populacao"`` selects
            the IBGE wording.
        run_id: The run to name; ``None`` uses each row's own ``run_id``.
        accessed: Access date to print; ``None`` (default) is today.

    Returns:
        A :class:`~omnisus.Citation` whose ``text`` is the Portuguese paragraph.

    Examples:
        >>> from datetime import date
        >>> import omnisus as sus
        >>> row = {
        ...     "dataset": "sim_obitos",
        ...     "source_uri": "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/DORR2023.dbc",
        ...     "source_sha256": "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
        ...     "run_id": "cap2",
        ... }
        >>> print(sus.citation_from_publications([row], snapshot_id=7, accessed=date(2026, 9, 23)).text)
        SIM — Declarações de Óbito (sim_obitos), arquivo DORR2023.dbc (...), SHA-256 15b5...baf0, acessado em 2026-09-23 pelo DATASUS. Importado com omnisus ..., lake snapshot 7, execução cap2.
    """
    accessed_on = accessed or datetime.now(UTC).date()
    rows = [dict(row) for row in publications]
    if dataset:
        for row in rows:
            row.setdefault("dataset", dataset)
    if dataset == "ibge_populacao" or (rows and "url" in rows[0] and "source_uri" not in rows[0]):
        paragraphs = [_ibge_paragraph(row, snapshot_id, accessed_on) for row in rows]
    else:
        paragraphs = [_ftp_paragraph(row, snapshot_id, accessed_on, run_id) for row in rows]
    if not paragraphs:
        paragraphs = [
            (
                f"Nenhuma publicação ativa encontrada"
                f"{f' para {dataset}' if dataset else ''}. "
                f"Importado com omnisus {__version__}, lake snapshot {snapshot_id}"
                f"{f', execução {run_id}' if run_id else ''}."
            )
        ]
    return Citation(
        text="\n\n".join(paragraphs),
        snapshot_id=snapshot_id,
        omnisus=__version__,
        dataset=dataset,
        run_id=run_id,
        publications=tuple(rows),
    )

The newest snapshot of the lake, to pin a query or a citation to it.

Parameters:

Name Type Description Default
lake Lake | LakeReader

An open :class:~omnisus.Lake or :class:~omnisus.LakeReader.

required

Returns:

Type Description
int

The snapshot id; pass it to LakeReader(snapshot_id=...) and :func:cite.

Raises:

Type Description
LookupError

the lake has no snapshots yet.

Examples:

>>> import omnisus as sus
>>> with sus.LakeReader() as lake:
...     snapshot = sus.latest_snapshot_id(lake)
Source code in src/omnisus/research.py
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
def latest_snapshot_id(lake: Lake | LakeReader) -> int:
    """The newest snapshot of the lake, to pin a query or a citation to it.

    Args:
        lake: An open :class:`~omnisus.Lake` or :class:`~omnisus.LakeReader`.

    Returns:
        The snapshot id; pass it to ``LakeReader(snapshot_id=...)`` and :func:`cite`.

    Raises:
        LookupError: the lake has no snapshots yet.

    Examples:
        >>> import omnisus as sus
        >>> with sus.LakeReader() as lake:  # doctest: +SKIP
        ...     snapshot = sus.latest_snapshot_id(lake)
    """
    snaps = lake.snapshots()
    if not snaps:
        raise LookupError("lake has no snapshots")
    raw = snaps[-1]["snapshot_id"]
    if isinstance(raw, int) and not isinstance(raw, bool):
        return raw
    raise TypeError(f"snapshot_id must be int, got {type(raw).__name__}")

The leftmost digits of an IBGE municipality code, to join bases on them.

DATASUS bases use 6 digits and IBGE uses 7 (with a check digit). This keeps the leftmost digits and never pads or invents a check digit.

Parameters:

Name Type Description Default
value object

A municipality code, e.g. "3550308" or 355030.

required
digits int

6 (default) or 7.

6

Returns:

Type Description
str | None

The key as text, or None for None or a blank value.

Raises:

Type Description
ValueError

digits is not 6 or 7.

Examples:

>>> import omnisus as sus
>>> sus.municipality_join_key("3550308")
'355030'
>>> sus.municipality_join_key(" ") is None
True
Source code in src/omnisus/research.py
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
def municipality_join_key(value: object, *, digits: int = 6) -> str | None:
    """The leftmost digits of an IBGE municipality code, to join bases on them.

    DATASUS bases use 6 digits and IBGE uses 7 (with a check digit). This keeps the
    leftmost ``digits`` and never pads or invents a check digit.

    Args:
        value: A municipality code, e.g. ``"3550308"`` or ``355030``.
        digits: ``6`` (default) or ``7``.

    Returns:
        The key as text, or ``None`` for ``None`` or a blank value.

    Raises:
        ValueError: ``digits`` is not 6 or 7.

    Examples:
        >>> import omnisus as sus
        >>> sus.municipality_join_key("3550308")
        '355030'
        >>> sus.municipality_join_key(" ") is None
        True
    """
    _require_join_digits(digits)
    if value is None:
        return None
    text = str(value).strip()
    if not text:
        return None
    return text[:digits]

SQL for :func:municipality_join_key on a lake column.

Parameters:

Name Type Description Default
column str

Column name, e.g. "codmunres".

required
digits int

6 (default) or 7.

6

Returns:

Type Description
str

A DuckDB expression to use in SELECT or JOIN ... ON.

Raises:

Type Description
ValueError

digits is not 6 or 7.

Examples:

>>> import omnisus as sus
>>> sus.municipality_join_key_sql("codmunres")
'left(trim(CAST("codmunres" AS VARCHAR)), 6)'
Source code in src/omnisus/research.py
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
def municipality_join_key_sql(column: str, *, digits: int = 6) -> str:
    """SQL for :func:`municipality_join_key` on a lake column.

    Args:
        column: Column name, e.g. ``"codmunres"``.
        digits: ``6`` (default) or ``7``.

    Returns:
        A DuckDB expression to use in ``SELECT`` or ``JOIN ... ON``.

    Raises:
        ValueError: ``digits`` is not 6 or 7.

    Examples:
        >>> import omnisus as sus
        >>> sus.municipality_join_key_sql("codmunres")
        'left(trim(CAST("codmunres" AS VARCHAR)), 6)'
    """
    _require_join_digits(digits)
    ident = quote_identifier(column)
    return f"left(trim(CAST({ident} AS VARCHAR)), {digits})"

Registry

datasets() lists every curated FTP dataset; products() covers every importer family except SIGTAP: the FTP datasets and the two families outside the registry (ibge_populacao, cnes_master). It states, per family, the scope fields, accepted policies, how an interrupted run is reconciled and whether available() applies. Year rules for IBGE remain in omnisus.sources.ibge.products (CENSUS_YEARS, ESTIMATE_UNAVAILABLE_YEARS); an estimate is importable only as its latest edition, so there is no floor year.

Identity and location of one DATASUS-FTP dataset.

Holds only identity and location. Anything behavioural belongs in an importer module, never as a flag here.

Source code in src/omnisus/sources/datasus_ftp/datasets.py
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
@dataclass(frozen=True, kw_only=True)
class Dataset:
    """Identity and location of one DATASUS-FTP dataset.

    Holds *only* identity and location. Anything behavioural
    belongs in an importer module, never as a flag here.
    """

    name: str
    """Registry key = lake table = YAML stem = CLI name (``sim_obitos``)."""

    prefix: str
    """DATASUS filename prefix: ``DO``, ``DN``, ``RD``, ``BI``, ``ATD``…"""

    ftp_dir: str
    """Directory on ``ftp.datasus.gov.br`` holding this dataset's files."""

    cadence: Cadence
    """How DATASUS publishes files. Decides filename shape (upstream fact)."""

    partition_by: tuple[str, ...]
    """How the lake table is laid out (our storage policy).

    Applied by ``Lake.ingest`` via ``SET PARTITIONED BY`` when the table is
    created, so files land under ``ano=…/uf=…/``. Each scope is already one
    ``(uf, ano[, mes])``, so partitioning costs nothing at write time and buys
    query pruning."""

    coverage: tuple[YM, YM | None]
    """``(first, last)`` published; ``last=None`` means ongoing.

    Provisional until the Tier 3 probe validates it against the server.
    """

    dictionary: Path | None = None
    """Frictionless YAML. ``None`` -> packaged ``dicionarios/<name>.yaml``."""

    geography: Literal["state", "national"] = "state"
    """Source coverage. National SINAN filenames use PREFIX + BR + YY."""

    year_digits: Literal[2, 4] = 4
    """Year width in a yearly state file name: 4 (``DORR2023.dbc``) or 2 (SIM CID-9's
    ``DORRR79.DBC``). Monthly and national names always carry two digits."""

    national_code: str = "BR"
    """What a national file name carries between prefix and year: ``BR`` for SINAN
    (``CHAGBR23.dbc``), nothing for the SIM subsets (``DOEXT23.dbc``)."""

    prelim_dir: str | None = None
    """Directory where DATASUS publishes this dataset's preliminary files,
    when it has one. Same filenames as ``ftp_dir``; a file is in one or the
    other, never both."""

    @property
    def monthly(self) -> bool:
        return self.cadence == "monthly"

    def directories(self) -> dict[Release, str]:
        """Every FTP directory this dataset is published in, final first.

        Returns:
            ``{"final": path}``, plus ``"prelim"`` when DATASUS publishes preliminary
            files elsewhere.

        Examples:
            >>> import omnisus as sus
            >>> sus.resolve("sim_obitos").directories()["prelim"]
            '/dissemin/publicos/SIM/PRELIM/DORES'
        """
        dirs: dict[Release, str] = {"final": self.ftp_dir}
        if self.prelim_dir is not None:
            dirs["prelim"] = self.prelim_dir
        return dirs

cadence instance-attribute

How DATASUS publishes files. Decides filename shape (upstream fact).

coverage instance-attribute

(first, last) published; last=None means ongoing.

Provisional until the Tier 3 probe validates it against the server.

dictionary = None class-attribute instance-attribute

Frictionless YAML. None -> packaged dicionarios/<name>.yaml.

ftp_dir instance-attribute

Directory on ftp.datasus.gov.br holding this dataset's files.

geography = 'state' class-attribute instance-attribute

Source coverage. National SINAN filenames use PREFIX + BR + YY.

name instance-attribute

Registry key = lake table = YAML stem = CLI name (sim_obitos).

national_code = 'BR' class-attribute instance-attribute

What a national file name carries between prefix and year: BR for SINAN (CHAGBR23.dbc), nothing for the SIM subsets (DOEXT23.dbc).

partition_by instance-attribute

How the lake table is laid out (our storage policy).

Applied by Lake.ingest via SET PARTITIONED BY when the table is created, so files land under ano=…/uf=…/. Each scope is already one (uf, ano[, mes]), so partitioning costs nothing at write time and buys query pruning.

prefix instance-attribute

DATASUS filename prefix: DO, DN, RD, BI, ATD…

prelim_dir = None class-attribute instance-attribute

Directory where DATASUS publishes this dataset's preliminary files, when it has one. Same filenames as ftp_dir; a file is in one or the other, never both.

year_digits = 4 class-attribute instance-attribute

Year width in a yearly state file name: 4 (DORR2023.dbc) or 2 (SIM CID-9's DORRR79.DBC). Monthly and national names always carry two digits.

directories()

Every FTP directory this dataset is published in, final first.

Returns:

Type Description
dict[Release, str]

{"final": path}, plus "prelim" when DATASUS publishes preliminary

dict[Release, str]

files elsewhere.

Examples:

>>> import omnisus as sus
>>> sus.resolve("sim_obitos").directories()["prelim"]
'/dissemin/publicos/SIM/PRELIM/DORES'
Source code in src/omnisus/sources/datasus_ftp/datasets.py
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
def directories(self) -> dict[Release, str]:
    """Every FTP directory this dataset is published in, final first.

    Returns:
        ``{"final": path}``, plus ``"prelim"`` when DATASUS publishes preliminary
        files elsewhere.

    Examples:
        >>> import omnisus as sus
        >>> sus.resolve("sim_obitos").directories()["prelim"]
        '/dissemin/publicos/SIM/PRELIM/DORES'
    """
    dirs: dict[Release, str] = {"final": self.ftp_dir}
    if self.prelim_dir is not None:
        dirs["prelim"] = self.prelim_dir
    return dirs

The :class:Dataset for a name; a Dataset value passes through.

A value passing through untouched is how an uncurated dataset reaches the pipeline.

Parameters:

Name Type Description Default
dataset str | Dataset

Dataset name, e.g. "sim_obitos", or a Dataset.

required

Returns:

Type Description
Dataset

The registry row: name, FTP directories, geography, monthly or yearly.

Raises:

Type Description
ValueError

the name is not in the registry.

Examples:

>>> import omnisus as sus
>>> d = sus.resolve("sih_aih_reduzida")
>>> (d.geography, d.monthly)
('state', True)
Source code in src/omnisus/sources/datasus_ftp/datasets.py
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
def resolve(dataset: str | Dataset) -> Dataset:
    """The :class:`Dataset` for a name; a ``Dataset`` value passes through.

    A value passing through untouched is how an uncurated dataset reaches the
    pipeline.

    Args:
        dataset: Dataset name, e.g. ``"sim_obitos"``, or a ``Dataset``.

    Returns:
        The registry row: name, FTP directories, geography, monthly or yearly.

    Raises:
        ValueError: the name is not in the registry.

    Examples:
        >>> import omnisus as sus
        >>> d = sus.resolve("sih_aih_reduzida")
        >>> (d.geography, d.monthly)
        ('state', True)
    """
    if isinstance(dataset, Dataset):
        return dataset
    try:
        return REGISTRY[dataset]
    except KeyError:
        raise ValueError(f"unknown dataset: {dataset!r}") from None

Every DATASUS FTP dataset the package curates, in registry order.

Returns:

Type Description
Dataset

class:~omnisus.Dataset values; .name is what :func:~omnisus.load

...

and the other functions accept.

Examples:

>>> import omnisus as sus
>>> "sim_obitos" in [d.name for d in sus.datasets()]
True
Source code in src/omnisus/products.py
45
46
47
48
49
50
51
52
53
54
55
56
57
def datasets() -> tuple[Dataset, ...]:
    """Every DATASUS FTP dataset the package curates, in registry order.

    Returns:
        :class:`~omnisus.Dataset` values; ``.name`` is what :func:`~omnisus.load`
        and the other functions accept.

    Examples:
        >>> import omnisus as sus
        >>> "sim_obitos" in [d.name for d in sus.datasets()]
        True
    """
    return tuple(REGISTRY.values())

Every importer family except SIGTAP, stated once: the FTP registry, IBGE and CNES master.

Facts, not forms. Year rules for IBGE stay in sources.ibge.products (CENSUS_YEARS, ESTIMATE_UNAVAILABLE_YEARS); an estimate is importable only as its latest edition, checked live at import, so there is no floor year.

Product dataclass

One importer family and what it supports.

Source code in src/omnisus/products.py
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
@dataclass(frozen=True)
class Product:
    """One importer family and what it supports."""

    name: str
    """Registry key or importer name: ``sim_obitos``, ``ibge_populacao``, ``cnes_master``."""

    dataset: Dataset | None
    """The FTP registry row, or ``None`` for IBGE and CNES master."""

    scope_fields: tuple[str, ...]
    """What identifies one import unit: ``("uf", "ano"[, "mes"])``, ``("ano",)``
    for national datasets, ``("product", "ano")`` for IBGE, empty for CNES master."""

    policies: tuple[str, ...]
    """Accepted ``ImportPolicy`` values; families outside the manifest append only."""

    reconcile_by: ReconcileBy
    """How an interrupted run is reconciled: ``Lake.publications(run_id=)``, the
    returned ``publication_id`` in ``ibge_population_manifest``, or a rerun
    (CNES master is an idempotent upsert)."""

    inventory: bool
    """Whether :func:`omnisus.available` can list the source's scopes."""

dataset instance-attribute

The FTP registry row, or None for IBGE and CNES master.

inventory instance-attribute

Whether :func:omnisus.available can list the source's scopes.

name instance-attribute

Registry key or importer name: sim_obitos, ibge_populacao, cnes_master.

policies instance-attribute

Accepted ImportPolicy values; families outside the manifest append only.

reconcile_by instance-attribute

How an interrupted run is reconciled: Lake.publications(run_id=), the returned publication_id in ibge_population_manifest, or a rerun (CNES master is an idempotent upsert).

scope_fields instance-attribute

What identifies one import unit: ("uf", "ano"[, "mes"]), ("ano",) for national datasets, ("product", "ano") for IBGE, empty for CNES master.

datasets()

Every DATASUS FTP dataset the package curates, in registry order.

Returns:

Type Description
Dataset

class:~omnisus.Dataset values; .name is what :func:~omnisus.load

...

and the other functions accept.

Examples:

>>> import omnisus as sus
>>> "sim_obitos" in [d.name for d in sus.datasets()]
True
Source code in src/omnisus/products.py
45
46
47
48
49
50
51
52
53
54
55
56
57
def datasets() -> tuple[Dataset, ...]:
    """Every DATASUS FTP dataset the package curates, in registry order.

    Returns:
        :class:`~omnisus.Dataset` values; ``.name`` is what :func:`~omnisus.load`
        and the other functions accept.

    Examples:
        >>> import omnisus as sus
        >>> "sim_obitos" in [d.name for d in sus.datasets()]
        True
    """
    return tuple(REGISTRY.values())

products()

Every importer family except SIGTAP: the FTP datasets, IBGE population, CNES master.

Returns:

Name Type Description
One Product

class:~omnisus.Product per family, with its scope fields, accepted

...

policies and how an interrupted run is reconciled.

Examples:

>>> import omnisus as sus
>>> [(p.name, p.reconcile_by) for p in sus.products()[-2:]]
[('ibge_populacao', 'publication_id'), ('cnes_master', 'rerun')]
Source code in src/omnisus/products.py
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
def products() -> tuple[Product, ...]:
    """Every importer family except SIGTAP: the FTP datasets, IBGE population, CNES master.

    Returns:
        One :class:`~omnisus.Product` per family, with its scope fields, accepted
        policies and how an interrupted run is reconciled.

    Examples:
        >>> import omnisus as sus
        >>> [(p.name, p.reconcile_by) for p in sus.products()[-2:]]
        [('ibge_populacao', 'publication_id'), ('cnes_master', 'rerun')]
    """
    ftp = tuple(
        Product(
            d.name,
            d,
            ("ano",)
            if d.geography == "national"
            else ("uf", "ano", "mes")
            if d.monthly
            else ("uf", "ano"),
            POLICIES,
            "run_id",
            True,
        )
        for d in datasets()
    )
    return (
        *ftp,
        Product("ibge_populacao", None, ("product", "ano"), ("append",), "publication_id", False),
        Product("cnes_master", None, (), ("append",), "rerun", False),
    )

One importer family and what it supports.

Source code in src/omnisus/products.py
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
@dataclass(frozen=True)
class Product:
    """One importer family and what it supports."""

    name: str
    """Registry key or importer name: ``sim_obitos``, ``ibge_populacao``, ``cnes_master``."""

    dataset: Dataset | None
    """The FTP registry row, or ``None`` for IBGE and CNES master."""

    scope_fields: tuple[str, ...]
    """What identifies one import unit: ``("uf", "ano"[, "mes"])``, ``("ano",)``
    for national datasets, ``("product", "ano")`` for IBGE, empty for CNES master."""

    policies: tuple[str, ...]
    """Accepted ``ImportPolicy`` values; families outside the manifest append only."""

    reconcile_by: ReconcileBy
    """How an interrupted run is reconciled: ``Lake.publications(run_id=)``, the
    returned ``publication_id`` in ``ibge_population_manifest``, or a rerun
    (CNES master is an idempotent upsert)."""

    inventory: bool
    """Whether :func:`omnisus.available` can list the source's scopes."""

dataset instance-attribute

The FTP registry row, or None for IBGE and CNES master.

inventory instance-attribute

Whether :func:omnisus.available can list the source's scopes.

name instance-attribute

Registry key or importer name: sim_obitos, ibge_populacao, cnes_master.

policies instance-attribute

Accepted ImportPolicy values; families outside the manifest append only.

reconcile_by instance-attribute

How an interrupted run is reconciled: Lake.publications(run_id=), the returned publication_id in ibge_population_manifest, or a rerun (CNES master is an idempotent upsert).

scope_fields instance-attribute

What identifies one import unit: ("uf", "ano"[, "mes"]), ("ano",) for national datasets, ("product", "ano") for IBGE, empty for CNES master.

Errors

Lake transaction errors are imported from omnisus.lake. A CommitOutcomeUnknown means the COMMIT raised and the handle is unusable; inspect the catalog before retrying. The FTP runner wraps transaction state failures in the top-level ImportAbortedError with partial progress. CatalogAttachError means the catalog could not be opened at all: .stage says which statement failed, and for a remote catalog the message never carries the connection string.

Bases: RuntimeError

Partial progress is known, but the remaining inputs need inspection.

Source code in src/omnisus/sources/_base.py
167
168
169
170
171
172
173
174
175
176
177
178
179
class ImportAbortedError(RuntimeError):
    """Partial progress is known, but the remaining inputs need inspection."""

    def __init__(
        self,
        report: ImportReport,
        unresolved: tuple[tuple[int, ScopeKey], ...],
    ) -> None:
        self.report = report
        self.unresolved = unresolved
        super().__init__(
            f"import aborted: {len(unresolved)} input(s) unresolved; inspect before retry"
        )

Bases: RuntimeError

The DuckLake catalog could not be opened, so there is no lake to use.

stage names the statement that failed: "install" (the ducklake extension), "attach" (the catalog itself) or "set_option"; the message adds the DuckDB error class. For a remote catalog the DuckDB error is not chained — its text can echo the connection string — so stage and class are the whole diagnostic. For a local catalog it is chained, so the file-level reason stays visible.

Source code in src/omnisus/lake/connection.py
15
16
17
18
19
20
21
22
23
24
25
26
27
28
class CatalogAttachError(RuntimeError):
    """The DuckLake catalog could not be opened, so there is no lake to use.

    ``stage`` names the statement that failed: ``"install"`` (the ducklake
    extension), ``"attach"`` (the catalog itself) or ``"set_option"``; the
    message adds the DuckDB error class. For a remote catalog the DuckDB error
    is not chained — its text can echo the connection string — so stage and
    class are the whole diagnostic. For a local catalog it is chained, so the
    file-level reason stays visible.
    """

    def __init__(self, stage: str, duckdb_error: str) -> None:
        self.stage = stage
        super().__init__(f"could not open DuckLake catalog: {duckdb_error} during {stage}")

Bases: RuntimeError

The handle cannot safely continue its managed transaction workflow.

Source code in src/omnisus/lake/_transactions.py
10
11
class TransactionStateError(RuntimeError):
    """The handle cannot safely continue its managed transaction workflow."""

Bases: TransactionStateError

COMMIT raised; do not infer rollback or retry the write automatically.

Source code in src/omnisus/lake/_transactions.py
14
15
class CommitOutcomeUnknown(TransactionStateError):  # noqa: N818 - mandated public API
    """COMMIT raised; do not infer rollback or retry the write automatically."""

Bases: Exception

The remote directory does not exist, or access was denied (550).

Terminal — never retried. Distinct from an existing but empty directory, which returns an empty :class:Listing.

Source code in src/omnisus/sources/datasus_ftp/inventory.py
42
43
44
45
46
47
class FtpPathNotFound(Exception):  # noqa: N818
    """The remote directory does not exist, or access was denied (550).

    Terminal — never retried. Distinct from an existing but empty directory,
    which returns an empty :class:`Listing`.
    """

Bases: Exception

The server could not be reached within the retry budget.

Source code in src/omnisus/sources/datasus_ftp/inventory.py
50
51
class FtpUnavailable(Exception):  # noqa: N818
    """The server could not be reached within the retry budget."""