| 799 | namespace internal { |
| 800 | |
| 801 | void receive(Socket socket) |
| 802 | { |
| 803 | StreamingRequestDecoder* decoder = new StreamingRequestDecoder(); |
| 804 | |
| 805 | const size_t size = 80 * 1024; |
| 806 | char* data = new char[size]; |
| 807 | |
| 808 | Future<Nothing> recv_loop = process::loop( |
| 809 | None(), |
| 810 | [=] { |
| 811 | return socket.recv(data, size); |
| 812 | }, |
| 813 | [=](size_t length) -> Future<ControlFlow<Nothing>> { |
| 814 | if (length == 0) { |
| 815 | return Break(); // EOF. |
| 816 | } |
| 817 | |
| 818 | // Decode as much of the data as possible into HTTP requests. |
| 819 | const deque<Request*> requests = decoder->decode(data, length); |
| 820 | |
| 821 | if (requests.empty() && decoder->failed()) { |
| 822 | return Failure("Decoder error"); |
| 823 | } |
| 824 | |
| 825 | if (!requests.empty()) { |
| 826 | // Get the peer address to augment the requests. |
| 827 | Try<Address> address = socket.peer(); |
| 828 | |
| 829 | if (address.isError()) { |
| 830 | return Failure("Failed to get peer address: " + address.error()); |
| 831 | } |
| 832 | |
| 833 | foreach (Request* request, requests) { |
| 834 | request->client = address.get(); |
| 835 | process_manager->handle(socket, request); |
| 836 | } |
| 837 | } |
| 838 | |
| 839 | return Continue(); |
| 840 | }); |
| 841 | |
| 842 | recv_loop.onAny([=](const Future<Nothing> f) { |
| 843 | if (f.isFailed()) { |
| 844 | Try<Address> peer = socket.peer(); |
| 845 | |
| 846 | LOG(WARNING) |
| 847 | << "Failed to recv on socket " << socket.get() << " to peer '" |
| 848 | << (peer.isSome() ? stringify(peer.get()) : "unknown") |
| 849 | << "': " << f.failure(); |
| 850 | } |
| 851 | |
| 852 | socket_manager->close(socket); |
| 853 | delete[] data; |
| 854 | delete decoder; |
| 855 | }); |
| 856 | } |
| 857 | |
| 858 | } // namespace internal { |
no test coverage detected