新增Twitch IRC协议相关的全套实现,包括: 1. 基础工具类:令牌处理、消息格式化、速率限制、消息去重 2. 核心适配器组件:IRC解析器、消息归一化、外发消息处理 3. API客户端:Helix API封装、认证提供者 4. 配置与部署:配置校验、设置向导 5. 辅助功能:配对管理、健康检查、目标解析等
92 lines
3.1 KiB
Python
92 lines
3.1 KiB
Python
from __future__ import annotations
|
|
|
|
from typing import Any
|
|
|
|
from yuxi.channels.models import DeliveryResult
|
|
|
|
from .send import format_privmsg_line, format_action_line
|
|
|
|
|
|
class TwitchOutboundAdapter:
|
|
def __init__(self, writer: Any, rate_limiter: Any, silent: bool = False):
|
|
self._writer = writer
|
|
self._rate_limiter = rate_limiter
|
|
self._silent = silent
|
|
|
|
@property
|
|
def silent(self) -> bool:
|
|
return self._silent
|
|
|
|
@silent.setter
|
|
def silent(self, value: bool) -> None:
|
|
self._silent = value
|
|
|
|
def _fmt(self, target: str, text: str) -> str:
|
|
return format_action_line(target, text) if self._silent else format_privmsg_line(target, text)
|
|
|
|
@staticmethod
|
|
def split_text(text: str, target: str, prefix_len: int | None = None) -> list[str]:
|
|
if prefix_len is None:
|
|
prefix_len = len(f"PRIVMSG {target} :")
|
|
available = max(200, 510 - prefix_len)
|
|
if available < 1:
|
|
available = 200
|
|
|
|
chunks: list[str] = []
|
|
if not text:
|
|
return [""]
|
|
remaining = text
|
|
|
|
while remaining:
|
|
encoded = remaining.encode("utf-8")
|
|
if len(encoded) <= available:
|
|
chunks.append(remaining)
|
|
break
|
|
|
|
cut = _find_utf8_cut(encoded, available)
|
|
split_pos = len(encoded[:cut].decode("utf-8", errors="replace"))
|
|
|
|
nl = remaining.rfind("\n", 0, split_pos)
|
|
if nl > split_pos * 0.5:
|
|
split_pos = nl + 1
|
|
else:
|
|
sp = remaining.rfind(" ", 0, split_pos)
|
|
if sp > split_pos * 0.5:
|
|
split_pos = sp + 1
|
|
|
|
chunk = remaining[:split_pos].rstrip()
|
|
chunks.append(chunk)
|
|
remaining = remaining[split_pos:].lstrip()
|
|
|
|
return chunks
|
|
|
|
async def send_chunks(self, target: str, chunks: list[str]) -> DeliveryResult:
|
|
for chunk in chunks:
|
|
if self._rate_limiter and not await self._rate_limiter.acquire():
|
|
return DeliveryResult(success=False, error="rate_limit_exceeded")
|
|
self._writer.write(self._fmt(target, chunk).encode("utf-8") + b"\r\n")
|
|
await self._writer.drain()
|
|
return DeliveryResult(success=True)
|
|
|
|
async def send_text(self, target: str, text: str) -> DeliveryResult:
|
|
chunks = self.split_text(text, target)
|
|
return await self.send_chunks(target, chunks)
|
|
|
|
async def send_single_line(self, target: str, text: str) -> DeliveryResult:
|
|
if self._rate_limiter and not await self._rate_limiter.acquire():
|
|
return DeliveryResult(success=False, error="rate_limit_exceeded")
|
|
self._writer.write(self._fmt(target, text).encode("utf-8") + b"\r\n")
|
|
await self._writer.drain()
|
|
return DeliveryResult(success=True)
|
|
|
|
|
|
def _find_utf8_cut(encoded: bytes, byte_limit: int) -> int:
|
|
cut = byte_limit
|
|
while cut > 0 and (encoded[cut - 1] & 0xC0) == 0x80:
|
|
cut -= 1
|
|
while cut > 0 and (encoded[cut - 1] & 0xC0) == 0xC0:
|
|
cut -= 1
|
|
if cut == 0:
|
|
cut = max(1, byte_limit)
|
|
return cut
|