Skip to main content

graphforge_storage/
io_stats.rs

1//! Process-global I/O counters for storage reads and staged rewrite commits.
2//! The read counters cover [`read_edges`](crate::catalog::read_edges),
3//! [`read_edges_filtered`](crate::catalog::read_edges_filtered), and
4//! [`read_nodes`](crate::catalog::read_nodes).
5//!
6//! # Why
7//! These prove the M15 T1 criterion (#767): with the adjacency index present,
8//! variable-length traversal must not scan the full edge file — it issues only
9//! an `edge_id`-filtered read whose materialized row count is proportional to
10//! the traversed neighborhood, independent of the total edge count. A scan
11//! over the index path can then assert `edge_full_reads == 0`, while the
12//! scan-build baseline shows `edge_full_rows >= total_edges`.
13//!
14//! # Semantics
15//! - A **full read** is one decode of a whole file: `read_edges` / `read_nodes`,
16//!   plus the [`read_edges_filtered`](crate::catalog::read_edges_filtered)
17//!   *fallback* (a requested id set covering more than half the file reads it
18//!   whole, then trims in memory — it is a full scan and is counted as one, so
19//!   it cannot hide behind the filtered API).
20//! - A **filtered read** is the predicate-pushdown path of
21//!   [`read_edges_filtered`](crate::catalog::read_edges_filtered); its row count
22//!   is the rows actually materialized after row-group and row-filter pruning.
23//! - `rows` count rows actually returned to the caller. A missing or empty file
24//!   still counts as one read of zero rows (the reader was invoked).
25//! - A **rewrite commit** is one successful, non-empty
26//!   [`RewriteBatch`](crate::RewriteBatch) commit, regardless of how many files
27//!   were staged in that batch.
28//!
29//! # Caveats
30//! Counters are process-global and aggregate across threads *and* queries
31//! (`read_nodes` is also used by the writer/mutator). They are advisory
32//! instrumentation, not per-query state. A test that asserts on them must
33//! [`reset`] immediately before the measured operation and keep that section
34//! single-threaded — the counters cannot attribute concurrent work. Each
35//! increment is one relaxed atomic add, negligible against a Parquet decode.
36
37use std::sync::atomic::{AtomicU64, Ordering};
38
39/// Table family involved in one filtered topology read.
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41#[doc(hidden)]
42pub enum FilteredReadTable {
43    /// Relationship topology.
44    Edge,
45    /// Node topology.
46    Node,
47}
48
49/// Storage strategy used for one filtered topology read.
50#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51#[doc(hidden)]
52pub enum FilteredReadStrategy {
53    /// Exact row ordinals derived from a proven dense node-id layout.
54    DenseRowSelection,
55    /// Conservative row-group pruning plus a Parquet row predicate.
56    RowGroupPredicate,
57    /// More than half the file was requested, so the whole file was read.
58    FullFallback,
59}
60
61/// Aggregate-only pruning work for one physical filtered-read attempt.
62#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63#[doc(hidden)]
64pub struct FilteredReadPruning {
65    /// Strategy used by the attempt.
66    pub strategy: FilteredReadStrategy,
67    /// Row groups whose metadata was considered.
68    pub row_groups_considered: u64,
69    /// Row groups retained for decoding.
70    pub row_groups_selected: u64,
71    /// Key-column pages whose metadata was considered.
72    pub pages_considered: u64,
73    /// Key-column pages containing at least one selected row.
74    pub pages_selected: u64,
75    /// Exact row ordinals selected before the membership guard.
76    pub exact_rows_selected: u64,
77    /// Dense selection was unavailable because its metadata contract failed.
78    pub metadata_fallbacks: u64,
79    /// Dense output validation failed and triggered a conservative retry.
80    pub validation_fallbacks: u64,
81}
82
83/// Optional observer for attributing filtered-read work to a physical operator.
84///
85/// The normal storage API installs no observer. Traversal diagnostics use this
86/// hook to distinguish concurrent hops without recording requested ids, paths,
87/// or graph contents.
88#[doc(hidden)]
89pub trait FilteredReadObserver: Send + Sync {
90    /// A physical Parquet read is about to be opened.
91    fn read_started(&self, table: FilteredReadTable);
92
93    /// Rows evaluated by the Parquet predicate/page-index path.
94    fn rows_scanned(&self, table: FilteredReadTable, rows: u64);
95
96    /// A physical read completed, with its returned row count and whether it
97    /// used the full-read fallback.
98    fn read_completed(&self, table: FilteredReadTable, rows: u64, full: bool);
99
100    /// A started physical read failed before completion.
101    fn read_failed(&self, table: FilteredReadTable);
102
103    /// Aggregate pruning work for a completed physical attempt.
104    fn pruning(&self, _table: FilteredReadTable, _pruning: FilteredReadPruning) {}
105}
106
107static EDGE_FULL_READS: AtomicU64 = AtomicU64::new(0);
108static EDGE_FULL_ROWS: AtomicU64 = AtomicU64::new(0);
109static EDGE_FILTERED_READS: AtomicU64 = AtomicU64::new(0);
110static EDGE_FILTERED_ROWS: AtomicU64 = AtomicU64::new(0);
111static NODE_FULL_READS: AtomicU64 = AtomicU64::new(0);
112static NODE_FULL_ROWS: AtomicU64 = AtomicU64::new(0);
113static NODE_FILTERED_READS: AtomicU64 = AtomicU64::new(0);
114static NODE_FILTERED_ROWS: AtomicU64 = AtomicU64::new(0);
115static EDGE_SCANNED_ROWS: AtomicU64 = AtomicU64::new(0);
116static NODE_SCANNED_ROWS: AtomicU64 = AtomicU64::new(0);
117static NODE_DENSE_ROW_SELECTION_READS: AtomicU64 = AtomicU64::new(0);
118static NODE_ROW_GROUP_PREDICATE_READS: AtomicU64 = AtomicU64::new(0);
119static NODE_ROW_GROUPS_CONSIDERED: AtomicU64 = AtomicU64::new(0);
120static NODE_ROW_GROUPS_SELECTED: AtomicU64 = AtomicU64::new(0);
121static NODE_PAGES_CONSIDERED: AtomicU64 = AtomicU64::new(0);
122static NODE_PAGES_SELECTED: AtomicU64 = AtomicU64::new(0);
123static NODE_EXACT_ROWS_SELECTED: AtomicU64 = AtomicU64::new(0);
124static NODE_METADATA_FALLBACKS: AtomicU64 = AtomicU64::new(0);
125static NODE_VALIDATION_FALLBACKS: AtomicU64 = AtomicU64::new(0);
126static REWRITE_COMMITS: AtomicU64 = AtomicU64::new(0);
127
128/// A point-in-time copy of the process-global I/O counters. Difference two
129/// snapshots — or [`reset`] then [`snapshot`] — to attribute work to a region.
130#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
131pub struct IoSnapshot {
132    /// Full edge-file reads: `read_edges` plus the filtered-read fallback.
133    pub edge_full_reads: u64,
134    /// Rows returned by those full edge reads.
135    pub edge_full_rows: u64,
136    /// `edge_id`-filtered edge reads that took the predicate-pushdown path.
137    pub edge_filtered_reads: u64,
138    /// Rows materialized by those filtered reads (post-pruning).
139    pub edge_filtered_rows: u64,
140    /// Full node-file reads: `read_nodes` plus the filtered-read fallback.
141    pub node_full_reads: u64,
142    /// Rows returned by those full node reads.
143    pub node_full_rows: u64,
144    /// `node_id`-filtered node reads that took the predicate-pushdown path
145    /// (`read_nodes_filtered`, #838).
146    pub node_filtered_reads: u64,
147    /// Rows materialized by those filtered node reads (post-pruning).
148    pub node_filtered_rows: u64,
149    /// Edge rows the pushdown predicate actually evaluated — i.e. rows in the
150    /// data pages the page index did **not** skip. The decode-cost proxy: for a
151    /// clustered (localized) id set this is a few pages regardless of total file
152    /// size; for a scattered set it approaches the whole file. (#838)
153    pub edge_scanned_rows: u64,
154    /// Node rows the pushdown predicate evaluated (pages not page-index-skipped).
155    pub node_scanned_rows: u64,
156    /// Node reads that used exact dense row selection.
157    pub node_dense_row_selection_reads: u64,
158    /// Node reads that used conservative row-group and predicate pruning.
159    pub node_row_group_predicate_reads: u64,
160    /// Node row groups whose metadata was considered for pruning.
161    pub node_row_groups_considered: u64,
162    /// Node row groups retained for decoding.
163    pub node_row_groups_selected: u64,
164    /// Node-id pages considered by exact dense selection.
165    pub node_pages_considered: u64,
166    /// Node-id pages containing selected rows.
167    pub node_pages_selected: u64,
168    /// Exact node row ordinals selected before the membership guard.
169    pub node_exact_rows_selected: u64,
170    /// Dense selection attempts rejected by the metadata contract.
171    pub node_metadata_fallbacks: u64,
172    /// Dense selection attempts rejected by post-read validation.
173    pub node_validation_fallbacks: u64,
174    /// Successful non-empty [`RewriteBatch`](crate::RewriteBatch) commits.
175    /// This counts persistence cycles, not the number of files in a batch.
176    pub rewrite_commits: u64,
177}
178
179/// Capture the current process-global counters.
180#[must_use]
181pub fn snapshot() -> IoSnapshot {
182    IoSnapshot {
183        edge_full_reads: EDGE_FULL_READS.load(Ordering::Relaxed),
184        edge_full_rows: EDGE_FULL_ROWS.load(Ordering::Relaxed),
185        edge_filtered_reads: EDGE_FILTERED_READS.load(Ordering::Relaxed),
186        edge_filtered_rows: EDGE_FILTERED_ROWS.load(Ordering::Relaxed),
187        node_full_reads: NODE_FULL_READS.load(Ordering::Relaxed),
188        node_full_rows: NODE_FULL_ROWS.load(Ordering::Relaxed),
189        node_filtered_reads: NODE_FILTERED_READS.load(Ordering::Relaxed),
190        node_filtered_rows: NODE_FILTERED_ROWS.load(Ordering::Relaxed),
191        edge_scanned_rows: EDGE_SCANNED_ROWS.load(Ordering::Relaxed),
192        node_scanned_rows: NODE_SCANNED_ROWS.load(Ordering::Relaxed),
193        node_dense_row_selection_reads: NODE_DENSE_ROW_SELECTION_READS.load(Ordering::Relaxed),
194        node_row_group_predicate_reads: NODE_ROW_GROUP_PREDICATE_READS.load(Ordering::Relaxed),
195        node_row_groups_considered: NODE_ROW_GROUPS_CONSIDERED.load(Ordering::Relaxed),
196        node_row_groups_selected: NODE_ROW_GROUPS_SELECTED.load(Ordering::Relaxed),
197        node_pages_considered: NODE_PAGES_CONSIDERED.load(Ordering::Relaxed),
198        node_pages_selected: NODE_PAGES_SELECTED.load(Ordering::Relaxed),
199        node_exact_rows_selected: NODE_EXACT_ROWS_SELECTED.load(Ordering::Relaxed),
200        node_metadata_fallbacks: NODE_METADATA_FALLBACKS.load(Ordering::Relaxed),
201        node_validation_fallbacks: NODE_VALIDATION_FALLBACKS.load(Ordering::Relaxed),
202        rewrite_commits: REWRITE_COMMITS.load(Ordering::Relaxed),
203    }
204}
205
206/// Reset every counter to zero. Call immediately before a measured operation.
207pub fn reset() {
208    for c in [
209        &EDGE_FULL_READS,
210        &EDGE_FULL_ROWS,
211        &EDGE_FILTERED_READS,
212        &EDGE_FILTERED_ROWS,
213        &NODE_FULL_READS,
214        &NODE_FULL_ROWS,
215        &NODE_FILTERED_READS,
216        &NODE_FILTERED_ROWS,
217        &EDGE_SCANNED_ROWS,
218        &NODE_SCANNED_ROWS,
219        &NODE_DENSE_ROW_SELECTION_READS,
220        &NODE_ROW_GROUP_PREDICATE_READS,
221        &NODE_ROW_GROUPS_CONSIDERED,
222        &NODE_ROW_GROUPS_SELECTED,
223        &NODE_PAGES_CONSIDERED,
224        &NODE_PAGES_SELECTED,
225        &NODE_EXACT_ROWS_SELECTED,
226        &NODE_METADATA_FALLBACKS,
227        &NODE_VALIDATION_FALLBACKS,
228        &REWRITE_COMMITS,
229    ] {
230        c.store(0, Ordering::Relaxed);
231    }
232}
233
234pub(crate) fn record_edge_full_read(rows: u64) {
235    EDGE_FULL_READS.fetch_add(1, Ordering::Relaxed);
236    EDGE_FULL_ROWS.fetch_add(rows, Ordering::Relaxed);
237}
238
239pub(crate) fn record_edge_filtered_read(rows: u64) {
240    EDGE_FILTERED_READS.fetch_add(1, Ordering::Relaxed);
241    EDGE_FILTERED_ROWS.fetch_add(rows, Ordering::Relaxed);
242}
243
244pub(crate) fn record_node_full_read(rows: u64) {
245    NODE_FULL_READS.fetch_add(1, Ordering::Relaxed);
246    NODE_FULL_ROWS.fetch_add(rows, Ordering::Relaxed);
247}
248
249pub(crate) fn record_node_filtered_read(rows: u64) {
250    NODE_FILTERED_READS.fetch_add(1, Ordering::Relaxed);
251    NODE_FILTERED_ROWS.fetch_add(rows, Ordering::Relaxed);
252}
253
254pub(crate) fn record_edge_scanned(rows: u64) {
255    EDGE_SCANNED_ROWS.fetch_add(rows, Ordering::Relaxed);
256}
257
258pub(crate) fn record_node_scanned(rows: u64) {
259    NODE_SCANNED_ROWS.fetch_add(rows, Ordering::Relaxed);
260}
261
262pub(crate) fn record_node_pruning(pruning: FilteredReadPruning) {
263    match pruning.strategy {
264        FilteredReadStrategy::DenseRowSelection => {
265            NODE_DENSE_ROW_SELECTION_READS.fetch_add(1, Ordering::Relaxed);
266        }
267        FilteredReadStrategy::RowGroupPredicate => {
268            NODE_ROW_GROUP_PREDICATE_READS.fetch_add(1, Ordering::Relaxed);
269        }
270        FilteredReadStrategy::FullFallback => {}
271    }
272    NODE_ROW_GROUPS_CONSIDERED.fetch_add(pruning.row_groups_considered, Ordering::Relaxed);
273    NODE_ROW_GROUPS_SELECTED.fetch_add(pruning.row_groups_selected, Ordering::Relaxed);
274    NODE_PAGES_CONSIDERED.fetch_add(pruning.pages_considered, Ordering::Relaxed);
275    NODE_PAGES_SELECTED.fetch_add(pruning.pages_selected, Ordering::Relaxed);
276    NODE_EXACT_ROWS_SELECTED.fetch_add(pruning.exact_rows_selected, Ordering::Relaxed);
277    NODE_METADATA_FALLBACKS.fetch_add(pruning.metadata_fallbacks, Ordering::Relaxed);
278    NODE_VALIDATION_FALLBACKS.fetch_add(pruning.validation_fallbacks, Ordering::Relaxed);
279}
280
281pub(crate) fn record_rewrite_commit() {
282    REWRITE_COMMITS.fetch_add(1, Ordering::Relaxed);
283}