diff --git a/backend/migrations/versions/c4e8a1f7d9b2_custom_emoji.py b/backend/migrations/versions/c4e8a1f7d9b2_custom_emoji.py
new file mode 100644
index 0000000..10515cf
--- /dev/null
+++ b/backend/migrations/versions/c4e8a1f7d9b2_custom_emoji.py
@@ -0,0 +1,42 @@
+"""custom emoji documents
+
+Revision ID: c4e8a1f7d9b2
+Revises: a9c3e7f1d2b4
+Create Date: 2026-05-31 12:00:00.000000
+
+"""
+
+from collections.abc import Sequence
+
+import sqlalchemy as sa
+from alembic import op
+from sqlalchemy.dialects import postgresql
+
+revision: str = "c4e8a1f7d9b2"
+down_revision: str | None = "a9c3e7f1d2b4"
+branch_labels: str | Sequence[str] | None = None
+depends_on: str | Sequence[str] | None = None
+
+
+def upgrade() -> None:
+ op.create_table(
+ "custom_emoji",
+ sa.Column("custom_emoji_id", sa.BigInteger(), nullable=False),
+ sa.Column("storage_key", sa.String(), nullable=True),
+ sa.Column("file_size", sa.BigInteger(), nullable=True),
+ sa.Column("mime", sa.String(), nullable=True),
+ sa.Column("kind", sa.String(), nullable=True),
+ sa.Column("downloaded", sa.Boolean(), nullable=False),
+ sa.Column("raw", postgresql.JSONB(astext_type=sa.Text()), nullable=False),
+ sa.Column(
+ "first_seen_at",
+ sa.DateTime(timezone=True),
+ server_default=sa.text("now()"),
+ nullable=False,
+ ),
+ sa.PrimaryKeyConstraint("custom_emoji_id"),
+ )
+
+
+def downgrade() -> None:
+ op.drop_table("custom_emoji")
diff --git a/backend/migrations/versions/d5f9b2c8e3a1_dialogs.py b/backend/migrations/versions/d5f9b2c8e3a1_dialogs.py
new file mode 100644
index 0000000..7bcc949
--- /dev/null
+++ b/backend/migrations/versions/d5f9b2c8e3a1_dialogs.py
@@ -0,0 +1,36 @@
+"""dialogs
+
+Revision ID: d5f9b2c8e3a1
+Revises: c4e8a1f7d9b2
+Create Date: 2026-05-31 18:00:00.000000
+
+"""
+
+from collections.abc import Sequence
+
+import sqlalchemy as sa
+from alembic import op
+
+revision: str = "d5f9b2c8e3a1"
+down_revision: str | None = "c4e8a1f7d9b2"
+branch_labels: str | Sequence[str] | None = None
+depends_on: str | Sequence[str] | None = None
+
+
+def upgrade() -> None:
+ op.create_table(
+ "dialogs",
+ sa.Column("account_id", sa.BigInteger(), nullable=False),
+ sa.Column("chat_id", sa.BigInteger(), nullable=False),
+ sa.Column(
+ "updated_at",
+ sa.DateTime(timezone=True),
+ server_default=sa.text("now()"),
+ nullable=False,
+ ),
+ sa.PrimaryKeyConstraint("account_id", "chat_id"),
+ )
+
+
+def downgrade() -> None:
+ op.drop_table("dialogs")
diff --git a/backend/src/api/app.py b/backend/src/api/app.py
index 420da78..99a839b 100644
--- a/backend/src/api/app.py
+++ b/backend/src/api/app.py
@@ -11,12 +11,16 @@ from starlette.applications import Starlette
from api.auth import BearerAuthMiddleware
from api.mcp.server import mcp
+from api.realtime import hub
from api.routers import (
accounts,
+ analytics,
annotations,
avatars,
backfill,
chats,
+ custom_emoji,
+ events,
folders,
media,
peers,
@@ -39,7 +43,10 @@ mcp_app = mcp.http_app(path="/")
@asynccontextmanager
async def lifespan(app_: Starlette) -> AsyncGenerator[None]:
+ pool = await container.get(asyncpg.Pool)
+ await hub.start(pool)
yield
+ await hub.stop()
await app_.state.dishka_container.close()
@@ -59,6 +66,7 @@ async def health(pool: FromDishka[asyncpg.Pool]) -> dict[str, bool]:
app.include_router(accounts.router)
+app.include_router(analytics.router)
app.include_router(policy.router)
app.include_router(folders.router)
app.include_router(backfill.router)
@@ -66,8 +74,10 @@ app.include_router(search.router)
app.include_router(chats.router)
app.include_router(media.router)
app.include_router(avatars.router)
+app.include_router(custom_emoji.router)
app.include_router(social.router)
app.include_router(presence.router)
+app.include_router(events.router)
app.include_router(peers.router)
app.include_router(annotations.router)
app.include_router(watches.router)
diff --git a/backend/src/api/realtime.py b/backend/src/api/realtime.py
new file mode 100644
index 0000000..130a7db
--- /dev/null
+++ b/backend/src/api/realtime.py
@@ -0,0 +1,123 @@
+import asyncio
+import json
+import logging
+from typing import Any
+
+import asyncpg
+
+from utils.env import env
+from utils.events import BG_EVENTS_CHANNEL
+from utils.read import chats as chats_read
+from utils.read import presence as presence_read
+
+logger = logging.getLogger(__name__)
+
+QUEUE_MAXSIZE = 256
+
+
+class Subscriber:
+ def __init__(self, account_id: int, chat_id: int | None) -> None:
+ self.account_id = account_id
+ self.chat_id = chat_id
+ self.queue: asyncio.Queue[dict[str, Any]] = asyncio.Queue(maxsize=QUEUE_MAXSIZE)
+
+
+class EventHub:
+ def __init__(self) -> None:
+ self._subscribers: set[Subscriber] = set()
+ self._pool: asyncpg.Pool | None = None
+ self._conn: asyncpg.Connection | None = None
+ self._tasks: set[asyncio.Task] = set()
+
+ def subscribe(self, account_id: int, chat_id: int | None) -> Subscriber:
+ sub = Subscriber(account_id, chat_id)
+ self._subscribers.add(sub)
+ return sub
+
+ def unsubscribe(self, sub: Subscriber) -> None:
+ self._subscribers.discard(sub)
+
+ async def start(self, pool: asyncpg.Pool) -> None:
+ self._pool = pool
+ conn = await asyncpg.connect(dsn=env.db.connection_url)
+ await conn.add_listener(BG_EVENTS_CHANNEL, self._on_notify)
+ self._conn = conn
+ logger.info("Realtime hub listening on %s", BG_EVENTS_CHANNEL)
+
+ async def stop(self) -> None:
+ for task in self._tasks:
+ task.cancel()
+ if self._conn is not None:
+ await self._conn.close()
+ self._conn = None
+
+ def _on_notify(
+ self, _conn: asyncpg.Connection, _pid: int, _channel: str, payload: str
+ ) -> None:
+ task = asyncio.create_task(self._dispatch(payload))
+ self._tasks.add(task)
+ task.add_done_callback(self._tasks.discard)
+
+ async def _dispatch(self, payload: str) -> None:
+ try:
+ event = json.loads(payload)
+ except json.JSONDecodeError:
+ return
+ account_id = event.get("account_id")
+ chat_id = event.get("chat_id")
+ targets = [
+ sub
+ for sub in self._subscribers
+ if sub.account_id == account_id
+ and (sub.chat_id is None or sub.chat_id == chat_id)
+ ]
+ if not targets:
+ return
+ frame = await self._build_frame(event)
+ if frame is None:
+ return
+ for sub in targets:
+ try:
+ sub.queue.put_nowait(frame)
+ except asyncio.QueueFull:
+ logger.warning("Dropping event for slow subscriber")
+
+ async def _build_frame( # noqa: PLR0911
+ self, event: dict[str, Any]
+ ) -> dict[str, Any] | None:
+ if self._pool is None:
+ return None
+ kind = event.get("kind")
+ account_id = event["account_id"]
+ if kind in {"message", "edit", "reaction"}:
+ view = await chats_read.get_message(
+ self._pool, account_id, event["chat_id"], event["message_id"]
+ )
+ if view is None:
+ return None
+ return {"type": kind, "message": view.model_dump(mode="json")}
+ if kind == "delete":
+ return {
+ "type": "delete",
+ "chat_id": event.get("chat_id"),
+ "message_ids": event.get("message_ids", []),
+ }
+ if kind == "presence":
+ sample = await presence_read.current_presence(
+ self._pool, account_id, event["chat_id"]
+ )
+ return {
+ "type": "presence",
+ "peer_id": event["chat_id"],
+ "sample": sample.model_dump(mode="json") if sample else None,
+ }
+ if kind == "receipt":
+ return {
+ "type": "receipt",
+ "chat_id": event["chat_id"],
+ "read_up_to": event["message_id"],
+ }
+ return None
+
+
+hub = EventHub()
diff --git a/backend/src/api/routers/analytics.py b/backend/src/api/routers/analytics.py
new file mode 100644
index 0000000..7cdd57c
--- /dev/null
+++ b/backend/src/api/routers/analytics.py
@@ -0,0 +1,30 @@
+from typing import Annotated
+
+import asyncpg
+from dishka.integrations.fastapi import DishkaRoute, FromDishka
+from fastapi import APIRouter, Query
+
+from utils.read import analytics
+from utils.read.models import ResponseStats, VolumeBucket
+
+router = APIRouter(prefix="/api/analytics", tags=["analytics"], route_class=DishkaRoute)
+
+AccountId = Annotated[int, Query()]
+ChatId = Annotated[int, Query()]
+
+
+@router.get("/volume")
+async def volume(
+ pool: FromDishka[asyncpg.Pool],
+ account_id: AccountId,
+ chat_id: ChatId,
+ days: Annotated[int, Query()] = 90,
+) -> list[VolumeBucket]:
+ return await analytics.message_volume(pool, account_id, chat_id, days=days)
+
+
+@router.get("/response-time")
+async def response_time(
+ pool: FromDishka[asyncpg.Pool], account_id: AccountId, chat_id: ChatId
+) -> ResponseStats:
+ return await analytics.response_stats(pool, account_id, chat_id)
diff --git a/backend/src/api/routers/backfill.py b/backend/src/api/routers/backfill.py
index c49f6e2..26e49e3 100644
--- a/backend/src/api/routers/backfill.py
+++ b/backend/src/api/routers/backfill.py
@@ -24,6 +24,10 @@ class FetchMediaRequest(BaseModel):
message_id: int
+class SyncDialogsRequest(BaseModel):
+ account_id: int
+
+
class EnqueueResponse(BaseModel):
job_id: int
@@ -78,6 +82,14 @@ async def enqueue_fetch_media(
return EnqueueResponse(job_id=job_id)
+@router.post("/dialogs/sync", status_code=201)
+async def enqueue_sync_dialogs(
+ pool: FromDishka[asyncpg.Pool], body: SyncDialogsRequest
+) -> EnqueueResponse:
+ job_id = await enqueue(pool, body.account_id, "sync_dialogs", {})
+ return EnqueueResponse(job_id=job_id)
+
+
@router.get("/jobs")
async def list_jobs(
pool: FromDishka[asyncpg.Pool],
diff --git a/backend/src/api/routers/chats.py b/backend/src/api/routers/chats.py
index 996f774..364b604 100644
--- a/backend/src/api/routers/chats.py
+++ b/backend/src/api/routers/chats.py
@@ -13,7 +13,9 @@ from utils.read.models import (
MessageVersionView,
MessageView,
Page,
+ PinnedView,
)
+from utils.read.pinned import get_pinned
router = APIRouter(prefix="/api", tags=["chats"], route_class=DishkaRoute)
@@ -55,6 +57,13 @@ async def chat_history(
)
+@router.get("/chats/{chat_id}/pinned")
+async def chat_pinned(
+ pool: FromDishka[asyncpg.Pool], chat_id: int, account_id: AccountId
+) -> PinnedView | None:
+ return await get_pinned(pool, account_id, chat_id)
+
+
@router.post("/chats/{chat_id}/enrich")
async def enrich_chat(
pool: FromDishka[asyncpg.Pool], chat_id: int, body: EnrichRequest
diff --git a/backend/src/api/routers/custom_emoji.py b/backend/src/api/routers/custom_emoji.py
new file mode 100644
index 0000000..456c9b5
--- /dev/null
+++ b/backend/src/api/routers/custom_emoji.py
@@ -0,0 +1,35 @@
+from typing import Annotated
+
+import asyncpg
+from dishka.integrations.fastapi import DishkaRoute, FromDishka
+from fastapi import APIRouter, HTTPException, Query
+from fastapi.responses import FileResponse
+
+from utils.jobs import enqueue
+from utils.read.custom_emoji import current_custom_emoji
+from utils.storage import ContentAddressedStorage
+
+router = APIRouter(
+ prefix="/api/custom-emoji", tags=["custom-emoji"], route_class=DishkaRoute
+)
+
+
+@router.get("/{custom_emoji_id}")
+async def serve_custom_emoji(
+ pool: FromDishka[asyncpg.Pool],
+ storage: FromDishka[ContentAddressedStorage],
+ custom_emoji_id: int,
+ account_id: Annotated[int, Query()],
+) -> FileResponse:
+ emoji = await current_custom_emoji(pool, custom_emoji_id)
+ if emoji is None or not emoji.downloaded or emoji.storage_key is None:
+ await enqueue(
+ pool, account_id, "fetch_custom_emoji", {"custom_emoji_id": custom_emoji_id}
+ )
+ raise HTTPException(
+ status_code=409, detail="custom emoji not downloaded; fetching"
+ )
+ return FileResponse(
+ storage.url(emoji.storage_key),
+ media_type=emoji.mime or "application/octet-stream",
+ )
diff --git a/backend/src/api/routers/events.py b/backend/src/api/routers/events.py
new file mode 100644
index 0000000..35cdcba
--- /dev/null
+++ b/backend/src/api/routers/events.py
@@ -0,0 +1,42 @@
+import asyncio
+import json
+from collections.abc import AsyncGenerator
+from typing import Annotated
+
+from fastapi import APIRouter, Query
+from fastapi.responses import StreamingResponse
+
+from api.realtime import Subscriber, hub
+
+router = APIRouter(prefix="/api/events", tags=["events"])
+
+HEARTBEAT_SECONDS = 15
+
+AccountId = Annotated[int, Query()]
+ChatId = Annotated[int | None, Query()]
+
+
+async def _stream(sub: Subscriber) -> AsyncGenerator[str]:
+ try:
+ yield ": connected\n\n"
+ while True:
+ try:
+ frame = await asyncio.wait_for(
+ sub.queue.get(), timeout=HEARTBEAT_SECONDS
+ )
+ except TimeoutError:
+ yield ": keepalive\n\n"
+ continue
+ yield f"event: {frame['type']}\ndata: {json.dumps(frame)}\n\n"
+ finally:
+ hub.unsubscribe(sub)
+
+
+@router.get("")
+async def events(account_id: AccountId, chat_id: ChatId = None) -> StreamingResponse:
+ sub = hub.subscribe(account_id, chat_id)
+ return StreamingResponse(
+ _stream(sub),
+ media_type="text/event-stream",
+ headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
+ )
diff --git a/backend/src/api/routers/presence.py b/backend/src/api/routers/presence.py
index 38c2ed6..67d3b8a 100644
--- a/backend/src/api/routers/presence.py
+++ b/backend/src/api/routers/presence.py
@@ -36,6 +36,13 @@ async def presence_history(
)
+@router.get("/current")
+async def current_presence(
+ pool: FromDishka[asyncpg.Pool], account_id: AccountId, peer_id: PeerId
+) -> PresenceSample | None:
+ return await presence.current_presence(pool, account_id, peer_id)
+
+
@router.get("/hourly")
async def presence_hourly(
pool: FromDishka[asyncpg.Pool],
diff --git a/backend/src/userbot/handlers/deletes.py b/backend/src/userbot/handlers/deletes.py
index 4d9dd60..6927914 100644
--- a/backend/src/userbot/handlers/deletes.py
+++ b/backend/src/userbot/handlers/deletes.py
@@ -3,6 +3,7 @@ from pyrogram.types import Message
from userbot import PyroClient
from userbot.modules.capture import repository
from userbot.modules.capture.repository import CHANNEL_ID_THRESHOLD
+from utils.events import notify_bg_event
@PyroClient.on_deleted_messages()
@@ -21,8 +22,12 @@ async def on_deleted_messages(client: PyroClient, messages: list[Message]) -> No
channels.setdefault(chat_id, []).append(message.id)
if box:
await repository.mark_deleted_box(ctx.pool, ctx.account_id, box)
+ await notify_bg_event(ctx.pool, "delete", ctx.account_id, message_ids=box)
for chat_id, ids in channels.items():
await repository.mark_deleted_channel(ctx.pool, ctx.account_id, chat_id, ids)
+ await notify_bg_event(
+ ctx.pool, "delete", ctx.account_id, chat_id=chat_id, message_ids=ids
+ )
handlers = on_deleted_messages.handlers
diff --git a/backend/src/userbot/handlers/edits.py b/backend/src/userbot/handlers/edits.py
index 792041a..bf5bab1 100644
--- a/backend/src/userbot/handlers/edits.py
+++ b/backend/src/userbot/handlers/edits.py
@@ -5,6 +5,7 @@ from userbot.modules.capture import repository
from userbot.modules.capture.chat_meta import meta_from_chat
from userbot.modules.capture.message import sender_id
from userbot.modules.media import capture_media, media_unique_id, self_destruct_ttl
+from utils.events import notify_bg_event
@PyroClient.on_edited_message()
@@ -32,8 +33,13 @@ async def on_edited_message(client: PyroClient, message: Message) -> None:
has_media=message.media is not None,
is_self_destruct=self_destruct_ttl(message) is not None,
)
- if changed and message.media is not None:
+ if not changed:
+ return
+ if message.media is not None:
await capture_media(client, message, ctx, chat_id, message.id, toggles)
+ await notify_bg_event(
+ ctx.pool, "edit", ctx.account_id, chat_id=chat_id, message_id=message.id
+ )
handlers = on_edited_message.handlers
diff --git a/backend/src/userbot/handlers/messages.py b/backend/src/userbot/handlers/messages.py
index e5061e0..30e6f9a 100644
--- a/backend/src/userbot/handlers/messages.py
+++ b/backend/src/userbot/handlers/messages.py
@@ -5,6 +5,7 @@ from userbot.modules.capture import capture_message
from userbot.modules.capture.chat_meta import meta_from_chat
from userbot.modules.stt import is_transcribable
from userbot.modules.stt.gate import safe_transcribe
+from utils.events import notify_bg_event
@PyroClient.on_message()
@@ -18,6 +19,9 @@ async def on_message(client: PyroClient, message: Message) -> None:
if not toggles.messages:
return
await capture_message(client, message, ctx, toggles)
+ await notify_bg_event(
+ ctx.pool, "message", ctx.account_id, chat_id=meta.chat_id, message_id=message.id
+ )
if (
toggles.stt
and is_transcribable(message)
diff --git a/backend/src/userbot/handlers/presence.py b/backend/src/userbot/handlers/presence.py
index 34c4ffa..0fd681a 100644
--- a/backend/src/userbot/handlers/presence.py
+++ b/backend/src/userbot/handlers/presence.py
@@ -4,6 +4,7 @@ from pyrogram.types import User
from userbot import PyroClient
from userbot.modules.presence import repository
+from utils.events import notify_bg_event
@PyroClient.on_user_status()
@@ -22,6 +23,7 @@ async def on_user_status(client: PyroClient, user: User) -> None:
str(user.raw),
)
await ctx.watches.on_status(user.id, is_online=user.status.name.lower() == "online")
+ await notify_bg_event(ctx.pool, "presence", ctx.account_id, chat_id=user.id)
handlers = on_user_status.handlers
diff --git a/backend/src/userbot/handlers/raw/reactions.py b/backend/src/userbot/handlers/raw/reactions.py
index 7bda931..d0a209c 100644
--- a/backend/src/userbot/handlers/raw/reactions.py
+++ b/backend/src/userbot/handlers/raw/reactions.py
@@ -3,6 +3,7 @@ from pyrogram import raw, utils
from userbot import PyroClient
from userbot.modules.capture import repository
from userbot.modules.capture.chat_meta import meta_from_peer
+from utils.events import notify_bg_event
HANDLES = (raw.types.UpdateMessageReactions,)
@@ -47,3 +48,10 @@ async def handle(
await repository.sync_reactions(
ctx.pool, ctx.account_id, meta.chat_id, update.msg_id, current
)
+ await notify_bg_event(
+ ctx.pool,
+ "reaction",
+ ctx.account_id,
+ chat_id=meta.chat_id,
+ message_id=update.msg_id,
+ )
diff --git a/backend/src/userbot/handlers/raw/read_receipts.py b/backend/src/userbot/handlers/raw/read_receipts.py
index 5990e20..d2619c9 100644
--- a/backend/src/userbot/handlers/raw/read_receipts.py
+++ b/backend/src/userbot/handlers/raw/read_receipts.py
@@ -4,6 +4,7 @@ from pyrogram import raw, utils
from userbot import PyroClient
from userbot.modules.read_receipts import repository
+from utils.events import notify_bg_event
HANDLES = (raw.types.UpdateReadHistoryOutbox,)
@@ -24,3 +25,6 @@ async def handle(
update.max_id,
str(update),
)
+ await notify_bg_event(
+ ctx.pool, "receipt", ctx.account_id, chat_id=chat_id, message_id=update.max_id
+ )
diff --git a/backend/src/userbot/modules/custom_emoji/__init__.py b/backend/src/userbot/modules/custom_emoji/__init__.py
new file mode 100644
index 0000000..e69de29
diff --git a/backend/src/userbot/modules/custom_emoji/repository.py b/backend/src/userbot/modules/custom_emoji/repository.py
new file mode 100644
index 0000000..ac3a3a7
--- /dev/null
+++ b/backend/src/userbot/modules/custom_emoji/repository.py
@@ -0,0 +1,34 @@
+import asyncpg
+
+_GET = """
+SELECT downloaded FROM custom_emoji WHERE custom_emoji_id = $1
+"""
+
+_UPSERT_DOWNLOADED = """
+INSERT INTO custom_emoji
+ (custom_emoji_id, storage_key, file_size, mime, kind, downloaded, raw)
+VALUES ($1, $2, $3, $4, $5, true, '{}'::jsonb)
+ON CONFLICT (custom_emoji_id) DO UPDATE SET
+ storage_key = EXCLUDED.storage_key,
+ file_size = EXCLUDED.file_size,
+ mime = EXCLUDED.mime,
+ kind = EXCLUDED.kind,
+ downloaded = true
+"""
+
+
+async def is_downloaded(pool: asyncpg.Pool, custom_emoji_id: int) -> bool:
+ return bool(await pool.fetchval(_GET, custom_emoji_id))
+
+
+async def upsert_downloaded( # noqa: PLR0913
+ pool: asyncpg.Pool,
+ custom_emoji_id: int,
+ storage_key: str,
+ file_size: int | None,
+ mime: str | None,
+ kind: str,
+) -> None:
+ await pool.execute(
+ _UPSERT_DOWNLOADED, custom_emoji_id, storage_key, file_size, mime, kind
+ )
diff --git a/backend/src/userbot/modules/jobs/handlers/__init__.py b/backend/src/userbot/modules/jobs/handlers/__init__.py
index 7aad1cd..40b0c76 100644
--- a/backend/src/userbot/modules/jobs/handlers/__init__.py
+++ b/backend/src/userbot/modules/jobs/handlers/__init__.py
@@ -2,8 +2,18 @@ from userbot.modules.jobs.handlers import (
backfill,
enrich_chat,
fetch_avatar,
+ fetch_custom_emoji,
fetch_media,
+ sync_dialogs,
transcribe,
)
-__all__ = ["backfill", "enrich_chat", "fetch_avatar", "fetch_media", "transcribe"]
+__all__ = [
+ "backfill",
+ "enrich_chat",
+ "fetch_avatar",
+ "fetch_custom_emoji",
+ "fetch_media",
+ "sync_dialogs",
+ "transcribe",
+]
diff --git a/backend/src/userbot/modules/jobs/handlers/backfill.py b/backend/src/userbot/modules/jobs/handlers/backfill.py
index 9aa8b94..1cd27e8 100644
--- a/backend/src/userbot/modules/jobs/handlers/backfill.py
+++ b/backend/src/userbot/modules/jobs/handlers/backfill.py
@@ -1,3 +1,5 @@
+from pyrogram.errors import PeerIdInvalid
+
from userbot.modules.capture import capture_message
from userbot.modules.jobs.context import JobContext
from userbot.modules.jobs.registry import register
@@ -23,11 +25,15 @@ async def backfill(ctx: JobContext) -> None:
max_id = (ctx.job.cursor or {}).get("max_id", 0)
processed = ctx.job.progress.get("processed", 0)
kwargs = {"max_id": max_id} if max_id else {}
- async for message in client.get_chat_history(chat_id, **kwargs):
- await capture_message(client, message, capture, toggles)
- processed += 1
- if processed % SAVE_EVERY == 0:
- next_max = message.id - 1
- await ctx.save_cursor({"max_id": next_max})
- await ctx.report_progress({"processed": processed, "max_id": next_max})
+ try:
+ async for message in client.get_chat_history(chat_id, **kwargs):
+ await capture_message(client, message, capture, toggles)
+ processed += 1
+ if processed % SAVE_EVERY == 0:
+ next_max = message.id - 1
+ await ctx.save_cursor({"max_id": next_max})
+ await ctx.report_progress({"processed": processed, "max_id": next_max})
+ except PeerIdInvalid:
+ await ctx.report_progress({"processed": processed, "error": "peer_id_invalid"})
+ return
await ctx.report_progress({"processed": processed, "done": True})
diff --git a/backend/src/userbot/modules/jobs/handlers/fetch_custom_emoji.py b/backend/src/userbot/modules/jobs/handlers/fetch_custom_emoji.py
new file mode 100644
index 0000000..96c903c
--- /dev/null
+++ b/backend/src/userbot/modules/jobs/handlers/fetch_custom_emoji.py
@@ -0,0 +1,45 @@
+from io import BytesIO
+
+from pyrogram.types import Sticker
+
+from userbot.modules.custom_emoji.repository import is_downloaded, upsert_downloaded
+from userbot.modules.jobs.context import JobContext
+from userbot.modules.jobs.registry import register
+
+
+def _kind(sticker: Sticker) -> str:
+ if sticker.is_animated:
+ return "animated"
+ if sticker.is_video:
+ return "video"
+ return "static"
+
+
+@register("fetch_custom_emoji")
+async def fetch_custom_emoji(ctx: JobContext) -> None:
+ client = ctx.client
+ if client is None:
+ return
+ capture = getattr(client, "capture", None)
+ if capture is None:
+ return
+ custom_emoji_id = int(ctx.job.params["custom_emoji_id"])
+ if await is_downloaded(ctx.pool, custom_emoji_id):
+ return
+ stickers = await client.get_custom_emoji_stickers([str(custom_emoji_id)])
+ if not stickers:
+ return
+ sticker = stickers[0]
+ buffer = await client.download_media(sticker.file_id, in_memory=True)
+ if not isinstance(buffer, BytesIO):
+ return
+ data = buffer.getvalue()
+ storage_key = capture.storage.put(data)
+ await upsert_downloaded(
+ ctx.pool,
+ custom_emoji_id,
+ storage_key,
+ len(data),
+ sticker.mime_type,
+ _kind(sticker),
+ )
diff --git a/backend/src/userbot/modules/jobs/handlers/sync_dialogs.py b/backend/src/userbot/modules/jobs/handlers/sync_dialogs.py
new file mode 100644
index 0000000..746b39a
--- /dev/null
+++ b/backend/src/userbot/modules/jobs/handlers/sync_dialogs.py
@@ -0,0 +1,108 @@
+from datetime import UTC, datetime
+
+from pyrogram import Client
+from pyrogram.errors import BadRequest, Forbidden
+from pyrogram.types import Chat, User
+
+from userbot.modules.avatars import note_avatar
+from userbot.modules.capture.context import CaptureContext
+from userbot.modules.groups.repository import insert_chat_history
+from userbot.modules.jobs.context import JobContext
+from userbot.modules.jobs.registry import register
+from userbot.modules.profiles.parse import snapshot_from_chat, snapshot_from_high_level
+from userbot.modules.profiles.repository import write_profile
+
+SAVE_EVERY = 100
+USERS_BATCH = 200
+
+_UPSERT_DIALOG = """
+INSERT INTO dialogs (account_id, chat_id) VALUES ($1, $2)
+ON CONFLICT (account_id, chat_id) DO UPDATE SET updated_at = now()
+"""
+
+
+async def _save_private(ctx: CaptureContext, chat: Chat, chat_id: int) -> bool:
+ fields, photo_file_id, photo_unique_id = snapshot_from_chat(chat)
+ await write_profile(ctx.pool, ctx.account_id, chat_id, fields, str(chat))
+ if photo_file_id and photo_unique_id:
+ await note_avatar(
+ ctx.pool, ctx.account_id, chat_id, "peer", photo_unique_id, photo_file_id
+ )
+ return bool(fields.first_name or fields.last_name or fields.username)
+
+
+async def _enrich_users(client: Client, ctx: CaptureContext, ids: list[int]) -> None:
+ for start in range(0, len(ids), USERS_BATCH):
+ batch = ids[start : start + USERS_BATCH]
+ try:
+ result = await client.get_users(batch)
+ except (BadRequest, Forbidden):
+ continue
+ users = result if isinstance(result, list) else [result]
+ for user in users:
+ if not isinstance(user, User):
+ continue
+ fields, photo_file_id, photo_unique_id = snapshot_from_high_level(user)
+ await write_profile(ctx.pool, ctx.account_id, user.id, fields, str(user))
+ if photo_file_id and photo_unique_id:
+ await note_avatar(
+ ctx.pool,
+ ctx.account_id,
+ user.id,
+ "peer",
+ photo_unique_id,
+ photo_file_id,
+ )
+
+
+async def _save_group(ctx: CaptureContext, chat: Chat, chat_id: int) -> None:
+ photo = chat.photo
+ photo_unique_id = photo.big_photo_unique_id if photo else None
+ photo_file_id = photo.big_file_id if photo else None
+ await insert_chat_history(
+ ctx.pool,
+ ctx.account_id,
+ chat_id,
+ 0,
+ "meta",
+ chat.title,
+ photo_unique_id,
+ None,
+ datetime.now(UTC),
+ str(chat),
+ )
+ if photo_file_id and photo_unique_id:
+ await note_avatar(
+ ctx.pool, ctx.account_id, chat_id, "chat", photo_unique_id, photo_file_id
+ )
+
+
+@register("sync_dialogs")
+async def sync_dialogs(ctx: JobContext) -> None:
+ client = ctx.client
+ if client is None:
+ return
+ capture = getattr(client, "capture", None)
+ if capture is None:
+ return
+ processed = ctx.job.progress.get("processed", 0)
+ nameless: list[int] = []
+ async for dialog in client.get_dialogs():
+ chat = dialog.chat
+ if chat is None or chat.id is None:
+ continue
+ chat_id = chat.id
+ try:
+ if chat_id > 0:
+ if not await _save_private(capture, chat, chat_id):
+ nameless.append(chat_id)
+ else:
+ await _save_group(capture, chat, chat_id)
+ except (BadRequest, Forbidden):
+ pass
+ await ctx.pool.execute(_UPSERT_DIALOG, ctx.account_id, chat_id)
+ processed += 1
+ if processed % SAVE_EVERY == 0:
+ await ctx.report_progress({"processed": processed})
+ await _enrich_users(client, capture, nameless)
+ await ctx.report_progress({"processed": processed, "done": True})
diff --git a/backend/src/userbot/modules/media/downloader.py b/backend/src/userbot/modules/media/downloader.py
index 095e833..fa70418 100644
--- a/backend/src/userbot/modules/media/downloader.py
+++ b/backend/src/userbot/modules/media/downloader.py
@@ -20,11 +20,20 @@ _MEDIA_ATTRS = (
)
+_WEB_PAGE_ATTRS = ("photo", "video", "animation", "document", "audio")
+
+
def media_object(message: Message) -> tuple[str | None, Any]:
for attr in _MEDIA_ATTRS:
obj = getattr(message, attr, None)
if obj is not None:
return attr, obj
+ web_page = getattr(message, "web_page", None)
+ if web_page is not None:
+ for attr in _WEB_PAGE_ATTRS:
+ obj = getattr(web_page, attr, None)
+ if obj is not None:
+ return attr, obj
return None, None
@@ -70,7 +79,8 @@ async def capture_media( # noqa: PLR0913
file_size = existing["file_size"]
downloaded = True
else:
- buffer = await client.download_media(message, in_memory=True)
+ target = message if getattr(message, kind or "", None) is obj else obj
+ buffer = await client.download_media(target, in_memory=True)
if isinstance(buffer, BytesIO):
data = buffer.getvalue()
storage_key = ctx.storage.put(data)
diff --git a/backend/src/userbot/modules/profiles/parse.py b/backend/src/userbot/modules/profiles/parse.py
index 88b01bb..7625c32 100644
--- a/backend/src/userbot/modules/profiles/parse.py
+++ b/backend/src/userbot/modules/profiles/parse.py
@@ -1,7 +1,7 @@
from dataclasses import dataclass
from pyrogram import Client, raw
-from pyrogram.types import User
+from pyrogram.types import Chat, User
@dataclass(frozen=True)
@@ -54,3 +54,18 @@ def snapshot_from_high_level(
is_deleted_account=bool(user.is_deleted),
)
return fields, photo_file_id, photo_unique_id
+
+
+def snapshot_from_chat(chat: Chat) -> tuple[ProfileFields, str | None, str | None]:
+ photo = chat.photo
+ photo_unique_id = photo.big_photo_unique_id if photo else None
+ photo_file_id = photo.big_file_id if photo else None
+ fields = ProfileFields(
+ first_name=chat.first_name,
+ last_name=chat.last_name,
+ username=chat.username,
+ phone=None,
+ photo_unique_id=photo_unique_id,
+ is_deleted_account=False,
+ )
+ return fields, photo_file_id, photo_unique_id
diff --git a/backend/src/userbot/runner.py b/backend/src/userbot/runner.py
index 793bb8a..ed958f5 100644
--- a/backend/src/userbot/runner.py
+++ b/backend/src/userbot/runner.py
@@ -13,6 +13,7 @@ from userbot import PyroClient
from userbot.modules.capture import CaptureContext, build_capture_context
from userbot.modules.jobs import JobConsumer
from utils.env import env
+from utils.jobs import enqueue
from utils.logging import logger, setup_logging
from utils.read.watches import WATCHES_CHANGED_CHANNEL
from utils.storage import ContentAddressedStorage
@@ -72,6 +73,17 @@ async def _setup_capture(
logger.info("[green]Capture context ready.[/]")
+async def _enqueue_sync_dialogs(pool: asyncpg.Pool, account_id: int) -> None:
+ existing = await pool.fetchval(
+ "SELECT 1 FROM jobs WHERE account_id = $1 AND kind = 'sync_dialogs' "
+ "AND status IN ('pending', 'running') LIMIT 1",
+ account_id,
+ )
+ if existing is None:
+ await enqueue(pool, account_id, "sync_dialogs", {})
+ logger.info("[green]Queued sync_dialogs.[/]")
+
+
async def _listen_changes(
clients: list[PyroClient], tasks: set[asyncio.Task]
) -> asyncpg.Connection:
@@ -134,6 +146,7 @@ async def runner() -> None:
await _setup_capture(pool, client, account_id, storage)
consumer = JobConsumer(client, pool, account_id)
consumer_tasks.append(asyncio.create_task(consumer.run()))
+ await _enqueue_sync_dialogs(pool, account_id)
if clients:
listen_conn = await _listen_changes(clients, reload_tasks)
diff --git a/backend/src/utils/db/models.py b/backend/src/utils/db/models.py
index 678b8a1..e4d7cac 100644
--- a/backend/src/utils/db/models.py
+++ b/backend/src/utils/db/models.py
@@ -477,3 +477,18 @@ class Alert(SQLModel, table=True):
DateTime(timezone=True), nullable=False, server_default=func.now()
)
)
+
+
+class Dialog(SQLModel, table=True):
+ __tablename__ = "dialogs"
+
+ account_id: int = Field(primary_key=True)
+ chat_id: int = Field(primary_key=True)
+ updated_at: datetime = Field(
+ sa_column=Column(
+ DateTime(timezone=True),
+ nullable=False,
+ server_default=func.now(),
+ onupdate=func.now(),
+ )
+ )
diff --git a/backend/src/utils/events.py b/backend/src/utils/events.py
new file mode 100644
index 0000000..816c815
--- /dev/null
+++ b/backend/src/utils/events.py
@@ -0,0 +1,29 @@
+import json
+from typing import Literal
+
+import asyncpg
+
+BG_EVENTS_CHANNEL = "bg_events"
+
+EventKind = Literal["message", "edit", "delete", "reaction", "presence", "receipt"]
+
+
+async def notify_bg_event( # noqa: PLR0913
+ pool: asyncpg.Pool,
+ kind: EventKind,
+ account_id: int,
+ *,
+ chat_id: int | None = None,
+ message_id: int | None = None,
+ message_ids: list[int] | None = None,
+) -> None:
+ payload: dict[str, object] = {"kind": kind, "account_id": account_id}
+ if chat_id is not None:
+ payload["chat_id"] = chat_id
+ if message_id is not None:
+ payload["message_id"] = message_id
+ if message_ids is not None:
+ payload["message_ids"] = message_ids
+ await pool.execute(
+ "SELECT pg_notify($1, $2)", BG_EVENTS_CHANNEL, json.dumps(payload)
+ )
diff --git a/backend/src/utils/read/accounts.py b/backend/src/utils/read/accounts.py
index 6bb6a56..365c12c 100644
--- a/backend/src/utils/read/accounts.py
+++ b/backend/src/utils/read/accounts.py
@@ -3,6 +3,12 @@ import asyncpg
from utils.read.models import AccountView
+async def self_user_id(pool: asyncpg.Pool, account_id: int) -> int | None:
+ return await pool.fetchval(
+ "SELECT tg_user_id FROM accounts WHERE account_id = $1", account_id
+ )
+
+
async def list_accounts(pool: asyncpg.Pool) -> list[AccountView]:
rows = await pool.fetch(
"SELECT account_id, label, phone, tg_user_id, is_active FROM accounts "
diff --git a/backend/src/utils/read/analytics.py b/backend/src/utils/read/analytics.py
new file mode 100644
index 0000000..2a615c4
--- /dev/null
+++ b/backend/src/utils/read/analytics.py
@@ -0,0 +1,69 @@
+from datetime import UTC, datetime, timedelta
+
+import asyncpg
+
+from utils.read.accounts import self_user_id
+from utils.read.models import ResponseStats, VolumeBucket
+
+
+async def message_volume(
+ pool: asyncpg.Pool, account_id: int, chat_id: int, *, days: int = 90
+) -> list[VolumeBucket]:
+ self_id = await self_user_id(pool, account_id)
+ rows = await pool.fetch(
+ "SELECT date_trunc('day', date) AS bucket, count(*) AS total, "
+ "count(*) FILTER (WHERE sender_id = $3) AS outgoing "
+ "FROM messages "
+ "WHERE account_id = $1 AND chat_id = $2 AND date >= $4 "
+ "GROUP BY bucket ORDER BY bucket",
+ account_id,
+ chat_id,
+ self_id,
+ datetime.now(UTC) - timedelta(days=days),
+ )
+ return [
+ VolumeBucket(
+ bucket=row["bucket"],
+ total=row["total"],
+ outgoing=row["outgoing"],
+ incoming=row["total"] - row["outgoing"],
+ )
+ for row in rows
+ ]
+
+
+async def response_stats(
+ pool: asyncpg.Pool, account_id: int, chat_id: int
+) -> ResponseStats:
+ self_id = await self_user_id(pool, account_id)
+ rows = await pool.fetch(
+ "WITH ordered AS ("
+ "SELECT sender_id, date, "
+ "lag(sender_id) OVER w AS prev_sender, "
+ "lag(date) OVER w AS prev_date "
+ "FROM messages "
+ "WHERE account_id = $1 AND chat_id = $2 AND sender_id IS NOT NULL "
+ "WINDOW w AS (ORDER BY date, message_id)), "
+ "resp AS ("
+ "SELECT (sender_id = $3) AS is_mine, "
+ "EXTRACT(EPOCH FROM (date - prev_date)) AS secs "
+ "FROM ordered "
+ "WHERE prev_sender IS NOT NULL AND prev_sender <> sender_id) "
+ "SELECT is_mine, count(*) AS n, "
+ "percentile_cont(0.5) WITHIN GROUP (ORDER BY secs) AS median_secs "
+ "FROM resp GROUP BY is_mine",
+ account_id,
+ chat_id,
+ self_id,
+ )
+ stats = ResponseStats(
+ mine_median_seconds=None, mine_count=0, their_median_seconds=None, their_count=0
+ )
+ for row in rows:
+ if row["is_mine"]:
+ stats.mine_median_seconds = row["median_secs"]
+ stats.mine_count = row["n"]
+ else:
+ stats.their_median_seconds = row["median_secs"]
+ stats.their_count = row["n"]
+ return stats
diff --git a/backend/src/utils/read/chats.py b/backend/src/utils/read/chats.py
index 1c25b40..fc67feb 100644
--- a/backend/src/utils/read/chats.py
+++ b/backend/src/utils/read/chats.py
@@ -1,5 +1,6 @@
import asyncpg
+from utils.read.accounts import self_user_id
from utils.read.message_view import build_message_view, load_raw, media_ref_from
from utils.read.models import (
ChatListItem,
@@ -8,6 +9,7 @@ from utils.read.models import (
MessageView,
Page,
)
+from utils.read.read_receipts import read_up_to
_MESSAGE_COLS = (
"chat_id, message_id, date, sender_id, text, has_media, is_self_destruct, "
@@ -52,35 +54,43 @@ async def list_chats(
pool: asyncpg.Pool, account_id: int, page: Page
) -> list[ChatListItem]:
rows = await pool.fetch(
- "SELECT m.chat_id, count(*) AS message_count, max(m.date) AS last_date, "
+ "WITH ids AS ("
+ "SELECT DISTINCT chat_id FROM messages WHERE account_id = $1 "
+ "UNION SELECT chat_id FROM dialogs WHERE account_id = $1), "
+ "agg AS (SELECT chat_id, count(*) AS message_count, max(date) AS last_date "
+ "FROM messages WHERE account_id = $1 GROUP BY chat_id) "
+ "SELECT ids.chat_id, COALESCE(agg.message_count, 0) AS message_count, "
+ "agg.last_date AS last_date, "
"(SELECT p.first_name FROM peers p "
- "WHERE p.account_id = $1 AND p.peer_id = m.chat_id) AS first_name, "
+ "WHERE p.account_id = $1 AND p.peer_id = ids.chat_id) AS first_name, "
"(SELECT p.last_name FROM peers p "
- "WHERE p.account_id = $1 AND p.peer_id = m.chat_id) AS last_name, "
+ "WHERE p.account_id = $1 AND p.peer_id = ids.chat_id) AS last_name, "
"(SELECT p.username FROM peers p "
- "WHERE p.account_id = $1 AND p.peer_id = m.chat_id) AS username, "
+ "WHERE p.account_id = $1 AND p.peer_id = ids.chat_id) AS username, "
"(SELECT ch.title FROM chat_history ch "
- "WHERE ch.account_id = $1 AND ch.chat_id = m.chat_id "
+ "WHERE ch.account_id = $1 AND ch.chat_id = ids.chat_id "
"AND ch.title IS NOT NULL ORDER BY ch.ts DESC LIMIT 1) AS group_title, "
"EXISTS (SELECT 1 FROM avatars a "
- "WHERE a.account_id = $1 AND a.owner_id = m.chat_id) AS has_avatar, "
- "(SELECT COALESCE((p.raw->>'is_bot')::bool, (p.raw->>'bot')::bool, false) "
- "FROM peers p WHERE p.account_id = $1 AND p.peer_id = m.chat_id) AS is_bot, "
+ "WHERE a.account_id = $1 AND a.owner_id = ids.chat_id) AS has_avatar, "
+ "(SELECT COALESCE((p.raw->>'is_bot')::bool, (p.raw->>'bot')::bool, "
+ "p.raw->>'type' = 'ChatType.BOT', false) "
+ "FROM peers p WHERE p.account_id = $1 AND p.peer_id = ids.chat_id) AS is_bot, "
"(SELECT COALESCE((p.raw->>'is_contact')::bool, (p.raw->>'contact')::bool, "
"false) FROM peers p "
- "WHERE p.account_id = $1 AND p.peer_id = m.chat_id) AS is_contact, "
- "(SELECT ch.raw->'chat'->>'type' = 'ChatType.CHANNEL' FROM chat_history ch "
- "WHERE ch.account_id = $1 AND ch.chat_id = m.chat_id "
- "AND ch.raw->'chat'->>'type' IS NOT NULL "
+ "WHERE p.account_id = $1 AND p.peer_id = ids.chat_id) AS is_contact, "
+ "(SELECT COALESCE(ch.raw->'chat'->>'type', ch.raw->>'type') "
+ "= 'ChatType.CHANNEL' FROM chat_history ch "
+ "WHERE ch.account_id = $1 AND ch.chat_id = ids.chat_id "
+ "AND COALESCE(ch.raw->'chat'->>'type', ch.raw->>'type') IS NOT NULL "
"ORDER BY ch.ts DESC LIMIT 1) AS is_broadcast, "
"(SELECT lm.text FROM messages lm "
- "WHERE lm.account_id = $1 AND lm.chat_id = m.chat_id "
+ "WHERE lm.account_id = $1 AND lm.chat_id = ids.chat_id "
"ORDER BY lm.date DESC, lm.message_id DESC LIMIT 1) AS last_text, "
"(SELECT lm.sender_id FROM messages lm "
- "WHERE lm.account_id = $1 AND lm.chat_id = m.chat_id "
+ "WHERE lm.account_id = $1 AND lm.chat_id = ids.chat_id "
"ORDER BY lm.date DESC, lm.message_id DESC LIMIT 1) AS last_sender_id "
- "FROM messages m WHERE m.account_id = $1 "
- "GROUP BY m.chat_id ORDER BY last_date DESC LIMIT $2 OFFSET $3",
+ "FROM ids LEFT JOIN agg ON agg.chat_id = ids.chat_id "
+ "ORDER BY last_date DESC NULLS LAST, ids.chat_id DESC LIMIT $2 OFFSET $3",
account_id,
page.capped_limit,
page.offset,
@@ -146,9 +156,43 @@ async def get_chat_history(
else:
views.append(_build_album(members, media_by_key))
index = end
+ await _apply_read_status(pool, account_id, chat_id, views)
return views
+async def get_message(
+ pool: asyncpg.Pool, account_id: int, chat_id: int, message_id: int
+) -> MessageView | None:
+ row = await pool.fetchrow(
+ f"SELECT {_MESSAGE_COLS} FROM messages " # noqa: S608
+ "WHERE account_id = $1 AND chat_id = $2 AND message_id = $3",
+ account_id,
+ chat_id,
+ message_id,
+ )
+ if row is None:
+ return None
+ media_by_key = await _media_map(pool, account_id, [row])
+ raw = load_raw(row["raw"])
+ view = build_message_view(row, raw, _single_media(row, raw, media_by_key))
+ await _apply_read_status(pool, account_id, chat_id, [view])
+ return view
+
+
+async def _apply_read_status(
+ pool: asyncpg.Pool, account_id: int, chat_id: int, views: list[MessageView]
+) -> None:
+ self_id = await self_user_id(pool, account_id)
+ if self_id is None:
+ return
+ marker = await read_up_to(pool, account_id, chat_id)
+ if marker is None:
+ return
+ for view in views:
+ if view.sender_id == self_id and view.message_id <= marker:
+ view.read = True
+
+
def _build_album(
members: list[tuple[asyncpg.Record, dict]],
media_by_key: dict[tuple[int, int], asyncpg.Record],
diff --git a/backend/src/utils/read/custom_emoji.py b/backend/src/utils/read/custom_emoji.py
new file mode 100644
index 0000000..313a7bc
--- /dev/null
+++ b/backend/src/utils/read/custom_emoji.py
@@ -0,0 +1,15 @@
+import asyncpg
+
+from utils.read.models import CustomEmojiRef
+
+_GET = """
+SELECT storage_key, downloaded, mime, kind FROM custom_emoji
+WHERE custom_emoji_id = $1
+"""
+
+
+async def current_custom_emoji(
+ pool: asyncpg.Pool, custom_emoji_id: int
+) -> CustomEmojiRef | None:
+ row = await pool.fetchrow(_GET, custom_emoji_id)
+ return CustomEmojiRef(**dict(row)) if row else None
diff --git a/backend/src/utils/read/media.py b/backend/src/utils/read/media.py
index 5aabe2e..101333f 100644
--- a/backend/src/utils/read/media.py
+++ b/backend/src/utils/read/media.py
@@ -1,5 +1,6 @@
import asyncpg
+from utils.read.message_view import load_raw
from utils.read.models import MediaVersionView, MediaView
_MEDIA_COLS = (
@@ -9,6 +10,45 @@ _MEDIA_COLS = (
_VERSION_COLS = "id, kind, storage_key, file_size, mime, observed_at"
+_WEB_PAGE_MEDIA_KINDS = ("photo", "video", "animation", "document", "audio")
+
+
+async def _web_page_media_stub(
+ pool: asyncpg.Pool, account_id: int, chat_id: int, message_id: int
+) -> MediaView | None:
+ row = await pool.fetchrow(
+ "SELECT date, raw FROM messages "
+ "WHERE account_id = $1 AND chat_id = $2 AND message_id = $3",
+ account_id,
+ chat_id,
+ message_id,
+ )
+ if row is None:
+ return None
+ web_page = load_raw(row["raw"]).get("web_page")
+ if not isinstance(web_page, dict):
+ return None
+ kind = next(
+ (k for k in _WEB_PAGE_MEDIA_KINDS if isinstance(web_page.get(k), dict)), None
+ )
+ if kind is None:
+ return None
+ obj = web_page[kind]
+ return MediaView(
+ id=0,
+ account_id=account_id,
+ chat_id=chat_id,
+ message_id=message_id,
+ kind=kind,
+ storage_key=None,
+ file_size=obj.get("file_size"),
+ mime=obj.get("mime_type"),
+ ttl_seconds=None,
+ downloaded=False,
+ extracted_text=None,
+ created_at=row["date"],
+ )
+
async def get_media(pool: asyncpg.Pool, media_id: int) -> MediaView | None:
row = await pool.fetchrow(
@@ -28,7 +68,9 @@ async def get_message_media(
chat_id,
message_id,
)
- return MediaView(**dict(row)) if row else None
+ if row is not None:
+ return MediaView(**dict(row))
+ return await _web_page_media_stub(pool, account_id, chat_id, message_id)
async def get_media_versions(
diff --git a/backend/src/utils/read/models.py b/backend/src/utils/read/models.py
index cd16aa1..fc2c401 100644
--- a/backend/src/utils/read/models.py
+++ b/backend/src/utils/read/models.py
@@ -141,6 +141,13 @@ class ServiceView(BaseModel):
duration: int | None = None
+class PinnedView(BaseModel):
+ message_id: int
+ text: str | None = None
+ media_kind: str | None = None
+ sender_name: str | None = None
+
+
class StickerView(BaseModel):
emoji: str | None = None
set_name: str | None = None
@@ -178,6 +185,7 @@ class MessageView(BaseModel):
sticker: StickerView | None = None
is_sticker: bool = False
is_animated_emoji: bool = False
+ read: bool = False
class MessageVersionView(BaseModel):
@@ -217,6 +225,13 @@ class AvatarRef(BaseModel):
mime: str | None
+class CustomEmojiRef(BaseModel):
+ storage_key: str | None
+ downloaded: bool
+ mime: str | None
+ kind: str | None
+
+
class CallbackView(BaseModel):
position: int
label: str | None
@@ -256,6 +271,20 @@ class PresenceHourly(BaseModel):
last_seen: datetime | None
+class VolumeBucket(BaseModel):
+ bucket: datetime
+ total: int
+ outgoing: int
+ incoming: int
+
+
+class ResponseStats(BaseModel):
+ mine_median_seconds: float | None
+ mine_count: int
+ their_median_seconds: float | None
+ their_count: int
+
+
class PeerView(BaseModel):
peer_id: int
first_name: str | None
diff --git a/backend/src/utils/read/pinned.py b/backend/src/utils/read/pinned.py
new file mode 100644
index 0000000..60f88f4
--- /dev/null
+++ b/backend/src/utils/read/pinned.py
@@ -0,0 +1,32 @@
+import asyncpg
+
+from utils.read.message_view import _media_kind, _peer_name, load_raw
+from utils.read.models import PinnedView
+
+_PINNED_SQL = (
+ "SELECT raw FROM messages "
+ "WHERE account_id = $1 AND chat_id = $2 "
+ "AND raw->>'service' LIKE '%PINNED_MESSAGE%' "
+ "ORDER BY date DESC, message_id DESC LIMIT 1"
+)
+
+
+async def get_pinned(
+ pool: asyncpg.Pool, account_id: int, chat_id: int
+) -> PinnedView | None:
+ row = await pool.fetchrow(_PINNED_SQL, account_id, chat_id)
+ if row is None:
+ return None
+ pinned = load_raw(row["raw"]).get("pinned_message")
+ if not isinstance(pinned, dict):
+ return None
+ message_id = pinned.get("id")
+ if message_id is None:
+ return None
+ sender = pinned.get("from_user")
+ return PinnedView(
+ message_id=message_id,
+ text=pinned.get("text") or pinned.get("caption"),
+ media_kind=_media_kind(pinned),
+ sender_name=_peer_name(sender) if isinstance(sender, dict) else None,
+ )
diff --git a/backend/src/utils/read/presence.py b/backend/src/utils/read/presence.py
index a90be62..5e61233 100644
--- a/backend/src/utils/read/presence.py
+++ b/backend/src/utils/read/presence.py
@@ -33,6 +33,19 @@ async def presence_history( # noqa: PLR0913
return [PresenceSample(**dict(row)) for row in rows]
+async def current_presence(
+ pool: asyncpg.Pool, account_id: int, peer_id: int
+) -> PresenceSample | None:
+ row = await pool.fetchrow(
+ "SELECT peer_id, ts, status, last_online_date, next_offline_date "
+ "FROM presence WHERE account_id = $1 AND peer_id = $2 "
+ "ORDER BY ts DESC LIMIT 1",
+ account_id,
+ peer_id,
+ )
+ return PresenceSample(**dict(row)) if row is not None else None
+
+
async def presence_hourly(
pool: asyncpg.Pool,
account_id: int,
diff --git a/backend/src/utils/read/read_receipts.py b/backend/src/utils/read/read_receipts.py
new file mode 100644
index 0000000..ae9fd38
--- /dev/null
+++ b/backend/src/utils/read/read_receipts.py
@@ -0,0 +1,10 @@
+import asyncpg
+
+
+async def read_up_to(pool: asyncpg.Pool, account_id: int, chat_id: int) -> int | None:
+ return await pool.fetchval(
+ "SELECT max(message_id) FROM read_receipts "
+ "WHERE account_id = $1 AND chat_id = $2 AND kind = 'read'",
+ account_id,
+ chat_id,
+ )
diff --git a/frontend/LICENSE b/frontend/LICENSE
new file mode 100644
index 0000000..8713e28
--- /dev/null
+++ b/frontend/LICENSE
@@ -0,0 +1,689 @@
+Beavergram frontend
+Copyright (C) 2026 Beavergram contributors
+
+This program is free software: you can redistribute it and/or modify
+it under the terms of the GNU General Public License as published by
+the Free Software Foundation, either version 3 of the License, or
+(at your option) any later version.
+
+Portions of the UI (styles, icon font, and component structure) are
+ported from Telegram A (telegram-tt), Copyright (C) Telegram-tt
+contributors, also licensed under GPL-3.0:
+https://github.com/Ajaxy/telegram-tt
+
+----------------------------------------------------------------------
+
+ GNU GENERAL PUBLIC LICENSE
+ Version 3, 29 June 2007
+
+ Copyright (C) 2007 Free Software Foundation, Inc.
+ Everyone is permitted to copy and distribute verbatim copies
+ of this license document, but changing it is not allowed.
+
+ Preamble
+
+ The GNU General Public License is a free, copyleft license for
+software and other kinds of works.
+
+ The licenses for most software and other practical works are designed
+to take away your freedom to share and change the works. By contrast,
+the GNU General Public License is intended to guarantee your freedom to
+share and change all versions of a program--to make sure it remains free
+software for all its users. We, the Free Software Foundation, use the
+GNU General Public License for most of our software; it applies also to
+any other work released this way by its authors. You can apply it to
+your programs, too.
+
+ When we speak of free software, we are referring to freedom, not
+price. Our General Public Licenses are designed to make sure that you
+have the freedom to distribute copies of free software (and charge for
+them if you wish), that you receive source code or can get it if you
+want it, that you can change the software or use pieces of it in new
+free programs, and that you know you can do these things.
+
+ To protect your rights, we need to prevent others from denying you
+these rights or asking you to surrender the rights. Therefore, you have
+certain responsibilities if you distribute copies of the software, or if
+you modify it: responsibilities to respect the freedom of others.
+
+ For example, if you distribute copies of such a program, whether
+gratis or for a fee, you must pass on to the recipients the same
+freedoms that you received. You must make sure that they, too, receive
+or can get the source code. And you must show them these terms so they
+know their rights.
+
+ Developers that use the GNU GPL protect your rights with two steps:
+(1) assert copyright on the software, and (2) offer you this License
+giving you legal permission to copy, distribute and/or modify it.
+
+ For the developers' and authors' protection, the GPL clearly explains
+that there is no warranty for this free software. For both users' and
+authors' sake, the GPL requires that modified versions be marked as
+changed, so that their problems will not be attributed erroneously to
+authors of previous versions.
+
+ Some devices are designed to deny users access to install or run
+modified versions of the software inside them, although the manufacturer
+can do so. This is fundamentally incompatible with the aim of
+protecting users' freedom to change the software. The systematic
+pattern of such abuse occurs in the area of products for individuals to
+use, which is precisely where it is most unacceptable. Therefore, we
+have designed this version of the GPL to prohibit the practice for those
+products. If such problems arise substantially in other domains, we
+stand ready to extend this provision to those domains in future versions
+of the GPL, as needed to protect the freedom of users.
+
+ Finally, every program is threatened constantly by software patents.
+States should not allow patents to restrict development and use of
+software on general-purpose computers, but in those that do, we wish to
+avoid the special danger that patents applied to a free program could
+make it effectively proprietary. To prevent this, the GPL assures that
+patents cannot be used to render the program non-free.
+
+ The precise terms and conditions for copying, distribution and
+modification follow.
+
+ TERMS AND CONDITIONS
+
+ 0. Definitions.
+
+ "This License" refers to version 3 of the GNU General Public License.
+
+ "Copyright" also means copyright-like laws that apply to other kinds of
+works, such as semiconductor masks.
+
+ "The Program" refers to any copyrightable work licensed under this
+License. Each licensee is addressed as "you". "Licensees" and
+"recipients" may be individuals or organizations.
+
+ To "modify" a work means to copy from or adapt all or part of the work
+in a fashion requiring copyright permission, other than the making of an
+exact copy. The resulting work is called a "modified version" of the
+earlier work or a work "based on" the earlier work.
+
+ A "covered work" means either the unmodified Program or a work based
+on the Program.
+
+ To "propagate" a work means to do anything with it that, without
+permission, would make you directly or secondarily liable for
+infringement under applicable copyright law, except executing it on a
+computer or modifying a private copy. Propagation includes copying,
+distribution (with or without modification), making available to the
+public, and in some countries other activities as well.
+
+ To "convey" a work means any kind of propagation that enables other
+parties to make or receive copies. Mere interaction with a user through
+a computer network, with no transfer of a copy, is not conveying.
+
+ An interactive user interface displays "Appropriate Legal Notices"
+to the extent that it includes a convenient and prominently visible
+feature that (1) displays an appropriate copyright notice, and (2)
+tells the user that there is no warranty for the work (except to the
+extent that warranties are provided), that licensees may convey the
+work under this License, and how to view a copy of this License. If
+the interface presents a list of user commands or options, such as a
+menu, a prominent item in the list meets this criterion.
+
+ 1. Source Code.
+
+ The "source code" for a work means the preferred form of the work
+for making modifications to it. "Object code" means any non-source
+form of a work.
+
+ A "Standard Interface" means an interface that either is an official
+standard defined by a recognized standards body, or, in the case of
+interfaces specified for a particular programming language, one that
+is widely used among developers working in that language.
+
+ The "System Libraries" of an executable work include anything, other
+than the work as a whole, that (a) is included in the normal form of
+packaging a Major Component, but which is not part of that Major
+Component, and (b) serves only to enable use of the work with that
+Major Component, or to implement a Standard Interface for which an
+implementation is available to the public in source code form. A
+"Major Component", in this context, means a major essential component
+(kernel, window system, and so on) of the specific operating system
+(if any) on which the executable work runs, or a compiler used to
+produce the work, or an object code interpreter used to run it.
+
+ The "Corresponding Source" for a work in object code form means all
+the source code needed to generate, install, and (for an executable
+work) run the object code and to modify the work, including scripts to
+control those activities. However, it does not include the work's
+System Libraries, or general-purpose tools or generally available free
+programs which are used unmodified in performing those activities but
+which are not part of the work. For example, Corresponding Source
+includes interface definition files associated with source files for
+the work, and the source code for shared libraries and dynamically
+linked subprograms that the work is specifically designed to require,
+such as by intimate data communication or control flow between those
+subprograms and other parts of the work.
+
+ The Corresponding Source need not include anything that users
+can regenerate automatically from other parts of the Corresponding
+Source.
+
+ The Corresponding Source for a work in source code form is that
+same work.
+
+ 2. Basic Permissions.
+
+ All rights granted under this License are granted for the term of
+copyright on the Program, and are irrevocable provided the stated
+conditions are met. This License explicitly affirms your unlimited
+permission to run the unmodified Program. The output from running a
+covered work is covered by this License only if the output, given its
+content, constitutes a covered work. This License acknowledges your
+rights of fair use or other equivalent, as provided by copyright law.
+
+ You may make, run and propagate covered works that you do not
+convey, without conditions so long as your license otherwise remains
+in force. You may convey covered works to others for the sole purpose
+of having them make modifications exclusively for you, or provide you
+with facilities for running those works, provided that you comply with
+the terms of this License in conveying all material for which you do
+not control copyright. Those thus making or running the covered works
+for you must do so exclusively on your behalf, under your direction
+and control, on terms that prohibit them from making any copies of
+your copyrighted material outside their relationship with you.
+
+ Conveying under any other circumstances is permitted solely under
+the conditions stated below. Sublicensing is not allowed; section 10
+makes it unnecessary.
+
+ 3. Protecting Users' Legal Rights From Anti-Circumvention Law.
+
+ No covered work shall be deemed part of an effective technological
+measure under any applicable law fulfilling obligations under article
+11 of the WIPO copyright treaty adopted on 20 December 1996, or
+similar laws prohibiting or restricting circumvention of such
+measures.
+
+ When you convey a covered work, you waive any legal power to forbid
+circumvention of technological measures to the extent such circumvention
+is effected by exercising rights under this License with respect to
+the covered work, and you disclaim any intention to limit operation or
+modification of the work as a means of enforcing, against the work's
+users, your or third parties' legal rights to forbid circumvention of
+technological measures.
+
+ 4. Conveying Verbatim Copies.
+
+ You may convey verbatim copies of the Program's source code as you
+receive it, in any medium, provided that you conspicuously and
+appropriately publish on each copy an appropriate copyright notice;
+keep intact all notices stating that this License and any
+non-permissive terms added in accord with section 7 apply to the code;
+keep intact all notices of the absence of any warranty; and give all
+recipients a copy of this License along with the Program.
+
+ You may charge any price or no price for each copy that you convey,
+and you may offer support or warranty protection for a fee.
+
+ 5. Conveying Modified Source Versions.
+
+ You may convey a work based on the Program, or the modifications to
+produce it from the Program, in the form of source code under the
+terms of section 4, provided that you also meet all of these conditions:
+
+ a) The work must carry prominent notices stating that you modified
+ it, and giving a relevant date.
+
+ b) The work must carry prominent notices stating that it is
+ released under this License and any conditions added under section
+ 7. This requirement modifies the requirement in section 4 to
+ "keep intact all notices".
+
+ c) You must license the entire work, as a whole, under this
+ License to anyone who comes into possession of a copy. This
+ License will therefore apply, along with any applicable section 7
+ additional terms, to the whole of the work, and all its parts,
+ regardless of how they are packaged. This License gives no
+ permission to license the work in any other way, but it does not
+ invalidate such permission if you have separately received it.
+
+ d) If the work has interactive user interfaces, each must display
+ Appropriate Legal Notices; however, if the Program has interactive
+ interfaces that do not display Appropriate Legal Notices, your
+ work need not make them do so.
+
+ A compilation of a covered work with other separate and independent
+works, which are not by their nature extensions of the covered work,
+and which are not combined with it such as to form a larger program,
+in or on a volume of a storage or distribution medium, is called an
+"aggregate" if the compilation and its resulting copyright are not
+used to limit the access or legal rights of the compilation's users
+beyond what the individual works permit. Inclusion of a covered work
+in an aggregate does not cause this License to apply to the other
+parts of the aggregate.
+
+ 6. Conveying Non-Source Forms.
+
+ You may convey a covered work in object code form under the terms
+of sections 4 and 5, provided that you also convey the
+machine-readable Corresponding Source under the terms of this License,
+in one of these ways:
+
+ a) Convey the object code in, or embodied in, a physical product
+ (including a physical distribution medium), accompanied by the
+ Corresponding Source fixed on a durable physical medium
+ customarily used for software interchange.
+
+ b) Convey the object code in, or embodied in, a physical product
+ (including a physical distribution medium), accompanied by a
+ written offer, valid for at least three years and valid for as
+ long as you offer spare parts or customer support for that product
+ model, to give anyone who possesses the object code either (1) a
+ copy of the Corresponding Source for all the software in the
+ product that is covered by this License, on a durable physical
+ medium customarily used for software interchange, for a price no
+ more than your reasonable cost of physically performing this
+ conveying of source, or (2) access to copy the
+ Corresponding Source from a network server at no charge.
+
+ c) Convey individual copies of the object code with a copy of the
+ written offer to provide the Corresponding Source. This
+ alternative is allowed only occasionally and noncommercially, and
+ only if you received the object code with such an offer, in accord
+ with subsection 6b.
+
+ d) Convey the object code by offering access from a designated
+ place (gratis or for a charge), and offer equivalent access to the
+ Corresponding Source in the same way through the same place at no
+ further charge. You need not require recipients to copy the
+ Corresponding Source along with the object code. If the place to
+ copy the object code is a network server, the Corresponding Source
+ may be on a different server (operated by you or a third party)
+ that supports equivalent copying facilities, provided you maintain
+ clear directions next to the object code saying where to find the
+ Corresponding Source. Regardless of what server hosts the
+ Corresponding Source, you remain obligated to ensure that it is
+ available for as long as needed to satisfy these requirements.
+
+ e) Convey the object code using peer-to-peer transmission, provided
+ you inform other peers where the object code and Corresponding
+ Source of the work are being offered to the general public at no
+ charge under subsection 6d.
+
+ A separable portion of the object code, whose source code is excluded
+from the Corresponding Source as a System Library, need not be
+included in conveying the object code work.
+
+ A "User Product" is either (1) a "consumer product", which means any
+tangible personal property which is normally used for personal, family,
+or household purposes, or (2) anything designed or sold for incorporation
+into a dwelling. In determining whether a product is a consumer product,
+doubtful cases shall be resolved in favor of coverage. For a particular
+product received by a particular user, "normally used" refers to a
+typical or common use of that class of product, regardless of the status
+of the particular user or of the way in which the particular user
+actually uses, or expects or is expected to use, the product. A product
+is a consumer product regardless of whether the product has substantial
+commercial, industrial or non-consumer uses, unless such uses represent
+the only significant mode of use of the product.
+
+ "Installation Information" for a User Product means any methods,
+procedures, authorization keys, or other information required to install
+and execute modified versions of a covered work in that User Product from
+a modified version of its Corresponding Source. The information must
+suffice to ensure that the continued functioning of the modified object
+code is in no case prevented or interfered with solely because
+modification has been made.
+
+ If you convey an object code work under this section in, or with, or
+specifically for use in, a User Product, and the conveying occurs as
+part of a transaction in which the right of possession and use of the
+User Product is transferred to the recipient in perpetuity or for a
+fixed term (regardless of how the transaction is characterized), the
+Corresponding Source conveyed under this section must be accompanied
+by the Installation Information. But this requirement does not apply
+if neither you nor any third party retains the ability to install
+modified object code on the User Product (for example, the work has
+been installed in ROM).
+
+ The requirement to provide Installation Information does not include a
+requirement to continue to provide support service, warranty, or updates
+for a work that has been modified or installed by the recipient, or for
+the User Product in which it has been modified or installed. Access to a
+network may be denied when the modification itself materially and
+adversely affects the operation of the network or violates the rules and
+protocols for communication across the network.
+
+ Corresponding Source conveyed, and Installation Information provided,
+in accord with this section must be in a format that is publicly
+documented (and with an implementation available to the public in
+source code form), and must require no special password or key for
+unpacking, reading or copying.
+
+ 7. Additional Terms.
+
+ "Additional permissions" are terms that supplement the terms of this
+License by making exceptions from one or more of its conditions.
+Additional permissions that are applicable to the entire Program shall
+be treated as though they were included in this License, to the extent
+that they are valid under applicable law. If additional permissions
+apply only to part of the Program, that part may be used separately
+under those permissions, but the entire Program remains governed by
+this License without regard to the additional permissions.
+
+ When you convey a copy of a covered work, you may at your option
+remove any additional permissions from that copy, or from any part of
+it. (Additional permissions may be written to require their own
+removal in certain cases when you modify the work.) You may place
+additional permissions on material, added by you to a covered work,
+for which you have or can give appropriate copyright permission.
+
+ Notwithstanding any other provision of this License, for material you
+add to a covered work, you may (if authorized by the copyright holders of
+that material) supplement the terms of this License with terms:
+
+ a) Disclaiming warranty or limiting liability differently from the
+ terms of sections 15 and 16 of this License; or
+
+ b) Requiring preservation of specified reasonable legal notices or
+ author attributions in that material or in the Appropriate Legal
+ Notices displayed by works containing it; or
+
+ c) Prohibiting misrepresentation of the origin of that material, or
+ requiring that modified versions of such material be marked in
+ reasonable ways as different from the original version; or
+
+ d) Limiting the use for publicity purposes of names of licensors or
+ authors of the material; or
+
+ e) Declining to grant rights under trademark law for use of some
+ trade names, trademarks, or service marks; or
+
+ f) Requiring indemnification of licensors and authors of that
+ material by anyone who conveys the material (or modified versions of
+ it) with contractual assumptions of liability to the recipient, for
+ any liability that these contractual assumptions directly impose on
+ those licensors and authors.
+
+ All other non-permissive additional terms are considered "further
+restrictions" within the meaning of section 10. If the Program as you
+received it, or any part of it, contains a notice stating that it is
+governed by this License along with a term that is a further
+restriction, you may remove that term. If a license document contains
+a further restriction but permits relicensing or conveying under this
+License, you may add to a covered work material governed by the terms
+of that license document, provided that the further restriction does
+not survive such relicensing or conveying.
+
+ If you add terms to a covered work in accord with this section, you
+must place, in the relevant source files, a statement of the
+additional terms that apply to those files, or a notice indicating
+where to find the applicable terms.
+
+ Additional terms, permissive or non-permissive, may be stated in the
+form of a separately written license, or stated as exceptions;
+the above requirements apply either way.
+
+ 8. Termination.
+
+ You may not propagate or modify a covered work except as expressly
+provided under this License. Any attempt otherwise to propagate or
+modify it is void, and will automatically terminate your rights under
+this License (including any patent licenses granted under the third
+paragraph of section 11).
+
+ However, if you cease all violation of this License, then your
+license from a particular copyright holder is reinstated (a)
+provisionally, unless and until the copyright holder explicitly and
+finally terminates your license, and (b) permanently, if the copyright
+holder fails to notify you of the violation by some reasonable means
+prior to 60 days after the cessation.
+
+ Moreover, your license from a particular copyright holder is
+reinstated permanently if the copyright holder notifies you of the
+violation by some reasonable means, this is the first time you have
+received notice of violation of this License (for any work) from that
+copyright holder, and you cure the violation prior to 30 days after
+your receipt of the notice.
+
+ Termination of your rights under this section does not terminate the
+licenses of parties who have received copies or rights from you under
+this License. If your rights have been terminated and not permanently
+reinstated, you do not qualify to receive new licenses for the same
+material under section 10.
+
+ 9. Acceptance Not Required for Having Copies.
+
+ You are not required to accept this License in order to receive or
+run a copy of the Program. Ancillary propagation of a covered work
+occurring solely as a consequence of using peer-to-peer transmission
+to receive a copy likewise does not require acceptance. However,
+nothing other than this License grants you permission to propagate or
+modify any covered work. These actions infringe copyright if you do
+not accept this License. Therefore, by modifying or propagating a
+covered work, you indicate your acceptance of this License to do so.
+
+ 10. Automatic Licensing of Downstream Recipients.
+
+ Each time you convey a covered work, the recipient automatically
+receives a license from the original licensors, to run, modify and
+propagate that work, subject to this License. You are not responsible
+for enforcing compliance by third parties with this License.
+
+ An "entity transaction" is a transaction transferring control of an
+organization, or substantially all assets of one, or subdividing an
+organization, or merging organizations. If propagation of a covered
+work results from an entity transaction, each party to that
+transaction who receives a copy of the work also receives whatever
+licenses to the work the party's predecessor in interest had or could
+give under the previous paragraph, plus a right to possession of the
+Corresponding Source of the work from the predecessor in interest, if
+the predecessor has it or can get it with reasonable efforts.
+
+ You may not impose any further restrictions on the exercise of the
+rights granted or affirmed under this License. For example, you may
+not impose a license fee, royalty, or other charge for exercise of
+rights granted under this License, and you may not initiate litigation
+(including a cross-claim or counterclaim in a lawsuit) alleging that
+any patent claim is infringed by making, using, selling, offering for
+sale, or importing the Program or any portion of it.
+
+ 11. Patents.
+
+ A "contributor" is a copyright holder who authorizes use under this
+License of the Program or a work on which the Program is based. The
+work thus licensed is called the contributor's "contributor version".
+
+ A contributor's "essential patent claims" are all patent claims
+owned or controlled by the contributor, whether already acquired or
+hereafter acquired, that would be infringed by some manner, permitted
+by this License, of making, using, or selling its contributor version,
+but do not include claims that would be infringed only as a
+consequence of further modification of the contributor version. For
+purposes of this definition, "control" includes the right to grant
+patent sublicenses in a manner consistent with the requirements of
+this License.
+
+ Each contributor grants you a non-exclusive, worldwide, royalty-free
+patent license under the contributor's essential patent claims, to
+make, use, sell, offer for sale, import and otherwise run, modify and
+propagate the contents of its contributor version.
+
+ In the following three paragraphs, a "patent license" is any express
+agreement or commitment, however denominated, not to enforce a patent
+(such as an express permission to practice a patent or covenant not to
+sue for patent infringement). To "grant" such a patent license to a
+party means to make such an agreement or commitment not to enforce a
+patent against the party.
+
+ If you convey a covered work, knowingly relying on a patent license,
+and the Corresponding Source of the work is not available for anyone
+to copy, free of charge and under the terms of this License, through a
+publicly available network server or other readily accessible means,
+then you must either (1) cause the Corresponding Source to be so
+available, or (2) arrange to deprive yourself of the benefit of the
+patent license for this particular work, or (3) arrange, in a manner
+consistent with the requirements of this License, to extend the patent
+license to downstream recipients. "Knowingly relying" means you have
+actual knowledge that, but for the patent license, your conveying the
+covered work in a country, or your recipient's use of the covered work
+in a country, would infringe one or more identifiable patents in that
+country that you have reason to believe are valid.
+
+ If, pursuant to or in connection with a single transaction or
+arrangement, you convey, or propagate by procuring conveyance of, a
+covered work, and grant a patent license to some of the parties
+receiving the covered work authorizing them to use, propagate, modify
+or convey a specific copy of the covered work, then the patent license
+you grant is automatically extended to all recipients of the covered
+work and works based on it.
+
+ A patent license is "discriminatory" if it does not include within
+the scope of its coverage, prohibits the exercise of, or is
+conditioned on the non-exercise of one or more of the rights that are
+specifically granted under this License. You may not convey a covered
+work if you are a party to an arrangement with a third party that is
+in the business of distributing software, under which you make payment
+to the third party based on the extent of your activity of conveying
+the work, and under which the third party grants, to any of the
+parties who would receive the covered work from you, a discriminatory
+patent license (a) in connection with copies of the covered work
+conveyed by you (or copies made from those copies), or (b) primarily
+for and in connection with specific products or compilations that
+contain the covered work, unless you entered into that arrangement,
+or that patent license was granted, prior to 28 March 2007.
+
+ Nothing in this License shall be construed as excluding or limiting
+any implied license or other defenses to infringement that may
+otherwise be available to you under applicable patent law.
+
+ 12. No Surrender of Others' Freedom.
+
+ If conditions are imposed on you (whether by court order, agreement or
+otherwise) that contradict the conditions of this License, they do not
+excuse you from the conditions of this License. If you cannot convey a
+covered work so as to satisfy simultaneously your obligations under this
+License and any other pertinent obligations, then as a consequence you may
+not convey it at all. For example, if you agree to terms that obligate you
+to collect a royalty for further conveying from those to whom you convey
+the Program, the only way you could satisfy both those terms and this
+License would be to refrain entirely from conveying the Program.
+
+ 13. Use with the GNU Affero General Public License.
+
+ Notwithstanding any other provision of this License, you have
+permission to link or combine any covered work with a work licensed
+under version 3 of the GNU Affero General Public License into a single
+combined work, and to convey the resulting work. The terms of this
+License will continue to apply to the part which is the covered work,
+but the special requirements of the GNU Affero General Public License,
+section 13, concerning interaction through a network will apply to the
+combination as such.
+
+ 14. Revised Versions of this License.
+
+ The Free Software Foundation may publish revised and/or new versions of
+the GNU General Public License from time to time. Such new versions will
+be similar in spirit to the present version, but may differ in detail to
+address new problems or concerns.
+
+ Each version is given a distinguishing version number. If the
+Program specifies that a certain numbered version of the GNU General
+Public License "or any later version" applies to it, you have the
+option of following the terms and conditions either of that numbered
+version or of any later version published by the Free Software
+Foundation. If the Program does not specify a version number of the
+GNU General Public License, you may choose any version ever published
+by the Free Software Foundation.
+
+ If the Program specifies that a proxy can decide which future
+versions of the GNU General Public License can be used, that proxy's
+public statement of acceptance of a version permanently authorizes you
+to choose that version for the Program.
+
+ Later license versions may give you additional or different
+permissions. However, no additional obligations are imposed on any
+author or copyright holder as a result of your choosing to follow a
+later version.
+
+ 15. Disclaimer of Warranty.
+
+ THERE IS NO WARRANTY FOR THE PROGRAM, TO THE EXTENT PERMITTED BY
+APPLICABLE LAW. EXCEPT WHEN OTHERWISE STATED IN WRITING THE COPYRIGHT
+HOLDERS AND/OR OTHER PARTIES PROVIDE THE PROGRAM "AS IS" WITHOUT WARRANTY
+OF ANY KIND, EITHER EXPRESSED OR IMPLIED, INCLUDING, BUT NOT LIMITED TO,
+THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
+PURPOSE. THE ENTIRE RISK AS TO THE QUALITY AND PERFORMANCE OF THE PROGRAM
+IS WITH YOU. SHOULD THE PROGRAM PROVE DEFECTIVE, YOU ASSUME THE COST OF
+ALL NECESSARY SERVICING, REPAIR OR CORRECTION.
+
+ 16. Limitation of Liability.
+
+ IN NO EVENT UNLESS REQUIRED BY APPLICABLE LAW OR AGREED TO IN WRITING
+WILL ANY COPYRIGHT HOLDER, OR ANY OTHER PARTY WHO MODIFIES AND/OR CONVEYS
+THE PROGRAM AS PERMITTED ABOVE, BE LIABLE TO YOU FOR DAMAGES, INCLUDING ANY
+GENERAL, SPECIAL, INCIDENTAL OR CONSEQUENTIAL DAMAGES ARISING OUT OF THE
+USE OR INABILITY TO USE THE PROGRAM (INCLUDING BUT NOT LIMITED TO LOSS OF
+DATA OR DATA BEING RENDERED INACCURATE OR LOSSES SUSTAINED BY YOU OR THIRD
+PARTIES OR A FAILURE OF THE PROGRAM TO OPERATE WITH ANY OTHER PROGRAMS),
+EVEN IF SUCH HOLDER OR OTHER PARTY HAS BEEN ADVISED OF THE POSSIBILITY OF
+SUCH DAMAGES.
+
+ 17. Interpretation of Sections 15 and 16.
+
+ If the disclaimer of warranty and limitation of liability provided
+above cannot be given local legal effect according to their terms,
+reviewing courts shall apply local law that most closely approximates
+an absolute waiver of all civil liability in connection with the
+Program, unless a warranty or assumption of liability accompanies a
+copy of the Program in return for a fee.
+
+ END OF TERMS AND CONDITIONS
+
+ How to Apply These Terms to Your New Programs
+
+ If you develop a new program, and you want it to be of the greatest
+possible use to the public, the best way to achieve this is to make it
+free software which everyone can redistribute and change under these terms.
+
+ To do so, attach the following notices to the program. It is safest
+to attach them to the start of each source file to most effectively
+state the exclusion of warranty; and each file should have at least
+the "copyright" line and a pointer to where the full notice is found.
+
+
+ Copyright (C)
+
+ This program is free software: you can redistribute it and/or modify
+ it under the terms of the GNU General Public License as published by
+ the Free Software Foundation, either version 3 of the License, or
+ (at your option) any later version.
+
+ This program is distributed in the hope that it will be useful,
+ but WITHOUT ANY WARRANTY; without even the implied warranty of
+ MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+ GNU General Public License for more details.
+
+ You should have received a copy of the GNU General Public License
+ along with this program. If not, see .
+
+Also add information on how to contact you by electronic and paper mail.
+
+ If the program does terminal interaction, make it output a short
+notice like this when it starts in an interactive mode:
+
+ Copyright (C)
+ This program comes with ABSOLUTELY NO WARRANTY; for details type `show w'.
+ This is free software, and you are welcome to redistribute it
+ under certain conditions; type `show c' for details.
+
+The hypothetical commands `show w' and `show c' should show the appropriate
+parts of the General Public License. Of course, your program's commands
+might be different; for a GUI interface, you would use an "about box".
+
+ You should also get your employer (if you work as a programmer) or school,
+if any, to sign a "copyright disclaimer" for the program, if necessary.
+For more information on this, and how to apply and follow the GNU GPL, see
+.
+
+ The GNU General Public License does not permit incorporating your program
+into proprietary programs. If your program is a subroutine library, you
+may consider it more useful to permit linking proprietary applications with
+the library. If this is what you want to do, use the GNU Lesser General
+Public License instead of this License. But first, please read
+.
diff --git a/frontend/bun.lock b/frontend/bun.lock
index 243b321..4b63081 100644
--- a/frontend/bun.lock
+++ b/frontend/bun.lock
@@ -6,6 +6,8 @@
"name": "frontend",
"dependencies": {
"bits-ui": "^2.18.1",
+ "lottie-web": "^5.13.0",
+ "pako": "^2.1.0",
},
"devDependencies": {
"@biomejs/biome": "2.4.15",
@@ -16,6 +18,7 @@
"@tailwindcss/typography": "^0.5.19",
"@tailwindcss/vite": "^4.2.2",
"@types/node": "^25.9.1",
+ "@types/pako": "^2.0.4",
"sass": "^1.100.0",
"svelte": "^5.55.2",
"svelte-check": "^4.4.6",
@@ -193,6 +196,8 @@
"@types/node": ["@types/node@25.9.1", "", { "dependencies": { "undici-types": ">=7.24.0 <7.24.7" } }, "sha512-xfrlY7UD5rMJk3ZVJP8BNzS28J36YJg+xp+LPXV1TdWxr8uMH5A860QNxYDGQe/ylDSgjxE52Q9VnO7p75tJxg=="],
+ "@types/pako": ["@types/pako@2.0.4", "", {}, "sha512-VWDCbrLeVXJM9fihYodcLiIv0ku+AlOa/TQ1SvYOaBuyrSKgEcro95LJyIsJ4vSo6BXIxOKxiJAat04CmST9Fw=="],
+
"@types/trusted-types": ["@types/trusted-types@2.0.7", "", {}, "sha512-ScaPdn1dQczgbl0QFTeTOmVHFULt394XJgOQNoyVhZ6r2vLnMLJfBPd53SB52T/3G36VI1/g2MZaX0cwDuXsfw=="],
"acorn": ["acorn@8.16.0", "", { "bin": { "acorn": "bin/acorn" } }, "sha512-UVJyE9MttOsBQIDKw1skb9nAwQuR5wuGD3+82K6JgJlm/Y+KI92oNsMNGZCYdDsVtRHSak0pcV5Dno5+4jh9sw=="],
@@ -293,6 +298,8 @@
"locate-character": ["locate-character@3.0.0", "", {}, "sha512-SW13ws7BjaeJ6p7Q6CO2nchbYEc3X3J6WrmTTDto7yMPqVSZTUyY5Tjbid+Ab8gLnATtygYtiDIJGQRRn2ZOiA=="],
+ "lottie-web": ["lottie-web@5.13.0", "", {}, "sha512-+gfBXl6sxXMPe8tKQm7qzLnUy5DUPJPKIyRHwtpCpyUEYjHYRJC/5gjUvdkuO2c3JllrPtHXH5UJJK8LRYl5yQ=="],
+
"lru-cache": ["lru-cache@11.5.1", "", {}, "sha512-RPimw/7aMdv2oqRrxKwvZXcPfwBrn/JZ2xYcY9Hus/6LaS3VOAKVWKWgNLCFSiOm1ESXinjsDlidVU7JlnCN2A=="],
"lz-string": ["lz-string@1.5.0", "", { "bin": { "lz-string": "bin/bin.js" } }, "sha512-h5bgJWpxJNswbU7qCrV0tIKQCaS3blPDrqKWx+QxzuzL1zGUzij9XCWLrSLsJPu5t+eWA/ycetzYAO5IOMcWAQ=="],
@@ -315,6 +322,8 @@
"obug": ["obug@2.1.1", "", {}, "sha512-uTqF9MuPraAQ+IsnPf366RG4cP9RtUi7MLO1N3KEc+wb0a6yKpeL0lmk2IB1jY5KHPAlTc6T/JRdC/YqxHNwkQ=="],
+ "pako": ["pako@2.1.0", "", {}, "sha512-w+eufiZ1WuJYgPXbV/PO3NCMEc3xqylkKHzp8bxp1uW4qaSNQUkwmLLEc3kKsfz8lpV1F8Ht3U1Cm+9Srog2ug=="],
+
"path-key": ["path-key@3.1.1", "", {}, "sha512-ojmeN0qd+y0jszEtoY48r0Peq5dwMEkIlCOu6Q5f41lfkswXuKtYrhgoTpLnyIcHm24Uhqx+5Tqm2InSwLhE6Q=="],
"path-scurry": ["path-scurry@2.0.2", "", { "dependencies": { "lru-cache": "^11.0.0", "minipass": "^7.1.2" } }, "sha512-3O/iVVsJAPsOnpwWIeD+d6z/7PmqApyQePUtCndjatj/9I5LylHvt5qluFaBT3I5h3r1ejfR056c+FCv+NnNXg=="],
diff --git a/frontend/package.json b/frontend/package.json
index 70a2d73..a078a2a 100644
--- a/frontend/package.json
+++ b/frontend/package.json
@@ -8,6 +8,7 @@
"@tailwindcss/typography": "^0.5.19",
"@tailwindcss/vite": "^4.2.2",
"@types/node": "^25.9.1",
+ "@types/pako": "^2.0.4",
"sass": "^1.100.0",
"svelte": "^5.55.2",
"svelte-check": "^4.4.6",
@@ -31,6 +32,8 @@
"version": "0.0.1",
"private": true,
"dependencies": {
- "bits-ui": "^2.18.1"
+ "bits-ui": "^2.18.1",
+ "lottie-web": "^5.13.0",
+ "pako": "^2.1.0"
}
}
diff --git a/frontend/src/lib/api/custom-emoji.ts b/frontend/src/lib/api/custom-emoji.ts
new file mode 100644
index 0000000..863cf37
--- /dev/null
+++ b/frontend/src/lib/api/custom-emoji.ts
@@ -0,0 +1,70 @@
+import { accounts } from "$lib/stores/accounts.svelte";
+import { auth } from "$lib/stores/auth.svelte";
+
+const BASE = import.meta.env.VITE_API_BASE ?? "/api";
+const RETRY_DELAY = 2500;
+
+export interface CustomEmojiAsset {
+ mime: string;
+ url: string;
+}
+
+const ready = new Map();
+const missing = new Set();
+const inflight = new Map>();
+
+function authHeaders(): Record {
+ return auth.token ? { Authorization: `Bearer ${auth.token}` } : {};
+}
+
+function delay(ms: number): Promise {
+ return new Promise((resolve) => {
+ setTimeout(resolve, ms);
+ });
+}
+
+async function fetchEmoji(
+ account: number,
+ id: string,
+ key: string,
+ retry: boolean
+): Promise {
+ const url = `${BASE}/custom-emoji/${id}?account_id=${account}`;
+ const response = await fetch(url, { headers: authHeaders() });
+ if (response.ok) {
+ const blob = await response.blob();
+ const asset = { url: URL.createObjectURL(blob), mime: blob.type };
+ ready.set(key, asset);
+ return asset;
+ }
+ if (response.status === 409 && retry) {
+ await delay(RETRY_DELAY);
+ return fetchEmoji(account, id, key, false);
+ }
+ missing.add(key);
+ return null;
+}
+
+export function loadCustomEmoji(id: string): Promise {
+ const account = accounts.selectedId;
+ if (account === null) {
+ return Promise.resolve(null);
+ }
+ const key = `${account}:${id}`;
+ const cached = ready.get(key);
+ if (cached) {
+ return Promise.resolve(cached);
+ }
+ if (missing.has(key)) {
+ return Promise.resolve(null);
+ }
+ const existing = inflight.get(key);
+ if (existing) {
+ return existing;
+ }
+ const promise = fetchEmoji(account, id, key, true).finally(() => {
+ inflight.delete(key);
+ });
+ inflight.set(key, promise);
+ return promise;
+}
diff --git a/frontend/src/lib/api/endpoints.ts b/frontend/src/lib/api/endpoints.ts
index 5a8a294..20a6dd5 100644
--- a/frontend/src/lib/api/endpoints.ts
+++ b/frontend/src/lib/api/endpoints.ts
@@ -1,14 +1,25 @@
import { request } from "$lib/api/client";
import type {
Account,
+ CaptureToggles,
Chat,
Folder,
+ JobStatus,
JobView,
MediaVersion,
MediaView,
MessageVersion,
MessageView,
PeerView,
+ PinnedView,
+ PolicyChatKind,
+ PolicyCreate,
+ PolicyRecord,
+ PresenceHourly,
+ PresenceSample,
+ ResponseStats,
+ SearchHit,
+ VolumeBucket,
} from "$lib/api/types";
import { accounts } from "$lib/stores/accounts.svelte";
@@ -29,6 +40,43 @@ export function listFolders(): Promise {
return request("/folders", { account: true });
}
+export function listPolicies(): Promise {
+ return request("/policy", { account: true });
+}
+
+export function createPolicy(body: PolicyCreate): Promise {
+ return request("/policy", {
+ method: "POST",
+ body: { ...body, account_id: accounts.selectedId },
+ });
+}
+
+export function updatePolicy(
+ id: number,
+ toggles: CaptureToggles
+): Promise {
+ return request(`/policy/${id}`, {
+ method: "PUT",
+ body: toggles,
+ });
+}
+
+export function deletePolicy(id: number): Promise {
+ return request(`/policy/${id}`, { method: "DELETE" });
+}
+
+export function effectivePolicy(query: {
+ chat_id: number;
+ is_bot?: boolean;
+ is_contact?: boolean | null;
+ kind: PolicyChatKind;
+}): Promise {
+ return request("/policy/effective", {
+ account: true,
+ query,
+ });
+}
+
export function listMessages(
chatId: number,
options: Page & { include_deleted?: boolean } = {}
@@ -39,6 +87,55 @@ export function listMessages(
});
}
+export function getCurrentPresence(
+ peerId: number
+): Promise {
+ return request("/presence/current", {
+ account: true,
+ query: { peer_id: peerId },
+ });
+}
+
+export function getPresenceHistory(
+ peerId: number,
+ page: Page = {}
+): Promise {
+ return request("/presence", {
+ account: true,
+ query: { peer_id: peerId, ...page },
+ });
+}
+
+export function getPresenceHourly(peerId: number): Promise {
+ return request("/presence/hourly", {
+ account: true,
+ query: { peer_id: peerId },
+ });
+}
+
+export function getMessageVolume(
+ chatId: number,
+ days = 90
+): Promise {
+ return request("/analytics/volume", {
+ account: true,
+ query: { chat_id: chatId, days },
+ });
+}
+
+export function getResponseStats(chatId: number): Promise {
+ return request("/analytics/response-time", {
+ account: true,
+ query: { chat_id: chatId },
+ });
+}
+
+export function getPinned(chatId: number): Promise {
+ return request(`/chats/${chatId}/pinned`, {
+ account: true,
+ });
+}
+
export function listMessageVersions(
chatId: number,
messageId: number
@@ -83,6 +180,30 @@ export function getJob(jobId: number): Promise {
return request(`/jobs/${jobId}`, { account: true });
}
+export function listJobs(status?: JobStatus): Promise {
+ return request("/jobs", {
+ account: true,
+ query: status ? { status } : {},
+ });
+}
+
+export function enqueueBackfill(
+ chatId: number,
+ media: boolean
+): Promise<{ job_id: number }> {
+ return request<{ job_id: number }>("/backfill", {
+ method: "POST",
+ body: { account_id: accounts.selectedId, chat_id: chatId, media },
+ });
+}
+
+export function syncDialogs(): Promise<{ job_id: number }> {
+ return request<{ job_id: number }>("/dialogs/sync", {
+ method: "POST",
+ body: { account_id: accounts.selectedId },
+ });
+}
+
export function getMediaVersions(
chatId: number,
messageId: number
@@ -105,6 +226,16 @@ export function getMessageMedia(
});
}
+export function searchMessages(
+ query: string,
+ options: Page & { chat_id?: number } = {}
+): Promise {
+ return request("/search", {
+ account: true,
+ query: { query, ...options },
+ });
+}
+
export function fetchMedia(
chatId: number,
messageId: number
diff --git a/frontend/src/lib/api/types.ts b/frontend/src/lib/api/types.ts
index 44ce2e8..4232542 100644
--- a/frontend/src/lib/api/types.ts
+++ b/frontend/src/lib/api/types.ts
@@ -139,6 +139,13 @@ export interface ServiceView {
pinned_message_id: number | null;
}
+export interface PinnedView {
+ media_kind: string | null;
+ message_id: number;
+ sender_name: string | null;
+ text: string | null;
+}
+
export interface StickerView {
emoji: string | null;
height: number | null;
@@ -169,6 +176,7 @@ export interface MessageView {
poll: PollView | null;
quote: string | null;
reactions: ReactionCount[];
+ read: boolean;
reply: ReplyView | null;
sender_id: number | null;
service: ServiceView | null;
@@ -257,6 +265,20 @@ export interface PresenceHourly {
samples: number;
}
+export interface VolumeBucket {
+ bucket: string;
+ incoming: number;
+ outgoing: number;
+ total: number;
+}
+
+export interface ResponseStats {
+ mine_count: number;
+ mine_median_seconds: number | null;
+ their_count: number;
+ their_median_seconds: number | null;
+}
+
export interface PeerView {
first_name: string | null;
has_avatar: boolean;
@@ -323,6 +345,13 @@ export interface PolicyRecord extends CaptureToggles {
scope_type: PolicyScopeType;
}
+export type PolicyChatKind = "dm" | "group" | "channel";
+
+export interface PolicyCreate extends CaptureToggles {
+ scope_id?: number | null;
+ scope_type: PolicyScopeType;
+}
+
export interface Folder {
bots: boolean;
broadcasts: boolean;
@@ -373,3 +402,32 @@ export interface JobView {
started_at: string | null;
status: JobStatus;
}
+
+export interface LiveMessageEvent {
+ message: MessageView;
+ type: "message" | "edit" | "reaction";
+}
+
+export interface LiveDeleteEvent {
+ chat_id: number | null;
+ message_ids: number[];
+ type: "delete";
+}
+
+export interface LivePresenceEvent {
+ peer_id: number;
+ sample: PresenceSample | null;
+ type: "presence";
+}
+
+export interface LiveReceiptEvent {
+ chat_id: number;
+ read_up_to: number;
+ type: "receipt";
+}
+
+export type LiveEvent =
+ | LiveMessageEvent
+ | LiveDeleteEvent
+ | LivePresenceEvent
+ | LiveReceiptEvent;
diff --git a/frontend/src/lib/components/ActionMessage.svelte b/frontend/src/lib/components/ActionMessage.svelte
new file mode 100644
index 0000000..96515a6
--- /dev/null
+++ b/frontend/src/lib/components/ActionMessage.svelte
@@ -0,0 +1,138 @@
+
+
+