peer_established.cpp 38 KB

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