diff --git a/cmvr-es/devices/microphone/CMakeLists.txt b/cmvr-es/devices/microphone/CMakeLists.txt index 1312054b..9ee1a945 100644 --- a/cmvr-es/devices/microphone/CMakeLists.txt +++ b/cmvr-es/devices/microphone/CMakeLists.txt @@ -10,4 +10,4 @@ target_link_libraries(microphone cmvr_es::proto ) -add_library(cmvr_es::device::microphone ALIAS microphone) +add_library(cmvr_es::device::microphone ALIAS microphone) \ No newline at end of file diff --git a/cmvr-es/devices/microphone/abstract_microphone.h b/cmvr-es/devices/microphone/abstract_microphone.h index 87ebcc45..6a4cdcb7 100644 --- a/cmvr-es/devices/microphone/abstract_microphone.h +++ b/cmvr-es/devices/microphone/abstract_microphone.h @@ -6,6 +6,8 @@ #define CMVR_ES_ABSTRACT_MICROPHONE_H #pragma once +#include + #include "devices/abstract_device.h" #include "cmvr/config/microphone_config/microphone_config.pb.h" namespace cmvr::device{ @@ -22,6 +24,12 @@ namespace cmvr::device{ virtual void resume() {} virtual void setVolume(const int volume) {} virtual int getVolume() {return 0;} + virtual bool startStreaming() {return true;} + virtual void stopStreaming() {} + virtual void getEncodedFrame(AudioStreamFrameData& frame_data, size_t& index) {} + virtual bool getLatestEncodedFrame(AudioStreamFrameData& frame_data, size_t& next_index) { + return false; + } protected: MicrophoneState state_{}; diff --git a/cmvr-es/devices/microphone/ffmpeg_microphone/include/ffmpeg_microphone.h b/cmvr-es/devices/microphone/ffmpeg_microphone/include/ffmpeg_microphone.h index 6fc65308..0e17046d 100644 --- a/cmvr-es/devices/microphone/ffmpeg_microphone/include/ffmpeg_microphone.h +++ b/cmvr-es/devices/microphone/ffmpeg_microphone/include/ffmpeg_microphone.h @@ -3,26 +3,30 @@ #ifndef CMVR_ES_FFMPEG_MICROPHONE_H #define CMVR_ES_FFMPEG_MICROPHONE_H -#include +#include +#include +#include #include #include -#include -#include "microphone/abstract_microphone.h" +#include + #include + +#include "common/base/ring_buffer.h" +#include "microphone/abstract_microphone.h" #include "speaker/ffmpeg_speaker/include/ffmpeg_ptr.h" namespace cmvr::device { class ffmpegMicroPhone final : public AbstractMicrophone { public: - // 录制状态 enum class RecordingState { STOPPED, RECORDING, PAUSED - }; - public: - ffmpegMicroPhone(const config::FFMpegMicroPhoneConfig& cfg); + }; + + explicit ffmpegMicroPhone(const config::FFMpegMicroPhoneConfig& cfg); ~ffmpegMicroPhone() override; std::string typeName() const override { return "FFMpegMicroPhone"; } bool init() override; @@ -36,49 +40,52 @@ namespace cmvr::device { void setVolume(const int volume) override; int getVolume() override; + bool startStreaming() override; + void stopStreaming() override; + void getEncodedFrame(AudioStreamFrameData& frame_data, size_t& index) override; + bool getLatestEncodedFrame(AudioStreamFrameData& frame_data, size_t& next_index) override; private: - // 音频采集和编码线程 void audioThread(); - - // 初始化FFmpeg采集和编码 bool initFFmpeg(); - - // 关闭FFmpeg资源 void closeFFmpeg(); + bool startCapture_(bool write_file); + void stopCapture_(); + AudioStreamFormat currentStreamFormat_() const; + void pushEncodedPacket_(const AVPacket* packet); private: - // FFmpeg采集相关 AVFormatContext* input_fmt_ctx; AVCodecContext* input_codec_ctx; AVFrame* input_frame; - // FFmpeg编码相关 AVFormatContext* output_fmt_ctx; - AVOutputFormat* output_fmt; + const AVOutputFormat* output_fmt; AVStream* audio_st; AVCodecContext* audio_codec_ctx; - AVCodec* audio_codec; + const AVCodec* audio_codec; SwrContext* swr_ctx; AVFrame* audio_frame; std::string output_file; std::string format_name; - std::atomic is_recording; + std::atomic is_capturing_{false}; + std::atomic write_output_file_{false}; std::atomic is_paused; std::mutex pause_mutex; std::condition_variable pause_cv; - // 音频参数 int sample_rate_; int channels_; int64_t next_pts; - private: std::string input_device_; std::shared_ptr audio_thread_; std::mutex mtx_; + std::shared_ptr> stream_frame_buffer_; + int stream_count_ = 0; + size_t buffer_size_ = 256; config::FFMpegMicroPhoneConfig config_; }; diff --git a/cmvr-es/devices/microphone/ffmpeg_microphone/src/ffmpeg_microphone.cpp b/cmvr-es/devices/microphone/ffmpeg_microphone/src/ffmpeg_microphone.cpp index 82dac71d..da212f5b 100644 --- a/cmvr-es/devices/microphone/ffmpeg_microphone/src/ffmpeg_microphone.cpp +++ b/cmvr-es/devices/microphone/ffmpeg_microphone/src/ffmpeg_microphone.cpp @@ -1,143 +1,297 @@ #include "common/base/logging/logger.h" #include "../include/ffmpeg_microphone.h" + +#include +#include +#include +#include +#include + using namespace std; using namespace cmvr::device; -ffmpegMicroPhone::ffmpegMicroPhone(const config::FFMpegMicroPhoneConfig& cfg):input_fmt_ctx(nullptr), input_codec_ctx(nullptr), input_frame(nullptr), - output_fmt_ctx(nullptr), output_fmt(nullptr), audio_st(nullptr), - audio_codec_ctx(nullptr), audio_codec(nullptr), swr_ctx(nullptr), - audio_frame(nullptr), is_recording(false), is_paused(false), - sample_rate_(44100), channels_(2), config_(cfg) +namespace { +std::string avErrorToString(int err) { - id_ = config_.id(); - input_device_ = config_.input_device(); - channels_ = config_.channels(); - sample_rate_ = config_.samplerate(); - state_.volume = config_.volume(); - state_.is_initialized = false; - state_.is_running = false; - state_.is_error = false; - state_.error_message= ""; + char errbuf[AV_ERROR_MAX_STRING_SIZE] = {0}; + av_strerror(err, errbuf, sizeof(errbuf)); + return errbuf; } -ffmpegMicroPhone::~ffmpegMicroPhone() { +void logSupportedAudioInputFormats() +{ + static const std::set audio_input_names = { + "alsa", "pulse", "oss", "jack", "openal", "sndio", "lavfi", "dshow", "avfoundation" + }; + + std::ostringstream supported; + void* opaque = nullptr; + const AVInputFormat* fmt = nullptr; + bool first = true; + while ((fmt = av_demuxer_iterate(&opaque)) != nullptr) { + if (!fmt->name || audio_input_names.count(fmt->name) == 0) { + continue; + } + if (!first) { + supported << ", "; + } + supported << fmt->name; + if (fmt->long_name) { + supported << " (" << fmt->long_name << ")"; + } + first = false; + } + + CMVR_LOG(INFO) << "[ffmpegMicroPhone] Supported audio input formats: " + << (first ? "none" : supported.str()); +} +} + +ffmpegMicroPhone::ffmpegMicroPhone(const config::FFMpegMicroPhoneConfig& cfg) + : input_fmt_ctx(nullptr), + input_codec_ctx(nullptr), + input_frame(nullptr), + output_fmt_ctx(nullptr), + output_fmt(nullptr), + audio_st(nullptr), + audio_codec_ctx(nullptr), + audio_codec(nullptr), + swr_ctx(nullptr), + audio_frame(nullptr), + is_paused(false), + sample_rate_(44100), + channels_(2), + next_pts(0), + config_(cfg) +{ + id_ = config_.id(); + input_device_ = config_.input_device().empty() ? "default" : config_.input_device(); + channels_ = config_.channels() > 0 ? config_.channels() : 2; + sample_rate_ = config_.samplerate() > 0 ? config_.samplerate() : 44100; + state_.volume = config_.volume() > 0 ? config_.volume() : 100; + state_.is_initialized = false; + state_.is_running = false; + state_.is_recording = false; + state_.is_error = false; + state_.error_message = ""; + stream_frame_buffer_ = std::make_shared>(buffer_size_); +} + +ffmpegMicroPhone::~ffmpegMicroPhone() +{ stop(); } -bool ffmpegMicroPhone::init() { +bool ffmpegMicroPhone::init() +{ try { lock_guard lock(mtx_); + stream_frame_buffer_->clear(); state_.is_initialized = true; CMVR_LOG(INFO) << "[ffmpegMicroPhone] (init): Success, id=" << id_; return true; - } - catch (const exception& e) { - CMVR_LOG(ERROR) <<"[ffmpegMicroPhone] (init): Failed, id=" << id_ << ": " << e.what(); + } catch (const exception& e) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] (init): Failed, id=" << id_ << ": " << e.what(); state_.is_error = true; state_.error_message = e.what(); return false; } } -bool ffmpegMicroPhone::start() { - try { +bool ffmpegMicroPhone::start() +{ + lock_guard lock(mtx_); + state_.is_initialized = true; + CMVR_LOG(INFO) << "[ffmpegMicroPhone] (start): Success, id=" << id_; + return true; +} + +bool ffmpegMicroPhone::stop() +{ + { lock_guard lock(mtx_); - CMVR_LOG(INFO) << "[ffmpegMicroPhone] (start): Success, id=" << id_; - return true; + stream_count_ = 0; + state_.is_recording = false; } - catch (const exception& e) { - CMVR_LOG(ERROR) <<"[ffmpegMicroPhone] (start): Failed, id=" << id_ << ": " << e.what(); - state_.is_error = true; - state_.error_message = e.what(); - return false; - } -} - -bool ffmpegMicroPhone::stop() { - stopRecording(); + stopCapture_(); CMVR_LOG(INFO) << "[ffmpegMicroPhone] (stop): Success, id=" << id_; return true; } -void ffmpegMicroPhone::pause() { - try { - is_paused = true; - CMVR_LOG(INFO) << "[ffmpegMicroPhone] (pause): Success, id=" << id_; +void ffmpegMicroPhone::pause() +{ + is_paused = true; + state_.is_error = false; + CMVR_LOG(INFO) << "[ffmpegMicroPhone] (pause): Success, id=" << id_; +} + +void ffmpegMicroPhone::resume() +{ + is_paused = false; + pause_cv.notify_one(); + state_.is_error = false; + CMVR_LOG(INFO) << "[ffmpegMicroPhone] (resume): Success, id=" << id_; +} + +void ffmpegMicroPhone::getState(MicrophoneState& state) +{ + lock_guard lock(mtx_); + state = state_; +} + +void ffmpegMicroPhone::startRecording(const std::string& outputFilePath) +{ + lock_guard lock(mtx_); + if (state_.is_recording) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] (startRecording): microphone is already recording"; + return; } - catch (const exception& e) { - CMVR_LOG(ERROR) <<"[ffmpegMicroPhone] (pause): Failed, id=" << id_ << ": " << e.what(); + if (stream_count_ > 0 || is_capturing_) { state_.is_error = true; - state_.error_message = e.what(); - } -} - -void ffmpegMicroPhone::resume() { - try { - is_paused = false; - pause_cv.notify_one(); - CMVR_LOG(INFO) << "[ffmpegMicroPhone] (resume): Success, id=" << id_; - } - catch (const exception& e) { - CMVR_LOG(ERROR) <<"[ffmpegMicroPhone] (resume): Failed, id=" << id_ << ": " << e.what(); - state_.is_error = true; - state_.error_message = e.what(); - } -} - -void ffmpegMicroPhone::getState(MicrophoneState& state) { - state.is_initialized = state_.is_initialized; - state.is_running = state_.is_running; - state.is_recording = state_.is_recording; - state.is_error = state_.is_error; - state.volume = state_.volume; - state.error_message = state_.error_message; -} - -void ffmpegMicroPhone::startRecording(const std::string& outputFilePath) { - try { - if (is_recording) { - CMVR_LOG(ERROR) << "[ffmpegMicroPhone] (startRecording): microphone is already recording"; - return; - } - // 从文件路径中提取扩展名 - std::string extension; - size_t dotPos = outputFilePath.find_last_of('.'); - if (dotPos != std::string::npos) { - extension = outputFilePath.substr(dotPos + 1); - } else { - extension = "wav"; - CMVR_LOG(WARNING) << "[ffmpegMicroPhone] (StartRecord): No file extension found, using wav as default, id=" << id_; - } - - output_file = outputFilePath; - format_name = extension; - - if (!initFFmpeg()) { - CMVR_LOG(ERROR) <<"[ffmpegMicroPhone] (StartRecord): FFmpeg init error, id=" << id_; - state_.is_error = true; - state_.error_message = "FFmpeg init error"; - return; - } - - is_recording = true; - is_paused = false; - audio_thread_ = std::make_shared(&ffmpegMicroPhone::audioThread, this); - } - catch (const exception& e) { - CMVR_LOG(ERROR) <<"[ffmpegMicroPhone] (StartRecord): Failed, id=" << id_ << ": " << e.what(); - state_.is_error = true; - state_.error_message = e.what(); - } -} - -void ffmpegMicroPhone::stopRecording() { - if (!is_recording) - { - CMVR_LOG(WARNING) << "[ffmpegMicroPhone] (stopRecording): recording has not been started"; + state_.error_message = "cannot start recording while microphone streaming is active"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] " << state_.error_message; return; } - is_recording = false; + size_t dot_pos = outputFilePath.find_last_of('.'); + format_name = dot_pos == std::string::npos ? "wav" : outputFilePath.substr(dot_pos + 1); + output_file = outputFilePath; + + if (!startCapture_(true)) { + state_.is_error = true; + state_.error_message = "FFmpeg init error"; + return; + } + + state_.is_recording = true; +} + +void ffmpegMicroPhone::stopRecording() +{ + { + lock_guard lock(mtx_); + if (!state_.is_recording) { + return; + } + state_.is_recording = false; + } + stopCapture_(); + CMVR_LOG(INFO) << "[ffmpegMicroPhone] (stopRecording): Success, id=" << id_; +} + +void ffmpegMicroPhone::setVolume(const int volume) +{ + lock_guard lock(mtx_); + state_.volume = std::clamp(volume, 0, 100); +} + +int ffmpegMicroPhone::getVolume() +{ + lock_guard lock(mtx_); + return state_.volume; +} + +bool ffmpegMicroPhone::startStreaming() +{ + lock_guard lock(mtx_); + if (state_.is_recording) { + state_.is_error = true; + state_.error_message = "cannot start streaming while microphone recording is active"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] " << state_.error_message; + return false; + } + if (stream_count_ == 0) { + stream_frame_buffer_->clear(); + format_name = "pcm_s16le"; + output_file.clear(); + if (!startCapture_(false)) { + state_.is_error = true; + state_.error_message = "FFmpeg init error"; + return false; + } + } + ++stream_count_; + state_.is_running = true; + state_.is_error = false; + state_.error_message.clear(); + return true; +} + +void ffmpegMicroPhone::stopStreaming() +{ + bool should_stop = false; + { + lock_guard lock(mtx_); + if (stream_count_ > 0) { + --stream_count_; + } + should_stop = stream_count_ == 0 && !state_.is_recording; + } + if (should_stop) { + stopCapture_(); + } +} + +void ffmpegMicroPhone::getEncodedFrame(AudioStreamFrameData& frame_data, size_t& index) +{ + if (!stream_frame_buffer_) { + return; + } + auto frame = stream_frame_buffer_->pop(index); + if (frame.has_value()) { + frame_data = frame.value(); + } +} + +bool ffmpegMicroPhone::getLatestEncodedFrame(AudioStreamFrameData& frame_data, size_t& next_index) +{ + if (!stream_frame_buffer_) { + return false; + } + + const size_t head = stream_frame_buffer_->getHead(); + if (head == 0) { + return false; + } + + size_t latest_index = head - 1; + auto frame = stream_frame_buffer_->pop(latest_index); + if (!frame.has_value()) { + return false; + } + + frame_data = frame.value(); + next_index = latest_index; + return true; +} + +bool ffmpegMicroPhone::startCapture_(bool write_file) +{ + if (is_capturing_) { + return true; + } + write_output_file_ = write_file; + if (!initFFmpeg()) { + closeFFmpeg(); + return false; + } + is_capturing_ = true; + is_paused = false; + state_.is_running = true; + audio_thread_ = std::make_shared(&ffmpegMicroPhone::audioThread, this); + return true; +} + +void ffmpegMicroPhone::stopCapture_() +{ + if (!is_capturing_) { + return; + } + + is_capturing_ = false; + is_paused = false; + pause_cv.notify_all(); + if (audio_thread_) { if (audio_thread_->joinable()) { audio_thread_->join(); @@ -145,29 +299,30 @@ void ffmpegMicroPhone::stopRecording() { audio_thread_.reset(); } - if (output_fmt_ctx) { + if (write_output_file_ && output_fmt_ctx) { av_write_trailer(output_fmt_ctx); } - closeFFmpeg(); - CMVR_LOG(INFO) << "[ffmpegMicroPhone] (stopRecording): Success, id=" << id_; + write_output_file_ = false; + + lock_guard lock(mtx_); + state_.is_running = false; } -void ffmpegMicroPhone::setVolume(const int volume) { - state_.volume = volume; -} - -int ffmpegMicroPhone::getVolume() { - return state_.volume; -} - -bool ffmpegMicroPhone::initFFmpeg() { - int err; +bool ffmpegMicroPhone::initFFmpeg() +{ + int err = 0; avdevice_register_all(); - // 初始化输入设备 - AVInputFormat* input_fmt = av_find_input_format("pulse"); + logSupportedAudioInputFormats(); + + const char* input_format_name = "pulse"; + AVInputFormat* input_fmt = av_find_input_format(input_format_name); if (!input_fmt) { - CMVR_LOG(ERROR) << "Could not find pulse input format"; + input_format_name = "alsa"; + input_fmt = av_find_input_format(input_format_name); + } + if (!input_fmt) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not find pulse or alsa input format"; return false; } @@ -176,158 +331,148 @@ bool ffmpegMicroPhone::initFFmpeg() { av_dict_set(&input_options, "channels", std::to_string(channels_).c_str(), 0); av_dict_set(&input_options, "format", "s16", 0); - if ((err = avformat_open_input(&input_fmt_ctx, input_device_.c_str(), input_fmt, &input_options)) < 0) { - char errbuf[AV_ERROR_MAX_STRING_SIZE]; - av_strerror(err, errbuf, AV_ERROR_MAX_STRING_SIZE); - CMVR_LOG(ERROR) << "Could not open input device: " << errbuf; - av_dict_free(&input_options); + CMVR_LOG(INFO) << "[ffmpegMicroPhone] Open input device: format=" + << input_format_name << ", device=" << input_device_; + + err = avformat_open_input(&input_fmt_ctx, input_device_.c_str(), input_fmt, &input_options); + av_dict_free(&input_options); + if (err < 0) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not open input device: " << avErrorToString(err); return false; } - av_dict_free(&input_options); if ((err = avformat_find_stream_info(input_fmt_ctx, nullptr)) < 0) { - CMVR_LOG(ERROR) << "Could not find stream info"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not find stream info: " << avErrorToString(err); return false; } - // 查找音频流 - int audio_stream_index = av_find_best_stream(input_fmt_ctx, AVMEDIA_TYPE_AUDIO, -1, -1, nullptr, 0); + const int audio_stream_index = av_find_best_stream(input_fmt_ctx, AVMEDIA_TYPE_AUDIO, -1, -1, nullptr, 0); if (audio_stream_index < 0) { - CMVR_LOG(ERROR) << "Could not find audio stream"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not find audio stream"; return false; } - // 获取输入编解码器上下文 AVCodecParameters* codecpar = input_fmt_ctx->streams[audio_stream_index]->codecpar; const AVCodec* input_codec = avcodec_find_decoder(codecpar->codec_id); if (!input_codec) { - CMVR_LOG(ERROR) << "Could not find input codec"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not find input codec"; return false; } input_codec_ctx = avcodec_alloc_context3(input_codec); if (!input_codec_ctx) { - CMVR_LOG(ERROR) << "Could not allocate input codec context"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not allocate input codec context"; return false; } if ((err = avcodec_parameters_to_context(input_codec_ctx, codecpar)) < 0) { - CMVR_LOG(ERROR) << "Could not copy input codec parameters"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not copy input codec parameters"; return false; } if ((err = avcodec_open2(input_codec_ctx, input_codec, nullptr)) < 0) { - CMVR_LOG(ERROR) << "Could not open input codec"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not open input codec: " << avErrorToString(err); return false; } - // 初始化输出格式 - output_fmt = av_guess_format(format_name.c_str(), nullptr, nullptr); - if (!output_fmt) { - CMVR_LOG(ERROR) << "Unsupported format: " << format_name; - return false; + if (write_output_file_) { + output_fmt = av_guess_format(format_name.c_str(), nullptr, nullptr); + if (!output_fmt) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Unsupported format: " << format_name; + return false; + } + if ((err = avformat_alloc_output_context2(&output_fmt_ctx, nullptr, format_name.c_str(), output_file.c_str())) < 0) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not allocate output context: " << avErrorToString(err); + return false; + } + output_fmt = output_fmt_ctx->oformat; + audio_codec = avcodec_find_encoder(output_fmt->audio_codec); + } else { + audio_codec = avcodec_find_encoder_by_name("pcm_s16le"); } - if ((err = avformat_alloc_output_context2(&output_fmt_ctx, output_fmt, nullptr, output_file.c_str())) < 0) { - CMVR_LOG(ERROR) << "Could not allocate output context"; - return false; - } - - // 查找音频编码器 - audio_codec = avcodec_find_encoder(output_fmt->audio_codec); if (!audio_codec) { - CMVR_LOG(ERROR) << "Could not find audio encoder"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not find audio encoder"; return false; } - // 添加音频流 - audio_st = avformat_new_stream(output_fmt_ctx, nullptr); - if (!audio_st) { - CMVR_LOG(ERROR) << "Could not create audio stream"; - return false; + if (write_output_file_) { + audio_st = avformat_new_stream(output_fmt_ctx, nullptr); + if (!audio_st) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not create audio stream"; + return false; + } } audio_codec_ctx = avcodec_alloc_context3(audio_codec); if (!audio_codec_ctx) { - CMVR_LOG(ERROR) << "Could not allocate audio codec context"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not allocate audio codec context"; return false; } - // 设置音频编码参数 - audio_codec_ctx->sample_fmt = audio_codec->sample_fmts ? audio_codec->sample_fmts[0] : AV_SAMPLE_FMT_FLTP; + audio_codec_ctx->sample_fmt = audio_codec->sample_fmts ? audio_codec->sample_fmts[0] : AV_SAMPLE_FMT_S16; audio_codec_ctx->bit_rate = 128000; audio_codec_ctx->sample_rate = sample_rate_; audio_codec_ctx->channel_layout = channels_ == 1 ? AV_CH_LAYOUT_MONO : AV_CH_LAYOUT_STEREO; audio_codec_ctx->channels = channels_; - audio_st->time_base = {1, sample_rate_}; + audio_codec_ctx->time_base = {1, sample_rate_}; + if (audio_st) { + audio_st->time_base = {1, sample_rate_}; + } - // 打开音频编码器 AVDictionary* opt = nullptr; if ((err = avcodec_open2(audio_codec_ctx, audio_codec, &opt)) < 0) { - char errbuf[AV_ERROR_MAX_STRING_SIZE]; - av_strerror(err, errbuf, AV_ERROR_MAX_STRING_SIZE); - CMVR_LOG(ERROR) << "Could not open audio codec: " << errbuf; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not open audio codec: " << avErrorToString(err); av_dict_free(&opt); return false; } av_dict_free(&opt); - // 复制编码参数到流 - if (avcodec_parameters_from_context(audio_st->codecpar, audio_codec_ctx) < 0) { - CMVR_LOG(ERROR) << "Could not copy codec parameters to stream"; + if (audio_st && avcodec_parameters_from_context(audio_st->codecpar, audio_codec_ctx) < 0) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not copy codec parameters to stream"; return false; } - // 分配音频帧 input_frame = av_frame_alloc(); audio_frame = av_frame_alloc(); if (!input_frame || !audio_frame) { - CMVR_LOG(ERROR) << "Could not allocate audio frames"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not allocate audio frames"; return false; } audio_frame->format = audio_codec_ctx->sample_fmt; audio_frame->channel_layout = audio_codec_ctx->channel_layout; + audio_frame->channels = audio_codec_ctx->channels; audio_frame->sample_rate = audio_codec_ctx->sample_rate; audio_frame->nb_samples = audio_codec_ctx->frame_size ? audio_codec_ctx->frame_size : 1024; - if (av_frame_get_buffer(audio_frame, 0) < 0) { - CMVR_LOG(ERROR) << "Could not allocate audio frame buffers"; - return false; - } - // 创建重采样上下文 - swr_ctx = swr_alloc(); - if (!swr_ctx) { - CMVR_LOG(ERROR) << "Could not allocate resampler context"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not allocate audio frame buffers"; return false; } swr_ctx = swr_alloc_set_opts(nullptr, - audio_codec_ctx->channel_layout, - audio_codec_ctx->sample_fmt, - audio_codec_ctx->sample_rate, - av_get_default_channel_layout(channels_), - input_codec_ctx->sample_fmt, - input_codec_ctx->sample_rate, - 0, nullptr); + audio_codec_ctx->channel_layout, + audio_codec_ctx->sample_fmt, + audio_codec_ctx->sample_rate, + av_get_default_channel_layout(channels_), + input_codec_ctx->sample_fmt, + input_codec_ctx->sample_rate, + 0, + nullptr); if (!swr_ctx || swr_init(swr_ctx) < 0) { - CMVR_LOG(ERROR) << "Failed to initialize resampler"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Failed to initialize resampler"; return false; } - // 打开输出文件 - if (!(output_fmt->flags & AVFMT_NOFILE)) { + if (write_output_file_ && output_fmt_ctx && !(output_fmt->flags & AVFMT_NOFILE)) { if ((err = avio_open(&output_fmt_ctx->pb, output_file.c_str(), AVIO_FLAG_WRITE)) < 0) { - char errbuf[AV_ERROR_MAX_STRING_SIZE]; - av_strerror(err, errbuf, AV_ERROR_MAX_STRING_SIZE); - CMVR_LOG(ERROR) << "Could not open output file: " << errbuf; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not open output file: " << avErrorToString(err); return false; } } - // 写入文件头 - if (avformat_write_header(output_fmt_ctx, nullptr) < 0) { - CMVR_LOG(ERROR) << "Could not write output file header"; + if (write_output_file_ && output_fmt_ctx && avformat_write_header(output_fmt_ctx, nullptr) < 0) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not write output file header"; return false; } @@ -335,94 +480,111 @@ bool ffmpegMicroPhone::initFFmpeg() { return true; } -void ffmpegMicroPhone::audioThread() { +void ffmpegMicroPhone::audioThread() +{ AVPacket* input_pkt = av_packet_alloc(); AVPacket* output_pkt = av_packet_alloc(); if (!input_pkt || !output_pkt) { - CMVR_LOG(ERROR) << "Could not allocate packets"; + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not allocate packets"; + av_packet_free(&input_pkt); + av_packet_free(&output_pkt); return; } - while (is_recording) { + while (is_capturing_) { if (is_paused) { - std::unique_lock lock(pause_mutex); - pause_cv.wait(lock, [this] { return !is_paused || !is_recording; }); - if (!is_recording) break; + unique_lock lock(pause_mutex); + pause_cv.wait(lock, [this] { return !is_paused || !is_capturing_; }); continue; } - // 从输入设备读取数据 - int ret = av_read_frame(input_fmt_ctx, input_pkt); - if (ret < 0) { - if (ret == AVERROR(EAGAIN)) continue; + const int read_ret = av_read_frame(input_fmt_ctx, input_pkt); + if (read_ret < 0) { + if (read_ret == AVERROR(EAGAIN)) { + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + continue; + } break; } - // 解码音频帧 - ret = avcodec_send_packet(input_codec_ctx, input_pkt); + int ret = avcodec_send_packet(input_codec_ctx, input_pkt); + av_packet_unref(input_pkt); if (ret < 0) { - av_packet_unref(input_pkt); continue; } while (avcodec_receive_frame(input_codec_ctx, input_frame) == 0) { - if (input_frame->format != AV_SAMPLE_FMT_S16) { - CMVR_LOG(ERROR) << "Input frame format mismatch"; - return; - } - // 重采样 - int dst_nb_samples = av_rescale_rnd( + const int dst_nb_samples = static_cast(av_rescale_rnd( swr_get_delay(swr_ctx, sample_rate_) + input_frame->nb_samples, - sample_rate_, sample_rate_, AV_ROUND_UP - ); + sample_rate_, + sample_rate_, + AV_ROUND_UP)); - int ret = swr_convert( - swr_ctx, audio_frame->data, dst_nb_samples, // 使用计算出的样本数 - (const uint8_t**)input_frame->data, input_frame->nb_samples - ); - if (ret < 0) { - CMVR_LOG(ERROR) << "swr_convert failed"; - av_frame_unref(input_frame); + if (dst_nb_samples > audio_frame->nb_samples) { + av_frame_unref(audio_frame); + audio_frame->format = audio_codec_ctx->sample_fmt; + audio_frame->channel_layout = audio_codec_ctx->channel_layout; + audio_frame->channels = audio_codec_ctx->channels; + audio_frame->sample_rate = audio_codec_ctx->sample_rate; + audio_frame->nb_samples = dst_nb_samples; + const int buffer_ret = av_frame_get_buffer(audio_frame, 0); + if (buffer_ret < 0) { + CMVR_LOG(ERROR) << "[ffmpegMicroPhone] Could not grow audio frame buffer: " + << avErrorToString(buffer_ret); + break; + } + } + + if (av_frame_make_writable(audio_frame) < 0) { continue; } - audio_frame->nb_samples = ret; // 更新实际转换的样本数 - // 更准确的pts计算 - audio_frame->pts = av_rescale_q( - next_pts, (AVRational){1, sample_rate_}, audio_codec_ctx->time_base - ); - next_pts += audio_frame->nb_samples; // 使用实际样本数更新 + ret = swr_convert( + swr_ctx, + audio_frame->data, + dst_nb_samples, + const_cast(input_frame->data), + input_frame->nb_samples); + av_frame_unref(input_frame); + if (ret < 0) { + continue; + } + + audio_frame->nb_samples = ret; + audio_frame->pts = av_rescale_q(next_pts, AVRational{1, sample_rate_}, audio_codec_ctx->time_base); + next_pts += audio_frame->nb_samples; - // 编码音频帧 ret = avcodec_send_frame(audio_codec_ctx, audio_frame); if (ret < 0) { continue; } while (avcodec_receive_packet(audio_codec_ctx, output_pkt) == 0) { - // 写入数据包 - av_packet_rescale_ts(output_pkt, audio_codec_ctx->time_base, audio_st->time_base); - output_pkt->stream_index = audio_st->index; - - if (av_interleaved_write_frame(output_fmt_ctx, output_pkt) < 0) { - av_packet_unref(output_pkt); - break; + if (write_output_file_ && output_fmt_ctx && audio_st) { + av_packet_rescale_ts(output_pkt, audio_codec_ctx->time_base, audio_st->time_base); + output_pkt->stream_index = audio_st->index; + if (av_interleaved_write_frame(output_fmt_ctx, output_pkt) < 0) { + av_packet_unref(output_pkt); + break; + } } + pushEncodedPacket_(output_pkt); av_packet_unref(output_pkt); } } - av_packet_unref(input_pkt); } - // 刷新编码器 avcodec_send_frame(audio_codec_ctx, nullptr); while (avcodec_receive_packet(audio_codec_ctx, output_pkt) == 0) { - av_packet_rescale_ts(output_pkt, audio_codec_ctx->time_base, audio_st->time_base); - output_pkt->stream_index = audio_st->index; - - if (av_interleaved_write_frame(output_fmt_ctx, output_pkt) < 0) { - break; + if (write_output_file_ && output_fmt_ctx && audio_st) { + av_packet_rescale_ts(output_pkt, audio_codec_ctx->time_base, audio_st->time_base); + output_pkt->stream_index = audio_st->index; + if (av_interleaved_write_frame(output_fmt_ctx, output_pkt) < 0) { + av_packet_unref(output_pkt); + break; + } } + pushEncodedPacket_(output_pkt); av_packet_unref(output_pkt); } @@ -430,38 +592,71 @@ void ffmpegMicroPhone::audioThread() { av_packet_free(&output_pkt); } -void ffmpegMicroPhone::closeFFmpeg() { +AudioStreamFormat ffmpegMicroPhone::currentStreamFormat_() const +{ + if (!audio_codec_ctx) { + return AudioStreamFormat::UNKNOWN; + } + switch (audio_codec_ctx->codec_id) { + case AV_CODEC_ID_AAC: + return AudioStreamFormat::AAC; + case AV_CODEC_ID_MP3: + return AudioStreamFormat::MP3; + case AV_CODEC_ID_PCM_S16LE: + case AV_CODEC_ID_PCM_S16BE: + return AudioStreamFormat::PCM; + default: + return AudioStreamFormat::UNKNOWN; + } +} +void ffmpegMicroPhone::pushEncodedPacket_(const AVPacket* packet) +{ + if (!packet || !packet->data || packet->size <= 0 || stream_count_ <= 0 || !stream_frame_buffer_) { + return; + } + + AudioStreamFrameData frame; + frame.data.assign(packet->data, packet->data + packet->size); + frame.sample_rate = sample_rate_; + frame.channels = channels_; + frame.format = currentStreamFormat_(); + frame.codec = audio_codec_ctx && audio_codec_ctx->codec ? audio_codec_ctx->codec->name : "pcm_s16le"; + frame.pts = packet->pts; + frame.nb_samples = audio_frame ? audio_frame->nb_samples : 0; + stream_frame_buffer_->push(frame); +} + +void ffmpegMicroPhone::closeFFmpeg() +{ if (input_frame) { - av_frame_unref(input_frame); // 释放内部缓冲区(如果有) - av_frame_free(&input_frame); // 释放 AVFrame 结构体 + av_frame_free(&input_frame); input_frame = nullptr; } if (audio_frame) { - av_frame_unref(audio_frame); // 释放内部缓冲区(如果有) - av_frame_free(&audio_frame); // 释放 AVFrame 结构体 + av_frame_free(&audio_frame); audio_frame = nullptr; } - if (swr_ctx) swr_free(&swr_ctx); - + if (swr_ctx) { + swr_free(&swr_ctx); + } if (input_codec_ctx) { - avcodec_close(input_codec_ctx); avcodec_free_context(&input_codec_ctx); } - if (audio_codec_ctx) { - avcodec_close(audio_codec_ctx); avcodec_free_context(&audio_codec_ctx); } - if (input_fmt_ctx) { avformat_close_input(&input_fmt_ctx); } - if (output_fmt_ctx) { - if (!(output_fmt->flags & AVFMT_NOFILE)) { + if (output_fmt && !(output_fmt->flags & AVFMT_NOFILE)) { avio_closep(&output_fmt_ctx->pb); } avformat_free_context(output_fmt_ctx); + output_fmt_ctx = nullptr; } + output_fmt = nullptr; + audio_st = nullptr; + audio_codec = nullptr; } diff --git a/cmvr-es/devices/speaker/abstract_speaker.h b/cmvr-es/devices/speaker/abstract_speaker.h index bffaccde..ab24542d 100644 --- a/cmvr-es/devices/speaker/abstract_speaker.h +++ b/cmvr-es/devices/speaker/abstract_speaker.h @@ -21,6 +21,8 @@ namespace cmvr::device{ virtual int getVolume() const {return 0;} virtual void pause() {} virtual void resume() {} + virtual bool pushAudioFrame(const AudioStreamFrameData& frame_data) { return false; } + virtual void stopStreaming() {} protected: SpeakerState state_{}; diff --git a/cmvr-es/devices/speaker/ffmpeg_speaker/include/ffmpeg_speaker.h b/cmvr-es/devices/speaker/ffmpeg_speaker/include/ffmpeg_speaker.h index 998e526f..9d8e56ed 100644 --- a/cmvr-es/devices/speaker/ffmpeg_speaker/include/ffmpeg_speaker.h +++ b/cmvr-es/devices/speaker/ffmpeg_speaker/include/ffmpeg_speaker.h @@ -6,6 +6,7 @@ #include #include #include +#include #include "speaker/abstract_speaker.h" #include #include @@ -33,6 +34,8 @@ namespace cmvr::device { void pause() override; void resume() override; void getState(SpeakerState& state) override; + bool pushAudioFrame(const AudioStreamFrameData& frame_data) override; + void stopStreaming() override; void resetPlayState(); bool initPulseDevice_(); @@ -40,6 +43,11 @@ namespace cmvr::device { private: void decode_audio_(); void play_audio_(); + bool startStreamingPlayback_(const AudioStreamFrameData& frame_data); + bool pushPcmFrame_(const AudioStreamFrameData& frame_data); + bool decodeStreamFrame_(const AudioStreamFrameData& frame_data); + bool initStreamDecoder_(const AudioStreamFrameData& frame_data); + void releaseStreamDecoder_(); pa_simple* pulse_simple_ = nullptr; // 修改:PulseAudio 简单 API 句柄 pa_sample_spec sample_spec_{}; // 新增:PulseAudio 采样规格 @@ -64,6 +72,15 @@ namespace cmvr::device { // 音频时钟同步 std::atomic audio_pts_{0}; std::chrono::time_point playback_start_time_; + std::mutex stream_decode_mtx_; + AVCodecContext* stream_decoder_ctx_ = nullptr; + SwrContext* stream_swr_ctx_ = nullptr; + AVFrame* stream_decode_frame_ = nullptr; + AVPacket* stream_decode_packet_ = nullptr; + std::string stream_codec_; + int stream_input_sample_rate_ = 0; + int stream_input_channels_ = 0; + bool is_streaming_input_ = false; config::FFMpegSpeakerConfig config_; }; diff --git a/cmvr-es/devices/speaker/ffmpeg_speaker/src/ffmpeg_speaker.cpp b/cmvr-es/devices/speaker/ffmpeg_speaker/src/ffmpeg_speaker.cpp index 366c8ea1..ee19882c 100644 --- a/cmvr-es/devices/speaker/ffmpeg_speaker/src/ffmpeg_speaker.cpp +++ b/cmvr-es/devices/speaker/ffmpeg_speaker/src/ffmpeg_speaker.cpp @@ -4,6 +4,9 @@ // #include +#include +#include +#include #include "../include/ffmpeg_speaker.h" @@ -18,7 +21,7 @@ ffmpegSpeaker::ffmpegSpeaker(const config::FFMpegSpeakerConfig& cfg):config_(cfg try { id_ = config_.id(); memset(&sample_spec_, 0, sizeof(sample_spec_)); - // state_.volume = config_.volume(); + state_.volume = 100; } catch (const exception& e) { CMVR_LOG(ERROR) << "[ffmpegSpeaker] ([ffmpegSpeaker]): Failed to parse config: " << e.what(); @@ -76,6 +79,8 @@ bool ffmpegSpeaker::start() { void ffmpegSpeaker::resetPlayState() { + releaseStreamDecoder_(); + if (pulse_simple_) { pa_simple_flush(pulse_simple_, nullptr); pa_simple_free(pulse_simple_); @@ -89,6 +94,7 @@ void ffmpegSpeaker::resetPlayState() } state_.is_initialized = false; + is_streaming_input_ = false; audio_path_.clear(); CMVR_LOG(INFO) << "[ffmpegSpeaker] (resetPlayState): Success, id=" << id_; @@ -150,6 +156,7 @@ void ffmpegSpeaker::resume() { return; } state_.is_paused = false; + cv_pause_.notify_all(); CMVR_LOG(INFO) << "[ffmpegSpeaker] (resume): Success, id=" << id_; } catch (const exception& e) { @@ -531,6 +538,286 @@ void ffmpegSpeaker::play_audio_() { CMVR_LOG(INFO) << "[ffmpegSpeaker] Playback finished"; } +bool ffmpegSpeaker::pushAudioFrame(const AudioStreamFrameData& frame_data) +{ + if (frame_data.data.empty()) { + return true; + } + + bool need_start = false; + bool need_restart = false; + { + std::lock_guard lock(mtx_); + need_start = !state_.is_running; + need_restart = state_.is_running && is_streaming_input_ && + (frame_data.sample_rate != sample_rate_ || frame_data.channels != channels_); + + if (state_.is_running && !is_streaming_input_) { + CMVR_LOG(ERROR) << "[ffmpegSpeaker] cannot stream audio while file playback is active"; + return false; + } + } + + if (need_restart) { + stopStreaming(); + need_start = true; + } + + if (need_start && !startStreamingPlayback_(frame_data)) { + return false; + } + + if (frame_data.format == AudioStreamFormat::PCM || + frame_data.codec == "pcm_s16le" || + frame_data.codec == "pcm") { + return pushPcmFrame_(frame_data); + } + return decodeStreamFrame_(frame_data); +} + +void ffmpegSpeaker::stopStreaming() +{ + { + std::lock_guard lock(mtx_); + if (!is_streaming_input_) { + return; + } + state_.is_decoding = false; + state_.is_paused = false; + is_streaming_input_ = false; + } + + cv_pause_.notify_all(); + + if (play_thread_ && play_thread_->joinable()) { + play_thread_->join(); + play_thread_.reset(); + } + + resetPlayState(); +} + +bool ffmpegSpeaker::startStreamingPlayback_(const AudioStreamFrameData& frame_data) +{ + std::lock_guard lock(mtx_); + + sample_rate_ = frame_data.sample_rate > 0 ? frame_data.sample_rate : 44100; + channels_ = frame_data.channels > 0 ? frame_data.channels : 2; + + AudioFrame queued_frame; + while (audio_queue_.pop(queued_frame)) { + } + + if (!initPulseDevice_()) { + CMVR_LOG(ERROR) << "[ffmpegSpeaker] PulseAudio device initialization failed"; + return false; + } + + state_.is_running = true; + state_.is_decoding = true; + state_.is_paused = false; + is_streaming_input_ = true; + + play_thread_ = std::make_shared([this]() { + try { + this->play_audio_(); + } catch (const std::exception& e) { + CMVR_LOG(ERROR) << "[ffmpegSpeaker] Stream play thread exception: " << e.what(); + std::lock_guard lock(mtx_); + state_.is_running = false; + state_.is_decoding = false; + } + }); + + return true; +} + +bool ffmpegSpeaker::pushPcmFrame_(const AudioStreamFrameData& frame_data) +{ + const size_t sample_count = frame_data.data.size() / sizeof(int16_t); + if (sample_count == 0) { + return true; + } + + auto buffer = std::make_shared>(sample_count); + std::memcpy(buffer->data(), frame_data.data.data(), sample_count * sizeof(int16_t)); + + float volume_scale = 1.0f; + { + std::lock_guard lock(mtx_); + volume_scale = static_cast(state_.volume) / 100.0f; + } + + for (int16_t& sample : *buffer) { + const float scaled = static_cast(sample) * volume_scale; + sample = static_cast(std::clamp(scaled, -32768.f, 32767.f)); + } + + while (!audio_queue_.push(buffer)) { + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + std::lock_guard lock(mtx_); + if (!state_.is_running) { + return false; + } + } + return true; +} + +bool ffmpegSpeaker::decodeStreamFrame_(const AudioStreamFrameData& frame_data) +{ + std::lock_guard decode_lock(stream_decode_mtx_); + if (!stream_decoder_ctx_ || + stream_codec_ != frame_data.codec || + stream_input_sample_rate_ != frame_data.sample_rate || + stream_input_channels_ != frame_data.channels) { + releaseStreamDecoder_(); + if (!initStreamDecoder_(frame_data)) { + return false; + } + } + + av_packet_unref(stream_decode_packet_); + if (av_new_packet(stream_decode_packet_, static_cast(frame_data.data.size())) < 0) { + return false; + } + std::memcpy(stream_decode_packet_->data, frame_data.data.data(), frame_data.data.size()); + stream_decode_packet_->pts = frame_data.pts; + + int ret = avcodec_send_packet(stream_decoder_ctx_, stream_decode_packet_); + av_packet_unref(stream_decode_packet_); + if (ret < 0) { + return false; + } + + while ((ret = avcodec_receive_frame(stream_decoder_ctx_, stream_decode_frame_)) == 0) { + if (!stream_swr_ctx_) { + stream_swr_ctx_ = swr_alloc_set_opts(nullptr, + av_get_default_channel_layout(channels_), + AV_SAMPLE_FMT_S16, + sample_rate_, + av_get_default_channel_layout(stream_decode_frame_->channels), + static_cast(stream_decode_frame_->format), + stream_decode_frame_->sample_rate, + 0, + nullptr); + if (!stream_swr_ctx_ || swr_init(stream_swr_ctx_) < 0) { + CMVR_LOG(ERROR) << "[ffmpegSpeaker] Failed to init stream resampler"; + av_frame_unref(stream_decode_frame_); + return false; + } + } + + const int64_t max_samples64 = av_rescale_rnd( + swr_get_delay(stream_swr_ctx_, stream_decode_frame_->sample_rate) + stream_decode_frame_->nb_samples, + sample_rate_, + stream_decode_frame_->sample_rate, + AV_ROUND_UP); + const int max_samples = static_cast(std::min(max_samples64, INT_MAX)); + + auto buffer = std::make_shared>(max_samples * channels_); + uint8_t* out[] = {reinterpret_cast(buffer->data()), nullptr}; + const int out_samples = swr_convert(stream_swr_ctx_, + out, + max_samples, + const_cast(stream_decode_frame_->data), + stream_decode_frame_->nb_samples); + av_frame_unref(stream_decode_frame_); + if (out_samples <= 0) { + continue; + } + + buffer->resize(out_samples * channels_); + + float volume_scale = 1.0f; + { + std::lock_guard lock(mtx_); + volume_scale = static_cast(state_.volume) / 100.0f; + } + for (int16_t& sample : *buffer) { + const float scaled = static_cast(sample) * volume_scale; + sample = static_cast(std::clamp(scaled, -32768.f, 32767.f)); + } + + while (!audio_queue_.push(buffer)) { + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + std::lock_guard lock(mtx_); + if (!state_.is_running) { + return false; + } + } + } + + return ret == AVERROR(EAGAIN) || ret == AVERROR_EOF; +} + +bool ffmpegSpeaker::initStreamDecoder_(const AudioStreamFrameData& frame_data) +{ + const AVCodec* decoder = nullptr; + if (!frame_data.codec.empty()) { + decoder = avcodec_find_decoder_by_name(frame_data.codec.c_str()); + } + if (!decoder) { + AVCodecID codec_id = AV_CODEC_ID_NONE; + if (frame_data.format == AudioStreamFormat::AAC) { + codec_id = AV_CODEC_ID_AAC; + } else if (frame_data.format == AudioStreamFormat::MP3) { + codec_id = AV_CODEC_ID_MP3; + } + if (codec_id != AV_CODEC_ID_NONE) { + decoder = avcodec_find_decoder(codec_id); + } + } + if (!decoder) { + CMVR_LOG(ERROR) << "[ffmpegSpeaker] Unsupported stream audio codec: " << frame_data.codec; + return false; + } + + stream_decoder_ctx_ = avcodec_alloc_context3(decoder); + if (!stream_decoder_ctx_) { + return false; + } + stream_decoder_ctx_->sample_rate = frame_data.sample_rate > 0 ? frame_data.sample_rate : sample_rate_; + stream_decoder_ctx_->channels = frame_data.channels > 0 ? frame_data.channels : channels_; + stream_decoder_ctx_->channel_layout = + stream_decoder_ctx_->channels == 1 ? AV_CH_LAYOUT_MONO : AV_CH_LAYOUT_STEREO; + + if (avcodec_open2(stream_decoder_ctx_, decoder, nullptr) < 0) { + releaseStreamDecoder_(); + return false; + } + + stream_decode_frame_ = av_frame_alloc(); + stream_decode_packet_ = av_packet_alloc(); + if (!stream_decode_frame_ || !stream_decode_packet_) { + releaseStreamDecoder_(); + return false; + } + + stream_codec_ = frame_data.codec; + stream_input_sample_rate_ = frame_data.sample_rate; + stream_input_channels_ = frame_data.channels; + return true; +} + +void ffmpegSpeaker::releaseStreamDecoder_() +{ + if (stream_swr_ctx_) { + swr_free(&stream_swr_ctx_); + } + if (stream_decode_frame_) { + av_frame_free(&stream_decode_frame_); + } + if (stream_decode_packet_) { + av_packet_free(&stream_decode_packet_); + } + if (stream_decoder_ctx_) { + avcodec_free_context(&stream_decoder_ctx_); + } + stream_codec_.clear(); + stream_input_sample_rate_ = 0; + stream_input_channels_ = 0; +} + bool ffmpegSpeaker::initPulseDevice_() { // 验证参数 if (sample_rate_ <= 0 || channels_ <= 0) { diff --git a/cmvr-es/devices/state_define.h b/cmvr-es/devices/state_define.h index 11b60ad7..6f5e0bee 100644 --- a/cmvr-es/devices/state_define.h +++ b/cmvr-es/devices/state_define.h @@ -6,6 +6,7 @@ #define CMVR_ES_STATE_DEFINE_H #include +#include #include #include #include @@ -39,6 +40,24 @@ namespace cmvr::device{ } AGVState; + enum class AudioStreamFormat { + PCM = 0, + MP3 = 1, + AAC = 2, + WAV = 3, + UNKNOWN = 99 + }; + + struct AudioStreamFrameData { + std::vector data; + int sample_rate = 44100; + int channels = 2; + AudioStreamFormat format = AudioStreamFormat::PCM; + std::string codec = "pcm_s16le"; + int64_t pts = 0; + int nb_samples = 0; + }; + // ------------------------------------- robot ------------------------------------- typedef enum { FORWARD, BACKWARD, diff --git a/cmvr-es/service/grpc/include/grpc_microphone_service.h b/cmvr-es/service/grpc/include/grpc_microphone_service.h index 33420240..1e4eac04 100644 --- a/cmvr-es/service/grpc/include/grpc_microphone_service.h +++ b/cmvr-es/service/grpc/include/grpc_microphone_service.h @@ -22,6 +22,7 @@ namespace cmvr::service grpc::Status StopRecord(grpc::ServerContext* context, const api::StopMicRecordingCommand_Request* request,api::StopMicRecordingCommand_Feedback* response) override; grpc::Status PauseRecord(grpc::ServerContext* context, const api::PauseMicRecordingCommand_Request* request,api::PauseMicRecordingCommand_Feedback* response) override; grpc::Status ResumeRecord(grpc::ServerContext* context, const api::ResumeMicRecordingCommand_Request* request,api::ResumeMicRecordingCommand_Feedback* response) override; + grpc::Status StreamAudio(grpc::ServerContext* context, const api::StreamMicAudioCommand_Request* request, grpc::ServerWriter* writer) override; grpc::Status SetVolume(grpc::ServerContext* context, const api::SetMicPhoneVolumeCommand_Request* request,api::SetMicPhoneVolumeCommand_Feedback* response) override; grpc::Status GetVolume(grpc::ServerContext* context, const api::GetMicPhoneVolumeCommand_Request* request,api::GetMicPhoneVolumeCommand_Feedback* response) override; private: diff --git a/cmvr-es/service/grpc/include/grpc_speaker_service.h b/cmvr-es/service/grpc/include/grpc_speaker_service.h index f2f61256..2bf0e4c1 100644 --- a/cmvr-es/service/grpc/include/grpc_speaker_service.h +++ b/cmvr-es/service/grpc/include/grpc_speaker_service.h @@ -17,6 +17,7 @@ namespace cmvr::service { ~gRPCSpeakerServiceImpl() override = default; grpc::Status GetStatus(grpc::ServerContext* context, const api::GetSpeakerStateCommand_Request* request,api::GetSpeakerStateCommand_Feedback* response) override; grpc::Status PlayAudio(grpc::ServerContext* context, const api::PlayAudioCommand_Request* request,api::PlayAudioCommand_Feedback* response) override; + grpc::Status StreamAudio(grpc::ServerContext* context, grpc::ServerReader* reader, api::StreamSpeakerAudioCommand_Feedback* response) override; grpc::Status StopPlayback(grpc::ServerContext* context, const api::StopSpeakerCommand_Request* request,api::StopSpeakerCommand_Feedback* response) override; grpc::Status PausePlayback(grpc::ServerContext* context, const api::PauseSpeakerCommand_Request* request,api::PauseSpeakerCommand_Feedback* response) override; grpc::Status ResumePlayback(grpc::ServerContext* context, const api::ResumeSpeakerCommand_Request* request,api::ResumeSpeakerCommand_Feedback* response) override; diff --git a/cmvr-es/service/grpc/src/grpc_microphone_service.cpp b/cmvr-es/service/grpc/src/grpc_microphone_service.cpp index 8bab8602..92dc9fee 100644 --- a/cmvr-es/service/grpc/src/grpc_microphone_service.cpp +++ b/cmvr-es/service/grpc/src/grpc_microphone_service.cpp @@ -1,4 +1,6 @@ #include "common/base/logging/logger.h" +#include +#include // // Created by linbo on 2025/6/13. // Created by xtkuang on 2025/6/13. @@ -18,6 +20,31 @@ grpc::Status failResponse(ResponseT* response, const std::string& message) { setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); return grpc::Status::OK; } + +cmvr::api::AudioData_AudioFormat toProtoAudioFormat(AudioStreamFormat format) { + switch (format) { + case AudioStreamFormat::PCM: + return cmvr::api::AudioData_AudioFormat_PCM; + case AudioStreamFormat::MP3: + return cmvr::api::AudioData_AudioFormat_MP3; + case AudioStreamFormat::AAC: + return cmvr::api::AudioData_AudioFormat_AAC; + case AudioStreamFormat::WAV: + return cmvr::api::AudioData_AudioFormat_WAV; + default: + return cmvr::api::AudioData_AudioFormat_PCM; + } +} + +void fillProtoAudioData(const AudioStreamFrameData& frame, cmvr::api::AudioData* audio) { + audio->set_data(frame.data.data(), frame.data.size()); + audio->set_sample_rate(frame.sample_rate); + audio->set_channels(frame.channels); + audio->set_format(toProtoAudioFormat(frame.format)); + audio->set_codec(frame.codec); + audio->set_pts(frame.pts); + audio->set_nb_samples(frame.nb_samples); +} } gRPCMicroPhoneServiceImpl::gRPCMicroPhoneServiceImpl(): dmgr_(DeviceManager::getInstance()) {} @@ -151,6 +178,61 @@ grpc::Status gRPCMicroPhoneServiceImpl::ResumeRecord(grpc::ServerContext* contex } } +grpc::Status gRPCMicroPhoneServiceImpl::StreamAudio(grpc::ServerContext* context, + const api::StreamMicAudioCommand_Request* request, + grpc::ServerWriter* writer) { + const string dev_id = request->header().device_id(); + CMVR_LOG(INFO) << "[gRPCMicroPhoneServiceImpl] (StreamAudio): id=" << dev_id; + const auto dev = dmgr_.getDevice(dev_id); + if (!dev) { + api::StreamMicAudioCommand_Feedback feedback; + feedback.mutable_header()->set_success(false); + feedback.mutable_header()->set_error_message("Microphone device not found: " + dev_id); + setCurrentTimestamp(feedback.mutable_header()->mutable_timestamp()); + writer->Write(feedback); + return grpc::Status::OK; + } + + if (!dev->start()) { + api::StreamMicAudioCommand_Feedback feedback; + feedback.mutable_header()->set_success(false); + feedback.mutable_header()->set_error_message("Failed to start microphone: " + dev_id); + setCurrentTimestamp(feedback.mutable_header()->mutable_timestamp()); + writer->Write(feedback); + return grpc::Status::OK; + } + + if (!dev->startStreaming()) { + api::StreamMicAudioCommand_Feedback feedback; + feedback.mutable_header()->set_success(false); + feedback.mutable_header()->set_error_message("Failed to start microphone streaming: " + dev_id); + setCurrentTimestamp(feedback.mutable_header()->mutable_timestamp()); + writer->Write(feedback); + return grpc::Status::OK; + } + + size_t index = 0; + while (!context->IsCancelled()) { + AudioStreamFrameData frame; + dev->getEncodedFrame(frame, index); + if (frame.data.empty()) { + std::this_thread::sleep_for(std::chrono::milliseconds(2)); + continue; + } + + api::StreamMicAudioCommand_Feedback feedback; + feedback.mutable_header()->set_success(true); + setCurrentTimestamp(feedback.mutable_header()->mutable_timestamp()); + fillProtoAudioData(frame, feedback.mutable_audio()); + if (!writer->Write(feedback)) { + break; + } + } + + dev->stopStreaming(); + return grpc::Status::OK; +} + grpc::Status gRPCMicroPhoneServiceImpl::SetVolume(grpc::ServerContext* context, const api::SetMicPhoneVolumeCommand_Request* request, api::SetMicPhoneVolumeCommand_Feedback* response) { try { diff --git a/cmvr-es/service/grpc/src/grpc_speaker_service.cpp b/cmvr-es/service/grpc/src/grpc_speaker_service.cpp index 106ea389..756fab83 100644 --- a/cmvr-es/service/grpc/src/grpc_speaker_service.cpp +++ b/cmvr-es/service/grpc/src/grpc_speaker_service.cpp @@ -1,4 +1,5 @@ #include "common/base/logging/logger.h" +#include // // Created by xtkuang on 2025/6/10. // @@ -18,6 +19,33 @@ grpc::Status failResponse(ResponseT* response, const std::string& message) { setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); return grpc::Status::OK; } + +AudioStreamFormat fromProtoAudioFormat(cmvr::api::AudioData_AudioFormat format) { + switch (format) { + case cmvr::api::AudioData_AudioFormat_PCM: + return AudioStreamFormat::PCM; + case cmvr::api::AudioData_AudioFormat_MP3: + return AudioStreamFormat::MP3; + case cmvr::api::AudioData_AudioFormat_AAC: + return AudioStreamFormat::AAC; + case cmvr::api::AudioData_AudioFormat_WAV: + return AudioStreamFormat::WAV; + default: + return AudioStreamFormat::UNKNOWN; + } +} + +AudioStreamFrameData fromProtoAudioData(const cmvr::api::AudioData& audio) { + AudioStreamFrameData frame; + frame.data.assign(audio.data().begin(), audio.data().end()); + frame.sample_rate = audio.sample_rate() > 0 ? audio.sample_rate() : 44100; + frame.channels = audio.channels() > 0 ? audio.channels() : 2; + frame.format = fromProtoAudioFormat(audio.format()); + frame.codec = audio.codec(); + frame.pts = audio.pts(); + frame.nb_samples = audio.nb_samples(); + return frame; +} } gRPCSpeakerServiceImpl::gRPCSpeakerServiceImpl(): dmgr_(DeviceManager::getInstance()) {} @@ -81,6 +109,48 @@ grpc::Status gRPCSpeakerServiceImpl::PlayAudio(grpc::ServerContext* context, } } +grpc::Status gRPCSpeakerServiceImpl::StreamAudio(grpc::ServerContext* context, + grpc::ServerReader* reader, + api::StreamSpeakerAudioCommand_Feedback* response) { + try { + api::StreamSpeakerAudioCommand_Request request; + std::shared_ptr dev; + std::string dev_id; + + while (reader->Read(&request)) { + if (!dev) { + dev_id = request.header().device_id(); + CMVR_LOG(INFO) << "[gRPCSpeakerServiceImpl] (StreamAudio): id=" << dev_id; + dev = dmgr_.getDevice(dev_id); + if (!dev) { + return failResponse(response, "Speaker device not found: " + dev_id); + } + if (!dev->start()) { + return failResponse(response, "Failed to start speaker: " + dev_id); + } + } + + if (!dev->pushAudioFrame(fromProtoAudioData(request.audio()))) { + dev->stopStreaming(); + return failResponse(response, "Failed to push speaker audio frame: " + dev_id); + } + } + + if (dev) { + dev->stopStreaming(); + } + response->mutable_header()->set_success(true); + setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); + return grpc::Status::OK; + } + catch (const std::exception& e) { + response->mutable_header()->set_success(false); + response->mutable_header()->set_error_message(e.what()); + setCurrentTimestamp(response->mutable_header()->mutable_timestamp()); + return grpc::Status::OK; + } +} + grpc::Status gRPCSpeakerServiceImpl::StopPlayback(grpc::ServerContext* context, const api::StopSpeakerCommand_Request* request, api::StopSpeakerCommand_Feedback* response) { try { diff --git a/protos/cmvr/api/common.proto b/protos/cmvr/api/common.proto index accb9361..53ecd55b 100644 --- a/protos/cmvr/api/common.proto +++ b/protos/cmvr/api/common.proto @@ -29,6 +29,22 @@ message CommandHeader { } } +message AudioData { + enum AudioFormat { + PCM = 0; + MP3 = 1; + AAC = 2; + WAV = 3; + } + bytes data = 1; + int32 sample_rate = 2; + int32 channels = 3; + AudioFormat format = 4; + string codec = 5; + int64 pts = 6; + int32 nb_samples = 7; +} + message ConfigParam { string param_name = 1; diff --git a/protos/cmvr/api/microphone_command.proto b/protos/cmvr/api/microphone_command.proto index 4183415d..71a6f4ec 100644 --- a/protos/cmvr/api/microphone_command.proto +++ b/protos/cmvr/api/microphone_command.proto @@ -11,6 +11,7 @@ message MicState { int32 volume = 4; string error_message = 5; } + message GetMicStateCommand { message Request { CommandHeader.Request header = 1; @@ -21,6 +22,7 @@ message GetMicStateCommand { MicState state = 2; } } + message StartMicRecordingCommand { message Request { CommandHeader.Request header = 1; @@ -59,10 +61,20 @@ message ResumeMicRecordingCommand { } } +message StreamMicAudioCommand { + message Request { + CommandHeader.Request header = 1; + } + message Feedback { + CommandHeader.Feedback header = 1; + AudioData audio = 2; + } +} + message SetMicPhoneVolumeCommand { message Request { CommandHeader.Request header = 1; - int32 volume = 2; // 音量值(0 ~ 100) + int32 volume = 2; } message Feedback { CommandHeader.Feedback header = 1; @@ -75,7 +87,6 @@ message GetMicPhoneVolumeCommand { } message Feedback { CommandHeader.Feedback header = 1; - int32 volume = 2; // 当前音量值 + int32 volume = 2; } } - diff --git a/protos/cmvr/api/microphone_service.proto b/protos/cmvr/api/microphone_service.proto index fcf31b79..d9cea818 100644 --- a/protos/cmvr/api/microphone_service.proto +++ b/protos/cmvr/api/microphone_service.proto @@ -4,17 +4,14 @@ import "cmvr/api/microphone_command.proto"; package cmvr.api; - service MicPhoneService { - // 基本控制 rpc GetStatus(GetMicStateCommand.Request) returns (GetMicStateCommand.Feedback); rpc StartRecord(StartMicRecordingCommand.Request) returns (StartMicRecordingCommand.Feedback); rpc StopRecord(StopMicRecordingCommand.Request) returns (StopMicRecordingCommand.Feedback); rpc PauseRecord(PauseMicRecordingCommand.Request) returns (PauseMicRecordingCommand.Feedback); rpc ResumeRecord(ResumeMicRecordingCommand.Request) returns (ResumeMicRecordingCommand.Feedback); + rpc StreamAudio(StreamMicAudioCommand.Request) returns (stream StreamMicAudioCommand.Feedback); - // 音量控制 rpc SetVolume(SetMicPhoneVolumeCommand.Request) returns (SetMicPhoneVolumeCommand.Feedback); rpc GetVolume(GetMicPhoneVolumeCommand.Request) returns (GetMicPhoneVolumeCommand.Feedback); - -} \ No newline at end of file +} diff --git a/protos/cmvr/api/speaker_command.proto b/protos/cmvr/api/speaker_command.proto index f5a23b8d..6243b6d5 100644 --- a/protos/cmvr/api/speaker_command.proto +++ b/protos/cmvr/api/speaker_command.proto @@ -4,33 +4,15 @@ import "cmvr/api/common.proto"; package cmvr.api; - -// 音频数据格式(与麦克风类似,但用于输出) -message AudioData { - enum AudioFormat { - PCM = 0; - MP3 = 1; - AAC = 2; - WAV = 3; - } - bytes data = 1; // 音频二进制数据 - int32 sample_rate = 2; // 采样率(Hz) - int32 channels = 3; // 声道数 - AudioFormat format = 4; // 音频格式 - string codec = 5; // 编码方式 -} - -// 扬声器状态 message SpeakerState { - bool is_initialized = 1; // 是否已初始化 - bool is_running = 2; // 是否正在播放 + bool is_initialized = 1; + bool is_running = 2; bool is_decoding = 3; bool is_paused = 5; - int32 volume = 6; // 当前音量(0 ~ 100) - string error_message = 7; // 错误信息 + int32 volume = 6; + string error_message = 7; } -// 各种扬声器命令 message GetSpeakerStateCommand { message Request { CommandHeader.Request header = 1; @@ -47,7 +29,19 @@ message PlayAudioCommand { CommandHeader.Request header = 1; string audio_path = 2; } - message Feedback { CommandHeader.Feedback header = 1; } + message Feedback { + CommandHeader.Feedback header = 1; + } +} + +message StreamSpeakerAudioCommand { + message Request { + CommandHeader.Request header = 1; + AudioData audio = 2; + } + message Feedback { + CommandHeader.Feedback header = 1; + } } message StopSpeakerCommand { @@ -80,7 +74,7 @@ message ResumeSpeakerCommand { message SetSpeakerVolumeCommand { message Request { CommandHeader.Request header = 1; - int32 volume = 2; // 音量值(0 ~ 100) + int32 volume = 2; } message Feedback { CommandHeader.Feedback header = 1; @@ -93,7 +87,6 @@ message GetSpeakerVolumeCommand { } message Feedback { CommandHeader.Feedback header = 1; - int32 volume = 2; // 当前音量值 + int32 volume = 2; } } - diff --git a/protos/cmvr/api/speaker_service.proto b/protos/cmvr/api/speaker_service.proto index 9cbc77cd..6fc3d927 100644 --- a/protos/cmvr/api/speaker_service.proto +++ b/protos/cmvr/api/speaker_service.proto @@ -4,17 +4,14 @@ import "cmvr/api/speaker_command.proto"; package cmvr.api; - service SpeakerService { - // 基本控制 rpc GetStatus(GetSpeakerStateCommand.Request) returns (GetSpeakerStateCommand.Feedback); rpc PlayAudio(PlayAudioCommand.Request) returns (PlayAudioCommand.Feedback); + rpc StreamAudio(stream StreamSpeakerAudioCommand.Request) returns (StreamSpeakerAudioCommand.Feedback); rpc StopPlayback(StopSpeakerCommand.Request) returns (StopSpeakerCommand.Feedback); rpc PausePlayback(PauseSpeakerCommand.Request) returns (PauseSpeakerCommand.Feedback); rpc ResumePlayback(ResumeSpeakerCommand.Request) returns (ResumeSpeakerCommand.Feedback); - // 音量控制 rpc SetVolume(SetSpeakerVolumeCommand.Request) returns (SetSpeakerVolumeCommand.Feedback); rpc GetVolume(GetSpeakerVolumeCommand.Request) returns (GetSpeakerVolumeCommand.Feedback); - -} \ No newline at end of file +}