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 },
]));
}
}