diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index 4686a4414b..81aafa6dda 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -28,6 +28,26 @@ Five questions organize the experience: is my request still here; who is actuall working; did my correction or stop take effect; where is the checked result; and how do I come back after failure without starting the work again? +### Realtime Bot entry and recipient purpose + +A native Bot replacement is another entry to this conversation lifecycle. Its +realtime connection is independent of periodic Goal work. Entry and recipient +purpose are separate: ordinary project chat, direct conversation with an existing +Agent, and the persistent steward share mechanics but have different objectives +and grants. The [steward operational contract](capable-manager-semantic-handoff-v0.md#10-operational-contract) +orders transport isolation, ordinary DM/role choice, progress/media/permissions +and installed replacement qualification under S5. + +Ordinary project chat needs a shared Core conversation context whose workspace, +executor and audience are explicitly authorized, without a user-created Goal or +an automatic global-steward objective. This is a remaining entry requirement, +not a new shipped Session schema. Lark must not implement it by creating hidden +Goals, copying another host's sessions, or introducing an independent executor. +Explicit recipient selection uses permitted stable references; labels do not +confer grants. Switching the selected recipient affects future input, while +accepted work and returns retain their original Session, source and audience. +Stop targets the exact current request rather than every Agent behind a Bot. + ### Managed and attached are different execution relationships | Relationship | App promise | Required evidence and limit | diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md index 90cf619289..4740777d85 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md @@ -24,6 +24,20 @@ Goal/Todo、lease、quota、验收和 effect 各自保留原有权限归属。 体验围绕五个问题组织:请求还在吗;谁确实在工作;纠偏或停止生效了吗; 核验过的结果在哪里;失败后如何回来且不重新启动工作? +### 实时 Bot 入口与接收者职责 + +原生 Bot 替换是这套对话生命周期的另一个入口,实时连接独立于 Goal 周期工作。 +入口与接收者职责分开:普通项目对话、与既有 Agent 直接交流、持久管家共享机制, +但目标与授权不同。[管家运行契约](capable-manager-semantic-handoff-v0.zh-CN.md#10-运行契约) +在 S5 下排列传输隔离、普通私聊/角色选择、进度/媒体/权限及安装态替换验收。 + +普通项目对话需要共享 Core 会话上下文,明确授权工作区、executor 和受众, +无需用户创建 Goal,也不自动赋予全局管家目标。这是尚待实现的入口要求, +不是新增已发布 Session schema。Lark 不能通过创建隐藏 Goal、复制其它 host session +或另建 executor 来实现。显式接收者选择使用权限范围内的稳定引用,标签不授予权限。 +切换接收者影响未来输入;已受理工作和回报保留原 Session、来源和受众。 +停止针对当前精确请求,不停止 Bot 后面的全部 Agent。 + ### Managed 与 attached 是不同的执行关系 | 关系 | App 承诺 | 必需证据与限制 | diff --git a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md index 5b02283204..6e1b05a3f9 100644 --- a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md +++ b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md @@ -752,6 +752,42 @@ Target an ingress receipt within two seconds on a healthy local service, indepen Use existing service recovery and receipt pumps. No manager-specific business automation for each kind of request. Expose configuration and failures through the existing CLI, capability settings and manager conversation. Troubleshooting distinguishes model failure, tool/policy denial, state conflict, unreachable receiver and transport formatting/delivery failure. +**Realtime IM entry (S1/S5/S10):** qualify long-connection reception, durable +admission, host execution and visible return separately. While one answer is +blocked, another request or correction must receive bounded, truthful admission +feedback; Core queue/steering support alone does not qualify the Lark consumer. +Retain persistent Sessions, provisional progress, exact permission decisions +and explicit media availability through existing Chat/Turn/operation owners. +No parallel bridge ledger or scheduler is required. The bundled provider's +[readiness guide](../../../loopx/extensions/lark/docs/realtime-conversation-readiness.md) +records bounded cross-conversation dispatch and remaining replacement qualification; +private DM onboarding, ordinary non-Goal chat, streaming and media remain +unqualified until the pinned installed journey passes. + +**Product boundary:** the Bot is a realtime conversation entry, not another +steward. Ordinary project chat and direct conversation with a selected existing +Agent must remain useful without team decomposition or a new Goal/Todo. The +steward is an explicit recipient when the user needs persistent commitments, +coordination and acceptance. Share authenticated ingress, Session/Turn, +execution/progress/attachments, operation decisions and recovery; keep recipient +purpose, audience, transcript and workspace grants distinct. A channel must not +inherit the global steward objective or portfolio visibility just to obtain a +working executor. Shared changes belong to the +[conversation-entry RFC](app-conversation-and-async-inbox-v0.md), not a second Bot +execution or approval authority. Replacing a tenant-controlled App on another +machine requires a freshly authorized App; credentials, source transcripts and +permission bindings do not travel as a workstation backup. + +Deliver these through M1/M3 and the existing S5 journey, without adding a parallel +milestone or treating a periodic heartbeat as realtime transport: + +| Order | User-visible exit | Existing owner and qualification | +| --- | --- | --- | +| First | A slow role does not hold every other conversation on the same Bot; follow-ups remain ordered | Lark transport has bounded workers/buffering, fresh binding checks, reply/ACK readback and stop/drain evidence; Core retains Session admission and budgets | +| Next | A first DM and explicit role choice continue the intended conversation; busy work promptly reports durable admission or rejection | Chat Session/Turn and typed ingress own continuity, audience and queue/steering; ordinary chat must not require users to manufacture a Goal Topic, and role names alone grant no authority | +| Then | Progress, images/files and permission answers work for each advertised host | Existing event/attachment/operation owners; bounded provisional cards, explicit unsupported media, authenticated exact-operation callbacks and delivery-only recovery | +| Switch gate | The installed provider/host journey survives reconnect, duplicates, cancellation and unavailable delivery | Pin versions and run the actual entry/readback; retire an old bridge only after qualification, with one consumer owner per App and no copied credentials or sessions | + **Accepted queue preparation failures (S1/S10, A12/A22/A23):** an accepted request owns a terminal outcome even before an adapter starts. A missing runtime asset, invalid workspace or failed session restoration must settle the affected diff --git a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md index 306bf39807..50c067b5e7 100644 --- a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md +++ b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md @@ -632,6 +632,31 @@ M0 盘点真实字段和 producer;以下是迁移验收底线,不代表已 使用现有服务恢复和 receipt pump,不为每类请求创建管家业务 automation。配置/故障通过已有 CLI、capability settings、管家对话展示。诊断区分模型失败、工具/策略拒绝、状态冲突、接收方不可达、格式/传输失败。 +**实时 IM 入口(S1/S5/S10):** 分开验收长连接收信、持久准入、host 执行和可见回报。 +一个回答阻塞时,另一个请求或纠正仍须得到有界、真实的准入反馈;Core 已支持 queue/steering +不能替代 Lark consumer 的验证。持续 Session、临时进度、精确权限决策和明确的媒体可用性 +复用已有 Chat/Turn/operation owner,不另建 bridge 账本或 scheduler。 +bundled provider 的[就绪指南](../../../loopx/extensions/lark/docs/realtime-conversation-readiness.md) +记录有界跨会话处理和仍待完成的替代验收;私聊开通、普通非 Goal 对话、流式回显与媒体输入, +须等固定版本的已安装旅程通过后才能宣称就绪。 + +**产品边界:** Bot 是实时会话入口,不是另一个管家。普通项目对话和与选定既有 Agent +直接交流,无需团队拆解或新增 Goal/Todo;需要长期承诺、协调和验收时,管家是显式接收者。 +共享认证入口、Session/Turn、执行/进度/附件、operation 决策与恢复;接收者职责、受众、 +会话记录和工作区授权保持独立。不能为了获得可用 executor,就让渠道继承全局管家目标或 +portfolio 可见性。共享改动归[对话入口 RFC](app-conversation-and-async-inbox-v0.zh-CN.md), +不另建 Bot 执行或审批权威。跨主机替换受租户控制的 App 时,使用重新授权的新 App; +凭据、来源会话和权限绑定不作为装机备份迁移。 + +按 M1/M3 和已有 S5 旅程交付,不新增平行里程碑,不把周期 heartbeat 当实时传输: + +| 次序 | 用户可感知出口 | 既有 owner 与验收 | +| --- | --- | --- | +| 先做 | 同一 Bot 的慢角色不拖住其它会话,同会话追问保持顺序 | Lark 传输有界 worker/缓冲、执行前绑定复核、reply/ACK 回读和停止/drain 证据;Core 保留 Session 准入与预算 | +| 接着 | 首次私聊与显式角色选择接续正确会话,忙碌工作及时报告持久准入或拒绝 | Chat Session/Turn 与 typed ingress 拥有连续性、受众及 queue/steering;普通对话无需用户先造 Goal Topic,角色名本身不授予权限 | +| 再做 | 每个宣称支持的 host 都能显示进度、接收图/文件和处理权限答复 | 既有 event/attachment/operation owner;有界临时卡片、明确不支持的媒体、认证后的精确 operation 回调,以及只重试投递的恢复 | +| 切换门槛 | 安装态 provider/host 旅程经得住重连、重复事件、取消和投递不可用 | 固定版本,跑实际入口和回读;验收后才退役旧 bridge,每个 App 保持唯一 consumer,不复制凭据或 session | + ## 11. 规范性里程碑 以完整用户旅程交付,不按零散字段拆 PR。管家工程负责人维护 canonical Todo 和私有 incident→验收映射;PR 引用本 RFC 的里程碑及验收 ID。公开进度只含可公开结果。完成需要当前部署证据,不是合并 PR 数。 diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index 08022d753f..e294283bd8 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -74,12 +74,12 @@ P0 blocks correctness or continuity in the current user journey. P1 enables repe | **S2 Typed kernel and durable authority · P0/P1** | Effect/Todo/quota/recovery owners, TS migration and store candidates exist; writer cutover/provider promotion remain incomplete | Migrate one real transaction/recovery lifecycle at a time, with semantic counterexamples before cutover/deletion. R1 correctness precedes migration volume. Require real backend, concurrency/fence, ambiguous commit, retention/export recovery, bridge costs and D1–D3 evidence | | **S3 Goal planning and multi-Agent collaboration · P0/P1** | Vision/replan, peer frontiers, claim/lease, directory, manager_context and explicit continuation exist; general handoff/shared amendment remain incomplete | R2 proves peer dependency; R3 closes parallel joins, pipelines, help/review, continuation and automatic return; R4 delivers one intent-preserving amendment class. Cover cycles, invalidated inputs, rejection/deferral, lease transfer, competing bases and aggregate acceptance | | **S4 Runtime/host/daemon · P0/P1** | Attached/managed, Turn, broker, runtime connectors and Desktop repairs exist; registration does not establish executable capacity | Qualify multi-Turn supervision for one real supported combination; restart/cancel/drain/stop retain work and fence old executors. Then expand host parity, unique service-profile ownership, clean installation and upgrades; show unsupported adapter capabilities | -| **S5 Frontend, Lark and human interaction · P0/P1** | Local chat, settings, proposals and partial Goal Channel verticals exist; shared audience/session/work readback needs qualification | One journey spans settings, work graph, handoff, blockers, cost, corrections, artifacts and return. Shared typed projections; reconnect/repeated-click/stale/original-route cases. [Live team workspace](live-team-workspace-v0.md) makes exchange, revision and original-coordinator continuation visible. Its [Work-scale map track](live-team-workspace-v0.md#11-delivery-order-and-relationship-to-aggressive-r2-progress) draws each Goal's typed Todo relations first (W1), then live state and outputs on the same nodes. Then intelligent review, keyboard accessibility, bilingual terminology, actionable errors and offline degradation; interrupt only for actual decisions | +| **S5 Frontend, Lark and human interaction · P0/P1** | Local chat, settings, proposals and partial Goal Channel verticals exist; shared audience/session/work readback needs qualification | One journey spans settings, work graph, handoff, blockers, cost, corrections, artifacts and return. Shared typed projections; reconnect/repeated-click/stale/original-route cases. Realtime IM reuses Chat/Turn: isolate independent conversations on one listener, then qualify ordinary DM onboarding, explicit role selection, busy-session admission, provisional progress, media and exact permission callbacks through [the shared operational contract](capable-manager-semantic-handoff-v0.md#10-operational-contract). [Live team workspace](live-team-workspace-v0.md) makes exchange, revision and original-coordinator continuation visible. Its [Work-scale map track](live-team-workspace-v0.md#11-delivery-order-and-relationship-to-aggressive-r2-progress) draws each Goal's typed Todo relations first (W1), then live state and outputs on the same nodes. Then intelligent review, keyboard accessibility, bilingual terminology, actionable errors and offline degradation; interrupt only for actual decisions | | **S6 Materials, evidence, memory and learning · P1** | Authority registry, material lifecycle/frontier, decision context, reward memory and turn recall exist; direction baseline and parts of attribution remain proposed | Connect material revision→same-Agent read→decision reference→artifact/outcome. Expose expiry/revocation/source loss and forgetting policy. Handoff preserves decision-relevant summaries and authorized artifacts; qualify OpenViking/Obelisk as optional providers. Prove causal utility with controls, not relevance alone | | **S7 Budget, scheduling and fleet scale · P0 observation/P1–P2 expansion** | Quota/scheduler and partial usage aggregates exist; full provider cost, distributed reservations and hundred-Agent concurrency need evidence | Separate configured budget, admission, consumption and estimates; unknown is not zero and replay cannot double-charge. R7 pagination/bounded summaries and [complete-history transport](typescript-control-plane-migration-v0.md), including refresh/replay/single-debit evidence beyond the RPC limit; provider/host limits, fairness, backpressure, event wake and isolation; report registration/activity/throughput and cost per accepted outcome separately | | **S8 Capabilities, extensions and domain integration · P1/P2** | Capability catalog, extension lifecycle, hooks, engineering/research/content/office capabilities and computer-use contracts exist | First exercise the shared control plane with existing issue-fix/PR-review and material/research callers. Every provider has readiness/version/permissions/default-off/uninstall/rollback/isolation and real-entry evidence. New domain effects start with one simulated operation, not a marketplace or workflow DSL | | **S9 Identity, authority, privacy and trust · continuous P0/P1–P2 remote** | Public/private scope, capability gates, fencing and confirmation contracts belong to existing owners | R1/R3 cover sender/audience/artifact scope and stale authority; R6 authenticates tenant/Goal/actor/host, rotation/revocation and least privilege. Qualify credential custody, untrusted tool/document inputs, dependency supply chain, audit retention/deletion and vulnerability response through real paths; roles/messages/memory mint no write authority | -| **S10 Reliability, diagnostics and operations · P0/P1** | Recovery/canary, read-only diagnostics prototype and DSH event adapter exist; C0/C1, overhead and full operations qualification are open | Failure classification→observable state→recovery drill→regression prevention; process/storage/network/delivery failures and data growth. Accepted Chat requests must settle even when runtime preparation fails before dispatch; qualify missing runtime assets, stop races and recovery without replay under [the shared conversation operational contract](capable-manager-semantic-handoff-v0.md#10-operational-contract). A context/provider read returning after the stop wait must honor the persisted Turn and exact Session claim: no late provider dispatch or handoff, even when a fresh request has completed. This bounded GQ08 repair does not certify upstream interrupt fidelity or effects already admitted elsewhere. Use [bounded repair lookup and targeted diagnostics](../../../skills/loopx-self-repair/references/targeted-diagnostics.md) to reduce redundant reads above the provider boundary; measure backend-specific cold/warm reads, writes and lock waits separately. Freeze SLO/RPO/RTO/capacity/retention boundaries and measure before qualification. Runbooks include upgrade, restore, stop and human takeover; test counts do not prove recovery | +| **S10 Reliability, diagnostics and operations · P0/P1** | Recovery/canary, read-only diagnostics prototype and DSH event adapter exist; C0/C1, overhead and full operations qualification are open | Failure classification→observable state→recovery drill→regression prevention; process/storage/network/delivery failures and data growth. Accepted Chat requests must settle even when runtime preparation fails before dispatch; qualify missing runtime assets, stop races and recovery without replay under [the shared conversation operational contract](capable-manager-semantic-handoff-v0.md#10-operational-contract). A context/provider read returning after the stop wait must honor the persisted Turn and exact Session claim: no late provider dispatch or handoff, even when a fresh request has completed. This bounded GQ08 repair does not certify upstream interrupt fidelity or effects already admitted elsewhere. Use [bounded repair lookup and targeted diagnostics](../../../skills/loopx-self-repair/references/targeted-diagnostics.md) to reduce redundant reads above the provider boundary; measure backend-specific cold/warm reads, writes and lock waits separately. Freeze SLO/RPO/RTO/capacity/retention boundaries and measure before qualification. Lark transport backpressure and drain retain the App consumer lease; unpersisted buffered lines are not accepted work or completion receipts. Runbooks include upgrade, restore, stop and human takeover; test counts do not prove recovery | | **S11 Evaluation and scientific research · continuous P1/P2 research** | Benchmark toolkit, Explore, long-horizon portfolio and ten frontier-science tracks have designs/partial implementations | Pin native/passive/governed arms, model/harness/budget/task split and evaluator; report native scores, cost, failures, attention and uncertainty. Prioritize sequential evidence, continuation and stride; memory, formal kernel, curriculum/evolution, active experiments and multiscale state follow T01–T10 gates without automatic production treatment | | **S12 Release, developer experience and community governance · P0 hygiene/P1** | Install/source validation, registration, DCO/PR, test layers, contributor routes and bilingual docs exist | Qualify first work and upgrade/rollback from clean machines/release artifacts; host/OS support follows the release contract. Reduce localization/test/review effort for useful changes; preserve exact-head evidence, fixtures, compatibility, maintainer routing and contributor credit; retire duplicate protocols/stale evidence | | **S13 Adoption, ecosystem and sustainability · P1 discovery/P2 pilots** | Public adoption loop, showcases, licensing/governance and observer-first product contract exist; paid PMF is unproven | Gather independent first/repeat usage and exit reasons; reproducible cases and pilots with fixed budgets/acceptance/rollback. Retain reusable adapters/delivery guides. Account for model/compute/storage/support and maintenance costs; only repeated demand justifies commercial hosting/support/distribution decisions, with no invented SLA or open-source-term change | diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md index ccd3868d36..cd3a62ef1e 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md @@ -61,12 +61,12 @@ managed 与 attached 的工作对话都应能持续在 LoopX 中进行:沿用 | **S2 typed 内核与 durable authority · P0/P1** | Effect/Todo/quota/recovery owner、TS 事务迁移与 store 候选已存在;writer 和 provider 晋升仍未全闭合 | 每次迁移一个真实事务/恢复生命周期,先语义反例再切换/删除旧 owner;R1 正确性先于迁移数量。真实 backend、并发/fence、ambiguous commit、保留/导出恢复、bridge 成本及 D1–D3 资格 | | **S3 目标规划与 multi-Agent 协作 · P0/P1** | Vision/replan、peer frontier、claim/lease、directory、manager_context 和显式接续有基础;通用 handoff/共享修订未闭环 | R2 必须证明 peer 依赖;R3 完成并行汇合、流水线、求助/复核、接续、自动回报;R4 做一个保持 intent 的 amendment class。检查依赖环、输入失效、拒绝/延期、lease 转移、同基线竞争及 aggregate acceptance | | **S4 runtime/host/daemon · P0/P1** | attached/managed、Turn、broker、runtime connector 和 Desktop 修复存在;“registered”不等于可执行 | 选择一个真实合格组合完成多 Turn supervision;restart/cancel/drain/stop 不丢工作且旧 executor 被 fence。之后扩 host parity、service-profile 唯一 owner、干净安装与版本升级;按 adapter 能力显示不支持项 | -| **S5 前端、Lark 与人机交互 · P0/P1** | 本地对话、settings、proposal 和部分 Goal Channel vertical 已有;统一受众/会话/工作回读仍需资格 | 用一个团队旅程贯穿设置、工作图、handoff、阻塞、成本、修订、产物和回报;共享 typed projection,验证重连/重复点击/stale/原路反馈。再做 intelligent review、无障碍键盘流程、中英术语、错误可恢复和离线降级;只在真实决策处打断人;[团队实时工作区](live-team-workspace-v0.zh-CN.md)让交换、修订与原协调员继续推进可见;其[工作尺度地图分线](live-team-workspace-v0.zh-CN.md#11-交付顺序与激进推进-r2-的关系)先画出每个 Goal 的类型化 Todo 关系(W1),再在同一节点叠加实时状态与产出 | +| **S5 前端、Lark 与人机交互 · P0/P1** | 本地对话、settings、proposal 和部分 Goal Channel vertical 已有;统一受众/会话/工作回读仍需资格 | 用一个团队旅程贯穿设置、工作图、handoff、阻塞、成本、修订、产物和回报;共享 typed projection,验证重连/重复点击/stale/原路反馈。实时 IM 复用 Chat/Turn:先让同一 listener 下的独立会话互不阻塞,再按[共享运行契约](capable-manager-semantic-handoff-v0.zh-CN.md#10-运行契约)验收普通私聊开通、显式角色选择、忙碌会话准入、临时进度、媒体和精确权限回调。再做 intelligent review、无障碍键盘流程、中英术语、错误可恢复和离线降级;只在真实决策处打断人;[团队实时工作区](live-team-workspace-v0.zh-CN.md)让交换、修订与原协调员继续推进可见;其[工作尺度地图分线](live-team-workspace-v0.zh-CN.md#11-交付顺序与激进推进-r2-的关系)先画出每个 Goal 的类型化 Todo 关系(W1),再在同一节点叠加实时状态与产出 | | **S6 材料、证据、记忆与学习 · P1** | authority registry、material lifecycle/frontier、decision context、reward memory、turn recall 已有;方向基线和部分归因仍是提案 | 先打通“材料 revision→同 Agent 阅读→决策引用→产物/结果”;失效、撤销、来源消失与遗忘策略可回读。handoff 保存影响决策的摘要与授权 artifact;OpenViking/Obelisk 按可选 provider 资格化。utility 的因果收益另以对照证明,不把相关性当提升 | | **S7 预算、调度与 fleet 规模 · P0 观测/P1–P2 扩展** | quota/scheduler 与部分 usage aggregate 存在;全 provider 成本、分布式资源预留及百 Agent 并发尚需证据 | 先区分配置预算、准入、消耗与估算;未知成本不记零、重复事件不双记。R7 分页/有界摘要及[完整历史传输](typescript-control-plane-migration-v0.zh-CN.md),验收超出 RPC 上限后的写回/重放/单次扣记;provider/host 限流、公平性、背压、事件唤醒与失败隔离;分别报告注册数/活跃数/吞吐量和每个验收成果成本 | | **S8 能力、扩展与领域集成 · P1/P2** | 已有 capability catalog、extension 生命周期、hook、工程/研究/content/office 能力及 computer-use 合同 | 优先用现有 issue-fix/PR-review 和材料/研究 caller 检验共享控制面;每个 provider 带 readiness、版本、权限、默认关闭、卸载/回滚、失败隔离与真实入口证据。新 domain effect 从模拟单操作闭环开始,不先建市场或通用工作流 DSL | | **S9 身份、权限、隐私与信任 · P0 持续/P1–P2 远端** | public/private 边界、作用域、capability gate、fence 与确认合同分布在已有 owner | 随 R1/R3 验 sender/audience/artifact scope 和 stale authority;远端 R6 必须认证 tenant/Goal/actor/host、轮换撤销与最小权限。凭据保管、非可信工具/文档输入、依赖供应链、审计留存/删除及漏洞响应纳入真实路径;角色、消息或 memory 不铸造写权限 | -| **S10 可靠性、诊断与运行运营 · P0/P1** | recovery/canary、read-only diagnostics 原型及 DSH event adapter 已有;C0/C1、开销和完整运营资格仍未闭合 | 故障分类→可观察状态→恢复演练→防复发;覆盖进程/存储/网络/投递故障和数据增长。Chat 上下文或 provider 读取晚于停止等待返回时,按持久 Turn 和精确 Session claim 判断:即使新请求已完成,也不得再启动旧请求或交接迟到结果。这项有界 GQ08 修复不证明上游 interrupt 保真,也不取消其他 owner 已准入的效果;完整恢复仍遵循[共享对话运行契约](capable-manager-semantic-handoff-v0.md#10-operational-contract)。定义并冻结 SLO、RPO/RTO、容量/保留边界,实测后标 qualified;运行手册含升级、备份恢复、停止与人工接管,不以测试数代替恢复结果 | +| **S10 可靠性、诊断与运行运营 · P0/P1** | recovery/canary、read-only diagnostics 原型及 DSH event adapter 已有;C0/C1、开销和完整运营资格仍未闭合 | 故障分类→可观察状态→恢复演练→防复发;覆盖进程/存储/网络/投递故障和数据增长。Chat 上下文或 provider 读取晚于停止等待返回时,按持久 Turn 和精确 Session claim 判断:即使新请求已完成,也不得再启动旧请求或交接迟到结果。这项有界 GQ08 修复不证明上游 interrupt 保真,也不取消其他 owner 已准入的效果;完整恢复仍遵循[共享对话运行契约](capable-manager-semantic-handoff-v0.md#10-operational-contract)。定义并冻结 SLO、RPO/RTO、容量/保留边界,实测后标 qualified;Lark 传输背压与 drain 保留 App consumer lease;未持久化的缓冲消息不算已受理工作或完成回执。运行手册含升级、备份恢复、停止与人工接管,不以测试数代替恢复结果 | | **S11 评测与科学研究 · P1 持续/P2 研究** | benchmark toolkit、Explore、长程 portfolio 与十轨 frontier science 有设计/局部实现 | 固定 native/passive/governed arm、模型/harness/预算/task split 与 evaluator;报告原生分数、成本、失败、人工介入和不确定性。sequential evidence、continuation、stride 为早期研究;memory、formal kernel、curriculum/evolution、主动实验与多尺度状态按 T01–T10 分阶段,不自动影响生产 | | **S12 发布、开发体验与社区治理 · P0 卫生/P1** | 安装、源码验证、扩展注册、DCO/PR、测试层级、contributor route 与双语文档已存在 | 从干净机器/发布包验一条首次工作和一次升级/回滚;host/OS 支持以 release contract 为准。缩短合理改动的定位、测试和 review 成本;公开精确 head、可重复 fixture、兼容窗口、维护者路由和贡献归属,退休重复协议及过时证据 | | **S13 采用、生态与商业可持续性 · P1 发现/P2 试点** | 公开 adoption loop、showcase、license/governance、observer-first 产品合同已存在;付费 PMF 未证明 | 先收集真实独立首次使用/重复使用/退出原因,做可复现案例和有固定预算/验收/回滚的试点;沉淀 reusable adapter 与交付手册。核算模型/计算/存储/支持成本及维护负担;满足重复需求后再决策商业托管边界、支持等级和分发,不承诺 SLA 或擅改开源条款 | diff --git a/loopx/extensions/lark/docs/lark-event-inbox.md b/loopx/extensions/lark/docs/lark-event-inbox.md index 17b9314466..7c7d50bcac 100644 --- a/loopx/extensions/lark/docs/lark-event-inbox.md +++ b/loopx/extensions/lark/docs/lark-event-inbox.md @@ -158,6 +158,12 @@ Agent-labelled Topic for each route. Users send requests inside the matching Topic. A group-level message with more than one eligible Agent route is deliberately rejected as ambiguous instead of guessing an Agent from prose. +For interactive replacement and busy-listener qualification, use the +[native realtime conversation readiness guide](realtime-conversation-readiness.md). +The long-lived listener is independent of a Goal's periodic heartbeat; +independent conversations can progress concurrently, while a busy conversation +still waits for its current answer before admitting a follow-up. + For a periodic-report request, semantic activation belongs to the Agent. After reading an exact item, the Agent calls `loopx periodic-report request` with its `message_id`. The Lark adapter validates binding and addressing evidence only; diff --git a/loopx/extensions/lark/docs/owner-controlled-bots.md b/loopx/extensions/lark/docs/owner-controlled-bots.md new file mode 100644 index 0000000000..3648af0956 --- /dev/null +++ b/loopx/extensions/lark/docs/owner-controlled-bots.md @@ -0,0 +1,99 @@ +# Owner-controlled assistant and steward Bots + +An operator can plan two Lark applications: an assistant for ordinary project +conversations, and a steward for explicitly delegated long-term coordination. +They may share LoopX Core mechanics, but each application needs its own verified +identity, audience and workspace grants. This is a deployment and development +handoff; it does not certify the unfinished assistant DM journey. + +## Readiness boundary + +| Surface | Current evidence | Remaining qualification | +| --- | --- | --- | +| Existing group/Agent Topic and manager connections | Native provider routing, Inbox, Chat Session and verified replies | Installed application and host journey | +| Independent conversations on one App listener | Bounded transient dispatch with source/shared-Session ordering | Live slow-request and recovery journey | +| Ordinary project conversation and private Bot DM | Product requirement in the conversation RFC | Shared Core context, authenticated onboarding and installed continuity | +| Busy-conversation admission, progress, media and permission callbacks | Existing Core/host pieces and provider presentation | End-to-end behavior for the selected host; no adapter-wide parity claim | +| Two operator-owned applications | Separate profile/identity configuration is the intended boundary | Both Apps together, exact identity isolation and one consumer owner per App | + +See [realtime readiness](realtime-conversation-readiness.md) for what the current +transport repair proves. Starting a listener is not evidence that the remaining +rows pass. + +## Install and establish ownership + +1. Read the target machine's existing installation report. Choose one LoopX + installation owner using the [install guide](../../../../docs/guides/installing-loopx.md). + For unreleased development, use a clean source checkout and qualify the exact + candidate before installing it. Record both source and installed revisions; + a public PR or passing source test does not identify the installed package. +2. Verify that the operator controls the tenant, developer console, application + administration and publication scope. A personal login alone does not prove + application ownership. Registration, login, consent and publication use that + operator's own account. +3. Create separate profiles, for example `personal-assistant` and + `personal-steward`. The existing App setup invokes + `lark-cli config init --new --name PROFILE --brand feishu --lang zh_cn`; + inspect the selected CLI's help and the existing settings wizard first. This + is an App registration action, not a read-only diagnostic. +4. Keep credentials in provider-private local storage. Verify each profile's + actual App identity and authorized user independently. Provider user ids can + be App-scoped; do not copy an old user id or grant into a different App. +5. Review permissions against the enabled journey. The current + [recommended scope bundle](../bot_scopes.py) includes group management and + document comments as well as conversation features. It is not evidence of a + minimal assistant-only permission set. Verify message events, readback and + any enabled card/media operations against the actual provider version. +6. Start with an owner-only audience and an explicit workspace. One listener + owner per App consumes `im.message.receive_v1`; do not run a bridge and native + listener concurrently for that same App. Two different Apps still need an + installed isolation check. Record the service owner and restart procedure. + +Keep configuration facts, transcripts, SDK/session databases and verification +URLs private. Moving to new Apps does not require importing a previous machine's +registry, Goals, Todos, credentials or model-session history. + +## Recipient purpose and shared implementation + +| Recipient | Default purpose | Access boundary | +| --- | --- | --- | +| Assistant | Ordinary authorized project chat; explicitly selected existing Agent when available | Selected workspace and conversation audience | +| Steward | Persistent commitments, decomposition, coordination and acceptance | Explicitly granted portfolio; an empty installation has no inherited commitments | + +The assistant must not borrow the global manager objective or create a hidden +Goal to bypass the current Goal/manager Chat context restriction. Implement the +authorized project-conversation context through the existing typed Core owner +in the [conversation RFC](../../../../docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md). +Reuse Chat Session/Turn, admission, operation receipts, progress and host +execution. Provider transport must not add a parallel scheduler, approval store, +task ledger or model runner. + +Separate durable admission from long model execution to improve a busy +conversation. A receipt must truthfully distinguish received, accepted, queued, +running and terminal result. Recipient changes affect subsequent input; accepted +work returns to its original Session and audience. A stop or permission answer +must resolve one exact request/operation and preserve existing host policy. + +## Installed acceptance and handoff record + +Exercise both Apps with synthetic content on the intended desktop and phone: + +- First message and follow-up preserve the intended Session; unauthorized users + and cross-App bindings fail closed. +- A blocked assistant request does not block an independent steward request; + another assistant input receives bounded, truthful queue/steering feedback. +- Progress stays provisional. Image/file input reaches the selected host or + reports unsupported. A permission callback binds the approving user, exact + operation and expiry; it never silently elevates policy. +- Exact stop, duplicate events, reconnect and lost reply readback do not restart + a model request, replay an effect or dispatch cancelled waiting input. +- Restart and login recovery use the same verified profiles and listener owner. + Record rollback to a previously qualified package without reusing old grants. + +For each case record `passed`, `failed` or `not_run`, the source/installed +revision, actual host capability, evidence category and next repair. Keep raw +provider evidence private. The +[steward operational contract](../../../../docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md#10-operational-contract) +owns the broader acceptance. Deployment handoff, product qualification and +optional state inheritance are separate completion facts. A new deployment +does not automatically retire a service on a different App or machine. diff --git a/loopx/extensions/lark/docs/realtime-conversation-readiness.md b/loopx/extensions/lark/docs/realtime-conversation-readiness.md new file mode 100644 index 0000000000..8f7c50b4ea --- /dev/null +++ b/loopx/extensions/lark/docs/realtime-conversation-readiness.md @@ -0,0 +1,114 @@ +# Native realtime conversation readiness + +LoopX's optional Lark provider receives messages through a long-lived +`im.message.receive_v1` consumer. The Chat server owns that listener; a Goal's +periodic heartbeat is a separate work trigger. Installation or a healthy +listener alone does not qualify an interactive conversation. + +For separate assistant and steward Apps, follow the +[owner-controlled Bot handoff](owner-controlled-bots.md). App setup, live product +qualification and optional state inheritance remain separate outcomes. + +## Existing path and its limits + +Configure a verified Bot profile and the intended group/Agent connections in +the existing workspace settings. A manager connection and worker Topics share +the same bundled transport. Exact worker Topics select their bound Agent; +ambiguous group messages must not guess a recipient from their text. App +credentials remain in the provider's local profile, outside Goal state. + +The existing runtime checks a provider-ready marker, reuses the bound Chat +Session, persists the source in the Inbox, and verifies the reply before ACK. +Received-message reactions, terminal replies, confirmation cards and +listener health are distinct observations. A private inbox configuration is +published with the shared atomic private-JSON writer: concurrent callers use +independent owner-private temporaries, and a failed publication preserves the +previous configuration and removes its temporary. + +For already-enabled Lark listeners, independent conversations now run through +up to four transport handlers. The dispatch queue counts at most 64 messages +including active handlers; provider frame and pipe buffers remain provider-owned. A source stays FIFO across binding revisions; sources bound +to the same Session also share a scheduling fence. One busy source does not +consume workers by waiting on a Session lock. Capacity applies backpressure +instead of dropping messages or growing an unbounded executor queue. Actual +processing reads fresh bindings and still uses Core admission/budget rules. +These limits bound transport work, not permission to run four models. + +Normal provider rotation drains observed messages before releasing the App +consumer lease. Stop or reader failure discards waiting transient lines without +Inbox ACK; active handlers retain ownership until they settle. This buffer is +not a durable queue or an RPO guarantee: only the existing Inbox/Chat path owns +persisted source/admission receipts. A handler failure surfaces through existing +service recovery, without retrying its model call in the transport helper. +Inactive profiles still start no consumer or dispatch workers. + +Same-conversation handlers remain serialized through terminal answer/reply. +Core Session queue/steering support therefore does not yet prove that a busy +Lark conversation can admit a correction promptly. The shipped group and Topic +setup is also not qualification of private Bot DM onboarding, ordinary non-Goal +Codex chat, token streaming or attachment delivery to a model. Deterministic +native-path checks use real Inbox files and reply readback with provider/model +doubles; they do not qualify an installed live provider or host. + +The dispatch helper stays in the Python Lark extension because it only schedules +that bundled provider's transient stream and invokes its existing handler. It +owns no shared domain policy, persistence, approvals or model execution. Shared +conversation/admission/progress changes belong to the existing typed Core owner. + +## Product boundary + +A Bot is the realtime entry to ordinary project chat, a selected existing Agent, +or an explicitly selected steward. The steward owns long-term commitments, +coordination and acceptance; ordinary chat should not inherit that objective or +portfolio access. Reuse Session/Turn and host capabilities, while preserving +recipient purpose, audience and workspace grants. The current Goal/manager +setup is the existing supported path, not proof that plain project chat is ready. + +## Interactive experience to retain + +[Claude-to-IM](https://github.com/op7418/Claude-to-IM) is a useful reference for +channel ergonomics. Retain persistent conversation continuity, visible +processing feedback, bounded streaming presentation, explicit permission +answers, media input and actionable recovery. Qualify each host adapter: +support in an IM adapter does not establish incremental model output or tool +approval callbacks in every SDK/runtime. + +Reuse Core Chat Session/Turn, the existing event/progress projection, attachment +handling and canonical operation/proposal receipts. Lark presents those facts +and authenticated user responses; it must not add a second task ledger, +approval store, scheduler or model-execution authority. Missing evidence is not +permission to restart a request in a fresh model thread, elevate host policy +or acknowledge work as completed. + +## Qualification before switching + +The acceptance owner is the steward RFC's +[operational contract](../../../../docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md#10-operational-contract), +under roadmap S1/S5/S10. Exercise a pinned installed package through the real +provider and host, with synthetic content: + +- A first request and follow-up retain the exact intended Session and audience. +- While one model request is blocked, a second request, a correction and a + distinct role's request receive truthful, bounded admission/queue feedback. +- A visible progress/partial card remains provisional; only the terminal + canonical result qualifies completion. Lost card readback retries delivery, + without another model request or replaying an effect. +- An image/file reaches the authorized host with bounded private lifetime; + unsupported media is explicitly unavailable, rather than silently omitted. +- A permission decision binds the exact operation, audience, expiry and + approving user. An unverified callback never broadens execution permission. +- Reconnect/restart, duplicate events, unavailable provider and stop preserve + source identity and accepted work, with no duplicate answer or late dispatch. + +Report pass, failure and untested separately for transport, admission, host +execution and user-visible delivery. Synthetic fixtures do not certify live +adoption. Do not transfer old sessions, credentials, bindings or permissions. +An existing bridge stays a separate operator-controlled service until the +replacement has passed its actual journey; setup must not silently start a +second client for the same App. + +The [official Lark SDK](https://github.com/larksuite/node-sdk/blob/main/README.md#subscribing-to-events-using-long-connection-mode) +describes WebSocket events, the provider processing deadline and non-broadcast +delivery across clients. A provider receive ACK, an Inbox ACK and a completed +model Turn are separate lifecycle facts. Qualify the selected SDK/provider +version; the transport deadline is not a promise of model response latency. diff --git a/loopx/extensions/lark/goal_topic_dispatch.py b/loopx/extensions/lark/goal_topic_dispatch.py new file mode 100644 index 0000000000..45cd8432d5 --- /dev/null +++ b/loopx/extensions/lark/goal_topic_dispatch.py @@ -0,0 +1,118 @@ +"""Bounded transient dispatch for one provider stream, not Chat admission authority.""" + +from collections import deque +from collections.abc import Callable +from concurrent.futures import ThreadPoolExecutor +import threading + + +class ProfileEventDispatch: + """Keep conversation order without letting one answer block every audience. + + A lane is only a transport scheduling key. The handler must read fresh + binding/authority and persist through the existing Inbox/Chat owners. + Pending lines are transient, just like the consumer's pipe buffer. + """ + + def __init__( + self, + handle: Callable[[str], None], + stop: threading.Event, + *, + workers: int = 4, + capacity: int = 64, + ) -> None: + if workers < 1 or capacity < workers: + raise ValueError("dispatch capacity must cover its positive worker limit") + self._handle = handle + self._stop = stop + self._workers = workers + self._capacity = capacity + self._condition = threading.Condition() + self._waiting: dict[str, deque[tuple[tuple[str, ...], str]]] = {} + self._ready: deque[str] = deque() + self._active: dict[str, tuple[str, ...]] = {} + self._pending = 0 + self._closing = False + self._cancelled = False + self._failure: Exception | None = None + self._executor = ThreadPoolExecutor( + max_workers=workers, thread_name_prefix="loopx-lark-dispatch" + ) + + def submit(self, lane: tuple[str, ...], line: str) -> bool: + with self._condition: + while self._pending >= self._capacity and not self._closing: + if self._stop.is_set(): + return False + self._condition.wait(0.1) + if self._failure is not None: + raise self._failure + if self._closing or self._stop.is_set(): + return False + source = lane[0] + queue = self._waiting.setdefault(source, deque()) + queue.append((lane, line)) + self._pending += 1 + if source not in self._active and source not in self._ready: + self._ready.append(source) + self._start_ready() + return True + + def _start_ready(self) -> None: + if self._stop.is_set(): + self._discard_waiting() + # A stable source keeps FIFO across rebinds; an optional Session key + # also fences different sources that share the same conversation owner. + for _ in range(len(self._ready)): + if self._cancelled or len(self._active) >= self._workers: + break + source = self._ready.popleft() + lane, line = self._waiting[source][0] + occupied = {key for active in self._active.values() for key in active} + if occupied.intersection(lane): + self._ready.append(source) + continue + self._waiting[source].popleft() + self._active[source] = lane + self._executor.submit(self._run, source, line) + + def _discard_waiting(self) -> None: + self._cancelled = True + self._pending -= sum(len(queue) for queue in self._waiting.values()) + self._waiting.clear() + self._ready.clear() + + def _run(self, source: str, line: str) -> None: + try: + self._handle(line) + except Exception as exc: + with self._condition: + if self._failure is None: + self._failure = exc + self._closing = True + self._discard_waiting() + finally: + with self._condition: + self._pending -= 1 + self._active.pop(source) + queue = self._waiting.get(source) + if queue: + self._ready.append(source) + else: + self._waiting.pop(source, None) + self._start_ready() + self._condition.notify_all() + + def close(self, *, cancel_pending: bool = False) -> None: + """Retain ownership until active handlers settle; never replay a handler.""" + + with self._condition: + self._closing = True + while self._pending: + if cancel_pending or self._stop.is_set(): + self._discard_waiting() + self._condition.wait(0.1) + self._executor.shutdown(wait=True) + if self._failure is not None: + raise self._failure diff --git a/loopx/extensions/lark/goal_topic_runtime.py b/loopx/extensions/lark/goal_topic_runtime.py index 470ede0087..51628b490a 100644 --- a/loopx/extensions/lark/goal_topic_runtime.py +++ b/loopx/extensions/lark/goal_topic_runtime.py @@ -44,6 +44,7 @@ ) from .goal_channel_targets import goal_channel_target_for_name from .goal_topic_connections import decide_lark_topic_event +from .goal_topic_dispatch import ProfileEventDispatch from .inbox_reply import CommandRunner, reply_lark_event_inbox from .manager_reply_delivery import ( load_delivery as _load_manager_delivery, @@ -69,6 +70,7 @@ manager_part_delivery_readback, ) from .manager_reply_format import repair_manager_reply_text +from .private_json import write_private_json_atomic from .outbound import LarkOutboundTextError, safe_lark_plain_text_fallback from .inbox_reactions import ( _create_reaction, @@ -297,6 +299,100 @@ def _event_payloads(stdout: Any) -> list[Mapping[str, Any]]: return events +def _profile_event_route( + *, + profile: str, + snapshot: Mapping[str, Any], + event: Mapping[str, Any], +) -> tuple[ + tuple[Mapping[str, Any], dict[str, Mapping[str, Any]], dict[str, Any]] | None, str +]: + """Select the existing profile/target/Topic scope; never grant turn authority.""" + + profile_config = _active_profile_configs(snapshot).get(profile) + if profile_config is None: + return None, "inactive" + target_payload = snapshot.get("target_payload") + target_payload = target_payload if isinstance(target_payload, Mapping) else {} + binding_payloads = snapshot.get("binding_payloads") + binding_payloads = binding_payloads if isinstance(binding_payloads, Mapping) else {} + active_target_refs = { + str(binding.get("target_ref") or "") + for goal_id, payload in binding_payloads.items() + if isinstance(payload, Mapping) + for binding in bindings_for_goal(payload, str(goal_id)) + if binding.get("enabled") is True + } + chat_id = str(event.get("chat_id") or "") + target_match = _target_for_profile_chat( + target_payload, + profile=profile, + chat_id=chat_id, + bot_app_id=str(profile_config.get("bot_app_id") or ""), + active_target_refs=active_target_refs, + root_id=str(event.get("root_id") or ""), + binding_payloads=binding_payloads, + ) + if target_match is None: + return None, "target_unmatched" + target_ref, _target = target_match + routed_event = dict(event) + root_id = str(routed_event.get("root_id") or "") + if not MESSAGE_ID_PATTERN.fullmatch(root_id) and not has_manager_binding( + binding_payloads, target_ref + ): + candidate_roots = _topic_roots_for_target( + binding_payloads, + target_ref=target_ref, + ) + if len(candidate_roots) > 1: + return None, "topic_context_ambiguous" + if not candidate_roots: + return None, "topic_context_missing" + routed_event["root_id"] = candidate_roots[0] + return ( + target_payload, + _binding_payloads_for_target(binding_payloads, target_ref=target_ref), + routed_event, + ), "matched" + + +def _profile_event_lane( + *, + profile: str, + snapshot: Mapping[str, Any], + event: Mapping[str, Any], + runtime_root: str | Path, +) -> tuple[str, ...]: + """Only serialize transport work; actual ingress re-reads its authority.""" + + scoped, _status = _profile_event_route( + profile=profile, snapshot=snapshot, event=event + ) + if scoped is not None: + targets, bindings, routed_event = scoped + decision = decide_lark_topic_event( + target_payload=targets, + binding_payloads=bindings, + event=routed_event, + runtime_root=runtime_root, + ) + route = decision.get("route") or {} + if route: + if route.get("conversation_kind") == "manager": + source = "manager." + _opaque_digest(route.get("app_ref"), event.get("chat_id")) + else: + # Compatibility direct_session opens/resumes this exact channel. + source = "topic." + _opaque_digest( + route.get("app_ref"), route.get("target_ref"), + route.get("topic_root_message_id"), + ) + if route.get("session_id"): + return (source, "session." + str(route["session_id"])) + return (source,) + return ("unmatched." + _opaque_digest(profile, event.get("chat_id")),) + + def poll_lark_goal_topic_profile_once( *, profile: str, @@ -342,62 +438,25 @@ def poll_lark_goal_topic_profile_once( "replied_count": 0, } - target_payload = snapshot.get("target_payload") - target_payload = target_payload if isinstance(target_payload, Mapping) else {} - binding_payloads = snapshot.get("binding_payloads") - binding_payloads = binding_payloads if isinstance(binding_payloads, Mapping) else {} - active_target_refs = { - str(binding.get("target_ref") or "") - for goal_id, payload in binding_payloads.items() - if isinstance(payload, Mapping) - for binding in bindings_for_goal(payload, str(goal_id)) - if binding.get("enabled") is True - } events = _event_payloads(result.get("stdout")) replied_count = 0 event_statuses: list[str] = [] event_reasons: list[str | None] = [] for event in events: - chat_id = str(event.get("chat_id") or "") - target_match = _target_for_profile_chat( - target_payload, + scoped, rejection = _profile_event_route( profile=profile, - chat_id=chat_id, - bot_app_id=str(profile_config.get("bot_app_id") or ""), - active_target_refs=active_target_refs, - root_id=str(event.get("root_id") or ""), - binding_payloads=binding_payloads, + snapshot=snapshot, + event=event, ) - if target_match is None: - event_statuses.append("target_unmatched") + if scoped is None: + event_statuses.append(rejection) event_reasons.append(None) continue - target_ref, _target = target_match - routed_event = dict(event) - root_id = str(routed_event.get("root_id") or "") - if not MESSAGE_ID_PATTERN.fullmatch(root_id) and not has_manager_binding( - binding_payloads, target_ref - ): - candidate_roots = _topic_roots_for_target( - binding_payloads, - target_ref=target_ref, - ) - if len(candidate_roots) > 1: - event_statuses.append("topic_context_ambiguous") - event_reasons.append(None) - continue - if not candidate_roots: - event_statuses.append("topic_context_missing") - event_reasons.append(None) - continue - routed_event["root_id"] = candidate_roots[0] + target_payload, binding_payloads, routed_event = scoped try: event_result = process_lark_goal_topic_event( target_payload=target_payload, - binding_payloads=_binding_payloads_for_target( - binding_payloads, - target_ref=target_ref, - ), + binding_payloads=binding_payloads, event=routed_event, runtime_root=runtime_root, goal_contexts=( @@ -529,44 +588,26 @@ def stop_consumer() -> None: replied_count = 0 provider_ready = False exit_reason: str | None = None - try: - stdout = process.stdout - if stdout is None: - return { - "ok": False, - "status": "stream_failed", - "event_count": 0, - "replied_count": 0, - } - for line in stdout: - if stop.is_set(): - break - stripped = line.strip() - if stripped.startswith(_EVENT_READY_PREFIX): - provider_ready = True - if health_sink is not None: - health_sink({"status": "listening", "error_code": None}) - continue - if stripped.startswith("[event] exited "): - match = _EVENT_EXIT_REASON.search(stripped) - exit_reason = match.group(1) if match else None - continue - if stripped.startswith(_EVENT_DIAGNOSTIC_PREFIX): - continue - result = poll_lark_goal_topic_profile_once( - profile=profile, - snapshot=snapshot_provider(), - runtime_root=runtime_root, - answer=answer, - consume_runner=lambda _args, payload=line: { - "returncode": 0, - "stdout": payload, - "stderr": "", - }, - provider_runner=provider_runner, - reply_runner=reply_runner, - proposal_deliverer=proposal_deliverer, - ) + result_lock = threading.Lock() + + def handle_event(line: str) -> None: + nonlocal event_count, replied_count, provider_ready + + result = poll_lark_goal_topic_profile_once( + profile=profile, + snapshot=snapshot_provider(), + runtime_root=runtime_root, + answer=answer, + consume_runner=lambda _args, payload=line: { + "returncode": 0, + "stdout": payload, + "stderr": "", + }, + provider_runner=provider_runner, + reply_runner=reply_runner, + proposal_deliverer=proposal_deliverer, + ) + with result_lock: if int(result.get("event_count") or 0) and not provider_ready: # A provider event is stronger readiness evidence than a # diagnostic marker and protects compatibility with providers @@ -605,18 +646,60 @@ def stop_consumer() -> None: ), flush=True, ) + + dispatch = ProfileEventDispatch(handle_event, stop) + stream_failed = True + try: + stdout = process.stdout + if stdout is None: + return { + "ok": False, + "status": "stream_failed", + "event_count": 0, + "replied_count": 0, + } + for line in stdout: + if stop.is_set(): + break + stripped = line.strip() + if stripped.startswith(_EVENT_READY_PREFIX): + provider_ready = True + if health_sink is not None: + health_sink({"status": "listening", "error_code": None}) + continue + if stripped.startswith("[event] exited "): + match = _EVENT_EXIT_REASON.search(stripped) + exit_reason = match.group(1) if match else None + continue + if stripped.startswith(_EVENT_DIAGNOSTIC_PREFIX): + continue + # Batch-shaped provider records are split before scheduling, so an + # unrelated conversation in a batch can progress independently too. + for event in _event_payloads(line): + lane = _profile_event_lane( + profile=profile, + snapshot=snapshot_provider(), + event=event, + runtime_root=runtime_root, + ) + if not dispatch.submit(lane, json.dumps(event, ensure_ascii=False)): + break + stream_failed = False finally: - watcher_done.set() - if process.poll() is None: - process.terminate() try: - returncode = process.wait(timeout=3) - except subprocess.TimeoutExpired: - process.kill() - returncode = process.wait(timeout=3) - watcher.join(timeout=1) - if callback_stream is not None: - callback_stream.close() + dispatch.close(cancel_pending=stream_failed) + finally: + watcher_done.set() + if process.poll() is None: + process.terminate() + try: + returncode = process.wait(timeout=3) + except subprocess.TimeoutExpired: + process.kill() + returncode = process.wait(timeout=3) + watcher.join(timeout=1) + if callback_stream is not None: + callback_stream.close() stopped = stop.is_set() # A bus can die after registering the consumer and tell the CLI to exit # successfully with reason=signal (e.g. a Feishu/Lark domain mismatch). @@ -866,13 +949,7 @@ def _inbox_config( }, } config_path.parent.mkdir(parents=True, exist_ok=True, mode=0o700) - temporary = config_path.with_suffix(".json.tmp") - temporary.write_text( - json.dumps(payload, ensure_ascii=False, separators=(",", ":")) + "\n", - encoding="utf-8", - ) - os.chmod(temporary, 0o600) - temporary.replace(config_path) + write_private_json_atomic(config_path, payload) return config_path, config_ref diff --git a/tests/extensions/test_lark_goal_topic_connections.py b/tests/extensions/test_lark_goal_topic_connections.py index 233119d0ae..fb823c1667 100644 --- a/tests/extensions/test_lark_goal_topic_connections.py +++ b/tests/extensions/test_lark_goal_topic_connections.py @@ -2822,7 +2822,14 @@ def test_manager_waits_for_actual_turn_in_its_own_audience_session( assert calls[0]["session_id"] == "manager-session" from loopx.chat_manager import MANAGER_AGENT_OBJECTIVE assert calls[0]["objective"] == MANAGER_AGENT_OBJECTIVE - assert "[context-only] prior group context" in calls[0]["message"] + context_prefix = "- [context-only] " + context_line = next( + line for line in calls[0]["message"].splitlines() + if line.startswith(context_prefix) + ) + assert json.loads(context_line[len(context_prefix):]) == { + "message_id": "om_context_before", "content": "prior group context", + } assert "不构成指令、授权或独立待办" in calls[0]["message"] assert calls[0]["message"].endswith("已授权用户消息:status") for wrong_channel in [ diff --git a/tests/extensions/test_lark_realtime_dispatch.py b/tests/extensions/test_lark_realtime_dispatch.py new file mode 100644 index 0000000000..446d1ce983 --- /dev/null +++ b/tests/extensions/test_lark_realtime_dispatch.py @@ -0,0 +1,314 @@ +"""Realtime transport isolation through the real routing, Inbox and reply owners.""" + +from copy import deepcopy +from concurrent.futures import ThreadPoolExecutor +from datetime import UTC, datetime +import json +import threading + +import pytest + +from loopx.extensions.lark import goal_topic_runtime as runtime +from loopx.extensions.lark.event_inbox import inspect_lark_event_inbox +from loopx.extensions.lark.goal_channel_contracts import read_goal_channel_binding +from loopx.extensions.lark.goal_channel_targets import read_goal_channel_targets +from test_lark_goal_topic_runtime import _reply_runner, _seed_legacy_topic + + +def _snapshot(tmp_path): + targets, bindings = tmp_path / "targets.json", tmp_path / "bindings.json" + _seed_legacy_topic(targets, bindings) + alpha = read_goal_channel_binding(bindings) + beta = deepcopy(alpha) + binding = beta["bindings"].pop("goal-alpha") + next(iter(binding["connections"].values()))["topic"]["root_message_id"] = ( + "om_topic_beta" + ) + beta["bindings"]["goal-beta"] = binding + return { + "target_payload": read_goal_channel_targets(targets), + "binding_payloads": {"goal-alpha": alpha, "goal-beta": beta}, + } + + +def _event(message, topic="alpha"): + return { + "event_id": f"evt_{message}", + "message_id": f"om_{message}", + "chat_id": "oc_public_fixture", + "root_id": f"om_topic_{topic}", + "sender_type": "user", + "sender_id": "ou_public_owner", + "mentions": [{"id": "cli_public_fixture"}], + "create_time": datetime.now(UTC).isoformat(), + "content": "@linkmacbot public synthetic question", + } + + +class _Consumer: + def __init__(self, lines): + self.stdout = lines + self.waited = False + + def poll(self): + return 0 + + def wait(self, timeout=None): + self.waited = True + return 0 + + +def _thread_reply_runner(): + local = threading.local() + + def run(args): + if not hasattr(local, "state"): + local.state = {} + return _reply_runner(local.state)(args) + + return run + + +def _start(tmp_path, snapshot, lines, answer, stop=None, health=None): + consumer = _Consumer(lines) + options = dict( + profile="mew", + snapshot_provider=lambda: snapshot, + stop=stop or threading.Event(), + runtime_root=tmp_path / "runtime", + answer=answer, + process_factory=lambda _args: consumer, + provider_runner=None, + reply_runner=_thread_reply_runner(), + health_sink=health.append if health is not None else None, + ) + return consumer, options + + +def _processed(tmp_path): + return sum( + inspect_lark_event_inbox(project=tmp_path / "runtime", config_path=path)[ + "processed_count" + ] + for path in (tmp_path / "runtime/.loopx/config/lark-goal-topics").glob("*.json") + ) + + +@pytest.mark.parametrize("batch", [False, True]) +def test_slow_role_does_not_block_another_role_but_its_followup_stays_ordered( + tmp_path, + batch, +): + snapshot = _snapshot(tmp_path) + active, release, beta_done, followup = [threading.Event() for _ in range(4)] + health = [] + order = [] + + def answer(route, _text): + message = route["message_id"] + order.append(message) + if message == "om_alpha_first": + active.set() + assert release.wait(5) + elif message == "om_alpha_followup": + followup.set() + else: + beta_done.set() + return "Public synthetic reply" + + def lines(): + yield "[event] ready event_key=im.message.receive_v1\n" + yield json.dumps(_event("alpha_first")) + assert active.wait(5) + events = [_event("alpha_followup"), _event("beta_first", "beta")] + if batch: + yield json.dumps(events) + else: + yield from map(json.dumps, events) + yield "[event] exited (reason: timeout)\n" + + consumer, options = _start(tmp_path, snapshot, lines(), answer, health=health) + with ThreadPoolExecutor(max_workers=1) as executor: + future = executor.submit(runtime.stream_lark_goal_topic_profile, **options) + try: + assert beta_done.wait(3), "An unrelated role waited for the slow answer" + assert not followup.is_set() + assert not future.done(), ( + "Consumer ownership ended before active work settled" + ) + assert not consumer.waited + finally: + release.set() + result = future.result(timeout=8) + assert result["event_count"] == result["replied_count"] == 3 + assert result["status"] == "stream_ended" + assert order == ["om_alpha_first", "om_beta_first", "om_alpha_followup"] + assert _processed(tmp_path) == 3 + assert sum(update.get("event_count", 0) for update in health) == 3 + assert "public synthetic question" not in json.dumps(health) + + +@pytest.mark.parametrize("disable", [False, True]) +def test_waiting_message_uses_fresh_binding_and_stop_never_acknowledges_it( + tmp_path, + disable, +): + snapshot = _snapshot(tmp_path) + active, release, queued = [threading.Event() for _ in range(3)] + stop = threading.Event() + answers = [] + health = [] + + def answer(route, _text): + answers.append(route["message_id"]) + active.set() + assert release.wait(5) + return "Public synthetic reply" + + def lines(): + yield "[event] ready event_key=im.message.receive_v1\n" + yield json.dumps(_event("alpha_first")) + assert active.wait(5) + yield json.dumps(_event("alpha_waiting")) + queued.set() + assert release.wait(5) + yield "[event] exited (reason: timeout)\n" + + _consumer, options = _start(tmp_path, snapshot, lines(), answer, stop, health) + with ThreadPoolExecutor(max_workers=1) as executor: + future = executor.submit(runtime.stream_lark_goal_topic_profile, **options) + try: + assert queued.wait(5) + if disable: + next( + iter( + snapshot["binding_payloads"]["goal-alpha"]["bindings"][ + "goal-alpha" + ]["connections"].values() + ) + )["enabled"] = False + else: + stop.set() + finally: + release.set() + result = future.result(timeout=8) + assert answers == ["om_alpha_first"] + assert result["replied_count"] == 1 + assert _processed(tmp_path) == 1 + assert not list((tmp_path / "runtime").rglob("om_alpha_waiting.json")) + if disable: + assert any( + update.get("last_event_status") == "ignored" + and update.get("last_event_reason") == "topic_mismatch" + for update in health + ) + else: + assert result["status"] == "stopped" + + +def test_explicit_session_groups_topics_but_different_sessions_are_independent( + tmp_path, monkeypatch +): + snapshot = _snapshot(tmp_path) + monkeypatch.setattr( + runtime, + "decide_lark_topic_event", + lambda **kwargs: { + "route": { + "session_id": "session-public", + "app_ref": "mew", + "target_ref": "fixture", + "topic_root_message_id": kwargs["event"]["root_id"], + } + }, + ) + keys = [ + runtime._profile_event_lane( + profile="mew", + snapshot=snapshot, + event=_event("message", topic), + runtime_root=tmp_path, + ) + for topic in ("alpha", "beta") + ] + assert keys[0][0] != keys[1][0] + assert keys[0][1] == keys[1][1] == "session.session-public" + monkeypatch.setattr( + runtime, + "decide_lark_topic_event", + lambda **kwargs: { + "route": { + "session_id": "session-new", + "app_ref": "mew", + "target_ref": "fixture", + "topic_root_message_id": kwargs["event"]["root_id"], + } + }, + ) + updated = runtime._profile_event_lane( + profile="mew", + snapshot=snapshot, + event=_event("message", "beta"), + runtime_root=tmp_path, + ) + assert updated[0] == keys[1][0] + assert updated[1] != keys[1][1] + + +def test_inactive_profile_keeps_the_native_listener_and_dispatch_off( + tmp_path, monkeypatch +): + def unexpected(*_args, **_kwargs): + raise AssertionError( + "An inactive profile must allocate no consumer or dispatcher" + ) + + monkeypatch.setattr(runtime, "ProfileEventDispatch", unexpected) + result = runtime.stream_lark_goal_topic_profile( + profile="mew", + snapshot_provider=lambda: {}, + stop=threading.Event(), + runtime_root=tmp_path / "runtime", + answer=unexpected, + process_factory=unexpected, + ) + assert result["status"] == "inactive" + assert result["event_count"] == result["replied_count"] == 0 + assert not (tmp_path / "runtime").exists() + + +def test_reader_failure_drains_active_work_and_cleans_up_without_accepting_buffered_input( + tmp_path, +): + snapshot = _snapshot(tmp_path) + active, release = threading.Event(), threading.Event() + answers = [] + + def answer(route, _text): + answers.append(route["message_id"]) + active.set() + assert release.wait(5) + return "Public synthetic reply" + + def lines(): + yield "[event] ready event_key=im.message.receive_v1\n" + yield json.dumps(_event("alpha_first")) + assert active.wait(5) + yield json.dumps(_event("alpha_waiting")) + raise OSError("synthetic reader failure") + + consumer, options = _start(tmp_path, snapshot, lines(), answer) + with ThreadPoolExecutor(max_workers=1) as executor: + future = executor.submit(runtime.stream_lark_goal_topic_profile, **options) + try: + with pytest.raises(TimeoutError): + future.result(timeout=0.1) + assert not consumer.waited + finally: + release.set() + with pytest.raises(OSError, match="synthetic reader failure"): + future.result(timeout=8) + assert consumer.waited + assert answers == ["om_alpha_first"] + assert _processed(tmp_path) == 1 + assert not list((tmp_path / "runtime").rglob("om_alpha_waiting.json")) diff --git a/tests/extensions/test_lark_topic_dispatch_bounds.py b/tests/extensions/test_lark_topic_dispatch_bounds.py new file mode 100644 index 0000000000..e392bfdaac --- /dev/null +++ b/tests/extensions/test_lark_topic_dispatch_bounds.py @@ -0,0 +1,135 @@ +"""Bounded transport workers are not a second admission or replay owner.""" + +from concurrent.futures import ThreadPoolExecutor, TimeoutError +import threading + +import pytest + +from loopx.extensions.lark.goal_topic_dispatch import ProfileEventDispatch + + +def test_one_busy_source_and_shared_session_leave_capacity_for_an_independent_role(): + active, release, independent, shared = [threading.Event() for _ in range(4)] + observed = [] + + def handle(line): + observed.append(line) + if line == "slow": + active.set() + assert release.wait(5) + elif line == "other": + independent.set() + elif line == "shared": + shared.set() + + dispatch = ProfileEventDispatch(handle, threading.Event(), workers=2, capacity=8) + try: + assert dispatch.submit(("topic-a", "session-a"), "slow") + assert active.wait(5) + # Rebinds retain source FIFO; a different source sharing the Session + # must wait without using a worker blocked on a per-Session lock. + assert dispatch.submit(("topic-a", "session-new"), "followup") + assert dispatch.submit(("topic-b", "session-a"), "shared") + assert dispatch.submit(("topic-c", "session-c"), "other") + assert independent.wait(3) + assert not shared.is_set() + assert "followup" not in observed + finally: + release.set() + dispatch.close() + assert sorted(observed) == ["followup", "other", "shared", "slow"] + assert observed[:2] == ["slow", "other"] + + +@pytest.mark.parametrize("stop_requested", [False, True]) +def test_capacity_backpressures_and_stop_discards_only_waiting_transport_lines( + stop_requested, +): + active, release, submitting = [threading.Event() for _ in range(3)] + stop = threading.Event() + observed = [] + + def handle(line): + observed.append(line) + if line == "active": + active.set() + assert release.wait(5) + + dispatch = ProfileEventDispatch(handle, stop, workers=1, capacity=2) + try: + assert dispatch.submit(("source",), "active") + assert active.wait(5) + assert dispatch.submit(("source",), "waiting") + + def submit(): + submitting.set() + return dispatch.submit(("other",), "next") + + with ThreadPoolExecutor(max_workers=1) as executor: + future = executor.submit(submit) + try: + assert submitting.wait(5) + with pytest.raises(TimeoutError): + future.result(timeout=0.1) + if stop_requested: + stop.set() + assert future.result(timeout=2) is False + else: + release.set() + assert future.result(timeout=3) is True + finally: + release.set() + finally: + release.set() + dispatch.close() + assert ( + observed == ["active"] + if stop_requested + else sorted(observed) == ["active", "next", "waiting"] + ) + + +def test_handler_failure_cancels_waiting_lines_and_surfaces_without_replay(): + active, release = threading.Event(), threading.Event() + observed = [] + + def handle(line): + observed.append(line) + active.set() + assert release.wait(5) + raise OSError("synthetic transport handler failure") + + dispatch = ProfileEventDispatch(handle, threading.Event(), workers=1, capacity=2) + try: + assert dispatch.submit(("source",), "first") + assert active.wait(5) + assert dispatch.submit(("source",), "waiting") + finally: + release.set() + with pytest.raises(OSError, match="synthetic transport handler failure"): + dispatch.close() + assert observed == ["first"] + + +def test_unplanned_reader_failure_keeps_active_ownership_without_dispatching_waiting_lines(): + active, release = threading.Event(), threading.Event() + observed = [] + + def handle(line): + observed.append(line) + active.set() + assert release.wait(5) + + dispatch = ProfileEventDispatch(handle, threading.Event(), workers=1, capacity=2) + assert dispatch.submit(("source",), "active") + assert active.wait(5) + assert dispatch.submit(("source",), "waiting") + with ThreadPoolExecutor(max_workers=1) as executor: + closing = executor.submit(dispatch.close, cancel_pending=True) + try: + with pytest.raises(TimeoutError): + closing.result(timeout=0.1) + finally: + release.set() + closing.result(timeout=3) + assert observed == ["active"] diff --git a/tests/extensions/test_lark_topic_inbox_config.py b/tests/extensions/test_lark_topic_inbox_config.py new file mode 100644 index 0000000000..eae514edcb --- /dev/null +++ b/tests/extensions/test_lark_topic_inbox_config.py @@ -0,0 +1,95 @@ +"""Concurrent transport callers must publish a complete owner-private inbox config.""" + +import json +import os +import stat +import threading +from concurrent.futures import ThreadPoolExecutor +from pathlib import Path + +import pytest + +from loopx.extensions.lark import goal_topic_runtime as runtime + + +def _options(tmp_path): + return { + "runtime_root": tmp_path, + "route": { + "app_ref": "public-profile", + "target_ref": "public-target", + "topic_root_message_id": "om_public_topic", + "conversation_kind": "manager", + }, + "target_payload": { + "schema_version": "loopx_goal_channel_provider_targets_v0", + "targets": { + "public-target": { + "name": "public-target", + "provider": "lark", + "enabled": True, + "channel": {"chat_id": "oc_public_chat"}, + "identity": { + "sender_profile": "public-profile", + "bot_display_name": "Public Bot", + "bot_app_id": "cli_public_app", + "bot_open_id": "ou_public_bot", + }, + } + }, + }, + } + + +def test_concurrent_topic_config_publish_uses_private_independent_temporaries( + tmp_path, monkeypatch +): + options = _options(tmp_path) + barrier = threading.Barrier(2) + observed = [] + original_os_replace = os.replace + + def before_publish(source): + source = Path(source) + observed.append((source.name, stat.S_IMODE(source.stat().st_mode))) + barrier.wait(timeout=5) + + def os_replace(source, destination): + before_publish(source) + return original_os_replace(source, destination) + + monkeypatch.setattr(os, "replace", os_replace) + with ThreadPoolExecutor(max_workers=2) as executor: + futures = [executor.submit(runtime._inbox_config, **options) for _ in range(2)] + results = [future.result(timeout=8) for future in futures] + + assert results[0] == results[1] + config_path, config_ref = results[0] + assert config_path == tmp_path / config_ref + payload = json.loads(config_path.read_text()) + assert payload["topic_root_message_id"] == "om_public_topic" + assert payload["material_review"]["enabled"] is True + assert payload["reply"]["received_reaction_policy"] == "retain" + assert stat.S_IMODE(config_path.stat().st_mode) == 0o600 + assert len({name for name, _ in observed}) == 2 + assert all(mode == 0o600 for _, mode in observed) + assert list(config_path.parent.glob("*.tmp")) == [] + + +def test_failed_topic_config_publish_preserves_previous_config_and_cleans_temporary( + tmp_path, monkeypatch +): + options = _options(tmp_path) + config_path, _ = runtime._inbox_config(**options) + previous = config_path.read_bytes() + + def fail_publish(*_args, **_kwargs): + raise OSError("synthetic publish failure") + + monkeypatch.setattr(os, "replace", fail_publish) + changed = {**options, "route": {**options["route"], "conversation_kind": "worker"}} + with pytest.raises(OSError, match="synthetic publish failure"): + runtime._inbox_config(**changed) + assert config_path.read_bytes() == previous + assert stat.S_IMODE(config_path.stat().st_mode) == 0o600 + assert list(config_path.parent.glob("*.tmp")) == []