host_runtime.cpp 87 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027
  1. #include <algorithm>
  2. #include <array>
  3. #include <bit>
  4. #include <cstdarg>
  5. #include <cstdio>
  6. #include <limits>
  7. #include <new>
  8. #include <utility>
  9. #include "../../core/logging/log.h"
  10. #include "../../middleware/bap/activity_message/activity_entity_slot_request_parser.h"
  11. #include "../../middleware/bap/activity_message/client_authoritative_data.h"
  12. #include "../../state/activity/mission/runtime.h"
  13. #include "../../state/activity/runtime.h"
  14. #include "../../state/activity_sdk/format.h"
  15. #include "../../state/activity_sdk/runtime.h"
  16. #include "host_runtime_internal.h"
  17. namespace sunrise::server::activity::host {
  18. namespace detail {
  19. struct MissionInputRecord final {
  20. MissionInputEvent view{};
  21. middleware::bap::activity_message::sense_update::DecodedPacket sense{};
  22. ClientMessageSnapshot clientMessage{};
  23. bool hasSense{};
  24. bool hasClientMessage{};
  25. };
  26. SRWLOCK g_lock{SRWLOCK_INIT};
  27. std::array<Instance, kInstanceCapacity> g_instances{};
  28. std::vector<PendingInput> g_pending{};
  29. std::array<Event, kEventCapacity> g_events{};
  30. std::vector<MissionInputRecord> g_missionInputs{};
  31. std::array<IncidentRecord, kIncidentHistoryCapacity> g_incidents{};
  32. std::array<ClientMessageRecord, kClientMessageHistoryCapacity> g_clientMessages{};
  33. std::array<ClientMessageDetail, kClientMessageDetailCapacity> g_clientMessageDetails{};
  34. std::size_t g_pendingRead{};
  35. std::size_t g_queuedIngress{};
  36. std::size_t g_queuedControls{};
  37. std::size_t g_eventStart{};
  38. std::size_t g_eventCount{};
  39. std::size_t g_clientMessageStart{};
  40. std::size_t g_clientMessageCount{};
  41. std::size_t g_clientMessageDetailStart{};
  42. std::size_t g_clientMessageDetailCount{};
  43. std::uint64_t g_sequence{};
  44. std::uint64_t g_scriptableReservationGeneration{1};
  45. std::uint64_t g_scriptableReservationSequence{};
  46. std::uint64_t g_missionInputSequence{};
  47. std::uint64_t g_eventGeneration{1};
  48. std::uint64_t g_clientMessageSequence{};
  49. std::uint64_t g_touch{};
  50. std::uint64_t g_droppedIngress{};
  51. std::uint64_t g_droppedIncidents{};
  52. std::uint64_t g_refusedControls{};
  53. std::uint64_t g_refusedIncidents{};
  54. std::uint64_t g_overwrittenEvents{};
  55. std::uint64_t g_overwrittenIncidents{};
  56. std::uint64_t g_overwrittenClientMessages{};
  57. /** @return True when the value stays inside the client's jump table. */
  58. [[nodiscard]] bool lifetime_allowed(std::uint8_t value) noexcept {
  59. return value <= kMaximumLifetimeState;
  60. }
  61. /** @return True when the shared codec accepts every bounded outer incident field. */
  62. [[nodiscard]] bool
  63. incident_allowed(const middleware::bap::activity_message::incident::Incident& incident) noexcept {
  64. return middleware::bap::activity_message::incident::outer_valid(incident);
  65. }
  66. /** Advances a diagnostic counter without making zero look like no event. */
  67. [[nodiscard]] std::uint64_t next_nonzero(std::uint64_t value) noexcept {
  68. return value == (std::numeric_limits<std::uint64_t>::max)() ? 1 : value + 1;
  69. }
  70. /** Finds one exact instance while the runtime lock is held. */
  71. [[nodiscard]] Instance* find_instance(const state::activity::SessionBinding& binding) noexcept {
  72. for (Instance& instance : g_instances) {
  73. if (instance.occupied && same_binding(instance.view.binding, binding)) {
  74. return &instance;
  75. }
  76. }
  77. return nullptr;
  78. }
  79. /** @return True when this exact binding already has an operator request waiting to reduce. */
  80. [[nodiscard]] bool has_queued_control(const state::activity::SessionBinding& binding) noexcept {
  81. for (std::size_t index = g_pendingRead; index < g_pending.size(); ++index) {
  82. const PendingInput& pending = g_pending[index];
  83. if (pending.kind == PendingKind::authControl
  84. && same_binding(pending.control.binding, binding)) {
  85. return true;
  86. }
  87. if (pending.kind == PendingKind::incidentControl
  88. && same_binding(pending.incidentControl.binding, binding)) {
  89. return true;
  90. }
  91. if (pending.kind == PendingKind::scriptableControl
  92. && same_binding(pending.scriptableControl.binding, binding)) {
  93. return true;
  94. }
  95. }
  96. return false;
  97. }
  98. /** A refused input is the one loss the host cannot see later, so it says so when it happens. */
  99. void report_ingress_drop(PendingKind kind, const char* reason) noexcept {
  100. std::array<char, core::log::kLineCapacity> line{};
  101. const int written =
  102. std::snprintf(line.data(),
  103. line.size(),
  104. "ev=activity stage=ingress result=drop reason=%s kind=%u queued=%zu",
  105. reason,
  106. static_cast<unsigned>(kind),
  107. g_pending.size());
  108. if (written <= 0) {
  109. return;
  110. }
  111. const auto length = static_cast<std::size_t>(written) < line.size()
  112. ? static_cast<std::size_t>(written)
  113. : line.size() - 1;
  114. core::log::write(core::log::Channel::server, core::log::Level::debug, {line.data(), length});
  115. }
  116. /** Appends one owned reducer row while the runtime lock is held. @return False when full. */
  117. bool append_pending(PendingInput&& pending) noexcept {
  118. if (g_pending.size() == g_pending.max_size()) {
  119. report_ingress_drop(pending.kind, "queue_full");
  120. return false;
  121. }
  122. try {
  123. g_pending.push_back(std::move(pending));
  124. } catch (const std::bad_alloc&) {
  125. report_ingress_drop(pending.kind, "no_memory");
  126. return false;
  127. }
  128. return true;
  129. }
  130. /** Allocates one diagnostic instance without evicting a binding active in this service slice. */
  131. [[nodiscard]] Instance* ensure_instance(const state::activity::SessionBinding& binding) noexcept {
  132. if (Instance* const current = find_instance(binding); current != nullptr) {
  133. return current;
  134. }
  135. Instance* selected = nullptr;
  136. for (Instance& instance : g_instances) {
  137. if (!instance.occupied) {
  138. selected = &instance;
  139. break;
  140. }
  141. if (!instance.view.active && !instance.view.outputPending
  142. && !instance.view.scriptableReservationPending && instance.view.incidentsPending == 0
  143. && (selected == nullptr || instance.lastTouched < selected->lastTouched)) {
  144. selected = &instance;
  145. }
  146. }
  147. if (selected == nullptr) {
  148. return nullptr;
  149. }
  150. clear_instance(*selected);
  151. selected->occupied = true;
  152. selected->view.binding = binding;
  153. selected->view.lifetimeState = kDefaultLifetimeState;
  154. return selected;
  155. }
  156. /** Finds one retained outbound incident revision while the runtime lock is held. */
  157. [[nodiscard]] IncidentRecord* find_incident(const state::activity::SessionBinding& binding,
  158. std::uint64_t revision) noexcept {
  159. for (IncidentRecord& record : g_incidents) {
  160. if (record.sequence != 0 && record.outbound && record.revision == revision
  161. && same_binding(record.binding, binding)) {
  162. return &record;
  163. }
  164. }
  165. return nullptr;
  166. }
  167. /** Reserves a full incident record without discarding an unstaged operator event. */
  168. [[nodiscard]] IncidentRecord* reserve_incident_record() noexcept {
  169. IncidentRecord* selected = nullptr;
  170. for (IncidentRecord& record : g_incidents) {
  171. if (record.sequence == 0) {
  172. return &record;
  173. }
  174. const bool evictable = !record.outbound || record.transportStages != 0
  175. || record.status == IncidentStatus::canceled;
  176. if (evictable && (selected == nullptr || record.sequence < selected->sequence)) {
  177. selected = &record;
  178. }
  179. }
  180. if (selected != nullptr) {
  181. ++g_overwrittenIncidents;
  182. }
  183. return selected;
  184. }
  185. /** Copies one incident summary into the chronological event history. */
  186. void fill_incident_event(Event& event,
  187. const middleware::bap::activity_message::incident::Incident& incident,
  188. std::uint64_t revision) noexcept {
  189. event.incidentRevision = revision;
  190. event.incidentTarget = incident.primaryTarget;
  191. event.incidentExtraTargets = incident.extraTargetCount;
  192. event.incidentSelectorBytes = incident.selectorLength;
  193. event.incidentPayloadBytes = incident.payloadLength;
  194. }
  195. /** Moves one instance to the newest eviction position. */
  196. void touch(Instance& instance) noexcept {
  197. g_touch = next_nonzero(g_touch);
  198. instance.lastTouched = g_touch;
  199. }
  200. /** @return True when this event kind enters the ordered mission-input feed. */
  201. [[nodiscard]] bool mission_input_kind(EventKind kind) noexcept {
  202. return kind == EventKind::senseUpdate || kind == EventKind::incidentReceived
  203. || kind == EventKind::clientStateChanged || kind == EventKind::entitySlotsRequested
  204. || kind == EventKind::clientMessageReceived;
  205. }
  206. /**
  207. * @return True when the feed reserved room for one more accepted row. The feed holds every
  208. * accepted row until its program commits it, so it has no row count of its own. Only the
  209. * allocator can refuse.
  210. */
  211. [[nodiscard]] bool reserve_mission_input_slot() noexcept {
  212. if (g_missionInputs.size() < g_missionInputs.capacity()) {
  213. return true;
  214. }
  215. // Grow a page at a time so the append that spends a durable sequence cannot reallocate.
  216. try {
  217. g_missionInputs.reserve(g_missionInputs.capacity() + kMissionInputReadPageSize);
  218. } catch (const std::bad_alloc&) {
  219. return false;
  220. }
  221. return true;
  222. }
  223. /** Reports one client input the mission-input feed could not store. */
  224. void report_mission_input_refusal(const Event& event) noexcept {
  225. std::array<char, core::log::kLineCapacity> line{};
  226. const int written =
  227. std::snprintf(line.data(),
  228. line.size(),
  229. "ev=mission_input result=allocation_refused session=0x%llX binding_rev=%llu "
  230. "kind=%u retained=%zu",
  231. static_cast<unsigned long long>(event.binding.sessionId),
  232. static_cast<unsigned long long>(event.binding.createdRevision),
  233. static_cast<unsigned>(event.kind),
  234. g_missionInputs.size());
  235. if (written > 0) {
  236. core::log::write(
  237. core::log::Channel::server,
  238. core::log::Level::debug,
  239. {line.data(), (std::min)(static_cast<std::size_t>(written), line.size() - 1)});
  240. }
  241. }
  242. /** Assigns one exact binding's ordered client mission-input sequence. */
  243. void stamp_mission_sequence(Event& event) noexcept {
  244. if (!mission_input_kind(event.kind)) {
  245. return;
  246. }
  247. Instance* const instance = find_instance(event.binding);
  248. if (instance == nullptr) {
  249. return;
  250. }
  251. // The durable sequence is what makes a row owed, so refuse before spending it.
  252. if (!reserve_mission_input_slot()) {
  253. report_mission_input_refusal(event);
  254. return;
  255. }
  256. std::uint64_t sequence = 0;
  257. if (state::activity::mission::issue_input_sequence(event.binding, sequence)) {
  258. instance->missionSequence = sequence;
  259. event.missionSequence = sequence;
  260. }
  261. }
  262. /** Appends one event in oldest-to-newest ring order. */
  263. void append_event(Event& event) noexcept {
  264. stamp_mission_sequence(event);
  265. g_sequence = next_nonzero(g_sequence);
  266. event.sequence = g_sequence;
  267. std::size_t index = (g_eventStart + g_eventCount) % g_events.size();
  268. if (g_eventCount == g_events.size()) {
  269. index = g_eventStart;
  270. g_eventStart = (g_eventStart + 1) % g_events.size();
  271. ++g_overwrittenEvents;
  272. } else {
  273. ++g_eventCount;
  274. }
  275. g_events[index] = event;
  276. }
  277. /** Appends framing metadata without entering or draining the reducer queue. */
  278. [[nodiscard]] std::uint64_t append_client_message(ClientMessageRecord record) noexcept {
  279. g_clientMessageSequence = next_nonzero(g_clientMessageSequence);
  280. record.sequence = g_clientMessageSequence;
  281. std::size_t index = (g_clientMessageStart + g_clientMessageCount) % g_clientMessages.size();
  282. if (g_clientMessageCount == g_clientMessages.size()) {
  283. index = g_clientMessageStart;
  284. g_clientMessageStart = (g_clientMessageStart + 1) % g_clientMessages.size();
  285. ++g_overwrittenClientMessages;
  286. } else {
  287. ++g_clientMessageCount;
  288. }
  289. g_clientMessages[index] = record;
  290. return record.sequence;
  291. }
  292. /** Appends one bounded decode in oldest-to-newest ring order. */
  293. void append_client_message_detail(
  294. std::uint64_t sequence,
  295. std::uint32_t messageType,
  296. const middleware::bap::activity_message::sense_update::DecodedPacket* sense) noexcept {
  297. std::size_t index =
  298. (g_clientMessageDetailStart + g_clientMessageDetailCount) % g_clientMessageDetails.size();
  299. if (g_clientMessageDetailCount == g_clientMessageDetails.size()) {
  300. index = g_clientMessageDetailStart;
  301. g_clientMessageDetailStart =
  302. (g_clientMessageDetailStart + 1) % g_clientMessageDetails.size();
  303. } else {
  304. ++g_clientMessageDetailCount;
  305. }
  306. ClientMessageDetail& detail = g_clientMessageDetails[index];
  307. detail = {};
  308. detail.sequence = sequence;
  309. detail.messageType = messageType;
  310. if (sense != nullptr) {
  311. detail.sense = *sense;
  312. detail.hasSenseDecode = true;
  313. }
  314. }
  315. /** Cancels one committed output while the runtime lock is held. */
  316. void cancel_output(Instance& instance, std::uint64_t now) noexcept {
  317. if (!instance.view.outputPending) {
  318. return;
  319. }
  320. Event event{};
  321. event.binding = instance.view.binding;
  322. event.tick = now;
  323. if (instance.view.outputKind == OutputKind::incident) {
  324. IncidentRecord* const record =
  325. find_incident(instance.view.binding, instance.view.incidentRevision);
  326. if (record != nullptr) {
  327. record->status = IncidentStatus::canceled;
  328. fill_incident_event(event, record->incident, record->revision);
  329. }
  330. event.kind = EventKind::incidentCanceled;
  331. instance.view.incidentsPending = 0;
  332. } else if (instance.view.outputKind == OutputKind::scriptableOverride) {
  333. event.kind = EventKind::scriptableOverrideCanceled;
  334. event.scriptableRevision = instance.view.scriptableRevision;
  335. instance.pendingScriptable = {};
  336. } else {
  337. event.stateRevision = instance.view.stateRevision;
  338. event.lifetimeState = instance.view.lifetimeState;
  339. event.kind = EventKind::authStateCanceled;
  340. }
  341. instance.view.outputPending = false;
  342. instance.view.outputKind = OutputKind::none;
  343. instance.view.outputStatus = OutputStatus::canceled;
  344. append_event(event);
  345. instance.view.lastEventSequence = g_sequence;
  346. }
  347. /** @return True when one decoded object has the exact retained observation key. */
  348. [[nodiscard]] bool same_sense_key(
  349. const SenseObservationKey& key,
  350. const middleware::bap::activity_message::sense_update::DecodedObject& object) noexcept {
  351. return key.registryKey == object.registryKey && key.objectTag == object.objectTag
  352. && key.slotType == object.slotType && key.slotIndex == object.slotIndex
  353. && key.senseSchema == object.senseSchema && key.schemaRow == object.schemaRow;
  354. }
  355. /** @return True when the packet carries this key at or after the selected object. */
  356. [[nodiscard]] bool
  357. packet_has_sense_key(const middleware::bap::activity_message::sense_update::DecodedPacket& packet,
  358. const SenseObservationKey& key,
  359. std::size_t first) noexcept {
  360. for (std::size_t index = first; index < packet.objectCount; ++index) {
  361. if (same_sense_key(key, packet.objects[index])) {
  362. return true;
  363. }
  364. }
  365. return false;
  366. }
  367. /** @return True when every retained object and value came from one complete decode. */
  368. [[nodiscard]] bool complete_sense_observation_input(const SenseInput& input) noexcept {
  369. namespace sense = middleware::bap::activity_message::sense_update;
  370. const sense::DecodedPacket& packet = input.decoded;
  371. if (input.sourceGeneration == 0 || input.clientMessageSequence == 0
  372. || input.verdict != state::activity::receipts::Verdict::framed
  373. || input.decodeStatus != sense::DecodeStatus::complete
  374. || packet.status != sense::DecodeStatus::complete || packet.objectsTruncated
  375. || packet.valuesTruncated || packet.groupsSkipped != 0 || input.groupsSkipped != 0
  376. || packet.objectCount > packet.objects.size() || packet.valueCount > packet.values.size()
  377. || packet.objectsDecoded != packet.objectCount || packet.objectsSeen != packet.objectCount
  378. || input.groupsSeen != packet.groupsSeen || input.groupsDecoded != packet.groupsDecoded
  379. || input.objectsSeen != packet.objectsSeen
  380. || input.objectsDecoded != packet.objectsDecoded) {
  381. return false;
  382. }
  383. std::size_t expectedValue = 0;
  384. for (std::size_t index = 0; index < packet.objectCount; ++index) {
  385. const sense::DecodedObject& object = packet.objects[index];
  386. if (object.status != sense::ObjectStatus::decoded || !object.hasGeneration
  387. || object.firstValue != expectedValue
  388. || object.valueCount > packet.valueCount - expectedValue) {
  389. return false;
  390. }
  391. expectedValue += object.valueCount;
  392. }
  393. return expectedValue == packet.valueCount;
  394. }
  395. /** Retains one accepted client input independently from panel and output events. */
  396. void append_mission_input(
  397. const Event& event,
  398. const middleware::bap::activity_message::sense_update::DecodedPacket* sense,
  399. const ClientMessageSnapshot* clientMessage = nullptr) noexcept {
  400. if (event.missionSequence == 0) {
  401. return;
  402. }
  403. g_missionInputSequence = next_nonzero(g_missionInputSequence);
  404. MissionInputRecord record{};
  405. record.view.event = event;
  406. record.view.sequence = g_missionInputSequence;
  407. if (record.view.event.sequence == 0) {
  408. record.view.event.sequence = g_missionInputSequence;
  409. }
  410. if (sense != nullptr) {
  411. record.sense = *sense;
  412. record.hasSense = true;
  413. }
  414. if (clientMessage != nullptr) {
  415. record.clientMessage = *clientMessage;
  416. record.hasClientMessage = true;
  417. }
  418. // The slot was reserved when the durable sequence was issued, so this never reallocates.
  419. g_missionInputs.push_back(std::move(record));
  420. }
  421. /** Mixes one fixed-width value into the local scene change guard. */
  422. void mix_scene_fingerprint(std::uint64_t& fingerprint, std::uint64_t value) noexcept {
  423. constexpr std::uint64_t kPrime = 1'099'511'628'211ULL;
  424. for (std::uint8_t shift = 0; shift < 64; shift += 8) {
  425. fingerprint ^= (value >> shift) & 0xFFU;
  426. fingerprint *= kPrime;
  427. }
  428. }
  429. /** @return A run-local change guard over one decoded object and every retained value. */
  430. [[nodiscard]] std::uint64_t
  431. scene_fingerprint(const middleware::bap::activity_message::sense_update::DecodedObject& object,
  432. std::span<const middleware::bap::activity_message::sense_update::DecodedValue>
  433. values) noexcept {
  434. constexpr std::uint64_t kOffset = 14'695'981'039'346'656'037ULL;
  435. std::uint64_t fingerprint = kOffset;
  436. mix_scene_fingerprint(fingerprint, object.objectRow);
  437. mix_scene_fingerprint(fingerprint, object.slotRow);
  438. mix_scene_fingerprint(fingerprint, object.schemaRow);
  439. mix_scene_fingerprint(fingerprint, object.generationPlusOne);
  440. mix_scene_fingerprint(fingerprint, object.deltaBits);
  441. mix_scene_fingerprint(fingerprint, object.hasGeneration ? 1U : 0U);
  442. mix_scene_fingerprint(fingerprint, values.size());
  443. for (const auto& value : values) {
  444. mix_scene_fingerprint(fingerprint, value.unsignedValue);
  445. mix_scene_fingerprint(fingerprint, static_cast<std::uint64_t>(value.signedValue));
  446. mix_scene_fingerprint(fingerprint, std::bit_cast<std::uint32_t>(value.realValue));
  447. mix_scene_fingerprint(fingerprint, value.schemaRow);
  448. mix_scene_fingerprint(fingerprint, value.fieldRow);
  449. mix_scene_fingerprint(fingerprint, value.occurrence);
  450. mix_scene_fingerprint(fingerprint, value.bitOffset);
  451. mix_scene_fingerprint(fingerprint, value.fieldOrdinal);
  452. mix_scene_fingerprint(fingerprint, value.width);
  453. mix_scene_fingerprint(fingerprint, static_cast<std::uint8_t>(value.kind));
  454. mix_scene_fingerprint(fingerprint, value.present ? 1U : 0U);
  455. }
  456. return fingerprint;
  457. }
  458. /** Writes one bounded type-43 diagnostic event. */
  459. void report_scene_sense(const char* format, ...) noexcept {
  460. std::array<char, core::log::kLineCapacity> line{};
  461. va_list arguments;
  462. va_start(arguments, format);
  463. const int written = std::vsnprintf(line.data(), line.size(), format, arguments);
  464. va_end(arguments);
  465. if (written <= 0) {
  466. return;
  467. }
  468. const auto length = static_cast<std::size_t>(written) < line.size()
  469. ? static_cast<std::size_t>(written)
  470. : line.size() - 1;
  471. core::log::write(core::log::Channel::server, core::log::Level::debug, {line.data(), length});
  472. }
  473. /** @return The trace row for one exact ClientRef, or null when it has not been seen. */
  474. [[nodiscard]] SceneSenseTraceRecord* find_scene_trace_record(
  475. SceneSenseTrace& trace,
  476. const middleware::bap::activity_message::sense_update::DecodedObject& object) noexcept {
  477. for (SceneSenseTraceRecord& record : trace.records) {
  478. if (record.occupied && same_sense_key(record.key, object)) {
  479. return &record;
  480. }
  481. }
  482. return nullptr;
  483. }
  484. /** @return One unused trace row, or null when the bounded table is full. */
  485. [[nodiscard]] SceneSenseTraceRecord* reserve_scene_trace_record(SceneSenseTrace& trace) noexcept {
  486. for (SceneSenseTraceRecord& record : trace.records) {
  487. if (!record.occupied) {
  488. return &record;
  489. }
  490. }
  491. return nullptr;
  492. }
  493. /** Reports one typed scalar with its exact reflected rows and wire position. */
  494. void report_scene_value(
  495. const Instance& instance,
  496. const SenseInput& input,
  497. const middleware::bap::activity_message::sense_update::DecodedObject& object,
  498. const middleware::bap::activity_message::sense_update::DecodedValue& value,
  499. std::size_t valueIndex) noexcept {
  500. using ValueKind = middleware::bap::activity_message::sense_update::ValueKind;
  501. constexpr const char* kPrefix =
  502. "ev=scene_sense kind=value activity=%d session=0x%llX binding_rev=%llu "
  503. "source_gen=%llu msg_seq=%llu key=0x%08X tag=0x%08X type=%u index=%u "
  504. "sense_schema=0x%08X gen_plus_one=%u value_index=%zu schema_row=%u "
  505. "field_row=%u ordinal=%u occurrence=%u bit=%u width=%u present=%u";
  506. std::array<char, core::log::kLineCapacity> prefix{};
  507. const int written =
  508. std::snprintf(prefix.data(),
  509. prefix.size(),
  510. kPrefix,
  511. static_cast<int>(instance.view.binding.destination.activityIndex),
  512. static_cast<unsigned long long>(instance.view.binding.sessionId),
  513. static_cast<unsigned long long>(instance.view.binding.createdRevision),
  514. static_cast<unsigned long long>(input.sourceGeneration),
  515. static_cast<unsigned long long>(input.clientMessageSequence),
  516. object.registryKey,
  517. object.objectTag,
  518. static_cast<unsigned>(object.slotType),
  519. static_cast<unsigned>(object.slotIndex),
  520. object.senseSchema,
  521. object.generationPlusOne,
  522. valueIndex,
  523. value.schemaRow,
  524. value.fieldRow,
  525. static_cast<unsigned>(value.fieldOrdinal),
  526. value.occurrence,
  527. value.bitOffset,
  528. static_cast<unsigned>(value.width),
  529. value.present ? 1U : 0U);
  530. if (written <= 0 || static_cast<std::size_t>(written) >= prefix.size()) {
  531. return;
  532. }
  533. if (!value.present) {
  534. report_scene_sense("%s domain=%s value=absent",
  535. prefix.data(),
  536. value.kind == ValueKind::unsignedInteger ? "uint"
  537. : value.kind == ValueKind::signedInteger ? "int"
  538. : value.kind == ValueKind::boolean ? "bool"
  539. : "real32");
  540. return;
  541. }
  542. switch (value.kind) {
  543. case ValueKind::unsignedInteger:
  544. report_scene_sense("%s domain=uint value=0x%llX",
  545. prefix.data(),
  546. static_cast<unsigned long long>(value.unsignedValue));
  547. break;
  548. case ValueKind::signedInteger:
  549. report_scene_sense("%s domain=int raw=0x%llX value=%lld",
  550. prefix.data(),
  551. static_cast<unsigned long long>(value.unsignedValue),
  552. static_cast<long long>(value.signedValue));
  553. break;
  554. case ValueKind::boolean:
  555. report_scene_sense("%s domain=bool raw=0x%llX value=%u",
  556. prefix.data(),
  557. static_cast<unsigned long long>(value.unsignedValue),
  558. value.unsignedValue != 0 ? 1U : 0U);
  559. break;
  560. case ValueKind::real32:
  561. report_scene_sense("%s domain=real32 raw=0x%llX value_bits=0x%08X value=%.9g",
  562. prefix.data(),
  563. static_cast<unsigned long long>(value.unsignedValue),
  564. std::bit_cast<std::uint32_t>(value.realValue),
  565. static_cast<double>(value.realValue));
  566. break;
  567. }
  568. }
  569. /** Reports complete changed type-43 objects without changing retained activity state. */
  570. void trace_scene_sense(Instance& instance, const SenseInput& input) noexcept {
  571. namespace sense = middleware::bap::activity_message::sense_update;
  572. if (!core::log::accepts(core::log::Channel::server, core::log::Level::debug)) {
  573. return;
  574. }
  575. SceneSenseTrace& trace = instance.sceneSenseTrace;
  576. if (trace.sourceGeneration != input.sourceGeneration) {
  577. trace = {};
  578. trace.sourceGeneration = input.sourceGeneration;
  579. }
  580. const sense::DecodedPacket& packet = input.decoded;
  581. bool hasScene = false;
  582. const std::size_t retainedObjectCount = (std::min)(packet.objectCount, packet.objects.size());
  583. for (std::size_t index = 0; index < retainedObjectCount; ++index) {
  584. if (packet.objects[index].slotType
  585. == static_cast<std::uint8_t>(state::activity_sdk::format::kAuthoredSceneSlotType)) {
  586. hasScene = true;
  587. break;
  588. }
  589. }
  590. if (!complete_sense_observation_input(input)) {
  591. if (!trace.incompleteReported && (hasScene || packet.objectsTruncated)) {
  592. trace.incompleteReported = true;
  593. report_scene_sense(
  594. "ev=scene_sense kind=packet result=skip reason=incomplete activity=%d "
  595. "session=0x%llX binding_rev=%llu source_gen=%llu msg_seq=%llu "
  596. "status=%s objects_seen=%u objects_kept=%zu values_kept=%zu "
  597. "objects_truncated=%u values_truncated=%u",
  598. static_cast<int>(instance.view.binding.destination.activityIndex),
  599. static_cast<unsigned long long>(instance.view.binding.sessionId),
  600. static_cast<unsigned long long>(instance.view.binding.createdRevision),
  601. static_cast<unsigned long long>(input.sourceGeneration),
  602. static_cast<unsigned long long>(input.clientMessageSequence),
  603. sense::decode_status_name(packet.status),
  604. packet.objectsSeen,
  605. packet.objectCount,
  606. packet.valueCount,
  607. packet.objectsTruncated ? 1U : 0U,
  608. packet.valuesTruncated ? 1U : 0U);
  609. }
  610. return;
  611. }
  612. for (std::size_t index = 0; index < packet.objectCount; ++index) {
  613. const sense::DecodedObject& object = packet.objects[index];
  614. if (object.slotType
  615. != static_cast<std::uint8_t>(state::activity_sdk::format::kAuthoredSceneSlotType)
  616. || packet_has_sense_key(packet,
  617. {object.registryKey,
  618. object.objectTag,
  619. object.senseSchema,
  620. object.schemaRow,
  621. object.slotIndex,
  622. object.slotType},
  623. index + 1)) {
  624. continue;
  625. }
  626. const std::span values(packet.values.data() + object.firstValue, object.valueCount);
  627. const std::uint64_t fingerprint = scene_fingerprint(object, values);
  628. SceneSenseTraceRecord* record = find_scene_trace_record(trace, object);
  629. const bool known = record != nullptr;
  630. if (known && record->fingerprint == fingerprint
  631. && record->generationPlusOne == object.generationPlusOne
  632. && record->valueCount == object.valueCount
  633. && record->hasGeneration == object.hasGeneration) {
  634. continue;
  635. }
  636. if (record == nullptr) {
  637. record = reserve_scene_trace_record(trace);
  638. }
  639. if (record == nullptr) {
  640. if (!trace.capacityReported) {
  641. trace.capacityReported = true;
  642. report_scene_sense(
  643. "ev=scene_sense kind=packet result=skip reason=trace_capacity activity=%d "
  644. "session=0x%llX binding_rev=%llu source_gen=%llu capacity=%zu",
  645. static_cast<int>(instance.view.binding.destination.activityIndex),
  646. static_cast<unsigned long long>(instance.view.binding.sessionId),
  647. static_cast<unsigned long long>(instance.view.binding.createdRevision),
  648. static_cast<unsigned long long>(input.sourceGeneration),
  649. trace.records.size());
  650. }
  651. continue;
  652. }
  653. record->key = {object.registryKey,
  654. object.objectTag,
  655. object.senseSchema,
  656. object.schemaRow,
  657. object.slotIndex,
  658. object.slotType};
  659. record->fingerprint = fingerprint;
  660. record->generationPlusOne = object.generationPlusOne;
  661. record->valueCount = object.valueCount;
  662. record->hasGeneration = object.hasGeneration;
  663. record->occupied = true;
  664. report_scene_sense("ev=scene_sense kind=object change=%s activity=%d session=0x%llX "
  665. "binding_rev=%llu source_gen=%llu msg_seq=%llu key=0x%08X tag=0x%08X "
  666. "object_row=%u type=%u index=%u slot_row=%u sense_schema=0x%08X "
  667. "schema_row=%u gen_plus_one=%u has_gen=%u delta_bits=%u values=%u",
  668. known ? "update" : "new",
  669. static_cast<int>(instance.view.binding.destination.activityIndex),
  670. static_cast<unsigned long long>(instance.view.binding.sessionId),
  671. static_cast<unsigned long long>(instance.view.binding.createdRevision),
  672. static_cast<unsigned long long>(input.sourceGeneration),
  673. static_cast<unsigned long long>(input.clientMessageSequence),
  674. object.registryKey,
  675. object.objectTag,
  676. object.objectRow,
  677. static_cast<unsigned>(object.slotType),
  678. static_cast<unsigned>(object.slotIndex),
  679. object.slotRow,
  680. object.senseSchema,
  681. object.schemaRow,
  682. object.generationPlusOne,
  683. object.hasGeneration ? 1U : 0U,
  684. object.deltaBits,
  685. object.valueCount);
  686. for (std::size_t valueIndex = 0; valueIndex < values.size(); ++valueIndex) {
  687. report_scene_value(instance, input, object, values[valueIndex], valueIndex);
  688. }
  689. }
  690. }
  691. /** Appends one observation and its complete owned value range. */
  692. [[nodiscard]] bool append_sense_observation(
  693. SenseObservationSnapshot& output,
  694. SenseObservation observation,
  695. std::span<const middleware::bap::activity_message::sense_update::DecodedValue>
  696. values) noexcept {
  697. if (output.observationCount == output.observations.size()
  698. || values.size() > output.values.size() - output.valueCount) {
  699. return false;
  700. }
  701. observation.firstValue = static_cast<std::uint32_t>(output.valueCount);
  702. observation.valueCount = static_cast<std::uint32_t>(values.size());
  703. output.observations[output.observationCount++] = observation;
  704. std::copy(values.begin(), values.end(), output.values.begin() + output.valueCount);
  705. output.valueCount += values.size();
  706. return true;
  707. }
  708. /** Replaces only keys present in one complete packet and keeps every omitted key. */
  709. [[nodiscard]] bool
  710. retain_sense_observations(Instance& instance, const SenseInput& input, std::uint64_t now) noexcept {
  711. namespace sense = middleware::bap::activity_message::sense_update;
  712. if (!complete_sense_observation_input(input)) {
  713. return false;
  714. }
  715. const sense::DecodedPacket& packet = input.decoded;
  716. SenseObservationSnapshot next{};
  717. next.revision = next_nonzero(instance.senseObservations.revision);
  718. next.sourceGeneration = input.sourceGeneration;
  719. for (std::size_t index = 0; index < packet.objectCount; ++index) {
  720. const sense::DecodedObject& object = packet.objects[index];
  721. const SenseObservationKey key{object.registryKey,
  722. object.objectTag,
  723. object.senseSchema,
  724. object.schemaRow,
  725. object.slotIndex,
  726. object.slotType};
  727. if (packet_has_sense_key(packet, key, index + 1)) {
  728. continue;
  729. }
  730. SenseObservation observation{};
  731. observation.binding = input.binding;
  732. observation.key = key;
  733. observation.sequence = next.revision;
  734. observation.tick = now;
  735. observation.sourceGeneration = input.sourceGeneration;
  736. observation.clientMessageSequence = input.clientMessageSequence;
  737. observation.generationPlusOne = object.generationPlusOne;
  738. observation.hasGeneration = object.hasGeneration;
  739. const std::span values(packet.values.data() + object.firstValue, object.valueCount);
  740. if (!append_sense_observation(next, observation, values)) {
  741. return false;
  742. }
  743. }
  744. if (instance.senseObservations.sourceGeneration == input.sourceGeneration) {
  745. const SenseObservationSnapshot& current = instance.senseObservations;
  746. for (std::size_t index = 0; index < current.observationCount; ++index) {
  747. const SenseObservation& observation = current.observations[index];
  748. if (packet_has_sense_key(packet, observation.key, 0)
  749. || observation.firstValue > current.valueCount
  750. || observation.valueCount > current.valueCount - observation.firstValue) {
  751. continue;
  752. }
  753. const std::span values(current.values.data() + observation.firstValue,
  754. observation.valueCount);
  755. static_cast<void>(append_sense_observation(next, observation, values));
  756. }
  757. }
  758. instance.senseObservations = next;
  759. instance.view.senseObservationCount = static_cast<std::uint32_t>(next.observationCount);
  760. instance.view.senseObservationValueCount = static_cast<std::uint32_t>(next.valueCount);
  761. instance.view.senseObservationRevision = next.revision;
  762. instance.view.senseObservationSourceGeneration = next.sourceGeneration;
  763. return true;
  764. }
  765. /** Merges complete squad deltas without discarding fields absent from a later report. */
  766. void retain_squad_sense(Instance& instance, const SenseInput& input) noexcept {
  767. namespace sense = middleware::bap::activity_message::sense_update;
  768. namespace squadSense = middleware::bap::activity_message::squad_sense;
  769. const sense::DecodedPacket& packet = input.decoded;
  770. if (input.sourceGeneration == 0 || input.sourceGeneration < instance.squadSenseSourceGeneration
  771. || packet.status == sense::DecodeStatus::malformed || packet.valuesTruncated
  772. || packet.objectsTruncated || packet.objectCount > packet.objects.size()
  773. || packet.valueCount > packet.values.size()) {
  774. return;
  775. }
  776. if (input.sourceGeneration != instance.squadSenseSourceGeneration) {
  777. instance.squadSense.clear();
  778. instance.squadSenseSourceGeneration = input.sourceGeneration;
  779. }
  780. for (const sense::DecodedObject& object : std::span(packet.objects).first(packet.objectCount)) {
  781. if (object.slotType != squadSense::kSlotType || object.senseSchema != squadSense::kSchema
  782. || object.status != sense::ObjectStatus::decoded || !object.hasGeneration) {
  783. continue;
  784. }
  785. auto found = std::ranges::find_if(instance.squadSense, [&](const SquadSenseRecord& record) {
  786. return same_sense_key(record.key, object);
  787. });
  788. squadSense::State merged =
  789. found == instance.squadSense.end() ? squadSense::State{} : found->state;
  790. // An uninitialized replica cannot replace the squad's recovery state.
  791. if (!squadSense::merge(merged, object, std::span(packet.values).first(packet.valueCount))
  792. || !merged.valid) {
  793. continue;
  794. }
  795. if (found != instance.squadSense.end()) {
  796. found->state = merged;
  797. } else if (instance.squadSense.size() < kScriptableGuardCapacity) {
  798. const SenseObservationKey key{object.registryKey,
  799. object.objectTag,
  800. object.senseSchema,
  801. object.schemaRow,
  802. object.slotIndex,
  803. object.slotType};
  804. try {
  805. instance.squadSense.push_back({key, merged});
  806. } catch (const std::bad_alloc&) {
  807. return;
  808. }
  809. }
  810. }
  811. }
  812. /** Applies one copied msg-6 decode summary. */
  813. void apply_sense(const SenseInput& input, std::uint64_t now) noexcept {
  814. Instance* const instance = find_instance(input.binding);
  815. if (instance == nullptr || !instance->view.active) {
  816. ++g_droppedIngress;
  817. return;
  818. }
  819. touch(*instance);
  820. ++instance->view.senseCount;
  821. trace_scene_sense(*instance, input);
  822. retain_squad_sense(*instance, input);
  823. static_cast<void>(retain_sense_observations(*instance, input, now));
  824. Event event{};
  825. event.binding = input.binding;
  826. event.tick = now;
  827. event.kind = EventKind::senseUpdate;
  828. event.epochFirst = input.epochFirst;
  829. event.epochSecond = input.epochSecond;
  830. event.payloadBytes = input.payloadBytes;
  831. event.peerHeardMask = input.peerHeardMask;
  832. event.tailBits = input.tailBits;
  833. event.consumedBits = input.consumedBits;
  834. event.firstGroupBits = input.firstGroupBits;
  835. event.firstRegistryKey = input.firstRegistryKey;
  836. event.groupsSeen = input.groupsSeen;
  837. event.groupsDecoded = input.groupsDecoded;
  838. event.groupsSkipped = input.groupsSkipped;
  839. event.objectsSeen = input.objectsSeen;
  840. event.objectsDecoded = input.objectsDecoded;
  841. event.firstSlotIndex = input.firstSlotIndex;
  842. event.firstSlotType = input.firstSlotType;
  843. // The event names one slot only: the first decoded ClientRef of the packet.
  844. if (input.decoded.objectCount != 0) {
  845. event.slotObjectTag = input.decoded.objects.front().objectTag;
  846. event.slotSenseSchema = input.decoded.objects.front().senseSchema;
  847. }
  848. event.senseDecodeStatus = input.decodeStatus;
  849. event.hasFirstObject = input.hasFirstObject;
  850. event.stateRevision = instance->view.stateRevision;
  851. event.sourceGeneration = input.sourceGeneration;
  852. event.clientMessageSequence = input.clientMessageSequence;
  853. event.lifetimeState = instance->view.lifetimeState;
  854. event.verdict = input.verdict;
  855. append_event(event);
  856. append_mission_input(event, complete_sense_observation_input(input) ? &input.decoded : nullptr);
  857. instance->view.lastEventSequence = g_sequence;
  858. }
  859. /** Retains one outer-valid client incident for inspection and parsed-field replay. */
  860. void apply_incident(const IncidentInput& input, std::uint64_t now) noexcept {
  861. Instance* const instance = find_instance(input.binding);
  862. if (instance == nullptr || !instance->view.active) {
  863. ++g_droppedIngress;
  864. ++g_droppedIncidents;
  865. return;
  866. }
  867. touch(*instance);
  868. ++instance->view.incidentsReceived;
  869. Event event{};
  870. event.binding = input.binding;
  871. event.tick = now;
  872. event.kind = EventKind::incidentReceived;
  873. event.sourceGeneration = input.sourceGeneration;
  874. event.clientMessageSequence = input.clientMessageSequence;
  875. event.payloadBytes = input.payloadBytes;
  876. fill_incident_event(event, input.incident, 0);
  877. event.hasPlayerTrigger = input.hasPlayerTrigger;
  878. if (input.hasPlayerTrigger) {
  879. event.playerTriggerRegistryKey = input.playerTrigger.registryKey;
  880. event.playerTriggerSlotType = input.playerTrigger.slotType;
  881. event.playerTriggerSlotIndex = input.playerTrigger.slotIndex;
  882. event.playerTriggerResolvedObjectId = input.playerTrigger.resolvedObjectId;
  883. }
  884. event.hasCinematic = input.hasCinematic;
  885. if (input.hasCinematic) {
  886. event.cinematicRegistryKey = input.cinematic.registryKey;
  887. event.cinematicSlotType = input.cinematic.slotType;
  888. event.cinematicSlotIndex = input.cinematic.slotIndex;
  889. event.cinematicRuntimeObjectId = input.cinematic.runtimeObjectId;
  890. event.cinematicEventValue = input.cinematic.eventValue;
  891. event.cinematicSignal = input.cinematicSignal;
  892. }
  893. append_event(event);
  894. append_mission_input(event, nullptr);
  895. instance->view.lastEventSequence = g_sequence;
  896. IncidentRecord* const record = reserve_incident_record();
  897. if (record == nullptr) {
  898. ++g_droppedIncidents;
  899. return;
  900. }
  901. *record = {};
  902. record->binding = input.binding;
  903. record->incident = input.incident;
  904. record->sequence = g_sequence;
  905. record->tick = now;
  906. record->lastSourceGeneration = input.sourceGeneration;
  907. record->clientMessageSequence = input.clientMessageSequence;
  908. record->payloadBytes = input.payloadBytes;
  909. record->status = IncidentStatus::received;
  910. }
  911. /** Retains one safe committed client State change in Host and ordered mission histories. */
  912. void apply_client_state_change(const ClientStateChangeInput& input, std::uint64_t now) noexcept {
  913. Instance* const instance = find_instance(input.binding);
  914. if (instance == nullptr || !instance->view.active) {
  915. ++g_droppedIngress;
  916. return;
  917. }
  918. touch(*instance);
  919. Event event{};
  920. event.binding = input.binding;
  921. event.tick = now;
  922. event.kind = EventKind::clientStateChanged;
  923. event.sourceGeneration = input.sourceGeneration;
  924. event.clientMessageSequence = input.clientMessageSequence;
  925. event.payloadBytes = input.payloadBytes;
  926. event.activityStateRevision = input.state.activityStateRevision;
  927. event.membershipRevision = input.state.membershipRevision;
  928. event.regionIndex = input.state.region.index;
  929. event.regionSliceSetHash = input.state.region.hash;
  930. event.currentRegionIndex = input.state.currentRegion.index;
  931. event.clientStateHasCurrentRegion = input.state.hasCurrentRegion;
  932. event.heldRegionIndex = input.state.heldRegion;
  933. event.teleportSliceSetIndex = input.state.teleportSliceSetIndex;
  934. event.teleportSliceSetHash = input.state.teleportSliceSetHash;
  935. event.spawnState = input.state.spawnState;
  936. event.teleportState = input.state.teleportState;
  937. event.clientStateHasRegion = input.state.hasRegion;
  938. event.clientStateHasSpawn = input.state.hasSpawn;
  939. event.clientStateHasTeleport = input.state.hasTeleport;
  940. append_event(event);
  941. append_mission_input(event, nullptr);
  942. instance->view.lastEventSequence = g_sequence;
  943. }
  944. /** Publishes one committed simulation-entity slot request without deriving readiness state. */
  945. void apply_entity_slots_requested(const EntitySlotsRequestedInput& input,
  946. std::uint64_t now) noexcept {
  947. Instance* const instance = find_instance(input.binding);
  948. if (instance == nullptr || !instance->view.active) {
  949. ++g_droppedIngress;
  950. return;
  951. }
  952. touch(*instance);
  953. Event event{};
  954. event.binding = input.binding;
  955. event.tick = now;
  956. event.kind = EventKind::entitySlotsRequested;
  957. event.stateRevision = instance->view.stateRevision;
  958. event.sourceGeneration = input.sourceGeneration;
  959. event.clientMessageSequence = input.clientMessageSequence;
  960. event.clientMessageType = middleware::bap::activity_message::entity_slot_request::kMessageType;
  961. event.requestedEntitySlots = input.requestedCount;
  962. append_event(event);
  963. append_mission_input(event, nullptr);
  964. instance->view.lastEventSequence = g_sequence;
  965. }
  966. /** Publishes one safe generic client envelope into the ordered mission-input feed. */
  967. void apply_client_message(const ClientMessageMissionInput& input, std::uint64_t now) noexcept {
  968. Instance* const instance = find_instance(input.binding);
  969. if (instance == nullptr || !instance->view.active) {
  970. ++g_droppedIngress;
  971. return;
  972. }
  973. touch(*instance);
  974. Event event{};
  975. event.binding = input.binding;
  976. event.tick = now;
  977. event.kind = EventKind::clientMessageReceived;
  978. event.stateRevision = instance->view.stateRevision;
  979. event.sourceGeneration = input.sourceGeneration;
  980. event.clientMessageSequence = input.clientMessageSequence;
  981. event.payloadBytes = input.payloadBytes;
  982. event.peerHeardMask = input.peerHeardMask;
  983. event.consumedBits = input.consumedBits;
  984. event.clientMessageType = input.messageType;
  985. event.clientMessageStatus = input.status;
  986. stamp_mission_sequence(event);
  987. ClientMessageSnapshot snapshot{};
  988. snapshot.messageType = input.messageType;
  989. snapshot.status = input.status;
  990. append_mission_input(event, nullptr, &snapshot);
  991. }
  992. /** Applies one operator transition only when its output slot is free. */
  993. void apply_control(const ControlRequest& request, std::uint64_t now) noexcept {
  994. Event event{};
  995. event.binding = request.binding;
  996. event.tick = now;
  997. event.lifetimeState = request.lifetimeState;
  998. Instance* const instance = find_instance(request.binding);
  999. if (instance == nullptr || !instance->view.active) {
  1000. event.kind = EventKind::operatorRefused;
  1001. ++g_refusedControls;
  1002. } else if (instance->view.outputPending) {
  1003. event.kind = EventKind::operatorRefused;
  1004. event.stateRevision = instance->view.stateRevision;
  1005. ++g_refusedControls;
  1006. } else if (instance->view.stateRevision == (std::numeric_limits<std::uint64_t>::max)()) {
  1007. event.kind = EventKind::operatorRefused;
  1008. event.stateRevision = instance->view.stateRevision;
  1009. ++g_refusedControls;
  1010. } else {
  1011. touch(*instance);
  1012. ++instance->view.stateRevision;
  1013. instance->view.lifetimeState = request.lifetimeState;
  1014. instance->view.lastOutputAttemptTick = 0;
  1015. instance->view.lastOutputSourceGeneration = 0;
  1016. instance->view.outputAttempts = 0;
  1017. instance->view.outputStatus = OutputStatus::pending;
  1018. instance->view.outputKind = OutputKind::authState;
  1019. instance->view.outputPending = true;
  1020. event.kind = EventKind::authStateCommitted;
  1021. event.stateRevision = instance->view.stateRevision;
  1022. }
  1023. append_event(event);
  1024. if (instance != nullptr) {
  1025. instance->view.lastEventSequence = g_sequence;
  1026. }
  1027. }
  1028. /** Commits one operator incident into the per-generation ordered output history. */
  1029. void apply_incident_control(const IncidentRequest& request, std::uint64_t now) noexcept {
  1030. Event event{};
  1031. event.binding = request.binding;
  1032. event.tick = now;
  1033. event.kind = EventKind::incidentRefused;
  1034. fill_incident_event(event, request.incident, 0);
  1035. Instance* const instance = find_instance(request.binding);
  1036. IncidentRecord* record = nullptr;
  1037. if (instance == nullptr || !instance->view.active || instance->view.outputPending
  1038. || instance->view.incidentRevision == (std::numeric_limits<std::uint64_t>::max)()) {
  1039. ++g_refusedControls;
  1040. ++g_refusedIncidents;
  1041. } else if ((record = reserve_incident_record()) == nullptr) {
  1042. ++g_refusedControls;
  1043. ++g_refusedIncidents;
  1044. } else {
  1045. touch(*instance);
  1046. ++instance->view.incidentRevision;
  1047. ++instance->view.incidentsQueued;
  1048. ++instance->view.incidentsPending;
  1049. instance->view.outputPending = true;
  1050. instance->view.outputKind = OutputKind::incident;
  1051. instance->view.outputStatus = OutputStatus::pending;
  1052. instance->view.lastOutputAttemptTick = 0;
  1053. instance->view.lastOutputSourceGeneration = 0;
  1054. instance->view.outputAttempts = 0;
  1055. *record = {};
  1056. record->binding = request.binding;
  1057. record->incident = request.incident;
  1058. record->revision = instance->view.incidentRevision;
  1059. record->tick = now;
  1060. record->status = IncidentStatus::queued;
  1061. record->outbound = true;
  1062. event.kind = EventKind::incidentQueued;
  1063. fill_incident_event(event, request.incident, record->revision);
  1064. }
  1065. append_event(event);
  1066. if (record != nullptr && event.kind == EventKind::incidentQueued) {
  1067. record->sequence = g_sequence;
  1068. }
  1069. if (instance != nullptr) {
  1070. instance->view.lastEventSequence = g_sequence;
  1071. }
  1072. }
  1073. } // namespace detail
  1074. using namespace detail;
  1075. /** Queues one owned msg-6 prefix for the Activity Host service. */
  1076. bool submit_sense(const SenseInput& input) noexcept {
  1077. if (!state::activity::binding_matches(input.binding) || input.sourceGeneration == 0) {
  1078. return false;
  1079. }
  1080. AcquireSRWLockExclusive(&g_lock);
  1081. PendingInput pending{};
  1082. pending.kind = PendingKind::sense;
  1083. pending.sense = input;
  1084. if (!append_pending(std::move(pending))) {
  1085. ++g_droppedIngress;
  1086. ReleaseSRWLockExclusive(&g_lock);
  1087. return false;
  1088. }
  1089. ++g_queuedIngress;
  1090. ReleaseSRWLockExclusive(&g_lock);
  1091. return true;
  1092. }
  1093. /** Queues one owned, outer-valid client msg 19 for the Activity Host service. */
  1094. bool submit_incident(const IncidentInput& input) noexcept {
  1095. if (!state::activity::binding_matches(input.binding) || !incident_allowed(input.incident)) {
  1096. return false;
  1097. }
  1098. AcquireSRWLockExclusive(&g_lock);
  1099. PendingInput pending{};
  1100. pending.kind = PendingKind::incident;
  1101. pending.incident = input;
  1102. if (!append_pending(std::move(pending))) {
  1103. ++g_droppedIngress;
  1104. ++g_droppedIncidents;
  1105. ReleaseSRWLockExclusive(&g_lock);
  1106. return false;
  1107. }
  1108. ++g_queuedIngress;
  1109. ReleaseSRWLockExclusive(&g_lock);
  1110. return true;
  1111. }
  1112. /** Queues one post-commit client State after-image for the ordered Host service slice. */
  1113. bool submit_client_state_change(const ClientStateChangeInput& input) noexcept {
  1114. // A report that moved no region, spawn or teleport field is the client's settle, sent once
  1115. // spawn-in completes. The mission surface needs it to time the opening line, so only a
  1116. // malformed report is refused, never a material-less one.
  1117. if (!state::activity::binding_matches(input.binding) || input.sourceGeneration == 0
  1118. || input.clientMessageSequence == 0 || !input.state.committed
  1119. || (input.state.hasRegion && input.state.region.index < 0)
  1120. || input.state.activityStateRevision == state::activity::kInvalidRevision) {
  1121. return false;
  1122. }
  1123. AcquireSRWLockExclusive(&g_lock);
  1124. PendingInput pending{};
  1125. pending.kind = PendingKind::clientStateChange;
  1126. pending.clientStateChange = input;
  1127. if (!append_pending(std::move(pending))) {
  1128. ++g_droppedIngress;
  1129. ReleaseSRWLockExclusive(&g_lock);
  1130. return false;
  1131. }
  1132. ++g_queuedIngress;
  1133. ReleaseSRWLockExclusive(&g_lock);
  1134. return true;
  1135. }
  1136. /** Queues one committed msg-20 simulation-entity slot request. */
  1137. bool submit_entity_slots_requested(const EntitySlotsRequestedInput& input) noexcept {
  1138. if (!state::activity::binding_matches(input.binding) || input.sourceGeneration == 0
  1139. || input.clientMessageSequence == 0 || input.requestedCount <= 0) {
  1140. return false;
  1141. }
  1142. AcquireSRWLockExclusive(&g_lock);
  1143. PendingInput pending{};
  1144. pending.kind = PendingKind::entitySlotsRequested;
  1145. pending.entitySlotsRequested = input;
  1146. if (!append_pending(std::move(pending))) {
  1147. ++g_droppedIngress;
  1148. ReleaseSRWLockExclusive(&g_lock);
  1149. return false;
  1150. }
  1151. ++g_queuedIngress;
  1152. ReleaseSRWLockExclusive(&g_lock);
  1153. return true;
  1154. }
  1155. /** Retains one owned client message without entering the reducer queue. */
  1156. std::uint64_t record_client_message(
  1157. const ClientMessageInput& input,
  1158. const middleware::bap::activity_message::sense_update::DecodedPacket* sense) noexcept {
  1159. if (!state::activity::binding_matches(input.binding)) {
  1160. return 0;
  1161. }
  1162. ClientMessageRecord record{};
  1163. record.binding.sessionId = input.binding.sessionId;
  1164. record.binding.createdRevision = input.binding.createdRevision;
  1165. record.authoritative = input.authoritative;
  1166. record.tick = GetTickCount64();
  1167. record.sourceGeneration = input.sourceGeneration;
  1168. record.payloadFingerprint = input.payloadFingerprint;
  1169. record.messageType = input.messageType;
  1170. record.payloadBytes = input.payloadBytes;
  1171. record.peerHeardMask = input.peerHeardMask;
  1172. record.consumedBits = input.consumedBits;
  1173. record.status = input.status;
  1174. record.hasPayloadFingerprint = input.hasPayloadFingerprint;
  1175. record.hasAuthoritative = input.hasAuthoritative;
  1176. AcquireSRWLockExclusive(&g_lock);
  1177. const std::uint64_t sequence = append_client_message(record);
  1178. if (sense != nullptr) {
  1179. append_client_message_detail(sequence, input.messageType, sense);
  1180. }
  1181. ReleaseSRWLockExclusive(&g_lock);
  1182. return sequence;
  1183. }
  1184. /** Queues one owned client envelope without a richer typed mission reducer. */
  1185. bool submit_client_message(const ClientMessageInput& input,
  1186. std::uint64_t clientMessageSequence) noexcept {
  1187. namespace activity_message = middleware::bap::activity_message;
  1188. namespace communication = activity_message::wire_schema::communication;
  1189. communication::ActivityCommunicationRoute route{};
  1190. const bool executableRoute =
  1191. state::activity_sdk::executable_communication_route(input.messageType, route);
  1192. if (!state::activity::binding_matches(input.binding) || input.sourceGeneration == 0
  1193. || clientMessageSequence == 0 || !executableRoute
  1194. || route.ingressDelivery != communication::IngressDeliveryPolicy::protocolHostInput
  1195. || (route.ingressClass != communication::IngressClass::nativeMetadataOnly
  1196. && route.ingressClass != communication::IngressClass::nativeParsed)
  1197. || input.messageType == activity_message::sense_update::kMessageType
  1198. || input.messageType == activity_message::incident::kMessageType
  1199. || input.messageType == activity_message::client_authoritative_data::kMessageType) {
  1200. return false;
  1201. }
  1202. AcquireSRWLockExclusive(&g_lock);
  1203. PendingInput pending{};
  1204. pending.kind = PendingKind::clientMessage;
  1205. pending.clientMessage.binding = input.binding;
  1206. pending.clientMessage.sourceGeneration = input.sourceGeneration;
  1207. pending.clientMessage.clientMessageSequence = clientMessageSequence;
  1208. pending.clientMessage.messageType = input.messageType;
  1209. pending.clientMessage.payloadBytes = input.payloadBytes;
  1210. pending.clientMessage.peerHeardMask = input.peerHeardMask;
  1211. pending.clientMessage.consumedBits = input.consumedBits;
  1212. pending.clientMessage.status = input.status;
  1213. if (!append_pending(std::move(pending))) {
  1214. ++g_droppedIngress;
  1215. ReleaseSRWLockExclusive(&g_lock);
  1216. return false;
  1217. }
  1218. ++g_queuedIngress;
  1219. ReleaseSRWLockExclusive(&g_lock);
  1220. return true;
  1221. }
  1222. /** Copies one retained scalar decode by its ingress sequence. */
  1223. bool client_message_detail(std::uint64_t sequence, ClientMessageDetail& output) noexcept {
  1224. output = {};
  1225. if (sequence == 0) {
  1226. return false;
  1227. }
  1228. AcquireSRWLockShared(&g_lock);
  1229. bool found = false;
  1230. for (std::size_t offset = g_clientMessageDetailCount; offset != 0; --offset) {
  1231. const std::size_t index =
  1232. (g_clientMessageDetailStart + offset - 1) % g_clientMessageDetails.size();
  1233. if (g_clientMessageDetails[index].sequence == sequence) {
  1234. output = g_clientMessageDetails[index];
  1235. found = true;
  1236. break;
  1237. }
  1238. }
  1239. ReleaseSRWLockShared(&g_lock);
  1240. return found;
  1241. }
  1242. /** Queues one operator Auth-state transition for the exact activity generation. */
  1243. bool request_auth_state(const state::activity::SessionBinding& binding,
  1244. std::uint8_t lifetimeState) noexcept {
  1245. if (!lifetime_allowed(lifetimeState) || !state::activity::binding_matches(binding)) {
  1246. return false;
  1247. }
  1248. AcquireSRWLockExclusive(&g_lock);
  1249. const Instance* const instance = find_instance(binding);
  1250. if (has_queued_control(binding)
  1251. || (instance != nullptr
  1252. && (instance->view.outputPending || instance->view.scriptableReservationPending))) {
  1253. ++g_refusedControls;
  1254. ReleaseSRWLockExclusive(&g_lock);
  1255. return false;
  1256. }
  1257. PendingInput pending{};
  1258. pending.kind = PendingKind::authControl;
  1259. pending.control = {binding, lifetimeState};
  1260. if (!append_pending(std::move(pending))) {
  1261. ++g_refusedControls;
  1262. ReleaseSRWLockExclusive(&g_lock);
  1263. return false;
  1264. }
  1265. ++g_queuedControls;
  1266. ReleaseSRWLockExclusive(&g_lock);
  1267. return true;
  1268. }
  1269. /** Queues one outer-valid operator msg 19 for the exact activity generation. */
  1270. bool request_incident(
  1271. const state::activity::SessionBinding& binding,
  1272. const middleware::bap::activity_message::incident::Incident& incident) noexcept {
  1273. if (!incident_allowed(incident) || !state::activity::binding_matches(binding)) {
  1274. return false;
  1275. }
  1276. AcquireSRWLockExclusive(&g_lock);
  1277. const Instance* const instance = find_instance(binding);
  1278. if (has_queued_control(binding)
  1279. || (instance != nullptr
  1280. && (instance->view.outputPending || instance->view.scriptableReservationPending))) {
  1281. ++g_refusedControls;
  1282. ++g_refusedIncidents;
  1283. ReleaseSRWLockExclusive(&g_lock);
  1284. return false;
  1285. }
  1286. PendingInput pending{};
  1287. pending.kind = PendingKind::incidentControl;
  1288. pending.incidentControl = {binding, incident};
  1289. if (!append_pending(std::move(pending))) {
  1290. ++g_refusedControls;
  1291. ++g_refusedIncidents;
  1292. ReleaseSRWLockExclusive(&g_lock);
  1293. return false;
  1294. }
  1295. ++g_queuedControls;
  1296. ReleaseSRWLockExclusive(&g_lock);
  1297. return true;
  1298. }
  1299. /** Applies queued client and operator events on the server service slice. */
  1300. void service(std::uint64_t now) noexcept {
  1301. std::array<state::activity::SessionBinding, state::activity::kSessionCapacity> bindings{};
  1302. std::size_t bindingCount = 0;
  1303. static_cast<void>(state::activity::snapshot_retained_bindings(bindings, bindingCount));
  1304. AcquireSRWLockExclusive(&g_lock);
  1305. for (Instance& instance : g_instances) {
  1306. bool active = false;
  1307. for (std::size_t index = 0; index < bindingCount; ++index) {
  1308. if (same_binding(instance.view.binding, bindings[index])) {
  1309. active = true;
  1310. break;
  1311. }
  1312. }
  1313. instance.view.active = instance.occupied && active;
  1314. if (instance.occupied && !instance.view.active && instance.view.outputPending) {
  1315. cancel_output(instance, now);
  1316. }
  1317. if (instance.occupied && !instance.view.active
  1318. && instance.view.scriptableReservationPending) {
  1319. instance.scriptableReservation = {};
  1320. instance.view.scriptableReservedRevision = 0;
  1321. instance.view.scriptableReservationPending = false;
  1322. }
  1323. }
  1324. for (std::size_t index = 0; index < bindingCount; ++index) {
  1325. Instance* const instance = ensure_instance(bindings[index]);
  1326. if (instance != nullptr) {
  1327. instance->view.active = true;
  1328. touch(*instance);
  1329. }
  1330. }
  1331. while (g_pendingRead < g_pending.size()) {
  1332. PendingInput pending = std::move(g_pending[g_pendingRead]);
  1333. ++g_pendingRead;
  1334. if (pending.kind == PendingKind::discardedControl) {
  1335. continue;
  1336. }
  1337. if (pending.kind == PendingKind::authControl) {
  1338. --g_queuedControls;
  1339. apply_control(pending.control, now);
  1340. } else if (pending.kind == PendingKind::incidentControl) {
  1341. --g_queuedControls;
  1342. apply_incident_control(pending.incidentControl, now);
  1343. } else if (pending.kind == PendingKind::scriptableControl) {
  1344. --g_queuedControls;
  1345. apply_scriptable_control(pending.scriptableControl, now);
  1346. } else if (pending.kind == PendingKind::incident) {
  1347. --g_queuedIngress;
  1348. apply_incident(pending.incident, now);
  1349. } else if (pending.kind == PendingKind::clientStateChange) {
  1350. --g_queuedIngress;
  1351. apply_client_state_change(pending.clientStateChange, now);
  1352. } else if (pending.kind == PendingKind::entitySlotsRequested) {
  1353. --g_queuedIngress;
  1354. apply_entity_slots_requested(pending.entitySlotsRequested, now);
  1355. } else if (pending.kind == PendingKind::clientMessage) {
  1356. --g_queuedIngress;
  1357. apply_client_message(pending.clientMessage, now);
  1358. } else {
  1359. --g_queuedIngress;
  1360. apply_sense(pending.sense, now);
  1361. }
  1362. }
  1363. g_pending.clear();
  1364. g_pendingRead = 0;
  1365. ReleaseSRWLockExclusive(&g_lock);
  1366. }
  1367. /** Copies the latest complete diagnostic view. */
  1368. void snapshot(DiagnosticsSnapshot& output) noexcept {
  1369. output = {};
  1370. AcquireSRWLockShared(&g_lock);
  1371. for (const Instance& instance : g_instances) {
  1372. if (instance.occupied && output.instanceCount < output.instances.size()) {
  1373. output.instances[output.instanceCount] = instance.view;
  1374. ++output.instanceCount;
  1375. }
  1376. }
  1377. for (std::size_t index = 0; index < g_eventCount; ++index) {
  1378. output.events[index] = g_events[(g_eventStart + index) % g_events.size()];
  1379. }
  1380. output.eventCount = g_eventCount;
  1381. for (const IncidentRecord& record : g_incidents) {
  1382. if (record.sequence != 0 && output.incidentCount < output.incidents.size()) {
  1383. output.incidents[output.incidentCount++] = record;
  1384. }
  1385. }
  1386. std::sort(output.incidents.begin(),
  1387. output.incidents.begin() + static_cast<std::ptrdiff_t>(output.incidentCount),
  1388. [](const IncidentRecord& left, const IncidentRecord& right) noexcept {
  1389. return left.sequence < right.sequence;
  1390. });
  1391. for (std::size_t index = 0; index < g_clientMessageCount; ++index) {
  1392. output.clientMessages[index] =
  1393. g_clientMessages[(g_clientMessageStart + index) % g_clientMessages.size()];
  1394. }
  1395. output.clientMessageCount = g_clientMessageCount;
  1396. output.droppedIngress = g_droppedIngress;
  1397. output.droppedIncidents = g_droppedIncidents;
  1398. output.queuedControls = g_queuedControls;
  1399. output.refusedControls = g_refusedControls;
  1400. output.refusedIncidents = g_refusedIncidents;
  1401. output.overwrittenEvents = g_overwrittenEvents;
  1402. output.overwrittenIncidents = g_overwrittenIncidents;
  1403. output.overwrittenClientMessages = g_overwrittenClientMessages;
  1404. ReleaseSRWLockShared(&g_lock);
  1405. }
  1406. /** Copies one exact activity generation without exposing the Host lock. */
  1407. bool instance_snapshot(const state::activity::SessionBinding& binding,
  1408. InstanceSnapshot& output) noexcept {
  1409. output = {};
  1410. AcquireSRWLockShared(&g_lock);
  1411. const Instance* const instance = find_instance(binding);
  1412. const bool found = instance != nullptr;
  1413. if (found) {
  1414. output = instance->view;
  1415. }
  1416. ReleaseSRWLockShared(&g_lock);
  1417. return found;
  1418. }
  1419. /** Reads the current feed position without replaying retained history. */
  1420. EventCursor current_event_cursor() noexcept {
  1421. AcquireSRWLockShared(&g_lock);
  1422. const EventCursor cursor{g_eventGeneration, g_sequence};
  1423. ReleaseSRWLockShared(&g_lock);
  1424. return cursor;
  1425. }
  1426. /** Copies retained events after one cursor and reports reset or overwrite gaps. */
  1427. void read_events_after(EventCursor after, EventRead& output) noexcept {
  1428. output = {};
  1429. AcquireSRWLockShared(&g_lock);
  1430. output.cursor = {g_eventGeneration, g_sequence};
  1431. output.reset = after.generation != g_eventGeneration;
  1432. if (g_eventCount != 0) {
  1433. const std::uint64_t afterSequence = output.reset ? 0 : after.sequence;
  1434. const std::uint64_t oldestSequence = g_events[g_eventStart].sequence;
  1435. std::size_t first = 0;
  1436. if (afterSequence < oldestSequence) {
  1437. const std::uint64_t retainedPredecessor = oldestSequence - 1;
  1438. if (afterSequence < retainedPredecessor) {
  1439. output.gap = true;
  1440. output.missed = retainedPredecessor - afterSequence;
  1441. }
  1442. } else {
  1443. while (first < g_eventCount
  1444. && g_events[(g_eventStart + first) % g_events.size()].sequence
  1445. <= afterSequence) {
  1446. ++first;
  1447. }
  1448. }
  1449. for (; first < g_eventCount; ++first) {
  1450. output.events[output.count++] = g_events[(g_eventStart + first) % g_events.size()];
  1451. }
  1452. }
  1453. ReleaseSRWLockShared(&g_lock);
  1454. }
  1455. /** Reads the current accepted-mission-input position without replaying retained history. */
  1456. MissionInputCursor current_mission_input_cursor() noexcept {
  1457. AcquireSRWLockShared(&g_lock);
  1458. const MissionInputCursor cursor{g_eventGeneration, g_missionInputSequence};
  1459. ReleaseSRWLockShared(&g_lock);
  1460. return cursor;
  1461. }
  1462. /** @return True when durable mission State still owes this retained accepted row. */
  1463. [[nodiscard]] bool mission_input_owed(const MissionInputRecord& record) noexcept {
  1464. state::activity::mission::InputSequenceSnapshot cursors{};
  1465. // A faulted binding owes nothing. Its program is skipped, so it never commits, and its rows
  1466. // would otherwise be retained for the life of the process.
  1467. if (!state::activity::mission::input_sequence_snapshot(record.view.event.binding, cursors)
  1468. || cursors.faulted) {
  1469. return false;
  1470. }
  1471. const std::uint64_t sequence = record.view.event.missionSequence;
  1472. return sequence > cursors.committed && sequence <= cursors.issued;
  1473. }
  1474. /** Drops every retained row that durable mission State no longer owes. */
  1475. void retire_settled_mission_inputs() noexcept {
  1476. static_cast<void>(std::erase_if(g_missionInputs, [](const MissionInputRecord& record) noexcept {
  1477. return !mission_input_owed(record);
  1478. }));
  1479. }
  1480. /** Copies accepted client mission inputs after one cursor. */
  1481. void read_mission_inputs_after(MissionInputCursor after, MissionInputRead& output) noexcept {
  1482. output = {};
  1483. AcquireSRWLockExclusive(&g_lock);
  1484. output.reset = after.generation != g_eventGeneration;
  1485. const std::uint64_t afterSequence = output.reset ? 0 : after.sequence;
  1486. retire_settled_mission_inputs();
  1487. output.cursor = {g_eventGeneration, afterSequence};
  1488. if (!g_missionInputs.empty()) {
  1489. const std::uint64_t retainedPredecessor = g_missionInputs.front().view.sequence - 1;
  1490. if (afterSequence < retainedPredecessor) {
  1491. output.gap = true;
  1492. output.missed = retainedPredecessor - afterSequence;
  1493. }
  1494. for (const MissionInputRecord& record : g_missionInputs) {
  1495. if (record.view.sequence <= afterSequence) {
  1496. continue;
  1497. }
  1498. if (output.count == output.events.size()) {
  1499. break;
  1500. }
  1501. output.events[output.count++] = record.view;
  1502. }
  1503. if (output.count != 0) {
  1504. output.cursor.sequence = output.events[output.count - 1].sequence;
  1505. }
  1506. } else if (afterSequence < g_missionInputSequence) {
  1507. output.gap = true;
  1508. output.missed = g_missionInputSequence - afterSequence;
  1509. output.cursor.sequence = g_missionInputSequence;
  1510. }
  1511. ReleaseSRWLockExclusive(&g_lock);
  1512. }
  1513. /** Copies the exact Sense values owned by one retained accepted-input row. */
  1514. bool mission_input_sense_snapshot(std::uint64_t sequence,
  1515. SenseObservationSnapshot& output) noexcept {
  1516. namespace sense = middleware::bap::activity_message::sense_update;
  1517. output = {};
  1518. if (sequence == 0) {
  1519. return false;
  1520. }
  1521. AcquireSRWLockShared(&g_lock);
  1522. const MissionInputRecord* selected = nullptr;
  1523. for (std::size_t offset = g_missionInputs.size(); offset != 0; --offset) {
  1524. const MissionInputRecord& candidate = g_missionInputs[offset - 1];
  1525. if (candidate.view.sequence == sequence) {
  1526. selected = &candidate;
  1527. break;
  1528. }
  1529. }
  1530. bool copied = selected != nullptr && selected->hasSense
  1531. && selected->view.event.kind == EventKind::senseUpdate;
  1532. if (copied) {
  1533. const Event& event = selected->view.event;
  1534. const sense::DecodedPacket& packet = selected->sense;
  1535. copied = packet.status == sense::DecodeStatus::complete && !packet.objectsTruncated
  1536. && !packet.valuesTruncated && packet.objectCount <= packet.objects.size()
  1537. && packet.valueCount <= packet.values.size();
  1538. if (copied) {
  1539. output.revision = event.sequence;
  1540. output.sourceGeneration = event.sourceGeneration;
  1541. }
  1542. for (std::size_t index = 0; copied && index < packet.objectCount; ++index) {
  1543. const sense::DecodedObject& object = packet.objects[index];
  1544. if (object.firstValue > packet.valueCount
  1545. || object.valueCount > packet.valueCount - object.firstValue) {
  1546. copied = false;
  1547. break;
  1548. }
  1549. SenseObservation observation{};
  1550. observation.binding = event.binding;
  1551. observation.key = {object.registryKey,
  1552. object.objectTag,
  1553. object.senseSchema,
  1554. object.schemaRow,
  1555. object.slotIndex,
  1556. object.slotType};
  1557. observation.sequence = event.sequence;
  1558. observation.tick = event.tick;
  1559. observation.sourceGeneration = event.sourceGeneration;
  1560. observation.clientMessageSequence = event.clientMessageSequence;
  1561. observation.generationPlusOne = object.generationPlusOne;
  1562. observation.hasGeneration = object.hasGeneration;
  1563. copied = append_sense_observation(
  1564. output, observation, {packet.values.data() + object.firstValue, object.valueCount});
  1565. }
  1566. }
  1567. if (!copied) {
  1568. output = {};
  1569. }
  1570. ReleaseSRWLockShared(&g_lock);
  1571. return copied;
  1572. }
  1573. /** Copies one exact generic client-message snapshot owned by the mission-input feed. */
  1574. bool mission_input_client_message_snapshot(std::uint64_t sequence,
  1575. ClientMessageSnapshot& output) noexcept {
  1576. output = {};
  1577. if (sequence == 0) {
  1578. return false;
  1579. }
  1580. AcquireSRWLockShared(&g_lock);
  1581. const MissionInputRecord* selected = nullptr;
  1582. for (std::size_t offset = g_missionInputs.size(); offset != 0; --offset) {
  1583. const MissionInputRecord& candidate = g_missionInputs[offset - 1];
  1584. if (candidate.view.sequence == sequence) {
  1585. selected = &candidate;
  1586. break;
  1587. }
  1588. }
  1589. const bool copied = selected != nullptr && selected->hasClientMessage
  1590. && selected->view.event.kind == EventKind::clientMessageReceived;
  1591. if (copied) {
  1592. output = selected->clientMessage;
  1593. }
  1594. ReleaseSRWLockShared(&g_lock);
  1595. return copied;
  1596. }
  1597. /** Copies initialized recovery state for one exact ActivityClient generation. */
  1598. bool snapshot_squad_sense(const state::activity::SessionBinding& binding,
  1599. std::uint64_t sourceGeneration,
  1600. const SenseObservationKey& key,
  1601. middleware::bap::activity_message::squad_sense::State& output) noexcept {
  1602. output = {};
  1603. AcquireSRWLockShared(&g_lock);
  1604. const Instance* const instance = find_instance(binding);
  1605. bool found = false;
  1606. if (instance != nullptr && instance->squadSenseSourceGeneration == sourceGeneration) {
  1607. for (const SquadSenseRecord& record : instance->squadSense) {
  1608. if (record.key.registryKey == key.registryKey && record.key.objectTag == key.objectTag
  1609. && record.key.senseSchema == key.senseSchema && record.key.slotType == key.slotType
  1610. && record.key.slotIndex == key.slotIndex && record.state.valid) {
  1611. output = record.state;
  1612. found = true;
  1613. break;
  1614. }
  1615. }
  1616. }
  1617. ReleaseSRWLockShared(&g_lock);
  1618. return found;
  1619. }
  1620. /** Copies the latest complete Sense observations for one exact activity generation. */
  1621. bool snapshot_sense_observations(const state::activity::SessionBinding& binding,
  1622. SenseObservationSnapshot& output) noexcept {
  1623. output = {};
  1624. AcquireSRWLockShared(&g_lock);
  1625. const Instance* const instance = find_instance(binding);
  1626. const bool found = instance != nullptr;
  1627. if (found) {
  1628. output = instance->senseObservations;
  1629. }
  1630. ReleaseSRWLockShared(&g_lock);
  1631. return found;
  1632. }
  1633. /** Selects one exact reflected scalar without changing the retained snapshot. */
  1634. SenseScalarStatus select_sense_scalar(const SenseObservationSnapshot& snapshot,
  1635. const state::activity::SessionBinding& binding,
  1636. const SenseScalarIdentity& identity,
  1637. SenseScalarSample& output) noexcept {
  1638. namespace sense = middleware::bap::activity_message::sense_update;
  1639. output = {};
  1640. if (binding.sessionId == state::activity::kAbsentSessionId
  1641. || binding.createdRevision == state::activity::kInvalidRevision
  1642. || identity.object.schemaRow == sense::kAbsentRuntimeRow
  1643. || identity.fieldSchemaRow == sense::kAbsentRuntimeRow
  1644. || identity.fieldRow == sense::kAbsentRuntimeRow
  1645. || identity.fieldSchemaRow != identity.object.schemaRow) {
  1646. return SenseScalarStatus::invalidIdentity;
  1647. }
  1648. if (snapshot.sourceGeneration == 0 || snapshot.observationCount > snapshot.observations.size()
  1649. || snapshot.valueCount > snapshot.values.size()) {
  1650. return SenseScalarStatus::invalidSnapshot;
  1651. }
  1652. const auto sameKey = [&identity](const SenseObservationKey& candidate) noexcept {
  1653. const SenseObservationKey& expected = identity.object;
  1654. return candidate.registryKey == expected.registryKey
  1655. && candidate.objectTag == expected.objectTag
  1656. && candidate.senseSchema == expected.senseSchema
  1657. && candidate.schemaRow == expected.schemaRow
  1658. && candidate.slotIndex == expected.slotIndex
  1659. && candidate.slotType == expected.slotType;
  1660. };
  1661. bool found = false;
  1662. for (std::size_t observationIndex = 0; observationIndex < snapshot.observationCount;
  1663. ++observationIndex) {
  1664. const SenseObservation& observation = snapshot.observations[observationIndex];
  1665. if (!same_binding(observation.binding, binding) || !sameKey(observation.key)) {
  1666. continue;
  1667. }
  1668. if (observation.sourceGeneration != snapshot.sourceGeneration
  1669. || observation.firstValue > snapshot.valueCount
  1670. || observation.valueCount > snapshot.valueCount - observation.firstValue) {
  1671. output = {};
  1672. return SenseScalarStatus::invalidSnapshot;
  1673. }
  1674. for (std::size_t valueIndex = 0; valueIndex < observation.valueCount; ++valueIndex) {
  1675. const sense::DecodedValue& value = snapshot.values[observation.firstValue + valueIndex];
  1676. if (value.schemaRow != identity.fieldSchemaRow || value.fieldRow != identity.fieldRow
  1677. || value.fieldOrdinal != identity.fieldOrdinal
  1678. || value.occurrence != identity.occurrence || value.kind != identity.kind) {
  1679. continue;
  1680. }
  1681. if (found) {
  1682. output = {};
  1683. return SenseScalarStatus::ambiguous;
  1684. }
  1685. output.binding = observation.binding;
  1686. output.identity = identity;
  1687. output.value = value;
  1688. output.observationRevision = observation.sequence;
  1689. output.tick = observation.tick;
  1690. output.sourceGeneration = observation.sourceGeneration;
  1691. output.clientMessageSequence = observation.clientMessageSequence;
  1692. output.generationPlusOne = observation.generationPlusOne;
  1693. output.hasGeneration = observation.hasGeneration;
  1694. found = true;
  1695. }
  1696. }
  1697. return found ? SenseScalarStatus::ready : SenseScalarStatus::notFound;
  1698. }
  1699. /** @return True when two selected scalars have the same presence and raw typed value. */
  1700. bool same_sense_scalar_value(const SenseScalarSample& left,
  1701. const SenseScalarSample& right) noexcept {
  1702. if (left.value.kind != right.value.kind || left.value.present != right.value.present) {
  1703. return false;
  1704. }
  1705. if (!left.value.present) {
  1706. return true;
  1707. }
  1708. return left.value.unsignedValue == right.value.unsignedValue
  1709. && left.value.signedValue == right.value.signedValue
  1710. && std::bit_cast<std::uint32_t>(left.value.realValue)
  1711. == std::bit_cast<std::uint32_t>(right.value.realValue);
  1712. }
  1713. /** @return True when both samples name the same ActivityClient and reported object generations. */
  1714. bool same_sense_scalar_generations(const SenseScalarSample& left,
  1715. const SenseScalarSample& right) noexcept {
  1716. return left.sourceGeneration != 0 && left.sourceGeneration == right.sourceGeneration
  1717. && left.hasGeneration && right.hasGeneration
  1718. && left.generationPlusOne == right.generationPlusOne;
  1719. }
  1720. /** @return Stable text for one exact scalar-selection result. */
  1721. const char* sense_scalar_status_name(SenseScalarStatus status) noexcept {
  1722. switch (status) {
  1723. case SenseScalarStatus::ready:
  1724. return "ready";
  1725. case SenseScalarStatus::invalidIdentity:
  1726. return "invalid_identity";
  1727. case SenseScalarStatus::invalidSnapshot:
  1728. return "invalid_snapshot";
  1729. case SenseScalarStatus::notFound:
  1730. return "not_found";
  1731. case SenseScalarStatus::ambiguous:
  1732. return "ambiguous";
  1733. }
  1734. return "invalid_status";
  1735. }
  1736. /** Reads the committed Auth state for one exact activity generation. */
  1737. bool auth_state(const state::activity::SessionBinding& binding, AuthState& output) noexcept {
  1738. output = {};
  1739. AcquireSRWLockShared(&g_lock);
  1740. const Instance* const instance = find_instance(binding);
  1741. const bool found = instance != nullptr;
  1742. if (found) {
  1743. output.revision = instance->view.stateRevision;
  1744. output.lifetimeState = instance->view.lifetimeState;
  1745. }
  1746. ReleaseSRWLockShared(&g_lock);
  1747. return found;
  1748. }
  1749. /** Reads the one pending incident occupying the exact instance's serialized output slot. */
  1750. bool pending_incident(const state::activity::SessionBinding& binding,
  1751. PendingIncident& output) noexcept {
  1752. output = {};
  1753. AcquireSRWLockShared(&g_lock);
  1754. const Instance* const instance = find_instance(binding);
  1755. const bool pending = instance != nullptr && instance->view.active
  1756. && instance->view.outputPending
  1757. && instance->view.outputKind == OutputKind::incident;
  1758. if (pending) {
  1759. const IncidentRecord* const record =
  1760. find_incident(binding, instance->view.incidentRevision);
  1761. if (record != nullptr && record->status != IncidentStatus::canceled) {
  1762. output.incident = record->incident;
  1763. output.revision = record->revision;
  1764. }
  1765. }
  1766. ReleaseSRWLockShared(&g_lock);
  1767. return pending && output.revision != 0;
  1768. }
  1769. /** Cancels the exact instance's raw incident before any later transport staging. */
  1770. bool cancel_pending_incident(const state::activity::SessionBinding& binding) noexcept {
  1771. AcquireSRWLockExclusive(&g_lock);
  1772. Instance* const instance = find_instance(binding);
  1773. const bool canceled = instance != nullptr && instance->view.outputPending
  1774. && instance->view.outputKind == OutputKind::incident;
  1775. if (canceled) {
  1776. cancel_output(*instance, GetTickCount64());
  1777. }
  1778. ReleaseSRWLockExclusive(&g_lock);
  1779. return canceled;
  1780. }
  1781. /** Records one refused BAP attempt without clearing the committed output. */
  1782. void note_auth_attempt(const state::activity::SessionBinding& binding,
  1783. std::uint64_t sourceGeneration,
  1784. std::uint64_t revision,
  1785. std::uint8_t lifetimeState,
  1786. OutputStatus status) noexcept {
  1787. if (revision == 0 || status == OutputStatus::idle || status == OutputStatus::pending
  1788. || status == OutputStatus::transportStaged || status == OutputStatus::canceled) {
  1789. return;
  1790. }
  1791. AcquireSRWLockExclusive(&g_lock);
  1792. Instance* const instance = find_instance(binding);
  1793. if (instance != nullptr && instance->view.outputPending
  1794. && instance->view.outputKind == OutputKind::authState
  1795. && instance->view.stateRevision == revision
  1796. && instance->view.lifetimeState == lifetimeState) {
  1797. instance->view.lastOutputAttemptTick = GetTickCount64();
  1798. instance->view.lastOutputSourceGeneration = sourceGeneration;
  1799. ++instance->view.outputAttempts;
  1800. instance->view.outputStatus = status;
  1801. }
  1802. ReleaseSRWLockExclusive(&g_lock);
  1803. }
  1804. /** Marks one committed Auth revision as staged into the transport output queue. */
  1805. void note_auth_transport_staged(const state::activity::SessionBinding& binding,
  1806. std::uint64_t sourceGeneration,
  1807. std::uint64_t revision,
  1808. std::uint8_t lifetimeState) noexcept {
  1809. if (revision == 0) {
  1810. return;
  1811. }
  1812. AcquireSRWLockExclusive(&g_lock);
  1813. Instance* const instance = find_instance(binding);
  1814. if (instance != nullptr && instance->view.outputPending
  1815. && instance->view.outputKind == OutputKind::authState
  1816. && instance->view.stateRevision == revision
  1817. && instance->view.lifetimeState == lifetimeState) {
  1818. instance->view.transportRevision = revision;
  1819. instance->view.lastOutputAttemptTick = GetTickCount64();
  1820. instance->view.lastOutputSourceGeneration = sourceGeneration;
  1821. ++instance->view.outputAttempts;
  1822. instance->view.outputStatus = OutputStatus::transportStaged;
  1823. instance->view.outputPending = false;
  1824. instance->view.outputKind = OutputKind::none;
  1825. Event event{};
  1826. event.binding = binding;
  1827. event.tick = instance->view.lastOutputAttemptTick;
  1828. event.kind = EventKind::authStateTransportStaged;
  1829. event.sourceGeneration = sourceGeneration;
  1830. event.stateRevision = revision;
  1831. event.lifetimeState = lifetimeState;
  1832. append_event(event);
  1833. instance->view.lastEventSequence = g_sequence;
  1834. }
  1835. ReleaseSRWLockExclusive(&g_lock);
  1836. }
  1837. /** Records one failed incident attempt without clearing the serialized output slot. */
  1838. void note_incident_attempt(const state::activity::SessionBinding& binding,
  1839. std::uint64_t sourceGeneration,
  1840. std::uint64_t revision,
  1841. IncidentStatus status) noexcept {
  1842. if (revision == 0
  1843. || (status != IncidentStatus::encodeFailed && status != IncidentStatus::frameRefused)) {
  1844. return;
  1845. }
  1846. AcquireSRWLockExclusive(&g_lock);
  1847. Instance* const instance = find_instance(binding);
  1848. IncidentRecord* const record = find_incident(binding, revision);
  1849. if (instance != nullptr && record != nullptr && instance->view.outputPending
  1850. && instance->view.outputKind == OutputKind::incident
  1851. && instance->view.incidentRevision == revision) {
  1852. record->lastAttemptTick = GetTickCount64();
  1853. record->lastSourceGeneration = sourceGeneration;
  1854. ++record->attempts;
  1855. record->status = status;
  1856. instance->view.lastOutputAttemptTick = record->lastAttemptTick;
  1857. instance->view.lastOutputSourceGeneration = sourceGeneration;
  1858. ++instance->view.outputAttempts;
  1859. instance->view.outputStatus = OutputStatus::frameRefused;
  1860. }
  1861. ReleaseSRWLockExclusive(&g_lock);
  1862. }
  1863. /** Marks one retained incident revision as staged into a matching transport output queue. */
  1864. void note_incident_transport_staged(const state::activity::SessionBinding& binding,
  1865. std::uint64_t sourceGeneration,
  1866. std::uint64_t revision) noexcept {
  1867. if (revision == 0) {
  1868. return;
  1869. }
  1870. AcquireSRWLockExclusive(&g_lock);
  1871. Instance* const instance = find_instance(binding);
  1872. IncidentRecord* const record = find_incident(binding, revision);
  1873. if (instance != nullptr && record != nullptr && instance->view.outputPending
  1874. && instance->view.outputKind == OutputKind::incident
  1875. && instance->view.incidentRevision == revision) {
  1876. record->lastAttemptTick = GetTickCount64();
  1877. record->lastSourceGeneration = sourceGeneration;
  1878. ++record->attempts;
  1879. ++record->transportStages;
  1880. record->status = IncidentStatus::transportStaged;
  1881. instance->view.incidentTransportRevision = revision;
  1882. instance->view.lastOutputAttemptTick = record->lastAttemptTick;
  1883. instance->view.lastOutputSourceGeneration = sourceGeneration;
  1884. ++instance->view.outputAttempts;
  1885. instance->view.outputStatus = OutputStatus::transportStaged;
  1886. instance->view.outputPending = false;
  1887. instance->view.outputKind = OutputKind::none;
  1888. instance->view.incidentsPending = 0;
  1889. Event event{};
  1890. event.binding = binding;
  1891. event.tick = record->lastAttemptTick;
  1892. event.kind = EventKind::incidentTransportStaged;
  1893. event.sourceGeneration = sourceGeneration;
  1894. fill_incident_event(event, record->incident, revision);
  1895. append_event(event);
  1896. instance->view.lastEventSequence = g_sequence;
  1897. }
  1898. ReleaseSRWLockExclusive(&g_lock);
  1899. }
  1900. /** Clears every instance, queue, and diagnostic counter. */
  1901. void reset() noexcept {
  1902. AcquireSRWLockExclusive(&g_lock);
  1903. g_eventGeneration = next_nonzero(g_eventGeneration);
  1904. for (Instance& instance : g_instances) {
  1905. clear_instance(instance);
  1906. }
  1907. std::vector<PendingInput>{}.swap(g_pending);
  1908. for (Event& event : g_events) {
  1909. event = {};
  1910. }
  1911. std::vector<MissionInputRecord>{}.swap(g_missionInputs);
  1912. for (IncidentRecord& incident : g_incidents) {
  1913. incident = {};
  1914. }
  1915. for (ClientMessageRecord& message : g_clientMessages) {
  1916. message = {};
  1917. }
  1918. SecureZeroMemory(g_clientMessageDetails.data(), sizeof(g_clientMessageDetails));
  1919. g_pendingRead = 0;
  1920. g_queuedIngress = 0;
  1921. g_queuedControls = 0;
  1922. g_eventStart = 0;
  1923. g_eventCount = 0;
  1924. g_clientMessageStart = 0;
  1925. g_clientMessageCount = 0;
  1926. g_clientMessageDetailStart = 0;
  1927. g_clientMessageDetailCount = 0;
  1928. g_sequence = 0;
  1929. g_scriptableReservationGeneration = next_nonzero(g_scriptableReservationGeneration);
  1930. g_scriptableReservationSequence = 0;
  1931. g_missionInputSequence = 0;
  1932. g_clientMessageSequence = 0;
  1933. g_touch = 0;
  1934. g_droppedIngress = 0;
  1935. g_droppedIncidents = 0;
  1936. g_refusedControls = 0;
  1937. g_refusedIncidents = 0;
  1938. g_overwrittenEvents = 0;
  1939. g_overwrittenIncidents = 0;
  1940. g_overwrittenClientMessages = 0;
  1941. ReleaseSRWLockExclusive(&g_lock);
  1942. }
  1943. } // namespace sunrise::server::activity::host