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_letterwhereis evaluated over themessageroot (the JSON payload aftermessage_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_keymakes the feed's event time the claim identity: a reconnect that re-delivers the same event withinidempotency_ttlcreates no new run. Without it (or on a feed without stable ids) every matched message fires a fresh run: documented at-least-once behaviour.throttlecaps 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 startRehearse offline
rflow test price-alert --fixture fixtures/ticker-message.jsonThe fixture is the extracted message object, so the where filter and the
notification run without touching the internet.