11 std::shared_ptr<std::atomic_bool>
cancelled = std::make_shared<std::atomic_bool>(
false);
14TaskContext::TaskContext(std::shared_ptr<std::atomic_bool> cancelled, std::function<
void(
float)> progress)
15 : cancelled_(
std::move(cancelled)), progress_(
std::move(progress)) {}
18 return cancelled_ && cancelled_->load(std::memory_order_relaxed);
22 if (progress_) progress_(std::clamp(progress, 0.0F, 1.0F));
29 std::lock_guard lock(mutex_);
31 for (
auto& [
id, task] : tasks_) {
33 task->cancelled->store(
true, std::memory_order_relaxed);
36 workReady_.notify_all();
37 if (worker_.joinable()) worker_.join();
41 if (
name.empty() || !work)
42 return eve::editing::failed<TaskId>(Status::Rejected,
RuleId(
"editing.task.invalid"),
43 "Task name and worker are required");
44 auto task = std::make_shared<TaskRecord>();
45 task->snapshot.name = std::move(
name);
46 task->work = std::move(work);
49 std::lock_guard lock(mutex_);
51 return eve::editing::failed<TaskId>(Status::Rejected,
RuleId(
"editing.task.stopping"),
52 "Task service is stopping");
53 task->snapshot.id =
TaskId(
"editor.task." + std::to_string(++sequence_));
54 id = task->snapshot.id;
55 tasks_.emplace(
id, task);
56 queue_.push_back(std::move(task));
58 workReady_.notify_one();
59 return eve::editing::applied<TaskId>(
id);
63 std::lock_guard lock(mutex_);
64 auto found = tasks_.find(task);
65 if (
found == tasks_.end())
66 return eve::editing::failed<void>(Status::NotFound,
RuleId(
"editing.task.not-found"),
67 "Task does not exist");
68 auto& record = *
found->second;
71 return eve::editing::failed<void>(Status::Conflict,
RuleId(
"editing.task.already-finished"),
72 "Finished tasks cannot be cancelled");
73 record.cancelled->store(
true, std::memory_order_relaxed);
75 workReady_.notify_all();
76 return eve::editing::applied<void>();
80 std::lock_guard lock(mutex_);
81 auto found = tasks_.find(task);
82 if (
found == tasks_.end())
83 return eve::editing::failed<TaskSnapshot>(Status::NotFound,
RuleId(
"editing.task.not-found"),
84 "Task does not exist");
85 return eve::editing::applied<TaskSnapshot>(
found->second->snapshot);
89 std::unique_lock lock(mutex_);
90 const bool idle = idle_.wait_for(lock, timeout, [
this] {
91 if (!queue_.empty())
return false;
92 return std::none_of(tasks_.begin(), tasks_.end(),
93 [](
const auto& entry) { return entry.second->snapshot.state == TaskState::Running; });
97 RuleId(
"editing.task.wait-timeout"),
98 "Task service did not become idle before the deadline");
99 return eve::editing::applied<void>();
102void TaskService::workerLoop() {
104 std::shared_ptr<TaskRecord> task;
106 std::unique_lock lock(mutex_);
107 workReady_.wait(lock, [
this] {
return stopping_ || !queue_.empty(); });
108 if (stopping_ && queue_.empty())
return;
109 task = queue_.front();
111 if (task->cancelled->load(std::memory_order_relaxed)) {
121 TaskContext
context(task->cancelled, [
this, weak = std::weak_ptr<TaskRecord>(task)](
float value) {
122 if (auto record = weak.lock()) {
123 std::lock_guard lock(mutex_);
124 record->snapshot.progress = value;
128 }
catch (
const std::exception&
error) {
129 outcome.status = Status::Failed;
131 RuleId(
"editor.task.exception"),
132 DiagnosticSeverity::Error,
error.what()));
134 outcome.status = Status::Failed;
136 RuleId(
"editor.task.exception"),
137 DiagnosticSeverity::Error,
138 "Unknown worker exception"));
142 std::lock_guard lock(mutex_);
143 task->snapshot.output = std::move(outcome.output);
144 task->snapshot.diagnostics = std::move(outcome.diagnostics);
145 task->snapshot.progress = 1.0F;
146 if (task->cancelled->load(std::memory_order_relaxed) || outcome.status == Status::Cancelled)
148 else if (outcome.status == Status::Applied || outcome.status == Status::NoOp)
const VegetationPresetContext & context
Move-only operation result carrying either a value or Status.
UUID-backed identifier adapter for legacy textual boundaries.
void reportProgress(float progress) const
Publish normalized progress in the inclusive range 0..1.
bool isCancellationRequested() const
True after cancellation was requested.
Result< TaskId > submit(std::string name, Work work)
Queue background work and return its stable task identity.
Result< void > cancel(const TaskId &task)
Request cooperative cancellation of queued or running work.
TaskService()
Start the editor background worker.
std::function< TaskOutcome(const TaskContext &)> Work
Result< void > waitIdle(std::chrono::milliseconds timeout)
Wait until no task is queued or running; intended for shutdown and tests.
~TaskService()
Cooperatively cancel pending work and join the worker.
Result< TaskSnapshot > snapshot(const TaskId &task) const
Return a thread-safe task snapshot.
StrongId< RuleIdTag > RuleId
StrongId< TaskIdTag > TaskId
Diagnostic ruleDiagnostic(eve::DiagnosticCode code, RuleId rule, DiagnosticSeverity severity, std::string message)
Build a common diagnostic carrying an open editing rule identity.
std::shared_ptr< std::atomic_bool > cancelled