Skip to content

Commit bf652e7

Browse files
committed
Merge branch 'cboulay/resolve_over_tcp' into cboulay/meta_resolve_fixes
# Conflicts: # testing/CMakeLists.txt
2 parents 1d2b724 + 2a7ad05 commit bf652e7

10 files changed

Lines changed: 335 additions & 4 deletions

‎cmake/SourceFiles.cmake‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,8 @@ set(lslsources
2929
src/portable_archive/portable_oarchive.hpp
3030
src/resolver_impl.cpp
3131
src/resolver_impl.h
32+
src/resolve_attempt_tcp.cpp
33+
src/resolve_attempt_tcp.h
3234
src/resolve_attempt_udp.cpp
3335
src/resolve_attempt_udp.h
3436
src/sample.cpp

‎src/api_config.cpp‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -307,6 +307,7 @@ void api_config::load(INI &pt) {
307307
// read the [lab] settings
308308
known_peers_ = parse_set(pt.get("lab.KnownPeers", "{}"));
309309
session_id_ = pt.get("lab.SessionID", "default");
310+
resolve_over_tcp_ = pt.get("lab.ResolveOverTCP", false);
310311

311312
// read the [tuning] settings
312313
use_protocol_version_ = std::min(

‎src/api_config.h‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -188,6 +188,17 @@ class api_config {
188188
*/
189189
const std::vector<std::string> &known_peers() const { return known_peers_; }
190190

191+
/**
192+
* @brief Whether to additionally resolve streams by probing TCP data ports directly.
193+
*
194+
* When enabled, a resolve also connects to every port in the BasePort..BasePort+PortRange
195+
* range on loopback and on each KnownPeer and requests stream info over TCP. Because TCP is a
196+
* symmetric, connection-oriented protocol, this is robust against the stateful-firewall issues
197+
* that can block UDP discovery, but it is VERY slow (one connection attempt per port per host,
198+
* and closed/filtered remote ports cost a full connect timeout each). Disabled by default.
199+
*/
200+
bool resolve_over_tcp() const { return resolve_over_tcp_; }
201+
191202
// === tuning parameters ===
192203

193204
/// The network protocol version to use.
@@ -298,6 +309,7 @@ class api_config {
298309
std::string listen_address_;
299310
std::vector<std::string> known_peers_;
300311
std::string session_id_;
312+
bool resolve_over_tcp_;
301313
// tuning parameters
302314
int use_protocol_version_;
303315
double watchdog_time_threshold_;

‎src/resolve_attempt_tcp.cpp‎

Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
1+
#include "resolve_attempt_tcp.h"
2+
#include "common.h"
3+
#include "resolver_impl.h"
4+
#include "stream_info_impl.h"
5+
#include <algorithm>
6+
#include <asio/buffer.hpp>
7+
#include <asio/connect.hpp>
8+
#include <asio/read.hpp>
9+
#include <asio/write.hpp>
10+
#include <exception>
11+
#include <loguru.hpp>
12+
#include <mutex>
13+
#include <utility>
14+
15+
using namespace lsl;
16+
using err_t = const asio::error_code &;
17+
18+
/// Maximum number of TCP probes in flight at once (one socket each).
19+
static const std::size_t MAX_CONCURRENT_PROBES = 16;
20+
21+
resolve_attempt_tcp::resolve_attempt_tcp(asio::io_context &io,
22+
const std::vector<tcp::endpoint> &targets, const std::string &query, resolver_impl &resolver,
23+
double cancel_after)
24+
: io_(io), resolver_(resolver), cancel_after_(cancel_after), cancelled_(false),
25+
targets_(targets), next_target_(0), active_workers_(0), cancel_timer_(io) {
26+
// Precompute the query message. Unlike the UDP query there is no return port or query id:
27+
// the reply comes back on this very connection and the TCP shortinfo responder replies with
28+
// the bare shortinfo message (no query-id prefix).
29+
query_msg_ = "LSL:shortinfo\r\n" + query + "\r\n";
30+
// register ourselves as a candidate for cancellation
31+
register_at(&resolver);
32+
}
33+
34+
resolve_attempt_tcp::~resolve_attempt_tcp() { unregister_from_all(); }
35+
36+
void resolve_attempt_tcp::begin() {
37+
// arm the cancel timer
38+
if (cancel_after_ != FOREVER) {
39+
cancel_timer_.expires_after(timeout_sec(cancel_after_));
40+
cancel_timer_.async_wait([shared_this = shared_from_this(), this](err_t err) {
41+
if (!err) do_cancel();
42+
});
43+
}
44+
// launch the worker pool
45+
active_workers_ = std::min(MAX_CONCURRENT_PROBES, targets_.size());
46+
worker_sockets_.resize(active_workers_);
47+
std::size_t workers = active_workers_;
48+
for (std::size_t w = 0; w < workers; w++) probe_next(w);
49+
}
50+
51+
void resolve_attempt_tcp::cancel() {
52+
post(io_, [shared_this = shared_from_this()]() { shared_this->do_cancel(); });
53+
}
54+
55+
void resolve_attempt_tcp::probe_next(std::size_t worker) {
56+
if (cancelled_) return;
57+
if (next_target_ >= targets_.size()) {
58+
// this worker is out of endpoints; once all workers are done, drop the cancel timer so
59+
// the attempt can finish (and the resolver can schedule a fresh burst)
60+
if (--active_workers_ == 0) cancel_timer_.cancel();
61+
return;
62+
}
63+
const tcp::endpoint ep = targets_[next_target_++];
64+
auto sock = std::make_shared<tcp_socket>(io_);
65+
worker_sockets_[worker] = sock;
66+
67+
auto self = shared_from_this();
68+
sock->async_connect(ep, [self, this, worker, ep, sock](err_t err) {
69+
if (cancelled_) return;
70+
if (err) {
71+
// closed / filtered / unreachable port: move on to the next target
72+
probe_next(worker);
73+
return;
74+
}
75+
// connected: send the shortinfo query
76+
asio::async_write(*sock, asio::buffer(query_msg_),
77+
[self, this, worker, ep, sock](err_t err, std::size_t /*unused*/) {
78+
if (cancelled_) return;
79+
if (err) {
80+
probe_next(worker);
81+
return;
82+
}
83+
// read the reply until the outlet closes the connection (EOF)
84+
auto reply = std::make_shared<std::string>();
85+
asio::async_read(*sock, asio::dynamic_buffer(*reply),
86+
[self, this, worker, ep, sock, reply](err_t err, std::size_t /*unused*/) {
87+
if (cancelled_) return;
88+
if (!err || err == asio::error::eof) handle_reply(*reply, ep);
89+
probe_next(worker);
90+
});
91+
});
92+
});
93+
}
94+
95+
void resolve_attempt_tcp::handle_reply(const std::string &reply, const tcp::endpoint &ep) {
96+
if (reply.empty()) return; // no match: the responder closed without replying
97+
try {
98+
stream_info_impl info;
99+
info.from_shortinfo_message(reply);
100+
// The shortinfo XML carries the ports but not the host address; fill it in from the
101+
// endpoint we connected to (mirrors the UDP resolver using the reply's source address).
102+
const std::string addr = ep.address().to_string();
103+
std::string uid = info.uid();
104+
{
105+
std::lock_guard<std::mutex> lock(resolver_.results_mut_);
106+
auto it = resolver_.results_.find(uid);
107+
if (it == resolver_.results_.end())
108+
it = resolver_.results_.emplace(uid, std::make_pair(info, lsl_clock())).first;
109+
else
110+
it->second.second = lsl_clock();
111+
auto &stored_info = it->second.first;
112+
if (ep.address().is_v4()) {
113+
if (stored_info.v4address().empty()) stored_info.v4address(addr);
114+
} else {
115+
if (stored_info.v6address().empty()) stored_info.v6address(addr);
116+
}
117+
}
118+
// prepone the next cancellation check so a hit cancels the remaining (slow) probes
119+
if (resolver_.check_cancellation_criteria()) resolver_.cancel_ongoing_resolve();
120+
} catch (std::exception &e) {
121+
LOG_F(WARNING, "resolve_attempt_tcp: could not parse a reply from %s: %s",
122+
ep.address().to_string().c_str(), e.what());
123+
}
124+
}
125+
126+
void resolve_attempt_tcp::do_cancel() {
127+
try {
128+
cancelled_ = true;
129+
for (auto &sock : worker_sockets_)
130+
if (sock && sock->is_open()) sock->close();
131+
cancel_timer_.cancel();
132+
} catch (std::exception &e) {
133+
LOG_F(WARNING, "Unexpected error while trying to cancel a resolve_attempt_tcp: %s",
134+
e.what());
135+
}
136+
}

‎src/resolve_attempt_tcp.h‎

Lines changed: 89 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,89 @@
1+
#ifndef RESOLVE_ATTEMPT_TCP_H
2+
#define RESOLVE_ATTEMPT_TCP_H
3+
4+
#include "cancellation.h"
5+
#include "socket_utils.h"
6+
#include <asio/io_context.hpp>
7+
#include <asio/ip/tcp.hpp>
8+
#include <asio/steady_timer.hpp>
9+
#include <cstddef>
10+
#include <memory>
11+
#include <string>
12+
#include <vector>
13+
14+
using asio::ip::tcp;
15+
using err_t = const asio::error_code &;
16+
17+
namespace lsl {
18+
class resolver_impl;
19+
20+
using steady_timer = asio::basic_waitable_timer<asio::chrono::steady_clock,
21+
asio::wait_traits<asio::chrono::steady_clock>, asio::io_context::executor_type>;
22+
23+
/**
24+
* An asynchronous resolve attempt that probes a set of TCP endpoints directly.
25+
*
26+
* For each endpoint it opens a TCP connection, sends an `LSL:shortinfo` query, and parses the
27+
* outlet's reply into a stream_info that is stored in the shared resolver results. Unlike the UDP
28+
* resolver this requires one connection per endpoint, so probes run through a small pool of
29+
* concurrent workers. Used as a firewall-robust (but slow) discovery fallback, gated behind the
30+
* `lab.ResolveOverTCP` config option.
31+
*/
32+
class resolve_attempt_tcp final : public cancellable_obj,
33+
public std::enable_shared_from_this<resolve_attempt_tcp> {
34+
public:
35+
/**
36+
* Instantiate and set up a new TCP resolve attempt.
37+
*
38+
* @param io The io_context that will run the async operations.
39+
* @param targets The TCP endpoints to probe (host x port-range).
40+
* @param query The query string the outlet must match to reply.
41+
* @param resolver The resolver whose results container is populated.
42+
* @param cancel_after Time after which the attempt is automatically cancelled.
43+
*/
44+
resolve_attempt_tcp(asio::io_context &io, const std::vector<tcp::endpoint> &targets,
45+
const std::string &query, resolver_impl &resolver, double cancel_after = 5.0);
46+
47+
/// Destructor.
48+
~resolve_attempt_tcp() final;
49+
50+
/// Start probing asynchronously.
51+
void begin();
52+
53+
/// Cancel operations asynchronously and destructively.
54+
void cancel() override;
55+
56+
private:
57+
/// Probe the next not-yet-taken target endpoint using the given worker's socket slot.
58+
void probe_next(std::size_t worker);
59+
60+
/// Parse a reply received from a given endpoint into the resolver results.
61+
void handle_reply(const std::string &reply, const tcp::endpoint &ep);
62+
63+
/// Cancel the outstanding operations.
64+
void do_cancel();
65+
66+
/// the IO service that executes our actions
67+
asio::io_context &io_;
68+
/// the resolver associated with this attempt
69+
resolver_impl &resolver_;
70+
/// the timeout for giving up
71+
double cancel_after_;
72+
/// whether the operation has been cancelled
73+
bool cancelled_;
74+
/// the endpoints to probe
75+
std::vector<tcp::endpoint> targets_;
76+
/// the message we send ("LSL:shortinfo\r\n<query>\r\n")
77+
std::string query_msg_;
78+
/// index of the next target to hand to a worker
79+
std::size_t next_target_;
80+
/// number of workers that still have endpoints to probe
81+
std::size_t active_workers_;
82+
/// per-worker current socket (kept so a cancel can close in-flight connects)
83+
std::vector<std::shared_ptr<tcp_socket>> worker_sockets_;
84+
/// timer to schedule the cancel action
85+
steady_timer cancel_timer_;
86+
};
87+
} // namespace lsl
88+
89+
#endif

‎src/resolver_impl.cpp‎

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,12 @@
11
#include "resolver_impl.h"
22
#include "api_config.h"
3+
#include "resolve_attempt_tcp.h"
34
#include "resolve_attempt_udp.h"
45
#include "socket_utils.h"
56
#include "stream_info_impl.h"
67
#include <asio/io_context.hpp>
78
#include <asio/ip/basic_resolver.hpp>
9+
#include <asio/ip/tcp.hpp>
810
#include <asio/ip/udp.hpp>
911
#include <algorithm>
1012
#include <exception>
@@ -75,6 +77,29 @@ resolver_impl::resolver_impl()
7577
if (cfg_->allow_ipv4()) {
7678
udp_protocols_.push_back(udp::v4());
7779
}
80+
81+
// build the TCP probe targets (loopback + known peers, each x port range) if enabled
82+
if (cfg_->resolve_over_tcp()) {
83+
uint16_t base = cfg_->base_port();
84+
uint16_t range = cfg_->port_range();
85+
LOG_F(WARNING,
86+
"lab.ResolveOverTCP is enabled: discovery will also TCP-probe ports %u-%u on loopback "
87+
"and on every KnownPeer. This is robust behind firewalls but VERY slow, especially "
88+
"across the internet.",
89+
base, static_cast<unsigned>(base + range - 1));
90+
std::vector<std::string> hosts{"127.0.0.1"};
91+
if (cfg_->allow_ipv6()) hosts.emplace_back("::1");
92+
for (const auto &peer : cfg_->known_peers()) hosts.push_back(peer);
93+
tcp::resolver tcp_resolver(*io_);
94+
for (const auto &host : hosts) {
95+
try {
96+
for (const auto &res : tcp_resolver.resolve(host, std::to_string(base))) {
97+
for (uint16_t p = base; p < base + range; p++)
98+
tcp_endpoints_.emplace_back(res.endpoint().address(), p);
99+
}
100+
} catch (std::exception &) {}
101+
}
102+
}
78103
}
79104

80105
void check_query(const std::string &query) {
@@ -220,6 +245,17 @@ void resolver_impl::next_resolve_wave() {
220245
// delay the next multicast wave
221246
wave_timer_timeout += cfg_->unicast_min_rtt();
222247
}
248+
// fire a TCP probe burst if enabled and a previous one isn't still in flight
249+
if (cfg_->resolve_over_tcp() && !tcp_endpoints_.empty() && tcp_attempt_.expired()) {
250+
try {
251+
auto attempt = std::make_shared<resolve_attempt_tcp>(
252+
*io_, tcp_endpoints_, query_, *this, cfg_->unicast_max_rtt());
253+
tcp_attempt_ = attempt;
254+
attempt->begin();
255+
} catch (std::exception &e) {
256+
LOG_F(WARNING, "Could not start a TCP resolve attempt: %s", e.what());
257+
}
258+
}
223259
wave_timer_.expires_after(timeout_sec(wave_timer_timeout));
224260
wave_timer_.async_wait([this](err_t err) {
225261
if (err != asio::error::operation_aborted) next_resolve_wave();

‎src/resolver_impl.h‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ using err_t = const asio::error_code &;
2222

2323
namespace lsl {
2424
class api_config;
25+
class resolve_attempt_tcp;
2526

2627
using steady_timer = asio::basic_waitable_timer<asio::chrono::steady_clock, asio::wait_traits<asio::chrono::steady_clock>, asio::io_context::executor_type>;
2728

@@ -132,6 +133,7 @@ class resolver_impl final : public cancellable_registry {
132133

133134
private:
134135
friend class resolve_attempt_udp;
136+
friend class resolve_attempt_tcp;
135137

136138
/// This function starts a new wave of resolves.
137139
void next_resolve_wave();
@@ -158,6 +160,8 @@ class resolver_impl final : public cancellable_registry {
158160
std::vector<udp::endpoint> mcast_endpoints_;
159161
/// the list of per-host UDP endpoints under consideration
160162
std::vector<udp::endpoint> ucast_endpoints_;
163+
/// the list of per-host TCP endpoints to probe when lab.ResolveOverTCP is enabled
164+
std::vector<tcp::endpoint> tcp_endpoints_;
161165

162166
// things related to cancellation
163167
/// if set, no more resolves can be started (destructively cancelled).
@@ -188,6 +192,8 @@ class resolver_impl final : public cancellable_registry {
188192
io_context_p io_;
189193
/// a thread that runs background IO if we are performing a resolve_continuous
190194
std::shared_ptr<std::thread> background_io_;
195+
/// the currently in-flight TCP probe attempt (if any); used to avoid overlapping bursts
196+
std::weak_ptr<resolve_attempt_tcp> tcp_attempt_;
191197
/// the overall timeout for a query
192198
steady_timer resolve_timeout_expired_;
193199
/// a timer that fires when a new wave should be scheduled

‎src/tcp_server.cpp‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -544,10 +544,11 @@ void client_session::handle_read_query_outcome(err_t err) {
544544
auto serv = serv_.lock();
545545
if (!serv) return;
546546
if (serv->info_->matches_query(query)) {
547-
// matches: reply (otherwise just close the stream)
547+
// matches: reply (otherwise just close the stream). Keep the session (and thus the
548+
// socket) alive until the shortinfo has been sent; the session is then destroyed,
549+
// closing the socket so the client sees EOF and knows the reply is complete.
548550
async_write(sock_, asio::buffer(serv->shortinfo_msg_),
549-
[serv](err_t /*unused*/, std::size_t /*unused*/) {
550-
/* keep the tcp_server alive until the shortinfo is sent completely*/
551+
[shared_this = shared_from_this(), serv](err_t /*unused*/, std::size_t /*unused*/) {
551552
});
552553
} else {
553554
DLOG_F(INFO, "%p got a shortinfo query response for the wrong query", this);

0 commit comments

Comments
 (0)