| |
| |
| |
| |
| |
| |
| |
|
|
| import { |
| clusterNewsCore, |
| analyzeCorrelationsCore, |
| type NewsItemCore, |
| type ClusteredEventCore, |
| type PredictionMarketCore, |
| type MarketDataCore, |
| type CorrelationSignalCore, |
| type SourceType, |
| type StreamSnapshot, |
| } from '@/services/analysis-core'; |
|
|
| |
| interface ClusterMessage { |
| type: 'cluster'; |
| id: string; |
| items: NewsItemCore[]; |
| sourceTiers: Record<string, number>; |
| } |
|
|
| interface CorrelationMessage { |
| type: 'correlation'; |
| id: string; |
| clusters: ClusteredEventCore[]; |
| predictions: PredictionMarketCore[]; |
| markets: MarketDataCore[]; |
| sourceTypes: Record<string, SourceType>; |
| } |
|
|
| interface ResetMessage { |
| type: 'reset'; |
| } |
|
|
| type WorkerMessage = ClusterMessage | CorrelationMessage | ResetMessage; |
|
|
| interface ClusterResult { |
| type: 'cluster-result'; |
| id: string; |
| clusters: ClusteredEventCore[]; |
| } |
|
|
| interface CorrelationResult { |
| type: 'correlation-result'; |
| id: string; |
| signals: CorrelationSignalCore[]; |
| } |
|
|
| |
| let previousSnapshot: StreamSnapshot | null = null; |
| const recentSignalKeys = new Set<string>(); |
|
|
| function isRecentDuplicate(key: string): boolean { |
| return recentSignalKeys.has(key); |
| } |
|
|
| function markSignalSeen(key: string): void { |
| recentSignalKeys.add(key); |
| setTimeout(() => recentSignalKeys.delete(key), 30 * 60 * 1000); |
| } |
|
|
| |
| self.onmessage = (event: MessageEvent<WorkerMessage>) => { |
| const message = event.data; |
|
|
| switch (message.type) { |
| case 'cluster': { |
| |
| const items = message.items.map(item => ({ |
| ...item, |
| pubDate: new Date(item.pubDate), |
| })); |
|
|
| const getSourceTier = (source: string): number => message.sourceTiers[source] ?? 4; |
| const clusters = clusterNewsCore(items, getSourceTier); |
|
|
| const result: ClusterResult = { |
| type: 'cluster-result', |
| id: message.id, |
| clusters, |
| }; |
| self.postMessage(result); |
| break; |
| } |
|
|
| case 'correlation': { |
| |
| const clusters = message.clusters.map(cluster => ({ |
| ...cluster, |
| firstSeen: new Date(cluster.firstSeen), |
| lastUpdated: new Date(cluster.lastUpdated), |
| allItems: cluster.allItems.map(item => ({ |
| ...item, |
| pubDate: new Date(item.pubDate), |
| })), |
| })); |
|
|
| const getSourceType = (source: string): SourceType => message.sourceTypes[source] ?? 'unknown'; |
|
|
| const { signals, snapshot } = analyzeCorrelationsCore( |
| clusters, |
| message.predictions, |
| message.markets, |
| previousSnapshot, |
| getSourceType, |
| isRecentDuplicate, |
| markSignalSeen |
| ); |
|
|
| previousSnapshot = snapshot; |
|
|
| const result: CorrelationResult = { |
| type: 'correlation-result', |
| id: message.id, |
| signals, |
| }; |
| self.postMessage(result); |
| break; |
| } |
|
|
| case 'reset': { |
| previousSnapshot = null; |
| recentSignalKeys.clear(); |
| break; |
| } |
| } |
| }; |
|
|
| |
| self.postMessage({ type: 'ready' }); |
|
|