Files
muzick/frontend/src/services/vibeSession.ts
T
kami 78f5feea11 fix(vibe): let a session survive a phone losing its connection
A locked screen drops the radio, changes network or dozes, and one
request fails. Both halves of an advance treated that as the session's
fault. The outbox kept the failed event at its head with the comment
that a later retry would pick it up, but nothing ever triggered one, so
it sat there while the caller was rejected. advanceVibe then cleared the
prefetched future and paused, destroying a plan that was still valid.

Transient failures now hold the outbox entry unsettled and resend the
same event id on a backoff, so no duplicate feedback reaches the
director. The retry around the serve sits on the serve alone: serving a
version is idempotent, while replaying the whole advance would report a
second outcome for a track heard once. Each wait ends early when the
browser says the network is back, which is the moment that matters when
a screen unlocks. A drop that outlives every retry leaves the future
intact to carry on from.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-08 23:28:50 +04:00

591 lines
23 KiB
TypeScript

import axios from 'axios';
import type { Track } from '../types';
import { usePlaybackStore, type VibeAdvanceReason } from '../store/usePlaybackStore';
import { useVibeStore } from '../store/useVibeStore';
import {
vibeService,
type DurableVibeSessionResponse,
type VibeCalendarContext,
type VibeEventType,
type VibePlanItem,
} from './vibeService';
import { trackService } from './trackService';
export interface StartedVibeSession {
status: 'complete' | 'exhausted' | 'failed';
tracks: Track[];
}
let startInFlight: Promise<StartedVibeSession> | null = null;
let advanceInFlight: Promise<void> | null = null;
let materialTail: Promise<void> = Promise.resolve();
// A seeded Vibe plays its seed first — asking for a vibe "from this track" and
// getting a different track is the surprise. The seed sits in front of the
// durable plan without being part of it, so the first advance must consume it
// locally instead of reporting feedback and serving the next item.
// ponytail: no 'completed' event is sent for the seed. The listener chose it
// explicitly; the director already has that signal from the session's seed id.
let seedPendingTrackId: string | null = null;
interface PendingEvent {
sessionId: string;
input: Parameters<typeof vibeService.event>[1];
retried: boolean;
settled: boolean;
/** Consecutive failures that were the network's fault rather than the server's. */
transientAttempts: number;
resolve: (response: DurableVibeSessionResponse) => void;
reject: (error: unknown) => void;
}
// A phone with a locked screen drops its radio, changes network, or dozes, and
// one request fails. Waiting through it is right: the plan is still valid and
// the event id is stable, so the same event is simply sent again. Roughly two
// minutes of waiting in total before giving up on the listener's behalf.
const TRANSIENT_BACKOFF_MS = [1_000, 2_000, 5_000, 10_000, 20_000, 30_000, 30_000, 30_000];
// The event ledger deduplicates client_event_id. Keep an event in this ordered
// outbox until the server acknowledges it so a transient failure never turns a
// retry into a second listener action.
const eventOutbox: PendingEvent[] = [];
let flushingOutbox = false;
function serializeMaterial<T>(operation: () => Promise<T>): Promise<T> {
const result = materialTail.then(operation, operation);
materialTail = result.then(() => undefined, () => undefined);
return result;
}
function newEventId(): string {
if (typeof crypto !== 'undefined' && typeof crypto.randomUUID === 'function') return crypto.randomUUID();
// UUID v4-shaped fallback for older embedded webviews. The server only uses
// this as an idempotency key, not as a source of entropy.
return 'xxxxxxxx-xxxx-4xxx-yxxx-xxxxxxxxxxxx'.replace(/[xy]/g, (letter) => {
const value = Math.floor(Math.random() * 16);
return (letter === 'x' ? value : (value & 0x3) | 0x8).toString(16);
});
}
function localCalendarContext(): VibeCalendarContext {
const now = new Date();
let timeZone: string | undefined;
try {
timeZone = Intl.DateTimeFormat().resolvedOptions().timeZone || undefined;
} catch {
// Some embedded players omit Intl time-zone support. The coarse calendar
// fields still provide useful, non-identifying context.
}
return {
localHour: now.getHours(),
weekday: now.getDay(),
month: now.getMonth() + 1,
timeZone,
};
}
function isPlayable(track: Track): boolean {
return !['HIDDEN', 'MISSING', 'DELETED'].includes(track.state);
}
async function hydrateItem(item: VibePlanItem | null): Promise<Track | null> {
if (!item) return null;
try {
const track = await trackService.getTrack(item.track_id);
return isPlayable(track) ? track : null;
} catch {
// A plan can outlive a hidden/deleted file. Never substitute another item
// for this ordinal: keeping the remaining order is safer than a mismatch.
return null;
}
}
async function hydratePreview(items: VibePlanItem[]): Promise<Track[]> {
const uniqueIds = [...new Set(items.map((item) => item.track_id))];
const loaded = await Promise.all(uniqueIds.map(async (id) => {
try {
const track = await trackService.getTrack(id);
return [id, isPlayable(track) ? track : null] as const;
} catch {
return [id, null] as const;
}
}));
const byId = new Map(loaded.filter((entry): entry is readonly [string, Track] => entry[1] !== null));
return items.flatMap((item) => {
const track = byId.get(item.track_id);
return track ? [track] : [];
});
}
/**
* Two recordings of one song — a cover, a remaster, another artist's version —
* are distinct track ids but read as a duplicate in one sitting. The title is the
* key; remixes and live cuts name themselves in the title, so they survive.
* ponytail: title string match, no normalisation beyond case and edges. Add
* feat./punctuation stripping only if real duplicates keep getting through.
*/
const songKey = (track: Track) => (track.title || track.id).trim().toLowerCase();
function dedupeSongs(tracks: Track[], seen = new Set<string>()): Track[] {
return tracks.filter((track) => {
const key = songKey(track);
if (seen.has(key)) return false;
seen.add(key);
return true;
});
}
/** Replace only the queue after the currently playing Vibe track. */
function replaceUnplayedQueue(preview: Track[]): void {
const playback = usePlaybackStore.getState();
const current = playback.currentTrack;
const queueIndex = current
? (playback.currentIndex >= 0 && playback.queue[playback.currentIndex]?.id === current.id
? playback.currentIndex
: playback.queue.findIndex((track) => track.id === current.id))
: -1;
const history = queueIndex >= 0
? playback.queue.slice(0, queueIndex + 1)
: current ? [current] : [];
// While the seed plays, the durable cursor's own track sits between it and the
// preview. A `preview` list never contains that served item, so keep it — but
// only while it is still the cursor, never after it has been retired.
const next = playback.queue[queueIndex + 1];
const served = queueIndex >= 0
&& playback.queue[queueIndex]?.id === seedPendingTrackId
&& next
&& useVibeStore.getState().currentPlanItem?.track_id === next.id
? [next]
: [];
const future = dedupeSongs(preview, new Set([...history, ...served].map(songKey)));
playback.setVibeQueue([...history, ...served, ...future]);
}
function isCurrentVibeOwner(sessionId: string): boolean {
return useVibeStore.getState().activeSessionId === sessionId
&& usePlaybackStore.getState().queueOwner === 'vibe';
}
async function reconcilePreview(sessionId: string, response: DurableVibeSessionResponse): Promise<Track[]> {
if (!isCurrentVibeOwner(sessionId)) return [];
const preview = await hydratePreview(response.preview);
if (!isCurrentVibeOwner(sessionId) || !useVibeStore.getState().setPlan(response.planVersion, preview)) return [];
useVibeStore.getState().setProfile(response.state);
replaceUnplayedQueue(preview);
return preview;
}
async function serveNextCurrent(sessionId: string, version: number): Promise<DurableVibeSessionResponse> {
// A concurrent device or a feedback replan can make a version stale between
// the event response and /next. A stale response has an uncommitted `now`;
// refresh once with its latest version before admitting a track to playback.
let response = await vibeService.next(sessionId, version);
if (response.now?.committed) return response;
if (!response.planVersion || response.planVersion === version) return response;
response = await vibeService.next(sessionId, response.planVersion);
return response;
}
async function serveNextPlayable(
sessionId: string,
version: number,
): Promise<{ response: DurableVibeSessionResponse; now: Track; preview: Track[] } | null> {
return resolvePlayableResponse(sessionId, await serveNextCurrent(sessionId, version));
}
/**
* Serving a version is idempotent, so a dropped connection here costs nothing
* but the wait. The feedback event has already been acknowledged by this point,
* which is why the retry sits around the serve alone: replaying the whole
* advance would send a second event for a track the listener heard once.
*/
async function serveNextPlayableRetrying(
sessionId: string,
version: number,
): Promise<{ response: DurableVibeSessionResponse; now: Track; preview: Track[] } | null> {
for (let attempt = 0; ; attempt++) {
try {
return await serveNextPlayable(sessionId, version);
} catch (error) {
if (!isTransientEventError(error) || attempt >= TRANSIENT_BACKOFF_MS.length - 1) throw error;
await waitBeforeRetry(attempt);
}
}
}
async function advanceResponsePastUnplayable(
sessionId: string,
response: DurableVibeSessionResponse,
): Promise<DurableVibeSessionResponse> {
if (!response.now?.committed || !response.planVersion) return response;
const unplayable = {
eventId: newEventId(),
planVersionId: response.now.plan_version_id,
ordinal: response.now.ordinal,
trackId: response.now.track_id,
};
try {
return await vibeService.advancePastUnplayable(sessionId, response.planVersion, unplayable);
} catch (error) {
// A response may have been lost after the server committed the advance.
// Retry the same event id so it returns the same replacement rather than
// consuming another future item.
if (isSessionTerminalError(error)) throw error;
return vibeService.advancePastUnplayable(sessionId, response.planVersion, unplayable);
}
}
async function resolvePlayableResponse(
sessionId: string,
initialResponse: DurableVibeSessionResponse,
): Promise<{ response: DurableVibeSessionResponse; now: Track; preview: Track[] } | null> {
let response = initialResponse;
// A durable plan can reference a file which has since become hidden. Commit
// past such entries but never load one into the player. An unplayable
// advancement already returns and commits its replacement, so process that
// response directly: asking ordinary /next again would replay the original
// served cursor rather than advancing through consecutive hidden entries.
for (let attempts = 0; attempts < 20; attempts++) {
if (!response.now?.committed || !response.planVersion) {
if (!response.planVersion) return null;
response = await serveNextCurrent(sessionId, response.planVersion);
continue;
}
let now = await hydrateItem(response.now);
let preview = await hydratePreview(response.preview);
if (now) return { response, now, preview };
const advanced = await advanceResponsePastUnplayable(sessionId, response);
if (!advanced.planVersion) return null;
// The unplayable transition may itself race a feedback replan. Its stale
// response did not advance the old revision, so version-serve the current
// revision normally rather than treating a preview item as committed.
if (!advanced.now?.committed) {
response = await serveNextCurrent(sessionId, advanced.planVersion);
continue;
}
response = advanced;
}
return null;
}
function deactivateBrokenSession(): void {
seedPendingTrackId = null;
const playback = usePlaybackStore.getState();
playback.setVibeAdvanceHandler(null);
useVibeStore.getState().reset();
playback.pause();
playback.setQueue([]);
playback.setCurrentTrack(null);
}
function isSessionTerminalError(error: unknown): boolean {
return axios.isAxiosError(error) && [401, 404, 409].includes(error.response?.status ?? 0);
}
export function vibeErrorMessage(error: unknown): string {
if (!axios.isAxiosError(error)) return 'Could not refresh this Vibe. Please try again.';
switch (error.response?.status) {
case 400: return 'Vibe needs a valid user identity.';
case 404: return 'This Vibe session is no longer available.';
case 409: return 'This Vibe session has already ended or was replaced.';
default: return 'Could not refresh this Vibe. Please try again.';
}
}
async function sendEvent(
type: VibeEventType,
trackId?: string,
positionMs?: number,
durationMs?: number,
): Promise<DurableVibeSessionResponse | null> {
const sessionId = useVibeStore.getState().activeSessionId;
if (!sessionId) return null;
const input = {
eventId: newEventId(),
type,
trackId,
occurredAt: new Date().toISOString(),
positionMs,
durationMs,
};
return new Promise<DurableVibeSessionResponse>((resolve, reject) => {
eventOutbox.push({ sessionId, input, retried: false, settled: false, transientAttempts: 0, resolve, reject });
void flushEventOutbox();
});
}
function retryableEventError(error: unknown): boolean {
return !isSessionTerminalError(error);
}
/**
* A failure the request never survived to reach an opinion about: no response
* at all (offline, timeout, DNS), or a server that is momentarily unable rather
* than refusing. These say nothing about the session, so they must not be
* allowed to discard a valid plan.
*/
function isTransientEventError(error: unknown): boolean {
if (!axios.isAxiosError(error)) return false;
if (!error.response) return true;
return error.response.status === 429 || error.response.status >= 500;
}
/** Wait out a backoff, but come back early the moment the network returns. */
function waitBeforeRetry(attempt: number): Promise<void> {
const delay = TRANSIENT_BACKOFF_MS[Math.min(attempt, TRANSIENT_BACKOFF_MS.length - 1)];
return new Promise((resolve) => {
let settled = false;
const finish = () => {
if (settled) return;
settled = true;
clearTimeout(timer);
window.removeEventListener('online', finish);
resolve();
};
const timer = setTimeout(finish, delay);
window.addEventListener('online', finish);
});
}
async function flushEventOutbox(): Promise<void> {
if (flushingOutbox) return;
flushingOutbox = true;
try {
while (eventOutbox.length > 0) {
const entry = eventOutbox[0];
try {
const response = await vibeService.event(entry.sessionId, entry.input);
eventOutbox.shift();
entry.settled = true;
entry.resolve(response);
} catch (error) {
// Retry once immediately using the exact same client event id. After
// that leave it at the head for a later retry, rather than discarding
// the idempotency key or allowing newer material events to overtake it.
if (!entry.retried && retryableEventError(error)) {
entry.retried = true;
continue;
}
// The network failed, not the session. Hold the entry unsettled and
// keep trying: rejecting here is what used to end the Vibe whenever a
// locked phone lost its connection for a moment.
if (isTransientEventError(error) && entry.transientAttempts < TRANSIENT_BACKOFF_MS.length) {
const attempt = entry.transientAttempts;
entry.transientAttempts += 1;
await waitBeforeRetry(attempt);
continue;
}
// A session that is gone/ended can never acknowledge this event. Do
// not let an irrecoverable old-session entry block a later session.
if (isSessionTerminalError(error)) eventOutbox.shift();
if (!entry.settled) {
entry.settled = true;
entry.reject(error);
}
return;
}
}
} finally {
flushingOutbox = false;
}
}
/**
* Send a non-navigation event. Material feedback reconciles the future before
* it resolves, so no old prefetch remains after Keep or an implicit update.
*/
export async function reportVibeEvent(
type: VibeEventType,
trackId?: string,
positionMs?: number,
durationMs?: number,
): Promise<void> {
return serializeMaterial(async () => {
const sessionId = useVibeStore.getState().activeSessionId;
if (!sessionId || !isCurrentVibeOwner(sessionId)) return;
try {
const response = await sendEvent(type, trackId, positionMs, durationMs);
// Even a duplicate material event can acknowledge a canonical
// replacement revision (replanned=false). Reconcile every valid
// revision so a response lost after its original replan cannot leave a
// stale locally-prefetched future behind.
if (response?.planVersion !== null && response?.planVersion !== undefined) {
await reconcilePreview(sessionId, response);
}
} catch (error) {
if (isSessionTerminalError(error)) deactivateBrokenSession();
throw error;
}
});
}
/** Advance only after the prior track's durable outcome has produced a new plan. */
export function advanceVibe(reason: VibeAdvanceReason): Promise<void> {
if (advanceInFlight) return advanceInFlight;
advanceInFlight = serializeMaterial(async () => {
const vibe = useVibeStore.getState();
const current = usePlaybackStore.getState().currentTrack;
if (!vibe.activeSessionId || !current || !isCurrentVibeOwner(vibe.activeSessionId)) return;
// The seed is not a plan item — step off it without touching the cursor.
if (seedPendingTrackId && current.id === seedPendingTrackId) {
seedPendingTrackId = null;
usePlaybackStore.getState().advance();
// A dislike still has to reach the director; a completed seed carries no
// information the session's seed id does not already hold.
if (reason !== 'completed') void sendEvent(reason, current.id).catch(() => undefined);
return;
}
try {
const feedback = await sendEvent(reason, current.id);
if (!feedback?.planVersion) {
usePlaybackStore.getState().pause();
return;
}
const served = await serveNextPlayableRetrying(vibe.activeSessionId, feedback.planVersion);
if (!served || served.response.sessionId !== vibe.activeSessionId || !isCurrentVibeOwner(vibe.activeSessionId)) {
// Never play an uncommitted or unresolvable plan item. The user can
// retry from the page after the director publishes another revision.
replaceUnplayedQueue([]);
usePlaybackStore.getState().pause();
return;
}
if (!useVibeStore.getState().setPlan(served.response.planVersion, served.preview)) return;
useVibeStore.getState().setProfile(served.response.state);
useVibeStore.getState().setCurrentPlanItem(served.response.now);
replaceUnplayedQueue([served.now, ...served.preview]);
usePlaybackStore.getState().advance();
} catch (error) {
// A network failure that outlived every retry leaves the plan valid and
// the prefetched future worth keeping, so the listener can carry on from
// the page once they are back on a connection.
if (isTransientEventError(error)) {
usePlaybackStore.getState().pause();
throw error;
}
// Clearing the future is deliberate: carrying on with stale prefetches
// after a rejected feedback/replan would violate the plan boundary.
replaceUnplayedQueue([]);
if (isSessionTerminalError(error)) deactivateBrokenSession();
else usePlaybackStore.getState().pause();
throw error;
}
}).finally(() => { advanceInFlight = null; });
return advanceInFlight;
}
/**
* A stream can fail after its track metadata was successfully hydrated. This
* advances the exact durable cursor through the explicit unplayable protocol,
* rather than treating it as ordinary feedback and allowing a replan to hide
* the failure.
*/
export function advancePastUnplayableVibeTrack(trackId: string): Promise<void> {
if (advanceInFlight) return advanceInFlight;
advanceInFlight = serializeMaterial(async () => {
const vibe = useVibeStore.getState();
const playback = usePlaybackStore.getState();
const currentItem = vibe.currentPlanItem;
if (!vibe.activeSessionId || !isCurrentVibeOwner(vibe.activeSessionId)) return;
// The seed has no durable cursor to advance — a seed that will not stream is
// simply stepped over, leaving the plan's first item to play next.
if (seedPendingTrackId === trackId) {
seedPendingTrackId = null;
playback.advance();
return;
}
if (!currentItem || currentItem.track_id !== trackId) return;
try {
const advanced = await advanceResponsePastUnplayable(vibe.activeSessionId, {
sessionId: vibe.activeSessionId,
planVersion: vibe.planVersion,
now: currentItem,
preview: [],
state: {},
replanned: false,
replanReason: null,
});
const served = await resolvePlayableResponse(vibe.activeSessionId, advanced);
if (!served || served.response.sessionId !== vibe.activeSessionId || !isCurrentVibeOwner(vibe.activeSessionId)) {
replaceUnplayedQueue([]);
playback.pause();
return;
}
if (!useVibeStore.getState().setPlan(served.response.planVersion, served.preview)) return;
useVibeStore.getState().setProfile(served.response.state);
useVibeStore.getState().setCurrentPlanItem(served.response.now);
replaceUnplayedQueue([served.now, ...served.preview]);
playback.advance();
} catch (error) {
replaceUnplayedQueue([]);
if (isSessionTerminalError(error)) deactivateBrokenSession();
else playback.pause();
throw error;
}
}).finally(() => { advanceInFlight = null; });
return advanceInFlight;
}
function installVibeAdvanceHandler(): void {
usePlaybackStore.getState().setVibeAdvanceHandler((reason) => {
void advanceVibe(reason).catch(() => undefined);
});
}
/** Start, version-serve and hydrate the first durable Vibe track. */
export async function startVibeSession(seed: Track): Promise<StartedVibeSession> {
if (startInFlight) return startInFlight;
startInFlight = serializeMaterial<StartedVibeSession>(async () => {
const started = await vibeService.start(seed.id, localCalendarContext());
if (!started.planVersion) return { status: 'exhausted', tracks: [] };
const served = await serveNextPlayable(started.sessionId, started.planVersion);
if (!served) return { status: 'exhausted', tracks: [] };
const vibe = useVibeStore.getState();
// A newly started session has its own revision sequence. Drop the old
// local revision before admitting revision 1 from this new session.
vibe.reset();
vibe.setInitialBatchStatus('loading');
vibe.setActiveSession({ sessionId: started.sessionId, seedTrackId: seed.id });
vibe.setCenterTrack(seed);
vibe.setPlan(served.response.planVersion, served.preview);
vibe.setProfile(served.response.state);
vibe.setCurrentPlanItem(served.response.now);
vibe.setInitialBatchStatus('idle');
const playback = usePlaybackStore.getState();
const seedFirst = isPlayable(seed) && seed.id !== served.now.id;
seedPendingTrackId = seedFirst ? seed.id : null;
const queue = dedupeSongs(
seedFirst ? [seed, served.now, ...served.preview] : [served.now, ...served.preview]
);
playback.setVibeQueue(queue);
playback.playTrack(queue[0]);
installVibeAdvanceHandler();
return { status: 'complete', tracks: queue };
});
try {
return await startInFlight;
} finally {
startInFlight = null;
}
}
export async function endVibeSession(): Promise<void> {
return serializeMaterial(async () => {
const sessionId = useVibeStore.getState().activeSessionId;
try {
if (sessionId) await vibeService.end(sessionId);
} finally {
seedPendingTrackId = null;
const playback = usePlaybackStore.getState();
playback.setVibeAdvanceHandler(null);
useVibeStore.getState().reset();
playback.pause();
playback.setQueue([]);
playback.setCurrentTrack(null);
}
});
}