Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,18 @@ without private bindings do not add provider authentication; configured aliases
reuse one request-scoped verified identity observation. Lark HTTP composition
resides in the extension, while the typed binding owner remains provider-neutral.

The same Settings read model includes Core bindings for locally installed
conversation transports. Listener status comes from the exact host-composed
provider's optional `health_snapshot()` returning a content-free `status`, or
from the existing Lark listener observation for default profiles. This local
hook must return cached observations without network calls. Only existing
listener labels (`starting`, `listening`, `retrying`, `stopped`, `standby`,
`inactive`) are projected; missing, malformed or failed observations become
`unknown`, rendered as connection unconfirmed. Raw provider fields and exception
messages are excluded. Listening is neither a delivery receipt nor task
acceptance. Native binding/composition/HTTP regressions qualify the source read
model, while each provider's installed liveness and result return remain separate.

Synthetic product previews: [desktop](../../assets/personal-workspace/private-project-conversations.png),
[narrow](../../assets/personal-workspace/private-project-conversations-narrow.png),
[revoked workspace](../../assets/personal-workspace/private-project-workspace-revoked.png).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,14 @@ provider 读回确认回复;发生没有 receipt 的不确定写入时不盲
时,同一请求内的别名检查与连接共用一次已验证身份观测。飞书 HTTP 组合位于
extension,typed binding owner 继续保持 provider-neutral。

同一设置读模型也展示本机已安装 transport 的 Core binding。监听状态取确切宿主组合
provider 的可选 `health_snapshot()`,返回仅含状态的观测;默认 profile 使用已有
Lark listener 观测。该本地 hook 只读取缓存,不发起网络请求。仅投影已有监听标签
`starting`、`listening`、`retrying`、`stopped`、`standby`、`inactive`;缺失、格式
错误或读取失败均为 `unknown`,界面显示“连接尚未确认”。不返回 provider 其它字段或
异常原文。监听就绪不等于送达回执或任务验收。原生 binding、宿主组合和 HTTP 回归
验证源码读模型;各 provider 的安装态监听及结果回传仍分别验收。

![合成私聊工作区设置](../../assets/personal-workspace/private-project-conversations.png)
![窄屏私聊设置](../../assets/personal-workspace/private-project-conversations-narrow.png)
![工作区撤权读回](../../assets/personal-workspace/private-project-workspace-revoked.png)
Expand Down
24 changes: 23 additions & 1 deletion loopx/extensions/lark/private_conversation_api.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
"""Loopback setup companion for native, owner-only external project Chat."""
from __future__ import annotations
from collections.abc import Mapping
import logging
from pathlib import Path
from typing import Any

Expand All @@ -9,6 +11,26 @@
class PrivateConversationRequestMixin:
server: Any

def _private_listener_status(self, transport_ref: str, lark_health: Mapping[str, Any]) -> str:
"""Project observations from the exact provider, never infer liveness."""
composed = getattr(self.server, "conversation_transports", None)
provider = composed.transports.get(transport_ref) if composed is not None else None
if provider is not None:
snapshot = getattr(provider, "health_snapshot", None)
try:
health = snapshot() if callable(snapshot) else None
except Exception as exc:
# Provider exceptions may include credentials. A failed read
# must not borrow another listener's possibly stale status.
logging.getLogger(__name__).warning("Conversation listener health unavailable: %s", type(exc).__name__)
health = None
else:
health = lark_health.get(transport_ref)
status = health.get("status") if isinstance(health, Mapping) else None
# Reuse the existing listener presentation vocabulary; project no
# other provider fields. The UI already labels unknown as unconfirmed.
return status if isinstance(status, str) and status in ("starting", "listening", "retrying", "stopped", "standby", "inactive") else "unknown"

def _private_conversations(self) -> None:
bindings = self.server.runtime_controller.project_contexts.conversation_bindings
current = bindings.read()
Expand All @@ -26,7 +48,7 @@ def _private_conversations(self) -> None:
"agent_candidates": [{key: item.get(key) for key in ["session_id", "goal_id", "agent_id", "executor_endpoint_id"]}
for item in bindings.agent_candidates(row["binding_id"])],
"agent_targets": row.get("agent_targets", []),
"listener_status": health.get(row["transport_ref"], {}).get("status", "starting"),
"listener_status": self._private_listener_status(row["transport_ref"], health),
**deliveries.get(row["binding_id"], {"pending_count": 0, "recovery_count": 0})}
for row in current["bindings"]]})

Expand Down
63 changes: 63 additions & 0 deletions tests/test_lark_private_conversations.py
Original file line number Diff line number Diff line change
Expand Up @@ -553,6 +553,69 @@ def request(path, body=None):
runtime.close()


def test_http_listener_health_belongs_to_each_bound_transport(ordinary): # noqa: F811
from loopx.chat_server import ChatHTTPServer, ChatRequestHandler
from loopx.capabilities.native_chat.transports import ChatConversationTransports
from test_chat_transport_composition import Transport

store, runtime, provider, lark = connect(ordinary)
external = Transport("external-owner")
composed = ChatConversationTransports(observe_default=lark.bindings.observe, transports=[external])
bindings = ChatConversationBindings(root=store.root, project_contexts=runtime.project_contexts,
observe=composed.observe)
runtime.project_contexts.conversation_bindings = bindings
bindings.configure(transport_ref=external.transport_ref,
project_ref=runtime.project_contexts.available()[0]["project_ref"], executor_endpoint_id="codex")
original = bindings.read()
server = ChatHTTPServer(("127.0.0.1", 0), ChatRequestHandler)
server.chat_store, server.runtime_controller, server.verbose = store, runtime, False
server.lark_private_conversations, server.conversation_transports = lark, composed
# A stale Lark row must not mask the current locally installed provider.
server.lark_goal_topic_runtime = SimpleNamespace(health_snapshot=lambda: {
"notes-app": {"status": "listening"}, "steward-app": {"status": "starting"},
"external-owner": {"status": "stopped"}}, close=lambda: None)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()

def snapshot():
connection = http.client.HTTPConnection(*server.server_address, timeout=15)
try:
connection.request("GET", "/api/chat/lark/private-conversations")
response = connection.getresponse()
assert response.status == 200
return {row["app_ref"]: row for row in json.loads(response.read())["connections"]}
finally:
connection.close()

def unavailable():
raise OSError("provider credential must not be exposed")

try:
for state in ["starting", "listening", "retrying", "stopped", "standby", "inactive"]:
external.health_snapshot = lambda: {"status": state, "credential": "private-provider-data"}
rows = snapshot()
assert rows["external-owner"]["listener_status"] == state
assert rows["notes-app"]["listener_status"] == "listening"
assert rows["steward-app"]["listener_status"] == "starting"
assert "private-provider-data" not in json.dumps(rows)
for malformed in [None, [], {"status": "private-provider-data"}, {"status": []}]:
external.health_snapshot = lambda: malformed
assert snapshot()["external-owner"]["listener_status"] == "unknown"
external.health_snapshot = unavailable
assert snapshot()["external-owner"]["listener_status"] == "unknown"
del external.health_snapshot
assert snapshot()["external-owner"]["listener_status"] == "unknown"
del server.conversation_transports
server.lark_goal_topic_runtime = SimpleNamespace(health_snapshot=lambda: {}, close=lambda: None)
assert all(row["listener_status"] == "unknown" for row in snapshot().values())
assert bindings.read() == original and not store.list_sessions() and not provider.writes
finally:
server.shutdown()
server.server_close()
thread.join(timeout=2)
runtime.close()


def test_bot_only_group_continues_after_user_logout_but_rejects_private_and_app_drift(ordinary): # noqa: F811
store, runtime, provider, transport = connect_group(ordinary)
try:
Expand Down
Loading