# -*- coding: utf-8 -*- # Copyright (c) 2025 relakkes@gmail.com # # 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, # 作品的发布时间(爬虫侧:小红书叫 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, # 优先给本地缓存地址:远程地址带签名、会过期(实测隔天即 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)