# -*- coding: utf-8 -*- # Copyright (c) 2025 relakkes@gmail.com # # This file is part of MediaCrawler project. # Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/monitor/ingest.py # GitHub: https://github.com/NanmiCoder # Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1 # # 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则: # 1. 不得用于任何商业用途。 # 2. 使用时应遵守目标平台的使用条款和robots.txt规则。 # 3. 不得进行大规模爬取或对平台造成运营干扰。 # 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。 # 5. 不得用于任何非法或不当的用途。 # # 详细许可条款请参阅项目根目录下的LICENSE文件。 # 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。 """Turn one run's crawled jsonl into snapshots and change events. Pure-ish and offline testable: give it a directory of jsonl files, a run row and a session, and it does the diffing. No network, no browser. Correctness notes that drive the code below: * Counts arrive as strings and may be abbreviated ("1.2万", "3亿"). A value that cannot be parsed is stored as NULL, never 0 -- 0 would forge a large negative delta on the next comparison. * The comment endpoint has no time-sort, so only the platform's top-N window is ever visible. A comment we have not seen before is therefore split into "posted since last run" vs "seen for the first time", rather than claiming the former always. * A bad cookie does not make the crawler exit non-zero; it exits 0 having fetched nothing. That is detected here as a suspected auth failure. """ import json import re from dataclasses import dataclass, field from pathlib import Path from typing import Any, Dict, List, Optional from sqlalchemy import func, select from sqlalchemy.ext.asyncio import AsyncSession from tools.time_util import get_current_timestamp from .platforms import PLATFORM_XHS from .models import ( EVENT_AUTH_FAILURE, EVENT_METRIC_DELTA, EVENT_NEW_COMMENT_POSTED, EVENT_NEW_COMMENT_SEEN, EVENT_NEW_NOTE, EVENT_NO_DATA, EVENT_RUN_FAILED, MonitorComment, MonitorEvent, MonitorNote, MonitorNoteMetric, MonitorRun, MonitorTask, RUN_FAILED, RUN_PARTIAL, RUN_SUCCESS, ) _COUNT_UNITS = { "": 1, "万": 10_000, "w": 10_000, "W": 10_000, "k": 1_000, "K": 1_000, "亿": 100_000_000, } _COUNT_RE = re.compile(r"^([\d.]+)\s*([万wWkK亿]?)$") # Metric fields shared by the snapshot table and the delta comparison. _METRIC_FIELDS = ("liked_count", "comment_count", "collected_count", "share_count") def parse_count(value: Any) -> Optional[int]: """Parse an XHS interaction count into an int, or None if unintelligible. Handles plain numbers, thousands separators, and the Chinese abbreviations the platform actually returns ("1.2万" -> 12000, "3亿" -> 300000000). """ if value is None or isinstance(value, bool): return None if isinstance(value, int): return value if isinstance(value, float): return int(value) text = str(value).strip().replace(",", "").replace(" ", "") if not text: return None match = _COUNT_RE.match(text) if not match: return None try: number = float(match.group(1)) except ValueError: return None return int(number * _COUNT_UNITS.get(match.group(2), 1)) # Windows reports hard process failures as NTSTATUS values, which surface in the # UI as meaningless large integers (e.g. 3221225794 = 0xC0000142). Translating # the ones we actually see saves the reader a hex-decoding detour. _WINDOWS_EXIT_REASONS = { 0xC0000005: "进程访问冲突 (ACCESS_VIOLATION)", 0xC00000FD: "栈溢出 (STACK_OVERFLOW)", 0xC000013A: "进程被中断(控制台关闭或 Ctrl+C)", 0xC0000142: "进程初始化失败 (STATUS_DLL_INIT_FAILED),属启动环境异常,重启服务后重试", 0xC0000409: "栈缓冲区溢出 (STACK_BUFFER_OVERRUN)", } def describe_exit_code(code: int) -> str: """Render an exit code so a human can act on it.""" unsigned = code & 0xFFFFFFFF if code < 0 else code reason = _WINDOWS_EXIT_REASONS.get(unsigned) if reason: return f"Crawler exited with code {code} (0x{unsigned:08X}): {reason}" return f"Crawler exited with code {code}" @dataclass class IngestResult: status: str notes_fetched: int = 0 comments_fetched: int = 0 new_notes: int = 0 new_comments: int = 0 is_baseline: bool = False error: Optional[str] = None events: List[str] = field(default_factory=list) def _read_jsonl(path: Path) -> List[Dict[str, Any]]: """Read a jsonl file, skipping blank or malformed lines.""" records: List[Dict[str, Any]] = [] if not path.exists(): return records with path.open("r", encoding="utf-8") as handle: for line in handle: line = line.strip() if not line: continue try: item = json.loads(line) except json.JSONDecodeError: continue if isinstance(item, dict): records.append(item) return records def find_run_files( out_dir: Path, platform: str = PLATFORM_XHS ) -> tuple[List[Path], List[Path]]: """Locate the contents/comments jsonl files a run produced. The crawler writes ``{save_data_path}/{platform}/jsonl/{type}_{item}_{date}.jsonl``. Glob rather than reconstructing the name: both the crawler type and the date are runtime-dependent. Returns lists because a crawl crossing midnight produces one file per day. """ jsonl_dir = out_dir / platform / "jsonl" if not jsonl_dir.is_dir(): return [], [] return ( sorted(jsonl_dir.glob("*_contents_*.jsonl")), sorted(jsonl_dir.glob("*_comments_*.jsonl")), ) async def _emit( session: AsyncSession, run: MonitorRun, event_type: str, title: str, *, severity: str = "info", target_kind: str = "", target_id: str = "", payload: Optional[Dict[str, Any]] = None, ) -> None: session.add( MonitorEvent( task_id=run.task_id, run_id=run.id, type=event_type, severity=severity, target_kind=target_kind, target_id=target_id, title=title, payload_json=json.dumps(payload or {}, ensure_ascii=False), created_at=get_current_timestamp(), ) ) async def _previous_run_started_at( session: AsyncSession, task_id: int, run_id: int ) -> Optional[int]: """Started-at of the most recent earlier successful run, in ms.""" return await session.scalar( select(MonitorRun.started_at) .where( MonitorRun.task_id == task_id, MonitorRun.id != run_id, MonitorRun.status.in_((RUN_SUCCESS, RUN_PARTIAL)), MonitorRun.started_at.is_not(None), # Same reasoning as _count_prior_successes: an empty run is a useless # reference point for "was this comment posted since last time?". MonitorRun.notes_fetched > 0, ) .order_by(MonitorRun.id.desc()) .limit(1) ) # How far back to look for proof that the stored login still works. _AUTH_PROOF_WINDOW_MS = 6 * 60 * 60 * 1000 async def _another_task_succeeded_recently(session: AsyncSession, task_id: int) -> bool: """Whether a different task fetched data recently, proving the login is valid.""" since = get_current_timestamp() - _AUTH_PROOF_WINDOW_MS count = await session.scalar( select(func.count()) .select_from(MonitorRun) .where( MonitorRun.task_id != task_id, MonitorRun.status == RUN_SUCCESS, MonitorRun.started_at.is_not(None), MonitorRun.started_at >= since, ) ) return bool(count) async def _count_prior_successes(session: AsyncSession, task_id: int, run_id: int) -> int: return ( await session.scalar( select(func.count()) .select_from(MonitorRun) .where( MonitorRun.task_id == task_id, MonitorRun.id != run_id, MonitorRun.status.in_((RUN_SUCCESS, RUN_PARTIAL)), # A run that fetched nothing established no baseline. Without this # check the first run that actually works after a failed one looks # like a flood of newly discovered works. MonitorRun.notes_fetched > 0, ) ) ) or 0 async def _ingest_notes( session: AsyncSession, run: MonitorRun, records: List[Dict[str, Any]], is_baseline: bool, ) -> int: """Upsert notes, write metric snapshots, and emit new-note/delta events.""" now = get_current_timestamp() new_count = 0 for record in records: note_id = record.get("note_id") if not note_id: continue note = await session.scalar( select(MonitorNote).where( MonitorNote.task_id == run.task_id, MonitorNote.note_id == note_id, ) ) title = (record.get("title") or "")[:500] raw_images = record.get("image_list") or "" cover = raw_images.split(",")[0] if raw_images else "" if note is None: note = MonitorNote( task_id=run.task_id, note_id=note_id, title=title, note_url=record.get("note_url") or "", cover=cover, creator_hash=record.get("creator_hash") or "", creator_name=record.get("nickname") or "", source_kind=record.get("type") or "", published_at=_as_int(record.get("time")), first_seen_run_id=run.id, first_seen_at=now, last_seen_run_id=run.id, last_seen_at=now, ) session.add(note) new_count += 1 if not is_baseline: await _emit( session, run, EVENT_NEW_NOTE, f"新作品:{title or note_id}", target_kind="note", target_id=note_id, payload={"note_id": note_id, "title": title}, ) else: # Only refresh descriptive fields; seen-tracking is updated below. if title: note.title = title # 封面地址**带签名、会过期**,所以每轮都用最新的覆盖它。原先只在首次入库 # 时写一次,结果旧作品的封面地址烂在库里 —— 隔天开始全是 403,而且再怎么 # 重跑也修不回来。落盘那份由 covers.cache_pending 负责(网络操作不在本模块)。 if cover: note.cover = cover # 昵称也要刷:作者改昵称是常事,只在首次入库写一次会一直显示旧的。 if record.get("nickname"): note.creator_name = record["nickname"] note.last_seen_run_id = run.id note.last_seen_at = now await _snapshot_metrics(session, run, note_id, record, now, is_baseline) return new_count def _as_int(value: Any) -> Optional[int]: try: return int(value) except (TypeError, ValueError): return None async def _snapshot_metrics( session: AsyncSession, run: MonitorRun, note_id: str, record: Dict[str, Any], now: int, is_baseline: bool, ) -> None: """Write this run's metric snapshot and report any change vs the previous one.""" previous = await session.scalar( select(MonitorNoteMetric) .where( MonitorNoteMetric.task_id == run.task_id, MonitorNoteMetric.note_id == note_id, MonitorNoteMetric.run_id != run.id, ) .order_by(MonitorNoteMetric.run_id.desc()) .limit(1) ) parsed = {name: parse_count(record.get(name)) for name in _METRIC_FIELDS} session.add( MonitorNoteMetric( task_id=run.task_id, note_id=note_id, run_id=run.id, captured_at=now, liked_count=parsed["liked_count"], comment_count=parsed["comment_count"], collected_count=parsed["collected_count"], share_count=parsed["share_count"], raw_liked_count=str(record.get("liked_count") or ""), raw_comment_count=str(record.get("comment_count") or ""), raw_collected_count=str(record.get("collected_count") or ""), raw_share_count=str(record.get("share_count") or ""), ) ) if previous is None or is_baseline: return deltas = {} for name in _METRIC_FIELDS: old, new = getattr(previous, name), parsed[name] # A None on either side means the value was unparseable; skip rather # than report a bogus change. if old is None or new is None or old == new: continue deltas[name] = {"from": old, "to": new, "delta": new - old} if deltas: summary = "、".join( f"{_metric_label(name)} {info['from']}→{info['to']}" for name, info in deltas.items() ) await _emit( session, run, EVENT_METRIC_DELTA, f"互动数据变化:{summary}", target_kind="note", target_id=note_id, payload={"note_id": note_id, "deltas": deltas}, ) def _metric_label(name: str) -> str: return { "liked_count": "点赞", "comment_count": "评论", "collected_count": "收藏", "share_count": "分享", }.get(name, name) async def _ingest_comments( session: AsyncSession, run: MonitorRun, records: List[Dict[str, Any]], is_baseline: bool, previous_run_started_at: Optional[int], ) -> int: """Upsert comments and emit events for ones never seen before.""" now = get_current_timestamp() new_count = 0 for record in records: comment_id = record.get("comment_id") note_id = record.get("note_id") if not comment_id or not note_id: continue exists = await session.scalar( select(MonitorComment.id).where( MonitorComment.task_id == run.task_id, MonitorComment.note_id == note_id, MonitorComment.comment_id == comment_id, ) ) if exists is not None: continue create_time = _as_int(record.get("create_time")) session.add( MonitorComment( task_id=run.task_id, note_id=note_id, comment_id=comment_id, content=(record.get("content") or "")[:2000], nickname=record.get("nickname") or "", creator_hash=record.get("creator_hash") or "", create_time=create_time, like_count=parse_count(record.get("like_count")), sub_comment_count=_as_int(record.get("sub_comment_count")) or 0, parent_comment_id=record.get("parent_comment_id") or "", first_seen_run_id=run.id, first_seen_at=now, ) ) new_count += 1 if is_baseline: continue # Without a time-sorted comment API we can only observe the top-N window, # so distinguish a genuinely new comment from one that just surfaced. posted = ( create_time is not None and previous_run_started_at is not None and create_time > previous_run_started_at ) await _emit( session, run, EVENT_NEW_COMMENT_POSTED if posted else EVENT_NEW_COMMENT_SEEN, f"{'新评论' if posted else '新出现评论'}:{(record.get('content') or '')[:60]}", target_kind="note", target_id=note_id, payload={ "note_id": note_id, "comment_id": comment_id, "create_time": create_time, "nickname": record.get("nickname") or "", }, ) return new_count async def ingest_run( session: AsyncSession, run: MonitorRun, task: MonitorTask, out_dir: Path, ) -> IngestResult: """Ingest one finished run and return what changed. Sets ``run.status``, ``run.is_baseline`` and the counters on the run row. On a failed or untrustworthy run nothing is diffed -- the "seen" sets only ever grow, so a partial run must never be allowed to look like deletions. """ # A non-zero exit is a genuine crash: trust nothing this run produced. if run.exit_code not in (0, None): run.status = RUN_FAILED run.error_message = describe_exit_code(run.exit_code) await _emit( session, run, EVENT_RUN_FAILED, f"采集进程异常退出(code={run.exit_code})", severity="error", payload={"exit_code": run.exit_code, "detail": run.error_message}, ) return IngestResult(status=RUN_FAILED, error=run.error_message) contents_paths, comment_paths = find_run_files(out_dir, task.platform) contents = [record for path in contents_paths for record in _read_jsonl(path)] comments = [record for path in comment_paths for record in _read_jsonl(path)] run.notes_fetched = len(contents) run.comments_fetched = len(comments) # A bad cookie does NOT fail the process: XHS cookie login is never validated, # so an unauthenticated session just returns zero notes with exit 0 -- and # usually does not even create an output file. Treating that as "the creator # posted nothing" would silently hide login outages, which is exactly what # monitoring exists to catch. if not contents: run.status = RUN_PARTIAL # Blaming the cookie is only honest if nothing else is authenticating. # A sibling task that just succeeded proves the login works, so the # fault is with this target (bad/expired per-creator token, an empty # account, or a page-structure change). if await _another_task_succeeded_recently(session, run.task_id): run.error_message = ( "Crawler produced no notes for this target, but other tasks " "succeeded recently, so the login is probably fine" ) await _emit( session, run, EVENT_NO_DATA, "本次未抓到任何作品:其他任务近期采集正常,登录态应该没问题,请检查该目标是否有效", severity="warning", payload={"out_dir": str(out_dir)}, ) else: run.error_message = "Crawler produced no notes; the login cookie may have expired" await _emit( session, run, EVENT_AUTH_FAILURE, "疑似登录态失效:本次未抓到任何作品,请检查 Cookie", severity="error", payload={"out_dir": str(out_dir)}, ) return IngestResult( status=RUN_PARTIAL, error=run.error_message, comments_fetched=len(comments), ) is_baseline = await _count_prior_successes(session, run.task_id, run.id) == 0 run.is_baseline = is_baseline run.status = RUN_SUCCESS run.error_message = None previous_started_at = ( None if is_baseline else await _previous_run_started_at(session, run.task_id, run.id) ) result = IngestResult( status=RUN_SUCCESS, notes_fetched=len(contents), comments_fetched=len(comments), is_baseline=is_baseline, ) result.new_notes = await _ingest_notes(session, run, contents, is_baseline) if task.enable_comments: result.new_comments = await _ingest_comments( session, run, comments, is_baseline, previous_started_at ) run.new_notes = result.new_notes run.new_comments = result.new_comments return result