Core: Data¶
All dataset, preprocessing, representation, and discovery machinery.
Current layout¶
datasets/- raw source adapters and dataset/source builders.discovery/- signal profiles, canonical entities, and provisional hypotheses for cross-vehicle alignment.preprocessing/- explicit representations, views, segments, materialization, PyG packing, temporal streams, scaler config, vocab config, and graph transforms.datamodule/- training-time loaders and batching policy.state.py- process-local dataset state for reuse within one Python process.
What to read next¶
graphids.core.data¶
data ¶
Data-layer public API.
The runtime datamodules are imported lazily so preprocessing and discovery can be used without importing optional training dependencies.
CANBusTemporalSource
dataclass
¶
CANBusTemporalSource(name: str, lake_root: str | None = None, val_fraction: float = 0.2, representation_cfg: TemporalRepresentationCfg = TemporalRepresentationCfg(), vocab_scope: str = 'train', train_source_mode: str = 'mixed', val_warmup_events: int = 0, test_warmup_events: int = 0)
Catalog to train/val/test CAN TemporalData cache builder.
DatasetState
dataclass
¶
Ready-to-serve train/val/test splits.
clear_cache ¶
get_or_build ¶
Return cached DatasetState for dataset.
Source code in graphids/core/data/state.py
datamodule ¶
DataModule primitives for temporal datasets.
TemporalDataModule ¶
TemporalDataModule(dataset, batch_size: int = 256, *, batch_mode: str = 'events', stream_lanes: int | None = None, chunk_size: int | None = None, num_workers: int = 0, pin_memory: bool = False, persistent_workers: bool = False)
Bases: LightningDataModule
Serve temporal event streams with PyG's TemporalDataLoader.
Source code in graphids/core/data/datamodule/temporal.py
stream_lanes ¶
Lane-batched temporal event iteration.
TemporalStreamLaneBatch
dataclass
¶
TemporalStreamLaneBatch(src: Tensor, dst: Tensor, t: Tensor, msg: Tensor, y: Tensor, attack_type: Tensor, event_id: Tensor, is_scored: Tensor, stream_id: Tensor, valid_mask: Tensor, lane_reset: Tensor, lane_stream_end: Tensor, reset_after: Tensor)
A fixed-lane batch of independent temporal streams.
Event tensors are shaped [lanes, chunk_size, ...]. valid_mask marks
real events; padded slots carry neutral values and must not affect loss or
metrics.
TemporalStreamLaneLoader ¶
Deterministically pack independent streams into fixed temporal lanes.
Source code in graphids/core/data/datamodule/stream_lanes.py
temporal ¶
Lightning data module for temporal PyG event streams.
TemporalDataModule ¶
TemporalDataModule(dataset, batch_size: int = 256, *, batch_mode: str = 'events', stream_lanes: int | None = None, chunk_size: int | None = None, num_workers: int = 0, pin_memory: bool = False, persistent_workers: bool = False)
Bases: LightningDataModule
Serve temporal event streams with PyG's TemporalDataLoader.
Source code in graphids/core/data/datamodule/temporal.py
datasets ¶
CANBusTemporalSource
dataclass
¶
CANBusTemporalSource(name: str, lake_root: str | None = None, val_fraction: float = 0.2, representation_cfg: TemporalRepresentationCfg = TemporalRepresentationCfg(), vocab_scope: str = 'train', train_source_mode: str = 'mixed', val_warmup_events: int = 0, test_warmup_events: int = 0)
Catalog to train/val/test CAN TemporalData cache builder.
can_bus ¶
CAN bus temporal dataset adapter and row schema.
CANBusTemporalSource
dataclass
¶
CANBusTemporalSource(name: str, lake_root: str | None = None, val_fraction: float = 0.2, representation_cfg: TemporalRepresentationCfg = TemporalRepresentationCfg(), vocab_scope: str = 'train', train_source_mode: str = 'mixed', val_warmup_events: int = 0, test_warmup_events: int = 0)
Catalog to train/val/test CAN TemporalData cache builder.
infer_attack_type ¶
Infer the attack code from filename/path substrings.
Source code in graphids/core/data/datasets/can_bus.py
load_can_rows ¶
Load, normalize, and parse raw CAN CSVs from source dirs.
Source code in graphids/core/data/datasets/can_bus.py
parse_payload ¶
Hex payload to byte_0..7 plus Shannon entropy.
Source code in graphids/core/data/datasets/can_bus.py
discovery ¶
Signal profile artifacts for CAN cache builds.
build_signal_profiles ¶
Aggregate raw CAN rows into one profile per vehicle/arbitration ID.
Source code in graphids/core/data/discovery/hypotheses.py
initialize_hypotheses ¶
Create empty editable mapping rows for profile review.
Source code in graphids/core/data/discovery/hypotheses.py
hypotheses ¶
Signal profile artifacts written beside graph caches.
build_signal_profiles ¶
Aggregate raw CAN rows into one profile per vehicle/arbitration ID.
Source code in graphids/core/data/discovery/hypotheses.py
initialize_hypotheses ¶
Create empty editable mapping rows for profile review.
Source code in graphids/core/data/discovery/hypotheses.py
preprocessing ¶
Core temporal preprocessing.
add_temporal_split_masks ¶
add_temporal_split_masks(table: DataFrame, *, split_name: str, warmup_events: int = 0) -> pl.DataFrame
Attach split identity plus warmup/scoring masks to an event table.
Source code in graphids/core/data/preprocessing/temporal.py
assert_temporal_splits_disjoint ¶
Raise if any split shares event ids with another split.
Source code in graphids/core/data/preprocessing/temporal.py
build_temporal_event_table ¶
build_temporal_event_table(df: DataFrame, *, id_col: str = 'node_id', raw_id_col: str = 'arb_id', timestamp_col: str = 'timestamp', num_unknown_buckets: int = 128) -> pl.DataFrame
Build one temporal event per row, preserving stream provenance.
The first row in each stream is represented as a self-event
(src_id == dst_id) so event counts match normalized row counts.
Source code in graphids/core/data/preprocessing/temporal.py
62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 | |
prepare_temporal_eval_table ¶
prepare_temporal_eval_table(table: DataFrame, *, split_name: str = 'test', warmup_events: int = 0) -> pl.DataFrame
Prepare validation/test-style full streams with warmup/scoring masks.
Source code in graphids/core/data/preprocessing/temporal.py
split_temporal_train_val_tables ¶
split_temporal_train_val_tables(table: DataFrame, *, val_fraction: float, val_warmup_events: int = 0) -> tuple[pl.DataFrame, pl.DataFrame]
Chronologically split each stream into train and validation intervals.
Source code in graphids/core/data/preprocessing/temporal.py
temporal_to_pyg ¶
Pack a temporal event table into PyG TemporalData tensors.
Source code in graphids/core/data/preprocessing/temporal.py
representations ¶
Data representation configs used by preprocessing.
scaler ¶
Per-column feature scalers for tensor-based graph preprocessing.
temporal ¶
Temporal event materialization for normalized CAN rows.
add_temporal_split_masks ¶
add_temporal_split_masks(table: DataFrame, *, split_name: str, warmup_events: int = 0) -> pl.DataFrame
Attach split identity plus warmup/scoring masks to an event table.
Source code in graphids/core/data/preprocessing/temporal.py
assert_temporal_splits_disjoint ¶
Raise if any split shares event ids with another split.
Source code in graphids/core/data/preprocessing/temporal.py
build_temporal_event_table ¶
build_temporal_event_table(df: DataFrame, *, id_col: str = 'node_id', raw_id_col: str = 'arb_id', timestamp_col: str = 'timestamp', num_unknown_buckets: int = 128) -> pl.DataFrame
Build one temporal event per row, preserving stream provenance.
The first row in each stream is represented as a self-event
(src_id == dst_id) so event counts match normalized row counts.
Source code in graphids/core/data/preprocessing/temporal.py
62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 | |
mark_split_start_self_events ¶
Remove transition features that cross into the start of a split.
Source code in graphids/core/data/preprocessing/temporal.py
mark_terminal_reset ¶
Ensure the final event in each sliced stream resets downstream state.
Source code in graphids/core/data/preprocessing/temporal.py
prepare_temporal_eval_table ¶
prepare_temporal_eval_table(table: DataFrame, *, split_name: str = 'test', warmup_events: int = 0) -> pl.DataFrame
Prepare validation/test-style full streams with warmup/scoring masks.
Source code in graphids/core/data/preprocessing/temporal.py
split_temporal_train_val_tables ¶
split_temporal_train_val_tables(table: DataFrame, *, val_fraction: float, val_warmup_events: int = 0) -> tuple[pl.DataFrame, pl.DataFrame]
Chronologically split each stream into train and validation intervals.
Source code in graphids/core/data/preprocessing/temporal.py
temporal_to_pyg ¶
Pack a temporal event table into PyG TemporalData tensors.
Source code in graphids/core/data/preprocessing/temporal.py
vocab ¶
Vocabulary scan, digest, persist, and load primitives.
load_vocab ¶
Return (entries, digest) from a persisted vocab file.
persist_vocab ¶
Atomic write and return the digest.
Source code in graphids/core/data/preprocessing/vocab.py
scan_arb_ids ¶
Sorted unique arb_id across every CSV under source_dirs.
Source code in graphids/core/data/preprocessing/vocab.py
vocab_digest ¶
SHA256 over (id, index) pairs sorted by index.