Spaces:
Running
Running
Download channel_msg.py from izuemon/any-env-code: direct link, hf CLI and curl.
- Browser
- Download file 30.1 kB
-
https://huggingface.co/spaces/izuemon/any-env-code/resolve/main/channel_msg.py
- Command line
-
hf download hf://spaces/izuemon/any-env-code/channel_msg.py
-
curl -L -o channel_msg.py https://huggingface.co/spaces/izuemon/any-env-code/resolve/main/channel_msg.py
30.1 kB
| """ | |
| ChannelWorks メッセージ受信ライブラリ | |
| 使い方: | |
| from channel_auth import ChannelWorksClient | |
| from channel_messenger import ChannelMessenger | |
| client = ChannelWorksClient( | |
| room_id="...", | |
| email="you@example.com", | |
| password="your-password", | |
| log_level=ChannelWorksClient.LOG_ERROR, | |
| ) | |
| client.login() | |
| client.start_auto_touch() | |
| def on_push(wrapper): | |
| entity = wrapper.get("entity") or {} | |
| print("受信:", entity.get("plainText")) | |
| messenger = ChannelMessenger( | |
| client=client, | |
| channel_slug="...", # channel_id を自動取得したい場合 | |
| # manager_id="...", channel_id="...", # 直接指定してもよい | |
| group_path="/groups/グループ番号", | |
| log_level=ChannelMessenger.LOG_ERROR, | |
| on_push=on_push, | |
| ) | |
| messenger.start() | |
| try: | |
| messenger.wait() | |
| except KeyboardInterrupt: | |
| pass | |
| finally: | |
| messenger.stop() | |
| client.stop_auto_touch() | |
| log_level: | |
| 0 = 何も出力しない | |
| 1 = エラー・警告のみ (デフォルト) | |
| 2 = デバッグログも出力 (送受信パケット等をすべて出力) | |
| """ | |
| import json | |
| import threading | |
| import time | |
| from datetime import datetime, timezone, timedelta | |
| from typing import Any, Callable, Dict, Optional, Tuple | |
| import requests | |
| import websocket | |
| from channel_auth import ChannelWorksClient | |
| # --------------------------------------------------------------------------- | |
| # 定数 | |
| # --------------------------------------------------------------------------- | |
| DEFAULT_GROUP_PATH = "/groups/574628" | |
| CHANNEL_API_URL_TEMPLATE = "https://api.channel.works/desk/channels/{room_id}" | |
| CHANNEL_LOOKUP_URL_TEMPLATE = "https://api.channel.works/desk/channels/{slug}" | |
| WS_URL = ( | |
| "wss://desk-ws.channel.io/socket.io/" | |
| "?platform=web&EIO=4&transport=websocket" | |
| ) | |
| ACCOUNT_WS_URL = ( | |
| "wss://account-ws.channel.io/socket.io/" | |
| "?platform=web&EIO=4&transport=websocket" | |
| ) | |
| JST = timezone(timedelta(hours=9)) | |
| # --------------------------------------------------------------------------- | |
| # モジュールレベルユーティリティ(ログに依存しない純粋関数) | |
| # --------------------------------------------------------------------------- | |
| def _timestamp_to_text(value: Any) -> str: | |
| if not isinstance(value, (int, float)): | |
| return str(value) | |
| if value < 100_000_000_000: | |
| return str(value) | |
| dt_utc = datetime.fromtimestamp(value / 1000, tz=timezone.utc) | |
| dt_jst = dt_utc.astimezone(JST) | |
| return ( | |
| f"{dt_jst:%Y-%m-%d %H:%M:%S.%f} JST " | |
| f"(UTC {dt_utc:%Y-%m-%d %H:%M:%S.%f})" | |
| ) | |
| def _format_json_with_timestamps(value: Any, key: Optional[str] = None) -> Any: | |
| if isinstance(value, dict): | |
| return {k: _format_json_with_timestamps(v, k) for k, v in value.items()} | |
| if isinstance(value, list): | |
| return [_format_json_with_timestamps(v) for v in value] | |
| if key and key.lower().endswith("at") and isinstance(value, (int, float)): | |
| return {"epoch_ms": value, "formatted": _timestamp_to_text(value)} | |
| return value | |
| def _get_nested(data: Any, *path: str) -> Any: | |
| current = data | |
| for key in path: | |
| if not isinstance(current, dict) or key not in current: | |
| return None | |
| current = current[key] | |
| return current | |
| def _find_first_key(data: Any, key: str) -> Any: | |
| if isinstance(data, dict): | |
| if key in data: | |
| return data[key] | |
| for value in data.values(): | |
| result = _find_first_key(value, key) | |
| if result is not None: | |
| return result | |
| elif isinstance(data, list): | |
| for value in data: | |
| result = _find_first_key(value, key) | |
| if result is not None: | |
| return result | |
| return None | |
| def _parse_socketio_event(packet: str) -> Tuple[Optional[str], Optional[list]]: | |
| if not packet.startswith("42"): | |
| return None, None | |
| payload_text = packet[2:] | |
| if payload_text.startswith("/"): | |
| try: | |
| namespace, payload_text = payload_text.split(",", 1) | |
| except ValueError: | |
| return None, None | |
| else: | |
| namespace = "/" | |
| try: | |
| payload = json.loads(payload_text) | |
| except json.JSONDecodeError: | |
| return namespace, None | |
| if not isinstance(payload, list): | |
| return namespace, None | |
| return namespace, payload | |
| # --------------------------------------------------------------------------- | |
| # WebSocket ラッパ(送受信をデバッグログするだけの薄いラッパ) | |
| # --------------------------------------------------------------------------- | |
| class _LoggingWS: | |
| def __init__(self, ws: websocket.WebSocket, label: str, messenger: "ChannelMessenger") -> None: | |
| self._ws = ws | |
| self._label = label | |
| self._messenger = messenger | |
| self._send_count = 0 | |
| self._recv_count = 0 | |
| def send(self, payload: Any) -> None: | |
| self._send_count += 1 | |
| text = ( | |
| payload.decode("utf-8", errors="replace") | |
| if isinstance(payload, (bytes, bytearray)) | |
| else payload | |
| ) | |
| self._messenger._log_block_debug( | |
| f"{self._label}-SEND", | |
| f"#{self._send_count} ({len(text)} bytes)", | |
| text, | |
| ) | |
| self._ws.send(payload) | |
| def recv(self) -> Optional[str]: | |
| packet = self._ws.recv() | |
| self._recv_count += 1 | |
| if packet is None: | |
| self._messenger._log_debug( | |
| f"{self._label}-RECV", | |
| f"#{self._recv_count} (None / 切断)", | |
| ) | |
| return None | |
| text = ( | |
| packet.decode("utf-8", errors="replace") | |
| if isinstance(packet, (bytes, bytearray)) | |
| else packet | |
| ) | |
| self._messenger._log_block_debug( | |
| f"{self._label}-RECV", | |
| f"#{self._recv_count} ({len(text)} chars)", | |
| text, | |
| ) | |
| return text | |
| def settimeout(self, t: float) -> None: | |
| self._ws.settimeout(t) | |
| def close(self) -> None: | |
| try: | |
| self._ws.close() | |
| finally: | |
| self._messenger._log_debug(self._label, "WebSocket をクローズしました。") | |
| def __getattr__(self, name: str) -> Any: | |
| return getattr(self._ws, name) | |
| # --------------------------------------------------------------------------- | |
| # 本体 | |
| # --------------------------------------------------------------------------- | |
| class ChannelMessenger: | |
| # ログレベル | |
| LOG_NONE = 0 | |
| LOG_ERROR = 1 | |
| LOG_DEBUG = 2 | |
| def __init__( | |
| self, | |
| client: ChannelWorksClient, | |
| channel_slug: Optional[str] = None, | |
| manager_id: Optional[str] = None, | |
| channel_id: Optional[str] = None, | |
| group_path: str = DEFAULT_GROUP_PATH, | |
| log_level: int = LOG_ERROR, | |
| heartbeat_interval: float = 30.0, | |
| heartbeat_response_timeout: float = 5.0, | |
| token_check_interval: float = 1.0, | |
| request_timeout: float = 30.0, | |
| channel_recv_timeout: float = 3.0, | |
| account_recv_timeout: float = 1.0, | |
| on_push: Optional[Callable[[Dict[str, Any]], None]] = None, | |
| ) -> None: | |
| if client is None: | |
| raise ValueError("client (ChannelWorksClient) は必須です。") | |
| self.client = client | |
| self.channel_slug = channel_slug | |
| self.manager_id = manager_id | |
| self.channel_id = channel_id | |
| self.group_path = group_path | |
| self.log_level = log_level | |
| self.heartbeat_interval = heartbeat_interval | |
| self.heartbeat_response_timeout = heartbeat_response_timeout | |
| self.token_check_interval = token_check_interval | |
| self.request_timeout = request_timeout | |
| self.channel_recv_timeout = channel_recv_timeout | |
| self.account_recv_timeout = account_recv_timeout | |
| self.on_push = on_push | |
| self._lock = threading.RLock() | |
| self._stop_event = threading.Event() | |
| self._channel_thread: Optional[threading.Thread] = None | |
| self._account_thread: Optional[threading.Thread] = None | |
| # x-account の最新値。channel/account 両ワーカーがこの dict を参照する。 | |
| self._token_holder: Dict[str, Optional[str]] = {"token": None} | |
| # ------------------------------------------------------------------ | |
| # ログ | |
| # ------------------------------------------------------------------ | |
| def _now_str() -> str: | |
| return datetime.now(JST).strftime("%Y-%m-%d %H:%M:%S.%f")[:-3] | |
| def _log_debug(self, tag: str, msg: str) -> None: | |
| if self.log_level >= self.LOG_DEBUG: | |
| print(f"[{self._now_str()}] [{tag}] {msg}", flush=True) | |
| def _log_error(self, tag: str, msg: str) -> None: | |
| if self.log_level >= self.LOG_ERROR: | |
| print(f"[{self._now_str()}] [{tag}] {msg}", flush=True) | |
| def _log_block_debug(self, tag: str, title: str, body: str) -> None: | |
| if self.log_level < self.LOG_DEBUG: | |
| return | |
| ts = self._now_str() | |
| print(f"\n[{ts}] [{tag}] ===== {title} =====", flush=True) | |
| print(body, flush=True) | |
| print(f"[{ts}] [{tag}] ===== /{title} =====", flush=True) | |
| def set_log_level(self, level: int) -> None: | |
| self.log_level = level | |
| # ------------------------------------------------------------------ | |
| # 公開 API | |
| # ------------------------------------------------------------------ | |
| def start(self) -> None: | |
| """channel / account の各 WebSocket をバックグラウンドで起動する。""" | |
| with self._lock: | |
| if self._is_running(): | |
| self._log_debug("MAIN", "ChannelMessenger は既に起動中です。") | |
| return | |
| tokens = self.client.get_tokens() or {} | |
| required = ("ch-session-1", "ch-veil-id", "x-account", "x-account-refresh") | |
| missing = [k for k in required if not tokens.get(k)] | |
| if missing: | |
| raise RuntimeError( | |
| "client から必要なトークンが取得できません。先に login() を実行してください: " | |
| + ", ".join(missing) | |
| ) | |
| if self.manager_id is None or self.channel_id is None: | |
| self._resolve_ids(tokens) | |
| self._token_holder["token"] = tokens["x-account"] | |
| self._stop_event.clear() | |
| self._channel_thread = threading.Thread( | |
| target=self._channel_worker, | |
| args=(tokens["x-account"], self.channel_id, self.manager_id), | |
| name="channel-ws", | |
| daemon=True, | |
| ) | |
| self._account_thread = threading.Thread( | |
| target=self._account_worker, | |
| name="account-ws", | |
| daemon=True, | |
| ) | |
| self._channel_thread.start() | |
| self._account_thread.start() | |
| self._log_debug( | |
| "MAIN", | |
| f"起動: manager_id={self.manager_id} channel_id={self.channel_id}", | |
| ) | |
| def wait(self, poll_interval: float = 0.5) -> None: | |
| """両ワーカーが終了するまでブロックする。""" | |
| try: | |
| while self._is_running(): | |
| time.sleep(poll_interval) | |
| except KeyboardInterrupt: | |
| self._log_debug("MAIN", "KeyboardInterrupt を受信") | |
| raise | |
| def stop(self, join_timeout: float = 5.0) -> None: | |
| """ワーカーを停止する。""" | |
| self._stop_event.set() | |
| for thread in (self._channel_thread, self._account_thread): | |
| if thread and thread.is_alive(): | |
| thread.join(timeout=join_timeout) | |
| self._channel_thread = None | |
| self._account_thread = None | |
| self._log_debug("MAIN", "ChannelMessenger を停止しました。") | |
| def is_running(self) -> bool: | |
| return self._is_running() | |
| def close(self) -> None: | |
| self.stop() | |
| def __enter__(self) -> "ChannelMessenger": | |
| return self | |
| def __exit__(self, exc_type, exc, tb) -> None: | |
| self.close() | |
| # ------------------------------------------------------------------ | |
| # 内部: 状態 | |
| # ------------------------------------------------------------------ | |
| def _is_running(self) -> bool: | |
| return any( | |
| t is not None and t.is_alive() | |
| for t in (self._channel_thread, self._account_thread) | |
| ) | |
| # ------------------------------------------------------------------ | |
| # 内部: HTTP / ID 解決 | |
| # ------------------------------------------------------------------ | |
| def _resolve_ids(self, tokens: Dict[str, str]) -> None: | |
| room_id = getattr(self.client, "room_id", None) | |
| if not room_id: | |
| raise RuntimeError("client.room_id が取得できません。") | |
| session = requests.Session() | |
| session.cookies.update( | |
| { | |
| "ch-session-1": tokens.get("ch-session-1") or "", | |
| "ch-veil-id": tokens.get("ch-veil-id") or "", | |
| "x-account": tokens.get("x-account") or "", | |
| "x-account-refresh": tokens.get("x-account-refresh") or "", | |
| } | |
| ) | |
| room_url = CHANNEL_API_URL_TEMPLATE.format(room_id=room_id) | |
| if self.manager_id is None: | |
| room_json = self._http_get_json(session, room_url, tokens) | |
| manager_id = _get_nested(room_json, "manager", "id") | |
| if manager_id is None: | |
| manager = _find_first_key(room_json, "manager") | |
| if isinstance(manager, dict): | |
| manager_id = manager.get("id") | |
| if manager_id is None: | |
| raise RuntimeError( | |
| f"{room_url} の JSON から manager.id を取得できませんでした。" | |
| ) | |
| self.manager_id = str(manager_id) | |
| if self.channel_id is None: | |
| if not self.channel_slug: | |
| raise RuntimeError( | |
| "channel_id を自動解決するには channel_slug を指定してください。" | |
| ) | |
| lookup_url = CHANNEL_LOOKUP_URL_TEMPLATE.format(slug=self.channel_slug) | |
| channel_json = self._http_get_json(session, lookup_url, tokens) | |
| channel_id = _get_nested(channel_json, "channel", "id") | |
| if channel_id is None: | |
| channel_id = _find_first_key(channel_json, "id") | |
| if channel_id is None: | |
| raise RuntimeError( | |
| f"{lookup_url} の JSON から channel.id を取得できませんでした。" | |
| ) | |
| self.channel_id = str(channel_id) | |
| def _http_get_json( | |
| self, session: requests.Session, url: str, tokens: Dict[str, str] | |
| ) -> Any: | |
| self._log_debug("HTTP", f"GET {url}") | |
| response = session.get( | |
| url, | |
| headers={ | |
| "Accept": "application/json", | |
| "Accept-Language": "ja", | |
| "Referer": "https://channel.works/", | |
| "Origin": "https://channel.works", | |
| "User-Agent": ( | |
| "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " | |
| "AppleWebKit/537.36 (KHTML, like Gecko) " | |
| "Chrome/152.0.0.0 Safari/537.36" | |
| ), | |
| "x-account": tokens["x-account"], | |
| }, | |
| timeout=self.request_timeout, | |
| ) | |
| self._log_debug("HTTP", f"-> status={response.status_code} ({url})") | |
| response.raise_for_status() | |
| try: | |
| data = response.json() | |
| except ValueError as exc: | |
| raise RuntimeError( | |
| f"JSON を取得できませんでした: {url}\n" | |
| f"status={response.status_code}\n" | |
| f"body={response.text[:500]}" | |
| ) from exc | |
| self._log_block_debug( | |
| "HTTP-BODY", | |
| url, | |
| json.dumps(data, ensure_ascii=False), | |
| ) | |
| return data | |
| # ------------------------------------------------------------------ | |
| # 内部: WebSocket 接続 | |
| # ------------------------------------------------------------------ | |
| def _open_channel_ws(self, session_token: str, channel_id: str) -> _LoggingWS: | |
| raw = websocket.create_connection( | |
| WS_URL, | |
| timeout=self.request_timeout, | |
| origin="https://desk.channel.io", | |
| host="desk-ws.channel.io", | |
| ) | |
| raw.settimeout(self.channel_recv_timeout) | |
| self._log_debug("WS", "Channel WebSocket 接続完了") | |
| ws = _LoggingWS(raw, label="CHANNEL-WS", messenger=self) | |
| deadline = time.monotonic() + self.request_timeout | |
| connect_sent = False | |
| while time.monotonic() < deadline: | |
| try: | |
| packet = ws.recv() | |
| except websocket.WebSocketTimeoutException: | |
| continue | |
| if packet is None: | |
| raise RuntimeError("Channel WebSocket が切断されました。") | |
| if packet.startswith("0"): | |
| self._log_debug("SOCKET", "Channel Engine.IO open 受信") | |
| self._send_channel_connect(ws, session_token, channel_id) | |
| connect_sent = True | |
| continue | |
| if packet == "2": | |
| self._send_pong(ws) | |
| continue | |
| if connect_sent and packet.startswith("40/desk/channel,"): | |
| self._log_debug("SOCKET", "Channel connect 応答 (sid) 受信") | |
| self._send_join_group(ws, self.group_path) | |
| return ws | |
| self._log_debug("SOCKET", f"Channel 接続待機中の受信: {packet!r}") | |
| raise TimeoutError("Channel connect 応答を受信できませんでした。") | |
| def _open_account_ws(self, session_token: str) -> _LoggingWS: | |
| raw = websocket.create_connection( | |
| ACCOUNT_WS_URL, | |
| timeout=self.request_timeout, | |
| origin="https://desk.channel.io", | |
| host="account-ws.channel.io", | |
| ) | |
| raw.settimeout(self.account_recv_timeout) | |
| self._log_debug("WS", "Account WebSocket 接続完了") | |
| ws = _LoggingWS(raw, label="ACCOUNT-WS", messenger=self) | |
| deadline = time.monotonic() + self.request_timeout | |
| while time.monotonic() < deadline: | |
| try: | |
| packet = ws.recv() | |
| except websocket.WebSocketTimeoutException: | |
| continue | |
| if packet is None: | |
| raise RuntimeError("Account WebSocket が切断されました。") | |
| if packet.startswith("0"): | |
| self._log_debug("SOCKET", "Account Engine.IO open 受信") | |
| self._send_account_connect(ws, session_token) | |
| return ws | |
| if packet == "2": | |
| self._send_pong(ws) | |
| continue | |
| self._log_debug("SOCKET", f"Account 接続待機中の受信: {packet!r}") | |
| raise TimeoutError("Account Engine.IO open パケットを受信できませんでした。") | |
| # ------------------------------------------------------------------ | |
| # 内部: Socket.IO パケット送信 | |
| # ------------------------------------------------------------------ | |
| def _send_join_group(self, ws: _LoggingWS, group_path: str) -> None: | |
| message = '42/desk/channel,' + json.dumps( | |
| ["join", group_path], | |
| ensure_ascii=False, | |
| separators=(",", ":"), | |
| ) | |
| self._log_debug("SOCKET", f"join 送信: {group_path}") | |
| ws.send(message) | |
| def _send_channel_connect( | |
| self, ws: _LoggingWS, session_token: str, channel_id: str | |
| ) -> None: | |
| message = '40/desk/channel,' + json.dumps( | |
| {"jwt": session_token, "channelId": channel_id}, | |
| ensure_ascii=False, | |
| separators=(",", ":"), | |
| ) | |
| self._log_debug("SOCKET", f"channel connect 送信: channelId={channel_id}") | |
| ws.send(message) | |
| def _send_account_connect(self, ws: _LoggingWS, session_token: str) -> None: | |
| message = '40/desk/account,' + json.dumps( | |
| {"jwt": session_token}, | |
| ensure_ascii=False, | |
| separators=(",", ":"), | |
| ) | |
| self._log_debug("SOCKET", "account connect 送信") | |
| ws.send(message) | |
| def _send_refresh(self, ws: _LoggingWS, session_token: str) -> None: | |
| message = '42/desk/account,0' + json.dumps( | |
| ["refresh", {"jwt": session_token}], | |
| ensure_ascii=False, | |
| separators=(",", ":"), | |
| ) | |
| self._log_debug("SOCKET", "refresh 送信(account WS)") | |
| ws.send(message) | |
| def _send_channel_refresh( | |
| self, ws: _LoggingWS, session_token: str, channel_id: str | |
| ) -> None: | |
| message = '42/desk/channel,3' + json.dumps( | |
| ["refresh", {"jwt": session_token, "channelId": channel_id}], | |
| ensure_ascii=False, | |
| separators=(",", ":"), | |
| ) | |
| self._log_debug( | |
| "SOCKET", f"refresh 送信(channel WS, channelId={channel_id})" | |
| ) | |
| ws.send(message) | |
| def _send_heartbeat(self, ws: _LoggingWS, manager_id: str) -> None: | |
| message = '42/desk/channel,' + json.dumps( | |
| ["heartbeat", manager_id], | |
| ensure_ascii=False, | |
| separators=(",", ":"), | |
| ) | |
| self._log_debug("SOCKET", f"heartbeat 送信: manager_id={manager_id}") | |
| ws.send(message) | |
| def _send_pong(self, ws: _LoggingWS) -> None: | |
| self._log_debug("SOCKET", "Engine.IO ping(2) 受信 -> pong(3) 送信") | |
| ws.send("3") | |
| # ------------------------------------------------------------------ | |
| # 内部: push 処理 | |
| # ------------------------------------------------------------------ | |
| def _handle_push(self, payload: list) -> None: | |
| if len(payload) < 2: | |
| return | |
| wrapper = payload[1] | |
| if not isinstance(wrapper, dict): | |
| return | |
| if self.log_level >= self.LOG_DEBUG: | |
| self._log_block_debug( | |
| "PUSH", | |
| f"type={wrapper.get('type')} event={wrapper.get('event')}", | |
| json.dumps( | |
| _format_json_with_timestamps(wrapper), | |
| ensure_ascii=False, | |
| ), | |
| ) | |
| if self.on_push is not None: | |
| try: | |
| self.on_push(wrapper) | |
| except Exception as exc: | |
| self._log_error("PUSH", f"on_push コールバックで例外: {exc!r}") | |
| # ------------------------------------------------------------------ | |
| # 内部: channel ワーカー | |
| # ------------------------------------------------------------------ | |
| def _channel_worker( | |
| self, session_token: str, channel_id: str, manager_id: str | |
| ) -> None: | |
| try: | |
| ws = self._open_channel_ws(session_token, channel_id) | |
| except Exception as exc: | |
| self._log_error("CHANNEL", f"接続失敗: {exc!r}") | |
| return | |
| current_token = session_token | |
| last_token_check = time.monotonic() | |
| ready_received = False | |
| heartbeat_phase = "waiting" # "waiting" / "first_sent" / "second_sent" | |
| phase_started_at = 0.0 | |
| try: | |
| while not self._stop_event.is_set(): | |
| now = time.monotonic() | |
| if ready_received: | |
| if heartbeat_phase == "waiting": | |
| if now - phase_started_at >= self.heartbeat_interval: | |
| self._send_heartbeat(ws, manager_id) | |
| heartbeat_phase = "first_sent" | |
| phase_started_at = now | |
| elif heartbeat_phase == "first_sent": | |
| if now - phase_started_at >= self.heartbeat_response_timeout: | |
| self._log_debug( | |
| "SOCKET", | |
| f"1回目 heartbeat 応答なし ({self.heartbeat_response_timeout}s) → 再送", | |
| ) | |
| self._send_heartbeat(ws, manager_id) | |
| phase_started_at = now | |
| elif heartbeat_phase == "second_sent": | |
| if now - phase_started_at >= self.heartbeat_response_timeout: | |
| self._log_debug( | |
| "SOCKET", | |
| f"2回目 heartbeat 応答なし ({self.heartbeat_response_timeout}s) → 再送", | |
| ) | |
| self._send_heartbeat(ws, manager_id) | |
| phase_started_at = now | |
| if now - last_token_check >= self.token_check_interval: | |
| new_token = self._token_holder.get("token") | |
| if new_token and new_token != current_token: | |
| self._log_debug( | |
| "AUTH", "x-account 変更検出 → Channel refresh 送信" | |
| ) | |
| self._send_channel_refresh(ws, new_token, channel_id) | |
| current_token = new_token | |
| last_token_check = now | |
| try: | |
| packet = ws.recv() | |
| except websocket.WebSocketTimeoutException: | |
| continue | |
| if packet is None: | |
| raise RuntimeError("Channel WebSocket が切断されました。") | |
| if packet == "2": | |
| self._send_pong(ws) | |
| continue | |
| if packet == "3": | |
| self._log_debug("SOCKET", "Channel Engine.IO pong(3) 受信") | |
| continue | |
| namespace, payload = _parse_socketio_event(packet) | |
| if namespace is None or payload is None: | |
| continue | |
| if namespace == "/desk/channel" and payload: | |
| event_name = payload[0] | |
| self._log_debug("SOCKET", f"event={event_name}") | |
| if event_name == "ready": | |
| ready_received = True | |
| heartbeat_phase = "waiting" | |
| phase_started_at = time.monotonic() | |
| elif event_name == "heartbeat": | |
| if heartbeat_phase == "first_sent": | |
| self._log_debug( | |
| "SOCKET", "1回目の応答受信 → 2回目を送信" | |
| ) | |
| self._send_heartbeat(ws, manager_id) | |
| heartbeat_phase = "second_sent" | |
| phase_started_at = time.monotonic() | |
| elif heartbeat_phase == "second_sent": | |
| self._log_debug( | |
| "SOCKET", "2回目の応答受信 → 次サイクル待機" | |
| ) | |
| heartbeat_phase = "waiting" | |
| phase_started_at = time.monotonic() | |
| elif event_name == "create": | |
| self._log_block_debug("CREATE", "raw packet", packet) | |
| elif event_name == "push": | |
| self._handle_push(payload) | |
| elif event_name == "expired": | |
| self._log_error( | |
| "CHANNEL", "expired を受信 → 接続を終了します。" | |
| ) | |
| break | |
| except Exception as exc: | |
| self._log_error("CHANNEL", f"エラーで終了: {exc!r}") | |
| finally: | |
| ws.close() | |
| # ------------------------------------------------------------------ | |
| # 内部: account ワーカー | |
| # ------------------------------------------------------------------ | |
| def _account_worker(self) -> None: | |
| try: | |
| ws = self._open_account_ws(self._token_holder["token"] or "") | |
| except Exception as exc: | |
| self._log_error("ACCOUNT", f"接続失敗: {exc!r}") | |
| return | |
| last_token_check = time.monotonic() | |
| try: | |
| while not self._stop_event.is_set(): | |
| now = time.monotonic() | |
| if now - last_token_check >= self.token_check_interval: | |
| current_tokens = self.client.get_tokens() or {} | |
| current_session_token = current_tokens.get("x-account") | |
| if ( | |
| current_session_token | |
| and current_session_token != self._token_holder["token"] | |
| ): | |
| self._log_debug( | |
| "AUTH", "x-account 変更検出 → refresh 送信(account WS)" | |
| ) | |
| self._send_refresh(ws, current_session_token) | |
| self._token_holder["token"] = current_session_token | |
| last_token_check = now | |
| try: | |
| packet = ws.recv() | |
| except websocket.WebSocketTimeoutException: | |
| continue | |
| if packet is None: | |
| raise RuntimeError("Account WebSocket が切断されました。") | |
| if packet == "2": | |
| self._send_pong(ws) | |
| continue | |
| if packet == "3": | |
| self._log_debug("SOCKET", "Account Engine.IO pong(3) 受信") | |
| continue | |
| namespace, payload = _parse_socketio_event(packet) | |
| if namespace is None or payload is None: | |
| continue | |
| self._log_debug( | |
| "SOCKET", f"Account event: ns={namespace}, payload={payload!r}" | |
| ) | |
| except Exception as exc: | |
| self._log_error("ACCOUNT", f"エラーで終了: {exc!r}") | |
| finally: | |
| ws.close() | |