diff --git a/Dockerfile b/Dockerfile index 1b623a8c9..fd23caaee 100644 --- a/Dockerfile +++ b/Dockerfile @@ -15,8 +15,6 @@ RUN apt-get update && apt-get install build-essential wget -y \ && rm Miniconda3-latest-Linux-x86_64.sh \ && conda config --add channels conda-forge \ && conda update -y --override-channels -c conda-forge conda \ - && conda tos accept --override-channels --channel https://repo.anaconda.com/pkgs/main \ - && conda tos accept --override-channels --channel https://repo.anaconda.com/pkgs/r \ && conda install -y --override-channels -c conda-forge conda-pack COPY requirements.yml requirements.txt requirements-dev.txt ./ diff --git a/pychunkedgraph/graph/chunkedgraph.py b/pychunkedgraph/graph/chunkedgraph.py index 210bff50b..3916fc94e 100644 --- a/pychunkedgraph/graph/chunkedgraph.py +++ b/pychunkedgraph/graph/chunkedgraph.py @@ -1020,3 +1020,8 @@ def get_earliest_timestamp(self): _, timestamp = self.client.read_log_entry(op_id) if timestamp is not None: return timestamp - timedelta(milliseconds=500) + # no ops: the ingest-completion boundary stamped during the root-layer build + stamped = self.meta.custom_data.get("earliest_ts") + if stamped is not None: + return datetime.datetime.fromisoformat(stamped) + return datetime.datetime.fromtimestamp(0, tz=datetime.timezone.utc) diff --git a/pychunkedgraph/ingest/create/abstract_layers.py b/pychunkedgraph/ingest/create/abstract_layers.py index 529a6846f..6994bbcc4 100644 --- a/pychunkedgraph/ingest/create/abstract_layers.py +++ b/pychunkedgraph/ingest/create/abstract_layers.py @@ -63,15 +63,25 @@ def add_layer( graph, _, _, graph_ids = flatgraph.build_gt_graph(edge_ids, make_directed=True) ccs = flatgraph.connected_components(graph) print("ccs", len(ccs)) + ts = get_valid_timestamp(time_stamp) _write_connected_components( cg, layer_id, parent_coords, ccs, graph_ids, - get_valid_timestamp(time_stamp), + ts, n_threads > 1, ) + + # Stamp the post-ingest boundary meshing reads to split initial from edited roots. + # ts is the explicit cell timestamp shared by every root just written; +500ms (the + # same guard get_earliest_timestamp puts below the first op) lifts the boundary + # strictly above them. + if layer_id == cg.meta.layer_count: + boundary = ts + datetime.timedelta(milliseconds=500) + cg.meta.custom_data["earliest_ts"] = boundary.isoformat() + cg.update_meta(cg.meta, overwrite=True) return f"{layer_id}_{'_'.join(map(str, parent_coords))}" @@ -166,7 +176,6 @@ def _write_connected_components( ccs_with_node_ids, node_layer_d_shared, time_stamp, - use_threads=use_threads, ) return @@ -192,9 +201,7 @@ def _write_components_helper(args): _write(cg, layer_id, parent_coords, ccs, node_layer_d_shared, time_stamp) -def _write( - cg, layer_id, parent_coords, ccs, node_layer_d_shared, time_stamp, use_threads=True -): +def _write(cg, layer_id, parent_coords, ccs, node_layer_d_shared, time_stamp): parent_layer_ids = range(layer_id, cg.meta.layer_count + 1) cc_connections = {l: [] for l in parent_layer_ids} for node_ids in ccs: @@ -217,7 +224,7 @@ def _write( reserved_parent_ids = cg.id_client.create_node_ids( parent_chunk_id, size=len(cc_connections[parent_layer_id]), - root_chunk=parent_layer_id == cg.meta.layer_count and use_threads, + root_chunk=parent_layer_id == cg.meta.layer_count, ) for i_cc, node_ids in enumerate(cc_connections[parent_layer_id]): diff --git a/pychunkedgraph/ingest/simple_tests.py b/pychunkedgraph/ingest/simple_tests.py new file mode 100644 index 000000000..46a13b108 --- /dev/null +++ b/pychunkedgraph/ingest/simple_tests.py @@ -0,0 +1,176 @@ +# pylint: disable=invalid-name, missing-function-docstring, broad-exception-caught + +""" +Some sanity tests to ensure chunkedgraph was created properly. +""" + +from datetime import datetime, timezone +import numpy as np + +from pychunkedgraph.graph import ChunkedGraph +from pychunkedgraph.graph import attributes + + +def family(cg: ChunkedGraph): + np.random.seed(42) + n_chunks = 100 + n_segments_per_chunk = 200 + timestamp = datetime.now(timezone.utc) + + node_ids = [] + for layer in range(2, cg.meta.layer_count - 1): + for _ in range(n_chunks): + c_x = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][0]) + c_y = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][1]) + c_z = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][2]) + chunk_id = cg.get_chunk_id(layer=layer, x=c_x, y=c_y, z=c_z) + max_segment_id = cg.get_segment_id(cg.id_client.get_max_node_id(chunk_id)) + if max_segment_id < 10: + continue + + segment_ids = np.random.randint(1, max_segment_id, n_segments_per_chunk) + for segment_id in segment_ids: + node_ids.append( + cg.get_node_id(np.uint64(segment_id), np.uint64(chunk_id)) + ) + + rows = cg.client.read_nodes( + node_ids=node_ids, end_time=timestamp, properties=attributes.Hierarchy.Parent + ) + valid_node_ids = [] + non_valid_node_ids = [] + for k in rows.keys(): + if len(rows[k]) > 0: + valid_node_ids.append(k) + else: + non_valid_node_ids.append(k) + + parents = cg.get_parents(valid_node_ids, time_stamp=timestamp) + children_dict = cg.get_children(parents) + for child, parent in zip(valid_node_ids, parents): + assert child in children_dict[parent] + print("success") + + +def existence(cg: ChunkedGraph): + np.random.seed(42) + layer = 2 + n_chunks = 100 + n_segments_per_chunk = 200 + timestamp = datetime.now(timezone.utc) + node_ids = [] + for _ in range(n_chunks): + c_x = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][0]) + c_y = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][1]) + c_z = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][2]) + chunk_id = cg.get_chunk_id(layer=layer, x=c_x, y=c_y, z=c_z) + max_segment_id = cg.get_segment_id(cg.id_client.get_max_node_id(chunk_id)) + if max_segment_id < 10: + continue + + segment_ids = np.random.randint(1, max_segment_id, n_segments_per_chunk) + for segment_id in segment_ids: + node_ids.append(cg.get_node_id(np.uint64(segment_id), np.uint64(chunk_id))) + + rows = cg.client.read_nodes( + node_ids=node_ids, end_time=timestamp, properties=attributes.Hierarchy.Parent + ) + valid_node_ids = [] + non_valid_node_ids = [] + for k in rows.keys(): + if len(rows[k]) > 0: + valid_node_ids.append(k) + else: + non_valid_node_ids.append(k) + + roots = [] + try: + roots = cg.get_roots(valid_node_ids) + assert len(roots) == len(valid_node_ids) + print("success") + except Exception as e: + print(f"Something went wrong: {e}") + print("At least one node failed. Checking nodes one by one:") + + if len(roots) != len(valid_node_ids): + log_dict = {} + success_dict = {} + for node_id in valid_node_ids: + try: + _ = cg.get_root(node_id, time_stamp=timestamp) + print(f"Success: {node_id} from chunk {cg.get_chunk_id(node_id)}") + success_dict[node_id] = True + except Exception as e: + print(f"{node_id} - chunk {cg.get_chunk_id(node_id)} failed: {e}") + success_dict[node_id] = False + t_id = node_id + while t_id is not None: + last_working_chunk = cg.get_chunk_id(t_id) + t_id = cg.get_parent(t_id) + + layer = cg.get_chunk_layer(last_working_chunk) + print(f"Failed on layer {layer} in chunk {last_working_chunk}") + log_dict[node_id] = last_working_chunk + if log_dict: # diagnostics above are informational; the suite must still fail + raise AssertionError(f"{len(log_dict)} nodes failed the existence check") + + +def cross_edges(cg: ChunkedGraph): + np.random.seed(42) + layer = 2 + n_chunks = 10 + n_segments_per_chunk = 200 + timestamp = datetime.now(timezone.utc) + node_ids = [] + for _ in range(n_chunks): + c_x = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][0]) + c_y = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][1]) + c_z = np.random.randint(0, cg.meta.layer_chunk_bounds[layer][2]) + chunk_id = cg.get_chunk_id(layer=layer, x=c_x, y=c_y, z=c_z) + max_segment_id = cg.get_segment_id(cg.id_client.get_max_node_id(chunk_id)) + if max_segment_id < 10: + continue + + segment_ids = np.random.randint(1, max_segment_id, n_segments_per_chunk) + for segment_id in segment_ids: + node_ids.append(cg.get_node_id(np.uint64(segment_id), np.uint64(chunk_id))) + + rows = cg.client.read_nodes( + node_ids=node_ids, end_time=timestamp, properties=attributes.Hierarchy.Parent + ) + valid_node_ids = [] + non_valid_node_ids = [] + for k in rows.keys(): + if len(rows[k]) > 0: + valid_node_ids.append(k) + else: + non_valid_node_ids.append(k) + + cc_edges = cg.get_atomic_cross_edges(valid_node_ids) + edges = [np.concatenate(list(v.values())) for v in cc_edges.values() if len(v)] + if not edges: + print("no cross edges in the sampled chunks; nothing to check") + return + cc_ids = np.unique(np.concatenate(edges)) + + roots = cg.get_roots(cc_ids) + root_dict = dict(zip(cc_ids, roots)) + root_dict_vec = np.vectorize(root_dict.get) + + for k in cc_edges: + if len(cc_edges[k]) == 0: + continue + local_ids = np.unique(np.concatenate(list(cc_edges[k].values()))) + assert len(np.unique(root_dict_vec(local_ids))) + print("success") + + +def run_all(cg: ChunkedGraph): + print("Running family tests:") + family(cg) + + print("\nRunning existence tests:") + existence(cg) + + print("\nRunning cross_edges tests:") + cross_edges(cg) diff --git a/pychunkedgraph/pipeline/README.md b/pychunkedgraph/pipeline/README.md index 226bacde4..6ee1ffcb7 100644 --- a/pychunkedgraph/pipeline/README.md +++ b/pychunkedgraph/pipeline/README.md @@ -38,6 +38,7 @@ typically need tuning between layers). Nothing auto-advances. | Path | Responsibility | |---|---| +| `__init__.py` | `run_and_exit(main)`: every entrypoint's `__main__` guard hands `main` to it, which leaves through `os._exit` so a client's non-daemon thread never stalls the pod, as CAVEpipelines' image contract requires. | | `grid.py` | Fixed-seed permutation: maps a batch's contiguous index window to *scattered* chunk coords so concurrent workers spread Bigtable row-key load instead of hot-spotting one tablet. Deterministic + invertible. | | `exit_codes.py` | Map success / transient / non-transient failure to the Job `podFailurePolicy`. | | `lock.py` | Per-chunk Bigtable claim/done cell (atomic CAS). One effective writer per chunk, token-fenced; a dead holder's claim expires so a retry re-claims; already-`done` chunks are skipped. Used by ingest. | @@ -74,7 +75,7 @@ The Job template sets these on each pod: | `PCG_LAYER` | layer being built | | `PCG_PERM_SEED` | permutation seed — **same across all pods and retries of a run** | | `PCG_BATCH_SIZE` | `B`, chunks per index | -| `PCG_N_THREADS` | parallel sub-workers inside a parent-chunk build (default 1) | +| `PCG_N_PROCESSES` | parallel sub-workers inside a parent-chunk build (default 1); above 1, pools size from the node's cores | | `PCG_LOCK_EXPIRY_SCALE` | (ingest) scales the per-layer claim TTL; default 1 | | `PCG_LOCK_POLL_SEC` / `PCG_HELD_MAX_WAIT_SEC` | (ingest) poll interval / max wait before deferring a held chunk | | `PCG_MESH_CACHE` | (meshing) `0` disables the mesh task's cloud cache; default on | @@ -95,8 +96,13 @@ re-claims only the unfinished. Meshing needs no lock — it overwrites shards id budget); the index is retried. - **Transient failure**: counts toward the bounded per-index retry budget; the batch retries (ingest skips done chunks). -- **Non-transient failure** (`FatalChunkError`, exit 42): the index is failed fast and - recorded for inspection. A batch finishes all its chunks before choosing an exit code. +- **Non-transient failure** (exit 42): the index is failed fast and recorded for inspection. + Ingest raises `FatalChunkError`; meshing reads a bug or bad input off the exception type, + so one deterministic fault cannot burn every retry an index has. A batch finishes all its + chunks before choosing an exit code. +- **Root verify** (ingest): after the root chunk is built, the pod runs the hierarchy + sanity suite (`ingest.simple_tests`). A failed check fails the pod without re-opening the + chunk, so re-submitting the root layer re-runs only this check, never the build. ## Testing diff --git a/pychunkedgraph/pipeline/__init__.py b/pychunkedgraph/pipeline/__init__.py index f85f01871..6b5757f38 100644 --- a/pychunkedgraph/pipeline/__init__.py +++ b/pychunkedgraph/pipeline/__init__.py @@ -6,3 +6,29 @@ scattered chunk coords, and processes each chunk. Self-contained and branch-portable (main + pcgv3). Run a workload as a module, e.g. ``python -m pychunkedgraph.pipeline.ingest``. """ + +import os +import sys +import traceback + + +def run_and_exit(main) -> None: + """Run a pipeline entrypoint, then ``os._exit``; every container and one-shot entry goes + through this. A normal return can hang joining a client's non-daemon thread at exit, and + the pod then stalls until SIGKILL. The exit code is main()'s return (0 if None), the + SystemExit code, or 1 with a printed traceback on any unhandled error.""" + try: + code = main() or 0 + except SystemExit as exc: # argparse / explicit exit + code = exc.code + if isinstance(code, str): + print(code, file=sys.stderr) + code = 1 + elif not isinstance(code, int): + code = 0 if code is None else 1 + except BaseException: # report and exit non-zero, never hang + traceback.print_exc() + code = 1 + sys.stdout.flush() + sys.stderr.flush() + os._exit(code) diff --git a/pychunkedgraph/pipeline/ingest/__main__.py b/pychunkedgraph/pipeline/ingest/__main__.py index 77719db2b..a2e2b5fae 100644 --- a/pychunkedgraph/pipeline/ingest/__main__.py +++ b/pychunkedgraph/pipeline/ingest/__main__.py @@ -1,8 +1,7 @@ """Container entrypoint: ``python -m pychunkedgraph.pipeline.ingest``.""" -import sys - +from .. import run_and_exit from .worker import main if __name__ == "__main__": - sys.exit(main()) + run_and_exit(main) diff --git a/pychunkedgraph/pipeline/ingest/setup.py b/pychunkedgraph/pipeline/ingest/setup.py index 3b683af1c..eb518aa43 100644 --- a/pychunkedgraph/pipeline/ingest/setup.py +++ b/pychunkedgraph/pipeline/ingest/setup.py @@ -19,6 +19,7 @@ from ...graph.client import BackendClientInfo from ...graph.client.bigtable import BigTableConfig from ...graph.meta import ChunkedGraphMeta, DataSource, GraphConfig +from .. import run_and_exit # Predetermined mount path of the dataset yaml (the chart mounts the dataset # ConfigMap here); overridable for local/testing. @@ -63,4 +64,4 @@ def main() -> None: if __name__ == "__main__": - main() + run_and_exit(main) diff --git a/pychunkedgraph/pipeline/ingest/worker.py b/pychunkedgraph/pipeline/ingest/worker.py index c44751560..afaa9fb6c 100644 --- a/pychunkedgraph/pipeline/ingest/worker.py +++ b/pychunkedgraph/pipeline/ingest/worker.py @@ -11,6 +11,7 @@ import time from datetime import timedelta +from ...ingest import simple_tests from .. import lock from ..exit_codes import FatalChunkError from ..worker import run @@ -122,8 +123,19 @@ def _process_one(table, cg, layer, coord, config, opts) -> str: return "transient" +def _verify_root(cg, layer) -> None: + """Once the root chunk is built, run the hierarchy sanity suite. + + The chunk is already marked done, so a re-submitted root layer re-runs only this + check, never the build. + """ + if layer != cg.meta.layer_count: + return + simple_tests.run_all(cg) + + def main() -> int: - return run(make_processor) + return run(make_processor, finalize=_verify_root) if __name__ == "__main__": diff --git a/pychunkedgraph/pipeline/meshing/__main__.py b/pychunkedgraph/pipeline/meshing/__main__.py index bbb100ccd..46f03d8ab 100644 --- a/pychunkedgraph/pipeline/meshing/__main__.py +++ b/pychunkedgraph/pipeline/meshing/__main__.py @@ -1,8 +1,7 @@ """Container entrypoint: ``python -m pychunkedgraph.pipeline.meshing``.""" -import sys - +from .. import run_and_exit from .worker import main if __name__ == "__main__": - sys.exit(main()) + run_and_exit(main) diff --git a/pychunkedgraph/pipeline/meshing/meta.py b/pychunkedgraph/pipeline/meshing/meta.py index 2969ea19a..fe9052a5e 100644 --- a/pychunkedgraph/pipeline/meshing/meta.py +++ b/pychunkedgraph/pipeline/meshing/meta.py @@ -5,7 +5,8 @@ required — the helper that applies it (``setup_mesh_meta``) does not substitute defaults for missing yaml entries. The only optional field is :attr:`dynamic_mesh_dir`, which is graph-id-derived and filled in by -:meth:`with_graph_id` when omitted from the yaml. +:meth:`with_graph_id` when omitted from the yaml. The mesh chunk_size is +derived (CG CHUNK_SIZE / per-axis downsample at ``mip``), not configured. Example yaml block:: @@ -14,13 +15,12 @@ mip: 0 max_layer: 6 max_error: 40 - chunk_size: [512, 512, 256] minishard_bits: {2: 1, 3: 3, 4: 6, 5: 9, 6: 12} # dynamic_mesh_dir: my_custom_dir # optional; default "dynamic_" """ from dataclasses import asdict, dataclass, replace -from typing import Dict, List, Optional +from typing import Dict, Optional @dataclass @@ -31,7 +31,6 @@ class MeshConfig: mip: int max_layer: int max_error: int - chunk_size: List[int] minishard_bits: Dict[int, int] dynamic_mesh_dir: Optional[str] = None @@ -54,8 +53,6 @@ def from_dict(cls, d: Dict) -> "MeshConfig": kwargs["minishard_bits"] = { int(k): int(v) for k, v in kwargs["minishard_bits"].items() } - if "chunk_size" in kwargs: - kwargs["chunk_size"] = [int(x) for x in kwargs["chunk_size"]] return cls(**kwargs) def with_graph_id(self, graph_id: str) -> "MeshConfig": diff --git a/pychunkedgraph/pipeline/meshing/setup.py b/pychunkedgraph/pipeline/meshing/setup.py index fb8951ddd..8753e9000 100644 --- a/pychunkedgraph/pipeline/meshing/setup.py +++ b/pychunkedgraph/pipeline/meshing/setup.py @@ -22,13 +22,15 @@ import argparse import logging +from datetime import datetime, timezone from os import environ -import numpy as np import yaml from ...graph.chunkedgraph import ChunkedGraph from ...meshing.meshgen import get_draco_encoding_settings_for_chunk +from ...meshing.meshgen_utils import get_mesh_block_shape_for_mip +from .. import run_and_exit from .meta import MeshConfig logger = logging.getLogger(__name__) @@ -39,45 +41,19 @@ def derive_initial_ts(cg: ChunkedGraph) -> int: - """Unix-seconds timestamp of a root id sampled from the dataset center. - - ``mesh.initial_ts`` is the threshold ``segregate_node_ids`` (see - ``meshing/manifest/utils.py``) uses to classify root ids as initial - vs post-ingest. It must sit above the last initial-ingest commit - and below any post-ingest commit. Picking a root id near the - volume center and using its commit timestamp satisfies both bounds - for any graph that completed initial ingest. - - Walks shells outward from the center of the L2 chunk grid (L1 - shares L2's coordinate grid) and returns the timestamp of the root - of the first SV found. + """Unix-seconds boundary for ``mesh.initial_ts`` (see ``segregate_node_ids``). + + ``get_earliest_timestamp`` returns the first edit, or — pre-edit — the + ingest-completion boundary stamped during the root-layer build. ``+1`` makes the + second-granularity threshold strictly above the last initial root (the check + is ``<`` and ``int()`` truncates). """ - hi = np.asarray(cg.meta.layer_chunk_bounds[2]) - center = hi // 2 - for r in range(int(hi.max()) + 1): - box = ( - np.array( - np.meshgrid( - np.arange(-r, r + 1), - np.arange(-r, r + 1), - np.arange(-r, r + 1), - indexing="ij", - ) - ) - .reshape(3, -1) - .T + earliest = cg.get_earliest_timestamp() + if earliest <= datetime.fromtimestamp(0, tz=timezone.utc): + raise RuntimeError( + "derive_initial_ts: no operations and no ingest earliest_ts stamped" ) - shell = box[np.max(np.abs(box), axis=1) == r] - coords = np.unique(np.clip(center + shell, 0, hi - 1), axis=0) - for c in coords: - chunk_id = cg.get_chunk_id(layer=1, x=int(c[0]), y=int(c[1]), z=int(c[2])) - svs = list(cg.range_read_chunk(chunk_id)) - if svs: - sv = svs[len(svs) // 2] - root = cg.get_root(sv) - ts = cg.get_node_timestamps(np.array([root]), return_numpy=False)[0] - return int(ts.timestamp()) - raise RuntimeError("derive_initial_ts: no SVs found anywhere in the volume") + return int(earliest.timestamp()) + 1 def setup_mesh_meta( @@ -97,6 +73,11 @@ def setup_mesh_meta( Returns the mesh meta dict persisted into bigtable. """ cfg = mesh_config.with_graph_id(cg.graph_id) + n_scales = len(cg.meta.ws_cv.info["scales"]) + if not 0 <= cfg.mip < n_scales: + raise ValueError( + f"mesh_config.mip {cfg.mip} exceeds watershed scales (available 0..{n_scales - 1})" + ) existing_mesh = cg.meta.custom_data.get("mesh", {}) existing_ts = existing_mesh.get("initial_ts") initial_ts = int(existing_ts) if existing_ts is not None else derive_initial_ts(cg) @@ -122,11 +103,12 @@ def setup_mesh_meta( for layer, bits in cfg.minishard_bits.items() if layer <= cfg.max_layer } + mesh_chunk_size = get_mesh_block_shape_for_mip(cg, 2, cfg.mip) mesh_spec = { "@type": "neuroglancer_legacy_mesh", "spatial_index": None, "mip": int(cfg.mip), - "chunk_size": list(cfg.chunk_size), + "chunk_size": [int(x) for x in mesh_chunk_size], "sharding": sharding, } cg.meta.ws_cv.mesh.meta.info = mesh_spec @@ -177,4 +159,4 @@ def main() -> None: if __name__ == "__main__": - main() + run_and_exit(main) diff --git a/pychunkedgraph/pipeline/meshing/worker.py b/pychunkedgraph/pipeline/meshing/worker.py index 91f4e2207..d986b3dcf 100644 --- a/pychunkedgraph/pipeline/meshing/worker.py +++ b/pychunkedgraph/pipeline/meshing/worker.py @@ -2,18 +2,31 @@ Mirrors ``meshing_sqs.MeshTask.execute``. Idempotent (overwrites shards), so it needs no per-chunk lock. Plugged into the generic ``pipeline.worker`` harness. ``mip`` -comes from the mesh meta written by setup; a per-chunk failure counts transient so -the batch retries it. +comes from the mesh meta written by setup; an infra failure counts transient so the +batch retries it, while a bug or bad input fails the index fast. """ import logging import os +import numpy as np + from ...meshing import meshgen from ..worker import run logger = logging.getLogger(__name__) +#: Failures no retry can clear: a bug or bad input, never infrastructure. Everything +#: else is transient, so a preemption or an RPC timeout still gets the batch's retries. +FATAL_ERRORS = ( + TypeError, + ValueError, + AttributeError, + KeyError, + IndexError, + AssertionError, +) + def make_processor(cg, layer, env): """Build the mesh per-chunk processor for this batch.""" @@ -22,7 +35,11 @@ def make_processor(cg, layer, env): cache = os.environ.get("PCG_MESH_CACHE", "1") != "0" def process_one(coord): - chunk_id = int(cg.get_chunk_id(layer=layer, x=coord[0], y=coord[1], z=coord[2])) + # np.uint64, as the legacy task passes: numpy 1.26 reads a Python int as + # int64, and `int64 | uint64` has no safe common type, so every id read raises. + chunk_id = np.uint64( + cg.get_chunk_id(layer=layer, x=coord[0], y=coord[1], z=coord[2]) + ) try: if layer == 2: meshgen.chunk_initial_mesh_task( @@ -33,6 +50,9 @@ def process_one(coord): cg.graph_id, chunk_id, mip, cache=cache ) return "ok" + except FATAL_ERRORS: + logger.exception(f"fatal mesh failure on chunk {layer}_{tuple(coord)}") + return "fatal" except Exception: logger.exception(f"mesh failure on chunk {layer}_{tuple(coord)}") return "transient" diff --git a/pychunkedgraph/pipeline/worker.py b/pychunkedgraph/pipeline/worker.py index f25aceac5..8cbf619d4 100644 --- a/pychunkedgraph/pipeline/worker.py +++ b/pychunkedgraph/pipeline/worker.py @@ -28,8 +28,11 @@ def layer_bounds(cg, layer: int): return cg.meta.layer_chunk_bounds[layer] -def run(make_processor) -> int: - """Run one batch index for the configured layer; returns a process exit code.""" +def run(make_processor, finalize=None) -> int: + """Run one batch index for the configured layer; returns a process exit code. + + ``finalize(cg, layer)`` runs only after a fully successful batch; its failure fails the + pod (FATAL) without re-opening any chunk, so it is safe to re-run.""" logging.basicConfig(level=NOTE) env = { "graph_id": os.environ["PCG_GRAPH_ID"], @@ -37,7 +40,7 @@ def run(make_processor) -> int: "seed": int(os.environ["PCG_PERM_SEED"]), "batch_size": int(os.environ["PCG_BATCH_SIZE"]), "index": int(os.environ["JOB_COMPLETION_INDEX"]), - "n_threads": int(os.environ.get("PCG_N_THREADS", 1)), + "n_threads": int(os.environ.get("PCG_N_PROCESSES", 1)), } layer, index = env["layer"], env["index"] @@ -68,4 +71,10 @@ def run(make_processor) -> int: return TRANSIENT if fatal: return FATAL + if finalize: + try: + finalize(cg, layer) + except Exception: + logger.exception("finalize failed") + return FATAL return SUCCESS diff --git a/pychunkedgraph/tests/test_uncategorized.py b/pychunkedgraph/tests/test_uncategorized.py index ddbb2cf74..7f636f680 100644 --- a/pychunkedgraph/tests/test_uncategorized.py +++ b/pychunkedgraph/tests/test_uncategorized.py @@ -280,19 +280,17 @@ def test_build_single_across_edge(self, gen_graph): assert test_ace in atomic_cross_edge_d[2] assert len(children) == 1 and to_label(cg, 1, 1, 0, 0, 0) in children - # Check for the one Level 3 node that should have been created. This one combines the two - # connected components of Level 2 - # to_label(cg, 3, 0, 0, 0, 1) - assert serialize_uint64(to_label(cg, 3, 0, 0, 0, 1)) in res.rows - - attr = attributes.Hierarchy.Child - row = res.rows[serialize_uint64(to_label(cg, 3, 0, 0, 0, 1))].cells["0"] - children = attr.deserialize(row[attr.key][0].value) - assert ( - len(children) == 2 - and to_label(cg, 2, 0, 0, 0, 1) in children - and to_label(cg, 2, 1, 0, 0, 1) in children - ) + # The one Level 3 node combines the two connected components of Level 2 + root = cg.get_root(to_label(cg, 1, 0, 0, 0, 0)) + assert root == cg.get_root(to_label(cg, 1, 1, 0, 0, 0)) + assert cg.get_chunk_layer(root) == 3 + assert serialize_uint64(root) in res.rows + + children = cg.get_children(root) + assert len(children) == 2 and set(children) == { + cg.get_parent(to_label(cg, 1, 0, 0, 0, 0)), + cg.get_parent(to_label(cg, 1, 1, 0, 0, 0)), + } # Make sure there are not any more entries in the table # include counters, meta and version rows @@ -388,19 +386,18 @@ def test_build_single_edge_and_single_across_edge(self, gen_graph): assert test_ace in atomic_cross_edge_d[2] assert len(children) == 1 and to_label(cg, 1, 1, 0, 0, 0) in children - # Check for the one Level 3 node that should have been created. This one combines the two - # connected components of Level 2 - # to_label(cg, 3, 0, 0, 0, 1) - assert serialize_uint64(to_label(cg, 3, 0, 0, 0, 1)) in res.rows - row = res.rows[serialize_uint64(to_label(cg, 3, 0, 0, 0, 1))].cells["0"] - column = attributes.Hierarchy.Child - children = column.deserialize(row[column.key][0].value) + # The one Level 3 node combines the two connected components of Level 2 + root = cg.get_root(to_label(cg, 1, 0, 0, 0, 0)) + assert root == cg.get_root(to_label(cg, 1, 0, 0, 0, 1)) + assert root == cg.get_root(to_label(cg, 1, 1, 0, 0, 0)) + assert cg.get_chunk_layer(root) == 3 + assert serialize_uint64(root) in res.rows - assert ( - len(children) == 2 - and to_label(cg, 2, 0, 0, 0, 1) in children - and to_label(cg, 2, 1, 0, 0, 1) in children - ) + children = cg.get_children(root) + assert len(children) == 2 and set(children) == { + cg.get_parent(to_label(cg, 1, 0, 0, 0, 0)), + cg.get_parent(to_label(cg, 1, 1, 0, 0, 0)), + } # Make sure there are not any more entries in the table # include counters, meta and version rows @@ -436,8 +433,12 @@ def test_build_big_graph(self, gen_graph): assert serialize_uint64(to_label(cg, 1, 0, 0, 0, 0)) in res.rows assert serialize_uint64(to_label(cg, 1, 7, 7, 7, 0)) in res.rows - assert serialize_uint64(to_label(cg, 5, 0, 0, 0, 1)) in res.rows - assert serialize_uint64(to_label(cg, 5, 0, 0, 0, 2)) in res.rows + root_a = cg.get_root(to_label(cg, 1, 0, 0, 0, 0)) + root_z = cg.get_root(to_label(cg, 1, 7, 7, 7, 0)) + assert root_a != root_z + assert cg.get_chunk_layer(root_a) == cg.get_chunk_layer(root_z) == 5 + assert serialize_uint64(root_a) in res.rows + assert serialize_uint64(root_z) in res.rows @pytest.mark.timeout(30) def test_double_chunk_creation(self, gen_graph): @@ -501,15 +502,15 @@ def test_double_chunk_creation(self, gen_graph): assert cg.get_chunk_layer(cg.get_root(to_label(cg, 1, 0, 0, 0, 2))) == 4 assert cg.get_chunk_layer(cg.get_root(to_label(cg, 1, 1, 0, 0, 1))) == 4 - root_seg_ids = [ - cg.get_segment_id(cg.get_root(to_label(cg, 1, 0, 0, 0, 1))), - cg.get_segment_id(cg.get_root(to_label(cg, 1, 0, 0, 0, 2))), - cg.get_segment_id(cg.get_root(to_label(cg, 1, 1, 0, 0, 1))), + svs = [ + to_label(cg, 1, 0, 0, 0, 1), + to_label(cg, 1, 0, 0, 0, 2), + to_label(cg, 1, 1, 0, 0, 1), ] - - assert 4 in root_seg_ids - assert 5 in root_seg_ids - assert 6 in root_seg_ids + roots = [cg.get_root(sv) for sv in svs] + assert len(set(roots)) == 3 + for sv, root in zip(svs, roots): + assert cg.get_parent(sv) in cg.get_children(root) class TestGraphSimpleQueries: @@ -524,141 +525,75 @@ class TestGraphSimpleQueries: @pytest.mark.timeout(30) def test_get_parent_and_children(self, gen_graph_simplequerytest): cg = gen_graph_simplequerytest - - children10000 = cg.get_children(to_label(cg, 1, 0, 0, 0, 0)) - children11000 = cg.get_children(to_label(cg, 1, 1, 0, 0, 0)) - children11001 = cg.get_children(to_label(cg, 1, 1, 0, 0, 1)) - children12000 = cg.get_children(to_label(cg, 1, 2, 0, 0, 0)) - - parent10000 = cg.get_parent( - to_label(cg, 1, 0, 0, 0, 0), - ) - parent11000 = cg.get_parent( - to_label(cg, 1, 1, 0, 0, 0), - ) - parent11001 = cg.get_parent( - to_label(cg, 1, 1, 0, 0, 1), - ) - parent12000 = cg.get_parent( - to_label(cg, 1, 2, 0, 0, 0), - ) - - children20001 = cg.get_children(to_label(cg, 2, 0, 0, 0, 1)) - children21001 = cg.get_children(to_label(cg, 2, 1, 0, 0, 1)) - children22001 = cg.get_children(to_label(cg, 2, 2, 0, 0, 1)) - - parent20001 = cg.get_parent( - to_label(cg, 2, 0, 0, 0, 1), - ) - parent21001 = cg.get_parent( - to_label(cg, 2, 1, 0, 0, 1), - ) - parent22001 = cg.get_parent( - to_label(cg, 2, 2, 0, 0, 1), - ) - - children30001 = cg.get_children(to_label(cg, 3, 0, 0, 0, 1)) - # children30002 = cg.get_children(to_label(cg, 3, 0, 0, 0, 2)) - children31001 = cg.get_children(to_label(cg, 3, 1, 0, 0, 1)) - - parent30001 = cg.get_parent( - to_label(cg, 3, 0, 0, 0, 1), - ) - # parent30002 = cg.get_parent(to_label(cg, 3, 0, 0, 0, 2), ) - parent31001 = cg.get_parent( - to_label(cg, 3, 1, 0, 0, 1), - ) - - children40001 = cg.get_children(to_label(cg, 4, 0, 0, 0, 1)) - children40002 = cg.get_children(to_label(cg, 4, 0, 0, 0, 2)) - - parent40001 = cg.get_parent( - to_label(cg, 4, 0, 0, 0, 1), - ) - parent40002 = cg.get_parent( - to_label(cg, 4, 0, 0, 0, 2), - ) + sv_a = to_label(cg, 1, 0, 0, 0, 0) + sv_b0 = to_label(cg, 1, 1, 0, 0, 0) + sv_b1 = to_label(cg, 1, 1, 0, 0, 1) + sv_c = to_label(cg, 1, 2, 0, 0, 0) # (non-existing) Children of L1 - assert np.array_equal(children10000, []) is True - assert np.array_equal(children11000, []) is True - assert np.array_equal(children11001, []) is True - assert np.array_equal(children12000, []) is True + for sv in (sv_a, sv_b0, sv_b1, sv_c): + assert np.array_equal(cg.get_children(sv), []) is True # Parent of L1 - assert parent10000 == to_label(cg, 2, 0, 0, 0, 1) - assert parent11000 == to_label(cg, 2, 1, 0, 0, 1) - assert parent11001 == to_label(cg, 2, 1, 0, 0, 1) - assert parent12000 == to_label(cg, 2, 2, 0, 0, 1) + l2_a, l2_b, l2_c = cg.get_parent(sv_a), cg.get_parent(sv_b0), cg.get_parent(sv_c) + assert cg.get_parent(sv_b1) == l2_b + assert [cg.get_chunk_layer(node) for node in (l2_a, l2_b, l2_c)] == [2, 2, 2] + assert len({l2_a, l2_b, l2_c}) == 3 # Children of L2 - assert len(children20001) == 1 and to_label(cg, 1, 0, 0, 0, 0) in children20001 - assert ( - len(children21001) == 2 - and to_label(cg, 1, 1, 0, 0, 0) in children21001 - and to_label(cg, 1, 1, 0, 0, 1) in children21001 - ) - assert len(children22001) == 1 and to_label(cg, 1, 2, 0, 0, 0) in children22001 - - # Parent of L2 - assert parent20001 == to_label(cg, 4, 0, 0, 0, 1) - assert parent21001 == to_label(cg, 3, 0, 0, 0, 1) - assert parent22001 == to_label(cg, 3, 1, 0, 0, 1) + children_l2_a = cg.get_children(l2_a) + children_l2_b = cg.get_children(l2_b) + children_l2_c = cg.get_children(l2_c) + assert len(children_l2_a) == 1 and sv_a in children_l2_a + assert len(children_l2_b) == 2 and sv_b0 in children_l2_b and sv_b1 in children_l2_b + assert len(children_l2_c) == 1 and sv_c in children_l2_c + + # Parent of L2: A skips to the root, B and C each get a node in their own L3 chunk + root_a = cg.get_parent(l2_a) + l3_b, l3_c = cg.get_parent(l2_b), cg.get_parent(l2_c) + assert cg.get_chunk_layer(root_a) == 4 + assert cg.get_chunk_layer(l3_b) == cg.get_chunk_layer(l3_c) == 3 + assert l3_b != l3_c # Children of L3 - assert len(children30001) == 1 and len(children31001) == 1 - assert to_label(cg, 2, 1, 0, 0, 1) in children30001 - assert to_label(cg, 2, 2, 0, 0, 1) in children31001 + children_l3_b = cg.get_children(l3_b) + children_l3_c = cg.get_children(l3_c) + assert len(children_l3_b) == 1 and l2_b in children_l3_b + assert len(children_l3_c) == 1 and l2_c in children_l3_c # Parent of L3 - assert parent30001 == parent31001 - assert ( - parent30001 == to_label(cg, 4, 0, 0, 0, 1) - and parent20001 == to_label(cg, 4, 0, 0, 0, 2) - ) or ( - parent30001 == to_label(cg, 4, 0, 0, 0, 2) - and parent20001 == to_label(cg, 4, 0, 0, 0, 1) - ) + root_bc = cg.get_parent(l3_b) + assert cg.get_parent(l3_c) == root_bc + assert cg.get_chunk_layer(root_bc) == 4 + assert root_bc != root_a # Children of L4 - assert parent10000 in children40001 - assert parent21001 in children40002 and parent22001 in children40002 + assert l2_a in cg.get_children(root_a) + children_root_bc = cg.get_children(root_bc) + assert l3_b in children_root_bc and l3_c in children_root_bc # (non-existing) Parent of L4 - assert parent40001 is None - assert parent40002 is None + assert cg.get_parent(root_a) is None + assert cg.get_parent(root_bc) is None - children2_separate = cg.get_children( - [ - to_label(cg, 2, 0, 0, 0, 1), - to_label(cg, 2, 1, 0, 0, 1), - to_label(cg, 2, 2, 0, 0, 1), - ] - ) + children2_separate = cg.get_children([l2_a, l2_b, l2_c]) assert len(children2_separate) == 3 - assert to_label(cg, 2, 0, 0, 0, 1) in children2_separate and np.all( - np.isin(children2_separate[to_label(cg, 2, 0, 0, 0, 1)], children20001) + assert l2_a in children2_separate and np.all( + np.isin(children2_separate[l2_a], children_l2_a) ) - assert to_label(cg, 2, 1, 0, 0, 1) in children2_separate and np.all( - np.isin(children2_separate[to_label(cg, 2, 1, 0, 0, 1)], children21001) + assert l2_b in children2_separate and np.all( + np.isin(children2_separate[l2_b], children_l2_b) ) - assert to_label(cg, 2, 2, 0, 0, 1) in children2_separate and np.all( - np.isin(children2_separate[to_label(cg, 2, 2, 0, 0, 1)], children22001) + assert l2_c in children2_separate and np.all( + np.isin(children2_separate[l2_c], children_l2_c) ) - children2_combined = cg.get_children( - [ - to_label(cg, 2, 0, 0, 0, 1), - to_label(cg, 2, 1, 0, 0, 1), - to_label(cg, 2, 2, 0, 0, 1), - ], - flatten=True, - ) + children2_combined = cg.get_children([l2_a, l2_b, l2_c], flatten=True) assert ( len(children2_combined) == 4 - and np.all(np.isin(children20001, children2_combined)) - and np.all(np.isin(children21001, children2_combined)) - and np.all(np.isin(children22001, children2_combined)) + and np.all(np.isin(children_l2_a, children2_combined)) + and np.all(np.isin(children_l2_b, children2_combined)) + and np.all(np.isin(children_l2_c, children2_combined)) ) @pytest.mark.timeout(30) @@ -680,13 +615,10 @@ def test_get_root(self, gen_graph_simplequerytest): with pytest.raises(Exception): cg.get_root(0) - assert ( - root10000 == to_label(cg, 4, 0, 0, 0, 1) - and root11000 == root11001 == root12000 == to_label(cg, 4, 0, 0, 0, 2) - ) or ( - root10000 == to_label(cg, 4, 0, 0, 0, 2) - and root11000 == root11001 == root12000 == to_label(cg, 4, 0, 0, 0, 1) - ) + assert root11000 == root11001 == root12000 + assert root10000 != root11000 + assert cg.get_chunk_layer(root10000) == cg.get_chunk_layer(root11000) == 4 + assert cg.get_parent(root10000) is None and cg.get_parent(root11000) is None @pytest.mark.timeout(30) def test_get_subgraph_nodes(self, gen_graph_simplequerytest): diff --git a/requirements.yml b/requirements.yml index 59ff911cd..c7e985e04 100644 --- a/requirements.yml +++ b/requirements.yml @@ -1,6 +1,10 @@ name: pychunkedgraph channels: - conda-forge + # never resolve against Anaconda's `defaults`: its channels are the only ones + # gated on a terms-of-service acceptance, and the plugin serving that gate + # breaks whenever conda itself is updated from conda-forge. + - nodefaults dependencies: - python==3.11.4 - pip