载入中...
搜索中...
未找到
Production.cpp
浏览该文件的文档.
2
3#include "common/Json.h"
4#include <algorithm>
5#include <cmath>
6#include <exception>
7#include <functional>
8#include <iomanip>
9#include <limits>
10#include <set>
11#include <sstream>
12#include <utility>
13
14namespace eve::production {
15namespace {
16
17template <class T>
18eve::Result<T> productionBindingFailure(eve::DiagnosticCode code, std::string message, std::string path = {}) {
20 eve::Diagnostic::error(code, message, path, {}, "production.squirrel"));
21}
22
23bool terminal(TaskState state) {
25}
26
27struct ScaledWork {
28 std::int64_t appliedNanoseconds = 0;
29 std::uint32_t remainderPermille = 0;
30};
31
32eve::Result<ScaledWork> scaleWork(std::int64_t effortNanoseconds, std::uint32_t efficiencyPermille,
33 std::uint32_t remainderPermille) {
34 if (effortNanoseconds < 0 || efficiencyPermille == 0 || remainderPermille >= 1000)
35 return productionBindingFailure<ScaledWork>(eve::DiagnosticCode::InvalidArgument,
36 "invalid production work scaling input");
37 const auto whole = effortNanoseconds / 1000;
38 const auto tail = effortNanoseconds % 1000;
39 if (whole > std::numeric_limits<std::int64_t>::max() / efficiencyPermille)
40 return productionBindingFailure<ScaledWork>(eve::DiagnosticCode::InvalidArgument,
41 "scaled production work overflows duration");
42 const auto tailNumerator = tail * static_cast<std::int64_t>(efficiencyPermille) + remainderPermille;
43 const auto tailApplied = tailNumerator / 1000;
44 const auto wholeApplied = whole * static_cast<std::int64_t>(efficiencyPermille);
45 if (wholeApplied > std::numeric_limits<std::int64_t>::max() - tailApplied)
46 return productionBindingFailure<ScaledWork>(eve::DiagnosticCode::InvalidArgument,
47 "scaled production work overflows duration");
49 {wholeApplied + tailApplied, static_cast<std::uint32_t>(tailNumerator % 1000)});
50}
51
52} // namespace
53
54WorkQueue::WorkQueue(eve::PersistentId instanceId) : instanceId_(instanceId) {}
55
56std::string_view taskStateName(TaskState state) {
57 switch (state) {
58 case TaskState::Queued: return "queued";
59 case TaskState::Running: return "running";
60 case TaskState::Paused: return "paused";
61 case TaskState::ReadyToSettle: return "ready_to_settle";
62 case TaskState::SettlementFailed: return "settlement_failed";
63 case TaskState::Completed: return "completed";
64 case TaskState::Cancelled: return "cancelled";
65 case TaskState::Failed: return "failed";
66 }
67 return "unknown";
68}
69
71 switch (kind) {
72 case ProductionEventKind::Enqueued: return "enqueued";
73 case ProductionEventKind::Started: return "started";
74 case ProductionEventKind::Paused: return "paused";
75 case ProductionEventKind::Resumed: return "resumed";
76 case ProductionEventKind::ReadyToSettle: return "ready_to_settle";
77 case ProductionEventKind::SettlementFailed: return "settlement_failed";
78 case ProductionEventKind::Completed: return "completed";
79 case ProductionEventKind::Cancelled: return "cancelled";
80 case ProductionEventKind::Failed: return "failed";
81 }
82 return "unknown";
83}
84
85void WorkQueue::emit(ProductionEventKind kind, const ProductionTask& task, std::string_view reason) {
86 events_.push_back(
87 {{nextEventSequence_++, std::string(reason)}, tick_, kind, task.id, task.owner, task.kind, task.product,
88 task.correlationId});
89 if (const auto next = revision_.incremented()) revision_ = *next;
90}
91
92int WorkQueue::slotCount(std::string_view owner) const {
93 const auto it = std::lower_bound(slots_.begin(), slots_.end(), owner,
94 [](const auto& entry, std::string_view key) { return entry.first < key; });
95 return it != slots_.end() && it->first == owner ? it->second : 1;
96}
97
98int WorkQueue::runningCount(std::string_view owner) const {
99 return static_cast<int>(std::count_if(tasks_.begin(), tasks_.end(), [&owner](const auto& task) {
100 return task->owner == owner && task->state == TaskState::Running;
101 }));
102}
103
104int WorkQueue::resourceCapacity(std::string_view owner, std::string_view resource) const {
105 if (resource == "slot") return slotCount(owner);
106 const auto key = std::pair(std::string(owner), std::string(resource));
107 const auto it = std::lower_bound(resources_.begin(), resources_.end(), key,
108 [](const auto& entry, const auto& wanted) { return entry.first < wanted; });
109 return it != resources_.end() && it->first == key ? it->second : 0;
110}
111
112bool WorkQueue::resourcesAvailable(const ProductionTask& task) const {
113 for (const auto& requirement : task.resources) {
114 int used = 0;
115 for (const auto& running : tasks_) {
116 if (running->owner != task.owner || running->state != TaskState::Running) continue;
117 for (const auto& held : running->resources)
118 if (held.resource == requirement.resource) used += held.units;
119 }
120 if (used + requirement.units > resourceCapacity(task.owner, requirement.resource)) return false;
121 }
122 return true;
123}
124
126 const auto taskRef = find(taskId);
127 if (!taskRef)
128 return productionBindingFailure<std::string>(eve::DiagnosticCode::NotFound,
129 "work task was not found", "taskId");
130 const auto& task = taskRef->get();
131 if (!task.dependencies.prerequisites.empty()) {
132 bool anyCompleted = false;
133 for (const auto& dependencyId : task.dependencies.prerequisites) {
134 const auto dependency = find(dependencyId);
135 if (!dependency)
136 return eve::Result<std::string>::success("missing_prerequisite:" + dependencyId);
137 if (dependency->get().state == TaskState::Completed) {
138 anyCompleted = true;
139 continue;
140 }
143 std::string(terminal(dependency->get().state) ? "prerequisite_failed:" : "waiting_for:") +
144 dependencyId);
145 }
146 if (task.dependencies.mode == TaskDependencyMode::AnyOf && !anyCompleted)
147 return eve::Result<std::string>::success("waiting_for_any_prerequisite");
148 }
149 for (const auto& definition : task.dependencies.requiredDefinitions) {
150 if (!definitionAvailable(task.owner, definition))
152 "requires_definition:" + definition.reference.id().format() + "@" +
153 std::to_string(definition.generation.value()));
154 }
155 for (const auto& tag : task.dependencies.requiredTags)
156 if (!tagAvailable(task.owner, tag))
157 return eve::Result<std::string>::success("requires_tag:" + tag);
158 for (const auto& blockerId : task.dependencies.blocks) {
159 const auto blocker = find(blockerId);
160 if (!blocker)
161 return eve::Result<std::string>::success("missing_blocker:" + blockerId);
162 if (!terminal(blocker->get().state))
163 return eve::Result<std::string>::success("blocked_by:" + blockerId);
164 }
166 return eve::Result<std::string>::success("stock_target_satisfied");
167 if (!resourcesAvailable(task)) return eve::Result<std::string>::success("resource_unavailable");
168 return eve::Result<std::string>::success(std::string{});
169}
170
171bool WorkQueue::dependenciesSatisfied(const ProductionTask& task) const {
172 auto result = blockReason(task.id);
173 return result.ok() && result.value().empty();
174}
175
176void WorkQueue::scheduleAll() {
177 std::set<std::string> owners;
178 for (const auto& task : tasks_) owners.insert(task->owner);
179 for (const auto& owner : owners) schedule(owner);
180}
181
182void WorkQueue::schedule(std::string_view owner) {
183 int available = slotCount(owner) - runningCount(owner);
184 const auto strategy = schedulerStrategy(owner);
185 while (available-- > 0) {
186 ProductionTask* best = nullptr;
187 for (auto& candidate : tasks_) {
188 if (candidate->owner != owner || candidate->state != TaskState::Queued ||
189 !dependenciesSatisfied(*candidate))
190 continue;
191 if (!best) {
192 best = candidate.get();
193 continue;
194 }
195 const auto remaining = [](const ProductionTask& value) {
196 return value.duration.nanoseconds() - value.progress.nanoseconds();
197 };
198 const bool preferred = strategy == SchedulerStrategy::Fifo
199 ? candidate->enqueueSequence < best->enqueueSequence
201 ? remaining(*candidate) < remaining(*best) ||
202 (remaining(*candidate) == remaining(*best) &&
203 candidate->enqueueSequence < best->enqueueSequence)
204 : candidate->priority > best->priority ||
205 (candidate->priority == best->priority && candidate->enqueueSequence < best->enqueueSequence);
206 if (preferred) best = candidate.get();
207 }
208 if (!best) break;
209 best->state = TaskState::Running;
210 if (best->reservationState == ReservationState::Reserved)
211 best->reservationState = ReservationState::Started;
213 }
214}
215
216eve::Result<std::string> WorkQueue::enqueue(std::string_view owner, std::string_view kind, std::string_view product,
217 eve::Value context, double duration, int priority) {
218 auto converted = eve::Duration::fromSeconds(duration);
219 if (!converted) return eve::Result<std::string>::failure(converted.status());
221 request.owner = std::string(owner);
222 request.kind = std::string(kind);
223 request.product = std::string(product);
224 request.context = std::move(context);
225 request.duration = std::move(converted).takeValue();
226 request.priority = priority;
227 request.settlementRequired = false;
228 return enqueue(std::move(request));
229}
230
232 const auto& owner = request.owner;
233 const auto& kind = request.kind;
234 const auto& product = request.product;
235 if (owner.empty() || kind.empty() || product.empty())
236 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
237 "work task owner, kind and product are required");
238 if (!request.context.isObject() || !request.reservation.isObject())
239 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
240 "work task context and reservation must be Value objects");
241 if (request.duration.nanoseconds() <= 0)
242 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
243 "work task duration must be positive", "duration");
244 if (request.batchSize == 0 || request.efficiencyPermille == 0)
245 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
246 "batch size and efficiency must be positive", "batchSize");
247 if (request.repeat.totalCycles == 0 || request.repeat.maintainStockTarget < -1 ||
248 request.repeat.observedStock < 0 ||
249 (request.random.drawCount > 0 && request.random.stream.empty()))
250 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
251 "invalid repeat, stock or random provenance metadata");
252 std::set<std::string> resourceNames;
253 for (const auto& requirement : request.resources) {
254 if (requirement.resource.empty() || requirement.units <= 0 ||
255 !resourceNames.insert(requirement.resource).second)
256 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
257 "resource requirements must be named, positive and unique",
258 "resources");
259 }
260 std::set<std::string> dependencyIds;
261 const auto validateDependencies = [this, &dependencyIds](const std::vector<std::string>& values,
262 std::string_view path) -> eve::Result<void> {
263 for (const auto& id : values) {
264 if (id.empty() || !dependencyIds.insert(id).second)
265 return productionBindingFailure<void>(eve::DiagnosticCode::InvalidArgument,
266 "production dependency ids must be non-empty and unique",
267 std::string(path));
268 if (!find(id))
269 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound,
270 "production dependency does not exist", std::string(path));
271 }
273 };
274 auto prerequisites = validateDependencies(request.dependencies.prerequisites, "dependencies.prerequisites");
275 if (!prerequisites) return eve::Result<std::string>::failure(prerequisites.status());
276 auto blockers = validateDependencies(request.dependencies.blocks, "dependencies.blocks");
277 if (!blockers) return eve::Result<std::string>::failure(blockers.status());
278 std::set<std::pair<std::string, std::uint64_t>> definitionRequirements;
279 for (const auto& definition : request.dependencies.requiredDefinitions) {
280 const auto key = std::pair(definition.reference.id().format(), definition.generation.value());
281 if (!definition.isValid() || !definitionRequirements.insert(key).second)
282 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
283 "required definitions must be valid and unique",
284 "dependencies.requiredDefinitions");
285 }
286 std::set<std::string> tagRequirements;
287 for (const auto& tag : request.dependencies.requiredTags) {
288 if (tag.empty() || !tagRequirements.insert(tag).second)
289 return productionBindingFailure<std::string>(eve::DiagnosticCode::InvalidArgument,
290 "required tags must be non-empty and unique",
291 "dependencies.requiredTags");
292 }
293 auto task = std::make_unique<ProductionTask>();
294 std::ostringstream id;
295 id << "task-" << std::setw(16) << std::setfill('0') << nextTaskId_;
296 task->id = id.str();
297 task->owner = std::string(owner);
298 task->kind = std::string(kind);
299 task->product = std::string(product);
300 task->context = std::move(request.context);
301 task->duration = request.duration;
302 task->priority = request.priority;
303 task->definition = std::move(request.definition);
304 task->reservation = std::move(request.reservation);
305 if (const auto* reservation = task->reservation.getIf<eve::Value::Object>();
306 reservation != nullptr && !reservation->empty())
308 task->dependencies = std::move(request.dependencies);
309 task->resources = std::move(request.resources);
310 task->termination = request.termination;
311 task->random = std::move(request.random);
312 task->repeat = request.repeat;
313 task->correlationId = request.correlationId.empty() ? task->id : std::move(request.correlationId);
314 task->batchSize = request.batchSize;
315 task->efficiencyPermille = request.efficiencyPermille;
316 task->settlementRequired = request.settlementRequired;
317 task->enqueueSequence = nextEnqueueSequence_++;
318 const std::string result = task->id;
319 tasks_.push_back(std::move(task));
320 ++nextTaskId_;
321 emit(ProductionEventKind::Enqueued, *tasks_.back());
322 schedule(owner);
324}
325
326ProductionTask* WorkQueue::mutableFind(std::string_view taskId) {
327 const auto it =
328 std::find_if(tasks_.begin(), tasks_.end(), [&taskId](const auto& task) { return task->id == taskId; });
329 return it == tasks_.end() ? nullptr : it->get();
330}
331
333 const auto* task = mutableFind(taskId);
334 return task == nullptr ? eve::OptionalRef<const ProductionTask>{}
335 : eve::OptionalRef<const ProductionTask>(std::cref(*task));
336}
337
339 const auto it =
340 std::find_if(tasks_.begin(), tasks_.end(), [&taskId](const auto& task) { return task->id == taskId; });
341 return it == tasks_.end() ? eve::OptionalRef<const ProductionTask>{}
342 : eve::OptionalRef<const ProductionTask>(std::cref(*it->get()));
343}
344
346 auto* task = mutableFind(taskId);
347 if (task == nullptr)
348 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
349 if (task->state != TaskState::Queued && task->state != TaskState::Running)
350 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict,
351 "only queued or running tasks can be paused", "taskId");
352 task->state = TaskState::Paused;
353 emit(ProductionEventKind::Paused, *task);
354 schedule(task->owner);
356}
357
359 auto* task = mutableFind(taskId);
360 if (task == nullptr)
361 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
362 if (task->state != TaskState::Paused)
363 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict, "only paused tasks can be resumed",
364 "taskId");
365 task->state = TaskState::Queued;
366 emit(ProductionEventKind::Resumed, *task);
367 schedule(task->owner);
369}
370
371eve::Result<void> WorkQueue::cancel(std::string_view taskId, std::string_view reason) {
372 auto* task = mutableFind(taskId);
373 if (task == nullptr)
374 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
375 if (terminal(task->state))
376 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict, "terminal work tasks cannot be cancelled",
377 "taskId");
378 task->refundPermille = refundFor(*task, task->termination.cancellation);
380 task->reason = std::string(reason);
381 emit(ProductionEventKind::Cancelled, *task, reason);
382 scheduleAll();
384}
385
386eve::Result<void> WorkQueue::fail(std::string_view taskId, std::string_view reason) {
387 auto* task = mutableFind(taskId);
388 if (task == nullptr)
389 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
390 if (terminal(task->state))
391 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict, "terminal work tasks cannot fail",
392 "taskId");
393 task->refundPermille = refundFor(*task, task->termination.failure);
394 task->state = TaskState::Failed;
395 task->reason = std::string(reason);
396 emit(ProductionEventKind::Failed, *task, reason);
397 scheduleAll();
399}
400
402 auto* task = mutableFind(taskId);
403 if (task == nullptr)
404 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
405 if ((!task->settlement.settlementId.empty() && task->settlement.settlementId == receipt.settlementId) ||
406 (!task->lastSettlementId.empty() && task->lastSettlementId == receipt.settlementId))
409 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict, "work task is not ready to settle", "taskId");
410 if (receipt.settlementId.empty() || !receipt.payload.isObject())
411 return productionBindingFailure<void>(eve::DiagnosticCode::InvalidArgument,
412 "settlement requires an id and object payload", "receipt");
413 task->settlement = std::move(receipt);
416 task->reason.clear();
419 continueAfterCycle(*task);
420 scheduleAll();
422}
423
424eve::Result<void> WorkQueue::failSettlement(std::string_view taskId, std::string_view reason) {
425 auto* task = mutableFind(taskId);
426 if (task == nullptr)
427 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
428 if (task->state != TaskState::ReadyToSettle)
429 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict, "work task is not ready to settle", "taskId");
430 if (reason.empty())
431 return productionBindingFailure<void>(eve::DiagnosticCode::InvalidArgument, "settlement failure needs a reason");
433 task->reason = std::string(reason);
434 emit(ProductionEventKind::SettlementFailed, *task, reason);
436}
437
439 auto* task = mutableFind(taskId);
440 if (task == nullptr)
441 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
443 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict, "work task has no failed settlement", "taskId");
445 task->reason.clear();
448}
449
450std::uint32_t WorkQueue::refundFor(const ProductionTask& task, RefundPolicy policy) {
451 if (policy == RefundPolicy::None) return 0;
452 if (policy == RefundPolicy::Full) return 1000;
453 const auto duration = task.duration.nanoseconds();
454 if (duration <= 0 || task.progress >= task.duration) return 0;
455 const auto remaining = duration - task.progress.nanoseconds();
456 std::int64_t residue = 0;
457 std::uint32_t permille = 0;
458 for (std::uint32_t index = 0; index < 1000; ++index) {
459 if (residue >= duration - remaining) {
460 residue -= duration - remaining;
461 ++permille;
462 } else {
463 residue += remaining;
464 }
465 }
466 return std::min<std::uint32_t>(permille, 1000);
467}
468
469void WorkQueue::completeIfReady(ProductionTask& task) {
470 if (task.progress < task.duration) return;
471 task.progress = task.duration;
472 task.workRemainderPermille = 0;
473 if (task.reservationState == ReservationState::Started)
474 task.reservationState = ReservationState::Consumed;
475 if (task.settlementRequired) {
476 task.state = TaskState::ReadyToSettle;
478 } else {
479 task.settlement.settlementId = "automatic:" + task.id;
480 task.state = TaskState::Completed;
482 continueAfterCycle(task);
483 }
484}
485
487 auto* task = mutableFind(taskId);
488 if (task == nullptr)
489 return productionBindingFailure<void>(eve::DiagnosticCode::NotFound, "work task was not found", "taskId");
490 if (task->state != TaskState::Running)
491 return productionBindingFailure<void>(eve::DiagnosticCode::Conflict,
492 "work contributions require a running task", "taskId");
493 if (contribution.contributor.empty() || contribution.effort.nanoseconds() <= 0 ||
494 contribution.efficiencyPermille == 0)
495 return productionBindingFailure<void>(eve::DiagnosticCode::InvalidArgument,
496 "work contribution must have a contributor, effort and efficiency");
497 auto scaled = scaleWork(contribution.effort.nanoseconds(), contribution.efficiencyPermille,
498 task->workRemainderPermille);
499 if (!scaled) return eve::Result<void>::failure(scaled.status());
500 const auto remaining = task->duration.nanoseconds() - task->progress.nanoseconds();
501 task->workRemainderPermille = scaled.value().remainderPermille;
502 task->progress = eve::Duration::fromNanoseconds(
503 task->progress.nanoseconds() + std::min(scaled.value().appliedNanoseconds, remaining));
504 completeIfReady(*task);
505 if (task->state != TaskState::Running) scheduleAll();
507}
508
510 if (step.tick <= tick_)
512 "production step tick must be strictly newer"));
513 if (step.delta.nanoseconds() < 0)
515 "production step duration must be non-negative"));
516
517 struct PendingProgress {
518 ProductionTask* task = nullptr;
519 std::int64_t appliedNanoseconds = 0;
520 std::uint32_t remainderPermille = 0;
521 };
522 std::vector<PendingProgress> pending;
523 pending.reserve(tasks_.size());
524 const std::int64_t delta = step.delta.nanoseconds();
525 for (auto& task : tasks_) {
526 if (task->state != TaskState::Running) continue;
527 auto scaled = scaleWork(delta, task->efficiencyPermille, task->workRemainderPermille);
528 if (!scaled) return eve::Result<void>::failure(scaled.status());
529 const auto remaining = task->duration.nanoseconds() - task->progress.nanoseconds();
530 pending.push_back({task.get(), std::min(scaled.value().appliedNanoseconds, remaining),
531 scaled.value().remainderPermille});
532 }
533
534 tick_ = step.tick;
535 bool completedAny = false;
536 for (const auto& update : pending) {
537 update.task->workRemainderPermille = update.remainderPermille;
538 update.task->progress = eve::Duration::fromNanoseconds(
539 update.task->progress.nanoseconds() + update.appliedNanoseconds);
540 if (update.task->progress >= update.task->duration) {
541 completeIfReady(*update.task);
542 completedAny = true;
543 }
544 }
545 if (completedAny) scheduleAll();
547}
548
549eve::Result<void> WorkQueue::setSlotCount(std::string_view owner, int slots) {
550 if (owner.empty() || slots < 0)
551 return productionBindingFailure<void>(eve::DiagnosticCode::InvalidArgument,
552 "owner must be non-empty and slots must be non-negative");
553 auto it = std::lower_bound(slots_.begin(), slots_.end(), owner,
554 [](const auto& entry, std::string_view key) { return entry.first < key; });
555 if (it != slots_.end() && it->first == owner)
556 it->second = slots;
557 else
558 slots_.insert(it, {std::string(owner), slots});
559 schedule(owner);
561}
562
563eve::Result<void> WorkQueue::setResourceCapacity(std::string_view owner, std::string_view resource, int capacity) {
564 if (owner.empty() || resource.empty() || resource == "slot" || capacity < 0)
565 return productionBindingFailure<void>(eve::DiagnosticCode::InvalidArgument,
566 "resource capacity needs owner, non-slot name and non-negative value");
567 const auto key = std::pair(std::string(owner), std::string(resource));
568 auto it = std::lower_bound(resources_.begin(), resources_.end(), key,
569 [](const auto& entry, const auto& wanted) { return entry.first < wanted; });
570 if (it != resources_.end() && it->first == key)
571 it->second = capacity;
572 else
573 resources_.insert(it, {key, capacity});
574 schedule(owner);
576}
577
578int WorkQueue::taskCount() const { return static_cast<int>(tasks_.size()); }
579
581 if (index < 0 || static_cast<size_t>(index) >= tasks_.size()) return {};
582 return std::ref(*tasks_[static_cast<size_t>(index)]);
583}
584
586 if (index < 0 || static_cast<size_t>(index) >= tasks_.size()) return {};
587 return std::cref(*tasks_[static_cast<size_t>(index)]);
588}
589
590int WorkQueue::ownerTaskCount(std::string_view owner) const {
591 return static_cast<int>(
592 std::count_if(tasks_.begin(), tasks_.end(), [&owner](const auto& t) { return t->owner == owner; }));
593}
594
596 if (index < 0) return {};
597 for (auto& task : tasks_)
598 if (task->owner == owner && index-- == 0) return std::ref(*task);
599 return {};
600}
601
603 if (index < 0) return {};
604 for (const auto& task : tasks_)
605 if (task->owner == owner && index-- == 0) return std::cref(*task);
606 return {};
607}
608
609int WorkQueue::eventCount() const { return static_cast<int>(events_.size()); }
610
612 if (index < 0 || static_cast<size_t>(index) >= events_.size()) return {};
613 return std::ref(events_[static_cast<size_t>(index)]);
614}
615
617 if (index < 0 || static_cast<size_t>(index) >= events_.size()) return {};
618 return std::cref(events_[static_cast<size_t>(index)]);
619}
620
621void WorkQueue::clearEvents() { events_.clear(); }
622
623} // namespace eve::production
double value
float duration
std::vector< QuestEvent > pending
int priority
std::map< std::string, Var > values
std::string message
DiagnosticCode code
const GltfImportRequest & request
std::uint32_t capacity
std::uint32_t key
TokenKind kind
std::vector< std::string > tail
const std::string * tag
std::vector< std::weak_ptr< DeviceBytes > > resources
Definition OnnxGpgpu.cpp:55
std::string path
Definition PlayHost.cpp:110
std::string id
Definition PlayHost.cpp:108
std::string taskId
std::uint32_t remainderPermille
std::int64_t appliedNanoseconds
float t
std::string resource
LocalPageCacheEntry slots[ShadowConfig::kLocalSlots]
std::map< Cell, int > best
std::vector< ecs::EntityHandle > schedule
The round's activation queue, resolved against the target battle.
std::vector< UnitCandidate > units
float step
Definition TreeMesh.cpp:314
uint32_t index
const VegetationPresetContext & context
static Diagnostic error(DiagnosticCode code, std::string message, std::string path={}, DiagnosticDetails details={}, std::string source={})
Construct an error diagnostic with the standard error severity.
Definition Diagnostic.h:125
constexpr std::int64_t nanoseconds() const noexcept
Return the exact signed nanosecond representation.
Definition Time.h:69
static Result< Duration > fromSeconds(double seconds)
Convert finite seconds to the nearest nanosecond.
Definition Time.cpp:16
static constexpr Duration fromNanoseconds(std::int64_t nanoseconds) noexcept
Construct an exact duration from nanoseconds.
Definition Time.h:55
Move-only operation result carrying either a value or Status.
Definition Result.h:155
static Result success(T value)
Construct a successful result owning value.
Definition Result.h:164
static Result failure(Status status)
Construct a failed result from a structured status.
Definition Result.h:175
static Status success(StatusCode code=StatusCode::Ok)
Construct a successful status with an explicit non-error outcome.
Definition Status.h:81
The canonical owning dynamic value used by data-facing protocols.
Definition Value.h:31
std::map< std::string, Value > Object
Definition Value.h:34
const T * getIf() const noexcept
Return a typed pointer, or nullptr when the kind differs.
Definition Value.h:180
bool isObject() const noexcept
Return true when this value is an object.
Definition Value.h:99
Id128 public API.
Definition Identity.h:112
constexpr std::optional< StrongUint64 > incremented() const noexcept
Returns the next value, or empty instead of unsigned wraparound.
eve::Result< void > settle(std::string_view taskId, ProductionSettlementReceipt receipt)
Atomically marks a ready task settled and retains its idempotency receipt.
eve::Result< std::string > blockReason(std::string_view taskId) const
Returns the deterministic reason a task cannot start, or an empty string when runnable.
bool definitionAvailable(std::string_view owner, const eve::definition::DefinitionHandle &definition) const
Reports whether an exact definition incarnation is available to an owner.
bool tagAvailable(std::string_view owner, std::string_view tag) const
Reports whether an exact tag is available to an owner.
eve::OptionalRef< const ProductionTask > find(std::string_view taskId)
Finds a retained task by stable ID.
eve::Result< void > contribute(std::string_view taskId, WorkContribution contribution)
Applies deterministic external work after efficiency scaling.
int taskCount() const
Returns retained task count in enqueue order.
eve::Result< void > failSettlement(std::string_view taskId, std::string_view reason)
Records a retryable domain settlement failure without losing completed work.
eve::OptionalRef< const ProductionTask > ownerTaskAt(std::string_view owner, int index)
Returns an owner's retained task by enqueue index, or an empty borrowed reference.
SchedulerStrategy schedulerStrategy(std::string_view owner) const
Returns an owner's queue ordering strategy, defaulting to Priority.
eve::Result< void > resume(std::string_view taskId)
Returns a paused task to deterministic scheduling.
eve::Result< void > cancel(std::string_view taskId, std::string_view reason="cancelled")
Cancels a non-terminal task with a structured outcome.
int eventCount() const
Returns retained event count.
eve::Result< void > setSlotCount(std::string_view owner, int slots)
Sets an owner's parallel slot count; zero prevents new tasks from running.
eve::Result< void > retrySettlement(std::string_view taskId)
Returns a failed settlement to the ready state for an explicit retry.
int resourceCapacity(std::string_view owner, std::string_view resource) const
Returns named capacity, defaulting slot to slotCount and other resources to zero.
eve::OptionalRef< const ProductionTask > taskAt(int index)
Returns a retained task by enqueue index, or an empty borrowed reference.
eve::Result< void > setResourceCapacity(std::string_view owner, std::string_view resource, int capacity)
Sets named running capacity for an owner; zero disables that capability.
int runningCount(std::string_view owner) const
Returns the number of currently running tasks for an owner.
eve::Result< void > advance(const eve::SimulationStep &step)
Applies one injected deterministic simulation step.
eve::Result< std::string > enqueue(std::string_view owner, std::string_view kind, std::string_view product, eve::Value context, double duration, int priority=0)
Enqueues a task and returns its stable ID.
int ownerTaskCount(std::string_view owner) const
Returns retained task count for an owner.
WorkQueue(const WorkQueue &)=delete
int slotCount(std::string_view owner) const
Returns an owner's slot count, defaulting to one.
eve::Result< void > pause(std::string_view taskId)
Pauses a queued or running task, or returns NotFound/Conflict.
eve::Result< void > fail(std::string_view taskId, std::string_view reason="failed")
Marks a non-terminal task failed with a structured outcome.
eve::OptionalRef< ProductionEvent > eventAt(int index)
Returns an event by sequence index, or an empty borrowed reference.
void clearEvents()
Clears retained events without resetting sequence numbering.
TaskState
Lifecycle state of a production task.
Definition Production.h:30
ProductionEventKind
Kind of a deterministic production lifecycle event.
Definition Production.h:33
std::string_view eventKindName(ProductionEventKind kind)
Returns the stable lowercase name of an event kind.
std::string_view taskStateName(TaskState state)
Returns the stable lowercase name of a task state.
RefundPolicy
Reservation refund rule applied when work terminates before settlement.
Definition Production.h:74
DiagnosticCode
Stable machine-readable diagnostic codes.
Definition Diagnostic.h:47
std::optional< std::reference_wrapper< T > > OptionalRef
Optional borrowed reference; it does not extend T's lifetime.
Definition BorrowedRef.h:24
One deterministic fixed-step emitted by SimulationClock.
Definition Time.h:158
std::vector< std::string > prerequisites
Definition Production.h:61
std::vector< std::string > blocks
Definition Production.h:62
std::vector< std::string > requiredTags
Definition Production.h:64
std::vector< eve::definition::DefinitionHandle > requiredDefinitions
Definition Production.h:63
Validated input used to create a domain-neutral production task.
Definition Production.h:137
Owning proof that a domain consumer published one task exactly once.
Definition Production.h:46
A subject-agnostic continuous task retained for audit and save games.
Definition Production.h:158
ReservationState reservationState
Definition Production.h:169
eve::definition::DefinitionHandle definition
Definition Production.h:167
ProductionDependencies dependencies
Definition Production.h:173
ProductionRepeatPolicy repeat
Definition Production.h:177
ProductionTerminationPolicy termination
Definition Production.h:175
std::vector< ProductionResourceRequirement > resources
Definition Production.h:174
OutputRandomProvenance random
Definition Production.h:176
ProductionSettlementReceipt settlement
Definition Production.h:171
Deterministic externally supplied work contribution.
Definition Production.h:130