diff --git a/src/polymarket_insider_tracker/alerter/formatter.py b/src/polymarket_insider_tracker/alerter/formatter.py index 9f69642..b38fe8b 100644 --- a/src/polymarket_insider_tracker/alerter/formatter.py +++ b/src/polymarket_insider_tracker/alerter/formatter.py @@ -190,10 +190,11 @@ class AlertFormatter: wallet_age_str = "" if assessment.fresh_wallet_signal: age_hours = assessment.fresh_wallet_signal.wallet_profile.age_hours - if age_hours < 1: - wallet_age_str = f" (Age: {int(age_hours * 60)}m)" - else: - wallet_age_str = f" (Age: {age_hours:.0f}h)" + if age_hours is not None: + if age_hours < 1: + wallet_age_str = f" (Age: {int(age_hours * 60)}m)" + else: + wallet_age_str = f" (Age: {age_hours:.0f}h)" fields: list[dict[str, object]] = [ { @@ -282,10 +283,11 @@ class AlertFormatter: wallet_line = f"*Wallet:* `{wallet_short}`" if assessment.fresh_wallet_signal: age_hours = assessment.fresh_wallet_signal.wallet_profile.age_hours - if age_hours < 1: - wallet_line += f" \\(Age: {int(age_hours * 60)}m\\)" - else: - wallet_line += f" \\(Age: {age_hours:.0f}h\\)" + if age_hours is not None: + if age_hours < 1: + wallet_line += f" \\(Age: {int(age_hours * 60)}m\\)" + else: + wallet_line += f" \\(Age: {age_hours:.0f}h\\)" lines.append(wallet_line) # Risk score @@ -366,10 +368,11 @@ class AlertFormatter: wallet_line = f"Wallet: {wallet_short}" if assessment.fresh_wallet_signal: age_hours = assessment.fresh_wallet_signal.wallet_profile.age_hours - if age_hours < 1: - wallet_line += f" (Age: {int(age_hours * 60)}m)" - else: - wallet_line += f" (Age: {age_hours:.0f}h)" + if age_hours is not None: + if age_hours < 1: + wallet_line += f" (Age: {int(age_hours * 60)}m)" + else: + wallet_line += f" (Age: {age_hours:.0f}h)" lines.append(wallet_line) # Risk diff --git a/src/polymarket_insider_tracker/alerter/history.py b/src/polymarket_insider_tracker/alerter/history.py index 7aee6c3..8c27647 100644 --- a/src/polymarket_insider_tracker/alerter/history.py +++ b/src/polymarket_insider_tracker/alerter/history.py @@ -362,7 +362,7 @@ class AlertHistory: start.timestamp(), end.timestamp(), ) - return count + return int(count) async def cleanup_old_alerts(self) -> int: """Remove alerts older than retention period. @@ -393,4 +393,4 @@ class AlertHistory: # Note: Individual alert records will expire via TTL # Wallet/market indexes will also expire via TTL logger.info(f"Cleaned up {removed} old alert references") - return removed + return int(removed) diff --git a/src/polymarket_insider_tracker/detector/scorer.py b/src/polymarket_insider_tracker/detector/scorer.py index 9e523d3..07691e5 100644 --- a/src/polymarket_insider_tracker/detector/scorer.py +++ b/src/polymarket_insider_tracker/detector/scorer.py @@ -268,7 +268,7 @@ class RiskScorer: """ key = f"{self._key_prefix}{wallet_address}:{market_id}" deleted = await self._redis.delete(key) - return deleted > 0 + return int(deleted) > 0 async def assess_batch(self, bundles: list[SignalBundle]) -> list[RiskAssessment]: """Assess multiple trade bundles. diff --git a/src/polymarket_insider_tracker/ingestor/clob_client.py b/src/polymarket_insider_tracker/ingestor/clob_client.py index d1643c6..86795fd 100644 --- a/src/polymarket_insider_tracker/ingestor/clob_client.py +++ b/src/polymarket_insider_tracker/ingestor/clob_client.py @@ -287,7 +287,8 @@ class ClobClient: try: response = self._client.get_midpoint(token_id) - return response.get("mid") + mid = response.get("mid") + return str(mid) if mid is not None else None except Exception as e: logger.warning("Failed to get midpoint for %s: %s", token_id, e) return None @@ -307,7 +308,8 @@ class ClobClient: try: response = self._client.get_price(token_id, side=side) - return response.get("price") + price = response.get("price") + return str(price) if price is not None else None except Exception as e: logger.warning("Failed to get %s price for %s: %s", side, token_id, e) return None @@ -321,7 +323,7 @@ class ClobClient: try: self._rate_limiter.acquire_sync() result = self._client.get_ok() - return result == "OK" + return str(result) == "OK" except Exception as e: logger.error("Health check failed: %s", e) return False @@ -334,7 +336,8 @@ class ClobClient: """ try: self._rate_limiter.acquire_sync() - return self._client.get_server_time() + result = self._client.get_server_time() + return int(result) if result is not None else None except Exception as e: logger.error("Failed to get server time: %s", e) return None diff --git a/src/polymarket_insider_tracker/ingestor/metadata_sync.py b/src/polymarket_insider_tracker/ingestor/metadata_sync.py index 8aa9293..52e87cd 100644 --- a/src/polymarket_insider_tracker/ingestor/metadata_sync.py +++ b/src/polymarket_insider_tracker/ingestor/metadata_sync.py @@ -351,7 +351,7 @@ class MarketMetadataSync: """ key = f"{self._key_prefix}{condition_id}" deleted = await self._redis.delete(key) - return deleted > 0 + return int(deleted) > 0 async def force_sync(self) -> None: """Force an immediate sync of all markets. diff --git a/src/polymarket_insider_tracker/ingestor/publisher.py b/src/polymarket_insider_tracker/ingestor/publisher.py index 217ee8f..8e8842f 100644 --- a/src/polymarket_insider_tracker/ingestor/publisher.py +++ b/src/polymarket_insider_tracker/ingestor/publisher.py @@ -186,9 +186,10 @@ class EventPublisher: The entry ID assigned by Redis. """ data = _serialize_trade_event(event) + # redis-py typing expects broader dict type than dict[str, str] entry_id = await self._redis.xadd( self._stream_name, - data, + data, # type: ignore[arg-type] maxlen=self._max_len, ) # entry_id may be bytes or str @@ -213,7 +214,8 @@ class EventPublisher: pipe = self._redis.pipeline() for event in events: data = _serialize_trade_event(event) - pipe.xadd(self._stream_name, data, maxlen=self._max_len) + # redis-py typing expects broader dict type than dict[str, str] + pipe.xadd(self._stream_name, data, maxlen=self._max_len) # type: ignore[arg-type] results = await pipe.execute() @@ -384,7 +386,8 @@ class EventPublisher: """ if not entry_ids: return 0 - return await self._redis.xack(self._stream_name, group_name, *entry_ids) + result = await self._redis.xack(self._stream_name, group_name, *entry_ids) + return int(result) async def get_stream_info(self) -> dict[str, Any]: """Get information about the stream. @@ -404,7 +407,8 @@ class EventPublisher: Returns: Number of entries in the stream. """ - return await self._redis.xlen(self._stream_name) + result = await self._redis.xlen(self._stream_name) + return int(result) async def trim_stream(self, max_len: int | None = None) -> int: """Trim the stream to a maximum length. @@ -416,4 +420,5 @@ class EventPublisher: Number of entries removed. """ length = max_len or self._max_len - return await self._redis.xtrim(self._stream_name, maxlen=length) + result = await self._redis.xtrim(self._stream_name, maxlen=length) + return int(result) diff --git a/src/polymarket_insider_tracker/profiler/funding.py b/src/polymarket_insider_tracker/profiler/funding.py index c8ef298..a688c2b 100644 --- a/src/polymarket_insider_tracker/profiler/funding.py +++ b/src/polymarket_insider_tracker/profiler/funding.py @@ -232,6 +232,7 @@ class FundingTracer: ) # Get logs with Transfer event filtering by recipient + # Note: web3 typing is overly restrictive for block params logs = await w3.eth.get_logs( { "address": AsyncWeb3.to_checksum_address(token_address), @@ -240,8 +241,8 @@ class FundingTracer: None, # from (any) padded_to, # to (target address) ], - "fromBlock": from_block, - "toBlock": to_block, + "fromBlock": from_block, # type: ignore[typeddict-item] + "toBlock": to_block, # type: ignore[typeddict-item] } ) @@ -317,7 +318,7 @@ class FundingTracer: origin_type="error", ) else: - chains[addr.lower()] = result # type: ignore[assignment] + chains[addr.lower()] = result return chains diff --git a/src/polymarket_insider_tracker/storage/repos.py b/src/polymarket_insider_tracker/storage/repos.py index 4c0308e..a9b7bdf 100644 --- a/src/polymarket_insider_tracker/storage/repos.py +++ b/src/polymarket_insider_tracker/storage/repos.py @@ -208,20 +208,20 @@ class WalletRepository: await self.session.execute(stmt) except Exception: # Fall back to SQLite upsert for testing - stmt = sqlite_insert(WalletProfileModel).values(**values, created_at=now) - stmt = stmt.on_conflict_do_update( + sqlite_stmt = sqlite_insert(WalletProfileModel).values(**values, created_at=now) + sqlite_stmt = sqlite_stmt.on_conflict_do_update( index_elements=["address"], set_={ - "nonce": stmt.excluded.nonce, - "first_seen_at": stmt.excluded.first_seen_at, - "is_fresh": stmt.excluded.is_fresh, - "matic_balance": stmt.excluded.matic_balance, - "usdc_balance": stmt.excluded.usdc_balance, - "analyzed_at": stmt.excluded.analyzed_at, - "updated_at": stmt.excluded.updated_at, + "nonce": sqlite_stmt.excluded.nonce, + "first_seen_at": sqlite_stmt.excluded.first_seen_at, + "is_fresh": sqlite_stmt.excluded.is_fresh, + "matic_balance": sqlite_stmt.excluded.matic_balance, + "usdc_balance": sqlite_stmt.excluded.usdc_balance, + "analyzed_at": sqlite_stmt.excluded.analyzed_at, + "updated_at": sqlite_stmt.excluded.updated_at, }, ) - await self.session.execute(stmt) + await self.session.execute(sqlite_stmt) await self.session.flush() return dto @@ -238,7 +238,8 @@ class WalletRepository: result = await self.session.execute( delete(WalletProfileModel).where(WalletProfileModel.address == address.lower()) ) - return result.rowcount > 0 + # SQLAlchemy Result does have rowcount but typing doesn't reflect it + return (result.rowcount or 0) > 0 # type: ignore[attr-defined] async def mark_stale(self, address: str) -> bool: """Mark a wallet profile as stale (soft delete). @@ -257,7 +258,8 @@ class WalletRepository: .where(WalletProfileModel.address == address.lower()) .values(analyzed_at=stale_time, updated_at=datetime.now(UTC)) ) - return result.rowcount > 0 + # SQLAlchemy Result does have rowcount but typing doesn't reflect it + return (result.rowcount or 0) > 0 # type: ignore[attr-defined] class FundingRepository: @@ -478,12 +480,12 @@ class RelationshipRepository: await self.session.execute(stmt) except Exception: # Fall back to SQLite upsert for testing - stmt = sqlite_insert(WalletRelationshipModel).values(**values) - stmt = stmt.on_conflict_do_update( + sqlite_stmt = sqlite_insert(WalletRelationshipModel).values(**values) + sqlite_stmt = sqlite_stmt.on_conflict_do_update( index_elements=["wallet_a", "wallet_b", "relationship_type"], - set_={"confidence": stmt.excluded.confidence}, + set_={"confidence": sqlite_stmt.excluded.confidence}, ) - await self.session.execute(stmt) + await self.session.execute(sqlite_stmt) await self.session.flush() return dto @@ -506,4 +508,5 @@ class RelationshipRepository: WalletRelationshipModel.relationship_type == relationship_type, ) ) - return result.rowcount > 0 + # SQLAlchemy Result does have rowcount but typing doesn't reflect it + return (result.rowcount or 0) > 0 # type: ignore[attr-defined]