| 4702 | // User of the stream has to forward the SS's responses to the returned promise stream, if it is set |
| 4703 | template <class Request, bool P> |
| 4704 | Optional<TSSDuplicateStreamData<REPLYSTREAM_TYPE(Request)>> |
| 4705 | maybeDuplicateTSSStreamFragment(Request& req, QueueModel* model, RequestStream<Request, P> const* ssStream) { |
| 4706 | if (model) { |
| 4707 | Optional<TSSEndpointData> tssData = model->getTssData(ssStream->getEndpoint().token.first()); |
| 4708 | |
| 4709 | if (tssData.present()) { |
| 4710 | CODE_PROBE(true, "duplicating stream to TSS"); |
| 4711 | resetReply(req); |
| 4712 | // FIXME: optimize to avoid creating new netNotifiedQueueWithAcknowledgements for each stream duplication |
| 4713 | RequestStream<Request> tssRequestStream(tssData.get().endpoint); |
| 4714 | ReplyPromiseStream<REPLYSTREAM_TYPE(Request)> tssReplyStream = tssRequestStream.getReplyStream(req); |
| 4715 | PromiseStream<REPLYSTREAM_TYPE(Request)> ssDuplicateReplyStream; |
| 4716 | TSSDuplicateStreamData<REPLYSTREAM_TYPE(Request)> streamData(ssDuplicateReplyStream); |
| 4717 | model->addActor.send(tssStreamComparison(req, streamData, tssReplyStream, tssData.get())); |
| 4718 | return Optional<TSSDuplicateStreamData<REPLYSTREAM_TYPE(Request)>>(streamData); |
| 4719 | } |
| 4720 | } |
| 4721 | return Optional<TSSDuplicateStreamData<REPLYSTREAM_TYPE(Request)>>(); |
| 4722 | } |
| 4723 | |
| 4724 | // Streams all of the KV pairs in a target key range into a ParallelStream fragment |
| 4725 | ACTOR Future<Void> getRangeStreamFragment(Reference<TransactionState> trState, |
no test coverage detected