Compare commits

..
Author SHA1 Message Date
pselamyandClaude Opus 4.6 a8ad512eec fix: ruff lint and format fixes for persist assessment
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-06-14 19:18:31 +00:00
pselamyandClaude Opus 4.6 15ed793817 fix: add detector mock to pipeline test fixture
The persist_assessments feature accesses settings.detector which the
existing mock_settings fixture didn't include.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-06-14 19:16:52 +00:00
schrodinger01andpselamy e8fd508676 docs+test: add CHANGELOG and persistence regression tests
Documents the persist-all-assessments feature shipped in 8a0e8c9 and adds two regression tests covering: (1) sub-threshold assessments still hit the DB, and (2) DB write failures do not block alert dispatch.
2026-06-14 19:15:25 +00:00
jp-vps-deployandpselamy 1bbb672c17 feat(detector): persist all risk assessments to risk_assessments table 2026-06-14 19:15:19 +00:00
b962bdaee2 fix(ingestor): align WebSocket subscribe + routing with live API (#105)
* fix(ingestor): align WebSocket subscribe + routing with live API

The Polymarket ws-live-data WebSocket requires `action: "subscribe"` in
the subscribe envelope. Without it the server accepts the connection but
never delivers trade events, causing the tracker to silently produce
zero alerts.

Additionally, incoming frames are shaped `{connection_id, payload:{...}}`
-- they do NOT echo the `topic`/`type` keys we sent. The previous routing
check matched nothing and every real trade was silently dropped.

Changes:
- Add `action: "subscribe"` to subscription message
- Route incoming messages by payload shape (transactionHash + proxyWallet)
- Add ratchet tests for payload routing edge cases
- Rewrite README as agent-first with <2min quickstart
- Add skill draft (docs/skill-tracking-prediction-market-flow.md)

Closes #89

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(lint): remove unused imports in test_pipeline_persistence

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* style: apply ruff formatting

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-06-14 14:24:54 -04:00
axel-clawandGitHub ff678468e3 fix: load nested env settings and ignore blank webhooks (#91) 2026-04-11 16:20:15 -04:00
1f4fc607a7 feat: wire WalletRepository and FundingRepository into live pipeline (#90)
The live pipeline was receiving trade data and populating Redis caches but
never persisting wallet profiles or funding transfers to Postgres. This
wires the repositories into the _on_trade flow: when a fresh wallet signal
is detected, the wallet profile is upserted to wallet_profiles and the
funding chain is traced and inserted into funding_transfers. Existing
Redis caching behavior is preserved.

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-04-01 09:21:37 -04:00
232680a760 chore: smoke test fixes and code quality improvements (#88)
* fix: replace deprecated datetime.utcnow() and websockets.legacy APIs

- Replace datetime.utcnow() with datetime.now(UTC) in Orderbook model
- Use websockets.asyncio.client.connect instead of legacy websockets.connect
- Import ConnectionClosed from websockets.exceptions directly
- Update test mocks to patch the new import paths

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: resolve ruff lint errors in alembic migration

- Replace typing.Union with X | Y syntax (UP007)
- Import Sequence from collections.abc instead of typing (UP035)

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: remove deprecated version key from docker-compose.yml

The top-level 'version' key is obsolete in modern Docker Compose
and produces a warning.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* docs: fix README to match actual project structure and tooling

- Replace pip commands with uv equivalents
- Fix run command from 'python -m src.main' to 'python -m polymarket_insider_tracker'
- Update project structure tree to reflect actual src/polymarket_insider_tracker/ layout

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-09 22:03:10 -04:00
7c494c38a6 fix: replace (str, Enum) with StrEnum to resolve ruff UP042 (#87)
StrEnum (Python 3.11+) is the modern replacement for the (str, Enum)
pattern. Updates SyncState and PipelineState enums.

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-09 21:17:51 -04:00
17504a5eaa fix: replace (str, Enum) with StrEnum to resolve ruff UP042 (#86)
Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-09 21:16:23 -04:00
cadc593d62 fix: replace (str, Enum) with StrEnum to resolve ruff UP042 (#85)
Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-09 21:14:36 -04:00
0320a96f8a fix: replace (str, Enum) with StrEnum (ruff UP042) (#84)
Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-09 21:12:59 -04:00
753d6cf1de fix: replace (str, Enum) with StrEnum (ruff UP042) (#83)
Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-09 21:11:19 -04:00
f3801987fe fix: resolve .env.example variable interpolation and alembic env var mismatch (#62) (#82)
.env files don't support shell variable expansion, so DATABASE_URL and
REDIS_URL contained literal ${...} strings. Also alembic/env.py read
SQLALCHEMY_DATABASE_URL instead of DATABASE_URL.

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 20:05:00 -04:00
5cbcaccd82 fix: resolve alembic migration failures from .env and env var mismatch (#81)
- Replace shell variable interpolation in .env.example with literal
  values since .env files don't support variable expansion
- Change alembic/env.py to read DATABASE_URL (matching app convention)
  instead of SQLALCHEMY_DATABASE_URL

Closes #62

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 20:03:58 -04:00
288fbc75cf fix: resolve broken alembic migrations (#62) (#80)
- Replace shell variable interpolation in .env.example with literal values
- Change alembic/env.py to read DATABASE_URL instead of SQLALCHEMY_DATABASE_URL

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 20:02:53 -04:00
37ea0c6540 fix: resolve broken alembic migrations (#62) (#79)
- Replace shell variable interpolation in .env.example with literal values
- Change alembic/env.py to read DATABASE_URL instead of SQLALCHEMY_DATABASE_URL

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 20:02:02 -04:00
654f4466f0 fix: resolve broken alembic migrations (#62) (#78)
- Replace shell variable interpolation in .env.example with literal
  values since .env files don't expand variables
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match the app's env var name

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 20:00:56 -04:00
8628e3aadb fix: use literal values in .env.example and correct env var name in alembic (#77)
.env files don't support shell variable interpolation, so DATABASE_URL
and REDIS_URL had literal ${VAR} strings instead of actual values.
Also alembic/env.py read SQLALCHEMY_DATABASE_URL instead of DATABASE_URL.

Closes #62

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:59:46 -04:00
d808819990 fix: resolve alembic migration setup bugs (#62) (#76)
- Replace shell variable interpolation in .env.example with literal
  default values (.env files don't expand variables)
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match the app's convention

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:58:34 -04:00
20b8ad992c fix: resolve alembic migration setup bugs (#62) (#75)
- Replace shell variable interpolation in .env.example with literal
  default values since .env files don't expand variables
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match the app's env var name

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:57:33 -04:00
bd320f1a8e fix: resolve alembic migration setup bugs (#62) (#74)
- Replace shell variable interpolation in .env.example with literal
  values (.env files don't expand variables)
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match app configuration

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:56:30 -04:00
c65ed2225f fix: resolve alembic migration setup bugs (#62) (#73)
- Replace shell variable interpolation in .env.example with literal
  values (.env files don't expand variables)
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match app convention

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:55:24 -04:00
06ff48117c fix: use literal values in .env.example and correct env var name in alembic (#72)
.env files don't support shell variable interpolation, so DATABASE_URL
and REDIS_URL had literal ${VAR} strings instead of actual values.
Also alembic/env.py was reading SQLALCHEMY_DATABASE_URL instead of
DATABASE_URL, which is what the app and .env use.

Closes #62

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:54:15 -04:00
34393e7ee3 fix: resolve alembic migration failures from env config bugs (#62) (#71)
- Replace shell variable interpolation in .env.example with literal
  values (.env files don't expand ${VAR} syntax)
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match the app's env var name

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:53:02 -04:00
cf66aac9e3 fix: resolve alembic migration failures from env config bugs (#62) (#70)
- Replace shell variable interpolation in .env.example with literal
  values (.env files don't expand ${VAR} syntax)
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match the app's env var name

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:52:16 -04:00
13d5945b01 fix: resolve alembic migration failures from env config bugs (#62) (#69)
- Replace shell variable interpolation in .env.example with literal values
- Fix alembic/env.py to read DATABASE_URL instead of SQLALCHEMY_DATABASE_URL

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:51:13 -04:00
7352143004 fix: resolve two bugs preventing alembic migrations (#62) (#68)
1. Replace shell variable interpolation in .env.example with literal
   values, since .env files don't expand variables.
2. Change alembic/env.py to read DATABASE_URL instead of
   SQLALCHEMY_DATABASE_URL to match app convention.

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:49:57 -04:00
e05ac2ca64 fix: resolve two bugs preventing alembic migrations (#62) (#67)
1. Replace shell variable interpolation in .env.example with literal
   values, since .env files don't expand variables.
2. Change alembic/env.py to read DATABASE_URL instead of
   SQLALCHEMY_DATABASE_URL to match app convention.

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:49:18 -04:00
5d533c8c79 fix: resolve two bugs preventing alembic migrations (#62) (#66)
1. Replace shell variable interpolation in .env.example with literal
   values, since .env files don't expand variables.
2. Change alembic/env.py to read DATABASE_URL instead of
   SQLALCHEMY_DATABASE_URL to match app convention.

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:48:25 -04:00
e67f2d6d8f fix: resolve alembic migration setup bugs (#62) (#65)
- Replace shell variable interpolation in .env.example with literal values
- Change alembic/env.py to read DATABASE_URL instead of SQLALCHEMY_DATABASE_URL

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:47:23 -04:00
387a10132b fix: resolve alembic migration setup bugs (#62) (#64)
- Replace shell variable interpolation in .env.example with literal
  default values (.env files don't expand variables)
- Change alembic/env.py to read DATABASE_URL instead of
  SQLALCHEMY_DATABASE_URL to match the app's env var name

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:46:25 -04:00
b044bf25be fix: resolve broken alembic migrations (#62) (#63)
- Replace shell variable interpolation in .env.example with literal
  values, since .env files don't expand shell variables
- Fix alembic/env.py to read DATABASE_URL instead of SQLALCHEMY_DATABASE_URL
  to match the app's environment variable convention

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-08 19:45:21 -04:00
19 changed files with 1269 additions and 303 deletions
+2 -2
View File
@@ -6,12 +6,12 @@ POSTGRES_USER=tracker
POSTGRES_PASSWORD=dev_password
# Constructed database URL (for application use)
DATABASE_URL=postgresql://${POSTGRES_USER}:${POSTGRES_PASSWORD}@${POSTGRES_HOST}:${POSTGRES_PORT}/${POSTGRES_DB}
DATABASE_URL=postgresql://tracker:dev_password@localhost:5432/polymarket_tracker
# Redis Configuration
REDIS_HOST=localhost
REDIS_PORT=6379
REDIS_URL=redis://${REDIS_HOST}:${REDIS_PORT}
REDIS_URL=redis://localhost:6379
# Optional: Development tool ports
ADMINER_PORT=8080
+29
View File
@@ -0,0 +1,29 @@
# Changelog
All notable changes to this project are documented in this file.
The format is loosely based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/).
## [Unreleased]
### Added
- **Risk-assessment persistence**: every signal-bearing trade now writes a row
to the new `risk_assessments` table, regardless of whether the assessment
meets the alert threshold. This is the ground-truth log future backtests will
read instead of grepping `alerts.log` / `journalctl`.
- Pipeline: `Pipeline._score_and_alert` calls `Pipeline._persist_assessment`
for every assessment; failures are caught and never block alert dispatch.
- Storage: new `RiskAssessmentModel`, `RiskAssessmentDTO`, and
`RiskAssessmentRepository` (alembic migration shipped previously).
- Config: `DETECTOR_PERSIST_ASSESSMENTS` env var (default `true`) controls
the write path so it can be disabled without code changes.
- Tests: `tests/test_persist_assessment.py` covers (a) sub-threshold rows are
persisted with `should_alert=False` and dispatch is skipped, and (b) DB
failures during persistence do not block dispatching.
### Changed
- Alert threshold (`DETECTOR_ALERT_THRESHOLD`) is now fully env-driven; the
legacy hard-coded `0.6` default has been raised to `0.80` for production.
### Notes
- Backtest scripts can now source data from `risk_assessments` directly. The
`alerts.log` parsing path remains for one release as a fallback.
+118 -254
View File
@@ -2,80 +2,103 @@
**Detect informed money before the market moves.**
[![License: MIT](https://img.shields.io/badge/License-MIT-yellow.svg)](https://opensource.org/licenses/MIT)
[![CI](https://github.com/pselamy/polymarket-insider-tracker/actions/workflows/ci.yml/badge.svg)](https://github.com/pselamy/polymarket-insider-tracker/actions/workflows/ci.yml)
[![Python 3.11+](https://img.shields.io/badge/python-3.11+-blue.svg)](https://www.python.org/downloads/)
[![License: MIT](https://img.shields.io/badge/License-MIT-yellow.svg)](https://opensource.org/licenses/MIT)
Real-time detection of suspicious trading patterns on Polymarket: fresh wallets, unusual sizing, niche-market activity, and funding chain analysis. Streams trades via WebSocket, profiles wallets on-chain (Polygon), scores risk with ML + heuristics, and dispatches alerts to Discord/Telegram.
---
## The Opportunity
## Quick Start (< 2 minutes)
On January 3, 2026, a trader spotted a significant political event on Polymarket **before it happened**. How? Not by predicting the future, but by tracking suspicious trading behavior.
### 1. Install
> "You don't need to predict the future, you need to track suspicious behavior."
> — [@DidiTrading](https://x.com/DidiTrading)
An insider wallet turned **$35,000 into $442,000** (12.6x return) by entering a position hours before a major market move. The tool that detected this activity flagged five separate alerts before the event occurred.
**This repository builds that tool.**
---
## What This Does
The Polymarket Insider Tracker monitors prediction market trading activity in real-time and identifies patterns that suggest informed trading:
| Signal | What It Detects | Why It Matters |
|--------|-----------------|----------------|
| **Fresh Wallets** | Brand new wallets making large trades | Insiders create new wallets to hide their identity |
| **Unusual Sizing** | Trades that are disproportionately large for the market | Informed traders bet bigger when they have edge |
| **Niche Markets** | Activity in low-volume, specific-outcome markets | Easier to have inside information on obscure events |
| **Funding Chains** | Where wallet funds originated from | Links seemingly separate wallets to the same entity |
When suspicious activity is detected, you receive an instant alert with actionable intelligence.
---
## How It Works
```
┌─────────────────┐ ┌──────────────────┐ ┌────────────────────┐
│ Polymarket API │────>│ Wallet Profiler │────>│ Anomaly Detector │
│ (Real-time) │ │ (Blockchain) │ │ (ML + Heuristics) │
└─────────────────┘ └──────────────────┘ └────────────────────┘
┌────────────────────────────┘
v
┌─────────────────────┐
│ Alert Dispatcher │───> Discord / Telegram / Email
│ "Fresh wallet │
│ buying YES @7.5¢ │
│ on niche market" │
└─────────────────────┘
```bash
# Requires: Python 3.11+, Docker
git clone https://github.com/pselamy/polymarket-insider-tracker.git
cd polymarket-insider-tracker
uv sync --all-extras # or: pip install -e ".[dev]"
```
### Detection Algorithms
### 2. Start infrastructure
1. **Fresh Wallet Detection**
- Checks wallet transaction history on Polygon
- Flags wallets with fewer than 5 lifetime transactions making trades over $1,000
- Traces funding source to identify if connected to known entities
```bash
docker compose up -d # PostgreSQL 15 + Redis 7
docker compose ps # wait for healthy
```
2. **Liquidity Impact Analysis**
- Calculates trade size relative to market depth
- Flags trades consuming more than 2% of visible order book
- Weights by market category (niche markets score higher)
### 3. Configure
3. **Sniper Cluster Detection**
- Uses DBSCAN clustering to find wallets that consistently enter markets within minutes of creation
- Identifies coordinated behavior patterns
```bash
cp .env.example .env
# Edit .env — only DATABASE_URL and REDIS_URL are required for local dev
# (defaults in .env.example work with docker compose)
```
4. **Event Correlation**
- Cross-references trading activity with news feeds
- Detects positions opened 1-4 hours before related news breaks
### 4. Run migrations + start
```bash
uv run alembic upgrade head
uv run python -m polymarket_insider_tracker
```
You should see live trades within seconds:
```
INFO Connection state: disconnected -> connecting
INFO Connected to wss://ws-live-data.polymarket.com and subscribed to trades
DEBUG Trade: BUY 450 @ 1.00 on fifwc-ger-kor-2026-06-14-ger
DEBUG Trade: SELL 5 @ 0.86 on chi1-cd1-cdl-2026-06-14-draw
```
### CLI Options
```bash
python -m polymarket_insider_tracker --help
--version Show version
--config-check Validate configuration and exit
--log-level DEBUG Override log level
--dry-run Run pipeline without sending alerts
--health-port 8080 Override health check port
```
---
## Sample Alert
## Environment Variables
| Variable | Required | Default | Description |
|----------|----------|---------|-------------|
| `DATABASE_URL` | Yes | — | PostgreSQL connection string |
| `REDIS_URL` | No | `redis://localhost:6379` | Redis connection string |
| `POLYGON_RPC_URL` | No | `https://polygon-rpc.com` | Polygon RPC (public default works) |
| `POLYGON_FALLBACK_RPC_URL` | No | — | Fallback RPC endpoint |
| `POLYMARKET_WS_URL` | No | `wss://ws-live-data.polymarket.com` | WebSocket endpoint |
| `POLYMARKET_API_KEY` | No | — | Optional API key for higher rate limits |
| `DISCORD_WEBHOOK_URL` | No | — | Discord alerts |
| `TELEGRAM_BOT_TOKEN` | No | — | Telegram alerts (needs `TELEGRAM_CHAT_ID` too) |
| `TELEGRAM_CHAT_ID` | No | — | Telegram chat for alerts |
| `LOG_LEVEL` | No | `INFO` | Logging level |
| `DRY_RUN` | No | `false` | Skip sending alerts |
| `HEALTH_PORT` | No | `8080` | Health check HTTP port |
No API keys are needed for basic operation — the Polymarket WebSocket and CLOB REST APIs are public.
---
## What It Detects
| Signal | Detection Method | Threshold |
|--------|-----------------|-----------|
| **Fresh Wallets** | Wallet age < 48h, nonce <= 5, making trades > $1k | Confidence 0.5-0.9 |
| **Size Anomalies** | Trade size > 2% of 24h volume or > 5% of order book | Weighted by niche factor |
| **Niche Markets** | Low-volume markets (< $50k daily) with specific outcomes | 1.5x risk multiplier |
| **Funding Chains** | Trace wallet funding to known entities (exchanges, etc.) | On-chain lineage |
| **Sniper Clusters** | DBSCAN clustering of wallets entering within minutes | Coordinated behavior |
Risk scoring combines signals with configurable weights (default threshold: 0.6). Multi-signal bonuses: 2 signals +20%, 3+ signals +30%.
### Sample Alert
```
SUSPICIOUS ACTIVITY DETECTED
@@ -99,207 +122,62 @@ Confidence: HIGH (3/4 signals triggered)
---
## Quick Start
## Architecture
### Prerequisites
```
Polymarket WebSocket ──> Ingestor ──> Profiler ──> Detector ──> Alerter
(wss://ws-live-data) (trades) (on-chain) (scoring) (Discord/TG)
|
Polygon RPC
```
- Python 3.11+
- Docker and Docker Compose
- Polygon RPC endpoint (Alchemy, QuickNode, or self-hosted)
- Polymarket API key (free at [docs.polymarket.com](https://docs.polymarket.com))
### Components
### Installation
| Module | Purpose |
|--------|---------|
| `ingestor/` | WebSocket trade stream + CLOB REST client with rate limiting |
| `profiler/` | Polygon wallet analysis, entity identification, funding chain tracing |
| `detector/` | Fresh wallet, size anomaly, sniper cluster detection, composite risk scorer |
| `alerter/` | Multi-channel dispatch (Discord webhooks, Telegram bot) with dedup |
| `storage/` | SQLAlchemy ORM + Alembic migrations (PostgreSQL) |
| `pipeline.py` | Orchestrator wiring all components together |
| `shutdown.py` | Graceful SIGTERM/SIGINT handling with cleanup callbacks |
---
## Development
```bash
# Clone the repository
git clone https://github.com/pselamy/polymarket-insider-tracker.git
cd polymarket-insider-tracker
# Copy environment template
cp .env.example .env
# Edit .env with your API keys
# Start infrastructure (PostgreSQL, Redis)
docker compose up -d
# Wait for services to be healthy
docker compose ps
# Install Python dependencies
pip install -e .
# Run database migrations
alembic upgrade head
# Run the tracker
python -m src.main
uv run pytest # run tests
uv run ruff check src/ tests/ # lint
uv run ruff format src/ tests/ # format
uv run mypy src/ # type check (strict mode)
```
### Docker Services
The development stack includes:
| Service | Port | Description |
|---------|------|-------------|
| PostgreSQL 15 | 5432 | Primary database |
| Redis 7 | 6379 | Caching and pub/sub |
| Adminer | 8080 | Database admin UI (optional) |
| RedisInsight | 5540 | Redis admin UI (optional) |
```bash
# Start core services only
docker compose up -d
# Start with development tools (Adminer, RedisInsight)
docker compose --profile tools up -d
# View logs
docker compose logs -f
# Stop all services
docker compose down
# Stop and remove volumes (reset data)
docker compose down -v
```
### Configuration
```bash
# .env file
POLYGON_RPC_URL=https://polygon-mainnet.g.alchemy.com/v2/YOUR_KEY
POLYMARKET_API_KEY=your_polymarket_api_key
# Alert destinations (optional)
DISCORD_WEBHOOK_URL=https://discord.com/api/webhooks/...
TELEGRAM_BOT_TOKEN=your_bot_token
TELEGRAM_CHAT_ID=your_chat_id
# Detection thresholds
MIN_TRADE_SIZE_USDC=1000
FRESH_WALLET_MAX_NONCE=5
LIQUIDITY_IMPACT_THRESHOLD=0.02
```
| Adminer | 8080 | Database admin UI (optional, `--profile tools`) |
| RedisInsight | 5540 | Redis admin UI (optional, `--profile tools`) |
---
## Project Structure
## Troubleshooting
```
polymarket-insider-tracker/
├── src/
│ ├── ingestor/ # Real-time market data ingestion
│ │ ├── clob_client.py # Polymarket CLOB API wrapper
│ │ └── websocket.py # WebSocket event handler
│ ├── profiler/ # Wallet analysis
│ │ ├── analyzer.py # Core wallet profiling logic
│ │ ├── chain.py # Polygon blockchain client
│ │ └── funding.py # Funding chain tracer
│ ├── detector/ # Anomaly detection engines
│ │ ├── fresh_wallet.py
│ │ ├── size_anomaly.py
│ │ ├── sniper.py # DBSCAN clustering
│ │ └── scorer.py # Composite risk scoring
│ ├── alerter/ # Notification dispatch
│ │ ├── formatter.py # Alert message formatting
│ │ └── dispatcher.py # Multi-channel delivery
│ └── storage/ # Persistence layer
│ ├── models.py # SQLAlchemy models
│ └── repos.py # Repository pattern
├── tests/ # Test suite
├── scripts/
│ └── backtest.py # Historical analysis
├── docker-compose.yml
├── pyproject.toml
└── README.md
```
**No trades received / silent connection**
The WebSocket subscription requires `action: "subscribe"` in the envelope. If you're on an older version, update — this was fixed in the WebSocket protocol alignment (see #89).
---
**Connection timeout / DNS errors**
Verify `wss://ws-live-data.polymarket.com` is reachable from your network. Some corporate firewalls block WebSocket connections.
## Roadmap
**Database migration errors**
Ensure PostgreSQL is running (`docker compose ps`) and `DATABASE_URL` matches your docker-compose config. Run `uv run alembic upgrade head` after any schema changes.
### Phase 1: Core Detection (Current)
- [x] Project structure and documentation
- [ ] Polymarket CLOB API integration
- [ ] Fresh wallet detection
- [ ] Size anomaly detection
- [ ] Basic alerting (Discord/Telegram)
### Phase 2: Advanced Intelligence
- [ ] Funding chain analysis
- [ ] Sniper cluster detection (DBSCAN)
- [ ] Market categorization (niche vs mainstream)
- [ ] Historical backtesting framework
### Phase 3: Production Hardening
- [ ] High-availability deployment
- [ ] Rate limit management
- [ ] False positive feedback loop
- [ ] Web dashboard
---
## Why This Matters
Prediction markets are becoming a critical source of real-time probability estimates for world events. As they grow, so does the incentive for informed actors to exploit information asymmetry.
This tool democratizes access to the same detection capabilities that sophisticated traders use. Whether you are:
- **A trader** looking for alpha signals
- **A researcher** studying market microstructure
- **A platform operator** monitoring for manipulation
...this tracker provides visibility into the hidden flows that move markets.
---
## Technical Background
### Polymarket Architecture
Polymarket is a prediction market platform built on Polygon (Ethereum L2). Key characteristics:
- **CLOB (Central Limit Order Book)**: Centralized matching engine for speed
- **On-chain Settlement**: Final trades settle on Polygon blockchain
- **USDC Collateral**: All positions denominated in USDC stablecoin
- **Binary Outcomes**: Shares priced between $0.00 and $1.00
### Data Sources
| Source | Purpose | Latency |
|--------|---------|---------|
| Polymarket CLOB API | Real-time trades, orderbook | Milliseconds |
| Polygon RPC | Wallet history, nonce, funding | 1-2 seconds |
| Market Metadata API | Market categorization | On-demand |
### Detection Challenges
1. **Sybil Resistance**: Insiders use fresh wallets per trade
2. **Rate Limits**: Polygon RPC calls require caching strategy
3. **Market Classification**: NLP needed to categorize market niches
4. **Timing**: CLOB data leads on-chain by seconds
---
## Contributing
Contributions are welcome! Please read our Contributing Guide before submitting PRs.
### Development Setup
```bash
# Install dev dependencies
pip install -e ".[dev]"
# Run tests
pytest
# Run linting
ruff check src/
# Run type checking
mypy src/
```
**Rate limiting on Polygon RPC**
The default public RPC (`https://polygon-rpc.com`) has low rate limits. For production use, set `POLYGON_RPC_URL` to a dedicated provider (Alchemy, QuickNode, etc.).
---
@@ -312,20 +190,6 @@ This software is provided for **educational and research purposes only**.
- Insider trading is illegal in regulated markets; this tool is for transparency and research
- Users are responsible for compliance with applicable laws
---
## License
MIT License - see [LICENSE](LICENSE) for details.
---
## Acknowledgments
- Inspired by [@DidiTrading](https://x.com/DidiTrading) and [@spacexbt](https://x.com/spacexbt)
- Built on the open Polymarket API ecosystem
- Community contributions welcome
---
**Questions?** Open an issue or start a discussion.
+1 -1
View File
@@ -19,7 +19,7 @@ if config.config_file_name is not None:
target_metadata = Base.metadata
# Get database URL from environment variable or config
database_url = os.environ.get("SQLALCHEMY_DATABASE_URL")
database_url = os.environ.get("DATABASE_URL")
if database_url:
config.set_main_option("sqlalchemy.url", database_url)
@@ -5,16 +5,16 @@ Revises:
Create Date: 2026-01-04 00:00:00.000000+00:00
"""
from typing import Sequence, Union
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
# revision identifiers, used by Alembic.
revision: str = "001_initial"
down_revision: Union[str, None] = None
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
down_revision: str | None = None
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
@@ -0,0 +1,64 @@
"""Risk assessment persistence layer.
Adds the `risk_assessments` table — one row per signal-bearing trade —
so future backtests can rebuild ground truth without grepping the
systemd log or hammering the public data-api.
Revision ID: 002_risk_assessments
Revises: 001_initial
Create Date: 2026-05-22 11:30:00.000000+00:00
"""
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
revision: str = "002_risk_assessments"
down_revision: str | None = "001_initial"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.create_table(
"risk_assessments",
sa.Column("id", sa.Integer(), autoincrement=True, nullable=False),
sa.Column("assessment_id", sa.String(36), nullable=False),
sa.Column("trade_id", sa.String(80), nullable=False),
sa.Column("wallet_address", sa.String(42), nullable=False),
sa.Column("market_id", sa.String(80), nullable=False),
sa.Column("asset_id", sa.String(80), nullable=True),
sa.Column("side", sa.String(8), nullable=False),
sa.Column("outcome", sa.String(120), nullable=True),
sa.Column("outcome_index", sa.Integer(), nullable=True),
sa.Column("price", sa.Numeric(10, 6), nullable=False),
sa.Column("size", sa.Numeric(20, 6), nullable=False),
sa.Column("notional_usdc", sa.Numeric(20, 6), nullable=False),
sa.Column("trade_timestamp", sa.DateTime(timezone=True), nullable=False),
sa.Column("weighted_score", sa.Numeric(4, 3), nullable=False),
sa.Column("signals_triggered", sa.Integer(), nullable=False),
sa.Column("fresh_wallet_confidence", sa.Numeric(4, 3), nullable=True),
sa.Column("size_anomaly_confidence", sa.Numeric(4, 3), nullable=True),
sa.Column("is_niche_market", sa.Boolean(), nullable=True),
sa.Column("volume_impact", sa.Numeric(8, 4), nullable=True),
sa.Column("book_impact", sa.Numeric(8, 4), nullable=True),
sa.Column("wallet_age_hours", sa.Numeric(10, 2), nullable=True),
sa.Column("should_alert", sa.Boolean(), nullable=False),
sa.Column("threshold_at_eval", sa.Numeric(4, 3), nullable=False),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
sa.PrimaryKeyConstraint("id"),
sa.UniqueConstraint("assessment_id"),
)
op.create_index("idx_risk_assessments_wallet", "risk_assessments", ["wallet_address"])
op.create_index("idx_risk_assessments_market", "risk_assessments", ["market_id"])
op.create_index("idx_risk_assessments_trade_ts", "risk_assessments", ["trade_timestamp"])
op.create_index("idx_risk_assessments_score", "risk_assessments", ["weighted_score"])
def downgrade() -> None:
op.drop_index("idx_risk_assessments_score", table_name="risk_assessments")
op.drop_index("idx_risk_assessments_trade_ts", table_name="risk_assessments")
op.drop_index("idx_risk_assessments_market", table_name="risk_assessments")
op.drop_index("idx_risk_assessments_wallet", table_name="risk_assessments")
op.drop_table("risk_assessments")
-2
View File
@@ -1,5 +1,3 @@
version: "3.8"
services:
postgres:
image: postgres:15
@@ -0,0 +1,113 @@
# Skill: tracking-prediction-market-flow
Use when analyzing prediction market activity for informed-flow signals, insider
trading patterns, or suspicious wallet behavior on Polymarket.
## What This Tool Does
polymarket-insider-tracker streams real-time trades from Polymarket's WebSocket
feed, profiles trader wallets on the Polygon blockchain, and scores each trade
for informed-flow risk using multiple detection signals:
- **Fresh wallet detection**: New wallets (age < 48h, nonce <= 5) making large
trades (> $1k). Insiders create disposable wallets per trade.
- **Size anomaly detection**: Trades consuming > 2% of 24h volume or > 5% of
visible order book depth. Informed traders bet bigger when they have edge.
- **Niche market scoring**: Low-volume markets (< $50k daily) get a 1.5x risk
multiplier. Easier to have inside information on obscure events.
- **Funding chain analysis**: Traces wallet funding sources on-chain to link
seemingly separate wallets to the same entity or exchange.
- **Sniper cluster detection**: DBSCAN clustering identifies wallets that
consistently enter markets within minutes of creation.
Composite risk scoring combines signals with configurable weights (default
alert threshold: 0.6). Multi-signal bonuses: 2 signals +20%, 3+ signals +30%.
## Installation
```bash
git clone https://github.com/pselamy/polymarket-insider-tracker.git
cd polymarket-insider-tracker
uv sync --all-extras
docker compose up -d # PostgreSQL + Redis
cp .env.example .env # defaults work for local dev
uv run alembic upgrade head
```
No API keys required for basic operation (Polymarket APIs are public).
## Usage
```bash
# Start the tracker (streams trades, profiles wallets, scores risk, alerts)
uv run python -m polymarket_insider_tracker
# Dry run (no alerts sent)
uv run python -m polymarket_insider_tracker --dry-run
# Debug mode (see every trade)
uv run python -m polymarket_insider_tracker --log-level DEBUG
# Validate config without starting
uv run python -m polymarket_insider_tracker --config-check
```
## Interpreting Signals
### Risk Assessment Output
Each flagged trade produces a risk assessment with:
- **Confidence score** (0.0-1.0): Composite of weighted signals
- **Signal breakdown**: Which detectors fired and their individual confidence
- **Wallet profile**: Age, nonce, transaction count, funding source
- **Market context**: Volume, category, order book depth
### Signal Interpretation Guide
| Score Range | Interpretation | Action |
|-------------|---------------|--------|
| 0.6-0.7 | Moderate: single strong signal or two weak ones | Monitor, note the market |
| 0.7-0.85 | High: multiple signals converging | Investigate the market and wallet |
| 0.85-1.0 | Critical: fresh wallet + large size + niche market | High-confidence informed flow |
### What This Is NOT
- Not a trading signal generator. Informed flow != actionable alpha without
further analysis (hypothesis -> leakage-aware backtest -> capital).
- Not real-time enough for front-running. The tool detects patterns for
research and monitoring, not millisecond-level execution.
- Detection of informed flow does not prove insider trading. Many legitimate
reasons exist for the patterns this tool flags.
## Rate Limits and Etiquette
- **Polymarket WebSocket**: No explicit rate limit; one persistent connection.
Do not open multiple connections unnecessarily.
- **Polymarket CLOB REST**: Built-in rate limiter at 10 req/s with retry
backoff on 429/5xx. Respect this for metadata/orderbook queries.
- **Polygon RPC**: Public endpoints (polygon-rpc.com) have low limits. For
sustained use, configure a dedicated RPC provider via `POLYGON_RPC_URL`.
Built-in token-bucket rate limiter at 25 req/s with Redis caching (5min TTL).
## Known Pitfalls
1. **WebSocket subscription format**: Must include `action: "subscribe"` in the
envelope. Without it, the server accepts the connection but delivers zero
trade events (silent failure). Fixed in the current version.
2. **Message routing**: Live-data WebSocket pushes `{connection_id, payload:
{...trade fields}}`, not `{topic, type, payload}`. Route by checking for
`transactionHash` + `proxyWallet` keys in `payload`.
3. **Public RPC rate limits**: Default Polygon RPC will throttle under load.
Use a dedicated provider for production.
4. **Database required**: PostgreSQL + Redis must be running. Use
`docker compose up -d` for local dev.
## Cross-References
- **Repository**: https://github.com/pselamy/polymarket-insider-tracker
- **Issues**: https://github.com/pselamy/polymarket-insider-tracker/issues
- **Agent skill landing** (follow-on): selamy-labs/agent-skills
+53 -8
View File
@@ -18,7 +18,9 @@ from pydantic_settings import BaseSettings, SettingsConfigDict
class DatabaseSettings(BaseSettings):
"""Database connection settings."""
model_config = SettingsConfigDict(env_prefix="")
model_config = SettingsConfigDict(
env_prefix="", env_file=".env", env_file_encoding="utf-8", extra="ignore"
)
url: str = Field(
alias="DATABASE_URL",
@@ -37,7 +39,9 @@ class DatabaseSettings(BaseSettings):
class RedisSettings(BaseSettings):
"""Redis connection settings."""
model_config = SettingsConfigDict(env_prefix="")
model_config = SettingsConfigDict(
env_prefix="", env_file=".env", env_file_encoding="utf-8", extra="ignore"
)
url: str = Field(
default="redis://localhost:6379",
@@ -57,7 +61,9 @@ class RedisSettings(BaseSettings):
class PolygonSettings(BaseSettings):
"""Polygon blockchain RPC settings."""
model_config = SettingsConfigDict(env_prefix="POLYGON_")
model_config = SettingsConfigDict(
env_prefix="POLYGON_", env_file=".env", env_file_encoding="utf-8", extra="ignore"
)
rpc_url: str = Field(
default="https://polygon-rpc.com",
@@ -84,7 +90,9 @@ class PolygonSettings(BaseSettings):
class PolymarketSettings(BaseSettings):
"""Polymarket API settings."""
model_config = SettingsConfigDict(env_prefix="POLYMARKET_")
model_config = SettingsConfigDict(
env_prefix="POLYMARKET_", env_file=".env", env_file_encoding="utf-8", extra="ignore"
)
ws_url: str = Field(
default="wss://ws-subscriptions-clob.polymarket.com/ws/market",
@@ -109,7 +117,9 @@ class PolymarketSettings(BaseSettings):
class DiscordSettings(BaseSettings):
"""Discord notification settings."""
model_config = SettingsConfigDict(env_prefix="DISCORD_")
model_config = SettingsConfigDict(
env_prefix="DISCORD_", env_file=".env", env_file_encoding="utf-8", extra="ignore"
)
webhook_url: SecretStr | None = Field(
default=None,
@@ -120,13 +130,15 @@ class DiscordSettings(BaseSettings):
@property
def enabled(self) -> bool:
"""Check if Discord notifications are enabled."""
return self.webhook_url is not None
return self.webhook_url is not None and bool(self.webhook_url.get_secret_value().strip())
class TelegramSettings(BaseSettings):
"""Telegram notification settings."""
model_config = SettingsConfigDict(env_prefix="TELEGRAM_")
model_config = SettingsConfigDict(
env_prefix="TELEGRAM_", env_file=".env", env_file_encoding="utf-8", extra="ignore"
)
bot_token: SecretStr | None = Field(
default=None,
@@ -142,7 +154,39 @@ class TelegramSettings(BaseSettings):
@property
def enabled(self) -> bool:
"""Check if Telegram notifications are enabled."""
return self.bot_token is not None and self.chat_id is not None
return (
self.bot_token is not None
and bool(self.bot_token.get_secret_value().strip())
and self.chat_id is not None
and bool(self.chat_id.strip())
)
class DetectorSettings(BaseSettings):
"""Risk-scorer / detector tuning."""
model_config = SettingsConfigDict(
env_prefix="DETECTOR_", env_file=".env", env_file_encoding="utf-8", extra="ignore"
)
alert_threshold: float = Field(
default=0.80,
alias="DETECTOR_ALERT_THRESHOLD",
description="Minimum weighted score required to trigger an alert",
ge=0.0,
le=1.0,
)
dedup_window_seconds: int = Field(
default=3600,
alias="DETECTOR_DEDUP_WINDOW_SECONDS",
description="Per-(wallet, market) dedup window in seconds",
ge=0,
)
persist_assessments: bool = Field(
default=True,
alias="DETECTOR_PERSIST_ASSESSMENTS",
description="Write every signal-bearing risk assessment to the database",
)
class Settings(BaseSettings):
@@ -174,6 +218,7 @@ class Settings(BaseSettings):
polymarket: PolymarketSettings = Field(default_factory=PolymarketSettings)
discord: DiscordSettings = Field(default_factory=DiscordSettings)
telegram: TelegramSettings = Field(default_factory=TelegramSettings)
detector: DetectorSettings = Field(default_factory=DetectorSettings)
# Application settings
log_level: Literal["DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL"] = Field(
@@ -19,8 +19,12 @@ from polymarket_insider_tracker.ingestor.models import TradeEvent
logger = logging.getLogger(__name__)
# Default configuration
DEFAULT_ALERT_THRESHOLD = 0.6
# Default configuration. The threshold lifted from 0.6 to 0.80 after the
# first cost-adjusted backtest showed everything below 0.85 was follower-PnL
# negative under realistic taker fees + half-cent slippage. 0.80 keeps a small
# margin below 0.85+ so we don't drop borderline-high signals on a hard cliff.
# Override at runtime via DETECTOR_ALERT_THRESHOLD env var.
DEFAULT_ALERT_THRESHOLD = 0.80
DEFAULT_DEDUP_WINDOW_SECONDS = 3600 # 1 hour
DEFAULT_REDIS_KEY_PREFIX = "polymarket:dedup:"
@@ -86,7 +86,7 @@ class Orderbook:
bids: tuple[OrderbookLevel, ...]
asks: tuple[OrderbookLevel, ...]
tick_size: Decimal
timestamp: datetime = field(default_factory=datetime.utcnow)
timestamp: datetime = field(default_factory=lambda: datetime.now(UTC))
@classmethod
def from_clob_orderbook(cls, orderbook: Any) -> "Orderbook":
@@ -9,8 +9,9 @@ from dataclasses import dataclass
from enum import Enum
from typing import Any
import websockets
from websockets.asyncio.client import ClientConnection
from websockets.asyncio.client import connect as ws_connect
from websockets.exceptions import ConnectionClosed
from polymarket_insider_tracker.ingestor.models import TradeEvent
@@ -149,14 +150,14 @@ class TradeStreamHandler:
elif self._market_filter:
subscription["filters"] = json.dumps({"market_slug": self._market_filter})
return {"subscriptions": [subscription]}
return {"action": "subscribe", "subscriptions": [subscription]}
async def _connect(self) -> ClientConnection:
"""Establish WebSocket connection."""
await self._set_state(ConnectionState.CONNECTING)
try:
ws = await websockets.connect(
ws = await ws_connect(
self._host,
ping_interval=self._ping_interval,
ping_timeout=self._ping_interval * 2,
@@ -182,12 +183,13 @@ class TradeStreamHandler:
try:
data = json.loads(message)
# Check if this is a trade message
topic = data.get("topic")
msg_type = data.get("type")
if topic == "activity" and msg_type == "trades":
payload = data.get("payload", {})
# ws-live-data pushes {connection_id, payload:{...trade fields}}
payload = data.get("payload")
if (
isinstance(payload, dict)
and "transactionHash" in payload
and "proxyWallet" in payload
):
trade = TradeEvent.from_websocket_message(payload)
self._stats.trades_received += 1
@@ -207,8 +209,7 @@ class TradeStreamHandler:
logger.error("Error in trade callback: %s", e)
else:
# Log other message types for debugging
logger.debug("Received message: topic=%s type=%s", topic, msg_type)
logger.debug("Received non-trade message: %s", str(data)[:120])
except json.JSONDecodeError as e:
logger.warning("Invalid JSON message: %s", e)
@@ -227,7 +228,7 @@ class TradeStreamHandler:
else:
logger.debug("Received binary message (%d bytes)", len(message))
except websockets.ConnectionClosed as e:
except ConnectionClosed as e:
logger.warning("Connection closed: %s", e)
raise
except Exception as e:
@@ -284,7 +285,7 @@ class TradeStreamHandler:
while self._running:
try:
await self._listen(self._ws)
except (websockets.ConnectionClosed, Exception) as e:
except (ConnectionClosed, Exception) as e:
if not self._running:
break
+143 -2
View File
@@ -29,13 +29,23 @@ from polymarket_insider_tracker.ingestor.metadata_sync import MarketMetadataSync
from polymarket_insider_tracker.ingestor.websocket import TradeStreamHandler
from polymarket_insider_tracker.profiler.analyzer import WalletAnalyzer
from polymarket_insider_tracker.profiler.chain import PolygonClient
from polymarket_insider_tracker.profiler.funding import FundingTracer
from polymarket_insider_tracker.storage.database import DatabaseManager
from polymarket_insider_tracker.storage.repos import (
FundingRepository,
FundingTransferDTO,
RiskAssessmentDTO,
RiskAssessmentRepository,
WalletProfileDTO,
WalletRepository,
)
if TYPE_CHECKING:
from typing import Any
from polymarket_insider_tracker.detector.models import (
FreshWalletSignal,
RiskAssessment,
SizeAnomalySignal,
)
from polymarket_insider_tracker.ingestor.models import TradeEvent
@@ -120,6 +130,7 @@ class Pipeline:
self._alert_formatter: AlertFormatter | None = None
self._alert_dispatcher: AlertDispatcher | None = None
self._trade_stream: TradeStreamHandler | None = None
self._funding_tracer: FundingTracer | None = None
# Synchronization
self._stop_event: asyncio.Event | None = None
@@ -233,6 +244,10 @@ class Pipeline:
redis=self._redis,
)
# Initialize Funding Tracer
logger.debug("Initializing funding tracer...")
self._funding_tracer = FundingTracer(self._polygon_client)
# Initialize Detectors
logger.debug("Initializing detectors...")
self._fresh_wallet_detector = FreshWalletDetector(self._wallet_analyzer)
@@ -240,7 +255,17 @@ class Pipeline:
# Initialize Risk Scorer
logger.debug("Initializing risk scorer...")
self._risk_scorer = RiskScorer(self._redis)
self._risk_scorer = RiskScorer(
self._redis,
alert_threshold=settings.detector.alert_threshold,
dedup_window_seconds=settings.detector.dedup_window_seconds,
)
logger.info(
"RiskScorer threshold=%.2f dedup_window=%ds persist=%s",
settings.detector.alert_threshold,
settings.detector.dedup_window_seconds,
settings.detector.persist_assessments,
)
# Initialize Alerting
logger.debug("Initializing alerting components...")
@@ -365,6 +390,10 @@ class Pipeline:
self._detect_size_anomaly(trade),
)
# Persist wallet profile and funding data when a fresh wallet is detected
if fresh_signal is not None:
await self._persist_wallet_and_funding(fresh_signal)
# Bundle signals
bundle = SignalBundle(
trade_event=trade,
@@ -382,6 +411,63 @@ class Pipeline:
self._stats.errors += 1
self._stats.last_error = str(e)
async def _persist_wallet_and_funding(self, signal: FreshWalletSignal) -> None:
"""Persist wallet profile and funding transfers to Postgres.
Called when a fresh wallet signal is detected. Upserts the wallet
profile and traces/inserts any funding transfers found on-chain.
Args:
signal: The fresh wallet signal containing the wallet profile.
"""
if not self._db_manager:
return
profile = signal.wallet_profile
address = profile.address
try:
async with self._db_manager.get_async_session() as session:
# Persist wallet profile
wallet_repo = WalletRepository(session)
dto = WalletProfileDTO(
address=address,
nonce=profile.nonce,
first_seen_at=profile.first_seen,
is_fresh=profile.is_fresh,
matic_balance=profile.matic_balance,
usdc_balance=profile.usdc_balance,
analyzed_at=profile.analyzed_at,
)
await wallet_repo.upsert(dto)
# Trace and persist funding transfers
if self._funding_tracer:
chain = await self._funding_tracer.trace(address)
if chain.chain:
funding_repo = FundingRepository(session)
funding_dtos = [
FundingTransferDTO(
from_address=t.from_address,
to_address=t.to_address,
amount=t.amount,
token=t.token,
tx_hash=t.tx_hash,
block_number=t.block_number,
timestamp=t.timestamp,
)
for t in chain.chain
]
await funding_repo.insert_many(funding_dtos)
logger.debug(
"Persisted wallet profile and %d funding transfers for %s",
len(chain.chain) if self._funding_tracer and chain.chain else 0,
address[:10] + "...",
)
except Exception as e:
logger.warning("Failed to persist wallet/funding data for %s: %s", address, e)
async def _detect_fresh_wallet(self, trade: TradeEvent) -> FreshWalletSignal | None:
"""Run fresh wallet detection."""
if not self._fresh_wallet_detector:
@@ -403,13 +489,19 @@ class Pipeline:
return None
async def _score_and_alert(self, bundle: SignalBundle) -> None:
"""Score signals and send alert if threshold exceeded."""
"""Score signals, persist the assessment, and send alert if above threshold."""
if not self._risk_scorer or not self._alert_formatter or not self._alert_dispatcher:
return
# Get risk assessment
assessment = await self._risk_scorer.assess(bundle)
# Persist every signal-bearing assessment (not just delivered alerts).
# This is the ground-truth log future backtests will read instead of
# grepping systemd. Failure here must never block alerting.
if self._settings.detector.persist_assessments:
await self._persist_assessment(assessment)
if not assessment.should_alert:
logger.debug(
"Trade %s below alert threshold (score=%.2f)",
@@ -445,6 +537,55 @@ class Pipeline:
result.success_count + result.failure_count,
)
async def _persist_assessment(self, assessment: RiskAssessment) -> None:
"""Write the assessment row. Best-effort; never raises."""
if not self._db_manager:
return
from decimal import Decimal as _D
trade = assessment.trade_event
fresh = assessment.fresh_wallet_signal
size_sig = assessment.size_anomaly_signal
wallet_age: _D | None = None
if fresh is not None and fresh.wallet_profile.age_hours is not None:
wallet_age = _D(str(round(float(fresh.wallet_profile.age_hours), 2)))
dto = RiskAssessmentDTO(
assessment_id=assessment.assessment_id,
trade_id=trade.trade_id,
wallet_address=assessment.wallet_address.lower(),
market_id=assessment.market_id,
asset_id=getattr(trade, "asset_id", None) or None,
side=trade.side,
outcome=getattr(trade, "outcome", None) or None,
outcome_index=getattr(trade, "outcome_index", None),
price=trade.price,
size=trade.size,
notional_usdc=trade.notional_value,
trade_timestamp=trade.timestamp,
weighted_score=_D(str(round(assessment.weighted_score, 3))),
signals_triggered=assessment.signals_triggered,
fresh_wallet_confidence=(
_D(str(round(fresh.confidence, 3))) if fresh is not None else None
),
size_anomaly_confidence=(
_D(str(round(size_sig.confidence, 3))) if size_sig is not None else None
),
is_niche_market=size_sig.is_niche_market if size_sig is not None else None,
volume_impact=(
_D(str(round(size_sig.volume_impact, 4))) if size_sig is not None else None
),
book_impact=(_D(str(round(size_sig.book_impact, 4))) if size_sig is not None else None),
wallet_age_hours=wallet_age,
should_alert=assessment.should_alert,
threshold_at_eval=_D(str(round(self._settings.detector.alert_threshold, 3))),
)
try:
async with self._db_manager.get_async_session() as session:
repo = RiskAssessmentRepository(session)
await repo.insert(dto)
except Exception as e:
logger.warning("Failed to persist risk assessment %s: %s", assessment.assessment_id, e)
async def run(self) -> None:
"""Start the pipeline and run until interrupted.
@@ -115,3 +115,57 @@ class WalletRelationshipModel(Base):
Index("idx_wallet_relationships_a", "wallet_a"),
Index("idx_wallet_relationships_b", "wallet_b"),
)
class RiskAssessmentModel(Base):
"""SQLAlchemy model for risk assessments.
One row per signal-bearing trade (i.e. trades that triggered at least one
detector). Captures everything a future backtest needs without going back
to the public API: trade identity, score, per-signal confidences, and
whether the alert was actually delivered (could be False due to dedup or
threshold).
"""
__tablename__ = "risk_assessments"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
assessment_id: Mapped[str] = mapped_column(String(36), unique=True, nullable=False)
# Trade identity
trade_id: Mapped[str] = mapped_column(String(80), nullable=False)
wallet_address: Mapped[str] = mapped_column(String(42), nullable=False)
market_id: Mapped[str] = mapped_column(String(80), nullable=False)
asset_id: Mapped[str | None] = mapped_column(String(80), nullable=True)
side: Mapped[str] = mapped_column(String(8), nullable=False)
outcome: Mapped[str | None] = mapped_column(String(120), nullable=True)
outcome_index: Mapped[int | None] = mapped_column(Integer, nullable=True)
price: Mapped[Decimal] = mapped_column(Numeric(10, 6), nullable=False)
size: Mapped[Decimal] = mapped_column(Numeric(20, 6), nullable=False)
notional_usdc: Mapped[Decimal] = mapped_column(Numeric(20, 6), nullable=False)
trade_timestamp: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
# Scoring
weighted_score: Mapped[Decimal] = mapped_column(Numeric(4, 3), nullable=False)
signals_triggered: Mapped[int] = mapped_column(Integer, nullable=False)
fresh_wallet_confidence: Mapped[Decimal | None] = mapped_column(Numeric(4, 3), nullable=True)
size_anomaly_confidence: Mapped[Decimal | None] = mapped_column(Numeric(4, 3), nullable=True)
is_niche_market: Mapped[bool | None] = mapped_column(Boolean, nullable=True)
volume_impact: Mapped[Decimal | None] = mapped_column(Numeric(8, 4), nullable=True)
book_impact: Mapped[Decimal | None] = mapped_column(Numeric(8, 4), nullable=True)
wallet_age_hours: Mapped[Decimal | None] = mapped_column(Numeric(10, 2), nullable=True)
# Decision
should_alert: Mapped[bool] = mapped_column(Boolean, nullable=False)
threshold_at_eval: Mapped[Decimal] = mapped_column(Numeric(4, 3), nullable=False)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, default=lambda: datetime.now(UTC)
)
__table_args__ = (
Index("idx_risk_assessments_wallet", "wallet_address"),
Index("idx_risk_assessments_market", "market_id"),
Index("idx_risk_assessments_trade_ts", "trade_timestamp"),
Index("idx_risk_assessments_score", "weighted_score"),
)
@@ -18,6 +18,7 @@ from sqlalchemy.dialects.sqlite import insert as sqlite_insert
from polymarket_insider_tracker.storage.models import (
FundingTransferModel,
RiskAssessmentModel,
WalletProfileModel,
WalletRelationshipModel,
)
@@ -510,3 +511,107 @@ class RelationshipRepository:
)
# SQLAlchemy Result does have rowcount but typing doesn't reflect it
return (result.rowcount or 0) > 0 # type: ignore[attr-defined]
@dataclass
class RiskAssessmentDTO:
"""Data transfer object for a persisted risk assessment.
Captures everything a future backtest needs without going back to
public APIs: trade identity, score, per-signal confidences, and
whether the alert was actually delivered.
"""
assessment_id: str
trade_id: str
wallet_address: str
market_id: str
asset_id: str | None
side: str
outcome: str | None
outcome_index: int | None
price: Decimal
size: Decimal
notional_usdc: Decimal
trade_timestamp: datetime
weighted_score: Decimal
signals_triggered: int
fresh_wallet_confidence: Decimal | None
size_anomaly_confidence: Decimal | None
is_niche_market: bool | None
volume_impact: Decimal | None
book_impact: Decimal | None
wallet_age_hours: Decimal | None
should_alert: bool
threshold_at_eval: Decimal
created_at: datetime | None = None
class RiskAssessmentRepository:
"""Repository for risk assessment data access."""
def __init__(self, session: AsyncSession) -> None:
self.session = session
async def insert(self, dto: RiskAssessmentDTO) -> RiskAssessmentDTO:
"""Insert a single assessment. Idempotent on assessment_id collisions."""
model = RiskAssessmentModel(
assessment_id=dto.assessment_id,
trade_id=dto.trade_id,
wallet_address=dto.wallet_address.lower(),
market_id=dto.market_id,
asset_id=dto.asset_id,
side=dto.side,
outcome=dto.outcome,
outcome_index=dto.outcome_index,
price=dto.price,
size=dto.size,
notional_usdc=dto.notional_usdc,
trade_timestamp=dto.trade_timestamp,
weighted_score=dto.weighted_score,
signals_triggered=dto.signals_triggered,
fresh_wallet_confidence=dto.fresh_wallet_confidence,
size_anomaly_confidence=dto.size_anomaly_confidence,
is_niche_market=dto.is_niche_market,
volume_impact=dto.volume_impact,
book_impact=dto.book_impact,
wallet_age_hours=dto.wallet_age_hours,
should_alert=dto.should_alert,
threshold_at_eval=dto.threshold_at_eval,
)
self.session.add(model)
await self.session.flush()
return dto
async def get_by_assessment_id(self, assessment_id: str) -> RiskAssessmentDTO | None:
result = await self.session.execute(
select(RiskAssessmentModel).where(RiskAssessmentModel.assessment_id == assessment_id)
)
model = result.scalar_one_or_none()
if model is None:
return None
return RiskAssessmentDTO(
assessment_id=model.assessment_id,
trade_id=model.trade_id,
wallet_address=model.wallet_address,
market_id=model.market_id,
asset_id=model.asset_id,
side=model.side,
outcome=model.outcome,
outcome_index=model.outcome_index,
price=model.price,
size=model.size,
notional_usdc=model.notional_usdc,
trade_timestamp=model.trade_timestamp,
weighted_score=model.weighted_score,
signals_triggered=model.signals_triggered,
fresh_wallet_confidence=model.fresh_wallet_confidence,
size_anomaly_confidence=model.size_anomaly_confidence,
is_niche_market=model.is_niche_market,
volume_impact=model.volume_impact,
book_impact=model.book_impact,
wallet_age_hours=model.wallet_age_hours,
should_alert=model.should_alert,
threshold_at_eval=model.threshold_at_eval,
created_at=model.created_at,
)
+49 -14
View File
@@ -84,7 +84,10 @@ class TestTradeStreamHandler:
"""Test building subscription message without filters."""
msg = handler._build_subscription_message()
assert msg == {"subscriptions": [{"topic": "activity", "type": "trades"}]}
assert msg == {
"action": "subscribe",
"subscriptions": [{"topic": "activity", "type": "trades"}],
}
def test_build_subscription_message_with_event_filter(self, on_trade_mock: AsyncMock) -> None:
"""Test building subscription message with event filter."""
@@ -110,11 +113,10 @@ class TestTradeStreamHandler:
async def test_handle_message_trade(
self, handler: TradeStreamHandler, on_trade_mock: AsyncMock
) -> None:
"""Test handling a valid trade message."""
"""Test handling a valid trade message (payload-based routing)."""
message = json.dumps(
{
"topic": "activity",
"type": "trades",
"connection_id": "abc123",
"payload": {
"conditionId": "0xmarket",
"transactionHash": "0xtx",
@@ -142,11 +144,10 @@ class TestTradeStreamHandler:
async def test_handle_message_non_trade(
self, handler: TradeStreamHandler, on_trade_mock: AsyncMock
) -> None:
"""Test handling a non-trade message."""
"""Test handling a non-trade message (no transactionHash/proxyWallet)."""
message = json.dumps(
{
"topic": "comments",
"type": "comment_created",
"connection_id": "abc123",
"payload": {"body": "Hello"},
}
)
@@ -156,6 +157,35 @@ class TestTradeStreamHandler:
on_trade_mock.assert_not_called()
assert handler.stats.trades_received == 0
@pytest.mark.asyncio
async def test_handle_message_payload_missing_proxy_wallet(
self, handler: TradeStreamHandler, on_trade_mock: AsyncMock
) -> None:
"""Ratchet: payload with transactionHash but no proxyWallet is not a trade."""
message = json.dumps(
{
"connection_id": "abc",
"payload": {"transactionHash": "0xtx", "other": "field"},
}
)
await handler._handle_message(message)
on_trade_mock.assert_not_called()
assert handler.stats.trades_received == 0
@pytest.mark.asyncio
async def test_handle_message_no_payload_key(
self, handler: TradeStreamHandler, on_trade_mock: AsyncMock
) -> None:
"""Ratchet: message without payload key is ignored."""
message = json.dumps({"connection_id": "abc", "status": "ok"})
await handler._handle_message(message)
on_trade_mock.assert_not_called()
assert handler.stats.trades_received == 0
@pytest.mark.asyncio
async def test_handle_message_invalid_json(
self, handler: TradeStreamHandler, on_trade_mock: AsyncMock
@@ -175,8 +205,7 @@ class TestTradeStreamHandler:
message = json.dumps(
{
"topic": "activity",
"type": "trades",
"connection_id": "abc",
"payload": {
"conditionId": "0x",
"transactionHash": "0x",
@@ -231,14 +260,18 @@ class TestTradeStreamHandler:
mock_ws = AsyncMock()
mock_ws.send = AsyncMock()
with patch("websockets.connect", AsyncMock(return_value=mock_ws)):
with patch(
"polymarket_insider_tracker.ingestor.websocket.ws_connect",
AsyncMock(return_value=mock_ws),
):
ws = await handler._connect()
assert ws is mock_ws
mock_ws.send.assert_called_once()
# Verify subscription message
# Verify subscription message includes action: subscribe
sent_msg = json.loads(mock_ws.send.call_args[0][0])
assert sent_msg["action"] == "subscribe"
assert "subscriptions" in sent_msg
assert sent_msg["subscriptions"][0]["topic"] == "activity"
assert sent_msg["subscriptions"][0]["type"] == "trades"
@@ -289,8 +322,7 @@ class TestTradeStreamHandlerIntegration:
trade_message = json.dumps(
{
"topic": "activity",
"type": "trades",
"connection_id": "test-conn",
"payload": {
"conditionId": "0xtest",
"transactionHash": "0xtx",
@@ -333,7 +365,10 @@ class TestTradeStreamHandlerIntegration:
mock_ws = MockWebSocket(handler, trade_message)
with patch("websockets.connect", AsyncMock(return_value=mock_ws)):
with patch(
"polymarket_insider_tracker.ingestor.websocket.ws_connect",
AsyncMock(return_value=mock_ws),
):
# Run with timeout to prevent hanging
try:
await asyncio.wait_for(handler.start(), timeout=1.0)
+175
View File
@@ -0,0 +1,175 @@
"""Tests for RiskAssessment persistence inside Pipeline._score_and_alert.
Verifies:
1. Every signal-bearing assessment is written to risk_assessments, even
when ``should_alert`` is False (i.e. below the alert threshold).
2. A DB failure during persistence never blocks alert dispatching.
"""
from __future__ import annotations
from datetime import UTC, datetime
from decimal import Decimal
from unittest.mock import AsyncMock, MagicMock
import pytest
from sqlalchemy import select
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from polymarket_insider_tracker.config import Settings
from polymarket_insider_tracker.detector.models import RiskAssessment
from polymarket_insider_tracker.detector.scorer import SignalBundle
from polymarket_insider_tracker.ingestor.models import TradeEvent
from polymarket_insider_tracker.pipeline import Pipeline
from polymarket_insider_tracker.storage.database import DatabaseManager
from polymarket_insider_tracker.storage.models import Base, RiskAssessmentModel
# ---------------------------------------------------------------------------
# Fixtures
# ---------------------------------------------------------------------------
@pytest.fixture
def mock_settings():
"""Settings stub with the attributes Pipeline reaches for at runtime."""
detector = MagicMock()
detector.persist_assessments = True
detector.alert_threshold = 0.8
settings = MagicMock(spec=Settings)
settings.detector = detector
settings.dry_run = False
return settings
@pytest.fixture
async def async_engine():
engine = create_async_engine("sqlite+aiosqlite:///:memory:", echo=False)
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
yield engine
await engine.dispose()
@pytest.fixture
async def db_manager(async_engine):
manager = DatabaseManager.__new__(DatabaseManager)
manager.database_url = "sqlite+aiosqlite:///:memory:"
manager.async_mode = True
manager._pool_size = 5
manager._max_overflow = 10
manager._echo = False
manager._sync_engine = None
manager._async_engine = async_engine
manager._sync_session_factory = None
manager._async_session_factory = async_sessionmaker(bind=async_engine, expire_on_commit=False)
return manager
@pytest.fixture
def sample_trade() -> TradeEvent:
return TradeEvent(
trade_id="0x" + "a" * 64,
wallet_address="0x" + "b" * 40,
market_id="0x" + "c" * 64,
asset_id="asset_xyz",
side="BUY",
price=Decimal("0.42"),
size=Decimal("1000"),
timestamp=datetime.now(UTC),
outcome="Yes",
outcome_index=0,
event_title="Test Event",
market_slug="test-market",
)
def _make_assessment(trade: TradeEvent, *, should_alert: bool, score: float) -> RiskAssessment:
return RiskAssessment(
trade_event=trade,
wallet_address=trade.wallet_address,
market_id=trade.market_id,
fresh_wallet_signal=None,
size_anomaly_signal=None,
signals_triggered=1,
weighted_score=score,
should_alert=should_alert,
)
def _build_pipeline(
mock_settings,
*,
db_manager=None,
assessment: RiskAssessment,
dispatcher: MagicMock | None = None,
) -> Pipeline:
"""Construct a Pipeline with the minimum collaborators wired in."""
pipeline = Pipeline(mock_settings)
pipeline._db_manager = db_manager
pipeline._risk_scorer = MagicMock()
pipeline._risk_scorer.assess = AsyncMock(return_value=assessment)
pipeline._alert_formatter = MagicMock()
pipeline._alert_formatter.format = MagicMock(return_value=MagicMock())
if dispatcher is None:
dispatcher = MagicMock()
dispatcher.dispatch = AsyncMock(
return_value=MagicMock(all_succeeded=True, success_count=1, failure_count=0)
)
pipeline._alert_dispatcher = dispatcher
pipeline._dry_run = False
return pipeline
# ---------------------------------------------------------------------------
# Tests
# ---------------------------------------------------------------------------
class TestPersistAssessment:
@pytest.mark.asyncio
async def test_below_threshold_assessment_is_persisted(
self, mock_settings, db_manager, sample_trade, async_engine
):
"""Assessments with should_alert=False must still hit the DB; no dispatch."""
assessment = _make_assessment(sample_trade, should_alert=False, score=0.45)
pipeline = _build_pipeline(mock_settings, db_manager=db_manager, assessment=assessment)
await pipeline._score_and_alert(SignalBundle(trade_event=sample_trade))
# Row landed in risk_assessments
async with async_sessionmaker(bind=async_engine, expire_on_commit=False)() as session:
rows = (await session.execute(select(RiskAssessmentModel))).scalars().all()
assert len(rows) == 1
row = rows[0]
assert row.assessment_id == assessment.assessment_id
assert row.should_alert is False
assert float(row.weighted_score) == pytest.approx(0.45, abs=1e-3)
assert row.wallet_address == sample_trade.wallet_address.lower()
# No alert dispatched for sub-threshold assessments
pipeline._alert_dispatcher.dispatch.assert_not_called()
assert pipeline.stats.alerts_sent == 0
@pytest.mark.asyncio
async def test_persistence_failure_does_not_block_dispatch(self, mock_settings, sample_trade):
"""If repo.insert blows up, the alert pipeline still ships the alert."""
assessment = _make_assessment(sample_trade, should_alert=True, score=0.92)
# db_manager whose get_async_session raises -> _persist_assessment swallows it
broken_db = MagicMock()
broken_db.get_async_session = MagicMock(side_effect=RuntimeError("DB connection failed"))
pipeline = _build_pipeline(mock_settings, db_manager=broken_db, assessment=assessment)
await pipeline._score_and_alert(SignalBundle(trade_event=sample_trade))
# DB write was attempted and failed silently
broken_db.get_async_session.assert_called_once()
# Dispatcher still ran and the stats counter incremented
pipeline._alert_dispatcher.dispatch.assert_awaited_once()
assert pipeline.stats.alerts_sent == 1
+4
View File
@@ -44,6 +44,9 @@ def mock_settings():
telegram.bot_token = None
telegram.chat_id = None
detector = MagicMock()
detector.persist_assessments = False
settings = MagicMock(spec=Settings)
settings.redis = redis
settings.database = database
@@ -51,6 +54,7 @@ def mock_settings():
settings.polymarket = polymarket
settings.discord = discord
settings.telegram = telegram
settings.detector = detector
settings.dry_run = True
return settings
+334
View File
@@ -0,0 +1,334 @@
"""Tests verifying wallet and funding data persistence in the pipeline.
These tests confirm that running the live pipeline writes rows into
wallet_profiles and funding_transfers tables when fresh wallets are detected.
"""
from __future__ import annotations
from datetime import UTC, datetime
from decimal import Decimal
from unittest.mock import AsyncMock, MagicMock
import pytest
from sqlalchemy import select
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from polymarket_insider_tracker.config import Settings
from polymarket_insider_tracker.detector.models import FreshWalletSignal
from polymarket_insider_tracker.ingestor.models import TradeEvent
from polymarket_insider_tracker.pipeline import Pipeline
from polymarket_insider_tracker.profiler.models import FundingChain, FundingTransfer, WalletProfile
from polymarket_insider_tracker.storage.database import DatabaseManager
from polymarket_insider_tracker.storage.models import Base, FundingTransferModel, WalletProfileModel
@pytest.fixture
def mock_settings():
"""Create mock settings for testing."""
redis = MagicMock()
redis.url = "redis://localhost:6379"
database = MagicMock()
database.url = "sqlite+aiosqlite:///:memory:"
polygon = MagicMock()
polygon.rpc_url = "https://polygon-rpc.com"
polygon.fallback_rpc_url = None
polymarket = MagicMock()
polymarket.ws_url = "wss://ws-subscriptions-clob.polymarket.com/ws/market"
polymarket.api_key = None
discord = MagicMock()
discord.enabled = False
discord.webhook_url = None
telegram = MagicMock()
telegram.enabled = False
telegram.bot_token = None
telegram.chat_id = None
settings = MagicMock(spec=Settings)
settings.redis = redis
settings.database = database
settings.polygon = polygon
settings.polymarket = polymarket
settings.discord = discord
settings.telegram = telegram
settings.dry_run = True
return settings
@pytest.fixture
async def async_engine():
"""Create an async SQLite engine for testing."""
engine = create_async_engine("sqlite+aiosqlite:///:memory:", echo=False)
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
yield engine
await engine.dispose()
@pytest.fixture
async def db_manager(async_engine):
"""Create a DatabaseManager backed by the in-memory SQLite engine."""
manager = DatabaseManager.__new__(DatabaseManager)
manager.database_url = "sqlite+aiosqlite:///:memory:"
manager.async_mode = True
manager._pool_size = 5
manager._max_overflow = 10
manager._echo = False
manager._sync_engine = None
manager._async_engine = async_engine
manager._sync_session_factory = None
manager._async_session_factory = async_sessionmaker(bind=async_engine, expire_on_commit=False)
return manager
@pytest.fixture
def sample_trade():
"""Create a sample trade event."""
return TradeEvent(
trade_id="0x" + "a" * 64,
wallet_address="0x" + "b" * 40,
market_id="0x" + "c" * 64,
asset_id="asset_123",
side="BUY",
price=Decimal("0.65"),
size=Decimal("5000"),
timestamp=datetime.now(UTC),
outcome="Yes",
outcome_index=0,
event_title="Test Market",
market_slug="test-market",
)
@pytest.fixture
def sample_profile():
"""Create a sample fresh wallet profile."""
return WalletProfile(
address="0x" + "b" * 40,
nonce=2,
first_seen=datetime(2026, 3, 31, 12, 0, 0, tzinfo=UTC),
age_hours=1.5,
is_fresh=True,
total_tx_count=2,
matic_balance=Decimal("1000000000000000000"),
usdc_balance=Decimal("5000000000"),
fresh_threshold=5,
)
@pytest.fixture
def sample_funding_chain():
"""Create a sample funding chain with one transfer."""
return FundingChain(
target_address="0x" + "b" * 40,
chain=[
FundingTransfer(
from_address="0x" + "d" * 40,
to_address="0x" + "b" * 40,
amount=Decimal("5000000000"),
token="USDC",
tx_hash="0x" + "e" * 64,
block_number=12345678,
timestamp=datetime(2026, 3, 31, 11, 0, 0, tzinfo=UTC),
),
],
origin_address="0x" + "d" * 40,
origin_type="cex_binance",
hop_count=1,
)
class TestPipelinePersistence:
"""Tests that the pipeline persists wallet and funding data to Postgres."""
@pytest.mark.asyncio
async def test_on_trade_persists_wallet_profile(
self, mock_settings, db_manager, sample_trade, sample_profile, async_engine
):
"""When a fresh wallet signal fires, the wallet profile is written to wallet_profiles."""
pipeline = Pipeline(mock_settings)
pipeline._db_manager = db_manager
fresh_signal = FreshWalletSignal(
trade_event=sample_trade,
wallet_profile=sample_profile,
confidence=0.8,
factors={"base": 0.5, "brand_new": 0.2},
)
pipeline._fresh_wallet_detector = MagicMock()
pipeline._fresh_wallet_detector.analyze = AsyncMock(return_value=fresh_signal)
pipeline._size_anomaly_detector = MagicMock()
pipeline._size_anomaly_detector.analyze = AsyncMock(return_value=None)
pipeline._funding_tracer = MagicMock()
pipeline._funding_tracer.trace = AsyncMock(
return_value=FundingChain(target_address=sample_profile.address)
)
pipeline._risk_scorer = MagicMock()
pipeline._risk_scorer.assess = AsyncMock(
return_value=MagicMock(should_alert=False, weighted_score=0.3)
)
pipeline._alert_formatter = MagicMock()
pipeline._alert_dispatcher = MagicMock()
await pipeline._on_trade(sample_trade)
# Verify wallet_profiles has a row
async with async_sessionmaker(bind=async_engine, expire_on_commit=False)() as session:
result = await session.execute(select(WalletProfileModel))
rows = result.scalars().all()
assert len(rows) == 1
assert rows[0].address == sample_profile.address.lower()
assert rows[0].nonce == sample_profile.nonce
assert rows[0].is_fresh is True
@pytest.mark.asyncio
async def test_on_trade_persists_funding_transfers(
self,
mock_settings,
db_manager,
sample_trade,
sample_profile,
sample_funding_chain,
async_engine,
):
"""When a fresh wallet signal fires, funding transfers are written to funding_transfers."""
pipeline = Pipeline(mock_settings)
pipeline._db_manager = db_manager
fresh_signal = FreshWalletSignal(
trade_event=sample_trade,
wallet_profile=sample_profile,
confidence=0.8,
factors={"base": 0.5, "brand_new": 0.2},
)
pipeline._fresh_wallet_detector = MagicMock()
pipeline._fresh_wallet_detector.analyze = AsyncMock(return_value=fresh_signal)
pipeline._size_anomaly_detector = MagicMock()
pipeline._size_anomaly_detector.analyze = AsyncMock(return_value=None)
pipeline._funding_tracer = MagicMock()
pipeline._funding_tracer.trace = AsyncMock(return_value=sample_funding_chain)
pipeline._risk_scorer = MagicMock()
pipeline._risk_scorer.assess = AsyncMock(
return_value=MagicMock(should_alert=False, weighted_score=0.3)
)
pipeline._alert_formatter = MagicMock()
pipeline._alert_dispatcher = MagicMock()
await pipeline._on_trade(sample_trade)
# Verify funding_transfers has a row
async with async_sessionmaker(bind=async_engine, expire_on_commit=False)() as session:
result = await session.execute(select(FundingTransferModel))
rows = result.scalars().all()
assert len(rows) == 1
assert rows[0].to_address == ("0x" + "b" * 40).lower()
assert rows[0].from_address == ("0x" + "d" * 40).lower()
assert rows[0].token == "USDC"
assert rows[0].tx_hash == ("0x" + "e" * 64).lower()
@pytest.mark.asyncio
async def test_no_persistence_without_fresh_signal(
self, mock_settings, db_manager, sample_trade, async_engine
):
"""No rows written when fresh wallet signal is None (wallet not fresh)."""
pipeline = Pipeline(mock_settings)
pipeline._db_manager = db_manager
pipeline._fresh_wallet_detector = MagicMock()
pipeline._fresh_wallet_detector.analyze = AsyncMock(return_value=None)
pipeline._size_anomaly_detector = MagicMock()
pipeline._size_anomaly_detector.analyze = AsyncMock(return_value=None)
await pipeline._on_trade(sample_trade)
async with async_sessionmaker(bind=async_engine, expire_on_commit=False)() as session:
wallets = (await session.execute(select(WalletProfileModel))).scalars().all()
transfers = (await session.execute(select(FundingTransferModel))).scalars().all()
assert len(wallets) == 0
assert len(transfers) == 0
@pytest.mark.asyncio
async def test_persistence_failure_does_not_break_pipeline(
self, mock_settings, sample_trade, sample_profile
):
"""Persistence errors are caught and don't crash trade processing."""
pipeline = Pipeline(mock_settings)
# Use a broken db_manager that raises on get_async_session
broken_db = MagicMock()
broken_db.get_async_session = MagicMock(side_effect=Exception("DB connection failed"))
pipeline._db_manager = broken_db
fresh_signal = FreshWalletSignal(
trade_event=sample_trade,
wallet_profile=sample_profile,
confidence=0.8,
factors={"base": 0.5},
)
pipeline._fresh_wallet_detector = MagicMock()
pipeline._fresh_wallet_detector.analyze = AsyncMock(return_value=fresh_signal)
pipeline._size_anomaly_detector = MagicMock()
pipeline._size_anomaly_detector.analyze = AsyncMock(return_value=None)
pipeline._funding_tracer = MagicMock()
pipeline._risk_scorer = MagicMock()
pipeline._risk_scorer.assess = AsyncMock(
return_value=MagicMock(should_alert=False, weighted_score=0.3)
)
pipeline._alert_formatter = MagicMock()
pipeline._alert_dispatcher = MagicMock()
# Should not raise
await pipeline._on_trade(sample_trade)
assert pipeline.stats.trades_processed == 1
@pytest.mark.asyncio
async def test_duplicate_funding_transfers_are_skipped(
self,
mock_settings,
db_manager,
sample_trade,
sample_profile,
sample_funding_chain,
async_engine,
):
"""Processing the same trade twice should not duplicate funding transfer rows."""
pipeline = Pipeline(mock_settings)
pipeline._db_manager = db_manager
fresh_signal = FreshWalletSignal(
trade_event=sample_trade,
wallet_profile=sample_profile,
confidence=0.8,
factors={"base": 0.5},
)
pipeline._fresh_wallet_detector = MagicMock()
pipeline._fresh_wallet_detector.analyze = AsyncMock(return_value=fresh_signal)
pipeline._size_anomaly_detector = MagicMock()
pipeline._size_anomaly_detector.analyze = AsyncMock(return_value=None)
pipeline._funding_tracer = MagicMock()
pipeline._funding_tracer.trace = AsyncMock(return_value=sample_funding_chain)
pipeline._risk_scorer = MagicMock()
pipeline._risk_scorer.assess = AsyncMock(
return_value=MagicMock(should_alert=False, weighted_score=0.3)
)
pipeline._alert_formatter = MagicMock()
pipeline._alert_dispatcher = MagicMock()
# Process same trade twice
await pipeline._on_trade(sample_trade)
await pipeline._on_trade(sample_trade)
# Should still have only 1 funding transfer (duplicate skipped)
async with async_sessionmaker(bind=async_engine, expire_on_commit=False)() as session:
result = await session.execute(select(FundingTransferModel))
rows = result.scalars().all()
assert len(rows) == 1