Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 0 additions & 2 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -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 ./
Expand Down
5 changes: 5 additions & 0 deletions pychunkedgraph/graph/chunkedgraph.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
19 changes: 13 additions & 6 deletions pychunkedgraph/ingest/create/abstract_layers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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))}"


Expand Down Expand Up @@ -166,7 +176,6 @@ def _write_connected_components(
ccs_with_node_ids,
node_layer_d_shared,
time_stamp,
use_threads=use_threads,
)
return

Expand All @@ -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:
Expand All @@ -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]):
Expand Down
176 changes: 176 additions & 0 deletions pychunkedgraph/ingest/simple_tests.py
Original file line number Diff line number Diff line change
@@ -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)
12 changes: 9 additions & 3 deletions pychunkedgraph/pipeline/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down Expand Up @@ -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 |
Expand All @@ -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

Expand Down
26 changes: 26 additions & 0 deletions pychunkedgraph/pipeline/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
5 changes: 2 additions & 3 deletions pychunkedgraph/pipeline/ingest/__main__.py
Original file line number Diff line number Diff line change
@@ -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)
3 changes: 2 additions & 1 deletion pychunkedgraph/pipeline/ingest/setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -63,4 +64,4 @@ def main() -> None:


if __name__ == "__main__":
main()
run_and_exit(main)
14 changes: 13 additions & 1 deletion pychunkedgraph/pipeline/ingest/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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__":
Expand Down
5 changes: 2 additions & 3 deletions pychunkedgraph/pipeline/meshing/__main__.py
Original file line number Diff line number Diff line change
@@ -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)
Loading
Loading