feat: Implement RSS matrix fan-in with per-source artifacts and deterministic merge (#454)

Implements #436: RSS matrix fan-in script with per-source artifacts and deterministic merge. Includes docs from #439: contracts, runbook, and ADR for matrix crawl fan-in. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>

Juan Manuel Servera committed Jun 13, 2026 at 17:06 UTC 711870d70d9b89259fecae558ae37e732c82012c
8 files changed +2932
data/archive/2026-W21-techcrunch.json new
+410
@@ -0,0 +1,410 @@
1 +{
2 + "week": "2026-W21",
3 + "source": "techcrunch",
4 + "crawled_at": "2026-05-21T12:39:56.134051Z",
5 + "articles": [
6 + {
7 + "title": "Scammers are abusing an internal Microsoft account to send spam links",
8 + "url": "https://techcrunch.com/2026/05/21/scammers-are-abusing-an-internal-microsoft-account-to-send-spam/",
9 + "published_at": "2026-05-21T11:42:57Z",
10 + "categories": [
11 + "Security",
12 + "cyberattacks",
13 + "cybersecurity",
14 + "Microsoft",
15 + "phishing",
16 + "scam"
17 + ],
18 + "summary": "The loophole allows spammers and scammers to send emails from a legitimate Microsoft email address typically used for sending genuine account alerts.",
19 + "github_links": [],
20 + "entities": [
21 + "Scammers",
22 + "Microsoft"
23 + ],
24 + "relevance_score": 0.2
25 + },
26 + {
27 + "title": "Beauty booking startup Fresha hits $1 billion valuation with KKR backing",
28 + "url": "https://techcrunch.com/2026/05/21/booking-platform-fresha-announces-80m-investment-unicorn-valuation/",
29 + "published_at": "2026-05-21T11:00:00Z",
30 + "categories": [
31 + "Startups",
32 + "Venture",
33 + "SaaS"
34 + ],
35 + "summary": "Beauty and wellness booking marketplace Fresha says it has raised $80 million investment from KKR’s Next Generation Technology Growth fund, KKR's growth equity arm.",
36 + "github_links": [],
37 + "entities": [
38 + "Beauty",
39 + "Fresha",
40 + "KKR"
41 + ],
42 + "relevance_score": 0.4
43 + },
44 + {
45 + "title": "Truecaller gets into the eSIM business to diversify its revenue streams",
46 + "url": "https://techcrunch.com/2026/05/20/truecaller-gets-into-the-esim-business-to-diversify-its-revenue-streams/",
47 + "published_at": "2026-05-21T06:00:00Z",
48 + "categories": [
49 + "Apps",
50 + "esim",
51 + "travel",
52 + "Truecaller"
53 + ],
54 + "summary": "The company said its plans will range from 1 GB over 7 days to 20 GB over 30 days. Initially, the launch will make the eSIM product available in 29 countries.",
55 + "github_links": [],
56 + "entities": [
57 + "Truecaller"
58 + ],
59 + "relevance_score": 0.2
60 + },
61 + {
62 + "title": "General Catalyst just led a $63M bet on India’s travel payments market",
63 + "url": "https://techcrunch.com/2026/05/20/indian-travel-fintech-scapia-more-than-doubles-valuation-to-over-500m-in-a-year/",
64 + "published_at": "2026-05-21T05:52:38Z",
65 + "categories": [
66 + "Fintech",
67 + "Startups",
68 + "General Catalyst",
69 + "Peak XV Partners",
70 + "Scapia",
71 + "Z47"
72 + ],
73 + "summary": "Scapia, an Indian startup that combines travel booking with co-branded credit cards and mobile payments, said the deal doubles its valuation.",
74 + "github_links": [],
75 + "entities": [
76 + "General",
77 + "Catalyst",
78 + "India"
79 + ],
80 + "relevance_score": 0.6
81 + },
82 + {
83 + "title": "Imperagen raises £5 million to use quantum physics, AI on enzyme engineering",
84 + "url": "https://techcrunch.com/2026/05/20/imperagen-raises-5-million-to-redefine-enzyme-engineering/",
85 + "published_at": "2026-05-21T04:00:00Z",
86 + "categories": [
87 + "Startups",
88 + "Biotech & Health",
89 + "Venture",
90 + "biotech"
91 + ],
92 + "summary": "Biotech company Imperagen announced on Thursday a £5 million ($6.7 million) seed round led by PXN Ventures, with participation from IQ Capital and Northern Gritstone.",
93 + "github_links": [],
94 + "entities": [
95 + "Imperagen",
96 + "AI"
97 + ],
98 + "relevance_score": 0.8
99 + },
100 + {
101 + "title": "Jensen Huang says he’s found a ‘brand new’ $200B market for Nvidia",
102 + "url": "https://techcrunch.com/2026/05/20/jensen-huang-says-hes-found-a-brand-new-200b-market-for-nvidia/",
103 + "published_at": "2026-05-21T00:28:31Z",
104 + "categories": [
105 + "AI",
106 + "Enterprise",
107 + "TC",
108 + "cpus",
109 + "nvidia"
110 + ],
111 + "summary": "The next big thing for Nvidia will be CPUs for AI agents, $200 billion worth, CEO Jensen Huang predicts.",
112 + "github_links": [],
113 + "entities": [
114 + "Jensen",
115 + "Huang",
116 + "Nvidia"
117 + ],
118 + "relevance_score": 0.4
119 + },
120 + {
121 + "title": "Anthropic says it’s about to have its first profitable quarter",
122 + "url": "https://techcrunch.com/2026/05/20/anthropic-says-its-about-to-have-its-first-profitable-quarter/",
123 + "published_at": "2026-05-21T00:21:21Z",
124 + "categories": [
125 + "AI",
126 + "Anthropic",
127 + "Claude",
128 + "OpenAI"
129 + ],
130 + "summary": "Anthropic has told its investors that it will more than double revenue to around $10.9 billion in its second quarter.",
131 + "github_links": [],
132 + "entities": [
133 + "Anthropic"
134 + ],
135 + "relevance_score": 0.2
136 + },
137 + {
138 + "title": "The SpaceX IPO filing is filled with AI bets, Starship dreams, and Elon Musk at the center",
139 + "url": "https://techcrunch.com/2026/05/20/the-spacex-ipo-filing-ai-bets-starship-dreams-elon-musk/",
140 + "published_at": "2026-05-20T23:03:02Z",
141 + "categories": [
142 + "Space",
143 + "Transportation",
144 + "Elon Musk",
145 + "IPOs",
146 + "SpaceX",
147 + "rockets",
148 + "xAI",
149 + "Starlink"
150 + ],
151 + "summary": "SpaceX has finally made the contents of its IPO filing public, weeks ahead of what is expected to be the largest IPO ever and one that will make Musk the CEO, CTO, and chairman of the board.",
152 + "github_links": [],
153 + "entities": [
154 + "SpaceX",
155 + "IPO",
156 + "AI",
157 + "Starship",
158 + "Elon",
159 + "Musk"
160 + ],
161 + "relevance_score": 0.2
162 + },
163 + {
164 + "title": "Clouted wants to take the guesswork out of making short videos go viral",
165 + "url": "https://techcrunch.com/2026/05/20/clouted-wants-to-take-the-guesswork-out-of-making-short-videos-go-viral/",
166 + "published_at": "2026-05-20T22:30:45Z",
167 + "categories": [
168 + "AI",
169 + "Media & Entertainment",
170 + "Startups",
171 + "Clouted",
172 + "Marketing",
173 + "slow ventures",
174 + "social media"
175 + ],
176 + "summary": "The video clipping startup raised a $7 million seed round led by Slow Ventures.",
177 + "github_links": [],
178 + "entities": [
179 + "Clouted"
180 + ],
181 + "relevance_score": 0.4
182 + },
183 + {
184 + "title": "xAI burned $6.4B last year — SpaceX’s IPO filing shows why the spending is far from over",
185 + "url": "https://techcrunch.com/2026/05/20/xai-burned-6-4b-last-year-spacexs-ipo-filing-shows-why-the-spending-is-far-from-over/",
186 + "published_at": "2026-05-20T22:26:08Z",
187 + "categories": [
188 + "AI",
189 + "Space",
190 + "Elon Musk",
191 + "Grok",
192 + "SpaceX",
193 + "spacex ipo",
194 + "X"
195 + ],
196 + "summary": "SpaceX's IPO filing reveals xAI lost $6.4 billion in 2025 while planning a massive Grok expansion — offering the first public look at Elon Musk's AI financials and more details about his ambitions.",
197 + "github_links": [],
198 + "entities": [
199 + "SpaceX",
200 + "IPO"
201 + ],
202 + "relevance_score": 0.2
203 + },
204 + {
205 + "title": "Nvidia posts another record quarter, reveals $43B of holdings in startups",
206 + "url": "https://techcrunch.com/2026/05/20/nvidia-posts-another-record-quarter-reveals-43-billion-of-holdings-in-startups/",
207 + "published_at": "2026-05-20T22:03:51Z",
208 + "categories": [
209 + "AI",
210 + "earnings",
211 + "Jensen Huang",
212 + "nvidia"
213 + ],
214 + "summary": "Nvidia announced another record revenue figure after market close on Wednesday, but forecasted that revenue growth would slow in the following quarter.",
215 + "github_links": [],
216 + "entities": [
217 + "Nvidia"
218 + ],
219 + "relevance_score": 0.4
220 + },
221 + {
222 + "title": "Musk’s xAI is being sued over its data center generators — now it’s buying $2.8B more",
223 + "url": "https://techcrunch.com/2026/05/20/musks-xai-is-being-sued-over-its-data-center-generators-now-its-buying-2-8b-more/",
224 + "published_at": "2026-05-20T21:55:49Z",
225 + "categories": [
226 + "AI",
227 + "Climate",
228 + "air pollution",
229 + "data centers",
230 + "Elon Musk",
231 + "natural gas",
232 + "SpaceX",
233 + "spacex ipo",
234 + "spacexai",
235 + "xAI"
236 + ],
237 + "summary": "Elon Musk's xAI said it will buy $2.8 billion worth of natural gas turbines over the next three years, according to SpaceX's IPO filing.",
238 + "github_links": [],
239 + "entities": [
240 + "Musk"
241 + ],
242 + "relevance_score": 0.2
243 + },
244 + {
245 + "title": "Anthropic will pay xAI $1.25B per month for compute",
246 + "url": "https://techcrunch.com/2026/05/20/anthropic-will-pay-xai-1-25-billion-per-month-for-compute/",
247 + "published_at": "2026-05-20T21:29:22Z",
248 + "categories": [
249 + "AI",
250 + "Anthropic",
251 + "colossus",
252 + "data centers",
253 + "SpaceX",
254 + "xAI"
255 + ],
256 + "summary": "Elon Musk's xAI surprised the AI world when it made a deal to sell compute to Anthropic. Now we know how much it's worth.",
257 + "github_links": [],
258 + "entities": [
259 + "Anthropic"
260 + ],
261 + "relevance_score": 0.4
262 + },
263 + {
264 + "title": "Sam Altman makes ‘mic drop’ offer to every Y Combinator startup",
265 + "url": "https://techcrunch.com/2026/05/20/sam-altman-makes-mic-drop-offer-to-every-y-combinator-startup/",
266 + "published_at": "2026-05-20T21:23:02Z",
267 + "categories": [
268 + "Startups",
269 + "Venture",
270 + "OpenAI",
271 + "sam altman",
272 + "Y Combinator"
273 + ],
274 + "summary": "Altman offered to have OpenAI invest in every single startup in this Y Combinator class: tokens for equity.",
275 + "github_links": [],
276 + "entities": [
277 + "Sam",
278 + "Altman",
279 + "Combinator"
280 + ],
281 + "relevance_score": 0.4
282 + },
283 + {
284 + "title": "You don’t need to be an AI startup to raise. Lucra has $20M to prove it.",
285 + "url": "https://techcrunch.com/video/you-dont-need-to-be-an-ai-startup-to-raise-lucra-has-20m-to-prove-it/",
286 + "published_at": "2026-05-20T21:21:22Z",
287 + "categories": [
288 + "Startups",
289 + "AI startup",
290 + "ark invest",
291 + "Cathie Wood",
292 + "Equity podcast",
293 + "Lucra",
294 + "startup fundraising",
295 + "venture capital"
296 + ],
297 + "summary": "Slapping &#8220;AI&#8221; on your&#160;startup’s&#160;pitch deck is&#160;basically table&#160;stakes right now. When a founder&#160;raised $20 million from Cathie Wood&#8217;s ARK Invest&#160;for an eSports&#160;gamification&#160;loyalty startup without those two letters in the spotlight, it got us wondering how the conversation even started&#160;—&#160;especially&#160;when ARK had already been burned by a company&#160;operating&#160;in the same space.&#160; On this episode of TechCrunch&#821...",
298 + "github_links": [],
299 + "entities": [
300 + "You",
301 + "AI",
302 + "Lucra"
303 + ],
304 + "relevance_score": 0.6
305 + },
306 + {
307 + "title": "Microsoft’s carbon-removal plans aren’t dead after all",
308 + "url": "https://techcrunch.com/2026/05/20/microsofts-carbon-removal-plans-arent-dead-after-all/",
309 + "published_at": "2026-05-20T20:30:24Z",
310 + "categories": [
311 + "Climate",
312 + "biogas",
313 + "carbon credits",
314 + "carbon removal",
315 + "Exclusive",
316 + "Microsoft"
317 + ],
318 + "summary": "Microsoft is responsible for over 90% of the carbon-removal market, and reports suggested the company was pausing purchases entirely. This new deal should help assuage the fears of CDR startups.",
319 + "github_links": [],
320 + "entities": [
321 + "Microsoft"
322 + ],
323 + "relevance_score": 0.2
324 + },
325 + {
326 + "title": "OpenAI claims it solved an 80-year-old math problem — for real this time",
327 + "url": "https://techcrunch.com/2026/05/20/openai-claims-it-solved-an-80-year-old-math-problem-for-real-this-time/",
328 + "published_at": "2026-05-20T20:28:27Z",
329 + "categories": [
330 + "AI",
331 + "ChatGPT",
332 + "erdos problems",
333 + "OpenAI",
334 + "reasoning models"
335 + ],
336 + "summary": "OpenAI claims its reasoning model disproved a geometry conjecture unsolved since 1946 — and this time, the mathematicians who exposed its last embarrassing claim are backing it up.",
337 + "github_links": [],
338 + "entities": [
339 + "OpenAI"
340 + ],
341 + "relevance_score": 0.6
342 + },
343 + {
344 + "title": "IrisGo, a startup backed by Andrew Ng, looks to become the AI desktop buddy you never knew you needed",
345 + "url": "https://techcrunch.com/2026/05/20/irisgo-a-startup-backed-by-andrew-ng-looks-to-become-the-ai-desktop-buddy-you-never-knew-you-needed/",
346 + "published_at": "2026-05-20T19:47:20Z",
347 + "categories": [
348 + "AI",
349 + "andrew ng",
350 + "google brain",
351 + "IrisGo"
352 + ],
353 + "summary": "Initially billed as an \"AI butler,\" Iris watches what happens on a user's desktop and automatically learns how to do tasks for them, its co-founder says.",
354 + "github_links": [],
355 + "entities": [
356 + "IrisGo",
357 + "Andrew",
358 + "Ng",
359 + "AI"
360 + ],
361 + "relevance_score": 0.4
362 + },
363 + {
364 + "title": "Tesla’s Full Self-Driving software is creeping into Europe",
365 + "url": "https://techcrunch.com/2026/05/20/teslas-full-self-driving-software-is-creeping-into-europe/",
366 + "published_at": "2026-05-20T18:32:06Z",
367 + "categories": [
368 + "Transportation",
369 + "autonomous vehicles",
370 + "Europe",
371 + "EVs",
372 + "Tesla",
373 + "Tesla FSD"
374 + ],
375 + "summary": "First came the Netherlands, now it's Lithuania. And more European countries appear to be in the queue for Tesla's driver-assistance system.",
376 + "github_links": [],
377 + "entities": [
378 + "Tesla",
379 + "Full",
380 + "Self",
381 + "Driving",
382 + "Europe"
383 + ],
384 + "relevance_score": 0.0
385 + },
386 + {
387 + "title": "Airbnb gets into hotels, expands AI for host onboarding and customer support",
388 + "url": "https://techcrunch.com/2026/05/20/airbnb-gets-into-hotels-expands-ai-for-host-onboarding-and-customer-support/",
389 + "published_at": "2026-05-20T18:14:04Z",
390 + "categories": [
391 + "Apps",
392 + "Airbnb",
393 + "customer support",
394 + "hotel bookings"
395 + ],
396 + "summary": "Airbnb will soon let you book luggage storage and car rental services on its app.",
397 + "github_links": [],
398 + "entities": [
399 + "Airbnb",
400 + "AI"
401 + ],
402 + "relevance_score": 0.4
403 + }
404 + ],
405 + "metadata": {
406 + "total_articles": 20,
407 + "relevant_articles": 12,
408 + "github_links_found": 0
409 + }
410 +}
\ No newline at end of file
data/archive/2026-W22-techcrunch.json new
+420
@@ -0,0 +1,420 @@
1 +{
2 + "week": "2026-W22",
3 + "source": "techcrunch",
4 + "crawled_at": "2026-05-25T11:55:51.494764Z",
5 + "articles": [
6 + {
7 + "title": "Everyone is navigating AI security in real time — even Google",
8 + "url": "https://techcrunch.com/2026/05/24/everyone-is-navigating-ai-security-in-real-time-even-google/",
9 + "published_at": "2026-05-24T21:39:21Z",
10 + "categories": [
11 + "AI",
12 + "TC"
13 + ],
14 + "summary": "We're in the transition period -- all of us.",
15 + "github_links": [],
16 + "entities": [
17 + "Everyone",
18 + "AI",
19 + "Google"
20 + ],
21 + "relevance_score": 0.2
22 + },
23 + {
24 + "title": "Xreal, Google’s smartglasses partner, thinks it has finally mastered this notoriously tricky industry",
25 + "url": "https://techcrunch.com/2026/05/24/xreal-googles-smartglasses-partner-thinks-it-has-finally-mastered-this-notoriously-tricky-industry/",
26 + "published_at": "2026-05-24T19:00:00Z",
27 + "categories": [
28 + "Hardware",
29 + "Google",
30 + "Google I/O",
31 + "AI",
32 + "SMART Glasses",
33 + "XReal"
34 + ],
35 + "summary": "Chi Xu, the founder and CEO of XREAL, thinks the smart glasses business has finally reached a turning point.",
36 + "github_links": [],
37 + "entities": [
38 + "Xreal",
39 + "Google"
40 + ],
41 + "relevance_score": 0.2
42 + },
43 + {
44 + "title": "6 kitchen gadgets that make adulting feel easier",
45 + "url": "https://techcrunch.com/2026/05/24/6-kitchen-gadgets-that-make-adulting-feel-easier/",
46 + "published_at": "2026-05-24T17:00:00Z",
47 + "categories": [
48 + "Hardware",
49 + "Gadgets",
50 + "Cooking Gadgets",
51 + "Kitchen gadgets",
52 + "robot chef",
53 + "evergreens"
54 + ],
55 + "summary": "From a robot stirring your soup to a bread machine that kneads your dough, here are 6 gadgets that may make you feel like you’ve won adulthood.",
56 + "github_links": [],
57 + "entities": [],
58 + "relevance_score": 0.0
59 + },
60 + {
61 + "title": "TechCrunch Mobility: Robotaxi reality check",
62 + "url": "https://techcrunch.com/2026/05/24/techcrunch-mobility-robotaxi-reality-check/",
63 + "published_at": "2026-05-24T16:05:00Z",
64 + "categories": [
65 + "Transportation",
66 + "Elon Musk",
67 + "SpaceX",
68 + "Tesla",
69 + "xAI",
70 + "Waymo",
71 + "robotaxi",
72 + "techcrunch mobility"
73 + ],
74 + "summary": "Welcome back to TechCrunch Mobility — your central hub for news and insights on the future of transportation.",
75 + "github_links": [],
76 + "entities": [
77 + "TechCrunch",
78 + "Mobility",
79 + "Robotaxi"
80 + ],
81 + "relevance_score": 0.2
82 + },
83 + {
84 + "title": "I tried Amazon’s Bee wearable and am both intrigued and slightly creeped out",
85 + "url": "https://techcrunch.com/2026/05/24/i-tried-amazons-bee-wearable-and-am-both-intrigued-and-slightly-creeped-out/",
86 + "published_at": "2026-05-24T15:00:00Z",
87 + "categories": [
88 + "AI",
89 + "Gadgets",
90 + "Amazon",
91 + "AI hardware",
92 + "Bee wearable"
93 + ],
94 + "summary": "Like other AI wearables, Amazon's Bee offers an odd combination of convenience and privacy anxiety.",
95 + "github_links": [],
96 + "entities": [
97 + "Amazon",
98 + "Bee"
99 + ],
100 + "relevance_score": 0.2
101 + },
102 + {
103 + "title": "The Dreamie alarm clock got me to stop using my phone in bed",
104 + "url": "https://techcrunch.com/2026/05/24/the-dreamie-alarm-clock-got-me-to-stop-using-my-phone-in-bed/",
105 + "published_at": "2026-05-24T13:00:00Z",
106 + "categories": [
107 + "Hardware",
108 + "dreamie"
109 + ],
110 + "summary": "What sets Dreamie apart from all of the other fancy alarm clocks is laughably simple: It can play podcasts.",
111 + "github_links": [],
112 + "entities": [
113 + "Dreamie"
114 + ],
115 + "relevance_score": 0.0
116 + },
117 + {
118 + "title": "SolarSquare in talks to raise up to $60M as India’s rooftop solar market draws major VC interest",
119 + "url": "https://techcrunch.com/2026/05/23/solarsquare-in-talks-to-raise-up-to-60m-as-indias-rooftop-solar-market-draws-major-vc-interest/",
120 + "published_at": "2026-05-23T20:03:13Z",
121 + "categories": [
122 + "Climate",
123 + "Startups",
124 + "b capital",
125 + "Elevation Capital",
126 + "Exclusive",
127 + "lightspeed venture partners",
128 + "SolarSquare"
129 + ],
130 + "summary": "SolarSquare could be valued at up to $500 million in the financing expected to close next month.",
131 + "github_links": [],
132 + "entities": [
133 + "SolarSquare",
134 + "India",
135 + "VC"
136 + ],
137 + "relevance_score": 0.6
138 + },
139 + {
140 + "title": "These special phone and app features can help protect you from spyware",
141 + "url": "https://techcrunch.com/2026/05/23/you-dont-have-to-click-anything-to-get-hacked-anymore-heres-how-to-fight-back/",
142 + "published_at": "2026-05-23T16:00:00Z",
143 + "categories": [
144 + "Security",
145 + "Android",
146 + "Apple",
147 + "Google",
148 + "hackers",
149 + "hacking",
150 + "WhatsApp",
151 + "Spyware",
152 + "Meta",
153 + "cybersecurity",
154 + "NSO Group",
155 + "Intellexa",
156 + "Paragon Solutions"
157 + ],
158 + "summary": "Apple, Meta, and Google offer special security modes that provide your devices more secure against targeted spyware attacks. Here are how those modes work, what they do, and how to switch them on.",
159 + "github_links": [],
160 + "entities": [],
161 + "relevance_score": 0.4
162 + },
163 + {
164 + "title": "Ferrari is using IBM’s AI to create F1 superfans",
165 + "url": "https://techcrunch.com/2026/05/23/ferrari-is-using-ai-to-create-f1-superfans/",
166 + "published_at": "2026-05-23T15:08:00Z",
167 + "categories": [
168 + "Transportation",
169 + "AI",
170 + "IBM",
171 + "Apps",
172 + "formula one",
173 + "Exclusive",
174 + "f1"
175 + ],
176 + "summary": "IBM and Scuderia Ferrari HP take TechCrunch inside how they are redefining the fan experience.",
177 + "github_links": [],
178 + "entities": [
179 + "Ferrari",
180 + "IBM",
181 + "AI",
182 + "F1"
183 + ],
184 + "relevance_score": 0.2
185 + },
186 + {
187 + "title": "Nuclear startup Deep Fission says it’s going public, again, and I have questions",
188 + "url": "https://techcrunch.com/2026/05/23/nuclear-startup-deep-fission-says-its-going-public-again-and-i-have-questions/",
189 + "published_at": "2026-05-23T14:50:00Z",
190 + "categories": [
191 + "Climate",
192 + "IPO",
193 + "nuclear power",
194 + "Deep Fission",
195 + "nuclear fission"
196 + ],
197 + "summary": "Deep Fission is seeking an IPO that could raise $157 million, though investors may have trouble buying the nuclear startup's story.",
198 + "github_links": [],
199 + "entities": [
200 + "Nuclear",
201 + "Deep",
202 + "Fission"
203 + ],
204 + "relevance_score": 0.4
205 + },
206 + {
207 + "title": "Elon Musk has given up on solar power (on Earth)",
208 + "url": "https://techcrunch.com/2026/05/23/elon-musk-has-given-up-on-solar-power-on-earth/",
209 + "published_at": "2026-05-23T13:00:00Z",
210 + "categories": [
211 + "AI",
212 + "Climate",
213 + "Analysis",
214 + "data centers",
215 + "Elon Musk",
216 + "Solar Power",
217 + "SpaceX",
218 + "Tesla",
219 + "xAI"
220 + ],
221 + "summary": "Elon Muks's xAI has gone all in on natural gas, while SpaceX is obsessed with orbital data centers. What happened to the \"solar-electric economy\" he promised?",
222 + "github_links": [],
223 + "entities": [
224 + "Elon",
225 + "Musk",
226 + "Earth"
227 + ],
228 + "relevance_score": 0.2
229 + },
230 + {
231 + "title": "Peec, one of Berlin’s rising startups, more than doubled annualized revenue in months to $10M, sources say",
232 + "url": "https://techcrunch.com/2026/05/23/peec-one-of-berlins-rising-startups-more-than-doubled-annualized-revenue-in-months-to-10m-sources-say/",
233 + "published_at": "2026-05-23T07:01:00Z",
234 + "categories": [
235 + "Startups",
236 + "Venture",
237 + "search marketing",
238 + "Antler",
239 + "peec ai"
240 + ],
241 + "summary": "Peec, which helps brands track their presence in AI searches, offers proof of a key trend among European startups.",
242 + "github_links": [],
243 + "entities": [
244 + "Peec",
245 + "Berlin"
246 + ],
247 + "relevance_score": 0.4
248 + },
249 + {
250 + "title": "AI is being used to resurrect the voices of dead pilots",
251 + "url": "https://techcrunch.com/2026/05/22/ai-is-being-used-to-resurrect-the-voices-of-dead-pilots/",
252 + "published_at": "2026-05-22T23:03:33Z",
253 + "categories": [
254 + "AI",
255 + "Transportation",
256 + "crash",
257 + "flight",
258 + "In Brief",
259 + "NTSB"
260 + ],
261 + "summary": "People used AI on a spectrogram image of cockpit recordings to reconstruct them, forcing the NTSB to temporarily block access to its docket system.",
262 + "github_links": [],
263 + "entities": [
264 + "AI"
265 + ],
266 + "relevance_score": 0.2
267 + },
268 + {
269 + "title": "SpaceX launches Starship V3 for the first time, but loses booster on return",
270 + "url": "https://techcrunch.com/2026/05/22/spacex-launches-starship-v3-for-the-first-time-but-loses-booster-on-return/",
271 + "published_at": "2026-05-22T22:55:57Z",
272 + "categories": [
273 + "Space",
274 + "TC",
275 + "SpaceX",
276 + "Starlink",
277 + "Starship",
278 + "starship v3"
279 + ],
280 + "summary": "The company had a mostly successful first launch of its upgraded Starship V3, which it needs to power its many ambitious goals in the years to come.",
281 + "github_links": [],
282 + "entities": [
283 + "SpaceX",
284 + "Starship",
285 + "V3"
286 + ],
287 + "relevance_score": 0.0
288 + },
289 + {
290 + "title": "Blue Origin cleared to fly New Glenn mega-rocket after April mishap",
291 + "url": "https://techcrunch.com/2026/05/22/blue-origin-cleared-to-fly-new-glenn-mega-rocket-after-april-mishap/",
292 + "published_at": "2026-05-22T21:37:17Z",
293 + "categories": [
294 + "Space",
295 + "Blue Origin",
296 + "In Brief",
297 + "new glenn"
298 + ],
299 + "summary": "Jeff Bezos' rocket company confirmed an engine failure led to the loss of an AST SpaceMobile satellite last month, but offered little detail.",
300 + "github_links": [],
301 + "entities": [
302 + "Blue",
303 + "Origin",
304 + "Glenn",
305 + "April"
306 + ],
307 + "relevance_score": 0.4
308 + },
309 + {
310 + "title": "Google goes for the glitter with disco-ball icons: ‘Are y’all sure you still want this?’",
311 + "url": "https://techcrunch.com/2026/05/22/google-goes-for-the-glitter-with-disco-ball-icons-are-yall-sure-you-still-want-this/",
312 + "published_at": "2026-05-22T21:02:36Z",
313 + "categories": [
314 + "AI",
315 + "Apps",
316 + "ai icons",
317 + "Android",
318 + "disco ball",
319 + "Google",
320 + "icons",
321 + "PIXEL",
322 + "Spotify"
323 + ],
324 + "summary": "You can now disco ball-ify your entire Pixel home screen, says Google.",
325 + "github_links": [],
326 + "entities": [
327 + "Google"
328 + ],
329 + "relevance_score": 0.2
330 + },
331 + {
332 + "title": "How VCs and founders use inflated ‘ARR’ to crown AI startups",
333 + "url": "https://techcrunch.com/2026/05/22/how-vcs-and-founders-use-inflated-arr-to-kingmake-ai-startups/",
334 + "published_at": "2026-05-22T20:40:48Z",
335 + "categories": [
336 + "AI",
337 + "Startups",
338 + "Venture",
339 + "annual recurring revenue",
340 + "Exclusive",
341 + "Valuations"
342 + ],
343 + "summary": "Some AI startups are stretching traditional revenue metrics when talking about progress publicly. And their investors are fully aware.",
344 + "github_links": [],
345 + "entities": [
346 + "VCs",
347 + "ARR",
348 + "AI"
349 + ],
350 + "relevance_score": 0.4
351 + },
352 + {
353 + "title": "Kash Patel’s clothing brand website shut down after reports it was hacked",
354 + "url": "https://techcrunch.com/2026/05/22/kash-patels-clothing-brand-website-shut-down-after-reports-it-was-hacked/",
355 + "published_at": "2026-05-22T16:28:06Z",
356 + "categories": [
357 + "Security",
358 + "cybersecurity",
359 + "FBI",
360 + "hackers",
361 + "hacking",
362 + "In Brief",
363 + "infostealer",
364 + "Kash Patel",
365 + "malware"
366 + ],
367 + "summary": "According to users on X, the website was hijacked by hackers in an attempt to trick visitors into installing malware.",
368 + "github_links": [],
369 + "entities": [
370 + "Kash",
371 + "Patel"
372 + ],
373 + "relevance_score": 0.0
374 + },
375 + {
376 + "title": "Apple says Epic lawsuit shouldn’t reshape App Store rules for all developers",
377 + "url": "https://techcrunch.com/2026/05/22/apple-says-epic-lawsuit-shouldnt-reshape-app-store-rules-for-all-developers/",
378 + "published_at": "2026-05-22T16:27:57Z",
379 + "categories": [
380 + "Apps",
381 + "Commerce",
382 + "Apple",
383 + "Court",
384 + "Epic Games",
385 + "lawsuit"
386 + ],
387 + "summary": "Apple is asking the Supreme Court to narrow the App Store injunction won by Epic Games and overturn the court’s contempt ruling over external payment fees.",
388 + "github_links": [],
389 + "entities": [
390 + "Apple",
391 + "Epic",
392 + "App",
393 + "Store"
394 + ],
395 + "relevance_score": 0.2
396 + },
397 + {
398 + "title": "Spotify’s AI bet: more of everything, less of what you want",
399 + "url": "https://techcrunch.com/2026/05/22/spotifys-ai-bet-more-of-everything-less-of-what-you-want/",
400 + "published_at": "2026-05-22T16:18:10Z",
401 + "categories": [
402 + "Apps",
403 + "AI audio",
404 + "Spotify"
405 + ],
406 + "summary": "Spotify has released a bunch of AI-powered tools that nudge users to create more content. It can be a bit much.",
407 + "github_links": [],
408 + "entities": [
409 + "Spotify",
410 + "AI"
411 + ],
412 + "relevance_score": 0.2
413 + }
414 + ],
415 + "metadata": {
416 + "total_articles": 20,
417 + "relevant_articles": 6,
418 + "github_links_found": 0
419 + }
420 +}
\ No newline at end of file
docs/decisions/adr-matrix-crawl-fan-in.md new
+119
@@ -0,0 +1,119 @@
1 +# ADR: Matrix Crawl Fan-In Architecture
2 +
3 +**Status:** Accepted
4 +**Date:** 2026-06-13
5 +**Decision makers:** Leela (Lead/Architect)
6 +**Related issues:** #331, #333, #356, #435, #436, #437, #438, #439
7 +
8 +---
9 +
10 +## Context
11 +
12 +SquadScope collects weekly GitHub repository signals and external news via RSS, then runs AI-powered analysis to produce trend summaries. As the system grows (more sources, more repos, longer analysis), the question arose: should we parallelize crawling with GitHub Actions matrix jobs?
13 +
14 +Investigation (documented in the [PRD](../processed/PRD-matrix-crawl-map-reduce-analysis.md)) revealed:
15 +
16 +1. RSS collection is already fast (~1s for 5 feeds) — matrix overhead would exceed actual work time.
17 +2. GitHub API crawling is rate-limit-constrained, not compute-constrained — parallelism risks quota exhaustion without proportional speedup.
18 +3. Analysis duration (28+ minutes) and context growth (112k+ tokens) are the actual bottlenecks.
19 +4. Downstream consumers expect canonical artifact formats regardless of collection topology.
20 +
21 +---
22 +
23 +## Decision
24 +
25 +We adopt a **staged, evidence-gated architecture** with three key design choices:
26 +
27 +### 1. Keep monolithic crawl as default; matrix is opt-in
28 +
29 +**Choice:** The default crawl topology remains monolithic (single GitHub crawl process + in-process RSS fetching). Matrix fan-out is enabled only when measurable triggers fire.
30 +
31 +**Alternatives considered:**
32 +- **Always-on matrix:** Rejected. Adds runner setup cost, artifact I/O overhead, cache merge complexity, and rate-limit instability for no measured benefit at current scale.
33 +- **Hybrid from day one:** Deferred. Premature complexity without evidence of need.
34 +
35 +**Triggers for RSS matrix:**
36 +- RSS p95 > 60 seconds
37 +- Source count > 10
38 +- Source requires independent credentials or isolation
39 +
40 +**Triggers for GitHub matrix:**
41 +- Shard experiment proves ≥25% speedup
42 +- No more than 10% API-call growth
43 +- Zero secondary-rate-limit regression
44 +
45 +### 2. Deterministic fan-in with strict contracts
46 +
47 +**Choice:** All matrix legs produce per-shard artifacts with schema validation, shared run context, and checksums. A fan-in job deterministically merges them into canonical payloads that are byte-stable for the same inputs.
48 +
49 +**Alternatives considered:**
50 +- **Append-only merge (no ordering):** Rejected. Non-deterministic output breaks caching, diffing, and downstream idempotency.
51 +- **Last-writer-wins:** Rejected. Data loss risk when shards overlap.
52 +- **Event-stream merge:** Over-engineered for batch workflows. Deferred.
53 +
54 +**Key contract properties:**
55 +- Deterministic ordering (repos by `full_name`, articles by `(source_id, url)`)
56 +- Idempotent: same inputs → same outputs
57 +- Fail-closed on required artifacts; graceful degradation on optional sources
58 +- Schema-versioned with forward-compatible rejection
59 +
60 +### 3. Map/reduce as analysis experiment, not crawl optimization
61 +
62 +**Choice:** Map/reduce is positioned as an analysis-quality tool (reducing LLM context, improving citations) — not as a crawl-speed mechanism. It runs in dry-run mode only until evidence validates quality parity.
63 +
64 +**Alternatives considered:**
65 +- **Map/reduce for crawl parallelism:** Rejected. Crawl is I/O-bound (API rate limits), not compute-bound. Parallelism doesn't help.
66 +- **Immediate production map/reduce:** Rejected. Must prove quality parity with single-pass analysis before promoting output.
67 +- **Vector DB / embedding approach:** Non-goal for MVP. No external paid infrastructure.
68 +
69 +---
70 +
71 +## Consequences
72 +
73 +### Positive
74 +
75 +- **No premature complexity:** Pipeline stays simple until evidence justifies change.
76 +- **Safe experimentation:** Dry-run and candidate-only modes allow testing without risk.
77 +- **Downstream stability:** Canonical artifact contracts isolate consumers from topology changes.
78 +- **Measurable rollout:** Clear acceptance criteria prevent subjective "good enough" decisions.
79 +- **Deterministic validation:** Byte-stable output enables automated regression testing.
80 +
81 +### Negative
82 +
83 +- **Delayed parallelism:** If source count grows rapidly, we'll need to implement matrix mode reactively rather than having it pre-built.
84 +- **Experiment overhead:** Shard experiments require dedicated runs comparing baseline vs. variant.
85 +- **Contract maintenance:** Schema versioning adds development overhead for artifact format changes.
86 +
87 +### Risks and mitigations
88 +
89 +| Risk | Mitigation |
90 +|------|-----------|
91 +| Triggers never fire, matrix code rots | Review triggers quarterly; remove dead code if unused for 6 months |
92 +| Fan-in non-determinism edge cases | Automated QA gates (#438) with fixture-based determinism tests |
93 +| Map/reduce quality never matches single-pass | Keep single-pass as default; map/reduce stays experimental until proven |
94 +| Schema version proliferation | Strict one-version-active policy; old schemas rejected immediately |
95 +
96 +---
97 +
98 +## Implementation Roadmap
99 +
100 +| Phase | Issue | Status | Description |
101 +|-------|-------|--------|-------------|
102 +| 0 | #333 | In progress | Define fan-in validation path and run-context schema |
103 +| 1 | #435 | Planned | GitHub shard experiment with guardrails |
104 +| 1 | #436 | Planned | RSS per-source artifacts + deterministic merge |
105 +| 1 | #437 | Planned | Observability metrics wiring |
106 +| 2 | #438 | Planned | Automated QA gates |
107 +| 2 | #439 | This PR | Documentation package (contracts, runbook, ADR) |
108 +| 3 | — | Future | Evidence-based rollout decision (enable or archive) |
109 +
110 +---
111 +
112 +## References
113 +
114 +- PRD: [Matrix Crawl and Map/Reduce Analysis](../processed/PRD-matrix-crawl-map-reduce-analysis.md)
115 +- Fan-in contracts: [docs/matrix-crawl-fan-in-contracts.md](../matrix-crawl-fan-in-contracts.md)
116 +- Runbook: [docs/matrix-crawl-runbook.md](../matrix-crawl-runbook.md)
117 +- Issue #356: Triage meta-issue for matrix crawl & map/reduce
118 +- Issue #331: Define map/reduce analysis promotion path
119 +- Issue #333: Define crawl matrix readiness and fan-in validation path
docs/matrix-crawl-fan-in-contracts.md new
+227
@@ -0,0 +1,227 @@
1 +# Matrix Crawl Fan-In Contracts
2 +
3 +**Status:** Reference specification
4 +**Related issues:** #333, #435, #436, #437, #438, #439
5 +**PRD:** [docs/processed/PRD-matrix-crawl-map-reduce-analysis.md](processed/PRD-matrix-crawl-map-reduce-analysis.md)
6 +
7 +---
8 +
9 +## Overview
10 +
11 +This document specifies the fan-in contracts for SquadScope's matrix crawl architecture. Fan-in is the merge step that combines per-shard or per-source artifacts into canonical downstream payloads. These contracts ensure deterministic, reproducible output regardless of whether collection ran as a monolithic process or as parallel matrix legs.
12 +
13 +---
14 +
15 +## Shared Run Context Schema
16 +
17 +Every crawl leg, fan-in validator, mapper, and reducer receives the same immutable run context. No matrix leg may compute its own time window from the local wall clock.
18 +
19 +```json
20 +{
21 + "schema_version": "run_context_v1",
22 + "run_id": "<week>-<sha256-prefix>",
23 + "week": "2026-W23",
24 + "since": "2026-06-01T00:00:00Z",
25 + "until": "2026-06-08T00:00:00Z",
26 + "source_config_checksum": "sha256:<hex>",
27 + "topic_config_checksum": "sha256:<hex>",
28 + "code_sha": "<sha256-of-relevant-source-files>",
29 + "created_at": "2026-06-05T17:42:56Z"
30 +}
31 +```
32 +
33 +### Required fields
34 +
35 +| Field | Type | Description |
36 +|-------|------|-------------|
37 +| `schema_version` | string | Must be `run_context_v1`. Fan-in rejects mismatches. |
38 +| `run_id` | string | Stable identifier: `{week}-{sha256-prefix}`. |
39 +| `week` | string | ISO week in `YYYY-WNN` format. |
40 +| `since` | string (ISO-8601) | Inclusive start of the collection window. |
41 +| `until` | string (ISO-8601) | Exclusive end of the collection window. |
42 +| `source_config_checksum` | string | SHA-256 of the source configuration file. |
43 +| `topic_config_checksum` | string | SHA-256 of `squadscope.topic.yml`. |
44 +| `code_sha` | string | SHA-256 of relevant pipeline source files. |
45 +| `created_at` | string (ISO-8601) | When the run context was generated. |
46 +
47 +---
48 +
49 +## Per-Source RSS Artifact Schema
50 +
51 +Each RSS matrix leg emits exactly one artifact per source:
52 +
53 +```json
54 +{
55 + "schema_version": 2,
56 + "source_id": "techcrunch",
57 + "run_context": { "...shared run context..." },
58 + "status": "success | partial | failed",
59 + "articles": [
60 + {
61 + "url": "https://...",
62 + "title": "...",
63 + "published": "2026-06-03T10:00:00Z",
64 + "source": "techcrunch",
65 + "relevance_score": 0.85,
66 + "github_urls": ["https://github.com/..."],
67 + "entities": ["Company A", "Project B"]
68 + }
69 + ],
70 + "metrics": {
71 + "articles_fetched": 20,
72 + "articles_relevant": 12,
73 + "fetch_duration_ms": 1200,
74 + "retries": 0
75 + },
76 + "error": null,
77 + "checksum": "sha256:<hex-of-articles-array>"
78 +}
79 +```
80 +
81 +### Validation rules
82 +
83 +- `schema_version` must equal `2` (current canonical version).
84 +- `run_context.week` must match the fan-in job's expected week.
85 +- `run_context.source_config_checksum` must match across all legs.
86 +- `checksum` must match the SHA-256 of the serialized `articles` array (sorted by URL, deterministic JSON).
87 +- `status: "failed"` artifacts carry no `articles` and must include an `error` object.
88 +
89 +---
90 +
91 +## GitHub Crawl Shard Artifact Schema
92 +
93 +Each GitHub shard leg (when enabled via experiment) emits:
94 +
95 +```json
96 +{
97 + "schema_version": "github_shard_v1",
98 + "shard_id": "new_repos | trending_repos | topic_primary | topic_secondary",
99 + "run_context": { "...shared run context..." },
100 + "repositories": [
101 + {
102 + "full_name": "owner/repo",
103 + "stars": 1234,
104 + "stars_gained": 45,
105 + "description": "...",
106 + "language": "Python",
107 + "topics": ["ai", "ml"],
108 + "created_at": "2026-01-15T...",
109 + "pushed_at": "2026-06-02T..."
110 + }
111 + ],
112 + "api_metrics": {
113 + "calls_made": 112,
114 + "cache_hits": 45,
115 + "cache_misses": 67,
116 + "secondary_rate_limit_events": 0,
117 + "search_api_remaining": 24
118 + },
119 + "checksum": "sha256:<hex>"
120 +}
121 +```
122 +
123 +---
124 +
125 +## Fan-In Merge Rules
126 +
127 +### Ordering guarantees
128 +
129 +1. **Deterministic repository ordering:** Repositories are sorted by `full_name` (lexicographic, case-sensitive). Ties are not possible since `full_name` is unique.
130 +2. **Deterministic article ordering:** Articles are sorted by `(source_id, url)` tuple.
131 +3. **Stable output:** Same input artifacts → byte-identical canonical output (excluding the `merged_at` timestamp, which is excluded from checksum computation).
132 +
133 +### Deduplication
134 +
135 +- **Repositories:** Deduplicated by `full_name`. If the same repo appears in multiple shards, the entry with the highest `stars_gained` is kept.
136 +- **Articles:** Deduplicated by normalized URL (scheme + host + path, query params stripped). First occurrence by source priority order wins.
137 +
138 +### Idempotency
139 +
140 +- Rerunning fan-in with the same input artifacts produces identical output.
141 +- Fan-in does not mutate input artifacts.
142 +- Reruns do not double-count articles, repos, retries, or restored stale artifacts.
143 +- The `run_id` in the run context serves as the idempotency key; the same `run_id` always maps to the same canonical output given the same inputs.
144 +
145 +### Failure policy
146 +
147 +| Scenario | Behavior |
148 +|----------|----------|
149 +| Required GitHub shard missing or invalid | **Fail closed.** No canonical output produced. Pipeline halts before analysis. |
150 +| Optional RSS source failed | **Degrade.** Produce canonical output if minimum-source policy passes (≥3/5 sources succeed). Record explicit warnings. |
151 +| Schema version mismatch | **Reject artifact.** Treat as missing. Apply required/optional rules above. |
152 +| Window/checksum mismatch | **Reject artifact.** Log validation error with details. |
153 +| All legs failed | **Fail closed.** Emit diagnostics artifact only. |
154 +
155 +### Minimum-source policy
156 +
157 +For RSS fan-in, the merge proceeds if:
158 +- At least 60% of configured sources report `status: "success"` or `status: "partial"`.
159 +- All `status: "partial"` sources still have ≥1 valid article.
160 +- The canonical output includes a `warnings` array documenting degraded sources.
161 +
162 +---
163 +
164 +## Canonical Output Schema
165 +
166 +Fan-in produces exactly two canonical artifacts consumed by downstream analysis:
167 +
168 +### `data/raw/{week}.json` (GitHub)
169 +
170 +The existing monolithic format, produced identically whether from one crawl process or merged shards.
171 +
172 +### `data/raw/{week}-external-news.json` (RSS/News)
173 +
174 +```json
175 +{
176 + "schema_version": 2,
177 + "run_context": { "..." },
178 + "sources": {
179 + "techcrunch": { "status": "success", "articles_count": 12 },
180 + "nvidia_blog": { "status": "success", "articles_count": 5 },
181 + "huggingface": { "status": "failed", "error": "timeout" }
182 + },
183 + "articles": [ "...merged, deduplicated, sorted..." ],
184 + "metrics": {
185 + "total_articles": 39,
186 + "total_relevant": 23,
187 + "sources_succeeded": 4,
188 + "sources_failed": 1
189 + },
190 + "warnings": ["huggingface: fetch timeout after 15s"],
191 + "merged_at": "2026-06-05T18:00:00Z",
192 + "checksum": "sha256:<hex-of-articles>"
193 +}
194 +```
195 +
196 +---
197 +
198 +## Contract Versioning
199 +
200 +- Schema versions use `<domain>_v<N>` format (e.g., `run_context_v1`, `github_shard_v1`).
201 +- Breaking changes increment the version number.
202 +- Fan-in rejects artifacts with unknown or mismatched schema versions.
203 +- The `schema_version` field is required in every artifact; omission is treated as a validation failure.
204 +
205 +---
206 +
207 +## Downstream Consumers
208 +
209 +The following components depend on canonical fan-in output and must NOT need to know whether collection was monolithic or matrix-based:
210 +
211 +- `scripts/correlate.py` — correlation analysis
212 +- `scripts/render_press_context.py` — press context rendering
213 +- `scripts/generate_content.py` — AI analysis
214 +- `scripts/map_reduce_dry_run.py` — map/reduce dry-run
215 +- `scripts/analysis_gate.py` — quality gate validation
216 +- Publishing workflow steps
217 +
218 +---
219 +
220 +## References
221 +
222 +- PRD: [Matrix Crawl and Map/Reduce Analysis](processed/PRD-matrix-crawl-map-reduce-analysis.md)
223 +- Issue #333: Define crawl matrix readiness and fan-in validation path
224 +- Issue #435: Run GitHub crawl shard experiment
225 +- Issue #436: Implement RSS matrix fan-in
226 +- Issue #437: Wire observability metrics
227 +- Issue #438: Automated QA gates
docs/matrix-crawl-runbook.md new
+227
@@ -0,0 +1,227 @@
1 +# Matrix Crawl Operator Runbook
2 +
3 +**Status:** Operational reference
4 +**Audience:** Pipeline operators, on-call engineers
5 +**Related issues:** #435, #436, #437, #438, #439
6 +
7 +---
8 +
9 +## Overview
10 +
11 +This runbook covers how to trigger, monitor, and troubleshoot the SquadScope matrix crawl and map/reduce dry-run pipeline. The default crawl topology is monolithic (non-matrix). Matrix fan-out is opt-in and experimental.
12 +
13 +---
14 +
15 +## 1. Triggering a Crawl Run
16 +
17 +### Scheduled (default)
18 +
19 +The `crawl-and-publish.yml` workflow runs automatically every Monday at 06:53 UTC via cron schedule.
20 +
21 +### Manual dispatch
22 +
23 +Use the GitHub Actions UI or CLI:
24 +
25 +```bash
26 +# Normal monolithic crawl
27 +gh workflow run crawl-and-publish.yml
28 +
29 +# Dry-run mode (no publishing, safe for experiments)
30 +gh workflow run crawl-and-publish.yml \
31 + -f run_mode=dry-run
32 +
33 +# Map/reduce dry-run analysis path
34 +gh workflow run crawl-and-publish.yml \
35 + -f run_mode=dry-run \
36 + -f analysis_path=map-reduce-dry-run
37 +
38 +# Force refresh all sources (ignore same-day cache)
39 +gh workflow run crawl-and-publish.yml \
40 + -f source_refresh_policy=force-refresh
41 +
42 +# Restore a specific past week
43 +gh workflow run crawl-and-publish.yml \
44 + -f run_mode=restore \
45 + -f rebuild_week=2026-W23
46 +```
47 +
48 +### Run modes
49 +
50 +| Mode | Publishes? | Creates release? | Use case |
51 +|------|:----------:|:----------------:|----------|
52 +| `normal` | Yes | If gated | Weekly production run |
53 +| `dry-run` | No | No | Testing changes, experiments |
54 +| `candidate-only` | No | No | Map/reduce validation |
55 +| `restore` | Yes | Optional | Rebuilding a past week |
56 +| `force-replace` | Yes | Yes | Explicit re-publication |
57 +
58 +### Analysis paths
59 +
60 +| Path | Description | Publishes? |
61 +|------|-------------|:----------:|
62 +| `single-pass` | Current monolithic AI analysis (default) | Yes |
63 +| `map-reduce-dry-run` | Experimental map/reduce scaffolding | Never |
64 +
65 +---
66 +
67 +## 2. Monitoring a Run
68 +
69 +### Key artifacts to check
70 +
71 +| Artifact | Location | Purpose |
72 +|----------|----------|---------|
73 +| `crawl-cache` | Actions artifact | GitHub API response cache |
74 +| `raw-data` | Actions artifact | Canonical crawl payloads |
75 +| `crawl-snapshots` | Actions artifact | Star/trending snapshots |
76 +| Rerun mode summary | `data/diagnostics/rerun-mode.json` | Mode validation output |
77 +| External news | `data/raw/{week}-external-news.json` | RSS crawl result |
78 +| Map/reduce candidates | `data/candidates/map-reduce/` | Dry-run output (if enabled) |
79 +
80 +### Metrics to watch (per #437)
81 +
82 +- **Crawl duration:** GitHub crawl p95 target < 6 minutes
83 +- **API calls:** Typical range 440–460 per run
84 +- **Search API remaining:** Should stay ≥ 24/30 after crawl
85 +- **Secondary rate limits:** Target: 0 events per run
86 +- **RSS fetch time:** Target < 5s total for all sources
87 +- **Source success rate:** Target ≥ 4/5 sources
88 +
89 +### Checking run status
90 +
91 +```bash
92 +# List recent workflow runs
93 +gh run list --workflow=crawl-and-publish.yml --limit 5
94 +
95 +# View a specific run
96 +gh run view <run-id>
97 +
98 +# Download artifacts for inspection
99 +gh run download <run-id> -n raw-data
100 +```
101 +
102 +---
103 +
104 +## 3. Troubleshooting
105 +
106 +### GitHub crawl failures
107 +
108 +| Symptom | Likely cause | Action |
109 +|---------|-------------|--------|
110 +| Secondary rate limit (HTTP 403 with `retry-after`) | Too many concurrent API calls | Check if shard experiment is running; reduce parallelism |
111 +| Search API quota exhausted | Excessive search queries | Wait for quota reset (resets per-minute); check for duplicate queries |
112 +| Crawl timeout (>10 min) | Network issues or API degradation | Retry; check [githubstatus.com](https://githubstatus.com) |
113 +| Cache miss storm | Config/code change invalidated cache | Expected on first run after changes; subsequent runs will rebuild cache |
114 +| Empty results | Token permission issue | Verify `GITHUB_TOKEN` has `contents: read` and is not expired |
115 +
116 +### RSS/News crawl failures
117 +
118 +| Symptom | Likely cause | Action |
119 +|---------|-------------|--------|
120 +| Single source timeout | Upstream feed slow/down | Check source URL manually; feed will retry once by default |
121 +| All sources failed | Network/DNS issue on runner | Check runner connectivity; retry the run |
122 +| Schema validation failure | Feed format changed | Check `scripts/techcrunch_crawler.py` parsing logic against current feed |
123 +| Deduplication anomaly | URL normalization issue | Check `GITHUB_URL_RE` and URL stripping logic |
124 +
125 +### Fan-in validation failures
126 +
127 +| Symptom | Likely cause | Action |
128 +|---------|-------------|--------|
129 +| Schema version mismatch | Mixed artifact versions | Ensure all legs use same code SHA (check `code_sha` in run context) |
130 +| Checksum mismatch | Non-deterministic serialization | Check for floating-point ordering or timestamp injection in articles |
131 +| Window mismatch | Leg computed own time | Verify all legs receive shared run context, not local `date` calls |
132 +| Missing required artifact | GitHub shard crashed | Check individual shard job logs; fix and re-run |
133 +| Minimum-source policy failed | ≥3 RSS sources down | Verify upstream feeds; consider temporary source list override |
134 +
135 +### Map/reduce dry-run failures
136 +
137 +| Symptom | Likely cause | Action |
138 +|---------|-------------|--------|
139 +| Missing raw JSON input | Crawl step didn't complete | Ensure crawl job succeeded before analysis |
140 +| Mapper schema violation | Contract change | Compare mapper output against `MAP_SCHEMA` in `scripts/map_reduce_dry_run.py` |
141 +| QA gate failure | Quality regression | Check `data/candidates/map-reduce/qa-report.json` for specific failures |
142 +| Token budget exceeded | Input growth | Check preflight vs actual token counts; may need input slicing |
143 +
144 +---
145 +
146 +## 4. Operational Procedures
147 +
148 +### Enabling RSS matrix mode (when triggers fire)
149 +
150 +Prerequisites (per PRD triggers):
151 +- RSS p95 > 60 seconds, OR
152 +- Source count > 10, OR
153 +- Source needs independent credentials/isolation
154 +
155 +Steps:
156 +1. Create a feature branch
157 +2. Modify workflow to add matrix strategy with `fail-fast: false`
158 +3. Each leg runs one source, uploads per-source artifact
159 +4. Add fan-in job that downloads all, validates, and merges
160 +5. Test with `dry-run` mode first
161 +6. Monitor metrics for 3+ runs before enabling for production
162 +
163 +### Running a GitHub shard experiment (#435)
164 +
165 +1. Trigger with `run_mode=dry-run` and shard configuration
166 +2. Run baseline (monolithic) and shard variant on same week window
167 +3. Compare: wall-clock time, API calls, rate-limit events, output stability
168 +4. Document results in experiment report
169 +5. Acceptance criteria: ≥25% speedup, ≤10% API growth, 0 rate-limit regression
170 +
171 +### Recovering from a failed run
172 +
173 +```bash
174 +# 1. Check what failed
175 +gh run view <run-id> --log-failed
176 +
177 +# 2. If crawl cache is stale, force refresh
178 +gh workflow run crawl-and-publish.yml \
179 + -f source_refresh_policy=force-refresh
180 +
181 +# 3. If a specific week needs rebuilding
182 +gh workflow run crawl-and-publish.yml \
183 + -f run_mode=restore \
184 + -f rebuild_week=2026-W23
185 +
186 +# 4. If analysis failed but crawl succeeded, re-run from artifacts
187 +gh run rerun <run-id> --failed
188 +```
189 +
190 +### Validating fan-in locally
191 +
192 +```bash
193 +# Run the deterministic map/reduce dry-run locally
194 +python scripts/map_reduce_dry_run.py \
195 + --raw-json data/raw/2026-W23.json \
196 + --output-dir data/candidates/map-reduce/ \
197 + --current-datetime "2026-06-05T17:42:56Z" \
198 + --run-id "local-test"
199 +
200 +# Verify canonical output is byte-stable
201 +python scripts/map_reduce_dry_run.py \
202 + --raw-json data/raw/2026-W23.json \
203 + --output-dir /tmp/mr-verify/ \
204 + --current-datetime "2026-06-05T17:42:56Z" \
205 + --run-id "local-test"
206 +
207 +diff data/candidates/map-reduce/ /tmp/mr-verify/
208 +```
209 +
210 +---
211 +
212 +## 5. Escalation Path
213 +
214 +1. **Self-serve:** Check this runbook and workflow logs
215 +2. **Team:** Tag the issue with `squad:bender` (crawler) or `squad:fry` (tests/QA)
216 +3. **Architecture:** Tag `squad:leela` for design decisions or contract changes
217 +4. **External:** GitHub API issues → check [githubstatus.com](https://githubstatus.com); RSS feed issues → check upstream provider status
218 +
219 +---
220 +
221 +## References
222 +
223 +- Workflow: [`.github/workflows/crawl-and-publish.yml`](../.github/workflows/crawl-and-publish.yml)
224 +- Fan-in contracts: [`docs/matrix-crawl-fan-in-contracts.md`](matrix-crawl-fan-in-contracts.md)
225 +- PRD: [`docs/processed/PRD-matrix-crawl-map-reduce-analysis.md`](processed/PRD-matrix-crawl-map-reduce-analysis.md)
226 +- Operator guide: [`docs/operator-guide.md`](operator-guide.md)
227 +- Pipeline validation: [`docs/pipeline-validation.md`](pipeline-validation.md)
scripts/archive/calibrate_hype_risk.py new
+273
@@ -0,0 +1,273 @@
1 +#!/usr/bin/env python3
2 +"""Calibrate hype risk scoring using accumulated momentum data.
3 +
4 +Reads momentum data from data/metrics/{topic}/momentum-*.json, compares
5 +hype risk predictions vs actual outcomes, and outputs a calibration report
6 +with recommended threshold adjustments.
7 +
8 +Usage:
9 + python scripts/calibrate_hype_risk.py [--topic ai-ml] [--output calibration.json]
10 +"""
11 +
12 +from __future__ import annotations
13 +
14 +import argparse
15 +import json
16 +import sys
17 +from datetime import datetime
18 +from pathlib import Path
19 +from typing import Any
20 +
21 +from scripts import topic_paths
22 +
23 +
24 +def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
25 + parser = argparse.ArgumentParser(
26 + description="Calibrate hype risk scoring model"
27 + )
28 + parser.add_argument(
29 + "--topic",
30 + default=None,
31 + help="Topic ID for path resolution.",
32 + )
33 + parser.add_argument(
34 + "--output",
35 + default=None,
36 + help="Output path for calibration report JSON.",
37 + )
38 + return parser.parse_args(argv)
39 +
40 +
41 +def load_momentum_files(metrics_directory: Path) -> list[dict[str, Any]]:
42 + """Load all momentum-*.json files from a metrics directory."""
43 + files = sorted(metrics_directory.glob("momentum-*.json"))
44 + results = []
45 + for f in files:
46 + try:
47 + with open(f, encoding="utf-8") as fh:
48 + data = json.load(fh)
49 + results.append(data)
50 + except (json.JSONDecodeError, OSError):
51 + continue
52 + return results
53 +
54 +
55 +def load_hype_risk_files(analyzed_directory: Path) -> list[dict[str, Any]]:
56 + """Load all *-hype-risk.json or hype risk assessment files."""
57 + files = sorted(analyzed_directory.glob("*hype*risk*.json"))
58 + results = []
59 + for f in files:
60 + try:
61 + with open(f, encoding="utf-8") as fh:
62 + data = json.load(fh)
63 + results.append(data)
64 + except (json.JSONDecodeError, OSError):
65 + continue
66 + return results
67 +
68 +
69 +def build_actual_outcomes(momentum_data: list[dict[str, Any]]) -> dict[str, str]:
70 + """Build repo -> actual outcome mapping from momentum data."""
71 + outcomes: dict[str, str] = {}
72 + for report in momentum_data:
73 + for repo in report.get("tracked_repos", []):
74 + repo_name = repo.get("repo", "")
75 + classification = repo.get("classification", "")
76 + if repo_name and classification:
77 + outcomes[repo_name] = classification
78 + return outcomes
79 +
80 +
81 +def build_predictions(
82 + hype_risk_data: list[dict[str, Any]],
83 + correlation_data: list[dict[str, Any]] | None = None,
84 +) -> dict[str, str]:
85 + """Build repo -> predicted risk level mapping from hype risk assessments."""
86 + predictions: dict[str, str] = {}
87 + for data in hype_risk_data:
88 + assessments = data.get("assessments", [])
89 + for assessment in assessments:
90 + repo = assessment.get("repo", "")
91 + risk = assessment.get("hype_risk", "")
92 + if repo and risk:
93 + predictions[repo] = risk
94 + return predictions
95 +
96 +
97 +def risk_to_expected_outcome(risk: str) -> str | None:
98 + """Map risk level to expected momentum outcome."""
99 + if risk in ("high",):
100 + return "faded"
101 + if risk in ("low", "very_low"):
102 + return "sustained"
103 + return None
104 +
105 +
106 +def compute_calibration(
107 + predictions: dict[str, str],
108 + actuals: dict[str, str],
109 +) -> dict[str, Any]:
110 + """Compare predictions vs actuals and compute accuracy by category."""
111 + accuracy_by_category: dict[str, dict[str, int]] = {}
112 + total_samples = 0
113 +
114 + for repo, risk_level in predictions.items():
115 + if repo not in actuals:
116 + continue
117 + expected = risk_to_expected_outcome(risk_level)
118 + if expected is None:
119 + continue
120 +
121 + actual = actuals[repo]
122 + total_samples += 1
123 +
124 + if risk_level not in accuracy_by_category:
125 + accuracy_by_category[risk_level] = {"predicted": 0, "correct": 0}
126 +
127 + accuracy_by_category[risk_level]["predicted"] += 1
128 + if expected == actual:
129 + accuracy_by_category[risk_level]["correct"] += 1
130 +
131 + for cat in accuracy_by_category.values():
132 + predicted = cat["predicted"]
133 + cat["accuracy"] = round(cat["correct"] / predicted, 4) if predicted > 0 else 0.0
134 +
135 + return {
136 + "samples": total_samples,
137 + "accuracy_by_category": accuracy_by_category,
138 + }
139 +
140 +
141 +def generate_recommendations(
142 + calibration: dict[str, Any],
143 + actuals: dict[str, str],
144 +) -> list[dict[str, Any]]:
145 + """Generate threshold adjustment recommendations based on calibration."""
146 + recommendations = []
147 + accuracy_by_cat = calibration.get("accuracy_by_category", {})
148 +
149 + high_stats = accuracy_by_cat.get("high", {})
150 + if high_stats.get("predicted", 0) > 0:
151 + high_acc = high_stats.get("accuracy", 0)
152 + if high_acc < 0.7:
153 + recommendations.append({
154 + "parameter": "high_risk_decay_threshold",
155 + "current": 0.5,
156 + "recommended": 0.6,
157 + "reason": (
158 + f"High-risk accuracy is {high_acc:.0%}, below 70% target. "
159 + "Raise decay threshold to reduce false positives."
160 + ),
161 + })
162 +
163 + low_stats = accuracy_by_cat.get("low", {})
164 + if low_stats.get("predicted", 0) > 0:
165 + low_acc = low_stats.get("accuracy", 0)
166 + if low_acc < 0.7:
167 + recommendations.append({
168 + "parameter": "sustained_threshold_weeks",
169 + "current": 2,
170 + "recommended": 3,
171 + "reason": (
172 + f"Low-risk (sustained) accuracy is {low_acc:.0%}. "
173 + "Extend observation window to improve confidence."
174 + ),
175 + })
176 +
177 + total_sustained = sum(1 for v in actuals.values() if v == "sustained")
178 + total_faded = sum(1 for v in actuals.values() if v == "faded")
179 + if total_sustained + total_faded > 0:
180 + sustained_ratio = total_sustained / (total_sustained + total_faded)
181 + if sustained_ratio > 0.7:
182 + recommendations.append({
183 + "parameter": "press_correlation_confidence_floor",
184 + "current": 0.4,
185 + "recommended": 0.5,
186 + "reason": (
187 + f"Sustained ratio is {sustained_ratio:.0%}, suggesting most "
188 + "press-correlated repos maintain growth. Raise confidence "
189 + "floor to only flag truly risky repos."
190 + ),
191 + })
192 +
193 + if not recommendations:
194 + recommendations.append({
195 + "parameter": "no_changes",
196 + "current": None,
197 + "recommended": None,
198 + "reason": "Calibration shows acceptable accuracy. No adjustments needed.",
199 + })
200 +
201 + return recommendations
202 +
203 +
204 +def run_calibration(
205 + topic_id: str | None = None,
206 + output_path: str | None = None,
207 +) -> dict[str, Any]:
208 + """Main calibration logic."""
209 + metrics_directory = topic_paths.metrics_dir(topic_id)
210 + analyzed_directory = topic_paths.analyzed_dir(topic_id)
211 +
212 + momentum_data = load_momentum_files(metrics_directory)
213 + if not momentum_data:
214 + print(f"No momentum data found in {metrics_directory}", file=sys.stderr)
215 + report = {
216 + "calibration_date": datetime.now().strftime("%Y-%m-%d"),
217 + "samples": 0,
218 + "accuracy_by_category": {},
219 + "recommended_adjustments": [],
220 + }
221 + if output_path:
222 + out = Path(output_path)
223 + out.parent.mkdir(parents=True, exist_ok=True)
224 + with open(out, "w", encoding="utf-8") as f:
225 + json.dump(report, f, indent=2, ensure_ascii=False)
226 + f.write("\n")
227 + return report
228 +
229 + hype_risk_data = load_hype_risk_files(analyzed_directory)
230 +
231 + actuals = build_actual_outcomes(momentum_data)
232 + predictions = build_predictions(hype_risk_data)
233 +
234 + calibration = compute_calibration(predictions, actuals)
235 + recommendations = generate_recommendations(calibration, actuals)
236 +
237 + report = {
238 + "calibration_date": datetime.now().strftime("%Y-%m-%d"),
239 + "samples": calibration["samples"],
240 + "accuracy_by_category": calibration["accuracy_by_category"],
241 + "recommended_adjustments": recommendations,
242 + }
243 +
244 + if output_path:
245 + out = Path(output_path)
246 + else:
247 + metrics_directory.mkdir(parents=True, exist_ok=True)
248 + out = metrics_directory / "calibration-report.json"
249 +
250 + out.parent.mkdir(parents=True, exist_ok=True)
251 + with open(out, "w", encoding="utf-8") as f:
252 + json.dump(report, f, indent=2, ensure_ascii=False)
253 + f.write("\n")
254 +
255 + print(f"Calibration report: {calibration['samples']} samples")
256 + for cat, stats in calibration["accuracy_by_category"].items():
257 + print(f" [{cat}] {stats['correct']}/{stats['predicted']} ({stats['accuracy']:.0%})")
258 + print(f"Wrote report to {out}")
259 +
260 + return report
261 +
262 +
263 +def main(argv: list[str] | None = None) -> dict[str, Any]:
264 + """CLI entry point."""
265 + args = parse_args(argv)
266 + return run_calibration(
267 + topic_id=args.topic,
268 + output_path=args.output,
269 + )
270 +
271 +
272 +if __name__ == "__main__":
273 + main()
scripts/rss_fan_in.py new
+665
@@ -0,0 +1,665 @@
1 +#!/usr/bin/env python3
2 +"""RSS matrix fan-in: per-source artifact emission and deterministic merge.
3 +
4 +This module implements the fan-in mechanism for matrix-based RSS crawling.
5 +Each RSS source emits a per-source artifact with shared run context, and
6 +the merge step combines them deterministically into the canonical
7 +external-news artifact consumed by downstream analysis.
8 +
9 +The default in-process RSS crawl path remains unchanged. This fan-in path
10 +activates only when the matrix crawl mode is explicitly enabled.
11 +
12 +Usage (emit per-source artifact):
13 + python -m scripts.rss_fan_in emit \
14 + --source techcrunch \
15 + --articles articles.json \
16 + --output artifacts/techcrunch.json \
17 + --run-context run-context.json
18 +
19 +Usage (merge per-source artifacts):
20 + python -m scripts.rss_fan_in merge \
21 + --artifacts-dir artifacts/ \
22 + --output data/raw/general/2026-W24-external-news.json \
23 + --run-context run-context.json
24 +
25 +References:
26 + - Issue #436: Implement RSS matrix fan-in
27 + - Issue #356: Matrix Crawl & Map/Reduce PRD
28 + - Issue #333: Make canonical artifacts matrix-ready
29 +"""
30 +
31 +from __future__ import annotations
32 +
33 +import argparse
34 +import hashlib
35 +import json
36 +import sys
37 +from datetime import UTC, datetime
38 +from pathlib import Path
39 +from typing import Any
40 +
41 +from scripts.techcrunch_crawler import (
42 + CANONICAL_SCHEMA_VERSION,
43 + artifact_checksum,
44 + dedupe_articles,
45 + iso_timestamp,
46 + load_source_configs,
47 + schema_checksum,
48 + source_config_checksum,
49 + source_content_checksum,
50 + validate_canonical_output,
51 + week_slug,
52 +)
53 +
54 +# Per-source artifact schema version (tracks independently of canonical)
55 +SOURCE_ARTIFACT_SCHEMA_VERSION = 1
56 +
57 +
58 +class FanInValidationError(Exception):
59 + """Raised when fan-in validation detects an unrecoverable error."""
60 +
61 + pass
62 +
63 +
64 +class FanInWarning:
65 + """Represents a non-fatal fan-in issue that allows merge to proceed."""
66 +
67 + def __init__(self, source_id: str, category: str, message: str) -> None:
68 + self.source_id = source_id
69 + self.category = category
70 + self.message = message
71 +
72 + def to_dict(self) -> dict[str, str]:
73 + return {
74 + "source_id": self.source_id,
75 + "category": self.category,
76 + "message": self.message,
77 + }
78 +
79 +
80 +# ---------------------------------------------------------------------------
81 +# Run Context
82 +# ---------------------------------------------------------------------------
83 +
84 +
85 +def build_run_context(
86 + *,
87 + run_id: str,
88 + week: str,
89 + crawl_window: dict[str, str],
90 + source_config_checksum_value: str,
91 + schema_checksum_value: str,
92 + sources_requested: list[str],
93 + required_sources: list[str] | None = None,
94 + optional_sources: list[str] | None = None,
95 + started_at: str | None = None,
96 + crawler_code_sha: str | None = None,
97 +) -> dict[str, Any]:
98 + """Build the shared run context distributed to all matrix jobs."""
99 + all_requested = sorted(sources_requested)
100 + required = sorted(required_sources or all_requested)
101 + optional = sorted(optional_sources or [])
102 + return {
103 + "schema_version": SOURCE_ARTIFACT_SCHEMA_VERSION,
104 + "run_id": run_id,
105 + "week": week,
106 + "crawl_window": crawl_window,
107 + "source_config_checksum": source_config_checksum_value,
108 + "schema_checksum": schema_checksum_value,
109 + "sources_requested": all_requested,
110 + "required_sources": required,
111 + "optional_sources": optional,
112 + "started_at": started_at or iso_timestamp(datetime.now(UTC)),
113 + "crawler_code_sha": crawler_code_sha or "",
114 + }
115 +
116 +
117 +def validate_run_context(ctx: dict[str, Any]) -> None:
118 + """Validate run context structure."""
119 + required_keys = {
120 + "schema_version",
121 + "run_id",
122 + "week",
123 + "crawl_window",
124 + "source_config_checksum",
125 + "schema_checksum",
126 + "sources_requested",
127 + "required_sources",
128 + "started_at",
129 + }
130 + missing = sorted(required_keys - set(ctx))
131 + if missing:
132 + raise FanInValidationError(f"Run context missing keys: {missing}")
133 + if ctx["schema_version"] != SOURCE_ARTIFACT_SCHEMA_VERSION:
134 + raise FanInValidationError(
135 + f"Run context schema_version mismatch: expected {SOURCE_ARTIFACT_SCHEMA_VERSION}, "
136 + f"got {ctx['schema_version']}"
137 + )
138 + window = ctx.get("crawl_window")
139 + if not isinstance(window, dict) or "since" not in window or "until" not in window:
140 + raise FanInValidationError("Run context crawl_window must have 'since' and 'until'")
141 +
142 +
143 +# ---------------------------------------------------------------------------
144 +# Per-Source Artifact
145 +# ---------------------------------------------------------------------------
146 +
147 +
148 +def build_source_artifact(
149 + *,
150 + source_id: str,
151 + articles: list[dict[str, Any]],
152 + status: dict[str, Any],
153 + run_context: dict[str, Any],
154 + crawled_at: datetime | None = None,
155 +) -> dict[str, Any]:
156 + """Build a per-source artifact for one RSS source's crawl results.
157 +
158 + Each per-source artifact contains enough context for the fan-in merge
159 + to validate provenance, detect staleness, and produce deterministic output.
160 + """
161 + validate_run_context(run_context)
162 + now = crawled_at or datetime.now(UTC)
163 +
164 + # Sort articles deterministically
165 + sorted_articles = sorted(
166 + articles,
167 + key=lambda a: (
168 + a.get("published_at", ""),
169 + a.get("url", ""),
170 + a.get("title", ""),
171 + ),
172 + reverse=True,
173 + )
174 +
175 + content_checksum = source_content_checksum(source_id, sorted_articles)
176 + relevant_count = sum(1 for a in sorted_articles if a.get("relevance_score", 0) >= 0.4)
177 +
178 + artifact: dict[str, Any] = {
179 + "source_artifact_schema_version": SOURCE_ARTIFACT_SCHEMA_VERSION,
180 + "source_id": source_id,
181 + "crawled_at": iso_timestamp(now),
182 + "run_context": {
183 + "run_id": run_context["run_id"],
184 + "week": run_context["week"],
185 + "crawl_window": run_context["crawl_window"],
186 + "source_config_checksum": run_context["source_config_checksum"],
187 + "schema_checksum": run_context["schema_checksum"],
188 + "started_at": run_context["started_at"],
189 + "crawler_code_sha": run_context.get("crawler_code_sha", ""),
190 + },
191 + "status": status,
192 + "metrics": {
193 + "total_articles": len(sorted_articles),
194 + "relevant_articles": relevant_count,
195 + "content_checksum": content_checksum,
196 + },
197 + "articles": sorted_articles,
198 + }
199 +
200 + # Compute artifact-level checksum over deterministic content
201 + artifact["artifact_checksum"] = _source_artifact_checksum(artifact)
202 + return artifact
203 +
204 +
205 +def _source_artifact_checksum(artifact: dict[str, Any]) -> str:
206 + """Compute a checksum for the per-source artifact content."""
207 + payload = {
208 + "source_id": artifact["source_id"],
209 + "run_context": artifact["run_context"],
210 + "articles": artifact["articles"],
211 + }
212 + serialized = json.dumps(payload, sort_keys=True, separators=(",", ":"), ensure_ascii=False)
213 + return hashlib.sha256(serialized.encode("utf-8")).hexdigest()
214 +
215 +
216 +def validate_source_artifact(artifact: dict[str, Any]) -> None:
217 + """Validate a per-source artifact structure."""
218 + if not isinstance(artifact, dict):
219 + raise FanInValidationError("Source artifact must be a JSON object")
220 + if artifact.get("source_artifact_schema_version") != SOURCE_ARTIFACT_SCHEMA_VERSION:
221 + raise FanInValidationError(
222 + f"Source artifact schema version mismatch: expected {SOURCE_ARTIFACT_SCHEMA_VERSION}, "
223 + f"got {artifact.get('source_artifact_schema_version')}"
224 + )
225 + required = {"source_id", "crawled_at", "run_context", "status", "metrics", "articles", "artifact_checksum"}
226 + missing = sorted(required - set(artifact))
227 + if missing:
228 + raise FanInValidationError(f"Source artifact missing keys: {missing}")
229 + # Verify checksum integrity
230 + expected = _source_artifact_checksum(artifact)
231 + if artifact["artifact_checksum"] != expected:
232 + raise FanInValidationError(
233 + f"Source artifact checksum mismatch for {artifact.get('source_id', '?')}: "
234 + f"expected {expected}, got {artifact['artifact_checksum']}"
235 + )
236 +
237 +
238 +# ---------------------------------------------------------------------------
239 +# Fan-In Merge
240 +# ---------------------------------------------------------------------------
241 +
242 +
243 +def validate_fan_in_compatibility(
244 + artifacts: list[dict[str, Any]],
245 + run_context: dict[str, Any],
246 +) -> list[FanInWarning]:
247 + """Validate that all per-source artifacts are compatible for merge.
248 +
249 + Returns warnings for non-fatal issues. Raises FanInValidationError for
250 + unrecoverable problems (schema mismatch, window mismatch, etc.).
251 + """
252 + warnings: list[FanInWarning] = []
253 +
254 + if not artifacts:
255 + raise FanInValidationError("No source artifacts provided for fan-in merge")
256 +
257 + validate_run_context(run_context)
258 +
259 + for artifact in artifacts:
260 + validate_source_artifact(artifact)
261 + ctx = artifact["run_context"]
262 + source_id = artifact["source_id"]
263 +
264 + # Schema version must match
265 + if ctx.get("schema_checksum") != run_context["schema_checksum"]:
266 + raise FanInValidationError(
267 + f"Schema checksum mismatch for source '{source_id}': "
268 + f"artifact has {ctx.get('schema_checksum')}, "
269 + f"run context has {run_context['schema_checksum']}"
270 + )
271 +
272 + # Crawl window must match
273 + if ctx.get("crawl_window") != run_context["crawl_window"]:
274 + raise FanInValidationError(
275 + f"Crawl window mismatch for source '{source_id}': "
276 + f"artifact has {ctx.get('crawl_window')}, "
277 + f"run context has {run_context['crawl_window']}"
278 + )
279 +
280 + # Source config checksum must match
281 + if ctx.get("source_config_checksum") != run_context["source_config_checksum"]:
282 + raise FanInValidationError(
283 + f"Source config checksum mismatch for source '{source_id}': "
284 + f"artifact has {ctx.get('source_config_checksum')}, "
285 + f"run context has {run_context['source_config_checksum']}"
286 + )
287 +
288 + # Run ID must match
289 + if ctx.get("run_id") != run_context["run_id"]:
290 + raise FanInValidationError(
291 + f"Run ID mismatch for source '{source_id}': "
292 + f"artifact has {ctx.get('run_id')}, "
293 + f"run context has {run_context['run_id']}"
294 + )
295 +
296 + # Check for required sources
297 + provided_sources = {a["source_id"] for a in artifacts}
298 + required_sources = set(run_context.get("required_sources", []))
299 + missing_required = sorted(required_sources - provided_sources)
300 + if missing_required:
301 + raise FanInValidationError(
302 + f"Missing required source artifacts: {missing_required}"
303 + )
304 +
305 + # Check for optional missing sources (warning, not error)
306 + optional_sources = set(run_context.get("optional_sources", []))
307 + missing_optional = sorted(optional_sources - provided_sources)
308 + for source_id in missing_optional:
309 + warnings.append(FanInWarning(
310 + source_id=source_id,
311 + category="missing_optional_source",
312 + message=f"Optional source '{source_id}' artifact not found",
313 + ))
314 +
315 + # Duplicate source check
316 + source_ids = [a["source_id"] for a in artifacts]
317 + seen: set[str] = set()
318 + for sid in source_ids:
319 + if sid in seen:
320 + raise FanInValidationError(f"Duplicate source artifact for '{sid}'")
321 + seen.add(sid)
322 +
323 + return warnings
324 +
325 +
326 +def merge_source_artifacts(
327 + artifacts: list[dict[str, Any]],
328 + run_context: dict[str, Any],
329 + *,
330 + merged_at: datetime | None = None,
331 +) -> tuple[dict[str, Any], list[FanInWarning]]:
332 + """Deterministically merge per-source artifacts into the canonical output.
333 +
334 + The merge is deterministic: given the same set of per-source artifacts
335 + and run context, it always produces the same canonical output (minus
336 + the crawled_at timestamp which is excluded from the checksum).
337 +
338 + Returns (canonical_output, warnings).
339 + """
340 + warnings = validate_fan_in_compatibility(artifacts, run_context)
341 + now = merged_at or datetime.now(UTC)
342 +
343 + # Sort artifacts by source_id for deterministic processing
344 + sorted_artifacts = sorted(artifacts, key=lambda a: a["source_id"])
345 +
346 + # Collect all articles from all sources
347 + all_articles: list[dict[str, Any]] = []
348 + source_statuses: list[dict[str, Any]] = []
349 + source_provenance: list[dict[str, Any]] = []
350 + errors: list[dict[str, str]] = []
351 +
352 + for artifact in sorted_artifacts:
353 + source_id = artifact["source_id"]
354 + status = artifact.get("status", {})
355 +
356 + all_articles.extend(artifact.get("articles", []))
357 + source_statuses.append(status)
358 +
359 + if not status.get("success", False):
360 + warnings.append(FanInWarning(
361 + source_id=source_id,
362 + category="source_failure",
363 + message=f"Source '{source_id}' reported failure: {status.get('error_message', 'unknown')}",
364 + ))
365 + if status.get("error_class") or status.get("error_message"):
366 + errors.append({
367 + "source": source_id,
368 + "error_class": status.get("error_class", "Unknown"),
369 + "error": status.get("error_message", "unknown error"),
370 + })
371 +
372 + provenance_entry = {
373 + "source_id": source_id,
374 + "action": "matrix_fan_in",
375 + "artifact_checksum": artifact["artifact_checksum"],
376 + "content_checksum": artifact["metrics"]["content_checksum"],
377 + "original_run_id": run_context["run_id"],
378 + "original_crawled_at": artifact["crawled_at"],
379 + "evaluated_at": iso_timestamp(now),
380 + "date": now.astimezone(UTC).date().isoformat(),
381 + "week": run_context["week"],
382 + "crawl_window": run_context["crawl_window"],
383 + "source_config_checksum": run_context["source_config_checksum"],
384 + "schema_checksum": run_context["schema_checksum"],
385 + "reasons": [],
386 + }
387 + source_provenance.append(provenance_entry)
388 +
389 + # Build reuse summary entries for fan-in sources
390 + reuse_summary = [
391 + {
392 + "source": artifact["source_id"],
393 + "action": "matrix_fan_in",
394 + "reused": False,
395 + "refreshed": True,
396 + "reasons": [],
397 + }
398 + for artifact in sorted_artifacts
399 + ]
400 +
401 + # Build canonical output using the shared build_output function
402 + requested_sources = sorted(run_context.get("sources_requested", []))
403 + succeeded = sorted(a["source_id"] for a in sorted_artifacts if a["status"].get("success"))
404 + failed = sorted(a["source_id"] for a in sorted_artifacts if not a["status"].get("success"))
405 +
406 + output = build_canonical_merged_output(
407 + articles=all_articles,
408 + crawled_at=now,
409 + run_context=run_context,
410 + source_statuses=source_statuses,
411 + source_provenance=source_provenance,
412 + reuse_summary=reuse_summary,
413 + errors=errors,
414 + requested_sources=requested_sources,
415 + succeeded_sources=succeeded,
416 + failed_sources=failed,
417 + )
418 +
419 + return output, warnings
420 +
421 +
422 +def build_canonical_merged_output(
423 + *,
424 + articles: list[dict[str, Any]],
425 + crawled_at: datetime,
426 + run_context: dict[str, Any],
427 + source_statuses: list[dict[str, Any]],
428 + source_provenance: list[dict[str, Any]],
429 + reuse_summary: list[dict[str, Any]],
430 + errors: list[dict[str, str]],
431 + requested_sources: list[str],
432 + succeeded_sources: list[str],
433 + failed_sources: list[str],
434 +) -> dict[str, Any]:
435 + """Build the canonical merged external-news artifact from fan-in results.
436 +
437 + This uses the same schema as the non-matrix path to ensure downstream
438 + compatibility.
439 + """
440 + # Deduplicate articles (same logic as non-matrix path)
441 + deduped_articles, dedupe_count = dedupe_articles(articles)
442 + relevant = [a for a in deduped_articles if a.get("relevance_score", 0) >= 0.4]
443 +
444 + all_github_links: set[str] = set()
445 + for a in deduped_articles:
446 + all_github_links.update(a.get("github_links", []))
447 +
448 + output: dict[str, Any] = {
449 + "schema_version": CANONICAL_SCHEMA_VERSION,
450 + "week": run_context["week"],
451 + "source": "external_news",
452 + "crawled_at": iso_timestamp(crawled_at),
453 + "crawl_window": run_context["crawl_window"],
454 + "articles": deduped_articles,
455 + "metadata": {
456 + "run_id": run_context["run_id"],
457 + "source_count": len(requested_sources),
458 + "source_config_checksum": run_context["source_config_checksum"],
459 + "schema_checksum": run_context["schema_checksum"],
460 + "sources_requested": sorted(requested_sources),
461 + "sources_succeeded": sorted(succeeded_sources),
462 + "sources_failed": sorted(failed_sources),
463 + "source_status": sorted(source_statuses, key=lambda s: s.get("source", "")),
464 + "source_reuse_summary": sorted(reuse_summary, key=lambda s: s.get("source", "")),
465 + "source_artifact_provenance": sorted(source_provenance, key=lambda s: s.get("source_id", "")),
466 + "sources_with_articles": dict(sorted(
467 + {str(a.get("source", "unknown")): 0 for a in deduped_articles}.items()
468 + )),
469 + "total_articles": len(deduped_articles),
470 + "relevant_articles": len(relevant),
471 + "github_links_found": len(all_github_links),
472 + "dedupe_count": dedupe_count,
473 + "errors": sorted(errors, key=lambda e: e.get("source", "")),
474 + "fan_in_mode": "matrix",
475 + "crawler_code_sha": run_context.get("crawler_code_sha", ""),
476 + },
477 + }
478 +
479 + # Compute per-source article counts
480 + from collections import Counter
481 +
482 + by_source = Counter(str(a.get("source", "unknown")) for a in deduped_articles)
483 + output["metadata"]["sources_with_articles"] = dict(sorted(by_source.items()))
484 +
485 + # Compute and set artifact checksum
486 + output["metadata"]["artifact_checksum"] = artifact_checksum(output)
487 +
488 + # Fill in provenance checksums that reference the merged artifact
489 + for entry in output["metadata"]["source_artifact_provenance"]:
490 + if not entry.get("artifact_checksum"):
491 + entry["artifact_checksum"] = output["metadata"]["artifact_checksum"]
492 +
493 + validate_canonical_output(output)
494 + return output
495 +
496 +
497 +# ---------------------------------------------------------------------------
498 +# CLI
499 +# ---------------------------------------------------------------------------
500 +
501 +
502 +def cmd_emit(args: argparse.Namespace) -> int:
503 + """Emit a per-source artifact from crawl results."""
504 + run_context = json.loads(Path(args.run_context).read_text(encoding="utf-8"))
505 + validate_run_context(run_context)
506 +
507 + articles = json.loads(Path(args.articles).read_text(encoding="utf-8"))
508 + if not isinstance(articles, list):
509 + print("ERROR: articles file must contain a JSON array", file=sys.stderr)
510 + return 1
511 +
512 + status = json.loads(Path(args.status).read_text(encoding="utf-8")) if args.status else {
513 + "source": args.source,
514 + "success": True,
515 + }
516 +
517 + artifact = build_source_artifact(
518 + source_id=args.source,
519 + articles=articles,
520 + status=status,
521 + run_context=run_context,
522 + )
523 +
524 + out_path = Path(args.output)
525 + out_path.parent.mkdir(parents=True, exist_ok=True)
526 + with open(out_path, "w", encoding="utf-8") as f:
527 + json.dump(artifact, f, indent=2, ensure_ascii=False)
528 +
529 + print(f"Emitted per-source artifact for '{args.source}' → {out_path}", file=sys.stderr)
530 + return 0
531 +
532 +
533 +def cmd_merge(args: argparse.Namespace) -> int:
534 + """Merge per-source artifacts into canonical output."""
535 + run_context = json.loads(Path(args.run_context).read_text(encoding="utf-8"))
536 + validate_run_context(run_context)
537 +
538 + artifacts_dir = Path(args.artifacts_dir)
539 + if not artifacts_dir.is_dir():
540 + print(f"ERROR: artifacts directory not found: {artifacts_dir}", file=sys.stderr)
541 + return 1
542 +
543 + # Load all per-source artifact files
544 + artifact_files = sorted(artifacts_dir.glob("*.json"))
545 + if not artifact_files:
546 + print(f"ERROR: no .json artifacts found in {artifacts_dir}", file=sys.stderr)
547 + return 1
548 +
549 + artifacts: list[dict[str, Any]] = []
550 + for path in artifact_files:
551 + try:
552 + data = json.loads(path.read_text(encoding="utf-8"))
553 + except (json.JSONDecodeError, OSError) as exc:
554 + print(f"ERROR: failed to read artifact {path}: {exc}", file=sys.stderr)
555 + return 1
556 + # Skip non-source-artifact files (e.g., run-context.json in same dir)
557 + if not isinstance(data, dict) or "source_artifact_schema_version" not in data:
558 + continue
559 + artifacts.append(data)
560 +
561 + if not artifacts:
562 + print(f"ERROR: no valid source artifacts found in {artifacts_dir}", file=sys.stderr)
563 + return 1
564 +
565 + try:
566 + output, warnings = merge_source_artifacts(artifacts, run_context)
567 + except FanInValidationError as exc:
568 + print(f"ERROR: fan-in validation failed: {exc}", file=sys.stderr)
569 + return 1
570 +
571 + for w in warnings:
572 + print(f"WARNING [{w.category}] {w.source_id}: {w.message}", file=sys.stderr)
573 +
574 + out_path = Path(args.output)
575 + out_path.parent.mkdir(parents=True, exist_ok=True)
576 + with open(out_path, "w", encoding="utf-8") as f:
577 + json.dump(output, f, indent=2, ensure_ascii=False)
578 +
579 + total = output["metadata"]["total_articles"]
580 + relevant = output["metadata"]["relevant_articles"]
581 + dedupe = output["metadata"]["dedupe_count"]
582 + sources = len(artifacts)
583 + print(
584 + f"Merged {total} articles from {sources} sources "
585 + f"({relevant} relevant, {dedupe} deduped) → {out_path}",
586 + file=sys.stderr,
587 + )
588 + return 0
589 +
590 +
591 +def cmd_validate(args: argparse.Namespace) -> int:
592 + """Validate per-source artifacts against run context without merging."""
593 + run_context = json.loads(Path(args.run_context).read_text(encoding="utf-8"))
594 + validate_run_context(run_context)
595 +
596 + artifacts_dir = Path(args.artifacts_dir)
597 + artifact_files = sorted(artifacts_dir.glob("*.json"))
598 + artifacts: list[dict[str, Any]] = []
599 + for path in artifact_files:
600 + try:
601 + data = json.loads(path.read_text(encoding="utf-8"))
602 + except (json.JSONDecodeError, OSError) as exc:
603 + print(f"ERROR: failed to read {path}: {exc}", file=sys.stderr)
604 + return 1
605 + if isinstance(data, dict) and "source_artifact_schema_version" in data:
606 + artifacts.append(data)
607 +
608 + if not artifacts:
609 + print(f"ERROR: no valid source artifacts in {artifacts_dir}", file=sys.stderr)
610 + return 1
611 +
612 + try:
613 + warnings = validate_fan_in_compatibility(artifacts, run_context)
614 + except FanInValidationError as exc:
615 + print(f"FAIL: {exc}", file=sys.stderr)
616 + return 1
617 +
618 + for w in warnings:
619 + print(f"WARNING [{w.category}] {w.source_id}: {w.message}", file=sys.stderr)
620 +
621 + print(f"OK: {len(artifacts)} source artifacts validated successfully", file=sys.stderr)
622 + return 0
623 +
624 +
625 +def main(argv: list[str] | None = None) -> int:
626 + parser = argparse.ArgumentParser(
627 + description="RSS matrix fan-in: per-source artifact emission and deterministic merge"
628 + )
629 + subparsers = parser.add_subparsers(dest="command")
630 +
631 + # emit subcommand
632 + emit_parser = subparsers.add_parser("emit", help="Emit a per-source artifact")
633 + emit_parser.add_argument("--source", required=True, help="Source ID (e.g., techcrunch)")
634 + emit_parser.add_argument("--articles", required=True, help="Path to articles JSON array")
635 + emit_parser.add_argument("--status", default=None, help="Path to source status JSON (optional)")
636 + emit_parser.add_argument("--run-context", required=True, help="Path to shared run context JSON")
637 + emit_parser.add_argument("--output", required=True, help="Output path for per-source artifact")
638 +
639 + # merge subcommand
640 + merge_parser = subparsers.add_parser("merge", help="Merge per-source artifacts")
641 + merge_parser.add_argument("--artifacts-dir", required=True, help="Directory containing per-source artifacts")
642 + merge_parser.add_argument("--run-context", required=True, help="Path to shared run context JSON")
643 + merge_parser.add_argument("--output", required=True, help="Output path for merged canonical artifact")
644 +
645 + # validate subcommand
646 + validate_parser = subparsers.add_parser("validate", help="Validate artifacts without merging")
647 + validate_parser.add_argument("--artifacts-dir", required=True, help="Directory containing per-source artifacts")
648 + validate_parser.add_argument("--run-context", required=True, help="Path to shared run context JSON")
649 +
650 + args = parser.parse_args(argv)
651 + if not args.command:
652 + parser.print_help()
653 + return 1
654 +
655 + if args.command == "emit":
656 + return cmd_emit(args)
657 + elif args.command == "merge":
658 + return cmd_merge(args)
659 + elif args.command == "validate":
660 + return cmd_validate(args)
661 + return 1
662 +
663 +
664 +if __name__ == "__main__":
665 + sys.exit(main())
tests/test_rss_fan_in.py new
+591
@@ -0,0 +1,591 @@
1 +"""Tests for RSS matrix fan-in: per-source artifacts and deterministic merge.
2 +
3 +Covers:
4 +- Per-source artifact emission and validation
5 +- Run context building and validation
6 +- Deterministic merge producing canonical output
7 +- Fan-in validation: schema mismatch, window mismatch, missing sources, duplicates
8 +- Partial optional-source failures with warnings
9 +- Deduplication across sources
10 +- Fixture-based determinism proof
11 +"""
12 +
13 +from __future__ import annotations
14 +
15 +import json
16 +from datetime import UTC, datetime, timedelta
17 +from pathlib import Path
18 +from typing import Any
19 +
20 +import pytest
21 +
22 +from scripts.rss_fan_in import (
23 + SOURCE_ARTIFACT_SCHEMA_VERSION,
24 + FanInValidationError,
25 + FanInWarning,
26 + build_canonical_merged_output,
27 + build_run_context,
28 + build_source_artifact,
29 + merge_source_artifacts,
30 + validate_fan_in_compatibility,
31 + validate_run_context,
32 + validate_source_artifact,
33 +)
34 +from scripts.techcrunch_crawler import (
35 + CANONICAL_SCHEMA_VERSION,
36 + iso_timestamp,
37 + schema_checksum,
38 + source_config_checksum,
39 + week_slug,
40 +)
41 +
42 +
43 +# --- Fixtures ---
44 +
45 +NOW = datetime(2026, 6, 13, 12, 0, 0, tzinfo=UTC)
46 +SINCE = datetime(2026, 6, 6, 0, 0, 0, tzinfo=UTC)
47 +UNTIL = datetime(2026, 6, 13, 0, 0, 0, tzinfo=UTC)
48 +WEEK = "2026-W24"
49 +RUN_ID = "test-run-12345"
50 +CONFIG_CHECKSUM = "abc123def456"
51 +SCHEMA_CHECKSUM_VALUE = schema_checksum()
52 +
53 +
54 +def _make_run_context(**overrides: Any) -> dict[str, Any]:
55 + ctx = {
56 + "schema_version": SOURCE_ARTIFACT_SCHEMA_VERSION,
57 + "run_id": RUN_ID,
58 + "week": WEEK,
59 + "crawl_window": {"since": iso_timestamp(SINCE), "until": iso_timestamp(UNTIL)},
60 + "source_config_checksum": CONFIG_CHECKSUM,
61 + "schema_checksum": SCHEMA_CHECKSUM_VALUE,
62 + "sources_requested": ["techcrunch", "github_blog", "nvidia_blog"],
63 + "required_sources": ["techcrunch", "github_blog"],
64 + "optional_sources": ["nvidia_blog"],
65 + "started_at": iso_timestamp(NOW),
66 + "crawler_code_sha": "sha256-test",
67 + }
68 + ctx.update(overrides)
69 + return ctx
70 +
71 +
72 +def _make_article(source: str, title: str = "Test Article", url: str = "") -> dict[str, Any]:
73 + return {
74 + "source": source,
75 + "title": title,
76 + "url": url or f"https://example.com/{source}/{title.lower().replace(' ', '-')}",
77 + "published_at": iso_timestamp(NOW - timedelta(hours=2)),
78 + "categories": ["AI", "Open Source"],
79 + "summary": "A test article about AI and machine learning frameworks.",
80 + "github_links": ["https://github.com/org/repo"],
81 + "entities": ["TestCo"],
82 + "relevance_score": 0.8,
83 + }
84 +
85 +
86 +def _make_status(source: str, success: bool = True) -> dict[str, Any]:
87 + status: dict[str, Any] = {
88 + "source": source,
89 + "host": f"{source}.example.com",
90 + "started_at": iso_timestamp(NOW),
91 + "ended_at": iso_timestamp(NOW + timedelta(seconds=1)),
92 + "duration_seconds": 1.0,
93 + "timeout_seconds": 15,
94 + "attempts": 1,
95 + "total_articles": 3 if success else 0,
96 + "relevant_articles": 2 if success else 0,
97 + "github_links_found": 1 if success else 0,
98 + "success": success,
99 + "error_class": "" if success else "ConnectionError",
100 + "error_message": "" if success else "Connection refused",
101 + }
102 + return status
103 +
104 +
105 +def _make_source_artifact(
106 + source_id: str,
107 + run_context: dict[str, Any] | None = None,
108 + articles: list[dict[str, Any]] | None = None,
109 + success: bool = True,
110 +) -> dict[str, Any]:
111 + ctx = run_context or _make_run_context()
112 + arts = articles if articles is not None else [
113 + _make_article(source_id, f"Article {i}") for i in range(3)
114 + ]
115 + status = _make_status(source_id, success=success)
116 + return build_source_artifact(
117 + source_id=source_id,
118 + articles=arts,
119 + status=status,
120 + run_context=ctx,
121 + crawled_at=NOW,
122 + )
123 +
124 +
125 +# --- Run Context Tests ---
126 +
127 +
128 +class TestRunContext:
129 + def test_build_run_context(self) -> None:
130 + ctx = build_run_context(
131 + run_id=RUN_ID,
132 + week=WEEK,
133 + crawl_window={"since": iso_timestamp(SINCE), "until": iso_timestamp(UNTIL)},
134 + source_config_checksum_value=CONFIG_CHECKSUM,
135 + schema_checksum_value=SCHEMA_CHECKSUM_VALUE,
136 + sources_requested=["techcrunch", "github_blog"],
137 + required_sources=["techcrunch"],
138 + optional_sources=["github_blog"],
139 + started_at=iso_timestamp(NOW),
140 + crawler_code_sha="sha-test",
141 + )
142 + assert ctx["run_id"] == RUN_ID
143 + assert ctx["week"] == WEEK
144 + assert ctx["sources_requested"] == ["github_blog", "techcrunch"]
145 + assert ctx["required_sources"] == ["techcrunch"]
146 + assert ctx["optional_sources"] == ["github_blog"]
147 +
148 + def test_validate_run_context_valid(self) -> None:
149 + ctx = _make_run_context()
150 + validate_run_context(ctx) # Should not raise
151 +
152 + def test_validate_run_context_missing_keys(self) -> None:
153 + ctx = _make_run_context()
154 + del ctx["run_id"]
155 + with pytest.raises(FanInValidationError, match="missing keys"):
156 + validate_run_context(ctx)
157 +
158 + def test_validate_run_context_bad_schema_version(self) -> None:
159 + ctx = _make_run_context(schema_version=99)
160 + with pytest.raises(FanInValidationError, match="schema_version mismatch"):
161 + validate_run_context(ctx)
162 +
163 + def test_validate_run_context_bad_crawl_window(self) -> None:
164 + ctx = _make_run_context(crawl_window={"only_since": "x"})
165 + with pytest.raises(FanInValidationError, match="crawl_window"):
166 + validate_run_context(ctx)
167 +
168 +
169 +# --- Per-Source Artifact Tests ---
170 +
171 +
172 +class TestSourceArtifact:
173 + def test_build_source_artifact_structure(self) -> None:
174 + ctx = _make_run_context()
175 + articles = [_make_article("techcrunch", f"Art {i}") for i in range(3)]
176 + status = _make_status("techcrunch")
177 +
178 + artifact = build_source_artifact(
179 + source_id="techcrunch",
180 + articles=articles,
181 + status=status,
182 + run_context=ctx,
183 + crawled_at=NOW,
184 + )
185 +
186 + assert artifact["source_artifact_schema_version"] == SOURCE_ARTIFACT_SCHEMA_VERSION
187 + assert artifact["source_id"] == "techcrunch"
188 + assert artifact["crawled_at"] == iso_timestamp(NOW)
189 + assert artifact["run_context"]["run_id"] == RUN_ID
190 + assert artifact["run_context"]["week"] == WEEK
191 + assert artifact["status"]["success"] is True
192 + assert artifact["metrics"]["total_articles"] == 3
193 + assert artifact["metrics"]["relevant_articles"] == 3
194 + assert "artifact_checksum" in artifact
195 + assert len(artifact["artifact_checksum"]) == 64 # SHA-256 hex
196 +
197 + def test_source_artifact_deterministic(self) -> None:
198 + ctx = _make_run_context()
199 + articles = [_make_article("techcrunch", f"Art {i}") for i in range(3)]
200 + status = _make_status("techcrunch")
201 +
202 + a1 = build_source_artifact(
203 + source_id="techcrunch", articles=articles,
204 + status=status, run_context=ctx, crawled_at=NOW,
205 + )
206 + a2 = build_source_artifact(
207 + source_id="techcrunch", articles=articles,
208 + status=status, run_context=ctx, crawled_at=NOW,
209 + )
210 + assert a1["artifact_checksum"] == a2["artifact_checksum"]
211 + assert a1["articles"] == a2["articles"]
212 +
213 + def test_source_artifact_different_order_same_checksum(self) -> None:
214 + """Articles in different order produce same checksum (sorted internally)."""
215 + ctx = _make_run_context()
216 + articles = [_make_article("techcrunch", f"Art {i}") for i in range(3)]
217 + status = _make_status("techcrunch")
218 +
219 + a1 = build_source_artifact(
220 + source_id="techcrunch", articles=articles,
221 + status=status, run_context=ctx, crawled_at=NOW,
222 + )
223 + a2 = build_source_artifact(
224 + source_id="techcrunch", articles=list(reversed(articles)),
225 + status=status, run_context=ctx, crawled_at=NOW,
226 + )
227 + assert a1["artifact_checksum"] == a2["artifact_checksum"]
228 +
229 + def test_validate_source_artifact_valid(self) -> None:
230 + artifact = _make_source_artifact("techcrunch")
231 + validate_source_artifact(artifact) # Should not raise
232 +
233 + def test_validate_source_artifact_bad_schema(self) -> None:
234 + artifact = _make_source_artifact("techcrunch")
235 + artifact["source_artifact_schema_version"] = 99
236 + with pytest.raises(FanInValidationError, match="schema version mismatch"):
237 + validate_source_artifact(artifact)
238 +
239 + def test_validate_source_artifact_tampered_checksum(self) -> None:
240 + artifact = _make_source_artifact("techcrunch")
241 + artifact["artifact_checksum"] = "tampered"
242 + with pytest.raises(FanInValidationError, match="checksum mismatch"):
243 + validate_source_artifact(artifact)
244 +
245 + def test_validate_source_artifact_missing_keys(self) -> None:
246 + artifact = _make_source_artifact("techcrunch")
247 + del artifact["metrics"]
248 + with pytest.raises(FanInValidationError, match="missing keys"):
249 + validate_source_artifact(artifact)
250 +
251 +
252 +# --- Fan-In Merge Tests ---
253 +
254 +
255 +class TestMerge:
256 + def test_merge_basic(self) -> None:
257 + ctx = _make_run_context()
258 + artifacts = [
259 + _make_source_artifact("techcrunch", ctx),
260 + _make_source_artifact("github_blog", ctx),
261 + ]
262 +
263 + output, warnings = merge_source_artifacts(artifacts, ctx, merged_at=NOW)
264 +
265 + assert output["schema_version"] == CANONICAL_SCHEMA_VERSION
266 + assert output["source"] == "external_news"
267 + assert output["week"] == WEEK
268 + assert output["crawl_window"] == ctx["crawl_window"]
269 + assert output["metadata"]["sources_requested"] == ["github_blog", "nvidia_blog", "techcrunch"]
270 + assert "techcrunch" in output["metadata"]["sources_succeeded"]
271 + assert "github_blog" in output["metadata"]["sources_succeeded"]
272 + assert output["metadata"]["fan_in_mode"] == "matrix"
273 + assert output["metadata"]["total_articles"] >= 0
274 + assert output["metadata"]["artifact_checksum"]
275 + # nvidia_blog is optional and missing → warning
276 + assert any(w.source_id == "nvidia_blog" for w in warnings)
277 +
278 + def test_merge_deterministic(self) -> None:
279 + """Same inputs always produce same output (excluding crawled_at)."""
280 + ctx = _make_run_context()
281 + artifacts = [
282 + _make_source_artifact("techcrunch", ctx),
283 + _make_source_artifact("github_blog", ctx),
284 + ]
285 +
286 + out1, _ = merge_source_artifacts(artifacts, ctx, merged_at=NOW)
287 + out2, _ = merge_source_artifacts(artifacts, ctx, merged_at=NOW)
288 +
289 + assert out1["metadata"]["artifact_checksum"] == out2["metadata"]["artifact_checksum"]
290 + assert out1["articles"] == out2["articles"]
291 + assert out1["metadata"]["source_artifact_provenance"] == out2["metadata"]["source_artifact_provenance"]
292 +
293 + def test_merge_different_artifact_order_same_result(self) -> None:
294 + """Order of input artifacts doesn't affect output."""
295 + ctx = _make_run_context()
296 + a1 = _make_source_artifact("techcrunch", ctx)
297 + a2 = _make_source_artifact("github_blog", ctx)
298 +
299 + out_ab, _ = merge_source_artifacts([a1, a2], ctx, merged_at=NOW)
300 + out_ba, _ = merge_source_artifacts([a2, a1], ctx, merged_at=NOW)
301 +
302 + assert out_ab["metadata"]["artifact_checksum"] == out_ba["metadata"]["artifact_checksum"]
303 +
304 + def test_merge_deduplicates_across_sources(self) -> None:
305 + """Articles with same URL from different sources are deduped."""
306 + ctx = _make_run_context()
307 + shared_url = "https://example.com/shared-article"
308 + art_tc = _make_article("techcrunch", "Shared Article", shared_url)
309 + art_gh = _make_article("github_blog", "Shared Article", shared_url)
310 +
311 + a1 = build_source_artifact(
312 + source_id="techcrunch",
313 + articles=[art_tc],
314 + status=_make_status("techcrunch"),
315 + run_context=ctx,
316 + crawled_at=NOW,
317 + )
318 + a2 = build_source_artifact(
319 + source_id="github_blog",
320 + articles=[art_gh],
321 + status=_make_status("github_blog"),
322 + run_context=ctx,
323 + crawled_at=NOW,
324 + )
325 +
326 + output, _ = merge_source_artifacts([a1, a2], ctx, merged_at=NOW)
327 + assert output["metadata"]["dedupe_count"] == 1
328 + # The merged article preserves both sources
329 + merged_article = output["articles"][0]
330 + assert "techcrunch" in merged_article["sources"]
331 + assert "github_blog" in merged_article["sources"]
332 +
333 + def test_merge_with_failed_optional_source(self) -> None:
334 + """Optional source failure produces warning but valid output."""
335 + ctx = _make_run_context(
336 + required_sources=["techcrunch"],
337 + optional_sources=["github_blog"],
338 + )
339 + a1 = _make_source_artifact("techcrunch", ctx)
340 + a2 = _make_source_artifact("github_blog", ctx, articles=[], success=False)
341 +
342 + output, warnings = merge_source_artifacts([a1, a2], ctx, merged_at=NOW)
343 +
344 + assert "github_blog" in output["metadata"]["sources_failed"]
345 + assert any(w.category == "source_failure" and w.source_id == "github_blog" for w in warnings)
346 + # Output is still valid
347 + assert output["metadata"]["artifact_checksum"]
348 +
349 +
350 +# --- Fan-In Validation Tests ---
351 +
352 +
353 +class TestValidation:
354 + def test_rejects_empty_artifacts(self) -> None:
355 + ctx = _make_run_context()
356 + with pytest.raises(FanInValidationError, match="No source artifacts"):
357 + validate_fan_in_compatibility([], ctx)
358 +
359 + def test_rejects_schema_mismatch(self) -> None:
360 + ctx = _make_run_context()
361 + artifact = _make_source_artifact("techcrunch", ctx)
362 + artifact["run_context"]["schema_checksum"] = "wrong"
363 + # Recompute checksum after tampering
364 + from scripts.rss_fan_in import _source_artifact_checksum
365 + artifact["artifact_checksum"] = _source_artifact_checksum(artifact)
366 +
367 + with pytest.raises(FanInValidationError, match="Schema checksum mismatch"):
368 + validate_fan_in_compatibility([artifact], ctx)
369 +
370 + def test_rejects_window_mismatch(self) -> None:
371 + ctx = _make_run_context()
372 + artifact = _make_source_artifact("techcrunch", ctx)
373 + artifact["run_context"]["crawl_window"] = {"since": "wrong", "until": "wrong"}
374 + from scripts.rss_fan_in import _source_artifact_checksum
375 + artifact["artifact_checksum"] = _source_artifact_checksum(artifact)
376 +
377 + with pytest.raises(FanInValidationError, match="Crawl window mismatch"):
378 + validate_fan_in_compatibility([artifact], ctx)
379 +
380 + def test_rejects_config_checksum_mismatch(self) -> None:
381 + ctx = _make_run_context()
382 + artifact = _make_source_artifact("techcrunch", ctx)
383 + artifact["run_context"]["source_config_checksum"] = "wrong"
384 + from scripts.rss_fan_in import _source_artifact_checksum
385 + artifact["artifact_checksum"] = _source_artifact_checksum(artifact)
386 +
387 + with pytest.raises(FanInValidationError, match="Source config checksum mismatch"):
388 + validate_fan_in_compatibility([artifact], ctx)
389 +
390 + def test_rejects_run_id_mismatch(self) -> None:
391 + ctx = _make_run_context()
392 + artifact = _make_source_artifact("techcrunch", ctx)
393 + artifact["run_context"]["run_id"] = "different-run"
394 + from scripts.rss_fan_in import _source_artifact_checksum
395 + artifact["artifact_checksum"] = _source_artifact_checksum(artifact)
396 +
397 + with pytest.raises(FanInValidationError, match="Run ID mismatch"):
398 + validate_fan_in_compatibility([artifact], ctx)
399 +
400 + def test_rejects_missing_required_sources(self) -> None:
401 + ctx = _make_run_context(required_sources=["techcrunch", "github_blog"])
402 + # Only provide techcrunch
403 + artifact = _make_source_artifact("techcrunch", ctx)
404 +
405 + with pytest.raises(FanInValidationError, match="Missing required source"):
406 + validate_fan_in_compatibility([artifact], ctx)
407 +
408 + def test_rejects_duplicate_sources(self) -> None:
409 + ctx = _make_run_context(required_sources=["techcrunch"])
410 + a1 = _make_source_artifact("techcrunch", ctx)
411 + a2 = _make_source_artifact("techcrunch", ctx)
412 +
413 + with pytest.raises(FanInValidationError, match="Duplicate source"):
414 + validate_fan_in_compatibility([a1, a2], ctx)
415 +
416 + def test_warns_missing_optional_source(self) -> None:
417 + ctx = _make_run_context(
418 + required_sources=["techcrunch"],
419 + optional_sources=["nvidia_blog"],
420 + )
421 + artifact = _make_source_artifact("techcrunch", ctx)
422 +
423 + warnings = validate_fan_in_compatibility([artifact], ctx)
424 + assert len(warnings) == 1
425 + assert warnings[0].source_id == "nvidia_blog"
426 + assert warnings[0].category == "missing_optional_source"
427 +
428 +
429 +# --- CLI Integration Tests ---
430 +
431 +
432 +class TestCLI:
433 + def test_emit_and_merge_roundtrip(self, tmp_path: Path) -> None:
434 + """Full roundtrip: emit per-source artifacts then merge them."""
435 + from scripts.rss_fan_in import main as fan_in_main
436 +
437 + ctx = _make_run_context(
438 + required_sources=["techcrunch", "github_blog"],
439 + optional_sources=[],
440 + )
441 + ctx_path = tmp_path / "run-context.json"
442 + ctx_path.write_text(json.dumps(ctx), encoding="utf-8")
443 +
444 + artifacts_dir = tmp_path / "artifacts"
445 + artifacts_dir.mkdir()
446 +
447 + # Emit two source artifacts
448 + for source_id in ["techcrunch", "github_blog"]:
449 + articles = [_make_article(source_id, f"Art {i}") for i in range(2)]
450 + articles_path = tmp_path / f"{source_id}-articles.json"
451 + articles_path.write_text(json.dumps(articles), encoding="utf-8")
452 +
453 + status = _make_status(source_id)
454 + status_path = tmp_path / f"{source_id}-status.json"
455 + status_path.write_text(json.dumps(status), encoding="utf-8")
456 +
457 + result = fan_in_main([
458 + "emit",
459 + "--source", source_id,
460 + "--articles", str(articles_path),
461 + "--status", str(status_path),
462 + "--run-context", str(ctx_path),
463 + "--output", str(artifacts_dir / f"{source_id}.json"),
464 + ])
465 + assert result == 0
466 +
467 + # Verify artifacts were created
468 + assert (artifacts_dir / "techcrunch.json").exists()
469 + assert (artifacts_dir / "github_blog.json").exists()
470 +
471 + # Merge
472 + merged_path = tmp_path / "merged.json"
473 + result = fan_in_main([
474 + "merge",
475 + "--artifacts-dir", str(artifacts_dir),
476 + "--run-context", str(ctx_path),
477 + "--output", str(merged_path),
478 + ])
479 + assert result == 0
480 + assert merged_path.exists()
481 +
482 + # Validate merged output
483 + merged = json.loads(merged_path.read_text(encoding="utf-8"))
484 + assert merged["schema_version"] == CANONICAL_SCHEMA_VERSION
485 + assert merged["source"] == "external_news"
486 + assert merged["metadata"]["fan_in_mode"] == "matrix"
487 + assert merged["metadata"]["total_articles"] == 4
488 + assert "techcrunch" in merged["metadata"]["sources_succeeded"]
489 + assert "github_blog" in merged["metadata"]["sources_succeeded"]
490 +
491 + def test_validate_command(self, tmp_path: Path) -> None:
492 + """Validate subcommand checks artifacts without merging."""
493 + from scripts.rss_fan_in import main as fan_in_main
494 +
495 + ctx = _make_run_context(required_sources=["techcrunch"], optional_sources=[])
496 + ctx_path = tmp_path / "run-context.json"
497 + ctx_path.write_text(json.dumps(ctx), encoding="utf-8")
498 +
499 + artifacts_dir = tmp_path / "artifacts"
500 + artifacts_dir.mkdir()
501 +
502 + artifact = _make_source_artifact("techcrunch", ctx)
503 + (artifacts_dir / "techcrunch.json").write_text(
504 + json.dumps(artifact), encoding="utf-8"
505 + )
506 +
507 + result = fan_in_main([
508 + "validate",
509 + "--artifacts-dir", str(artifacts_dir),
510 + "--run-context", str(ctx_path),
511 + ])
512 + assert result == 0
513 +
514 + def test_merge_fails_on_missing_required(self, tmp_path: Path) -> None:
515 + """Merge fails when required source is missing."""
516 + from scripts.rss_fan_in import main as fan_in_main
517 +
518 + ctx = _make_run_context(
519 + required_sources=["techcrunch", "github_blog"],
520 + optional_sources=[],
521 + )
522 + ctx_path = tmp_path / "run-context.json"
523 + ctx_path.write_text(json.dumps(ctx), encoding="utf-8")
524 +
525 + artifacts_dir = tmp_path / "artifacts"
526 + artifacts_dir.mkdir()
527 +
528 + # Only provide techcrunch
529 + artifact = _make_source_artifact("techcrunch", ctx)
530 + (artifacts_dir / "techcrunch.json").write_text(
531 + json.dumps(artifact), encoding="utf-8"
532 + )
533 +
534 + result = fan_in_main([
535 + "merge",
536 + "--artifacts-dir", str(artifacts_dir),
537 + "--run-context", str(ctx_path),
538 + "--output", str(tmp_path / "merged.json"),
539 + ])
540 + assert result == 1 # Fails due to missing required source
541 +
542 +
543 +# --- Determinism Proof ---
544 +
545 +
546 +class TestDeterminism:
547 + """Fixture-based proof that same inputs → same merged output."""
548 +
549 + def test_deterministic_merge_with_fixed_fixtures(self, tmp_path: Path) -> None:
550 + """Given fixed input artifacts, merge always produces identical output."""
551 + ctx = _make_run_context(
552 + required_sources=["techcrunch", "github_blog", "nvidia_blog"],
553 + optional_sources=[],
554 + )
555 +
556 + # Create deterministic articles
557 + articles_by_source = {}
558 + for source_id in ["techcrunch", "github_blog", "nvidia_blog"]:
559 + articles_by_source[source_id] = [
560 + {
561 + "source": source_id,
562 + "title": f"{source_id} Article {i}",
563 + "url": f"https://{source_id}.example.com/article-{i}",
564 + "published_at": "2026-06-12T10:00:00Z",
565 + "categories": ["AI"],
566 + "summary": f"Summary for {source_id} article {i}",
567 + "github_links": [],
568 + "entities": [],
569 + "relevance_score": 0.8,
570 + }
571 + for i in range(3)
572 + ]
573 +
574 + artifacts = [
575 + build_source_artifact(
576 + source_id=sid,
577 + articles=articles_by_source[sid],
578 + status=_make_status(sid),
579 + run_context=ctx,
580 + crawled_at=NOW,
581 + )
582 + for sid in ["techcrunch", "github_blog", "nvidia_blog"]
583 + ]
584 +
585 + # Merge 10 times and verify all produce identical output
586 + checksums = set()
587 + for _ in range(10):
588 + output, _ = merge_source_artifacts(artifacts, ctx, merged_at=NOW)
589 + checksums.add(output["metadata"]["artifact_checksum"])
590 +
591 + assert len(checksums) == 1, f"Non-deterministic merge: got {len(checksums)} distinct checksums"