File size: 30,058 Bytes
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
38fd3cc
 
 
 
cf15289
38fd3cc
 
 
 
 
 
 
cf15289
 
 
 
 
 
 
38fd3cc
 
 
 
 
0803c29
 
 
 
 
38fd3cc
 
cf15289
84432b4
cf15289
84432b4
cf15289
 
 
 
 
 
 
 
 
 
 
84432b4
 
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
84432b4
 
 
cf15289
84432b4
cf15289
 
84432b4
0803c29
cf15289
84432b4
 
 
cf15289
84432b4
 
 
 
 
 
cf15289
0803c29
84432b4
 
 
 
 
cf15289
84432b4
 
 
cf15289
 
 
 
84432b4
 
 
 
 
 
cf15289
0803c29
84432b4
 
 
 
 
7fced47
 
 
84432b4
 
 
 
cf15289
84432b4
 
 
 
 
 
cf15289
84432b4
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
38fd3cc
cf15289
 
38fd3cc
cf15289
 
38fd3cc
cf15289
 
 
 
 
7530cd8
cf15289
 
 
 
7530cd8
cf15289
 
 
 
 
7530cd8
38fd3cc
cf15289
 
38fd3cc
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
38fd3cc
cf15289
 
38fd3cc
cf15289
 
38fd3cc
cf15289
 
38fd3cc
cf15289
 
 
 
 
 
 
38fd3cc
 
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
0803c29
 
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
0803c29
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
38fd3cc
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
cae42e0
cf15289
 
 
 
 
 
 
 
 
 
 
38fd3cc
cf15289
 
 
0803c29
cf15289
 
 
7fced47
 
 
 
 
 
cf15289
 
 
 
7fced47
cf15289
 
7fced47
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
7fced47
cf15289
 
7fced47
 
 
 
 
 
cf15289
 
 
 
7fced47
cf15289
7fced47
cf15289
 
 
 
 
 
 
 
8b07c37
 
 
 
cf15289
 
38fd3cc
cf15289
 
 
 
 
 
 
 
 
 
84432b4
cf15289
 
 
 
 
 
 
 
38fd3cc
cf15289
 
 
 
 
cae42e0
cf15289
 
38fd3cc
cf15289
 
 
 
 
 
 
0803c29
cf15289
 
0803c29
cf15289
38fd3cc
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
0803c29
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
7fced47
cf15289
 
 
 
 
 
 
 
 
 
 
38fd3cc
cf15289
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
38fd3cc
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
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
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
"""
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()