Events & Subscriptions
Hydra App pushes real-time updates over gRPC server-streaming RPCs. This page catalogues what each stream emits; the Streaming guide covers how to consume them without losing events.
Every stream here is available over both transports: as a gRPC server-stream, or as a JSON-RPC subscription over a WebSocket. The event shapes are identical; only the field naming differs, per the wire encoding.
Available Subscriptions
- Subscribe Client Events - On-chain sync, blocks, transactions, balances (EventService)
- Subscribe Node Events - Channels, payments, peers, watchtowers (EventService)
- Subscribe Market Events - Public market data and orderbook updates
- Subscribe DEX Events - Your personal trading activity
- Subscribe Simple Swaps - Simple swap progress tracking
- Subscribe HTLC Events - On-chain HTLC lifecycle (EventService; added 2026-07-02)
This page is the event catalogue — what each stream emits. How to consume them (subscribe ordering, reconnect, dedupe, cancellation) is the Streaming guide.
Subscribe Client Events
The on-chain half of the node: chain sync, new blocks, wallet transactions, and balance changes for one network.
Service: EventService · Method: SubscribeClientEvents (server-streaming) · JSON-RPC namespace: event
Request:
| Field | Type | Required | Description |
|---|---|---|---|
network | Network | YES | The network to subscribe to. One stream per network |
Response: stream of ClientEvent. Each message sets exactly one arm of its event oneof:
| Arm | Payload | Meaning |
|---|---|---|
syncing | (empty) | The client has started syncing with the chain |
sync_progress | { current_block, target_block } | How far through the sync it is |
sync_error | { error, retry_in_ms } | Sync hit an error and will retry in retry_in_ms — not a terminal failure, and not a reason to tear down the stream |
synced | (empty) | Caught up to the chain tip. This is the "safe to act" signal |
new_block | { number } | A new block was produced |
transaction_update | { transaction: Transaction } | A wallet transaction was created, updated or confirmed — the full Transaction, not a delta |
transaction_removed | { txid } | A transaction left the mempool without confirming |
balance_update | { asset_id, balance: Balance } | That asset's balance changed. balance carries both the on-chain and off-chain breakdown |
balance_updateis the right way to track balancesIt fires on every change, on-chain and off-chain alike, and carries the whole
Balance— so pollingGetBalanceson a timer buys you nothing but load. Use the poll only to reconcile after a reconnect.
transaction_removedis not a failure event you can ignore. A transaction that leaves the mempool without confirming has not happened — anything you credited on the strength of it must be rolled back. It is how a dropped or replaced transaction surfaces.
Dedupe key: transaction.id for transaction events, (asset_id, network) for balance events.
import { EventServiceClient } from './proto/EventServiceClientPb'
import { SubscribeClientEventsRequest } from './proto/event_pb'
const events = new EventServiceClient('http://localhost:5003')
const request = new SubscribeClientEventsRequest()
request.setNetwork({ protocol: 1, id: '0a03cf40' })
const stream = events.subscribeClientEvents(request, {})
stream.on('data', (evt) => {
if (evt.hasSynced()) {
console.log('chain synced — safe to act')
} else if (evt.hasBalanceUpdate()) {
const u = evt.getBalanceUpdate()
console.log(`balance ${u.getAssetId()}:`, u.getBalance()?.toObject())
} else if (evt.hasTransactionRemoved()) {
console.warn('dropped from mempool:', evt.getTransactionRemoved().getTxid())
}
})
Subscribe Node Events
The off-chain half of the node: channel lifecycle, payment lifecycle, peer and watchtower connectivity for one network.
Service: EventService · Method: SubscribeNodeEvents (server-streaming) · JSON-RPC namespace: event
Request:
| Field | Type | Required | Description |
|---|---|---|---|
network | Network | YES | The network to subscribe to. One stream per network |
Response: stream of NodeEvent. Each message sets exactly one arm of its event oneof:
| Arm | Payload | Meaning |
|---|---|---|
syncing | (empty) | The node has started syncing its off-chain state |
sync_progress | { current_block, target_block } | How far through the sync it is |
sync_error | { error, retry_in_ms } | Sync hit an error and will retry in retry_in_ms |
synced | (empty) | Off-chain state is caught up |
channel_update | { channel: Channel } | A channel changed — the full Channel, every asset |
channel_closed | { channel_id } | A channel was fully closed and removed |
asset_channel_update | { channel_id, asset_id, asset_channel } | One asset inside a channel changed |
asset_channel_closed | { channel_id, asset_id } | One asset channel was closed; the channel itself may still hold others |
payment_update | { payment: Payment } | An off-chain payment was created, advanced or settled — the full Payment |
revealed_preimage | { preimage } | A hashlock preimage became known to this node, hex-encoded |
peer_connected / peer_disconnected | { node_id } | Peer connectivity |
watchtower_connected / watchtower_disconnected | { node_id } | Watchtower session connectivity |
channel_updateandasset_channel_updateare not duplicatesA channel holds several assets.
channel_updatecarries the whole channel and is the right thing to render a channel list from;asset_channel_updatenames one(channel_id, asset_id)and is what a per-asset view should follow. A single on-chain operation can produce both.
watchtower_disconnectedmeans your channels are undefended until it returns. It is worth alerting on, not just logging — see Watchtowers.
revealed_preimageis the atomic-swap signal on the channel rail, the wayClaimed.preimageis on the on-chain rail: once the preimage is known, the other leg can be settled.
Dedupe key: channel.id + status for channel events, payment.id for payment events, node_id for peer/watchtower events.
import { SubscribeNodeEventsRequest } from './proto/event_pb'
const request = new SubscribeNodeEventsRequest()
request.setNetwork({ protocol: 1, id: '0a03cf40' })
const stream = events.subscribeNodeEvents(request, {})
stream.on('data', (evt) => {
if (evt.hasChannelUpdate()) {
const ch = evt.getChannelUpdate().getChannel()
console.log(`channel ${ch.getId()} → ${ch.getStatus()}`)
} else if (evt.hasPaymentUpdate()) {
const p = evt.getPaymentUpdate().getPayment()
console.log(`payment ${p.getId()} → ${Object.keys(p.getStatus().toObject())[0]}`)
} else if (evt.hasWatchtowerDisconnected()) {
console.error('watchtower lost:', evt.getWatchtowerDisconnected().getNodeId())
}
})
On-chain HTLC transitions used to arrive here as
NodeEvent.HtlcUpdate. They moved to their ownSubscribeHtlcEventsstream on 2026-07-02 and are no longer emitted on this stream.
Subscribe Market Events
Stream real-time public market updates for a specific trading pair.
Service: OrderbookServiceMethod: SubscribeMarketEvents
Parameters:
| Name | Type | Required | Description |
|---|---|---|---|
base | OrderbookCurrency | YES | Base currency |
quote | OrderbookCurrency | YES | Quote currency |
Response: Stream of MarketEvent
MarketEvent Types
| Event | Description | Fields |
|---|---|---|
is_synced | Initial sync complete | bool - Always true when synced |
orderbook_update | Orderbook changed | OrderbookUpdate - Added/removed orders |
trade_update | Trade executed | Trade - Trade details |
daily_stats_update | 24h stats updated | MarketDailyStats - Volume, price stats |
candlestick_update | New candlestick data | CandlestickUpdate - OHLCV data |
venue_status | (2026-09-16) The hub stopped or resumed accepting orders | VenueStatus — see below |
market_info | (2026-09-16) This market's configuration changed | MarketInfo — full state, latest wins |
venue_status — the hub is in maintenance
(2026-09-16) Venue-wide, so it arrives on every market you are subscribed to at once. Exactly one of normal / maintenance is set; the reason lives inside the maintenance arm because it has no meaning while the venue trades normally.
| Field | Type | Description |
|---|---|---|
since | Timestamp | When the venue entered this state |
normal | VenueNormal | Trading normally (empty message) |
maintenance | VenueMaintenance | In maintenance |
VenueMaintenance carries reason (the operator's own words when they declared it, the hub's description of what it detected otherwise) and operator_declared. An operator's window closes placement. One the hub enters by itself (operator_declared false, a lost upstream stream say) usually clears on its own and keeps taking orders: the hub only stops cancelling the orders of clients it has lost sight of.
Your resting orders survive a maintenance window, and cancels keep working through it. What an operator's window stops is placement. Without this event the only way to learn the venue closed is to be refused, so a quoting bot that ignores it spends the window retrying placements that cannot succeed.
market_info — this market's configuration changed
(2026-09-16) Fee rates, placement floors, discount ceilings and lifecycle state. Full state, latest wins — replace your cached copy rather than merging into it.
Every one of those hot-reloads on the hub. The operator edits them with no restart and no disconnect, so a cached
MarketInfogoes stale silently and nothing else tells you. This event is the only push that does.
OrderbookUpdate
| Field | Type | Description |
|---|---|---|
updated_orders | map<string, LiquidityPosition> | Modified orders |
removed_orders | string[] | Removed order IDs |
Trade
| Field | Type | Description |
|---|---|---|
taker_order_id | string | Taker's order ID |
base_amount | DecimalString | Base amount traded |
quote_amount | DecimalString | Quote amount traded |
price | DecimalString | Execution price |
final_price | DecimalString | Price after fees |
timestamp | Timestamp | Trade time |
maker_order_side | OrderSide | BUY or SELL |
MarketDailyStats
| Field | Type | Description |
|---|---|---|
volatility | MarketVolatility | 24h price and volume data |
MarketVolatility:
| Field | Type | Description |
|---|---|---|
first_price | DecimalString | Opening price (24h ago) |
last_price | DecimalString | Current price |
high_price | DecimalString | Highest price (24h) |
low_price | DecimalString | Lowest price (24h) |
base_volume | DecimalString | 24h volume in base |
quote_volume | DecimalString | 24h volume in quote |
CandlestickUpdate
| Field | Type | Description |
|---|---|---|
interval | CandlestickInterval | Time interval |
candlestick | Candlestick | OHLCV data |
CandlestickInterval Enum:
ONE_MINUTE, THREE_MINUTES, FIVE_MINUTES, FIFTEEN_MINUTES, THIRTY_MINUTES, ONE_HOUR, TWO_HOURS, FOUR_HOURS, SIX_HOURS, EIGHT_HOURS, TWELVE_HOURS, ONE_DAY, THREE_DAYS, ONE_WEEK, ONE_MONTH
Candlestick:
| Field | Type | Description |
|---|---|---|
timestamp | Timestamp | Candle timestamp |
open | DecimalString | Opening price |
close | DecimalString | Closing price |
high | DecimalString | Highest price |
low | DecimalString | Lowest price |
base_volume | DecimalString | Volume in base |
quote_volume | DecimalString | Volume in quote |
Example Usage
import { OrderbookServiceClient } from './proto/OrderbookServiceClientPb'
import { SubscribeMarketEventsRequest } from './proto/orderbook_pb'
const client = new OrderbookServiceClient('http://localhost:5003')
const request = new SubscribeMarketEventsRequest()
request.setBase({
network: { protocol: 1, id: '0a03cf40' },
assetId: '0x0000000000000000000000000000000000000000000000000000000000000000'
})
request.setQuote({
network: { protocol: 2, id: '11155111' },
assetId: '0x0000000000000000000000000000000000000000'
})
const stream = client.subscribeMarketEvents(request, {})
stream.on('data', (event) => {
if (event.hasIsSynced()) {
console.log('✓ Market data synced')
}
if (event.hasOrderbookUpdate()) {
const update = event.getOrderbookUpdate()
const updatedOrders = update.getUpdatedOrdersMap()
const removedOrders = update.getRemovedOrdersList()
console.log('Orderbook updated:')
console.log(' Modified:', updatedOrders.size)
console.log(' Removed:', removedOrders.length)
}
if (event.hasTradeUpdate()) {
const trade = event.getTradeUpdate()
console.log('Trade executed:')
console.log(' Amount:', trade.getBaseAmount())
console.log(' Price:', trade.getPrice())
console.log(' Side:', trade.getMakerOrderSide() === 0 ? 'BUY' : 'SELL')
}
if (event.hasDailyStatsUpdate()) {
const stats = event.getDailyStatsUpdate()
const vol = stats.getVolatility()
console.log('24h Stats:')
console.log(' High:', vol.getHighPrice())
console.log(' Low:', vol.getLowPrice())
console.log(' Volume:', vol.getBaseVolume())
}
if (event.hasCandlestickUpdate()) {
const candleUpdate = event.getCandlestickUpdate()
const candle = candleUpdate.getCandlestick()
console.log('New candle:', candleUpdate.getInterval())
console.log(' O:', candle.getOpen())
console.log(' H:', candle.getHigh())
console.log(' L:', candle.getLow())
console.log(' C:', candle.getClose())
}
})
stream.on('error', (err) => {
console.error('Stream error:', err)
})
stream.on('end', () => {
console.log('Stream ended')
})
Go Example
package main
import (
"context"
"fmt"
"io"
"log"
pb "path/to/proto"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
conn, err := grpc.Dial("localhost:5003", grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
log.Fatalf("Failed to connect: %v", err)
}
defer conn.Close()
client := pb.NewOrderbookServiceClient(conn)
request := &pb.SubscribeMarketEventsRequest{
Base: &pb.OrderbookCurrency{
Network: &pb.Network{
Protocol: pb.Protocol_PROTOCOL_BITCOIN,
Id: "0a03cf40",
},
AssetId: "0x0000000000000000000000000000000000000000000000000000000000000000",
},
Quote: &pb.OrderbookCurrency{
Network: &pb.Network{
Protocol: pb.Protocol_PROTOCOL_EVM,
Id: "11155111",
},
AssetId: "0x0000000000000000000000000000000000000000",
},
}
stream, err := client.SubscribeMarketEvents(context.Background(), request)
if err != nil {
log.Fatalf("Failed to subscribe: %v", err)
}
for {
event, err := stream.Recv()
if err == io.EOF {
fmt.Println("Stream ended")
break
}
if err != nil {
log.Fatalf("Stream error: %v", err)
}
if event.GetIsSynced() {
fmt.Println("✓ Market data synced")
}
if update := event.GetOrderbookUpdate(); update != nil {
updatedOrders := update.GetUpdatedOrders()
removedOrders := update.GetRemovedOrders()
fmt.Println("Orderbook updated:")
fmt.Printf(" Modified: %d\n", len(updatedOrders))
fmt.Printf(" Removed: %d\n", len(removedOrders))
}
if trade := event.GetTradeUpdate(); trade != nil {
fmt.Println("Trade executed:")
fmt.Printf(" Amount: %s\n", trade.GetBaseAmount())
fmt.Printf(" Price: %s\n", trade.GetPrice())
side := "SELL"
if trade.GetMakerOrderSide() == pb.OrderSide_BUY {
side = "BUY"
}
fmt.Printf(" Side: %s\n", side)
}
if stats := event.GetDailyStatsUpdate(); stats != nil {
vol := stats.GetVolatility()
fmt.Println("24h Stats:")
fmt.Printf(" High: %s\n", vol.GetHighPrice())
fmt.Printf(" Low: %s\n", vol.GetLowPrice())
fmt.Printf(" Volume: %s\n", vol.GetBaseVolume())
}
if candleUpdate := event.GetCandlestickUpdate(); candleUpdate != nil {
candle := candleUpdate.GetCandlestick()
fmt.Printf("New candle: %v\n", candleUpdate.GetInterval())
fmt.Printf(" O: %s\n", candle.GetOpen())
fmt.Printf(" H: %s\n", candle.GetHigh())
fmt.Printf(" L: %s\n", candle.GetLow())
fmt.Printf(" C: %s\n", candle.GetClose())
}
}
}
Rust Example
use futures::stream::StreamExt;
use proto::orderbook_service_client::OrderbookServiceClient;
use proto::{Network, OrderbookCurrency, Protocol, SubscribeMarketEventsRequest};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let mut client = OrderbookServiceClient::connect("http://localhost:5003").await?;
let request = SubscribeMarketEventsRequest {
base: Some(OrderbookCurrency {
network: Some(Network {
protocol: Protocol::Bitcoin as i32,
id: "0a03cf40".to_string(),
}),
asset_id: "0x0000000000000000000000000000000000000000000000000000000000000000".to_string(),
}),
quote: Some(OrderbookCurrency {
network: Some(Network {
protocol: Protocol::Evm as i32,
id: "11155111".to_string(),
}),
asset_id: "0x0000000000000000000000000000000000000000".to_string(),
}),
};
let mut stream = client.subscribe_market_events(request).await?.into_inner();
while let Some(event) = stream.next().await {
match event {
Ok(event) => {
if event.is_synced {
println!("✓ Market data synced");
}
if let Some(update) = &event.orderbook_update {
println!("Orderbook updated:");
println!(" Modified: {}", update.updated_orders.len());
println!(" Removed: {}", update.removed_orders.len());
}
if let Some(trade) = &event.trade_update {
println!("Trade executed:");
println!(" Amount: {}", trade.base_amount);
println!(" Price: {}", trade.price);
let side = if trade.maker_order_side == 0 { "BUY" } else { "SELL" };
println!(" Side: {}", side);
}
if let Some(stats) = &event.daily_stats_update {
if let Some(vol) = &stats.volatility {
println!("24h Stats:");
println!(" High: {}", vol.high_price);
println!(" Low: {}", vol.low_price);
println!(" Volume: {}", vol.base_volume);
}
}
if let Some(candle_update) = &event.candlestick_update {
if let Some(candle) = &candle_update.candlestick {
println!("New candle: {:?}", candle_update.interval);
println!(" O: {}", candle.open);
println!(" H: {}", candle.high);
println!(" L: {}", candle.low);
println!(" C: {}", candle.close);
}
}
}
Err(e) => {
eprintln!("Stream error: {}", e);
break;
}
}
}
println!("Stream ended");
Ok(())
}
Subscribe DEX Events
Stream your personal trading activity including orders, balances, and swaps.
Service: OrderbookServiceMethod: SubscribeDexEvents
Parameters: None
Response: Stream of DexEvent
DexEvent Types
| Event | Description |
|---|---|
is_synced | Initial sync complete |
balance_update | Your balance changed |
order_update | Your order created/updated/completed/cancelled |
order_matched | Your order matched with counterparty |
swap_update | Swap progress update |
market_trade_update | Your market trade completed |
swap_trade_update | Your swap trade completed |
BalanceUpdate
| Field | Type | Description |
|---|---|---|
currency | OrderbookCurrency | Currency identifier |
balance | CurrencyBalance | Updated balance |
CurrencyBalance: the shape GetOrderbookBalances returns.
| Field | Type | Description |
|---|---|---|
sending | DecimalString | What a new order can send now |
receiving | DecimalString | What a new order can receive now |
pending_sending | DecimalString | Your side arriving from a deposit or channel open not confirmed yet |
pending_receiving | DecimalString | The hub's side arriving from a deposit or channel open not confirmed yet |
unavailable_sending | DecimalString | Your side owned but not spendable now: reserve, payments in flight, withdrawals, redemptions |
unavailable_receiving | DecimalString | The hub's side owned but not spendable now |
in_use_sending | DecimalString | Held by your open orders (sending) |
in_use_receiving | DecimalString | Held by your open orders (receiving) |
ineligible_sending | DecimalString | Your side free on channels too short for the venue to count |
ineligible_receiving | DecimalString | The hub's side free on those channels |
OrderUpdate
Four variants:
OrderCreated:
order_id- New order identifierorder- Complete order detailsreleased—OrderRelease, set when a limit order matched part of what it asked for at once and the rest could not rest, or its time in force lets it only take:orderis placed at what it matched, andreleasedsays what was not placed and why —below_best_fill,self_trade_preventedorimmediate_or_cancel. See Price priority
OrderUpdated:
order_id- Updated order identifierorder- Updated order details
OrderCompleted:
order_id- Completed order identifier
OrderCanceled:
order_id- Cancelled order identifierreason— (2026-09-16)CancelReason, why it left the book
A cancel you did not ask for is the hub telling you something about your own settlement state.
CANCEL_REASON_USER_REQUESTED(1) covers your own cancels and your shutdown.HUB_DISCONNECTED(2) /AUTH_DISCONNECTED(3) /NETWORK_REVOKED(4) mean your settlement node or auth session lapsed — reconnect, then re-place.CAPACITY_LOST(5) /CHANNEL_CLOSED(6) mean the backing channel no longer covers the order.FEE_REBASELINE(7) means a hub fee reload re-baselined the reservation past your capacity — re-place at the new rates.HOUSEKEEPING(8) is book hygiene (dust residue, a rolled-back match, boot reconciliation) and means nothing is wrong on your side.Three more say the book cancelled the order instead of letting it take or rest, and each wants a re-place at another price or size:
POST_ONLY_WOULD_CROSS(9) — a post-only order, or any limit order on aMARKET_STATE_POST_ONLYmarket, was put back on the book across an order on offer and would have taken it;SELF_TRADE_PREVENTED(10) — it would have traded against an order of yours, and the self-trade rule of the newer of the two cancelled it;BELOW_BEST_FILL(11) — what a limit order the hub put back on the book had left was below the smallest fill of an order on offer it would rest across (see Price priority). A new order is never cancelled with the last two: it is refused when it filled nothing, and placed at what it filled otherwise.
FILL_FAILED(12): a fill of an order that only takes (TIME_IN_FORCE_IOC,TIME_IN_FORCE_FOK) failed, and what it took there is not offered again; the fills that settled stand. Nothing to re-place unless you still want the rest.An order cancelled for any reason while a fill of it is still settling ends once that fill settles (it shows
pending_cancelmeanwhile); a fill that fails then gives nothing back to it.
CANCEL_REASON_UNSPECIFIED(0) comes only from a hub predating the field; treat it asUSER_REQUESTEDrather than dropping the update, or you lose the order's terminal state. A reason your client does not know yet reads the same way.
MatchedOrder
| Field | Type | Description |
|---|---|---|
own_order_id | string | Your order ID |
swap_route | SwapPath[] | Multi-hop swap route |
swap_role | SwapRole enum | Your role in this match — SWAP_ROLE_LAST_MAKER (1), SWAP_ROLE_INTERMEDIATE_MAKER (2), SWAP_ROLE_TAKER (3) |
There is no
is_takerboolean. The field isswap_role, and it distinguishes three roles, not two: the last maker in a route is the sole preimage generator, while an intermediate maker follows. An earlier revision of this page showedis_taker; the proto has carriedswap_rolesince the on-chain settlement work.
SWAP_ROLE_UNSPECIFIED(0) is always a protocol error. Note thatSWAP_ROLE_TAKERmoved from 0 to 3 on 2026-05-03 — re-map any persisted integers.
SwapPath:
| Field | Type | Description |
|---|---|---|
swap_id | string | Swap identifier |
first_hop | SwapHop | First hop details |
next_hops | SwapHop[] | Subsequent hops |
SwapHop:
| Field | Type | Description |
|---|---|---|
sending_currency | OrderbookCurrency | Currency sent |
receiving_currency | OrderbookCurrency | Currency received |
sending_amount | DecimalString | Amount sent |
receiving_amount | DecimalString | Amount received |
receiving_fee | DecimalString | Fee paid |
sending_onchain | OnchainSendSettlement? | This hop's sending leg settles on-chain: { counterparty_htlc_pubkey, timelock, confirmation_depth }. Unset = channel settlement |
receiving_onchain | OnchainRecvSettlement? | This hop's receiving leg settles on-chain: { confirmation_depth }. Unset = channel settlement |
A hop's two legs are decided independently, so one can be on-chain while the other rides a channel. See channel vs on-chain settlement for what each side has to do.
SwapUpdate
| Field | Type | Description |
|---|---|---|
order_id | string | Order identifier |
swap_id | string | Swap identifier |
progress | SwapProgress | Current progress |
SwapProgress:
| Field | Type | Description |
|---|---|---|
receiving_amount | DecimalString | Amount receiving |
paying_amount | DecimalString | Amount paying |
status | SwapStatus | Current status |
error | string | Error message (optional) |
SwapStatus Enum:
| Value | Description |
|---|---|
ORDER_MATCHED (0) | Order matched |
RECEIVING_INVOICE_CREATED (1) | Receiving invoice created |
PAYING_INVOICE_RECEIVED (2) | Payment invoice received |
PAYMENT_SENT (3) | Payment sent |
PAYMENT_RECEIVED (4) | Payment received |
PAYMENT_CLAIMED (5) | You claimed payment |
PAYMENT_CLAIMED_BY_COUNTERPARTY (6) | Counterparty claimed |
SWAP_COMPLETED (7) | Swap successful |
SWAP_FAILED (8) | Swap failed |
MarketTradeUpdate
| Field | Type | Description |
|---|---|---|
base | OrderbookCurrency | Base currency |
quote | OrderbookCurrency | Quote currency |
trade | ClientMarketTrade | Your trade details, including role |
One update per fill per role, and a routed swap fills on every market it crosses. A swap that routes
A → B → Csends you one update for each hop it touched, each naming that hop's market. If you were both the resting side and the crossing side of one fill, you get that fill twice — once withrole = TRADE_ROLE_MAKER, once withTRADE_ROLE_TAKER— which is what the trade-history RPCs also return. Key a live "my trades" mirror on(swap_id, order_id, role), not onswap_id.
SwapTradeUpdate
| Field | Type | Description |
|---|---|---|
from_currency | OrderbookCurrency | Source currency |
to_currency | OrderbookCurrency | Destination currency |
trade | ClientSwapTrade | Your swap trade details |
Sent to the taker only, once per swap. A fill your resting order made arrives as a
MarketTradeUpdatewithrole = TRADE_ROLE_MAKERand never as aSwapTradeUpdate— see Get All Swap Trades.
Example Usage
const request = new SubscribeDexEventsRequest()
const stream = client.subscribeDexEvents(request, {})
stream.on('data', (event) => {
const timestamp = event.getTimestamp()
if (event.hasBalanceUpdate()) {
const balance = event.getBalanceUpdate()
console.log('Balance updated:', {
currency: balance.getCurrency()?.getAssetId(),
available: balance.getBalance()?.getSending()
})
}
if (event.hasOrderUpdate()) {
const update = event.getOrderUpdate()
if (update.hasOrderCreated()) {
const created = update.getOrderCreated()
console.log('Order created:', created.getOrderId())
}
if (update.hasOrderCompleted()) {
const completed = update.getOrderCompleted()
console.log('Order completed:', completed.getOrderId())
}
if (update.hasOrderCanceled()) {
const canceled = update.getOrderCanceled()
console.log('Order cancelled:', canceled.getOrderId())
}
}
if (event.hasOrderMatched()) {
const matched = event.getOrderMatched()
console.log('Order matched:', {
orderId: matched.getOwnOrderId(),
swapRole: matched.getSwapRole(),
routes: matched.getSwapRouteList().length
})
}
if (event.hasSwapUpdate()) {
const swap = event.getSwapUpdate()
const progress = swap.getProgress()
console.log('Swap update:', {
swapId: swap.getSwapId(),
status: progress.getStatus(),
receiving: progress.getReceivingAmount(),
paying: progress.getPayingAmount()
})
if (progress.hasError()) {
console.error('Swap error:', progress.getError())
}
}
if (event.hasMarketTradeUpdate()) {
const trade = event.getMarketTradeUpdate()
console.log('Trade executed:', trade.getTrade())
}
})
Go Example
package main
import (
"context"
"fmt"
"io"
"log"
pb "path/to/proto"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
conn, err := grpc.Dial("localhost:5003", grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
log.Fatalf("Failed to connect: %v", err)
}
defer conn.Close()
client := pb.NewOrderbookServiceClient(conn)
request := &pb.SubscribeDexEventsRequest{}
stream, err := client.SubscribeDexEvents(context.Background(), request)
if err != nil {
log.Fatalf("Failed to subscribe: %v", err)
}
for {
event, err := stream.Recv()
if err == io.EOF {
fmt.Println("Stream ended")
break
}
if err != nil {
log.Fatalf("Stream error: %v", err)
}
timestamp := event.GetTimestamp()
if balance := event.GetBalanceUpdate(); balance != nil {
fmt.Printf("Balance updated:\n")
fmt.Printf(" Currency: %s\n", balance.GetCurrency().GetAssetId())
fmt.Printf(" Available: %s\n", balance.GetBalance().GetSending())
}
if update := event.GetOrderUpdate(); update != nil {
if created := update.GetOrderCreated(); created != nil {
fmt.Printf("Order created: %s\n", created.GetOrderId())
}
if completed := update.GetOrderCompleted(); completed != nil {
fmt.Printf("Order completed: %s\n", completed.GetOrderId())
}
if canceled := update.GetOrderCanceled(); canceled != nil {
fmt.Printf("Order cancelled: %s\n", canceled.GetOrderId())
}
}
if matched := event.GetOrderMatched(); matched != nil {
fmt.Printf("Order matched:\n")
fmt.Printf(" Order ID: %s\n", matched.GetOwnOrderId())
fmt.Printf(" Is Taker: %v\n", matched.GetIsTaker())
fmt.Printf(" Routes: %d\n", len(matched.GetSwapRoute()))
}
if swap := event.GetSwapUpdate(); swap != nil {
progress := swap.GetProgress()
fmt.Printf("Swap update:\n")
fmt.Printf(" Swap ID: %s\n", swap.GetSwapId())
fmt.Printf(" Status: %v\n", progress.GetStatus())
fmt.Printf(" Receiving: %s\n", progress.GetReceivingAmount())
fmt.Printf(" Paying: %s\n", progress.GetPayingAmount())
if progress.GetError() != "" {
fmt.Printf(" Error: %s\n", progress.GetError())
}
}
if trade := event.GetMarketTradeUpdate(); trade != nil {
fmt.Printf("Trade executed: %v\n", trade.GetTrade())
}
_ = timestamp // Use timestamp if needed
}
}
Rust Example
use futures::stream::StreamExt;
use proto::orderbook_service_client::OrderbookServiceClient;
use proto::SubscribeDexEventsRequest;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let mut client = OrderbookServiceClient::connect("http://localhost:5003").await?;
let request = SubscribeDexEventsRequest {};
let mut stream = client.subscribe_dex_events(request).await?.into_inner();
while let Some(event) = stream.next().await {
match event {
Ok(event) => {
let _timestamp = &event.timestamp;
if let Some(balance) = &event.balance_update {
println!("Balance updated:");
if let Some(currency) = &balance.currency {
println!(" Currency: {}", currency.asset_id);
}
if let Some(bal) = &balance.balance {
println!(" Available: {}", bal.sending);
}
}
if let Some(update) = &event.order_update {
if let Some(created) = &update.order_created {
println!("Order created: {}", created.order_id);
}
if let Some(completed) = &update.order_completed {
println!("Order completed: {}", completed.order_id);
}
if let Some(canceled) = &update.order_canceled {
println!("Order cancelled: {}", canceled.order_id);
}
}
if let Some(matched) = &event.order_matched {
println!("Order matched:");
println!(" Order ID: {}", matched.own_order_id);
println!(" Swap role: {:?}", matched.swap_role());
println!(" Routes: {}", matched.swap_route.len());
}
if let Some(swap) = &event.swap_update {
if let Some(progress) = &swap.progress {
println!("Swap update:");
println!(" Swap ID: {}", swap.swap_id);
println!(" Status: {:?}", progress.status);
println!(" Receiving: {}", progress.receiving_amount);
println!(" Paying: {}", progress.paying_amount);
if !progress.error.is_empty() {
println!(" Error: {}", progress.error);
}
}
}
if let Some(trade) = &event.market_trade_update {
println!("Trade executed: {:?}", trade.trade);
}
}
Err(e) => {
eprintln!("Stream error: {}", e);
break;
}
}
}
println!("Stream ended");
Ok(())
}
Subscribe Simple Swaps
Monitor progress of Simple Swap operations with automatic channel setup.
Service: SwapServiceMethod: SubscribeSimpleSwaps
Parameters: None
Response: Stream of SimpleSwapUpdate
SimpleSwapUpdate Structure
| Field | Type | Description |
|---|---|---|
timestamp | Timestamp | Update time |
simple_swap_id | string | Swap identifier |
update | - | One of many update types |
Update Types
Channel Setup Updates
FundingSendingChannel:
txid- Funding transaction IDchannel_id- Channel identifieramount- Funding amountis_opening- Whether opening new channel
LeasingReceivingChannel:
txid- Rental transaction IDchannel_id- Rented channel IDamount- Rental amount
DualFundingChannel:
txid- Dual-fund transaction IDchannel_id- Channel identifierself_amount- Your contributioncounterparty_amount- Peer's contributionis_opening- Whether opening new channel
Channel Ready Updates
SendingChannelReady:
channel_id- Ready channel ID
ReceivingChannelReady:
channel_id- Ready channel ID
DualFundChannelReady:
channel_id- Ready channel ID
Balance Updates
WaitingForBalances:
needed_sending- Required sending balanceneeded_receiving- Required receiving balance
BalancesReady:
- No fields - balances are sufficient
Order Updates
OrderCreated:
order_id- Created order identifier
OrderCompleted:
order_id- Completed order identifiersent_amount- Amount sentreceived_amount- Amount receivedunmatched_amount— (2026-09-17)SwapAmount, the part of the order that never filled
A completed swap order is not necessarily a filled one. A swap order fills one maker at a time and the matching engine refuses any fill below the market's per-fill minimum,
max(min_base_amount, min_quote_amount / price). Once the takeable depth at the quoted prices is gone, whatever tail is left under that floor cannot be filled at all and the order completes having moved less than it asked to — sosent_amountcan be below the amount you ordered, andreceived_amountbelow the quote, without anything having failed.Read the
SwapAmountvariant before the figure. Afromorder's shortfall is sending currency that stayed in your sending channel, andsent_amount+unmatched_amountis what the order set out to sell; atoorder's is receiving currency you asked for and did not get, pairing withreceived_amount. Zero on an order that filled in full.The price-change tolerance does not bound this. It is checked once, against the effective RATE, before the order is placed, and it aborts a swap whose quote went stale; it says nothing about how much of the order fills.
Withdrawal Updates
WithdrawingSendingFunds:
txids- Withdrawal transaction IDsself_amount- Your withdrawalcounterparty_amount- Peer's withdrawal
WithdrawingReceivingFunds:
txids- Withdrawal transaction IDsself_amount- Your withdrawalcounterparty_amount- Peer's withdrawal
WithdrawingDualFundedFunds:
txids- Withdrawal transaction IDssending_self_amount- Your sending withdrawalsending_counterparty_amount- Peer's sending withdrawalreceiving_self_amount- Your receiving withdrawalreceiving_counterparty_amount- Peer's receiving withdrawal
WithdrawingFundsViaLiquidityService: (added 2026-07-23)
is_sending_side-truefor the sending side's exit,falsefor the receiving sidefee- The exit fee, settled off-chain with the liquidity servicefee_payment_currency- Which asset pays it (SENDING/RECEIVING/NATIVE)
Emitted instead of the local
WithdrawingSendingFunds/WithdrawingReceivingFundsmilestone when that side'sexit_railisEXIT_RAIL_LIQUIDITY_SERVICE— the hub broadcasts the withdrawal and pays the gas, so there is no local transaction and notxidshere. This is what makes a channel exit possible with no native balance.
SendingFundsWithdrawn:
- No fields - withdrawal complete
ReceivingFundsWithdrawn:
- No fields - withdrawal complete
DualFundedFundsWithdrawn:
- No fields - withdrawal complete
Completion Updates
SimpleSwapCompleted:
- No fields - swap successful
SimpleSwapError:
error- Error message
On-chain HTLC Updates (added 2026-06-24)
Emitted only when a swap leg settles on-chain (opted via SwapRequest.settlement) instead of through a channel. Channel-only swaps never emit these. See the channel-vs-on-chain settlement model.
LockingOnchainHtlc:
txid- the sending leg's on-chain HTLC lock transaction, just broadcast
OnchainHtlcLocked:
- No fields - the sending leg's lock is reported and being watched
CounterpartyLockConfirming:
txid,current,required- the counterparty's inbound on-chain lock is gaining confirmations
CounterpartyLockObserved:
- No fields - the counterparty's inbound lock is observed and confirmed
ClaimingOnchainHtlc:
txid- the receiving leg's on-chain claim transaction, just broadcast
OnchainHtlcClaimed:
- No fields - the receiving leg's claim is observed (the swap secret is now public)
RefundingOnchainHtlc:
txid- the sending leg's on-chain HTLC is being refunded after timeout (the swap did not complete)
Example Usage
import { SwapServiceClient } from './proto/SwapServiceClientPb'
import { SubscribeSimpleSwapsRequest } from './proto/swap_pb'
const client = new SwapServiceClient('http://localhost:5003')
const request = new SubscribeSimpleSwapsRequest()
const stream = client.subscribeSimpleSwaps(request, {})
stream.on('data', (update) => {
const swapId = update.getSimpleSwapId()
const timestamp = update.getTimestamp()
// Channel setup
if (update.hasFundingSendingChannel()) {
const funding = update.getFundingSendingChannel()
console.log(`[${swapId}] Funding sending channel:`, {
txid: funding.getTxid(),
channelId: funding.getChannelId(),
amount: funding.getAmount(),
isOpening: funding.getIsOpening()
})
}
if (update.hasLeasingReceivingChannel()) {
const leasing = update.getLeasingReceivingChannel()
console.log(`[${swapId}] Leasing receiving channel:`, {
txid: leasing.getTxid(),
channelId: leasing.getChannelId(),
amount: leasing.getAmount()
})
}
if (update.hasDualFundingChannel()) {
const dualFund = update.getDualFundingChannel()
console.log(`[${swapId}] Dual-funding channel:`, {
txid: dualFund.getTxid(),
selfAmount: dualFund.getSelfAmount(),
counterpartyAmount: dualFund.getCounterpartyAmount()
})
}
// Channel ready
if (update.hasSendingChannelReady()) {
console.log(`[${swapId}] ✓ Sending channel ready`)
}
if (update.hasReceivingChannelReady()) {
console.log(`[${swapId}] ✓ Receiving channel ready`)
}
// Balance status
if (update.hasWaitingForBalances()) {
const waiting = update.getWaitingForBalances()
console.log(`[${swapId}] Waiting for balances:`, {
sending: waiting.getNeededSending(),
receiving: waiting.getNeededReceiving()
})
}
if (update.hasBalancesReady()) {
console.log(`[${swapId}] ✓ Balances ready`)
}
// Order execution
if (update.hasOrderCreated()) {
const order = update.getOrderCreated()
console.log(`[${swapId}] Order created:`, order.getOrderId())
}
if (update.hasOrderCompleted()) {
const completed = update.getOrderCompleted()
console.log(`[${swapId}] ✓ Order completed:`, {
sent: completed.getSentAmount(),
received: completed.getReceivedAmount()
})
}
// Withdrawals
if (update.hasWithdrawingReceivingFunds()) {
const withdrawing = update.getWithdrawingReceivingFunds()
console.log(`[${swapId}] Withdrawing funds:`, {
txids: withdrawing.getTxidsList(),
amount: withdrawing.getSelfAmount()
})
}
if (update.hasReceivingFundsWithdrawn()) {
console.log(`[${swapId}] ✓ Funds withdrawn`)
}
// Final status
if (update.hasSimpleSwapCompleted()) {
console.log(`[${swapId}] ✓✓✓ SWAP COMPLETED SUCCESSFULLY ✓✓✓`)
}
if (update.hasSimpleSwapError()) {
const error = update.getSimpleSwapError()
console.error(`[${swapId}] ✗ ERROR:`, error.getError())
}
})
stream.on('error', (err) => {
console.error('Stream error:', err)
})
Go Example
package main
import (
"context"
"fmt"
"io"
"log"
pb "path/to/proto"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
func main() {
conn, err := grpc.Dial("localhost:5003", grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
log.Fatalf("Failed to connect: %v", err)
}
defer conn.Close()
client := pb.NewSwapServiceClient(conn)
request := &pb.SubscribeSimpleSwapsRequest{}
stream, err := client.SubscribeSimpleSwaps(context.Background(), request)
if err != nil {
log.Fatalf("Failed to subscribe: %v", err)
}
for {
update, err := stream.Recv()
if err == io.EOF {
fmt.Println("Stream ended")
break
}
if err != nil {
log.Fatalf("Stream error: %v", err)
}
swapId := update.GetSimpleSwapId()
// Channel setup
if funding := update.GetFundingSendingChannel(); funding != nil {
fmt.Printf("[%s] Funding sending channel:\n", swapId)
fmt.Printf(" Txid: %s\n", funding.GetTxid())
fmt.Printf(" Channel ID: %s\n", funding.GetChannelId())
fmt.Printf(" Amount: %s\n", funding.GetAmount())
fmt.Printf(" Is Opening: %v\n", funding.GetIsOpening())
}
if leasing := update.GetLeasingReceivingChannel(); leasing != nil {
fmt.Printf("[%s] Leasing receiving channel:\n", swapId)
fmt.Printf(" Txid: %s\n", leasing.GetTxid())
fmt.Printf(" Channel ID: %s\n", leasing.GetChannelId())
fmt.Printf(" Amount: %s\n", leasing.GetAmount())
}
if dualFund := update.GetDualFundingChannel(); dualFund != nil {
fmt.Printf("[%s] Dual-funding channel:\n", swapId)
fmt.Printf(" Txid: %s\n", dualFund.GetTxid())
fmt.Printf(" Self Amount: %s\n", dualFund.GetSelfAmount())
fmt.Printf(" Counterparty Amount: %s\n", dualFund.GetCounterpartyAmount())
}
// Channel ready
if update.GetSendingChannelReady() != nil {
fmt.Printf("[%s] ✓ Sending channel ready\n", swapId)
}
if update.GetReceivingChannelReady() != nil {
fmt.Printf("[%s] ✓ Receiving channel ready\n", swapId)
}
// Balance status
if waiting := update.GetWaitingForBalances(); waiting != nil {
fmt.Printf("[%s] Waiting for balances:\n", swapId)
fmt.Printf(" Sending: %s\n", waiting.GetNeededSending())
fmt.Printf(" Receiving: %s\n", waiting.GetNeededReceiving())
}
if update.GetBalancesReady() != nil {
fmt.Printf("[%s] ✓ Balances ready\n", swapId)
}
// Order execution
if order := update.GetOrderCreated(); order != nil {
fmt.Printf("[%s] Order created: %s\n", swapId, order.GetOrderId())
}
if completed := update.GetOrderCompleted(); completed != nil {
fmt.Printf("[%s] ✓ Order completed:\n", swapId)
fmt.Printf(" Sent: %s\n", completed.GetSentAmount())
fmt.Printf(" Received: %s\n", completed.GetReceivedAmount())
}
// Withdrawals
if withdrawing := update.GetWithdrawingReceivingFunds(); withdrawing != nil {
fmt.Printf("[%s] Withdrawing funds:\n", swapId)
fmt.Printf(" Txids: %v\n", withdrawing.GetTxids())
fmt.Printf(" Amount: %s\n", withdrawing.GetSelfAmount())
}
if update.GetReceivingFundsWithdrawn() != nil {
fmt.Printf("[%s] ✓ Funds withdrawn\n", swapId)
}
// Final status
if update.GetSimpleSwapCompleted() != nil {
fmt.Printf("[%s] ✓✓✓ SWAP COMPLETED SUCCESSFULLY ✓✓✓\n", swapId)
}
if swapError := update.GetSimpleSwapError(); swapError != nil {
fmt.Printf("[%s] ✗ ERROR: %s\n", swapId, swapError.GetError())
}
}
}
Rust Example
use futures::stream::StreamExt;
use proto::swap_service_client::SwapServiceClient;
use proto::SubscribeSimpleSwapsRequest;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let mut client = SwapServiceClient::connect("http://localhost:5003").await?;
let request = SubscribeSimpleSwapsRequest {};
let mut stream = client.subscribe_simple_swaps(request).await?.into_inner();
while let Some(update) = stream.next().await {
match update {
Ok(update) => {
let swap_id = &update.simple_swap_id;
// Channel setup
if let Some(funding) = &update.funding_sending_channel {
println!("[{}] Funding sending channel:", swap_id);
println!(" Txid: {}", funding.txid);
println!(" Channel ID: {}", funding.channel_id);
println!(" Amount: {}", funding.amount);
println!(" Is Opening: {}", funding.is_opening);
}
if let Some(leasing) = &update.leasing_receiving_channel {
println!("[{}] Leasing receiving channel:", swap_id);
println!(" Txid: {}", leasing.txid);
println!(" Channel ID: {}", leasing.channel_id);
println!(" Amount: {}", leasing.amount);
}
if let Some(dual_fund) = &update.dual_funding_channel {
println!("[{}] Dual-funding channel:", swap_id);
println!(" Txid: {}", dual_fund.txid);
println!(" Self Amount: {}", dual_fund.self_amount);
println!(" Counterparty Amount: {}", dual_fund.counterparty_amount);
}
// Channel ready
if update.sending_channel_ready.is_some() {
println!("[{}] ✓ Sending channel ready", swap_id);
}
if update.receiving_channel_ready.is_some() {
println!("[{}] ✓ Receiving channel ready", swap_id);
}
// Balance status
if let Some(waiting) = &update.waiting_for_balances {
println!("[{}] Waiting for balances:", swap_id);
println!(" Sending: {}", waiting.needed_sending);
println!(" Receiving: {}", waiting.needed_receiving);
}
if update.balances_ready.is_some() {
println!("[{}] ✓ Balances ready", swap_id);
}
// Order execution
if let Some(order) = &update.order_created {
println!("[{}] Order created: {}", swap_id, order.order_id);
}
if let Some(completed) = &update.order_completed {
println!("[{}] ✓ Order completed:", swap_id);
println!(" Sent: {}", completed.sent_amount);
println!(" Received: {}", completed.received_amount);
}
// Withdrawals
if let Some(withdrawing) = &update.withdrawing_receiving_funds {
println!("[{}] Withdrawing funds:", swap_id);
println!(" Txids: {:?}", withdrawing.txids);
println!(" Amount: {}", withdrawing.self_amount);
}
if update.receiving_funds_withdrawn.is_some() {
println!("[{}] ✓ Funds withdrawn", swap_id);
}
// Final status
if update.simple_swap_completed.is_some() {
println!("[{}] ✓✓✓ SWAP COMPLETED SUCCESSFULLY ✓✓✓", swap_id);
}
if let Some(error) = &update.simple_swap_error {
eprintln!("[{}] ✗ ERROR: {}", swap_id, error.error);
}
}
Err(e) => {
eprintln!("Stream error: {}", e);
break;
}
}
}
println!("Stream ended");
Ok(())
}
Subscribe HTLC Events
Added 2026-07-02 (replaced the 2026-06-04 NodeEvent.HtlcUpdate variant).
Service: EventService · Method: SubscribeHtlcEvents (server-streaming) · JSON-RPC namespace: event
Subscribes to on-chain HTLC lifecycle changes for one network. The stream emits the full Htlc on each transition; the HTLC's status oneof conveys what happened:
Locked— the HTLC is on-chain. UseHtlc.lock_depthto tell confirmed (present) from mempool/unconfirmed (absent).Claimed { claim_txid, claim_depth?, preimage }— the recipient claimed it, revealing the 32-byte preimage (the swap's shared secret — critical for atomic-swap takers waiting on the maker's claim).Refunded { refund_txid, refund_depth? }— refunded to the sender after timelock expiry.
Request:
| Field | Type | Required | Description |
|---|---|---|---|
network | Network | YES | The network to subscribe to |
Dedupe key: (htlc.htlc_id, status-variant) — the same HTLC re-emits as it advances Locked → Claimed/Refunded, and may repeat within a status as depth grows.
import { EventServiceClient } from './proto/EventServiceClientPb'
import { SubscribeHtlcEventsRequest } from './proto/event_pb'
const events = new EventServiceClient('http://localhost:5003')
const request = new SubscribeHtlcEventsRequest()
request.setNetwork({ protocol: 1, id: '0a03cf40' })
const stream = events.subscribeHtlcEvents(request, {})
stream.on('data', (htlc) => {
if (htlc.hasClaimed()) {
// The preimage is now public — settle the other leg.
const preimage = htlc.getClaimed().getPreimage_asU8()
console.log(`HTLC ${htlc.getHtlcId()} claimed; preimage`, Buffer.from(preimage).toString('hex'))
} else if (htlc.hasRefunded()) {
console.log(`HTLC ${htlc.getHtlcId()} refunded`)
} else {
console.log(`HTLC ${htlc.getHtlcId()} locked (confirmed: ${htlc.hasLockDepth()})`)
}
})
To watch an HTLC the node did not create itself (e.g. a counterparty's inbound lock), register it first with
htlc.WatchHtlc— its transitions then arrive on this stream.
Common Patterns
Track specific swap
async function trackSwap(swapId: string): Promise<void> {
return new Promise((resolve, reject) => {
const stream = client.subscribeSimpleSwaps(new SubscribeSimpleSwapsRequest(), {})
stream.on('data', (update) => {
// Filter for our swap
if (update.getSimpleSwapId() !== swapId) return
if (update.hasSimpleSwapCompleted()) {
stream.cancel()
resolve()
}
if (update.hasSimpleSwapError()) {
stream.cancel()
reject(new Error(update.getSimpleSwapError()?.getError()))
}
})
stream.on('error', reject)
})
}
Go Example
func trackSwap(client pb.SwapServiceClient, swapId string) error {
request := &pb.SubscribeSimpleSwapsRequest{}
stream, err := client.SubscribeSimpleSwaps(context.Background(), request)
if err != nil {
return fmt.Errorf("failed to subscribe: %w", err)
}
for {
update, err := stream.Recv()
if err == io.EOF {
return fmt.Errorf("stream ended unexpectedly")
}
if err != nil {
return fmt.Errorf("stream error: %w", err)
}
// Filter for our swap
if update.GetSimpleSwapId() != swapId {
continue
}
if update.GetSimpleSwapCompleted() != nil {
return nil // Success
}
if swapError := update.GetSimpleSwapError(); swapError != nil {
return fmt.Errorf("swap error: %s", swapError.GetError())
}
}
}
Rust Example
async fn track_swap(
client: &mut SwapServiceClient<tonic::transport::Channel>,
swap_id: String,
) -> Result<(), Box<dyn std::error::Error>> {
let request = SubscribeSimpleSwapsRequest {};
let mut stream = client.subscribe_simple_swaps(request).await?.into_inner();
while let Some(update) = stream.next().await {
let update = update?;
// Filter for our swap
if update.simple_swap_id != swap_id {
continue;
}
if update.simple_swap_completed.is_some() {
return Ok(()); // Success
}
if let Some(error) = update.simple_swap_error {
return Err(format!("swap error: {}", error.error).into());
}
}
Err("stream ended unexpectedly".into())
}
Build price chart from market events
const priceHistory: number[] = []
stream.on('data', (event) => {
if (event.hasTradeUpdate()) {
const trade = event.getTradeUpdate()
const price = parseFloat(trade.getPrice())
priceHistory.push(price)
// Update chart
updatePriceChart(priceHistory)
}
})
Go Example
var priceHistory []float64
stream, err := client.SubscribeMarketEvents(context.Background(), request)
if err != nil {
log.Fatalf("Failed to subscribe: %v", err)
}
for {
event, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
log.Fatalf("Stream error: %v", err)
}
if trade := event.GetTradeUpdate(); trade != nil {
price, err := strconv.ParseFloat(trade.GetPrice(), 64)
if err != nil {
log.Printf("Failed to parse price: %v", err)
continue
}
priceHistory = append(priceHistory, price)
// Update chart
updatePriceChart(priceHistory)
}
}
Rust Example
let mut price_history: Vec<f64> = Vec::new();
let mut stream = client.subscribe_market_events(request).await?.into_inner();
while let Some(event) = stream.next().await {
match event {
Ok(event) => {
if let Some(trade) = &event.trade_update {
if let Ok(price) = trade.price.parse::<f64>() {
price_history.push(price);
// Update chart
update_price_chart(&price_history);
}
}
}
Err(e) => {
eprintln!("Stream error: {}", e);
break;
}
}
}
Monitor orderbook depth
const orderbookDepth = new Map<string, LiquidityPosition>()
stream.on('data', (event) => {
if (event.hasOrderbookUpdate()) {
const update = event.getOrderbookUpdate()
// Add/update orders
update.getUpdatedOrdersMap().forEach((position, orderId) => {
orderbookDepth.set(orderId, position)
})
// Remove orders
update.getRemovedOrdersList().forEach(orderId => {
orderbookDepth.delete(orderId)
})
console.log('Current depth:', orderbookDepth.size, 'orders')
}
})
Go Example
orderbookDepth := make(map[string]*pb.LiquidityPosition)
stream, err := client.SubscribeMarketEvents(context.Background(), request)
if err != nil {
log.Fatalf("Failed to subscribe: %v", err)
}
for {
event, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
log.Fatalf("Stream error: %v", err)
}
if update := event.GetOrderbookUpdate(); update != nil {
// Add/update orders
for orderId, position := range update.GetUpdatedOrders() {
orderbookDepth[orderId] = position
}
// Remove orders
for _, orderId := range update.GetRemovedOrders() {
delete(orderbookDepth, orderId)
}
fmt.Printf("Current depth: %d orders\n", len(orderbookDepth))
}
}
Rust Example
use std::collections::HashMap;
let mut orderbook_depth: HashMap<String, LiquidityPosition> = HashMap::new();
let mut stream = client.subscribe_market_events(request).await?.into_inner();
while let Some(event) = stream.next().await {
match event {
Ok(event) => {
if let Some(update) = &event.orderbook_update {
// Add/update orders
for (order_id, position) in &update.updated_orders {
orderbook_depth.insert(order_id.clone(), position.clone());
}
// Remove orders
for order_id in &update.removed_orders {
orderbook_depth.remove(order_id);
}
println!("Current depth: {} orders", orderbook_depth.len());
}
}
Err(e) => {
eprintln!("Stream error: {}", e);
break;
}
}
}
Best Practices
- Handle reconnections - Streams can disconnect, implement retry logic
- Filter events - Only process events relevant to your use case
- Limit subscriptions - Don't subscribe to too many markets simultaneously
- Clean up streams - Call
stream.cancel()when done - Buffer updates - Rate-limit UI updates to avoid overwhelming the interface
- Error handling - Always listen to
errorandendevents
Stream Lifecycle
// 1. Create subscription
const stream = client.subscribeMarketEvents(request, {})
// 2. Listen to events
stream.on('data', (event) => { /* handle */ })
stream.on('error', (err) => { /* handle */ })
stream.on('end', () => { /* reconnect */ })
// 3. Cancel when done
stream.cancel()
Go Example
// 1. Create subscription
stream, err := client.SubscribeMarketEvents(context.Background(), request)
if err != nil {
log.Fatalf("Failed to subscribe: %v", err)
}
// 2. Listen to events
for {
event, err := stream.Recv()
if err == io.EOF {
// Stream ended - reconnect if needed
break
}
if err != nil {
// Handle error
log.Printf("Stream error: %v", err)
break
}
// Handle event
// ...
}
// 3. Stream automatically closes when loop exits
// For manual cancellation, use context.WithCancel
Rust Example
// 1. Create subscription
let mut stream = client.subscribe_market_events(request).await?.into_inner();
// 2. Listen to events
while let Some(event) = stream.next().await {
match event {
Ok(event) => {
// Handle event
// ...
}
Err(e) => {
// Handle error
eprintln!("Stream error: {}", e);
break;
}
}
}
// 3. Stream automatically closes when dropped
// For manual cancellation, drop the stream or use tokio::select!