载入中...
搜索中...
未找到
ProductionSnapshot.cpp
浏览该文件的文档.
2
3#include <exception>
4#include <utility>
5
6namespace eve::production {
7namespace {
8
9eve::LogicalId productionSchema() {
10 const auto schema = eve::LogicalId::parse("production:queue");
11 if (!schema) std::terminate();
12 return *schema;
13}
14
15const eve::SnapshotMigrationChain& productionMigrations() {
16 static const eve::SnapshotMigrationChain chain = [] {
18 const auto registration =
19 result.add(productionSchema(), eve::SchemaVersion(0), eve::SchemaVersion(1),
21 const auto* object = payload.getIf<eve::Value::Object>();
22 if (!object)
24 eve::DiagnosticCode::ParseError, "production queue payload must be an object"));
25 eve::Value::Object migrated = *object;
26 if (!migrated.contains("version")) migrated.emplace("version", eve::Value(std::int64_t(1)));
27 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
28 });
29 if (!registration.ok()) std::terminate();
30 const auto settlementRegistration =
31 result.add(productionSchema(), eve::SchemaVersion(1), eve::SchemaVersion(2),
33 const auto* object = payload.getIf<eve::Value::Object>();
34 if (!object)
36 eve::DiagnosticCode::ParseError, "production queue payload must be an object"));
37 eve::Value::Object migrated = *object;
38 migrated["version"] = eve::Value(std::int64_t(2));
39 if (auto tasks = migrated.find("tasks"); tasks != migrated.end()) {
40 if (auto* array = tasks->second.getIf<eve::Value::Array>()) {
41 for (auto& item : *array) {
42 auto* task = item.getIf<eve::Value::Object>();
43 if (!task) continue;
44 task->try_emplace("definition", eve::Value(std::string{}));
45 task->try_emplace("definitionGeneration", eve::Value(std::string("0")));
46 task->try_emplace("reservation", eve::Value(eve::Value::Object{}));
47 const auto state = task->find("state");
48 const bool completed = state != task->end() &&
49 state->second.getIf<std::string>() != nullptr &&
50 *state->second.getIf<std::string>() == "completed";
51 const auto id = task->find("id");
52 const auto* taskId = id == task->end() ? nullptr : id->second.getIf<std::string>();
53 task->try_emplace("settlementId", eve::Value(
54 completed && taskId ? "legacy:" + *taskId : std::string{}));
55 task->try_emplace("settlementPayload", eve::Value(eve::Value::Object{}));
56 task->try_emplace("settlementRequired", eve::Value(false));
57 }
58 }
59 }
60 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
61 });
62 if (!settlementRegistration.ok()) std::terminate();
63 const auto dependencyRegistration =
64 result.add(productionSchema(), eve::SchemaVersion(2), eve::SchemaVersion(3),
66 const auto* object = payload.getIf<eve::Value::Object>();
67 if (!object)
69 eve::DiagnosticCode::ParseError, "production queue payload must be an object"));
70 eve::Value::Object migrated = *object;
71 migrated["version"] = eve::Value(std::int64_t(3));
72 if (auto tasks = migrated.find("tasks"); tasks != migrated.end()) {
73 if (auto* array = tasks->second.getIf<eve::Value::Array>()) {
74 for (auto& item : *array) {
75 auto* task = item.getIf<eve::Value::Object>();
76 if (!task) continue;
77 task->try_emplace("dependencyMode", eve::Value("all_of"));
78 task->try_emplace("prerequisites", eve::Value(eve::Value::Array{}));
79 task->try_emplace("blocks", eve::Value(eve::Value::Array{}));
80 }
81 }
82 }
83 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
84 });
85 if (!dependencyRegistration.ok()) std::terminate();
86 const auto executionRegistration =
87 result.add(productionSchema(), eve::SchemaVersion(3), eve::SchemaVersion(4),
89 const auto* object = payload.getIf<eve::Value::Object>();
90 if (!object)
92 eve::DiagnosticCode::ParseError, "production queue payload must be an object"));
93 eve::Value::Object migrated = *object;
94 migrated["version"] = eve::Value(std::int64_t(4));
95 migrated.try_emplace("resources", eve::Value(eve::Value::Array{}));
96 if (auto tasks = migrated.find("tasks"); tasks != migrated.end()) {
97 if (auto* array = tasks->second.getIf<eve::Value::Array>()) {
98 for (auto& item : *array) {
99 auto* task = item.getIf<eve::Value::Object>();
100 if (!task) continue;
101 task->try_emplace("batchSize", eve::Value(std::int64_t(1)));
102 task->try_emplace("efficiencyPermille", eve::Value(std::int64_t(1000)));
103 task->try_emplace("refundCancellation", eve::Value("full"));
104 task->try_emplace("refundFailure", eve::Value("full"));
105 task->try_emplace("refundPermille", eve::Value(std::int64_t(0)));
106 task->try_emplace("requirements", eve::Value(eve::Value::Array{}));
107 }
108 }
109 }
110 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
111 });
112 if (!executionRegistration.ok()) std::terminate();
113 const auto advancedRegistration =
114 result.add(productionSchema(), eve::SchemaVersion(4), eve::SchemaVersion(5),
116 const auto* object = payload.getIf<eve::Value::Object>();
117 if (!object)
119 eve::DiagnosticCode::ParseError, "production queue payload must be an object"));
120 eve::Value::Object migrated = *object;
121 migrated["version"] = eve::Value(std::int64_t(5));
122 migrated.try_emplace("schedulers", eve::Value(eve::Value::Array{}));
123 if (auto events = migrated.find("events"); events != migrated.end()) {
124 if (auto* array = events->second.getIf<eve::Value::Array>()) {
125 for (auto& item : *array) {
126 auto* event = item.getIf<eve::Value::Object>();
127 if (!event) continue;
128 const auto taskId = event->find("taskId");
129 const auto* id = taskId == event->end() ? nullptr : taskId->second.getIf<std::string>();
130 event->try_emplace("correlationId", eve::Value(id ? *id : std::string{}));
131 }
132 }
133 }
134 if (auto tasks = migrated.find("tasks"); tasks != migrated.end()) {
135 if (auto* array = tasks->second.getIf<eve::Value::Array>()) {
136 for (auto& item : *array) {
137 auto* task = item.getIf<eve::Value::Object>();
138 if (!task) continue;
139 const auto taskId = task->find("id");
140 const auto* id = taskId == task->end() ? nullptr : taskId->second.getIf<std::string>();
141 task->try_emplace("completedCycles", eve::Value(std::int64_t(0)));
142 task->try_emplace("continuous", eve::Value(false));
143 task->try_emplace("correlationId", eve::Value(id ? *id : std::string{}));
144 task->try_emplace("lastSettlementId", eve::Value(std::string{}));
145 task->try_emplace("maintainStockTarget", eve::Value(std::int64_t(-1)));
146 task->try_emplace("observedStock", eve::Value(std::int64_t(0)));
147 task->try_emplace("randomDrawCount", eve::Value(std::int64_t(0)));
148 task->try_emplace("randomDrawStart", eve::Value(std::string("0")));
149 task->try_emplace("randomSeed", eve::Value(std::string("0")));
150 task->try_emplace("randomStream", eve::Value(std::string{}));
151 task->try_emplace("totalCycles", eve::Value(std::int64_t(1)));
152 }
153 }
154 }
155 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
156 });
157 if (!advancedRegistration.ok()) std::terminate();
158 const auto requirementRegistration =
159 result.add(productionSchema(), eve::SchemaVersion(5), eve::SchemaVersion(6),
161 const auto* object = payload.getIf<eve::Value::Object>();
162 if (!object)
164 eve::DiagnosticCode::ParseError, "production queue payload must be an object"));
165 eve::Value::Object migrated = *object;
166 migrated["version"] = eve::Value(std::int64_t(6));
167 migrated.try_emplace("availableDefinitions", eve::Value(eve::Value::Array{}));
168 migrated.try_emplace("availableTags", eve::Value(eve::Value::Array{}));
169 if (auto tasks = migrated.find("tasks"); tasks != migrated.end()) {
170 if (auto* array = tasks->second.getIf<eve::Value::Array>()) {
171 for (auto& item : *array) {
172 auto* task = item.getIf<eve::Value::Object>();
173 if (!task) continue;
174 task->try_emplace("requiredDefinitions", eve::Value(eve::Value::Array{}));
175 task->try_emplace("requiredTags", eve::Value(eve::Value::Array{}));
176 const auto reservation = task->find("reservation");
177 const auto* reservationObject = reservation == task->end()
178 ? nullptr : reservation->second.getIf<eve::Value::Object>();
179 const bool reserved = reservationObject != nullptr && !reservationObject->empty();
180 std::string reservationState = "none";
181 if (reserved) {
182 const auto state = task->find("state");
183 const auto* stateName = state == task->end()
184 ? nullptr : state->second.getIf<std::string>();
185 reservationState = stateName != nullptr && *stateName == "queued"
186 ? "reserved" : stateName != nullptr &&
187 (*stateName == "ready_to_settle" ||
188 *stateName == "settlement_failed" || *stateName == "completed")
189 ? "consumed" : "started";
190 }
191 task->try_emplace("reservationReleaseId", eve::Value(std::string{}));
192 task->try_emplace("reservationReleasePayload",
194 task->try_emplace("reservationReleaseRefundPermille",
195 eve::Value(std::int64_t(0)));
196 task->try_emplace("reservationState", eve::Value(reservationState));
197 }
198 }
199 }
200 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
201 });
202 if (!requirementRegistration.ok()) std::terminate();
203 const auto remainderRegistration =
204 result.add(productionSchema(), eve::SchemaVersion(6), eve::SchemaVersion(7),
206 const auto* object = payload.getIf<eve::Value::Object>();
207 if (!object)
209 eve::DiagnosticCode::ParseError, "production queue payload must be an object"));
210 eve::Value::Object migrated = *object;
211 migrated["version"] = eve::Value(std::int64_t(7));
212 if (auto tasks = migrated.find("tasks"); tasks != migrated.end()) {
213 if (auto* array = tasks->second.getIf<eve::Value::Array>()) {
214 for (auto& item : *array) {
215 auto* task = item.getIf<eve::Value::Object>();
216 if (task)
217 task->try_emplace("workRemainderPermille", eve::Value(std::int64_t(0)));
218 }
219 }
220 }
221 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
222 });
223 if (!remainderRegistration.ok()) std::terminate();
224 return result;
225 }();
226 return chain;
227}
228
229template <class T>
230eve::Result<T> snapshotFailure(eve::DiagnosticCode code, std::string message) {
232}
233
234} // namespace
235
237 auto serialized = snapshot();
238 if (!serialized.ok()) return eve::Result<eve::SnapshotEnvelope>::failure(serialized.status());
239 auto payload = eve::Value::fromJson(std::move(serialized).takeValue());
241 return eve::makeSnapshotEnvelope("production.queue", productionSchema(), eve::SchemaVersion(7), instanceId_,
242 revision_, tick_, std::move(payload).takeValue(), hashProvider);
243}
244
246 const eve::SnapshotHashProvider& hashProvider) {
247 if (source.type != "production.queue" || source.schema != productionSchema())
248 return snapshotFailure<void>(eve::DiagnosticCode::InvalidArgument,
249 "snapshot does not belong to production::WorkQueue");
250 if (!instanceId_.isNil() && source.instanceId != instanceId_)
251 return snapshotFailure<void>(eve::DiagnosticCode::Conflict,
252 "snapshot instanceId does not match production::WorkQueue");
253 auto migrated = productionMigrations().migrate(source, eve::SchemaVersion(7), hashProvider);
254 if (!migrated.ok()) return eve::Result<void>::failure(migrated.status());
255 const auto& candidateEnvelope = migrated.value();
256 auto metadata = eve::validateSnapshotPayloadMetadata(candidateEnvelope.payload, candidateEnvelope.revision,
257 candidateEnvelope.tick);
258 if (!metadata.ok()) return eve::Result<void>::failure(metadata.status());
259 auto payload = candidateEnvelope.payload.toJson();
260 if (!payload.ok()) return eve::Result<void>::failure(payload.status());
261
262 WorkQueue candidate(instanceId_);
263 auto restored = candidate.restore(std::move(payload).takeValue());
264 if (!restored.ok()) return eve::Result<void>::failure(restored.status());
265 candidate.instanceId_ = candidateEnvelope.instanceId;
266 candidate.revision_ = candidateEnvelope.revision;
267 candidate.tick_ = candidateEnvelope.tick;
268 *this = std::move(candidate);
270}
271
273 auto value = snapshot(hashProvider);
274 if (!value.ok()) return eve::Result<std::string>::failure(value.status());
275 return std::move(value).andThen(
276 [](eve::SnapshotEnvelope&& envelope) { return eve::serializeSnapshotEnvelope(envelope); });
277}
278
280 auto source = eve::parseSnapshotEnvelope(json, hashProvider);
281 if (!source.ok()) return eve::Result<void>::failure(source.status());
282 return restoreSnapshot(std::move(source).takeValue(), hashProvider);
283}
284
285} // namespace eve::production
double value
Value::Object payload
std::string message
DiagnosticCode code
std::string taskId
std::string string
Battle::Events events
Json object
const UnitySourceAsset & source
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
Human-readable, scoped identifier in the form namespace:name.
Definition Identity.h:411
static std::optional< LogicalId > parse(std::string_view text)
Parses a scoped logical name.
Definition Identity.cpp:36
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
One directed payload migration step.
Definition Snapshot.h:145
Result< void > add(LogicalId schema, SchemaVersion from, SchemaVersion to, Migration migration)
Register one schema-local migration edge.
Definition Snapshot.cpp:196
The canonical owning dynamic value used by data-facing protocols.
Definition Value.h:31
static Result< Value > fromJson(std::string_view json)
Parse one strict JSON value into an owning Value.
Definition Value.cpp:57
std::map< std::string, Value > Object
Definition Value.h:34
std::vector< Value > Array
Definition Value.h:33
constexpr bool isNil() const noexcept
Returns whether this value is the all-zero nil ID.
Definition Identity.h:176
Generic multi-owner, multi-slot continuous-progress work queue.
Definition Production.h:200
eve::Result< std::string > snapshotEnvelopeJson(const eve::SnapshotHashProvider &hashProvider) const
Serializes the common production snapshot envelope.
eve::Result< std::string > snapshot() const
Serializes the complete queue as deterministic JSON.
eve::Result< void > restoreSnapshot(const eve::SnapshotEnvelope &snapshot, const eve::SnapshotHashProvider &hashProvider)
Restores a verified or migrated production envelope atomically.
eve::Result< void > restore(std::string_view json)
Transactionally restores a snapshot; failure preserves current state.
eve::Result< void > restoreSnapshotJson(std::string_view json, const eve::SnapshotHashProvider &hashProvider)
Parses and transactionally restores a common production envelope.
std::string stateName(OrderState state)
Returns the stable lowercase name of a state.
DiagnosticCode
Stable machine-readable diagnostic codes.
Definition Diagnostic.h:47
Result< void > validateSnapshotPayloadMetadata(const Value &payload, Revision revision, SimulationTick tick)
Validate optional payload copies of envelope revision and tick.
Definition Snapshot.cpp:80
Result< SnapshotEnvelope > parseSnapshotEnvelope(std::string_view json, const SnapshotHashProvider &hashProvider)
Parse and verify an envelope from canonical or compatible JSON text.
Definition Snapshot.cpp:190
detail::StrongUint64< detail::SchemaVersionTag > SchemaVersion
Persistent data-format version; not a runtime replacement generation.
Result< std::string > serializeSnapshotEnvelope(const SnapshotEnvelope &snapshot)
Serialize an envelope as deterministic compact JSON.
Definition Snapshot.cpp:184
Result< SnapshotEnvelope > makeSnapshotEnvelope(std::string type, LogicalId schema, SchemaVersion schemaVersion, PersistentId instanceId, Revision revision, SimulationTick tick, Value payload, const SnapshotHashProvider &hashProvider)
Construct and seal a snapshot envelope.
Definition Snapshot.cpp:103
std::function< Result< ContentId >(std::string_view canonicalInput)> SnapshotHashProvider
Injected content-digest implementation used by snapshots.
Definition Snapshot.h:36
Stable outer format shared by persistence and cross-process snapshots.
Definition Snapshot.h:46