Skip to content

Commit 4dce873

Browse files
committed
feat: add streaming parquet repair, validation, and sidecar management tools with corresponding tests
1 parent 200afcb commit 4dce873

23 files changed

Lines changed: 3386 additions & 53 deletions

.gitignore

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,8 @@ docs/tilelang_ports/mamba3_path_c_lowered.metal
6363
reports/raw/*
6464
vbgui/.DS_Store
6565
secrets/
66+
.env
67+
outputs/secrets/
6668
*.gh_tokens
6769
outputs/logs/*
6870
outputs/*.log
@@ -88,3 +90,33 @@ outputs/reindexed_commits/*
8890
outputs/non_strict_backup_20260626_165907/*
8991
outputs/profiles/*
9092
outputs/reindexed_pr/*
93+
outputs/conveyor/tmp/*
94+
outputs/conveyor/tmp/streaming_conveyor_aq5oisfq/aosp-frameworks-native/_src/METADATA
95+
outputs/nebius_smoke/parquet/*
96+
outputs/*.json
97+
outputs/nebius_smoke/*
98+
outputs/sidecar_audit_commits_*
99+
outputs/sidecar_audit_all_verified/sidecar_parquet_audit.json
100+
outputs/sidecar_audit_1024_verified/sidecar_parquet_audit.md
101+
outputs/sidecar_audit_1024_check/sidecar_parquet_audit.json
102+
outputs/sidecar_audit_1024_check/sidecar_parquet_audit.md
103+
outputs/sidecar_audit_1024_verified/sidecar_parquet_audit.json
104+
outputs/sidecar_audit_1k_2k_4k_verified/sidecar_parquet_audit.json
105+
outputs/sidecar_audit_1k_2k_4k_verified/sidecar_parquet_audit.md
106+
outputs/sidecar_audit_all_after_validity_fix/sidecar_parquet_audit.json
107+
outputs/sidecar_audit_all_after_validity_fix/sidecar_parquet_audit.md
108+
outputs/sidecar_audit_all_final_poststop_valid/sidecar_parquet_audit.json
109+
outputs/sidecar_audit_all_final_poststop_valid/sidecar_parquet_audit.md
110+
outputs/sidecar_audit_all_final_valid/sidecar_parquet_audit.json
111+
outputs/sidecar_audit_all_final_valid/sidecar_parquet_audit.md
112+
outputs/sidecar_audit_all_verified/sidecar_parquet_audit.md
113+
outputs/sidecar_audit_code_1k_2k_4k/sidecar_parquet_audit.json
114+
outputs/sidecar_audit_code_1k_2k_4k/sidecar_parquet_audit.md
115+
outputs/sidecar_audit_commit_pr_8k_verified/sidecar_parquet_audit.json
116+
outputs/sidecar_audit_commit_pr_8k_verified/sidecar_parquet_audit.md
117+
outputs/sidecar_audit_pr_1k_2k_4k_8k/sidecar_parquet_audit.json
118+
outputs/sidecar_audit_pr_1k_2k_4k_8k/sidecar_parquet_audit.md
119+
outputs/conveyor/progress_code.jsonl
120+
outputs/conveyor/progress_commits.jsonl
121+
outputs/conveyor/locks/code.lock
122+
outputs/conveyor/locks/commits.lock

cppmega_mlx/tokenizer/cpp_tokenizer.py

Lines changed: 69 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,9 @@
4040
# and tabs *inside* a literal carry meaning (e.g. ``"a b"``) and must NOT be
4141
# collapsed to ``<SPACE>``; doing so would change the program semantics and the
4242
# re-encoded ids. The scanner below tracks literal state and only collapses
43-
# whitespace runs that live OUTSIDE any literal.
43+
# whitespace runs that live OUTSIDE any literal. It also tracks comments so
44+
# natural-language apostrophes in docstrings (``can't``) are not mistaken for C++
45+
# char literals; whitespace inside comments is still normalized.
4446
_WS_RUN_RE = re.compile(r"[ \t]+")
4547
_NL_RUN_RE = re.compile(r"[\r\n]+")
4648
_NL_SENTINEL = "<NL>"
@@ -60,18 +62,37 @@ def normalize_whitespace_with_offsets(text: str) -> tuple[str, list[tuple[int, i
6062
sidecar channels with no off-by-N shift.
6163
6264
Whitespace INSIDE string/char/raw-string literals is preserved verbatim
63-
(identity offsets); only whitespace outside any literal is collapsed.
65+
(identity offsets); only whitespace outside any literal is collapsed. C++
66+
comments are not literals: quotes/apostrophes inside comments are copied as
67+
plain text, while comment whitespace is still collapsed to sentinels.
6468
"""
6569
chars: list[str] = []
6670
spans: list[tuple[int, int]] = []
6771
n = len(text)
6872
i = 0
69-
# Literal state: None=code, '"'=string, "'"=char, "raw"=raw-string.
73+
# State: None=code, '"'=string, "'"=char, "raw"=raw-string,
74+
# "line_comment"=// comment, "block_comment"=/* */ comment.
7075
state: str | None = None
7176
raw_delim = "" # closing delimiter for the active raw string: )<d>"
7277
while i < n:
7378
ch = text[i]
7479
if state is None:
80+
if ch == "/" and i + 1 < n and text[i + 1] == "/":
81+
chars.append("/")
82+
spans.append((i, i + 1))
83+
chars.append("/")
84+
spans.append((i + 1, i + 2))
85+
state = "line_comment"
86+
i += 2
87+
continue
88+
if ch == "/" and i + 1 < n and text[i + 1] == "*":
89+
chars.append("/")
90+
spans.append((i, i + 1))
91+
chars.append("*")
92+
spans.append((i + 1, i + 2))
93+
state = "block_comment"
94+
i += 2
95+
continue
7596
# Detect raw-string prefix: R"<delim>( ... )<delim>"
7697
if ch in ("R", "u", "U", "L") and _raw_string_opener_at(text, i):
7798
r_pos = text.index('R', i)
@@ -109,6 +130,51 @@ def normalize_whitespace_with_offsets(text: str) -> tuple[str, list[tuple[int, i
109130
chars.append(ch)
110131
spans.append((i, i + 1))
111132
i += 1
133+
elif state == "line_comment":
134+
if ch in "\r\n":
135+
j = i + 1
136+
while j < n and text[j] in "\r\n":
137+
j += 1
138+
_append_sentinel(chars, spans, _NL_SENTINEL, i, j)
139+
state = None
140+
i = j
141+
continue
142+
if ch in " \t":
143+
j = i + 1
144+
while j < n and text[j] in " \t":
145+
j += 1
146+
_append_sentinel(chars, spans, _SPACE_SENTINEL, i, j)
147+
i = j
148+
continue
149+
chars.append(ch)
150+
spans.append((i, i + 1))
151+
i += 1
152+
elif state == "block_comment":
153+
if text.startswith("*/", i):
154+
chars.append("*")
155+
spans.append((i, i + 1))
156+
chars.append("/")
157+
spans.append((i + 1, i + 2))
158+
state = None
159+
i += 2
160+
continue
161+
if ch in "\r\n":
162+
j = i + 1
163+
while j < n and text[j] in "\r\n":
164+
j += 1
165+
_append_sentinel(chars, spans, _NL_SENTINEL, i, j)
166+
i = j
167+
continue
168+
if ch in " \t":
169+
j = i + 1
170+
while j < n and text[j] in " \t":
171+
j += 1
172+
_append_sentinel(chars, spans, _SPACE_SENTINEL, i, j)
173+
i = j
174+
continue
175+
chars.append(ch)
176+
spans.append((i, i + 1))
177+
i += 1
112178
elif state == "raw":
113179
if text.startswith(raw_delim, i):
114180
for k in range(i, i + len(raw_delim)):

scripts/_verify_conveyor_live.py

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,79 @@
4646
}
4747
ALL_SIDE = [c for v in SIDE_CHANNELS.values() for c in v]
4848

49+
def extract_git_history_process_health():
50+
"""Report real JSONL writer duplicates, ignoring transient git subprocess forks.
51+
52+
On macOS a subprocess fork can briefly inherit the parent's command line
53+
before exec'ing git. A raw ps grouping by --output sees that as a duplicate
54+
extract_git_history.py, but it is parented by the real extract process and
55+
is not a JSONL writer. Treat only non-extract-parent processes as writers.
56+
"""
57+
try:
58+
out = subprocess.check_output(
59+
['ps', '-axo', 'pid,ppid,stat,etime,command'],
60+
text=True,
61+
)
62+
except Exception as e:
63+
return {'error': str(e)}
64+
65+
proc_re = re.compile(r'\s*(\d+)\s+(\d+)\s+(\S+)\s+(\S+)\s+(.*)')
66+
output_re = re.compile(r'--output\s+(\S+)')
67+
cmd_by_pid = {}
68+
rows = []
69+
for line in out.splitlines():
70+
m = proc_re.match(line)
71+
if not m:
72+
continue
73+
pid = int(m.group(1))
74+
ppid = int(m.group(2))
75+
cmd = m.group(5)
76+
cmd_by_pid[pid] = cmd
77+
if 'extract_git_history.py' not in cmd:
78+
continue
79+
out_m = output_re.search(cmd)
80+
rows.append({
81+
'pid': pid,
82+
'ppid': ppid,
83+
'stat': m.group(3),
84+
'etime': m.group(4),
85+
'output': out_m.group(1) if out_m else '',
86+
'line': line,
87+
})
88+
89+
raw_by_output = collections.defaultdict(list)
90+
writer_by_output = collections.defaultdict(list)
91+
ignored_children = []
92+
for row in rows:
93+
raw_by_output[row['output']].append(row)
94+
parent_cmd = cmd_by_pid.get(row['ppid'], '')
95+
if 'extract_git_history.py' in parent_cmd:
96+
ignored_children.append(row)
97+
else:
98+
writer_by_output[row['output']].append(row)
99+
100+
raw_dupes = {
101+
output: [row['line'] for row in group]
102+
for output, group in raw_by_output.items()
103+
if len(group) > 1
104+
}
105+
writer_dupes = {
106+
output: [row['line'] for row in group]
107+
for output, group in writer_by_output.items()
108+
if len(group) > 1
109+
}
110+
return {
111+
'raw_extract_outputs': len(raw_by_output),
112+
'raw_extract_procs': sum(len(group) for group in raw_by_output.values()),
113+
'raw_dupe_outputs': len(raw_dupes),
114+
'root_writer_outputs': len(writer_by_output),
115+
'root_writer_procs': sum(len(group) for group in writer_by_output.values()),
116+
'root_writer_dupe_outputs': len(writer_dupes),
117+
'fork_before_exec_children_ignored': len(ignored_children),
118+
'raw_dupe_examples': dict(list(raw_dupes.items())[:5]),
119+
'root_writer_dupe_examples': dict(list(writer_dupes.items())[:5]),
120+
}
121+
49122
def strip_for_clang(text: str) -> str:
50123
"""Strip leading language/platform header comments + special-token text; keep code."""
51124
text = SPECIAL_TEXT_RE.sub('', text)
@@ -497,6 +570,7 @@ def bucket_report(buckets):
497570
'total_samples': len(samples),
498571
'counts': dict(counts),
499572
'sampling_plan': sampling_plan,
573+
'process_health': extract_git_history_process_health(),
500574
'skipped_race': skipped,
501575
'per_type': {t: agg(t) for t in ('CODE', 'COMMIT', 'BUILD-FILE')},
502576
'dedup': dedup_info,

0 commit comments

Comments
 (0)