diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d6de6c36..6895172a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -196,6 +196,12 @@ jobs: esp_idf_version: v6.0.2 target: esp32s3 path: . + - name: 构建 Timing ESP Unity 测试固件 + uses: espressif/esp-idf-ci-action@e6f5c74232b1ccd4c97ed641f1e48553853f1fd5 # v1.2.0 + with: + esp_idf_version: v6.0.2 + target: esp32s3 + path: tests/board/timing_esp_unity - name: 构建 FATFS/WL + SQLite 存储 Profile uses: espressif/esp-idf-ci-action@e6f5c74232b1ccd4c97ed641f1e48553853f1fd5 # v1.2.0 with: diff --git a/components/voicelife_timing/CMakeLists.txt b/components/voicelife_timing/CMakeLists.txt index 114fb6c8..fb685975 100644 --- a/components/voicelife_timing/CMakeLists.txt +++ b/components/voicelife_timing/CMakeLists.txt @@ -1,4 +1,5 @@ idf_component_register( - SRCS "src/timing_task.cc" + SRCS "src/timing_task.cc" "src/timing_runtime.cc" INCLUDE_DIRS "include" + REQUIRES voicelife_contracts ) diff --git a/components/voicelife_timing/include/voicelife/timing/timing_runtime.h b/components/voicelife_timing/include/voicelife/timing/timing_runtime.h new file mode 100644 index 00000000..481aaa54 --- /dev/null +++ b/components/voicelife_timing/include/voicelife/timing/timing_runtime.h @@ -0,0 +1,121 @@ +#pragma once + +#include + +#include "voicelife/contracts/status.h" +#include "voicelife/timing/timing_task.h" + +namespace voicelife::timing { + +/// 为 Timing Runner 提供绝对墙上时钟的 Port。 +class ClockPort { + public: + /** @brief 允许通过接口指针安全释放实现。 */ + virtual ~ClockPort() = default; + + /** + * @brief 读取当前绝对墙上时钟。 + * @return 当前 trigger_at 时间域中的绝对时间。 + */ + virtual TriggerAt Now() = 0; +}; + +/// 单个一次性定时器的调度 Port。 +class OneShotTimerPort { + public: + /// 一次性 timer 到期时的轻量回调。 + using ExpiryCallback = std::function; + + /** @brief 允许通过接口指针安全释放实现。 */ + virtual ~OneShotTimerPort() = default; + + /** + * @brief 绑定一次性 timer 到期回调。 + * @param callback 仅用于轻量通知 Runner 的回调。 + * + * 仅允许在首次设置 timer 前绑定;运行中解除必须调用 ClearExpiryCallbackAndWait。 + */ + virtual void SetExpiryCallback(ExpiryCallback callback) = 0; + + /** + * @brief 解除到期回调并等待已经开始的回调执行完毕。 + * + * 返回前必须保证后续不会再开始旧回调;TimingTaskRuntime 析构依赖此所有权屏障。 + */ + virtual void ClearExpiryCallbackAndWait() = 0; + + /** + * @brief 从现在起按相对延时设置一次唤醒。 + * @param delay 非负微秒延时。 + * @return timer 设置结果;平台暂不可用与内部错误使用不同错误码。 + */ + virtual Status ArmAfter(std::chrono::microseconds delay) = 0; + + /** + * @brief 取消当前一次性唤醒。 + * @return timer 停止结果;原本未设置视为成功。 + */ + virtual Status Disarm() = 0; +}; + +/// 轻量通知普通 Runner 执行上下文的 Port。 +class RunnerWakePort { + public: + /** @brief 允许通过接口指针安全释放实现。 */ + virtual ~RunnerWakePort() = default; + + /** + * @brief 通知 Runner 处理命令或 timer 到期事件。 + * + * 实现必须在本 Port 借用期内保持 Runner 已绑定,并保证通知不可失败; + * 销毁方应先解除 timer callback,再解绑 Runner。 + */ + virtual void Notify() = 0; +}; + +/// 编排异步命令、到期推进和单个一次性唤醒的 Timing 运行时。 +class TimingTaskRuntime final : public TimingTaskService { + public: + /** + * @brief 绑定纯 Runner、墙上时钟、一次性 timer 和 Runner 唤醒 Port。 + * @param runner 独占 task 注册表的 Runner。 + * @param clock 提供绝对墙上时钟。 + * @param timer 管理唯一的一次性唤醒。 + * @param wake 轻量通知 Runner 执行上下文。 + * + * 四个依赖均为借用引用,必须比本对象活得更久;析构时会先解除 timer callback。 + */ + TimingTaskRuntime(InMemoryTimingTaskRunner& runner, ClockPort& clock, OneShotTimerPort& timer, + RunnerWakePort& wake); + + /** @brief 解除 timer 回调并等待在途回调结束;不会销毁借用的依赖。 */ + ~TimingTaskRuntime() override; + + /** + * @brief 接收一个一次性 task 注册命令并通知 Runner。 + * @param command 待 Runner 应用的注册命令。 + * @return 命令已被接收时返回 kAccepted。 + */ + CommandAcceptance RegisterTask(RegisterTaskCommand command) override; + + /** + * @brief 接收一个 task 取消命令并通知 Runner。 + * @param command 待 Runner 应用的取消命令。 + * @return 命令已被接收时返回 kAccepted。 + */ + CommandAcceptance CancelTask(CancelTaskCommand command) override; + + /** + * @brief 在 Runner 执行上下文中消费命令、推进到期批次并重设一次性唤醒。 + * @return timer 成功设置或停止时成功,否则保留平台错误分类。 + */ + Status ProcessWake(); + + private: + InMemoryTimingTaskRunner& runner_; + ClockPort& clock_; + OneShotTimerPort& timer_; + RunnerWakePort& wake_; +}; + +} // namespace voicelife::timing diff --git a/components/voicelife_timing/include/voicelife/timing/timing_task.h b/components/voicelife_timing/include/voicelife/timing/timing_task.h index e4702c0c..3bf544ce 100644 --- a/components/voicelife_timing/include/voicelife/timing/timing_task.h +++ b/components/voicelife_timing/include/voicelife/timing/timing_task.h @@ -4,8 +4,10 @@ #include #include #include +#include #include #include +#include #include #include @@ -90,9 +92,12 @@ struct CancelTaskCommand { CancelTaskResultCallback on_result; }; -/// 公开命令入口的同步接收结果;不表示命令已经被 Runner 应用。 +/// 公开命令入口的同步接收结果;accepted 不表示命令已经被 Runner 应用。 enum class CommandAcceptance { + /// 命令已进入 Runner 队列,或已在 callback 上下文中立即应用。 kAccepted, + /// 命令未进入 Runner 队列,调用方可以稍后重试。 + kUnavailable, }; /// Runner 推进一个已确定到期批次后的可观察结果。 @@ -110,23 +115,24 @@ struct RunDueTasksResult { */ bool CanTransition(TaskStatus from, TaskStatus to); -/// 日程模块异步提交一次性 task 注册和取消命令的公开接口。 +/// 日程模块提交一次性 task 注册和取消命令的公开接口。 +/// 外部调用异步排队;Runner callback 在所属 Runner 执行上下文中调用时立即应用。 class TimingTaskService { public: /** @brief 允许通过接口指针安全释放实现。 */ virtual ~TimingTaskService() = default; /** - * @brief 异步提交一个一次性 task。 + * @brief 提交一个一次性 task;外部调用异步排队,Runner callback 内调用立即应用。 * @param command 包含不透明 task_id、绝对到点时刻和运行时 callback 的命令。 - * @return kAccepted 表示命令已被接收,最终应用结果由 Runner 另行通知。 + * @return kAccepted 表示命令已接收;kUnavailable 表示未接收且可以重试。 */ virtual CommandAcceptance RegisterTask(RegisterTaskCommand command) = 0; /** - * @brief 异步提交一个 task 取消请求。 + * @brief 提交一个 task 取消请求;外部调用异步排队,Runner callback 内调用立即应用。 * @param command 包含待取消的不透明 task_id 和可选的最终结果回调。 - * @return kAccepted 表示命令已被接收;Runner 消费后通过 on_result 通知 cancelled 或 not found。 + * @return kAccepted 表示命令已接收;kUnavailable 表示未接收且可以重试。 */ virtual CommandAcceptance CancelTask(CancelTaskCommand command) = 0; }; @@ -134,19 +140,19 @@ class TimingTaskService { /** * @brief Host 可用的异步命令 Runner,独占内存 task 注册表。 * - * 公开命令入口只保存命令;注册表仅由 Runner 的消费和推进入口修改。 + * 外部命令入口只保存命令;Runner callback 内命令由同一 Runner 执行上下文立即应用。 */ class InMemoryTimingTaskRunner final : public TimingTaskService { public: /** - * @brief 异步接收一条注册命令。 + * @brief 接收一条注册命令;Runner callback 内调用时立即应用,其他上下文异步排队。 * @param command 待 Runner 应用的注册命令。 * @return 命令接收后固定返回 kAccepted。 */ CommandAcceptance RegisterTask(RegisterTaskCommand command) override; /** - * @brief 异步接收一条取消命令。 + * @brief 接收一条取消命令;Runner callback 内调用时立即应用,其他上下文异步排队。 * @param command 待 Runner 应用的取消命令。 * @return 命令接收后固定返回 kAccepted。 */ @@ -172,6 +178,12 @@ class InMemoryTimingTaskRunner final : public TimingTaskService { */ [[nodiscard]] std::optional NextWakeAt() const; + /** + * @brief 判断调用线程当前是否正在执行本 Runner 的 task callback。 + * @return 仅在本 Runner 调用 task callback 的动态范围内返回 true。 + */ + [[nodiscard]] bool IsInCallbackContext() const; + private: /// pending task 及其仅驻留内存的运行时回调。 struct PendingTask { @@ -179,14 +191,20 @@ class InMemoryTimingTaskRunner final : public TimingTaskService { TaskCallback callback; }; + using PendingTaskList = std::list; /// 按公开入口接收顺序保存的 Runner 命令。 using PendingCommand = std::variant; + void ApplyRegisterTask(RegisterTaskCommand command, TriggerAt applied_at); + void ApplyCancelTask(CancelTaskCommand command, TriggerAt applied_at); + std::deque commands_; std::vector used_task_ids_; - std::vector pending_tasks_; + PendingTaskList pending_tasks_; + std::unordered_map pending_tasks_by_id_; std::vector terminal_tasks_; bool processing_commands_ = false; + std::optional callback_applied_at_; }; } // namespace voicelife::timing diff --git a/components/voicelife_timing/src/timing_runtime.cc b/components/voicelife_timing/src/timing_runtime.cc new file mode 100644 index 00000000..7eafff2b --- /dev/null +++ b/components/voicelife_timing/src/timing_runtime.cc @@ -0,0 +1,44 @@ +#include "voicelife/timing/timing_runtime.h" + +#include +#include + +namespace voicelife::timing { + +TimingTaskRuntime::TimingTaskRuntime(InMemoryTimingTaskRunner& runner, ClockPort& clock, OneShotTimerPort& timer, + RunnerWakePort& wake) + : runner_(runner), clock_(clock), timer_(timer), wake_(wake) { + timer_.SetExpiryCallback([this] { wake_.Notify(); }); +} + +TimingTaskRuntime::~TimingTaskRuntime() { timer_.ClearExpiryCallbackAndWait(); } + +CommandAcceptance TimingTaskRuntime::RegisterTask(RegisterTaskCommand command) { + const auto callback_internal = runner_.IsInCallbackContext(); + const auto acceptance = runner_.RegisterTask(std::move(command)); + if (!callback_internal) { + wake_.Notify(); + } + return acceptance; +} + +CommandAcceptance TimingTaskRuntime::CancelTask(CancelTaskCommand command) { + const auto callback_internal = runner_.IsInCallbackContext(); + const auto acceptance = runner_.CancelTask(std::move(command)); + if (!callback_internal) { + wake_.Notify(); + } + return acceptance; +} + +Status TimingTaskRuntime::ProcessWake() { + const auto now = clock_.Now(); + const auto result = runner_.RunDueTasks(now); + if (!result.next_wake_at.has_value()) { + return timer_.Disarm(); + } + const auto arm_now = clock_.Now(); + return timer_.ArmAfter(std::max(std::chrono::microseconds::zero(), *result.next_wake_at - arm_now)); +} + +} // namespace voicelife::timing diff --git a/components/voicelife_timing/src/timing_task.cc b/components/voicelife_timing/src/timing_task.cc index b7de274c..b9b24833 100644 --- a/components/voicelife_timing/src/timing_task.cc +++ b/components/voicelife_timing/src/timing_task.cc @@ -5,6 +5,12 @@ namespace voicelife::timing { +namespace { + +thread_local InMemoryTimingTaskRunner* active_callback_runner = nullptr; + +} // namespace + std::optional TaskId::Create(std::string value) { if (value.empty()) { return std::nullopt; @@ -24,11 +30,19 @@ bool CanTransition(TaskStatus from, TaskStatus to) { } CommandAcceptance InMemoryTimingTaskRunner::RegisterTask(RegisterTaskCommand command) { + if (active_callback_runner == this) { + ApplyRegisterTask(std::move(command), *callback_applied_at_); + return CommandAcceptance::kAccepted; + } commands_.push_back(std::move(command)); return CommandAcceptance::kAccepted; } CommandAcceptance InMemoryTimingTaskRunner::CancelTask(CancelTaskCommand command) { + if (active_callback_runner == this) { + ApplyCancelTask(std::move(command), *callback_applied_at_); + return CommandAcceptance::kAccepted; + } commands_.push_back(std::move(command)); return CommandAcceptance::kAccepted; } @@ -48,61 +62,69 @@ size_t InMemoryTimingTaskRunner::ProcessPendingCommands(TriggerAt applied_at) { auto pending_command = std::move(commands_.front()); commands_.pop_front(); if (auto* cancel = std::get_if(&pending_command)) { - const auto task = - std::find_if(pending_tasks_.begin(), pending_tasks_.end(), - [&cancel](const auto& pending) { return pending.task.id == cancel->task_id; }); - if (task == pending_tasks_.end()) { - if (cancel->on_result) { - cancel->on_result(CancelTaskResult::kNotFound); - } - continue; - } - terminal_tasks_.reserve(terminal_tasks_.size() + 1); - task->task.status = TaskStatus::kCancelled; - task->task.updated_at = applied_at; - terminal_tasks_.push_back(std::move(task->task)); - pending_tasks_.erase(task); - if (cancel->on_result) { - cancel->on_result(CancelTaskResult::kCancelled); - } + ApplyCancelTask(std::move(*cancel), applied_at); continue; } - auto command = std::move(std::get(pending_command)); - if (std::find(used_task_ids_.begin(), used_task_ids_.end(), command.task_id.Value()) != used_task_ids_.end()) { - if (command.on_result) { - command.on_result(RegisterTaskResult::kDuplicate); - } - continue; + ApplyRegisterTask(std::move(std::get(pending_command)), applied_at); + } + return processed_count; +} + +void InMemoryTimingTaskRunner::ApplyRegisterTask(RegisterTaskCommand command, TriggerAt applied_at) { + if (std::find(used_task_ids_.begin(), used_task_ids_.end(), command.task_id.Value()) != used_task_ids_.end()) { + if (command.on_result) { + command.on_result(RegisterTaskResult::kDuplicate); } - PendingTask pending_task{ - .task = - { - .id = std::move(command.task_id), - .trigger_at = command.trigger_at, - .status = TaskStatus::kPending, - .created_at = applied_at, - .updated_at = applied_at, - }, - .callback = std::move(command.callback), - }; - auto used_task_id = pending_task.task.id.Value(); - pending_tasks_.reserve(pending_tasks_.size() + 1); - used_task_ids_.reserve(used_task_ids_.size() + 1); - const auto insertion_point = std::lower_bound(pending_tasks_.begin(), pending_tasks_.end(), pending_task, - [](const PendingTask& lhs, const PendingTask& rhs) { - if (lhs.task.trigger_at != rhs.task.trigger_at) { - return lhs.task.trigger_at < rhs.task.trigger_at; - } - return lhs.task.id.Value() < rhs.task.id.Value(); - }); - pending_tasks_.insert(insertion_point, std::move(pending_task)); - used_task_ids_.push_back(std::move(used_task_id)); + return; + } + PendingTask pending_task{ + .task = + { + .id = std::move(command.task_id), + .trigger_at = command.trigger_at, + .status = TaskStatus::kPending, + .created_at = applied_at, + .updated_at = applied_at, + }, + .callback = std::move(command.callback), + }; + auto used_task_id = pending_task.task.id.Value(); + used_task_ids_.reserve(used_task_ids_.size() + 1); + const auto insertion_point = std::lower_bound(pending_tasks_.begin(), pending_tasks_.end(), pending_task, + [](const PendingTask& lhs, const PendingTask& rhs) { + if (lhs.task.trigger_at != rhs.task.trigger_at) { + return lhs.task.trigger_at < rhs.task.trigger_at; + } + return lhs.task.id.Value() < rhs.task.id.Value(); + }); + const auto inserted_task = pending_tasks_.insert(insertion_point, std::move(pending_task)); + pending_tasks_by_id_.emplace(used_task_id, inserted_task); + used_task_ids_.push_back(std::move(used_task_id)); + if (command.on_result) { + command.on_result(RegisterTaskResult::kRegistered); + } +} + +void InMemoryTimingTaskRunner::ApplyCancelTask(CancelTaskCommand command, TriggerAt applied_at) { + const auto task_by_id = pending_tasks_by_id_.find(command.task_id.Value()); + if (task_by_id != pending_tasks_by_id_.end()) { + const auto task = task_by_id->second; + terminal_tasks_.reserve(terminal_tasks_.size() + 1); + task->task.status = TaskStatus::kCancelled; + task->task.updated_at = applied_at; + terminal_tasks_.push_back(std::move(task->task)); + pending_tasks_.erase(task); + pending_tasks_by_id_.erase(task_by_id); if (command.on_result) { - command.on_result(RegisterTaskResult::kRegistered); + command.on_result(CancelTaskResult::kCancelled); } + return; + } + + if (command.on_result) { + command.on_result(CancelTaskResult::kNotFound); } - return processed_count; } RunDueTasksResult InMemoryTimingTaskRunner::RunDueTasks(TriggerAt now) { @@ -110,20 +132,47 @@ RunDueTasksResult InMemoryTimingTaskRunner::RunDueTasks(TriggerAt now) { const auto due_end = std::upper_bound( pending_tasks_.begin(), pending_tasks_.end(), now, [](TriggerAt boundary, const PendingTask& pending) { return boundary < pending.task.trigger_at; }); - const auto processed_count = static_cast(due_end - pending_tasks_.begin()); - terminal_tasks_.reserve(terminal_tasks_.size() + processed_count); + std::vector due_task_ids; + due_task_ids.reserve(static_cast(std::distance(pending_tasks_.begin(), due_end))); for (auto due_task = pending_tasks_.begin(); due_task != due_end; ++due_task) { - due_task->task.status = TaskStatus::kExecuting; - due_task->task.updated_at = now; - due_task->callback(due_task->task.id, due_task->task.trigger_at); - due_task->task.status = TaskStatus::kCompleted; - due_task->task.updated_at = now; - terminal_tasks_.push_back(std::move(due_task->task)); + due_task_ids.push_back(due_task->task.id.Value()); + } + terminal_tasks_.reserve(terminal_tasks_.size() + due_task_ids.size()); + struct BatchGuard { + InMemoryTimingTaskRunner& runner; + InMemoryTimingTaskRunner* previous_callback_runner; + ~BatchGuard() { + runner.callback_applied_at_.reset(); + active_callback_runner = previous_callback_runner; + } + } batch_guard{*this, active_callback_runner}; + size_t processed_count = 0; + size_t skipped_count = 0; + for (const auto& due_task_id : due_task_ids) { + const auto due_task_by_id = pending_tasks_by_id_.find(due_task_id); + if (due_task_by_id == pending_tasks_by_id_.end()) { + ++skipped_count; + continue; + } + const auto due_task = due_task_by_id->second; + auto executing_task = std::move(*due_task); + pending_tasks_.erase(due_task); + pending_tasks_by_id_.erase(due_task_by_id); + executing_task.task.status = TaskStatus::kExecuting; + executing_task.task.updated_at = now; + callback_applied_at_ = now; + active_callback_runner = this; + executing_task.callback(executing_task.task.id, executing_task.task.trigger_at); + active_callback_runner = batch_guard.previous_callback_runner; + callback_applied_at_.reset(); + executing_task.task.status = TaskStatus::kCompleted; + executing_task.task.updated_at = now; + terminal_tasks_.push_back(std::move(executing_task.task)); + ++processed_count; } - pending_tasks_.erase(pending_tasks_.begin(), due_end); return { .processed_count = processed_count, - .skipped_count = 0, + .skipped_count = skipped_count, .next_wake_at = NextWakeAt(), }; } @@ -135,4 +184,6 @@ std::optional InMemoryTimingTaskRunner::NextWakeAt() const { return pending_tasks_.front().task.trigger_at; } +bool InMemoryTimingTaskRunner::IsInCallbackContext() const { return active_callback_runner == this; } + } // namespace voicelife::timing diff --git a/components/voicelife_timing_esp/CMakeLists.txt b/components/voicelife_timing_esp/CMakeLists.txt new file mode 100644 index 00000000..893d7c9d --- /dev/null +++ b/components/voicelife_timing_esp/CMakeLists.txt @@ -0,0 +1,6 @@ +idf_component_register( + SRCS "src/esp_timing_runtime.cc" + INCLUDE_DIRS "include" + REQUIRES voicelife_contracts voicelife_timing + PRIV_REQUIRES esp_timer freertos +) diff --git a/components/voicelife_timing_esp/include/voicelife/timing_esp/esp_timing_runtime.h b/components/voicelife_timing_esp/include/voicelife/timing_esp/esp_timing_runtime.h new file mode 100644 index 00000000..5e932252 --- /dev/null +++ b/components/voicelife_timing_esp/include/voicelife/timing_esp/esp_timing_runtime.h @@ -0,0 +1,64 @@ +#pragma once + +#include + +#include "voicelife/contracts/status.h" +#include "voicelife/timing/timing_runtime.h" + +namespace voicelife::timing_esp { + +using timing::CancelTaskCommand; +using timing::CommandAcceptance; +using timing::RegisterTaskCommand; +using timing::TimingTaskService; + +/** + * @brief ESP-IDF 平台上的 one-shot esp_timer 与 FreeRTOS Runner 运行时。 + * + * 外部 RegisterTask/CancelTask 通过 FreeRTOS Queue 交给普通 Runner task; + * esp_timer callback 只发送 Task Notification,不执行 task callback。 + */ +class EspTimingTaskRuntime final : public TimingTaskService { + public: + /** + * @brief 创建并启动 ESP Timing 运行时。 + * @return timer、Queue 和 Runner task 全部创建成功时返回实例,否则返回类型化失败。 + */ + static Result> Create(); + + /** + * @brief 停止 Runner 并释放 timer、Queue 和 Task Notification 资源。 + * + * 必须从 Runner task 之外销毁;销毁会等待已开始的 callback 与已接受命令处理完成。 + */ + ~EspTimingTaskRuntime() override; + + /** @brief 禁止拷贝构造。 */ + EspTimingTaskRuntime(const EspTimingTaskRuntime&) = delete; + /** @brief 禁止拷贝赋值。 */ + EspTimingTaskRuntime& operator=(const EspTimingTaskRuntime&) = delete; + + /** + * @brief 将一个一次性 task 注册命令提交给 ESP Runner。 + * @param command 待 Runner 应用的注册命令。 + * @return 命令进入 Runner 队列时返回 kAccepted;队列忙、已满或停机时返回 kUnavailable。 + */ + CommandAcceptance RegisterTask(RegisterTaskCommand command) override; + + /** + * @brief 将一个 task 取消命令提交给 ESP Runner。 + * @param command 待 Runner 应用的取消命令。 + * @return 命令进入 Runner 队列时返回 kAccepted;队列忙、已满或停机时返回 kUnavailable。 + */ + CommandAcceptance CancelTask(CancelTaskCommand command) override; + + private: + /// 封装 ESP-IDF timer、Queue 和 Runner task 的平台实现。 + class Impl; + + explicit EspTimingTaskRuntime(std::unique_ptr impl); + + std::unique_ptr impl_; +}; + +} // namespace voicelife::timing_esp diff --git a/components/voicelife_timing_esp/src/esp_timing_runtime.cc b/components/voicelife_timing_esp/src/esp_timing_runtime.cc new file mode 100644 index 00000000..21f2622b --- /dev/null +++ b/components/voicelife_timing_esp/src/esp_timing_runtime.cc @@ -0,0 +1,373 @@ +#include "voicelife/timing_esp/esp_timing_runtime.h" + +#include + +#ifdef ESP_PLATFORM + +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +namespace voicelife::timing_esp { + +using namespace timing; + +namespace { + +constexpr UBaseType_t kCommandQueueDepth = 32; +constexpr uint32_t kRunnerStackWords = 6144; +constexpr UBaseType_t kRunnerPriority = 4; +constexpr char kTag[] = "VoiceLifeTiming"; + +class EspSystemClock final : public ClockPort { + public: + TriggerAt Now() override { + return std::chrono::time_point_cast(std::chrono::system_clock::now()); + } +}; + +class EspOneShotTimer final : public OneShotTimerPort { + public: + ~EspOneShotTimer() override { + if (timer_ != nullptr) { + ClearExpiryCallbackAndWait(); + (void)esp_timer_delete(timer_); + } + if (quiesced_ != nullptr) { + vSemaphoreDelete(quiesced_); + } + } + + bool Initialize() { + quiesced_ = xSemaphoreCreateBinary(); + if (quiesced_ == nullptr) { + return false; + } + esp_timer_create_args_t args{}; + args.callback = &TimerEntry; + args.arg = this; + args.dispatch_method = ESP_TIMER_TASK; + args.name = "voicelife_timing"; + return esp_timer_create(&args, &timer_) == ESP_OK; + } + + void SetExpiryCallback(ExpiryCallback callback) override { callback_ = std::move(callback); } + + void ClearExpiryCallbackAndWait() override { + if (timer_ == nullptr || quiesced_callback_) { + callback_ = {}; + return; + } + quiescing_.store(true); + const auto stop_result = esp_timer_stop(timer_); + if (stop_result != ESP_OK && stop_result != ESP_ERR_INVALID_STATE) { + std::abort(); + } + while (xSemaphoreTake(quiesced_, 0) == pdTRUE) { + } + portENTER_CRITICAL(&barrier_lock_); + barrier_armed_ = true; + const auto barrier_result = esp_timer_start_once(timer_, 1); + if (barrier_result != ESP_OK) { + barrier_armed_ = false; + } + portEXIT_CRITICAL(&barrier_lock_); + if (barrier_result != ESP_OK || xSemaphoreTake(quiesced_, portMAX_DELAY) != pdTRUE) { + std::abort(); + } + while (callback_accessing_owner_.load()) { + taskYIELD(); + } + callback_ = {}; + quiesced_callback_ = true; + } + + Status ArmAfter(std::chrono::microseconds delay) override { + const auto stop_result = esp_timer_stop(timer_); + if (stop_result != ESP_OK && stop_result != ESP_ERR_INVALID_STATE) { + return Status::Error(ErrorCode::kInternal, "停止已有 ESP Timing timer 失败"); + } + const auto delay_us = delay.count() > 0 ? static_cast(delay.count()) : 1ULL; + const auto result = esp_timer_start_once(timer_, delay_us); + if (result == ESP_OK) { + return Status::Ok(); + } + return Status::Error(result == ESP_ERR_NO_MEM ? ErrorCode::kUnavailable : ErrorCode::kInternal, + "设置 ESP Timing timer 失败"); + } + + Status Disarm() override { + const auto result = esp_timer_stop(timer_); + if (result == ESP_OK || result == ESP_ERR_INVALID_STATE) { + return Status::Ok(); + } + return Status::Error(ErrorCode::kInternal, "停止 ESP Timing timer 失败"); + } + + private: + static void TimerEntry(void* argument) { + auto& timer = *static_cast(argument); + timer.callback_accessing_owner_.store(true); + if (timer.quiescing_.load()) { + bool completes_barrier = false; + portENTER_CRITICAL(&timer.barrier_lock_); + if (timer.barrier_armed_) { + timer.barrier_armed_ = false; + (void)esp_timer_stop(timer.timer_); + completes_barrier = true; + } + portEXIT_CRITICAL(&timer.barrier_lock_); + if (completes_barrier) { + xSemaphoreGive(timer.quiesced_); + } + timer.callback_accessing_owner_.store(false); + return; + } + if (timer.callback_) { + timer.callback_(); + } + timer.callback_accessing_owner_.store(false); + } + + esp_timer_handle_t timer_ = nullptr; + SemaphoreHandle_t quiesced_ = nullptr; + ExpiryCallback callback_; + std::atomic quiescing_{false}; + std::atomic callback_accessing_owner_{false}; + portMUX_TYPE barrier_lock_ = portMUX_INITIALIZER_UNLOCKED; + bool barrier_armed_ = false; + bool quiesced_callback_ = false; +}; + +class FreeRtosRunnerWake final : public RunnerWakePort { + public: + void SetRunnerTask(TaskHandle_t runner_task) { runner_task_ = runner_task; } + + void Notify() override { + configASSERT(runner_task_ != nullptr); + const auto result = xTaskNotifyGive(runner_task_); + configASSERT(result == pdPASS); + (void)result; + } + + private: + TaskHandle_t runner_task_ = nullptr; +}; + +using QueuedCommand = std::variant; + +} // namespace + +class EspTimingTaskRuntime::Impl { + public: + Impl() : runtime_(runner_, clock_, timer_, wake_) {} + + ~Impl() { Stop(); } + + bool Start() { + command_queue_ = xQueueCreate(kCommandQueueDepth, sizeof(QueuedCommand*)); + stopped_ = xSemaphoreCreateBinary(); + submission_mutex_ = xSemaphoreCreateMutex(); + if (command_queue_ == nullptr || stopped_ == nullptr || submission_mutex_ == nullptr || !timer_.Initialize()) { + return false; + } + if (xTaskCreate(&RunnerEntry, "voicelife_timing", kRunnerStackWords, this, kRunnerPriority, &runner_task_) != + pdPASS) { + runner_task_ = nullptr; + return false; + } + wake_.SetRunnerTask(runner_task_); + accepting_.store(true); + return true; + } + + CommandAcceptance RegisterTask(RegisterTaskCommand command) { + if (runner_.IsInCallbackContext()) { + return runner_.RegisterTask(std::move(command)); + } + return Enqueue(std::move(command)); + } + + CommandAcceptance CancelTask(CancelTaskCommand command) { + if (runner_.IsInCallbackContext()) { + return runner_.CancelTask(std::move(command)); + } + return Enqueue(std::move(command)); + } + + private: + static void RunnerEntry(void* argument) { + auto& self = *static_cast(argument); + self.Run(); + vTaskDelete(nullptr); + } + + template + CommandAcceptance Enqueue(Command command) { + if (submission_mutex_ == nullptr || xSemaphoreTake(submission_mutex_, 0) != pdTRUE) { + return CommandAcceptance::kUnavailable; + } + struct SubmissionGuard { + SemaphoreHandle_t mutex; + ~SubmissionGuard() { xSemaphoreGive(mutex); } + } submission_guard{submission_mutex_}; + if (!accepting_.load()) { + return CommandAcceptance::kUnavailable; + } + auto* queued_command = new (std::nothrow) QueuedCommand(std::move(command)); + if (queued_command == nullptr) { + return CommandAcceptance::kUnavailable; + } + if (xQueueSend(command_queue_, &queued_command, 0) != pdPASS) { + delete queued_command; + return CommandAcceptance::kUnavailable; + } + wake_.Notify(); + return CommandAcceptance::kAccepted; + } + + void Run() { + while (true) { + (void)ulTaskNotifyTake(pdTRUE, portMAX_DELAY); + if (stopping_.load()) { + DrainCommands(); + while (runner_.ProcessPendingCommands(clock_.Now()) != 0) { + } + timer_.ClearExpiryCallbackAndWait(); + break; + } + DrainCommands(); + while (!stopping_.load()) { + const auto status = runtime_.ProcessWake(); + if (status.ok()) { + break; + } + ESP_LOGE(kTag, "TIMING_TIMER_ERROR code=%d msg=%s", static_cast(status.code), + status.message.c_str()); + if (status.code == ErrorCode::kInternal) { + std::abort(); + } + vTaskDelay(1); + DrainCommands(); + } + } + xSemaphoreGive(stopped_); + } + + void DrainCommands() { + if (command_queue_ == nullptr) { + return; + } + QueuedCommand* command = nullptr; + while (xQueueReceive(command_queue_, &command, 0) == pdPASS) { + std::visit( + [this](auto&& value) { + using Command = std::decay_t; + if constexpr (std::is_same_v) { + runner_.RegisterTask(std::move(value)); + } else { + runner_.CancelTask(std::move(value)); + } + }, + std::move(*command)); + delete command; + } + } + + void Stop() { + if (runner_task_ != nullptr) { + (void)xSemaphoreTake(submission_mutex_, portMAX_DELAY); + accepting_.store(false); + xSemaphoreGive(submission_mutex_); + stopping_.store(true); + wake_.Notify(); + (void)xSemaphoreTake(stopped_, portMAX_DELAY); + runner_task_ = nullptr; + wake_.SetRunnerTask(nullptr); + } + DrainCommands(); + if (command_queue_ != nullptr) { + vQueueDelete(command_queue_); + command_queue_ = nullptr; + } + if (stopped_ != nullptr) { + vSemaphoreDelete(stopped_); + stopped_ = nullptr; + } + if (submission_mutex_ != nullptr) { + vSemaphoreDelete(submission_mutex_); + submission_mutex_ = nullptr; + } + } + + InMemoryTimingTaskRunner runner_; + EspSystemClock clock_; + EspOneShotTimer timer_; + FreeRtosRunnerWake wake_; + TimingTaskRuntime runtime_; + QueueHandle_t command_queue_ = nullptr; + SemaphoreHandle_t stopped_ = nullptr; + SemaphoreHandle_t submission_mutex_ = nullptr; + TaskHandle_t runner_task_ = nullptr; + std::atomic accepting_{false}; + std::atomic stopping_{false}; +}; + +Result> EspTimingTaskRuntime::Create() { + auto impl = std::make_unique(); + if (!impl->Start()) { + return Result>::Failure(ErrorCode::kUnavailable, + "ESP Timing 平台资源创建失败"); + } + return Result>::Success( + std::unique_ptr(new EspTimingTaskRuntime(std::move(impl)))); +} + +EspTimingTaskRuntime::EspTimingTaskRuntime(std::unique_ptr impl) : impl_(std::move(impl)) {} + +EspTimingTaskRuntime::~EspTimingTaskRuntime() = default; + +CommandAcceptance EspTimingTaskRuntime::RegisterTask(RegisterTaskCommand command) { + return impl_->RegisterTask(std::move(command)); +} + +CommandAcceptance EspTimingTaskRuntime::CancelTask(CancelTaskCommand command) { + return impl_->CancelTask(std::move(command)); +} + +} // namespace voicelife::timing_esp + +#else + +namespace voicelife::timing_esp { + +class EspTimingTaskRuntime::Impl {}; + +Result> EspTimingTaskRuntime::Create() { + return Result>::Failure(ErrorCode::kUnavailable, + "ESP Timing 仅在 ESP 平台可用"); +} + +EspTimingTaskRuntime::EspTimingTaskRuntime(std::unique_ptr impl) : impl_(std::move(impl)) {} + +EspTimingTaskRuntime::~EspTimingTaskRuntime() = default; + +CommandAcceptance EspTimingTaskRuntime::RegisterTask(RegisterTaskCommand) { return CommandAcceptance::kUnavailable; } + +CommandAcceptance EspTimingTaskRuntime::CancelTask(CancelTaskCommand) { return CommandAcceptance::kUnavailable; } + +} // namespace voicelife::timing_esp + +#endif diff --git a/components/voicelife_timing_esp/test/CMakeLists.txt b/components/voicelife_timing_esp/test/CMakeLists.txt new file mode 100644 index 00000000..8b0038d9 --- /dev/null +++ b/components/voicelife_timing_esp/test/CMakeLists.txt @@ -0,0 +1,4 @@ +idf_component_register( + SRCS "esp_timing_runtime_test.cc" + PRIV_REQUIRES unity voicelife_timing_esp +) diff --git a/components/voicelife_timing_esp/test/esp_timing_runtime_test.cc b/components/voicelife_timing_esp/test/esp_timing_runtime_test.cc new file mode 100644 index 00000000..892b4c76 --- /dev/null +++ b/components/voicelife_timing_esp/test/esp_timing_runtime_test.cc @@ -0,0 +1,93 @@ +#include "voicelife/timing_esp/esp_timing_runtime.h" + +#include +#include +#include +#include + +#include +#include +#include + +using namespace std::chrono_literals; +using voicelife::timing::CommandAcceptance; +using voicelife::timing::RegisterTaskCommand; +using voicelife::timing::RegisterTaskResult; +using voicelife::timing::TaskId; +using voicelife::timing::TriggerAt; +using voicelife::timing_esp::EspTimingTaskRuntime; + +namespace { + +TaskId RequireTaskId(std::string value) { + auto task_id = TaskId::Create(std::move(value)); + TEST_ASSERT_TRUE(task_id.has_value()); + return std::move(*task_id); +} + +RegisterTaskCommand Registration(std::string task_id, TriggerAt trigger_at) { + return { + .task_id = RequireTaskId(std::move(task_id)), + .trigger_at = trigger_at, + .callback = [](const TaskId&, TriggerAt) {}, + .on_result = [](RegisterTaskResult) {}, + }; +} + +} // namespace + +TEST_CASE("ESP Timing runtime creates and shuts down", "[timing_esp]") { + auto created = EspTimingTaskRuntime::Create(); + + TEST_ASSERT_TRUE(created.ok()); + TEST_ASSERT_TRUE(created.value.has_value()); + created.value->reset(); +} + +TEST_CASE("ESP Timing queue rejects work beyond its bounded capacity", "[timing_esp]") { + auto created = EspTimingTaskRuntime::Create(); + TEST_ASSERT_TRUE(created.ok()); + auto runtime = std::move(*created.value); + const auto trigger_at = + std::chrono::time_point_cast(std::chrono::system_clock::now()) + 1h; + size_t accepted = 0; + size_t unavailable = 0; + + const auto original_priority = uxTaskPriorityGet(nullptr); + vTaskPrioritySet(nullptr, configMAX_PRIORITIES - 1); + for (size_t index = 0; index < 64; ++index) { + const auto result = runtime->RegisterTask(Registration("queued-" + std::to_string(index), trigger_at)); + accepted += result == CommandAcceptance::kAccepted ? 1 : 0; + unavailable += result == CommandAcceptance::kUnavailable ? 1 : 0; + } + vTaskPrioritySet(nullptr, original_priority); + + TEST_ASSERT_GREATER_THAN(0, accepted); + TEST_ASSERT_GREATER_THAN(0, unavailable); +} + +TEST_CASE("ESP timer wakes the Runner before business callback execution", "[timing_esp]") { + auto created = EspTimingTaskRuntime::Create(); + TEST_ASSERT_TRUE(created.ok()); + auto runtime = std::move(*created.value); + auto* callback_done = xSemaphoreCreateBinary(); + TEST_ASSERT_NOT_NULL(callback_done); + std::atomic registration_applied{false}; + const auto trigger_at = + std::chrono::time_point_cast(std::chrono::system_clock::now()) + 20ms; + + const auto acceptance = runtime->RegisterTask({ + .task_id = RequireTaskId("timer-callback"), + .trigger_at = trigger_at, + .callback = [callback_done](const TaskId&, TriggerAt) { xSemaphoreGive(callback_done); }, + .on_result = + [®istration_applied](RegisterTaskResult result) { + registration_applied.store(result == RegisterTaskResult::kRegistered); + }, + }); + + TEST_ASSERT_EQUAL(static_cast(CommandAcceptance::kAccepted), static_cast(acceptance)); + TEST_ASSERT_EQUAL(pdTRUE, xSemaphoreTake(callback_done, pdMS_TO_TICKS(1000))); + TEST_ASSERT_TRUE(registration_applied.load()); + vSemaphoreDelete(callback_done); +} diff --git a/scripts/check_architecture.cmake b/scripts/check_architecture.cmake index 8c21a206..630af745 100644 --- a/scripts/check_architecture.cmake +++ b/scripts/check_architecture.cmake @@ -13,6 +13,7 @@ set(known_components voicelife_storage_fatfs voicelife_storage_sqlite voicelife_timing + voicelife_timing_esp voicelife_voice voicelife_linx voicelife_linx_esp @@ -98,8 +99,10 @@ assert_dependencies(voicelife_storage_sqlite PUBLIC voicelife_contracts voicelif assert_dependencies(voicelife_storage_sqlite PRIVATE sqlite3) assert_dependencies(voicelife_storage_fatfs PUBLIC voicelife_contracts) assert_dependencies(voicelife_storage_fatfs PRIVATE esp_partition fatfs) -assert_dependencies(voicelife_timing PUBLIC) +assert_dependencies(voicelife_timing PUBLIC voicelife_contracts) assert_dependencies(voicelife_timing PRIVATE) +assert_dependencies(voicelife_timing_esp PUBLIC voicelife_contracts voicelife_timing) +assert_dependencies(voicelife_timing_esp PRIVATE esp_timer freertos) assert_dependencies(voicelife_mcp PUBLIC voicelife_contracts) assert_dependencies(voicelife_mcp PRIVATE yyjson) assert_dependencies(voicelife_voice PUBLIC voicelife_contracts) diff --git a/tests/board/timing_esp_unity/.gitignore b/tests/board/timing_esp_unity/.gitignore new file mode 100644 index 00000000..81b24e74 --- /dev/null +++ b/tests/board/timing_esp_unity/.gitignore @@ -0,0 +1,4 @@ +/build/ +/dependencies.lock +/sdkconfig +/sdkconfig.old diff --git a/tests/board/timing_esp_unity/CMakeLists.txt b/tests/board/timing_esp_unity/CMakeLists.txt new file mode 100644 index 00000000..f1b90384 --- /dev/null +++ b/tests/board/timing_esp_unity/CMakeLists.txt @@ -0,0 +1,13 @@ +cmake_minimum_required(VERSION 3.22) + +list(APPEND EXTRA_COMPONENT_DIRS + "${CMAKE_CURRENT_LIST_DIR}/../../../components/voicelife_contracts" + "${CMAKE_CURRENT_LIST_DIR}/../../../components/voicelife_timing" + "${CMAKE_CURRENT_LIST_DIR}/../../../components/voicelife_timing_esp" + "${CMAKE_CURRENT_LIST_DIR}/../../../components/voicelife_timing_esp/test" + "${CMAKE_CURRENT_LIST_DIR}/../../../third_party/yyjson" +) + +include($ENV{IDF_PATH}/tools/cmake/project.cmake) +idf_build_set_property(MINIMAL_BUILD ON) +project(voicelife_timing_esp_unity) diff --git a/tests/board/timing_esp_unity/main/CMakeLists.txt b/tests/board/timing_esp_unity/main/CMakeLists.txt new file mode 100644 index 00000000..381b7977 --- /dev/null +++ b/tests/board/timing_esp_unity/main/CMakeLists.txt @@ -0,0 +1,4 @@ +idf_component_register( + SRCS "main.cc" + REQUIRES test unity +) diff --git a/tests/board/timing_esp_unity/main/main.cc b/tests/board/timing_esp_unity/main/main.cc new file mode 100644 index 00000000..c174cee9 --- /dev/null +++ b/tests/board/timing_esp_unity/main/main.cc @@ -0,0 +1,3 @@ +#include + +extern "C" void app_main() { unity_run_all_tests(); } diff --git a/tests/board/timing_esp_unity/sdkconfig.defaults b/tests/board/timing_esp_unity/sdkconfig.defaults new file mode 100644 index 00000000..e476a181 --- /dev/null +++ b/tests/board/timing_esp_unity/sdkconfig.defaults @@ -0,0 +1,6 @@ +CONFIG_IDF_TARGET="esp32s3" +CONFIG_FREERTOS_UNICORE=y +CONFIG_FREERTOS_HZ=1000 +CONFIG_COMPILER_OPTIMIZATION_SIZE=y +CONFIG_LOG_DEFAULT_LEVEL_INFO=y +CONFIG_ESP_MAIN_TASK_STACK_SIZE=8192 diff --git a/tests/host/CMakeLists.txt b/tests/host/CMakeLists.txt index 70b7496d..f9e81663 100644 --- a/tests/host/CMakeLists.txt +++ b/tests/host/CMakeLists.txt @@ -78,7 +78,12 @@ add_voicelife_library(schedule voicelife_schedule target_include_directories(schedule PRIVATE "${ROOT_DIR}/components/voicelife_schedule/src") target_link_libraries(schedule PUBLIC contracts Threads::Threads) add_voicelife_library(timing voicelife_timing - "${ROOT_DIR}/components/voicelife_timing/src/timing_task.cc") + "${ROOT_DIR}/components/voicelife_timing/src/timing_task.cc" + "${ROOT_DIR}/components/voicelife_timing/src/timing_runtime.cc") +target_link_libraries(timing PUBLIC contracts) +add_voicelife_library(timing_esp voicelife_timing_esp + "${ROOT_DIR}/components/voicelife_timing_esp/src/esp_timing_runtime.cc") +target_link_libraries(timing_esp PUBLIC timing contracts) add_voicelife_library(mcp voicelife_mcp "${ROOT_DIR}/components/voicelife_mcp/src/mcp_server.cc" "${ROOT_DIR}/components/voicelife_mcp/src/mcp_json_writer.cc") @@ -210,6 +215,13 @@ target_link_libraries(timing_cancellation_test PRIVATE timing) add_voicelife_test(timing_due_tasks_test "unit;timing" timing_due_tasks_test.cc) target_link_libraries(timing_due_tasks_test PRIVATE timing) +add_voicelife_test(timing_runtime_test "unit;timing" timing_runtime_test.cc) +target_link_libraries(timing_runtime_test PRIVATE timing) + +add_voicelife_test(timing_esp_runtime_contract_test "unit;timing;esp;contract" + timing_esp_runtime_contract_test.cc) +target_link_libraries(timing_esp_runtime_contract_test PRIVATE timing_esp) + add_voicelife_test(fatfs_volume_contract_test "unit;storage;fatfs" "${ROOT_DIR}/components/voicelife_storage_fatfs/test/fatfs_volume_contract_test.cc") target_link_libraries(fatfs_volume_contract_test PRIVATE storage_fatfs) diff --git a/tests/host/timing_due_tasks_test.cc b/tests/host/timing_due_tasks_test.cc index 1651c9a6..6bb05f67 100644 --- a/tests/host/timing_due_tasks_test.cc +++ b/tests/host/timing_due_tasks_test.cc @@ -1,6 +1,7 @@ #include #include #include +#include #include #include @@ -153,6 +154,169 @@ void ContinuesDueBatchAfterEarlierCallbackReturns() { Check(result.processed_count == 2, "Runner should finish the due batch after an earlier callback returns"); } +void DefersTaskRegisteredByCallbackUntilNextRun() { + InMemoryTimingTaskRunner runner; + std::vector callback_order; + std::optional nested_registration_result; + bool registration_applied_inside_callback = false; + const auto trigger_at = TriggerAt{1h}; + runner.RegisterTask(Registration( + RequireTaskId("task-1"), trigger_at, + [&runner, &callback_order, &nested_registration_result, ®istration_applied_inside_callback, trigger_at]( + const TaskId& task_id, TriggerAt) { + callback_order.push_back(task_id.Value()); + runner.RegisterTask({ + .task_id = RequireTaskId("task-2"), + .trigger_at = trigger_at, + .callback = [&callback_order](const TaskId& nested_task_id, + TriggerAt) { callback_order.push_back(nested_task_id.Value()); }, + .on_result = + [&nested_registration_result](RegisterTaskResult result) { nested_registration_result = result; }, + }); + registration_applied_inside_callback = nested_registration_result == RegisterTaskResult::kRegistered; + })); + + const auto first_result = runner.RunDueTasks(trigger_at); + + Check(nested_registration_result == RegisterTaskResult::kRegistered, + "callback registration should apply immediately in the Runner context"); + Check(registration_applied_inside_callback, + "callback should observe the final registration result before it returns"); + Check(callback_order == std::vector{"task-1"}, + "task registered after batch formation should not join the current batch"); + Check(first_result.processed_count == 1, "current batch should process only its original task"); + Check(first_result.skipped_count == 0, "callback registration should not skip an original batch task"); + Check(first_result.next_wake_at == trigger_at, + "already-due task registered by callback should become the next wake time"); + + const auto second_result = runner.RunDueTasks(trigger_at); + + Check(callback_order == std::vector{"task-1", "task-2"}, + "task registered by callback should execute on the next run"); + Check(second_result.processed_count == 1, "next run should process the deferred task exactly once"); + Check(!second_result.next_wake_at.has_value(), "completed deferred task should leave no next wake time"); +} + +void SkipsLaterBatchTaskCancelledByCallback() { + InMemoryTimingTaskRunner runner; + std::vector callback_order; + std::optional nested_cancellation_result; + bool cancellation_applied_inside_callback = false; + const auto trigger_at = TriggerAt{1h}; + runner.RegisterTask(Registration( + RequireTaskId("task-a"), trigger_at, + [&runner, &callback_order, &nested_cancellation_result, &cancellation_applied_inside_callback]( + const TaskId& task_id, TriggerAt) { + callback_order.push_back(task_id.Value()); + runner.CancelTask({ + .task_id = RequireTaskId("task-b"), + .on_result = + [&nested_cancellation_result](CancelTaskResult result) { nested_cancellation_result = result; }, + }); + cancellation_applied_inside_callback = nested_cancellation_result == CancelTaskResult::kCancelled; + })); + runner.RegisterTask(Registration( + RequireTaskId("task-b"), trigger_at, + [&callback_order](const TaskId& task_id, TriggerAt) { callback_order.push_back(task_id.Value()); })); + runner.RegisterTask(Registration( + RequireTaskId("task-c"), trigger_at, + [&callback_order](const TaskId& task_id, TriggerAt) { callback_order.push_back(task_id.Value()); })); + + const auto result = runner.RunDueTasks(trigger_at); + + Check(nested_cancellation_result == CancelTaskResult::kCancelled, + "callback cancellation should find a later pending batch task"); + Check(cancellation_applied_inside_callback, + "callback should observe the final cancellation result before it returns"); + Check(callback_order == std::vector{"task-a", "task-c"}, + "cancelled later batch task should be skipped while following tasks continue"); + Check(result.processed_count == 2, "only callbacks that ran should count as processed"); + Check(result.skipped_count == 1, "cancelled batch task should count as skipped"); + Check(!result.next_wake_at.has_value(), "finished batch should leave no next wake time"); +} + +void DoesNotRollbackExecutingTaskWhenCallbackCancelsItself() { + InMemoryTimingTaskRunner runner; + std::vector callback_events; + std::optional cancellation_result; + const auto trigger_at = TriggerAt{1h}; + runner.RegisterTask(Registration( + RequireTaskId("task-1"), trigger_at, + [&runner, &callback_events, &cancellation_result](const TaskId& task_id, TriggerAt) { + callback_events.push_back("callback-start"); + runner.CancelTask({ + .task_id = RequireTaskId(task_id.Value()), + .on_result = [&cancellation_result](CancelTaskResult result) { cancellation_result = result; }, + }); + callback_events.push_back("callback-end"); + })); + + const auto first_result = runner.RunDueTasks(trigger_at); + const auto second_result = runner.RunDueTasks(trigger_at); + + Check(cancellation_result == CancelTaskResult::kNotFound, + "executing task should no longer be cancellable as pending"); + Check(callback_events == std::vector{"callback-start", "callback-end"}, + "self-cancellation should not interrupt or repeat the executing callback"); + Check(first_result.processed_count == 1, "executing task should complete normally"); + Check(first_result.skipped_count == 0, "executing task should not be reclassified as skipped"); + Check(second_result.processed_count == 0, "completed task should not run again after self-cancellation"); +} + +void KeepsLaterDueTaskRegisteredWhileCallbackRuns() { + InMemoryTimingTaskRunner runner; + std::optional next_wake_inside_callback; + const auto trigger_at = TriggerAt{1h}; + runner.RegisterTask(Registration(RequireTaskId("task-a"), trigger_at, + [&runner, &next_wake_inside_callback](const TaskId&, TriggerAt) { + next_wake_inside_callback = runner.NextWakeAt(); + })); + runner.RegisterTask(Registration(RequireTaskId("task-b"), trigger_at, [](const TaskId&, TriggerAt) {})); + + const auto result = runner.RunDueTasks(trigger_at); + + Check(next_wake_inside_callback == trigger_at, + "later due task should remain in the pending registry while an earlier callback runs"); + Check(result.processed_count == 2, "both registered batch tasks should complete"); +} + +void KeepsExternalRegistrationQueuedWhileCallbackRuns() { + InMemoryTimingTaskRunner runner; + std::optional external_registration_result; + bool result_visible_inside_callback = false; + const auto trigger_at = TriggerAt{1h}; + runner.RegisterTask( + Registration(RequireTaskId("task-1"), trigger_at, + [&runner, &external_registration_result, &result_visible_inside_callback, trigger_at]( + const TaskId&, TriggerAt) { + std::thread external_caller([&runner, &external_registration_result, trigger_at] { + runner.RegisterTask({ + .task_id = RequireTaskId("task-2"), + .trigger_at = trigger_at, + .callback = [](const TaskId&, TriggerAt) {}, + .on_result = [&external_registration_result]( + RegisterTaskResult result) { external_registration_result = result; }, + }); + }); + external_caller.join(); + result_visible_inside_callback = external_registration_result.has_value(); + })); + + const auto first_result = runner.RunDueTasks(trigger_at); + + Check(!result_visible_inside_callback, + "external registration should remain queued even while a Runner callback is active"); + Check(!external_registration_result.has_value(), + "external caller should wait for the next Runner command-consumption boundary"); + Check(first_result.processed_count == 1, "external registration should not alter the active batch"); + + const auto second_result = runner.RunDueTasks(trigger_at); + + Check(external_registration_result == RegisterTaskResult::kRegistered, + "next Runner cycle should apply the queued external registration"); + Check(second_result.processed_count == 1, "next Runner cycle should execute the newly due task"); +} + } // namespace int main() { @@ -162,5 +326,10 @@ int main() { DoesNotRepeatCompletedTaskOnLaterRun(); AppliesAcceptedCancellationBeforeFormingDueBatch(); ContinuesDueBatchAfterEarlierCallbackReturns(); + DefersTaskRegisteredByCallbackUntilNextRun(); + SkipsLaterBatchTaskCancelledByCallback(); + DoesNotRollbackExecutingTaskWhenCallbackCancelsItself(); + KeepsLaterDueTaskRegisteredWhileCallbackRuns(); + KeepsExternalRegistrationQueuedWhileCallbackRuns(); return 0; } diff --git a/tests/host/timing_esp_runtime_contract_test.cc b/tests/host/timing_esp_runtime_contract_test.cc new file mode 100644 index 00000000..bab5a541 --- /dev/null +++ b/tests/host/timing_esp_runtime_contract_test.cc @@ -0,0 +1,23 @@ +#include +#include +#include + +#include "support/test_support.h" +#include "voicelife/timing_esp/esp_timing_runtime.h" + +using voicelife::ErrorCode; +using voicelife::test::Check; +using voicelife::timing::TimingTaskService; +using voicelife::timing_esp::EspTimingTaskRuntime; + +static_assert(std::derived_from); +static_assert(std::has_virtual_destructor_v); + +int main() { + const auto created = EspTimingTaskRuntime::Create(); + + Check(created.status.code == ErrorCode::kUnavailable, + "Host must expose an explicit unavailable result instead of pretending ESP resources exist"); + Check(!created.value.has_value(), "failed ESP runtime creation must not return a partial runtime"); + return 0; +} diff --git a/tests/host/timing_runtime_test.cc b/tests/host/timing_runtime_test.cc new file mode 100644 index 00000000..6fccd953 --- /dev/null +++ b/tests/host/timing_runtime_test.cc @@ -0,0 +1,260 @@ +#include "voicelife/timing/timing_runtime.h" + +#include +#include +#include +#include +#include + +#include "support/test_support.h" + +using voicelife::ErrorCode; +using voicelife::Status; +using voicelife::test::Check; +using namespace std::chrono_literals; +using namespace voicelife::timing; + +namespace { + +TaskId RequireTaskId(std::string value) { + auto task_id = TaskId::Create(std::move(value)); + Check(task_id.has_value(), "test task id should be valid"); + return std::move(*task_id); +} + +class FixedClock final : public ClockPort { + public: + explicit FixedClock(TriggerAt now) : now_(now) {} + + TriggerAt Now() override { return now_; } + void SetNow(TriggerAt now) { now_ = now; } + + private: + TriggerAt now_; +}; + +class SequenceClock final : public ClockPort { + public: + explicit SequenceClock(std::vector readings) : readings_(std::move(readings)) {} + + TriggerAt Now() override { + Check(next_reading_ < readings_.size(), "runtime should not read the clock more often than expected"); + return readings_[next_reading_++]; + } + + private: + std::vector readings_; + size_t next_reading_ = 0; +}; + +class RecordingOneShotTimer final : public OneShotTimerPort { + public: + void SetExpiryCallback(ExpiryCallback callback) override { callback_ = std::move(callback); } + void ClearExpiryCallbackAndWait() override { + callback_ = {}; + ++clear_count; + } + Status ArmAfter(std::chrono::microseconds delay) override { + armed_delays.push_back(delay); + return arm_result; + } + Status Disarm() override { + ++disarm_count; + return disarm_result; + } + void Fire() { + Check(static_cast(callback_), "timer expiry callback should be configured"); + callback_(); + } + void FireIfBound() { + if (callback_) { + callback_(); + } + } + + std::vector armed_delays; + size_t disarm_count = 0; + size_t clear_count = 0; + Status arm_result = Status::Ok(); + Status disarm_result = Status::Ok(); + + private: + ExpiryCallback callback_; +}; + +class RecordingRunnerWake final : public RunnerWakePort { + public: + void Notify() override { ++notification_count; } + + size_t notification_count = 0; +}; + +RegisterTaskCommand Registration(std::string task_id, TriggerAt trigger_at) { + return { + .task_id = RequireTaskId(std::move(task_id)), + .trigger_at = trigger_at, + .callback = [](const TaskId&, TriggerAt) {}, + .on_result = [](RegisterTaskResult) {}, + }; +} + +void ArmsOnlyTheEarliestPendingTaskAfterRegistrationWake() { + InMemoryTimingTaskRunner runner; + FixedClock clock(TriggerAt{2h}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + TimingTaskRuntime runtime(runner, clock, timer, wake); + + runtime.RegisterTask(Registration("later", TriggerAt{5h})); + runtime.RegisterTask(Registration("earlier", TriggerAt{3h})); + + Check(wake.notification_count == 2, "each external registration should wake the Runner"); + runtime.ProcessWake(); + + Check(timer.armed_delays == std::vector{1h}, + "Runner should arm one timer for the earliest pending task"); + Check(timer.disarm_count == 0, "pending tasks should keep the one-shot timer armed"); +} + +void DisarmsTimerWhenNoTaskRemainsPending() { + InMemoryTimingTaskRunner runner; + FixedClock clock(TriggerAt{2h}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + TimingTaskRuntime runtime(runner, clock, timer, wake); + Check(timer.ArmAfter(1h).ok(), "test setup should arm the timer"); + + runtime.ProcessWake(); + + Check(timer.disarm_count == 1, "empty pending registry should disarm the one-shot timer"); + Check(timer.armed_delays == std::vector{1h}, + "empty pending registry should not arm another wake"); +} + +void ClampsOverdueDeferredTaskToZeroDelay() { + InMemoryTimingTaskRunner runner; + FixedClock clock(TriggerAt{2h}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + TimingTaskRuntime runtime(runner, clock, timer, wake); + runtime.RegisterTask({ + .task_id = RequireTaskId("task-1"), + .trigger_at = TriggerAt{2h}, + .callback = [&runtime](const TaskId&, + TriggerAt) { runtime.RegisterTask(Registration("task-2", TriggerAt{1h})); }, + .on_result = [](RegisterTaskResult) {}, + }); + + runtime.ProcessWake(); + + Check(timer.armed_delays == std::vector{0us}, + "overdue task deferred by the active batch should arm a zero delay"); +} + +void TimerExpiryOnlyWakesRunnerUntilNormalContextProcessesIt() { + InMemoryTimingTaskRunner runner; + FixedClock clock(TriggerAt{1h}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + TimingTaskRuntime runtime(runner, clock, timer, wake); + size_t callback_count = 0; + runtime.RegisterTask({ + .task_id = RequireTaskId("task-1"), + .trigger_at = TriggerAt{2h}, + .callback = [&callback_count](const TaskId&, TriggerAt) { ++callback_count; }, + .on_result = [](RegisterTaskResult) {}, + }); + runtime.ProcessWake(); + const auto notifications_before_expiry = wake.notification_count; + clock.SetNow(TriggerAt{2h}); + + timer.Fire(); + + Check(wake.notification_count == notifications_before_expiry + 1, "timer expiry should only notify the Runner"); + Check(callback_count == 0, "timer expiry callback should not execute Timing business callbacks"); + + runtime.ProcessWake(); + + Check(callback_count == 1, "normal Runner context should execute the due callback after wake"); +} + +void RefreshesClockAfterCallbacksBeforeArmingNextWake() { + InMemoryTimingTaskRunner runner; + SequenceClock clock({TriggerAt{2h}, TriggerAt{4h}}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + TimingTaskRuntime runtime(runner, clock, timer, wake); + runtime.RegisterTask(Registration("due", TriggerAt{2h})); + runtime.RegisterTask(Registration("became-overdue", TriggerAt{3h})); + + runtime.ProcessWake(); + + Check(timer.armed_delays == std::vector{0us}, + "Runner should recompute delay from the clock after callbacks finish"); +} + +void ClearsTimerCallbackBeforeRuntimeDependenciesCanOutliveIt() { + InMemoryTimingTaskRunner runner; + FixedClock clock(TriggerAt{1h}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + { TimingTaskRuntime runtime(runner, clock, timer, wake); } + + timer.FireIfBound(); + + Check(timer.clear_count == 1, "runtime destruction should wait while clearing the timer callback"); + Check(wake.notification_count == 0, "destroyed runtime should no longer receive timer expiry notifications"); +} + +void ResultCallbackRegistrationQueuesAndWakesTheNextRunnerTurn() { + InMemoryTimingTaskRunner runner; + FixedClock clock(TriggerAt{1h}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + TimingTaskRuntime runtime(runner, clock, timer, wake); + runtime.RegisterTask({ + .task_id = RequireTaskId("first"), + .trigger_at = TriggerAt{5h}, + .callback = [](const TaskId&, TriggerAt) {}, + .on_result = + [&runtime](RegisterTaskResult) { runtime.RegisterTask(Registration("from-result", TriggerAt{3h})); }, + }); + + runtime.ProcessWake(); + + Check(wake.notification_count == 2, "result callback registration should notify a later Runner turn"); + Check(timer.armed_delays == std::vector{4h}, + "current turn should arm from commands applied before its fixed command snapshot"); + + runtime.ProcessWake(); + + Check(timer.armed_delays == std::vector({4h, 2h}), + "next Runner turn should apply the command registered by the result callback"); +} + +void ReportsOneShotTimerArmFailure() { + InMemoryTimingTaskRunner runner; + FixedClock clock(TriggerAt{1h}); + RecordingOneShotTimer timer; + RecordingRunnerWake wake; + TimingTaskRuntime runtime(runner, clock, timer, wake); + runtime.RegisterTask(Registration("task", TriggerAt{2h})); + timer.arm_result = Status::Error(ErrorCode::kUnavailable, "timer busy"); + + const auto result = runtime.ProcessWake(); + Check(result.code == ErrorCode::kUnavailable, "Runner should preserve a platform timer arm failure category"); +} + +} // namespace + +int main() { + ArmsOnlyTheEarliestPendingTaskAfterRegistrationWake(); + DisarmsTimerWhenNoTaskRemainsPending(); + ClampsOverdueDeferredTaskToZeroDelay(); + TimerExpiryOnlyWakesRunnerUntilNormalContextProcessesIt(); + RefreshesClockAfterCallbacksBeforeArmingNextWake(); + ClearsTimerCallbackBeforeRuntimeDependenciesCanOutliveIt(); + ResultCallbackRegistrationQueuesAndWakesTheNextRunnerTurn(); + ReportsOneShotTimerArmFailure(); + return 0; +}