From e285ab3ee0bc4cf3cbf95935264f8bed8495cb6e Mon Sep 17 00:00:00 2001 From: romysaputrasihananda Date: Thu, 11 Jun 2026 01:52:53 +0700 Subject: [PATCH] feat: Telegram notifications for all 6 SSE events + hourly PnL summary - Handle position_opened/closed/modified and order_placed/cancelled/modified - Unified SsePayload struct (Option fields) for both position and order events - Hourly PnL summary task (send_pnl_summary every 3600s, skip first tick) - Add sim binary for testing TG notification flow without MT5 - Export telegram + helpers via lib.rs for bin access - Update README with Telegram setup and event table - Update .env.example with Telegram comment block Co-Authored-By: Claude Sonnet 4.6 --- .env.example | 4 +- README.md | 26 ++++++- src/bin/sim.rs | 51 +++++++++++++ src/lib.rs | 2 + src/live.rs | 201 ++++++++++++++++++++++++++++++++++++------------- 5 files changed, 230 insertions(+), 54 deletions(-) create mode 100644 src/bin/sim.rs create mode 100644 src/lib.rs diff --git a/.env.example b/.env.example index 8485a43..694b6c9 100644 --- a/.env.example +++ b/.env.example @@ -42,6 +42,8 @@ BACKTEST_CANDLES=50000 # LIVE_POLL_SECS=30 # ── Telegram Notifications ─────────────────────────────────────────────────── +# Get token from @BotFather, chat_id from @userinfobot or Telegram API. +# All three are optional — omit to disable notifications. # TELEGRAM_BOT_TOKEN=123456:ABC-xxxx # TELEGRAM_CHAT_ID=-100xxxxxxxxxx -# TELEGRAM_THREAD_ID=123 # optional, for forum topics/threads +# TELEGRAM_THREAD_ID=123 # forum topic/thread ID (optional) diff --git a/README.md b/README.md index bf0bf6f..9c42800 100644 --- a/README.md +++ b/README.md @@ -43,6 +43,22 @@ Polls MT5 every `LIVE_POLL_SECS` seconds. Places a pending limit order when a va State for each symbol is persisted in `.ares_state_{symbol}.json` so the bot survives restarts without orphaning orders. +### Telegram notifications + +Set `TELEGRAM_BOT_TOKEN` and `TELEGRAM_CHAT_ID` to receive real-time trade updates. Each trade lifecycle is tracked as a single edited message: + +| Event | Message | +|---|---| +| Pending order placed | 🟡 PENDING — entry / SL / TP / lot | +| Limit order filled | ⚡ FILLED | +| Position SL/TP modified | 🔄 MODIFIED | +| TP hit | ✅ TP HIT — P&L + balance | +| SL hit | ❌ SL HIT — P&L + balance | +| Order expired | ⏱ EXPIRED | +| Order cancelled externally | 🚫 CANCELLED | + +A 📊 PnL summary (today + all-time) is sent automatically every hour. + ### Multi-pair ```bash @@ -80,6 +96,9 @@ Copy `.env.example` to `.env` and adjust. All fields are optional except `MT5_BA | `TIMEOUT_CANDLES` | `0` | Force-close open trade after N bars (0 = disabled) | | `LIVE` | `false` | Set to `true` to enable live mode | | `LIVE_POLL_SECS` | `30` | Poll interval for live mode | +| `TELEGRAM_BOT_TOKEN` | — | Bot token from @BotFather (optional) | +| `TELEGRAM_CHAT_ID` | — | Target chat/group ID | +| `TELEGRAM_THREAD_ID` | — | Forum thread ID (optional) | ## Backtest results @@ -102,9 +121,12 @@ ares/ ├── src/ │ ├── main.rs # entry point, env parsing, mode dispatch │ ├── backtest.rs # walk-forward simulation engine -│ ├── live.rs # live trading loop +│ ├── live.rs # live trading loop + SSE position listener │ ├── detector.rs # momentum FVG detection -│ └── helpers.rs # shared math utilities +│ ├── telegram.rs # Telegram Bot API client +│ ├── helpers.rs # shared math utilities +│ └── bin/ +│ └── sim.rs # Telegram notification simulator (no MT5 needed) └── crates/ ├── domain/ # shared types (Candle, Symbol, Timeframe, Side…) └── mt5-client/ # async HTTP client for the MT5 bridge diff --git a/src/bin/sim.rs b/src/bin/sim.rs new file mode 100644 index 0000000..cec06f2 --- /dev/null +++ b/src/bin/sim.rs @@ -0,0 +1,51 @@ +/// Simulate: pending → filled → TP/SL result (single trade lifecycle) +use anyhow::Result; +use ares::telegram::TelegramConfig; +use reqwest::Client; +use tokio::time::{sleep, Duration}; + +#[tokio::main] +async fn main() -> Result<()> { + dotenvy::dotenv().ok(); + tracing_subscriber::fmt().with_max_level(tracing::Level::INFO).init(); + + let token = std::env::var("TELEGRAM_BOT_TOKEN").expect("TELEGRAM_BOT_TOKEN"); + let chat_id = std::env::var("TELEGRAM_CHAT_ID").expect("TELEGRAM_CHAT_ID").parse::()?; + let thread_id = std::env::var("TELEGRAM_THREAD_ID").ok().and_then(|s| s.parse::().ok()); + + let tg = TelegramConfig { token, chat_id, thread_id }; + let http = Client::new(); + + // ── 1. Pending order placed ─────────────────────────────────────────────── + tracing::info!("step 1: PENDING"); + let msg_id = tg.send(&http, + "🟡 PENDING\nXAUUSDm Long\nEntry: 3325.50\nSL: 3320.00 TP: 3333.75\nVol: 0.03 lot RR: 1.5", + ).await?; + tracing::info!(msg_id, "sent"); + + sleep(Duration::from_secs(3)).await; + + // ── 2. Limit order triggered / filled ──────────────────────────────────── + tracing::info!("step 2: FILLED"); + tg.edit(&http, msg_id, + "⚡ FILLED\nXAUUSDm Long\nEntry: 3325.50\nSL: 3320.00 TP: 3333.75\nVol: 0.03 lot", + ).await?; + tracing::info!(msg_id, "edited"); + + sleep(Duration::from_secs(3)).await; + + // ── 3a. TP hit ──────────────────────────────────────────────────────────── + tracing::info!("step 3: TP HIT"); + tg.edit(&http, msg_id, + "✅ TP HIT +12.68\nXAUUSDm Long\n3325.50 → 3333.75\nVol: 0.03 lot\nBal: $5012.68", + ).await?; + tracing::info!(msg_id, "edited TP HIT"); + + // uncomment to test SL hit instead: + // tg.edit(&http, msg_id, + // "❌ SL HIT -7.50\nXAUUSDm Long\n3325.50 → 3320.00\nVol: 0.03 lot\nBal: $4992.50", + // ).await?; + + println!("\nDone ✅ check Telegram"); + Ok(()) +} diff --git a/src/lib.rs b/src/lib.rs new file mode 100644 index 0000000..62e2ca6 --- /dev/null +++ b/src/lib.rs @@ -0,0 +1,2 @@ +pub mod helpers; +pub mod telegram; diff --git a/src/live.rs b/src/live.rs index d830903..0245057 100644 --- a/src/live.rs +++ b/src/live.rs @@ -95,21 +95,22 @@ impl PosState { fn clear(symbol: &str) { let _ = std::fs::remove_file(Self::path(symbol)); } } -// ── SSE position event ──────────────────────────────────────────────────────── +// ── SSE event payload (unified for positions + orders) ──────────────────────── #[derive(Debug, Deserialize)] -struct SsePosition { +struct SsePayload { ticket: u64, symbol: String, - #[serde(rename = "type")] - pos_type: u32, - volume: f64, - price_open: f64, - price_current: Option, - sl: f64, - tp: f64, - profit: f64, magic: u64, + #[serde(rename = "type")] + kind: Option, + volume: Option, + price_open: Option, // position open price + price: Option, // order limit price + price_current: Option, + sl: Option, + tp: Option, + profit: Option, } // ── entry point ─────────────────────────────────────────────────────────────── @@ -132,16 +133,39 @@ pub async fn run(mt5: &mt5_client::Mt5Client, cfg: &LiveConfig) -> Result<()> { let tf_mins = timeframe_minutes(cfg.timeframe); let expiry_dur = chrono::Duration::minutes(tf_mins * cfg.fvg_expiry_candles as i64); - // spawn SSE position listener if cfg.telegram.is_some() { let mt5_arc = Arc::new(mt5_client::Mt5Client::new(cfg.mt5_base_url.clone())); let http_arc = Arc::new(reqwest::Client::new()); - let symbol = cfg.symbol.clone(); - let tg = cfg.telegram.clone(); - let base_url = cfg.mt5_base_url.clone(); - tokio::spawn(async move { - sse_task(mt5_arc, http_arc, base_url, symbol, tg).await; - }); + + // SSE position listener + { + let mt5_c = Arc::clone(&mt5_arc); + let http_c = Arc::clone(&http_arc); + let symbol = cfg.symbol.clone(); + let tg = cfg.telegram.clone(); + let base_url = cfg.mt5_base_url.clone(); + tokio::spawn(async move { + sse_task(mt5_c, http_c, base_url, symbol, tg).await; + }); + } + + // hourly PnL summary + { + let mt5_c = Arc::clone(&mt5_arc); + let http_c = Arc::clone(&http_arc); + let symbol = cfg.symbol.clone(); + let tg = cfg.telegram.clone().unwrap(); + tokio::spawn(async move { + let mut ticker = interval(Duration::from_secs(3600)); + ticker.tick().await; // skip first immediate tick + loop { + ticker.tick().await; + if let Err(e) = send_pnl_summary(&mt5_c, &http_c, &symbol, &tg).await { + tracing::warn!(symbol = %symbol, "PnL summary error: {e:#}"); + } + } + }); + } } let http = reqwest::Client::new(); @@ -424,62 +448,67 @@ async fn handle_sse_event( data: &str, tg: &Option, ) { - let pos: SsePosition = match serde_json::from_str(data) { + let payload: SsePayload = match serde_json::from_str(data) { Ok(p) => p, - Err(e) => { tracing::warn!("SSE parse error: {e:#}"); return; } + Err(e) => { tracing::warn!(%symbol, event_type, "SSE parse error: {e:#}"); return; } }; - if pos.magic != MAGIC || pos.symbol != symbol { return; } + if payload.magic != MAGIC || payload.symbol != symbol { return; } match event_type { - "position_opened" => on_position_opened(http, symbol, &pos, tg).await, - "position_closed" => on_position_closed(mt5, http, symbol, &pos, tg).await, - _ => {} + "position_opened" => on_position_opened(http, symbol, &payload, tg).await, + "position_closed" => on_position_closed(mt5, http, symbol, &payload, tg).await, + "position_modified" => on_position_modified(http, symbol, &payload, tg).await, + "order_placed" => tracing::info!(%symbol, ticket = payload.ticket, "order placed (SSE)"), + "order_cancelled" => on_order_cancelled(http, symbol, &payload, tg).await, + "order_modified" => on_order_modified(http, symbol, &payload, tg).await, + other => tracing::debug!(%symbol, "unknown SSE event: {other}"), } } async fn on_position_opened( - http: &reqwest::Client, - symbol: &str, - pos: &SsePosition, - tg: &Option, + http: &reqwest::Client, + symbol: &str, + payload: &SsePayload, + tg: &Option, ) { - tracing::info!(%symbol, ticket = pos.ticket, "position opened (SSE)"); + tracing::info!(%symbol, ticket = payload.ticket, "position opened (SSE)"); - // promote pending state → position state - let state = State::load(symbol); + let state = State::load(symbol); let tg_msg_id = state.as_ref().and_then(|s| s.tg_message_id); + let kind = payload.kind.unwrap_or(0); + let entry = payload.price_open.unwrap_or(0.0); + let sl = payload.sl.unwrap_or(0.0); + let tp = payload.tp.unwrap_or(0.0); + let volume = payload.volume.unwrap_or(0.0); let ps = PosState { - ticket: pos.ticket, + ticket: payload.ticket, tg_message_id: tg_msg_id, - side: if pos.pos_type == 0 { "Long".to_string() } else { "Short".to_string() }, - entry: pos.price_open, - sl: pos.sl, - tp: pos.tp, - volume: pos.volume, + side: if kind == 0 { "Long".to_string() } else { "Short".to_string() }, + entry, sl, tp, volume, }; let _ = ps.save(symbol); State::clear(symbol); if let (Some(tg), Some(msg_id)) = (tg, tg_msg_id) { - let side_str = &ps.side; let text = format!( "⚡ FILLED\n{} {}\nEntry: {:.5}\nSL: {:.5} TP: {:.5}\nVol: {:.2} lot", - symbol, side_str, pos.price_open, pos.sl, pos.tp, pos.volume, + symbol, ps.side, entry, sl, tp, volume, ); let _ = tg.edit(http, msg_id, &text).await; } } async fn on_position_closed( - mt5: &mt5_client::Mt5Client, - http: &reqwest::Client, - symbol: &str, - pos: &SsePosition, - tg: &Option, + mt5: &mt5_client::Mt5Client, + http: &reqwest::Client, + symbol: &str, + payload: &SsePayload, + tg: &Option, ) { - tracing::info!(%symbol, ticket = pos.ticket, profit = pos.profit, "position closed (SSE)"); + let profit = payload.profit.unwrap_or(0.0); + tracing::info!(%symbol, ticket = payload.ticket, profit, "position closed (SSE)"); let ps = PosState::load(symbol); PosState::clear(symbol); @@ -487,24 +516,94 @@ async fn on_position_closed( let Some(tg) = tg else { return }; let Some(msg_id) = ps.as_ref().and_then(|s| s.tg_message_id) else { return }; - let (icon, label) = if pos.profit >= 0.0 { ("✅", "TP HIT") } else { ("❌", "SL HIT") }; - let exit = pos.price_current.unwrap_or(0.0); - let entry = ps.as_ref().map(|s| s.entry).unwrap_or(pos.price_open); + let (icon, label) = if profit >= 0.0 { ("✅", "TP HIT") } else { ("❌", "SL HIT") }; + let exit = payload.price_current.unwrap_or(0.0); + let entry = ps.as_ref().map(|s| s.entry).unwrap_or(payload.price_open.unwrap_or(0.0)); + let volume = payload.volume.unwrap_or(ps.as_ref().map(|s| s.volume).unwrap_or(0.0)); let side_str = ps.as_ref().map(|s| s.side.as_str()).unwrap_or("?"); let acct_bal = mt5.account().await.ok().map(|a| a.balance).unwrap_or(Decimal::ZERO); let text = format!( "{} {} {:+.2}\n{} {}\n{:.5} → {:.5}\nVol: {:.2} lot\nBal: ${}", - icon, label, pos.profit, + icon, label, profit, symbol, side_str, entry, exit, - pos.volume, acct_bal, + volume, acct_bal, ); let _ = tg.edit(http, msg_id, &text).await; +} - // send PnL summary after close - let _ = send_pnl_summary(mt5, http, symbol, tg).await; +async fn on_position_modified( + http: &reqwest::Client, + symbol: &str, + payload: &SsePayload, + tg: &Option, +) { + tracing::info!(%symbol, ticket = payload.ticket, "position modified (SSE)"); + + if let Some(mut ps) = PosState::load(symbol) { + if let Some(sl) = payload.sl { ps.sl = sl; } + if let Some(tp) = payload.tp { ps.tp = tp; } + if let Some(v) = payload.volume { ps.volume = v; } + let _ = ps.save(symbol); + + if let (Some(tg), Some(msg_id)) = (tg, ps.tg_message_id) { + let text = format!( + "🔄 MODIFIED\n{} {}\nEntry: {:.5}\nSL: {:.5} TP: {:.5}\nVol: {:.2} lot", + symbol, ps.side, ps.entry, ps.sl, ps.tp, ps.volume, + ); + let _ = tg.edit(http, msg_id, &text).await; + } + } +} + +async fn on_order_cancelled( + http: &reqwest::Client, + symbol: &str, + payload: &SsePayload, + tg: &Option, +) { + tracing::info!(%symbol, ticket = payload.ticket, "order cancelled (SSE)"); + + let state = State::load(symbol); + if state.as_ref().map(|s| s.ticket) != Some(payload.ticket) { return; } + + if let (Some(tg), Some(msg_id)) = (tg, state.as_ref().and_then(|s| s.tg_message_id)) { + let side_str = state.as_ref().map(|s| s.side.as_str()).unwrap_or("?"); + let entry = state.as_ref().map(|s| s.entry).unwrap_or(0.0); + let text = format!( + "🚫 CANCELLED\n{} {}\nEntry: {:.5}", + symbol, side_str, entry, + ); + let _ = tg.edit(http, msg_id, &text).await; + } + State::clear(symbol); +} + +async fn on_order_modified( + http: &reqwest::Client, + symbol: &str, + payload: &SsePayload, + tg: &Option, +) { + tracing::info!(%symbol, ticket = payload.ticket, "order modified (SSE)"); + + if let Some(mut state) = State::load(symbol) { + if state.ticket != payload.ticket { return; } + if let Some(p) = payload.price { state.entry = p; } + if let Some(sl) = payload.sl { state.sl = sl; } + if let Some(tp) = payload.tp { state.tp = tp; } + let _ = state.save(symbol); + + if let (Some(tg), Some(msg_id)) = (tg, state.tg_message_id) { + let text = format!( + "🔄 ORDER MODIFIED\n{} {}\nEntry: {:.5}\nSL: {:.5} TP: {:.5}", + symbol, state.side, state.entry, state.sl, state.tp, + ); + let _ = tg.edit(http, msg_id, &text).await; + } + } } // ── PnL summary ───────────────────────────────────────────────────────────────