initial checkin

This commit is contained in:
Dave Horton
2024-04-10 16:30:41 -04:00
commit 39d2611cc3
13 changed files with 2215 additions and 0 deletions
+4
View File
@@ -0,0 +1,4 @@
module.exports = ({logger, makeService}) => {
require('./llm-voicebot')({logger, makeService});
};
+110
View File
@@ -0,0 +1,110 @@
const OpenAIAssistant = require('../utils/ai-assistants/llm-openai');
const {system_instructions} = require('../../data/settings.json');
const {processStreamingResponse, streamingResponseComplete} = require('../utils/process-streaming-response');
const service = ({logger, makeService}) => {
const assistant = new OpenAIAssistant({
logger,
model: process.env.OPENAI_MODEL || 'gpt-4-turbo',
name: process.env.BOT_NAME || 'jambonz-llm-voicebot',
instructions: system_instructions
});
assistant.init();
const svc = makeService({path: '/llm-voicebot'});
svc.on('session:new', (session, path) => {
session.locals = { ...session.locals,
logger: logger.child({call_sid: session.call_sid}),
deepgramOptions: {
endpointing: 350,
utteranceEndMs: 1000,
},
turns: 0,
says: 0,
textOffset: 0,
assistant
};
session.locals.logger.info({session, path}, `new incoming call: ${session.call_sid}`);
session
.on('/user-input-event', onUserInputEvent.bind(null, session))
.on('close', onClose.bind(null, session))
.on('error', onError.bind(null, session));
session
.answer()
.pause({length: 0.5})
.config({
recognizer: {
vendor: 'default',
language: 'default',
deepgramOptions: session.locals.deepgramOptions
},
bargeIn: {
enable: true,
input: ['speech'],
actionHook: '/user-input-event',
sticky: true,
}
})
.say({text: 'Hi there! You are speaking to chat GPT. What would you like to know?'})
.reply();
});
};
const onUserInputEvent = async(session, evt) => {
const {logger} = session.locals;
logger.info({evt}, 'got speech evt');
switch (evt.reason) {
case 'speechDetected':
handleUserUtterance(session, evt);
break;
case 'timeout':
break;
default:
session.reply();
break;
}
};
const handleUserUtterance = async(session, evt) => {
const {logger, assistant, thread} = session.locals;
const {speech} = evt;
const userMessage = speech.alternatives[0].transcript;
logger.info({utterance: speech.alternatives[0]}, 'handling user utterance');
if (!thread) {
const thread = await assistant.createThread(userMessage);
thread
.on('botStreamingResponse', (response) => {
processStreamingResponse(session, response);
})
.on('botStreamingCompleted', () => {
streamingResponseComplete(session);
});
session.locals.thread = thread;
}
else {
thread.addUserMessage(speech.alternatives[0].transcript);
}
session.reply();
};
const onClose = (session) => {
const {logger, thread} = session.locals;
logger.info('call ended');
if (thread) {
thread.close();
}
};
const onError = (session, err) => {
const {logger} = session.locals;
logger.error(err, 'Error in call');
};
module.exports = service;
+98
View File
@@ -0,0 +1,98 @@
const OpenAI = require('openai');
const openai = new OpenAI();
const Emitter = require('events');
/* TODO: in future, add support for other AI Assistants beyond OpenAI */
class OpenAIThread extends Emitter {
constructor(logger, assistant_id, thread, run) {
super();
this.logger = logger;
this.assistant_id = assistant_id;
this.thread = thread;
this.run = run;
this._attachListeners();
}
_attachListeners() {
this.run
.on('event', (evt) => {
if (evt.event === 'thread.run.completed') {
this.emit('botStreamingCompleted');
}
})
.on('textDelta', (delta, snapshot) => {
this.emit('botStreamingResponse', snapshot.value);
})
.on('connect', () => this.logger.info('connected'))
.on('end', () => this.logger.info('ended'));
}
async addUserMessage(userMessage) {
this.logger.info('adding user message');
await openai.beta.threads.messages.create(this.thread.id, {
role: 'user',
content: userMessage,
});
this.run = await openai.beta.threads.runs.stream(this.thread.id, {
assistant_id: this.assistant_id
});
this._attachListeners();
}
async close() {
await openai.beta.threads.del(this.thread.id);
}
}
class OpenAIAssistant {
constructor({logger, model, name, instructions}) {
this.logger = logger;
this.model = model;
this.name = name;
this.instructions = instructions;
}
async init() {
try {
this.assistant = await openai.beta.assistants.create({
model: this.model,
name: this.name,
instructions: this.instructions,
});
} catch (err) {
this.logger.error(err);
throw err;
}
}
async createThread(initialUserMessage) {
/* create a thread */
try {
this.logger.info({initialUserMessage}, `creating thread with message: ${initialUserMessage}`);
const thread = await openai.beta.threads.create({
messages: [
{
role: 'user',
content: initialUserMessage,
},
],
});
/* run the thread */
this.logger.info('running thread');
const run = await openai.beta.threads.runs.stream(thread.id, {
assistant_id: this.assistant.id,
});
const t = new OpenAIThread(this.logger, this.assistant.id, thread, run);
return t;
} catch (err) {
this.logger.error(err);
throw err;
}
}
}
module.exports = OpenAIAssistant;
+66
View File
@@ -0,0 +1,66 @@
const processStreamingResponse = async(session, text) => {
const {logger, says, textOffset} = session.locals;
let spoken = false;
/* trim the leading part of the text we have already said */
const trimmed = text.substring(textOffset);
/**
* when we have a new response, say the first sentence as soon as available to get it out,
* after that fall back to speaking the rest of the response in chunks of paragraphs
*/
if (says === 0) {
const pos = trimmed.indexOf('.');
if (-1 !== pos) {
const firstSentence = trimmed.substring(0, pos + 1);
logger.info(`speaking first sentence: ${firstSentence}`);
session
.say({text: firstSentence})
.send({execImmediate: true});
session.locals.says++;
session.locals.textOffset = pos + 1;
spoken = true;
session.locals.unsent = trimmed.substring(pos + 1);
}
}
if (says > 0) {
const pos = trimmed.indexOf('\n\n');
if (-1 !== pos) {
const paragraph = trimmed.substring(0, pos);
logger.info(`speaking paragraph: ${paragraph}`);
session
.say({text: paragraph})
.send({execImmediate: false});
session.locals.says++;
session.locals.textOffset += (pos + 2);
spoken = true;
session.locals.unsent = trimmed.substring(pos + 2);
}
}
if (!spoken) {
session.locals.unsent = trimmed;
}
};
const streamingResponseComplete = async(session) => {
const {logger, unsent} = session.locals;
if (unsent) {
logger.info('sending final unsent text');
session
.say({text: unsent})
.send({execImmediate: false});
}
session.locals.turns++;
session.locals.says = 0;
session.locals.textOffset = 0;
session.locals.unsent = '';
};
module.exports = {
processStreamingResponse,
streamingResponseComplete
};