/* * Copyright (C) 2014 Cloudius Systems, Ltd. */ #ifndef NET_NATIVE_STACK_IMPL_HH_ #define NET_NATIVE_STACK_IMPL_HH_ #include "core/reactor.hh" namespace net { template class native_server_socket_impl; template class native_connected_socket_impl; class native_network_stack; // native_server_socket_impl template class native_server_socket_impl : public server_socket_impl { typename Protocol::listener _listener; public: native_server_socket_impl(Protocol& proto, uint16_t port, listen_options opt); virtual future accept() override; }; template native_server_socket_impl::native_server_socket_impl(Protocol& proto, uint16_t port, listen_options opt) : _listener(proto.listen(port)) { } template future native_server_socket_impl::accept() { return _listener.accept().then([this] (typename Protocol::connection conn) { return make_ready_future( connected_socket(std::make_unique>(std::move(conn))), socket_address()); // FIXME: don't fake it }); } // native_connected_socket_impl template class native_connected_socket_impl : public connected_socket_impl { typename Protocol::connection _conn; class native_data_source_impl; class native_data_sink_impl; public: explicit native_connected_socket_impl(typename Protocol::connection conn) : _conn(std::move(conn)) {} virtual input_stream input() override; virtual output_stream output() override; }; template class native_connected_socket_impl::native_data_source_impl final : public data_source_impl { typename Protocol::connection& _conn; size_t _cur_frag = 0; bool _eof = false; packet _buf; public: explicit native_data_source_impl(typename Protocol::connection& conn) : _conn(conn) {} virtual future> get() override { if (_eof) { return make_ready_future>(temporary_buffer(0)); } if (_cur_frag != _buf.nr_frags()) { auto& f = _buf.fragments()[_cur_frag++]; return make_ready_future>( temporary_buffer(f.base, f.size, make_deleter(deleter(), [p = _buf.share()] () mutable {}))); } return _conn.wait_for_data().then([this] { _buf = _conn.read(); _cur_frag = 0; _eof = !_buf.len(); return get(); }); } }; template class native_connected_socket_impl::native_data_sink_impl final : public data_sink_impl { typename Protocol::connection& _conn; public: explicit native_data_sink_impl(typename Protocol::connection& conn) : _conn(conn) {} virtual future<> put(std::vector> data) override { std::vector frags; frags.reserve(data.size()); for (auto& e : data) { frags.push_back(fragment{e.get_write(), e.size()}); } return _conn.send(packet(std::move(frags), [tmp = std::move(data)] () mutable {})); } virtual future<> put(temporary_buffer data) override { return _conn.send(packet({data.get_write(), data.size()}, data.release())); } virtual future<> close() override { _conn.close_write(); return make_ready_future<>(); } }; template input_stream native_connected_socket_impl::input() { data_source ds(std::make_unique(_conn)); return input_stream(std::move(ds)); } template output_stream native_connected_socket_impl::output() { data_sink ds(std::make_unique(_conn)); return output_stream(std::move(ds), 8192); } } #endif /* NET_NATIVE_STACK_IMPL_HH_ */