Skip to main content

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}