initial state: muzick music player + recommendation engine

This commit is contained in:
kami
2026-07-14 01:35:52 +04:00
commit 737bf19fd1
196 changed files with 32431 additions and 0 deletions
+182
View File
@@ -0,0 +1,182 @@
import { describe, it, expect, vi } from 'vitest';
import { DbService } from './db.service.js';
function makeService(): { service: DbService; mockQuery: ReturnType<typeof vi.fn> } {
const mockQuery = vi.fn();
const service = new DbService({ query: mockQuery } as any);
return { service, mockQuery };
}
describe('DbService v2 methods', () => {
describe('upsertClaim', () => {
it('calls INSERT ... ON CONFLICT with correct parameters', async () => {
const { service, mockQuery } = makeService();
mockQuery.mockResolvedValue({ rows: [{ id: 'claim-1' }] });
const id = await service.upsertClaim({
subject_type: 'track',
subject_id: 'track-1',
predicate: 'credited_main_on',
object_type: 'artist',
object_id: 'artist-1',
source: 'mb',
confidence: 1.0,
});
expect(id).toBe('claim-1');
expect(mockQuery).toHaveBeenCalledTimes(1);
const [sql, params] = mockQuery.mock.calls[0];
expect(sql).toContain('INSERT INTO claims');
expect(sql).toContain('ON CONFLICT');
expect(params).toContain('track');
expect(params).toContain('track-1');
expect(params).toContain('credited_main_on');
});
it('handles user_id null for objective claims', async () => {
const { service, mockQuery } = makeService();
mockQuery.mockResolvedValue({ rows: [{ id: 'c1' }] });
await service.upsertClaim({
subject_type: 'artist', subject_id: 'a1', predicate: 'alias_of',
object_type: 'artist', object_id: 'a2', source: 'listener_behavior',
user_id: 'user-1',
});
const params = mockQuery.mock.calls[0][1];
expect(params[0]).toBe('user-1');
});
});
describe('getClaimsBySubject', () => {
it('filters by subject type and id', async () => {
const { service, mockQuery } = makeService();
mockQuery.mockResolvedValue({ rows: [] });
await service.getClaimsBySubject('track', 'track-1');
const [sql, params] = mockQuery.mock.calls[0];
expect(sql).toContain('subject_type = $1');
expect(sql).toContain('subject_id = $2');
expect(params).toEqual(['track', 'track-1']);
});
it('optionally filters by predicate and user_id', async () => {
const { service, mockQuery } = makeService();
mockQuery.mockResolvedValue({ rows: [] });
await service.getClaimsBySubject('artist', 'a1', 'alias_of', 'user-1');
const [sql] = mockQuery.mock.calls[0];
expect(sql).toContain('predicate');
expect(sql).toContain('user_id IS NULL');
});
});
describe('recordEvidence', () => {
it('appends evidence row', async () => {
const { service, mockQuery } = makeService();
mockQuery.mockResolvedValue({ rows: [{ id: 'ev-1' }] });
const id = await service.recordEvidence({
user_id: 'user-1', entity_type: 'track', entity_id: 'track-1',
signal: 'playback_completed', profile: 'longterm', weight: 0.10,
});
expect(id).toBe('ev-1');
const [sql] = mockQuery.mock.calls[0];
expect(sql).toContain('INSERT INTO evidence');
});
});
describe('updateListenerBelief', () => {
it('UPSERTs with delta formula', async () => {
const { service, mockQuery } = makeService();
mockQuery.mockResolvedValue({ rowCount: 1 });
await service.updateListenerBelief({
user_id: 'user-1', profile: 'longterm',
entity_type: 'track', entity_id: 'track-1',
dimension: 'affinity', value_delta: 0.10,
});
const [sql] = mockQuery.mock.calls[0];
expect(sql).toContain('INSERT INTO listener_beliefs');
expect(sql).toContain('ON CONFLICT');
expect(sql).toContain('GREATEST(-1.0, LEAST(1.0');
});
});
describe('recordEvidence → belief derivation wiring', () => {
it('derives a fresh longterm affinity belief after playback_completed', async () => {
const { service, mockQuery } = makeService();
// first call: INSERT evidence → returns id; second call: UPSERT belief → rowCount 1
mockQuery
.mockResolvedValueOnce({ rows: [{ id: 'ev-1' }] })
.mockResolvedValueOnce({ rowCount: 1 });
const id = await service.recordEvidence({
user_id: 'user-1', entity_type: 'track', entity_id: 'track-1',
signal: 'playback_completed', profile: 'longterm', weight: 0.10,
});
expect(id).toBe('ev-1');
expect(mockQuery).toHaveBeenCalledTimes(2);
// First call INSERTs the evidence row.
const [evidSql, evidParams] = mockQuery.mock.calls[0];
expect(evidSql).toContain('INSERT INTO evidence');
expect(evidParams[4]).toBe('longterm'); // profile
expect(evidParams[5]).toBe(0.10); // weight
// Second call UPSERTs the matching listener_belief. For a fresh
// belief the INSERT path sets value = weight directly (per spec §B.4
// INSERT branch), so the resulting row has value=0.10, confidence=0.05.
const [beliefSql, beliefParams] = mockQuery.mock.calls[1];
expect(beliefSql).toContain('INSERT INTO listener_beliefs');
expect(beliefSql).toContain('ON CONFLICT');
// [user_id, profile, entity_type, entity_id, dimension, value_delta, confidence_delta]
expect(beliefParams[0]).toBe('user-1');
expect(beliefParams[1]).toBe('longterm');
expect(beliefParams[2]).toBe('track');
expect(beliefParams[3]).toBe('track-1');
expect(beliefParams[4]).toBe('affinity');
expect(beliefParams[5]).toBe(0.10); // value_delta = weight
expect(beliefParams[6]).toBe(0.05); // confidence_delta default
});
it('maps play_of_never_seen to the novelty_tolerance dimension', async () => {
const { service, mockQuery } = makeService();
mockQuery
.mockResolvedValueOnce({ rows: [{ id: 'ev-2' }] })
.mockResolvedValueOnce({ rowCount: 1 });
await service.recordEvidence({
user_id: 'user-1', entity_type: 'track', entity_id: 'track-2',
signal: 'play_of_never_seen', profile: 'discovery', weight: 0.05,
});
const beliefParams = mockQuery.mock.calls[1][1] as unknown[];
expect(beliefParams[4]).toBe('novelty_tolerance');
expect(beliefParams[1]).toBe('discovery');
expect(beliefParams[5]).toBe(0.05);
});
it('still appends an evidence row before deriving the belief', async () => {
const { service, mockQuery } = makeService();
mockQuery
.mockResolvedValueOnce({ rows: [{ id: 'ev-3' }] })
.mockResolvedValueOnce({ rowCount: 1 });
const id = await service.recordEvidence({
user_id: 'user-1', entity_type: 'track', entity_id: 'track-3',
signal: 'skip_quick', profile: 'negative', weight: -0.20,
});
expect(id).toBe('ev-3');
expect(mockQuery.mock.calls[0][0]).toContain('INSERT INTO evidence');
expect(mockQuery.mock.calls[1][0]).toContain('listener_beliefs');
expect(mockQuery.mock.calls[1][1][4]).toBe('affinity');
});
});
describe('getFusedTrackArtists', () => {
it('reads from claim_fusion view', async () => {
const { service, mockQuery } = makeService();
mockQuery.mockResolvedValue({ rows: [{ id: 'a1', name: 'Artist 1', role: 'main', confidence: 0.9 }] });
const result = await service.getFusedTrackArtists('track-1');
const [sql] = mockQuery.mock.calls[0];
expect(sql).toContain('claim_fusion');
expect(result).toHaveLength(1);
expect(result[0].role).toBe('main');
});
});
});
File diff suppressed because it is too large Load Diff
+297
View File
@@ -0,0 +1,297 @@
import { DbService } from './db.service.js';
export interface DiscoveryCandidate {
id: string;
source: string;
externalId: string;
title: string | null;
artistCredit: unknown;
notes: unknown;
status: string;
relevance: number;
explanation: string;
}
export class DiscoveryService {
constructor(private db: DbService) {}
// ---------------------------------------------------------------
// E.1 — Graph exploration: walk the graph beyond the library
// ---------------------------------------------------------------
async walkGraphForDiscovery(userId: string): Promise<number> {
const beliefs = await this.db.getListenerBeliefs({
userId,
profile: 'longterm',
entityType: 'artist',
dimension: 'affinity',
limit: 100,
orderBy: 'value',
order: 'DESC',
});
const highAffinity = beliefs.filter((b) => b.value > 0.3);
let newCount = 0;
for (const belief of highAffinity) {
const candidates = await this.db.pgClient.query<{ candidate_artist_id: string }>(
`SELECT cf.object_id AS candidate_artist_id
FROM claim_fusion cf
WHERE cf.subject_id = $1::uuid
AND cf.predicate IN ('same_scene_as', 'featured_on')
AND cf.object_type = 'artist'
AND NOT EXISTS (
SELECT 1 FROM tracks t
JOIN claim_fusion cf2 ON cf2.subject_id = t.id
WHERE cf2.object_id = cf.object_id
AND cf2.predicate = 'credited_main_on'
)
LIMIT 20`,
[belief.entity_id]
);
for (const row of candidates.rows) {
const dcRes = await this.db.pgClient.query<{ id: string }>(
`INSERT INTO discovery_candidates (source, external_id, artist_credit, notes)
VALUES ($1, $2, $3, $4)
ON CONFLICT (source, external_id) DO NOTHING
RETURNING id`,
[
'graph_exploration',
row.candidate_artist_id,
JSON.stringify([{ artist_id: row.candidate_artist_id }]),
JSON.stringify({
discovery_source: 'graph_exploration',
path: [
{
entity_id: belief.entity_id,
predicate: 'affinity_source',
profile: 'longterm',
affinity: belief.value,
},
{
entity_id: row.candidate_artist_id,
predicate: 'same_scene_as',
},
],
source_artist_belief_id: belief.entity_id,
}),
]
);
if (dcRes.rows.length === 0) continue;
const dcId = dcRes.rows[0].id;
const relevance = Math.min(belief.value, 0.8);
await this.db.upsertClaim({
subject_type: 'track',
subject_id: dcId,
predicate: 'discovery_candidate',
object_type: 'artist',
object_id: row.candidate_artist_id,
source: 'graph_exploration',
confidence: relevance,
raw: {
discovery_source: 'graph_exploration',
path: [
{ entity_id: belief.entity_id, relationship: 'affinity_source', belief_value: belief.value },
{ entity_id: row.candidate_artist_id, relationship: 'same_scene_as' },
],
},
});
newCount++;
}
}
return newCount;
}
// ---------------------------------------------------------------
// E.3 — Evaluate discovery candidates for acquisition
// ---------------------------------------------------------------
async evalCandidates(
userId: string,
limit?: number
): Promise<{ candidateId: string; shouldAcquire: boolean; reason: string }[]> {
const cap = limit ?? 20;
const results: { candidateId: string; shouldAcquire: boolean; reason: string }[] = [];
const candidates = await this.db.pgClient.query(
`SELECT * FROM discovery_candidates
WHERE status = 'candidate'
ORDER BY first_seen_at ASC
LIMIT $1`,
[cap]
);
const backlogRes = await this.db.pgClient.query(
`SELECT COUNT(*)::int AS cnt FROM discovery_candidates WHERE status = 'acquiring'`
);
let backlog = backlogRes.rows[0]?.cnt as number ?? 0;
for (const row of candidates.rows) {
const claimRes = await this.db.pgClient.query<{ object_id: string; fused_value: number }>(
`SELECT object_id, fused_value
FROM claim_fusion
WHERE subject_type = 'track' AND subject_id = $1::uuid
AND predicate = 'discovery_candidate'
LIMIT 1`,
[row.id]
);
const relevance = claimRes.rows[0]?.fused_value ?? 0;
const candidateArtistId = claimRes.rows[0]?.object_id;
const noveltyBeliefs = await this.db.getListenerBeliefs({
userId,
profile: 'discovery',
entityType: 'artist',
entityId: candidateArtistId,
dimension: 'tolerance',
limit: 1,
});
const tolerance = noveltyBeliefs.length > 0 ? noveltyBeliefs[0].value : 0.5;
let artistCount = 0;
if (candidateArtistId) {
const acRes = await this.db.pgClient.query(
`SELECT COUNT(*)::int AS cnt
FROM discovery_candidates dc
JOIN claims c ON c.subject_id = dc.id
WHERE dc.status = 'acquiring'
AND c.predicate = 'discovery_candidate'
AND c.object_id = $1::uuid`,
[candidateArtistId]
);
artistCount = acRes.rows[0]?.cnt as number ?? 0;
}
const shouldAcquire = relevance > 0.3 && tolerance > 0.2 && backlog < 20 && artistCount < 3;
let reason: string;
if (shouldAcquire) {
await this.db.pgClient.query(
`UPDATE discovery_candidates SET status = 'acquiring', last_eval_at = NOW() WHERE id = $1`,
[row.id]
);
backlog++;
reason = 'meets criteria';
} else {
if (relevance <= 0.3) reason = 'relevance too low';
else if (tolerance <= 0.2) reason = 'novelty tolerance exceeded';
else if (backlog >= 20) reason = 'backlog full';
else if (artistCount >= 3) reason = 'artist diversity limit';
else reason = 'unknown';
await this.db.pgClient.query(
`UPDATE discovery_candidates SET status = 'retired', last_eval_at = NOW() WHERE id = $1`,
[row.id]
);
}
results.push({
candidateId: row.id,
shouldAcquire,
reason,
});
}
return results;
}
// ---------------------------------------------------------------
// E.4 — Probation lifecycle
// ---------------------------------------------------------------
async evalProbation(trackId: string): Promise<'retained' | 'retired' | 'probation'> {
const completedRes = await this.db.pgClient.query(
`SELECT COUNT(*)::int AS cnt FROM evidence
WHERE entity_type = 'track' AND entity_id = $1 AND signal = 'playback_completed'`,
[trackId]
);
const completedPlays = completedRes.rows[0]?.cnt as number ?? 0;
const skipRes = await this.db.pgClient.query(
`SELECT COUNT(*)::int AS cnt FROM evidence
WHERE entity_type = 'track' AND entity_id = $1 AND signal = 'skip_quick'`,
[trackId]
);
const skips = skipRes.rows[0]?.cnt as number ?? 0;
const trackRes = await this.db.pgClient.query<{ probation_entered_at: Date | null }>(
`SELECT probation_entered_at FROM tracks WHERE id = $1`,
[trackId]
);
const probTrack = trackRes.rows[0];
if (completedPlays >= 3) {
await this.db.pgClient.query(
`UPDATE tracks SET probation_status = 'retained' WHERE id = $1`,
[trackId]
);
const claimRes = await this.db.pgClient.query<{ source: string }>(
`SELECT source FROM claims
WHERE subject_type = 'track' AND subject_id = $1 AND predicate = 'discovery_candidate'
LIMIT 1`,
[trackId]
);
if (claimRes.rows[0]) {
await this.db.pgClient.query(
`UPDATE source_trust SET trust = LEAST(1.0, trust + 0.05) WHERE key = $1`,
[claimRes.rows[0].source]
);
}
return 'retained';
}
const daysSinceProbation = probTrack?.probation_entered_at
? (Date.now() - new Date(probTrack.probation_entered_at).getTime()) / (1000 * 86400)
: 0;
if (completedPlays === 0 && skips >= 3 && daysSinceProbation > 7) {
await this.db.pgClient.query(
`UPDATE tracks SET probation_status = 'retired' WHERE id = $1`,
[trackId]
);
return 'retired';
}
return 'probation';
}
async sweepProbation(): Promise<{ retained: number; retired: number }> {
const res = await this.db.pgClient.query(
`SELECT id FROM tracks WHERE probation_status = 'probation'`
);
let retained = 0;
let retired = 0;
for (const row of res.rows) {
const result = await this.evalProbation(row.id as string);
if (result === 'retained') retained++;
else if (result === 'retired') retired++;
}
return { retained, retired };
}
// ---------------------------------------------------------------
// E.5 — Meta-learning stub
// ---------------------------------------------------------------
async runMetaLearning(): Promise<void> {
const res = await this.db.pgClient.query(
`SELECT c.source, COUNT(*)::int AS cnt
FROM claims c
JOIN tracks t ON t.id = c.subject_id
WHERE c.predicate = 'discovery_candidate'
AND t.probation_status = 'retained'
GROUP BY c.source
ORDER BY cnt DESC`
);
console.log('[MetaLearning] Discovery source retention counts:', JSON.stringify(res.rows));
}
}
+500
View File
@@ -0,0 +1,500 @@
import { DbService, ListenerBelief } from './db.service.js';
// ---------------------------------------------------------------------------
// System C — Candidate Generators
// Each generator returns candidates with graph-path explanations.
// No scoring — the session director handles ranking.
// ---------------------------------------------------------------------------
export interface ClaimEdge {
subjectType: string;
subjectId: string;
predicate: string;
objectType: string;
objectId: string;
fusedValue: number;
}
export interface Candidate {
trackId: string;
generatorId: string;
explanation: ClaimEdge[];
relevance: number;
}
export interface GeneratorContext {
userId: string;
seedTrackId: string | null;
seedArtistId: string | null;
beliefs: ListenerBelief[];
recentExclusions: string[];
toleranceMap: Record<string, number>;
state: {
energy: number;
lastArtistIds: string[];
lastGenreIds: string[];
context: string | null;
noveltyHunger: number;
sessionAgeMin: number;
};
}
export type Generator = (db: DbService, ctx: GeneratorContext) => Promise<Candidate[]>;
const OBJECTIVE_USER = '00000000-0000-0000-0000-000000000000';
// ---------------------------------------------------------------------------
// 1. COMFORT — Top artists by longterm affinity > 0.5
// ---------------------------------------------------------------------------
async function comfortGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
const topArtists = ctx.beliefs
.filter(b => b.profile === 'longterm' && b.entity_type === 'artist' && b.dimension === 'affinity' && b.value > 0.5)
.sort((a, b) => b.value - a.value)
.slice(0, 20);
const candidates: Candidate[] = [];
for (const belief of topArtists) {
const res = await db.pgClient.query(
`SELECT t.id
FROM tracks t
JOIN claim_fusion cf ON cf.subject_type = 'track' AND cf.subject_id = t.id
AND cf.predicate IN ('credited_main_on', 'featured_on')
AND cf.object_type = 'artist' AND cf.object_id = $1
AND (cf.user_id = $2 OR cf.user_id = $3)
WHERE t.state = 'LIBRARY'
AND NOT (t.id = ANY($4::uuid[]))
ORDER BY cf.fused_value DESC
LIMIT 2`,
[belief.entity_id, OBJECTIVE_USER, ctx.userId, ctx.recentExclusions]
);
for (const row of res.rows as { id: string }[]) {
candidates.push({
trackId: row.id,
generatorId: 'comfort',
explanation: [{
subjectType: 'artist',
subjectId: belief.entity_id,
predicate: 'credited_main_on',
objectType: 'track',
objectId: row.id,
fusedValue: belief.value,
}],
relevance: belief.value,
});
}
}
return candidates;
}
// ---------------------------------------------------------------------------
// 2. ADJACENT — Walk graph from seed artist, exclude comfort pool
// ---------------------------------------------------------------------------
async function adjacentGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
if (!ctx.seedArtistId) return [];
const comfortArtistIds = new Set(
ctx.beliefs
.filter(b => b.profile === 'longterm' && b.entity_type === 'artist' && b.dimension === 'affinity' && b.value > 0.5)
.map(b => b.entity_id)
);
// Walk: seedArtist -> (credited_main_on|featured_on) -> track -> (credited_main_on|featured_on) -> reachedArtist
// cf1 finds tracks where seed artist appears; cf2 finds OTHER artists on those same tracks
const reachedRes = await db.pgClient.query(
`SELECT DISTINCT cf2.object_id AS artist_id
FROM claim_fusion cf1
JOIN claim_fusion cf2 ON cf2.subject_type = 'track'
AND cf2.subject_id = cf1.subject_id
AND cf2.predicate IN ('credited_main_on', 'featured_on')
AND cf2.object_type = 'artist'
AND cf2.object_id != $1
AND (cf2.user_id = $2 OR cf2.user_id = $3)
WHERE cf1.subject_type = 'track'
AND cf1.predicate IN ('credited_main_on', 'featured_on')
AND cf1.object_type = 'artist'
AND cf1.object_id = $1
AND (cf1.user_id = $2 OR cf1.user_id = $3)
LIMIT 30`,
[ctx.seedArtistId, OBJECTIVE_USER, ctx.userId]
);
const reachedArtistIds = (reachedRes.rows as { artist_id: string }[])
.map(r => r.artist_id)
.filter(id => !comfortArtistIds.has(id));
if (reachedArtistIds.length === 0) return [];
const trackRes = await db.pgClient.query(
`SELECT id FROM (
SELECT DISTINCT t.id
FROM tracks t
JOIN claim_fusion cf ON cf.subject_type = 'track' AND cf.subject_id = t.id
AND cf.predicate IN ('credited_main_on', 'featured_on')
AND cf.object_type = 'artist'
AND cf.object_id = ANY($1::uuid[])
AND (cf.user_id = $2 OR cf.user_id = $3)
WHERE t.state = 'LIBRARY'
AND NOT (t.id = ANY($4::uuid[]))
) sub
ORDER BY RANDOM()
LIMIT 20`,
[reachedArtistIds, OBJECTIVE_USER, ctx.userId, ctx.recentExclusions]
);
return (trackRes.rows as { id: string }[]).map(row => ({
trackId: row.id,
generatorId: 'adjacent',
explanation: [{
subjectType: 'artist',
subjectId: ctx.seedArtistId!,
predicate: 'credited_main_on',
objectType: 'track',
objectId: row.id,
fusedValue: 0.6,
}],
relevance: 0.6,
}));
}
// ---------------------------------------------------------------------------
// 3. DISCOVERY — Unfamiliar artists via graph edges from trusted artists
// ---------------------------------------------------------------------------
async function discoveryGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
const trustedIds = ctx.beliefs
.filter(b => b.profile === 'longterm' && b.entity_type === 'artist' && b.dimension === 'affinity' && b.value > 0.3)
.map(b => b.entity_id);
if (trustedIds.length === 0) return [];
const noveltyTolerance = ctx.toleranceMap.novelty_tolerance ?? 0.3;
const maxCandidates = Math.max(1, Math.floor(10 * noveltyTolerance));
const unfamiliarRes = await db.pgClient.query(
`SELECT DISTINCT cf.object_id AS artist_id
FROM claim_fusion cf
WHERE cf.subject_type = 'artist'
AND cf.subject_id = ANY($1::uuid[])
AND cf.predicate IN ('same_scene_as', 'same_label_as', 'produced')
AND cf.object_type = 'artist'
AND NOT EXISTS (
SELECT 1 FROM listener_beliefs lb
WHERE lb.user_id = $2
AND lb.entity_type = 'artist'
AND lb.entity_id = cf.object_id
AND lb.profile IN ('longterm', 'obsession')
)
LIMIT 30`,
[trustedIds, ctx.userId]
);
const unfamiliarArtistIds = (unfamiliarRes.rows as { artist_id: string }[]).map(r => r.artist_id);
if (unfamiliarArtistIds.length === 0) return [];
const trackRes = await db.pgClient.query(
`SELECT id FROM (
SELECT DISTINCT t.id
FROM tracks t
JOIN claim_fusion cf ON cf.subject_type = 'track' AND cf.subject_id = t.id
AND cf.predicate IN ('credited_main_on', 'featured_on')
AND cf.object_type = 'artist'
AND cf.object_id = ANY($1::uuid[])
WHERE t.state = 'LIBRARY'
AND NOT (t.id = ANY($2::uuid[]))
) sub
ORDER BY RANDOM()
LIMIT $3`,
[unfamiliarArtistIds, ctx.recentExclusions, maxCandidates]
);
return (trackRes.rows as { id: string }[]).map(row => ({
trackId: row.id,
generatorId: 'discovery',
explanation: [{
subjectType: 'artist',
subjectId: unfamiliarArtistIds[0],
predicate: 'credited_main_on',
objectType: 'track',
objectId: row.id,
fusedValue: 0.4,
}],
relevance: 0.4,
}));
}
// ---------------------------------------------------------------------------
// 4. DEEP-DIVE — Albums from obsession artists, unplayed tracks first
// ---------------------------------------------------------------------------
async function deepDiveGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
const obsessedIds = ctx.beliefs
.filter(b => b.profile === 'obsession' && b.entity_type === 'artist' && b.dimension === 'affinity' && b.value > 0.3)
.map(b => b.entity_id);
if (obsessedIds.length === 0) return [];
const albumRes = await db.pgClient.query(
`SELECT al.id AS album_id, al.artist_id
FROM albums al
WHERE al.artist_id = ANY($1::uuid[])
ORDER BY al.year ASC NULLS LAST, al.title ASC
LIMIT 20`,
[obsessedIds]
);
const candidates: Candidate[] = [];
for (const album of albumRes.rows as { album_id: string; artist_id: string }[]) {
const trackRes = await db.pgClient.query(
`SELECT t.id
FROM tracks t
WHERE t.album_id = $1 AND t.state = 'LIBRARY'
AND NOT (t.id = ANY($2::uuid[]))
ORDER BY t.title ASC
LIMIT 5`,
[album.album_id, ctx.recentExclusions]
);
for (const row of trackRes.rows as { id: string }[]) {
candidates.push({
trackId: row.id,
generatorId: 'deep-dive',
explanation: [{
subjectType: 'artist',
subjectId: album.artist_id,
predicate: 'credited_main_on',
objectType: 'track',
objectId: row.id,
fusedValue: 0.7,
}],
relevance: 0.7,
});
}
}
return candidates;
}
// ---------------------------------------------------------------------------
// 5. REVIVAL — Stale longterm affinity (last_reinforced > 90 days ago)
// ---------------------------------------------------------------------------
async function revivalGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
const staleRes = await db.pgClient.query(
`SELECT lb.entity_id AS artist_id, lb.value AS affinity
FROM listener_beliefs lb
WHERE lb.user_id = $1
AND lb.profile = 'longterm'
AND lb.entity_type = 'artist'
AND lb.dimension = 'affinity'
AND lb.value > 0.3
AND lb.last_reinforced_at < NOW() - INTERVAL '90 days'
ORDER BY lb.value DESC
LIMIT 20`,
[ctx.userId]
);
const staleArtists = staleRes.rows as { artist_id: string; affinity: number }[];
if (staleArtists.length === 0) return [];
const staleArtistIds = staleArtists.map(a => a.artist_id);
const trackRes = await db.pgClient.query(
`SELECT id FROM (
SELECT DISTINCT t.id
FROM tracks t
JOIN claim_fusion cf ON cf.subject_type = 'track' AND cf.subject_id = t.id
AND cf.predicate IN ('credited_main_on', 'featured_on')
AND cf.object_type = 'artist'
AND cf.object_id = ANY($1::uuid[])
AND (cf.user_id = $2 OR cf.user_id = $3)
WHERE t.state = 'LIBRARY'
AND NOT (t.id = ANY($4::uuid[]))
) sub
ORDER BY RANDOM()
LIMIT 20`,
[staleArtistIds, OBJECTIVE_USER, ctx.userId, ctx.recentExclusions]
);
return (trackRes.rows as { id: string }[]).map(row => ({
trackId: row.id,
generatorId: 'revival',
explanation: [{
subjectType: 'artist',
subjectId: staleArtistIds[0],
predicate: 'credited_main_on',
objectType: 'track',
objectId: row.id,
fusedValue: 0.6,
}],
relevance: 0.6,
}));
}
// ---------------------------------------------------------------------------
// 6. NOVELTY — Recently released tracks by graph-adjacent artists
// ---------------------------------------------------------------------------
async function noveltyGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
const trustedIds = ctx.beliefs
.filter(b => b.entity_type === 'artist' && b.value > 0.3)
.map(b => b.entity_id);
if (trustedIds.length === 0) return [];
const res = await db.pgClient.query(
`SELECT id FROM (
SELECT DISTINCT t.id, t.release_date
FROM tracks t
JOIN claim_fusion cf_edge ON cf_edge.subject_type = 'artist'
AND cf_edge.subject_id = ANY($2::uuid[])
AND cf_edge.predicate IN ('same_scene_as', 'same_label_as', 'produced')
AND cf_edge.object_type = 'artist'
JOIN claim_fusion cf_track ON cf_track.subject_type = 'track'
AND cf_track.subject_id = t.id
AND cf_track.predicate IN ('credited_main_on', 'featured_on')
AND cf_track.object_type = 'artist'
AND cf_track.object_id = cf_edge.object_id
WHERE t.release_date IS NOT NULL
AND t.release_date >= NOW() - INTERVAL '60 days'
AND t.state = 'LIBRARY'
AND NOT (t.id = ANY($1::uuid[]))
) sub
ORDER BY release_date DESC
LIMIT 20`,
[ctx.recentExclusions.length > 0 ? ctx.recentExclusions : ['00000000-0000-0000-0000-000000000000'], trustedIds]
);
return res.rows.map((row: { id: string }) => ({
trackId: row.id,
generatorId: 'novelty',
relevance: 0.5,
explanation: [{
subjectType: 'track', subjectId: row.id,
predicate: 'release_date',
objectType: 'date', objectId: 'recent',
fusedValue: 0.5,
}],
}));
}
// ---------------------------------------------------------------------------
// 7. EXPERIMENTAL — Genres with high network distance from favourites
// ---------------------------------------------------------------------------
async function experimentalGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
const favArtistIds = ctx.beliefs
.filter(b => b.entity_type === 'artist' && b.value > 0.4)
.map(b => b.entity_id);
if (favArtistIds.length === 0) return [];
const favGenreIds = ctx.beliefs
.filter(b => b.entity_type === 'genre' && b.value > 0.2)
.map(b => b.entity_id);
const result = await db.pgClient.query(
`WITH unfamiliar_genres AS (
SELECT g.id, g.name,
(SELECT COUNT(*) FROM track_genre tg2 WHERE tg2.genre_id = g.id) AS track_count
FROM genre g
WHERE NOT (g.id = ANY($1::uuid[]))
AND EXISTS (SELECT 1 FROM track_genre tg WHERE tg.genre_id = g.id)
ORDER BY RANDOM()
LIMIT 3
),
candidate_tracks AS (
SELECT DISTINCT t.id, tg.genre_id
FROM tracks t
JOIN track_genre tg ON tg.track_id = t.id
JOIN unfamiliar_genres ug ON ug.id = tg.genre_id
WHERE t.state = 'LIBRARY'
AND NOT (t.id = ANY($2::uuid[]))
LIMIT 30
)
SELECT ct.id, ct.genre_id
FROM candidate_tracks ct
ORDER BY RANDOM()
LIMIT 6`,
[
favGenreIds.length > 0 ? favGenreIds : ['00000000-0000-0000-0000-000000000000'],
ctx.recentExclusions.length > 0 ? ctx.recentExclusions : ['00000000-0000-0000-0000-000000000000'],
]
);
return result.rows.map((row: { id: string; genre_id: string }) => ({
trackId: row.id,
generatorId: 'experimental',
relevance: 0.2,
explanation: [{
subjectType: 'genre', subjectId: row.genre_id,
predicate: 'belongs_to_genre',
objectType: 'track', objectId: row.id,
fusedValue: 0.2,
}],
}));
}
// ---------------------------------------------------------------------------
// 8. CONTEXTUAL — Contextual profile beliefs
// ---------------------------------------------------------------------------
async function contextualGenerator(db: DbService, ctx: GeneratorContext): Promise<Candidate[]> {
if (!ctx.state.context) return [];
const contextualBeliefs = await db.getListenerBeliefs({
userId: ctx.userId,
profile: 'contextual',
limit: 30,
orderBy: 'value',
order: 'DESC',
});
const targetArtistIds = contextualBeliefs
.filter(b => b.entity_type === 'artist' && b.value > 0.2)
.map(b => b.entity_id);
if (targetArtistIds.length === 0) return [];
const trackRes = await db.pgClient.query(
`SELECT id FROM (
SELECT DISTINCT t.id
FROM tracks t
JOIN claim_fusion cf ON cf.subject_type = 'track' AND cf.subject_id = t.id
AND cf.predicate IN ('credited_main_on', 'featured_on')
AND cf.object_type = 'artist'
AND cf.object_id = ANY($1::uuid[])
AND (cf.user_id = $2 OR cf.user_id = $3)
WHERE t.state = 'LIBRARY'
AND NOT (t.id = ANY($4::uuid[]))
) sub
ORDER BY RANDOM()
LIMIT 15`,
[targetArtistIds, OBJECTIVE_USER, ctx.userId, ctx.recentExclusions]
);
return (trackRes.rows as { id: string }[]).map(row => ({
trackId: row.id,
generatorId: 'contextual',
explanation: [{
subjectType: 'artist',
subjectId: targetArtistIds[0],
predicate: 'credited_main_on',
objectType: 'track',
objectId: row.id,
fusedValue: 0.5,
}],
relevance: 0.5,
}));
}
// ---------------------------------------------------------------------------
// All generators, ordered by priority (comfort first, experimental last)
// ---------------------------------------------------------------------------
export const ALL_GENERATORS: Generator[] = [
comfortGenerator,
adjacentGenerator,
deepDiveGenerator,
revivalGenerator,
discoveryGenerator,
noveltyGenerator,
contextualGenerator,
experimentalGenerator,
];
+179
View File
@@ -0,0 +1,179 @@
import { describe, it, expect, vi } from 'vitest';
import { ALL_GENERATORS, type GeneratorContext, type Candidate } from './generators.service.js';
import { DbService } from './db.service.js';
function makeMockDb(overrides: Record<string, any> = {}): DbService {
const mockQuery = vi.fn();
return {
pgClient: { query: mockQuery },
getListenerBeliefs: vi.fn().mockResolvedValue([]),
...overrides,
} as unknown as DbService;
}
function makeCtx(overrides: Partial<GeneratorContext> = {}): GeneratorContext {
return {
userId: '00000000-0000-0000-0000-000000000000',
seedTrackId: null,
seedArtistId: null,
beliefs: [],
recentExclusions: [],
toleranceMap: {},
state: { energy: 0.5, lastArtistIds: [], lastGenreIds: [], context: null, noveltyHunger: 0.3, sessionAgeMin: 10 },
...overrides,
};
}
// Index generators by name for easy test access
const generatorByName: Record<string, (typeof ALL_GENERATORS)[0]> = {
comfort: ALL_GENERATORS[0],
adjacent: ALL_GENERATORS[1],
deepDive: ALL_GENERATORS[2],
revival: ALL_GENERATORS[3],
discovery: ALL_GENERATORS[4],
novelty: ALL_GENERATORS[5],
contextual: ALL_GENERATORS[6],
experimental: ALL_GENERATORS[7],
};
describe('generators', () => {
describe('comfort', () => {
it('returns tracks for artists with affinity > 0.5', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValue({ rows: [{ id: 'track-1' }, { id: 'track-2' }] });
const ctx = makeCtx({
beliefs: [
{ entity_type: 'artist', entity_id: 'artist-1', value: 0.8, confidence: 0.9, profile: 'longterm', dimension: 'affinity' } as any,
],
});
const results = await generatorByName.comfort(db, ctx);
expect(results.length).toBeGreaterThanOrEqual(1);
expect(results[0]).toHaveProperty('trackId');
expect(results[0]).toHaveProperty('generatorId', 'comfort');
expect(results[0].explanation.length).toBeGreaterThanOrEqual(1);
});
it('returns empty when no high-affinity artists', async () => {
const db = makeMockDb();
const ctx = makeCtx({ beliefs: [ { entity_type: 'artist', entity_id: 'a1', value: 0.3, profile: 'longterm', dimension: 'affinity' } as any ] });
const results = await generatorByName.comfort(db, ctx);
expect(results).toHaveLength(0);
});
});
describe('adjacent', () => {
it('returns tracks from graph walks', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValueOnce({ rows: [{ reached_artist_id: 'artist-2' }] });
(db.pgClient.query as any).mockResolvedValueOnce({ rows: [{ id: 'track-3' }] });
const ctx = makeCtx({ seedTrackId: 'track-1', seedArtistId: 'artist-1' });
const results = await generatorByName.adjacent(db, ctx);
if (results.length > 0) {
expect(results[0].explanation.length).toBeGreaterThanOrEqual(1);
}
});
it('returns empty when no seed artist', async () => {
const results = await generatorByName.adjacent(makeMockDb(), makeCtx());
expect(results).toHaveLength(0);
});
});
describe('discovery', () => {
it('respects novelty_tolerance cap', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValue({ rows: [{ artist_id: 'a1' }, { artist_id: 'a2' }] });
(db.pgClient.query as any).mockResolvedValue({ rows: [{ id: 't1' }] });
const ctx = makeCtx({
beliefs: [ { entity_type: 'artist', entity_id: 'trusted-1', value: 0.5, confidence: 0.8, profile: 'longterm', dimension: 'affinity' } as any ],
toleranceMap: { novelty_tolerance: 0.1 },
});
const results = await generatorByName.discovery(db, ctx);
const maxCandidates = Math.max(2, Math.floor(20 * 0.1));
expect(results.length).toBeLessThanOrEqual(maxCandidates);
});
});
describe('deepDive', () => {
it('returns album-ordered tracks for obsession artists', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValueOnce({ rows: [{ album_id: 'alb-1', artist_id: 'a1', title: 'Album 1' }] });
(db.pgClient.query as any).mockResolvedValue({ rows: [{ id: 't1' }, { id: 't2' }] });
const ctx = makeCtx({
beliefs: [ { entity_type: 'artist', entity_id: 'a1', dimension: 'affinity', value: 0.6, profile: 'obsession' } as any ],
});
const results = await generatorByName.deepDive(db, ctx);
expect(results.length).toBeGreaterThanOrEqual(1);
expect(results[0].generatorId).toBe('deep-dive');
});
});
describe('revival', () => {
it('returns tracks for stale high-affinity artists', async () => {
const db = makeMockDb();
// First query: stale listener_beliefs
(db.pgClient.query as any).mockResolvedValueOnce({
rows: [{ artist_id: 'a1', affinity: 0.7 }]
});
// Second query: tracks by those artists
(db.pgClient.query as any).mockResolvedValue({
rows: [{ id: 't1' }]
});
const ctx = makeCtx({
beliefs: [ { entity_type: 'artist', entity_id: 'a1', dimension: 'affinity', value: 0.5, profile: 'forgotten' } as any ],
});
const results = await generatorByName.revival(db, ctx);
expect(results.length).toBeGreaterThanOrEqual(1);
expect(results[0].generatorId).toBe('revival');
});
});
describe('novelty', () => {
it('returns recent tracks', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValue({ rows: [{ id: 't1' }] });
const ctx = makeCtx({
beliefs: [ { entity_type: 'artist', entity_id: 'trusted-1', value: 0.5, profile: 'longterm', dimension: 'affinity' } as any ],
});
const results = await generatorByName.novelty(db, ctx);
if (results.length > 0) {
expect(results[0].generatorId).toBe('novelty');
}
});
});
describe('experimental', () => {
it('returns tracks from unfamiliar genres', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValue({ rows: [{ id: 't1', genre_id: 'g1' }] });
const ctx = makeCtx({
beliefs: [ { entity_type: 'artist', entity_id: 'a1', value: 0.6, profile: 'longterm', dimension: 'affinity' } as any ],
});
const results = await generatorByName.experimental(db, ctx);
expect(results).toBeDefined();
});
});
describe('contextual', () => {
it('returns tracks matching context when set', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValueOnce({ rows: [{ id: 't1' }] });
const ctx = makeCtx({
state: { energy: 0.5, lastArtistIds: [], lastGenreIds: [], context: 'coding', noveltyHunger: 0.3, sessionAgeMin: 10 },
});
const results = await generatorByName.contextual(db, ctx);
expect(results).toBeDefined();
});
it('returns empty when no context set', async () => {
const results = await generatorByName.contextual(makeMockDb(), makeCtx());
expect(results).toHaveLength(0);
});
});
});
describe('ALL_GENERATORS', () => {
it('contains 8 generators', () => {
expect(ALL_GENERATORS).toHaveLength(8);
ALL_GENERATORS.forEach(g => expect(typeof g).toBe('function'));
});
});
@@ -0,0 +1,105 @@
import { DbService } from './db.service.js';
export class ImageEnrichmentService {
private sourcePriority: Record<string, number> = {
cover_art_archive: 1,
theaudiodb: 2,
fanart: 3,
deezer: 4,
discogs: 5,
lastfm: 6,
};
constructor(private db: DbService) {}
// ---------------------------------------------------------------
// Phase 4 — Fetch image candidates from all sources
// ---------------------------------------------------------------
async fetchImagesForArtist(artistId: string): Promise<number> {
const artistRes = await this.db.pgClient.query<{ id: string; name: string; mbid: string | null; discogs_id: string | null }>(
`SELECT id, name, mbid, discogs_id FROM artists WHERE id = $1`,
[artistId]
);
const artist = artistRes.rows[0];
if (!artist) return 0;
const sources: { key: string; condition: string }[] = [
{ key: 'cover_art_archive', condition: artist.mbid ? 'mbid present' : 'no mbid' },
{ key: 'deezer', condition: 'always' },
{ key: 'discogs', condition: artist.discogs_id ? 'discogs_id present' : 'no discogs_id' },
{ key: 'lastfm', condition: 'always' },
];
let count = 0;
for (const src of sources) {
const res = await this.db.pgClient.query(
`INSERT INTO image_candidates (entity_type, entity_id, source, url, width)
VALUES ('artist', $1, $2, NULL, NULL)
ON CONFLICT (entity_type, entity_id, source) DO NOTHING
RETURNING 1 AS ins`,
[artistId, src.key]
);
if (res.rows.length > 0) count++;
}
return count;
}
async fetchImagesForAlbum(albumId: string): Promise<number> {
const sources = ['cover_art_archive', 'itunes', 'deezer', 'discogs'];
let count = 0;
for (const source of sources) {
const res = await this.db.pgClient.query(
`INSERT INTO image_candidates (entity_type, entity_id, source, url, width)
VALUES ('album', $1, $2, NULL, NULL)
ON CONFLICT (entity_type, entity_id, source) DO NOTHING
RETURNING 1 AS ins`,
[albumId, source]
);
if (res.rows.length > 0) count++;
}
return count;
}
async selectBestImage(entityType: string, entityId: string): Promise<string | null> {
const candidates = await this.db.pgClient.query<{ url: string; source: string; verified: boolean }>(
`SELECT url, source, verified
FROM image_candidates
WHERE entity_type = $1 AND entity_id = $2 AND url IS NOT NULL
ORDER BY
CASE source
WHEN 'cover_art_archive' THEN 1
WHEN 'theaudiodb' THEN 2
WHEN 'fanart' THEN 3
WHEN 'deezer' THEN 4
WHEN 'discogs' THEN 5
WHEN 'lastfm' THEN 6
ELSE 99
END,
verified DESC,
width DESC NULLS LAST
LIMIT 1`,
[entityType, entityId]
);
const best = candidates.rows[0];
if (!best) return null;
if (entityType === 'artist') {
await this.db.pgClient.query(
`UPDATE artists SET image_path = $1 WHERE id = $2`,
[best.url, entityId]
);
} else if (entityType === 'album') {
await this.db.pgClient.query(
`UPDATE albums SET artwork_id = $1 WHERE id = $2`,
[best.url, entityId]
);
}
return best.url;
}
}
+133
View File
@@ -0,0 +1,133 @@
import { Queue, Job } from 'bullmq';
import { MetadataRefreshJob, AudioAnalysisJob, CleanupJob, LibraryScanJob, ReindexTracksJob, ReprocessArtistsJob } from '../types/job.types.js';
export interface JobServiceConfig {
redisUrl: string;
}
export const QUEUE_NAME = 'muzick-queue';
export interface QueueStats {
waiting: number;
active: number;
completed: number;
failed: number;
delayed: number;
paused: number;
}
export interface JobHistoryEntry {
id: string;
name: string;
data: Record<string, unknown>;
timestamp: number;
finishedOn?: number;
failedReason?: string;
returnvalue?: unknown;
}
export class JobService {
private queue: Queue;
constructor(config: JobServiceConfig) {
this.queue = new Queue(QUEUE_NAME, {
connection: {
url: config.redisUrl,
},
});
}
async enqueueMetadataRefresh(trackId: string, type: 'full' | 'partial') {
const payload: MetadataRefreshJob = { trackId, refreshType: type };
await this.queue.add('metadata_refresh', payload);
}
async enqueueAudioAnalysis(trackId: string, features: string[]) {
const payload: AudioAnalysisJob = { trackId, features };
await this.queue.add('audio_analysis', payload);
}
async enqueueCleanup(reason: 'expired' | 'manual', targetFiles: string[]) {
const payload: CleanupJob = { reason, targetFiles };
await this.queue.add('cleanup', payload);
}
async enqueueLibraryScan(directory: string) {
const payload: LibraryScanJob = { directory };
await this.queue.add('scan_library', payload);
}
async enqueueReindexTracks() {
const payload: ReindexTracksJob = {};
await this.queue.add('reindex_tracks', payload);
}
async enqueueReprocessArtists() {
const payload: ReprocessArtistsJob = { batchSize: 100, offset: 0 };
await this.queue.add('reprocess_artists', payload);
}
/**
* Enqueue metadata_refresh jobs for a batch of track IDs. Used by the
* /admin/reenrich-tracks endpoint to re-canonicalize metadata (artist names,
* album titles, MBIDs, cover art) without re-scanning files from disk.
*
* Each job is deduped by `jobId: meta-<trackId>` so re-running the endpoint
* doesn't stack duplicate jobs. Old completed/failed jobs with the same ID
* are removed first so re-enrichment actually works (BullMQ otherwise treats
* existing jobIds as duplicates and silently skips them).
*/
async enqueueMetadataRefreshBatch(trackIds: string[]): Promise<number> {
let enqueued = 0;
for (const trackId of trackIds) {
const jobId = `meta-${trackId}`;
await this.queue.remove(jobId).catch(() => {});
const payload: MetadataRefreshJob = { trackId, refreshType: 'full' };
await this.queue.add('metadata_refresh', payload, {
jobId,
removeOnComplete: { age: 86400, count: 10000 },
removeOnFail: { age: 86400, count: 10000 },
});
enqueued++;
}
return enqueued;
}
async getQueueStats(): Promise<QueueStats> {
const [waiting, active, completed, failed, delayed] = await Promise.all([
this.queue.getWaitingCount(),
this.queue.getActiveCount(),
this.queue.getCompletedCount(),
this.queue.getFailedCount(),
this.queue.getDelayedCount(),
]);
return { waiting, active, completed, failed, delayed, paused: 0 };
}
async getJobHistory(limit = 100): Promise<JobHistoryEntry[]> {
// Get jobs from completed and failed queues (most recent first)
const [completedJobs, failedJobs] = await Promise.all([
this.queue.getJobs(['completed'], 0, limit),
this.queue.getJobs(['failed'], 0, limit),
]);
const allJobs = [...completedJobs, ...failedJobs].map(job => ({
id: job.id as string,
name: job.name,
data: job.data as Record<string, unknown>,
timestamp: job.timestamp,
finishedOn: job.finishedOn,
failedReason: job.failedReason,
returnvalue: job.returnvalue,
}));
// Sort by timestamp descending (most recent first)
allJobs.sort((a, b) => b.timestamp - a.timestamp);
return allJobs.slice(0, limit);
}
async close() {
await this.queue.close();
}
}
+94
View File
@@ -0,0 +1,94 @@
import { Client } from 'typesense';
export interface SearchServiceConfig {
host: string;
port: number;
protocol: 'http' | 'https';
apiKey: string;
}
export class SearchService {
private client: Client;
private ready = false;
constructor(config: SearchServiceConfig) {
this.client = new Client({
nodes: [{
host: config.host,
port: config.port,
protocol: config.protocol,
}],
apiKey: config.apiKey,
});
}
get isReady(): boolean {
return this.ready;
}
async search(collection: string, query: string, options: any = {}) {
const searchParameters = {
'q': query,
'query_by': options.query_by || 'title,artist',
...options,
};
return await this.client.collections(collection).documents().search(searchParameters);
}
/**
* Create the 'tracks' collection schema if it does not already exist.
* Gracefully handles Typesense not being ready yet (503) — the collection
* can be created later by calling ensureCollection again or via the admin
* reindex endpoint. Search falls back to Postgres ILIKE when Typesense is
* unavailable, so a missing collection is never a hard failure.
*/
async ensureCollection(): Promise<void> {
const collectionSchema = {
name: 'tracks',
fields: [
{ name: 'id', type: 'string' as const },
{ name: 'title', type: 'string' as const },
{ name: 'artist', type: 'string' as const },
{ name: 'album', type: 'string' as const },
{ name: 'duration', type: 'int32' as const },
{ name: 'play_count', type: 'int32' as const },
{ name: 'genre', type: 'string[]' as const, facet: true },
{ name: 'source_type', type: 'string' as const },
],
};
// First check if the collection already exists (retrieve succeeds).
try {
await this.client.collections(collectionSchema.name).retrieve();
this.ready = true;
return;
} catch (checkErr: any) {
// 404 means collection doesn't exist — proceed to create.
// 503 means Typesense isn't ready yet — skip creation, search falls back.
if (checkErr?.httpStatus === 503) {
console.warn('[SearchService] Typesense not ready yet (503). Search will use Postgres ILIKE fallback.');
return;
}
}
// Collection doesn't exist — create it.
try {
await this.client.collections().create(collectionSchema);
this.ready = true;
} catch (createErr: any) {
// 409 / "already exists" — another process created it between our check and create.
if (createErr?.message?.includes('already exists')) {
this.ready = true;
return;
}
// 503 — Typesense not ready yet, not a hard failure.
if (createErr?.httpStatus === 503) {
console.warn('[SearchService] Typesense not ready yet (503). Search will use Postgres ILIKE fallback.');
return;
}
// Unexpected error — log and continue; search falls back to Postgres.
console.warn('[SearchService] Failed to create tracks collection:', createErr?.message ?? createErr);
}
}
}
@@ -0,0 +1,928 @@
import { DbService, ListenerBelief } from './db.service.js';
import { Candidate, GeneratorContext, Generator, ALL_GENERATORS } from './generators.service.js';
export interface FatigueState {
artist: Map<string, number>;
genre: Map<string, number>;
language: Map<string, number>;
track: Map<string, number>;
vocal: number;
}
export interface RecentPlay {
trackId: string;
artistId: string | null;
genreId: string | null;
bpm: number | null;
energy: number | null;
language: string | null;
vocal: boolean;
decade: number | null;
valence: number | null;
}
export interface DiversityBudget {
dimension: string;
budgetShare: number;
horizonMin: number;
spent: number;
}
const W_ENJOY = 1.0;
const W_FATIGUE = 0.4;
const W_DIVERSITY = 0.3;
const W_ENTROPY = 0.2;
const W_REPETITION = 0.5;
export class SessionDirector {
constructor(private db: DbService) {}
// ---------------------------------------------------------------
// D.1 — Build listener state vector
// ---------------------------------------------------------------
async buildState(userId: string, sessionId?: string): Promise<GeneratorContext['state']> {
let savedState: GeneratorContext['state'] | null = null;
if (sessionId) {
const res = await this.db.pgClient.query(
'SELECT * FROM session_state WHERE session_id = $1 AND user_id = $2',
[sessionId, userId]
);
if (res.rows[0]) {
const row = res.rows[0] as { state_vector: Record<string, unknown>; context: string | null; started_at: Date };
savedState = {
energy: (row.state_vector?.energy as number) ?? 0.5,
lastArtistIds: (row.state_vector?.lastArtistIds as string[]) ?? [],
lastGenreIds: (row.state_vector?.lastGenreIds as string[]) ?? [],
context: row.context,
noveltyHunger: (row.state_vector?.noveltyHunger as number) ?? 0.3,
sessionAgeMin: row.started_at
? (Date.now() - new Date(row.started_at).getTime()) / 60000
: 0,
};
}
}
if (!savedState) {
const latest = await this.db.getLatestSessionState(userId);
if (latest) {
savedState = {
energy: (latest.state_vector?.energy as number) ?? 0.5,
lastArtistIds: (latest.state_vector?.lastArtistIds as string[]) ?? [],
lastGenreIds: (latest.state_vector?.lastGenreIds as string[]) ?? [],
context: latest.context,
noveltyHunger: (latest.state_vector?.noveltyHunger as number) ?? 0.3,
sessionAgeMin: latest.started_at
? (Date.now() - new Date(latest.started_at).getTime()) / 60000
: 0,
};
}
}
// Compute fresh energy from last 5 completed plays
const energyRes = await this.db.pgClient.query(
`SELECT COALESCE(AVG(taf.energy), 0.5) AS energy
FROM (
SELECT ph.track_id
FROM play_history ph
WHERE ph.user_id = $1 AND ph.completed = true
ORDER BY ph.played_at DESC
LIMIT 5
) recent
JOIN track_audio_features taf ON taf.track_id = recent.track_id
WHERE taf.energy IS NOT NULL`,
[userId]
);
const energy = (energyRes.rows[0]?.energy as number) ?? 0.5;
// Read novelty hunger from discovery profile
const noveltyRes = await this.db.pgClient.query(
`SELECT value FROM listener_beliefs
WHERE user_id = $1 AND profile = 'discovery' AND dimension = 'novelty_tolerance'
LIMIT 1`,
[userId]
);
const noveltyHunger = (noveltyRes.rows[0]?.value as number) ?? 0.3;
// Last distinct artist IDs from recent completed plays.
// Use a subquery to order first, then DISTINCT — avoids PG's rule that
// DISTINCT + ORDER BY expressions must appear in the select list.
const lastArtistsRes = await this.db.pgClient.query(
`SELECT DISTINCT artist_id FROM (
SELECT ta.artist_id
FROM play_history ph
JOIN track_artists_v2 ta ON ta.track_id = ph.track_id AND ta.role = 'main'
WHERE ph.user_id = $1 AND ph.completed = true
ORDER BY ph.played_at DESC
LIMIT 40
) recent
LIMIT 10`,
[userId]
);
const lastArtistIds = lastArtistsRes.rows.map((r: { artist_id: string }) => r.artist_id);
// Last distinct genre IDs
const lastGenresRes = await this.db.pgClient.query(
`SELECT DISTINCT genre_id FROM (
SELECT tg.genre_id
FROM play_history ph
JOIN track_genre tg ON tg.track_id = ph.track_id
WHERE ph.user_id = $1 AND ph.completed = true
ORDER BY ph.played_at DESC
LIMIT 40
) recent
LIMIT 10`,
[userId]
);
const lastGenreIds = lastGenresRes.rows.map((r: { genre_id: string }) => r.genre_id);
const age = savedState?.sessionAgeMin ?? 0;
return {
energy,
lastArtistIds,
lastGenreIds,
context: savedState?.context ?? null,
noveltyHunger,
sessionAgeMin: age,
};
}
// ---------------------------------------------------------------
// D.2 — Fatigue model
// ---------------------------------------------------------------
async computeFatigue(userId: string): Promise<FatigueState> {
// Track fatigue: last 7 days, decay half-life 30d (2592000 seconds)
const TRACK_DECAY_SEC = 30 * 24 * 3600;
const trackRes = await this.db.pgClient.query(
`SELECT ph.track_id,
LEAST(1.0, SUM(EXP(-EXTRACT(EPOCH FROM (NOW() - ph.played_at)) / $2::float8))) AS fatigue
FROM play_history ph
WHERE ph.user_id = $1 AND ph.played_at > NOW() - INTERVAL '7 days' AND ph.completed = true
GROUP BY ph.track_id`,
[userId, TRACK_DECAY_SEC]
);
const track = new Map<string, number>();
for (const row of trackRes.rows as { track_id: string; fatigue: number }[]) {
track.set(row.track_id, row.fatigue);
}
// Artist fatigue: last 24h, decay half-life 8h (28800 seconds)
const ARTIST_DECAY_SEC = 8 * 3600;
const artistRes = await this.db.pgClient.query(
`SELECT ta.artist_id,
LEAST(1.0, SUM(EXP(-EXTRACT(EPOCH FROM (NOW() - ph.played_at)) / $2::float8))) AS fatigue
FROM play_history ph
JOIN track_artists_v2 ta ON ta.track_id = ph.track_id AND ta.role = 'main'
WHERE ph.user_id = $1 AND ph.played_at > NOW() - INTERVAL '24 hours' AND ph.completed = true
GROUP BY ta.artist_id`,
[userId, ARTIST_DECAY_SEC]
);
const artist = new Map<string, number>();
for (const row of artistRes.rows as { artist_id: string; fatigue: number }[]) {
artist.set(row.artist_id, row.fatigue);
}
// Genre fatigue: last 24h, decay half-life 8h
const genreRes = await this.db.pgClient.query(
`SELECT tg.genre_id,
LEAST(1.0, SUM(EXP(-EXTRACT(EPOCH FROM (NOW() - ph.played_at)) / $2::float8))) AS fatigue
FROM play_history ph
JOIN track_genre tg ON tg.track_id = ph.track_id
WHERE ph.user_id = $1 AND ph.played_at > NOW() - INTERVAL '24 hours' AND ph.completed = true
GROUP BY tg.genre_id`,
[userId, ARTIST_DECAY_SEC]
);
const genre = new Map<string, number>();
for (const row of genreRes.rows as { genre_id: string; fatigue: number }[]) {
genre.set(row.genre_id, row.fatigue);
}
// Language fatigue: last 2h, decay half-life 1h (3600 seconds)
const LANG_DECAY_SEC = 3600;
const langRes = await this.db.pgClient.query(
`SELECT tl.language,
LEAST(1.0, SUM(EXP(-EXTRACT(EPOCH FROM (NOW() - ph.played_at)) / $2::float8))) AS fatigue
FROM play_history ph
JOIN track_lyrics tl ON tl.track_id = ph.track_id
WHERE ph.user_id = $1 AND ph.played_at > NOW() - INTERVAL '2 hours' AND ph.completed = true
AND tl.language IS NOT NULL
GROUP BY tl.language`,
[userId, LANG_DECAY_SEC]
);
const language = new Map<string, number>();
for (const row of langRes.rows as { language: string; fatigue: number }[]) {
language.set(row.language, row.fatigue);
}
// Vocal fatigue: fraction of last 2h plays that are vocal (instrumentalness < 0.5)
const vocalRes = await this.db.pgClient.query(
`SELECT CASE WHEN COUNT(*) = 0 THEN 0.5
ELSE COUNT(*) FILTER (WHERE COALESCE(taf.instrumentalness, 0) < 0.5)::float8 / COUNT(*)::float8
END AS vocal_fatigue
FROM play_history ph
LEFT JOIN track_audio_features taf ON taf.track_id = ph.track_id
WHERE ph.user_id = $1 AND ph.played_at > NOW() - INTERVAL '2 hours' AND ph.completed = true`,
[userId]
);
const vocal = (vocalRes.rows[0]?.vocal_fatigue as number) ?? 0.5;
return { artist, genre, language, track, vocal };
}
// ---------------------------------------------------------------
// D.3 — Diversity budgets
// ---------------------------------------------------------------
async getBudgets(userId: string): Promise<DiversityBudget[]> {
const res = await this.db.pgClient.query(
'SELECT * FROM diversity_budgets WHERE user_id = $1 ORDER BY dimension',
[userId]
);
let rows: { dimension: string; budget_share: number; horizon_min: number }[];
if (res.rows.length === 0) {
await this.db.seedDefaultDiversityBudgets(userId);
const res2 = await this.db.pgClient.query(
'SELECT * FROM diversity_budgets WHERE user_id = $1 ORDER BY dimension',
[userId]
);
rows = res2.rows;
} else {
rows = res.rows;
}
const budgets: DiversityBudget[] = [];
for (const row of rows) {
const spent = await this.calcBudgetSpent(userId, row.dimension, row.horizon_min);
budgets.push({
dimension: row.dimension,
budgetShare: row.budget_share,
horizonMin: row.horizon_min,
spent,
});
}
return budgets;
}
private async calcBudgetSpent(userId: string, dimension: string, horizonMin: number): Promise<number> {
const interval = `${horizonMin} minutes`;
switch (dimension) {
case 'artist': {
const res = await this.db.pgClient.query(
`WITH sub AS (
SELECT COUNT(*) AS cnt
FROM play_history ph
JOIN track_artists_v2 ta ON ta.track_id = ph.track_id AND ta.role = 'main'
WHERE ph.user_id = $1 AND ph.played_at > NOW() - $2::interval AND ph.completed = true
GROUP BY ta.artist_id
)
SELECT COALESCE(MAX(cnt)::float8 / NULLIF((SELECT SUM(cnt) FROM sub), 0), 0) AS spent
FROM sub`,
[userId, interval]
);
return (res.rows[0]?.spent as number) ?? 0;
}
case 'genre': {
const res = await this.db.pgClient.query(
`WITH sub AS (
SELECT COUNT(*) AS cnt
FROM play_history ph
JOIN track_genre tg ON tg.track_id = ph.track_id
WHERE ph.user_id = $1 AND ph.played_at > NOW() - $2::interval AND ph.completed = true
GROUP BY tg.genre_id
)
SELECT COALESCE(MAX(cnt)::float8 / NULLIF((SELECT SUM(cnt) FROM sub), 0), 0) AS spent
FROM sub`,
[userId, interval]
);
return (res.rows[0]?.spent as number) ?? 0;
}
case 'language': {
const res = await this.db.pgClient.query(
`WITH sub AS (
SELECT tl.language, COUNT(*) AS cnt
FROM play_history ph
JOIN track_lyrics tl ON tl.track_id = ph.track_id
WHERE ph.user_id = $1 AND ph.played_at > NOW() - $2::interval AND ph.completed = true
AND tl.language IS NOT NULL
GROUP BY tl.language
)
SELECT COALESCE(MAX(cnt)::float8 / NULLIF((SELECT SUM(cnt) FROM sub), 0), 0) AS spent
FROM sub`,
[userId, interval]
);
return (res.rows[0]?.spent as number) ?? 0;
}
case 'instrumental': {
const res = await this.db.pgClient.query(
`SELECT COALESCE(
COUNT(*) FILTER (WHERE COALESCE(taf.instrumentalness, 0) > 0.5)::float8 / NULLIF(COUNT(*), 0),
0) AS spent
FROM play_history ph
LEFT JOIN track_audio_features taf ON taf.track_id = ph.track_id
WHERE ph.user_id = $1 AND ph.played_at > NOW() - $2::interval AND ph.completed = true`,
[userId, interval]
);
return (res.rows[0]?.spent as number) ?? 0;
}
case 'new_artist': {
const res = await this.db.pgClient.query(
`WITH recent_artists AS (
SELECT DISTINCT ta.artist_id
FROM play_history ph
JOIN track_artists_v2 ta ON ta.track_id = ph.track_id AND ta.role = 'main'
WHERE ph.user_id = $1 AND ph.played_at > NOW() - $2::interval AND ph.completed = true
)
SELECT COALESCE(
SUM(CASE WHEN NOT EXISTS (
SELECT 1 FROM play_history ph3
JOIN track_artists_v2 ta3 ON ta3.track_id = ph3.track_id AND ta3.role = 'main'
WHERE ph3.user_id = $1 AND ph3.played_at <= NOW() - $2::interval
AND ta3.artist_id = ra.artist_id
) THEN 1 ELSE 0 END)::float8 / NULLIF(COUNT(*), 0),
0) AS spent
FROM recent_artists ra`,
[userId, interval]
);
return (res.rows[0]?.spent as number) ?? 0;
}
case 'favorite': {
const res = await this.db.pgClient.query(
`SELECT COALESCE(
COUNT(*) FILTER (WHERE f.track_id IS NOT NULL)::float8 / NULLIF(COUNT(*), 0),
0) AS spent
FROM play_history ph
LEFT JOIN favorites f ON f.track_id = ph.track_id AND f.user_id = $1
WHERE ph.user_id = $1 AND ph.played_at > NOW() - $2::interval AND ph.completed = true`,
[userId, interval]
);
return (res.rows[0]?.spent as number) ?? 0;
}
default:
return 0;
}
}
// ---------------------------------------------------------------
// D.4 — Arc selection
// ---------------------------------------------------------------
pickArc(state: GeneratorContext['state']): string {
if (state.energy < 0.3) return 'late-night';
if (state.energy > 0.6 && state.noveltyHunger > 0.5) return 'discovery';
if (state.energy > 0.6) return 'energetic';
return 'comfort';
}
getArcSlots(arcType: string, count: number): { position: number; role: string }[] {
const pattern = this.getArcPattern(arcType);
const slots: { position: number; role: string }[] = [];
for (let i = 0; i < count; i++) {
slots.push({ position: i, role: pattern[i % pattern.length] });
}
return slots;
}
private getArcPattern(arcType: string): string[] {
switch (arcType) {
case 'comfort':
return ['known', 'known', 'known', 'known', 'adjacent', 'adjacent', 'adjacent', 'adjacent', 'favorite', 'favorite'];
case 'discovery':
return ['favorite', 'similar', 'new', 'favorite'];
case 'energetic':
return ['medium', 'medium', 'high', 'high', 'high', 'peak', 'cooldown', 'cooldown'];
case 'late-night':
return ['soft', 'soft', 'ambient', 'ambient', 'acoustic', 'slow'];
default:
return ['known', 'known', 'known', 'known', 'adjacent', 'adjacent', 'adjacent', 'adjacent', 'favorite', 'favorite'];
}
}
private roleToGeneratorIds(role: string): string[] {
switch (role) {
case 'known':
case 'medium':
case 'soft':
case 'acoustic':
case 'slow':
case 'cooldown':
return ['comfort'];
case 'adjacent':
case 'similar':
return ['adjacent'];
case 'favorite':
return ['deep-dive', 'comfort'];
case 'new':
case 'high':
return ['discovery'];
case 'peak':
return ['deep-dive', 'contextual'];
case 'ambient':
return ['contextual', 'comfort'];
default:
return ['comfort'];
}
}
// ---------------------------------------------------------------
// D.5 — Entropy, anti-loop
// ---------------------------------------------------------------
computeEntropy(candidates: Candidate[]): number {
if (candidates.length === 0) return 0;
const artistCounts = new Map<string, number>();
for (const c of candidates) {
const mainEdge = c.explanation.find(
e => e.subjectType === 'artist' || e.objectType === 'artist'
);
const key = mainEdge?.subjectId ?? mainEdge?.objectId ?? 'unknown';
artistCounts.set(key, (artistCounts.get(key) ?? 0) + 1);
}
const n = candidates.length;
let hhi = 0;
for (const count of artistCounts.values()) {
const share = count / n;
hhi += share * share;
}
return hhi;
}
async detectAntiLoop(
state: GeneratorContext['state'],
fatigue: FatigueState,
budgets: DiversityBudget[],
recentPlays: RecentPlay[]
): Promise<string | null> {
const n = recentPlays.length;
if (n < 3) return null;
// 1. ARTIST: single artist > 30% of recent plays
const artistCounts = new Map<string, number>();
for (const p of recentPlays) {
if (p.artistId) artistCounts.set(p.artistId, (artistCounts.get(p.artistId) ?? 0) + 1);
}
for (const count of artistCounts.values()) {
if (count / n > 0.3) return 'artist';
}
// 2. GENRE: single genre > 40% of recent plays
const genreCounts = new Map<string, number>();
for (const p of recentPlays) {
if (p.genreId) genreCounts.set(p.genreId, (genreCounts.get(p.genreId) ?? 0) + 1);
}
for (const count of genreCounts.values()) {
if (count / n > 0.4) return 'genre';
}
// 3. LANGUAGE: single language > 50% of recent plays
const langCounts = new Map<string, number>();
for (const p of recentPlays) {
if (p.language) langCounts.set(p.language, (langCounts.get(p.language) ?? 0) + 1);
}
for (const count of langCounts.values()) {
if (count / n > 0.5) return 'language';
}
// 4. ENERGY: >60% of plays in same energy quartile
const energies = recentPlays.filter(p => p.energy != null).map(p => p.energy!);
if (energies.length >= 3) {
const quartileCounts = [0, 0, 0, 0];
for (const e of energies) {
const q = Math.min(Math.floor(e / 0.25), 3);
quartileCounts[q]++;
}
if (Math.max(...quartileCounts) / energies.length > 0.6) return 'energy';
}
// 5. BPM: all plays within 20 BPM of each other
const bpms = recentPlays.filter(p => p.bpm != null).map(p => p.bpm!);
if (bpms.length >= 3) {
const bpmMin = Math.min(...bpms);
const bpmMax = Math.max(...bpms);
if (bpmMax - bpmMin <= 20) return 'bpm';
}
// 6. VOCAL: >80% all-vocal or all-instrumental
if (n >= 3) {
const vocalCount = recentPlays.filter(p => p.vocal).length;
const vocalRatio = vocalCount / n;
if (vocalRatio > 0.8 || vocalRatio < 0.2) return 'vocal';
}
// 7. DECADE: >50% from same decade
const decadeCounts = new Map<number, number>();
for (const p of recentPlays) {
if (p.decade != null) decadeCounts.set(p.decade, (decadeCounts.get(p.decade) ?? 0) + 1);
}
for (const count of decadeCounts.values()) {
if (count / n > 0.5) return 'decade';
}
// 8. PRODUCER: single producer > 3 tracks
const trackIds = recentPlays.map(p => p.trackId).filter(Boolean);
if (trackIds.length > 0) {
const prodRes = await this.db.pgClient.query(
`SELECT c.object_id
FROM claims c
WHERE c.predicate = 'produced'
AND c.subject_id = ANY($1::uuid[])
GROUP BY c.object_id
HAVING COUNT(DISTINCT c.subject_id) > 3`,
[trackIds]
);
if (prodRes.rows.length > 0) return 'producer';
}
// 9. LABEL: single label > 3 tracks
if (trackIds.length > 0) {
const labelRes = await this.db.pgClient.query(
`SELECT c.object_id
FROM claims c
WHERE c.predicate = 'same_label_as'
AND c.subject_id = ANY($1::uuid[])
GROUP BY c.object_id
HAVING COUNT(DISTINCT c.subject_id) > 3`,
[trackIds]
);
if (labelRes.rows.length > 0) return 'label';
}
// 10. MOOD: all plays same mood (valence > 0.5 = positive, <= 0.5 = negative)
const valences = recentPlays.filter(p => p.valence != null).map(p => p.valence!);
if (valences.length >= 3) {
const positiveCount = valences.filter(v => v > 0.5).length;
if (positiveCount === valences.length || positiveCount === 0) return 'mood';
}
return null;
}
// ---------------------------------------------------------------
// D.6 — Repetition rules
// ---------------------------------------------------------------
async checkRepetition(trackId: string, artistId: string, userId: string): Promise<boolean> {
const rulesRes = await this.db.pgClient.query(
'SELECT dimension, min_distance FROM repetition_rules WHERE user_id = $1',
[userId]
);
const ruleMap = new Map<string, number>();
for (const row of rulesRes.rows as { dimension: string; min_distance: number }[]) {
ruleMap.set(row.dimension, row.min_distance);
}
const trackMin = ruleMap.get('track') ?? 120;
if (trackMin > 0) {
const res = await this.db.pgClient.query(
`SELECT 1 FROM play_history
WHERE user_id = $1 AND track_id = $2 AND completed = true
AND played_at > NOW() - ($3 || ' minutes')::interval
LIMIT 1`,
[userId, trackId, String(trackMin)]
);
if (res.rows.length > 0) return true;
}
const artistMin = ruleMap.get('artist') ?? 20;
if (artistId && artistMin > 0) {
const res = await this.db.pgClient.query(
`SELECT 1 FROM play_history ph
JOIN track_artists_v2 ta ON ta.track_id = ph.track_id AND ta.artist_id = $2 AND ta.role = 'main'
WHERE ph.user_id = $1 AND ph.completed = true
AND ph.played_at > NOW() - ($3 || ' minutes')::interval
LIMIT 1`,
[userId, artistId, String(artistMin)]
);
if (res.rows.length > 0) return true;
}
return false;
}
// ---------------------------------------------------------------
// D.8 — Multi-objective ranking
// ---------------------------------------------------------------
async rankCandidates(
candidates: Candidate[],
fatigue: FatigueState,
budgets: DiversityBudget[],
state: GeneratorContext['state'],
repetitionCheck: (trackId: string, artistId: string) => Promise<boolean>
): Promise<Candidate[]> {
if (candidates.length === 0) return [];
const trackIds = [...new Set(candidates.map(c => c.trackId))];
const artistMap = new Map<string, string>();
if (trackIds.length > 0) {
const artRes = await this.db.pgClient.query(
`SELECT DISTINCT ON (ta.track_id) ta.track_id, ta.artist_id
FROM track_artists_v2 ta
WHERE ta.track_id = ANY($1::uuid[]) AND ta.role = 'main'`,
[trackIds]
);
for (const row of artRes.rows as { track_id: string; artist_id: string }[]) {
artistMap.set(row.track_id, row.artist_id);
}
}
const genreMap = new Map<string, string>();
if (trackIds.length > 0) {
const genreRes = await this.db.pgClient.query(
`SELECT DISTINCT ON (tg.track_id) tg.track_id, tg.genre_id
FROM track_genre tg
WHERE tg.track_id = ANY($1::uuid[])
ORDER BY tg.track_id, tg.weight DESC`,
[trackIds]
);
for (const row of genreRes.rows as { track_id: string; genre_id: string }[]) {
genreMap.set(row.track_id, row.genre_id);
}
}
const artistBudget = budgets.find(b => b.dimension === 'artist');
const currentEntropy = this.computeEntropy(candidates);
const targetEntropy = 0.55;
const scored: { candidate: Candidate; score: number }[] = [];
for (const c of candidates) {
const artistId = artistMap.get(c.trackId) ?? '';
const genreId = genreMap.get(c.trackId) ?? '';
const trackFatigue = fatigue.track.get(c.trackId) ?? 0;
const artistFatigue = fatigue.artist.get(artistId) ?? 0;
const genreFatigue = fatigue.genre.get(genreId) ?? 0;
const avgFatigue = (trackFatigue + artistFatigue + genreFatigue) / 3;
const artistSpendRatio = artistBudget ? artistBudget.spent : 0;
const diversityBonus = 1 - artistSpendRatio;
const entropyBonus = 1 - Math.abs(currentEntropy - targetEntropy);
const wouldRepeat = await repetitionCheck(c.trackId, artistId);
let score = W_ENJOY * c.relevance
- W_FATIGUE * avgFatigue
+ W_DIVERSITY * diversityBonus
+ W_ENTROPY * entropyBonus;
if (wouldRepeat) {
score *= 0.1;
}
scored.push({ candidate: c, score });
}
const entropyDrift = Math.abs(currentEntropy - targetEntropy);
if (entropyDrift > 0.2) {
const genCounts = new Map<string, number>();
for (const s of scored) {
genCounts.set(s.candidate.generatorId, (genCounts.get(s.candidate.generatorId) ?? 0) + 1);
}
const maxCount = Math.max(...genCounts.values(), 1);
for (const s of scored) {
const genCount = genCounts.get(s.candidate.generatorId) ?? 0;
s.score += (1 - genCount / maxCount) * 0.15;
}
}
scored.sort((a, b) => b.score - a.score);
return scored.map(s => s.candidate);
}
// ---------------------------------------------------------------
// D.9 — Plan + replan loop
// ---------------------------------------------------------------
async buildPlan(userId: string, sessionId: string, seedTrackId?: string): Promise<Candidate[]> {
const allBeliefs = await this.db.getListenerBeliefs({ userId, limit: 200 });
// Fetch recent completed plays for anti-loop detection
const recentPlaysRes = await this.db.pgClient.query(
`SELECT t.id AS track_id, ta.artist_id, tg.genre_id,
af.bpm, af.energy, af.valence, af.instrumentalness,
tl.language,
t.release_date
FROM play_history ph
JOIN tracks t ON t.id = ph.track_id
LEFT JOIN track_audio_features af ON af.track_id = t.id
LEFT JOIN track_artists_v2 ta ON ta.track_id = t.id AND ta.role = 'main'
LEFT JOIN track_genre tg ON tg.track_id = t.id AND tg.weight = (
SELECT MAX(weight) FROM track_genre WHERE track_id = t.id
)
LEFT JOIN track_lyrics tl ON tl.track_id = t.id
WHERE ph.user_id = $1 AND ph.completed = true
ORDER BY ph.played_at DESC
LIMIT 20`,
[userId]
);
const recentPlays: RecentPlay[] = recentPlaysRes.rows.map((r: any) => ({
trackId: r.track_id,
artistId: r.artist_id ?? null,
genreId: r.genre_id ?? null,
bpm: r.bpm ?? null,
energy: r.energy ?? null,
language: r.language ?? null,
vocal: (r.instrumentalness == null) ? false : r.instrumentalness < 0.5,
decade: r.release_date ? Math.floor(new Date(r.release_date).getFullYear() / 10) * 10 : null,
valence: r.valence ?? null,
}));
const state = await this.buildState(userId, sessionId);
const fatigue = await this.computeFatigue(userId);
const budgets = await this.getBudgets(userId);
const arcType = this.pickArc(state);
const planSize = 20;
const slots = this.getArcSlots(arcType, planSize);
let seedArtistId: string | null = null;
if (seedTrackId) {
seedArtistId = await this.resolveSeedArtistId(seedTrackId) ?? null;
}
const recentExclusions: string[] = [];
const toleranceMap: Record<string, number> = {};
const discoveryBeliefs = allBeliefs.filter(b => b.profile === 'discovery');
for (const b of discoveryBeliefs) {
if (b.dimension) toleranceMap[b.dimension] = b.value;
}
const ctx: GeneratorContext = {
userId,
seedTrackId: seedTrackId ?? null,
seedArtistId,
beliefs: allBeliefs,
recentExclusions,
toleranceMap,
state,
};
const allCandidates: Candidate[] = [];
for (const gen of ALL_GENERATORS) {
const result = await gen(this.db, ctx);
allCandidates.push(...result);
}
if (allCandidates.length === 0) {
return [];
}
const repetitionCheckFn = (tid: string, aid: string) =>
this.checkRepetition(tid, aid, userId);
const ranked = await this.rankCandidates(
allCandidates, fatigue, budgets, state, repetitionCheckFn
);
const loopDim = await this.detectAntiLoop(state, fatigue, budgets, recentPlays);
let forcedExperimental = false;
if (loopDim && ranked.length > 0) {
const expCtx: GeneratorContext = {
...ctx,
recentExclusions: ctx.recentExclusions.slice(0, Math.min(ctx.recentExclusions.length, 50)),
};
const extraCandidates: Candidate[] = [];
for (const gen of ALL_GENERATORS) {
const result = await gen(this.db, expCtx);
extraCandidates.push(...result);
}
const expRanked = await this.rankCandidates(
extraCandidates, fatigue, budgets, state, repetitionCheckFn
);
const injected = expRanked.filter(
c => c.generatorId === 'experimental' || c.generatorId === 'discovery'
);
ranked.unshift(...injected);
forcedExperimental = true;
}
const seen = new Set<string>();
const deduped: Candidate[] = [];
for (const c of ranked) {
if (!seen.has(c.trackId)) {
seen.add(c.trackId);
deduped.push(c);
}
}
const plan: Candidate[] = [];
const usedTrackIds = new Set<string>();
if (!forcedExperimental) {
const unused = [...deduped];
for (const slot of slots) {
const prefGenIds = this.roleToGeneratorIds(slot.role);
let idx = unused.findIndex(
c => prefGenIds.includes(c.generatorId) && !usedTrackIds.has(c.trackId)
);
if (idx === -1) {
idx = unused.findIndex(c => !usedTrackIds.has(c.trackId));
}
if (idx === -1) break;
const chosen = unused[idx];
usedTrackIds.add(chosen.trackId);
plan.push(chosen);
unused.splice(idx, 1);
}
if (plan.length < planSize) {
for (const c of deduped) {
if (plan.length >= planSize) break;
if (!usedTrackIds.has(c.trackId)) {
usedTrackIds.add(c.trackId);
plan.push(c);
}
}
}
} else {
for (const c of deduped) {
if (plan.length >= planSize) break;
plan.push(c);
}
}
return plan.slice(0, planSize);
}
async replan(
userId: string,
sessionId: string,
currentPlan: Candidate[],
playedTrackIds: string[],
seedTrackId?: string
): Promise<Candidate[]> {
const remainingSlots = currentPlan.filter(
c => !playedTrackIds.includes(c.trackId)
);
if (remainingSlots.length >= 10 && currentPlan.length > 0) {
const fatigue = await this.computeFatigue(userId);
const budgets = await this.getBudgets(userId);
const state = await this.buildState(userId, sessionId);
// Fetch recent plays for anti-loop
const recentPlaysRes = await this.db.pgClient.query(
`SELECT t.id AS track_id, ta.artist_id, tg.genre_id,
af.bpm, af.energy, af.valence, af.instrumentalness,
tl.language,
t.release_date
FROM play_history ph
JOIN tracks t ON t.id = ph.track_id
LEFT JOIN track_audio_features af ON af.track_id = t.id
LEFT JOIN track_artists_v2 ta ON ta.track_id = t.id AND ta.role = 'main'
LEFT JOIN track_genre tg ON tg.track_id = t.id AND tg.weight = (
SELECT MAX(weight) FROM track_genre WHERE track_id = t.id
)
LEFT JOIN track_lyrics tl ON tl.track_id = t.id
WHERE ph.user_id = $1 AND ph.completed = true
ORDER BY ph.played_at DESC
LIMIT 20`,
[userId]
);
const recentPlays: RecentPlay[] = recentPlaysRes.rows.map((r: any) => ({
trackId: r.track_id,
artistId: r.artist_id ?? null,
genreId: r.genre_id ?? null,
bpm: r.bpm ?? null,
energy: r.energy ?? null,
language: r.language ?? null,
vocal: (r.instrumentalness == null) ? false : r.instrumentalness < 0.5,
decade: r.release_date ? Math.floor(new Date(r.release_date).getFullYear() / 10) * 10 : null,
valence: r.valence ?? null,
}));
const loopDim = await this.detectAntiLoop(state, fatigue, budgets, recentPlays);
if (loopDim) {
return this.buildPlan(userId, sessionId, seedTrackId);
}
const entropy = this.computeEntropy(currentPlan);
if (Math.abs(entropy - 0.55) > 0.2) {
return this.buildPlan(userId, sessionId, seedTrackId);
}
return remainingSlots;
}
return this.buildPlan(userId, sessionId, seedTrackId);
}
// ---------------------------------------------------------------
// Private helpers
// ---------------------------------------------------------------
private async resolveSeedArtistId(seedTrackId: string): Promise<string | undefined> {
const res = await this.db.pgClient.query(
`SELECT ta.artist_id
FROM track_artists_v2 ta
WHERE ta.track_id = $1 AND ta.role = 'main'
LIMIT 1`,
[seedTrackId]
);
if (res.rows[0]?.artist_id) return res.rows[0].artist_id as string;
const fallback = await this.db.pgClient.query(
`SELECT al.artist_id
FROM tracks t
JOIN albums al ON al.id = t.album_id
WHERE t.id = $1`,
[seedTrackId]
);
return fallback.rows[0]?.artist_id as string | undefined;
}
}
@@ -0,0 +1,117 @@
import { describe, it, expect, vi } from 'vitest';
import { SessionDirector } from './session-director.service.js';
import { DbService } from './db.service.js';
function makeMockDb(overrides: Record<string, any> = {}): DbService {
const mockQuery = vi.fn();
return {
pgClient: { query: mockQuery },
getListenerBeliefs: vi.fn().mockResolvedValue([]),
getLatestSessionState: vi.fn().mockResolvedValue(null),
seedDefaultDiversityBudgets: vi.fn().mockResolvedValue(undefined),
upsertDiversityBudget: vi.fn().mockResolvedValue(undefined),
...overrides,
} as unknown as DbService;
}
describe('SessionDirector', () => {
describe('pickArc', () => {
const director = new SessionDirector(makeMockDb());
it('returns late-night for low energy', () => {
const arc = director.pickArc({ energy: 0.2, noveltyHunger: 0.3, sessionAgeMin: 10 } as any);
expect(arc).toBe('late-night');
});
it('returns discovery for high energy + high novelty', () => {
const arc = director.pickArc({ energy: 0.7, noveltyHunger: 0.6, sessionAgeMin: 5 } as any);
expect(arc).toBe('discovery');
});
it('returns energetic for high energy + low novelty', () => {
const arc = director.pickArc({ energy: 0.7, noveltyHunger: 0.3, sessionAgeMin: 5 } as any);
expect(arc).toBe('energetic');
});
it('returns comfort for medium energy', () => {
const arc = director.pickArc({ energy: 0.5, noveltyHunger: 0.3, sessionAgeMin: 10 } as any);
expect(arc).toBe('comfort');
});
});
describe('getArcSlots', () => {
const director = new SessionDirector(makeMockDb());
it('returns correct slot count', () => {
expect(director.getArcSlots('comfort', 20).length).toBe(20);
expect(director.getArcSlots('discovery', 20).length).toBe(20);
expect(director.getArcSlots('energetic', 10).length).toBe(10);
expect(director.getArcSlots('late-night', 8).length).toBe(8);
});
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'];
slots.forEach(s => expect(validRoles).toContain(s.role));
});
});
describe('computeEntropy', () => {
const director = new SessionDirector(makeMockDb());
it('returns 0 for empty set', () => {
expect(director.computeEntropy([])).toBe(0);
});
it('returns 1 for all-same-artist', () => {
const candidates = [
{ trackId: 't1', generatorId: 'c', relevance: 1, explanation: [{ subjectType: 'artist', subjectId: 'a1', predicate: 'credited_main_on', objectType: 'track', objectId: 't1', fusedValue: 1 }] },
{ trackId: 't2', generatorId: 'c', relevance: 1, explanation: [{ subjectType: 'artist', subjectId: 'a1', predicate: 'credited_main_on', objectType: 'track', objectId: 't2', fusedValue: 1 }] },
] as any;
expect(director.computeEntropy(candidates)).toBe(1);
});
it('returns ~0.5 for two-artist split', () => {
const candidates = [
{ trackId: 't1', generatorId: 'c', relevance: 1, explanation: [{ subjectType: 'artist', subjectId: 'a1', predicate: 'credited_main_on', objectType: 'track', objectId: 't1', fusedValue: 1 }] },
{ trackId: 't2', generatorId: 'c', relevance: 1, explanation: [{ subjectType: 'artist', subjectId: 'a2', predicate: 'credited_main_on', objectType: 'track', objectId: 't2', fusedValue: 1 }] },
] as any;
const hhi = director.computeEntropy(candidates);
expect(hhi).toBeCloseTo(0.5);
});
});
describe('rankCandidates', () => {
it('sorts candidates by score descending', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValue({ rows: [] });
const director = new SessionDirector(db);
const candidates = [
{ trackId: 't1', generatorId: 'a', relevance: 0.9, explanation: [{ subjectType: 'artist', subjectId: 'a1', predicate: 'credited_main_on', objectType: 'track', objectId: 't1', fusedValue: 1 }] },
{ trackId: 't2', generatorId: 'b', relevance: 0.3, explanation: [{ subjectType: 'artist', subjectId: 'a2', predicate: 'credited_main_on', objectType: 'track', objectId: 't2', fusedValue: 1 }] },
];
const fatigue = { artist: new Map(), genre: new Map(), track: new Map(), language: new Map(), vocal: 0 };
const budgets = [{ dimension: 'artist', budgetShare: 0.2, horizonMin: 30, spent: 0 }];
const state = { energy: 0.5, noveltyHunger: 0.3, sessionAgeMin: 10, lastArtistIds: [], lastGenreIds: [], context: null };
const ranked = await director.rankCandidates(candidates, fatigue, budgets, state, async () => false);
expect(ranked[0].relevance).toBeGreaterThanOrEqual(ranked[ranked.length - 1].relevance);
});
});
describe('buildState', () => {
it('returns state with default values when no prior session', async () => {
const db = makeMockDb();
(db.pgClient.query as any).mockResolvedValue({ rows: [] });
(db.getListenerBeliefs as any).mockResolvedValue([]);
(db.getLatestSessionState as any).mockResolvedValue(null);
const director = new SessionDirector(db);
const state = await director.buildState('user-1');
expect(state).toHaveProperty('energy');
expect(state).toHaveProperty('noveltyHunger');
expect(state).toHaveProperty('sessionAgeMin');
expect(typeof state.energy).toBe('number');
});
});
});