diff --git a/docs/pinned-comment-plan.md b/docs/pinned-comment-plan.md new file mode 100644 index 00000000..a6babef4 --- /dev/null +++ b/docs/pinned-comment-plan.md @@ -0,0 +1,460 @@ +# Pinned Comment Feature Plan + +## Overview + +Allow the streamer or a delegated moderator to pin an existing chat message. The pinned comment appears prominently in the chat UI and auto-expires when a TTL is reached or the stream ends. + +## Design Decisions + +- **Pinnable content**: Existing chat messages only (referenced by AT-URI) +- **Who can pin**: Streamer + mods with new `message.pin` permission +- **Expiration**: Optional TTL (`expiresAt` datetime) + auto-clear at stream end (client-side) +- **Storage pattern**: Record in streamer's AT Protocol repo, modeled after `place.stream.chat.gate` +- **Single active pin**: Only one pinned comment at a time. Creating a new pin replaces the previous one. +- **Viewer dismiss**: Any viewer can hide the pin from their own view (local state only). The pin remains for everyone else. + +--- + +## 1. New Lexicon: `place.stream.chat.pinnedRecord` + +**File**: `lexicons/place/stream/chat/pinnedRecord.json` + +```json +{ + "lexicon": 1, + "id": "place.stream.chat.pinnedRecord", + "defs": { + "main": { + "type": "record", + "key": "tid", + "description": "Record pinning a chat message for prominent display.", + "record": { + "type": "object", + "required": ["pinnedMessage", "createdAt"], + "properties": { + "pinnedMessage": { + "type": "string", + "format": "at-uri", + "description": "AT-URI of the pinned chat message." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When this pin was created." + }, + "expiresAt": { + "type": "string", + "format": "datetime", + "description": "Optional expiration time. If set, the pin is considered inactive after this time." + } + } + } + } + } +} +``` + +--- + +## 2. Pinned Record View + Defs + +### 2a. Add `pinnedRecordView` to chat defs + +**File**: `lexicons/place/stream/chat/defs.json` (add new definition) + +```json +"pinnedRecordView": { + "type": "object", + "description": "View of a pinned chat record with hydrated message data.", + "required": ["uri", "cid", "record", "indexedAt"], + "properties": { + "uri": { "type": "string", "format": "at-uri" }, + "cid": { "type": "string", "format": "cid" }, + "record": { "type": "ref", "ref": "place.stream.chat.pinnedRecord" }, + "indexedAt": { "type": "string", "format": "datetime" }, + "pinnedBy": { "type": "ref", "ref": "app.bsky.actor.defs#profileViewBasic" }, + "message": { "type": "ref", "ref": "place.stream.chat.defs#messageView" } + } +} +``` + +This gives us a hydrated view that includes the full message data, so the frontend doesn't need to look it up from the chat index. + +### 2b. Go struct for bus delivery + +The bus will deliver a `PinnedRecordView` (not just the raw record) so the frontend can render the message text directly. Similar to how `ChatGate` records are published raw but the frontend processes them. + +--- + +## 3. New Permission: `message.pin` + +### 3a. Add to lexicon enum + +**File**: `lexicons/place/stream/moderation/permission.json` + +Add `"message.pin"` to the permissions enum: + +```json +"enum": ["ban", "hide", "livestream.manage", "message.pin"] +``` + +### 3b. Register in Go permission system + +**File**: `pkg/moderation/permissions.go` + +```go +const PermissionMessagePin = "message.pin" + +var ActionPermissions = map[string]string{ + // ... existing entries ... + "createPin": PermissionMessagePin, + "deletePin": PermissionMessagePin, +} +``` + +### 3c. Update frontend permissions + +**File**: `js/components/src/streamplace-store/moderation.tsx` + +Add `canPin: boolean` to `ModerationPermissions` interface. Derive from permissions array: + +```ts +canPin: isOwner || permissions.includes("message.pin"), +``` + +--- + +## 4. New RPC Procedures + +### 4a. Lexicon: `place.stream.moderation.createPin` + +**File**: `lexicons/place/stream/moderation/createPin.json` + +```json +{ + "lexicon": 1, + "id": "place.stream.moderation.createPin", + "defs": { + "main": { + "type": "procedure", + "description": "Pin a chat message on behalf of a streamer. Requires 'message.pin' permission. Creates a place.stream.chat.pinnedRecord in the streamer's repo, replacing any existing pin.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamer", "messageUri"], + "properties": { + "streamer": { "type": "string", "format": "did" }, + "messageUri": { "type": "string", "format": "at-uri" }, + "expiresAt": { "type": "string", "format": "datetime" } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "cid"], + "properties": { + "uri": { "type": "string", "format": "at-uri" }, + "cid": { "type": "string", "format": "cid" } + } + } + }, + "errors": [ + { "name": "Unauthorized" }, + { "name": "Forbidden" }, + { "name": "SessionNotFound" } + ] + } + } +} +``` + +### 4b. Lexicon: `place.stream.moderation.deletePin` + +**File**: `lexicons/place/stream/moderation/deletePin.json` + +Same pattern as `deleteGate`: + +- Input: `streamer` (DID), `pinUri` (at-uri) +- Output: empty +- Errors: Unauthorized, Forbidden, SessionNotFound + +### 4c. Go handlers + +**File**: `pkg/spxrpc/place_stream_moderation.go` + +Add `handlePlaceStreamModerationCreatePin`: + +1. Validate input (DID, AT-URI) +2. `GetDelegatedModerationContext(ctx, input.Streamer, "createPin")` +3. Before creating: list existing `place.stream.chat.pinnedRecord` records in streamer's repo, delete any existing ones (single-pin semantics) +4. Build `streamplace.ChatPinnedRecord` struct +5. Create via `com.atproto.repo.createRecord` on streamer's repo +6. Audit log +7. Return URI + CID + +Add `handlePlaceStreamModerationDeletePin`: + +1. Validate input +2. `GetDelegatedModerationContext(ctx, input.Streamer, "deletePin")` +3. Extract rkey, delete record via `com.atproto.repo.deleteRecord` +4. Audit log + +### 4d. Register handlers + +In the XRPC server setup where `createGate`/`deleteGate` are registered, add `createPin` and `deletePin`. + +--- + +## 5. Backend: Model + DB + +**File**: `pkg/model/pinned_record.go` + +```go +type PinnedRecord struct { + RKey string `gorm:"primaryKey;column:rkey"` + CID string `gorm:"column:cid"` + RepoDID string `gorm:"column:repo_did"` + Repo *Repo `gorm:"foreignKey:DID;references:RepoDID"` + PinnedMessage string `gorm:"column:pinned_message"` + ExpiresAt *time.Time `gorm:"column:expires_at"` + CreatedAt time.Time `gorm:"column:created_at"` +} +``` + +Methods: + +- `CreatePinnedRecord(ctx, pin)` - insert +- `GetPinnedRecord(ctx, rkey)` - single lookup +- `DeletePinnedRecord(ctx, rkey)` - delete by rkey +- `GetActivePinnedRecord(ctx, streamerDID)` - returns the most recent non-expired pin for a streamer +- `DeleteAllPinnedRecords(ctx, streamerDID)` - bulk delete (called before creating new pin) + +Add `PinnedRecord{}` to the AutoMigrate list in `pkg/model/model.go`. + +--- + +## 6. Constants + Firehose + +### 6a. Constants + +**File**: `pkg/constants/constants.go` + +Add: `var PLACE_STREAM_CHAT_PINNED_RECORD = "place.stream.chat.pinnedRecord"` + +### 6b. Sync (create/update) + +**File**: `pkg/atproto/sync.go` - `handleCreateUpdate` + +Add case `*streamplace.ChatPinnedRecord`: + +- Sync bluesky repo +- Delete existing pinned records for this streamer (single-pin enforcement at DB level) +- Create new `PinnedRecord` model entry +- Build a hydrated view (resolve the pinned message, include author info) and publish to bus on streamer's channel + +### 6c. Firehose (delete) + +**File**: `pkg/atproto/firehose.go` - `EvtKindDeleteRecord` + +Add handling for `constants.PLACE_STREAM_CHAT_PINNED_RECORD`: + +- `DeletePinnedRecord(ctx, rkey)` +- Publish deletion marker to bus + +--- + +## 7. Bus / WebSocket Delivery + +### Create event + +Publish a `PinnedRecordView`-like object to the streamer's bus channel: + +```json +{ + "$type": "place.stream.chat.defs#pinnedRecordView", + "uri": "at://...", + "cid": "...", + "record": { + "pinnedMessage": "at://...", + "createdAt": "...", + "expiresAt": "..." + }, + "indexedAt": "...", + "message": { + /* hydrated messageView with author, text, facets, etc */ + } +} +``` + +### Delete event + +Publish a deletion marker: + +```json +{ + "$type": "place.stream.chat.pinnedRecord", + "deleted": true, + "rkey": "..." +} +``` + +--- + +## 8. Frontend State + +### LivestreamState additions + +**File**: `js/components/src/livestream-store/livestream-state.tsx` + +```ts +pinnedComment: PinnedRecordView | null; +``` + +Where `PinnedRecordView` is a type with the hydrated message + pin metadata. + +### Store init + +**File**: `js/components/src/livestream-store/livestream-store.tsx` + +Add `pinnedComment: null` to initial state. + +### Chat hooks + +**File**: `js/components/src/livestream-store/chat.tsx` + +Add: + +- `usePinnedComment()` - selector: `state.pinnedComment` +- `usePinChatMessage()` - hook to pin a message: + - If streamer: direct `com.atproto.repo.createRecord` for `place.stream.chat.pinnedRecord` + - If mod: call `place.stream.moderation.createPin` XRPC + - On success, the bus event will update state automatically +- `useUnpinChatMessage()` - hook to unpin: + - If streamer: direct `com.atproto.repo.deleteRecord` + - If mod: call `place.stream.moderation.deletePin` XRPC + +### WebSocket consumer + +**File**: `js/components/src/livestream-store/websocket-consumer.tsx` + +Add handler for `place.stream.chat.defs#pinnedRecordView`: + +- Set `state.pinnedComment` to the hydrated view + +Add handler for `place.stream.chat.pinnedRecord` with `deleted: true`: + +- Set `state.pinnedComment` to `null` + +### Expiration + +In the Chat component, `useEffect` checks `pinnedComment.record.expiresAt`. If current time passes it, clear `pinnedComment` locally. Also clear when `livestream.record.endedAt` is set (stream ended). + +--- + +## 9. Frontend UI + +### Pinned comment as persistent stream notification + +The pinned comment uses the existing `StreamNotificationProvider` + `streamNotificationManager` system. + +**File**: `js/components/src/components/stream-notification/pinned-comment-notification.tsx` (new) + +A custom render function passed to `streamNotificationManager.show()` with `duration: 0` (manual dismiss only), making it a persistent notification that sits at the top of the notification stack. Other temporary notifications (teleport, etc.) will appear above/below it and dismiss independently. + +The render function receives `(isExiting, onDismiss, startTime)` and renders: + +- Pin icon + author name (with chatProfile color) + message text (with facets rendered) +- Two dismiss actions: + - **Unpin** (X icon, only visible if `canPin`): calls `useUnpinChatMessage()`, removes the pin globally, then calls `onDismiss("user")` to dismiss the notification + - **Hide** (eye-off icon, visible to all viewers): calls `onDismiss("user")` to dismiss the notification locally only. Does NOT call the server - the pin remains for other viewers +- Does NOT link to or scroll to the original message (chat history is limited) +- Auto-dismisses when `expiresAt` is reached (a `useEffect` watching the pinned comment state calls `streamNotificationManager.requestDismiss("pinned-comment", "auto")`) +- Auto-dismisses when stream ends + +**Integration point**: The `TeleportWatcher` / livestream provider already manages notification lifecycle. The pinned comment notification is managed similarly - when `state.pinnedComment` changes in the Zustand store, a hook or effect calls `streamNotificationManager.show()` to create/update the notification, or `streamNotificationManager.hide()` to remove it. + +The notification ID should be fixed (e.g., `"pinned-comment"`) so that replacing a pin updates the same notification rather than creating a new one. + +### Pin action in mod menu + +**File**: `js/components/src/components/chat/mod-view.tsx` + +Add "Pin this message" in the moderation actions group (visible when `canPin` is true): + +```tsx +{ + modPermissions.canPin && ( + { + /* pin message */ + }} + > + Pin this message + + ); +} +``` + +--- + +## 10. Moderator Panel + +**File**: `js/components/src/components/dashboard/moderator-panel.tsx` + +Add `"message.pin"` to the list of assignable permissions when adding/editing moderators. + +--- + +## 11. Code Generation + +After creating lexicon JSON files, run `lexgen` to generate Go types: + +- `pkg/streamplace/chatpinnedrecord.go` (auto-generated) +- Update generated TypeScript types in the `streamplace` package + +--- + +## File Change Summary + +| File | Action | +| ---------------------------------------------------------------------------------- | ------------------------------------------------------- | +| `lexicons/place/stream/chat/pinnedRecord.json` | **New** | +| `lexicons/place/stream/chat/defs.json` | **Edit** - add `pinnedRecordView` | +| `lexicons/place/stream/moderation/createPin.json` | **New** | +| `lexicons/place/stream/moderation/deletePin.json` | **New** | +| `lexicons/place/stream/moderation/permission.json` | **Edit** - add `message.pin` | +| `pkg/constants/constants.go` | **Edit** | +| `pkg/moderation/permissions.go` | **Edit** | +| `pkg/model/pinned_record.go` | **New** | +| `pkg/model/model.go` | **Edit** | +| `pkg/atproto/sync.go` | **Edit** | +| `pkg/atproto/firehose.go` | **Edit** | +| `pkg/spxrpc/place_stream_moderation.go` | **Edit** | +| `js/components/src/streamplace-store/moderation.tsx` | **Edit** | +| `js/components/src/streamplace-store/block.tsx` | **Edit** - add pin/unpin hooks | +| `js/components/src/livestream-store/livestream-state.tsx` | **Edit** | +| `js/components/src/livestream-store/livestream-store.tsx` | **Edit** | +| `js/components/src/livestream-store/chat.tsx` | **Edit** | +| `js/components/src/livestream-store/websocket-consumer.tsx` | **Edit** | +| `js/components/src/components/stream-notification/pinned-comment-notification.tsx` | **New** - notification render function | +| `js/components/src/livestream-provider/index.tsx` | **Edit** - manage pinned comment notification lifecycle | +| `js/components/src/components/chat/mod-view.tsx` | **Edit** - add pin action | +| `js/components/src/components/dashboard/moderator-panel.tsx` | **Edit** | + +--- + +## Implementation Order + +1. Lexicons (pinnedRecord, createPin, deletePin, permission, defs) +2. Code generation (lexgen) +3. Backend: constants, permissions, model, DB migration +4. Backend: sync/firehose integration +5. Backend: XRPC handlers + registration +6. Frontend: state types, store init, hooks +7. Frontend: WebSocket consumer handlers +8. Frontend: pinned-comment notification render function + lifecycle management +9. Frontend: mod-view pin action +10. Frontend: moderator panel permission diff --git a/js/components/src/components/chat/mod-view.tsx b/js/components/src/components/chat/mod-view.tsx index cdcaffae..4d0431be 100644 --- a/js/components/src/components/chat/mod-view.tsx +++ b/js/components/src/components/chat/mod-view.tsx @@ -17,6 +17,7 @@ import { ChatMessageViewHydrated } from "streamplace"; import { useDeleteChatMessage, useLivestreamStore, + usePinChatMessage, } from "../../livestream-store"; import { useStreamplaceStore } from "../../streamplace-store"; import { formatHandle, formatHandleWithAt } from "../../utils/format-handle"; @@ -55,6 +56,7 @@ export const ModView = forwardRef(() => { let [messageRemoved, setMessageRemoved] = useState(false); let { createBlock, isLoading: isBlockLoading } = useCreateBlockRecord(); let { createHideChat, isLoading: isHideLoading } = useCreateHideChatRecord(); + const pinChatMessage = usePinChatMessage(); const setReportModalOpen = usePlayerStore((x) => x.setReportModalOpen); const setReportSubject = usePlayerStore((x) => x.setReportSubject); @@ -91,6 +93,7 @@ export const ModView = forwardRef(() => { message && agent?.did && ((modPermissions.canHide && message.author.did !== streamerDID) || + (modPermissions.canPin && message.author.did !== streamerDID) || (modPermissions.canBan && message.author.did !== agent.did && message.author.did !== streamerDID)) @@ -124,6 +127,7 @@ export const ModView = forwardRef(() => { setMessageRemoved={setMessageRemoved} createHideChat={createHideChat} createBlock={createBlock} + pinChatMessage={pinChatMessage} toast={toast} setReportModalOpen={setReportModalOpen} setReportSubject={setReportSubject} @@ -148,6 +152,11 @@ interface ModViewContentProps { setMessageRemoved: (removed: boolean) => void; createHideChat: (uri: string, streamerDID?: string) => Promise; createBlock: (did: string, streamerDID?: string) => Promise; + pinChatMessage: ( + messageUri: string, + streamerDID: string, + expiresAt?: string, + ) => Promise; toast: ReturnType; setReportModalOpen: (open: boolean) => void; setReportSubject: (subject: any) => void; @@ -166,6 +175,7 @@ function ModViewContent({ setMessageRemoved, createHideChat, createBlock, + pinChatMessage, toast, setReportModalOpen, setReportSubject, @@ -223,6 +233,27 @@ function ModViewContent({ )} + {modPermissions.canPin && message.author.did !== streamerDID && ( + { + if (!streamerDID) return; + pinChatMessage(message.uri, streamerDID) + .then(() => { + toast.show("Comment pinned", "", { duration: 3 }); + onOpenChange?.(false); + }) + .catch((e) => { + toast.show( + "Error pinning comment", + e instanceof Error ? e.message : "Failed to pin", + { duration: 5 }, + ); + }); + }} + > + Pin this message + + )} {modPermissions.canBan && agent?.did && message.author.did !== agent.did && diff --git a/js/components/src/components/dashboard/moderator-panel.tsx b/js/components/src/components/dashboard/moderator-panel.tsx index 8916610b..f13ce7d2 100644 --- a/js/components/src/components/dashboard/moderator-panel.tsx +++ b/js/components/src/components/dashboard/moderator-panel.tsx @@ -372,6 +372,7 @@ function AddModeratorDialog({ ban: false, hide: false, "livestream.manage": false, + "message.pin": false, }); const [error, setError] = useState(null); const toast = useToast(); @@ -380,7 +381,12 @@ function AddModeratorDialog({ useEffect(() => { if (!visible) { setModeratorDID(""); - setPermissions({ ban: false, hide: false, "livestream.manage": false }); + setPermissions({ + ban: false, + hide: false, + "livestream.manage": false, + "message.pin": false, + }); setError(null); } }, [visible]); @@ -401,7 +407,12 @@ function AddModeratorDialog({ const selectedPermissions = Object.entries(permissions) .filter(([_, enabled]) => enabled) - .map(([perm]) => perm) as ("ban" | "hide" | "livestream.manage")[]; + .map(([perm]) => perm) as ( + | "ban" + | "hide" + | "livestream.manage" + | "message.pin" + )[]; if (selectedPermissions.length === 0) { setError("Please select at least one permission"); diff --git a/js/components/src/components/ui/resizeable.tsx b/js/components/src/components/ui/resizeable.tsx index a59d07ed..fc44e9e8 100644 --- a/js/components/src/components/ui/resizeable.tsx +++ b/js/components/src/components/ui/resizeable.tsx @@ -26,8 +26,8 @@ const AnimatedView = Animated.createAnimatedComponent(View); const { height: SCREEN_HEIGHT } = Dimensions.get("window"); const TIMING_CONFIG = { - duration: 300, - easing: Easing.inOut(Easing.quad), + duration: 400, + easing: Easing.out(Easing.quad), }; type ResizableChatSheetProps = { diff --git a/js/components/src/livestream-store/chat.tsx b/js/components/src/livestream-store/chat.tsx index de064166..2ece439b 100644 --- a/js/components/src/livestream-store/chat.tsx +++ b/js/components/src/livestream-store/chat.tsx @@ -437,3 +437,95 @@ export const useReportChatMessage = () => { }; export const reduceChat = reduceChatIncremental; + +export const usePinChatMessage = () => { + const agent = usePDSAgent(); + const store = getStoreFromContext(); + + return async ( + messageUri: string, + streamerDID: string, + expiresAt?: string, + ) => { + if (!agent || !agent.did) { + throw new Error("No PDS agent or user DID found"); + } + + // 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, + createdAt: new Date().toISOString(), + ...(expiresAt ? { expiresAt } : {}), + }; + + const result = await agent.com.atproto.repo.createRecord({ + repo: streamerDID, + collection: "place.stream.chat.pinnedRecord", + record, + }); + return result; + } + + // Otherwise, use delegated moderation endpoint + const result = await agent.place.stream.moderation.createPin({ + streamer: streamerDID, + messageUri, + ...(expiresAt ? { expiresAt } : {}), + }); + return result; + }; +}; + +export const useUnpinChatMessage = () => { + const agent = usePDSAgent(); + const store = getStoreFromContext(); + + return async (pinUri: string, streamerDID: string) => { + if (!agent || !agent.did) { + throw new Error("No PDS agent or user DID found"); + } + + // If streamer, delete directly + if (agent.did === streamerDID) { + const rkey = pinUri.split("/").pop(); + if (!rkey) { + throw new Error("Invalid pin URI"); + } + + await agent.com.atproto.repo.deleteRecord({ + repo: streamerDID, + collection: "place.stream.chat.pinnedRecord", + rkey, + }); + // Optimistically clear the pinned comment + store.setState({ pinnedComment: null }); + return; + } + + // Otherwise, use delegated moderation endpoint + await agent.place.stream.moderation.deletePin({ + streamer: streamerDID, + pinUri, + }); + // Optimistically clear the pinned comment + store.setState({ pinnedComment: null }); + }; +}; diff --git a/js/components/src/livestream-store/livestream-state.tsx b/js/components/src/livestream-store/livestream-state.tsx index 9f7a4ba6..2a14b6a7 100644 --- a/js/components/src/livestream-store/livestream-state.tsx +++ b/js/components/src/livestream-store/livestream-state.tsx @@ -2,6 +2,7 @@ import { AppBskyActorDefs } from "@atproto/api"; import { ChatMessageViewHydrated, LivestreamViewHydrated, + PinnedRecordViewHydrated, PlaceStreamDefs, PlaceStreamLiveTeleport, PlaceStreamModerationPermission, @@ -28,6 +29,7 @@ export interface LivestreamState { setActiveTeleportUri: (uri: string | null) => void; websocketConnected: boolean; hasReceivedSegment: boolean; + pinnedComment: PinnedRecordViewHydrated | null; moderationPermissions: PlaceStreamModerationPermission.Record[]; setModerationPermissions: ( permissions: PlaceStreamModerationPermission.Record[], diff --git a/js/components/src/livestream-store/livestream-store.tsx b/js/components/src/livestream-store/livestream-store.tsx index 97e4ddb7..61b205be 100644 --- a/js/components/src/livestream-store/livestream-store.tsx +++ b/js/components/src/livestream-store/livestream-store.tsx @@ -27,6 +27,7 @@ export const makeLivestreamStore = (): StoreApi => { setActiveTeleportUri: (uri) => set({ activeTeleportUri: uri }), websocketConnected: false, hasReceivedSegment: false, + pinnedComment: null, moderationPermissions: [], setModerationPermissions: (perms) => set({ moderationPermissions: perms }), localLivestreamURI: null, @@ -60,6 +61,9 @@ export const useHandleWebsocketMessages = () => { export const useChat = () => useLivestreamStore((x) => x.chat); +export const usePinnedComment = () => + useLivestreamStore((x) => x.pinnedComment); + export const useProfile = () => useLivestreamStore((x) => x.profile); export const useViewers = () => useLivestreamStore((x) => x.viewers); diff --git a/js/components/src/livestream-store/websocket-consumer.tsx b/js/components/src/livestream-store/websocket-consumer.tsx index 1b188f34..0cc3ba6a 100644 --- a/js/components/src/livestream-store/websocket-consumer.tsx +++ b/js/components/src/livestream-store/websocket-consumer.tsx @@ -2,6 +2,7 @@ import { AppBskyActorDefs } from "@atproto/api"; import { ChatMessageViewHydrated, LivestreamViewHydrated, + PinnedRecordViewHydrated, PlaceStreamChatDefs, PlaceStreamChatGate, PlaceStreamChatMessage, @@ -123,6 +124,20 @@ export const handleWebSocketMessages = ( pendingHides: newPendingHides, }; state = reduceChat(state, [], [], [hiddenMessageUri]); + } else if (PlaceStreamChatDefs.isPinnedRecordView(message)) { + const pinnedView = message as PinnedRecordViewHydrated; + state = { + ...state, + pinnedComment: pinnedView, + }; + } else if ( + (message as any).$type === "place.stream.chat.pinnedRecord" && + (message as any).deleted === true + ) { + state = { + ...state, + pinnedComment: null, + }; } else if (PlaceStreamLiveTeleport.isRecord(message)) { const teleportRecord = message as PlaceStreamLiveTeleport.Record; state = { diff --git a/js/components/src/streamplace-store/moderation.tsx b/js/components/src/streamplace-store/moderation.tsx index 7169b8d2..34f071f9 100644 --- a/js/components/src/streamplace-store/moderation.tsx +++ b/js/components/src/streamplace-store/moderation.tsx @@ -6,6 +6,7 @@ import { usePDSAgent } from "./xrpc"; export interface ModerationPermissions { canBan: boolean; canHide: boolean; + canPin: boolean; canManageLivestream: boolean; isOwner: boolean; isLoading: boolean; @@ -177,6 +178,7 @@ export function useCanModerate( return { canBan: isOwner || permissions.includes("ban"), canHide: isOwner || permissions.includes("hide"), + canPin: isOwner || permissions.includes("message.pin"), canManageLivestream: isOwner || permissions.includes("livestream.manage"), isOwner, isLoading, diff --git a/js/components/src/streamplace-store/moderator-management.tsx b/js/components/src/streamplace-store/moderator-management.tsx index 3888da20..15ce5e96 100644 --- a/js/components/src/streamplace-store/moderator-management.tsx +++ b/js/components/src/streamplace-store/moderator-management.tsx @@ -81,7 +81,7 @@ export function useListModerators(): ListModeratorsResult { interface AddModeratorParams { moderatorDID: string; - permissions: ("ban" | "hide" | "livestream.manage")[]; + permissions: ("ban" | "hide" | "livestream.manage" | "message.pin")[]; expirationTime?: string; // ISO 8601 datetime string } diff --git a/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md index 382fe718..bb73679b 100644 --- a/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md +++ b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md @@ -29,6 +29,27 @@ description: Reference for the place.stream.chat.defs lexicon --- + + +### `pinnedRecordView` + +**Type:** `object` + +View of a pinned chat record with hydrated message data. + +**Properties:** + +| Name | Type | Req'd | Description | Constraints | +| ----------- | ------------------------------------------------------------------------------------------------------------------------------------------------ | ----- | ----------- | ------------------ | +| `uri` | `string` | ✅ | | Format: `at-uri` | +| `cid` | `string` | ✅ | | Format: `cid` | +| `record` | [`place.stream.chat.pinnedRecord`](/lex-reference/place-stream-chat-pinnedrecord) | ✅ | | | +| `indexedAt` | `string` | ✅ | | Format: `datetime` | +| `pinnedBy` | [`app.bsky.actor.defs#profileViewBasic`](https://github.com/bluesky-social/atproto/tree/main/lexicons/app/bsky/actor/defs.json#profileViewBasic) | ❌ | | | +| `message` | [`#messageView`](#messageview) | ❌ | | | + +--- + ## Lexicon Source ```json @@ -81,6 +102,37 @@ description: Reference for the place.stream.chat.defs lexicon } } } + }, + "pinnedRecordView": { + "type": "object", + "description": "View of a pinned chat record with hydrated message data.", + "required": ["uri", "cid", "record", "indexedAt"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri" + }, + "cid": { + "type": "string", + "format": "cid" + }, + "record": { + "type": "ref", + "ref": "place.stream.chat.pinnedRecord" + }, + "indexedAt": { + "type": "string", + "format": "datetime" + }, + "pinnedBy": { + "type": "ref", + "ref": "app.bsky.actor.defs#profileViewBasic" + }, + "message": { + "type": "ref", + "ref": "#messageView" + } + } } } } diff --git a/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-pinnedrecord.md b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-pinnedrecord.md new file mode 100644 index 00000000..8935e957 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-pinnedrecord.md @@ -0,0 +1,65 @@ +--- +title: place.stream.chat.pinnedRecord +description: Reference for the place.stream.chat.pinnedRecord lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `record` + +Record pinning a chat message for prominent display. + +**Record Key:** `tid` + +**Record Properties:** + +| Name | Type | Req'd | Description | Constraints | +| --------------- | -------- | ----- | --------------------------------------------------------------------------------- | ------------------ | +| `pinnedMessage` | `string` | ✅ | AT-URI of the pinned chat message. | Format: `at-uri` | +| `createdAt` | `string` | ✅ | When this pin was created. | Format: `datetime` | +| `expiresAt` | `string` | ❌ | Optional expiration time. If set, the pin is considered inactive after this time. | Format: `datetime` | + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.chat.pinnedRecord", + "defs": { + "main": { + "type": "record", + "key": "tid", + "description": "Record pinning a chat message for prominent display.", + "record": { + "type": "object", + "required": ["pinnedMessage", "createdAt"], + "properties": { + "pinnedMessage": { + "type": "string", + "format": "at-uri", + "description": "AT-URI of the pinned chat message." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When this pin was created." + }, + "expiresAt": { + "type": "string", + "format": "datetime", + "description": "Optional expiration time. If set, the pin is considered inactive after this time." + } + } + } + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-createpin.md b/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-createpin.md new file mode 100644 index 00000000..5c1ce46e --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-createpin.md @@ -0,0 +1,123 @@ +--- +title: place.stream.moderation.createPin +description: Reference for the place.stream.moderation.createPin lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Pin a chat message on behalf of a streamer. Requires 'message.pin' permission. Creates a place.stream.chat.pinnedRecord in the streamer's repo, replacing any existing pin. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ------------ | -------- | ----- | -------------------------------------- | ------------------ | +| `streamer` | `string` | ✅ | The DID of the streamer. | Format: `did` | +| `messageUri` | `string` | ✅ | The AT-URI of the chat message to pin. | Format: `at-uri` | +| `expiresAt` | `string` | ❌ | Optional expiration time for this pin. | Format: `datetime` | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ----- | -------- | ----- | ---------------------------------------- | ---------------- | +| `uri` | `string` | ✅ | The AT-URI of the created pinned record. | Format: `at-uri` | +| `cid` | `string` | ✅ | The CID of the created pinned record. | Format: `cid` | + +**Possible Errors:** + +- `Unauthorized`: The request lacks valid authentication credentials. +- `Forbidden`: The caller does not have permission to pin messages for this streamer. +- `SessionNotFound`: The streamer's OAuth session could not be found or is invalid. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.moderation.createPin", + "defs": { + "main": { + "type": "procedure", + "description": "Pin a chat message on behalf of a streamer. Requires 'message.pin' permission. Creates a place.stream.chat.pinnedRecord in the streamer's repo, replacing any existing pin.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamer", "messageUri"], + "properties": { + "streamer": { + "type": "string", + "format": "did", + "description": "The DID of the streamer." + }, + "messageUri": { + "type": "string", + "format": "at-uri", + "description": "The AT-URI of the chat message to pin." + }, + "expiresAt": { + "type": "string", + "format": "datetime", + "description": "Optional expiration time for this pin." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "cid"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri", + "description": "The AT-URI of the created pinned record." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "The CID of the created pinned record." + } + } + } + }, + "errors": [ + { + "name": "Unauthorized", + "description": "The request lacks valid authentication credentials." + }, + { + "name": "Forbidden", + "description": "The caller does not have permission to pin messages for this streamer." + }, + { + "name": "SessionNotFound", + "description": "The streamer's OAuth session could not be found or is invalid." + } + ] + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-deletepin.md b/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-deletepin.md new file mode 100644 index 00000000..49c4b83c --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-deletepin.md @@ -0,0 +1,101 @@ +--- +title: place.stream.moderation.deletePin +description: Reference for the place.stream.moderation.deletePin lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `procedure` + +Unpin a pinned chat message on behalf of a streamer. Requires 'message.pin' permission. Deletes the place.stream.chat.pinnedRecord from the streamer's repo. + +**Parameters:** _(None defined)_ + +**Input:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +| Name | Type | Req'd | Description | Constraints | +| ---------- | -------- | ----- | ------------------------------------------ | ---------------- | +| `streamer` | `string` | ✅ | The DID of the streamer. | Format: `did` | +| `pinUri` | `string` | ✅ | The AT-URI of the pinned record to delete. | Format: `at-uri` | + +**Output:** + +- **Encoding:** `application/json` +- **Schema:** + +**Schema Type:** `object` + +_(No properties defined)_ +**Possible Errors:** + +- `Unauthorized`: The request lacks valid authentication credentials. +- `Forbidden`: The caller does not have permission to unpin messages for this streamer. +- `SessionNotFound`: The streamer's OAuth session could not be found or is invalid. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.moderation.deletePin", + "defs": { + "main": { + "type": "procedure", + "description": "Unpin a pinned chat message on behalf of a streamer. Requires 'message.pin' permission. Deletes the place.stream.chat.pinnedRecord from the streamer's repo.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamer", "pinUri"], + "properties": { + "streamer": { + "type": "string", + "format": "did", + "description": "The DID of the streamer." + }, + "pinUri": { + "type": "string", + "format": "at-uri", + "description": "The AT-URI of the pinned record to delete." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "properties": {} + } + }, + "errors": [ + { + "name": "Unauthorized", + "description": "The request lacks valid authentication credentials." + }, + { + "name": "Forbidden", + "description": "The caller does not have permission to unpin messages for this streamer." + }, + { + "name": "SessionNotFound", + "description": "The streamer's OAuth session could not be found or is invalid." + } + ] + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-permission.md b/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-permission.md index 986eafb7..c45f3a9b 100644 --- a/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-permission.md +++ b/js/docs/src/content/docs/lex-reference/moderation/place-stream-moderation-permission.md @@ -52,7 +52,7 @@ Record granting moderation permissions to a user for this streamer's content. "type": "array", "items": { "type": "string", - "enum": ["ban", "hide", "livestream.manage"] + "enum": ["ban", "hide", "livestream.manage", "message.pin"] }, "description": "Array of permissions granted to this moderator. 'ban' covers blocks/bans (with optional expiration), 'hide' covers message gates, 'livestream.manage' allows updating livestream metadata." }, diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index dbcdff4a..f9b0ae1c 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -979,6 +979,96 @@ } } }, + "/xrpc/place.stream.moderation.createPin": { + "post": { + "summary": "Pin a chat message on behalf of a streamer. Requires 'message.pin' permission. Creates a place.stream.chat.pinnedRecord in the streamer's repo, replacing any existing pin.", + "operationId": "place.stream.moderation.createPin", + "tags": ["place.stream.moderation"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "uri": { + "type": "string", + "description": "The AT-URI of the created pinned record.", + "format": "uri" + }, + "cid": { + "type": "string", + "description": "The CID of the created pinned record.", + "format": "cid" + } + }, + "required": ["uri", "cid"] + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "Unauthorized" + }, + { + "const": "Forbidden" + }, + { + "const": "SessionNotFound" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "streamer": { + "type": "string", + "description": "The DID of the streamer.", + "format": "did" + }, + "messageUri": { + "type": "string", + "description": "The AT-URI of the chat message to pin.", + "format": "uri" + }, + "expiresAt": { + "type": "string", + "description": "Optional expiration time for this pin.", + "format": "date-time" + } + }, + "required": ["streamer", "messageUri"] + } + } + } + } + } + }, "/xrpc/place.stream.moderation.deleteBlock": { "post": { "summary": "Delete a block (unban) on behalf of a streamer. Requires 'ban' permission. Deletes an app.bsky.graph.block record from the streamer's repository.", @@ -1125,6 +1215,79 @@ } } }, + "/xrpc/place.stream.moderation.deletePin": { + "post": { + "summary": "Unpin a pinned chat message on behalf of a streamer. Requires 'message.pin' permission. Deletes the place.stream.chat.pinnedRecord from the streamer's repo.", + "operationId": "place.stream.moderation.deletePin", + "tags": ["place.stream.moderation"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": {} + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "Unauthorized" + }, + { + "const": "Forbidden" + }, + { + "const": "SessionNotFound" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "streamer": { + "type": "string", + "description": "The DID of the streamer.", + "format": "did" + }, + "pinUri": { + "type": "string", + "description": "The AT-URI of the pinned record to delete.", + "format": "uri" + } + }, + "required": ["streamer", "pinUri"] + } + } + } + } + } + }, "/xrpc/place.stream.moderation.updateLivestream": { "post": { "summary": "Update livestream metadata on behalf of a streamer. Requires 'livestream.manage' permission. Updates a place.stream.livestream record in the streamer's repository.", diff --git a/js/streamplace/src/useful-types.ts b/js/streamplace/src/useful-types.ts index 2533a434..8ddf81c4 100644 --- a/js/streamplace/src/useful-types.ts +++ b/js/streamplace/src/useful-types.ts @@ -1,6 +1,7 @@ import { PlaceStreamChatDefs, PlaceStreamChatMessage, + PlaceStreamChatPinnedRecord, PlaceStreamLivestream, } from "./lexicons"; @@ -13,3 +14,8 @@ export interface ChatMessageViewHydrated extends PlaceStreamChatDefs.MessageView { record: PlaceStreamChatMessage.Record; } + +export interface PinnedRecordViewHydrated + extends PlaceStreamChatDefs.PinnedRecordView { + record: PlaceStreamChatPinnedRecord.Record; +} diff --git a/lexicons/place/stream/chat/defs.json b/lexicons/place/stream/chat/defs.json index da097e62..aa1dcd75 100644 --- a/lexicons/place/stream/chat/defs.json +++ b/lexicons/place/stream/chat/defs.json @@ -36,6 +36,25 @@ } } } + }, + "pinnedRecordView": { + "type": "object", + "description": "View of a pinned chat record with hydrated message data.", + "required": ["uri", "cid", "record", "indexedAt"], + "properties": { + "uri": { "type": "string", "format": "at-uri" }, + "cid": { "type": "string", "format": "cid" }, + "record": { "type": "ref", "ref": "place.stream.chat.pinnedRecord" }, + "indexedAt": { "type": "string", "format": "datetime" }, + "pinnedBy": { + "type": "ref", + "ref": "app.bsky.actor.defs#profileViewBasic" + }, + "message": { + "type": "ref", + "ref": "#messageView" + } + } } } } diff --git a/lexicons/place/stream/chat/pinnedRecord.json b/lexicons/place/stream/chat/pinnedRecord.json new file mode 100644 index 00000000..f43ac525 --- /dev/null +++ b/lexicons/place/stream/chat/pinnedRecord.json @@ -0,0 +1,32 @@ +{ + "lexicon": 1, + "id": "place.stream.chat.pinnedRecord", + "defs": { + "main": { + "type": "record", + "key": "tid", + "description": "Record pinning a chat message for prominent display.", + "record": { + "type": "object", + "required": ["pinnedMessage", "createdAt"], + "properties": { + "pinnedMessage": { + "type": "string", + "format": "at-uri", + "description": "AT-URI of the pinned chat message." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When this pin was created." + }, + "expiresAt": { + "type": "string", + "format": "datetime", + "description": "Optional expiration time. If set, the pin is considered inactive after this time." + } + } + } + } + } +} diff --git a/lexicons/place/stream/moderation/createPin.json b/lexicons/place/stream/moderation/createPin.json new file mode 100644 index 00000000..b7be8c50 --- /dev/null +++ b/lexicons/place/stream/moderation/createPin.json @@ -0,0 +1,67 @@ +{ + "lexicon": 1, + "id": "place.stream.moderation.createPin", + "defs": { + "main": { + "type": "procedure", + "description": "Pin a chat message on behalf of a streamer. Requires 'message.pin' permission. Creates a place.stream.chat.pinnedRecord in the streamer's repo, replacing any existing pin.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamer", "messageUri"], + "properties": { + "streamer": { + "type": "string", + "format": "did", + "description": "The DID of the streamer." + }, + "messageUri": { + "type": "string", + "format": "at-uri", + "description": "The AT-URI of the chat message to pin." + }, + "expiresAt": { + "type": "string", + "format": "datetime", + "description": "Optional expiration time for this pin." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["uri", "cid"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri", + "description": "The AT-URI of the created pinned record." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "The CID of the created pinned record." + } + } + } + }, + "errors": [ + { + "name": "Unauthorized", + "description": "The request lacks valid authentication credentials." + }, + { + "name": "Forbidden", + "description": "The caller does not have permission to pin messages for this streamer." + }, + { + "name": "SessionNotFound", + "description": "The streamer's OAuth session could not be found or is invalid." + } + ] + } + } +} diff --git a/lexicons/place/stream/moderation/deletePin.json b/lexicons/place/stream/moderation/deletePin.json new file mode 100644 index 00000000..50537fd0 --- /dev/null +++ b/lexicons/place/stream/moderation/deletePin.json @@ -0,0 +1,50 @@ +{ + "lexicon": 1, + "id": "place.stream.moderation.deletePin", + "defs": { + "main": { + "type": "procedure", + "description": "Unpin a pinned chat message on behalf of a streamer. Requires 'message.pin' permission. Deletes the place.stream.chat.pinnedRecord from the streamer's repo.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["streamer", "pinUri"], + "properties": { + "streamer": { + "type": "string", + "format": "did", + "description": "The DID of the streamer." + }, + "pinUri": { + "type": "string", + "format": "at-uri", + "description": "The AT-URI of the pinned record to delete." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "properties": {} + } + }, + "errors": [ + { + "name": "Unauthorized", + "description": "The request lacks valid authentication credentials." + }, + { + "name": "Forbidden", + "description": "The caller does not have permission to unpin messages for this streamer." + }, + { + "name": "SessionNotFound", + "description": "The streamer's OAuth session could not be found or is invalid." + } + ] + } + } +} diff --git a/lexicons/place/stream/moderation/permission.json b/lexicons/place/stream/moderation/permission.json index ee579693..68a3d3ce 100644 --- a/lexicons/place/stream/moderation/permission.json +++ b/lexicons/place/stream/moderation/permission.json @@ -19,7 +19,7 @@ "type": "array", "items": { "type": "string", - "enum": ["ban", "hide", "livestream.manage"] + "enum": ["ban", "hide", "livestream.manage", "message.pin"] }, "description": "Array of permissions granted to this moderator. 'ban' covers blocks/bans (with optional expiration), 'hide' covers message gates, 'livestream.manage' allows updating livestream metadata." }, diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index 8e7ca0be..66b71a27 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -340,6 +340,20 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt } } + if collection.String() == constants.PLACE_STREAM_CHAT_PINNED_RECORD { + log.Debug(ctx, "deleting pinned record", "userDID", evt.Repo, "rkey", rkey.String()) + err := atsync.Model.DeletePinnedRecord(ctx, rkey.String()) + if err != nil { + log.Error(ctx, "failed to delete pinned record", "err", err) + } + deletedPin := map[string]any{ + "$type": constants.PLACE_STREAM_CHAT_PINNED_RECORD, + "rkey": rkey.String(), + "deleted": true, + } + go atsync.Bus.Publish(evt.Repo, deletedPin) + } + default: log.Error(ctx, "unexpected record op kind") } diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 13610a55..a739a695 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -217,6 +217,51 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD } go atsync.Bus.Publish(userDID, streamplaceGate) + case *streamplace.ChatPinnedRecord: + repo, err := atsync.SyncBlueskyRepoCached(ctx, userDID) + if err != nil { + return fmt.Errorf("failed to sync bluesky repo: %w", err) + } + if r == nil { + return nil + } + log.Debug(ctx, "creating pinned record", "userDID", userDID, "pinnedMessage", rec.PinnedMessage) + // Delete existing pinned records for this streamer (single-pin semantics) + err = atsync.Model.DeleteAllPinnedRecords(ctx, userDID) + if err != nil { + log.Error(ctx, "failed to delete existing pinned records", "err", err) + } + // Parse optional expiresAt + var expiresAt *time.Time + if rec.ExpiresAt != nil { + t, err := time.Parse(time.RFC3339, *rec.ExpiresAt) + if err == nil { + expiresAt = &t + } + } + pin := &model.PinnedRecord{ + RKey: rkey.String(), + RepoDID: userDID, + PinnedMessage: rec.PinnedMessage, + CID: cid, + CreatedAt: now, + Repo: repo, + ExpiresAt: expiresAt, + } + err = atsync.Model.CreatePinnedRecord(ctx, pin) + if err != nil { + return fmt.Errorf("failed to create pinned record: %w", err) + } + pin, err = atsync.Model.GetPinnedRecord(ctx, rkey.String()) + if err != nil { + return fmt.Errorf("failed to get pinned record after we just saved it: %w", err) + } + pinnedView, err := pin.ToStreamplacePinnedRecord() + if err != nil { + return fmt.Errorf("failed to convert pinned record: %w", err) + } + go atsync.Bus.Publish(userDID, pinnedView) + case *streamplace.ChatProfile: repo, err := atsync.SyncBlueskyRepoCached(ctx, userDID) if err != nil { diff --git a/pkg/constants/constants.go b/pkg/constants/constants.go index 56cb1511..4bebb298 100644 --- a/pkg/constants/constants.go +++ b/pkg/constants/constants.go @@ -13,6 +13,7 @@ var APP_BSKY_FEED_POST = "app.bsky.feed.post" // var APP_BSKY_GRAPH_BLOCK = "app.bsky.graph.block" //nolint:all var APP_BSKY_ACTOR_PROFILE = "app.bsky.actor.profile" //nolint:all var PLACE_STREAM_CHAT_GATE = "place.stream.chat.gate" //nolint:all +var PLACE_STREAM_CHAT_PINNED_RECORD = "place.stream.chat.pinnedRecord" //nolint:all var PLACE_STREAM_DEFAULT_METADATA = "place.stream.metadata.configuration" //nolint:all var PLACE_STREAM_LIVE_RECOMMENDATIONS = "place.stream.live.recommendations" //nolint:all diff --git a/pkg/gen/gen.go b/pkg/gen/gen.go index 6d8fbbf8..30efab4e 100644 --- a/pkg/gen/gen.go +++ b/pkg/gen/gen.go @@ -26,6 +26,7 @@ func main() { streamplace.ChatMessage_ReplyRef{}, streamplace.ServerSettings{}, streamplace.ChatGate{}, + streamplace.ChatPinnedRecord{}, streamplace.MultistreamTarget{}, streamplace.BroadcastOrigin{}, streamplace.BroadcastSyndication{}, diff --git a/pkg/model/model.go b/pkg/model/model.go index 5fd9e944..5571182f 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -85,6 +85,12 @@ type Model interface { GetGate(ctx context.Context, rkey string) (*Gate, error) GetUserGates(ctx context.Context, userDID string) ([]*Gate, error) + CreatePinnedRecord(ctx context.Context, pin *PinnedRecord) error + DeletePinnedRecord(ctx context.Context, rkey string) error + DeleteAllPinnedRecords(ctx context.Context, streamerDID string) error + GetPinnedRecord(ctx context.Context, rkey string) (*PinnedRecord, error) + GetActivePinnedRecord(ctx context.Context, streamerDID string) (*PinnedRecord, error) + CreateChatProfile(ctx context.Context, profile *ChatProfile) error GetChatProfile(ctx context.Context, repoDID string) (*ChatProfile, error) @@ -179,6 +185,7 @@ func MakeDB(dbURL string) (Model, error) { ChatMessage{}, ChatProfile{}, Gate{}, + PinnedRecord{}, ServerSettings{}, Labeler{}, Label{}, diff --git a/pkg/model/pinned_record.go b/pkg/model/pinned_record.go new file mode 100644 index 00000000..90add2f7 --- /dev/null +++ b/pkg/model/pinned_record.go @@ -0,0 +1,73 @@ +package model + +import ( + "context" + "errors" + "time" + + "gorm.io/gorm" + "stream.place/streamplace/pkg/streamplace" +) + +type PinnedRecord struct { + RKey string `gorm:"primaryKey;column:rkey"` + CID string `gorm:"column:cid"` + RepoDID string `json:"repoDID" gorm:"column:repo_did"` + Repo *Repo `json:"repo,omitempty" gorm:"foreignKey:DID;references:RepoDID"` + PinnedMessage string `gorm:"column:pinned_message" json:"pinnedMessage"` + ExpiresAt *time.Time `gorm:"column:expires_at" json:"expiresAt"` + CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` +} + +func (p *PinnedRecord) ToStreamplacePinnedRecord() (*streamplace.ChatPinnedRecord, error) { + rec := &streamplace.ChatPinnedRecord{ + LexiconTypeID: "place.stream.chat.pinnedRecord", + PinnedMessage: p.PinnedMessage, + CreatedAt: p.CreatedAt.UTC().Format(time.RFC3339), + } + if p.ExpiresAt != nil { + s := p.ExpiresAt.UTC().Format(time.RFC3339) + rec.ExpiresAt = &s + } + return rec, nil +} + +func (m *DBModel) CreatePinnedRecord(ctx context.Context, pin *PinnedRecord) error { + return m.DB.Create(pin).Error +} + +func (m *DBModel) GetPinnedRecord(ctx context.Context, rkey string) (*PinnedRecord, error) { + var pin PinnedRecord + err := m.DB.Preload("Repo").Where("rkey = ?", rkey).First(&pin).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + return &pin, nil +} + +func (m *DBModel) DeletePinnedRecord(ctx context.Context, rkey string) error { + return m.DB.Where("rkey = ?", rkey).Delete(&PinnedRecord{}).Error +} + +func (m *DBModel) DeleteAllPinnedRecords(ctx context.Context, streamerDID string) error { + return m.DB.Where("repo_did = ?", streamerDID).Delete(&PinnedRecord{}).Error +} + +func (m *DBModel) GetActivePinnedRecord(ctx context.Context, streamerDID string) (*PinnedRecord, error) { + var pin PinnedRecord + now := time.Now() + err := m.DB.Preload("Repo"). + Where("repo_did = ? AND (expires_at IS NULL OR expires_at > ?)", streamerDID, now). + Order("created_at DESC"). + First(&pin).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + return &pin, nil +} diff --git a/pkg/moderation/permissions.go b/pkg/moderation/permissions.go index 81b0a9d1..92f2b8e7 100644 --- a/pkg/moderation/permissions.go +++ b/pkg/moderation/permissions.go @@ -17,6 +17,7 @@ const ( PermissionBan = "ban" PermissionHide = "hide" PermissionLivestreamManage = "livestream.manage" + PermissionMessagePin = "message.pin" ) // ActionPermissions maps moderation actions to required permissions @@ -26,6 +27,8 @@ var ActionPermissions = map[string]string{ "createGate": PermissionHide, "deleteGate": PermissionHide, "updateLivestream": PermissionLivestreamManage, + "createPin": PermissionMessagePin, + "deletePin": PermissionMessagePin, } // PermissionChecker validates moderation permissions diff --git a/pkg/spxrpc/place_stream_moderation.go b/pkg/spxrpc/place_stream_moderation.go index c6bb21a8..a70ff81d 100644 --- a/pkg/spxrpc/place_stream_moderation.go +++ b/pkg/spxrpc/place_stream_moderation.go @@ -360,3 +360,127 @@ func (s *Server) logAudit(ctx context.Context, streamerDID, moderatorDID, action return s.statefulDB.CreateAuditLog(ctx, auditLog) } + +func (s *Server) handlePlaceStreamModerationCreatePin(ctx context.Context, input *streamplace.ModerationCreatePin_Input) (*streamplace.ModerationCreatePin_Output, error) { + // Validate input + if err := validateDID(input.Streamer); err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, fmt.Sprintf("invalid streamer DID: %v", err)) + } + if err := validateATURI(input.MessageUri); err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, fmt.Sprintf("invalid messageUri: %v", err)) + } + + // Get delegated moderation context (validates OAuth, permission, and returns client) + modCtx, err := s.GetDelegatedModerationContext(ctx, input.Streamer, "createPin") + if err != nil { + 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 + pinnedRecord := &streamplace.ChatPinnedRecord{ + LexiconTypeID: "place.stream.chat.pinnedRecord", + PinnedMessage: input.MessageUri, + CreatedAt: time.Now().UTC().Format(time.RFC3339), + ExpiresAt: input.ExpiresAt, + } + + createInput := comatproto.RepoCreateRecord_Input{ + Collection: constants.PLACE_STREAM_CHAT_PINNED_RECORD, + Record: &lexutil.LexiconTypeDecoder{Val: pinnedRecord}, + Repo: input.Streamer, + } + createOutput := comatproto.RepoCreateRecord_Output{} + + err = modCtx.StreamerClient.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.createRecord", map[string]any{}, createInput, &createOutput) + if err != nil { + log.Error(ctx, "failed to create pinned record", "err", err) + if auditErr := s.logAudit(ctx, input.Streamer, modCtx.ModeratorDID, "createPin", input.MessageUri, "", "", false, err.Error()); auditErr != nil { + log.Error(ctx, "failed to create audit log", "error", auditErr) + } + return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("failed to create pinned record: %v", err)) + } + + // Log successful audit entry + if err := s.logAudit(ctx, input.Streamer, modCtx.ModeratorDID, "createPin", input.MessageUri, "", createOutput.Uri, true, ""); err != nil { + log.Error(ctx, "failed to create audit log", "error", err) + } + + return &streamplace.ModerationCreatePin_Output{ + Uri: createOutput.Uri, + Cid: createOutput.Cid, + }, nil +} + +func (s *Server) handlePlaceStreamModerationDeletePin(ctx context.Context, input *streamplace.ModerationDeletePin_Input) (*streamplace.ModerationDeletePin_Output, error) { + // Validate input + if err := validateDID(input.Streamer); err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, fmt.Sprintf("invalid streamer DID: %v", err)) + } + if err := validateATURI(input.PinUri); err != nil { + return nil, echo.NewHTTPError(http.StatusBadRequest, fmt.Sprintf("invalid pinUri: %v", err)) + } + + // Get delegated moderation context (validates OAuth, permission, and returns client) + modCtx, err := s.GetDelegatedModerationContext(ctx, input.Streamer, "deletePin") + if err != nil { + return nil, err + } + + // Parse pinUri to extract rkey + rkey, err := extractRKey(input.PinUri) + if err != nil { + log.Error(ctx, "failed to extract rkey from pinUri", "uri", input.PinUri, "err", err) + return nil, echo.NewHTTPError(http.StatusBadRequest, "invalid pinUri format") + } + + // Delete pinned record from streamer's repo + deleteInput := comatproto.RepoDeleteRecord_Input{ + Collection: constants.PLACE_STREAM_CHAT_PINNED_RECORD, + Rkey: rkey, + Repo: input.Streamer, + } + deleteOutput := comatproto.RepoDeleteRecord_Output{} + + err = modCtx.StreamerClient.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.deleteRecord", map[string]any{}, deleteInput, &deleteOutput) + if err != nil { + log.Error(ctx, "failed to delete pinned record", "err", err) + if auditErr := s.logAudit(ctx, input.Streamer, modCtx.ModeratorDID, "deletePin", input.PinUri, "", "", false, err.Error()); auditErr != nil { + log.Error(ctx, "failed to create audit log", "error", auditErr) + } + return nil, echo.NewHTTPError(http.StatusInternalServerError, fmt.Sprintf("failed to delete pinned record: %v", err)) + } + + // Log successful audit entry + if err := s.logAudit(ctx, input.Streamer, modCtx.ModeratorDID, "deletePin", input.PinUri, "", "", true, ""); err != nil { + log.Error(ctx, "failed to create audit log", "error", err) + } + + return &streamplace.ModerationDeletePin_Output{}, nil +} diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index 7367ebc7..346c0a55 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -296,8 +296,10 @@ func (s *Server) RegisterHandlersPlaceStream(e *echo.Echo) error { e.POST("/xrpc/place.stream.live.stopLivestream", s.HandlePlaceStreamLiveStopLivestream) e.POST("/xrpc/place.stream.moderation.createBlock", s.HandlePlaceStreamModerationCreateBlock) e.POST("/xrpc/place.stream.moderation.createGate", s.HandlePlaceStreamModerationCreateGate) + e.POST("/xrpc/place.stream.moderation.createPin", s.HandlePlaceStreamModerationCreatePin) e.POST("/xrpc/place.stream.moderation.deleteBlock", s.HandlePlaceStreamModerationDeleteBlock) e.POST("/xrpc/place.stream.moderation.deleteGate", s.HandlePlaceStreamModerationDeleteGate) + e.POST("/xrpc/place.stream.moderation.deletePin", s.HandlePlaceStreamModerationDeletePin) e.POST("/xrpc/place.stream.moderation.updateLivestream", s.HandlePlaceStreamModerationUpdateLivestream) e.POST("/xrpc/place.stream.multistream.createTarget", s.HandlePlaceStreamMultistreamCreateTarget) e.POST("/xrpc/place.stream.multistream.deleteTarget", s.HandlePlaceStreamMultistreamDeleteTarget) @@ -627,6 +629,24 @@ func (s *Server) HandlePlaceStreamModerationCreateGate(c echo.Context) error { return c.JSON(200, out) } +func (s *Server) HandlePlaceStreamModerationCreatePin(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamModerationCreatePin") + defer span.End() + + var body placestream.ModerationCreatePin_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestream.ModerationCreatePin_Output + var handleErr error + // func (s *Server) handlePlaceStreamModerationCreatePin(ctx context.Context,body *placestream.ModerationCreatePin_Input) (*placestream.ModerationCreatePin_Output, error) + out, handleErr = s.handlePlaceStreamModerationCreatePin(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) HandlePlaceStreamModerationDeleteBlock(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamModerationDeleteBlock") defer span.End() @@ -663,6 +683,24 @@ func (s *Server) HandlePlaceStreamModerationDeleteGate(c echo.Context) error { return c.JSON(200, out) } +func (s *Server) HandlePlaceStreamModerationDeletePin(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamModerationDeletePin") + defer span.End() + + var body placestream.ModerationDeletePin_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *placestream.ModerationDeletePin_Output + var handleErr error + // func (s *Server) handlePlaceStreamModerationDeletePin(ctx context.Context,body *placestream.ModerationDeletePin_Input) (*placestream.ModerationDeletePin_Output, error) + out, handleErr = s.handlePlaceStreamModerationDeletePin(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) HandlePlaceStreamModerationUpdateLivestream(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamModerationUpdateLivestream") defer span.End() diff --git a/pkg/streamplace/cbor_gen.go b/pkg/streamplace/cbor_gen.go index 0e88d5fb..197bd3b9 100644 --- a/pkg/streamplace/cbor_gen.go +++ b/pkg/streamplace/cbor_gen.go @@ -3634,6 +3634,228 @@ func (t *ChatGate) UnmarshalCBOR(r io.Reader) (err error) { return nil } +func (t *ChatPinnedRecord) MarshalCBOR(w io.Writer) error { + if t == nil { + _, err := w.Write(cbg.CborNull) + return err + } + + cw := cbg.NewCborWriter(w) + fieldCount := 4 + + if t.ExpiresAt == nil { + fieldCount-- + } + + if _, err := cw.Write(cbg.CborEncodeMajorType(cbg.MajMap, uint64(fieldCount))); err != nil { + return err + } + + // t.LexiconTypeID (string) (string) + if len("$type") > 1000000 { + return xerrors.Errorf("Value in field \"$type\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("$type"))); err != nil { + return err + } + if _, err := cw.WriteString(string("$type")); err != nil { + return err + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("place.stream.chat.pinnedRecord"))); err != nil { + return err + } + if _, err := cw.WriteString(string("place.stream.chat.pinnedRecord")); err != nil { + return err + } + + // t.CreatedAt (string) (string) + if len("createdAt") > 1000000 { + return xerrors.Errorf("Value in field \"createdAt\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("createdAt"))); err != nil { + return err + } + if _, err := cw.WriteString(string("createdAt")); err != nil { + return err + } + + if len(t.CreatedAt) > 1000000 { + return xerrors.Errorf("Value in field t.CreatedAt was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.CreatedAt))); err != nil { + return err + } + if _, err := cw.WriteString(string(t.CreatedAt)); err != nil { + return err + } + + // t.ExpiresAt (string) (string) + if t.ExpiresAt != nil { + + if len("expiresAt") > 1000000 { + return xerrors.Errorf("Value in field \"expiresAt\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("expiresAt"))); err != nil { + return err + } + if _, err := cw.WriteString(string("expiresAt")); err != nil { + return err + } + + if t.ExpiresAt == nil { + if _, err := cw.Write(cbg.CborNull); err != nil { + return err + } + } else { + if len(*t.ExpiresAt) > 1000000 { + return xerrors.Errorf("Value in field t.ExpiresAt was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(*t.ExpiresAt))); err != nil { + return err + } + if _, err := cw.WriteString(string(*t.ExpiresAt)); err != nil { + return err + } + } + } + + // t.PinnedMessage (string) (string) + if len("pinnedMessage") > 1000000 { + return xerrors.Errorf("Value in field \"pinnedMessage\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("pinnedMessage"))); err != nil { + return err + } + if _, err := cw.WriteString(string("pinnedMessage")); err != nil { + return err + } + + if len(t.PinnedMessage) > 1000000 { + return xerrors.Errorf("Value in field t.PinnedMessage was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.PinnedMessage))); err != nil { + return err + } + if _, err := cw.WriteString(string(t.PinnedMessage)); err != nil { + return err + } + return nil +} + +func (t *ChatPinnedRecord) UnmarshalCBOR(r io.Reader) (err error) { + *t = ChatPinnedRecord{} + + cr := cbg.NewCborReader(r) + + maj, extra, err := cr.ReadHeader() + if err != nil { + return err + } + defer func() { + if err == io.EOF { + err = io.ErrUnexpectedEOF + } + }() + + if maj != cbg.MajMap { + return fmt.Errorf("cbor input should be of type map") + } + + if extra > cbg.MaxLength { + return fmt.Errorf("ChatPinnedRecord: map struct too large (%d)", extra) + } + + n := extra + + nameBuf := make([]byte, 13) + for i := uint64(0); i < n; i++ { + nameLen, ok, err := cbg.ReadFullStringIntoBuf(cr, nameBuf, 1000000) + if err != nil { + return err + } + + if !ok { + // Field doesn't exist on this type, so ignore it + if err := cbg.ScanForLinks(cr, func(cid.Cid) {}); err != nil { + return err + } + continue + } + + switch string(nameBuf[:nameLen]) { + // t.LexiconTypeID (string) (string) + case "$type": + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.LexiconTypeID = string(sval) + } + // t.CreatedAt (string) (string) + case "createdAt": + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.CreatedAt = string(sval) + } + // t.ExpiresAt (string) (string) + case "expiresAt": + + { + b, err := cr.ReadByte() + if err != nil { + return err + } + if b != cbg.CborNull[0] { + if err := cr.UnreadByte(); err != nil { + return err + } + + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.ExpiresAt = (*string)(&sval) + } + } + // t.PinnedMessage (string) (string) + case "pinnedMessage": + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.PinnedMessage = string(sval) + } + + default: + // Field doesn't exist on this type, so ignore it + if err := cbg.ScanForLinks(r, func(cid.Cid) {}); err != nil { + return err + } + } + } + + return nil +} func (t *MultistreamTarget) MarshalCBOR(w io.Writer) error { if t == nil { _, err := w.Write(cbg.CborNull) diff --git a/pkg/streamplace/chatdefs.go b/pkg/streamplace/chatdefs.go index 6670cff7..0f2aedac 100644 --- a/pkg/streamplace/chatdefs.go +++ b/pkg/streamplace/chatdefs.go @@ -54,3 +54,15 @@ func (t *ChatDefs_MessageView_ReplyTo) UnmarshalJSON(b []byte) error { return nil } } + +// ChatDefs_PinnedRecordView is a "pinnedRecordView" in the place.stream.chat.defs schema. +// +// View of a pinned chat record with hydrated message data. +type ChatDefs_PinnedRecordView struct { + Cid string `json:"cid" cborgen:"cid"` + IndexedAt string `json:"indexedAt" cborgen:"indexedAt"` + Message *ChatDefs_MessageView `json:"message,omitempty" cborgen:"message,omitempty"` + PinnedBy *appbsky.ActorDefs_ProfileViewBasic `json:"pinnedBy,omitempty" cborgen:"pinnedBy,omitempty"` + Record *ChatPinnedRecord `json:"record" cborgen:"record"` + Uri string `json:"uri" cborgen:"uri"` +} diff --git a/pkg/streamplace/chatpinnedRecord.go b/pkg/streamplace/chatpinnedRecord.go new file mode 100644 index 00000000..b61cddfd --- /dev/null +++ b/pkg/streamplace/chatpinnedRecord.go @@ -0,0 +1,23 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.chat.pinnedRecord + +package streamplace + +import ( + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +func init() { + lexutil.RegisterType("place.stream.chat.pinnedRecord", &ChatPinnedRecord{}) +} + +type ChatPinnedRecord struct { + LexiconTypeID string `json:"$type" cborgen:"$type,const=place.stream.chat.pinnedRecord"` + // createdAt: When this pin was created. + CreatedAt string `json:"createdAt" cborgen:"createdAt"` + // expiresAt: Optional expiration time. If set, the pin is considered inactive after this time. + ExpiresAt *string `json:"expiresAt,omitempty" cborgen:"expiresAt,omitempty"` + // pinnedMessage: AT-URI of the pinned chat message. + PinnedMessage string `json:"pinnedMessage" cborgen:"pinnedMessage"` +} diff --git a/pkg/streamplace/moderationcreatePin.go b/pkg/streamplace/moderationcreatePin.go new file mode 100644 index 00000000..15d742d9 --- /dev/null +++ b/pkg/streamplace/moderationcreatePin.go @@ -0,0 +1,39 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.moderation.createPin + +package streamplace + +import ( + "context" + + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +// ModerationCreatePin_Input is the input argument to a place.stream.moderation.createPin call. +type ModerationCreatePin_Input struct { + // expiresAt: Optional expiration time for this pin. + ExpiresAt *string `json:"expiresAt,omitempty" cborgen:"expiresAt,omitempty"` + // messageUri: The AT-URI of the chat message to pin. + MessageUri string `json:"messageUri" cborgen:"messageUri"` + // streamer: The DID of the streamer. + Streamer string `json:"streamer" cborgen:"streamer"` +} + +// ModerationCreatePin_Output is the output of a place.stream.moderation.createPin call. +type ModerationCreatePin_Output struct { + // cid: The CID of the created pinned record. + Cid string `json:"cid" cborgen:"cid"` + // uri: The AT-URI of the created pinned record. + Uri string `json:"uri" cborgen:"uri"` +} + +// ModerationCreatePin calls the XRPC method "place.stream.moderation.createPin". +func ModerationCreatePin(ctx context.Context, c lexutil.LexClient, input *ModerationCreatePin_Input) (*ModerationCreatePin_Output, error) { + var out ModerationCreatePin_Output + if err := c.LexDo(ctx, lexutil.Procedure, "application/json", "place.stream.moderation.createPin", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/pkg/streamplace/moderationdeletePin.go b/pkg/streamplace/moderationdeletePin.go new file mode 100644 index 00000000..28654d62 --- /dev/null +++ b/pkg/streamplace/moderationdeletePin.go @@ -0,0 +1,33 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.moderation.deletePin + +package streamplace + +import ( + "context" + + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +// ModerationDeletePin_Input is the input argument to a place.stream.moderation.deletePin call. +type ModerationDeletePin_Input struct { + // pinUri: The AT-URI of the pinned record to delete. + PinUri string `json:"pinUri" cborgen:"pinUri"` + // streamer: The DID of the streamer. + Streamer string `json:"streamer" cborgen:"streamer"` +} + +// ModerationDeletePin_Output is the output of a place.stream.moderation.deletePin call. +type ModerationDeletePin_Output struct { +} + +// ModerationDeletePin calls the XRPC method "place.stream.moderation.deletePin". +func ModerationDeletePin(ctx context.Context, c lexutil.LexClient, input *ModerationDeletePin_Input) (*ModerationDeletePin_Output, error) { + var out ModerationDeletePin_Output + if err := c.LexDo(ctx, lexutil.Procedure, "application/json", "place.stream.moderation.deletePin", nil, input, &out); err != nil { + return nil, err + } + + return &out, nil +} -- 2.51.2 From a4adbcb8a32362d86fa72e4340a86dd817e877cf Mon Sep 17 00:00:00 2001 From: "Natalie B." <22222885+espeon@users.noreply.github.com> Date: Fri, 20 Mar 2026 18:35:00 -0500 Subject: [PATCH 2/7] add dropdown trigger for more options --- .../src/components/chat/mod-view.tsx | 69 ++++++++++++++----- pkg/atproto/labeler_firehose.go | 4 ++ pkg/spxrpc/place_stream_moderation.go | 19 +++++ 3 files changed, 73 insertions(+), 19 deletions(-) diff --git a/js/components/src/components/chat/mod-view.tsx b/js/components/src/components/chat/mod-view.tsx index 4d0431be..afdc1976 100644 --- a/js/components/src/components/chat/mod-view.tsx +++ b/js/components/src/components/chat/mod-view.tsx @@ -26,6 +26,9 @@ import { DropdownMenu, DropdownMenuGroup, DropdownMenuItem, + DropdownMenuSub, + DropdownMenuSubContent, + DropdownMenuSubTrigger, DropdownMenuTrigger, layout, ResponsiveDropdownMenuContent, @@ -234,25 +237,53 @@ function ModViewContent({ )} {modPermissions.canPin && message.author.did !== streamerDID && ( - { - if (!streamerDID) return; - pinChatMessage(message.uri, streamerDID) - .then(() => { - toast.show("Comment pinned", "", { duration: 3 }); - onOpenChange?.(false); - }) - .catch((e) => { - toast.show( - "Error pinning comment", - e instanceof Error ? e.message : "Failed to pin", - { duration: 5 }, - ); - }); - }} - > - Pin this message - + + + Pin this message + + + { + if (!streamerDID) return; + pinChatMessage(message.uri, streamerDID) + .then(() => { + toast.show("Comment pinned", "", { duration: 3 }); + onOpenChange?.(false); + }) + .catch((e) => { + toast.show( + "Error pinning comment", + e instanceof Error ? e.message : "Failed to pin", + { duration: 5 }, + ); + }); + }} + > + Pin indefinitely + + { + if (!streamerDID) return; + const expiresAt = new Date(); + expiresAt.setHours(expiresAt.getHours() + 1); // Set expiration to 1 hour from now + pinChatMessage(message.uri, streamerDID, expiresAt) + .then(() => { + toast.show("Comment pinned", "", { duration: 3 }); + onOpenChange?.(false); + }) + .catch((e) => { + toast.show( + "Error pinning comment", + e instanceof Error ? e.message : "Failed to pin", + { duration: 5 }, + ); + }); + }} + > + Pin for 1 hour + + + )} {modPermissions.canBan && agent?.did && diff --git a/pkg/atproto/labeler_firehose.go b/pkg/atproto/labeler_firehose.go index 9ca51fdb..470f7405 100644 --- a/pkg/atproto/labeler_firehose.go +++ b/pkg/atproto/labeler_firehose.go @@ -182,6 +182,10 @@ func (atsync *ATProtoSynchronizer) StartLabelerFirehoseRetry(ctx context.Context log.Error(ctx, "failed to get chat message for label", "err", err) continue } + if msg == nil { + log.Debug(ctx, "chat message not found for label, skipping", "uri", l.URI) + continue + } chatView, err := msg.ToStreamplaceMessageView() if err != nil { log.Error(ctx, "failed to convert chat message to streamplace message view", "err", err) diff --git a/pkg/spxrpc/place_stream_moderation.go b/pkg/spxrpc/place_stream_moderation.go index a70ff81d..f621f8e3 100644 --- a/pkg/spxrpc/place_stream_moderation.go +++ b/pkg/spxrpc/place_stream_moderation.go @@ -376,6 +376,25 @@ 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 a7dbdb6be9bac1326ed897023a307dc08ee48ddd Mon Sep 17 00:00:00 2001 From: "Natalie B." <22222885+espeon@users.noreply.github.com> Date: Fri, 20 Mar 2026 23:58:48 -0500 Subject: [PATCH 3/7] Add pinned chat message feature Add UI pin controls (preset durations), a pinned comment notification banner, websocket/server support to broadcast the active pin, and DB/ lexicon/schema changes to record the pinner and TTL --- docs/pinned-comment-plan.md | 460 ------------------ .../src/components/chat/mod-view.tsx | 110 +++-- .../stream-notification/pin-notification.tsx | 135 +++++ js/components/src/lib/stream-notifications.ts | 28 ++ .../src/livestream-provider/index.tsx | 38 +- .../chat/place-stream-chat-defs.md | 18 +- .../chat/place-stream-chat-pinnedrecord.md | 6 + .../lex-reference/place-stream-livestream.md | 9 +- lexicons/place/stream/chat/defs.json | 2 +- lexicons/place/stream/chat/pinnedRecord.json | 5 + lexicons/place/stream/livestream.json | 3 +- pkg/api/websocket.go | 53 ++ pkg/atproto/sync.go | 51 +- pkg/model/pinned_record.go | 22 + pkg/streamplace/cbor_gen.go | 59 ++- pkg/streamplace/chatdefs.go | 13 +- pkg/streamplace/chatpinnedRecord.go | 2 + pkg/streamplace/streamlivestream.go | 8 + 18 files changed, 484 insertions(+), 538 deletions(-) delete mode 100644 docs/pinned-comment-plan.md create mode 100644 js/components/src/components/stream-notification/pin-notification.tsx diff --git a/docs/pinned-comment-plan.md b/docs/pinned-comment-plan.md deleted file mode 100644 index a6babef4..00000000 --- a/docs/pinned-comment-plan.md +++ /dev/null @@ -1,460 +0,0 @@ -# Pinned Comment Feature Plan - -## Overview - -Allow the streamer or a delegated moderator to pin an existing chat message. The pinned comment appears prominently in the chat UI and auto-expires when a TTL is reached or the stream ends. - -## Design Decisions - -- **Pinnable content**: Existing chat messages only (referenced by AT-URI) -- **Who can pin**: Streamer + mods with new `message.pin` permission -- **Expiration**: Optional TTL (`expiresAt` datetime) + auto-clear at stream end (client-side) -- **Storage pattern**: Record in streamer's AT Protocol repo, modeled after `place.stream.chat.gate` -- **Single active pin**: Only one pinned comment at a time. Creating a new pin replaces the previous one. -- **Viewer dismiss**: Any viewer can hide the pin from their own view (local state only). The pin remains for everyone else. - ---- - -## 1. New Lexicon: `place.stream.chat.pinnedRecord` - -**File**: `lexicons/place/stream/chat/pinnedRecord.json` - -```json -{ - "lexicon": 1, - "id": "place.stream.chat.pinnedRecord", - "defs": { - "main": { - "type": "record", - "key": "tid", - "description": "Record pinning a chat message for prominent display.", - "record": { - "type": "object", - "required": ["pinnedMessage", "createdAt"], - "properties": { - "pinnedMessage": { - "type": "string", - "format": "at-uri", - "description": "AT-URI of the pinned chat message." - }, - "createdAt": { - "type": "string", - "format": "datetime", - "description": "When this pin was created." - }, - "expiresAt": { - "type": "string", - "format": "datetime", - "description": "Optional expiration time. If set, the pin is considered inactive after this time." - } - } - } - } - } -} -``` - ---- - -## 2. Pinned Record View + Defs - -### 2a. Add `pinnedRecordView` to chat defs - -**File**: `lexicons/place/stream/chat/defs.json` (add new definition) - -```json -"pinnedRecordView": { - "type": "object", - "description": "View of a pinned chat record with hydrated message data.", - "required": ["uri", "cid", "record", "indexedAt"], - "properties": { - "uri": { "type": "string", "format": "at-uri" }, - "cid": { "type": "string", "format": "cid" }, - "record": { "type": "ref", "ref": "place.stream.chat.pinnedRecord" }, - "indexedAt": { "type": "string", "format": "datetime" }, - "pinnedBy": { "type": "ref", "ref": "app.bsky.actor.defs#profileViewBasic" }, - "message": { "type": "ref", "ref": "place.stream.chat.defs#messageView" } - } -} -``` - -This gives us a hydrated view that includes the full message data, so the frontend doesn't need to look it up from the chat index. - -### 2b. Go struct for bus delivery - -The bus will deliver a `PinnedRecordView` (not just the raw record) so the frontend can render the message text directly. Similar to how `ChatGate` records are published raw but the frontend processes them. - ---- - -## 3. New Permission: `message.pin` - -### 3a. Add to lexicon enum - -**File**: `lexicons/place/stream/moderation/permission.json` - -Add `"message.pin"` to the permissions enum: - -```json -"enum": ["ban", "hide", "livestream.manage", "message.pin"] -``` - -### 3b. Register in Go permission system - -**File**: `pkg/moderation/permissions.go` - -```go -const PermissionMessagePin = "message.pin" - -var ActionPermissions = map[string]string{ - // ... existing entries ... - "createPin": PermissionMessagePin, - "deletePin": PermissionMessagePin, -} -``` - -### 3c. Update frontend permissions - -**File**: `js/components/src/streamplace-store/moderation.tsx` - -Add `canPin: boolean` to `ModerationPermissions` interface. Derive from permissions array: - -```ts -canPin: isOwner || permissions.includes("message.pin"), -``` - ---- - -## 4. New RPC Procedures - -### 4a. Lexicon: `place.stream.moderation.createPin` - -**File**: `lexicons/place/stream/moderation/createPin.json` - -```json -{ - "lexicon": 1, - "id": "place.stream.moderation.createPin", - "defs": { - "main": { - "type": "procedure", - "description": "Pin a chat message on behalf of a streamer. Requires 'message.pin' permission. Creates a place.stream.chat.pinnedRecord in the streamer's repo, replacing any existing pin.", - "input": { - "encoding": "application/json", - "schema": { - "type": "object", - "required": ["streamer", "messageUri"], - "properties": { - "streamer": { "type": "string", "format": "did" }, - "messageUri": { "type": "string", "format": "at-uri" }, - "expiresAt": { "type": "string", "format": "datetime" } - } - } - }, - "output": { - "encoding": "application/json", - "schema": { - "type": "object", - "required": ["uri", "cid"], - "properties": { - "uri": { "type": "string", "format": "at-uri" }, - "cid": { "type": "string", "format": "cid" } - } - } - }, - "errors": [ - { "name": "Unauthorized" }, - { "name": "Forbidden" }, - { "name": "SessionNotFound" } - ] - } - } -} -``` - -### 4b. Lexicon: `place.stream.moderation.deletePin` - -**File**: `lexicons/place/stream/moderation/deletePin.json` - -Same pattern as `deleteGate`: - -- Input: `streamer` (DID), `pinUri` (at-uri) -- Output: empty -- Errors: Unauthorized, Forbidden, SessionNotFound - -### 4c. Go handlers - -**File**: `pkg/spxrpc/place_stream_moderation.go` - -Add `handlePlaceStreamModerationCreatePin`: - -1. Validate input (DID, AT-URI) -2. `GetDelegatedModerationContext(ctx, input.Streamer, "createPin")` -3. Before creating: list existing `place.stream.chat.pinnedRecord` records in streamer's repo, delete any existing ones (single-pin semantics) -4. Build `streamplace.ChatPinnedRecord` struct -5. Create via `com.atproto.repo.createRecord` on streamer's repo -6. Audit log -7. Return URI + CID - -Add `handlePlaceStreamModerationDeletePin`: - -1. Validate input -2. `GetDelegatedModerationContext(ctx, input.Streamer, "deletePin")` -3. Extract rkey, delete record via `com.atproto.repo.deleteRecord` -4. Audit log - -### 4d. Register handlers - -In the XRPC server setup where `createGate`/`deleteGate` are registered, add `createPin` and `deletePin`. - ---- - -## 5. Backend: Model + DB - -**File**: `pkg/model/pinned_record.go` - -```go -type PinnedRecord struct { - RKey string `gorm:"primaryKey;column:rkey"` - CID string `gorm:"column:cid"` - RepoDID string `gorm:"column:repo_did"` - Repo *Repo `gorm:"foreignKey:DID;references:RepoDID"` - PinnedMessage string `gorm:"column:pinned_message"` - ExpiresAt *time.Time `gorm:"column:expires_at"` - CreatedAt time.Time `gorm:"column:created_at"` -} -``` - -Methods: - -- `CreatePinnedRecord(ctx, pin)` - insert -- `GetPinnedRecord(ctx, rkey)` - single lookup -- `DeletePinnedRecord(ctx, rkey)` - delete by rkey -- `GetActivePinnedRecord(ctx, streamerDID)` - returns the most recent non-expired pin for a streamer -- `DeleteAllPinnedRecords(ctx, streamerDID)` - bulk delete (called before creating new pin) - -Add `PinnedRecord{}` to the AutoMigrate list in `pkg/model/model.go`. - ---- - -## 6. Constants + Firehose - -### 6a. Constants - -**File**: `pkg/constants/constants.go` - -Add: `var PLACE_STREAM_CHAT_PINNED_RECORD = "place.stream.chat.pinnedRecord"` - -### 6b. Sync (create/update) - -**File**: `pkg/atproto/sync.go` - `handleCreateUpdate` - -Add case `*streamplace.ChatPinnedRecord`: - -- Sync bluesky repo -- Delete existing pinned records for this streamer (single-pin enforcement at DB level) -- Create new `PinnedRecord` model entry -- Build a hydrated view (resolve the pinned message, include author info) and publish to bus on streamer's channel - -### 6c. Firehose (delete) - -**File**: `pkg/atproto/firehose.go` - `EvtKindDeleteRecord` - -Add handling for `constants.PLACE_STREAM_CHAT_PINNED_RECORD`: - -- `DeletePinnedRecord(ctx, rkey)` -- Publish deletion marker to bus - ---- - -## 7. Bus / WebSocket Delivery - -### Create event - -Publish a `PinnedRecordView`-like object to the streamer's bus channel: - -```json -{ - "$type": "place.stream.chat.defs#pinnedRecordView", - "uri": "at://...", - "cid": "...", - "record": { - "pinnedMessage": "at://...", - "createdAt": "...", - "expiresAt": "..." - }, - "indexedAt": "...", - "message": { - /* hydrated messageView with author, text, facets, etc */ - } -} -``` - -### Delete event - -Publish a deletion marker: - -```json -{ - "$type": "place.stream.chat.pinnedRecord", - "deleted": true, - "rkey": "..." -} -``` - ---- - -## 8. Frontend State - -### LivestreamState additions - -**File**: `js/components/src/livestream-store/livestream-state.tsx` - -```ts -pinnedComment: PinnedRecordView | null; -``` - -Where `PinnedRecordView` is a type with the hydrated message + pin metadata. - -### Store init - -**File**: `js/components/src/livestream-store/livestream-store.tsx` - -Add `pinnedComment: null` to initial state. - -### Chat hooks - -**File**: `js/components/src/livestream-store/chat.tsx` - -Add: - -- `usePinnedComment()` - selector: `state.pinnedComment` -- `usePinChatMessage()` - hook to pin a message: - - If streamer: direct `com.atproto.repo.createRecord` for `place.stream.chat.pinnedRecord` - - If mod: call `place.stream.moderation.createPin` XRPC - - On success, the bus event will update state automatically -- `useUnpinChatMessage()` - hook to unpin: - - If streamer: direct `com.atproto.repo.deleteRecord` - - If mod: call `place.stream.moderation.deletePin` XRPC - -### WebSocket consumer - -**File**: `js/components/src/livestream-store/websocket-consumer.tsx` - -Add handler for `place.stream.chat.defs#pinnedRecordView`: - -- Set `state.pinnedComment` to the hydrated view - -Add handler for `place.stream.chat.pinnedRecord` with `deleted: true`: - -- Set `state.pinnedComment` to `null` - -### Expiration - -In the Chat component, `useEffect` checks `pinnedComment.record.expiresAt`. If current time passes it, clear `pinnedComment` locally. Also clear when `livestream.record.endedAt` is set (stream ended). - ---- - -## 9. Frontend UI - -### Pinned comment as persistent stream notification - -The pinned comment uses the existing `StreamNotificationProvider` + `streamNotificationManager` system. - -**File**: `js/components/src/components/stream-notification/pinned-comment-notification.tsx` (new) - -A custom render function passed to `streamNotificationManager.show()` with `duration: 0` (manual dismiss only), making it a persistent notification that sits at the top of the notification stack. Other temporary notifications (teleport, etc.) will appear above/below it and dismiss independently. - -The render function receives `(isExiting, onDismiss, startTime)` and renders: - -- Pin icon + author name (with chatProfile color) + message text (with facets rendered) -- Two dismiss actions: - - **Unpin** (X icon, only visible if `canPin`): calls `useUnpinChatMessage()`, removes the pin globally, then calls `onDismiss("user")` to dismiss the notification - - **Hide** (eye-off icon, visible to all viewers): calls `onDismiss("user")` to dismiss the notification locally only. Does NOT call the server - the pin remains for other viewers -- Does NOT link to or scroll to the original message (chat history is limited) -- Auto-dismisses when `expiresAt` is reached (a `useEffect` watching the pinned comment state calls `streamNotificationManager.requestDismiss("pinned-comment", "auto")`) -- Auto-dismisses when stream ends - -**Integration point**: The `TeleportWatcher` / livestream provider already manages notification lifecycle. The pinned comment notification is managed similarly - when `state.pinnedComment` changes in the Zustand store, a hook or effect calls `streamNotificationManager.show()` to create/update the notification, or `streamNotificationManager.hide()` to remove it. - -The notification ID should be fixed (e.g., `"pinned-comment"`) so that replacing a pin updates the same notification rather than creating a new one. - -### Pin action in mod menu - -**File**: `js/components/src/components/chat/mod-view.tsx` - -Add "Pin this message" in the moderation actions group (visible when `canPin` is true): - -```tsx -{ - modPermissions.canPin && ( - { - /* pin message */ - }} - > - Pin this message - - ); -} -``` - ---- - -## 10. Moderator Panel - -**File**: `js/components/src/components/dashboard/moderator-panel.tsx` - -Add `"message.pin"` to the list of assignable permissions when adding/editing moderators. - ---- - -## 11. Code Generation - -After creating lexicon JSON files, run `lexgen` to generate Go types: - -- `pkg/streamplace/chatpinnedrecord.go` (auto-generated) -- Update generated TypeScript types in the `streamplace` package - ---- - -## File Change Summary - -| File | Action | -| ---------------------------------------------------------------------------------- | ------------------------------------------------------- | -| `lexicons/place/stream/chat/pinnedRecord.json` | **New** | -| `lexicons/place/stream/chat/defs.json` | **Edit** - add `pinnedRecordView` | -| `lexicons/place/stream/moderation/createPin.json` | **New** | -| `lexicons/place/stream/moderation/deletePin.json` | **New** | -| `lexicons/place/stream/moderation/permission.json` | **Edit** - add `message.pin` | -| `pkg/constants/constants.go` | **Edit** | -| `pkg/moderation/permissions.go` | **Edit** | -| `pkg/model/pinned_record.go` | **New** | -| `pkg/model/model.go` | **Edit** | -| `pkg/atproto/sync.go` | **Edit** | -| `pkg/atproto/firehose.go` | **Edit** | -| `pkg/spxrpc/place_stream_moderation.go` | **Edit** | -| `js/components/src/streamplace-store/moderation.tsx` | **Edit** | -| `js/components/src/streamplace-store/block.tsx` | **Edit** - add pin/unpin hooks | -| `js/components/src/livestream-store/livestream-state.tsx` | **Edit** | -| `js/components/src/livestream-store/livestream-store.tsx` | **Edit** | -| `js/components/src/livestream-store/chat.tsx` | **Edit** | -| `js/components/src/livestream-store/websocket-consumer.tsx` | **Edit** | -| `js/components/src/components/stream-notification/pinned-comment-notification.tsx` | **New** - notification render function | -| `js/components/src/livestream-provider/index.tsx` | **Edit** - manage pinned comment notification lifecycle | -| `js/components/src/components/chat/mod-view.tsx` | **Edit** - add pin action | -| `js/components/src/components/dashboard/moderator-panel.tsx` | **Edit** | - ---- - -## Implementation Order - -1. Lexicons (pinnedRecord, createPin, deletePin, permission, defs) -2. Code generation (lexgen) -3. Backend: constants, permissions, model, DB migration -4. Backend: sync/firehose integration -5. Backend: XRPC handlers + registration -6. Frontend: state types, store init, hooks -7. Frontend: WebSocket consumer handlers -8. Frontend: pinned-comment notification render function + lifecycle management -9. Frontend: mod-view pin action -10. Frontend: moderator panel permission diff --git a/js/components/src/components/chat/mod-view.tsx b/js/components/src/components/chat/mod-view.tsx index afdc1976..dc52daa6 100644 --- a/js/components/src/components/chat/mod-view.tsx +++ b/js/components/src/components/chat/mod-view.tsx @@ -237,53 +237,69 @@ function ModViewContent({ )} {modPermissions.canPin && message.author.did !== streamerDID && ( - - - Pin this message - - - { - if (!streamerDID) return; - pinChatMessage(message.uri, streamerDID) - .then(() => { - toast.show("Comment pinned", "", { duration: 3 }); - onOpenChange?.(false); - }) - .catch((e) => { - toast.show( - "Error pinning comment", - e instanceof Error ? e.message : "Failed to pin", - { duration: 5 }, - ); - }); - }} - > - Pin indefinitely - - { - if (!streamerDID) return; - const expiresAt = new Date(); - expiresAt.setHours(expiresAt.getHours() + 1); // Set expiration to 1 hour from now - pinChatMessage(message.uri, streamerDID, expiresAt) - .then(() => { - toast.show("Comment pinned", "", { duration: 3 }); - onOpenChange?.(false); - }) - .catch((e) => { - toast.show( - "Error pinning comment", - e instanceof Error ? e.message : "Failed to pin", - { duration: 5 }, - ); - }); - }} - > - Pin for 1 hour - - - + + + + Pin this message + + + + { + if (!streamerDID) return; + pinChatMessage(message.uri, streamerDID) + .then(() => { + toast.show("Comment pinned", "", { duration: 3 }); + onOpenChange?.(false); + }) + .catch((e) => { + toast.show( + "Error pinning comment", + e instanceof Error ? e.message : "Failed to pin", + { duration: 5 }, + ); + }); + }} + > + Until stream end + + {[5, 10, 15, 30, 60].map((minutes) => ( + { + if (!streamerDID) return; + const expiresAt = new Date( + Date.now() + minutes * 60 * 1000, + ); + pinChatMessage( + message.uri, + streamerDID, + expiresAt.toISOString(), + ) + .then(() => { + toast.show("Comment pinned", "", { duration: 3 }); + onOpenChange?.(false); + }) + .catch((e) => { + toast.show( + "Error pinning comment", + e instanceof Error + ? e.message + : "Failed to pin", + { duration: 5 }, + ); + }); + }} + > + + {minutes < 60 ? `${minutes} min` : "1 hour"} + + + ))} + + + + )} {modPermissions.canBan && agent?.did && diff --git a/js/components/src/components/stream-notification/pin-notification.tsx b/js/components/src/components/stream-notification/pin-notification.tsx new file mode 100644 index 00000000..a56766b3 --- /dev/null +++ b/js/components/src/components/stream-notification/pin-notification.tsx @@ -0,0 +1,135 @@ +import { EyeOff, Pin, X } from "lucide-react-native"; +import { useEffect, useState } from "react"; +import { Linking, Pressable, View } from "react-native"; +import { PinnedRecordViewHydrated, PlaceStreamChatProfile } from "streamplace"; +import { + Text, + useCanModerate, + useLivestreamStore, + useTheme, + zero, +} from "../../"; +import { RichtextSegment, segmentize } from "../../lib/facet"; +import { formatHandleWithAt } from "../../utils/format-handle"; + +const getRgbColor = (color?: PlaceStreamChatProfile.Color) => + color ? `rgb(${color.red}, ${color.green}, ${color.blue})` : undefined; + +function renderSegment(segment: RichtextSegment, index: number) { + if (segment.features && segment.features.length > 0) { + const ftr = segment.features[0]; + if (ftr.$type === "app.bsky.richtext.facet#link") { + return ( + Linking.openURL((ftr as any).uri || "")} + > + {segment.text} + + ); + } + if (ftr.$type === "app.bsky.richtext.facet#mention") { + return ( + + {segment.text} + + ); + } + } + return {segment.text}; +} + +export function PinnedCommentNotification({ + pinnedComment, + onDismiss, + onUnpin, +}: { + pinnedComment: PinnedRecordViewHydrated; + onDismiss: () => void; + onUnpin: () => void; +}) { + const z = useTheme(); + const message = pinnedComment.message; + const pinnedByColor = (pinnedComment.pinnedBy as any)?.color || "#bebebe"; + const record = pinnedComment.record; + + const currentStreamer = useLivestreamStore((state) => state.profile?.did); + + console.log("checking if we can mod", currentStreamer); + + const canActuallyPin = useCanModerate(currentStreamer)?.canPin; + + const messageRecord = message?.record as any; + + const [expiresAt] = useState( + record.expiresAt ? new Date(record.expiresAt) : null, + ); + + useEffect(() => { + if (!expiresAt) return; + const remaining = expiresAt.getTime() - Date.now(); + if (remaining <= 0) { + onDismiss(); + return; + } + const timeout = setTimeout(onDismiss, remaining); + return () => clearTimeout(timeout); + }, [expiresAt, onDismiss]); + + const authorName = message ? formatHandleWithAt(message.author) : "unknown"; + const authorColor = getRgbColor((message as any)?.chatProfile?.color); + + const segments = messageRecord + ? segmentize(messageRecord.text, messageRecord.facets) + : []; + + return ( + + + + + + + + + {authorName} + + {segments.map((seg, i) => renderSegment(seg, i))} + + + + {canActuallyPin && ( + + + + )} + + + + + + + ); +} diff --git a/js/components/src/lib/stream-notifications.ts b/js/components/src/lib/stream-notifications.ts index e05c4d56..4bbf40ad 100644 --- a/js/components/src/lib/stream-notifications.ts +++ b/js/components/src/lib/stream-notifications.ts @@ -1,8 +1,36 @@ import React from "react"; +import { PinnedRecordViewHydrated } from "streamplace"; import { streamNotification } from "../components/stream-notification"; +import { PinnedCommentNotification } from "../components/stream-notification/pin-notification"; import { TeleportNotification } from "../components/stream-notification/teleport-notification"; export const StreamNotifications = { + pinnedComment: (params: { + pinnedComment: PinnedRecordViewHydrated; + onDismiss?: (reason?: "user" | "auto") => void; + onUnpin?: () => void; + }) => { + streamNotification.show({ + id: "pinned-comment", + render: (isExiting, onDismiss) => { + return React.createElement(PinnedCommentNotification, { + pinnedComment: params.pinnedComment, + onDismiss: () => onDismiss("user"), + onUnpin: () => { + params.onUnpin?.(); + onDismiss("user"); + }, + }); + }, + duration: 0, // manually dismissed or auto-dismissed by TTL + onDismiss: params.onDismiss, + }); + }, + + pinnedCommentDismiss: () => { + streamNotification.hide("pinned-comment"); + }, + teleport: (params: { targetHandle: string; targetDID: string; diff --git a/js/components/src/livestream-provider/index.tsx b/js/components/src/livestream-provider/index.tsx index 1a752cab..e40e2249 100644 --- a/js/components/src/livestream-provider/index.tsx +++ b/js/components/src/livestream-provider/index.tsx @@ -6,6 +6,8 @@ import { LivestreamContext, makeLivestreamStore, useLivestreamStore, + usePinnedComment, + useUnpinChatMessage, } from "../livestream-store"; import { useDID, usePDSAgent } from "../streamplace-store"; import { useLivestreamWebsocket } from "./websocket"; @@ -134,6 +136,39 @@ export function TeleportWatcher({ return <>; } +export function PinnedCommentWatcher() { + const pinnedComment = usePinnedComment(); + const streamerDID = useLivestreamStore((state) => state.profile?.did); + const unpinChatMessage = useUnpinChatMessage(); + const prevPinnedRef = useRef(null); + + // Show/hide notification when pinned comment changes + useEffect(() => { + const currentUri = pinnedComment?.uri ?? null; + if (currentUri === prevPinnedRef.current) return; + prevPinnedRef.current = currentUri; + + if (pinnedComment) { + StreamNotifications.pinnedComment({ + pinnedComment, + onDismiss: () => { + // local dismiss + }, + onUnpin: () => { + if (!streamerDID) return; + unpinChatMessage(pinnedComment.uri, streamerDID).catch((e) => { + console.error("Failed to unpin:", e); + }); + }, + }); + } else { + StreamNotifications.pinnedCommentDismiss(); + } + }, [pinnedComment, streamerDID, unpinChatMessage]); + + return <>; +} + export function LivestreamPoller({ children, src, @@ -143,12 +178,11 @@ export function LivestreamPoller({ src: string; onTeleport?: (targetHandle: string, targetDID: string) => void; }) { - // Websocket watcher is a sibling instead of a parent to avoid - // re-rendering when the websocket does stuff return ( <> + {children} ); diff --git a/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md index bb73679b..1d4a71ca 100644 --- a/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md +++ b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-defs.md @@ -39,14 +39,14 @@ View of a pinned chat record with hydrated message data. **Properties:** -| Name | Type | Req'd | Description | Constraints | -| ----------- | ------------------------------------------------------------------------------------------------------------------------------------------------ | ----- | ----------- | ------------------ | -| `uri` | `string` | ✅ | | Format: `at-uri` | -| `cid` | `string` | ✅ | | Format: `cid` | -| `record` | [`place.stream.chat.pinnedRecord`](/lex-reference/place-stream-chat-pinnedrecord) | ✅ | | | -| `indexedAt` | `string` | ✅ | | Format: `datetime` | -| `pinnedBy` | [`app.bsky.actor.defs#profileViewBasic`](https://github.com/bluesky-social/atproto/tree/main/lexicons/app/bsky/actor/defs.json#profileViewBasic) | ❌ | | | -| `message` | [`#messageView`](#messageview) | ❌ | | | +| Name | Type | Req'd | Description | Constraints | +| ----------- | --------------------------------------------------------------------------------- | ----- | ----------- | ------------------ | +| `uri` | `string` | ✅ | | Format: `at-uri` | +| `cid` | `string` | ✅ | | Format: `cid` | +| `record` | [`place.stream.chat.pinnedRecord`](/lex-reference/place-stream-chat-pinnedrecord) | ✅ | | | +| `indexedAt` | `string` | ✅ | | Format: `datetime` | +| `pinnedBy` | [`place.stream.chat.profile`](/lex-reference/place-stream-chat-profile) | ❌ | | | +| `message` | [`#messageView`](#messageview) | ❌ | | | --- @@ -126,7 +126,7 @@ View of a pinned chat record with hydrated message data. }, "pinnedBy": { "type": "ref", - "ref": "app.bsky.actor.defs#profileViewBasic" + "ref": "place.stream.chat.profile" }, "message": { "type": "ref", diff --git a/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-pinnedrecord.md b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-pinnedrecord.md index 8935e957..fc88032a 100644 --- a/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-pinnedrecord.md +++ b/js/docs/src/content/docs/lex-reference/chat/place-stream-chat-pinnedrecord.md @@ -22,6 +22,7 @@ Record pinning a chat message for prominent display. | Name | Type | Req'd | Description | Constraints | | --------------- | -------- | ----- | --------------------------------------------------------------------------------- | ------------------ | | `pinnedMessage` | `string` | ✅ | AT-URI of the pinned chat message. | Format: `at-uri` | +| `pinnedBy` | `string` | ❌ | DID of the user who pinned the message. | Format: `did` | | `createdAt` | `string` | ✅ | When this pin was created. | Format: `datetime` | | `expiresAt` | `string` | ❌ | Optional expiration time. If set, the pin is considered inactive after this time. | Format: `datetime` | @@ -47,6 +48,11 @@ Record pinning a chat message for prominent display. "format": "at-uri", "description": "AT-URI of the pinned chat message." }, + "pinnedBy": { + "type": "string", + "format": "did", + "description": "DID of the user who pinned the message." + }, "createdAt": { "type": "string", "format": "datetime", diff --git a/js/docs/src/content/docs/lex-reference/place-stream-livestream.md b/js/docs/src/content/docs/lex-reference/place-stream-livestream.md index 2825bf4c..db8ee166 100644 --- a/js/docs/src/content/docs/lex-reference/place-stream-livestream.md +++ b/js/docs/src/content/docs/lex-reference/place-stream-livestream.md @@ -123,9 +123,9 @@ Record announcing a livestream is happening **Properties:** -| Name | Type | Req'd | Description | Constraints | -| ------------ | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----- | ----------- | ----------- | -| `livestream` | Union of:
  [`#livestreamView`](#livestreamview)
  [`#viewerCount`](#viewercount)
  [`#teleportArrival`](#teleportarrival)
  [`#teleportCanceled`](#teleportcanceled)
  [`place.stream.defs#blockView`](/lex-reference/place-stream-defs#blockview)
  [`place.stream.defs#renditions`](/lex-reference/place-stream-defs#renditions)
  [`place.stream.defs#rendition`](/lex-reference/place-stream-defs#rendition)
  [`place.stream.chat.defs#messageView`](/lex-reference/place-stream-chat-defs#messageview) | ✅ | | | +| Name | Type | Req'd | Description | Constraints | +| ------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----- | ----------- | ----------- | +| `livestream` | Union of:
  [`#livestreamView`](#livestreamview)
  [`#viewerCount`](#viewercount)
  [`#teleportArrival`](#teleportarrival)
  [`#teleportCanceled`](#teleportcanceled)
  [`place.stream.defs#blockView`](/lex-reference/place-stream-defs#blockview)
  [`place.stream.defs#renditions`](/lex-reference/place-stream-defs#renditions)
  [`place.stream.defs#rendition`](/lex-reference/place-stream-defs#rendition)
  [`place.stream.chat.defs#messageView`](/lex-reference/place-stream-chat-defs#messageview)
  [`place.stream.chat.defs#pinnedRecordView`](/lex-reference/place-stream-chat-defs#pinnedrecordview) | ✅ | | | --- @@ -309,7 +309,8 @@ Record announcing a livestream is happening "place.stream.defs#blockView", "place.stream.defs#renditions", "place.stream.defs#rendition", - "place.stream.chat.defs#messageView" + "place.stream.chat.defs#messageView", + "place.stream.chat.defs#pinnedRecordView" ] } } diff --git a/lexicons/place/stream/chat/defs.json b/lexicons/place/stream/chat/defs.json index aa1dcd75..5f9c495d 100644 --- a/lexicons/place/stream/chat/defs.json +++ b/lexicons/place/stream/chat/defs.json @@ -48,7 +48,7 @@ "indexedAt": { "type": "string", "format": "datetime" }, "pinnedBy": { "type": "ref", - "ref": "app.bsky.actor.defs#profileViewBasic" + "ref": "place.stream.chat.profile" }, "message": { "type": "ref", diff --git a/lexicons/place/stream/chat/pinnedRecord.json b/lexicons/place/stream/chat/pinnedRecord.json index f43ac525..32b5b33e 100644 --- a/lexicons/place/stream/chat/pinnedRecord.json +++ b/lexicons/place/stream/chat/pinnedRecord.json @@ -15,6 +15,11 @@ "format": "at-uri", "description": "AT-URI of the pinned chat message." }, + "pinnedBy": { + "type": "string", + "format": "did", + "description": "DID of the user who pinned the message." + }, "createdAt": { "type": "string", "format": "datetime", diff --git a/lexicons/place/stream/livestream.json b/lexicons/place/stream/livestream.json index e02f79a7..19b73cf3 100644 --- a/lexicons/place/stream/livestream.json +++ b/lexicons/place/stream/livestream.json @@ -162,7 +162,8 @@ "place.stream.defs#blockView", "place.stream.defs#renditions", "place.stream.defs#rendition", - "place.stream.chat.defs#messageView" + "place.stream.chat.defs#messageView", + "place.stream.chat.defs#pinnedRecordView" ] } } diff --git a/pkg/api/websocket.go b/pkg/api/websocket.go index 418a406c..67ad07f6 100644 --- a/pkg/api/websocket.go +++ b/pkg/api/websocket.go @@ -259,6 +259,59 @@ func (a *StreamplaceAPI) HandleWebsocket(ctx context.Context) httprouter.Handle } }() + // get the latest active pinned message for the repo + go func() { + pin, err := a.Model.GetActivePinnedRecord(ctx, repoDID) + if err != nil { + log.Error(ctx, "could not get pinned record", "error", err) + return + } + if pin != nil { + prv, err := pin.ToStreamplacePinnedRecordView() + if err != nil { + log.Error(ctx, "could not convert pinned record to streamplace view", "error", err) + return + } + // look up the original message, pinner + msg, err := a.Model.GetChatMessage(prv.Record.PinnedMessage) + if err != nil { + log.Error(ctx, "failed to get pinned message", err) + return + } + // if the message was deleted, treat as no pinned message + if msg != nil && msg.DeletedAt != nil { + log.Log(ctx, "pinned message was deleted, skipping", "uri", msg.URI) + return + } + // if no pinned by, use the repo owner as the pinner + if prv.Record.PinnedBy == nil { + prv.Record.PinnedBy = &repoDID + } + profile, err := a.Model.GetChatProfile(ctx, *prv.Record.PinnedBy) + if err != nil { + log.Error(ctx, "failed to get chat profile", err) + return + } + if msg != nil { + msgView, err := msg.ToStreamplaceMessageView() + if err != nil { + log.Error(ctx, "failed to convert chat message: %w", err) + return + } + prv.Message = msgView + } + if profile != nil { + profileView, err := profile.ToStreamplaceChatProfile() + if err != nil { + log.Error(ctx, "failed to convert chat profile: %w", err) + return + } + prv.PinnedBy = profileView + } + initialBurst <- prv + } + }() + go func() { teleports, err := a.Model.GetActiveTeleportsToRepo(repoDID) if err != nil { diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index a739a695..53f93032 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -226,11 +226,10 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD return nil } log.Debug(ctx, "creating pinned record", "userDID", userDID, "pinnedMessage", rec.PinnedMessage) - // Delete existing pinned records for this streamer (single-pin semantics) - err = atsync.Model.DeleteAllPinnedRecords(ctx, userDID) - if err != nil { - log.Error(ctx, "failed to delete existing pinned records", "err", err) - } + // err = atsync.Model.DeleteAllPinnedRecords(ctx, userDID) + // if err != nil { + // log.Error(ctx, "failed to delete existing pinned records", "err", err) + // } // Parse optional expiresAt var expiresAt *time.Time if rec.ExpiresAt != nil { @@ -239,12 +238,27 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD expiresAt = &t } } + // serialise createdAt + createdAt, err := time.Parse(time.RFC3339, rec.CreatedAt) + if err != nil { + return fmt.Errorf("failed to parse createdAt: %w", err) + } + + var pinnedBy string + if rec.PinnedBy == nil { + pinnedBy = userDID + } else { + pinnedBy = *rec.PinnedBy + } + pin := &model.PinnedRecord{ RKey: rkey.String(), RepoDID: userDID, PinnedMessage: rec.PinnedMessage, + PinnedBy: pinnedBy, + IndexedAt: &now, CID: cid, - CreatedAt: now, + CreatedAt: createdAt, Repo: repo, ExpiresAt: expiresAt, } @@ -256,10 +270,33 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD if err != nil { return fmt.Errorf("failed to get pinned record after we just saved it: %w", err) } - pinnedView, err := pin.ToStreamplacePinnedRecord() + pinnedView, err := pin.ToStreamplacePinnedRecordView() if err != nil { return fmt.Errorf("failed to convert pinned record: %w", err) } + // look up the original message, pinner + msg, err := atsync.Model.GetChatMessage(pinnedView.Record.PinnedMessage) + if err != nil { + return fmt.Errorf("failed to get chat message: %w", err) + } + profile, err := atsync.Model.GetChatProfile(ctx, pinnedBy) + if err != nil { + return fmt.Errorf("failed to get chat profile: %w", err) + } + if msg != nil { + msgView, err := msg.ToStreamplaceMessageView() + if err != nil { + return fmt.Errorf("failed to convert chat message: %w", err) + } + pinnedView.Message = msgView + } + if profile != nil { + profileView, err := profile.ToStreamplaceChatProfile() + if err != nil { + return fmt.Errorf("failed to convert chat profile: %w", err) + } + pinnedView.PinnedBy = profileView + } go atsync.Bus.Publish(userDID, pinnedView) case *streamplace.ChatProfile: diff --git a/pkg/model/pinned_record.go b/pkg/model/pinned_record.go index 90add2f7..10789d9a 100644 --- a/pkg/model/pinned_record.go +++ b/pkg/model/pinned_record.go @@ -15,6 +15,8 @@ type PinnedRecord struct { RepoDID string `json:"repoDID" gorm:"column:repo_did"` Repo *Repo `json:"repo,omitempty" gorm:"foreignKey:DID;references:RepoDID"` PinnedMessage string `gorm:"column:pinned_message" json:"pinnedMessage"` + PinnedBy string `gorm:"column:pinned_by" json:"pinnedBy"` + IndexedAt *time.Time `gorm:"column:indexed_at" json:"indexedAt"` ExpiresAt *time.Time `gorm:"column:expires_at" json:"expiresAt"` CreatedAt time.Time `gorm:"column:created_at" json:"createdAt"` } @@ -32,6 +34,26 @@ func (p *PinnedRecord) ToStreamplacePinnedRecord() (*streamplace.ChatPinnedRecor return rec, nil } +func (p *PinnedRecord) ToStreamplacePinnedRecordView() (*streamplace.ChatDefs_PinnedRecordView, error) { + pr := &streamplace.ChatPinnedRecord{ + LexiconTypeID: "place.stream.chat.pinnedRecord", + PinnedMessage: p.PinnedMessage, + CreatedAt: p.CreatedAt.UTC().Format(time.RFC3339), + } + if p.ExpiresAt != nil { + s := p.ExpiresAt.UTC().Format(time.RFC3339) + pr.ExpiresAt = &s + } + rec := &streamplace.ChatDefs_PinnedRecordView{ + LexiconTypeID: "place.stream.chat.defs#pinnedRecordView", + Record: pr, + Cid: p.CID, + IndexedAt: p.CreatedAt.UTC().Format(time.RFC3339Nano), + // message, pinnedby not included, will fill in later + Uri: "at://" + p.RepoDID + "/place.stream.chat.pinnedRecord/" + p.RKey, + } + return rec, nil +} func (m *DBModel) CreatePinnedRecord(ctx context.Context, pin *PinnedRecord) error { return m.DB.Create(pin).Error } diff --git a/pkg/streamplace/cbor_gen.go b/pkg/streamplace/cbor_gen.go index 197bd3b9..6d7456dc 100644 --- a/pkg/streamplace/cbor_gen.go +++ b/pkg/streamplace/cbor_gen.go @@ -3641,12 +3641,16 @@ func (t *ChatPinnedRecord) MarshalCBOR(w io.Writer) error { } cw := cbg.NewCborWriter(w) - fieldCount := 4 + fieldCount := 5 if t.ExpiresAt == nil { fieldCount-- } + if t.PinnedBy == nil { + fieldCount-- + } + if _, err := cw.Write(cbg.CborEncodeMajorType(cbg.MajMap, uint64(fieldCount))); err != nil { return err } @@ -3670,6 +3674,38 @@ func (t *ChatPinnedRecord) MarshalCBOR(w io.Writer) error { return err } + // t.PinnedBy (string) (string) + if t.PinnedBy != nil { + + if len("pinnedBy") > 1000000 { + return xerrors.Errorf("Value in field \"pinnedBy\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("pinnedBy"))); err != nil { + return err + } + if _, err := cw.WriteString(string("pinnedBy")); err != nil { + return err + } + + if t.PinnedBy == nil { + if _, err := cw.Write(cbg.CborNull); err != nil { + return err + } + } else { + if len(*t.PinnedBy) > 1000000 { + return xerrors.Errorf("Value in field t.PinnedBy was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(*t.PinnedBy))); err != nil { + return err + } + if _, err := cw.WriteString(string(*t.PinnedBy)); err != nil { + return err + } + } + } + // t.CreatedAt (string) (string) if len("createdAt") > 1000000 { return xerrors.Errorf("Value in field \"createdAt\" was too long") @@ -3802,6 +3838,27 @@ func (t *ChatPinnedRecord) UnmarshalCBOR(r io.Reader) (err error) { t.LexiconTypeID = string(sval) } + // t.PinnedBy (string) (string) + case "pinnedBy": + + { + b, err := cr.ReadByte() + if err != nil { + return err + } + if b != cbg.CborNull[0] { + if err := cr.UnreadByte(); err != nil { + return err + } + + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.PinnedBy = (*string)(&sval) + } + } // t.CreatedAt (string) (string) case "createdAt": diff --git a/pkg/streamplace/chatdefs.go b/pkg/streamplace/chatdefs.go index 0f2aedac..bbb89b04 100644 --- a/pkg/streamplace/chatdefs.go +++ b/pkg/streamplace/chatdefs.go @@ -59,10 +59,11 @@ func (t *ChatDefs_MessageView_ReplyTo) UnmarshalJSON(b []byte) error { // // View of a pinned chat record with hydrated message data. type ChatDefs_PinnedRecordView struct { - Cid string `json:"cid" cborgen:"cid"` - IndexedAt string `json:"indexedAt" cborgen:"indexedAt"` - Message *ChatDefs_MessageView `json:"message,omitempty" cborgen:"message,omitempty"` - PinnedBy *appbsky.ActorDefs_ProfileViewBasic `json:"pinnedBy,omitempty" cborgen:"pinnedBy,omitempty"` - Record *ChatPinnedRecord `json:"record" cborgen:"record"` - Uri string `json:"uri" cborgen:"uri"` + LexiconTypeID string `json:"$type" cborgen:"$type,const=place.stream.chat.defs#pinnedRecordView"` + Cid string `json:"cid" cborgen:"cid"` + IndexedAt string `json:"indexedAt" cborgen:"indexedAt"` + Message *ChatDefs_MessageView `json:"message,omitempty" cborgen:"message,omitempty"` + PinnedBy *ChatProfile `json:"pinnedBy,omitempty" cborgen:"pinnedBy,omitempty"` + Record *ChatPinnedRecord `json:"record" cborgen:"record"` + Uri string `json:"uri" cborgen:"uri"` } diff --git a/pkg/streamplace/chatpinnedRecord.go b/pkg/streamplace/chatpinnedRecord.go index b61cddfd..7c013214 100644 --- a/pkg/streamplace/chatpinnedRecord.go +++ b/pkg/streamplace/chatpinnedRecord.go @@ -18,6 +18,8 @@ type ChatPinnedRecord struct { CreatedAt string `json:"createdAt" cborgen:"createdAt"` // expiresAt: Optional expiration time. If set, the pin is considered inactive after this time. ExpiresAt *string `json:"expiresAt,omitempty" cborgen:"expiresAt,omitempty"` + // pinnedBy: DID of the user who pinned the message. + PinnedBy *string `json:"pinnedBy,omitempty" cborgen:"pinnedBy,omitempty"` // pinnedMessage: AT-URI of the pinned chat message. PinnedMessage string `json:"pinnedMessage" cborgen:"pinnedMessage"` } diff --git a/pkg/streamplace/streamlivestream.go b/pkg/streamplace/streamlivestream.go index aef26c79..11aefb67 100644 --- a/pkg/streamplace/streamlivestream.go +++ b/pkg/streamplace/streamlivestream.go @@ -73,6 +73,7 @@ type Livestream_StreamplaceAnything_Livestream struct { Defs_Renditions *Defs_Renditions Defs_Rendition *Defs_Rendition ChatDefs_MessageView *ChatDefs_MessageView + ChatDefs_PinnedRecordView *ChatDefs_PinnedRecordView } func (t *Livestream_StreamplaceAnything_Livestream) MarshalJSON() ([]byte, error) { @@ -108,6 +109,10 @@ func (t *Livestream_StreamplaceAnything_Livestream) MarshalJSON() ([]byte, error t.ChatDefs_MessageView.LexiconTypeID = "place.stream.chat.defs#messageView" return json.Marshal(t.ChatDefs_MessageView) } + if t.ChatDefs_PinnedRecordView != nil { + t.ChatDefs_PinnedRecordView.LexiconTypeID = "place.stream.chat.defs#pinnedRecordView" + return json.Marshal(t.ChatDefs_PinnedRecordView) + } return nil, fmt.Errorf("can not marshal empty union as JSON") } @@ -142,6 +147,9 @@ func (t *Livestream_StreamplaceAnything_Livestream) UnmarshalJSON(b []byte) erro case "place.stream.chat.defs#messageView": t.ChatDefs_MessageView = new(ChatDefs_MessageView) return json.Unmarshal(b, t.ChatDefs_MessageView) + case "place.stream.chat.defs#pinnedRecordView": + t.ChatDefs_PinnedRecordView = new(ChatDefs_PinnedRecordView) + return json.Unmarshal(b, t.ChatDefs_PinnedRecordView) default: return nil } -- 2.51.2 From 0e30341ff4229cefbee8281c213cdce9e6b81dae Mon Sep 17 00:00:00 2001 From: "Natalie B." <22222885+espeon@users.noreply.github.com> Date: Mon, 23 Mar 2026 18:27:14 -0500 Subject: [PATCH 4/7] change over to aturi --- js/components/src/components/chat/mod-view.tsx | 11 ++++++----- js/components/src/components/ui/dropdown.tsx | 5 +++-- pkg/atproto/sync.go | 4 ++-- pkg/model/model.go | 4 ++-- pkg/model/pinned_record.go | 12 ++++++------ 5 files changed, 19 insertions(+), 17 deletions(-) diff --git a/js/components/src/components/chat/mod-view.tsx b/js/components/src/components/chat/mod-view.tsx index dc52daa6..09151ad1 100644 --- a/js/components/src/components/chat/mod-view.tsx +++ b/js/components/src/components/chat/mod-view.tsx @@ -97,9 +97,7 @@ export const ModView = forwardRef(() => { agent?.did && ((modPermissions.canHide && message.author.did !== streamerDID) || (modPermissions.canPin && message.author.did !== streamerDID) || - (modPermissions.canBan && - message.author.did !== agent.did && - message.author.did !== streamerDID)) + modPermissions.canBan) ); return ( @@ -236,10 +234,13 @@ function ModViewContent({ )} - {modPermissions.canPin && message.author.did !== streamerDID && ( + {modPermissions.canPin && ( - + Pin this message diff --git a/js/components/src/components/ui/dropdown.tsx b/js/components/src/components/ui/dropdown.tsx index 43de0712..a1ea5656 100644 --- a/js/components/src/components/ui/dropdown.tsx +++ b/js/components/src/components/ui/dropdown.tsx @@ -74,7 +74,7 @@ export const DropdownMenuSubTrigger = forwardRef< inset?: boolean; children?: React.ReactNode; } ->(({ inset, children, subMenuTitle, ...props }, ref) => { +>(({ inset, children, subMenuTitle, style, ...props }, ref) => { const { icons } = useTheme(); const { open } = DropdownMenuPrimitive.useSubContext(); const Icon = @@ -96,6 +96,7 @@ export const DropdownMenuSubTrigger = forwardRef< layout.flex.alignCenter, p[2], pr[8], + style, ]} > {children} @@ -513,7 +514,7 @@ export const DropdownMenuGroup = forwardRef< const { theme } = useTheme(); const { inset, title, children, ...rest } = props; return ( - + {title && ( {title} diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 53f93032..e53877fc 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -252,7 +252,7 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD } pin := &model.PinnedRecord{ - RKey: rkey.String(), + Uri: aturi.String(), RepoDID: userDID, PinnedMessage: rec.PinnedMessage, PinnedBy: pinnedBy, @@ -266,7 +266,7 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD if err != nil { return fmt.Errorf("failed to create pinned record: %w", err) } - pin, err = atsync.Model.GetPinnedRecord(ctx, rkey.String()) + pin, err = atsync.Model.GetPinnedRecord(ctx, pin.Uri) if err != nil { return fmt.Errorf("failed to get pinned record after we just saved it: %w", err) } diff --git a/pkg/model/model.go b/pkg/model/model.go index 5571182f..b524b507 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -86,9 +86,9 @@ type Model interface { GetUserGates(ctx context.Context, userDID string) ([]*Gate, error) CreatePinnedRecord(ctx context.Context, pin *PinnedRecord) error - DeletePinnedRecord(ctx context.Context, rkey string) error + DeletePinnedRecord(ctx context.Context, uri string) error DeleteAllPinnedRecords(ctx context.Context, streamerDID string) error - GetPinnedRecord(ctx context.Context, rkey string) (*PinnedRecord, error) + GetPinnedRecord(ctx context.Context, uri string) (*PinnedRecord, error) GetActivePinnedRecord(ctx context.Context, streamerDID string) (*PinnedRecord, error) CreateChatProfile(ctx context.Context, profile *ChatProfile) error diff --git a/pkg/model/pinned_record.go b/pkg/model/pinned_record.go index 10789d9a..67152709 100644 --- a/pkg/model/pinned_record.go +++ b/pkg/model/pinned_record.go @@ -10,7 +10,7 @@ import ( ) type PinnedRecord struct { - RKey string `gorm:"primaryKey;column:rkey"` + Uri string `gorm:"primaryKey;column:uri"` CID string `gorm:"column:cid"` RepoDID string `json:"repoDID" gorm:"column:repo_did"` Repo *Repo `json:"repo,omitempty" gorm:"foreignKey:DID;references:RepoDID"` @@ -50,7 +50,7 @@ func (p *PinnedRecord) ToStreamplacePinnedRecordView() (*streamplace.ChatDefs_Pi Cid: p.CID, IndexedAt: p.CreatedAt.UTC().Format(time.RFC3339Nano), // message, pinnedby not included, will fill in later - Uri: "at://" + p.RepoDID + "/place.stream.chat.pinnedRecord/" + p.RKey, + Uri: p.Uri, } return rec, nil } @@ -58,9 +58,9 @@ func (m *DBModel) CreatePinnedRecord(ctx context.Context, pin *PinnedRecord) err return m.DB.Create(pin).Error } -func (m *DBModel) GetPinnedRecord(ctx context.Context, rkey string) (*PinnedRecord, error) { +func (m *DBModel) GetPinnedRecord(ctx context.Context, uri string) (*PinnedRecord, error) { var pin PinnedRecord - err := m.DB.Preload("Repo").Where("rkey = ?", rkey).First(&pin).Error + err := m.DB.Preload("Repo").Where("uri = ?", uri).First(&pin).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } @@ -70,8 +70,8 @@ func (m *DBModel) GetPinnedRecord(ctx context.Context, rkey string) (*PinnedReco return &pin, nil } -func (m *DBModel) DeletePinnedRecord(ctx context.Context, rkey string) error { - return m.DB.Where("rkey = ?", rkey).Delete(&PinnedRecord{}).Error +func (m *DBModel) DeletePinnedRecord(ctx context.Context, uri string) error { + return m.DB.Where("uri = ?", uri).Delete(&PinnedRecord{}).Error } func (m *DBModel) DeleteAllPinnedRecords(ctx context.Context, streamerDID string) error { -- 2.51.2 From 3403fdce56510cf5809b9b3955d1b22a0dbd70ca Mon Sep 17 00:00:00 2001 From: "Natalie B." <22222885+espeon@users.noreply.github.com> Date: Mon, 23 Mar 2026 18:35:07 -0500 Subject: [PATCH 5/7] wire pin expiry up --- js/components/src/livestream-provider/index.tsx | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/js/components/src/livestream-provider/index.tsx b/js/components/src/livestream-provider/index.tsx index e40e2249..716e6e07 100644 --- a/js/components/src/livestream-provider/index.tsx +++ b/js/components/src/livestream-provider/index.tsx @@ -4,6 +4,7 @@ import { deleteTeleport } from "../lib/slash-commands/teleport"; import { StreamNotifications } from "../lib/stream-notifications"; import { LivestreamContext, + getStoreFromContext, makeLivestreamStore, useLivestreamStore, usePinnedComment, @@ -140,6 +141,7 @@ export function PinnedCommentWatcher() { const pinnedComment = usePinnedComment(); const streamerDID = useLivestreamStore((state) => state.profile?.did); const unpinChatMessage = useUnpinChatMessage(); + const store = getStoreFromContext(); const prevPinnedRef = useRef(null); // Show/hide notification when pinned comment changes @@ -152,7 +154,7 @@ export function PinnedCommentWatcher() { StreamNotifications.pinnedComment({ pinnedComment, onDismiss: () => { - // local dismiss + store.setState({ pinnedComment: null }); }, onUnpin: () => { if (!streamerDID) return; -- 2.51.2 From 3db5430c8b168f7188368078caaac42c4e50289a Mon Sep 17 00:00:00 2001 From: "Natalie B." <22222885+espeon@users.noreply.github.com> Date: Mon, 23 Mar 2026 18:44:38 -0500 Subject: [PATCH 6/7] rework chat popout and add notif provider --- js/app/src/screens/chat-popout.native.tsx | 5 +- js/app/src/screens/chat-popout.tsx | 101 +++++++++++++--------- 2 files changed, 63 insertions(+), 43 deletions(-) diff --git a/js/app/src/screens/chat-popout.native.tsx b/js/app/src/screens/chat-popout.native.tsx index 118b6a7e..99f14366 100644 --- a/js/app/src/screens/chat-popout.native.tsx +++ b/js/app/src/screens/chat-popout.native.tsx @@ -4,6 +4,7 @@ import { ChatBox, LivestreamProvider, PlayerProvider, + StreamNotificationProvider, Text, useKeyboard, useLivestreamInfo, @@ -171,7 +172,9 @@ export function PopoutChatInner({ user }: { user: string }) { - + + + {profile && } diff --git a/js/app/src/screens/chat-popout.tsx b/js/app/src/screens/chat-popout.tsx index d6fe5a30..382da0b5 100644 --- a/js/app/src/screens/chat-popout.tsx +++ b/js/app/src/screens/chat-popout.tsx @@ -3,6 +3,7 @@ import { ChatBox, LivestreamProvider, PlayerProvider, + StreamNotificationProvider, Text, tokens, usePlayerStore, @@ -21,6 +22,7 @@ interface ChatPopoutParams { reverse?: string; hideAfter?: string; hideChatBox?: string; + showNotifications?: string; } export default function PopoutChat({ route }) { @@ -63,54 +65,69 @@ export function PopoutChatInner({ params }: { params: ChatPopoutParams }) { ? parseInt(params.hideAfter, 10) : undefined; const hideChatBox = params.hideChatBox === "true"; + const showNotifications = params.showNotifications === "true"; useEffect(() => { setSrc(params.user); }, [params.user, setSrc]); + const chat = ( + + ); + return ( - - - - {!hideChatBox && - (profile ? ( - ( - - )} - /> - ) : ( - openLoginModal({ name: "ChatPopout" })} - style={[ - zero.layout.flex.row, - zero.layout.flex.center, - zero.gap.all[4], - zero.r.xl, - { - padding: 18, - backgroundColor: "rgba(255, 255, 255, 0.1)", - maxWidth: tokens.breakpoints.sm, - }, - ]} - > - Log in or sign up to chat - - - ))} - + + {showNotifications ? ( + + {chat} + + ) : ( + chat + )} + {!hideChatBox && + (profile ? ( + ( + + )} + /> + ) : ( + openLoginModal({ name: "ChatPopout" })} + style={[ + zero.layout.flex.row, + zero.layout.flex.center, + zero.gap.all[4], + zero.r.xl, + { + padding: 18, + backgroundColor: "rgba(255, 255, 255, 0.1)", + maxWidth: tokens.breakpoints.sm, + }, + ]} + > + Log in or sign up to chat + + + ))} ); } -- 2.51.2 From c214920b290580951d0285d4e289dd6889d00e8c Mon Sep 17 00:00:00 2001 From: "Natalie B." <22222885+espeon@users.noreply.github.com> Date: Mon, 23 Mar 2026 18:49:31 -0500 Subject: [PATCH 7/7] iso time --- pkg/model/pinned_record.go | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/pkg/model/pinned_record.go b/pkg/model/pinned_record.go index 67152709..c047ee88 100644 --- a/pkg/model/pinned_record.go +++ b/pkg/model/pinned_record.go @@ -5,6 +5,7 @@ import ( "errors" "time" + "github.com/bluesky-social/indigo/util" "gorm.io/gorm" "stream.place/streamplace/pkg/streamplace" ) @@ -25,10 +26,10 @@ func (p *PinnedRecord) ToStreamplacePinnedRecord() (*streamplace.ChatPinnedRecor rec := &streamplace.ChatPinnedRecord{ LexiconTypeID: "place.stream.chat.pinnedRecord", PinnedMessage: p.PinnedMessage, - CreatedAt: p.CreatedAt.UTC().Format(time.RFC3339), + CreatedAt: p.CreatedAt.UTC().Format(util.ISO8601), } if p.ExpiresAt != nil { - s := p.ExpiresAt.UTC().Format(time.RFC3339) + s := p.ExpiresAt.UTC().Format(util.ISO8601) rec.ExpiresAt = &s } return rec, nil @@ -38,10 +39,10 @@ func (p *PinnedRecord) ToStreamplacePinnedRecordView() (*streamplace.ChatDefs_Pi pr := &streamplace.ChatPinnedRecord{ LexiconTypeID: "place.stream.chat.pinnedRecord", PinnedMessage: p.PinnedMessage, - CreatedAt: p.CreatedAt.UTC().Format(time.RFC3339), + CreatedAt: p.CreatedAt.UTC().Format(util.ISO8601), } if p.ExpiresAt != nil { - s := p.ExpiresAt.UTC().Format(time.RFC3339) + s := p.ExpiresAt.UTC().Format(util.ISO8601) pr.ExpiresAt = &s } rec := &streamplace.ChatDefs_PinnedRecordView{