Skip to content
5 changes: 5 additions & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,7 @@ set(POLYMARKET_CLIENT_SOURCES
src/websocket_client_transport.cpp
src/websocket_resilience.cpp
src/websocket_market_data.cpp
src/websocket_user_data.cpp
src/market_fetcher.cpp
src/market_discovery.cpp
src/market_time.cpp
Expand All @@ -132,6 +133,10 @@ set(POLYMARKET_CLIENT_SOURCES
src/orderbook_runtime_stream.cpp
src/orderbook_subscription.cpp
src/orderbook_stream.cpp
src/user_stream.cpp
src/user_stream_runtime.cpp
src/user_stream_runtime_events.cpp
src/user_stream_subscription.cpp
src/order_amounts.cpp
src/clob_order_execution.cpp
src/clob_order_submission.cpp
Expand Down
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ Reusable C++20 client for Polymarket: REST, WebSocket streaming, and order signi
- **REST**: market discovery, orderbook/price queries, auth key management, and trading endpoints.
- **Transport controls**: configurable libcurl timeouts, keepalive, connection reuse, proxy/user-agent, request metrics, and cumulative stats.
- **WebSocket**: orderbook streaming via IXWebSocket with reconnect, subscription replay, typed callbacks, and backpressure counters.
- **User stream**: authenticated CLOB user channel (`UserStream`) with typed order/trade events, per-market or all-market subscriptions, plus gap callbacks (invalidate state) and recovery callbacks (reconcile via REST once the subscription is restored).
- **Signing**: CLOB V2 EIP-712 order signing (secp256k1, keccak).
- **Decimal math**: shared scaled-integer conversion for trading amounts.
- **Structured errors**: opt-in `Result<T>` APIs with typed SDK error classification.
Expand Down Expand Up @@ -135,6 +136,7 @@ int main() {
- `rest_example`: fetch markets from CLOB REST
- `sign_example`: sign a dummy order (requires `PRIVATE_KEY`)
- `ws_example`: connect to Polymarket WS and subscribe to orderbook agg
- `user_stream_example`: stream your own order and trade events (requires `PRIVATE_KEY`; optional condition IDs as arguments)
- `uma_oracle_watch`: stream UMA adapter lifecycle events over Polygon JSON-RPC
- `condition_resolution_watch`: stream Conditional Tokens resolution/redemption events
- `evm_event_indexer_example`: persistent HTTP catch-up + live WS indexer with a cursor file
Expand Down
3 changes: 3 additions & 0 deletions cmake/PolymarketExamples.cmake
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,9 @@ if(POLYMARKET_CLIENT_BUILD_EXAMPLES)
add_executable(ws_example examples/ws_example.cpp)
target_link_libraries(ws_example PRIVATE polymarket::client)

add_executable(user_stream_example examples/user_stream_example.cpp)
target_link_libraries(user_stream_example PRIVATE polymarket::client)

add_executable(uma_oracle_watch examples/uma_oracle_watch.cpp)
target_link_libraries(uma_oracle_watch PRIVATE polymarket::client)

Expand Down
10 changes: 10 additions & 0 deletions cmake/PolymarketTests.cmake
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,16 @@ if(POLYMARKET_CLIENT_BUILD_TESTS)
target_link_libraries(test_orderbook_stream PRIVATE polymarket::client)
add_test(NAME test_orderbook_stream COMMAND test_orderbook_stream)

add_executable(test_user_stream
tests/test_user_stream.cpp
tests/test_user_stream_connection.cpp
tests/test_user_stream_failures.cpp)
target_include_directories(test_user_stream PRIVATE
${CMAKE_CURRENT_SOURCE_DIR}/src
${CMAKE_CURRENT_SOURCE_DIR}/tests)
target_link_libraries(test_user_stream PRIVATE polymarket::client)
add_test(NAME test_user_stream COMMAND test_user_stream)

add_executable(test_orderbook_owner_reset tests/test_orderbook_owner_reset.cpp)
target_include_directories(test_orderbook_owner_reset PRIVATE
${CMAKE_CURRENT_SOURCE_DIR}/tests)
Expand Down
88 changes: 88 additions & 0 deletions examples/user_stream_example.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
#include "clob_client.hpp"
#include "user_stream.hpp"

#include <chrono>
#include <cstdlib>
#include <iostream>
#include <thread>

int main(int argc, char **argv)
{
using namespace polymarket;

const char *pk_env = std::getenv("PRIVATE_KEY");
const char *key_env = std::getenv("POLY_API_KEY");
const char *secret_env = std::getenv("POLY_API_SECRET");
const char *passphrase_env = std::getenv("POLY_API_PASSPHRASE");
const bool has_api_credentials = key_env && secret_env && passphrase_env;
if (!pk_env && !has_api_credentials)
{
std::cout << "PRIVATE_KEY or POLY_API_KEY/POLY_API_SECRET/POLY_API_PASSPHRASE "
"not set; skipping user stream example.\n";
return 0;
}

try
{
ApiCredentials credentials;
if (has_api_credentials)
{
credentials = {key_env, secret_env, passphrase_env};
}
else
{
// Derive L2 API credentials from the signer (no orders are placed).
ClobClient client("https://clob.polymarket.com", 137, pk_env);
credentials = client.create_or_derive_api_key();
}

UserStream stream(Config{}, credentials);
stream.on_order([](const UserOrderEvent &order)
{ std::cout << "[order] " << order.type << ' ' << order.side << ' '
<< order.size_matched << '/' << order.original_size
<< " @ " << order.price << " id=" << order.id << '\n'; });
stream.on_trade([](const UserTradeEvent &trade)
{ std::cout << "[trade] " << trade.status << ' ' << trade.side << ' '
<< trade.size << " @ " << trade.price << " id=" << trade.id << '\n'; });
stream.on_error([](const std::string &error)
{ std::cerr << "[error] " << error << '\n'; });
stream.on_stream_gap([]
{ std::cout << "[gap] events may have been missed; local state is stale\n"; });
stream.on_stream_recovered([]
{ std::cout << "[recovered] subscription restored; reconcile via REST\n"; });

// Optional condition IDs narrow the stream; none means all markets.
if (argc > 1)
{
for (int i = 1; i < argc; ++i)
stream.subscribe(argv[i]);
}
else
{
stream.subscribe_all_markets();
}

if (!stream.connect())
{
std::cerr << "failed to connect user stream\n";
return 1;
}

const char *seconds_env = std::getenv("USER_STREAM_SECONDS");
const auto deadline = std::chrono::steady_clock::now() +
std::chrono::seconds(seconds_env ? std::atoi(seconds_env) : 60);
while (std::chrono::steady_clock::now() < deadline && !stream.authentication_failed())
std::this_thread::sleep_for(std::chrono::milliseconds(100));
stream.stop();
if (stream.authentication_failed())
return 1;
std::cout << "orders=" << stream.order_events()
<< " trades=" << stream.trade_events() << '\n';
}
catch (const std::exception &e)
{
std::cerr << "Error: " << e.what() << "\n";
return 1;
}
return 0;
}
1 change: 1 addition & 0 deletions include/types.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -262,6 +262,7 @@ namespace polymarket
std::string clob_ws_url = "wss://ws-subscriptions-clob.polymarket.com/ws/market";
std::string gamma_api_url = "https://gamma-api.polymarket.com";
std::string rtds_ws_url = "wss://ws-live-data.polymarket.com";
std::string clob_user_ws_url = "wss://ws-subscriptions-clob.polymarket.com/ws/user";

// Trading parameters
double trigger_combined = 0.98;
Expand Down
137 changes: 137 additions & 0 deletions include/user_stream.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
#pragma once

#include "clob_types.hpp"
#include "order_signer.hpp"
#include "types.hpp"

#include <cstdint>
#include <functional>
#include <memory>
#include <optional>
#include <string>
#include <vector>

namespace polymarket
{
namespace detail
{
class UserStreamRuntime;
}

// Order lifecycle event from the authenticated CLOB user channel.
struct UserOrderEvent
{
std::string id;
std::string owner;
std::string market; // condition ID
std::string asset_id; // token ID
std::string side; // BUY or SELL
std::string original_size;
std::string size_matched;
std::string price;
std::string type; // PLACEMENT, UPDATE, or CANCELLATION
std::string status;
std::string order_type;
std::string maker_address;
std::string order_owner;
std::vector<std::string> associate_trades;
std::string outcome;
std::string created_at;
std::string expiration;
std::string timestamp;
};

// Trade lifecycle event (MATCHED, MINED, CONFIRMED, RETRYING, FAILED).
struct UserTradeEvent
{
std::string id;
std::string taker_order_id;
std::string market; // condition ID
std::string asset_id; // token ID
std::string side;
std::string size;
std::string price;
std::string status;
std::string owner;
std::string fee_rate_bps;
std::string match_time;
std::string last_update;
std::string timestamp;
std::string trade_owner;
std::string maker_address;
std::string transaction_hash;
std::optional<uint32_t> bucket_index;
std::vector<MakerOrder> maker_orders;
std::string trader_side; // TAKER or MAKER
std::string outcome;
};

using UserOrderCallback = std::function<void(const UserOrderEvent &order)>;
using UserTradeCallback = std::function<void(const UserTradeEvent &trade)>;
// Fired immediately whenever events may have been missed (every connect,
// disconnect, queue overflow, or invalid payload). Treat local order and
// trade state as stale; do not reconcile here, because on reconnect it
// runs before the subscription is restored and the server does not replay
// events missed in between.
using UserStreamGapCallback = std::function<void()>;
// Fired once the authenticated subscription has been resent after a gap.
// Reconcile order and trade state via REST (`ClobClient::get_open_orders`,
// `get_trades`) here and merge it with events delivered from this point.
using UserStreamRecoveredCallback = std::function<void()>;
// Fired when the server rejects the session (close code 1008, e.g. invalid
// API credentials). The stream stops reconnecting and run() returns; call
// connect() to retry.
using UserStreamErrorCallback = std::function<void(const std::string &error)>;

// Authenticated user-channel stream. Subscriptions are replayed with the
// API credentials after every reconnect.
class UserStream
{
public:
UserStream(const Config &config, const ApiCredentials &credentials);
~UserStream();

UserStream(const UserStream &) = delete;
UserStream &operator=(const UserStream &) = delete;

// Receive events for every market the API key trades.
void subscribe_all_markets();
// Receive events only for these condition IDs. Ignored for markets
// already covered by subscribe_all_markets().
void subscribe(const std::vector<std::string> &condition_ids);
void subscribe(const std::string &condition_id);
void unsubscribe(const std::vector<std::string> &condition_ids);
void unsubscribe(const std::string &condition_id);
void unsubscribe_all();

bool is_subscribed_to_all_markets() const;
std::vector<std::string> subscribed_markets() const;

// Callbacks
void on_order(UserOrderCallback callback);
void on_trade(UserTradeCallback callback);
void on_stream_gap(UserStreamGapCallback callback);
void on_stream_recovered(UserStreamRecoveredCallback callback);
void on_error(UserStreamErrorCallback callback);

// Connection
bool connect();
void disconnect();
bool is_connected() const;
bool authentication_failed() const;

// Run event loop (blocking)
void run();

// Stop
void stop();

// Statistics
uint64_t order_events() const;
uint64_t trade_events() const;

private:
std::shared_ptr<detail::UserStreamRuntime> runtime_;
};

} // namespace polymarket
2 changes: 2 additions & 0 deletions include/websocket_client.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ namespace polymarket
using OnMessageCallback = std::function<void(const std::string &)>;
using OnConnectCallback = std::function<void()>;
using OnDisconnectCallback = std::function<void()>;
using OnCloseCallback = std::function<void(uint16_t code, const std::string &reason)>;
using OnErrorCallback = std::function<void(const std::string &)>;
using OnSequencedMessageCallback = std::function<void(const std::string &, uint64_t)>;
using OnStreamGapCallback = std::function<void(uint64_t)>;
Expand Down Expand Up @@ -92,6 +93,7 @@ namespace polymarket
void on_typed_message(OnTypedMessageCallback callback);
void on_connect(OnConnectCallback callback);
void on_disconnect(OnDisconnectCallback callback);
void on_close(OnCloseCallback callback);
void on_error(OnErrorCallback callback);
void on_stream_gap(OnStreamGapCallback callback);

Expand Down
Loading
Loading