provozni zaloha
This commit is contained in:
@@ -22,8 +22,9 @@ from pathlib import Path
|
||||
from zoneinfo import ZoneInfo
|
||||
|
||||
from croniter import croniter
|
||||
from db import get_db, init_db, log_operation
|
||||
from db import log_operation
|
||||
from random_times import compute_fire_times, random_cfg_from_row
|
||||
import store
|
||||
|
||||
WORKSPACE = Path(__file__).resolve().parent.parent.parent.parent
|
||||
DEFAULT_DB_PATH = WORKSPACE / "db" / "reminders.sqlite"
|
||||
@@ -60,50 +61,17 @@ 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), "schedule_type": "at"} for r in rows]
|
||||
return [{**row, "schedule_type": "at"} for row in store.due_at(conn, since, until)]
|
||||
|
||||
|
||||
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:
|
||||
for row in store.enabled_cron(conn):
|
||||
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:
|
||||
if not store.is_fire_delivered(conn, row["id"], row["schedule_id"], "cron", fire_iso):
|
||||
due.append({
|
||||
"id": row["id"],
|
||||
"text": row["text"],
|
||||
@@ -116,17 +84,8 @@ def _due_cron(conn, now: datetime) -> list[dict]:
|
||||
|
||||
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:
|
||||
for row in store.enabled_random(conn):
|
||||
cfg = random_cfg_from_row(row)
|
||||
try:
|
||||
fires = compute_fire_times(now.date(), row["text"], cfg)
|
||||
@@ -136,15 +95,7 @@ def _due_random(conn, now: datetime) -> list[dict]:
|
||||
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:
|
||||
if not store.is_fire_delivered(conn, row["id"], row["schedule_id"], "random", fire_iso):
|
||||
due.append({
|
||||
"id": row["id"],
|
||||
"text": row["text"],
|
||||
@@ -156,22 +107,12 @@ def _due_random(conn, now: datetime) -> list[dict]:
|
||||
|
||||
|
||||
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_prague().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),
|
||||
)
|
||||
delivered_at = _now_prague().isoformat(timespec="seconds") if status == "delivered" else None
|
||||
store.record_fire(conn, reminder_id, schedule_id, schedule_type, fire_time, status, delivered_at, error)
|
||||
|
||||
|
||||
def main() -> None:
|
||||
if not DB_PATH.exists():
|
||||
init_db(DB_PATH)
|
||||
|
||||
conn = get_db(DB_PATH)
|
||||
try:
|
||||
with store.connection(DB_PATH) as conn:
|
||||
now = _now_prague()
|
||||
due = _due_at(conn, now) + _due_cron(conn, now) + _due_random(conn, now)
|
||||
if not due:
|
||||
@@ -194,8 +135,6 @@ def main() -> None:
|
||||
|
||||
_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