Skip to content
Draft
Show file tree
Hide file tree
Changes from all 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
150 changes: 131 additions & 19 deletions src/client.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,8 @@ constexpr timeval beaconCleanInterval{180, 0};
// special interval to attempt to reconnect to disconnected name servers
constexpr timeval tcpNSCheckInterval{10, 0};

constexpr timeval dnsRecheckInterval{10, 0};

// searchSequenceID in CMD_SEARCH is redundant.
// So we use a static value and instead rely on IDs for individual PVs
constexpr uint32_t search_seq{0x66696e64}; // "find"
Expand Down Expand Up @@ -217,13 +219,18 @@ void Channel::disconnect(const std::shared_ptr<Channel>& self)
name.c_str());

} else if(context->state==ContextImpl::Running) { // reconnect to specific server
if(!forcedServerHostname.empty()) {
forcedServer.setAddress(forcedServerHostname.c_str(), forcedServer.port());
log_info_printf(io, "Forced server re-resolved for '%s': %s\n",
name.c_str(), forcedServer.tostring().c_str());
}

conn = Connection::build(context, forcedServer, true);

conn->pending[cid] = self;
state = Connecting;

conn->createChannels();

}
}

Expand Down Expand Up @@ -382,6 +389,9 @@ std::shared_ptr<Channel> Channel::build(const std::shared_ptr<ContextImpl>& cont

} else { // bypass search and connect so a specific server
chan->forcedServer = forceServer;
if(isHostname(server)) {
chan->forcedServerHostname = server;
}
chan->conn = Connection::build(context, forceServer);

chan->conn->pending[chan->cid] = chan;
Expand Down Expand Up @@ -564,6 +574,8 @@ ContextImpl::ContextImpl(const Config& conf, const evbase& tcp_loop)
event_new(tcp_loop.base, -1, EV_TIMEOUT|EV_PERSIST, &ContextImpl::cacheCleanS, this))
,nsChecker(__FILE__, __LINE__,
event_new(tcp_loop.base, -1, EV_TIMEOUT|EV_PERSIST, &ContextImpl::onNSCheckS, this))
,dnsRecheckTimer(__FILE__, __LINE__,
event_new(tcp_loop.base, -1, EV_TIMEOUT|EV_PERSIST, &ContextImpl::onDNSRecheckS, this))
{
searchBuckets.resize(nBuckets);

Expand Down Expand Up @@ -614,8 +626,15 @@ ContextImpl::ContextImpl(const Config& conf, const evbase& tcp_loop)
if(isucast && ep.addr.family()==AF_INET && bcasts.find(ep.addr)!=bcasts.end())
isucast = false;

log_info_printf(io, "Searching to %s%s\n", std::string(SB()<<ep).c_str(), (isucast?" unicast":""));
searchDest.emplace_back(ep, isucast);
std::string hostname;
auto hit = effective.addressHostnames.find(addr);
if(hit != effective.addressHostnames.end())
hostname = hit->second;

log_info_printf(io, "Searching to %s%s%s\n", std::string(SB()<<ep).c_str(),
(isucast?" unicast":""),
(hostname.empty()?"":(std::string(" hostname=")+hostname).c_str()));
searchDest.emplace_back(ep, isucast, std::move(hostname));
}

for(auto& addr : effective.nameServers) {
Expand All @@ -624,10 +643,20 @@ ContextImpl::ContextImpl(const Config& conf, const evbase& tcp_loop)
saddr.setAddress(addr.c_str(), effective.tcp_port);
}catch(std::runtime_error& e) {
log_err_printf(setup, "%s Ignoring...\n", e.what());
continue;
}

log_info_printf(io, "Searching to TCP %s\n", saddr.tostring().c_str());
nameServers.emplace_back(saddr, nullptr);
std::string nsHostname;
auto hit = effective.nameServerHostnames.find(addr);
if(hit != effective.nameServerHostnames.end())
nsHostname = hit->second;

log_warn_printf(io, "TRACE nameServers loop: addr='%s' hostnameMapSize=%zu found=%d nsHostname='%s'\n",
addr.c_str(), effective.nameServerHostnames.size(), (int)(hit != effective.nameServerHostnames.end()), nsHostname.c_str());

log_info_printf(io, "Searching to TCP %s%s\n", saddr.tostring().c_str(),
(nsHostname.empty()?"":(std::string(" hostname=")+nsHostname).c_str()));
nameServers.push_back({saddr, nullptr, std::move(nsHostname)});
}

if(searchDest.empty() && nameServers.empty())
Expand Down Expand Up @@ -673,20 +702,27 @@ ContextImpl::~ContextImpl() {}

void ContextImpl::startNS()
{
if(nameServers.empty()) // vector size const after ctor, contents remain mutable
bool hasHostnames = std::any_of(searchDest.begin(), searchDest.end(),
[](const SearchDest& sd){ return !sd.hostname.empty(); });

if(nameServers.empty() && !hasHostnames)
return;

tcp_loop.call([this]() {
// start connections to name servers
tcp_loop.call([this, hasHostnames]() {
for(auto& ns : nameServers) {
const auto& serv = ns.first;
ns.second = Connection::build(shared_from_this(), serv);
ns.second->nameserver = true;
log_debug_printf(io, "Connecting to nameserver %s\n", ns.second->peerName.c_str());
ns.conn = Connection::build(shared_from_this(), ns.addr);
ns.conn->nameserver = true;
log_debug_printf(io, "Connecting to nameserver %s%s%s\n", ns.conn->peerName.c_str(),
(ns.hostname.empty()?"":" hostname="), ns.hostname.c_str());
}

if(event_add(nsChecker.get(), &tcpNSCheckInterval))
log_err_printf(setup, "Error enabling TCP search reconnect timer\n%s", "");
if(!nameServers.empty()) {
if(event_add(nsChecker.get(), &tcpNSCheckInterval))
log_err_printf(setup, "Error enabling TCP search reconnect timer\n%s", "");
}

if(event_add(dnsRecheckTimer.get(), &dnsRecheckInterval))
log_err_printf(setup, "Error enabling DNS recheck timer\n%s", "");
});
}

Expand All @@ -705,6 +741,7 @@ void ContextImpl::close()
(void)event_del(searchRx6.get());
(void)event_del(beaconCleaner.get());
(void)event_del(cacheCleaner.get());
(void)event_del(dnsRecheckTimer.get());

auto conns(std::move(connByAddr));
// explicitly break ref. loop of channel cache
Expand Down Expand Up @@ -1224,7 +1261,7 @@ void ContextImpl::tickSearch(SearchKind kind, bool poked)
pport[0] = pport[1] = 0;

for(auto& pair : nameServers) {
auto& serv = pair.second;
auto& serv = pair.conn;

if(!serv->ready || !serv->connection())
continue;
Expand Down Expand Up @@ -1323,12 +1360,18 @@ void ContextImpl::tickBeaconCleanS(evutil_socket_t fd, short evt, void *raw)
void ContextImpl::onNSCheck()
{
for(auto& ns : nameServers) {
if(ns.second && ns.second->state != ConnBase::Disconnected) // hold-off, connecting, or connected
if(!ns.hostname.empty())
continue; // hostname entries owned by onDNSRecheck

if(ns.conn && ns.conn->state != ConnBase::Disconnected)
continue;

ns.second = Connection::build(shared_from_this(), ns.first);
ns.second->nameserver = true;
log_debug_printf(io, "Reconnecting nameserver %s\n", ns.second->peerName.c_str());
// drop old conn first so its dtor's connByAddr.erase(peerAddr) runs
// before build() inserts the fresh entry (same addr -> would erase it)
ns.conn.reset();
ns.conn = Connection::build(shared_from_this(), ns.addr);
ns.conn->nameserver = true;
log_debug_printf(io, "Reconnecting nameserver %s\n", ns.conn->peerName.c_str());
}
}

Expand All @@ -1341,6 +1384,75 @@ void ContextImpl::onNSCheckS(evutil_socket_t fd, short evt, void *raw)
}
}

void ContextImpl::onDNSRecheck()
{
for(auto& sd : searchDest) {
if(sd.hostname.empty())
continue;

SockAddr resolved;
try {
resolved.setAddress(sd.hostname.c_str(), sd.dest.addr.port());
} catch(std::exception& e) {
log_warn_printf(io, "DNS resolution failed for search dest '%s': %s\n",
sd.hostname.c_str(), e.what());
continue;
}

if(resolved != sd.dest.addr) {
log_info_printf(io, "Search dest %s re-resolved: %s -> %s\n",
sd.hostname.c_str(),
sd.dest.addr.tostring().c_str(),
resolved.tostring().c_str());
sd.dest.addr = resolved;
}
}

for(auto& ns : nameServers) {
if(ns.hostname.empty())
continue;

SockAddr resolved;
try {
resolved.setAddress(ns.hostname.c_str(), ns.addr.port());
} catch(std::exception& e) {
log_warn_printf(io, "DNS resolution failed for nameserver '%s': %s\n",
ns.hostname.c_str(), e.what());
continue;
}

bool ipChanged = (resolved != ns.addr);
bool connDown = (!ns.conn || ns.conn->state == ConnBase::Disconnected);

if(ipChanged) {
log_info_printf(io, "Nameserver %s re-resolved: %s -> %s\n",
ns.hostname.c_str(), ns.addr.tostring().c_str(),
resolved.tostring().c_str());
ns.addr = resolved;
}

if(ipChanged || connDown) {
// drop old conn first so its dtor's connByAddr.erase(peerAddr) runs
// before build() inserts the fresh entry
ns.conn.reset();
ns.conn = Connection::build(shared_from_this(), ns.addr);
ns.conn->nameserver = true;
log_debug_printf(io, "Reconnecting nameserver %s (%s)%s\n",
ns.conn->peerName.c_str(), ns.hostname.c_str(),
ipChanged ? " after DNS change" : "");
}
}
}

void ContextImpl::onDNSRecheckS(evutil_socket_t fd, short evt, void *raw)
{
try {
static_cast<ContextImpl*>(raw)->onDNSRecheck();
}catch(std::exception& e){
log_exc_printf(io, "Unhandled error in DNS recheck timer callback: %s\n", e.what());
}
}

void ContextImpl::cacheClean(const std::string& name, Context::cacheAction action)
{
auto next(chanByName.begin()),
Expand Down
17 changes: 14 additions & 3 deletions src/clientimpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,7 @@ struct Channel {

// channel created with .server() to bypass normal search process
SockAddr forcedServer;
std::string forcedServerHostname;

// when state==Searching, number of repetitions
size_t nSearch = 0u;
Expand Down Expand Up @@ -275,10 +276,12 @@ struct ContextImpl : public std::enable_shared_from_this<ContextImpl>

// search destination address and whether to set the unicast flag
struct SearchDest {
const SockEndpoint dest;
SockEndpoint dest;
const bool isucast;
bool lastSuccess = true;
SearchDest(SockEndpoint dest, bool isu) :dest(dest), isucast(isu) {}
std::string hostname;
SearchDest(SockEndpoint dest, bool isu, std::string hostname = {})
:dest(dest), isucast(isu), hostname(std::move(hostname)) {}
};

std::vector<SearchDest> searchDest;
Expand All @@ -299,7 +302,12 @@ struct ContextImpl : public std::enable_shared_from_this<ContextImpl>

std::map<SockAddr, std::weak_ptr<Connection>> connByAddr;

std::vector<std::pair<SockAddr, std::shared_ptr<Connection>>> nameServers;
struct NameServerEntry {
SockAddr addr;
std::shared_ptr<Connection> conn;
std::string hostname;
};
std::vector<NameServerEntry> nameServers;

evbase tcp_loop;
const evevent searchRx4, searchRx6;
Expand All @@ -315,6 +323,7 @@ struct ContextImpl : public std::enable_shared_from_this<ContextImpl>
const evevent beaconCleaner;
const evevent cacheCleaner;
const evevent nsChecker;
const evevent dnsRecheckTimer;

INST_COUNTER(ClientContextImpl);

Expand Down Expand Up @@ -345,6 +354,8 @@ struct ContextImpl : public std::enable_shared_from_this<ContextImpl>
static void cacheCleanS(evutil_socket_t fd, short evt, void *raw);
void onNSCheck();
static void onNSCheckS(evutil_socket_t fd, short evt, void *raw);
void onDNSRecheck();
static void onDNSRecheckS(evutil_socket_t fd, short evt, void *raw);
};

struct Context::Pvt {
Expand Down
31 changes: 22 additions & 9 deletions src/config.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,8 @@ namespace {
constexpr double tmoScale = 4.0/3.0; // 40 second idle timeout / 30 configured

void split_addr_into(const char* name, std::vector<std::string>& out, const std::string& inp,
uint16_t defaultPort, bool required=false)
uint16_t defaultPort, bool required=false,
std::map<std::string, std::string>* hostnameMap=nullptr)
{
size_t pos=0u;

Expand All @@ -166,7 +167,16 @@ void split_addr_into(const char* name, std::vector<std::string>& out, const std:
SockEndpoint ep(temp);
if(ep.addr.port()==0)
ep.addr.setPort(defaultPort);
out.push_back(SB()<<ep);
auto resolved = (SB()<<ep).str();
out.push_back(resolved);

log_warn_printf(config, "TRACE split_addr_into: token='%s' resolved='%s' isHostname=%d hostnameMap=%p\n",
temp.c_str(), resolved.c_str(), (int)isHostname(temp), (void*)hostnameMap);
if(hostnameMap && isHostname(temp)) {
(*hostnameMap)[resolved] = temp;
log_warn_printf(config, "TRACE hostnameMap[%s] = %s (size now %zu)\n",
resolved.c_str(), temp.c_str(), hostnameMap->size());
}

} catch(std::exception& e){
if(required)
Expand Down Expand Up @@ -572,17 +582,20 @@ void _fromDefs(Config& self, const std::map<std::string, std::string>& defs, boo
log_warn_printf(clientsetup, "%s invalid integer : %s", pickone.name.c_str(), e.what());
}
}
if(self.tcp_port==0u && !self.nameServers.empty()) {
log_warn_printf(clientsetup, "ignoring EPICS_PVA_SERVER_PORT=%d\n", 0);
self.tcp_port = 5075;
}

if(pickone({"EPICS_PVA_ADDR_LIST"})) {
split_addr_into(pickone.name.c_str(), self.addressList, pickone.val, self.udp_port);
split_addr_into(pickone.name.c_str(), self.addressList, pickone.val, self.udp_port,
false, &self.addressHostnames);
}

if(pickone({"EPICS_PVA_NAME_SERVERS"})) {
split_addr_into(pickone.name.c_str(), self.nameServers, pickone.val, self.tcp_port);
auto nameServersName(pickone.name);
auto nameServersVal(pickone.val);
if(self.tcp_port==0u) {
log_warn_printf(clientsetup, "ignoring EPICS_PVA_SERVER_PORT=%d\n", 0);
self.tcp_port = 5075;
}
split_addr_into(nameServersName.c_str(), self.nameServers, nameServersVal, self.tcp_port,
false, &self.nameServerHostnames);
}

if(pickone({"EPICS_PVA_AUTO_ADDR_LIST"})) {
Expand Down
10 changes: 10 additions & 0 deletions src/pvxs/client.h
Original file line number Diff line number Diff line change
Expand Up @@ -1026,6 +1026,16 @@ struct PVXS_API Config {
//! @since 0.2.0
std::vector<std::string> nameServers;

//! Maps resolved IP string -> original hostname for entries in addressList.
//! Populated automatically when addressList entries are hostnames.
//! @since NEXT
std::map<std::string, std::string> addressHostnames;

//! Maps resolved IP string -> original hostname for entries in nameServers.
//! Populated automatically when nameServers entries are hostnames.
//! @since NEXT
std::map<std::string, std::string> nameServerHostnames;

//! UDP port to bind. Default is 5076. May be zero, cf. Server::config() to find allocated port.
unsigned short udp_port = 5076;
//! Default TCP port for name servers
Expand Down
Loading
Loading