forked from retoor/devplacepy
feat: add online presence tracking with last_seen column and configurable timeout
Add `last_seen` column to users table with index, implement `set_last_seen` and `get_online_users` database functions, expose presence config env vars (`PRESENCE_TIMEOUT_SECONDS`, `PRESENCE_ONLINE_LIMIT`, `PRESENCE_ONLINE_MARGIN_SECONDS`), include `last_seen` in follow list responses, and update profile docs to mention online indicator.
This commit is contained in:
@@ -2,8 +2,7 @@
|
||||
|
||||
import logging
|
||||
from collections import OrderedDict
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Optional
|
||||
from typing import Any
|
||||
|
||||
from fastapi import WebSocket
|
||||
|
||||
@@ -15,14 +14,12 @@ DELIVERED_CAP = 4000
|
||||
class ConnectionManager:
|
||||
def __init__(self) -> None:
|
||||
self._connections: dict[str, set[WebSocket]] = {}
|
||||
self._last_seen: dict[str, str] = {}
|
||||
self._delivered: "OrderedDict[str, bool]" = OrderedDict()
|
||||
|
||||
def register(self, user_uid: str, websocket: WebSocket) -> bool:
|
||||
sockets = self._connections.setdefault(user_uid, set())
|
||||
was_offline = len(sockets) == 0
|
||||
sockets.add(websocket)
|
||||
self._last_seen.pop(user_uid, None)
|
||||
logger.info(
|
||||
"messaging socket registered for %s (sockets=%d)", user_uid, len(sockets)
|
||||
)
|
||||
@@ -35,7 +32,6 @@ class ConnectionManager:
|
||||
sockets.discard(websocket)
|
||||
if not sockets:
|
||||
self._connections.pop(user_uid, None)
|
||||
self._last_seen[user_uid] = datetime.now(timezone.utc).isoformat()
|
||||
logger.info("messaging socket removed for %s (now offline)", user_uid)
|
||||
return True
|
||||
logger.debug(
|
||||
@@ -43,9 +39,6 @@ class ConnectionManager:
|
||||
)
|
||||
return False
|
||||
|
||||
def is_online(self, user_uid: str) -> bool:
|
||||
return bool(self._connections.get(user_uid))
|
||||
|
||||
def has_connections(self) -> bool:
|
||||
return bool(self._connections)
|
||||
|
||||
@@ -61,9 +54,6 @@ class ConnectionManager:
|
||||
def was_delivered(self, message_uid: str) -> bool:
|
||||
return message_uid in self._delivered
|
||||
|
||||
def last_seen(self, user_uid: str) -> Optional[str]:
|
||||
return self._last_seen.get(user_uid)
|
||||
|
||||
def sockets_for(self, user_uid: str) -> list[WebSocket]:
|
||||
return list(self._connections.get(user_uid, ()))
|
||||
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Optional
|
||||
|
||||
from devplacepy.config import (
|
||||
PRESENCE_ONLINE_LIMIT,
|
||||
PRESENCE_ONLINE_MARGIN_SECONDS,
|
||||
PRESENCE_TIMEOUT_SECONDS,
|
||||
PRESENCE_WRITE_SECONDS,
|
||||
)
|
||||
from devplacepy.database import get_online_users, set_last_seen
|
||||
|
||||
_last_write: dict[str, float] = {}
|
||||
|
||||
|
||||
def touch(user_uid: str) -> None:
|
||||
if not user_uid:
|
||||
return
|
||||
now = time.monotonic()
|
||||
if now - _last_write.get(user_uid, 0.0) < PRESENCE_WRITE_SECONDS:
|
||||
return
|
||||
_last_write[user_uid] = now
|
||||
set_last_seen(user_uid, datetime.now(timezone.utc).isoformat())
|
||||
|
||||
|
||||
def seconds_since(last_seen: Optional[str]) -> Optional[float]:
|
||||
if not last_seen:
|
||||
return None
|
||||
try:
|
||||
seen = datetime.fromisoformat(last_seen)
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
if seen.tzinfo is None:
|
||||
seen = seen.replace(tzinfo=timezone.utc)
|
||||
return (datetime.now(timezone.utc) - seen).total_seconds()
|
||||
|
||||
|
||||
def is_online(user: Optional[dict]) -> bool:
|
||||
if not user:
|
||||
return False
|
||||
elapsed = seconds_since(user.get("last_seen"))
|
||||
return elapsed is not None and elapsed < PRESENCE_TIMEOUT_SECONDS
|
||||
|
||||
|
||||
def stays_online(elapsed: Optional[float], was_online: bool) -> bool:
|
||||
if elapsed is None:
|
||||
return False
|
||||
grace = PRESENCE_TIMEOUT_SECONDS + PRESENCE_ONLINE_MARGIN_SECONDS
|
||||
return elapsed < (grace if was_online else PRESENCE_TIMEOUT_SECONDS)
|
||||
|
||||
|
||||
def _cutoff_iso(seconds: int) -> str:
|
||||
return (datetime.now(timezone.utc) - timedelta(seconds=seconds)).isoformat()
|
||||
|
||||
|
||||
def online_cutoff_iso() -> str:
|
||||
return _cutoff_iso(PRESENCE_TIMEOUT_SECONDS)
|
||||
|
||||
|
||||
def sort_by_username(rows: list) -> list:
|
||||
return sorted(rows, key=lambda row: (row.get("username") or "").lower())
|
||||
|
||||
|
||||
def online_users(limit: int = PRESENCE_ONLINE_LIMIT) -> list:
|
||||
return sort_by_username(get_online_users(online_cutoff_iso(), limit))
|
||||
|
||||
|
||||
def online_candidates(limit: int = PRESENCE_ONLINE_LIMIT) -> list:
|
||||
return sort_by_username(
|
||||
get_online_users(
|
||||
_cutoff_iso(PRESENCE_TIMEOUT_SECONDS + PRESENCE_ONLINE_MARGIN_SECONDS), limit
|
||||
)
|
||||
)
|
||||
@@ -0,0 +1,118 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
from typing import Optional
|
||||
|
||||
from devplacepy.config import PRESENCE_ONLINE_LIMIT
|
||||
from devplacepy.database import db, get_users_by_uids
|
||||
from devplacepy.services import presence
|
||||
from devplacepy.services.base import BaseService
|
||||
from devplacepy.services.pubsub import publish as pubsub_publish
|
||||
from devplacepy.services.pubsub.hub import pubsub
|
||||
|
||||
ROSTER_TOPIC = "public.presence.roster"
|
||||
TOPIC_PATTERN = re.compile(r"^public\.presence\.(?P<uid>[A-Za-z0-9_-]{1,128})$")
|
||||
|
||||
|
||||
class PresenceRelayService(BaseService):
|
||||
title = "Presence relay"
|
||||
description = (
|
||||
"Single source of truth for LIVE online status. Each tick it recomputes one online "
|
||||
"set with hysteresis - a user is online at the timeout window but only drops after an "
|
||||
"extra grace margin - and drives BOTH the per-user avatar dots (public.presence.{uid}) "
|
||||
"and the shared feed roster (public.presence.roster) from that one set, so they never "
|
||||
"disagree. It publishes ONLY on change (a real online<->offline transition or a first-seen "
|
||||
"subscriber), reads due users in one batched query per tick, and does nothing when idle. "
|
||||
"Runs on the service lock owner where every subscriber converges."
|
||||
)
|
||||
default_enabled = True
|
||||
|
||||
def __init__(self):
|
||||
super().__init__(name="presence_relay", interval_seconds=2)
|
||||
self._online: set[str] = set()
|
||||
self._published: dict[str, bool] = {}
|
||||
self._roster_uids: Optional[frozenset] = None
|
||||
|
||||
async def run_once(self) -> None:
|
||||
if "users" not in db.tables:
|
||||
return
|
||||
dot_pairs: list[tuple[str, str]] = []
|
||||
roster_subscribed = False
|
||||
for entry in pubsub.topics():
|
||||
topic = entry["topic"]
|
||||
if not entry["subscribers"] or "*" in topic:
|
||||
continue
|
||||
if topic == ROSTER_TOPIC:
|
||||
roster_subscribed = True
|
||||
continue
|
||||
match = TOPIC_PATTERN.match(topic)
|
||||
if match is not None:
|
||||
dot_pairs.append((topic, match.group("uid")))
|
||||
|
||||
roster_rows = presence.online_candidates() if roster_subscribed else []
|
||||
rows: dict[str, dict] = {}
|
||||
dot_uids = {uid for _, uid in dot_pairs}
|
||||
if dot_uids:
|
||||
rows.update(get_users_by_uids(list(dot_uids)))
|
||||
for row in roster_rows:
|
||||
rows[row["uid"]] = row
|
||||
|
||||
prev = self._online
|
||||
online: set[str] = set()
|
||||
for uid, row in rows.items():
|
||||
elapsed = presence.seconds_since(row.get("last_seen"))
|
||||
if presence.stays_online(elapsed, uid in prev):
|
||||
online.add(uid)
|
||||
self._online = online
|
||||
|
||||
await self._publish_dots(dot_pairs, online, rows)
|
||||
await self._publish_roster(roster_subscribed, roster_rows, online)
|
||||
|
||||
async def _publish_dots(self, dot_pairs, online, rows) -> None:
|
||||
active = {topic for topic, _ in dot_pairs}
|
||||
self._published = {t: v for t, v in self._published.items() if t in active}
|
||||
published = 0
|
||||
for topic, uid in dot_pairs:
|
||||
is_on = uid in online
|
||||
if self._published.get(topic) == is_on:
|
||||
continue
|
||||
row = rows.get(uid) or {}
|
||||
published += await pubsub_publish(
|
||||
topic, {"online": is_on, "last_seen": row.get("last_seen")}
|
||||
)
|
||||
self._published[topic] = is_on
|
||||
if published:
|
||||
self.log(f"pushed {published} presence change(s)")
|
||||
|
||||
async def _publish_roster(self, subscribed, roster_rows, online) -> None:
|
||||
if not subscribed:
|
||||
self._roster_uids = None
|
||||
return
|
||||
users = [row for row in roster_rows if row["uid"] in online][:PRESENCE_ONLINE_LIMIT]
|
||||
uids = frozenset(row["uid"] for row in users)
|
||||
if uids == self._roster_uids:
|
||||
return
|
||||
self._roster_uids = uids
|
||||
await pubsub_publish(
|
||||
ROSTER_TOPIC,
|
||||
{
|
||||
"count": len(users),
|
||||
"users": [
|
||||
{
|
||||
"uid": row["uid"],
|
||||
"username": row["username"],
|
||||
"avatar_seed": row.get("avatar_seed") or row["username"],
|
||||
}
|
||||
for row in users
|
||||
],
|
||||
},
|
||||
)
|
||||
self.log(f"roster changed: {len(users)} online")
|
||||
|
||||
def collect_metrics(self) -> dict:
|
||||
return {
|
||||
"online": len(self._online),
|
||||
"tracked_dots": len(self._published),
|
||||
}
|
||||
Reference in New Issue
Block a user