From 562d641f9c55881effdc8aabe612620c057e2e0b Mon Sep 17 00:00:00 2001 From: Codex Date: Fri, 18 Sep 2026 02:59:18 +0200 Subject: [PATCH] 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 --- .../src/sentiment_engine/schemas/processed.py | 1 + .../src/sentiment_engine/signal/velocity.py | 23 +++++++++++++++---- 2 files changed, 19 insertions(+), 5 deletions(-) diff --git a/sentiment_engine/src/sentiment_engine/schemas/processed.py b/sentiment_engine/src/sentiment_engine/schemas/processed.py index c70b71f..2a41293 100644 --- a/sentiment_engine/src/sentiment_engine/schemas/processed.py +++ b/sentiment_engine/src/sentiment_engine/schemas/processed.py @@ -123,6 +123,7 @@ class ProcessedItem(BaseModel): processed_ts: float = Field(..., description="Processing completion timestamp") processing_latency_ms: float = Field(..., description="Total processing time") 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]]: return ( diff --git a/sentiment_engine/src/sentiment_engine/signal/velocity.py b/sentiment_engine/src/sentiment_engine/signal/velocity.py index d57158c..7749657 100644 --- a/sentiment_engine/src/sentiment_engine/signal/velocity.py +++ b/sentiment_engine/src/sentiment_engine/signal/velocity.py @@ -39,7 +39,7 @@ class VelocityComputer: def __init__(self): self.settings = get_settings() 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) self._hype_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: """Remove entries older than cutoff""" - while window and window[0]["ts"] < cutoff: - window.popleft() + # 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() + 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: """Compute weighted mention score: polarity_confidence × source_credibility × engagement_score""" @@ -166,7 +179,7 @@ class VelocityComputer: return 0.0 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: return 0.0 @@ -180,7 +193,7 @@ class VelocityComputer: return 0.0 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 if len(counts) >= 3: