Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 5 additions & 20 deletions src/api/MessageVideoTaskHandler.cc
Original file line number Diff line number Diff line change
Expand Up @@ -109,17 +109,7 @@ VideoTask::MsgQuerySwitchSend MessageVideoTaskHandler::Handle(VideoTask::MsgQuer
VideoTask::MsgSaveOrUpdateSend MessageVideoTaskHandler::Handle(VideoTask::MsgSaveOrUpdateRecv&& data,
std::error_condition& errc) {
VideoTask::MsgSaveOrUpdateSend retData{};
errc = task_config_.ModifyTaskParam(data.channelId, data.algorithmId, data.taskConfig);
if (util::ErrorEnum::Success != errc) {
return retData;
}

errc = task_config_.ModifyTaskStrategy(data.channelId, data.algorithmId, data.scheduleId);
if (util::ErrorEnum::Success != errc) {
return retData;
}

errc = task_config_.SwitchTask(data.channelId, data.algorithmId, true);
errc = task_config_.SaveOrUpdateTask(data.channelId, data.algorithmId, data.taskConfig, data.scheduleId);
return retData;
}

Expand Down Expand Up @@ -237,14 +227,8 @@ VideoTask::MsgApplyParamsBatchSend MessageVideoTaskHandler::Handle(VideoTask::Ms
VideoTask::MsgApplyParamsBatchSend retData{};

for (const auto& channelId : data.targetChannelIds) {
auto params = data.taskConfig;
std::error_condition ret = task_config_.ModifyTaskParam(channelId, data.algorithmId, params);
if (util::ErrorEnum::Success == ret) {
ret = task_config_.ModifyTaskStrategy(channelId, data.algorithmId, data.scheduleId);
}
if (util::ErrorEnum::Success == ret) {
ret = task_config_.SwitchTask(channelId, data.algorithmId, true);
}
std::error_condition ret =
task_config_.SaveOrUpdateTask(channelId, data.algorithmId, data.taskConfig, data.scheduleId);
if (util::ErrorEnum::Success != ret) {
MsgResultInfo failedEl;
failedEl.id = channelId;
Expand Down Expand Up @@ -290,7 +274,8 @@ namespace {
std::error_condition qerr = queue_status.actionStatus;
action.statusCode = std::to_string(qerr.value());
action.statusDesc = qerr.message();
action.statusDescKey = "api.error." + util::ErrorEnumName(static_cast<util::ErrorEnum>(qerr.value()));
action.statusDescKey =
"api.error." + util::ErrorEnumName(static_cast<util::ErrorEnum>(qerr.value()));
action.holdCount = queue_status.queueStatus.holdCount;
action.alarmCount = queue_status.alarmCount;
action.insertCount = queue_status.queueStatus.insertCount;
Expand Down
56 changes: 1 addition & 55 deletions src/flow/stream/StreamViewerEncoder.cc
Original file line number Diff line number Diff line change
Expand Up @@ -2,24 +2,13 @@

#include "flow/stream/StreamViewerEncoder.h"

#include <algorithm>

#include "media/PixelFormat.h"
#include "media/VideoUtil.h"
#include "mem/IDeviceContext.h"
#include "service/detail/ServiceRegistry.h"
#include "service/media/IVideoFrameTransform.h"
#include "util/Log.h"

namespace cosmo {
namespace {
size_t MinStartupKeyFrameSize(int width, int height) {
const auto pixels =
static_cast<size_t>(std::max(width, 1)) * static_cast<size_t>(std::max(height, 1));
return std::max<size_t>(1024, pixels / 4096);
}
} // namespace

StreamViewerEncoder::~StreamViewerEncoder() {
async_queue_.Stop();
if (encoder_) {
Expand Down Expand Up @@ -56,39 +45,9 @@ bool StreamViewerEncoder::OpenEncoder() {
return false;
}

startup_small_key_frame_count_ = 0;
startup_key_frame_accepted_ = false;
return true;
}

bool StreamViewerEncoder::ContainsKeyFrame(const uint8_t* data, size_t size) const {
if (!data || size == 0) {
return false;
}

size_t offset = 0;
while (offset < size) {
const size_t nal_size = media::SeparateHVideoFrame(data + offset, size - offset);
if (nal_size == 0) {
break;
}
if (media::GetFrameType(video_type_, data + offset, nal_size) == media::HFrameType::I) {
return true;
}
offset += nal_size;
}

return false;
}

bool StreamViewerEncoder::IsSmallStartupKeyFrame(const uint8_t* data, size_t size) const {
if (startup_key_frame_accepted_) {
return false;
}

return ContainsKeyFrame(data, size) && size < MinStartupKeyFrameSize(width_, height_);
}

bool StreamViewerEncoder::HandFrame(VideoFramePtr frame) {
if (encoder_) {
debug_info_.recvFrames += 1;
Expand Down Expand Up @@ -122,20 +81,7 @@ void StreamViewerEncoder::ProcFrame(VideoFramePtr frame) {
if (!packet) {
return;
}
if (IsSmallStartupKeyFrame(packet->GetData(), packet->GetSize())) {
startup_small_key_frame_count_ += 1;
LOG_WARN("Drop small startup keyframe size:{} min:{} count:{}", packet->GetSize(),
MinStartupKeyFrameSize(width_, height_), startup_small_key_frame_count_);
LOG_WARN("Reopen encoder after small startup keyframe, width:{} height:{}", width_, height_);
encoder_.reset();
OpenEncoder();
return;
}
if (packet && video_pusher_) {
if (ContainsKeyFrame(packet->GetData(), packet->GetSize())) {
startup_key_frame_accepted_ = true;
}
startup_small_key_frame_count_ = 0;
if (video_pusher_) {
debug_info_.sendFrames += 1;
video_pusher_->PushFrame(packet->GetData(), packet->GetSize());
}
Expand Down
4 changes: 0 additions & 4 deletions src/flow/stream/StreamViewerEncoder.h
Original file line number Diff line number Diff line change
Expand Up @@ -36,17 +36,13 @@ class StreamViewerEncoder {

private:
bool OpenEncoder();
bool ContainsKeyFrame(const uint8_t* data, size_t size) const;
bool IsSmallStartupKeyFrame(const uint8_t* data, size_t size) const;
void ProcFrame(VideoFramePtr frame);

std::mutex mtx_;
MsgOverviewDebugInfo debug_info_;
media::VideoCodecType video_type_{media::VideoCodecType::kInvalid};
int width_{0};
int height_{0};
int startup_small_key_frame_count_{0};
bool startup_key_frame_accepted_{false};
std::shared_ptr<media::VideoEncoder> encoder_ = nullptr;
RtmpStreamPusherPtr video_pusher_{nullptr};
AsyncQueue<VideoFramePtr> async_queue_;
Expand Down
7 changes: 7 additions & 0 deletions src/service/camera/ICameraTaskConfig.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,13 @@ class ICameraTaskConfig {
const std::string& algorithmId,
cosmo::MsgTaskConfig& params) = 0;

/// Atomically validate, persist, and enable a camera-algorithm task.
/// Validation failures must leave the task list and task configuration unchanged.
virtual cosmo::util::ErrorEnum SaveOrUpdateTask(const std::string& cameraId,
const std::string& algorithmId,
const cosmo::MsgTaskConfig& params,
const std::string& scheduleId) = 0;

/// Query current dynamic parameters for a camera–algorithm task.
/// @param cameraId Camera identifier.
/// @param algorithmId Algorithm identifier.
Expand Down
3 changes: 3 additions & 0 deletions src/service/camera/impl/CameraServiceImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,8 @@ class CameraServiceImpl : public ICameraService {

util::ErrorEnum ModifyTaskParam(const std::string& cameraId, const std::string& algorithmId,
MsgTaskConfig& params) override;
util::ErrorEnum SaveOrUpdateTask(const std::string& cameraId, const std::string& algorithmId,
const MsgTaskConfig& params, const std::string& scheduleId) override;
util::ErrorEnum QueryTaskParam(const std::string& cameraId, const std::string& algorithmId,
std::vector<MsgDynamicKeyValue>& params) override;
util::ErrorEnum ModifyTaskArea(const std::string& cameraId, const std::string& algorithmId,
Expand Down Expand Up @@ -185,6 +187,7 @@ class CameraServiceImpl : public ICameraService {
void LoadCameraTaskList(CameraEntityPtr camera);
void SaveCameraTaskList(const CameraEntityPtr& camera);
util::ErrorEnum MakeCameraTask(const CameraEntityPtr& camera, CameraTaskPtr task);
util::ErrorEnum CheckTaskStartResource() const;
void PrepareCameraTaskOverview(const CameraEntityPtr& camera, CameraTaskPtr task);
void SwitchCameraTask(const CameraEntityPtr& camera, CameraTaskPtr task);
void SwitchCameraTaskAsync(CameraEntityPtr camera, CameraTaskPtr task);
Expand Down
166 changes: 142 additions & 24 deletions src/service/camera/impl/CameraTaskConfig.cc
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,34 @@ namespace {
// Task parameter / area / strategy / switch operations
// ============================================================

util::ErrorEnum CameraServiceImpl::CheckTaskStartResource() const {
#ifdef COSMO_NN_USE_SOPHON_BACKEND
if (!ServiceRegistry::Instance().Get<IConfigReadService>().GetResourceLimit()) {
return util::ErrorEnum::Success;
}

size_t packet_total = 0, packet_proc = 0, packet_discard = 0, continues_discard_sec = 0;
ServiceRegistry::Instance().Get<ITaskQuery>().PacketStatus(packet_total, packet_proc, packet_discard,
continues_discard_sec);
double discard_percent = 0.0;
if (packet_total > 0) {
discard_percent = static_cast<double>(packet_discard) / static_cast<double>(packet_total);
}
auto gpu_info = ServiceRegistry::Instance().Get<IDeviceInfoService>().GetGpuUtilization();
std::vector<GpuMemSnapshot> devs;
devs.reserve(gpu_info.gpudevusage.size());
std::transform(gpu_info.gpudevusage.begin(), gpu_info.gpudevusage.end(), std::back_inserter(devs),
[](const auto& d) { return GpuMemSnapshot{d.gpumemtotal, d.gpumemavailable}; });
auto score = CalcCustomScore(gpu_info.gpuusage, gpu_info.gpumemtotal, gpu_info.gpumemavailable, devs,
discard_percent, continues_discard_sec);
if (score > 100.0) {
LOG_WARN("Task start rejected by resource guard, score:{}", score);
return util::ErrorEnum::ResourceLimit;
}
#endif
return util::ErrorEnum::Success;
}

std::string CameraServiceImpl::GetChannelName(const std::string& channelId) const {
std::shared_lock<std::shared_mutex> lock(mtx_);
auto it = std::find_if(cameras_.begin(), cameras_.end(), [&](const CameraEntityPtr& camera) {
Expand All @@ -58,6 +86,103 @@ std::string CameraServiceImpl::GetChannelName(const std::string& channelId) cons
return "";
}

util::ErrorEnum CameraServiceImpl::SaveOrUpdateTask(const std::string& cameraId,
const std::string& algorithmId,
const MsgTaskConfig& params,
const std::string& scheduleId) {
auto camera = GetCamera(cameraId);
if (!camera) {
LOG_INFO("{} Not Exist", cameraId);
return util::ErrorEnum::CameraNotExist;
}

std::lock_guard<std::mutex> command_lock(camera->command_mtx_);
if (camera->deleting_) {
LOG_INFO("{} Is Being Deleted", cameraId);
return util::ErrorEnum::CameraNotExist;
}

std::string scheduleName;
if (!ServiceRegistry::Instance().Get<IScheduleService>().Exist(scheduleId, scheduleName)) {
return util::ErrorEnum::TimeTemplateNotExist;
}

bool needs_enable = true;
{
std::shared_lock<std::shared_mutex> lock(camera->task_mtx_);
auto it = std::find_if(camera->tasks_.begin(), camera->tasks_.end(),
[&](const CameraTaskPtr& cfg) { return cfg->algorithm_code_ == algorithmId; });
if (it != camera->tasks_.end()) {
if (!(*it)->task_ || !(*it)->task_->IsReady()) {
LOG_WARN("[{}/{}] SaveOrUpdate skipped because task unit is not ready", cameraId,
algorithmId);
return util::ErrorEnum::TaskCreateFailed;
}
needs_enable = !(*it)->is_enabled_.load(std::memory_order_acquire);
} else if (camera->tasks_.size() >= camera->max_task_count_) {
return util::ErrorEnum::TaskTooMuch;
}
}

if (needs_enable) {
auto resource_ret = CheckTaskStartResource();
if (util::ErrorEnum::Success != resource_ret) {
LOG_WARN("[{}/{}] SaveOrUpdate rejected before mutation: {}", cameraId, algorithmId,
static_cast<uint32_t>(resource_ret));
return resource_ret;
}
}

CameraTaskPtr task;
bool start_task = false;
{
std::lock_guard<std::shared_mutex> lock(camera->task_mtx_);
auto it = std::find_if(camera->tasks_.begin(), camera->tasks_.end(),
[&](const CameraTaskPtr& cfg) { return cfg->algorithm_code_ == algorithmId; });
if (it != camera->tasks_.end()) {
task = *it;
MsgTaskConfig previous;
previous.params = task->task_->GetParams();
task->task_->GetArea(previous.areas, previous.shieldedAreas);

auto ret = task->task_->SetParams(params);
if (util::ErrorEnum::Success != ret) {
auto rollback_ret = task->task_->SetParams(previous);
if (util::ErrorEnum::Success != rollback_ret) {
LOG_ERRO("[{}/{}] SaveOrUpdate parameter rollback failed: {}", cameraId, algorithmId,
static_cast<uint32_t>(rollback_ret));
}
return ret;
}
} else {
task = std::make_shared<CameraTask>();
task->algorithm_code_ = algorithmId;
auto ret = MakeCameraTask(camera, task);
if (util::ErrorEnum::Success != ret) {
return ret;
}
ret = task->task_->SetParams(params);
if (util::ErrorEnum::Success != ret) {
return ret;
}
camera->tasks_.push_back(task);
}

const bool was_enabled = task->is_enabled_.load(std::memory_order_acquire);
task->data_.taskConfig = params;
task->schedule_id_ = scheduleId;
task->schedule_name_ = scheduleName;
task->is_enabled_.store(true, std::memory_order_release);
SaveCameraTaskList(camera);
start_task = !was_enabled;
}

if (start_task) {
SwitchCameraTaskAsync(camera, task);
}
return util::ErrorEnum::Success;
}

util::ErrorEnum CameraServiceImpl::ModifyTaskParam(const std::string& cameraId,
const std::string& algorithmId, MsgTaskConfig& params) {
return WithCamera(cameraId, [&](const CameraEntityPtr& camera) {
Expand Down Expand Up @@ -214,30 +339,7 @@ bool CameraServiceImpl::ScheduleInUse(const std::string& scheduleId) {

util::ErrorEnum CameraServiceImpl::SwitchTask(const std::string& cameraId, const std::string& algorithmId,
bool enable) {
// Authorization removed: no longer reject task start due to auth/expiry/limit (matches old 23461e05)
#ifdef COSMO_NN_USE_SOPHON_BACKEND
if (enable) {
if (ServiceRegistry::Instance().Get<IConfigReadService>().GetResourceLimit()) {
size_t packet_total = 0, packet_proc = 0, packet_discard = 0, continues_discard_sec = 0;
ServiceRegistry::Instance().Get<ITaskQuery>().PacketStatus(packet_total, packet_proc,
packet_discard, continues_discard_sec);
double discard_percent = 0.0;
if (packet_total > 0) {
discard_percent = static_cast<double>(packet_discard) / static_cast<double>(packet_total);
}
auto gpu_info = ServiceRegistry::Instance().Get<IDeviceInfoService>().GetGpuUtilization();
std::vector<GpuMemSnapshot> devs;
devs.reserve(gpu_info.gpudevusage.size());
std::transform(gpu_info.gpudevusage.begin(), gpu_info.gpudevusage.end(), std::back_inserter(devs),
[](const auto& d) { return GpuMemSnapshot{d.gpumemtotal, d.gpumemavailable}; });
auto score = CalcCustomScore(gpu_info.gpuusage, gpu_info.gpumemtotal, gpu_info.gpumemavailable,
devs, discard_percent, continues_discard_sec);
if (score > 100.0) {
return util::ErrorEnum::ResourceLimit;
}
}
}
#endif
// Authorization checks were removed; resource admission remains enforced for real start transitions.
auto camera = GetCamera(cameraId);
if (!camera) {
LOG_INFO("{} Not Exist", cameraId);
Expand All @@ -249,6 +351,22 @@ util::ErrorEnum CameraServiceImpl::SwitchTask(const std::string& cameraId, const
return util::ErrorEnum::CameraNotExist;
}

{
std::shared_lock<std::shared_mutex> lock(camera->task_mtx_);
auto it = std::find_if(camera->tasks_.begin(), camera->tasks_.end(),
[&](const CameraTaskPtr& cfg) { return cfg->algorithm_code_ == algorithmId; });
if (it != camera->tasks_.end() && (*it)->is_enabled_.load(std::memory_order_acquire) == enable) {
return util::ErrorEnum::Success;
}
}

if (enable) {
auto resource_ret = CheckTaskStartResource();
if (util::ErrorEnum::Success != resource_ret) {
return resource_ret;
}
}

CameraTaskPtr taskToSwitch = nullptr;
{
std::lock_guard<std::shared_mutex> lock(camera->task_mtx_);
Expand Down
4 changes: 4 additions & 0 deletions test/mock/MockCameraService.h
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,10 @@ class MockCameraService : public cosmo::service::ICameraService {
MAKE_MOCK3(ModifyTaskParam,
cosmo::util::ErrorEnum(const std::string&, const std::string&, cosmo::MsgTaskConfig&),
override);
MAKE_MOCK4(SaveOrUpdateTask,
cosmo::util::ErrorEnum(const std::string&, const std::string&, const cosmo::MsgTaskConfig&,
const std::string&),
override);
MAKE_MOCK3(QueryTaskParam,
cosmo::util::ErrorEnum(const std::string&, const std::string&,
std::vector<cosmo::MsgDynamicKeyValue>&),
Expand Down
Loading
Loading