Skip to content

Commit b368bf8

Browse files
committed
feat(audio): 原生模块支持流式采集(采集线程 + ThreadSafeFunction)
Phase 2 前置:同步阻塞采集无法用于真实链路,新增生产采集入口。 - capture_session:WASAPI 激活/初始化/读包原语(pimpl 隐藏 COM),供同步与流式 两路复用;含 isProcessLoopbackSupported 能力探测(真实尝试激活而非查版本号, 避免兼容性设置伪造版本号导致误判) - streaming_capture:采集线程按 chunkDurationMs 聚块回调,Stop 阻塞等待线程退出 以保证无回调泄漏 - addon:新增 startCapture/stopCapture/isProcessLoopbackSupported,块经 ThreadSafeFunction 投递到 JS 线程并拷贝为 ArrayBuffer - loopback_capture 精简为复用 CaptureSession(去重 WASAPI 样板代码) 实测 spike-streaming:5 个 1s 块周期准确、样本数 16000 一致、元数据正确、 stopCapture 后零泄漏。
1 parent d41773d commit b368bf8

8 files changed

Lines changed: 753 additions & 207 deletions

File tree

‎client/native/process-audio/binding.gyp‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,9 @@
55
"sources": [
66
"src/addon.cc",
77
"src/window_finder.cc",
8-
"src/loopback_capture.cc"
8+
"src/capture_session.cc",
9+
"src/loopback_capture.cc",
10+
"src/streaming_capture.cc"
911
],
1012
"include_dirs": [
1113
"<!@(node -p \"require('node-addon-api').include\")"

‎client/native/process-audio/src/addon.cc‎

Lines changed: 136 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,19 @@
11
// N-API 绑定入口
22
//
3-
// @ai-context: Phase 1 spike 第一步仅暴露窗口/进程解析,用于验证编译链路
4-
// 与"窗口 → PID → 应用根进程"回溯是否可靠;WASAPI 进程环回采集在验证
5-
// 通过后追加。
3+
// @ai-context: 导出三类能力——窗口/进程解析、能力探测、采集。
4+
// 采集分两个入口:captureToWav(同步阻塞,spike 验证用)与
5+
// startCapture/stopCapture(采集线程 + ThreadSafeFunction 流式回调,生产用)。
6+
// @ai-context: 单实例语义——IPC 层同时只有一路采集,重复 startCapture 报错
7+
// 而非默默覆盖,避免采集线程泄漏。
68

79
#include <napi.h>
810

11+
#include <memory>
12+
#include <utility>
13+
14+
#include "capture_session.h"
915
#include "loopback_capture.h"
16+
#include "streaming_capture.h"
1017
#include "window_finder.h"
1118

1219
namespace {
@@ -104,10 +111,136 @@ Napi::Value CaptureToWav(const Napi::CallbackInfo& info) {
104111
return out;
105112
}
106113

114+
/** isProcessLoopbackSupported(): 能力探测(真实尝试激活一次) */
115+
Napi::Value IsSupported(const Napi::CallbackInfo& info) {
116+
return Napi::Boolean::New(info.Env(), process_audio::IsProcessLoopbackSupported());
117+
}
118+
119+
// ================================================================
120+
// 流式采集(生产入口)
121+
// ================================================================
122+
123+
/** 投递给 JS 线程的负载:一个音频块或一次错误 */
124+
struct StreamEvent {
125+
bool is_error = false;
126+
std::string error;
127+
process_audio::StreamingChunk chunk;
128+
};
129+
130+
/** 全局单例采集器与其线程安全回调句柄 */
131+
std::unique_ptr<process_audio::StreamingCapture> g_capture;
132+
Napi::ThreadSafeFunction g_tsfn;
133+
134+
/** 在 JS 线程消费采集线程投递的事件 */
135+
void DispatchStreamEvent(Napi::Env env, Napi::Function callback, StreamEvent* event) {
136+
if (env != nullptr && callback != nullptr) {
137+
Napi::Object payload = Napi::Object::New(env);
138+
if (event->is_error) {
139+
payload.Set("error", Napi::String::New(env, event->error));
140+
} else {
141+
const size_t bytes = event->chunk.samples.size() * sizeof(float);
142+
// 拷贝一份给 JS:采集线程的 vector 回调后即释放,不能共享内存
143+
Napi::ArrayBuffer buffer = Napi::ArrayBuffer::New(env, bytes);
144+
if (bytes > 0) {
145+
std::memcpy(buffer.Data(), event->chunk.samples.data(), bytes);
146+
}
147+
payload.Set("audioBuffer", buffer);
148+
payload.Set("sampleRate", Napi::Number::New(env, event->chunk.sample_rate));
149+
payload.Set("channels", Napi::Number::New(env, event->chunk.channels));
150+
payload.Set("durationMs", Napi::Number::New(env, event->chunk.duration_ms));
151+
}
152+
callback.Call({payload});
153+
}
154+
delete event;
155+
}
156+
157+
/**
158+
* startCapture({ pid, sampleRate, channels, chunkDurationMs }, cb):
159+
* 启动采集线程,每聚成一块就回调 cb({audioBuffer, sampleRate, channels,
160+
* durationMs});致命错误回调 cb({error})。
161+
*
162+
* 会话打开在采集线程内完成,故"目标不可采"这类失败以 error 回调形式
163+
* 异步上报,调用方应据此触发降级。
164+
*/
165+
Napi::Value StartCapture(const Napi::CallbackInfo& info) {
166+
Napi::Env env = info.Env();
167+
if (info.Length() < 2 || !info[0].IsObject() || !info[1].IsFunction()) {
168+
Napi::TypeError::New(env, "startCapture(options, callback) 参数不完整")
169+
.ThrowAsJavaScriptException();
170+
return env.Undefined();
171+
}
172+
if (g_capture && g_capture->running()) {
173+
Napi::Error::New(env, "采集已在进行中,请先 stopCapture")
174+
.ThrowAsJavaScriptException();
175+
return env.Undefined();
176+
}
177+
178+
Napi::Object opts = info[0].As<Napi::Object>();
179+
process_audio::StreamingOptions options;
180+
options.root_pid = opts.Has("pid")
181+
? opts.Get("pid").As<Napi::Number>().Uint32Value() : 0;
182+
options.sample_rate = opts.Has("sampleRate")
183+
? opts.Get("sampleRate").As<Napi::Number>().Uint32Value() : 16000;
184+
options.channels = opts.Has("channels")
185+
? opts.Get("channels").As<Napi::Number>().Uint32Value() : 1;
186+
options.chunk_duration_ms = opts.Has("chunkDurationMs")
187+
? opts.Get("chunkDurationMs").As<Napi::Number>().Uint32Value() : 5000;
188+
189+
g_tsfn = Napi::ThreadSafeFunction::New(
190+
env, info[1].As<Napi::Function>(), "process-audio-capture", 0, 1);
191+
192+
g_capture = std::make_unique<process_audio::StreamingCapture>();
193+
const std::string err = g_capture->Start(
194+
options,
195+
[](process_audio::StreamingChunk&& chunk) {
196+
auto* event = new StreamEvent();
197+
event->chunk = std::move(chunk);
198+
// 采集线程不能直接碰 JS,统一经 ThreadSafeFunction 投递
199+
if (g_tsfn.BlockingCall(event, DispatchStreamEvent) != napi_ok) {
200+
delete event;
201+
}
202+
},
203+
[](const std::string& message) {
204+
auto* event = new StreamEvent();
205+
event->is_error = true;
206+
event->error = message;
207+
if (g_tsfn.BlockingCall(event, DispatchStreamEvent) != napi_ok) {
208+
delete event;
209+
}
210+
});
211+
212+
Napi::Object result = Napi::Object::New(env);
213+
if (!err.empty()) {
214+
g_tsfn.Release();
215+
g_capture.reset();
216+
result.Set("ok", Napi::Boolean::New(env, false));
217+
result.Set("error", Napi::String::New(env, err));
218+
return result;
219+
}
220+
result.Set("ok", Napi::Boolean::New(env, true));
221+
result.Set("error", Napi::String::New(env, ""));
222+
return result;
223+
}
224+
225+
/** stopCapture(): 停止采集并等待线程退出(幂等) */
226+
Napi::Value StopCapture(const Napi::CallbackInfo& info) {
227+
Napi::Env env = info.Env();
228+
if (g_capture) {
229+
g_capture->Stop();
230+
g_capture.reset();
231+
g_tsfn.Release();
232+
}
233+
return Napi::Boolean::New(env, true);
234+
}
235+
107236
Napi::Object Init(Napi::Env env, Napi::Object exports) {
108237
exports.Set("listAudioWindows", Napi::Function::New(env, ListAudioWindows));
109238
exports.Set("resolveRootPid", Napi::Function::New(env, ResolveRootPid));
110239
exports.Set("captureToWav", Napi::Function::New(env, CaptureToWav));
240+
exports.Set("isProcessLoopbackSupported",
241+
Napi::Function::New(env, IsSupported));
242+
exports.Set("startCapture", Napi::Function::New(env, StartCapture));
243+
exports.Set("stopCapture", Napi::Function::New(env, StopCapture));
111244
return exports;
112245
}
113246

0 commit comments

Comments
 (0)