Files
MediaCrawler/api/monitor/service.py
T
butubb 06718a1351
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(monitor): 抖音接入博主监控
上游爬虫本身不缺抖音能力(三模式、四项指标、二级评论都与小红书对等、指标还是同名同列),
缺的全在监控层的适配。这次把「平台之间不一样」的管子集中到一个新模块,再把散落的
xhs 硬编码接上去。

* 新增 api/monitor/adapters.py:产物目录名、jsonl 字段别名、目标链接形态与正则、
  通知链接模板。不放进 platforms.py 是因为那个模块被 describe_all() 整个序列化进
  /api/config/platforms 交给前端,塞进正则和目录名会让爬虫内部细节漏进 API 载荷。
  代价是两个注册表可能漂移,用一条测试钉住「声明接通就必须有适配器」。
* 两个必须知道的坑,都在这版里处理掉了:
  1) 抖音的平台 id 是 dy,而 store 把产物写在 douyin/ 下(store/douyin/_store_impl.py:47)。
     不改就是 ingest 一个文件都读不到 —— 不报错,只是 0 条,然后被冒充成「疑似登录失效」。
  2) 抖音的作品没有 note_id(叫 aweme_id)、评论也用 aweme_id 指作品。ingest 第一步是
     `if not note_id: continue`,不映射就逐条全丢。
  另外抖音顶层评论的 parent_comment_id 是字符串 "0",归一成空串,免得前端多出悬空的父节点。
* 顺带把「东西抓到了、只是没落在期望目录里」单独识别出来。这类故障的现象和登录失效
  一模一样,按登录失效报会把人指去查完全错误的方向。
* 修两个既有 bug(今天只有小红书所以无害,加抖音就踩响):
  - service.py update_task 换目标时漏传 task.platform,回落到默认小红书
  - scheduler.py 取 cookie 没传 platform,抖音任务会读着小红书那份 cookie 不动
* 行为变更(已与用户确认):cookie 闸门改成「没 cookie 且没开 CDP」才跳过。
  CDP 模式下登录态来自被接管的浏览器,粘不粘 cookie 由不得它决定;不放行的话,
  选了「接管已有 Chrome」却没粘 cookie 的用户会看到任务永远不触发,而且不报错。
  副作用是开启了 CDP 的小红书任务也不再被该闸门拦住 —— 语义上是对的。
* 目标输入框的示例链接与措辞改由能力矩阵提供(notes_label 抖音说「作品」、小红书说
  「笔记」;「建议只填纯 ID」是小红书专属劝告,抖音链接不带令牌,不再显示)。

测试 +22 条(858 通过),其中最关键的是「抖音作品/评论不被静默丢弃」与「产物目录名
不等于平台 id」两条 —— 都是把最难查的失败模式钉死在回归网里。

注意:抖音这条路的**端到端尚未验证**,需要一份可用的抖音登录态(CDP 那台 Chrome 里
登录,或导出一份 cookie)。单测覆盖的是解析与入库,真实抓取还没跑过。
2026-10-10 14:55:28 +08:00

830 lines
30 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
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,
MonitorEvent,
MonitorNote,
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 _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 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]] = []
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
result.append(
{
"task_id": note.task_id,
"note_id": note.note_id,
"title": note.title,
"note_url": note.note_url,
# 按博主分组用。creator_hash 是唯一稳定的创作者标识(原始 user_id
# 被爬虫刻意匿名化了),creator_name 是已脱敏的昵称。
"creator_hash": note.creator_hash,
"creator_name": note.creator_name,
# 优先给本地缓存地址:远程地址带签名、会过期(实测隔天即 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,
"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", ""),
}
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)