fix(federation): Queue PDUs during remote joins #2156

Merged
Aranjedeath merged 1 commit from eleboucher/continuwuity:fix/remote-join-pdu-queue into main 2026-08-22 18:43:24 +00:00
Contributor

Queues federation PDUs received during remote joins until the room state is committed, preventing them from being dropped before promotion.

Fixes: #2063

Pull request checklist:

  • This pull request targets the main branch, and the branch is named something other than
    main.
  • I have written an appropriate pull request title and my description is clear.
  • I understand I am responsible for the contents of this pull request.
  • I have followed the contributing guidelines:
<!-- In order to help reviewers know what your pull request does at a glance, you should ensure that 1. Your PR title is a short, single sentence describing what you changed 2. You have described in more detail what you have changed, why you have changed it, what the intended effect is, and why you think this will be beneficial to the project. If you have made any potentially strange/questionable design choices, but didn't feel they'd benefit from code comments, please don't mention them here - after opening your pull request, go to "files changed", and click on the "+" symbol in the line number gutter, and attach comments to the lines that you think would benefit from some clarification. --> Queues federation PDUs received during remote joins until the room state is committed, preventing them from being dropped before promotion. <!-- Example: This pull request allows us to warp through time and space ten times faster than before by double-inverting the warp drive with hyperheated jump fluid, both making the drive faster and more efficient. This resolves the common issue where we have to wait more than 10 milliseconds to engage, use, and disengage the warp drive when travelling between galaxies. --> Fixes: #2063 <!-- Fixes: #... --> <!-- Uncomment the above line(s) if your pull request fixes an issue or closes another pull request by superseding it. Replace `#...` with the issue/pr number, such as `#123`. --> **Pull request checklist:** <!-- You need to complete these before your PR can be considered. If you aren't sure about some, feel free to ask for clarification in #dev:continuwuity.org. --> - [x] This pull request targets the `main` branch, and the branch is named something other than `main`. - [x] I have written an appropriate pull request title and my description is clear. - [x] I understand I am responsible for the contents of this pull request. - I have followed the [contributing guidelines][c1]: - [x] My contribution follows the [code style][c2], if applicable. - [x] I ran [pre-commit checks][c1pc] before opening/drafting this pull request. - [x] I have [tested my contribution][c1t] (or proof-read it for documentation-only changes) myself, if applicable. This includes ensuring code compiles. - [x] My commit messages follow the [commit message format][c1cm] and are descriptive. <!-- Notes on these requirements: - While not required, we encourage you to sign your commits with GPG or SSH to attest the authenticity of your changes. - While we allow LLM-assisted contributions, we do not appreciate contributions that are low quality, which is typical of machine-generated contributions that have not had a lot of love and care from a human. Please do not open a PR if all you have done is asked ChatGPT to tidy up the codebase with a +-100,000 diff. - In the case of code style violations, reviewers may leave review comments/change requests indicating what the ideal change would look like. For example, a reviewer may suggest you lower a log level, or use `match` instead of `if/else` etc. - In the case of code style violations, pre-commit check failures, minor things like typos/spelling errors, and in some cases commit format violations, reviewers may modify your branch directly, typically by making changes and adding a commit. Particularly in the latter case, a reviewer may rebase your commits to squash "spammy" ones (like "fix", "fix", "actually fix"), and reword commit messages that don't satisfy the format. - Pull requests MUST pass the `Checks` CI workflows to be capable of being merged. This can only be bypassed in exceptional circumstances. If your CI flakes, let us know in matrix:r/dev:continuwuity.org. - Pull requests have to be based on the latest `main` commit before being merged. If the main branch changes while you're making your changes, you should make sure you rebase on main before opening a PR. Your branch will be rebased on main before it is merged if it has fallen behind. - We typically only do fast-forward merges, so your entire commit log will be included. Once in main, it's difficult to get out cleanly, so put on your best dress, smile for the cameras! --> [c1]: https://forgejo.ellis.link/continuwuation/continuwuity/src/branch/main/CONTRIBUTING.md [c2]: https://forgejo.ellis.link/continuwuation/continuwuity/src/branch/main/docs/development/code_style.mdx [c1pc]: https://forgejo.ellis.link/continuwuation/continuwuity/src/branch/main/CONTRIBUTING.md#pre-commit-checks [c1t]: https://forgejo.ellis.link/continuwuation/continuwuity/src/branch/main/CONTRIBUTING.md#running-tests-locally [c1cm]: https://forgejo.ellis.link/continuwuation/continuwuity/src/branch/main/CONTRIBUTING.md#commit-messages
fix(federation): Queue PDUs during remote joins
Some checks failed
Auto Labeler / Apply labels based on changed files (pull_request_target) Successful in 3s
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Failing after 10s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m16s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m21s
e32df87c35
chore: Add changelog for remote join PDUs
Some checks failed
Checks / Changelog / Check changelog is added (pull_request_target) Has been cancelled
Documentation / Build and Deploy Documentation (pull_request) Has been cancelled
Checks / Prek / Pre-commit & Formatting (pull_request) Has been cancelled
Checks / Prek / Check changed files (pull_request) Has been cancelled
Checks / Prek / Clippy and Cargo Tests (pull_request) Has been cancelled
56ded53345
eleboucher force-pushed fix/remote-join-pdu-queue from 56ded53345
Some checks failed
Checks / Changelog / Check changelog is added (pull_request_target) Has been cancelled
Documentation / Build and Deploy Documentation (pull_request) Has been cancelled
Checks / Prek / Pre-commit & Formatting (pull_request) Has been cancelled
Checks / Prek / Check changed files (pull_request) Has been cancelled
Checks / Prek / Clippy and Cargo Tests (pull_request) Has been cancelled
to 9a0c936927
All checks were successful
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m17s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 7s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m21s
2026-08-20 14:56:25 +00:00
Compare
nex requested changes 2026-08-20 15:29:14 +00:00
Dismissed
@ -0,0 +4,4 @@
impl super::Service {
pub fn begin_remote_join(&self, room_id: &RoomId) {
let (tx, mut rx) = tokio::sync::mpsc::channel(100);
Owner

The channel buffer capacity could probably be customisable

The channel buffer capacity could probably be customisable
eleboucher marked this conversation as resolved
@ -0,0 +43,4 @@
}
pub fn cancel_remote_join(&self, room_id: &RoomId) {
self.joining_rooms.write().remove(room_id);
Owner

What happens if the join fails? I think this'll make the queue task leak, since the channel doesn't get closed, no?

What happens if the join fails? I think this'll make the queue task leak, since the channel doesn't get closed, no?
Author
Contributor

if they join i made them closed and discard the queue

if they join i made them closed and discard the queue
Owner

Are you sure? I think just removing the room ID from the map (called by the defer! in join_remote_room if any error is propagated or when the function ends) will cause the rx.recv().await call in begin_remote_join to hang forever if the recv channel is never closed and it never receives the Complete enum, leaving the task dangling with no way to interrupt it either. Are you sure you don't need to close the channel or send it super::PendingJoinPdu::Complete in cancel_remote_join?

Are you sure? I think just removing the room ID from the map (called by the defer! in `join_remote_room` if any error is propagated or when the function ends) will cause the `rx.recv().await` call in `begin_remote_join` to hang forever if the recv channel is never closed and it never receives the Complete enum, leaving the task dangling with no way to interrupt it either. Are you sure you don't need to close the channel or send it `super::PendingJoinPdu::Complete` in `cancel_remote_join`?
Author
Contributor

remove drops the queue’s only Sender, so rx.recv().await returns None and the task exits without processing buffered PDUs.

remove drops the queue’s only Sender, so rx.recv().await returns None and the task exits without processing buffered PDUs.
Owner

Oh yes, right you are

Oh yes, [right you are](https://docs.rs/tokio/latest/tokio/sync/mpsc/index.html#disconnection)
nex marked this conversation as resolved
@ -0,0 +62,4 @@
};
let pdu = super::PendingJoinPdu::Pdu(
origin.to_owned(),
RawJsonValue::from_string(pdu.get().to_owned()).expect("raw PDU JSON is valid"),
Owner

This seems like one hell of a roundtrip, can't pdu just be cloned?

This seems like one hell of a roundtrip, can't `pdu` just be cloned?
eleboucher marked this conversation as resolved
@ -300,6 +301,8 @@ impl Service {
) -> Result {
// public so the admin command force-join-room-remotely works
info!("Joining {room_id} over federation.");
self.services.event_handler.begin_remote_join(room_id);
Owner

There's no point starting the queue until just before we execute send_join (since it's at that point the remote server may have persisted our join)

There's no point starting the queue until just before we execute `send_join` (since it's at that point the remote server may have persisted our join)
eleboucher marked this conversation as resolved
@ -677,1 +680,4 @@
self.services.sync.wake_all_joined(room_id).await;
self.services
.event_handler
.process_pending_join_pdus(room_id)
Owner

How does this work if we receive PDUs after the join finishes (meaning the channel is no longer receiving) while the queue is still processing? Won't the live incoming pdus and queued pdus race each other for the persistence lock?

How does this work if we receive PDUs after the join finishes (meaning the channel is no longer receiving) while the queue is still processing? Won't the live incoming pdus and queued pdus race each other for the persistence lock?
Author
Contributor

Both queued and live PDUs acquire the same per-room federation mutex, so they cannot process concurrently.

Both queued and live PDUs acquire the same per-room federation mutex, so they cannot process concurrently.
Owner

They can't process concurrently, but they could process out of order I think?

They can't process concurrently, but they could process out of order I think?
Owner

Oh I see, the federation mutex is held for the entire queue processing duration

Oh I see, the federation mutex is held for the entire queue processing duration
nex marked this conversation as resolved
nex requested changes 2026-08-20 16:04:42 +00:00
Dismissed
@ -0,0 +11,4 @@
};
self.services.server.runtime().spawn(async move {
let mut pdus = Vec::new();
Owner

There probably also needs to be a limit to how many PDUs we're willing to queue up to avoid malicious servers sending us many huge transactions to exhaust memory while joining

There probably also needs to be a limit to how many PDUs we're willing to queue up to avoid malicious servers sending us many huge transactions to exhaust memory while joining
Author
Contributor

made it configurable max 50

made it configurable max 50
nex marked this conversation as resolved
eleboucher force-pushed fix/remote-join-pdu-queue from 9a0c936927
All checks were successful
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m17s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 7s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m21s
to 2dffa36e52
All checks were successful
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 8s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m20s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m23s
2026-08-20 16:14:02 +00:00
Compare
eleboucher requested review from nex 2026-08-20 16:20:30 +00:00
feat(complement): Simplify local test runs
Some checks failed
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 8s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Clippy and Cargo Tests (pull_request) Has been cancelled
Checks / Prek / Pre-commit & Formatting (pull_request) Has been cancelled
22bbce25b5
eleboucher force-pushed fix/remote-join-pdu-queue from 22bbce25b5
Some checks failed
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 8s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Clippy and Cargo Tests (pull_request) Has been cancelled
Checks / Prek / Pre-commit & Formatting (pull_request) Has been cancelled
to 2dffa36e52
All checks were successful
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 8s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m20s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m23s
2026-08-20 20:39:48 +00:00
Compare
eleboucher force-pushed fix/remote-join-pdu-queue from 2dffa36e52
All checks were successful
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 8s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m20s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m23s
to 3c835ec86f
All checks were successful
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 7s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m18s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m22s
2026-08-21 06:25:48 +00:00
Compare
nex requested review from Jade 2026-08-22 16:17:54 +00:00
nex approved these changes 2026-08-22 18:26:34 +00:00
eleboucher force-pushed fix/remote-join-pdu-queue from 3c835ec86f
All checks were successful
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 7s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m18s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m22s
to df95abf5ab
Some checks failed
Documentation / Build and Deploy Documentation (pull_request) Has been skipped
Checks / Changelog / Check changelog is added (pull_request_target) Successful in 7s
Checks / Prek / Check changed files (pull_request) Successful in 6s
Checks / Prek / Pre-commit & Formatting (pull_request) Successful in 1m14s
Update flake hashes / update-flake-hashes (pull_request) Successful in 1m26s
Checks / Prek / Clippy and Cargo Tests (pull_request) Successful in 8m5s
Documentation / Build and Deploy Documentation (push) Successful in 1m3s
Checks / Prek / Check changed files (push) Successful in 6s
Checks / Prek / Pre-commit & Formatting (push) Successful in 1m15s
Release Docker Image / Build linux-amd64 (release) (push) Failing after 3m5s
Checks / Prek / Clippy and Cargo Tests (push) Successful in 9m39s
Release Docker Image / Build linux-arm64 (release) (push) Successful in 12m41s
Release Docker Image / Create Multi-arch Release Manifest (push) Has been skipped
Release Docker Image / Build linux-amd64 (max-perf) (push) Has been skipped
Release Docker Image / Build linux-arm64 (max-perf) (push) Has been skipped
Release Docker Image / Create Max-Perf Manifest (push) Has been skipped
Release Docker Image / Mirror Images (push) Has been skipped
Release Docker Image / Release Binaries (push) Has been skipped
2026-08-22 18:29:55 +00:00
Compare
ginger approved these changes 2026-08-22 18:49:16 +00:00
Sign in to join this conversation.
No reviewers
No milestone
No project
No assignees
4 participants
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set.

Reference
continuwuation/continuwuity!2156
No description provided.