diff --git a/plugins/platforms/line/adapter.py b/plugins/platforms/line/adapter.py index 2129680804..8cb0cdf06d 100644 --- a/plugins/platforms/line/adapter.py +++ b/plugins/platforms/line/adapter.py @@ -1,18 +1,12 @@ -"""LINE Messaging API platform adapter for Hermes Agent. - -An aiohttp webhook server accepts signature-verified LINE events and relays them -through ``BasePlatformAdapter``. Design highlights: +"""LINE Messaging API adapter: aiohttp webhook (signature-verified) → BasePlatformAdapter. * Reply token preferred (free, single-use, ~60s TTL), metered Push as fallback. -* Optional slow-LLM postback button: past ``slow_response_threshold`` (default 45s) - the reply token is burned on a Template Buttons bubble; tapping it yields a fresh - free token carrying the cached answer (PENDING → READY → DELIVERED, ERROR on - cancel). Threshold 0 disables the button. -* Three allowlists: users (U…), groups (C…), rooms (R…); ``LINE_ALLOW_ALL_USERS`` - is a dev-only escape hatch. -* Media via public HTTPS: LINE takes no binary uploads, so local files are served - under ``/line/media//`` (allowed-roots guard); ``LINE_PUBLIC_URL`` - overrides host:port behind tunnels/proxies or wildcard binds. +* Slow-LLM postback button past ``slow_response_threshold`` (45s; 0 disables): the reply + token is burned on a Template Buttons bubble; the tap yields a fresh free token that + delivers the cached answer (PENDING → READY → DELIVERED, ERROR on cancel). +* Three allowlists (users U…, groups C…, rooms R…); ``LINE_ALLOW_ALL_USERS`` is dev-only. +* Media via public HTTPS only: local files served under ``/line/media//`` + (allowed-roots guard); ``LINE_PUBLIC_URL`` overrides host:port behind tunnels/wildcards. * ≤5 message objects per call; text chunked at 4500 chars (bubble hard limit 5000). """ @@ -39,20 +33,13 @@ from typing import Any, Callable, Dict, List, Optional, Set, Tuple from urllib.parse import quote as _urlquote from gateway.platforms._shared import get_scoped_secret as _get_scoped_secret - - -logger = logging.getLogger(__name__) - -# Plugin discovery imports adapter.py late enough that gateway is already loaded, -# so gateway internals are imported eagerly here. from gateway.platforms.base import ( gateway_trust_env, BasePlatformAdapter, MessageEvent, MessageType, SendResult, cache_audio_from_bytes, cache_document_from_bytes, cache_image_from_bytes, cache_video_from_bytes, ) from gateway.config import Platform - -# --- Constants --- +logger = logging.getLogger(__name__) LINE_REPLY_URL = "https://api.line.me/v2/bot/message/reply" LINE_PUSH_URL = "https://api.line.me/v2/bot/message/push" @@ -60,13 +47,10 @@ LINE_LOADING_URL = "https://api.line.me/v2/bot/chat/loading/start" LINE_CONTENT_URL_FMT = "https://api-data.line.me/v2/bot/message/{message_id}/content" LINE_BOT_INFO_URL = "https://api.line.me/v2/bot/info" -# LINE Messaging API hard limits -LINE_PER_BUBBLE_CHARS = 5000 +LINE_PER_BUBBLE_CHARS = 5000 # LINE hard limit LINE_SAFE_BUBBLE_CHARS = 4500 # conservative chunking limit LINE_MAX_MESSAGES_PER_CALL = 5 LINE_REPLY_TOKEN_TTL_SECONDS = 50 # below LINE's ~60s - -# Webhook hardening WEBHOOK_BODY_MAX_BYTES = 1_048_576 # 1 MiB — webhooks are tiny JSON DEFAULT_WEBHOOK_PORT = 8646 DEFAULT_WEBHOOK_PATH = "/line/webhook" @@ -80,14 +64,12 @@ DEFAULT_HOST = None # Bind hosts LINE's servers could never fetch media from → public URL required. _WILDCARD_HOSTS = frozenset({"0.0.0.0", "::", ""}) -# Slow-LLM postback button defaults -DEFAULT_SLOW_RESPONSE_THRESHOLD = 45.0 # seconds; 0 disables +DEFAULT_SLOW_RESPONSE_THRESHOLD = 45.0 # seconds; 0 disables the postback button DEFAULT_PENDING_REPLY_TEXT = "🤔 Still thinking. Tap below to fetch the answer when it's ready." DEFAULT_BUTTON_LABEL = "Get answer" DEFAULT_DELIVERED_TEXT = "Already replied ✅" DEFAULT_INTERRUPTED_TEXT = "Run was interrupted before completion." -# Media defaults MEDIA_TOKEN_TTL_SECONDS = 1800 # 30 minutes; LINE caches the URL aggressively LINE_IMAGE_MAX_BYTES = 10 * 1024 * 1024 # 10 MB per LINE docs LINE_AV_MAX_BYTES = 200 * 1024 * 1024 # 200 MB for voice/video @@ -107,15 +89,16 @@ _FALLBACK_PNG_PREVIEW = bytes.fromhex( ) -# --- Markdown stripping (URL-preserving) --- - -_MD_LINK_RE = re.compile(r"\[([^\]]+)\]\((https?://[^\s)]+)\)") -_MD_BOLD_RE = re.compile(r"\*\*(.+?)\*\*") -_MD_ITAL_RE = re.compile(r"(? str: @@ -123,13 +106,8 @@ def strip_markdown_preserving_urls(text: str) -> str: tappable (LINE auto-links bare URLs only). Code-block content is kept.""" if not text: return text - text = _MD_CODE_BLOCK_RE.sub(lambda m: m.group(1).rstrip("\n"), text) - text = _MD_CODE_INLINE_RE.sub(r"\1", text) - text = _MD_LINK_RE.sub(lambda m: f"{m.group(1)} ({m.group(2)})", text) - text = _MD_BOLD_RE.sub(r"\1", text) - text = _MD_ITAL_RE.sub(r"\1", text) - text = _MD_HEADING_RE.sub("", text) - text = _MD_BULLET_RE.sub("• ", text) + for pattern, repl in _MD_STRIP_RULES: + text = pattern.sub(repl, text) return text @@ -162,8 +140,6 @@ def split_for_line(text: str, max_chars: int = LINE_SAFE_BUBBLE_CHARS) -> List[s return chunks -# --- Webhook signature verification --- - def verify_line_signature(body: bytes, signature: str, channel_secret: str) -> bool: """Verify ``X-Line-Signature``: base64(HMAC-SHA256(secret, raw body)), constant-time.""" if not signature or not channel_secret or body is None: @@ -173,14 +149,13 @@ def verify_line_signature(body: bytes, signature: str, channel_secret: str) -> b expected = base64.b64encode(digest).decode("utf-8") except Exception: return False - # Compare as bytes: compare_digest raises TypeError on non-ASCII str, and - # the signature is a raw request header. + # Bytes: compare_digest raises TypeError on non-ASCII str, and the header is raw. return hmac.compare_digest(expected.encode(), signature.encode()) -# --- Cache state machine — slow-LLM postback flow --- - class State(enum.Enum): + """Slow-LLM postback cache states.""" + PENDING = "pending" # button sent, LLM still running READY = "ready" # response cached, waiting for postback tap DELIVERED = "delivered" @@ -200,7 +175,7 @@ _KEEP = object() # sentinel: leave entry.payload untouched class RequestCache: - """In-memory cache for slow-LLM postback retrieval.""" + """In-memory cache for slow-LLM postback retrieval (PENDING → READY|ERROR → DELIVERED).""" def __init__(self) -> None: self._entries: Dict[str, _CacheEntry] = {} @@ -232,8 +207,6 @@ class RequestCache: self._transition(request_id, {State.READY, State.ERROR}, State.DELIVERED) -# --- Inbound dedup --- - class _MessageDeduplicator: """Bounded LRU of LINE webhook event IDs to ignore at-least-once retries.""" @@ -246,16 +219,13 @@ class _MessageDeduplicator: return False if event_id in self._seen: return True - if len(self._seen) >= self._max: - # Drop the oldest 10% so we don't trim on every insert. + if len(self._seen) >= self._max: # drop the oldest 10% so we don't trim every insert cutoff = sorted(self._seen.values())[len(self._seen) // 10 or 1] self._seen = {k: v for k, v in self._seen.items() if v > cutoff} self._seen[event_id] = time.time() return False -# --- Source / chat-id resolution --- - # LINE source type → (id key, normalized chat_type) _SOURCE_KINDS = {"group": ("groupId", "group"), "room": ("roomId", "room"), "user": ("userId", "dm")} @@ -269,12 +239,7 @@ def _resolve_chat(source: Dict[str, Any]) -> Tuple[str, str]: def _allowed_for_source( - source: Dict[str, Any], - *, - allow_all: bool, - user_ids: Set[str], - group_ids: Set[str], - room_ids: Set[str], + source: Dict[str, Any], *, allow_all: bool, user_ids: Set[str], group_ids: Set[str], room_ids: Set[str], ) -> bool: """Three-list gate: users, groups, rooms.""" if allow_all: @@ -283,8 +248,6 @@ def _allowed_for_source( return bool(sid) and sid in {"dm": user_ids, "group": group_ids, "room": room_ids}[chat_type] -# --- LINE Reply / Push HTTP client --- - class _LineClient: """Thin aiohttp wrapper around the LINE Messaging API (no ``line-bot-sdk`` dependency).""" @@ -316,8 +279,7 @@ class _LineClient: if not chat_id or not chat_id.startswith("U"): return import aiohttp # noqa: F401 — ImportError must escape the swallow-all below - # LINE caps loadingSeconds in 5-step increments, max 60. - clamped = max(5, min(60, (seconds // 5) * 5 or 5)) + clamped = max(5, min(60, (seconds // 5) * 5 or 5)) # LINE: 5-step increments, max 60 try: async with self._session(5.0) as session: await session.post( @@ -349,8 +311,6 @@ class _LineClient: return None -# --- Message builders --- - def _text_message(text: str) -> Dict[str, Any]: """Build a LINE text message object, capped to per-bubble max.""" if len(text) > LINE_PER_BUBBLE_CHARS: @@ -364,9 +324,7 @@ def _text_messages(content: str) -> List[Dict[str, Any]]: return [_text_message(c) for c in chunks][:LINE_MAX_MESSAGES_PER_CALL] -def build_postback_button_message( - text: str, button_label: str, request_id: str -) -> Dict[str, Any]: +def build_postback_button_message(text: str, button_label: str, request_id: str) -> Dict[str, Any]: """Slow-LLM postback bubble. Template Buttons stay tappable from history (Quick Reply chips vanish on the next message). LINE limits: text ≤160, altText ≤400.""" truncated = text if len(text) <= 160 else text[:157] + "..." @@ -377,11 +335,7 @@ def build_postback_button_message( "data": json.dumps({"action": "show_response", "request_id": request_id}), "displayText": button_label[:300] or "Get answer", } - return { - "type": "template", - "altText": alt, - "template": {"type": "buttons", "text": truncated, "actions": [action]}, - } + return {"type": "template", "altText": alt, "template": {"type": "buttons", "text": truncated, "actions": [action]}} # Gateway busy-ack prefixes (interrupting / queued / steered / background review); @@ -390,24 +344,16 @@ _SYSTEM_BYPASS_PREFIXES: Tuple[str, ...] = ("⚡ Interrupting", "⏳ Queued", " def _is_system_bypass(content: str) -> bool: - if not content: - return False - return any(content.startswith(p) for p in _SYSTEM_BYPASS_PREFIXES) + return bool(content) and any(content.startswith(p) for p in _SYSTEM_BYPASS_PREFIXES) -# --- Configuration helpers --- - def _csv_set(value: str) -> Set[str]: - if not value: - return set() - return {x.strip() for x in value.split(",") if x.strip()} + return {x.strip() for x in (value or "").split(",") if x.strip()} def _truthy_env(name: str, default: bool = False) -> bool: v = os.getenv(name) - if v is None: - return default - return v.strip().lower() in {"1", "true", "yes", "on"} + return default if v is None else v.strip().lower() in {"1", "true", "yes", "on"} def _credentials(config) -> Tuple[str, str]: @@ -440,8 +386,6 @@ _OUTBOUND_MEDIA = { _INBOUND_MEDIA_EXT = {"image": ".jpg", "audio": ".m4a", "video": ".mp4", "file": ".bin"} -# --- Adapter --- - class LineAdapter(BasePlatformAdapter): """LINE Messaging API gateway adapter (no message editing → REQUIRES_EDIT_FINALIZE stays False).""" @@ -457,22 +401,16 @@ class LineAdapter(BasePlatformAdapter): return _csv_set(os.getenv(env, "")) | set(extra.get(key, [])) self.channel_access_token, self.channel_secret = _credentials(config) - - # Webhook server. Host default ``None`` → dual-stack bind (see DEFAULT_HOST); - # ``LINE_HOST``/extra.host pin an address, empty string collapses to None. + # Host ``None`` → dual-stack bind (see DEFAULT_HOST); empty string collapses to None. self.webhook_host = env_or("LINE_HOST", "host", DEFAULT_HOST) or DEFAULT_HOST self.webhook_port = _coerce(int, env_or("LINE_PORT", "port", DEFAULT_WEBHOOK_PORT), DEFAULT_WEBHOOK_PORT) self.webhook_path = extra.get("webhook_path", DEFAULT_WEBHOOK_PATH) - - # Public base URL — required for media when the bind isn't publicly reachable. + # Required for media when the bind isn't publicly reachable. self.public_base_url = (env_or("LINE_PUBLIC_URL", "public_url") or "").rstrip("/") - - # Three-allowlist gating self.allow_all = _truthy_env("LINE_ALLOW_ALL_USERS", bool(extra.get("allow_all_users", False))) self.allowed_users = allowlist("LINE_ALLOWED_USERS", "allowed_users") self.allowed_groups = allowlist("LINE_ALLOWED_GROUPS", "allowed_groups") self.allowed_rooms = allowlist("LINE_ALLOWED_ROOMS", "allowed_rooms") - # Slow-LLM postback button threshold + user-overridable copy threshold = env_or("LINE_SLOW_RESPONSE_THRESHOLD", "slow_response_threshold", DEFAULT_SLOW_RESPONSE_THRESHOLD) self.slow_response_threshold = _coerce(float, threshold, DEFAULT_SLOW_RESPONSE_THRESHOLD) @@ -480,27 +418,18 @@ class LineAdapter(BasePlatformAdapter): self.button_label = env_or("LINE_BUTTON_LABEL", "button_label", DEFAULT_BUTTON_LABEL) self.delivered_text = env_or("LINE_DELIVERED_TEXT", "delivered_text", DEFAULT_DELIVERED_TEXT) self.interrupted_text = env_or("LINE_INTERRUPTED_TEXT", "interrupted_text", DEFAULT_INTERRUPTED_TEXT) - # Runtime state self._client: Optional[_LineClient] = None - self._app = None # aiohttp.web.Application - self._runner = None # aiohttp.web.AppRunner - self._site = None # aiohttp.web.TCPSite + self._app = self._runner = self._site = None # aiohttp web.Application / AppRunner / TCPSite self._reply_tokens: Dict[str, Tuple[str, float]] = {} # chat_id → (token, expiry) self._cache = RequestCache() self._dedup = _MessageDeduplicator() self._bot_user_id: Optional[str] = None self._lock_key: Optional[str] = None - - # Media state self._media_tokens: Dict[str, Tuple[str, float]] = {} # token → (path, expiry) self._media_temp_paths: Set[str] = set() self._media_ttl = MEDIA_TOKEN_TTL_SECONDS - - # One outstanding postback button per chat: chat_id → cache request_id. - self._pending_buttons: Dict[str, str] = {} - - # --- Connection lifecycle --- + self._pending_buttons: Dict[str, str] = {} # one outstanding button per chat: chat_id → request_id def _fail(self, code: str, detail: str, *, retryable: bool = False) -> bool: """Record a fatal connect error and return False.""" @@ -537,8 +466,7 @@ class LineAdapter(BasePlatformAdapter): self._app = web.Application(client_max_size=WEBHOOK_BODY_MAX_BYTES) self._app.router.add_post(self.webhook_path, self._handle_webhook) - # Public health probe — useful for tunnel/proxy verification. - self._app.router.add_get(f"{self.webhook_path}/health", self._handle_health) + self._app.router.add_get(f"{self.webhook_path}/health", self._handle_health) # tunnel/proxy probe self._app.router.add_get(f"{DEFAULT_MEDIA_PATH_PREFIX}/{{token}}/{{filename}}", self._handle_media) # Plugin-registered routes must be wired before AppRunner.setup() freezes the router. self._wire_plugin_handlers(self._app) @@ -592,17 +520,13 @@ class LineAdapter(BasePlatformAdapter): pass self._lock_key = None - # --- Webhook handlers --- - async def _handle_health(self, request) -> Any: from aiohttp import web return web.json_response({"status": "ok", "platform": "line"}) async def _handle_webhook(self, request) -> Any: from aiohttp import web - - # Explicit body cap: aiohttp's client_max_size only covers some body modes. - try: + try: # explicit body cap: aiohttp's client_max_size only covers some body modes body = await request.read() except Exception as exc: logger.debug("LINE: read failed: %s", exc) @@ -627,13 +551,9 @@ class LineAdapter(BasePlatformAdapter): event_type = event.get("type") source = event.get("source") or {} webhook_event_id = event.get("webhookEventId", "") or "" - - # LINE webhooks may be re-delivered. - if webhook_event_id and self._dedup.is_duplicate(webhook_event_id): + if webhook_event_id and self._dedup.is_duplicate(webhook_event_id): # at-least-once redelivery logger.debug("LINE: ignoring duplicate webhook event %s", webhook_event_id) return - - # Self-echo filter. if self._bot_user_id and source.get("userId", "") == self._bot_user_id: return if not _allowed_for_source( @@ -665,8 +585,7 @@ class LineAdapter(BasePlatformAdapter): media_types: List[str] = [] if msg_type == "text": text = msg.get("text", "") or "" - elif msg_type in _INBOUND_MEDIA_EXT: - # Fetch the binary, cache it, surface a vision-friendly local path. + elif msg_type in _INBOUND_MEDIA_EXT: # fetch, cache, surface a vision-friendly local path local_path, media_type = await self._download_media( message_id, msg_type, filename=msg.get("fileName") or msg.get("file_name") ) @@ -681,9 +600,7 @@ class LineAdapter(BasePlatformAdapter): text = f"[location: {msg.get('title', '')} {msg.get('address', '')}]".strip() else: text = f"[unsupported message type: {msg_type}]" - - # Best-effort typing indicator (DM only). - if chat_type == "dm" and self._client: + if chat_type == "dm" and self._client: # best-effort typing indicator (DM only) asyncio.create_task(self._client.loading(chat_id)) source_obj = self.build_source( chat_id=chat_id, chat_type=chat_type, user_id=user_id, user_name=user_id, chat_name=chat_id @@ -715,6 +632,7 @@ class LineAdapter(BasePlatformAdapter): def _settle() -> None: self._cache.mark_delivered(request_id) self._pending_buttons.pop(chat_id, None) + if entry.state is State.READY: messages = _text_messages(str(entry.payload or "")) try: @@ -734,8 +652,7 @@ class LineAdapter(BasePlatformAdapter): _settle() except Exception as exc: logger.warning("LINE: postback ERROR reply failed: %s", exc) - elif entry.state in (State.DELIVERED, State.PENDING): - # DELIVERED → "already replied"; PENDING → re-issue the wait notice. + elif entry.state in (State.DELIVERED, State.PENDING): # "already replied" / re-issue wait notice text = self.delivered_text if entry.state is State.DELIVERED else self.pending_text try: await self._client.reply(reply_token, [_text_message(text)]) @@ -766,15 +683,13 @@ class LineAdapter(BasePlatformAdapter): logger.warning("LINE: failed to cache %s payload: %s", msg_type, exc) return None, "" - # --- Outbound send (text) --- - async def send( self, chat_id: str, content: str, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None ) -> SendResult: if not self._client: return SendResult(success=False, error="LINE adapter not connected") - # With a PENDING postback button outstanding, cache the response for the tap — - # except system busy-acks, which must land as visible bubbles. + # A PENDING postback button caches the response for the tap — except system + # busy-acks, which must land as visible bubbles. pending_rid = self._pending_buttons.get(chat_id) if pending_rid and not _is_system_bypass(content): self._cache.set_ready(pending_rid, content) @@ -786,10 +701,7 @@ class LineAdapter(BasePlatformAdapter): def _consume_reply_token(self, chat_id: str) -> Tuple[str, bool]: """Consume a stashed reply token if present and unexpired → ``(token, used_reply)``.""" - entry = self._reply_tokens.pop(chat_id, None) - if not entry: - return "", False - token, expires_at = entry + token, expires_at = self._reply_tokens.pop(chat_id, None) or ("", 0.0) if not token or time.time() >= expires_at: return "", False return token, True @@ -808,19 +720,14 @@ class LineAdapter(BasePlatformAdapter): """Strip Markdown that LINE can't render. URLs are preserved.""" return strip_markdown_preserving_urls(content) - # --- Slow-LLM postback button — driven by _keep_typing --- - async def _keep_typing(self, chat_id: str, *args, **kwargs) -> None: - """Wrap the base typing heartbeat; fire the postback button at threshold.""" + """Wrap the base typing heartbeat; fire the slow-LLM postback button at threshold.""" if self.slow_response_threshold <= 0 or not self._client or not chat_id: await super()._keep_typing(chat_id, *args, **kwargs) return async def _fire_postback() -> None: - try: - await asyncio.sleep(self.slow_response_threshold) - except asyncio.CancelledError: - raise + await asyncio.sleep(self.slow_response_threshold) # Only fire while a usable reply token remains (the agent responding # consumes it) and no button is already outstanding. if chat_id not in self._reply_tokens or chat_id in self._pending_buttons: @@ -838,6 +745,7 @@ class LineAdapter(BasePlatformAdapter): except Exception as exc: logger.warning("LINE: postback button send failed: %s", exc) self._pending_buttons.pop(chat_id, None) + post_task = asyncio.create_task(_fire_postback()) try: await super()._keep_typing(chat_id, *args, **kwargs) @@ -856,13 +764,10 @@ class LineAdapter(BasePlatformAdapter): if rid: self._cache.set_error(rid, self.interrupted_text) - # --- Outbound media (image / voice / video) --- - def _register_media(self, file_path: str, *, cleanup: bool = False) -> str: - """Register a local file for HTTPS serving; return the URL token.""" + """Register a local file for HTTPS serving (evicting expired tokens); return the URL token.""" now = time.time() - for token in list(self._media_tokens.keys()): - path, exp = self._media_tokens[token] + for token, (path, exp) in list(self._media_tokens.items()): if now > exp: self._media_tokens.pop(token, None) if path in self._media_temp_paths: @@ -884,8 +789,7 @@ class LineAdapter(BasePlatformAdapter): # guard should have fired); fall back to localhost so the URL is well-formed. host = "127.0.0.1" if self._missing_public_url() else self.webhook_host base = f"https://{host}" if self.webhook_port == 443 else f"https://{host}:{self.webhook_port}" - safe_name = _urlquote(filename, safe="") - return f"{base}{DEFAULT_MEDIA_PATH_PREFIX}/{token}/{safe_name}" + return f"{base}{DEFAULT_MEDIA_PATH_PREFIX}/{token}/{_urlquote(filename, safe='')}" def _serve_file(self, path: Path) -> str: """Register ``path`` for serving and return its public URL.""" @@ -893,9 +797,7 @@ class LineAdapter(BasePlatformAdapter): def _missing_public_url(self) -> bool: """True when no LINE_PUBLIC_URL is set and the bind host is wildcard/dual-stack ``None``.""" - if self.public_base_url: - return False - return self.webhook_host is None or self.webhook_host in _WILDCARD_HOSTS + return not self.public_base_url and (self.webhook_host is None or self.webhook_host in _WILDCARD_HOSTS) def _check_media_file(self, kind: str, file_path: str) -> Tuple[Optional[Path], Optional[SendResult]]: """Shared preflight for send_image_file/send_voice/send_video → ``(path, error)``.""" @@ -912,10 +814,8 @@ class LineAdapter(BasePlatformAdapter): return path, None async def _handle_media(self, request) -> Any: - """Serve a registered local file for LINE's media URLs. - Defence-in-depth: the resolved path is rechecked against allowed roots - (tempdir, ``/tmp`` → ``/private/tmp`` on macOS, ``HERMES_HOME``). - """ + """Serve a registered local file for LINE's media URLs. Defence-in-depth: the resolved + path is rechecked against allowed roots (tempdir, ``/tmp``→``/private/tmp`` on macOS, HERMES_HOME).""" from aiohttp import web token = request.match_info["token"] entry = self._media_tokens.get(token) @@ -971,7 +871,6 @@ class LineAdapter(BasePlatformAdapter): path, err = self._check_media_file("video", video_path) if err: return err - # LINE requires previewImageUrl: use the supplied preview, else a stdlib 1×1 PNG. if preview_path and Path(preview_path).is_file(): preview_url = self._serve_file(Path(preview_path)) @@ -981,8 +880,7 @@ class LineAdapter(BasePlatformAdapter): tmp.write(_FALLBACK_PNG_PREVIEW) tmp.flush() tmp.close() - preview_token = self._register_media(tmp.name, cleanup=True) - preview_url = self._media_url(preview_token, "preview.png") + preview_url = self._media_url(self._register_media(tmp.name, cleanup=True), "preview.png") except Exception: _unlink_quietly(tmp.name) raise @@ -994,10 +892,8 @@ class LineAdapter(BasePlatformAdapter): async def _send_messages( self, chat_id: str, messages: List[Dict[str, Any]], *, force_push: bool = False, text: bool = False ) -> SendResult: - """Send built message objects, batched at 5/call: reply token first, then push. - ``text`` selects the text-send contract: reply success reports the token - as message_id and a push failure is logged at error level. - """ + """Send built message objects, batched at 5/call: reply token first, then push. ``text`` + selects the text contract: reply success reports the token as message_id; push failure logs at error.""" if not self._client: return SendResult(success=False, error="LINE adapter not connected") if not messages: @@ -1020,9 +916,7 @@ class LineAdapter(BasePlatformAdapter): if text: logger.error("LINE: push send failed: %s", exc) return SendResult(success=False, error=str(exc)) - - # Subsequent batches: always push (reply token is single-use). - while rest: + while rest: # subsequent batches: always push (reply token is single-use) batch = rest[:LINE_MAX_MESSAGES_PER_CALL] rest = rest[LINE_MAX_MESSAGES_PER_CALL:] try: @@ -1040,8 +934,6 @@ def _unlink_quietly(path: str) -> None: pass -# --- Plugin entry-point hooks --- - def check_requirements() -> bool: """Plugin gate: require credentials AND aiohttp at runtime.""" if not (_get_scoped_secret("LINE_CHANNEL_ACCESS_TOKEN") and _get_scoped_secret("LINE_CHANNEL_SECRET")): @@ -1092,8 +984,7 @@ async def _standalone_send( if not token or not chat_id: return {"error": "LINE standalone send: missing token or chat_id"} messages = _text_messages(message or "") or [_text_message("")] - if media_files: - # Tell the recipient media was generated but not delivered. + if media_files: # tell the recipient media was generated but not delivered messages.append(_text_message(f"[{len(media_files)} attachment(s) generated; not deliverable from cron]")) messages = messages[:LINE_MAX_MESSAGES_PER_CALL] try: @@ -1129,6 +1020,7 @@ def interactive_setup() -> None: return if value: _set_env(var, value) + _prompt("LINE_CHANNEL_ACCESS_TOKEN", "Channel access token", secret=True) _prompt("LINE_CHANNEL_SECRET", "Channel secret", secret=True) _prompt("LINE_PUBLIC_URL", "Public HTTPS base URL (optional, e.g. https://my-tunnel.example.com)") @@ -1154,8 +1046,7 @@ def register(ctx) -> None: standalone_sender_fn=_standalone_send, allowed_users_env="LINE_ALLOWED_USERS", allow_all_env="LINE_ALLOW_ALL_USERS", - # LINE per-bubble cap is 5000; smart-chunker uses 4500. - max_message_length=LINE_SAFE_BUBBLE_CHARS, + max_message_length=LINE_SAFE_BUBBLE_CHARS, # per-bubble cap is 5000; smart-chunker uses 4500 emoji="💚", pii_safe=False, allow_update_command=True,