Are you an LLM? Read llms.txt for a summary of the docs, or llms-full.txt for the full context.
Skip to content

Price alert stream

Template: price-alert-stream ยท category: monitoring ยท risk: monitor_only

Subscribe to a public price WebSocket feed, filter by symbol and a price threshold, and notify a channel, as an ordinary, durable rflow workflow. No networks, no signer: neither embedded engine boots.

This is the reference use of the trigger.stream WebSocket trigger.

Honest scope

Durable automation, not HFT. A matched message claims a run through rflow's normal path (throttle, rate limit, journal, recovery). It is not a low-latency trading engine, and a stream without a stable message id is not exactly-once: a reconnect can re-deliver.

When to use it

  • watch a market and page an operator when a symbol crosses a level
  • prove a streaming intake path before wiring reads/sends behind it
  • turn an exchange/oracle/vendor feed into normal, replayable rflow runs

How it works

# recipe: partial
rflow_version: 1
name: price-alert-stream
 
# no networks: and no signer: - neither embedded engine boots
config:
  port: 3940
  db_connection: ${DATABASE_URL}
 
notifications:
  channels:
    ops:
      telegram:
        bot_token: ${TG_BOT_TOKEN}
        chat_id: ${TG_CHAT_ID}
 
workflows:
  price-alert:
    trigger:
      stream: 
        websocket: 
          url: ${STREAM_WS_URL}          # e.g. wss://stream.binance.com:9443/ws #
          subscribe: 
            method: SUBSCRIBE
            params: ["ethusdt@ticker"] 
            id: 1
          where: "${{ message.s == 'ETHUSDT' and wei(message.c, 8) > wei('3500', 8) }}"
          idempotency_key: "${{ message.E }}"
          idempotency_ttl: 1h
          heartbeat: 30s
          reconnect: 
            min_backoff: 1s
            max_backoff: 60s
          throttle: 
            threshold: 
              count: 1
              window: 10m
            cooldown: 10m
    steps:
      - id: alert
        notify:
          channel: ops
          message: "ETHUSDT over 3500: last price ${{ trigger.args.c }}"
    on_failure: dead_letter
  • where is evaluated over the message root (the JSON payload after message_json_path, default $); only matching messages fire. wei(message.c, 8) parses the decimal-string price to an exact integer so the comparison is precise.
  • idempotency_key makes the feed's event time the claim identity: a reconnect that re-delivers the same event within idempotency_ttl creates no new run. Without it (or on a feed without stable ids) every matched message fires a fresh run: documented at-least-once behaviour.
  • throttle caps a chatty feed to one alert per 10 minutes while the price stays above the level.

Reliability + redaction

The client connects with a timeout, reconnects with bounded exponential backoff on any drop, pings on the heartbeat interval, and drops any frame over max_message_bytes (default 1 MiB) with a warning. The run payload is { args: <message>, stream: { url_host, received_at } }. A token in the URL is redacted to the host only and never lands in a payload or a log line.

Generate it

rflow new --template price-alert-stream
cd price-alert-stream
cp .env.example .env    # STREAM_WS_URL, TG_BOT_TOKEN, TG_CHAT_ID
docker compose up -d
rflow validate && rflow start

Rehearse offline

rflow test price-alert --fixture fixtures/ticker-message.json

The fixture is the extracted message object, so the where filter and the notification run without touching the internet.