新增 IRC 渠道的完整扩展实现,包括客户端连接管理、消息收发与流式处理、安全策略与配对验证、消息去重与规范化、状态监控与健康探测、配置模型与账户管理、错误处理等模块。
368 lines
14 KiB
Python
368 lines
14 KiB
Python
import asyncio
|
|
import logging
|
|
import time
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
|
|
from yuxi.channel.extensions.irc.client import (
|
|
IRCFatalError,
|
|
connect_irc_client,
|
|
graceful_disconnect,
|
|
handle_nick_collision,
|
|
handle_ping,
|
|
identify_nickserv,
|
|
join_channel,
|
|
read_lines,
|
|
send_line,
|
|
)
|
|
from yuxi.channel.extensions.irc.dedupe import IrcDedupeCache
|
|
from yuxi.channel.extensions.irc.protocol import (
|
|
parse_ctcp,
|
|
parse_irc_line,
|
|
parse_irc_prefix,
|
|
sanitize_irc_text,
|
|
)
|
|
from yuxi.channel.extensions.irc.types import IrcInboundMessage, ResolvedIrcAccount
|
|
from yuxi.channel.message.models import GroupContext, MessageType, PeerInfo, UnifiedMessage
|
|
from yuxi.channel.routing.models import PeerKind
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
RECONNECT_BASE_DELAY = 2.0
|
|
RECONNECT_MAX_DELAY = 120.0
|
|
RECONNECT_BACKOFF_FACTOR = 2.0
|
|
RECONNECT_JITTER = 0.5
|
|
|
|
_CHANTYPES = {"#", "&"}
|
|
|
|
|
|
async def _handle_ctcp_query(
|
|
writer: asyncio.StreamWriter,
|
|
parsed: object,
|
|
ctcp_cmd: str,
|
|
ctcp_args: str | None,
|
|
) -> None:
|
|
sender_nick = parse_irc_prefix(parsed.prefix).nick if parsed.prefix else None
|
|
if not sender_nick:
|
|
return
|
|
|
|
match ctcp_cmd:
|
|
case "VERSION":
|
|
reply = "ForcePilot IRC Bot (https://forcepilot.ai)"
|
|
case "PING":
|
|
reply = ctcp_args or str(int(time.time()))
|
|
case "TIME":
|
|
reply = datetime.now(datetime.timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC")
|
|
case "FINGER":
|
|
reply = "ForcePilot AI Bot"
|
|
case "SOURCE":
|
|
reply = "https://github.com/ForcePilot"
|
|
case _:
|
|
return
|
|
|
|
send_line(writer, f"NOTICE {sender_nick} :\x01{ctcp_cmd} {reply}\x01")
|
|
await writer.drain()
|
|
|
|
|
|
def _parse_server_time(tags: dict[str, str | bool]) -> datetime | None:
|
|
time_val = tags.get("time")
|
|
if not isinstance(time_val, str):
|
|
return None
|
|
try:
|
|
if time_val.endswith("Z"):
|
|
time_val = time_val[:-1] + "+00:00"
|
|
return datetime.fromisoformat(time_val)
|
|
except ValueError:
|
|
return None
|
|
|
|
|
|
def create_irc_inbound_message(
|
|
parsed: object,
|
|
current_nick: str,
|
|
) -> IrcInboundMessage | None:
|
|
raw_text = parsed.trailing or ""
|
|
ctcp_cmd, ctcp_args = parse_ctcp(raw_text)
|
|
|
|
if ctcp_cmd == "ACTION":
|
|
text = f"/me {ctcp_args or ''}"
|
|
elif ctcp_cmd is not None:
|
|
return None
|
|
else:
|
|
text = sanitize_irc_text(raw_text)
|
|
|
|
if not text.strip():
|
|
return None
|
|
|
|
raw_target = parsed.params[0] if parsed.params else ""
|
|
is_group = any(raw_target.startswith(c) for c in _CHANTYPES)
|
|
|
|
prefix_info = parse_irc_prefix(parsed.prefix) if parsed.prefix else None
|
|
sender_nick = prefix_info.nick if prefix_info else "unknown"
|
|
sender_user = prefix_info.user if prefix_info else None
|
|
sender_host = prefix_info.host if prefix_info else None
|
|
|
|
if sender_nick.lower() == current_nick.lower():
|
|
return None
|
|
|
|
target = raw_target if is_group else sender_nick
|
|
|
|
server_time = _parse_server_time(parsed.tags)
|
|
account = parsed.tags.get("account")
|
|
account = account if isinstance(account, str) else None
|
|
msg_id = parsed.tags.get("msgid")
|
|
msg_id = msg_id if isinstance(msg_id, str) else None
|
|
|
|
return IrcInboundMessage(
|
|
message_id=str(uuid.uuid4()),
|
|
target=target,
|
|
raw_target=raw_target,
|
|
sender_nick=sender_nick,
|
|
sender_user=sender_user,
|
|
sender_host=sender_host,
|
|
text=text,
|
|
timestamp=time.monotonic(),
|
|
is_group=is_group,
|
|
server_time=server_time,
|
|
account=account,
|
|
msg_id=msg_id,
|
|
)
|
|
|
|
|
|
def irc_inbound_to_unified(msg: IrcInboundMessage, account: ResolvedIrcAccount) -> UnifiedMessage:
|
|
sender_id = msg.sender_nick
|
|
if msg.sender_user and msg.sender_host:
|
|
sender_id = f"{msg.sender_nick}!{msg.sender_user}@{msg.sender_host}"
|
|
|
|
peer = PeerInfo(
|
|
kind=PeerKind.DIRECT if not msg.is_group else PeerKind.GROUP,
|
|
id=msg.target,
|
|
display_name=msg.sender_nick,
|
|
username=msg.sender_nick,
|
|
)
|
|
|
|
group = None
|
|
if msg.is_group:
|
|
group = GroupContext(
|
|
id=msg.target,
|
|
name=msg.target,
|
|
)
|
|
|
|
return UnifiedMessage(
|
|
msg_id=msg.message_id,
|
|
channel_type="irc",
|
|
account_id=account.account_id,
|
|
content=msg.text,
|
|
sender=peer,
|
|
message_type=MessageType.TEXT,
|
|
group=group,
|
|
timestamp=datetime.now(),
|
|
raw_payload={
|
|
"sender_id": sender_id,
|
|
"raw_target": msg.raw_target,
|
|
"is_group": msg.is_group,
|
|
},
|
|
conversation_label=msg.sender_nick if not msg.is_group else msg.target,
|
|
surface="irc",
|
|
originating_channel="irc",
|
|
)
|
|
|
|
|
|
async def monitor_irc_provider(
|
|
account: ResolvedIrcAccount,
|
|
cancel_event: asyncio.Event,
|
|
status_sink,
|
|
on_inbound,
|
|
client_ref: list | None = None,
|
|
dedupe_cache: IrcDedupeCache | None = None,
|
|
) -> None:
|
|
attempt = 0
|
|
while not cancel_event.is_set():
|
|
try:
|
|
await _monitor_irc_session(account, cancel_event, status_sink, on_inbound, client_ref, dedupe_cache)
|
|
except IRCFatalError as e:
|
|
logger.error("IRC fatal error: %s", e)
|
|
except (ConnectionError, OSError) as e:
|
|
logger.warning("IRC connection error: %s", e)
|
|
except asyncio.CancelledError:
|
|
break
|
|
|
|
if cancel_event.is_set():
|
|
break
|
|
|
|
delay = min(
|
|
RECONNECT_BASE_DELAY * (RECONNECT_BACKOFF_FACTOR ** attempt),
|
|
RECONNECT_MAX_DELAY,
|
|
)
|
|
delay += RECONNECT_JITTER * (time.monotonic() % 1.0)
|
|
attempt += 1
|
|
logger.info("IRC reconnecting in %.1fs (attempt %d)", delay, attempt)
|
|
try:
|
|
await asyncio.wait_for(cancel_event.wait(), timeout=delay)
|
|
break
|
|
except asyncio.TimeoutError:
|
|
pass
|
|
|
|
status_sink(connected=False)
|
|
|
|
|
|
async def _monitor_irc_session(
|
|
account: ResolvedIrcAccount,
|
|
cancel_event: asyncio.Event,
|
|
status_sink,
|
|
on_inbound,
|
|
client_ref: list | None = None,
|
|
dedupe_cache: IrcDedupeCache | None = None,
|
|
) -> None:
|
|
client = await connect_irc_client(account)
|
|
if client_ref is not None:
|
|
client_ref[0] = client
|
|
status_sink(connected=True)
|
|
|
|
current_nick = account.nick
|
|
|
|
try:
|
|
async for line in read_lines(client.reader, cancel_event):
|
|
parsed = parse_irc_line(line)
|
|
match parsed.command:
|
|
case "PING":
|
|
await handle_ping(client.writer, parsed.trailing)
|
|
case "001":
|
|
if account.nickserv_password:
|
|
await identify_nickserv(client.writer, account.nickserv_password)
|
|
logger.info("NickServ IDENTIFY sent for %s", current_nick)
|
|
for channel in account.config.channels:
|
|
await join_channel(client.writer, channel)
|
|
logger.info("Joined channel %s", channel)
|
|
case "PRIVMSG":
|
|
raw_text = parsed.trailing or ""
|
|
ctcp_cmd, ctcp_args = parse_ctcp(raw_text)
|
|
if ctcp_cmd in ("VERSION", "PING", "TIME", "FINGER", "SOURCE"):
|
|
await _handle_ctcp_query(client.writer, parsed, ctcp_cmd, ctcp_args)
|
|
continue
|
|
is_echo = parsed.tags.get("msgid") and parsed.tags.get("account") == account.nick
|
|
if is_echo:
|
|
logger.debug("Echo-message ignored")
|
|
continue
|
|
dedupe_key = f"{account.account_id}:{parsed.prefix}:{parsed.trailing}"
|
|
if dedupe_cache is not None and dedupe_key in dedupe_cache:
|
|
continue
|
|
if dedupe_cache is not None:
|
|
dedupe_cache.add(dedupe_key)
|
|
inbound_msg = create_irc_inbound_message(parsed, current_nick)
|
|
if inbound_msg:
|
|
await on_inbound(inbound_msg)
|
|
case "NOTICE":
|
|
notice_text = sanitize_irc_text(parsed.trailing or "")
|
|
source = parsed.params[0] if parsed.params else ""
|
|
prefix_info = parse_irc_prefix(parsed.prefix) if parsed.prefix else None
|
|
sender_nick = prefix_info.nick if prefix_info else source
|
|
|
|
is_server_notice = (
|
|
prefix_info is None
|
|
or prefix_info.server is not None
|
|
or (sender_nick or "").lower() in (
|
|
"nickserv", "chanserv", "authserv", "memoserv", "operserv", "hostserv"
|
|
)
|
|
)
|
|
|
|
if is_server_notice and notice_text:
|
|
um = UnifiedMessage(
|
|
msg_id=str(uuid.uuid4()),
|
|
channel_type="irc",
|
|
account_id=account.account_id,
|
|
content=f"[IRC Notice] {notice_text}",
|
|
sender=PeerInfo(
|
|
kind=PeerKind.DIRECT,
|
|
id=sender_nick,
|
|
display_name=sender_nick,
|
|
username=sender_nick,
|
|
),
|
|
message_type=MessageType.SYSTEM,
|
|
timestamp=datetime.now(),
|
|
surface="irc",
|
|
originating_channel="irc",
|
|
)
|
|
await on_inbound(um)
|
|
else:
|
|
logger.debug("NOTICE from %s: %s", source, notice_text[:100])
|
|
case "KICK":
|
|
kicked_nick = parsed.params[1] if len(parsed.params) > 1 else ""
|
|
if kicked_nick.lower() == current_nick.lower():
|
|
logger.warning(
|
|
"Kicked from %s by %s: %s",
|
|
parsed.params[0] if parsed.params else "?",
|
|
parse_irc_prefix(parsed.prefix).nick if parsed.prefix else "server",
|
|
parsed.trailing or "",
|
|
)
|
|
status_sink(connected=True, warning=f"Kicked from {parsed.params[0]}")
|
|
case "005":
|
|
for param in parsed.params[:-1] if len(parsed.params) > 1 else []:
|
|
if "=" in param:
|
|
key, val = param.split("=", 1)
|
|
if key == "CHANTYPES":
|
|
global _CHANTYPES
|
|
_CHANTYPES = set(val)
|
|
logger.debug("IRC ISUPPORT CHANTYPES=%s", val)
|
|
elif key == "CASEMAPPING":
|
|
logger.debug("IRC ISUPPORT CASEMAPPING=%s", val)
|
|
elif key == "UTF8ONLY":
|
|
logger.debug("IRC ISUPPORT UTF8ONLY")
|
|
case "INVITE":
|
|
inviter = parse_irc_prefix(parsed.prefix).nick if parsed.prefix else "unknown"
|
|
channel = parsed.params[1] if len(parsed.params) > 1 else ""
|
|
logger.info("Invited to %s by %s", channel, inviter)
|
|
if channel:
|
|
await join_channel(client.writer, channel)
|
|
logger.info("Auto-joined invited channel %s", channel)
|
|
case "TOPIC":
|
|
channel = parsed.params[0] if parsed.params else ""
|
|
topic = parsed.trailing or ""
|
|
logger.info("Topic for %s: %s", channel, topic[:100])
|
|
case "ACCOUNT":
|
|
nick = parse_irc_prefix(parsed.prefix).nick if parsed.prefix else "unknown"
|
|
account_name = parsed.params[0] if parsed.params else ""
|
|
logger.debug("ACCOUNT %s -> %s", nick, account_name)
|
|
case "AWAY":
|
|
nick = parse_irc_prefix(parsed.prefix).nick if parsed.prefix else "unknown"
|
|
is_away = parsed.trailing is not None
|
|
logger.debug("AWAY %s away=%s", nick, is_away)
|
|
case "CHGHOST":
|
|
nick = parse_irc_prefix(parsed.prefix).nick if parsed.prefix else "unknown"
|
|
new_user = parsed.params[0] if parsed.params else ""
|
|
new_host = parsed.params[1] if len(parsed.params) > 1 else ""
|
|
logger.debug("CHGHOST %s -> %s@%s", nick, new_user, new_host)
|
|
case "471" | "473" | "474" | "475":
|
|
channel = parsed.params[1] if len(parsed.params) > 1 else "?"
|
|
logger.warning(
|
|
"IRC JOIN failed for %s: %s (%s)",
|
|
channel, parsed.trailing, parsed.command,
|
|
)
|
|
status_sink(connected=True, warning=f"JOIN {channel} failed: {parsed.trailing}")
|
|
case "401" | "403" | "404" | "405":
|
|
target = parsed.params[1] if len(parsed.params) > 1 else "?"
|
|
logger.warning(
|
|
"IRC send failed to %s: %s (%s)",
|
|
target, parsed.trailing, parsed.command,
|
|
)
|
|
case "433" | "436":
|
|
current_nick = await handle_nick_collision(client.writer, account, current_nick)
|
|
client.current_nick = current_nick
|
|
case "432" | "464" | "465":
|
|
raise IRCFatalError(f"IRC fatal error {parsed.command}: {parsed.trailing}")
|
|
case "ERROR":
|
|
raise IRCFatalError(f"Server sent ERROR: {parsed.trailing or 'unknown'}")
|
|
case "NICK":
|
|
if parsed.prefix:
|
|
prefix_info = parse_irc_prefix(parsed.prefix)
|
|
if prefix_info.nick and prefix_info.nick.lower() == current_nick.lower():
|
|
current_nick = parsed.params[0] if parsed.params else current_nick
|
|
client.current_nick = current_nick
|
|
logger.info("Nick changed to %s", current_nick)
|
|
except asyncio.CancelledError:
|
|
pass
|
|
except IRCFatalError:
|
|
raise
|
|
finally:
|
|
await graceful_disconnect(client)
|
|
status_sink(connected=False)
|