any-env-code / channel_msg.py
izuemon's picture
Rename get_msg.py to channel_msg.py
35956c7 verified
Raw History Blame Contribute Delete
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}
# ------------------------------------------------------------------
# ログ
# ------------------------------------------------------------------
@staticmethod
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()