Add email qq and signal
This commit is contained in:
@@ -1,7 +1,7 @@
|
||||
"""Communication channels for EvoScientist.
|
||||
|
||||
This module provides an extensible interface for different messaging channels
|
||||
(iMessage, Telegram, Discord, Slack, WeChat, DingTalk, Feishu) to communicate with the EvoScientist agent.
|
||||
(iMessage, Telegram, Discord, Slack, WeChat, DingTalk, Feishu, Email, QQ, Signal) to communicate with the EvoScientist agent.
|
||||
"""
|
||||
|
||||
from .base import Channel, RawIncoming, IncomingMessage, OutgoingMessage, chunk_text
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
"""Email channel for EvoScientist.
|
||||
|
||||
Uses IMAP polling for inbound + SMTP for outbound. Pure Python, no extra deps.
|
||||
|
||||
Usage in config:
|
||||
channel_enabled = "email"
|
||||
email_imap_host = "imap.gmail.com"
|
||||
email_smtp_host = "smtp.gmail.com"
|
||||
...
|
||||
"""
|
||||
|
||||
from .channel import EmailChannel, EmailConfig
|
||||
from ..channel_manager import register_channel, _parse_csv
|
||||
|
||||
__all__ = ["EmailChannel", "EmailConfig"]
|
||||
|
||||
|
||||
def create_from_config(config) -> EmailChannel:
|
||||
allowed = _parse_csv(config.email_allowed_senders)
|
||||
return EmailChannel(EmailConfig(
|
||||
imap_host=config.email_imap_host,
|
||||
imap_port=config.email_imap_port,
|
||||
imap_username=config.email_imap_username,
|
||||
imap_password=config.email_imap_password,
|
||||
imap_mailbox=config.email_imap_mailbox,
|
||||
imap_use_ssl=config.email_imap_use_ssl,
|
||||
smtp_host=config.email_smtp_host,
|
||||
smtp_port=config.email_smtp_port,
|
||||
smtp_username=config.email_smtp_username,
|
||||
smtp_password=config.email_smtp_password,
|
||||
smtp_use_tls=config.email_smtp_use_tls,
|
||||
from_address=config.email_from_address,
|
||||
poll_interval=config.email_poll_interval,
|
||||
mark_seen=config.email_mark_seen,
|
||||
max_body_chars=config.email_max_body_chars,
|
||||
subject_prefix=config.email_subject_prefix,
|
||||
allowed_senders=allowed,
|
||||
))
|
||||
|
||||
|
||||
register_channel("email", create_from_config)
|
||||
@@ -0,0 +1,366 @@
|
||||
"""Email channel implementation using IMAP + SMTP."""
|
||||
|
||||
import asyncio
|
||||
import email as email_lib
|
||||
import email.utils
|
||||
import html
|
||||
import imaplib
|
||||
import logging
|
||||
import re
|
||||
import smtplib
|
||||
import ssl
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from email import encoders
|
||||
from email.header import decode_header, make_header
|
||||
from email.message import EmailMessage
|
||||
from email.mime.base import MIMEBase
|
||||
from email.mime.multipart import MIMEMultipart
|
||||
from email.mime.text import MIMEText
|
||||
from email.utils import parseaddr
|
||||
from pathlib import Path
|
||||
|
||||
from ..base import Channel, RawIncoming, ChannelError
|
||||
from ..capabilities import EMAIL as EMAIL_CAPS
|
||||
from ..mixins import PollingMixin
|
||||
from ..config import BaseChannelConfig
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _decode_hdr(raw: str) -> str:
|
||||
try:
|
||||
return str(make_header(decode_header(raw))) if raw else ""
|
||||
except Exception:
|
||||
return raw or ""
|
||||
|
||||
|
||||
def _strip_html(text: str) -> str:
|
||||
text = re.sub(r"<br\s*/?>", "\n", text, flags=re.I)
|
||||
text = re.sub(r"<p[^>]*>", "\n", text, flags=re.I)
|
||||
text = re.sub(r"</p>", "\n", text, flags=re.I)
|
||||
text = re.sub(r"<[^>]+>", "", text)
|
||||
return html.unescape(text).strip()
|
||||
|
||||
|
||||
@dataclass
|
||||
class EmailConfig(BaseChannelConfig):
|
||||
imap_host: str = ""
|
||||
imap_port: int = 993
|
||||
imap_username: str = ""
|
||||
imap_password: str = ""
|
||||
imap_mailbox: str = "INBOX"
|
||||
imap_use_ssl: bool = True
|
||||
smtp_host: str = ""
|
||||
smtp_port: int = 587
|
||||
smtp_username: str = ""
|
||||
smtp_password: str = ""
|
||||
smtp_use_tls: bool = True
|
||||
from_address: str = ""
|
||||
poll_interval: int = 30
|
||||
mark_seen: bool = True
|
||||
max_body_chars: int = 12000
|
||||
subject_prefix: str = "Re: "
|
||||
text_chunk_limit: int = 4096
|
||||
|
||||
|
||||
class EmailChannel(Channel, PollingMixin):
|
||||
"""Email channel using IMAP polling + SMTP."""
|
||||
|
||||
name = "email"
|
||||
|
||||
capabilities = EMAIL_CAPS
|
||||
_non_retryable_patterns = ("auth", "login", "credential")
|
||||
|
||||
def __init__(self, config: EmailConfig):
|
||||
super().__init__(config)
|
||||
self._imap: imaplib.IMAP4_SSL | imaplib.IMAP4 | None = None
|
||||
|
||||
async def start(self) -> None:
|
||||
cfg = self.config
|
||||
if not cfg.imap_host or not cfg.imap_username:
|
||||
raise ChannelError("Email imap_host and imap_username are required")
|
||||
loop = asyncio.get_event_loop()
|
||||
await loop.run_in_executor(None, self._connect_imap)
|
||||
self._running = True
|
||||
logger.info(f"Email channel started (IMAP: {cfg.imap_host}, poll {cfg.poll_interval}s)")
|
||||
await self._start_polling()
|
||||
|
||||
def _connect_imap(self) -> None:
|
||||
cfg = self.config
|
||||
try:
|
||||
if cfg.imap_use_ssl:
|
||||
self._imap = imaplib.IMAP4_SSL(cfg.imap_host, cfg.imap_port, ssl_context=ssl.create_default_context())
|
||||
else:
|
||||
self._imap = imaplib.IMAP4(cfg.imap_host, cfg.imap_port)
|
||||
self._imap.login(cfg.imap_username, cfg.imap_password)
|
||||
self._imap.select(cfg.imap_mailbox)
|
||||
except Exception as e:
|
||||
raise ChannelError(f"IMAP failed: {e}")
|
||||
|
||||
def _reconnect_imap(self) -> None:
|
||||
try:
|
||||
if self._imap:
|
||||
self._imap.noop()
|
||||
return
|
||||
except Exception:
|
||||
pass
|
||||
self._connect_imap()
|
||||
|
||||
async def _poll_once(self) -> None:
|
||||
loop = asyncio.get_event_loop()
|
||||
messages = await loop.run_in_executor(None, self._fetch_unseen)
|
||||
for m in messages:
|
||||
await self._process_email(m)
|
||||
|
||||
def _fetch_unseen(self) -> list[dict]:
|
||||
self._reconnect_imap()
|
||||
results = []
|
||||
try:
|
||||
st, data = self._imap.search(None, "UNSEEN")
|
||||
if st != "OK":
|
||||
return []
|
||||
for mid in data[0].split()[-20:]:
|
||||
st, msg_data = self._imap.fetch(mid, "(RFC822)")
|
||||
if st != "OK":
|
||||
continue
|
||||
msg = email_lib.message_from_bytes(msg_data[0][1])
|
||||
from_name, from_addr = parseaddr(msg.get("From", ""))
|
||||
body = self._extract_body(msg)
|
||||
if len(body) > self.config.max_body_chars:
|
||||
body = body[:self.config.max_body_chars] + "\n[...truncated]"
|
||||
# Extract attachments and inline images
|
||||
attachments = []
|
||||
if msg.is_multipart():
|
||||
for part in msg.walk():
|
||||
content_disp = part.get("Content-Disposition") or ""
|
||||
content_type = part.get_content_type() or ""
|
||||
is_attachment = "attachment" in content_disp.lower()
|
||||
is_inline_image = (
|
||||
"inline" in content_disp.lower()
|
||||
and content_type.startswith("image/")
|
||||
)
|
||||
# Also detect non-text parts with a filename but no
|
||||
# Content-Disposition header (common for PDFs, docs,
|
||||
# etc. sent by some email clients).
|
||||
is_named_file = (
|
||||
not is_attachment
|
||||
and not is_inline_image
|
||||
and part.get_filename()
|
||||
and not content_type.startswith("multipart/")
|
||||
and not content_type.startswith("text/")
|
||||
)
|
||||
if is_attachment or is_inline_image or is_named_file:
|
||||
filename = part.get_filename() or "attachment"
|
||||
filename = _decode_hdr(filename)
|
||||
payload_data = part.get_payload(decode=True)
|
||||
if payload_data:
|
||||
from ..base import MAX_ATTACHMENT_BYTES, MEDIA_DIR
|
||||
if len(payload_data) > MAX_ATTACHMENT_BYTES:
|
||||
attachments.append({"annotation": f"[attachment: {filename} - too large ({len(payload_data)} bytes)]"})
|
||||
else:
|
||||
MEDIA_DIR.mkdir(parents=True, exist_ok=True)
|
||||
local_path = MEDIA_DIR / f"email_{mid.decode()}_{filename}"
|
||||
local_path.write_bytes(payload_data)
|
||||
label = "inline-image" if is_inline_image else "attachment"
|
||||
attachments.append({"path": str(local_path), "annotation": f"[{label}: {local_path}]"})
|
||||
if self.config.mark_seen:
|
||||
self._imap.store(mid, "+FLAGS", "\\Seen")
|
||||
results.append({
|
||||
"from_addr": from_addr, "from_name": _decode_hdr(from_name),
|
||||
"subject": _decode_hdr(msg.get("Subject", "")), "body": body,
|
||||
"message_id": msg.get("Message-ID", ""), "date": msg.get("Date", ""),
|
||||
"references": msg.get("References", ""), "attachments": attachments,
|
||||
})
|
||||
except Exception as e:
|
||||
logger.error(f"IMAP fetch: {e}")
|
||||
return results
|
||||
|
||||
def _extract_body(self, msg) -> str:
|
||||
if msg.is_multipart():
|
||||
for part in msg.walk():
|
||||
ct = part.get_content_type()
|
||||
if ct == "text/plain":
|
||||
return self._decode_payload(part)
|
||||
for part in msg.walk():
|
||||
if part.get_content_type() == "text/html":
|
||||
return _strip_html(self._decode_payload(part))
|
||||
return "[no text content]"
|
||||
text = self._decode_payload(msg)
|
||||
return _strip_html(text) if msg.get_content_type() == "text/html" else text
|
||||
|
||||
@staticmethod
|
||||
def _decode_payload(part) -> str:
|
||||
payload = part.get_payload(decode=True)
|
||||
if not payload:
|
||||
return ""
|
||||
charset = part.get_content_charset() or "utf-8"
|
||||
return payload.decode(charset, errors="replace")
|
||||
|
||||
async def _process_email(self, m: dict) -> None:
|
||||
subject = m["subject"]
|
||||
text = f"[邮件] 主题: {subject}\n\n{m['body']}" if subject else m["body"]
|
||||
try:
|
||||
ts = email_lib.utils.parsedate_to_datetime(m["date"])
|
||||
except Exception:
|
||||
ts = datetime.now()
|
||||
# Process attachments
|
||||
media_paths: list[str] = []
|
||||
annotations: list[str] = []
|
||||
for att in m.get("attachments", []):
|
||||
if att.get("path"):
|
||||
media_paths.append(att["path"])
|
||||
if att.get("annotation"):
|
||||
annotations.append(att["annotation"])
|
||||
await self._enqueue_raw(RawIncoming(
|
||||
sender_id=m["from_addr"], chat_id=m["from_addr"], text=text, timestamp=ts,
|
||||
message_id=m["message_id"],
|
||||
media_files=media_paths,
|
||||
content_annotations=annotations,
|
||||
metadata={"chat_id": m["from_addr"], "subject": subject,
|
||||
"original_message_id": m["message_id"], "references": m["references"], "backend": "email"},
|
||||
))
|
||||
|
||||
# ── Send ──────────────────────────────────────────────────────
|
||||
|
||||
def _is_ready(self) -> bool:
|
||||
return bool(self.config.smtp_host)
|
||||
|
||||
async def _send_chunk(self, chat_id, formatted_text, raw_text, reply_to, metadata):
|
||||
loop = asyncio.get_event_loop()
|
||||
try:
|
||||
await loop.run_in_executor(
|
||||
None, self._smtp_send_html, chat_id, formatted_text, raw_text, metadata or {},
|
||||
)
|
||||
except Exception as e:
|
||||
err_str = str(e).lower()
|
||||
# Only fall back to plain text for format-related errors, not server rejections
|
||||
if any(code in err_str for code in ("550", "553", "554", "auth", "rejected")):
|
||||
raise
|
||||
logger.warning(f"HTML email failed ({e}), falling back to plain text")
|
||||
await loop.run_in_executor(
|
||||
None, self._smtp_send, chat_id, raw_text, metadata or {},
|
||||
)
|
||||
|
||||
def _smtp_send(self, to: str, content: str, meta: dict) -> None:
|
||||
cfg = self.config
|
||||
from_addr = cfg.from_address or cfg.smtp_username
|
||||
logger.debug(f"SMTP plain send: from={from_addr} to={to}")
|
||||
msg = EmailMessage()
|
||||
orig_subj = meta.get("subject", "")
|
||||
msg["Subject"] = f"{cfg.subject_prefix}{orig_subj}" if orig_subj and not orig_subj.lower().startswith("re:") else (orig_subj or "EvoScientist Reply")
|
||||
msg["From"] = from_addr
|
||||
msg["To"] = to
|
||||
orig_id = meta.get("original_message_id", "")
|
||||
if orig_id:
|
||||
msg["In-Reply-To"] = orig_id
|
||||
msg["References"] = f"{meta.get('references', '')} {orig_id}".strip()
|
||||
msg.set_content(content)
|
||||
try:
|
||||
if cfg.smtp_use_tls:
|
||||
srv = smtplib.SMTP(cfg.smtp_host, cfg.smtp_port, timeout=30)
|
||||
srv.starttls()
|
||||
else:
|
||||
srv = smtplib.SMTP_SSL(cfg.smtp_host, cfg.smtp_port, context=ssl.create_default_context(), timeout=30)
|
||||
srv.login(cfg.smtp_username, cfg.smtp_password)
|
||||
srv.sendmail(from_addr, [to], msg.as_string())
|
||||
srv.quit()
|
||||
except Exception as e:
|
||||
logger.error(f"SMTP send failed: from={from_addr} to={to} error={e}")
|
||||
raise RuntimeError(f"SMTP: {e}")
|
||||
|
||||
def _smtp_send_html(self, to: str, html_content: str, plain_content: str, meta: dict) -> None:
|
||||
"""Send an email with both HTML and plain-text parts."""
|
||||
cfg = self.config
|
||||
from_addr = cfg.from_address or cfg.smtp_username
|
||||
logger.debug(f"SMTP HTML send: from={from_addr} to={to}")
|
||||
msg = MIMEMultipart("alternative")
|
||||
orig_subj = meta.get("subject", "")
|
||||
msg["Subject"] = f"{cfg.subject_prefix}{orig_subj}" if orig_subj and not orig_subj.lower().startswith("re:") else (orig_subj or "EvoScientist Reply")
|
||||
msg["From"] = from_addr
|
||||
msg["To"] = to
|
||||
orig_id = meta.get("original_message_id", "")
|
||||
if orig_id:
|
||||
msg["In-Reply-To"] = orig_id
|
||||
msg["References"] = f"{meta.get('references', '')} {orig_id}".strip()
|
||||
msg.attach(MIMEText(plain_content, "plain", "utf-8"))
|
||||
msg.attach(MIMEText(html_content, "html", "utf-8"))
|
||||
try:
|
||||
if cfg.smtp_use_tls:
|
||||
srv = smtplib.SMTP(cfg.smtp_host, cfg.smtp_port, timeout=30)
|
||||
srv.starttls()
|
||||
else:
|
||||
srv = smtplib.SMTP_SSL(cfg.smtp_host, cfg.smtp_port, context=ssl.create_default_context(), timeout=30)
|
||||
srv.login(cfg.smtp_username, cfg.smtp_password)
|
||||
srv.sendmail(from_addr, [to], msg.as_string())
|
||||
srv.quit()
|
||||
except Exception as e:
|
||||
logger.error(f"SMTP HTML send failed: from={from_addr} to={to} error={e}")
|
||||
raise RuntimeError(f"SMTP HTML: {e}")
|
||||
|
||||
# ── Media send (email attachment) ─────────────────────────────
|
||||
|
||||
async def _send_media_impl(
|
||||
self,
|
||||
recipient: str,
|
||||
file_path: str,
|
||||
caption: str = "",
|
||||
metadata: dict | None = None,
|
||||
) -> bool:
|
||||
"""Send a file as an email attachment via SMTP."""
|
||||
loop = asyncio.get_event_loop()
|
||||
await loop.run_in_executor(
|
||||
None, self._smtp_send_attachment, recipient, file_path, caption, metadata or {},
|
||||
)
|
||||
return True
|
||||
|
||||
def _smtp_send_attachment(self, to: str, file_path: str, caption: str, meta: dict) -> None:
|
||||
"""Send an email with a file attachment."""
|
||||
cfg = self.config
|
||||
from_addr = cfg.from_address or cfg.smtp_username
|
||||
logger.debug(f"SMTP attachment send: from={from_addr} to={to} file={file_path}")
|
||||
msg = MIMEMultipart()
|
||||
orig_subj = meta.get("subject", "")
|
||||
msg["Subject"] = f"{cfg.subject_prefix}{orig_subj}" if orig_subj and not orig_subj.lower().startswith("re:") else (orig_subj or "EvoScientist Reply")
|
||||
msg["From"] = from_addr
|
||||
msg["To"] = to
|
||||
orig_id = meta.get("original_message_id", "")
|
||||
if orig_id:
|
||||
msg["In-Reply-To"] = orig_id
|
||||
msg["References"] = f"{meta.get('references', '')} {orig_id}".strip()
|
||||
|
||||
# Text body
|
||||
if caption:
|
||||
msg.attach(MIMEText(caption, "plain", "utf-8"))
|
||||
|
||||
# Attachment
|
||||
path = Path(file_path)
|
||||
part = MIMEBase("application", "octet-stream")
|
||||
part.set_payload(path.read_bytes())
|
||||
encoders.encode_base64(part)
|
||||
part.add_header("Content-Disposition", f"attachment; filename={path.name}")
|
||||
msg.attach(part)
|
||||
|
||||
try:
|
||||
if cfg.smtp_use_tls:
|
||||
srv = smtplib.SMTP(cfg.smtp_host, cfg.smtp_port, timeout=30)
|
||||
srv.starttls()
|
||||
else:
|
||||
srv = smtplib.SMTP_SSL(cfg.smtp_host, cfg.smtp_port, context=ssl.create_default_context(), timeout=30)
|
||||
srv.login(cfg.smtp_username, cfg.smtp_password)
|
||||
srv.sendmail(from_addr, [to], msg.as_string())
|
||||
srv.quit()
|
||||
except Exception as e:
|
||||
logger.error(f"SMTP attachment send failed: from={from_addr} to={to} error={e}")
|
||||
raise RuntimeError(f"SMTP attachment: {e}")
|
||||
|
||||
async def _cleanup(self) -> None:
|
||||
await self._stop_polling()
|
||||
if self._imap:
|
||||
try:
|
||||
self._imap.close()
|
||||
self._imap.logout()
|
||||
except Exception:
|
||||
pass
|
||||
self._imap = None
|
||||
logger.info("Email channel stopped")
|
||||
@@ -0,0 +1,75 @@
|
||||
"""Email credential validation."""
|
||||
|
||||
import imaplib
|
||||
import logging
|
||||
import smtplib
|
||||
import ssl
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def validate_email_imap(
|
||||
host: str, port: int, username: str, password: str,
|
||||
use_ssl: bool = True,
|
||||
) -> tuple[bool, str]:
|
||||
"""Validate IMAP credentials.
|
||||
|
||||
Returns:
|
||||
Tuple of (is_valid, message).
|
||||
"""
|
||||
if not host or not username or not password:
|
||||
return False, "host, username, and password are required"
|
||||
|
||||
import asyncio
|
||||
loop = asyncio.get_event_loop()
|
||||
|
||||
def _check():
|
||||
try:
|
||||
if use_ssl:
|
||||
ctx = ssl.create_default_context()
|
||||
conn = imaplib.IMAP4_SSL(host, port, ssl_context=ctx)
|
||||
else:
|
||||
conn = imaplib.IMAP4(host, port)
|
||||
conn.login(username, password)
|
||||
conn.logout()
|
||||
return True, "IMAP credentials valid"
|
||||
except imaplib.IMAP4.error as e:
|
||||
return False, f"IMAP auth failed: {e}"
|
||||
except Exception as e:
|
||||
return False, f"IMAP error: {e}"
|
||||
|
||||
return await loop.run_in_executor(None, _check)
|
||||
|
||||
|
||||
async def validate_email_smtp(
|
||||
host: str, port: int, username: str, password: str,
|
||||
use_tls: bool = True,
|
||||
) -> tuple[bool, str]:
|
||||
"""Validate SMTP credentials.
|
||||
|
||||
Returns:
|
||||
Tuple of (is_valid, message).
|
||||
"""
|
||||
if not host or not username or not password:
|
||||
return False, "host, username, and password are required"
|
||||
|
||||
import asyncio
|
||||
loop = asyncio.get_event_loop()
|
||||
|
||||
def _check():
|
||||
try:
|
||||
if use_tls:
|
||||
server = smtplib.SMTP(host, port, timeout=10)
|
||||
server.starttls()
|
||||
else:
|
||||
ctx = ssl.create_default_context()
|
||||
server = smtplib.SMTP_SSL(host, port, context=ctx, timeout=10)
|
||||
server.login(username, password)
|
||||
server.quit()
|
||||
return True, "SMTP credentials valid"
|
||||
except smtplib.SMTPAuthenticationError as e:
|
||||
return False, f"SMTP auth failed: {e}"
|
||||
except Exception as e:
|
||||
return False, f"SMTP error: {e}"
|
||||
|
||||
return await loop.run_in_executor(None, _check)
|
||||
@@ -0,0 +1,124 @@
|
||||
"""Email channel server.
|
||||
|
||||
Standalone script to run the Email channel with CLI options.
|
||||
|
||||
Usage:
|
||||
python -m EvoScientist.channels.email.serve --imap-host HOST --imap-username USER --imap-password PASS --smtp-host HOST --smtp-username USER --smtp-password PASS --from-address ADDR [OPTIONS]
|
||||
|
||||
Examples:
|
||||
# Basic usage
|
||||
python -m EvoScientist.channels.email.serve --imap-host imap.gmail.com --imap-username bot@gmail.com --imap-password PASS --smtp-host smtp.gmail.com --smtp-username bot@gmail.com --smtp-password PASS --from-address bot@gmail.com
|
||||
|
||||
# With allowed senders and custom poll interval
|
||||
python -m EvoScientist.channels.email.serve --imap-host imap.gmail.com --imap-username bot@gmail.com --imap-password PASS --smtp-host smtp.gmail.com --smtp-username bot@gmail.com --smtp-password PASS --from-address bot@gmail.com --allow user@example.com --poll-interval 60
|
||||
|
||||
# With agent and thinking
|
||||
python -m EvoScientist.channels.email.serve --imap-host imap.gmail.com --imap-username bot@gmail.com --imap-password PASS --smtp-host smtp.gmail.com --smtp-username bot@gmail.com --smtp-password PASS --from-address bot@gmail.com --agent --thinking
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
|
||||
from .channel import EmailChannel, EmailConfig
|
||||
from ..bus import MessageBus
|
||||
from ..standalone import run_standalone
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.DEBUG,
|
||||
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
|
||||
datefmt="%H:%M:%S",
|
||||
)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def parse_args():
|
||||
"""Parse command line arguments."""
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Email channel server",
|
||||
formatter_class=argparse.RawDescriptionHelpFormatter,
|
||||
)
|
||||
parser.add_argument(
|
||||
"--imap-host",
|
||||
required=True,
|
||||
help="IMAP server hostname",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--imap-username",
|
||||
required=True,
|
||||
help="IMAP username",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--imap-password",
|
||||
required=True,
|
||||
help="IMAP password",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--smtp-host",
|
||||
required=True,
|
||||
help="SMTP server hostname",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--smtp-username",
|
||||
required=True,
|
||||
help="SMTP username",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--smtp-password",
|
||||
required=True,
|
||||
help="SMTP password",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--from-address",
|
||||
required=True,
|
||||
help="From email address for outgoing messages",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--allow",
|
||||
action="append",
|
||||
dest="allowed_senders",
|
||||
help="Allowed sender (email address). Can be used multiple times.",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--poll-interval",
|
||||
type=int,
|
||||
default=30,
|
||||
help="IMAP poll interval in seconds (default: 30)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--agent",
|
||||
action="store_true",
|
||||
help="Use EvoScientist agent as handler (default: echo)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--thinking",
|
||||
action="store_true",
|
||||
help="Send thinking content as intermediate messages (requires --agent)",
|
||||
)
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def main():
|
||||
"""Entry point."""
|
||||
args = parse_args()
|
||||
|
||||
config = EmailConfig(
|
||||
imap_host=args.imap_host,
|
||||
imap_username=args.imap_username,
|
||||
imap_password=args.imap_password,
|
||||
smtp_host=args.smtp_host,
|
||||
smtp_username=args.smtp_username,
|
||||
smtp_password=args.smtp_password,
|
||||
from_address=args.from_address,
|
||||
allowed_senders=set(args.allowed_senders) if args.allowed_senders else None,
|
||||
poll_interval=args.poll_interval,
|
||||
)
|
||||
|
||||
send_thinking = args.thinking and args.agent
|
||||
bus = MessageBus()
|
||||
channel = EmailChannel(config)
|
||||
|
||||
run_standalone(channel, bus, use_agent=args.agent, send_thinking=send_thinking)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,26 @@
|
||||
"""QQ channel for EvoScientist.
|
||||
|
||||
Uses the official qq-botpy SDK for WebSocket connection.
|
||||
|
||||
Usage in config:
|
||||
channel_enabled = "qq"
|
||||
qq_app_id = "your_app_id"
|
||||
qq_app_secret = "your_app_secret"
|
||||
"""
|
||||
|
||||
from .channel import QQChannel, QQConfig
|
||||
from ..channel_manager import register_channel, _parse_csv
|
||||
|
||||
__all__ = ["QQChannel", "QQConfig"]
|
||||
|
||||
|
||||
def create_from_config(config) -> QQChannel:
|
||||
allowed = _parse_csv(config.qq_allowed_senders)
|
||||
return QQChannel(QQConfig(
|
||||
app_id=config.qq_app_id,
|
||||
app_secret=config.qq_app_secret,
|
||||
allowed_senders=allowed,
|
||||
))
|
||||
|
||||
|
||||
register_channel("qq", create_from_config)
|
||||
@@ -0,0 +1,258 @@
|
||||
"""QQ channel implementation using botpy SDK."""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from collections import deque
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
|
||||
from ..base import Channel, RawIncoming, ChannelError
|
||||
from ..capabilities import QQ as QQ_CAPS
|
||||
from ..config import BaseChannelConfig
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
try:
|
||||
import botpy
|
||||
from botpy.message import C2CMessage, GroupMessage
|
||||
|
||||
QQ_AVAILABLE = True
|
||||
except ImportError:
|
||||
QQ_AVAILABLE = False
|
||||
botpy = None
|
||||
C2CMessage = None
|
||||
GroupMessage = None
|
||||
|
||||
|
||||
@dataclass
|
||||
class QQConfig(BaseChannelConfig):
|
||||
app_id: str = ""
|
||||
app_secret: str = ""
|
||||
text_chunk_limit: int = 4096
|
||||
|
||||
|
||||
def _make_bot_class(channel: "QQChannel") -> "type[botpy.Client]":
|
||||
"""Create a botpy Client subclass bound to the given channel."""
|
||||
intents = botpy.Intents(public_messages=True, direct_message=True)
|
||||
|
||||
class _Bot(botpy.Client):
|
||||
def __init__(self):
|
||||
super().__init__(intents=intents)
|
||||
|
||||
async def on_ready(self):
|
||||
logger.info(f"QQ bot ready: {self.robot.name}")
|
||||
|
||||
async def on_c2c_message_create(self, message: "C2CMessage"):
|
||||
await channel._on_msg(message, "c2c")
|
||||
|
||||
async def on_group_at_message_create(self, message: "GroupMessage"):
|
||||
await channel._on_msg(message, "group")
|
||||
|
||||
return _Bot
|
||||
|
||||
|
||||
class QQChannel(Channel):
|
||||
"""QQ channel using botpy SDK."""
|
||||
|
||||
name = "qq"
|
||||
|
||||
capabilities = QQ_CAPS
|
||||
_ready_attrs = ("_client", "_running")
|
||||
_non_retryable_patterns = ()
|
||||
_mention_pattern = r"@\S+\s*"
|
||||
_mention_strip_count = 1
|
||||
|
||||
def __init__(self, config: QQConfig):
|
||||
super().__init__(config)
|
||||
self._client: "botpy.Client | None" = None
|
||||
self._bot_task: asyncio.Task | None = None
|
||||
self._processed_ids: deque = deque(maxlen=1000)
|
||||
self._msg_seq: dict[str, int] = {} # msg_id -> next seq number
|
||||
self._msg_seq_order: deque = deque(maxlen=500)
|
||||
|
||||
# ── Lifecycle ─────────────────────────────────────────────────
|
||||
|
||||
async def start(self) -> None:
|
||||
if not QQ_AVAILABLE:
|
||||
raise ChannelError("QQ SDK not installed. Run: pip install qq-botpy")
|
||||
if not self.config.app_id or not self.config.app_secret:
|
||||
raise ChannelError("QQ app_id and app_secret are required")
|
||||
self._running = True
|
||||
BotClass = _make_bot_class(self)
|
||||
self._client = BotClass()
|
||||
self._bot_task = asyncio.create_task(self._run_bot())
|
||||
logger.info("QQ channel starting...")
|
||||
|
||||
async def _run_bot(self) -> None:
|
||||
try:
|
||||
await self._client.start(appid=self.config.app_id, secret=self.config.app_secret)
|
||||
except Exception as e:
|
||||
logger.error(f"QQ auth failed: {e}")
|
||||
self._running = False
|
||||
|
||||
# ── Incoming ──────────────────────────────────────────────────
|
||||
|
||||
async def _on_msg(self, message, msg_type: str) -> None:
|
||||
try:
|
||||
if message.id in self._processed_ids:
|
||||
return
|
||||
self._processed_ids.append(message.id)
|
||||
|
||||
author = message.author
|
||||
content = (message.content or "").strip()
|
||||
|
||||
if msg_type == "c2c":
|
||||
sender_id = str(getattr(author, "user_openid", ""))
|
||||
chat_id = sender_id
|
||||
else:
|
||||
sender_id = str(getattr(author, "member_openid", ""))
|
||||
chat_id = str(getattr(message, "group_openid", ""))
|
||||
|
||||
# Handle attachments (images, files, audio, video)
|
||||
annotations: list[str] = []
|
||||
media_paths: list[str] = []
|
||||
attachments = getattr(message, "attachments", None) or []
|
||||
for att in attachments:
|
||||
url = getattr(att, "url", "") or ""
|
||||
filename = getattr(att, "filename", "attachment") or "attachment"
|
||||
content_type = getattr(att, "content_type", "") or ""
|
||||
if url:
|
||||
local, ann = await self._download_attachment(
|
||||
url, f"qq_{filename}",
|
||||
)
|
||||
if local:
|
||||
media_paths.append(local)
|
||||
if ann:
|
||||
annotations.append(ann)
|
||||
else:
|
||||
annotations.append(f"[{content_type or 'attachment'}: {filename}]")
|
||||
|
||||
if not content and not media_paths and not annotations:
|
||||
return
|
||||
|
||||
await self._enqueue_raw(RawIncoming(
|
||||
sender_id=sender_id,
|
||||
chat_id=chat_id,
|
||||
text=content,
|
||||
media_files=media_paths,
|
||||
content_annotations=annotations,
|
||||
timestamp=datetime.now(),
|
||||
message_id=message.id,
|
||||
is_group=(msg_type == "group"),
|
||||
was_mentioned=True,
|
||||
metadata={
|
||||
"chat_id": chat_id,
|
||||
"msg_type": msg_type,
|
||||
"event_id": message.id,
|
||||
"backend": "qq",
|
||||
},
|
||||
))
|
||||
except Exception as e:
|
||||
logger.error(f"Error handling QQ message: {e}")
|
||||
|
||||
# ── Send ──────────────────────────────────────────────────────
|
||||
|
||||
def _next_msg_seq(self, msg_id: str) -> int:
|
||||
"""Return the next msg_seq for *msg_id* and increment the counter."""
|
||||
seq = self._msg_seq.get(msg_id, 1)
|
||||
self._msg_seq[msg_id] = seq + 1
|
||||
if msg_id not in set(self._msg_seq_order):
|
||||
self._msg_seq_order.append(msg_id)
|
||||
if len(self._msg_seq_order) > 500:
|
||||
oldest = self._msg_seq_order.popleft()
|
||||
self._msg_seq.pop(oldest, None)
|
||||
return seq
|
||||
|
||||
async def _send_chunk(self, chat_id, formatted_text, raw_text, reply_to, metadata):
|
||||
if not self._client:
|
||||
raise ChannelError("QQ client not initialized")
|
||||
msg_type = (metadata or {}).get("msg_type", "c2c")
|
||||
msg_id = (metadata or {}).get("event_id", "")
|
||||
seq = self._next_msg_seq(msg_id)
|
||||
if msg_type == "group":
|
||||
await self._client.api.post_group_message(
|
||||
group_openid=chat_id, msg_type=0,
|
||||
content=raw_text, msg_id=msg_id, msg_seq=seq,
|
||||
)
|
||||
else:
|
||||
await self._client.api.post_c2c_message(
|
||||
openid=chat_id, msg_type=0,
|
||||
content=raw_text, msg_id=msg_id, msg_seq=seq,
|
||||
)
|
||||
|
||||
# _send_typing_action: inherited no-op (QQ Bot API has no typing indicator)
|
||||
|
||||
# ── Media send ────────────────────────────────────────────────
|
||||
|
||||
# qq-botpy file_type constants: 1=image, 2=video, 3=audio
|
||||
_FILE_TYPE_MAP = {
|
||||
".jpg": 1, ".jpeg": 1, ".png": 1, ".gif": 1, ".webp": 1, ".bmp": 1,
|
||||
".mp4": 2, ".mov": 2, ".avi": 2,
|
||||
".mp3": 3, ".ogg": 3, ".m4a": 3, ".wav": 3, ".silk": 3,
|
||||
}
|
||||
|
||||
async def _send_media_impl(
|
||||
self,
|
||||
recipient: str,
|
||||
file_path: str,
|
||||
caption: str = "",
|
||||
metadata: dict | None = None,
|
||||
) -> bool:
|
||||
"""Send a media file through QQ Bot API.
|
||||
|
||||
Uses post_group_file / post_c2c_file with a URL. Local files
|
||||
without a public URL are not supported — falls back to a text hint.
|
||||
"""
|
||||
if not self._client:
|
||||
raise ChannelError("QQ client not initialized")
|
||||
|
||||
from pathlib import Path
|
||||
chat_id = self._resolve_media_chat_id(recipient, metadata)
|
||||
msg_type = (metadata or {}).get("msg_type", "c2c")
|
||||
ext = Path(file_path).suffix.lower()
|
||||
file_type = self._FILE_TYPE_MAP.get(ext, 1) # default to image
|
||||
|
||||
# qq-botpy file API requires a URL, not a local path
|
||||
is_url = file_path.startswith("http://") or file_path.startswith("https://")
|
||||
if not is_url:
|
||||
# Fallback: send text hint for local files
|
||||
name = Path(file_path).name
|
||||
hint = f"[文件] {name}" + (f"\n{caption}" if caption else "")
|
||||
await self._send_chunk(chat_id, hint, hint, None, metadata or {})
|
||||
return True
|
||||
|
||||
try:
|
||||
if msg_type == "group":
|
||||
await self._client.api.post_group_file(
|
||||
group_openid=chat_id,
|
||||
file_type=file_type,
|
||||
url=file_path,
|
||||
srv_send_msg=True,
|
||||
)
|
||||
else:
|
||||
await self._client.api.post_c2c_file(
|
||||
openid=chat_id,
|
||||
file_type=file_type,
|
||||
url=file_path,
|
||||
srv_send_msg=True,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"QQ media send failed: {e}")
|
||||
return False
|
||||
|
||||
if caption:
|
||||
await self._send_chunk(chat_id, caption, caption, None, metadata or {})
|
||||
return True
|
||||
|
||||
# ── Cleanup ───────────────────────────────────────────────────
|
||||
|
||||
async def _cleanup(self) -> None:
|
||||
self._running = False
|
||||
if self._bot_task:
|
||||
self._bot_task.cancel()
|
||||
try:
|
||||
await self._bot_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._client = None
|
||||
logger.info("QQ channel stopped")
|
||||
@@ -0,0 +1,37 @@
|
||||
"""QQ Bot credential validation."""
|
||||
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
QQ_TOKEN_URL = "https://bots.qq.com/app/getAppAccessToken"
|
||||
|
||||
|
||||
async def validate_qq(
|
||||
app_id: str,
|
||||
app_secret: str,
|
||||
) -> tuple[bool, str]:
|
||||
"""Validate QQ Bot credentials by fetching an access token.
|
||||
|
||||
Returns:
|
||||
Tuple of (is_valid, message).
|
||||
"""
|
||||
if not app_id or not app_secret:
|
||||
return False, "app_id and app_secret are required"
|
||||
|
||||
try:
|
||||
import httpx
|
||||
except ImportError:
|
||||
return False, "httpx not installed"
|
||||
|
||||
body = {"appId": app_id, "clientSecret": app_secret}
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient() as client:
|
||||
resp = await client.post(QQ_TOKEN_URL, json=body, timeout=10)
|
||||
data = resp.json()
|
||||
if data.get("access_token"):
|
||||
return True, "QQ Bot credentials valid"
|
||||
return False, f"Error: {data.get('message', data)}"
|
||||
except Exception as e:
|
||||
return False, f"Error: {e}"
|
||||
@@ -0,0 +1,87 @@
|
||||
"""QQ channel server.
|
||||
|
||||
Standalone script to run the QQ channel with CLI options.
|
||||
|
||||
Usage:
|
||||
python -m EvoScientist.channels.qq.serve --app-id ID --app-secret SECRET [OPTIONS]
|
||||
|
||||
Examples:
|
||||
# Basic usage
|
||||
python -m EvoScientist.channels.qq.serve --app-id ID --app-secret SECRET
|
||||
|
||||
# Sandbox mode with allowed senders
|
||||
python -m EvoScientist.channels.qq.serve --app-id ID --app-secret SECRET --allow user123
|
||||
|
||||
# With agent and thinking
|
||||
python -m EvoScientist.channels.qq.serve --app-id ID --app-secret SECRET --agent --thinking
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
|
||||
from .channel import QQChannel, QQConfig
|
||||
from ..bus import MessageBus
|
||||
from ..standalone import run_standalone
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.DEBUG,
|
||||
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
|
||||
datefmt="%H:%M:%S",
|
||||
)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def parse_args():
|
||||
"""Parse command line arguments."""
|
||||
parser = argparse.ArgumentParser(
|
||||
description="QQ channel server",
|
||||
formatter_class=argparse.RawDescriptionHelpFormatter,
|
||||
)
|
||||
parser.add_argument(
|
||||
"--app-id",
|
||||
required=True,
|
||||
help="QQ bot app ID",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--app-secret",
|
||||
required=True,
|
||||
help="QQ bot app secret",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--allow",
|
||||
action="append",
|
||||
dest="allowed_senders",
|
||||
help="Allowed sender (QQ user ID). Can be used multiple times.",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--agent",
|
||||
action="store_true",
|
||||
help="Use EvoScientist agent as handler (default: echo)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--thinking",
|
||||
action="store_true",
|
||||
help="Send thinking content as intermediate messages (requires --agent)",
|
||||
)
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def main():
|
||||
"""Entry point."""
|
||||
args = parse_args()
|
||||
|
||||
config = QQConfig(
|
||||
app_id=args.app_id,
|
||||
app_secret=args.app_secret,
|
||||
allowed_senders=set(args.allowed_senders) if args.allowed_senders else None,
|
||||
)
|
||||
|
||||
send_thinking = args.thinking and args.agent
|
||||
bus = MessageBus()
|
||||
channel = QQChannel(config)
|
||||
|
||||
run_standalone(channel, bus, use_agent=args.agent, send_thinking=send_thinking)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,27 @@
|
||||
"""Signal channel for EvoScientist.
|
||||
|
||||
Uses signal-cli in JSON RPC mode — no public IP needed.
|
||||
|
||||
Usage in config:
|
||||
channel_enabled = "signal"
|
||||
signal_phone_number = "+1234567890"
|
||||
"""
|
||||
|
||||
from .channel import SignalChannel, SignalConfig
|
||||
from ..channel_manager import register_channel, _parse_csv
|
||||
|
||||
__all__ = ["SignalChannel", "SignalConfig"]
|
||||
|
||||
|
||||
def create_from_config(config) -> SignalChannel:
|
||||
allowed = _parse_csv(config.signal_allowed_senders)
|
||||
return SignalChannel(SignalConfig(
|
||||
phone_number=config.signal_phone_number,
|
||||
cli_path=config.signal_cli_path,
|
||||
config_dir=config.signal_config_dir or None,
|
||||
rpc_port=config.signal_rpc_port,
|
||||
allowed_senders=allowed,
|
||||
))
|
||||
|
||||
|
||||
register_channel("signal", create_from_config)
|
||||
@@ -0,0 +1,425 @@
|
||||
"""Signal channel implementation using signal-cli JSON RPC."""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
import subprocess
|
||||
from collections import deque
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from ..base import Channel, RawIncoming, ChannelError
|
||||
from ..capabilities import SIGNAL as SIGNAL_CAPS
|
||||
from ..config import BaseChannelConfig
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@dataclass
|
||||
class SignalConfig(BaseChannelConfig):
|
||||
phone_number: str = ""
|
||||
cli_path: str = "signal-cli"
|
||||
config_dir: str | None = None
|
||||
rpc_port: int = 7583
|
||||
text_chunk_limit: int = 4096
|
||||
|
||||
|
||||
class SignalChannel(Channel):
|
||||
"""Signal channel using signal-cli JSON RPC."""
|
||||
|
||||
name = "signal"
|
||||
|
||||
capabilities = SIGNAL_CAPS
|
||||
_non_retryable_patterns = ("unregistered", "auth")
|
||||
|
||||
def __init__(self, config: SignalConfig):
|
||||
super().__init__(config)
|
||||
self._reader: asyncio.StreamReader | None = None
|
||||
self._writer: asyncio.StreamWriter | None = None
|
||||
self._rpc_id = 0
|
||||
self._daemon_proc = None
|
||||
# Cache message_id → sender for reaction targetAuthor (bounded)
|
||||
self._msg_senders: dict[str, str] = {}
|
||||
self._msg_senders_order: deque = deque(maxlen=200)
|
||||
|
||||
async def start(self) -> None:
|
||||
if not self.config.phone_number:
|
||||
raise ChannelError("Signal phone_number is required")
|
||||
|
||||
# Try to start signal-cli daemon if not already running
|
||||
await self._ensure_daemon()
|
||||
|
||||
# Connect to JSON RPC socket
|
||||
await self._connect()
|
||||
|
||||
self._running = True
|
||||
logger.info(f"Signal channel started (phone: {self.config.phone_number})")
|
||||
|
||||
# Listen for incoming messages in background task
|
||||
# (start() must return so that run() can iterate receive())
|
||||
self._listen_task = asyncio.create_task(self._listen_loop())
|
||||
|
||||
async def _ensure_daemon(self) -> None:
|
||||
"""Start signal-cli daemon if not already running."""
|
||||
try:
|
||||
reader, writer = await asyncio.wait_for(
|
||||
asyncio.open_connection("localhost", self.config.rpc_port),
|
||||
timeout=2,
|
||||
)
|
||||
writer.close()
|
||||
await writer.wait_closed()
|
||||
logger.info("signal-cli daemon already running")
|
||||
return
|
||||
except (ConnectionRefusedError, asyncio.TimeoutError, OSError):
|
||||
pass
|
||||
|
||||
# Start daemon
|
||||
cmd = [self.config.cli_path, "-u", self.config.phone_number]
|
||||
if self.config.config_dir:
|
||||
cmd.extend(["--config", self.config.config_dir])
|
||||
cmd.extend(["daemon", "--tcp",
|
||||
f"localhost:{self.config.rpc_port}", "--no-receive-stdout"])
|
||||
|
||||
logger.info(f"Starting signal-cli daemon: {' '.join(cmd)}")
|
||||
try:
|
||||
self._daemon_proc = subprocess.Popen(
|
||||
cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
|
||||
)
|
||||
except FileNotFoundError:
|
||||
raise ChannelError(
|
||||
f"signal-cli not found at '{self.config.cli_path}'. "
|
||||
"Install: https://github.com/AsamK/signal-cli"
|
||||
)
|
||||
|
||||
# Wait for daemon to be ready
|
||||
for _ in range(30):
|
||||
await asyncio.sleep(1)
|
||||
try:
|
||||
reader, writer = await asyncio.open_connection(
|
||||
"localhost", self.config.rpc_port,
|
||||
)
|
||||
writer.close()
|
||||
await writer.wait_closed()
|
||||
logger.info("signal-cli daemon started")
|
||||
return
|
||||
except (ConnectionRefusedError, OSError):
|
||||
continue
|
||||
|
||||
raise ChannelError("signal-cli daemon failed to start within 30s")
|
||||
|
||||
async def _connect(self) -> None:
|
||||
"""Connect to signal-cli JSON RPC socket."""
|
||||
try:
|
||||
self._reader, self._writer = await asyncio.open_connection(
|
||||
"localhost", self.config.rpc_port,
|
||||
)
|
||||
except Exception as e:
|
||||
raise ChannelError(f"Cannot connect to signal-cli: {e}")
|
||||
|
||||
async def _listen_loop(self) -> None:
|
||||
"""Listen for incoming JSON RPC notifications."""
|
||||
while self._running and self._reader:
|
||||
try:
|
||||
line = await self._reader.readline()
|
||||
if not line:
|
||||
break
|
||||
data = json.loads(line.decode())
|
||||
await self._handle_rpc(data)
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
except Exception as e:
|
||||
logger.error(f"Signal listen error: {e}")
|
||||
# Reconnect
|
||||
if self._running:
|
||||
await asyncio.sleep(2)
|
||||
try:
|
||||
await self._connect()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
async def _handle_rpc(self, data: dict) -> None:
|
||||
"""Handle a JSON RPC message from signal-cli."""
|
||||
method = data.get("method", "")
|
||||
|
||||
if method != "receive":
|
||||
return
|
||||
|
||||
params = data.get("params", {})
|
||||
envelope = params.get("envelope", {})
|
||||
source = envelope.get("source") or envelope.get("sourceUuid") or ""
|
||||
source_number = envelope.get("sourceNumber") or source
|
||||
source_name = envelope.get("sourceName") or ""
|
||||
timestamp = envelope.get("timestamp", 0)
|
||||
|
||||
# Ignore messages from self
|
||||
if source_number == self.config.phone_number or source == self.config.phone_number:
|
||||
logger.debug("Ignoring message from self")
|
||||
return
|
||||
|
||||
# Data message (text)
|
||||
data_msg = envelope.get("dataMessage", {})
|
||||
if data_msg:
|
||||
text = data_msg.get("message", "")
|
||||
group_info = data_msg.get("groupInfo", {})
|
||||
is_group = bool(group_info)
|
||||
chat_id = group_info.get("groupId", source_number) if is_group else source_number
|
||||
msg_ts = data_msg.get("timestamp", timestamp)
|
||||
|
||||
media_paths: list[str] = []
|
||||
annotations: list[str] = []
|
||||
_VOICE_TYPES = {"audio/aac", "audio/ogg", "audio/mp4", "audio/mpeg", "audio/opus"}
|
||||
attachments = data_msg.get("attachments", [])
|
||||
for att in attachments:
|
||||
att_size = att.get("size", 0)
|
||||
att_name = att.get("filename", "attachment")
|
||||
att_file = att.get("file") # signal-cli provides local path
|
||||
content_type = att.get("contentType", "")
|
||||
is_voice = content_type in _VOICE_TYPES or att.get("voiceNote", False)
|
||||
media_label = "voice" if is_voice else "attachment"
|
||||
if att_file:
|
||||
from pathlib import Path as _Path
|
||||
att_path = _Path(att_file)
|
||||
if att_path.exists():
|
||||
from ..base import MAX_ATTACHMENT_BYTES
|
||||
if att_path.stat().st_size > MAX_ATTACHMENT_BYTES:
|
||||
annotations.append(f"[{media_label}: {att_name} - too large ({att_path.stat().st_size} bytes)]")
|
||||
else:
|
||||
local = self._media_path(f"signal_{att_name}")
|
||||
import shutil
|
||||
shutil.copy2(str(att_path), str(local))
|
||||
media_paths.append(str(local))
|
||||
annotations.append(f"[{media_label}: {local}]")
|
||||
else:
|
||||
annotations.append(f"[{media_label}: {att_name} - file not found]")
|
||||
elif att_size:
|
||||
too_large = self._check_attachment_size(att_size, att_name)
|
||||
if too_large:
|
||||
annotations.append(too_large)
|
||||
else:
|
||||
annotations.append(f"[{media_label}: {att_name}]")
|
||||
|
||||
if not text and not media_paths and not annotations:
|
||||
if not attachments:
|
||||
return
|
||||
# Had attachments but none downloaded successfully
|
||||
if not annotations:
|
||||
text = "[attachment]"
|
||||
|
||||
try:
|
||||
ts = datetime.fromtimestamp(msg_ts / 1000) if msg_ts else datetime.now()
|
||||
except (ValueError, TypeError, OSError):
|
||||
ts = datetime.now()
|
||||
|
||||
was_mentioned = not is_group # DMs always pass
|
||||
if is_group:
|
||||
mentions = data_msg.get("mentions", [])
|
||||
for m in mentions:
|
||||
if m.get("uuid") == self.config.phone_number or m.get("number") == self.config.phone_number:
|
||||
was_mentioned = True
|
||||
break
|
||||
|
||||
# Cache message_id → sender for reaction targetAuthor
|
||||
self._cache_msg_sender(str(msg_ts), source_number)
|
||||
|
||||
logger.info("Signal message from %s: %s", source_number, text[:50] if text else "[media]")
|
||||
await self._enqueue_raw(RawIncoming(
|
||||
sender_id=source_number,
|
||||
chat_id=chat_id,
|
||||
text=text,
|
||||
content_annotations=annotations,
|
||||
media_files=media_paths,
|
||||
timestamp=ts,
|
||||
message_id=str(msg_ts),
|
||||
is_group=is_group,
|
||||
was_mentioned=was_mentioned,
|
||||
metadata={
|
||||
"chat_id": chat_id,
|
||||
"source_name": source_name,
|
||||
"sender_id": source_number,
|
||||
"backend": "signal",
|
||||
},
|
||||
))
|
||||
|
||||
# ── Typing indicator ────────────────────────────────────────────
|
||||
|
||||
async def _send_typing_action(self, chat_id: str) -> None:
|
||||
"""Send typing indicator via signal-cli JSON RPC."""
|
||||
params: dict[str, Any] = {
|
||||
"account": self.config.phone_number,
|
||||
}
|
||||
if self._is_group_id(chat_id):
|
||||
params["groupId"] = chat_id
|
||||
else:
|
||||
params["recipient"] = [chat_id]
|
||||
try:
|
||||
await self._rpc_call("sendTyping", params)
|
||||
except Exception:
|
||||
pass # typing indicator is best-effort
|
||||
|
||||
# ── ACK reaction ─────────────────────────────────────────────
|
||||
|
||||
def _cache_msg_sender(self, message_id: str, sender: str) -> None:
|
||||
"""Store message_id → sender mapping for reaction targetAuthor."""
|
||||
if len(self._msg_senders) >= 200:
|
||||
oldest = self._msg_senders_order.popleft()
|
||||
self._msg_senders.pop(oldest, None)
|
||||
self._msg_senders[message_id] = sender
|
||||
self._msg_senders_order.append(message_id)
|
||||
|
||||
async def _send_ack_reaction(self, chat_id: str, message_id: str, emoji: str = "👀") -> None:
|
||||
"""Send an acknowledgment reaction via signal-cli sendReaction."""
|
||||
target_author = self._msg_senders.get(message_id, "")
|
||||
if not target_author:
|
||||
return # cannot send reaction without knowing the original sender
|
||||
try:
|
||||
params: dict[str, Any] = {
|
||||
"account": self.config.phone_number,
|
||||
"emoji": emoji,
|
||||
"targetAuthor": target_author,
|
||||
"targetTimestamp": int(message_id),
|
||||
}
|
||||
if self._is_group_id(chat_id):
|
||||
params["groupId"] = chat_id
|
||||
else:
|
||||
params["recipient"] = [chat_id]
|
||||
await self._rpc_call("sendReaction", params)
|
||||
except Exception as e:
|
||||
logger.debug(f"Signal ack reaction failed: {e}")
|
||||
|
||||
async def _remove_ack_reaction(self, chat_id: str, message_id: str, emoji: str = "👀") -> None:
|
||||
"""Remove ACK reaction via signal-cli sendReaction --remove."""
|
||||
target_author = self._msg_senders.get(message_id, "")
|
||||
if not target_author:
|
||||
return
|
||||
try:
|
||||
params: dict[str, Any] = {
|
||||
"account": self.config.phone_number,
|
||||
"emoji": emoji,
|
||||
"targetAuthor": target_author,
|
||||
"targetTimestamp": int(message_id),
|
||||
"remove": True,
|
||||
}
|
||||
if self._is_group_id(chat_id):
|
||||
params["groupId"] = chat_id
|
||||
else:
|
||||
params["recipient"] = [chat_id]
|
||||
await self._rpc_call("sendReaction", params)
|
||||
except Exception as e:
|
||||
logger.debug(f"Signal remove ACK reaction failed: {e}")
|
||||
|
||||
# ── Send ──────────────────────────────────────────────────────
|
||||
|
||||
@staticmethod
|
||||
def _is_group_id(chat_id: str) -> bool:
|
||||
"""Return True if *chat_id* looks like a Signal group ID.
|
||||
|
||||
Group IDs are base64-encoded strings (e.g. ``"aB3d...=="``).
|
||||
Individual recipients are either phone numbers (``"+1234..."``)
|
||||
or UUIDs (``"817ab5e9-..."``) — neither of which is a group.
|
||||
"""
|
||||
return not chat_id.startswith("+") and "-" not in chat_id
|
||||
|
||||
def _is_ready(self) -> bool:
|
||||
return self._writer is not None and not self._writer.is_closing()
|
||||
|
||||
async def _rpc_call(self, method: str, params: dict) -> dict | None:
|
||||
"""Send a JSON RPC call to signal-cli."""
|
||||
if not self._writer:
|
||||
return None
|
||||
|
||||
self._rpc_id += 1
|
||||
request = {
|
||||
"jsonrpc": "2.0",
|
||||
"id": self._rpc_id,
|
||||
"method": method,
|
||||
"params": params,
|
||||
}
|
||||
line = json.dumps(request) + "\n"
|
||||
self._writer.write(line.encode())
|
||||
await self._writer.drain()
|
||||
return None # We don't wait for response in this simple impl
|
||||
|
||||
async def _send_chunk(
|
||||
self, chat_id, formatted_text, raw_text, reply_to, metadata,
|
||||
):
|
||||
# Determine if group or individual
|
||||
params: dict[str, Any] = {
|
||||
"message": raw_text,
|
||||
"account": self.config.phone_number,
|
||||
}
|
||||
|
||||
if self._is_group_id(chat_id):
|
||||
params["groupId"] = chat_id
|
||||
else:
|
||||
params["recipient"] = [chat_id]
|
||||
|
||||
await self._rpc_call("send", params)
|
||||
|
||||
# ── Mention stripping ────────────────────────────────────────────
|
||||
|
||||
def _strip_mention(self, text: str) -> str:
|
||||
"""Strip bot mention from Signal messages.
|
||||
|
||||
Signal mentions are embedded as special objects that reference
|
||||
the phone number. The text contains a placeholder character (U+FFFC)
|
||||
at the mention position.
|
||||
"""
|
||||
phone = self.config.phone_number
|
||||
if phone:
|
||||
# Remove phone number if directly mentioned as text
|
||||
text = re.sub(rf"@?{re.escape(phone)}\s*", "", text).strip()
|
||||
# Remove Unicode Object Replacement Character used as mention placeholder
|
||||
text = text.replace("\uFFFC", "").strip()
|
||||
return text
|
||||
|
||||
# ── Media send ────────────────────────────────────────────────
|
||||
|
||||
async def _send_media_impl(
|
||||
self,
|
||||
recipient: str,
|
||||
file_path: str,
|
||||
caption: str = "",
|
||||
metadata: dict | None = None,
|
||||
) -> bool:
|
||||
"""Send a media file via signal-cli JSON RPC.
|
||||
|
||||
Uses the "send" RPC method with the attachments parameter.
|
||||
"""
|
||||
chat_id = self._resolve_media_chat_id(recipient, metadata)
|
||||
params: dict[str, Any] = {
|
||||
"account": self.config.phone_number,
|
||||
"attachments": [file_path],
|
||||
}
|
||||
if caption:
|
||||
params["message"] = caption
|
||||
|
||||
if self._is_group_id(chat_id):
|
||||
params["groupId"] = chat_id
|
||||
else:
|
||||
params["recipient"] = [chat_id]
|
||||
|
||||
await self._rpc_call("send", params)
|
||||
return True
|
||||
|
||||
# ── Cleanup ───────────────────────────────────────────────────
|
||||
|
||||
async def _cleanup(self) -> None:
|
||||
if hasattr(self, "_listen_task") and self._listen_task:
|
||||
self._listen_task.cancel()
|
||||
self._listen_task = None
|
||||
if self._writer:
|
||||
self._writer.close()
|
||||
try:
|
||||
await self._writer.wait_closed()
|
||||
except Exception:
|
||||
pass
|
||||
self._writer = None
|
||||
self._reader = None
|
||||
if self._daemon_proc:
|
||||
self._daemon_proc.terminate()
|
||||
self._daemon_proc = None
|
||||
logger.info("Signal channel stopped")
|
||||
@@ -0,0 +1,39 @@
|
||||
"""Signal credential validation."""
|
||||
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def validate_signal(
|
||||
phone_number: str,
|
||||
cli_path: str = "signal-cli",
|
||||
rpc_port: int = 7583,
|
||||
) -> tuple[bool, str]:
|
||||
"""Validate Signal setup by checking signal-cli availability.
|
||||
|
||||
Returns:
|
||||
Tuple of (is_valid, message).
|
||||
"""
|
||||
import asyncio
|
||||
import subprocess
|
||||
|
||||
if not phone_number:
|
||||
return False, "phone_number is required"
|
||||
|
||||
# Check signal-cli binary
|
||||
loop = asyncio.get_event_loop()
|
||||
def _check():
|
||||
try:
|
||||
result = subprocess.run(
|
||||
[cli_path, "--version"], capture_output=True, text=True, timeout=5,
|
||||
)
|
||||
if result.returncode == 0:
|
||||
return True, f"signal-cli {result.stdout.strip()}"
|
||||
return False, "signal-cli returned error"
|
||||
except FileNotFoundError:
|
||||
return False, f"signal-cli not found at '{cli_path}'"
|
||||
except Exception as e:
|
||||
return False, f"Error: {e}"
|
||||
|
||||
return await loop.run_in_executor(None, _check)
|
||||
@@ -0,0 +1,99 @@
|
||||
"""Signal channel server.
|
||||
|
||||
Standalone script to run the Signal channel with CLI options.
|
||||
|
||||
Usage:
|
||||
python -m EvoScientist.channels.signal.serve --phone-number NUMBER [OPTIONS]
|
||||
|
||||
Examples:
|
||||
# Basic usage
|
||||
python -m EvoScientist.channels.signal.serve --phone-number +1234567890
|
||||
|
||||
# With custom signal-cli path and allowed senders
|
||||
python -m EvoScientist.channels.signal.serve --phone-number +1234567890 --cli-path /usr/local/bin/signal-cli --allow +9876543210
|
||||
|
||||
# With agent and thinking
|
||||
python -m EvoScientist.channels.signal.serve --phone-number +1234567890 --agent --thinking
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
|
||||
from .channel import SignalChannel, SignalConfig
|
||||
from ..bus import MessageBus
|
||||
from ..standalone import run_standalone
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.DEBUG,
|
||||
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
|
||||
datefmt="%H:%M:%S",
|
||||
)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def parse_args():
|
||||
"""Parse command line arguments."""
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Signal channel server",
|
||||
formatter_class=argparse.RawDescriptionHelpFormatter,
|
||||
)
|
||||
parser.add_argument(
|
||||
"--phone-number",
|
||||
required=True,
|
||||
help="Signal phone number (e.g. +1234567890)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--cli-path",
|
||||
default="signal-cli",
|
||||
help="Path to signal-cli binary (default: signal-cli)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--config-dir",
|
||||
help="signal-cli config directory",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--rpc-port",
|
||||
type=int,
|
||||
default=7583,
|
||||
help="signal-cli JSON RPC port (default: 7583)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--allow",
|
||||
action="append",
|
||||
dest="allowed_senders",
|
||||
help="Allowed sender (phone number). Can be used multiple times.",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--agent",
|
||||
action="store_true",
|
||||
help="Use EvoScientist agent as handler (default: echo)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--thinking",
|
||||
action="store_true",
|
||||
help="Send thinking content as intermediate messages (requires --agent)",
|
||||
)
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def main():
|
||||
"""Entry point."""
|
||||
args = parse_args()
|
||||
|
||||
config = SignalConfig(
|
||||
phone_number=args.phone_number,
|
||||
cli_path=args.cli_path,
|
||||
config_dir=args.config_dir,
|
||||
rpc_port=args.rpc_port,
|
||||
allowed_senders=set(args.allowed_senders) if args.allowed_senders else None,
|
||||
)
|
||||
|
||||
send_thinking = args.thinking and args.agent
|
||||
bus = MessageBus()
|
||||
channel = SignalChannel(config)
|
||||
|
||||
run_standalone(channel, bus, use_agent=args.agent, send_thinking=send_thinking)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -1322,6 +1322,9 @@ def _step_channels(config: EvoScientistConfig) -> dict[str, object]:
|
||||
("feishu", "Feishu", [("feishu_app_id", "App ID"), ("feishu_app_secret", "App Secret")], "aiohttp", "feishu"),
|
||||
("dingtalk", "DingTalk", [("dingtalk_client_id", "Client ID (AppKey)"), ("dingtalk_client_secret", "Client Secret (AppSecret)")], "aiohttp", "dingtalk"),
|
||||
("wechat", "WeChat", [("wechat_wecom_corp_id", "WeCom Corp ID"), ("wechat_wecom_agent_id", "WeCom Agent ID"), ("wechat_wecom_secret", "WeCom Secret")], "aiohttp", "wechat"),
|
||||
("email", "Email", [("email_imap_host", "IMAP host"), ("email_imap_username", "IMAP username"), ("email_imap_password", "IMAP password"), ("email_smtp_host", "SMTP host"), ("email_smtp_username", "SMTP username"), ("email_smtp_password", "SMTP password"), ("email_from_address", "From address")], None, None),
|
||||
("qq", "QQ", [("qq_app_id", "App ID"), ("qq_app_secret", "App Secret")], "botpy", "qq"),
|
||||
("signal", "Signal", [("signal_phone_number", "Phone number (E.164)")], None, None),
|
||||
("imessage", "iMessage", [], None, None), # handled via _setup_imessage()
|
||||
]
|
||||
|
||||
@@ -1546,6 +1549,28 @@ def _probe_channel(
|
||||
_val("dingtalk_client_secret"),
|
||||
_val("dingtalk_proxy") or None,
|
||||
)
|
||||
elif ch_name == "email":
|
||||
from ..channels.email.probe import validate_email_imap
|
||||
return await validate_email_imap(
|
||||
_val("email_imap_host"),
|
||||
int(_val("email_imap_port", "993")),
|
||||
_val("email_imap_username"),
|
||||
_val("email_imap_password"),
|
||||
_val("email_imap_use_ssl", "True").lower() not in ("false", "0", "no"),
|
||||
)
|
||||
elif ch_name == "qq":
|
||||
from ..channels.qq.probe import validate_qq
|
||||
return await validate_qq(
|
||||
_val("qq_app_id"),
|
||||
_val("qq_app_secret"),
|
||||
)
|
||||
elif ch_name == "signal":
|
||||
from ..channels.signal.probe import validate_signal
|
||||
return await validate_signal(
|
||||
_val("signal_phone_number"),
|
||||
_val("signal_cli_path", "signal-cli"),
|
||||
int(_val("signal_rpc_port", "7583")),
|
||||
)
|
||||
else:
|
||||
return True, "No probe available"
|
||||
|
||||
|
||||
@@ -85,7 +85,7 @@ class EvoScientistConfig:
|
||||
show_thinking: bool = True
|
||||
|
||||
# Channel Settings
|
||||
channel_enabled: str = "" # "imessage" | "telegram" | "discord" | "slack" | "wechat" | "dingtalk" | "feishu" | "" (comma-separated for multiple)
|
||||
channel_enabled: str = "" # "imessage" | "telegram" | "discord" | "slack" | "wechat" | "dingtalk" | "feishu" | "email" | "qq" | "signal" | "" (comma-separated for multiple)
|
||||
channel_send_thinking: bool = True # forward thinking to any channel
|
||||
require_mention: str = "group" # "always" | "group" | "off"
|
||||
text_chunk_limit: int = 0 # 0 = use capability default
|
||||
|
||||
@@ -48,12 +48,14 @@ telegram = ["python-telegram-bot>=21.0"]
|
||||
discord = ["discord.py>=2.3"]
|
||||
slack = ["slack-sdk>=3.27", "aiohttp>=3.9"]
|
||||
wechat = ["pycryptodome>=3.20"]
|
||||
qq = ["qq-botpy>=1.0"]
|
||||
all-channels = [
|
||||
"python-telegram-bot>=21.0",
|
||||
"discord.py>=2.3",
|
||||
"aiohttp>=3.9",
|
||||
"slack-sdk>=3.27",
|
||||
"pycryptodome>=3.20",
|
||||
"qq-botpy>=1.0",
|
||||
]
|
||||
|
||||
[project.urls]
|
||||
|
||||
Reference in New Issue
Block a user