diff --git a/connectd/connectd.c b/connectd/connectd.c index 98eb1f13b1d4..f1e08db71aeb 100644 --- a/connectd/connectd.c +++ b/connectd/connectd.c @@ -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 */ @@ -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, @@ -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; diff --git a/connectd/connectd_wire.csv b/connectd/connectd_wire.csv index 03dccc94718c..9b9670704735 100644 --- a/connectd/connectd_wire.csv +++ b/connectd/connectd_wire.csv @@ -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, diff --git a/connectd/multiplex.c b/connectd/multiplex.c index 16b8de428863..b5fec97b0d90 100644 --- a/connectd/multiplex.c +++ b/connectd/multiplex.c @@ -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) @@ -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); @@ -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, @@ -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! */ @@ -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; } @@ -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, diff --git a/lightningd/connect_control.c b/lightningd/connect_control.c index 54b3d46c1180..99e417118a1b 100644 --- a/lightningd/connect_control.c +++ b/lightningd/connect_control.c @@ -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, diff --git a/lightningd/lightningd.c b/lightningd/lightningd.c index 6dbc2a0473b1..52d9bb92c65e 100644 --- a/lightningd/lightningd.c +++ b/lightningd/lightningd.c @@ -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; diff --git a/lightningd/lightningd.h b/lightningd/lightningd.h index 87a2851ceff4..75759e21efa5 100644 --- a/lightningd/lightningd.h +++ b/lightningd/lightningd.h @@ -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 */ diff --git a/lightningd/options.c b/lightningd/options.c index 6e3979b46850..67fb7da4044a 100644 --- a/lightningd/options.c +++ b/lightningd/options.c @@ -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, diff --git a/tests/test_gossip.py b/tests/test_gossip.py index 7391ba76dbb0..0fefb264b880 100644 --- a/tests/test_gossip.py +++ b/tests/test_gossip.py @@ -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): + """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),