File size: 3,438 Bytes
dbb1bf9 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 | /**
* Web Worker for heavy computational tasks (clustering & correlation analysis).
* Runs O(n²) Jaccard clustering and correlation detection off the main thread.
*
* All core logic is imported from src/services/analysis-core.ts
* to maintain a single source of truth.
*/
import {
clusterNewsCore,
analyzeCorrelationsCore,
type NewsItemCore,
type ClusteredEventCore,
type PredictionMarketCore,
type MarketDataCore,
type CorrelationSignalCore,
type SourceType,
type StreamSnapshot,
} from '@/services/analysis-core';
// Message types for worker communication
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[];
}
// Worker-local state (persists between messages)
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);
}
// Worker message handler
self.onmessage = (event: MessageEvent<WorkerMessage>) => {
const message = event.data;
switch (message.type) {
case 'cluster': {
// Deserialize dates (they come as strings over postMessage)
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': {
// Deserialize dates in clusters
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;
}
}
};
// Signal that worker is ready
self.postMessage({ type: 'ready' });
|