tvsignals-to-tg/app/pipeline.py
Artemii Peretiachenko cdbda8fea3 Replace TradingView polling with a local LTF/FVG scanner that posts Telegram cards and Heryon webhooks.
Per-strategy Heryon accounts, reversal captions with previous-trade path on charts, and Telegram replies chained by ticker.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-30 20:37:16 +02:00

181 lines
6.4 KiB
Python

from __future__ import annotations
import asyncio
import logging
from app.binance import (
chart_kline_limit,
chart_right_pad,
fetch_klines,
get_tick_size,
interval_timedelta,
to_binance_interval,
to_binance_symbol,
)
from app.chart import render_setup_chart
from app.config import Settings
from app.formatter import format_caption
from app.heryon import build_heryon_payload, send_heryon, to_tv_perp_ticker
from app.indicators.common import IndicatorSignal, format_px, realized_pnl_pct
from app.models import SignalPayload
from app.state import ScannerStore
from app.telegram import TelegramError, send_message, send_photo, telegram_message_id
logger = logging.getLogger(__name__)
def to_telegram_payload(signal: IndicatorSignal, ticker: str) -> SignalPayload:
tick = get_tick_size(ticker)
return SignalPayload(
ticker=ticker,
action=signal.side,
entry_price=format_px(signal.entry, tick),
current_price=format_px(signal.close, tick),
stop_loss_price=format_px(signal.sl, tick),
take_profit_1_price=format_px(signal.tp1, tick),
take_profit_2_price=format_px(signal.tp2, tick) or None,
take_profit_3_price=format_px(signal.tp3, tick) or None,
visual_timeframe=signal.visual_timeframe,
signal_sequence=1,
signal_time=signal.bar_open_ts,
)
async def deliver_telegram(
settings: Settings,
signal: SignalPayload,
message_thread_id: int,
*,
store: ScannerStore | None = None,
) -> None:
"""Fetch chart + post to Telegram (used by inbound webhook and scanner)."""
caption = format_caption(signal)
photo: bytes | None = None
chart_error: str | None = None
owns_store = store is None
db = store or ScannerStore(settings.scanner_state_path)
reply_id = db.get_last_tg_message(message_thread_id, signal.ticker)
try:
symbol = to_binance_symbol(signal.ticker)
interval = to_binance_interval(signal.visual_timeframe)
limit = chart_kline_limit(interval)
if signal.prev_signal_time and signal.prev_signal_time > 0:
bar_sec = interval_timedelta(interval).total_seconds()
if bar_sec > 0:
span = int((signal.signal_time - signal.prev_signal_time) / bar_sec) + 24
limit = min(1500, max(limit, span))
df = await fetch_klines(symbol, interval, limit=limit)
# LTF 15m: chart shows TP1 + remainder green (like FVG); caption still has TP2/TP3.
chart_tp2 = None if interval == "15m" else signal.take_profit_2_price
chart_tp3 = None if interval == "15m" else signal.take_profit_3_price
photo = await asyncio.to_thread(
render_setup_chart,
df,
ticker=signal.ticker,
action=signal.action.value,
entry=signal.entry_price,
stop_loss=signal.stop_loss_price,
tp1=signal.take_profit_1_price,
tp2=chart_tp2,
tp3=chart_tp3,
timeframe=signal.visual_timeframe,
current_price=(
signal.current_price if signal.signal_sequence > 1 else None
),
signal_time=signal.signal_time,
right_pad=chart_right_pad(interval),
prev_entry=signal.prev_entry_price,
prev_side=signal.prev_side,
prev_signal_time=signal.prev_signal_time,
)
except Exception as exc: # noqa: BLE001 — fallback to text-only post
chart_error = str(exc)
logger.exception("Chart generation failed, falling back to text-only: %s", exc)
try:
if photo is not None:
response = await send_photo(
settings,
photo=photo,
caption=caption,
message_thread_id=message_thread_id,
reply_to_message_id=reply_id,
)
else:
response = await send_message(
settings,
text=caption,
message_thread_id=message_thread_id,
reply_to_message_id=reply_id,
)
message_id = telegram_message_id(response)
if message_id is not None:
db.set_last_tg_message(message_thread_id, signal.ticker, message_id)
logger.info(
"Delivered %s: %s %s seq=%s thread=%s chart_error=%s",
"photo" if photo is not None else "text",
signal.ticker,
signal.action.value,
signal.signal_sequence,
message_thread_id,
chart_error,
)
except TelegramError as exc:
logger.exception("Telegram delivery failed: %s", exc)
except Exception as exc: # noqa: BLE001
logger.exception("Unexpected delivery error: %s", exc)
finally:
if owns_store:
db.close()
async def deliver_generated(
settings: Settings,
signal: IndicatorSignal,
symbol: str,
message_thread_id: int,
*,
store: ScannerStore | None = None,
) -> None:
ticker = to_tv_perp_ticker(symbol)
telegram_payload = to_telegram_payload(signal, ticker)
if store is not None:
previous = store.get_last_open(signal.strategy_id, symbol)
if previous is not None:
prev_side, prev_entry, prev_ts = previous
if prev_side != signal.side:
tick = get_tick_size(ticker)
telegram_payload = telegram_payload.model_copy(
update={
"is_reversal": True,
"realized_pnl_pct": realized_pnl_pct(
prev_side, prev_entry, signal.entry
),
"prev_side": prev_side,
"prev_entry_price": format_px(prev_entry, tick),
"prev_signal_time": prev_ts or None,
}
)
heryon_payload = build_heryon_payload(signal, symbol=symbol, settings=settings)
await deliver_telegram(
settings,
telegram_payload,
message_thread_id,
store=store,
)
if store is not None:
store.set_last_open(
signal.strategy_id,
symbol,
signal.side,
signal.entry,
signal.bar_open_ts,
)
try:
await send_heryon(settings, heryon_payload)
except Exception as exc: # noqa: BLE001
logger.exception(
"Heryon send failed nonce=%s: %s", heryon_payload.get("nonce"), exc
)