This commit is contained in:
shadowcz007
2024-06-02 19:53:07 +08:00
parent 77bfb08d76
commit d549a5eb6a
3 changed files with 307 additions and 43 deletions
+12 -2
View File
@@ -1,4 +1,3 @@
#
import os
import subprocess
import importlib.util
@@ -17,7 +16,7 @@ python = sys.executable
try:
sys.stdout.isatty()
except:
print('#fix sys.stdout.isatty')
# print('#fix sys.stdout.isatty')
sys.stdout.isatty = lambda: False
llama_port=None
@@ -580,6 +579,17 @@ async def mixlab_app_handler(request):
async def mixlab_live_handler(request):
html_file = os.path.join(current_path, "web/live.html")
if os.path.exists(html_file):
live_server_path=os.path.join(current_path, "nodes/vad-websockets-perclient.py")
import threading
import subprocess
def run_vad_script():
subprocess.run([python, live_server_path])
thread = threading.Thread(target=run_vad_script)
thread.start()
with open(html_file, 'r', encoding='utf-8', errors='ignore') as f:
html_data = f.read()
return web.Response(text=html_data, content_type='text/html')
+91
View File
@@ -0,0 +1,91 @@
import asyncio
import websockets
import webrtcvad
import sys
from faster_whisper import WhisperModel
import numpy as np
import io
# pip install webrtcvad-wheels
# pip install websockets
# pip install faster-whisper
vad = webrtcvad.Vad(3)
class AudioStream:
def __init__(self) -> None:
self.sample_rate = 16000
self.frame_size = 320
self.bytes_per_sample = 2
self.idle_cut = (self.sample_rate/2)/self.frame_size # chunk audio if no voice for 0.5 seconds
self.last_voice_activity = {}
self.audio_buffers = {}
def convert_buffer_size(self, audio_frame):
while len(audio_frame) < (self.frame_size * self.bytes_per_sample):
audio_frame = audio_frame + b'\x00'
return audio_frame
def manage_client_idle(self, client_id):
if client_id not in self.last_voice_activity:
self.last_voice_activity[client_id] = 0
return self.last_voice_activity[client_id]
def voice_activity_detection(self, audio_frame, client_id):
idle_time = self.manage_client_idle(client_id)
converted_frame = self.convert_buffer_size(audio_frame)
is_speech = vad.is_speech(converted_frame, sample_rate=self.sample_rate)
if is_speech:
self.last_voice_activity[client_id] = 0
if client_id not in self.audio_buffers:
self.audio_buffers[client_id] = []
self.audio_buffers[client_id].append(converted_frame)
return "1"
else:
if idle_time == self.idle_cut:
self.last_voice_activity[client_id] = 0
return "X"
else:
self.last_voice_activity[client_id] += 1
return "_"
audiostream = AudioStream()
model_size = "faster-distil-medium.en"
whisper_model = WhisperModel(model_size, device="cuda", compute_type="float16")
async def handler(websocket, path):
client_id = id(websocket)
print(f"WebSocket connection established for client {client_id} from {path}")
try:
async for message in websocket:
is_active = audiostream.voice_activity_detection(message, client_id)
if is_active == "X": # Voice activity stopped
audio_buffer = audiostream.audio_buffers.pop(client_id, None)
if audio_buffer:
audio_data = b''.join(audio_buffer)
audio_data = np.frombuffer(audio_data, dtype=np.int16)
with io.BytesIO() as wav_buffer:
wav_buffer.write(audio_data.tobytes())
wav_buffer.seek(0)
segments, info = whisper_model.transcribe(wav_buffer, beam_size=5, language="en", condition_on_previous_text=False)
transcript = " ".join(segment.text for segment in segments)
await websocket.send(transcript)
else:
await websocket.send(is_active)
except websockets.exceptions.ConnectionClosed:
print(f"WebSocket connection closed for client {client_id}")
except Exception as e:
print(f"Error occurred: {e}")
finally:
await websocket.close()
async def main():
PORT = 5000
async with websockets.serve(handler, 'localhost', PORT):
print(f"WebSocket server started at ws://localhost:{PORT}")
await asyncio.Future()
asyncio.run(main())
+204 -41
View File
@@ -1,46 +1,209 @@
<!-- https://stackoverflow.com/questions/67118642/audiocontext-getusermedia-and-websockets-audio-streaming -->
<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>Live</title>
<script src="https://cdn.jsdelivr.net/npm/livekit-client/dist/livekit-client.umd.min.js"></script>
</head>
<html>
<body>
<script>
const {
Participant,
RemoteParticipant,
RemoteTrack,
RemoteTrackPublication,
Room,
RoomEvent,
VideoPresets
} = LivekitClient
// creates a new room with options
const room = new Room({
// automatically manage subscribed video quality
adaptiveStream: true,
// optimize publishing bandwidth and CPU for published tracks
dynacast: true,
// default capture settings
videoCaptureDefaults: {
resolution: VideoPresets.h720.resolution,
},
});
// pre-warm connection, this can be called as early as your page is loaded
room.prepareConnection(url, token);
</script>
<div class='message'>Welcome!</div>
<button onclick='startRecording()'>Start recording</button>
<button onclick='stopRecording()'>Stop recording</button>
<br />
<div>WebSocket: <span id="webSocketStatus">Not Connected</span></div>
</body>
</html>
</html>
<script>
//================= CONFIG =================
// Global Variables
// let websocket_uri = 'wss://edca-36-68-8-204.ap.ngrok.io';
let websocket_uri = 'ws://127.0.0.1:5000';
let bufferSize = 512,
AudioContext,
context,
processor,
input,
globalStream,
websocket;
// Initialize WebSocket
initWebSocket();
function downsampleBuffer(buffer, sampleRate, outSampleRate) {
if (outSampleRate == sampleRate) {
return buffer;
}
if (outSampleRate > sampleRate) {
throw 'downsampling rate show be smaller than original sample rate';
}
var sampleRateRatio = sampleRate / outSampleRate;
var newLength = Math.round(buffer.length / sampleRateRatio);
var result = new Int16Array(newLength);
var offsetResult = 0;
var offsetBuffer = 0;
while (offsetResult < result.length) {
var nextOffsetBuffer = Math.round((offsetResult + 1) * sampleRateRatio);
var accum = 0,
count = 0;
for (var i = offsetBuffer; i < nextOffsetBuffer && i < buffer.length; i++) {
accum += buffer[i];
count++;
}
result[offsetResult] = Math.min(1, accum / count) * 0x7fff;
offsetResult++;
offsetBuffer = nextOffsetBuffer;
}
return result.buffer;
} // closes function downsampleBuffer()
//================= RECORDING =================
// Define the AudioWorkletProcessor in a Blob and add it to the AudioWorklet.
const processorCode = `
class MyProcessor extends AudioWorkletProcessor {
process(inputs, outputs, parameters) {
const input = inputs[0];
const output = outputs[0];
if (input && input[0]) {
const left = input[0];
const left16 = this.downsampleBuffer(left, sampleRate, 16000);
this.port.postMessage(left16);
}
return true;
}
downsampleBuffer(buffer, sampleRate, outSampleRate) {
if (outSampleRate === sampleRate) {
return buffer;
}
const sampleRateRatio = sampleRate / outSampleRate;
const newLength = Math.round(buffer.length / sampleRateRatio);
const result = new Int16Array(newLength);
let offsetResult = 0;
let offsetBuffer = 0;
while (offsetResult < result.length) {
const nextOffsetBuffer = Math.round((offsetResult + 1) * sampleRateRatio);
let accum = 0, count = 0;
for (let i = offsetBuffer; i < nextOffsetBuffer && i < buffer.length; i++) {
accum += buffer[i];
count++;
}
result[offsetResult] = Math.min(1, accum / count) * 0x7FFF;
offsetResult++;
offsetBuffer = nextOffsetBuffer;
}
return result;
}
}
registerProcessor('my-processor', MyProcessor);
`;
const blob = new Blob([processorCode], { type: 'application/javascript' });
const url = URL.createObjectURL(blob);
async function startRecording() {
streamStreaming = true;
AudioContext = window.AudioContext || window.webkitAudioContext;
context = new AudioContext({ latencyHint: 'interactive' });
await context.audioWorklet.addModule(url);
const processor = new AudioWorkletNode(context, 'my-processor');
processor.connect(context.destination);
context.resume();
processor.port.onmessage = (event) => {
websocket.send(event.data);
};
const handleSuccess = function (stream) {
globalStream = stream;
const input = context.createMediaStreamSource(stream);
input.connect(processor);
};
navigator.mediaDevices.getUserMedia({ audio: true, video: false }).then(handleSuccess);
}
function stopRecording() {
streamStreaming = false;
if (globalStream) {
let track = globalStream.getTracks()[0];
track.stop();
}
if (input) {
input.disconnect();
}
if (processor) {
processor.disconnect();
}
if (context) {
context.close().then(function () {
input = null;
processor = null;
context = null;
AudioContext = null;
});
}
}
let lastActivityTime = Date.now();
const idleCut = 5000; // Replace with the appropriate idle time in milliseconds
let isSpeaking = false;
function handleVoiceActivity(result) {
if (result === "1" && isSpeaking == false) {
lastActivityTime = Date.now();
console.log("Speaking");
isSpeaking = true;
} else if (result === "X") {
// console.log("Stopped speaking for too long, sending signal");
// send();
} else if (result === "_") {
let currentTime = Date.now();
if (currentTime - lastActivityTime > idleCut && isSpeaking) {
console.log("Stopped speaking for too long, sending signal");
isSpeaking = false;
send();
}
}
}
function send() {
// Implement the send logic here, e.g., making an HTTP request
console.log("Signal sent");
}
function initWebSocket() {
// Create WebSocket
websocket = new WebSocket(websocket_uri);
//console.log("Websocket created...");
// WebSocket Definitions: executed when triggered webSocketStatus
websocket.onopen = function () {
console.log("connected to server");
//websocket.send("CONNECTED TO YOU");
document.getElementById("webSocketStatus").innerHTML = 'Connected';
}
websocket.onclose = function (e) {
console.log("connection closed (" + e.code + ")");
document.getElementById("webSocketStatus").innerHTML = 'Not Connected';
}
websocket.onmessage = function (e) {
// console.log("message received: " + e.data);
// console.log(e.data);
let result = e.data;
document.querySelector('.message').innerHTML = result;
handleVoiceActivity(result)
}
} // closes function initWebSocket()
</script>