Skip to content

fix(replication): respond to kGetObj over existing client socket - #112

Open
g-husam wants to merge 15 commits into
mainfrom
fix/transfer-service-get-ssrf
Open

g-husam wants to merge 15 commits into
mainfrom
fix/transfer-service-get-ssrf

Conversation

@g-husam

@g-husam g-husam commented Sep 30, 2026 •

Copy link
Copy Markdown
Collaborator

Fixes b/565101350

Problem

Any peer that can reach the listener could send a kGetObj header that told the service which host to connect to and which file to send there:

ProcessIncomingData              ← header from any peer that can reach the listener
  HandleGetObjRequest
    task_queue_.push(RespondToGetTask(header.dest_address, header.source_obj_id))
      ExecuteRespondToGetTask     (later, on any outbound worker)
        GetConnectionFromPool(header.dest_address)   → connect()   ← SSRF
        send kRespondToGetObj + file(header.source_obj_id)

Fix

The service now sends the response back on the socket the request came in on. It no longer reads dest_address, so a peer can no longer make it call connect().

 ProcessIncomingData (inbound worker)
   HandleGetObjRequest
-    send kAck
-    task_queue_.push(RespondToGetTask)
-      ...later, outbound worker:
-        GetConnectionFromPool(header.dest_address) → connect()
-        send kRespondToGetObj + payload, wait kAck
+    RespondToGetTask.Execute(this)        // inline, same socket
+      send kRespondToGetObj + payload (or kError), wait kAck
sequenceDiagram
    participant R as Requester (AsyncGet worker)
    participant S as Responder (inbound worker)
    R->>S: kGetObj
    alt object readable
        S-->>R: kRespondToGetObj + payload
        R->>S: kAck
    else cannot open / not found
        S-->>R: kError
    end
Loading

Behavior change: the second hop is gone

AsyncGet is still asynchronous. It returns a std::future, and an outbound worker picks the task up from the queue. The change is on the responder. Before, the inbound worker sent kAck and queued a RespondToGetTask, and a second worker opened a connection back to the requester. Now the inbound worker that read the request sends the reply itself.

This is acceptable, and better, because:

  • The socket is busy either way. The requester blocks on that connection until the response arrives. Handing the reply to another worker only adds latency: one more queue hop and a connection-pool setup before the first byte.
  • Natural backpressure. The inbound pool size caps the number of in-flight responses. Before, they could pile up in task_queue_ without limit.
  • Errors reach the requester. They arrive as kError on the same socket. Before, they were lost when the callback connection could not be opened.
  • Same cost as Put. One worker handles one transfer from start to finish, as ExecutePutTask already does.
  • No deadlock. A dedicated epoll_thread_pool_ handles inbound events, so two peers that flood each other with Gets can still serve each other's responses (BidirectionalConcurrentGetAndPut). Initialize(threads=…) now sizes each pool, so a service runs 2 × threads worker threads (the Python default of 16 gives 32). The initialize() docstring says so. The trade-off is unchanged from Put: an inbound worker is held for the whole transfer it serves, and socket calls have no timeout.

Hardening (separate commits)

  • The listener only accepts requests (kPutObj, kGetObj). Replies that arrive there (kRespondToGetObj, kAck, kError) and unknown message types close the connection. Before, an unsolicited kRespondToGetObj could write any file, and an unsolicited kError could fail any in-flight local task whose task_id it named. The dispatch switch has no default, so -Wswitch (part of -Wall, which this build does not enable yet) flags a new MessageType that is not handled.
  • Remotely triggered work stays out of the local task bookkeeping. HandleGetObjRequest registered the RespondToGetTask in pending_tasks_ under the peer's task_id, and the peer learns our ids from the requests we send it. A kGetObj that named an in-flight local Get's id overwrote that Get's entry and dropped its promise (std::future_error: broken_promise). The responder now logs its own outcome and never touches pending_tasks_.
  • The requester uses its own task_id and dest_obj_id when it completes the task and writes the payload. Before, it used the values from the response, so a peer could redirect the write or leave the future pending forever.
  • ObjInfoHeader is zero-initialized. Only the first byte of each string field was; the rest of the 2249-byte header went on the wire as whatever the worker's stack held before, in every kError, kAck and kGetObj header and in the address fields of kRespondToGetObj.
  • Objects that cannot be opened (empty file, directory, name too long) now get a kError reply. Before, the exception dropped the connection and the requester hung.
  • When a receiver cannot create the destination file, it reads and discards the payload (RecvAndDiscard) and replies kError. This keeps the connection in sync, which matters because the pool never replaces a closed connection.
  • A failed inbound transfer removes its <dest>.tmp. obj_size is the peer's number, so a kPutObj that announced a large size and then stopped left a sparse file of that apparent size behind, and every failed Get left its destination's .tmp.
  • Shutdown() calls shutdown(2) on every accepted socket before joining the workers, so a stalled peer cannot hang it. It destroys the connection pools only after the workers are joined: a worker still inside a Get or Put returns its connection to the pool through a raw pointer, and before, that pool was already gone (heap-use-after-free under ASAN, reached by the queued-tasks shutdown test below).
  • Shutdown() no longer lets still-queued tasks create new connection pools. ThreadPool::stop() runs every queued task; before, such a task re-created a pool after the pools were shut down, and if the peer accepted the connection but never answered, shutdown() hung forever (pre-existing; also affected Put).

Wire protocol

This changes the kGetObj exchange on the wire. That is fine: the transport only runs within one training job, and every rank runs the same build, so old and new versions never talk to each other. The on-disk format is unchanged. Bump the minor version on release.

Tests

tests/replication/transfer_service/transfer_service_p2p_test.cpp, net_util_test.cpp:

Group Shows
GetObjDoesNotConnectToDestAddress, SpoofedDestAddressTest/* (12 addresses, including the metadata IP, broadcast, and bad ports) the service never connects out, whatever dest_address says; the response arrives on the request socket
SequentialGetRequestsOnOneSocket…, ConcurrentGetRequestsBeyondThreadCount… (8 clients, 2 workers), GetSucceedsWhenRequesterAdvertisesUnreachableIp happy paths, including more clients than workers
ResponderSurvivesClientDisconnectMidResponse (16 MiB) / …ClosingInsteadOfAck, both with a single worker a peer that leaves early only fails its own request
InvalidGetRequestTest/*, GetRequestForUnopenableObject…, GetRequestWithUnterminatedAddressFields… kError carries the request's task_id; the socket still works afterwards
ReplyHeadersCarryNoBytesAfterTerminators a kError and a kRespondToGetObj reply, read off a raw socket, are all zero after each field's terminator
MalformedGetResponseTest/* (zero size, wrong type, bare ACK, truncated payload, close), GetWritesToLocalDestinationNotToPeerNamedPath the requester's future always resolves; a peer cannot choose where the file is written
PutToUncreatableDestination…, Put/GetFailsCleanlyAndReusesConnection…, RecvAndDiscard_* the payload is discarded, kError is sent, and the same pooled connection works afterwards
TruncatedPutLeavesNoTemporaryFile, Put/GetFailsInRenameWhenTargetExistAsADirectory, MalformedGetResponseTest/truncated_payload a transfer that fails after the .tmp was created leaves no file behind
UnsolicitedRespondToGetObj…, UnsolicitedAckOnListener…, UnsolicitedErrorOnListenerDoesNotFailLocalTask, GetRequestNamingLocalTaskIdDoesNotTouchLocalTask, UnknownMessageType… a reply on the listener closes the connection and writes no file; neither a reply nor a kGetObj that names an in-flight task's task_id touches that task, which completes on its own socket
ShutdownWhile… ×2, ShutdownWithQueuedTasksAgainstSilentPeerCompletes Shutdown() does not crash or hang, even with tasks still queued against a peer that never answers

Every new test was first run against the code before its commit and failed there. At the branch head, the full transfer_service_test suite (119 tests) passed 3 runs in a row, and the new and timing-sensitive tests passed --gtest_repeat=10. pytest tests/replication (46 tests) passed against the rebuilt extension. MLFLogSinkTest.* was excluded locally (pre-existing, unrelated). An ASAN build (-fsanitize=address) of the P2P suite runs the four *Shutdown* tests clean for 10 iterations; it reported the pool use-after-free before that commit. A gcov build of the suite at the branch head shows every executable line this PR adds is executed, except: the two last-resort catch blocks, the sockfd < 0 guard, the epoll_ctl(ADD) failure path and the kRespondToGetObj header-send failure, which need fault injection; and the shutdown guard in GetOrCreateConnectionPool, which is only reachable for a peer that has no pool yet when Shutdown() runs.

Out of scope

  • Send/receive timeouts: a stalled peer still holds a worker until shutdown.
  • Peer authentication, and restricting source_obj_id/dest_obj_id to a base directory.
  • Having ScopedConnection drop broken connections instead of returning them to the pool.
  • GetOrCreateConnectionPool holds the write lock while it connects, so one unreachable peer stalls every other outbound task during that time. Pre-existing.
  • ExecuteGetTask logs, but tolerates, a task_id mismatch in a response.
  • RespondToGetTask now runs inline; folding it into a plain method is a refactor for later.
  • Enabling -Wall (so -Wswitch is active) and an ASAN run in CI.
  • No test covers a Put or Get to a peer that refuses connections (ConnectionPool::Initialize() returning false). Pre-existing.
  • A peer address with a non-numeric port ("host:abc") makes std::stoi throw inside the worker; the packaged_task in ThreadPool::enqueue swallows it and the future never resolves. Pre-existing; fix with absl::SimpleAtoi.

Follow-up hardening for the b/565101350 fix (respond to kGetObj over the
request's own connection instead of connecting back to header.dest_address).

Wire-protocol change (intentional): kGetObj is no longer acknowledged with
kAck; the responder streams kRespondToGetObj + payload on the request socket
and the requester replies kAck there. Old and new builds do not interoperate
for Gets. This is acceptable because the transport is scoped to a single
training job in which all ranks run the same build; the on-disk checkpoint
format is unchanged. Bump the minor version at release.

- Shutdown(): skip the promise-less RespondToGetTask entries in
  pending_tasks_ (was a null dereference whenever a remote Get was in
  flight), and shutdown(2) every accepted client socket before joining the
  thread pools so a responder blocked on a stalled peer cannot hang
  Shutdown(). Accepted sockets are tracked in client_fds_ and closed once.
- ExecuteRespondToGetTask(): BufferObject throws for empty files,
  directories and unreadable paths; the old "nullptr -> kError" branch was
  dead and the escaping exception orphaned the connection and hung the
  requester forever. Catch it and send kError on the request socket.
- ProcessIncomingData(): close the connection on an unsolicited
  kRespondToGetObj (it was an arbitrary file write sink), on an unknown
  message type, and on any handler exception, instead of leaving the fd
  registered but never re-armed.
- ExecuteGetTask(): report a failure (and shut the pooled connection down
  so it is not reused mid-payload) if the destination BufferObject cannot be
  created.
- Document that `threads` sizes each of the two pools, why two pools are
  required, and that RespondToGetTask must run synchronously on a
  non-owning fd.

Tests: strengthen GetObjDoesNotConnectToDestAddress (poll the decoy
listener instead of an instantaneous accept()), add a 12-case
SpoofedDestAddressTest, unterminated address fields, unopenable object ->
kError on the same socket, unsolicited kRespondToGetObj / unknown type ->
connection closed, Shutdown() while a responder waits for ACK and while a
peer is not reading (16 MiB, 4 KiB SO_RCVBUF), and the ported
GetSucceedsWhenRequesterAdvertisesUnreachableIp.
…ET/PUT failures

Follow-up hardening for the same-socket kGetObj response path.

* HandleDataReceive: when the destination cannot be created, drain the
  payload the peer is already streaming (new RecvAndDiscard) and reply
  kError instead of letting the exception escape. Pooled connections are
  never replaced once closed, so unread payload bytes left on the stream
  would have poisoned the connection for every later task on it.
* ExecuteGetTask: complete the task and write the payload using the local
  task's task_id and dest_obj_id. The values echoed by the peer were used
  before, so a peer could leave the AsyncGet future pending forever or
  choose where the payload was written. Drops the shutdown() stopgap,
  which would have permanently poisoned a pooled connection.
* HandleGetObjRequest: use the non-throwing std::filesystem::exists so
  that stat errors other than ENOENT (e.g. ENAMETOOLONG for an id that
  fills its field) answer kError instead of dropping the connection.
* net_util_test: RecvAndDiscard byte accounting, zero-length no-op,
  multi-chunk payloads and EOF reporting.
* transfer_service_p2p_test:
  - sequential and over-subscribed concurrent raw GETs (8 clients on 2
    inbound workers) are all served in order with the right task_id;
  - the responder survives a client that disconnects mid-response or
    closes instead of ACKing (single inbound worker, so a pinned worker
    would be visible);
  - invalid GET requests (empty, unknown and unterminated source_obj_id)
    get kError and the same socket then serves a valid GET;
  - uncreatable destinations are drained and answered with kError over a
    raw Put, AsyncPut and AsyncGet, and the single pooled connection is
    reused afterwards;
  - a fake peer returning malformed GET responses (zero size, wrong type,
    bare ACK, truncated payload, close) never hangs the requester's
    future, and a peer naming a different dest_obj_id cannot redirect
    where the payload is written.

All waits are bounded so a regression fails the test instead of hanging
the suite. Full suite green on 3 runs; the new and timing-sensitive tests
pass --gtest_repeat=10.
Comment-only change to the code and tests added in the previous three commits: shorter declarative sentences, no clefts or roundabout phrasing. The only non-comment change is a shorter warning message in ExecuteGetTask for a response whose task_id does not match the request.
…started

ThreadPool::stop() runs every task that is still queued before the workers exit. Shutdown() has already shut down and cleared the connection pools by then, so a queued Put or Get re-created a pool that nothing would ever shut down. If the peer completed the TCP handshake but never answered (for example, a wedged rank), the task blocked in recv() forever and Shutdown() never returned, which hangs the Python thread that called shutdown().

GetOrCreateConnectionPool now returns nullptr once running_ is false, so drained tasks fail fast with 'Failed to get connection'. This predates the kGetObj change and affects Put as well.

New test ShutdownWithQueuedTasksAgainstSilentPeerCompletes reproduces the hang (10 s timeout before the fix) and passes 5/5 with it. Full transfer_service_test 114/114; pytest tests/replication 46 passed.
The four failure branches before the ACK (bad obj_size, cannot create the buffer, receive failed, rename failed) all ended with the same three steps: resolve the local Get future if there is one, reply kError, return. A local fail() lambda now does that, so each branch is one line. The separate drain-failure log is gone: if the drain fails the peer is gone, and SendErrorResponse already logs the failed kError send.

No behavior change. Error messages now carry the obj_size and the rename errno text; the prefixes tests match on are unchanged. Full transfer_service_test 114/114, the failure-path tests 5/5 under --gtest_repeat, pytest tests/replication 46 passed.
…h switch

Move the switch on header.type into DispatchMessage. Every enumerator
returns, so the switch has no default and -Wswitch reports a new
MessageType that is not handled. The type byte comes from the network
and can hold any uint8_t value, so the unknown-type check stays, but
after the switch where it cannot mask that warning.

No behavior change. UnknownMessageTypeClosesConnection covers the
unknown-type path.
@g-husam
g-husam force-pushed the fix/transfer-service-get-ssrf branch from 4c1b795 to 64d30bf Compare October 6, 2026 14:28
@g-husam
g-husam requested a review from Leahlijuan October 6, 2026 14:38
@g-husam
g-husam marked this pull request as ready for review October 6, 2026 14:38
…listener

The listener only receives requests. Every reply (kRespondToGetObj, kAck,
kError) is read by the worker that owns the request socket, so a reply
arriving on the listener is a protocol violation. Treat all three the
same way: log and close the connection.

Before, kAck was logged and the connection kept, and kError called
ReportResult(task_id, false). That let any peer that could reach the
listener fail an in-flight local task by naming its task_id; the
responder learns that id from the request header.

Tests: UnsolicitedAckOnListenerClosesConnection and
UnsolicitedErrorOnListenerDoesNotFailLocalTask. The second one forges a
kError for a real in-flight Get and checks that the Get still completes
on the request socket. Both fail against the previous code.
@github-actions

github-actions Bot commented Oct 6, 2026

Copy link
Copy Markdown

Python Code Coverage Summary

Code Coverage

Package Line Rate Branch Rate Health
src.ml_flashpoint 100% 100% ✔
src.ml_flashpoint.adapter 100% 100% ✔
src.ml_flashpoint.adapter.megatron 97% 95% ✔
src.ml_flashpoint.adapter.nemo 98% 94% ✔
src.ml_flashpoint.adapter.pytorch 99% 92% ✔
src.ml_flashpoint.checkpoint_object_manager 93% 93% ➖
src.ml_flashpoint.core 95% 92% ✔
src.ml_flashpoint.replication 83% 83% ❌
Summary 95% (2396 / 2524) 92% (573 / 624) ➖

Minimum allowed line rate is 90%

@github-actions

github-actions Bot commented Oct 6, 2026

Copy link
Copy Markdown

C++ Code Coverage Summary

Code Coverage

Package Line Rate Branch Rate Health
src.ml_flashpoint.checkpoint_object_manager.buffer_object 94% 56% ✔
src.ml_flashpoint.replication.transfer_service 85% 47% ➖
Summary 87% (990 / 1140) 48% (784 / 1623) ➖

Minimum allowed line rate is 80%

@g-husam
g-husam enabled auto-merge (squash) October 6, 2026 19:11
… joined

Shutdown() cleared connection_pools_ before it stopped the thread pools.
A worker that was still inside a Get or Put held a ScopedConnection,
which returns its socket to the pool through a raw pointer when it goes
out of scope. With the pool already destroyed, that release was a
heap-use-after-free.

Shutdown() still shuts the pools down first, so blocked workers wake up
and give up, but it now destroys the pools only after both thread pools
have joined their workers.

Verified with an ASAN build of transfer_service_p2p_test.cpp
(-fsanitize=address): ShutdownWithQueuedTasksAgainstSilentPeerCompletes
reports the use-after-free in ConnectionPool::ReleaseConnection before
this change; all four *Shutdown* tests are clean for 10 iterations after
it. Without ASAN the bug has no deterministic symptom, so there is no
new plain test.
HandleGetObjRequest registered the RespondToGetTask in pending_tasks_
under the task_id from the request header, and ExecuteRespondToGetTask
reported its outcome through ReportResult with that id. Both are the
peer's choice, and the peer learns our task ids from the kGetObj
requests we send it. A peer that sent kGetObj naming the id of a local
Get in flight overwrote that task's entry, which dropped its promise:
the caller's future failed with std::future_error (broken_promise) and
the real reply later found no task to complete. The same path also
broke a service that fetched an object from itself.

RespondToGetTask now stays out of pending_tasks_ entirely. It logs its
outcome from its own metric container through LogTaskResult, which is
the same log line ReportResult writes, so the timing logs keep their
format (TimestampsAreRecorded still parses three of them). The
null-promise guard in Shutdown() goes away with the only null entry.

Test: GetRequestNamingLocalTaskIdDoesNotTouchLocalTask sends a kGetObj
for an existing object naming an in-flight Get's task_id, checks that
the request is served, that the Get's future stays pending, and that
the real reply then completes it. Fails against the previous code with
"Broken promise".
The header constructor set only the first byte of each char array.
snprintf writes up to the terminator, and SendAll ships all of
kHeaderSize, so the bytes after each terminator went on the wire as
whatever the worker's stack held before: pointers, and fragments of
headers from other peers. Every kError, kAck and kGetObj header and the
address fields of every kRespondToGetObj header leaked this way.

Default member initializers now zero the whole struct. It stays packed,
and its size does not change.

Test: ReplyHeadersCarryNoBytesAfterTerminators reads a kError and a
kRespondToGetObj reply off a raw socket and checks that every byte after
each field's terminator is zero. Fails against the previous code on all
five fields of the kError reply and on both address fields of the
kRespondToGetObj reply.
…fails

HandleDataReceive writes an incoming object to <dest>.tmp and renames it
at the end. When the receive or the rename failed, or when BufferObject's
constructor threw after it had created the file, the temporary file was
left behind. obj_size comes from the peer, so a kPutObj that announces a
large size and then stops leaves a sparse file of that apparent size on
disk, and every failed Get left its destination's .tmp behind.

The failure paths after the file is created now close and remove it.
The rename error text is captured before the removal so errno is not
overwritten.

Tests: TruncatedPutLeavesNoTemporaryFile (kPutObj announcing 16 MiB,
16 bytes sent, then EOF; expects kError and no file). The cleanup
lines in PutFailsInRenameWhenTargetExistAsADirectory,
GetFailsInRenameWhenTargetExistAsADirectory and MalformedGetResponseTest
become assertions that the temporary file is gone. All of these fail
against the previous code.
…pools

- DispatchMessage: the comment implied the compiler reports a new
  MessageType by default. -Wswitch is part of -Wall, which this build
  does not enable yet. Say so.
- transfer_service.h: state the trade-off of the inbound pool (a worker
  is held for the whole transfer it serves and socket calls have no
  timeout) and that the service runs 2 * threads workers.
- bindings.cpp: the initialize() docstring says the same, since Python
  callers only see the single `threads` argument.
- net_util_test.cpp: use EXPECT_TRUE instead of ASSERT_TRUE inside the
  writer thread, as gtest requires off the main thread.

No behavior change.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant