MCPcopy Create free account
hub / github.com/apache/brpc / SendTwoMessagesOnServerExtraStream

Function SendTwoMessagesOnServerExtraStream

test/brpc_streaming_rpc_unittest.cpp:134–188  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

132};
133
134static void* SendTwoMessagesOnServerExtraStream(void* arg) {
135 auto* state = static_cast<BatchStreamFeedbackRaceState*>(arg);
136 const brpc::StreamId sid = state->server_extra_stream_id;
137
138 // Wait until server-side stream is connected.
139 const int64_t connect_deadline_us = butil::gettimeofday_us() + 2 * 1000 * 1000L;
140 bool connected = false;
141 while (butil::gettimeofday_us() < connect_deadline_us) {
142 brpc::SocketUniquePtr ptr;
143 if (brpc::Socket::Address(sid, &ptr) == 0) {
144 brpc::Stream* s = static_cast<brpc::Stream*>(ptr->conn());
145 if (s->_host_socket != NULL && s->_connected) {
146 connected = true;
147 break;
148 }
149 }
150 usleep(1000);
151 }
152
153 if (!connected) {
154 state->server_first_write_rc.store(ETIMEDOUT, std::memory_order_relaxed);
155 state->server_second_write_rc.store(ETIMEDOUT, std::memory_order_relaxed);
156 state->server_write_done.store(true, std::memory_order_release);
157 return NULL;
158 }
159
160 // 1) Send a payload exactly equal to max_buf_size(64).
161 {
162 std::string payload(64, 'a');
163 butil::IOBuf out;
164 out.append(payload);
165 state->server_first_write_rc.store(brpc::StreamWrite(sid, out), std::memory_order_relaxed);
166 }
167
168 // 2) Then send another byte. This write should become writable only after
169 // client sends FEEDBACK with consumed_size >= 64.
170 const int64_t write_deadline_us = butil::gettimeofday_us() + 2 * 1000 * 1000L;
171 int rc = -1;
172 while (butil::gettimeofday_us() < write_deadline_us) {
173 butil::IOBuf out;
174 out.append("b", 1);
175 rc = brpc::StreamWrite(sid, out);
176 if (rc == 0) {
177 break;
178 }
179 if (rc != EAGAIN) {
180 break;
181 }
182 const timespec duetime = butil::milliseconds_from_now(100);
183 (void)brpc::StreamWait(sid, &duetime);
184 }
185 state->server_second_write_rc.store(rc, std::memory_order_relaxed);
186 state->server_write_done.store(true, std::memory_order_release);
187 return NULL;
188}
189
190class MyServiceWithBatchStream : public test::EchoService {
191public:

Callers

nothing calls this directly

Calls 6

StreamWriteFunction · 0.85
milliseconds_from_nowFunction · 0.85
StreamWaitFunction · 0.85
gettimeofday_usFunction · 0.70
storeMethod · 0.45
appendMethod · 0.45

Tested by

no test coverage detected