|
# retoor <retoor@molodetz.nl>
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
|
|
from devplacepy.database import db
|
|
from devplacepy.services.base import BaseService
|
|
from devplacepy.services.pubsub import publish as pubsub_publish
|
|
|
|
from . import store
|
|
|
|
BATCH_LIMIT = 500
|
|
SWEEP_SECONDS = 60
|
|
|
|
|
|
def battle_topic(war_uid: str) -> str:
|
|
return f"public.battle.{war_uid}"
|
|
|
|
|
|
class OpinionWarService(BaseService):
|
|
title = "Opinion Wars"
|
|
description = (
|
|
"Relays persisted battle events onto the pub/sub bus for live cards, resolves "
|
|
"ended wars that nobody viewed, and sends the fight-cooldown-ready "
|
|
"notifications. Runs on the service lock owner, where every pub/sub "
|
|
"subscriber converges; the event table stays the source of truth."
|
|
)
|
|
default_enabled = True
|
|
|
|
def __init__(self):
|
|
super().__init__(name="opinionwar", interval_seconds=2)
|
|
self._watermark = 0
|
|
self._primed = False
|
|
self._last_sweep = 0.0
|
|
|
|
def _max_id(self) -> int:
|
|
if "opinion_war_events" not in db.tables:
|
|
return 0
|
|
rows = list(db.query("SELECT MAX(id) AS max_id FROM opinion_war_events"))
|
|
value = rows[0]["max_id"] if rows else None
|
|
return int(value or 0)
|
|
|
|
async def _relay_events(self) -> None:
|
|
if "opinion_war_events" not in db.tables:
|
|
return
|
|
if not self._primed:
|
|
self._watermark = self._max_id()
|
|
self._primed = True
|
|
return
|
|
rows = list(
|
|
db.query(
|
|
"SELECT * FROM opinion_war_events WHERE id > :wm "
|
|
"ORDER BY id ASC LIMIT :lim",
|
|
wm=self._watermark,
|
|
lim=BATCH_LIMIT,
|
|
)
|
|
)
|
|
if not rows:
|
|
return
|
|
delivered = 0
|
|
for row in rows:
|
|
delivered += await pubsub_publish(
|
|
battle_topic(row["war_uid"]), store._event_dict(dict(row))
|
|
)
|
|
self._watermark = max(row["id"] for row in rows)
|
|
if delivered:
|
|
self.log(f"relayed {len(rows)} battle event(s), {delivered} live frame(s)")
|
|
|
|
def _sweep(self) -> None:
|
|
now = time.monotonic()
|
|
if now - self._last_sweep < SWEEP_SECONDS:
|
|
return
|
|
self._last_sweep = now
|
|
resolved = store.resolve_due_wars()
|
|
if resolved:
|
|
self.log(f"resolved {resolved} ended war(s)")
|
|
notified = 0
|
|
for fighter in store.cooldown_ready_fighters():
|
|
if store.notify_cooldown_ready(fighter):
|
|
notified += 1
|
|
if notified:
|
|
self.log(f"sent {notified} fight-ready notification(s)")
|
|
|
|
async def run_once(self) -> None:
|
|
await self._relay_events()
|
|
self._sweep()
|
|
|
|
def collect_metrics(self) -> dict:
|
|
return {"watermark": self._watermark, "primed": self._primed}
|