载入中...
搜索中...
未找到
Transaction.cpp
浏览该文件的文档.
2
3#include "common/Json.h"
5
6#include <simplesquirrel/simplesquirrel.hpp>
7
8#include <algorithm>
9#include <exception>
10#include <iomanip>
11#include <limits>
12#include <sstream>
13#include <utility>
14
15namespace eve::transaction {
16namespace {
17
18std::string quote(const std::string& value) {
19 std::ostringstream out;
20 out << '"';
21 for (unsigned char c : value) {
22 switch (c) {
23 case '"': out << "\\\""; break;
24 case '\\': out << "\\\\"; break;
25 case '\b': out << "\\b"; break;
26 case '\f': out << "\\f"; break;
27 case '\n': out << "\\n"; break;
28 case '\r': out << "\\r"; break;
29 case '\t': out << "\\t"; break;
30 default:
31 if (c < 0x20)
32 out << "\\u" << std::hex << std::setw(4) << std::setfill('0') << int(c) << std::dec;
33 else
34 out << static_cast<char>(c);
35 }
36 }
37 return out.str() + '"';
38}
39
40std::string canonicalJson(const eve::json::Value& value) {
41 if (value.isNull()) return "null";
42 if (value.isBool()) return value.asBool() ? "true" : "false";
43 if (value.isNumber()) return value.asString();
44 if (value.isString()) return quote(value.asString());
45 if (value.isArray()) {
46 std::string out = "[";
47 for (size_t i = 0; i < value.size(); ++i) {
48 if (i) out += ',';
49 out += canonicalJson(value.at(i));
50 }
51 return out + ']';
52 }
53 if (value.isObject()) {
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 return {};
64}
65
66eve::Value planBindingValue(const Plan* plan) {
67 if (!plan) return eve::Value();
69 {"id", plan->id()},
70 {"identity", plan->identity().isNil() ? std::string{} : plan->identity().format()},
71 {"state", stateName(plan->state())},
72 {"correlation", plan->correlation()},
73 {"causation", plan->causation()},
74 });
75}
76
77eve::LogicalId transactionSchema() {
78 const auto schema = eve::LogicalId::parse("transaction:ledger");
79 if (!schema) std::terminate();
80 return *schema;
81}
82
83const eve::SnapshotMigrationChain& transactionMigrations() {
84 static const eve::SnapshotMigrationChain chain = [] {
86 const auto registration =
87 result.add(transactionSchema(), eve::SchemaVersion(0), eve::SchemaVersion(1),
89 const auto* object = payload.getIf<eve::Value::Object>();
90 if (!object)
92 eve::DiagnosticCode::ParseError, "transaction ledger payload must be an object"));
93 eve::Value::Object migrated = *object;
94 if (!migrated.contains("version")) migrated.emplace("version", eve::Value(std::int64_t(1)));
95 return eve::Result<eve::Value>::success(eve::Value(std::move(migrated)));
96 });
97 if (!registration.ok()) std::terminate();
98 return result;
99 }();
100 return chain;
101}
102
103bool parseU64(const eve::json::Value& value, uint64_t& out) {
104 if (!value.isString()) return false;
105 try {
106 size_t used = 0;
107 out = std::stoull(value.asString(), &used);
108 return used == value.asString().size();
109 } catch (...) {
110 return false;
111 }
112}
113
114bool parseState(const std::string& text, State& state) {
115 if (text == "open")
117 else if (text == "validated")
119 else if (text == "committed")
121 else if (text == "rolled_back")
123 else if (text == "failed")
125 else
126 return false;
127 return true;
128}
129
130} // namespace
131
132namespace {
133
134template <class Call>
135eve::Result<void> invokeParticipant(Call&& call, std::string_view phase, const ITransactionParticipant& participant) {
136 try {
137 return std::forward<Call>(call)();
138 } catch (const std::exception& exception) {
141 std::string(phase) + " threw for participant " + std::string(participant.name()) + ": " + exception.what(),
142 "transaction." + std::string(phase)));
143 } catch (...) {
146 std::string(phase) + " threw an unknown exception for participant " + std::string(participant.name()),
147 "transaction." + std::string(phase)));
148 }
149}
150
151void appendFailure(std::vector<eve::Diagnostic>& diagnostics, const eve::Status& status, std::string_view phase,
152 size_t index, const ITransactionParticipant& participant) {
153 const auto& original = status.diagnostics();
154 diagnostics.insert(diagnostics.end(), original.begin(), original.end());
155 diagnostics.push_back(eve::Diagnostic::error(
156 eve::DiagnosticCode::Failed, std::string(phase) + " failed for participant " + std::string(participant.name()),
157 "transaction." + std::string(phase) + "[" + std::to_string(index) + "]"));
158}
159
160eve::Result<TransactionReceipt> coordinatorFailure(eve::StatusCode code, std::vector<eve::Diagnostic> diagnostics) {
161 if (diagnostics.empty())
162 diagnostics.push_back(
163 eve::Diagnostic::error(eve::DiagnosticCode::Failed, "transaction participant lifecycle failed"));
164 return eve::Result<TransactionReceipt>::failure(eve::Status(code, std::move(diagnostics)));
165}
166
167} // namespace
168
169Ledger::Ledger(eve::PersistentId instanceId, eve::UuidEntropySource transactionEntropy, eve::UuidClock transactionClock)
170 : instanceId_(instanceId),
171 transactionEntropy_(std::move(transactionEntropy)),
172 transactionClock_(std::move(transactionClock)) {
173 if (transactionEntropy_) transactionIdGenerator_.emplace(transactionEntropy_, transactionClock_);
174}
175
177 std::span<ITransactionParticipant*> participants) const {
178 if (context.transactionId().empty())
179 return coordinatorFailure(eve::StatusCode::Rejected,
181 "transaction id must not be empty", "transactionId")});
182 if (participants.empty())
183 return coordinatorFailure(
186 "transaction requires at least one participant", "participants")});
187 for (size_t i = 0; i < participants.size(); ++i) {
188 if (!participants[i])
189 return coordinatorFailure(
191 {eve::Diagnostic::error(eve::DiagnosticCode::InvalidArgument, "participant must not be null",
192 "participants[" + std::to_string(i) + "]")});
193 for (size_t previous = 0; previous < i; ++previous)
194 if (participants[previous] == participants[i])
195 return coordinatorFailure(eve::StatusCode::Rejected,
197 "a participant may occur only once in a transaction",
198 "participants[" + std::to_string(i) + "]")});
199 }
200
201 TransactionReceipt receipt;
202 receipt.identity = context.identity();
203 receipt.transactionId = context.transactionId();
204 receipt.correlationId = context.correlationId();
205 receipt.causationId = context.causationId();
206 receipt.participantCount = participants.size();
207
208 std::vector<size_t> prepared;
209 prepared.reserve(participants.size());
210 for (size_t i = 0; i < participants.size(); ++i) {
211 auto result = invokeParticipant([&] { return participants[i]->prepare(context); }, "prepare", *participants[i]);
212 if (!result) {
213 std::vector<eve::Diagnostic> diagnostics;
214 appendFailure(diagnostics, result.status(), "prepare", i, *participants[i]);
215 bool cleanupFailed = false;
216 for (auto it = prepared.rbegin(); it != prepared.rend(); ++it) {
217 auto rollback = invokeParticipant([&] { return participants[*it]->rollback(context); }, "rollback",
218 *participants[*it]);
219 if (!rollback) {
220 cleanupFailed = true;
221 appendFailure(diagnostics, rollback.status(), "rollback", *it, *participants[*it]);
222 }
223 }
224 return coordinatorFailure(cleanupFailed ? eve::StatusCode::Failed : result.code(), std::move(diagnostics));
225 }
226 prepared.push_back(i);
227 }
228 receipt.preparedCount = prepared.size();
229
230 std::vector<size_t> committed;
231 committed.reserve(prepared.size());
232 for (size_t i : prepared) {
233 auto result = invokeParticipant([&] { return participants[i]->commit(context); }, "commit", *participants[i]);
234 if (!result) {
235 std::vector<eve::Diagnostic> diagnostics;
236 appendFailure(diagnostics, result.status(), "commit", i, *participants[i]);
237 bool cleanupFailed = false;
238
239 // A failed commit is contractually still uncommitted. Roll it and
240 // every later prepared participant back before compensating effects
241 // that were already made observable by earlier commits.
242 for (auto it = prepared.rbegin(); it != prepared.rend(); ++it) {
243 if (std::find(committed.begin(), committed.end(), *it) != committed.end()) continue;
244 auto rollback = invokeParticipant([&] { return participants[*it]->rollback(context); }, "rollback",
245 *participants[*it]);
246 if (!rollback) {
247 cleanupFailed = true;
248 appendFailure(diagnostics, rollback.status(), "rollback", *it, *participants[*it]);
249 }
250 }
251 for (auto it = committed.rbegin(); it != committed.rend(); ++it) {
252 auto compensation = invokeParticipant([&] { return participants[*it]->compensate(context); },
253 "compensation", *participants[*it]);
254 if (!compensation) {
255 cleanupFailed = true;
256 appendFailure(diagnostics, compensation.status(), "compensation", *it, *participants[*it]);
257 }
258 }
259 return coordinatorFailure(cleanupFailed ? eve::StatusCode::Failed : result.code(), std::move(diagnostics));
260 }
261 committed.push_back(i);
262 }
263
264 receipt.committedCount = committed.size();
267}
268
270 std::span<ITransactionParticipant*> participants) const {
271 if (context.transactionId().empty())
272 return coordinatorFailure(eve::StatusCode::Rejected,
274 "transaction id must not be empty", "transactionId")});
275 if (participants.empty())
276 return coordinatorFailure(
279 "compensation requires at least one participant", "participants")});
280 for (size_t i = 0; i < participants.size(); ++i) {
281 if (!participants[i])
282 return coordinatorFailure(
284 {eve::Diagnostic::error(eve::DiagnosticCode::InvalidArgument, "participant must not be null",
285 "participants[" + std::to_string(i) + "]")});
286 for (size_t previous = 0; previous < i; ++previous)
287 if (participants[previous] == participants[i])
288 return coordinatorFailure(eve::StatusCode::Rejected,
290 "a participant may occur only once in compensation",
291 "participants[" + std::to_string(i) + "]")});
292 }
293
294 TransactionReceipt receipt;
295 receipt.identity = context.identity();
296 receipt.transactionId = context.transactionId();
297 receipt.correlationId = context.correlationId();
298 receipt.causationId = context.causationId();
299 receipt.participantCount = participants.size();
300 std::vector<eve::Diagnostic> diagnostics;
301 for (size_t offset = 0; offset < participants.size(); ++offset) {
302 const size_t index = participants.size() - 1 - offset;
303 auto result = invokeParticipant([&] { return participants[index]->compensate(context); }, "compensation",
304 *participants[index]);
305 if (!result) {
306 appendFailure(diagnostics, result.status(), "compensation", index, *participants[index]);
307 continue;
308 }
309 ++receipt.compensatedCount;
310 }
311 if (!diagnostics.empty()) {
313 return coordinatorFailure(eve::StatusCode::Failed, std::move(diagnostics));
314 }
317}
318
319std::string stateName(State state) {
320 switch (state) {
321 case State::Open: return "open";
322 case State::Validated: return "validated";
323 case State::Committed: return "committed";
324 case State::RolledBack: return "rolled_back";
325 case State::Failed: return "failed";
326 }
327 return "unknown";
328}
329
330Plan::Plan(std::string id, std::string correlation, std::string causation, eve::TransactionId identity,
331 eve::UuidEntropySource operationEntropy, eve::UuidClock operationClock)
332 : id_(std::move(id)), identity_(identity), correlation_(std::move(correlation)), causation_(std::move(causation)) {
333 if (operationEntropy) operationIdGenerator_.emplace(std::move(operationEntropy), std::move(operationClock));
334 emit("opened");
335}
336
337const std::string& Plan::id() const { return id_; }
338const eve::TransactionId& Plan::identity() const noexcept { return identity_; }
339State Plan::state() const { return state_; }
340const std::string& Plan::correlation() const { return correlation_; }
341const std::string& Plan::causation() const { return causation_; }
342
343void Plan::emit(const std::string& type, const std::string& operationId, const std::string& detail) {
344 Event event;
345 event.sequence = nextEvent_++;
346 event.transactionId = id_;
347 event.operationId = operationId;
348 event.type = type;
349 event.detail = detail;
350 event.transactionIdentity = identity_;
351 if (const auto parsed = eve::OperationId::parse(operationId)) event.operationIdentity = *parsed;
352 events_.push_back(std::move(event));
353}
354
355eve::Result<eve::OperationId> Plan::stage(const std::string& kind, const std::string& target,
356 const std::string& payloadJson, eve::OperationId operationId) {
357 const auto failure = [](eve::DiagnosticCode code, std::string message) {
359 };
360 error_.clear();
361 if (state_ != State::Open) return failure(eve::DiagnosticCode::Conflict, "transaction plan is frozen");
362 if (kind.empty()) return failure(eve::DiagnosticCode::InvalidArgument, "operation kind must not be empty");
363 auto document = eve::json::Document::parse(payloadJson);
364 if (!document.valid()) return failure(eve::DiagnosticCode::ParseError, "payload must be valid JSON");
365 if (nextOperation_ == std::numeric_limits<uint64_t>::max())
366 return failure(eve::DiagnosticCode::Failed, "operation sequence exhausted");
367
368 if (operationId.isNil()) {
369 if (!operationIdGenerator_)
371 "canonical operation identity requires an injected UUID entropy source");
372 const auto generated = operationIdGenerator_->generate();
373 if (!generated) return failure(eve::DiagnosticCode::Failed, "canonical operation identity generation failed");
374 operationId = eve::OperationId::fromUuid(*generated);
375 }
376 if (findOperation(operationId))
377 return failure(eve::DiagnosticCode::Conflict, "operation identity already exists in this plan");
378
380 operation.identity = operationId;
381 operation.id = operationId.format();
382 operation.kind = kind;
383 operation.target = target;
384 operation.payload = canonicalJson(document.root());
385 operations_.push_back(std::move(operation));
386 ++nextOperation_;
387 emit("staged", operations_.back().id);
389}
390
391Operation* findMutable(std::deque<Operation>& operations, const std::string& id) {
392 auto it = std::find_if(operations.begin(), operations.end(), [&id](const Operation& op) { return op.id == id; });
393 return it == operations.end() ? nullptr : &*it;
394}
395
396Operation* findMutable(std::deque<Operation>& operations, eve::OperationId id) {
397 auto it = std::find_if(operations.begin(), operations.end(),
398 [id](const Operation& op) { return !id.isNil() && op.identity == id; });
399 return it == operations.end() ? nullptr : &*it;
400}
401
402eve::Result<void> Plan::markValid(const std::string& operationId) {
403 if (state_ != State::Open) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction is not open", {}, {}, "transaction"));
404 Operation* operation = findMutable(operations_, operationId);
405 if (!operation) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::NotFound, "operation was not found", {}, {}, "transaction"));
406 operation->checked = true;
407 operation->valid = true;
408 operation->error.clear();
409 emit("operation_valid", operationId);
411}
412
413eve::Result<void> Plan::markValid(eve::OperationId operationId) {
414 if (state_ != State::Open) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction is not open", {}, {}, "transaction"));
415 Operation* operation = findMutable(operations_, operationId);
416 if (!operation) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::NotFound, "operation was not found", {}, {}, "transaction"));
417 operation->checked = true;
418 operation->valid = true;
419 operation->error.clear();
420 emit("operation_valid", operation->id);
422}
423
424eve::Result<void> Plan::markInvalid(const std::string& operationId, const std::string& error) {
425 if (state_ != State::Open) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction is not open", {}, {}, "transaction"));
426 Operation* operation = findMutable(operations_, operationId);
427 if (!operation) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::NotFound, "operation was not found", {}, {}, "transaction"));
428 operation->checked = true;
429 operation->valid = false;
430 operation->error = error.empty() ? "invalid operation" : error;
431 emit("operation_invalid", operationId, operation->error);
433}
434
435eve::Result<void> Plan::markInvalid(eve::OperationId operationId, const std::string& error) {
436 if (state_ != State::Open) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction is not open", {}, {}, "transaction"));
437 Operation* operation = findMutable(operations_, operationId);
438 if (!operation) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::NotFound, "operation was not found", {}, {}, "transaction"));
439 operation->checked = true;
440 operation->valid = false;
441 operation->error = error.empty() ? "invalid operation" : error;
442 emit("operation_invalid", operation->id, operation->error);
444}
445
446eve::Result<void> Plan::validate() {
447 error_.clear();
448 if (state_ != State::Open) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction is not open", {}, {}, "transaction"));
449 if (operations_.empty()) {
450 error_ = "transaction has no operations";
451 emit("validation_failed", {}, error_);
453 }
454 for (const auto& operation : operations_) {
455 if (!operation.checked || !operation.valid) {
456 error_ = operation.checked ? operation.error : "operation was not validated: " + operation.id;
457 emit("validation_failed", operation.id, error_);
459 }
460 }
461 state_ = State::Validated;
462 emit("validated");
464}
465
466eve::Result<void> Plan::commit() {
467 if (state_ != State::Validated) return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction is not validated", {}, {}, "transaction"));
468 state_ = State::Committed;
469 emit("committed");
471}
472
473eve::Result<void> Plan::rollback(const std::string& reason) {
474 if (state_ != State::Open && state_ != State::Validated)
475 return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction cannot be rolled back in its current state", {}, {}, "transaction"));
476 state_ = State::RolledBack;
477 error_ = reason;
478 emit("rolled_back", {}, reason);
480}
481
482eve::Result<void> Plan::fail(const std::string& error) {
483 if (state_ == State::Committed || state_ == State::RolledBack || state_ == State::Failed)
484 return eve::Result<void>::failure(eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transaction is already terminal", {}, {}, "transaction"));
485 state_ = State::Failed;
486 error_ = error.empty() ? "transaction failed" : error;
487 emit("failed", {}, error_);
489}
490
491const std::string& Plan::error() const { return error_; }
492int Plan::operationCount() const { return static_cast<int>(operations_.size()); }
493const Operation* Plan::operationAt(int index) const {
494 return index >= 0 && static_cast<size_t>(index) < operations_.size() ? &operations_[static_cast<size_t>(index)]
495 : nullptr;
496}
497const Operation* Plan::findOperation(const std::string& operationId) const {
498 auto it = std::find_if(operations_.begin(), operations_.end(),
499 [&operationId](const Operation& operation) { return operation.id == operationId; });
500 return it == operations_.end() ? nullptr : &*it;
501}
502const Operation* Plan::findOperation(eve::OperationId operationId) const {
503 auto it = std::find_if(operations_.begin(), operations_.end(), [operationId](const Operation& operation) {
504 return !operationId.isNil() && operation.identity == operationId;
505 });
506 return it == operations_.end() ? nullptr : &*it;
507}
508int Plan::eventCount() const { return static_cast<int>(events_.size()); }
509const Event* Plan::eventAt(int index) const {
510 return index >= 0 && static_cast<size_t>(index) < events_.size() ? &events_[static_cast<size_t>(index)] : nullptr;
511}
512
513std::string Plan::snapshotJson() const {
514 std::ostringstream out;
515 out << "{\"id\":" << quote(id_);
516 if (!identity_.isNil()) out << ",\"identity\":" << quote(identity_.format());
517 out << ",\"state\":" << quote(stateName(state_)) << ",\"correlation\":" << quote(correlation_)
518 << ",\"causation\":" << quote(causation_) << ",\"error\":" << quote(error_)
519 << ",\"nextOperation\":" << quote(std::to_string(nextOperation_))
520 << ",\"nextEvent\":" << quote(std::to_string(nextEvent_)) << ",\"operations\":[";
521 for (size_t i = 0; i < operations_.size(); ++i) {
522 if (i) out << ',';
523 const auto& op = operations_[i];
524 out << "{\"id\":" << quote(op.id);
525 if (!op.identity.isNil()) out << ",\"identity\":" << quote(op.identity.format());
526 out << ",\"kind\":" << quote(op.kind) << ",\"target\":" << quote(op.target) << ",\"payload\":" << op.payload
527 << ",\"checked\":" << (op.checked ? "true" : "false") << ",\"valid\":" << (op.valid ? "true" : "false")
528 << ",\"error\":" << quote(op.error) << '}';
529 }
530 out << "],\"events\":[";
531 for (size_t i = 0; i < events_.size(); ++i) {
532 if (i) out << ',';
533 const auto& event = events_[i];
534 out << "{\"sequence\":" << quote(std::to_string(event.sequence))
535 << ",\"transactionId\":" << quote(event.transactionId) << ",\"operationId\":" << quote(event.operationId)
536 << ",\"type\":" << quote(event.type) << ",\"detail\":" << quote(event.detail) << '}';
537 }
538 return out.str() + "]}";
539}
540
541eve::Result<Plan*> Ledger::create(const std::string& correlation, const std::string& causation,
542 eve::TransactionId identity) {
543 const auto failure = [](eve::DiagnosticCode code, std::string message) {
545 };
546 if (identity.isNil()) {
547 if (!transactionIdGenerator_)
549 "canonical transaction identity requires an injected UUID entropy source");
550 const auto generated = transactionIdGenerator_->generate();
551 if (!generated) return failure(eve::DiagnosticCode::Failed, "canonical transaction identity generation failed");
552 identity = eve::TransactionId::fromUuid(*generated);
553 }
554 if (find(identity)) return failure(eve::DiagnosticCode::Conflict, "transaction identity already exists");
555 plans_.push_back(std::unique_ptr<Plan>(
556 new Plan(identity.format(), correlation, causation, identity, transactionEntropy_, transactionClock_)));
558}
559
560Plan* Ledger::find(const std::string& transactionId) {
561 auto it = std::find_if(plans_.begin(), plans_.end(),
562 [&transactionId](const auto& plan) { return plan->id() == transactionId; });
563 return it == plans_.end() ? nullptr : it->get();
564}
565Plan* Ledger::find(eve::TransactionId transactionId) {
566 if (transactionId.isNil()) return nullptr;
567 auto it = std::find_if(plans_.begin(), plans_.end(),
568 [transactionId](const auto& plan) { return plan->identity() == transactionId; });
569 return it == plans_.end() ? nullptr : it->get();
570}
571int Ledger::count() const { return static_cast<int>(plans_.size()); }
572Plan* Ledger::at(int index) {
573 return index >= 0 && static_cast<size_t>(index) < plans_.size() ? plans_[static_cast<size_t>(index)].get()
574 : nullptr;
575}
576
577std::string Ledger::snapshotJson() const {
578 std::ostringstream out;
579 out << "{\"version\":1,\"nextTransaction\":" << quote(std::to_string(nextTransaction_)) << ",\"plans\":[";
580 for (size_t i = 0; i < plans_.size(); ++i) {
581 if (i) out << ',';
582 out << plans_[i]->snapshotJson();
583 }
584 return out.str() + "]}";
585}
586
587eve::Result<void> Ledger::restore(std::string_view json) {
588 const auto failure = [](eve::DiagnosticCode code, std::string message) {
589 return eve::Result<void>::failure(eve::Diagnostic::error(code, message, "transaction.ledger"));
590 };
591 auto document = eve::json::Document::parse(std::string(json));
592 auto root = document.root();
593 uint64_t restoredNext = 0;
594 if (!document.valid() || !root.isObject() || root.getInt("version") != 1 ||
595 !parseU64(root.get("nextTransaction"), restoredNext) || restoredNext == 0 || !root.get("plans").isArray()) {
596 return failure(eve::DiagnosticCode::ParseError, "invalid transaction ledger snapshot");
597 }
598 std::vector<std::unique_ptr<Plan>> restored;
599 uint64_t previousNumeric = 0;
600 std::vector<std::string> restoredPlanIds;
601 auto plans = root.get("plans");
602 for (size_t i = 0; i < plans.size(); ++i) {
603 auto value = plans.at(i);
604 State state = State::Open;
605 uint64_t nextOperation = 0, nextEvent = 0;
606 if (!value.isObject() || !parseState(value.getString("state"), state) ||
607 !parseU64(value.get("nextOperation"), nextOperation) || nextOperation == 0 ||
608 !parseU64(value.get("nextEvent"), nextEvent) || nextEvent == 0 || !value.get("operations").isArray() ||
609 !value.get("events").isArray()) {
610 return failure(eve::DiagnosticCode::ParseError, "invalid transaction plan at index " + std::to_string(i));
611 }
612 const std::string id = value.getString("id");
613 const std::string identityText = value.getString("identity");
614 eve::TransactionId identity;
615 bool canonicalIdentity = false;
616 if (!identityText.empty()) {
617 const auto parsed = eve::TransactionId::parse(identityText);
618 if (!parsed || parsed->isNil() || id != parsed->format()) {
619 return failure(eve::DiagnosticCode::ParseError,
620 "invalid canonical transaction id at index " + std::to_string(i));
621 }
622 identity = *parsed;
623 canonicalIdentity = true;
624 } else if (const auto parsed = eve::TransactionId::parse(id); parsed && !parsed->isNil()) {
625 identity = *parsed;
626 canonicalIdentity = true;
627 }
628
629 uint64_t numeric = 0;
630 if (!canonicalIdentity) {
631 const std::string prefix = "transaction-";
632 try {
633 numeric = std::stoull(id.substr(prefix.size()));
634 } catch (...) {
635 numeric = 0;
636 }
637 if (!id.starts_with(prefix) || numeric <= previousNumeric || numeric >= restoredNext) {
638 return failure(eve::DiagnosticCode::ParseError, "invalid transaction id at index " + std::to_string(i));
639 }
640 }
641 if (id.empty() || std::find(restoredPlanIds.begin(), restoredPlanIds.end(), id) != restoredPlanIds.end()) {
642 return failure(eve::DiagnosticCode::Conflict, "duplicate transaction id at index " + std::to_string(i));
643 }
644 restoredPlanIds.push_back(id);
645 auto plan = std::unique_ptr<Plan>(new Plan(id, value.getString("correlation"), value.getString("causation"),
646 identity, transactionEntropy_, transactionClock_));
647 plan->state_ = state;
648 plan->error_ = value.getString("error");
649 plan->operations_.clear();
650 plan->events_.clear();
651 auto operations = value.get("operations");
652 std::vector<std::string> restoredOperationIds;
653 for (size_t j = 0; j < operations.size(); ++j) {
654 auto opValue = operations.at(j);
655 if (!opValue.isObject() || opValue.getString("id").empty() || opValue.getString("kind").empty() ||
656 !opValue.get("payload") || !opValue.get("checked").isBool() || !opValue.get("valid").isBool()) {
657 return failure(eve::DiagnosticCode::ParseError, "invalid operation in plan " + id);
658 }
659 Operation op;
660 op.id = opValue.getString("id");
661 const std::string operationIdentityText = opValue.getString("identity");
662 if (!operationIdentityText.empty()) {
663 const auto parsed = eve::OperationId::parse(operationIdentityText);
664 if (!parsed || parsed->isNil() || op.id != parsed->format()) {
665 return failure(eve::DiagnosticCode::ParseError, "invalid canonical operation id in plan " + id);
666 }
667 op.identity = *parsed;
668 } else if (const auto parsed = eve::OperationId::parse(op.id); parsed && !parsed->isNil()) {
669 op.identity = *parsed;
670 }
671 if (std::find(restoredOperationIds.begin(), restoredOperationIds.end(), op.id) !=
672 restoredOperationIds.end()) {
673 return failure(eve::DiagnosticCode::Conflict, "duplicate operation id in plan " + id);
674 }
675 restoredOperationIds.push_back(op.id);
676 op.kind = opValue.getString("kind");
677 op.target = opValue.getString("target");
678 op.payload = canonicalJson(opValue.get("payload"));
679 op.checked = opValue.getBool("checked");
680 op.valid = opValue.getBool("valid");
681 op.error = opValue.getString("error");
682 plan->operations_.push_back(std::move(op));
683 }
684 uint64_t previousEvent = 0;
685 auto events = value.get("events");
686 for (size_t j = 0; j < events.size(); ++j) {
687 auto eventValue = events.at(j);
688 Event event;
689 if (!eventValue.isObject() || !parseU64(eventValue.get("sequence"), event.sequence) ||
690 event.sequence <= previousEvent || eventValue.getString("transactionId") != id ||
691 eventValue.getString("type").empty()) {
692 return failure(eve::DiagnosticCode::ParseError, "invalid event in plan " + id);
693 }
694 event.transactionId = id;
695 event.operationId = eventValue.getString("operationId");
696 event.type = eventValue.getString("type");
697 event.detail = eventValue.getString("detail");
698 event.transactionIdentity = identity;
699 if (const auto parsed = eve::OperationId::parse(event.operationId)) event.operationIdentity = *parsed;
700 previousEvent = event.sequence;
701 plan->events_.push_back(std::move(event));
702 }
703 if (!plan->events_.empty() && plan->events_.back().sequence >= nextEvent) {
704 return failure(eve::DiagnosticCode::ParseError, "invalid event allocator in plan " + id);
705 }
706 plan->nextOperation_ = nextOperation;
707 plan->nextEvent_ = nextEvent;
708 if (!canonicalIdentity) previousNumeric = numeric;
709 restored.push_back(std::move(plan));
710 }
711 plans_ = std::move(restored);
712 nextTransaction_ = restoredNext;
714}
715
716eve::Result<eve::SnapshotEnvelope> Ledger::snapshot(const eve::SnapshotHashProvider& hashProvider) const {
717 auto payload = eve::Value::fromJson(snapshotJson());
719 return eve::makeSnapshotEnvelope("transaction.ledger", transactionSchema(), eve::SchemaVersion(1), instanceId_,
720 revision_, tick_, std::move(payload).takeValue(), hashProvider);
721}
722
723eve::Result<void> Ledger::restoreSnapshot(const eve::SnapshotEnvelope& source,
724 const eve::SnapshotHashProvider& hashProvider) {
725 if (source.type != "transaction.ledger" || source.schema != transactionSchema())
727 "snapshot does not belong to transaction::Ledger"));
728 if (!instanceId_.isNil() && source.instanceId != instanceId_)
730 eve::DiagnosticCode::Conflict, "snapshot instanceId does not match transaction::Ledger"));
731
732 auto migrated = transactionMigrations().migrate(source, eve::SchemaVersion(1), hashProvider);
733 if (!migrated.ok()) return eve::Result<void>::failure(migrated.status());
734 const auto& candidateEnvelope = migrated.value();
735 auto metadata = eve::validateSnapshotPayloadMetadata(candidateEnvelope.payload, candidateEnvelope.revision,
736 candidateEnvelope.tick);
737 if (!metadata.ok()) return eve::Result<void>::failure(metadata.status());
738 auto payload = candidateEnvelope.payload.toJson();
739 if (!payload.ok()) return eve::Result<void>::failure(payload.status());
740
741 Ledger candidate(instanceId_, transactionEntropy_, transactionClock_);
742 auto restored = candidate.restore(std::move(payload).takeValue());
743 if (!restored.ok()) return eve::Result<void>::failure(restored.status());
744 candidate.instanceId_ = candidateEnvelope.instanceId;
745 candidate.revision_ = candidateEnvelope.revision;
746 candidate.tick_ = candidateEnvelope.tick;
747 *this = std::move(candidate);
749}
750
751eve::Result<std::string> Ledger::snapshotEnvelopeJson(const eve::SnapshotHashProvider& hashProvider) const {
752 auto value = snapshot(hashProvider);
753 if (!value.ok()) return eve::Result<std::string>::failure(value.status());
754 return std::move(value).andThen(
755 [](eve::SnapshotEnvelope&& envelope) { return eve::serializeSnapshotEnvelope(envelope); });
756}
757
758eve::Result<void> Ledger::restoreSnapshotJson(std::string_view json, const eve::SnapshotHashProvider& hashProvider) {
759 auto source = eve::parseSnapshotEnvelope(json, hashProvider);
760 if (!source.ok()) return eve::Result<void>::failure(source.status());
761 return restoreSnapshot(std::move(source).takeValue(), hashProvider);
762}
763
764Ledger* Transaction::newLedger() {
765 auto* module = Transaction::create();
766 module->ledgers_.push_back(std::make_unique<Ledger>());
767 return module->ledgers_.back().get();
768}
769
771
772void Transaction::expose(ssq::Table& table) {
773 const HSQUIRRELVM vm = table.getHandle();
774 auto operation =
775 table.addClass<Operation>("TransactionOperation", std::function<Operation*()>([] { return nullptr; }), false);
776 operation.addFunc("getId", [](Operation* o) { return o ? o->id : std::string{}; });
777 operation.addFunc("getKind", [](Operation* o) { return o ? o->kind : std::string{}; });
778 operation.addFunc("getTarget", [](Operation* o) { return o ? o->target : std::string{}; });
779 operation.addFunc("getPayload", [](Operation* o) { return o ? o->payload : std::string{}; });
780 operation.addFunc("isChecked", [](Operation* o) { return o && o->checked; });
781 operation.addFunc("isValid", [](Operation* o) { return o && o->valid; });
782 operation.addFunc("getError", [](Operation* o) { return o ? o->error : std::string{}; });
783
784 auto event = table.addClass<Event>("TransactionEvent", std::function<Event*()>([] { return nullptr; }), false);
785 event.addFunc("getSequence", [](Event* e) { return e ? static_cast<int64_t>(e->sequence) : int64_t{0}; });
786 event.addFunc("getTransactionId", [](Event* e) { return e ? e->transactionId : std::string{}; });
787 event.addFunc("getOperationId", [](Event* e) { return e ? e->operationId : std::string{}; });
788 event.addFunc("getType", [](Event* e) { return e ? e->type : std::string{}; });
789 event.addFunc("getDetail", [](Event* e) { return e ? e->detail : std::string{}; });
790
791 auto plan = table.addClass<Plan>("TransactionPlan", std::function<Plan*()>([] { return nullptr; }), false);
792 plan.addFunc("getId", &Plan::id);
793 plan.addFunc("getState", [](Plan* p) { return p ? stateName(p->state()) : std::string{}; });
794 plan.addFunc("getCorrelation", &Plan::correlation);
795 plan.addFunc("getCausation", &Plan::causation);
796 // Project the canonical Result API into the common Squirrel result-table
797 // schema. The script surface retains the short domain names, while the
798 // failure status and diagnostics remain observable to the caller.
799 plan.addFunc("markValid", [vm](Plan* p, const std::string& id) {
800 if (!p)
803 "transaction plan must not be null", "plan", {}, "transaction")));
804 return eve::script::projectResult(vm, p->markValid(id));
805 });
806 plan.addFunc("markInvalid", [vm](Plan* p, const std::string& id, const std::string& error) {
807 if (!p)
810 "transaction plan must not be null", "plan", {}, "transaction")));
811 return eve::script::projectResult(vm, p->markInvalid(id, error));
812 });
813 plan.addFunc("validate", [vm](Plan* p) {
814 if (!p)
817 "transaction plan must not be null", "plan", {}, "transaction")));
818 return eve::script::projectResult(vm, p->validate());
819 });
820 plan.addFunc("commit", [vm](Plan* p) {
821 if (!p)
824 "transaction plan must not be null", "plan", {}, "transaction")));
825 return eve::script::projectResult(vm, p->commit());
826 });
827 plan.addFunc("rollback", [vm](Plan* p, const std::string& reason) {
828 if (!p)
831 "transaction plan must not be null", "plan", {}, "transaction")));
832 return eve::script::projectResult(vm, p->rollback(reason));
833 });
834 plan.addFunc("fail", [vm](Plan* p, const std::string& error) {
835 if (!p)
838 "transaction plan must not be null", "plan", {}, "transaction")));
839 return eve::script::projectResult(vm, p->fail(error));
840 });
841 plan.addFunc("getError", &Plan::error);
842 plan.addFunc("operationCount", &Plan::operationCount);
843 plan.addFunc("operationAt",
844 [](Plan* p, int i) -> Operation* { return p ? const_cast<Operation*>(p->operationAt(i)) : nullptr; });
845 plan.addFunc("findOperation", [](Plan* p, const std::string& id) -> Operation* {
846 return p ? const_cast<Operation*>(p->findOperation(id)) : nullptr;
847 });
848 plan.addFunc("eventCount", &Plan::eventCount);
849 plan.addFunc("eventAt", [](Plan* p, int i) -> Event* { return p ? const_cast<Event*>(p->eventAt(i)) : nullptr; });
850 plan.addFunc("snapshotJson", &Plan::snapshotJson);
851 plan.addFunc("stage", [vm](Plan* value, const std::string& kind, const std::string& target,
852 const std::string& payloadJson, const std::string& operationIdentity) {
853 if (!value)
855 vm,
857 "transaction plan must not be null",
858 "plan", {}, "transaction")),
859 [](eve::OperationId id) { return eve::Value(id.isNil() ? std::string{} : id.format()); });
860 eve::OperationId identity;
861 if (!operationIdentity.empty()) {
862 const auto parsed = eve::OperationId::parse(operationIdentity);
863 if (!parsed)
865 vm,
867 eve::DiagnosticCode::InvalidArgument, "operation identity must be canonical UUID text",
868 "operationIdentity", {}, "transaction")),
869 [](eve::OperationId id) { return eve::Value(id.isNil() ? std::string{} : id.format()); });
870 identity = *parsed;
871 }
873 vm, value->stage(kind, target, payloadJson, identity),
874 [](eve::OperationId id) { return eve::Value(id.isNil() ? std::string{} : id.format()); });
875 });
876
877 auto ledger = table.addClass<Ledger>("TransactionLedger", std::function<Ledger*()>([] { return nullptr; }), false);
878 // Bind the compatibility spelling through a typed lambda so Ledger's
879 // canonical TransactionId overload is never considered by deduction.
880 ledger.addFunc("find",
881 [](Ledger* ledger, const std::string& id) -> Plan* { return ledger ? ledger->find(id) : nullptr; });
882 ledger.addFunc("count", &Ledger::count);
883 ledger.addFunc("at", &Ledger::at);
884 ledger.addFunc("snapshotJson", &Ledger::snapshotJson);
885 ledger.addFunc("restore", [vm](Ledger* value, const std::string& json) {
886 if (!value)
889 "transaction ledger must not be null", "ledger",
890 {}, "transaction")));
891 return eve::script::projectResult(vm, value->restore(json));
892 });
893 ledger.addFunc("create", [vm](Ledger* value, const std::string& correlation, const std::string& causation,
894 const std::string& transactionIdentity) {
895 if (!value)
899 "transaction ledger must not be null", "ledger", {}, "transaction")),
900 [](Plan* plan) { return planBindingValue(plan); });
901 eve::TransactionId identity;
902 if (!transactionIdentity.empty()) {
903 const auto parsed = eve::TransactionId::parse(transactionIdentity);
904 if (!parsed)
906 vm,
908 eve::DiagnosticCode::InvalidArgument, "transaction identity must be canonical UUID text",
909 "transactionIdentity", {}, "transaction")),
910 [](Plan* plan) { return planBindingValue(plan); });
911 identity = *parsed;
912 }
913 return eve::script::projectResult(vm, value->create(correlation, causation, identity),
914 [](Plan* plan) { return planBindingValue(plan); });
915 });
916
917 auto cls = table.addClass(name, Transaction::create, false);
918 expose(cls);
919}
920
921void Transaction::expose(ssq::Class& cls) {
922 cls.addFunc("getName", &Transaction::getName);
923 cls.addFunc("newLedger", [](Transaction*) { return Transaction::newLedger(); });
924}
925
926} // namespace eve::transaction
LogicalId target
ActionParameterOperation operation
double value
Value::Object payload
int root
Definition AnimSmr.cpp:119
float phase
Definition CaveMesh.cpp:58
struct SQVM * HSQUIRRELVM
glm::vec4 p[6]
HSQUIRRELVM vm
Definition ECS.cpp:20
HSQOBJECT cls
Definition ECS.cpp:21
std::string message
DiagnosticCode code
wgpu::PopErrorScopeStatus status
std::int32_t c
std::string text
TokenKind kind
size_t offset
char quote
std::string name
#define Module_IMPL(ModuleName, newExpr)
Definition Module.h:26
graphics::Canvas * previous
std::string error
Definition Package.cpp:60
std::string id
Definition PlayHost.cpp:108
int detail
std::string string
The single Squirrel projection for common Result, Status and Value.
Battle::Events events
Json object
uint32_t index
const UnitySourceAsset & source
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
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
Structured status and zero or more diagnostics for an operation.
Definition Status.h:68
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
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
Id128 public API.
Definition Identity.h:112
std::string format() const
Formats the ID in lower-case canonical UUID text.
Definition Identity.h:173
constexpr bool isNil() const noexcept
Returns whether this value is the all-zero nil ID.
Definition Identity.h:176
static std::optional< Id128 > parse(std::string_view text) noexcept
Parses canonical UUID text.
Definition Identity.h:151
static constexpr Id128 fromUuid(const Id128< OtherTag > &value) noexcept
Re-tags an existing UUID at an explicit domain boundary.
Definition Identity.h:136
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
Result< TransactionReceipt > execute(const TransactionContext &context, std::span< ITransactionParticipant * > participants) const
Atomically execute borrowed participants for one context.
Result< TransactionReceipt > compensate(const TransactionContext &context, std::span< ITransactionParticipant * > participants) const
Compensate already committed participants in reverse dependency order.
Deterministic owner and identifier allocator for transaction plans.
Ledger(eve::PersistentId instanceId={}, eve::UuidEntropySource transactionEntropy={}, eve::UuidClock transactionClock={})
Creates an empty ledger with an optional persistent identity.
eve::Result< void > restore(std::string_view json)
Transactionally restores a raw ledger snapshot.
Immutable-after-validation operation plan for script-coordinated work.
Immutable context shared by every participant of one transaction.
Definition Transaction.h:30
Script module factory for generic transaction ledgers.
std::string stateName(OrderState state)
Returns the stable lowercase name of a state.
std::unordered_map< std::string, SkillDefinition > & table()
Definition Skill.cpp:65
ssq::Table projectResult(HSQUIRRELVM vm, Result< void > &&result)
Consume and project a void native Result using the common schema.
std::string stateName(State state)
Returns the stable lowercase name of a transaction state.
State
Lifecycle state of a generic transaction plan.
Operation * findMutable(std::deque< Operation > &operations, const std::string &id)
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
StatusCode
Stable outcome category for an operation.
Definition Status.h:27
std::function< std::chrono::system_clock::time_point()> UuidClock
Clock callback used by the injectable UUIDv7 generator.
Definition Identity.h:469
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
Result< std::string > serializeSnapshotEnvelope(const SnapshotEnvelope &snapshot)
Serialize an envelope as deterministic compact JSON.
Definition Snapshot.cpp:184
std::function< bool(std::span< std::uint8_t > bytes)> UuidEntropySource
Entropy callback used by UUIDv7 generation.
Definition Identity.h:466
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
Deterministically ordered transaction lifecycle event.
One inert, script-interpreted operation in a transaction plan.
eve::OperationId identity
Canonical operation identity; nil for legacy string-only plans.
Observable summary returned after an atomic participant commit.
Definition Transaction.h:97
eve::TransactionId identity
Canonical transaction identity; nil for legacy string contexts.
Definition Transaction.h:99