apps/web/src/lib/flowSelector.server.ts

1import { createGateway } from '@ai-sdk/gateway'
2import { propagateAttributes, startObservation } from '@langfuse/tracing'
3import { confirmedHandoff } from '@local/ai/components/workflows'
4import { experimental_evaluate as evaluate } from 'ai'
5import {
6	changeQuestion,
7	entryQuestion,
8	type Flow,
9	routineContinuation,
10	selectDestination,
11} from './flowSelector'
12
13export async function selectFlow(
14	input: {
15		message: string
16		history: { role: 'user' | 'assistant'; text: string }[]
17		page: string
18		active: Flow | null
19		chatId: string
20	},
21	signal?: AbortSignal,
22) {
23	if (
24		input.active &&
25		(routineContinuation(input.message) ||
26			confirmedHandoff(input.history, input.message))
27	)
28		return selectDestination('continue', 1, input.active)
29	const history =
30		input.history.at(-1)?.role === 'user' &&
31		input.history.at(-1)?.text === input.message
32			? input.history.slice(0, -1)
33			: input.history
34	const span = propagateAttributes(
35		{
36			sessionId: input.chatId,
37			traceName: 'jev-flow-selector',
38			tags: ['flow-selector'],
39		},
40		() =>
41			startObservation(
42				'jev-flow-selector',
43				{
44					model: 'typesafe-ai/jev',
45					input: { active: input.active, historyTurns: Math.min(history.length, 6) },
46				},
47				{ asType: 'generation' },
48			),
49	)
50	try {
51		const result = await evaluate({
52			model: createGateway({
53				apiKey: process.env.AI_GATEWAY_API_KEY || process.env.AI_GATEWAY,
54			}).evaluationModel('typesafe-ai/jev'),
55			state: {
56				message: input.message,
57				page: input.page,
58				activeWorkflow: input.active,
59				history: history.slice(-6).map((t) => ({
60					...t,
61					text:
62						t.text.length > 4000
63							? `${t.text.slice(0, 1000)}\n[Middle omitted]\n${t.text.slice(-3000)}`
64							: t.text,
65				})),
66			},
67			questions: { destination: input.active ? changeQuestion : entryQuestion },
68			maxRetries: 0,
69			abortSignal: signal
70				? AbortSignal.any([signal, AbortSignal.timeout(1500)])
71				: AbortSignal.timeout(1500),
72			providerOptions: { gateway: { zeroDataRetention: true } },
73		})
74		const a = result.answers.destination
75		const decision = selectDestination(
76			a.choice,
77			Object.entries(a.probabilities ?? {}).find(
78				([choice]) => choice === a.choice,
79			)?.[1] ?? 0,
80			input.active,
81		)
82		span.update({
83			output: { decision, probabilities: a.probabilities },
84			usageDetails: {
85				input: result.usage.inputTokens ?? 0,
86				output: result.usage.outputTokens ?? 0,
87			},
88		})
89		return decision
90	} catch {
91		span.update({
92			level: 'WARNING',
93			statusMessage: 'Selector unavailable; existing workflow retained',
94		})
95		return selectDestination('uncertain', 0, input.active)
96	} finally {
97		span.end()
98	}
99}
100