1"""Remote Job Feed Aggregator - one search across five public remote-job feeds."""
2
3from __future__ import annotations
4
5import asyncio
6from datetime import datetime, timedelta, timezone
7from typing import Any
8
9import httpx
10from apify import Actor
11
12from . import feeds
13
14USER_AGENT = (
15 'Mozilla/5.0 (compatible; RemoteJobFeedAggregator/0.1; '
16 '+https://apify.com/store) apify-actor'
17)
18RETRY_STATUS = {429, 500, 502, 503, 504}
19MAX_ATTEMPTS = 3
20
21
22async def get(client: httpx.AsyncClient, url: str) -> tuple[httpx.Response | None, str | None]:
23 delay = 1.0
24 for attempt in range(1, MAX_ATTEMPTS + 1):
25 try:
26 response = await client.get(url)
27 except httpx.HTTPError as exc:
28 if attempt == MAX_ATTEMPTS:
29 return None, f'network error: {type(exc).__name__}'
30 await asyncio.sleep(delay)
31 delay *= 2
32 continue
33 if response.status_code in RETRY_STATUS:
34 if attempt == MAX_ATTEMPTS:
35 return None, f'HTTP {response.status_code} after {MAX_ATTEMPTS} attempts'
36 await asyncio.sleep(delay)
37 delay *= 2
38 continue
39 if response.status_code >= 400:
40 return None, f'HTTP {response.status_code}'
41 return response, None
42 return None, 'exhausted retries'
43
44
45async def collect_remotive(client, pages, delay, include_description):
46 items, errors = [], []
47 response, error = await get(client, 'https://remotive.com/api/remote-jobs')
48 if error:
49 errors.append(f'remotive: {error}')
50 return items, errors
51 try:
52 items = feeds.parse_remotive(response.json(), include_description)
53 except ValueError:
54 errors.append('remotive: invalid JSON')
55 return items, errors
56
57
58async def collect_himalayas(client, pages, delay, include_description):
59 items, errors, cursor, offset = [], [], None, 0
60 for page in range(pages):
61
62 url = 'https://himalayas.app/jobs/api?limit=20'
63 url += f'&cursor={cursor}' if cursor else f'&offset={offset}'
64 response, error = await get(client, url)
65 if error:
66 errors.append(f'himalayas: {error}')
67 break
68 try:
69 payload = response.json()
70 except ValueError:
71 errors.append('himalayas: invalid JSON')
72 break
73 batch = feeds.parse_himalayas(payload, include_description)
74 items.extend(batch)
75 cursor = payload.get('nextCursor')
76 offset += 20
77 if not batch or (not cursor and offset >= (payload.get('totalCount') or 0)):
78 break
79 if page < pages - 1:
80 await asyncio.sleep(delay)
81 return items, errors
82
83
84async def collect_arbeitnow(client, pages, delay, include_description):
85 items, errors = [], []
86 for page in range(1, pages + 1):
87 response, error = await get(
88 client, f'https://www.arbeitnow.com/api/job-board-api?page={page}'
89 )
90 if error:
91 errors.append(f'arbeitnow: {error}')
92 break
93 try:
94 payload = response.json()
95 except ValueError:
96 errors.append('arbeitnow: invalid JSON')
97 break
98 batch = feeds.parse_arbeitnow(payload, include_description)
99 items.extend(batch)
100 if not batch:
101 break
102 if page < pages:
103 await asyncio.sleep(delay)
104 return items, errors
105
106
107async def collect_wwr(client, pages, delay, include_description):
108 items, errors = [], []
109 categories = feeds.WWR_CATEGORIES[: max(1, min(pages * 2, len(feeds.WWR_CATEGORIES)))]
110 for index, category in enumerate(categories):
111 response, error = await get(
112 client, f'https://weworkremotely.com/categories/{category}.rss'
113 )
114 if error:
115 errors.append(f'weworkremotely/{category}: {error}')
116 continue
117 items.extend(feeds.parse_wwr_rss(response.text, category, include_description))
118 if index < len(categories) - 1:
119 await asyncio.sleep(delay)
120 return items, errors
121
122
123async def collect_hn(client, pages, delay, include_description):
124 """Find the newest 'Ask HN: Who is hiring?' story, then read its comments."""
125 items, errors = [], []
126 response, error = await get(
127 client,
128 'https://hn.algolia.com/api/v1/search?tags=story,author_whoishiring'
129 '&query=who%20is%20hiring&hitsPerPage=1&restrictSearchableAttributes=title',
130 )
131 if error:
132 errors.append(f'hackernews: {error}')
133 return items, errors
134 try:
135 hits = response.json().get('hits') or []
136 except ValueError:
137 errors.append('hackernews: invalid JSON')
138 return items, errors
139 if not hits:
140 errors.append('hackernews: no "Who is hiring" thread found')
141 return items, errors
142 story_id = hits[0].get('objectID')
143 await asyncio.sleep(delay)
144 for page in range(pages):
145 response, error = await get(
146 client,
147 f'https://hn.algolia.com/api/v1/search_by_date?tags=comment,story_{story_id}'
148 f'&hitsPerPage=100&page={page}',
149 )
150 if error:
151 errors.append(f'hackernews: {error}')
152 break
153 try:
154 payload = response.json()
155 except ValueError:
156 errors.append('hackernews: invalid JSON')
157 break
158 batch = feeds.parse_hackernews(payload, include_description)
159 items.extend(batch)
160 if page + 1 >= (payload.get('nbPages') or 1) or not batch:
161 break
162 if page < pages - 1:
163 await asyncio.sleep(delay)
164 return items, errors
165
166
167COLLECTORS = {
168 'remotive': collect_remotive,
169 'himalayas': collect_himalayas,
170 'arbeitnow': collect_arbeitnow,
171 'weworkremotely': collect_wwr,
172 'hackernews': collect_hn,
173}
174
175
176def matches(item: dict[str, Any], cfg: dict[str, Any]) -> bool:
177 title = (item.get('title') or '').lower()
178 if cfg['exclude'] and any(k in title for k in cfg['exclude']):
179 return False
180 if cfg['keywords']:
181 haystack = ' '.join(
182 str(x).lower()
183 for x in (
184 item.get('title'),
185 item.get('companyName'),
186 item.get('category'),
187 item.get('seniority'),
188 item.get('_excerpt'),
189 item.get('descriptionText'),
190 ' '.join(item.get('tags') or []),
191 )
192 if x
193 )
194 if not any(k in haystack for k in cfg['keywords']):
195 return False
196 if cfg['locations']:
197 haystack = ' '.join(
198 str(x).lower()
199 for x in ([item.get('locationText')] + list(item.get('locationRestrictions') or []))
200 if x
201 )
202 if not any(k in haystack for k in cfg['locations']):
203 return False
204 if cfg['cutoff'] is not None and item.get('postedAt'):
205 try:
206 if datetime.fromisoformat(item['postedAt']) < cfg['cutoff']:
207 return False
208 except ValueError:
209 pass
210 return True
211
212
213async def main() -> None:
214 async with Actor:
215 actor_input = await Actor.get_input() or {}
216
217 sources = [s for s in (actor_input.get('sources') or list(feeds.SOURCES)) if s in COLLECTORS]
218 if not sources:
219 sources = ['remotive', 'himalayas', 'arbeitnow', 'weworkremotely']
220 max_items = int(actor_input.get('maxItems') or 200)
221 pages = int(actor_input.get('maxPagesPerSource') or 3)
222 delay = float(actor_input.get('requestDelaySeconds') or 1)
223 include_description = bool(actor_input.get('includeDescription', False))
224 deduplicate = bool(actor_input.get('deduplicate', True))
225 posted_within = int(actor_input.get('postedWithinDays') or 0)
226
227 cfg = {
228 'keywords': [k.lower() for k in (actor_input.get('keywords') or []) if k],
229 'exclude': [k.lower() for k in (actor_input.get('excludeKeywords') or []) if k],
230 'locations': [k.lower() for k in (actor_input.get('locationKeywords') or []) if k],
231 'cutoff': (
232 datetime.now(timezone.utc) - timedelta(days=posted_within)
233 if posted_within > 0
234 else None
235 ),
236 }
237
238 Actor.log.info('Querying %d feed(s): %s', len(sources), ', '.join(sources))
239
240 collected: list[dict[str, Any]] = []
241 errors: list[str] = []
242 per_source: dict[str, int] = {}
243 timeout = httpx.Timeout(60.0, connect=15.0)
244 headers = {'User-Agent': USER_AGENT, 'Accept': 'application/json, application/xml, */*'}
245
246 async with httpx.AsyncClient(
247 timeout=timeout, headers=headers, follow_redirects=True
248 ) as client:
249 for index, source in enumerate(sources):
250 if index:
251 await asyncio.sleep(delay)
252 try:
253 items, source_errors = await COLLECTORS[source](
254 client, pages, delay, include_description
255 )
256 except Exception as exc:
257 Actor.log.exception('Feed %s failed', source)
258 errors.append(f'{source}: unexpected {type(exc).__name__}')
259 continue
260 errors.extend(source_errors)
261 per_source[source] = len(items)
262 collected.extend(items)
263 Actor.log.info('%s -> %d job(s) returned', source, len(items))
264
265 kept = [item for item in collected if matches(item, cfg)]
266 Actor.log.info('%d of %d job(s) match the filters', len(kept), len(collected))
267
268 if deduplicate:
269 merged: dict[str, dict[str, Any]] = {}
270 for item in kept:
271 key = feeds.dedup_key(item)
272 if not key:
273 continue
274 if key in merged:
275 if item['source'] not in merged[key]['alsoSeenOn']:
276 merged[key]['alsoSeenOn'].append(item['source'])
277 else:
278 merged[key] = item
279 kept = list(merged.values())
280 Actor.log.info('%d job(s) after deduplication', len(kept))
281
282 kept.sort(key=lambda x: x.get('postedAt') or '', reverse=True)
283 kept = kept[:max_items]
284
285 now = datetime.now(timezone.utc).isoformat()
286 for item in kept:
287 item.pop('_excerpt', None)
288 item['scrapedAt'] = now
289
290 if kept:
291 await Actor.push_data(kept)
292 else:
293 Actor.log.warning(
294 'No job matched. Try fewer keywords, more sources, or a larger '
295 '"Max pages per source".'
296 )
297
298 await Actor.set_value('RUN_SUMMARY', {
299 'sources': sources,
300 'returnedPerSource': per_source,
301 'matched': len(kept),
302 'errors': errors,
303 'finishedAt': now,
304 })
305 Actor.log.info('Done. Pushed %d job(s). %d feed error(s).', len(kept), len(errors))