载入中...
搜索中...
未找到
EditingTaskService.cpp
浏览该文件的文档.
2
3#include <algorithm>
4#include <exception>
5
6namespace eve::editing {
7
9 TaskSnapshot snapshot;
10 Work work;
11 std::shared_ptr<std::atomic_bool> cancelled = std::make_shared<std::atomic_bool>(false);
12};
13
14TaskContext::TaskContext(std::shared_ptr<std::atomic_bool> cancelled, std::function<void(float)> progress)
15 : cancelled_(std::move(cancelled)), progress_(std::move(progress)) {}
16
18 return cancelled_ && cancelled_->load(std::memory_order_relaxed);
19}
20
21void TaskContext::reportProgress(float progress) const {
22 if (progress_) progress_(std::clamp(progress, 0.0F, 1.0F));
23}
24
25TaskService::TaskService() : worker_([this] { workerLoop(); }) {}
26
28 {
29 std::lock_guard lock(mutex_);
30 stopping_ = true;
31 for (auto& [id, task] : tasks_) {
32 (void)id;
33 task->cancelled->store(true, std::memory_order_relaxed);
34 }
35 }
36 workReady_.notify_all();
37 if (worker_.joinable()) worker_.join();
38}
39
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);
47 TaskId id;
48 {
49 std::lock_guard lock(mutex_);
50 if (stopping_)
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));
57 }
58 workReady_.notify_one();
59 return eve::editing::applied<TaskId>(id);
60}
61
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;
69 if (record.snapshot.state == TaskState::Succeeded || record.snapshot.state == TaskState::Failed ||
70 record.snapshot.state == TaskState::Cancelled)
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);
74 if (record.snapshot.state == TaskState::Queued) record.snapshot.state = TaskState::Cancelled;
75 workReady_.notify_all();
76 return eve::editing::applied<void>();
77}
78
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);
86}
87
88Result<void> TaskService::waitIdle(std::chrono::milliseconds timeout) {
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; });
94 });
95 if (!idle)
96 return eve::editing::failed<void>(Status::Failed, eve::DiagnosticCode::Failed,
97 RuleId("editing.task.wait-timeout"),
98 "Task service did not become idle before the deadline");
99 return eve::editing::applied<void>();
100}
101
102void TaskService::workerLoop() {
103 while (true) {
104 std::shared_ptr<TaskRecord> task;
105 {
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();
110 queue_.pop_front();
111 if (task->cancelled->load(std::memory_order_relaxed)) {
112 task->snapshot.state = TaskState::Cancelled;
113 idle_.notify_all();
114 continue;
115 }
116 task->snapshot.state = TaskState::Running;
117 }
118
119 TaskOutcome outcome;
120 try {
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;
125 }
126 });
127 outcome = task->work(context);
128 } catch (const std::exception& error) {
129 outcome.status = Status::Failed;
130 outcome.diagnostics.push_back(ruleDiagnostic(eve::DiagnosticCode::Failed,
131 RuleId("editor.task.exception"),
132 DiagnosticSeverity::Error, error.what()));
133 } catch (...) {
134 outcome.status = Status::Failed;
135 outcome.diagnostics.push_back(ruleDiagnostic(eve::DiagnosticCode::Failed,
136 RuleId("editor.task.exception"),
137 DiagnosticSeverity::Error,
138 "Unknown worker exception"));
139 }
140
141 {
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)
147 task->snapshot.state = TaskState::Cancelled;
148 else if (outcome.status == Status::Applied || outcome.status == Status::NoOp)
149 task->snapshot.state = TaskState::Succeeded;
150 else
151 task->snapshot.state = TaskState::Failed;
152 task->work = {};
153 }
154 idle_.notify_all();
155 }
156}
157
158} // namespace eve::editing
double value
std::string name
std::string error
Definition Package.cpp:60
std::string id
Definition PlayHost.cpp:108
bool found
const VegetationPresetContext & context
Move-only operation result carrying either a value or Status.
Definition Result.h:155
UUID-backed identifier adapter for legacy textual boundaries.
Definition Identity.h:314
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
Definition EditingIds.h:60
StrongId< TaskIdTag > TaskId
Definition EditingIds.h:70
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