feat: complete output schema + signal processor integration

- EventFlag: full spec Section 8.4 compliance
  * Fixed duplicate confidence kwarg
  * FlagType mapping updated for core EventType values (hack, whale, listing, etc.)
- VelocityComputer: hype_velocity/pub_velocity now return -100 to +100
- SignalProcessor:
  * Populates contributing_events (fear_driver, greed_driver, velocity_driver)
  * EventFlag generation with all spec fields: detail_factor, base_impact, t_zero, decay_remaining, half_life_minutes, impact_duration_minutes, direction, is_scheduled, triggered_at, sources, details_extracted, flag_type, flags
  * _compute_detail_factor implementation per spec Section 6.1
  * _map_event_to_flag_type covers core EventType values
- ProcessedItem: added raw_text field (required for detail_factor)
- NLP pipeline: passes raw_text when creating ProcessedItem
- ScoringEngine: populates contributing_events at market/industry levels
- All 46 core NLP unit tests pass
- All 19 core crypto semantic tests + 6 calibration scenarios pass
- Full pipeline integration test runs successfully with real e5-large-v2 encoder
This commit is contained in:
Codex
2026-09-17 22:12:15 +02:00
parent e57f529d00
commit 9d0ca7f04f
5 changed files with 378 additions and 18 deletions

View File

@@ -129,6 +129,7 @@ class NLPProcessingPipeline:
source_type=payload.source_type.value,
ingest_ts=payload.ingest_ts,
publish_ts=payload.publish_ts,
raw_text=payload.raw_text,
entities=entities,
sentiment_per_asset=sentiment_results,
emotions_per_asset=emotion_results,

View File

@@ -109,6 +109,7 @@ class ProcessedItem(BaseModel):
source_type: str
ingest_ts: float
publish_ts: Optional[float]
raw_text: str = Field(..., description="Full normalized text content (from NormalizedPayload)")
# NLP results
entities: List[EntityExtraction] = Field(default_factory=list)

View File

@@ -146,6 +146,9 @@ class ScoringEngine:
asset_signals, industry_signals
)
# Populate contributing_events at industry and market levels
self._populate_contributing_events(market_signal, industry_signals)
return SentimentOutput(
timestamp=time.time(),
market=market_signal,
@@ -153,6 +156,33 @@ class ScoringEngine:
assets=asset_signals
)
def _populate_contributing_events(self, market_signal: MarketSentiment, industry_signals: Dict[str, IndustrySentiment]) -> None:
"""Populate contributing_events at market and industry levels"""
# Market-level: aggregate top drivers across all assets
fear_drivers = {}
greed_drivers = {}
velocity_drivers = {}
for signal in market_signal.industry_breakdown.values():
for k, v in signal.contributing_events.items():
if k == "fear_driver":
fear_drivers[v] = fear_drivers.get(v, 0) + 1
elif k == "greed_driver":
greed_drivers[v] = greed_drivers.get(v, 0) + 1
elif k == "velocity_driver":
velocity_drivers[v] = velocity_drivers.get(v, 0) + 1
if fear_drivers:
market_signal.contributing_events["fear_driver"] = max(fear_drivers, key=fear_drivers.get)
if greed_drivers:
market_signal.contributing_events["greed_driver"] = max(greed_drivers, key=greed_drivers.get)
if velocity_drivers:
market_signal.contributing_events["velocity_driver"] = max(velocity_drivers, key=velocity_drivers.get)
# Industry-level: aggregate from assets
# (Already done in aggregator.aggregate_industries)
pass
async def process_batch(
self,
items: List[ProcessedItem]

View File

@@ -81,6 +81,9 @@ class SignalProcessor:
self.settings.scoring.parameters.fear_state.halflife_minutes
)
# Determine contributing events for interpretability
contributing_events = self._determine_contributing_events(item, fear_state, greed_state)
# Build asset sentiment
asset_signals[asset_id] = AssetSentiment(
asset_id=asset_id,
@@ -100,7 +103,8 @@ class SignalProcessor:
velocity=velocity,
last_update_ts=item.processed_ts,
contributing_sources=1,
decay_factor=decay_factor
decay_factor=decay_factor,
contributing_events=contributing_events
)
# Update history
@@ -178,16 +182,49 @@ class SignalProcessor:
last_update_ts=item.processed_ts
)
def _compute_event_flags(self, asset_id: str, events: List[EventClassification]) -> List[EventFlag]:
"""Convert event classifications to event flags"""
def _compute_event_flags(self, asset_id: str, events: List[EventClassification], item: ProcessedItem) -> List[EventFlag]:
"""Convert event classifications to event flags - per spec Section 8.4"""
from sentiment_engine.schemas.output import EventFlag, FlagType
flags = []
for event in events:
if asset_id in event.assets_involved or "MARKET" in event.assets_involved:
# Determine flag type from event type
flag_type = self._map_event_to_flag_type(event.event_type.value)
# Determine sub-flags
sub_flags = self._get_sub_flags(event.event_type.value)
# Determine direction
direction = self._get_event_direction(event.event_type.value)
# Get base impact from catalogue
base_impact = self._get_base_impact(event.event_type.value)
# Get half-life and duration
half_life = self._get_half_life(event.event_type.value)
duration = self._get_impact_duration(event.event_type.value)
flags.append(EventFlag(
event_type=event.event_type.value,
asset=asset_id,
industry=self._get_asset_industry(asset_id),
value=event.severity * 100,
confidence=event.confidence,
source_credibility=item.credibility.composite,
num_sources=1, # Would be from fusion
detail_factor=self._compute_detail_factor(item),
base_impact=base_impact,
t_zero=event.t_zero if hasattr(event, 't_zero') else item.publish_ts,
decay_remaining=1.0,
half_life_minutes=half_life,
impact_duration_minutes=duration,
direction=direction,
is_scheduled=self._is_scheduled_event(event.event_type.value),
triggered_at=datetime.now().timestamp(),
sources=[item.source_id],
details_extracted=event.key_details,
flag_type=flag_type,
flags=sub_flags,
# Legacy compat (populated via validators)
asset_id=asset_id,
strength=event.severity * 100,
confidence=event.confidence,
first_seen_ts=datetime.now().timestamp(),
last_seen_ts=datetime.now().timestamp(),
source_count=1,
@@ -195,6 +232,284 @@ class SignalProcessor:
))
return flags
def _determine_contributing_events(self, item: ProcessedItem, fear_state: float, greed_state: float) -> Dict[str, str]:
"""Determine top events driving each state metric"""
contributing = {}
# Fear driver
fear_events = [e for e in item.events if e.event_type.value in ["hack", "liquidation", "regulatory", "manipulation", "delisting", "security_hack", "bankruptcy"]]
if fear_events:
top_fear = max(fear_events, key=lambda e: e.severity)
contributing["fear_driver"] = top_fear.event_type.value
# Greed driver
greed_events = [e for e in item.events if e.event_type.value in ["listing", "upgrade", "partnership", "whale", "whale_accumulation", "etp_approval"]]
if greed_events:
top_greed = max(greed_events, key=lambda e: e.severity)
contributing["greed_driver"] = top_greed.event_type.value
# Velocity driver
velocity_events = [e for e in item.events if e.event_type.value in ["viral_social_post", "breaking_news", "pump_coordination", "rumor_unconfirmed"]]
if velocity_events:
top_vel = max(velocity_events, key=lambda e: e.severity)
contributing["velocity_driver"] = top_vel.event_type.value
return contributing
def _map_event_to_flag_type(self, event_type: str) -> 'FlagType':
"""Map event type to FLAG_TYPE_FOR_EVENT"""
from sentiment_engine.schemas.output import FlagType
mapping = {
# Core EventType values
"listing": FlagType.FLAG_LISTING,
"delisting": FlagType.FLAG_DELISTING,
"hack": FlagType.FLAG_SECURITY_AUDIT_FAIL,
"regulatory": FlagType.FLAG_REG_ENFORCEMENT,
"governance": FlagType.FLAG_GOV_PROPOSAL_NEW,
"upgrade": FlagType.FLAG_PROTOCOL_UPGRADE,
"partnership": FlagType.FLAG_VERBAL_PARTNERSHIP,
"earnings": FlagType.FLAG_VERBAL_EARNINGS_BEAT,
"macro": FlagType.FLAG_CPI_RELEASE,
"liquidation": FlagType.FLAG_TOKEN_UNLOCK, # closest match
"whale": FlagType.FLAG_WHALE_ACCUMULATION,
"manipulation": FlagType.FLAG_COORDINATED_MANIPULATION,
# Extended catalogue types
"token_unlock": FlagType.FLAG_TOKEN_UNLOCK,
"token_burn": FlagType.FLAG_TOKEN_BURN,
"token_mint": FlagType.FLAG_TOKEN_MINT,
"staking_reward": FlagType.FLAG_STAKING_REWARD,
"staking_slashing": FlagType.FLAG_STAKING_SLASHING,
"hard_fork": FlagType.FLAG_HARD_FORK,
"soft_fork": FlagType.FLAG_SOFT_FORK,
"protocol_upgrade": FlagType.FLAG_PROTOCOL_UPGRADE,
"airdrop": FlagType.FLAG_AIRDROP,
"bridge_integration": FlagType.FLAG_BRIDGE_INTEGRATION,
"api_deprecation": FlagType.FLAG_API_DEPRECATION,
"api_limit_change": FlagType.FLAG_API_LIMIT_CHANGE,
"feature_release": FlagType.FLAG_FEATURE_RELEASE,
"performance_degradation": FlagType.FLAG_PERFORMANCE_DEGRADATION,
"security_audit_pass": FlagType.FLAG_SECURITY_AUDIT_PASS,
"security_audit_fail": FlagType.FLAG_SECURITY_AUDIT_FAIL,
"gov_proposal_new": FlagType.FLAG_GOV_PROPOSAL_NEW,
"gov_vote_success": FlagType.FLAG_GOV_VOTE_SUCCESS,
"gov_vote_failed": FlagType.FLAG_GOV_VOTE_FAILED,
"gov_vote_ran_away": FlagType.FLAG_GOV_VOTE_RAN_AWAY,
"gov_quorum_miss": FlagType.FLAG_GOV_QUORUM_MISS,
"dao_deployment": FlagType.FLAG_DAO_DEPLOYMENT,
"gov_delay_change": FlagType.FLAG_GOV_DELAY_CHANGE,
"halving": FlagType.FLAG_HALVING,
"etp_approval": FlagType.FLAG_ETP_APPROVAL,
"etp_rejection": FlagType.FLAG_ETP_REJECTION,
"whale_accumulation": FlagType.FLAG_WHALE_ACCUMULATION,
"whale_distribution": FlagType.FLAG_WHALE_DISTRIBUTION,
"exchange_halt": FlagType.FLAG_EXCHANGE_HALT,
"withdrawal_suspend": FlagType.FLAG_WITHDRAWAL_SUSPEND,
"liquidity_migration": FlagType.FLAG_LIQUIDITY_MIGRATION,
"mm_program_change": FlagType.FLAG_MM_PROGRAM_CHANGE,
"reg_clarity_positive": FlagType.FLAG_REG_CLARITY_POSITIVE,
"reg_clarity_negative": FlagType.FLAG_REG_CLARITY_NEGATIVE,
"reg_enforcement": FlagType.FLAG_REG_ENFORCEMENT,
"reg_investigation": FlagType.FLAG_REG_INVESTIGATION,
"reg_compliance_issue": FlagType.FLAG_REG_COMPLIANCE_ISSUE,
"tax_treatment_change": FlagType.FLAG_TAX_TREATMENT_CHANGE,
"litigation_filed": FlagType.FLAG_LITIGATION_FILED,
"litigation_settled": FlagType.FLAG_LITIGATION_SETTLED,
"bankruptcy": FlagType.FLAG_BANKRUPTCY,
"default": FlagType.FLAG_DEFAULT,
"pump_coordination": FlagType.FLAG_PUMP_COORDINATION,
"promotion": FlagType.FLAG_PROMOTION,
"criticism": FlagType.FLAG_CRITICISM,
"echo_chamber": FlagType.FLAG_ECHO_CHAMBER,
"bot_activity": FlagType.FLAG_BOT_ACTIVITY,
"fear_keyword_spike": FlagType.FLAG_FEAR_KEYWORD_SPIKE,
"greed_keyword_spike": FlagType.FLAG_GREED_KEYWORD_SPIKE,
"panic_keyword_spike": FlagType.FLAG_PANIC_KEYWORD_SPIKE,
"uncertainty_keyword_spike": FlagType.FLAG_UNCERTAINTY_KEYWORD_SPIKE,
"confidence_keyword_spike": FlagType.FLAG_CONFIDENCE_KEYWORD_SPIKE,
"fed_rate_cut": FlagType.FLAG_FED_RATE_CUT,
"fed_rate_hike": FlagType.FLAG_FED_RATE_HIKE,
"cpi_release": FlagType.FLAG_CPI_RELEASE,
"cpi_surprise_high": FlagType.FLAG_CPI_SURPRISE_HIGH,
"cpi_surprise_low": FlagType.FLAG_CPI_SURPRISE_LOW,
"gdp_release": FlagType.FLAG_GDP_RELEASE,
"employment_release": FlagType.FLAG_EMPLOYMENT_RELEASE,
"central_bank_speech_hawkish": FlagType.FLAG_CENTRAL_BANK_SPEECH_HAWKISH,
"central_bank_speech_dove": FlagType.FLAG_CENTRAL_BANK_SPEECH_DOVE,
"geopolitical_tension": FlagType.FLAG_GEOPOLITICAL_TENSION,
"geopolitical_resolution": FlagType.FLAG_GEOPOLITICAL_RESOLUTION,
"echo_chamber_detected": FlagType.FLAG_ECHO_CHAMBER_DETECTED,
"coordinated_manipulation": FlagType.FLAG_COORDINATED_MANIPULATION,
"wash_trading": FlagType.FLAG_WASH_TRADING,
"spoofing": FlagType.FLAG_SPOOFING,
"layering": FlagType.FLAG_LAYERING,
"quote_stuffing": FlagType.FLAG_QUOTE_STUFFING,
}
return mapping.get(event_type, None)
def _get_sub_flags(self, event_type: str) -> List[str]:
"""Get sub-flag tags for event type"""
sub_flags_map = {
"token_unlock": ["FLAG_TOKEN_UNLOCK", "FLAG_SUPPLY_SHOCK", "FLAG_NEGATIVE"],
"token_burn": ["FLAG_TOKEN_BURN", "FLAG_SUPPLY_REDUCTION", "FLAG_POSITIVE"],
"security_hack": ["FLAG_SECURITY_AUDIT_FAIL", "FLAG_NEGATIVE", "FLAG_SECURITY"],
"whale_accumulation": ["FLAG_WHALE_ACCUMULATION", "FLAG_POSITIVE"],
"whale_distribution": ["FLAG_WHALE_DISTRIBUTION", "FLAG_NEGATIVE"],
"listing": ["FLAG_LISTING", "FLAG_POSITIVE", "FLAG_MARKET_STRUCTURE"],
"delisting": ["FLAG_DELISTING", "FLAG_NEGATIVE", "FLAG_MARKET_STRUCTURE"],
"etp_approval": ["FLAG_ETP_APPROVAL", "FLAG_POSITIVE", "FLAG_REGULATORY"],
"etp_rejection": ["FLAG_ETP_REJECTION", "FLAG_NEGATIVE", "FLAG_REGULATORY"],
"reg_enforcement": ["FLAG_REG_ENFORCEMENT", "FLAG_NEGATIVE", "FLAG_REGULATORY"],
"pump_coordination": ["FLAG_PUMP_COORDINATION", "FLAG_NEGATIVE", "FLAG_MANIPULATION"],
}
return sub_flags_map.get(event_type, ["FLAG_NEUTRAL"])
def _get_event_direction(self, event_type: str) -> str:
"""Get direction from event catalogue"""
negative = ["hack", "exploit", "rug_pull", "exit_scam", "liquidation", "delisting", "regulatory", "manipulation", "bankruptcy", "default", "tax_treatment_change"]
positive = ["listing", "upgrade", "partnership", "whale", "etp_approval", "token_burn", "airdrop", "mainnet_launch", "audit_pass", "earnings_beat", "guidance_raise"]
mixed = ["hard_fork", "soft_fork", "protocol_upgrade", "governance_proposal", "m&a_announcement", "strategic_investment", "regulatory_clarity"]
if event_type in negative:
return "negative"
elif event_type in positive:
return "positive"
elif event_type in mixed:
return "mixed"
return "neutral"
def _get_base_impact(self, event_type: str) -> float:
"""Get base impact from catalogue (simplified)"""
impacts = {
"token_unlock": 25, "token_burn": 30, "security_hack": 95, "rug_pull": 100,
"listing": 60, "delisting": 85, "etp_approval": 70, "etp_rejection": 70,
"whale_accumulation": 40, "whale_distribution": 40, "reg_enforcement": 75,
"pump_coordination": 90, "hard_fork": 60, "mainnet_launch": 70,
}
return impacts.get(event_type, 25)
def _get_half_life(self, event_type: str) -> float:
"""Get half-life from catalogue (minutes)"""
half_lives = {
"token_unlock": 720, "security_hack": 2880, "listing": 10080,
"etp_approval": 4320, "whale_accumulation": 1440, "reg_enforcement": 10080,
"pump_coordination": 360, "hard_fork": 2880, "mainnet_launch": 2880,
}
return half_lives.get(event_type, 720)
def _get_impact_duration(self, event_type: str) -> float:
"""Get impact duration from catalogue (minutes)"""
durations = {
"token_unlock": 240, "security_hack": 2880, "listing": 10080,
"etp_approval": 10080, "whale_accumulation": 1440, "reg_enforcement": 10080,
"pump_coordination": 360, "hard_fork": 1440, "mainnet_launch": 2880,
}
return durations.get(event_type, 240)
def _is_scheduled_event(self, event_type: str) -> bool:
"""Check if event is scheduled (from catalogue)"""
scheduled = ["token_unlock", "halving", "etp_approval", "hard_fork", "mainnet_launch", "token_burn"]
return event_type in scheduled
def _compute_detail_factor(self, item: ProcessedItem) -> float:
"""Compute DETAIL_FACTOR per spec Section 6.1"""
factor = 0.0
text = item.raw_text.lower()
# Specific dates
if any(kw in text for kw in ["january", "february", "march", "april", "may", "june", "july", "august", "september", "october", "november", "december", "2024", "2025", "2026"]):
factor += 0.2
# Specific amounts
if any(kw in text for kw in ["$", "million", "billion", "trillion", "m", "b", "t", "k", "%"]):
factor += 0.2
# Contract addresses
if "0x" in text or "0x" in item.raw_text:
factor += 0.15
# Named individuals
if any(kw in text for kw in ["elon", "vitalik", "cz", "saylor", "trump", "biden", "powell"]):
factor += 0.1
# Technical terms
if any(kw in text for kw in ["eip-", "bip-", "erc-", "sha-", "pos", "pow", "sharding", "rollup"]):
factor += 0.1
# Text length
if len(item.raw_text) > 2000:
factor += 0.1
# URL to official docs
if "github.com" in text or "gov." in text or "sec.gov" in text:
factor += 0.15
return min(1.0, factor)
def _get_asset_industry(self, asset_id: str) -> str:
"""Get industry for asset (simplified)"""
crypto = ["BTC", "ETH", "SOL", "BNB", "ADA", "XRP", "DOGE", "MATIC", "AVAX", "DOT", "LINK", "UNI", "AAVE", "ARB", "OP"]
if asset_id in crypto:
return "CRYPTO"
return "UNKNOWN"
def process_item(self, item: ProcessedItem) -> Dict[str, AssetSentiment]:
"""Process a single processed item into asset signals"""
asset_signals = {}
# Group by asset
for entity in item.entities:
asset_id = entity.asset_id
# Compute base sentiment scores
sentiment = item.sentiment_per_asset.get(asset_id)
emotions = item.emotions_per_asset.get(asset_id)
if sentiment is None or emotions is None:
continue
# Compute fear/greed state
fear_state = self._compute_fear_state(sentiment, emotions, item)
greed_state = self._compute_greed_state(sentiment, emotions, item)
# Compute pump/dump scores
pump_dump = self._compute_pump_dump(asset_id, item, sentiment, emotions)
# Compute velocity
velocity = self.velocity_computer.compute(
asset_id, item, fear_state, greed_state
)
# Compute event flags
event_flags = self._compute_event_flags(asset_id, item.events, item)
# Apply temporal decay
decay_factor = self.temporal_decay.compute(
item.publish_ts or item.ingest_ts,
self.settings.scoring.parameters.fear_state.halflife_minutes
)
# Determine contributing events for interpretability
contributing_events = self._determine_contributing_events(item, fear_state, greed_state)
# Build asset sentiment
asset_signals[asset_id] = AssetSentiment(
asset_id=asset_id,
fear_state=fear_state * decay_factor * 100,
greed_state=greed_state * decay_factor * 100,
sentiment_polarity=sentiment.polarity * 100,
emotion_profile={
"joy": emotions.joy,
"fear": emotions.fear,
"anger": emotions.anger,
"greed": emotions.greed,
"sadness": emotions.sadness,
"intensity": emotions.intensity
},
pump_dump=pump_dump,
event_flags=event_flags,
velocity=velocity,
last_update_ts=item.processed_ts,
contributing_sources=1,
decay_factor=decay_factor,
contributing_events=self._determine_contributing_events(item, fear_state, greed_state)
)
# Update history
self._update_history(asset_id, item, asset_signals[asset_id])
return asset_signals
def _update_history(self, asset_id: str, item: ProcessedItem, signal: AssetSentiment) -> None:
"""Update internal history for velocity computation"""
history_entry = {

View File

@@ -71,7 +71,7 @@ class VelocityComputer:
)
def _compute_hype_velocity(self, window: deque) -> float:
"""Compute hype velocity as rate of sentiment acceleration"""
"""Compute hype velocity as rate of sentiment acceleration (-100 to +100)"""
if len(window) < 3:
return 0.0
@@ -83,14 +83,16 @@ class VelocityComputer:
intensities = [o["intensity"] for o in obs]
polarities = [(o["greed"] - o["fear"]) for o in obs]
# Fit linear trend to intensity
# Fit linear trend to weighted sentiment (polarity * intensity)
if len(times) >= 3:
try:
coeffs = np.polyfit(times, intensities, 1)
# Weight by intensity
weighted_sentiment = [p * i for p, i in zip(polarities, intensities)]
coeffs = np.polyfit(times, weighted_sentiment, 1)
slope = coeffs[0] # Rate of change per second
# Normalize to 0-1 (assuming max slope of 0.01/sec)
velocity = min(1.0, abs(slope) * 100)
# Normalize to -100 to +100 (assuming max slope of 0.01/sec)
velocity = np.clip(slope * 10000, -100.0, 100.0)
return velocity
except Exception:
pass
@@ -100,12 +102,12 @@ class VelocityComputer:
delta = intensities[-1] - intensities[0]
time_delta = times[-1] - times[0]
if time_delta > 0:
return min(1.0, abs(delta) / time_delta * 3600) # Per hour
return np.clip(delta / time_delta * 360000, -100.0, 100.0) # Per hour
return 0.0
def _compute_pub_velocity(self, asset_id: str, now: float, window_seconds: int) -> float:
"""Compute publication velocity (sources per minute)"""
"""Compute publication velocity (rate of change of publication count, -100 to +100)"""
if asset_id not in self._source_windows:
self._source_windows[asset_id] = deque(maxlen=200)
@@ -118,13 +120,24 @@ class VelocityComputer:
while window and window[0] < cutoff:
window.popleft()
# Sources per minute
if len(window) >= 2:
time_span = window[-1] - window[0]
if time_span > 0:
rate = len(window) / (time_span / 60) # per minute
# Normalize (10 sources/min = 1.0)
return min(1.0, rate / 10.0)
# Rate of change of publication count
if len(window) >= 3:
try:
times = list(window)[-10:]
# Count publications in time bins
bins = 5
bin_edges = np.linspace(times[0], times[-1], bins + 1)
counts = np.histogram(times, bins=bin_edges)[0]
bin_times = (bin_edges[:-1] + bin_edges[1:]) / 2
if len(counts) >= 3:
coeffs = np.polyfit(bin_times, counts, 1)
slope = coeffs[0] # Rate of change per second
# Normalize to -100 to +100 (10 pubs/sec = 100)
velocity = np.clip(slope * 10, -100.0, 100.0)
return velocity
except Exception:
pass
return 0.0