载入中...
搜索中...
未找到
GameEvent.cpp
浏览该文件的文档.
2
4
5#include "common/Json.h"
7
8#include <simplesquirrel/simplesquirrel.hpp>
9
10#include <algorithm>
11#include <cmath>
12#include <iomanip>
13#include <limits>
14#include <sstream>
15#include <string_view>
16
17namespace eve::game_event {
18namespace {
19
20eve::schema::FieldDefinition schemaField(std::string name, eve::schema::ValueType type, bool required = true) {
22 result.name = std::move(name);
23 result.type = type;
24 result.required = required;
25 return result;
26}
27
28eve::schema::SchemaDefinition eventStreamSchemaDefinition() {
30 event.additionalProperties = false;
31 event.fields = {
32 schemaField("eventId", eve::schema::ValueType::String),
33 schemaField("sequence", eve::schema::ValueType::String),
34 schemaField("type", eve::schema::ValueType::String),
35 schemaField("source", eve::schema::ValueType::String),
36 schemaField("subject", eve::schema::ValueType::String),
37 schemaField("causationKind", eve::schema::ValueType::String),
38 schemaField("causation", eve::schema::ValueType::String),
39 schemaField("correlationKind", eve::schema::ValueType::String),
40 schemaField("correlation", eve::schema::ValueType::String),
41 schemaField("schemaId", eve::schema::ValueType::String),
42 schemaField("schemaVersion", eve::schema::ValueType::String),
43 schemaField("tick", eve::schema::ValueType::String),
44 schemaField("flags", eve::schema::ValueType::Integer),
45 schemaField("payload", eve::schema::ValueType::Any),
46 };
49 eventNode.objectSchema = std::make_shared<const eve::schema::SchemaDefinition>(std::move(event));
50 auto events = schemaField("events", eve::schema::ValueType::Array);
51 events.itemSchema = std::make_shared<const eve::schema::SchemaNode>(std::move(eventNode));
52
54 schema.id = "game_event:stream";
55 schema.version = 1;
56 schema.title = "Game Event Stream Snapshot";
57 schema.description = "Current persistent payload owned by GameEventLog.";
58 schema.additionalProperties = false;
59 schema.fields = {
60 schemaField("version", eve::schema::ValueType::Integer),
61 schemaField("revision", eve::schema::ValueType::String),
62 schemaField("nextSequence", eve::schema::ValueType::String),
63 std::move(events),
64 };
65 return schema;
66}
67
68eve::Result<void> validateCurrentEventStream(std::string_view json) {
69 constexpr auto id = "game_event:stream";
71 auto registration = eve::schema::SchemaRegistry::registerVersioned(eventStreamSchemaDefinition());
72 if (!registration.ok()) return eve::Result<void>::failure(registration.status());
73 }
74 const auto errors = eve::schema::SchemaRegistry::validate(id, 1, std::string(json));
75 if (errors.empty()) return eve::Result<void>::success();
76 const auto& error = errors.front();
79 eve::DiagnosticDetails{{"schemaId", id}, {"schemaVersion", "1"}, {"validationCode", error.code}},
80 "game_event.snapshot.schema"));
81}
82
83std::string quote(const std::string& value) {
84 std::ostringstream out;
85 out << '"';
86 for (unsigned char c : value) {
87 switch (c) {
88 case '"': out << "\\\""; break;
89 case '\\': out << "\\\\"; break;
90 case '\b': out << "\\b"; break;
91 case '\f': out << "\\f"; break;
92 case '\n': out << "\\n"; break;
93 case '\r': out << "\\r"; break;
94 case '\t': out << "\\t"; break;
95 default:
96 if (c < 0x20)
97 out << "\\u" << std::hex << std::setw(4) << std::setfill('0') << int(c) << std::dec;
98 else
99 out << static_cast<char>(c);
100 }
101 }
102 return out.str() + '"';
103}
104
105bool parseU64(const eve::json::Value& value, uint64_t& out) {
106 if (!value.isString()) return false;
107 try {
108 size_t used = 0;
109 out = std::stoull(value.asString(), &used);
110 return used == value.asString().size();
111 } catch (...) {
112 return false;
113 }
114}
115
116bool parseI64(const eve::json::Value& value, int64_t& out) {
117 if (!value.isString()) return false;
118 try {
119 const std::string text = value.asString();
120 size_t used = 0;
121 out = std::stoll(text, &used);
122 return used == text.size();
123 } catch (...) {
124 return false;
125 }
126}
127
128std::string canonicalJson(const eve::json::Value& value) {
129 if (value.isNull()) return "null";
130 if (value.isBool()) return value.asBool() ? "true" : "false";
131 if (value.isNumber()) return value.asString();
132 if (value.isString()) return quote(value.asString());
133 if (value.isArray()) {
134 std::string out = "[";
135 for (size_t i = 0; i < value.size(); ++i) {
136 if (i) out += ',';
137 out += canonicalJson(value.at(i));
138 }
139 return out + ']';
140 }
141 if (value.isObject()) {
142 auto keys = value.keys();
143 std::sort(keys.begin(), keys.end());
144 std::string out = "{";
145 for (size_t i = 0; i < keys.size(); ++i) {
146 if (i) out += ',';
147 out += quote(keys[i]) + ':' + canonicalJson(value.get(keys[i].c_str()));
148 }
149 return out + '}';
150 }
151 return {};
152}
153
154const char* correlationKindName(CorrelationId::Kind kind) {
155 switch (kind) {
156 case CorrelationId::Kind::None: return "none";
157 case CorrelationId::Kind::Id: return "id";
158 case CorrelationId::Kind::Legacy: return "legacy";
159 }
160 return "none";
161}
162
163const char* causationKindName(CausationRef::Kind kind) {
164 switch (kind) {
165 case CausationRef::Kind::None: return "none";
166 case CausationRef::Kind::Event: return "event";
167 case CausationRef::Kind::Command: return "command";
168 case CausationRef::Kind::Legacy: return "legacy";
169 }
170 return "none";
171}
172
173CorrelationId correlationFromSnapshot(const std::string& kind, const std::string& text) {
174 if (kind == "none") return {};
175 if (kind == "id") {
176 if (const auto parsed = EventId::parse(text)) return CorrelationId::fromEventId(*parsed);
177 return {};
178 }
180}
181
182CausationRef causationFromSnapshot(const std::string& kind, const std::string& text) {
183 if (kind == "none") return {};
184 if (kind == "event") {
185 if (const auto parsed = EventId::parse(text)) return CausationRef::fromEventId(*parsed);
186 return {};
187 }
188 if (kind == "command") {
189 if (const auto parsed = CommandId::parse(text)) return CausationRef::fromCommandId(*parsed);
190 return {};
191 }
193}
194
195eve::LogicalId eventStreamSchema() {
196 const auto schema = eve::LogicalId::parse("game_event:stream");
197 if (!schema) std::terminate();
198 return *schema;
199}
200
201const eve::SnapshotMigrationChain& eventStreamMigrations() {
202 static const eve::SnapshotMigrationChain chain = [] {
204 const auto registration =
205 result.add(eventStreamSchema(), eve::SchemaVersion(0), eve::SchemaVersion(1),
207 const auto* object = payload.getIf<eve::Value::Object>();
208 if (!object)
210 eve::DiagnosticCode::ParseError, "event stream payload must be an object"));
212 });
213 if (!registration.ok()) std::terminate();
214 return result;
215 }();
216 return chain;
217}
218
219} // namespace
220
221EventConsumer::EventConsumer(GameEventLog* stream, EventSequence sequence)
222 : stream_(stream), nextSequence_(sequence.value() == 0 ? EventSequence(1) : sequence) {}
223
224int EventConsumer::read(int maxCount) {
225 batch_.clear();
226 if (!stream_ || maxCount <= 0) return 0;
227 if (nextSequence_.value() < stream_->oldestSequence().value()) nextSequence_ = stream_->oldestSequence();
228 for (const auto& event : stream_->events_) {
229 if (event.sequence.value() < nextSequence_.value()) continue;
230 if (static_cast<int>(batch_.size()) >= maxCount) break;
231 batch_.push_back(&event);
232 nextSequence_ = EventSequence(event.sequence.value() + 1);
233 }
234 return static_cast<int>(batch_.size());
235}
236
237int EventConsumer::batchCount() const { return static_cast<int>(batch_.size()); }
238
239const GameEvent* EventConsumer::batchAt(int index) const {
240 return index >= 0 && static_cast<size_t>(index) < batch_.size() ? batch_[static_cast<size_t>(index)] : nullptr;
241}
242
243EventSequence EventConsumer::position() const { return nextSequence_; }
244
245void EventConsumer::seek(EventSequence sequence) {
246 nextSequence_ = sequence.value() == 0 ? EventSequence(1) : sequence;
247 batch_.clear();
248}
249
250void EventConsumer::seek(uint64_t sequence) { seek(EventSequence(sequence)); }
251
252GameEventLog::GameEventLog(eve::UuidEntropySource entropy, eve::UuidClock clock)
253 : eventIdGenerator_(std::in_place, std::move(entropy), std::move(clock)) {}
254
256 const auto failure = [](eve::DiagnosticCode code, std::string message) {
258 };
259 if (envelope.type.empty()) return failure(eve::DiagnosticCode::InvalidArgument, "event type must not be empty");
260 if (!envelope.schemaId.isValid())
261 return failure(eve::DiagnosticCode::InvalidArgument, "event schemaId must be a valid logical ID");
262 if (envelope.schemaVersion.isZero())
263 return failure(eve::DiagnosticCode::InvalidArgument, "event schemaVersion must be non-zero");
264 auto payloadDocument = eve::json::Document::parse(envelope.payload);
265 if (!payloadDocument.valid()) return failure(eve::DiagnosticCode::ParseError, "payload must be valid JSON");
266 const auto next = nextSequence_.incremented();
267 if (!next) return failure(eve::DiagnosticCode::Failed, "event sequence exhausted");
268 if (!revision_.incremented()) return failure(eve::DiagnosticCode::Failed, "event stream revision exhausted");
269
270 if (envelope.eventId.isNil()) {
271 if (eventIdGenerator_) {
272 const auto generated = eventIdGenerator_->generate();
273 if (!generated) return failure(eve::DiagnosticCode::Failed, "event ID generation failed");
274 envelope.eventId = EventId::fromUuid(*generated);
275 } else {
277 "eventId is required when no UUIDv7 generator is configured");
278 }
279 }
280 if (!envelope.eventId.isNil()) {
281 for (const auto& existing : events_) {
282 if (existing.eventId == envelope.eventId)
283 return failure(eve::DiagnosticCode::Conflict, "eventId already exists in this stream");
284 }
285 }
286
287 envelope.sequence = nextSequence_;
288 envelope.payload = canonicalJson(payloadDocument.root());
289 events_.push_back(std::move(envelope));
290 nextSequence_ = *next;
291 revision_ = *revision_.incremented();
292 if (events_.back().tick > snapshotTick_) snapshotTick_ = events_.back().tick;
293 return eve::Result<EventSequence>::success(events_.back().sequence);
294}
295
297 const auto it =
298 std::lower_bound(events_.begin(), events_.end(), sequence.value(),
299 [](const GameEvent& event, uint64_t value) { return event.sequence.value() < value; });
300 return it != events_.end() && it->sequence.value() == sequence.value() ? &*it : nullptr;
301}
302
304
305int GameEventLog::queryType(const std::string& type) {
306 query_.clear();
307 for (const auto& event : events_)
308 if (event.type == type) query_.push_back(&event);
309 return static_cast<int>(query_.size());
310}
311
312int GameEventLog::querySource(const std::string& source) {
313 query_.clear();
314 for (const auto& event : events_)
315 if (event.source == source) query_.push_back(&event);
316 return static_cast<int>(query_.size());
317}
318
319int GameEventLog::querySubject(const std::string& subject) {
320 query_.clear();
321 for (const auto& event : events_)
322 if (event.subject == subject) query_.push_back(&event);
323 return static_cast<int>(query_.size());
324}
325
327 query_.clear();
328 for (const auto& event : events_)
329 if (event.sequence.value() >= firstSequence.value()) query_.push_back(&event);
330 return static_cast<int>(query_.size());
331}
332
333int GameEventLog::querySequence(uint64_t firstSequence) { return querySequence(EventSequence(firstSequence)); }
334
336 return index >= 0 && static_cast<size_t>(index) < query_.size() ? query_[static_cast<size_t>(index)] : nullptr;
337}
338
339int GameEventLog::size() const { return static_cast<int>(events_.size()); }
340
342 return index >= 0 && static_cast<size_t>(index) < events_.size() ? &events_[static_cast<size_t>(index)] : nullptr;
343}
344
346 consumers_.push_back(std::unique_ptr<EventConsumer>(new EventConsumer(this, firstSequence)));
347 return consumers_.back().get();
348}
349
350EventConsumer* GameEventLog::newConsumer(uint64_t firstSequence) { return newConsumer(EventSequence(firstSequence)); }
351
353 while (!events_.empty() && events_.front().sequence.value() < firstSequence.value()) events_.pop_front();
354 query_.clear();
355 for (auto& consumer : consumers_) consumer->batch_.clear();
356}
357
358void GameEventLog::clearBefore(uint64_t firstSequence) { clearBefore(EventSequence(firstSequence)); }
359
361 events_.clear();
362 query_.clear();
363 for (auto& consumer : consumers_) consumer->batch_.clear();
364}
365
367 clear();
368 nextSequence_ = EventSequence(1);
369 revision_ = eve::Revision::zero();
370 snapshotTick_ = eve::SimulationTick::zero();
371 for (auto& consumer : consumers_) consumer->seek(1);
372}
373
374std::string GameEventLog::snapshotJson() const {
375 std::ostringstream out;
376 out << "{\"version\":2,\"revision\":" << quote(std::to_string(revision_.value()))
377 << ",\"nextSequence\":" << quote(std::to_string(nextSequence_.value())) << ",\"events\":[";
378 bool first = true;
379 for (const auto& event : events_) {
380 if (!first) out << ',';
381 first = false;
382 out << "{\"eventId\":" << quote(event.eventId.format())
383 << ",\"sequence\":" << quote(std::to_string(event.sequence.value())) << ",\"type\":" << quote(event.type)
384 << ",\"source\":" << quote(event.source) << ",\"subject\":" << quote(event.subject)
385 << ",\"causationKind\":" << quote(causationKindName(event.causation.kind()))
386 << ",\"causation\":" << quote(event.causation.format())
387 << ",\"correlationKind\":" << quote(correlationKindName(event.correlation.kind()))
388 << ",\"correlation\":" << quote(event.correlation.format())
389 << ",\"schemaId\":" << quote(event.schemaId.format())
390 << ",\"schemaVersion\":" << quote(std::to_string(event.schemaVersion.value()))
391 << ",\"tick\":" << quote(std::to_string(event.tick.value())) << ",\"flags\":" << event.flags
392 << ",\"payload\":" << event.payload << '}';
393 }
394 return out.str() + "]}";
395}
396
398 std::string error;
399 auto document = eve::json::Document::parse(std::string(json), &error);
400 const auto root = document.root();
401 const auto invalid = [&](eve::DiagnosticCode code, std::string message) {
403 };
404 const int snapshotVersion = root.getInt("version");
405 if (!document.valid() || !root.isObject() || (snapshotVersion != 1 && snapshotVersion != 2)) {
406 return invalid(eve::DiagnosticCode::ParseError, error.empty() ? "invalid event stream snapshot" : error);
407 }
408 if (snapshotVersion == 2) {
409 auto schemaValidation = validateCurrentEventStream(json);
410 if (!schemaValidation.ok()) return schemaValidation;
411 }
412 uint64_t restoredNext = 0;
413 if (!parseU64(root.get("nextSequence"), restoredNext) || restoredNext == 0) {
414 return invalid(eve::DiagnosticCode::ParseError, "invalid nextSequence");
415 }
416 uint64_t restoredRevision = restoredNext - 1;
417 if (!root.get("revision").isNull() && !parseU64(root.get("revision"), restoredRevision))
418 return invalid(eve::DiagnosticCode::ParseError, "invalid revision");
419 const auto values = root.get("events");
420 if (!values.isArray()) {
421 return invalid(eve::DiagnosticCode::ParseError, "events must be an array");
422 }
423 std::deque<GameEvent> restored;
424 uint64_t previous = 0;
425 std::vector<std::string> restoredEventIds;
426 for (size_t i = 0; i < values.size(); ++i) {
427 const auto value = values.at(i);
428 GameEvent event;
429 int64_t restoredTick = 0;
430 uint64_t restoredSequence = 0;
431 if (!value.isObject() || !parseU64(value.get("sequence"), restoredSequence) || restoredSequence <= previous ||
432 !parseI64(value.get("tick"), restoredTick) || !value.get("flags").isNumber() || !value.get("payload")) {
433 return invalid(eve::DiagnosticCode::ParseError, "invalid event at index " + std::to_string(i));
434 }
435 event.sequence = EventSequence(restoredSequence);
436 if (restoredTick < 0) {
437 return invalid(eve::DiagnosticCode::ParseError, "invalid tick at index " + std::to_string(i));
438 }
439 event.tick = SimulationTick(static_cast<uint64_t>(restoredTick));
440 event.type = value.getString("type");
441 if (event.type.empty()) {
442 return invalid(eve::DiagnosticCode::ParseError, "event type must not be empty");
443 }
444 const std::string eventIdText = value.getString("eventId");
445 if (eventIdText.empty()) {
446 event.eventId = EventId::nil();
447 } else {
448 const auto parsedId = EventId::parse(eventIdText);
449 if (!parsedId)
450 return invalid(eve::DiagnosticCode::ParseError, "invalid eventId at index " + std::to_string(i));
451 event.eventId = *parsedId;
452 if (!event.eventId.isNil()) {
453 const auto formattedId = event.eventId.format();
454 if (std::find(restoredEventIds.begin(), restoredEventIds.end(), formattedId) != restoredEventIds.end())
455 return invalid(eve::DiagnosticCode::Conflict,
456 "duplicate non-nil eventId at index " + std::to_string(i));
457 restoredEventIds.push_back(formattedId);
458 }
459 }
460 event.source = value.getString("source");
461 event.subject = value.getString("subject");
462 const std::string causationText = value.getString("causation");
463 const std::string correlationText = value.getString("correlation");
464 const std::string causationKind =
465 snapshotVersion == 1 ? (causationText.empty() ? "none" : "legacy") : value.getString("causationKind");
466 const std::string correlationKind =
467 snapshotVersion == 1 ? (correlationText.empty() ? "none" : "legacy") : value.getString("correlationKind");
468 if (causationKind != "none" && causationKind != "event" && causationKind != "command" &&
469 causationKind != "legacy")
470 return invalid(eve::DiagnosticCode::ParseError, "invalid causation kind at index " + std::to_string(i));
471 if (correlationKind != "none" && correlationKind != "id" && correlationKind != "legacy")
472 return invalid(eve::DiagnosticCode::ParseError, "invalid correlation kind at index " + std::to_string(i));
473 if ((causationKind == "none" && !causationText.empty()) || (causationKind != "none" && causationText.empty()) ||
474 (correlationKind == "none" && !correlationText.empty()) ||
475 (correlationKind != "none" && correlationText.empty()))
476 return invalid(eve::DiagnosticCode::ParseError,
477 "causal reference value/kind mismatch at index " + std::to_string(i));
478 event.causation = causationFromSnapshot(causationKind, causationText);
479 event.correlation = correlationFromSnapshot(correlationKind, correlationText);
480 if ((causationKind == "event" || causationKind == "command") &&
482 return invalid(eve::DiagnosticCode::ParseError, "invalid causation ID at index " + std::to_string(i));
483 if (correlationKind == "id" && event.correlation.kind() == CorrelationId::Kind::None)
484 return invalid(eve::DiagnosticCode::ParseError, "invalid correlation ID at index " + std::to_string(i));
485 const std::string schemaIdText = value.getString("schemaId");
486 if (schemaIdText.empty()) {
487 event.schemaId = eve::LogicalId::fromParts("game_event", event.type).value_or(eve::LogicalId());
488 } else {
489 const auto parsedSchema = eve::LogicalId::parse(schemaIdText);
490 if (!parsedSchema)
491 return invalid(eve::DiagnosticCode::ParseError, "invalid schemaId at index " + std::to_string(i));
492 event.schemaId = *parsedSchema;
493 }
494 uint64_t restoredSchemaVersion = snapshotVersion == 1 ? 1 : 0;
495 if (snapshotVersion == 2 &&
496 (!parseU64(value.get("schemaVersion"), restoredSchemaVersion) || restoredSchemaVersion == 0)) {
497 return invalid(eve::DiagnosticCode::ParseError, "invalid schemaVersion at index " + std::to_string(i));
498 }
499 event.schemaVersion = SchemaVersion(restoredSchemaVersion);
500 const double flags = value.getDouble("flags", -1.0);
501 if (flags < 0 || flags > std::numeric_limits<uint32_t>::max() || std::floor(flags) != flags) {
502 return invalid(eve::DiagnosticCode::ParseError, "invalid flags at index " + std::to_string(i));
503 }
504 event.flags = static_cast<uint32_t>(flags);
505 event.payload = canonicalJson(value.get("payload"));
506 previous = event.sequence.value();
507 restored.push_back(std::move(event));
508 }
509 if (!restored.empty() && restored.back().sequence.value() >= restoredNext) {
510 return invalid(eve::DiagnosticCode::ParseError, "nextSequence must exceed retained events");
511 }
512 events_ = std::move(restored);
513 nextSequence_ = EventSequence(restoredNext);
514 revision_ = eve::Revision(restoredRevision);
515 snapshotTick_ = eve::SimulationTick::zero();
516 for (const auto& event : events_)
517 if (event.tick > snapshotTick_) snapshotTick_ = event.tick;
518 query_.clear();
519 consumers_.clear();
521}
522
526 return eve::makeSnapshotEnvelope("game_event.stream", eventStreamSchema(), eve::SchemaVersion(1), instanceId_,
527 revision_, snapshotTick_, std::move(payload).takeValue(), hashProvider);
528}
529
531 const eve::SnapshotHashProvider& hashProvider) {
532 if (source.type != "game_event.stream" || source.schema != eventStreamSchema())
534 eve::DiagnosticCode::InvalidArgument, "snapshot does not belong to game_event::GameEventLog"));
535 if (!instanceId_.isNil() && source.instanceId != instanceId_)
537 eve::DiagnosticCode::Conflict, "snapshot instanceId does not match game_event::GameEventLog"));
538 auto migrated = eventStreamMigrations().migrate(source, eve::SchemaVersion(1), hashProvider);
539 if (!migrated.ok()) return eve::Result<void>::failure(migrated.status());
540 auto metadata = eve::validateSnapshotPayloadMetadata(migrated.value().payload, migrated.value().revision,
541 migrated.value().tick);
542 if (!metadata.ok()) return eve::Result<void>::failure(metadata.status());
543 auto payload = migrated.value().payload.toJson();
544 if (!payload.ok()) return eve::Result<void>::failure(payload.status());
545
546 GameEventLog candidate(instanceId_);
547 candidate.eventIdGenerator_ = eventIdGenerator_;
548 auto restored = candidate.restore(std::move(payload).takeValue());
549 if (!restored.ok()) return eve::Result<void>::failure(restored.status());
550 if (candidate.snapshotTick_ != migrated.value().tick)
552 eve::DiagnosticCode::Conflict, "event stream payload tick disagrees with snapshot envelope"));
553 candidate.instanceId_ = migrated.value().instanceId;
554 candidate.revision_ = migrated.value().revision;
555 candidate.snapshotTick_ = migrated.value().tick;
556 // The UUID generator is a runtime capability, not persisted state. Keep
557 // it across the atomic payload replacement and invalidate old cursors.
558 eventIdGenerator_ = std::move(candidate.eventIdGenerator_);
559 instanceId_ = candidate.instanceId_;
560 revision_ = candidate.revision_;
561 snapshotTick_ = candidate.snapshotTick_;
562 nextSequence_ = candidate.nextSequence_;
563 events_ = std::move(candidate.events_);
564 query_.clear();
565 consumers_.clear();
567}
568
570 auto value = snapshot(hashProvider);
571 if (!value.ok()) return eve::Result<std::string>::failure(value.status());
572 return std::move(value).andThen(
573 [](eve::SnapshotEnvelope&& envelope) { return eve::serializeSnapshotEnvelope(envelope); });
574}
575
577 const eve::SnapshotHashProvider& hashProvider) {
578 auto source = eve::parseSnapshotEnvelope(json, hashProvider);
579 if (!source.ok()) return eve::Result<void>::failure(source.status());
580 return restoreSnapshot(std::move(source).takeValue(), hashProvider);
581}
582
583EventSequence GameEventLog::oldestSequence() const {
584 return events_.empty() ? nextSequence_ : events_.front().sequence;
585}
586
588 auto* module = GameEventModule::create();
589 module->streams_.push_back(std::make_unique<GameEventLog>());
590 return module->streams_.back().get();
591}
592
594 GameEventModule::expose);
595GameEventModule* GameEventModule::create() {
596 auto* existing = ModuleManager::find(name);
597 if (existing) return static_cast<GameEventModule*>(existing);
598 auto* module = new GameEventModule();
600 return module;
601}
602const char* GameEventModule::name = "GameEvent";
603
604void GameEventModule::expose(ssq::Table& table) {
605 auto envelope = table.addClass<GameEvent>(
606 "GameEventRecord", std::function<GameEvent*()>([]() -> GameEvent* { return nullptr; }), false);
607 envelope.addFunc("getEventId", [](GameEvent* e) { return e ? e->eventId.format() : std::string{}; });
608 envelope.addFunc("getSequence",
609 [](GameEvent* e) { return e ? static_cast<int64_t>(e->sequence.value()) : int64_t{0}; });
610 envelope.addFunc("getType", [](GameEvent* e) { return e ? e->type : std::string{}; });
611 envelope.addFunc("getSource", [](GameEvent* e) { return e ? e->source : std::string{}; });
612 envelope.addFunc("getSubject", [](GameEvent* e) { return e ? e->subject : std::string{}; });
613 envelope.addFunc("getCausation", [](GameEvent* e) { return e ? e->causation.format() : std::string{}; });
614 envelope.addFunc("getCausationKind", [](GameEvent* e) {
615 return e ? std::string(causationKindName(e->causation.kind())) : std::string{};
616 });
617 envelope.addFunc("getCorrelation", [](GameEvent* e) { return e ? e->correlation.format() : std::string{}; });
618 envelope.addFunc("getCorrelationKind", [](GameEvent* e) {
619 return e ? std::string(correlationKindName(e->correlation.kind())) : std::string{};
620 });
621 envelope.addFunc("getSchemaId", [](GameEvent* e) { return e ? e->schemaId.format() : std::string{}; });
622 envelope.addFunc("getSchemaVersion",
623 [](GameEvent* e) { return e ? static_cast<int64_t>(e->schemaVersion.value()) : int64_t{0}; });
624 envelope.addFunc("getTick", [](GameEvent* e) { return e ? static_cast<int64_t>(e->tick.value()) : int64_t{0}; });
625 envelope.addFunc("getFlags", [](GameEvent* e) { return e ? static_cast<int64_t>(e->flags) : int64_t{0}; });
626 envelope.addFunc("getPayload", [](GameEvent* e) { return e ? e->payload : std::string{}; });
627
628 auto consumer = table.addClass<EventConsumer>(
629 "EventConsumer", std::function<EventConsumer*()>([]() -> EventConsumer* { return nullptr; }), false);
630 consumer.addFunc("read", &EventConsumer::read);
631 consumer.addFunc("batchCount", &EventConsumer::batchCount);
632 consumer.addFunc("batchAt", [](EventConsumer* c, int index) -> GameEvent* {
633 return c ? const_cast<GameEvent*>(c->batchAt(index)) : nullptr;
634 });
635 consumer.addFunc("position",
636 [](EventConsumer* c) { return c ? static_cast<int64_t>(c->position().value()) : int64_t{0}; });
637 consumer.addFunc("seek", [](EventConsumer* c, int64_t sequence) {
638 if (c) c->seek(sequence > 0 ? uint64_t(sequence) : 1);
639 });
640
641 auto stream = table.addClass<GameEventLog>(
642 "GameEventLog", std::function<GameEventLog*()>([]() -> GameEventLog* { return nullptr; }), false);
643 const HSQUIRRELVM vm = table.getHandle();
644 stream.addFunc("append", [vm](GameEventLog* s, const std::string& eventId, const std::string& type,
645 const std::string& source, const std::string& subject, const std::string& causation,
646 const std::string& correlation, int64_t tick, int64_t flags,
647 const std::string& payload) {
648 if (!s) {
650 vm,
652 "event stream must not be null", "stream")),
653 [](EventSequence) { return eve::Value::null(); });
654 }
655 if (tick < 0 || flags < 0 || static_cast<std::uint64_t>(flags) > std::numeric_limits<std::uint32_t>::max()) {
657 vm,
660 "event tick and flags must be non-negative and in range", "envelope")),
661 [](EventSequence) { return eve::Value::null(); });
662 }
663 GameEvent envelope;
664 if (!eventId.empty()) {
665 const auto parsed = EventId::parse(eventId);
666 if (!parsed) {
668 vm,
670 "eventId must be a UUID", "eventId")),
671 [](EventSequence) { return eve::Value::null(); });
672 }
673 envelope.eventId = *parsed;
674 }
675 envelope.type = type;
676 envelope.source = source;
677 envelope.subject = subject;
678 if (!causation.empty()) {
679 const auto parsed = EventId::parse(causation);
680 envelope.causation = parsed ? CausationRef::fromEventId(*parsed) : CausationRef::fromLegacy(causation);
681 }
682 if (!correlation.empty()) {
683 const auto parsed = EventId::parse(correlation);
684 envelope.correlation =
685 parsed ? CorrelationId::fromEventId(*parsed) : CorrelationId::fromLegacy(correlation);
686 }
687 const auto schema = eve::LogicalId::fromParts("game_event", type);
688 if (schema) envelope.schemaId = *schema;
689 envelope.schemaVersion = SchemaVersion(1);
690 envelope.tick = SimulationTick(static_cast<std::uint64_t>(tick));
691 envelope.flags = static_cast<std::uint32_t>(flags);
692 envelope.payload = payload;
693 return eve::script::projectResult(vm, s->append(std::move(envelope)), [](EventSequence sequence) {
694 return eve::Value::integer(static_cast<std::int64_t>(sequence.value()));
695 });
696 });
697 stream.addFunc("find", [](GameEventLog* s, int64_t sequence) -> GameEvent* {
698 return s && sequence > 0 ? const_cast<GameEvent*>(s->find(uint64_t(sequence))) : nullptr;
699 });
700 stream.addFunc("queryType", &GameEventLog::queryType);
701 stream.addFunc("querySource", &GameEventLog::querySource);
702 stream.addFunc("querySubject", &GameEventLog::querySubject);
703 stream.addFunc("querySequence", [](GameEventLog* s, int64_t sequence) {
704 return s ? s->querySequence(sequence > 0 ? uint64_t(sequence) : 1) : 0;
705 });
706 stream.addFunc("queryAt", [](GameEventLog* s, int index) -> GameEvent* {
707 return s ? const_cast<GameEvent*>(s->queryAt(index)) : nullptr;
708 });
709 stream.addFunc("size", &GameEventLog::size);
710 stream.addFunc("at", [](GameEventLog* s, int index) -> GameEvent* {
711 return s ? const_cast<GameEvent*>(s->at(index)) : nullptr;
712 });
713 stream.addFunc("newConsumer", [](GameEventLog* s, int64_t sequence) {
714 return s ? s->newConsumer(sequence > 0 ? uint64_t(sequence) : 1) : nullptr;
715 });
716 stream.addFunc("clearBefore", [](GameEventLog* s, int64_t sequence) {
717 if (s) s->clearBefore(sequence > 0 ? uint64_t(sequence) : 1);
718 });
719 stream.addFunc("clear", &GameEventLog::clear);
720 stream.addFunc("reset", &GameEventLog::reset);
721 stream.addFunc("snapshotJson", &GameEventLog::snapshotJson);
722 stream.addFunc("restore", [vm = table.getHandle()](GameEventLog* stream, const std::string& json) {
723 if (!stream)
726 "game event log must not be null")));
727 return eve::script::projectResult(vm, stream->restore(json));
728 });
729
730 auto cls = table.addClass(name, GameEventModule::create, false);
731 expose(cls);
732}
733
734void GameEventModule::expose(ssq::Class& cls) {
735 cls.addFunc("getName", &GameEventModule::getName);
736 cls.addFunc("newLog", [](GameEventModule*) { return GameEventModule::newLog(); });
737}
738
739} // namespace eve::game_event
double value
Value::Object payload
int root
Definition AnimSmr.cpp:119
int subject
Definition AnimSmr.cpp:163
const std::string & s
struct SQVM * HSQUIRRELVM
HSQUIRRELVM vm
Definition ECS.cpp:20
HSQOBJECT cls
Definition ECS.cpp:21
std::map< std::string, Var > values
std::string message
DiagnosticCode code
std::int32_t c
std::int32_t first
std::string text
TokenKind kind
char quote
bool required
std::string name
graphics::Canvas * previous
std::unique_ptr< gpgpu::Sequence > sequence
Definition OnnxGpgpu.cpp:43
std::string error
Definition Package.cpp:60
std::string string
The single Squirrel projection for common Result, Status and Value.
Battle::Events events
SimulationTick tick
uint32_t index
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 > fromParts(std::string_view namespaceName, std::string_view name)
Builds a logical ID from separate namespace and name components.
Definition Identity.cpp:48
static std::optional< LogicalId > parse(std::string_view text)
Parses a scoped logical name.
Definition Identity.cpp:36
bool isValid() const noexcept
Returns whether this value contains a valid namespace and name.
Definition Identity.h:433
Module *(* creator_t)()
Definition Module.h:57
static Module * find(const char *name)
Finds a live module instance by name, or nullptr. @ownership Borrowed; ModuleManager owns the returne...
Definition Module.cpp:175
static void insert(const char *name, Module *inst)
Stores a module instance under a name (used by Module_IMPL).
Definition Module.cpp:194
virtual std::string getName() const =0
Returns the name.
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 Value null()
Compatibility factory for a null value.
Definition Value.h:102
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
static constexpr Id128 nil() noexcept
Returns the all-zero nil value.
Definition Identity.h:142
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
constexpr std::optional< StrongUint64 > incremented() const noexcept
Returns the next value, or empty instead of unsigned wraparound.
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.
constexpr bool isZero() const noexcept
Returns whether this value is zero.
static CausationRef fromEventId(EventId value)
Creates a causation reference to an event.
Definition GameEvent.h:105
static CausationRef fromLegacy(std::string value)
Creates a compatibility-only arbitrary string causation.
Definition GameEvent.h:125
Kind kind() const noexcept
Returns the tagged causation kind.
Definition GameEvent.h:133
static CausationRef fromCommandId(CommandId value)
Creates a causation reference to a command.
Definition GameEvent.h:115
Kind kind() const noexcept
Returns the tagged correlation kind.
Definition GameEvent.h:74
static CorrelationId fromLegacy(std::string value)
Creates a compatibility-only arbitrary string correlation.
Definition GameEvent.h:66
static CorrelationId fromEventId(EventId value)
Creates a canonical correlation ID from an EventId.
Definition GameEvent.h:58
Independent sequence cursor used to consume an GameEventLog in batches.
Definition GameEvent.h:224
int read(int maxCount)
Reads at most maxCount events and advances past the returned batch.
int batchCount() const
Returns the number of events in the most recent batch.
Deterministic in-memory stream of topic-neutral event envelopes.
Definition GameEvent.h:261
const GameEvent * queryAt(int index) const
Returns one result from the latest query, or nullptr.
int queryType(const std::string &type)
Queries all retained events of a type in sequence order.
const GameEvent * at(int index) const
Returns a retained event by stream order, or nullptr.
eve::Result< eve::SnapshotEnvelope > snapshot(const eve::SnapshotHashProvider &hashProvider) const
Captures the stream payload in the common snapshot envelope.
void clearBefore(EventSequence firstSequence)
Removes retained events with sequence lower than firstSequence.
int size() const
Returns the number of retained events.
eve::Result< std::string > snapshotEnvelopeJson(const eve::SnapshotHashProvider &hashProvider) const
Serializes the common event-stream snapshot envelope.
eve::Result< EventSequence > append(GameEvent envelope)
Appends a typed-metadata envelope and assigns stream-local order.
const GameEvent * find(EventSequence sequence) const
Returns an event by exact sequence, or nullptr. @ownership Borrowed from this stream; callers must no...
EventConsumer * newConsumer(EventSequence firstSequence=EventSequence(1))
Creates an independently positioned consumer cursor.
void clear()
Removes retained events while preserving the next sequence number.
std::string snapshotJson() const
Exports stream state as deterministic compact JSON.
int querySource(const std::string &source)
Queries all retained events from a source in sequence order.
eve::Result< void > restoreSnapshotJson(std::string_view json, const eve::SnapshotHashProvider &hashProvider)
Parses and transactionally restores a common event-stream envelope.
int querySubject(const std::string &subject)
Queries all retained events concerning a subject in sequence order.
eve::Result< void > restoreSnapshot(const eve::SnapshotEnvelope &snapshot, const eve::SnapshotHashProvider &hashProvider)
Restores a verified or migrated stream envelope atomically.
int querySequence(EventSequence firstSequence)
Queries retained events whose sequence is at least firstSequence.
eve::Result< void > restore(std::string_view json)
Transactionally restores stream state from a snapshot.
void reset()
Removes all state and restarts sequence allocation at one.
Script module factory for generic GameEventLog objects.
Definition GameEvent.h:409
static GameEventLog * newLog()
Allocates a module-owned event stream. @ownership Borrowed from the manager-owned GameEvent module; t...
GameEventModule()=default
Constructs a GameEventModule.
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
static eve::Result< SchemaRegistrationStatus > registerVersioned(const SchemaDefinition &definition)
Registers a new exact (id, version) entry without replacing it.
static std::vector< ValidationError > validate(const std::string &schemaId, const std::string &json)
Validates JSON text against the highest registered schema version.
static const SchemaDefinition * resolve(const std::string &schemaId, int schemaVersion)
Resolves one exact (schemaId, schemaVersion) pair, or nullptr.
eve::SimulationTick SimulationTick
Definition GameEvent.h:29
eve::EventSequence EventSequence
Definition GameEvent.h:28
eve::SchemaVersion SchemaVersion
Definition GameEvent.h:31
ModuleRegister GameEventModule_register("GameEvent",(ModuleManager::creator_t)(GameEventModule::create), GameEventModule::expose)
std::unordered_map< std::string, SkillDefinition > & table()
Definition Skill.cpp:65
ValueType
JSON-compatible value kinds understood by a schema field.
Definition SchemaTypes.h:16
ssq::Table projectResult(HSQUIRRELVM vm, Result< void > &&result)
Consume and project a void native Result using the common schema.
std::vector< DiagnosticDetail > DiagnosticDetails
Owning collection of diagnostic details with stable insertion order.
Definition Diagnostic.h:88
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
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
detail::StrongUint64< detail::RevisionTag > Revision
Monotonic content/state revision used for optimistic-concurrency checks.
Definition Revision.h:13
detail::StrongUint64< detail::EventSequenceTag > EventSequence
Stream-local event ordering value; it is not a global event identity.
std::function< Result< ContentId >(std::string_view canonicalInput)> SnapshotHashProvider
Injected content-digest implementation used by snapshots.
Definition Snapshot.h:36
EVENGINE_API_FOUNDATION public API.
Definition Module.h:202
Stable outer format shared by persistence and cross-process snapshots.
Definition Snapshot.h:46
Envelope metadata shared by event domains.
Definition GameEvent.h:183
std::string type
Compatibility event name and diagnostic topic.
Definition GameEvent.h:189
EventSequence sequence
Sequence assigned by the owning stream; it is never global identity.
Definition GameEvent.h:187
EventId eventId
Stable event identity; nil is allowed only for legacy transient events.
Definition GameEvent.h:185
std::string payload
Canonical serialized payload at the stream persistence boundary.
Definition GameEvent.h:207
SchemaVersion schemaVersion
Version of the payload schema, not a stream or runtime generation.
Definition GameEvent.h:201
SchemaId schemaId
Logical schema identifier for the typed payload codec.
Definition GameEvent.h:199
CorrelationId correlation
Business workflow correlation reference.
Definition GameEvent.h:197
CausationRef causation
Direct causation reference, normally an event or command ID.
Definition GameEvent.h:195
Metadata and validation constraints for one object member.
Definition SchemaTypes.h:77
Versioned runtime schema for a JSON object.
Definition SchemaTypes.h:93
std::vector< FieldDefinition > fields
Definition SchemaTypes.h:99
Recursive node in the deliberately small Eve Schema language.
Definition SchemaTypes.h:47
ValueType type
The value kind accepted by this node.
Definition SchemaTypes.h:49
std::shared_ptr< const SchemaDefinition > objectSchema
Inline object shape for an object node.
Definition SchemaTypes.h:57