Refactored timers module to use ctx.state and fixed new timers never firing.
CI / Formatting (push) Successful in 13s
CI / Linting (push) Successful in 14s
CI / Tests (Python 3.12) (push) Successful in 27s
CI / Tests (Python 3.13) (push) Successful in 26s
CI / Tests (Python 3.14) (push) Successful in 25s
CI / Type Checking (push) Successful in 24s
CI / Spelling (push) Successful in 14s
CI / Formatting (push) Successful in 13s
CI / Linting (push) Successful in 14s
CI / Tests (Python 3.12) (push) Successful in 27s
CI / Tests (Python 3.13) (push) Successful in 26s
CI / Tests (Python 3.14) (push) Successful in 25s
CI / Type Checking (push) Successful in 24s
CI / Spelling (push) Successful in 14s
This commit is contained in:
@@ -35,7 +35,7 @@ from .commands import (
|
|||||||
settimermessage,
|
settimermessage,
|
||||||
)
|
)
|
||||||
from .routes import timer_list_page
|
from .routes import timer_list_page
|
||||||
from .scheduler import TimerScheduler, clear_scheduler, get_scheduler, set_scheduler
|
from .scheduler import TimerScheduler
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"addtimer",
|
"addtimer",
|
||||||
@@ -91,7 +91,7 @@ async def setup(ctx: ModuleContext) -> None:
|
|||||||
ctx.logger.info(f"Loaded {count} timer(s) from database.")
|
ctx.logger.info(f"Loaded {count} timer(s) from database.")
|
||||||
|
|
||||||
sched.start(ctx)
|
sched.start(ctx)
|
||||||
set_scheduler(sched)
|
ctx.state["scheduler"] = sched
|
||||||
|
|
||||||
|
|
||||||
@on_teardown
|
@on_teardown
|
||||||
@@ -102,5 +102,7 @@ async def teardown(ctx: ModuleContext) -> None:
|
|||||||
|
|
||||||
:param ctx: Module context.
|
:param ctx: Module context.
|
||||||
"""
|
"""
|
||||||
await get_scheduler().stop(ctx)
|
scheduler: TimerScheduler | None = ctx.state.get("scheduler")
|
||||||
clear_scheduler()
|
if scheduler is not None:
|
||||||
|
await scheduler.stop(ctx)
|
||||||
|
ctx.state["scheduler"] = None
|
||||||
|
|||||||
@@ -46,7 +46,7 @@ async def count_chat_message(ctx: EventContext[ChatEvent]) -> None:
|
|||||||
ctx.logger.debug("Skipping chat count for %s: %s.", event.message_id, reason)
|
ctx.logger.debug("Skipping chat count for %s: %s.", event.message_id, reason)
|
||||||
return
|
return
|
||||||
|
|
||||||
counted_ids = get_scheduler().counted_ids
|
counted_ids = get_scheduler(ctx.module).counted_ids
|
||||||
ctx.logger.debug(
|
ctx.logger.debug(
|
||||||
"Counting message %s from %s for %d timer(s).",
|
"Counting message %s from %s for %d timer(s).",
|
||||||
event.message_id,
|
event.message_id,
|
||||||
@@ -73,7 +73,7 @@ async def handle_visibility_update(ctx: EventContext[VisibilityUpdateEvent]) ->
|
|||||||
return
|
return
|
||||||
|
|
||||||
affected = set(event.message_ids)
|
affected = set(event.message_ids)
|
||||||
counted_ids = get_scheduler().counted_ids
|
counted_ids = get_scheduler(ctx.module).counted_ids
|
||||||
|
|
||||||
ctx.logger.debug(
|
ctx.logger.debug(
|
||||||
"Removing %d hidden message(s) from chat counts.",
|
"Removing %d hidden message(s) from chat counts.",
|
||||||
|
|||||||
@@ -155,7 +155,7 @@ async def settimermessage(ctx: CommandContext) -> None:
|
|||||||
ctx.logger.info(f"Timer {display} message updated by {ctx.user.display_name}.")
|
ctx.logger.info(f"Timer {display} message updated by {ctx.user.display_name}.")
|
||||||
await ctx.owncast_client.send_message(f"Message set for {display}.")
|
await ctx.owncast_client.send_message(f"Message set for {display}.")
|
||||||
if row["enabled"]:
|
if row["enabled"]:
|
||||||
get_scheduler().reschedule()
|
get_scheduler(ctx.module).reschedule()
|
||||||
|
|
||||||
|
|
||||||
@on_command("settimerinterval", requires_moderator=True)
|
@on_command("settimerinterval", requires_moderator=True)
|
||||||
@@ -202,7 +202,7 @@ async def settimerinterval(ctx: CommandContext) -> None:
|
|||||||
f"Interval for {display} set to {normalized} ({interval_type})."
|
f"Interval for {display} set to {normalized} ({interval_type})."
|
||||||
)
|
)
|
||||||
if row["enabled"]:
|
if row["enabled"]:
|
||||||
get_scheduler().reschedule()
|
get_scheduler(ctx.module).reschedule()
|
||||||
|
|
||||||
|
|
||||||
@on_command("settimerlines", requires_moderator=True)
|
@on_command("settimerlines", requires_moderator=True)
|
||||||
@@ -249,7 +249,7 @@ async def settimerlines(ctx: CommandContext) -> None:
|
|||||||
f"Minimum chat lines for {display} set to {label}."
|
f"Minimum chat lines for {display} set to {label}."
|
||||||
)
|
)
|
||||||
if row["enabled"]:
|
if row["enabled"]:
|
||||||
get_scheduler().reschedule()
|
get_scheduler(ctx.module).reschedule()
|
||||||
|
|
||||||
|
|
||||||
@on_command("enabletimer", requires_moderator=True)
|
@on_command("enabletimer", requires_moderator=True)
|
||||||
@@ -293,12 +293,12 @@ async def enabletimer(ctx: CommandContext) -> None:
|
|||||||
)
|
)
|
||||||
|
|
||||||
# Start tracking chat lines for this timer.
|
# Start tracking chat lines for this timer.
|
||||||
get_scheduler().counted_ids[row["id"]] = set()
|
get_scheduler(ctx.module).counted_ids[row["id"]] = set()
|
||||||
ctx.logger.debug("Started chat line tracking for timer %s.", display)
|
ctx.logger.debug("Started chat line tracking for timer %s.", display)
|
||||||
|
|
||||||
ctx.logger.info(f"Timer {display} enabled by {ctx.user.display_name}.")
|
ctx.logger.info(f"Timer {display} enabled by {ctx.user.display_name}.")
|
||||||
await ctx.owncast_client.send_message(f"Timer {display} enabled.")
|
await ctx.owncast_client.send_message(f"Timer {display} enabled.")
|
||||||
get_scheduler().reschedule()
|
get_scheduler(ctx.module).reschedule()
|
||||||
|
|
||||||
|
|
||||||
@on_command("disabletimer", requires_moderator=True)
|
@on_command("disabletimer", requires_moderator=True)
|
||||||
@@ -333,12 +333,12 @@ async def disabletimer(ctx: CommandContext) -> None:
|
|||||||
)
|
)
|
||||||
|
|
||||||
# Stop tracking chat lines for this timer.
|
# Stop tracking chat lines for this timer.
|
||||||
get_scheduler().counted_ids.pop(row["id"], None)
|
get_scheduler(ctx.module).counted_ids.pop(row["id"], None)
|
||||||
ctx.logger.debug("Stopped chat line tracking for timer %s.", display)
|
ctx.logger.debug("Stopped chat line tracking for timer %s.", display)
|
||||||
|
|
||||||
ctx.logger.info(f"Timer {display} disabled by {ctx.user.display_name}.")
|
ctx.logger.info(f"Timer {display} disabled by {ctx.user.display_name}.")
|
||||||
await ctx.owncast_client.send_message(f"Timer {display} disabled.")
|
await ctx.owncast_client.send_message(f"Timer {display} disabled.")
|
||||||
get_scheduler().reschedule()
|
get_scheduler(ctx.module).reschedule()
|
||||||
|
|
||||||
|
|
||||||
@on_command("deletetimer", requires_moderator=True)
|
@on_command("deletetimer", requires_moderator=True)
|
||||||
@@ -364,11 +364,11 @@ async def deletetimer(ctx: CommandContext) -> None:
|
|||||||
await ctx.storage.execute("DELETE FROM timers WHERE id = ?", (row["id"],))
|
await ctx.storage.execute("DELETE FROM timers WHERE id = ?", (row["id"],))
|
||||||
|
|
||||||
# Stop tracking chat lines.
|
# Stop tracking chat lines.
|
||||||
get_scheduler().counted_ids.pop(row["id"], None)
|
get_scheduler(ctx.module).counted_ids.pop(row["id"], None)
|
||||||
|
|
||||||
ctx.logger.info(f"Timer {display} deleted by {ctx.user.display_name}.")
|
ctx.logger.info(f"Timer {display} deleted by {ctx.user.display_name}.")
|
||||||
await ctx.owncast_client.send_message(f"Timer {display} deleted.")
|
await ctx.owncast_client.send_message(f"Timer {display} deleted.")
|
||||||
get_scheduler().reschedule()
|
get_scheduler(ctx.module).reschedule()
|
||||||
|
|
||||||
|
|
||||||
@on_command("listtimers", requires_moderator=True, cooldown=15)
|
@on_command("listtimers", requires_moderator=True, cooldown=15)
|
||||||
|
|||||||
@@ -146,17 +146,19 @@ def _next_fire_time(row: aiosqlite.Row, now: datetime) -> datetime | None:
|
|||||||
last_fired = row["last_fired_at"]
|
last_fired = row["last_fired_at"]
|
||||||
last_fired_dt = datetime.fromisoformat(last_fired) if last_fired else None
|
last_fired_dt = datetime.fromisoformat(last_fired) if last_fired else None
|
||||||
|
|
||||||
|
# Never fired before: immediately due.
|
||||||
|
if last_fired_dt is None:
|
||||||
|
return now
|
||||||
|
|
||||||
if row["interval_type"] == IntervalType.SIMPLE:
|
if row["interval_type"] == IntervalType.SIMPLE:
|
||||||
interval_secs = _duration_to_seconds(row["interval_value"])
|
interval_secs = _duration_to_seconds(row["interval_value"])
|
||||||
if interval_secs <= 0:
|
if interval_secs <= 0:
|
||||||
return None
|
return None
|
||||||
anchor = last_fired_dt if last_fired_dt is not None else now
|
return last_fired_dt + timedelta(seconds=interval_secs)
|
||||||
return anchor + timedelta(seconds=interval_secs)
|
|
||||||
|
|
||||||
# Cron timer: wait for the next scheduled tick.
|
# Cron timer: wait for the next scheduled tick.
|
||||||
try:
|
try:
|
||||||
anchor = last_fired_dt if last_fired_dt is not None else now
|
return next(CronSim(row["interval_value"], last_fired_dt))
|
||||||
return next(CronSim(row["interval_value"], anchor))
|
|
||||||
except (ValueError, KeyError, CronSimError):
|
except (ValueError, KeyError, CronSimError):
|
||||||
return None
|
return None
|
||||||
|
|
||||||
@@ -377,31 +379,14 @@ async def _compute_next_delay(ctx: ModuleContext) -> tuple[float | None, str | N
|
|||||||
return soonest_delay, soonest_display
|
return soonest_delay, soonest_display
|
||||||
|
|
||||||
|
|
||||||
# Module-level singleton for the scheduler instance.
|
def get_scheduler(ctx: ModuleContext) -> TimerScheduler:
|
||||||
_scheduler: TimerScheduler | None = None
|
"""Return the scheduler stored in the module context's state.
|
||||||
|
|
||||||
|
|
||||||
def get_scheduler() -> TimerScheduler:
|
|
||||||
"""Return the active scheduler instance.
|
|
||||||
|
|
||||||
|
:param ctx: The module context.
|
||||||
:return: The active TimerScheduler.
|
:return: The active TimerScheduler.
|
||||||
:raises RuntimeError: If the scheduler has not been initialized.
|
:raises RuntimeError: If the scheduler has not been initialized.
|
||||||
"""
|
"""
|
||||||
if _scheduler is None:
|
scheduler = ctx.state.get("scheduler")
|
||||||
|
if not isinstance(scheduler, TimerScheduler):
|
||||||
raise RuntimeError("TimerScheduler is not initialized.")
|
raise RuntimeError("TimerScheduler is not initialized.")
|
||||||
return _scheduler
|
return scheduler
|
||||||
|
|
||||||
|
|
||||||
def set_scheduler(scheduler: TimerScheduler) -> None:
|
|
||||||
"""Set the active scheduler instance.
|
|
||||||
|
|
||||||
:param scheduler: The TimerScheduler to install.
|
|
||||||
"""
|
|
||||||
global _scheduler
|
|
||||||
_scheduler = scheduler
|
|
||||||
|
|
||||||
|
|
||||||
def clear_scheduler() -> None:
|
|
||||||
"""Clear the active scheduler instance."""
|
|
||||||
global _scheduler
|
|
||||||
_scheduler = None
|
|
||||||
|
|||||||
Reference in New Issue
Block a user