helix_im/module/trait_impl.rs
1use super::*;
2
3#[derive(serde::Deserialize)]
4struct RuntimeIdentityPatch {
5 auth_user_id: Option<String>,
6 company_id: Option<String>,
7}
8
9impl Module for ImModule {
10 fn name(&self) -> &'static str {
11 "helix-im"
12 }
13
14 /// 路由判定:IM 模块处理所有 IM Inbound 帧和 IM 命令。
15 /// PortReply / Timer 不走此方法(通过 corr_map 定向路由)。
16 ///
17 /// MV3-G02e 草稿族(`im_save_draft` / `im_query_draft`)在此与下面的 `handle`
18 /// **同源放行**(`crate::draft::is_draft_command`):accepts 与 handle 各写一份命令名
19 /// 正是「query 族被 accepts 闸静默丢弃」旧事故的成因,本族不重复该反模式。
20 fn accepts(&self, tick: &Tick) -> bool {
21 if let Tick::Command(cmd) = tick {
22 if cmd.name.as_ref() == crate::recent_history::COMMAND
23 || crate::draft::is_draft_command(cmd.name.as_ref())
24 {
25 return true;
26 }
27 }
28 acl::from_tick::accepts_tick(tick)
29 }
30
31 /// 处理一个 Tick,零 I/O,零 await。
32 fn handle(&mut self, tick: &Tick, now_ms: u64, out: &mut EffectSink) -> Result<(), CoreError> {
33 let start = out.as_slice().len();
34 self.render_scope.observe_reply(tick);
35 let result = self.handle_scoped_tick(tick, now_ms, out);
36 // 包括提前 return 和错误路径在内,任何跨端 Emit 都必须经过同一个业务闸门。
37 self.render_scope.guard_effects(&self.config, start, out);
38 result
39 }
40
41 /// 启动:加载持久频道归属后才允许其渲染;I/O 仍全部由端口执行。
42 fn on_start(&mut self, out: &mut EffectSink) -> Result<(), CoreError> {
43 let start = out.as_slice().len();
44 let result = self.start_scoped_module(out);
45 self.render_scope.guard_effects(&self.config, start, out);
46 result
47 }
48
49 /// 停止生命周期也不允许已失效公司投影穿过出口。
50 fn on_stop(&mut self, out: &mut EffectSink) -> Result<(), CoreError> {
51 let start = out.as_slice().len();
52 let result = self.stop_scoped_module(out);
53 self.render_scope.guard_effects(&self.config, start, out);
54 result
55 }
56}
57
58impl ImModule {
59 /// 执行业务 Tick;最终输出由 Module 实现统一审查。
60 fn handle_scoped_tick(
61 &mut self,
62 tick: &Tick,
63 now_ms: u64,
64 out: &mut EffectSink,
65 ) -> Result<(), CoreError> {
66 if let Tick::PortReply {
67 corr,
68 outcome: helix_core::tick::PortOutcome::Ok(reply),
69 } = tick
70 {
71 if matches!(
72 self.state.corr_map.get(corr),
73 Some(
74 CorrelationContext::DialogListQuery { .. }
75 | CorrelationContext::ChannelViewSnapshotQuery { .. }
76 | CorrelationContext::ChannelSyncPage { .. }
77 )
78 ) {
79 self.render_scope.observe_rows(reply.0.as_ref());
80 }
81 }
82 self.diagnostics.now_ms = now_ms;
83 let effect_start = out.as_slice().len();
84 let corr_floor = self.next_corr;
85 let recovery_reply = self.observe_sync_reply(tick);
86 match tick {
87 Tick::Inbound(bytes) => {
88 // 入口保留旧容错:坏 JSON inbound 是可丢 WS 噪音,不升级为 CoreError。
89 let Ok(frame) = crate::ws::WsFrame::parse(bytes.as_bytes()) else {
90 return Ok(());
91 };
92 self.dispatch_ws_frame(&frame, now_ms, out)?;
93 }
94
95 Tick::Command(cmd) => match cmd.name.as_ref() {
96 RUNTIME_IDENTITY_COMMAND => {
97 let identity: RuntimeIdentityPatch =
98 serde_json::from_slice(cmd.payload.as_ref()).map_err(|error| {
99 CoreError::ModuleError {
100 module: self.name(),
101 source: Box::new(crate::error::ImError::Parse(error.to_string())),
102 }
103 })?;
104 let auth_user_id = identity
105 .auth_user_id
106 .unwrap_or_else(|| self.config.auth_user_id.clone());
107 let company_id = identity
108 .company_id
109 .unwrap_or_else(|| self.config.company_id.clone());
110 if auth_user_id != self.config.auth_user_id {
111 self.render_scope = Default::default();
112 }
113 if auth_user_id != self.config.auth_user_id
114 || company_id != self.config.company_id
115 {
116 self.scope_hydration.reset();
117 // 原始 HTTP 读体可能不带频道/公司字段;切换身份时显式结束旧请求,
118 // A→B→A 后的迟到回包也不能重新晋升为当前公司的投影。
119 let mut cancelled_reads: Vec<_> = self
120 .state
121 .corr_map
122 .iter()
123 .filter_map(|(corr, context)| match context {
124 CorrelationContext::OutboundReadReply { req_id, .. }
125 | CorrelationContext::OutboundMembersByIds { req_id }
126 | CorrelationContext::OutboundContactCandidates { req_id } => {
127 Some((*corr, req_id.clone()))
128 }
129 _ => None,
130 })
131 .collect();
132 cancelled_reads.sort_unstable_by_key(|(corr, _)| corr.raw());
133 for (corr, req_id) in cancelled_reads {
134 self.state.corr_map.remove(&corr);
135 out.push(crate::read_relay::emit_read_error(&req_id, "CANCELLED"));
136 }
137 self.finish_sync_observation(
138 crate::sync_observation::SyncResult::Cancelled,
139 );
140 self.cancel_category(out);
141 self.cancel_forward_details(None, "CANCELLED", out);
142 self.cancel_forwards(out)
143 .map_err(|e| CoreError::ModuleError {
144 module: self.name(),
145 source: Box::new(e),
146 })?;
147 self.state.reset_message_v3_identity();
148 }
149 self.config.auth_user_id = auth_user_id;
150 self.config.company_id = company_id;
151 // Native may deliver the trusted runtime identity after the WS hello.
152 // The identity reset intentionally invalidates A's recovery window, but a
153 // live B transport must immediately open its own window; otherwise every
154 // startup sync commit is misclassified as ordinary online work and emits
155 // one channel-list event per channel.
156 if self.state.connection_id.is_some()
157 && !self.config.auth_user_id.is_empty()
158 && !self
159 .state
160 .recovery_session
161 .is_active_for(self.config.auth_user_id.as_str())
162 {
163 self.state
164 .recovery_session
165 .begin(self.config.auth_user_id.as_str());
166 }
167 }
168 name if crate::query::forward_detail::is_command(name) => {
169 self.forward_detail_command(name, cmd.payload.as_ref(), out)
170 .map_err(|e| CoreError::ModuleError {
171 module: self.name(),
172 source: Box::new(e),
173 })?;
174 }
175 "im_send_message" => {
176 self.handle_send_message(cmd.payload.as_ref(), now_ms, out)
177 .map_err(|e| CoreError::ModuleError {
178 module: self.name(),
179 source: Box::new(e),
180 })?;
181 }
182 "im_retry_upload" => {
183 crate::send::retry_upload::handle_retry_upload(self, cmd.payload.as_ref(), out)
184 .map_err(|e| CoreError::ModuleError {
185 module: self.name(),
186 source: Box::new(e),
187 })?;
188 }
189 "im_retry_send" => {
190 crate::send::retry_send::handle_retry_send(
191 self,
192 cmd.payload.as_ref(),
193 now_ms,
194 out,
195 )
196 .map_err(|e| CoreError::ModuleError {
197 module: self.name(),
198 source: Box::new(e),
199 })?;
200 }
201 // 连接控制(04 文档表 B `IMWS:reconnect`):helix-im 是 sans-IO,不持 transport——
202 // 只 emit `im:net:reconnect_requested` 控制信号,由 driver 观察后重连 NativeTransport
203 // (连接生命周期归 driver)。payload 无须解析(无参控制信号)。
204 "im_reconnect" => {
205 out.push(crate::acl::to_effect::emit_reconnect_requested());
206 }
207 "im_sync_channels" => {
208 crate::sync::explicit::handle(self, cmd.payload.as_ref(), out).map_err(
209 |e| CoreError::ModuleError {
210 module: self.name(),
211 source: Box::new(e),
212 },
213 )?;
214 }
215 "im_create_posts" => {
216 crate::forward::start(self, cmd.payload.as_ref(), now_ms, out).map_err(
217 |e| CoreError::ModuleError {
218 module: self.name(),
219 source: Box::new(e),
220 },
221 )?;
222 }
223 // MV3-G02e 草稿域:Effect 仅 Persist、零领域事件、结构化 Command Result。
224 // 必须排在 outbound / query 两个 guard 之前认领——草稿既不是 HTTP outbound,
225 // 也不在 query 闭集内,落到 `_` 分支就会退回「unhandled command」死代码态。
226 crate::recent_history::COMMAND => {
227 crate::recent_history::handle(self, cmd.payload.as_ref(), out).map_err(
228 |e| CoreError::ModuleError {
229 module: self.name(),
230 source: Box::new(e),
231 },
232 )?;
233 }
234 name if crate::draft::is_draft_command(name) => {
235 crate::draft::handle_command(self, name, cmd.payload.as_ref(), now_ms, out)
236 .map_err(|e| CoreError::ModuleError {
237 module: self.name(),
238 source: Box::new(e),
239 })?;
240 }
241 // 分类接龙独立authority,不复用文字接龙字段或确认规则。
242 name if crate::category_chain::is_command(name) => {
243 self.handle_category_command(name, cmd.payload.as_ref(), now_ms, out)
244 .map_err(|e| CoreError::ModuleError {
245 module: self.name(),
246 source: Box::new(e),
247 })?;
248 }
249 // 文字接龙命令拥有独立的 correlation/authority barrier,不能落入通用 fire-and-forget。
250 name if crate::chain::is_command(name) => {
251 self.handle_chain_command(name, cmd.payload.as_ref(), out)
252 .map_err(|e| CoreError::ModuleError {
253 module: self.name(),
254 source: Box::new(e),
255 })?;
256 }
257 // Tier 1 outbound 自驱动命令(read/revoke/leave/schedule/cancel/create/makeTopic):
258 // 纯 HTTP-fire,业务契约在 crate::commands(endpoint 锚现网真源)。
259 name if crate::commands::is_outbound(name) => {
260 if name == "im_create_schedule" {
261 match self.begin_schedule_media(cmd.payload.as_ref(), now_ms, out) {
262 Ok(true) => return Ok(()),
263 Ok(false) => {}
264 Err(source) => {
265 // 校验失败也必须结束关联等待,不能让界面等到 Host 超时。
266 if let Ok(payload) = serde_json::from_slice::<serde_json::Value>(
267 cmd.payload.as_ref(),
268 ) {
269 if let Some(req_id) = payload["req_id"].as_str() {
270 out.push(crate::read_relay::emit_read_error(
271 req_id,
272 "SCHEDULE_CONTENT_INVALID",
273 ));
274 return Ok(());
275 }
276 }
277 return Err(CoreError::ModuleError {
278 module: self.name(),
279 source: Box::new(source),
280 });
281 }
282 }
283 }
284 let corr = self.alloc_corr_internal();
285 let topic_request = (name == "im_make_topic")
286 .then(|| make_topic_request_context(cmd.payload.as_ref()))
287 .flatten();
288 let effects = match crate::commands::handle_outbound(
289 name,
290 cmd.payload.as_ref(),
291 self.config.api_base_url.as_str(),
292 // 第二网关 base(spec06 方案 A):vote/score 命令走它,既有命令忽略。
293 self.config.default_api_base_url.as_str(),
294 self.state.connection_id.as_deref(),
295 corr,
296 ) {
297 Ok(effects) => effects,
298 Err(error) if name == "im_make_topic" => {
299 if let Some((root_message_id, req_id, _)) = topic_request.as_ref() {
300 out.push(crate::acl::to_effect::emit_topic_operation_status(
301 req_id,
302 root_message_id,
303 "engine-rejected",
304 None,
305 Some("invalid-request"),
306 true,
307 ));
308 return Ok(());
309 }
310 return Err(CoreError::ModuleError {
311 module: self.name(),
312 source: Box::new(error),
313 });
314 }
315 Err(error) => {
316 return Err(CoreError::ModuleError {
317 module: self.name(),
318 source: Box::new(error),
319 });
320 }
321 };
322 // spec06 缺陷A:读族注册回灌上下文(PortReply 透传 `im:read:result{req_id,body}`;写族
323 // fire-and-forget 不注册)。S7:byIds 成员快照走 OutboundMembersByIds(额外 emit im:channel:members)。
324 if crate::commands::is_read(name) {
325 let explicit_req_id = crate::read_relay::read_req_id(cmd.payload.as_ref());
326 let increment_hydration = name == "im_channel_load_increment_by_channel_id";
327 if increment_hydration
328 && (self.config.auth_user_id.is_empty()
329 || self.config.company_id.is_empty())
330 {
331 return Err(CoreError::ModuleError {
332 module: self.name(),
333 source: Box::new(crate::ImError::Parse(
334 "increment hydration requires RuntimeAuth account and tenant"
335 .to_string(),
336 )),
337 });
338 }
339 // 读族只有明确 correlation 才允许回灌;hydration 兼容旧 driver 的同时,
340 // 由显式 transport req_id 决定是否发布 sender-only channelIncrement。
341 let emit_channel_increment =
342 increment_hydration && explicit_req_id.is_some();
343 let req_id = match explicit_req_id {
344 Some(req_id) => req_id,
345 None if increment_hydration => {
346 format!("increment-hydration-{}", corr.raw())
347 }
348 None => {
349 // 无 correlation 的旧读请求仍执行 HTTP,但回包不得进入任何 UI 投影。
350 for effect in effects {
351 out.push(effect);
352 }
353 return Ok(());
354 }
355 };
356 {
357 let ctx = if increment_hydration {
358 CorrelationContext::OutboundIncrementHydration {
359 req_id,
360 emit_channel_increment,
361 scope_channel: None,
362 }
363 } else if name == "im_channels_members_by_ids" {
364 CorrelationContext::OutboundMembersByIds { req_id }
365 } else if name == "im_user_candidates" {
366 CorrelationContext::OutboundContactCandidates { req_id }
367 } else if name == "im_get_posts_after_index" {
368 // G11h initial-window 独立于 locate/exact:HTTP 成功必须先过 durable read-back。
369 let args = serde_json::from_slice::<serde_json::Value>(
370 cmd.payload.as_ref(),
371 )
372 .map_err(|error| {
373 CoreError::ModuleError {
374 module: self.name(),
375 source: Box::new(crate::ImError::Parse(error.to_string())),
376 }
377 })?;
378 if let Some(request) =
379 crate::outbound::posts::read::initial_window_request(
380 &args, name,
381 )
382 .map_err(|error| {
383 CoreError::ModuleError {
384 module: self.name(),
385 source: Box::new(error),
386 }
387 })?
388 {
389 CorrelationContext::OutboundInitialWindow {
390 req_id: request.req_id,
391 post_id: request.post_id,
392 page_size: request.page_size,
393 }
394 } else {
395 let channel_id = args
396 .get("channel_id")
397 .and_then(serde_json::Value::as_str)
398 .map(str::to_owned);
399 CorrelationContext::OutboundReadReply {
400 req_id,
401 command: name.to_string(),
402 channel_id,
403 }
404 }
405 } else if name == "im_get_posts" {
406 // G11g exact 查询独立于旧 locate 终态;locate 只保留兼容入参,不改变按 id 读回。
407 let requested_ids = serde_json::from_slice::<serde_json::Value>(
408 cmd.payload.as_ref(),
409 )
410 .ok()
411 .and_then(|value| {
412 crate::outbound::posts::read::exact_post_ids(&value, name).ok()
413 })
414 .unwrap_or_default();
415 CorrelationContext::OutboundExactPosts {
416 req_id,
417 requested_ids,
418 }
419 } else if name == "im_get_replies" || name == "im_get_reply_branch" {
420 // G12 回复族冻结请求语义,PortReply 只生成 typed thread event。
421 let request =
422 crate::render_ready_replies::ReplyProjectionRequest::from_command(
423 name,
424 cmd.payload.as_ref(),
425 corr.raw(),
426 self.config.auth_user_id.as_str(),
427 )
428 .ok_or_else(|| CoreError::ModuleError {
429 module: self.name(),
430 source: Box::new(crate::ImError::Parse(
431 "reply command requires a valid req_id and payload"
432 .to_string(),
433 )),
434 })?;
435 if request.mode
436 == crate::render_ready_replies::ReplyProjectionMode::Snapshot
437 && !request.root_hint.is_empty()
438 {
439 self.state
440 .reply_projection_revisions
441 .insert(request.root_hint.clone(), request.revision);
442 }
443 CorrelationContext::OutboundReplies { request }
444 } else if name == "im_post_read_list" {
445 CorrelationContext::OutboundPostReaders { req_id }
446 } else if name == "im_channel_load_post_pinned" {
447 let channel_id = crate::query::pinned_projection::parse_channel_id(
448 cmd.payload.as_ref(),
449 )
450 .map_err(|error| CoreError::ModuleError {
451 module: self.name(),
452 source: Box::new(error),
453 })?;
454 let account_id = self.config.auth_user_id.clone();
455 let projection_key =
456 crate::query::pinned_projection::projection_key(
457 account_id.as_str(),
458 channel_id,
459 )
460 .map_err(|error| {
461 CoreError::ModuleError {
462 module: self.name(),
463 source: Box::new(error),
464 }
465 })?;
466 let epoch = self
467 .state
468 .pinned_projection_epochs
469 .get(&channel_id)
470 .copied()
471 .unwrap_or(0);
472 CorrelationContext::OutboundPinnedReply {
473 req_id,
474 account_id,
475 channel_id,
476 projection_key,
477 epoch,
478 }
479 } else {
480 let channel_id = serde_json::from_slice::<serde_json::Value>(
481 cmd.payload.as_ref(),
482 )
483 .ok()
484 .and_then(|value| {
485 value
486 .get("channel_id")
487 .and_then(serde_json::Value::as_str)
488 .map(str::to_owned)
489 });
490 CorrelationContext::OutboundReadReply {
491 req_id,
492 command: name.to_string(),
493 channel_id,
494 }
495 };
496 self.state.corr_map.insert(corr, ctx);
497 }
498 } else if name == "im_create_schedule" {
499 if let Ok(args) =
500 serde_json::from_slice::<serde_json::Value>(cmd.payload.as_ref())
501 {
502 if let Some(channel_id) = args
503 .get("channel_id")
504 .and_then(serde_json::Value::as_str)
505 .and_then(ChannelId::from_str)
506 {
507 let request_id = args
508 .get("req_id")
509 .and_then(serde_json::Value::as_str)
510 .filter(|value| !value.is_empty())
511 .map(str::to_string);
512 if let Some(request_id) = request_id.as_ref() {
513 self.state
514 .pending_schedule_requests
515 .insert(channel_id, request_id.clone());
516 }
517 self.state.corr_map.insert(
518 corr,
519 CorrelationContext::OutboundScheduleCreate {
520 channel_id,
521 request_id,
522 },
523 );
524 }
525 }
526 } else if name == "im_cancel_schedule" {
527 if let Ok(args) =
528 serde_json::from_slice::<serde_json::Value>(cmd.payload.as_ref())
529 {
530 if let Some(channel_id) = args
531 .get("channel_id")
532 .and_then(serde_json::Value::as_str)
533 .and_then(ChannelId::from_str)
534 {
535 let request_id = args
536 .get("req_id")
537 .and_then(serde_json::Value::as_str)
538 .filter(|value| !value.is_empty())
539 .map(str::to_string);
540 if let Some(request_id) = request_id.as_ref() {
541 self.state
542 .pending_schedule_cancel_requests
543 .insert(channel_id, request_id.clone());
544 }
545 self.state.corr_map.insert(
546 corr,
547 CorrelationContext::OutboundScheduleCancel {
548 channel_id,
549 request_id,
550 },
551 );
552 }
553 }
554 } else if name == "im_create_channel" {
555 if let Ok(args) =
556 serde_json::from_slice::<serde_json::Value>(cmd.payload.as_ref())
557 {
558 let request_id = args
559 .get("req_id")
560 .and_then(serde_json::Value::as_str)
561 .filter(|value| !value.is_empty())
562 .map(str::to_string);
563 if let Some(members) = args
564 .get("user_ids")
565 .and_then(serde_json::Value::as_array)
566 .map(|user_ids| {
567 user_ids
568 .iter()
569 .filter_map(serde_json::Value::as_str)
570 .map(|user_id| serde_json::json!({ "id": user_id }))
571 .collect::<Vec<_>>()
572 })
573 .filter(|members| !members.is_empty())
574 {
575 self.state.corr_map.insert(
576 corr,
577 CorrelationContext::OutboundChannelCreate {
578 members,
579 request_id,
580 },
581 );
582 }
583 }
584 } else if name == "im_make_topic" {
585 if let Some((root_message_id, req_id, display_name)) = topic_request {
586 self.state.corr_map.insert(
587 corr,
588 CorrelationContext::OutboundMakeTopic {
589 root_message_id: root_message_id.clone(),
590 req_id: req_id.clone(),
591 display_name,
592 },
593 );
594 out.push(crate::acl::to_effect::emit_topic_operation_status(
595 &req_id,
596 &root_message_id,
597 "remote-pending",
598 None,
599 None,
600 false,
601 ));
602 }
603 } else if name == "im_channel_change_info"
604 || name == "im_channel_change_notice"
605 || name == "im_channel_change_display_name"
606 || name == "im_channel_change_orient"
607 || name == "im_channel_change_permission"
608 {
609 if let Ok(value) =
610 serde_json::from_slice::<serde_json::Value>(cmd.payload.as_ref())
611 {
612 if let Some(channel_id) = value
613 .get("channel_id")
614 .and_then(|id| id.as_str())
615 .and_then(ChannelId::from_str)
616 {
617 let causation_id = value
618 .get("req_id")
619 .and_then(serde_json::Value::as_str)
620 .filter(|value| !value.is_empty())
621 .map(str::to_string);
622 self.state.corr_map.insert(
623 corr,
624 CorrelationContext::OutboundChannelSettings {
625 channel_id,
626 causation_id,
627 },
628 );
629 }
630 }
631 } else if name == "im_channel_leave" {
632 if let Some(channel_id) =
633 serde_json::from_slice::<serde_json::Value>(cmd.payload.as_ref())
634 .ok()
635 .and_then(|value| {
636 value
637 .get("channel_id")
638 .and_then(|id| id.as_str())
639 .map(str::to_string)
640 })
641 .and_then(|id| ChannelId::from_str(&id))
642 {
643 self.state
644 .corr_map
645 .insert(corr, CorrelationContext::OutboundLeave { channel_id });
646 }
647 }
648 for eff in effects {
649 out.push(eff);
650 }
651 }
652 // 投影/状态命令:最近消息由 helix-im 统一编排 local-first/远端 fallback。
653 name if crate::query::is_query(name) => {
654 self.handle_query_command(name, cmd.payload.as_ref(), out)
655 .map_err(|e| CoreError::ModuleError {
656 module: self.name(),
657 source: Box::new(e),
658 })?;
659 }
660 _ => tracing::debug!("im: unhandled command '{}'", cmd.name),
661 },
662
663 Tick::PortReply { corr, outcome } => {
664 let result = self.handle_port_reply(*corr, outcome, now_ms, out);
665 if recovery_reply && result.is_err() {
666 self.finish_sync_observation(crate::sync_observation::SyncResult::Failed);
667 }
668 result.map_err(|e| CoreError::ModuleError {
669 module: self.name(),
670 source: Box::new(e),
671 })?;
672 }
673
674 Tick::PortProgress { corr, progress } => {
675 self.handle_file_upload_progress(*corr, *progress, out)
676 .map_err(|e| CoreError::ModuleError {
677 module: self.name(),
678 source: Box::new(e),
679 })?;
680 }
681
682 Tick::Timer(id) => {
683 if self.forward_detail_timeout(*id, out) {
684 return Ok(());
685 }
686 if self
687 .forward_timeout(*id, out)
688 .map_err(|e| CoreError::ModuleError {
689 module: self.name(),
690 source: Box::new(e),
691 })?
692 {
693 return Ok(());
694 }
695 if self
696 .handle_category_timeout(*id, out)
697 .map_err(|e| CoreError::ModuleError {
698 module: self.name(),
699 source: Box::new(e),
700 })?
701 {
702 return Ok(());
703 }
704 if self.schedule_media_timeout(*id, out) {
705 return Ok(());
706 }
707 self.handle_timer(*id, now_ms, out)
708 .map_err(|e| CoreError::ModuleError {
709 module: self.name(),
710 source: Box::new(e),
711 })?;
712 }
713
714 Tick::Connected(_transport) => {
715 // A4:连接建立 ≠ 可用。现网 Go 服务端连上后**单向下发 hello 帧**;
716 // 必须等收到 hello(→ Connected + connectionId)才能 resync——
717 // 否则握手前发 sync/notify 会被 Go 401/403(缺 connectionId 身份头)。
718 // 故此处只置 Connecting,resync 推迟到 hello Inbound 分支(handle_hello)。
719 // Native reader 与 activation tick 来自两个异步源;hello 可能先入队。
720 // connectionId 是服务端握手的权威凭据,晚到的 transport tick 不得把
721 // 已握手状态从 Connected 降回 Connecting。
722 if self.state.connection_id.is_none() {
723 self.state.conn = ConnState::Connecting;
724 }
725 // Do not publish a runtime transition here: the server hello
726 // establishes the recovery epoch, and PersistOk is the only
727 // completion boundary for a V2 frame.
728 self.state.recovery_session.phase = crate::sync_session::RecoveryPhase::Comparing;
729 tracing::debug!("helix-im: transport connected, awaiting hello handshake");
730 }
731
732 Tick::Disconnected(_transport) => {
733 self.scope_hydration.reset();
734 self.finish_sync_observation(crate::sync_observation::SyncResult::Cancelled);
735 self.state.increment_pull = None;
736 self.state.corr_map.retain(|_, context| {
737 !matches!(
738 context,
739 CorrelationContext::IncrementPullHttp
740 | CorrelationContext::IncrementPullPersist
741 | CorrelationContext::OutboundIncrementHydration {
742 scope_channel: Some(_),
743 ..
744 }
745 )
746 });
747 self.state.conn = ConnState::Disconnected;
748 self.state.recovery_session.phase = crate::sync_session::RecoveryPhase::Failed;
749 self.state.reset_transport_query_session();
750 // 清除所有 channel 的 inflight_sync(连接断开,所有在途 Http 不会回报)
751 // 若不清除,重连后 proactive resync 会被 B1 守卫全部跳过
752 for ch in self.state.channels.values_mut() {
753 ch.inflight_sync = None;
754 ch.last_sync_from_seq = None;
755 }
756 // 断线清理(语义保真 = 旧 corr_to_sync.clear() + continuation_pending.clear()):
757 // 1. 清除 SyncPull(在途 sync Http corr,断线后不会再收到 PortReply)。
758 // 2. ChannelPersist 本体**保留**、仅把 wants_continuation 置 false——
759 // 对应的 Effect::Persist 是本地 SQLite 写(driver spawn_blocking 执行),
760 // 其 PortReply 与 WS 断线物理无关,断线后极可能照常完成并回报,
761 // 留着可继续推进 cursor(HX-C008 cursor 单调推进);只是不再续拉 sync。
762 // 旧代码 continuation_pending.clear() 仅清续拉意图、故意不清 corr_to_channel
763 // 本体(module.rs 路由1 仍能推 cursor),此处逐字节保真。
764 // 3. OptimisticSend / ScanCursors / ScanChannelProjections 不动(本地回报与连接无关)。
765 for ctx in self.state.corr_map.values_mut() {
766 if let CorrelationContext::ChannelPersist {
767 wants_continuation, ..
768 } = ctx
769 {
770 *wants_continuation = false;
771 }
772 }
773 self.state
774 .corr_map
775 .retain(|_, ctx| !matches!(ctx, CorrelationContext::SyncPull { .. }));
776 tracing::debug!("helix-im: disconnected, cleared inflight sync state");
777 }
778 }
779 self.observe_sync_tick(&out.as_slice()[effect_start..], corr_floor);
780 Ok(())
781 }
782
783 /// 启动:先 Scan `channel_event_cursor` 与 `channel` 投影,过滤终态后才触发 proactive sync。
784 fn start_scoped_module(&mut self, out: &mut EffectSink) -> Result<(), CoreError> {
785 self.diagnose(crate::diagnostics::Observation {
786 event: "login_sync_started",
787 stage: "store",
788 result: "started",
789 ..Default::default()
790 });
791 self.start_lifecycle(out)
792 .map_err(|e| CoreError::ModuleError {
793 module: self.name(),
794 source: Box::new(e),
795 })?;
796 self.state.media_recovery_ready = false;
797 self.state.pending_media_recovery_compensation = None;
798 let corr = self.alloc_corr_internal();
799 out.push(helix_core::Effect::Persist {
800 corr,
801 ops: vec![helix_core::effect::StorageOp::Scan(
802 helix_core::effect::ScanSpec {
803 table: "pending_media",
804 limit: None,
805 filter: None,
806 order_by: &[],
807 },
808 )],
809 });
810 self.state.pending_media_rehydrate_corr = Some(corr);
811 Ok(())
812 }
813
814 /// 停止:取消所有 timer + 关闭连接
815 fn stop_scoped_module(&mut self, out: &mut EffectSink) -> Result<(), CoreError> {
816 self.finish_sync_observation(crate::sync_observation::SyncResult::Cancelled);
817 self.cancel_category(out);
818 self.stop_schedule_media(out);
819 self.cancel_forward_details(None, "CANCELLED", out);
820 self.cancel_forwards(out)
821 .map_err(|e| CoreError::ModuleError {
822 module: self.name(),
823 source: Box::new(e),
824 })?;
825 self.diagnose(crate::diagnostics::Observation {
826 event: "runtime_stopped",
827 result: "interrupted",
828 reason: "runtime_stopped",
829 ..Default::default()
830 });
831 self.stop_lifecycle(out)
832 .map_err(|e| CoreError::ModuleError {
833 module: self.name(),
834 source: Box::new(e),
835 })
836 }
837}
838
839/// 冻结 make-topic 的 root、request 与用户确认标题,供异步 authority/persist/readback 串联。
840fn make_topic_request_context(payload: &[u8]) -> Option<(String, String, String)> {
841 let value = serde_json::from_slice::<serde_json::Value>(payload).ok()?;
842 let root_message_id = value
843 .get("root_id")
844 .and_then(serde_json::Value::as_str)
845 .unwrap_or_default()
846 .to_string();
847 let req_id = value
848 .get("req_id")
849 .and_then(serde_json::Value::as_str)
850 .filter(|id| !id.is_empty())?
851 .to_string();
852 let display_name = value
853 .get("display_name")
854 .and_then(serde_json::Value::as_str)
855 .filter(|name| !name.is_empty())
856 .unwrap_or("话题")
857 .to_string();
858 Some((root_message_id, req_id, display_name))
859}