antecedent.data

Data constructors and conversion probes (Arrow / float64 columns).

  1"""Data constructors and conversion probes (Arrow / float64 columns)."""
  2
  3from __future__ import annotations
  4
  5from dataclasses import dataclass
  6from typing import Any, Mapping, Sequence
  7
  8import numpy as np
  9from numpy.typing import NDArray
 10
 11from ._data import as_columns, as_multi_env_columns, to_f64
 12from ._native import ArrowLoadInfo, load_float64_arrow_c_columns, load_float64_columns
 13
 14
 15@dataclass(frozen=True)
 16class EventFrame:
 17    """Irregular event marks + timestamps; aligned via ``align_interval_ns`` before temporal algos."""
 18
 19    names: list[str]
 20    columns: list[NDArray[np.float64]]
 21    event_times_ns: NDArray[np.int64]
 22    align_interval_ns: int
 23
 24
 25@dataclass(frozen=True)
 26class PanelFrame:
 27    """Multi-unit time series sharing one schema.
 28
 29    Discovery: ``JPCMCIPlus`` (multi-env context) or pooled-units ``PCMCI`` /
 30    ``PCMCIPlus`` / ``LPCMCI``. Estimation uses stacked cluster-HAC SE
 31    (frequentist) or ``BayesianTemporalGcomp`` when ``inference=Bayesian``.
 32    """
 33
 34    names: list[str]
 35    unit_columns: list[list[NDArray[np.float64]]]
 36    unit_ids: list[int]
 37
 38
 39@dataclass(frozen=True)
 40class MultiEnvFrame:
 41    """Multi-environment series (J-PCMCI+ discovery)."""
 42
 43    names: list[str]
 44    env_columns: list[list[NDArray[np.float64]]]
 45
 46
 47def event(
 48    data: Mapping[str, Any] | Any,
 49    event_times_ns: Sequence[int] | NDArray[np.int64],
 50    *,
 51    align_interval_ns: int,
 52) -> EventFrame:
 53    """Build an [`EventFrame`] for ``analyze`` (duration-bin align → temporal path)."""
 54    if align_interval_ns <= 0:
 55        raise ValueError("align_interval_ns must be > 0")
 56    names, columns = as_columns(data)
 57    times = np.asarray(event_times_ns, dtype=np.int64)
 58    if times.ndim != 1:
 59        raise ValueError(f"event_times_ns must be 1-d, got shape {times.shape}")
 60    if len(times) != len(columns[0]):
 61        raise ValueError(
 62            f"event_times_ns length {len(times)} != column length {len(columns[0])}"
 63        )
 64    return EventFrame(
 65        names=names,
 66        columns=columns,
 67        event_times_ns=times,
 68        align_interval_ns=int(align_interval_ns),
 69    )
 70
 71
 72def panel(
 73    units: Sequence[Mapping[str, Any] | Any] | Mapping[Any, Mapping[str, Any] | Any],
 74) -> PanelFrame:
 75    """Build a [`PanelFrame`] from a sequence of unit frames or ``{unit_id: frame}``."""
 76    if isinstance(units, Mapping):
 77        ids = [int(k) for k in units.keys()]
 78        frames = list(units.values())
 79    else:
 80        frames = list(units)
 81        ids = list(range(len(frames)))
 82    if not frames:
 83        raise ValueError("panel needs ≥1 unit")
 84    names, first = as_columns(frames[0])
 85    unit_columns = [first]
 86    for i, frame in enumerate(frames[1:], start=1):
 87        n, cols = as_columns(frame)
 88        if n != names:
 89            raise ValueError(
 90                f"unit {i} column names {n!r} do not match unit 0 {names!r}"
 91            )
 92        unit_columns.append(cols)
 93    return PanelFrame(names=names, unit_columns=unit_columns, unit_ids=ids)
 94
 95
 96def multi_env(
 97    envs: Sequence[Mapping[str, Any] | Any],
 98) -> MultiEnvFrame:
 99    """Build a [`MultiEnvFrame`] (sequence of environment frames for J-PCMCI+)."""
100    names, env_columns = as_multi_env_columns(envs)
101    return MultiEnvFrame(names=names, env_columns=env_columns)
102
103
104__all__ = [
105    "ArrowLoadInfo",
106    "EventFrame",
107    "MultiEnvFrame",
108    "PanelFrame",
109    "event",
110    "load_float64_arrow_c_columns",
111    "load_float64_columns",
112    "multi_env",
113    "panel",
114    "to_f64",
115]
class ArrowLoadInfo:

Result of the conversion probe (same Arrow→tabular path as analyze/discover).

diagnostic_count
column_names

Schema names after library-owned ingestion (proves the batch was parsed).

bytes_borrowed
row_count
bytes_copied
column_count
@dataclass(frozen=True)
class EventFrame:
16@dataclass(frozen=True)
17class EventFrame:
18    """Irregular event marks + timestamps; aligned via ``align_interval_ns`` before temporal algos."""
19
20    names: list[str]
21    columns: list[NDArray[np.float64]]
22    event_times_ns: NDArray[np.int64]
23    align_interval_ns: int

Irregular event marks + timestamps; aligned via align_interval_ns before temporal algos.

names: list[str]
columns: list[NDArray[typing.Any]]
event_times_ns: NDArray[typing.Any]
align_interval_ns: int
@dataclass(frozen=True)
class MultiEnvFrame:
40@dataclass(frozen=True)
41class MultiEnvFrame:
42    """Multi-environment series (J-PCMCI+ discovery)."""
43
44    names: list[str]
45    env_columns: list[list[NDArray[np.float64]]]

Multi-environment series (J-PCMCI+ discovery).

names: list[str]
env_columns: list[list[NDArray[typing.Any]]]
@dataclass(frozen=True)
class PanelFrame:
26@dataclass(frozen=True)
27class PanelFrame:
28    """Multi-unit time series sharing one schema.
29
30    Discovery: ``JPCMCIPlus`` (multi-env context) or pooled-units ``PCMCI`` /
31    ``PCMCIPlus`` / ``LPCMCI``. Estimation uses stacked cluster-HAC SE
32    (frequentist) or ``BayesianTemporalGcomp`` when ``inference=Bayesian``.
33    """
34
35    names: list[str]
36    unit_columns: list[list[NDArray[np.float64]]]
37    unit_ids: list[int]

Multi-unit time series sharing one schema.

Discovery: JPCMCIPlus (multi-env context) or pooled-units PCMCI / PCMCIPlus / LPCMCI. Estimation uses stacked cluster-HAC SE (frequentist) or BayesianTemporalGcomp when inference=Bayesian.

names: list[str]
unit_columns: list[list[NDArray[typing.Any]]]
unit_ids: list[int]
def event( data: Union[Mapping[str, Any], Any], event_times_ns: Union[Sequence[int], NDArray[Any]], *, align_interval_ns: int) -> EventFrame:
48def event(
49    data: Mapping[str, Any] | Any,
50    event_times_ns: Sequence[int] | NDArray[np.int64],
51    *,
52    align_interval_ns: int,
53) -> EventFrame:
54    """Build an [`EventFrame`] for ``analyze`` (duration-bin align → temporal path)."""
55    if align_interval_ns <= 0:
56        raise ValueError("align_interval_ns must be > 0")
57    names, columns = as_columns(data)
58    times = np.asarray(event_times_ns, dtype=np.int64)
59    if times.ndim != 1:
60        raise ValueError(f"event_times_ns must be 1-d, got shape {times.shape}")
61    if len(times) != len(columns[0]):
62        raise ValueError(
63            f"event_times_ns length {len(times)} != column length {len(columns[0])}"
64        )
65    return EventFrame(
66        names=names,
67        columns=columns,
68        event_times_ns=times,
69        align_interval_ns=int(align_interval_ns),
70    )

Build an [EventFrame] for analyze (duration-bin align → temporal path).

def load_float64_arrow_c_columns(names, columns):

Load float64 columns from Arrow C Data Interface exporters (PyArrow / __arrow_c_array__).

Prefers zero-copy borrow of contiguous float64 value buffers.

def load_float64_columns(names, columns):

Conversion probe: NumPy → Arrow → library-owned tabular storage.

Shares the same ingestion path as analyze* / discover_*. The loaded table is not retained across the FFI boundary; call analysis APIs with the original NumPy columns.

def multi_env( envs: Sequence[Union[Mapping[str, Any], Any]]) -> MultiEnvFrame:
 97def multi_env(
 98    envs: Sequence[Mapping[str, Any] | Any],
 99) -> MultiEnvFrame:
100    """Build a [`MultiEnvFrame`] (sequence of environment frames for J-PCMCI+)."""
101    names, env_columns = as_multi_env_columns(envs)
102    return MultiEnvFrame(names=names, env_columns=env_columns)

Build a [MultiEnvFrame] (sequence of environment frames for J-PCMCI+).

def panel( units: Union[Sequence[Union[Mapping[str, Any], Any]], Mapping[Any, Union[Mapping[str, Any], Any]]]) -> PanelFrame:
73def panel(
74    units: Sequence[Mapping[str, Any] | Any] | Mapping[Any, Mapping[str, Any] | Any],
75) -> PanelFrame:
76    """Build a [`PanelFrame`] from a sequence of unit frames or ``{unit_id: frame}``."""
77    if isinstance(units, Mapping):
78        ids = [int(k) for k in units.keys()]
79        frames = list(units.values())
80    else:
81        frames = list(units)
82        ids = list(range(len(frames)))
83    if not frames:
84        raise ValueError("panel needs ≥1 unit")
85    names, first = as_columns(frames[0])
86    unit_columns = [first]
87    for i, frame in enumerate(frames[1:], start=1):
88        n, cols = as_columns(frame)
89        if n != names:
90            raise ValueError(
91                f"unit {i} column names {n!r} do not match unit 0 {names!r}"
92            )
93        unit_columns.append(cols)
94    return PanelFrame(names=names, unit_columns=unit_columns, unit_ids=ids)

Build a [PanelFrame] from a sequence of unit frames or {unit_id: frame}.

def to_f64(arr: Any) -> NDArray[numpy.float64]:
79def to_f64(arr: Any) -> NDArray[np.float64]:
80    a = np.asarray(arr, dtype=np.float64)
81    if a.ndim != 1:
82        raise ValueError(f"expected 1-d column, got shape {a.shape}")
83    if a.dtype == object:
84        raise TypeError("object-dtype columns are not supported")
85    return a