mirror of
https://github.com/vicliu624/trail-mate.git
synced 2026-07-20 02:21:09 +00:00
Align UI storage event runtime contracts
This commit is contained in:
@@ -360,6 +360,20 @@ if(BUILD_TESTING)
|
||||
add_test(NAME trailmate_map_tile_async_runtime_smoke
|
||||
COMMAND trailmate_map_tile_async_runtime_smoke)
|
||||
|
||||
add_executable(trailmate_ui_storage_event_runtime_contract_smoke
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/core_sys/tests/test_ui_storage_event_runtime_contract.cpp"
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/ui_map_runtime/src/map_tiles/map_tile_async_runtime.cpp"
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/ui_map_runtime/src/map_tiles/map_tile_render_queue.cpp")
|
||||
target_include_directories(trailmate_ui_storage_event_runtime_contract_smoke
|
||||
PRIVATE
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/core_sys/include"
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/core_gps/include"
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/ui_map_runtime/include")
|
||||
target_compile_features(trailmate_ui_storage_event_runtime_contract_smoke
|
||||
PRIVATE cxx_std_17)
|
||||
add_test(NAME trailmate_ui_storage_event_runtime_contract_smoke
|
||||
COMMAND trailmate_ui_storage_event_runtime_contract_smoke)
|
||||
|
||||
add_executable(trailmate_runtime_map_presentation_adapters_smoke
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/ui_shared/tests/test_runtime_map_presentation_adapters.cpp"
|
||||
"${TRAIL_MATE_REPO_ROOT}/modules/ui_shared/src/ui/presentation_sources/runtime_map_workspace_source.cpp"
|
||||
|
||||
@@ -0,0 +1,440 @@
|
||||
#pragma once
|
||||
|
||||
#include "sys/runtime_async.h"
|
||||
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
|
||||
namespace gps
|
||||
{
|
||||
namespace runtime
|
||||
{
|
||||
|
||||
struct TrackPoint
|
||||
{
|
||||
double latitude = 0.0;
|
||||
double longitude = 0.0;
|
||||
double altitude = 0.0;
|
||||
uint32_t timestamp_ms = 0;
|
||||
bool has_altitude = false;
|
||||
};
|
||||
|
||||
enum class TrackCommandKind : uint8_t
|
||||
{
|
||||
StartNewTrack,
|
||||
StopTrack,
|
||||
AppendPoint,
|
||||
Flush,
|
||||
ListTracks,
|
||||
};
|
||||
|
||||
enum class TrackEventKind : uint8_t
|
||||
{
|
||||
Started,
|
||||
Stopped,
|
||||
PointBuffered,
|
||||
FlushSucceeded,
|
||||
FlushFailed,
|
||||
ListReady,
|
||||
Failed,
|
||||
};
|
||||
|
||||
enum class TrackRecorderStatus : uint8_t
|
||||
{
|
||||
Idle,
|
||||
Starting,
|
||||
Recording,
|
||||
Flushing,
|
||||
Stopping,
|
||||
Stopped,
|
||||
Error,
|
||||
Recovering,
|
||||
};
|
||||
|
||||
struct TrackPointBatch
|
||||
{
|
||||
const TrackPoint* points = nullptr;
|
||||
std::size_t count = 0;
|
||||
};
|
||||
|
||||
struct TrackCommand
|
||||
{
|
||||
uint32_t command_id = 0;
|
||||
TrackCommandKind kind = TrackCommandKind::AppendPoint;
|
||||
uint32_t track_id = 0;
|
||||
TrackPointBatch point_batch{};
|
||||
uint32_t deadline_ms = 0;
|
||||
uint32_t created_at_ms = 0;
|
||||
};
|
||||
|
||||
struct TrackEvent
|
||||
{
|
||||
TrackEventKind kind = TrackEventKind::Failed;
|
||||
uint32_t command_id = 0;
|
||||
uint32_t track_id = 0;
|
||||
uint32_t timestamp_ms = 0;
|
||||
int32_t error = 0;
|
||||
};
|
||||
|
||||
class TrackFlushPolicy
|
||||
{
|
||||
public:
|
||||
virtual ~TrackFlushPolicy() = default;
|
||||
|
||||
virtual bool shouldFlush(std::size_t buffer_size, uint32_t now_ms) const = 0;
|
||||
virtual bool isCritical(const TrackCommand& command) const = 0;
|
||||
};
|
||||
|
||||
class DefaultTrackFlushPolicy : public TrackFlushPolicy
|
||||
{
|
||||
public:
|
||||
bool shouldFlush(std::size_t buffer_size, uint32_t now_ms) const override
|
||||
{
|
||||
(void)now_ms;
|
||||
return buffer_size >= 8;
|
||||
}
|
||||
|
||||
bool isCritical(const TrackCommand& command) const override
|
||||
{
|
||||
return command.kind == TrackCommandKind::StopTrack ||
|
||||
command.kind == TrackCommandKind::Flush;
|
||||
}
|
||||
};
|
||||
|
||||
template <std::size_t N>
|
||||
class TrackPointBuffer
|
||||
{
|
||||
public:
|
||||
bool append(const TrackPoint& point)
|
||||
{
|
||||
if (count_ >= N)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
points_[count_++] = point;
|
||||
return true;
|
||||
}
|
||||
|
||||
TrackPointBatch takeBatch(const TrackFlushPolicy& policy, uint32_t now_ms)
|
||||
{
|
||||
if (!policy.shouldFlush(count_, now_ms))
|
||||
{
|
||||
return {};
|
||||
}
|
||||
batch_count_ = count_;
|
||||
for (std::size_t i = 0; i < count_; ++i)
|
||||
{
|
||||
batch_[i] = points_[i];
|
||||
}
|
||||
count_ = 0;
|
||||
return TrackPointBatch{batch_, batch_count_};
|
||||
}
|
||||
|
||||
TrackPointBatch takeAll()
|
||||
{
|
||||
batch_count_ = count_;
|
||||
for (std::size_t i = 0; i < count_; ++i)
|
||||
{
|
||||
batch_[i] = points_[i];
|
||||
}
|
||||
count_ = 0;
|
||||
return TrackPointBatch{batch_, batch_count_};
|
||||
}
|
||||
|
||||
std::size_t size() const
|
||||
{
|
||||
return count_;
|
||||
}
|
||||
|
||||
private:
|
||||
TrackPoint points_[N]{};
|
||||
TrackPoint batch_[N]{};
|
||||
std::size_t count_ = 0;
|
||||
std::size_t batch_count_ = 0;
|
||||
};
|
||||
|
||||
class TrackStateMachine
|
||||
{
|
||||
public:
|
||||
TrackRecorderStatus state() const
|
||||
{
|
||||
return state_;
|
||||
}
|
||||
|
||||
void transition(const TrackEvent& event)
|
||||
{
|
||||
switch (event.kind)
|
||||
{
|
||||
case TrackEventKind::Started:
|
||||
state_ = TrackRecorderStatus::Recording;
|
||||
break;
|
||||
case TrackEventKind::Stopped:
|
||||
state_ = TrackRecorderStatus::Stopped;
|
||||
break;
|
||||
case TrackEventKind::FlushSucceeded:
|
||||
state_ = TrackRecorderStatus::Recording;
|
||||
break;
|
||||
case TrackEventKind::FlushFailed:
|
||||
case TrackEventKind::Failed:
|
||||
state_ = TrackRecorderStatus::Error;
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
last_error_ = event.error;
|
||||
}
|
||||
|
||||
void setState(TrackRecorderStatus state)
|
||||
{
|
||||
state_ = state;
|
||||
}
|
||||
|
||||
int32_t lastError() const
|
||||
{
|
||||
return last_error_;
|
||||
}
|
||||
|
||||
private:
|
||||
TrackRecorderStatus state_ = TrackRecorderStatus::Idle;
|
||||
int32_t last_error_ = 0;
|
||||
};
|
||||
|
||||
class ITrackFileAdapter
|
||||
{
|
||||
public:
|
||||
virtual ~ITrackFileAdapter() = default;
|
||||
|
||||
virtual bool open(uint32_t track_id) = 0;
|
||||
virtual bool append(uint32_t track_id, const TrackPoint* points, std::size_t count) = 0;
|
||||
virtual bool flush(uint32_t track_id) = 0;
|
||||
virtual bool close(uint32_t track_id) = 0;
|
||||
virtual std::size_t list(uint32_t* track_ids, std::size_t capacity) = 0;
|
||||
};
|
||||
|
||||
class ITrackEventSink
|
||||
{
|
||||
public:
|
||||
virtual ~ITrackEventSink() = default;
|
||||
|
||||
virtual bool publish(const TrackEvent& event) = 0;
|
||||
};
|
||||
|
||||
class TrackStorageWorker
|
||||
{
|
||||
public:
|
||||
TrackStorageWorker(ITrackFileAdapter& files,
|
||||
sys::runtime::IBusArbiter& bus,
|
||||
ITrackEventSink& events,
|
||||
TrackFlushPolicy& policy)
|
||||
: files_(files), bus_(bus), events_(events), policy_(policy)
|
||||
{
|
||||
}
|
||||
|
||||
bool submit(const TrackCommand& command)
|
||||
{
|
||||
if (pending_)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
pending_command_ = command;
|
||||
pending_ = true;
|
||||
return true;
|
||||
}
|
||||
|
||||
void tick(uint32_t now_ms)
|
||||
{
|
||||
if (!pending_)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
sys::runtime::BusAcquireRequest request{};
|
||||
request.command_id = pending_command_.command_id;
|
||||
request.deadline_ms = pending_command_.deadline_ms;
|
||||
request.policy = policy_.isCritical(pending_command_)
|
||||
? sys::runtime::BusAccessPolicy::DurableCommit
|
||||
: sys::runtime::BusAccessPolicy::BackgroundWorkerBounded;
|
||||
const sys::runtime::BusAcquireResult acquire = bus_.acquire(request);
|
||||
if (acquire.status != sys::runtime::BusAcquireStatus::Acquired)
|
||||
{
|
||||
publish(TrackEventKind::FlushFailed, now_ms, -10);
|
||||
pending_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
bool ok = false;
|
||||
switch (pending_command_.kind)
|
||||
{
|
||||
case TrackCommandKind::StartNewTrack:
|
||||
ok = files_.open(pending_command_.track_id);
|
||||
publish(ok ? TrackEventKind::Started : TrackEventKind::Failed,
|
||||
now_ms,
|
||||
ok ? 0 : -11);
|
||||
break;
|
||||
case TrackCommandKind::StopTrack:
|
||||
ok = files_.flush(pending_command_.track_id) &&
|
||||
files_.close(pending_command_.track_id);
|
||||
publish(ok ? TrackEventKind::Stopped : TrackEventKind::Failed,
|
||||
now_ms,
|
||||
ok ? 0 : -12);
|
||||
break;
|
||||
case TrackCommandKind::AppendPoint:
|
||||
case TrackCommandKind::Flush:
|
||||
ok = files_.append(pending_command_.track_id,
|
||||
pending_command_.point_batch.points,
|
||||
pending_command_.point_batch.count);
|
||||
if (ok)
|
||||
{
|
||||
ok = files_.flush(pending_command_.track_id);
|
||||
}
|
||||
publish(ok ? TrackEventKind::FlushSucceeded : TrackEventKind::FlushFailed,
|
||||
now_ms,
|
||||
ok ? 0 : -13);
|
||||
break;
|
||||
case TrackCommandKind::ListTracks:
|
||||
{
|
||||
uint32_t scratch[1]{};
|
||||
(void)files_.list(scratch, 0);
|
||||
publish(TrackEventKind::ListReady, now_ms, 0);
|
||||
ok = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
bus_.release(acquire.token);
|
||||
pending_ = false;
|
||||
}
|
||||
|
||||
bool busy() const
|
||||
{
|
||||
return pending_;
|
||||
}
|
||||
|
||||
private:
|
||||
void publish(TrackEventKind kind, uint32_t now_ms, int32_t error)
|
||||
{
|
||||
TrackEvent event{};
|
||||
event.kind = kind;
|
||||
event.command_id = pending_command_.command_id;
|
||||
event.track_id = pending_command_.track_id;
|
||||
event.timestamp_ms = now_ms;
|
||||
event.error = error;
|
||||
(void)events_.publish(event);
|
||||
}
|
||||
|
||||
ITrackFileAdapter& files_;
|
||||
sys::runtime::IBusArbiter& bus_;
|
||||
ITrackEventSink& events_;
|
||||
TrackFlushPolicy& policy_;
|
||||
TrackCommand pending_command_{};
|
||||
bool pending_ = false;
|
||||
};
|
||||
|
||||
template <std::size_t N>
|
||||
class TrackRuntime
|
||||
{
|
||||
public:
|
||||
TrackRuntime(TrackPointBuffer<N>& points,
|
||||
TrackStateMachine& states,
|
||||
TrackFlushPolicy& policy,
|
||||
TrackStorageWorker& worker,
|
||||
ITrackEventSink& events)
|
||||
: points_(points), states_(states), policy_(policy), worker_(worker), events_(events)
|
||||
{
|
||||
}
|
||||
|
||||
bool startNewTrack(uint32_t track_id, uint32_t now_ms)
|
||||
{
|
||||
states_.setState(TrackRecorderStatus::Starting);
|
||||
return submit(TrackCommandKind::StartNewTrack, track_id, {}, now_ms);
|
||||
}
|
||||
|
||||
bool stopTrack(uint32_t now_ms)
|
||||
{
|
||||
states_.setState(TrackRecorderStatus::Stopping);
|
||||
return submit(TrackCommandKind::StopTrack, active_track_id_, points_.takeAll(), now_ms);
|
||||
}
|
||||
|
||||
bool appendPoint(const TrackPoint& point, uint32_t now_ms)
|
||||
{
|
||||
if (!points_.append(point))
|
||||
{
|
||||
publish(TrackEventKind::Failed, now_ms, -20);
|
||||
return false;
|
||||
}
|
||||
publish(TrackEventKind::PointBuffered, now_ms, 0);
|
||||
TrackPointBatch batch = points_.takeBatch(policy_, now_ms);
|
||||
if (batch.count == 0)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
return submit(TrackCommandKind::AppendPoint, active_track_id_, batch, now_ms);
|
||||
}
|
||||
|
||||
bool listTracks(uint32_t now_ms)
|
||||
{
|
||||
return submit(TrackCommandKind::ListTracks, active_track_id_, {}, now_ms);
|
||||
}
|
||||
|
||||
void tick(uint32_t now_ms)
|
||||
{
|
||||
worker_.tick(now_ms);
|
||||
}
|
||||
|
||||
void handle(const TrackEvent& event)
|
||||
{
|
||||
states_.transition(event);
|
||||
last_event_ = event;
|
||||
if (event.kind == TrackEventKind::Started)
|
||||
{
|
||||
active_track_id_ = event.track_id;
|
||||
}
|
||||
}
|
||||
|
||||
TrackEvent lastEvent() const
|
||||
{
|
||||
return last_event_;
|
||||
}
|
||||
|
||||
private:
|
||||
bool submit(TrackCommandKind kind,
|
||||
uint32_t track_id,
|
||||
TrackPointBatch batch,
|
||||
uint32_t now_ms)
|
||||
{
|
||||
TrackCommand command{};
|
||||
command.command_id = next_command_id_++;
|
||||
command.kind = kind;
|
||||
command.track_id = track_id;
|
||||
command.point_batch = batch;
|
||||
command.created_at_ms = now_ms;
|
||||
if (kind == TrackCommandKind::StartNewTrack)
|
||||
{
|
||||
active_track_id_ = track_id;
|
||||
}
|
||||
return worker_.submit(command);
|
||||
}
|
||||
|
||||
void publish(TrackEventKind kind, uint32_t now_ms, int32_t error)
|
||||
{
|
||||
TrackEvent event{};
|
||||
event.kind = kind;
|
||||
event.track_id = active_track_id_;
|
||||
event.timestamp_ms = now_ms;
|
||||
event.error = error;
|
||||
last_event_ = event;
|
||||
(void)events_.publish(event);
|
||||
}
|
||||
|
||||
TrackPointBuffer<N>& points_;
|
||||
TrackStateMachine& states_;
|
||||
TrackFlushPolicy& policy_;
|
||||
TrackStorageWorker& worker_;
|
||||
ITrackEventSink& events_;
|
||||
TrackEvent last_event_{};
|
||||
uint32_t active_track_id_ = 0;
|
||||
uint32_t next_command_id_ = 1;
|
||||
};
|
||||
|
||||
} // namespace runtime
|
||||
} // namespace gps
|
||||
@@ -0,0 +1,261 @@
|
||||
#pragma once
|
||||
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <cstring>
|
||||
|
||||
namespace sys
|
||||
{
|
||||
namespace runtime
|
||||
{
|
||||
|
||||
enum class NoticeCategory : uint8_t
|
||||
{
|
||||
General,
|
||||
ChatDelivery,
|
||||
Storage,
|
||||
Tracker,
|
||||
Protocol,
|
||||
};
|
||||
|
||||
enum class NoticeSeverity : uint8_t
|
||||
{
|
||||
Info,
|
||||
Success,
|
||||
Warning,
|
||||
Error,
|
||||
};
|
||||
|
||||
enum class FeedbackEventKind : uint8_t
|
||||
{
|
||||
Posted,
|
||||
Presented,
|
||||
Dropped,
|
||||
Deduped,
|
||||
};
|
||||
|
||||
struct NoticeIntent
|
||||
{
|
||||
static constexpr std::size_t kMaxMessageBytes = 192;
|
||||
|
||||
uint32_t notice_id = 0;
|
||||
NoticeCategory category = NoticeCategory::General;
|
||||
NoticeSeverity severity = NoticeSeverity::Info;
|
||||
char message[kMaxMessageBytes]{};
|
||||
uint32_t duration_ms = 0;
|
||||
uint32_t dedupe_key = 0;
|
||||
uint32_t created_at_ms = 0;
|
||||
};
|
||||
|
||||
inline void setNoticeMessage(NoticeIntent& intent, const char* message)
|
||||
{
|
||||
std::strncpy(intent.message, message ? message : "", NoticeIntent::kMaxMessageBytes - 1);
|
||||
intent.message[NoticeIntent::kMaxMessageBytes - 1] = '\0';
|
||||
}
|
||||
|
||||
struct FeedbackEvent
|
||||
{
|
||||
FeedbackEventKind kind = FeedbackEventKind::Dropped;
|
||||
uint32_t notice_id = 0;
|
||||
uint32_t timestamp_ms = 0;
|
||||
int32_t error = 0;
|
||||
};
|
||||
|
||||
class FeedbackPolicy
|
||||
{
|
||||
public:
|
||||
virtual ~FeedbackPolicy() = default;
|
||||
|
||||
virtual bool shouldShow(const NoticeIntent& intent) const = 0;
|
||||
virtual bool dedupe(const NoticeIntent& previous, const NoticeIntent& next) const = 0;
|
||||
virtual uint32_t duration(const NoticeIntent& intent) const = 0;
|
||||
};
|
||||
|
||||
class DefaultFeedbackPolicy : public FeedbackPolicy
|
||||
{
|
||||
public:
|
||||
bool shouldShow(const NoticeIntent& intent) const override
|
||||
{
|
||||
return intent.message[0] != '\0';
|
||||
}
|
||||
|
||||
bool dedupe(const NoticeIntent& previous, const NoticeIntent& next) const override
|
||||
{
|
||||
if (next.dedupe_key != 0 && previous.dedupe_key == next.dedupe_key)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
return previous.category == next.category && previous.severity == next.severity &&
|
||||
std::strcmp(previous.message, next.message) == 0;
|
||||
}
|
||||
|
||||
uint32_t duration(const NoticeIntent& intent) const override
|
||||
{
|
||||
if (intent.duration_ms != 0)
|
||||
{
|
||||
return intent.duration_ms;
|
||||
}
|
||||
return intent.severity == NoticeSeverity::Error ? 2200 : 1400;
|
||||
}
|
||||
};
|
||||
|
||||
class IFeedbackPresenter
|
||||
{
|
||||
public:
|
||||
virtual ~IFeedbackPresenter() = default;
|
||||
|
||||
virtual bool present(const NoticeIntent& intent) = 0;
|
||||
};
|
||||
|
||||
class IFeedbackEventSink
|
||||
{
|
||||
public:
|
||||
virtual ~IFeedbackEventSink() = default;
|
||||
|
||||
virtual bool publish(const FeedbackEvent& event) = 0;
|
||||
};
|
||||
|
||||
template <std::size_t N>
|
||||
class FeedbackQueue
|
||||
{
|
||||
public:
|
||||
bool enqueue(const NoticeIntent& intent)
|
||||
{
|
||||
if (count_ >= N)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
notices_[(head_ + count_) % N] = intent;
|
||||
++count_;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool popReady(uint32_t now_ms, NoticeIntent& out)
|
||||
{
|
||||
(void)now_ms;
|
||||
if (count_ == 0)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
out = notices_[head_];
|
||||
head_ = (head_ + 1) % N;
|
||||
--count_;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool dedupe(const NoticeIntent& intent, const FeedbackPolicy& policy)
|
||||
{
|
||||
for (std::size_t i = 0; i < count_; ++i)
|
||||
{
|
||||
NoticeIntent& queued = notices_[(head_ + i) % N];
|
||||
if (policy.dedupe(queued, intent))
|
||||
{
|
||||
queued = intent;
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
std::size_t size() const
|
||||
{
|
||||
return count_;
|
||||
}
|
||||
|
||||
private:
|
||||
NoticeIntent notices_[N]{};
|
||||
std::size_t head_ = 0;
|
||||
std::size_t count_ = 0;
|
||||
};
|
||||
|
||||
template <std::size_t N>
|
||||
class FeedbackRuntime
|
||||
{
|
||||
public:
|
||||
FeedbackRuntime(FeedbackQueue<N>& queue,
|
||||
FeedbackPolicy& policy,
|
||||
IFeedbackPresenter& presenter,
|
||||
IFeedbackEventSink& events)
|
||||
: queue_(queue), policy_(policy), presenter_(presenter), events_(events)
|
||||
{
|
||||
}
|
||||
|
||||
bool post(NoticeIntent intent)
|
||||
{
|
||||
if (!policy_.shouldShow(intent))
|
||||
{
|
||||
publish(FeedbackEventKind::Dropped, intent.notice_id, intent.created_at_ms, -1);
|
||||
return false;
|
||||
}
|
||||
|
||||
if (intent.notice_id == 0)
|
||||
{
|
||||
intent.notice_id = next_notice_id_++;
|
||||
}
|
||||
intent.duration_ms = policy_.duration(intent);
|
||||
|
||||
if (queue_.dedupe(intent, policy_))
|
||||
{
|
||||
publish(FeedbackEventKind::Deduped, intent.notice_id, intent.created_at_ms, 0);
|
||||
return true;
|
||||
}
|
||||
|
||||
const bool queued = queue_.enqueue(intent);
|
||||
publish(queued ? FeedbackEventKind::Posted : FeedbackEventKind::Dropped,
|
||||
intent.notice_id,
|
||||
intent.created_at_ms,
|
||||
queued ? 0 : -2);
|
||||
return queued;
|
||||
}
|
||||
|
||||
void handle(const FeedbackEvent& event)
|
||||
{
|
||||
last_event_ = event;
|
||||
}
|
||||
|
||||
std::size_t drainToUi(uint32_t now_ms)
|
||||
{
|
||||
std::size_t presented = 0;
|
||||
NoticeIntent intent{};
|
||||
while (queue_.popReady(now_ms, intent))
|
||||
{
|
||||
if (presenter_.present(intent))
|
||||
{
|
||||
++presented;
|
||||
publish(FeedbackEventKind::Presented, intent.notice_id, now_ms, 0);
|
||||
}
|
||||
else
|
||||
{
|
||||
publish(FeedbackEventKind::Dropped, intent.notice_id, now_ms, -3);
|
||||
}
|
||||
}
|
||||
return presented;
|
||||
}
|
||||
|
||||
FeedbackEvent lastEvent() const
|
||||
{
|
||||
return last_event_;
|
||||
}
|
||||
|
||||
private:
|
||||
void publish(FeedbackEventKind kind, uint32_t notice_id, uint32_t timestamp_ms, int32_t error)
|
||||
{
|
||||
FeedbackEvent event{};
|
||||
event.kind = kind;
|
||||
event.notice_id = notice_id;
|
||||
event.timestamp_ms = timestamp_ms;
|
||||
event.error = error;
|
||||
last_event_ = event;
|
||||
(void)events_.publish(event);
|
||||
}
|
||||
|
||||
FeedbackQueue<N>& queue_;
|
||||
FeedbackPolicy& policy_;
|
||||
IFeedbackPresenter& presenter_;
|
||||
IFeedbackEventSink& events_;
|
||||
FeedbackEvent last_event_{};
|
||||
uint32_t next_notice_id_ = 1;
|
||||
};
|
||||
|
||||
} // namespace runtime
|
||||
} // namespace sys
|
||||
@@ -0,0 +1,406 @@
|
||||
#pragma once
|
||||
|
||||
#include "sys/runtime_async.h"
|
||||
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <cstring>
|
||||
|
||||
namespace sys
|
||||
{
|
||||
namespace runtime
|
||||
{
|
||||
|
||||
enum class PersistencePolicyMode : uint8_t
|
||||
{
|
||||
DebouncedSave,
|
||||
BatchSave,
|
||||
ImmediateCriticalSave,
|
||||
DropDuplicateSave,
|
||||
};
|
||||
|
||||
enum class PersistenceEventKind : uint8_t
|
||||
{
|
||||
SaveQueued,
|
||||
SaveStarted,
|
||||
SaveSucceeded,
|
||||
SaveFailed,
|
||||
SaveCoalesced,
|
||||
SaveDropped,
|
||||
};
|
||||
|
||||
struct StoreSnapshot
|
||||
{
|
||||
const uint8_t* data = nullptr;
|
||||
std::size_t len = 0;
|
||||
bool valid = false;
|
||||
};
|
||||
|
||||
struct StoreStorageResult
|
||||
{
|
||||
bool ok = false;
|
||||
std::size_t bytes = 0;
|
||||
int32_t error = 0;
|
||||
};
|
||||
|
||||
struct PersistenceCommand
|
||||
{
|
||||
uint32_t command_id = 0;
|
||||
const char* store_key = nullptr;
|
||||
PersistencePolicyMode policy = PersistencePolicyMode::DebouncedSave;
|
||||
uint32_t deadline_ms = 0;
|
||||
uint32_t created_at_ms = 0;
|
||||
};
|
||||
|
||||
struct PersistenceEvent
|
||||
{
|
||||
PersistenceEventKind kind = PersistenceEventKind::SaveFailed;
|
||||
const char* store_key = nullptr;
|
||||
uint32_t command_id = 0;
|
||||
uint32_t timestamp_ms = 0;
|
||||
int32_t error = 0;
|
||||
};
|
||||
|
||||
class PersistencePolicy
|
||||
{
|
||||
public:
|
||||
virtual ~PersistencePolicy() = default;
|
||||
|
||||
virtual PersistencePolicyMode modeFor(const char* store_key) const = 0;
|
||||
virtual uint32_t delayFor(PersistencePolicyMode policy) const = 0;
|
||||
virtual BusAccessPolicy busPolicyFor(PersistencePolicyMode policy) const = 0;
|
||||
};
|
||||
|
||||
class DefaultPersistencePolicy : public PersistencePolicy
|
||||
{
|
||||
public:
|
||||
PersistencePolicyMode modeFor(const char* store_key) const override
|
||||
{
|
||||
(void)store_key;
|
||||
return PersistencePolicyMode::DebouncedSave;
|
||||
}
|
||||
|
||||
uint32_t delayFor(PersistencePolicyMode policy) const override
|
||||
{
|
||||
switch (policy)
|
||||
{
|
||||
case PersistencePolicyMode::ImmediateCriticalSave:
|
||||
return 0;
|
||||
case PersistencePolicyMode::BatchSave:
|
||||
return 500;
|
||||
case PersistencePolicyMode::DropDuplicateSave:
|
||||
case PersistencePolicyMode::DebouncedSave:
|
||||
default:
|
||||
return 150;
|
||||
}
|
||||
}
|
||||
|
||||
BusAccessPolicy busPolicyFor(PersistencePolicyMode policy) const override
|
||||
{
|
||||
return policy == PersistencePolicyMode::ImmediateCriticalSave
|
||||
? BusAccessPolicy::DurableCommit
|
||||
: BusAccessPolicy::BackgroundWorkerBounded;
|
||||
}
|
||||
};
|
||||
|
||||
class IStoreSnapshotProvider
|
||||
{
|
||||
public:
|
||||
virtual ~IStoreSnapshotProvider() = default;
|
||||
|
||||
virtual StoreSnapshot snapshot(const char* store_key) = 0;
|
||||
};
|
||||
|
||||
class IStoreStorageAdapter
|
||||
{
|
||||
public:
|
||||
virtual ~IStoreStorageAdapter() = default;
|
||||
|
||||
virtual StoreStorageResult write(const char* store_key,
|
||||
const uint8_t* bytes,
|
||||
std::size_t len) = 0;
|
||||
virtual StoreStorageResult read(const char* store_key,
|
||||
uint8_t* bytes,
|
||||
std::size_t capacity,
|
||||
std::size_t& out_len) = 0;
|
||||
};
|
||||
|
||||
class IPersistenceEventSink
|
||||
{
|
||||
public:
|
||||
virtual ~IPersistenceEventSink() = default;
|
||||
|
||||
virtual bool publish(const PersistenceEvent& event) = 0;
|
||||
};
|
||||
|
||||
template <std::size_t N>
|
||||
class DirtyStoreRegistry
|
||||
{
|
||||
public:
|
||||
bool markDirty(const char* store_key,
|
||||
PersistencePolicyMode policy,
|
||||
uint32_t due_ms)
|
||||
{
|
||||
if (store_key == nullptr || store_key[0] == '\0')
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
for (std::size_t i = 0; i < count_; ++i)
|
||||
{
|
||||
if (sameKey(records_[i].store_key, store_key))
|
||||
{
|
||||
records_[i].policy = policy;
|
||||
records_[i].due_ms = due_ms;
|
||||
records_[i].dirty = true;
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
if (count_ >= N)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
records_[count_++] = DirtyStoreRecord{store_key, due_ms, policy, true};
|
||||
return true;
|
||||
}
|
||||
|
||||
bool takeDue(uint32_t now_ms, PersistenceCommand& out)
|
||||
{
|
||||
for (std::size_t i = 0; i < count_; ++i)
|
||||
{
|
||||
DirtyStoreRecord& record = records_[i];
|
||||
if (!record.dirty || static_cast<int32_t>(record.due_ms - now_ms) > 0)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
out.store_key = record.store_key;
|
||||
out.policy = record.policy;
|
||||
out.created_at_ms = now_ms;
|
||||
record.dirty = false;
|
||||
compact();
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
bool hasPending(const char* store_key) const
|
||||
{
|
||||
for (std::size_t i = 0; i < count_; ++i)
|
||||
{
|
||||
if (records_[i].dirty && sameKey(records_[i].store_key, store_key))
|
||||
{
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
std::size_t size() const
|
||||
{
|
||||
return count_;
|
||||
}
|
||||
|
||||
private:
|
||||
struct DirtyStoreRecord
|
||||
{
|
||||
const char* store_key = nullptr;
|
||||
uint32_t due_ms = 0;
|
||||
PersistencePolicyMode policy = PersistencePolicyMode::DebouncedSave;
|
||||
bool dirty = false;
|
||||
};
|
||||
|
||||
static bool sameKey(const char* lhs, const char* rhs)
|
||||
{
|
||||
return lhs != nullptr && rhs != nullptr && std::strcmp(lhs, rhs) == 0;
|
||||
}
|
||||
|
||||
void compact()
|
||||
{
|
||||
DirtyStoreRecord kept[N]{};
|
||||
std::size_t kept_count = 0;
|
||||
for (std::size_t i = 0; i < count_; ++i)
|
||||
{
|
||||
if (records_[i].dirty)
|
||||
{
|
||||
kept[kept_count++] = records_[i];
|
||||
}
|
||||
}
|
||||
for (std::size_t i = 0; i < kept_count; ++i)
|
||||
{
|
||||
records_[i] = kept[i];
|
||||
}
|
||||
count_ = kept_count;
|
||||
}
|
||||
|
||||
DirtyStoreRecord records_[N]{};
|
||||
std::size_t count_ = 0;
|
||||
};
|
||||
|
||||
class PersistenceWorker
|
||||
{
|
||||
public:
|
||||
PersistenceWorker(IStoreSnapshotProvider& snapshots,
|
||||
IStoreStorageAdapter& storage,
|
||||
IBusArbiter& bus,
|
||||
IPersistenceEventSink& events,
|
||||
PersistencePolicy& policy)
|
||||
: snapshots_(snapshots), storage_(storage), bus_(bus), events_(events), policy_(policy)
|
||||
{
|
||||
}
|
||||
|
||||
bool submit(const PersistenceCommand& command)
|
||||
{
|
||||
if (pending_)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
pending_command_ = command;
|
||||
pending_ = true;
|
||||
return true;
|
||||
}
|
||||
|
||||
void tick(uint32_t now_ms)
|
||||
{
|
||||
if (!pending_)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
PersistenceEvent started{};
|
||||
started.kind = PersistenceEventKind::SaveStarted;
|
||||
started.store_key = pending_command_.store_key;
|
||||
started.command_id = pending_command_.command_id;
|
||||
started.timestamp_ms = now_ms;
|
||||
(void)events_.publish(started);
|
||||
|
||||
BusAcquireRequest request{};
|
||||
request.policy = policy_.busPolicyFor(pending_command_.policy);
|
||||
request.command_id = pending_command_.command_id;
|
||||
request.deadline_ms = pending_command_.deadline_ms;
|
||||
const BusAcquireResult acquire = bus_.acquire(request);
|
||||
if (acquire.status != BusAcquireStatus::Acquired)
|
||||
{
|
||||
publishResult(PersistenceEventKind::SaveFailed, now_ms, -11);
|
||||
pending_ = false;
|
||||
return;
|
||||
}
|
||||
|
||||
StoreSnapshot snapshot = snapshots_.snapshot(pending_command_.store_key);
|
||||
StoreStorageResult result{};
|
||||
if (snapshot.valid)
|
||||
{
|
||||
result = storage_.write(pending_command_.store_key, snapshot.data, snapshot.len);
|
||||
}
|
||||
else
|
||||
{
|
||||
result.error = -12;
|
||||
}
|
||||
bus_.release(acquire.token);
|
||||
|
||||
publishResult(result.ok ? PersistenceEventKind::SaveSucceeded
|
||||
: PersistenceEventKind::SaveFailed,
|
||||
now_ms,
|
||||
result.ok ? 0 : result.error);
|
||||
pending_ = false;
|
||||
}
|
||||
|
||||
bool busy() const
|
||||
{
|
||||
return pending_;
|
||||
}
|
||||
|
||||
private:
|
||||
void publishResult(PersistenceEventKind kind, uint32_t now_ms, int32_t error)
|
||||
{
|
||||
PersistenceEvent event{};
|
||||
event.kind = kind;
|
||||
event.store_key = pending_command_.store_key;
|
||||
event.command_id = pending_command_.command_id;
|
||||
event.timestamp_ms = now_ms;
|
||||
event.error = error;
|
||||
(void)events_.publish(event);
|
||||
}
|
||||
|
||||
IStoreSnapshotProvider& snapshots_;
|
||||
IStoreStorageAdapter& storage_;
|
||||
IBusArbiter& bus_;
|
||||
IPersistenceEventSink& events_;
|
||||
PersistencePolicy& policy_;
|
||||
PersistenceCommand pending_command_{};
|
||||
bool pending_ = false;
|
||||
};
|
||||
|
||||
template <std::size_t N>
|
||||
class PersistenceRuntime
|
||||
{
|
||||
public:
|
||||
PersistenceRuntime(DirtyStoreRegistry<N>& registry,
|
||||
PersistenceWorker& worker,
|
||||
IPersistenceEventSink& events,
|
||||
PersistencePolicy& policy)
|
||||
: registry_(registry), worker_(worker), events_(events), policy_(policy)
|
||||
{
|
||||
}
|
||||
|
||||
bool markDirty(const char* store_key, uint32_t now_ms)
|
||||
{
|
||||
const PersistencePolicyMode mode = policy_.modeFor(store_key);
|
||||
const bool marked = registry_.markDirty(store_key, mode, now_ms + policy_.delayFor(mode));
|
||||
PersistenceEvent event{};
|
||||
event.kind = marked ? PersistenceEventKind::SaveQueued
|
||||
: PersistenceEventKind::SaveDropped;
|
||||
event.store_key = store_key;
|
||||
event.timestamp_ms = now_ms;
|
||||
(void)events_.publish(event);
|
||||
return marked;
|
||||
}
|
||||
|
||||
bool requestSave(const char* store_key,
|
||||
PersistencePolicyMode policy,
|
||||
uint32_t now_ms)
|
||||
{
|
||||
return registry_.markDirty(store_key, policy, now_ms + policy_.delayFor(policy));
|
||||
}
|
||||
|
||||
void tick(uint32_t now_ms)
|
||||
{
|
||||
if (!worker_.busy())
|
||||
{
|
||||
PersistenceCommand command{};
|
||||
if (registry_.takeDue(now_ms, command))
|
||||
{
|
||||
command.command_id = next_command_id_++;
|
||||
if (!worker_.submit(command))
|
||||
{
|
||||
(void)registry_.markDirty(command.store_key,
|
||||
command.policy,
|
||||
now_ms + policy_.delayFor(command.policy));
|
||||
}
|
||||
}
|
||||
}
|
||||
worker_.tick(now_ms);
|
||||
}
|
||||
|
||||
void handle(const PersistenceEvent& event)
|
||||
{
|
||||
last_event_ = event;
|
||||
}
|
||||
|
||||
PersistenceEvent lastEvent() const
|
||||
{
|
||||
return last_event_;
|
||||
}
|
||||
|
||||
private:
|
||||
DirtyStoreRegistry<N>& registry_;
|
||||
PersistenceWorker& worker_;
|
||||
IPersistenceEventSink& events_;
|
||||
PersistencePolicy& policy_;
|
||||
PersistenceEvent last_event_{};
|
||||
uint32_t next_command_id_ = 1;
|
||||
};
|
||||
|
||||
} // namespace runtime
|
||||
} // namespace sys
|
||||
@@ -270,6 +270,164 @@ class IBusArbiter
|
||||
virtual StorageHealthState health() const = 0;
|
||||
};
|
||||
|
||||
class IBusAdapter
|
||||
{
|
||||
public:
|
||||
virtual ~IBusAdapter() = default;
|
||||
|
||||
virtual bool tryAcquire(uint32_t timeout_ms) = 0;
|
||||
virtual void release() = 0;
|
||||
virtual uint32_t nowMs() const = 0;
|
||||
virtual uint32_t owner() const = 0;
|
||||
};
|
||||
|
||||
class BusPolicyStrategy
|
||||
{
|
||||
public:
|
||||
virtual ~BusPolicyStrategy() = default;
|
||||
|
||||
virtual BusAccessPolicy select(const RuntimeCommand& command) const = 0;
|
||||
virtual uint32_t timeoutFor(BusAccessPolicy policy) const = 0;
|
||||
};
|
||||
|
||||
class DefaultBusPolicyStrategy : public BusPolicyStrategy
|
||||
{
|
||||
public:
|
||||
BusAccessPolicy select(const RuntimeCommand& command) const override
|
||||
{
|
||||
if (command.priority == RuntimePriority::Realtime ||
|
||||
command.priority == RuntimePriority::Interactive)
|
||||
{
|
||||
return BusAccessPolicy::InteractiveWorkerBounded;
|
||||
}
|
||||
if (command.kind == RuntimeCommandKind::TrackStop ||
|
||||
command.kind == RuntimeCommandKind::TrackFlush)
|
||||
{
|
||||
return BusAccessPolicy::DurableCommit;
|
||||
}
|
||||
return BusAccessPolicy::BackgroundWorkerBounded;
|
||||
}
|
||||
|
||||
uint32_t timeoutFor(BusAccessPolicy policy) const override
|
||||
{
|
||||
switch (policy)
|
||||
{
|
||||
case BusAccessPolicy::UiNeverBlock:
|
||||
return 0;
|
||||
case BusAccessPolicy::InteractiveWorkerBounded:
|
||||
return 2;
|
||||
case BusAccessPolicy::BackgroundWorkerBounded:
|
||||
return 25;
|
||||
case BusAccessPolicy::DurableCommit:
|
||||
return 150;
|
||||
case BusAccessPolicy::RecoveryExclusive:
|
||||
return 500;
|
||||
default:
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
class StorageBusArbiter : public IBusArbiter
|
||||
{
|
||||
public:
|
||||
StorageBusArbiter(IBusAdapter& adapter, BusPolicyStrategy& policy)
|
||||
: adapter_(adapter), policy_(policy)
|
||||
{
|
||||
}
|
||||
|
||||
BusAcquireResult acquire(const BusAcquireRequest& request) override
|
||||
{
|
||||
const uint32_t start_ms = adapter_.nowMs();
|
||||
uint32_t timeout_ms = policy_.timeoutFor(request.policy);
|
||||
if (request.deadline_ms != 0)
|
||||
{
|
||||
const uint32_t remaining =
|
||||
static_cast<int32_t>(request.deadline_ms - start_ms) > 0
|
||||
? request.deadline_ms - start_ms
|
||||
: 0;
|
||||
if (remaining < timeout_ms)
|
||||
{
|
||||
timeout_ms = remaining;
|
||||
}
|
||||
}
|
||||
|
||||
const bool acquired = adapter_.tryAcquire(timeout_ms);
|
||||
const uint32_t end_ms = adapter_.nowMs();
|
||||
|
||||
BusAcquireResult result{};
|
||||
result.status = acquired ? BusAcquireStatus::Acquired
|
||||
: (timeout_ms == 0 ? BusAcquireStatus::Busy
|
||||
: BusAcquireStatus::TimedOut);
|
||||
result.token.resource = request.resource;
|
||||
result.token.owner = request.command_id;
|
||||
result.token.acquired_ms = acquired ? end_ms : 0;
|
||||
result.token.valid = acquired;
|
||||
result.diagnostics.resource = request.resource;
|
||||
result.diagnostics.owner = adapter_.owner();
|
||||
result.diagnostics.command_id = request.command_id;
|
||||
result.diagnostics.wait_ms = end_ms - start_ms;
|
||||
result.diagnostics.policy = request.policy;
|
||||
|
||||
updateHealth(result.status, end_ms);
|
||||
return result;
|
||||
}
|
||||
|
||||
void release(const BusAccessToken& token) override
|
||||
{
|
||||
if (!token.valid)
|
||||
{
|
||||
return;
|
||||
}
|
||||
adapter_.release();
|
||||
consecutive_timeouts_ = 0;
|
||||
if (health_.status == StorageHealthStatus::Slow ||
|
||||
health_.status == StorageHealthStatus::Recovering)
|
||||
{
|
||||
health_.status = StorageHealthStatus::Healthy;
|
||||
health_.last_error = 0;
|
||||
health_.last_transition_ms = adapter_.nowMs();
|
||||
}
|
||||
}
|
||||
|
||||
StorageHealthState health() const override
|
||||
{
|
||||
return health_;
|
||||
}
|
||||
|
||||
BusAccessPolicy selectPolicy(const RuntimeCommand& command) const
|
||||
{
|
||||
return policy_.select(command);
|
||||
}
|
||||
|
||||
private:
|
||||
void updateHealth(BusAcquireStatus status, uint32_t now_ms)
|
||||
{
|
||||
if (status == BusAcquireStatus::Acquired)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
health_.last_transition_ms = now_ms;
|
||||
if (status == BusAcquireStatus::Unavailable)
|
||||
{
|
||||
health_.status = StorageHealthStatus::Unavailable;
|
||||
health_.last_error = -3;
|
||||
return;
|
||||
}
|
||||
|
||||
++consecutive_timeouts_;
|
||||
health_.last_error = status == BusAcquireStatus::TimedOut ? -2 : -1;
|
||||
health_.status = consecutive_timeouts_ >= 3 ? StorageHealthStatus::Degraded
|
||||
: StorageHealthStatus::Slow;
|
||||
}
|
||||
|
||||
IBusAdapter& adapter_;
|
||||
BusPolicyStrategy& policy_;
|
||||
StorageHealthState health_{};
|
||||
uint8_t consecutive_timeouts_ = 0;
|
||||
};
|
||||
|
||||
class IPlatformStorageAdapter
|
||||
{
|
||||
public:
|
||||
|
||||
@@ -0,0 +1,461 @@
|
||||
#pragma once
|
||||
|
||||
#include "sys/feedback_runtime.h"
|
||||
#include "sys/persistence_runtime.h"
|
||||
#include "sys/runtime_async.h"
|
||||
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <cassert>
|
||||
|
||||
namespace sys
|
||||
{
|
||||
namespace runtime
|
||||
{
|
||||
|
||||
class FakeClock
|
||||
{
|
||||
public:
|
||||
uint32_t now() const
|
||||
{
|
||||
return now_ms_;
|
||||
}
|
||||
|
||||
void advance(uint32_t ms)
|
||||
{
|
||||
now_ms_ += ms;
|
||||
}
|
||||
|
||||
private:
|
||||
uint32_t now_ms_ = 0;
|
||||
};
|
||||
|
||||
class FakeCommandQueue : public ICommandQueue
|
||||
{
|
||||
public:
|
||||
bool enqueue(const RuntimeCommand& command) override
|
||||
{
|
||||
return queue_.enqueue(command);
|
||||
}
|
||||
|
||||
std::size_t cancel(uint32_t dedupe_key) override
|
||||
{
|
||||
return queue_.cancel(dedupe_key);
|
||||
}
|
||||
|
||||
bool popReady(uint32_t now_ms, RuntimeCommand& out) override
|
||||
{
|
||||
return queue_.popReady(now_ms, out);
|
||||
}
|
||||
|
||||
std::size_t size() const
|
||||
{
|
||||
return queue_.size();
|
||||
}
|
||||
|
||||
private:
|
||||
FixedCommandQueue<32> queue_{};
|
||||
};
|
||||
|
||||
class FakeEventBus : public IEventSink,
|
||||
public IPersistenceEventSink,
|
||||
public IFeedbackEventSink
|
||||
{
|
||||
public:
|
||||
bool publish(const RuntimeEvent& event) override
|
||||
{
|
||||
return runtime_events_.publish(event);
|
||||
}
|
||||
|
||||
bool publish(const PersistenceEvent& event) override
|
||||
{
|
||||
if (persistence_count_ >= kMaxEvents)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
persistence_events_[persistence_count_++] = event;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool publish(const FeedbackEvent& event) override
|
||||
{
|
||||
if (feedback_count_ >= kMaxEvents)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
feedback_events_[feedback_count_++] = event;
|
||||
return true;
|
||||
}
|
||||
|
||||
std::size_t drain()
|
||||
{
|
||||
std::size_t count = 0;
|
||||
RuntimeEvent event{};
|
||||
while (runtime_events_.pop(event))
|
||||
{
|
||||
++count;
|
||||
}
|
||||
count += persistence_count_;
|
||||
count += feedback_count_;
|
||||
persistence_count_ = 0;
|
||||
feedback_count_ = 0;
|
||||
return count;
|
||||
}
|
||||
|
||||
std::size_t persistenceCount() const
|
||||
{
|
||||
return persistence_count_;
|
||||
}
|
||||
|
||||
std::size_t feedbackCount() const
|
||||
{
|
||||
return feedback_count_;
|
||||
}
|
||||
|
||||
const PersistenceEvent& persistenceEvent(std::size_t index) const
|
||||
{
|
||||
return persistence_events_[index];
|
||||
}
|
||||
|
||||
const FeedbackEvent& feedbackEvent(std::size_t index) const
|
||||
{
|
||||
return feedback_events_[index];
|
||||
}
|
||||
|
||||
private:
|
||||
static constexpr std::size_t kMaxEvents = 32;
|
||||
|
||||
FixedEventSink<32> runtime_events_{};
|
||||
PersistenceEvent persistence_events_[kMaxEvents]{};
|
||||
FeedbackEvent feedback_events_[kMaxEvents]{};
|
||||
std::size_t persistence_count_ = 0;
|
||||
std::size_t feedback_count_ = 0;
|
||||
};
|
||||
|
||||
class FakeStorageBackend : public IPlatformStorageAdapter,
|
||||
public IStoreSnapshotProvider,
|
||||
public IStoreStorageAdapter
|
||||
{
|
||||
public:
|
||||
void scriptDelay(const char* operation, uint32_t ms)
|
||||
{
|
||||
(void)operation;
|
||||
delay_ms_ = ms;
|
||||
}
|
||||
|
||||
void scriptFailure(const char* operation, int32_t error)
|
||||
{
|
||||
(void)operation;
|
||||
fail_ = true;
|
||||
error_ = error;
|
||||
}
|
||||
|
||||
PlatformStorageResult read(const PlatformStorageReadRequest& request) override
|
||||
{
|
||||
last_command_id_ = request.command_id;
|
||||
return platformResult(request.capacity);
|
||||
}
|
||||
|
||||
PlatformStorageResult write(const PlatformStorageWriteRequest& request) override
|
||||
{
|
||||
last_command_id_ = request.command_id;
|
||||
return platformResult(request.len);
|
||||
}
|
||||
|
||||
PlatformStorageResult list(const PlatformStorageListRequest& request) override
|
||||
{
|
||||
last_command_id_ = request.command_id;
|
||||
return platformResult(0);
|
||||
}
|
||||
|
||||
PlatformStorageResult flush(const PlatformStorageFlushRequest& request) override
|
||||
{
|
||||
last_command_id_ = request.command_id;
|
||||
return platformResult(0);
|
||||
}
|
||||
|
||||
StoreSnapshot snapshot(const char* store_key) override
|
||||
{
|
||||
(void)store_key;
|
||||
StoreSnapshot snapshot{};
|
||||
snapshot.data = snapshot_;
|
||||
snapshot.len = snapshot_len_;
|
||||
snapshot.valid = !fail_;
|
||||
return snapshot;
|
||||
}
|
||||
|
||||
StoreStorageResult write(const char* store_key,
|
||||
const uint8_t* bytes,
|
||||
std::size_t len) override
|
||||
{
|
||||
(void)store_key;
|
||||
(void)bytes;
|
||||
StoreStorageResult result{};
|
||||
result.ok = !fail_;
|
||||
result.bytes = result.ok ? len : 0;
|
||||
result.error = result.ok ? 0 : error_;
|
||||
++write_count_;
|
||||
return result;
|
||||
}
|
||||
|
||||
StoreStorageResult read(const char* store_key,
|
||||
uint8_t* bytes,
|
||||
std::size_t capacity,
|
||||
std::size_t& out_len) override
|
||||
{
|
||||
(void)store_key;
|
||||
(void)bytes;
|
||||
out_len = fail_ ? 0 : capacity;
|
||||
StoreStorageResult result{};
|
||||
result.ok = !fail_;
|
||||
result.bytes = out_len;
|
||||
result.error = result.ok ? 0 : error_;
|
||||
return result;
|
||||
}
|
||||
|
||||
uint32_t delayMs() const
|
||||
{
|
||||
return delay_ms_;
|
||||
}
|
||||
|
||||
uint32_t lastCommandId() const
|
||||
{
|
||||
return last_command_id_;
|
||||
}
|
||||
|
||||
std::size_t writeCount() const
|
||||
{
|
||||
return write_count_;
|
||||
}
|
||||
|
||||
private:
|
||||
PlatformStorageResult platformResult(std::size_t bytes)
|
||||
{
|
||||
PlatformStorageResult result{};
|
||||
result.ok = !fail_;
|
||||
result.bytes = result.ok ? bytes : 0;
|
||||
result.error = result.ok ? 0 : error_;
|
||||
return result;
|
||||
}
|
||||
|
||||
uint8_t snapshot_[4] = {1, 2, 3, 4};
|
||||
std::size_t snapshot_len_ = sizeof(snapshot_);
|
||||
uint32_t delay_ms_ = 0;
|
||||
uint32_t last_command_id_ = 0;
|
||||
std::size_t write_count_ = 0;
|
||||
bool fail_ = false;
|
||||
int32_t error_ = -1;
|
||||
};
|
||||
|
||||
class FakeBusArbiter : public IBusArbiter
|
||||
{
|
||||
public:
|
||||
void scriptAcquire(BusAcquireStatus status)
|
||||
{
|
||||
scripted_status_ = status;
|
||||
}
|
||||
|
||||
BusAcquireResult acquire(const BusAcquireRequest& request) override
|
||||
{
|
||||
++acquire_count_;
|
||||
last_request_ = request;
|
||||
BusAcquireResult result{};
|
||||
result.status = scripted_status_;
|
||||
result.token.valid = scripted_status_ == BusAcquireStatus::Acquired;
|
||||
result.token.resource = request.resource;
|
||||
result.token.owner = request.command_id;
|
||||
result.token.acquired_ms = now_ms_;
|
||||
result.diagnostics.resource = request.resource;
|
||||
result.diagnostics.command_id = request.command_id;
|
||||
result.diagnostics.policy = request.policy;
|
||||
if (scripted_status_ != BusAcquireStatus::Acquired)
|
||||
{
|
||||
health_.status = StorageHealthStatus::Slow;
|
||||
health_.last_error = -1;
|
||||
health_.last_transition_ms = now_ms_;
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
void release(const BusAccessToken& token) override
|
||||
{
|
||||
if (token.valid)
|
||||
{
|
||||
++release_count_;
|
||||
}
|
||||
}
|
||||
|
||||
StorageHealthState health() const override
|
||||
{
|
||||
return health_;
|
||||
}
|
||||
|
||||
BusAcquireRequest diagnostics() const
|
||||
{
|
||||
return last_request_;
|
||||
}
|
||||
|
||||
void setNow(uint32_t now_ms)
|
||||
{
|
||||
now_ms_ = now_ms;
|
||||
}
|
||||
|
||||
std::size_t acquireCount() const
|
||||
{
|
||||
return acquire_count_;
|
||||
}
|
||||
|
||||
std::size_t releaseCount() const
|
||||
{
|
||||
return release_count_;
|
||||
}
|
||||
|
||||
private:
|
||||
BusAcquireStatus scripted_status_ = BusAcquireStatus::Acquired;
|
||||
BusAcquireRequest last_request_{};
|
||||
StorageHealthState health_{};
|
||||
uint32_t now_ms_ = 0;
|
||||
std::size_t acquire_count_ = 0;
|
||||
std::size_t release_count_ = 0;
|
||||
};
|
||||
|
||||
class FakeUiOwner : public IUiEffectSink
|
||||
{
|
||||
public:
|
||||
void assertNoBlockingCalls() const
|
||||
{
|
||||
assert(blocking_calls_ == 0);
|
||||
}
|
||||
|
||||
bool apply(const RuntimeUiEffect& effect) override
|
||||
{
|
||||
last_effect_ = effect;
|
||||
++effect_count_;
|
||||
return true;
|
||||
}
|
||||
|
||||
void tick()
|
||||
{
|
||||
++tick_count_;
|
||||
}
|
||||
|
||||
void recordBlockingCall()
|
||||
{
|
||||
++blocking_calls_;
|
||||
}
|
||||
|
||||
std::size_t effectCount() const
|
||||
{
|
||||
return effect_count_;
|
||||
}
|
||||
|
||||
private:
|
||||
RuntimeUiEffect last_effect_{};
|
||||
std::size_t effect_count_ = 0;
|
||||
std::size_t tick_count_ = 0;
|
||||
std::size_t blocking_calls_ = 0;
|
||||
};
|
||||
|
||||
class FakeFeedbackPresenter : public IFeedbackPresenter
|
||||
{
|
||||
public:
|
||||
bool present(const NoticeIntent& intent) override
|
||||
{
|
||||
if (count_ >= kMaxNotices)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
captured_[count_++] = intent;
|
||||
return true;
|
||||
}
|
||||
|
||||
std::size_t captured() const
|
||||
{
|
||||
return count_;
|
||||
}
|
||||
|
||||
const NoticeIntent& capturedNotice(std::size_t index) const
|
||||
{
|
||||
return captured_[index];
|
||||
}
|
||||
|
||||
private:
|
||||
static constexpr std::size_t kMaxNotices = 16;
|
||||
NoticeIntent captured_[kMaxNotices]{};
|
||||
std::size_t count_ = 0;
|
||||
};
|
||||
|
||||
class RuntimeHarness
|
||||
{
|
||||
public:
|
||||
FakeClock& clock()
|
||||
{
|
||||
return clock_;
|
||||
}
|
||||
|
||||
FakeCommandQueue& commands()
|
||||
{
|
||||
return commands_;
|
||||
}
|
||||
|
||||
FakeEventBus& events()
|
||||
{
|
||||
return events_;
|
||||
}
|
||||
|
||||
FakeStorageBackend& storage()
|
||||
{
|
||||
return storage_;
|
||||
}
|
||||
|
||||
FakeBusArbiter& bus()
|
||||
{
|
||||
return bus_;
|
||||
}
|
||||
|
||||
FakeUiOwner& ui()
|
||||
{
|
||||
return ui_;
|
||||
}
|
||||
|
||||
FakeFeedbackPresenter& feedbackPresenter()
|
||||
{
|
||||
return feedback_;
|
||||
}
|
||||
|
||||
void advance(uint32_t ms)
|
||||
{
|
||||
clock_.advance(ms);
|
||||
bus_.setNow(clock_.now());
|
||||
}
|
||||
|
||||
std::size_t drainUi()
|
||||
{
|
||||
ui_.tick();
|
||||
return events_.drain();
|
||||
}
|
||||
|
||||
void runUntilIdle(uint32_t step_ms = 10, uint32_t max_steps = 64)
|
||||
{
|
||||
for (uint32_t i = 0; i < max_steps; ++i)
|
||||
{
|
||||
advance(step_ms);
|
||||
if (drainUi() == 0)
|
||||
{
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private:
|
||||
FakeClock clock_{};
|
||||
FakeCommandQueue commands_{};
|
||||
FakeEventBus events_{};
|
||||
FakeStorageBackend storage_{};
|
||||
FakeBusArbiter bus_{};
|
||||
FakeUiOwner ui_{};
|
||||
FakeFeedbackPresenter feedback_{};
|
||||
};
|
||||
|
||||
} // namespace runtime
|
||||
} // namespace sys
|
||||
@@ -40,6 +40,39 @@ class FakeUiEffectSink final : public sys::runtime::IUiEffectSink
|
||||
sys::runtime::RuntimeUiEffect last_effect{};
|
||||
};
|
||||
|
||||
class FakeBusAdapter final : public sys::runtime::IBusAdapter
|
||||
{
|
||||
public:
|
||||
bool acquire_ok = true;
|
||||
uint32_t now_ms = 100;
|
||||
uint32_t last_timeout_ms = 0;
|
||||
int acquire_count = 0;
|
||||
int release_count = 0;
|
||||
|
||||
bool tryAcquire(uint32_t timeout_ms) override
|
||||
{
|
||||
last_timeout_ms = timeout_ms;
|
||||
++acquire_count;
|
||||
now_ms += timeout_ms;
|
||||
return acquire_ok;
|
||||
}
|
||||
|
||||
void release() override
|
||||
{
|
||||
++release_count;
|
||||
}
|
||||
|
||||
uint32_t nowMs() const override
|
||||
{
|
||||
return now_ms;
|
||||
}
|
||||
|
||||
uint32_t owner() const override
|
||||
{
|
||||
return 77;
|
||||
}
|
||||
};
|
||||
|
||||
void test_priority_pop_order()
|
||||
{
|
||||
sys::runtime::FixedCommandQueue<4> queue;
|
||||
@@ -223,6 +256,45 @@ void test_event_to_ui_effect_bridge()
|
||||
assert(out.event_id == event.event_id);
|
||||
}
|
||||
|
||||
void test_storage_bus_arbiter_uses_policy_timeout()
|
||||
{
|
||||
FakeBusAdapter adapter;
|
||||
sys::runtime::DefaultBusPolicyStrategy policy;
|
||||
sys::runtime::StorageBusArbiter arbiter(adapter, policy);
|
||||
|
||||
sys::runtime::BusAcquireRequest request{};
|
||||
request.resource = 3;
|
||||
request.command_id = 9;
|
||||
request.policy = sys::runtime::BusAccessPolicy::InteractiveWorkerBounded;
|
||||
|
||||
const sys::runtime::BusAcquireResult result = arbiter.acquire(request);
|
||||
assert(result.status == sys::runtime::BusAcquireStatus::Acquired);
|
||||
assert(result.token.valid);
|
||||
assert(result.token.owner == 9);
|
||||
assert(adapter.last_timeout_ms == 2);
|
||||
assert(result.diagnostics.owner == 77);
|
||||
|
||||
arbiter.release(result.token);
|
||||
assert(adapter.release_count == 1);
|
||||
}
|
||||
|
||||
void test_storage_bus_arbiter_reports_degraded_after_timeouts()
|
||||
{
|
||||
FakeBusAdapter adapter;
|
||||
adapter.acquire_ok = false;
|
||||
sys::runtime::DefaultBusPolicyStrategy policy;
|
||||
sys::runtime::StorageBusArbiter arbiter(adapter, policy);
|
||||
|
||||
sys::runtime::BusAcquireRequest request{};
|
||||
request.policy = sys::runtime::BusAccessPolicy::BackgroundWorkerBounded;
|
||||
|
||||
assert(arbiter.acquire(request).status == sys::runtime::BusAcquireStatus::TimedOut);
|
||||
assert(arbiter.health().status == sys::runtime::StorageHealthStatus::Slow);
|
||||
assert(arbiter.acquire(request).status == sys::runtime::BusAcquireStatus::TimedOut);
|
||||
assert(arbiter.acquire(request).status == sys::runtime::BusAcquireStatus::TimedOut);
|
||||
assert(arbiter.health().status == sys::runtime::StorageHealthStatus::Degraded);
|
||||
}
|
||||
|
||||
} // namespace
|
||||
|
||||
int main()
|
||||
@@ -234,5 +306,7 @@ int main()
|
||||
test_runtime_facade_submit_tick();
|
||||
test_runtime_facade_dedupe_cancel_policy();
|
||||
test_event_to_ui_effect_bridge();
|
||||
test_storage_bus_arbiter_uses_policy_timeout();
|
||||
test_storage_bus_arbiter_reports_degraded_after_timeouts();
|
||||
return 0;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,266 @@
|
||||
#include "gps/track_runtime.h"
|
||||
#include "sys/feedback_runtime.h"
|
||||
#include "sys/persistence_runtime.h"
|
||||
#include "sys/runtime_harness.h"
|
||||
#include "ui_map_runtime/map_tiles/map_tile_async_runtime.h"
|
||||
|
||||
#include <cassert>
|
||||
#include <cstddef>
|
||||
|
||||
namespace
|
||||
{
|
||||
|
||||
class MapCommandSink final : public ui::map_tiles::IMapTileCommandSink
|
||||
{
|
||||
public:
|
||||
bool enqueue(const ui::map_tiles::LoadTileCommand& command) override
|
||||
{
|
||||
if (count >= 4)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
commands[count++] = command;
|
||||
return true;
|
||||
}
|
||||
|
||||
std::size_t cancelGeneration(uint32_t generation) override
|
||||
{
|
||||
std::size_t removed = 0;
|
||||
ui::map_tiles::LoadTileCommand kept[4]{};
|
||||
std::size_t kept_count = 0;
|
||||
for (std::size_t i = 0; i < count; ++i)
|
||||
{
|
||||
if (commands[i].runtime.generation == generation)
|
||||
{
|
||||
++removed;
|
||||
}
|
||||
else
|
||||
{
|
||||
kept[kept_count++] = commands[i];
|
||||
}
|
||||
}
|
||||
for (std::size_t i = 0; i < kept_count; ++i)
|
||||
{
|
||||
commands[i] = kept[i];
|
||||
}
|
||||
count = kept_count;
|
||||
return removed;
|
||||
}
|
||||
|
||||
ui::map_tiles::LoadTileCommand commands[4]{};
|
||||
std::size_t count = 0;
|
||||
};
|
||||
|
||||
class MapUiSink final : public ui::map_tiles::IMapTileUiSink
|
||||
{
|
||||
public:
|
||||
bool applyTile(const ui::map_tiles::MapTileEvent& event) override
|
||||
{
|
||||
last = event;
|
||||
++count;
|
||||
return true;
|
||||
}
|
||||
|
||||
ui::map_tiles::MapTileEvent last{};
|
||||
std::size_t count = 0;
|
||||
};
|
||||
|
||||
class TrackEvents final : public gps::runtime::ITrackEventSink
|
||||
{
|
||||
public:
|
||||
bool publish(const gps::runtime::TrackEvent& event) override
|
||||
{
|
||||
if (count >= 8)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
events[count++] = event;
|
||||
return true;
|
||||
}
|
||||
|
||||
gps::runtime::TrackEvent events[8]{};
|
||||
std::size_t count = 0;
|
||||
};
|
||||
|
||||
class TrackFiles final : public gps::runtime::ITrackFileAdapter
|
||||
{
|
||||
public:
|
||||
bool open(uint32_t track_id) override
|
||||
{
|
||||
active_track = track_id;
|
||||
++open_count;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool append(uint32_t track_id,
|
||||
const gps::runtime::TrackPoint* points,
|
||||
std::size_t count) override
|
||||
{
|
||||
assert(track_id == active_track);
|
||||
assert(points != nullptr || count == 0);
|
||||
appended += count;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool flush(uint32_t track_id) override
|
||||
{
|
||||
assert(track_id == active_track);
|
||||
++flush_count;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool close(uint32_t track_id) override
|
||||
{
|
||||
assert(track_id == active_track);
|
||||
++close_count;
|
||||
return true;
|
||||
}
|
||||
|
||||
std::size_t list(uint32_t* track_ids, std::size_t capacity) override
|
||||
{
|
||||
if (track_ids && capacity > 0)
|
||||
{
|
||||
track_ids[0] = active_track;
|
||||
return 1;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
uint32_t active_track = 0;
|
||||
std::size_t appended = 0;
|
||||
std::size_t open_count = 0;
|
||||
std::size_t flush_count = 0;
|
||||
std::size_t close_count = 0;
|
||||
};
|
||||
|
||||
ui::map_tiles::MapTileRef tileRef()
|
||||
{
|
||||
ui::map_tiles::MapTileRef ref{};
|
||||
ref.layer = ui::map_tiles::MapTileLayer::Osm;
|
||||
ref.z = 12;
|
||||
ref.x = 100;
|
||||
ref.y = 200;
|
||||
return ref;
|
||||
}
|
||||
|
||||
void test_map_tile_runtime_contract()
|
||||
{
|
||||
MapCommandSink commands;
|
||||
ui::map_tiles::MapTileAsyncRuntime async(commands);
|
||||
ui::map_tiles::MapTileStateMachine states;
|
||||
MapUiSink ui;
|
||||
ui::map_tiles::MapTileRuntime runtime(async, states, &ui);
|
||||
|
||||
ui::map_tiles::MapViewportPlan plan{};
|
||||
plan.generation = 7;
|
||||
plan.tile_count = 1;
|
||||
plan.tiles[0] = tileRef();
|
||||
assert(runtime.requestVisibleTiles(plan, 10) == 1);
|
||||
assert(commands.count == 1);
|
||||
|
||||
ui::map_tiles::MapTileEvent ready{};
|
||||
ready.kind = ui::map_tiles::MapTileEventKind::Ready;
|
||||
ready.generation = 7;
|
||||
ready.tile = tileRef();
|
||||
ui::map_tiles::MapTileRenderQueue render_queue;
|
||||
assert(runtime.handle(ready, render_queue));
|
||||
assert(ui.count == 1);
|
||||
assert(runtime.snapshot().ready_count == 1);
|
||||
}
|
||||
|
||||
void test_persistence_runtime_contract()
|
||||
{
|
||||
sys::runtime::RuntimeHarness harness;
|
||||
sys::runtime::DirtyStoreRegistry<4> registry;
|
||||
sys::runtime::DefaultPersistencePolicy policy;
|
||||
sys::runtime::PersistenceWorker worker(harness.storage(),
|
||||
harness.storage(),
|
||||
harness.bus(),
|
||||
harness.events(),
|
||||
policy);
|
||||
sys::runtime::PersistenceRuntime<4> runtime(registry,
|
||||
worker,
|
||||
harness.events(),
|
||||
policy);
|
||||
|
||||
assert(runtime.markDirty("nodes", 0));
|
||||
runtime.tick(100);
|
||||
assert(harness.storage().writeCount() == 0);
|
||||
runtime.tick(150);
|
||||
assert(harness.storage().writeCount() == 1);
|
||||
assert(harness.bus().acquireCount() == 1);
|
||||
assert(harness.events().persistenceCount() >= 3);
|
||||
}
|
||||
|
||||
void test_feedback_runtime_contract()
|
||||
{
|
||||
sys::runtime::RuntimeHarness harness;
|
||||
sys::runtime::FeedbackQueue<4> queue;
|
||||
sys::runtime::DefaultFeedbackPolicy policy;
|
||||
sys::runtime::FeedbackRuntime<4> runtime(queue,
|
||||
policy,
|
||||
harness.feedbackPresenter(),
|
||||
harness.events());
|
||||
|
||||
sys::runtime::NoticeIntent sent{};
|
||||
sent.category = sys::runtime::NoticeCategory::ChatDelivery;
|
||||
sent.severity = sys::runtime::NoticeSeverity::Success;
|
||||
sys::runtime::setNoticeMessage(sent, "Sent");
|
||||
sent.dedupe_key = 42;
|
||||
sent.created_at_ms = 10;
|
||||
assert(runtime.post(sent));
|
||||
assert(runtime.post(sent));
|
||||
assert(runtime.drainToUi(20) == 1);
|
||||
assert(harness.feedbackPresenter().captured() == 1);
|
||||
assert(harness.feedbackPresenter().capturedNotice(0).duration_ms == 1400);
|
||||
}
|
||||
|
||||
void test_track_runtime_contract()
|
||||
{
|
||||
sys::runtime::RuntimeHarness harness;
|
||||
TrackFiles files;
|
||||
TrackEvents events;
|
||||
gps::runtime::DefaultTrackFlushPolicy policy;
|
||||
gps::runtime::TrackStorageWorker worker(files, harness.bus(), events, policy);
|
||||
gps::runtime::TrackPointBuffer<8> points;
|
||||
gps::runtime::TrackStateMachine states;
|
||||
gps::runtime::TrackRuntime<8> runtime(points, states, policy, worker, events);
|
||||
|
||||
assert(runtime.startNewTrack(123, 0));
|
||||
runtime.tick(1);
|
||||
assert(events.count == 1);
|
||||
runtime.handle(events.events[0]);
|
||||
assert(states.state() == gps::runtime::TrackRecorderStatus::Recording);
|
||||
|
||||
gps::runtime::TrackPoint point{};
|
||||
point.latitude = 1.0;
|
||||
point.longitude = 2.0;
|
||||
for (int i = 0; i < 8; ++i)
|
||||
{
|
||||
point.timestamp_ms = static_cast<uint32_t>(i);
|
||||
assert(runtime.appendPoint(point, 10 + static_cast<uint32_t>(i)));
|
||||
}
|
||||
runtime.tick(30);
|
||||
assert(files.appended == 8);
|
||||
assert(files.flush_count == 1);
|
||||
}
|
||||
|
||||
void test_runtime_harness_keeps_ui_drain_separate()
|
||||
{
|
||||
sys::runtime::RuntimeHarness harness;
|
||||
harness.advance(5);
|
||||
harness.ui().assertNoBlockingCalls();
|
||||
assert(harness.drainUi() == 0);
|
||||
}
|
||||
|
||||
} // namespace
|
||||
|
||||
int main()
|
||||
{
|
||||
test_map_tile_runtime_contract();
|
||||
test_persistence_runtime_contract();
|
||||
test_feedback_runtime_contract();
|
||||
test_track_runtime_contract();
|
||||
test_runtime_harness_keeps_ui_drain_separate();
|
||||
return 0;
|
||||
}
|
||||
@@ -42,8 +42,11 @@ struct LoadTileCommand
|
||||
MapTileRef tile{};
|
||||
};
|
||||
|
||||
struct MapTileAsyncEvent
|
||||
using MapTileEventKind = MapTileAsyncEventKind;
|
||||
|
||||
class MapTileEvent
|
||||
{
|
||||
public:
|
||||
MapTileAsyncEventKind kind = MapTileAsyncEventKind::Failed;
|
||||
uint32_t command_id = 0;
|
||||
uint32_t generation = 0;
|
||||
@@ -53,6 +56,37 @@ struct MapTileAsyncEvent
|
||||
int32_t error = 0;
|
||||
};
|
||||
|
||||
using MapTileAsyncEvent = MapTileEvent;
|
||||
|
||||
struct MapTileDecodeInput
|
||||
{
|
||||
const uint8_t* payload = nullptr;
|
||||
std::size_t payload_size = 0;
|
||||
MapTileFormat format = MapTileFormat::Unknown;
|
||||
};
|
||||
|
||||
struct MapTileDecodeResult
|
||||
{
|
||||
bool ok = false;
|
||||
int32_t error = 0;
|
||||
};
|
||||
|
||||
class IMapTileDecoder
|
||||
{
|
||||
public:
|
||||
virtual ~IMapTileDecoder() = default;
|
||||
|
||||
virtual MapTileDecodeResult decode(const MapTileDecodeInput& input) = 0;
|
||||
};
|
||||
|
||||
class IMapTileUiSink
|
||||
{
|
||||
public:
|
||||
virtual ~IMapTileUiSink() = default;
|
||||
|
||||
virtual bool applyTile(const MapTileEvent& event) = 0;
|
||||
};
|
||||
|
||||
class IMapTileCommandSink
|
||||
{
|
||||
public:
|
||||
@@ -83,6 +117,57 @@ class IMapTileWorkerBackend
|
||||
MapTileFormat& out_format) = 0;
|
||||
};
|
||||
|
||||
struct MapTileStateSnapshot
|
||||
{
|
||||
uint32_t active_generation = 0;
|
||||
std::size_t ready_count = 0;
|
||||
std::size_t failed_count = 0;
|
||||
int32_t last_error = 0;
|
||||
};
|
||||
|
||||
class MapTileStateMachine
|
||||
{
|
||||
public:
|
||||
void transition(const MapTileEvent& event)
|
||||
{
|
||||
active_generation_ = event.generation;
|
||||
last_error_ = event.error;
|
||||
if (event.kind == MapTileAsyncEventKind::Ready)
|
||||
{
|
||||
++ready_count_;
|
||||
}
|
||||
else if (event.kind == MapTileAsyncEventKind::Failed ||
|
||||
event.kind == MapTileAsyncEventKind::ResourceBusy)
|
||||
{
|
||||
++failed_count_;
|
||||
}
|
||||
}
|
||||
|
||||
void cancelGeneration(uint32_t generation)
|
||||
{
|
||||
if (active_generation_ == generation)
|
||||
{
|
||||
active_generation_ = 0;
|
||||
}
|
||||
}
|
||||
|
||||
MapTileStateSnapshot snapshot() const
|
||||
{
|
||||
MapTileStateSnapshot snapshot{};
|
||||
snapshot.active_generation = active_generation_;
|
||||
snapshot.ready_count = ready_count_;
|
||||
snapshot.failed_count = failed_count_;
|
||||
snapshot.last_error = last_error_;
|
||||
return snapshot;
|
||||
}
|
||||
|
||||
private:
|
||||
uint32_t active_generation_ = 0;
|
||||
std::size_t ready_count_ = 0;
|
||||
std::size_t failed_count_ = 0;
|
||||
int32_t last_error_ = 0;
|
||||
};
|
||||
|
||||
class MapTileAsyncRuntime
|
||||
{
|
||||
public:
|
||||
@@ -91,6 +176,7 @@ class MapTileAsyncRuntime
|
||||
|
||||
uint32_t activeGeneration() const;
|
||||
std::size_t requestVisibleTiles(const MapViewportPlan& plan, uint32_t now_ms);
|
||||
std::size_t cancelGeneration(uint32_t generation);
|
||||
bool handleEvent(const MapTileAsyncEvent& event, MapTileRenderQueue& render_queue);
|
||||
|
||||
private:
|
||||
@@ -105,6 +191,49 @@ class MapTileAsyncRuntime
|
||||
uint32_t next_command_id_ = 1;
|
||||
};
|
||||
|
||||
class MapTileRuntime
|
||||
{
|
||||
public:
|
||||
MapTileRuntime(MapTileAsyncRuntime& runtime,
|
||||
MapTileStateMachine& state_machine,
|
||||
IMapTileUiSink* ui_sink = nullptr)
|
||||
: runtime_(runtime), state_machine_(state_machine), ui_sink_(ui_sink)
|
||||
{
|
||||
}
|
||||
|
||||
std::size_t requestVisibleTiles(const MapViewportPlan& plan, uint32_t now_ms)
|
||||
{
|
||||
return runtime_.requestVisibleTiles(plan, now_ms);
|
||||
}
|
||||
|
||||
std::size_t cancelGeneration(uint32_t generation)
|
||||
{
|
||||
state_machine_.cancelGeneration(generation);
|
||||
return runtime_.cancelGeneration(generation);
|
||||
}
|
||||
|
||||
bool handle(const MapTileEvent& event, MapTileRenderQueue& render_queue)
|
||||
{
|
||||
state_machine_.transition(event);
|
||||
const bool accepted = runtime_.handleEvent(event, render_queue);
|
||||
if (accepted && ui_sink_)
|
||||
{
|
||||
(void)ui_sink_->applyTile(event);
|
||||
}
|
||||
return accepted;
|
||||
}
|
||||
|
||||
MapTileStateSnapshot snapshot() const
|
||||
{
|
||||
return state_machine_.snapshot();
|
||||
}
|
||||
|
||||
private:
|
||||
MapTileAsyncRuntime& runtime_;
|
||||
MapTileStateMachine& state_machine_;
|
||||
IMapTileUiSink* ui_sink_ = nullptr;
|
||||
};
|
||||
|
||||
class MapTileWorker
|
||||
{
|
||||
public:
|
||||
|
||||
@@ -60,6 +60,15 @@ std::size_t MapTileAsyncRuntime::requestVisibleTiles(const MapViewportPlan& plan
|
||||
return queued;
|
||||
}
|
||||
|
||||
std::size_t MapTileAsyncRuntime::cancelGeneration(uint32_t generation)
|
||||
{
|
||||
if (active_generation_ == generation)
|
||||
{
|
||||
active_generation_ = 0;
|
||||
}
|
||||
return commands_.cancelGeneration(generation);
|
||||
}
|
||||
|
||||
bool MapTileAsyncRuntime::handleEvent(const MapTileAsyncEvent& event, MapTileRenderQueue& render_queue)
|
||||
{
|
||||
if (event.generation != active_generation_)
|
||||
|
||||
@@ -1,29 +1,20 @@
|
||||
#include "ui/runtime/ui_feedback.h"
|
||||
|
||||
#include "lvgl.h"
|
||||
#include "sys/feedback_runtime.h"
|
||||
#include "ui/widgets/system_notification.h"
|
||||
|
||||
#include <atomic>
|
||||
#include <cstddef>
|
||||
#include <cstring>
|
||||
|
||||
namespace ui::feedback
|
||||
{
|
||||
namespace
|
||||
{
|
||||
|
||||
constexpr size_t kMaxNoticeTextBytes = 192;
|
||||
constexpr size_t kNoticeQueueCapacity = 8;
|
||||
constexpr uint32_t kDrainPeriodMs = 20;
|
||||
|
||||
struct PostedNotice
|
||||
{
|
||||
char text[kMaxNoticeTextBytes]{};
|
||||
uint32_t duration_ms = 3000;
|
||||
Severity severity = Severity::Info;
|
||||
bool hide = false;
|
||||
};
|
||||
|
||||
class LvglSystemNotificationPresenter final : public IFeedbackPresenter
|
||||
{
|
||||
public:
|
||||
@@ -47,11 +38,12 @@ LvglSystemNotificationPresenter s_default_presenter;
|
||||
IFeedbackPresenter* s_presenter = &s_default_presenter;
|
||||
bool s_ready = false;
|
||||
lv_timer_t* s_drain_timer = nullptr;
|
||||
|
||||
PostedNotice s_queue[kNoticeQueueCapacity]{};
|
||||
size_t s_queue_head = 0;
|
||||
size_t s_queue_count = 0;
|
||||
std::atomic_flag s_queue_lock = ATOMIC_FLAG_INIT;
|
||||
std::atomic<uint32_t> s_hide_requests{0};
|
||||
sys::runtime::FeedbackEvent s_last_feedback_event{};
|
||||
|
||||
IFeedbackPresenter& active_presenter();
|
||||
void ensure_presenter_ready();
|
||||
|
||||
class QueueLock
|
||||
{
|
||||
@@ -72,43 +64,71 @@ class QueueLock
|
||||
QueueLock& operator=(const QueueLock&) = delete;
|
||||
};
|
||||
|
||||
void copy_notice_text(char* out, size_t out_len, const char* text)
|
||||
sys::runtime::NoticeSeverity to_runtime_severity(Severity severity)
|
||||
{
|
||||
if (!out || out_len == 0)
|
||||
switch (severity)
|
||||
{
|
||||
return;
|
||||
case Severity::Success:
|
||||
return sys::runtime::NoticeSeverity::Success;
|
||||
case Severity::Warning:
|
||||
return sys::runtime::NoticeSeverity::Warning;
|
||||
case Severity::Error:
|
||||
return sys::runtime::NoticeSeverity::Error;
|
||||
case Severity::Info:
|
||||
default:
|
||||
return sys::runtime::NoticeSeverity::Info;
|
||||
}
|
||||
std::strncpy(out, text ? text : "", out_len - 1);
|
||||
out[out_len - 1] = '\0';
|
||||
}
|
||||
|
||||
void enqueue_notice(const PostedNotice& notice)
|
||||
Severity from_runtime_severity(sys::runtime::NoticeSeverity severity)
|
||||
{
|
||||
QueueLock lock;
|
||||
if (s_queue_count == kNoticeQueueCapacity)
|
||||
switch (severity)
|
||||
{
|
||||
s_queue_head = (s_queue_head + 1) % kNoticeQueueCapacity;
|
||||
--s_queue_count;
|
||||
case sys::runtime::NoticeSeverity::Success:
|
||||
return Severity::Success;
|
||||
case sys::runtime::NoticeSeverity::Warning:
|
||||
return Severity::Warning;
|
||||
case sys::runtime::NoticeSeverity::Error:
|
||||
return Severity::Error;
|
||||
case sys::runtime::NoticeSeverity::Info:
|
||||
default:
|
||||
return Severity::Info;
|
||||
}
|
||||
|
||||
const size_t tail = (s_queue_head + s_queue_count) % kNoticeQueueCapacity;
|
||||
s_queue[tail] = notice;
|
||||
++s_queue_count;
|
||||
}
|
||||
|
||||
bool pop_notice(PostedNotice& out)
|
||||
class RuntimeFeedbackEventSink final : public sys::runtime::IFeedbackEventSink
|
||||
{
|
||||
QueueLock lock;
|
||||
if (s_queue_count == 0)
|
||||
public:
|
||||
bool publish(const sys::runtime::FeedbackEvent& event) override
|
||||
{
|
||||
return false;
|
||||
s_last_feedback_event = event;
|
||||
return true;
|
||||
}
|
||||
};
|
||||
|
||||
out = s_queue[s_queue_head];
|
||||
s_queue_head = (s_queue_head + 1) % kNoticeQueueCapacity;
|
||||
--s_queue_count;
|
||||
return true;
|
||||
}
|
||||
class RuntimeFeedbackPresenter final : public sys::runtime::IFeedbackPresenter
|
||||
{
|
||||
public:
|
||||
bool present(const sys::runtime::NoticeIntent& intent) override
|
||||
{
|
||||
ensure_presenter_ready();
|
||||
NoticeIntent ui_intent{};
|
||||
ui_intent.text = intent.message;
|
||||
ui_intent.duration_ms = intent.duration_ms;
|
||||
ui_intent.severity = from_runtime_severity(intent.severity);
|
||||
active_presenter().show_notice(ui_intent);
|
||||
return true;
|
||||
}
|
||||
};
|
||||
|
||||
sys::runtime::FeedbackQueue<kNoticeQueueCapacity> s_feedback_queue;
|
||||
sys::runtime::DefaultFeedbackPolicy s_feedback_policy;
|
||||
RuntimeFeedbackEventSink s_feedback_events;
|
||||
RuntimeFeedbackPresenter s_feedback_presenter;
|
||||
sys::runtime::FeedbackRuntime<kNoticeQueueCapacity> s_feedback_runtime(s_feedback_queue,
|
||||
s_feedback_policy,
|
||||
s_feedback_presenter,
|
||||
s_feedback_events);
|
||||
|
||||
IFeedbackPresenter& active_presenter()
|
||||
{
|
||||
@@ -131,21 +151,13 @@ void ensure_presenter_ready()
|
||||
|
||||
void drain_queued_notices()
|
||||
{
|
||||
PostedNotice payload{};
|
||||
while (pop_notice(payload))
|
||||
QueueLock lock;
|
||||
(void)s_feedback_runtime.drainToUi(lv_tick_get());
|
||||
const uint32_t hide_count = s_hide_requests.exchange(0, std::memory_order_acq_rel);
|
||||
if (hide_count > 0)
|
||||
{
|
||||
ensure_presenter_ready();
|
||||
if (payload.hide)
|
||||
{
|
||||
active_presenter().hide_notice();
|
||||
continue;
|
||||
}
|
||||
|
||||
NoticeIntent intent{};
|
||||
intent.text = payload.text;
|
||||
intent.duration_ms = payload.duration_ms;
|
||||
intent.severity = payload.severity;
|
||||
active_presenter().show_notice(intent);
|
||||
active_presenter().hide_notice();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -194,13 +206,15 @@ bool is_ready()
|
||||
|
||||
bool show_notice(const NoticeIntent& intent)
|
||||
{
|
||||
PostedNotice payload{};
|
||||
copy_notice_text(payload.text, sizeof(payload.text), intent.text);
|
||||
payload.duration_ms = intent.duration_ms;
|
||||
payload.severity = intent.severity;
|
||||
payload.hide = false;
|
||||
enqueue_notice(payload);
|
||||
return true;
|
||||
sys::runtime::NoticeIntent runtime_intent{};
|
||||
sys::runtime::setNoticeMessage(runtime_intent, intent.text);
|
||||
runtime_intent.category = sys::runtime::NoticeCategory::General;
|
||||
runtime_intent.severity = to_runtime_severity(intent.severity);
|
||||
runtime_intent.duration_ms = intent.duration_ms;
|
||||
runtime_intent.created_at_ms = lv_tick_get();
|
||||
|
||||
QueueLock lock;
|
||||
return s_feedback_runtime.post(runtime_intent);
|
||||
}
|
||||
|
||||
bool show_notice(const char* text, uint32_t duration_ms)
|
||||
@@ -214,9 +228,7 @@ bool show_notice(const char* text, uint32_t duration_ms)
|
||||
|
||||
bool hide_notice()
|
||||
{
|
||||
PostedNotice payload{};
|
||||
payload.hide = true;
|
||||
enqueue_notice(payload);
|
||||
s_hide_requests.fetch_add(1, std::memory_order_acq_rel);
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user