Preserve ordered huddle TTS reconnect replay

Signed-off-by: John Tennant <jtennant@squareup.com>
This commit is contained in:
John Tennant
2026-07-29 12:20:49 -04:00
committed by John Tennant
parent c3e03bb84d
commit ed87750b40
6 changed files with 63 additions and 21 deletions
+4 -4
View File
@@ -195,8 +195,8 @@ pub(crate) fn load_from_path(path: &Path) -> Result<TtsSettings, String> {
let value: serde_json::Value = serde_json::from_slice(&bytes)
.map_err(|error| format!("text-to-speech settings are not valid JSON: {error}"))?;
// Pre-V1 experiment builds stored an unversioned, incompatible shape.
// Migrate it to deterministic V1 defaults instead of carrying it over.
// Unversioned settings are incompatible with the V1 schema. Use
// deterministic V1 defaults rather than interpreting ambiguous fields.
if value.get("version").is_none() {
return Ok(TtsSettings::default());
}
@@ -211,8 +211,8 @@ pub(crate) fn load_from_path(path: &Path) -> Result<TtsSettings, String> {
));
}
// Early experiments used one bare Pocket `voiceId`. Preserve the toggle
// and qualify that value into the ordered cross-backend preference schema.
// Legacy V1 settings may contain one bare Pocket `voiceId`. Preserve the
// toggle and qualify it into the ordered cross-backend preference schema.
if value.get("voicePreferences").is_none() {
let legacy_voice = value
.get("voiceId")
@@ -132,21 +132,25 @@ export function useTtsSubscription(
const seenOrder: string[] = [];
const MAX_SEEN_EVENTS = 5000;
relayClient
.subscribeLive(buildHuddleTtsLiveFilter(ephemeralChannelId), (event) => {
if (disposed) return;
// Dedup by event ID (covers reconnect replay).
if (seenEventIds.has(event.id)) return;
seenEventIds.add(event.id);
seenOrder.push(event.id);
if (seenOrder.length > MAX_SEEN_EVENTS) {
const oldest = seenOrder.shift();
if (oldest !== undefined) seenEventIds.delete(oldest);
}
.subscribeLive(
buildHuddleTtsLiveFilter(ephemeralChannelId),
(event) => {
if (disposed) return;
// Dedup by event ID (covers reconnect replay).
if (seenEventIds.has(event.id)) return;
seenEventIds.add(event.id);
seenOrder.push(event.id);
if (seenOrder.length > MAX_SEEN_EVENTS) {
const oldest = seenOrder.shift();
if (oldest !== undefined) seenEventIds.delete(oldest);
}
// Preserve arrival order while the initial authoritative membership
// lookup is pending. A failed lookup clears this buffer fail-closed.
initialMembershipGate.push(event);
})
// Preserve arrival order while the initial authoritative membership
// lookup is pending. A failed lookup clears this buffer fail-closed.
initialMembershipGate.push(event);
},
{ replayMissedHistory: true },
)
.then((dispose) => {
if (disposed) {
void dispose();
+4 -1
View File
@@ -445,8 +445,9 @@ export class RelayClient {
async subscribeLive(
filter: RelaySubscriptionFilter,
onEvent: (event: RelayEvent) => void,
options?: { replayMissedHistory?: boolean },
) {
return this.subscribe(filter, onEvent);
return this.subscribe(filter, onEvent, options);
}
async subscribeToChannelMentionEvents(
@@ -600,6 +601,7 @@ export class RelayClient {
private async subscribe(
filter: RelaySubscriptionFilter,
onEvent: (event: RelayEvent) => void,
options?: { replayMissedHistory?: boolean },
) {
await this.ensureConnected();
@@ -621,6 +623,7 @@ export class RelayClient {
mode: "live",
filter,
onEvent,
replayMissedHistory: options?.replayMissedHistory,
resolveReady,
});
@@ -58,6 +58,7 @@ type LiveSubscription = {
mode: "live";
filter: RelaySubscriptionFilter;
onEvent: (event: RelayEvent) => void;
replayMissedHistory?: boolean;
resolveReady?: () => void;
lastSeenCreatedAt?: number;
closedRetryAttempt?: number;
@@ -5,6 +5,7 @@ import {
buildReconnectReplayFilter,
replayLiveSubscriptions,
REPLAY_BATCH_SIZE,
shouldPageReconnectReplay,
} from "./relayReconnectReplay.ts";
import { buildChannelFilter } from "./relayChannelFilters.ts";
@@ -113,6 +114,32 @@ test("reconnect replay caps large steady-state limits", () => {
});
});
test("reconnect replay preserves the live-only zero-history contract", () => {
const filter = {
kinds: [9],
"#h": ["channel-1"],
limit: 0,
};
assert.deepEqual(replayFilter(filter, 123), {
kinds: [9],
"#h": ["channel-1"],
limit: 0,
since: 123,
});
});
test("missed-history replay is explicit for live-only subscriptions", () => {
const filter = {
kinds: [9],
"#h": ["channel-1"],
limit: 0,
};
assert.equal(shouldPageReconnectReplay(filter), false);
assert.equal(shouldPageReconnectReplay(filter, true), true);
});
test("reconnect replay keeps the stricter existing since window", () => {
const filter = {
kinds: [9],
@@ -69,7 +69,11 @@ export function buildReconnectReplayFilter(
return replayFilter;
}
export function shouldPageReconnectReplay(filter: RelaySubscriptionFilter) {
export function shouldPageReconnectReplay(
filter: RelaySubscriptionFilter,
replayMissedHistory = false,
) {
if (replayMissedHistory) return true;
return (
filter.limit > 0 &&
Array.isArray(filter["#h"]) &&
@@ -176,7 +180,10 @@ export async function replayLiveSubscriptions({
);
const shouldPageReplay =
replaySince !== undefined &&
shouldPageReconnectReplay(subscription.filter);
shouldPageReconnectReplay(
subscription.filter,
subscription.replayMissedHistory,
);
return { subId, subscription, replaySince, shouldPageReplay };
});