/*
 * Copyright (c) 2024 Huawei Device Co., Ltd.
 * Licensed under the Apache License, Version 2.0 (the "License");
 * you may not use this file except in compliance with the License.
 * You may obtain a copy of the License at
 *
 *     http://www.apache.org/licenses/LICENSE-2.0
 *
 * Unless required by applicable law or agreed to in writing, software
 * distributed under the License is distributed on an "AS IS" BASIS,
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 * See the License for the specific language governing permissions and
 * limitations under the License.
 */

#include "dmic_dev.h"

#include <condition_variable>
#include <ctime>
#include <mutex>
#include <string>
#include <thread>

#include "daudio_constants.h"
#include "daudio_errorcode.h"
#include "daudio_hidumper.h"
#include "daudio_hisysevent.h"
#include "daudio_hitrace.h"
#include "daudio_log.h"
#include "daudio_radar.h"
#include "daudio_source_manager.h"
#include "daudio_util.h"

#undef DH_LOG_TAG
#define DH_LOG_TAG "DMicDev"

namespace OHOS {
namespace DistributedHardware {
static constexpr size_t DATA_QUEUE_EXT_SIZE = 20;
void DMicDev::OnEngineTransEvent(const AVTransEvent &event)
{
    if (event.type == EventType::EVENT_START_SUCCESS) {
        OnStateChange(DATA_OPENED);
    } else if ((event.type == EventType::EVENT_STOP_SUCCESS) ||
        (event.type == EventType::EVENT_CHANNEL_CLOSED) ||
        (event.type == EventType::EVENT_START_FAIL)) {
        OnStateChange(DATA_CLOSED);
    }
}

void DMicDev::OnEngineTransMessage(const std::shared_ptr<AVTransMessage> &message)
{
    CHECK_NULL_VOID(message);
    DHLOGI("On Engine message, type : %{public}s.", GetEventNameByType(message->type_).c_str());
    DAudioSourceManager::GetInstance().HandleDAudioNotify(message->dstDevId_, message->dstDevId_,
        message->type_, message->content_);
}

void DMicDev::OnEngineTransDataAvailable(const std::shared_ptr<AudioData> &audioData)
{
    std::lock_guard<std::mutex> lock(ringbufferMutex_);
    CHECK_NULL_VOID(ringBuffer_);
    CHECK_NULL_VOID(audioData);
    CHECK_AND_RETURN_LOG(ringBuffer_->RingBufferInsert(audioData->Data(),
        static_cast<int32_t>(audioData->Capacity())) != DH_SUCCESS, "RingBufferInsert failed.");
    DHLOGD("Ringbuffer insert one");
    int64_t timestamp = audioData->GetPts();
    std::lock_guard<std::mutex> timeLock(ptsMutex_);
    ptsMap_[frameInIndex_] = timestamp;
    frameInIndex_++;
    if (frameInIndex_ == indexFlag_) {
        timestamp = audioData->GetPtsSpecial();
        ptsMap_[frameInIndex_] = timestamp;
        DHLOGD("audioData->GetPtsSpecial(): %{public}" PRId64", index: %{public}" PRIu64, timestamp, frameInIndex_);
        frameInIndex_ = 0;
    }
}

void DMicDev::ReadFromRingbuffer()
{
    std::shared_ptr<AudioData> sendData = std::make_shared<AudioData>(frameSize_);
    bool canRead = false;
    while (isRingbufferOn_.load()) {
        CHECK_NULL_VOID(ringBuffer_);
        canRead = false;
        {
            std::lock_guard<std::mutex> lock(ringbufferMutex_);
            if (ringBuffer_->CanBufferReadLen(frameSize_)) {
                canRead = true;
            }
        }
        if (!canRead) {
            DHLOGD("Can not read from ringbuffer.");
            std::this_thread::sleep_for(std::chrono::milliseconds(RINGBUFFER_WAIT_SECONDS));
            continue;
        }
        {
            std::lock_guard<std::mutex> lock(ringbufferMutex_);
            if (ringBuffer_->RingBufferGetData(sendData->Data(), sendData->Capacity()) != DH_SUCCESS) {
                DHLOGE("Read ringbuffer failed.");
                continue;
            }
        }
        SendToProcess(sendData);
    }
}

void DMicDev::SendToProcess(const std::shared_ptr<AudioData> &audioData)
{
    DHLOGD("On Engine Data available");
    CHECK_NULL_VOID(audioData);
    int64_t pts = 0;
    {
        std::lock_guard<std::mutex> lock(ptsMutex_);
        auto iter = ptsMap_.find(frameOutIndex_);
        if (iter != ptsMap_.end()) {
            pts = iter->second;
            audioData->SetPts(pts);
            ptsMap_.erase(iter);
        } else {
            DHLOGI("iter == ptsMap_.end()");
        }
        frameOutIndex_++;
        if (frameOutIndex_ == frameOutIndexFlag_) {
            frameOutIndex_ = 0;
        }
    }
    DHLOGD("Set timestamp %{public}" PRId64 " for frame index %{public}" PRIu64, audioData->GetPts(), frameOutIndex_);
    framnum_++;
    DHLOGD("current frame index: %{public}" PRIu64, framnum_);
    if (echoCannelOn_) {
#ifdef ECHO_CANNEL_ENABLE
        CHECK_NULL_VOID(echoManager_);
        echoManager_->OnMicDataReceived(audioData);
#endif
    } else {
        OnDecodeTransDataDone(audioData);
    }
}

int32_t DMicDev::InitReceiverEngine(IAVEngineProvider *providerPtr)
{
    DHLOGI("InitReceiverEngine enter.");
    if (micTrans_ == nullptr) {
        micTrans_ = std::make_shared<AVTransReceiverTransport>(devId_, shared_from_this());
    }
    int32_t ret = micTrans_->InitEngine(providerPtr);
    if (ret != DH_SUCCESS) {
        DHLOGE("Mic dev initialize av receiver adapter failed.");
        return ret;
    }
    return DH_SUCCESS;
}

int32_t DMicDev::InitSenderEngine(IAVEngineProvider *providerPtr)
{
    DHLOGI("InitReceiverEngine enter.");
    return DH_SUCCESS;
}

int32_t DMicDev::InitCtrlTrans()
{
    DHLOGI("InitCtrlTrans enter");
    if (micCtrlTrans_ == nullptr) {
        micCtrlTrans_ = std::make_shared<DaudioSourceCtrlTrans>(devId_,
            SESSIONNAME_MIC_SOURCE, SESSIONNAME_MIC_SINK, shared_from_this());
    }
    int32_t ret = micCtrlTrans_->SetUp(shared_from_this());
    CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS, ret, "Mic ctrl SetUp failed.");
    ret = micCtrlTrans_->Start();
    CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS, ret, "Mic ctrl Start failed.");
    return ret;
}

void DMicDev::OnCtrlTransEvent(const AVTransEvent &event)
{
    if (event.type == EventType::EVENT_START_SUCCESS) {
        OnStateChange(DATA_OPENED);
    } else if ((event.type == EventType::EVENT_STOP_SUCCESS) ||
        (event.type == EventType::EVENT_CHANNEL_CLOSED) ||
        (event.type == EventType::EVENT_START_FAIL)) {
        OnStateChange(DATA_CLOSED);
    }
}

void DMicDev::OnCtrlTransMessage(const std::shared_ptr<AVTransMessage> &message)
{
    CHECK_NULL_VOID(message);
    DHLOGI("On Engine message, type : %{public}s.", GetEventNameByType(message->type_).c_str());
    DAudioSourceManager::GetInstance().HandleDAudioNotify(message->dstDevId_, message->dstDevId_,
        message->type_, message->content_);
}

int32_t DMicDev::EnableDevice(const int32_t dhId, const std::string &capability)
{
    DHLOGI("Enable IO device, device pin: %{public}d.", dhId);
    int32_t ret = DAudioHdiHandler::GetInstance().RegisterAudioDevice(devId_, dhId, capability, shared_from_this());
    if (ret != DH_SUCCESS) {
        DHLOGE("Register device failed, ret: %{public}d.", ret);
        DAudioHisysevent::GetInstance().SysEventWriteFault(DAUDIO_REGISTER_FAIL, devId_, std::to_string(dhId), ret,
            "daudio register device failed.");
        return ret;
    }
    dhId_ = dhId;
    GetCodecCaps(capability);
    return DH_SUCCESS;
}

void DMicDev::AddToVec(std::vector<AudioCodecType> &container, const AudioCodecType value)
{
    auto it = std::find(container.begin(), container.end(), value);
    if (it == container.end()) {
        container.push_back(value);
    }
}

void DMicDev::GetCodecCaps(const std::string &capability)
{
    auto pos = capability.find(AAC);
    if (pos != std::string::npos) {
        AddToVec(codec_, AudioCodecType::AUDIO_CODEC_AAC_EN);
        DHLOGI("Daudio codec cap: AAC");
    }
    pos = capability.find(OPUS);
    if (pos != std::string::npos) {
        AddToVec(codec_, AudioCodecType::AUDIO_CODEC_OPUS);
        DHLOGI("Daudio codec cap: OPUS");
    }
}

int32_t DMicDev::DisableDevice(const int32_t dhId)
{
    DHLOGI("Disable IO device, device pin: %{public}d.", dhId);
    int32_t ret = DAudioHdiHandler::GetInstance().UnRegisterAudioDevice(devId_, dhId);
    if (ret != DH_SUCCESS) {
        DHLOGE("UnRegister failed, ret: %{public}d.", ret);
        DAudioHisysevent::GetInstance().SysEventWriteFault(DAUDIO_UNREGISTER_FAIL, devId_, std::to_string(dhId), ret,
            "daudio unregister device failed.");
        return ret;
    }
    return DH_SUCCESS;
}

bool DMicDev::IsMimeSupported(const AudioCodecType coder)
{
    auto iter = std::find(codec_.begin(), codec_.end(), coder);
    if (iter == codec_.end()) {
        DHLOGI("devices have no cap: %{public}d", static_cast<int>(coder));
        return false;
    }
    return true;
}

int32_t DMicDev::CreateStream(const int32_t streamId)
{
    DHLOGI("Open stream of mic device streamId: %{public}d.", streamId);
    std::shared_ptr<IAudioEventCallback> cbObj = audioEventCallback_.lock();
    CHECK_NULL_RETURN(cbObj, ERR_DH_AUDIO_NULLPTR);

    cJSON *jParam = cJSON_CreateObject();
    CHECK_NULL_RETURN(jParam, ERR_DH_AUDIO_NULLPTR);
    cJSON_AddStringToObject(jParam, KEY_DH_ID, std::to_string(dhId_).c_str());
    if (triggerFirstTokenId_ > 0) {
        cJSON_AddNumberToObject(jParam, KEY_TRIGGER_FIRST_TOKENID, static_cast<double>(triggerFirstTokenId_));
        DHLOGI("[MultiUserTrigger] DMicDev::CreateStream add triggerFirstTokenId=%{public}s to event",
            GetAnonyString(std::to_string(triggerFirstTokenId_)).c_str());
    }
    char *jsonData = cJSON_PrintUnformatted(jParam);
    if (jsonData == nullptr) {
        cJSON_Delete(jParam);
        DHLOGE("Failed to create JSON data.");
        return ERR_DH_AUDIO_NULLPTR;
    }
    std::string jsonDataStr(jsonData);
    AudioEvent event(AudioEventType::OPEN_MIC, jsonDataStr);
    cbObj->NotifyEvent(event);
    DAudioHisysevent::GetInstance().SysEventWriteBehavior(DAUDIO_OPEN, devId_, std::to_string(dhId_),
        "daudio mic device open success.");
    streamId_ = streamId;
    cJSON_Delete(jParam);
    cJSON_free(jsonData);
    DaudioRadar::GetInstance().ReportMicOpen("Start", MicOpen::CREATE_STREAM,
        BizState::BIZ_STATE_START, DH_SUCCESS);
    return DH_SUCCESS;
}

int32_t DMicDev::DestroyStream(const int32_t streamId)
{
    DHLOGI("Close stream of mic device streamId: %{public}d.", streamId);
    std::shared_ptr<IAudioEventCallback> cbObj = audioEventCallback_.lock();
    CHECK_NULL_RETURN(cbObj, ERR_DH_AUDIO_NULLPTR);

    cJSON *jParam = cJSON_CreateObject();
    CHECK_NULL_RETURN(jParam, ERR_DH_AUDIO_NULLPTR);
    cJSON_AddStringToObject(jParam, KEY_DH_ID, std::to_string(dhId_).c_str());
    char *jsonData = cJSON_PrintUnformatted(jParam);
    if (jsonData == nullptr) {
        cJSON_Delete(jParam);
        DHLOGE("Failed to create JSON data.");
        return ERR_DH_AUDIO_NULLPTR;
    }
    std::string jsonDataStr(jsonData);
    AudioEvent event(AudioEventType::CLOSE_MIC, jsonDataStr);
    cbObj->NotifyEvent(event);
    DAudioHisysevent::GetInstance().SysEventWriteBehavior(DAUDIO_CLOSE, devId_, std::to_string(dhId_),
        "daudio mic device close success.");
    cJSON_Delete(jParam);
    cJSON_free(jsonData);
    curPort_ = 0;
    DaudioRadar::GetInstance().ReportMicClose("DestroyStream", MicClose::DESTROY_STREAM,
        BizState::BIZ_STATE_START, DH_SUCCESS);
    return DH_SUCCESS;
}

int32_t DMicDev::SetParameters(const int32_t streamId, const AudioParamHDF &param)
{
    DHLOGD("Set mic parameters {samplerate: %{public}d, channelmask: %{public}d, format: %{public}d, "
        "period: %{public}d, framesize: %{public}d, ext{%{public}s}}.", param.sampleRate,
        param.channelMask, param.bitFormat, param.period, param.frameSize, param.ext.c_str());
    if (param.capturerFlags == MMAP_MODE && param.period != MMAP_NORMAL_PERIOD && param.period != MMAP_VOIP_PERIOD) {
        DHLOGE("The period is invalid : %{public}" PRIu32, param.period);
        return ERR_DH_AUDIO_SA_PARAM_INVALID;
    }
    curPort_ = dhId_;
    paramHDF_ = param;

    ParseTriggerFirstTokenIdFromExt(param.ext);

    param_.comParam.sampleRate = paramHDF_.sampleRate;
    param_.comParam.channelMask = paramHDF_.channelMask;
    param_.comParam.bitFormat = paramHDF_.bitFormat;
    param_.comParam.frameSize = paramHDF_.frameSize;
    if (paramHDF_.streamUsage == StreamUsage::STREAM_USAGE_VOICE_COMMUNICATION) {
        param_.captureOpts.sourceType = SOURCE_TYPE_VOICE_COMMUNICATION;
    } else {
        param_.captureOpts.sourceType = SOURCE_TYPE_MIC;
    }
    param_.captureOpts.capturerFlags = paramHDF_.capturerFlags;
    if (paramHDF_.capturerFlags == MMAP_MODE) {
        lowLatencyHalfSize_ = LOW_LATENCY_JITTER_TIME_MS / paramHDF_.period;
        lowLatencyMaxfSize_ = LOW_LATENCY_JITTER_MAX_TIME_MS / paramHDF_.period;
    }
    param_.comParam.codecType = AudioCodecType::AUDIO_CODEC_AAC;
    if (paramHDF_.streamUsage == StreamUsage::STREAM_USAGE_VOICE_COMMUNICATION &&
        IsMimeSupported(AudioCodecType::AUDIO_CODEC_OPUS)) {
        param_.comParam.codecType = AudioCodecType::AUDIO_CODEC_OPUS;
    } else if (IsMimeSupported(AudioCodecType::AUDIO_CODEC_AAC_EN)) {
        param_.comParam.codecType = AudioCodecType::AUDIO_CODEC_AAC_EN;
    }
    DHLOGI("codecType : %{public}d.", static_cast<int>(param_.comParam.codecType));
    return DH_SUCCESS;
}

void DMicDev::ParseTriggerFirstTokenIdFromExt(const std::string &ext)
{
    if (!ext.empty()) {
        std::string extStr = ext;
        std::string tokenKey = TRIGGER_FIRST_TOKENID_PREFIX;
        size_t pos = extStr.find(tokenKey);
        if (pos != std::string::npos) {
            std::string tokenIdStr = extStr.substr(pos + tokenKey.length());
            size_t endPos = tokenIdStr.find_first_of(EXT_PARAM_DELIMITERS);
            if (endPos != std::string::npos) {
                tokenIdStr = tokenIdStr.substr(0, endPos);
            }
            errno = 0;
            char *endPtr = nullptr;
            unsigned long tokenIdVal = strtoul(tokenIdStr.c_str(), &endPtr, DECIMAL_BASE);
            if (endPtr != tokenIdStr.c_str() && *endPtr == '\0' && errno == 0) {
                triggerFirstTokenId_ = static_cast<uint32_t>(tokenIdVal);
                DHLOGI("[MultiUserTrigger] DMicDev::SetParameters parsed triggerFirstTokenId=%{public}s from ext",
                    GetAnonyString(std::to_string(triggerFirstTokenId_)).c_str());
            }
        }
    }
}

int32_t DMicDev::NotifyEvent(const int32_t streamId, const AudioEvent &event)
{
    DHLOGD("Notify mic event, type: %{public}d.", event.type);
    std::shared_ptr<IAudioEventCallback> cbObj = audioEventCallback_.lock();
    CHECK_NULL_RETURN(cbObj, ERR_DH_AUDIO_NULLPTR);
    switch (event.type) {
        case AudioEventType::AUDIO_START:
            curStatus_ = AudioStatus::STATUS_START;
            isExistedEmpty_.store(false);
            isStartStatus_.store(true);
            break;
        case AudioEventType::AUDIO_STOP:
            curStatus_ = AudioStatus::STATUS_STOP;
            isExistedEmpty_.store(false);
            break;
        default:
            break;
    }
    AudioEvent audioEvent(event.type, event.content);
    cbObj->NotifyEvent(audioEvent);
    return DH_SUCCESS;
}

int32_t DMicDev::SetUp()
{
    DHLOGI("Set up mic device.");
    CHECK_NULL_RETURN(micTrans_, ERR_DH_AUDIO_NULLPTR);
    int32_t ret = micTrans_->SetUp(param_, param_, shared_from_this(), CAP_MIC);
    if (ret != DH_SUCCESS) {
        DHLOGE("Mic trans set up failed. ret: %{public}d.", ret);
        return ret;
    }
    frameSize_ = static_cast<int32_t>(param_.comParam.frameSize);
    {
        std::lock_guard<std::mutex> lock(ringbufferMutex_);
        ringBuffer_ = std::make_unique<DaudioRingBuffer>();
        ringBuffer_->RingBufferInit(frameData_);
        CHECK_NULL_RETURN(frameData_, ERR_DH_AUDIO_NULLPTR);
    }
    isRingbufferOn_.store(true);
    ringbufferThread_ = std::thread([ptr = shared_from_this()]() {
        if (!ptr) {
            return;
        }
        ptr->ReadFromRingbuffer();
    });
    echoCannelOn_ = false;
    DHLOGI("echoCannelOn_: %{public}d", echoCannelOn_);
#ifdef ECHO_CANNEL_ENABLE
    if (echoCannelOn_ && echoManager_ == nullptr) {
        echoManager_ = std::make_shared<DAudioEchoCannelManager>();
    }
    AudioCommonParam info;
    info.sampleRate = param_.comParam.sampleRate;
    info.channelMask = param_.comParam.channelMask;
    info.bitFormat = param_.comParam.bitFormat;
    info.frameSize = param_.comParam.frameSize;
    if (echoManager_ != nullptr) {
        echoManager_->SetUp(info, shared_from_this());
    }
#endif
    DumpFileUtil::OpenDumpFile(DUMP_SERVER_PARA, DUMP_DAUDIO_MIC_READ_FROM_BUF_NAME, &dumpFileCommn_);
    DumpFileUtil::OpenDumpFile(DUMP_SERVER_PARA, DUMP_DAUDIO_LOWLATENCY_MIC_FROM_BUF_NAME, &dumpFileFast_);
    return DH_SUCCESS;
}

int32_t DMicDev::Start()
{
    DHLOGI("Start mic device.");
    CHECK_NULL_RETURN(micTrans_, ERR_DH_AUDIO_NULLPTR);
    int32_t ret = micTrans_->Start();
    DaudioRadar::GetInstance().ReportMicOpenProgress("Start", MicOpen::TRANS_START, ret);
    if (ret != DH_SUCCESS) {
        DHLOGE("Mic trans start failed, ret: %{public}d.", ret);
        return ret;
    }
    std::unique_lock<std::mutex> lck(channelWaitMutex_);
    auto status = channelWaitCond_.wait_for(lck, std::chrono::seconds(CHANNEL_WAIT_SECONDS),
        [this]() { return isTransReady_.load(); });
    if (!status) {
        DHLOGE("Wait channel open timeout(%{public}ds).", CHANNEL_WAIT_SECONDS);
        return ERR_DH_AUDIO_SA_WAIT_TIMEOUT;
    }
    isOpened_.store(true);
    return DH_SUCCESS;
}

int32_t DMicDev::Pause()
{
    DHLOGI("Not support.");
    return DH_SUCCESS;
}

int32_t DMicDev::Restart()
{
    DHLOGI("Not surpport.");
    return DH_SUCCESS;
}

int32_t DMicDev::Stop()
{
    DHLOGI("Stop mic device.");
    CHECK_NULL_RETURN(micTrans_, DH_SUCCESS);
    isOpened_.store(false);
    isTransReady_.store(false);
    int32_t ret = micTrans_->Stop();
    if (ret != DH_SUCCESS) {
        DHLOGE("Stop mic trans failed, ret: %{public}d.", ret);
    }
    DaudioRadar::GetInstance().ReportMicCloseProgress("Stop", MicClose::STOP_TRANS, ret);
    CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS, ret, "micTrans stop failed, ret: %{public}d.", ret);
#ifdef ECHO_CANNEL_ENABLE
    CHECK_NULL_RETURN(echoManager_, DH_SUCCESS);
    ret = echoManager_->Stop();
    CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS, ret, "Echo manager stop failed, ret: %{public}d.", ret);
#endif
    return DH_SUCCESS;
}

int32_t DMicDev::Release()
{
    DHLOGI("Release mic device.");
    if (ashmem_ != nullptr) {
        ashmem_->UnmapAshmem();
        ashmem_->CloseAshmem();
        ashmem_ = nullptr;
        DHLOGI("UnInit ashmem success.");
    }
    AVsyncDeintAshmem();
    avSyncParam_.fd = -1;
    avSyncParam_.sharedMemLen = 0;
    DHLOGI("AVsync mode closed");
    {
        std::lock_guard<std::mutex> lock(dataQueueMtx_);
        dataQueue_.clear();
    }
    if (micCtrlTrans_ != nullptr) {
        int32_t res = micCtrlTrans_->Release();
        CHECK_AND_RETURN_RET_LOG(res != DH_SUCCESS, res, "Mic ctrl Release failed.");
    }
    CHECK_NULL_RETURN(micTrans_, DH_SUCCESS);
    int32_t ret = micTrans_->Release();
    DaudioRadar::GetInstance().ReportMicCloseProgress("Release", MicClose::RELEASE_TRANS, ret);
    if (ret != DH_SUCCESS) {
        DHLOGE("Release mic trans failed, ret: %{public}d.", ret);
        return ret;
    }
    isRingbufferOn_.store(false);
    if (ringbufferThread_.joinable()) {
        ringbufferThread_.join();
    }
    {
        std::lock_guard<std::mutex> lock(ringbufferMutex_);
        ringBuffer_ = nullptr;
        if (frameData_ != nullptr) {
            delete[] frameData_;
            frameData_ = nullptr;
        }
    }
#ifdef ECHO_CANNEL_ENABLE
    if (echoManager_ != nullptr) {
        echoManager_->Release();
        echoManager_ = nullptr;
    }
#endif
    DumpFileUtil::CloseDumpFile(&dumpFileCommn_);
    DumpFileUtil::CloseDumpFile(&dumpFileFast_);
    return DH_SUCCESS;
}

bool DMicDev::IsOpened()
{
    return isOpened_.load();
}

int32_t DMicDev::WriteStreamData(const int32_t streamId, std::shared_ptr<AudioData> &data)
{
    (void)streamId;
    (void)data;
    return DH_SUCCESS;
}

int32_t DMicDev::ReadTimeStampFromAVsync(int64_t &timePts)
{
    CHECK_AND_RETURN_RET_LOG(!IsAVsync() || avsyncAshmem_ == nullptr, DH_SUCCESS,
        "Ashmem is nullptr or IsAVsync is false.");
    auto syncData = avsyncAshmem_->ReadFromAshmem(avSyncParam_.sharedMemLen, 0);
    CHECK_AND_RETURN_RET_LOG(syncData == nullptr, ERR_DH_AUDIO_FAILED, "syncData is null.");
    AVsyncShareData *readSyncShareData = reinterpret_cast<AVsyncShareData *>(const_cast<void *>(syncData));
    while (!readSyncShareData->lock && IsAVsync()) {
        DHLOGE("AVsync lock is false");
        syncData = avsyncAshmem_->ReadFromAshmem(avSyncParam_.sharedMemLen, 0);
        CHECK_AND_RETURN_RET_LOG(syncData == nullptr, ERR_DH_AUDIO_FAILED, "syncData is null.");
        readSyncShareData = reinterpret_cast<AVsyncShareData *>(const_cast<void *>(syncData));
    }
    readSyncShareData->lock = 0;
    bool ret = avsyncAshmem_->WriteToAshmem(static_cast<void*>(readSyncShareData), sizeof(AVsyncShareData), 0);
    CHECK_AND_RETURN_RET_LOG(!ret, ERR_DH_AUDIO_FAILED, "Write avsync data failed.");
    timePts = static_cast<int64_t>(readSyncShareData->video_current_pts);
    readSyncShareData->lock = 1;
    ret = avsyncAshmem_->WriteToAshmem(static_cast<void*>(readSyncShareData), sizeof(AVsyncShareData), 0);
    CHECK_AND_RETURN_RET_LOG(!ret, ERR_DH_AUDIO_FAILED, "Write avsync data failed.");
    DHLOGD("readSyncShareData->audio_current_pts: %{public}" PRId64", readSyncShareData->audio_update_clock: "
        "%{public}" PRId64, readSyncShareData->audio_current_pts, readSyncShareData->audio_update_clock);
    return DH_SUCCESS;
}

int32_t DMicDev::WriteTimeStampToAVsync(const int64_t timePts)
{
    CHECK_AND_RETURN_RET_LOG(!IsAVsync() || avsyncAshmem_ == nullptr, DH_SUCCESS,
        "Ashmem is nullptr or IsAVsync is false.");
    auto syncData = avsyncAshmem_->ReadFromAshmem(avSyncParam_.sharedMemLen, 0);
    CHECK_AND_RETURN_RET_LOG(syncData == nullptr, ERR_DH_AUDIO_FAILED, "syncData is null.");
    AVsyncShareData *readSyncShareData = reinterpret_cast<AVsyncShareData *>(const_cast<void *>(syncData));
    while (!readSyncShareData->lock && IsAVsync()) {
        DHLOGE("AVsync lock is false");
        syncData = avsyncAshmem_->ReadFromAshmem(avSyncParam_.sharedMemLen, 0);
        CHECK_AND_RETURN_RET_LOG(syncData == nullptr, ERR_DH_AUDIO_FAILED, "syncData is null.");
        readSyncShareData = reinterpret_cast<AVsyncShareData *>(const_cast<void *>(syncData));
    }
    readSyncShareData->lock = 0;
    bool ret = avsyncAshmem_->WriteToAshmem(static_cast<void*>(readSyncShareData), sizeof(AVsyncShareData), 0);
    CHECK_AND_RETURN_RET_LOG(!ret, ERR_DH_AUDIO_FAILED, "Write avsync data failed.");
    struct timespec time = {0, 0};
    clock_gettime(CLOCK_REALTIME, &time);
    int64_t updatePts = static_cast<int64_t>(time.tv_sec) * TIME_CONVERSION_STOU +
        static_cast<int64_t>(time.tv_nsec) / TIME_CONVERSION_NTOU;
    readSyncShareData->lock = 1;
    readSyncShareData->audio_current_pts = static_cast<uint64_t>(timePts);
    readSyncShareData->audio_update_clock = static_cast<uint64_t>(updatePts);
    ret = avsyncAshmem_->WriteToAshmem(static_cast<void*>(readSyncShareData), sizeof(AVsyncShareData), 0);
    CHECK_AND_RETURN_RET_LOG(!ret, ERR_DH_AUDIO_FAILED, "Write avsync data failed.");
    DHLOGD("readSyncShareData->audio_current_pts: %{public}" PRId64", readSyncShareData->audio_update_clock: "
        "%{public}" PRId64, readSyncShareData->audio_current_pts, readSyncShareData->audio_update_clock);
    return DH_SUCCESS;
}

uint32_t DMicDev::GetQueSize()
{
    return dataQueue_.size();
}

bool DMicDev::IsAVsync()
{
    bool isAVsync = false;
    std::lock_guard<std::mutex> lock(avSyncMutex_);
    isAVsync = avSyncParam_.isAVsync;
    return isAVsync;
}

int32_t DMicDev::AVsyncMacthScene(std::shared_ptr<AudioData> &data)
{
    std::lock_guard<std::mutex> lock(dataQueueMtx_);
    if (GetQueSize() < scene_ && isStartStatus_.load()) {
        int64_t videoPts = 0;
        int64_t audioPts = 0;
        int32_t ret = ReadTimeStampFromAVsync(videoPts);
        CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS, ERR_DH_AUDIO_FAILED, "ReadTimeStampFromAVsync failed.");
        if (GetQueSize() != 0) {
            CHECK_NULL_RETURN(dataQueue_.front(), ERR_DH_AUDIO_FAILED);
            audioPts = dataQueue_.front()->GetPts();
        }
        int64_t diff = static_cast<int64_t>(audioPts/TIME_CONVERSION_NTOU - videoPts/TIME_CONVERSION_STOU);
        DHLOGI("diff: %{public}" PRId64", videoPts: %{public}" PRId64
            ", audioPts: %{public}" PRId64, diff, videoPts, audioPts);
        if (diff > DADUIO_TIME_DIFF_MAX || GetQueSize() == 0) {
            data = std::make_shared<AudioData>(param_.comParam.frameSize);
        } else {
            isStartStatus_.store(false);
            data = dataQueue_.front();
            dataQueue_.pop_front();
        }
        return DH_SUCCESS;
    }
    isStartStatus_.store(false);
    if (GetQueSize() == 0) {
        isExistedEmpty_.store(true);
        DHLOGD("Data queue is empty");
        data = std::make_shared<AudioData>(param_.comParam.frameSize);
    } else {
        data = dataQueue_.front();
        dataQueue_.pop_front();
    }
    return DH_SUCCESS;
}

int32_t DMicDev::GetAudioDataFromQueue(std::shared_ptr<AudioData> &data)
{
    if (IsAVsync()) {
        int32_t ret = AVsyncMacthScene(data);
        if (ret != DH_SUCCESS) {
            DHLOGD("AVsyncMacthScene failed, insert empty frame");
            data = std::make_shared<AudioData>(param_.comParam.frameSize);
        }
    } else {
        std::lock_guard<std::mutex> lock(dataQueueMtx_);
        if (GetQueSize() == 0) {
            isExistedEmpty_.store(true);
            DHLOGD("Data queue is empty");
            data = std::make_shared<AudioData>(param_.comParam.frameSize);
        } else {
            data = dataQueue_.front();
            dataQueue_.pop_front();
        }
    }
    return DH_SUCCESS;
}

int32_t DMicDev::ReadStreamData(const int32_t streamId, std::shared_ptr<AudioData> &data)
{
    int64_t startTime = GetNowTimeUs();
    if (curStatus_ != AudioStatus::STATUS_START) {
        DHLOGE("Distributed audio is not starting status.");
        return ERR_DH_AUDIO_FAILED;
    }
    int32_t ret = GetAudioDataFromQueue(data);
    CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS || data == nullptr, ERR_DH_AUDIO_NULLPTR, "GetAudioData failed");
    ret = WriteTimeStampToAVsync(data->GetPts());
    CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS, ERR_DH_AUDIO_FAILED, "WriteTimeStampToAVsync failed");
    DHLOGD("Read stream data audioPts: %{public}" PRId64, data->GetPts());
    DumpFileUtil::WriteDumpFile(dumpFileCommn_, static_cast<void *>(data->Data()), data->Size());
    int64_t endTime = GetNowTimeUs();
    if (IsOutDurationRange(startTime, endTime, lastReadStartTime_)) {
        DHLOGE("This time read data spend: %{public}" PRId64" us, The interval of read data this time and "
            "the last time: %{public}" PRId64" us", endTime - startTime, startTime - lastReadStartTime_);
    }
    lastReadStartTime_ = startTime;
    return DH_SUCCESS;
}

int32_t DMicDev::ReadMmapPosition(const int32_t streamId, uint64_t &frames, CurrentTimeHDF &time)
{
    DHLOGD("Read mmap position. frames: %{public}" PRIu64", tvsec: %{public}" PRId64", tvNSec:%{public}" PRId64,
        writeNum_, writeTvSec_, writeTvNSec_);
    frames = writeNum_;
    time.tvSec = writeTvSec_;
    time.tvNSec = writeTvNSec_;
    return DH_SUCCESS;
}

int32_t DMicDev::RefreshAshmemInfo(const int32_t streamId,
    int32_t fd, int32_t ashmemLength, int32_t lengthPerTrans)
{
    DHLOGD("RefreshAshmemInfo: fd:%{public}d, ashmemLength: %{public}d, lengthPerTrans: %{public}d",
        fd, ashmemLength, lengthPerTrans);
    if (param_.captureOpts.capturerFlags == MMAP_MODE) {
        DHLOGD("DMic dev low-latency mode");
        if (ashmem_ != nullptr) {
            return DH_SUCCESS;
        }
        if (ashmemLength < ASHMEM_MAX_LEN) {
            ashmem_ = sptr<Ashmem>(new Ashmem(fd, ashmemLength));
            ashmemLength_ = ashmemLength;
            lengthPerTrans_ = lengthPerTrans;
            DHLOGD("Create ashmem success. fd:%{public}d, ashmem length: %{public}d, lengthPreTrans: %{public}d",
                fd, ashmemLength_, lengthPerTrans_);
            bool mapRet = ashmem_->MapReadAndWriteAshmem();
            if (!mapRet) {
                DHLOGE("Mmap ashmem failed.");
                return ERR_DH_AUDIO_NULLPTR;
            }
        }
    }
    return DH_SUCCESS;
}

int32_t DMicDev::MmapStart()
{
    CHECK_NULL_RETURN(ashmem_, ERR_DH_AUDIO_NULLPTR);
    std::lock_guard<std::mutex> lock(writeAshmemMutex_);
    frameIndex_ = 0;
    startTime_ = 0;
    isEnqueueRunning_.store(true);
    enqueueDataThread_ = std::thread([ptr = shared_from_this()]() {
        if (!ptr) {
            return;
        }
        ptr->EnqueueThread();
    });
    if (pthread_setname_np(enqueueDataThread_.native_handle(), ENQUEUE_THREAD) != DH_SUCCESS) {
        DHLOGE("Enqueue data thread setname failed.");
    }
    return DH_SUCCESS;
}

void DMicDev::EnqueueThread()
{
    writeIndex_ = 0;
    writeNum_ = 0;
    int64_t timeIntervalns = static_cast<int64_t>(paramHDF_.period * AUDIO_NS_PER_SECOND / AUDIO_MS_PER_SECOND);
    DHLOGD("Enqueue thread start, lengthPerWrite length: %{public}d, interval: %{public}d.", lengthPerTrans_,
        paramHDF_.period);
    FillJitterQueue();
    while (ashmem_ != nullptr && isEnqueueRunning_.load()) {
        int64_t timeOffset = UpdateTimeOffset(frameIndex_, timeIntervalns, startTime_);
        DHLOGD("Write frameIndex: %{public}" PRId64", timeOffset: %{public}" PRId64, frameIndex_, timeOffset);
        std::shared_ptr<AudioData> audioData = nullptr;
        {
            std::lock_guard<std::mutex> lock(dataQueueMtx_);
            if (dataQueue_.empty()) {
                DHLOGD("Data queue is Empty.");
                audioData = std::make_shared<AudioData>(param_.comParam.frameSize);
            } else {
                audioData = dataQueue_.front();
                dataQueue_.pop_front();
            }
            if (audioData == nullptr) {
                DHLOGD("The audioData is nullptr.");
                continue;
            }
            DumpFileUtil::WriteDumpFile(dumpFileFast_, static_cast<void *>(audioData->Data()), audioData->Size());
            bool writeRet = ashmem_->WriteToAshmem(audioData->Data(), audioData->Size(), writeIndex_);
            if (writeRet) {
                DHLOGD("Write to ashmem success! write index: %{public}d, writeLength: %{public}d.",
                    writeIndex_, lengthPerTrans_);
            } else {
                DHLOGE("Write data to ashmem failed.");
            }
        }
        writeIndex_ += lengthPerTrans_;
        if (writeIndex_ >= ashmemLength_) {
            writeIndex_ = 0;
        }
        writeNum_ += static_cast<uint64_t>(CalculateSampleNum(param_.comParam.sampleRate, paramHDF_.period));
        GetCurrentTime(writeTvSec_, writeTvNSec_);
        frameIndex_++;
        AbsoluteSleep(startTime_ + frameIndex_ * timeIntervalns - timeOffset);
    }
}

void DMicDev::FillJitterQueue()
{
    while (isEnqueueRunning_.load()) {
        {
            std::lock_guard<std::mutex> lock(dataQueueMtx_);
            if (paramHDF_.period == 0) {
                DHLOGE("DMicDev paramHDF_.period is zero");
                break;
            }
            if (dataQueue_.size() >= (LOW_LATENCY_JITTER_TIME_MS / paramHDF_.period)) {
                break;
            }
        }
        usleep(MMAP_WAIT_FRAME_US);
    }
    DHLOGD("Mic jitter data queue fill end.");
}

int32_t DMicDev::MmapStop()
{
    std::lock_guard<std::mutex> lock(writeAshmemMutex_);
    isEnqueueRunning_.store(false);
    if (enqueueDataThread_.joinable()) {
        enqueueDataThread_.join();
    }
    DHLOGI("Mic mmap stop end.");
    return DH_SUCCESS;
}

AudioParam DMicDev::GetAudioParam() const
{
    return param_;
}

int32_t DMicDev::NotifyHdfAudioEvent(const AudioEvent &event, const int32_t portId)
{
    int32_t ret = DAudioHdiHandler::GetInstance().NotifyEvent(devId_, portId, streamId_, event);
    if (ret != DH_SUCCESS) {
        DHLOGE("Notify event: %{public}d, result: %{public}s, streamId: %{public}d.",
            event.type, event.content.c_str(), streamId_);
    }
    return ret;
}

int32_t DMicDev::OnStateChange(const AudioEventType type)
{
    DHLOGD("On mic device state change, type: %{public}d", type);
    AudioEvent event;
    switch (type) {
        case AudioEventType::DATA_OPENED:
            isTransReady_.store(true);
            channelWaitCond_.notify_one();
            event.type = AudioEventType::MIC_OPENED;
            break;
        case AudioEventType::DATA_CLOSED:
            isTransReady_.store(false);
            event.type = AudioEventType::MIC_CLOSED;
            break;
        default:
            break;
    }
    event.content = GetCJsonString(KEY_DH_ID, std::to_string(dhId_).c_str());
    std::shared_ptr<IAudioEventCallback> cbObj = audioEventCallback_.lock();
    CHECK_NULL_RETURN(cbObj, ERR_DH_AUDIO_NULLPTR);
    cbObj->NotifyEvent(event);
    return DH_SUCCESS;
}

int32_t DMicDev::SendMessage(uint32_t type, std::string content, std::string dstDevId)
{
    DHLOGD("Send message to remote.");
    if (type != static_cast<uint32_t>(OPEN_MIC) && type != static_cast<uint32_t>(CLOSE_MIC) &&
        type != static_cast<uint32_t>(ENHANCE_PARAM_CHANGE)) {
        DHLOGE("Send message to remote. unsupported type: %{public}u", type);
        return ERR_DH_AUDIO_NULLPTR;
    }
    CHECK_NULL_RETURN(micCtrlTrans_, ERR_DH_AUDIO_NULLPTR);
    micCtrlTrans_->SendAudioEvent(type, content, dstDevId);
    return DH_SUCCESS;
}

int32_t DMicDev::OnDecodeTransDataDone(const std::shared_ptr<AudioData> &audioData)
{
    CHECK_NULL_RETURN(audioData, ERR_DH_AUDIO_NULLPTR);
    std::lock_guard<std::mutex> lock(dataQueueMtx_);
    dataQueSize_ = curStatus_ != AudioStatus::STATUS_START ?
        (param_.captureOpts.capturerFlags == MMAP_MODE ? lowLatencyHalfSize_ : scene_) :
        (param_.captureOpts.capturerFlags == MMAP_MODE ? lowLatencyMaxfSize_ : scene_ + scene_);
    if (isExistedEmpty_.load()) {
        dataQueSize_ = param_.captureOpts.capturerFlags == MMAP_MODE ? dataQueSize_ : DATA_QUEUE_EXT_SIZE;
    }
    uint64_t queueSize;
    while (dataQueue_.size() > dataQueSize_) {
        queueSize = static_cast<uint64_t>(dataQueue_.size());
        DHLOGI("Data queue overflow. buf current size: %{public}" PRIu64, queueSize);
        dataQueue_.pop_front();
    }
    std::shared_ptr<AudioData> writeAudioData = std::make_shared<AudioData>(param_.comParam.frameSize);
    if (memcpy_s(writeAudioData->Data(), writeAudioData->Capacity(), audioData->Data(), audioData->Capacity()) != EOK) {
        DHLOGE("Copy audio data failed");
    }
    writeAudioData->SetPts(audioData->GetPts());
    dataQueue_.push_back(writeAudioData);
    queueSize = static_cast<uint64_t>(dataQueue_.size());
    DHLOGD("Push new mic data, buf len: %{public}" PRIu64", audioPts: %{public}" PRId64,
        queueSize, audioData->GetPts());
    return DH_SUCCESS;
}

int32_t DMicDev::AVsyncRefreshAshmem(int32_t fd, int32_t ashmemLength)
{
    DHLOGD("AVsync mode: fd:%{public}d, ashmemLength: %{public}d", fd, ashmemLength);
    CHECK_AND_RETURN_RET_LOG(avsyncAshmem_ != nullptr, DH_SUCCESS, "AvsyncAshmem_ is not nullptr");
    if (ashmemLength >= static_cast<int32_t>(sizeof(AVsyncShareData)) && ashmemLength < ASHMEM_MAX_LEN && fd > 0) {
        avsyncAshmem_ = sptr<Ashmem>(new Ashmem(fd, ashmemLength));
        ashmemLength_ = ashmemLength;
        DHLOGD("Create ashmem success. fd:%{public}d, ashmem length: %{public}d", fd, ashmemLength_);
        bool mapRet = avsyncAshmem_->MapReadAndWriteAshmem();
        if (!mapRet) {
            DHLOGE("Mmap ashmem failed, cleaning up.");
            AVsyncDeintAshmem();
            avSyncParam_.fd = -1;
            avSyncParam_.sharedMemLen = 0;
            return ERR_DH_AUDIO_NULLPTR;
        }
    }
    return DH_SUCCESS;
}

void DMicDev::AVsyncDeintAshmem()
{
    if (avsyncAshmem_ != nullptr) {
        avsyncAshmem_->UnmapAshmem();
        avsyncAshmem_->CloseAshmem();
        avsyncAshmem_ = nullptr;
    }
}

int32_t DMicDev::UpdateWorkModeParam(const std::string &devId, const std::string &dhId, const AudioAsyncParam &param)
{
    std::lock_guard<std::mutex> lock(avSyncMutex_);
    avSyncParam_ = param;
    if (avSyncParam_.isAVsync) {
        auto ret = AVsyncRefreshAshmem(avSyncParam_.fd, avSyncParam_.sharedMemLen);
        CHECK_AND_RETURN_RET_LOG(ret != DH_SUCCESS, ret, "AVsyncRefreshAshmemInfo failed");
        scene_ = avSyncParam_.scene == static_cast<uint32_t>(AudioAVScene::BROADCAST) ?
            DATA_QUEUE_BROADCAST_SIZE : DATA_QUEUE_VIDEOCALL_SIZE;
        DHLOGI("AVsync mode opened");
    } else {
        AVsyncDeintAshmem();
        avSyncParam_.fd = -1;
        avSyncParam_.sharedMemLen = 0;
        DHLOGI("AVsync mode closed");
    }
    return DH_SUCCESS;
}
} // DistributedHardware
} // OHOS