Skip to content

data.wrangler

Turn candles into validated Bar streams. Makes you declare whether your timestamps mean the bar's open or its close.

wrangler

Wrangle user-supplied OHLCV rows into validated, time-ordered Bar tuples.

The single most dangerous bug in bar-based backtesting is the silent off-by-one-bar look-ahead: treating a bar's OPEN timestamp as if it were its CLOSE (or vice versa) hands the strategy one bar of the future. This module kills that bug at the API level — the caller MUST declare what the source timestamp means via stamp:

  • stamp="open": ts_event = ts and ts_init = ts + step
  • stamp="close": ts_init = ts and ts_event = ts - step

where step = unit x unit_number in nanoseconds. There is no default.

Other hard guarantees
  • NAIVE timestamps are rejected with an error naming the fix — never guessed.
  • Prices convert via str() -> Decimal (never float -> Decimal) and must land exactly on the instrument's tick grid (RowOffGridError with the offending row index otherwise). NaN/Infinity prices and negative volumes are rejected with the row context.
  • Only fixed-span intraday units (SECOND/MINUTE/HOUR) are accepted: a Globex trading day is 23 hours, so DAY-and-above bars belong to the session-aware calendar-resampling layer, not a fixed nanosecond step.
  • Output is sorted ascending by ts_init (stable), ready for a feed.

pandas is imported lazily — only bars_from_dataframe needs it.

RowOffGridError

RowOffGridError(price: Decimal, tick_size: Decimal, *, row: int, field: str)

Bases: OffGridError

An input price is off the tick grid; carries the offending row index.

Source code in src/topstep_backtest/data/wrangler.py
def __init__(self, price: Decimal, tick_size: Decimal, *, row: int, field: str) -> None:
    super().__init__(price, tick_size)
    self.row = row
    self.field = field
    # Re-point the message at the row context (OffGridError.__init__ set a
    # generic one); ``price``/``tick_size`` attributes remain intact.
    self.args = (f"row {row}: {field}={price} is not on the {tick_size} tick grid",)

row instance-attribute

row = row

field instance-attribute

field = field

args instance-attribute

args = (f'row {row}: {field}={price} is not on the {tick_size} tick grid',)

step_ns

step_ns(unit: AggregateBarUnit, unit_number: int) -> int

The bar step (open -> close span) in nanoseconds.

Supports SECOND / MINUTE / HOUR only. TICK bars have no fixed time span. DAY / WEEK / MONTH bars are session-scoped, not fixed-span: a Globex trading day runs 23 hours (18:00 ET -> 17:00 ET), so a fixed 86,400s step would mis-stamp every daily bar and mis-attribute its trading day. Session-aware daily bars arrive with the calendar-resampling layer — supply intraday bars here.

Source code in src/topstep_backtest/data/wrangler.py
def step_ns(unit: AggregateBarUnit, unit_number: int) -> int:
    """The bar step (open -> close span) in nanoseconds.

    Supports SECOND / MINUTE / HOUR only. TICK bars have no fixed time span.
    DAY / WEEK / MONTH bars are session-scoped, not fixed-span: a Globex
    trading day runs 23 hours (18:00 ET -> 17:00 ET), so a fixed 86,400s step
    would mis-stamp every daily bar and mis-attribute its trading day.
    Session-aware daily bars arrive with the calendar-resampling layer —
    supply intraday bars here.
    """
    if unit_number < 1:
        raise ValueError(f"unit_number must be >= 1, got {unit_number}")
    per_unit = _STEP_NS_PER_UNIT.get(unit)
    if per_unit is None:
        if unit in _SESSION_UNITS:
            raise ValueError(
                f"unsupported bar unit {unit!r}: a Globex trading day is 23 hours "
                "(18:00 ET -> 17:00 ET), so DAY-and-above bars have no fixed "
                "nanosecond step; session-aware daily bars arrive with the "
                "calendar-resampling layer — supply intraday (SECOND/MINUTE/HOUR) "
                "bars instead"
            )
        raise ValueError(
            f"unsupported bar unit {unit!r}: only SECOND/MINUTE/HOUR have a fixed time step"
        )
    return per_unit * unit_number

bars_from_records

bars_from_records(rows: Iterable[tuple[object, ...]], *, contract_id: str, spec: InstrumentSpec, unit: AggregateBarUnit, unit_number: int, stamp: Literal['open', 'close']) -> tuple[Bar, ...]

Build tick-grid-validated Bar objects from (ts, o, h, l, c, v) rows.

stamp declares what the source timestamp means — see module docstring. Rows are sorted ascending by the resulting ts_init (stable sort), so the output is feed-ready regardless of input order.

Source code in src/topstep_backtest/data/wrangler.py
def bars_from_records(
    rows: Iterable[tuple[object, ...]],
    *,
    contract_id: str,
    spec: InstrumentSpec,
    unit: AggregateBarUnit,
    unit_number: int,
    stamp: Literal["open", "close"],
) -> tuple[Bar, ...]:
    """Build tick-grid-validated ``Bar`` objects from ``(ts, o, h, l, c, v)`` rows.

    ``stamp`` declares what the source timestamp means — see module docstring.
    Rows are sorted ascending by the resulting ``ts_init`` (stable sort), so
    the output is feed-ready regardless of input order.
    """
    if stamp not in ("open", "close"):
        raise ValueError(f'stamp must be "open" or "close", got {stamp!r}')
    step = step_ns(unit, unit_number)
    bar_type = BarType(contract_id=contract_id, unit=unit, unit_number=unit_number)

    bars: list[Bar] = []
    for row_index, row in enumerate(rows):
        if len(row) != 6:
            raise ValueError(
                f"row {row_index}: expected 6 fields (ts, open, high, low, close, "
                f"volume), got {len(row)}"
            )
        ts = _ts_to_ns(row[0], row=row_index)
        prices: dict[str, Decimal] = {}
        for field, raw in zip(_OHLCV_FIELDS[:4], row[1:5], strict=True):
            price = _price_to_decimal(raw, row=row_index, field=field)
            if not is_on_grid(price, spec.tick_size):
                raise RowOffGridError(price, spec.tick_size, row=row_index, field=field)
            prices[field] = price
        volume = _volume_to_int(row[5], row=row_index)

        if stamp == "open":
            ts_event, ts_init = ts, ts + step
        else:
            ts_event, ts_init = ts - step, ts
        bars.append(
            Bar(
                bar_type=bar_type,
                ts_event=ts_event,
                ts_init=ts_init,
                open=prices["open"],
                high=prices["high"],
                low=prices["low"],
                close=prices["close"],
                volume=volume,
            )
        )
    bars.sort(key=lambda b: b.ts_init)
    return tuple(bars)

bars_from_dataframe

bars_from_dataframe(df: Any, *, contract_id: str, spec: InstrumentSpec, unit: AggregateBarUnit, unit_number: int, stamp: Literal['open', 'close']) -> tuple[Bar, ...]

Build Bar objects from a pandas DataFrame of OHLCV candles.

The timestamp may be a column (case-insensitive: timestamp / ts / time / datetime / ts_event) or the index — but the index is used ONLY when it is a DatetimeIndex or is named (case-insensitively) after one of those timestamp columns. A default RangeIndex (or any anonymous integer index) is rejected: its 0, 1, 2, ... would otherwise be read as epoch nanoseconds and every bar would silently land in 1970. OHLCV column names are matched case-insensitively. Everything else — stamping, naive-timestamp rejection, tick-grid enforcement, sorting — is delegated to bars_from_records.

Source code in src/topstep_backtest/data/wrangler.py
def bars_from_dataframe(
    df: Any,
    *,
    contract_id: str,
    spec: InstrumentSpec,
    unit: AggregateBarUnit,
    unit_number: int,
    stamp: Literal["open", "close"],
) -> tuple[Bar, ...]:
    """Build ``Bar`` objects from a pandas DataFrame of OHLCV candles.

    The timestamp may be a column (case-insensitive: timestamp / ts / time /
    datetime / ts_event) or the index — but the index is used ONLY when it is
    a ``DatetimeIndex`` or is named (case-insensitively) after one of those
    timestamp columns. A default ``RangeIndex`` (or any anonymous integer
    index) is rejected: its 0, 1, 2, ... would otherwise be read as epoch
    nanoseconds and every bar would silently land in 1970. OHLCV column names
    are matched case-insensitively. Everything else — stamping,
    naive-timestamp rejection, tick-grid enforcement, sorting — is delegated
    to ``bars_from_records``.
    """
    pd = _import_pandas()
    if not isinstance(df, pd.DataFrame):
        raise TypeError(f"expected a pandas DataFrame, got {type(df).__name__}")

    by_lower: dict[str, object] = {}
    for column in df.columns:
        by_lower.setdefault(str(column).lower(), column)

    ts_column = next((by_lower[a] for a in _TS_COLUMN_ALIASES if a in by_lower), None)
    if ts_column is not None:
        ts_values: list[object] = df[ts_column].tolist()
    else:
        index_name = str(df.index.name).lower() if df.index.name is not None else None
        if not isinstance(df.index, pd.DatetimeIndex) and (index_name not in _TS_COLUMN_ALIASES):
            raise ValueError(
                "no timestamp column found and the index "
                f"({type(df.index).__name__}) is neither a DatetimeIndex nor named "
                f"after one of the accepted timestamp columns "
                f"{list(_TS_COLUMN_ALIASES)}; refusing to interpret a plain "
                "integer index as epoch nanoseconds"
            )
        ts_values = df.index.tolist()

    missing = [f for f in _OHLCV_FIELDS if f not in by_lower]
    if missing:
        raise ValueError(
            f"DataFrame is missing required OHLCV columns {missing} "
            f"(case-insensitive); found columns {[str(c) for c in df.columns]}"
        )
    series: list[list[object]] = [df[by_lower[f]].tolist() for f in _OHLCV_FIELDS]

    rows = zip(ts_values, *series, strict=True)
    return bars_from_records(
        rows,
        contract_id=contract_id,
        spec=spec,
        unit=unit,
        unit_number=unit_number,
        stamp=stamp,
    )