Skip to content

Commit 509a5c8

Browse files
committed
Begin decoupling checksum logic
1 parent 8a35beb commit 509a5c8

6 files changed

Lines changed: 190 additions & 40 deletions

File tree

pyproject.toml

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,8 @@ dependencies = [
3131
"universal-pathlib>=0.3.3",
3232
"xarray>=2025.10.1",
3333
"zarr>=3.1.3",
34+
"crc32c_dist_rs",
35+
"google-crc32c>=1.5.0",
3436
]
3537

3638
[project.optional-dependencies]
@@ -80,7 +82,12 @@ docs = [
8082
required-version = ">=0.8.17"
8183

8284
[tool.uv.sources]
83-
crc32c_dist_rs = { path = "/home/brian_michell_tgs_com/source/crc32c_dist_rs/target/wheels/crc32c_dist_rs-0.1.4-cp313-cp313-manylinux_2_34_x86_64.whl" }
85+
crc32c_dist_rs = [
86+
{ path = "/home/brian_michell_tgs_com/source/crc32c_dist_rs/target/wheels/crc32c_dist_rs-0.1.4-cp311-cp311-manylinux_2_34_x86_64.whl", marker = "python_version == '3.11'" },
87+
{ path = "/home/brian_michell_tgs_com/source/crc32c_dist_rs/target/wheels/crc32c_dist_rs-0.1.4-cp312-cp312-manylinux_2_34_x86_64.whl", marker = "python_version == '3.12'" },
88+
{ path = "/home/brian_michell_tgs_com/source/crc32c_dist_rs/target/wheels/crc32c_dist_rs-0.1.4-cp313-cp313-manylinux_2_34_x86_64.whl", marker = "python_version == '3.13'" },
89+
]
90+
8491

8592
[tool.ruff]
8693
target-version = "py311"

src/mdio/segy/_workers.py

Lines changed: 4 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -32,9 +32,8 @@ def header_scan_worker(
3232
segy_file_kwargs: SegyFileArguments,
3333
trace_range: tuple[int, int],
3434
subset: list[str] | None = None,
35-
calculate_checksum: bool = False,
36-
) -> HeaderArray | tuple[HeaderArray, tuple[int, int, int]]:
37-
"""Header scan worker with optional checksum calculation.
35+
) -> HeaderArray:
36+
"""Header scan worker.
3837
3938
If SegyFile is not open, it can either accept a path string or a handle that was opened in
4039
a different context manager.
@@ -43,11 +42,9 @@ def header_scan_worker(
4342
segy_file_kwargs: Arguments to open SegyFile instance.
4443
trace_range: Tuple consisting of the trace ranges to read.
4544
subset: List of header names to filter and keep.
46-
calculate_checksum: If True, also calculate CRC32C for this trace range.
4745
4846
Returns:
49-
HeaderArray if calculate_checksum is False, otherwise tuple of (HeaderArray, checksum_info)
50-
where checksum_info is (byte_offset, crc32c, byte_length).
47+
HeaderArray parsed from SEG-Y library.
5148
"""
5249
segy_file = SegyFileWrapper(**segy_file_kwargs)
5350

@@ -73,33 +70,7 @@ def header_scan_worker(
7370
# (singleton) so we can concat and assign stuff later.
7471
trace_header = np.array(trace_header, dtype=new_dtype, ndmin=1)
7572

76-
if not calculate_checksum:
77-
return HeaderArray(trace_header) # wrap back so we can use aliases
78-
79-
# Verify checksum libraries are available (should have been checked earlier)
80-
if not is_checksum_available():
81-
logger.warning("Checksum calculation requested but libraries not available. Skipping checksum.")
82-
# Return headers without checksum info - this should not happen in practice
83-
# as the caller should check availability first
84-
return HeaderArray(trace_header)
85-
86-
# Calculate checksum from the raw bytes ALREADY IN MEMORY (NO ADDITIONAL I/O!)
87-
raw_bytes = traces.trace_buffer_array.tobytes()
88-
partial_crc32c = calculate_bytes_crc32c(raw_bytes)
89-
90-
# Calculate byte offset and length
91-
trace_header_size = segy_file.spec.trace.header.itemsize
92-
# sample_size = segy_file.spec.trace.sample.itemsize
93-
sample_size = 4 # This will always be a 4-byte float
94-
num_samples = len(segy_file.sample_labels)
95-
trace_size = trace_header_size + (num_samples * sample_size)
96-
97-
byte_offset = 3600 + start_trace * trace_size
98-
byte_length = len(raw_bytes)
99-
100-
checksum_info = (byte_offset, partial_crc32c, byte_length)
101-
102-
return HeaderArray(trace_header), checksum_info
73+
return HeaderArray(trace_header) # wrap back so we can use aliases
10374

10475

10576
def trace_worker( # noqa: PLR0913

src/mdio/segy/checksum.py

Lines changed: 71 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,9 +18,13 @@
1818
import os
1919
from typing import TYPE_CHECKING
2020
from typing import Any
21+
import numpy as np
22+
from segy import SegyFile
23+
from mdio.segy._raw_trace_wrapper import SegyFileRawTraceWrapper
2124

25+
from segy.arrays import HeaderArray
2226
if TYPE_CHECKING:
23-
from segy.arrays import HeaderArray
27+
from mdio.segy._workers import SegyFileArguments
2428

2529
logger = logging.getLogger(__name__)
2630

@@ -43,6 +47,72 @@
4347
DistributedCRC32C = None # type: ignore[assignment,misc]
4448

4549

50+
def header_scan_worker(
51+
segy_file_kwargs: "SegyFileArguments",
52+
trace_range: tuple[int, int],
53+
subset: list[str] | None = None,
54+
) -> HeaderArray | tuple[HeaderArray, tuple[int, int, int]]:
55+
"""Header scan worker with optional checksum calculation.
56+
57+
If SegyFile is not open, it can either accept a path string or a handle that was opened in
58+
a different context manager.
59+
60+
Args:
61+
segy_file_kwargs: Arguments to open SegyFile instance.
62+
trace_range: Tuple consisting of the trace ranges to read.
63+
subset: List of header names to filter and keep.
64+
calculate_checksum: If True, also calculate CRC32C for this trace range.
65+
66+
Returns:
67+
HeaderArray if calculate_checksum is False, otherwise tuple of (HeaderArray, checksum_info)
68+
where checksum_info is (byte_offset, crc32c, byte_length).
69+
"""
70+
print("Using header_scan_worker from checksum.py")
71+
segy_file = SegyFile(**segy_file_kwargs)
72+
73+
start_trace, end_trace = trace_range
74+
trace_indices = list(range(start_trace, end_trace))
75+
76+
# Read full trace data (header + samples) ONCE using the raw trace wrapper
77+
# This avoids duplicate I/O - we get both headers and raw bytes in a single read
78+
traces = SegyFileRawTraceWrapper(segy_file, trace_indices)
79+
80+
# Extract headers from the data we already read (no additional I/O)
81+
trace_header = traces.header
82+
83+
if subset is not None:
84+
# struct field selection needs a list, not a tuple; a subset is a tuple from the template.
85+
trace_header = trace_header[list(subset)]
86+
87+
# Get non-void fields from dtype and copy to new array for memory efficiency
88+
fields = trace_header.dtype.fields
89+
non_void_fields = [(name, dtype) for name, (dtype, _) in fields.items()]
90+
new_dtype = np.dtype(non_void_fields)
91+
92+
# Copy to non-padded memory, ndmin is to handle the case where there is 1 trace in block
93+
# (singleton) so we can concat and assign stuff later.
94+
trace_header = np.array(trace_header, dtype=new_dtype, ndmin=1)
95+
96+
97+
# Calculate checksum from the raw bytes ALREADY IN MEMORY (NO ADDITIONAL I/O!)
98+
raw_bytes = traces.trace_buffer_array.tobytes()
99+
partial_crc32c = calculate_bytes_crc32c(raw_bytes)
100+
101+
# Calculate byte offset and length
102+
trace_header_size = segy_file.spec.trace.header.itemsize
103+
# sample_size = segy_file.spec.trace.sample.itemsize
104+
sample_size = 4 # This will always be a 4-byte float
105+
num_samples = len(segy_file.sample_labels)
106+
trace_size = trace_header_size + (num_samples * sample_size)
107+
108+
byte_offset = 3600 + start_trace * trace_size
109+
byte_length = len(raw_bytes)
110+
111+
checksum_info = (byte_offset, partial_crc32c, byte_length)
112+
113+
return HeaderArray(trace_header), checksum_info
114+
115+
46116
def is_checksum_available() -> bool:
47117
"""Check if checksum libraries are available.
48118

src/mdio/segy/parsers.py

Lines changed: 65 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -16,10 +16,11 @@
1616
from tqdm.auto import tqdm
1717
from upath import UPath
1818

19-
from mdio.segy._workers import header_scan_worker
19+
# All imports from the `checksum` module are slated to be removed in the future
2020
from mdio.segy.checksum import create_distributed_crc32c
2121
from mdio.segy.checksum import finalize_distributed_checksum
2222
from mdio.segy.checksum import is_checksum_available
23+
from mdio.segy.checksum import should_calculate_checksum
2324

2425
if TYPE_CHECKING:
2526
from segy.arrays import HeaderArray
@@ -51,6 +52,14 @@ def parse_headers( # noqa: PLR0913
5152
HeaderArray if calculate_checksum is False.
5253
Tuple of (HeaderArray, combined_crc32c) if calculate_checksum is True.
5354
"""
55+
# Dynamically import the appropriate header_scan_worker based on checksum requirement
56+
if should_calculate_checksum() and is_checksum_available():
57+
from mdio.segy.checksum import header_scan_worker
58+
else:
59+
from mdio.segy._workers import header_scan_worker
60+
61+
# Type hint for the dynamically imported function
62+
header_scan_worker: Any
5463
# Initialize combiner only if checksum calculation is needed and available
5564
combiner: Any = None
5665
if calculate_checksum:
@@ -101,7 +110,6 @@ def parse_headers( # noqa: PLR0913
101110
repeat(segy_file_kwargs),
102111
trace_ranges,
103112
repeat(subset),
104-
repeat(calculate_checksum),
105113
)
106114

107115
if progress_bar is True:
@@ -124,4 +132,58 @@ def parse_headers( # noqa: PLR0913
124132
# Merge headers and return
125133
return np.concatenate(results)
126134

127-
return finalize_distributed_checksum(results, combiner)
135+
return finalize_distributed_checksum(results, combiner)
136+
137+
# def _parse_headers(
138+
# segy_file_kwargs: SegyFileArguments,
139+
# num_traces: int,
140+
# subset: list[str] | None = None,
141+
# block_size: int = 10000,
142+
# progress_bar: bool = True,
143+
# ) -> HeaderArray:
144+
# """Read and parse given `byte_locations` from SEG-Y file.
145+
146+
# Args:
147+
# segy_file_kwargs: SEG-Y file arguments.
148+
# num_traces: Total number of traces in the SEG-Y file.
149+
# subset: List of header names to filter and keep.
150+
# block_size: Number of traces to read for each block.
151+
# progress_bar: Enable or disable progress bar. Default is True.
152+
153+
# Returns:
154+
# HeaderArray. Keys are the index names, values are numpy arrays of parsed headers for the
155+
# current block. Array is of type byte_type except IBM32 which is mapped to FLOAT32.
156+
# """
157+
# trace_count = num_traces
158+
# n_blocks = int(ceil(trace_count / block_size))
159+
160+
# trace_ranges = []
161+
# for idx in range(n_blocks):
162+
# start, stop = idx * block_size, (idx + 1) * block_size
163+
# stop = min(stop, trace_count)
164+
165+
# trace_ranges.append((start, stop))
166+
167+
# num_cpus = int(os.getenv("MDIO__IMPORT__CPU_COUNT", default_cpus))
168+
# num_workers = min(n_blocks, num_cpus)
169+
170+
# tqdm_kw = {"unit": "block", "dynamic_ncols": True}
171+
# # For Unix async writes with s3fs/fsspec & multiprocessing, use 'spawn' instead of default
172+
# # 'fork' to avoid deadlocks on cloud stores. Slower but necessary. Default on Windows.
173+
# context = mp.get_context("spawn")
174+
# with ProcessPoolExecutor(num_workers, mp_context=context) as executor:
175+
# lazy_work = executor.map(header_scan_worker, repeat(segy_file_kwargs), trace_ranges, repeat(subset))
176+
177+
# if progress_bar is True:
178+
# lazy_work = tqdm(
179+
# iterable=lazy_work,
180+
# total=n_blocks,
181+
# desc="Scanning SEG-Y for geometry attributes",
182+
# **tqdm_kw,
183+
# )
184+
185+
# # This executes the lazy work.
186+
# headers: list[HeaderArray] = list(lazy_work)
187+
188+
# # Merge blocks before return
189+
# return np.concatenate(headers)

tests/integration/test_crc32c_checksum.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,7 @@ def test_corrupted_data_detects_changes(self, tmp_path: Path, monkeypatch: pytes
126126
"""Test that checksums detect data corruption using synthetic streamer data."""
127127
# Enable checksum calculation via environment variable
128128
monkeypatch.setenv("MDIO__IMPORT__RAW_HEADERS", "true")
129+
monkeypatch.setenv("MDIO__IMPORT__DO_CRC32C", "true")
129130

130131
# Use streamer configuration for this test
131132
test_conf = STREAMER_3D_CONF

uv.lock

Lines changed: 41 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

0 commit comments

Comments
 (0)