diff --git a/js/app/components/name-color-picker/name-color-picker.tsx b/js/app/components/name-color-picker/name-color-picker.tsx index 40f20b100..966b7cf74 100644 --- a/js/app/components/name-color-picker/name-color-picker.tsx +++ b/js/app/components/name-color-picker/name-color-picker.tsx @@ -1,4 +1,9 @@ -import { Button, formatHandleWithAt, zero } from "@streamplace/components"; +import { + Button, + formatHandleWithAt, + Text, + zero, +} from "@streamplace/components"; import { Palette, SwatchBook, X } from "lucide-react-native"; import { useEffect, useState } from "react"; import { @@ -6,24 +11,52 @@ import { Modal, Platform, Pressable, - Text, TouchableOpacity, View, } from "react-native"; import ColorPicker, { HueSlider, + InputWidget, Panel1, - Preview, Swatches, } from "reanimated-color-picker"; import { useStore } from "store"; import { useChatProfile, useUserProfile } from "store/hooks"; import { PlaceStreamChatProfile } from "streamplace"; +/** + * Returns black or white depending on which contrasts better against the given color + */ +function getContrastColor(color: string): string { + let r: number, g: number, b: number; + + if (color.startsWith("#")) { + const hex = color.replace("#", ""); + r = parseInt(hex.substring(0, 2), 16); + g = parseInt(hex.substring(2, 4), 16); + b = parseInt(hex.substring(4, 6), 16); + } else if (color.startsWith("rgb")) { + const match = color.match(/(\d+),\s*(\d+),\s*(\d+)/); + if (match) { + r = parseInt(match[1]); + g = parseInt(match[2]); + b = parseInt(match[3]); + } else { + return "#fff"; + } + } else { + return "#fff"; + } + + const luminance = (0.299 * r + 0.587 * g + 0.114 * b) / 255; + return luminance > 0.5 ? "#000" : "#fff"; +} + /** * Parses an RGB color string and returns an object with red, green, and blue values */ function parseRgbString(rgbString: string): PlaceStreamChatProfile.Color { + console.log(rgbString); if ( !rgbString || (!rgbString.startsWith("rgb(") && !rgbString.startsWith("rgba(")) @@ -45,6 +78,25 @@ function parseRgbString(rgbString: string): PlaceStreamChatProfile.Color { }; } +function rgbToHex(rgb: PlaceStreamChatProfile.Color) { + const hex = ( + (1 << 24) + + (rgb.red << 16) + + (rgb.green << 8) + + rgb.blue + ).toString(16); + return `#${hex.slice(-6)}`; +} + +// rgb(r, g, b) to hex +function cssRgbToHex(rgb: string) { + if (rgb.startsWith("#")) { + return rgb; + } + const parsed = parseRgbString(rgb); + return rgbToHex(parsed); +} + export function useNameColorPicker() { const [modalVisible, setModalVisible] = useState(false); const [tempColor, setTempColor] = useState("#bd6e86"); @@ -159,32 +211,77 @@ export function useNameColorPicker() { zero.bg.gray[800], zero.r.md, zero.p[3], - zero.mb[5], + zero.mb[3], zero.layout.flex.alignCenter, ]} > - - @{profile.handle} - - - Preview + + {formatHandleWithAt(profile)} )} {/* Color Picker */} - + setTempColor(result.rgb)} + onChangeJS={(result) => setTempColor(result.rgb)} > - - + + {cssRgbToHex(currentColor) !== cssRgbToHex(tempColor) && ( + + + {cssRgbToHex(currentColor)} + + + )} + + + @@ -250,47 +347,7 @@ export default function NameColorPicker({ text?: (color: string) => React.ReactNode; buttonProps?: any; }) { - const [modalVisible, setModalVisible] = useState(false); - const [tempColor, setTempColor] = useState("#bd6e86"); - const createChatProfileRecord = useStore( - (state) => state.createChatProfileRecord, - ); - const getChatProfileRecordFromPDS = useStore( - (state) => state.getChatProfileRecordFromPDS, - ); - const chatProfile = useChatProfile(); - const profile = useUserProfile(); - const isWeb = Platform.OS === "web"; - - const currentColor = chatProfile?.profile?.color - ? `rgb(${chatProfile.profile.color.red}, ${chatProfile.profile.color.green}, ${chatProfile.profile.color.blue})` - : "#bd6e86"; - - useEffect(() => { - if (profile?.did && !chatProfile?.profile) { - getChatProfileRecordFromPDS(); - } - setTempColor(currentColor); - }, [profile?.did, chatProfile?.profile?.color, currentColor]); - - const handleOpenModal = () => { - if (!isWeb) { - Keyboard.dismiss(); - } - setTempColor(currentColor); - setModalVisible(true); - }; - - const handleCloseModal = () => { - setModalVisible(false); - setTempColor(currentColor); // Reset to current color on cancel - }; - - const handleSaveColor = () => { - setModalVisible(false); - const parsed = parseRgbString(tempColor); - createChatProfileRecord(parsed.red, parsed.green, parsed.blue); - }; + const { currentColor, openModal, modal } = useNameColorPicker(); return ( @@ -298,7 +355,7 @@ export default function NameColorPicker({ variant="secondary" leftIcon={} style={[buttonProps?.style]} - onPress={handleOpenModal} + onPress={openModal} {...buttonProps} > @@ -306,149 +363,7 @@ export default function NameColorPicker({ - - - e.stopPropagation()} - > - {/* Header */} - - - - - Choose Color - - - - - - - - {/* User Preview */} - {profile?.handle && ( - - - {formatHandleWithAt(profile)} - - - Preview - - - )} - - {/* Color Picker */} - - setTempColor(result.rgb)} - > - - - - - - - - - - - - - - - - {/* Actions */} - - - - Cancel - - - - - Save Color - - - - - - + {modal} {children} -- 2.51.2 From 01bbebaf34c8d503bb52c33afdcddbdd100051c0 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Sun, 31 May 2026 15:28:30 -0500 Subject: [PATCH 02/28] feat: try delegating pins --- js/app/src/screens/chat-popout.tsx | 5 +- .../components/dashboard/moderator-panel.tsx | 47 +++++++++++++++++++ .../src/content/docs/features/chat-popout.md | 13 +++++ pkg/spxrpc/place_stream_moderation.go | 19 -------- 4 files changed, 63 insertions(+), 21 deletions(-) diff --git a/js/app/src/screens/chat-popout.tsx b/js/app/src/screens/chat-popout.tsx index 382da0b51..49622165c 100644 --- a/js/app/src/screens/chat-popout.tsx +++ b/js/app/src/screens/chat-popout.tsx @@ -22,6 +22,7 @@ interface ChatPopoutParams { reverse?: string; hideAfter?: string; hideChatBox?: string; + hidePinnedComments?: string; showNotifications?: string; } @@ -65,7 +66,7 @@ export function PopoutChatInner({ params }: { params: ChatPopoutParams }) { ? parseInt(params.hideAfter, 10) : undefined; const hideChatBox = params.hideChatBox === "true"; - const showNotifications = params.showNotifications === "true"; + const hidePinnedComments = params.hidePinnedComments === "true"; useEffect(() => { setSrc(params.user); @@ -88,7 +89,7 @@ export function PopoutChatInner({ params }: { params: ChatPopoutParams }) { { maxHeight: "100vh" }, ]} > - {showNotifications ? ( + {!hidePinnedComments ? ( {chat} diff --git a/js/components/src/components/dashboard/moderator-panel.tsx b/js/components/src/components/dashboard/moderator-panel.tsx index f13ce7d26..4366cb668 100644 --- a/js/components/src/components/dashboard/moderator-panel.tsx +++ b/js/components/src/components/dashboard/moderator-panel.tsx @@ -596,6 +596,53 @@ function AddModeratorDialog({ + + + setPermissions((p) => ({ + ...p, + "message.pin": !p["message.pin"], + })) + } + style={[ + layout.flex.row, + layout.flex.alignCenter, + p[3], + r.md, + bg.neutral[800], + borders.width.thin, + borders.color.neutral[700], + ]} + > + + {permissions["message.pin"] && ( + ✓ + )} + + + + Pin Messages + + + Pin and unpin chat messages + + + diff --git a/js/docs/src/content/docs/features/chat-popout.md b/js/docs/src/content/docs/features/chat-popout.md index b8b120d03..d071e85aa 100644 --- a/js/docs/src/content/docs/features/chat-popout.md +++ b/js/docs/src/content/docs/features/chat-popout.md @@ -63,3 +63,16 @@ Example: ``` https://stream.place/chat-popout/did:plc:abcdef123456...?hideChatBox=true ``` + +### `hidePinnedComments` + +Hides pinned comment notifications in the popout. + +- Type: `boolean` (values: `"true"`, `"false"`) +- Default: `"false"` + +Example: + +``` +https://stream.place/chat-popout/did:plc:abcdef123456...?hidePinnedComments=true +``` diff --git a/pkg/spxrpc/place_stream_moderation.go b/pkg/spxrpc/place_stream_moderation.go index f621f8e31..a70ff81dd 100644 --- a/pkg/spxrpc/place_stream_moderation.go +++ b/pkg/spxrpc/place_stream_moderation.go @@ -376,25 +376,6 @@ func (s *Server) handlePlaceStreamModerationCreatePin(ctx context.Context, input return nil, err } - // Check that the streamer has an active livestream - ls, err := s.model.GetLatestLivestreamForRepo(input.Streamer) - if err != nil { - log.Error(ctx, "failed to get livestream for streamer", "err", err) - return nil, echo.NewHTTPError(http.StatusInternalServerError, "failed to check livestream status") - } - if ls == nil || ls.Livestream == nil { - return nil, echo.NewHTTPError(http.StatusBadRequest, "no active livestream for this streamer") - } - lsRec, err := lexutil.CborDecodeValue(*ls.Livestream) - if err != nil { - log.Error(ctx, "failed to decode livestream record", "err", err) - return nil, echo.NewHTTPError(http.StatusInternalServerError, "failed to check livestream status") - } - lsTyped, ok := lsRec.(*streamplace.Livestream) - if ok && lsTyped.EndedAt != nil { - return nil, echo.NewHTTPError(http.StatusBadRequest, "livestream has ended, cannot pin comments") - } - // Delete existing pinned records for this streamer (single-pin semantics) var listOutput comatproto.RepoListRecords_Output err = modCtx.StreamerClient.Do(ctx, xrpc.Query, "application/json", "com.atproto.repo.listRecords", -- 2.51.2 From 6668993348bdf50efe0eec90225fbbc1a56640a1 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Sun, 31 May 2026 18:34:57 -0500 Subject: [PATCH 03/28] Allow multiple pinned messages in chat Remove the logic that automatically deletes existing pinned records when a new one is created, allowing historical pins to persist. --- js/components/src/livestream-store/chat.tsx | 16 ------------ pkg/spxrpc/place_stream_moderation.go | 29 +-------------------- 2 files changed, 1 insertion(+), 44 deletions(-) diff --git a/js/components/src/livestream-store/chat.tsx b/js/components/src/livestream-store/chat.tsx index da6a63ccd..9949c753c 100644 --- a/js/components/src/livestream-store/chat.tsx +++ b/js/components/src/livestream-store/chat.tsx @@ -479,22 +479,6 @@ export const usePinChatMessage = () => { // If streamer, create directly if (agent.did === streamerDID) { - // First delete any existing pinned records - const listResult = await agent.com.atproto.repo.listRecords({ - repo: streamerDID, - collection: "place.stream.chat.pinnedRecord", - }); - for (const rec of listResult.data.records) { - const rkey = rec.uri.split("/").pop(); - if (rkey) { - await agent.com.atproto.repo.deleteRecord({ - repo: streamerDID, - collection: "place.stream.chat.pinnedRecord", - rkey, - }); - } - } - const record = { $type: "place.stream.chat.pinnedRecord", pinnedMessage: messageUri, diff --git a/pkg/spxrpc/place_stream_moderation.go b/pkg/spxrpc/place_stream_moderation.go index a70ff81dd..2a806e5c2 100644 --- a/pkg/spxrpc/place_stream_moderation.go +++ b/pkg/spxrpc/place_stream_moderation.go @@ -376,34 +376,7 @@ func (s *Server) handlePlaceStreamModerationCreatePin(ctx context.Context, input return nil, err } - // Delete existing pinned records for this streamer (single-pin semantics) - var listOutput comatproto.RepoListRecords_Output - err = modCtx.StreamerClient.Do(ctx, xrpc.Query, "application/json", "com.atproto.repo.listRecords", - map[string]any{ - "repo": input.Streamer, - "collection": constants.PLACE_STREAM_CHAT_PINNED_RECORD, - }, nil, &listOutput) - if err != nil { - log.Error(ctx, "failed to list existing pinned records", "err", err) - } else { - for _, rec := range listOutput.Records { - rkey, delErr := extractRKey(rec.Uri) - if delErr != nil { - continue - } - deleteInput := comatproto.RepoDeleteRecord_Input{ - Collection: constants.PLACE_STREAM_CHAT_PINNED_RECORD, - Rkey: rkey, - Repo: input.Streamer, - } - deleteOutput := comatproto.RepoDeleteRecord_Output{} - if err := modCtx.StreamerClient.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.deleteRecord", map[string]any{}, deleteInput, &deleteOutput); err != nil { - log.Error(ctx, "failed to delete existing pinned record", "rkey", rkey, "err", err) - } - } - } - - // Create the pinned record + // Create the pinned record (old pins persist as history) pinnedRecord := &streamplace.ChatPinnedRecord{ LexiconTypeID: "place.stream.chat.pinnedRecord", PinnedMessage: input.MessageUri, -- 2.51.2 From ebabf5e7a96c208d4964e5e496fe0376f59c1330 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 2 Jun 2026 15:08:50 -0700 Subject: [PATCH 04/28] vod: add login from /upload --- js/app/src/router.tsx | 17 ++++++++++-- js/app/src/screens/upload.tsx | 49 +++++++++++++++++++++++++++++++++++ 2 files changed, 64 insertions(+), 2 deletions(-) diff --git a/js/app/src/router.tsx b/js/app/src/router.tsx index 743d1afbb..69c150318 100644 --- a/js/app/src/router.tsx +++ b/js/app/src/router.tsx @@ -275,6 +275,7 @@ export const UploadButton = () => { const isCompact = windowWidth <= 800; const { status: betaStatus, loading: betaLoading } = useBetaStatus("vod"); + const navigation = useNavigation(); if (!did) return null; if (betaLoading) return null; @@ -294,7 +295,13 @@ export const UploadButton = () => { to={{ screen: "HomeTab", params: { screen: "Upload" } }} style={{ marginRight: 10 }} > - @@ -317,7 +324,13 @@ export const UploadButton = () => { to={{ screen: "HomeTab", params: { screen: "Upload" } }} style={{ marginRight: 10 }} > - diff --git a/js/app/src/screens/upload.tsx b/js/app/src/screens/upload.tsx index ba8aae3d4..70e196a89 100644 --- a/js/app/src/screens/upload.tsx +++ b/js/app/src/screens/upload.tsx @@ -1,3 +1,4 @@ +import { useRoute } from "@react-navigation/native"; import { Admonition, Button, @@ -40,6 +41,8 @@ import { TextInput, useWindowDimensions, } from "react-native"; +import { useStore } from "store"; +import { useIsReady, useUserProfile } from "store/hooks"; import type { PlaceStreamLivestream, PlaceStreamVideo } from "streamplace"; import * as tus from "tus-js-client"; @@ -106,6 +109,10 @@ const BETA_FEATURE = "vod"; export default function UploadScreen() { const agent = usePDSAgent(); + const isReady = useIsReady(); + const userProfile = useUserProfile(); + const openLoginModal = useStore((state) => state.openLoginModal); + const route = useRoute(); const { status: betaStatus, loading: betaLoading, @@ -595,6 +602,48 @@ export default function UploadScreen() { // ── render ──────────────────────────────────────────────────────────────── + // Login gate: uploading needs an account. Wait for auth to resolve, then + // prompt logged-out users with a login button that returns them here. + if (!isReady) { + return ; + } + if (!userProfile) { + return ( + + + + Log in to upload videos + + + You need to be logged in to upload and manage your videos. + + + + + ); + } + // Beta gate: hold the upload UI until we know the account's access status. // While we're still resolving it for a logged-in user, show a spinner; if // they're not granted, swap the whole page for the request-access flow. -- 2.51.2 From 0c963c216cfcade19b4189d240578808a3b36b84 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 2 Jun 2026 15:11:45 -0700 Subject: [PATCH 05/28] app: fix at:// video links (route renamed VodPlayerDemo -> Video) The at:// deep-link resolver still navigated place.stream.video records to the "VodPlayerDemo" screen, which was renamed to "Video". Navigating to a non-existent route silently fell back to the homepage, so pasting an at:///place.stream.video/ URL never reached the player. Point it at the "Video" route. Co-Authored-By: Claude Opus 4.8 --- js/app/src/linking-config.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/js/app/src/linking-config.ts b/js/app/src/linking-config.ts index 0dc601d47..d51552745 100644 --- a/js/app/src/linking-config.ts +++ b/js/app/src/linking-config.ts @@ -65,7 +65,7 @@ function resolveAtUriNavigation( return { routes: [ { - name: "VodPlayerDemo", + name: "Video", params: { user: atUri.authority, tid: atUri.rkey }, }, ], -- 2.51.2 From 0203e9a4d10b2d922c1841268dd5234de399da86 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 2 Jun 2026 11:13:54 -0700 Subject: [PATCH 06/28] atproto: multi-relay firehose with cross-relay dedup Subscribe to several relays at once so losing any one no longer drops all incoming atproto data. --relay-host / SP_RELAY_HOST is now comma-separated; StartFirehose runs one resilient consumer per relay (reconnecting with capped backoff, never crashing the node on a single relay's outage) plus one liveness monitor, all sharing a deduper. Dedup keys on the commit CID, which is content-addressed and therefore identical no matter which relay forwarded it; identity events key on (did, time, handle) so the same broadcast collapses while genuine sequential changes pass through. The deduper is a two-generation rotating set: O(1), mutex-safe, bounded memory, expiry by dropping the old generation. A duplicate that slips past the window is harmless since downstream handlers are idempotent. Bonus: --relay-self / SP_RELAY_SELF (default on) appends our own PDS firehose so locally-published records are indexed immediately. Adds metrics streamplace_firehose_relays_connected{relay} and streamplace_firehose_events_deduped_total{kind}. Co-Authored-By: Claude Opus 4.8 --- pkg/atproto/firehose.go | 275 +++++++++++++++++------- pkg/atproto/firehose_dedup.go | 66 ++++++ pkg/atproto/firehose_dedup_test.go | 141 ++++++++++++ pkg/atproto/firehose_multirelay_test.go | 116 ++++++++++ pkg/config/config.go | 10 +- pkg/spmetrics/spmetrics.go | 18 ++ 6 files changed, 550 insertions(+), 76 deletions(-) create mode 100644 pkg/atproto/firehose_dedup.go create mode 100644 pkg/atproto/firehose_dedup_test.go create mode 100644 pkg/atproto/firehose_multirelay_test.go diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index a5d763774..c7d70546c 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -8,6 +8,7 @@ import ( "net/url" "runtime" "strings" + "sync/atomic" "time" comatproto "github.com/bluesky-social/indigo/api/atproto" @@ -28,6 +29,7 @@ import ( "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/model" notificationpkg "stream.place/streamplace/pkg/notifications" + "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/statedb" "slices" @@ -35,128 +37,253 @@ import ( "github.com/gorilla/websocket" ) +// dedupWindow bounds how long a key is remembered. Relays run within seconds of +// each other, so a few minutes covers all realistic skew. +const dedupWindow = 5 * time.Minute + type ATProtoSynchronizer struct { CLI *config.CLI Model model.Model StatefulDB *statedb.StatefulDB - LastSeen time.Time - LastEvent time.Time Noter notificationpkg.FirebaseNotifier Bus *bus.Bus PLCDirectory identity.Directory CachedPLCDirectory identity.Directory OATProxy *oatproxy.OATProxy + + // firehose liveness, written from every relay consumer concurrently + // (unix nanos). + lastSeen atomic.Int64 + lastEvent atomic.Int64 + + // cross-relay dedup, shared by every relay consumer. Initialized at the + // top of StartFirehose. + commitDedup *firehoseDeduper + identityDedup *firehoseDeduper +} + +func (atsync *ATProtoSynchronizer) markSeen() { + atsync.lastSeen.Store(time.Now().UnixNano()) +} + +func (atsync *ATProtoSynchronizer) markEvent(t time.Time) { + atsync.lastEvent.Store(t.UnixNano()) +} + +func sinceNanos(ns int64) time.Duration { + if ns == 0 { + return 0 + } + return time.Since(time.Unix(0, ns)) +} + +// relayHosts returns the ordered, de-duplicated list of relay websocket URLs to +// consume. It is the comma-separated --relay-host list, optionally with our own +// PDS firehose appended (so we always index records published to our built-in +// PDS, even before an external relay crawls them back to us). +func (atsync *ATProtoSynchronizer) relayHosts() []string { + seen := map[string]struct{}{} + var hosts []string + add := func(h string) { + h = strings.TrimSpace(h) + if h == "" { + return + } + if _, ok := seen[h]; ok { + return + } + seen[h] = struct{}{} + hosts = append(hosts, h) + } + for _, h := range strings.Split(atsync.CLI.RelayHost, ",") { + add(h) + } + if atsync.CLI.RelaySelf { + add(httpToWSURL(atsync.CLI.OwnPublicURL())) + } + return hosts } +func httpToWSURL(u string) string { + if strings.HasPrefix(u, "https://") { + return "wss://" + strings.TrimPrefix(u, "https://") + } + if strings.HasPrefix(u, "http://") { + return "ws://" + strings.TrimPrefix(u, "http://") + } + return u +} + +// StartFirehose subscribes to every configured relay (plus our own PDS, if +// enabled) concurrently, de-duplicating the overlapping event streams so each +// commit is indexed exactly once. It returns only when ctx is cancelled; a +// single relay flapping is logged and retried forever rather than taking the +// node down, since surviving relay outages is the whole point of running more +// than one. func (atsync *ATProtoSynchronizer) StartFirehose(ctx context.Context) error { - retryCount := 0 - retryWindow := time.Now() + ctx = log.WithLogValues(ctx, "func", "StartFirehose") + + relays := atsync.relayHosts() + if len(relays) == 0 { + return fmt.Errorf("no relay hosts configured") + } + + atsync.commitDedup = newFirehoseDeduper(dedupWindow) + atsync.identityDedup = newFirehoseDeduper(dedupWindow) + atsync.markSeen() + atsync.markEvent(time.Now()) + + log.Log(ctx, "starting firehose consumers", "relays", relays) + g, ctx := errgroup.WithContext(ctx) + for _, relay := range relays { + relay := relay + g.Go(func() error { + atsync.consumeRelay(ctx, relay) + return nil + }) + } + g.Go(func() error { + atsync.monitorFirehose(ctx) + return nil + }) + return g.Wait() +} + +// consumeRelay keeps a single relay's firehose connected, reconnecting with +// capped exponential backoff. It never returns an error: an unhealthy relay +// must not disturb the others (or crash the node), so it just loops until ctx +// is cancelled. +func (atsync *ATProtoSynchronizer) consumeRelay(ctx context.Context, relay string) { + ctx = log.WithLogValues(ctx, "relay", relay) + const ( + minBackoff = time.Second + maxBackoff = 30 * time.Second + ) + backoff := minBackoff for { if ctx.Err() != nil { - return nil + return + } + start := time.Now() + err := atsync.connectRelay(ctx, relay) + if ctx.Err() != nil { + return } - err := atsync.StartFirehoseRetry(ctx) if err != nil { - log.Error(ctx, "firehose error", "err", err) - - // Check if we're within the 1-minute window - now := time.Now() - if now.Sub(retryWindow) > time.Minute { - // Reset the counter if more than a minute has passed - retryCount = 1 - retryWindow = now - } else { - // Increment retry count if within the window - retryCount++ - if retryCount >= 3 { - log.Error(ctx, "firehose failed 3 times within a minute, crashing", "err", err) - return fmt.Errorf("firehose failed 3 times within a minute: %w", err) - } + log.Error(ctx, "relay firehose disconnected; reconnecting", "err", err, "backoff", backoff) + } else { + log.Warn(ctx, "relay firehose closed; reconnecting", "backoff", backoff) + } + // A connection that stayed healthy for a while earns a fresh backoff. + if time.Since(start) > time.Minute { + backoff = minBackoff + } + select { + case <-ctx.Done(): + return + case <-time.After(backoff): + } + if backoff < maxBackoff { + backoff *= 2 + if backoff > maxBackoff { + backoff = maxBackoff } } } } -func (atsync *ATProtoSynchronizer) StartFirehoseRetry(ctx context.Context) error { - ctx = log.WithLogValues(ctx, "func", "StartFirehose") - ctx, cancel := context.WithCancel(ctx) - defer cancel() - dialer := websocket.DefaultDialer - u, err := url.Parse(atsync.CLI.RelayHost) +// connectRelay dials one relay and pumps its firehose until the connection +// drops or ctx is cancelled. Event handlers are spawned on the parent ctx (not +// the per-connection one) so an in-flight commit keeps indexing across a +// reconnect — important because dedup has already claimed it, so no other relay +// will re-deliver it to us. +func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay string) error { + u, err := url.Parse(relay) if err != nil { - return fmt.Errorf("invalid relayHost URI: %w", err) + return fmt.Errorf("invalid relay URI %q: %w", relay, err) } u.Path = "xrpc/com.atproto.sync.subscribeRepos" - // if cursor != 0 { - // u.RawQuery = fmt.Sprintf("cursor=%d", cursor) - // } - con, _, err := dialer.Dial(u.String(), http.Header{ + + con, _, err := websocket.DefaultDialer.Dial(u.String(), http.Header{ "User-Agent": []string{aqhttp.UserAgent}, }) if err != nil { return fmt.Errorf("subscribing to firehose failed (dialing): %w", err) } + defer con.Close() + + spmetrics.FirehoseRelaysConnected.WithLabelValues(relay).Set(1) + defer spmetrics.FirehoseRelaysConnected.WithLabelValues(relay).Set(0) + + streamCtx, cancel := context.WithCancel(ctx) + defer cancel() rsc := &events.RepoStreamCallbacks{ RepoCommit: func(evt *comatproto.SyncSubscribeRepos_Commit) error { + atsync.markSeen() + if atsync.commitDedup.seen(evt.Commit.String()) { + spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("commit").Inc() + return nil + } go atsync.handleCommitEventOps(ctx, evt) return nil }, RepoIdentity: func(evt *comatproto.SyncSubscribeRepos_Identity) error { + atsync.markSeen() + if atsync.identityDedup.seen(identityDedupKey(evt)) { + spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("identity").Inc() + return nil + } go atsync.handleIdentityEventOps(ctx, evt) return nil }, Error: func(evt *events.ErrorFrame) error { - log.Error(ctx, "firehose error", "err", evt.Error, "message", evt.Message) cancel() - return fmt.Errorf("firehose error: %s", evt.Error) + return fmt.Errorf("firehose error: %s: %s", evt.Error, evt.Message) }, } - scheduler := parallel.NewScheduler( - 10, - 100, - atsync.CLI.RelayHost, - rsc.EventHandler, - ) - - log.Log(ctx, "starting firehose consumer", "relayHost", atsync.CLI.RelayHost) + scheduler := parallel.NewScheduler(10, 100, relay, rsc.EventHandler) - g, ctx := errgroup.WithContext(ctx) + log.Log(ctx, "connected to relay firehose") + return events.HandleRepoStream(streamCtx, con, scheduler, nil) +} - g.Go(func() error { - err := events.HandleRepoStream(ctx, con, scheduler, nil) - if err != nil { - log.Error(ctx, "firehose error", "err", err) - return err - } - return nil - }) +// identityDedupKey distinguishes a single identity broadcast (so the same one +// relayed by several relays collapses) from genuinely separate updates to the +// same DID (a later handle change must NOT be dropped). The originating Time is +// preserved as the event is relayed, so (did, time, handle) is stable for one +// broadcast yet distinct across real changes. seq is deliberately excluded — it +// is assigned per-relay and so differs for the very duplicates we want to fold. +func identityDedupKey(evt *comatproto.SyncSubscribeRepos_Identity) string { + handle := "" + if evt.Handle != nil { + handle = *evt.Handle + } + return evt.Did + "\x00" + evt.Time + "\x00" + handle +} - g.Go(func() error { - ticker := time.NewTicker(5 * time.Second) - defer ticker.Stop() - for { - select { - case <-ctx.Done(): - return nil - case <-ticker.C: - since := time.Since(atsync.LastEvent) - goroutines := runtime.NumGoroutine() - if since > 10*time.Second { - log.Warn(ctx, fmt.Sprintf("firehose is %s behind real time", since), "goroutines", goroutines) - } else { - log.Debug(ctx, fmt.Sprintf("firehose is %s behind real time", since), "goroutines", goroutines) - } - if time.Since(atsync.LastSeen) > 10*time.Second { - log.Warn(ctx, fmt.Sprintf("firehose dry; no new events for %s", time.Since(atsync.LastSeen))) - } +func (atsync *ATProtoSynchronizer) monitorFirehose(ctx context.Context) { + ticker := time.NewTicker(5 * time.Second) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + since := sinceNanos(atsync.lastEvent.Load()) + goroutines := runtime.NumGoroutine() + if since > 10*time.Second { + log.Warn(ctx, fmt.Sprintf("firehose is %s behind real time", since), "goroutines", goroutines) + } else { + log.Debug(ctx, fmt.Sprintf("firehose is %s behind real time", since), "goroutines", goroutines) + } + if dry := sinceNanos(atsync.lastSeen.Load()); dry > 10*time.Second { + log.Warn(ctx, fmt.Sprintf("firehose dry; no new events for %s", dry)) } } - }) - - return g.Wait() + } } var CollectionFilter = []string{ @@ -169,8 +296,6 @@ var CollectionFilter = []string{ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Commit) { ctx = log.WithLogValues(ctx, "event", "commit", "did", evt.Repo, "rev", evt.Rev, "seq", fmt.Sprintf("%d", evt.Seq), "func", "handleCommitEventOps") - now := time.Now() - atsync.LastSeen = now if evt.TooBig { log.Warn(ctx, "skipping tooBig events for now") @@ -208,7 +333,7 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt continue } opTime := aqt.Time() - atsync.LastEvent = opTime + atsync.markEvent(opTime) r, err := atsync.Model.GetRepo(evt.Repo) if err != nil { diff --git a/pkg/atproto/firehose_dedup.go b/pkg/atproto/firehose_dedup.go new file mode 100644 index 000000000..0d7921060 --- /dev/null +++ b/pkg/atproto/firehose_dedup.go @@ -0,0 +1,66 @@ +package atproto + +import ( + "sync" + "time" +) + +// firehoseDeduper collapses duplicate firehose events that arrive when we +// subscribe to more than one relay (and to our own PDS). The same repo commit +// is forwarded by every relay that carries it; its commit CID is content +// addressed, so it is byte-identical regardless of which relay delivered it. +// That makes the commit CID a perfect cross-relay dedup key. +// +// It is a two-generation rotating set: keys live in cur until a window elapses, +// then cur becomes prev and a fresh cur is started. A lookup checks both +// generations, so any key is remembered for at least one window and at most +// two. Expiry is free (drop the old map) and there are no per-entry timers. +// +// Relays run within seconds of each other, so a window of a few minutes +// absorbs all realistic skew. If a duplicate ever does slip past (e.g. a relay +// lagging longer than the window), the downstream handlers are idempotent, so +// the worst case is wasted work — never lost or corrupted data. +type firehoseDeduper struct { + mu sync.Mutex + cur map[string]struct{} + prev map[string]struct{} + window time.Duration + lastSwap time.Time +} + +func newFirehoseDeduper(window time.Duration) *firehoseDeduper { + return &firehoseDeduper{ + cur: make(map[string]struct{}), + prev: make(map[string]struct{}), + window: window, + lastSwap: time.Now(), + } +} + +// seen reports whether key has been observed within the dedup window. The first +// call for a key returns false and records it; subsequent calls return true +// until the key ages out. It is safe for concurrent use, and two callers racing +// on the same key are serialized so exactly one sees false. +func (d *firehoseDeduper) seen(key string) bool { + d.mu.Lock() + defer d.mu.Unlock() + + if time.Since(d.lastSwap) >= d.window { + d.prev = d.cur + d.cur = make(map[string]struct{}, len(d.prev)) + d.lastSwap = time.Now() + } + + if _, ok := d.cur[key]; ok { + return true + } + if _, ok := d.prev[key]; ok { + // Promote into cur so a steadily-recurring key doesn't age out from + // under us at a generation boundary. + d.cur[key] = struct{}{} + return true + } + + d.cur[key] = struct{}{} + return false +} diff --git a/pkg/atproto/firehose_dedup_test.go b/pkg/atproto/firehose_dedup_test.go new file mode 100644 index 000000000..df8c0003f --- /dev/null +++ b/pkg/atproto/firehose_dedup_test.go @@ -0,0 +1,141 @@ +package atproto + +import ( + "sync" + "sync/atomic" + "testing" + "time" + + "stream.place/streamplace/pkg/config" +) + +func TestFirehoseDeduperWithinWindow(t *testing.T) { + d := newFirehoseDeduper(time.Minute) + + if d.seen("cid-a") { + t.Fatal("first sighting of cid-a should not be a duplicate") + } + if !d.seen("cid-a") { + t.Fatal("second sighting of cid-a should be a duplicate") + } + if d.seen("cid-b") { + t.Fatal("first sighting of cid-b should not be a duplicate") + } + if !d.seen("cid-b") { + t.Fatal("second sighting of cid-b should be a duplicate") + } +} + +func TestFirehoseDeduperForgetsAfterTwoWindows(t *testing.T) { + const window = 20 * time.Millisecond + d := newFirehoseDeduper(window) + + if d.seen("cid-a") { + t.Fatal("first sighting should not be a duplicate") + } + + // One full window: cid-a rotates from cur into prev (still remembered). + time.Sleep(3 * window) + if d.seen("cid-b") { // triggers the rotation; prev={cid-a}, cur={cid-b} + t.Fatal("cid-b is new") + } + if !d.seen("cid-a") { + t.Fatal("cid-a should still be remembered after a single window (lives in prev)") + } + + // A second window with no sighting of cid-a drops it from both generations. + time.Sleep(3 * window) + if d.seen("cid-c") { // rotation: prev={cid-b,cid-a-promoted?}, cur={cid-c} + t.Fatal("cid-c is new") + } + time.Sleep(3 * window) + if d.seen("cid-d") { // another rotation drops the generation that held cid-a + t.Fatal("cid-d is new") + } + if d.seen("cid-a") { + t.Fatal("cid-a should have aged out after going unseen for two windows") + } +} + +// TestFirehoseDeduperConcurrent simulates the same commit arriving from many +// relays at once: exactly one caller must win (see it as new) so it is indexed +// exactly once. +func TestFirehoseDeduperConcurrent(t *testing.T) { + d := newFirehoseDeduper(time.Minute) + + const goroutines = 200 + var newSightings atomic.Int64 + var wg sync.WaitGroup + start := make(chan struct{}) + + for i := 0; i < goroutines; i++ { + wg.Add(1) + go func() { + defer wg.Done() + <-start + if !d.seen("the-one-true-commit") { + newSightings.Add(1) + } + }() + } + close(start) + wg.Wait() + + if got := newSightings.Load(); got != 1 { + t.Fatalf("expected exactly one caller to see the commit as new, got %d", got) + } +} + +func TestRelayHostsDedupesAndAppendsSelf(t *testing.T) { + cases := []struct { + name string + relayHost string + relaySelf bool + httpAddr string + want []string + }{ + { + name: "single relay", + relayHost: "wss://bsky.network", + want: []string{"wss://bsky.network"}, + }, + { + name: "comma separated with whitespace", + relayHost: "wss://relay1.example, wss://relay2.example ,wss://relay1.example", + want: []string{"wss://relay1.example", "wss://relay2.example"}, + }, + { + name: "appends self", + relayHost: "wss://bsky.network", + relaySelf: true, + httpAddr: "127.0.0.1:39000", + want: []string{"wss://bsky.network", "ws://127.0.0.1:39000"}, + }, + { + name: "self already listed is not duplicated", + relayHost: "wss://bsky.network,ws://127.0.0.1:39000", + relaySelf: true, + httpAddr: "127.0.0.1:39000", + want: []string{"wss://bsky.network", "ws://127.0.0.1:39000"}, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + atsync := &ATProtoSynchronizer{CLI: &config.CLI{ + RelayHost: tc.relayHost, + RelaySelf: tc.relaySelf, + HTTPAddr: tc.httpAddr, + }} + got := atsync.relayHosts() + if len(got) != len(tc.want) { + t.Fatalf("got %v, want %v", got, tc.want) + } + for i := range got { + if got[i] != tc.want[i] { + t.Fatalf("got %v, want %v", got, tc.want) + } + } + }) + } +} diff --git a/pkg/atproto/firehose_multirelay_test.go b/pkg/atproto/firehose_multirelay_test.go new file mode 100644 index 000000000..d402ff024 --- /dev/null +++ b/pkg/atproto/firehose_multirelay_test.go @@ -0,0 +1,116 @@ +package atproto + +import ( + "context" + "fmt" + "net/url" + "testing" + "time" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/bluesky-social/indigo/util" + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/bus" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/devenv" + "stream.place/streamplace/pkg/model" + "stream.place/streamplace/pkg/spmetrics" + "stream.place/streamplace/pkg/statedb" + "stream.place/streamplace/pkg/streamplace" +) + +// TestMultiRelayDedup subscribes to the same dev PDS firehose twice (via two +// distinct-but-equivalent URL spellings, so they aren't collapsed as identical +// relays). Every commit is therefore delivered twice, and the deduper must drop +// the second copy: the record is indexed exactly once and the dedup counter +// climbs. +func TestMultiRelayDedup(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + dev := devenv.WithDevEnv(t) + + parsed, err := url.Parse(dev.PDSURL) + require.NoError(t, err) + port := parsed.Port() + require.NotEmpty(t, port, "dev PDS URL should have a port") + // Two spellings of the same server: relayHosts() keeps them distinct, so we + // open two real connections to the one firehose. + relayHost := fmt.Sprintf("ws://localhost:%s,ws://127.0.0.1:%s", port, port) + + cli := config.CLI{ + BroadcasterHost: "example.com", + DBURL: ":memory:", + RelayHost: relayHost, + PLCURL: dev.PLCURL, + } + b := bus.NewBus() + cli.DataDir = t.TempDir() + mod, err := model.MakeDB(":memory:") + require.NoError(t, err) + state, err := statedb.MakeDB(context.Background(), &cli, nil, mod) + require.NoError(t, err) + atsync := &ATProtoSynchronizer{ + CLI: &cli, + StatefulDB: state, + Model: mod, + Bus: b, + PLCDirectory: dev.TestDirectory(), + } + + // Sanity-check the fan-out: two relays, no self. + require.Len(t, atsync.relayHosts(), 2) + + dedupedBefore := counterValue(t, spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("commit")) + + go func() { + _ = atsync.StartFirehose(ctx) + }() + + user := dev.CreateAccount(t) + msg := &streamplace.ChatMessage{ + LexiconTypeID: "place.stream.chat.message", + Text: "Hello from two relays!", + CreatedAt: time.Now().Add(-time.Second).Format(util.ISO8601), + Streamer: user.DID, + } + _, err = comatproto.RepoCreateRecord(ctx, user.XRPC, &comatproto.RepoCreateRecord_Input{ + Collection: "place.stream.chat.message", + Repo: user.DID, + Record: &lexutil.LexiconTypeDecoder{Val: msg}, + }) + require.NoError(t, err) + + // The message is indexed exactly once, even though two relays delivered it. + err = untilNoErrors(t, func() error { + messages, err := mod.MostRecentChatMessages(user.DID) + if err != nil { + return err + } + if len(messages) != 1 { + return fmt.Errorf("expected exactly 1 message, got %d", len(messages)) + } + return nil + }) + require.NoError(t, err) + + // And the duplicates the second relay delivered were dropped by the deduper. + err = untilNoErrors(t, func() error { + deduped := counterValue(t, spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("commit")) + if deduped <= dedupedBefore { + return fmt.Errorf("expected commit dedup counter to climb (before=%v now=%v)", dedupedBefore, deduped) + } + return nil + }) + require.NoError(t, err) +} + +func counterValue(t *testing.T, c prometheus.Counter) float64 { + t.Helper() + var m dto.Metric + require.NoError(t, c.Write(&m)) + return m.GetCounter().GetValue() +} diff --git a/pkg/config/config.go b/pkg/config/config.go index 9ec684889..9ad23b671 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -91,6 +91,7 @@ type CLI struct { PKCS11KeypairID string StreamerName string RelayHost string + RelaySelf bool Debug map[string]map[string]int AllowedStreams []string WideOpen bool @@ -497,11 +498,18 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }, &urfavecli.StringFlag{ Name: "relay-host", - Usage: "websocket url for relay firehose", + Usage: "comma-separated websocket url(s) for relay firehose(s); subscribing to several relays survives any one going down (duplicate events are deduped)", Value: "wss://bsky.network", Destination: &cli.RelayHost, Sources: urfavecli.EnvVars("SP_RELAY_HOST"), }, + &urfavecli.BoolFlag{ + Name: "relay-self", + Usage: "also subscribe to our own PDS firehose so records published here are indexed immediately", + Value: true, + Destination: &cli.RelaySelf, + Sources: urfavecli.EnvVars("SP_RELAY_SELF"), + }, &urfavecli.StringFlag{ Name: "color", Usage: "'true' to enable colorized logging, 'false' to disable", diff --git a/pkg/spmetrics/spmetrics.go b/pkg/spmetrics/spmetrics.go index 380c31aa5..1113bc27f 100644 --- a/pkg/spmetrics/spmetrics.go +++ b/pkg/spmetrics/spmetrics.go @@ -103,6 +103,24 @@ var LabelerFirehosesConnected = promauto.NewGaugeVec(prometheus.GaugeOpts{ Help: "number of currently connected labeler firehoses", }, []string{"labeler"}) +// FirehoseRelaysConnected is 1 while a relay's subscribeRepos websocket is +// connected and 0 while it is reconnecting, labeled by relay host. With +// multi-relay support this shows at a glance how many of the configured +// relays are currently feeding us. +var FirehoseRelaysConnected = promauto.NewGaugeVec(prometheus.GaugeOpts{ + Name: "streamplace_firehose_relays_connected", + Help: "1 if the relay's firehose websocket is currently connected, else 0", +}, []string{"relay"}) + +// FirehoseEventsDedupedTotal counts events dropped because the same commit +// (or identity update) already arrived from another relay. Labeled by event +// kind ("commit" / "identity"). A high count is expected and healthy — it is +// the redundant traffic we are paying for resilience. +var FirehoseEventsDedupedTotal = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "streamplace_firehose_events_deduped_total", + Help: "firehose events dropped as cross-relay duplicates, by kind", +}, []string{"kind"}) + // --- VOD processing --------------------------------------------------------- // VODProcessAttemptsTotal increments once per task dequeued for VOD -- 2.51.2 From 52057c8384fec241c038eced9655ae3f96f8ed81 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 2 Jun 2026 11:20:27 -0700 Subject: [PATCH 07/28] atproto: per-relay firehose cursor resume MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Persist each relay's high-water sequence number in the index DB (new RelayCursor table) so a reconnect or restart resumes from where we left off instead of re-tailing from live (which leaves a gap). The latest seq is kept in memory, advanced on every frame, and flushed every 5s plus once on shutdown to keep write pressure off the high-volume firehose. A relay with no stored cursor dials with none, so a fresh external relay tails from live rather than backfilling its entire history; once we've recorded progress we resume from it. This also stops our own PDS firehose (relay-self) from replaying its whole history on every reconnect — it replays once, then resumes. Cursor precision is best-effort: the parallel scheduler can surface frames out of sequence order, so on an unclean crash a few in-flight frames just below the high-water mark may be skipped on resume. That is safe — downstream handlers are idempotent, and with multiple relays plus a cold deduper after restart those commits are re-delivered and re-indexed. Co-Authored-By: Claude Opus 4.8 --- pkg/atproto/firehose.go | 33 ++++++++++- pkg/atproto/firehose_cursor.go | 85 +++++++++++++++++++++++++++++ pkg/atproto/firehose_cursor_test.go | 59 ++++++++++++++++++++ pkg/model/model.go | 4 ++ pkg/model/relay_cursor.go | 40 ++++++++++++++ 5 files changed, 219 insertions(+), 2 deletions(-) create mode 100644 pkg/atproto/firehose_cursor.go create mode 100644 pkg/atproto/firehose_cursor_test.go create mode 100644 pkg/model/relay_cursor.go diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index c7d70546c..2c9c92d25 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -7,6 +7,7 @@ import ( "net/http" "net/url" "runtime" + "strconv" "strings" "sync/atomic" "time" @@ -156,6 +157,27 @@ func (atsync *ATProtoSynchronizer) StartFirehose(ctx context.Context) error { // is cancelled. func (atsync *ATProtoSynchronizer) consumeRelay(ctx context.Context, relay string) { ctx = log.WithLogValues(ctx, "relay", relay) + + cursor := atsync.newRelayCursor(ctx, relay) + // Persist progress on a timer and once more on the way out, so a restart + // resumes near where we left off rather than re-tailing from live. + flushDone := make(chan struct{}) + go func() { + defer close(flushDone) + ticker := time.NewTicker(cursorFlushInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + cursor.flush(ctx) + return + case <-ticker.C: + cursor.flush(ctx) + } + } + }() + defer func() { <-flushDone }() + const ( minBackoff = time.Second maxBackoff = 30 * time.Second @@ -166,7 +188,7 @@ func (atsync *ATProtoSynchronizer) consumeRelay(ctx context.Context, relay strin return } start := time.Now() - err := atsync.connectRelay(ctx, relay) + err := atsync.connectRelay(ctx, relay, cursor) if ctx.Err() != nil { return } @@ -198,12 +220,17 @@ func (atsync *ATProtoSynchronizer) consumeRelay(ctx context.Context, relay strin // the per-connection one) so an in-flight commit keeps indexing across a // reconnect — important because dedup has already claimed it, so no other relay // will re-deliver it to us. -func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay string) error { +func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay string, cursor *relayCursor) error { u, err := url.Parse(relay) if err != nil { return fmt.Errorf("invalid relay URI %q: %w", relay, err) } u.Path = "xrpc/com.atproto.sync.subscribeRepos" + if seq, ok := cursor.param(); ok { + q := u.Query() + q.Set("cursor", strconv.FormatInt(seq, 10)) + u.RawQuery = q.Encode() + } con, _, err := websocket.DefaultDialer.Dial(u.String(), http.Header{ "User-Agent": []string{aqhttp.UserAgent}, @@ -222,6 +249,7 @@ func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay strin rsc := &events.RepoStreamCallbacks{ RepoCommit: func(evt *comatproto.SyncSubscribeRepos_Commit) error { atsync.markSeen() + cursor.observe(evt.Seq) if atsync.commitDedup.seen(evt.Commit.String()) { spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("commit").Inc() return nil @@ -231,6 +259,7 @@ func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay strin }, RepoIdentity: func(evt *comatproto.SyncSubscribeRepos_Identity) error { atsync.markSeen() + cursor.observe(evt.Seq) if atsync.identityDedup.seen(identityDedupKey(evt)) { spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("identity").Inc() return nil diff --git a/pkg/atproto/firehose_cursor.go b/pkg/atproto/firehose_cursor.go new file mode 100644 index 000000000..a79a4d694 --- /dev/null +++ b/pkg/atproto/firehose_cursor.go @@ -0,0 +1,85 @@ +package atproto + +import ( + "context" + "sync/atomic" + "time" + + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/model" +) + +// cursorFlushInterval bounds how often a relay's progress is written to the +// index DB. The firehose is high-volume, so we persist on a timer rather than +// per-event; between flushes the latest seq lives in memory and already covers +// in-process reconnects. +const cursorFlushInterval = 5 * time.Second + +// relayCursor tracks how far we've consumed one relay's firehose so we can +// resume after a disconnect or restart instead of re-tailing from live (which +// would leave a gap). It keeps the high-water sequence number in memory, +// updated on every frame, and persists it periodically and once on shutdown. +// +// Because the parallel scheduler may surface frames slightly out of sequence +// order, the persisted cursor is the highest seq observed — on an unclean crash +// a handful of in-flight frames just below it can be skipped on resume. That is +// safe here: downstream handlers are idempotent, and with several relays plus a +// cold deduper after restart, those commits get re-delivered and re-indexed. +type relayCursor struct { + host string + model model.Model + + latest atomic.Int64 // highest seq seen; 0 = nothing yet (tail from live) + flushed int64 // last persisted value; only the flush loop touches it +} + +func (atsync *ATProtoSynchronizer) newRelayCursor(ctx context.Context, host string) *relayCursor { + rc := &relayCursor{host: host, model: atsync.Model} + stored, err := atsync.Model.GetRelayCursor(host) + if err != nil { + log.Error(ctx, "failed to load relay cursor; tailing from live", "err", err) + return rc + } + if stored != nil { + rc.latest.Store(stored.Cursor) + rc.flushed = stored.Cursor + log.Log(ctx, "resuming relay from stored cursor", "cursor", stored.Cursor) + } + return rc +} + +// observe advances the high-water mark. Safe for concurrent callers (the +// scheduler runs several event workers). +func (rc *relayCursor) observe(seq int64) { + for { + cur := rc.latest.Load() + if seq <= cur { + return + } + if rc.latest.CompareAndSwap(cur, seq) { + return + } + } +} + +// param returns the cursor to dial with and whether to send one at all. With no +// progress yet we send none, so a fresh external relay tails from live instead +// of backfilling its entire history. +func (rc *relayCursor) param() (int64, bool) { + v := rc.latest.Load() + return v, v > 0 +} + +// flush persists the high-water mark if it has advanced since the last write. +// Only ever called from the single flush goroutine, so flushed is unsynchronized. +func (rc *relayCursor) flush(ctx context.Context) { + v := rc.latest.Load() + if v == rc.flushed { + return + } + if err := rc.model.UpsertRelayCursor(rc.host, v); err != nil { + log.Error(ctx, "failed to persist relay cursor", "err", err, "cursor", v) + return + } + rc.flushed = v +} diff --git a/pkg/atproto/firehose_cursor_test.go b/pkg/atproto/firehose_cursor_test.go new file mode 100644 index 000000000..94987fce3 --- /dev/null +++ b/pkg/atproto/firehose_cursor_test.go @@ -0,0 +1,59 @@ +package atproto + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/model" +) + +func TestRelayCursorResume(t *testing.T) { + mod, err := model.MakeDB(":memory:") + require.NoError(t, err) + atsync := &ATProtoSynchronizer{Model: mod} + ctx := context.Background() + const host = "wss://relay.example" + + // A fresh relay has no stored cursor, so we dial with none and tail live + // rather than backfilling the relay's whole history. + rc := atsync.newRelayCursor(ctx, host) + if _, ok := rc.param(); ok { + t.Fatal("fresh relay should not send a cursor") + } + + // observe tracks the high-water mark and never regresses on out-of-order + // frames (the parallel scheduler can surface them out of sequence). + rc.observe(100) + rc.observe(50) + rc.observe(120) + if v, ok := rc.param(); !ok || v != 120 { + t.Fatalf("expected cursor 120, got %d (ok=%v)", v, ok) + } + + // flush persists the high-water mark to the index DB. + rc.flush(ctx) + stored, err := mod.GetRelayCursor(host) + require.NoError(t, err) + require.NotNil(t, stored) + require.Equal(t, int64(120), stored.Cursor) + + // A second flush with no advance is a no-op, and the stored value is stable. + rc.observe(120) + rc.flush(ctx) + stored, err = mod.GetRelayCursor(host) + require.NoError(t, err) + require.Equal(t, int64(120), stored.Cursor) + + // A new cursor (as if the process restarted) resumes from the stored value. + resumed := atsync.newRelayCursor(ctx, host) + v, ok := resumed.param() + require.True(t, ok) + require.Equal(t, int64(120), v) + + // Cursors are independent per relay. + other := atsync.newRelayCursor(ctx, "wss://other.example") + if _, ok := other.param(); ok { + t.Fatal("a different relay must not inherit another relay's cursor") + } +} diff --git a/pkg/model/model.go b/pkg/model/model.go index 0ae5e4e07..058b7b5ce 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -102,6 +102,9 @@ type Model interface { GetLabeler(did string) (*Labeler, error) UpdateLabelerCursor(did string, cursor int64) error + GetRelayCursor(host string) (*RelayCursor, error) + UpsertRelayCursor(host string, cursor int64) error + CreateLabel(label *Label) error GetActiveLabels(uri string) ([]*comatproto.LabelDefs_Label, error) @@ -243,6 +246,7 @@ func MakeDB(dbURL string) (Model, error) { PinnedRecord{}, ServerSettings{}, Labeler{}, + RelayCursor{}, Label{}, BroadcastOrigin{}, MetadataConfiguration{}, diff --git a/pkg/model/relay_cursor.go b/pkg/model/relay_cursor.go new file mode 100644 index 000000000..aaa1db1c8 --- /dev/null +++ b/pkg/model/relay_cursor.go @@ -0,0 +1,40 @@ +package model + +import ( + "errors" + + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +// RelayCursor remembers how far we have consumed each relay's firehose, keyed by +// the relay's websocket URL. On reconnect or restart we resume from the stored +// sequence number instead of re-tailing from live (which would leave a gap) or +// replaying from the beginning. Cursors are per-relay because each relay +// assigns its own sequence numbers. +type RelayCursor struct { + Host string `gorm:"primaryKey;column:host"` + Cursor int64 `gorm:"column:cursor"` +} + +// GetRelayCursor returns the stored cursor for a relay, or nil if we have never +// recorded one (i.e. this is a fresh subscription). +func (m *DBModel) GetRelayCursor(host string) (*RelayCursor, error) { + var rc RelayCursor + err := m.DB.Where("host = ?", host).First(&rc).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + return &rc, nil +} + +// UpsertRelayCursor stores the latest consumed sequence number for a relay. +func (m *DBModel) UpsertRelayCursor(host string, cursor int64) error { + return m.DB.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "host"}}, + DoUpdates: clause.AssignmentColumns([]string{"cursor"}), + }).Create(&RelayCursor{Host: host, Cursor: cursor}).Error +} -- 2.51.2 From 8bf9410faca39e673aa4127668c6c6ad0f29e50a Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 2 Jun 2026 12:21:23 -0700 Subject: [PATCH 08/28] atproto: always self-index, drop the relay-self flag There's no case where we wouldn't want to index our own actions, so remove the --relay-self / SP_RELAY_SELF toggle and always append our own PDS firehose. It's gated only on having an HTTP listener address to dial, so it stays a no-op in tests that run without the server. Co-Authored-By: Claude Opus 4.8 --- pkg/atproto/firehose.go | 10 ++++++---- pkg/atproto/firehose_dedup_test.go | 8 ++------ pkg/config/config.go | 10 +--------- 3 files changed, 9 insertions(+), 19 deletions(-) diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index 2c9c92d25..d7884d883 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -79,9 +79,11 @@ func sinceNanos(ns int64) time.Duration { } // relayHosts returns the ordered, de-duplicated list of relay websocket URLs to -// consume. It is the comma-separated --relay-host list, optionally with our own -// PDS firehose appended (so we always index records published to our built-in -// PDS, even before an external relay crawls them back to us). +// consume: the comma-separated --relay-host list, plus our own PDS firehose so +// records published to our built-in PDS are always indexed immediately (rather +// than waiting for an external relay to crawl them back to us). We only add +// ourselves when we actually have an HTTP listener to dial — i.e. never in +// tests that run without the server. func (atsync *ATProtoSynchronizer) relayHosts() []string { seen := map[string]struct{}{} var hosts []string @@ -99,7 +101,7 @@ func (atsync *ATProtoSynchronizer) relayHosts() []string { for _, h := range strings.Split(atsync.CLI.RelayHost, ",") { add(h) } - if atsync.CLI.RelaySelf { + if atsync.CLI.HTTPAddr != "" { add(httpToWSURL(atsync.CLI.OwnPublicURL())) } return hosts diff --git a/pkg/atproto/firehose_dedup_test.go b/pkg/atproto/firehose_dedup_test.go index df8c0003f..c55cc9a52 100644 --- a/pkg/atproto/firehose_dedup_test.go +++ b/pkg/atproto/firehose_dedup_test.go @@ -90,12 +90,11 @@ func TestRelayHostsDedupesAndAppendsSelf(t *testing.T) { cases := []struct { name string relayHost string - relaySelf bool httpAddr string want []string }{ { - name: "single relay", + name: "single relay, no listener so no self", relayHost: "wss://bsky.network", want: []string{"wss://bsky.network"}, }, @@ -105,16 +104,14 @@ func TestRelayHostsDedupesAndAppendsSelf(t *testing.T) { want: []string{"wss://relay1.example", "wss://relay2.example"}, }, { - name: "appends self", + name: "appends self when we have a listener", relayHost: "wss://bsky.network", - relaySelf: true, httpAddr: "127.0.0.1:39000", want: []string{"wss://bsky.network", "ws://127.0.0.1:39000"}, }, { name: "self already listed is not duplicated", relayHost: "wss://bsky.network,ws://127.0.0.1:39000", - relaySelf: true, httpAddr: "127.0.0.1:39000", want: []string{"wss://bsky.network", "ws://127.0.0.1:39000"}, }, @@ -124,7 +121,6 @@ func TestRelayHostsDedupesAndAppendsSelf(t *testing.T) { t.Run(tc.name, func(t *testing.T) { atsync := &ATProtoSynchronizer{CLI: &config.CLI{ RelayHost: tc.relayHost, - RelaySelf: tc.relaySelf, HTTPAddr: tc.httpAddr, }} got := atsync.relayHosts() diff --git a/pkg/config/config.go b/pkg/config/config.go index 9ad23b671..6c50de9e2 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -91,7 +91,6 @@ type CLI struct { PKCS11KeypairID string StreamerName string RelayHost string - RelaySelf bool Debug map[string]map[string]int AllowedStreams []string WideOpen bool @@ -498,18 +497,11 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }, &urfavecli.StringFlag{ Name: "relay-host", - Usage: "comma-separated websocket url(s) for relay firehose(s); subscribing to several relays survives any one going down (duplicate events are deduped)", + Usage: "comma-separated websocket url(s) for relay firehose(s); subscribing to several relays survives any one going down (duplicate events are deduped). Our own PDS firehose is always included so locally-published records are indexed immediately", Value: "wss://bsky.network", Destination: &cli.RelayHost, Sources: urfavecli.EnvVars("SP_RELAY_HOST"), }, - &urfavecli.BoolFlag{ - Name: "relay-self", - Usage: "also subscribe to our own PDS firehose so records published here are indexed immediately", - Value: true, - Destination: &cli.RelaySelf, - Sources: urfavecli.EnvVars("SP_RELAY_SELF"), - }, &urfavecli.StringFlag{ Name: "color", Usage: "'true' to enable colorized logging, 'false' to disable", -- 2.51.2 From fca2aba6132f20730cc35beb8bac6272610b06ca Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 3 Jun 2026 11:11:14 -0700 Subject: [PATCH 09/28] go mod tidy --- go.mod | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/go.mod b/go.mod index 1bf710ea9..ad90ee773 100644 --- a/go.mod +++ b/go.mod @@ -57,6 +57,7 @@ require ( github.com/pion/webrtc/v4 v4.0.11 github.com/piprate/json-gold v0.5.0 github.com/prometheus/client_golang v1.23.2 + github.com/prometheus/client_model v0.6.2 github.com/rivo/uniseg v0.4.7 github.com/rs/cors v1.11.1 github.com/samber/slog-http v1.4.0 @@ -444,7 +445,6 @@ require ( github.com/polydawn/refmt v0.89.1-0.20221221234430-40501e09de1f // indirect github.com/polyfloyd/go-errorlint v1.8.0 // indirect github.com/pquerna/cachecontrol v0.0.0-20180517163645-1555304b9b35 // indirect - github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/common v0.66.1 // indirect github.com/prometheus/procfs v0.17.0 // indirect github.com/prometheus/statsd_exporter v0.27.1 // indirect -- 2.51.2 From 2cd03d96ca36e04ae46fe44f77b4b6c9d4de954e Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 3 Jun 2026 11:21:27 -0700 Subject: [PATCH 10/28] v0.11.1 --- js/app/package.json | 2 +- js/docs/package.json | 2 +- lerna.json | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/js/app/package.json b/js/app/package.json index 446a01f96..f31ae1d8d 100644 --- a/js/app/package.json +++ b/js/app/package.json @@ -1,7 +1,7 @@ { "name": "@streamplace/app", "main": "./src/entrypoint.tsx", - "version": "0.11.0", + "version": "0.11.1", "runtimeVersion": "0.10.0", "scripts": { "start": "npx expo start -c --port 38081", diff --git a/js/docs/package.json b/js/docs/package.json index 68b4129ee..02fafb599 100644 --- a/js/docs/package.json +++ b/js/docs/package.json @@ -1,7 +1,7 @@ { "name": "streamplace-docs", "type": "module", - "version": "0.11.0", + "version": "0.11.1", "scripts": { "dev": "astro dev --host 0.0.0.0 --port 38082", "start": "astro dev --host 0.0.0.0 --port 38082", diff --git a/lerna.json b/lerna.json index abc389737..ec6197c80 100644 --- a/lerna.json +++ b/lerna.json @@ -1,5 +1,5 @@ { "$schema": "node_modules/lerna/schemas/lerna-schema.json", - "version": "0.11.0", + "version": "0.11.1", "npmClient": "pnpm" } -- 2.51.2 From 6041a989095932e88c5d87bb7995948ee6b25c11 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 3 Jun 2026 13:46:28 -0700 Subject: [PATCH 11/28] atproto: self-subscribe with Host: ServerHost so origins self-index MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The node indexes its own server-repo records (most importantly place.stream.media.origin) by subscribing to its own firehose at the loopback OwnPublicURL(). The firehose endpoint dispatches between the broadcaster and server-repo firehoses by request Host (Server.isServerPDS). Dialing loopback sends Host: 127.0.0.1, so on any node where ServerHost != BroadcasterHost the self-subscription lands on the broadcaster firehose and the server repo's own commits are never replayed back to us. media_origins stays empty, and getVideoList's "can this node serve it" filter drops every row — so VODs upload and play but won't list. Single-node dev (ServerHost == BroadcasterHost) hides this because isServerPDS short-circuits to true. Force Host: ServerHost on the self-subscription so it gets the server-PDS firehose, matching what a single-node deployment gets implicitly. The dial target stays loopback; gorilla/websocket pulls the Host header out and uses it as the HTTP Host. Factor the self URL into selfRelayURL() so connectRelay can recognize that connection. Note: existing nodes carry a relay_cursors row for ws://127.0.0.1: advanced in the broadcaster firehose's seq space; delete it once after deploy so the server branch replays from seq 0 and backfills existing origins (UpsertMediaOrigin is idempotent). Co-Authored-By: Claude Opus 4.8 --- pkg/atproto/firehose.go | 34 +++++++++++++++++++++++++++++----- 1 file changed, 29 insertions(+), 5 deletions(-) diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index d7884d883..bcf4efae6 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -101,12 +101,23 @@ func (atsync *ATProtoSynchronizer) relayHosts() []string { for _, h := range strings.Split(atsync.CLI.RelayHost, ",") { add(h) } - if atsync.CLI.HTTPAddr != "" { - add(httpToWSURL(atsync.CLI.OwnPublicURL())) + if self := atsync.selfRelayURL(); self != "" { + add(self) } return hosts } +// selfRelayURL is the websocket URL of our own firehose listener, or "" +// when there's no HTTP listener to dial (e.g. tests run without the +// server). Used both to add ourselves to the relay set and to recognize +// that connection in connectRelay so we can present the right Host header. +func (atsync *ATProtoSynchronizer) selfRelayURL() string { + if atsync.CLI.HTTPAddr == "" { + return "" + } + return httpToWSURL(atsync.CLI.OwnPublicURL()) +} + func httpToWSURL(u string) string { if strings.HasPrefix(u, "https://") { return "wss://" + strings.TrimPrefix(u, "https://") @@ -234,9 +245,22 @@ func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay strin u.RawQuery = q.Encode() } - con, _, err := websocket.DefaultDialer.Dial(u.String(), http.Header{ - "User-Agent": []string{aqhttp.UserAgent}, - }) + header := http.Header{"User-Agent": []string{aqhttp.UserAgent}} + // Our own loopback listener serves both the broadcaster and the + // server-repo firehoses on one port and disambiguates them by the + // request Host (see Server.isServerPDS). Dialing loopback sends + // Host: 127.0.0.1, which lands us on the broadcaster firehose, so the + // server repo's own records — most importantly place.stream.media.origin, + // which getVideoList's "can this node serve it" filter depends on — + // would never get indexed locally. Force Host to ServerHost so the + // self-subscription gets the server-PDS firehose, matching what a + // single-node (ServerHost == BroadcasterHost) deployment gets for free. + // gorilla/websocket pulls the "Host" header out and uses it as the + // HTTP Host while still dialing the loopback address in u. + if relay == atsync.selfRelayURL() && atsync.CLI.ServerHost != "" { + header.Set("Host", atsync.CLI.ServerHost) + } + con, _, err := websocket.DefaultDialer.Dial(u.String(), header) if err != nil { return fmt.Errorf("subscribing to firehose failed (dialing): %w", err) } -- 2.51.2 From c58a714fd4b02a7cc744594d1fc4501a45c81bc0 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 3 Jun 2026 13:52:02 -0700 Subject: [PATCH 12/28] v0.11.2 --- js/app/package.json | 2 +- js/docs/package.json | 2 +- lerna.json | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/js/app/package.json b/js/app/package.json index f31ae1d8d..67b663f6a 100644 --- a/js/app/package.json +++ b/js/app/package.json @@ -1,7 +1,7 @@ { "name": "@streamplace/app", "main": "./src/entrypoint.tsx", - "version": "0.11.1", + "version": "0.11.2", "runtimeVersion": "0.10.0", "scripts": { "start": "npx expo start -c --port 38081", diff --git a/js/docs/package.json b/js/docs/package.json index 02fafb599..693501bd9 100644 --- a/js/docs/package.json +++ b/js/docs/package.json @@ -1,7 +1,7 @@ { "name": "streamplace-docs", "type": "module", - "version": "0.11.1", + "version": "0.11.2", "scripts": { "dev": "astro dev --host 0.0.0.0 --port 38082", "start": "astro dev --host 0.0.0.0 --port 38082", diff --git a/lerna.json b/lerna.json index ec6197c80..df61457d9 100644 --- a/lerna.json +++ b/lerna.json @@ -1,5 +1,5 @@ { "$schema": "node_modules/lerna/schemas/lerna-schema.json", - "version": "0.11.1", + "version": "0.11.2", "npmClient": "pnpm" } -- 2.51.2 From 624dfda4333fbb5046cfb2c35578906302395cbf Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 4 Jun 2026 09:56:13 -0700 Subject: [PATCH 13/28] login: always send to settings, no back --- .../settings/account-category-settings.tsx | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/js/app/components/settings/account-category-settings.tsx b/js/app/components/settings/account-category-settings.tsx index 5f040c9c7..cbc9170d6 100644 --- a/js/app/components/settings/account-category-settings.tsx +++ b/js/app/components/settings/account-category-settings.tsx @@ -55,14 +55,10 @@ export function AccountCategorySettings() { width="min" variant="secondary" onPress={() => { - if (navigation.canGoBack()) { - navigation.goBack(); - } else { - const params = convertNavigationParams({ - screen: "MainSettings", - }); - navigation.navigate(params.screen as any, params.params); - } + const params = convertNavigationParams({ + screen: "MainSettings", + }); + navigation.navigate(params.screen as any, params.params); }} > {tn("go-back")} -- 2.51.2 From 55aa678423636b5ffda0c2d78677262a898f6a28 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 4 Jun 2026 09:57:58 -0700 Subject: [PATCH 14/28] v0.11.3 --- js/app/package.json | 2 +- js/docs/package.json | 2 +- lerna.json | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/js/app/package.json b/js/app/package.json index 67b663f6a..c4220e420 100644 --- a/js/app/package.json +++ b/js/app/package.json @@ -1,7 +1,7 @@ { "name": "@streamplace/app", "main": "./src/entrypoint.tsx", - "version": "0.11.2", + "version": "0.11.3", "runtimeVersion": "0.10.0", "scripts": { "start": "npx expo start -c --port 38081", diff --git a/js/docs/package.json b/js/docs/package.json index 693501bd9..6dd49c8f3 100644 --- a/js/docs/package.json +++ b/js/docs/package.json @@ -1,7 +1,7 @@ { "name": "streamplace-docs", "type": "module", - "version": "0.11.2", + "version": "0.11.3", "scripts": { "dev": "astro dev --host 0.0.0.0 --port 38082", "start": "astro dev --host 0.0.0.0 --port 38082", diff --git a/lerna.json b/lerna.json index df61457d9..cb16fbfcc 100644 --- a/lerna.json +++ b/lerna.json @@ -1,5 +1,5 @@ { "$schema": "node_modules/lerna/schemas/lerna-schema.json", - "version": "0.11.2", + "version": "0.11.3", "npmClient": "pnpm" } -- 2.51.2 From 4c04849aa3d6af7ed746fc1a04351428b62b8b5b Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Thu, 4 Jun 2026 17:31:59 -0500 Subject: [PATCH 15/28] readd branding admin tool --- js/app/src/linking-config.ts | 4 ---- js/app/src/navigation-helper.tsx | 7 +++++-- js/app/src/navigation-types.ts | 6 +++--- js/app/src/shell.tsx | 6 ++++++ 4 files changed, 14 insertions(+), 9 deletions(-) diff --git a/js/app/src/linking-config.ts b/js/app/src/linking-config.ts index d51552745..351c4f155 100644 --- a/js/app/src/linking-config.ts +++ b/js/app/src/linking-config.ts @@ -136,7 +136,6 @@ export const SCREEN_PATHS = { PopoutChat: "chat-popout/:user", Embed: "embed/:user", InfoWidgetEmbed: "info-widget", - LegacyStream: "legacy/:user", DanmuOBS: "widgets/:user/danmu", PopoutStreamMonitor: "widgets/stream-monitor", PopoutInfoWidget: "widgets/info", @@ -228,7 +227,6 @@ export const streamplaceLinkingOptions: LinkingOptions + ); } -- 2.51.2 From 3dd826d8235c09c4e5762be1a46302ec4a94867a Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 4 Jun 2026 20:33:49 -0700 Subject: [PATCH 16/28] ci: cauterize docker push --- .gitlab-ci.yml | 155 ------------------------------------------------- 1 file changed, 155 deletions(-) diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index 228107503..14f61d3a7 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -16,66 +16,6 @@ workflow: variables: DOCKER_REPO: "public.ecr.aws/m4j3c0j7/streamplace" -build-bunny: - stage: build - interruptible: true - image: - name: gcr.io/kaniko-project/executor:v1.14.0-debug - entrypoint: [""] - timeout: 2 hours - script: - - mkdir -p /kaniko/.docker - - cp "$DOCKER_AUTH_CONFIG" /kaniko/.docker/config.json - - /kaniko/executor - --context "${CI_PROJECT_DIR}/docker" - --dockerfile "${CI_PROJECT_DIR}/docker/bunny.Dockerfile" - --destination "$DOCKER_REPO:bunny" - -builder: - stage: build - interruptible: true - allow_failure: true - image: - name: gcr.io/kaniko-project/executor:v1.14.0-debug - entrypoint: [""] - variables: - FQ_IMAGE_NAME: "$DOCKER_REPO:builder-$DOCKERFILE_HASH" - TARGET: builder - timeout: 2 hours - script: - - mkdir -p /kaniko/.docker - - cp "$DOCKER_AUTH_CONFIG" /kaniko/.docker/config.json - - /kaniko/executor - --build-arg TARGETARCH=amd64 - --build-arg DOCKERFILE_HASH=$DOCKERFILE_HASH - --cache=true - --context "${CI_PROJECT_DIR}" - --dockerfile "${CI_PROJECT_DIR}/docker/build.Dockerfile" - --destination "$FQ_IMAGE_NAME" - --target "$TARGET" - -builder-rebuild: - tags: - - tiny-job - stage: build - allow_failure: true - interruptible: true - image: - name: curlimages/curl:latest - entrypoint: [""] - needs: - - job: builder - timeout: 2 hours - script: - - echo "CI_PIPELINE_ID=$CI_PIPELINE_ID" - - echo "CI_PIPELINE_ID=$CI_PIPELINE_IID" - - echo "CI_PROJECT_ID=$CI_PROJECT_ID" - - > - curl - --data-raw '{"project_id": "'$CI_PROJECT_ID'", "pipeline_id": "'$CI_PIPELINE_ID'", "job_token": "'$CI_JOB_TOKEN'"}' - --fail-with-body - https://retry-gitlab.iameli.workers.dev/ - .build: stage: build interruptible: true @@ -125,7 +65,6 @@ leak-test-viewers: needs: - job: build-linux-amd64 artifacts: true - - job: build-bunny image: name: "$DOCKER_REPO:bunny" variables: @@ -147,7 +86,6 @@ leak-test-streamers: needs: - job: build-linux-amd64 artifacts: true - - job: build-bunny image: name: "$DOCKER_REPO:bunny" variables: @@ -180,96 +118,6 @@ test: reports: junit: test.xml -build-docker-mistserver: - stage: build - interruptible: true - needs: - - job: build-linux-amd64 - artifacts: true - image: - name: gcr.io/kaniko-project/executor:v1.14.0-debug - entrypoint: [""] - timeout: 2 hours - script: - - mkdir -p /kaniko/.docker - - cp "$DOCKER_AUTH_CONFIG" /kaniko/.docker/config.json - - /kaniko/executor - --build-arg TARGETARCH=amd64 - --build-arg STREAMPLACE_URL=$STREAMPLACE_URL_LINUX_AMD64 - --context "${CI_PROJECT_DIR}" - --dockerfile "${CI_PROJECT_DIR}/docker/mistserver.Dockerfile" - --destination "$DOCKER_REPO:$STREAMPLACE_BRANCH-mistserver" - --destination "$DOCKER_REPO:$STREAMPLACE_VERSION-mistserver" - -build-docker-amd64: - stage: build - interruptible: true - needs: - - job: build-linux-amd64 - artifacts: true - image: - name: gcr.io/kaniko-project/executor:v1.14.0-debug - entrypoint: [""] - timeout: 2 hours - script: - - mkdir -p /kaniko/.docker - - cp "$DOCKER_AUTH_CONFIG" /kaniko/.docker/config.json - - /kaniko/executor - --build-arg TARGETARCH=amd64 - --build-arg STREAMPLACE_URL=$STREAMPLACE_URL_LINUX_AMD64 - --context "${CI_PROJECT_DIR}" - --dockerfile "${CI_PROJECT_DIR}/docker/release.Dockerfile" - --destination "$DOCKER_REPO:$STREAMPLACE_BRANCH-amd64" - --destination "$DOCKER_REPO:$STREAMPLACE_VERSION-amd64" - -build-docker-arm64: - stage: build - interruptible: true - tags: - - linux - - arm64 - needs: - - job: build-linux-arm64 - artifacts: true - image: - name: gcr.io/kaniko-project/executor:v1.14.0-debug - entrypoint: [""] - timeout: 2 hours - script: - - mkdir -p /kaniko/.docker - - cp "$DOCKER_AUTH_CONFIG" /kaniko/.docker/config.json - - /kaniko/executor - --build-arg TARGETARCH=arm64 - --build-arg STREAMPLACE_URL=$STREAMPLACE_URL_LINUX_ARM64 - --context "${CI_PROJECT_DIR}" - --dockerfile "${CI_PROJECT_DIR}/docker/release.Dockerfile" - --destination "$DOCKER_REPO:$STREAMPLACE_BRANCH-arm64" - --destination "$DOCKER_REPO:$STREAMPLACE_VERSION-arm64" - --customPlatform=linux/arm64 - -build-docker-manifest: - stage: build - interruptible: true - image: - name: curlimages/curl:latest - entrypoint: [""] - needs: - - job: build-docker-arm64 - - job: build-docker-amd64 - - job: build-linux-amd64 # just for the environment variables - artifacts: true - before_script: - - curl -s -L https://github.com/estesp/manifest-tool/releases/download/v2.1.7/binaries-manifest-tool-2.1.7.tar.gz | tar xvz - - mv manifest-tool-linux-amd64 ./manifest-tool - - chmod +x ./manifest-tool - timeout: 2 hours - script: - - ./manifest-tool --docker-cfg "$DOCKER_AUTH_CONFIG" push from-args - --platforms linux/amd64,linux/arm64 - --template $DOCKER_REPO:$STREAMPLACE_VERSION-ARCH - --tags $STREAMPLACE_BRANCH - --target $DOCKER_REPO:$STREAMPLACE_VERSION - .build-android: stage: build interruptible: true @@ -371,9 +219,6 @@ build-desktop-darwin: - build-darwin-arm64 - build-android-debug - build-android-release - - build-docker-amd64 - - build-docker-arm64 - - build-docker-manifest - build-desktop-darwin - build-ios - test -- 2.51.2 From e53b86c9eae08cecddcb91dd22c453a5e9055658 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 4 Jun 2026 22:18:40 -0700 Subject: [PATCH 17/28] ci: build and publish docker images to ghcr.io from GitHub Actions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replaces the (now-removed) GitLab kaniko image builds with a GitHub Actions workflow that publishes to the free ghcr.io registry. Pushes on `next` and `v*` tags; pull requests build everything but skip the push, so the Dockerfiles stay validated without publishing. Images (all tags of ghcr.io/streamplace/streamplace): - :next / : / :latest release image, multi-arch amd64+arm64 - :next-mistserver / ... MistServer companion, amd64 only (upstream MistServer is x86-64) - :bunny leak-test fixture image, amd64 - :builder / :builder- cross-compile toolchain, amd64 only (clang/llvm/aptly/winehq are x86-64), pullable for complex local builds The release/mistserver images bake in the streamplace binary that the workflow just cross-compiled (uploaded between jobs as an artifact) rather than curl-ing it from a package registry, so they have no external download dependency. New docker/*.ghcr.Dockerfile files do this; the existing release.Dockerfile/mistserver.Dockerfile are left untouched for the GitLab no-push follow-up. The release image only runs `streamplace self-test` on the native arch, skipping it under arm64 QEMU emulation. Auth is GITHUB_TOKEN (no secrets needed). Builds validated locally with podman. Caveats worth a look before relying on it: - builder image is large; cached to a dedicated ghcr registry tag. - the binaries job cross-compiles both arches in one runner — watch disk on ubuntu-latest (Maximize build space step frees ~25GB). - overlaps compile work with build.yaml; could later be merged. Co-Authored-By: Claude Opus 4.8 --- .github/workflows/docker.yaml | 223 ++++++++++++++++++++++++++++++ docker/mistserver.ghcr.Dockerfile | 12 ++ docker/release.ghcr.Dockerfile | 24 ++++ 3 files changed, 259 insertions(+) create mode 100644 .github/workflows/docker.yaml create mode 100644 docker/mistserver.ghcr.Dockerfile create mode 100644 docker/release.ghcr.Dockerfile diff --git a/.github/workflows/docker.yaml b/.github/workflows/docker.yaml new file mode 100644 index 000000000..6c49b22e8 --- /dev/null +++ b/.github/workflows/docker.yaml @@ -0,0 +1,223 @@ +name: docker + +# Builds the streamplace container images and publishes them to ghcr.io. +# +# ghcr.io/streamplace/streamplace:next - release image (amd64+arm64) +# ghcr.io/streamplace/streamplace: - immutable release image +# ghcr.io/streamplace/streamplace:latest - newest tagged release +# ghcr.io/streamplace/streamplace:next-mistserver - MistServer companion (amd64) +# ghcr.io/streamplace/streamplace:bunny - leak-test fixture image (amd64) +# ghcr.io/streamplace/streamplace:builder - cross-compile toolchain (amd64) +# +# Pushes happen on `next` and on `v*` tags. Pull requests build everything but +# do not push, so the Dockerfiles stay validated without publishing. + +on: + pull_request: + push: + branches: + - next + tags: + - "v*" + +concurrency: + group: ${{ github.workflow }}-${{ github.head_ref || github.run_id }} + cancel-in-progress: true + +permissions: + contents: read + packages: write + +env: + REGISTRY_IMAGE: ghcr.io/streamplace/streamplace + +jobs: + # The cross-compilation toolchain image, published so it can be pulled for + # complex builds locally. amd64 only: the toolchain (clang/llvm, golangci, + # aptly, the winehq apt repo) is x86-64. Layer-cached to ghcr so an unchanged + # docker/build.Dockerfile is a fast cache hit. + builder: + name: builder image + runs-on: ubuntu-latest + steps: + - name: Maximize build space + run: + sudo rm -rf /usr/share/dotnet /usr/local/lib/android /opt/ghc + /opt/hostedtoolcache/CodeQL + + - name: Check out code + uses: actions/checkout@v4.1.7 + with: + fetch-depth: 0 + ref: ${{ github.event.pull_request.head.sha }} + + - uses: docker/setup-buildx-action@v3 + + - name: Log in to ghcr.io + uses: docker/login-action@v3 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ github.token }} + + - name: Compute tags + id: tags + run: | + short=${GITHUB_SHA::7} + echo "tags=${REGISTRY_IMAGE}:builder,${REGISTRY_IMAGE}:builder-${short}" >> "$GITHUB_OUTPUT" + + - name: Build and push builder + uses: docker/build-push-action@v6 + with: + context: . + file: docker/build.Dockerfile + target: builder + platforms: linux/amd64 + push: ${{ github.event_name != 'pull_request' }} + tags: ${{ steps.tags.outputs.tags }} + # The builder image is large; use a dedicated registry cache tag + # (no 10GB cap, unlike the GitHub Actions cache). + cache-from: + type=registry,ref=${{ env.REGISTRY_IMAGE }}:builder-buildcache + cache-to: + ${{ github.event_name != 'pull_request' && + format('type=registry,ref={0}:builder-buildcache,mode=max', + env.REGISTRY_IMAGE) || '' }} + + # Cross-compile the linux binaries inside the builder container (the same + # mechanism as build.yaml) and hand them to the image jobs as an artifact. + binaries: + name: linux binaries + runs-on: ubuntu-latest + outputs: + version: ${{ steps.version.outputs.version }} + steps: + - name: Maximize build space + run: + sudo rm -rf /usr/share/dotnet /usr/local/lib/android /opt/ghc + /opt/hostedtoolcache/CodeQL + + - name: Check out code + uses: actions/checkout@v4.1.7 + with: + fetch-depth: 0 + ref: ${{ github.event.pull_request.head.sha }} + + - name: Log in to ghcr.io (build-layer cache) + uses: docker/login-action@v3 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ github.token }} + + - name: Cross-compile linux amd64 + arm64 + run: | + sudo apt install podman -y + make in-container \ + BUILDER_TARGET=builder-no-darwin \ + DOCKER_PWD_MOUNT_PATH=/app \ + DOCKER_BUILD_OPTS="--layers --cache-to ghcr.io/streamplace/streamplace --cache-from ghcr.io/streamplace/streamplace" \ + DOCKER_OPTS="-e CI=true -e GITHUB_ACTION=true" \ + IN_CONTAINER_CMD="make linux-amd64 && make linux-arm64 && go run ./pkg/config/git/git.go -v > sp-version.txt" + + - name: Capture version + id: version + run: echo "version=$(cat sp-version.txt)" >> "$GITHUB_OUTPUT" + + - name: Upload binaries + uses: actions/upload-artifact@v4 + with: + name: streamplace-linux-binaries + path: | + build-linux-amd64/streamplace + build-linux-arm64/streamplace + retention-days: 1 + if-no-files-found: error + + # Assemble and publish the product images from the freshly-built binaries. + images: + name: ${{ matrix.image }} image + runs-on: ubuntu-latest + needs: binaries + strategy: + fail-fast: false + matrix: + include: + - image: release + file: docker/release.ghcr.Dockerfile + platforms: linux/amd64,linux/arm64 + suffix: "" + - image: mistserver + file: docker/mistserver.ghcr.Dockerfile + platforms: linux/amd64 + suffix: "-mistserver" + - image: bunny + file: docker/bunny.Dockerfile + platforms: linux/amd64 + suffix: "-bunny" + steps: + - name: Check out code + uses: actions/checkout@v4.1.7 + with: + fetch-depth: 0 + ref: ${{ github.event.pull_request.head.sha }} + + - name: Download binaries + uses: actions/download-artifact@v4 + with: + name: streamplace-linux-binaries + path: . + + - name: Stage build context + run: | + mkdir -p stage/build-linux-amd64 stage/build-linux-arm64 + cp build-linux-amd64/streamplace stage/build-linux-amd64/streamplace + cp build-linux-arm64/streamplace stage/build-linux-arm64/streamplace + cp docker/mistserver.json stage/mistserver.json + cp docker/bunny.sh stage/bunny.sh + + - uses: docker/setup-qemu-action@v3 + if: contains(matrix.platforms, 'arm64') + + - uses: docker/setup-buildx-action@v3 + + - name: Log in to ghcr.io + uses: docker/login-action@v3 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ github.token }} + + - name: Compute image tags + id: tags + env: + SUFFIX: ${{ matrix.suffix }} + VERSION: ${{ needs.binaries.outputs.version }} + run: | + set -euo pipefail + tags=("${REGISTRY_IMAGE}:${VERSION}${SUFFIX}") + if [ "${{ github.event_name }}" = "pull_request" ]; then + tags+=("${REGISTRY_IMAGE}:pr-${{ github.event.number }}${SUFFIX}") + elif [ "${{ github.ref_type }}" = "tag" ]; then + tags+=("${REGISTRY_IMAGE}:latest${SUFFIX}") + else + tags+=("${REGISTRY_IMAGE}:${{ github.ref_name }}${SUFFIX}") + fi + # bunny gets a stable, version-independent tag (its contents never change). + if [ "${{ matrix.image }}" = "bunny" ]; then + tags+=("${REGISTRY_IMAGE}:bunny") + fi + printf -v joined '%s,' "${tags[@]}" + echo "tags=${joined%,}" >> "$GITHUB_OUTPUT" + echo "tags: ${joined%,}" + + - name: Build and push ${{ matrix.image }} + uses: docker/build-push-action@v6 + with: + context: stage + file: ${{ matrix.file }} + platforms: ${{ matrix.platforms }} + push: ${{ github.event_name != 'pull_request' }} + tags: ${{ steps.tags.outputs.tags }} + cache-from: type=gha,scope=${{ matrix.image }} + cache-to: type=gha,mode=max,scope=${{ matrix.image }} diff --git a/docker/mistserver.ghcr.Dockerfile b/docker/mistserver.ghcr.Dockerfile new file mode 100644 index 000000000..3e5e2b2ab --- /dev/null +++ b/docker/mistserver.ghcr.Dockerfile @@ -0,0 +1,12 @@ +# amd64-only MistServer companion image published to ghcr.io by +# .github/workflows/docker.yaml. amd64 only because the upstream MistServer +# release (mistserver_64) is x86-64. The streamplace binary is baked in from the +# build context; see docker/release.ghcr.Dockerfile for the rationale. +FROM --platform=linux/amd64 ubuntu:24.04 +RUN apt-get update && apt-get install -y curl ca-certificates && rm -rf /var/lib/apt/lists/* +COPY build-linux-amd64/streamplace /usr/local/bin/streamplace +RUN chmod +x /usr/local/bin/streamplace +RUN cd /usr/bin && curl -L -o - https://r.mistserver.org/dl/mistserver_64V3.7.tar.gz | tar xzv +RUN mkdir -p /config +ADD mistserver.json /config/mistserver.json +CMD ["MistController", "-c", "/config/mistserver.json"] diff --git a/docker/release.ghcr.Dockerfile b/docker/release.ghcr.Dockerfile new file mode 100644 index 000000000..06f2e59a6 --- /dev/null +++ b/docker/release.ghcr.Dockerfile @@ -0,0 +1,24 @@ +# Multi-arch release image published to ghcr.io by .github/workflows/docker.yaml. +# +# Unlike docker/release.Dockerfile (which curls a prebuilt tarball from a +# package registry), the streamplace binary is baked in from the build context +# that the workflow just cross-compiled, so there is no external download +# dependency. buildx populates TARGETARCH per platform, selecting the matching +# binary that the workflow staged under build-linux-/. +ARG TARGETARCH +FROM --platform=linux/$TARGETARCH ubuntu:24.04 +RUN apt-get update && apt-get install -y curl ca-certificates && rm -rf /var/lib/apt/lists/* +ARG TARGETARCH +ARG BUILDARCH +COPY build-linux-${TARGETARCH}/streamplace /usr/local/bin/streamplace +# upload-artifact drops the executable bit, so restore it. Only self-test on the +# native arch — under QEMU emulation (e.g. arm64 built on an amd64 runner) the +# self-test is slow and flaky, so skip it there. +RUN chmod +x /usr/local/bin/streamplace \ + && if [ "$TARGETARCH" = "$BUILDARCH" ]; then \ + streamplace self-test; \ + else \ + echo "skipping self-test: $TARGETARCH image built under emulation on $BUILDARCH"; \ + fi +ENV SP_DATA_DIR=/var/lib/streamplace +CMD ["streamplace"] -- 2.51.2 From fba4d8b5b4c2da9c942fd08f3b99c17f75e48e6b Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Fri, 5 Jun 2026 11:23:33 -0500 Subject: [PATCH 18/28] Improve chat message reliability - Fix chat navigation path in AvatarButton - Add error toast when message submission fails - Implement rollback of optimistic chat messages on server error --- js/app/src/router.tsx | 4 +-- .../src/components/chat/chat-box.tsx | 6 +++- js/components/src/livestream-store/chat.tsx | 34 ++++++++++++++++--- 3 files changed, 35 insertions(+), 9 deletions(-) diff --git a/js/app/src/router.tsx b/js/app/src/router.tsx index 69c150318..d115b43ec 100644 --- a/js/app/src/router.tsx +++ b/js/app/src/router.tsx @@ -193,9 +193,7 @@ export const AvatarButton = () => { if (userProfile) { source = { uri: userProfile.avatar }; return ( - + { state = reduceChat(state, [localChat], [], []); store.setState(state); - await pdsAgent.com.atproto.repo.createRecord({ - repo: userDID, - collection: "place.stream.chat.message", - record, - }); + try { + await pdsAgent.com.atproto.repo.createRecord({ + repo: userDID, + collection: "place.stream.chat.message", + record, + }); + } catch (err) { + // Remove the optimistic message if the server call fails + const currentState = store.getState(); + const updatedIndex = { ...currentState.chatIndex }; + for (const [key, existingMsg] of Object.entries(updatedIndex)) { + if (existingMsg.uri === localChat.uri) { + delete updatedIndex[key]; + break; + } + } + store.setState({ + ...currentState, + chatIndex: updatedIndex, + chat: Object.keys(updatedIndex) + .sort((a, b) => { + const aTime = parseInt(a.split("-")[0], 10); + const bTime = parseInt(b.split("-")[0], 10); + return bTime - aTime; + }) + .map((key) => updatedIndex[key]), + }); + throw err; + } }; }; -- 2.51.2 From 2f5096acec7b538d81c2f220918c4d1b181d0242 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Thu, 4 Jun 2026 18:10:27 -0500 Subject: [PATCH 19/28] Shrink seek bar, add to embed --- js/app/src/screens/vod-embed.tsx | 74 +++++++++++++++++-- .../components/mobile-player/ui/seek-bar.tsx | 5 +- .../src/components/vod/vod-player.tsx | 19 ++--- 3 files changed, 76 insertions(+), 22 deletions(-) diff --git a/js/app/src/screens/vod-embed.tsx b/js/app/src/screens/vod-embed.tsx index e01fe15e7..c17f4a5e8 100644 --- a/js/app/src/screens/vod-embed.tsx +++ b/js/app/src/screens/vod-embed.tsx @@ -1,12 +1,18 @@ -import { VideoProvider, View, VodPlayer, zero } from "@streamplace/components"; +import { + LivestreamProvider, + VideoProvider, + View, + VodPlayer, + zero, +} from "@streamplace/components"; import { Redirect } from "components/aqlink"; +import { BottomControlBar } from "components/mobile/desktop-ui/bottom-controls"; +import { TopControlBar } from "components/mobile/desktop-ui/top-controls"; import { useEffect } from "react"; import { useStore } from "store"; -// Chrome-less VOD player for iframe embeds (/embed/:user/video/:tid). Mirrors -// the live EmbedScreen: hide the sidebar and render just the minimal -// (expo-video's native controls handle play/scrub) with no back -// button or surrounding metadata. +const { layout, h, w, position, px, py } = zero; + export default function VodEmbedScreen({ route, }: { @@ -29,9 +35,61 @@ export default function VodEmbedScreen({ return ( - - - + + + + + + {}} + embedded={true} + /> + + + + {}} + showChat={false} + /> + + + + + + ); } diff --git a/js/components/src/components/mobile-player/ui/seek-bar.tsx b/js/components/src/components/mobile-player/ui/seek-bar.tsx index eff71d862..2154e2e11 100644 --- a/js/components/src/components/mobile-player/ui/seek-bar.tsx +++ b/js/components/src/components/mobile-player/ui/seek-bar.tsx @@ -5,8 +5,7 @@ import { usePlayerStore } from "../../../player-store"; const TRACK_HEIGHT = 3; const THUMB_SIZE = 14; -// A tall, invisible touch strip so the thin bar is easy to grab/drag. -const TOUCH_HEIGHT = 28; +const TOUCH_HEIGHT = 16; // Native VOD scrub bar. The old implementation leaned on @rn-primitives/slider // driven by web pointer events (onPointerUp/onPointerEnter), which never fire @@ -64,7 +63,7 @@ export function SeekBar() { const bufferedPct = Math.max(0, Math.min(1, bufferedEnd / duration)) * 100; return ( - + diff --git a/js/components/src/components/vod/vod-player.tsx b/js/components/src/components/vod/vod-player.tsx index 8d37770bf..4361f31bf 100644 --- a/js/components/src/components/vod/vod-player.tsx +++ b/js/components/src/components/vod/vod-player.tsx @@ -4,24 +4,20 @@ import { PlayerProvider } from "../../player-store/player-provider"; import { usePlayerStore } from "../../player-store/player-store"; import Video from "../mobile-player/video"; -// A deliberately minimal VOD player: a PlayerProvider, the src pushed into the -// player store in "vod" mode, and the platform diff --git a/js/app/src/screens/vod-embed.tsx b/js/app/src/screens/vod-embed.tsx index c17f4a5e8..eb6dec832 100644 --- a/js/app/src/screens/vod-embed.tsx +++ b/js/app/src/screens/vod-embed.tsx @@ -6,13 +6,10 @@ import { zero, } from "@streamplace/components"; import { Redirect } from "components/aqlink"; -import { BottomControlBar } from "components/mobile/desktop-ui/bottom-controls"; -import { TopControlBar } from "components/mobile/desktop-ui/top-controls"; +import { DesktopUi } from "components/mobile/desktop-ui"; import { useEffect } from "react"; import { useStore } from "store"; -const { layout, h, w, position, px, py } = zero; - export default function VodEmbedScreen({ route, }: { @@ -37,56 +34,8 @@ export default function VodEmbedScreen({ - - - - {}} - embedded={true} - /> - - - - {}} - showChat={false} - /> - - - + + diff --git a/js/components/src/components/mobile-player/video.tsx b/js/components/src/components/mobile-player/video.tsx index 768d3b4c6..95ce73362 100644 --- a/js/components/src/components/mobile-player/video.tsx +++ b/js/components/src/components/mobile-player/video.tsx @@ -139,6 +139,7 @@ const VideoElement = forwardRef< VideoProps & { videoRef?: React.RefObject } >((props, ref) => { const x = usePlayerStore((x) => x); + const mode = usePlayerStore((x) => x.mode); const url = useStreamplaceStore((x) => x.url); const playerEvent = usePlayerStore((x) => x.playerEvent); const setMuteWasForced = usePlayerStore((x) => x.setMuteWasForced); @@ -204,7 +205,9 @@ const VideoElement = forwardRef< localVideoRef.current.play().catch((err) => { console.log("error playing video", err.name); if (err.name === "NotAllowedError") { - if (localVideoRef.current) { + if (mode === "vod") { + setAutoplayFailed(true); + } else if (localVideoRef.current) { console.log("Setting muted and retrying"); setMuted(true); localVideoRef.current.muted = true; diff --git a/js/components/src/components/vod/vod-player.tsx b/js/components/src/components/vod/vod-player.tsx index 4361f31bf..01d559634 100644 --- a/js/components/src/components/vod/vod-player.tsx +++ b/js/components/src/components/vod/vod-player.tsx @@ -2,20 +2,33 @@ import { useEffect } from "react"; import { View } from "react-native"; import { PlayerProvider } from "../../player-store/player-provider"; import { usePlayerStore } from "../../player-store/player-store"; +import { useMuted, useSetMuted } from "../../streamplace-store"; import Video from "../mobile-player/video"; export function VodPlayer({ src, objectFit = "contain", + embedded, + muted: mutedProp, + pictureInPictureEnabled, children, }: { src: string; objectFit?: "contain" | "cover"; + embedded?: boolean; + muted?: boolean; + pictureInPictureEnabled?: boolean; children?: React.ReactNode; }) { return ( - + {children} @@ -25,24 +38,53 @@ export function VodPlayer({ function VodPlayerInner({ src, objectFit, + embedded, + muted: mutedProp, + pictureInPictureEnabled, children, }: { src: string; objectFit: "contain" | "cover"; + embedded: boolean; + muted?: boolean; + pictureInPictureEnabled?: boolean; children?: React.ReactNode; }) { const setSrc = usePlayerStore((x) => x.setSrc); const setMode = usePlayerStore((x) => x.setMode); + const setEmbedded = usePlayerStore((x) => x.setEmbedded); const storeSrc = usePlayerStore((x) => x.src); + const setMuted = useSetMuted(); + const muted = useMuted(); + useEffect(() => { setMode("vod"); setSrc(src); }, [src, setMode, setSrc]); + useEffect(() => { + setEmbedded(embedded); + }, [embedded, setEmbedded]); + + useEffect(() => { + if (mutedProp !== undefined) { + let wasMuted: boolean | null = muted; + setTimeout(() => setMuted(mutedProp), 200); + return () => { + if (wasMuted !== null) setMuted(wasMuted); + }; + } + }, [mutedProp]); + return ( - {storeSrc === src ? ); -- 2.51.2 From b966eb0c53fdadb1e089ba6ba2feb896ca90e363 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Thu, 4 Jun 2026 19:45:39 -0500 Subject: [PATCH 21/28] Add visual play/pause indicator to desktop UI - Introduce `PlayPauseIndicator` component for state transition feedback - Remove manual play button from loading overlay - Refine control bar visibility logic to persist when paused --- js/app/components/mobile/desktop-ui.tsx | 60 +++++++------ js/app/components/mobile/desktop-ui/index.ts | 1 + .../mobile/desktop-ui/mute-overlay.tsx | 1 + .../desktop-ui/play-pause-indicator.tsx | 90 +++++++++++++++++++ .../mobile-player/ui/autoplay-button.tsx | 1 + .../ui/viewer-loading-overlay.tsx | 79 +++++----------- 6 files changed, 147 insertions(+), 85 deletions(-) create mode 100644 js/app/components/mobile/desktop-ui/play-pause-indicator.tsx diff --git a/js/app/components/mobile/desktop-ui.tsx b/js/app/components/mobile/desktop-ui.tsx index 7050fbbad..37dbf2ef7 100644 --- a/js/app/components/mobile/desktop-ui.tsx +++ b/js/app/components/mobile/desktop-ui.tsx @@ -1,4 +1,5 @@ import { + PlayerStatus, PlayerUI, PortalHost, Toast, @@ -26,6 +27,7 @@ import { MuteOverlay, TopControlBar, } from "./desktop-ui/index"; +import { PlayPauseIndicator } from "./desktop-ui/play-pause-indicator"; import { useResponsiveLayout } from "./useResponsiveLayout"; const { h, layout, position, w, px, py, r, p } = zero; @@ -72,6 +74,7 @@ export function DesktopUi({ const fullscreen = usePlayerStore((state) => state.fullscreen); const setFullscreen = usePlayerStore((state) => state.setFullscreen); const selectedRendition = usePlayerStore((state) => state.selectedRendition); + const status = usePlayerStore((state) => state.status); const safeAreaInsets = embedded ? { ...originalSafeAreaInsets, top: 0 } @@ -96,12 +99,13 @@ export function DesktopUi({ if (selectedRendition === "audio") return; if (ingest !== null) return; + if (status === PlayerStatus.PAUSE) return; fadeTimeout.current = setTimeout(() => { fadeOpacity.value = withTiming(0, { duration: 400 }); setIsControlsVisible(false); }, FADE_OUT_DELAY); - }, [fadeOpacity, selectedRendition, ingest]); + }, [fadeOpacity, selectedRendition, ingest, status]); const onPlayerHover = useCallback(() => { resetFadeTimer(); @@ -225,8 +229,8 @@ export function DesktopUi({ collapsable={false} > - + )} - - - - - - {isSelfAndNotLive && ( + + + + + {fullscreen && } ); diff --git a/js/app/components/mobile/desktop-ui/index.ts b/js/app/components/mobile/desktop-ui/index.ts index 91c372da3..3b8f0ba89 100644 --- a/js/app/components/mobile/desktop-ui/index.ts +++ b/js/app/components/mobile/desktop-ui/index.ts @@ -2,5 +2,6 @@ export { BottomControlBar } from "./bottom-controls"; export { KebabMenu } from "./kebab"; export { LiveBubble } from "./live-bubble"; export { MuteOverlay } from "./mute-overlay"; +export { PlayPauseIndicator } from "./play-pause-indicator"; export { TopControlBar } from "./top-controls"; export { VolumeSlider } from "./volume-slider"; diff --git a/js/app/components/mobile/desktop-ui/mute-overlay.tsx b/js/app/components/mobile/desktop-ui/mute-overlay.tsx index 09ed23ff8..565647daf 100644 --- a/js/app/components/mobile/desktop-ui/mute-overlay.tsx +++ b/js/app/components/mobile/desktop-ui/mute-overlay.tsx @@ -35,6 +35,7 @@ export function MuteOverlay() { h.percent[100], w.percent[100], ]} + pointerEvents="box-none" > { diff --git a/js/app/components/mobile/desktop-ui/play-pause-indicator.tsx b/js/app/components/mobile/desktop-ui/play-pause-indicator.tsx new file mode 100644 index 000000000..1827414dd --- /dev/null +++ b/js/app/components/mobile/desktop-ui/play-pause-indicator.tsx @@ -0,0 +1,90 @@ +import { PlayerStatus, usePlayerStore } from "@streamplace/components"; +import { Pause, Play } from "lucide-react-native"; +import { useCallback, useEffect, useRef } from "react"; +import Animated, { + useAnimatedStyle, + useSharedValue, + withSequence, + withTiming, +} from "react-native-reanimated"; + +export function PlayPauseIndicator() { + const status = usePlayerStore((x) => x.status); + const prevStatus = useRef(status); + const opacity = useSharedValue(0); + const scale = useSharedValue(0.5); + + const animate = useCallback( + (toPlaying: boolean) => { + opacity.value = 1; + opacity.value = withSequence( + withTiming(1, { duration: 100 }), + withTiming(1, { duration: 400 }), + withTiming(0, { duration: 300 }), + ); + if (toPlaying) { + scale.value = 0.8; + scale.value = withSequence( + withTiming(1, { duration: 100 }), + withTiming(1, { duration: 400 }), + withTiming(1.3, { duration: 300 }), + ); + } else { + scale.value = 1; + scale.value = withSequence( + withTiming(0.8, { duration: 100 }), + withTiming(0.8, { duration: 500 }), + ); + } + }, + [opacity, scale], + ); + + useEffect(() => { + if (status !== prevStatus.current) { + const wasPlaying = prevStatus.current === PlayerStatus.PLAYING; + const isPlaying = status === PlayerStatus.PLAYING; + const isPaused = status === PlayerStatus.PAUSE; + + if (isPlaying || isPaused) { + animate(isPlaying); + } + prevStatus.current = status; + } + }, [status, animate]); + + const animatedStyle = useAnimatedStyle(() => ({ + opacity: opacity.value, + transform: [{ scale: scale.value }], + })); + + const isPlaying = status === PlayerStatus.PLAYING; + + return ( + + {isPlaying ? ( + + ) : ( + + )} + + ); +} diff --git a/js/components/src/components/mobile-player/ui/autoplay-button.tsx b/js/components/src/components/mobile-player/ui/autoplay-button.tsx index 069558d39..8636e9306 100644 --- a/js/components/src/components/mobile-player/ui/autoplay-button.tsx +++ b/js/components/src/components/mobile-player/ui/autoplay-button.tsx @@ -50,6 +50,7 @@ export function AutoplayButton() { h.percent[100], w.percent[100], ]} + pointerEvents="box-none" > x.status); - const togglePlayPause = usePlayerStore((x) => x.togglePlayPause); - const { theme, zero: zt } = useTheme(); const opacity = useSharedValue(0); useEffect(() => { @@ -28,60 +18,33 @@ export function ViewerLoadingOverlay() { } }, [status, opacity]); - const animatedStyle = useAnimatedStyle(() => { - return { - opacity: opacity.value, - }; - }); + const animatedStyle = useAnimatedStyle(() => ({ + opacity: opacity.value, + })); if (status === PlayerStatus.PLAYING) { return ; } - if (status === PlayerStatus.SUSPEND) { - return null; // No overlay when stopped + if (status === PlayerStatus.SUSPEND || status === PlayerStatus.PAUSE) { + return null; } - const isPaused = status === PlayerStatus.PAUSE; - return ( - <> - - {!isPaused && } - - {isPaused && ( - - - - )} - + + + ); } -- 2.51.2 From cc7f9d2bcfd79cf2421183cd67442786fd2769f0 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Thu, 4 Jun 2026 20:31:07 -0500 Subject: [PATCH 22/28] show chat on live, show chrome on mobile web portrait --- js/app/components/mobile/player.tsx | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/js/app/components/mobile/player.tsx b/js/app/components/mobile/player.tsx index 3d0777caf..6329eee09 100644 --- a/js/app/components/mobile/player.tsx +++ b/js/app/components/mobile/player.tsx @@ -82,18 +82,18 @@ function PlayerWithProvider( onTeleport?: (targetHandle: string, targetDID: string) => void; }, ) { - let [showChat, setShowChat] = useState(props.mode === "live"); + let [showChat, setShowChat] = useState(props.mode != "vod"); if (props.mode === "vod") { showChat = false; } const { shouldShowChatSidePanel, chatPanelWidth } = useResponsiveLayout(); const chatVisible = shouldShowChatSidePanel && showChat; const { width: screenWidth, height: screenHeight } = useWindowDimensions(); - const { top: safeTop } = useSafeAreaInsets(); + let { top: safeTop } = useSafeAreaInsets(); const segDims = useSegmentDimensions(); const isPortrait = screenHeight > screenWidth; + // if the screen is portrait and video is landscaps const isPortraitLandscapeCase = - Platform.OS !== "web" && isPortrait && segDims.width > segDims.height && !shouldShowChatSidePanel && @@ -176,9 +176,7 @@ function PlayerWithProvider( chatSection = ( <> - {!showUnavailable && ( - - )} + ); } else if (shouldShowChatSidePanel) { @@ -233,7 +231,6 @@ function PlayerWithProvider( Back - {chatSection} ); @@ -263,7 +260,10 @@ function PlayerWithProvider( flex: 1, width: "100%", height: "100%", - paddingTop: isPortraitLandscapeCase ? 54 : undefined, + paddingTop: + isPortraitLandscapeCase && Platform.OS != "web" + ? 54 + : undefined, }, ]} > -- 2.51.2 From c35a8679dbcb50ed2d746a4b4ff31d8a9b2fca80 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Thu, 4 Jun 2026 22:15:54 -0500 Subject: [PATCH 23/28] Improve mobile VOD player layout and metadata - Add time display to VOD controls. - Update landscape video player sizing and positioning. - Redesign mobile metadata to include tag support and a scrollable action bar. - Add animation to the Like button. --- js/app/components/mobile/player.tsx | 7 +- js/app/components/mobile/ui.tsx | 4 +- .../components/mobile-player/ui/viewers.tsx | 1 - .../mobile-player/ui/vod-controls.tsx | 22 +++- .../src/components/vod/like-button.tsx | 37 ++++-- .../src/components/vod/vod-description.tsx | 43 +++++-- .../components/vod/vod-mobile-metadata.tsx | 105 +++++++++++------- 7 files changed, 156 insertions(+), 63 deletions(-) diff --git a/js/app/components/mobile/player.tsx b/js/app/components/mobile/player.tsx index 6329eee09..f176f999c 100644 --- a/js/app/components/mobile/player.tsx +++ b/js/app/components/mobile/player.tsx @@ -455,12 +455,17 @@ export function PlayerInner( aspectRatio: vodAspectRatio, // Cap so a landscape video fits (letterboxed via objectFit:contain) // instead of clipping; flexShrink:0 keeps it from collapsing. - maxHeight: winHeight * 0.7, + maxHeight: isLandscape ? winHeight : winHeight * 0.7, + // center the video horizontally when in landscape since it won't fill the full width + marginHorizontal: isLandscape + ? (winWidth - Math.min(winWidth, winHeight * vodAspectRatio)) / 2 + : 0, flexShrink: 0, }} > {videoContent} + {/* will get pushed below the video if landscape so probably fine? */} ); diff --git a/js/app/components/mobile/ui.tsx b/js/app/components/mobile/ui.tsx index 5c6d7107b..9243809e4 100644 --- a/js/app/components/mobile/ui.tsx +++ b/js/app/components/mobile/ui.tsx @@ -306,15 +306,15 @@ export function MobileUi({ {mode === "vod" && ( - + )} {chatSection} diff --git a/js/components/src/components/mobile-player/ui/viewers.tsx b/js/components/src/components/mobile-player/ui/viewers.tsx index 28f129f9b..9789445be 100644 --- a/js/components/src/components/mobile-player/ui/viewers.tsx +++ b/js/components/src/components/mobile-player/ui/viewers.tsx @@ -29,7 +29,6 @@ export function DehydratedViewers({ atoms.layout.flex.center, atoms.layout.flex.row, atoms.gap.all[2], - atoms.px[1], ]} > diff --git a/js/components/src/components/mobile-player/ui/vod-controls.tsx b/js/components/src/components/mobile-player/ui/vod-controls.tsx index 97388ee57..2531a0400 100644 --- a/js/components/src/components/mobile-player/ui/vod-controls.tsx +++ b/js/components/src/components/mobile-player/ui/vod-controls.tsx @@ -1,8 +1,8 @@ import { Gauge, Pause, Play } from "lucide-react-native"; import { Pressable } from "react-native"; +import { zero } from "../../.."; import { useLivestreamStore } from "../../../livestream-store"; import { PlayerStatus, usePlayerStore } from "../../../player-store"; - import { DropdownMenu, DropdownMenuGroup, @@ -15,6 +15,16 @@ import { View, } from "../../ui"; +function formatTime(seconds: number): string { + const s = Math.floor(seconds); + const h = Math.floor(s / 3600); + const m = Math.floor((s % 3600) / 60); + const sec = s % 60; + const pad = (n: number) => String(n).padStart(2, "0"); + if (h > 0) return `${h}:${pad(m)}:${pad(sec)}`; + return `${m}:${pad(sec)}`; +} + export function VodControls() { const mode = usePlayerStore((x) => x.mode); const status = usePlayerStore((x) => x.status); @@ -26,6 +36,9 @@ export function VodControls() { const renditions = mode === "vod" ? vodLevels : liveRenditions; const th = useTheme(); + const playTime = usePlayerStore((x) => x.playTime); + const duration = usePlayerStore((x) => x.duration); + if (mode !== "vod") return null; const isPlaying = status === PlayerStatus.PLAYING; @@ -36,7 +49,7 @@ export function VodControls() { flexDirection: "row", alignItems: "center", paddingHorizontal: 16, - paddingVertical: 4, + paddingVertical: 10, gap: 12, }} > @@ -55,6 +68,11 @@ export function VodControls() { /> )} + + {formatTime(playTime)} + / + {formatTime(duration)} + diff --git a/js/components/src/components/vod/like-button.tsx b/js/components/src/components/vod/like-button.tsx index ec6d86f33..b669e5f6c 100644 --- a/js/components/src/components/vod/like-button.tsx +++ b/js/components/src/components/vod/like-button.tsx @@ -1,7 +1,13 @@ +import { ThumbsUp } from "lucide-react-native"; import { useCallback, useEffect, useState } from "react"; import { TouchableOpacity, View } from "react-native"; +import Animated, { + useAnimatedStyle, + useSharedValue, + withSpring, +} from "react-native-reanimated"; import { useDID } from "../../streamplace-store/streamplace-store"; -import { gap, layout, p, useTheme } from "../../ui"; +import { gap, layout, useTheme } from "../../ui"; import { useCreateLike, useDeleteLike, @@ -21,6 +27,12 @@ export function LikeButton({ subjectUri }: { subjectUri: string }) { const deleteLike = useDeleteLike(); const { theme } = useTheme(); + const scale = useSharedValue(1); + + const animatedStyle = useAnimatedStyle(() => ({ + transform: [{ scale: scale.value }], + })); + const loadLikes = useCallback(async () => { try { const result = await getLikes(subjectUri, 50); @@ -40,6 +52,9 @@ export function LikeButton({ subjectUri }: { subjectUri: string }) { }, [loadLikes]); const toggleLike = useCallback(async () => { + scale.value = withSpring(1.05, { stiffness: 500, damping: 10 }, () => { + scale.value = withSpring(1, { stiffness: 500 }); + }); setLoading(true); try { if (userLiked && userLikeUri) { @@ -58,19 +73,19 @@ export function LikeButton({ subjectUri }: { subjectUri: string }) { } finally { setLoading(false); } - }, [userLiked, userLikeUri, subjectUri, createLike, deleteLike]); + }, [userLiked, userLikeUri, subjectUri, createLike, deleteLike, scale]); + + const heartColor = userLiked + ? theme.colors.background + : theme.colors.foreground; return ( - - - {userLiked ? "\u2764\uFE0F" : "\u2661"} - - - {likeCount} - + + + + + {likeCount} ); diff --git a/js/components/src/components/vod/vod-description.tsx b/js/components/src/components/vod/vod-description.tsx index 05f62255f..1c0137068 100644 --- a/js/components/src/components/vod/vod-description.tsx +++ b/js/components/src/components/vod/vod-description.tsx @@ -1,12 +1,9 @@ -import { View } from "react-native"; +import { ScrollView, View } from "react-native"; import type { PlaceStreamVideo } from "streamplace"; -import { mt, useTheme } from "../../ui"; +import { hexToRgba, mt, useTheme } from "../../ui"; import { useVideoStore } from "../../video-store/video-store"; import { Text } from "../ui/text"; -// The video's description, shown below the metadata. For now this is just the -// plain text below the video; a collapsible dropdown (and richtext facets) can -// come later. export function VodDescription() { const video = useVideoStore((x) => x.video); const { theme } = useTheme(); @@ -14,11 +11,43 @@ export function VodDescription() { if (!video) return null; const record = video.record as unknown as PlaceStreamVideo.Record; const description = record.description?.trim(); - if (!description) return null; + const tags = (record.tags as string[] | undefined) ?? []; return ( - {description} + {tags.length > 0 ? ( + + {tags.map((tag) => ( + + + {tag} + + + ))} + + ) : null} + {description ? ( + {description} + ) : null} ); } diff --git a/js/components/src/components/vod/vod-mobile-metadata.tsx b/js/components/src/components/vod/vod-mobile-metadata.tsx index 51c1a63f6..e99e6c39f 100644 --- a/js/components/src/components/vod/vod-mobile-metadata.tsx +++ b/js/components/src/components/vod/vod-mobile-metadata.tsx @@ -1,6 +1,5 @@ import { Image } from "expo-image"; -import { useWindowDimensions, View } from "react-native"; -import { zero } from "../.."; +import { ScrollView, useWindowDimensions, View } from "react-native"; import { useAuthor } from "../../hooks/useAuthor"; import { useAvatar } from "../../hooks/useAvatar"; import { useTitle } from "../../hooks/useTitle"; @@ -9,12 +8,10 @@ import { gap, layout, useTheme } from "../../ui"; import { useVideoStore } from "../../video-store/video-store"; import { Viewers } from "../mobile-player/ui/viewers"; import { ShareSheet } from "../share/sharesheet"; +import { Button } from "../ui"; import { Text } from "../ui/text"; import { LikeButton } from "./like-button"; -// Below this width we drop the avatar and the view count so the title, like -// and share controls don't crowd on small phones. Above it (tablets/desktop) -// the full row shows, matching what the desktop metadata bar used to carry. const NARROW_BREAKPOINT = 480; // rkeyFromAturi pulls the record key off an at:// URI @@ -49,25 +46,38 @@ export function VodMobileMetadata() { message: `Check out "${title || "this video"}" on Streamplace!`, }; - return ( + const handleText = ( + + {handle ? `@${handle}` : did} + + ); + + const actions = ( - {/* Left: avatar + title/author */} - - {wide && ( + + + {wide ? ( + + ) : null} + + + ); + + if (wide) { + return ( + + ) : null} - )} - - - {title || "Untitled"} - - - {handle ? `@${handle}` : did} - + + + + {title || "Untitled"} + + + + {handleText} + {actions} + + + ); + } - {/* Right: like + views + share */} - - {/* Use the server-canonical (DID-based) video.uri, not the store's - aturi — when the page is reached via a handle URL the aturi's - authority is the handle, and a like keyed on a handle subject - won't match the DID-keyed video record. */} - - {wide && } - - + return ( + + + {title || "Untitled"} + + {handleText} + + {actions} + ); } -- 2.51.2 From b8733006508f3ca8c8657cb312bcc4b2f4729fd1 Mon Sep 17 00:00:00 2001 From: Natalie Bridgers Date: Thu, 4 Jun 2026 15:52:08 -0500 Subject: [PATCH 24/28] take into account viewers and followers when generating live page --- pkg/model/follow.go | 25 ++++++ pkg/model/model.go | 1 + pkg/spxrpc/place_stream_live.go | 154 ++++++++++++++++++++++++++++---- pkg/spxrpc/spxrpc.go | 4 +- 4 files changed, 167 insertions(+), 17 deletions(-) diff --git a/pkg/model/follow.go b/pkg/model/follow.go index d911d239b..5450e9871 100644 --- a/pkg/model/follow.go +++ b/pkg/model/follow.go @@ -58,3 +58,28 @@ func (m *DBModel) GetUserFollowingUser(ctx context.Context, userDID, subjectDID } return &follow, result.Error } + +type followerCountRow struct { + SubjectDID string + Count int +} + +func (m *DBModel) CountFollowersBatch(ctx context.Context, dids []string) (map[string]int, error) { + if len(dids) == 0 { + return map[string]int{}, nil + } + var rows []followerCountRow + err := m.DB.Table("follows"). + Select("subject_did, COUNT(*) as count"). + Where("subject_did IN ?", dids). + Group("subject_did"). + Find(&rows).Error + if err != nil { + return nil, err + } + counts := make(map[string]int, len(rows)) + for _, r := range rows { + counts[r.SubjectDID] = r.Count + } + return counts, nil +} diff --git a/pkg/model/model.go b/pkg/model/model.go index 058b7b5ce..91f80686b 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -48,6 +48,7 @@ type Model interface { GetUserFollowing(ctx context.Context, userDID string) ([]Follow, error) GetUserFollowers(ctx context.Context, userDID string) ([]Follow, error) GetUserFollowingUser(ctx context.Context, userDID, subjectDID string) (*Follow, error) + CountFollowersBatch(ctx context.Context, dids []string) (map[string]int, error) DeleteFollow(ctx context.Context, userDID, rev string) error CreateFeedPost(ctx context.Context, post *FeedPost) error diff --git a/pkg/spxrpc/place_stream_live.go b/pkg/spxrpc/place_stream_live.go index 27f50b49a..b1c5fdb05 100644 --- a/pkg/spxrpc/place_stream_live.go +++ b/pkg/spxrpc/place_stream_live.go @@ -8,6 +8,7 @@ import ( "net/http" "net/url" "os" + "slices" "strconv" "time" @@ -72,11 +73,65 @@ func (s *Server) handlePlaceStreamLiveDenyTeleport(ctx context.Context, input *p } var replicationUpgrader = websocket.Upgrader{ - ReadBufferSize: 1024, - WriteBufferSize: 1024 * 1024 * 10, // 10MB - CheckOrigin: func(r *http.Request) bool { - return true - }, + CheckOrigin: func(r *http.Request) bool { return true }, +} + +const ( + scoreViewerWeight = 10 + scoreFollowerWeight = 1 + scoreFollowBoost = 1000 +) + +// getLiveStreamerScores returns a score map for all currently live streamers. +// Scores combine viewer count and follower count, with a boost for streamers +// the requesting user follows. Results are cached for 30s per userDID. +func (s *Server) getLiveStreamerScores(ctx context.Context, userDID string) (map[string]float64, error) { + cacheKey := userDID + if cacheKey == "" { + cacheKey = "_anon" + } + if cached, found := s.ScoreCache.Get(cacheKey); found { + return cached.(map[string]float64), nil + } + + segs, err := s.localDB.MostRecentSegments() + if err != nil { + return nil, err + } + dids := make([]string, len(segs)) + for i, seg := range segs { + dids[i] = seg.RepoDID + } + + followerCounts, err := s.model.CountFollowersBatch(ctx, dids) + if err != nil { + return nil, err + } + + var followedSet map[string]bool + if userDID != "" { + follows, err := s.model.GetUserFollowing(ctx, userDID) + if err == nil { + followedSet = make(map[string]bool, len(follows)) + for _, f := range follows { + followedSet[f.SubjectDID] = true + } + } + } + + scores := make(map[string]float64, len(dids)) + for _, did := range dids { + viewers := float64(s.bus.GetViewerCount(did)) + followers := float64(followerCounts[did]) + score := viewers*scoreViewerWeight + followers*scoreFollowerWeight + if followedSet != nil && followedSet[did] { + score += scoreFollowBoost + } + scores[did] = score + } + + s.ScoreCache.SetDefault(cacheKey, scores) + return scores, nil } func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before string, limit int, userDID string) (*placestream.LiveGetSegments_Output, error) { @@ -155,12 +210,86 @@ func (s *Server) handlePlaceStreamLiveGetSegments(ctx context.Context, before st } func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before string, limit int) (*placestream.LiveGetLiveUsers_Output, error) { - // Check cache first - cacheKey := fmt.Sprintf("live_users_%s_%d", before, limit) + // Extract echo context for query params + ec, _ := ctx.Value(echoContextKey).(echo.Context) + sortMode := "ranked" + if ec != nil { + if p := ec.QueryParam("sort"); p != "" { + sortMode = p + } + } + + // Get optional user DID for personalized scoring + var userDID string + sess, _ := oatproxy.GetOAuthSession(ctx) + if sess != nil { + userDID = sess.DID + } + + // Cache key includes sort mode and user for personalized results + cacheKey := fmt.Sprintf("live_users_%s_%s_%d", sortMode, userDID, limit) if cached, found := s.LiveUsersCache.Get(cacheKey); found { return cached.(*placestream.LiveGetLiveUsers_Output), nil } + if sortMode == "latest" { + return s.getLiveUsersLatest(ctx, before, limit, cacheKey) + } + return s.getLiveUsersRanked(ctx, limit, userDID, cacheKey) +} + +// getLiveUsersRanked returns live streams sorted by a score combining viewer +// count, follower count, and follow-relationship boost. +func (s *Server) getLiveUsersRanked(ctx context.Context, limit int, userDID string, cacheKey string) (*placestream.LiveGetLiveUsers_Output, error) { + scores, err := s.getLiveStreamerScores(ctx, userDID) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to compute stream scores") + } + + dids := make([]string, 0, len(scores)) + for did := range scores { + dids = append(dids, did) + } + slices.SortStableFunc(dids, func(a, dids_b string) int { + sa, sb := scores[a], scores[dids_b] + if sa > sb { + return -1 + } + if sa < sb { + return 1 + } + return 0 + }) + if len(dids) > limit { + dids = dids[:limit] + } + + ls, err := s.model.GetLatestLivestreams(limit, nil, dids) + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, "Failed to fetch livestreams") + } + + streams := make([]*placestream.Livestream_LivestreamView, len(ls)) + for i, l := range ls { + stream, err := l.ToLivestreamView() + if err != nil { + return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("Failed to convert livestream to streamplace livestream: %s", err)) + } + stream.ViewerCount = &placestream.Livestream_ViewerCount{ + LexiconTypeID: "place.stream.livestream#viewerCount", + Count: int64(s.bus.GetViewerCount(stream.Author.Did)), + } + streams[i] = stream + } + + liveUsers := &placestream.LiveGetLiveUsers_Output{Streams: streams} + s.LiveUsersCache.SetDefault(cacheKey, liveUsers) + return liveUsers, nil +} + +// getLiveUsersLatest returns live streams ordered by start time (newest first). +// Used for moderation tooling. +func (s *Server) getLiveUsersLatest(ctx context.Context, before string, limit int, cacheKey string) (*placestream.LiveGetLiveUsers_Output, error) { var beforeTime *time.Time if before != "" { parsedTime, err := time.Parse(time.RFC3339, before) @@ -183,27 +312,20 @@ func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before s } streams := make([]*placestream.Livestream_LivestreamView, len(ls)) - for i, l := range ls { stream, err := l.ToLivestreamView() if err != nil { return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("Failed to convert livestream to streamplace livestream: %s", err)) } - viewers := s.bus.GetViewerCount(stream.Author.Did) stream.ViewerCount = &placestream.Livestream_ViewerCount{ LexiconTypeID: "place.stream.livestream#viewerCount", - Count: int64(viewers), + Count: int64(s.bus.GetViewerCount(stream.Author.Did)), } streams[i] = stream } - liveUsers := &placestream.LiveGetLiveUsers_Output{ - Streams: streams, - } - - // Cache the result + liveUsers := &placestream.LiveGetLiveUsers_Output{Streams: streams} s.LiveUsersCache.SetDefault(cacheKey, liveUsers) - return liveUsers, nil } diff --git a/pkg/spxrpc/spxrpc.go b/pkg/spxrpc/spxrpc.go index 216b88497..479250316 100644 --- a/pkg/spxrpc/spxrpc.go +++ b/pkg/spxrpc/spxrpc.go @@ -36,6 +36,7 @@ type Server struct { OGImageCache *cache.Cache LiveUsersCache *cache.Cache GameSearchCache *cache.Cache + ScoreCache *cache.Cache ATSync *atproto.ATProtoSynchronizer statefulDB *statedb.StatefulDB bus *bus.Bus @@ -61,8 +62,9 @@ func NewServer(ctx context.Context, cli *config.CLI, model model.Model, stateful cli: cli, model: model, OGImageCache: cache.New(5*time.Minute, 10*time.Minute), - LiveUsersCache: cache.New(5*time.Second, 10*time.Second), + LiveUsersCache: cache.New(30*time.Second, 60*time.Second), GameSearchCache: cache.New(60*time.Second, 2*time.Minute), + ScoreCache: cache.New(30*time.Second, 60*time.Second), ATSync: atsync, statefulDB: statefulDB, bus: bus, -- 2.51.2 From e5757fe3f27e29914ae64639f31c364ca1690735 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Fri, 5 Jun 2026 10:18:03 -0700 Subject: [PATCH 25/28] player: default show chat --- js/app/components/mobile/player.tsx | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/js/app/components/mobile/player.tsx b/js/app/components/mobile/player.tsx index f176f999c..51558161e 100644 --- a/js/app/components/mobile/player.tsx +++ b/js/app/components/mobile/player.tsx @@ -82,7 +82,7 @@ function PlayerWithProvider( onTeleport?: (targetHandle: string, targetDID: string) => void; }, ) { - let [showChat, setShowChat] = useState(props.mode != "vod"); + let [showChat, setShowChat] = useState(true); if (props.mode === "vod") { showChat = false; } -- 2.51.2 From c5a10280c2531469337c27d2e5f19a735df4a058 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Fri, 5 Jun 2026 10:18:19 -0700 Subject: [PATCH 26/28] v0.11.4 --- js/app/package.json | 2 +- js/components/package.json | 2 +- js/docs/package.json | 2 +- lerna.json | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/js/app/package.json b/js/app/package.json index c4220e420..2365162d6 100644 --- a/js/app/package.json +++ b/js/app/package.json @@ -1,7 +1,7 @@ { "name": "@streamplace/app", "main": "./src/entrypoint.tsx", - "version": "0.11.3", + "version": "0.11.4", "runtimeVersion": "0.10.0", "scripts": { "start": "npx expo start -c --port 38081", diff --git a/js/components/package.json b/js/components/package.json index 6451cc4ba..46f0f2b8f 100644 --- a/js/components/package.json +++ b/js/components/package.json @@ -1,6 +1,6 @@ { "name": "@streamplace/components", - "version": "0.11.0", + "version": "0.11.4", "description": "Streamplace React (Native) Components", "main": "dist/index.js", "types": "src/index.tsx", diff --git a/js/docs/package.json b/js/docs/package.json index 6dd49c8f3..e2e430380 100644 --- a/js/docs/package.json +++ b/js/docs/package.json @@ -1,7 +1,7 @@ { "name": "streamplace-docs", "type": "module", - "version": "0.11.3", + "version": "0.11.4", "scripts": { "dev": "astro dev --host 0.0.0.0 --port 38082", "start": "astro dev --host 0.0.0.0 --port 38082", diff --git a/lerna.json b/lerna.json index cb16fbfcc..37197a1c4 100644 --- a/lerna.json +++ b/lerna.json @@ -1,5 +1,5 @@ { "$schema": "node_modules/lerna/schemas/lerna-schema.json", - "version": "0.11.3", + "version": "0.11.4", "npmClient": "pnpm" } -- 2.51.2 From 8c159cc0d5d0a219e4263859c4156055d290ccb2 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Fri, 5 Jun 2026 16:29:23 -0700 Subject: [PATCH 27/28] ci: restore linux build deps for release jobs (fixes deb-release 404) The release jobs (deb-release, release, npm-release, homebrew-release) download the linux .deb / .tar.gz files from the GitLab package registry, which are uploaded by build-linux-amd64 / build-linux-arm64. Those used to be ordered before the release jobs transitively, via the build-docker-* jobs' `needs`. When the build-docker-* jobs were removed (ECR push registry down), that transitive ordering went with them, so the release jobs could start before the .deb files were uploaded -> 404 on download. Depend on the linux build jobs directly in .release-template. Co-Authored-By: Claude Opus 4.8 --- .gitlab-ci.yml | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index 14f61d3a7..1a8aa7874 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -215,6 +215,12 @@ build-desktop-darwin: interruptible: true image: "$DOCKER_REPO:builder-$DOCKERFILE_HASH" needs: + # The linux build jobs upload the .deb / .tar.gz files that the release + # jobs download from the package registry, so they must finish first. + # They used to be pulled in transitively via the build-docker-* jobs; + # those were removed, so depend on them directly to avoid a download race. + - build-linux-amd64 + - build-linux-arm64 - build-darwin-amd64 - build-darwin-arm64 - build-android-debug -- 2.51.2 From 3e371ec808305bffc82f85eef5b02682665144e4 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 6 Jun 2026 14:06:07 -0700 Subject: [PATCH 28/28] v0.11.5 --- js/app/package.json | 2 +- js/docs/package.json | 2 +- lerna.json | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/js/app/package.json b/js/app/package.json index 2365162d6..a868bead7 100644 --- a/js/app/package.json +++ b/js/app/package.json @@ -1,7 +1,7 @@ { "name": "@streamplace/app", "main": "./src/entrypoint.tsx", - "version": "0.11.4", + "version": "0.11.5", "runtimeVersion": "0.10.0", "scripts": { "start": "npx expo start -c --port 38081", diff --git a/js/docs/package.json b/js/docs/package.json index e2e430380..00fdfb323 100644 --- a/js/docs/package.json +++ b/js/docs/package.json @@ -1,7 +1,7 @@ { "name": "streamplace-docs", "type": "module", - "version": "0.11.4", + "version": "0.11.5", "scripts": { "dev": "astro dev --host 0.0.0.0 --port 38082", "start": "astro dev --host 0.0.0.0 --port 38082", diff --git a/lerna.json b/lerna.json index 37197a1c4..770d3a03a 100644 --- a/lerna.json +++ b/lerna.json @@ -1,5 +1,5 @@ { "$schema": "node_modules/lerna/schemas/lerna-schema.json", - "version": "0.11.4", + "version": "0.11.5", "npmClient": "pnpm" }