fix(watchtogether): harden realtime sync, room lifecycle, and lobby UX (#273)

* fix(watchtogether): harden realtime sync, room lifecycle, and lobby UX

Remediates all findings from a deep review of the Watch Together feature.

Server:
- Serialize every websocket write (pong/error replies bypassed the write
  mutex, racing broadcasts on the same gorilla conn)
- Send room_closed with a reason on terminal connect failures so clients
  stop reconnecting to dead rooms
- Persist room state outside the service-wide mutex via a shared
  generation-CAS helper; drop ~450 lines of dead duplicated methods
- Measure transport latency from server-side ping/pong RTT (was one-way
  client-clock delta, poisoned by clock skew) and clamp the lead time
- Re-evaluate readiness when a waiting participant disconnects and add a
  30s waiting deadline that skips stragglers (activates ignoreWait)
- Guard the host-disconnect close timer against reconnect races
- Clamp buffering-report anchor moves; clear stale member sessions on
  selection change
- Janitor: evict empty live rooms and close rooms idle >24h
- Snapshot gains an additive members list with profile display names

Web:
- Surface terminal room errors (REST 404/410/403 and WS error codes) as
  closedReason instead of reconnecting forever on "Connecting..."
- Memoize the playback-sync hook and narrow VideoPlayer's video-listener
  effect deps to stop re-subscribing 13 listeners on every render
- Preserve invite-link destination through login/profile guards
- Lobby: terminal ended/missing-token states with CTAs, End-room confirm
  dialog, toast feedback via shared action helpers (dedup with player),
  participant list with guest Leave, mobile-visible connection status,
  document title, a11y labels/focus reveal, unified status dot component
- Delete dead useWatchTogetherRoom hook (345 lines, zero importers)

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(watchtogether): guard indexed access in join-page keyboard nav for noUncheckedIndexedAccess

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(watchtogether): reconcile CAS conflicts without regressing newer state; roll back unpersisted waiting-resume

Addresses Codex review on PR #273:
- persistRoomChangeLocked now undoes the failed writer's optimistic
  generation increment and only adopts the refreshed database row when it
  is at least as new as the local copy, so a stale conflict refresh can
  no longer overwrite a concurrent writer's newer in-memory state (and a
  failed write can no longer leave a phantom generation)
- maybeResumeFromWaitingLocked restores the waiting state and re-arms the
  deadline when the resume transition fails to persist, instead of
  broadcasting a resume the database never recorded

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Quick
2026-07-02 11:34:02 -04:00
committed by GitHub
co-authored by Claude Fable 5
parent 26089950af
commit b250dbb59b
22 changed files with 1392 additions and 1102 deletions
+58 -15
View File
@@ -7,6 +7,7 @@ import (
"io"
"net/http"
"sync"
"sync/atomic"
"time"
"github.com/Silo-Server/silo-server/internal/access"
@@ -110,9 +111,16 @@ type watchTogetherPingMessage struct {
ClientSentAt string `json:"client_sent_at"`
}
// watchTogetherRoomConn serializes every write to the underlying gorilla
// connection. gorilla/websocket does not support concurrent writers, and room
// broadcasts arrive from other members' goroutines, so all writes — including
// pong/error replies from the read loop — must go through this wrapper.
type watchTogetherRoomConn struct {
conn *websocket.Conn
writeMu sync.Mutex
// pingSentAtNano is the send time of the most recent protocol-level ping,
// used to measure round-trip latency when the pong arrives.
pingSentAtNano atomic.Int64
}
func (c *watchTogetherRoomConn) WriteJSON(v any) error {
@@ -130,9 +138,35 @@ func (c *watchTogetherRoomConn) Close() error {
func (c *watchTogetherRoomConn) WritePing() error {
c.writeMu.Lock()
defer c.writeMu.Unlock()
c.pingSentAtNano.Store(time.Now().UnixNano())
return writeWebSocketControl(c.conn, websocket.PingMessage, nil)
}
// TakePingSentAt returns and clears the send time of the last unanswered
// protocol-level ping.
func (c *watchTogetherRoomConn) TakePingSentAt() time.Time {
nano := c.pingSentAtNano.Swap(0)
if nano == 0 {
return time.Time{}
}
return time.Unix(0, nano)
}
func (c *watchTogetherRoomConn) WriteError(code, message string) {
_ = c.WriteJSON(map[string]string{
"type": "error",
"code": code,
"message": message,
})
}
func (c *watchTogetherRoomConn) writeRoomClosed(reason string) {
_ = c.WriteJSON(map[string]string{
"type": "room_closed",
"reason": reason,
})
}
func NewWatchTogetherHandler(
service *watchtogether.Service,
scopeResolver WatchTogetherScopeResolver,
@@ -653,19 +687,33 @@ func (h *WatchTogetherHandler) HandleRoomWebSocket(w http.ResponseWriter, r *htt
reg, snapshot, err := h.Service.Connect(ctx, roomID, claims.UserID, profileID, realtimeConn)
if err != nil {
// Terminal failures use room_closed so clients stop reconnecting
// instead of retrying a room that will never come back.
if errors.Is(err, watchtogether.ErrRoomNotFound) {
writeWebSocketError(conn, "not_found", "Room not found")
realtimeConn.writeRoomClosed("not_found")
} else if errors.Is(err, watchtogether.ErrRoomClosed) {
writeWebSocketError(conn, "gone", "Room is no longer active")
realtimeConn.writeRoomClosed("ended")
} else {
writeWebSocketError(conn, "internal_error", "Failed to connect room socket")
realtimeConn.WriteError("internal_error", "Failed to connect room socket")
}
return
}
defer h.Service.Disconnect(reg, false)
configureWebSocket(conn)
// Measure round-trip latency from protocol-level ping/pong on the server
// clock; client-reported timestamps are subject to clock skew and cannot
// be trusted for command scheduling.
conn.SetPongHandler(func(string) error {
_ = conn.SetReadDeadline(time.Now().Add(wsPingInterval + wsPongTimeout))
if sentAt := realtimeConn.TakePingSentAt(); !sentAt.IsZero() {
_ = h.Service.HandlePingForConnection(ctx, reg, claims.UserID, profileID, time.Since(sentAt).Milliseconds())
}
return nil
})
startWebSocketPingLoop(ctx, realtimeConn.WritePing)
// Prime an RTT sample right away instead of waiting for the first tick.
_ = realtimeConn.WritePing()
if err := realtimeConn.WriteJSON(map[string]any{
"type": "snapshot",
@@ -680,15 +728,15 @@ func (h *WatchTogetherHandler) HandleRoomWebSocket(w http.ResponseWriter, r *htt
return
}
if err := h.handleRoomClientMessage(ctx, conn, reg, claims.UserID, profileID, data); err != nil {
writeWebSocketError(conn, "bad_request", err.Error())
if err := h.handleRoomClientMessage(ctx, realtimeConn, reg, claims.UserID, profileID, data); err != nil {
realtimeConn.WriteError("bad_request", err.Error())
}
}
}
func (h *WatchTogetherHandler) handleRoomClientMessage(
ctx context.Context,
conn *websocket.Conn,
rc *watchTogetherRoomConn,
reg *watchtogether.Registration,
userID int,
profileID string,
@@ -768,16 +816,11 @@ func (h *WatchTogetherHandler) handleRoomClientMessage(
if err := json.Unmarshal(data, &msg); err != nil {
return err
}
clientSentAt, err := time.Parse(time.RFC3339Nano, msg.ClientSentAt)
if err != nil {
return errors.New("client_sent_at must be RFC3339Nano")
}
// The echoed timestamps only serve the client's clock-offset estimate;
// latency for command scheduling is measured server-side from
// protocol-level ping/pong.
now := time.Now().UTC()
pingMS := now.Sub(clientSentAt.UTC()).Milliseconds()
if pingMS > 0 {
_ = h.Service.HandlePingForConnection(ctx, reg, userID, profileID, pingMS)
}
return writeWebSocketJSON(conn, map[string]string{
return rc.WriteJSON(map[string]string{
"type": "pong",
"client_sent_at": msg.ClientSentAt,
"server_received_at": now.Format(time.RFC3339Nano),
+1 -1
View File
@@ -816,9 +816,9 @@ func NewRouter(deps Dependencies) chi.Router {
watchtogether.NewRepository(deps.DB),
deps.SessionMgr,
deps.FileRepo,
playbackHandler.CommandDispatcher,
watchtogether.NewCatalogSelectionResolver(detailSvc),
watchtogether.NewSuggestionRepository(deps.DB),
watchtogether.NewProfileNameResolver(deps.UserStoreProvider),
),
viewerResolver,
roomTokenService,
+42
View File
@@ -0,0 +1,42 @@
package watchtogether
import (
"context"
"strings"
"github.com/Silo-Server/silo-server/internal/userstore"
)
const fallbackMemberName = "Guest"
// ProfileNameResolver resolves a profile's display name for room member lists.
type ProfileNameResolver interface {
ProfileDisplayName(ctx context.Context, userID int, profileID string) string
}
type userStoreProfileNames struct {
provider userstore.UserStoreProvider
}
// NewProfileNameResolver adapts a UserStoreProvider into a ProfileNameResolver.
func NewProfileNameResolver(provider userstore.UserStoreProvider) ProfileNameResolver {
if provider == nil {
return nil
}
return &userStoreProfileNames{provider: provider}
}
func (r *userStoreProfileNames) ProfileDisplayName(ctx context.Context, userID int, profileID string) string {
if r == nil || r.provider == nil {
return fallbackMemberName
}
store, err := r.provider.ForUser(ctx, userID)
if err != nil || store == nil {
return fallbackMemberName
}
profile, err := store.GetProfile(ctx, profileID)
if err != nil || profile == nil || strings.TrimSpace(profile.Name) == "" {
return fallbackMemberName
}
return profile.Name
}
+31
View File
@@ -152,6 +152,37 @@ func (r *Repository) UpdateAnchor(
)
}
// ListIdleRoomIDs returns rooms that are still open but have seen no
// playback-anchor activity since the cutoff, so the janitor can close them.
func (r *Repository) ListIdleRoomIDs(ctx context.Context, cutoff time.Time, limit int) ([]string, error) {
if r == nil || r.pool == nil {
return nil, fmt.Errorf("watch together repository unavailable")
}
const query = `
SELECT id FROM watch_together_rooms
WHERE phase <> 'ended'
AND GREATEST(anchor_updated_at, created_at) < $1
LIMIT $2
`
rows, err := r.pool.Query(ctx, query, cutoff.UTC(), limit)
if err != nil {
return nil, fmt.Errorf("list idle watch together rooms: %w", err)
}
defer rows.Close()
var roomIDs []string
for rows.Next() {
var id string
if err := rows.Scan(&id); err != nil {
return nil, fmt.Errorf("scan idle watch together room id: %w", err)
}
roomIDs = append(roomIDs, id)
}
return roomIDs, rows.Err()
}
func (r *Repository) CloseRoom(ctx context.Context, roomID string, closedAt time.Time) (*Room, error) {
const query = `
UPDATE watch_together_rooms
File diff suppressed because it is too large Load Diff
+234 -51
View File
@@ -12,6 +12,9 @@ import (
type stubRepo struct {
room Room
// anchorErr, when set, is returned from UpdateAnchor to simulate a
// database failure.
anchorErr error
}
func (s *stubRepo) CreateRoom(_ context.Context, room Room) (*Room, error) {
@@ -31,6 +34,9 @@ func (s *stubRepo) GetRoomByJoinToken(context.Context, string) (*Room, error) {
room := s.room
return &room, nil
}
func (s *stubRepo) ListIdleRoomIDs(context.Context, time.Time, int) ([]string, error) {
return nil, nil
}
func (s *stubRepo) UpdatePolicy(_ context.Context, _ string, policy GuestControlPolicy, generation int64, expectedGeneration int64) (*Room, error) {
if s.room.Generation != expectedGeneration {
return nil, ErrRoomStateConflict
@@ -51,6 +57,9 @@ func (s *stubRepo) UpdateAnchor(
generation int64,
expectedGeneration int64,
) (*Room, error) {
if s.anchorErr != nil {
return nil, s.anchorErr
}
if s.room.Generation != expectedGeneration {
return nil, ErrRoomStateConflict
}
@@ -125,27 +134,6 @@ func (s *stubFiles) GetByID(context.Context, int) (*models.MediaFile, error) {
return &cp, nil
}
type dispatchedCommand struct {
sessionID string
name playback.CommandName
}
type stubDispatcher struct {
commands []dispatchedCommand
}
func (s *stubDispatcher) DispatchToSession(
command playback.CommandEnvelope,
_ time.Duration,
_ func(),
) playback.CommandDispatchResult {
s.commands = append(s.commands, dispatchedCommand{
sessionID: command.SessionID,
name: command.Name,
})
return playback.CommandDispatchResult{}
}
type stubConn struct{}
func (stubConn) WriteJSON(any) error { return nil }
@@ -204,8 +192,8 @@ func baseRoom(now time.Time) Room {
}
}
func newServiceForTest(now time.Time, repo *stubRepo, sessions *stubSessions, files *stubFiles, dispatcher *stubDispatcher, resolver WatchTogetherSelectionResolver) *Service {
service := NewService(repo, sessions, files, dispatcher, resolver, nil)
func newServiceForTest(now time.Time, repo *stubRepo, sessions *stubSessions, files *stubFiles, resolver WatchTogetherSelectionResolver) *Service {
service := NewService(repo, sessions, files, resolver, nil, nil)
service.hostDisconnectTTL = time.Hour
service.now = func() time.Time { return now }
service.rooms[repo.room.ID] = &liveRoom{
@@ -215,6 +203,14 @@ func newServiceForTest(now time.Time, repo *stubRepo, sessions *stubSessions, fi
return service
}
func registrationFor(roomID string, userID int, profileID string, conn RoomConnection) *Registration {
return &Registration{
roomID: roomID,
memberKey: buildMemberKey(userID, profileID),
connection: conn,
}
}
func stringPtr(value string) *string {
return &value
}
@@ -223,14 +219,12 @@ func TestGuestPlayPausePolicyStillRejectsGuestSeek(t *testing.T) {
now := time.Date(2026, 4, 9, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
repo.room.GuestControlPolicy = GuestControlPolicyGuestPlayPause
dispatcher := &stubDispatcher{}
conn := &recordingConn{}
service := newServiceForTest(
now,
repo,
&stubSessions{},
&stubFiles{file: &models.MediaFile{ID: 42, ContentID: "movie-1"}},
dispatcher,
nil,
)
service.rooms[repo.room.ID].members[buildMemberKey(8, "guest")] = &memberState{
@@ -241,20 +235,20 @@ func TestGuestPlayPausePolicyStillRejectsGuestSeek(t *testing.T) {
}
position := 120.0
_, err := service.HandleTransportRequest(context.Background(), repo.room.ID, 8, "guest", TransportRequest{
reg := registrationFor(repo.room.ID, 8, "guest", conn)
_, err := service.HandleTransportRequestForConnection(context.Background(), reg, 8, "guest", TransportRequest{
Action: TransportActionSeek,
PositionSeconds: &position,
IsPaused: false,
})
if !errors.Is(err, ErrTransportNotAllowed) {
t.Fatalf("HandleTransportRequest(guest seek) error = %v, want ErrTransportNotAllowed", err)
t.Fatalf("HandleTransportRequestForConnection(guest seek) error = %v, want ErrTransportNotAllowed", err)
}
}
func TestGuestDriftTriggersCorrection(t *testing.T) {
now := time.Date(2026, 4, 9, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
dispatcher := &stubDispatcher{}
conn := &recordingConn{}
service := newServiceForTest(
now,
@@ -266,7 +260,6 @@ func TestGuestDriftTriggersCorrection(t *testing.T) {
MediaFileID: 42,
}},
&stubFiles{file: &models.MediaFile{ID: 42, ContentID: "movie-1"}},
dispatcher,
nil,
)
service.rooms[repo.room.ID].members[buildMemberKey(8, "guest")] = &memberState{
@@ -276,13 +269,14 @@ func TestGuestDriftTriggersCorrection(t *testing.T) {
connection: conn,
}
_, err := service.HandleStateReport(context.Background(), repo.room.ID, 8, "guest", StateReport{
reg := registrationFor(repo.room.ID, 8, "guest", conn)
_, err := service.HandleStateReportForConnection(context.Background(), reg, 8, "guest", StateReport{
SessionID: "session-1",
PositionSeconds: 2,
IsPaused: false,
})
if err != nil {
t.Fatalf("HandleStateReport() error = %v", err)
t.Fatalf("HandleStateReportForConnection() error = %v", err)
}
if len(conn.payloads) == 0 {
@@ -297,7 +291,6 @@ func TestHostAttachKeepsRoomSelectionAnchor(t *testing.T) {
repo.room.IsPaused = true
repo.room.AnchorUpdatedAt = now
repo.room.Generation = 1
dispatcher := &stubDispatcher{}
conn := &recordingConn{}
service := newServiceForTest(
now,
@@ -311,7 +304,6 @@ func TestHostAttachKeepsRoomSelectionAnchor(t *testing.T) {
IsPaused: false,
}},
&stubFiles{file: &models.MediaFile{ID: 42, ContentID: "movie-1"}},
dispatcher,
nil,
)
service.rooms[repo.room.ID].members[buildMemberKey(7, "host")] = &memberState{
@@ -320,9 +312,10 @@ func TestHostAttachKeepsRoomSelectionAnchor(t *testing.T) {
connection: conn,
}
snapshot, err := service.AttachSession(context.Background(), repo.room.ID, 7, "host", "session-1")
reg := registrationFor(repo.room.ID, 7, "host", conn)
snapshot, err := service.AttachSessionForConnection(context.Background(), reg, 7, "host", "session-1")
if err != nil {
t.Fatalf("AttachSession() error = %v", err)
t.Fatalf("AttachSessionForConnection() error = %v", err)
}
if snapshot.AnchorPositionSeconds != 0 {
@@ -346,7 +339,6 @@ func TestHostAttachKeepsRoomSelectionAnchorEvenWhenGuestAttached(t *testing.T) {
repo.room.IsPaused = true
repo.room.AnchorUpdatedAt = now
repo.room.Generation = 1
dispatcher := &stubDispatcher{}
service := newServiceForTest(
now,
repo,
@@ -359,7 +351,6 @@ func TestHostAttachKeepsRoomSelectionAnchorEvenWhenGuestAttached(t *testing.T) {
IsPaused: false,
}},
&stubFiles{file: &models.MediaFile{ID: 42, ContentID: "movie-1"}},
dispatcher,
nil,
)
service.rooms[repo.room.ID].members[buildMemberKey(8, "guest")] = &memberState{
@@ -374,9 +365,10 @@ func TestHostAttachKeepsRoomSelectionAnchorEvenWhenGuestAttached(t *testing.T) {
connection: stubConn{},
}
snapshot, err := service.AttachSession(context.Background(), repo.room.ID, 7, "host", "host-session")
reg := registrationFor(repo.room.ID, 7, "host", stubConn{})
snapshot, err := service.AttachSessionForConnection(context.Background(), reg, 7, "host", "host-session")
if err != nil {
t.Fatalf("AttachSession() error = %v", err)
t.Fatalf("AttachSessionForConnection() error = %v", err)
}
if snapshot.AnchorPositionSeconds != 0 {
@@ -395,7 +387,6 @@ func TestAttachSessionAcceptsEpisodeContentID(t *testing.T) {
repo.room.IsPaused = true
repo.room.AnchorUpdatedAt = now
repo.room.Generation = 1
dispatcher := &stubDispatcher{}
service := newServiceForTest(
now,
repo,
@@ -412,7 +403,6 @@ func TestAttachSessionAcceptsEpisodeContentID(t *testing.T) {
ContentID: "series-1",
EpisodeID: "episode-19",
}},
dispatcher,
nil,
)
service.rooms[repo.room.ID].members[buildMemberKey(7, "host")] = &memberState{
@@ -421,9 +411,10 @@ func TestAttachSessionAcceptsEpisodeContentID(t *testing.T) {
connection: stubConn{},
}
snapshot, err := service.AttachSession(context.Background(), repo.room.ID, 7, "host", "host-session")
reg := registrationFor(repo.room.ID, 7, "host", stubConn{})
snapshot, err := service.AttachSessionForConnection(context.Background(), reg, 7, "host", "host-session")
if err != nil {
t.Fatalf("AttachSession() error = %v", err)
t.Fatalf("AttachSessionForConnection() error = %v", err)
}
if snapshot.AttachedSessionID != "host-session" {
@@ -437,7 +428,7 @@ func TestAttachSessionAcceptsEpisodeContentID(t *testing.T) {
func TestCreateRoomStartsInLobbyWithoutSelection(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{}
service := NewService(repo, &stubSessions{}, &stubFiles{}, &stubDispatcher{}, nil, nil)
service := NewService(repo, &stubSessions{}, &stubFiles{}, nil, nil, nil)
service.now = func() time.Time { return now }
room, err := service.CreateRoom(context.Background(), CreateRoomInput{
@@ -480,7 +471,6 @@ func TestHostCanSelectItemFromLobby(t *testing.T) {
repo,
&stubSessions{},
&stubFiles{},
&stubDispatcher{},
&stubSelectionResolver{resolved: &ResolvedSelection{
ContentID: "movie-2",
FileID: intPtr(55),
@@ -515,6 +505,41 @@ func TestHostCanSelectItemFromLobby(t *testing.T) {
}
}
func TestSelectItemClearsStaleMemberSessions(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
service := newServiceForTest(
now,
repo,
&stubSessions{},
&stubFiles{},
&stubSelectionResolver{resolved: &ResolvedSelection{ContentID: "movie-2"}},
)
guest := &memberState{
userID: 8,
profileID: "guest",
sessionID: "old-session",
isReady: true,
ignoreWait: true,
connection: stubConn{},
}
service.rooms[repo.room.ID].members[buildMemberKey(8, "guest")] = guest
_, err := service.SelectItem(context.Background(), "room-1", 7, "host", SelectItemInput{
ContentID: "movie-2",
})
if err != nil {
t.Fatalf("SelectItem() error = %v", err)
}
if guest.sessionID != "" {
t.Fatalf("guest session = %q, want cleared", guest.sessionID)
}
if guest.isReady || guest.ignoreWait {
t.Fatal("guest readiness flags should reset on new selection")
}
}
func TestGuestCannotSelectItem(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
@@ -523,7 +548,6 @@ func TestGuestCannotSelectItem(t *testing.T) {
repo,
&stubSessions{},
&stubFiles{},
&stubDispatcher{},
&stubSelectionResolver{resolved: &ResolvedSelection{ContentID: "movie-2"}},
)
@@ -543,7 +567,6 @@ func TestSelectItemRejectsInvalidSelection(t *testing.T) {
repo,
&stubSessions{},
&stubFiles{},
&stubDispatcher{},
&stubSelectionResolver{err: ErrInvalidSelection},
)
@@ -559,7 +582,6 @@ func TestAttachSessionEnforcesSelectedFileID(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
repo.room.SelectedFileID = intPtr(99)
dispatcher := &stubDispatcher{}
service := newServiceForTest(
now,
repo,
@@ -570,7 +592,6 @@ func TestAttachSessionEnforcesSelectedFileID(t *testing.T) {
MediaFileID: 42,
}},
&stubFiles{file: &models.MediaFile{ID: 42, ContentID: "movie-1"}},
dispatcher,
nil,
)
service.rooms[repo.room.ID].members[buildMemberKey(7, "host")] = &memberState{
@@ -579,9 +600,171 @@ func TestAttachSessionEnforcesSelectedFileID(t *testing.T) {
connection: stubConn{},
}
_, err := service.AttachSession(context.Background(), repo.room.ID, 7, "host", "host-session")
reg := registrationFor(repo.room.ID, 7, "host", stubConn{})
_, err := service.AttachSessionForConnection(context.Background(), reg, 7, "host", "host-session")
if !errors.Is(err, ErrSessionMismatch) {
t.Fatalf("AttachSession() error = %v, want ErrSessionMismatch", err)
t.Fatalf("AttachSessionForConnection() error = %v, want ErrSessionMismatch", err)
}
}
func TestDisconnectOfLastUnreadyMemberResumesWaitingRoom(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
repo.room.PlaybackState = RoomPlaybackStateWaiting
repo.room.IsPaused = true
repo.room.ResumeOnReady = true
service := newServiceForTest(now, repo, &stubSessions{}, &stubFiles{}, nil)
hostConn := &recordingConn{}
service.rooms[repo.room.ID].members[buildMemberKey(7, "host")] = &memberState{
userID: 7,
profileID: "host",
sessionID: "host-session",
isReady: true,
connection: hostConn,
}
guestConn := &recordingConn{}
service.rooms[repo.room.ID].members[buildMemberKey(8, "guest")] = &memberState{
userID: 8,
profileID: "guest",
sessionID: "guest-session",
isReady: false,
connection: guestConn,
}
service.Disconnect(registrationFor(repo.room.ID, 8, "guest", guestConn), false)
if repo.room.PlaybackState != RoomPlaybackStatePlaying {
t.Fatalf("playback state = %q, want %q", repo.room.PlaybackState, RoomPlaybackStatePlaying)
}
foundCommand := false
for _, payload := range hostConn.payloads {
if payload["type"] == "transport_command" {
foundCommand = true
}
}
if !foundCommand {
t.Fatal("expected a resume transport command for the remaining member")
}
}
func TestBufferingReportCannotTeleportRoomAnchor(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
service := newServiceForTest(now, repo, &stubSessions{}, &stubFiles{}, nil)
guestConn := &recordingConn{}
service.rooms[repo.room.ID].members[buildMemberKey(8, "guest")] = &memberState{
userID: 8,
profileID: "guest",
sessionID: "guest-session",
connection: guestConn,
}
// Anchor was 10s, 10s ago and playing: expected position is ~20s. A
// report claiming 500s must be clamped back to the expected position.
reg := registrationFor(repo.room.ID, 8, "guest", guestConn)
snapshot, err := service.HandleBufferingForConnection(context.Background(), reg, 8, "guest", StateReport{
SessionID: "guest-session",
PositionSeconds: 500,
IsPaused: false,
})
if err != nil {
t.Fatalf("HandleBufferingForConnection() error = %v", err)
}
if snapshot.PlaybackState != RoomPlaybackStateWaiting {
t.Fatalf("playback state = %q, want %q", snapshot.PlaybackState, RoomPlaybackStateWaiting)
}
if snapshot.AnchorPositionSeconds > 20.001 || snapshot.AnchorPositionSeconds < 19.999 {
t.Fatalf("anchor = %v, want ~20 (clamped)", snapshot.AnchorPositionSeconds)
}
}
func TestPingIsClampedToMaxTransportLead(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
service := newServiceForTest(now, repo, &stubSessions{}, &stubFiles{}, nil)
guestConn := &recordingConn{}
member := &memberState{
userID: 8,
profileID: "guest",
sessionID: "guest-session",
connection: guestConn,
}
service.rooms[repo.room.ID].members[buildMemberKey(8, "guest")] = member
reg := registrationFor(repo.room.ID, 8, "guest", guestConn)
if err := service.HandlePingForConnection(context.Background(), reg, 8, "guest", 3_600_000); err != nil {
t.Fatalf("HandlePingForConnection() error = %v", err)
}
if member.lastPingMS != maxTransportLead.Milliseconds() {
t.Fatalf("lastPingMS = %d, want %d", member.lastPingMS, maxTransportLead.Milliseconds())
}
}
func TestReadyPersistFailureKeepsWaitingState(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
repo.room.PlaybackState = RoomPlaybackStateWaiting
repo.room.IsPaused = true
repo.room.ResumeOnReady = true
repo.anchorErr = errors.New("database unavailable")
service := newServiceForTest(now, repo, &stubSessions{}, &stubFiles{}, nil)
hostConn := &recordingConn{}
service.rooms[repo.room.ID].members[buildMemberKey(7, "host")] = &memberState{
userID: 7,
profileID: "host",
sessionID: "host-session",
connection: hostConn,
}
reg := registrationFor(repo.room.ID, 7, "host", hostConn)
snapshot, err := service.HandleReadyForConnection(context.Background(), reg, 7, "host", StateReport{
SessionID: "host-session",
})
if err != nil {
t.Fatalf("HandleReadyForConnection() error = %v", err)
}
if snapshot.PlaybackState != RoomPlaybackStateWaiting {
t.Fatalf("playback state = %q, want %q (resume must not be announced when persistence failed)",
snapshot.PlaybackState, RoomPlaybackStateWaiting)
}
live := service.rooms[repo.room.ID]
if live.room.Generation != repo.room.Generation {
t.Fatalf("live generation = %d, want %d (failed write must not leave a phantom generation)",
live.room.Generation, repo.room.Generation)
}
if live.waitingTimer == nil {
t.Fatal("waiting deadline should stay armed so the resume is retried")
}
}
func TestStaleLiveConflictAdoptsDatabaseRow(t *testing.T) {
now := time.Date(2026, 4, 10, 12, 0, 20, 0, time.UTC)
repo := &stubRepo{room: baseRoom(now)}
service := newServiceForTest(now, repo, &stubSessions{}, &stubFiles{}, nil)
// The database row has moved ahead of the cached live copy.
repo.room.Generation = 5
hostConn := &recordingConn{}
service.rooms[repo.room.ID].members[buildMemberKey(7, "host")] = &memberState{
userID: 7,
profileID: "host",
sessionID: "host-session",
connection: hostConn,
}
reg := registrationFor(repo.room.ID, 7, "host", hostConn)
snapshot, err := service.HandleTransportRequestForConnection(context.Background(), reg, 7, "host", TransportRequest{
Action: TransportActionPause,
})
if err != nil {
t.Fatalf("HandleTransportRequestForConnection() error = %v", err)
}
if snapshot.Generation != 5 {
t.Fatalf("snapshot generation = %d, want 5 (conflict must adopt the newer database row)", snapshot.Generation)
}
}
+11
View File
@@ -74,6 +74,16 @@ type Room struct {
ClosedAt *time.Time
}
// MemberSummary describes one connected room member in a snapshot.
type MemberSummary struct {
UserID int `json:"user_id"`
ProfileID string `json:"profile_id"`
DisplayName string `json:"display_name"`
IsHost bool `json:"is_host"`
IsSelf bool `json:"is_self"`
Connected bool `json:"connected"`
}
type Snapshot struct {
RoomID string `json:"room_id"`
Phase RoomPhase `json:"phase"`
@@ -97,6 +107,7 @@ type Snapshot struct {
SelfIgnoreWait bool `json:"self_ignore_wait"`
AttachedSessionID string `json:"attached_session_id,omitempty"`
InvitePath string `json:"invite_path,omitempty"`
Members []MemberSummary `json:"members,omitempty"`
}
type RoomJoinResult struct {
+16 -2
View File
@@ -142,8 +142,21 @@ function ScrollRestorationManager() {
return null;
}
/**
* Builds a guard redirect target (e.g. "/login") that preserves the current
* location so the user returns to it after authenticating.
*/
function guardRedirectTarget(base: string, location: ReturnType<typeof useLocation>): string {
const destination = `${location.pathname}${location.search}`;
if (destination === "/" || destination === "") {
return base;
}
return `${base}?redirect=${encodeURIComponent(destination)}`;
}
function RequireAuth({ children }: { children: ReactNode }) {
const { user, loading, setupLoading } = useAuth();
const location = useLocation();
if (loading || setupLoading) {
return (
<div className="p-8" role="status" aria-live="polite">
@@ -152,7 +165,7 @@ function RequireAuth({ children }: { children: ReactNode }) {
</div>
);
}
if (!user) return <Navigate to="/login" replace />;
if (!user) return <Navigate to={guardRedirectTarget("/login", location)} replace />;
return <>{children}</>;
}
@@ -172,7 +185,8 @@ function SetupGate({ children }: { children: ReactNode }) {
function RequireProfile({ children }: { children: ReactNode }) {
const { profile } = useAuth();
if (!profile) return <Navigate to="/profiles" replace />;
const location = useLocation();
if (!profile) return <Navigate to={guardRedirectTarget("/profiles", location)} replace />;
return <>{children}</>;
}
@@ -0,0 +1,32 @@
export type WatchTogetherConnectionState = "disconnected" | "connecting" | "connected";
const connectionLabels: Record<WatchTogetherConnectionState, string> = {
connected: "Connected",
connecting: "Connecting…",
disconnected: "Disconnected",
};
/**
* Shared watch-together connection status label. Pair with
* `ConnectionStatusDot` so every surface uses the same palette + vocabulary.
*/
export function ConnectionStateLabel({ state }: { state: WatchTogetherConnectionState }) {
return <>{connectionLabels[state] ?? connectionLabels.disconnected}</>;
}
/** Shared watch-together connection indicator dot. */
export function ConnectionStatusDot({
state,
className = "h-2 w-2",
}: {
state: WatchTogetherConnectionState;
className?: string;
}) {
const color =
state === "connected"
? "bg-emerald-400"
: state === "connecting"
? "animate-pulse bg-amber-300"
: "bg-red-400";
return <span aria-hidden="true" className={`inline-block rounded-full ${color} ${className}`} />;
}
@@ -0,0 +1,27 @@
import { ConfirmDialog } from "@/components/ConfirmDialog";
/** Shared confirmation for ending a watch party (lobby page + in-player panel). */
export function EndWatchPartyDialog({
open,
onOpenChange,
onConfirm,
isPending,
}: {
open: boolean;
onOpenChange: (open: boolean) => void;
onConfirm: () => void;
isPending?: boolean;
}) {
return (
<ConfirmDialog
open={open}
onOpenChange={onOpenChange}
title="End watch party?"
description="End the watch party for everyone?"
confirmLabel="End Party"
variant="destructive"
onConfirm={onConfirm}
isPending={isPending}
/>
);
}
+11
View File
@@ -7,6 +7,15 @@ export type WatchTogetherPlaybackState = "idle" | "waiting" | "paused" | "playin
export type WatchTogetherSelectionMode = "host_pick" | "vote";
export type WatchTogetherTransportAction = "play" | "pause" | "seek";
export interface WatchTogetherRoomMember {
user_id: number;
profile_id: string;
display_name: string;
is_host: boolean;
is_self: boolean;
connected: boolean;
}
export interface WatchTogetherRoomSnapshot {
room_id: string;
phase: WatchTogetherRoomPhase;
@@ -30,6 +39,8 @@ export interface WatchTogetherRoomSnapshot {
self_ignore_wait: boolean;
attached_session_id?: string;
invite_path?: string;
/** Optional additive field; absent on older servers. Host is listed first. */
members?: WatchTogetherRoomMember[];
}
export interface WatchTogetherTransportCommand {
+79
View File
@@ -0,0 +1,79 @@
import { toast } from "sonner";
import {
buildWatchTogetherInviteUrl,
type GuestControlPolicy,
type WatchTogetherRoomSnapshot,
} from "@/lib/watchTogether";
async function copyTextToClipboard(text: string) {
if (typeof navigator !== "undefined" && navigator.clipboard?.writeText) {
await navigator.clipboard.writeText(text);
return;
}
// Fallback for insecure contexts / older browsers.
const textarea = document.createElement("textarea");
textarea.value = text;
textarea.setAttribute("readonly", "");
textarea.style.position = "fixed";
textarea.style.opacity = "0";
document.body.appendChild(textarea);
textarea.select();
try {
if (!document.execCommand("copy")) {
throw new Error("Copy command failed");
}
} finally {
textarea.remove();
}
}
/**
* Copies the room invite link to the clipboard with toast feedback.
* Returns false when the invite link is not available yet (no toast is shown),
* so callers can surface their own "not ready" message.
*/
export async function copyWatchTogetherInvite(
invitePath: string | null | undefined,
roomCode?: string | null,
): Promise<boolean> {
const inviteUrl = buildWatchTogetherInviteUrl(invitePath);
if (!inviteUrl) {
return false;
}
try {
await copyTextToClipboard(inviteUrl);
toast.success(`Invite copied. Room code ${roomCode ?? ""}`.trim());
} catch {
toast.error("Failed to copy invite link");
}
return true;
}
/** Applies a guest-control policy change with toast feedback. */
export async function setWatchTogetherGuestControl(
updatePolicy: (policy: GuestControlPolicy) => Promise<WatchTogetherRoomSnapshot | null>,
policy: GuestControlPolicy,
): Promise<void> {
try {
const nextRoom = await updatePolicy(policy);
if (nextRoom) {
toast.success(
nextRoom.guest_control_policy === "guest_play_pause"
? "Guests can now pause and resume"
: "Room is now host controlled",
);
}
} catch (error) {
toast.error(error instanceof Error ? error.message : "Failed to update room");
}
}
/** Ends the watch party with toast feedback. */
export async function endWatchTogetherRoom(closeRoom: () => Promise<void>): Promise<void> {
try {
await closeRoom();
toast.success("Room ended");
} catch (error) {
toast.error(error instanceof Error ? error.message : "Failed to end room");
}
}
+24
View File
@@ -13,6 +13,7 @@ const mocks = vi.hoisted(() => ({
useAuth: vi.fn(),
navigate: vi.fn(),
editorPin: "",
searchParams: new URLSearchParams(),
}));
vi.mock("react-router", async () => {
@@ -20,6 +21,7 @@ vi.mock("react-router", async () => {
return {
...actual,
useNavigate: () => mocks.navigate,
useSearchParams: () => [mocks.searchParams, vi.fn()],
};
});
@@ -152,6 +154,7 @@ describe("Profiles", () => {
mocks.useAuth.mockReset();
mocks.navigate.mockReset();
mocks.editorPin = "";
mocks.searchParams = new URLSearchParams();
mocks.useProfiles.mockReturnValue({
data: [makeProfile(), makeProfile({ id: "profile-2", name: "Guest" })],
@@ -246,4 +249,25 @@ describe("Profiles", () => {
expect(selectProfile).toHaveBeenCalledWith(lockedProfile, "verified-token");
expect(mocks.navigate).toHaveBeenCalledWith("/");
});
it("returns to the sanitized redirect target after selecting a profile", async () => {
mocks.searchParams = new URLSearchParams("redirect=%2Frooms%2Fabc%3Froom_token%3Dxyz");
await render(<Profiles />);
await click(findButton(container, "Main"));
expect(selectProfile).toHaveBeenCalled();
expect(mocks.navigate).toHaveBeenCalledWith("/rooms/abc?room_token=xyz");
});
it("ignores unsafe redirect targets when selecting a profile", async () => {
mocks.searchParams = new URLSearchParams("redirect=https%3A%2F%2Fevil.example");
await render(<Profiles />);
await click(findButton(container, "Main"));
expect(mocks.navigate).toHaveBeenCalledWith("/");
});
});
+12 -5
View File
@@ -1,5 +1,5 @@
import { useState } from "react";
import { useNavigate } from "react-router";
import { useNavigate, useSearchParams } from "react-router";
import { Lock, Plus } from "lucide-react";
import { toast } from "sonner";
@@ -13,6 +13,7 @@ import { useDocumentTitle } from "@/hooks/useDocumentTitle";
import { useAuth } from "@/hooks/useAuth";
import { useAvailableUserLibraries } from "@/hooks/queries/libraries";
import { useProfiles } from "@/hooks/queries/profiles";
import { sanitizeAuthRedirect } from "@/lib/authRedirect";
export default function Profiles() {
const { data: profiles = [], isLoading: profilesLoading, avatarUploadEnabled } = useProfiles();
@@ -21,11 +22,17 @@ export default function Profiles() {
const [pinProfile, setPinProfile] = useState<Profile | null>(null);
const { selectProfile, verifyProfilePin, logout } = useAuth();
const navigate = useNavigate();
const [searchParams] = useSearchParams();
const redirectTarget = sanitizeAuthRedirect(searchParams.get("redirect"));
useDocumentTitle("Profiles");
const isLoading = profilesLoading || librariesLoading;
function enterApp() {
navigate(redirectTarget ?? "/");
}
function handleSelect(profile: Profile) {
if (profile.has_pin) {
setPinProfile(profile);
@@ -33,7 +40,7 @@ export default function Profiles() {
}
selectProfile(profile);
navigate("/");
enterApp();
}
async function handleCreateSuccess(
@@ -44,7 +51,7 @@ export default function Profiles() {
) {
if (context.pin === "") {
selectProfile(profile);
navigate("/");
enterApp();
return;
}
@@ -52,7 +59,7 @@ export default function Profiles() {
const response = await verifyProfilePin(profile.id, context.pin);
if (response.valid && response.profile_token) {
selectProfile(profile, response.profile_token);
navigate("/");
enterApp();
return;
}
} catch {
@@ -153,7 +160,7 @@ export default function Profiles() {
onVerified={(profile, token) => {
setPinProfile(null);
selectProfile(profile, token);
navigate("/");
enterApp();
}}
verifyPin={verifyProfilePin}
/>
+106 -41
View File
@@ -1,4 +1,4 @@
import { useCallback, useEffect, useMemo, useState } from "react";
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { useNavigate, useSearchParams } from "react-router";
import { ApiClientError } from "@/api/client";
import { Button } from "@/components/ui/button";
@@ -24,6 +24,23 @@ function describeJoinError(error: unknown) {
return error instanceof Error ? error.message : "Failed to join room.";
}
const selectionModeOptions = [
{
value: "host_pick",
title: "Host Picks",
caption: "Host decides what to watch.",
},
{
value: "vote",
title: "Vote Together",
caption: "Everyone suggests and votes.",
},
] as const satisfies ReadonlyArray<{
value: WatchTogetherSelectionMode;
title: string;
caption: string;
}>;
export default function WatchTogetherJoin() {
useDocumentTitle("Watch Party");
const navigate = useNavigate();
@@ -34,8 +51,35 @@ export default function WatchTogetherJoin() {
const [creating, setCreating] = useState(false);
const [joining, setJoining] = useState(false);
const [error, setError] = useState<string | null>(null);
const modeButtonRefs = useRef<Array<HTMLButtonElement | null>>([]);
const hasInviteToken = token !== "";
const handleModeKeyDown = useCallback(
(event: React.KeyboardEvent<HTMLButtonElement>, index: number) => {
let nextIndex: number | null = null;
if (event.key === "ArrowRight" || event.key === "ArrowDown") {
nextIndex = (index + 1) % selectionModeOptions.length;
} else if (event.key === "ArrowLeft" || event.key === "ArrowUp") {
nextIndex = (index - 1 + selectionModeOptions.length) % selectionModeOptions.length;
} else if (event.key === "Home") {
nextIndex = 0;
} else if (event.key === "End") {
nextIndex = selectionModeOptions.length - 1;
}
if (nextIndex === null) {
return;
}
const nextOption = selectionModeOptions[nextIndex];
if (!nextOption) {
return;
}
event.preventDefault();
setSelectionMode(nextOption.value);
modeButtonRefs.current[nextIndex]?.focus();
},
[],
);
const joinRoom = useCallback(
async (input: { code?: string; join_token?: string }) => {
setJoining(true);
@@ -90,6 +134,25 @@ export default function WatchTogetherJoin() {
[hasInviteToken],
);
// While an invite token is being auto-joined (and hasn't failed yet), show a
// pending state instead of the full create/join form.
const autoJoinPending = hasInviteToken && !error;
if (autoJoinPending) {
return (
<div className="mx-auto flex w-full max-w-5xl flex-col items-center gap-4 px-6 py-24 text-center">
<div
aria-hidden="true"
className="border-foreground/20 border-t-foreground h-8 w-8 animate-spin rounded-full border-2"
/>
<h1 className="text-2xl font-semibold tracking-tight sm:text-3xl">{headline}</h1>
<p className="text-muted-foreground max-w-sm text-sm" role="status">
Hang tight — we're taking you into the room.
</p>
</div>
);
}
return (
<div className="mx-auto flex w-full max-w-5xl flex-col gap-10 px-6 py-10 sm:px-8">
<div className="max-w-2xl">
@@ -99,7 +162,7 @@ export default function WatchTogetherJoin() {
<h1 className="mt-3 text-3xl font-semibold tracking-tight sm:text-4xl">{headline}</h1>
<p className="text-muted-foreground mt-3 max-w-2xl text-sm leading-6">
Start an empty room, invite people in, then choose what everyone watches from the lobby.
Invite links still open the same flow and take guests into the room automatically.
Got an invite link? Just open it — you'll join automatically.
</p>
</div>
@@ -135,7 +198,10 @@ export default function WatchTogetherJoin() {
</div>
{error ? (
<div className="rounded-[8px] border border-red-500/30 bg-red-500/10 px-4 py-3 text-sm text-red-200">
<div
role="alert"
className="rounded-lg border border-red-500/30 bg-red-500/10 px-4 py-3 text-sm text-red-200"
>
{error}
</div>
) : null}
@@ -151,40 +217,39 @@ export default function WatchTogetherJoin() {
</div>
<div>
<label className="text-sm font-medium">Selection mode</label>
<div className="mt-2 flex gap-2" role="radiogroup">
<button
type="button"
role="radio"
aria-checked={selectionMode === "host_pick"}
onClick={() => setSelectionMode("host_pick")}
className={`flex-1 rounded-[8px] border px-3 py-2.5 text-left text-sm transition-colors ${
selectionMode === "host_pick"
? "border-foreground/50 bg-accent font-medium"
: "text-muted-foreground hover:bg-muted/60"
}`}
>
<div className="font-medium">Host Picks</div>
<div className="text-muted-foreground mt-0.5 text-xs">
Host decides what to watch.
</div>
</button>
<button
type="button"
role="radio"
aria-checked={selectionMode === "vote"}
onClick={() => setSelectionMode("vote")}
className={`flex-1 rounded-[8px] border px-3 py-2.5 text-left text-sm transition-colors ${
selectionMode === "vote"
? "border-foreground/50 bg-accent font-medium"
: "text-muted-foreground hover:bg-muted/60"
}`}
>
<div className="font-medium">Vote Together</div>
<div className="text-muted-foreground mt-0.5 text-xs">
Everyone suggests and votes.
</div>
</button>
<label id="watch-selection-mode-label" className="text-sm font-medium">
Selection mode
</label>
<div
className="mt-2 flex gap-2"
role="radiogroup"
aria-labelledby="watch-selection-mode-label"
>
{selectionModeOptions.map((option, index) => {
const selected = selectionMode === option.value;
return (
<button
key={option.value}
ref={(element) => {
modeButtonRefs.current[index] = element;
}}
type="button"
role="radio"
aria-checked={selected}
tabIndex={selected ? 0 : -1}
onClick={() => setSelectionMode(option.value)}
onKeyDown={(event) => handleModeKeyDown(event, index)}
className={`flex-1 rounded-lg border px-3 py-2.5 text-left text-sm transition-colors ${
selected
? "border-foreground/50 bg-accent font-medium"
: "text-muted-foreground hover:bg-muted/60"
}`}
>
<div className="font-medium">{option.title}</div>
<div className="text-muted-foreground mt-0.5 text-xs">{option.caption}</div>
</button>
);
})}
</div>
</div>
@@ -192,7 +257,7 @@ export default function WatchTogetherJoin() {
type="button"
onClick={() => void createRoom()}
disabled={creating || joining}
className="h-11 rounded-[8px] px-5"
className="h-11 rounded-lg px-5"
>
{creating ? "Creating..." : "Create Watch Party"}
</Button>
@@ -204,7 +269,7 @@ export default function WatchTogetherJoin() {
<div>
<h2 className="text-lg font-semibold">Join a room</h2>
<p className="text-foreground/55 mt-1 text-sm leading-6">
Enter a room code or retry an invite link to enter the existing party.
Have a room code from the host? Enter it here to join their party.
</p>
</div>
@@ -235,7 +300,7 @@ export default function WatchTogetherJoin() {
type="button"
onClick={() => void joinRoom({ code: code.trim().toUpperCase() })}
disabled={joining || creating || code.trim() === ""}
className="h-11 rounded-[8px] px-5"
className="h-11 rounded-lg px-5"
>
{joining ? "Joining..." : "Join Watch Party"}
</Button>
@@ -246,7 +311,7 @@ export default function WatchTogetherJoin() {
variant="outline"
onClick={() => void joinRoom({ join_token: token })}
disabled={joining || creating}
className="h-11 rounded-[8px] px-5"
className="h-11 rounded-lg px-5"
>
Retry Invite Link
</Button>
+147 -39
View File
@@ -1,10 +1,11 @@
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { useCallback, useEffect, useMemo, useRef, useState, type ReactNode } from "react";
import { useQuery } from "@tanstack/react-query";
import { useLocation, useParams, useSearchParams } from "react-router";
import { useLocation, useNavigate, useParams, useSearchParams } from "react-router";
import {
ChevronLeft,
ChevronRight,
Link,
Link as LinkIcon,
LogOut,
Play,
Search,
ShieldCheck,
@@ -15,10 +16,20 @@ import {
import { ApiClientError } from "@/api/client";
import { type BrowseItem } from "@/api/types";
import { Button } from "@/components/ui/button";
import {
ConnectionStateLabel,
ConnectionStatusDot,
} from "@/components/watchtogether/ConnectionStatusDot";
import { EndWatchPartyDialog } from "@/components/watchtogether/EndWatchPartyDialog";
import { fetchCatalogPage, createCatalogSearchState } from "@/hooks/queries/catalog";
import { useCatalogItemDetail } from "@/hooks/queries/catalogRead";
import { useSeasons, useSeasonEpisodes } from "@/hooks/queries/episodes";
import { buildWatchTogetherInviteUrl } from "@/lib/watchTogether";
import { useDocumentTitle } from "@/hooks/useDocumentTitle";
import {
copyWatchTogetherInvite,
endWatchTogetherRoom,
setWatchTogetherGuestControl,
} from "@/lib/watchTogetherActions";
import { storage } from "@/utils/storage";
import { useWatchPlaybackController } from "@/playback/watchPlaybackContext";
import { useWatchTogetherRoomConnection } from "@/player/hooks/useWatchTogetherRoomConnection";
@@ -28,6 +39,15 @@ import { decodeThumbhash } from "@/lib/thumbhash";
import { WatchTogetherSuggestionPanel } from "./WatchTogetherSuggestionPanel";
function describeRoomError(error: unknown) {
if (error === "not_found") {
return "Room not found.";
}
if (error === "ended") {
return "That watch party is no longer active.";
}
if (error === "forbidden") {
return "You don't have access to this watch party.";
}
if (error === "host_left" || error === "room_closed") {
return "The room has ended.";
}
@@ -296,20 +316,36 @@ function NowPlayingHero({
);
}
/* ─── Connection status dot ─── */
function StatusDot({ state }: { state: string }) {
const color =
state === "connected"
? "bg-green-400"
: state === "connecting"
? "bg-yellow-400 animate-pulse"
: "bg-red-400";
return <span className={`inline-block h-2 w-2 rounded-full ${color}`} />;
/* ─── Terminal full-body state (room ended, missing token, …) ─── */
function RoomTerminalState({
title,
description,
children,
}: {
title: string;
description: string;
children?: ReactNode;
}) {
return (
<div
role="alert"
className="mx-auto flex w-full max-w-5xl flex-col items-center gap-3 px-6 py-24 text-center"
>
<div className="bg-surface flex size-16 items-center justify-center rounded-2xl border border-white/10">
<Users className="text-muted-foreground size-7" />
</div>
<h1 className="mt-2 text-xl font-semibold tracking-tight">{title}</h1>
<p className="text-muted-foreground max-w-sm text-sm">{description}</p>
<div className="mt-3 flex items-center gap-2">{children}</div>
</div>
);
}
export default function WatchTogetherRoomPage() {
useDocumentTitle("Watch Party");
const { roomId } = useParams<{ roomId: string }>();
const location = useLocation();
const navigate = useNavigate();
const [searchParams] = useSearchParams();
const roomToken = searchParams.get("room_token");
const playbackController = useWatchPlaybackController();
@@ -326,6 +362,8 @@ export default function WatchTogetherRoomPage() {
const [candidate, setCandidate] = useState<BrowseItem | null>(null);
const [candidateContext, setCandidateContext] = useState<string | null>(null);
const [submitting, setSubmitting] = useState(false);
const [endConfirmOpen, setEndConfirmOpen] = useState(false);
const [ending, setEnding] = useState(false);
const [searchStep, setSearchStep] = useState<SearchStep>({ stage: "results" });
const lastAutoStartRevisionRef = useRef<number | null>(null);
const suppressAutoStartSelectionRef = useRef(
@@ -407,30 +445,36 @@ export default function WatchTogetherRoomPage() {
});
}, [playbackController, roomConnection.room, roomId, roomToken]);
const inviteUrl = buildWatchTogetherInviteUrl(roomConnection.room?.invite_path);
const hasInvite = Boolean(roomConnection.room?.invite_path);
const isHost = roomConnection.room?.self_can_manage_room === true;
const isVoteMode = roomConnection.room?.selection_mode === "vote";
const pendingAction: PendingAction = isVoteMode ? "suggest" : "select";
const roomError = !roomId || !roomToken ? "Room token is required." : roomConnection.closedReason;
const handleCopyInvite = useCallback(async () => {
if (!inviteUrl) {
return;
}
await navigator.clipboard.writeText(inviteUrl);
toast.success(`Invite copied. Room code ${roomConnection.room?.code ?? ""}`.trim());
}, [inviteUrl, roomConnection.room?.code]);
await copyWatchTogetherInvite(roomConnection.room?.invite_path, roomConnection.room?.code);
}, [roomConnection.room?.code, roomConnection.room?.invite_path]);
const handleTogglePolicy = useCallback(async () => {
const room = roomConnection.room;
if (!room) {
return;
}
await roomConnection.updatePolicy(
await setWatchTogetherGuestControl(
roomConnection.updatePolicy,
room.guest_control_policy === "guest_play_pause" ? "host_only" : "guest_play_pause",
);
}, [roomConnection]);
const handleEndRoom = useCallback(async () => {
setEnding(true);
try {
await endWatchTogetherRoom(roomConnection.closeRoom);
} finally {
setEnding(false);
setEndConfirmOpen(false);
}
}, [roomConnection]);
const handleConfirmCandidate = useCallback(async () => {
if (!candidate || !roomId) {
return;
@@ -452,6 +496,8 @@ export default function WatchTogetherRoomPage() {
await roomConnection.selectItem({
content_id: candidate.content_id,
});
setCandidate(null);
setCandidateContext(null);
}
} catch (error) {
toast.error(error instanceof Error ? error.message : "Failed to update room");
@@ -519,7 +565,29 @@ export default function WatchTogetherRoomPage() {
}, []);
if (!roomId || !roomToken) {
return <div className="mx-auto max-w-5xl px-6 py-10 text-sm">Room token is required.</div>;
return (
<RoomTerminalState
title="This invite link is incomplete"
description="The link is missing its access token, so the room can't be opened. Ask the host for a fresh invite, or join with a room code instead."
>
<Button type="button" onClick={() => navigate("/rooms/join")}>
Join with a code
</Button>
</RoomTerminalState>
);
}
if (roomConnection.closedReason) {
return (
<RoomTerminalState
title={describeRoomError(roomConnection.closedReason)}
description="Start a new watch party or join another room to keep watching together."
>
<Button type="button" onClick={() => navigate("/rooms/join")}>
Start a new party
</Button>
</RoomTerminalState>
);
}
const isPlaying = roomConnection.room?.phase === "playing";
@@ -538,40 +606,37 @@ export default function WatchTogetherRoomPage() {
Watch Party
</div>
<div className="mt-0.5 text-lg font-semibold tracking-tight">
{roomConnection.room?.code ?? roomId}
{roomConnection.room?.code ?? "…"}
</div>
</div>
<div className="bg-border hidden h-8 w-px sm:block" />
<div className="hidden items-center gap-4 sm:flex">
<div className="flex items-center gap-4">
<div className="flex items-center gap-2 text-sm">
<Users className="text-muted-foreground size-3.5" />
<span>{roomConnection.room?.member_count ?? 0}</span>
</div>
<div className="flex items-center gap-2 text-sm">
<StatusDot state={roomConnection.connectionState} />
<span className="text-muted-foreground">
{roomConnection.connectionState === "connected"
? "Connected"
: roomConnection.connectionState === "connecting"
? "Connecting..."
: "Disconnected"}
<ConnectionStatusDot state={roomConnection.connectionState} />
<span className="text-muted-foreground hidden sm:inline">
<ConnectionStateLabel state={roomConnection.connectionState} />
</span>
</div>
</div>
</div>
<div className="flex items-center gap-2">
{inviteUrl ? (
{hasInvite ? (
<Button
type="button"
variant="outline"
size="sm"
onClick={() => void handleCopyInvite()}
className="gap-1.5"
aria-label="Copy invite link"
>
<Link className="size-3.5" />
<LinkIcon className="size-3.5" />
<span className="hidden sm:inline">Invite</span>
</Button>
) : null}
@@ -595,19 +660,53 @@ export default function WatchTogetherRoomPage() {
type="button"
variant="destructive"
size="sm"
onClick={() => void roomConnection.closeRoom()}
onClick={() => setEndConfirmOpen(true)}
>
End
</Button>
</>
) : null}
{roomConnection.room && !isHost ? (
<Button
type="button"
variant="outline"
size="sm"
onClick={() => navigate("/rooms/join")}
className="gap-1.5"
>
<LogOut className="size-3.5" />
Leave
</Button>
) : null}
</div>
</div>
{/* ─── Error state ─── */}
{roomError ? (
<div className="rounded-xl border border-red-500/20 bg-red-500/10 px-5 py-4 text-sm text-red-300">
{describeRoomError(roomError)}
{/* ─── Participants (optional additive server field) ─── */}
{roomConnection.room?.members?.length ? (
<div className="glass-subtle flex flex-wrap items-center gap-2 rounded-xl px-4 py-3 sm:px-5">
{roomConnection.room.members.map((member, index) => (
<div
key={`${member.profile_id}-${index}`}
className="flex items-center gap-2 rounded-full border border-white/10 bg-white/[0.03] px-3 py-1.5 text-sm"
>
<span
aria-hidden="true"
className={`inline-block h-1.5 w-1.5 rounded-full ${
member.connected ? "bg-emerald-400" : "bg-white/25"
}`}
/>
<span className="font-medium">
{member.display_name}
{member.is_self ? " (you)" : ""}
</span>
<span className="sr-only">{member.connected ? "Connected" : "Disconnected"}</span>
{member.is_host ? (
<span className="text-muted-foreground rounded-full border border-white/10 px-1.5 py-0.5 text-[10px] font-semibold tracking-wide uppercase">
Host
</span>
) : null}
</div>
))}
</div>
) : null}
@@ -649,6 +748,7 @@ export default function WatchTogetherRoomPage() {
placeholder={
isVoteMode ? "Search to suggest something..." : "Search movies and series..."
}
aria-label="Search movies and series"
className="border-border bg-surface placeholder:text-muted-foreground h-12 w-full rounded-xl border py-3 pr-4 pl-11 text-sm shadow-sm transition-all duration-200 outline-none focus:border-white/30 focus:ring-1 focus:ring-white/10"
/>
{query && (
@@ -658,6 +758,7 @@ export default function WatchTogetherRoomPage() {
setQuery("");
setSearchStep({ stage: "results" });
}}
aria-label="Clear search"
className="text-muted-foreground hover:text-foreground absolute top-1/2 right-4 -translate-y-1/2 transition-colors"
>
<X className="size-4" />
@@ -882,6 +983,13 @@ export default function WatchTogetherRoomPage() {
/>
</>
) : null}
<EndWatchPartyDialog
open={endConfirmOpen}
onOpenChange={setEndConfirmOpen}
onConfirm={() => void handleEndRoom()}
isPending={ending}
/>
</div>
);
}
+11 -3
View File
@@ -60,6 +60,12 @@ function SuggestionCard({
type="button"
onClick={onVoteToggle}
disabled={isLoading}
aria-pressed={suggestion.voted_by_me}
aria-label={
suggestion.voted_by_me
? `Remove vote for ${suggestion.title}`
: `Vote for ${suggestion.title}`
}
className={`absolute top-2 right-2 z-10 flex items-center gap-1 rounded-full border px-2 py-1 text-[11px] font-semibold tabular-nums backdrop-blur-sm transition-all duration-200 ${
suggestion.voted_by_me
? "border-primary/40 bg-primary/20 text-primary"
@@ -72,12 +78,13 @@ function SuggestionCard({
{/* Host: Play overlay */}
{isHost ? (
<div className="pointer-events-none absolute inset-0 flex items-center justify-center bg-black/0 opacity-0 transition-all duration-200 group-hover/suggestion:bg-black/40 group-hover/suggestion:opacity-100">
<div className="pointer-events-none absolute inset-0 flex items-center justify-center bg-black/0 opacity-0 transition-all duration-200 group-focus-within/suggestion:bg-black/40 group-focus-within/suggestion:opacity-100 group-hover/suggestion:bg-black/40 group-hover/suggestion:opacity-100">
<button
type="button"
onClick={onPromote}
disabled={isLoading}
className="pointer-events-auto z-10 flex size-10 items-center justify-center rounded-full border border-white/30 bg-white/15 text-white backdrop-blur-sm transition-transform duration-200 group-hover/suggestion:scale-100"
aria-label={`Play ${suggestion.title} for everyone`}
className="pointer-events-auto z-10 flex size-10 items-center justify-center rounded-full border border-white/30 bg-white/15 text-white backdrop-blur-sm transition-transform duration-200 group-hover/suggestion:scale-100 focus-visible:opacity-100"
>
<Play className="size-4 fill-white" />
</button>
@@ -90,7 +97,8 @@ function SuggestionCard({
type="button"
onClick={onDelete}
disabled={isLoading}
className="absolute top-2 left-2 z-10 flex size-6 items-center justify-center rounded-full border border-white/15 bg-black/50 text-white/60 opacity-0 backdrop-blur-sm transition-all duration-200 group-hover/suggestion:opacity-100 hover:border-red-500/40 hover:bg-red-500/20 hover:text-red-300"
aria-label={`Remove suggestion ${suggestion.title}`}
className="absolute top-2 left-2 z-10 flex size-6 items-center justify-center rounded-full border border-white/15 bg-black/50 text-white/60 opacity-0 backdrop-blur-sm transition-all duration-200 group-focus-within/suggestion:opacity-100 group-hover/suggestion:opacity-100 hover:border-red-500/40 hover:bg-red-500/20 hover:text-red-300 focus-visible:opacity-100"
title="Remove suggestion"
>
<Trash2 className="size-3" />
+21 -40
View File
@@ -54,7 +54,11 @@ import type {
SubtitleMode,
} from "../types";
import { toMediaTime, toPlayerTime } from "../utils/mediaTimeline";
import { buildWatchTogetherInviteUrl } from "@/lib/watchTogether";
import {
copyWatchTogetherInvite,
endWatchTogetherRoom,
setWatchTogetherGuestControl,
} from "@/lib/watchTogetherActions";
import { toast } from "sonner";
// Reserved index for the in-progress live AI translation track. Sits well above
@@ -381,6 +385,7 @@ export function VideoPlayer({
});
const roomPlaybackActive = !!watchTogetherRoomId && !watchTogether.closedReason;
const roomSyncWaiting = watchTogether.room?.playback_state === "waiting";
const watchTogetherRoomActive = watchTogether.room !== null;
const showWatchTogetherNotice = useCallback((message: string, tone: "info" | "warning") => {
setNotice({
@@ -1407,10 +1412,7 @@ export function VideoPlayer({
setCurrentTime(toMediaTime(video.currentTime, streamOriginRef.current));
setAwaitingFirstFrame(false);
clearBuffering();
if (
watchTogether.room?.playback_state === "waiting" &&
watchTogetherSync.attachedSessionId === sessionId
) {
if (roomSyncWaiting && watchTogetherSync.attachedSessionId === sessionId) {
watchTogetherSync.reportReady();
}
};
@@ -1438,16 +1440,13 @@ export function VideoPlayer({
bufferingTimerRef.current = null;
}, 500);
}
if (watchTogether.room && watchTogetherSync.attachedSessionId === sessionId) {
if (watchTogetherRoomActive && watchTogetherSync.attachedSessionId === sessionId) {
watchTogetherSync.reportBuffering();
}
};
const onCanPlay = () => {
clearBuffering();
if (
watchTogether.room?.playback_state === "waiting" &&
watchTogetherSync.attachedSessionId === sessionId
) {
if (roomSyncWaiting && watchTogetherSync.attachedSessionId === sessionId) {
watchTogetherSync.reportReady();
}
};
@@ -1456,7 +1455,7 @@ export function VideoPlayer({
setAwaitingFirstFrame(false);
};
const onStalled = () => {
if (watchTogether.room && watchTogetherSync.attachedSessionId === sessionId) {
if (watchTogetherRoomActive && watchTogetherSync.attachedSessionId === sessionId) {
watchTogetherSync.reportBuffering();
}
};
@@ -1507,7 +1506,10 @@ export function VideoPlayer({
video.removeEventListener("error", onError);
video.removeEventListener("ended", onVideoEnded);
};
}, [pendingSeekTime, sessionId, watchTogether.room, watchTogetherSync]); // Listener behavior depends on pending seek reconciliation
// Listener behavior depends on pending seek reconciliation. Watch-together
// deps are intentionally narrowed to the primitives the handlers read so
// room snapshot churn doesn't re-subscribe every listener.
}, [pendingSeekTime, roomSyncWaiting, sessionId, watchTogetherRoomActive, watchTogetherSync]);
// Apply persisted volume on mount (separate from listener effect).
useEffect(() => {
@@ -2035,45 +2037,24 @@ export function VideoPlayer({
}, [onReturnFromPostRoll]);
const handleCopyWatchTogetherInvite = useCallback(async () => {
const inviteUrl = buildWatchTogetherInviteUrl(watchTogether.room?.invite_path);
if (!inviteUrl) {
const copied = await copyWatchTogetherInvite(
watchTogether.room?.invite_path,
watchTogether.room?.code,
);
if (!copied) {
showWatchTogetherNotice("Invite link is not ready yet.", "info");
return;
}
try {
await navigator.clipboard.writeText(inviteUrl);
toast.success(`Invite copied. Room code ${watchTogether.room?.code ?? ""}`.trim());
} catch {
toast.error("Failed to copy invite link");
}
}, [showWatchTogetherNotice, watchTogether.room]);
const handleToggleGuestControl = useCallback(
async (policy: "host_only" | "guest_play_pause") => {
try {
const nextRoom = await watchTogether.updatePolicy(policy);
if (nextRoom) {
toast.success(
nextRoom.guest_control_policy === "guest_play_pause"
? "Guests can now pause and resume"
: "Room is now host controlled",
);
}
} catch (error) {
toast.error(error instanceof Error ? error.message : "Failed to update room");
}
await setWatchTogetherGuestControl(watchTogether.updatePolicy, policy);
},
[watchTogether],
);
const handleEndRoom = useCallback(async () => {
try {
await watchTogether.closeRoom();
toast.success("Room ended");
} catch (error) {
toast.error(error instanceof Error ? error.message : "Failed to end room");
}
await endWatchTogetherRoom(watchTogether.closeRoom);
}, [watchTogether]);
// -- Render --
@@ -1,25 +1,21 @@
import { useState } from "react";
import {
ConnectionStateLabel,
ConnectionStatusDot,
type WatchTogetherConnectionState,
} from "@/components/watchtogether/ConnectionStatusDot";
import { EndWatchPartyDialog } from "@/components/watchtogether/EndWatchPartyDialog";
import type { GuestControlPolicy, WatchTogetherRoomSnapshot } from "@/lib/watchTogether";
interface WatchTogetherPanelProps {
room: WatchTogetherRoomSnapshot | null;
connectionState: "disconnected" | "connecting" | "connected";
connectionState: WatchTogetherConnectionState;
visible: boolean;
onCopyInvite: () => void;
onToggleGuestControl: (policy: GuestControlPolicy) => void;
onEndRoom: () => void;
}
function connectionLabel(connectionState: WatchTogetherPanelProps["connectionState"]) {
switch (connectionState) {
case "connected":
return "Live";
case "connecting":
return "Connecting";
default:
return "Offline";
}
}
export function WatchTogetherPanel({
room,
connectionState,
@@ -28,6 +24,7 @@ export function WatchTogetherPanel({
onToggleGuestControl,
onEndRoom,
}: WatchTogetherPanelProps) {
const [endConfirmOpen, setEndConfirmOpen] = useState(false);
const isHost = room?.self_can_manage_room === true;
const policy = room?.guest_control_policy ?? "host_only";
@@ -41,6 +38,8 @@ export function WatchTogetherPanel({
return (
<div
aria-hidden={!visible}
inert={!visible}
className={`absolute top-4 right-4 z-50 w-56 rounded-xl border border-white/10 bg-black/80 p-3 text-white shadow-2xl backdrop-blur-xl transition-opacity duration-300 ${
visible ? "opacity-100" : "pointer-events-none opacity-0"
}`}
@@ -48,16 +47,10 @@ export function WatchTogetherPanel({
{/* Label + status */}
<div className="flex items-center gap-2 text-[10px] font-semibold tracking-[0.16em] text-white/60 uppercase">
<span>Watch Party</span>
<span
className={`size-1.5 rounded-full ${
connectionState === "connected"
? "bg-emerald-400"
: connectionState === "connecting"
? "bg-amber-300"
: "bg-white/35"
}`}
/>
<span>{connectionLabel(connectionState)}</span>
<ConnectionStatusDot state={connectionState} className="size-1.5" />
<span>
<ConnectionStateLabel state={connectionState} />
</span>
</div>
{/* Code + viewers */}
@@ -85,6 +78,7 @@ export function WatchTogetherPanel({
<button
type="button"
onClick={onCopyInvite}
aria-label="Copy invite link"
className="rounded-md bg-white/12 px-2.5 py-1 text-[11px] font-medium text-white/90 transition-colors hover:bg-white/20"
>
Invite
@@ -100,13 +94,22 @@ export function WatchTogetherPanel({
</button>
<button
type="button"
onClick={onEndRoom}
onClick={() => setEndConfirmOpen(true)}
className="ml-auto rounded-md bg-red-500/25 px-2.5 py-1 text-[11px] font-medium text-red-100 transition-colors hover:bg-red-500/35"
>
End
</button>
</div>
) : null}
<EndWatchPartyDialog
open={endConfirmOpen}
onOpenChange={setEndConfirmOpen}
onConfirm={() => {
setEndConfirmOpen(false);
onEndRoom();
}}
/>
</div>
);
}
@@ -1,4 +1,11 @@
import { useCallback, useEffect, useRef, type MutableRefObject, type RefObject } from "react";
import {
useCallback,
useEffect,
useMemo,
useRef,
type MutableRefObject,
type RefObject,
} from "react";
import { toMediaTime } from "../utils/mediaTimeline";
import type { WatchTogetherRoomConnectionResult } from "./useWatchTogetherRoomConnection";
@@ -39,6 +46,8 @@ export function useWatchTogetherPlaybackSync({
const serverTimeOffsetMs = roomConnection.serverTimeOffsetMs;
const attachedSessionId = room?.attached_session_id ?? null;
const roomConnected = room !== null;
const roomPlaybackState = room?.playback_state ?? null;
const roomPhase = room?.phase ?? null;
const sendRoomMessage = roomConnection.sendRoomMessage;
const waitingStateRef = useRef<"idle" | "buffering" | "ready">("idle");
@@ -126,7 +135,7 @@ export function useWatchTogetherPlaybackSync({
!roomConnected ||
!sessionId ||
attachedSessionId !== sessionId ||
room?.playback_state !== "waiting" ||
roomPlaybackState !== "waiting" ||
waitingStateRef.current === "ready" ||
!video
) {
@@ -150,8 +159,8 @@ export function useWatchTogetherPlaybackSync({
[
attachedSessionId,
connectionState,
room,
roomConnected,
roomPlaybackState,
sendRoomMessage,
sessionId,
streamOriginRef,
@@ -167,7 +176,7 @@ export function useWatchTogetherPlaybackSync({
!roomConnected ||
!sessionId ||
attachedSessionId !== sessionId ||
room?.phase !== "playing" ||
roomPhase !== "playing" ||
waitingStateRef.current === "buffering" ||
!video
) {
@@ -191,8 +200,8 @@ export function useWatchTogetherPlaybackSync({
[
attachedSessionId,
connectionState,
room,
roomConnected,
roomPhase,
sendRoomMessage,
sessionId,
streamOriginRef,
@@ -200,10 +209,15 @@ export function useWatchTogetherPlaybackSync({
],
);
return {
attachedSessionId,
requestTransport,
reportReady,
reportBuffering,
};
// Stable identity so consumers (e.g. VideoPlayer's video-event-listener
// effect) don't re-run on every room snapshot.
return useMemo(
() => ({
attachedSessionId,
requestTransport,
reportReady,
reportBuffering,
}),
[attachedSessionId, requestTransport, reportReady, reportBuffering],
);
}
@@ -1,345 +0,0 @@
import {
useCallback,
useEffect,
useMemo,
useRef,
useState,
type MutableRefObject,
type RefObject,
} from "react";
import { getProfileToken } from "@/api/client";
import {
closeWatchTogetherRoom,
type GuestControlPolicy,
type WatchTogetherRoomSnapshot,
updateWatchTogetherRoomPolicy,
} from "@/lib/watchTogether";
import { usePlayerConfig } from "../context/PlayerConfigContext";
import { toMediaTime } from "../utils/mediaTimeline";
type ConnectionState = "disconnected" | "connecting" | "connected";
interface UseWatchTogetherRoomOptions {
roomId?: string | null;
roomToken?: string | null;
sessionId?: string | null;
videoRef: RefObject<HTMLVideoElement | null>;
streamOriginRef: MutableRefObject<number>;
playbackRealtimeConnected?: boolean;
}
interface UseWatchTogetherRoomResult {
connectionState: ConnectionState;
room: WatchTogetherRoomSnapshot | null;
closedReason: string | null;
requestTransport: (
action: "play" | "pause" | "seek",
positionSeconds: number,
isPaused: boolean,
) => void;
updatePolicy: (policy: GuestControlPolicy) => Promise<WatchTogetherRoomSnapshot | null>;
closeRoom: () => Promise<void>;
}
const reconnectDelays = [500, 1_000, 2_000, 5_000];
const stateReportIntervalMs = 1_500;
function buildRoomWebSocketUrl(
apiBaseUrl: string,
roomId: string,
roomToken: string,
accessToken: string | null,
profileId: string | null,
profileToken: string | null,
) {
if (typeof window === "undefined") {
return "";
}
const apiBase = apiBaseUrl.startsWith("http")
? apiBaseUrl
: new URL(apiBaseUrl, window.location.origin).toString();
const wsBase = apiBase.replace(/^http/, "ws");
const url = new URL(`${wsBase}/watch-together/rooms/${roomId}/ws`);
url.searchParams.set("room_token", roomToken);
if (accessToken) {
url.searchParams.set("token", accessToken);
}
if (profileId) {
url.searchParams.set("profile_id", profileId);
}
if (profileToken) {
url.searchParams.set("profile_token", profileToken);
}
return url.toString();
}
export function useWatchTogetherRoom({
roomId,
roomToken,
sessionId,
videoRef,
streamOriginRef,
playbackRealtimeConnected,
}: UseWatchTogetherRoomOptions): UseWatchTogetherRoomResult {
const config = usePlayerConfig();
const stateRoomId = roomId ?? null;
const [activeRoomId, setActiveRoomId] = useState<string | null>(stateRoomId);
const [connectionStateValue, setConnectionState] = useState<ConnectionState>("disconnected");
const [roomValue, setRoom] = useState<WatchTogetherRoomSnapshot | null>(null);
const [closedReasonValue, setClosedReason] = useState<string | null>(null);
const socketRef = useRef<WebSocket | null>(null);
const roomRef = useRef<WatchTogetherRoomSnapshot | null>(null);
const closedReasonRef = useRef<string | null>(null);
const playbackRealtimeConnectedRef = useRef(playbackRealtimeConnected ?? true);
const connectionState = activeRoomId === stateRoomId ? connectionStateValue : "disconnected";
const room = activeRoomId === stateRoomId ? roomValue : null;
const closedReason = activeRoomId === stateRoomId ? closedReasonValue : null;
useEffect(() => {
roomRef.current = room;
}, [room]);
useEffect(() => {
closedReasonRef.current = closedReason;
}, [closedReason]);
useEffect(() => {
playbackRealtimeConnectedRef.current = playbackRealtimeConnected ?? true;
}, [playbackRealtimeConnected]);
const websocketUrl = useMemo(() => {
if (!roomId || !roomToken) {
return null;
}
return buildRoomWebSocketUrl(
config.apiBaseUrl,
roomId,
roomToken,
config.getAccessToken(),
config.getProfileId(),
getProfileToken(),
);
}, [config, roomId, roomToken]);
const sendMessage = useCallback((message: Record<string, unknown>) => {
const socket = socketRef.current;
if (!socket || socket.readyState !== WebSocket.OPEN) {
return false;
}
socket.send(JSON.stringify(message));
return true;
}, []);
useEffect(() => {
if (!roomId || !roomToken || !websocketUrl) {
return;
}
let disposed = false;
let attempt = 0;
let reconnectTimer: number | null = null;
const scheduleReconnect = () => {
if (disposed || closedReasonRef.current) {
return;
}
const delay = reconnectDelays[Math.min(attempt, reconnectDelays.length - 1)];
attempt += 1;
reconnectTimer = window.setTimeout(connect, delay);
};
const connect = () => {
if (disposed) {
return;
}
setConnectionState("connecting");
let socket: WebSocket;
try {
socket = new WebSocket(websocketUrl);
} catch {
scheduleReconnect();
return;
}
socketRef.current = socket;
socket.addEventListener("open", () => {
if (disposed) {
socket.close();
return;
}
attempt = 0;
setActiveRoomId(stateRoomId);
setConnectionState("connected");
if (sessionId && playbackRealtimeConnectedRef.current) {
socket.send(JSON.stringify({ type: "attach_session", session_id: sessionId }));
}
});
socket.addEventListener("message", (event) => {
let message: Record<string, unknown>;
try {
message = JSON.parse(String(event.data)) as Record<string, unknown>;
} catch {
return;
}
switch (message.type) {
case "snapshot": {
const payload = message.room as WatchTogetherRoomSnapshot | undefined;
if (payload) {
setActiveRoomId(stateRoomId);
roomRef.current = payload;
setRoom(payload);
}
return;
}
case "room_closed":
setActiveRoomId(stateRoomId);
roomRef.current = null;
closedReasonRef.current =
typeof message.reason === "string" ? message.reason : "room_closed";
setRoom(null);
setClosedReason(closedReasonRef.current);
socket.close();
return;
case "pong":
return;
default:
if (message.type === "error") {
console.warn("[watch-together]", message.code ?? "error", message.message ?? "");
}
}
});
socket.addEventListener("close", () => {
setActiveRoomId(stateRoomId);
setConnectionState("disconnected");
if (socketRef.current === socket) {
socketRef.current = null;
}
scheduleReconnect();
});
socket.addEventListener("error", () => {
socket.close();
});
};
connect();
return () => {
disposed = true;
if (reconnectTimer !== null) {
window.clearTimeout(reconnectTimer);
}
if (socketRef.current) {
const socket = socketRef.current;
socketRef.current = null;
if (socket.readyState === WebSocket.OPEN || socket.readyState === WebSocket.CONNECTING) {
socket.close();
}
}
};
}, [roomId, roomToken, sessionId, stateRoomId, websocketUrl]);
useEffect(() => {
if (
!roomId ||
!sessionId ||
connectionState !== "connected" ||
playbackRealtimeConnected === false
) {
return;
}
sendMessage({ type: "attach_session", session_id: sessionId });
}, [connectionState, playbackRealtimeConnected, roomId, sendMessage, sessionId]);
useEffect(() => {
if (
!roomId ||
!sessionId ||
connectionState !== "connected" ||
playbackRealtimeConnected === false
) {
return;
}
const intervalId = window.setInterval(() => {
const currentRoom = roomRef.current;
const video = videoRef.current;
if (!video || currentRoom?.attached_session_id !== sessionId) {
return;
}
sendMessage({
type: "state_report",
session_id: sessionId,
position_seconds: toMediaTime(video.currentTime, streamOriginRef.current),
is_paused: video.paused,
});
}, stateReportIntervalMs);
return () => {
window.clearInterval(intervalId);
};
}, [
connectionState,
playbackRealtimeConnected,
roomId,
sendMessage,
sessionId,
streamOriginRef,
videoRef,
]);
const requestTransport = useCallback(
(action: "play" | "pause" | "seek", positionSeconds: number, isPaused: boolean) => {
if (!roomId) {
return;
}
sendMessage({
type: "transport_request",
action,
position_seconds: positionSeconds,
is_paused: isPaused,
});
},
[roomId, sendMessage],
);
const updatePolicy = useCallback(
async (policy: GuestControlPolicy) => {
if (!roomId) {
return null;
}
const response = await updateWatchTogetherRoomPolicy(roomId, policy);
setRoom(response.room);
return response.room;
},
[roomId],
);
const closeRoom = useCallback(async () => {
if (!roomId) {
return;
}
await closeWatchTogetherRoom(roomId);
}, [roomId]);
return {
connectionState,
room,
closedReason,
requestTransport,
updatePolicy,
closeRoom,
};
}
@@ -1,5 +1,5 @@
import { useCallback, useEffect, useRef, useState } from "react";
import { getAccessToken, getProfileToken } from "@/api/client";
import { ApiClientError, getAccessToken, getProfileToken } from "@/api/client";
import {
closeWatchTogetherRoom,
type CreateWatchTogetherSuggestionInput,
@@ -54,6 +54,26 @@ export interface WatchTogetherRoomConnectionResult {
const reconnectDelays = [500, 1_000, 2_000, 5_000];
const pingIntervalMs = 15_000;
/**
* Maps a room-fetch failure to a terminal closed reason. Returns null for
* transient errors (network, 5xx) so the reconnect loop keeps retrying.
*/
function closedReasonFromRoomFetchError(error: unknown): string | null {
if (!(error instanceof ApiClientError)) {
return null;
}
switch (error.status) {
case 404:
return "not_found";
case 410:
return "ended";
case 403:
return "forbidden";
default:
return null;
}
}
function buildRoomWebSocketUrl(
apiBaseUrl: string,
roomId: string,
@@ -104,6 +124,13 @@ export function useWatchTogetherRoomConnection({
closedReasonRef.current = closedReason;
}, [closedReason]);
// Sets the ref synchronously (before the state commit) so the reconnect
// loop can never fire between setState and the ref-syncing effect above.
const markClosed = useCallback((reason: string) => {
closedReasonRef.current = reason;
setClosedReason(reason);
}, []);
const sendRoomMessage = useCallback((message: Record<string, unknown>) => {
const socket = socketRef.current;
if (!socket || socket.readyState !== WebSocket.OPEN) {
@@ -140,7 +167,15 @@ export function useWatchTogetherRoomConnection({
}
setRoom(response.room);
})
.catch(() => {});
.catch((error: unknown) => {
if (cancelled) {
return;
}
const reason = closedReasonFromRoomFetchError(error);
if (reason) {
markClosed(reason);
}
});
void listWatchTogetherSuggestions(roomId, roomToken)
.then((response) => {
@@ -155,7 +190,7 @@ export function useWatchTogetherRoomConnection({
cancelled = true;
window.clearTimeout(resetClosedReasonTimer);
};
}, [roomId, roomToken]);
}, [markClosed, roomId, roomToken]);
useEffect(() => {
if (!roomId || !roomToken) {
@@ -259,7 +294,7 @@ export function useWatchTogetherRoomConnection({
}
case "room_closed":
setRoom(null);
setClosedReason(typeof message.reason === "string" ? message.reason : "room_closed");
markClosed(typeof message.reason === "string" ? message.reason : "room_closed");
socket.close();
return;
case "transport_command": {
@@ -289,10 +324,19 @@ export function useWatchTogetherRoomConnection({
}
}
return;
default:
if (message.type === "error") {
case "error": {
// Defense in depth: some server paths report terminal conditions
// as an error message instead of (or before) room_closed.
if (message.code === "not_found" || message.code === "gone") {
setRoom(null);
markClosed(message.code === "gone" ? "ended" : "not_found");
socket.close();
} else {
console.warn("[watch-together]", message.code ?? "error", message.message ?? "");
}
return;
}
default:
}
});
@@ -332,7 +376,7 @@ export function useWatchTogetherRoomConnection({
socket.close();
}
};
}, [roomId, roomToken]);
}, [markClosed, roomId, roomToken]);
const updatePolicy = useCallback(
async (policy: GuestControlPolicy) => {