49 lines
1.1 KiB
JavaScript
49 lines
1.1 KiB
JavaScript
const config = require('./config');
|
|
|
|
const Redis = require('redis');
|
|
const redis = {};
|
|
|
|
redis.client = Redis.createClient({
|
|
url: `redis://default@${config.REDIS_HOST}:${config.REDIS_PORT}`
|
|
});
|
|
|
|
redis.connect = async () => {
|
|
await redis.client.connect().catch(console.error);
|
|
}
|
|
|
|
redis.quit = async () => {
|
|
await redis.client.quit();
|
|
}
|
|
|
|
redis.ensureGroup = async (stream, group) => {
|
|
try {
|
|
await redis.client.xGroupCreate(stream, group, "0", {MKSTREAM: true});
|
|
|
|
console.log(`Created consumer group: ${group} (${stream})`);
|
|
} catch (err) {
|
|
if (err) {
|
|
if (err.message.includes('BUSYGROUP')) {
|
|
console.log(`Consumer group already exists: ${group} (${stream})`);
|
|
} else {
|
|
console.error(err);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
redis.sendTranscriptEvent = async (streamKey, device_id, text) => {
|
|
try {
|
|
await redis.client.xAdd(streamKey, "*", {
|
|
text,
|
|
device_id: device_id,
|
|
});
|
|
|
|
console.log("Transcript added to stream");
|
|
} catch (err) {
|
|
console.error(err);
|
|
}
|
|
}
|
|
|
|
module.exports = {
|
|
redis
|
|
} |