Skip to content

Commit d60cc08

Browse files
committed
chg: [crawler] add interactive browser crawler capture. user can open a live browser to crawl content
1 parent 8fd997b commit d60cc08

8 files changed

Lines changed: 597 additions & 7 deletions

File tree

bin/crawlers/Crawler.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -392,6 +392,7 @@ def compute(self, capture): # TODO ADD FUNCTION TO MANUALLY IMPORT ???
392392
if crawlers.is_domain_correlation_cache(self.original_domain.id):
393393
crawlers.save_domain_correlation_cache(self.original_domain.was_up(), domain)
394394

395+
crawlers.release_interactive_session_by_capture(capture.uuid, status='completed')
395396
task.remove()
396397
self.root_item = None
397398

bin/lib/crawlers.py

Lines changed: 304 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2318,6 +2318,310 @@ def can_launch_forum_crawler_account():
23182318
return get_nb_running_forum_crawler_accounts() < get_forum_crawler_max_accounts()
23192319

23202320

2321+
#### INTERACTIVE CRAWLER SESSIONS ####
2322+
2323+
INTERACTIVE_SESSION_TTL = 3600
2324+
INTERACTIVE_SESSION_META_TTL = 3600
2325+
INTERACTIVE_ACTIVE_STATES = {'starting', 'ready', 'finishing'}
2326+
INTERACTIVE_FINAL_STATES = {'completed', 'expired', 'error', 'closed'}
2327+
2328+
def get_max_interactive_crawler():
2329+
nb = r_cache.hget('crawler:lacus', 'max_interactive_crawler')
2330+
if not nb:
2331+
nb = r_db.hget('crawler:lacus', 'max_interactive_crawler')
2332+
if not nb:
2333+
nb = 1
2334+
save_max_interactive_crawler(nb)
2335+
else:
2336+
r_cache.hset('crawler:lacus', 'max_interactive_crawler', int(nb))
2337+
return int(nb)
2338+
2339+
def save_max_interactive_crawler(nb):
2340+
r_db.hset('crawler:lacus', 'max_interactive_crawler', int(nb))
2341+
r_cache.hset('crawler:lacus', 'max_interactive_crawler', int(nb))
2342+
2343+
def api_set_max_interactive_crawler(data):
2344+
nb = data.get('nb', 1)
2345+
try:
2346+
nb = int(nb)
2347+
if nb < 0:
2348+
nb = 0
2349+
except (TypeError, ValueError):
2350+
return {'error': 'Invalid number of interactive crawler sessions'}, 400
2351+
save_max_interactive_crawler(nb)
2352+
return nb, 200
2353+
2354+
def _cleanup_interactive_task_capture(task_uuid=None, capture_uuid=None):
2355+
if capture_uuid:
2356+
capture = CrawlerCapture(capture_uuid)
2357+
if capture.exists():
2358+
capture.delete()
2359+
if task_uuid:
2360+
task = CrawlerTask(task_uuid)
2361+
if task.exists():
2362+
task.delete()
2363+
2364+
def cleanup_stale_interactive_sessions(now=None):
2365+
if now is None:
2366+
now = int(time.time())
2367+
for session_uuid, launch_time in r_cache.zrange('crawler:interactive:sessions', 0, -1, withscores=True):
2368+
session = InteractiveCrawlerSession(session_uuid)
2369+
if not session.exists():
2370+
r_cache.srem('crawler:interactive:active', session_uuid)
2371+
r_cache.zrem('crawler:interactive:sessions', session_uuid)
2372+
continue
2373+
status = session.get_status()
2374+
if status in INTERACTIVE_FINAL_STATES:
2375+
continue
2376+
if now - int(launch_time) > INTERACTIVE_SESSION_TTL:
2377+
session.expire()
2378+
2379+
def get_nb_active_interactive_sessions():
2380+
cleanup_stale_interactive_sessions()
2381+
return r_cache.scard('crawler:interactive:active')
2382+
2383+
def get_interactive_usage():
2384+
return {'active': get_nb_active_interactive_sessions(), 'max': get_max_interactive_crawler()}
2385+
2386+
def get_interactive_session_by_capture(capture_uuid):
2387+
session_uuid = r_cache.hget('crawler:interactive:captures', capture_uuid)
2388+
if session_uuid:
2389+
return InteractiveCrawlerSession(session_uuid)
2390+
for candidate in r_cache.zrange('crawler:interactive:sessions', 0, -1):
2391+
session = InteractiveCrawlerSession(candidate)
2392+
if session.get_capture_uuid() == capture_uuid:
2393+
return session
2394+
return None
2395+
2396+
def release_interactive_session_by_capture(capture_uuid, status='completed'):
2397+
session = get_interactive_session_by_capture(capture_uuid)
2398+
if session and session.exists():
2399+
session.release(status=status)
2400+
2401+
def get_active_interactive_sessions():
2402+
cleanup_stale_interactive_sessions()
2403+
sessions = []
2404+
for session_uuid in r_cache.smembers('crawler:interactive:active'):
2405+
session = InteractiveCrawlerSession(session_uuid)
2406+
if session.exists():
2407+
sessions.append(session.get_meta())
2408+
return sorted(sessions, key=lambda m: m.get('launch_time', 0))
2409+
2410+
def get_user_active_interactive_session(user_id):
2411+
cleanup_stale_interactive_sessions()
2412+
session_uuid = r_cache.hget('crawler:interactive:users', user_id)
2413+
if session_uuid:
2414+
session = InteractiveCrawlerSession(session_uuid)
2415+
if session.is_active():
2416+
return session
2417+
r_cache.hdel('crawler:interactive:users', user_id)
2418+
return None
2419+
2420+
def reserve_interactive_session(user_id, url, task_uuid=None):
2421+
cleanup_stale_interactive_sessions()
2422+
if get_user_active_interactive_session(user_id):
2423+
return None, {'error': 'User already has an active interactive session'}, 409
2424+
max_sessions = get_max_interactive_crawler()
2425+
if max_sessions <= 0:
2426+
return None, {'error': 'Interactive crawler sessions are disabled'}, 403
2427+
session_uuid = gen_uuid()
2428+
launch_time = int(time.time())
2429+
if r_cache.hget('crawler:interactive:users', user_id):
2430+
return None, {'error': 'User already has an active interactive session'}, 409
2431+
if r_cache.scard('crawler:interactive:active') >= max_sessions:
2432+
return None, {'error': 'No interactive crawler slots available'}, 429
2433+
r_cache.hset(f'crawler:interactive:session:{session_uuid}', mapping={'user': user_id, 'url': url, 'status': 'starting', 'launch_time': launch_time})
2434+
if task_uuid:
2435+
r_cache.hset(f'crawler:interactive:session:{session_uuid}', 'task_uuid', task_uuid)
2436+
r_cache.expire(f'crawler:interactive:session:{session_uuid}', INTERACTIVE_SESSION_META_TTL)
2437+
r_cache.hset('crawler:interactive:users', user_id, session_uuid)
2438+
r_cache.sadd('crawler:interactive:active', session_uuid)
2439+
r_cache.zadd('crawler:interactive:sessions', {session_uuid: launch_time})
2440+
return InteractiveCrawlerSession(session_uuid), None, 200
2441+
2442+
2443+
def _remote_headed_response_to_meta(response):
2444+
if response is None:
2445+
return {}
2446+
if isinstance(response, dict):
2447+
return response
2448+
meta = {}
2449+
for field in ('uuid', 'status', 'raw_status', 'finish_requested', 'view_url', 'created_at', 'expires_at', 'error'):
2450+
if hasattr(response, field):
2451+
meta[field] = getattr(response, field)
2452+
return meta
2453+
2454+
def refresh_interactive_session_status(session):
2455+
capture_uuid = session.get_capture_uuid()
2456+
if not capture_uuid:
2457+
return session.get_meta()
2458+
try:
2459+
lacus = get_lacus()
2460+
remote = _remote_headed_response_to_meta(lacus.get_remote_headed_session(capture_uuid))
2461+
if remote.get('status'):
2462+
session.set('remote_status', remote['status'])
2463+
if remote.get('raw_status') is not None:
2464+
session.set('remote_raw_status', remote['raw_status'])
2465+
if remote.get('finish_requested') is not None:
2466+
session.set('finish_requested', str(remote['finish_requested']))
2467+
if remote.get('view_url'):
2468+
session.set('remote_url', remote['view_url'])
2469+
if session.get_status() == 'starting':
2470+
session.set('status', 'ready')
2471+
if remote.get('expires_at'):
2472+
session.set('expires_at', remote['expires_at'])
2473+
if remote.get('error'):
2474+
session.set('error', remote['error'])
2475+
session.release(status='error')
2476+
capture_status = lacus.get_capture_status(capture_uuid)
2477+
session.set('capture_status', int(capture_status))
2478+
except Exception as e:
2479+
session.set('last_status_error', str(e))
2480+
return session.get_meta()
2481+
2482+
class InteractiveCrawlerSession:
2483+
def __init__(self, session_uuid):
2484+
self.uuid = session_uuid
2485+
2486+
def exists(self):
2487+
return r_cache.exists(f'crawler:interactive:session:{self.uuid}')
2488+
2489+
def get(self, field):
2490+
return r_cache.hget(f'crawler:interactive:session:{self.uuid}', field)
2491+
2492+
def set(self, field, value):
2493+
return r_cache.hset(f'crawler:interactive:session:{self.uuid}', field, value)
2494+
2495+
def get_user(self):
2496+
return self.get('user')
2497+
2498+
def get_status(self):
2499+
return self.get('status') or 'unknown'
2500+
2501+
def is_active(self):
2502+
return self.exists() and self.get_status() in INTERACTIVE_ACTIVE_STATES
2503+
2504+
def get_capture_uuid(self):
2505+
return self.get('capture_uuid')
2506+
2507+
def get_task_uuid(self):
2508+
return self.get('task_uuid')
2509+
2510+
def get_meta(self):
2511+
meta = r_cache.hgetall(f'crawler:interactive:session:{self.uuid}')
2512+
meta['uuid'] = self.uuid
2513+
try:
2514+
meta['launch_time'] = int(meta.get('launch_time', 0))
2515+
except (TypeError, ValueError):
2516+
meta['launch_time'] = 0
2517+
return meta
2518+
2519+
def release(self, status='completed'):
2520+
user = self.get_user()
2521+
task_uuid = self.get_task_uuid()
2522+
capture_uuid = self.get_capture_uuid()
2523+
self.set('status', status)
2524+
self.set('end_time', int(time.time()))
2525+
r_cache.expire(f'crawler:interactive:session:{self.uuid}', INTERACTIVE_SESSION_META_TTL)
2526+
r_cache.srem('crawler:interactive:active', self.uuid)
2527+
r_cache.zrem('crawler:interactive:sessions', self.uuid)
2528+
if capture_uuid:
2529+
r_cache.hdel('crawler:interactive:captures', capture_uuid)
2530+
if user:
2531+
r_cache.hdel('crawler:interactive:users', user)
2532+
if status in {'error', 'expired', 'closed'}:
2533+
_cleanup_interactive_task_capture(task_uuid=task_uuid, capture_uuid=capture_uuid)
2534+
2535+
def expire(self):
2536+
self.release(status='expired')
2537+
2538+
def api_start_interactive_capture(data, user_org, user_id):
2539+
task, resp = api_parse_task_dict_basic(data, user_id)
2540+
if resp != 200:
2541+
return task, resp
2542+
if task.get('urls'):
2543+
return {'error': 'Interactive capture accepts only one URL'}, 400
2544+
task['depth_limit'] = 0
2545+
filter_local_ips_error = api_validate_global_urls(url=task.get('url'))
2546+
if filter_local_ips_error:
2547+
return filter_local_ips_error
2548+
session, error, code = reserve_interactive_session(user_id, task['url'])
2549+
if error:
2550+
return error, code
2551+
try:
2552+
task_uuid = create_task(task['url'], depth=0, har=task['har'], screenshot=task['screenshot'], proxy=task['proxy'], tags=task['tags'], parent='interactive', priority=90, external=True)
2553+
if not task_uuid:
2554+
session.release(status='error')
2555+
return {'error': 'Aborted by Crawler'}, 400
2556+
session.set('task_uuid', task_uuid)
2557+
capture_uuid = session.uuid
2558+
lacus = get_lacus()
2559+
returned_uuid = lacus.enqueue(url=task['url'], depth=0, proxy=task['proxy'], with_favicon=True, force=True, uuid=capture_uuid, remote_headfull=True, general_timeout_in_sec=90)
2560+
capture_uuid = returned_uuid or capture_uuid
2561+
session.set('capture_uuid', capture_uuid)
2562+
r_cache.hset('crawler:interactive:captures', capture_uuid, session.uuid)
2563+
create_capture(capture_uuid, task_uuid)
2564+
CrawlerTask(task_uuid).start()
2565+
refresh_interactive_session_status(session)
2566+
return session.get_meta(), 200
2567+
except Exception as e:
2568+
session.set('error', str(e))
2569+
session.release(status='error')
2570+
return {'error': 'Unable to start interactive capture', 'details': str(e)}, 502
2571+
2572+
def api_get_interactive_session(session_uuid, user_id, is_admin=False):
2573+
session = InteractiveCrawlerSession(session_uuid)
2574+
if not session.exists():
2575+
return {'error': 'Unknown interactive session'}, 404
2576+
if not is_admin and session.get_user() != user_id:
2577+
return {'error': 'Forbidden'}, 403
2578+
return refresh_interactive_session_status(session), 200
2579+
2580+
def api_finish_interactive_session(session_uuid, user_id):
2581+
session = InteractiveCrawlerSession(session_uuid)
2582+
if not session.exists():
2583+
return {'error': 'Unknown interactive session'}, 404
2584+
if session.get_user() != user_id:
2585+
return {'error': 'Forbidden'}, 403
2586+
session.set('status', 'finishing')
2587+
capture_uuid = session.get_capture_uuid()
2588+
task_uuid = session.get_task_uuid()
2589+
try:
2590+
lacus = get_lacus()
2591+
remote = _remote_headed_response_to_meta(lacus.finish_remote_headed_session(capture_uuid))
2592+
if remote.get('status'):
2593+
session.set('remote_status', remote['status'])
2594+
if remote.get('finish_requested') is not None:
2595+
session.set('finish_requested', str(remote['finish_requested']))
2596+
if remote.get('view_url'):
2597+
session.set('remote_url', remote['view_url'])
2598+
except Exception as e:
2599+
session.set('error', str(e))
2600+
if capture_uuid and task_uuid:
2601+
refresh_interactive_session_status(session)
2602+
return session.get_meta(), 200
2603+
2604+
def api_admin_close_interactive_session(session_uuid):
2605+
session = InteractiveCrawlerSession(session_uuid)
2606+
if not session.exists():
2607+
return {'error': 'Unknown interactive session'}, 404
2608+
capture_uuid = session.get_capture_uuid()
2609+
if capture_uuid:
2610+
try:
2611+
lacus = get_lacus()
2612+
remote = _remote_headed_response_to_meta(lacus.finish_remote_headed_session(capture_uuid))
2613+
if remote.get('status'):
2614+
session.set('remote_status', remote['status'])
2615+
if remote.get('finish_requested') is not None:
2616+
session.set('finish_requested', str(remote['finish_requested']))
2617+
if remote.get('view_url'):
2618+
session.set('remote_url', remote['view_url'])
2619+
except Exception as e:
2620+
session.set('error', str(e))
2621+
session.release(status='closed')
2622+
return session.get_meta(), 200
2623+
2624+
23212625
#### CRAWLER CAPTURE ####
23222626

23232627
def get_nb_crawler_captures():

0 commit comments

Comments
 (0)