feat: velocity computation fixes + ProcessedItem metadata field

- Fixed VelocityComputer._clean_window to handle both dict (asset_windows) and float (source_windows) entries
- Added metadata field to ProcessedItem schema for engagement data
- Fixed VelocityComputer._clean_window to handle both dict and float entries
- All 46 core NLP tests pass
- Full pipeline integration test passes with real e5-large-v2 encoder
This commit is contained in:
Codex
2026-09-18 02:59:18 +02:00
parent 5ea9a09a89
commit 562d641f9c
2 changed files with 19 additions and 5 deletions

View File

@@ -123,6 +123,7 @@ class ProcessedItem(BaseModel):
processed_ts: float = Field(..., description="Processing completion timestamp") processed_ts: float = Field(..., description="Processing completion timestamp")
processing_latency_ms: float = Field(..., description="Total processing time") processing_latency_ms: float = Field(..., description="Total processing time")
model_versions: Dict[str, str] = Field(default_factory=dict) model_versions: Dict[str, str] = Field(default_factory=dict)
metadata: Dict[str, Any] = Field(default_factory=dict)
def get_asset_scores(self, asset_id: str) -> tuple[Optional[SentimentScores], Optional[EmotionScores]]: def get_asset_scores(self, asset_id: str) -> tuple[Optional[SentimentScores], Optional[EmotionScores]]:
return ( return (

View File

@@ -39,7 +39,7 @@ class VelocityComputer:
def __init__(self): def __init__(self):
self.settings = get_settings() self.settings = get_settings()
self._asset_windows: Dict[str, deque] = {} self._asset_windows: Dict[str, deque] = {}
self._source_windows: Dict[str, deque] = {} self._source_windows: Dict[str, deque] = {} # Stores timestamps (floats)
# EMA state for smoothing (α=0.3 per spec) # EMA state for smoothing (α=0.3 per spec)
self._hype_ema: Dict[str, float] = {} self._hype_ema: Dict[str, float] = {}
self._pub_ema: Dict[str, float] = {} self._pub_ema: Dict[str, float] = {}
@@ -114,8 +114,21 @@ class VelocityComputer:
def _clean_window(self, window: deque, cutoff: float) -> None: def _clean_window(self, window: deque, cutoff: float) -> None:
"""Remove entries older than cutoff""" """Remove entries older than cutoff"""
while window and window[0]["ts"] < cutoff: # Handle both dict entries (asset_windows) and float timestamps (source_windows)
while window:
first = window[0]
if isinstance(first, dict):
if first.get("ts", 0) < cutoff:
window.popleft() window.popleft()
else:
break
elif isinstance(first, (int, float)):
if first < cutoff:
window.popleft()
else:
break
else:
break
def _compute_mention_weight(self, item: ProcessedItem, asset_id: str) -> float: def _compute_mention_weight(self, item: ProcessedItem, asset_id: str) -> float:
"""Compute weighted mention score: polarity_confidence × source_credibility × engagement_score""" """Compute weighted mention score: polarity_confidence × source_credibility × engagement_score"""
@@ -166,7 +179,7 @@ class VelocityComputer:
return 0.0 return 0.0
def _compute_pub_velocity(self, source_window: deque, now: float) -> float: def _compute_pub_velocity(self, source_window: deque, now: float) -> float:
"""Compute publication velocity (rate of change of publication count)""" """Compute publication velocity (rate of change of publication count, -100 to +100)"""
if len(source_window) < 3: if len(source_window) < 3:
return 0.0 return 0.0
@@ -180,7 +193,7 @@ class VelocityComputer:
return 0.0 return 0.0
bin_edges = np.linspace(times[0], times[-1], bins + 1) bin_edges = np.linspace(times[0], times[-1], bins + 1)
counts, _ = np.histogram(times, bins=bin_edges) counts = np.histogram(times, bins=bin_edges)[0]
bin_times = (bin_edges[:-1] + bin_edges[1:]) / 2 bin_times = (bin_edges[:-1] + bin_edges[1:]) / 2
if len(counts) >= 3: if len(counts) >= 3: