Skip to content

Contracts

Market-data envelope

MarketDataWorkerBase stamps every output with MarketDataEnvelope.Stamp (openspec pubsub-topics, "Envelope on a market data event"):

Extension Value
exchange the venue config, else the worker's default venue
marketworld the world config: act (default) or hyp.<name>
marketuniverse the marketuniverse config, default ALL
clocktype the clocktype config, else wall for act and historical for hyp.<name>
eventtime, ce-time Outgoing.EventTime (the exchange's time), else the input's eventtime, else its ce-time, else now
walltime now
scenarioid, runid, portfolioworld, strategyworld never set, even when the input carries them

OHLCV candle

Topic act\|hyp.<name>.exchange.<venue>.ohlcv.candle.<interval>.<PAIR>
ce-type com.virtufin.exchange.<venue>.ohlcv.candle
ce-subject currency_pair/<PAIR>
Time eventtime = ce-time = windowEnd

Only closed candles are published. <interval> is the venue's own interval string (1s, 1m, 5m, 1h, 1d, 1w, 1M, ...), the same in the topic and in the payload. A candle at any standard interval is market data.

{
  "symbol": "BTCUSDT",
  "interval": "1m",
  "open": 62000.5,
  "high": 62100,
  "low": 61950.25,
  "close": 62050.12345678,
  "volume": 12.5,
  "windowStart": "2026-10-09T12:00:00.000Z",
  "windowEnd": "2026-10-09T12:00:59.999Z"
}

Encode and decode through CandleJson, never by hand:

var payload = CandleJson.Encode(new Candle("1m", closed));
if (CandleJson.TryDecode(input.Type, payload, out var candle)) { /* candle.Closed, candle.Interval */ }

TryDecode returns false for any other ce-type and throws for a malformed candle, so a malformed candle becomes an error response instead of a silent skip.

Writing a venue worker

public sealed class MyCandlesWorker : MarketDataWorkerBase
{
    public MyCandlesWorker() : base(new Uri("urn:virtufin:worker:mycandles"), "mycandles.response", "myvenue")
    {
        Bind(new MyKlineFolder()
            .FromCloudEvents(DecodeOrSkip)               // Decoded<T>.Skip for acks and open candles
            .ToCloudEvents((ce, candle) =>
            [
                new Outgoing(
                    CandleJson.Encode(candle),
                    CandleTopics.Type(Venue(ce)),
                    CandleTopics.Topic(MarketDataConfig.World(ce), Venue(ce), candle.Interval, candle.Closed.Symbol.Value),
                    CandleTopics.Subject(candle.Closed.Symbol.Value)) { EventTime = candle.Closed.WindowEnd },
            ]));
    }
}