From fe13798c99033ed7f2c0079f68a1d1b4dd1b6b15 Mon Sep 17 00:00:00 2001 From: kami Date: Sun, 2 Aug 2026 00:49:27 +0400 Subject: [PATCH] feat(vibe): plan musical arcs and callbacks --- .../src/routes/vibe-sessions.routes.test.ts | 17 + backend/src/routes/vibe-sessions.routes.ts | 6 + backend/src/services/generators.service.ts | 14 + .../src/services/session-director.service.ts | 519 +++++++++++++++++- backend/src/services/session-director.test.ts | 256 ++++++++- .../vibe-session-coordinator.service.test.ts | 61 +- .../vibe-session-coordinator.service.ts | 79 ++- 7 files changed, 910 insertions(+), 42 deletions(-) diff --git a/backend/src/routes/vibe-sessions.routes.test.ts b/backend/src/routes/vibe-sessions.routes.test.ts index a9300b5..f739f4d 100644 --- a/backend/src/routes/vibe-sessions.routes.test.ts +++ b/backend/src/routes/vibe-sessions.routes.test.ts @@ -68,6 +68,23 @@ describe('durable Vibe session routes', () => { await app.close(); }); + it('reserves track_served for the authoritative /next operation', async () => { + const { app, coordinator } = await appWithCoordinator(); + const result = await app.inject({ + method: 'POST', url: `/v2/vibe/sessions/${SESSION_ID}/events`, + payload: { + type: 'track_served', + trackId: TRACK_ID, + payload: { planVersionId: '33333333-3333-4333-8333-333333333333', ordinal: 0 }, + }, + }); + + expect(result.statusCode).toBe(400); + expect(result.json()).toEqual({ error: 'track_served is reserved for the server /next operation' }); + expect(coordinator.appendEvent).not.toHaveBeenCalled(); + await app.close(); + }); + it('returns a lifecycle conflict when an initial plan race ends or replaces the session', async () => { const { app, coordinator } = await appWithCoordinator(); coordinator.start.mockRejectedValueOnce( diff --git a/backend/src/routes/vibe-sessions.routes.ts b/backend/src/routes/vibe-sessions.routes.ts index a477340..6e2b326 100644 --- a/backend/src/routes/vibe-sessions.routes.ts +++ b/backend/src/routes/vibe-sessions.routes.ts @@ -101,6 +101,12 @@ export default async function vibeSessionsRoutes( if (!body || !VIBE_EVENT_TYPES.includes(body.type as typeof VIBE_EVENT_TYPES[number])) { return reply.code(400).send({ error: 'type must be a supported Vibe event type' }); } + // Delivery is an authoritative state transition performed only by /next. + // Accepting this event from the public ledger endpoint would let a client + // fabricate exposure rows and consume the server-side surprise budget. + if (body.type === 'track_served') { + return reply.code(400).send({ error: 'track_served is reserved for the server /next operation' }); + } if (body.eventId !== undefined && !validUuid(body.eventId)) { return reply.code(400).send({ error: 'eventId must be a UUID' }); } diff --git a/backend/src/services/generators.service.ts b/backend/src/services/generators.service.ts index 1077563..8612b5f 100644 --- a/backend/src/services/generators.service.ts +++ b/backend/src/services/generators.service.ts @@ -20,6 +20,20 @@ export interface Candidate { generatorId: string; explanation: ClaimEdge[]; relevance: number; + /** + * Filled by the session director after sequence planning. Generators remain + * deliberately unaware of slots and objectives, while the durable plan can + * retain why this particular candidate won its position. + */ + plan?: { + slotRole: string; + score: number; + scoreBreakdown: Record; + explanation: Record; + /** Revision-level policy and constraint evidence, copied into the durable + * objective snapshot by the coordinator. */ + objective?: Record; + }; } export interface GeneratorContext { diff --git a/backend/src/services/session-director.service.ts b/backend/src/services/session-director.service.ts index 4e07647..e04c621 100644 --- a/backend/src/services/session-director.service.ts +++ b/backend/src/services/session-director.service.ts @@ -22,6 +22,8 @@ export interface RecentPlay { vocal: boolean | null; decade: number | null; valence: number | null; + acousticness?: number | null; + instrumentalness?: number | null; albumId?: string | null; producerIds?: string[]; labelIds?: string[]; @@ -40,6 +42,8 @@ export interface CandidateConstraintMetadata { energy?: number; bpm?: number; valence?: number; + acousticness?: number; + instrumentalness?: number; decade?: number; producerIds?: string[]; labelIds?: string[]; @@ -87,6 +91,46 @@ export interface PlanBuildOptions { retainedPlan?: Candidate[]; } +export interface ArcRange { + min?: number; + max?: number; + /** Preferred delta from the immediately preceding track. */ + maxDelta?: number; +} + +export interface CallbackToken { + id: string; + phase: 'anchor' | 'return'; + /** The return may use the anchor artist, genre, or simply a familiar item. */ + theme: 'artist' | 'genre' | 'favorite'; + minSeparation: number; + maxSeparation: number; +} + +export interface SurpriseDirective { + /** A surprise is bounded and must be followed by a recovery anchor. */ + recoveryRole: string; + maxPerPlan: number; + maxPerHour: number; +} + +export interface ArcSlot { + position: number; + role: string; + /** Optional only for legacy callers; getArcSlots always supplies targets. */ + targets?: { + energy?: ArcRange; + tempo?: ArcRange; + valence?: ArcRange; + acousticness?: ArcRange; + instrumentality?: ArcRange; + /** 0 familiar, 1 exploratory. Derived from source and known metadata. */ + novelty?: ArcRange; + }; + callback?: CallbackToken; + surprise?: SurpriseDirective; +} + const W_ENJOY = 1.0; const W_FATIGUE = 0.4; const W_DIVERSITY = 0.3; @@ -208,12 +252,106 @@ function metadataFromRecentPlay(play: RecentPlay): CandidateConstraintMetadata { energy: play.energy ?? undefined, bpm: play.bpm ?? undefined, valence: play.valence ?? undefined, + acousticness: play.acousticness ?? undefined, + instrumentalness: play.instrumentalness ?? undefined, decade: play.decade ?? undefined, producerIds: play.producerIds ?? [], labelIds: play.labelIds ?? [], }; } +function candidateNovelty(candidate: Candidate, metadata: CandidateConstraintMetadata | undefined): number { + if (metadata?.favorite) return 0.05; + if (metadata?.newArtist) return 0.8; + if (candidate.generatorId === 'discovery') return 0.7; + if (candidate.generatorId === 'adjacent') return 0.45; + if (candidate.generatorId === 'deep-dive') return 0.2; + return 0.3; +} + +function rangeScore(value: number | undefined, range: ArcRange | undefined): { score: number; known: boolean } { + if (!range) return { score: 1, known: true }; + if (value === undefined || !Number.isFinite(value)) return { score: 0.5, known: false }; + const min = range.min ?? -Infinity; + const max = range.max ?? Infinity; + if (value >= min && value <= max) return { score: 1, known: true }; + const distance = value < min ? min - value : value - max; + // A target miss is a soft degradation. It is never an eligibility failure. + return { score: Math.max(0, 1 - distance / 0.5), known: true }; +} + +/** + * Unknown analysis must remain playable, but it cannot displace a candidate + * whose measured audio features satisfy this slot. This is intentionally + * limited to absolute targets: `maxDelta` is a transition preference, not an + * independent eligibility rule. + */ +function matchesMeasuredArcTargets( + metadata: CandidateConstraintMetadata | undefined, + slot: ArcSlot, +): boolean { + const targets = slot.targets ?? {}; + const dimensions: Array<[number | undefined, ArcRange | undefined]> = [ + [metadata?.energy, targets.energy], + [metadata?.bpm, targets.tempo], + [metadata?.valence, targets.valence], + [metadata?.acousticness, targets.acousticness], + [metadata?.instrumentalness, targets.instrumentality], + ]; + const absolute = dimensions.filter(([, range]) => range && (range.min !== undefined || range.max !== undefined)); + if (absolute.length === 0) return false; + return absolute.every(([value, range]) => { + if (value === undefined || !Number.isFinite(value)) return false; + return value >= (range!.min ?? -Infinity) && value <= (range!.max ?? Infinity); + }); +} + +function hasAbsoluteArcTargets(slot: ArcSlot): boolean { + const targets = slot.targets ?? {}; + return [targets.energy, targets.tempo, targets.valence, targets.acousticness, targets.instrumentality] + .some(range => range && (range.min !== undefined || range.max !== undefined)); +} + +/** A local DJ transition is useful even with partial analysis: unknown values + * lower confidence, whereas known values are evaluated smoothly. */ +export function scoreArcTransition( + candidate: Candidate, + candidateMetadata: CandidateConstraintMetadata | undefined, + previousMetadata: CandidateConstraintMetadata | undefined, + slot: ArcSlot, +): { score: number; confidence: number; targetScore: number; transitionScore: number } { + const targetsForSlot = slot.targets ?? {}; + const targets = [ + rangeScore(candidateMetadata?.energy, targetsForSlot.energy), + rangeScore(candidateMetadata?.bpm, targetsForSlot.tempo), + rangeScore(candidateMetadata?.valence, targetsForSlot.valence), + rangeScore(candidateMetadata?.acousticness, targetsForSlot.acousticness), + rangeScore(candidateMetadata?.instrumentalness, targetsForSlot.instrumentality), + rangeScore(candidateNovelty(candidate, candidateMetadata), targetsForSlot.novelty), + ]; + const targetScore = targets.reduce((total, item) => total + item.score, 0) / targets.length; + const knownTargets = targets.filter(item => item.known).length; + const deltas: Array<[number | undefined, number | undefined, ArcRange | undefined, number]> = [ + [candidateMetadata?.energy, previousMetadata?.energy, targetsForSlot.energy, 0.3], + [candidateMetadata?.bpm, previousMetadata?.bpm, targetsForSlot.tempo, 35], + [candidateMetadata?.valence, previousMetadata?.valence, targetsForSlot.valence, 0.3], + [candidateMetadata?.acousticness, previousMetadata?.acousticness, targetsForSlot.acousticness, 0.3], + [candidateMetadata?.instrumentalness, previousMetadata?.instrumentalness, targetsForSlot.instrumentality, 0.3], + ]; + const knownDeltas = deltas.filter(([current, previous]) => current !== undefined && previous !== undefined); + const transitionScore = knownDeltas.length === 0 ? 0.5 : knownDeltas.reduce((total, [current, previous, target, defaultMaxDelta]) => { + const maxDelta = target?.maxDelta ?? defaultMaxDelta; + return total + Math.max(0, 1 - Math.abs(current! - previous!) / maxDelta); + }, 0) / knownDeltas.length; + const confidence = (knownTargets + knownDeltas.length) / (targets.length + deltas.length); + return { + score: targetScore * 0.65 + transitionScore * 0.35, + confidence, + targetScore, + transitionScore, + }; +} + /** * The constraint layer is intentionally pure. Ranking supplies its candidate * order; this layer chooses a feasible sequence and returns every soft rule it @@ -222,7 +360,7 @@ function metadataFromRecentPlay(play: RecentPlay): CandidateConstraintMetadata { */ export function selectConstrainedSequence(params: { candidates: Candidate[]; - slots: { position: number; role: string }[]; + slots: ArcSlot[]; metadata: Map; budgets: DiversityBudget[]; roleToGeneratorIds: (role: string) => string[]; @@ -249,6 +387,10 @@ export function selectConstrainedSequence(params: { const softDimensions = ['artist', 'genre', 'language']; const planLength = retainedPlan.length + slots.length; const loopedValueSet = new Set([...loopedValues, ...(loopedValue ? [loopedValue] : [])]); + let previousMetadata = retainedPlan.length > 0 + ? metadata.get(retainedPlan[retainedPlan.length - 1].trackId) + : undefined; + const callbackAnchors = new Map(); // Candidate ranking is already stable. Build each role's preferred pool // once, preserving that order, rather than sorting the whole pool for every // slot in a long plan. @@ -322,26 +464,105 @@ export function selectConstrainedSequence(params: { increment('album', valuesForDimension(metadata.get(retained.trackId), 'album')); } + const hardConstraints = { + artistMaxPerPlan: MAX_ARTIST_PER_PLAN, + albumMaxPer40Tracks: MAX_ALBUM_PER_40_TRACKS, + uniqueTracks: true, + }; + for (let position = 0; position < slots.length; position++) { - const preferred = roleToGeneratorIds(slots[position].role); - const ordered = orderedForRole(slots[position].role); + let slot = slots[position]; let chosen: Candidate | undefined; + let chosenTransition: ReturnType | undefined; + let chosenCallbackScore = 0; + let chosenSurpriseScore = 0; + let chosenSelectionScore = 0; + let chosenRankScore = 0; let relaxedSoft = false; let relaxedArc = false; + let downgradedSurpriseGenerators: Set | undefined; + + /** A surprise is only admitted if its declared recovery role still has a + * hard-feasible anchor after the surprise itself has consumed its caps. */ + const canReserveRecovery = (surpriseCandidate: Candidate, surpriseMetadata: CandidateConstraintMetadata | undefined) => { + if (!slot.surprise) return true; + const recoveryGenerators = new Set(roleToGeneratorIds(slot.surprise.recoveryRole)); + const surpriseArtist = valueForDimension(surpriseMetadata, 'artist'); + const surpriseAlbum = valueForDimension(surpriseMetadata, 'album'); + return candidates.some(recovery => { + if (recovery.trackId === surpriseCandidate.trackId || selectedIds.has(recovery.trackId)) return false; + if (!recoveryGenerators.has(recovery.generatorId)) return false; + const recoveryMetadata = metadata.get(recovery.trackId); + const artist = valueForDimension(recoveryMetadata, 'artist'); + const album = valueForDimension(recoveryMetadata, 'album'); + const prospectiveArtistCount = artist && artist === surpriseArtist ? 1 : 0; + const prospectiveAlbumCount = album && album === surpriseAlbum ? 1 : 0; + return (!artist || count('artist', artist) + prospectiveArtistCount < MAX_ARTIST_PER_PLAN) + && (!album || count('album', album) + prospectiveAlbumCount < MAX_ALBUM_PER_40_TRACKS); + }); + }; + + // A surprise whose recovery cannot be reserved is downgraded before + // selection. This prevents an unrecoverable high-risk pick from being + // written to a revision merely because the later recovery slot is empty. + if (slot.surprise) { + const surpriseGenerators = new Set(roleToGeneratorIds(slot.role)); + const reservable = candidates.some(candidate => + !selectedIds.has(candidate.trackId) + && surpriseGenerators.has(candidate.generatorId) + && canReserveRecovery(candidate, metadata.get(candidate.trackId)), + ); + if (!reservable) { + relaxations.push({ + constraint: 'surprise_recovery', + stage: 'arc_precision', + reason: 'no hard-feasible recovery anchor remained for the requested surprise', + }); + downgradedSurpriseGenerators = surpriseGenerators; + slot = { + ...slot, + role: slot.surprise.recoveryRole, + surprise: undefined, + targets: { ...(slot.targets ?? {}), novelty: { max: 0.35 } }, + }; + } + } for (let pass = 0; pass < 3 && !chosen; pass++) { // pass 0: all soft constraints + arc role; pass 1: soft budgets/loop; // pass 2: permit an arc-source fallback. Hard sequence caps remain. const relaxSoft = pass >= 1; const relaxArc = pass >= 2; - for (const candidate of ordered) { + const preferred = roleToGeneratorIds(slot.role); + const ordered = orderedForRole(slot.role); + let best: { candidate: Candidate; transition: ReturnType; callbackScore: number; surpriseScore: number; score: number; rankScore: number } | undefined; + const fittingArcCandidateExists = hasAbsoluteArcTargets(slot) && ordered.some(candidate => { + if (selectedIds.has(candidate.trackId)) return false; + const m = metadata.get(candidate.trackId); + const artist = valueForDimension(m, 'artist'); + const album = valueForDimension(m, 'album'); + return (!artist || count('artist', artist) < MAX_ARTIST_PER_PLAN) + && (!album || count('album', album) < MAX_ALBUM_PER_40_TRACKS) + && matchesMeasuredArcTargets(m, slot); + }); + for (const [candidateIndex, candidate] of ordered.entries()) { if (selectedIds.has(candidate.trackId)) continue; + // Do not relabel the rejected high-risk item as a recovery merely + // because source relaxation is allowed for the downgraded slot. + if (downgradedSurpriseGenerators?.has(candidate.generatorId)) continue; const m = metadata.get(candidate.trackId); const artist = valueForDimension(m, 'artist'); const album = valueForDimension(m, 'album'); if (artist && count('artist', artist) >= MAX_ARTIST_PER_PLAN) continue; if (album && count('album', album) >= MAX_ALBUM_PER_40_TRACKS) continue; - if (!relaxArc && !preferred.includes(candidate.generatorId)) continue; + // A surprise is never silently satisfied by an arbitrary comfort + // fallback. If no explainable exploratory source is available the + // planner reports a shortened/degraded sequence instead. + if ((!relaxArc || slot.surprise) && !preferred.includes(candidate.generatorId)) continue; + // Arc measurements are an enforceable preference when the catalog + // gives us a hard-feasible measured fit. Sparse analysis falls back + // gracefully and is documented as an arc-precision relaxation below. + if (fittingArcCandidateExists && !matchesMeasuredArcTargets(m, slot)) continue; if (!relaxSoft && !explicitIntent) { if (softDimensions.some(d => exceedsUpperLimit(d, valuesForDimension(m, d)))) continue; @@ -352,30 +573,134 @@ export function selectConstrainedSequence(params: { if (!deficits.some(d => valuesForDimension(m, d).includes('true'))) continue; } } - chosen = candidate; + const transition = scoreArcTransition(candidate, m, previousMetadata, slot); + const anchor = slot.callback ? callbackAnchors.get(slot.callback.id) : undefined; + let callbackScore = 0; + if (slot.callback?.phase === 'return' && anchor) { + const separation = position - anchor.position; + const themeMatches = slot.callback.theme === 'favorite' + ? m?.favorite === true + : slot.callback.theme === 'artist' + ? !!m?.artistId && m.artistId === anchor.metadata.artistId + : !!m?.genreId && m.genreId === anchor.metadata.genreId; + // The token is preference-only: a hard diversity/repetition rule may + // still make a literal return impossible, in which case familiarity + // provides a graceful callback rather than a failed plan. + const inWindow = separation >= slot.callback.minSeparation && separation <= slot.callback.maxSeparation; + callbackScore = inWindow && themeMatches ? 1 : (m?.favorite ? 0.35 : 0); + } + const surpriseScore = slot.surprise + ? (candidate.generatorId === 'discovery' || candidate.generatorId === 'adjacent' || m?.newArtist ? 1 : 0) + : 0; + if (slot.surprise && !canReserveRecovery(candidate, m)) continue; + // Keep ranking meaningful, but let the sequence controller prefer a + // fitting continuation over the next independent highest-ranked song. + const rankScore = 1 - candidateIndex / Math.max(1, ordered.length); + const score = rankScore * 0.45 + transition.score * 0.4 + callbackScore * 0.1 + surpriseScore * 0.05; + if (!best || score > best.score || (score === best.score && candidate.trackId.localeCompare(best.candidate.trackId) < 0)) { + best = { candidate, transition, callbackScore, surpriseScore, score, rankScore }; + } + } + if (best) { + chosen = best.candidate; + chosenTransition = best.transition; + chosenCallbackScore = best.callbackScore; + chosenSurpriseScore = best.surpriseScore; + chosenSelectionScore = best.score; + chosenRankScore = best.rankScore; relaxedSoft = relaxSoft; relaxedArc = relaxArc; - break; } } if (!chosen) break; if (relaxedSoft && !relaxations.some(r => r.stage === 'soft_budget')) { relaxations.push({ constraint: loopDimension ?? 'diversity_budget', stage: 'soft_budget', reason: 'eligible inventory could not satisfy projected soft constraints' }); } - if (relaxedArc && !relaxations.some(r => r.stage === 'arc_precision')) { + if ((hasAbsoluteArcTargets(slot) && !matchesMeasuredArcTargets(metadata.get(chosen.trackId), slot)) + && !relaxations.some(r => r.stage === 'arc_precision' && r.constraint === 'arc_precision')) { + relaxations.push({ constraint: 'arc_precision', stage: 'arc_precision', reason: 'no hard-feasible measured candidate satisfied the requested arc targets' }); + } + if (relaxedArc && !relaxations.some(r => r.constraint === 'arc_source')) { relaxations.push({ constraint: 'arc_source', stage: 'arc_precision', reason: 'no hard-feasible candidate matched the requested arc slot' }); } - selected.push(chosen); + const chosenMetadata = metadata.get(chosen.trackId); + const decorated: Candidate = { + ...chosen, + plan: { + slotRole: slot.role, + score: chosenSelectionScore || (chosenTransition ? chosenTransition.score : chosen.relevance), + scoreBreakdown: { + relevance: chosen.relevance, + rankingPositionScore: chosenRankScore, + transition: chosenTransition?.transitionScore ?? 0.5, + arcTarget: chosenTransition?.targetScore ?? 0.5, + transitionConfidence: chosenTransition?.confidence ?? 0, + callback: chosenCallbackScore, + surprise: chosenSurpriseScore, + selection: chosenSelectionScore, + weights: { rank: 0.45, transition: 0.4, callback: 0.1, surprise: 0.05 }, + }, + explanation: { + arcRole: slot.role, + targets: slot.targets, + callback: slot.callback ? { + id: slot.callback.id, + phase: slot.callback.phase, + theme: slot.callback.theme, + minSeparation: slot.callback.minSeparation, + maxSeparation: slot.callback.maxSeparation, + matchScore: chosenCallbackScore, + } : null, + surprise: slot.surprise ? { + recoveryRole: slot.surprise.recoveryRole, + maxPerPlan: slot.surprise.maxPerPlan, + maxPerHour: slot.surprise.maxPerHour, + selectionScore: chosenSurpriseScore, + } : null, + arcPrecision: { + measuredFit: matchesMeasuredArcTargets(chosenMetadata, slot), + enforcedBecauseFitExists: hasAbsoluteArcTargets(slot), + }, + }, + }, + }; + selected.push(decorated); selectedIds.add(chosen.trackId); - const m = metadata.get(chosen.trackId); + const m = chosenMetadata; for (const dimension of ['artist', 'album', 'genre', 'language', ...lowerDimensions]) { increment(dimension, valuesForDimension(m, dimension)); } + previousMetadata = chosenMetadata; + if (slot.callback?.phase === 'anchor' && chosenMetadata) { + callbackAnchors.set(slot.callback.id, { metadata: chosenMetadata, position }); + } } if (selected.length < slots.length) { relaxations.push({ constraint: 'freshness', stage: 'freshness', reason: 'hard exclusions and sequence caps left too few eligible candidates' }); } - return { plan: selected, relaxations }; + const policy = { + selection: 'rank_transition_callback_surprise_v1', + weights: { rank: 0.45, transition: 0.4, callback: 0.1, surprise: 0.05 }, + arcTargets: 'enforced_when_a_hard_feasible_measured_fit_exists', + surprise: 'requires_reserved_recovery_anchor', + }; + const constraints = { + hard: hardConstraints, + softBudgetDimensions: [...softDimensions, ...lowerDimensions], + loopDimension: loopDimension ?? null, + }; + const objective = { policy, constraints, relaxations }; + return { + plan: selected.map(candidate => ({ + ...candidate, + plan: candidate.plan ? { + ...candidate.plan, + explanation: { ...candidate.plan.explanation, policy, constraints, relaxations }, + objective, + } : candidate.plan, + })), + relaxations, + }; } export class SessionDirector { @@ -723,27 +1048,77 @@ export class SessionDirector { return 'comfort'; } - getArcSlots(arcType: string, count: number): { position: number; role: string }[] { + getArcSlots(arcType: string, count: number): ArcSlot[] { const pattern = this.getArcPattern(arcType); - const slots: { position: number; role: string }[] = []; + const slots: ArcSlot[] = []; for (let i = 0; i < count; i++) { - slots.push({ position: i, role: pattern[i % pattern.length] }); + const template = pattern[i % pattern.length]; + // A repeated template must not replay an old callback token. The token + // identity is scoped to the cycle so each return has one clear anchor. + const cycle = Math.floor(i / pattern.length); + slots.push({ + ...template, + position: i, + callback: template.callback ? { ...template.callback, id: `${template.callback.id}:${cycle}` } : undefined, + surprise: template.surprise && cycle === 0 ? template.surprise : undefined, + }); } return slots; } - private getArcPattern(arcType: string): string[] { + private getArcPattern(arcType: string): Omit[] { + const familiar = { novelty: { max: 0.35 } }; + const adjacent = { novelty: { min: 0.2, max: 0.65 } }; + const newDiscovery = { novelty: { min: 0.55, max: 1 } }; + const callbackAnchor: CallbackToken = { + id: 'comfort-theme', phase: 'anchor', theme: 'artist', minSeparation: 2, maxSeparation: 6, + }; + const callbackReturn: CallbackToken = { + ...callbackAnchor, phase: 'return', + }; + const surprise: SurpriseDirective = { recoveryRole: 'favorite', maxPerPlan: 1, maxPerHour: 1 }; + const energeticSurprise: SurpriseDirective = { recoveryRole: 'cooldown', maxPerPlan: 1, maxPerHour: 1 }; switch (arcType) { case 'comfort': - return ['known', 'known', 'known', 'known', 'adjacent', 'adjacent', 'adjacent', 'adjacent', 'favorite', 'favorite']; + return [ + { role: 'known', targets: { energy: { min: 0.3, max: 0.65, maxDelta: 0.2 }, tempo: { maxDelta: 25 }, ...familiar } }, + { role: 'known', targets: { energy: { min: 0.3, max: 0.65, maxDelta: 0.18 }, tempo: { maxDelta: 22 }, ...familiar } }, + { role: 'adjacent', targets: { energy: { min: 0.35, max: 0.7, maxDelta: 0.22 }, tempo: { maxDelta: 28 }, ...adjacent } }, + { role: 'favorite', targets: { energy: { min: 0.3, max: 0.7, maxDelta: 0.25 }, valence: { maxDelta: 0.3 }, ...familiar }, callback: callbackAnchor }, + { role: 'adjacent', targets: { energy: { min: 0.35, max: 0.75, maxDelta: 0.25 }, ...adjacent } }, + { role: 'favorite', targets: { energy: { min: 0.3, max: 0.7, maxDelta: 0.25 }, ...familiar }, callback: callbackReturn }, + { role: 'surprise', targets: { energy: { min: 0.25, max: 0.75, maxDelta: 0.3 }, ...newDiscovery }, surprise }, + { role: 'favorite', targets: { energy: { min: 0.25, max: 0.65, maxDelta: 0.25 }, acousticness: { min: 0.1 }, ...familiar } }, + ]; case 'discovery': - return ['favorite', 'similar', 'new', 'favorite']; + return [ + { role: 'favorite', targets: { energy: { maxDelta: 0.25 }, ...familiar }, callback: callbackAnchor }, + { role: 'similar', targets: { energy: { maxDelta: 0.25 }, tempo: { maxDelta: 30 }, ...adjacent } }, + { role: 'surprise', targets: { energy: { maxDelta: 0.3 }, ...newDiscovery }, surprise }, + { role: 'favorite', targets: { energy: { maxDelta: 0.25 }, ...familiar }, callback: callbackReturn }, + ]; case 'energetic': - return ['medium', 'medium', 'high', 'high', 'high', 'peak', 'cooldown', 'cooldown']; + return [ + { role: 'medium', targets: { energy: { min: 0.45, max: 0.7, maxDelta: 0.18 }, tempo: { min: 95, max: 145, maxDelta: 25 }, ...familiar } }, + { role: 'medium', targets: { energy: { min: 0.5, max: 0.75, maxDelta: 0.15 }, tempo: { min: 100, max: 155, maxDelta: 20 }, ...adjacent } }, + { role: 'high', targets: { energy: { min: 0.65, max: 0.9, maxDelta: 0.22 }, tempo: { min: 110, max: 175, maxDelta: 28 }, ...adjacent } }, + { role: 'high', targets: { energy: { min: 0.7, max: 0.95, maxDelta: 0.18 }, tempo: { min: 115, max: 180, maxDelta: 25 }, ...newDiscovery } }, + { role: 'peak', targets: { energy: { min: 0.8, max: 1, maxDelta: 0.2 }, tempo: { min: 120, max: 190, maxDelta: 30 }, valence: { min: 0.45, maxDelta: 0.3 }, ...newDiscovery } }, + { role: 'surprise', targets: { energy: { min: 0.7, max: 1, maxDelta: 0.25 }, tempo: { maxDelta: 35 }, ...newDiscovery }, surprise: energeticSurprise }, + { role: 'cooldown', targets: { energy: { min: 0.45, max: 0.75, maxDelta: 0.3 }, tempo: { maxDelta: 35 }, ...familiar } }, + { role: 'cooldown', targets: { energy: { min: 0.35, max: 0.65, maxDelta: 0.2 }, acousticness: { min: 0.1 }, ...familiar } }, + ]; case 'late-night': - return ['soft', 'soft', 'ambient', 'ambient', 'acoustic', 'slow']; + return [ + { role: 'soft', targets: { energy: { max: 0.4, maxDelta: 0.15 }, tempo: { max: 115, maxDelta: 20 }, acousticness: { min: 0.2 }, ...familiar } }, + { role: 'soft', targets: { energy: { max: 0.38, maxDelta: 0.12 }, tempo: { max: 110, maxDelta: 18 }, ...familiar } }, + { role: 'ambient', targets: { energy: { max: 0.3, maxDelta: 0.15 }, tempo: { max: 100, maxDelta: 18 }, instrumentality: { min: 0.35 }, ...adjacent } }, + { role: 'ambient', targets: { energy: { max: 0.28, maxDelta: 0.1 }, acousticness: { min: 0.25 }, instrumentality: { min: 0.35 }, ...adjacent } }, + { role: 'acoustic', targets: { energy: { max: 0.45, maxDelta: 0.2 }, acousticness: { min: 0.35 }, ...familiar }, callback: callbackAnchor }, + { role: 'slow', targets: { energy: { max: 0.4, maxDelta: 0.15 }, tempo: { max: 105, maxDelta: 18 }, ...familiar }, callback: callbackReturn }, + ]; default: - return ['known', 'known', 'known', 'known', 'adjacent', 'adjacent', 'adjacent', 'adjacent', 'favorite', 'favorite']; + return this.getArcPattern('comfort'); } } @@ -766,6 +1141,8 @@ export class SessionDirector { return ['discovery']; case 'peak': return ['deep-dive', 'contextual']; + case 'surprise': + return ['discovery', 'adjacent']; case 'ambient': return ['contextual', 'comfort']; default: @@ -1075,7 +1452,7 @@ export class SessionDirector { WHERE ph.user_id = $1 AND ph.completed = true AND old_artist.artist_id = artist.artist_id ) AS new_artist, - taf.energy, taf.bpm, taf.valence, + taf.energy, taf.bpm, taf.valence, taf.acousticness, taf.instrumentalness, ${lineageIdsSql('produced', 'artist')} AS producer_ids, ${lineageIdsSql('same_label_as', 'artist')} AS label_ids FROM tracks t @@ -1107,6 +1484,8 @@ export class SessionDirector { energy: (row.energy as number | null) ?? undefined, bpm: (row.bpm as number | null) ?? undefined, valence: (row.valence as number | null) ?? undefined, + acousticness: (row.acousticness as number | null) ?? undefined, + instrumentalness: (row.instrumentalness as number | null) ?? undefined, decade: row.release_date ? Math.floor(new Date(row.release_date as string).getFullYear() / 10) * 10 : undefined, producerIds: Array.isArray(row.producer_ids) ? row.producer_ids as string[] : [], @@ -1223,6 +1602,48 @@ export class SessionDirector { ); } + /** + * The surprise budget is consumed when a surprise is actually delivered to + * playback, not when it is merely present in a preview. Replans copy the + * unserved tail into a new immutable revision, so counting `slot_role` + * rows would charge the same unserved surprise once per revision. The + * `track_served` event is the durable, exactly-once delivery boundary; its + * plan-version payload lets us identify the precise immutable item that + * was exposed. + */ + private async getSessionSurpriseUsage(sessionId: string, userId: string): Promise<{ session: number; hour: number }> { + const res = await this.db.pgClient.query( + `WITH surprise_exposures AS ( + SELECT DISTINCT e.id, e.occurred_at + FROM vibe_events e + JOIN vibe_plan_items item + ON item.plan_version_id = CASE + WHEN e.payload->>'planVersionId' ~* '^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$' + THEN (e.payload->>'planVersionId')::uuid + END + AND item.ordinal = CASE + WHEN (e.payload->>'ordinal') ~ '^(0|[1-9][0-9]{0,8})$' + OR ((e.payload->>'ordinal') ~ '^1[0-9]{9}$' AND (e.payload->>'ordinal') <= '2147483647') + THEN (e.payload->>'ordinal')::integer + END + AND item.track_id = e.track_id + JOIN vibe_plan_versions version + ON version.id = item.plan_version_id + AND version.session_id = e.session_id + WHERE e.session_id = $1 + AND e.user_id = $2 + AND e.type = 'track_served' + AND item.slot_role = 'surprise' + ) + SELECT COUNT(*)::int AS session_count, + COUNT(*) FILTER (WHERE occurred_at > NOW() - INTERVAL '1 hour')::int AS hour_count + FROM surprise_exposures`, + [sessionId, userId], + ); + const row = res.rows[0] as { session_count?: number | string; hour_count?: number | string } | undefined; + return { session: Number(row?.session_count ?? 0), hour: Number(row?.hour_count ?? 0) }; + } + // --------------------------------------------------------------- // D.9 — Plan + replan loop // --------------------------------------------------------------- @@ -1251,7 +1672,7 @@ export class SessionDirector { // Fetch recent completed plays for anti-loop detection const recentPlaysRes = await this.db.pgClient.query( `SELECT t.id AS track_id, t.album_id, artist.artist_id, tg.genre_id, - af.bpm, af.energy, af.valence, af.instrumentalness, + af.bpm, af.energy, af.valence, af.acousticness, af.instrumentalness, tl.language, t.release_date, ${lineageIdsSql('produced', 'artist')} AS producer_ids, @@ -1283,6 +1704,8 @@ export class SessionDirector { vocal: (r.instrumentalness == null) ? null : r.instrumentalness < 0.5, decade: r.release_date ? Math.floor(new Date(r.release_date).getFullYear() / 10) * 10 : null, valence: r.valence ?? null, + acousticness: r.acousticness ?? null, + instrumentalness: r.instrumentalness ?? null, albumId: r.album_id ?? null, producerIds: r.producer_ids ?? [], labelIds: r.label_ids ?? [], @@ -1295,7 +1718,9 @@ export class SessionDirector { const arcType = this.pickArc(state); const planSize = PLAN_SIZE; - const slots = this.getArcSlots(arcType, Math.max(0, planSize - retainedPlan.length)); + const rawSlots = this.getArcSlots(arcType, Math.max(0, planSize - retainedPlan.length)); + const surpriseUsage = await this.getSessionSurpriseUsage(sessionId, userId); + const slots = this.applySurpriseBudget(rawSlots, surpriseUsage, state.sessionAgeMin); let seedArtistId: string | null = null; if (seedTrackId) { @@ -1337,6 +1762,21 @@ export class SessionDirector { if (allCandidates.length === 0) { return []; } + const availableGeneratorIds = new Set(allCandidates.map(candidate => candidate.generatorId)); + // A budgeted surprise without an explainable discovery candidate becomes + // its declared recovery anchor before ranking. This is a deterministic + // degradation, not a random or mislabeled fallback. + const effectiveSlots = slots.map(slot => { + const hasSurpriseSource = !slot.surprise + || this.roleToGeneratorIds(slot.role).some(id => availableGeneratorIds.has(id)); + if (hasSurpriseSource || !slot.surprise) return slot; + return { + ...slot, + role: slot.surprise.recoveryRole, + surprise: undefined, + targets: { ...(slot.targets ?? {}), novelty: { max: 0.35 } }, + }; + }); const repetitionState = await this.buildRepetitionState(userId); const candidateArtistMap = await this.loadArtistMap( @@ -1374,7 +1814,7 @@ export class SessionDirector { ]); const constrained = this.constrainedSequence({ candidates: deduped, - slots, + slots: effectiveSlots, metadata, budgets, roleToGeneratorIds: role => this.roleToGeneratorIds(role), @@ -1394,6 +1834,37 @@ export class SessionDirector { return constrained.plan.slice(0, Math.max(0, planSize - retainedPlan.length)); } + /** Apply the durable delivery budget to fresh arc slots. Retained entries + * are intentionally absent here: they have not consumed anything until a + * `track_served` event exists for their immutable revision. */ + private applySurpriseBudget( + rawSlots: ArcSlot[], + surpriseUsage: { session: number; hour: number }, + sessionAgeMin: number, + ): ArcSlot[] { + // One surprise is allowed in a rolling hour. The session capacity grows + // slowly (one per listening hour), preserving a recovery anchor after + // every admitted surprise instead of filling each replan with novelty. + const sessionSurpriseLimit = Math.max(1, Math.ceil(Math.max(0, sessionAgeMin) / 60)); + let plannedSurprises = 0; + return rawSlots.map(slot => { + if (!slot.surprise) return slot; + const allowed = plannedSurprises < slot.surprise.maxPerPlan + && surpriseUsage.hour + plannedSurprises < slot.surprise.maxPerHour + && surpriseUsage.session + plannedSurprises < sessionSurpriseLimit; + if (allowed) { + plannedSurprises++; + return slot; + } + return { + ...slot, + role: slot.surprise.recoveryRole, + surprise: undefined, + targets: { ...(slot.targets ?? {}), novelty: { max: 0.35 } }, + }; + }); + } + async replan( userId: string, sessionId: string, diff --git a/backend/src/services/session-director.test.ts b/backend/src/services/session-director.test.ts index 227f417..581035d 100644 --- a/backend/src/services/session-director.test.ts +++ b/backend/src/services/session-director.test.ts @@ -1,5 +1,5 @@ import { describe, it, expect, vi } from 'vitest'; -import { mergeUniquePlan, selectConstrainedSequence, SessionDirector } from './session-director.service.js'; +import { mergeUniquePlan, scoreArcTransition, selectConstrainedSequence, SessionDirector } from './session-director.service.js'; import { DbService } from './db.service.js'; import { ALL_GENERATORS } from './generators.service.js'; @@ -122,9 +122,261 @@ describe('SessionDirector', () => { it('has valid role names', () => { const slots = director.getArcSlots('comfort', 20); - const validRoles = ['known', 'adjacent', 'favorite', 'similar', 'new', 'medium', 'high', 'peak', 'cooldown', 'soft', 'ambient', 'acoustic', 'slow']; + const validRoles = ['known', 'adjacent', 'favorite', 'similar', 'new', 'medium', 'high', 'peak', 'cooldown', 'soft', 'ambient', 'acoustic', 'slow', 'surprise']; slots.forEach(s => expect(validRoles).toContain(s.role)); }); + + it('creates measurable targets, callbacks, and one bounded surprise in the first arc cycle', () => { + const slots = director.getArcSlots('comfort', 20); + expect(slots.every(slot => slot.targets && Object.keys(slot.targets).length > 0)).toBe(true); + expect(slots.filter(slot => slot.surprise)).toHaveLength(1); + const anchor = slots.find(slot => slot.callback?.phase === 'anchor'); + const callback = slots.find(slot => slot.callback?.phase === 'return'); + expect(anchor?.callback?.id).toBe(callback?.callback?.id); + expect(anchor?.callback?.minSeparation).toBeGreaterThan(0); + }); + }); + + describe('durable surprise delivery accounting', () => { + it('counts only exact, served surprise plan-item exposures', async () => { + const sessionId = '00000000-0000-4000-8000-000000000001'; + const userId = '00000000-0000-4000-8000-000000000002'; + const otherSessionId = '00000000-0000-4000-8000-000000000003'; + const revisionOneId = '00000000-0000-4000-8000-000000000011'; + const revisionTwoId = '00000000-0000-4000-8000-000000000012'; + const otherRevisionId = '00000000-0000-4000-8000-000000000013'; + const now = new Date(); + const recentAt = new Date(now.getTime() - 5 * 60 * 1000); + const expiredAt = new Date(now.getTime() - 61 * 60 * 1000); + const retainedTrackId = '00000000-0000-4000-8000-000000000021'; + const servedTrackId = '00000000-0000-4000-8000-000000000022'; + const oldTrackId = '00000000-0000-4000-8000-000000000023'; + + // This mirrors the three tables involved in the query. The unserved + // retained row exists in both immutable revisions, but has no ledger + // event and therefore must not consume a surprise budget. + const versions = [ + { id: revisionOneId, sessionId }, + { id: revisionTwoId, sessionId }, + { id: otherRevisionId, sessionId: otherSessionId }, + ]; + const items = [ + { planVersionId: revisionOneId, ordinal: 6, trackId: retainedTrackId, slotRole: 'surprise' }, + { planVersionId: revisionTwoId, ordinal: 6, trackId: retainedTrackId, slotRole: 'surprise' }, + { planVersionId: revisionTwoId, ordinal: 7, trackId: servedTrackId, slotRole: 'surprise' }, + { planVersionId: revisionTwoId, ordinal: 8, trackId: oldTrackId, slotRole: 'surprise' }, + { planVersionId: revisionTwoId, ordinal: 9, trackId: '00000000-0000-4000-8000-000000000024', slotRole: 'favorite' }, + { planVersionId: otherRevisionId, ordinal: 7, trackId: servedTrackId, slotRole: 'surprise' }, + ]; + const events = [ + // The exact revision-two association counts once. + { id: '00000000-0000-4000-8000-000000000031', sessionId, userId, type: 'track_served', trackId: servedTrackId, payload: { planVersionId: revisionTwoId, ordinal: 7 }, occurredAt: recentAt }, + // A valid historical exposure remains in the session total but falls + // out of the rolling 60-minute counter. + { id: '00000000-0000-4000-8000-000000000032', sessionId, userId, type: 'track_served', trackId: oldTrackId, payload: { planVersionId: revisionTwoId, ordinal: 8 }, occurredAt: expiredAt }, + { id: '00000000-0000-4000-8000-000000000033', sessionId, userId, type: 'track_served', trackId: servedTrackId, payload: { planVersionId: revisionOneId, ordinal: 7 }, occurredAt: recentAt }, // wrong ordinal + { id: '00000000-0000-4000-8000-000000000034', sessionId, userId, type: 'track_served', trackId: retainedTrackId, payload: { planVersionId: revisionTwoId, ordinal: 7 }, occurredAt: recentAt }, // wrong track + { id: '00000000-0000-4000-8000-000000000035', sessionId, userId, type: 'track_served', trackId: servedTrackId, payload: { planVersionId: '00000000-0000-4000-8000-000000000014', ordinal: 7 }, occurredAt: recentAt }, // wrong version + { id: '00000000-0000-4000-8000-000000000036', sessionId: otherSessionId, userId, type: 'track_served', trackId: servedTrackId, payload: { planVersionId: otherRevisionId, ordinal: 7 }, occurredAt: recentAt }, // wrong session + { id: '00000000-0000-4000-8000-000000000037', sessionId, userId, type: 'track_finished', trackId: servedTrackId, payload: { planVersionId: revisionTwoId, ordinal: 7 }, occurredAt: recentAt }, + // Legacy/corrupt payloads must neither cast-fail nor claim an actual + // surprise exposure when somebody writes directly to the event ledger. + { id: '00000000-0000-4000-8000-000000000038', sessionId, userId, type: 'track_served', trackId: servedTrackId, payload: { planVersionId: 'not-a-uuid', ordinal: 7 }, occurredAt: recentAt }, + { id: '00000000-0000-4000-8000-000000000039', sessionId, userId, type: 'track_served', trackId: servedTrackId, payload: { planVersionId: revisionTwoId, ordinal: 'not-an-integer' }, occurredAt: recentAt }, + { id: '00000000-0000-4000-8000-000000000040', sessionId, userId, type: 'track_served', trackId: servedTrackId, payload: { planVersionId: revisionTwoId, ordinal: '999999999999999999999999999999999999' }, occurredAt: recentAt }, + ]; + const exposureIds: string[] = []; + const db = makeMockDb(); + (db.pgClient.query as any).mockImplementation((sql: string, params: unknown[]) => { + // Faithfully evaluate the query's joins against the in-memory rows; + // do not treat merely planned items as delivered exposure. + expect(params).toEqual([sessionId, userId]); + const matching = events.filter(event => { + const item = items.find(candidate => candidate.planVersionId === event.payload.planVersionId + && candidate.ordinal === event.payload.ordinal + && candidate.trackId === event.trackId); + const version = item && versions.find(candidate => candidate.id === item.planVersionId); + return event.sessionId === sessionId + && event.userId === userId + && event.type === 'track_served' + && item?.slotRole === 'surprise' + && version?.sessionId === event.sessionId; + }); + exposureIds.push(...new Set(matching.map(event => event.id))); + return Promise.resolve({ + rows: [{ + session_count: new Set(matching.map(event => event.id)).size, + hour_count: new Set(matching.filter(event => event.occurredAt > new Date(Date.now() - 60 * 60 * 1000)).map(event => event.id)).size, + }], + }); + }); + + const usage = await (new SessionDirector(db) as any).getSessionSurpriseUsage(sessionId, userId); + + expect(usage).toEqual({ session: 2, hour: 1 }); + expect(exposureIds).toEqual([ + '00000000-0000-4000-8000-000000000031', + '00000000-0000-4000-8000-000000000032', + ]); + const [sql, params] = (db.pgClient.query as any).mock.calls[0] as [string, unknown[]]; + expect(params).toEqual([sessionId, userId]); + expect(sql).toContain('SELECT DISTINCT e.id, e.occurred_at'); + expect(sql).toContain("e.payload->>'planVersionId' ~*"); + expect(sql).toContain("THEN (e.payload->>'planVersionId')::uuid"); + expect(sql).toContain("e.payload->>'ordinal') ~ '^(0|[1-9][0-9]{0,8})$'"); + expect(sql).toContain("THEN (e.payload->>'ordinal')::integer"); + expect(sql).toContain('AND item.track_id = e.track_id'); + expect(sql).toContain('AND version.session_id = e.session_id'); + expect(sql).toContain("e.type = 'track_served'"); + expect(sql).toContain("item.slot_role = 'surprise'"); + expect(sql).toContain('WHERE e.session_id = $1'); + expect(sql).toContain('AND e.user_id = $2'); + expect(sql).toContain("occurred_at > NOW() - INTERVAL '1 hour'"); + }); + }); + + describe('transition-aware sequence scoring', () => { + it('prefers a smooth, on-arc candidate and treats missing analysis as lower confidence', () => { + const slot = { + position: 0, + role: 'high', + targets: { energy: { min: 0.7, max: 0.9, maxDelta: 0.2 }, tempo: { min: 120, max: 160, maxDelta: 25 } }, + }; + const previous = { energy: 0.72, bpm: 132 }; + const smooth = scoreArcTransition(candidate('smooth'), { energy: 0.78, bpm: 140 }, previous, slot); + const abrupt = scoreArcTransition(candidate('abrupt'), { energy: 0.15, bpm: 72 }, previous, slot); + const unknown = scoreArcTransition(candidate('unknown'), {}, previous, slot); + + expect(smooth.score).toBeGreaterThan(abrupt.score); + expect(unknown.confidence).toBeLessThan(smooth.confidence); + expect(unknown.score).toBeGreaterThan(0); + }); + + it('selects a callback inside its separation window while preserving hard caps', () => { + const candidates = ['anchor', 'bridge-a', 'bridge-b', 'return', 'other'].map(id => ({ ...candidate(id), generatorId: 'comfort' })); + const metadata = new Map([ + ['anchor', { artistId: 'theme', albumId: 'a1', favorite: true, energy: 0.5 }], + ['bridge-a', { artistId: 'a2', albumId: 'a2', energy: 0.5 }], + ['bridge-b', { artistId: 'a3', albumId: 'a3', energy: 0.5 }], + ['return', { artistId: 'theme', albumId: 'a4', favorite: true, energy: 0.5 }], + ['other', { artistId: 'a4', albumId: 'a5', favorite: false, energy: 0.5 }], + ]); + const token = { id: 'theme', theme: 'artist' as const, minSeparation: 2, maxSeparation: 4 }; + const result = selectConstrainedSequence({ + candidates, + metadata, + budgets: [], + roleToGeneratorIds: () => ['comfort'], + slots: [ + { position: 0, role: 'known', targets: {}, callback: { ...token, phase: 'anchor' as const } }, + { position: 1, role: 'known', targets: {} }, + { position: 2, role: 'known', targets: {} }, + { position: 3, role: 'favorite', targets: {}, callback: { ...token, phase: 'return' as const } }, + ], + }); + expect(result.plan.map(item => item.trackId)).toEqual(['anchor', 'bridge-a', 'bridge-b', 'return']); + expect(result.plan[3].plan?.scoreBreakdown.callback).toBe(1); + }); + + it('plans an energetic rise through peak and cooldown when measured candidates exist', () => { + const db = makeMockDb(); + const director = new SessionDirector(db); + const slots = director.getArcSlots('energetic', 8); + const entries = [ + ['medium-1', 'comfort', 0.55, 110], ['medium-2', 'comfort', 0.62, 125], + ['high-1', 'discovery', 0.72, 135], ['high-2', 'discovery', 0.8, 150], + ['peak', 'deep-dive', 0.9, 165], ['surprise', 'discovery', 0.85, 155], + ['cooldown-1', 'comfort', 0.62, 130], ['cooldown-2', 'comfort', 0.5, 110], + ] as const; + const candidates = entries.map(([trackId, generatorId]) => ({ ...candidate(trackId), generatorId })); + const metadata: Map = new Map(entries.map(([trackId, generatorId, energy, bpm], index) => [trackId, { + artistId: `artist-${index}`, albumId: `album-${index}`, energy, bpm, + valence: 0.6, acousticness: trackId.startsWith('cooldown') ? 0.3 : 0.1, + favorite: String(trackId).startsWith('medium') || String(trackId).startsWith('cooldown'), + newArtist: generatorId === 'discovery', + }])); + const result = selectConstrainedSequence({ + candidates, slots, metadata, budgets: [], + roleToGeneratorIds: role => (director as any).roleToGeneratorIds(role), + }); + + expect(result.plan.map(item => item.trackId)).toEqual(entries.map(([trackId]) => trackId)); + expect(result.plan.map(item => metadata.get(item.trackId)?.energy)) + .toEqual([0.55, 0.62, 0.72, 0.8, 0.9, 0.85, 0.62, 0.5]); + expect(result.relaxations).toEqual([]); + }); + + it('keeps the familiar-new-familiar discovery callback intact', () => { + const db = makeMockDb(); + const director = new SessionDirector(db); + const candidates = [ + { ...candidate('favorite-anchor'), generatorId: 'deep-dive' }, + { ...candidate('adjacent'), generatorId: 'adjacent' }, + { ...candidate('new'), generatorId: 'discovery' }, + { ...candidate('favorite-return'), generatorId: 'deep-dive' }, + ]; + const metadata = new Map([ + ['favorite-anchor', { artistId: 'theme', albumId: 'a1', favorite: true, energy: 0.5 }], + ['adjacent', { artistId: 'bridge', albumId: 'a2', energy: 0.55 }], + ['new', { artistId: 'new', albumId: 'a3', newArtist: true, energy: 0.6 }], + ['favorite-return', { artistId: 'theme', albumId: 'a4', favorite: true, energy: 0.55 }], + ]); + const result = selectConstrainedSequence({ + candidates, slots: director.getArcSlots('discovery', 4), metadata, budgets: [], + roleToGeneratorIds: role => (director as any).roleToGeneratorIds(role), + }); + + expect(result.plan.map(item => item.trackId)) + .toEqual(['favorite-anchor', 'adjacent', 'new', 'favorite-return']); + expect(result.plan[3].plan?.explanation.callback).toMatchObject({ phase: 'return', matchScore: 1 }); + }); + + it('downgrades a surprise deterministically when no preferred recovery anchor is feasible', () => { + const db = makeMockDb(); + const director = new SessionDirector(db); + const candidates = [ + { ...candidate('favorite-anchor'), generatorId: 'deep-dive' }, + { ...candidate('adjacent'), generatorId: 'adjacent' }, + { ...candidate('unrecoverable-surprise'), generatorId: 'discovery' }, + // This may fill the downgraded favourite slot only through the + // explicitly persisted arc-source relaxation; it cannot reserve a + // recovery for the surprise because it is not a favourite source. + { ...candidate('fallback'), generatorId: 'contextual' }, + ]; + const metadata = new Map(candidates.map((item, index) => [item.trackId, { + artistId: `artist-${index}`, albumId: `album-${index}`, + favorite: item.trackId === 'favorite-anchor', newArtist: item.trackId === 'unrecoverable-surprise', energy: 0.5, + }])); + const result = selectConstrainedSequence({ + candidates, slots: director.getArcSlots('discovery', 3), metadata, budgets: [], + roleToGeneratorIds: role => (director as any).roleToGeneratorIds(role), + }); + + expect(result.plan.map(item => item.trackId)).toEqual(['favorite-anchor', 'adjacent', 'fallback']); + expect(result.plan[2].plan?.slotRole).toBe('favorite'); + expect(result.relaxations).toContainEqual(expect.objectContaining({ constraint: 'surprise_recovery' })); + }); + + it('falls back safely with sparse audio analysis and records arc precision', () => { + const db = makeMockDb(); + const director = new SessionDirector(db); + const result = selectConstrainedSequence({ + candidates: [ + { ...candidate('unknown-analysis'), generatorId: 'comfort' }, + { ...candidate('wrong-energy'), generatorId: 'comfort' }, + ], + slots: director.getArcSlots('energetic', 1), + metadata: new Map([ + ['unknown-analysis', { artistId: 'a1', albumId: 'x1' }], + ['wrong-energy', { artistId: 'a2', albumId: 'x2', energy: 0.1, bpm: 70 }], + ]), + budgets: [], roleToGeneratorIds: role => (director as any).roleToGeneratorIds(role), + }); + + expect(result.plan.map(item => item.trackId)).toEqual(['unknown-analysis']); + expect(result.relaxations).toContainEqual(expect.objectContaining({ constraint: 'arc_precision' })); + expect(result.plan[0].plan?.explanation.arcPrecision).toMatchObject({ measuredFit: false }); + }); }); describe('computeEntropy', () => { diff --git a/backend/src/services/vibe-session-coordinator.service.test.ts b/backend/src/services/vibe-session-coordinator.service.test.ts index a1f5e4f..ecf149e 100644 --- a/backend/src/services/vibe-session-coordinator.service.test.ts +++ b/backend/src/services/vibe-session-coordinator.service.test.ts @@ -54,6 +54,7 @@ function setup() { const director = { buildPlan: vi.fn().mockResolvedValue([{ trackId: TRACK_ID, generatorId: 'comfort', relevance: 0.8, explanation: [] }]), buildState: vi.fn().mockResolvedValue({ energy: 0.5, noveltyHunger: 0.3 }), + replan: vi.fn().mockResolvedValue([{ trackId: TRACK_ID, generatorId: 'comfort', relevance: 0.8, explanation: [] }]), } as any; return { db, director, coordinator: new VibeSessionCoordinator(db, director) }; } @@ -80,6 +81,31 @@ describe('VibeSessionCoordinator', () => { .toEqual(['session_started']); }); + it('persists a director-selected arc role and explainable sequence score', async () => { + const { db, director, coordinator } = setup(); + (director.buildPlan as any).mockResolvedValueOnce([{ + trackId: TRACK_ID, generatorId: 'discovery', relevance: 0.7, explanation: [{ predicate: 'near' }], + plan: { + slotRole: 'surprise', score: 0.82, + scoreBreakdown: { relevance: 0.7, transition: 0.9, arcTarget: 0.85 }, + explanation: { arcRole: 'surprise', surprise: { recoveryRole: 'favorite' } }, + }, + }]); + + await coordinator.start('user-1', {}); + + expect(db.publishVibePlan).toHaveBeenCalledWith(expect.objectContaining({ + items: [expect.objectContaining({ + slot_role: 'surprise', score: 0.82, + score_breakdown: expect.objectContaining({ transition: 0.9 }), + explanation: expect.objectContaining({ + paths: [{ predicate: 'near' }], + planner: expect.objectContaining({ arcRole: 'surprise' }), + }), + })], + })); + }); + it('returns the canonical replacement on an idempotent material-event retry without replanning', async () => { const { db, coordinator } = setup(); (db.recordVibeEvent as any).mockResolvedValueOnce({ @@ -118,13 +144,42 @@ describe('VibeSessionCoordinator', () => { const { db, director, coordinator } = setup(); const response = await coordinator.appendEvent('user-1', SESSION_ID, { type: 'completed', trackId: TRACK_ID }); - expect(director.buildPlan).toHaveBeenCalledWith('user-1', SESSION_ID, TRACK_ID); + expect(director.replan).toHaveBeenCalledWith( + 'user-1', SESSION_ID, expect.any(Array), [TRACK_ID], TRACK_ID, + { excludedTrackIds: new Set() }, + ); expect(db.publishVibePlan).toHaveBeenCalledWith(expect.objectContaining({ sessionId: SESSION_ID, reason: 'feedback:completed', items: [expect.objectContaining({ committed: false })], })); expect(response).toMatchObject({ replanned: true, replanReason: 'feedback:completed', planVersion: 2 }); }); + it('passes the durable unserved callback tail into the retention-aware replan', async () => { + const { db, director, coordinator } = setup(); + (db.getVibePlan as any).mockResolvedValue({ + ...plan(), + items: [{ + ...plan().items[0], + slot_role: 'favorite', + explanation: { + paths: [], + planner: { arcRole: 'favorite', callback: { id: 'theme:0', phase: 'return' } }, + }, + }], + }); + + await coordinator.appendEvent('user-1', SESSION_ID, { type: 'completed', trackId: 'other-track' }); + + expect(director.replan).toHaveBeenCalledWith( + 'user-1', SESSION_ID, + [expect.objectContaining({ + trackId: TRACK_ID, + plan: expect.objectContaining({ slotRole: 'favorite', explanation: expect.objectContaining({ arcRole: 'favorite' }) }), + })], + ['other-track'], 'other-track', { excludedTrackIds: new Set() }, + ); + }); + it('keeps the durable seed excluded when feedback supplies a different local replan anchor', async () => { const { db, director, coordinator } = setup(); const seedTrackId = '44444444-4444-4444-8444-444444444444'; @@ -132,9 +187,11 @@ describe('VibeSessionCoordinator', () => { await coordinator.appendEvent('user-1', SESSION_ID, { type: 'completed', trackId: TRACK_ID }); - expect(director.buildPlan).toHaveBeenCalledWith( + expect(director.replan).toHaveBeenCalledWith( 'user-1', SESSION_ID, + expect.any(Array), + [TRACK_ID], TRACK_ID, { excludedTrackIds: new Set([seedTrackId]) }, ); diff --git a/backend/src/services/vibe-session-coordinator.service.ts b/backend/src/services/vibe-session-coordinator.service.ts index c52f846..39779d7 100644 --- a/backend/src/services/vibe-session-coordinator.service.ts +++ b/backend/src/services/vibe-session-coordinator.service.ts @@ -1,5 +1,6 @@ import { DbService, VibeEvent, VibePlan, VibeSession } from './db.service.js'; import { SessionDirector } from './session-director.service.js'; +import { Candidate } from './generators.service.js'; /** * This is deliberately a narrow bridge between the durable Vibe ledger and @@ -67,7 +68,7 @@ export interface VibeSessionResponse { export class VibeSessionCoordinator { constructor( private readonly db: DbService, - private readonly director: Pick, + private readonly director: Pick, ) {} async start(userId: string, input: StartVibeSessionInput): Promise { @@ -111,15 +112,18 @@ export class VibeSessionCoordinator { policyVersion, intent: input.intent ?? null, horizonTracks: candidates.length, + ...(candidates[0]?.plan?.objective ?? {}), }, items: candidates.map((candidate, ordinal) => ({ ordinal, track_id: candidate.trackId, - slot_role: ordinal === 0 ? 'next' : null, + slot_role: candidate.plan?.slotRole ?? (ordinal === 0 ? 'next' : null), candidate_source: candidate.generatorId, - score: candidate.relevance, - score_breakdown: { relevance: candidate.relevance }, - explanation: candidate.explanation, + score: candidate.plan?.score ?? candidate.relevance, + score_breakdown: candidate.plan?.scoreBreakdown ?? { relevance: candidate.relevance }, + explanation: candidate.plan + ? { paths: candidate.explanation, planner: candidate.plan.explanation } + : candidate.explanation, committed: false, })), }); @@ -216,11 +220,24 @@ export class VibeSessionCoordinator { // feedback tracks, the seed is not necessarily present in the event // ledger, so carry it explicitly into every replacement request. const seedTrackId = session.seed_track_id ?? undefined; - const candidates = seedTrackId - ? await this.director.buildPlan(userId, sessionId, input.trackId ?? seedTrackId, { - excludedTrackIds: new Set([seedTrackId]), - }) - : await this.director.buildPlan(userId, sessionId, input.trackId); + // Revisions retain the durable, unserved queue tail rather than building + // an unrelated plan after every signal. Besides reducing churn, this + // preserves a valid callback/recovery pair that has already been shown + // to the client while allowing the director to refill under the same + // hard caps and current feedback state. + const current = await this.db.getVibePlan(sessionId, userId); + const retained = current ? this.unservedCandidates(current) : []; + const excludedTrackIds = new Set([ + ...(seedTrackId ? [seedTrackId] : []), + ]); + const candidates = await this.director.replan( + userId, + sessionId, + retained, + input.trackId ? [input.trackId] : [], + input.trackId ?? seedTrackId, + { excludedTrackIds }, + ); const reason = `feedback:${input.type}`; const plan = await this.db.publishVibePlan({ sessionId, @@ -232,15 +249,18 @@ export class VibeSessionCoordinator { feedbackEventId: result.event.id, feedbackType: input.type, horizonTracks: candidates.length, + ...(candidates[0]?.plan?.objective ?? {}), }, items: candidates.map((candidate, ordinal) => ({ ordinal, track_id: candidate.trackId, - slot_role: ordinal === 0 ? 'next' : null, + slot_role: candidate.plan?.slotRole ?? (ordinal === 0 ? 'next' : null), candidate_source: candidate.generatorId, - score: candidate.relevance, - score_breakdown: { relevance: candidate.relevance }, - explanation: candidate.explanation, + score: candidate.plan?.score ?? candidate.relevance, + score_breakdown: candidate.plan?.scoreBreakdown ?? { relevance: candidate.relevance }, + explanation: candidate.plan + ? { paths: candidate.explanation, planner: candidate.plan.explanation } + : candidate.explanation, committed: false, })), }); @@ -303,6 +323,37 @@ export class VibeSessionCoordinator { }; } + /** Reconstruct the planner envelope from the durable revision. Older + * revisions stored only paths, so they remain valid retention inputs. */ + private unservedCandidates(plan: VibePlan): Candidate[] { + return plan.items + .filter(item => !item.committed) + .map(item => { + const stored = item.explanation; + const hasPlanner = !!stored && !Array.isArray(stored) && typeof stored === 'object' + && 'planner' in stored; + const object = hasPlanner ? stored as { paths?: Candidate['explanation']; planner?: Record } : undefined; + const planner = object?.planner; + return { + trackId: item.track_id, + generatorId: item.candidate_source, + relevance: item.score, + explanation: object?.paths ?? (Array.isArray(stored) ? stored : []), + plan: planner ? { + slotRole: item.slot_role ?? 'retained', + score: item.score, + scoreBreakdown: item.score_breakdown, + explanation: planner, + objective: { + policy: planner.policy, + constraints: planner.constraints, + relaxations: planner.relaxations, + }, + } : undefined, + }; + }); + } + private mapLifecycleError(error: unknown): Error { if (error instanceof Error && (error.message.includes('Cannot record a new event for') || error.message.includes('Cannot resume ') || error.message.includes('Cannot publish a plan for'))) { return new VibeSessionLifecycleError(error.message);