feat(vibe): reconcile mutable session previews in playback
This commit is contained in:
@@ -1,25 +1,417 @@
|
||||
import axios from 'axios';
|
||||
import type { Track } from '../types';
|
||||
import { usePlaybackStore } from '../store/usePlaybackStore';
|
||||
import { usePlaybackStore, type VibeAdvanceReason } from '../store/usePlaybackStore';
|
||||
import { useVibeStore } from '../store/useVibeStore';
|
||||
import { fetchNextBatch, vibeService, type VibeBatchStatus } from './vibeService';
|
||||
|
||||
export const INITIAL_VIBE_BATCH_SIZE = 5;
|
||||
import {
|
||||
vibeService,
|
||||
type DurableVibeSessionResponse,
|
||||
type VibeEventType,
|
||||
type VibePlanItem,
|
||||
} from './vibeService';
|
||||
import { trackService } from './trackService';
|
||||
|
||||
export interface StartedVibeSession {
|
||||
status: VibeBatchStatus;
|
||||
status: 'complete' | 'exhausted' | 'failed';
|
||||
tracks: Track[];
|
||||
}
|
||||
|
||||
let startInFlight: Promise<StartedVibeSession> | null = null;
|
||||
let advanceInFlight: Promise<void> | null = null;
|
||||
let materialTail: Promise<void> = Promise.resolve();
|
||||
|
||||
interface PendingEvent {
|
||||
sessionId: string;
|
||||
input: Parameters<typeof vibeService.event>[1];
|
||||
retried: boolean;
|
||||
settled: boolean;
|
||||
resolve: (response: DurableVibeSessionResponse) => void;
|
||||
reject: (error: unknown) => void;
|
||||
}
|
||||
|
||||
// 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 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] : [];
|
||||
});
|
||||
}
|
||||
|
||||
/** 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] : [];
|
||||
const seen = new Set(history.map((track) => track.id));
|
||||
const future = preview.filter((track) => !seen.has(track.id));
|
||||
playback.setVibeQueue([...history, ...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 [];
|
||||
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));
|
||||
}
|
||||
|
||||
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 {
|
||||
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 401: return 'Vibe needs a trusted local user identity. Set MUZICK_VIBE_USER_ID and try again.';
|
||||
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, resolve, reject });
|
||||
void flushEventOutbox();
|
||||
});
|
||||
}
|
||||
|
||||
function retryableEventError(error: unknown): boolean {
|
||||
return !isSessionTerminalError(error);
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
// 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;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Start a V2 plan and immediately hand its first recommendations to playback.
|
||||
* Keeping this in one place prevents entry points from accidentally replacing a
|
||||
* generated Vibe queue with a normal browse queue.
|
||||
* 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;
|
||||
|
||||
try {
|
||||
const feedback = await sendEvent(reason, current.id);
|
||||
if (!feedback?.planVersion) {
|
||||
usePlaybackStore.getState().pause();
|
||||
return;
|
||||
}
|
||||
const served = await serveNextPlayable(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().setCurrentPlanItem(served.response.now);
|
||||
replaceUnplayedQueue([served.now, ...served.preview]);
|
||||
usePlaybackStore.getState().advance();
|
||||
} catch (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 || !currentItem || currentItem.track_id !== trackId || !isCurrentVibeOwner(vibe.activeSessionId)) 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().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 = beginVibeSession(seed);
|
||||
startInFlight = serializeMaterial<StartedVibeSession>(async () => {
|
||||
const started = await vibeService.start(seed.id);
|
||||
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.setCurrentPlanItem(served.response.now);
|
||||
vibe.setInitialBatchStatus('idle');
|
||||
|
||||
const playback = usePlaybackStore.getState();
|
||||
playback.setVibeQueue([served.now, ...served.preview]);
|
||||
playback.playTrack(served.now);
|
||||
installVibeAdvanceHandler();
|
||||
return { status: 'complete', tracks: [served.now, ...served.preview] };
|
||||
});
|
||||
try {
|
||||
return await startInFlight;
|
||||
} finally {
|
||||
@@ -27,25 +419,18 @@ export async function startVibeSession(seed: Track): Promise<StartedVibeSession>
|
||||
}
|
||||
}
|
||||
|
||||
async function beginVibeSession(seed: Track): Promise<StartedVibeSession> {
|
||||
const { sessionId } = await vibeService.start(seed.id);
|
||||
const vibe = useVibeStore.getState();
|
||||
// Do not replace a working Vibe until the new session has produced a usable
|
||||
// initial batch. This also keeps the page prefetcher attached to the old
|
||||
// session while this request is in flight.
|
||||
const result = await fetchNextBatch(INITIAL_VIBE_BATCH_SIZE, sessionId);
|
||||
if (result.tracks.length === 0) return result;
|
||||
|
||||
vibe.setInitialBatchStatus('loading');
|
||||
vibe.setActiveSession({ sessionId, seedTrackId: seed.id });
|
||||
vibe.setSeedTrackId(seed.id);
|
||||
vibe.setCenterTrack(seed);
|
||||
vibe.setBuffer(result.tracks);
|
||||
vibe.setInitialBatchStatus('idle');
|
||||
|
||||
const playback = usePlaybackStore.getState();
|
||||
playback.setQueue(result.tracks);
|
||||
playback.playTrack(result.tracks[0]);
|
||||
|
||||
return result;
|
||||
export async function endVibeSession(): Promise<void> {
|
||||
return serializeMaterial(async () => {
|
||||
const sessionId = useVibeStore.getState().activeSessionId;
|
||||
try {
|
||||
if (sessionId) await vibeService.end(sessionId);
|
||||
} finally {
|
||||
const playback = usePlaybackStore.getState();
|
||||
playback.setVibeAdvanceHandler(null);
|
||||
useVibeStore.getState().reset();
|
||||
playback.pause();
|
||||
playback.setQueue([]);
|
||||
playback.setCurrentTrack(null);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user