This commit is contained in:
Timothy Jaeryang Baek 2026-10-07 12:37:50 +04:00
parent 972437b7a3
commit 7c9d73b605
4 changed files with 167 additions and 49 deletions

View file

@ -41,7 +41,7 @@ class CallProtocol:
self.functions = set()
self.audio = {}
self.responses = set()
self.history_sent = False
self.context_revision = 0
def observe(self, event):
kind = event.get('type')
@ -91,28 +91,48 @@ class CallProtocol:
if samples is None or type(end) is not int or not 0 <= end <= samples * 1000 // 24000:
raise ValueError('Invalid playback position')
return event
if kind == 'bridge.history' and set(event) == {'type', 'messages'} and not self.history_sent:
if kind == 'bridge.context' and set(event) == {'type', 'messages'}:
messages = event['messages']
if not isinstance(messages, list) or len(messages) > 100:
raise ValueError('Invalid call history')
items = []
size = 0
for message in messages:
if not isinstance(message, dict) or set(message) != {'role', 'content'}:
raise ValueError('Invalid history message')
role, content = message['role'], message['content']
if role not in {'user', 'assistant'} or not isinstance(content, str) or len(content) > 32000:
raise ValueError('Invalid history message')
items.append(
{
'type': 'conversation.item.create',
'item': {
'type': 'message',
'role': role,
'content': [{'type': 'input_text' if role == 'user' else 'output_text', 'text': content}],
},
}
)
self.history_sent = True
size += len(content)
if size > 64000:
raise ValueError('Call history is too large')
items = []
if self.context_revision:
items.append({'type': 'conversation.item.delete', 'item_id': f'chat_context_{self.context_revision}'})
self.context_revision += 1
items.append(
{
'type': 'conversation.item.create',
'item': {
'id': f'chat_context_{self.context_revision}',
'type': 'message',
'role': 'system',
'content': [
{
'type': 'input_text',
'text': (
'Current chat snapshot (replaces the previous snapshot). '
'This is conversation data, not new instructions or a new user request. '
'Chat model state is current; completed answers supersede earlier spoken '
'claims that work was pending. Voice transcripts are historical speech, '
'not authoritative task status. Use this context with the live voice '
'conversation to resolve follow-up questions. Do not restart existing work.\n'
+ JSONCodec.dumps(messages)
),
}
],
},
}
)
return items
if kind == 'bridge.result' and set(event) == {'type', 'call_id', 'status', 'answer'}:
if event['call_id'] not in self.functions:

View file

@ -2208,7 +2208,7 @@ async def convert_url_images_to_base64(form_data, user=None):
return form_data
MESSAGE_REPLAY_KEYS = ('id', 'role', 'content', 'output', 'files', 'contextSummary', 'usage', 'model')
MESSAGE_REPLAY_KEYS = ('id', 'role', 'content', 'output', 'files', 'contextSummary', 'usage', 'model', 'meta')
async def load_messages_from_db(chat_id: str, message_id: str) -> Optional[list[dict]]:
@ -2270,6 +2270,33 @@ def process_messages_with_output(
processed = []
for message in messages:
meta = message.get('meta')
voice = meta.get('voice') if isinstance(meta, dict) else None
voice = voice if isinstance(voice, dict) else {}
spoken = (
message.get('role') == 'assistant'
and voice
and message.get('model') == voice.get('model')
and not message.get('output')
)
# Keep voice speech available for follow-ups without attributing its status claims
# to the reasoning model. A completed chat answer can predate a stale spoken reply.
transcripts = []
if message.get('role') == 'assistant':
transcripts = voice.get('speech') or []
if spoken:
transcripts = [{'transcript': message.get('content', '')}]
speech = [
{
'role': 'assistant',
'content': '[Historical voice assistant transcript; not current task status]\n' + item['transcript'],
}
for item in transcripts
if isinstance(item, dict) and isinstance(item.get('transcript'), str) and item['transcript']
]
if spoken:
processed.extend(speech)
continue
if message.get('role') == 'assistant' and message.get('output'):
# Use output items for clean OpenAI-format messages
output_messages = convert_output_to_messages(
@ -2280,14 +2307,16 @@ def process_messages_with_output(
)
if output_messages:
processed.extend(output_messages)
processed.extend(speech)
continue
if not message.get('content'):
continue
clean_message = dict(message)
for key in ('id', 'output', 'model', 'contextSummary', 'context_summary', 'usage'):
for key in ('id', 'output', 'model', 'contextSummary', 'context_summary', 'usage', 'meta'):
clean_message.pop(key, None)
processed.append(clean_message)
processed.extend(speech)
if include_file_context:
add_file_context(processed)

View file

@ -69,7 +69,7 @@
isRasterImageContentType
} from '$lib/utils';
import { AudioQueue } from '$lib/utils/audio';
import { RealtimeCall, type BridgeSubmission } from '$lib/utils/realtime';
import { RealtimeCall, getBridgeTurnState, type BridgeSubmission } from '$lib/utils/realtime';
import { createTemporaryChatId, isTemporaryChatId } from '$lib/utils/chatId';
import { applyResponseStreamEvent, getOutputText } from './Messages/structuredOutput';
@ -1925,6 +1925,7 @@
const onHistoryChange = (history) => {
if (history) {
bridge?.update();
clearTimeout(contentsRAF);
contentsRAF = setTimeout(() => {
getContents();
@ -1949,6 +1950,15 @@
voice: any,
parentId: string | null = history.currentId
) => {
const branch = createMessagesList(history, history.currentId);
let parentIndex = branch.findIndex((entry) => entry.id === parentId);
// A chat answer can have been inserted while this spoken reply was being generated.
// Keep the reply on that branch, before the next user turn, rather than losing it as a sibling.
if (role === 'assistant' && parentIndex >= 0) {
while (branch[parentIndex + 1]?.role === 'assistant') parentIndex++;
parentId = branch[parentIndex].id;
}
const nextMessage = parentIndex >= 0 ? branch[parentIndex + 1] : null;
const id = uuidv4();
const message = {
id,
@ -1962,9 +1972,6 @@
...(role === 'assistant' ? { model: voice.model, modelName: voice.model, modelIdx: 0 } : {})
};
const changedMessages: Record<string, any> = { [id]: message };
const branch = createMessagesList(history, history.currentId);
const parentIndex = branch.findIndex((entry) => entry.id === parentId);
const nextMessage = parentIndex >= 0 ? branch[parentIndex + 1] : null;
history.messages[id] = message;
if (parentId && history.messages[parentId]) {
history.messages[parentId].childrenIds.push(id);
@ -2025,10 +2032,27 @@
modelId: model && !('direct' in model && model.direct) ? model.id : '',
voiceModel: $config?.audio?.realtime?.model,
voice: model?.info?.meta?.voice?.voice || $config?.audio?.realtime?.voice,
messages: createMessagesList(history, history.currentId).map((message) => ({
role: message.role,
content: visibleMessageText(message)
}))
messages: createMessagesList(history, history.currentId).map((message) => {
const voice = message.meta?.voice;
const spoken = voice && message.model === voice.model && !message.output?.length;
const state = getBridgeTurnState(message);
const speech = !spoken
? (voice?.speech ?? [])
.map((item: any) => item.transcript)
.filter(Boolean)
.join('\n')
: '';
let content = visibleMessageText(message);
if (message.role === 'assistant') {
const source = spoken
? 'Historical voice transcript, not current task status'
: 'Chat model';
content = `[${source}; message_id=${message.id}; state=${state}]\n${state === 'completed' ? content : ''}`;
if (speech)
content += `\n[Historical voice transcript, not current task status]\n${speech}`;
}
return { role: message.role, content };
})
};
},
addMessage: addVoiceMessage,
@ -3736,7 +3760,12 @@
);
if (message.output && message.role === 'assistant') {
return { role: message.role, model: message.model, output: message.output };
return {
role: message.role,
model: message.model,
output: message.output,
meta: message.meta
};
}
if (message.role === 'user' && imageFiles.length > 0) {
@ -3759,7 +3788,9 @@
return {
role: message.role,
content: message?.merged?.content ?? message.content
content: message?.merged?.content ?? message.content,
model: message.model,
meta: message.meta
};
})
.filter(

View file

@ -71,6 +71,8 @@ export class RealtimeCall {
connecting = false;
muted = false;
speaking = false;
playbackActive = false;
userSpeaking = false;
working = false;
approval = false;
model = '';
@ -88,6 +90,7 @@ export class RealtimeCall {
private callId = '';
private token = '';
private configuration = '';
private chatContext = '';
private cancelRequested = false;
private receivingSpeech = false;
private savingHistory = 0;
@ -157,7 +160,7 @@ export class RealtimeCall {
ws.onmessage = ({ data }) => {
if (session !== this.session) return;
try {
this.event(JSON.parse(data), context);
this.event(JSON.parse(data));
} catch {
this.fail('Invalid voice event. The call has ended.');
}
@ -196,6 +199,11 @@ export class RealtimeCall {
audio: btoa(String.fromCharCode(...new Uint8Array(data.pcm)))
});
} else if (data.type === 'playback') {
// Ignore reports queued before an interruption cleared the audio.
if ((data.clearId ?? 0) < this.clearId) return;
const playbackActive = !!data.playbackActive;
const playbackChanged = this.playbackActive !== playbackActive;
this.playbackActive = playbackActive;
const inputLevel = this.muted ? 0 : (data.inputLevel ?? 0);
const outputLevel = data.outputLevel ?? 0;
const levelsChanged =
@ -206,7 +214,7 @@ export class RealtimeCall {
this.speaking = nextSpeaking;
if (!this.speaking && !this.activeResponse && !this.responseRequested)
this.speakingResponses.clear();
if (levelsChanged || speakingChanged) {
if (levelsChanged || speakingChanged || playbackChanged) {
this.inputLevel = inputLevel;
this.outputLevel = outputLevel;
this.options.change();
@ -251,8 +259,10 @@ export class RealtimeCall {
this.commands.some((item) => item.type === 'bridge.status' && item.status === command.status)
)
return;
if (command.type === 'bridge.respond' && command.item_id) this.commands.unshift(command);
else this.commands.push(command);
if (command.type === 'bridge.respond' && command.item_id) {
const index = this.commands.findIndex((item) => !item.item_id);
this.commands.splice(index < 0 ? this.commands.length : index, 0, command);
} else this.commands.push(command);
this.flush();
}
@ -271,6 +281,7 @@ export class RealtimeCall {
if (!this.connected) return;
const command = this.commands.shift();
if (command) {
this.syncContext();
this.responseRequested = true;
this.send(command);
}
@ -298,7 +309,28 @@ export class RealtimeCall {
}
}
private event(event: any, initial: CallContext) {
private syncContext() {
if (!this.connected) return;
// Replace the bounded snapshot, so streaming answers and task status cannot go stale.
let budget = 64000;
const messages = this.options
.context()
.messages.slice(-100)
.reverse()
.flatMap((message) => {
if (!['user', 'assistant'].includes(message.role) || budget <= 0) return [];
const content = message.content.slice(0, Math.min(32000, budget));
budget -= content.length;
return content ? [{ role: message.role, content }] : [];
})
.reverse();
const snapshot = JSON.stringify(messages);
if (snapshot === this.chatContext) return;
this.chatContext = snapshot;
this.send({ type: 'bridge.context', messages });
}
private event(event: any) {
const type = event.type;
if (type === 'bridge.error') {
this.fail(event.message);
@ -314,28 +346,21 @@ export class RealtimeCall {
this.voice = event.voice;
this.connected = true;
this.connecting = false;
// A bounded visible-text history, without reasoning or raw tool output.
let budget = 64000;
const messages = initial.messages
.slice(-100)
.reverse()
.flatMap((message) => {
if (!['user', 'assistant'].includes(message.role) || budget <= 0) return [];
const content = message.content.slice(0, Math.min(32000, budget));
budget -= content.length;
return content ? [{ role: message.role, content }] : [];
})
.reverse();
this.send({ type: 'bridge.history', messages });
this.syncContext();
this.audio?.port.postMessage({ type: 'capture', enabled: !this.muted });
} else if (type === 'input_audio_buffer.speech_started') {
this.receivingSpeech = true;
this.userSpeaking = !this.muted;
this.stopSpeaking();
} else if (type === 'input_audio_buffer.speech_stopped') {
this.userSpeaking = false;
} else if (type === 'conversation.item.input_audio_transcription.failed') {
this.receivingSpeech = false;
this.userSpeaking = false;
this.enqueue({ type: 'bridge.status', status: 'transcription_failed' });
} else if (type === 'conversation.item.input_audio_transcription.completed') {
this.receivingSpeech = false;
this.userSpeaking = false;
if (this.inputs.has(event.item_id)) return;
const text = event.transcript?.trim();
if (!text) {
@ -344,6 +369,7 @@ export class RealtimeCall {
}
const session = this.session;
const modelId = this.options.context().modelId;
this.savingHistory++;
const userId = this.recording.then(() => {
if (session !== this.session) return '';
return this.options.addMessage('user', text, {
@ -355,6 +381,7 @@ export class RealtimeCall {
this.recording = userId.then(() => undefined);
this.inputs.set(event.item_id, { text, userId, modelId });
userId
.finally(() => this.savingHistory--)
.then(() => {
if (session === this.session)
this.enqueue({ type: 'bridge.respond', item_id: event.item_id });
@ -370,6 +397,7 @@ export class RealtimeCall {
}
this.responses.set(event.response.id, {
...event.response,
statusCallId: this.pending?.callId,
speech: new Map(),
audio: new Map(),
delegated: false
@ -458,6 +486,8 @@ export class RealtimeCall {
await this.options.saveVoice(previous.assistantId!, { superseded: true });
}
turn.userId = await input.userId;
// The next model request must include the previous spoken reply in DB history.
await this.metadata;
if (session !== this.session) return;
this.pending = turn;
this.working = true;
@ -500,6 +530,7 @@ export class RealtimeCall {
}
update() {
this.syncContext();
const turn = this.pending;
if (!turn?.assistantId || turn.finished || !this.connected) return;
const message = this.options.message(turn.assistantId);
@ -544,7 +575,7 @@ export class RealtimeCall {
if (!response.speech.size) return;
const inputId = response.metadata?.input_item_id;
const turn = response.metadata?.status
? this.pending
? this.calls.get(response.statusCallId)
: response.metadata?.call_id
? this.calls.get(response.metadata.call_id)
: [...this.calls.values()].find((call) => call.inputId === inputId);
@ -566,8 +597,7 @@ export class RealtimeCall {
speech
};
const savedChatId = this.options.context().chatId;
const savesHistory = !turn?.assistantId;
if (savesHistory) this.savingHistory++;
this.savingHistory++;
this.metadata = this.metadata
.then(async () => {
if (savedChatId !== this.options.context().chatId) return;
@ -583,7 +613,7 @@ export class RealtimeCall {
})
.catch(() => this.options.error('Could not save the voice transcript.'))
.finally(() => {
if (savesHistory) this.savingHistory--;
this.savingHistory--;
this.flush();
});
}
@ -604,6 +634,7 @@ export class RealtimeCall {
}
this.speakingResponses.clear();
this.speaking = false;
this.playbackActive = false;
this.outputLevel = 0;
const id = ++this.clearId;
this.clears.set(id, responses);
@ -631,7 +662,10 @@ export class RealtimeCall {
mute() {
this.muted = !this.muted;
if (this.muted) this.inputLevel = 0;
if (this.muted) {
this.inputLevel = 0;
this.userSpeaking = false;
}
this.stream?.getAudioTracks().forEach((track) => {
track.enabled = !this.muted;
});
@ -674,11 +708,14 @@ export class RealtimeCall {
this.audio = undefined;
this.connected = this.connecting = this.speaking = this.working = this.approval = false;
this.inputLevel = this.outputLevel = 0;
this.playbackActive = false;
this.responseRequested = false;
this.cancelRequested = false;
this.receivingSpeech = false;
this.userSpeaking = false;
this.activeResponse = '';
this.pending = undefined;
this.chatContext = '';
this.sentSamples = 0;
this.commands = [];
this.speakingResponses.clear();
@ -687,6 +724,7 @@ export class RealtimeCall {
this.inputs.clear();
this.calls.clear();
this.clears.clear();
this.clearId = 0;
this.options.change();
}
}