File size: 14,984 Bytes
7146e75
 
 
 
 
6b68c4d
7146e75
 
b4ca58d
 
 
 
 
 
 
 
f4874f3
 
 
 
 
 
 
6b68c4d
 
 
 
 
 
 
 
b4ca58d
 
 
7146e75
f4874f3
 
7146e75
 
 
 
 
 
b4ca58d
 
 
 
7146e75
b4ca58d
 
 
 
 
 
 
 
 
7146e75
 
 
b4ca58d
 
 
7146e75
 
 
f4874f3
7146e75
 
f4874f3
 
 
 
 
 
 
 
 
b4ca58d
 
 
 
f4874f3
7146e75
 
f4874f3
 
 
 
 
 
 
 
b4ca58d
f4874f3
 
 
 
 
 
 
 
 
b4ca58d
 
f4874f3
b4ca58d
 
 
 
 
 
f4874f3
b4ca58d
 
 
f4874f3
b4ca58d
 
 
 
 
 
 
 
 
 
 
 
 
f4874f3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
b4ca58d
f4874f3
b4ca58d
 
 
 
 
 
 
f4874f3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
6b68c4d
f4874f3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
7146e75
 
 
b4ca58d
 
 
f4874f3
 
 
7146e75
 
 
f4874f3
 
 
7146e75
 
6b68c4d
f4874f3
 
 
7146e75
f4874f3
 
 
7146e75
 
 
f4874f3
 
 
 
 
 
 
 
 
 
 
 
6b68c4d
f4874f3
 
7146e75
f4874f3
7146e75
 
b4ca58d
f4874f3
 
b4ca58d
 
 
 
6b68c4d
 
 
b4ca58d
f4874f3
 
 
 
b4ca58d
 
 
 
 
 
 
 
 
 
 
 
 
 
f4874f3
b4ca58d
 
 
 
 
 
f4874f3
b4ca58d
 
 
 
 
 
 
 
 
 
 
 
 
f4874f3
b4ca58d
 
 
 
 
 
 
 
 
 
 
 
7146e75
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
// 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.

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,
  };

  // 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);
  }

  /**
   * 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));

    const bridge = getServer();
    if (!bridge || typeof bridge.turn !== 'function') {
      throw fail(
        'dispatchTurn',
        new Error('no host bridge: the standalone stage has nothing to synthesise with')
      );
    }

    performance.mark('turn:dispatch');
    const dispatchedAt = now();
    pendingDispatchAt = dispatchedAt;
    setThinking(true);
    emit('turn-start', { text: greeting ? null : text, speed, greeting });

    try {
      const directive = greeting
        ? await bridge.greeting()
        : await bridge.turn({ text: String(text ?? ''), speed });
      performance.mark('turn:response');
      const responseMs = Math.round(now() - dispatchedAt);

      // 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.
      if (directive === undefined || directive === null) {
        throw new Error('the host returned nothing for this turn - see its log');
      }
      if (directive.error) throw new Error(directive.error);
      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;
      state.lastError = null;
      emit('turn', {
        turnId: state.lastTurnId,
        subtitle: state.lastSubtitle,
        speed: state.lastSpeed,
        timings: state.lastStageTimings,
        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 });
    },

    /**
     * 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;
    },
  };
}