Skip to content
Closed
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions envoy/event/dispatcher.h
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,14 @@ class Dispatcher : public DispatcherBase, public ScopeTracker {
Network::TransportSocketPtr&& transport_socket,
const Network::ConnectionSocket::OptionsSharedPtr& options) PURE;

/**
* Register an internal listener manager for this dispatcher.
*/
virtual void
registerInternalListenerManager(Network::InternalListenerManager& internal_listener_manager) PURE;

virtual Network::InternalListenerManagerOptRef getInternalListenerManager() PURE;

/**
* Creates an async DNS resolver. The resolver should only be used on the thread that runs this
* dispatcher.
Expand Down
11 changes: 11 additions & 0 deletions envoy/network/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,17 @@ envoy_cc_library(
],
)

envoy_cc_library(
name = "client_connection_manager",
hdrs = ["client_connection_manager.h"],
deps = [
":address_interface",
":connection_interface",
":listen_socket_interface",
":transport_socket_interface",
],
)

envoy_cc_library(
name = "connection_handler_interface",
hdrs = ["connection_handler.h"],
Expand Down
1 change: 1 addition & 0 deletions envoy/network/client_connection_manager.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
// NOLINT(namespace-envoy)
Comment thread
lambdai marked this conversation as resolved.
Outdated
26 changes: 26 additions & 0 deletions envoy/network/client_connection_manager.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
#pragma once

#include "envoy/network/address.h"
#include "envoy/network/connection.h"
#include "envoy/network/listen_socket.h"
#include "envoy/network/transport_socket.h"

namespace Envoy {
namespace Network {

class ClientConnectionFactory {
Comment thread
lambdai marked this conversation as resolved.
Outdated
public:
virtual ~ClientConnectionFactory() = default;
std::string category() { return "network.connection"; }
virtual std::string name() PURE;

virtual Network::ClientConnectionPtr
createClientConnection(Event::Dispatcher& dispatcher,
Network::Address::InstanceConstSharedPtr address,
Network::Address::InstanceConstSharedPtr source_address,
Network::TransportSocketPtr&& transport_socket,
const Network::ConnectionSocket::OptionsSharedPtr& options) PURE;
};

} // namespace Network
} // namespace Envoy
44 changes: 44 additions & 0 deletions envoy/network/listener.h
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,16 @@ class UdpListenerConfig {

using UdpListenerConfigOptRef = OptRef<UdpListenerConfig>;

/**
* Configuration for an internal listener.
*/
class InternalListenerConfig {
public:
virtual ~InternalListenerConfig() = default;
};

using InternalListenerConfigOptRef = OptRef<InternalListenerConfig>;

/**
* A configuration for an individual listener.
*/
Expand Down Expand Up @@ -184,6 +194,11 @@ class ListenerConfig {
*/
virtual UdpListenerConfigOptRef udpListenerConfig() PURE;

/**
* @return the internal configuration for the listener IFF it is an internal listener.
*/
virtual InternalListenerConfigOptRef internalListenerConfig() PURE;

/**
* @return traffic direction of the listener.
*/
Expand Down Expand Up @@ -426,6 +441,35 @@ class UdpListener : public virtual Listener {
};

using UdpListenerPtr = std::unique_ptr<UdpListener>;
class InternalListenerCallbacks {
public:
virtual ~InternalListenerCallbacks() = default;

/**
* Called when a new connection is accepted.
* @param socket supplies the socket that is moved into the callee.
*/
virtual void onAccept(ConnectionSocketPtr&& socket) PURE;

virtual Event::Dispatcher& dispatcher() PURE;
};
using InternalListenerCallbacksOptRef =
absl::optional<std::reference_wrapper<InternalListenerCallbacks>>;

class InternalListener {};

using InternalListenerPtr = std::unique_ptr<InternalListener>;
using InternalListenerOptRef = absl::optional<std::reference_wrapper<InternalListener>>;

class InternalListenerManager {
Comment thread
lambdai marked this conversation as resolved.
public:
virtual ~InternalListenerManager() = default;
virtual InternalListenerCallbacksOptRef
findByAddress(const Address::InstanceConstSharedPtr& listen_address) PURE;
};

using InternalListenerManagerOptRef =
absl::optional<std::reference_wrapper<InternalListenerManager>>;

/**
* Handles delivering datagrams to the correct worker.
Expand Down
2 changes: 2 additions & 0 deletions source/common/event/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ envoy_cc_library(
"//conditions:default": "posix",
}),
deps = [
"//envoy/network:client_connection_manager",
"//source/common/config:utility_lib",
":dispatcher_includes",
":libevent_scheduler_lib",
":real_time_system_lib",
Expand Down
26 changes: 24 additions & 2 deletions source/common/event/dispatcher_impl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -9,13 +9,15 @@
#include "envoy/api/api.h"
#include "envoy/common/scope_tracker.h"
#include "envoy/config/overload/v3/overload.pb.h"
#include "envoy/network/client_connection_manager.h"
#include "envoy/network/listen_socket.h"
#include "envoy/network/listener.h"

#include "source/common/buffer/buffer_impl.h"
#include "source/common/common/assert.h"
#include "source/common/common/lock_guard.h"
#include "source/common/common/thread.h"
#include "source/common/config/utility.h"
#include "source/common/event/file_event_impl.h"
#include "source/common/event/libevent_scheduler.h"
#include "source/common/event/scaled_range_timer_manager_impl.h"
Expand Down Expand Up @@ -153,8 +155,28 @@ DispatcherImpl::createClientConnection(Network::Address::InstanceConstSharedPtr
Network::TransportSocketPtr&& transport_socket,
const Network::ConnectionSocket::OptionsSharedPtr& options) {
ASSERT(isThreadSafe());
return std::make_unique<Network::ClientConnectionImpl>(*this, address, source_address,
std::move(transport_socket), options);

auto address_type_name = [](const Network::Address::InstanceConstSharedPtr& addr) {
ASSERT(addr != nullptr);
switch (addr->type()) {
Comment thread
lambdai marked this conversation as resolved.
Outdated
// TODO: create IP and
case Network::Address::Type::Ip:
case Network::Address::Type::Pipe:
return "does-not-exist";
case Network::Address::Type::EnvoyInternal:
return "EnvoyInternal";
}
};
// TODO: register and find by address type instead of name.
auto factory = Config::Utility::getFactoryByName<Network::ClientConnectionFactory>(
address_type_name(address));
if (factory == nullptr) {
// get rid of this once the ip and pipe factory is offered.
return std::make_unique<Network::ClientConnectionImpl>(*this, address, source_address,
std::move(transport_socket), options);
}
return factory->createClientConnection(*this, address, source_address,
std::move(transport_socket), options);
}

Network::DnsResolverSharedPtr DispatcherImpl::createDnsResolver(
Expand Down
12 changes: 12 additions & 0 deletions source/common/event/dispatcher_impl.h
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,17 @@ class DispatcherImpl : Logger::Loggable<Logger::Id::main>,
Network::Address::InstanceConstSharedPtr source_address,
Network::TransportSocketPtr&& transport_socket,
const Network::ConnectionSocket::OptionsSharedPtr& options) override;

void registerInternalListenerManager(
Network::InternalListenerManager& internal_listener_manager) override {
ASSERT(!internal_listener_manager_.has_value());
internal_listener_manager_ = internal_listener_manager;
}

Network::InternalListenerManagerOptRef getInternalListenerManager() override {
return internal_listener_manager_;
}

Network::DnsResolverSharedPtr createDnsResolver(
const std::vector<Network::Address::InstanceConstSharedPtr>& resolvers,
const envoy::config::core::v3::DnsResolverOptions& dns_resolver_options) override;
Expand Down Expand Up @@ -175,6 +186,7 @@ class DispatcherImpl : Logger::Loggable<Logger::Id::main>,
MonotonicTime approximate_monotonic_time_;
WatchdogRegistrationPtr watchdog_registration_;
const ScaledRangeTimerManagerPtr scaled_timer_manager_;
Network::InternalListenerManagerOptRef internal_listener_manager_;
};

} // namespace Event
Expand Down
19 changes: 19 additions & 0 deletions source/extensions/io_socket/user_space/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ envoy_cc_extension(
name = "config",
srcs = ["config.h"],
deps = [
":client_connection_factory",
":io_handle_impl_lib",
],
)
Expand Down Expand Up @@ -57,3 +58,21 @@ envoy_cc_library(
"//source/common/network:default_socket_interface_lib",
],
)

envoy_cc_library(
name = "client_connection_factory",
srcs = [
"client_connection_factory.cc",
],
hdrs = [
"client_connection_factory.h",
],
deps = [
":io_handle_impl_lib",
"//envoy/network:client_connection_manager",
"//envoy/network:connection_interface",
"//envoy/registry",
"//source/common/network:connection_lib",
"//source/common/network:listen_socket_lib",
],
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
#include "source/extensions/io_socket/user_space/client_connection_factory.h"

#include "envoy/registry/registry.h"

#include "source/common/network/address_impl.h"
#include "source/common/network/connection_impl.h"
#include "source/common/network/listen_socket_impl.h"
#include "source/extensions/io_socket/user_space/io_handle_impl.h"

namespace Envoy {

namespace Extensions {
namespace IoSocket {
namespace UserSpace {

Network::ClientConnectionPtr InternalClientConnectionFactory::createClientConnection(
Event::Dispatcher& dispatcher, Network::Address::InstanceConstSharedPtr address,
Network::Address::InstanceConstSharedPtr source_address,
Network::TransportSocketPtr&& transport_socket,
const Network::ConnectionSocket::OptionsSharedPtr& options) {

Network::IoHandlePtr io_handle_client;
Network::IoHandlePtr io_handle_server;

std::tie(io_handle_client, io_handle_server) =
Extensions::IoSocket::UserSpace::IoHandleFactory::createIoHandlePair();

auto client_conn = std::make_unique<Network::ClientConnectionImpl>(
dispatcher,
std::make_unique<Network::ConnectionSocketImpl>(std::move(io_handle_client), source_address,
address),
source_address, std::move(transport_socket), options);
// TODO(lambdai): refactor nested if.
auto internal_listener_manager = dispatcher.getInternalListenerManager();
if (internal_listener_manager.has_value()) {
// It's either in main thread or the worker is not yet started.
auto internal_listener = internal_listener_manager.value().get().findByAddress(address);
if (internal_listener.has_value()) {
auto original_address = address;
// if (options != nullptr) {
// for (const auto& opt : *options) {
// auto* internal_opt = dynamic_cast<const Network::InternalSocketOptionImpl*>(opt.get());
// if (internal_opt != nullptr) {
// original_address = internal_opt->original_remote_address_;
// }
// }
// }
auto accepted_socket = std::make_unique<Network::AcceptedSocketImpl>(
std::move(io_handle_server), original_address, source_address);
// TODO: also check if disabled
internal_listener.value().get().onAccept(std::move(accepted_socket));
FANCY_LOG(debug, "lambdai: find internal listener {} ", address->asStringView());
} else {
FANCY_LOG(debug, "lambdai: cannot find internal listener {} ", address->asStringView());
// injected error into client_conn;
io_handle_server->close();
}
} else {
FANCY_LOG(debug, "lambdai: cannot find internal listener {} ", address->asStringView());
// injected error into client_conn;
io_handle_server->close();
}
return client_conn;
}
REGISTER_FACTORY(InternalClientConnectionFactory, Network::ClientConnectionFactory);

} // namespace UserSpace
} // namespace IoSocket
} // namespace Extensions
} // namespace Envoy
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
#pragma once

#include <memory>
#include <string>

#include "envoy/common/pure.h"
#include "envoy/network/client_connection_manager.h"
#include "envoy/network/connection.h"

namespace Envoy {

namespace Extensions {
namespace IoSocket {
namespace UserSpace {

class InternalClientConnectionFactory : public Network::ClientConnectionFactory {
public:
std::string name() override { return "EnvoyInternal"; }
Network::ClientConnectionPtr
createClientConnection(Event::Dispatcher& dispatcher,
Network::Address::InstanceConstSharedPtr address,
Network::Address::InstanceConstSharedPtr source_address,
Network::TransportSocketPtr&& transport_socket,
const Network::ConnectionSocket::OptionsSharedPtr& options) override;
};

} // namespace UserSpace
} // namespace IoSocket
} // namespace Extensions
} // namespace Envoy
3 changes: 3 additions & 0 deletions source/server/admin/admin.h
Original file line number Diff line number Diff line change
Expand Up @@ -360,6 +360,9 @@ class AdminImpl : public Admin,
Network::UdpListenerConfigOptRef udpListenerConfig() override {
return Network::UdpListenerConfigOptRef();
}
Network::InternalListenerConfigOptRef internalListenerConfig() override {
return Network::InternalListenerConfigOptRef();
}
envoy::config::core::v3::TrafficDirection direction() const override {
return envoy::config::core::v3::UNSPECIFIED;
}
Expand Down
3 changes: 3 additions & 0 deletions source/server/listener_impl.h
Original file line number Diff line number Diff line change
Expand Up @@ -313,6 +313,9 @@ class ListenerImpl final : public Network::ListenerConfig,
return udp_listener_config_ != nullptr ? *udp_listener_config_
: Network::UdpListenerConfigOptRef();
}
Network::InternalListenerConfigOptRef internalListenerConfig() override {
return Network::InternalListenerConfigOptRef();
}
Network::ConnectionBalancer& connectionBalancer() override { return *connection_balancer_; }
ResourceLimit& openConnections() override { return *open_connections_; }
const std::vector<AccessLog::InstanceSharedPtr>& accessLogs() const override {
Expand Down
1 change: 1 addition & 0 deletions test/common/event/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ envoy_cc_test(
"//source/common/event:deferred_task",
"//source/common/event:dispatcher_includes",
"//source/common/event:dispatcher_lib",
"//source/common/network:address_lib",
"//source/common/stats:isolated_store_lib",
"//test/mocks:common_lib",
"//test/mocks/server:watch_dog_mocks",
Expand Down
Loading