Compare commits

...
40 changed files with 1136 additions and 154 deletions
@@ -632,6 +632,8 @@ describe('Deferred phone verification gate', () => {
await configurePhoneGate({deferred_phone_gate_window_hours: 0.0001});
const outsideWindow = await createGuildWithInvite(harness);
await addFillerMember(outsideWindow.inviteCode);
const beforeJoin = await readFlags(subject.userId);
expect(beforeJoin & DEFERRED_PHONE_ON_COMMUNITY_JOIN).not.toBe(0);
await createBuilder(harness, subject.token).post(`/invites/${outsideWindow.inviteCode}`).expect(200).execute();
const flags = await readFlags(subject.userId);
@@ -0,0 +1,108 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {ChannelTypes} from '@fluxer/constants/src/ChannelConstants';
import {describe, expect, it} from 'vitest';
import {
type ChannelID,
createChannelID,
createMessageID,
createUserID,
type MessageID,
type UserID,
} from '../../../BrandedTypes';
import type {ChannelRow} from '../../../database/types/ChannelTypes';
import type {IGatewayService} from '../../../infrastructure/IGatewayService';
import type {UserCacheService} from '../../../infrastructure/UserCacheService';
import type {RequestCache} from '../../../middleware/RequestCacheMiddleware';
import {Channel} from '../../../models/Channel';
import type {IUserRepository} from '../../../user/IUserRepository';
import {MessageProcessingService} from './MessageProcessingService';
const CHANNEL_ID = createChannelID(1532860318772891648n);
const AUTHOR_ID = createUserID(1471426754353995881n);
const RECIPIENT_ID = createUserID(1485344055661987728n);
const MESSAGE_ID = createMessageID(1546325276953149440n);
function dmChannelRow(lastMessageId: MessageID | null): ChannelRow {
return {
channel_id: CHANNEL_ID,
guild_id: null,
type: ChannelTypes.DM,
name: null,
topic: null,
icon_hash: null,
url: null,
parent_id: null,
position: null,
owner_id: null,
recipient_ids: new Set<UserID>([AUTHOR_ID, RECIPIENT_ID]),
nsfw: null,
content_warning_level: null,
content_warning_text: null,
rate_limit_per_user: null,
bitrate: null,
user_limit: null,
voice_connection_limit: null,
rtc_region: null,
last_message_id: lastMessageId,
last_pin_timestamp: null,
permission_overwrites: null,
nicks: null,
soft_deleted: false,
indexed_at: null,
version: 0,
};
}
function buildService(): {service: MessageProcessingService; opened: Array<Channel>} {
const opened: Array<Channel> = [];
const userRepository = {
isDmChannelOpen: async (userId: UserID, _channelId: ChannelID) => userId === AUTHOR_ID,
openPrivateChannelForUser: async (_userId: UserID, channel: Channel) => {
opened.push(channel);
},
} as unknown as IUserRepository;
const userCacheService = {
getUserPartialResponses: async (userIds: Array<UserID>) =>
new Map(userIds.map((userId) => [userId, {id: userId.toString()}])),
} as unknown as UserCacheService;
const gatewayService = {
dispatchPresence: async () => {},
} as unknown as IGatewayService;
const service = new MessageProcessingService(
undefined as never,
userRepository,
userCacheService,
gatewayService,
undefined as never,
undefined as never,
);
return {service, opened};
}
describe('MessageProcessingService.updateDMRecipients', () => {
it('snapshots the new message id when the in-request channel is stale', async () => {
const {service, opened} = buildService();
await service.updateDMRecipients({
channel: new Channel(dmChannelRow(null)),
channelId: CHANNEL_ID,
messageId: MESSAGE_ID,
requestCache: {} as RequestCache,
});
expect(opened).toHaveLength(1);
expect(opened[0].lastMessageId).toBe(MESSAGE_ID);
});
it('keeps a newer last message id already present on the channel', async () => {
const {service, opened} = buildService();
const newer = createMessageID(MESSAGE_ID + 10n);
await service.updateDMRecipients({
channel: new Channel(dmChannelRow(newer)),
channelId: CHANNEL_ID,
messageId: MESSAGE_ID,
requestCache: {} as RequestCache,
});
expect(opened).toHaveLength(1);
expect(opened[0].lastMessageId).toBe(newer);
});
});
@@ -9,7 +9,7 @@ import type {GatewayChannelMention, IGatewayService} from '../../../infrastructu
import type {UserCacheService} from '../../../infrastructure/UserCacheService';
import {Logger} from '../../../Logger';
import type {RequestCache} from '../../../middleware/RequestCacheMiddleware';
import type {Channel} from '../../../models/Channel';
import {Channel} from '../../../models/Channel';
import type {Message} from '../../../models/Message';
import type {User} from '../../../models/User';
import type {ReadStateService} from '../../../read_state/ReadStateService';
@@ -33,6 +33,13 @@ interface MentionProcessingResult {
mentionChannels: Array<GatewayChannelMention>;
}
function channelWithLastMessageId(channel: Channel, messageId: MessageID): Channel {
if (channel.lastMessageId != null && channel.lastMessageId >= messageId) {
return channel;
}
return new Channel({...channel.toRow(), last_message_id: messageId});
}
export class MessageProcessingService {
constructor(
private channelRepository: IChannelRepositoryAggregate,
@@ -64,10 +71,12 @@ export class MessageProcessingService {
async updateDMRecipients({
channel,
channelId,
messageId,
requestCache,
}: {
channel: Channel;
channelId: ChannelID;
messageId: MessageID;
requestCache: RequestCache;
}): Promise<void> {
if (channel.guildId || channel.type !== ChannelTypes.DM) return;
@@ -76,11 +85,12 @@ export class MessageProcessingService {
const openStates = await this.batchCheckDmChannelOpen(recipientIds, channelId);
const closedRecipients = openStates.filter((state) => !state.isOpen);
if (closedRecipients.length === 0) return;
const snapshotChannel = channelWithLastMessageId(channel, messageId);
await Promise.all(
closedRecipients.map((state) =>
this.openDmAndDispatch({
recipientId: state.recipientId,
channel,
channel: snapshotChannel,
requestCache,
}),
),
@@ -980,7 +980,7 @@ export class MessageSendService {
await this.settlePostCreateWork(messageId, [
{
step: 'update_dm_recipients',
promise: this.deps.processingService.updateDMRecipients({channel, channelId, requestCache}),
promise: this.deps.processingService.updateDMRecipients({channel, channelId, messageId, requestCache}),
},
{
step: 'process_message_after_creation',
@@ -1,6 +1,10 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
let cachedEnabled = false;
// Null until a policy has actually been read. A deferral is only ever granted
// while the gate is on, so reading "not known yet" as "off" revokes it from every
// account holding one and answers their next request with the phone requirement
// the deferral exists to hold back.
let cachedEnabled: boolean | null = null;
export function resolveDeferredPhoneGateEnabled(policy: {
deferred_phone_gate_enabled: boolean;
@@ -9,7 +13,7 @@ export function resolveDeferredPhoneGateEnabled(policy: {
return policy.deferred_phone_gate_enabled && !policy.single_community_enabled;
}
export function getCachedDeferredPhoneGateEnabled(): boolean {
export function getCachedDeferredPhoneGateEnabled(): boolean | null {
return cachedEnabled;
}
+1 -1
View File
@@ -134,7 +134,7 @@ function suppressDeferredPhoneFlags(rawFlags: number): number {
if ((rawFlags & DEFERRED_PHONE_ON_COMMUNITY_JOIN) === 0) {
return rawFlags;
}
if (!getCachedDeferredPhoneGateEnabled()) {
if (getCachedDeferredPhoneGateEnabled() === false) {
return rawFlags & ~DEFERRED_PHONE_ON_COMMUNITY_JOIN;
}
return rawFlags & ~DEFERRABLE_PHONE_FLAGS;
@@ -0,0 +1,56 @@
// @vitest-environment happy-dom
// SPDX-License-Identifier: AGPL-3.0-or-later
import {useAntiShiftFloating} from '@app/features/app/hooks/useAntiShiftFloating';
import {act, createElement, useLayoutEffect, useState} from 'react';
import {createRoot, type Root} from 'react-dom/client';
import {afterEach, beforeEach, describe, expect, it} from 'vitest';
let host: HTMLDivElement;
let root: Root;
let target: HTMLDivElement;
function PortalLikeFloating({onReady}: {onReady: (isReady: boolean) => void}) {
const {setFloating, state} = useAntiShiftFloating(target, true, {placement: 'top'});
const [mounted, setMounted] = useState(false);
useLayoutEffect(() => {
setMounted(true);
}, []);
onReady(state.isReady);
return mounted ? createElement('div', {ref: setFloating, 'data-testid': 'floating'}) : null;
}
async function settle(): Promise<void> {
for (let attempt = 0; attempt < 20; attempt += 1) {
await act(async () => {
await new Promise((resolve) => requestAnimationFrame(() => resolve(null)));
await Promise.resolve();
});
}
}
beforeEach(() => {
host = document.createElement('div');
target = document.createElement('div');
document.body.append(host, target);
root = createRoot(host);
});
afterEach(() => {
act(() => {
root.unmount();
});
document.body.replaceChildren();
});
describe('useAntiShiftFloating', () => {
it('positions a floating element that mounts after the first commit', async () => {
const readyStates: Array<boolean> = [];
act(() => {
root.render(createElement(PortalLikeFloating, {onReady: (isReady) => readyStates.push(isReady)}));
});
await settle();
expect(readyStates.at(-1)).toBe(true);
expect(host.querySelector<HTMLElement>('[data-testid="floating"]')?.style.visibility).toBe('visible');
});
});
@@ -279,6 +279,11 @@ export function useAntiShiftFloating(
constrainHeight = false,
} = options;
const floatingRef = useRef<HTMLElement>(null);
const [floatingElement, setFloatingElement] = useState<HTMLElement | null>(null);
const setFloating = useCallback((element: HTMLElement | null) => {
floatingRef.current = element;
setFloatingElement(element);
}, []);
const [state, setState] = useState<FloatingState>(() => {
const {x, y} = target ? getInitialGuess(target, placement, offsetMainAxis, offsetCrossAxis) : {x: -9999, y: -9999};
return {x, y, isReady: false, offsetX: 0, offsetY: 0, placement};
@@ -481,7 +486,7 @@ export function useAntiShiftFloating(
};
}, []);
useLayoutEffect(() => {
if (!enabled || target == null || floatingRef.current == null) {
if (!enabled || target == null || floatingElement == null) {
setState((prev) => ({...prev, isReady: false, offsetX: 0, offsetY: 0}));
return;
}
@@ -489,7 +494,7 @@ export function useAntiShiftFloating(
setState((prev) => ({...prev, isReady: false, offsetX: 0, offsetY: 0}));
return;
}
const floating = floatingRef.current;
const floating = floatingElement;
updatePosition();
if (shouldAutoUpdate) {
const cleanupCallbacks = [
@@ -502,7 +507,7 @@ export function useAntiShiftFloating(
}
};
} else if (shouldObserveFloatingResize) {
cleanupRef.current = observeFloatingResize(floatingRef.current, updatePosition);
cleanupRef.current = observeFloatingResize(floating, updatePosition);
}
return () => {
positionRevisionRef.current += 1;
@@ -519,9 +524,18 @@ export function useAntiShiftFloating(
}
setState((prev) => ({...prev, isReady: false, offsetX: 0, offsetY: 0}));
};
}, [enabled, target, shouldAutoUpdate, shouldObserveFloatingResize, updatePosition, updatePositionNow]);
}, [
enabled,
target,
floatingElement,
shouldAutoUpdate,
shouldObserveFloatingResize,
updatePosition,
updatePositionNow,
]);
return {
ref: floatingRef,
setFloating,
state,
style: {
position: 'fixed' as const,
@@ -29,6 +29,8 @@ import {
resolveChannelMessagesWindowStatus,
selectChannelMessagesFillerVisible,
selectChannelMessagesSpacerHeight,
selectChannelMessagesTailGapId,
selectChannelMessagesTailProbeId,
selectChannelMessagesWindowBar,
} from '@app/features/messaging/state/ChannelMessagesLoadStateMachine';
import MessageEdit from '@app/features/messaging/state/MessageEdit';
@@ -66,7 +68,9 @@ import {clsx} from 'clsx';
import {runInAction} from 'mobx';
import {observer, useLocalObservable} from 'mobx-react-lite';
import type React from 'react';
import {useCallback, useEffect, useMemo, useRef} from 'react';
import {useCallback, useEffect, useMemo, useRef, useState} from 'react';
const TAIL_PROBE_MAX_ATTEMPTS = 2;
const MESSAGE_LIST_FOR_DESCRIPTOR = msg({
message: 'Message list for {channelName}',
@@ -167,6 +171,9 @@ export const Messages = observer(function Messages({
const scrollerContainerRef = useRef<HTMLDivElement | null>(null);
const lastStateSnapshotRef = useRef<MessagesStateSnapshot | null>(null);
const recoveryFetchChannelIdRef = useRef<string | null>(null);
const tailProbeKeyRef = useRef<string | null>(null);
const tailProbeAttemptsRef = useRef<{key: string; attempts: number} | null>(null);
const [settledTailProbeKey, setSettledTailProbeKey] = useState<string | null>(null);
interface MessageState extends MessagesStateSnapshot {
highlightedMessageId: string | null;
isAtBottom: boolean;
@@ -202,6 +209,17 @@ export const Messages = observer(function Messages({
});
const windowBar = selectChannelMessagesWindowBar(windowStatus);
const windowNeedsPage = windowStatus.needsPage;
const tailInput = {
status: windowStatus,
loading: safeMessages.loadingMore || safeMessages.probeLoading,
newestLoadedMessageId: safeMessages.last()?.id ?? null,
knownLatestMessageId: state.lastReadStateMessageId,
};
const tailWatermarkMessageId = state.lastReadStateMessageId;
const tailGapMessageId = selectChannelMessagesTailGapId(tailInput);
const tailProbeMessageId = selectChannelMessagesTailProbeId(tailInput);
const tailGapKey = tailGapMessageId == null ? null : `${tailGapMessageId}:${tailWatermarkMessageId}`;
const tailProbeKey = tailProbeMessageId == null ? null : `${tailProbeMessageId}:${tailWatermarkMessageId}`;
const canAutoAck = shouldAutoAck({
channelActive: allowAutoAck,
windowFocused: isWindowFocused,
@@ -454,6 +472,42 @@ export const Messages = observer(function Messages({
}
});
}, [channel.id, isGatewayConnected, selectedChannelId, windowNeedsPage, state.messageVersion]);
useEffect(() => {
if (!isGatewayConnected) {
tailProbeKeyRef.current = null;
return;
}
if (tailProbeMessageId == null || tailProbeKey == null || tailWatermarkMessageId == null) {
return;
}
if (selectedChannelId !== channel.id || tailProbeKeyRef.current === tailProbeKey) {
return;
}
tailProbeKeyRef.current = tailProbeKey;
void MessageCommands.fetchMessages(channel.id, null, tailProbeMessageId, MAX_MESSAGES_PER_CHANNEL, undefined, {
tailProbe: {
watermarkMessageId: tailWatermarkMessageId,
onSettled: (settlement) => {
if (settlement === 'applied') {
setSettledTailProbeKey(tailProbeKey);
return;
}
if (settlement === 'failed') {
const attempts =
tailProbeAttemptsRef.current?.key === tailProbeKey ? tailProbeAttemptsRef.current.attempts : 0;
if (attempts >= TAIL_PROBE_MAX_ATTEMPTS) {
setSettledTailProbeKey(tailProbeKey);
return;
}
tailProbeAttemptsRef.current = {key: tailProbeKey, attempts: attempts + 1};
}
if (tailProbeKeyRef.current === tailProbeKey) {
tailProbeKeyRef.current = null;
}
},
},
});
}, [channel.id, isGatewayConnected, selectedChannelId, tailProbeKey, tailProbeMessageId, tailWatermarkMessageId]);
useMessageListKeyboardNavigation({
containerRef: scrollManager.ref,
channelId: channel.id,
@@ -488,10 +542,11 @@ export const Messages = observer(function Messages({
}, []);
useEffect(() => {
if (!canAutoAck || !state.isAtBottom || !state.messages?.ready) return;
if (tailGapKey != null && settledTailProbeKey !== tailGapKey) return;
if (ReadStates.hasUnread(channel.id)) {
ReadStateCommands.ackWithStickyUnread(channel.id);
}
}, [canAutoAck, state.isAtBottom, state.messages?.ready, channel.id]);
}, [canAutoAck, state.isAtBottom, state.messages?.ready, tailGapKey, settledTailProbeKey, channel.id]);
useEffect(() => {
return () => {
const readState = ReadStates.getIfExists(channel.id);
@@ -23,6 +23,7 @@ import ReadStates from '@app/features/read_state/state/ReadStates';
import QuickSwitcher from '@app/features/search/state/QuickSwitcher';
import Nagbar from '@app/features/ui/state/Nagbar';
import UserGuildSettings from '@app/features/user/state/UserGuildSettings';
import UserProfile from '@app/features/user/state/UserProfile';
import MediaEngine from '@app/features/voice/engine/MediaEngineFacade';
import {FAVORITES_GUILD_ID} from '@fluxer/constants/src/AppConstants';
@@ -51,6 +52,7 @@ export function handleGuildCreate(data: GuildReadyData, _context: GatewayHandler
GuildCount.handleGuildCreate(data);
if (!isSync) {
MemberSidebar.handleGuildCreate(data.id);
UserProfile.handleGuildCreate();
}
if (!data.unavailable) {
Channels.handleGuildCreate(data);
@@ -20,6 +20,7 @@ import MentionFeed from '@app/features/notification/state/MentionFeed';
import Permission from '@app/features/permissions/state/Permission';
import Presence from '@app/features/presence/state/Presence';
import QuickSwitcher from '@app/features/search/state/QuickSwitcher';
import UserProfile from '@app/features/user/state/UserProfile';
import MediaEngine from '@app/features/voice/engine/MediaEngineFacade';
import Webhooks from '@app/features/webhook/state/Webhooks';
import type {Guild} from '@fluxer/schema/src/domains/guild/GuildResponseSchemas';
@@ -52,5 +53,6 @@ export function handleGuildDelete(data: GuildDeletePayload, _context: GatewayHan
Messages.handleGuildUnavailable(data.id, data.unavailable ?? false);
Messages.handleCleanup();
MentionFeed.handleGuildDelete(data.id);
UserProfile.handleGuildDelete(data.unavailable);
QuickSwitcher.recomputeIfOpen();
}
@@ -6,6 +6,7 @@ import GuildMembers from '@app/features/member/state/GuildMembers';
import MemberSearch from '@app/features/member/state/MemberSearch';
import Permission from '@app/features/permissions/state/Permission';
import Presence from '@app/features/presence/state/Presence';
import UserProfile from '@app/features/user/state/UserProfile';
import Users from '@app/features/user/state/Users';
import type {GuildMemberData} from '@fluxer/schema/src/domains/guild/GuildMemberSchemas';
import type {User} from '@fluxer/schema/src/domains/user/UserResponseSchemas';
@@ -21,4 +22,5 @@ export function handleGuildMemberAdd(data: GuildMemberAddPayload, _context: Gate
GuildVerification.handleGuildMemberUpdate(data.guild_id);
Presence.handleGuildMemberAdd(data.guild_id, data.user.id);
MemberSearch.handleMemberAdd(data.guild_id, data.user.id);
UserProfile.handleGuildMemberAdd(data.user.id);
}
@@ -5,6 +5,7 @@ import GuildMembers from '@app/features/member/state/GuildMembers';
import MemberSearch from '@app/features/member/state/MemberSearch';
import Permission from '@app/features/permissions/state/Permission';
import Presence from '@app/features/presence/state/Presence';
import UserProfile from '@app/features/user/state/UserProfile';
import Users from '@app/features/user/state/Users';
import type {UserPartial} from '@fluxer/schema/src/domains/user/UserResponseSchemas';
@@ -19,4 +20,5 @@ export function handleGuildMemberRemove(data: GuildMemberRemovePayload, _context
GuildMembers.handleMemberRemove(data.guild_id, data.user.id);
Permission.handleGuildMemberUpdate(data.user.id);
Presence.handleGuildMemberRemove(data.guild_id, data.user.id);
UserProfile.handleGuildMemberRemove(data.user.id);
}
@@ -271,7 +271,7 @@ function SelectionToolbarSurface({
}),
[editor],
);
const {ref, state, style, updatePosition} = useAntiShiftFloating(virtualReference, true, {
const {setFloating, state, style, updatePosition} = useAntiShiftFloating(virtualReference, true, {
placement: 'top',
offsetMainAxis: 8,
shouldAutoUpdate: true,
@@ -381,7 +381,7 @@ function SelectionToolbarSurface({
>
<flx-lexical-selection-formatting-toolbar
ref={(element) => {
(ref as React.MutableRefObject<HTMLElement | null>).current = element;
setFloating(element);
toolbarRef.current = element;
}}
className={flxElementClassName(styles.toolbar)}
@@ -23,6 +23,7 @@ import {MessageEditFailedModal} from '@app/features/messaging/components/alerts/
import {MessageEditTooQuickModal} from '@app/features/messaging/components/alerts/MessageEditTooQuickModal';
import type {Message as MessageModel} from '@app/features/messaging/models/MessagingMessage';
import type {JumpOptions} from '@app/features/messaging/state/ChannelMessages';
import {selectChannelMessagesTailProbeOutcome} from '@app/features/messaging/state/ChannelMessagesLoadStateMachine';
import MessageEdit from '@app/features/messaging/state/MessageEdit';
import MessageEditMobile from '@app/features/messaging/state/MessageEditMobile';
import MessageQueue from '@app/features/messaging/state/MessageQueue';
@@ -115,6 +116,19 @@ export interface JumpToMessageOptions {
interface FetchMessagesOptions {
throwOnError?: boolean;
tailProbe?: TailProbeContext;
}
export type TailProbeSettlement = 'applied' | 'retry' | 'failed';
export interface TailProbeContext {
watermarkMessageId: string;
onSettled: (settlement: TailProbeSettlement) => void;
}
interface TailProbeEpoch {
loadGeneration: number;
jumpTicket: number;
}
interface MessagePageState {
@@ -142,11 +156,12 @@ function makeFetchKey(
): string {
const SEP = '\x1f';
const throwOnError = options?.throwOnError ? '1' : '0';
const tailProbe = options?.tailProbe ? `1.${options.tailProbe.watermarkMessageId}` : '0';
if (!jump) {
return `${channelId}${SEP}${before ?? ''}${SEP}${after ?? ''}${SEP}${limit}${SEP}${throwOnError}`;
return `${channelId}${SEP}${before ?? ''}${SEP}${after ?? ''}${SEP}${limit}${SEP}${throwOnError}${SEP}${tailProbe}`;
}
return (
`${channelId}${SEP}${before ?? ''}${SEP}${after ?? ''}${SEP}${limit}${SEP}${throwOnError}${SEP}` +
`${channelId}${SEP}${before ?? ''}${SEP}${after ?? ''}${SEP}${limit}${SEP}${throwOnError}${SEP}${tailProbe}${SEP}` +
`${jump.present ? '1' : '0'}${SEP}${jump.messageId ?? ''}${SEP}${jump.offset ?? 0}${SEP}` +
`${jump.flash ? '1' : '0'}${SEP}${jump.returnToMessageId ?? ''}${SEP}` +
`${jump.returnChannelId ?? ''}${SEP}${jump.returnGuildId ?? ''}${SEP}${jump.jumpType ?? ''}`
@@ -244,6 +259,7 @@ function handleMessageFetchSuccess(
messages: Array<WireMessage>,
pageState: MessagePageState,
jump?: JumpOptions,
tailProbe?: TailProbeContext,
): void {
Messages.handleLoadMessagesSuccess({
channelId,
@@ -254,11 +270,13 @@ function handleMessageFetchSuccess(
hasMoreAfter: pageState.hasMoreAfter,
cached: false,
jump,
tailProbe: tailProbe != null,
});
ReadStates.handleLoadMessages({
channelId,
isAfter: pageState.isAfter,
messages,
tailProbeWatermarkId: tailProbe?.watermarkMessageId ?? null,
});
MessageReferences.handleMessagesFetchSuccess(channelId, messages);
void requestMissingGuildMembers(channelId, messages);
@@ -374,6 +392,30 @@ function applyMessageFetchCacheHit(
}
}
function isTailProbeApplicable(channelId: string, probe: TailProbeEpoch, after: string | null): boolean {
const current = Messages.getMessages(channelId);
const outcome = selectChannelMessagesTailProbeOutcome({
probeGeneration: probe.loadGeneration,
currentGeneration: current.loadGeneration,
probeJumpTicket: probe.jumpTicket,
currentJumpTicket: current.jumpTicket,
ready: current.ready,
hasMoreAfter: current.hasMoreAfter,
anchorMessageId: after,
newestLoadedMessageId: current.last()?.id ?? null,
});
if (outcome === 'apply') {
return true;
}
logger.debug(`Discarding tail probe for channel ${channelId} (${outcome})`);
Messages.handleTailProbeSettled({channelId});
return false;
}
function settleTailProbe(options: FetchMessagesOptions | undefined, settlement: TailProbeSettlement): void {
options?.tailProbe?.onSettled(settlement);
}
export async function fetchMessages(
channelId: string,
before: string | null,
@@ -392,13 +434,16 @@ export async function fetchMessages(
switch (preflightDecision.type) {
case 'useInFlightRequest':
logger.debug(`Using in-flight fetchMessages for channel ${channelId} (deduped)`);
settleTailProbe(options, 'retry');
return inFlight as Promise<Array<WireMessage>>;
case 'blockForGate':
logger.debug(`Skipping message fetch for gated channel ${channelId}`);
Messages.handleLoadMessagesBlocked({channelId});
settleTailProbe(options, 'retry');
return [];
case 'useCache':
applyMessageFetchCacheHit(channelId, preflightDecision.cacheHit, before, after, limit, jump);
settleTailProbe(options, 'retry');
return [];
case 'startFetch':
break;
@@ -409,19 +454,38 @@ export async function fetchMessages(
forceFailure: DeveloperOptions.forceFailMessageLoads,
});
if (executionDecision.type === 'simulateFailure') {
if (options?.tailProbe != null) {
settleTailProbe(options, 'failed');
return [];
}
return handleForcedMessageLoadFailure(channelId, jump);
}
Messages.handleLoadMessages({channelId, jump});
Messages.handleLoadMessages({channelId, jump, tailProbe: options?.tailProbe != null});
let probeEpoch: TailProbeEpoch | null = null;
if (options?.tailProbe) {
const started = Messages.getMessages(channelId);
probeEpoch = {loadGeneration: started.loadGeneration, jumpTicket: started.jumpTicket};
}
try {
const timeStart = Date.now();
logger.debug(`Fetching messages for channel ${channelId}`);
const messages = await requestChannelMessages(channelId, before, after, limit, jump);
if (probeEpoch != null && !isTailProbeApplicable(channelId, probeEpoch, after)) {
settleTailProbe(options, 'retry');
return [];
}
const pageState = calculateMessagePageState(channelId, before, after, limit, messages, jump);
logger.info(`Fetched ${messages.length} messages for channel ${channelId}, took ${Date.now() - timeStart}ms`);
handleMessageFetchSuccess(channelId, messages, pageState, jump);
handleMessageFetchSuccess(channelId, messages, pageState, jump, options?.tailProbe);
settleTailProbe(options, 'applied');
return messages;
} catch (error) {
logger.error(`Failed to fetch messages for channel ${channelId}:`, error);
if (probeEpoch != null) {
Messages.handleTailProbeSettled({channelId});
settleTailProbe(options, 'failed');
return [];
}
Messages.handleLoadMessagesFailure({channelId});
if (options?.throwOnError) {
throw error;
@@ -5,6 +5,7 @@ import {UploadingAttachment} from '@app/features/messaging/models/UploadingAttac
import {resolveChannelIncomingMessageDecision} from '@app/features/messaging/state/ChannelIncomingMessageStateMachine';
import {resolveChannelMessagesLoadDecision} from '@app/features/messaging/state/ChannelMessagesLoadStateMachine';
import MessageReactions from '@app/features/messaging/state/MessageReactions';
import {mergeAscendingById} from '@app/features/messaging/utils/MessagePaginationUtils';
import SelectedChannel from '@app/features/navigation/state/SelectedChannel';
import type {JumpType} from '@fluxer/constants/src/JumpConstants';
import {JumpTypes} from '@fluxer/constants/src/JumpConstants';
@@ -48,6 +49,7 @@ interface LoadCompleteOptions {
hasMoreBefore?: boolean;
hasMoreAfter?: boolean;
cached?: boolean;
tailProbe?: boolean;
}
type MessageInput = Message | WireMessage;
@@ -281,6 +283,8 @@ class MessageBufferSegment {
}
}
let nextLoadGeneration = 0;
export class ChannelMessages {
private static readonly channelCache = new Map<string, ChannelMessages>();
private static readonly maxChannelsInMemory = 50;
@@ -292,6 +296,8 @@ export class ChannelMessages {
jumpDestinationId: string | null = null;
jumpDestinationOffset = 0;
jumpTicket = 1;
loadGeneration = 0;
probeLoading = false;
hasJumped = false;
landedAtLiveEdge = false;
jumpHighlight = true;
@@ -605,9 +611,9 @@ export class ChannelMessages {
}, true);
}
merge(records: Array<Message>, prepend = false, clearBuffer = false): ChannelMessages {
merge(records: Array<Message>, prepend = false, clearBuffer = false, ordered = false): ChannelMessages {
return this.cloneAnd((draft) => {
draft.mergeInto(records, prepend, clearBuffer);
draft.mergeInto(records, prepend, clearBuffer, ordered);
}, true);
}
@@ -822,8 +828,18 @@ export class ChannelMessages {
return this.cloneAnd({ready: true, cached: true}).merge([hydrateMessage(this, wire, 'preserve')]);
}
beginProbeLoad(): ChannelMessages {
return this.cloneAnd({loadGeneration: ++nextLoadGeneration, probeLoading: true});
}
endProbeLoad(): ChannelMessages {
if (!this.probeLoading) return this;
return this.cloneAnd({probeLoading: false});
}
beginLoad(jump?: JumpOptions): ChannelMessages {
return this.cloneAnd({
loadGeneration: ++nextLoadGeneration,
loadingMore: true,
hasJumped: jump != null,
landedAtLiveEdge: jump?.present ?? false,
@@ -845,6 +861,7 @@ export class ChannelMessages {
hasMoreBefore = false,
hasMoreAfter = false,
cached = false,
tailProbe = false,
} = options;
const records = [...windowMessages].reverse().map((m) => hydrateMessage(this, m, 'empty'));
const loadDecision = resolveChannelMessagesLoadDecision({
@@ -862,26 +879,32 @@ export class ChannelMessages {
next = next.merge(unsent);
}
} else {
next = this.merge(records, loadDecision.prepend, true);
if (loadDecision.trimBottom) {
next = this.merge(records, loadDecision.prepend, true, true);
if (!tailProbe && loadDecision.trimBottom) {
next = next.trimToWindow(true, false);
} else if (loadDecision.trimTop) {
} else if (!tailProbe && loadDecision.trimTop) {
next = next.trimToWindow(false, true);
}
}
const jumpPatch = tailProbe
? {}
: {
jumpType: jump?.jumpType ?? JumpTypes.ANIMATED,
jumpHighlight: jump?.flash ?? false,
hasJumped: jump != null,
landedAtLiveEdge: jump?.present ?? false,
jumpDestinationId: jump?.messageId ?? null,
jumpDestinationOffset: jump && jump.messageId != null && jump.offset != null ? jump.offset : 0,
jumpTicket: jump ? next.jumpTicket + 1 : next.jumpTicket,
jumpReturnMessageId: jump?.returnToMessageId ?? null,
jumpReturnChannelId: jump?.returnToMessageId ? (jump.returnChannelId ?? this.channelId) : null,
jumpReturnGuildId: jump?.returnToMessageId ? (jump.returnGuildId ?? null) : null,
};
next = next.cloneAnd({
ready: true,
loadingMore: false,
jumpType: jump?.jumpType ?? JumpTypes.ANIMATED,
jumpHighlight: jump?.flash ?? false,
hasJumped: jump != null,
landedAtLiveEdge: jump?.present ?? false,
jumpDestinationId: jump?.messageId ?? null,
jumpDestinationOffset: jump && jump.messageId != null && jump.offset != null ? jump.offset : 0,
jumpTicket: jump ? next.jumpTicket + 1 : next.jumpTicket,
jumpReturnMessageId: jump?.returnToMessageId ?? null,
jumpReturnChannelId: jump?.returnToMessageId ? (jump.returnChannelId ?? this.channelId) : null,
jumpReturnGuildId: jump?.returnToMessageId ? (jump.returnGuildId ?? null) : null,
probeLoading: false,
...jumpPatch,
hasMoreBefore: loadDecision.preserveHasMoreBefore ? next.hasMoreBefore : hasMoreBefore,
hasMoreAfter: loadDecision.preserveHasMoreAfter ? next.hasMoreAfter : hasMoreAfter,
cached,
@@ -914,7 +937,7 @@ export class ChannelMessages {
this.messageIndex = {};
}
private mergeInto(incoming: Array<Message>, prepend = false, clearSideBuffer = false): void {
private mergeInto(incoming: Array<Message>, prepend = false, clearSideBuffer = false, ordered = false): void {
const newItems: Array<Message> = [];
for (const msg of incoming) {
const existing = this.messageIndex[msg.id];
@@ -937,7 +960,11 @@ export class ChannelMessages {
buffer.clear();
}
if (newItems.length === 0) return;
this.messageList = prepend ? newItems.concat(this.messageList) : this.messageList.concat(newItems);
if (prepend) {
this.messageList = newItems.concat(this.messageList);
return;
}
this.messageList = ordered ? mergeAscendingById(this.messageList, newItems) : this.messageList.concat(newItems);
}
private cloneAnd(
@@ -956,6 +983,8 @@ export class ChannelMessages {
clone.jumpDestinationId = this.jumpDestinationId;
clone.jumpDestinationOffset = this.jumpDestinationOffset;
clone.jumpTicket = this.jumpTicket;
clone.loadGeneration = this.loadGeneration;
clone.probeLoading = this.probeLoading;
clone.hasJumped = this.hasJumped;
clone.landedAtLiveEdge = this.landedAtLiveEdge;
clone.jumpHighlight = this.jumpHighlight;
@@ -978,6 +1007,8 @@ export class ChannelMessages {
clone.jumpDestinationOffset =
patch.jumpDestinationOffset !== undefined ? patch.jumpDestinationOffset : this.jumpDestinationOffset;
clone.jumpTicket = patch.jumpTicket !== undefined ? patch.jumpTicket : this.jumpTicket;
clone.loadGeneration = patch.loadGeneration !== undefined ? patch.loadGeneration : this.loadGeneration;
clone.probeLoading = 'probeLoading' in patch ? !!patch.probeLoading : this.probeLoading;
clone.hasJumped = 'hasJumped' in patch ? !!patch.hasJumped : this.hasJumped;
clone.landedAtLiveEdge = 'landedAtLiveEdge' in patch ? !!patch.landedAtLiveEdge : this.landedAtLiveEdge;
clone.jumpHighlight = 'jumpHighlight' in patch ? !!patch.jumpHighlight : this.jumpHighlight;
@@ -12,6 +12,9 @@ import {
selectChannelMessagesFillerVisible,
selectChannelMessagesLoadDecision,
selectChannelMessagesSpacerHeight,
selectChannelMessagesTailGapId,
selectChannelMessagesTailProbeId,
selectChannelMessagesTailProbeOutcome,
selectChannelMessagesWindowBar,
selectChannelMessagesWindowStatus,
transitionChannelMessagesLoadSnapshot,
@@ -389,3 +392,129 @@ describe('selectChannelMessagesFillerVisible', () => {
expect(streamCases).toBeGreaterThan(0);
});
});
describe('selectChannelMessagesTailProbeId', () => {
const ID = {older: '1519773906704011264', newer: '1519773906708205568'};
function tailInput(
overrides: Partial<ChannelMessagesWindowInput> = {},
tail: {
loading?: boolean;
newestLoadedMessageId?: string | null;
knownLatestMessageId?: string | null;
} = {},
) {
const windowInput: ChannelMessagesWindowInput = {
ready: true,
loading: false,
failed: false,
messageCount: 12,
hasMoreBefore: true,
hasMoreAfter: false,
...overrides,
};
return {
status: resolveChannelMessagesWindowStatus(windowInput),
loading: tail.loading ?? windowInput.loading,
newestLoadedMessageId: tail.newestLoadedMessageId === undefined ? ID.older : tail.newestLoadedMessageId,
knownLatestMessageId: tail.knownLatestMessageId === undefined ? ID.newer : tail.knownLatestMessageId,
};
}
it('probes from the newest loaded message when the live edge is known to be ahead', () => {
expect(selectChannelMessagesTailProbeId(tailInput())).toBe(ID.older);
});
it('stays silent for a window that already advertises a newer page', () => {
expect(selectChannelMessagesTailProbeId(tailInput({hasMoreAfter: true}))).toBeNull();
});
it('stays silent for a failed window so a failed probe cannot re-arm itself', () => {
expect(selectChannelMessagesTailProbeId(tailInput({failed: true}))).toBeNull();
});
it('stays silent while a load is in flight', () => {
expect(selectChannelMessagesTailProbeId(tailInput({loading: true}))).toBeNull();
});
it('stays silent for a window that has not loaded a page yet', () => {
expect(selectChannelMessagesTailProbeId(tailInput({ready: false}))).toBeNull();
});
it('never probes an empty window, so a channel with no readable history cannot arm it', () => {
expect(
selectChannelMessagesTailProbeId(
tailInput({messageCount: 0, hasMoreBefore: false}, {newestLoadedMessageId: null}),
),
).toBeNull();
});
it('stays silent when the watermark matches or trails the newest loaded message', () => {
expect(selectChannelMessagesTailProbeId(tailInput({}, {knownLatestMessageId: ID.older}))).toBeNull();
expect(
selectChannelMessagesTailProbeId(
tailInput({}, {newestLoadedMessageId: ID.newer, knownLatestMessageId: ID.older}),
),
).toBeNull();
expect(selectChannelMessagesTailProbeId(tailInput({}, {knownLatestMessageId: null}))).toBeNull();
});
it('keeps reporting the tail gap while the probe it armed is in flight', () => {
const inFlight = tailInput({loading: true});
expect(selectChannelMessagesTailProbeId(inFlight)).toBeNull();
expect(selectChannelMessagesTailGapId(inFlight)).toBe(ID.older);
});
it('reports no tail gap once the newest loaded message reaches the watermark', () => {
expect(selectChannelMessagesTailGapId(tailInput({}, {newestLoadedMessageId: ID.newer}))).toBeNull();
});
it('still reports the gap on a deep window so auto-ack cannot run past it', () => {
expect(selectChannelMessagesTailGapId(tailInput({messageCount: 200}))).toBe(ID.older);
expect(selectChannelMessagesTailProbeId(tailInput({messageCount: 200}))).toBe(ID.older);
});
});
describe('selectChannelMessagesTailProbeOutcome', () => {
const ID = {older: '1519773906704011264', newer: '1519773906708205568'};
function outcomeInput(overrides: Partial<Parameters<typeof selectChannelMessagesTailProbeOutcome>[0]> = {}) {
return {
probeGeneration: 4,
currentGeneration: 4,
probeJumpTicket: 2,
currentJumpTicket: 2,
ready: true,
hasMoreAfter: false,
anchorMessageId: ID.older,
newestLoadedMessageId: ID.older,
...overrides,
};
}
it('applies a probe onto the window it was issued against', () => {
expect(selectChannelMessagesTailProbeOutcome(outcomeInput())).toBe('apply');
});
it('discards without touching shared state once another load owns the window', () => {
expect(selectChannelMessagesTailProbeOutcome(outcomeInput({currentGeneration: 5}))).toBe('discard');
});
it('releases a probe whose window a jump has taken out of the ready state', () => {
expect(selectChannelMessagesTailProbeOutcome(outcomeInput({ready: false}))).toBe('release');
expect(selectChannelMessagesTailProbeOutcome(outcomeInput({currentGeneration: 9, ready: false}))).toBe('discard');
});
it('releases a probe whose tail moved under it', () => {
expect(selectChannelMessagesTailProbeOutcome(outcomeInput({newestLoadedMessageId: ID.newer}))).toBe('release');
expect(selectChannelMessagesTailProbeOutcome(outcomeInput({newestLoadedMessageId: null}))).toBe('release');
});
it('releases a probe whose window started advertising a newer page', () => {
expect(selectChannelMessagesTailProbeOutcome(outcomeInput({hasMoreAfter: true}))).toBe('release');
});
it('discards a probe whose window a cached jump took over without starting a load', () => {
expect(selectChannelMessagesTailProbeOutcome(outcomeInput({currentJumpTicket: 3}))).toBe('discard');
});
});
@@ -1,5 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {compare as compareSnowflakes} from '@fluxer/snowflake/src/SnowflakeUtils';
import {assign, getInitialSnapshot, type SnapshotFrom, setup, transition} from 'xstate';
export interface ChannelMessagesLoadInput {
@@ -281,6 +282,47 @@ export function selectChannelMessagesSpacerHeight(status: ChannelMessagesWindowS
return status.olderPageAvailable || status.newerPageAvailable ? fillerHeight : 0;
}
export interface ChannelMessagesTailInput {
status: ChannelMessagesWindowStatus;
loading: boolean;
newestLoadedMessageId: string | null;
knownLatestMessageId: string | null;
}
export function selectChannelMessagesTailGapId(input: ChannelMessagesTailInput): string | null {
if (input.status.phase !== 'stream' || input.status.retryVisible || input.status.newerPageAvailable) return null;
if (input.newestLoadedMessageId == null || input.knownLatestMessageId == null) return null;
if (compareSnowflakes(input.knownLatestMessageId, input.newestLoadedMessageId) <= 0) return null;
return input.newestLoadedMessageId;
}
export function selectChannelMessagesTailProbeId(input: ChannelMessagesTailInput): string | null {
if (input.loading) return null;
return selectChannelMessagesTailGapId(input);
}
export interface ChannelMessagesTailProbeResultInput {
probeGeneration: number;
currentGeneration: number;
probeJumpTicket: number;
currentJumpTicket: number;
ready: boolean;
hasMoreAfter: boolean;
anchorMessageId: string | null;
newestLoadedMessageId: string | null;
}
export type ChannelMessagesTailProbeOutcome = 'apply' | 'release' | 'discard';
export function selectChannelMessagesTailProbeOutcome(
input: ChannelMessagesTailProbeResultInput,
): ChannelMessagesTailProbeOutcome {
if (input.currentGeneration !== input.probeGeneration) return 'discard';
if (input.currentJumpTicket !== input.probeJumpTicket) return 'discard';
if (!input.ready || input.hasMoreAfter || input.newestLoadedMessageId !== input.anchorMessageId) return 'release';
return 'apply';
}
export function selectChannelMessagesWindowBar(status: ChannelMessagesWindowStatus): ChannelMessagesWindowBar {
if (status.phase === 'placeholder') return 'none';
if (status.phase === 'retry' || status.retryVisible) return 'retry';
@@ -449,9 +449,9 @@ class Messages {
}
@action
handleLoadMessages(action: {channelId: string; jump?: JumpOptions}): boolean {
handleLoadMessages(action: {channelId: string; jump?: JumpOptions; tailProbe?: boolean}): boolean {
const messages = ChannelMessages.getOrCreate(action.channelId);
this.commitMessages(messages.beginLoad(action.jump));
this.commitMessages(action.tailProbe ? messages.beginProbeLoad() : messages.beginLoad(action.jump));
this.notifyChange();
return false;
}
@@ -505,6 +505,7 @@ class Messages {
hasMoreBefore?: boolean;
hasMoreAfter?: boolean;
cached?: boolean;
tailProbe?: boolean;
messages: Array<WireMessage>;
}): boolean {
const messages = ChannelMessages.getOrCreate(action.channelId).applyLoadedWindow({
@@ -515,6 +516,7 @@ class Messages {
hasMoreBefore: action.hasMoreBefore,
hasMoreAfter: action.hasMoreAfter,
cached: action.cached,
tailProbe: action.tailProbe,
});
this.commitMessages(messages);
this.notifyChange();
@@ -529,6 +531,15 @@ class Messages {
return false;
}
@action
handleTailProbeSettled(action: {channelId: string}): boolean {
const messages = ChannelMessages.get(action.channelId);
if (!messages?.probeLoading) return false;
this.commitMessages(messages.endProbeLoad());
this.notifyChange();
return true;
}
@action
handleLoadMessagesBlocked(action: {channelId: string}): boolean {
const messages = ChannelMessages.getOrCreate(action.channelId);
@@ -1,7 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {describe, expect, it} from 'vitest';
import {calculateAroundPaginationState, getAroundWindowCounts} from './MessagePaginationUtils';
import {calculateAroundPaginationState, getAroundWindowCounts, mergeAscendingById} from './MessagePaginationUtils';
describe('MessagePaginationUtils', () => {
it('splits around windows with the newer side receiving the extra item for even limits', () => {
@@ -42,3 +42,30 @@ describe('MessagePaginationUtils', () => {
expect(state.hasMoreAfter).toBe(false);
});
});
describe('mergeAscendingById', () => {
const ID = {a: '1519773906704011264', b: '1519773906708205568', c: '1519773906712399872'};
it('appends when every incoming message is newer than the tail', () => {
const existing = [{id: ID.a}];
expect(mergeAscendingById(existing, [{id: ID.b}, {id: ID.c}])).toEqual([{id: ID.a}, {id: ID.b}, {id: ID.c}]);
});
it('returns the existing list untouched for an empty page', () => {
const existing = [{id: ID.a}];
expect(mergeAscendingById(existing, [])).toBe(existing);
});
it('interleaves a recovered message that is older than a live message already at the tail', () => {
expect(mergeAscendingById([{id: ID.a}, {id: ID.c}], [{id: ID.b}])).toEqual([{id: ID.a}, {id: ID.b}, {id: ID.c}]);
});
it('keeps the whole page in order when it spans a live message already at the tail', () => {
expect(mergeAscendingById([{id: ID.a}, {id: ID.b}], [{id: ID.c}])).toEqual([{id: ID.a}, {id: ID.b}, {id: ID.c}]);
expect(mergeAscendingById([{id: ID.b}], [{id: ID.a}, {id: ID.c}])).toEqual([{id: ID.a}, {id: ID.b}, {id: ID.c}]);
});
it('appends onto an empty window', () => {
expect(mergeAscendingById([], [{id: ID.a}])).toEqual([{id: ID.a}]);
});
});
@@ -1,5 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {compare as compareSnowflakes} from '@fluxer/snowflake/src/SnowflakeUtils';
export interface AroundWindowCounts {
newer: number;
older: number;
@@ -49,3 +51,32 @@ export function calculateAroundPaginationState(input: AroundPaginationStateInput
isAtKnownLatest,
};
}
export function mergeAscendingById<T extends {id: string}>(existing: Array<T>, incoming: Array<T>): Array<T> {
if (incoming.length === 0) return existing;
const tail = existing[existing.length - 1];
if (tail == null || compareSnowflakes(incoming[0].id, tail.id) > 0) {
return existing.concat(incoming);
}
const merged: Array<T> = [];
let left = 0;
let right = 0;
while (left < existing.length && right < incoming.length) {
if (compareSnowflakes(incoming[right].id, existing[left].id) < 0) {
merged.push(incoming[right]);
right += 1;
} else {
merged.push(existing[left]);
left += 1;
}
}
while (left < existing.length) {
merged.push(existing[left]);
left += 1;
}
while (right < incoming.length) {
merged.push(incoming[right]);
right += 1;
}
return merged;
}
@@ -102,6 +102,48 @@ describe('ReadStates unread invariant', () => {
expect(ReadStates.getVisualUnreadMessageId(channelId)).toBeNull();
});
it('lowers a watermark its own probe finds nothing behind', () => {
const {channelId, state} = seedReadChannel();
state.lastMessageId = ID.newer;
loadedMessages.push({id: ID.ack, author: {id: 'someone'}});
ReadStates.handleLoadMessages({channelId, isAfter: true, messages: [], tailProbeWatermarkId: ID.newer});
expect(ReadStates.lastMessageId(channelId)).toBe(ID.ack);
expect(ReadStates.hasUnread(channelId)).toBe(false);
});
it('keeps a watermark that is ahead when an ordinary after page comes back empty', () => {
const {channelId, state} = seedReadChannel();
state.lastMessageId = ID.newer;
loadedMessages.push({id: ID.ack, author: {id: 'someone'}});
ReadStates.handleLoadMessages({channelId, isAfter: true, messages: []});
expect(ReadStates.lastMessageId(channelId)).toBe(ID.newer);
});
it('keeps a watermark that advanced while its own probe was in flight', () => {
const {channelId, state} = seedReadChannel();
state.lastMessageId = ID.newer;
loadedMessages.push({id: ID.ack, author: {id: 'someone'}});
ReadStates.handleLoadMessages({channelId, isAfter: true, messages: [], tailProbeWatermarkId: ID.ack});
expect(ReadStates.lastMessageId(channelId)).toBe(ID.newer);
});
it('keeps a watermark that is ahead when the page was not an after page', () => {
const {channelId, state} = seedReadChannel();
state.lastMessageId = ID.newer;
loadedMessages.push({id: ID.ack, author: {id: 'someone'}});
ReadStates.handleLoadMessages({channelId, messages: [], tailProbeWatermarkId: ID.newer});
expect(ReadStates.lastMessageId(channelId)).toBe(ID.newer);
});
it('keeps a watermark that is ahead when the window is not at the live edge', () => {
const {channelId, state} = seedReadChannel();
hasNewestMessages = false;
state.lastMessageId = ID.newer;
loadedMessages.push({id: ID.ack, author: {id: 'someone'}});
ReadStates.handleLoadMessages({channelId, isAfter: true, messages: [], tailProbeWatermarkId: ID.newer});
expect(ReadStates.lastMessageId(channelId)).toBe(ID.newer);
});
it('anchors the divider when a window is loaded whose ack sits outside it', () => {
const {channelId, state} = seedReadChannel();
state.lastMessageId = ID.newer;
@@ -439,7 +439,12 @@ class ReadStates {
});
}
handleLoadMessages(action: {channelId: string; isAfter?: boolean; messages: Array<WireMessage>}): void {
handleLoadMessages(action: {
channelId: string;
isAfter?: boolean;
messages: Array<WireMessage>;
tailProbeWatermarkId?: string | null;
}): void {
const state = this.get(action.channelId);
state.messagesLoaded = true;
const messages = Messages.getMessages(action.channelId);
@@ -448,6 +453,17 @@ class ReadStates {
state.lastMessageId = newestMessage.id;
}
const landedOnNewestWindow = messages.hasNewestMessages();
if (
action.isAfter &&
action.tailProbeWatermarkId != null &&
action.tailProbeWatermarkId === state.lastMessageId &&
action.messages.length === 0 &&
landedOnNewestWindow &&
newestMessage != null &&
isNewerMessageId(state.lastMessageId, newestMessage.id)
) {
state.lastMessageId = newestMessage.id;
}
const landedOnAck = state.ackMessageId != null && messages.jumpDestinationId === state.ackMessageId;
if (state.hasUnread() || landedOnNewestWindow || landedOnAck) {
state.rebuild();
@@ -428,6 +428,7 @@
max-inline-size: none;
overflow: visible;
white-space: nowrap;
unicode-bidi: plaintext;
user-select: text;
-webkit-user-select: text;
}
@@ -505,6 +506,7 @@
font-weight: 400;
color: var(--message-timestamp-color);
line-height: var(--message-line-height);
unicode-bidi: isolate;
}
.messageTimestamp {
@@ -557,6 +559,7 @@
font-size: var(--message-timestamp-compact-font-size);
font-weight: 500;
line-height: var(--message-line-height);
unicode-bidi: isolate;
}
.messageTimestampHover {
@@ -948,6 +951,7 @@
white-space: nowrap;
max-width: 30%;
vertical-align: baseline;
unicode-bidi: plaintext;
font-weight: 500;
color: color-mix(
in srgb,
@@ -0,0 +1,58 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {readFileSync} from 'node:fs';
import {fileURLToPath} from 'node:url';
import {describe, expect, it} from 'vitest';
const MESSAGE_CSS = readFileSync(fileURLToPath(new URL('./Message.module.css', import.meta.url)), 'utf8');
const ISOLATED_NAMES = ['.messageUsername', '.repliedUsername'];
const ISOLATED_TIMESTAMPS = [
'.messageTimestamp',
'.messageTimestampCompact',
'.messageTimestampHover',
'.messageTimestampCompactHover',
];
function rules(css: string): Array<{selector: string; declarations: string}> {
const source = css.replace(/\/\*[\s\S]*?\*\//g, '');
const out: Array<{selector: string; declarations: string}> = [];
const chain: Array<string> = [];
let buffer = '';
for (const char of source) {
if (char === '{') {
chain.push(buffer.split(/[;}]/).pop()?.trim().replace(/\s+/g, ' ') ?? '');
buffer = '';
} else if (char === '}') {
if (chain.length > 0) {
out.push({selector: chain.join(' '), declarations: buffer});
chain.pop();
}
buffer = '';
} else {
buffer += char;
}
}
return out;
}
function declaredBidi(className: string): Array<string> {
const selectorPattern = new RegExp(`\\${className}(?![\\w-])`);
const out: Array<string> = [];
for (const rule of rules(MESSAGE_CSS)) {
if (!rule.selector.split(',').some((selector) => selectorPattern.test(selector))) continue;
const match = rule.declarations.match(/unicode-bidi\s*:\s*([^;]+)/);
if (match) out.push(match[1].trim());
}
return out;
}
describe('message bidi isolation', () => {
it.each(ISOLATED_NAMES)('resolves %s in its own bidi paragraph', (className) => {
expect(declaredBidi(className)).toContain('plaintext');
});
it.each(ISOLATED_TIMESTAMPS)('keeps %s out of a neighbouring right-to-left run', (className) => {
expect(declaredBidi(className)).toContain('isolate');
});
});
@@ -0,0 +1,23 @@
// @vitest-environment happy-dom
// SPDX-License-Identifier: AGPL-3.0-or-later
import {scopePortalHostToDocument} from '@app/features/ui/overlay/PortalHostContext';
import {describe, expect, it} from 'vitest';
describe('scopePortalHostToDocument', () => {
it('keeps a host that belongs to the scoped document', () => {
const host = document.createElement('div');
expect(scopePortalHostToDocument(host, document)).toBe(host);
});
it('drops a host owned by another document', () => {
const popoutDocument = document.implementation.createHTMLDocument('popout');
const popoutHost = popoutDocument.createElement('div');
expect(scopePortalHostToDocument(popoutHost, document)).toBeNull();
expect(scopePortalHostToDocument(popoutHost, popoutDocument)).toBe(popoutHost);
});
it('passes a missing host through', () => {
expect(scopePortalHostToDocument(null, document)).toBeNull();
});
});
@@ -39,3 +39,8 @@ export function usePortalHost(): PortalHostElement {
export function resolvePortalHost(host: PortalHostElement): HTMLElement {
return host ?? document.body;
}
export function scopePortalHostToDocument(host: PortalHostElement, scopeDocument: Document): PortalHostElement {
if (host == null) return null;
return host.ownerDocument === scopeDocument ? host : null;
}
@@ -4,7 +4,7 @@ import Accessibility from '@app/features/accessibility/state/Accessibility';
import {useAntiShiftFloating} from '@app/features/app/hooks/useAntiShiftFloating';
import {shouldDisableAutofocusOnMobile} from '@app/features/platform/utils/AutofocusUtils';
import * as PopoutCommands from '@app/features/ui/commands/PopoutCommands';
import {usePortalHost} from '@app/features/ui/overlay/PortalHostContext';
import {scopePortalHostToDocument, usePortalHost} from '@app/features/ui/overlay/PortalHostContext';
import {type Popout, PopoutKeyContext, type PopoutReferenceRect} from '@app/features/ui/popover';
import {PopoutResizePositionContext} from '@app/features/ui/popover/PopoutResizePositionContext';
import {getPopoutFocusManagerInsideElements} from '@app/features/ui/popover/PopoverFocusManagerUtils';
@@ -186,6 +186,7 @@ const PopoutItem: React.FC<PopoutItemProps> = observer(
);
const {
ref: popoutRef,
setFloating,
state,
style,
beginManualPositioning,
@@ -206,7 +207,7 @@ const PopoutItem: React.FC<PopoutItemProps> = observer(
useLayoutEffect(() => {
focusRefs.setReference(target);
}, [focusRefs, target]);
const mergedPopoutRef = useMergeRefs([popoutRef, focusRefs.setFloating]);
const mergedPopoutRef = useMergeRefs([setFloating, focusRefs.setFloating]);
const prefersReducedMotion = Accessibility.useReducedMotion;
const [isVisible, setIsVisible] = useState(true);
const [targetInDOM, setTargetInDOM] = useState(() => ownerDocument.contains(target));
@@ -426,10 +427,11 @@ interface PopoutsProps {
export const Popouts: React.FC<PopoutsProps> = observer(({ownerDocument}) => {
const prevPopoutKeysRef = useRef<Set<string>>(new Set());
const portalHost = usePortalHost();
const activePortalHost = usePortalHost();
let scopeDocument = document;
if (portalHost != null) scopeDocument = portalHost.ownerDocument;
if (activePortalHost != null) scopeDocument = activePortalHost.ownerDocument;
if (ownerDocument != null) scopeDocument = ownerDocument;
const portalHost = scopePortalHostToDocument(activePortalHost, scopeDocument);
const popouts = PopoutState.getPopouts(scopeDocument);
const topPopout = popouts.length ? popouts[popouts.length - 1] : null;
const needsBackdrop = Boolean(topPopout && !topPopout.disableBackdrop);
@@ -24,26 +24,30 @@ class UserProfile {
}
handleGatewayReady(): void {
Object.values(this.profileTimeouts).forEach(clearTimeout);
this.profiles = {};
this.profileTimeouts = {};
this.clearAllProfiles();
}
handleProfileInvalidate(userId: string, guildId?: string): void {
const targetGuildId = guildId ?? ME;
this.clearProfileTimeout(userId, targetGuildId);
const userProfiles = this.profiles[userId];
if (!userProfiles) return;
const {[targetGuildId]: _, ...remainingGuildProfiles} = userProfiles;
if (Object.keys(remainingGuildProfiles).length === 0) {
const {[userId]: __, ...remainingProfiles} = this.profiles;
this.profiles = remainingProfiles;
} else {
this.profiles = {
...this.profiles,
[userId]: remainingGuildProfiles,
};
}
this.removeProfile(userId, targetGuildId);
}
handleGuildMemberAdd(userId: string): void {
this.clearUserProfiles(userId);
}
handleGuildMemberRemove(userId: string): void {
this.clearUserProfiles(userId);
}
handleGuildCreate(): void {
this.clearAllProfiles();
}
handleGuildDelete(unavailable?: boolean): void {
if (unavailable) return;
this.clearAllProfiles();
}
handleProfileCreate(profile: Profile): void {
@@ -82,6 +86,37 @@ class UserProfile {
this.profileTimeouts = updatedTimeouts;
}
private clearAllProfiles(): void {
Object.values(this.profileTimeouts).forEach(clearTimeout);
this.profiles = {};
this.profileTimeouts = {};
}
private clearUserProfiles(userId: string): void {
const userProfiles = this.profiles[userId];
if (!userProfiles) return;
for (const guildId of Object.keys(userProfiles)) {
this.clearProfileTimeout(userId, guildId);
}
const {[userId]: _, ...remainingProfiles} = this.profiles;
this.profiles = remainingProfiles;
}
private removeProfile(userId: string, guildId: string): void {
const userProfiles = this.profiles[userId];
if (!userProfiles) return;
const {[guildId]: _, ...remainingGuildProfiles} = userProfiles;
if (Object.keys(remainingGuildProfiles).length === 0) {
const {[userId]: __, ...remainingProfiles} = this.profiles;
this.profiles = remainingProfiles;
} else {
this.profiles = {
...this.profiles,
[userId]: remainingGuildProfiles,
};
}
}
private createTimeoutKey(userId: string, guildId: string): string {
return `${userId}:${guildId}`;
}
@@ -101,23 +136,8 @@ class UserProfile {
this.clearProfileTimeout(userId, guildId);
const timeout = setTimeout(() => {
runInAction(() => {
const userProfiles = this.profiles[userId];
if (!userProfiles) {
const {[timeoutKey]: _, ...remainingTimeouts} = this.profileTimeouts;
this.profileTimeouts = remainingTimeouts;
return;
}
const {[guildId]: _, ...remainingGuildProfiles} = userProfiles;
if (Object.keys(remainingGuildProfiles).length === 0) {
const {[userId]: __, ...remainingProfiles} = this.profiles;
this.profiles = remainingProfiles;
} else {
this.profiles = {
...this.profiles,
[userId]: remainingGuildProfiles,
};
}
const {[timeoutKey]: ___, ...remainingTimeouts} = this.profileTimeouts;
this.removeProfile(userId, guildId);
const {[timeoutKey]: _, ...remainingTimeouts} = this.profileTimeouts;
this.profileTimeouts = remainingTimeouts;
});
}, PROFILE_TIMEOUT_MS);
@@ -0,0 +1,66 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {Profile} from '@app/features/user/models/Profile';
import {afterEach, describe, expect, it, vi} from 'vitest';
vi.mock('@app/features/auth/state/Authentication', () => ({default: {currentUserId: null}}));
const {default: UserProfile} = await import('@app/features/user/state/UserProfile');
const GUILD_ID = '1000';
const TARGET_ID = '2000';
function cache(userId: string = TARGET_ID, guildId?: string): void {
UserProfile.handleProfileCreate({
userId,
guildId: guildId ?? null,
mutualGuilds: [{id: GUILD_ID, nick: null}],
} as unknown as Profile);
}
describe('UserProfile cache invalidation', () => {
afterEach(() => {
UserProfile.handleGatewayReady();
});
it('drops every cached scope for a member removed from a community', () => {
cache();
cache(TARGET_ID, GUILD_ID);
expect(UserProfile.getProfile(TARGET_ID)?.mutualGuilds).toHaveLength(1);
UserProfile.handleGuildMemberRemove(TARGET_ID);
expect(UserProfile.getProfile(TARGET_ID)).toBeNull();
expect(UserProfile.getProfile(TARGET_ID, GUILD_ID)).toBeNull();
});
it('drops every cached scope for a member who joined a community', () => {
cache();
UserProfile.handleGuildMemberAdd(TARGET_ID);
expect(UserProfile.getProfile(TARGET_ID)).toBeNull();
});
it('keeps other users cached when one member is removed', () => {
cache();
cache('3000');
UserProfile.handleGuildMemberRemove(TARGET_ID);
expect(UserProfile.getProfile(TARGET_ID)).toBeNull();
expect(UserProfile.getProfile('3000')).not.toBeNull();
});
it('drops cached profiles when the viewer leaves a community', () => {
cache();
UserProfile.handleGuildDelete(false);
expect(UserProfile.getProfile(TARGET_ID)).toBeNull();
});
it('keeps cached profiles when a community only goes unavailable', () => {
cache();
UserProfile.handleGuildDelete(true);
expect(UserProfile.getProfile(TARGET_ID)).not.toBeNull();
});
it('drops cached profiles when the viewer joins a community', () => {
cache();
UserProfile.handleGuildCreate();
expect(UserProfile.getProfile(TARGET_ID)).toBeNull();
});
});
@@ -4,7 +4,7 @@ title: HTTP API
description: Request format, body representations, shared headers, cross-origin policy, and the error objects.
---
The Fluxer HTTP API is the set of routes a client calls to read and change data. Every route is published under `/v1`, and `/v1` is the only version. The base URL comes from `endpoints.api` in the [instance discovery document](/http-api/instance/#get-instance-discovery), which is served unversioned at `/.well-known/fluxer`.
The Fluxer HTTP API is the set of routes a client calls to read and change data. Every route is published under `/v1`, and `/v1` is the only version. The base URL comes from the [instance discovery document](/http-api/instance/#get-instance-discovery), which is served unversioned at `/.well-known/fluxer`. A third-party client reads `endpoints.api_public` there, and the first-party web application reads `endpoints.api_client`.
Every route is also mounted at the root, so a path resolves with or without the prefix. A client MUST use the `/v1` form. [Download stored object](/http-api/downloads/#download-stored-object) is the one exception, and it resolves at the root alone.
@@ -57,7 +57,7 @@ Each value is an absolute URL supplied by the operator. A value can be a bare or
| gift | string | Gift link base URL |
| webapp | string | Web application base URL |
<sup>1</sup> Both fields are the same configured client API endpoint, while `api_public` is the separately configured public API endpoint
<sup>1</sup> Both fields are the same configured client API endpoint, which the first-party web application reads. `api_public` is the separately configured public API endpoint, which a bot, a library, or any other third-party client reads. A deployment can point the two at one origin, and the instance Fluxer hosts does not
<sup>2</sup> The value is the configured Gateway endpoint and its scheme is `ws` or `wss` as the operator configured it
+3 -1
View File
@@ -43,9 +43,11 @@ A client that knows only a Fluxer origin reads endpoint discovery first. `GET /.
GET https://example.com/.well-known/fluxer
```
The instance Fluxer hosts answers discovery at `https://fluxer.app/.well-known/fluxer`. That origin is the one thing a client is given. Every base URL below it still comes from the response.
It returns the [instance discovery object](/http-api/instance/#instance-discovery-object). Every base URL a client uses comes from the [instance endpoints object](/http-api/instance/#instance-endpoints-object) inside it. A client MUST read every base URL from that response, and it MUST NOT derive one from the origin it was given or assume an official Fluxer domain.
Take `endpoints.api` from that response, then send a credential in the `Authorization` header.
Take the base URL for the kind of client being built, then send a credential in the `Authorization` header. A bot, a library, or any other third-party client takes `endpoints.api_public`. `endpoints.api_client` is the endpoint the first-party web application uses, and `endpoints.api` repeats it.
```text
GET https://api.example.com/v1/users/@me
@@ -17,7 +17,7 @@ By the end you have a Fluxer instance on a hostname you own, with the web app, H
| Requirement | Value |
| --- | --- |
| Host | Any machine that runs Docker with Linux containers. A Linux server, or a Mac or Windows PC with Docker Desktop |
| Runtime | Docker Engine 24 or newer with the Compose plugin v2.20.2 or newer. Docker Desktop ships both |
| Runtime | Docker Engine 24 or newer, or Podman 5 or newer, with the Compose plugin v2.20.2 or newer |
| Processor | Intel, AMD and Apple Silicon all work. Every image ships for `amd64` and `arm64`, and Docker pulls the matching one |
| CPU | 2 vCPU minimum, 4 vCPU for a small active community. No service sets a CPU limit, so every container sees all of them |
| RAM | 8 GB minimum with the shipped defaults, 16 GB for a small active community. 4 GB needs the [low-memory overrides](/operator/configuration/#resources) |
@@ -159,6 +159,7 @@ The run prints one line per phase and ends with the URL to open. It takes severa
| `--update` | Upgrade: record, back up, refresh, pull, recreate, verify |
| `--rollback` | Put back the images and stack files of the newest record |
| `--allow-root` | Permit running as root |
| `--engine <command>` | Container engine to drive. Default `docker`, or `podman` when `docker` is absent |
Run `--dry-run` first to see what a set of flags does. The PowerShell script takes the same flags as `-Domain`, `-Email` and so on, and has no `-AllowRoot`. [Upgrading](/operator/upgrading/#run-the-upgrade) has the flags that only `--update` and `--rollback` read, and the exit codes.
+35 -10
View File
@@ -40,6 +40,7 @@ param(
[string]$Domain = '',
[string]$Email = '',
[string]$Dir = '',
[string]$Engine = '',
[string]$Ref = '',
[string]$ImageTag = 'v1',
[string]$Tls = 'bundled',
@@ -65,6 +66,7 @@ $FluxerRawBase = 'https://raw.githubusercontent.com/fluxerapp/fluxer'
$FluxerStackPath = 'deploy/self-hosting'
$FluxerHealthPath = '/_health'
$FluxerInitService = 'seaweedfs-init'
$FluxerMinimumPodmanVersion = '5.0.0'
$FluxerMinimumEngineVersion = '24.0.0'
$FluxerMinimumComposeVersion = '2.20.2'
$FluxerReadyTimeoutSeconds = 600
@@ -194,6 +196,8 @@ function Show-FluxerUsage {
Write-FluxerLine 'Options:'
Write-FluxerLine ' -Domain <host> Hostname the instance answers on. Prompted when absent.'
Write-FluxerLine ' -Email <address> Address written as FLUXER_VAPID_EMAIL. Prompted when absent.'
Write-FluxerLine ' -Engine <command> Container engine to drive. Default: docker, or podman when'
Write-FluxerLine ' docker is absent.'
Write-FluxerLine ' -Dir <path> Working directory. Default: the fluxer folder in the home'
Write-FluxerLine ' directory, or the working directory itself when -Update or'
Write-FluxerLine ' -Rollback is given and it holds an instance.'
@@ -256,6 +260,9 @@ function Test-FluxerVersionAtLeast([string]$Found, [string]$Minimum) {
# is kept rather than discarded. The two streams stay apart because every caller reads Text as
# data, and one warning line on stderr would be one more line of JSON, one more image reference or
# one more container ID.
$script:FluxerEngine = 'docker'
$script:FluxerEngineLabel = 'Docker Engine'
function Invoke-FluxerCapture([string[]]$CommandArgs) {
$previous = $ErrorActionPreference
$ErrorActionPreference = 'Continue'
@@ -263,7 +270,7 @@ function Invoke-FluxerCapture([string[]]$CommandArgs) {
$err = @()
$code = 0
try {
$output = & docker @CommandArgs 2>&1
$output = & $script:FluxerEngine @CommandArgs 2>&1
$code = $LASTEXITCODE
foreach ($item in @($output)) {
if ($item -is [System.Management.Automation.ErrorRecord]) {
@@ -283,7 +290,7 @@ function Invoke-FluxerDocker([string[]]$CommandArgs) {
$ErrorActionPreference = 'Continue'
$code = 0
try {
& docker @CommandArgs
& $script:FluxerEngine @CommandArgs
$code = $LASTEXITCODE
} finally {
$ErrorActionPreference = $previous
@@ -295,7 +302,7 @@ function Invoke-FluxerDocker([string[]]$CommandArgs) {
# it, so the bytes go to the file through the process itself.
function Invoke-FluxerDockerToFile([string[]]$CommandArgs, [string]$OutFile, [string]$WorkingDir) {
$errorFile = "$OutFile.stderr"
$process = Start-Process -FilePath 'docker' -ArgumentList $CommandArgs -RedirectStandardOutput $OutFile -RedirectStandardError $errorFile -WorkingDirectory $WorkingDir -NoNewWindow -Wait -PassThru
$process = Start-Process -FilePath $script:FluxerEngine -ArgumentList $CommandArgs -RedirectStandardOutput $OutFile -RedirectStandardError $errorFile -WorkingDirectory $WorkingDir -NoNewWindow -Wait -PassThru
$code = $process.ExitCode
if (Test-Path -LiteralPath $errorFile) {
Remove-Item -LiteralPath $errorFile -Force
@@ -330,22 +337,40 @@ function Invoke-FluxerPreflight {
if ($PSVersionTable.PSVersion.Major -lt 6) {
[System.Net.ServicePointManager]::SecurityProtocol = [System.Net.SecurityProtocolType]::Tls12
}
if ($null -eq (Get-Command docker -ErrorAction SilentlyContinue)) {
Stop-Fluxer 'docker is not on PATH. Install Docker Desktop with the WSL2 backend.' $FluxerExitPrerequisite
if ($Engine.Length -gt 0) {
$script:FluxerEngine = $Engine
} elseif ($null -ne (Get-Command docker -ErrorAction SilentlyContinue)) {
$script:FluxerEngine = 'docker'
} elseif ($null -ne (Get-Command podman -ErrorAction SilentlyContinue)) {
$script:FluxerEngine = 'podman'
} else {
Stop-Fluxer 'Neither docker nor podman is on PATH. Install Docker Desktop with the WSL2 backend, or pass -Engine with the command that drives your containers.' $FluxerExitPrerequisite
}
if ($null -eq (Get-Command $script:FluxerEngine -ErrorAction SilentlyContinue)) {
Stop-Fluxer "$($script:FluxerEngine) is not on PATH. Pass -Engine with the command that drives your containers." $FluxerExitPrerequisite
}
# Every engine that speaks the Docker CLI reports itself as "<name> version X",
# and Podman's 5.x is not Docker's 5.x, so the floor has to belong to whichever
# one answered.
$reported = Invoke-FluxerCapture @('--version')
$minimumEngine = $FluxerMinimumEngineVersion
if ($reported.Code -eq 0 -and $reported.Text -match '^\s*podman\s+version') {
$script:FluxerEngineLabel = 'Podman'
$minimumEngine = $FluxerMinimumPodmanVersion
}
$engine = Invoke-FluxerCapture @('version', '--format', '{{.Server.Version}}')
if ($engine.Code -ne 0) {
Stop-Fluxer 'The Docker daemon does not answer. Start Docker Desktop and run this script again.' $FluxerExitPrerequisite
Stop-Fluxer "$($script:FluxerEngine) does not answer. Start it and run this script again." $FluxerExitPrerequisite
}
$compose = Invoke-FluxerCapture @('compose', 'version', '--short')
if ($compose.Code -ne 0) {
Stop-Fluxer 'The Docker Compose v2 plugin is missing. Docker Desktop ships it.' $FluxerExitPrerequisite
}
if (-not (Test-FluxerVersionAtLeast $engine.Text $FluxerMinimumEngineVersion)) {
Stop-Fluxer "Docker Engine $($engine.Text) is older than $FluxerMinimumEngineVersion. Compose is $($compose.Text)." $FluxerExitPrerequisite
if (-not (Test-FluxerVersionAtLeast $engine.Text $minimumEngine)) {
Stop-Fluxer "$($script:FluxerEngineLabel) $($engine.Text) is older than $minimumEngine. Compose is $($compose.Text)." $FluxerExitPrerequisite
}
if (-not (Test-FluxerVersionAtLeast $compose.Text $FluxerMinimumComposeVersion)) {
Stop-Fluxer "Docker Compose $($compose.Text) is older than $FluxerMinimumComposeVersion. Engine is $($engine.Text)." $FluxerExitPrerequisite
Stop-Fluxer "Compose $($compose.Text) is older than $FluxerMinimumComposeVersion. $($script:FluxerEngineLabel) is $($engine.Text)." $FluxerExitPrerequisite
}
$osType = Invoke-FluxerCapture @('info', '--format', '{{.OSType}}')
if ($osType.Code -ne 0) {
@@ -354,7 +379,7 @@ function Invoke-FluxerPreflight {
if ($osType.Text -ne 'linux') {
Stop-Fluxer "Docker runs $($osType.Text) containers. Switch Docker Desktop to Linux containers." $FluxerExitPrerequisite
}
Write-FluxerLine "Docker Engine $($engine.Text) and Compose $($compose.Text) are ready."
Write-FluxerLine "$($script:FluxerEngineLabel) $($engine.Text) and Compose $($compose.Text) are ready."
if (-not (Test-FluxerWsl2)) {
Write-FluxerLine 'WSL 2 is not present. Docker Desktop with the WSL2 backend is the supported Windows setup.'
}
+104 -60
View File
@@ -53,6 +53,10 @@ export LC_ALL
FLUXER_RAW_BASE='https://raw.githubusercontent.com/fluxerapp/fluxer'
FLUXER_STACK_PATH='deploy/self-hosting'
FLUXER_MIN_ENGINE='24.0.0'
# Podman numbers its releases on its own scale, so the Docker Engine floor says
# nothing about it. 5.0 is the first release that runs this stack's compose file
# as written, healthcheck conditions and all.
FLUXER_MIN_PODMAN='5.0.0'
FLUXER_MIN_COMPOSE='2.20.2'
# Both overlays this script downloads use the !override tag, which Compose learned
# in 2.24.4. A stack that loads neither runs on the lower minimum, so the higher
@@ -174,6 +178,8 @@ Usage: sh install.sh --domain <host> --email <address> [options]
Options:
--domain <host> Hostname the instance answers on. Prompted when absent.
--email <address> Address the operator reads. Prompted when absent.
--engine <command> Container engine to drive. Default docker, or podman
when docker is absent.
--dir <path> Working directory. Default ~/fluxer, or the working
directory itself when --update or --rollback is given
and it holds an instance.
@@ -184,7 +190,7 @@ Options:
--edge-bind <addr:port> Plain HTTP bind under --tls proxy. Default 127.0.0.1:8080.
--non-interactive Never prompt. A missing required value is an error.
--dry-run Print the plan. Change nothing.
--no-start Write everything and skip docker compose up -d.
--no-start Write everything and skip $fluxer_engine compose up -d.
--update Upgrade: record, back up, refresh, pull, recreate, verify.
--rollback Restore the images and stack files of the last record.
--backup-dir <path> Where records go. Default <dir>/backups.
@@ -247,6 +253,9 @@ trap fluxer_on_signal TERM
opt_domain=''
opt_email=''
opt_dir=''
opt_engine=''
fluxer_engine='docker'
fluxer_engine_label='Docker Engine'
opt_ref=''
opt_image_tag='v1'
opt_tls='bundled'
@@ -273,6 +282,11 @@ while [ $# -gt 0 ]; do
opt_email=$2
shift 2
;;
--engine)
fluxer_take_value '--engine' $# "${2:-}"
opt_engine=$2
shift 2
;;
--dir)
fluxer_take_value '--dir' $# "${2:-}"
opt_dir=$2
@@ -467,29 +481,59 @@ fluxer_version_ge() {
fluxer_engine_version=''
fluxer_compose_version=''
# Every engine that speaks the Docker CLI reports itself as "<name> version X".
# Docker says "Docker version 28.0.0, build ...", Podman says "podman version
# 5.8.4", and the podman-docker shim answers as Podman under the name docker.
# Reading the name as well as the number is what lets the floor below belong to
# the engine actually in use, because Podman's 5.x is not Docker's 5.x.
fluxer_resolve_engine() {
if [ -n "$opt_engine" ]; then
fluxer_engine=$opt_engine
elif command -v docker >/dev/null 2>&1; then
fluxer_engine='docker'
elif command -v podman >/dev/null 2>&1; then
fluxer_engine='podman'
else
fluxer_fail 2 "Neither docker nor podman is installed. $(fluxer_docker_hint)"
fi
if ! command -v "$fluxer_engine" >/dev/null 2>&1; then
fluxer_fail 2 "$fluxer_engine is not installed. Pass --engine with the command that drives your containers."
fi
}
fluxer_preflight() {
if [ "$opt_allow_root" -eq 0 ] && [ "$(id -u)" -eq 0 ]; then
fluxer_fail 2 'Running as root. Use an account in the docker group, or pass --allow-root.'
fi
if ! command -v docker >/dev/null 2>&1; then
fluxer_fail 2 "Docker is not installed. $(fluxer_docker_hint)"
fluxer_resolve_engine
if ! $fluxer_engine compose version >/dev/null 2>&1; then
fluxer_fail 2 "$fluxer_engine has no compose subcommand. $(fluxer_docker_hint)"
fi
if ! docker compose version >/dev/null 2>&1; then
fluxer_fail 2 "The Docker Compose v2 plugin is missing. $(fluxer_docker_hint)"
fi
fluxer_engine_version=$(docker --version 2>/dev/null | sed -n 's/^Docker version \([0-9][0-9.]*\).*/\1/p')
fluxer_compose_version=$(docker compose version --short 2>/dev/null | sed -n 's/^v\{0,1\}\([0-9][0-9.]*\).*/\1/p')
fluxer_engine_report=$($fluxer_engine --version 2>/dev/null)
fluxer_engine_kind=$(printf '%s\n' "$fluxer_engine_report" | sed -n 's/^\([A-Za-z][A-Za-z]*\) version .*/\1/p' | tr 'A-Z' 'a-z')
fluxer_engine_version=$(printf '%s\n' "$fluxer_engine_report" | sed -n 's/^[A-Za-z][A-Za-z]* version v\{0,1\}\([0-9][0-9.]*\).*/\1/p')
fluxer_compose_version=$($fluxer_engine compose version --short 2>/dev/null | sed -n 's/^v\{0,1\}\([0-9][0-9.]*\).*/\1/p')
if [ -z "$fluxer_engine_version" ] || [ -z "$fluxer_compose_version" ]; then
fluxer_fail 2 'Cannot read the Docker Engine and Compose versions.'
fluxer_fail 2 "Cannot read the engine and Compose versions from $fluxer_engine. It answered \"$fluxer_engine_report\" to --version."
fi
if ! fluxer_version_ge "$fluxer_engine_version" "$FLUXER_MIN_ENGINE"; then
fluxer_fail 2 "Docker Engine $fluxer_engine_version with Compose $fluxer_compose_version. Fluxer needs Engine $FLUXER_MIN_ENGINE or newer."
case $fluxer_engine_kind in
podman)
fluxer_engine_label='Podman'
fluxer_min_engine=$FLUXER_MIN_PODMAN
;;
*)
fluxer_engine_label='Docker Engine'
fluxer_min_engine=$FLUXER_MIN_ENGINE
;;
esac
if ! fluxer_version_ge "$fluxer_engine_version" "$fluxer_min_engine"; then
fluxer_fail 2 "$fluxer_engine_label $fluxer_engine_version with Compose $fluxer_compose_version. Fluxer needs $fluxer_engine_label $fluxer_min_engine or newer."
fi
if ! fluxer_version_ge "$fluxer_compose_version" "$FLUXER_MIN_COMPOSE"; then
fluxer_fail 2 "Docker Engine $fluxer_engine_version with Compose $fluxer_compose_version. Fluxer needs Compose $FLUXER_MIN_COMPOSE or newer."
fluxer_fail 2 "$fluxer_engine_label $fluxer_engine_version with Compose $fluxer_compose_version. Fluxer needs Compose $FLUXER_MIN_COMPOSE or newer."
fi
if ! docker info >/dev/null 2>&1; then
fluxer_fail 2 'The Docker daemon does not answer. Start Docker and run this again.'
if ! $fluxer_engine info >/dev/null 2>&1; then
fluxer_fail 2 "$fluxer_engine does not answer. Start it and run this again."
fi
if ! command -v curl >/dev/null 2>&1; then
fluxer_fail 2 "curl is not installed. $(fluxer_host_tool_hint curl)"
@@ -497,7 +541,7 @@ fluxer_preflight() {
if ! command -v openssl >/dev/null 2>&1; then
fluxer_fail 2 "openssl is not installed. $(fluxer_host_tool_hint openssl)"
fi
fluxer_say "Docker Engine $fluxer_engine_version with Compose $fluxer_compose_version."
fluxer_say "$fluxer_engine_label $fluxer_engine_version with Compose $fluxer_compose_version."
}
fluxer_valid_domain() {
@@ -735,9 +779,9 @@ fluxer_check_volumes() {
if [ -z "$fluxer_project" ]; then
return 0
fi
docker volume ls -q --filter "label=com.docker.compose.project=$fluxer_project" > "$fluxer_scratch/volumes" 2>/dev/null || return 0
$fluxer_engine volume ls -q --filter "label=com.docker.compose.project=$fluxer_project" > "$fluxer_scratch/volumes" 2>/dev/null || return 0
if [ -s "$fluxer_scratch/volumes" ]; then
fluxer_fail 3 "Docker already holds volumes for the $fluxer_project project. They hold the secrets of an earlier install, and this .env does not open them. Run install.sh --update in the directory that holds that instance, or remove the volumes with docker volume rm before you install again."
fluxer_fail 3 "Docker already holds volumes for the $fluxer_project project. They hold the secrets of an earlier install, and this .env does not open them. Run install.sh --update in the directory that holds that instance, or remove the volumes with $fluxer_engine volume rm before you install again."
fi
}
@@ -845,9 +889,9 @@ fluxer_write_env() {
}
fluxer_stack_ready() {
docker compose ps -aq > "$fluxer_scratch/ids" 2>/dev/null || return 1
$fluxer_engine compose ps -aq > "$fluxer_scratch/ids" 2>/dev/null || return 1
[ -s "$fluxer_scratch/ids" ] || return 1
xargs docker inspect --format "$FLUXER_INSPECT_FORMAT" < "$fluxer_scratch/ids" > "$fluxer_scratch/state" 2>/dev/null || return 1
xargs $fluxer_engine inspect --format "$FLUXER_INSPECT_FORMAT" < "$fluxer_scratch/ids" > "$fluxer_scratch/state" 2>/dev/null || return 1
fluxer_ready=1
fluxer_init_done=0
while read -r fluxer_service fluxer_status fluxer_health fluxer_code; do
@@ -1078,11 +1122,11 @@ fluxer_require_compose_files() {
fluxer_compose_count=$((${fluxer_compose_count:-0} + 1))
[ ! -e "$fluxer_path" ] || continue
if fluxer_stack_files | grep -qxF "$fluxer_name"; then
fluxer_fail 2 "COMPOSE_FILE from $fluxer_compose_from names $fluxer_name and $fluxer_path is not there, so every docker compose command in $opt_dir fails and this run stops before it changes anything. This script downloads $fluxer_name, and an instance set up before it existed does not hold that file yet. Put it in place and run this again:
fluxer_fail 2 "COMPOSE_FILE from $fluxer_compose_from names $fluxer_name and $fluxer_path is not there, so every $fluxer_engine compose command in $opt_dir fails and this run stops before it changes anything. This script downloads $fluxer_name, and an instance set up before it existed does not hold that file yet. Put it in place and run this again:
curl -fsSL --proto '=https' --tlsv1.2 -o $fluxer_path $FLUXER_RAW_BASE/$opt_ref/$FLUXER_STACK_PATH/$fluxer_name
Leave the COMPOSE_FILE line as it is. Without $fluxer_name the edge container binds 80 and 443 and requests its own certificate."
fi
fluxer_fail 2 "COMPOSE_FILE from $fluxer_compose_from names $fluxer_name and $fluxer_path is not there, so every docker compose command in $opt_dir fails. This script does not download $fluxer_name. Put that file back, or take it out of the COMPOSE_FILE line."
fluxer_fail 2 "COMPOSE_FILE from $fluxer_compose_from names $fluxer_name and $fluxer_path is not there, so every $fluxer_engine compose command in $opt_dir fails. This script does not download $fluxer_name. Put that file back, or take it out of the COMPOSE_FILE line."
done
if [ "$fluxer_compose_count" -gt 1 ] &&
! fluxer_version_ge "$fluxer_compose_version" "$FLUXER_MIN_COMPOSE_OVERLAY"; then
@@ -1102,7 +1146,7 @@ Leave the COMPOSE_FILE line as it is. Without $fluxer_name the edge container bi
# command.
fluxer_compose_query() {
fluxer_query_status=0
docker compose config "$1" > "$fluxer_scratch/compose-out" 2> "$fluxer_scratch/compose-err" || fluxer_query_status=$?
$fluxer_engine compose config "$1" > "$fluxer_scratch/compose-out" 2> "$fluxer_scratch/compose-err" || fluxer_query_status=$?
return "$fluxer_query_status"
}
@@ -1151,14 +1195,14 @@ fluxer_indent_file() {
}
fluxer_running_image_ids() {
if ! docker compose ps -aq > "$fluxer_scratch/containers" 2> "$fluxer_scratch/ps-err"; then
fluxer_fail 2 "docker compose ps failed in $opt_dir, so what is running cannot be read and the record would name no image ID at all. Compose printed:
if ! $fluxer_engine compose ps -aq > "$fluxer_scratch/containers" 2> "$fluxer_scratch/ps-err"; then
fluxer_fail 2 "$fluxer_engine compose ps failed in $opt_dir, so what is running cannot be read and the record would name no image ID at all. Compose printed:
$(fluxer_indent_file "$fluxer_scratch/ps-err" ' ')"
fi
[ -s "$fluxer_scratch/containers" ] || return 0
if ! xargs docker inspect --format '{{.Config.Image}} {{.Image}}' \
if ! xargs $fluxer_engine inspect --format '{{.Config.Image}} {{.Image}}' \
< "$fluxer_scratch/containers" > "$fluxer_scratch/inspected" 2> "$fluxer_scratch/inspect-err"; then
fluxer_fail 2 "docker inspect failed in $opt_dir, so the image ID under each reference cannot be read and a rollback would have nothing to go back to. Docker printed:
fluxer_fail 2 "$fluxer_engine inspect failed in $opt_dir, so the image ID under each reference cannot be read and a rollback would have nothing to go back to. Docker printed:
$(fluxer_indent_file "$fluxer_scratch/inspect-err" ' ')"
fi
sort -u "$fluxer_scratch/inspected"
@@ -1187,11 +1231,11 @@ fluxer_recorded_id_for() {
fluxer_record_state() {
fluxer_running_image_ids > "$fluxer_scratch/running"
if ! fluxer_compose_images > "$fluxer_scratch/refs"; then
fluxer_fail 2 "docker compose config --images failed in $opt_dir, so the running version cannot be recorded. Compose printed:
fluxer_fail 2 "$fluxer_engine compose config --images failed in $opt_dir, so the running version cannot be recorded. Compose printed:
$(fluxer_compose_error ' ')"
fi
if [ ! -s "$fluxer_scratch/refs" ]; then
fluxer_fail 2 "docker compose config --images returned nothing in $opt_dir, so the running version cannot be recorded. The stack files there declare no service with an image."
fluxer_fail 2 "$fluxer_engine compose config --images returned nothing in $opt_dir, so the running version cannot be recorded. The stack files there declare no service with an image."
fi
fluxer_prepare_record
printf '%s\n' "$(fluxer_env_value FLUXER_IMAGE_TAG)" > "$fluxer_record/$FLUXER_TAG_FILE"
@@ -1225,9 +1269,9 @@ fluxer_save_current_files() {
}
fluxer_postgres_running() {
fluxer_pg_id=$(docker compose ps -q postgres 2>/dev/null | head -n 1)
fluxer_pg_id=$($fluxer_engine compose ps -q postgres 2>/dev/null | head -n 1)
[ -n "$fluxer_pg_id" ] || return 1
[ "$(docker inspect --format '{{.State.Status}}' "$fluxer_pg_id" 2>/dev/null)" = 'running' ]
[ "$($fluxer_engine inspect --format '{{.State.Status}}' "$fluxer_pg_id" 2>/dev/null)" = 'running' ]
}
# Step 2 of an upgrade, and the part that gates the rest.
@@ -1255,13 +1299,13 @@ fluxer_postgres_running() {
fluxer_dump_postgres() {
if ! fluxer_postgres_running; then
fluxer_say 'Postgres is not running. Starting it for the dump.'
if ! docker compose up -d --wait postgres; then
if ! $fluxer_engine compose up -d --wait postgres; then
fluxer_fail 7 "Postgres does not start in $opt_dir, so no dump can be taken."
fi
fi
fluxer_dump_path="$fluxer_record/$FLUXER_DUMP_FILE"
fluxer_say 'Dumping the database.'
if ! docker compose exec -T postgres pg_dump -U fluxer -d fluxer --format=custom > "$fluxer_dump_path"; then
if ! $fluxer_engine compose exec -T postgres pg_dump -U fluxer -d fluxer --format=custom > "$fluxer_dump_path"; then
rm -f "$fluxer_dump_path"
fluxer_fail 7 'pg_dump failed. The instance is untouched.'
fi
@@ -1282,7 +1326,7 @@ fluxer_volume_error() {
}
fluxer_volume_size_kb() {
docker run --rm -v "$1:/data:ro" "$FLUXER_HELPER_IMAGE" du -sk /data 2> "$fluxer_scratch/volume-err" |
$fluxer_engine run --rm -v "$1:/data:ro" "$FLUXER_HELPER_IMAGE" du -sk /data 2> "$fluxer_scratch/volume-err" |
awk 'NR==1 {print $1}'
}
@@ -1307,7 +1351,7 @@ fluxer_copy_volumes() {
while read -r fluxer_volume; do
[ -n "$fluxer_volume" ] || continue
fluxer_full="${fluxer_project}_${fluxer_volume}"
if ! docker volume inspect "$fluxer_full" >/dev/null 2>&1; then
if ! $fluxer_engine volume inspect "$fluxer_full" >/dev/null 2>&1; then
fluxer_fail 7 "The volume $fluxer_full does not exist, so the uploads cannot be copied. The stack declares that volume, so this is a project name other than $fluxer_project or a volume that was removed. Pass --no-volume-backup to take the database dump alone when the uploads of this instance live somewhere this script cannot reach."
fi
fluxer_size=$(fluxer_volume_size_kb "$fluxer_full")
@@ -1331,21 +1375,21 @@ $(fluxer_volume_error ' ')" ;;
return 0
fi
fluxer_say 'Stopping the stack for a consistent copy of the uploads.'
if ! docker compose stop; then
fluxer_fail 7 "docker compose stop failed in $opt_dir."
if ! $fluxer_engine compose stop; then
fluxer_fail 7 "$fluxer_engine compose stop failed in $opt_dir."
fi
while read -r fluxer_volume; do
[ -n "$fluxer_volume" ] || continue
fluxer_full="${fluxer_project}_${fluxer_volume}"
fluxer_say "Copying $fluxer_full."
if ! docker run --rm -v "$fluxer_full:/data:ro" -v "$fluxer_record:/backup" "$FLUXER_HELPER_IMAGE" tar czf "/backup/$fluxer_volume.tgz" -C /data .; then
docker compose up -d --remove-orphans || true
if ! $fluxer_engine run --rm -v "$fluxer_full:/data:ro" -v "$fluxer_record:/backup" "$FLUXER_HELPER_IMAGE" tar czf "/backup/$fluxer_volume.tgz" -C /data .; then
$fluxer_engine compose up -d --remove-orphans || true
fluxer_fail 7 "Copying $fluxer_full failed. The stack is started again on the images it was running."
fi
done < "$fluxer_scratch/copy-volumes"
fluxer_say 'Starting the stack again before the upgrade continues.'
if ! docker compose up -d --remove-orphans; then
fluxer_fail 7 "docker compose up -d failed in $opt_dir after the copy. Read docker compose logs there."
if ! $fluxer_engine compose up -d --remove-orphans; then
fluxer_fail 7 "$fluxer_engine compose up -d failed in $opt_dir after the copy. Read $fluxer_engine compose logs there."
fi
}
@@ -1420,7 +1464,7 @@ fluxer_compose_services() {
fluxer_restart_mounted() {
[ -s "$fluxer_scratch/restart" ] || return 0
if ! fluxer_compose_services > "$fluxer_scratch/services"; then
fluxer_fail 6 "docker compose config --services failed in $opt_dir, so the services that mount a refreshed file cannot be restarted. Compose printed:
fluxer_fail 6 "$fluxer_engine compose config --services failed in $opt_dir, so the services that mount a refreshed file cannot be restarted. Compose printed:
$(fluxer_compose_error ' ')"
fi
while read -r fluxer_file fluxer_service; do
@@ -1430,8 +1474,8 @@ $(fluxer_compose_error ' ')"
continue
fi
fluxer_say "Restarting $fluxer_service, because $fluxer_file is mounted into it and up -d does not reload a mounted file."
if ! docker compose restart "$fluxer_service"; then
fluxer_fail 6 "docker compose restart $fluxer_service failed in $opt_dir."
if ! $fluxer_engine compose restart "$fluxer_service"; then
fluxer_fail 6 "$fluxer_engine compose restart $fluxer_service failed in $opt_dir."
fi
done < "$fluxer_scratch/restart"
}
@@ -1464,7 +1508,7 @@ fluxer_newest_record() {
fluxer_verify_stack() {
fluxer_say 'Waiting for every service to report ready.'
if ! fluxer_wait_ready; then
fluxer_fail 6 "The stack is not ready after $FLUXER_READY_TIMEOUT seconds. Read docker compose logs in $opt_dir."
fluxer_fail 6 "The stack is not ready after $FLUXER_READY_TIMEOUT seconds. Read $fluxer_engine compose logs in $opt_dir."
fi
fluxer_domain_value=$(fluxer_env_value FLUXER_DOMAIN)
if [ -n "$fluxer_domain_value" ]; then
@@ -1499,7 +1543,7 @@ fluxer_plan_update() {
fi
fluxer_running_image_ids > "$fluxer_scratch/running"
if ! fluxer_compose_images > "$fluxer_scratch/refs"; then
fluxer_say ' refusal docker compose config --images fails here, and step 1 of the upgrade reads that list'
fluxer_say ' refusal $fluxer_engine compose config --images fails here, and step 1 of the upgrade reads that list'
fluxer_say ' compose said'
fluxer_compose_error ' '
if [ -s "$fluxer_scratch/plan-secrets" ]; then
@@ -1560,7 +1604,7 @@ fluxer_plan_update() {
fluxer_say " restart $fluxer_service, because $fluxer_file changes and a mounted file survives up -d"
done < "$fluxer_scratch/restart"
fi
fluxer_say ' commands docker compose pull, docker compose up -d'
fluxer_say ' commands $fluxer_engine compose pull, $fluxer_engine compose up -d'
fluxer_plan_footer
fluxer_say 'Drop --dry-run to run this.'
}
@@ -1586,7 +1630,7 @@ fluxer_plan_rollback() {
[ -n "$fluxer_ref" ] || continue
if [ "$fluxer_id" = '-' ]; then
fluxer_say " $fluxer_ref was not recorded with an ID"
elif docker image inspect --format '{{.Id}}' "$fluxer_id" >/dev/null 2>&1; then
elif $fluxer_engine image inspect --format '{{.Id}}' "$fluxer_id" >/dev/null 2>&1; then
fluxer_say " $fluxer_ref back to $fluxer_id"
else
fluxer_say " $fluxer_ref is gone from this host, so $fluxer_id cannot come back"
@@ -1616,8 +1660,8 @@ fluxer_run_update() {
# By hand:
# docker compose pull
fluxer_say 'Pulling images.'
if ! docker compose pull; then
fluxer_fail 4 "docker compose pull failed in $opt_dir. The stack files are refreshed and the instance still runs the old images."
if ! $fluxer_engine compose pull; then
fluxer_fail 4 "$fluxer_engine compose pull failed in $opt_dir. The stack files are refreshed and the instance still runs the old images."
fi
# Recreates the containers whose image or configuration changed and leaves
# the rest running.
@@ -1639,8 +1683,8 @@ fluxer_run_update() {
# Gateway, so the hostname returns errors for a minute or two after this
# call. The api healthcheck allows 90 seconds before it counts a failure.
fluxer_say 'Recreating the stack.'
if ! docker compose up -d --remove-orphans; then
fluxer_fail 6 "docker compose up -d failed in $opt_dir. Read docker compose logs there."
if ! $fluxer_engine compose up -d --remove-orphans; then
fluxer_fail 6 "$fluxer_engine compose up -d failed in $opt_dir. Read $fluxer_engine compose logs there."
fi
fluxer_restart_mounted
fluxer_verify_stack
@@ -1721,18 +1765,18 @@ fluxer_run_rollback() {
while read -r fluxer_ref fluxer_id; do
[ -n "$fluxer_ref" ] || continue
[ "$fluxer_id" != '-' ] || continue
if ! docker image inspect --format '{{.Id}}' "$fluxer_id" >/dev/null 2>&1; then
if ! $fluxer_engine image inspect --format '{{.Id}}' "$fluxer_id" >/dev/null 2>&1; then
fluxer_say "$fluxer_ref is gone from this host, so it keeps the image it has now."
continue
fi
if ! docker image tag "$fluxer_id" "$fluxer_ref"; then
if ! $fluxer_engine image tag "$fluxer_id" "$fluxer_ref"; then
fluxer_fail 3 "Cannot put $fluxer_id back on $fluxer_ref."
fi
fluxer_moved=1
done < "$fluxer_rollback_dir/$FLUXER_IMAGES_FILE"
fi
if [ "$fluxer_moved" -eq 0 ]; then
fluxer_fail 2 "Nothing in $fluxer_rollback_dir can be put back. The recorded tag is the one in .env and every recorded image has been removed from this host, which a docker image prune does."
fluxer_fail 2 "Nothing in $fluxer_rollback_dir can be put back. The recorded tag is the one in .env and every recorded image has been removed from this host, which a $fluxer_engine image prune does."
fi
while read -r fluxer_file; do
@@ -1746,8 +1790,8 @@ fluxer_run_rollback() {
fluxer_say "Stack files in $opt_dir are the ones the record holds."
fluxer_say 'Recreating the stack.'
if ! docker compose up -d --remove-orphans; then
fluxer_fail 6 "docker compose up -d failed in $opt_dir. Read docker compose logs there."
if ! $fluxer_engine compose up -d --remove-orphans; then
fluxer_fail 6 "$fluxer_engine compose up -d failed in $opt_dir. Read $fluxer_engine compose logs there."
fi
# The mounted file came back from the record, so its service restarts.
# Comparing it first would save one restart and cost the reader a reason.
@@ -1820,19 +1864,19 @@ fluxer_write_env
fluxer_say "Wrote $opt_dir/.env, readable by you alone."
if [ "$opt_no_start" -eq 1 ]; then
fluxer_say "Start the instance with docker compose up -d in $opt_dir."
fluxer_say "Start the instance with $fluxer_engine compose up -d in $opt_dir."
fluxer_say 'Open it and create the first admin account. Finish the setup wizard in the same sitting.'
fluxer_say "Secrets live in $opt_dir/.env. Back that file up."
exit 0
fi
fluxer_say 'Starting the stack.'
if ! docker compose up -d; then
fluxer_fail 6 "docker compose up -d failed in $opt_dir. Read docker compose logs there."
if ! $fluxer_engine compose up -d; then
fluxer_fail 6 "$fluxer_engine compose up -d failed in $opt_dir. Read $fluxer_engine compose logs there."
fi
fluxer_say 'Waiting for every service to report ready. This takes several minutes on the first start, which pulls eighteen images.'
if ! fluxer_wait_ready; then
fluxer_fail 6 "The stack is not ready after $FLUXER_READY_TIMEOUT seconds. Read docker compose logs in $opt_dir."
fluxer_fail 6 "The stack is not ready after $FLUXER_READY_TIMEOUT seconds. Read $fluxer_engine compose logs in $opt_dir."
fi
fluxer_probe "$opt_domain"
+26 -4
View File
@@ -1,7 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use crate::ast::{AlertType, ListItem, Node, ParserFlags, TableAlignment};
use crate::constants::{MAX_AST_NODES, MAX_LINE_LENGTH};
use crate::constants::{CODE_FENCE_LENGTH, MAX_AST_NODES, MAX_LINE_LENGTH};
use crate::links::{has_open_inline_code, has_valid_code_fence_language};
use crate::normalize::{normalize_nodes, replace_trailing_whitespace_with_newline};
use crate::parser::{MarkdownParser, ParseError, RuntimeState};
@@ -157,7 +157,7 @@ pub(crate) fn parse_block(state: &mut RuntimeState<'_>) -> Result<BlockParseResu
}
if ParserFlags::has(state.parser.flags(), ParserFlags::ALLOW_CODE_BLOCKS) {
if let Some(fence_pos) = line.find("```") {
if let Some(fence_pos) = find_unescaped_fence(&line) {
let starts_with_fence =
starts_with(trimmed, "```") && fence_pos == line.len() - trimmed.len();
if starts_with_fence {
@@ -459,7 +459,7 @@ fn parse_code_block(lines: &[Line], current: usize) -> Result<Option<CodeBlockRe
while fence_length < trimmed_line.len() && byte_at(trimmed_line, fence_length) == b'`' {
fence_length += 1;
}
if fence_length < 3 {
if fence_length < CODE_FENCE_LENGTH {
return Ok(None);
}
let language_part = &trimmed_line[fence_length..];
@@ -1197,7 +1197,7 @@ pub(crate) fn opens_code_block_midline(lines: &[Line], index: usize, flags: u32)
return false;
}
let line = &lines[index].text;
let Some(fence_pos) = line.find("```") else {
let Some(fence_pos) = find_unescaped_fence(line) else {
return false;
};
let trimmed = trim_start(line);
@@ -1267,6 +1267,28 @@ fn alert_type(label: &str) -> Option<AlertType> {
}
}
fn find_unescaped_fence(line: &str) -> Option<usize> {
let mut search = 0usize;
while let Some(offset) = line[search..].find("```") {
let fence_pos = search + offset;
if !is_escaped_fence(line, fence_pos) {
return Some(fence_pos);
}
search = fence_pos + 1;
}
None
}
fn is_escaped_fence(line: &str, fence_pos: usize) -> bool {
let mut backslashes = 0usize;
let mut index = fence_pos;
while index > 0 && byte_at(line, index - 1) == b'\\' {
backslashes += 1;
index -= 1;
}
backslashes % 2 == 1
}
fn slice_lines_from_fence(lines: &[Line], current: usize, fence_pos: usize) -> Vec<Line> {
let mut out = Vec::with_capacity(lines.len() - current);
out.push(Line {
@@ -1,5 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
pub const CODE_FENCE_LENGTH: usize = 3;
pub const MAX_AST_NODES: usize = 10_000;
pub const MAX_INLINE_DEPTH: usize = 10;
pub const MAX_LINES: usize = 10_000;
+17 -1
View File
@@ -1,7 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use crate::ast::{EmojiKind, Node, ParserFlags, ParserResult};
use crate::constants::{MAX_INLINE_DEPTH, MAX_LINE_LENGTH};
use crate::constants::{CODE_FENCE_LENGTH, MAX_INLINE_DEPTH, MAX_LINE_LENGTH};
use crate::links;
use crate::normalize::{
combine_adjacent_text, compact_empty_text_nodes, flatten_top_level_formatting,
@@ -125,6 +125,14 @@ fn parse_inline_with_context(
position += 1 + emoji.len;
continue;
}
if next == b'`' {
let run = backtick_run_length(text, position + 1);
if run >= CODE_FENCE_LENGTH {
accumulated.extend(std::iter::repeat_n(b'`', run));
position += 1 + run;
continue;
}
}
if is_escapable_character(next)
|| is_ordered_list_marker_dot_escape(text, position)
|| is_word_dot_escape(text, position)
@@ -349,6 +357,14 @@ fn regional_indicator_letter(ch: char) -> Option<char> {
char::from_u32(u32::from(b'a') + codepoint - 0x1f1e6)
}
fn backtick_run_length(text: &str, start: usize) -> usize {
let mut run = 0usize;
while start + run < text.len() && byte_at(text, start + run) == b'`' {
run += 1;
}
run
}
fn is_ordered_list_marker_dot_escape(text: &str, backslash_position: usize) -> bool {
if backslash_position + 2 >= text.len()
|| byte_at(text, backslash_position + 1) != b'.'
@@ -202,3 +202,46 @@ fn unterminated_midline_fence_after_a_line_stays_text() {
json!([{"type": "Text", "content": "hello\nfoo```bar with no closing fence"}])
);
}
#[test]
fn escaped_fence_stays_text_on_one_line() {
assert_eq!(
parse("\\```hello```"),
json!([{"type": "Text", "content": "```hello```"}])
);
}
#[test]
fn escaped_fence_stays_text_across_lines() {
assert_eq!(
parse("\\```rust\nfn main() {}\n```"),
json!([{"type": "Text", "content": "```rust\nfn main() {}\n```"}])
);
}
#[test]
fn escaped_midline_fence_stays_text() {
assert_eq!(
parse("label\\```rust\nfn main() {}\n```"),
json!([{"type": "Text", "content": "label```rust\nfn main() {}\n```"}])
);
}
#[test]
fn escaped_backslash_before_fence_still_opens_a_code_block() {
assert_eq!(
parse("\\\\```hello```"),
json!([
{"type": "Text", "content": "\\"},
{"type": "CodeBlock", "content": "hello"}
])
);
}
#[test]
fn longer_escaped_fence_stays_text() {
assert_eq!(
parse("\\````hello````"),
json!([{"type": "Text", "content": "````hello````"}])
);
}