Api

Events & Subscriptions

Real-time streaming updates for market data, orders, and swaps

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

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:

FieldTypeRequiredDescription
networkNetworkYESThe network to subscribe to. One stream per network

Response: stream of ClientEvent. Each message sets exactly one arm of its event oneof:

ArmPayloadMeaning
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_update is the right way to track balances

It fires on every change, on-chain and off-chain alike, and carries the whole Balance — so polling GetBalances on a timer buys you nothing but load. Use the poll only to reconcile after a reconnect.

transaction_removed is 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:

FieldTypeRequiredDescription
networkNetworkYESThe network to subscribe to. One stream per network

Response: stream of NodeEvent. Each message sets exactly one arm of its event oneof:

ArmPayloadMeaning
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_update and asset_channel_update are not duplicates

A channel holds several assets. channel_update carries the whole channel and is the right thing to render a channel list from; asset_channel_update names one (channel_id, asset_id) and is what a per-asset view should follow. A single on-chain operation can produce both.

watchtower_disconnected means your channels are undefended until it returns. It is worth alerting on, not just logging — see Watchtowers.

revealed_preimage is the atomic-swap signal on the channel rail, the way Claimed.preimage is 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 own SubscribeHtlcEvents stream 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:

NameTypeRequiredDescription
baseOrderbookCurrencyYESBase currency
quoteOrderbookCurrencyYESQuote currency

Response: Stream of MarketEvent

MarketEvent Types

EventDescriptionFields
is_syncedInitial sync completebool - Always true when synced
orderbook_updateOrderbook changedOrderbookUpdate - Added/removed orders
trade_updateTrade executedTrade - Trade details
daily_stats_update24h stats updatedMarketDailyStats - Volume, price stats
candlestick_updateNew candlestick dataCandlestickUpdate - OHLCV data
venue_status(2026-09-16) The hub stopped or resumed accepting ordersVenueStatus — see below
market_info(2026-09-16) This market's configuration changedMarketInfo — 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.

FieldTypeDescription
sinceTimestampWhen the venue entered this state
normalVenueNormalTrading normally (empty message)
maintenanceVenueMaintenanceIn 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 MarketInfo goes stale silently and nothing else tells you. This event is the only push that does.

OrderbookUpdate

FieldTypeDescription
updated_ordersmap<string, LiquidityPosition>Modified orders
removed_ordersstring[]Removed order IDs

Trade

FieldTypeDescription
taker_order_idstringTaker's order ID
base_amountDecimalStringBase amount traded
quote_amountDecimalStringQuote amount traded
priceDecimalStringExecution price
final_priceDecimalStringPrice after fees
timestampTimestampTrade time
maker_order_sideOrderSideBUY or SELL

MarketDailyStats

FieldTypeDescription
volatilityMarketVolatility24h price and volume data

MarketVolatility:

FieldTypeDescription
first_priceDecimalStringOpening price (24h ago)
last_priceDecimalStringCurrent price
high_priceDecimalStringHighest price (24h)
low_priceDecimalStringLowest price (24h)
base_volumeDecimalString24h volume in base
quote_volumeDecimalString24h volume in quote

CandlestickUpdate

FieldTypeDescription
intervalCandlestickIntervalTime interval
candlestickCandlestickOHLCV 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:

FieldTypeDescription
timestampTimestampCandle timestamp
openDecimalStringOpening price
closeDecimalStringClosing price
highDecimalStringHighest price
lowDecimalStringLowest price
base_volumeDecimalStringVolume in base
quote_volumeDecimalStringVolume 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

EventDescription
is_syncedInitial sync complete
balance_updateYour balance changed
order_updateYour order created/updated/completed/cancelled
order_matchedYour order matched with counterparty
swap_updateSwap progress update
market_trade_updateYour market trade completed
swap_trade_updateYour swap trade completed

BalanceUpdate

FieldTypeDescription
currencyOrderbookCurrencyCurrency identifier
balanceCurrencyBalanceUpdated balance

CurrencyBalance: the shape GetOrderbookBalances returns.

FieldTypeDescription
sendingDecimalStringWhat a new order can send now
receivingDecimalStringWhat a new order can receive now
pending_sendingDecimalStringYour side arriving from a deposit or channel open not confirmed yet
pending_receivingDecimalStringThe hub's side arriving from a deposit or channel open not confirmed yet
unavailable_sendingDecimalStringYour side owned but not spendable now: reserve, payments in flight, withdrawals, redemptions
unavailable_receivingDecimalStringThe hub's side owned but not spendable now
in_use_sendingDecimalStringHeld by your open orders (sending)
in_use_receivingDecimalStringHeld by your open orders (receiving)
ineligible_sendingDecimalStringYour side free on channels too short for the venue to count
ineligible_receivingDecimalStringThe hub's side free on those channels

OrderUpdate

Four variants:

OrderCreated:

  • order_id - New order identifier
  • order - Complete order details
  • released — 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: order is placed at what it matched, and released says what was not placed and why — below_best_fill, self_trade_prevented or immediate_or_cancel. See Price priority

OrderUpdated:

  • order_id - Updated order identifier
  • order - Updated order details

OrderCompleted:

  • order_id - Completed order identifier

OrderCanceled:

  • order_id - Cancelled order identifier
  • reason — (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 a MARKET_STATE_POST_ONLY market, 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_cancel meanwhile); a fill that fails then gives nothing back to it.

CANCEL_REASON_UNSPECIFIED (0) comes only from a hub predating the field; treat it as USER_REQUESTED rather 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

FieldTypeDescription
own_order_idstringYour order ID
swap_routeSwapPath[]Multi-hop swap route
swap_roleSwapRole enumYour role in this match — SWAP_ROLE_LAST_MAKER (1), SWAP_ROLE_INTERMEDIATE_MAKER (2), SWAP_ROLE_TAKER (3)

There is no is_taker boolean. The field is swap_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 showed is_taker; the proto has carried swap_role since the on-chain settlement work.

SWAP_ROLE_UNSPECIFIED (0) is always a protocol error. Note that SWAP_ROLE_TAKER moved from 0 to 3 on 2026-05-03 — re-map any persisted integers.

SwapPath:

FieldTypeDescription
swap_idstringSwap identifier
first_hopSwapHopFirst hop details
next_hopsSwapHop[]Subsequent hops

SwapHop:

FieldTypeDescription
sending_currencyOrderbookCurrencyCurrency sent
receiving_currencyOrderbookCurrencyCurrency received
sending_amountDecimalStringAmount sent
receiving_amountDecimalStringAmount received
receiving_feeDecimalStringFee paid
sending_onchainOnchainSendSettlement?This hop's sending leg settles on-chain: { counterparty_htlc_pubkey, timelock, confirmation_depth }. Unset = channel settlement
receiving_onchainOnchainRecvSettlement?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

FieldTypeDescription
order_idstringOrder identifier
swap_idstringSwap identifier
progressSwapProgressCurrent progress

SwapProgress:

FieldTypeDescription
receiving_amountDecimalStringAmount receiving
paying_amountDecimalStringAmount paying
statusSwapStatusCurrent status
errorstringError message (optional)

SwapStatus Enum:

ValueDescription
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

FieldTypeDescription
baseOrderbookCurrencyBase currency
quoteOrderbookCurrencyQuote currency
tradeClientMarketTradeYour 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 → C sends 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 with role = TRADE_ROLE_MAKER, once with TRADE_ROLE_TAKER — which is what the trade-history RPCs also return. Key a live "my trades" mirror on (swap_id, order_id, role), not on swap_id.

SwapTradeUpdate

FieldTypeDescription
from_currencyOrderbookCurrencySource currency
to_currencyOrderbookCurrencyDestination currency
tradeClientSwapTradeYour swap trade details

Sent to the taker only, once per swap. A fill your resting order made arrives as a MarketTradeUpdate with role = TRADE_ROLE_MAKER and never as a SwapTradeUpdate — 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

FieldTypeDescription
timestampTimestampUpdate time
simple_swap_idstringSwap identifier
update-One of many update types

Update Types

Channel Setup Updates

FundingSendingChannel:

  • txid - Funding transaction ID
  • channel_id - Channel identifier
  • amount - Funding amount
  • is_opening - Whether opening new channel

LeasingReceivingChannel:

  • txid - Rental transaction ID
  • channel_id - Rented channel ID
  • amount - Rental amount

DualFundingChannel:

  • txid - Dual-fund transaction ID
  • channel_id - Channel identifier
  • self_amount - Your contribution
  • counterparty_amount - Peer's contribution
  • is_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 balance
  • needed_receiving - Required receiving balance

BalancesReady:

  • No fields - balances are sufficient

Order Updates

OrderCreated:

  • order_id - Created order identifier

OrderCompleted:

  • order_id - Completed order identifier
  • sent_amount - Amount sent
  • received_amount - Amount received
  • unmatched_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 — so sent_amount can be below the amount you ordered, and received_amount below the quote, without anything having failed.

Read the SwapAmount variant before the figure. A from order's shortfall is sending currency that stayed in your sending channel, and sent_amount + unmatched_amount is what the order set out to sell; a to order's is receiving currency you asked for and did not get, pairing with received_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 IDs
  • self_amount - Your withdrawal
  • counterparty_amount - Peer's withdrawal

WithdrawingReceivingFunds:

  • txids - Withdrawal transaction IDs
  • self_amount - Your withdrawal
  • counterparty_amount - Peer's withdrawal

WithdrawingDualFundedFunds:

  • txids - Withdrawal transaction IDs
  • sending_self_amount - Your sending withdrawal
  • sending_counterparty_amount - Peer's sending withdrawal
  • receiving_self_amount - Your receiving withdrawal
  • receiving_counterparty_amount - Peer's receiving withdrawal

WithdrawingFundsViaLiquidityService: (added 2026-07-23)

  • is_sending_side - true for the sending side's exit, false for the receiving side
  • fee - The exit fee, settled off-chain with the liquidity service
  • fee_payment_currency - Which asset pays it (SENDING / RECEIVING / NATIVE)

Emitted instead of the local WithdrawingSendingFunds / WithdrawingReceivingFunds milestone when that side's exit_rail is EXIT_RAIL_LIQUIDITY_SERVICE — the hub broadcasts the withdrawal and pays the gas, so there is no local transaction and no txids here. 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. Use Htlc.lock_depth to 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:

FieldTypeRequiredDescription
networkNetworkYESThe 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

  1. Handle reconnections - Streams can disconnect, implement retry logic
  2. Filter events - Only process events relevant to your use case
  3. Limit subscriptions - Don't subscribe to too many markets simultaneously
  4. Clean up streams - Call stream.cancel() when done
  5. Buffer updates - Rate-limit UI updates to avoid overwhelming the interface
  6. Error handling - Always listen to error and end events

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!

← Back to API Reference | Next: Error Codes →


Copyright © 2025