feat(notify): 通知拆成「新作品」与「异常」两个开关,异常默认开
问题:cookie 过期导致任务失败,但没有任何通知。查下来不是代码问题 —— notify_enabled 在两个任务上都是 False,而它默认就是关的,事件(run_failed) 也确实生成了,只卡在最后一道闸门。 但那个默认值是错的。代码里的理由是「一条任务列表都推到一个群会很快变吵,所以默认静默」, 这个理由对新作品成立(可能每轮都有),对失败不成立:一次登录态失效意味着这个任务事实上 已经死了,而你不会知道,直到某天发现数据停在几周前。最该被告知的就是这种情况。 现在拆开: - notify_enabled —— 推送新作品,可能每轮都有,默认关 - notify_failures —— 推送异常(登录失效/运行失败/没抓到数据),默认开 事件按开关过滤(build_run_message):只勾了「新作品」的任务不该因为一次失败被推消息, 反之亦然,否则拆开开关就没有意义。已有任务由 _ensure_columns 补上 notify_failures=1, 所以会自动开始收到异常推送。 列名 notify_enabled 是历史遗留(它早先是唯一的通知开关),语义已收窄为「新作品」, 用注释写明,不做列重命名 —— 那需要单独的迁移,不值为一个内部工具做。
This commit is contained in:
@@ -107,9 +107,15 @@ class MonitorTask(MonitorBase):
|
|||||||
max_comments_count: Mapped[int] = mapped_column(Integer, nullable=False, default=50)
|
max_comments_count: Mapped[int] = mapped_column(Integer, nullable=False, default=50)
|
||||||
run_timeout_seconds: Mapped[int] = mapped_column(Integer, nullable=False, default=3600)
|
run_timeout_seconds: Mapped[int] = mapped_column(Integer, nullable=False, default=3600)
|
||||||
|
|
||||||
# Push notifications are opt-in per task. A task list that all pushes to one
|
# 通知分成两类,因为它们的性质完全不同:
|
||||||
# webhook turns noisy fast, so silence is the default.
|
#
|
||||||
|
# * `notify_enabled` —— **推送新作品**。可能每轮都有,一条任务列表都推到同一个群
|
||||||
|
# 会很快变吵,所以默认关。(列名是历史遗留:它早先是唯一的通知开关。)
|
||||||
|
# * `notify_failures` —— **推送异常**(登录失效 / 运行失败 / 没抓到数据)。频率低,
|
||||||
|
# 而且一旦发生就意味着这个任务从此**默默采不到任何东西**,你会一直不知道,
|
||||||
|
# 直到某天发现数据停在几周前。这正是最该被告知的情况,所以默认**开**。
|
||||||
notify_enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
|
notify_enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
|
||||||
|
notify_failures: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
|
||||||
|
|
||||||
# Scheduler state. Persisted so the schedule survives an API restart.
|
# Scheduler state. Persisted so the schedule survives an API restart.
|
||||||
next_run_at: Mapped[Optional[int]] = mapped_column(BigInteger, index=True)
|
next_run_at: Mapped[Optional[int]] = mapped_column(BigInteger, index=True)
|
||||||
|
|||||||
+17
-3
@@ -107,14 +107,27 @@ async def build_run_message(
|
|||||||
task: MonitorTask,
|
task: MonitorTask,
|
||||||
run: MonitorRun,
|
run: MonitorRun,
|
||||||
) -> Optional[str]:
|
) -> Optional[str]:
|
||||||
"""Compose one markdown summary for a finished run, or None if nothing to say."""
|
"""Compose one markdown summary for a finished run, or None if nothing to say.
|
||||||
|
|
||||||
|
事件按开关过滤:只勾了「新作品」的任务,不该因为一次失败被推消息,反之亦然 ——
|
||||||
|
否则拆开这两个开关就没有意义了。
|
||||||
|
"""
|
||||||
|
allowed = []
|
||||||
|
if task.notify_enabled:
|
||||||
|
allowed.append(EVENT_NEW_NOTE)
|
||||||
|
if task.notify_failures:
|
||||||
|
allowed.extend([EVENT_AUTH_FAILURE, EVENT_RUN_FAILED, EVENT_NO_DATA])
|
||||||
|
|
||||||
|
if not allowed:
|
||||||
|
return None
|
||||||
|
|
||||||
events = list(
|
events = list(
|
||||||
(
|
(
|
||||||
await session.scalars(
|
await session.scalars(
|
||||||
select(MonitorEvent)
|
select(MonitorEvent)
|
||||||
.where(
|
.where(
|
||||||
MonitorEvent.run_id == run.id,
|
MonitorEvent.run_id == run.id,
|
||||||
MonitorEvent.type.in_(NOTIFIABLE_EVENT_TYPES),
|
MonitorEvent.type.in_(allowed),
|
||||||
)
|
)
|
||||||
.order_by(MonitorEvent.id)
|
.order_by(MonitorEvent.id)
|
||||||
)
|
)
|
||||||
@@ -176,7 +189,8 @@ async def notify_run(session: AsyncSession, task: MonitorTask, run: MonitorRun)
|
|||||||
Returns the message that was sent, or None. Never raises.
|
Returns the message that was sent, or None. Never raises.
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
if not task.notify_enabled:
|
# 两个开关是分开的:只开「异常」不该因为新作品而发消息,反之亦然。
|
||||||
|
if not (task.notify_enabled or task.notify_failures):
|
||||||
return None
|
return None
|
||||||
|
|
||||||
webhook_url = await get_webhook_url(session)
|
webhook_url = await get_webhook_url(session)
|
||||||
|
|||||||
@@ -183,6 +183,7 @@ async def create_task(session: AsyncSession, payload: Dict[str, Any]) -> Monitor
|
|||||||
max_comments_count=payload.get("max_comments_count") or defaults["max_comments_count"],
|
max_comments_count=payload.get("max_comments_count") or defaults["max_comments_count"],
|
||||||
run_timeout_seconds=payload.get("run_timeout_seconds", 3600),
|
run_timeout_seconds=payload.get("run_timeout_seconds", 3600),
|
||||||
notify_enabled=payload.get("notify_enabled", False),
|
notify_enabled=payload.get("notify_enabled", False),
|
||||||
|
notify_failures=payload.get("notify_failures", True),
|
||||||
next_run_at=schedule.next_occurrence(
|
next_run_at=schedule.next_occurrence(
|
||||||
mode=schedule_mode,
|
mode=schedule_mode,
|
||||||
interval_minutes=interval_minutes,
|
interval_minutes=interval_minutes,
|
||||||
@@ -240,6 +241,7 @@ async def update_task(session: AsyncSession, task_id: int, payload: Dict[str, An
|
|||||||
"max_comments_count",
|
"max_comments_count",
|
||||||
"run_timeout_seconds",
|
"run_timeout_seconds",
|
||||||
"notify_enabled",
|
"notify_enabled",
|
||||||
|
"notify_failures",
|
||||||
):
|
):
|
||||||
if field in payload and payload[field] is not None:
|
if field in payload and payload[field] is not None:
|
||||||
setattr(task, field, payload[field])
|
setattr(task, field, payload[field])
|
||||||
@@ -706,6 +708,7 @@ async def list_tasks(
|
|||||||
"max_comments_count": task.max_comments_count,
|
"max_comments_count": task.max_comments_count,
|
||||||
"run_timeout_seconds": task.run_timeout_seconds,
|
"run_timeout_seconds": task.run_timeout_seconds,
|
||||||
"notify_enabled": task.notify_enabled,
|
"notify_enabled": task.notify_enabled,
|
||||||
|
"notify_failures": task.notify_failures,
|
||||||
"next_run_at": task.next_run_at,
|
"next_run_at": task.next_run_at,
|
||||||
"last_run_at": task.last_run_at,
|
"last_run_at": task.last_run_at,
|
||||||
"last_status": task.last_status,
|
"last_status": task.last_status,
|
||||||
|
|||||||
@@ -60,9 +60,10 @@ class MonitorTaskCreate(BaseModel):
|
|||||||
max_comments_count: Optional[int] = Field(default=None, ge=1, le=500)
|
max_comments_count: Optional[int] = Field(default=None, ge=1, le=500)
|
||||||
run_timeout_seconds: int = Field(default=3600, ge=60, le=86400)
|
run_timeout_seconds: int = Field(default=3600, ge=60, le=86400)
|
||||||
enabled: bool = True
|
enabled: bool = True
|
||||||
# Push a WeCom summary for runs that failed or found new works. Opt-in per
|
# 两类通知分开:新作品可能每轮都有(默认关,避免刷屏),
|
||||||
# task so a single webhook does not get flooded.
|
# 异常频率低且意味着任务已经停止工作(默认开,否则你会一直不知道)。
|
||||||
notify_enabled: bool = False
|
notify_enabled: bool = False
|
||||||
|
notify_failures: bool = True
|
||||||
# Raw pasted values: full URLs or bare ids, in either form.
|
# Raw pasted values: full URLs or bare ids, in either form.
|
||||||
targets: List[str] = Field(min_length=1)
|
targets: List[str] = Field(min_length=1)
|
||||||
|
|
||||||
@@ -98,6 +99,7 @@ class MonitorTaskUpdate(BaseModel):
|
|||||||
max_comments_count: Optional[int] = Field(default=None, ge=1, le=500)
|
max_comments_count: Optional[int] = Field(default=None, ge=1, le=500)
|
||||||
run_timeout_seconds: Optional[int] = Field(default=None, ge=60, le=86400)
|
run_timeout_seconds: Optional[int] = Field(default=None, ge=60, le=86400)
|
||||||
notify_enabled: Optional[bool] = None
|
notify_enabled: Optional[bool] = None
|
||||||
|
notify_failures: Optional[bool] = None
|
||||||
# When present, replaces the whole target list.
|
# When present, replaces the whole target list.
|
||||||
targets: Optional[List[str]] = None
|
targets: Optional[List[str]] = None
|
||||||
|
|
||||||
|
|||||||
@@ -118,6 +118,8 @@ export function TaskEditorDialog({ open, onOpenChange, task }: TaskEditorDialogP
|
|||||||
const [enableComments, setEnableComments] = useState(true)
|
const [enableComments, setEnableComments] = useState(true)
|
||||||
const [maxComments, setMaxComments] = useState('50')
|
const [maxComments, setMaxComments] = useState('50')
|
||||||
const [notifyEnabled, setNotifyEnabled] = useState(false)
|
const [notifyEnabled, setNotifyEnabled] = useState(false)
|
||||||
|
// 异常推送默认开:失败意味着这个任务从此默默采不到东西,而你不会知道。
|
||||||
|
const [notifyFailures, setNotifyFailures] = useState(true)
|
||||||
const [targets, setTargets] = useState('')
|
const [targets, setTargets] = useState('')
|
||||||
|
|
||||||
// Reset the form whenever the dialog is (re)opened. For a new task the
|
// Reset the form whenever the dialog is (re)opened. For a new task the
|
||||||
@@ -142,6 +144,7 @@ export function TaskEditorDialog({ open, onOpenChange, task }: TaskEditorDialogP
|
|||||||
String(task?.max_comments_count ?? settings?.values['collect.default_max_comments'] ?? 50),
|
String(task?.max_comments_count ?? settings?.values['collect.default_max_comments'] ?? 50),
|
||||||
)
|
)
|
||||||
setNotifyEnabled(task?.notify_enabled ?? false)
|
setNotifyEnabled(task?.notify_enabled ?? false)
|
||||||
|
setNotifyFailures(task?.notify_failures ?? true)
|
||||||
setTargets(task ? task.targets.map((t) => t.raw_value || t.external_id).join('\n') : '')
|
setTargets(task ? task.targets.map((t) => t.raw_value || t.external_id).join('\n') : '')
|
||||||
}, [open, task, settings])
|
}, [open, task, settings])
|
||||||
|
|
||||||
@@ -185,6 +188,7 @@ export function TaskEditorDialog({ open, onOpenChange, task }: TaskEditorDialogP
|
|||||||
run_timeout_seconds: 3600,
|
run_timeout_seconds: 3600,
|
||||||
enabled: true,
|
enabled: true,
|
||||||
notify_enabled: notifyEnabled,
|
notify_enabled: notifyEnabled,
|
||||||
|
notify_failures: notifyFailures,
|
||||||
targets: targetList,
|
targets: targetList,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -455,24 +459,55 @@ export function TaskEditorDialog({ open, onOpenChange, task }: TaskEditorDialogP
|
|||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
<div className="flex items-start gap-2 rounded-md border border-cyber-border-subtle bg-cyber-bg-tertiary/40 p-3">
|
<div className="rounded-md border border-cyber-border-subtle bg-cyber-bg-tertiary/40 p-3 space-y-3">
|
||||||
<Checkbox
|
{/* 两类通知的性质完全不同,所以分成两个开关:
|
||||||
id="notify-enabled"
|
异常低频且意味着任务已经停止工作 —— 默认开;
|
||||||
checked={notifyEnabled}
|
新作品可能每轮都有 —— 默认关,否则会刷屏。 */}
|
||||||
onCheckedChange={(checked) => setNotifyEnabled(checked === true)}
|
<div className="flex items-start gap-2">
|
||||||
/>
|
<Checkbox
|
||||||
<div className="space-y-0.5">
|
id="notify-failures"
|
||||||
<label
|
checked={notifyFailures}
|
||||||
htmlFor="notify-enabled"
|
onCheckedChange={(checked) => setNotifyFailures(checked === true)}
|
||||||
className="text-xs font-mono text-cyber-text-primary cursor-pointer"
|
/>
|
||||||
>
|
<div className="space-y-0.5">
|
||||||
推送企业微信通知
|
<label
|
||||||
</label>
|
htmlFor="notify-failures"
|
||||||
<p className="text-[10px] font-mono text-cyber-text-muted">
|
className="text-xs font-mono text-cyber-text-primary cursor-pointer"
|
||||||
仅在本任务**采集失败 / 登录态失效**或**发现新作品**时推送,
|
>
|
||||||
一轮只发一条汇总。需先在监控页配置 Webhook 地址。
|
推送异常通知(建议保持开启)
|
||||||
</p>
|
</label>
|
||||||
|
<p className="text-[10px] font-mono text-cyber-text-muted leading-relaxed">
|
||||||
|
登录态失效、采集进程失败、一篇都没抓到时推送。
|
||||||
|
<span className="text-cyber-neon-cyan">
|
||||||
|
异常意味着这个任务从此默默采不到任何东西
|
||||||
|
</span>
|
||||||
|
—— 关掉的话你不会知道,直到某天发现数据停在几周前。
|
||||||
|
</p>
|
||||||
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
|
<div className="flex items-start gap-2 border-t border-cyber-border-subtle pt-3">
|
||||||
|
<Checkbox
|
||||||
|
id="notify-enabled"
|
||||||
|
checked={notifyEnabled}
|
||||||
|
onCheckedChange={(checked) => setNotifyEnabled(checked === true)}
|
||||||
|
/>
|
||||||
|
<div className="space-y-0.5">
|
||||||
|
<label
|
||||||
|
htmlFor="notify-enabled"
|
||||||
|
className="text-xs font-mono text-cyber-text-primary cursor-pointer"
|
||||||
|
>
|
||||||
|
推送新作品通知
|
||||||
|
</label>
|
||||||
|
<p className="text-[10px] font-mono text-cyber-text-muted leading-relaxed">
|
||||||
|
发现新作品时推送。监控多个博主时可能每轮都有,容易刷屏,所以默认关闭。
|
||||||
|
</p>
|
||||||
|
</div>
|
||||||
|
</div>
|
||||||
|
|
||||||
|
<p className="border-t border-cyber-border-subtle pt-2 text-[10px] font-mono text-cyber-text-muted">
|
||||||
|
两者都是一轮只发一条汇总。需先在右上角「系统设置」里配置企业微信 Webhook 地址。
|
||||||
|
</p>
|
||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
|
|||||||
@@ -56,8 +56,10 @@ export interface MonitorTask {
|
|||||||
enable_comments: boolean
|
enable_comments: boolean
|
||||||
max_comments_count: number
|
max_comments_count: number
|
||||||
run_timeout_seconds: number
|
run_timeout_seconds: number
|
||||||
/** Opt-in per task so one webhook does not get flooded. */
|
/** 推送**新作品**。可能每轮都有,默认关以免刷屏。 */
|
||||||
notify_enabled: boolean
|
notify_enabled: boolean
|
||||||
|
/** 推送**异常**(登录失效/运行失败/没抓到数据)。默认开。 */
|
||||||
|
notify_failures: boolean
|
||||||
/** Epoch milliseconds. */
|
/** Epoch milliseconds. */
|
||||||
next_run_at: number | null
|
next_run_at: number | null
|
||||||
last_run_at: number | null
|
last_run_at: number | null
|
||||||
@@ -258,6 +260,7 @@ export interface TaskCreatePayload {
|
|||||||
run_timeout_seconds: number
|
run_timeout_seconds: number
|
||||||
enabled: boolean
|
enabled: boolean
|
||||||
notify_enabled: boolean
|
notify_enabled: boolean
|
||||||
|
notify_failures: boolean
|
||||||
targets: string[]
|
targets: string[]
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user