Skip to content

Write a Custom Connector

C++ extension point

Connectors are written against the C++ engine. Once a connector is built into lrvx, every binding (Python, Node.js, Codon, embedded JS) gets access to it automatically — no further work per language. If you just want to use an existing connector from Python or Node.js, see the Python / Node.js bindings docs instead.

Connect lrvx to a new exchange or data source.

Overview

A connector:

  1. Connects to an exchange API (WebSocket, REST, FIX, etc.)
  2. Parses incoming messages
  3. Converts to lrvx event types
  4. Emits events to the engine

1. Implement IExchangeConnector

Header: lrvx/connector/abstract_exchange_connector.h

#include "lrvx/connector/abstract_exchange_connector.h"
#include "lrvx/book/events/trade_event.h"
#include "lrvx/book/events/book_update_event.h"

class MyExchangeConnector : public lrvx::IExchangeConnector
{
public:
  MyExchangeConnector(const std::string& symbol, lrvx::SymbolId symbolId)
    : _symbol(symbol), _symbolId(symbolId) {}

  std::string exchangeId() const override {
    return "myexchange";
  }

  void start() override {
    _running = true;
    _thread = std::thread(&MyExchangeConnector::run, this);
  }

  void stop() override {
    _running = false;
    if (_thread.joinable()) {
      _thread.join();
    }
  }

private:
  void run() {
    // Connect to exchange
    connect();

    while (_running) {
      // Receive and parse messages
      auto msg = receiveMessage();
      handleMessage(msg);
    }

    disconnect();
  }

  void handleMessage(const Message& msg) {
    if (msg.type == MessageType::Trade) {
      handleTrade(msg);
    } else if (msg.type == MessageType::BookUpdate) {
      handleBookUpdate(msg);
    }
  }

  void handleTrade(const Message& msg) {
    lrvx::TradeEvent ev;
    ev.trade.symbol = _symbolId;
    ev.trade.price = lrvx::Price::fromDouble(msg.price);
    ev.trade.quantity = lrvx::Quantity::fromDouble(msg.qty);
    ev.trade.isBuy = msg.side == "buy";
    ev.trade.exchangeTsNs = msg.timestamp;
    ev.recvNs = lrvx::nowNsMonotonic();

    // Emit via base class method
    emitTrade(ev);
  }

  void handleBookUpdate(const Message& msg) {
    // For book updates, create event directly (simplified)
    // In practice, you might use a pool for large events

    lrvx::BookUpdateEvent ev(/* pmr allocator */);
    ev.update.symbol = _symbolId;
    ev.update.type = lrvx::BookUpdateType::SNAPSHOT;

    for (const auto& bid : msg.bids) {
      ev.update.bids.push_back({
        lrvx::Price::fromDouble(bid.price),
        lrvx::Quantity::fromDouble(bid.qty)
      });
    }
    for (const auto& ask : msg.asks) {
      ev.update.asks.push_back({
        lrvx::Price::fromDouble(ask.price),
        lrvx::Quantity::fromDouble(ask.qty)
      });
    }

    ev.recvNs = lrvx::nowNsMonotonic();
    ev.sourceExchange = _exchangeId;  // from registry->registerExchange(...)

    emitBookUpdate(ev);
  }

  std::string _symbol;
  lrvx::SymbolId _symbolId;
  std::atomic<bool> _running{false};
  std::thread _thread;
};

2. Wire Callbacks

Connect the connector to your buses:

#include "lrvx/book/bus/trade_bus.h"
#include "lrvx/book/bus/book_update_bus.h"
#include "lrvx/util/memory/pool.h"

// Create buses
auto tradeBus = std::make_unique<TradeBus>();
auto bookBus = std::make_unique<BookUpdateBus>();

// Create pool for book events
pool::Pool<BookUpdateEvent, 128> bookPool;

// Create connector
auto connector = std::make_shared<MyExchangeConnector>("BTCUSD", symbolId);

// Wire callbacks
connector->setCallbacks(
  // Book update callback
  [&bookBus, &bookPool](const BookUpdateEvent& ev) {
    // Acquire from pool for variable-size events
    if (auto handle = bookPool.acquire()) {
      (*handle)->update = ev.update;
      (*handle)->recvNs = ev.recvNs;
      bookBus->publish(std::move(handle));
    }
  },
  // Trade callback
  [&tradeBus](const TradeEvent& ev) {
    tradeBus->publish(ev);
  }
);

3. Register with Engine

std::vector<std::shared_ptr<IExchangeConnector>> connectors;
connectors.push_back(connector);

std::vector<std::unique_ptr<ISubsystem>> subsystems;
subsystems.push_back(std::move(tradeBus));
subsystems.push_back(std::move(bookBus));
// ... add strategies, etc.

EngineConfig config{};
Engine engine(config, std::move(subsystems), std::move(connectors));
engine.start();

4. Use ConnectorFactory (Optional)

For dynamic connector creation:

#include "lrvx/connector/connector_factory.h"

// Register factory
ConnectorFactory::instance().registerConnector("myexchange",
  [registry](const std::string& symbol) {
    auto symbolId = registry->getSymbolId("myexchange", symbol);
    return std::make_shared<MyExchangeConnector>(symbol, *symbolId);
  }
);

// Create connector
auto conn = ConnectorFactory::instance().createConnector("myexchange", "BTCUSD");

5. Best Practices

Thread Safety

  • Connector runs its own thread(s)
  • emitTrade() and emitBookUpdate() are thread-safe
  • Callbacks may execute on connector thread

Timestamps

Capture timestamps at the right points:

void handleMessage(const RawMessage& raw) {
  MonoNanos recvNs = nowNsMonotonic();  // Capture immediately

  // Parse message...
  auto parsed = parse(raw);

  TradeEvent ev;
  ev.trade.exchangeTsNs = parsed.exchangeTimestamp;  // From exchange
  ev.recvNs = recvNs;                                 // When we received it
  ev.publishTsNs = nowNsMonotonic();                  // Right before emit

  emitTrade(ev);
}

Source exchange

Every BookUpdateEvent must also name the venue it came from:

// In the constructor -- registerExchange is idempotent, and an id resolved
// lazily on the first frame is InvalidExchangeId for whatever ran before it.
_exchangeId = registry->registerExchange(exchangeId());

// On every book event
ev.sourceExchange = _exchangeId;

CompositeBookMatrix::onBookUpdate returns immediately for an update whose sourceExchange is out of range, so a connector that leaves the field at InvalidExchangeId contributes nothing to the cross-venue book — it is simply absent from best bid/ask, with no error anywhere. recvNs is the other half of the same contract: checkStaleness() skips any venue whose lastUpdateNs is still zero, so an unstamped feed that freezes keeps being quoted forever.

Error Handling

A log line is not a health signal: whoever supervises the feed is in another thread, usually another subsystem, and reads setErrorCallbacks. Report the close as well as recovering from it — the in-tree connectors all honour the same contract, described in Connectors.

void run() {
  while (_running) {
    try {
      if (!_connected) {
        connect();
      }
      auto msg = receiveMessage();
      handleMessage(msg);
    } catch (const ConnectionError& e) {
      LRVX_LOG("Connection lost: " << e.what());
      _connected = false;
      emitDisconnect(e.what());  // tell the supervisor, then recover
      reconnect();
    }
  }
}

Sequence Numbers

Track exchange sequence numbers for gap detection. Detecting a gap is half the job: drop the frame (never apply a delta onto a book known to be behind), ask for a snapshot, and raise the event — a supervisor that cannot see the hole keeps trading off a book that has quietly stopped updating.

void handleBookUpdate(const Message& msg) {
  if (_lastSeq > 0 && msg.seq != _lastSeq + 1) {
    LRVX_LOG("Gap detected: " << _lastSeq << " -> " << msg.seq);
    emitSequenceGap(_lastSeq + 1, msg.seq);
    _lastSeq = 0;
    requestSnapshot();  // Request full book snapshot
    return;             // and drop this delta
  }
  _lastSeq = msg.seq;

  // ... create and emit event
}

Staleness

A socket that stays open while the data stops is the failure the disconnect path cannot catch. Stamp each symbol as its data arrives and override pollFeedHealth(), which the connector's owner calls on its own cadence:

void handleBookUpdate(const Message& msg) {
  markFeedActivity(resolveSymbolId(msg.symbol), nowMonoNanos());
  // ... create and emit event
}

void pollFeedHealth(MonoNanos now) override {
  checkStaleFeeds(now, _config.staleDataTimeoutMs);  // 0 disables the check
}

Pooled Book Events

For high-frequency book updates, use pooled events:

class MyExchangeConnector : public IExchangeConnector
{
  pool::Pool<BookUpdateEvent, 64> _bookPool;

  void handleBookUpdate(const Message& msg) {
    auto evOpt = _bookPool.acquire();
    if (!evOpt) {
      LRVX_LOG("Pool exhausted, dropping book update");
      return;
    }

    auto& ev = *evOpt;
    // ... populate ev

    emitBookUpdate(*ev);
    // Handle is moved to bus, returns to pool when all consumers done
  }
};

6. Complete Example

#include "lrvx/connector/abstract_exchange_connector.h"
#include "lrvx/book/events/trade_event.h"
#include "lrvx/book/events/book_update_event.h"
#include "lrvx/util/memory/pool.h"
#include "lrvx/util/base/time.h"
#include "lrvx/log/log.h"

#include <thread>
#include <atomic>

namespace myexchange
{

using namespace lrvx;

class MyConnector : public IExchangeConnector
{
public:
  MyConnector(SymbolId symbol, TradeBus& tradeBus, BookUpdateBus& bookBus)
    : _symbol(symbol)
    , _tradeBus(tradeBus)
    , _bookBus(bookBus)
  {}

  std::string exchangeId() const override { return "myexchange"; }

  void start() override {
    if (_running.exchange(true)) return;
    _thread = std::thread(&MyConnector::run, this);
  }

  void stop() override {
    if (!_running.exchange(false)) return;
    if (_thread.joinable()) _thread.join();
  }

private:
  void run() {
    LRVX_LOG("[myexchange] Connecting...");
    // ... connect to exchange

    while (_running) {
      // ... receive and process messages
    }

    LRVX_LOG("[myexchange] Disconnected");
  }

  void onTradeMessage(/* ... */) {
    TradeEvent ev;
    ev.trade.symbol = _symbol;
    ev.trade.price = Price::fromDouble(/* ... */);
    ev.trade.quantity = Quantity::fromDouble(/* ... */);
    ev.trade.isBuy = /* ... */;
    ev.trade.exchangeTsNs = /* ... */;
    ev.recvNs = nowNsMonotonic();

    _tradeBus.publish(ev);
  }

  void onBookMessage(/* ... */) {
    auto evOpt = _bookPool.acquire();
    if (!evOpt) return;

    auto& ev = *evOpt;
    ev->update.symbol = _symbol;
    ev->update.type = BookUpdateType::SNAPSHOT;
    // ... populate bids/asks

    ev->recvNs = nowNsMonotonic();
    _bookBus.publish(std::move(ev));
  }

  SymbolId _symbol;
  TradeBus& _tradeBus;
  BookUpdateBus& _bookBus;

  pool::Pool<BookUpdateEvent, 64> _bookPool;

  std::atomic<bool> _running{false};
  std::thread _thread;
};

}  // namespace myexchange

See Also