Skip to content

fetch

champalimaud.fetch

Every network fetch in the project.

Writes the config data files that champalimaud.load reads. FETCHERS maps each dataset of champalimaud.load.DATASETS to the function that writes it. The fetches need a saved CAVE token, and the census needs CODEX_API_TOKEN if it is missing.

  1. cell types and other tables from Codex
  2. the connections of every proofread neuron, from CAVE
  3. skeletons, one SWC file per cell, from Codex
  4. every synapse of one example cell, from CAVE

FETCHERS = {'census': lambda: download_codex(Sources.CODEX_CENSUS, config.CENSUS), 'visual_types': lambda: download_codex('visual_neuron_types'), 'columns': lambda: download_codex('column_assignment'), 'classification': lambda: download_codex('classification'), 'connections': fetch_connections, 'skeletons': fetch_skeletons, 'example_synapses': fetch_example_synapses} module-attribute

The function that writes each dataset, by the name in champalimaud.load.DATASETS, in the order to run them.

Inventory

Bases: NamedTuple

What the server reports about its materializations.

Source code in champalimaud/fetch.py
class Inventory(NamedTuple):
    """What the server reports about its materializations."""

    versions: list
    version: int
    tables: list

Sources

Remote services and the ids of the products fetched from them.

Every id, URL, and limit below is a choice; the comment on each line says why.

Source code in champalimaud/fetch.py
class Sources:
    """Remote services and the ids of the products fetched from them.

    Every id, URL, and limit below is a choice; the comment on each
    line says why.
    """

    # The only CAVE datastack this account can read;
    # flywire_fafb_production returns 403.
    # Frozen at materialization 783 (the 2024 public release).
    STACK = "flywire_fafb_public"
    # One finished SWC per root id, in the public bucket behind Codex.
    # lod1 is their healed skeleton at materialization 783, which
    # matches STACK; the whole set is the sk_lod1_783_healed.zip on the
    # Codex download page.
    # Nothing in this path names a materialization, unlike the archive
    # under skeletons/fafb/archives/783/,
    # so a republish for a new one would silently change the geometry.
    # Confirmed reachable on 2026-09-29;
    # the Codex FAQ does not document it.
    SKELETON_URL = (
        "https://storage.googleapis.com/flywire-data/codex/skeletons/fafb/lod1"
    )
    # Codex account token, a different login from CAVE, from the account
    # page on codex.flywire.ai.
    # It lives in the gitignored .env at the repo root;
    # direnv exports it.
    # Needed only when config.CENSUS is missing.
    CODEX_ENV = "CODEX_API_TOKEN"
    CODEX_HOST = "https://codex.flywire.ai/api/download_resource"
    # The census; download_codex fetches any of CODEX_PRODUCTS on
    # demand.
    CODEX_CENSUS = "consolidated_cell_types"
    # Every table download_resource serves for fafb, checked by name on
    # 2026-09-29 (all gzipped CSV).
    # The skeletons are not in here;
    # they are one SWC per root id under SKELETON_URL.
    CODEX_PRODUCTS = (
        "consolidated_cell_types",
        "classification",
        "names",
        "neurons",
        "cell_stats",
        "coordinates",
        "labels",
        "processed_labels",
        "visual_neuron_types",
        "column_assignment",
        "connectivity_tags",
        "connections_princeton",
        "connections_princeton_no_threshold",
        "synapse_table",
    )
    # Codex also serves other datasets (hemibrain, ...); this work is
    # FAFB.
    CODEX_DATASET = "fafb"
    # The synapse set behind the Codex release (one row per synapse,
    # with a size and a neuropil).
    # Every connection count comes from here.
    # Views are not listed by get_tables.
    # On the cells compared (six LC10a) it reproduces every Codex
    # connection to proofread neurons with identical counts,
    # which TRANSMITTER_TABLE does not;
    # counting its rows is what the release's thresholds of five and
    # ten synapses refer to.
    SYNAPSE_VIEW = "synapses_v3_neuropil_v6_merge_view"
    # The unfiltered Buhmann predictions, with the per-transmitter
    # probabilities that the view lacks.
    # No connection count comes from it.
    TRANSMITTER_TABLE = "synapses_nt_v1"
    # Root ids per query.
    # query_view silently stops at ROW_CAP rows,
    # so ids go in batches and a batch that hits the cap is split in
    # half and retried.
    # Fifty is arbitrary; a few hundred also works for LC-sized cells,
    # while one id per call is slow.
    BATCH = 50
    ROW_CAP = 500_000
    # Seconds a query may take before it fails and is retried.
    # One batch of fifty cells answers in about ten seconds,
    # and the largest queries a few times that;
    # five minutes is arbitrary but far above both.
    QUERY_TIMEOUT = 300

TimeoutAdapter

Bases: HTTPAdapter

Gives every request a timeout.

requests has none by default, and a connection that dies mid-query (a laptop suspended during a long fetch) then waits forever.

Source code in champalimaud/fetch.py
class TimeoutAdapter(HTTPAdapter):
    """Gives every request a timeout.

    requests has none by default, and a connection that dies mid-query
    (a laptop suspended during a long fetch) then waits forever.
    """

    def send(self, request, *args, **kwargs):
        """Send a request, with a default timeout."""
        if kwargs.get("timeout") is None:
            kwargs["timeout"] = Sources.QUERY_TIMEOUT
        return super().send(request, *args, **kwargs)

send(request, *args, **kwargs)

Send a request, with a default timeout.

Source code in champalimaud/fetch.py
def send(self, request, *args, **kwargs):
    """Send a request, with a default timeout."""
    if kwargs.get("timeout") is None:
        kwargs["timeout"] = Sources.QUERY_TIMEOUT
    return super().send(request, *args, **kwargs)

batch_connections(outputs, inputs, proofread)

Count the pair weights of one batch of cells.

Parameters:

Name Type Description Default
outputs DataFrame

Synapses whose presynaptic cell is in the batch.

required
inputs DataFrame

Synapses whose postsynaptic cell is in the batch. Those from a proofread cell are dropped, since that cell's own batch counts them.

required
proofread Series

Root ids of every proofread neuron.

required

Returns:

Type Description
DataFrame

Columns pre_pt_root_id, post_pt_root_id, and weight (UInt32), one row per pair of cells, self-contacts dropped.

Examples:

>>> import polars as pl
>>> outputs = pl.DataFrame(
...     {"pre_pt_root_id": [2, 2, 2], "post_pt_root_id": [1, 9, 2]}
... )
>>> inputs = pl.DataFrame(
...     {"pre_pt_root_id": [1, 9], "post_pt_root_id": [2, 2]}
... )
>>> batch_connections(outputs, inputs, pl.Series([1, 2])).rows()
[(2, 1, 1), (2, 9, 1), (9, 2, 1)]
Source code in champalimaud/fetch.py
def batch_connections(
    outputs: pl.DataFrame, inputs: pl.DataFrame, proofread: pl.Series
) -> pl.DataFrame:
    """Count the pair weights of one batch of cells.

    Parameters
    ----------
    outputs : polars.DataFrame
        Synapses whose presynaptic cell is in the batch.
    inputs : polars.DataFrame
        Synapses whose postsynaptic cell is in the batch.
        Those from a proofread cell are dropped, since that cell's own
        batch counts them.
    proofread : polars.Series
        Root ids of every proofread neuron.

    Returns
    -------
    polars.DataFrame
        Columns ``pre_pt_root_id``, ``post_pt_root_id``, and
        ``weight`` (UInt32), one row per pair of cells, self-contacts
        dropped.

    Examples
    --------
    >>> import polars as pl
    >>> outputs = pl.DataFrame(
    ...     {"pre_pt_root_id": [2, 2, 2], "post_pt_root_id": [1, 9, 2]}
    ... )
    >>> inputs = pl.DataFrame(
    ...     {"pre_pt_root_id": [1, 9], "post_pt_root_id": [2, 2]}
    ... )
    >>> batch_connections(outputs, inputs, pl.Series([1, 2])).rows()
    [(2, 1, 1), (2, 9, 1), (9, 2, 1)]
    """
    inputs_from_fragments = inputs.filter(
        ~pl.col("pre_pt_root_id").is_in(proofread.implode())
    )
    return (
        pl.concat([outputs, inputs_from_fragments])
        .filter(pl.col("pre_pt_root_id") != pl.col("post_pt_root_id"))
        .group_by("pre_pt_root_id", "post_pt_root_id")
        .agg(weight=pl.len().cast(pl.UInt32))
        .sort("pre_pt_root_id", "post_pt_root_id")
    )

cave_materialize()

Create the CAVE materialize client for Sources.STACK.

Every request carries a timeout, from TimeoutAdapter. Needs a saved CAVE token.

Returns:

Type Description
MaterializationClient

The client.

Source code in champalimaud/fetch.py
def cave_materialize():
    """Create the CAVE materialize client for `Sources.STACK`.

    Every request carries a timeout, from `TimeoutAdapter`.
    Needs a saved CAVE token.

    Returns
    -------
    caveclient.materializationengine.MaterializationClient
        The client.
    """
    # caveclient types CAVEclient() as its global client, which lacks
    # `materialize`; the attribute exists on the datastack client.
    client = CAVEclient(Sources.STACK)
    materialize = client.materialize  # ty: ignore[unresolved-attribute]
    materialize.session.mount("https://", TimeoutAdapter())
    return materialize

check_inventory(m)

Read the materialization versions and the table names.

Parameters:

Name Type Description Default
m MaterializationClient

As cave_materialize returns it.

required

Returns:

Type Description
Inventory

The versions the server lists, the version in use, and the table names.

Examples:

>>> class Stub:
...     version = 783
...     def get_versions(self):
...         return [783]
...     def get_tables(self):
...         return ["nuclei_v1"]
>>> check_inventory(Stub())
Inventory(versions=[783], version=783, tables=['nuclei_v1'])
Source code in champalimaud/fetch.py
def check_inventory(m) -> Inventory:
    """Read the materialization versions and the table names.

    Parameters
    ----------
    m : caveclient.materializationengine.MaterializationClient
        As `cave_materialize` returns it.

    Returns
    -------
    Inventory
        The versions the server lists, the version in use, and the
        table names.

    Examples
    --------
    >>> class Stub:
    ...     version = 783
    ...     def get_versions(self):
    ...         return [783]
    ...     def get_tables(self):
    ...         return ["nuclei_v1"]
    >>> check_inventory(Stub())
    Inventory(versions=[783], version=783, tables=['nuclei_v1'])
    """
    return Inventory(
        versions=m.get_versions(), version=m.version, tables=m.get_tables()
    )

codex_url(product, api_token)

Build the Codex download URL for one product.

Parameters:

Name Type Description Default
product str

One of Sources.CODEX_PRODUCTS.

required
api_token str

The Codex account token.

required

Returns:

Type Description
str

The URL, in Sources.CODEX_DATASET.

Examples:

>>> codex_url("classification", "<token>")
'https://codex.flywire.ai/api/download_resource?data_product=classification&dataset=fafb&api_token=<token>'
Source code in champalimaud/fetch.py
def codex_url(product: str, api_token: str) -> str:
    """Build the Codex download URL for one product.

    Parameters
    ----------
    product : str
        One of `Sources.CODEX_PRODUCTS`.
    api_token : str
        The Codex account token.

    Returns
    -------
    str
        The URL, in `Sources.CODEX_DATASET`.

    Examples
    --------
    >>> codex_url("classification", "<token>")
    'https://codex.flywire.ai/api/download_resource?data_product=classification&dataset=fafb&api_token=<token>'
    """
    url = Sources.CODEX_HOST
    url += f"?data_product={product}"
    url += f"&dataset={Sources.CODEX_DATASET}"
    url += f"&api_token={api_token}"
    return url

download_codex(product, dest=None)

Write one Codex product to a file.

Needs Sources.CODEX_ENV in the environment.

Parameters:

Name Type Description Default
product str

One of Sources.CODEX_PRODUCTS.

required
dest Path

Where to write; by default config.CODEX/<product>.csv.gz.

None

Returns:

Type Description
Path

The file written.

Raises:

Type Description
SystemExit

When Sources.CODEX_ENV is not set.

Source code in champalimaud/fetch.py
def download_codex(product: str, dest=None):
    """Write one Codex product to a file.

    Needs `Sources.CODEX_ENV` in the environment.

    Parameters
    ----------
    product : str
        One of `Sources.CODEX_PRODUCTS`.
    dest : pathlib.Path, optional
        Where to write; by default ``config.CODEX/<product>.csv.gz``.

    Returns
    -------
    pathlib.Path
        The file written.

    Raises
    ------
    SystemExit
        When `Sources.CODEX_ENV` is not set.
    """
    token = os.environ.get(Sources.CODEX_ENV, "").strip()
    if not token:
        msg = f"set {Sources.CODEX_ENV} (codex.flywire.ai account token)"
        raise SystemExit(msg)
    if dest is None:
        dest = config.CODEX / f"{product}.csv.gz"
    dest.parent.mkdir(parents=True, exist_ok=True)
    urllib.request.urlretrieve(codex_url(product, token), dest)
    return dest

drop_self_edges(syn)

Remove synapses whose pre and post root ids are the same cell.

Parameters:

Name Type Description Default
syn DataFrame

Has pre_pt_root_id and post_pt_root_id.

required

Returns:

Type Description
DataFrame

The rows of syn with different ids at the two ends.

Examples:

>>> import pandas as pd
>>> syn = pd.DataFrame(
...     {"pre_pt_root_id": [1, 2], "post_pt_root_id": [1, 3]}
... )
>>> drop_self_edges(syn)["pre_pt_root_id"].tolist()
[2]
Source code in champalimaud/fetch.py
def drop_self_edges(syn):
    """Remove synapses whose pre and post root ids are the same cell.

    Parameters
    ----------
    syn : pandas.DataFrame
        Has ``pre_pt_root_id`` and ``post_pt_root_id``.

    Returns
    -------
    pandas.DataFrame
        The rows of `syn` with different ids at the two ends.

    Examples
    --------
    >>> import pandas as pd
    >>> syn = pd.DataFrame(
    ...     {"pre_pt_root_id": [1, 2], "post_pt_root_id": [1, 3]}
    ... )
    >>> drop_self_edges(syn)["pre_pt_root_id"].tolist()
    [2]
    """
    return syn[syn["pre_pt_root_id"] != syn["post_pt_root_id"]]

fetch_connection_shard(ids, proofread, dest)

Fetch the connections of some cells and write their pair weights.

Parameters:

Name Type Description Default
ids list of int

The cells of this batch.

required
proofread Series

Root ids of every proofread neuron.

required
dest Path

Parquet file to write; it is written under a temporary name first, so a fetch that dies halfway leaves no truncated file.

required

Returns:

Type Description
int

The number of pairs written.

Raises:

Type Description
RuntimeError

From query_synapses, at once. Any other failure is retried up to five times, with smaller queries from the third attempt.

Source code in champalimaud/fetch.py
def fetch_connection_shard(
    ids: list[int], proofread: pl.Series, dest: Path
) -> int:
    """Fetch the connections of some cells and write their pair weights.

    Parameters
    ----------
    ids : list of int
        The cells of this batch.
    proofread : polars.Series
        Root ids of every proofread neuron.
    dest : pathlib.Path
        Parquet file to write; it is written under a temporary name
        first, so a fetch that dies halfway leaves no truncated file.

    Returns
    -------
    int
        The number of pairs written.

    Raises
    ------
    RuntimeError
        From `query_synapses`, at once.
        Any other failure is retried up to five times, with smaller
        queries from the third attempt.
    """
    if not hasattr(CLIENTS, "materialize"):
        CLIENTS.materialize = cave_materialize()
    # A gateway timeout (502, 503) usually means the query was too
    # heavy for the server, so the third attempt asks for ten cells at
    # a time and the fourth for one.
    frames: dict[str, pl.DataFrame] = {}
    for attempt in range(1, 6):
        size = len(ids) if attempt <= 2 else 10 if attempt == 3 else 1
        try:
            frames = {}
            for column in ("pre_pt_root_id", "post_pt_root_id"):
                answers = [
                    query_synapses(
                        CLIENTS.materialize, column, ids[i : i + size]
                    )
                    for i in range(0, len(ids), size)
                ]
                rows = pd.concat(answers, ignore_index=True)
                # Built from numpy so an empty answer still has integer
                # columns.
                frames[column] = pl.DataFrame(
                    {
                        name: rows[name].to_numpy(dtype="int64")
                        for name in ("pre_pt_root_id", "post_pt_root_id")
                    }
                )
            break
        except RuntimeError:
            raise
        except Exception as e:
            # A dropped connection or a server hiccup;
            # if every attempt fails,
            # a rerun resumes from the batches on disk.
            if attempt == 5:
                raise
            print(
                f"retry {attempt} after {type(e).__name__}: {e}",
                file=sys.stderr,
            )
            time.sleep(15 * attempt)
    pairs = batch_connections(
        frames["pre_pt_root_id"], frames["post_pt_root_id"], proofread
    )
    # Written under a temporary name so a fetch that dies halfway
    # leaves no truncated file for the next run to skip.
    part = dest.with_name(dest.name + ".part")
    pairs.write_parquet(part)
    part.replace(dest)
    return pairs.height

fetch_connections(workers=4)

Fetch every proofread cell's connections and merge them.

The connections go in batches of Sources.BATCH cells to config.CONNECTIONS_PARTS, and a rerun skips the batches that are there; the merge writes config.CONNECTIONS and deletes the folder.

Parameters:

Name Type Description Default
workers int

Batches fetched at once.

4

Raises:

Type Description
SystemExit

When any batch failed; rerun to fetch the missing ones.

Source code in champalimaud/fetch.py
def fetch_connections(workers: int = 4) -> None:
    """Fetch every proofread cell's connections and merge them.

    The connections go in batches of `Sources.BATCH` cells to
    ``config.CONNECTIONS_PARTS``, and a rerun skips the batches that
    are there; the merge writes ``config.CONNECTIONS`` and deletes the
    folder.

    Parameters
    ----------
    workers : int, default 4
        Batches fetched at once.

    Raises
    ------
    SystemExit
        When any batch failed; rerun to fetch the missing ones.
    """
    proofread = load_proofread_ids()["root_id"].sort()
    ids = proofread.to_list()
    batches = [
        ids[i : i + Sources.BATCH] for i in range(0, len(ids), Sources.BATCH)
    ]
    config.CONNECTIONS_PARTS.mkdir(parents=True, exist_ok=True)

    def shard_path(k: int) -> Path:
        return config.CONNECTIONS_PARTS / f"part-{k:05d}.parquet"

    todo = [k for k in range(len(batches)) if not shard_path(k).exists()]
    print(
        f"{len(batches) - len(todo)} of {len(batches)} batches on disk; "
        f"fetching {len(todo)} with {workers} workers",
        flush=True,
    )
    failed = []
    started = time.time()
    with ThreadPoolExecutor(max_workers=workers) as pool:
        futures = {
            pool.submit(
                fetch_connection_shard, batches[k], proofread, shard_path(k)
            ): k
            for k in todo
        }
        for done, future in enumerate(as_completed(futures), start=1):
            k = futures[future]
            try:
                pairs = future.result()
            except Exception as e:  # noqa: BLE001
                # Keep going; the failed batches are reported and a
                # rerun fetches them.
                failed.append(k)
                print(
                    f"batch {k + 1} failed: {type(e).__name__}: {e}",
                    flush=True,
                )
                continue
            minutes = (time.time() - started) / 60
            print(
                f"batch {k + 1}/{len(batches)}: {pairs:,} pairs; "
                f"{done}/{len(todo)} fetched in {minutes:.0f} min",
                flush=True,
            )
    if failed:
        msg = (
            f"{len(failed)} batches failed (first: {sorted(failed)[:10]}); "
            f"rerun fetch_connections() to fetch them"
        )
        raise SystemExit(msg)
    merged = config.CONNECTIONS.with_name(config.CONNECTIONS.name + ".part")
    (
        pl.scan_parquet(config.CONNECTIONS_PARTS / "part-*.parquet")
        .sort("pre_pt_root_id", "post_pt_root_id")
        .sink_parquet(merged, row_group_size=1_000_000)
    )
    merged.replace(config.CONNECTIONS)
    shutil.rmtree(config.CONNECTIONS_PARTS)
    total = (
        pl.scan_parquet(config.CONNECTIONS).select(pl.len()).collect().item()
    )
    print(f"wrote {config.CONNECTIONS}: {total:,} pairs")

fetch_example_synapses(cell_type=EXAMPLE_TYPE)

Write every synapse of the first cell of a type.

The rows are those of Sources.SYNAPSE_VIEW, the table that fetch_connections counts, with the synapses of the cell onto itself dropped.

Parameters:

Name Type Description Default
cell_type str

The primary_type whose first census cell is fetched.

`EXAMPLE_TYPE`

Returns:

Type Description
Path

config.EXAMPLE_SYNAPSES, with columns cell (the root id fetched), id, pre_pt_root_id, and post_pt_root_id.

Source code in champalimaud/fetch.py
def fetch_example_synapses(cell_type: str = EXAMPLE_TYPE) -> Path:
    """Write every synapse of the first cell of a type.

    The rows are those of `Sources.SYNAPSE_VIEW`, the table that
    `fetch_connections` counts, with the synapses of the cell onto
    itself dropped.

    Parameters
    ----------
    cell_type : str, default `EXAMPLE_TYPE`
        The ``primary_type`` whose first census cell is fetched.

    Returns
    -------
    pathlib.Path
        ``config.EXAMPLE_SYNAPSES``, with columns ``cell`` (the root
        id fetched), ``id``, ``pre_pt_root_id``, and
        ``post_pt_root_id``.
    """
    root_id = root_ids_of_type(load_census(), cell_type)[0]
    m = cave_materialize()
    rows = pd.concat(
        [
            query_synapses(m, column, [root_id])
            for column in ("pre_pt_root_id", "post_pt_root_id")
        ],
        ignore_index=True,
    )
    rows = drop_self_edges(rows).drop_duplicates("id")
    config.EXAMPLE_SYNAPSES.parent.mkdir(parents=True, exist_ok=True)
    part = config.EXAMPLE_SYNAPSES.with_name(
        config.EXAMPLE_SYNAPSES.name + ".part"
    )
    pl.from_pandas(rows).with_columns(cell=pl.lit(root_id)).select(
        "cell", "id", "pre_pt_root_id", "post_pt_root_id"
    ).write_parquet(part)
    part.replace(config.EXAMPLE_SYNAPSES)
    print(f"wrote {len(rows):,} synapses of {root_id} to {part.parent}")
    return config.EXAMPLE_SYNAPSES

fetch_skeletons(types=SKELETON_TYPES)

Write the first skeletons of each type to config.SKELETONS.

Parameters:

Name Type Description Default
types sequence of str

The primary_type values to fetch; the first SKELETONS_PER_TYPE cells of each.

`SKELETON_TYPES`

Returns:

Type Description
list of pathlib.Path

The <root_id>.swc files written; a skeleton that fails to download is skipped.

Source code in champalimaud/fetch.py
def fetch_skeletons(types=SKELETON_TYPES):
    """Write the first skeletons of each type to ``config.SKELETONS``.

    Parameters
    ----------
    types : sequence of str, default `SKELETON_TYPES`
        The ``primary_type`` values to fetch; the first
        `SKELETONS_PER_TYPE` cells of each.

    Returns
    -------
    list of pathlib.Path
        The ``<root_id>.swc`` files written; a skeleton that fails to
        download is skipped.
    """
    census = load_census()
    config.SKELETONS.mkdir(parents=True, exist_ok=True)
    written = []
    for cell_type in types:
        for root_id in root_ids_of_type(census, cell_type)[
            :SKELETONS_PER_TYPE
        ]:
            path = config.SKELETONS / f"{root_id}.swc"
            # Download under a temporary name so a fetch that dies
            # halfway leaves no truncated SWC behind to be read as data.
            part = path.with_name(path.name + ".part")
            try:
                urllib.request.urlretrieve(skeleton_url(root_id), part)
            except Exception as e:  # noqa: BLE001
                # One missing skeleton should not stop the rest.
                part.unlink(missing_ok=True)
                print(
                    f"skip {cell_type} {root_id}: {type(e).__name__}: {e}",
                    file=sys.stderr,
                )
                continue
            part.replace(path)
            written.append(path)
    print(f"wrote {len(written)} skeletons to {config.SKELETONS}")
    return written

ids_in_nuclei(m, root_ids)

Count the root ids that nuclei_v1 has a row for.

Parameters:

Name Type Description Default
m MaterializationClient

As cave_materialize returns it.

required
root_ids sequence of int

The ids to look up.

required

Returns:

Type Description
int

How many of them are in the table.

Examples:

>>> import pandas as pd
>>> class Stub:
...     def query_table(self, table, **kwargs):
...         return pd.DataFrame({"pt_root_id": [1, 1, 2]})
>>> ids_in_nuclei(Stub(), [1, 2, 3])
2
Source code in champalimaud/fetch.py
def ids_in_nuclei(m, root_ids: Sequence[int]) -> int:
    """Count the root ids that ``nuclei_v1`` has a row for.

    Parameters
    ----------
    m : caveclient.materializationengine.MaterializationClient
        As `cave_materialize` returns it.
    root_ids : sequence of int
        The ids to look up.

    Returns
    -------
    int
        How many of them are in the table.

    Examples
    --------
    >>> import pandas as pd
    >>> class Stub:
    ...     def query_table(self, table, **kwargs):
    ...         return pd.DataFrame({"pt_root_id": [1, 1, 2]})
    >>> ids_in_nuclei(Stub(), [1, 2, 3])
    2
    """
    hit = m.query_table(
        "nuclei_v1",
        filter_in_dict={"pt_root_id": list(root_ids)},
        select_columns=["pt_root_id"],
    )
    return int(hit["pt_root_id"].nunique())

query_synapses(m, column, ids)

Query the synapse rows with a column in some root ids.

A query that hits Sources.ROW_CAP is truncated without an error, so it is split in half until each answer is under the cap.

Parameters:

Name Type Description Default
m MaterializationClient

As cave_materialize returns it.

required
column (pre_pt_root_id, post_pt_root_id)

The column to test.

"pre_pt_root_id"
ids list of int

Root ids to look for in column.

required

Returns:

Type Description
DataFrame

Columns id, pre_pt_root_id, and post_pt_root_id; CAVE hands the rows back as a pandas frame.

Raises:

Type Description
RuntimeError

When one root id alone has at least Sources.ROW_CAP synapses in column.

Source code in champalimaud/fetch.py
def query_synapses(m, column: str, ids: list[int]) -> pd.DataFrame:
    """Query the synapse rows with a column in some root ids.

    A query that hits `Sources.ROW_CAP` is truncated without an error,
    so it is split in half until each answer is under the cap.

    Parameters
    ----------
    m : caveclient.materializationengine.MaterializationClient
        As `cave_materialize` returns it.
    column : {"pre_pt_root_id", "post_pt_root_id"}
        The column to test.
    ids : list of int
        Root ids to look for in `column`.

    Returns
    -------
    pandas.DataFrame
        Columns ``id``, ``pre_pt_root_id``, and ``post_pt_root_id``;
        CAVE hands the rows back as a pandas frame.

    Raises
    ------
    RuntimeError
        When one root id alone has at least `Sources.ROW_CAP`
        synapses in `column`.
    """
    rows = m.query_view(
        Sources.SYNAPSE_VIEW,
        filter_in_dict={column: ids},
        select_columns=["id", "pre_pt_root_id", "post_pt_root_id"],
    )
    # A full page means the server truncated the answer; there is no
    # error, so the row count is the only signal.
    if len(rows) >= Sources.ROW_CAP and len(ids) == 1:
        msg = (
            f"root id {ids[0]} has at least {Sources.ROW_CAP:,} synapses "
            f"in {column}; the answer is truncated"
        )
        raise RuntimeError(msg)
    if len(rows) >= Sources.ROW_CAP:
        mid = len(ids) // 2
        first = query_synapses(m, column, ids[:mid])
        second = query_synapses(m, column, ids[mid:])
        rows = pd.concat([first, second], ignore_index=True)
    return rows.loc[:, ["id", "pre_pt_root_id", "post_pt_root_id"]]

sample_tables(m)

Query a few rows of the tables every fetch reads.

Parameters:

Name Type Description Default
m MaterializationClient

As cave_materialize returns it.

required

Returns:

Type Description
nuclei, view_row, synapse : pandas.DataFrame

Three rows of nuclei_v1 (a cell body, with its position split into x, y, and z), one row of Sources.SYNAPSE_VIEW, and one row of Sources.TRANSMITTER_TABLE.

Examples:

>>> import pandas as pd
>>> class Stub:
...     def query_table(self, table, **kwargs):
...         return pd.DataFrame({"table": [table]})
...     def query_view(self, view, **kwargs):
...         return pd.DataFrame({"table": [view]})
>>> nuclei, view_row, synapse = sample_tables(Stub())
>>> nuclei["table"].tolist()
['nuclei_v1']
Source code in champalimaud/fetch.py
def sample_tables(m) -> tuple[pd.DataFrame, pd.DataFrame, pd.DataFrame]:
    """Query a few rows of the tables every fetch reads.

    Parameters
    ----------
    m : caveclient.materializationengine.MaterializationClient
        As `cave_materialize` returns it.

    Returns
    -------
    nuclei, view_row, synapse : pandas.DataFrame
        Three rows of ``nuclei_v1`` (a cell body, with its position
        split into x, y, and z), one row of `Sources.SYNAPSE_VIEW`,
        and one row of `Sources.TRANSMITTER_TABLE`.

    Examples
    --------
    >>> import pandas as pd
    >>> class Stub:
    ...     def query_table(self, table, **kwargs):
    ...         return pd.DataFrame({"table": [table]})
    ...     def query_view(self, view, **kwargs):
    ...         return pd.DataFrame({"table": [view]})
    >>> nuclei, view_row, synapse = sample_tables(Stub())
    >>> nuclei["table"].tolist()
    ['nuclei_v1']
    """
    nuclei = m.query_table(
        "nuclei_v1",
        select_columns=["pt_root_id", "pt_position"],
        limit=3,
        split_positions=True,
    )
    view_row = m.query_view(Sources.SYNAPSE_VIEW, limit=1)
    synapse = m.query_table(Sources.TRANSMITTER_TABLE, limit=1)
    return nuclei, view_row, synapse

skeleton_url(root_id)

Build the Codex SWC URL for one root id.

This is the only skeleton source besides the bulk zip on their download page.

Parameters:

Name Type Description Default
root_id int

The cell.

required

Returns:

Type Description
str

The URL under Sources.SKELETON_URL.

Examples:

>>> skeleton_url(123).rsplit("/", 1)[1]
'123.swc'
Source code in champalimaud/fetch.py
def skeleton_url(root_id) -> str:
    """Build the Codex SWC URL for one root id.

    This is the only skeleton source besides the bulk zip on their
    download page.

    Parameters
    ----------
    root_id : int
        The cell.

    Returns
    -------
    str
        The URL under `Sources.SKELETON_URL`.

    Examples
    --------
    >>> skeleton_url(123).rsplit("/", 1)[1]
    '123.swc'
    """
    return f"{Sources.SKELETON_URL}/{root_id}.swc"