| 132 | }; |
| 133 | |
| 134 | static 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 | |
| 190 | class MyServiceWithBatchStream : public test::EchoService { |
| 191 | public: |
nothing calls this directly
no test coverage detected