Repository navigation
Conversation
7bb1e32 to
2d9f696
Compare
reivilibre
left a comment
There was a problem hiding this comment.
Partial review owing to time
|
|
||
| Unlike the timeline-based path, we don't trust the remote server to tell us the | ||
| state at the event. Instead we fill in the event's state DAG by fetching any | ||
| missing `prev_state_events`, persist them as outliers, and then calculate the |
There was a problem hiding this comment.
A little confusion here about persisting them as outliers.
For the prev_state_events, won't we also fetch all their prev_state_events recursively and thus have the state for those (so they won't be conventional outliers, as we have state)?
Or is the point that: whilst we have the state events that the events themselves claim, we don't resolve it against other branches and don't retrospectively work out the 'current state' at the event?
There was a problem hiding this comment.
It's complicated. Outliers have two properties here:
- They don't have state attached to them
- They do not connect to the timeline (until you scrollback enough and then "de-outlier" them)
Note that the reason why we don't have state attached to them is because they do not connect to the timeline: the lack of timeline connection is the primary thing that makes an outlier, the lack of state is a consequence of this.
In this framing, outliers are the right choice for state DAG events. They do not connect to the timeline, but they do have state attached to them by virtue of them being part of the state DAG.
| async def _compute_event_context_with_maybe_missing_prevs_state_dag( | ||
| self, dest: str, event: MSC4242Event | ||
| ) -> EventContext: | ||
| """Build an EventContext for a pulled MSC4242 State DAG event whose |
There was a problem hiding this comment.
| """Build an EventContext for a pulled MSC4242 State DAG event whose | |
| """Build an EventContext for a non-outlier pulled MSC4242 State DAG event whose |
| # TODO(kegan): we need a staging area else we'll OOM well before this point | ||
| max_state_dag_events = 100_000 |
There was a problem hiding this comment.
I'm not going to block on this at this time, but have you had any thought about what we 'should' do here, for rooms with unbearably large numbers of state events?
Is the idea just that we should buffer them in the DB whilst we download them & make it cheap enough to persist large numbers of events occasionally?
It feels a shame that we have to download the whole gap before letting the DB start crunching through them. Would it be sensible to more or less flip this into a depth-first walk and start persisting one branch at a time to get the latency down? (Not that I'm saying to do this now, but it would help me to understand.)
There was a problem hiding this comment.
There's no good answer here unfortunately. There is a tradeoff between:
- tolerating long lasting partitions where you've missed an enormous amount of content
- tolerating malicious servers who flood you with an enormous amount of content
The max here is sufficiently high that I believe you'll only hit it in the malice case. In the honest case, it's safe to buffer them in the DB to continue downloading them and to then process them in causal order. We can in theory try to negotiate a "fork point/meet/greatest lower bound/most recent common ancestor" and then walk from there, but this has to tolerate malicious servers who can purposefully degrade this optimisation. In practice, CRDT folks shrug and just do set reconciliation, totally ignoring the shape of the DAG (see RBSR, Rateless IBLT), it's a research topic to do better here.
The flip-side approach here would be to set a limit to how much divergence we're willing to tolerate e.g. if you've really missed 100k+ events, perhaps you should just re-join the room at this point.
Hopefully this clarifies some things.
| if event.is_state(): | ||
| back_set = {event.event_id: event} | ||
| else: | ||
| prev_state_events = await self._get_events_from_remote( |
There was a problem hiding this comment.
It appears that if the remote returns 404 or another error for these events, those events are silently dropped from the back set.
It's a little bit surprising not to acknowledge this in some way.
- If the remote returns none of the events, the back set is empty and we don't iterate, so exit by returning
[]. Seems this actually turns out OK in the end, but feels like an accident. - If the remote returns some of the events, it seems like we will do a bunch of work (outbound
/gmerequests...) to fetch events, then eventually we will fail auth. It does look like we persist some of those events though. How much of this is intentional? Would be good to spell out the sad path explicitly in the comments
There was a problem hiding this comment.
You make some good points.. I'd rather just nuke the code than handle this. This code pre-dates the MSC changing to mention that the receiver MUST walk the prev_state_events of message events as the first hop if you pass a message event to /get_missing_events (the motivation here is partly for RTT/bandwidth, a single /gme where the receiver includes these events is superior to many /event request to then walk back from the prev_state_events, so in theory I think we can just nuke this else branch entirely.
If /get_missing_events is called with state_dag: true and latest_events contains a message event (not a state event) then the response MUST include the prev_state_events for that message event. This ensures that a server can perform a single /get_missing_events request to optimistically fill in the state DAG. Without this, a server would need to make additional /event requests to fetch the prev_state_events in the message event and then use those state events in the /get_missing_events request.
| # remember which events we're querying for. If we don't make forward progress | ||
| # we'll bail. | ||
| before_back_set = set(back_set) | ||
| max_events_per_req = limit * pow(2, iteration) # 8x1, 8x2, 8x4, ... |
There was a problem hiding this comment.
This should be capped. After only 29 iterations this will be 2^32 and may exceed the integer size at the remote.
There's also no way we actually want to accept that many events in one batch, is there?
| remote_events_map = { | ||
| ev.event_id: ev | ||
| for ev in remote_events | ||
| if ev.room_id == room_id and supports_msc4242_state_dag(ev) |
There was a problem hiding this comment.
The and supports_msc4242_state_dag(ev) part of this check seems dubious. The room version is taken from our local store, where we already know room_id to represent a state DAG room. So the check can never fail.
The room ID check is correct/necessary though
There was a problem hiding this comment.
You'd be right, but Python demands it. remote_events is list[EventBase] and it's supports_msc4242_state_dag which coerces them to be MSC4242Event which we need in order to access .prev_state_events.
| remote_event_ids, | ||
| ) | ||
| seen_event_ids.update(seen_remotes) | ||
| unseen_remotes = set(remote_events_map).difference(seen_event_ids) |
There was a problem hiding this comment.
suspicious not to use remote_event_ids here. I suspect that's intentional, but the comment narrative walks me into surprise
| # all unseen events must be returned | ||
| missed_events.update( |
There was a problem hiding this comment.
I think it would be good to put a comment on what each of missed and unseen and seen means (around R1834 I guess). When reaching this line I am finding myself confused because my intuition of missed seems to have been wrong.
| back_set = {event.event_id: event} | ||
| else: | ||
| back_set = { | ||
| event_id: missed_events[event_id] for event_id in new_back_set | ||
| } |
There was a problem hiding this comment.
I'm not seeing back_set be mutated before this line, so unsure about the reason for the before_back_set copy.
Given that back_set == before_back_set for all of this loop body until this line, it seems it might be better to remove the copy so I don't have to understand what the 'difference' is
| async def _compute_event_context_with_maybe_missing_prevs_state_dag( | ||
| self, dest: str, event: MSC4242Event | ||
| ) -> EventContext: | ||
| """Build an EventContext for a pulled MSC4242 State DAG event whose |
There was a problem hiding this comment.
Are we missing logic to accept /send-sent events, since this only seems to apply to pulled events (as docstring says)?
Co-authored-by: Eric Eastwood <madlittlemods@gmail.com>
Co-authored-by: Eric Eastwood <erice@element.io>
Co-authored-by: Eric Eastwood <erice@element.io>
2d9f696 to
02eeef1
Compare
Process inbound /send requests in a state DAG room.
This is more complex than you might think, since you may be sent events you do not know the
prev_state_eventsfor. You need to then hit/get_missing_eventsto fill the gap then process in causal order, and have to handle all the edge cases where the remote server is malicious.Split out from #19425
Part of a series of 5x PRs to land the federation part of MSC4242 (#19718, #20127, #20133, #20194, inbound-pulls (this PR)).
Working this PR has surfaced an issue with out of band invites not working so this isn't the final PR to land state DAGs unfortunately, there will need to be a 6th PR.
Unfortunately, the diff is large but unavoidably so. Around 70% of the diff is tests.
Pull Request Checklist
EventStoretoEventWorkerStore.".code blocks.