diff --git a/lib/tasks/gather.js b/lib/tasks/gather.js index db4a0ff9..2d3d23fa 100644 --- a/lib/tasks/gather.js +++ b/lib/tasks/gather.js @@ -901,6 +901,29 @@ class TaskGather extends SttTask { this.playComplete = true; } + /* a deepgram is_final result finalizes the audio in [start, start + duration]; if that window + covers the last unprocessed interim word then the word has been finalized or dropped (e.g. an + empty final after line noise), so it is no longer pending and must not hold off UtteranceEnd */ + _clearFinalizedUnprocessedWord(evt) { + if (!evt.is_final || !this._dgTimeOfLastUnprocessedWord) return false; + if (typeof evt.start !== 'number' || typeof evt.duration !== 'number') return false; + if (evt.start + evt.duration < this._dgTimeOfLastUnprocessedWord) return false; + this.logger.debug(`Gather:_onTranscription - deepgram finalized audio through ${evt.start + evt.duration}, ` + + `clearing unprocessed word end time ${this._dgTimeOfLastUnprocessedWord}`); + this._dgTimeOfLastUnprocessedWord = null; + return true; + } + + /* UtteranceEnd was deferred for a word that deepgram has now finalized as empty (dropped); with no + new words deepgram will not send another UtteranceEnd, so return the buffer now */ + _resolveDeferredUtteranceEnd() { + if (this.resolved || this.killed || this._bufferedTranscripts.length === 0) return; + this.logger.info('Gather:_onTranscription - deferred UtteranceEnd satisfied, return buffered transcript'); + const evt = this.consolidateTranscripts(this._bufferedTranscripts, 1, this.language, this.vendor); + this._bufferedTranscripts = []; + this._resolve('speech', evt); + } + _onTranscription(cs, ep, evt, fsEvent) { // check if we are in graceful shutdown mode if (ep.gracefulShutdownResolver) { @@ -928,6 +951,7 @@ class TaskGather extends SttTask { // eslint-disable-next-line max-len if (utteranceTime && this._dgTimeOfLastUnprocessedWord && utteranceTime < this._dgTimeOfLastUnprocessedWord && utteranceTime != -1) { this.logger.info('Gather:_onTranscription - got UtteranceEnd with unprocessed words, continue listening'); + this._dgUtteranceEndDeferred = true; } else { this.logger.info('Gather:_onTranscription - got UtteranceEnd from deepgram, return buffered transcript'); @@ -942,6 +966,14 @@ class TaskGather extends SttTask { this.logger.debug('Gather:_onTranscription - discarding Metadata event from deepgram'); return; } + if (this.vendor === 'deepgram' && this._clearFinalizedUnprocessedWord(evt) && this._dgUtteranceEndDeferred) { + /* a final with words means the caller is still talking; deepgram will send a new UtteranceEnd */ + this._dgUtteranceEndDeferred = false; + if (!evt.channel?.alternatives?.[0]?.transcript) { + /* resolve once the rest of this method has handled this final */ + setImmediate(() => this._resolveDeferredUtteranceEnd()); + } + } evt = this.normalizeTranscription(evt, this.vendor, 1, this.language, this.shortUtterance, this.data.recognizer.punctuation); diff --git a/test/unit/gather-deepgram-utterance-end.test.js b/test/unit/gather-deepgram-utterance-end.test.js new file mode 100644 index 00000000..9affec9d --- /dev/null +++ b/test/unit/gather-deepgram-utterance-end.test.js @@ -0,0 +1,131 @@ +const test = require('node:test'); +const assert = require('node:assert'); + +const TaskGather = require('../../lib/tasks/gather'); + +const logger = {debug() {}, info() {}, error() {}, warn() {}}; + +/* a continuous-asr deepgram gather; the constructor is not run, only what _onTranscription touches is set */ +const makeGather = () => { + const resolved = []; + const gather = Object.create(TaskGather.prototype); + const {normalizeTranscription, consolidateTranscripts} = require('../../lib/utils/transcription-utils')(logger); + const cs = {emit() {}, calculateSttLatency() {}, callGone: false}; + Object.assign(gather, { + logger, + normalizeTranscription, + consolidateTranscripts, + vendor: 'deepgram', + language: 'en-US', + data: {recognizer: {}}, + isContinuousAsr: true, + asrTimeout: 2000, + _bufferedTranscripts: [], + _sonioxTranscripts: [], + cs, + eventIsForOurBug: () => true, + doesVendorContinueListeningAfterFinalTranscript: () => true, + _clearTimer: () => false, + _startTimer() {}, + _startAsrTimer() {}, + _clearAsrTimer() {}, + _resolve(reason, evt) { + this.resolved = true; + resolved.push({reason, transcript: evt.alternatives[0].transcript}); + } + }); + const fsEvent = {getHeader: () => undefined}; + /* a deferred UtteranceEnd resolves on setImmediate, after the final has been handled */ + const send = async(evt) => { + gather._onTranscription(cs, {}, evt, fsEvent); + await new Promise((resolve) => setImmediate(resolve)); + }; + return {gather, send, resolved}; +}; + +const results = ({start, duration, is_final, speech_final = false, words = []}) => ({ + type: 'Results', + start, + duration, + is_final, + speech_final, + channel: { + alternatives: [{ + transcript: words.map((w) => w.word).join(' '), + confidence: 0.9, + words + }] + } +}); + +const firstUtterance = results({ + start: 0, duration: 3, is_final: true, + words: [{word: 'pay', start: 0.5, end: 1.0}, {word: 'my', start: 1.1, end: 1.4}, {word: 'bill', start: 1.5, end: 2.9}] +}); +const strayInterim = results({start: 3, duration: 1.5, is_final: false, words: [{word: 'uh', start: 4.0, end: 4.2}]}); +const utteranceEnd = {type: 'UtteranceEnd', channel: [0, 1], last_word_end: 2.9}; + +test('an empty final covering a stray interim word lets UtteranceEnd return the buffer', async() => { + const {gather, send, resolved} = makeGather(); + await send(firstUtterance); + await send(strayInterim); + assert.strictEqual(gather._dgTimeOfLastUnprocessedWord, 4.2); + + await send(results({start: 3, duration: 2, is_final: true})); + assert.strictEqual(gather._dgTimeOfLastUnprocessedWord, null); + + await send(utteranceEnd); + assert.deepStrictEqual(resolved, [{reason: 'speech', transcript: 'pay my bill'}]); +}); + +test('an empty final that ends before the interim word keeps UtteranceEnd waiting', async() => { + const {gather, send, resolved} = makeGather(); + await send(firstUtterance); + await send(strayInterim); + + /* deepgram finalized [3, 4.0) but the word ending at 4.2 falls in the next segment */ + await send(results({start: 3, duration: 1.0, is_final: true})); + assert.strictEqual(gather._dgTimeOfLastUnprocessedWord, 4.2); + + await send(utteranceEnd); + assert.deepStrictEqual(resolved, []); +}); + +test('an empty final arriving after a deferred UtteranceEnd returns the buffer', async() => { + const {send, resolved} = makeGather(); + await send(firstUtterance); + await send(strayInterim); + + await send(utteranceEnd); + assert.deepStrictEqual(resolved, []); + + /* deepgram sends no second UtteranceEnd until new speech, so this final must end the gather */ + await send(results({start: 3, duration: 2, is_final: true})); + assert.deepStrictEqual(resolved, [{reason: 'speech', transcript: 'pay my bill'}]); +}); + +test('a final with words after a deferred UtteranceEnd keeps listening for the next UtteranceEnd (#1088)', async() => { + const {send, resolved} = makeGather(); + await send(firstUtterance); + await send(results({start: 3, duration: 1.5, is_final: false, words: [{word: 'please', start: 3.8, end: 4.2}]})); + + await send(utteranceEnd); + assert.deepStrictEqual(resolved, []); + + /* the caller may still be talking, so the words alone must not end the gather */ + await send(results({start: 3, duration: 1.5, is_final: true, words: [{word: 'please', start: 3.8, end: 4.2}]})); + assert.deepStrictEqual(resolved, []); + + await send({...utteranceEnd, last_word_end: 4.2}); + assert.deepStrictEqual(resolved, [{reason: 'speech', transcript: 'pay my bill please'}]); +}); + +test('a deferred UtteranceEnd is not resolved by a final that ends before the pending word', async() => { + const {send, resolved} = makeGather(); + await send(firstUtterance); + await send(strayInterim); + + await send(utteranceEnd); + await send(results({start: 3, duration: 1.0, is_final: true})); + assert.deepStrictEqual(resolved, []); +});