WolfDavid's picture
feat(02-06): analyze/translate/languageInfo on the shared turn loop; tokens on the turn event
969caf4
Raw History Blame
19.9 kB
// THIS MODULE IS TRANSPORT-AGNOSTIC AND IS THE ONLY IMPLEMENTATION OF THE TURN LOOP.
// Both avatar.js (inline gr.HTML) and avatar-iframe.js (postMessage) import it.
// Do NOT add turn behaviour to either transport file - add it here, or the iframe
// fallback silently loses the feature. tests/test_transport_seam.py enforces this.
//
// The renderer is reachable only through the seven-method stagePort. The host bridge is
// reachable only through getServer(), which is a GETTER rather than a value so this
// module never holds a host object and can be constructed before the bridge exists.
//
// Push-to-talk landed in wave 4. Note that mic.js and asr.js are imported HERE and
// constructed lazily by this module rather than being handed in by a transport: mic
// capture and ASR run in the parent document under BOTH transports, and the iframe
// transport has no AudioContext of its own to lend. Constructing them here means the
// two transport files needed zero lines of change to gain push-to-talk, which is the
// strongest possible form of the guarantee the seam exists to give. Both are still
// injectable through the factory so a harness can substitute a different model.
//
// The turn itself landed in wave 5, in this file and nowhere else, so the same holds:
// dispatchTurn, replay and requestSlower work under both transports because there is
// exactly one implementation of each. The bridge object getServer() returns exposes
// the host-registered functions (turn, greeting) as async methods; this module calls
// them with ONE argument each, because the host's bridge packs multiple arguments into
// a list and the far side would receive that list as a single positional.
//
// Plan 01-11: every entry point below that a tap can reach - dispatchTurn, replay,
// requestSlower, startListening - calls stagePort.unlockAudio() as its FIRST statement,
// before any await. The facade wraps these methods synchronously, so that statement runs
// inside the gesture's own call stack; anything after an await does not. Browsers that
// gate audio behind a gesture (iOS Safari; Chromium in a cross-origin embed) refuse a
// resume() issued after the server round trip, which is exactly where the only other
// resume() lives (audio-queue.js), and the owner's phone was silent for that reason.
//
// Plan 02-06: analyze / translate / languageInfo are host round trips implemented HERE so
// both transports get them; the host page renders tokens and caches translations (D-17);
// this module holds no per-line cache. Every bridge call packs ONE payload object (the
// host's bridge would turn two arguments into a list), and every result goes through the
// same two checks dispatchTurn applies: `undefined` (the host's client swallowed an HTTP
// error) and `{error}` (the far side refused the request without raising).
import { createAsr } from './asr.js';
import { createMic, isHallucination, REJECT } from './mic.js';
/** VOICEVOX speedScale for the "Slower" re-read. Divides every phoneme length. */
export const SLOWER_SPEED = 0.75;
/**
* @param {object} opts
* @param {object} opts.stagePort the renderer boundary
* @param {Function} [opts.emit] the facade's event bus
* @param {Function} [opts.getServer] returns the host bridge, or null when there is none
* @param {object} [opts.mic] overrides the mic built from avatar/mic.js
* @param {object} [opts.asr] overrides the ASR built from avatar/asr.js
* @param {object} [opts.asrOptions] model/dtype/device overrides for the default ASR
* @param {object} [opts.micOptions] capture overrides for the default mic
*/
export function createTurnLoop({
stagePort,
emit = () => {},
getServer = () => null,
mic = null,
asr = null,
asrOptions = {},
micOptions = {},
} = {}) {
if (!stagePort) throw new Error('createTurnLoop: a stagePort is required');
// The turn loop's own slice of __debug. The facade merges this object; it does not
// own it, and this module does not own the facade's. Every key is initialised here
// rather than on first use, because the parity suite compares __debug KEY SETS across
// transports and a lazily-added key would make that comparison time-dependent.
const state = {
thinking: false,
listening: false,
speaking: false,
lastTurnId: null,
replayCount: 0,
turnCount: 0,
// The number that matters: dispatch (Enter, click or mic release) to the first
// scheduled audio sample, in milliseconds, for the most recent turn.
lastTurnMs: null,
lastReplayMs: null,
lastStageTimings: null,
lastSubtitle: null,
lastSpeed: null,
lastError: null,
asrTier: null,
asrModel: null,
lastTranscript: null,
micRejectedCount: 0,
micLastRejectReason: null,
// Plan 02-06: the language bridge. Tokens ride on every directive; the learner's own
// lines are tokenised through analyze(); translations are counted and timed here and
// cached by the host page, never by this module (D-17).
analyzeCount: 0,
lastAnalyzeMs: null,
lastTokenCount: null,
lastTokens: null,
translateCount: 0,
lastTranslateMs: null,
lastTranslation: null,
lastTranslateError: null,
languageInfo: null,
};
// Set when a turn or a replay has been dispatched and its speech-start has not yet
// been observed. speech-start is the event that closes the "thinking" window and
// stamps lastTurnMs, so it is measured where the event arrives - the same place under
// both transports - rather than guessed at from the speak() promise.
let pendingDispatchAt = null;
let pendingReplayAt = null;
const now = () => performance.now();
function setThinking(value) {
const on = !!value;
state.thinking = on;
stagePort.setThinking(on);
return on;
}
function setListening(value) {
const on = !!value;
state.listening = on;
stagePort.setListening(on);
return on;
}
const micInstance =
mic ||
createMic({
emit,
onListening: setListening,
// Push-to-talk exists to make acoustic feedback impossible, so the mic refuses to
// open while the avatar is thinking or speaking. mic.js adds the 200 ms tail after
// speech-end on top of this.
isBusy: () => state.thinking || state.speaking,
...micOptions,
});
const asrInstance = asr || createAsr({ emit, ...asrOptions });
/**
* The facade calls this for every event it fans out, under both transports - the
* iframe transport forwards the frame's events through the same bus. It is how this
* module learns that speech started or ended without holding a reference to the
* audio path, which belongs to the stage.
*/
function observe(name, data) {
if (name === 'speech-start') {
state.speaking = true;
if (pendingDispatchAt !== null) {
state.lastTurnMs = Math.round(now() - pendingDispatchAt);
pendingDispatchAt = null;
performance.mark('turn:speech-start');
try {
performance.measure('turn:dispatch-to-speech', 'turn:dispatch', 'turn:speech-start');
} catch {
/* a mark was cleared; the number is already in lastTurnMs */
}
// The thinking pose clears HERE, at speech-start, not when the response arrives:
// decode and scheduling still sit between the two, and the face must not go idle
// while the visitor is still waiting to hear something.
setThinking(false);
emit('latency', {
turnId: state.lastTurnId,
lastTurnMs: state.lastTurnMs,
timings: state.lastStageTimings,
speed: state.lastSpeed,
});
}
if (pendingReplayAt !== null) {
state.lastReplayMs = Math.round(now() - pendingReplayAt);
pendingReplayAt = null;
performance.mark('replay:speech-start');
}
} else if (name === 'speech-end') {
state.speaking = false;
micInstance.noteSpeechEnd();
} else if (name === 'asr-tier' && data) {
state.asrTier = data.tier ?? null;
state.asrModel = data.model ?? null;
}
}
function busyReason() {
if (state.thinking) return 'the avatar is still thinking about the last turn';
if (state.speaking) return 'the avatar is still speaking';
return null;
}
function fail(where, err) {
const message = String(err?.message ?? err);
state.lastError = message;
emit('error', { message, where });
return err instanceof Error ? err : new Error(message);
}
/**
* The host bridge, or a throw that names the missing function. The standalone stage
* (asr-harness.html, stage.html) has no bridge at all; a host that forgot to register
* a server function has a bridge without the method. Both are the same failure to a
* caller.
*/
function requireBridge(name, without) {
const bridge = getServer();
if (!bridge || typeof bridge[name] !== 'function') {
throw new Error(`no host bridge: the standalone stage has ${without}`);
}
return bridge;
}
/**
* The two checks every bridge result needs. The host's client swallows an HTTP error
* into `undefined`, and the far side answers a bad request with {error} rather than
* raising, so both are checked before a result is trusted.
*/
function checkResult(result, what) {
if (result === undefined || result === null) {
throw new Error(`the host returned nothing for ${what} - see its log`);
}
if (result.error) throw new Error(result.error);
return result;
}
/**
* One turn: text in, speech out. Resolves at speech-end with a summary of the turn.
*
* Order matters and is asserted statically: the thinking pose engages BEFORE the
* first await, because switching the pose when the response arrives would forfeit
* the entire latency the pose exists to cover.
*
* @param {string} text what the avatar should say back
* @param {object} [opts]
* @param {number} [opts.speed=1.0] VOICEVOX speedScale; SLOWER_SPEED for the re-read
* @param {boolean} [opts.greeting] ignore text and speak the server's fixed greeting
*/
async function dispatchTurn(text, { speed = 1.0, greeting = false } = {}) {
stagePort.unlockAudio(); // first, synchronously: still inside the tap that got us here
const busy = busyReason();
if (busy) throw fail('dispatchTurn', new Error(busy));
let bridge;
try {
bridge = requireBridge(greeting ? 'greeting' : 'turn', 'nothing to synthesise with');
} catch (err) {
throw fail('dispatchTurn', err);
}
performance.mark('turn:dispatch');
const dispatchedAt = now();
pendingDispatchAt = dispatchedAt;
setThinking(true);
emit('turn-start', { text: greeting ? null : text, speed, greeting });
try {
const directive = checkResult(
greeting ? await bridge.greeting() : await bridge.turn({ text: String(text ?? ''), speed }),
'this turn'
);
performance.mark('turn:response');
const responseMs = Math.round(now() - dispatchedAt);
if (!directive.audio_url || !Array.isArray(directive.timeline)) {
throw new Error('the host returned a directive with no audio or no timeline');
}
state.turnCount += 1;
state.lastTurnId = directive.turn_id ?? null;
state.lastSubtitle = directive.subtitle ?? null;
state.lastSpeed = directive.speed ?? speed;
state.lastStageTimings = directive.timings ?? null;
// Tokens are optional on the wire (an older host, or an analysis that failed and
// yielded []): the turn still speaks, the host page just has nothing to make tappable.
state.lastTokens = Array.isArray(directive.tokens) ? directive.tokens : [];
state.lastTokenCount = state.lastTokens.length;
state.lastError = null;
emit('turn', {
turnId: state.lastTurnId,
subtitle: state.lastSubtitle,
speed: state.lastSpeed,
timings: state.lastStageTimings,
tokens: state.lastTokens,
responseMs,
greeting,
});
// Resolves at speech-end. speech-start arrives through observe() on the way.
const played = await stagePort.speak({
audioUrl: directive.audio_url,
timeline: directive.timeline,
subtitle: directive.subtitle,
expression: directive.expression,
turnId: directive.turn_id,
});
return {
turnId: state.lastTurnId,
subtitle: state.lastSubtitle,
speed: state.lastSpeed,
timings: state.lastStageTimings,
responseMs,
lastTurnMs: state.lastTurnMs,
duration: played?.duration ?? null,
};
} catch (err) {
pendingDispatchAt = null;
setThinking(false);
throw fail('dispatchTurn', err);
}
}
return {
state,
getServer,
observe,
mic: micInstance,
asr: asrInstance,
setThinking,
setListening,
dispatchTurn,
/**
* Re-play the cached directive. There is deliberately no network access of any
* kind in here: the deployed suite asserts zero requests with a browser request
* listener, and a cache miss must fail rather than quietly re-download. The stage
* keeps the decoded AudioBuffer and the timeline it last spoke.
*/
async replay() {
stagePort.unlockAudio();
const busy = busyReason();
if (busy) throw fail('replay', new Error(busy));
if (state.turnCount === 0) throw fail('replay', new Error('nothing has been said yet'));
state.replayCount += 1;
performance.mark('replay:dispatch');
pendingReplayAt = now();
emit('replay', { turnId: state.lastTurnId, subtitle: state.lastSubtitle });
try {
return await stagePort.replayCached();
} catch (err) {
pendingReplayAt = null;
throw fail('replay', err);
}
},
/**
* Re-synthesise the last utterance at SLOWER_SPEED. This IS a host round trip and
* must be: the timeline has to be rebuilt from the re-synthesised query, because
* speedScale divides every phoneme and a timeline scaled here would drift by exactly
* the speed ratio against the new audio.
*/
async requestSlower() {
stagePort.unlockAudio(); // dispatchTurn unlocks too; this covers the early throw below
if (!state.lastSubtitle) {
throw fail('requestSlower', new Error('nothing to slow down yet - say something first'));
}
return dispatchTurn(state.lastSubtitle, { speed: SLOWER_SPEED });
},
/**
* Tokenise a learner line through the host's language core (plan 02-06). No audio is
* involved, so no unlockAudio() and no busy check: the host page calls this for the
* text the learner typed or the transcript ASR produced, then renders the tokens.
*
* @param {string} text
* @returns {Promise<{tokens: object[], timings: object|null}>}
*/
async analyze(text) {
if (typeof text !== 'string' || !text.trim()) {
throw fail('analyze', new Error('analyze: text is required'));
}
const t0 = now();
try {
const bridge = requireBridge('analyze', 'no language core to analyze with');
const result = checkResult(await bridge.analyze({ text }), 'analyze');
state.analyzeCount += 1;
state.lastAnalyzeMs = Math.round(now() - t0);
return {
tokens: Array.isArray(result.tokens) ? result.tokens : [],
timings: result.timings ?? null,
};
} catch (err) {
throw fail('analyze', err);
}
},
/**
* Translate one line to English on the host's CPU (D-15). The host page caches the
* answer per line for the session (D-17) - this module deliberately does not, so the
* cache has exactly one owner. `lineId` is echoed back so the page can file it.
*
* @param {string} text
* @param {string|null} [lineId]
* @returns {Promise<{text: string, lineId: string|null, timings: object|null, ms: number}>}
*/
async translate(text, lineId = null) {
if (typeof text !== 'string' || !text.trim()) {
throw fail('translate', new Error('translate: text is required'));
}
const t0 = now();
try {
const bridge = requireBridge('translate', 'no translator to translate with');
const result = checkResult(await bridge.translate({ text, line_id: lineId }), 'translate');
const ms = Math.round(now() - t0);
state.translateCount += 1;
state.lastTranslateMs = ms;
state.lastTranslation = result.text ?? null;
state.lastTranslateError = null;
return {
text: result.text,
lineId: result.line_id ?? lineId,
timings: result.timings ?? null,
ms,
};
} catch (err) {
state.lastTranslateError = String(err?.message ?? err);
throw fail('translate', err);
}
},
/**
* The host's language-asset facts (sizes, warm-up measurement, container env), read
* once through the bridge and kept on the debug slice so a deployed probe can assert
* on them after the fact.
*/
async languageInfo() {
try {
const bridge = requireBridge('language_info', 'no language core to describe');
const info = checkResult(await bridge.language_info(), 'languageInfo');
state.languageInfo = info;
return info;
} catch (err) {
throw fail('languageInfo', err);
}
},
/**
* pointerdown on the push-to-talk control. The host binds the control; the behaviour
* is here so both transports get it from one implementation.
*
* @returns {Promise<boolean>} whether capture actually started
*/
async startListening() {
// The stage's context, not the mic's: the reply to what is about to be said
// arrives seconds after this pointerdown, and the pointerdown is the only gesture.
stagePort.unlockAudio();
const started = await micInstance.start();
if (!started) {
state.micRejectedCount = micInstance.__debug.rejectedCount;
state.micLastRejectReason = micInstance.__debug.lastRejectReason;
}
return started;
},
/**
* pointerup. Gate -> ASR -> 'transcript'.
*
* A gated-out push emits NOTHING - not an empty transcript, which every downstream
* consumer would then have to special-case - and returns null.
*
* @returns {Promise<string|null>} the transcript, or null when nothing survived
*/
async stopListening() {
const utterance = await micInstance.stop();
state.micRejectedCount = micInstance.__debug.rejectedCount;
state.micLastRejectReason = micInstance.__debug.lastRejectReason;
if (!utterance.ok) return null;
let result;
try {
result = await asrInstance.transcribe(utterance.samples, utterance.sampleRate);
} catch (err) {
fail('stopListening', err);
return null;
}
state.asrTier = asrInstance.getTier();
state.asrModel = asrInstance.getModel();
// Second line of defence: the audio passed the gate but Whisper still produced
// subtitle boilerplate. Only short pushes are eligible - see mic.js.
if (!result.text || isHallucination(result.text, utterance.durationMs)) {
micInstance.noteTranscriptRejected(
result.text ? REJECT.HALLUCINATION : REJECT.NO_AUDIO
);
state.micRejectedCount = micInstance.__debug.rejectedCount;
state.micLastRejectReason = micInstance.__debug.lastRejectReason;
return null;
}
state.lastTranscript = result.text;
emit('transcript', {
text: result.text,
durationMs: Math.round(utterance.durationMs),
tier: result.tier,
gated: false,
});
return result.text;
},
};
}