Files
silo-server/internal/watchtogether/repository.go
b250dbb59b 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>
2026-07-02 11:34:02 -04:00

317 lines
7.5 KiB
Go

package watchtogether
import (
"context"
"errors"
"fmt"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
var ErrRoomNotFound = errors.New("watch together room not found")
var ErrRoomStateConflict = errors.New("watch together room state conflict")
type Repository struct {
pool *pgxpool.Pool
}
func NewRepository(pool *pgxpool.Pool) *Repository {
return &Repository{pool: pool}
}
func (r *Repository) CreateRoom(ctx context.Context, room Room) (*Room, error) {
if r == nil || r.pool == nil {
return nil, fmt.Errorf("watch together repository unavailable")
}
const query = `
INSERT INTO watch_together_rooms (
id, code, join_token, host_user_id, host_profile_id,
phase, playback_state, resume_on_ready, selection_mode, selection_revision,
selected_content_id, selected_file_id, selected_library_id,
guest_control_policy,
anchor_position_seconds, is_paused, anchor_updated_at,
generation, created_at, closed_at
) VALUES (
$1, $2, $3, $4, $5,
$6, $7, $8, $9, $10,
$11, $12, $13,
$14,
$15, $16, $17,
$18, $19, $20
)
RETURNING ` + roomColumns
created, err := scanRoom(r.pool.QueryRow(
ctx,
query,
room.ID,
room.Code,
room.JoinToken,
room.HostUserID,
room.HostProfileID,
room.Phase,
room.PlaybackState,
room.ResumeOnReady,
room.SelectionMode,
room.SelectionRevision,
room.SelectedContentID,
room.SelectedFileID,
room.SelectedLibraryID,
room.GuestControlPolicy,
room.AnchorPositionSeconds,
room.IsPaused,
room.AnchorUpdatedAt,
room.Generation,
room.CreatedAt,
room.ClosedAt,
))
if err != nil {
return nil, err
}
return created, nil
}
func (r *Repository) GetRoomByID(ctx context.Context, roomID string) (*Room, error) {
return r.getRoom(ctx, `SELECT `+roomColumns+` FROM watch_together_rooms WHERE id = $1`, roomID)
}
func (r *Repository) GetRoomByCode(ctx context.Context, code string) (*Room, error) {
return r.getRoom(
ctx,
`SELECT `+roomColumns+` FROM watch_together_rooms WHERE LOWER(code) = LOWER($1)`,
strings.TrimSpace(code),
)
}
func (r *Repository) GetRoomByJoinToken(ctx context.Context, joinToken string) (*Room, error) {
return r.getRoom(
ctx,
`SELECT `+roomColumns+` FROM watch_together_rooms WHERE join_token = $1`,
strings.TrimSpace(joinToken),
)
}
func (r *Repository) UpdatePolicy(
ctx context.Context,
roomID string,
policy GuestControlPolicy,
generation int64,
expectedGeneration int64,
) (*Room, error) {
const query = `
UPDATE watch_together_rooms
SET guest_control_policy = $2,
generation = $3
WHERE id = $1
AND generation = $4
AND phase <> 'ended'
RETURNING ` + roomColumns
return r.scanConditionalUpdate(ctx, query, roomID, policy, generation, expectedGeneration)
}
func (r *Repository) UpdateAnchor(
ctx context.Context,
roomID string,
positionSeconds float64,
isPaused bool,
playbackState RoomPlaybackState,
resumeOnReady bool,
anchorUpdatedAt time.Time,
generation int64,
expectedGeneration int64,
) (*Room, error) {
const query = `
UPDATE watch_together_rooms
SET anchor_position_seconds = $2,
is_paused = $3,
playback_state = $4,
resume_on_ready = $5,
anchor_updated_at = $6,
generation = $7
WHERE id = $1
AND generation = $8
AND phase = 'playing'
RETURNING ` + roomColumns
return r.scanConditionalUpdate(
ctx,
query,
roomID,
positionSeconds,
isPaused,
playbackState,
resumeOnReady,
anchorUpdatedAt.UTC(),
generation,
expectedGeneration,
)
}
// 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
SET phase = $2,
closed_at = $3
WHERE id = $1
AND phase <> 'ended'
RETURNING ` + roomColumns
return r.scanConditionalUpdate(ctx, query, roomID, RoomPhaseEnded, closedAt.UTC())
}
func (r *Repository) UpdateSelection(
ctx context.Context,
roomID string,
selection SelectItemInput,
phase RoomPhase,
playbackState RoomPlaybackState,
resumeOnReady bool,
anchorPosition float64,
isPaused bool,
anchorUpdatedAt time.Time,
selectionRevision int64,
generation int64,
expectedGeneration int64,
) (*Room, error) {
const query = `
UPDATE watch_together_rooms
SET phase = $2,
playback_state = $3,
resume_on_ready = $4,
selected_content_id = $5,
selected_file_id = $6,
selected_library_id = $7,
anchor_position_seconds = $8,
is_paused = $9,
anchor_updated_at = $10,
selection_revision = $11,
generation = $12
WHERE id = $1
AND generation = $13
AND phase <> 'ended'
RETURNING ` + roomColumns
return r.scanConditionalUpdate(
ctx,
query,
roomID,
phase,
playbackState,
resumeOnReady,
selection.ContentID,
selection.FileID,
selection.LibraryID,
anchorPosition,
isPaused,
anchorUpdatedAt.UTC(),
selectionRevision,
generation,
expectedGeneration,
)
}
const roomColumns = `
id, code, join_token, host_user_id, host_profile_id,
phase, playback_state, resume_on_ready, selection_mode, selection_revision,
selected_content_id, selected_file_id, selected_library_id,
guest_control_policy,
anchor_position_seconds, is_paused, anchor_updated_at,
generation, created_at, closed_at
`
func (r *Repository) getRoom(ctx context.Context, query string, arg string) (*Room, error) {
if r == nil || r.pool == nil {
return nil, fmt.Errorf("watch together repository unavailable")
}
return scanRoom(r.pool.QueryRow(ctx, query, arg))
}
func scanRoom(row pgx.Row) (*Room, error) {
var room Room
if err := row.Scan(
&room.ID,
&room.Code,
&room.JoinToken,
&room.HostUserID,
&room.HostProfileID,
&room.Phase,
&room.PlaybackState,
&room.ResumeOnReady,
&room.SelectionMode,
&room.SelectionRevision,
&room.SelectedContentID,
&room.SelectedFileID,
&room.SelectedLibraryID,
&room.GuestControlPolicy,
&room.AnchorPositionSeconds,
&room.IsPaused,
&room.AnchorUpdatedAt,
&room.Generation,
&room.CreatedAt,
&room.ClosedAt,
); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrRoomNotFound
}
return nil, fmt.Errorf("scan watch together room: %w", err)
}
return &room, nil
}
func (r *Repository) scanConditionalUpdate(ctx context.Context, query string, args ...any) (*Room, error) {
if r == nil || r.pool == nil {
return nil, fmt.Errorf("watch together repository unavailable")
}
room, err := scanRoom(r.pool.QueryRow(ctx, query, args...))
if !errors.Is(err, ErrRoomNotFound) {
return room, err
}
roomID, _ := args[0].(string)
existing, lookupErr := r.GetRoomByID(ctx, roomID)
if lookupErr != nil {
return nil, lookupErr
}
if existing.Phase == RoomPhaseEnded {
return nil, ErrRoomClosed
}
return nil, ErrRoomStateConflict
}