Spaces:
Running
Running
Update get_msg.py
Browse files- get_msg.py +239 -100
get_msg.py
CHANGED
|
@@ -20,11 +20,18 @@ WS_URL = (
|
|
| 20 |
"wss://desk-ws.channel.io/socket.io/"
|
| 21 |
"?platform=web&EIO=4&transport=websocket"
|
| 22 |
)
|
|
|
|
|
|
|
|
|
|
|
|
|
| 23 |
|
| 24 |
HEARTBEAT_INTERVAL = 30.0
|
| 25 |
TOKEN_CHECK_INTERVAL = 1.0
|
| 26 |
REQUEST_TIMEOUT = 30.0
|
| 27 |
|
|
|
|
|
|
|
|
|
|
| 28 |
JST = timezone(timedelta(hours=9))
|
| 29 |
|
| 30 |
|
|
@@ -49,12 +56,12 @@ def log_block(tag: str, title: str, body: str) -> None:
|
|
| 49 |
class LoggingWS:
|
| 50 |
"""websocket.WebSocket をラップして送受信を必ずログ出力する。"""
|
| 51 |
|
| 52 |
-
def __init__(self, ws: websocket.WebSocket) -> None:
|
| 53 |
self._ws = ws
|
|
|
|
| 54 |
self._send_count = 0
|
| 55 |
self._recv_count = 0
|
| 56 |
|
| 57 |
-
# --- 送信 ---
|
| 58 |
def send(self, payload: str | bytes) -> None:
|
| 59 |
self._send_count += 1
|
| 60 |
text = (
|
|
@@ -63,18 +70,17 @@ class LoggingWS:
|
|
| 63 |
else payload
|
| 64 |
)
|
| 65 |
log_block(
|
| 66 |
-
"
|
| 67 |
f"#{self._send_count} ({len(text)} bytes)",
|
| 68 |
text,
|
| 69 |
)
|
| 70 |
self._ws.send(payload)
|
| 71 |
|
| 72 |
-
# --- 受信 ---
|
| 73 |
def recv(self) -> str | bytes | None:
|
| 74 |
packet = self._ws.recv()
|
| 75 |
self._recv_count += 1
|
| 76 |
if packet is None:
|
| 77 |
-
log("
|
| 78 |
return None
|
| 79 |
text = (
|
| 80 |
packet.decode("utf-8", errors="replace")
|
|
@@ -82,7 +88,7 @@ class LoggingWS:
|
|
| 82 |
else packet
|
| 83 |
)
|
| 84 |
log_block(
|
| 85 |
-
"
|
| 86 |
f"#{self._recv_count} ({len(text)} chars)",
|
| 87 |
text,
|
| 88 |
)
|
|
@@ -92,10 +98,9 @@ class LoggingWS:
|
|
| 92 |
try:
|
| 93 |
self._ws.close()
|
| 94 |
finally:
|
| 95 |
-
log(
|
| 96 |
|
| 97 |
def __getattr__(self, name: str) -> Any:
|
| 98 |
-
# 未定義の属性は元の WebSocket に委譲
|
| 99 |
return getattr(self._ws, name)
|
| 100 |
|
| 101 |
|
|
@@ -158,7 +163,6 @@ def get_nested(data: Any, *path: str) -> Any:
|
|
| 158 |
|
| 159 |
|
| 160 |
def find_first_key(data: Any, key: str) -> Any:
|
| 161 |
-
"""レスポンス形式が多少変わっても id を拾えるように再帰検索する。"""
|
| 162 |
if isinstance(data, dict):
|
| 163 |
if key in data:
|
| 164 |
return data[key]
|
|
@@ -219,10 +223,8 @@ def fetch_ids(tokens: dict[str, str]) -> tuple[str, str]:
|
|
| 219 |
def timestamp_to_text(value: Any) -> str:
|
| 220 |
if not isinstance(value, (int, float)):
|
| 221 |
return str(value)
|
| 222 |
-
|
| 223 |
if value < 100_000_000_000:
|
| 224 |
return str(value)
|
| 225 |
-
|
| 226 |
dt_utc = datetime.fromtimestamp(value / 1000, tz=timezone.utc)
|
| 227 |
dt_jst = dt_utc.astimezone(JST)
|
| 228 |
return (
|
|
@@ -233,17 +235,10 @@ def timestamp_to_text(value: Any) -> str:
|
|
| 233 |
|
| 234 |
def format_json_with_timestamps(value: Any, key: str | None = None) -> Any:
|
| 235 |
if isinstance(value, dict):
|
| 236 |
-
return {
|
| 237 |
-
k: format_json_with_timestamps(v, k)
|
| 238 |
-
for k, v in value.items()
|
| 239 |
-
}
|
| 240 |
if isinstance(value, list):
|
| 241 |
return [format_json_with_timestamps(v) for v in value]
|
| 242 |
-
if (
|
| 243 |
-
key
|
| 244 |
-
and key.lower().endswith("at")
|
| 245 |
-
and isinstance(value, (int, float))
|
| 246 |
-
):
|
| 247 |
return {
|
| 248 |
"epoch_ms": value,
|
| 249 |
"formatted": timestamp_to_text(value),
|
|
@@ -279,7 +274,6 @@ def print_create_message(payload: list[Any]) -> None:
|
|
| 279 |
return
|
| 280 |
|
| 281 |
entity = entity_wrapper.get("entity", {})
|
| 282 |
-
refers = entity_wrapper.get("refers", {})
|
| 283 |
|
| 284 |
if isinstance(entity, dict):
|
| 285 |
log("CREATE", f"chatKey : {entity.get('chatKey')}")
|
|
@@ -327,7 +321,6 @@ def print_create_message(payload: list[Any]) -> None:
|
|
| 327 |
# Socket.IO / Engine.IO
|
| 328 |
# ---------------------------------------------------------------------------
|
| 329 |
def parse_socketio_event(packet: str) -> tuple[str | None, list[Any] | None]:
|
| 330 |
-
"""Socket.IO の 42... パケットを namespace と JSON payload に分解する。"""
|
| 331 |
log("PARSE", f"packet={packet!r}")
|
| 332 |
if not packet.startswith("42"):
|
| 333 |
log("PARSE", "-> Socket.IO イベント(42)ではありません。")
|
|
@@ -366,10 +359,7 @@ def send_channel_connect(
|
|
| 366 |
message = (
|
| 367 |
'40/desk/channel,'
|
| 368 |
+ json.dumps(
|
| 369 |
-
{
|
| 370 |
-
"jwt": session_token,
|
| 371 |
-
"channelId": channel_id,
|
| 372 |
-
},
|
| 373 |
ensure_ascii=False,
|
| 374 |
separators=(",", ":"),
|
| 375 |
)
|
|
@@ -378,6 +368,22 @@ def send_channel_connect(
|
|
| 378 |
ws.send(message)
|
| 379 |
|
| 380 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 381 |
def send_refresh(
|
| 382 |
ws: LoggingWS,
|
| 383 |
session_token: str,
|
|
@@ -390,7 +396,7 @@ def send_refresh(
|
|
| 390 |
separators=(",", ":"),
|
| 391 |
)
|
| 392 |
)
|
| 393 |
-
log("SOCKET", "refresh を送信します。")
|
| 394 |
ws.send(message)
|
| 395 |
|
| 396 |
|
|
@@ -408,7 +414,6 @@ def send_heartbeat(ws: LoggingWS, manager_id: str) -> None:
|
|
| 408 |
|
| 409 |
|
| 410 |
def send_pong(ws: LoggingWS) -> None:
|
| 411 |
-
"""Engine.IO ping(2) に対する pong(3)。"""
|
| 412 |
log("SOCKET", "Engine.IO ping(2) を受信 -> pong(3) を返します。")
|
| 413 |
ws.send("3")
|
| 414 |
|
|
@@ -418,7 +423,7 @@ def wait_engineio_open(
|
|
| 418 |
session_token: str,
|
| 419 |
channel_id: str,
|
| 420 |
) -> None:
|
| 421 |
-
log("SOCKET", "Engine.IO open 待機開始")
|
| 422 |
deadline = time.monotonic() + REQUEST_TIMEOUT
|
| 423 |
|
| 424 |
while time.monotonic() < deadline:
|
|
@@ -429,10 +434,10 @@ def wait_engineio_open(
|
|
| 429 |
continue
|
| 430 |
|
| 431 |
if packet is None:
|
| 432 |
-
raise RuntimeError("WebSocketが切断されました。")
|
| 433 |
|
| 434 |
if packet.startswith("0"):
|
| 435 |
-
log("SOCKET", "Engine.IO open パケットを受信しました。")
|
| 436 |
send_channel_connect(ws, session_token, channel_id)
|
| 437 |
log("SOCKET", "desk/channel に接続しました。")
|
| 438 |
return
|
|
@@ -443,7 +448,144 @@ def wait_engineio_open(
|
|
| 443 |
|
| 444 |
log("SOCKET", f"接続待機中に受信: {packet!r}")
|
| 445 |
|
| 446 |
-
raise TimeoutError("Engine.IO open パケットを受信できませんでした。")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 447 |
|
| 448 |
|
| 449 |
# ---------------------------------------------------------------------------
|
|
@@ -459,6 +601,7 @@ def run() -> None:
|
|
| 459 |
log("MAIN", f"CHANNEL_API_URL={CHANNEL_API_URL}")
|
| 460 |
log("MAIN", f"CHANNEL_LOOKUP_URL={CHANNEL_LOOKUP_URL}")
|
| 461 |
log("MAIN", f"WS_URL={WS_URL}")
|
|
|
|
| 462 |
|
| 463 |
client = ChannelWorksClient(
|
| 464 |
room_id=room_id,
|
|
@@ -489,8 +632,7 @@ def run() -> None:
|
|
| 489 |
missing = [key for key in required if not tokens.get(key)]
|
| 490 |
if missing:
|
| 491 |
raise RuntimeError(
|
| 492 |
-
"必要なトークンが取得できませんでした: "
|
| 493 |
-
+ ", ".join(missing)
|
| 494 |
)
|
| 495 |
log("AUTH", f"トークン取得 OK({len(tokens)} 件)")
|
| 496 |
|
|
@@ -500,83 +642,80 @@ def run() -> None:
|
|
| 500 |
log("HTTP", f"channel.id = {channel_id}")
|
| 501 |
|
| 502 |
last_session_token = tokens["x-account"]
|
|
|
|
| 503 |
|
| 504 |
-
|
|
|
|
| 505 |
WS_URL,
|
| 506 |
timeout=3.0,
|
| 507 |
origin="https://desk.channel.io",
|
| 508 |
host="desk-ws.channel.io",
|
| 509 |
)
|
| 510 |
-
log("WS", "WebSocket 接続完了")
|
| 511 |
-
|
| 512 |
-
|
| 513 |
-
try:
|
| 514 |
-
wait_engineio_open(
|
| 515 |
-
ws,
|
| 516 |
-
last_session_token,
|
| 517 |
-
channel_id,
|
| 518 |
-
)
|
| 519 |
-
|
| 520 |
-
last_heartbeat = time.monotonic()
|
| 521 |
-
last_token_check = time.monotonic()
|
| 522 |
-
|
| 523 |
-
log("SOCKET", "待機中...")
|
| 524 |
-
|
| 525 |
-
while True:
|
| 526 |
-
now = time.monotonic()
|
| 527 |
-
|
| 528 |
-
if now - last_heartbeat >= HEARTBEAT_INTERVAL:
|
| 529 |
-
send_heartbeat(ws, manager_id)
|
| 530 |
-
last_heartbeat = now
|
| 531 |
|
| 532 |
-
|
| 533 |
-
|
| 534 |
-
|
| 535 |
-
|
| 536 |
-
|
| 537 |
-
|
| 538 |
-
|
| 539 |
-
|
| 540 |
-
|
| 541 |
-
|
| 542 |
-
|
| 543 |
-
|
| 544 |
-
|
| 545 |
-
|
| 546 |
-
|
| 547 |
-
|
| 548 |
-
|
| 549 |
-
|
| 550 |
-
|
| 551 |
-
|
| 552 |
-
|
| 553 |
-
|
| 554 |
-
|
| 555 |
-
|
| 556 |
-
|
| 557 |
-
|
| 558 |
-
|
| 559 |
-
|
| 560 |
-
|
| 561 |
-
|
| 562 |
-
|
| 563 |
-
|
| 564 |
-
|
| 565 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 566 |
|
| 567 |
-
|
| 568 |
-
if namespace is None or payload is None:
|
| 569 |
-
continue
|
| 570 |
|
| 571 |
-
|
| 572 |
-
|
| 573 |
-
|
|
|
|
|
|
|
|
|
|
| 574 |
|
| 575 |
-
|
| 576 |
-
|
|
|
|
| 577 |
|
| 578 |
-
|
| 579 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 580 |
|
| 581 |
finally:
|
| 582 |
client.stop_auto_touch()
|
|
|
|
| 20 |
"wss://desk-ws.channel.io/socket.io/"
|
| 21 |
"?platform=web&EIO=4&transport=websocket"
|
| 22 |
)
|
| 23 |
+
ACCOUNT_WS_URL = (
|
| 24 |
+
"wss://account-ws.channel.io/socket.io/"
|
| 25 |
+
"?platform=web&EIO=4&transport=websocket"
|
| 26 |
+
)
|
| 27 |
|
| 28 |
HEARTBEAT_INTERVAL = 30.0
|
| 29 |
TOKEN_CHECK_INTERVAL = 1.0
|
| 30 |
REQUEST_TIMEOUT = 30.0
|
| 31 |
|
| 32 |
+
# account 側は短いタイムアウトで recv し、トークン監視を素早く回す
|
| 33 |
+
ACCOUNT_RECV_TIMEOUT = 1.0
|
| 34 |
+
|
| 35 |
JST = timezone(timedelta(hours=9))
|
| 36 |
|
| 37 |
|
|
|
|
| 56 |
class LoggingWS:
|
| 57 |
"""websocket.WebSocket をラップして送受信を必ずログ出力する。"""
|
| 58 |
|
| 59 |
+
def __init__(self, ws: websocket.WebSocket, label: str = "WS") -> None:
|
| 60 |
self._ws = ws
|
| 61 |
+
self._label = label
|
| 62 |
self._send_count = 0
|
| 63 |
self._recv_count = 0
|
| 64 |
|
|
|
|
| 65 |
def send(self, payload: str | bytes) -> None:
|
| 66 |
self._send_count += 1
|
| 67 |
text = (
|
|
|
|
| 70 |
else payload
|
| 71 |
)
|
| 72 |
log_block(
|
| 73 |
+
f"{self._label}-SEND",
|
| 74 |
f"#{self._send_count} ({len(text)} bytes)",
|
| 75 |
text,
|
| 76 |
)
|
| 77 |
self._ws.send(payload)
|
| 78 |
|
|
|
|
| 79 |
def recv(self) -> str | bytes | None:
|
| 80 |
packet = self._ws.recv()
|
| 81 |
self._recv_count += 1
|
| 82 |
if packet is None:
|
| 83 |
+
log(f"{self._label}-RECV", f"#{self._recv_count} (None / 切断)")
|
| 84 |
return None
|
| 85 |
text = (
|
| 86 |
packet.decode("utf-8", errors="replace")
|
|
|
|
| 88 |
else packet
|
| 89 |
)
|
| 90 |
log_block(
|
| 91 |
+
f"{self._label}-RECV",
|
| 92 |
f"#{self._recv_count} ({len(text)} chars)",
|
| 93 |
text,
|
| 94 |
)
|
|
|
|
| 98 |
try:
|
| 99 |
self._ws.close()
|
| 100 |
finally:
|
| 101 |
+
log(self._label, "WebSocket をクローズしました。")
|
| 102 |
|
| 103 |
def __getattr__(self, name: str) -> Any:
|
|
|
|
| 104 |
return getattr(self._ws, name)
|
| 105 |
|
| 106 |
|
|
|
|
| 163 |
|
| 164 |
|
| 165 |
def find_first_key(data: Any, key: str) -> Any:
|
|
|
|
| 166 |
if isinstance(data, dict):
|
| 167 |
if key in data:
|
| 168 |
return data[key]
|
|
|
|
| 223 |
def timestamp_to_text(value: Any) -> str:
|
| 224 |
if not isinstance(value, (int, float)):
|
| 225 |
return str(value)
|
|
|
|
| 226 |
if value < 100_000_000_000:
|
| 227 |
return str(value)
|
|
|
|
| 228 |
dt_utc = datetime.fromtimestamp(value / 1000, tz=timezone.utc)
|
| 229 |
dt_jst = dt_utc.astimezone(JST)
|
| 230 |
return (
|
|
|
|
| 235 |
|
| 236 |
def format_json_with_timestamps(value: Any, key: str | None = None) -> Any:
|
| 237 |
if isinstance(value, dict):
|
| 238 |
+
return {k: format_json_with_timestamps(v, k) for k, v in value.items()}
|
|
|
|
|
|
|
|
|
|
| 239 |
if isinstance(value, list):
|
| 240 |
return [format_json_with_timestamps(v) for v in value]
|
| 241 |
+
if key and key.lower().endswith("at") and isinstance(value, (int, float)):
|
|
|
|
|
|
|
|
|
|
|
|
|
| 242 |
return {
|
| 243 |
"epoch_ms": value,
|
| 244 |
"formatted": timestamp_to_text(value),
|
|
|
|
| 274 |
return
|
| 275 |
|
| 276 |
entity = entity_wrapper.get("entity", {})
|
|
|
|
| 277 |
|
| 278 |
if isinstance(entity, dict):
|
| 279 |
log("CREATE", f"chatKey : {entity.get('chatKey')}")
|
|
|
|
| 321 |
# Socket.IO / Engine.IO
|
| 322 |
# ---------------------------------------------------------------------------
|
| 323 |
def parse_socketio_event(packet: str) -> tuple[str | None, list[Any] | None]:
|
|
|
|
| 324 |
log("PARSE", f"packet={packet!r}")
|
| 325 |
if not packet.startswith("42"):
|
| 326 |
log("PARSE", "-> Socket.IO イベント(42)ではありません。")
|
|
|
|
| 359 |
message = (
|
| 360 |
'40/desk/channel,'
|
| 361 |
+ json.dumps(
|
| 362 |
+
{"jwt": session_token, "channelId": channel_id},
|
|
|
|
|
|
|
|
|
|
| 363 |
ensure_ascii=False,
|
| 364 |
separators=(",", ":"),
|
| 365 |
)
|
|
|
|
| 368 |
ws.send(message)
|
| 369 |
|
| 370 |
|
| 371 |
+
def send_account_connect(
|
| 372 |
+
ws: LoggingWS,
|
| 373 |
+
session_token: str,
|
| 374 |
+
) -> None:
|
| 375 |
+
message = (
|
| 376 |
+
'40/desk/account,'
|
| 377 |
+
+ json.dumps(
|
| 378 |
+
{"jwt": session_token},
|
| 379 |
+
ensure_ascii=False,
|
| 380 |
+
separators=(",", ":"),
|
| 381 |
+
)
|
| 382 |
+
)
|
| 383 |
+
log("SOCKET", "account connect 送信")
|
| 384 |
+
ws.send(message)
|
| 385 |
+
|
| 386 |
+
|
| 387 |
def send_refresh(
|
| 388 |
ws: LoggingWS,
|
| 389 |
session_token: str,
|
|
|
|
| 396 |
separators=(",", ":"),
|
| 397 |
)
|
| 398 |
)
|
| 399 |
+
log("SOCKET", "refresh を送信します(account WS)。")
|
| 400 |
ws.send(message)
|
| 401 |
|
| 402 |
|
|
|
|
| 414 |
|
| 415 |
|
| 416 |
def send_pong(ws: LoggingWS) -> None:
|
|
|
|
| 417 |
log("SOCKET", "Engine.IO ping(2) を受信 -> pong(3) を返します。")
|
| 418 |
ws.send("3")
|
| 419 |
|
|
|
|
| 423 |
session_token: str,
|
| 424 |
channel_id: str,
|
| 425 |
) -> None:
|
| 426 |
+
log("SOCKET", "Channel Engine.IO open 待機開始")
|
| 427 |
deadline = time.monotonic() + REQUEST_TIMEOUT
|
| 428 |
|
| 429 |
while time.monotonic() < deadline:
|
|
|
|
| 434 |
continue
|
| 435 |
|
| 436 |
if packet is None:
|
| 437 |
+
raise RuntimeError("Channel WebSocketが切断されました。")
|
| 438 |
|
| 439 |
if packet.startswith("0"):
|
| 440 |
+
log("SOCKET", "Channel Engine.IO open パケットを受信しました。")
|
| 441 |
send_channel_connect(ws, session_token, channel_id)
|
| 442 |
log("SOCKET", "desk/channel に接続しました。")
|
| 443 |
return
|
|
|
|
| 448 |
|
| 449 |
log("SOCKET", f"接続待機中に受信: {packet!r}")
|
| 450 |
|
| 451 |
+
raise TimeoutError("Channel Engine.IO open パケットを受信できませんでした。")
|
| 452 |
+
|
| 453 |
+
|
| 454 |
+
def wait_account_engineio_open(
|
| 455 |
+
ws: LoggingWS,
|
| 456 |
+
session_token: str,
|
| 457 |
+
) -> None:
|
| 458 |
+
log("SOCKET", "Account Engine.IO open 待機開始")
|
| 459 |
+
deadline = time.monotonic() + REQUEST_TIMEOUT
|
| 460 |
+
|
| 461 |
+
while time.monotonic() < deadline:
|
| 462 |
+
try:
|
| 463 |
+
packet = ws.recv()
|
| 464 |
+
except websocket.WebSocketTimeoutException:
|
| 465 |
+
continue
|
| 466 |
+
|
| 467 |
+
if packet is None:
|
| 468 |
+
raise RuntimeError("Account WebSocketが切断されました。")
|
| 469 |
+
|
| 470 |
+
if packet.startswith("0"):
|
| 471 |
+
log("SOCKET", "Account Engine.IO open パケットを受信しました。")
|
| 472 |
+
send_account_connect(ws, session_token)
|
| 473 |
+
log("SOCKET", "desk/account に接続しました。")
|
| 474 |
+
return
|
| 475 |
+
|
| 476 |
+
if packet == "2":
|
| 477 |
+
send_pong(ws)
|
| 478 |
+
continue
|
| 479 |
+
|
| 480 |
+
log("SOCKET", f"Account 接続待機中に受信: {packet!r}")
|
| 481 |
+
|
| 482 |
+
raise TimeoutError("Account Engine.IO open パケットを受信できませんでした。")
|
| 483 |
+
|
| 484 |
+
|
| 485 |
+
# ---------------------------------------------------------------------------
|
| 486 |
+
# 各 WebSocket の受信ループ
|
| 487 |
+
# ---------------------------------------------------------------------------
|
| 488 |
+
def channel_loop(
|
| 489 |
+
ws: LoggingWS,
|
| 490 |
+
session_token: str,
|
| 491 |
+
channel_id: str,
|
| 492 |
+
manager_id: str,
|
| 493 |
+
stop_event: threading.Event,
|
| 494 |
+
) -> None:
|
| 495 |
+
wait_engineio_open(ws, session_token, channel_id)
|
| 496 |
+
last_heartbeat = time.monotonic()
|
| 497 |
+
|
| 498 |
+
log("SOCKET", "Channel 待機中...")
|
| 499 |
+
|
| 500 |
+
while not stop_event.is_set():
|
| 501 |
+
now = time.monotonic()
|
| 502 |
+
|
| 503 |
+
if now - last_heartbeat >= HEARTBEAT_INTERVAL:
|
| 504 |
+
send_heartbeat(ws, manager_id)
|
| 505 |
+
last_heartbeat = now
|
| 506 |
+
|
| 507 |
+
try:
|
| 508 |
+
packet = ws.recv()
|
| 509 |
+
except websocket.WebSocketTimeoutException:
|
| 510 |
+
continue
|
| 511 |
+
|
| 512 |
+
if packet is None:
|
| 513 |
+
raise RuntimeError("Channel WebSocketが切断されました。")
|
| 514 |
+
|
| 515 |
+
if packet == "2":
|
| 516 |
+
send_pong(ws)
|
| 517 |
+
continue
|
| 518 |
+
|
| 519 |
+
if packet == "3":
|
| 520 |
+
log("SOCKET", "Channel Engine.IO pong(3) を受信(無視)")
|
| 521 |
+
continue
|
| 522 |
+
|
| 523 |
+
namespace, payload = parse_socketio_event(packet)
|
| 524 |
+
if namespace is None or payload is None:
|
| 525 |
+
continue
|
| 526 |
+
|
| 527 |
+
if namespace == "/desk/channel" and payload:
|
| 528 |
+
event_name = payload[0]
|
| 529 |
+
log("SOCKET", f"event={event_name}")
|
| 530 |
+
if event_name == "create":
|
| 531 |
+
print_create_message(payload)
|
| 532 |
+
|
| 533 |
+
|
| 534 |
+
def account_loop(
|
| 535 |
+
ws: LoggingWS,
|
| 536 |
+
token_holder: dict[str, str],
|
| 537 |
+
client: ChannelWorksClient,
|
| 538 |
+
stop_event: threading.Event,
|
| 539 |
+
) -> None:
|
| 540 |
+
wait_account_engineio_open(ws, token_holder["token"])
|
| 541 |
+
last_token_check = time.monotonic()
|
| 542 |
+
|
| 543 |
+
log("SOCKET", "Account 待機中...")
|
| 544 |
+
|
| 545 |
+
while not stop_event.is_set():
|
| 546 |
+
now = time.monotonic()
|
| 547 |
+
|
| 548 |
+
if now - last_token_check >= TOKEN_CHECK_INTERVAL:
|
| 549 |
+
current_tokens = client.get_tokens() or {}
|
| 550 |
+
current_session_token = current_tokens.get("x-account")
|
| 551 |
+
|
| 552 |
+
if (
|
| 553 |
+
current_session_token
|
| 554 |
+
and current_session_token != token_holder["token"]
|
| 555 |
+
):
|
| 556 |
+
log(
|
| 557 |
+
"AUTH",
|
| 558 |
+
"x-account の変更を検出しました(Account WS に refresh 送信)。",
|
| 559 |
+
)
|
| 560 |
+
send_refresh(ws, current_session_token)
|
| 561 |
+
token_holder["token"] = current_session_token
|
| 562 |
+
|
| 563 |
+
last_token_check = now
|
| 564 |
+
|
| 565 |
+
try:
|
| 566 |
+
packet = ws.recv()
|
| 567 |
+
except websocket.WebSocketTimeoutException:
|
| 568 |
+
continue
|
| 569 |
+
|
| 570 |
+
if packet is None:
|
| 571 |
+
raise RuntimeError("Account WebSocketが切断されました。")
|
| 572 |
+
|
| 573 |
+
if packet == "2":
|
| 574 |
+
send_pong(ws)
|
| 575 |
+
continue
|
| 576 |
+
|
| 577 |
+
if packet == "3":
|
| 578 |
+
log("SOCKET", "Account Engine.IO pong(3) を受信(無視)")
|
| 579 |
+
continue
|
| 580 |
+
|
| 581 |
+
namespace, payload = parse_socketio_event(packet)
|
| 582 |
+
if namespace is None or payload is None:
|
| 583 |
+
continue
|
| 584 |
+
|
| 585 |
+
log(
|
| 586 |
+
"SOCKET",
|
| 587 |
+
f"Account event: namespace={namespace}, payload={payload!r}",
|
| 588 |
+
)
|
| 589 |
|
| 590 |
|
| 591 |
# ---------------------------------------------------------------------------
|
|
|
|
| 601 |
log("MAIN", f"CHANNEL_API_URL={CHANNEL_API_URL}")
|
| 602 |
log("MAIN", f"CHANNEL_LOOKUP_URL={CHANNEL_LOOKUP_URL}")
|
| 603 |
log("MAIN", f"WS_URL={WS_URL}")
|
| 604 |
+
log("MAIN", f"ACCOUNT_WS_URL={ACCOUNT_WS_URL}")
|
| 605 |
|
| 606 |
client = ChannelWorksClient(
|
| 607 |
room_id=room_id,
|
|
|
|
| 632 |
missing = [key for key in required if not tokens.get(key)]
|
| 633 |
if missing:
|
| 634 |
raise RuntimeError(
|
| 635 |
+
"必要なトークンが取得できませんでした: " + ", ".join(missing)
|
|
|
|
| 636 |
)
|
| 637 |
log("AUTH", f"トークン取得 OK({len(tokens)} 件)")
|
| 638 |
|
|
|
|
| 642 |
log("HTTP", f"channel.id = {channel_id}")
|
| 643 |
|
| 644 |
last_session_token = tokens["x-account"]
|
| 645 |
+
token_holder: dict[str, str] = {"token": last_session_token}
|
| 646 |
|
| 647 |
+
# ---- Channel 用 WebSocket ----
|
| 648 |
+
channel_raw_ws = websocket.create_connection(
|
| 649 |
WS_URL,
|
| 650 |
timeout=3.0,
|
| 651 |
origin="https://desk.channel.io",
|
| 652 |
host="desk-ws.channel.io",
|
| 653 |
)
|
| 654 |
+
log("WS", "Channel WebSocket 接続完了")
|
| 655 |
+
channel_ws = LoggingWS(channel_raw_ws, label="CHANNEL-WS")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 656 |
|
| 657 |
+
# ---- Account 用 WebSocket ----
|
| 658 |
+
account_raw_ws = websocket.create_connection(
|
| 659 |
+
ACCOUNT_WS_URL,
|
| 660 |
+
timeout=ACCOUNT_RECV_TIMEOUT,
|
| 661 |
+
origin="https://desk.channel.io",
|
| 662 |
+
host="account-ws.channel.io",
|
| 663 |
+
)
|
| 664 |
+
log("WS", "Account WebSocket 接続完了")
|
| 665 |
+
account_ws = LoggingWS(account_raw_ws, label="ACCOUNT-WS")
|
| 666 |
+
|
| 667 |
+
stop_event = threading.Event()
|
| 668 |
+
|
| 669 |
+
def run_channel() -> None:
|
| 670 |
+
try:
|
| 671 |
+
channel_loop(
|
| 672 |
+
channel_ws,
|
| 673 |
+
token_holder["token"],
|
| 674 |
+
channel_id,
|
| 675 |
+
manager_id,
|
| 676 |
+
stop_event,
|
| 677 |
+
)
|
| 678 |
+
except Exception as exc:
|
| 679 |
+
log("CHANNEL", f"エラーで終了: {exc!r}")
|
| 680 |
+
stop_event.set()
|
| 681 |
+
|
| 682 |
+
def run_account() -> None:
|
| 683 |
+
try:
|
| 684 |
+
account_loop(account_ws, token_holder, client, stop_event)
|
| 685 |
+
except Exception as exc:
|
| 686 |
+
log("ACCOUNT", f"エラーで終了: {exc!r}")
|
| 687 |
+
stop_event.set()
|
| 688 |
+
|
| 689 |
+
channel_thread = threading.Thread(
|
| 690 |
+
target=run_channel, name="channel-ws", daemon=True
|
| 691 |
+
)
|
| 692 |
+
account_thread = threading.Thread(
|
| 693 |
+
target=run_account, name="account-ws", daemon=True
|
| 694 |
+
)
|
| 695 |
+
channel_thread.start()
|
| 696 |
+
account_thread.start()
|
| 697 |
|
| 698 |
+
log("MAIN", "両 WebSocket スレッド起動。待機中...")
|
|
|
|
|
|
|
| 699 |
|
| 700 |
+
try:
|
| 701 |
+
while not stop_event.is_set():
|
| 702 |
+
stop_event.wait(0.5)
|
| 703 |
+
except KeyboardInterrupt:
|
| 704 |
+
log("MAIN", "KeyboardInterrupt を受信しました。")
|
| 705 |
+
stop_event.set()
|
| 706 |
|
| 707 |
+
# スレッド終了を待つ
|
| 708 |
+
channel_thread.join(timeout=5)
|
| 709 |
+
account_thread.join(timeout=5)
|
| 710 |
|
| 711 |
+
try:
|
| 712 |
+
channel_ws.close()
|
| 713 |
+
except Exception:
|
| 714 |
+
pass
|
| 715 |
+
try:
|
| 716 |
+
account_ws.close()
|
| 717 |
+
except Exception:
|
| 718 |
+
pass
|
| 719 |
|
| 720 |
finally:
|
| 721 |
client.stop_auto_touch()
|