Repository navigation
Conversation
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
force-pushed
the
fix/transfer-service-get-ssrf
branch
from
October 6, 2026 14:28
4c1b795 to
64d30bf
Compare
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.
Python Code Coverage Summary
Minimum allowed line rate is |
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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Fixes b/565101350
Problem
Any peer that can reach the listener could send a
kGetObjheader that told the service which host to connect to and which file to send there: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 callconnect().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 endBehavior change: the second hop is gone
AsyncGetis still asynchronous. It returns astd::future, and an outbound worker picks the task up from the queue. The change is on the responder. Before, the inbound worker sentkAckand queued aRespondToGetTask, 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:
task_queue_without limit.kErroron the same socket. Before, they were lost when the callback connection could not be opened.ExecutePutTaskalready does.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 runs2 × threadsworker threads (the Python default of 16 gives 32). Theinitialize()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)
kPutObj,kGetObj). Replies that arrive there (kRespondToGetObj,kAck,kError) and unknown message types close the connection. Before, an unsolicitedkRespondToGetObjcould write any file, and an unsolicitedkErrorcould fail any in-flight local task whosetask_idit named. The dispatchswitchhas nodefault, so-Wswitch(part of-Wall, which this build does not enable yet) flags a newMessageTypethat is not handled.HandleGetObjRequestregistered theRespondToGetTaskinpending_tasks_under the peer'stask_id, and the peer learns our ids from the requests we send it. AkGetObjthat 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 touchespending_tasks_.task_idanddest_obj_idwhen 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.ObjInfoHeaderis 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 everykError,kAckandkGetObjheader and in the address fields ofkRespondToGetObj.kErrorreply. Before, the exception dropped the connection and the requester hung.RecvAndDiscard) and replieskError. This keeps the connection in sync, which matters because the pool never replaces a closed connection.<dest>.tmp.obj_sizeis the peer's number, so akPutObjthat 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()callsshutdown(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
kGetObjexchange 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:GetObjDoesNotConnectToDestAddress,SpoofedDestAddressTest/*(12 addresses, including the metadata IP, broadcast, and bad ports)dest_addresssays; the response arrives on the request socketSequentialGetRequestsOnOneSocket…,ConcurrentGetRequestsBeyondThreadCount…(8 clients, 2 workers),GetSucceedsWhenRequesterAdvertisesUnreachableIpResponderSurvivesClientDisconnectMidResponse(16 MiB) /…ClosingInsteadOfAck, both with a single workerInvalidGetRequestTest/*,GetRequestForUnopenableObject…,GetRequestWithUnterminatedAddressFields…kErrorcarries the request'stask_id; the socket still works afterwardsReplyHeadersCarryNoBytesAfterTerminatorskErrorand akRespondToGetObjreply, read off a raw socket, are all zero after each field's terminatorMalformedGetResponseTest/*(zero size, wrong type, bare ACK, truncated payload, close),GetWritesToLocalDestinationNotToPeerNamedPathPutToUncreatableDestination…,Put/GetFailsCleanlyAndReusesConnection…,RecvAndDiscard_*kErroris sent, and the same pooled connection works afterwardsTruncatedPutLeavesNoTemporaryFile,Put/GetFailsInRenameWhenTargetExistAsADirectory,MalformedGetResponseTest/truncated_payload.tmpwas created leaves no file behindUnsolicitedRespondToGetObj…,UnsolicitedAckOnListener…,UnsolicitedErrorOnListenerDoesNotFailLocalTask,GetRequestNamingLocalTaskIdDoesNotTouchLocalTask,UnknownMessageType…kGetObjthat names an in-flight task'stask_idtouches that task, which completes on its own socketShutdownWhile…×2,ShutdownWithQueuedTasksAgainstSilentPeerCompletesShutdown()does not crash or hang, even with tasks still queued against a peer that never answersEvery new test was first run against the code before its commit and failed there. At the branch head, the full
transfer_service_testsuite (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-resortcatchblocks, thesockfd < 0guard, theepoll_ctl(ADD)failure path and thekRespondToGetObjheader-send failure, which need fault injection; and the shutdown guard inGetOrCreateConnectionPool, which is only reachable for a peer that has no pool yet whenShutdown()runs.Out of scope
source_obj_id/dest_obj_idto a base directory.ScopedConnectiondrop broken connections instead of returning them to the pool.GetOrCreateConnectionPoolholds the write lock while it connects, so one unreachable peer stalls every other outbound task during that time. Pre-existing.ExecuteGetTasklogs, but tolerates, atask_idmismatch in a response.RespondToGetTasknow runs inline; folding it into a plain method is a refactor for later.-Wall(so-Wswitchis active) and an ASAN run in CI.ConnectionPool::Initialize()returning false). Pre-existing."host:abc") makesstd::stoithrow inside the worker; thepackaged_taskinThreadPool::enqueueswallows it and the future never resolves. Pre-existing; fix withabsl::SimpleAtoi.