- 后端
- 即时通讯
【免费下载链接】dendrite
Dendrite is a second-generation Matrix homeserver written in Go!
导读
Dendrite 是使用 Go 编写的第二代 Matrix 家服务器(homeserver),而 Sync API 是其客户端同步体系的"心脏"——它专门负责处理客户端的/sync请求,让客户端能够获得完整会话历史、增量事件更新和房间状态。本文以 syncapi/README.md 为骨架,结合 syncapi/syncapi.go、syncapi/sync/requestpool.go、syncapi/streams/stream_pdu.go 等源码实现,系统讲解 Sync API 的职责范围、房间选择与事件计算逻辑、状态增量模型以及已知限制,读完本文你将能理解 Dendrite/sync的完整数据流与演进脉络。
一、Sync API 的定位与核心职责
Sync API 是 Matrix 客户端-服务器 API 中/sync端点的服务端实现。它从房间服务器(room server)的输出日志中获取数据,负责把"房间发生了什么事"以客户端可消费的形式推送回去。
1.1 当前版本已实现的能力
根据 README 与源码,当前 Sync API 具备以下能力:
- 为携带合法
access_token的用户返回合法的/sync响应。请求经过 syncapi/routing/routing.go 中的MakeAuthAPI认证后进入处理流程。 - 完整同步(complete sync):当请求未提供
since值时,返回完整同步结果,其中包含有效的next_batch令牌;响应覆盖用户被邀请(invited)或已加入(joined)的所有房间,对已加入房间包含完整当前房间状态以及时间线中最新的 20(硬编码)条事件。 - 增量同步(incremental sync):当提供
since值时,被邀请、加入、离开房间等成员关系变化都能在/sync响应中正确反映。 - 大状态增量:对于非常大的状态增量,房间的
state区块会正确填充为时间线起点处的房间状态。 - 加入房间:当
/sync将客户端状态转换为 "joined" 时,按规范包含完整当前房间状态。 - 只唤醒需要唤醒的用户流(详见第三节的 Notifier 机制)。
- 遵循
timeout查询参数。
1.2 从源码看组件装配
组件入口在 syncapi/syncapi.go,AddPublicRoutes完成了:
- 初始化 sync 数据库连接(
storage.NewSyncServerDatasource); - 创建 Typing 缓存与 Notifier;
- 构建九大流提供器(
streams.NewSyncStreamProviders),见 syncapi/streams/streams.go; - 启动多个 JetStream 消费者(房间事件、presence、key change、account data、notification data、typing、send-to-device、receipt 等),见 syncapi/consumers/roomserver.go;
- 通过
routing.Setup注册 HTTP 路由。
二、处理一个/sync请求的两大问题
当服务器收到/sync请求时,需要解决两个问题:
- 计算出要返回哪些房间给客户端;
- 对每个房间,计算出要返回哪些事件给客户端。
2.1 房间选择逻辑(借鉴 Synapse)
"选择哪些房间"的逻辑参考了 Synapse 的实现(README 引用了 Synapse v0.19.3 的synapse/handlers/sync.py),具体步骤为:
- 获取该用户当前已加入的房间列表;
- 获取该用户在给定流位置与当前时刻之间的成员关系变化列表;
- 对每个有成员关系变化的房间:
- 检查房间是否为"新加入"(仅检查 join 事件不足够,因为允许重复 join;判断"新加入"后需下发完整房间状态,且
limited恒为 true); - 检查用户当前是否仍被邀请,是则加入
invited区块; - 检查用户当前是否已离开/被封禁,是则加入
archived区块;
- 检查房间是否为"新加入"(仅检查 join 事件不足够,因为允许重复 join;判断"新加入"后需下发完整房间状态,且
- 追加已加入的房间(joined room list)。
在源码中,这一步体现在 syncapi/streams/stream_pdu.go 的IncrementalSync:调用snapshot.GetStateDeltas(或GetStateDeltasForFullStateSync)获取StateDelta列表,每个StateDelta携带RoomID、StateEvents、NewlyJoined、Membership等字段(见 syncapi/types/types.go),随后按spec.Join、spec.Peek、spec.Leave/spec.Ban分别填充响应的join、peek、leave区块(stream_pdu.go)。
2.2 事件选择逻辑:为什么不能照搬 Synapse
对于每个房间,/sync响应返回最近的时间线事件与时间线起点处的房间状态。"选择哪些事件"的逻辑并非完全基于 Synapse 代码,因为 Synapse 在计算房间状态方面存在已知缺陷。服务器为了知道返回哪些事件,需要在房间历史的多个时间点计算房间状态。
README 用一个 15 事件的房间来说明(字母是状态事件,数字是时间线事件,'表示状态更新):
index 0 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 (基于1的索引,StreamPosition(0) 代表没有事件) timeline [A, B, C, D, 1, 2, 3, D', 4, D'', 5, B', D''', D'''', 6]该房间的当前状态为:[A, B', C, D'''']。
- 若以
?since=14&limit=5请求,只返回 1 条时间线事件(最近的一条):
15 [ 6 ]- 若以
?since=9&limit=5请求,返回 5 条时间线事件(最近的 5 条):
11 12 13 14 15 [5, B', D''', D'''', 6]房间在时间线起点的状态可用两种方式表示:
- 从索引 0 开始的
full_state:[A, B, C, D''](即 0-11 之间的状态); - 从索引 9 开始的部分状态:
[D''](即 9-11 之间的状态)。
2.3 状态事件推进的正确方式
服务器基于状态冲突解决算法推进状态事件(例如从D'推进到D'')。你可能认为只需为每个状态事件的(event type, state_key)元组更新条目即可推进当前状态,但这种推进方式可能与状态冲突解决算法计算出的状态发生分歧。例如,若同一状态键出现两个"同时"更新——即在事件图中处于同一深度的两次更新——则状态冲突解决算法的最终结果可能与事件在时间线中出现的顺序不一致。
状态事件的正确推进由房间服务器输出的OutputRoomEvents中的AddsStateEventIDs与RemovesStateEventIDs表示。这在消费者代码中得到印证:syncapi/consumers/roomserver.go的onNewRoomEvent会基于msg.AddsStateEventIDs与msg.RemovesStateEventIDs调用db.WriteEvent持久化事件与状态增量(roomserver.go)。
三、用户流唤醒机制:Notifier
README 提到"只唤醒需要唤醒的用户流"。对应实现为 syncapi/notifier/notifier.go:
- Notifier 内部维护
roomIDToJoinedUsers(房间→已加入用户集合)与roomIDToPeekingDevices(房间→peeking 设备集合),以及userDeviceStreams(用户→设备→UserDeviceStream); OnNewEvent会根据事件所在房间的 joined 用户列表、peeking 设备列表以及成员事件的目标用户,决定唤醒哪些UserDeviceStream(notifier.go);- 唤醒前先更新当前流位置
currPos,以"先更新位置再唤醒"的顺序避免竞态导致事件遗漏(代码注释明确说明这一设计动机); GetListener返回一个可等待的UserDeviceStreamListener,/sync长轮询正是在此等待新数据(notifier.go)。
长轮询编排在 syncapi/sync/requestpool.go 的OnIncomingSyncRequest中:
shouldReturnImmediately判断:初始同步(无 since)、timeout=0、或full_state=true时立即返回(requestpool.go);- 否则注册监听器并
select等待:客户端取消(Context.Done)、超时(timer.C)或 notifier 唤醒(GetNotifyChannel); - 唤醒后重新计算当前位置,进入下一轮处理;若增量同步没有产生任何更新,则把
since前移到当前流位置继续阻塞等待,避免对"安静的客户端"返回无意义的 no-op 响应(requestpool.go)。
3.1 九大流与StreamingToken
一个/sync响应不仅仅包含房间事件。当前StreamingToken由 9 个独立流位置构成(syncapi/types/types.go):
| 流字段 | 含义 |
|---|---|
PDUPosition | 房间事件(PDU)流 |
TypingPosition | 输入状态(typing)流 |
ReceiptPosition | 已读回执(receipt)流 |
SendToDevicePosition | send-to-device 消息流 |
InvitePosition | 邀请事件流 |
AccountDataPosition | 账户数据流 |
DeviceListPosition | 设备列表变更流 |
NotificationDataPosition | 未读通知数据流 |
PresencePosition | 在线状态(presence)流 |
StreamingToken序列化为s{PDU}_{Typing}_{Receipt}_{SendToDevice}_{Invite}_{AccountData}_{DeviceList}_{NotificationData}_{Presence}形式的next_batch令牌(types.go)。对应的 9 个流提供器均实现StreamProvider接口(CompleteSync/IncrementalSync/Advance/LatestPosition,模板见 syncapi/streams/template_stream.go),在 requestpool.go 中按序调用并汇总成NextBatch。
四、从输出日志到同步数据库:数据消费管线
房间服务器通过 JetStream 输出日志发布事件,Sync API 的消费者负责把事件写入同步数据库并通知 notifier。核心消费者是 syncapi/consumers/roomserver.go 的OutputRoomEventConsumer,其onMessage按OutputEvent.Type分派(roomserver.go):
OutputTypeNewRoomEvent:新房间事件,忽略非自证的红action事件(自 redaction 事件会先忽略,等待OutputTypeRedactedEvent校验后入库),然后调用onNewRoomEvent;OutputTypeOldRoomEvent:旧事件回填(backfill),写入时excludeFromSync标记为ev.StateKey() != nil,避免用旧状态误导客户端(roomserver.go);OutputTypeNewInviteEvent/OutputTypeRetireInviteEvent:邀请的登记与撤销,通知目标用户;OutputTypeNewPeek/OutputTypeRetirePeek:peek 登记/撤销;OutputTypeRedactedEvent:对事件做 redaction 处理并作为普通事件通知客户端;OutputTypePurgeRoom:从同步库中彻底清除整个房间。
onNewRoomEvent的关键链路(roomserver.go):
- 通过
NeededStateEventIDs计算缺失的状态事件,先在本地库查询,再通过rsAPI.QueryEventsByID向房间服务器补充; updateStateEvent为状态事件补全prev_content等 unsigned 字段;db.WriteEvent写入事件并生成 PDU 流位置;writeFTS若启用全文检索则索引事件(消息正文、房间名、话题,m.replace编辑会删除原文索引,见 roomserver.go);pduStream.Advance(pduPos)推进流位置,notifier.OnNewEvent唤醒相关用户。
五、过滤与时间线限制(当前实现)
5.1limit硬编码为 20
README 明确说明:Filters are not honoured or implemented. The limit for each room is hard-coded to 20.——即每个房间的时间线事件数固定为 20。源码中对应常量:
const defaultSyncTimeout = time.Duration(0) const DefaultTimelineLimit = 20见 syncapi/sync/request.go。
5.2 当前已支持的过滤能力
需要说明的是,仓库当前实现已经比 README 撰写时的状态有所演进:虽然 README 列出的已知问题清单中"Filters are not honoured",但当前源码已实现过滤器解析与部分应用:
- 请求解析在 syncapi/sync/request.go:支持
timeout、since、full_state、filter(既支持内联 JSON 过滤器,也支持按filter_id从数据库加载已保存的过滤器); - 完整同步时会为 account data 提升默认 limit(
filter.AccountData.Limit = math.MaxInt32),避免客户端等待数据"涓流式"到达而出现怪异行为(request.go); - 支持
lazy_load_members成员惰性加载:PDUStreamProvider中的lazyLoadMembers只下发时间线中出现用户的成员事件,并配合lazyLoadCache缓存,减少大房间状态体积(stream_pdu.go); - 支持忽略用户(ignored users):
addIgnoredUsersToFilter将m.ignored_user_list中的用户注入NotSenders过滤条件(stream_pdu.go)。
5.3 默认过滤器
synctypes.DefaultFilter()(syncapi/synctypes/filter.go)提供了与 Synapse 对齐的默认值:
EventFormat默认"client"(可选"federation",校验见 filter.go);EventFilter.Limit默认 10;RoomFilter.IncludeLeave默认false;StateFilter.LazyLoadMembers默认false。
六、配置项:sync_api与全文检索
Sync API 的配置定义在 setup/config/config_syncapi.go:
type SyncAPI struct { Matrix *Global `yaml:"-"` Database DatabaseOptions `yaml:"database,omitempty"` RealIPHeader string `yaml:"real_ip_header"` Fulltext Fulltext `yaml:"search"` }database:sync 数据库连接字符串,单数据库模式下默认file:syncapi.db(见Defaults);real_ip_header:指定用于获取真实客户端 IP 的 HTTP 头(如X-Real-IP),在 requestpool.go 的updateLastSeen中使用;search:全文检索配置(Fulltext结构体):enabled(默认 false)、index_path(默认./searchindex)、in_memory(仅测试用)、language(默认"en",用于分词分析)。启用后/search端点可用,否则返回 501(见 routing.go)。
配置示例(dendrite-sample.yaml):
sync_api: # real_ip_header: X-Real-IP search: enabled: false index_path: "./searchindex" language: "en"七、完整同步与增量同步的执行路径
7.1 完整同步(Complete Sync)
PDUStreamProvider.CompleteSync(stream_pdu.go):
- 以
StreamPosition(0)为起点,LatestPosition为终点,Backwards: true构造Range——完整同步从最新事件往前取,保证返回房间内最新的事件; - 查询用户加入的所有房间(
RoomIDsWithMembership(..., spec.Join)),并通过notifier.JoinedUsers失效懒加载缓存中的相关用户; - 对每个房间调用
getJoinResponseForCompleteSync:取当前状态(CurrentState,可排除时间线中已出现的状态事件以避免重复)、计算prev_batch(基于PositionInTopology得到回翻拓扑位置并Decrement)、填充Timeline与State(stream_pdu.go); - 同时处理 peek 房间(
PeeksInRange)写入Rooms.Peek。
7.2 增量同步(Incremental Sync)
PDUStreamProvider.IncrementalSync(stream_pdu.go):
- 根据
since与当前位置构造Range(Backwards: from > to); - 若
WantFullState则用GetStateDeltasForFullStateSync,否则用GetStateDeltas获取状态增量; - 对
NewlyJoined房间反转范围(从最新往回取,受过滤限制),并调用addRoomDeltaToResponse填充响应; - 对每个 delta 更新
newPos,返回给请求池作为新的next_batch。
7.3timeout与full_state语义
getTimeout(request.go)将timeout毫秒数转换为time.Duration;full_state非空且非"false"即视为wantFullState(request.go),并直接导致shouldReturnImmediately返回 true(requestpool.go)。
八、历史可见性(History Visibility)
虽然 README 的"已知问题"中写道m.room.history_visibility不被遵守、一律按shared处理,但当前仓库已实现完整的历史可见性过滤:
- 优先级映射:
world_readable(0) <shared(1) <invited(2) <joined(3),见 syncapi/internal/history_visibility.go; applyHistoryVisibilityFilter会先取时间线中状态事件的当前状态,构造alwaysIncludeIDs,确保最新状态事件不被过滤掉,再调用internal.ApplyHistoryVisibilityFilter(stream_pdu.go);- 每次计算耗时通过 Prometheus 直方图
dendrite_syncapi_calculateHistoryVisibility_duration_millis暴露(history_visibility.go); - 完整同步中对 joined 房间应用该过滤,peek 房间暂不应用(代码注释标明 TODO,见 stream_pdu.go)。
九、其他流提供器的补充职责
除 PDU 流外,各流提供器各司其职,共同构成完整同步响应:
- Invite 流(stream_invite.go):
InviteEventsInRange取范围内的邀请事件构造invite_state(含unsigned.invite_room_state的部分房间状态);增量同步中还会为"被撤销的邀请"生成伪leave事件(使用sha256派生伪 event ID),但仅当用户并非房间现有成员时才下发,避免误将用户移出房间(stream_invite.go); - AccountData 流(stream_accountdata.go):通过
GetAccountDataInRange定位发生变化的 data type,再调用userAPI.QueryAccountData取内容;完整同步时非 joined 房间的 account data 会被跳过;全局 account data 放入顶层account_data,房间级放入对应 join 房间的account_data; - SendToDevice 流(stream_sendtodevice.go):取
SendToDeviceUpdatesForSync的结果写入to_device区块,并跳过忽略用户的发送方;请求池在处理前会先CleanSendToDeviceUpdates清理旧消息,避免同一消息重复下发(requestpool.go); - Presence/Receipt/Typing/DeviceList/NotificationData 流:分别由对应消费者推进位置并由对应
StreamProvider在完整/增量同步时填充presence、ephemeral、device_lists、unread_notifications等区块(各消费者在 syncapi/syncapi.go 中启动)。
十、已知问题与局限(README 原述)
README 明确列出了该实现当前存在的已知问题,理解这些局限有助于判断 Sync API 的能力边界:
m.room.history_visibility不被遵守:一律按shared处理(注:此条为 README 撰写时的状态,当前仓库已实现历史可见性过滤,见上文第八节);- 所有 ephemeral 事件未实现(presence、typing、receipts);
- Account data(用户级与房间级)未实现;
to_device消息未实现;- 通过
prev_batch的回翻分页(back-pagination)未实现; limited标志可能撒谎;- Filters 不被遵守或实现,每个房间的
limit硬编码为 20; full_state查询参数未实现;set_presence查询参数未实现;- "忽略的用户"不会被忽略;
- 被 redaction 的事件仍会发送给客户端;
- 联邦邀请(若存在)无法工作,因为它们不是"真实"事件、不会出现在正确的表中;
invite_state未实现(原因同上);- 当前实现在提供非常旧的
since令牌时扩展性不佳; - 若客户端发送重复的 join 事件(本应为 no-op),整个当前房间状态可能被重新发送给客户端。
这些条目中,部分能力(如 ephemeral 事件、account data、to_device、过滤器、history visibility、ignored users)在当前仓库中已由上文提到的流提供器与过滤器逻辑实现,README 作为历史文档仍保留原始清单;部分限制(如prev_batch回翻分页依赖GetEventsInTopologicalRange实现于/messages端点,完整同步limited的准确性)仍是后续优化的方向。
10.1 从"简单索引"到可扩展性
README 说明:该版本 Sync API 使用非常简单的索引来计算不同时间点的房间状态,当since非常旧或请求full_state时效率低下,因为状态增量会非常大;索引在一定程度上缓解了问题,但未来应采用更好的数据结构。这一表述与GetStateDeltas/GetStateDeltasForFullStateSync的数据库查询路径一致,提醒部署者:对活跃大房间保持合理的同步频率(避免长时间不同步后一次性全量拉取)是缓解该问题的有效手段。
十一、快速上手:验证 Sync API
Dendrite 以 monolith 模式运行时(go build ./cmd/dendrite),Sync API 随整个家服务器一起启动。可通过以下方式验证:
- 使用任意 Matrix 客户端(如 Element Web)登录,观察
/sync请求; - 直接调用接口测试完整同步:
curl -H "Authorization: Bearer <access_token>" \ "https://<your-domain>/_matrix/client/v3/sync?timeout=30000"- 不传
since得到完整同步(含next_batch); - 携带
next_batch作为since即得到增量同步; timeout=30000表示长轮询 30 秒(毫秒单位);filter={"room":{"timeline":{"limit":50}}}可调整时间线条数(当前实现中过滤器的应用范围以源码为准)。
总结
Dendrite 的 Sync API 是一个以"房间选择 + 事件选择"为核心、以 9 路流位置为骨架、以 JetStream 消费与 Notifier 唤醒为动力的同步引擎。理解它需要抓住三条主线:数据从哪来(room server 输出日志 → 消费者 → sync 数据库)、响应怎么算(状态增量 + 时间线 + 状态起点)、请求怎么等(Notifier 长轮询 + 流位置令牌)。README 的已知问题清单则勾勒出该模块的演进边界——部分问题已在当前源码中解决,其余(如非常旧since令牌下的扩展性)仍是值得关注的设计课题。若需深入,可继续阅读 syncapi/README.md、syncapi/sync/requestpool.go、syncapi/streams/stream_pdu.go 与 syncapi/consumers/roomserver.go。
- 后端
- 即时通讯
【免费下载链接】dendrite
Dendrite is a second-generation Matrix homeserver written in Go!
相关推荐
自主托管的Firefox Sync服务器:Run-Your-Own Firefox Sync Server
自主托管的Firefox Sync服务器:Run Your Own Firefox Sync Server ! CircleCI Build Status ht
后端新手必看:hgnetv2_b4.ssld_stage1_in22k_in1k环境配置与Python代码示例
新手必看:hgnetv2_b4.ssld_stage1_in22k_in1k环境配置与Python代码示例 hgnetv2_b4.ssld_stage1_in2
深入解析Shairport Sync与Avahi:零配置网络服务发现原理
深入解析Shairport Sync与Avahi:零配置网络服务发现原理 Shairport Sync作为一款功能强大的AirPlay音频接收器,其核心功能之一
音视频
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考