import { getDb } from "../index"; import { ClientAudioBufferDoc, ClientAudioBufferType, CreateAudioBufferInput, } from "./clientsAudioBuffers.types"; import { Collection, InsertOneResult, DeleteResult } from "mongodb"; const getClientAudioBuffersCollection = (): Collection => { return getDb().collection("clientAudioBuffers"); }; /** * Deletes all audio buffers associated with a specific callSid. * Useful for clearing context between different tickets in the same call. */ const deleteAudioBuffersByCallSid = async ( callSid: string, ): Promise => { return getClientAudioBuffersCollection().deleteMany({ callSid }); }; const insertClientAudioBuffer = async ( audioBufferInput: CreateAudioBufferInput, ): Promise> => { const collection = getClientAudioBuffersCollection(); const realtimeDate = new Date(); const ttlMs = audioBufferInput.ttlMs ?? 3600_000; const expiresAt = new Date(realtimeDate.getTime() + ttlMs); const inputFieldsForDocument = { callSid: audioBufferInput.callSid, clientId: audioBufferInput.clientId, bucketStartMs: audioBufferInput.bucketStartMs, bucketDurationMs: audioBufferInput.bucketDurationMs, data: audioBufferInput.data, }; const clientAudioBufferMongoDocument: ClientAudioBufferType = { ...inputFieldsForDocument, codec: "mulaw", sampleRateHz: 8000, channels: 1, createdAt: realtimeDate, expiresAt: expiresAt, }; return collection.insertOne(clientAudioBufferMongoDocument); }; const findClientAudioBuffersByCallSidAndRange = async ( callSid: string, startMs?: number, endMs?: number, ): Promise => { const collection = getClientAudioBuffersCollection(); if (startMs === undefined) { return collection.find({ callSid }).toArray(); } const endMsToUse = endMs ?? Date.now(); return collection .find({ callSid, bucketStartMs: { $gte: startMs, $lte: endMsToUse, }, }) .toArray(); }; const getAudioChunksByCallSid = async ( callSid: string, limit?: number, ): Promise => { const query = getClientAudioBuffersCollection().find({ callSid }); if (limit) { query.sort({ bucketStartMs: -1 }).limit(limit); } else { query.sort({ bucketStartMs: 1 }); } const docs = await query.toArray(); if (!docs.length) { return []; } const orderedDocs = limit ? docs.reverse() : docs; return orderedDocs.map((chunk) => { const data: any = chunk.data; if (Buffer.isBuffer(data)) return data; if (data?.buffer) return Buffer.from(data.buffer); return Buffer.from(data); }); }; const getMergedAudioByCallSid = async ( callSid: string, limit: number, ): Promise => { const buffers = await getAudioChunksByCallSid(callSid, limit); if (!buffers.length) { return Buffer.alloc(0); } return Buffer.concat(buffers); }; export { insertClientAudioBuffer, getClientAudioBuffersCollection, findClientAudioBuffersByCallSidAndRange, getAudioChunksByCallSid, getMergedAudioByCallSid, deleteAudioBuffersByCallSid, };