Fixes to sending entity updates and entity rpcs within an environment set up for cross host entity migration

Signed-off-by: kberg-amzn <karlberg@amazon.com>
This commit is contained in:
kberg-amzn
2021-09-28 19:25:04 -07:00
parent 0a829f9661
commit 02bc89cd92
21 changed files with 185 additions and 134 deletions
@@ -10,7 +10,6 @@
#include <Multiplayer/NetworkEntity/EntityReplication/EntityReplicator.h>
#include <Source/NetworkEntity/EntityReplication/PropertyPublisher.h>
#include <Source/NetworkEntity/EntityReplication/PropertySubscriber.h>
#include <Source/AutoGen/Multiplayer.AutoPackets.h>
#include <Multiplayer/IMultiplayer.h>
#include <Multiplayer/Components/NetBindComponent.h>
#include <Multiplayer/EntityDomains/IEntityDomain.h>
@@ -107,10 +106,34 @@ namespace Multiplayer
}
}
void EntityReplicationManager::SendUpdates(AZ::TimeMs hostTimeMs)
void EntityReplicationManager::SendUpdates()
{
m_frameTimeMs = AZ::GetElapsedTimeMs();
SendEntityUpdates(hostTimeMs);
{
EntityReplicatorList toSendList = GenerateEntityUpdateList();
AZLOG
(
NET_ReplicationInfo,
"Sending %zd updates from %s to %s",
toSendList.size(),
GetNetworkEntityManager()->GetHostId().GetString().c_str(),
GetRemoteHostId().GetString().c_str()
);
// Prep a replication record for send, at this point, everything needs to be sent
for (EntityReplicator* replicator : toSendList)
{
replicator->GetPropertyPublisher()->PrepareSerialization();
}
// While our to send list is not empty, build up another packet to send
do
{
SendEntityUpdateMessages(toSendList);
} while (!toSendList.empty());
}
SendEntityRpcs(m_deferredRpcMessagesReliable, true);
SendEntityRpcs(m_deferredRpcMessagesUnreliable, false);
@@ -130,65 +153,6 @@ namespace Multiplayer
);
}
void EntityReplicationManager::SendEntityUpdatesPacketHelper
(
AZ::TimeMs hostTimeMs,
EntityReplicatorList& toSendList,
uint32_t maxPayloadSize,
AzNetworking::IConnection& connection
)
{
uint32_t pendingPacketSize = 0;
EntityReplicatorList replicatorUpdatedList;
MultiplayerPackets::EntityUpdates entityUpdatePacket;
entityUpdatePacket.SetHostTimeMs(hostTimeMs);
entityUpdatePacket.SetHostFrameId(GetNetworkTime()->GetHostFrameId());
// Serialize everything
while (!toSendList.empty())
{
EntityReplicator* replicator = toSendList.front();
NetworkEntityUpdateMessage updateMessage(replicator->GenerateUpdatePacket());
const uint32_t nextMessageSize = updateMessage.GetEstimatedSerializeSize();
// Check if we are over our limits
const bool payloadFull = (pendingPacketSize + nextMessageSize > maxPayloadSize);
const bool capacityReached = (entityUpdatePacket.GetEntityMessages().size() >= entityUpdatePacket.GetEntityMessages().capacity());
const bool largeEntityDetected = (payloadFull && replicatorUpdatedList.empty());
if (capacityReached || (payloadFull && !largeEntityDetected))
{
break;
}
pendingPacketSize += nextMessageSize;
entityUpdatePacket.ModifyEntityMessages().push_back(updateMessage);
replicatorUpdatedList.push_back(replicator);
toSendList.pop_front();
if (largeEntityDetected)
{
AZLOG_WARN("\n\n*******************************");
AZLOG_WARN
(
"Serializing extremely large entity (%u) - MaxPayload: %d NeededSize %d",
aznumeric_cast<uint32_t>(replicator->GetEntityHandle().GetNetEntityId()),
maxPayloadSize,
nextMessageSize
);
AZLOG_WARN("*******************************");
break;
}
}
const AzNetworking::PacketId sentId = connection.SendUnreliablePacket(entityUpdatePacket);
// Update the sent things with the packet id
for (EntityReplicator* replicator : replicatorUpdatedList)
{
replicator->GetPropertyPublisher()->FinalizeSerialization(sentId);
}
}
EntityReplicationManager::EntityReplicatorList EntityReplicationManager::GenerateEntityUpdateList()
{
if (m_replicationWindow == nullptr)
@@ -260,76 +224,92 @@ namespace Multiplayer
return toSendList;
}
void EntityReplicationManager::SendEntityUpdates(AZ::TimeMs hostTimeMs)
void EntityReplicationManager::SendEntityUpdateMessages(EntityReplicatorList& replicatorList)
{
EntityReplicatorList toSendList = GenerateEntityUpdateList();
AZLOG
(
NET_ReplicationInfo,
"Sending %zd updates from %s to %s",
toSendList.size(),
GetNetworkEntityManager()->GetHostId().GetString().c_str(),
GetRemoteHostId().GetString().c_str()
);
// prep a replication record for send, at this point, everything needs to be sent
for (EntityReplicator* replicator : toSendList)
uint32_t pendingPacketSize = 0;
EntityReplicatorList replicatorUpdatedList;
NetworkEntityUpdateVector entityUpdates;
// Serialize everything
while (!replicatorList.empty())
{
replicator->GetPropertyPublisher()->PrepareSerialization();
EntityReplicator* replicator = replicatorList.front();
NetworkEntityUpdateMessage updateMessage(replicator->GenerateUpdatePacket());
const uint32_t nextMessageSize = updateMessage.GetEstimatedSerializeSize();
// Check if we are over our limits
const bool payloadFull = (pendingPacketSize + nextMessageSize > m_maxPayloadSize);
const bool capacityReached = (entityUpdates.size() >= entityUpdates.capacity());
const bool largeEntityDetected = (payloadFull && replicatorUpdatedList.empty());
if (capacityReached || (payloadFull && !largeEntityDetected))
{
break;
}
pendingPacketSize += nextMessageSize;
entityUpdates.push_back(updateMessage);
replicatorUpdatedList.push_back(replicator);
replicatorList.pop_front();
if (largeEntityDetected)
{
AZLOG_WARN
(
"Serializing extremely large entity (%u) - MaxPayload: %d NeededSize %d",
aznumeric_cast<uint32_t>(replicator->GetEntityHandle().GetNetEntityId()),
m_maxPayloadSize,
nextMessageSize
);
break;
}
}
// While our to send list is not empty, build up another packet to send
do
const AzNetworking::PacketId sentId = m_replicationWindow->SendEntityUpdateMessages(entityUpdates);
// Update the sent things with the packet id
for (EntityReplicator* replicator : replicatorUpdatedList)
{
SendEntityUpdatesPacketHelper(hostTimeMs, toSendList, m_maxPayloadSize, m_connection);
} while (!toSendList.empty());
replicator->FinalizeSerialization(sentId);
}
}
void EntityReplicationManager::SendEntityRpcs(RpcMessages& deferredRpcs, bool reliable)
void EntityReplicationManager::SendEntityRpcs(RpcMessages& rpcMessages, bool reliable)
{
while (!deferredRpcs.empty())
while (!rpcMessages.empty())
{
MultiplayerPackets::EntityRpcs entityRpcsPacket;
NetworkEntityRpcVector entityRpcs;
uint32_t pendingPacketSize = 0;
while (!deferredRpcs.empty())
while (!rpcMessages.empty())
{
NetworkEntityRpcMessage& message = deferredRpcs.front();
NetworkEntityRpcMessage& message = rpcMessages.front();
const uint32_t nextRpcSize = message.GetEstimatedSerializeSize();
if ((pendingPacketSize + nextRpcSize) > m_maxPayloadSize)
{
// We're over our limit, break and send an Rpc packet
if (entityRpcsPacket.GetEntityRpcs().size() == 0)
if (entityRpcs.size() == 0)
{
AZLOG(NET_Replicator, "Encountered an RPC that is above our MTU, message will be segmented (object size %u, max allowed size %u)", nextRpcSize, m_maxPayloadSize);
entityRpcsPacket.ModifyEntityRpcs().push_back(message);
deferredRpcs.pop_front();
entityRpcs.push_back(message);
rpcMessages.pop_front();
}
break;
}
pendingPacketSize += nextRpcSize;
if (entityRpcsPacket.GetEntityRpcs().full())
if (entityRpcs.full())
{
// Packet was full, send what we've accumulated so far
AZLOG(NET_Replicator, "We've hit our RPC message limit (RPC count %u, packet size %u)", aznumeric_cast<uint32_t>(entityRpcsPacket.GetEntityRpcs().size()), pendingPacketSize);
AZLOG(NET_Replicator, "We've hit our RPC message limit (RPC count %u, packet size %u)", aznumeric_cast<uint32_t>(entityRpcs.size()), pendingPacketSize);
break;
}
entityRpcsPacket.ModifyEntityRpcs().push_back(message);
deferredRpcs.pop_front();
entityRpcs.push_back(message);
rpcMessages.pop_front();
}
if (reliable)
{
m_connection.SendReliablePacket(entityRpcsPacket);
}
else
{
m_connection.SendUnreliablePacket(entityRpcsPacket);
}
m_replicationWindow->SendEntityRpcs(entityRpcs, reliable);
}
}
@@ -474,7 +454,7 @@ namespace Multiplayer
}
// @nt: TODO - delete once dropped RPC problem fixed
void EntityReplicationManager::AddAutonomousEntityReplicatorCreatedHandle(AZ::Event<NetEntityId>::Handler& handler)
void EntityReplicationManager::AddAutonomousEntityReplicatorCreatedHandler(AZ::Event<NetEntityId>::Handler& handler)
{
handler.Connect(m_autonomousEntityReplicatorCreated);
}
@@ -79,7 +79,7 @@ namespace Multiplayer
void AddDeferredRpcMessage(NetworkEntityRpcMessage& rpcMessage);
void AddAutonomousEntityReplicatorCreatedHandle(AZ::Event<NetEntityId>::Handler& handler);
void AddAutonomousEntityReplicatorCreatedHandler(AZ::Event<NetEntityId>::Handler& handler);
bool HandleEntityMigration(AzNetworking::IConnection* invokingConnection, EntityMigrationMessage& message);
bool HandleEntityDeleteMessage(EntityReplicator* entityReplicator, const AzNetworking::IPacketHeader& packetHeader, const NetworkEntityUpdateMessage& updateMessage);
@@ -117,10 +117,8 @@ namespace Multiplayer
using EntityReplicatorList = AZStd::deque<EntityReplicator*>;
EntityReplicatorList GenerateEntityUpdateList();
void SendEntityUpdatesPacketHelper(AZ::TimeMs hostTimeMs, EntityReplicatorList& toSendList, uint32_t maxPayloadSize, AzNetworking::IConnection& connection);
void SendEntityUpdates(AZ::TimeMs hostTimeMs);
void SendEntityRpcs(RpcMessages& deferredRpcs, bool reliable);
void SendEntityUpdateMessages(EntityReplicatorList& replicatorList);
void SendEntityRpcs(RpcMessages& rpcMessages, bool reliable);
void MigrateEntityInternal(NetEntityId entityId);
void OnEntityExitDomain(const ConstNetworkEntityHandle& entityHandle);
@@ -14,7 +14,6 @@
#include <Multiplayer/NetworkEntity/NetworkEntityRpcMessage.h>
#include <Multiplayer/NetworkEntity/EntityReplication/EntityReplicator.h>
#include <Multiplayer/NetworkEntity/EntityReplication/EntityReplicationManager.h>
#include <Source/AutoGen/Multiplayer.AutoPackets.h>
#include <Source/NetworkEntity/NetworkEntityAuthorityTracker.h>
#include <Source/NetworkEntity/NetworkEntityTracker.h>
#include <Source/NetworkEntity/EntityReplication/PropertyPublisher.h>
@@ -495,6 +494,11 @@ namespace Multiplayer
return updateMessage;
}
void EntityReplicator::FinalizeSerialization(AzNetworking::PacketId sentId)
{
m_propertyPublisher->FinalizeSerialization(sentId);
}
void EntityReplicator::DeferRpcMessage(NetworkEntityRpcMessage& entityRpcMessage)
{
// Received rpc metrics, log rpc sent, number of bytes, and the componentId/rpcId for bandwidth metrics
@@ -336,7 +336,6 @@ namespace Multiplayer
case PropertyPublisher::EntityReplicatorState::Deleting:
{
AZ_Assert(m_serializationPhase == PropertyPublisher::EntityReplicatorSerializationPhase::Prepared, "Unexpected serialization phase");
FinalizeDeleteEntityRecord(sentId);
}
break;