载入中...
搜索中...
未找到
ProductionPersistence.cpp
浏览该文件的文档.
2
3#include "common/Json.h"
4#include <algorithm>
5#include <exception>
6#include <functional>
7#include <iomanip>
8#include <set>
9#include <sstream>
10#include <utility>
11
12namespace eve::production {
13namespace {
14
15template <class T>
16eve::Result<T> persistenceFailure(eve::DiagnosticCode code, std::string message, std::string path = {}) {
18 eve::Diagnostic::error(code, message, path, {}, "production.persistence"));
19}
20
21std::string quote(std::string_view value) {
22 std::ostringstream out;
23 out << '"';
24 for (unsigned char c : value) {
25 switch (c) {
26 case '"': out << "\\\""; break;
27 case '\\': out << "\\\\"; break;
28 case '\n': out << "\\n"; break;
29 case '\r': out << "\\r"; break;
30 case '\t': out << "\\t"; break;
31 default:
32 if (c < 0x20)
33 out << "\\u" << std::hex << std::setw(4) << std::setfill('0') << int(c) << std::dec;
34 else
35 out << static_cast<char>(c);
36 }
37 }
38 return out.str() + '"';
39}
40
41std::string canonicalJson(const eve::json::Value& value) {
42 if (value.isNull()) return "null";
43 if (value.isBool()) return value.asBool() ? "true" : "false";
44 if (value.isNumber()) return value.asString();
45 if (value.isString()) return quote(value.asString());
46 if (value.isArray()) {
47 std::string out = "[";
48 for (size_t i = 0; i < value.size(); ++i) {
49 if (i) out += ',';
50 out += canonicalJson(value.at(i));
51 }
52 return out + ']';
53 }
54 auto keys = value.keys();
55 std::sort(keys.begin(), keys.end());
56 std::string out = "{";
57 for (size_t i = 0; i < keys.size(); ++i) {
58 if (i) out += ',';
59 out += quote(keys[i]) + ':' + canonicalJson(value.get(keys[i].c_str()));
60 }
61 return out + '}';
62}
63
64bool parseU64(const eve::json::Value& value, uint64_t& result) {
65 if (!value.isString()) return false;
66 try {
67 size_t used = 0;
68 result = std::stoull(value.asString(), &used);
69 return used == value.asString().size();
70 } catch (...) {
71 return false;
72 }
73}
74
75bool parseI64(const eve::json::Value& value, std::int64_t& result) {
76 if (!value.isString()) return false;
77 try {
78 size_t used = 0;
79 result = std::stoll(value.asString(), &used);
80 return used == value.asString().size();
81 } catch (...) {
82 return false;
83 }
84}
85
86bool parseDuration(const eve::json::Value& nanoseconds, const eve::json::Value& legacySeconds, eve::Duration& result) {
87 if (!nanoseconds.isNull()) {
88 std::int64_t value = 0;
89 if (!parseI64(nanoseconds, value)) return false;
91 return true;
92 }
93 if (!legacySeconds.isNumber()) return false;
94 auto converted = eve::Duration::fromSeconds(legacySeconds.asDouble());
95 if (!converted) return false;
96 result = std::move(converted).takeValue();
97 return true;
98}
99
100bool parseState(const std::string& name, TaskState& state) {
101 if (name == "queued")
103 else if (name == "running")
105 else if (name == "paused")
107 else if (name == "ready_to_settle")
109 else if (name == "settlement_failed")
111 else if (name == "completed")
113 else if (name == "cancelled")
115 else if (name == "failed")
117 else
118 return false;
119 return true;
120}
121
122bool parseEventKind(const std::string& name, ProductionEventKind& kind) {
123 if (name == "enqueued")
125 else if (name == "started")
127 else if (name == "paused")
129 else if (name == "resumed")
131 else if (name == "ready_to_settle")
133 else if (name == "settlement_failed")
135 else if (name == "completed")
137 else if (name == "cancelled")
139 else if (name == "failed")
141 else
142 return false;
143 return true;
144}
145
146} // namespace
147
149 std::ostringstream out;
150 const auto refundName = [](RefundPolicy policy) {
151 if (policy == RefundPolicy::None) return "none";
152 if (policy == RefundPolicy::Proportional) return "proportional";
153 return "full";
154 };
155 out << "{\"version\":7,\"revision\":" << quote(std::to_string(revision_.value()))
156 << ",\"nextEnqueueSequence\":" << quote(std::to_string(nextEnqueueSequence_))
157 << ",\"nextEventSequence\":" << quote(std::to_string(nextEventSequence_))
158 << ",\"nextTaskId\":" << quote(std::to_string(nextTaskId_))
159 << ",\"tick\":" << quote(std::to_string(tick_.value())) << ",\"slots\":[";
160 for (size_t i = 0; i < slots_.size(); ++i) {
161 if (i) out << ',';
162 out << "{\"owner\":" << quote(slots_[i].first) << ",\"value\":" << slots_[i].second << '}';
163 }
164 out << "],\"resources\":[";
165 for (size_t i = 0; i < resources_.size(); ++i) {
166 if (i) out << ',';
167 out << "{\"capacity\":" << resources_[i].second << ",\"owner\":" << quote(resources_[i].first.first)
168 << ",\"resource\":" << quote(resources_[i].first.second) << '}';
169 }
170 out << "],\"availableDefinitions\":[";
171 for (size_t i = 0; i < availableDefinitions_.size(); ++i) {
172 if (i) out << ',';
173 out << "{\"definition\":" << quote(std::get<1>(availableDefinitions_[i]))
174 << ",\"generation\":" << quote(std::to_string(std::get<2>(availableDefinitions_[i])))
175 << ",\"owner\":" << quote(std::get<0>(availableDefinitions_[i])) << '}';
176 }
177 out << "],\"availableTags\":[";
178 for (size_t i = 0; i < availableTags_.size(); ++i) {
179 if (i) out << ',';
180 out << "{\"owner\":" << quote(availableTags_[i].first)
181 << ",\"tag\":" << quote(availableTags_[i].second) << '}';
182 }
183 out << "],\"events\":[";
184 for (size_t i = 0; i < events_.size(); ++i) {
185 if (i) out << ',';
186 const auto& event = events_[i];
187 out << "{\"correlationId\":" << quote(event.correlationId)
188 << ",\"kind\":" << quote(eventKindName(event.kind)) << ",\"owner\":" << quote(event.owner)
189 << ",\"product\":" << quote(event.product) << ",\"reason\":" << quote(event.reason)
190 << ",\"sequence\":" << quote(std::to_string(event.sequence)) << ",\"taskId\":" << quote(event.taskId)
191 << ",\"taskKind\":" << quote(event.taskKind) << ",\"tick\":" << quote(std::to_string(event.tick.value()))
192 << '}';
193 }
194 out << "],\"schedulers\":[";
195 for (size_t i = 0; i < schedulerStrategies_.size(); ++i) {
196 if (i) out << ',';
197 const auto strategy = schedulerStrategies_[i].second;
198 const auto name = strategy == SchedulerStrategy::Fifo ? "fifo" :
199 strategy == SchedulerStrategy::ShortestRemaining ? "shortest_remaining" : "priority";
200 out << "{\"owner\":" << quote(schedulerStrategies_[i].first)
201 << ",\"strategy\":" << quote(name) << '}';
202 }
203 out << "],\"tasks\":[";
204 const auto writeIds = [&out](const std::vector<std::string>& ids) {
205 out << '[';
206 for (size_t index = 0; index < ids.size(); ++index) {
207 if (index) out << ',';
208 out << quote(ids[index]);
209 }
210 out << ']';
211 };
212 for (size_t i = 0; i < tasks_.size(); ++i) {
213 if (i) out << ',';
214 const auto& t = *tasks_[i];
215 auto context = t.context.toJson();
216 if (!context.ok()) return eve::Result<std::string>::failure(context.status());
217 auto reservation = t.reservation.toJson();
218 if (!reservation.ok()) return eve::Result<std::string>::failure(reservation.status());
219 auto releasePayload = t.reservationRelease.reservation.toJson();
220 if (!releasePayload.ok()) return eve::Result<std::string>::failure(releasePayload.status());
221 auto settlementPayload = t.settlement.payload.toJson();
222 if (!settlementPayload.ok()) return eve::Result<std::string>::failure(settlementPayload.status());
223 out << "{\"batchSize\":" << t.batchSize << ",\"blocks\":";
224 writeIds(t.dependencies.blocks);
225 out << ",\"completedCycles\":" << t.completedCycles
226 << ",\"continuous\":" << (t.repeat.continuous ? "true" : "false")
227 << ",\"context\":" << std::move(context).takeValue()
228 << ",\"correlationId\":" << quote(t.correlationId)
229 << ",\"definition\":" << quote(t.definition.isValid() ? t.definition.reference.id().format() : "")
230 << ",\"definitionGeneration\":" << quote(std::to_string(t.definition.generation.value()))
231 << ",\"dependencyMode\":"
232 << quote(t.dependencies.mode == TaskDependencyMode::AllOf ? "all_of" : "any_of")
233 << ",\"durationNs\":" << quote(std::to_string(t.duration.nanoseconds()))
234 << ",\"efficiencyPermille\":" << t.efficiencyPermille
235 << ",\"enqueueSequence\":" << quote(std::to_string(t.enqueueSequence)) << ",\"id\":" << quote(t.id)
236 << ",\"kind\":" << quote(t.kind) << ",\"owner\":" << quote(t.owner) << ",\"priority\":" << t.priority
237 << ",\"maintainStockTarget\":" << t.repeat.maintainStockTarget
238 << ",\"lastSettlementId\":" << quote(t.lastSettlementId)
239 << ",\"observedStock\":" << t.repeat.observedStock
240 << ",\"product\":" << quote(t.product) << ",\"prerequisites\":";
241 writeIds(t.dependencies.prerequisites);
242 out << ",\"randomDrawCount\":" << t.random.drawCount
243 << ",\"randomDrawStart\":" << quote(std::to_string(t.random.drawStart))
244 << ",\"randomSeed\":" << quote(std::to_string(t.random.seed))
245 << ",\"randomStream\":" << quote(t.random.stream)
246 << ",\"refundCancellation\":" << quote(refundName(t.termination.cancellation))
247 << ",\"refundFailure\":" << quote(refundName(t.termination.failure))
248 << ",\"refundPermille\":" << t.refundPermille << ",\"requiredDefinitions\":[";
249 for (size_t requirementIndex = 0; requirementIndex < t.dependencies.requiredDefinitions.size();
250 ++requirementIndex) {
251 if (requirementIndex) out << ',';
252 const auto& definition = t.dependencies.requiredDefinitions[requirementIndex];
253 out << "{\"definition\":" << quote(definition.reference.id().format())
254 << ",\"generation\":" << quote(std::to_string(definition.generation.value())) << '}';
255 }
256 out << "],\"requiredTags\":";
257 writeIds(t.dependencies.requiredTags);
258 out << ",\"requirements\":[";
259 for (size_t requirementIndex = 0; requirementIndex < t.resources.size(); ++requirementIndex) {
260 if (requirementIndex) out << ',';
261 out << "{\"resource\":" << quote(t.resources[requirementIndex].resource)
262 << ",\"units\":" << t.resources[requirementIndex].units << '}';
263 }
264 out
265 << ']'
266 << ",\"progressNs\":" << quote(std::to_string(t.progress.nanoseconds()))
267 << ",\"workRemainderPermille\":" << t.workRemainderPermille
268 << ",\"reason\":" << quote(t.reason) << ",\"reservation\":"
269 << std::move(reservation).takeValue()
270 << ",\"reservationReleaseId\":" << quote(t.reservationRelease.releaseId)
271 << ",\"reservationReleasePayload\":" << std::move(releasePayload).takeValue()
272 << ",\"reservationReleaseRefundPermille\":" << t.reservationRelease.refundPermille
273 << ",\"reservationState\":" << quote(
274 t.reservationState == ReservationState::Reserved ? "reserved" :
275 t.reservationState == ReservationState::Started ? "started" :
276 t.reservationState == ReservationState::Consumed ? "consumed" :
277 t.reservationState == ReservationState::Released ? "released" : "none")
278 << ",\"settlementId\":" << quote(t.settlement.settlementId)
279 << ",\"settlementPayload\":" << std::move(settlementPayload).takeValue()
280 << ",\"settlementRequired\":" << (t.settlementRequired ? "true" : "false")
281 << ",\"state\":" << quote(taskStateName(t.state))
282 << ",\"totalCycles\":" << t.repeat.totalCycles << '}';
283 }
284 return eve::Result<std::string>::success(out.str() + "]}");
285}
286
287eve::Result<void> WorkQueue::restore(std::string_view json) {
288 std::string error;
289 auto doc = eve::json::Document::parse(std::string(json), &error);
290 if (!doc.valid() || !doc.root().isObject()) {
291 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
292 error.empty() ? "snapshot must be an object" : error);
293 }
294 WorkQueue candidate(instanceId_);
295 const auto root = doc.root();
296 const int version = root.getInt("version");
297 if ((version != 1 && version != 2 && version != 3 && version != 4 && version != 5 && version != 6 &&
298 version != 7) ||
299 !parseU64(root.get("nextTaskId"), candidate.nextTaskId_) ||
300 !parseU64(root.get("nextEnqueueSequence"), candidate.nextEnqueueSequence_) ||
301 !parseU64(root.get("nextEventSequence"), candidate.nextEventSequence_) || candidate.nextTaskId_ == 0 ||
302 candidate.nextEnqueueSequence_ == 0 || candidate.nextEventSequence_ == 0 || !root.get("slots").isArray() ||
303 !root.get("events").isArray() || !root.get("tasks").isArray() ||
304 (version >= 4 && !root.get("resources").isArray()) ||
305 (version >= 5 && !root.get("schedulers").isArray()) ||
306 (version >= 6 && (!root.get("availableDefinitions").isArray() || !root.get("availableTags").isArray()))) {
307 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid snapshot counters or arrays");
308 }
309 if (!root.get("tick").isNull()) {
310 uint64_t tick = 0;
311 if (!parseU64(root.get("tick"), tick)) {
312 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid snapshot tick", "tick");
313 }
314 candidate.tick_ = eve::SimulationTick(tick);
315 }
316 if (!root.get("revision").isNull()) {
317 uint64_t revision = 0;
318 if (!parseU64(root.get("revision"), revision)) {
319 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid snapshot revision",
320 "revision");
321 }
322 candidate.revision_ = eve::Revision(revision);
323 } else {
324 candidate.revision_ = eve::Revision(candidate.nextEventSequence_ - 1);
325 }
326 for (size_t i = 0; i < root.get("slots").size(); ++i) {
327 auto value = root.get("slots").at(i);
328 if (!value.isObject() || !value.get("owner").isString() || !value.get("value").isNumber()) {
329 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid slot entry", "slots");
330 }
331 int slots = value.get("value").asInt();
332 auto setSlots = candidate.setSlotCount(value.get("owner").asString(), slots);
333 if (!setSlots.ok()) return eve::Result<void>::failure(setSlots.status());
334 }
335 if (version >= 4) {
336 for (size_t i = 0; i < root.get("resources").size(); ++i) {
337 const auto value = root.get("resources").at(i);
338 if (!value.isObject() || !value.get("owner").isString() || !value.get("resource").isString() ||
339 !value.get("capacity").isNumber())
340 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
341 "invalid resource capacity", "resources");
342 auto result = candidate.setResourceCapacity(value.get("owner").asString(),
343 value.get("resource").asString(),
344 value.get("capacity").asInt());
345 if (!result) return eve::Result<void>::failure(result.status());
346 }
347 }
348 if (version >= 5) {
349 for (size_t i = 0; i < root.get("schedulers").size(); ++i) {
350 const auto value = root.get("schedulers").at(i);
351 if (!value.isObject() || !value.get("owner").isString() || !value.get("strategy").isString())
352 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
353 "invalid scheduler strategy", "schedulers");
354 SchedulerStrategy strategy;
355 const auto name = value.get("strategy").asString();
356 if (name == "priority") strategy = SchedulerStrategy::Priority;
357 else if (name == "fifo") strategy = SchedulerStrategy::Fifo;
358 else if (name == "shortest_remaining") strategy = SchedulerStrategy::ShortestRemaining;
359 else return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
360 "unknown scheduler strategy", "schedulers.strategy");
361 auto result = candidate.setSchedulerStrategy(value.get("owner").asString(), strategy);
362 if (!result) return eve::Result<void>::failure(result.status());
363 }
364 }
365 if (version >= 6) {
366 for (size_t i = 0; i < root.get("availableDefinitions").size(); ++i) {
367 const auto value = root.get("availableDefinitions").at(i);
368 uint64_t generation = 0;
369 if (!value.isObject() || !value.get("owner").isString() ||
370 !value.get("definition").isString() || !parseU64(value.get("generation"), generation))
371 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
372 "invalid available definition", "availableDefinitions");
373 auto reference = eve::DefinitionRef::parse(value.get("definition").asString());
374 if (!reference || generation == 0)
375 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
376 "invalid available definition handle", "availableDefinitions");
378 std::move(reference).takeValue(), eve::Generation(generation)};
379 auto result = candidate.setDefinitionAvailable(value.get("owner").asString(), handle, true);
380 if (!result) return eve::Result<void>::failure(result.status());
381 }
382 for (size_t i = 0; i < root.get("availableTags").size(); ++i) {
383 const auto value = root.get("availableTags").at(i);
384 if (!value.isObject() || !value.get("owner").isString() || !value.get("tag").isString())
385 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
386 "invalid available tag", "availableTags");
387 auto result = candidate.setTagAvailable(value.get("owner").asString(), value.get("tag").asString(), true);
388 if (!result) return eve::Result<void>::failure(result.status());
389 }
390 }
391 std::set<std::string> ids;
392 for (size_t i = 0; i < root.get("tasks").size(); ++i) {
393 auto value = root.get("tasks").at(i);
394 auto task = std::make_unique<ProductionTask>();
395 uint64_t sequence = 0;
397 if (!value.isObject() || !value.get("id").isString() || !value.get("owner").isString() ||
398 !value.get("kind").isString() || !value.get("product").isString() ||
399 (!value.get("progressNs").isString() && !value.get("progress").isNumber()) ||
400 (!value.get("durationNs").isString() && !value.get("duration").isNumber()) ||
401 !value.get("priority").isNumber() || !value.get("state").isString() || !value.get("reason").isString() ||
402 !parseU64(value.get("enqueueSequence"), sequence) || !parseState(value.get("state").asString(), state) ||
403 value.get("id").asString().empty() || value.get("owner").asString().empty() ||
404 value.get("kind").asString().empty() || value.get("product").asString().empty() ||
405 !ids.insert(value.get("id").asString()).second) {
406 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid task entry", "tasks");
407 }
408 task->id = value.get("id").asString();
409 task->owner = value.get("owner").asString();
410 task->kind = value.get("kind").asString();
411 task->product = value.get("product").asString();
412 auto context = eve::Value::fromJson(canonicalJson(value.get("context")));
413 if (!context.ok()) return eve::Result<void>::failure(context.status());
414 if (!context.value().isObject())
415 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
416 "work task context must be an object", "tasks.context");
417 task->context = std::move(context).takeValue();
418 if (version >= 2) {
419 if (!value.get("definition").isString() || !value.get("definitionGeneration").isString() ||
420 !value.get("settlementId").isString() || !value.get("settlementRequired").isBool())
421 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
422 "invalid task production metadata", "tasks");
423 uint64_t generation = 0;
424 if (!parseU64(value.get("definitionGeneration"), generation))
425 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
426 "invalid definition generation", "tasks.definitionGeneration");
427 const std::string definition = value.get("definition").asString();
428 if (!definition.empty()) {
429 auto reference = eve::DefinitionRef::parse(definition);
430 if (!reference || generation == 0)
431 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
432 "invalid pinned definition", "tasks.definition");
433 task->definition = {std::move(reference).takeValue(), eve::Generation(generation)};
434 } else if (generation != 0) {
435 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
436 "definition generation has no definition", "tasks.definition");
437 }
438 auto reservation = eve::Value::fromJson(canonicalJson(value.get("reservation")));
439 auto settlementPayload = eve::Value::fromJson(canonicalJson(value.get("settlementPayload")));
440 if (!reservation || !settlementPayload || !reservation.value().isObject() ||
441 !settlementPayload.value().isObject())
442 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
443 "invalid reservation or settlement payload", "tasks");
444 task->reservation = std::move(reservation).takeValue();
445 task->settlement.settlementId = value.get("settlementId").asString();
446 task->settlement.payload = std::move(settlementPayload).takeValue();
447 task->settlementRequired = value.get("settlementRequired").asBool();
448 if (state == TaskState::Completed && task->settlement.settlementId.empty())
449 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
450 "settled task has no receipt", "tasks.settlementId");
451 } else if (state == TaskState::Completed) {
452 task->settlement.settlementId = "legacy:" + task->id;
453 task->settlementRequired = false;
454 }
455 if (version >= 3) {
456 if (!value.get("dependencyMode").isString() || !value.get("prerequisites").isArray() ||
457 !value.get("blocks").isArray())
458 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
459 "invalid production dependency metadata", "tasks");
460 const std::string mode = value.get("dependencyMode").asString();
461 if (mode == "all_of")
462 task->dependencies.mode = TaskDependencyMode::AllOf;
463 else if (mode == "any_of")
464 task->dependencies.mode = TaskDependencyMode::AnyOf;
465 else
466 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
467 "invalid production dependency mode", "tasks.dependencyMode");
468 const auto parseIds = [](const eve::json::Value& source,
469 std::vector<std::string>& target) -> bool {
470 std::set<std::string> unique;
471 for (size_t index = 0; index < source.size(); ++index) {
472 const auto entry = source.at(index);
473 if (!entry.isString() || entry.asString().empty() ||
474 !unique.insert(entry.asString()).second)
475 return false;
476 target.push_back(entry.asString());
477 }
478 return true;
479 };
480 if (!parseIds(value.get("prerequisites"), task->dependencies.prerequisites) ||
481 !parseIds(value.get("blocks"), task->dependencies.blocks))
482 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
483 "invalid production dependency ids", "tasks");
484 }
485 if (version >= 4) {
486 if (!value.get("batchSize").isNumber() || !value.get("efficiencyPermille").isNumber() ||
487 !value.get("refundCancellation").isString() || !value.get("refundFailure").isString() ||
488 !value.get("refundPermille").isNumber() || !value.get("requirements").isArray())
489 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
490 "invalid production execution metadata", "tasks");
491 const auto parseRefund = [](std::string_view text, RefundPolicy& output) {
492 if (text == "none") output = RefundPolicy::None;
493 else if (text == "full") output = RefundPolicy::Full;
494 else if (text == "proportional") output = RefundPolicy::Proportional;
495 else return false;
496 return true;
497 };
498 const int batchSize = value.get("batchSize").asInt();
499 const int efficiency = value.get("efficiencyPermille").asInt();
500 const int refund = value.get("refundPermille").asInt();
501 if (batchSize <= 0 || efficiency <= 0 || refund < 0 || refund > 1000 ||
502 !parseRefund(value.get("refundCancellation").asString(), task->termination.cancellation) ||
503 !parseRefund(value.get("refundFailure").asString(), task->termination.failure))
504 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
505 "invalid production batch, efficiency or refund", "tasks");
506 task->batchSize = static_cast<std::uint32_t>(batchSize);
507 task->efficiencyPermille = static_cast<std::uint32_t>(efficiency);
508 task->refundPermille = static_cast<std::uint32_t>(refund);
509 std::set<std::string> resourceNames;
510 for (size_t index = 0; index < value.get("requirements").size(); ++index) {
511 const auto requirement = value.get("requirements").at(index);
512 if (!requirement.isObject() || !requirement.get("resource").isString() ||
513 !requirement.get("units").isNumber() || requirement.get("resource").asString().empty() ||
514 requirement.get("units").asInt() <= 0 ||
515 !resourceNames.insert(requirement.get("resource").asString()).second)
516 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
517 "invalid task resource requirement", "tasks.requirements");
518 task->resources.push_back(
519 {requirement.get("resource").asString(), requirement.get("units").asInt()});
520 }
521 }
522 if (version >= 5) {
523 if (!value.get("completedCycles").isNumber() || !value.get("continuous").isBool() ||
524 !value.get("correlationId").isString() || !value.get("lastSettlementId").isString() ||
525 !value.get("maintainStockTarget").isNumber() ||
526 !value.get("observedStock").isNumber() || !value.get("randomDrawCount").isNumber() ||
527 !value.get("randomDrawStart").isString() || !value.get("randomSeed").isString() ||
528 !value.get("randomStream").isString() || !value.get("totalCycles").isNumber())
529 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
530 "invalid advanced production metadata", "tasks");
531 uint64_t randomSeed = 0;
532 uint64_t randomDrawStart = 0;
533 const int completedCycles = value.get("completedCycles").asInt();
534 const int drawCount = value.get("randomDrawCount").asInt();
535 const int totalCycles = value.get("totalCycles").asInt();
536 const int stockTarget = value.get("maintainStockTarget").asInt();
537 const int observedStock = value.get("observedStock").asInt();
538 if (completedCycles < 0 || drawCount < 0 || totalCycles <= 0 || stockTarget < -1 || observedStock < 0 ||
539 !parseU64(value.get("randomSeed"), randomSeed) ||
540 !parseU64(value.get("randomDrawStart"), randomDrawStart) ||
541 (drawCount > 0 && value.get("randomStream").asString().empty()))
542 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
543 "invalid advanced production values", "tasks");
544 task->completedCycles = static_cast<std::uint32_t>(completedCycles);
545 task->repeat.continuous = value.get("continuous").asBool();
546 task->repeat.totalCycles = static_cast<std::uint32_t>(totalCycles);
547 task->repeat.maintainStockTarget = stockTarget;
548 task->repeat.observedStock = observedStock;
549 task->correlationId = value.get("correlationId").asString();
550 task->lastSettlementId = value.get("lastSettlementId").asString();
551 task->random = {value.get("randomStream").asString(), randomSeed, randomDrawStart,
552 static_cast<std::uint32_t>(drawCount)};
553 } else {
554 task->correlationId = task->id;
555 }
556 if (version >= 6) {
557 if (!value.get("requiredDefinitions").isArray() || !value.get("requiredTags").isArray() ||
558 !value.get("reservationReleaseId").isString() ||
559 !value.get("reservationReleaseRefundPermille").isNumber() ||
560 !value.get("reservationState").isString())
561 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
562 "invalid definition or tag requirements", "tasks");
563 std::set<std::pair<std::string, std::uint64_t>> definitions;
564 for (size_t index = 0; index < value.get("requiredDefinitions").size(); ++index) {
565 const auto requirement = value.get("requiredDefinitions").at(index);
566 uint64_t generation = 0;
567 if (!requirement.isObject() || !requirement.get("definition").isString() ||
568 !parseU64(requirement.get("generation"), generation))
569 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
570 "invalid required definition", "tasks.requiredDefinitions");
571 const auto name = requirement.get("definition").asString();
573 if (!reference || generation == 0 || !definitions.insert({name, generation}).second)
574 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
575 "invalid required definition handle", "tasks.requiredDefinitions");
576 task->dependencies.requiredDefinitions.push_back(
577 {std::move(reference).takeValue(), eve::Generation(generation)});
578 }
579 std::set<std::string> tags;
580 for (size_t index = 0; index < value.get("requiredTags").size(); ++index) {
581 const auto requirement = value.get("requiredTags").at(index);
582 if (!requirement.isString() || requirement.asString().empty() ||
583 !tags.insert(requirement.asString()).second)
584 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
585 "invalid required tag", "tasks.requiredTags");
586 task->dependencies.requiredTags.push_back(requirement.asString());
587 }
588 const auto reservationState = value.get("reservationState").asString();
589 if (reservationState == "none") task->reservationState = ReservationState::None;
590 else if (reservationState == "reserved") task->reservationState = ReservationState::Reserved;
591 else if (reservationState == "started") task->reservationState = ReservationState::Started;
592 else if (reservationState == "consumed") task->reservationState = ReservationState::Consumed;
593 else if (reservationState == "released") task->reservationState = ReservationState::Released;
594 else return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
595 "invalid reservation state", "tasks.reservationState");
596 auto releasePayload = eve::Value::fromJson(canonicalJson(value.get("reservationReleasePayload")));
597 const int releaseRefund = value.get("reservationReleaseRefundPermille").asInt();
598 if (!releasePayload || !releasePayload.value().isObject() || releaseRefund < 0 || releaseRefund > 1000)
599 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
600 "invalid reservation release", "tasks.reservationRelease");
601 task->reservationRelease = {value.get("reservationReleaseId").asString(),
602 std::move(releasePayload).takeValue(),
603 static_cast<std::uint32_t>(releaseRefund)};
604 if ((task->reservationState == ReservationState::Released) !=
605 !task->reservationRelease.releaseId.empty())
606 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
607 "reservation release state and receipt disagree",
608 "tasks.reservationReleaseId");
609 }
610 if (version >= 7) {
611 if (!value.get("workRemainderPermille").isNumber())
612 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
613 "invalid production work remainder", "tasks.workRemainderPermille");
614 const int remainder = value.get("workRemainderPermille").asInt();
615 if (remainder < 0 || remainder >= 1000)
616 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
617 "production work remainder must be between 0 and 999",
618 "tasks.workRemainderPermille");
619 task->workRemainderPermille = static_cast<std::uint32_t>(remainder);
620 }
621 if (!parseDuration(value.get("durationNs"), value.get("duration"), task->duration) ||
622 !parseDuration(value.get("progressNs"), value.get("progress"), task->progress)) {
623 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid task duration",
624 "tasks.duration");
625 }
626 task->priority = value.get("priority").asInt();
627 task->state = state;
628 task->enqueueSequence = sequence;
629 task->reason = value.get("reason").asString();
630 if (task->duration.nanoseconds() <= 0 || task->progress.nanoseconds() < 0 || task->progress > task->duration) {
631 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid task progress",
632 "tasks.progress");
633 }
634 candidate.tasks_.push_back(std::move(task));
635 }
636 for (const auto& task : candidate.tasks_) {
637 for (const auto& dependencyId : task->dependencies.prerequisites) {
638 if (dependencyId == task->id || !ids.contains(dependencyId))
639 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
640 "invalid prerequisite reference", "tasks.prerequisites");
641 }
642 for (const auto& blockerId : task->dependencies.blocks) {
643 if (blockerId == task->id || !ids.contains(blockerId))
644 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
645 "invalid blocker reference", "tasks.blocks");
646 }
647 }
648 std::set<std::string> visiting;
649 std::set<std::string> visited;
650 std::function<bool(const ProductionTask&)> acyclic = [&](const ProductionTask& task) {
651 if (visited.contains(task.id)) return true;
652 if (!visiting.insert(task.id).second) return false;
653 for (const auto& dependencyId : task.dependencies.prerequisites) {
654 const auto dependency = candidate.find(dependencyId);
655 if (!dependency || !acyclic(dependency->get())) return false;
656 }
657 visiting.erase(task.id);
658 visited.insert(task.id);
659 return true;
660 };
661 for (const auto& task : candidate.tasks_)
662 if (!acyclic(*task))
663 return persistenceFailure<void>(eve::DiagnosticCode::ParseError,
664 "production dependency graph contains a cycle",
665 "tasks.prerequisites");
666 uint64_t previousEventSequence = 0;
667 for (size_t i = 0; i < root.get("events").size(); ++i) {
668 auto value = root.get("events").at(i);
669 ProductionEvent event;
670 if (!value.isObject() || !value.get("kind").isString() || !value.get("owner").isString() ||
671 !value.get("product").isString() || !value.get("reason").isString() || !value.get("taskId").isString() ||
672 !value.get("taskKind").isString() || !parseU64(value.get("sequence"), event.sequence) ||
673 !parseEventKind(value.get("kind").asString(), event.kind) || event.sequence <= previousEventSequence ||
674 event.sequence >= candidate.nextEventSequence_) {
675 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid event entry", "events");
676 }
677 previousEventSequence = event.sequence;
678 event.owner = value.get("owner").asString();
679 event.product = value.get("product").asString();
680 event.reason = value.get("reason").asString();
681 event.taskId = value.get("taskId").asString();
682 event.taskKind = value.get("taskKind").asString();
683 event.correlationId = version >= 5 && value.get("correlationId").isString()
684 ? value.get("correlationId").asString() : event.taskId;
685 if (!value.get("tick").isNull()) {
686 uint64_t tick = 0;
687 if (!parseU64(value.get("tick"), tick)) {
688 return persistenceFailure<void>(eve::DiagnosticCode::ParseError, "invalid event tick",
689 "events.tick");
690 }
691 event.tick = eve::SimulationTick(tick);
692 }
693 candidate.events_.push_back(std::move(event));
694 }
695 *this = std::move(candidate);
697}
698
700 tasks_.clear();
701 events_.clear();
702 slots_.clear();
703 resources_.clear();
704 schedulerStrategies_.clear();
705 availableDefinitions_.clear();
706 availableTags_.clear();
707 nextTaskId_ = nextEnqueueSequence_ = nextEventSequence_ = 1;
708 revision_ = eve::Revision::zero();
710}
711
712} // namespace eve::production
LogicalId target
double value
int root
Definition AnimSmr.cpp:119
std::string output
AuthorityStoreHandleRef reference
Definition Authority.cpp:24
std::string message
DiagnosticCode code
glm::uvec4 ids
std::int32_t second
std::int32_t c
std::int32_t first
std::string text
TokenKind kind
char quote
std::string name
std::unique_ptr< gpgpu::Sequence > sequence
Definition OnnxGpgpu.cpp:43
std::string error
Definition Package.cpp:60
std::uint64_t revision
std::string path
Definition PlayHost.cpp:110
std::uint32_t generation
PrimitiveHandle handle
float t
LocalPageCacheEntry slots[ShadowConfig::kLocalSlots]
SimulationTick tick
float size
Definition TreeMesh.cpp:156
std::set< std::string > visiting
uint32_t index
const UnitySourceAsset & source
const VegetationPresetContext & context
const std::map< std::string, Value > & definitions
static Result< DefinitionRef > parse(std::string_view text)
Parse a namespace:name logical definition identifier.
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
Signed, fixed-resolution simulation duration in nanoseconds.
Definition Time.h:43
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
static Result< Value > fromJson(std::string_view json)
Parse one strict JSON value into an owning Value.
Definition Value.cpp:57
constexpr std::uint64_t value() const noexcept
Returns the underlying value at an explicit protocol boundary.
static constexpr StrongUint64 zero() noexcept
Returns the zero value for this strong type.
static Document parse(const std::string &text, std::string *error=nullptr)
Parse.
Definition Json.cpp:581
EVENGINE_API_FOUNDATION public API.
Definition Json.h:38
double asDouble(double fallback=0.0) const
As double.
Definition Json.cpp:377
bool isNull() const
True when null.
Definition Json.cpp:325
bool isNumber() const
True when number.
Definition Json.cpp:327
Generic multi-owner, multi-slot continuous-progress work queue.
Definition Production.h:200
eve::OptionalRef< const ProductionTask > find(std::string_view taskId)
Finds a retained task by stable ID.
eve::Result< std::string > snapshot() const
Serializes the complete queue as deterministic JSON.
eve::Result< void > setSchedulerStrategy(std::string_view owner, SchedulerStrategy strategy)
Selects an owner's deterministic queue ordering strategy.
eve::Result< void > setDefinitionAvailable(std::string_view owner, const eve::definition::DefinitionHandle &definition, bool available)
Publishes or withdraws one generation-qualified definition fact for an owner.
eve::Result< void > setTagAvailable(std::string_view owner, std::string_view tag, bool available)
Publishes or withdraws one exact tag fact for an owner.
void clear()
Clears tasks, owner settings, events, and stable counters.
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 > restore(std::string_view json)
Transactionally restores a snapshot; failure preserves current state.
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.
TaskState
Lifecycle state of a production task.
Definition Production.h:30
SchedulerStrategy
Stable owner-local queue ordering strategy.
Definition Production.h:87
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
detail::StrongUint64< detail::GenerationTag > Generation
Registry/object replacement generation used to reject stale handles.
Definition Generation.h:13
detail::StrongUint64< detail::SimulationTickTag > SimulationTick
Deterministic simulation time step; it is not wall-clock time.
Definition Time.h:31
detail::StrongUint64< detail::RevisionTag > Revision
Monotonic content/state revision used for optimistic-concurrency checks.
Definition Revision.h:13
A generation-qualified reference to one definition incarnation.
Deterministically sequenced production lifecycle event.
Definition Production.h:189
A subject-agnostic continuous task retained for audit and save games.
Definition Production.h:158