分组是从作品推出来的(按作品的 creator_hash 归组),于是没有作品的博主根本 不进列表:目标加了、资料采到了、粉丝数就躺在库里,界面上什么都看不见。 而这类博主恰恰是最该看见的 —— 还在涨粉,只是最近没发东西。藏起来正好藏反了。 线上就有一个:5 个目标里 3 个没作品,那 3 个连同已采到的粉丝数一起消失了。 改成 **账号快照 ∪ 作品** 两个来源: * 有快照没作品 → 一个 0 篇的组,备注和粉丝数照常显示,组里写「暂无作品」; * 有作品没快照 → 一个没有账号指标的组(小红书那条路不产生快照)。 /notes 因此多返回一份 `creators`,而不是让前端从作品里推 —— 作品推不出上面 第一类人。组头改读它,`MonitorNote` 上那几个字段降级成「顺着作品问作者」用。
1093 lines
41 KiB
Python
1093 lines
41 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/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
|
||
|
||
|
||
async def list_notes(
|
||
session: AsyncSession,
|
||
task_id: Optional[int] = None,
|
||
only_new: bool = False,
|
||
limit: int = 200,
|
||
platform: Optional[str] = None,
|
||
) -> List[Dict[str, Any]]:
|
||
"""Tracked notes with their latest metrics and change vs the previous run."""
|
||
query = select(MonitorNote).order_by(MonitorNote.last_seen_at.desc()).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 []
|
||
|
||
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)
|