From 71201d574d42bbf83a5d773cc92552d614b4dafc Mon Sep 17 00:00:00 2001 From: Jasmine Cha Date: Thu, 30 Jul 2026 06:07:00 +0000 Subject: [PATCH] mctpd: Support MCTP Discovery Notify command - Defer EID assignment/discovery to event loop to keep D-Bus unblocked (ack immediately). - Tear down existing peer state to allow clean re-enumeration on re-registration. - Add per-link 32-entry LRU rate limiter (5 req/sec per physical address) to prevent memory exhaustion attacks. - Add unit tests validating Discovery Notify workflows. Assisted-by: Antigravity:Gemini-Next Signed-off-by: Jasmine Cha --- src/mctpd.c | 259 ++++++++++++++++++++++++++++++++++- tests/mctp_test_utils.py | 10 ++ tests/test_mctpd.py | 180 ++++++++++++++++++++++++ tests/test_mctpd_endpoint.py | 15 +- 4 files changed, 459 insertions(+), 5 deletions(-) diff --git a/src/mctpd.c b/src/mctpd.c index a1d27282..32472d63 100644 --- a/src/mctpd.c +++ b/src/mctpd.c @@ -140,6 +140,20 @@ enum discovery_state { DISCOVERY_UNDISCOVERED, }; +struct discovery_notify_ctx { + struct ctx *ctx; + dest_phys phys; + struct discovery_notify_ctx *next; +}; + +#define MAX_LINK_RATE_LIMIT_ENTRIES 32 + +struct phys_rate_limit_entry { + dest_phys phys; + uint64_t last_discovery_time_us; + uint32_t discovery_count; +}; + struct link { enum discovery_state discovered; bool published; @@ -153,6 +167,9 @@ struct link { sd_bus_slot *slot_busowner; sd_event_source *role_defer; + struct phys_rate_limit_entry rate_limits[MAX_LINK_RATE_LIMIT_ENTRIES]; + size_t num_rate_limit_entries; + struct ctx *ctx; }; @@ -314,6 +331,8 @@ struct ctx { // bus owner/bridge polling interval in usecs for // checking endpoint's accessibility. uint64_t endpoint_poll; + struct discovery_notify_ctx *pending_discoveries; + sd_event_source *discovery_defer; // interface configuration (from config file), to be matched and // applied on new interface events @@ -364,6 +383,10 @@ static int add_peer(struct ctx *ctx, const dest_phys *dest, mctp_eid_t eid, static int add_peer_from_addr(struct ctx *ctx, const struct sockaddr_mctp_ext *addr, struct peer **ret_peer); +static int change_peer_eid(struct peer *peer, mctp_eid_t new_eid); +static int endpoint_assign_eid(struct ctx *ctx, sd_bus_error *berr, + const dest_phys *dest, struct peer **ret_peer, + mctp_eid_t static_eid, bool assign_bridge); static int remove_peer(struct peer *peer); static int remove_bridged_peers(struct peer *bridge); static int query_peer_properties(struct peer *peer); @@ -750,9 +773,15 @@ static int reply_message(struct ctx *ctx, int sd, const void *resp, if (reply_addr.smctp_addr.s_addr == 0 || reply_addr.smctp_addr.s_addr == 0xff) { - bug_warn("reply_message can't take EID %d", - reply_addr.smctp_addr.s_addr); - return -EPROTO; + /* Validate physical address metadata before falling back to physical send */ + if (addr && addr->smctp_ifindex > 0 && + addr->smctp_halen <= sizeof(addr->smctp_haddr)) { + return reply_message_phys(ctx, sd, resp, resp_len, + addr); + } + warnx("reply_message: EID %d specified without valid physical address info", + reply_addr.smctp_addr.s_addr); + return -EINVAL; } len = mctp_ops.mctp.sendto(sd, resp, resp_len, 0, @@ -1329,6 +1358,213 @@ handle_control_endpoint_discovery(struct ctx *ctx, int sd, return reply_message_phys(ctx, sd, resp, sizeof(*resp), addr); } +/* + * Rate limiting is evaluated per physical endpoint (MAC / I2C address) using a bounded + * LRU table on the link to isolate noisy endpoints on shared multi-drop buses (e.g. I2C/SMBus) + * without risking unbounded memory allocation. + */ +static bool link_check_rate_limit(struct link *link_data, const dest_phys *phys) +{ + struct phys_rate_limit_entry *entry = NULL; + struct timespec ts; + uint64_t now_us; + size_t oldest_idx = 0; + uint64_t oldest_time = UINT64_MAX; + + if (clock_gettime(CLOCK_MONOTONIC, &ts) < 0) + return true; + + now_us = (uint64_t)ts.tv_sec * 1000000ULL + ts.tv_nsec / 1000ULL; + + for (size_t i = 0; i < link_data->num_rate_limit_entries; i++) { + if (match_phys(&link_data->rate_limits[i].phys, phys)) { + entry = &link_data->rate_limits[i]; + break; + } + if (link_data->rate_limits[i].last_discovery_time_us < + oldest_time) { + oldest_time = + link_data->rate_limits[i].last_discovery_time_us; + oldest_idx = i; + } + } + + if (!entry) { + if (link_data->num_rate_limit_entries < + MAX_LINK_RATE_LIMIT_ENTRIES) { + entry = &link_data->rate_limits + [link_data->num_rate_limit_entries++]; + } else { + entry = &link_data->rate_limits[oldest_idx]; + } + memset(entry, 0, sizeof(*entry)); + entry->phys = *phys; + entry->last_discovery_time_us = now_us; + entry->discovery_count = 0; + } + + if (now_us - entry->last_discovery_time_us > 1000000ULL) { + entry->last_discovery_time_us = now_us; + entry->discovery_count = 0; + } + + if (entry->discovery_count >= 5) + return false; + + entry->discovery_count++; + return true; +} + +static bool is_discovery_pending(struct ctx *ctx, int ifindex, + const dest_phys *phys) +{ + struct discovery_notify_ctx *d; + + for (d = ctx->pending_discoveries; d; d = d->next) { + if (d->phys.ifindex == ifindex && match_phys(&d->phys, phys)) + return true; + } + return false; +} + +static void process_deferred_discovery(struct ctx *ctx, + struct discovery_notify_ctx *dctx) +{ + struct peer *peer = NULL; + int rc; + + peer = find_peer_by_phys(ctx, &dctx->phys); + if (peer) { + if (ctx->verbose) { + warnx("Discovery Notify received for existing peer %s (EID %d); tearing down for re-discovery", + peer_tostr(peer), peer->eid); + } + remove_peer(peer); + peer = NULL; + } + + rc = endpoint_assign_eid(ctx, NULL, &dctx->phys, &peer, 0, true); + + if (rc == 0) { + if (peer && peer->pool_size > 0) { + endpoint_allocate_eids(peer); + } + } else if (ctx->verbose) { + warnx("Deferred EID assignment failed for %s: %s", + dest_phys_tostr(&dctx->phys), strerror(-rc)); + } +} + +static int deferred_discoveries_cb(sd_event_source *s, void *userdata) +{ + struct ctx *ctx = userdata; + struct discovery_notify_ctx *dctx; + + (void)s; + sd_event_source_unref(ctx->discovery_defer); + ctx->discovery_defer = NULL; + + while (ctx->pending_discoveries) { + dctx = ctx->pending_discoveries; + ctx->pending_discoveries = dctx->next; + + process_deferred_discovery(ctx, dctx); + free(dctx); + } + + return 0; +} + +static int handle_control_discovery_notify(struct ctx *ctx, int sd, + const struct sockaddr_mctp_ext *addr, + const uint8_t *buf, + const size_t buf_size) +{ + struct mctp_ctrl_resp_discovery_notify respi = { 0 }, *resp = &respi; + struct discovery_notify_ctx *dctx = NULL; + struct mctp_ctrl_msg_hdr *req = NULL; + struct link *link_data; + dest_phys phys = { 0 }; + int rc; + + if (buf_size < sizeof(*req)) { + warnx("short Discovery Notify message"); + return -ENOMSG; + } + req = (void *)buf; + + link_data = mctp_nl_get_link_userdata(ctx->nl, addr->smctp_ifindex); + if (!link_data) { + bug_warn("unconfigured interface %d", addr->smctp_ifindex); + return -ENOENT; + } + + if (link_data->role != ENDPOINT_ROLE_BUS_OWNER) { + if (ctx->verbose) { + warnx("Ignoring Discovery Notify on interface %d: not a bus owner", + addr->smctp_ifindex); + } + return 0; + } + + phys.ifindex = addr->smctp_ifindex; + phys.hwaddr_len = addr->smctp_halen; + if (addr->smctp_halen > 0 && addr->smctp_halen <= sizeof(phys.hwaddr)) { + memcpy(phys.hwaddr, addr->smctp_haddr, addr->smctp_halen); + } + + if (!link_check_rate_limit(link_data, &phys)) { + if (ctx->verbose) { + warnx("Discovery Notify rate limit exceeded for %s", + dest_phys_tostr(&phys)); + } + return 0; + } + + /* Respond immediately over physical socket to acknowledge Discovery Notify */ + mctp_ctrl_msg_hdr_init_resp(&respi.ctrl_hdr, *req); + resp->completion_code = MCTP_CTRL_CC_SUCCESS; + + rc = reply_message_phys(ctx, sd, resp, sizeof(*resp), addr); + if (rc < 0) { + warnx("Failed to send Discovery Notify response on ifindex %d: %s", + addr->smctp_ifindex, strerror(-rc)); + return rc; + } + + /* Deduplicate incoming Discovery Notify for both existing and new physical endpoints */ + if (is_discovery_pending(ctx, addr->smctp_ifindex, &phys)) { + return 0; + } + + dctx = calloc(1, sizeof(*dctx)); + if (!dctx) { + return -ENOMEM; + } + + dctx->ctx = ctx; + dctx->phys = phys; + + dctx->next = ctx->pending_discoveries; + ctx->pending_discoveries = dctx; + + if (!ctx->discovery_defer) { + rc = sd_event_add_defer(ctx->event, &ctx->discovery_defer, + deferred_discoveries_cb, ctx); + if (rc < 0) { + if (ctx->verbose) { + warnx("Failed to defer EID assignment event for %s", + dest_phys_tostr(&phys)); + } + ctx->pending_discoveries = dctx->next; + free(dctx); + return rc; + } + } + + return 0; +} + static int handle_control_unsupported(struct ctx *ctx, int sd, const struct sockaddr_mctp_ext *addr, const uint8_t *buf, const size_t buf_size) @@ -1423,6 +1659,10 @@ static int cb_listen_control_msg(sd_event_source *s, int sd, uint32_t revents, rc = handle_control_endpoint_discovery(ctx, sd, &addr, buf, buf_size); break; + case MCTP_CTRL_CMD_DISCOVERY_NOTIFY: + rc = handle_control_discovery_notify(ctx, sd, &addr, buf, + buf_size); + break; default: if (ctx->verbose) { warnx("Ignoring unsupported command code 0x%02x", @@ -2216,6 +2456,8 @@ static int remove_peer(struct peer *peer) static void free_peers(struct ctx *ctx) { + struct discovery_notify_ctx *dctx, *tmp_dctx; + for (size_t i = 0; i < ctx->num_peers; i++) { struct peer *peer = ctx->peers[i]; free(peer->message_types); @@ -2231,6 +2473,17 @@ static void free_peers(struct ctx *ctx) } free(ctx->peers); + + sd_event_source_disable_unref(ctx->discovery_defer); + ctx->discovery_defer = NULL; + + dctx = ctx->pending_discoveries; + while (dctx) { + tmp_dctx = dctx->next; + free(dctx); + dctx = tmp_dctx; + } + ctx->pending_discoveries = NULL; } /* Returns -EEXIST if the new_eid is already used */ diff --git a/tests/mctp_test_utils.py b/tests/mctp_test_utils.py index 338c1282..d6c9757f 100644 --- a/tests/mctp_test_utils.py +++ b/tests/mctp_test_utils.py @@ -1,3 +1,13 @@ +import trio + + +async def wait_until(predicate, timeout=2.0, interval=0.005): + """Poll predicate every `interval` seconds until truthy or `timeout` expires.""" + with trio.fail_after(timeout): + while not predicate(): + await trio.sleep(interval) + + async def mctpd_mctp_base_iface_obj(dbus): obj = await dbus.get_proxy_object( 'au.com.codeconstruct.MCTP1', '/au/com/codeconstruct/mctp1' diff --git a/tests/test_mctpd.py b/tests/test_mctpd.py index 78537947..df54f351 100644 --- a/tests/test_mctpd.py +++ b/tests/test_mctpd.py @@ -9,6 +9,7 @@ mctpd_mctp_endpoint_common_obj, mctpd_mctp_endpoint_control_obj, mctpd_mctp_base_iface_obj, + wait_until, ) from mctpenv import ( Endpoint, @@ -2339,3 +2340,182 @@ async def test_iface_config_match_path_none(dbus, sysnet, nursery): res = await mctpd.stop_mctpd() assert res == 0 + + +async def test_discovery_notify_bus_owner_success(dbus, mctpd): + """Test Discovery Notify processing when mctpd is in Bus Owner role. + + When an endpoint issues Discovery Notify (0x0D), mctpd immediately + acknowledges over physical socket and defers EID assignment. + """ + ep = mctpd.network.endpoints[0] + + # Endpoint has no EID yet + assert ep.eid is None or ep.eid == 0 + + # Send Discovery Notify (Command Code 0x0D, Request bit set, IID 1) + cmd = MCTPControlCommand(True, 1, 0x0D) + rsp = await ep.send_control(mctpd.network.mctp_socket, cmd) + + # Expect immediate ACK with MCTP_CTRL_CC_SUCCESS (0x00) and IID 1 (0x01 0x0d 0x00) + assert rsp.hex(' ') == '01 0d 00' + + # Wait until deferred EID assignment executes and endpoint receives an EID + await wait_until(lambda: ep.eid is not None and ep.eid != 0) + + # Verify endpoint object created on D-Bus and EID assigned + assert await mctpd_mctp_endpoint_control_obj( + dbus, f"/au/com/codeconstruct/mctp1/networks/1/endpoints/{ep.eid}" + ) + + # Verify neighbour and route entries created in kernel + assert len(mctpd.system.neighbours) == 1 + assert mctpd.system.neighbours[0].lladdr == ep.lladdr + assert mctpd.system.neighbours[0].eid == ep.eid + assert len(mctpd.system.routes) == 1 + + +async def test_discovery_notify_deduplication(dbus, mctpd): + """Verify deduplication when multiple Discovery Notify requests arrive in rapid succession.""" + ep = mctpd.network.endpoints[0] + + # Send two Discovery Notify requests back-to-back before yielding to event loop + cmd1 = MCTPControlCommand(True, 1, 0x0D) + cmd2 = MCTPControlCommand(True, 2, 0x0D) + + rsp1 = await ep.send_control(mctpd.network.mctp_socket, cmd1) + rsp2 = await ep.send_control(mctpd.network.mctp_socket, cmd2) + + assert rsp1.hex(' ') == '01 0d 00' + assert rsp2.hex(' ') == '02 0d 00' + + # Wait until single EID assignment completes + await wait_until(lambda: ep.eid is not None and ep.eid != 0) + + # Ensure valid D-Bus object path + assert await mctpd_mctp_endpoint_control_obj( + dbus, f"/au/com/codeconstruct/mctp1/networks/1/endpoints/{ep.eid}" + ) + + +async def test_discovery_notify_existing_endpoint(dbus, mctpd, routed_ep): + """Verify Discovery Notify on an already assigned endpoint triggers teardown and re-discovery.""" + ep = routed_ep + + cmd = MCTPControlCommand(True, 3, 0x0D) + rsp = await ep.send_control(mctpd.network.mctp_socket, cmd) + assert rsp.hex(' ') == '03 0d 00' + + # Wait until re-discovery and EID assignment completes + await wait_until(lambda: ep.eid is not None and ep.eid != 0) + + # Endpoint is published on D-Bus + assert await mctpd_mctp_endpoint_control_obj( + dbus, f"/au/com/codeconstruct/mctp1/networks/1/endpoints/{ep.eid}" + ) + + +async def test_discovery_notify_short_message(mctpd): + """Verify truncated Discovery Notify message is rejected without crash.""" + ep = mctpd.network.endpoints[0] + + # Send 1-byte raw message (< sizeof(struct mctp_ctrl_msg_hdr)) + # mctpd should drop without reply (timing out on response wait) + addr = MCTPSockAddr(ep.iface.net, ep.eid or 0, 0, 0x80) + if mctpd.network.mctp_socket.addr_ext: + addr.set_ext(ep.iface.ifindex, ep.lladdr) + + await mctpd.network.mctp_socket.send(addr, bytes([0x80])) + + # Verify no endpoint was assigned or created from truncated message + await wait_until(lambda: ep.eid is None or ep.eid == 0) + assert ep.eid is None or ep.eid == 0 + + +async def test_discovery_notify_cleanup_on_shutdown(dbus, sysnet, nursery): + """Verify free_peers cleans up pending discoveries safely on shutdown.""" + mctpd = MctpdWrapper(dbus, sysnet) + await mctpd.start_mctpd(nursery) + + ep = mctpd.network.endpoints[0] + + # Send Discovery Notify + cmd = MCTPControlCommand(True, 1, 0x0D) + rsp = await ep.send_control(mctpd.network.mctp_socket, cmd) + assert rsp.hex(' ') == '01 0d 00' + + # Immediately stop mctpd before deferred callback completes + res = await mctpd.stop_mctpd() + assert res == 0 + + +async def test_discovery_notify_rate_limit(dbus, mctpd): + """Verify rate-limiting drops excessive Discovery Notify requests exceeding 5 per second per endpoint.""" + ep = mctpd.network.endpoints[0] + + # Send 5 valid Discovery Notify requests from first endpoint + for i in range(1, 6): + cmd = MCTPControlCommand(True, i, 0x0D) + rsp = await ep.send_control(mctpd.network.mctp_socket, cmd) + assert rsp.hex(' ') == f'{i:02x} 0d 00' + + # 6th request from first endpoint within same second should be rate-limited and dropped + cmd6 = MCTPControlCommand(True, 6, 0x0D) + with trio.move_on_after(0.2) as scope: + await ep.send_control(mctpd.network.mctp_socket, cmd6) + assert scope.cancelled_caught + + # Verify that a second endpoint on the same interface is NOT rate-limited (per-endpoint isolation) + ep2 = Endpoint(ep.iface, bytes([0x22]), types=[0, 1]) + mctpd.network.add_endpoint(ep2) + cmd_ep2 = MCTPControlCommand(True, 1, 0x0D) + rsp_ep2 = await ep2.send_control(mctpd.network.mctp_socket, cmd_ep2) + assert rsp_ep2.hex(' ') == '01 0d 00' + + +async def test_discovery_notify_uart_success(dbus, mctpd): + """Test Discovery Notify over UART / serial interface with addressless transport. + + Simulates an MCTP-over-UART (DSP0253 / PhysicalBinding.SERIAL) endpoint sending + Discovery Notify (0x0D) with an empty hardware address (bytes([])), verifying + that mctpd assigns an EID, updates the routing table, and publishes the D-Bus + endpoint object without creating invalid neighbour entries. + """ + net = 1 + # Create point-to-point UART / serial interface (bytes([]) hardware address) + uart_iface = mctpd.system.Interface( + 'mctpserial0', + 2, + net, + bytes([]), + 68, + 254, + True, + phys_binding=PhysicalBinding.SERIAL, + ) + await mctpd.system.add_interface(uart_iface) + await mctpd.system.add_address(mctpd.system.Address(uart_iface, 8)) + + # Add an unassigned remote endpoint connected over the UART interface + ep = Endpoint(uart_iface, bytes([]), types=[0, 1]) + mctpd.network.add_endpoint(ep) + assert ep.eid is None or ep.eid == 0 + + # Send Discovery Notify (0x0D) from the UART endpoint + cmd = MCTPControlCommand(True, 1, 0x0D) + rsp = await ep.send_control(mctpd.network.mctp_socket, cmd) + assert rsp.hex(' ') == '01 0d 00' + + # Wait until deferred EID assignment completes on UART endpoint + await wait_until(lambda: ep.eid is not None and ep.eid != 0) + + # Verify that the endpoint received the expected EID assignment + assert await mctpd_mctp_endpoint_control_obj( + dbus, f"/au/com/codeconstruct/mctp1/networks/1/endpoints/{ep.eid}" + ) + + # Verify routing table has route via UART interface and no invalid neighbour entry + route = mctpd.system.lookup_route(net, ep.eid) + assert route is not None + assert route.iface == uart_iface + assert not any(n.eid == ep.eid for n in mctpd.system.neighbours) diff --git a/tests/test_mctpd_endpoint.py b/tests/test_mctpd_endpoint.py index f29a967a..89730f46 100644 --- a/tests/test_mctpd_endpoint.py +++ b/tests/test_mctpd_endpoint.py @@ -1,8 +1,9 @@ -import pytest import asyncdbus +import pytest +import trio from mctp_test_utils import ( - mctpd_mctp_iface_control_obj, mctpd_mctp_endpoint_control_obj, + mctpd_mctp_iface_control_obj, ) from mctpenv import ( Endpoint, @@ -218,3 +219,13 @@ async def test_simple(self, dbus, mctpd): mctpd.network.mctp_socket, MCTPControlCommand(True, 0, 0x0B) ) assert rsp.hex(' ') == '00 0b 05' + + async def test_discovery_notify_ignored(self, dbus, mctpd): + """Discovery Notify command on endpoint interface (not bus owner) is ignored.""" + bo = mctpd.network.endpoints[0] + + # Send Discovery Notify (0x0D) to mctpd running in endpoint role + cmd = MCTPControlCommand(True, 1, 0x0D) + with trio.move_on_after(0.5) as scope: + await bo.send_control(mctpd.network.mctp_socket, cmd) + assert scope.cancelled_caught