peer_established.cpp 39 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853
  1. #include <Windows.h>
  2. #include <algorithm>
  3. #include <array>
  4. #include <cstdio>
  5. #include <memory>
  6. #include <new>
  7. #include "../../../middleware/encoding/bit_reader.h"
  8. #include "../../../middleware/encoding/bit_writer.h"
  9. #include "../../../middleware/gameplay/external/common_state.h"
  10. #include "../../../middleware/gameplay/external/control_state_codec.h"
  11. #include "../../../middleware/gameplay/peer/connect_messages.h"
  12. #include "../../../middleware/gameplay/peer/established_packet.h"
  13. #include "../../../middleware/gameplay/peer/packet_fragments.h"
  14. #include "../../../middleware/gameplay/peer/reliable_assembly.h"
  15. #include "../../bap/runtime.h"
  16. #include "../gameplay_log.h"
  17. #include "../group/group_host.h"
  18. #include "peer_packet_fragments.h"
  19. #include "peer_transport_internal.h"
  20. namespace sunrise::server::gameplay::peer {
  21. namespace {
  22. namespace gp = state::gameplay;
  23. namespace wire = middleware::gameplay::peer;
  24. namespace bits = middleware::encoding::bits;
  25. namespace fragments = wire::packet_fragments;
  26. /** Delay sentinel used until a round trip has been measured. */
  27. constexpr std::uint16_t kDelaySentinel = 1023;
  28. /** Smallest head-minus-cursor the peer accepts. This host keeps at most one packet in flight. */
  29. constexpr std::uint8_t kMinimumHeadCursor = 1;
  30. /** One packet cannot report more delivered messages than this. */
  31. constexpr std::size_t kMessageReportCapacity = 8;
  32. /**
  33. * Milliseconds between two resends of the same queue. The peer discards a packet more than 128
  34. * sequences ahead of its window, so this host must not send faster than the peer does.
  35. */
  36. constexpr std::uint64_t kResendInterval = 250;
  37. /** Reflected root 0x80806AE6 after its lane-presence bit. */
  38. constexpr std::size_t kPlayerSnapshotBits = 1373;
  39. /** Eight unique logical packet failures bound the replay sample. */
  40. constexpr std::size_t kRejectedPacketLimit = 8;
  41. /** A rejected logical packet can span every native transport fragment. */
  42. constexpr std::size_t kRejectedPacketCapacity = fragments::kMaximumPacketBytes;
  43. /** A 256-byte hex chunk leaves room for event fields in the 1024-byte log line. */
  44. constexpr std::size_t kRejectedHexChunk = 256;
  45. struct RejectedPacket final {
  46. gp::entity_identity::Source source{};
  47. std::array<std::byte, kRejectedPacketCapacity> bytes{};
  48. std::size_t size{};
  49. };
  50. SRWLOCK g_rejectedPacketLock{SRWLOCK_INIT};
  51. std::array<RejectedPacket, kRejectedPacketLimit> g_rejectedPackets{};
  52. std::size_t g_rejectedPacketCount{};
  53. wire::EstablishedPacket g_rejectedPacketHeader{};
  54. /**
  55. * Logs complete payload bytes in numbered chunks without truncating a packet.
  56. * @param capture
  57. * Process-local capture number.
  58. * @param payload Complete decrypted datagram.
  59. */
  60. void log_rejected_hex(std::size_t capture, std::span<const std::byte> payload) noexcept {
  61. std::array<char, core::log::kLineCapacity> line{};
  62. for (std::size_t offset = 0; offset < payload.size(); offset += kRejectedHexChunk) {
  63. const auto count = (std::min)(kRejectedHexChunk, payload.size() - offset);
  64. const int prefix = std::snprintf(
  65. line.data(),
  66. line.size(),
  67. "ev=gameplay stage=rejected_packet_hex capture=%zu offset=%zu bytes=%zu hex=",
  68. capture,
  69. offset,
  70. count);
  71. if (prefix <= 0 || static_cast<std::size_t>(prefix) >= line.size()) {
  72. return;
  73. }
  74. auto length = static_cast<std::size_t>(prefix);
  75. if (!core::log::append_hex(line, length, payload.subspan(offset, count))) {
  76. return;
  77. }
  78. core::log::write(
  79. core::log::Channel::server, core::log::Level::debug, {line.data(), length});
  80. }
  81. }
  82. /**
  83. * Retains a bounded replay sample only while server debug logging is enabled.
  84. * @param source Source already admitted for this decode.
  85. * @param payload Complete decrypted datagram.
  86. * @param externalOffset Start of the external handler in bits.
  87. * @param laneOffset Start of the rejected entity lane in bits.
  88. * @param stoppedOffset Reader position at failure in bits.
  89. */
  90. void log_rejected_entity_packet(const gp::entity_identity::Source& source,
  91. std::span<const std::byte> payload,
  92. std::size_t externalOffset,
  93. std::size_t laneOffset,
  94. std::size_t stoppedOffset) noexcept {
  95. if (!core::log::accepts(core::log::Channel::server, core::log::Level::debug) || payload.empty()
  96. || payload.size() > kRejectedPacketCapacity) {
  97. return;
  98. }
  99. AcquireSRWLockExclusive(&g_rejectedPacketLock);
  100. bool duplicate = false;
  101. for (std::size_t index = 0; index < g_rejectedPacketCount; ++index) {
  102. const auto& prior = g_rejectedPackets[index];
  103. if (prior.source == source && prior.size == payload.size()
  104. && std::equal(payload.begin(), payload.end(), prior.bytes.begin())) {
  105. duplicate = true;
  106. break;
  107. }
  108. }
  109. if (duplicate || g_rejectedPacketCount == kRejectedPacketLimit) {
  110. ReleaseSRWLockExclusive(&g_rejectedPacketLock);
  111. return;
  112. }
  113. auto& saved = g_rejectedPackets[g_rejectedPacketCount];
  114. saved.source = source;
  115. saved.size = payload.size();
  116. std::copy(payload.begin(), payload.end(), saved.bytes.begin());
  117. const auto capture = ++g_rejectedPacketCount;
  118. const bool headerValid = wire::decode_established(payload, true, g_rejectedPacketHeader);
  119. const auto sequence = g_rejectedPacketHeader.ack.outboundHead;
  120. const auto hasSequence = g_rejectedPacketHeader.ack.outboundHeadPresent;
  121. const auto guard = g_rejectedPacketHeader.connectionSequenceLow2;
  122. ReleaseSRWLockExclusive(&g_rejectedPacketLock);
  123. report(core::log::Level::debug,
  124. "ev=gameplay stage=rejected_packet capture=%zu reason=lane2 bytes=%zu "
  125. "external_bit=%zu lane_bit=%zu stopped_bit=%zu header=%u sequence=%u has_sequence=%u "
  126. "guard=%u",
  127. capture,
  128. payload.size(),
  129. externalOffset,
  130. laneOffset,
  131. stoppedOffset,
  132. headerValid ? 1U : 0U,
  133. static_cast<unsigned>(sequence),
  134. hasSequence ? 1U : 0U,
  135. static_cast<unsigned>(guard));
  136. report(core::log::Level::debug,
  137. "ev=gameplay stage=rejected_packet_source capture=%zu activity=0x%llX revision=%llu "
  138. "client_generation=%llu group=0x%llX peer=%llu channel=%llu view=%llu "
  139. "address=0x%08X port=%u local_port=%u local_sequence=%u remote_sequence=%u",
  140. capture,
  141. static_cast<unsigned long long>(source.activitySessionId),
  142. static_cast<unsigned long long>(source.activityRevision),
  143. static_cast<unsigned long long>(source.activityClientGeneration),
  144. static_cast<unsigned long long>(source.groupSessionId),
  145. static_cast<unsigned long long>(source.peerGeneration),
  146. static_cast<unsigned long long>(source.channelGeneration),
  147. static_cast<unsigned long long>(source.viewGeneration),
  148. source.address,
  149. static_cast<unsigned>(source.port),
  150. static_cast<unsigned>(source.localPort),
  151. source.localConnectionSequence,
  152. source.remoteConnectionSequence);
  153. log_rejected_hex(capture, payload);
  154. }
  155. /** Packet ordinals share the receive ring's half-range ordering. */
  156. std::uint64_t packet_ordinal(const gp::PeerLink& peer, std::uint16_t sequence) noexcept {
  157. if (!peer.ringInitialized) {
  158. return gp::kPacketSequenceModulus + sequence;
  159. }
  160. const auto forward =
  161. (sequence + gp::kPacketSequenceModulus - peer.receiveHead) % gp::kPacketSequenceModulus;
  162. return forward < gp::kPacketSequenceHalf
  163. ? peer.receiveOrdinal + forward
  164. : peer.receiveOrdinal - (gp::kPacketSequenceModulus - forward);
  165. }
  166. /**
  167. * Records one received packet sequence in the acknowledgement history.
  168. * @param peer Peer receiving the packet.
  169. * @param sequence Sequence the packet published.
  170. */
  171. void record_sequence(gp::PeerLink& peer, std::uint16_t sequence) noexcept {
  172. if (!peer.ringInitialized) {
  173. peer.receiveOrdinal = packet_ordinal(peer, sequence);
  174. peer.ringInitialized = true;
  175. peer.receiveHead = sequence;
  176. peer.received = {};
  177. return;
  178. }
  179. // Add the modulus before subtracting. A bare difference is signed and goes negative on a wrap.
  180. const std::uint16_t advance = static_cast<std::uint16_t>(
  181. (sequence + gp::kPacketSequenceModulus - peer.receiveHead) % gp::kPacketSequenceModulus);
  182. if (advance == 0 || advance >= gp::kPacketSequenceHalf) {
  183. // A repeat or an older packet leaves the published history alone.
  184. return;
  185. }
  186. std::array<bool, gp::kAckHistory> shifted{};
  187. for (std::size_t index = 0; index < shifted.size(); ++index) {
  188. // Entry `index` is the packet `index + 1` before the new head, so the old head lands at
  189. // `advance - 1`. Anything newer than the old head and older than this packet was skipped.
  190. if (index + 1 < advance) {
  191. continue;
  192. }
  193. if (index + 1 == advance) {
  194. shifted[index] = true;
  195. continue;
  196. }
  197. const std::size_t source = index - advance;
  198. shifted[index] = source < peer.received.size() && peer.received[source];
  199. }
  200. peer.received = shifted;
  201. peer.receiveOrdinal += advance;
  202. peer.receiveHead = sequence;
  203. }
  204. /**
  205. * Applies one reassembled reliable message.
  206. * @param peer Peer that sent it, held under the lock.
  207. * @param message Reassembled message and its inner header.
  208. */
  209. void apply_message(gp::PeerLink& peer, const wire::AssembledMessage& message) noexcept {
  210. if (message.id == static_cast<std::uint8_t>(wire::ConnectId::establish)
  211. && peer.stage == gp::PeerStage::connecting) {
  212. // The reliable establish is what moves a connected peer past the out-of-band pair.
  213. peer.stage = gp::PeerStage::connected;
  214. }
  215. }
  216. /**
  217. * Clears the send queue once the peer acknowledges the packet that carried it.
  218. * @param peer Peer whose acknowledgement arrived, held under the lock.
  219. * @param ack Acknowledgement state the packet published.
  220. * @return True when this acknowledgement emptied the queue.
  221. */
  222. bool apply_acknowledgement(gp::PeerLink& peer, const wire::AckState& ack) noexcept {
  223. if (!peer.outbound.awaitingAcknowledgement
  224. || !wire::acknowledgement_covers(ack, peer.outbound.sentInPacket)) {
  225. return false;
  226. }
  227. // The peer has the packet, so every fragment in it is delivered. The next sequence is kept
  228. // because message sequences continue across messages.
  229. for (gp::OutboundFragment& fragment : peer.outbound.fragments) {
  230. fragment = {};
  231. }
  232. peer.outbound.count = 0;
  233. peer.outbound.awaitingAcknowledgement = false;
  234. return true;
  235. }
  236. /** Common state retained from one complete external frame. */
  237. struct ParsedExternal {
  238. middleware::gameplay::external::CommonState common{};
  239. middleware::gameplay::external::SimulationEventBatch lane0{};
  240. middleware::gameplay::external::ControlStateBatch lane1{};
  241. middleware::gameplay::external::EntityBatch entities{};
  242. bool commonPresent{};
  243. };
  244. /** Exact external component that refused a frame. */
  245. enum class ExternalReadResult : std::uint8_t {
  246. accepted,
  247. prefix,
  248. common,
  249. lane0,
  250. lane1,
  251. lane2,
  252. lane3,
  253. filler,
  254. };
  255. using middleware::gameplay::external::read_flag;
  256. /** Reads or skips the receive-only player lane without retaining its local snapshot. */
  257. [[nodiscard]] bool read_player_lane(bits::Reader& reader) noexcept {
  258. bool present = false;
  259. bool trailingList = false;
  260. return read_flag(reader, present)
  261. && (!present || (reader.skip(kPlayerSnapshotBits) && read_flag(reader, trailingList)));
  262. }
  263. /** Reads common, channels 0 to 3, and filler. */
  264. [[nodiscard]] ExternalReadResult
  265. read_external(std::span<const std::byte> payload,
  266. std::size_t bitOffset,
  267. const gp::entity_identity::Source& source,
  268. const middleware::gameplay::external::Lane0Codec& lane0,
  269. const Lane0Transport& lane0Transport,
  270. const middleware::gameplay::external::TypePayloadCodec& entities,
  271. ParsedExternal& output) noexcept {
  272. output.commonPresent = false;
  273. bits::Reader reader(payload);
  274. const std::unique_ptr<ParsedExternal> candidateStorage(new (std::nothrow) ParsedExternal{});
  275. if (!candidateStorage) {
  276. return ExternalReadResult::lane2;
  277. }
  278. ParsedExternal& candidate = *candidateStorage;
  279. bool lanePresent = false;
  280. bool externalPresent = false;
  281. if (!reader.skip(bitOffset) || !read_flag(reader, externalPresent) || !externalPresent) {
  282. return ExternalReadResult::prefix;
  283. }
  284. if (!read_flag(reader, candidate.commonPresent)
  285. || (candidate.commonPresent
  286. && !middleware::gameplay::external::read_common_state(reader, candidate.common))) {
  287. return ExternalReadResult::common;
  288. }
  289. output.commonPresent = candidate.commonPresent;
  290. output.common = candidate.common;
  291. if (lane0Transport.write != nullptr) {
  292. if (!middleware::gameplay::external::read_simulation_event_lane(
  293. reader, lane0Transport.payloadCodec, candidate.lane0)) {
  294. return ExternalReadResult::lane0;
  295. }
  296. } else if (lane0.read != nullptr) {
  297. if (!lane0.read(lane0.context, reader)) {
  298. return ExternalReadResult::lane0;
  299. }
  300. } else if (!read_flag(reader, lanePresent) || lanePresent) {
  301. return ExternalReadResult::lane0;
  302. }
  303. if (!middleware::gameplay::external::read_control_state_lane(
  304. reader,
  305. middleware::gameplay::external::control_state_payload_codec(),
  306. candidate.lane1)) {
  307. return ExternalReadResult::lane1;
  308. }
  309. const auto lane2Offset = payload.size() * 8U - reader.remaining_bits();
  310. if (g_entityTransport.read != nullptr
  311. ? !g_entityTransport.read(g_entityTransport.context, source, reader, candidate.entities)
  312. : !middleware::gameplay::external::read_entity_batch(
  313. reader, entities, candidate.entities)) {
  314. log_rejected_entity_packet(
  315. source, payload, bitOffset, lane2Offset, payload.size() * 8U - reader.remaining_bits());
  316. return ExternalReadResult::lane2;
  317. }
  318. if (!read_player_lane(reader)) {
  319. return ExternalReadResult::lane3;
  320. }
  321. wire::FillerTrailer filler{};
  322. if (!wire::read_filler_and_padding(reader, filler) || reader.remaining_bits() != 0) {
  323. return ExternalReadResult::filler;
  324. }
  325. output = candidate;
  326. return ExternalReadResult::accepted;
  327. }
  328. /** @return Stable log name for one external read result. */
  329. [[nodiscard]] const char* external_result_name(ExternalReadResult result) noexcept {
  330. // One name per ExternalReadResult, in declaration order.
  331. constexpr std::array<const char*, 8> names = {
  332. "accepted", "prefix", "common", "lane0", "lane1", "lane2", "lane3", "filler"};
  333. const auto index = static_cast<std::size_t>(result);
  334. return index < names.size() ? names[index] : "unknown";
  335. }
  336. /** Queues message 44 after the peer transaction accepts the initial common root. */
  337. void queue_common_request(const state::gameplay::Endpoint& from,
  338. std::uint64_t groupSessionId,
  339. const state::activity::SessionBinding& binding,
  340. std::uint64_t ownerGeneration,
  341. std::uint8_t requested) noexcept {
  342. if (!sunrise::server::bap::request_replication_epoch(binding, ownerGeneration, requested)) {
  343. return;
  344. }
  345. AcquireSRWLockExclusive(&g_lock);
  346. gp::PeerLink* const peer = find_locked(from);
  347. if (peer != nullptr && peer->externalGroupSessionId == groupSessionId
  348. && peer->commonReconciler.owner_generation() == ownerGeneration) {
  349. static_cast<void>(peer->commonReconciler.commit_request());
  350. }
  351. ReleaseSRWLockExclusive(&g_lock);
  352. }
  353. /** Writes the external handler body after both reliable queues; true when all of it fit. */
  354. [[nodiscard]] bool write_external(const gp::PeerLink& peer, bits::Writer& writer) noexcept {
  355. middleware::gameplay::external::CommonState common{};
  356. const auto viewPhase = peer.viewReceptor.phase();
  357. if (peer.stage != gp::PeerStage::connected
  358. || (viewPhase != gp::external::view_receptor::Phase::provisional
  359. && viewPhase != gp::external::view_receptor::Phase::accepted)
  360. || peer.externalGroupSessionId == 0 || !peer.commonReconciler.outbound_common(common)
  361. || (g_lane0Transport.write == nullptr && g_lane0Codec.write == nullptr)) {
  362. return false;
  363. }
  364. const bool commonPresent = !peer.commonCommitted;
  365. if (!writer.write(commonPresent ? 1U : 0U, 1)
  366. || (commonPresent && !middleware::gameplay::external::write_common_state(writer, common))
  367. || (g_lane0Transport.write != nullptr
  368. ? !g_lane0Transport.write(g_lane0Transport.context,
  369. peer.externalGroupSessionId,
  370. peer.nextExternalTransmission,
  371. writer)
  372. : !g_lane0Codec.write(g_lane0Codec.context, writer))
  373. || !writer.write(0, 1) || !writer.write(0, 1) || !writer.write(0, 1) || !writer.write(1, 1)
  374. || !writer.write(0, 1)) {
  375. return false;
  376. }
  377. return true;
  378. }
  379. /** Builds and sends one ACK packet from a peer copy taken under the lock. */
  380. [[nodiscard]] bool send_acknowledgement(const gp::PeerLink& peer) noexcept {
  381. wire::AckState ack{};
  382. ack.outboundHead = peer.outboundHead;
  383. ack.outboundHeadPresent = peer.outboundHeadPresent;
  384. // The peer subtracts this from the decoded sequence to place its receive window. A zero
  385. // collapses that window and the peer discards every packet.
  386. ack.headMinusCursor = kMinimumHeadCursor;
  387. ack.receiveHead = peer.receiveHead;
  388. ack.ringInitialized = peer.ringInitialized;
  389. ack.received = peer.received;
  390. // No round trip is timed, so the delay field carries its sentinel.
  391. ack.delay = kDelaySentinel;
  392. std::array<std::byte, kReplyCapacity> buffer{};
  393. bits::Writer writer(buffer);
  394. const std::uint8_t guard = wire::connection_sequence_low2(peer.localConnectionSequence);
  395. std::size_t size = 0;
  396. // Only the 32-byte queue carries this host's messages; the 6-byte queue stays empty.
  397. middleware::gameplay::external::CommonState externalCommon{};
  398. const auto viewPhase = peer.viewReceptor.phase();
  399. const bool external = (viewPhase == gp::external::view_receptor::Phase::provisional
  400. || viewPhase == gp::external::view_receptor::Phase::accepted)
  401. && peer.externalGroupSessionId != 0
  402. && peer.commonReconciler.outbound_common(externalCommon);
  403. if (!wire::write_head_and_ack(writer, guard, ack) || !wire::write_queue(writer, peer.outbound)
  404. || !wire::write_empty_queue(writer)) {
  405. return false;
  406. }
  407. if (external) {
  408. if (!writer.write(1, 1) || !write_external(peer, writer) || !writer.write(0, 1)) {
  409. return false;
  410. }
  411. } else if (!wire::write_absent_filler(writer)) {
  412. return false;
  413. }
  414. if (!writer.finish(size)) {
  415. return false;
  416. }
  417. return send_transport(peer.endpoint, {buffer.data(), size});
  418. }
  419. } // namespace
  420. /** Consumes one established packet. */
  421. void consume_established(const gp::Endpoint& from,
  422. std::span<const std::byte> payload,
  423. std::uint64_t now) noexcept {
  424. PacketInput input{};
  425. if (!prepare_packet_input(from, payload, now, input)) {
  426. return;
  427. }
  428. payload = input.payload;
  429. wire::EstablishedPacket packet{};
  430. if (!wire::decode_established(payload, true, packet)) {
  431. report(core::log::Level::debug, "ev=gameplay stage=packet result=drop reason=grammar");
  432. return;
  433. }
  434. std::array<std::uint8_t, kMessageReportCapacity> delivered{};
  435. std::size_t deliveredCount = 0;
  436. unsigned stage = 0;
  437. bool queueCleared = false;
  438. std::uint16_t clearedPacket = 0;
  439. std::uint64_t externalGroupSessionId = 0;
  440. // The reliable window never resynchronises, so a stalled queue is only visible as a refused
  441. // record against the sequence it is still waiting for.
  442. std::size_t largeDropped = 0;
  443. std::uint16_t largeNext = 0;
  444. std::uint16_t largeFirst = 0;
  445. bool peerFound = false;
  446. bool guardAccepted = false;
  447. bool externalExpected = false;
  448. bool externalValid = true;
  449. const char* externalFailure = "none";
  450. const std::unique_ptr<ParsedExternal> externalStorage(new (std::nothrow) ParsedExternal{});
  451. if (!externalStorage) {
  452. return;
  453. }
  454. ParsedExternal& external = *externalStorage;
  455. state::activity::SessionBinding commonBinding{};
  456. std::uint64_t commonOwnerGeneration = 0;
  457. std::uint8_t commonRequestedGeneration = 0;
  458. bool commonRequest = false;
  459. gp::external::common_reconciler::Reconciler commonCandidate{};
  460. bool commonCandidatePresent = false;
  461. DisplacedExternals completed{};
  462. std::size_t completedCount = 0;
  463. std::uint8_t expectedGuard = 0;
  464. std::array<wire::AssembledMessage, kMessageReportCapacity> bodies{};
  465. std::array<bool, kMessageReportCapacity> deferredView{};
  466. gp::entity_identity::Source ingress{};
  467. AcquireSRWLockExclusive(&g_lock);
  468. gp::PeerLink* peer = find_locked(from);
  469. if (peer != nullptr) {
  470. ingress = entity_source(*peer);
  471. expectedGuard = wire::connection_sequence_low2(peer->remoteConnectionSequence);
  472. guardAccepted =
  473. packet.connectionSequenceLow2 == expectedGuard && packet_channel_matches(*peer, input);
  474. }
  475. if (guardAccepted) {
  476. if (packet.ack.outboundHeadPresent) {
  477. record_sequence(*peer, packet.ack.outboundHead);
  478. }
  479. clearedPacket = peer->outbound.sentInPacket;
  480. queueCleared = apply_acknowledgement(*peer, packet.ack);
  481. peer->acknowledgementOwed = true;
  482. peer->lastTick = now;
  483. largeDropped = wire::accept_records(packet.large, peer->large);
  484. largeNext = peer->large.nextSequence;
  485. largeFirst = packet.large.count == 0 ? 0 : packet.large.records[0].sequence;
  486. wire::accept_records(packet.small, peer->small);
  487. wire::AssembledMessage message{};
  488. while (wire::drain_message(peer->large, message)) {
  489. apply_message(*peer, message);
  490. if (deliveredCount < delivered.size()) {
  491. delivered[deliveredCount] = message.id;
  492. bodies[deliveredCount] = message;
  493. ++deliveredCount;
  494. }
  495. }
  496. while (wire::drain_message(peer->small, message)) {
  497. apply_message(*peer, message);
  498. if (deliveredCount < delivered.size()) {
  499. delivered[deliveredCount] = message.id;
  500. bodies[deliveredCount] = message;
  501. ++deliveredCount;
  502. }
  503. }
  504. stage = static_cast<unsigned>(peer->stage);
  505. }
  506. ReleaseSRWLockExclusive(&g_lock);
  507. // Reliable controls establish the view used by this packet's external payload.
  508. for (std::size_t index = 0; index < deliveredCount; ++index) {
  509. report(core::log::Level::info,
  510. "ev=gameplay stage=message result=ok id=%u peerstage=%u",
  511. static_cast<unsigned>(delivered[index]),
  512. stage);
  513. // The connect establish belongs to this layer and apply_message already took it, so
  514. // handing it to the group layer would only report it as undecoded on every connection.
  515. const wire::AssembledMessage& body = bodies[index];
  516. if (body.id == static_cast<std::uint8_t>(wire::ConnectId::establish)) {
  517. continue;
  518. }
  519. // Group handling runs outside the lock because answering takes it again.
  520. bits::Reader reader({body.bytes.data(), gp::kReassemblyCapacity});
  521. namespace viewWire = middleware::gameplay::group;
  522. if (body.id == viewWire::kViewMessageId) {
  523. auto stageReader = reader;
  524. viewWire::ViewEstablishment transition{};
  525. if (stageReader.skip(body.bodyBitOffset) && viewWire::read_view(stageReader, transition)
  526. && transition.kind == 5) {
  527. // Stage 5 may depend on common state carried later in this same packet.
  528. deferredView[index] = true;
  529. continue;
  530. }
  531. }
  532. if (reader.skip(body.bodyBitOffset) && !group::consume(from, body.id, reader, now)) {
  533. report(core::log::Level::debug,
  534. "ev=gameplay stage=message result=undecoded id=%u",
  535. static_cast<unsigned>(body.id));
  536. }
  537. }
  538. AcquireSRWLockExclusive(&g_lock);
  539. peer = find_locked(from);
  540. const bool sameChannel = peer != nullptr && peer->peerGeneration == ingress.peerGeneration
  541. && peer->channelGeneration == ingress.channelGeneration
  542. && peer->localConnectionSequence == ingress.localConnectionSequence
  543. && peer->remoteConnectionSequence == ingress.remoteConnectionSequence;
  544. guardAccepted = guardAccepted && sameChannel;
  545. if (peer != nullptr) {
  546. peerFound = true;
  547. expectedGuard = wire::connection_sequence_low2(peer->remoteConnectionSequence);
  548. guardAccepted =
  549. guardAccepted && sameChannel && packet.connectionSequenceLow2 == expectedGuard;
  550. externalExpected = peer->viewReceptor.accepts_inbound_entities();
  551. if (guardAccepted && externalExpected) {
  552. externalGroupSessionId = peer->externalGroupSessionId;
  553. const auto source = entity_source(*peer);
  554. const auto ordinal = packet_ordinal(*peer, packet.ack.outboundHead);
  555. const ExternalReadResult externalRead = read_external(payload,
  556. packet.externalBitOffset,
  557. source,
  558. g_lane0Codec,
  559. g_lane0Transport,
  560. g_entityCodec,
  561. external);
  562. externalValid = externalRead == ExternalReadResult::accepted;
  563. externalFailure = external_result_name(externalRead);
  564. if (external.commonPresent) {
  565. commonCandidate = peer->commonReconciler;
  566. const auto result = commonCandidate.observe(external.common);
  567. const bool commonValid =
  568. result == gp::external::common_reconciler::ObserveResult::initialAccepted
  569. || result
  570. == gp::external::common_reconciler::ObserveResult::
  571. awaitingRequestedGeneration
  572. || result == gp::external::common_reconciler::ObserveResult::ready;
  573. if (!commonValid) {
  574. externalValid = false;
  575. externalFailure = "reconcile";
  576. }
  577. commonRequest =
  578. result == gp::external::common_reconciler::ObserveResult::initialAccepted
  579. && commonCandidate.pending_request(commonRequestedGeneration);
  580. if (commonRequest) {
  581. commonBinding = peer->activityBinding;
  582. commonOwnerGeneration = commonCandidate.owner_generation();
  583. }
  584. commonCandidatePresent = commonValid;
  585. if (commonValid) {
  586. static_cast<void>(commonCandidate.qualify_entities(
  587. &external.common, ordinal, packet.ack.outboundHeadPresent));
  588. peer->commonReconciler = commonCandidate;
  589. }
  590. }
  591. if (externalValid) {
  592. if (!commonCandidatePresent) {
  593. commonCandidate = peer->commonReconciler;
  594. }
  595. const bool currentEntityEpoch = commonCandidate.qualify_entities(
  596. external.commonPresent ? &external.common : nullptr,
  597. ordinal,
  598. packet.ack.outboundHeadPresent);
  599. if (external.entities.recordPresent && !currentEntityEpoch) {
  600. externalValid = false;
  601. externalFailure = "entity_epoch";
  602. } else {
  603. external.entities.hasAllocationEpoch = currentEntityEpoch;
  604. external.entities.allocationEpoch = commonCandidate.requested_generation();
  605. external.entities.allocationDomain = commonCandidate.allocation_domain();
  606. commonCandidatePresent = true;
  607. }
  608. }
  609. middleware::gameplay::external::EntityBaselineMutation entityMutation{};
  610. if (externalValid && external.entities.recordPresent
  611. && g_entityTransport.prepare != nullptr) {
  612. externalValid = g_entityTransport.commit != nullptr
  613. && g_entityTransport.prepare(g_entityTransport.context,
  614. source,
  615. external.entities,
  616. packet.ack.outboundHead,
  617. packet.ack.outboundHeadPresent,
  618. ordinal,
  619. entityMutation);
  620. if (!externalValid) {
  621. externalFailure = "lane2_prepare";
  622. }
  623. }
  624. if (externalValid && externalGroupSessionId != 0
  625. && g_lane0Transport.accepted != nullptr) {
  626. externalValid = g_lane0Transport.accepted(
  627. g_lane0Transport.context, externalGroupSessionId, external.lane0);
  628. if (!externalValid) {
  629. externalFailure = "lane0_accept";
  630. }
  631. }
  632. // The peer lock keeps a prepared baseline stable until commit.
  633. if (externalValid && externalGroupSessionId != 0 && external.entities.recordPresent) {
  634. externalValid = g_entityTransport.commit != nullptr
  635. ? g_entityTransport.commit(
  636. g_entityTransport.context, source, entityMutation)
  637. : g_entityAccepted == nullptr
  638. || g_entityAccepted(g_entityAcceptedContext,
  639. externalGroupSessionId,
  640. external.entities);
  641. if (!externalValid) {
  642. externalFailure = "lane2_accept";
  643. }
  644. }
  645. if (externalValid && commonCandidatePresent) {
  646. peer->commonReconciler = commonCandidate;
  647. }
  648. if (externalValid && external.entities.recordPresent
  649. && g_entityTransport.observed != nullptr) {
  650. external.entities.ignoredRecordMask = entityMutation.ignoredRecordMask;
  651. g_entityTransport.observed(g_entityTransport.context,
  652. source,
  653. external.entities,
  654. packet.ack.outboundHead,
  655. packet.ack.outboundHeadPresent,
  656. ordinal,
  657. now);
  658. }
  659. }
  660. }
  661. // Authenticated transport receipts remain valid when an inbound application lane is refused.
  662. if (guardAccepted) {
  663. for (auto& contribution : peer->externalContributions) {
  664. if (!contribution.occupied) {
  665. continue;
  666. }
  667. const wire::AckOutcome outcome =
  668. wire::acknowledgement_outcome(packet.ack, contribution.packetSequence);
  669. if (outcome == wire::AckOutcome::unresolved) {
  670. continue;
  671. }
  672. completed[completedCount++] = {
  673. contribution.groupSessionId, contribution.transmissionId, outcome};
  674. if (outcome == wire::AckOutcome::received && contribution.commonPresent
  675. && contribution.viewGeneration == peer->viewGeneration) {
  676. peer->commonCommitted = true;
  677. }
  678. contribution = {};
  679. }
  680. }
  681. ReleaseSRWLockExclusive(&g_lock);
  682. if (!peerFound) {
  683. return;
  684. }
  685. if (!guardAccepted) {
  686. report(core::log::Level::debug,
  687. "ev=gameplay stage=packet result=drop reason=%s got=%u expect=%u",
  688. sameChannel ? "channel_low2" : "channel_changed",
  689. static_cast<unsigned>(packet.connectionSequenceLow2),
  690. static_cast<unsigned>(expectedGuard));
  691. return;
  692. }
  693. if (!externalValid) {
  694. report(core::log::Level::debug,
  695. "ev=gameplay stage=external result=drop reason=%s group=0x%016llX",
  696. externalFailure,
  697. static_cast<unsigned long long>(externalGroupSessionId));
  698. }
  699. if (commonRequest && externalGroupSessionId != 0) {
  700. queue_common_request(from,
  701. externalGroupSessionId,
  702. commonBinding,
  703. commonOwnerGeneration,
  704. commonRequestedGeneration);
  705. }
  706. for (std::size_t index = 0; index < deliveredCount; ++index) {
  707. if (!deferredView[index]) {
  708. continue;
  709. }
  710. const auto& body = bodies[index];
  711. bits::Reader reader({body.bytes.data(), gp::kReassemblyCapacity});
  712. if (reader.skip(body.bodyBitOffset)) {
  713. static_cast<void>(group::consume(from, body.id, reader, now));
  714. }
  715. }
  716. notify_external_outcomes(completed, completedCount);
  717. if (queueCleared) {
  718. report(core::log::Level::info,
  719. "ev=gameplay stage=sendqueue result=cleared packet=%u base=%u entries=%u",
  720. static_cast<unsigned>(clearedPacket),
  721. static_cast<unsigned>(packet.ack.receiveHead),
  722. static_cast<unsigned>(packet.ack.reportedCount));
  723. }
  724. report(core::log::Level::debug,
  725. "ev=gameplay stage=packet result=ok seq=%u base=%u entries=%u large=%u small=%u "
  726. "first=%u next=%u drop=%zu",
  727. static_cast<unsigned>(packet.ack.outboundHead),
  728. static_cast<unsigned>(packet.ack.receiveHead),
  729. static_cast<unsigned>(packet.ack.reportedCount),
  730. static_cast<unsigned>(packet.large.count),
  731. static_cast<unsigned>(packet.small.count),
  732. static_cast<unsigned>(largeFirst),
  733. static_cast<unsigned>(largeNext),
  734. largeDropped);
  735. }
  736. /** Sends any owed acknowledgement. */
  737. void service(std::uint64_t now) noexcept {
  738. std::array<gp::PeerLink, gp::kAssociationCapacity> owed{};
  739. std::size_t count = 0;
  740. AcquireSRWLockExclusive(&g_lock);
  741. for (gp::PeerLink& peer : g_peers) {
  742. // An unacknowledged send queue keeps the packet going out until the peer confirms it.
  743. // Every packet burns one sequence, so the resend is paced.
  744. const bool resendDue = peer.outbound.count != 0 && now - peer.lastSend >= kResendInterval;
  745. const bool due = peer.acknowledgementOwed || resendDue;
  746. if (peer.stage == gp::PeerStage::absent || !due) {
  747. continue;
  748. }
  749. middleware::gameplay::external::CommonState externalCommon{};
  750. const auto viewPhase = peer.viewReceptor.phase();
  751. const bool external = (viewPhase == gp::external::view_receptor::Phase::provisional
  752. || viewPhase == gp::external::view_receptor::Phase::accepted)
  753. && peer.externalGroupSessionId != 0
  754. && peer.commonReconciler.outbound_common(externalCommon);
  755. const std::uint16_t nextPacket =
  756. static_cast<std::uint16_t>((peer.outboundHead + 1) % gp::kPacketSequenceModulus);
  757. if (external
  758. && peer.externalContributions[nextPacket % peer.externalContributions.size()]
  759. .occupied) {
  760. continue;
  761. }
  762. peer.acknowledgementOwed = false;
  763. peer.lastSend = now;
  764. // Only the first send of the current contents is stamped. A resend carries the same
  765. // fragments, so re-stamping would move the target past what the peer can acknowledge.
  766. if (peer.outbound.count != 0 && !peer.outbound.awaitingAcknowledgement) {
  767. peer.outbound.sentInPacket =
  768. static_cast<std::uint16_t>((peer.outboundHead + 1) % gp::kPacketSequenceModulus);
  769. peer.outbound.awaitingAcknowledgement = true;
  770. }
  771. // The packet sequence advances here so the copy carries the value it will publish.
  772. peer.outboundHead = nextPacket;
  773. peer.outboundHeadPresent = true;
  774. if (external) {
  775. ++peer.nextExternalTransmission;
  776. if (peer.nextExternalTransmission == 0) {
  777. ++peer.nextExternalTransmission;
  778. }
  779. auto& contribution =
  780. peer.externalContributions[nextPacket % peer.externalContributions.size()];
  781. contribution.transmissionId = peer.nextExternalTransmission;
  782. contribution.groupSessionId = peer.externalGroupSessionId;
  783. contribution.viewGeneration = peer.viewGeneration;
  784. contribution.packetSequence = nextPacket;
  785. contribution.commonPresent = !peer.commonCommitted;
  786. contribution.lane0Present = true;
  787. contribution.occupied = true;
  788. }
  789. peer.lastTick = now;
  790. owed[count] = peer;
  791. ++count;
  792. }
  793. ReleaseSRWLockExclusive(&g_lock);
  794. for (std::size_t index = 0; index < count; ++index) {
  795. if (send_acknowledgement(owed[index])) {
  796. continue;
  797. }
  798. report(core::log::Level::debug, "ev=gameplay stage=ack result=fail");
  799. DisplacedExternals displaced{};
  800. std::size_t displacedCount = 0;
  801. AcquireSRWLockExclusive(&g_lock);
  802. gp::PeerLink* const peer = find_locked(owed[index].endpoint);
  803. if (peer != nullptr && peer->peerGeneration == owed[index].peerGeneration) {
  804. peer->acknowledgementOwed = true;
  805. if (peer->outbound.awaitingAcknowledgement
  806. && peer->outbound.sentInPacket == owed[index].outboundHead) {
  807. peer->outbound.awaitingAcknowledgement = false;
  808. }
  809. auto& reserved = peer->externalContributions[owed[index].outboundHead
  810. % peer->externalContributions.size()];
  811. // The packet never left, so the stake is displaced like any other lost contribution.
  812. if (reserved.occupied
  813. && reserved.transmissionId == owed[index].nextExternalTransmission) {
  814. displaced[displacedCount++] = {
  815. reserved.groupSessionId, reserved.transmissionId, wire::AckOutcome::unresolved};
  816. reserved = {};
  817. }
  818. }
  819. ReleaseSRWLockExclusive(&g_lock);
  820. notify_external_outcomes(displaced, displacedCount);
  821. }
  822. }
  823. } // namespace sunrise::server::gameplay::peer