Skip to content

Commit ed5ea53

Browse files
authored
Merge pull request #36 from cacheMon/claude/fervent-noether-nrf2fq
Fix correctness bugs in BPF prober and trace pipeline
2 parents 3ca07b8 + 796b3ba commit ed5ea53

6 files changed

Lines changed: 381 additions & 171 deletions

File tree

.github/scripts/bpf_smoke.py

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,31 @@ def main():
5959
("kretprobe", "vfs_read", "trace_vfs_read_ret"),
6060
("kprobe", "vfs_write", "trace_vfs_write"),
6161
("kretprobe", "vfs_write", "trace_vfs_write_ret"),
62+
# fsync de-dup pair: the kretprobe clears the nested-call marker.
63+
("kprobe", "vfs_fsync", "trace_vfs_fsync"),
64+
("kretprobe", "vfs_fsync", "trace_vfs_fsync_ret"),
65+
("kprobe", "vfs_fsync_range", "trace_vfs_fsync_range"),
6266
]
67+
68+
# Symbol-conditional probes. These validate that the pt_regs-unwrapping
69+
# *_x64 variants and the DIO direction entry probes were compiled in and
70+
# attach — a guard mismatch would otherwise pass CI (BPF() load succeeds
71+
# without them) and abort the tracer at startup instead.
72+
conditional_probes = [
73+
(b"__x64_sys_mremap", [("kprobe", "__x64_sys_mremap", "trace_mremap_entry_x64"),
74+
("kretprobe", "__x64_sys_mremap", "trace_mremap_ret")]),
75+
(b"__x64_sys_openat", [("kprobe", "__x64_sys_openat", "trace_openat_entry_x64")]),
76+
(b"__x64_sys_io_uring_enter", [("kprobe", "__x64_sys_io_uring_enter", "trace_io_uring_enter_x64")]),
77+
(b"iomap_dio_rw", [("kprobe", "iomap_dio_rw", "trace_dio_entry_iomap"),
78+
("kretprobe", "iomap_dio_rw", "trace_dio_return")]),
79+
(b"__blockdev_direct_IO", [("kprobe", "__blockdev_direct_IO", "trace_dio_entry_blockdev")]),
80+
]
81+
for symbol, symbol_probes in conditional_probes:
82+
if BPF.get_kprobe_functions(symbol):
83+
probes.extend(symbol_probes)
84+
else:
85+
print(f"SKIP: {symbol.decode()} not present on this kernel")
86+
6387
for kind, event, fn in probes:
6488
if kind == "kprobe":
6589
b.attach_kprobe(event=event, fn_name=fn)

src/tracer/FlagMapper.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -697,7 +697,7 @@ def format_vfs_flags(self, op_name, flags):
697697
str: Human-readable flags, or an empty string when the operation
698698
does not define flag semantics for the current value.
699699
"""
700-
if op_name in {"OPEN", "READ", "WRITE", "CLOSE", "FSYNC", "READDIR"}:
700+
if op_name in {"OPEN", "READ", "WRITE", "CLOSE", "FSYNC", "FDATASYNC", "READDIR"}:
701701
return self.format_fs_flags(flags)
702702
if op_name == "MKDIR":
703703
return self.format_mode_flags(flags)

src/tracer/IOTracer.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1110,10 +1110,10 @@ def trace(self):
11101110
while remaining > 0 and self.running:
11111111
sleep_time = min(0.1, remaining)
11121112
time.sleep(sleep_time)
1113-
1113+
11141114
current = time.time()
11151115
remaining = end_time - current # type: ignore
1116-
1116+
11171117
if self.verbose and int(current) % 10 == 0 and int(current) > int(current - sleep_time):
11181118
elapsed = current - start
11191119
logger("info", f"Progress: {elapsed:.1f}s/{duration_target}s") # type: ignore
@@ -1123,7 +1123,7 @@ def trace(self):
11231123
# Run indefinitely until Ctrl+C
11241124
while self.running:
11251125
time.sleep(0.1)
1126-
1126+
11271127
if self.verbose:
11281128
current = time.time()
11291129
if int(current) % 30 == 0: # Every 30 seconds

src/tracer/KernelProbeTracker.py

Lines changed: 17 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -174,10 +174,13 @@ def attach_probes(self):
174174
# Capture the user-provided filename before the kernel resolves it.
175175
# Must be registered BEFORE vfs_open so the path is staged in time.
176176
# Uses the same fallback chain as the kretprobe.
177+
# NOTE: __x64_sys_* wrappers take a single pt_regs* holding the
178+
# user registers, so they need the *_x64 probe variants that
179+
# unwrap it; reading PARM1-4 directly there yields garbage.
177180
if BPF.get_kprobe_functions(b'do_sys_openat2'):
178181
self.add_kprobe("do_sys_openat2", "trace_do_sys_openat2_entry")
179182
elif BPF.get_kprobe_functions(b'__x64_sys_openat'):
180-
self.add_kprobe("__x64_sys_openat", "trace_do_sys_openat2_entry")
183+
self.add_kprobe("__x64_sys_openat", "trace_openat_entry_x64")
181184
else:
182185
self.add_kprobe("sys_openat", "trace_do_sys_openat2_entry")
183186
self.add_kprobe("vfs_open", "trace_vfs_open")
@@ -191,6 +194,9 @@ def attach_probes(self):
191194
else:
192195
self.add_kretprobe("sys_openat", "trace_sys_openat_ret")
193196
self.add_kprobe("vfs_fsync", "trace_vfs_fsync")
197+
# Return probe clears the marker that suppresses the duplicate
198+
# event from the nested vfs_fsync -> vfs_fsync_range call.
199+
self.add_kretprobe("vfs_fsync", "trace_vfs_fsync_ret")
194200
self.add_kprobe("ksys_sync", "trace_ksys_sync")
195201
self.add_kprobe("vfs_fsync_range", "trace_vfs_fsync_range")
196202
self.add_kprobe("__fput", "trace_fput")
@@ -200,9 +206,10 @@ def attach_probes(self):
200206
self.add_kretprobe("do_mmap", "trace_mmap_ret")
201207
self.add_kprobe("__vm_munmap", "trace_munmap")
202208

203-
# mremap probes — kernel may export the arch wrapper or the generic symbol
209+
# mremap probes — kernel may export the arch wrapper or the generic
210+
# symbol. The wrapper needs the pt_regs-unwrapping variant.
204211
if BPF.get_kprobe_functions(b'__x64_sys_mremap'):
205-
self.add_kprobe("__x64_sys_mremap", "trace_mremap_entry")
212+
self.add_kprobe("__x64_sys_mremap", "trace_mremap_entry_x64")
206213
self.add_kretprobe("__x64_sys_mremap", "trace_mremap_ret")
207214
elif BPF.get_kprobe_functions(b'sys_mremap'):
208215
self.add_kprobe("sys_mremap", "trace_mremap_entry")
@@ -251,12 +258,16 @@ def attach_probes(self):
251258
# else:
252259
# logger("warning", "filemap_fault not available - mmap I/O tracking disabled")
253260

254-
# Direct I/O probe for bypass detection (return probe only - no latency tracking)
261+
# Direct I/O probes. The entry probe stages the I/O direction from
262+
# the iov_iter (the return value alone cannot distinguish a read
263+
# from a write); the return probe emits the completion event.
255264
if BPF.get_kprobe_functions(b'iomap_dio_rw'):
265+
self.add_kprobe("iomap_dio_rw", "trace_dio_entry_iomap")
256266
self.add_kretprobe("iomap_dio_rw", "trace_dio_return")
257267
if self.developer_mode:
258268
logger("info", "Direct I/O tracing enabled via iomap_dio_rw")
259269
elif BPF.get_kprobe_functions(b'__blockdev_direct_IO'):
270+
self.add_kprobe("__blockdev_direct_IO", "trace_dio_entry_blockdev")
260271
self.add_kretprobe("__blockdev_direct_IO", "trace_dio_return")
261272
if self.developer_mode:
262273
logger("info", "Direct I/O tracing enabled via __blockdev_direct_IO")
@@ -376,7 +387,8 @@ def attach_probes(self):
376387
if self.developer_mode:
377388
logger("info", "io_uring tracing enabled via __io_uring_enter")
378389
elif BPF.get_kprobe_functions(b'__x64_sys_io_uring_enter'):
379-
self.add_kprobe("__x64_sys_io_uring_enter", "trace_io_uring_enter")
390+
# Syscall wrapper: needs the pt_regs-unwrapping variant.
391+
self.add_kprobe("__x64_sys_io_uring_enter", "trace_io_uring_enter_x64")
380392
if self.developer_mode:
381393
logger("info", "io_uring tracing enabled via __x64_sys_io_uring_enter")
382394
elif BPF.get_kprobe_functions(b'__sys_io_uring_enter'):

src/tracer/PathResolver.py

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -220,23 +220,29 @@ def resolve_path(self, inode: int, pid: int | None = None, filename: str | None
220220
def cleanup_old_cache(self):
221221
"""
222222
Remove old entries from cache to prevent memory bloat.
223-
223+
224224
Removes:
225225
- Process entries older than cache_timeout * 10 seconds
226226
- Limits inode cache to 5000 most recent entries
227+
228+
Thread-safety: this runs on the polling thread (from the perf-buffer
229+
callbacks via cache maintenance), the same thread that mutates these
230+
dicts. Iteration still works on list() snapshots and removals
231+
tolerate missing entries as defense in depth.
227232
"""
228233
current_time = time.time()
229-
234+
230235
# Clean up process cache
231236
pids_to_remove = []
232-
for pid, last_time in self.last_update.items():
237+
for pid, last_time in list(self.last_update.items()):
233238
if current_time - last_time > self.cache_timeout * 10:
234239
pids_to_remove.append(pid)
235-
240+
236241
for pid in pids_to_remove:
237242
self.pid_to_files.pop(pid, None)
238243
self.last_update.pop(pid, None)
239-
244+
245+
240246
# Optionally limit inode cache size
241247
if len(self.inode_to_path) > 10000:
242248
# Keep only the most recent 5000 entries

0 commit comments

Comments
 (0)