载入中...
搜索中...
未找到
PixelWorldTransport.cpp
浏览该文件的文档.
2
3#include <algorithm>
4#include <limits>
5#include <string>
6
8namespace {
9
10bool valid(PixelChunkTransportConfig config) {
11 return config.maximumInFlightTransfers > 0 && config.maximumInFlightTransfers <= 1024 &&
12 config.maximumChunksPerPart > 0 && config.maximumChunksPerPart <= 1024 &&
13 config.maximumPartsPerTransfer > 0 && config.maximumPartsPerTransfer <= 65'536 &&
14 config.maximumBufferedTransfers > 0 && config.maximumBufferedTransfers <= 1024;
15}
16
17bool sameMetadata(const PixelChunkTransferPart& left, const PixelChunkTransferPart& right) {
18 return left.streamId == right.streamId && left.transferId == right.transferId &&
19 left.partCount == right.partCount && left.interest == right.interest &&
20 left.batch.catalogFingerprint == right.batch.catalogFingerprint &&
21 left.batch.sourceSeed == right.batch.sourceSeed &&
22 left.batch.sourceRevision == right.batch.sourceRevision &&
23 left.batch.sourceTick == right.batch.sourceTick &&
24 left.batch.sourceLastEditSequence == right.batch.sourceLastEditSequence &&
25 left.batch.fullResync == right.batch.fullResync;
26}
27
28} // namespace
29
30ReliablePixelChunkSender::ReliablePixelChunkSender(std::uint64_t streamId,
31 PixelChunkTransportConfig config)
32 : streamId_(streamId), config_(config) {}
33
34eve::Result<ReliablePixelChunkSender> ReliablePixelChunkSender::create(
35 std::uint64_t streamId, PixelChunkTransportConfig config) {
36 if (streamId == 0 || !valid(config))
38 eve::Diagnostic::error(eve::DiagnosticCode::InvalidArgument, "stream id or transport policy is invalid",
39 "config", {}, "pixelworld.reliable-transport"));
41 ReliablePixelChunkSender(streamId, config));
42}
43
44eve::Result<std::vector<PixelChunkTransferPart>> ReliablePixelChunkSender::capture(
46 if (inFlight_.size() >= config_.maximumInFlightTransfers)
49 "maximumInFlightTransfers", {}, "pixelworld.reliable-transport"));
50 if (nextTransferId_ == std::numeric_limits<std::uint64_t>::max())
53 "transferId", {}, "pixelworld.reliable-transport"));
54 auto candidateCursor = cursor_;
55 auto captured = candidateCursor.capture(source, interest);
56 if (!captured.ok())
57 return eve::Result<std::vector<PixelChunkTransferPart>>::failure(captured.status());
58 auto update = std::move(captured).takeValue();
59 if (captured_ && !update.fullResync && update.batch.chunks.empty() &&
60 update.batch.sourceRevision == lastEnqueuedRevision_) {
61 cursor_ = std::move(candidateCursor);
63 }
64
65 const std::size_t chunkCount = update.batch.chunks.size();
66 const std::size_t partCount = std::max<std::size_t>(
67 1, (chunkCount + config_.maximumChunksPerPart - 1) / config_.maximumChunksPerPart);
68 if (partCount > config_.maximumPartsPerTransfer)
70 eve::Diagnostic::error(eve::DiagnosticCode::PreconditionViolation, "captured update exceeds part budget",
71 "maximumPartsPerTransfer", {}, "pixelworld.reliable-transport"));
72 std::vector<PixelChunkTransferPart> parts;
73 parts.reserve(partCount);
74 for (std::size_t index = 0; index < partCount; ++index) {
76 part.streamId = streamId_;
77 part.transferId = nextTransferId_;
78 part.partIndex = std::uint32_t(index);
79 part.partCount = std::uint32_t(partCount);
80 part.interest = interest;
81 part.batch.catalogFingerprint = update.batch.catalogFingerprint;
82 part.batch.sourceSeed = update.batch.sourceSeed;
83 part.batch.sourceRevision = update.batch.sourceRevision;
84 part.batch.sourceTick = update.batch.sourceTick;
85 part.batch.sourceLastEditSequence = update.batch.sourceLastEditSequence;
86 part.batch.fullResync = update.batch.fullResync;
87 const std::size_t begin = index * config_.maximumChunksPerPart;
88 const std::size_t end = std::min(chunkCount, begin + config_.maximumChunksPerPart);
89 if (begin < end)
90 part.batch.chunks.assign(update.batch.chunks.begin() + std::ptrdiff_t(begin),
91 update.batch.chunks.begin() + std::ptrdiff_t(end));
92 parts.push_back(std::move(part));
93 }
94 inFlight_.emplace(nextTransferId_, parts);
95 cursor_ = std::move(candidateCursor);
96 ++nextTransferId_;
97 captured_ = true;
98 lastEnqueuedRevision_ = update.batch.sourceRevision;
99 return eve::Result<std::vector<PixelChunkTransferPart>>::success(std::move(parts));
100}
101
102std::vector<PixelChunkTransferPart> ReliablePixelChunkSender::pendingParts() const {
103 std::vector<PixelChunkTransferPart> result;
104 for (const auto& [transfer, parts] : inFlight_) {
105 (void)transfer;
106 result.insert(result.end(), parts.begin(), parts.end());
107 }
108 return result;
109}
110
111eve::Result<PixelChunkAckReceipt> ReliablePixelChunkSender::acknowledge(
113 if (ack.streamId != streamId_)
115 eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "ACK belongs to another stream", "streamId", {},
116 "pixelworld.reliable-transport"));
117 if (ack.acknowledgedThrough <= acknowledgedThrough_)
119 {0, 0, acknowledgedThrough_});
120 if (ack.acknowledgedThrough >= nextTransferId_)
122 eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "ACK advances beyond the sent window",
123 "acknowledgedThrough", {}, "pixelworld.reliable-transport"));
124 const auto acknowledged = inFlight_.find(ack.acknowledgedThrough);
125 if (acknowledged == inFlight_.end() || acknowledged->second.empty() ||
126 acknowledged->second.front().batch.sourceRevision != ack.appliedRevision)
128 eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "ACK revision does not match the transfer",
129 "appliedRevision", {}, "pixelworld.reliable-transport"));
130 PixelChunkAckReceipt receipt;
131 for (auto iterator = inFlight_.begin(); iterator != inFlight_.end() &&
132 iterator->first <= ack.acknowledgedThrough;) {
133 for (const auto& part : iterator->second)
134 receipt.tombstonesReleased += std::uint32_t(std::count_if(
135 part.batch.chunks.begin(), part.batch.chunks.end(),
136 [](const auto& chunk) { return chunk.removed; }));
137 iterator = inFlight_.erase(iterator);
138 ++receipt.transfersReleased;
139 }
140 acknowledgedThrough_ = ack.acknowledgedThrough;
141 receipt.acknowledgedThrough = acknowledgedThrough_;
143}
144
145std::size_t ReliablePixelChunkSender::inFlightTransferCount() const noexcept {
146 return inFlight_.size();
147}
148
149std::size_t ReliablePixelChunkSender::retainedTombstoneCount() const noexcept {
150 std::size_t count = 0;
151 for (const auto& [transfer, parts] : inFlight_) {
152 (void)transfer;
153 for (const auto& part : parts)
154 count += std::count_if(part.batch.chunks.begin(), part.batch.chunks.end(),
155 [](const auto& chunk) { return chunk.removed; });
156 }
157 return count;
158}
159
160ReliablePixelChunkReceiver::ReliablePixelChunkReceiver(std::uint64_t streamId,
162 : streamId_(streamId), config_(config) {}
163
165 std::uint64_t streamId, PixelChunkTransportConfig config) {
166 if (streamId == 0 || !valid(config))
168 eve::Diagnostic::error(eve::DiagnosticCode::InvalidArgument, "stream id or transport policy is invalid",
169 "config", {}, "pixelworld.reliable-transport"));
171 ReliablePixelChunkReceiver(streamId, config));
172}
173
176 if (part.streamId != streamId_ || part.transferId == 0 || part.partCount == 0 ||
177 part.partCount > config_.maximumPartsPerTransfer || part.partIndex >= part.partCount ||
178 part.batch.chunks.size() > config_.maximumChunksPerPart)
180 eve::Diagnostic::error(eve::DiagnosticCode::InvalidArgument, "transfer part metadata exceeds policy",
181 "part", {}, "pixelworld.reliable-transport"));
183 if (part.transferId < expectedTransferId_) {
184 receipt.duplicate = true;
185 receipt.ack = {streamId_, expectedTransferId_ - 1, replica.revision()};
186 receipt.partsBuffered = std::uint32_t(bufferedPartCount());
188 }
189 auto found = buffered_.find(part.transferId);
190 if (found == buffered_.end()) {
191 if (buffered_.size() >= config_.maximumBufferedTransfers)
193 eve::DiagnosticCode::PreconditionViolation, "out-of-order receive window is full",
194 "maximumBufferedTransfers", {}, "pixelworld.reliable-transport"));
195 Assembly assembly;
196 assembly.partCount = part.partCount;
197 assembly.parts.resize(part.partCount);
198 found = buffered_.emplace(part.transferId, std::move(assembly)).first;
199 }
200 Assembly& assembly = found->second;
201 if (assembly.partCount != part.partCount)
203 eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "transfer part count changed", "partCount", {},
204 "pixelworld.reliable-transport"));
205 auto& slot = assembly.parts[part.partIndex];
206 if (slot) {
207 if (*slot != part)
209 eve::Diagnostic::error(eve::DiagnosticCode::Conflict, "duplicate part payload conflicts", "part", {},
210 "pixelworld.reliable-transport"));
211 receipt.duplicate = true;
212 } else {
213 slot = std::make_unique<PixelChunkTransferPart>(std::move(part));
214 }
215
216 while (true) {
217 auto next = buffered_.find(expectedTransferId_);
218 if (next == buffered_.end()) break;
219 const bool complete = std::all_of(next->second.parts.begin(), next->second.parts.end(),
220 [](const auto& value) { return bool(value); });
221 if (!complete) break;
222 const PixelChunkTransferPart& first = *next->second.parts.front();
224 combined.chunks.clear();
225 for (const auto& value : next->second.parts) {
226 if (!sameMetadata(first, *value))
228 eve::DiagnosticCode::Conflict, "transfer parts disagree on authority metadata", "part.batch", {},
229 "pixelworld.reliable-transport"));
230 combined.chunks.insert(combined.chunks.end(), value->batch.chunks.begin(),
231 value->batch.chunks.end());
232 }
233 auto applied = replica.applyChunkBatch(combined, replica.revision());
234 if (!applied.ok()) {
235 buffered_.erase(next);
236 return eve::Result<PixelChunkReceiveReceipt>::failure(applied.status());
237 }
238 buffered_.erase(next);
239 ++expectedTransferId_;
240 ++receipt.transfersApplied;
241 }
242 receipt.ack = {streamId_, expectedTransferId_ - 1, replica.revision()};
243 receipt.partsBuffered = std::uint32_t(bufferedPartCount());
245}
246
248 std::size_t count = 0;
249 for (const auto& [transfer, assembly] : buffered_) {
250 (void)transfer;
251 count += std::count_if(assembly.parts.begin(), assembly.parts.end(),
252 [](const auto& value) { return bool(value); });
253 }
254 return count;
255}
256
257} // namespace eve::pixelworld_streaming
double value
std::vector< eve::artifact::PartView > parts
HexVec3 left
HexVec3 right
std::int32_t first
bool valid
float begin
bool found
std::uint32_t count
uint32_t index
const UnitySourceAsset & source
static Diagnostic error(DiagnosticCode code, std::string message, std::string path={}, DiagnosticDetails details={}, std::string source={})
Construct an error diagnostic with the standard error severity.
Definition Diagnostic.h:125
Move-only operation result carrying either a value or Status.
Definition Result.h:155
static Result success(T value)
Construct a successful result owning value.
Definition Result.h:164
static Result failure(Status status)
Construct a failed result from a structured status.
Definition Result.h:175
Sparse, chunked, deterministic 2D falling-material world.
Definition PixelWorld.h:224
eve::Result< PixelChunkApplyReceipt > applyChunkBatch(const PixelChunkBatch &batch, std::uint64_t expectedRevision)
Transactionally apply a canonical authoritative Chunk correction batch.
std::uint64_t revision() const noexcept
Revision.
Ordered atomic receiver with duplicate suppression and bounded out-of-order buffering.
eve::Result< PixelChunkReceiveReceipt > receive(PixelChunkTransferPart part, eve::pixelworld::PixelWorld &replica)
Buffer one part and atomically apply every newly contiguous complete transfer.
std::size_t bufferedPartCount() const noexcept
Count parts retained while waiting for gaps or completion.
static eve::Result< ReliablePixelChunkReceiver > create(std::uint64_t streamId, PixelChunkTransportConfig config={})
Validate policy and create a receiver for one nonzero stream session.
Reliable sender retaining owning transfer parts and tombstones until application ACK.
Owning authoritative Chunk correction with source world metadata.
Definition PixelWorld.h:89
bool fullResync
Replace the replica projection instead of applying an incremental correction.
Definition PixelWorld.h:96
std::vector< PixelChunkSnapshot > chunks
Definition PixelWorld.h:97
eve::SimulationTick sourceTick
Definition PixelWorld.h:93
Inclusive finite rectangle expressed in Chunk coordinates.
Definition PixelWorld.h:104
Counters returned after accepting a cumulative ACK.
Receiver outcome after one possibly duplicate or out-of-order part.
Cumulative application ACK proving transfers were committed to replica authority.
One independently retransmittable part of an atomic Chunk transfer.
Bounded application-level reliability policy layered over any byte transport.