diff --git a/src/mrib.c b/src/mrib.c index 4b84fa5..2953472 100644 --- a/src/mrib.c +++ b/src/mrib.c @@ -173,15 +173,23 @@ static size_t mrib_find(int ifindex) return i; } +// Build the forwarding-filter for a (group, source) pair from all users +static void mrib_filter_build(struct mrib_iface *iface, + const struct in6_addr *group, const struct in6_addr *source, + mrib_filter *filter) +{ + struct mrib_user *user; + list_for_each_entry(user, &iface->users, head) + if (user->cb_newsource) + user->cb_newsource(user, group, source, filter); +} + // Notify all users of a new multicast source static void mrib_notify_newsource(struct mrib_iface *iface, const struct in6_addr *group, const struct in6_addr *source) { mrib_filter filter = 0; - struct mrib_user *user; - list_for_each_entry(user, &iface->users, head) - if (user->cb_newsource) - user->cb_newsource(user, group, source, &filter); + mrib_filter_build(iface, group, source, &filter); char groupbuf[INET6_ADDRSTRLEN], sourcebuf[INET6_ADDRSTRLEN]; inet_ntop(AF_INET6, group, groupbuf, sizeof(groupbuf)); @@ -203,6 +211,29 @@ static void mrib_notify_newsource(struct mrib_iface *iface, } } +// Re-evaluate the forwarding state of all known sources of a group +// (e.g. after a downstream membership change). Sources are usually +// detected before the first downstream member subscribes, in which +// case their multicast route is installed as a placeholder without +// outgoing interfaces and would only be re-evaluated after the +// regular lifetime-based cleanup otherwise. +void mrib_refresh_user(struct mrib_user *user, const struct in6_addr *group) +{ + struct mrib_iface *iface = user->iface; + if (!iface) + return; + + struct mrib_route *c; + list_for_each_entry(c, &iface->routes, head) { + if (!IN6_ARE_ADDR_EQUAL(&c->group, group)) + continue; + + mrib_filter filter = 0; + mrib_filter_build(iface, &c->group, &c->source, &filter); + mrib_set(&c->group, &c->source, iface, filter, 0); + } +} + // Calculate IGMP-checksum static uint16_t igmp_checksum(const uint16_t *buf, size_t len) { diff --git a/src/mrib.h b/src/mrib.h index b3dd72b..ce8ae05 100644 --- a/src/mrib.h +++ b/src/mrib.h @@ -80,6 +80,9 @@ int mrib_flush(struct mrib_user *user, const struct in6_addr *group, uint8_t gro // Add interface to filter int mrib_filter_add(mrib_filter *filter, struct mrib_user *user); +// Re-evaluate forwarding state of all known sources of a group +void mrib_refresh_user(struct mrib_user *user, const struct in6_addr *group); + // Send IGMP-packet int mrib_send_igmp(struct mrib_querier *querier, struct igmpv3_query *igmp, size_t len, const struct sockaddr_in *dest); diff --git a/src/proxy.c b/src/proxy.c index 80e7b58..6544b4a 100644 --- a/src/proxy.c +++ b/src/proxy.c @@ -35,6 +35,7 @@ struct proxy_downlink { struct querier_user_iface iface; struct mrib_user mrib; struct client client; + struct proxy *proxy; enum proxy_flags flags; }; @@ -92,8 +93,13 @@ static void proxy_trigger(struct querier_user_iface *user, const struct in6_addr bool include, const struct in6_addr *sources, size_t len) { struct proxy_downlink *iface = container_of(user, struct proxy_downlink, iface); - if (proxy_match_scope(iface->flags, group)) + if (proxy_match_scope(iface->flags, group)) { client_set(&iface->client, group, include, sources, len); + + // Re-evaluate the forwarding state of all known sources of this + // group, as source detection may have happened before membership + mrib_refresh_user(&iface->proxy->mrib, group); + } } // Remove proxy with given name @@ -181,6 +187,8 @@ int proxy_set(int uplink, const int downlinks[], size_t downlinks_cnt, enum prox if (!downlink) goto err; + downlink->proxy = proxy; + if (client_init(&downlink->client, uplink)) goto downlink_err3;