* 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>
317 lines
7.5 KiB
Go
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
|
|
}
|