From 0bb8dd2084a8098b76a671970967126ca9f6bd26 Mon Sep 17 00:00:00 2001 From: Michael Schimmel Date: Wed, 4 Feb 2026 12:22:57 +0100 Subject: [PATCH] WASAPI: added CaptureSST.pas --- WASAPI/CaptureSTT.pas | 435 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 435 insertions(+) create mode 100644 WASAPI/CaptureSTT.pas diff --git a/WASAPI/CaptureSTT.pas b/WASAPI/CaptureSTT.pas new file mode 100644 index 0000000..c2c8343 --- /dev/null +++ b/WASAPI/CaptureSTT.pas @@ -0,0 +1,435 @@ +unit CaptureSTT; + +interface + +uses + Winapi.Windows, Winapi.Messages, System.SysUtils, System.Classes, Winapi.MMSystem, + System.SyncObjs, System.Math, System.Net.HttpClient, System.Net.HttpClientComponent, + System.Net.Mime, System.JSON, System.IOUtils, System.Generics.Collections; + +type + TTranscriptionTask = record + TaskID: string; + SequenceIndex: Int64; + end; + + TTranscriptionResultEvent = reference to procedure(const AText: string); + TLogEvent = reference to procedure(const AMsg: string); + + (* Polling worker for async API tasks *) + TTranscriptionWorker = class(TThread) + strict private + fQueue: TThreadedQueue; + fClient: TNetHTTPClient; + fOnResult: TTranscriptionResultEvent; + procedure ProcessTask(const ATask: TTranscriptionTask); + protected + procedure Execute; override; + public + constructor Create(AOnResult: TTranscriptionResultEvent); + destructor Destroy; override; + procedure Enqueue(const ATaskID: string; AIndex: Int64); + end; + + (* Audio capture with silence detection and VAD *) + TAudioCaptureThread = class(TThread) + strict private + fWaveIn: HWAVEIN; + fHeaders: array[0..1] of TWaveHdr; + fBuffers: array[0..1] of array of Byte; + fEvent: THandle; + fCurrentRms: Integer; + fIsSilent: Boolean; + fDbThreshold: Double; + fSilenceCounter: Integer; + fRequiredSilencePackets: Integer; + fActiveFragment: TMemoryStream; + fTotalSegmentsSent: Int64; + fForceUploadFlag: Integer; + fWorker: TTranscriptionWorker; + fOnLog: TLogEvent; + fUrlBase: string; + fRequestParams: TStrings; + + procedure PrepareBuffers; + procedure UnprepareBuffers; + function CalculateRms(ABuffer: PByte; ALength: Cardinal): Integer; + procedure ProcessSilence(ARms: Integer; ABuffer: PByte; ALength: Cardinal); + procedure WriteWavHeader(AStream: TStream; ADataSize: Cardinal); + procedure TrimSilentStart(AOffset: Int64); + procedure SyncUploadAndEnqueue(AStream: TMemoryStream); + procedure PerformUpload(AStream: TMemoryStream); + protected + procedure Execute; override; + public + constructor Create(AWorker: TTranscriptionWorker; const AUrl: string; AParams: TStrings; AOnLog: TLogEvent); + destructor Destroy; override; + procedure ForceUpload; + property CurrentRms: Integer read fCurrentRms; + property DbThreshold: Double read fDbThreshold write fDbThreshold; + property RequiredSilencePackets: Integer read fRequiredSilencePackets write fRequiredSilencePackets; + end; + +implementation + +{ TTranscriptionWorker } + +constructor TTranscriptionWorker.Create(AOnResult: TTranscriptionResultEvent); +begin + inherited Create(False); + fOnResult := AOnResult; + fQueue := TThreadedQueue.Create(100, 1000, 100); + fClient := TNetHTTPClient.Create(nil); +end; + +destructor TTranscriptionWorker.Destroy; +begin + Terminate; + fQueue.DoShutDown; + fClient.Free; + fQueue.Free; + inherited; +end; + +procedure TTranscriptionWorker.Enqueue(const ATaskID: string; AIndex: Int64); +var + task: TTranscriptionTask; +begin + task.TaskID := ATaskID; + task.SequenceIndex := AIndex; + fQueue.PushItem(task); +end; + +procedure TTranscriptionWorker.Execute; +var + task: TTranscriptionTask; +begin + while not Terminated do + begin + if (fQueue.PopItem(task) = TWaitResult.wrSignaled) and (not Terminated) then + ProcessTask(task); + end; +end; + +procedure TTranscriptionWorker.ProcessTask(const ATask: TTranscriptionTask); +var + resp: IHTTPResponse; + jsonResult: TJSONObject; + status, transcriptionText, sText: string; + resVal, segmentsVal: TJSONValue; + segmentsArray: TJSONArray; + i: Integer; + isDone: Boolean; +begin + isDone := False; + repeat + if Terminated then exit; + try + resp := fClient.Get('http://minerva.lan:8000/task/' + ATask.TaskID); + if (resp.StatusCode = 200) then + begin + jsonResult := TJSONObject.ParseJSONValue(resp.ContentAsString(TEncoding.UTF8)) as TJSONObject; + if Assigned(jsonResult) then + try + if jsonResult.TryGetValue('status', status) then + begin + if (status = 'completed') then + begin + if jsonResult.TryGetValue('result', resVal) and (resVal is TJSONObject) and + TJSONObject(resVal).TryGetValue('segments', segmentsVal) then + begin + segmentsArray := segmentsVal as TJSONArray; + transcriptionText := ''; + for i := 0 to segmentsArray.Count - 1 do + if (segmentsArray.Items[i] as TJSONObject).TryGetValue('text', sText) then + transcriptionText := transcriptionText + sText; + + sText := transcriptionText.Trim; + if Assigned(fOnResult) then + TThread.Queue(nil, procedure begin fOnResult(sText); end); + isDone := True; + end; + end + else if (status = 'failed') or (status = 'error') then + isDone := True; + end; + finally + jsonResult.Free; + end; + end; + except + isDone := True; + end; + if not isDone then Sleep(500); + until isDone; +end; + +{ TAudioCaptureThread } + +constructor TAudioCaptureThread.Create(AWorker: TTranscriptionWorker; const AUrl: string; AParams: TStrings; AOnLog: TLogEvent); +var + format: TWaveFormatEx; +begin + inherited Create(True); + fWorker := AWorker; + fUrlBase := AUrl; + fOnLog := AOnLog; + fRequestParams := TStringList.Create; + if Assigned(AParams) then + fRequestParams.Assign(AParams); + + fEvent := CreateEvent(nil, False, False, nil); + fActiveFragment := TMemoryStream.Create; + fDbThreshold := -35.0; + fRequiredSilencePackets := 30; + fIsSilent := True; + + FillChar(format, sizeof(format), 0); + format.wFormatTag := WAVE_FORMAT_PCM; + format.nChannels := 1; + format.nSamplesPerSec := 16000; + format.wBitsPerSample := 16; + format.nBlockAlign := 2; + format.nAvgBytesPerSec := 32000; + + if waveInOpen(@fWaveIn, WAVE_MAPPER, @format, fEvent, 0, CALLBACK_EVENT) <> MMSYSERR_NOERROR then + raise Exception.Create('waveInOpen failed'); + + PrepareBuffers; +end; + +destructor TAudioCaptureThread.Destroy; +begin + waveInStop(fWaveIn); + UnprepareBuffers; + waveInClose(fWaveIn); + CloseHandle(fEvent); + fActiveFragment.Free; + fRequestParams.Free; + inherited; +end; + +procedure TAudioCaptureThread.ForceUpload; +begin + TInterlocked.Exchange(fForceUploadFlag, 1); +end; + +procedure TAudioCaptureThread.SyncUploadAndEnqueue(AStream: TMemoryStream); +var + client: TNetHTTPClient; + formData: TMultipartFormData; + resp: IHTTPResponse; + jsonResp: TJSONObject; + taskID, urlParams, fullUrl: string; + i: Integer; +begin + client := TNetHTTPClient.Create(nil); + formData := TMultipartFormData.Create; + try + AStream.Position := 0; + formData.AddStream('file', AStream, False, 'fragment.wav', 'audio/wav'); + + urlParams := ''; + for i := 0 to fRequestParams.Count - 1 do + begin + if (urlParams <> '') then urlParams := urlParams + '&'; + urlParams := urlParams + fRequestParams[i]; + end; + + fullUrl := fUrlBase; + if urlParams <> '' then + fullUrl := fullUrl + '?' + urlParams; + + resp := client.Post(fullUrl, formData); + + if (resp.StatusCode = 200) then + begin + jsonResp := TJSONObject.ParseJSONValue(resp.ContentAsString(TEncoding.UTF8)) as TJSONObject; + try + if Assigned(jsonResp) and jsonResp.TryGetValue('identifier', taskID) then + fWorker.Enqueue(taskID, fTotalSegmentsSent); + finally + jsonResp.Free; + end; + end; + finally + formData.Free; + client.Free; + AStream.Free; + end; +end; + +procedure TAudioCaptureThread.PerformUpload(AStream: TMemoryStream); +begin + inc(fTotalSegmentsSent); + SyncUploadAndEnqueue(AStream); +end; + +procedure TAudioCaptureThread.WriteWavHeader(AStream: TStream; ADataSize: Cardinal); +type + TWavHeader = packed record + RIFF: array[0..3] of AnsiChar; + FileSize: Cardinal; + WAVE: array[0..3] of AnsiChar; + fmt: array[0..3] of AnsiChar; + FormatSize: Cardinal; + FormatTag: Word; + Channels: Word; + SamplesPerSec: Cardinal; + AvgBytesPerSec: Cardinal; + BlockAlign: Word; + BitsPerSample: Word; + DataMark: array[0..3] of AnsiChar; + DataSize: Cardinal; + end; +var + header: TWavHeader; +begin + header.RIFF := 'RIFF'; + header.FileSize := ADataSize + sizeof(TWavHeader) - 8; + header.WAVE := 'WAVE'; + header.fmt := 'fmt '; + header.FormatSize := 16; + header.FormatTag := 1; + header.Channels := 1; + header.SamplesPerSec := 16000; + header.BitsPerSample := 16; + header.BlockAlign := 2; + header.AvgBytesPerSec := 32000; + header.DataMark := 'data'; + header.DataSize := ADataSize; + AStream.WriteBuffer(header, sizeof(TWavHeader)); +end; + +function TAudioCaptureThread.CalculateRms(ABuffer: PByte; ALength: Cardinal): Integer; +var + i, count: Integer; + pSamples: PSmallInt; + sum, sampleVal: Double; +begin + sum := 0; + pSamples := PSmallInt(ABuffer); + count := ALength div 2; + for i := 0 to count - 1 do + begin + sampleVal := pSamples^; + sum := sum + (sampleVal * sampleVal); + inc(pSamples); + end; + if (count > 0) then Result := Round(Sqrt(sum / count)) else Result := 0; +end; + +procedure TAudioCaptureThread.TrimSilentStart(AOffset: Int64); +var + newSize: Int64; +begin + newSize := fActiveFragment.Size - AOffset; + if (newSize > 0) then + Move(PByte(fActiveFragment.Memory)[AOffset], fActiveFragment.Memory^, newSize); + fActiveFragment.Size := newSize; + fActiveFragment.Position := newSize; +end; + +procedure TAudioCaptureThread.ProcessSilence(ARms: Integer; ABuffer: PByte; ALength: Cardinal); +var + uploadStream: TMemoryStream; + currentDb: Double; + maxIdleSize, keepSize, sendSize: Int64; + forceRequested: Boolean; +begin + forceRequested := TInterlocked.CompareExchange(fForceUploadFlag, 0, 1) = 1; + fActiveFragment.WriteBuffer(ABuffer^, ALength); + if (ARms > 0) then currentDb := 20 * Log10(ARms / 32768) else currentDb := -100.0; + + if (currentDb > fDbThreshold) then + begin + if fIsSilent then + begin + fIsSilent := False; + if Assigned(fOnLog) then fOnLog('Activity detected...'); + end; + fSilenceCounter := 0; + end + else + begin + inc(fSilenceCounter); + if fIsSilent then + begin + maxIdleSize := Int64(fRequiredSilencePackets) * ALength; + if (fActiveFragment.Size > maxIdleSize) then + TrimSilentStart(fActiveFragment.Size - maxIdleSize); + end + else if (fSilenceCounter >= fRequiredSilencePackets) or forceRequested then + begin + fIsSilent := True; + keepSize := (Int64(fRequiredSilencePackets) div 2) * ALength; + if forceRequested then keepSize := 0; + + sendSize := fActiveFragment.Size - keepSize; + if (sendSize > 0) then + begin + uploadStream := TMemoryStream.Create; + WriteWavHeader(uploadStream, sendSize); + fActiveFragment.Position := 0; + uploadStream.CopyFrom(fActiveFragment, sendSize); + PerformUpload(uploadStream); + TrimSilentStart(sendSize); + end; + fSilenceCounter := 0; + if Assigned(fOnLog) then fOnLog('Fragment uploaded.'); + end; + end; +end; + +procedure TAudioCaptureThread.PrepareBuffers; +var + i: Integer; + bufferSize: Cardinal; +begin + bufferSize := 3200; + for i := 0 to High(fBuffers) do + begin + SetLength(fBuffers[i], bufferSize); + fHeaders[i].lpData := PAnsiChar(@fBuffers[i][0]); + fHeaders[i].dwBufferLength := bufferSize; + fHeaders[i].dwFlags := 0; + waveInPrepareHeader(fWaveIn, @fHeaders[i], sizeof(TWaveHdr)); + waveInAddBuffer(fWaveIn, @fHeaders[i], sizeof(TWaveHdr)); + end; +end; + +procedure TAudioCaptureThread.UnprepareBuffers; +var + i: Integer; +begin + waveInReset(fWaveIn); + for i := 0 to High(fHeaders) do + waveInUnprepareHeader(fWaveIn, @fHeaders[i], sizeof(TWaveHdr)); +end; + +procedure TAudioCaptureThread.Execute; +var + i: Integer; + rmsValue: Integer; +begin + waveInStart(fWaveIn); + while (not Terminated) do + begin + if (WaitForSingleObject(fEvent, 100) = WAIT_OBJECT_0) then + begin + for i := 0 to High(fHeaders) do + begin + if (not Terminated) and ((fHeaders[i].dwFlags and WHDR_DONE) <> 0) then + begin + rmsValue := CalculateRms(PByte(fHeaders[i].lpData), fHeaders[i].dwBytesRecorded); + TInterlocked.Exchange(fCurrentRms, rmsValue); + ProcessSilence(rmsValue, PByte(fHeaders[i].lpData), fHeaders[i].dwBytesRecorded); + fHeaders[i].dwFlags := fHeaders[i].dwFlags and not WHDR_DONE; + waveInAddBuffer(fWaveIn, @fHeaders[i], sizeof(TWaveHdr)); + end; + end; + end; + end; +end; + +end.