Compare commits
20
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
aeb79c6286 | ||
|
|
897e8d36e8 | ||
|
|
459815c0cc | ||
|
|
773c67db8f | ||
|
|
0c15159a5f | ||
|
|
19c52a9425 | ||
|
|
c9a1717e12 | ||
|
|
ca95d47dc8 | ||
|
|
aaaa11a912 | ||
|
|
b6fa3da540 | ||
|
|
15adb30bc7 | ||
|
|
faf9566332 | ||
|
|
b0aeb81caf | ||
|
|
04f78f3f01 | ||
|
|
3fbd31fbff | ||
|
|
a5dc880af1 | ||
|
|
9b9845fd14 | ||
|
|
106048999f | ||
|
|
6afc46734a | ||
|
|
9ee1919985 |
@@ -14,10 +14,10 @@ jobs:
|
|||||||
validate:
|
validate:
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
steps:
|
steps:
|
||||||
|
- uses: actions/checkout@v4
|
||||||
- run: python3 --version
|
- run: python3 --version
|
||||||
- run: node --version
|
- run: node --version
|
||||||
- run: echo "Runner works!"
|
- run: echo "Runner works!"
|
||||||
- run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
|
||||||
|
|
||||||
- name: Python syntax check
|
- name: Python syntax check
|
||||||
run: |
|
run: |
|
||||||
@@ -45,8 +45,8 @@ print('✅ config.yaml.example valid')
|
|||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
needs: [validate]
|
needs: [validate]
|
||||||
steps:
|
steps:
|
||||||
|
- uses: actions/checkout@v4
|
||||||
- run: echo "🚀 Deploy tag $(echo $GITHUB_REF_NAME)"
|
- run: echo "🚀 Deploy tag $(echo $GITHUB_REF_NAME)"
|
||||||
- run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
|
||||||
|
|
||||||
- name: Deploy to Tanko (canary)
|
- name: Deploy to Tanko (canary)
|
||||||
run: ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 jerome@192.168.68.122 "cd /root && bash -s" < scripts/deploy.sh --ct=tanko --mode=native "$GITHUB_REF_NAME" 2>&1 || echo "⚠️ Tanko deploy skipped"
|
run: ssh -o StrictHostKeyChecking=no -o ConnectTimeout=10 jerome@192.168.68.122 "cd /root && bash -s" < scripts/deploy.sh --ct=tanko --mode=native "$GITHUB_REF_NAME" 2>&1 || echo "⚠️ Tanko deploy skipped"
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ jobs:
|
|||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Python syntax
|
- name: Python syntax
|
||||||
run: |
|
run: |
|
||||||
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null
|
python3 -m py_compile plugins/platforms/zulip/adapter.py 2>/dev/null
|
||||||
@@ -37,7 +37,7 @@ jobs:
|
|||||||
environment: canary
|
environment: canary
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to Tanko
|
- name: Deploy to Tanko
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
@@ -66,7 +66,7 @@ jobs:
|
|||||||
environment: production
|
environment: production
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to all Hermes agents
|
- name: Deploy to all Hermes agents
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
@@ -97,7 +97,7 @@ jobs:
|
|||||||
environment: agent-zero
|
environment: agent-zero
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
run: git clone --depth 1 http://abiba-bot:mipjoq-tybbox-2ciHru@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
run: git clone --depth 1 http://abiba-bot:***REMOVED***@192.168.68.17:3000/SyslogSolution/zulip-platform-plugins.git .
|
||||||
- name: Deploy to kagentz
|
- name: Deploy to kagentz
|
||||||
env:
|
env:
|
||||||
TAG: ${{ github.ref_name }}
|
TAG: ${{ github.ref_name }}
|
||||||
|
|||||||
@@ -203,6 +203,22 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
logger.error("Zulip POST %s network error: %s", path, exc)
|
logger.error("Zulip POST %s network error: %s", path, exc)
|
||||||
return {"result": "error", "msg": str(exc)}
|
return {"result": "error", "msg": str(exc)}
|
||||||
|
|
||||||
|
async def _api_patch(self, path: str, payload: Dict[str, Any]) -> Dict[str, Any]:
|
||||||
|
import aiohttp
|
||||||
|
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
||||||
|
try:
|
||||||
|
async with self._session.patch(
|
||||||
|
url, data=payload, auth=self._auth(),
|
||||||
|
timeout=aiohttp.ClientTimeout(total=15),
|
||||||
|
) as resp:
|
||||||
|
data = await resp.json()
|
||||||
|
if resp.status >= 400:
|
||||||
|
logger.debug("Zulip PATCH %s -> %s: %s", path, resp.status, str(data)[:200])
|
||||||
|
return data
|
||||||
|
except (aiohttp.ClientError, asyncio.TimeoutError) as exc:
|
||||||
|
logger.debug("Zulip PATCH %s network error: %s", path, exc)
|
||||||
|
return {"result": "error", "msg": str(exc)}
|
||||||
|
|
||||||
async def _api_delete(self, path: str, params: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
async def _api_delete(self, path: str, params: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
||||||
import aiohttp
|
import aiohttp
|
||||||
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
url = f"{self._site}/api/v1/{path.lstrip('/')}"
|
||||||
@@ -326,6 +342,23 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
|
|
||||||
return SendResult(success=True, message_id=str(last_id) if last_id else None)
|
return SendResult(success=True, message_id=str(last_id) if last_id else None)
|
||||||
|
|
||||||
|
async def edit_message(
|
||||||
|
self,
|
||||||
|
chat_id: str,
|
||||||
|
message_id: str,
|
||||||
|
content: str,
|
||||||
|
**kwargs,
|
||||||
|
) -> SendResult:
|
||||||
|
"""Edit a previously-sent message. Used by the gateway stream consumer
|
||||||
|
for progressive message updates during agent streaming."""
|
||||||
|
formatted = self.format_message(content)
|
||||||
|
data = await self._api_patch(f"messages/{message_id}", {"content": formatted})
|
||||||
|
if data.get("result") != "success":
|
||||||
|
msg = str(data.get("msg", "edit failed"))
|
||||||
|
logger.debug("Zulip edit_message(%s) -> %s", message_id, msg)
|
||||||
|
return SendResult(success=False, error=msg)
|
||||||
|
return SendResult(success=True, message_id=message_id)
|
||||||
|
|
||||||
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
||||||
kind, to, topic = _parse_target(chat_id, self._default_topic)
|
kind, to, topic = _parse_target(chat_id, self._default_topic)
|
||||||
if kind == "direct":
|
if kind == "direct":
|
||||||
|
|||||||
@@ -43,7 +43,6 @@ import time
|
|||||||
import uuid
|
import uuid
|
||||||
from collections import deque
|
from collections import deque
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from pathlib import Path
|
|
||||||
from typing import Any, Dict, List, Optional, Tuple
|
from typing import Any, Dict, List, Optional, Tuple
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -68,6 +67,7 @@ DEFAULT_STREAM = "agent-hub"
|
|||||||
DEFAULT_ALL_BOTS_USER_ID = 1
|
DEFAULT_ALL_BOTS_USER_ID = 1
|
||||||
DEFAULT_POLL_INTERVAL = 3.0
|
DEFAULT_POLL_INTERVAL = 3.0
|
||||||
MAX_ZULIP_MESSAGE = 10000
|
MAX_ZULIP_MESSAGE = 10000
|
||||||
|
TRUNCATION_NOTICE = "\n\n[...truncated at Zulip limit]"
|
||||||
ECHO_TAG_PREFIX = "hermes-zulip-"
|
ECHO_TAG_PREFIX = "hermes-zulip-"
|
||||||
RECONNECT_BACKOFF = [2, 5, 10, 30, 60]
|
RECONNECT_BACKOFF = [2, 5, 10, 30, 60]
|
||||||
DEDUP_WINDOW = 300 # 5 minutes
|
DEDUP_WINDOW = 300 # 5 minutes
|
||||||
@@ -121,7 +121,7 @@ def _truncate(text: str, limit: int = MAX_ZULIP_MESSAGE) -> str:
|
|||||||
"""Truncate to Zulip's message limit with notice."""
|
"""Truncate to Zulip's message limit with notice."""
|
||||||
if len(text) <= limit:
|
if len(text) <= limit:
|
||||||
return text
|
return text
|
||||||
return text[:limit] + "\n\n[...truncated at Zulip limit]"
|
return text[:limit - len(TRUNCATION_NOTICE)] + TRUNCATION_NOTICE
|
||||||
|
|
||||||
|
|
||||||
def _parse_zulip_timestamp(ts: Any) -> datetime:
|
def _parse_zulip_timestamp(ts: Any) -> datetime:
|
||||||
@@ -882,116 +882,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
reply_to_topic=topic,
|
reply_to_topic=topic,
|
||||||
)
|
)
|
||||||
|
|
||||||
# ------------------------------------------------------------------
|
|
||||||
# Media / attachment extraction (vision support)
|
|
||||||
# ------------------------------------------------------------------
|
|
||||||
|
|
||||||
async def _extract_inline_media(
|
|
||||||
self, content: str, media_urls: List[str], media_types: List[str],
|
|
||||||
) -> None:
|
|
||||||
"""Download Zulip ``/user_uploads/...`` attachments referenced in a
|
|
||||||
message so vision/transcription tools can read them locally.
|
|
||||||
|
|
||||||
Zulip renders uploaded files as ``<img src="/user_uploads/...">``
|
|
||||||
in HTML. After ``_strip_html()`` the markdown-style references
|
|
||||||
survive in the plain text. This method finds those references,
|
|
||||||
downloads the files, and caches them into Hermes media storage.
|
|
||||||
"""
|
|
||||||
for rel in re.findall(r"\]\((/user_uploads/[^)]+)\)", content):
|
|
||||||
dl_url = f"{self._site}{rel}"
|
|
||||||
try:
|
|
||||||
resp = await self._http_client.get(
|
|
||||||
dl_url,
|
|
||||||
headers={"Authorization": self._auth_header},
|
|
||||||
timeout=30.0,
|
|
||||||
)
|
|
||||||
if resp.status_code >= 400:
|
|
||||||
logger.debug(
|
|
||||||
"[%s] Inline media %s -> %d",
|
|
||||||
self.name, rel, resp.status_code,
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
data = resp.content
|
|
||||||
mime = (
|
|
||||||
resp.headers.get("content-type", "application/octet-stream")
|
|
||||||
.split(";")[0]
|
|
||||||
.strip()
|
|
||||||
)
|
|
||||||
except Exception:
|
|
||||||
logger.debug("[%s] Failed to download inline media %s", self.name, rel)
|
|
||||||
continue
|
|
||||||
fname = rel.rsplit("/", 1)[-1] or "file"
|
|
||||||
await self._cache_attachment(data, fname, mime, media_urls, media_types)
|
|
||||||
|
|
||||||
async def _extract_event_attachments(
|
|
||||||
self, message: Dict[str, Any], media_urls: List[str], media_types: List[str],
|
|
||||||
) -> None:
|
|
||||||
"""Download files from Zulip's ``message.attachments`` field (UI file
|
|
||||||
uploads that Zulip sends as structured event data rather than inline
|
|
||||||
markdown links)."""
|
|
||||||
atts = message.get("attachments") or []
|
|
||||||
for att in atts:
|
|
||||||
path_id = att.get("path_id")
|
|
||||||
fname = att.get("name") or "file"
|
|
||||||
if not path_id:
|
|
||||||
continue
|
|
||||||
dl_url = f"{self._site}/api/v1/user_uploads/{path_id}"
|
|
||||||
try:
|
|
||||||
resp = await self._http_client.get(
|
|
||||||
dl_url,
|
|
||||||
headers={"Authorization": self._auth_header},
|
|
||||||
timeout=30.0,
|
|
||||||
)
|
|
||||||
if resp.status_code >= 400:
|
|
||||||
logger.debug(
|
|
||||||
"[%s] Event attachment %s -> %d",
|
|
||||||
self.name, fname, resp.status_code,
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
data = resp.content
|
|
||||||
mime = (
|
|
||||||
resp.headers.get("content-type", "application/octet-stream")
|
|
||||||
.split(";")[0]
|
|
||||||
.strip()
|
|
||||||
)
|
|
||||||
except Exception:
|
|
||||||
logger.debug("[%s] Failed to download event attachment %s", self.name, fname)
|
|
||||||
continue
|
|
||||||
await self._cache_attachment(data, fname, mime, media_urls, media_types)
|
|
||||||
|
|
||||||
async def _cache_attachment(
|
|
||||||
self, data: bytes, fname: str, mime: str,
|
|
||||||
media_urls: List[str], media_types: List[str],
|
|
||||||
) -> None:
|
|
||||||
"""Cache a downloaded attachment into Hermes media storage so the
|
|
||||||
vision_analyze / ocr tools can access it via a local file path."""
|
|
||||||
try:
|
|
||||||
from gateway.platforms.base import (
|
|
||||||
cache_image_from_bytes,
|
|
||||||
cache_audio_from_bytes,
|
|
||||||
cache_document_from_bytes,
|
|
||||||
)
|
|
||||||
except ImportError:
|
|
||||||
logger.warning("[%s] cache_*_from_bytes not available in gateway", self.name)
|
|
||||||
return
|
|
||||||
ext = Path(fname).suffix
|
|
||||||
try:
|
|
||||||
if mime.startswith("image/"):
|
|
||||||
media_urls.append(cache_image_from_bytes(data, ext or ".png"))
|
|
||||||
media_types.append(mime)
|
|
||||||
elif mime.startswith("audio/"):
|
|
||||||
media_urls.append(cache_audio_from_bytes(data, ext or ".ogg"))
|
|
||||||
media_types.append(mime)
|
|
||||||
else:
|
|
||||||
media_urls.append(cache_document_from_bytes(data, fname))
|
|
||||||
media_types.append(mime)
|
|
||||||
except Exception as exc:
|
|
||||||
logger.warning("[%s] Failed to cache attachment %s: %s", self.name, fname, exc)
|
|
||||||
|
|
||||||
# ------------------------------------------------------------------
|
|
||||||
# Message routing
|
|
||||||
# ------------------------------------------------------------------
|
|
||||||
|
|
||||||
async def _route_message(
|
async def _route_message(
|
||||||
self,
|
self,
|
||||||
text: str,
|
text: str,
|
||||||
@@ -1022,21 +912,7 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
chat_id_alt=str(reply_to_id) if reply_to_id else None,
|
chat_id_alt=str(reply_to_id) if reply_to_id else None,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Gen 5: Extract media/attachments so vision tools can read them
|
# Derive message type
|
||||||
media_urls: List[str] = []
|
|
||||||
media_types: List[str] = []
|
|
||||||
await self._extract_inline_media(text, media_urls, media_types)
|
|
||||||
await self._extract_event_attachments(raw, media_urls, media_types)
|
|
||||||
|
|
||||||
# Derive message type — upgrade to PHOTO/VOICe/DOCUMENT when media present
|
|
||||||
if media_types:
|
|
||||||
if any(m.startswith("image/") for m in media_types):
|
|
||||||
message_type = MessageType.PHOTO
|
|
||||||
elif any(m.startswith("audio/") for m in media_types):
|
|
||||||
message_type = MessageType.VOICE
|
|
||||||
else:
|
|
||||||
message_type = MessageType.DOCUMENT
|
|
||||||
else:
|
|
||||||
message_type = MessageType.TEXT
|
message_type = MessageType.TEXT
|
||||||
|
|
||||||
message_event = MessageEvent(
|
message_event = MessageEvent(
|
||||||
@@ -1046,8 +922,6 @@ class ZulipAdapter(BasePlatformAdapter):
|
|||||||
message_id=msg_id,
|
message_id=msg_id,
|
||||||
raw_message=raw,
|
raw_message=raw,
|
||||||
timestamp=_parse_zulip_timestamp(timestamp),
|
timestamp=_parse_zulip_timestamp(timestamp),
|
||||||
media_urls=media_urls or None,
|
|
||||||
media_types=media_types or None,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
# Store reply metadata so send() can respond to the right thread
|
# Store reply metadata so send() can respond to the right thread
|
||||||
|
|||||||
Reference in New Issue
Block a user