diff --git a/backend/alembic/versions/f9eff917e5b3_room_membership_last_read_at_for_unread_.py b/backend/alembic/versions/f9eff917e5b3_room_membership_last_read_at_for_unread_.py new file mode 100644 index 0000000..bfa7e71 --- /dev/null +++ b/backend/alembic/versions/f9eff917e5b3_room_membership_last_read_at_for_unread_.py @@ -0,0 +1,32 @@ +"""room membership last_read_at for unread indicators + +Revision ID: f9eff917e5b3 +Revises: f0f6e494454a +Create Date: 2026-08-16 20:14:32.449934 + +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = 'f9eff917e5b3' +down_revision: Union[str, Sequence[str], None] = 'f0f6e494454a' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + """Upgrade schema.""" + # ### commands auto generated by Alembic - please adjust! ### + op.add_column('room_memberships', sa.Column('last_read_at', sa.DateTime(timezone=True), server_default=sa.text('now()'), nullable=False)) + # ### end Alembic commands ### + + +def downgrade() -> None: + """Downgrade schema.""" + # ### commands auto generated by Alembic - please adjust! ### + op.drop_column('room_memberships', 'last_read_at') + # ### end Alembic commands ### diff --git a/backend/app/models/membership.py b/backend/app/models/membership.py index 5d8f289..90fbcfd 100644 --- a/backend/app/models/membership.py +++ b/backend/app/models/membership.py @@ -26,6 +26,13 @@ class RoomMembership(Base): joined_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now(), nullable=False ) + # A server_default (not an app-code default) so every membership-creation + # call site (create_room, join_room, add_member) gets a sane starting + # point automatically: joining counts as being caught up as of then, not + # retroactively unread for the room's entire prior history. + last_read_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), server_default=func.now(), nullable=False + ) room = relationship("Room", back_populates="memberships") user = relationship("User") diff --git a/backend/app/routers/rooms.py b/backend/app/routers/rooms.py index 700d1f8..c2eab0a 100644 --- a/backend/app/routers/rooms.py +++ b/backend/app/routers/rooms.py @@ -62,6 +62,7 @@ from app.services.room_service import ( list_member_rooms, list_open_rooms, list_room_members, + mark_room_read, remove_member, transfer_ownership, update_room, @@ -138,8 +139,9 @@ async def list_my_rooms_endpoint( owner_id=room.owner_id, created_at=room.created_at, role=role, + has_unread=has_unread, ) - for room, role in rooms + for room, role, has_unread in rooms ] @@ -206,6 +208,16 @@ async def leave_room_endpoint( raise HTTPException(status_code=404, detail="Not a member of this room") +@router.post("/{room_id}/read", status_code=204) +async def mark_room_read_endpoint( + room_id: uuid.UUID, + current_user: User = Depends(get_current_user), + db: AsyncSession = Depends(get_db), +): + await require_room_member(room_id, current_user, db) + await mark_room_read(db, room_id, current_user.id) + + def _member_status(user: User, online_ids: set[uuid.UUID]) -> str: # appear_offline always wins, regardless of actual connection -- that's # the whole point of the override (lurking in a room undetected). diff --git a/backend/app/schemas/room.py b/backend/app/schemas/room.py index 18b2837..83227ed 100644 --- a/backend/app/schemas/room.py +++ b/backend/app/schemas/room.py @@ -35,6 +35,10 @@ class RoomListItem(RoomRead): class MyRoomItem(RoomRead): role: RoomRole + # Whether this room has a message newer than the caller's last_read_at -- + # computed by the router/service, not a stored column on Room itself + # (it's inherently per-viewer, unlike everything else on RoomRead). + has_unread: bool class RoomMemberRead(BaseModel): diff --git a/backend/app/services/message_events.py b/backend/app/services/message_events.py index 41eaacb..ddd81a0 100644 --- a/backend/app/services/message_events.py +++ b/backend/app/services/message_events.py @@ -12,7 +12,12 @@ from app.ws.presence import Presence async def _notify_offline_members( - db: AsyncSession, presence: Presence, room_id: uuid.UUID, sender: User, message: Message + db: AsyncSession, + broadcaster: Broadcaster, + presence: Presence, + room_id: uuid.UUID, + sender: User, + message: Message, ) -> None: result = await db.execute( select(RoomMembership.user_id).where(RoomMembership.room_id == room_id) @@ -26,6 +31,17 @@ async def _notify_offline_members( if not offline_ids: return + # This is also exactly the right audience for "give this room an unread + # dot": presence.connected_user_ids(room_id) means "has this room's + # channel joined right now" -- which the client only does while the tab + # is genuinely foregrounded (see useChatSocket.ts's visibility-gated + # join/leave), so a backgrounded-but-open room correctly lands here too, + # not just rooms that aren't open at all. + for user_id in offline_ids: + await broadcaster.publish_to_user( + user_id, {"type": "unread_update", "room_id": str(room_id)} + ) + room = await db.get(Room, room_id) if message.content: body = f"{sender.username}: {message.content}"[:120] @@ -81,7 +97,7 @@ async def broadcast_new_message( trigger identical fan-out/push/event behavior.""" payload = await _message_payload(db, message, sender.username) await broadcaster.publish(room_id, payload) - await _notify_offline_members(db, presence, room_id, sender, message) + await _notify_offline_members(db, broadcaster, presence, room_id, sender, message) await dispatch_event(db, "message.created", room_id, payload) diff --git a/backend/app/services/room_service.py b/backend/app/services/room_service.py index 5436b3f..722e320 100644 --- a/backend/app/services/room_service.py +++ b/backend/app/services/room_service.py @@ -1,6 +1,6 @@ import uuid -from sqlalchemy import delete, select +from sqlalchemy import delete, func, select from sqlalchemy.exc import IntegrityError from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import selectinload @@ -79,14 +79,25 @@ async def list_open_rooms(db: AsyncSession, user_id: uuid.UUID) -> list[tuple[Ro ] -async def list_member_rooms(db: AsyncSession, user_id: uuid.UUID) -> list[tuple[Room, RoomRole]]: +async def list_member_rooms( + db: AsyncSession, user_id: uuid.UUID +) -> list[tuple[Room, RoomRole, bool]]: + last_message_at = ( + select(func.max(Message.created_at)) + .where(Message.room_id == Room.id) + .correlate(Room) + .scalar_subquery() + ) result = await db.execute( - select(Room, RoomMembership.role) + select(Room, RoomMembership.role, RoomMembership.last_read_at, last_message_at) .join(RoomMembership, RoomMembership.room_id == Room.id) .where(RoomMembership.user_id == user_id) .order_by(Room.created_at) ) - return [(room, role) for room, role in result.all()] + return [ + (room, role, last_message_at is not None and last_message_at > last_read_at) + for room, role, last_read_at, last_message_at in result.all() + ] async def get_room(db: AsyncSession, room_id: uuid.UUID) -> Room: @@ -244,6 +255,12 @@ async def transfer_ownership( return room +async def mark_room_read(db: AsyncSession, room_id: uuid.UUID, user_id: uuid.UUID) -> None: + membership = await _get_membership(db, room_id, user_id) + membership.last_read_at = func.now() + await db.commit() + + async def leave_room(db: AsyncSession, room_id: uuid.UUID, user_id: uuid.UUID) -> None: membership = await _get_membership(db, room_id, user_id) if membership.role == RoomRole.owner: diff --git a/backend/app/ws/chat.py b/backend/app/ws/chat.py index ea1d8b8..b62b991 100644 --- a/backend/app/ws/chat.py +++ b/backend/app/ws/chat.py @@ -21,6 +21,7 @@ from app.services.message_service import ( edit_message, toggle_reaction, ) +from app.services.room_service import mark_room_read router = APIRouter(tags=["ws"]) @@ -162,6 +163,15 @@ async def chat_endpoint(websocket: WebSocket, db: AsyncSession = Depends(get_db) message = await create_message( db, envelope.room_id, user.id, envelope.content, image_id, file_id ) + # Sending implies having seen the room as of now -- without + # this, GET /rooms/mine would show the sender's own room as + # unread the instant they send into it (last_read_at isn't + # otherwise bumped until the frontend's own message echo + # triggers a mark-read call, which is a real but avoidable + # race). Deliberately not done in create_message() itself: + # the incoming-webhook path also calls it, and a webhook's + # attributed sender may not actually be watching. + await mark_room_read(db, envelope.room_id, user.id) await broadcast_new_message(db, broadcaster, presence, envelope.room_id, message, user) elif envelope.type == "edit": diff --git a/backend/tests/test_unread.py b/backend/tests/test_unread.py new file mode 100644 index 0000000..3af8b9e --- /dev/null +++ b/backend/tests/test_unread.py @@ -0,0 +1,175 @@ +import uuid + +from tests.conftest import register_and_login + + +def _unique(prefix: str) -> str: + return f"{prefix}-{uuid.uuid4().hex[:8]}" + + +def _recv(ws) -> dict: + """Reads the next frame, discarding member_updated presence-change + broadcasts -- another connection going online/offline is real, expected + noise these tests aren't about.""" + while True: + msg = ws.receive_json() + if msg.get("type") != "member_updated": + return msg + + +def _register_ws(ws_client, username: str) -> dict: + from app.schemas.user import UserCreate + from app.services.auth_service import register_user + + async def _seed(): + async with ws_client.session_factory() as session: + await register_user( + session, + UserCreate(username=username, email=f"{username}@example.com", password="password123"), + ) + + ws_client.portal.call(_seed) + resp = ws_client.post( + "/api/auth/login", json={"username_or_email": username, "password": "password123"} + ) + assert resp.status_code == 200, resp.text + return resp.json() + + +def _has_unread(rooms: list[dict], room_id: str) -> bool: + return next(r for r in rooms if r["id"] == room_id)["has_unread"] + + +def _send_and_sync(ws, room_id: str, content: str) -> dict: + """Sends a message and waits for its ack, then a sync barrier: the WS + handler processes frames strictly sequentially, so an ack for a second, + idempotent "join" only arrives once the message frame's *full* handling + -- including the offline-member notify step this feature hooks into -- + has actually finished. Without this, the message's own ack (itself just + a mid-handler side effect, not the handler's return) proves nothing + about whether _notify_offline_members has run yet, and closing the + sender's socket right after that ack can cancel that still-in-flight + work (mirrors the same "sync barrier" pattern in test_broadcast.py).""" + ws.send_json({"type": "message", "room_id": room_id, "content": content}) + message = ws.receive_json() + ws.send_json({"type": "join", "room_id": room_id}) + assert ws.receive_json()["type"] == "joined" + return message + + +def test_message_marks_room_unread_and_notifies_offline_member(ws_client_factory): + instance1 = ws_client_factory() + instance2 = ws_client_factory() + + alice = _register_ws(instance1, _unique("alice")) + room = instance1.post("/api/rooms", json={"name": _unique("general")}).json() + + bob = _register_ws(instance2, _unique("bob")) + instance2.post(f"/api/rooms/{room['id']}/join") + + # Bob is connected (so he can receive the per-user unread_update signal) + # but never joins this room's channel -- exactly the "room isn't open" + # case this feature exists for. + with instance2.websocket_connect("/ws/chat") as bob_ws: + with instance1.websocket_connect("/ws/chat") as alice_ws: + alice_ws.send_json({"type": "join", "room_id": room["id"]}) + assert alice_ws.receive_json()["type"] == "joined" + message = _send_and_sync(alice_ws, room["id"], "hi bob") + assert message["type"] == "message" + + update = _recv(bob_ws) + assert update == {"type": "unread_update", "room_id": room["id"]} + + bob_rooms = instance2.get("/api/rooms/mine").json() + assert _has_unread(bob_rooms, room["id"]) is True + + alice_rooms = instance1.get("/api/rooms/mine").json() + assert _has_unread(alice_rooms, room["id"]) is False + + +def test_no_unread_signal_for_member_with_room_joined(ws_client_factory): + instance1 = ws_client_factory() + instance2 = ws_client_factory() + + alice = _register_ws(instance1, _unique("alice")) + room = instance1.post("/api/rooms", json={"name": _unique("general")}).json() + + bob = _register_ws(instance2, _unique("bob")) + instance2.post(f"/api/rooms/{room['id']}/join") + + with instance2.websocket_connect("/ws/chat") as bob_ws: + bob_ws.send_json({"type": "join", "room_id": room["id"]}) + assert bob_ws.receive_json()["type"] == "joined" + + with instance1.websocket_connect("/ws/chat") as alice_ws: + alice_ws.send_json({"type": "join", "room_id": room["id"]}) + assert alice_ws.receive_json()["type"] == "joined" + message = _send_and_sync(alice_ws, room["id"], "hi bob") + assert message["type"] == "message" + + # Bob has this room's channel joined, so he gets the normal message + # broadcast, not an unread_update -- he's actively watching. Alice's + # sync-barrier "joined" ack (from _send_and_sync) is private to her + # own connection, not broadcast, so bob sees nothing further here. + # + # Note: this only proves the real-time *signal* is suppressed for a + # joined member. Persisted has_unread (GET /rooms/mine) is a + # separate, client-driven mechanism (see test_mark_read_endpoint_ + # clears_unread) -- being joined to the channel doesn't by itself + # advance last_read_at server-side; the frontend does that + # explicitly whenever a message arrives while the room is both + # joined and genuinely visible. + update = _recv(bob_ws) + assert update["type"] == "message" + + +def test_mark_read_endpoint_clears_unread(ws_client_factory): + instance1 = ws_client_factory() + instance2 = ws_client_factory() + + _register_ws(instance1, _unique("alice")) + room = instance1.post("/api/rooms", json={"name": _unique("general")}).json() + + _register_ws(instance2, _unique("bob")) + instance2.post(f"/api/rooms/{room['id']}/join") + + with instance1.websocket_connect("/ws/chat") as alice_ws: + alice_ws.send_json({"type": "join", "room_id": room["id"]}) + assert alice_ws.receive_json()["type"] == "joined" + message = _send_and_sync(alice_ws, room["id"], "hi bob") + assert message["type"] == "message" + + assert _has_unread(instance2.get("/api/rooms/mine").json(), room["id"]) is True + + resp = instance2.post(f"/api/rooms/{room['id']}/read") + assert resp.status_code == 204 + + assert _has_unread(instance2.get("/api/rooms/mine").json(), room["id"]) is False + + +def test_joining_room_does_not_retroactively_mark_history_unread(ws_client_factory): + instance1 = ws_client_factory() + instance2 = ws_client_factory() + + _register_ws(instance1, _unique("alice")) + room = instance1.post("/api/rooms", json={"name": _unique("general")}).json() + + with instance1.websocket_connect("/ws/chat") as alice_ws: + alice_ws.send_json({"type": "join", "room_id": room["id"]}) + assert alice_ws.receive_json()["type"] == "joined" + message = _send_and_sync(alice_ws, room["id"], "before bob joins") + assert message["type"] == "message" + + _register_ws(instance2, _unique("bob")) + instance2.post(f"/api/rooms/{room['id']}/join") + + assert _has_unread(instance2.get("/api/rooms/mine").json(), room["id"]) is False + + +async def test_mark_read_requires_room_membership(client, db_session): + await register_and_login(client, db_session, username=_unique("alice")) + room = (await client.post("/api/rooms", json={"name": _unique("general")})).json() + + await register_and_login(client, db_session, username=_unique("bob")) + resp = await client.post(f"/api/rooms/{room['id']}/read") + assert resp.status_code == 403 diff --git a/frontend/src/api/rooms.ts b/frontend/src/api/rooms.ts index cb0d2a5..deca755 100644 --- a/frontend/src/api/rooms.ts +++ b/frontend/src/api/rooms.ts @@ -59,6 +59,10 @@ export function listRoomAttachments(roomId: string): Promise { return apiFetch(`/api/rooms/${roomId}/attachments`) } +export function markRoomRead(roomId: string): Promise { + return apiFetch(`/api/rooms/${roomId}/read`, { method: 'POST' }) +} + export function addRoomMember(roomId: string, userId: string): Promise { return apiFetch(`/api/rooms/${roomId}/members`, { method: 'POST', diff --git a/frontend/src/components/ChatPane.tsx b/frontend/src/components/ChatPane.tsx index 7a33a2e..797736c 100644 --- a/frontend/src/components/ChatPane.tsx +++ b/frontend/src/components/ChatPane.tsx @@ -1,6 +1,6 @@ import { useCallback, useEffect, useState } from 'react' import { NetworkError } from '../api/client' -import { getRoomMessages } from '../api/rooms' +import { getRoomMessages, markRoomRead } from '../api/rooms' import type { ChatSocketHandle } from '../ws/useChatSocket' import type { ChatMessageEnvelope, Message, MyRoomItem, RoomMember, ServerEnvelope } from '../types' import { Composer } from './Composer' @@ -15,9 +15,19 @@ interface ChatPaneProps { onToggleInfo: () => void infoOpen: boolean socket: ChatSocketHandle + onRoomRead: (roomId: string) => void } -export function ChatPane({ room, members, isMobile, onBack, onToggleInfo, infoOpen, socket }: ChatPaneProps) { +export function ChatPane({ + room, + members, + isMobile, + onBack, + onToggleInfo, + infoOpen, + socket, + onRoomRead, +}: ChatPaneProps) { const [history, setHistory] = useState([]) const [live, setLive] = useState([]) const [wsError, setWsError] = useState(null) @@ -56,6 +66,20 @@ export function ChatPane({ room, members, isMobile, onBack, onToggleInfo, infoOp return () => socket.leaveRoom(room.id) }, [socket, room.id]) + const markRead = useCallback(() => { + // Live check, not a cached ref -- same reasoning as joinRoom's in + // useChatSocket.ts: a backgrounded-but-open tab must keep accumulating + // unread rather than auto-marking-read the instant a message arrives + // somewhere it can't actually be seen. + if (document.visibilityState !== 'visible') return + onRoomRead(room.id) + markRoomRead(room.id).catch(() => { + // Best-effort -- an unread dot lagging by one message isn't worth + // surfacing an error for; the next successful mark-read call (or a + // future refreshRooms()) resyncs it. + }) + }, [room.id, onRoomRead]) + useEffect( () => // The socket is shared across every room this tab visits, so a @@ -74,8 +98,10 @@ export function ChatPane({ room, members, isMobile, onBack, onToggleInfo, infoOp // while this socket wasn't in the room's channel, so resync // instead of trusting whatever's already in state. refreshHistory() + markRead() } else if (envelope.type === 'message' && envelope.room_id === room.id) { setLive((prev) => [...prev, envelope]) + markRead() } else if (envelope.type === 'message_update' && envelope.room_id === room.id) { setHistory((prev) => prev.map((m) => @@ -98,7 +124,7 @@ export function ChatPane({ room, members, isMobile, onBack, onToggleInfo, infoOp setWsError(envelope.detail) } }), - [socket, room.id, refreshHistory], + [socket, room.id, refreshHistory, markRead], ) const connected = socket.connected diff --git a/frontend/src/components/RoomRow.css b/frontend/src/components/RoomRow.css index c30be41..0972a65 100644 --- a/frontend/src/components/RoomRow.css +++ b/frontend/src/components/RoomRow.css @@ -50,3 +50,11 @@ text-overflow: ellipsis; margin-top: 2px; } + +.room-row-unread-dot { + flex: none; + width: 8px; + height: 8px; + border-radius: 50%; + background: var(--ds-accent); +} diff --git a/frontend/src/components/RoomRow.tsx b/frontend/src/components/RoomRow.tsx index 6c2cfae..06e0ada 100644 --- a/frontend/src/components/RoomRow.tsx +++ b/frontend/src/components/RoomRow.tsx @@ -32,6 +32,7 @@ export function RoomRow({ room, colorIndex, active }: RoomRowProps) { {room.description &&
{room.description}
} + {room.has_unread && !active && } ) } diff --git a/frontend/src/pages/ChatShellPage.tsx b/frontend/src/pages/ChatShellPage.tsx index 7eb7a51..4be92da 100644 --- a/frontend/src/pages/ChatShellPage.tsx +++ b/frontend/src/pages/ChatShellPage.tsx @@ -59,13 +59,18 @@ export function ChatShellPage() { const socket = useChatSocketContext() + const setRoomUnread = useCallback((id: string, hasUnread: boolean) => { + setRooms((prev) => prev.map((r) => (r.id === id ? { ...r, has_unread: hasUnread } : r))) + }, []) + useEffect( () => socket.subscribe((envelope) => { if (envelope.type === 'room_added') refreshRooms() else if (envelope.type === 'member_updated' && envelope.room_id === roomId) refreshMembers() + else if (envelope.type === 'unread_update') setRoomUnread(envelope.room_id, true) }), - [socket, refreshRooms, refreshMembers, roomId], + [socket, refreshRooms, refreshMembers, roomId, setRoomUnread], ) useEffect(() => { @@ -110,6 +115,7 @@ export function ChatShellPage() { onToggleInfo={() => setInfoOpen((v) => !v)} infoOpen={infoOpen} socket={socket} + onRoomRead={(id) => setRoomUnread(id, false)} /> ) : ( !isMobile && ( diff --git a/frontend/src/types.ts b/frontend/src/types.ts index a639328..3b023fd 100644 --- a/frontend/src/types.ts +++ b/frontend/src/types.ts @@ -37,6 +37,7 @@ export interface RoomListItem extends Room { export interface MyRoomItem extends Room { role: RoomRole + has_unread: boolean } export interface RoomMember { @@ -138,6 +139,11 @@ export interface ChatMemberUpdatedEnvelope { user_id: string } +export interface ChatUnreadUpdateEnvelope { + type: 'unread_update' + room_id: string +} + export type ServerEnvelope = | ChatMessageEnvelope | ChatMessageUpdateEnvelope @@ -146,6 +152,7 @@ export type ServerEnvelope = | ChatErrorEnvelope | ChatRoomAddedEnvelope | ChatMemberUpdatedEnvelope + | ChatUnreadUpdateEnvelope export interface AdminUser { id: string