Files
MediaCrawler/api/monitor/service.py
T
butubb 20e672834c
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(monitor): 博主的粉丝数,以及给作品起备注
两件都是「作品栏里把这东西认出来」的延伸:

* **账号级指标**:作品列表只会说「这条涨了多少赞」,说不了「这个人整个
  账号的粉丝在涨还是在掉」。抖音的资料接口本来就有粉丝数/总获赞/作品数,
  每轮顺手记一条快照(`monitor_creator_stat`,粒度 = 任务×博主×轮次,
  和作品指标同形)。组头显示最近一条。

  快照在「一条作品都没采到」的早退**之前**落:作品列表被风控挡住的那一轮,
  正是「粉丝还在涨、但新作品没在发现」最该被看见的时刻。

* **作品备注**:博主备注回答「这个账号是谁」,这条回答「这条我要盯着」。
  一个博主底下常常只有一两件值得盯的作品,所以不能合并成一条。键取
  (platform, note_id),跨任务共用一份。

两边都守住同一条口径:**不知道就是 null,不写成 0** —— 0 在趋势图上是一条
砸到底的线,和「还没采到」是两回事。
2026-10-10 18:09:03 +08:00

991 lines
37 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: Sequence[int],
creator_hashes: Sequence[str],
) -> Dict[tuple, "MonitorCreatorStat"]:
"""``(task_id, creator_hash) -> 最近一条``账号级快照。
按 run_id 而不是 captured_at 取「最近」:和作品指标用的是同一个口径,两者放一起
看才不会出现「作品数据来自第 8 轮、粉丝数来自第 9 轮」这种对不上的情况。
"""
if not task_ids or not creator_hashes:
return {}
rows = (
await session.scalars(
select(MonitorCreatorStat)
.where(
MonitorCreatorStat.task_id.in_(list(task_ids)),
MonitorCreatorStat.creator_hash.in_(list(creator_hashes)),
)
.order_by(MonitorCreatorStat.run_id.desc())
)
).all()
latest: Dict[tuple, MonitorCreatorStat] = {}
for row in rows:
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 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)