mirror of
https://github.com/jambonz/freeswitch-modules.git
synced 2026-01-25 02:08:27 +00:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e4a27ae133 | |||
| 19f20bf0e7 | |||
| b019a634bd | |||
| f0b304b8a1 | |||
| be3714465b | |||
| b495dba126 |
@@ -173,13 +173,12 @@ public:
|
||||
char *metadata,
|
||||
const char* awsAccessKeyId,
|
||||
const char* awsSecretAccessKey,
|
||||
const char* awsSessionToken,
|
||||
responseHandler_t responseHandler,
|
||||
errorHandler_t errorHandler) :
|
||||
m_bot(bot), m_alias(alias), m_region(region), m_sessionId(sessionId), m_finished(false), m_finishing(false), m_packets(0),
|
||||
m_pStream(nullptr), m_bPlayDone(false), m_bDiscardAudio(false)
|
||||
{
|
||||
Aws::String key(awsAccessKeyId);
|
||||
Aws::String secret(awsSecretAccessKey);
|
||||
Aws::String awsLocale(locale);
|
||||
Aws::Client::ClientConfiguration config;
|
||||
config.region = region;
|
||||
@@ -190,8 +189,11 @@ public:
|
||||
for (int i = 4; i < 20; i++) keySnippet[i] = 'x';
|
||||
keySnippet[19] = '\0';
|
||||
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "GStreamer %p ACCESS_KEY_ID %s\n", this, keySnippet);
|
||||
if (*awsAccessKeyId && *awsSecretAccessKey) {
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "GStreamer %p ACCESS_KEY_ID %s\n", this, keySnippet);
|
||||
if (*awsAccessKeyId && *awsSecretAccessKey && *awsSessionToken) {
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "using AWS creds %s %s %s\n", awsAccessKeyId, awsSecretAccessKey, awsSessionToken);
|
||||
m_client = Aws::MakeUnique<LexRuntimeV2Client>(ALLOC_TAG, AWSCredentials(awsAccessKeyId, awsSecretAccessKey, awsSessionToken), config);
|
||||
} else if (*awsAccessKeyId && *awsSecretAccessKey) {
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "using AWS creds %s %s\n", awsAccessKeyId, awsSecretAccessKey);
|
||||
m_client = Aws::MakeUnique<LexRuntimeV2Client>(ALLOC_TAG, AWSCredentials(awsAccessKeyId, awsSecretAccessKey), config);
|
||||
}
|
||||
@@ -540,7 +542,7 @@ static void *SWITCH_THREAD_FUNC lex_thread(switch_thread_t *thread, void *obj) {
|
||||
bool ok = true;
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "lex_thread: starting cb %p\n", (void *) cb);
|
||||
GStreamer* pStreamer = new GStreamer(cb->sessionId, cb->bot, cb->alias, cb->region, cb->locale,
|
||||
cb->intent, cb->metadata, cb->awsAccessKeyId, cb->awsSecretAccessKey,
|
||||
cb->intent, cb->metadata, cb->awsAccessKeyId, cb->awsSecretAccessKey, cb->awsSessionToken,
|
||||
cb->responseHandler, cb->errorHandler);
|
||||
if (!pStreamer) {
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "lex_thread: Error allocating streamer\n");
|
||||
@@ -641,6 +643,7 @@ extern "C" {
|
||||
memset(cb, sizeof(cb), 0);
|
||||
const char* awsAccessKeyId = switch_channel_get_variable(channel, "AWS_ACCESS_KEY_ID");
|
||||
const char* awsSecretAccessKey = switch_channel_get_variable(channel, "AWS_SECRET_ACCESS_KEY");
|
||||
const char* awsSessionToken = switch_channel_get_variable(channel, "AWS_SESSION_TOKEN");
|
||||
|
||||
if (!hasDefaultCredentials && (!awsAccessKeyId || !awsSecretAccessKey)) {
|
||||
switch_log_printf(SWITCH_CHANNEL_SESSION_LOG(session), SWITCH_LOG_ERROR,
|
||||
@@ -654,6 +657,7 @@ extern "C" {
|
||||
if (awsAccessKeyId && awsSecretAccessKey) {
|
||||
strncpy(cb->awsAccessKeyId, awsAccessKeyId, 128);
|
||||
strncpy(cb->awsSecretAccessKey, awsSecretAccessKey, 128);
|
||||
if (awsSessionToken) strncpy(cb->awsSessionToken, awsSessionToken, 1024);
|
||||
}
|
||||
else {
|
||||
strncpy(cb->awsAccessKeyId, std::getenv("AWS_ACCESS_KEY_ID"), 128);
|
||||
|
||||
@@ -30,6 +30,7 @@ struct cap_cb {
|
||||
char sessionId[256];
|
||||
char awsAccessKeyId[128];
|
||||
char awsSecretAccessKey[128];
|
||||
char awsSessionToken[1024];
|
||||
SpeexResamplerState *resampler;
|
||||
void* streamer;
|
||||
responseHandler_t responseHandler;
|
||||
|
||||
@@ -49,12 +49,11 @@ public:
|
||||
const char* region,
|
||||
const char* awsAccessKeyId,
|
||||
const char* awsSecretAccessKey,
|
||||
const char* awsSessionToken,
|
||||
responseHandler_t responseHandler
|
||||
) : m_sessionId(sessionId), m_bugname(bugname), m_finished(false), m_interim(interim), m_finishing(false), m_connected(false), m_connecting(false),
|
||||
m_packets(0), m_responseHandler(responseHandler), m_pStream(nullptr),
|
||||
m_audioBuffer(320 * (samples_per_second == 8000 ? 1 : 2), 15) {
|
||||
Aws::String key(awsAccessKeyId);
|
||||
Aws::String secret(awsSecretAccessKey);
|
||||
Aws::Client::ClientConfiguration config;
|
||||
if (region != nullptr && strlen(region) > 0) config.region = region;
|
||||
char keySnippet[20];
|
||||
@@ -63,8 +62,10 @@ public:
|
||||
for (int i = 4; i < 20; i++) keySnippet[i] = 'x';
|
||||
keySnippet[19] = '\0';
|
||||
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "GStreamer %p ACCESS_KEY_ID %s, region %s\n", this, keySnippet, region);
|
||||
if (*awsAccessKeyId && *awsSecretAccessKey) {
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "GStreamer %p ACCESS_KEY_ID %s, region %s\n", this, keySnippet, region);
|
||||
if (*awsAccessKeyId && *awsSecretAccessKey && *awsSessionToken) {
|
||||
m_client = Aws::MakeUnique<TranscribeStreamingServiceClient>(ALLOC_TAG, AWSCredentials(awsAccessKeyId, awsSecretAccessKey, awsSessionToken), config);
|
||||
} else if (*awsAccessKeyId && *awsSecretAccessKey) {
|
||||
m_client = Aws::MakeUnique<TranscribeStreamingServiceClient>(ALLOC_TAG, AWSCredentials(awsAccessKeyId, awsSecretAccessKey), config);
|
||||
}
|
||||
else {
|
||||
@@ -320,8 +321,8 @@ static void *SWITCH_THREAD_FUNC aws_transcribe_thread(switch_thread_t *thread, v
|
||||
struct cap_cb *cb = (struct cap_cb *) obj;
|
||||
bool ok = true;
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "transcribe_thread: starting cb %p\n", (void *) cb);
|
||||
GStreamer* pStreamer = new GStreamer(cb->sessionId, cb->bugname, cb->channels, cb->lang, cb->interim, cb->samples_per_second, cb->region, cb->awsAccessKeyId, cb->awsSecretAccessKey,
|
||||
cb->responseHandler);
|
||||
GStreamer* pStreamer = new GStreamer(cb->sessionId, cb->bugname, cb->channels, cb->lang, cb->interim, cb->samples_per_second,
|
||||
cb->region, cb->awsAccessKeyId, cb->awsSecretAccessKey, cb->awsSessionToken, cb->responseHandler);
|
||||
if (!pStreamer) {
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "transcribe_thread: Error allocating streamer\n");
|
||||
return nullptr;
|
||||
@@ -408,6 +409,7 @@ extern "C" {
|
||||
memset(cb, sizeof(cb), 0);
|
||||
const char* awsAccessKeyId = switch_channel_get_variable(channel, "AWS_ACCESS_KEY_ID");
|
||||
const char* awsSecretAccessKey = switch_channel_get_variable(channel, "AWS_SECRET_ACCESS_KEY");
|
||||
const char* awsSessionToken = switch_channel_get_variable(channel, "AWS_SESSION_TOKEN");
|
||||
const char* awsRegion = switch_channel_get_variable(channel, "AWS_REGION");
|
||||
cb->channels = channels;
|
||||
LanguageCode code = LanguageCodeMapper::GetLanguageCodeForName(lang);
|
||||
@@ -419,12 +421,12 @@ extern "C" {
|
||||
strncpy(cb->sessionId, switch_core_session_get_uuid(session), MAX_SESSION_ID);
|
||||
strncpy(cb->bugname, bugname, MAX_BUG_LEN);
|
||||
|
||||
if (awsAccessKeyId && awsSecretAccessKey && awsRegion) {
|
||||
if (awsRegion) strncpy(cb->region, awsRegion, MAX_REGION);
|
||||
if (awsAccessKeyId && awsSecretAccessKey) {
|
||||
switch_log_printf(SWITCH_CHANNEL_SESSION_LOG(session), SWITCH_LOG_DEBUG, "Using channel vars for aws authentication\n");
|
||||
strncpy(cb->awsAccessKeyId, awsAccessKeyId, 128);
|
||||
strncpy(cb->awsSecretAccessKey, awsSecretAccessKey, 128);
|
||||
strncpy(cb->region, awsRegion, MAX_REGION);
|
||||
|
||||
if (awsSessionToken) strncpy(cb->awsSessionToken, awsSessionToken, 1024);
|
||||
}
|
||||
else if (std::getenv("AWS_ACCESS_KEY_ID") &&
|
||||
std::getenv("AWS_SECRET_ACCESS_KEY") &&
|
||||
|
||||
@@ -28,6 +28,7 @@ struct cap_cb {
|
||||
char sessionId[MAX_SESSION_ID+1];
|
||||
char awsAccessKeyId[128];
|
||||
char awsSecretAccessKey[128];
|
||||
char awsSessionToken[1024];
|
||||
uint32_t channels;
|
||||
SpeexResamplerState *resampler;
|
||||
void* streamer;
|
||||
|
||||
@@ -304,7 +304,7 @@ public:
|
||||
|
||||
void finish() {
|
||||
if (m_finished) return;
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "GStreamer::finish - calling StopContinuousRecognitionAsync (%p)\n", this);
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_INFO, "GStreamer::finish - calling StopContinuousRecognitionAsync (%p)\n", this);
|
||||
m_finished = true;
|
||||
m_recognizer->StopContinuousRecognitionAsync().get();
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_INFO, "GStreamer::finish - recognition has completed (%p)\n", this);
|
||||
|
||||
@@ -179,23 +179,55 @@ extern "C" {
|
||||
if (a->flushed) return;
|
||||
bool fireEvent = false;
|
||||
CircularBuffer_t *cBuffer = (CircularBuffer_t *) a->circularBuffer;
|
||||
size_t total_bytes_to_process;
|
||||
|
||||
auto audioData = e.Result->GetAudioData();
|
||||
auto bytes_received = audioData->size();
|
||||
// Buffer to hold combined data if there is unprocessed byte from the last call.
|
||||
std::unique_ptr<uint8_t[]> combinedData;
|
||||
if (a->has_last_byte) {
|
||||
a->has_last_byte = false; // We'll handle the last_byte now, so toggle the flag off
|
||||
|
||||
// Allocate memory for the new data array
|
||||
combinedData.reset(new uint8_t[bytes_received + 1]);
|
||||
|
||||
// Prepend the last byte from previous call
|
||||
combinedData[0] = a->last_byte;
|
||||
|
||||
// Copy the new data following the prepended byte
|
||||
memcpy(combinedData.get() + 1, audioData->data(), bytes_received);
|
||||
|
||||
total_bytes_to_process = bytes_received + 1;
|
||||
} else {
|
||||
// Allocate memory for the new data array
|
||||
combinedData.reset(new uint8_t[bytes_received]);
|
||||
memcpy(combinedData.get(), audioData->data(), bytes_received);
|
||||
total_bytes_to_process = bytes_received;
|
||||
}
|
||||
|
||||
// If we now have an odd total, save the last byte for next time
|
||||
auto data = combinedData.get();
|
||||
if ((total_bytes_to_process % sizeof(int16_t)) != 0) {
|
||||
a->last_byte = data[total_bytes_to_process - 1];
|
||||
a->has_last_byte = true;
|
||||
total_bytes_to_process--;
|
||||
}
|
||||
|
||||
if (a->file) {
|
||||
fwrite(audioData->data(), 1, audioData->size(), a->file);
|
||||
fwrite(data, 1, total_bytes_to_process, a->file);
|
||||
}
|
||||
|
||||
/**
|
||||
* this sort of reinterpretation can be dangerous as a general rule, but in this case we know that the data
|
||||
* is 16-bit PCM, so it's safe to do this and its much faster than copying the data byte by byte
|
||||
*/
|
||||
const uint16_t* begin = reinterpret_cast<const uint16_t*>(audioData->data());
|
||||
const uint16_t* end = reinterpret_cast<const uint16_t*>(audioData->data() + audioData->size());
|
||||
const uint16_t* begin = reinterpret_cast<const uint16_t*>(data);
|
||||
const uint16_t* end = reinterpret_cast<const uint16_t*>(data + total_bytes_to_process);
|
||||
|
||||
/* lock as briefly as possible */
|
||||
switch_mutex_lock(a->mutex);
|
||||
if (cBuffer->capacity() - cBuffer->size() < audioData->size()) {
|
||||
cBuffer->set_capacity(cBuffer->size() + std::max( audioData->size(), (size_t)BUFFER_SIZE));
|
||||
if (cBuffer->capacity() - cBuffer->size() < total_bytes_to_process) {
|
||||
cBuffer->set_capacity(cBuffer->size() + std::max( total_bytes_to_process, (size_t)BUFFER_SIZE));
|
||||
}
|
||||
cBuffer->insert(cBuffer->end(), begin, end);
|
||||
switch_mutex_unlock(a->mutex);
|
||||
|
||||
@@ -82,6 +82,9 @@ static switch_status_t a_speech_feed_tts(switch_speech_handle_t *sh, char *text,
|
||||
a->draining = 0;
|
||||
a->reads = 0;
|
||||
a->flushed = 0;
|
||||
a->response_code = 0;
|
||||
a->err_msg = NULL;
|
||||
a->has_last_byte = 0;
|
||||
|
||||
return azure_speech_feed_tts(a, text, flags);
|
||||
}
|
||||
|
||||
@@ -32,6 +32,9 @@ typedef struct azure_data {
|
||||
SpeexResamplerState *resampler;
|
||||
void *circularBuffer;
|
||||
switch_mutex_t *mutex;
|
||||
|
||||
int has_last_byte;
|
||||
uint8_t last_byte;
|
||||
} azure_t;
|
||||
|
||||
#endif
|
||||
@@ -84,6 +84,8 @@ static switch_status_t d_speech_feed_tts(switch_speech_handle_t *sh, char *text,
|
||||
deepgram_t *d = createOrRetrievePrivateData(sh);
|
||||
d->draining = 0;
|
||||
d->reads = 0;
|
||||
d->response_code = 0;
|
||||
d->err_msg = NULL;
|
||||
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_INFO, "d_speech_feed_tts\n");
|
||||
|
||||
|
||||
@@ -106,6 +106,8 @@ static switch_status_t ell_speech_feed_tts(switch_speech_handle_t *sh, char *tex
|
||||
elevenlabs_t *el = createOrRetrievePrivateData(sh);
|
||||
el->draining = 0;
|
||||
el->reads = 0;
|
||||
el->response_code = 0;
|
||||
el->err_msg = NULL;
|
||||
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_INFO, "ell_speech_feed_tts\n");
|
||||
|
||||
|
||||
@@ -404,12 +404,14 @@ extern "C" {
|
||||
jambonz::AudioPipe *pAudioPipe = static_cast<jambonz::AudioPipe *>(tech_pvt->pAudioPipe);
|
||||
if (pAudioPipe) reaper(tech_pvt);
|
||||
destroy_tech_pvt(tech_pvt);
|
||||
switch_mutex_unlock(tech_pvt->mutex);
|
||||
switch_log_printf(SWITCH_CHANNEL_SESSION_LOG(session), SWITCH_LOG_DEBUG, "(%u) jb_transcribe_session_stop, bug removed\n", id);
|
||||
return SWITCH_STATUS_SUCCESS;
|
||||
}
|
||||
else {
|
||||
switch_log_printf(SWITCH_CHANNEL_SESSION_LOG(session), SWITCH_LOG_DEBUG, "jb_transcribe_session_stop: race condition, previous close completed\n");
|
||||
}
|
||||
switch_mutex_unlock(tech_pvt->mutex);
|
||||
}
|
||||
return SWITCH_STATUS_FALSE;
|
||||
}
|
||||
|
||||
@@ -100,7 +100,7 @@ static switch_status_t start_capture(switch_core_session_t *session, switch_medi
|
||||
return status;
|
||||
}
|
||||
switch_channel_set_private(channel, bugname, bug);
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "added media bug for jb transcribe\n");
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "added media bug for jb transcribe: %s\n", bugname);
|
||||
|
||||
return SWITCH_STATUS_SUCCESS;
|
||||
}
|
||||
@@ -146,7 +146,7 @@ SWITCH_STANDARD_API(jb_transcribe_function)
|
||||
if ((lsession = switch_core_session_locate(argv[0]))) {
|
||||
if (!strcasecmp(argv[1], "stop")) {
|
||||
char *bugname = argc > 2 ? argv[2] : MY_BUG_NAME;
|
||||
switch_log_printf(SWITCH_CHANNEL_SESSION_LOG(session), SWITCH_LOG_INFO, "stop transcribing\n");
|
||||
switch_log_printf(SWITCH_CHANNEL_SESSION_LOG(session), SWITCH_LOG_INFO, "stop transcribing %s\n", bugname);
|
||||
status = do_stop(lsession, bugname);
|
||||
} else if (!strcasecmp(argv[1], "start")) {
|
||||
char* lang = argv[2];
|
||||
|
||||
@@ -99,6 +99,8 @@ static switch_status_t p_speech_feed_tts(switch_speech_handle_t *sh, char *text,
|
||||
playht_t *p = createOrRetrievePrivateData(sh);
|
||||
p->draining = 0;
|
||||
p->reads = 0;
|
||||
p->response_code = 0;
|
||||
p->err_msg = NULL;
|
||||
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "p_speech_feed_tts\n");
|
||||
|
||||
|
||||
@@ -76,7 +76,8 @@ static CURL* createEasyHandle(void) {
|
||||
|
||||
// set connect timeout to 3 seconds and total timeout to 109 seconds
|
||||
curl_easy_setopt(easy, CURLOPT_CONNECTTIMEOUT_MS, 3000L);
|
||||
curl_easy_setopt(easy, CURLOPT_TIMEOUT, 10L);
|
||||
//For long text, PlayHT took more than 20 seconds to complete the download.
|
||||
curl_easy_setopt(easy, CURLOPT_TIMEOUT, 60L);
|
||||
|
||||
return easy ;
|
||||
}
|
||||
|
||||
@@ -82,6 +82,8 @@ static switch_status_t d_speech_feed_tts(switch_speech_handle_t *sh, char *text,
|
||||
rimelabs_t *d = createOrRetrievePrivateData(sh);
|
||||
d->draining = 0;
|
||||
d->reads = 0;
|
||||
d->response_code = 0;
|
||||
d->err_msg = NULL;
|
||||
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_INFO, "d_speech_feed_tts\n");
|
||||
|
||||
|
||||
@@ -90,6 +90,8 @@ static switch_status_t w_speech_feed_tts(switch_speech_handle_t *sh, char *text,
|
||||
whisper_t *w = createOrRetrievePrivateData(sh);
|
||||
w->draining = 0;
|
||||
w->reads = 0;
|
||||
w->response_code = 0;
|
||||
w->err_msg = NULL;
|
||||
|
||||
switch_log_printf(SWITCH_CHANNEL_LOG, SWITCH_LOG_DEBUG, "w_speech_feed_tts\n");
|
||||
|
||||
|
||||
Reference in New Issue
Block a user