Skip to content
Merged
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
4 changes: 4 additions & 0 deletions connectd/connectd.c
Original file line number Diff line number Diff line change
Expand Up @@ -1695,6 +1695,7 @@ static void connect_init(struct daemon *daemon, const u8 *msg)
struct wireaddr *announceable;
char *tor_password;
bool dev_disconnect, dev_throttle_gossip, dev_limit_connections_inflight;
u32 dev_gossip_cpu_budget;
char *errstr;

/* Fields which require allocation are allocated off daemon */
Expand All @@ -1717,6 +1718,7 @@ static void connect_init(struct daemon *daemon, const u8 *msg)
&daemon->dev_no_ping_timer,
&daemon->dev_handshake_no_reply,
&dev_throttle_gossip,
&dev_gossip_cpu_budget,
&daemon->dev_no_reconnect,
&daemon->dev_fast_reconnect,
&dev_limit_connections_inflight,
Expand Down Expand Up @@ -1791,6 +1793,8 @@ static void connect_init(struct daemon *daemon, const u8 *msg)
daemon->incoming_stream_limit = 1000;
daemon->cpu_budget_usec_limit = 1500;
}
if (dev_gossip_cpu_budget)
daemon->cpu_budget_usec_limit = dev_gossip_cpu_budget;

if (dev_limit_connections_inflight)
daemon->max_connect_in_flight = 1;
Expand Down
1 change: 1 addition & 0 deletions connectd/connectd_wire.csv
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ msgdata,connectd_init,dev_no_ping_timer,bool,
# Allow incoming connections, but don't talk.
msgdata,connectd_init,dev_noreply,bool,
msgdata,connectd_init,dev_throttle_gossip,bool,
msgdata,connectd_init,dev_gossip_cpu_budget,u32,
msgdata,connectd_init,dev_no_reconnect,bool,
msgdata,connectd_init,dev_fast_reconnect,bool,
msgdata,connectd_init,dev_limit_connections_inflight,bool,
Expand Down
56 changes: 38 additions & 18 deletions connectd/multiplex.c
Original file line number Diff line number Diff line change
Expand Up @@ -637,6 +637,15 @@ static void maybe_reset_usage_window(struct peer *peer)
peer->gs.cpu_usec_this_second = 0;
}

/* Each peer gets its "fair share" of our query-answering CPU. No floor:
* that would let enough peers claim more than the whole budget between
* them. A small share only slows a peer's queries, nothing else. */
static u64 peer_cpu_budget(const struct peer *peer)
{
return peer->daemon->cpu_budget_usec_limit
/ peer_htable_count(peer->daemon->peers);
}

/* usage/limit, in usec: e.g. usage == 3*limit means "3 seconds worth
* of quota burned in one go." */
static u64 usec_over_budget(u64 usage, u64 limit)
Expand All @@ -662,12 +671,15 @@ static u64 maybe_throttle_usec(struct peer *peer, bool *warned, const char *dire
if (usage1 <= limit1 && usage2 <= limit2)
return 0;

status_unusual_once(warned,
CI_UNEXPECTED
"Throttling %s peer %s: too much %s",
direction,
fmt_node_id(tmpctx, &peer->id),
usage1 > limit1 ? "traffic" : "CPU");
/* A peer hitting its budget is expected (that's what it's for), so
* don't make noise about it. */
if (!*warned) {
status_debug("Throttling %s peer %s: too much %s",
direction,
fmt_node_id(tmpctx, &peer->id),
usage1 > limit1 ? "traffic" : "CPU");
*warned = true;
}

need = usec_over_budget(usage1, limit1);
need2 = usec_over_budget(usage2, limit2);
Expand All @@ -691,12 +703,11 @@ static const u8 *maybe_gossip_msg(const tal_t *ctx, struct peer *peer)
u32 timestamp;
const u8 **msgs;
u64 cpu_budget, wait_usec;
bool answering;

maybe_reset_usage_window(peer);

/* Each peer gets its "fair share" of our query-answering CPU */
cpu_budget = peer->daemon->cpu_budget_usec_limit
/ peer_htable_count(peer->daemon->peers);
cpu_budget = peer_cpu_budget(peer);

wait_usec = maybe_throttle_usec(peer, &peer->gs.throttle_warned, "outgoing",
peer->gs.bytes_this_second, peer->daemon->gossip_stream_limit,
Expand All @@ -717,11 +728,16 @@ static const u8 *maybe_gossip_msg(const tal_t *ctx, struct peer *peer)

/* This can return more than one: it's the expensive part (gossmap
* walks, checksum/timestamp lookups), so it's what we charge for
* cpu_usec_this_second above. */
* cpu_usec_this_second above. But only while we're answering a
* query: we're called every time the output queue drains, and the
* rest of the time it's a no-op. */
answering = peer->scid_queries || peer->range_scids;
query_start = time_mono();
msgs = maybe_create_query_responses(tmpctx, peer, gossmap);
peer->gs.cpu_usec_this_second
+= time_to_usec(timemono_between(time_mono(), query_start));
if (answering)
peer->gs.cpu_usec_this_second
+= time_to_usec(timemono_between(time_mono(),
query_start));
if (tal_count(msgs) > 0) {
/* We return the first one for immediate sending, and queue
* others for future. We add all the lengths now though! */
Expand Down Expand Up @@ -1526,10 +1542,16 @@ static struct io_plan *read_body_from_peer_done(struct io_conn *peer_conn,

/* If we swallow this, just try again. */
handled_start = time_mono();
/* Count the CPU time for local messages (esp. gossip queries) */
if (handle_message_locally(peer, decrypted)) {
peer->gs.cpu_usec_this_second
+= time_to_usec(timemono_between(time_mono(), handled_start));
/* Only gossip queries count against the CPU budget: answering
* them walks the gossmap for the peer. The rest is cheap,
* ratelimited on its own (onion messages), or handed to
* another daemon (gossip, custom messages). */
if (type == WIRE_QUERY_CHANNEL_RANGE
|| type == WIRE_QUERY_SHORT_CHANNEL_IDS)
peer->gs.cpu_usec_this_second
+= time_to_usec(timemono_between(time_mono(),
handled_start));
/* Make sure to update peer->peer_in_lastmsg so we blame correct msg! */
goto out;
}
Expand Down Expand Up @@ -1642,9 +1664,7 @@ static struct io_plan *read_hdr_from_peer(struct io_conn *peer_conn,

maybe_reset_usage_window(peer);

/* Each peer gets its "fair share" of our local-message CPU */
cpu_budget = peer->daemon->cpu_budget_usec_limit
/ peer_htable_count(peer->daemon->peers);
cpu_budget = peer_cpu_budget(peer);

wait_usec = maybe_throttle_usec(peer, &peer->throttle_warned, "incoming",
peer->gs.bytes_rcvd_this_second, peer->daemon->incoming_stream_limit,
Expand Down
1 change: 1 addition & 0 deletions lightningd/connect_control.c
Original file line number Diff line number Diff line change
Expand Up @@ -760,6 +760,7 @@ int connectd_init(struct lightningd *ld)
ld->dev_no_ping_timer,
ld->dev_handshake_no_reply,
ld->dev_throttle_gossip,
ld->dev_gossip_cpu_budget,
!ld->reconnect,
ld->dev_fast_reconnect,
ld->dev_limit_connections_inflight,
Expand Down
1 change: 1 addition & 0 deletions lightningd/lightningd.c
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,7 @@ static struct lightningd *new_lightningd(const tal_t *ctx)
ld->dev_fast_gossip = false;
ld->dev_fast_gossip_prune = false;
ld->dev_throttle_gossip = false;
ld->dev_gossip_cpu_budget = 0;
ld->dev_suppress_gossip = false;
ld->dev_fast_reconnect = false;
ld->dev_reject_closing_fee = false;
Expand Down
1 change: 1 addition & 0 deletions lightningd/lightningd.h
Original file line number Diff line number Diff line change
Expand Up @@ -314,6 +314,7 @@ struct lightningd {
bool dev_fast_gossip;
bool dev_fast_gossip_prune;
bool dev_throttle_gossip;
u32 dev_gossip_cpu_budget;
bool dev_suppress_gossip;

/* How long to aim for low-priority commitment closes */
Expand Down
4 changes: 4 additions & 0 deletions lightningd/options.c
Original file line number Diff line number Diff line change
Expand Up @@ -928,6 +928,10 @@ static void dev_register_opts(struct lightningd *ld)
opt_set_bool,
&ld->dev_throttle_gossip,
"Throttle gossip right down, for testing");
clnopt_witharg("--dev-gossip-cpu-budget", OPT_DEV|OPT_SHOWINT,
opt_set_u32, opt_show_u32,
&ld->dev_gossip_cpu_budget,
"Total CPU usec per second for answering gossip queries (0: default)");
clnopt_noarg("--dev-limit-connections-inflight", OPT_DEV,
opt_set_bool,
&ld->dev_limit_connections_inflight,
Expand Down
26 changes: 26 additions & 0 deletions tests/test_gossip.py
Original file line number Diff line number Diff line change
Expand Up @@ -2330,6 +2330,32 @@ def test_gossip_query_channel_range_cpu_throttle(node_factory, chainparams):
l2.daemon.wait_for_log(r'Throttling outgoing peer .*: too much CPU')


def test_gossip_cpu_throttle_only_queries(node_factory):
Comment thread
nGoline marked this conversation as resolved.
"""Only answering gossip queries counts against the CPU budget: a
peer sending us ordinary messages (here pings) must not get its reads
throttled, even with a tiny CPU budget."""
# A tiny CPU budget, but the normal traffic limits, so only CPU can
# trip the throttle.
l1 = node_factory.get_node(options={'dev-gossip-cpu-budget': 10})

# ping with num_pong_bytes 0, no padding.
ping = '0012' + '0000' + '0000'
out = subprocess.run(['devtools/gossipwith',
'--no-gossip',
'--hex',
'--network={}'.format(TEST_NETWORK),
'--filter=19',
'--max-messages=200',
'--timeout-after=30',
'{}@localhost:{}'.format(l1.info['id'], l1.port)]
+ [ping] * 200,
timeout=TIMEOUT, stdout=subprocess.PIPE).stdout.split()

# Every ping was answered, so all 200 really were read.
assert len(out) == 200
assert not l1.daemon.is_in_log('Throttling incoming peer')


def test_generate_gossip_store(node_factory):
l1 = node_factory.get_node(start=False)
chans = [GenChannel(0, 1),
Expand Down
Loading