Files
butubb f6ddc46d62
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(notify): 通知拆成「新作品」与「异常」两个开关,异常默认开
问题:cookie 过期导致任务失败,但没有任何通知。查下来不是代码问题 ——
notify_enabled 在两个任务上都是 False,而它默认就是关的,事件(run_failed)
也确实生成了,只卡在最后一道闸门。

但那个默认值是错的。代码里的理由是「一条任务列表都推到一个群会很快变吵,所以默认静默」,
这个理由对新作品成立(可能每轮都有),对失败不成立:一次登录态失效意味着这个任务事实上
已经死了,而你不会知道,直到某天发现数据停在几周前。最该被告知的就是这种情况。

现在拆开:
- notify_enabled  —— 推送新作品,可能每轮都有,默认关
- notify_failures —— 推送异常(登录失效/运行失败/没抓到数据),默认开

事件按开关过滤(build_run_message):只勾了「新作品」的任务不该因为一次失败被推消息,
反之亦然,否则拆开开关就没有意义。已有任务由 _ensure_columns 补上 notify_failures=1,
所以会自动开始收到异常推送。

列名 notify_enabled 是历史遗留(它早先是唯一的通知开关),语义已收窄为「新作品」,
用注释写明,不做列重命名 —— 那需要单独的迁移,不值为一个内部工具做。
2026-10-09 13:37:02 +08:00

394 lines
19 KiB
Python

# -*- coding: utf-8 -*-
# Copyright (c) 2025 [email protected]
#
# This file is part of MediaCrawler project.
# Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/monitor/models.py
# GitHub: https://github.com/NanmiCoder
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
#
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
# 1. 不得用于任何商业用途。
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
# 3. 不得进行大规模爬取或对平台造成运营干扰。
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
# 5. 不得用于任何非法或不当的用途。
#
# 详细许可条款请参阅项目根目录下的LICENSE文件。
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
"""Monitoring layer data model.
Lives in its own SQLite database (``data/monitor.db``) with its own declarative
Base, deliberately separate from the crawler's ``database/models.py``. The
crawler's DB store overwrites ``liked_count`` and friends in place on every
re-crawl, so it cannot answer "how did this note change?". These tables keep the
history the crawler throws away.
All timestamps are epoch **milliseconds** (BigInteger), matching the project's
own ``tools.time_util.get_current_timestamp()`` convention. Using ints
throughout avoids naive/aware datetime mixing bugs.
"""
from typing import Optional
from sqlalchemy import (
BigInteger,
Boolean,
ForeignKey,
Integer,
String,
Text,
UniqueConstraint,
)
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship
class MonitorBase(DeclarativeBase):
"""Declarative base for the monitoring database."""
# Run statuses
RUN_PENDING = "pending"
RUN_RUNNING = "running"
RUN_SUCCESS = "success"
RUN_PARTIAL = "partial"
RUN_FAILED = "failed"
RUN_TIMEOUT = "timeout"
RUN_INTERRUPTED = "interrupted"
# Event types
EVENT_NEW_NOTE = "new_note"
EVENT_NEW_COMMENT_POSTED = "new_comment_posted"
EVENT_NEW_COMMENT_SEEN = "new_comment_seen"
EVENT_METRIC_DELTA = "metric_delta"
EVENT_RUN_FAILED = "run_failed"
EVENT_AUTH_FAILURE = "suspected_auth_failure"
# A run that completed cleanly yet fetched nothing, where the login is provably
# fine because another task just succeeded with it. The target, not the cookie,
# is what needs looking at.
EVENT_NO_DATA = "no_data_found"
# Task modes. One subprocess handles exactly one crawler type, so a task is
# either creator-driven or note-driven -- never both.
MODE_CREATOR = "creator"
MODE_NOTE = "note"
class MonitorTask(MonitorBase):
"""One monitored schedule: a set of targets, plus when to run them."""
__tablename__ = "monitor_task"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
name: Mapped[str] = mapped_column(String(200), nullable=False)
platform: Mapped[str] = mapped_column(String(32), nullable=False, default="xhs")
mode: Mapped[str] = mapped_column(String(16), nullable=False)
enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
interval_minutes: Mapped[int] = mapped_column(Integer, nullable=False, default=360)
# How the task is scheduled. `interval` is the original "every N minutes" and
# stays the default; `daily` and `weekly` fire at chosen clock times instead
# (the arithmetic lives in schedule.py).
#
# The clock fields are comma-separated text rather than a child table: they
# are a handful of small integers, always read as a whole, and a table would
# buy nothing but joins.
schedule_mode: Mapped[str] = mapped_column(String(16), nullable=False, default="interval")
# 0-23, e.g. "9,12,18". Empty in interval mode.
schedule_hours: Mapped[str] = mapped_column(String(96), nullable=False, default="")
# 0-6 with Monday = 0, matching Python's date.weekday(). Weekly mode only.
schedule_days: Mapped[str] = mapped_column(String(32), nullable=False, default="")
# Minute past the hour, shared by every time in the schedule.
schedule_minute: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
# Crawl window knobs, mirrored onto each run's CLI flags.
max_notes_count: Mapped[int] = mapped_column(Integer, nullable=False, default=20)
enable_comments: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
max_comments_count: Mapped[int] = mapped_column(Integer, nullable=False, default=50)
run_timeout_seconds: Mapped[int] = mapped_column(Integer, nullable=False, default=3600)
# 通知分成两类,因为它们的性质完全不同:
#
# * `notify_enabled` —— **推送新作品**。可能每轮都有,一条任务列表都推到同一个群
# 会很快变吵,所以默认关。(列名是历史遗留:它早先是唯一的通知开关。)
# * `notify_failures` —— **推送异常**(登录失效 / 运行失败 / 没抓到数据)。频率低,
# 而且一旦发生就意味着这个任务从此**默默采不到任何东西**,你会一直不知道,
# 直到某天发现数据停在几周前。这正是最该被告知的情况,所以默认**开**。
notify_enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
notify_failures: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
# Scheduler state. Persisted so the schedule survives an API restart.
next_run_at: Mapped[Optional[int]] = mapped_column(BigInteger, index=True)
last_run_at: Mapped[Optional[int]] = mapped_column(BigInteger)
last_status: Mapped[str] = mapped_column(String(32), nullable=False, default="idle")
last_error: Mapped[Optional[str]] = mapped_column(Text)
# Lets the UI answer "why did I not get a push for this run?".
last_notified_at: Mapped[Optional[int]] = mapped_column(BigInteger)
created_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
updated_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
targets: Mapped[list["MonitorTarget"]] = relationship(
back_populates="task",
cascade="all, delete-orphan",
lazy="selectin",
)
class MonitorTarget(MonitorBase):
"""One watched creator or note belonging to a task.
``external_id`` is the stable identity (XHS user_id / note_id). It is kept
separate from ``xsec_token`` on purpose: tokens expire within weeks, so
treating a tokenised URL as the primary key would make every long-running
task fail eventually.
"""
__tablename__ = "monitor_target"
__table_args__ = (
UniqueConstraint("task_id", "kind", "external_id", name="uq_monitor_target"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
task_id: Mapped[int] = mapped_column(
ForeignKey("monitor_task.id", ondelete="CASCADE"), nullable=False, index=True
)
kind: Mapped[str] = mapped_column(String(16), nullable=False)
external_id: Mapped[str] = mapped_column(String(128), nullable=False)
xsec_token: Mapped[str] = mapped_column(String(512), nullable=False, default="")
xsec_source: Mapped[str] = mapped_column(String(64), nullable=False, default="")
raw_value: Mapped[str] = mapped_column(Text, nullable=False, default="")
label: Mapped[str] = mapped_column(String(200), nullable=False, default="")
enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
created_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
task: Mapped["MonitorTask"] = relationship(back_populates="targets")
class MonitorRun(MonitorBase):
"""One subprocess execution. The run history in the UI is this table."""
__tablename__ = "monitor_run"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
task_id: Mapped[int] = mapped_column(
ForeignKey("monitor_task.id", ondelete="CASCADE"), nullable=False, index=True
)
trigger: Mapped[str] = mapped_column(String(16), nullable=False, default="scheduled")
status: Mapped[str] = mapped_column(String(16), nullable=False, default=RUN_PENDING, index=True)
phase: Mapped[str] = mapped_column(String(16), nullable=False)
# Where this run's jsonl landed. Each run gets its own directory because the
# crawler's file writer names output by date only.
save_data_path: Mapped[str] = mapped_column(Text, nullable=False, default="")
queued_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
not_before: Mapped[int] = mapped_column(BigInteger, nullable=False, default=0)
started_at: Mapped[Optional[int]] = mapped_column(BigInteger)
finished_at: Mapped[Optional[int]] = mapped_column(BigInteger)
# BigInteger, not Integer: Windows reports failures as unsigned 32-bit
# NTSTATUS values (0xC0000142 = 3221225794), which overflow MySQL's signed
# INT. SQLite's dynamic typing hid this until the data was migrated.
exit_code: Mapped[Optional[int]] = mapped_column(BigInteger)
notes_fetched: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
comments_fetched: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
new_notes: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
new_comments: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
# The very first successful run of a task establishes the baseline: every
# note is "new" at that point, so emitting events would be pure noise.
is_baseline: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
# Window actually used, so the UI can be honest that comments are the top N
# in the platform's own ordering rather than a complete set.
max_comments_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
error_message: Mapped[Optional[str]] = mapped_column(Text)
class MonitorNote(MonitorBase):
"""A note ever seen by a task, plus when it was first/last seen.
Grain is (task, note) so the same note tracked by two tasks stays independent.
"""
__tablename__ = "monitor_note"
__table_args__ = (
UniqueConstraint("task_id", "note_id", name="uq_monitor_note"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
task_id: Mapped[int] = mapped_column(
ForeignKey("monitor_task.id", ondelete="CASCADE"), nullable=False, index=True
)
note_id: Mapped[str] = mapped_column(String(128), nullable=False, index=True)
title: Mapped[str] = mapped_column(Text, nullable=False, default="")
note_url: Mapped[str] = mapped_column(Text, nullable=False, default="")
cover: Mapped[str] = mapped_column(Text, nullable=False, default="")
# 创作者匿名哈希。爬虫刻意不落原始 user_id(见 tools/user_hash.py),
# 所以这是唯一稳定的创作者标识 —— 按博主分组就靠它。
creator_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="")
# 创作者昵称,**已由爬虫脱敏**(张***三 这种)。存的是脱敏后的值,与项目一贯的
# 匿名化姿态一致;不存的话分组只能显示一串哈希,根本认不出是谁。
creator_name: Mapped[str] = mapped_column(String(200), nullable=False, default="")
source_kind: Mapped[str] = mapped_column(String(16), nullable=False, default="")
published_at: Mapped[Optional[int]] = mapped_column(BigInteger)
first_seen_run_id: Mapped[Optional[int]] = mapped_column(Integer)
first_seen_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
last_seen_run_id: Mapped[Optional[int]] = mapped_column(Integer)
last_seen_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
class MonitorNoteMetric(MonitorBase):
"""One metric snapshot per (task, note, run) -- the time series.
Raw strings are kept alongside the parsed integers so a mis-parsed "1.2万"
can always be audited after the fact.
"""
__tablename__ = "monitor_note_metric"
__table_args__ = (
UniqueConstraint("task_id", "note_id", "run_id", name="uq_note_metric"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
task_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True)
note_id: Mapped[str] = mapped_column(String(128), nullable=False, index=True)
run_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True)
captured_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
# NULL (not 0) when the platform value could not be parsed: storing 0 would
# forge a large negative delta on the next comparison.
liked_count: Mapped[Optional[int]] = mapped_column(Integer)
comment_count: Mapped[Optional[int]] = mapped_column(Integer)
collected_count: Mapped[Optional[int]] = mapped_column(Integer)
share_count: Mapped[Optional[int]] = mapped_column(Integer)
raw_liked_count: Mapped[str] = mapped_column(String(64), nullable=False, default="")
raw_comment_count: Mapped[str] = mapped_column(String(64), nullable=False, default="")
raw_collected_count: Mapped[str] = mapped_column(String(64), nullable=False, default="")
raw_share_count: Mapped[str] = mapped_column(String(64), nullable=False, default="")
class MonitorComment(MonitorBase):
"""A comment ever seen by a task.
The (task, note, comment) uniqueness gives idempotent dedup across runs for
free -- re-running the same crawl cannot double-count.
"""
__tablename__ = "monitor_comment"
__table_args__ = (
UniqueConstraint("task_id", "note_id", "comment_id", name="uq_monitor_comment"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
task_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True)
note_id: Mapped[str] = mapped_column(String(128), nullable=False, index=True)
comment_id: Mapped[str] = mapped_column(String(128), nullable=False)
content: Mapped[str] = mapped_column(Text, nullable=False, default="")
nickname: Mapped[str] = mapped_column(String(200), nullable=False, default="")
creator_hash: Mapped[str] = mapped_column(String(64), nullable=False, default="")
# Platform-stated publish time. Used to distinguish a genuinely new comment
# from one that merely entered the visible top-N window this run.
create_time: Mapped[Optional[int]] = mapped_column(BigInteger)
like_count: Mapped[Optional[int]] = mapped_column(Integer)
sub_comment_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
parent_comment_id: Mapped[str] = mapped_column(String(128), nullable=False, default="")
first_seen_run_id: Mapped[Optional[int]] = mapped_column(Integer)
first_seen_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
class MonitorEvent(MonitorBase):
"""Append-only change feed. This is what the dashboard reads."""
__tablename__ = "monitor_event"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
task_id: Mapped[int] = mapped_column(Integer, nullable=False, index=True)
run_id: Mapped[Optional[int]] = mapped_column(Integer, index=True)
type: Mapped[str] = mapped_column(String(32), nullable=False, index=True)
severity: Mapped[str] = mapped_column(String(16), nullable=False, default="info")
target_kind: Mapped[str] = mapped_column(String(16), nullable=False, default="")
target_id: Mapped[str] = mapped_column(String(128), nullable=False, default="")
title: Mapped[str] = mapped_column(Text, nullable=False, default="")
payload_json: Mapped[str] = mapped_column(Text, nullable=False, default="{}")
created_at: Mapped[int] = mapped_column(BigInteger, nullable=False, index=True)
is_read: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
class MonitorSetting(MonitorBase):
"""Key/value store. Holds the XHS cookie for unattended runs."""
__tablename__ = "monitor_setting"
key: Mapped[str] = mapped_column(String(64), primary_key=True)
value: Mapped[str] = mapped_column(Text, nullable=False, default="")
updated_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
class AuthSession(MonitorBase):
"""A WebUI login session.
Only the SHA-256 of the token is stored, never the token itself -- a leaked
database therefore does not hand over live sessions. This mirrors the
existing posture of never returning the XHS cookie or webhook value.
A stateful table (rather than a signed stateless token) is what makes "log
out" and "password changed" take effect immediately.
"""
__tablename__ = "auth_session"
token_hash: Mapped[str] = mapped_column(String(64), primary_key=True)
created_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
expires_at: Mapped[int] = mapped_column(BigInteger, nullable=False, index=True)
last_seen_at: Mapped[int] = mapped_column(BigInteger, nullable=False)
SETTING_AUTH_PASSWORD_HASH = "auth_password_hash"
SETTING_AUTH_PASSWORD_UPDATED_AT = "auth_password_updated_at"
# Settings are namespaced by scope: `platform.<p>.<name>` for values each
# platform keeps its own copy of, `system.<name>` for values shared across all of
# them. Key builders live in settings.py.
SETTING_WECOM_WEBHOOK = "system.wecom_webhook"
# Pre-namespacing keys, kept only so the startup migration can find and move
# them. Nothing should read these directly.
LEGACY_SETTING_KEY_RENAMES = {
# Pre-batch-2 flat keys.
"xhs_cookie": "platform.xhs.cookie",
"xhs_cookie_updated_at": "platform.xhs.cookie_updated_at",
"xhs_cookie_last_ok_at": "platform.xhs.cookie_last_ok_at",
"wecom_webhook": "system.wecom_webhook",
# Batch-2 keys, before settings gained a scope. Those values belonged to
# Xiaohongshu because it was the only platform, so they migrate to its scope;
# the two scheduling keys were always instance-wide.
"collect.default_interval_minutes": "platform.xhs.default_interval_minutes",
"collect.default_max_notes": "platform.xhs.default_max_notes",
"collect.default_max_comments": "platform.xhs.default_max_comments",
"collect.enable_sub_comments": "platform.xhs.enable_sub_comments",
"collect.crawl_sleep_sec": "platform.xhs.crawl_sleep_sec",
"collect.active_hours_start": "system.active_hours_start",
"collect.active_hours_end": "system.active_hours_end",
"proxy.enable_ip_proxy": "platform.xhs.enable_ip_proxy",
"proxy.provider": "platform.xhs.proxy_provider",
"proxy.pool_count": "platform.xhs.proxy_pool_count",
"proxy.static_proxy_url": "platform.xhs.static_proxy_url",
}
# utf8mb4 is forced on every table rather than left to the schema default: this
# deployment's MySQL server *and* the target database both default to latin1,
# which would mangle or reject Chinese text. Setting it per table means it holds
# regardless of what the schema default happens to be.
#
# Must run after every model is declared, hence the end of the module.
for _table in MonitorBase.metadata.tables.values():
_table.kwargs["mysql_charset"] = "utf8mb4"
_table.kwargs["mysql_collate"] = "utf8mb4_unicode_ci"