Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
75 changes: 75 additions & 0 deletions examples/example-arbfeed.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
import os
import time
import logging
from pythonjsonlogger import jsonlogger

from flumine import Flumine, clients
from flumine.events.events import TerminationEvent
from flumine.worker import BackgroundWorker
from strategies.arbsignal import ArbSignalStrategy
from workers.arbfeed import poll_arb_feed
from workers.arbscanner import ArbScanner

"""
External signal feed polled by a flumine worker.

Runs offline by default: without X402_WALLET_KEY the scanner returns
bundled SAMPLE data (not real market data) and makes no network calls.
Setting X402_WALLET_KEY opts in to the paid live feed (x402, USDC on
Base, capped by max_usd_per_call per request, eth-account required).

cd examples && python example-arbfeed.py

To gate real markets swap the SimulatedClient for a BetfairClient and
subscribe the strategy to a stream, see examples/tennisexample.py
"""

logger = logging.getLogger()

custom_format = "%(asctime) %(levelname) %(message)"
log_handler = logging.StreamHandler()
formatter = jsonlogger.JsonFormatter(custom_format)
formatter.converter = time.gmtime
log_handler.setFormatter(formatter)
logger.addHandler(log_handler)
logger.setLevel(logging.INFO)

scanner = ArbScanner(max_usd_per_call=0.02)
if scanner.demo:
logger.info("Arb feed in demo mode, using SAMPLE data (no network calls)")
else:
logger.warning("Arb feed in LIVE mode, each poll pays up to $0.02 USDC")

client = clients.SimulatedClient()
framework = Flumine(client=client)

strategy = ArbSignalStrategy(
name="arbsignal",
context={"min_net_yield_c": 1.0, "max_signal_age": 120},
)
framework.add_strategy(strategy)

framework.add_worker(
BackgroundWorker(
framework,
poll_arb_feed,
func_kwargs={"scanner": scanner, "limit": 10, "mode": "opportunities"},
interval=60 if scanner.demo else int(os.environ.get("ARB_POLL_INTERVAL", 300)),
)
)


def stop(context: dict, flumine) -> None:
# demo only: terminate once the feed has been processed
for s in flumine.strategies:
for signal in s.active_signals():
logger.info(
"Active signal",
extra={"pair": signal.get("event") or signal.get("pair")},
)
flumine.handler_queue.put(TerminationEvent(flumine))


framework.add_worker(BackgroundWorker(framework, stop, interval=None, start_delay=3))

framework.run()
178 changes: 178 additions & 0 deletions examples/resources/kalshi_predictit_arb_sample.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,178 @@
{
"_notice": "SAMPLE DATA for offline demo/tests only. Fictional pairs and prices, NOT real market data.",
"q": "",
"mode": "all",
"opportunities": [
{
"event": "SAMPLE State XX Senate winner? - Democratic",
"match_score": 0.74,
"kalshi": {
"ticker": "SAMPLE-SENATEXX-26-D",
"yes_bid_c": 59,
"yes_ask_c": 60
},
"predictit": {
"market": "SAMPLE Which party will win the XX Senate race?",
"contract": "Democratic",
"yes_ask_c": 68,
"no_ask_c": 31,
"position_limit_usd": 850
},
"best_direction": {
"direction": "YES_KALSHI__NO_PREDICTIT",
"kalshi_price_c": 60,
"predictit_price_c": 31,
"stake_c": 91,
"raw_spread_c": 9,
"net_if_yes_wins_c": 3.1,
"net_if_no_wins_c": 5.2,
"net_yield_c": 3.1,
"net_yield_pct": 3.4,
"fees": {
"kalshi_taker_c": 2,
"predictit_profit_fee_c": 6.9,
"predictit_withdrawal_drag_c": 3.45
}
},
"days_to_settlement": 120,
"roc_annualized_pct": 10.3,
"executable": true,
"liquidity": {
"vwap_buy_c": { "100": 60, "500": 60, "1000": 61 },
"book_depth_usd": 12000.0
},
"slippage_ok_100": true
},
{
"event": "SAMPLE District YY-01 House winner? - Republican",
"match_score": 0.69,
"kalshi": {
"ticker": "SAMPLE-HOUSEYY01-26-R",
"yes_bid_c": 44,
"yes_ask_c": 45
},
"predictit": {
"market": "SAMPLE YY-01 House race",
"contract": "Republican",
"yes_ask_c": 47,
"no_ask_c": 47,
"position_limit_usd": 850
},
"best_direction": {
"direction": "YES_KALSHI__NO_PREDICTIT",
"kalshi_price_c": 45,
"predictit_price_c": 47,
"stake_c": 92,
"raw_spread_c": 8,
"net_if_yes_wins_c": 1.2,
"net_if_no_wins_c": 2.4,
"net_yield_c": 1.2,
"net_yield_pct": 1.3,
"fees": {
"kalshi_taker_c": 2,
"predictit_profit_fee_c": 5.3,
"predictit_withdrawal_drag_c": 2.65
}
},
"days_to_settlement": 400,
"roc_annualized_pct": 1.2,
"executable": true,
"liquidity": {
"vwap_buy_c": { "100": 45, "500": 46, "1000": 46 },
"book_depth_usd": 4000.0
},
"slippage_ok_100": true
},
{
"event": "SAMPLE State ZZ governor party winner? - Democratic",
"match_score": 0.66,
"kalshi": {
"ticker": "SAMPLE-GOVPARTYZZ-26-D",
"yes_bid_c": 38,
"yes_ask_c": 39
},
"predictit": {
"market": "SAMPLE ZZ Governor race",
"contract": "Democratic",
"yes_ask_c": 41,
"no_ask_c": 41,
"position_limit_usd": 850
},
"best_direction": {
"direction": "YES_KALSHI__NO_PREDICTIT",
"kalshi_price_c": 39,
"predictit_price_c": 41,
"stake_c": 80,
"raw_spread_c": 20,
"net_if_yes_wins_c": 0.4,
"net_if_no_wins_c": 1.1,
"net_yield_c": 0.4,
"net_yield_pct": 0.5,
"fees": {
"kalshi_taker_c": 2,
"predictit_profit_fee_c": 5.9,
"predictit_withdrawal_drag_c": 2.95
}
},
"days_to_settlement": 400,
"roc_annualized_pct": 0.5,
"executable": false,
"liquidity": {
"vwap_buy_c": { "100": 39, "500": 40, "1000": 40 },
"book_depth_usd": 2500.0
},
"slippage_ok_100": true
},
{
"event": "SAMPLE State WW Senate winner? - Republican",
"match_score": 0.71,
"kalshi": {
"ticker": "SAMPLE-SENATEWW-26-R",
"yes_bid_c": 80,
"yes_ask_c": 81
},
"predictit": {
"market": "SAMPLE WW Senate race",
"contract": "Republican",
"yes_ask_c": 83,
"no_ask_c": 10,
"position_limit_usd": 850
},
"best_direction": {
"direction": "YES_KALSHI__NO_PREDICTIT",
"kalshi_price_c": 81,
"predictit_price_c": 10,
"stake_c": 91,
"raw_spread_c": 9,
"net_if_yes_wins_c": -2.0,
"net_if_no_wins_c": 4.3,
"net_yield_c": -2.0,
"net_yield_pct": -2.2,
"fees": {
"kalshi_taker_c": 2,
"predictit_profit_fee_c": 8.1,
"predictit_withdrawal_drag_c": 4.05
}
},
"days_to_settlement": 400,
"roc_annualized_pct": -2.0,
"executable": false,
"liquidity": {
"vwap_buy_c": { "100": 81, "500": 82, "1000": 82 },
"book_depth_usd": 2500.0
},
"slippage_ok_100": false
},
{
"event": "SAMPLE record with missing fields (client must tolerate this)"
}
],
"stats": {
"pairs_evaluated": 5,
"executable": 2,
"kalshi_markets_scanned": 5,
"predictit_contracts_scanned": 5,
"duration_ms": 0
},
"fetched_at": "2026-01-01T00:00:00.000Z"
}
63 changes: 63 additions & 0 deletions examples/strategies/arbsignal.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
import logging
import time

from flumine import BaseStrategy

logger = logging.getLogger(__name__)


class ArbSignalStrategy(BaseStrategy):
"""
Example strategy gated on an external signal feed, signals
are stored in context["arb_signals"] by the callback in
examples/workers/arbfeed.py
- Only processes OPEN markets whilst a fresh executable
signal >= context["min_net_yield_c"] exists
- Logs each new signal, no orders are placed
"""

def check_market_book(self, market, market_book):
if market_book.status == "OPEN" and self.active_signals():
return True

def process_market_book(self, market, market_book):
seen = market.context.setdefault("arb_signals_seen", set())
for signal in self.active_signals():
key = _label(signal)
if key not in seen:
seen.add(key)
logger.info(
"Arb signal active",
extra={
"market_id": market.market_id,
"pair": key,
"net_yield_c": _net_yield_c(signal),
},
)

def active_signals(self) -> list:
updated = self.context.get("arb_signals_updated")
max_age = self.context.get("max_signal_age", 120)
if updated is None or time.time() - updated > max_age:
return []
min_yield = self.context.get("min_net_yield_c", 1.0)
return [
s
for s in self.context.get("arb_signals", [])
if _label(s)
and s.get("executable") is True
and _net_yield_c(s) >= min_yield
]


def _label(signal: dict) -> str:
# live feed rows carry "event", older rows "pair", rows with neither are skipped
return str(signal.get("event") or signal.get("pair") or "").strip()


def _net_yield_c(signal: dict) -> float:
# tolerate missing fields, feed schema is not guaranteed
try:
return float(signal["best_direction"]["net_yield_c"])
except (KeyError, TypeError, ValueError):
return 0.0
60 changes: 60 additions & 0 deletions examples/workers/arbfeed.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
import logging
import time

import requests

from flumine.events.events import CustomEvent
from .arbscanner import ArbScanner, ArbScannerError, is_executable

logger = logging.getLogger(__name__)

"""
Worker polling an external signal feed (here a Kalshi <-> PredictIt
arbitrage scanner, see examples/workers/arbscanner.py) and passing
the result to the main thread via a CustomEvent, the callback then
stores the signals in each strategy context:

framework.add_worker(
BackgroundWorker(
framework,
poll_arb_feed,
func_kwargs={"scanner": ArbScanner(), "q": "senate"},
interval=60,
)
)

Without X402_WALLET_KEY the scanner returns bundled SAMPLE data and
makes no network calls, see examples/example-arbfeed.py
"""


def poll_arb_feed(
context: dict,
flumine,
scanner: ArbScanner,
q: str = None,
limit: int = 10,
mode: str = "opportunities",
) -> None:
try:
response = scanner.scan(q=q, limit=limit, mode=mode)
except (ArbScannerError, requests.RequestException, ValueError) as e:
logger.warning("poll_arb_feed error", extra={"error": str(e)})
return
context["polls"] = context.get("polls", 0) + 1
flumine.handler_queue.put(CustomEvent(response, callback))


def callback(flumine, event):
# executed on the main thread, safe to update strategy context
response = event.event
signals = response["opportunities"]
logger.info(
"Arb feed update: %s signals, %s executable%s",
len(signals),
len([s for s in signals if is_executable(s)]),
" (SAMPLE DATA)" if response.get("sample_data") else "",
)
for strategy in flumine.strategies:
strategy.context["arb_signals"] = signals
strategy.context["arb_signals_updated"] = time.time()
Loading