Files
MediaCrawler/api/monitor/service.py
T
butubb d52b6942ac
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
fix(monitor): 作品栏只显示每个博主最新的 N 条(库里一条不删)
原来 max_notes_count **只**被当作每轮拉取的条数,落库时从不删任何东西,列表也不按
博主截断。于是超出的作品会一直堆着(21、22、23…),而掉出采集窗口的那些**再也采不到**
—— 指标冻在最后一次,看上去却和正在跟踪的作品一模一样。界面上那句「仅显示每个博主
最新的 N 条作品」是假的,后端根本没做这个截断。

现在界面按窗口显示:库是账本,界面是窗口。趋势图、报表、**导出**读的仍然是全量 ——
导出是数据,少几行就是在丢东西,所以 `windowed` 默认关,只有 /monitor/notes 打开。

窗口按**发布时间**排序,不是「最后一次见到」:降级路径(作品列表被风控挡住)会把所有
已知作品都刷一遍,那样每个人的 last_seen_at 都一样,按它开窗等于没开。解析不出日期的
作品**永不隐藏** —— 因为读不出日期就让一条作品消失,比多显示一条糟得多。

只对博主模式开窗;笔记模式的目标本身就是一件作品。
2026-10-11 08:30:10 +08:00

1160 lines
44 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- 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/service.py
# GitHub: https://github.com/NanmiCoder
# Non-commercial learning license 1.1
#
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
# 1. 不得用于任何商业用途。
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
# 3. 不得进行大规模爬取或对平台造成运营干扰。
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
# 5. 不得用于任何非法或不当的用途。
#
# 详细许可条款请参阅项目根目录下的LICENSE文件。
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
"""Task CRUD and dashboard queries for the monitoring layer."""
import asyncio
from typing import Any, Dict, List, Optional, Sequence
from urllib.parse import parse_qs, urlparse
from sqlalchemy import delete, func, select
from sqlalchemy.ext.asyncio import AsyncSession
from tools.time_util import get_current_timestamp
from . import adapters, app_settings, covers, platforms, schedule
from .db import get_session
from .platforms import PLATFORM_XHS
from .models import (
MODE_CREATOR,
MODE_NOTE,
MonitorComment,
MonitorCreatorAlias,
MonitorCreatorStat,
MonitorEvent,
MonitorNote,
MonitorNoteAlias,
MonitorNoteMetric,
MonitorRun,
MonitorTarget,
MonitorTask,
RUN_SUCCESS,
RUN_PARTIAL,
)
from .runner import execute_task
# Keep strong references to in-flight manual runs; asyncio only holds weak ones,
# so without this a run can be garbage collected mid-flight.
_background_runs: set[asyncio.Task] = set()
MIN_INTERVAL_MINUTES = 30
MAX_INTERVAL_MINUTES = 7 * 24 * 60
# Changing any of these invalidates the pending run slot.
SCHEDULE_FIELDS = {
"interval_minutes",
"schedule_mode",
"schedule_hours",
"schedule_days",
"schedule_minute",
}
def _require_clock_fields(mode: str, hours: list, days: list) -> None:
"""A clock schedule with no clock time can never fire.
The create schema already rejects that shape, but the update path merges
partial fields and therefore has no schema-level view of the result -- so the
check lives here, where both paths meet.
"""
if mode in schedule.CLOCK_MODES and not hours:
raise ValueError("按钟点调度至少要选一个时间")
if mode == schedule.MODE_WEEKLY and not days:
raise ValueError("按周调度至少要选一个星期")
# 各平台的链接形态、id 形状、短链域名都在 adapters.py —— 那里是「平台之间不一样」
# 的东西的唯一出处,所以这里不再留任何平台字面量。
class TargetParseError(ValueError):
"""Raised when a pasted monitoring target cannot be understood."""
def parse_target_input(
value: str, mode: str, platform: str = PLATFORM_XHS
) -> Dict[str, str]:
"""Parse a pasted creator/note value into a stable id plus a refreshable token.
Accepts either a full URL (with or without ``xsec_token``) or a bare id.
Storing the id separately from the token is what keeps a long-running task
alive: tokens expire, ids do not.
URL shapes are platform-specific and come from ``adapters``. A platform with
no adapter is rejected here as well as at task creation -- parsing a Douyin
link as if it were a Xiaohongshu one would be worse than refusing it.
"""
if not adapters.has_adapter(platform):
raise TargetParseError(f"暂不支持解析该平台({platform})的目标链接")
spec = adapters.adapter(platform)
raw = (value or "").strip()
if not raw:
raise TargetParseError("Empty target")
creator_mode = mode == MODE_CREATOR
expected = "博主主页" if creator_mode else "笔记"
external_id = ""
if raw.startswith("http") or "/" in raw:
if any(host in raw for host in spec.short_link_hosts):
# 短链要联网跳一次才知道指向谁,而这里没有网络可跳。明确拒绝好过存一个
# 解析不出 id 的值 —— 那会变成一个永远抓不到东西、还不报错的任务。
raise TargetParseError(f"{expected}短链无法解析,请粘贴完整链接:{raw}")
patterns = spec.creator_url_res if creator_mode else spec.note_url_res
for pattern in patterns:
match = pattern.search(raw)
if match:
external_id = match.group(1)
break
if not external_id:
raise TargetParseError(f"无法从链接中解析出{expected} ID:{raw}")
elif (spec.creator_bare_re if creator_mode else spec.note_bare_re).match(raw):
external_id = raw
else:
raise TargetParseError(f"无法识别的目标:{raw}")
params = parse_qs(urlparse(raw).query) if raw.startswith("http") else {}
return {
"external_id": external_id,
"xsec_token": (params.get("xsec_token") or [""])[0],
"xsec_source": (params.get("xsec_source") or [""])[0],
"raw_value": raw,
}
async def platform_task_ids(session: AsyncSession, platform: str) -> List[int]:
"""Ids of the tasks belonging to a platform.
Note/comment/event tables carry no platform column -- they hang off a task --
so scoping a query to a platform means scoping it to that task set.
"""
return list(
await session.scalars(select(MonitorTask.id).where(MonitorTask.platform == platform))
)
# ---------------------------------------------------------------------------
# Task CRUD
# ---------------------------------------------------------------------------
async def create_task(session: AsyncSession, payload: Dict[str, Any]) -> MonitorTask:
mode = payload["mode"]
if mode not in (MODE_CREATOR, MODE_NOTE):
raise ValueError(f"Unsupported mode: {mode}")
# Rejects unknown platforms and, more importantly, platforms whose crawler
# exists upstream but whose monitoring is not wired up -- accepting those
# would create a task that can never produce data.
platform = payload.get("platform") or PLATFORM_XHS
platforms.ensure_runnable(platform)
now = get_current_timestamp()
# Fall back to the configured defaults for anything the caller left out, so
# the Settings page actually governs new tasks.
defaults = await app_settings.defaults(session, platform)
interval_minutes = payload.get("interval_minutes") or defaults["interval_minutes"]
schedule_mode = payload.get("schedule_mode") or schedule.MODE_INTERVAL
schedule_hours = list(payload.get("schedule_hours") or [])
schedule_days = list(payload.get("schedule_days") or [])
schedule_minute = int(payload.get("schedule_minute") or 0)
_require_clock_fields(schedule_mode, schedule_hours, schedule_days)
task = MonitorTask(
name=payload["name"],
platform=platform,
mode=mode,
enabled=payload.get("enabled", True),
interval_minutes=interval_minutes,
schedule_mode=schedule_mode,
schedule_hours=schedule.format_hours(schedule_hours),
schedule_days=schedule.format_days(schedule_days),
schedule_minute=schedule_minute,
max_notes_count=payload.get("max_notes_count") or defaults["max_notes_count"],
enable_comments=payload.get("enable_comments", True),
max_comments_count=payload.get("max_comments_count") or defaults["max_comments_count"],
run_timeout_seconds=payload.get("run_timeout_seconds", 3600),
notify_enabled=payload.get("notify_enabled", False),
notify_failures=payload.get("notify_failures", True),
next_run_at=schedule.next_occurrence(
mode=schedule_mode,
interval_minutes=interval_minutes,
hours=schedule_hours,
days=schedule_days,
minute=schedule_minute,
after_ms=now,
),
last_status="idle",
created_at=now,
updated_at=now,
)
session.add(task)
await session.flush()
seen: set[str] = set()
for value in payload.get("targets", []):
parsed = parse_target_input(value, mode, platform)
if parsed["external_id"] in seen:
continue
seen.add(parsed["external_id"])
session.add(
MonitorTarget(
task_id=task.id,
kind=mode,
external_id=parsed["external_id"],
xsec_token=parsed["xsec_token"],
xsec_source=parsed["xsec_source"],
raw_value=parsed["raw_value"],
label=parsed["external_id"],
enabled=True,
created_at=now,
)
)
await session.flush()
return task
async def update_task(session: AsyncSession, task_id: int, payload: Dict[str, Any]) -> MonitorTask:
task = await session.get(MonitorTask, task_id)
if task is None:
raise ValueError(f"Task {task_id} not found")
was_enabled = task.enabled
for field in (
"name",
"enabled",
"interval_minutes",
"schedule_mode",
"schedule_minute",
"max_notes_count",
"enable_comments",
"max_comments_count",
"run_timeout_seconds",
"notify_enabled",
"notify_failures",
):
if field in payload and payload[field] is not None:
setattr(task, field, payload[field])
# The clock lists are stored as comma-separated text, so they cannot go
# through the generic loop above.
if payload.get("schedule_hours") is not None:
task.schedule_hours = schedule.format_hours(payload["schedule_hours"])
if payload.get("schedule_days") is not None:
task.schedule_days = schedule.format_days(payload["schedule_days"])
# Replacing targets resets the baseline implicitly: a note set that now
# includes new ids will simply report them as new on the next run.
if payload.get("targets") is not None:
await session.execute(delete(MonitorTarget).where(MonitorTarget.task_id == task_id))
now = get_current_timestamp()
seen: set[str] = set()
for value in payload["targets"]:
parsed = parse_target_input(value, task.mode, task.platform)
if parsed["external_id"] in seen:
continue
seen.add(parsed["external_id"])
session.add(
MonitorTarget(
task_id=task_id,
kind=task.mode,
external_id=parsed["external_id"],
xsec_token=parsed["xsec_token"],
xsec_source=parsed["xsec_source"],
raw_value=parsed["raw_value"],
label=parsed["external_id"],
enabled=True,
created_at=now,
)
)
# Any change to when the task runs invalidates the pending slot, so recompute
# it from the merged state rather than working out which field moved.
# Re-enabling counts as a change too: otherwise a task switched off for a
# month comes back holding a next_run_at a month in the past and fires the
# instant it is saved.
if (SCHEDULE_FIELDS & set(payload)) or (task.enabled and not was_enabled):
hours = schedule.parse_hours(task.schedule_hours)
days = schedule.parse_days(task.schedule_days)
_require_clock_fields(task.schedule_mode, hours, days)
task.next_run_at = schedule.next_occurrence(
mode=task.schedule_mode,
interval_minutes=task.interval_minutes,
hours=hours,
days=days,
minute=task.schedule_minute,
after_ms=get_current_timestamp(),
)
task.updated_at = get_current_timestamp()
await session.flush()
return task
async def delete_task(session: AsyncSession, task_id: int) -> None:
task = await session.get(MonitorTask, task_id)
if task is None:
raise ValueError(f"Task {task_id} not found")
await session.delete(task)
def trigger_manual_run(task_id: int) -> None:
"""Fire a run in the background and return immediately.
A crawl takes minutes, so the HTTP request must not wait for it. The UI
follows progress through the logs WebSocket and the run history.
"""
task = asyncio.create_task(execute_task(task_id, trigger="manual"))
_background_runs.add(task)
task.add_done_callback(_background_runs.discard)
# ---------------------------------------------------------------------------
# Dashboard queries
# ---------------------------------------------------------------------------
async def _creator_alias_map(session: AsyncSession) -> Dict[tuple, str]:
"""``(platform, creator_hash) -> 备注``。
整体读一次再在内存里取,而不是每条作品查一次 —— 作品列表动辄几十条。备注本身很少,
表不会大。
"""
rows = (await session.scalars(select(MonitorCreatorAlias))).all()
return {(row.platform, row.creator_hash): row.alias for row in rows if row.alias}
async def set_creator_alias(
session: AsyncSession, platform: str, creator_hash: str, alias: str
) -> None:
"""给博主起备注;传空串就是删掉这条备注(界面上的"清空")。"""
alias = (alias or "").strip()[:128]
existing = await session.scalar(
select(MonitorCreatorAlias).where(
MonitorCreatorAlias.platform == platform,
MonitorCreatorAlias.creator_hash == creator_hash,
)
)
if not alias:
if existing is not None:
await session.delete(existing)
return
if existing is None:
session.add(
MonitorCreatorAlias(
platform=platform,
creator_hash=creator_hash,
alias=alias,
updated_at=get_current_timestamp(),
)
)
return
existing.alias = alias
existing.updated_at = get_current_timestamp()
async def _latest_successful_run_id(session: AsyncSession, task_id: int) -> Optional[int]:
return await session.scalar(
select(MonitorRun.id)
.where(
MonitorRun.task_id == task_id,
MonitorRun.status.in_((RUN_SUCCESS, RUN_PARTIAL)),
)
.order_by(MonitorRun.id.desc())
.limit(1)
)
def _delta(current: Optional[int], previous: Optional[int]) -> Optional[int]:
if current is None or previous is None:
return None
return current - previous
async def _note_alias_map(session: AsyncSession) -> Dict[tuple, str]:
"""``(platform, note_id) -> 作品备注``。和博主备注一个道理,整体读一次。"""
rows = (await session.scalars(select(MonitorNoteAlias))).all()
return {(row.platform, row.note_id): row.alias for row in rows if row.alias}
async def set_note_alias(
session: AsyncSession, platform: str, note_id: str, alias: str
) -> None:
"""给作品起备注;空串就是删掉这条备注。"""
alias = (alias or "").strip()[:128]
existing = await session.scalar(
select(MonitorNoteAlias).where(
MonitorNoteAlias.platform == platform,
MonitorNoteAlias.note_id == note_id,
)
)
if not alias:
if existing is not None:
await session.delete(existing)
return
if existing is None:
session.add(
MonitorNoteAlias(
platform=platform,
note_id=note_id,
alias=alias,
updated_at=get_current_timestamp(),
)
)
return
existing.alias = alias
existing.updated_at = get_current_timestamp()
async def _latest_creator_stats(
session: AsyncSession,
task_ids: Optional[Sequence[int]] = None,
creator_hashes: Optional[Sequence[str]] = None,
) -> Dict[tuple, "MonitorCreatorStat"]:
"""``(task_id, creator_hash) -> 最近一条``账号级快照。
两个参数都是**可选过滤**:传 None 就是不限。``list_notes`` 两个都给(只要手里
这批作品涉及的博主),``list_creators`` 只给任务(要的是全部博主,包括一条作品
都没有的那些)。
按 run_id 而不是 captured_at 取「最近」:和作品指标用的是同一个口径,两者放一起
看才不会出现「作品数据来自第 8 轮、粉丝数来自第 9 轮」这种对不上的情况。
"""
if task_ids is not None and not task_ids:
return {}
if creator_hashes is not None and not creator_hashes:
return {}
query = select(MonitorCreatorStat).order_by(MonitorCreatorStat.run_id.desc())
if task_ids is not None:
query = query.where(MonitorCreatorStat.task_id.in_(list(task_ids)))
if creator_hashes is not None:
query = query.where(MonitorCreatorStat.creator_hash.in_(list(creator_hashes)))
latest: Dict[tuple, MonitorCreatorStat] = {}
for row in (await session.scalars(query)).all():
latest.setdefault((row.task_id, row.creator_hash), row) # 已按 run_id 倒序
return latest
def _window_order(note: MonitorNote) -> tuple:
"""开窗排序:按**发布时间**倒序。
采集端 ``aweme/post`` 拿的就是最新 N 条,界面的窗口必须和它口径一致 —— 否则会出现
「显示着的这条其实已经采不到了」。
**时间不知道的排最前(永不隐藏)**:因为解析不出日期就让一条作品从列表里消失,比
多显示一条糟得多。
"""
return (note.published_at is None, note.published_at or 0)
async def _within_creator_window(
session: AsyncSession, notes: List[MonitorNote]
) -> List[MonitorNote]:
"""每个博主只留**最新 N 条**,N 取任务上的 ``max_notes_count``。
**库里一条都不删。** 作品栏是「我正在盯的窗口」,库是账本:掉出窗口的作品采集端就
再也取不到了(它取的是最新 N 条),指标会冻在最后一次而**它看上去和正在跟踪的
作品一模一样** —— 那才是真正会骗人的地方。趋势图、报表、导出读的仍然是全量。
只对**博主模式**开窗:笔记模式的 ``max_notes_count`` 管的是「每条目标拉多少」,
而它的目标本身就是一件作品,不存在「一个博主的第 21 条」。
"""
task_ids = {note.task_id for note in notes}
tasks = {
row.id: row
for row in (
await session.scalars(
select(MonitorTask).where(MonitorTask.id.in_(list(task_ids)))
)
).all()
}
by_creator: Dict[tuple, List[MonitorNote]] = {}
for note in notes:
by_creator.setdefault((note.task_id, note.creator_hash), []).append(note)
kept: List[MonitorNote] = []
for (note_task_id, _creator_hash), group in by_creator.items():
task = tasks.get(note_task_id)
if task is None or task.mode != MODE_CREATOR:
kept.extend(group)
continue
group.sort(key=_window_order, reverse=True)
kept.extend(group[: task.max_notes_count])
return kept
async def list_notes(
session: AsyncSession,
task_id: Optional[int] = None,
only_new: bool = False,
limit: int = 200,
platform: Optional[str] = None,
*,
windowed: bool = False,
) -> List[Dict[str, Any]]:
"""Tracked notes with their latest metrics and change vs the previous run.
``windowed`` 打开后只返回每个博主最新 ``max_notes_count`` 条(见
``_within_creator_window``)。**默认关闭**:导出要的是全量账本,而默认开启的话,
忘了传这个参数的那条路会静默少数据。
"""
query = select(MonitorNote).order_by(MonitorNote.last_seen_at.desc())
# 开窗时**不能**在这里 limit:得先把每个博主的窗开出来,顺序反了的话,全局上限会把
# 某个博主最新的那几条挤掉,剩下的还都是别人的。
if not windowed:
query = query.limit(limit)
if task_id is not None:
query = query.where(MonitorNote.task_id == task_id)
if platform is not None:
scoped = await platform_task_ids(session, platform)
if not scoped:
return []
query = query.where(MonitorNote.task_id.in_(scoped))
notes = list((await session.scalars(query)).all())
if not notes:
return []
if windowed:
notes = await _within_creator_window(session, notes)
# 开完窗再截全局上限。窗口内的作品本来就不多,这里的 200 只是兜底。
notes = notes[:limit]
if not notes:
return []
note_ids = [note.note_id for note in notes]
# Fetch every snapshot for these notes in one go and gather the two most
# recent per note, rather than issuing two queries per note.
snapshots = list(
(
await session.scalars(
select(MonitorNoteMetric)
.where(MonitorNoteMetric.note_id.in_(note_ids))
.order_by(MonitorNoteMetric.note_id, MonitorNoteMetric.run_id.desc())
)
).all()
)
by_note: Dict[str, List[MonitorNoteMetric]] = {}
for snapshot in snapshots:
by_note.setdefault(snapshot.note_id, []).append(snapshot)
latest_run_ids: Dict[int, Optional[int]] = {}
result: List[Dict[str, Any]] = []
# 博主的备注。键是 (platform, creator_hash) —— 作品的 platform 挂在它的任务上。
task_platform = {
row.id: row.platform
for row in (
await session.execute(
select(MonitorTask.id, MonitorTask.platform).where(
MonitorTask.id.in_({note.task_id for note in notes})
)
)
).all()
}
aliases = await _creator_alias_map(session)
note_aliases = await _note_alias_map(session)
# 账号级指标(粉丝 / 总获赞 / 作品数)。**挂在作品上一起返回**,因为界面上就是按博主
# 归组显示的 —— 让前端为了一个组头再发一轮请求没道理。同一个博主的所有作品拿到的是
# 同一条(键里带 task_id,所以跨任务不会串)。没有的(小红书那条路不产生它)就是 null,
# 前端据此整块不显示,而不是显示一个 0。
creator_stats = await _latest_creator_stats(
session,
[note.task_id for note in notes],
[note.creator_hash for note in notes if note.creator_hash],
)
for note in notes:
series = by_note.get(note.note_id, [])
current = series[0] if series else None
previous = series[1] if len(series) > 1 else None
if only_new:
if note.task_id not in latest_run_ids:
latest_run_ids[note.task_id] = await _latest_successful_run_id(session, note.task_id)
if note.first_seen_run_id != latest_run_ids[note.task_id]:
continue
stat = creator_stats.get((note.task_id, note.creator_hash))
result.append(
{
"task_id": note.task_id,
"note_id": note.note_id,
"title": note.title,
"note_url": note.note_url,
# 作品的发布时间(爬虫侧:小红书叫 time、抖音叫 create_time)。
# 和 first_seen_at 不是一回事 —— 那是**我们第一次看到它**的时间;把一个
# 早就存在的作品加进监控时,两者能差好几个月。可能为 null(平台没给,
# 或者值解析不出来),所以前端要能显示成「—」。
"published_at": note.published_at,
# 按博主分组用。creator_hash 是唯一稳定的创作者标识(原始 user_id
# 被爬虫刻意匿名化了),creator_name 是昵称本身 —— 本仓库关掉了脱敏
# (见 config.MASK_NICKNAME),所以就是原文。
"creator_hash": note.creator_hash,
"creator_name": note.creator_name,
# 人自己起的备注,界面上优先显示它 —— 昵称认不出是谁,哈希更认不出。
"creator_alias": aliases.get(
(task_platform.get(note.task_id, ""), note.creator_hash), ""
),
# 这条作品自己的备注。和博主备注是两件事:博主备注回答"这是谁",它回答
# "这条我要盯着"。
"note_alias": note_aliases.get(
(task_platform.get(note.task_id, ""), note.note_id), ""
),
# 博主账号级指标 —— 作品列表给不了的东西。三个值都可能为 null(平台没采
# 到、或者这条作品来自不产生它的数据源),前端据此整块不画。
"creator_fans": stat.fans if stat else None,
"creator_total_favorited": stat.total_favorited if stat else None,
"creator_works": stat.works_count if stat else None,
"creator_stats_at": stat.captured_at if stat else None,
# 优先给本地缓存地址:远程地址带签名、会过期(实测隔天即 403),
# 本地那份不会。没有缓存时才退回远程,至少让图先显示出来。
"cover": covers.cover_url(note.note_id, note.cover),
"first_seen_at": note.first_seen_at,
"last_seen_at": note.last_seen_at,
"is_new": note.first_seen_run_id == latest_run_ids.get(note.task_id),
"metrics": {
"liked_count": current.liked_count if current else None,
"comment_count": current.comment_count if current else None,
"collected_count": current.collected_count if current else None,
"share_count": current.share_count if current else None,
},
"deltas": {
"liked_count": _delta(
current.liked_count if current else None,
previous.liked_count if previous else None,
),
"comment_count": _delta(
current.comment_count if current else None,
previous.comment_count if previous else None,
),
"collected_count": _delta(
current.collected_count if current else None,
previous.collected_count if previous else None,
),
"share_count": _delta(
current.share_count if current else None,
previous.share_count if previous else None,
),
},
"snapshot_count": len(series),
}
)
return result
async def list_creators(
session: AsyncSession,
task_id: Optional[int] = None,
platform: Optional[str] = None,
) -> List[Dict[str, Any]]:
"""作品栏里要显示的**博主** —— **包括一条作品都没有的**。
分组原先是从作品推出来的(按作品的 creator_hash 归组),于是没有作品的博主根本
不会出现在列表里:目标加了、资料也采到了、粉丝数就躺在库里,界面上什么都看不见。
而「这个号在涨粉、只是最近没发作品」恰恰是最该看见的一种情况 —— 藏起来正好藏反了。
所以来源换成 **账号快照 ∪ 作品**:
* 有快照没作品 → 一个 0 篇的组,粉丝数照常显示;
* 有作品没快照 → 一个没有账号指标的组(小红书那条路不产生快照,就是这种情况)。
``creator_alias`` 从作品备注那张表来;``note_count`` / ``last_activity_at`` 用来
排序,让最近还在动的博主排在前面。
"""
scope: Optional[List[int]] = None
if task_id is not None:
scope = [task_id]
elif platform is not None:
scope = await platform_task_ids(session, platform)
if not scope:
return []
# 作品一侧:谁有作品、有几篇、最后一次是什么时候。
work_query = (
select(
MonitorNote.task_id,
MonitorNote.creator_hash,
func.count().label("note_count"),
func.max(MonitorNote.last_seen_at).label("last_seen_at"),
func.max(MonitorNote.creator_name).label("creator_name"),
)
.where(MonitorNote.creator_hash != "")
.group_by(MonitorNote.task_id, MonitorNote.creator_hash)
)
if scope is not None:
work_query = work_query.where(MonitorNote.task_id.in_(scope))
work: Dict[tuple, Dict[str, Any]] = {}
for row in (await session.execute(work_query)).all():
work[(row.task_id, row.creator_hash)] = {
"note_count": row.note_count,
"last_seen_at": row.last_seen_at,
"creator_name": row.creator_name or "",
}
stats = await _latest_creator_stats(session, task_ids=scope)
if not work and not stats:
return []
# 备注是按 (platform, creator_hash) 存的,所以要知道每个博主属于哪个平台。
involved = {key[0] for key in set(work) | set(stats)}
task_platform = {
row.id: row.platform
for row in (
await session.execute(
select(MonitorTask.id, MonitorTask.platform).where(
MonitorTask.id.in_(list(involved))
)
)
).all()
}
aliases = await _creator_alias_map(session)
result: List[Dict[str, Any]] = []
for key in set(work) | set(stats):
row_task, creator_hash = key
work_row = work.get(key)
stat = stats.get(key)
result.append(
{
"task_id": row_task,
"creator_hash": creator_hash,
# 昵称优先取作品的(那是界面上本来就在用的),快照的兜底 —— 没有作品
# 的博主只剩快照这一个来源。
"creator_name": (work_row or {}).get("creator_name")
or (stat.nickname if stat else ""),
"creator_alias": aliases.get(
(task_platform.get(row_task, ""), creator_hash), ""
),
"note_count": (work_row or {}).get("note_count", 0),
"creator_fans": stat.fans if stat else None,
"creator_total_favorited": stat.total_favorited if stat else None,
"creator_works": stat.works_count if stat else None,
"creator_stats_at": stat.captured_at if stat else None,
# 排序用:作品最近出现的时间,或者账号指标的采集时间,取晚的那个。
"last_activity_at": max(
(work_row or {}).get("last_seen_at") or 0,
stat.captured_at if stat else 0,
),
}
)
result.sort(key=lambda row: row["last_activity_at"], reverse=True)
return result
async def note_series(session: AsyncSession, note_id: str, task_id: Optional[int] = None) -> List[Dict[str, Any]]:
"""Metric time series for one note."""
query = (
select(MonitorNoteMetric)
.where(MonitorNoteMetric.note_id == note_id)
.order_by(MonitorNoteMetric.run_id)
)
if task_id is not None:
query = query.where(MonitorNoteMetric.task_id == task_id)
return [
{
"run_id": row.run_id,
"captured_at": row.captured_at,
"liked_count": row.liked_count,
"comment_count": row.comment_count,
"collected_count": row.collected_count,
"share_count": row.share_count,
}
for row in (await session.scalars(query)).all()
]
async def _note_meta_map(
session: AsyncSession, note_ids: List[str]
) -> Dict[str, Dict[str, Any]]:
"""Look up note title/cover/url for a set of note ids.
Fetched as one query and joined in Python rather than as a SQL join: the
comment table has no foreign key to the note table (both are keyed by the
platform's note id, per task), and a single IN() is easier to follow here.
"""
if not note_ids:
return {}
rows = (
await session.scalars(select(MonitorNote).where(MonitorNote.note_id.in_(set(note_ids))))
).all()
return {
row.note_id: {
"note_title": row.title,
"note_cover": covers.cover_url(row.note_id, row.cover),
"note_url": row.note_url,
"published_at": row.published_at,
"task_id": row.task_id,
# 博主维度也带上,评论流才能按 博主 -> 作品 -> 评论 三级展开。
"creator_hash": row.creator_hash,
"creator_name": row.creator_name,
}
for row in rows
}
async def list_comments(
session: AsyncSession,
task_id: Optional[int] = None,
note_id: Optional[str] = None,
limit: int = 200,
platform: Optional[str] = None,
) -> List[Dict[str, Any]]:
"""Comments, each carrying the note it belongs to.
The note association is the point: without it a comment stream is unreadable,
since a bare note_id tells the operator nothing.
"""
query = select(MonitorComment).order_by(MonitorComment.first_seen_at.desc()).limit(limit)
if task_id is not None:
query = query.where(MonitorComment.task_id == task_id)
if note_id is not None:
query = query.where(MonitorComment.note_id == note_id)
if platform is not None:
scoped = await platform_task_ids(session, platform)
if not scoped:
return []
query = query.where(MonitorComment.task_id.in_(scoped))
comments = list((await session.scalars(query)).all())
meta = await _note_meta_map(session, [row.note_id for row in comments])
return [
{
"task_id": row.task_id,
"note_id": row.note_id,
"comment_id": row.comment_id,
"content": row.content,
"nickname": row.nickname,
"create_time": row.create_time,
"like_count": row.like_count,
"sub_comment_count": row.sub_comment_count,
"first_seen_at": row.first_seen_at,
"note_title": meta.get(row.note_id, {}).get("note_title", ""),
"note_cover": meta.get(row.note_id, {}).get("note_cover", ""),
"note_url": meta.get(row.note_id, {}).get("note_url", ""),
"note_creator_hash": meta.get(row.note_id, {}).get("creator_hash", ""),
"note_creator_name": meta.get(row.note_id, {}).get("creator_name", ""),
"note_published_at": meta.get(row.note_id, {}).get("published_at"),
}
for row in comments
]
async def comment_note_groups(
session: AsyncSession, task_id: Optional[int] = None, platform: Optional[str] = None
) -> List[Dict[str, Any]]:
"""Notes that have comments, newest first, with their comment counts.
Feeds the comment filter dropdown: the operator picks a work by title, so
the counts need to be visible before choosing.
"""
scoped_ids: Optional[List[int]] = None
if platform is not None:
scoped_ids = await platform_task_ids(session, platform)
if not scoped_ids:
return []
count_query = select(
MonitorComment.note_id, func.count().label("comment_count")
).group_by(MonitorComment.note_id)
if task_id is not None:
count_query = count_query.where(MonitorComment.task_id == task_id)
if scoped_ids is not None:
count_query = count_query.where(MonitorComment.task_id.in_(scoped_ids))
counts = {row.note_id: row.comment_count for row in (await session.execute(count_query)).all()}
if not counts:
return []
latest_query = (
select(MonitorComment.note_id, func.max(MonitorComment.first_seen_at).label("latest"))
.where(MonitorComment.note_id.in_(set(counts)))
.group_by(MonitorComment.note_id)
)
if task_id is not None:
latest_query = latest_query.where(MonitorComment.task_id == task_id)
if scoped_ids is not None:
latest_query = latest_query.where(MonitorComment.task_id.in_(scoped_ids))
latest = {row.note_id: row.latest for row in (await session.execute(latest_query)).all()}
meta = await _note_meta_map(session, list(counts))
groups = [
{
"note_id": note_id,
"note_title": meta.get(note_id, {}).get("note_title", ""),
"note_cover": meta.get(note_id, {}).get("note_cover", ""),
"note_url": meta.get(note_id, {}).get("note_url", ""),
"creator_hash": meta.get(note_id, {}).get("creator_hash", ""),
"creator_name": meta.get(note_id, {}).get("creator_name", ""),
"comment_count": count,
"latest_at": latest.get(note_id, 0),
}
for note_id, count in counts.items()
]
groups.sort(key=lambda group: group["latest_at"], reverse=True)
return groups
async def list_events(
session: AsyncSession,
task_id: Optional[int] = None,
event_type: Optional[str] = None,
since_id: Optional[int] = None,
limit: int = 200,
platform: Optional[str] = None,
) -> List[Dict[str, Any]]:
query = select(MonitorEvent).order_by(MonitorEvent.id.desc()).limit(limit)
if task_id is not None:
query = query.where(MonitorEvent.task_id == task_id)
if event_type is not None:
query = query.where(MonitorEvent.type == event_type)
if since_id is not None:
query = query.where(MonitorEvent.id > since_id)
if platform is not None:
scoped = await platform_task_ids(session, platform)
if not scoped:
return []
query = query.where(MonitorEvent.task_id.in_(scoped))
return [
{
"id": row.id,
"task_id": row.task_id,
"run_id": row.run_id,
"type": row.type,
"severity": row.severity,
"target_kind": row.target_kind,
"target_id": row.target_id,
"title": row.title,
"created_at": row.created_at,
"is_read": row.is_read,
}
for row in (await session.scalars(query)).all()
]
async def list_runs(session: AsyncSession, task_id: int, limit: int = 50) -> List[Dict[str, Any]]:
rows = (
await session.scalars(
select(MonitorRun)
.where(MonitorRun.task_id == task_id)
.order_by(MonitorRun.id.desc())
.limit(limit)
)
).all()
return [
{
"id": row.id,
"task_id": row.task_id,
"status": row.status,
"trigger": row.trigger,
"queued_at": row.queued_at,
"started_at": row.started_at,
"finished_at": row.finished_at,
"exit_code": row.exit_code,
"notes_fetched": row.notes_fetched,
"comments_fetched": row.comments_fetched,
"new_notes": row.new_notes,
"new_comments": row.new_comments,
"is_baseline": row.is_baseline,
"max_comments_count": row.max_comments_count,
"error_message": row.error_message,
}
for row in rows
]
async def list_tasks(
session: AsyncSession, platform: Optional[str] = None
) -> List[Dict[str, Any]]:
query = select(MonitorTask).order_by(MonitorTask.id)
if platform is not None:
query = query.where(MonitorTask.platform == platform)
tasks = list((await session.scalars(query)).all())
if not tasks:
return []
counts = dict(
(
await session.execute(
select(MonitorTarget.task_id, func.count())
.group_by(MonitorTarget.task_id)
)
).all()
)
unread = dict(
(
await session.execute(
select(MonitorEvent.task_id, func.count())
.where(MonitorEvent.is_read.is_(False))
.group_by(MonitorEvent.task_id)
)
).all()
)
return [
{
"id": task.id,
"name": task.name,
"platform": task.platform,
"mode": task.mode,
"enabled": task.enabled,
"interval_minutes": task.interval_minutes,
**_schedule_fields(task),
"max_notes_count": task.max_notes_count,
"enable_comments": task.enable_comments,
"max_comments_count": task.max_comments_count,
"run_timeout_seconds": task.run_timeout_seconds,
"notify_enabled": task.notify_enabled,
"notify_failures": task.notify_failures,
"next_run_at": task.next_run_at,
"last_run_at": task.last_run_at,
"last_status": task.last_status,
"last_error": task.last_error,
"last_notified_at": task.last_notified_at,
"target_count": counts.get(task.id, 0),
"targets": [
{"id": t.id, "external_id": t.external_id, "raw_value": t.raw_value, "enabled": t.enabled}
for t in task.targets
],
"unread_events": unread.get(task.id, 0),
}
for task in tasks
]
def _schedule_fields(task: MonitorTask) -> Dict[str, Any]:
"""The schedule columns, plus the sentence the task list renders.
The label is composed here rather than in the frontend so the list and the
editor cannot drift on what a given schedule means.
"""
hours = schedule.parse_hours(task.schedule_hours)
days = schedule.parse_days(task.schedule_days)
return {
"schedule_mode": task.schedule_mode,
"schedule_hours": hours,
"schedule_days": days,
"schedule_minute": task.schedule_minute,
"schedule_label": schedule.describe(
mode=task.schedule_mode,
interval_minutes=task.interval_minutes,
hours=hours,
days=days,
minute=task.schedule_minute,
),
}
async def overview(session: AsyncSession, platform: Optional[str] = None) -> Dict[str, Any]:
"""Headline numbers for the dashboard tiles, scoped to one platform."""
now = get_current_timestamp()
day_ago = now - 24 * 60 * 60 * 1000
# Nothing but the task table carries a platform column, so the other counts
# are scoped through the platform's task ids.
scoped: Optional[List[int]] = None
if platform is not None:
scoped = await platform_task_ids(session, platform)
def by_task(stmt, column):
return stmt if scoped is None else stmt.where(column.in_(scoped))
task_count = select(func.count()).select_from(MonitorTask)
if platform is not None:
task_count = task_count.where(MonitorTask.platform == platform)
enabled_count = select(func.count()).select_from(MonitorTask).where(
MonitorTask.enabled.is_(True)
)
if platform is not None:
enabled_count = enabled_count.where(MonitorTask.platform == platform)
return {
"platform": platform,
"tasks": await session.scalar(task_count) or 0,
"enabled_tasks": await session.scalar(enabled_count) or 0,
"notes": await session.scalar(
by_task(select(func.count()).select_from(MonitorNote), MonitorNote.task_id)
)
or 0,
"comments": await session.scalar(
by_task(select(func.count()).select_from(MonitorComment), MonitorComment.task_id)
)
or 0,
"events_24h": await session.scalar(
by_task(
select(func.count())
.select_from(MonitorEvent)
.where(MonitorEvent.created_at >= day_ago),
MonitorEvent.task_id,
)
)
or 0,
"unread_events": await session.scalar(
by_task(
select(func.count())
.select_from(MonitorEvent)
.where(MonitorEvent.is_read.is_(False)),
MonitorEvent.task_id,
)
)
or 0,
"running_runs": await session.scalar(
by_task(
select(func.count())
.select_from(MonitorRun)
.where(MonitorRun.status == "running"),
MonitorRun.task_id,
)
)
or 0,
}
async def mark_events_read(session: AsyncSession, task_id: Optional[int] = None) -> int:
query = select(MonitorEvent).where(MonitorEvent.is_read.is_(False))
if task_id is not None:
query = query.where(MonitorEvent.task_id == task_id)
rows = list((await session.scalars(query)).all())
for row in rows:
row.is_read = True
return len(rows)