diff --git a/.github/workflows/ubuntu-smoke.yml b/.github/workflows/ubuntu-smoke.yml index f28ee05..0e2353d 100644 --- a/.github/workflows/ubuntu-smoke.yml +++ b/.github/workflows/ubuntu-smoke.yml @@ -140,6 +140,12 @@ jobs: - name: Run protocol_v1_named_pipe_bridge_smoke run: ./build-linux/protocol_v1_named_pipe_bridge_smoke --self-test + - name: Build telegram_archive_parser_smoke + run: cmake --build build-linux --target telegram_archive_parser_smoke -j + + - name: Run telegram_archive_parser_smoke + run: ./build-linux/telegram_archive_parser_smoke + - name: Build market_data_subscription_contract_test run: cmake --build build-linux --target market_data_subscription_contract_test -j diff --git a/.github/workflows/windows-smoke.yml b/.github/workflows/windows-smoke.yml index 58a5475..d6aa60e 100644 --- a/.github/workflows/windows-smoke.yml +++ b/.github/workflows/windows-smoke.yml @@ -71,6 +71,9 @@ jobs: - name: Build Bridge Protocol v1 named-pipe smoke example run: cmake --build build-windows --config Debug --target protocol_v1_named_pipe_bridge_smoke + - name: Build Telegram archive parser smoke example + run: cmake --build build-windows --config Debug --target telegram_archive_parser_smoke + - name: Run MetaTrader file tests shell: pwsh run: | @@ -92,6 +95,7 @@ jobs: .\build-windows\Debug\named_pipe_bridge_smoke.exe --self-test .\build-windows\Debug\protocol_v1_bridge_smoke.exe --self-test .\build-windows\Debug\protocol_v1_named_pipe_bridge_smoke.exe --self-test + .\build-windows\Debug\telegram_archive_parser_smoke.exe - name: Test MetaEditor compile smoke script shell: pwsh diff --git a/CMakeLists.txt b/CMakeLists.txt index 776366e..f59e6c3 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -738,6 +738,44 @@ if(OPTIONX_BUILD_EXAMPLES) ) endif() + add_executable(telegram_signal_bridge_smoke examples/telegram_signal_bridge_smoke.cpp) + target_compile_features(telegram_signal_bridge_smoke PRIVATE cxx_std_17) + target_include_directories(telegram_signal_bridge_smoke PRIVATE + ${EXAMPLE_INCLUDE_DIRS} + ${EXAMPLE_DEPS_INCLUDE_DIRS} + ) + target_link_directories(telegram_signal_bridge_smoke PRIVATE ${EXAMPLE_LIBRARY_DIRS}) + target_compile_definitions( + telegram_signal_bridge_smoke PRIVATE + ${EXAMPLE_DEFINES} + LOGIT_BASE_PATH="${LOGIT_BASE_PATH_FWD}" + ) + if(MINGW) + target_compile_options(telegram_signal_bridge_smoke PRIVATE -Wa,-mbig-obj) + elseif(MSVC) + target_compile_options(telegram_signal_bridge_smoke PRIVATE /bigobj) + endif() + target_link_libraries(telegram_signal_bridge_smoke PRIVATE ${EXAMPLE_LIBS} optionx_cpp) + + add_executable(telegram_archive_parser_smoke examples/telegram_archive_parser_smoke.cpp) + target_compile_features(telegram_archive_parser_smoke PRIVATE cxx_std_17) + target_include_directories(telegram_archive_parser_smoke PRIVATE + ${EXAMPLE_INCLUDE_DIRS} + ${EXAMPLE_DEPS_INCLUDE_DIRS} + ) + target_link_directories(telegram_archive_parser_smoke PRIVATE ${EXAMPLE_LIBRARY_DIRS}) + target_compile_definitions( + telegram_archive_parser_smoke PRIVATE + ${EXAMPLE_DEFINES} + LOGIT_BASE_PATH="${LOGIT_BASE_PATH_FWD}" + ) + if(MINGW) + target_compile_options(telegram_archive_parser_smoke PRIVATE -Wa,-mbig-obj) + elseif(MSVC) + target_compile_options(telegram_archive_parser_smoke PRIVATE /bigobj) + endif() + target_link_libraries(telegram_archive_parser_smoke PRIVATE ${EXAMPLE_LIBS} optionx_cpp) + add_executable(protocol_v1_bridge_smoke examples/protocol_v1_bridge_smoke.cpp) target_compile_features(protocol_v1_bridge_smoke PRIVATE cxx_std_17) @@ -888,6 +926,9 @@ if(OPTIONX_BUILD_TESTS) bridge_host_test metatrader_paths_test telegram_dto_test + telegram_signal_parser_test + telegram_signal_bridge_test + telegram_worker_source_test ) if(OPTIONX_LIGHTWEIGHT_BRIDGE_SMOKE_TESTS) list(APPEND OPTIONX_LIGHTWEIGHT_TESTS diff --git a/examples/README.md b/examples/README.md index 46c10c0..b0650af 100644 --- a/examples/README.md +++ b/examples/README.md @@ -55,6 +55,14 @@ Currently maintained examples: - `protocol_v1_named_pipe_bridge_smoke.cpp` starts Bridge Protocol v1 over a local named pipe and, on Windows, can run `--self-test` with a local pipe client. +- `telegram_signal_bridge_smoke.cpp` demonstrates the Telegram parser and live + bridge lifecycle with a deterministic in-memory source. The production + source boundary can be backed by `tg-client-stdio`; this example needs no + Telegram credentials. +- `telegram_archive_parser_smoke.cpp` replays exported Telegram raw-message + records through the parser and keeps executable signals, outcomes and + diagnostics separate. It uses deterministic fixtures and needs no Telegram + credentials. - `metatrader_file_bridge_smoke.cpp` runs the C++ side of the MetaTrader Common\Files bridge against a temporary command/event layout. - `metatrader_file_command_writer_smoke.cpp` demonstrates the C++ command-writer diff --git a/examples/telegram_archive_parser_smoke.cpp b/examples/telegram_archive_parser_smoke.cpp new file mode 100644 index 0000000..4179164 --- /dev/null +++ b/examples/telegram_archive_parser_smoke.cpp @@ -0,0 +1,81 @@ +/// \file telegram_archive_parser_smoke.cpp +/// \brief Demonstrates replaying exported Telegram messages through the parser. + +#include + +#include +#include +#include +#include +#include + +namespace { + +optionx::bridges::telegram::TelegramRawMessage make_message( + std::int64_t message_id, + std::string text, + std::int64_t reply_to_message_id = 0); + +std::vector demo_archive(); + +void print_parsed( + const optionx::bridges::telegram::TelegramParsedMessage& parsed); + +} // namespace + +int main() { + // In production these records arrive one at a time from + // tg-client-stdio::WorkerClient::stream_messages(). Keeping the replay + // loop source-independent makes parser tests and backtests deterministic. + optionx::bridges::telegram::TelegramSignalParser parser; + for (const auto& raw : demo_archive()) { + print_parsed(parser.parse(raw)); + } + return 0; +} + +namespace { + +optionx::bridges::telegram::TelegramRawMessage make_message( + const std::int64_t message_id, + std::string text, + const std::int64_t reply_to_message_id) { + optionx::bridges::telegram::TelegramRawMessage message; + message.chat_id = "-1001234567890"; + message.chat_title = "Demo archive"; + message.message_id = message_id; + message.date_ms = 1800000000000 + message_id * 1000; + message.reply_to_message_id = reply_to_message_id; + message.text = std::move(text); + return message; +} + +std::vector demo_archive() { + return { + make_message(100, "EURUSD BUY 5m"), + make_message(101, "EURUSD WIN", 100), + make_message(102, "USDJPY BUY 2 weeks"), + }; +} + +void print_parsed( + const optionx::bridges::telegram::TelegramParsedMessage& parsed) { + std::cout << "message=" << parsed.raw.message_identity() + << " signals=" << parsed.signals.size() + << " outcomes=" << parsed.outcomes.size() + << " diagnostics=" << parsed.diagnostics.size() << '\n'; + for (const auto& signal : parsed.signals) { + std::cout << " signal=" << signal.symbol + << " direction=" << optionx::to_str(signal.order_type) + << " duration=" << signal.duration << '\n'; + } + for (const auto& outcome : parsed.outcomes) { + std::cout << " outcome=" << outcome.symbol + << " result=" << static_cast(outcome.result) << '\n'; + } + for (const auto& diagnostic : parsed.diagnostics) { + std::cout << " diagnostic=" << diagnostic.code << '\n'; + } +} + +} // namespace diff --git a/examples/telegram_signal_bridge_smoke.cpp b/examples/telegram_signal_bridge_smoke.cpp new file mode 100644 index 0000000..86b8a4f --- /dev/null +++ b/examples/telegram_signal_bridge_smoke.cpp @@ -0,0 +1,113 @@ +/// \file telegram_signal_bridge_smoke.cpp +/// \brief Demonstrates Telegram parsing and bridge dispatch with a fake source. + +#include + +#include +#include +#include +#include + +namespace { + +class DemoMessageSource final + : public optionx::bridges::telegram::TelegramMessageSource { +public: + bool start(message_callback_t on_message, error_callback_t on_error) override; + void stop() noexcept override; + void emit(const optionx::bridges::telegram::TelegramRawMessage& message); + +private: + message_callback_t m_on_message; +}; + +std::unique_ptr +make_config(); + +optionx::bridges::telegram::TelegramRawMessage make_message( + std::string text, + std::int64_t message_id); + +} // namespace + +int main() { + // The production source can be an adapter around tg-client-stdio. A fake + // source keeps this example deterministic and runnable without Telegram + // credentials or an authorized session. + auto source = std::make_shared(); + optionx::bridges::telegram::TelegramSignalBridge bridge(source); + if (!bridge.configure(make_config())) { + std::cerr << "failed to configure Telegram bridge\n"; + return 1; + } + + std::int64_t next_signal_id = 0; + bridge.on_signal_id() = [&next_signal_id]() { + return ++next_signal_id; + }; + bridge.on_status_update() = [](const optionx::BridgeStatusUpdate& update) { + std::cout << "status=" << optionx::to_str(update.status); + if (!update.message.empty()) { + std::cout << " message=" << update.message; + } + std::cout << '\n'; + }; + bridge.on_signal_report() = [](const optionx::BridgeSignalReport& report) { + std::cout << "report=" << report.reason_code + << " status=" << optionx::to_str(report.status) << '\n'; + }; + bridge.on_trade_signal() = [](std::unique_ptr signal) { + std::cout << "signal=" << signal->symbol + << " direction=" << optionx::to_str(signal->order_type) + << " duration=" << signal->duration + << " id=" << signal->signal_id << '\n'; + }; + + bridge.run(); + source->emit(make_message("EURUSD BUY 5m", 1)); + source->emit(make_message("EURUSD BUY 5m", 1)); + bridge.shutdown(); + return 0; +} + +namespace { + +bool DemoMessageSource::start(message_callback_t on_message, error_callback_t on_error) { + (void)on_error; + m_on_message = std::move(on_message); + return true; +} + +void DemoMessageSource::stop() noexcept { + m_on_message = {}; +} + +void DemoMessageSource::emit( + const optionx::bridges::telegram::TelegramRawMessage& message) { + if (m_on_message) { + m_on_message(message); + } +} + +std::unique_ptr +make_config() { + auto config = std::make_unique< + optionx::bridges::telegram::TelegramSignalBridgeConfig>(); + config->bridge_id = 9001; + config->fixed_amount = 1.0; + return config; +} + +optionx::bridges::telegram::TelegramRawMessage make_message( + std::string text, + const std::int64_t message_id) { + optionx::bridges::telegram::TelegramRawMessage message; + message.chat_id = "demo-signals"; + message.chat_title = "Demo Telegram Signals"; + message.message_id = message_id; + message.date_ms = 1800000000000; + message.text = std::move(text); + return message; +} + +} // namespace diff --git a/guides/bridge-examples.md b/guides/bridge-examples.md index 93669e0..56646a0 100644 --- a/guides/bridge-examples.md +++ b/guides/bridge-examples.md @@ -15,6 +15,7 @@ families are grouped by external protocol or adapter contract, not by transport. | BotBinary/BinaryBot | `optionx_cpp/bridges/bot_binary.hpp` | `BotBinaryBridgeConfig` | `examples/bot_binary_bridge_smoke.cpp` | Compatibility intake for BotBinary `request=...` HTTP URLs and file-signal filenames. | | BotBinary command helper | `optionx_cpp/bridges/bot_binary.hpp` | none | `examples/bot_binary_command_builder_smoke.cpp` | Formatter/parser helper for legacy BotBinary command strings. | | Legacy trading pipe | `optionx_cpp/bridges/legacy_trading.hpp` | `LegacyTradingBridgeConfig` | `examples/named_pipe_bridge_smoke.cpp` | Compatibility bridge for the older named-pipe JSON trading protocol. | +| Telegram signal bridge | `optionx_cpp/bridges/telegram.hpp` | `TelegramSignalBridgeConfig` | `examples/telegram_signal_bridge_smoke.cpp` | Deterministic parser, source boundary, signal callback, duplicate report, and shutdown lifecycle without Telegram credentials. | ## Choosing A Bridge @@ -29,6 +30,10 @@ families are grouped by external protocol or adapter contract, not by transport. BotBinary command string or file-signal filename. - Use the legacy trading pipe only for old clients that already speak that named-pipe JSON format. +- Use the Telegram signal bridge when a user-client source must turn channel + messages into normalized signals. The production source is expected to be a + `tg-client-stdio` adapter; the parser and bridge can be tested independently + with a fake source. Compatibility bridges should convert their external payload into `TradeSignal` callbacks and reports. They do not need to expose Bridge Protocol v1 endpoints diff --git a/guides/bridge-taxonomy.md b/guides/bridge-taxonomy.md index b1369f7..546e15e 100644 --- a/guides/bridge-taxonomy.md +++ b/guides/bridge-taxonomy.md @@ -28,11 +28,18 @@ OptionX protocol together. | TradingView extension | `optionx_cpp/bridges/trading_view.hpp` | Adapter for payloads emitted by `browser_extensions/tradingview-alert-extension`. | HTTP. | | BinaryBot/BotBinary | `optionx_cpp/bridges/bot_binary.hpp` | Compatibility bridge and formatter/parser helpers for observed BinaryBot-compatible command strings. | HTTP `request=...`, file-signal name. | | Legacy trading pipe | `optionx_cpp/bridges/legacy_trading.hpp` | Compatibility bridge for the older named-pipe JSON trading protocol. | Named pipe. | +| Telegram signal bridge | `optionx_cpp/bridges/telegram.hpp` | User-client message parser and live signal adapter. The Telegram worker/session remains an external source boundary. | stdio worker source, with source adapters kept outside the parser. | All families converge internally on OptionX DTOs such as `TradeSignal`, `TradeRequest`, account snapshots and bridge callbacks. The public wire format does not need to be the same for every family. +Telegram is a source adapter family rather than a transport-only family. The +public bridge consumes `TelegramMessageSource` callbacks, while authorization, +proxy handling, dialog discovery and historical export belong to the +`tg-client-stdio` worker/supervisor layer. Historical export remains a separate +archive capability and is not added to `BaseBridge`. + For practical embedding of the native HTTP/WebSocket server bridge, see `guides/protocol-v1-bridge-runtime.md`. For runnable bridge entry points, see `examples/README.md`. diff --git a/guides/telegram-bridge-design.md b/guides/telegram-bridge-design.md index c4d97c2..bed09d4 100644 --- a/guides/telegram-bridge-design.md +++ b/guides/telegram-bridge-design.md @@ -1,7 +1,8 @@ # Telegram Bridge Design -This document captures the intended direction for a future Telegram signal -bridge. It is a design note, not a committed public API. +This document captures the architecture and current boundaries of the Telegram +signal bridge. The public C++ DTO/parser/bridge layer is implemented, while +the authorized Telegram worker adapter remains a separate integration step. ## Problem Shape @@ -53,6 +54,13 @@ newline-delimited JSON over stdin/stdout. Keep the protocol versioned so a framed transport can replace JSONL later if media bytes ever need to cross the stdio boundary. +The current OptionX bridge does not own a Telethon process. It consumes the +`TelegramMessageSource` interface, so fake sources can exercise parsing and +lifecycle without credentials. `TelegramWorkerMessageSource` now binds that +interface to a WorkerClient-shaped process adapter; the parent repository will +pin the concrete `tg-client-stdio::WorkerClient` only after its stacked worker +PRs are merged. + ## Stdio Protocol Envelope Every JSONL record should use an envelope so responses, long-running exports, @@ -107,7 +115,7 @@ Keep the live bridge, archive export, and parsing separate. ### Worker Client -`TelegramWorkerClient` owns the process/session protocol. It should expose +`tg-client-stdio::WorkerClient` owns the process/session protocol. It exposes operations such as: - `auth.status`; @@ -131,6 +139,11 @@ message events into `TradeSignal` callbacks and signal reports. It should not expose historical export through `BaseBridge::run()` or `process()`. Bridge lifecycle remains live-intake lifecycle. +The current bridge also applies a bounded identity-based dedupe cache. Parser +diagnostics, duplicate messages, allocator failures and callback failures are +reported through `BridgeSignalReport`; they do not silently become accepted +signals. + ### Archive Source Historical export is a separate capability. The first implementation can be @@ -308,17 +321,30 @@ Proxy config should support at least SOCKS5 and HTTP where the underlying Telegram client library supports them. Proxy failures must be distinct from authorization failures. -## First PR Sequence +## Implementation Status And Next Steps + +Completed without an authorized Telegram session: + +1. `tg-client-stdio` worker protocol for dialogs, streaming export, + live listen/stop, auth status/code/2FA and HTTP/SOCKS proxy configuration. +2. C++ worker supervisor and typed raw-message export DTOs in the standalone + worker repository. +3. OptionX raw/parsed Telegram DTOs, deterministic regex parser and + source-independent `TelegramSignalBridge`. +4. Fake-source unit coverage and a runnable no-credentials bridge example. + +Current no-credentials examples include a live fake-source bridge smoke and a +deterministic archive/parser replay. The latter uses the same raw-message shape +that `messages.export` streams, while keeping parser outcomes separate from +executable signals. + +Next steps: -1. Refactor `telegram-monitoring-tool` into a non-interactive worker command - with the JSONL envelope above, preserving the current interactive CLI as a - thin wrapper if needed. -2. Add worker operations for `dialogs.list`, `messages.export` and - `messages.listen`. -3. Add C++ protocol DTOs and a small worker client/supervisor in OptionX. -4. Add `TelegramSignalParser` with pure text fixture tests. -5. Add `TelegramSignalBridge` live intake using the parser. -6. Add archive/parser example for historical backtest fixture generation. +1. Merge and pin the worker repository's supervisor/archive PRs. +2. Pin the merged worker repository in an OptionX consumer and run the adapter + against the mock worker process. +3. Perform the first real authorization, proxy and live-channel check with an + operator-provided Telegram session. -Keep the first OptionX PR focused on DTOs, parser and docs if the worker is not -ready yet. +OCR/vision remains a separate optional provider and should not block the text +parser or the first authorized-session test. diff --git a/include/optionx_cpp/bridges/telegram.hpp b/include/optionx_cpp/bridges/telegram.hpp index 8aeaf51..962a744 100644 --- a/include/optionx_cpp/bridges/telegram.hpp +++ b/include/optionx_cpp/bridges/telegram.hpp @@ -7,5 +7,9 @@ #include "bridges/telegram/TelegramParsedMessage.hpp" #include "bridges/telegram/TelegramRawMessage.hpp" +#include "bridges/telegram/TelegramSignalParser.hpp" +#include "bridges/telegram/TelegramSignalBridgeConfig.hpp" +#include "bridges/telegram/TelegramSignalBridge.hpp" +#include "bridges/telegram/TelegramWorkerMessageSource.hpp" #endif // OPTIONX_HEADER_BRIDGES_TELEGRAM_HPP_INCLUDED diff --git a/include/optionx_cpp/bridges/telegram/TelegramRawMessage.hpp b/include/optionx_cpp/bridges/telegram/TelegramRawMessage.hpp index 706a19c..d30c303 100644 --- a/include/optionx_cpp/bridges/telegram/TelegramRawMessage.hpp +++ b/include/optionx_cpp/bridges/telegram/TelegramRawMessage.hpp @@ -73,7 +73,6 @@ namespace optionx::bridges::telegram { /// \brief Serializes the DTO using worker-compatible field names. nlohmann::json to_json() const { - validate(); return nlohmann::json{ {"chat_id", chat_id}, {"chat_title", chat_title}, diff --git a/include/optionx_cpp/bridges/telegram/TelegramSignalBridge.hpp b/include/optionx_cpp/bridges/telegram/TelegramSignalBridge.hpp new file mode 100644 index 0000000..ae3bc09 --- /dev/null +++ b/include/optionx_cpp/bridges/telegram/TelegramSignalBridge.hpp @@ -0,0 +1,439 @@ +#pragma once +#ifndef OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_BRIDGE_HPP_INCLUDED +#define OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_BRIDGE_HPP_INCLUDED + +/// \file TelegramSignalBridge.hpp +/// \brief Converts Telegram message events into OptionX TradeSignal objects. + +#include "bridges/BaseBridge.hpp" +#include "bridges/detail/BridgeTradeSignalValidation.hpp" +#include "bridges/telegram/TelegramSignalBridgeConfig.hpp" + +#include +#include +#include +#include +#include +#include +#include + +namespace optionx::bridges::telegram { + + /// \class TelegramMessageSource + /// \brief Minimal live-message source boundary used by the Telegram bridge. + /// + /// A concrete adapter may be backed by tg-client-stdio, a test fixture, or + /// another Telegram client. The source must stop invoking callbacks before + /// `stop()` returns. + class TelegramMessageSource { + public: + using message_callback_t = std::function; + using error_callback_t = std::function; + + virtual ~TelegramMessageSource() = default; + virtual bool start(message_callback_t on_message, + error_callback_t on_error) = 0; + virtual void stop() noexcept = 0; + }; + + /// \class TelegramSignalBridge + /// \brief Publishes executable signals parsed from live Telegram messages. + class TelegramSignalBridge final : public BaseBridge { + private: + struct RuntimeState { + std::mutex mutex; + bridge_status_callback_t status_callback; + BaseBridge::trade_signal_callback_t trade_signal_callback; + BaseBridge::signal_report_callback_t signal_report_callback; + BaseBridge::signal_id_allocator_t signal_id_allocator; + std::shared_ptr source; + std::deque dedupe_order; + std::unordered_set dedupe_keys; + bool running = false; + }; + + public: + explicit TelegramSignalBridge( + std::shared_ptr source = {}) + : m_state(std::make_shared()), + m_source(std::move(source)) {} + + ~TelegramSignalBridge() override { + shutdown(); + } + + /// \brief Replaces the live source while the bridge is stopped. + bool set_message_source(std::shared_ptr source) { + std::lock_guard lock(m_state->mutex); + if (m_state->running) { + return false; + } + m_source = std::move(source); + return true; + } + + bool configure(std::unique_ptr config) override { + if (!config) { + return false; + } + const auto* typed = dynamic_cast(config.get()); + if (!typed) { + config->dispatch_callbacks(false, "Invalid Telegram signal bridge config type."); + return false; + } + auto next_config = std::make_shared(*typed); + const auto validation = next_config->validate(); + config->dispatch_callbacks(validation.first, validation.second); + if (!validation.first) { + return false; + } + std::lock_guard lock(m_config_mutex); + m_config = std::move(next_config); + return true; + } + + bridge_status_callback_t& on_status_update() override { + return m_state->status_callback; + } + + trade_signal_callback_t& on_trade_signal() override { + return m_state->trade_signal_callback; + } + + signal_report_callback_t& on_signal_report() override { + return m_state->signal_report_callback; + } + + signal_id_allocator_t& on_signal_id() override { + return m_state->signal_id_allocator; + } + + void update_account_info(const AccountInfoUpdate& info) override { + (void)info; + } + + void run() override { + const auto config = get_config(); + if (!config) { + notify_status(BridgeStatus::SERVER_START_FAILED, + "Telegram bridge is not configured."); + return; + } + auto source = get_source(); + if (!source) { + notify_status(BridgeStatus::SERVER_START_FAILED, + "Telegram bridge has no message source."); + return; + } + if (!get_signal_id_allocator()) { + notify_status(BridgeStatus::SERVER_START_FAILED, + "Telegram bridge requires a signal ID allocator."); + return; + } + + { + std::lock_guard lock(m_state->mutex); + if (m_state->running) { + return; + } + m_state->running = true; + m_state->source = source; + m_state->dedupe_keys.clear(); + m_state->dedupe_order.clear(); + } + + try { + const auto parser = TelegramSignalParser(config->parser); + const bool started = source->start( + [state = m_state, config, parser](const TelegramRawMessage& raw) { + process_message(state, *config, parser, raw); + }, + [state = m_state](const std::string& message) { + notify_status(state, BridgeStatus::CONNECTION_ERROR, message); + }); + if (!started) { + set_running(false); + notify_status(BridgeStatus::SERVER_START_FAILED, + "Telegram message source failed to start."); + return; + } + notify_status(BridgeStatus::SERVER_STARTED, {}); + } + catch (const std::exception& error) { + set_running(false); + notify_status(BridgeStatus::SERVER_START_FAILED, error.what()); + } + catch (...) { + set_running(false); + notify_status(BridgeStatus::SERVER_START_FAILED, + "Telegram message source threw an unknown exception."); + } + } + + void shutdown() override { + std::shared_ptr source; + bool was_running = false; + { + std::lock_guard lock(m_state->mutex); + was_running = m_state->running; + m_state->running = false; + source = m_state->source; + m_state->source.reset(); + } + if (source) { + try { + source->stop(); + } + catch (...) { + notify_status(BridgeStatus::CONNECTION_ERROR, + "Telegram message source threw during stop."); + } + } + if (was_running) { + notify_status(BridgeStatus::SERVER_STOPPED, {}); + } + } + + private: + std::shared_ptr get_config() const { + std::lock_guard lock(m_config_mutex); + return m_config; + } + + std::shared_ptr get_source() const { + std::lock_guard lock(m_state->mutex); + return m_source; + } + + BaseBridge::signal_id_allocator_t get_signal_id_allocator() const { + std::lock_guard lock(m_state->mutex); + return m_state->signal_id_allocator; + } + + void set_running(const bool running) { + std::lock_guard lock(m_state->mutex); + m_state->running = running; + if (!running) { + m_state->source.reset(); + } + } + + static void notify_status( + const std::shared_ptr& state, + const BridgeStatus status, + const std::string& message) { + bridge_status_callback_t callback; + { + std::lock_guard lock(state->mutex); + callback = state->status_callback; + } + if (callback) { + try { + callback({status, {}, message}); + } + catch (...) { + } + } + } + + void notify_status(const BridgeStatus status, const std::string& message) const { + notify_status(m_state, status, message); + } + + static void emit_report( + const std::shared_ptr& state, + BridgeSignalReport report) { + signal_report_callback_t callback; + { + std::lock_guard lock(state->mutex); + callback = state->signal_report_callback; + } + if (callback) { + try { + callback(report); + } + catch (...) { + } + } + } + + static std::string make_dedupe_key( + const TelegramRawMessage& raw, + const TelegramParsedSignal& parsed, + const std::size_t index) { + return raw.message_identity() + "|" + std::to_string(index) + "|" + + parsed.symbol + "|" + optionx::to_str(parsed.order_type) + "|" + + optionx::to_str(parsed.option_type) + "|" + + std::to_string(parsed.duration) + "|" + + std::to_string(parsed.expiry_time); + } + + static void process_message( + const std::shared_ptr& state, + const TelegramSignalBridgeConfig& config, + const TelegramSignalParser& parser, + const TelegramRawMessage& raw) { + try { + raw.validate(); + const auto parsed = parser.parse(raw); + for (const auto& diagnostic : parsed.diagnostics) { + BridgeSignalReport report; + report.bridge_id = config.bridge_id; + report.bridge_type = BridgeType::TELEGRAM_SIGNAL; + report.status = BridgeSignalReportStatus::INVALID; + report.reason_code = diagnostic.code; + report.message = diagnostic.message; + report.event_id = raw.message_identity(); + report.raw_payload = raw.to_json(); + report.context = { + {"offset", diagnostic.offset}, + {"length", diagnostic.length}, + }; + emit_report(state, std::move(report)); + } + + for (std::size_t index = 0; index < parsed.signals.size(); ++index) { + const auto& parsed_signal = parsed.signals[index]; + auto signal = std::make_unique(); + signal->bridge_id = config.bridge_id; + signal->symbol = parsed_signal.symbol; + signal->order_type = parsed_signal.order_type; + signal->option_type = parsed_signal.option_type; + signal->duration = parsed_signal.duration; + signal->expiry_time = parsed_signal.expiry_time; + signal->signal_name = parsed_signal.signal_name; + signal->comment = raw.text; + signal->amount = config.fixed_amount; + const auto dedupe_key = make_dedupe_key(raw, parsed_signal, index); + signal->unique_hash = dedupe_key; + + detail::validate_executable_trade_signal( + *signal, "Telegram signal", true); + + BaseBridge::signal_id_allocator_t allocator; + BaseBridge::trade_signal_callback_t callback; + bool duplicate = false; + { + std::lock_guard lock(state->mutex); + if (!state->running) { + return; + } + if (state->dedupe_keys.find(dedupe_key) != state->dedupe_keys.end()) { + duplicate = true; + } + else { + state->dedupe_keys.insert(dedupe_key); + state->dedupe_order.push_back(dedupe_key); + while (state->dedupe_order.size() > config.dedupe_cache_size) { + state->dedupe_keys.erase(state->dedupe_order.front()); + state->dedupe_order.pop_front(); + } + allocator = state->signal_id_allocator; + callback = state->trade_signal_callback; + } + } + if (duplicate) { + BridgeSignalReport report; + report.bridge_id = config.bridge_id; + report.bridge_type = BridgeType::TELEGRAM_SIGNAL; + report.status = BridgeSignalReportStatus::DUPLICATE; + report.reason_code = "duplicate_message"; + report.message = "Telegram signal was already dispatched."; + report.event_id = raw.message_identity(); + report.dedupe_key = dedupe_key; + report.symbol = parsed_signal.symbol; + report.signal_name = parsed_signal.signal_name; + report.raw_payload = raw.to_json(); + emit_report(state, std::move(report)); + continue; + } + + try { + signal->signal_id = allocator(); + if (signal->signal_id == 0) { + throw std::runtime_error("Telegram signal ID allocator returned zero."); + } + } + catch (const std::exception& error) { + { + std::lock_guard lock(state->mutex); + state->dedupe_keys.erase(dedupe_key); + } + emit_report(state, BridgeSignalReport{ + config.bridge_id, + BridgeType::TELEGRAM_SIGNAL, + BridgeSignalReportStatus::INTAKE_ERROR, + "signal_id_allocation_failed", + error.what(), + {}, + raw.message_identity(), + dedupe_key, + parsed_signal.symbol, + parsed_signal.signal_name, + {}, + raw.to_json(), + {}, + {}, + 0, + raw.date_ms, + }); + continue; + } + if (callback) { + try { + callback(std::move(signal)); + } + catch (...) { + emit_report(state, BridgeSignalReport{ + config.bridge_id, + BridgeType::TELEGRAM_SIGNAL, + BridgeSignalReportStatus::INTAKE_ERROR, + "trade_signal_callback_failed", + "Telegram trade signal callback threw.", + {}, + raw.message_identity(), + dedupe_key, + parsed_signal.symbol, + parsed_signal.signal_name, + {}, + raw.to_json(), + {}, + {}, + 0, + raw.date_ms, + }); + } + } + } + } + catch (const std::exception& error) { + emit_report(state, BridgeSignalReport{ + config.bridge_id, + BridgeType::TELEGRAM_SIGNAL, + BridgeSignalReportStatus::INVALID, + "telegram_message_parse_failed", + error.what(), + {}, + raw.message_identity(), + {}, + {}, + {}, + {}, + raw.to_json(), + {}, + {}, + 0, + raw.date_ms, + }); + } + } + + std::shared_ptr m_state; + mutable std::mutex m_config_mutex; + std::shared_ptr m_config; + std::shared_ptr m_source; + }; + +} // namespace optionx::bridges::telegram + +#endif // OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_BRIDGE_HPP_INCLUDED diff --git a/include/optionx_cpp/bridges/telegram/TelegramSignalBridgeConfig.hpp b/include/optionx_cpp/bridges/telegram/TelegramSignalBridgeConfig.hpp new file mode 100644 index 0000000..0578f01 --- /dev/null +++ b/include/optionx_cpp/bridges/telegram/TelegramSignalBridgeConfig.hpp @@ -0,0 +1,122 @@ +#pragma once +#ifndef OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_BRIDGE_CONFIG_HPP_INCLUDED +#define OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_BRIDGE_CONFIG_HPP_INCLUDED + +/// \file TelegramSignalBridgeConfig.hpp +/// \brief Configuration for the Telegram signal bridge. + +#include "data/bridge.hpp" +#include "bridges/telegram/TelegramSignalParser.hpp" + +#include +#include +#include + +namespace optionx::bridges::telegram { + + /// \class TelegramSignalBridgeConfig + /// \brief Parser and dispatch settings for Telegram live signal intake. + class TelegramSignalBridgeConfig final : public IBridgeConfig { + public: + BridgeId bridge_id = 0; + double fixed_amount = 0.0; + std::size_t dedupe_cache_size = 4096; + TelegramParserConfig parser = TelegramSignalParser::default_config(); + + void to_json(nlohmann::json& j) const override { + j = nlohmann::json{ + {"bridge_id", bridge_id}, + {"fixed_amount", fixed_amount}, + {"dedupe_cache_size", dedupe_cache_size}, + {"use_chat_title_as_signal_name", parser.use_chat_title_as_signal_name}, + {"signal_rules", nlohmann::json::array()}, + {"outcome_rules", nlohmann::json::array()}, + }; + for (const auto& rule : parser.signal_rules) { + j["signal_rules"].push_back({ + {"name", rule.name}, + {"pattern", rule.pattern}, + {"symbol_group", rule.symbol_group}, + {"direction_group", rule.direction_group}, + {"expiry_group", rule.expiry_group}, + {"unit_group", rule.unit_group}, + {"option_type", rule.option_type}, + }); + } + for (const auto& rule : parser.outcome_rules) { + j["outcome_rules"].push_back({ + {"name", rule.name}, + {"pattern", rule.pattern}, + {"symbol_group", rule.symbol_group}, + {"result_group", rule.result_group}, + }); + } + } + + void from_json(const nlohmann::json& j) override { + bridge_id = j.value("bridge_id", bridge_id); + fixed_amount = j.value("fixed_amount", fixed_amount); + dedupe_cache_size = j.value("dedupe_cache_size", dedupe_cache_size); + parser.use_chat_title_as_signal_name = j.value( + "use_chat_title_as_signal_name", + parser.use_chat_title_as_signal_name); + + if (j.contains("signal_rules")) { + parser.signal_rules.clear(); + for (const auto& item : j.at("signal_rules")) { + TelegramSignalRule rule; + rule.name = item.value("name", ""); + rule.pattern = item.at("pattern").get(); + rule.symbol_group = item.value("symbol_group", 1u); + rule.direction_group = item.value("direction_group", 2u); + rule.expiry_group = item.value("expiry_group", 3u); + rule.unit_group = item.value("unit_group", 4u); + rule.option_type = item.value("option_type", OptionType::SPRINT); + parser.signal_rules.push_back(std::move(rule)); + } + } + if (j.contains("outcome_rules")) { + parser.outcome_rules.clear(); + for (const auto& item : j.at("outcome_rules")) { + TelegramOutcomeRule rule; + rule.name = item.value("name", ""); + rule.pattern = item.at("pattern").get(); + rule.symbol_group = item.value("symbol_group", 1u); + rule.result_group = item.value("result_group", 2u); + parser.outcome_rules.push_back(std::move(rule)); + } + } + } + + std::pair validate() const override { + if (!std::isfinite(fixed_amount) || fixed_amount <= 0.0) { + return {false, "Telegram fixed_amount must be positive and finite."}; + } + if (dedupe_cache_size == 0) { + return {false, "Telegram dedupe_cache_size must be positive."}; + } + try { + (void)TelegramSignalParser(parser); + } + catch (const std::exception& error) { + return {false, std::string("Invalid Telegram parser rules: ") + error.what()}; + } + return {true, {}}; + } + + std::unique_ptr clone_unique() const override { + return std::make_unique(*this); + } + + std::shared_ptr clone_shared() const override { + return std::make_shared(*this); + } + + BridgeType bridge_type() const override { + return BridgeType::TELEGRAM_SIGNAL; + } + }; + +} // namespace optionx::bridges::telegram + +#endif // OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_BRIDGE_CONFIG_HPP_INCLUDED diff --git a/include/optionx_cpp/bridges/telegram/TelegramSignalParser.hpp b/include/optionx_cpp/bridges/telegram/TelegramSignalParser.hpp new file mode 100644 index 0000000..0d56b66 --- /dev/null +++ b/include/optionx_cpp/bridges/telegram/TelegramSignalParser.hpp @@ -0,0 +1,376 @@ +#pragma once +#ifndef OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_PARSER_HPP_INCLUDED +#define OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_PARSER_HPP_INCLUDED + +/// \file TelegramSignalParser.hpp +/// \brief Deterministic regex parser for Telegram signal and outcome messages. + +#include "bridges/telegram/TelegramParsedMessage.hpp" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace optionx::bridges::telegram { + + /// \struct TelegramSignalRule + /// \brief Capture-group mapping and defaults for one signal regex. + struct TelegramSignalRule { + std::string name; + std::string pattern; + std::size_t symbol_group = 1; + std::size_t direction_group = 2; + std::size_t expiry_group = 3; + std::size_t unit_group = 4; + OptionType option_type = OptionType::SPRINT; + }; + + /// \struct TelegramOutcomeRule + /// \brief Capture-group mapping for one outcome regex. + struct TelegramOutcomeRule { + std::string name; + std::string pattern; + std::size_t symbol_group = 1; + std::size_t result_group = 2; + }; + + /// \struct TelegramParserConfig + /// \brief Source-independent deterministic parser configuration. + struct TelegramParserConfig { + std::vector signal_rules; + std::vector outcome_rules; + bool use_chat_title_as_signal_name = true; + }; + + /// \class TelegramSignalParser + /// \brief Parses raw text while preserving outcomes and diagnostics. + class TelegramSignalParser final { + public: + /// \brief Returns conservative rules for common binary signal text. + static TelegramParserConfig default_config() { + TelegramParserConfig config; + config.signal_rules.push_back({ + "pair-direction-expiry", + R"(\b([A-Z]{3,6}(?:[/_-]?[A-Z]{3,6})?)\s+(BUY|SELL|CALL|PUT)\b(?:\s+(\d{1,5})\s*(s|sec|secs|m|min|mins|h|hr|hour|hours))?)", + 1, + 2, + 3, + 4, + OptionType::SPRINT, + }); + config.outcome_rules.push_back({ + "pair-outcome", + R"(\b([A-Z]{3,6}(?:[/_-]?[A-Z]{3,6})?)\s+(WIN|LOSS|REFUND|DRAW)\b)", + 1, + 2, + }); + return config; + } + + explicit TelegramSignalParser( + TelegramParserConfig config = default_config()) + : m_config(std::move(config)) { + for (const auto& rule : m_config.signal_rules) { + m_signal_rules.push_back({ + rule, + std::regex(rule.pattern, std::regex::ECMAScript | std::regex::icase), + }); + } + for (const auto& rule : m_config.outcome_rules) { + m_outcome_rules.push_back({ + rule, + std::regex(rule.pattern, std::regex::ECMAScript | std::regex::icase), + }); + } + } + + /// \brief Parses one message into signals, outcomes and diagnostics. + TelegramParsedMessage parse(const TelegramRawMessage& raw) const { + raw.validate(); + TelegramParsedMessage result; + result.raw = raw; + parse_signals(raw, result); + parse_outcomes(raw, result); + return result; + } + + private: + struct CompiledSignalRule { + TelegramSignalRule spec; + std::regex expression; + }; + + struct CompiledOutcomeRule { + TelegramOutcomeRule spec; + std::regex expression; + }; + + struct SignalCandidate { + std::size_t begin = 0; + std::size_t end = 0; + TelegramParsedSignal signal; + }; + + static std::string upper(std::string value) { + std::transform(value.begin(), value.end(), value.begin(), [](unsigned char ch) { + return static_cast(std::toupper(ch)); + }); + return value; + } + + static bool has_group( + const std::smatch& match, + const std::size_t group) { + return group != 0 && group < match.size() && match[group].matched; + } + + static std::uint32_t parse_duration( + const std::string& value, + const std::string& unit) { + std::uint64_t amount = 0; + try { + std::size_t consumed = 0; + amount = std::stoull(value, &consumed); + if (consumed != value.size() || amount == 0) { + throw std::invalid_argument("expiry amount is invalid"); + } + } + catch (const std::exception&) { + throw std::invalid_argument("expiry amount is invalid"); + } + + const auto normalized_unit = upper(unit); + std::uint64_t multiplier = 0; + if (normalized_unit == "S" || normalized_unit == "SEC" || + normalized_unit == "SECS") { + multiplier = 1; + } + else if (normalized_unit == "M" || normalized_unit == "MIN" || + normalized_unit == "MINS") { + multiplier = 60; + } + else if (normalized_unit == "H" || normalized_unit == "HR" || + normalized_unit == "HOUR" || normalized_unit == "HOURS") { + multiplier = 60 * 60; + } + else { + throw std::invalid_argument("expiry unit is unsupported"); + } + if (amount > (std::numeric_limits::max)() / multiplier) { + throw std::invalid_argument("expiry duration is too large"); + } + return static_cast(amount * multiplier); + } + + TelegramParsedSignal build_signal( + const TelegramRawMessage& raw, + const TelegramSignalRule& rule, + const std::smatch& match, + TelegramParsedMessage& result) const { + if (!has_group(match, rule.symbol_group) || + !has_group(match, rule.direction_group)) { + throw std::invalid_argument("signal rule has no symbol or direction"); + } + + TelegramParsedSignal signal; + signal.source_message_identity = raw.message_identity(); + signal.symbol = upper(match[rule.symbol_group].str()); + if (!optionx::to_enum(upper(match[rule.direction_group].str()), signal.order_type)) { + throw std::invalid_argument("signal direction is unsupported"); + } + signal.option_type = rule.option_type; + signal.signal_name = m_config.use_chat_title_as_signal_name + ? raw.chat_title + : rule.name; + signal.raw_text = match.str(); + + if (has_group(match, rule.expiry_group)) { + if (!has_group(match, rule.unit_group)) { + throw std::invalid_argument("signal expiry has no unit"); + } + signal.duration = parse_duration( + match[rule.expiry_group].str(), + match[rule.unit_group].str()); + } + return signal; + } + + static bool same_signal( + const TelegramParsedSignal& left, + const TelegramParsedSignal& right) { + return left.symbol == right.symbol && + left.order_type == right.order_type && + left.option_type == right.option_type && + left.duration == right.duration && + left.expiry_time == right.expiry_time && + left.signal_name == right.signal_name; + } + + void parse_signals( + const TelegramRawMessage& raw, + TelegramParsedMessage& result) const { + std::vector candidates; + for (const auto& compiled : m_signal_rules) { + for (std::sregex_iterator it(raw.text.begin(), raw.text.end(), compiled.expression); + it != std::sregex_iterator(); ++it) { + const auto& match = *it; + try { + candidates.push_back({ + static_cast(match.position()), + static_cast(match.position() + match.length()), + build_signal(raw, compiled.spec, match, result), + }); + } + catch (const std::exception& error) { + result.diagnostics.push_back({ + "invalid_signal_match", + error.what(), + static_cast(match.position()), + static_cast(match.length()), + }); + } + } + } + + std::sort(candidates.begin(), candidates.end(), [](const auto& left, const auto& right) { + if (left.begin != right.begin) { + return left.begin < right.begin; + } + return left.end < right.end; + }); + + std::vector rejected(candidates.size(), false); + bool unmatched_expiry = false; + static const std::regex trailing_expiry( + R"(^\s+(?:\d+\s*)?(?:s|sec|secs|m|min|mins|h|hr|hour|hours|d|day|days|w|week|weeks|tick|ticks)\b)", + std::regex::ECMAScript | std::regex::icase); + for (std::size_t i = 0; i < candidates.size(); ++i) { + if (candidates[i].signal.duration != 0) { + continue; + } + const auto tail = raw.text.substr( + candidates[i].end, + std::min(raw.text.size() - candidates[i].end, 48)); + if (std::regex_search(tail, trailing_expiry)) { + rejected[i] = true; + unmatched_expiry = true; + } + } + bool ambiguous = false; + for (std::size_t i = 0; i < candidates.size(); ++i) { + for (std::size_t j = i + 1; j < candidates.size(); ++j) { + if (candidates[j].begin >= candidates[i].end) { + break; + } + if (same_signal(candidates[i].signal, candidates[j].signal)) { + rejected[j] = true; + } + else { + rejected[i] = true; + rejected[j] = true; + ambiguous = true; + } + } + } + if (ambiguous) { + result.diagnostics.push_back({ + "ambiguous_overlapping_signal", + "overlapping signal rules produced different meanings", + 0, + raw.text.size(), + }); + } + for (std::size_t i = 0; i < candidates.size(); ++i) { + if (!rejected[i]) { + result.signals.push_back(std::move(candidates[i].signal)); + } + } + + if (unmatched_expiry) { + result.diagnostics.push_back({ + "unmatched_expiry", + "signal expiry was not recognized", + 0, + raw.text.size(), + }); + } + else if (result.signals.empty() && !raw.text.empty()) { + static const std::regex unsupported_expiry( + R"(\b(?:BUY|SELL|CALL|PUT)\b[^\r\n]{0,32}\b(?:\d+\s*)?(?:d|day|days|w|week|weeks|tick|ticks)\b)", + std::regex::ECMAScript | std::regex::icase); + std::smatch match; + if (std::regex_search(raw.text, match, unsupported_expiry)) { + result.diagnostics.push_back({ + "unmatched_expiry", + "signal expiry was not recognized", + static_cast(match.position()), + static_cast(match.length()), + }); + } + } + } + + static TelegramOutcomeResult outcome_from_token(const std::string& value) { + const auto token = upper(value); + if (token == "WIN") { + return TelegramOutcomeResult::WIN; + } + if (token == "LOSS") { + return TelegramOutcomeResult::LOSS; + } + if (token == "REFUND") { + return TelegramOutcomeResult::REFUND; + } + return TelegramOutcomeResult::UNKNOWN; + } + + void parse_outcomes( + const TelegramRawMessage& raw, + TelegramParsedMessage& result) const { + for (const auto& compiled : m_outcome_rules) { + for (std::sregex_iterator it(raw.text.begin(), raw.text.end(), compiled.expression); + it != std::sregex_iterator(); ++it) { + const auto& match = *it; + try { + if (!has_group(match, compiled.spec.result_group)) { + throw std::invalid_argument("outcome rule has no result"); + } + TelegramParsedOutcome outcome; + outcome.source_message_identity = raw.message_identity(); + outcome.reply_to_message_identity = raw.reply_to_message_identity(); + if (has_group(match, compiled.spec.symbol_group)) { + outcome.symbol = upper(match[compiled.spec.symbol_group].str()); + } + outcome.result = outcome_from_token( + match[compiled.spec.result_group].str()); + outcome.signal_name = raw.chat_title; + outcome.raw_text = match.str(); + result.outcomes.push_back(std::move(outcome)); + } + catch (const std::exception& error) { + result.diagnostics.push_back({ + "invalid_outcome_match", + error.what(), + static_cast(match.position()), + static_cast(match.length()), + }); + } + } + } + } + + TelegramParserConfig m_config; + std::vector m_signal_rules; + std::vector m_outcome_rules; + }; + +} // namespace optionx::bridges::telegram + +#endif // OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_SIGNAL_PARSER_HPP_INCLUDED diff --git a/include/optionx_cpp/bridges/telegram/TelegramWorkerMessageSource.hpp b/include/optionx_cpp/bridges/telegram/TelegramWorkerMessageSource.hpp new file mode 100644 index 0000000..c4eb5fa --- /dev/null +++ b/include/optionx_cpp/bridges/telegram/TelegramWorkerMessageSource.hpp @@ -0,0 +1,137 @@ +#pragma once +#ifndef OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_WORKER_MESSAGE_SOURCE_HPP_INCLUDED +#define OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_WORKER_MESSAGE_SOURCE_HPP_INCLUDED + +/// \file TelegramWorkerMessageSource.hpp +/// \brief Adapter from a tg-client-stdio-style worker client to Telegram bridge input. + +#include "bridges/telegram/TelegramSignalBridge.hpp" + +#include +#include +#include +#include + +namespace optionx::bridges::telegram { + + /// \struct TelegramWorkerSourceConfig + /// \brief Chat and topic selection for one worker live listener. + struct TelegramWorkerSourceConfig { + std::vector chats; + std::vector topic_ids; + }; + + /// \class TelegramWorkerMessageSource + /// \brief Binds a WorkerClient-like object to TelegramMessageSource. + /// + /// The worker type is a template so this header does not force OptionX to + /// include or link a particular worker repository. It is compatible with + /// `tg_client_stdio::WorkerClient` and deterministic test doubles that + /// provide the same `start_listening` and `stop_listening` operations. + template + class TelegramWorkerMessageSource final : public TelegramMessageSource { + public: + TelegramWorkerMessageSource( + WorkerClient& worker, + TelegramWorkerSourceConfig config) + : m_worker(worker), + m_config(std::move(config)) {} + + bool start(message_callback_t on_message, + error_callback_t on_error) override { + if (m_started || m_config.chats.empty() || !on_message) { + return false; + } + m_on_message = std::move(on_message); + m_on_error = std::move(on_error); + try { + if (!m_worker.start_listening( + m_config.chats, + [this](const auto& record) { handle_record(record); }, + m_config.topic_ids)) { + clear_callbacks(); + return false; + } + m_started = true; + return true; + } + catch (const std::exception& error) { + report_error(error.what()); + clear_callbacks(); + return false; + } + catch (...) { + report_error("Telegram worker listener failed to start."); + clear_callbacks(); + return false; + } + } + + void stop() noexcept override { + if (!m_started) { + clear_callbacks(); + return; + } + try { + (void)m_worker.stop_listening(); + } + catch (const std::exception& error) { + report_error(error.what()); + } + catch (...) { + report_error("Telegram worker listener failed to stop."); + } + m_started = false; + clear_callbacks(); + } + + private: + void handle_record(const nlohmann::json& record) { + if (record.value("message_type", "") == "error") { + const auto payload = record.value("payload", nlohmann::json::object()); + report_error(payload.value("message", "Telegram worker live error.")); + return; + } + if (record.value("operation", "") != "message.received") { + return; + } + try { + const auto payload = record.at("payload"); + const auto raw = TelegramRawMessage::from_json(payload.at("message")); + if (m_on_message) { + m_on_message(raw); + } + } + catch (const std::exception& error) { + report_error(error.what()); + } + catch (...) { + report_error("Telegram worker message record was invalid."); + } + } + + void report_error(const std::string& message) { + if (m_on_error) { + try { + m_on_error(message); + } + catch (...) { + } + } + } + + void clear_callbacks() noexcept { + m_on_message = {}; + m_on_error = {}; + } + + WorkerClient& m_worker; + TelegramWorkerSourceConfig m_config; + message_callback_t m_on_message; + error_callback_t m_on_error; + bool m_started = false; + }; + +} // namespace optionx::bridges::telegram + +#endif // OPTIONX_HEADER_BRIDGES_TELEGRAM_TELEGRAM_WORKER_MESSAGE_SOURCE_HPP_INCLUDED diff --git a/include/optionx_cpp/data/trading/enums.hpp b/include/optionx_cpp/data/trading/enums.hpp index 04a27c7..fa5b51c 100644 --- a/include/optionx_cpp/data/trading/enums.hpp +++ b/include/optionx_cpp/data/trading/enums.hpp @@ -100,7 +100,8 @@ namespace optionx { METATRADER_FILE_TRANSPORT, ///< MetaTrader common-files JSON-RPC bridge transport. BRIDGE_PROTOCOL_V1_HTTP_WEBSOCKET, ///< Bridge Protocol v1 HTTP/WebSocket server. BRIDGE_PROTOCOL_V1_NAMED_PIPE, ///< Bridge Protocol v1 named-pipe server. - BOT_BINARY ///< BotBinary/BinaryBot compatibility bridge. + BOT_BINARY, ///< BotBinary/BinaryBot compatibility bridge. + TELEGRAM_SIGNAL ///< Telegram user-client signal bridge. }; /// \brief Converts BridgeType to its string representation. @@ -115,7 +116,8 @@ namespace optionx { "METATRADER_FILE_TRANSPORT", "BRIDGE_PROTOCOL_V1_HTTP_WEBSOCKET", "BRIDGE_PROTOCOL_V1_NAMED_PIPE", - "BOT_BINARY" + "BOT_BINARY", + "TELEGRAM_SIGNAL" }; return utils::enum_string_or_unknown(str_data, static_cast(value)); } @@ -133,6 +135,7 @@ namespace optionx { {"BRIDGE_PROTOCOL_V1_HTTP_WEBSOCKET", BridgeType::BRIDGE_PROTOCOL_V1_HTTP_WEBSOCKET}, {"BRIDGE_PROTOCOL_V1_NAMED_PIPE", BridgeType::BRIDGE_PROTOCOL_V1_NAMED_PIPE}, {"BOT_BINARY", BridgeType::BOT_BINARY}, + {"TELEGRAM_SIGNAL", BridgeType::TELEGRAM_SIGNAL}, {"BINARYBOT", BridgeType::BOT_BINARY}, {"BOTBINARY", BridgeType::BOT_BINARY} }; diff --git a/tests/telegram_signal_bridge_test.cpp b/tests/telegram_signal_bridge_test.cpp new file mode 100644 index 0000000..aa3b2e5 --- /dev/null +++ b/tests/telegram_signal_bridge_test.cpp @@ -0,0 +1,108 @@ +#include + +#include "optionx_cpp/bridges/telegram.hpp" + +#include +#include +#include + +namespace { + +class FakeMessageSource final + : public optionx::bridges::telegram::TelegramMessageSource { +public: + bool start(message_callback_t on_message, error_callback_t on_error) override { + (void)on_error; + m_on_message = std::move(on_message); + started = true; + return true; + } + + void stop() noexcept override { + stopped = true; + m_on_message = {}; + } + + void emit(optionx::bridges::telegram::TelegramRawMessage message) { + if (m_on_message) { + m_on_message(message); + } + } + + bool started = false; + bool stopped = false; + +private: + message_callback_t m_on_message; +}; + +optionx::bridges::telegram::TelegramRawMessage make_message() { + optionx::bridges::telegram::TelegramRawMessage message; + message.chat_id = "-10042"; + message.chat_title = "Signals"; + message.message_id = 123; + message.date_ms = 1800000000000; + message.text = "EURUSD BUY 5m"; + return message; +} + +std::unique_ptr config() { + auto value = std::make_unique< + optionx::bridges::telegram::TelegramSignalBridgeConfig>(); + value->bridge_id = 17; + value->fixed_amount = 1.0; + return value; +} + +} // namespace + +TEST(TelegramSignalBridge, PublishesParsedSignalAndReportsDuplicate) { + auto source = std::make_shared(); + optionx::bridges::telegram::TelegramSignalBridge bridge(source); + ASSERT_TRUE(bridge.configure(config())); + + std::vector> signals; + std::vector reports; + std::int64_t next_signal_id = 100; + bridge.on_signal_id() = [&] { return ++next_signal_id; }; + bridge.on_trade_signal() = [&](std::unique_ptr signal) { + signals.push_back(std::move(signal)); + }; + bridge.on_signal_report() = [&](const auto& report) { + reports.push_back(report); + }; + + bridge.run(); + ASSERT_TRUE(source->started); + source->emit(make_message()); + source->emit(make_message()); + + ASSERT_EQ(signals.size(), 1u); + EXPECT_EQ(signals.front()->signal_id, 101); + EXPECT_EQ(signals.front()->bridge_id, 17); + EXPECT_EQ(signals.front()->symbol, "EURUSD"); + EXPECT_EQ(signals.front()->duration, 300u); + ASSERT_EQ(reports.size(), 1u); + EXPECT_EQ(reports.front().status, + optionx::BridgeSignalReportStatus::DUPLICATE); + EXPECT_EQ(reports.front().reason_code, "duplicate_message"); + + bridge.shutdown(); + EXPECT_TRUE(source->stopped); +} + +TEST(TelegramSignalBridge, RejectsInvalidConfigurationBeforeStartingSource) { + auto source = std::make_shared(); + optionx::bridges::telegram::TelegramSignalBridge bridge(source); + auto invalid = config(); + invalid->fixed_amount = 0.0; + + EXPECT_FALSE(bridge.configure(std::move(invalid))); + bridge.run(); + EXPECT_FALSE(source->started); +} + +int main(int argc, char** argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +} diff --git a/tests/telegram_signal_parser_test.cpp b/tests/telegram_signal_parser_test.cpp new file mode 100644 index 0000000..5261a39 --- /dev/null +++ b/tests/telegram_signal_parser_test.cpp @@ -0,0 +1,71 @@ +#include + +#include "optionx_cpp/bridges/telegram.hpp" + +namespace { + +optionx::bridges::telegram::TelegramRawMessage message_with_text(std::string text) { + optionx::bridges::telegram::TelegramRawMessage message; + message.chat_id = "-10042"; + message.chat_title = "Morning signals"; + message.message_id = 123; + message.date_ms = 1800000000000; + message.text = std::move(text); + return message; +} + +} // namespace + +TEST(TelegramSignalParser, ParsesExecutableSignalAndNormalizesExpiry) { + const auto parsed = optionx::bridges::telegram::TelegramSignalParser().parse( + message_with_text("EURUSD BUY 5m")); + + ASSERT_EQ(parsed.signals.size(), 1u); + EXPECT_EQ(parsed.signals[0].symbol, "EURUSD"); + EXPECT_EQ(parsed.signals[0].order_type, optionx::OrderType::BUY); + EXPECT_EQ(parsed.signals[0].option_type, optionx::OptionType::SPRINT); + EXPECT_EQ(parsed.signals[0].duration, 300u); + EXPECT_EQ(parsed.signals[0].source_message_identity, "telegram:-10042:0:123"); + EXPECT_TRUE(parsed.outcomes.empty()); +} + +TEST(TelegramSignalParser, ParsesMultipleSignalsAndOutcomeSeparately) { + const auto parsed = optionx::bridges::telegram::TelegramSignalParser().parse( + message_with_text("EURUSD BUY 5m\nGBPUSD SELL 10m\nEURUSD WIN")); + + ASSERT_EQ(parsed.signals.size(), 2u); + EXPECT_EQ(parsed.signals[0].symbol, "EURUSD"); + EXPECT_EQ(parsed.signals[1].symbol, "GBPUSD"); + ASSERT_EQ(parsed.outcomes.size(), 1u); + EXPECT_EQ(parsed.outcomes[0].symbol, "EURUSD"); + EXPECT_EQ(parsed.outcomes[0].result, + optionx::bridges::telegram::TelegramOutcomeResult::WIN); +} + +TEST(TelegramSignalParser, RejectsUnsupportedExpiryUnitFailClosed) { + const auto parsed = optionx::bridges::telegram::TelegramSignalParser().parse( + message_with_text("USDJPY BUY 2 weeks")); + + EXPECT_TRUE(parsed.signals.empty()); + ASSERT_EQ(parsed.diagnostics.size(), 1u); + EXPECT_EQ(parsed.diagnostics.front().code, "unmatched_expiry"); +} + +TEST(TelegramSignalParser, RejectsConflictingOverlappingRules) { + auto config = optionx::bridges::telegram::TelegramParserConfig{}; + config.signal_rules = { + {"sprint", R"((EURUSD)\s+(BUY))", 1, 2, 0, 0, optionx::OptionType::SPRINT}, + {"classic", R"((EURUSD)\s+(BUY))", 1, 2, 0, 0, optionx::OptionType::CLASSIC}, + }; + optionx::bridges::telegram::TelegramSignalParser parser(std::move(config)); + const auto parsed = parser.parse(message_with_text("EURUSD BUY")); + + EXPECT_TRUE(parsed.signals.empty()); + ASSERT_EQ(parsed.diagnostics.size(), 1u); + EXPECT_EQ(parsed.diagnostics.front().code, "ambiguous_overlapping_signal"); +} + +int main(int argc, char** argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +} diff --git a/tests/telegram_worker_source_test.cpp b/tests/telegram_worker_source_test.cpp new file mode 100644 index 0000000..c422c5f --- /dev/null +++ b/tests/telegram_worker_source_test.cpp @@ -0,0 +1,104 @@ +#include + +#include "optionx_cpp/bridges/telegram.hpp" + +#include +#include +#include + +namespace { + +class FakeWorkerClient { +public: + using handler_t = std::function; + + bool start_listening( + const std::vector& chats, + handler_t handler, + const std::vector& topics) { + selected_chats = chats; + selected_topics = topics; + m_handler = std::move(handler); + return true; + } + + bool stop_listening() { + stopped = true; + m_handler = {}; + return true; + } + + void emit_message() { + m_handler(nlohmann::json{ + {"message_type", "event"}, + {"operation", "message.received"}, + {"request_id", 0}, + {"payload", { + {"message", { + {"chat_id", "-10042"}, + {"chat_title", "Signals"}, + {"topic_id", "7"}, + {"message_id", 12}, + {"date_ms", 1800000000000LL}, + {"text", "EURUSD BUY 5m"}, + {"media", nlohmann::json::array()}, + }} + }} + }); + } + + void emit_error() { + m_handler(nlohmann::json{ + {"message_type", "error"}, + {"request_id", 0}, + {"payload", {{"message", "listener failed"}}}, + }); + } + + std::vector selected_chats; + std::vector selected_topics; + bool stopped = false; + +private: + handler_t m_handler; +}; + +} // namespace + +TEST(TelegramWorkerMessageSource, AdaptsLiveRecordsAndErrors) { + FakeWorkerClient worker; + optionx::bridges::telegram::TelegramWorkerSourceConfig config; + config.chats = {"-10042"}; + config.topic_ids = {"7"}; + optionx::bridges::telegram::TelegramWorkerMessageSource source(worker, config); + + std::vector messages; + std::vector errors; + ASSERT_TRUE(source.start( + [&](const auto& message) { messages.push_back(message); }, + [&](const auto& error) { errors.push_back(error); })); + worker.emit_message(); + worker.emit_error(); + + ASSERT_EQ(messages.size(), 1u); + EXPECT_EQ(messages.front().message_identity(), "telegram:-10042:7:12"); + ASSERT_EQ(errors.size(), 1u); + EXPECT_EQ(errors.front(), "listener failed"); + EXPECT_EQ(worker.selected_chats, config.chats); + EXPECT_EQ(worker.selected_topics, config.topic_ids); + + source.stop(); + EXPECT_TRUE(worker.stopped); +} + +TEST(TelegramWorkerMessageSource, RejectsEmptyChatSelection) { + FakeWorkerClient worker; + optionx::bridges::telegram::TelegramWorkerMessageSource source( + worker, {}); + EXPECT_FALSE(source.start([](const auto&) {}, [](const auto&) {})); +} + +int main(int argc, char** argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +}