mission_script_runtime_feed.cpp 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749
  1. /**
  2. * The two Host feeds, the queue of accepted rows, and the VM callback each row reaches.
  3. * Every function here runs under the mission runtime lock its caller already holds.
  4. */
  5. #include <Windows.h>
  6. #include <algorithm>
  7. #include <array>
  8. #include <cstddef>
  9. #include <cstdint>
  10. #include <cstdio>
  11. #include <limits>
  12. #include <new>
  13. #include <string_view>
  14. #include <vector>
  15. #include "../../../core/logging/log.h"
  16. #include "../../../state/activity/mission/runtime.h"
  17. #include "../../../state/activity/runtime.h"
  18. #include "../host_runtime.h"
  19. #include "mission_script_region.h"
  20. #include "mission_script_runtime_internal.h"
  21. #include "mission_script_vm.h"
  22. namespace sunrise::server::activity::mission {
  23. namespace {
  24. /** One queued mission event and the retry bookkeeping the drain uses. */
  25. struct PendingMissionEvent final {
  26. host::Event event{};
  27. host::SenseObservationSnapshot sense{};
  28. host::ClientMessageSnapshot clientMessage{};
  29. std::uint64_t firstAttempt{};
  30. std::uint64_t nextAttempt{};
  31. std::uint32_t attempts{};
  32. bool missionSequenceObserved{};
  33. bool senseAvailable{};
  34. bool clientMessageAvailable{};
  35. bool occupied{};
  36. };
  37. /** The queue grows with real accepted input and is drained in the tick it arrives. */
  38. std::vector<PendingMissionEvent> g_pendingMissionEvents{};
  39. host::EventCursor g_eventCursor{};
  40. host::MissionInputCursor g_missionInputCursor{};
  41. void clear_pending_event(PendingMissionEvent& pending) noexcept {
  42. SecureZeroMemory(&pending, sizeof(pending));
  43. }
  44. /** @return True for the host events that report an output's progress, not an input. */
  45. [[nodiscard]] bool delivery_lifecycle_event(host::EventKind kind) noexcept {
  46. switch (kind) {
  47. case host::EventKind::authStateCommitted:
  48. case host::EventKind::authStateTransportStaged:
  49. case host::EventKind::authStateCanceled:
  50. case host::EventKind::incidentQueued:
  51. case host::EventKind::incidentTransportStaged:
  52. case host::EventKind::incidentCanceled:
  53. case host::EventKind::incidentRefused:
  54. case host::EventKind::scriptableOverrideCommitted:
  55. case host::EventKind::scriptableOverrideTransportStaged:
  56. case host::EventKind::scriptableOverrideCanceled:
  57. case host::EventKind::operatorRefused:
  58. return true;
  59. default:
  60. return false;
  61. }
  62. }
  63. /**
  64. * @return True for a row that arrives on the ordered mission-input feed and owns a sequence.
  65. * A host-state row must never answer true. It would consume a mission-input sequence it does not
  66. * own, which faults the binding on the next real input.
  67. */
  68. [[nodiscard]] bool host_feed_row(host::EventKind kind) noexcept {
  69. switch (kind) {
  70. case host::EventKind::timerElapsed:
  71. case host::EventKind::effectResult:
  72. case host::EventKind::phaseEntered:
  73. case host::EventKind::triggerEntered:
  74. case host::EventKind::triggerExited:
  75. case host::EventKind::squadState:
  76. case host::EventKind::squadProvoked:
  77. case host::EventKind::entitySpawned:
  78. case host::EventKind::entityDied:
  79. case host::EventKind::sceneFinished:
  80. case host::EventKind::objectiveProgress:
  81. case host::EventKind::sessionJoined:
  82. case host::EventKind::sessionLeft:
  83. case host::EventKind::playerTrigger:
  84. case host::EventKind::cinematicStarted:
  85. case host::EventKind::cinematicSkipRequested:
  86. case host::EventKind::cinematicTerminated:
  87. case host::EventKind::actorPathState:
  88. case host::EventKind::damageState:
  89. case host::EventKind::deviceState:
  90. case host::EventKind::regionChanged:
  91. case host::EventKind::objectState:
  92. case host::EventKind::fireteamState:
  93. case host::EventKind::objectInteracted:
  94. case host::EventKind::ghostLinkState:
  95. return false;
  96. default:
  97. return true;
  98. }
  99. }
  100. /** True when the event may reach a callback for this instance's ActivityClient generation. */
  101. [[nodiscard]] bool eligible_event(const RuntimeInstance& instance,
  102. const host::Event& event) noexcept {
  103. if (event.attemptGeneration != 0 && event.attemptGeneration != instance.attempt.generation) {
  104. return false;
  105. }
  106. if (event.kind == host::EventKind::timerElapsed) {
  107. return true;
  108. }
  109. if (event.sourceGeneration != instance.view.activityClientGeneration) {
  110. return false;
  111. }
  112. return event.kind == host::EventKind::clientStateChanged
  113. || event.kind == host::EventKind::incidentReceived
  114. || event.kind == host::EventKind::clientMessageReceived
  115. || event.kind == host::EventKind::effectResult
  116. || event.kind == host::EventKind::phaseEntered
  117. || event.kind == host::EventKind::triggerEntered
  118. || event.kind == host::EventKind::triggerExited
  119. || event.kind == host::EventKind::squadState
  120. || event.kind == host::EventKind::squadProvoked
  121. || event.kind == host::EventKind::entitySpawned
  122. || event.kind == host::EventKind::entityDied
  123. || event.kind == host::EventKind::sceneFinished
  124. || event.kind == host::EventKind::objectiveProgress
  125. || event.kind == host::EventKind::entitySlotsRequested
  126. || event.kind == host::EventKind::sessionJoined
  127. || event.kind == host::EventKind::sessionLeft
  128. || event.kind == host::EventKind::playerTrigger
  129. || event.kind == host::EventKind::cinematicStarted
  130. || event.kind == host::EventKind::actorPathState
  131. || event.kind == host::EventKind::damageState
  132. || event.kind == host::EventKind::deviceState
  133. || event.kind == host::EventKind::regionChanged
  134. || event.kind == host::EventKind::objectState
  135. || event.kind == host::EventKind::fireteamState
  136. || event.kind == host::EventKind::objectInteracted
  137. || event.kind == host::EventKind::ghostLinkState
  138. || event.kind == host::EventKind::cinematicSkipRequested
  139. || event.kind == host::EventKind::cinematicTerminated
  140. || delivery_lifecycle_event(event.kind) || event.has_sense_observations();
  141. }
  142. /** Faults the instance unless the ordered mission input arrives with no gap, starting at one. */
  143. [[nodiscard]] bool validate_mission_sequence(RuntimeInstance& instance,
  144. const host::Event& event) noexcept {
  145. if (event.kind != host::EventKind::senseUpdate
  146. && event.kind != host::EventKind::incidentReceived
  147. && event.kind != host::EventKind::clientStateChanged
  148. && event.kind != host::EventKind::entitySlotsRequested
  149. && event.kind != host::EventKind::clientMessageReceived) {
  150. return true;
  151. }
  152. if (instance.lastMissionSequence == 0) {
  153. if (event.missionSequence == 1) {
  154. return true;
  155. }
  156. fault_instance(instance, "activity mission input did not start at sequence one");
  157. log_line(core::log::Level::warn, &instance, "events", "initial_binding_gap");
  158. return false;
  159. }
  160. const std::uint64_t expected =
  161. instance.lastMissionSequence == (std::numeric_limits<std::uint64_t>::max)()
  162. ? 1
  163. : instance.lastMissionSequence + 1;
  164. if (event.missionSequence == expected) {
  165. return true;
  166. }
  167. fault_instance(instance, "activity mission input sequence has a gap");
  168. log_line(core::log::Level::warn, &instance, "events", "binding_gap");
  169. return false;
  170. }
  171. /** One free queue row, growing the queue by one when none is free; null when it cannot grow. */
  172. [[nodiscard]] PendingMissionEvent* free_pending_event() noexcept {
  173. for (PendingMissionEvent& pending : g_pendingMissionEvents) {
  174. if (!pending.occupied) {
  175. return &pending;
  176. }
  177. }
  178. if (g_pendingMissionEvents.size() == g_pendingMissionEvents.max_size()) {
  179. return nullptr;
  180. }
  181. try {
  182. g_pendingMissionEvents.emplace_back();
  183. } catch (const std::bad_alloc&) {
  184. return nullptr;
  185. }
  186. return &g_pendingMissionEvents.back();
  187. }
  188. /** Copies one accepted input and its values into a queue row; faults the instance when full. */
  189. [[nodiscard]] bool queue_mission_event(RuntimeInstance* instance,
  190. const host::MissionInputEvent& input) noexcept {
  191. PendingMissionEvent* const pending = free_pending_event();
  192. if (pending == nullptr) {
  193. if (instance != nullptr) {
  194. fault_instance(*instance, "mission event queue allocation failed");
  195. clear_pending_events(instance->view.binding);
  196. }
  197. log_line(core::log::Level::warn, instance, "events", "allocation_failed");
  198. return false;
  199. }
  200. clear_pending_event(*pending);
  201. if (input.event.has_sense_observations()) {
  202. pending->senseAvailable =
  203. host::mission_input_sense_snapshot(input.sequence, pending->sense);
  204. }
  205. if (input.event.kind == host::EventKind::clientMessageReceived) {
  206. pending->clientMessageAvailable =
  207. host::mission_input_client_message_snapshot(input.sequence, pending->clientMessage);
  208. }
  209. pending->nextAttempt = 0;
  210. pending->event = input.event;
  211. pending->occupied = true;
  212. return true;
  213. }
  214. /** @return True when the left sequence comes before the right one across the wrap. */
  215. [[nodiscard]] constexpr bool mission_sequence_precedes(std::uint64_t left,
  216. std::uint64_t right) noexcept {
  217. // Sequences wrap, so ordering holds only inside half the 64-bit range.
  218. constexpr std::uint64_t halfRange = std::uint64_t{1} << 63U;
  219. return left != right && right - left < halfRange;
  220. }
  221. static_assert(mission_sequence_precedes((std::numeric_limits<std::uint64_t>::max)(), 1));
  222. static_assert(!mission_sequence_precedes(1, (std::numeric_limits<std::uint64_t>::max)()));
  223. /** @return True when the durable cursor already committed this retained input row. */
  224. [[nodiscard]] bool mission_sequence_committed(std::uint64_t sequence,
  225. std::uint64_t committed) noexcept {
  226. return committed != 0
  227. && (sequence == committed || mission_sequence_precedes(sequence, committed));
  228. }
  229. /** True when the same binding still holds a queued row with an earlier mission sequence. */
  230. [[nodiscard]] bool has_earlier_pending_event(const PendingMissionEvent& selected) noexcept {
  231. for (const PendingMissionEvent& pending : g_pendingMissionEvents) {
  232. if (pending.occupied && &pending != &selected
  233. && same_binding(pending.event.binding, selected.event.binding)
  234. && mission_sequence_precedes(pending.event.missionSequence,
  235. selected.event.missionSequence)) {
  236. return true;
  237. }
  238. }
  239. return false;
  240. }
  241. /** TODO: no caller. Decide whether the drain gate retires on this or on `binding_matches`. */
  242. [[nodiscard]] bool host_binding_active(const state::activity::SessionBinding& binding) noexcept {
  243. host::InstanceSnapshot snapshot{};
  244. return host::instance_snapshot(binding, snapshot) && snapshot.active;
  245. }
  246. /** Dispatches queued rows in mission-sequence order and retires those no callback can take. */
  247. void drain_pending_mission_events(std::uint64_t now) noexcept {
  248. bool progressed = false;
  249. do {
  250. progressed = false;
  251. for (PendingMissionEvent& pending : g_pendingMissionEvents) {
  252. if (!pending.occupied || now < pending.nextAttempt
  253. || has_earlier_pending_event(pending)) {
  254. continue;
  255. }
  256. RuntimeInstance* const instance = find_instance(pending.event.binding);
  257. if (instance == nullptr) {
  258. if (!state::activity::binding_matches(pending.event.binding)) {
  259. clear_pending_event(pending);
  260. progressed = true;
  261. }
  262. continue;
  263. }
  264. if (instance->programStatus != ProgramStatus::loaded) {
  265. clear_pending_event(pending);
  266. progressed = true;
  267. continue;
  268. }
  269. if (instance->startPending) {
  270. continue;
  271. }
  272. if (instance->timerPending) {
  273. continue;
  274. }
  275. if (mission_sequence_committed(pending.event.missionSequence,
  276. instance->lastMissionSequence)) {
  277. clear_pending_event(pending);
  278. progressed = true;
  279. continue;
  280. }
  281. if (!pending.missionSequenceObserved) {
  282. if (!validate_mission_sequence(*instance, pending.event)) {
  283. clear_pending_events(instance->view.binding);
  284. progressed = true;
  285. break;
  286. }
  287. pending.missionSequenceObserved = true;
  288. }
  289. if (!eligible_event(*instance, pending.event)) {
  290. if (!commit_mission_state(
  291. *instance, instance->missionStarted, pending.event.missionSequence)) {
  292. clear_pending_events(instance->view.binding);
  293. progressed = true;
  294. break;
  295. }
  296. clear_pending_event(pending);
  297. progressed = true;
  298. continue;
  299. }
  300. if (pending.event.kind == host::EventKind::senseUpdate && !pending.senseAvailable) {
  301. fault_instance(*instance, "accepted mission Sense values were unavailable");
  302. clear_pending_events(instance->view.binding);
  303. log_line(core::log::Level::warn, instance, "events", "sense_unavailable");
  304. progressed = true;
  305. break;
  306. }
  307. if (pending.event.kind == host::EventKind::clientMessageReceived
  308. && !pending.clientMessageAvailable) {
  309. fault_instance(*instance,
  310. "accepted mission client-message values were unavailable");
  311. clear_pending_events(instance->view.binding);
  312. log_line(core::log::Level::warn, instance, "events", "client_message_unavailable");
  313. progressed = true;
  314. break;
  315. }
  316. lua_vm::Intent intent{};
  317. if (instance->deliveryStage != DeliveryStage::idle
  318. || lua_vm::pending_intent(instance->vm, intent)) {
  319. continue;
  320. }
  321. const bool firstAttempt = pending.attempts == 0;
  322. if (firstAttempt) {
  323. pending.firstAttempt = now;
  324. }
  325. ++pending.attempts;
  326. const host::SenseObservationSnapshot* const sense =
  327. pending.event.kind == host::EventKind::senseUpdate && pending.senseAvailable
  328. ? &pending.sense
  329. : nullptr;
  330. const host::ClientMessageSnapshot* const clientMessage =
  331. pending.event.kind == host::EventKind::clientMessageReceived
  332. && pending.clientMessageAvailable
  333. ? &pending.clientMessage
  334. : nullptr;
  335. static_cast<void>(
  336. dispatch_event(*instance, pending.event, sense, clientMessage, firstAttempt, now));
  337. clear_pending_event(pending);
  338. progressed = true;
  339. }
  340. } while (progressed);
  341. }
  342. /** @return True when one exact accepted sequence is already retained locally or in this read. */
  343. [[nodiscard]] bool input_sequence_retained(const state::activity::SessionBinding& binding,
  344. std::uint64_t sequence,
  345. const host::MissionInputRead& inputs) noexcept {
  346. for (const PendingMissionEvent& pending : g_pendingMissionEvents) {
  347. if (pending.occupied && same_binding(pending.event.binding, binding)
  348. && pending.event.missionSequence == sequence) {
  349. return true;
  350. }
  351. }
  352. for (std::size_t index = 0; index < inputs.count; ++index) {
  353. if (same_binding(inputs.events[index].event.binding, binding)
  354. && inputs.events[index].event.missionSequence == sequence) {
  355. return true;
  356. }
  357. }
  358. return false;
  359. }
  360. /** @return True when every accepted but uncommitted sequence is present in this read or local
  361. * queue. */
  362. [[nodiscard]] bool
  363. outstanding_input_interval_complete(const state::activity::SessionBinding& binding,
  364. const mission_state::InputSequenceSnapshot& state,
  365. const host::MissionInputRead& inputs) noexcept {
  366. if (state.issued < state.committed) {
  367. return false;
  368. }
  369. const std::uint64_t outstanding = state.issued - state.committed;
  370. if (outstanding > g_pendingMissionEvents.size() + inputs.count) {
  371. return false;
  372. }
  373. std::uint64_t sequence = state.committed;
  374. for (std::uint64_t index = 0; index < outstanding; ++index) {
  375. ++sequence;
  376. if (!input_sequence_retained(binding, sequence, inputs)) {
  377. return false;
  378. }
  379. }
  380. return true;
  381. }
  382. /** Faults only retained bindings whose durable uncommitted interval is provably incomplete. */
  383. void reconcile_input_feed_loss(const host::MissionInputRead& inputs) noexcept {
  384. host::DiagnosticsSnapshot hostState{};
  385. host::snapshot(hostState);
  386. for (std::size_t index = 0; index < hostState.instanceCount; ++index) {
  387. const host::InstanceSnapshot& hostInstance = hostState.instances[index];
  388. if (!state::activity::binding_matches(hostInstance.binding)) {
  389. continue;
  390. }
  391. mission_state::InputSequenceSnapshot inputState{};
  392. if (!mission_state::input_sequence_snapshot(hostInstance.binding, inputState)
  393. || inputState.faulted
  394. || outstanding_input_interval_complete(hostInstance.binding, inputState, inputs)) {
  395. continue;
  396. }
  397. mission_state::Snapshot snapshot{};
  398. const mission_state::Status status =
  399. mission_state::fault_input_feed(hostInstance.binding, snapshot);
  400. RuntimeInstance* const instance = find_instance(hostInstance.binding);
  401. if (status != mission_state::Status::ready) {
  402. log_line(core::log::Level::warn,
  403. instance,
  404. "events",
  405. mission_state::status_name(status),
  406. "reason=feed_gap_fault_refused");
  407. continue;
  408. }
  409. if (instance != nullptr) {
  410. accept_mission_state(*instance, snapshot);
  411. lua_vm::fault(instance->vm, "accepted mission input feed lost a row");
  412. instance->programStatus = ProgramStatus::programError;
  413. }
  414. log_line(core::log::Level::warn, instance, "events", "feed_gap_faulted");
  415. clear_pending_events(hostInstance.binding);
  416. }
  417. }
  418. /** Reads one page of the ordered feed and queues every row not yet committed. */
  419. [[nodiscard]] bool consume_mission_input_page() noexcept {
  420. host::MissionInputRead inputs{};
  421. host::read_mission_inputs_after(g_missionInputCursor, inputs);
  422. if (inputs.reset) {
  423. g_missionInputCursor = {inputs.cursor.generation, 0};
  424. reconcile_input_feed_loss(inputs);
  425. log_line(core::log::Level::warn, nullptr, "events", "mission_feed_reset");
  426. }
  427. if (inputs.gap) {
  428. log_line(core::log::Level::warn, nullptr, "events", "mission_feed_gap");
  429. reconcile_input_feed_loss(inputs);
  430. }
  431. for (std::size_t index = 0; index < inputs.count; ++index) {
  432. const host::MissionInputEvent& input = inputs.events[index];
  433. RuntimeInstance* const instance = find_instance(input.event.binding);
  434. mission_state::InputSequenceSnapshot inputState{};
  435. const bool hasInputState =
  436. mission_state::input_sequence_snapshot(input.event.binding, inputState);
  437. if ((instance == nullptr && !state::activity::binding_matches(input.event.binding))
  438. || (instance != nullptr && instance->programStatus != ProgramStatus::loaded)
  439. || (hasInputState
  440. && (inputState.faulted
  441. || mission_sequence_committed(input.event.missionSequence,
  442. inputState.committed)))) {
  443. g_missionInputCursor.generation = inputs.cursor.generation;
  444. g_missionInputCursor.sequence = input.sequence;
  445. continue;
  446. }
  447. if (instance != nullptr
  448. && mission_sequence_committed(input.event.missionSequence,
  449. instance->lastMissionSequence)) {
  450. g_missionInputCursor.generation = inputs.cursor.generation;
  451. g_missionInputCursor.sequence = input.sequence;
  452. continue;
  453. }
  454. if (instance != nullptr) {
  455. lua_vm::Snapshot diagnostics{};
  456. lua_vm::snapshot(instance->vm, diagnostics);
  457. if (diagnostics.faulted) {
  458. g_missionInputCursor.generation = inputs.cursor.generation;
  459. g_missionInputCursor.sequence = input.sequence;
  460. continue;
  461. }
  462. }
  463. if (!queue_mission_event(instance, input)) {
  464. return false;
  465. }
  466. g_missionInputCursor.generation = inputs.cursor.generation;
  467. g_missionInputCursor.sequence = input.sequence;
  468. }
  469. return inputs.count == host::kMissionInputReadPageSize;
  470. }
  471. } // namespace
  472. /** Retires every queued mission event that belongs to one binding. */
  473. void clear_pending_events(const state::activity::SessionBinding& binding) noexcept {
  474. for (PendingMissionEvent& pending : g_pendingMissionEvents) {
  475. if (pending.occupied && same_binding(pending.event.binding, binding)) {
  476. clear_pending_event(pending);
  477. }
  478. }
  479. }
  480. /** Retires every queued mission event that belongs to one binding. */
  481. void clear_all_pending_events() noexcept {
  482. std::vector<PendingMissionEvent>{}.swap(g_pendingMissionEvents);
  483. }
  484. /** Keeps accepted values but removes every VM- and ActivityClient-generation-local decision. */
  485. void reset_pending_events_for_reattach(const state::activity::SessionBinding& binding) noexcept {
  486. for (PendingMissionEvent& pending : g_pendingMissionEvents) {
  487. if (!pending.occupied || !same_binding(pending.event.binding, binding)) {
  488. continue;
  489. }
  490. pending.firstAttempt = 0;
  491. pending.nextAttempt = 0;
  492. pending.attempts = 0;
  493. pending.missionSequenceObserved = false;
  494. }
  495. }
  496. /** Retires rows as soon as authoritative State no longer owns their exact binding. */
  497. void retire_unbound_pending_events() noexcept {
  498. for (PendingMissionEvent& pending : g_pendingMissionEvents) {
  499. if (pending.occupied && !state::activity::binding_matches(pending.event.binding)) {
  500. clear_pending_event(pending);
  501. }
  502. }
  503. }
  504. /** @return Queued rows still owed to one exact binding. */
  505. std::size_t pending_event_count(const state::activity::SessionBinding& binding) noexcept {
  506. std::size_t count = 0;
  507. for (const PendingMissionEvent& pending : g_pendingMissionEvents) {
  508. count += pending.occupied && same_binding(pending.event.binding, binding) ? 1U : 0U;
  509. }
  510. return count;
  511. }
  512. /** @return True when one queued accepted row is still owed to this instance. */
  513. bool has_pending_host_input(const RuntimeInstance& instance) noexcept {
  514. return std::any_of(g_pendingMissionEvents.begin(),
  515. g_pendingMissionEvents.end(),
  516. [&instance](const PendingMissionEvent& pending) noexcept {
  517. return pending.occupied
  518. && same_binding(pending.event.binding, instance.view.binding);
  519. });
  520. }
  521. /** Points both Host cursors at the current feed heads and replays retained accepted inputs. */
  522. void reset_feed_cursors() noexcept {
  523. g_eventCursor = host::current_event_cursor();
  524. g_missionInputCursor = host::current_mission_input_cursor();
  525. // Durable per-binding cursors suppress callbacks that already committed; an evicted
  526. // predecessor still faults as a gap.
  527. g_missionInputCursor.sequence = 0;
  528. }
  529. /** Clears both Host cursors. */
  530. void clear_feed_cursors() noexcept {
  531. g_eventCursor = {};
  532. g_missionInputCursor = {};
  533. }
  534. /**
  535. * Runs one event through the VM, commits what it changed, and faults on a script failure.
  536. * @param sense Values owned by a Sense row, or null.
  537. * @param clientMessage Envelope snapshot owned by a client-message row, or null.
  538. * @param firstAttempt True on the first delivery attempt, which is the one that counts it.
  539. * @return The VM call status; `inactive` when no callback could take the event.
  540. */
  541. lua_vm::CallStatus dispatch_event(RuntimeInstance& instance,
  542. const host::Event& event,
  543. const host::SenseObservationSnapshot* sense,
  544. const host::ClientMessageSnapshot* clientMessage,
  545. bool firstAttempt,
  546. std::uint64_t now) noexcept {
  547. if (instance.programStatus == ProgramStatus::missing
  548. && event.kind == host::EventKind::senseUpdate && sense != nullptr) {
  549. push_squad_edges(instance, *sense);
  550. return lua_vm::CallStatus::inactive;
  551. }
  552. if (instance.programStatus != ProgramStatus::loaded || !eligible_event(instance, event)) {
  553. return lua_vm::CallStatus::inactive;
  554. }
  555. if (firstAttempt) {
  556. ++instance.eventsSeen;
  557. instance.lastEventSequence = event.sequence;
  558. }
  559. // The three region numbers the script is about to read. A report restates only the leg it
  560. // moved, so `pending` and `current` read -1 on most reports and `held` is the one that says
  561. // where the client is standing.
  562. if (firstAttempt && event.kind == host::EventKind::clientStateChanged) {
  563. std::array<char, 64> legs{};
  564. const int written =
  565. std::snprintf(legs.data(),
  566. legs.size(),
  567. "held=%d pending=%d current=%d",
  568. event.heldRegionIndex,
  569. event.clientStateHasRegion ? event.regionIndex : -1,
  570. event.clientStateHasCurrentRegion ? event.currentRegionIndex : -1);
  571. if (written > 0) {
  572. log_line(core::log::Level::debug,
  573. &instance,
  574. "client_state",
  575. "legs",
  576. {legs.data(), static_cast<std::size_t>(written)});
  577. }
  578. }
  579. instance.dispatchAttemptGeneration = event.attemptGeneration;
  580. instance.dispatchInputSequence = event.missionSequence;
  581. const lua_vm::CallStatus status = lua_vm::dispatch(instance.vm, event, clientMessage, now);
  582. if (event.kind == host::EventKind::clientStateChanged) {
  583. // A pending-region report can name the next slice while the player still holds the old
  584. // one, so the held region wins.
  585. if (event.heldRegionIndex >= 0) {
  586. instance.activeRegion = event.heldRegionIndex;
  587. } else if (event.currentRegionIndex >= 0) {
  588. instance.activeRegion = event.currentRegionIndex;
  589. }
  590. host::Event changed{};
  591. if (firstAttempt && make_region_changed(event, changed)) {
  592. push_script_event(instance, changed);
  593. }
  594. }
  595. if (firstAttempt && event.kind == host::EventKind::incidentReceived) {
  596. push_player_trigger(instance, event);
  597. push_cinematic(instance, event);
  598. }
  599. if (event.kind == host::EventKind::senseUpdate && sense != nullptr) {
  600. if (firstAttempt) {
  601. observe_player_life(instance, *sense);
  602. publish_fireteam_life(now);
  603. }
  604. push_trigger_edges(instance, *sense);
  605. push_ghost_edges(instance, *sense);
  606. push_object_interaction_edges(instance, *sense);
  607. push_damage_edges(instance, *sense);
  608. push_actor_path_edges(instance, *sense);
  609. push_squad_edges(instance, *sense);
  610. push_combatant_damage_edges(instance, *sense);
  611. push_device_edges(instance, *sense);
  612. push_scene_edges(instance, *sense);
  613. push_objective_edges(instance, *sense);
  614. }
  615. instance.dispatchAttemptGeneration = 0;
  616. instance.dispatchInputSequence = 0;
  617. note_vm_status(instance, "event", lua_vm::status_name(status));
  618. if (firstAttempt && instance.eventsSeen == 1) {
  619. log_line(core::log::Level::info,
  620. &instance,
  621. "dispatch",
  622. lua_vm::status_name(status),
  623. event.kind == host::EventKind::senseUpdate ? "detail=sense"
  624. : event.kind == host::EventKind::clientStateChanged ? "detail=client_state"
  625. : event.kind == host::EventKind::incidentReceived ? "detail=incident"
  626. : event.kind == host::EventKind::entitySlotsRequested
  627. ? "detail=entity_slots_requested"
  628. : event.kind == host::EventKind::timerElapsed ? "detail=timer"
  629. : event.kind == host::EventKind::effectResult ? "detail=effect_result"
  630. : "detail=client_message");
  631. }
  632. const std::uint64_t nextInputSequence =
  633. host_feed_row(event.kind) ? event.missionSequence : instance.lastMissionSequence;
  634. if (status == lua_vm::CallStatus::committed) {
  635. if (!commit_mission_state(instance, true, nextInputSequence)) {
  636. return lua_vm::CallStatus::scriptError;
  637. }
  638. ++instance.eventsCommitted;
  639. lua_vm::Snapshot diagnostics{};
  640. lua_vm::snapshot(instance.vm, diagnostics);
  641. if (event.kind == host::EventKind::clientStateChanged
  642. || diagnostics.stateRevision != instance.lastLoggedRevision) {
  643. instance.lastLoggedRevision = diagnostics.stateRevision;
  644. log_line(core::log::Level::debug,
  645. &instance,
  646. "event",
  647. "committed",
  648. event.kind == host::EventKind::senseUpdate ? "detail=sense"
  649. : event.kind == host::EventKind::clientStateChanged ? "detail=client_state"
  650. : event.kind == host::EventKind::incidentReceived ? "detail=incident"
  651. : event.kind == host::EventKind::entitySlotsRequested
  652. ? "detail=entity_slots_requested"
  653. : event.kind == host::EventKind::timerElapsed ? "detail=timer"
  654. : event.kind == host::EventKind::effectResult ? "detail=effect_result"
  655. : "detail=client_message");
  656. }
  657. return status;
  658. }
  659. if (status == lua_vm::CallStatus::noHandler) {
  660. return commit_mission_state(instance, true, nextInputSequence)
  661. ? status
  662. : lua_vm::CallStatus::scriptError;
  663. }
  664. if (status == lua_vm::CallStatus::inactive) {
  665. return status;
  666. }
  667. lua_vm::Snapshot diagnostics{};
  668. lua_vm::snapshot(instance.vm, diagnostics);
  669. log_line(core::log::Level::warn,
  670. &instance,
  671. "event",
  672. lua_vm::status_name(status),
  673. {},
  674. diagnostics.lastError.data());
  675. persist_mission_fault(instance);
  676. return status;
  677. }
  678. /** Reads new host events, advances delivery, and queues only the lifecycle rows for scripts. */
  679. void consume_delivery_events(std::uint64_t now) noexcept {
  680. host::EventRead events{};
  681. host::read_events_after(g_eventCursor, events);
  682. g_eventCursor = events.cursor;
  683. if (events.reset) {
  684. log_line(core::log::Level::warn, nullptr, "delivery", "event_feed_reset");
  685. return;
  686. }
  687. if (events.gap) {
  688. log_line(core::log::Level::warn, nullptr, "delivery", "event_feed_gap");
  689. }
  690. for (std::size_t index = 0; index < events.count; ++index) {
  691. RuntimeInstance* const instance = find_instance(events.events[index].binding);
  692. if (instance != nullptr) {
  693. observe_delivery_event(*instance, events.events[index], now);
  694. // Sense, client-state, incident and client-message rows reach the script through the
  695. // ordered mission-input feed. Only the delivery lifecycle rows belong in this queue.
  696. if (delivery_lifecycle_event(events.events[index].kind)) {
  697. push_script_event(*instance, events.events[index]);
  698. }
  699. }
  700. }
  701. }
  702. /**
  703. * Reads the ordered mission-input feed and drains the queue. The feed is unbounded and one read
  704. * copies at most a page, so a burst is consumed in the tick it arrives instead of a page a tick.
  705. */
  706. void consume_mission_inputs(std::uint64_t now) noexcept {
  707. drain_pending_mission_events(now);
  708. while (consume_mission_input_page()) {
  709. drain_pending_mission_events(now);
  710. }
  711. drain_pending_mission_events(now);
  712. }
  713. } // namespace sunrise::server::activity::mission