diff --git a/modules/core_sys/include/sys/runtime_async.h b/modules/core_sys/include/sys/runtime_async.h index c1b9e9eb..13e77d0a 100644 --- a/modules/core_sys/include/sys/runtime_async.h +++ b/modules/core_sys/include/sys/runtime_async.h @@ -46,6 +46,15 @@ enum class RuntimeEventKind : uint8_t FeedbackDropped, }; +enum class RuntimeStatus : uint8_t +{ + Idle, + Queued, + Running, + Degraded, + Failed, +}; + enum class RuntimePriority : uint8_t { Realtime = 0, @@ -63,6 +72,18 @@ enum class RuntimeCancelPolicy : uint8_t DropIfStale, }; +struct RuntimeIntent +{ + RuntimeCommandKind kind = RuntimeCommandKind::Unknown; + uint32_t origin = 0; + uint32_t generation = 0; + uint32_t dedupe_key = 0; + uint32_t deadline_ms = 0; + uint32_t submitted_at_ms = 0; + RuntimeCancelPolicy cancel_policy = RuntimeCancelPolicy::None; + RuntimePriority priority_hint = RuntimePriority::Normal; +}; + enum class BusAccessPolicy : uint8_t { UiNeverBlock, @@ -112,6 +133,15 @@ struct RuntimeEvent int32_t error = 0; }; +struct RuntimeState +{ + RuntimeStatus status = RuntimeStatus::Idle; + RuntimeCommandKind active_kind = RuntimeCommandKind::Unknown; + uint32_t active_command_id = 0; + uint32_t last_event_id = 0; + int32_t last_error = 0; +}; + struct BusAcquireRequest { uint32_t resource = 0; @@ -153,6 +183,83 @@ struct StorageHealthState uint32_t last_transition_ms = 0; }; +struct RuntimeRetryDecision +{ + bool retry = false; + uint32_t delay_ms = 0; +}; + +struct PlatformStorageReadRequest +{ + uint32_t command_id = 0; + const char* path = nullptr; + uint8_t* buffer = nullptr; + std::size_t capacity = 0; +}; + +struct PlatformStorageWriteRequest +{ + uint32_t command_id = 0; + const char* path = nullptr; + const uint8_t* data = nullptr; + std::size_t len = 0; + bool durable = false; +}; + +struct PlatformStorageListRequest +{ + uint32_t command_id = 0; + const char* path = nullptr; +}; + +struct PlatformStorageFlushRequest +{ + uint32_t command_id = 0; + uint32_t handle = 0; +}; + +struct PlatformStorageResult +{ + bool ok = false; + std::size_t bytes = 0; + int32_t error = 0; +}; + +struct RuntimeUiEffect +{ + RuntimeEventKind kind = RuntimeEventKind::Unknown; + uint32_t event_id = 0; + uint32_t command_id = 0; + int32_t error = 0; +}; + +class ICommandQueue +{ + public: + virtual ~ICommandQueue() = default; + + virtual bool enqueue(const RuntimeCommand& command) = 0; + virtual std::size_t cancel(uint32_t dedupe_key) = 0; + virtual bool popReady(uint32_t now_ms, RuntimeCommand& out) = 0; +}; + +class IEventSink +{ + public: + virtual ~IEventSink() = default; + + virtual bool publish(const RuntimeEvent& event) = 0; +}; + +class IActiveWorker +{ + public: + virtual ~IActiveWorker() = default; + + virtual void tick(uint32_t now_ms) = 0; + virtual bool submit(const RuntimeCommand& command) = 0; +}; + class IBusArbiter { public: @@ -163,6 +270,25 @@ class IBusArbiter virtual StorageHealthState health() const = 0; }; +class IPlatformStorageAdapter +{ + public: + virtual ~IPlatformStorageAdapter() = default; + + virtual PlatformStorageResult read(const PlatformStorageReadRequest& request) = 0; + virtual PlatformStorageResult write(const PlatformStorageWriteRequest& request) = 0; + virtual PlatformStorageResult list(const PlatformStorageListRequest& request) = 0; + virtual PlatformStorageResult flush(const PlatformStorageFlushRequest& request) = 0; +}; + +class IUiEffectSink +{ + public: + virtual ~IUiEffectSink() = default; + + virtual bool apply(const RuntimeUiEffect& effect) = 0; +}; + class IUiOwnerGuard { public: @@ -172,6 +298,190 @@ class IUiOwnerGuard virtual void recordForbiddenBlockingCall(RuntimeCommandKind kind) = 0; }; +class RuntimePolicyStrategy +{ + public: + virtual ~RuntimePolicyStrategy() = default; + + virtual RuntimePriority selectPriority(const RuntimeIntent& intent) const = 0; + virtual BusAccessPolicy selectBusPolicy(const RuntimeCommand& command) const = 0; + virtual RuntimeRetryDecision selectRetry(const RuntimeCommand& command, + const PlatformStorageResult& result) const = 0; +}; + +class DefaultRuntimePolicyStrategy : public RuntimePolicyStrategy +{ + public: + RuntimePriority selectPriority(const RuntimeIntent& intent) const override + { + return intent.priority_hint; + } + + BusAccessPolicy selectBusPolicy(const RuntimeCommand& command) const override + { + if (command.priority == RuntimePriority::Realtime || + command.priority == RuntimePriority::Interactive) + { + return BusAccessPolicy::InteractiveWorkerBounded; + } + if (command.priority == RuntimePriority::Idle || + command.priority == RuntimePriority::Background) + { + return BusAccessPolicy::BackgroundWorkerBounded; + } + return BusAccessPolicy::BackgroundWorkerBounded; + } + + RuntimeRetryDecision selectRetry(const RuntimeCommand& command, + const PlatformStorageResult& result) const override + { + RuntimeRetryDecision decision{}; + if (!result.ok && command.cancel_policy != RuntimeCancelPolicy::DropIfStale) + { + decision.retry = true; + decision.delay_ms = 250; + } + return decision; + } +}; + +class RuntimeFacade +{ + public: + RuntimeFacade(ICommandQueue& commands, + IActiveWorker& worker, + IEventSink& events, + RuntimeState& state, + RuntimePolicyStrategy& policy) + : commands_(commands), worker_(worker), events_(events), state_(state), policy_(policy) + { + } + + bool submit(const RuntimeIntent& intent) + { + RuntimeCommand command{}; + command.command_id = next_command_id_++; + command.kind = intent.kind; + command.priority = policy_.selectPriority(intent); + command.cancel_policy = intent.cancel_policy; + command.created_at_ms = intent.submitted_at_ms; + command.deadline_ms = intent.deadline_ms; + command.generation = intent.generation; + command.dedupe_key = intent.dedupe_key; + command.origin = intent.origin; + + if (command.cancel_policy == RuntimeCancelPolicy::CancelByDedupeKey && + command.dedupe_key != 0) + { + (void)commands_.cancel(command.dedupe_key); + } + + const bool queued = commands_.enqueue(command); + RuntimeEvent event{}; + event.event_id = next_event_id_++; + event.kind = queued ? RuntimeEventKind::CommandQueued : RuntimeEventKind::CommandFailed; + event.command_id = command.command_id; + event.timestamp_ms = intent.submitted_at_ms; + event.generation = command.generation; + event.error = queued ? 0 : -1; + (void)publishEvent(event); + state_.status = queued ? RuntimeStatus::Queued : RuntimeStatus::Failed; + state_.active_kind = intent.kind; + state_.active_command_id = command.command_id; + state_.last_event_id = event.event_id; + state_.last_error = event.error; + return queued; + } + + void tick(uint32_t now_ms) + { + RuntimeCommand command{}; + while (commands_.popReady(now_ms, command)) + { + state_.status = RuntimeStatus::Running; + state_.active_kind = command.kind; + state_.active_command_id = command.command_id; + + RuntimeEvent started{}; + started.event_id = next_event_id_++; + started.kind = RuntimeEventKind::CommandStarted; + started.command_id = command.command_id; + started.timestamp_ms = now_ms; + started.generation = command.generation; + (void)publishEvent(started); + state_.last_event_id = started.event_id; + + if (!worker_.submit(command)) + { + RuntimeEvent failed{}; + failed.event_id = next_event_id_++; + failed.kind = RuntimeEventKind::CommandFailed; + failed.command_id = command.command_id; + failed.timestamp_ms = now_ms; + failed.generation = command.generation; + failed.error = -1; + (void)publishEvent(failed); + state_.status = RuntimeStatus::Failed; + state_.last_event_id = failed.event_id; + state_.last_error = failed.error; + } + } + worker_.tick(now_ms); + } + + std::size_t drainEvents() + { + const std::size_t drained = pending_event_count_; + pending_event_count_ = 0; + return drained; + } + + private: + bool publishEvent(const RuntimeEvent& event) + { + const bool ok = events_.publish(event); + if (ok) + { + ++pending_event_count_; + } + return ok; + } + + ICommandQueue& commands_; + IActiveWorker& worker_; + IEventSink& events_; + RuntimeState& state_; + RuntimePolicyStrategy& policy_; + uint32_t next_command_id_ = 1; + uint32_t next_event_id_ = 1; + std::size_t pending_event_count_ = 0; +}; + +class RuntimeEventUiEffectBridge : public IEventSink +{ + public: + RuntimeEventUiEffectBridge(IEventSink& events, IUiEffectSink& effects) + : events_(events), effects_(effects) + { + } + + bool publish(const RuntimeEvent& event) override + { + const bool published = events_.publish(event); + RuntimeUiEffect effect{}; + effect.kind = event.kind; + effect.event_id = event.event_id; + effect.command_id = event.command_id; + effect.error = event.error; + (void)effects_.apply(effect); + return published; + } + + private: + IEventSink& events_; + IUiEffectSink& effects_; +}; + template class FixedRuntimeQueue { @@ -255,10 +565,10 @@ class FixedRuntimeQueue }; template -class FixedCommandQueue +class FixedCommandQueue : public ICommandQueue { public: - bool enqueue(const RuntimeCommand& command) + bool enqueue(const RuntimeCommand& command) override { return queue_.enqueue(command); } @@ -343,6 +653,11 @@ class FixedCommandQueue return true; } + bool popReady(uint32_t now_ms, RuntimeCommand& out) override + { + return popNext(now_ms, out); + } + std::size_t cancelGeneration(uint32_t generation) { return removeMatching(generation, 0, true); @@ -353,6 +668,11 @@ class FixedCommandQueue return removeMatching(0, dedupe_key, false); } + std::size_t cancel(uint32_t dedupe_key) override + { + return cancelDedupeKey(dedupe_key); + } + std::size_t size() const { return queue_.size(); @@ -404,5 +724,28 @@ class FixedCommandQueue template using FixedEventQueue = FixedRuntimeQueue; +template +class FixedEventSink : public IEventSink +{ + public: + bool publish(const RuntimeEvent& event) override + { + return queue_.enqueue(event); + } + + bool pop(RuntimeEvent& out) + { + return queue_.pop(out); + } + + std::size_t size() const + { + return queue_.size(); + } + + private: + FixedEventQueue queue_{}; +}; + } // namespace runtime } // namespace sys diff --git a/modules/core_sys/include/sys/storage_event_runtime.h b/modules/core_sys/include/sys/storage_event_runtime.h index 8f0f9794..2e84affc 100644 --- a/modules/core_sys/include/sys/storage_event_runtime.h +++ b/modules/core_sys/include/sys/storage_event_runtime.h @@ -1,5 +1,7 @@ #pragma once +#include "sys/runtime_async.h" + #include #include #include @@ -90,6 +92,38 @@ class LatestSnapshotStorageRuntime state_ = StorageWorkState::FailedPendingRetry; } + bool flushPending(IPlatformStorageAdapter& storage, + IEventSink& events, + const char* path, + uint32_t now_ms) + { + StorageWorkItem work{}; + if (!takeNext(work)) + { + return true; + } + + PlatformStorageWriteRequest request{}; + request.command_id = work.generation; + request.path = path; + request.data = work.data; + request.len = work.len; + request.durable = true; + const PlatformStorageResult result = storage.write(request); + complete(work.generation, result.ok); + + RuntimeEvent event{}; + event.event_id = work.generation; + event.kind = result.ok ? RuntimeEventKind::PersistenceSaved + : RuntimeEventKind::PersistenceFailed; + event.command_id = work.generation; + event.timestamp_ms = now_ms; + event.generation = work.generation; + event.error = result.error; + (void)events.publish(event); + return result.ok; + } + bool pending() const { return pending_; } bool busy() const { return state_ == StorageWorkState::InFlight; } StorageWorkState state() const { return state_; } diff --git a/modules/core_sys/tests/test_runtime_async.cpp b/modules/core_sys/tests/test_runtime_async.cpp index bd8172d5..65b08f5c 100644 --- a/modules/core_sys/tests/test_runtime_async.cpp +++ b/modules/core_sys/tests/test_runtime_async.cpp @@ -5,6 +5,41 @@ namespace { +class FakeWorker final : public sys::runtime::IActiveWorker +{ + public: + void tick(uint32_t now_ms) override + { + last_tick_ms = now_ms; + } + + bool submit(const sys::runtime::RuntimeCommand& command) override + { + submitted = true; + last_command = command; + return accept; + } + + bool accept = true; + bool submitted = false; + uint32_t last_tick_ms = 0; + sys::runtime::RuntimeCommand last_command{}; +}; + +class FakeUiEffectSink final : public sys::runtime::IUiEffectSink +{ + public: + bool apply(const sys::runtime::RuntimeUiEffect& effect) override + { + applied = true; + last_effect = effect; + return true; + } + + bool applied = false; + sys::runtime::RuntimeUiEffect last_effect{}; +}; + void test_priority_pop_order() { sys::runtime::FixedCommandQueue<4> queue; @@ -90,6 +125,104 @@ void test_event_queue_bounds() assert(out.event_id == 1); } +void test_runtime_facade_submit_tick() +{ + sys::runtime::FixedCommandQueue<4> commands; + sys::runtime::FixedEventSink<4> events; + sys::runtime::RuntimeState state; + sys::runtime::DefaultRuntimePolicyStrategy policy; + FakeWorker worker; + sys::runtime::RuntimeFacade facade(commands, worker, events, state, policy); + + sys::runtime::RuntimeIntent intent{}; + intent.kind = sys::runtime::RuntimeCommandKind::PersistenceSave; + intent.origin = 7; + intent.generation = 11; + intent.dedupe_key = 99; + intent.submitted_at_ms = 100; + intent.priority_hint = sys::runtime::RuntimePriority::Background; + intent.cancel_policy = sys::runtime::RuntimeCancelPolicy::CancelByDedupeKey; + + assert(facade.submit(intent)); + assert(state.status == sys::runtime::RuntimeStatus::Queued); + assert(commands.size() == 1); + assert(events.size() == 1); + assert(facade.drainEvents() == 1); + assert(facade.drainEvents() == 0); + + sys::runtime::RuntimeEvent queued{}; + assert(events.pop(queued)); + assert(queued.kind == sys::runtime::RuntimeEventKind::CommandQueued); + assert(queued.command_id == 1); + + facade.tick(120); + assert(worker.submitted); + assert(worker.last_tick_ms == 120); + assert(worker.last_command.kind == intent.kind); + assert(worker.last_command.origin == intent.origin); + assert(worker.last_command.generation == intent.generation); + assert(worker.last_command.dedupe_key == intent.dedupe_key); + assert(worker.last_command.priority == intent.priority_hint); + assert(state.status == sys::runtime::RuntimeStatus::Running); + + sys::runtime::RuntimeEvent started{}; + assert(events.pop(started)); + assert(started.kind == sys::runtime::RuntimeEventKind::CommandStarted); + assert(started.command_id == worker.last_command.command_id); + assert(facade.drainEvents() == 1); +} + +void test_runtime_facade_dedupe_cancel_policy() +{ + sys::runtime::FixedCommandQueue<4> commands; + sys::runtime::FixedEventSink<4> events; + sys::runtime::RuntimeState state; + sys::runtime::DefaultRuntimePolicyStrategy policy; + FakeWorker worker; + sys::runtime::RuntimeFacade facade(commands, worker, events, state, policy); + + sys::runtime::RuntimeIntent first{}; + first.kind = sys::runtime::RuntimeCommandKind::PersistenceSave; + first.dedupe_key = 44; + first.cancel_policy = sys::runtime::RuntimeCancelPolicy::CancelByDedupeKey; + first.submitted_at_ms = 1; + + sys::runtime::RuntimeIntent second = first; + second.submitted_at_ms = 2; + + assert(facade.submit(first)); + assert(facade.submit(second)); + assert(commands.size() == 1); + + facade.tick(10); + assert(worker.submitted); + assert(worker.last_command.command_id == 2); +} + +void test_event_to_ui_effect_bridge() +{ + sys::runtime::FixedEventSink<2> events; + FakeUiEffectSink effects; + sys::runtime::RuntimeEventUiEffectBridge bridge(events, effects); + + sys::runtime::RuntimeEvent event{}; + event.event_id = 7; + event.kind = sys::runtime::RuntimeEventKind::PersistenceFailed; + event.command_id = 3; + event.error = -9; + + assert(bridge.publish(event)); + assert(effects.applied); + assert(effects.last_effect.kind == event.kind); + assert(effects.last_effect.event_id == event.event_id); + assert(effects.last_effect.command_id == event.command_id); + assert(effects.last_effect.error == event.error); + + sys::runtime::RuntimeEvent out{}; + assert(events.pop(out)); + assert(out.event_id == event.event_id); +} + } // namespace int main() @@ -98,5 +231,8 @@ int main() test_dedupe_replace(); test_cancel_generation(); test_event_queue_bounds(); + test_runtime_facade_submit_tick(); + test_runtime_facade_dedupe_cancel_policy(); + test_event_to_ui_effect_bridge(); return 0; } diff --git a/modules/core_sys/tests/test_storage_event_runtime.cpp b/modules/core_sys/tests/test_storage_event_runtime.cpp index b8a8fe8b..46749182 100644 --- a/modules/core_sys/tests/test_storage_event_runtime.cpp +++ b/modules/core_sys/tests/test_storage_event_runtime.cpp @@ -10,6 +10,51 @@ using sys::runtime::StorageWorkState; namespace { +class FakeStorageAdapter final : public sys::runtime::IPlatformStorageAdapter +{ + public: + sys::runtime::PlatformStorageResult read( + const sys::runtime::PlatformStorageReadRequest& request) override + { + (void)request; + return {}; + } + + sys::runtime::PlatformStorageResult write( + const sys::runtime::PlatformStorageWriteRequest& request) override + { + writes += 1; + last_path = request.path; + last_len = request.len; + sys::runtime::PlatformStorageResult result{}; + result.ok = write_ok; + result.bytes = write_ok ? request.len : 0; + result.error = write_ok ? 0 : -7; + return result; + } + + sys::runtime::PlatformStorageResult list( + const sys::runtime::PlatformStorageListRequest& request) override + { + (void)request; + return {}; + } + + sys::runtime::PlatformStorageResult flush( + const sys::runtime::PlatformStorageFlushRequest& request) override + { + (void)request; + sys::runtime::PlatformStorageResult result{}; + result.ok = true; + return result; + } + + bool write_ok = true; + int writes = 0; + const char* last_path = nullptr; + std::size_t last_len = 0; +}; + void burst_updates_coalesce_to_latest_snapshot() { LatestSnapshotStorageRuntime<16> runtime; @@ -66,6 +111,46 @@ void failed_save_is_retried_with_same_generation() assert(std::memcmp(retry.data, blob, sizeof(blob)) == 0); } +void flush_pending_writes_through_storage_port_and_publishes_event() +{ + LatestSnapshotStorageRuntime<16> runtime; + FakeStorageAdapter storage; + sys::runtime::FixedEventSink<4> events; + const uint8_t blob[] = {1, 3, 5, 7}; + + assert(runtime.requestSave(12, blob, sizeof(blob))); + assert(runtime.flushPending(storage, events, "/nodes.bin", 50)); + assert(storage.writes == 1); + assert(storage.last_path != nullptr); + assert(std::strcmp(storage.last_path, "/nodes.bin") == 0); + assert(storage.last_len == sizeof(blob)); + assert(!runtime.pending()); + + sys::runtime::RuntimeEvent event{}; + assert(events.pop(event)); + assert(event.kind == sys::runtime::RuntimeEventKind::PersistenceSaved); + assert(event.timestamp_ms == 50); +} + +void failed_flush_keeps_latest_snapshot_pending_for_retry() +{ + LatestSnapshotStorageRuntime<16> runtime; + FakeStorageAdapter storage; + storage.write_ok = false; + sys::runtime::FixedEventSink<4> events; + const uint8_t blob[] = {2, 4, 6}; + + assert(runtime.requestSave(9, blob, sizeof(blob))); + assert(!runtime.flushPending(storage, events, "/nodes.bin", 60)); + assert(runtime.pending()); + assert(runtime.state() == StorageWorkState::FailedPendingRetry); + + sys::runtime::RuntimeEvent event{}; + assert(events.pop(event)); + assert(event.kind == sys::runtime::RuntimeEventKind::PersistenceFailed); + assert(event.error == -7); +} + } // namespace int main() @@ -73,5 +158,7 @@ int main() burst_updates_coalesce_to_latest_snapshot(); completion_ignores_stale_generation(); failed_save_is_retried_with_same_generation(); + flush_pending_writes_through_storage_port_and_publishes_event(); + failed_flush_keeps_latest_snapshot_pending_for_retry(); return 0; } diff --git a/modules/ui_map_runtime/include/ui_map_runtime/map_tiles/map_tile_async_runtime.h b/modules/ui_map_runtime/include/ui_map_runtime/map_tiles/map_tile_async_runtime.h index 3ce9585f..20c3f042 100644 --- a/modules/ui_map_runtime/include/ui_map_runtime/map_tiles/map_tile_async_runtime.h +++ b/modules/ui_map_runtime/include/ui_map_runtime/map_tiles/map_tile_async_runtime.h @@ -86,7 +86,8 @@ class IMapTileWorkerBackend class MapTileAsyncRuntime { public: - explicit MapTileAsyncRuntime(IMapTileCommandSink& commands); + explicit MapTileAsyncRuntime(IMapTileCommandSink& commands, + sys::runtime::RuntimePolicyStrategy* policy = nullptr); uint32_t activeGeneration() const; std::size_t requestVisibleTiles(const MapViewportPlan& plan, uint32_t now_ms); @@ -94,8 +95,12 @@ class MapTileAsyncRuntime private: static sys::runtime::RuntimePriority priorityFor(MapTileInteractionMode mode); + sys::runtime::RuntimeCommand commandFromIntent(const sys::runtime::RuntimeIntent& intent); IMapTileCommandSink& commands_; + sys::runtime::DefaultRuntimePolicyStrategy default_policy_{}; + sys::runtime::RuntimePolicyStrategy* policy_ = nullptr; + sys::runtime::RuntimeState state_{}; uint32_t active_generation_ = 0; uint32_t next_command_id_ = 1; }; @@ -107,16 +112,17 @@ class MapTileWorker sys::runtime::IBusArbiter& bus, IMapTileEventSink& events, uint8_t* scratch, - std::size_t scratch_size); + std::size_t scratch_size, + sys::runtime::RuntimePolicyStrategy* policy = nullptr); bool execute(const LoadTileCommand& command, uint32_t now_ms); private: - static sys::runtime::BusAccessPolicy busPolicyFor(sys::runtime::RuntimePriority priority); - IMapTileWorkerBackend& backend_; sys::runtime::IBusArbiter& bus_; IMapTileEventSink& events_; + sys::runtime::DefaultRuntimePolicyStrategy default_policy_{}; + sys::runtime::RuntimePolicyStrategy* policy_ = nullptr; uint8_t* scratch_ = nullptr; std::size_t scratch_size_ = 0; }; diff --git a/modules/ui_map_runtime/src/map_tiles/map_tile_async_runtime.cpp b/modules/ui_map_runtime/src/map_tiles/map_tile_async_runtime.cpp index 144d0dcc..46a34348 100644 --- a/modules/ui_map_runtime/src/map_tiles/map_tile_async_runtime.cpp +++ b/modules/ui_map_runtime/src/map_tiles/map_tile_async_runtime.cpp @@ -5,8 +5,9 @@ namespace ui namespace map_tiles { -MapTileAsyncRuntime::MapTileAsyncRuntime(IMapTileCommandSink& commands) - : commands_(commands) +MapTileAsyncRuntime::MapTileAsyncRuntime(IMapTileCommandSink& commands, + sys::runtime::RuntimePolicyStrategy* policy) + : commands_(commands), policy_(policy ? policy : &default_policy_) { } @@ -31,23 +32,29 @@ std::size_t MapTileAsyncRuntime::requestVisibleTiles(const MapViewportPlan& plan : MapViewportPlan::kMaxTiles; for (std::size_t i = 0; i < count; ++i) { + sys::runtime::RuntimeIntent intent{}; + intent.kind = sys::runtime::RuntimeCommandKind::MapTileLoad; + intent.priority_hint = priorityFor(plan.interaction_mode); + intent.cancel_policy = sys::runtime::RuntimeCancelPolicy::CancelByGeneration; + intent.submitted_at_ms = now_ms; + intent.generation = plan.generation; + intent.origin = static_cast(plan.interaction_mode); + intent.dedupe_key = + (static_cast(plan.tiles[i].z & 0xFF) << 24) ^ + ((plan.tiles[i].x & 0xFFFu) << 12) ^ + (plan.tiles[i].y & 0xFFFu) ^ + (static_cast(plan.tiles[i].layer) << 28); + LoadTileCommand command{}; command.tile = plan.tiles[i]; - command.runtime.command_id = next_command_id_++; - command.runtime.kind = sys::runtime::RuntimeCommandKind::MapTileLoad; - command.runtime.priority = priorityFor(plan.interaction_mode); - command.runtime.cancel_policy = sys::runtime::RuntimeCancelPolicy::CancelByGeneration; - command.runtime.created_at_ms = now_ms; - command.runtime.generation = plan.generation; - command.runtime.dedupe_key = - (static_cast(command.tile.z & 0xFF) << 24) ^ - ((command.tile.x & 0xFFFu) << 12) ^ - (command.tile.y & 0xFFFu) ^ - (static_cast(command.tile.layer) << 28); + command.runtime = commandFromIntent(intent); if (commands_.enqueue(command)) { ++queued; + state_.status = sys::runtime::RuntimeStatus::Queued; + state_.active_kind = command.runtime.kind; + state_.active_command_id = command.runtime.command_id; } } return queued; @@ -87,19 +94,41 @@ sys::runtime::RuntimePriority MapTileAsyncRuntime::priorityFor(MapTileInteractio : sys::runtime::RuntimePriority::Normal; } +sys::runtime::RuntimeCommand MapTileAsyncRuntime::commandFromIntent( + const sys::runtime::RuntimeIntent& intent) +{ + sys::runtime::RuntimeCommand command{}; + command.command_id = next_command_id_++; + command.kind = intent.kind; + command.priority = policy_->selectPriority(intent); + command.cancel_policy = intent.cancel_policy; + command.created_at_ms = intent.submitted_at_ms; + command.deadline_ms = intent.deadline_ms; + command.generation = intent.generation; + command.dedupe_key = intent.dedupe_key; + command.origin = intent.origin; + return command; +} + MapTileWorker::MapTileWorker(IMapTileWorkerBackend& backend, sys::runtime::IBusArbiter& bus, IMapTileEventSink& events, uint8_t* scratch, - std::size_t scratch_size) - : backend_(backend), bus_(bus), events_(events), scratch_(scratch), scratch_size_(scratch_size) + std::size_t scratch_size, + sys::runtime::RuntimePolicyStrategy* policy) + : backend_(backend), + bus_(bus), + events_(events), + policy_(policy ? policy : &default_policy_), + scratch_(scratch), + scratch_size_(scratch_size) { } bool MapTileWorker::execute(const LoadTileCommand& command, uint32_t now_ms) { sys::runtime::BusAcquireRequest request{}; - request.policy = busPolicyFor(command.runtime.priority); + request.policy = policy_->selectBusPolicy(command.runtime); request.command_id = command.runtime.command_id; request.deadline_ms = command.runtime.deadline_ms; @@ -146,12 +175,5 @@ bool MapTileWorker::execute(const LoadTileCommand& command, uint32_t now_ms) return ok; } -sys::runtime::BusAccessPolicy MapTileWorker::busPolicyFor(sys::runtime::RuntimePriority priority) -{ - return priority == sys::runtime::RuntimePriority::Interactive - ? sys::runtime::BusAccessPolicy::InteractiveWorkerBounded - : sys::runtime::BusAccessPolicy::BackgroundWorkerBounded; -} - } // namespace map_tiles } // namespace ui diff --git a/modules/ui_map_runtime/tests/test_map_tile_async_runtime.cpp b/modules/ui_map_runtime/tests/test_map_tile_async_runtime.cpp index 36982bc2..7fd4c9d7 100644 --- a/modules/ui_map_runtime/tests/test_map_tile_async_runtime.cpp +++ b/modules/ui_map_runtime/tests/test_map_tile_async_runtime.cpp @@ -143,6 +143,33 @@ class FakeEventSink final : public ui::map_tiles::IMapTileEventSink } }; +class FakePolicy final : public sys::runtime::RuntimePolicyStrategy +{ + public: + sys::runtime::RuntimePriority selectPriority( + const sys::runtime::RuntimeIntent& intent) const override + { + (void)intent; + return sys::runtime::RuntimePriority::Realtime; + } + + sys::runtime::BusAccessPolicy selectBusPolicy( + const sys::runtime::RuntimeCommand& command) const override + { + (void)command; + return sys::runtime::BusAccessPolicy::RecoveryExclusive; + } + + sys::runtime::RuntimeRetryDecision selectRetry( + const sys::runtime::RuntimeCommand& command, + const sys::runtime::PlatformStorageResult& result) const override + { + (void)command; + (void)result; + return {}; + } +}; + void test_generation_cancels_old_commands() { FakeCommandSink sink; @@ -236,6 +263,28 @@ void test_worker_success_publishes_ready() assert(events.events[0].payload_size == 3); } +void test_runtime_and_worker_use_policy_strategy() +{ + FakeCommandSink sink; + FakePolicy policy; + ui::map_tiles::MapTileAsyncRuntime runtime(sink, &policy); + + ui::map_tiles::MapViewportPlan plan{}; + plan.generation = 5; + plan.tile_count = 1; + plan.tiles[0] = make_tile(50); + assert(runtime.requestVisibleTiles(plan, 500) == 1); + assert(sink.commands[0].runtime.priority == sys::runtime::RuntimePriority::Realtime); + + FakeBackend backend; + FakeBusArbiter bus; + FakeEventSink events; + uint8_t scratch[8]{}; + ui::map_tiles::MapTileWorker worker(backend, bus, events, scratch, sizeof(scratch), &policy); + assert(worker.execute(sink.commands[0], 510)); + assert(bus.last_policy == sys::runtime::BusAccessPolicy::RecoveryExclusive); +} + } // namespace int main() @@ -244,5 +293,6 @@ int main() test_stale_event_is_ignored(); test_worker_busy_does_not_read_storage(); test_worker_success_publishes_ready(); + test_runtime_and_worker_use_policy_strategy(); return 0; }