migrace /remind na sqlite
This commit is contained in:
@@ -3,41 +3,37 @@
|
||||
# requires-python = ">=3.11"
|
||||
# dependencies = ["croniter", "pyyaml"]
|
||||
# ///
|
||||
"""Deterministic reminder sender.
|
||||
"""Deterministic reminder sender backed by SQLite.
|
||||
|
||||
Runs every minute from the nanobot user crontab (NOT through the agent).
|
||||
Reads reminder.yaml, finds reminders due this minute, sends each directly to
|
||||
Telegram via the Bot API, appends the delivery to reminder.log, and dedups via
|
||||
.reminder_state.json so each scheduled fire is delivered exactly once.
|
||||
|
||||
No LLM and no nanobot process involved on purpose -- see knowledge.md/history
|
||||
for why the previous agent-driven cron job spammed empty-output messages.
|
||||
Runs every minute from the nanobot user crontab.
|
||||
Reads reminders from SQLite, finds due fires, sends each directly to Telegram,
|
||||
logs delivery to reminder.log, and dedups via reminder_fires table.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
from datetime import datetime
|
||||
from datetime import date, datetime, timedelta
|
||||
from pathlib import Path
|
||||
from zoneinfo import ZoneInfo
|
||||
|
||||
import yaml
|
||||
from croniter import croniter
|
||||
from db import get_db, init_db, log_operation
|
||||
from random_times import compute_fire_times
|
||||
|
||||
WORKSPACE = Path(__file__).resolve().parent.parent.parent.parent # .../workspace
|
||||
REMINDER_YAML = WORKSPACE / "reminder.yaml"
|
||||
STATE_FILE = WORKSPACE / ".reminder_state.json"
|
||||
LOG_DIR = WORKSPACE / "log"
|
||||
LOG_FILE = LOG_DIR / "reminder.log"
|
||||
WORKSPACE = Path(__file__).resolve().parent.parent.parent.parent
|
||||
DEFAULT_DB_PATH = WORKSPACE / "db" / "reminders.sqlite"
|
||||
DB_PATH = Path(os.environ.get("REMIND_DB", str(DEFAULT_DB_PATH)))
|
||||
CONFIG = Path.home() / ".nanobot" / "config.json"
|
||||
|
||||
TZ = ZoneInfo("Europe/Prague")
|
||||
CHAT_ID = "8826147089" # Telegram user id (Martin); same target the old cron job used
|
||||
CHAT_ID = "8826147089"
|
||||
TOLERANCE_SECONDS = 60
|
||||
|
||||
|
||||
def _telegram_token() -> str:
|
||||
@@ -54,96 +50,165 @@ def _send_telegram(text: str) -> None:
|
||||
resp.read()
|
||||
|
||||
|
||||
def _load_state() -> dict:
|
||||
if STATE_FILE.exists():
|
||||
def _now() -> datetime:
|
||||
return datetime.now(TZ).replace(tzinfo=None)
|
||||
|
||||
|
||||
def _due_at(conn, now: datetime) -> list[dict]:
|
||||
"""Find due one-time reminders."""
|
||||
since = (now - timedelta(seconds=TOLERANCE_SECONDS)).isoformat(timespec="seconds")
|
||||
until = now.isoformat(timespec="seconds")
|
||||
rows = conn.execute(
|
||||
"""
|
||||
SELECT r.id, r.text, sa.id AS schedule_id, sa.at_datetime AS fire_time
|
||||
FROM reminders r
|
||||
JOIN schedule_at sa ON sa.reminder_id = r.id
|
||||
WHERE r.enabled = 1 AND r.deleted_at IS NULL
|
||||
AND sa.at_datetime > ?
|
||||
AND sa.at_datetime <= ?
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM reminder_fires rf
|
||||
WHERE rf.reminder_id = r.id AND rf.schedule_id = sa.id
|
||||
AND rf.schedule_type = 'at' AND rf.fire_time = sa.at_datetime
|
||||
AND rf.status = 'delivered'
|
||||
)
|
||||
""",
|
||||
(since, until),
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
def _due_cron(conn, now: datetime) -> list[dict]:
|
||||
"""Find due cron reminders."""
|
||||
rows = conn.execute(
|
||||
"""
|
||||
SELECT r.id, r.text, sc.id AS schedule_id, sc.cron_expr
|
||||
FROM reminders r
|
||||
JOIN schedule_cron sc ON sc.reminder_id = r.id
|
||||
WHERE r.enabled = 1 AND r.deleted_at IS NULL
|
||||
"""
|
||||
).fetchall()
|
||||
due = []
|
||||
for row in rows:
|
||||
prev = croniter(row["cron_expr"], now + timedelta(seconds=1)).get_prev(datetime)
|
||||
if 0 <= (now - prev).total_seconds() < TOLERANCE_SECONDS:
|
||||
fire_iso = prev.isoformat(timespec="seconds")
|
||||
already = conn.execute(
|
||||
"""
|
||||
SELECT 1 FROM reminder_fires
|
||||
WHERE reminder_id = ? AND schedule_id = ? AND schedule_type = 'cron'
|
||||
AND fire_time = ? AND status = 'delivered'
|
||||
""",
|
||||
(row["id"], row["schedule_id"], fire_iso),
|
||||
).fetchone()
|
||||
if not already:
|
||||
due.append({
|
||||
"id": row["id"],
|
||||
"text": row["text"],
|
||||
"schedule_id": row["schedule_id"],
|
||||
"fire_time": fire_iso,
|
||||
})
|
||||
return due
|
||||
|
||||
|
||||
def _due_random(conn, now: datetime) -> list[dict]:
|
||||
"""Find due random reminders."""
|
||||
rows = conn.execute(
|
||||
"""
|
||||
SELECT r.id, r.text, sr.id AS schedule_id, sr.times_per_day, sr.window_start, sr.window_end,
|
||||
sr.days_filter, sr.from_date, sr.until_date
|
||||
FROM reminders r
|
||||
JOIN schedule_random sr ON sr.reminder_id = r.id
|
||||
WHERE r.enabled = 1 AND r.deleted_at IS NULL
|
||||
"""
|
||||
).fetchall()
|
||||
due = []
|
||||
for row in rows:
|
||||
cfg = {
|
||||
"times_per_day": row["times_per_day"],
|
||||
"window": f"{_minutes_to_hhmm(row['window_start'])}-{_minutes_to_hhmm(row['window_end'])}",
|
||||
}
|
||||
if row["days_filter"]:
|
||||
cfg["days"] = row["days_filter"]
|
||||
if row["from_date"]:
|
||||
cfg["from"] = row["from_date"]
|
||||
if row["until_date"]:
|
||||
cfg["until"] = row["until_date"]
|
||||
try:
|
||||
data = json.loads(STATE_FILE.read_text(encoding="utf-8"))
|
||||
return data if isinstance(data, dict) else {}
|
||||
except Exception:
|
||||
return {}
|
||||
return {}
|
||||
fires = compute_fire_times(now.date(), row["text"], cfg)
|
||||
except ValueError as exc:
|
||||
print(f"remind_send: bad random config for {row['text']!r}: {exc}", file=sys.stderr)
|
||||
continue
|
||||
for ft in fires:
|
||||
if 0 <= (now - ft).total_seconds() < TOLERANCE_SECONDS:
|
||||
fire_iso = ft.isoformat(timespec="seconds")
|
||||
already = conn.execute(
|
||||
"""
|
||||
SELECT 1 FROM reminder_fires
|
||||
WHERE reminder_id = ? AND schedule_id = ? AND schedule_type = 'random'
|
||||
AND fire_time = ? AND status = 'delivered'
|
||||
""",
|
||||
(row["id"], row["schedule_id"], fire_iso),
|
||||
).fetchone()
|
||||
if not already:
|
||||
due.append({
|
||||
"id": row["id"],
|
||||
"text": row["text"],
|
||||
"schedule_id": row["schedule_id"],
|
||||
"fire_time": fire_iso,
|
||||
})
|
||||
return due
|
||||
|
||||
|
||||
def _key(text: str) -> str:
|
||||
return hashlib.sha1(text.encode("utf-8")).hexdigest()[:8]
|
||||
def _minutes_to_hhmm(total: int) -> str:
|
||||
return f"{total // 60:02d}:{total % 60:02d}"
|
||||
|
||||
|
||||
def _due_fire(item: dict, now: datetime) -> datetime | None:
|
||||
"""Most recent scheduled fire-time within the last 60s, or None."""
|
||||
fire: datetime | None = None
|
||||
|
||||
at_str = item.get("at")
|
||||
if at_str:
|
||||
at_time = datetime.fromisoformat(at_str).replace(tzinfo=None)
|
||||
if 0 <= (now - at_time).total_seconds() < 60:
|
||||
fire = at_time
|
||||
|
||||
for at_str in item.get("at_times", []):
|
||||
at_time = datetime.fromisoformat(at_str).replace(tzinfo=None)
|
||||
if 0 <= (now - at_time).total_seconds() < 60 and (fire is None or at_time > fire):
|
||||
fire = at_time
|
||||
|
||||
for expr in item.get("cron_exprs", []):
|
||||
prev = croniter(expr, now).get_prev(datetime)
|
||||
if 0 <= (now - prev).total_seconds() < 60 and (fire is None or prev > fire):
|
||||
fire = prev
|
||||
|
||||
random_cfg = item.get("random")
|
||||
if random_cfg:
|
||||
try:
|
||||
for ft in compute_fire_times(now.date(), (item.get("text") or "").strip(), random_cfg):
|
||||
if 0 <= (now - ft).total_seconds() < 60 and (fire is None or ft > fire):
|
||||
fire = ft
|
||||
except ValueError as exc: # malformed config: skip this reminder, keep others working
|
||||
print(f"remind_send: bad random config for {item.get('text')!r}: {exc}", file=sys.stderr)
|
||||
|
||||
return fire
|
||||
def _record_fire(conn, reminder_id: int, schedule_id: int, schedule_type: str, fire_time: str, status: str, error: str | None = None) -> None:
|
||||
now = _now().isoformat(timespec="seconds")
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO reminder_fires (reminder_id, schedule_id, schedule_type, fire_time, delivered_at, status, error_message)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(reminder_id, schedule_id, schedule_type, fire_time, now if status == "delivered" else None, status, error),
|
||||
)
|
||||
|
||||
|
||||
def main() -> None:
|
||||
if not REMINDER_YAML.exists():
|
||||
return
|
||||
if not DB_PATH.exists():
|
||||
init_db(DB_PATH)
|
||||
|
||||
data = yaml.safe_load(REMINDER_YAML.read_text(encoding="utf-8")) or {}
|
||||
now = datetime.now(TZ).replace(tzinfo=None)
|
||||
conn = get_db(DB_PATH)
|
||||
try:
|
||||
now = _now()
|
||||
due = _due_at(conn, now) + _due_cron(conn, now) + _due_random(conn, now)
|
||||
|
||||
state = _load_state()
|
||||
fresh: dict[str, str] = {}
|
||||
for fire in due:
|
||||
text = fire["text"]
|
||||
rid = fire["id"]
|
||||
sid = fire["schedule_id"]
|
||||
ft = fire["fire_time"]
|
||||
# Determine schedule_type from which query produced it
|
||||
# We can infer: if 'schedule_id' came from schedule_at, it's 'at'
|
||||
# But we don't have that info here. Let's look it up.
|
||||
st = conn.execute(
|
||||
"SELECT 'at' FROM schedule_at WHERE id = ? UNION ALL SELECT 'cron' FROM schedule_cron WHERE id = ? UNION ALL SELECT 'random' FROM schedule_random WHERE id = ?",
|
||||
(sid, sid, sid),
|
||||
).fetchone()
|
||||
schedule_type = st[0] if st else "unknown"
|
||||
|
||||
for item in data.get("reminders", []):
|
||||
text = (item.get("text") or "").strip()
|
||||
if not text:
|
||||
continue
|
||||
key = _key(text)
|
||||
last = state.get(key)
|
||||
try:
|
||||
_send_telegram(f"⏰ Reminder: {text}")
|
||||
except Exception as e:
|
||||
print(f"remind_send: delivery failed for {text!r}: {e}", file=sys.stderr)
|
||||
_record_fire(conn, rid, sid, schedule_type, ft, "failed", str(e))
|
||||
continue
|
||||
|
||||
fire = _due_fire(item, now)
|
||||
if fire is None:
|
||||
if last: # preserve dedup info for reminders not due this minute
|
||||
fresh[key] = last
|
||||
continue
|
||||
|
||||
fire_iso = fire.isoformat()
|
||||
if last == fire_iso: # this exact fire was already delivered
|
||||
fresh[key] = last
|
||||
continue
|
||||
|
||||
try:
|
||||
_send_telegram(f"⏰ Reminder: {text}")
|
||||
except Exception as e: # leave state untouched so next run retries
|
||||
print(f"remind_send: delivery failed for {text!r}: {e}", file=sys.stderr)
|
||||
if last:
|
||||
fresh[key] = last
|
||||
continue
|
||||
|
||||
ts = datetime.now(TZ).replace(tzinfo=None).isoformat(timespec="seconds")
|
||||
LOG_DIR.mkdir(parents=True, exist_ok=True)
|
||||
with LOG_FILE.open("a", encoding="utf-8") as f:
|
||||
f.write(f"{ts} {text}\n")
|
||||
fresh[key] = fire_iso
|
||||
|
||||
if fresh != state:
|
||||
STATE_FILE.write_text(json.dumps(fresh, indent=2, ensure_ascii=False), encoding="utf-8")
|
||||
_record_fire(conn, rid, sid, schedule_type, ft, "delivered")
|
||||
log_operation("DELIVER", rid, f'text="{text}"')
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
Reference in New Issue
Block a user