Compare commits

...
23 changed files with 348 additions and 371 deletions
+10 -4
View File
@@ -365,10 +365,16 @@ FLUXER_DISCOVERY_ENABLED=true
#FLUXER_ERLANG_SCHEDULERS_MAX=16
# In-flight request ceiling for the four services Compose forwards it to: the
# users and messages routers and their shards. The Rust built-in defaults are 192
# for messages, 320 for snowflakes and 64 elsewhere, and they govern every service
# Compose does not forward this to.
#FLUXER_SVC_MAX_CONCURRENT_REQUESTS=20
# users and messages routers and their shards. Leave it unset and each service
# uses its own built-in default, which is what the numbers below describe. Set it
# and the one value replaces the built-in default on all four, so size it for the
# busiest of them rather than for the smallest. The built-in defaults are 192 for
# messages, 320 for snowflakes and 64 elsewhere, and they govern every service
# Compose does not forward this to. A router holds a slot for the whole round
# trip to its shard, so this is a ceiling on requests in flight at once and not a
# rate: too low a value does not slow requests down, it rejects them, and the api
# turns that rejection into a 503.
#FLUXER_SVC_MAX_CONCURRENT_REQUESTS=192
# The api and the Rust services name their fixed Postgres statement shapes so the
# server can reuse their plans. Named prepared statements require a session that
+4 -4
View File
@@ -565,7 +565,7 @@ services:
<<: *fluxer-env
FLUXER_SVC_NAME: users
FLUXER_SVC_MODE: router
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
healthcheck: *fluxer-svc-healthcheck
depends_on:
nats: {condition: service_healthy}
@@ -583,7 +583,7 @@ services:
FLUXER_SVC_MODE: shard
FLUXER_SVC_SHARD_ID: "0"
FLUXER_POSTGRES_MAX_CONNECTIONS: "20"
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
healthcheck: *fluxer-svc-healthcheck
depends_on:
nats: {condition: service_healthy}
@@ -633,7 +633,7 @@ services:
<<: *fluxer-env
FLUXER_SVC_NAME: messages
FLUXER_SVC_MODE: router
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
healthcheck: *fluxer-svc-healthcheck
depends_on:
nats: {condition: service_healthy}
@@ -651,7 +651,7 @@ services:
FLUXER_SVC_MODE: shard
FLUXER_SVC_SHARD_ID: "0"
FLUXER_POSTGRES_MAX_CONNECTIONS: "20"
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
healthcheck: *fluxer-svc-healthcheck
depends_on:
nats: {condition: service_healthy}
@@ -229,14 +229,11 @@ export class AdminMessageService {
hitsPerPage: limit,
page: 1,
});
const messageEntries = result.hits.map((hit) => ({
channelId: createChannelID(BigInt(hit.channelId)),
messageId: createMessageID(BigInt(hit.id)),
}));
const resolvedMessages = await Promise.all(
messageEntries.map(({channelId, messageId}) => this.getMessageResponseForAdmin(channelId, messageId)),
);
const messageResponses = resolvedMessages.filter((message): message is MessageResponse => message !== null);
const messageResponses = await createMessageResponseDataService().buildMessages({
userId: createUserID(0n),
messages: result.messages,
access: await this.getMessageResponseAccessForAdmin(channelId),
});
const attachmentStatuses = await this.getAttachmentStatusesForMessages(messageResponses);
const priorReports = await this.getPriorReportsForMessages(messageResponses);
const adminMessages = messageResponses.map((message) =>
@@ -282,19 +279,6 @@ export class AdminMessageService {
});
}
private async getMessageResponseForAdmin(
channelId: ChannelID,
messageId: MessageID,
): Promise<MessageResponse | null> {
const access = await this.getMessageResponseAccessForAdmin(channelId);
return createMessageResponseDataService().getMessage({
userId: createUserID(0n),
channelId,
messageId,
access,
});
}
private async getPriorReportsForMessages(messages: Array<MessageResponse>): Promise<Map<string, Array<string>>> {
const authorIds = messages.map((message) => createUserID(BigInt(message.author.id)));
return this.deps.ncmecSubmissionService.getUserPriorReportIds(authorIds);
@@ -4,11 +4,10 @@ import {ChannelTypes, Permissions} from '@fluxer/constants/src/ChannelConstants'
import {UnknownMessageError} from '@fluxer/errors/src/domains/channel/UnknownMessageError';
import {FeatureTemporarilyDisabledError} from '@fluxer/errors/src/domains/core/FeatureTemporarilyDisabledError';
import type {MessageSearchRequest} from '@fluxer/schema/src/domains/message/MessageRequestSchemas';
import type {MessageResponse, MessageSearchResponse} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
import type {MessageSearchResponse} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
import {snowflakeToDate} from '@fluxer/snowflake/src/Snowflake';
import {AttachmentDecayService} from '../../../attachment/AttachmentDecayService';
import type {AttachmentID, ChannelID, MessageID, UserID} from '../../../BrandedTypes';
import {createChannelID, createMessageID} from '../../../BrandedTypes';
import type {UserCacheService} from '../../../infrastructure/UserCacheService';
import type {RequestCache} from '../../../middleware/RequestCacheMiddleware';
import type {Channel} from '../../../models/Channel';
@@ -192,29 +191,19 @@ export class MessageRetrievalService {
hitsPerPage,
page,
});
const messageEntries = result.hits.map((hit) => ({
channelId: createChannelID(BigInt(hit.channelId)),
messageId: createMessageID(BigInt(hit.id)),
}));
const access = {
sourceGuildId: channel.guildId,
messageHistoryCutoff: !hasReadHistory ? (authChannel.guild?.message_history_cutoff ?? null) : null,
canReadMessageHistory: hasReadHistory,
};
const responseDataService = createMessageResponseDataService();
const foundMessages = await Promise.all(
messageEntries.map(({channelId, messageId}) =>
responseDataService.getMessage({
userId,
channelId,
messageId,
access,
}),
),
const builtMessages = await createMessageResponseDataService().buildMessages({
userId,
messages: result.messages,
access,
});
const messageResponses = builtMessages.map(
({referenced_message: _referencedMessage, ...searchMessage}) => searchMessage,
);
const messageResponses = foundMessages
.filter((message): message is MessageResponse => message !== null)
.map(({referenced_message: _referencedMessage, ...searchMessage}) => searchMessage);
return {
channels: messageResponses.length > 0 ? [await this.mapSearchChannelResponse(channel, userId, requestCache)] : [],
messages: messageResponses,
@@ -156,7 +156,7 @@ export class GuildSearchService {
page,
cursor,
});
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result, userId, requestCache);
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result.messages, userId, requestCache);
return {
messages: mappedResponses.messages,
channels: mappedResponses.channels,
@@ -238,7 +238,7 @@ export class GuildSearchService {
page,
cursor,
});
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result, userId, requestCache);
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result.messages, userId, requestCache);
return {
messages: mappedResponses.messages,
channels: mappedResponses.channels,
@@ -253,7 +253,7 @@ export class GlobalSearchService {
page,
cursor,
});
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result, userId, requestCache);
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result.messages, userId, requestCache);
return {
messages: mappedResponses.messages,
channels: mappedResponses.channels,
@@ -1,20 +1,18 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {
MessageResponse,
MessageSearchResultsResponse,
} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
import type {MessageSearchResultsResponse} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
import type {UserID} from '../BrandedTypes';
import {createChannelID, createMessageID} from '../BrandedTypes';
import {createChannelID} from '../BrandedTypes';
import {mapChannelToResponse} from '../channel/ChannelMappers';
import type {IChannelRepository} from '../channel/IChannelRepository';
import {
createMessageResponseDataService,
messageResponseAccessForChannel,
} from '../channel/services/message/MessageResponseDataService';
import {createMessageResponseDataService} from '../channel/services/message/MessageResponseDataService';
import type {UserCacheService} from '../infrastructure/UserCacheService';
import type {RequestCache} from '../middleware/RequestCacheMiddleware';
import type {Channel} from '../models/Channel';
import type {Message} from '../models/Message';
import {mapWithConcurrency} from '../utils/ConcurrencyUtils';
const CHANNEL_LOOKUP_CONCURRENCY = 16;
export class MessageSearchResponseMapper {
constructor(
@@ -23,71 +21,46 @@ export class MessageSearchResponseMapper {
) {}
async mapSearchResultToResponses(
result: {
hits: Array<{
channelId: string;
id: string;
}>;
total: number;
},
messages: Array<Message>,
userId: UserID,
requestCache: RequestCache,
): Promise<{
messages: Array<MessageSearchResultsResponse['messages'][number]>;
channels: Array<MessageSearchResultsResponse['channels'][number]>;
}> {
const messageEntries = result.hits.map((hit) => ({
channelId: createChannelID(BigInt(hit.channelId)),
messageId: createMessageID(BigInt(hit.id)),
}));
const orderedChannelIds = new Set<string>();
for (const entry of messageEntries) {
orderedChannelIds.add(entry.channelId.toString());
}
const channels = await Promise.all(
Array.from(orderedChannelIds).map((channelId) =>
this.channelRepository.findUnique(createChannelID(BigInt(channelId))),
),
const orderedChannelIds = Array.from(new Set(messages.map((message) => message.channelId.toString())));
const channels = await mapWithConcurrency(orderedChannelIds, CHANNEL_LOOKUP_CONCURRENCY, (channelId) =>
this.channelRepository.findUnique(createChannelID(BigInt(channelId))),
);
const validChannels = channels.filter((channel): channel is Channel => channel !== null);
const channelById = new Map(validChannels.map((channel) => [channel.id.toString(), channel] as const));
const responseDataService = createMessageResponseDataService();
const messageResponsesWithEntries = await Promise.all(
messageEntries.map(async (entry) => {
const channel = channelById.get(entry.channelId.toString());
if (!channel) return null;
const message = await responseDataService.getMessage({
userId,
channelId: entry.channelId,
messageId: entry.messageId,
access: messageResponseAccessForChannel(channel),
});
return message ? {message, channelId: entry.channelId.toString()} : null;
}),
const channelById = new Map(
channels
.filter((channel): channel is Channel => channel !== null)
.map((channel) => [channel.id.toString(), channel] as const),
);
const validMessageResponseEntries = messageResponsesWithEntries.filter(
(entry): entry is {message: MessageResponse; channelId: string} => entry !== null,
);
const messageResponses = validMessageResponseEntries.map((entry) => {
const {referenced_message: _referencedMessage, ...searchMessage} = entry.message;
return searchMessage;
const renderableMessages = messages.filter((message) => channelById.has(message.channelId.toString()));
const messageResponses = await createMessageResponseDataService().buildMessagesForChannels({
userId,
messages: renderableMessages,
channelById,
});
const orderedResponseChannelIds = new Set(validMessageResponseEntries.map((entry) => entry.channelId));
const orderedChannels = Array.from(orderedResponseChannelIds)
const searchMessages = messageResponses.map(
({referenced_message: _referencedMessage, ...searchMessage}) => searchMessage,
);
const respondedChannelIds = new Set(searchMessages.map((message) => message.channel_id));
const orderedChannels = orderedChannelIds
.filter((channelId) => respondedChannelIds.has(channelId))
.map((channelId) => channelById.get(channelId))
.filter((channel): channel is Channel => channel !== undefined);
const channelResponses = await Promise.all(
orderedChannels.map((channel) =>
mapChannelToResponse({
channel,
currentUserId: userId,
userCacheService: this.userCacheService,
requestCache,
}),
),
const channelResponses = await mapWithConcurrency(orderedChannels, CHANNEL_LOOKUP_CONCURRENCY, (channel) =>
mapChannelToResponse({
channel,
currentUserId: userId,
userCacheService: this.userCacheService,
requestCache,
}),
);
return {
messages: messageResponses,
messages: searchMessages,
channels: channelResponses,
};
}
@@ -6,6 +6,7 @@ import {type ChannelID, createChannelID, createMessageID, type MessageID} from '
import type {IMessageRepository} from '../channel/repositories/IMessageRepository';
import {Logger} from '../Logger';
import type {Message} from '../models/Message';
import {mapWithConcurrency} from '../utils/ConcurrencyUtils';
import type {IMessageSearchService} from './IMessageSearchService';
import {deleteMessageSearchDocuments} from './MessageSearchIndexCleanup';
@@ -13,11 +14,16 @@ const RECONCILE_BATCH_SIZE = 250;
const MAX_RECONCILE_PAGES = 40;
const MAX_STALE_DELETE_ABSOLUTE = 250;
const MAX_STALE_DELETE_RATIO = 0.5;
const HIT_LOOKUP_CONCURRENCY = 32;
interface MessageLookupRepository {
readonly messages: Pick<IMessageRepository, 'getMessage'>;
}
interface MessageSearchLookupResult extends SearchResult<SearchableMessage> {
messages: Array<Message>;
}
interface SearchExistingMessagesParams {
searchService: IMessageSearchService;
messageRepository: MessageLookupRepository;
@@ -28,12 +34,21 @@ interface SearchExistingMessagesParams {
cursor?: Array<string>;
}
interface ValidatedHit {
hit: SearchableMessage;
message: Message | null;
}
interface ValidatedHits {
validHits: Array<SearchableMessage>;
validHits: Array<ValidatedHit>;
staleMessageIds: Array<MessageID>;
lookupErrorCount: number;
}
function resolvedMessages(validHits: Array<ValidatedHit>): Array<Message> {
return validHits.map((entry) => entry.message).filter((message): message is Message => message !== null);
}
export async function searchExistingMessages({
searchService,
messageRepository,
@@ -42,7 +57,7 @@ export async function searchExistingMessages({
hitsPerPage,
page,
cursor,
}: SearchExistingMessagesParams): Promise<SearchResult<SearchableMessage>> {
}: SearchExistingMessagesParams): Promise<MessageSearchLookupResult> {
const result = await searchService.searchMessages(query, filters, {
hitsPerPage,
page: cursor?.length ? undefined : page,
@@ -50,7 +65,7 @@ export async function searchExistingMessages({
});
const validated = await validateSearchHits(messageRepository, result.hits);
if (validated.staleMessageIds.length === 0) {
return result;
return {...result, messages: resolvedMessages(validated.validHits)};
}
if (cursor?.length) {
if (validated.lookupErrorCount === 0) {
@@ -58,8 +73,9 @@ export async function searchExistingMessages({
}
return {
...result,
hits: validated.validHits,
hits: validated.validHits.map((entry) => entry.hit),
total: Math.max(validated.validHits.length, result.total - validated.staleMessageIds.length),
messages: resolvedMessages(validated.validHits),
};
}
return reconcileOffsetSearchResult({
@@ -79,9 +95,9 @@ async function reconcileOffsetSearchResult({
filters,
hitsPerPage,
page,
}: Omit<SearchExistingMessagesParams, 'cursor'>): Promise<SearchResult<SearchableMessage>> {
}: Omit<SearchExistingMessagesParams, 'cursor'>): Promise<MessageSearchLookupResult> {
const requestedOffset = (page - 1) * hitsPerPage;
const pageHits: Array<SearchableMessage> = [];
const pageHits: Array<ValidatedHit> = [];
const staleMessageIds: Array<MessageID> = [];
let lookupErrorCount = 0;
let examinedCount = 0;
@@ -102,9 +118,9 @@ async function reconcileOffsetSearchResult({
lookupErrorCount += validated.lookupErrorCount;
examinedCount += result.hits.length;
staleMessageIds.push(...validated.staleMessageIds);
for (const hit of validated.validHits) {
for (const entry of validated.validHits) {
if (validTotal >= requestedOffset && pageHits.length < hitsPerPage) {
pageHits.push(hit);
pageHits.push(entry);
}
validTotal += 1;
}
@@ -121,8 +137,9 @@ async function reconcileOffsetSearchResult({
await deleteStaleSearchDocuments(searchService, staleMessageIds, examinedCount);
}
return {
hits: pageHits,
hits: pageHits.map((entry) => entry.hit),
total: Math.max(pageHits.length, corpusTotal - staleMessageIds.length),
messages: resolvedMessages(pageHits),
};
}
@@ -130,38 +147,36 @@ async function validateSearchHits(
messageRepository: MessageLookupRepository,
hits: Array<SearchableMessage>,
): Promise<ValidatedHits> {
const checked = await Promise.all(
hits.map(async (hit) => {
let channelId: ChannelID;
let messageId: MessageID;
try {
channelId = createChannelID(BigInt(hit.channelId));
messageId = createMessageID(BigInt(hit.id));
} catch (_invalidId) {
return {hit: null, staleMessageId: null, lookupError: false};
}
let message: Message | null;
try {
message = await messageRepository.messages.getMessage(channelId, messageId);
} catch (error) {
Logger.warn(
{error, messageId: hit.id, channelId: hit.channelId},
'Search read repair lookup failed; keeping document',
);
return {hit, staleMessageId: null, lookupError: true};
}
if (message && message.channelId.toString() === hit.channelId) {
return {hit, staleMessageId: null, lookupError: false};
}
return {hit: null, staleMessageId: messageId, lookupError: false};
}),
);
const validHits: Array<SearchableMessage> = [];
const checked = await mapWithConcurrency(hits, HIT_LOOKUP_CONCURRENCY, async (hit) => {
let channelId: ChannelID;
let messageId: MessageID;
try {
channelId = createChannelID(BigInt(hit.channelId));
messageId = createMessageID(BigInt(hit.id));
} catch (_invalidId) {
return {entry: null, staleMessageId: null, lookupError: false};
}
let message: Message | null;
try {
message = await messageRepository.messages.getMessage(channelId, messageId);
} catch (error) {
Logger.warn(
{error, messageId: hit.id, channelId: hit.channelId},
'Search read repair lookup failed; keeping document',
);
return {entry: {hit, message: null}, staleMessageId: null, lookupError: true};
}
if (message && message.channelId.toString() === hit.channelId) {
return {entry: {hit, message}, staleMessageId: null, lookupError: false};
}
return {entry: null, staleMessageId: messageId, lookupError: false};
});
const validHits: Array<ValidatedHit> = [];
const staleMessageIds: Array<MessageID> = [];
let lookupErrorCount = 0;
for (const item of checked) {
if (item.hit) {
validHits.push(item.hit);
if (item.entry) {
validHits.push(item.entry);
}
if (item.staleMessageId) {
staleMessageIds.push(item.staleMessageId);
@@ -125,12 +125,13 @@ describe('User Settings synced_preferences', () => {
});
test('reports TOO_LARGE for a snapshot inside the encoded-length bound but over the byte cap', async () => {
const account = await createTestAccount(harness);
const baseEntries = Math.floor((SYNCED_PREFERENCES_MAX_BYTES - 8192) / 1006);
const build = (tailLength: number) =>
encodeSyncedPreferences(
create(SyncedPreferencesSchema, {
localSpamOverrides: create(LocalUserSpamOverridesSchema, {
spammerUserIds: [
...Array.from({length: 260}, (_, i) => `${i}-${'x'.repeat(1000)}`),
...Array.from({length: baseEntries}, (_, i) => `${i}-${'x'.repeat(1000)}`),
'y'.repeat(tailLength),
],
}),
@@ -138,7 +139,13 @@ describe('User Settings synced_preferences', () => {
);
let tail = 1;
let encoded = build(tail);
while (encodedSyncedPreferencesByteLength(encoded) < SYNCED_PREFERENCES_MAX_BYTES + 1 && tail < 8000) {
while (encodedSyncedPreferencesByteLength(encoded) <= SYNCED_PREFERENCES_MAX_BYTES) {
tail += 512;
encoded = build(tail);
}
tail = Math.max(1, tail - 512);
encoded = build(tail);
while (encodedSyncedPreferencesByteLength(encoded) <= SYNCED_PREFERENCES_MAX_BYTES) {
tail += 1;
encoded = build(tail);
}
@@ -5,6 +5,7 @@ import {YOUTUBE_PROVIDER_NAME} from '@app/features/app/config/I18nDisplayConstan
import styles from '@app/features/channel/components/embeds/media/EmbedYouTube.module.css';
import {OverlayActionButton, OverlayPlayButton} from '@app/features/channel/components/embeds/media/MediaButtons';
import {useNearViewport} from '@app/features/messaging/hooks/useNearViewport';
import ActiveIframeEmbed from '@app/features/messaging/state/ActiveIframeEmbed';
import {openExternalUrlWithWarning} from '@app/features/messaging/utils/ExternalLinkUtils';
import * as ImageCacheUtils from '@app/features/messaging/utils/ImageCacheUtils';
import {
@@ -21,7 +22,7 @@ import {useLingui} from '@lingui/react/macro';
import {ArrowSquareOutIcon, PlayIcon} from '@phosphor-icons/react';
import {AnimatePresence, motion} from 'framer-motion';
import {observer} from 'mobx-react-lite';
import {type FC, useCallback, useEffect, useMemo, useState} from 'react';
import {type FC, useCallback, useEffect, useId, useMemo, useState} from 'react';
const PLAY_VIDEO_DESCRIPTOR = msg({
message: 'Play video',
@@ -147,7 +148,8 @@ const Thumbnail: FC<ThumbnailProps> = observer(
);
export const EmbedYouTube: FC<EmbedYouTubeProps> = observer(({embed, width = YOUTUBE_CONFIG.DEFAULT_WIDTH}) => {
const {i18n} = useLingui();
const [hasInteracted, setHasInteracted] = useState(false);
const embedId = useId();
const hasInteracted = ActiveIframeEmbed.isActive(embedId);
const posterSrc = embed.thumbnail?.proxy_url || '';
const {ref: visibilityRef, isNearViewport} = useNearViewport<HTMLDivElement>({rememberKey: posterSrc});
const [posterCacheAtMount] = useState(() => ({src: posterSrc, cached: ImageCacheUtils.hasImage(posterSrc)}));
@@ -171,10 +173,14 @@ export const EmbedYouTube: FC<EmbedYouTubeProps> = observer(({embed, width = YOU
);
return cleanup;
}, [isNearViewport, loadedPosterSrc, posterSrc]);
const handleInitialPlay = useCallback((event: React.MouseEvent | React.KeyboardEvent) => {
event.stopPropagation();
setHasInteracted(true);
}, []);
const handleInitialPlay = useCallback(
(event: React.MouseEvent | React.KeyboardEvent) => {
event.stopPropagation();
ActiveIframeEmbed.claim(embedId);
},
[embedId],
);
useEffect(() => () => ActiveIframeEmbed.release(embedId), [embedId]);
const handleOpenInNewTab = useCallback(
(event: React.MouseEvent | React.KeyboardEvent) => {
event.stopPropagation();
@@ -250,6 +256,7 @@ export const EmbedYouTube: FC<EmbedYouTubeProps> = observer(({embed, width = YOU
sandbox="allow-forms allow-modals allow-popups allow-popups-to-escape-sandbox allow-same-origin allow-scripts"
src={embedVideoUrl.toString()}
className={styles.iframe}
scrolling="no"
data-embed-media="true"
aria-label={embed.title || i18n._(VIDEO_DESCRIPTOR, {youtubeProviderName: YOUTUBE_PROVIDER_NAME})}
data-flx="channel.embeds.media.embed-you-tube.iframe"
@@ -1,50 +0,0 @@
/* SPDX-License-Identifier: AGPL-3.0-or-later */
.body {
display: flex;
flex-direction: column;
gap: 0.75rem;
}
.description {
color: var(--text-secondary);
margin: 0;
font-size: 0.875rem;
line-height: 1.5;
}
.list {
display: flex;
flex-direction: column;
gap: 0.5rem;
margin: 0;
padding: 0;
list-style: none;
}
.listItem {
color: var(--text-secondary);
font-size: 0.8125rem;
line-height: 1.4;
padding-inline-start: 0.75rem;
position: relative;
}
.listItem::before {
content: '•';
position: absolute;
left: 0;
color: var(--text-chat-muted);
}
.listItem strong {
color: var(--text-primary);
font-weight: 600;
}
.hint {
color: var(--text-chat-muted);
font-size: 0.75rem;
margin: 0;
line-height: 1.4;
}
@@ -1,116 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import * as Modal from '@app/features/app/components/dialogs/Modal';
import styles from '@app/features/channel/components/pickers/gif/FavoriteGifFirstTimePromptModal.module.css';
import FavoriteGif from '@app/features/expressions/state/FavoriteGif';
import {Button} from '@app/features/ui/button/Button';
import * as ModalCommands from '@app/features/ui/commands/ModalCommands';
import {formatUserSettingsPath} from '@app/features/user/components/settings_utils/SettingsConstants';
import {msg} from '@lingui/core/macro';
import {Trans, useLingui} from '@lingui/react/macro';
import {observer} from 'mobx-react-lite';
import {useRef} from 'react';
const HOW_SHOULD_WE_SAVE_YOUR_GIF_FAVORITES_DESCRIPTOR = msg({
message: 'How should we save your GIF favorites?',
comment: 'Confirmation prompt in the channel and chat favorite gif first time prompt modal.',
});
interface FavoriteGifFirstTimePromptModalProps {
onConfirm: () => void;
}
export const FavoriteGifFirstTimePromptModal = observer(function FavoriteGifFirstTimePromptModal({
onConfirm,
}: FavoriteGifFirstTimePromptModalProps) {
const {i18n} = useLingui();
const initialFocusRef = useRef<HTMLButtonElement | null>(null);
const mediaSettingsPath = formatUserSettingsPath(i18n, 'chat_settings', 'media');
const handleConfirm = () => {
FavoriteGif.setSaveGifFavoritesAsSavedMedia(false);
FavoriteGif.markFirstTimePromptSeen();
ModalCommands.pop();
onConfirm();
};
const handleUseSavedMedia = () => {
FavoriteGif.setSaveGifFavoritesAsSavedMedia(true);
FavoriteGif.markFirstTimePromptSeen();
ModalCommands.pop();
onConfirm();
};
const handleCancel = () => {
ModalCommands.pop();
};
return (
<Modal.Root
size="small"
centered
initialFocusRef={initialFocusRef}
data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.modal-root"
>
<Modal.Header
title={i18n._(HOW_SHOULD_WE_SAVE_YOUR_GIF_FAVORITES_DESCRIPTOR)}
onClose={handleCancel}
data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.modal-header"
/>
<Modal.Content data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.modal-content">
<div className={styles.body} data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.body">
<p
className={styles.description}
data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.description"
>
<Trans>
You can store starred GIFs as URL-only favorites or upload them to your saved media. Pick the one that
fits how you use them. You can change it any time in {mediaSettingsPath}.
</Trans>
</p>
<ul className={styles.list} data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.list">
<li
className={styles.listItem}
data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.list-item"
>
<Trans>
<strong data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.strong">
URL-only favorites (default)
</strong>
: synced across your devices, no upload, doesn't count against saved media. The original media may
disappear if its host removes it.
</Trans>
</li>
<li
className={styles.listItem}
data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.list-item--2"
>
<Trans>
<strong data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.strong--2">
Saved media
</strong>
: uploaded, taggable, searchable, and persistent, but counts against your saved media limit.
</Trans>
</li>
</ul>
<p className={styles.hint} data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.hint">
<Trans>We'll only ask once.</Trans>
</p>
</div>
</Modal.Content>
<Modal.Footer data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.modal-footer">
<Button
variant="secondary"
onClick={handleUseSavedMedia}
data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.button.use-saved-media"
>
<Trans>Use saved media</Trans>
</Button>
<Button
variant="primary"
onClick={handleConfirm}
ref={initialFocusRef}
data-flx="channel.pickers.gif.favorite-gif-first-time-prompt-modal.button.confirm"
>
<Trans>Use URL-only (recommended)</Trans>
</Button>
</Modal.Footer>
</Modal.Root>
);
});
@@ -0,0 +1,60 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {describe, expect, it} from 'vitest';
import {type FavoriteGifEntry, slimFavoriteGifEntry} from './FavoriteGifTypes';
function format(src: string, width: number) {
return {src, proxy_src: `https://media.test/external/sig/${src}`, width, height: Math.round(width * 0.84)};
}
const FAT: FavoriteGifEntry = {
url: 'https://klipy.com/gifs/reaction',
proxy_url: 'https://media.test/external/sig/tiny.webm',
width: 220,
height: 185,
media: {
nanowebm: format('nano.webm', 90),
tinywebm: format('tiny.webm', 220),
mediumwebm: format('medium.webm', 320),
webm: format('full.webm', 498),
gif: format('full.gif', 498),
},
content_type: 'video/webm',
placeholder: 'placeholder',
};
describe('slimFavoriteGifEntry', () => {
it('drops the media map so a favorite is cheap to sync', () => {
expect(slimFavoriteGifEntry(FAT).media).toEqual({});
});
it('promotes a preview that covers a 200px tile at 2x into the top-level fields', () => {
const slim = slimFavoriteGifEntry(FAT);
expect(slim.proxy_url).toBe(FAT.media.webm.proxy_src);
expect(slim.width).toBe(498);
expect(slim.content_type).toBe('video/webm');
});
it('keeps url and placeholder untouched', () => {
const slim = slimFavoriteGifEntry(FAT);
expect(slim.url).toBe(FAT.url);
expect(slim.placeholder).toBe(FAT.placeholder);
});
it('is idempotent so the synced roundtrip stays stable', () => {
const once = slimFavoriteGifEntry(FAT);
expect(slimFavoriteGifEntry(once)).toEqual(once);
});
it('leaves an entry that already carries no media alone', () => {
const urlOnly: FavoriteGifEntry = {...FAT, media: {}};
expect(slimFavoriteGifEntry(urlOnly)).toBe(urlOnly);
});
it('does not invent a preview when every format is unusable', () => {
const broken: FavoriteGifEntry = {...FAT, media: {webm: {src: '', proxy_src: '', width: 0, height: 0}}};
const slim = slimFavoriteGifEntry(broken);
expect(slim.media).toEqual({});
expect(slim.proxy_url).toBe(FAT.proxy_url);
});
});
@@ -105,4 +105,24 @@ export function pickCanonicalPreviewFormat(
return null;
}
export const STORED_PREVIEW_DEVICE_PIXEL_RATIO = 2;
export function slimFavoriteGifEntry(entry: FavoriteGifEntry): FavoriteGifEntry {
const best = pickBestPreviewFormat(entry.media, 'any', {
cssWidth: PREVIEW_TILE_CSS_WIDTH,
devicePixelRatio: STORED_PREVIEW_DEVICE_PIXEL_RATIO,
});
if (best == null) {
return Object.keys(entry.media).length === 0 ? entry : {...entry, media: {}};
}
return {
...entry,
proxy_url: best.format.proxy_src,
width: best.format.width,
height: best.format.height,
content_type: inferFormatContentType(best.key),
media: {},
};
}
export {inferFormatContentType};
@@ -3,7 +3,6 @@
import RuntimeConfig from '@app/features/app/state/RuntimeConfig';
import styles from '@app/features/channel/components/GifPicker.module.css';
import {safePause, safePlay, useGifVideoPool} from '@app/features/channel/components/GifVideoPool';
import {FavoriteGifFirstTimePromptModal} from '@app/features/channel/components/pickers/gif/FavoriteGifFirstTimePromptModal';
import type {GifPickerGridItemData} from '@app/features/channel/components/pickers/gif/GifPickerTypes';
import {PickerThumbnail} from '@app/features/channel/components/pickers/shared/PickerThumbnail';
import {usePooledVideo} from '@app/features/channel/components/pickers/shared/usePooledVideo';
@@ -24,7 +23,6 @@ import {isKeyboardActivationKey} from '@app/features/input/utils/KeyboardUtils';
import {decodeThumbHashDataURL} from '@app/features/messaging/utils/ThumbHashUtils';
import {ComponentBus} from '@app/features/platform/utils/ComponentBus';
import {remFromPx} from '@app/features/theme/layout/RemFromPx';
import {modal, push} from '@app/features/ui/commands/ModalCommands';
import FocusRing from '@app/features/ui/focus_ring/FocusRing';
import {Tooltip} from '@app/features/ui/tooltip/Tooltip';
import {msg} from '@lingui/core/macro';
@@ -400,17 +398,6 @@ export const GifPickerGridItem = observer(function GifPickerGridItem({
const handleFavoriteClick = (e: React.MouseEvent) => {
e.stopPropagation();
if (isFavoritePending) return;
if (!FavoriteGif.hasSeenFavoriteGifFirstTimePrompt && !isFavorited) {
push(
modal(() => (
<FavoriteGifFirstTimePromptModal
onConfirm={() => void performFavoriteToggle()}
data-flx="channel.pickers.gif.gif-picker-grid-item.handle-favorite-click.favorite-gif-first-time-prompt-modal"
/>
)),
);
return;
}
void performFavoriteToggle();
};
const favoriteTooltipText = (() => {
@@ -1,10 +1,12 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {
FavoriteGifEntry,
FavoriteGifMediaFormat,
import {
type FavoriteGifEntry,
type FavoriteGifMediaFormat,
slimFavoriteGifEntry,
} from '@app/features/channel/components/pickers/gif/FavoriteGifTypes';
import {makeSyncedField} from '@app/features/user/state/SyncedField';
import {FAVORITE_GIF_MAX_ENCODED_BYTES} from '@app/features/user/state/SyncedFieldBudget';
import type {FavoriteGifMediaFormat as FavoriteGifMediaFormatProto} from '@fluxer/schema/src/gen/fluxer/user/preferences/v1/pickers_pb';
import {FavoriteGifSettingsSchema} from '@fluxer/schema/src/gen/fluxer/user/preferences/v1/pickers_pb';
import {makeAutoObservable} from 'mobx';
@@ -52,8 +54,9 @@ class FavoriteGif {
field: 'favoriteGifs',
schema: FavoriteGifSettingsSchema,
persist: ['favoriteGifs', 'saveGifFavoritesAsSavedMedia', 'hasSeenFavoriteGifFirstTimePrompt'],
maxEncodedBytes: FAVORITE_GIF_MAX_ENCODED_BYTES,
toMessage: (s) => ({
entries: s.favoriteGifs.map((entry) => ({
entries: s.favoriteGifs.map(slimFavoriteGifEntry).map((entry) => ({
url: entry.url,
proxyUrl: entry.proxy_url,
width: entry.width,
@@ -66,15 +69,17 @@ class FavoriteGif {
seenFirstTimePrompt: s.hasSeenFavoriteGifFirstTimePrompt,
}),
applyMessage: (s, m) => {
s.favoriteGifs = m.entries.map((entry) => ({
url: entry.url,
proxy_url: entry.proxyUrl,
width: entry.width,
height: entry.height,
media: mediaFromProto(entry.media),
content_type: entry.contentType,
placeholder: entry.placeholder ? entry.placeholder : null,
}));
s.favoriteGifs = m.entries.map((entry) =>
slimFavoriteGifEntry({
url: entry.url,
proxy_url: entry.proxyUrl,
width: entry.width,
height: entry.height,
media: mediaFromProto(entry.media),
content_type: entry.contentType,
placeholder: entry.placeholder ? entry.placeholder : null,
}),
);
s.saveGifFavoritesAsSavedMedia = m.saveAsSavedMedia;
s.hasSeenFavoriteGifFirstTimePrompt = m.seenFirstTimePrompt;
},
@@ -110,10 +115,6 @@ class FavoriteGif {
setSaveGifFavoritesAsSavedMedia(value: boolean): void {
this.saveGifFavoritesAsSavedMedia = value;
}
markFirstTimePromptSeen(): void {
this.hasSeenFavoriteGifFirstTimePrompt = true;
}
}
export default new FavoriteGif();
@@ -1,6 +1,5 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {FavoriteGifFirstTimePromptModal} from '@app/features/channel/components/pickers/gif/FavoriteGifFirstTimePromptModal';
import * as FavoriteGifCommands from '@app/features/expressions/commands/FavoriteGifCommands';
import * as FavoriteMemeCommands from '@app/features/expressions/commands/FavoriteMemeCommands';
import {AddFavoriteMemeModal} from '@app/features/expressions/components/modals/AddFavoriteMemeModal';
@@ -106,17 +105,6 @@ export function useMediaFavorite({
});
}
};
if (!FavoriteGif.hasSeenFavoriteGifFirstTimePrompt && !isFavorited) {
ModalCommands.push(
modal(() => (
<FavoriteGifFirstTimePromptModal
onConfirm={performToggle}
data-flx="messaging.use-media-favorite.toggle-favorite.favorite-gif-first-time-prompt-modal"
/>
)),
);
return;
}
performToggle();
return;
}
@@ -0,0 +1,27 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {makeAutoObservable} from 'mobx';
class ActiveIframeEmbed {
activeId: string | null = null;
constructor() {
makeAutoObservable(this, {}, {autoBind: true});
}
isActive(id: string): boolean {
return this.activeId === id;
}
claim(id: string): void {
this.activeId = id;
}
release(id: string): void {
if (this.activeId === id) {
this.activeId = null;
}
}
}
export default new ActiveIframeEmbed();
@@ -1,6 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Logger} from '@app/features/platform/utils/AppLogger';
import {DEFAULT_SYNCED_FIELD_MAX_ENCODED_BYTES, decideOversizePush} from '@app/features/user/state/SyncedFieldBudget';
import {verifyRoundtripStability} from '@app/features/user/state/SyncedFieldRoundtrip';
import {
createSyncedFieldMachineSnapshot,
@@ -59,6 +60,7 @@ export async function makeSyncedField<
}
const isEnabled = (): boolean => config.enabled?.() ?? true;
let machine = createSyncedFieldMachineSnapshot();
let lastPreparedBytes: number | null = null;
let ownerUserId: string | null = null;
let observedUserId: string | null = null;
const suspend = (reason: SyncedFieldFailureReason, message: string, error?: unknown): void => {
@@ -130,13 +132,20 @@ export async function makeSyncedField<
return;
}
const encodedBytes = toBinary(schema, candidate).length;
if (encodedBytes > maxEncodedBytes) {
const decision = decideOversizePush({encodedBytes, maxEncodedBytes, lastPreparedBytes});
lastPreparedBytes = encodedBytes;
if (decision === 'drop') {
logger.error(
`${tag}: payload of ${encodedBytes} bytes exceeds per-field budget of ${maxEncodedBytes}; push dropped.`,
);
transitionOnly({type: 'sync.localAlreadySynced'});
return;
}
if (decision === 'push-shrinks') {
logger.warn(
`${tag}: payload of ${encodedBytes} bytes is over the per-field budget of ${maxEncodedBytes} but smaller than the previous candidate. Pushing so removals persist.`,
);
}
const stability = verifyRoundtripStability<T, M>({
schema,
store,
@@ -204,7 +213,7 @@ export async function makeSyncedField<
processCommands();
}
await Promise.resolve();
const maxEncodedBytes = config.maxEncodedBytes ?? 65_536;
const maxEncodedBytes = config.maxEncodedBytes ?? DEFAULT_SYNCED_FIELD_MAX_ENCODED_BYTES;
const UserSettings = (await import('@app/features/user/state/UserSettings')).default;
const SessionManager = (await import('@app/features/platform/state/AuthSession')).default;
const applyDefaultsForUserChange = (userId: string): void => {
@@ -0,0 +1,45 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {describe, expect, it} from 'vitest';
import {decideOversizePush} from './SyncedFieldBudget';
const MAX = 1000;
describe('decideOversizePush', () => {
it('accepts a payload inside the budget', () => {
expect(decideOversizePush({encodedBytes: 999, maxEncodedBytes: MAX, lastPreparedBytes: null})).toBe(
'within-budget',
);
expect(decideOversizePush({encodedBytes: MAX, maxEncodedBytes: MAX, lastPreparedBytes: null})).toBe(
'within-budget',
);
});
it('drops the first over-budget payload', () => {
expect(decideOversizePush({encodedBytes: 1001, maxEncodedBytes: MAX, lastPreparedBytes: null})).toBe('drop');
});
it('drops an over-budget payload that grows', () => {
expect(decideOversizePush({encodedBytes: 1200, maxEncodedBytes: MAX, lastPreparedBytes: 1100})).toBe('drop');
});
it('drops an over-budget payload that stays the same size', () => {
expect(decideOversizePush({encodedBytes: 1100, maxEncodedBytes: MAX, lastPreparedBytes: 1100})).toBe('drop');
});
it('pushes an over-budget payload that shrinks so removals persist', () => {
expect(decideOversizePush({encodedBytes: 1050, maxEncodedBytes: MAX, lastPreparedBytes: 1100})).toBe(
'push-shrinks',
);
});
it('lets a run of removals drain back under the budget', () => {
let last: number | null = null;
const decisions: Array<string> = [];
for (const bytes of [1300, 1200, 1100, 900]) {
decisions.push(decideOversizePush({encodedBytes: bytes, maxEncodedBytes: MAX, lastPreparedBytes: last}));
last = bytes;
}
expect(decisions).toEqual(['drop', 'push-shrinks', 'push-shrinks', 'within-budget']);
});
});
@@ -0,0 +1,17 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
export const DEFAULT_SYNCED_FIELD_MAX_ENCODED_BYTES = 65_536;
export const FAVORITE_GIF_MAX_ENCODED_BYTES = 384 * 1024;
export type OversizePushDecision = 'within-budget' | 'push-shrinks' | 'drop';
export function decideOversizePush(args: {
encodedBytes: number;
maxEncodedBytes: number;
lastPreparedBytes: number | null;
}): OversizePushDecision {
if (args.encodedBytes <= args.maxEncodedBytes) return 'within-budget';
if (args.lastPreparedBytes != null && args.encodedBytes < args.lastPreparedBytes) return 'push-shrinks';
return 'drop';
}
@@ -2,7 +2,11 @@
import {create} from '@bufbuild/protobuf';
import {MAX_GROUP_DM_OTHER_RECIPIENTS} from '@fluxer/constants/src/LimitConstants';
import {encodeSyncedPreferences, SyncedPreferencesSchema} from '@fluxer/schema/src/domains/user/SyncedPreferencesCodec';
import {
encodeSyncedPreferences,
SYNCED_PREFERENCES_MAX_ENCODED_LENGTH,
SyncedPreferencesSchema,
} from '@fluxer/schema/src/domains/user/SyncedPreferencesCodec';
import {
CreatePrivateChannelRequest,
CustomStatusPayload,
@@ -60,8 +64,12 @@ describe('UserSettingsUpdateRequest synced_preferences', () => {
false,
);
});
it('accepts a string at exactly the size cap', () => {
const atCap = 'A'.repeat(SYNCED_PREFERENCES_MAX_ENCODED_LENGTH);
expect(UserSettingsUpdateRequest.safeParse({synced_preferences: atCap}).success).toBe(true);
});
it('rejects strings exceeding the size cap', () => {
const oversized = 'A'.repeat(400000);
const oversized = 'A'.repeat(SYNCED_PREFERENCES_MAX_ENCODED_LENGTH + 1);
expect(UserSettingsUpdateRequest.safeParse({synced_preferences: oversized}).success).toBe(false);
});
});
@@ -11,7 +11,7 @@ export type {SyncedPreferences} from '@fluxer/schema/src/gen/fluxer/user/prefere
export {SyncedPreferencesSchema} from '@fluxer/schema/src/gen/fluxer/user/preferences/v1/preferences_pb';
export const EMPTY_SYNCED_PREFERENCES_ENCODED = '';
export const SYNCED_PREFERENCES_MAX_BYTES = 256 * 1024;
export const SYNCED_PREFERENCES_MAX_BYTES = 512 * 1024;
export const SYNCED_PREFERENCES_MAX_ENCODED_LENGTH = Math.ceil(SYNCED_PREFERENCES_MAX_BYTES / 3) * 4;
const BASE64_PATTERN = /^[A-Za-z0-9+/_-]*={0,2}$/;