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
1 change: 1 addition & 0 deletions components/ocs_scheduler/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ idf_component_register(
"async_task_scheduler.cpp"
"fanout_task.cpp"
"async_func_scheduler.cpp"
"event_func_scheduler.cpp"
"async_func.cpp"
"periodic_task_scheduler.cpp"
"constant_delay_estimator.cpp"
Expand Down
12 changes: 5 additions & 7 deletions components/ocs_scheduler/async_func_scheduler.h
Original file line number Diff line number Diff line change
Expand Up @@ -5,24 +5,22 @@

#pragma once

#include <functional>
#include <memory>
#include <vector>

#include "ocs_core/future.h"
#include "ocs_core/noncopyable.h"
#include "ocs_core/static_recursive_mutex.h"
#include "ocs_scheduler/ifunc_scheduler.h"
#include "ocs_scheduler/itask.h"
#include "ocs_system/iarena.h"

namespace ocs {
namespace scheduler {

class AsyncFuncScheduler : public ITask, private core::NonCopyable<> {
class AsyncFuncScheduler : public IFuncScheduler,
public ITask,
private core::NonCopyable<> {
public:
using FuturePtr = std::shared_ptr<core::Future>;
using Func = std::function<status::StatusCode()>;

//! Initialize.
//!
//! @params
Expand All @@ -38,7 +36,7 @@ class AsyncFuncScheduler : public ITask, private core::NonCopyable<> {
//!
//! @remarks
//! It is safe to call scheduler functions in @p func.
FuturePtr add(Func func);
FuturePtr add(Func func) override;

private:
const size_t max_event_count_ { 0 };
Expand Down
33 changes: 33 additions & 0 deletions components/ocs_scheduler/event_func_scheduler.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
/*
* SPDX-FileCopyrightText: 2026 Tendry Lab
* SPDX-License-Identifier: Apache-2.0
*/

#include "ocs_scheduler/event_func_scheduler.h"

namespace ocs {
namespace scheduler {

EventFuncScheduler::EventFuncScheduler(IFuncScheduler& func_scheduler,
EventGroupHandle_t handle,
EventBits_t event)
: event_(event)
, func_scheduler_(func_scheduler)
, handle_(handle) {
configASSERT(event_);
configASSERT(handle_);
}

EventFuncScheduler::FuturePtr EventFuncScheduler::add(Func func) {
auto future = func_scheduler_.add(func);
if (!future) {
return nullptr;
}

xEventGroupSetBits(handle_, event_);

return future;
}

} // namespace scheduler
} // namespace ocs
47 changes: 47 additions & 0 deletions components/ocs_scheduler/event_func_scheduler.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
/*
* SPDX-FileCopyrightText: 2026 Tendry Lab
* SPDX-License-Identifier: Apache-2.0
*/

#pragma once

#include "ocs_core/freertos.h"
#include "ocs_core/noncopyable.h"
#include "ocs_scheduler/ifunc_scheduler.h"

namespace ocs {
namespace scheduler {

//! Notify the func scheduler owner each time a func is scheduled.
//!
//! @remarks
//! Typical usage is to attach the func scheduler to the AsyncTaskScheduler, and to
//! use the received event to wake up the task scheduler, so the scheduled funcs are
//! handled without delay.
class EventFuncScheduler : public IFuncScheduler, private core::NonCopyable<> {
public:
//! Initialize.
//!
//! @params
//! - @p func_scheduler to schedule funcs.
//! - @p handle to post asynchronous events.
//! - @p event to post to the event group after each successfully scheduled func.
EventFuncScheduler(IFuncScheduler& func_scheduler,
EventGroupHandle_t handle,
EventBits_t event);

//! Schedule @p func and post the event.
//!
//! @remarks
//! The event isn't posted if @p func can't be scheduled.
FuturePtr add(Func func) override;

private:
const EventBits_t event_ { 0 };

IFuncScheduler& func_scheduler_;
EventGroupHandle_t handle_ { nullptr };
};

} // namespace scheduler
} // namespace ocs
36 changes: 36 additions & 0 deletions components/ocs_scheduler/ifunc_scheduler.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
/*
* SPDX-FileCopyrightText: 2026 Tendry Lab
* SPDX-License-Identifier: Apache-2.0
*/

#pragma once

#include <functional>
#include <memory>

#include "ocs_core/future.h"
#include "ocs_status/code.h"

namespace ocs {
namespace scheduler {

class IFuncScheduler {
public:
using FuturePtr = std::shared_ptr<core::Future>;
using Func = std::function<status::StatusCode()>;

//! Destroy.
virtual ~IFuncScheduler() = default;

//! Add @p func to be executed asynchronously.
//!
//! @remarks
//! It is safe to call scheduler functions in @p func.
//!
//! @return
//! Future to wait for the @p func result, nullptr if @p func can't be scheduled.
virtual FuturePtr add(Func func) = 0;
};

} // namespace scheduler
} // namespace ocs
1 change: 1 addition & 0 deletions components/ocs_scheduler/test/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ idf_component_register(
"test_async_task.cpp"
"test_async_task_scheduler.cpp"
"test_async_func_scheduler.cpp"
"test_event_func_scheduler.cpp"
"test_periodic_task_scheduler.cpp"

REQUIRES
Expand Down
70 changes: 70 additions & 0 deletions components/ocs_scheduler/test/test_event_func_scheduler.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
/*
* SPDX-FileCopyrightText: 2026 Tendry Lab
* SPDX-License-Identifier: Apache-2.0
*/

#include "unity.h"

#include "ocs_core/static_event_group.h"
#include "ocs_scheduler/async_func_scheduler.h"
#include "ocs_scheduler/event_func_scheduler.h"
#include "ocs_system/heap_arena.h"

namespace ocs {
namespace scheduler {

namespace {

system::HeapArena heap_arena;

} // namespace

TEST_CASE("Event func scheduler: event is posted when func is scheduled",
"[event_func_scheduler], [ocs_scheduler]") {
core::StaticEventGroup event_group;

const EventBits_t event = BIT(1);

AsyncFuncScheduler async_func_scheduler(heap_arena, 1);
EventFuncScheduler func_scheduler(async_func_scheduler, event_group.get(), event);

TEST_ASSERT_EQUAL(0, xEventGroupGetBits(event_group.get()));

auto future = func_scheduler.add([]() {
return status::StatusCode::NoData;
});
TEST_ASSERT_NOT_NULL(future);

TEST_ASSERT_EQUAL(event, xEventGroupGetBits(event_group.get()));

// The func is handled by the underlying scheduler.
TEST_ASSERT_EQUAL(status::StatusCode::OK, async_func_scheduler.run());
TEST_ASSERT_EQUAL(status::StatusCode::OK, future->wait());
TEST_ASSERT_EQUAL(status::StatusCode::NoData, future->code());
}

TEST_CASE("Event func scheduler: event isn't posted when func can't be scheduled",
"[event_func_scheduler], [ocs_scheduler]") {
core::StaticEventGroup event_group;

const EventBits_t event = BIT(1);

AsyncFuncScheduler async_func_scheduler(heap_arena, 1);
EventFuncScheduler func_scheduler(async_func_scheduler, event_group.get(), event);

TEST_ASSERT_NOT_NULL(func_scheduler.add([]() {
return status::StatusCode::OK;
}));

xEventGroupClearBits(event_group.get(), event);

// The underlying scheduler is full.
TEST_ASSERT_NULL(func_scheduler.add([]() {
return status::StatusCode::OK;
}));

TEST_ASSERT_EQUAL(0, xEventGroupGetBits(event_group.get()));
}

} // namespace scheduler
} // namespace ocs
Loading