Skip to content

Bar Aggregator

The Bar Aggregator system provides flexible bar generation from trade data with support for multiple bar types, timeframes, and custom aggregated structures.

Overview

The aggregator system is built around a policy-based design with zero-cost abstractions:

flowchart TB
    TE[TradeEvent] --> BA[BarAggregator]

    subgraph Policies
        direction LR
        Time[TimeBarPolicy]
        Tick[TickBarPolicy]
        Volume[VolumeBarPolicy]
        Renko[RenkoBarPolicy]
        Range[RangeBarPolicy]
    end

    BA -.-> Policies
    BA --> BE[BarEvent]

    BE --> BM[BarMatrix]
    BE --> Strategy[Strategy.onBar]

    BM --> Access["bars[symbol][timeframe][idx]"]

Quick Start

#include "lrvx/aggregator/bar_aggregator.h"
#include "lrvx/aggregator/bus/bar_bus.h"

// Create a 1-minute time bar aggregator
BarBus bus;
TimeBarAggregator aggregator(TimeBarPolicy(std::chrono::seconds(60)), &bus);

// Subscribe to bar events
bus.subscribe(&strategy);

// Start
bus.start();
aggregator.start();

// Feed trades
aggregator.onTrade(tradeEvent);

// Stop (flushes remaining bars)
aggregator.stop();
bus.stop();

Bar Types

Time Bars

Close after a fixed time interval.

TimeBarAggregator aggregator(TimeBarPolicy(std::chrono::seconds(60)), &bus);

Use cases: Traditional OHLCV charts, backtesting, most strategies.

A trade whose aligned bucket precedes the live bar is dropped and counted in lateTradeCount() rather than folded into the live bar; a trade that is merely out of order inside the live bucket is folded in as usual. See Bar types.

Tick Bars

Close after a fixed number of trades.

TickBarAggregator aggregator(TickBarPolicy(100), &bus);  // 100-tick bars

Use cases: High-frequency trading, eliminating time-based noise, volume-normalized analysis.

Volume Bars

Close after a fixed notional volume.

VolumeBarAggregator aggregator(VolumeBarPolicy::fromDouble(1000000.0), &bus);  // $1M bars

Use cases: Volume-weighted analysis, consistent information content per bar.

Renko Bars

Close when price moves by a fixed amount (brick size).

RenkoBarAggregator aggregator(RenkoBarPolicy::fromDouble(10.0), &bus);  // $10 bricks

Use cases: Trend following, noise elimination, support/resistance identification.

A brick closes at its boundary with the crossing trade counted in it, and the next brick opens at that boundary. A trade that jumps several brick widths fills in the bricks a continuous price path would have produced, bounded by RenkoBarPolicy::kMaxGapBricks; the bar that absorbs the remainder past the bound carries BarCloseReason::Gap. See Bar types.

Range Bars

Close when high-low range exceeds a threshold.

RangeBarAggregator aggregator(RangeBarPolicy::fromDouble(5.0), &bus);  // $5 range

Use cases: Volatility-based analysis, breakout detection.

Heikin-Ashi Bars

Smoothed candlesticks using averaged OHLC values:

HeikinAshiBarAggregator aggregator(HeikinAshiBarPolicy(std::chrono::seconds(60)), &bus);

Use cases: Trend identification, noise reduction, smoother signals.

Bar Structure

struct Bar {
  Price open, high, low, close;
  Volume volume;          // Notional volume (price * quantity)
  Volume buyVolume;       // Volume from buy trades (for delta calculation)
  Quantity tradeCount;    // Number of trades in bar
  TimePoint startTime;    // Bar open time
  TimePoint endTime;      // Bar close time
  BarCloseReason reason;  // Threshold, Gap or Forced; see reference/api/aggregator/bar.md
};

// Calculate delta (buy pressure - sell pressure)
Volume sellVolume = bar.volume - bar.buyVolume;
int64_t delta = bar.buyVolume.raw() - sellVolume.raw();

Volume's int64_t constructor is explicit, so a raw arithmetic result cannot be assigned to a Volume. Subtract the fixed-point values, or work in raw units as above.

Multi-Timeframe Analysis

Use MultiTimeframeAggregator to produce multiple timeframes from a single trade stream:

MultiTimeframeAggregator<4> aggregator(&bus);
aggregator.addTimeInterval(std::chrono::seconds(60));    // M1
aggregator.addTimeInterval(std::chrono::seconds(300));   // M5
aggregator.addTimeInterval(std::chrono::seconds(3600));  // H1
aggregator.addTickInterval(100);                          // 100-tick bars

bus.subscribe(&strategy);
bus.start();
aggregator.start();

// All trades go to all timeframes
aggregator.onTrade(trade);

Mixed Bar Types

You can mix time, tick, and volume bars in a single aggregator:

aggregator.addTimeInterval(std::chrono::seconds(60));
aggregator.addTickInterval(50);
aggregator.addVolumeInterval(100000.0);

Bar History: BarSeries

BarSeries is a ring buffer for storing bar history:

BarSeries<256> series;  // Last 256 bars

series.push(bar);

// Access (0 = newest, 1 = previous, etc.)
const Bar& latest = series[0];
const Bar& previous = series[1];

// Iteration (newest to oldest)
for (const auto& bar : series) {
  // ...
}

Multi-Symbol Multi-Timeframe: BarMatrix

BarMatrix provides O(1) access to bar history across symbols and timeframes:

A matrix holds MaxSymbols * MaxTimeframes * Depth bars, which is megabytes for any realistic configuration and about 40 MB with the default template arguments. Put it on the heap; a local declaration overflows the stack before the first line of the function body runs, and no compiler diagnostic warns about it. BarMatrix<...>::kStorageBytes gives the figure if you need it for a budget of your own, and an instantiation past LRVX_BAR_MATRIX_MAX_BYTES (256 MB by default) is a compile error.

// 256 symbols, 8 timeframes, 64 bars depth
auto matrix = std::make_unique<BarMatrix<256, 8, 64>>();

std::array<TimeframeId, 3> tfs = {timeframe::M1, timeframe::M5, timeframe::H1};
matrix->configure(tfs);

// Subscribe to receive bars
bus.subscribe(matrix.get());

// Access: matrix[symbol][timeframe][index]
const Bar* bar = matrix->bar(symbolId, timeframe::H1, 0);  // Latest H1 bar
const Bar* prev = matrix->bar(symbolId, timeframe::H1, 1); // Previous H1 bar

// Or by timeframe index
const Bar* bar = matrix->bar(symbolId, 0, 0);  // First configured timeframe

Warmup with Historical Data

std::vector<Bar> historicalBars = loadFromDatabase();
matrix->warmup(symbolId, timeframe::H1, historicalBars);

TimeframeId

TimeframeId encodes bar type and parameter:

// Presets
timeframe::M1   // 1 minute
timeframe::M5   // 5 minutes
timeframe::M15  // 15 minutes
timeframe::H1   // 1 hour
timeframe::H4   // 4 hours
timeframe::D1   // 1 day

// Custom
TimeframeId tf = TimeframeId::time(std::chrono::seconds(120));  // 2 minutes
TimeframeId tick = TimeframeId::tick(500);                       // 500 ticks
TimeframeId vol = TimeframeId::volume(1000000);                  // $1M volume

Custom Policies

Implement your own bar policy by satisfying the BarPolicy concept:

The concept requires four members: shouldClose, update, initBar, param(), plus a kBarType constant. All four functions must be noexcept. Omitting param() fails the constraint and BarAggregator<MyCustomPolicy> will not instantiate.

Two members are optional, declared as concepts in aggregator/aggregation_policy.h, and picked up by every aggregator at once — BarAggregator, MultiTimeframeAggregator, and the batch aggregators behind the C ABI and the Python bindings:

Member Concept Effect
bool isLate(const TradeEvent&, const Bar&) const noexcept DetectsLateTrades Returning true drops the trade instead of folding it in, and increments lateTradeCount(). Only TimeBarPolicy declares it.
template <typename Emit> void closeAndReopen(const TradeEvent&, Bar&, Emit&&) const ClosesAndReopens Replaces the default close (publish the bar as it stands, then initBar at the trade price). The policy publishes every bar itself through Emit and leaves the Bar& re-initialized as the one that opens next. Only RenkoBarPolicy declares it, because a brick's close price and the next brick's open are the same boundary and have to be computed together.

kBarType must be one of the existing BarType values — there is no BarType::Custom. Pick the value that best describes the closing rule (Tick, Volume, Range, ...).

struct MyCustomPolicy {
  static constexpr BarType kBarType = BarType::Volume;

  explicit MyCustomPolicy(uint64_t threshold) noexcept : _threshold(threshold) {}

  bool shouldClose(const TradeEvent& trade, const Bar& bar) const noexcept {
    // Your closing logic
    return bar.volume.raw() >= static_cast<int64_t>(_threshold);
  }

  void update(const TradeEvent& trade, Bar& bar) noexcept {
    updateOHLCV(trade, bar);  // Helper for the OHLCV update
    // Additional custom updates
  }

  void initBar(const TradeEvent& trade, Bar& bar) noexcept {
    initBarFromTrade(trade, bar);  // Helper for initialization
    // Additional initialization
  }

  // Required by the concept: the type parameter carried on every BarEvent.
  uint64_t param() const noexcept { return _threshold; }

 private:
  uint64_t _threshold;
};

// Use with BarAggregator
BarAggregator<MyCustomPolicy> aggregator(MyCustomPolicy{1000000}, &bus);

The initialization helper is initBarFromTrade(const TradeEvent&, Bar&). There is no initializeBar.

BarEvent

Bar events contain full bar data plus metadata:

struct BarEvent {
  using Listener = IMarketDataSubscriber;

  SymbolId symbol{};
  InstrumentType instrument = InstrumentType::Spot;
  BarType barType{};
  uint64_t barTypeParam{};  // interval in NANOSECONDS, tick count, volume threshold
  Bar bar{};

  uint64_t tickSequence = 0;  // internal, set by bus
};

barTypeParam is uint64_t, and for BarType::Time it carries nanoseconds, not seconds. A 1-minute bar has barTypeParam == 60'000'000'000.

Strategy Integration

Using BarStrategy Helper

class MyStrategy : public BarStrategy<4> {
public:
  // Strategy's constructors both require a const SymbolRegistry&.
  using BarStrategy::BarStrategy;

protected:
  void onSymbolBar(SymbolContext& ctx, const BarEvent& ev) override {
    // Access bars via helper methods
    auto* h1 = bar(timeframe::H1, 0);
    auto* h1_prev = bar(timeframe::H1, 1);

    // Or use optional-returning methods
    auto closeOpt = close(timeframe::H1, 0);
    if (closeOpt && *closeOpt > *close(timeframe::H1, 1)) {
      // H1 closed higher
    }
  }
};

Strategy::onBar is final — it maintains the per-symbol context and bar rings before dispatching. Override the protected onSymbolBar(SymbolContext&, const BarEvent&) instead.

Manual Integration

Only a subscriber that is not a Strategy overrides onBar directly:

class MyBarSubscriber : public IMarketDataSubscriber {
  void onBar(const BarEvent& ev) override {
    // barTypeParam is nanoseconds for Time bars
    if (ev.barType == BarType::Time && ev.barTypeParam == 60'000'000'000ULL) {
      // Handle 1-minute bars
    }
  }
};

Performance

Operation Complexity
Policy shouldClose() O(1), inlined
Symbol lookup O(1) via SymbolStateMap
Bar history access O(1) ring buffer
Timeframe lookup O(n), n ≤ 8
Bar push O(1) amortized

Benchmark results (GCC 14, LTO, Release):

  • TimeBarAggregator.onTrade: ~45ns/trade
  • MultiTimeframeAggregator (4 TF): ~15ns/timeframe (~60ns total)
  • MultiTimeframeAggregator (8 TF): ~11ns/timeframe (~88ns total)
  • BarMatrix random access: ~5ns

Implementation Notes

MultiTimeframeAggregator uses tag + switch dispatch for policy execution, chosen after benchmarking against alternatives:

Dispatch Method ns/trade (4 TF) Notes
Tag + Switch ~60ns Winner: fastest across GCC/Clang/MSVC
std::variant/visit ~72ns 17-20% slower
Function pointers ~85ns Indirect call overhead

The implementation uses:

  • PolicyTag enum for runtime type discrimination
  • PolicyStorage union for type-safe storage without vtables
  • Single heap allocation for all slots (avoids stack overflow with large MaxTimeframes)
  • LRVX_FORCE_INLINE on hot path methods

Files

File Description
aggregator/bar.h Bar struct, BarType, BarCloseReason
aggregator/timeframe.h TimeframeId, presets
aggregator/aggregation_policy.h BarPolicy concept
aggregator/bar_aggregator.h BarAggregator template
aggregator/multi_timeframe_aggregator.h Multi-TF aggregator
aggregator/bar_series.h Ring buffer for history
aggregator/bar_matrix.h Multi-symbol multi-TF storage
aggregator/events/bar_event.h BarEvent struct
aggregator/bus/bar_bus.h EventBus
aggregator/policies/*.h Time, Tick, Volume, Renko, Range policies

See Also