| 358 | if (not client.getReadableTransportStream().eventClose.addListener(closeListener)) |
| 359 | { |
| 360 | closeAsync(client); |
| 361 | } |
| 362 | return; |
| 363 | } |
| 364 | |
| 365 | // Register before starting the body stream, which can synchronously emit buffered body data. |
| 366 | struct AfterWrite |
| 367 | { |
| 368 | HttpAsyncServer& pself; |
| 369 | HttpConnection& client; |
| 370 | |
| 371 | void operator()() |
| 372 | { |
| 373 | SC_HTTP_ASSERT_RELEASE(client.response.getWritableStream().eventFinish.removeListener(*this)); |
| 374 | |
| 375 | // Determine if we should keep the connection alive |
| 376 | const bool underMaxRequests = |
| 377 | (pself.maxRequestsPerConnection == 0) or (client.requestCount + 1 < pself.maxRequestsPerConnection); |
| 378 | const bool shouldKeepAlive = client.response.getKeepAlive() and underMaxRequests and |
| 379 | not client.getReadableTransportStream().isEnded(); |
| 380 | |
| 381 | if (shouldKeepAlive and pself.state == State::Started) // We may get some after-writes after server stop |
| 382 | { |
| 383 | // Increment request count |
| 384 | client.requestCount++; |
| 385 | |
| 386 | // Reset request and response for next request |
| 387 | client.request.setHeaderMemory(client.getHeaderMemory()); |
| 388 | client.response.reset(); |
| 389 | |
| 390 | if (&client.getWritableTransportStream() == &client.writableSocketStream) |
| 391 | { |
| 392 | SC_HTTP_ASSERT_RELEASE(client.socket.isValid()); |
| 393 | Result writableRes = |
| 394 | client.writableSocketStream.init(client.buffersPool, *pself.eventLoop, client.socket); |
| 395 | SC_HTTP_TRUST_RESULT(writableRes); |
| 396 | } |
| 397 | else if (pself.transportReuse.isValid()) |
| 398 | { |
| 399 | Result reuseResult = pself.transportReuse(client); |
| 400 | if (not reuseResult) |
| 401 | { |
| 402 | if (pself.onError.isValid()) |
| 403 | { |
| 404 | pself.onError(reuseResult); |
| 405 | } |
| 406 | pself.closeAsync(client); |
| 407 | return; |
| 408 | } |
| 409 | } |
| 410 | |
| 411 | // Resume reading in any case to avoid deadlocking |
| 412 | client.getReadableTransportStream().resumeReading(); |
| 413 | |
| 414 | // Re-register for next request headers |
| 415 | EventDataListener dataListener{pself, client}; |
| 416 | SC_HTTP_ASSERT_RELEASE(client.getReadableTransportStream().eventData.addListener(dataListener)); |
| 417 | EventEndListener endListener{pself, client}; |
no test coverage detected