From b317387e0cae4762ecb42ff12f2149ba9b9abc81 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Wed, 13 May 2026 16:12:46 +0800 Subject: [PATCH] =?UTF-8?q?feat(nextcloudtalk):=20=E6=96=B0=E5=A2=9E?= =?UTF-8?q?=E6=8A=95=E7=A5=A8=E5=92=8C=E6=B6=88=E6=81=AF=E7=BD=AE=E9=A1=B6?= =?UTF-8?q?=E7=9B=B8=E5=85=B3=E5=8A=9F=E8=83=BD=EF=BC=8C=E4=BC=98=E5=8C=96?= =?UTF-8?q?=E4=BB=A3=E7=A0=81=E7=BB=93=E6=9E=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 新增close_poll、get_poll_results方法实现投票管理 2. 新增pin_message、unpin_message、list_pins实现消息置顶功能 3. 调整导入顺序优化代码整洁度 4. 修复去重逻辑的判断顺序问题 5. 清理多余的空行和导入冗余 --- .../adapters/nextcloudtalk/adapter.py | 162 ++++++++++++++++-- .../channels/adapters/nextcloudtalk/dedup.py | 4 +- .../adapters/nextcloudtalk/normalizer.py | 3 +- .../adapters/nextcloudtalk/secret_contract.py | 1 - 4 files changed, 151 insertions(+), 19 deletions(-) diff --git a/backend/package/yuxi/channels/adapters/nextcloudtalk/adapter.py b/backend/package/yuxi/channels/adapters/nextcloudtalk/adapter.py index ac661de6..ad9fc42a 100644 --- a/backend/package/yuxi/channels/adapters/nextcloudtalk/adapter.py +++ b/backend/package/yuxi/channels/adapters/nextcloudtalk/adapter.py @@ -11,11 +11,11 @@ from urllib.parse import urlparse from yuxi.channels.base import BaseChannelAdapter from yuxi.channels.capabilities import ChannelCapabilities -from yuxi.channels.meta import ChannelMeta -from yuxi.channels.infra.circuit_breaker import CircuitBreaker, CircuitBreakerOpenError from yuxi.channels.exceptions import ( ChannelAuthenticationError, ) +from yuxi.channels.infra.circuit_breaker import CircuitBreaker, CircuitBreakerOpenError +from yuxi.channels.meta import ChannelMeta from yuxi.channels.models import ( ChannelMessage, ChannelResponse, @@ -30,12 +30,12 @@ from yuxi.utils.datetime_utils import utc_now_naive from .client import NextcloudTalkClient, verify_hmac_signature from .config import ( - resolve_room_enabled, - resolve_room_system_prompt, - resolve_room_skills, - resolve_dm_system_prompt, resolve_dm_enabled, resolve_dm_skills, + resolve_dm_system_prompt, + resolve_room_enabled, + resolve_room_skills, + resolve_room_system_prompt, ) from .dedup import NextcloudTalkDedupGuard from .formatter import format_outbound @@ -47,15 +47,21 @@ from .pairing import send_pairing_challenge from .probe import probe_capabilities, resolve_api_base, resolve_api_version from .security import check_dm_policy, check_group_policy, check_mention_gate, collect_security_warnings from .send import ( - send_media, - send_reaction, - send_text, - send_reaction_delete, - get_reactions as _get_reactions, - edit_message as _send_edit, delete_message as _send_delete, ) -from .session import resolve_session_route, make_thread_key +from .send import ( + edit_message as _send_edit, +) +from .send import ( + get_reactions as _get_reactions, +) +from .send import ( + send_media, + send_reaction, + send_reaction_delete, + send_text, +) +from .session import make_thread_key, resolve_session_route logger = logging.getLogger(__name__) @@ -109,6 +115,11 @@ class NextcloudTalkAdapter(BaseChannelAdapter): media=True, block_streaming=True, polls=True, + approval=True, + vision=True, + pin=True, + unpin=True, + list_pins=True, ) meta = ChannelMeta( id="nextcloud-talk", @@ -790,6 +801,131 @@ class NextcloudTalkAdapter(BaseChannelAdapter): except Exception as e: return DeliveryResult(success=False, error=str(e)) + async def close_poll(self, chat_id: str, poll_id: str) -> DeliveryResult: + if not self._client: + return DeliveryResult(success=False, error="Channel not connected") + + async def _do_close(): + path = f"{self._api_base}/poll/{chat_id}/{poll_id}" + try: + result = await self._client.delete(path) + ocs = result.get("ocs", {}) + if ocs.get("meta", {}).get("status") == "ok": + return DeliveryResult(success=True, message_id=poll_id) + return DeliveryResult( + success=False, + error=ocs.get("meta", {}).get("message", "Failed to close poll"), + ) + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + try: + return await self._circuit_breaker.call(_do_close) + except CircuitBreakerOpenError: + return DeliveryResult(success=False, error="Circuit breaker open") + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + async def get_poll_results(self, chat_id: str, poll_id: str) -> DeliveryResult: + if not self._client: + return DeliveryResult(success=False, error="Channel not connected") + + async def _do_get(): + path = f"{self._api_base}/poll/{chat_id}/{poll_id}" + try: + result = await self._client.get(path) + ocs = result.get("ocs", {}) + if ocs.get("meta", {}).get("status") == "ok": + data = ocs.get("data", {}) + return DeliveryResult( + success=True, + metadata={ + "question": data.get("question", ""), + "options": data.get("options", []), + "votes": data.get("votes", {}), + "numVoters": data.get("numVoters", 0), + "status": data.get("status", "open"), + }, + ) + return DeliveryResult( + success=False, + error=ocs.get("meta", {}).get("message", "Failed to get poll results"), + ) + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + try: + return await self._circuit_breaker.call(_do_get) + except CircuitBreakerOpenError: + return DeliveryResult(success=False, error="Circuit breaker open") + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + async def pin_message(self, chat_id: str, msg_id: str) -> DeliveryResult: + if not self._client: + return DeliveryResult(success=False, error="Channel not connected") + + async def _do_pin(): + path = f"{self._api_base}/pin/{chat_id}/{msg_id}" + try: + result = await self._client.post(path) + ocs = result.get("ocs", {}) + if ocs.get("meta", {}).get("status") == "ok": + return DeliveryResult(success=True, message_id=msg_id) + return DeliveryResult( + success=False, + error=ocs.get("meta", {}).get("message", "Failed to pin message"), + ) + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + try: + return await self._circuit_breaker.call(_do_pin) + except CircuitBreakerOpenError: + return DeliveryResult(success=False, error="Circuit breaker open") + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + async def unpin_message(self, chat_id: str, msg_id: str) -> DeliveryResult: + if not self._client: + return DeliveryResult(success=False, error="Channel not connected") + + async def _do_unpin(): + path = f"{self._api_base}/pin/{chat_id}/{msg_id}" + try: + result = await self._client.delete(path) + ocs = result.get("ocs", {}) + if ocs.get("meta", {}).get("status") == "ok": + return DeliveryResult(success=True, message_id=msg_id) + return DeliveryResult( + success=False, + error=ocs.get("meta", {}).get("message", "Failed to unpin message"), + ) + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + try: + return await self._circuit_breaker.call(_do_unpin) + except CircuitBreakerOpenError: + return DeliveryResult(success=False, error="Circuit breaker open") + except Exception as e: + return DeliveryResult(success=False, error=str(e)) + + async def list_pins(self, chat_id: str) -> list[dict]: + if not self._client: + return [] + path = f"{self._api_base}/pin/{chat_id}" + try: + result = await self._client.get(path) + ocs = result.get("ocs", {}) + if ocs.get("meta", {}).get("status") == "ok": + data = ocs.get("data", []) + if isinstance(data, list): + return data + return [] + except Exception: + return [] + async def list_conversations(self) -> list[dict[str, Any]]: if not self._client: return [] diff --git a/backend/package/yuxi/channels/adapters/nextcloudtalk/dedup.py b/backend/package/yuxi/channels/adapters/nextcloudtalk/dedup.py index 735cd2bf..364b0796 100644 --- a/backend/package/yuxi/channels/adapters/nextcloudtalk/dedup.py +++ b/backend/package/yuxi/channels/adapters/nextcloudtalk/dedup.py @@ -108,10 +108,8 @@ class NextcloudTalkDedupGuard: if not token or not message_id: return False self._load_persisted() - key = self._make_key(token, message_id) - if key in self._committed: - return True self._gc() + key = self._make_key(token, message_id) return key in self._committed def stats(self) -> dict[str, int]: diff --git a/backend/package/yuxi/channels/adapters/nextcloudtalk/normalizer.py b/backend/package/yuxi/channels/adapters/nextcloudtalk/normalizer.py index 9e9b8da6..967bcca1 100644 --- a/backend/package/yuxi/channels/adapters/nextcloudtalk/normalizer.py +++ b/backend/package/yuxi/channels/adapters/nextcloudtalk/normalizer.py @@ -1,6 +1,6 @@ from __future__ import annotations -from datetime import datetime, UTC +from datetime import UTC, datetime from typing import Any from yuxi.channels.models import ( @@ -14,7 +14,6 @@ from yuxi.channels.models import ( MessageType, ) - _AS2_TYPE_MAP: dict[str, EventType] = { "Create": EventType.MESSAGE_RECEIVED, "Announce": EventType.MESSAGE_RECEIVED, diff --git a/backend/package/yuxi/channels/adapters/nextcloudtalk/secret_contract.py b/backend/package/yuxi/channels/adapters/nextcloudtalk/secret_contract.py index 4c250509..ca2f682a 100644 --- a/backend/package/yuxi/channels/adapters/nextcloudtalk/secret_contract.py +++ b/backend/package/yuxi/channels/adapters/nextcloudtalk/secret_contract.py @@ -4,7 +4,6 @@ import logging import os from typing import Any - logger = logging.getLogger(__name__) SECRET_TARGETS = [