fix(app): refetch the tail when a channel window falls behind (#2565)

This commit is contained in:
Hampus
2026-09-06 23:51:08 +02:00
committed by GitHub
parent f00c6ee47a
commit 977b6767cd
10 changed files with 475 additions and 27 deletions
@@ -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 {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();