izuemon commited on
Commit
cf15289
·
verified ·
1 Parent(s): 8b07c37

Update get_msg.py

Browse files
Files changed (1) hide show
  1. get_msg.py +643 -655
get_msg.py CHANGED
@@ -1,9 +1,51 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
  import json
2
- import os
3
  import threading
4
  import time
5
  from datetime import datetime, timezone, timedelta
6
- from typing import Any
7
 
8
  import requests
9
  import websocket
@@ -11,11 +53,14 @@ import websocket
11
  from channel_auth import ChannelWorksClient
12
 
13
 
14
- ROOM_ID = os.getenv("CW_ROOM_ID")
15
- CHANNEL_SLUG = os.getenv("CW_CHANNEL_SLUG")
 
 
 
 
 
16
 
17
- CHANNEL_API_URL = f"https://api.channel.works/desk/channels/{ROOM_ID}"
18
- CHANNEL_LOOKUP_URL = f"https://api.channel.works/desk/channels/{CHANNEL_SLUG}"
19
  WS_URL = (
20
  "wss://desk-ws.channel.io/socket.io/"
21
  "?platform=web&EIO=4&transport=websocket"
@@ -24,73 +69,121 @@ ACCOUNT_WS_URL = (
24
  "wss://account-ws.channel.io/socket.io/"
25
  "?platform=web&EIO=4&transport=websocket"
26
  )
27
- GROUP_PATH = os.getenv("CW_GROUP_ID", "/groups/574628")
28
-
29
- # heartbeat 応答を待つタイムアウト
30
- HEARTBEAT_RESPONSE_TIMEOUT = 5.0
31
- TOKEN_CHECK_INTERVAL = 1.0
32
- REQUEST_TIMEOUT = 30.0
33
- # heartbeat を送る間隔
34
- HEARTBEAT_INTERVAL = 30.0
35
-
36
- # recv 用のタイムアウト
37
- CHANNEL_RECV_TIMEOUT = 3.0
38
- ACCOUNT_RECV_TIMEOUT = 1.0
39
 
40
  JST = timezone(timedelta(hours=9))
41
 
42
-
43
  # ---------------------------------------------------------------------------
44
- # ログ用ヘルパ
45
  # ---------------------------------------------------------------------------
46
- def log(tag: str, message: str) -> None:
47
- now = datetime.now(JST).strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]
48
- print(f"[{now}] [{tag}] {message}", flush=True)
 
 
 
 
 
 
 
 
49
 
50
 
51
- def log_block(tag: str, title: str, body: str) -> None:
52
- now = datetime.now(JST).strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]
53
- print(f"\n[{now}] [{tag}] ===== {title} =====", flush=True)
54
- print(body, flush=True)
55
- print(f"[{now}] [{tag}] ===== /{title} =====", flush=True)
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
56
 
57
 
58
  # ---------------------------------------------------------------------------
59
- # WebSocket ラッパ
60
  # ---------------------------------------------------------------------------
61
- class LoggingWS:
62
- def __init__(self, ws: websocket.WebSocket, label: str = "WS") -> None:
63
  self._ws = ws
64
  self._label = label
 
65
  self._send_count = 0
66
  self._recv_count = 0
67
 
68
- def send(self, payload) -> None:
69
  self._send_count += 1
70
  text = (
71
  payload.decode("utf-8", errors="replace")
72
  if isinstance(payload, (bytes, bytearray))
73
  else payload
74
  )
75
- log_block(
76
  f"{self._label}-SEND",
77
  f"#{self._send_count} ({len(text)} bytes)",
78
  text,
79
  )
80
  self._ws.send(payload)
81
 
82
- def recv(self):
83
  packet = self._ws.recv()
84
  self._recv_count += 1
85
  if packet is None:
86
- log(f"{self._label}-RECV", f"#{self._recv_count} (None / 切断)")
 
 
 
87
  return None
88
  text = (
89
  packet.decode("utf-8", errors="replace")
90
  if isinstance(packet, (bytes, bytearray))
91
  else packet
92
  )
93
- log_block(
94
  f"{self._label}-RECV",
95
  f"#{self._recv_count} ({len(text)} chars)",
96
  text,
@@ -104,683 +197,578 @@ class LoggingWS:
104
  try:
105
  self._ws.close()
106
  finally:
107
- log(self._label, "WebSocket をクローズしました。")
108
 
109
  def __getattr__(self, name: str) -> Any:
110
  return getattr(self._ws, name)
111
 
112
 
113
  # ---------------------------------------------------------------------------
114
- # 環境変数
115
- # ---------------------------------------------------------------------------
116
- def require_env(name: str) -> str:
117
- value = os.getenv(name)
118
- if not value:
119
- raise RuntimeError(f"環境変数 {name} が設定されていません。")
120
- return value
121
-
122
-
123
- # ---------------------------------------------------------------------------
124
- # HTTP
125
- # ---------------------------------------------------------------------------
126
- def get_json(session: requests.Session, url: str, tokens) -> Any:
127
- log("HTTP", f"GET {url}")
128
- response = session.get(
129
- url,
130
- headers={
131
- "Accept": "application/json",
132
- "Accept-Language": "ja",
133
- "Referer": "https://channel.works/",
134
- "Origin": "https://channel.works",
135
- "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
136
- "AppleWebKit/537.36 (KHTML, like Gecko) "
137
- "Chrome/152.0.0.0 Safari/537.36",
138
- "x-account": tokens["x-account"],
139
- },
140
- timeout=REQUEST_TIMEOUT,
141
- )
142
- log("HTTP", f"-> status={response.status_code} ({url})")
143
- response.raise_for_status()
144
- try:
145
- data = response.json()
146
- except ValueError as exc:
147
- raise RuntimeError(
148
- f"JSONを取得できませんでした: {url}\n"
149
- f"status={response.status_code}\n"
150
- f"body={response.text[:500]}"
151
- ) from exc
152
- log_block("HTTP-BODY", url, json.dumps(data, ensure_ascii=False, indent=2))
153
- return data
154
-
155
-
156
- def get_nested(data: Any, *path: str) -> Any:
157
- current = data
158
- for key in path:
159
- if not isinstance(current, dict) or key not in current:
160
- return None
161
- current = current[key]
162
- return current
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]
169
- for value in data.values():
170
- result = find_first_key(value, key)
171
- if result is not None:
172
- return result
173
- elif isinstance(data, list):
174
- for value in data:
175
- result = find_first_key(value, key)
176
- if result is not None:
177
- return result
178
- return None
179
-
180
-
181
- def fetch_ids(tokens: dict[str, str]) -> tuple[str, str]:
182
- log("HTTP", "fetch_ids 開始")
183
- session = requests.Session()
184
- session.cookies.update(
185
- {
186
- "ch-session-1": tokens["ch-session-1"],
187
- "ch-veil-id": tokens["ch-veil-id"],
188
- "x-account": tokens["x-account"],
189
- "x-account-refresh": tokens["x-account-refresh"],
190
- }
191
- )
192
- room_json = get_json(session, CHANNEL_API_URL, tokens)
193
- manager_id = get_nested(room_json, "manager", "id")
194
- if manager_id is None:
195
- manager_id = find_first_key(room_json, "manager")
196
- if isinstance(manager_id, dict):
197
- manager_id = manager_id.get("id")
198
-
199
- channel_json = get_json(session, CHANNEL_LOOKUP_URL, tokens)
200
- channel_id = get_nested(channel_json, "channel", "id")
201
- if channel_id is None:
202
- channel_id = find_first_key(channel_json, "id")
203
-
204
- if manager_id is None:
205
- raise RuntimeError(
206
- f"{CHANNEL_API_URL} のJSONから manager.id を取得できませんでした。"
207
- )
208
- if channel_id is None:
209
- raise RuntimeError(
210
- f"{CHANNEL_LOOKUP_URL} のJSONから channel.id を取得できませんでした。"
211
- )
212
- log("HTTP", f"manager.id = {manager_id}")
213
- log("HTTP", f"channel.id = {channel_id}")
214
- return str(manager_id), str(channel_id)
215
-
216
-
217
- # ---------------------------------------------------------------------------
218
- # 表示整形
219
  # ---------------------------------------------------------------------------
220
- def timestamp_to_text(value: Any) -> str:
221
- if not isinstance(value, (int, float)):
222
- return str(value)
223
- if value < 100_000_000_000:
224
- return str(value)
225
- dt_utc = datetime.fromtimestamp(value / 1000, tz=timezone.utc)
226
- dt_jst = dt_utc.astimezone(JST)
227
- return (
228
- f"{dt_jst:%Y-%m-%d %H:%M:%S.%f} JST "
229
- f"(UTC {dt_utc:%Y-%m-%d %H:%M:%S.%f})"
230
- )
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
231
 
 
 
232
 
233
- def format_json_with_timestamps(value: Any, key: str | None = None) -> Any:
234
- if isinstance(value, dict):
235
- return {k: format_json_with_timestamps(v, k) for k, v in value.items()}
236
- if isinstance(value, list):
237
- return [format_json_with_timestamps(v) for v in value]
238
- if key and key.lower().endswith("at") and isinstance(value, (int, float)):
239
- return {"epoch_ms": value, "formatted": timestamp_to_text(value)}
240
- return value
241
-
242
 
243
- def _print_message_entity(tag: str, entity: dict[str, Any]) -> None:
244
- """メッセージ entity の中身を整形してログ出力する共通処理。"""
245
- for key in (
246
- "chatKey", "id", "mainKey", "channelId", "chatType", "chatId",
247
- "hubId", "personType", "personId", "language", "plainText",
248
- "writingType",
249
- ):
250
- log(tag, f"{key:<11}: {entity.get(key)}")
251
- for key in ("createdAt", "updatedAt"):
252
- if key in entity:
253
- log(tag, f"{key:<11}: {timestamp_to_text(entity[key])}")
254
- blocks = entity.get("blocks")
255
- if blocks:
256
- log(tag, "blocks:")
257
- for index, block in enumerate(blocks, 1):
258
- if isinstance(block, dict):
259
- log(
260
- tag,
261
- f" [{index}] type={block.get('type')} "
262
- f"value={block.get('value')!r}",
263
- )
264
- else:
265
- log(tag, f" [{index}] {block!r}")
266
-
267
-
268
- def print_push_message(payload: list[Any]) -> None:
269
- """push イベントの payload をパースして表示する。
270
-
271
- 想定フォーマット:
272
- 42/desk/channel,["push",{"event":"push","entity":{...},
273
- "type":"message","refers":{...}}]
274
- すなわち payload = ["push", {"event":..., "entity":..., "type":..., "refers":...}]
275
- """
276
- if len(payload) < 2:
277
- log_block("PUSH", "payload", json.dumps(payload, ensure_ascii=False, indent=2))
278
- return
279
-
280
- event_name = payload[0]
281
- wrapper = payload[1]
282
- log("PUSH", f"受信イベント: {event_name}")
283
-
284
- if not isinstance(wrapper, dict):
285
- log_block(
286
- "PUSH",
287
- "wrapper",
288
- json.dumps(
289
- format_json_with_timestamps(wrapper),
290
- ensure_ascii=False,
291
- indent=2,
292
- ),
293
- )
294
- return
295
-
296
- if "type" in wrapper:
297
- log("PUSH", f"type : {wrapper.get('type')}")
298
- if "event" in wrapper:
299
- log("PUSH", f"event : {wrapper.get('event')}")
300
-
301
- entity = wrapper.get("entity", {})
302
- if isinstance(entity, dict):
303
- _print_message_entity("PUSH", entity)
304
-
305
- refers = wrapper.get("refers")
306
- if isinstance(refers, dict):
307
- manager = refers.get("manager")
308
- if isinstance(manager, dict):
309
- log(
310
- "PUSH",
311
- "refers.manager: "
312
- f"id={manager.get('id')} name={manager.get('name')!r}",
313
  )
314
- channel = refers.get("channel")
315
- if isinstance(channel, dict):
316
- log(
317
- "PUSH",
318
- "refers.channel: "
319
- f"id={channel.get('id')} name={channel.get('name')!r}",
320
  )
321
- group = refers.get("group")
322
- if isinstance(group, dict):
323
- log(
324
- "PUSH",
325
- "refers.group : "
326
- f"id={group.get('id')} title={group.get('title')!r}",
327
  )
328
 
329
- log_block(
330
- "PUSH",
331
- "詳細JSON",
332
- json.dumps(
333
- format_json_with_timestamps(payload),
334
- ensure_ascii=False,
335
- indent=2,
336
- ),
337
- )
338
-
339
-
340
- # ---------------------------------------------------------------------------
341
- # Socket.IO / Engine.IO
342
- # ---------------------------------------------------------------------------
343
- def parse_socketio_event(packet: str) -> tuple[str | None, list[Any] | None]:
344
- log("PARSE", f"packet={packet!r}")
345
- if not packet.startswith("42"):
346
- log("PARSE", "-> Socket.IO イベント(42)ではありません。")
347
- return None, None
348
-
349
- payload_text = packet[2:]
350
- if payload_text.startswith("/"):
351
  try:
352
- namespace, payload_text = payload_text.split(",", 1)
353
- except ValueError:
354
- return None, None
355
- else:
356
- namespace = "/"
 
 
 
 
 
 
 
 
 
 
 
 
 
357
 
358
- try:
359
- payload = json.loads(payload_text)
360
- except json.JSONDecodeError:
361
- log("PARSE", f"-> JSON解析失敗: {payload_text!r}")
362
- return namespace, None
363
 
364
- if not isinstance(payload, list):
365
- return namespace, None
366
- return namespace, payload
367
 
 
 
368
 
369
- def send_channel_connect(ws: LoggingWS, session_token: str, channel_id: str) -> None:
370
- message = (
371
- '40/desk/channel,'
372
- + json.dumps(
373
- {"jwt": session_token, "channelId": channel_id},
374
- ensure_ascii=False,
375
- separators=(",", ":"),
376
  )
377
- )
378
- log("SOCKET", f"channel connect 送信: channelId={channel_id}")
379
- ws.send(message)
380
-
381
 
382
- def send_account_connect(ws: LoggingWS, session_token: str) -> None:
383
- message = (
384
- '40/desk/account,'
385
- + json.dumps(
386
- {"jwt": session_token},
387
- ensure_ascii=False,
388
- separators=(",", ":"),
 
 
 
 
 
 
 
 
 
389
  )
390
- )
391
- log("SOCKET", "account connect 送信")
392
- ws.send(message)
393
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
394
 
395
- def send_refresh(ws: LoggingWS, session_token: str) -> None:
396
- """account WS 用の refresh(channelId は送らない)"""
397
- message = (
398
- '42/desk/account,0'
399
- + json.dumps(
400
- ["refresh", {"jwt": session_token}],
401
- ensure_ascii=False,
402
- separators=(",", ":"),
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
403
  )
404
- )
405
- log("SOCKET", "refresh を送信します(account WS)。")
406
- ws.send(message)
407
-
408
-
409
- def send_channel_refresh(
410
- ws: LoggingWS, session_token: str, channel_id: str
411
- ) -> None:
412
- """channel WS 用の refresh(channelId を含める)"""
413
- message = (
414
- '42/desk/channel,3'
415
- + json.dumps(
416
- ["refresh", {"jwt": session_token, "channelId": channel_id}],
417
- ensure_ascii=False,
418
- separators=(",", ":"),
419
  )
420
- )
421
- log(
422
- "SOCKET",
423
- f"refresh を送信します(channel WS, channelId={channel_id})。",
424
- )
425
- ws.send(message)
426
-
427
-
428
- def send_heartbeat(ws: LoggingWS, manager_id: str) -> None:
429
- message = (
430
- '42/desk/channel,'
431
- + json.dumps(
432
- ["heartbeat", manager_id],
433
- ensure_ascii=False,
434
- separators=(",", ":"),
435
  )
436
- )
437
- log("SOCKET", f"heartbeat 送信: manager_id={manager_id}")
438
- ws.send(message)
439
-
440
-
441
- def send_pong(ws: LoggingWS) -> None:
442
- log("SOCKET", "Engine.IO ping(2) を受信 -> pong(3) を返します。")
443
- ws.send("3")
444
-
445
-
446
- # ---------------------------------------------------------------------------
447
- # 各 WS のセットアップ
448
- # ---------------------------------------------------------------------------
449
- def open_channel_ws(session_token: str, channel_id: str) -> LoggingWS:
450
- raw = websocket.create_connection(
451
- WS_URL,
452
- timeout=REQUEST_TIMEOUT,
453
- origin="https://desk.channel.io",
454
- host="desk-ws.channel.io",
455
- )
456
- raw.settimeout(CHANNEL_RECV_TIMEOUT)
457
- log("WS", "Channel WebSocket 接続完了")
458
- ws = LoggingWS(raw, label="CHANNEL-WS")
459
-
460
- # Engine.IO open を待ってから connect を送る
461
- deadline = time.monotonic() + REQUEST_TIMEOUT
462
- connect_sent = False
463
- while time.monotonic() < deadline:
464
- try:
465
- packet = ws.recv()
466
- except websocket.WebSocketTimeoutException:
467
- continue
468
- if packet is None:
469
- raise RuntimeError("Channel WebSocket が切断されました。")
470
- if packet.startswith("0"):
471
- log("SOCKET", "Channel Engine.IO open パケットを受信しました。")
472
- send_channel_connect(ws, session_token, channel_id)
473
- connect_sent = True
474
- continue
475
- if packet == "2":
476
- send_pong(ws)
477
- continue
478
- # connect 応答待ち: 40/desk/channel,{"sid":"..."}
479
- if connect_sent and packet.startswith("40/desk/channel,"):
480
- log("SOCKET", "Channel connect 応答 (sid) を受信しました。")
481
- send_join_group(ws, GROUP_PATH)
482
- return ws
483
- log("SOCKET", f"Channel 接続待機中に受信: {packet!r}")
484
- raise TimeoutError("Channel connect 応答を受信できませんでした。")
485
-
486
-
487
- def open_account_ws(session_token: str) -> LoggingWS:
488
- raw = websocket.create_connection(
489
- ACCOUNT_WS_URL,
490
- timeout=REQUEST_TIMEOUT,
491
- origin="https://desk.channel.io",
492
- host="account-ws.channel.io",
493
- )
494
- raw.settimeout(ACCOUNT_RECV_TIMEOUT)
495
- log("WS", "Account WebSocket 接続完了")
496
- ws = LoggingWS(raw, label="ACCOUNT-WS")
497
-
498
- deadline = time.monotonic() + REQUEST_TIMEOUT
499
- while time.monotonic() < deadline:
500
- try:
501
- packet = ws.recv()
502
- except websocket.WebSocketTimeoutException:
503
- continue
504
- if packet is None:
505
- raise RuntimeError("Account WebSocket が切断されました。")
506
- if packet.startswith("0"):
507
- log("SOCKET", "Account Engine.IO open パケットを受信しました。")
508
- send_account_connect(ws, session_token)
509
- return ws
510
- if packet == "2":
511
- send_pong(ws)
512
- continue
513
- log("SOCKET", f"Account 接続待機中に受信: {packet!r}")
514
- raise TimeoutError("Account Engine.IO open パケットを受信できませんでした。")
515
-
516
-
517
- # ---------------------------------------------------------------------------
518
- # 各スレッド
519
- # ---------------------------------------------------------------------------
520
- def channel_worker(
521
- session_token: str,
522
- channel_id: str,
523
- manager_id: str,
524
- stop_event: threading.Event,
525
- token_holder: dict[str, str],
526
- ) -> None:
527
- try:
528
- ws = open_channel_ws(session_token, channel_id)
529
- except Exception as exc:
530
- log("CHANNEL", f"接続失敗: {exc!r}")
531
- return
532
-
533
- current_token = session_token
534
- last_token_check = time.monotonic()
535
-
536
- ready_received = False
537
- # フェーズ: "waiting"(30秒待機中) / "first_sent" / "second_sent"
538
- heartbeat_phase = "waiting"
539
- phase_started_at = 0.0
540
-
541
- log("SOCKET", "Channel 待機中...")
542
- try:
543
- while not stop_event.is_set():
544
- now = time.monotonic()
545
-
546
- if ready_received:
547
- if heartbeat_phase == "waiting":
548
- # 30 秒待ってから 1 回目を送る
549
- if now - phase_started_at >= HEARTBEAT_INTERVAL:
550
- send_heartbeat(ws, manager_id)
551
- heartbeat_phase = "first_sent"
552
- phase_started_at = now
553
- elif heartbeat_phase == "first_sent":
554
- # 1 回目の応答が 5 秒以内に来なければ再送
555
- if now - phase_started_at >= HEARTBEAT_RESPONSE_TIMEOUT:
556
- log(
557
- "SOCKET",
558
- f"1回目 heartbeat の応答が "
559
- f"{HEARTBEAT_RESPONSE_TIMEOUT}s 以内に来ません。再送します。",
560
- )
561
- send_heartbeat(ws, manager_id)
562
- phase_started_at = now
563
- elif heartbeat_phase == "second_sent":
564
- # 2 回目の応答が 5 秒以内に来なければ再送
565
- if now - phase_started_at >= HEARTBEAT_RESPONSE_TIMEOUT:
566
- log(
567
- "SOCKET",
568
- f"2回目 heartbeat の応答が "
569
- f"{HEARTBEAT_RESPONSE_TIMEOUT}s 以内に来ません。再送します。",
570
- )
571
- send_heartbeat(ws, manager_id)
572
- phase_started_at = now
573
-
574
- # x-account の変更チェック(channel refresh 用)
575
- if now - last_token_check >= TOKEN_CHECK_INTERVAL:
576
- new_token = token_holder.get("token")
577
- if new_token and new_token != current_token:
578
- log("AUTH", "x-account 変更検出 -> Channel refresh 送信")
579
- send_channel_refresh(ws, new_token, channel_id)
580
- current_token = new_token
581
- last_token_check = now
582
 
 
 
 
583
  try:
584
  packet = ws.recv()
585
  except websocket.WebSocketTimeoutException:
586
  continue
587
-
588
  if packet is None:
589
  raise RuntimeError("Channel WebSocket が切断されました。")
590
- if packet == "2":
591
- send_pong(ws)
592
- continue
593
- if packet == "3":
594
- log("SOCKET", "Channel Engine.IO pong(3) 受信(無視)")
595
  continue
596
-
597
- namespace, payload = parse_socketio_event(packet)
598
- if namespace is None or payload is None:
599
  continue
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
600
 
601
- if namespace == "/desk/channel" and payload:
602
- event_name = payload[0]
603
- log("SOCKET", f"event={event_name}")
604
-
605
- if event_name == "ready":
606
- # ready 受信でサイクル開始(まず 30 秒待機)
607
- ready_received = True
608
- heartbeat_phase = "waiting"
609
- phase_started_at = time.monotonic()
610
-
611
- elif event_name == "heartbeat":
612
- if heartbeat_phase == "first_sent":
613
- # 1 回目の応答 → すぐ 2 回目を送る
614
- log("SOCKET", "1回目の応答を受信 → 2回目を送信します。")
615
- send_heartbeat(ws, manager_id)
616
- heartbeat_phase = "second_sent"
617
- phase_started_at = time.monotonic()
618
- elif heartbeat_phase == "second_sent":
619
- # 2 回目の応答 → 次のサイクルへ(30 秒待機)
620
- log("SOCKET", "2回目の応答を受信 → 次のサイクルまで待機します。")
621
- heartbeat_phase = "waiting"
622
- phase_started_at = time.monotonic()
623
-
624
- elif event_name == "create":
625
- # create はパースせず生パケットをそのままログ出力
626
- log_block("CREATE", "raw packet", packet)
627
-
628
- elif event_name == "push":
629
- # push をパースしてメッセージを表示
630
- print_push_message(payload)
631
-
632
- elif event_name == "expired":
633
- log("CHANNEL", "expired を受信しました。接続を終了します。")
634
- break
635
- except Exception as exc:
636
- log("CHANNEL", f"エラーで終了: {exc!r}")
637
- finally:
638
- ws.close()
639
-
640
-
641
- def account_worker(
642
- token_holder: dict[str, str],
643
- client: ChannelWorksClient,
644
- stop_event: threading.Event,
645
- ) -> None:
646
- try:
647
- ws = open_account_ws(token_holder["token"])
648
- except Exception as exc:
649
- log("ACCOUNT", f"接続失敗: {exc!r}")
650
- return
651
-
652
- last_token_check = time.monotonic()
653
- log("SOCKET", "Account 待機中...")
654
- try:
655
- while not stop_event.is_set():
656
- now = time.monotonic()
657
- if now - last_token_check >= TOKEN_CHECK_INTERVAL:
658
- current_tokens = client.get_tokens() or {}
659
- current_session_token = current_tokens.get("x-account")
660
- if (
661
- current_session_token
662
- and current_session_token != token_holder["token"]
663
- ):
664
- log("AUTH", "x-account 変更検出 -> refresh 送信")
665
- send_refresh(ws, current_session_token)
666
- token_holder["token"] = current_session_token
667
- last_token_check = now
668
-
669
  try:
670
  packet = ws.recv()
671
  except websocket.WebSocketTimeoutException:
672
  continue
673
-
674
  if packet is None:
675
  raise RuntimeError("Account WebSocket が切断されました。")
 
 
 
 
676
  if packet == "2":
677
- send_pong(ws)
678
  continue
679
- if packet == "3":
680
- log("SOCKET", "Account Engine.IO pong(3) 受信(無視)")
681
- continue
682
-
683
- namespace, payload = parse_socketio_event(packet)
684
- if namespace is None or payload is None:
685
- continue
686
- log("SOCKET", f"Account event: ns={namespace}, payload={payload!r}")
687
- except Exception as exc:
688
- log("ACCOUNT", f"エラーで終了: {exc!r}")
689
- finally:
690
- ws.close()
691
-
692
- def send_join_group(ws: LoggingWS, group_path: str) -> None:
693
- message = (
694
- '42/desk/channel,'
695
- + json.dumps(
696
  ["join", group_path],
697
  ensure_ascii=False,
698
  separators=(",", ":"),
699
  )
700
- )
701
- log("SOCKET", f"join 送信: {group_path}")
702
- ws.send(message)
703
- # ---------------------------------------------------------------------------
704
- # メイン
705
- # ---------------------------------------------------------------------------
706
- def run() -> None:
707
- log("MAIN", "起動")
708
- room_id = require_env("CW_ROOM_ID") if os.getenv("CW_ROOM_ID") else ROOM_ID
709
- email = require_env("CW_EMAIL")
710
- password = require_env("CW_PASSWORD")
711
 
712
- log("MAIN", f"room_id={room_id}")
713
- log("MAIN", f"WS_URL={WS_URL}")
714
- log("MAIN", f"ACCOUNT_WS_URL={ACCOUNT_WS_URL}")
 
 
 
 
 
 
 
715
 
716
- client = ChannelWorksClient(
717
- room_id=room_id,
718
- email=email,
719
- password=password,
720
- log_level=ChannelWorksClient.LOG_ERROR,
721
- touch_interval=600,
722
- timeout=REQUEST_TIMEOUT,
723
- )
724
 
725
- try:
726
- client.login()
727
- client.start_auto_touch()
728
- tokens = client.get_tokens()
729
- if not tokens:
730
- raise RuntimeError("get_tokens() が空でした。")
731
- required = (
732
- "ch-session-1",
733
- "ch-veil-id",
734
- "x-account",
735
- "x-account-refresh",
736
  )
737
- missing = [k for k in required if not tokens.get(k)]
738
- if missing:
739
- raise RuntimeError(
740
- "必要なトークンが取得できません: " + ", ".join(missing)
741
- )
742
 
743
- manager_id, channel_id = fetch_ids(tokens)
744
- token_holder = {"token": tokens["x-account"]}
745
-
746
- stop_event = threading.Event()
747
-
748
- # channel 側を先に起動
749
- channel_thread = threading.Thread(
750
- target=channel_worker,
751
- args=(
752
- token_holder["token"],
753
- channel_id,
754
- manager_id,
755
- stop_event,
756
- token_holder,
757
- ),
758
- name="channel-ws",
759
- daemon=True,
760
  )
761
- account_thread = threading.Thread(
762
- target=account_worker,
763
- args=(token_holder, client, stop_event),
764
- name="account-ws",
765
- daemon=True,
766
  )
767
- channel_thread.start()
768
- account_thread.start()
769
 
770
- log("MAIN", "両スレッド���動。待機中...")
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
771
  try:
772
- while channel_thread.is_alive() or account_thread.is_alive():
773
- time.sleep(0.5)
774
- except KeyboardInterrupt:
775
- log("MAIN", "KeyboardInterrupt 受信")
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
776
  finally:
777
- stop_event.set()
778
- channel_thread.join(timeout=5)
779
- account_thread.join(timeout=5)
780
- finally:
781
- client.stop_auto_touch()
782
- log("AUTH", "auto_touch 停止")
 
 
 
 
 
783
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
784
 
785
- if __name__ == "__main__":
786
- run()
 
1
+ """
2
+ ChannelWorks メッセージ受信ライブラリ
3
+
4
+ 使い方:
5
+ from channel_auth import ChannelWorksClient
6
+ from channel_messenger import ChannelMessenger
7
+
8
+ client = ChannelWorksClient(
9
+ room_id="...",
10
+ email="you@example.com",
11
+ password="your-password",
12
+ log_level=ChannelWorksClient.LOG_ERROR,
13
+ )
14
+ client.login()
15
+ client.start_auto_touch()
16
+
17
+ def on_push(wrapper):
18
+ entity = wrapper.get("entity") or {}
19
+ print("受信:", entity.get("plainText"))
20
+
21
+ messenger = ChannelMessenger(
22
+ client=client,
23
+ channel_slug="...", # channel_id を自動取得したい場合
24
+ # manager_id="...", channel_id="...", # 直接指定してもよい
25
+ group_path="/groups/グループ番号",
26
+ log_level=ChannelMessenger.LOG_ERROR,
27
+ on_push=on_push,
28
+ )
29
+ messenger.start()
30
+ try:
31
+ messenger.wait()
32
+ except KeyboardInterrupt:
33
+ pass
34
+ finally:
35
+ messenger.stop()
36
+ client.stop_auto_touch()
37
+
38
+ log_level:
39
+ 0 = 何も出力しない
40
+ 1 = エラー・警告のみ (デフォルト)
41
+ 2 = デバッグログも出力 (送受信パケット等をすべて出力)
42
+ """
43
+
44
  import json
 
45
  import threading
46
  import time
47
  from datetime import datetime, timezone, timedelta
48
+ from typing import Any, Callable, Dict, Optional, Tuple
49
 
50
  import requests
51
  import websocket
 
53
  from channel_auth import ChannelWorksClient
54
 
55
 
56
+ # ---------------------------------------------------------------------------
57
+ # 定数
58
+ # ---------------------------------------------------------------------------
59
+ DEFAULT_GROUP_PATH = "/groups/574628"
60
+
61
+ CHANNEL_API_URL_TEMPLATE = "https://api.channel.works/desk/channels/{room_id}"
62
+ CHANNEL_LOOKUP_URL_TEMPLATE = "https://api.channel.works/desk/channels/{slug}"
63
 
 
 
64
  WS_URL = (
65
  "wss://desk-ws.channel.io/socket.io/"
66
  "?platform=web&EIO=4&transport=websocket"
 
69
  "wss://account-ws.channel.io/socket.io/"
70
  "?platform=web&EIO=4&transport=websocket"
71
  )
 
 
 
 
 
 
 
 
 
 
 
 
72
 
73
  JST = timezone(timedelta(hours=9))
74
 
75
+
76
  # ---------------------------------------------------------------------------
77
+ # モジュールレベルユーティリティ(ログに依存しない純粋関数)
78
  # ---------------------------------------------------------------------------
79
+ def _timestamp_to_text(value: Any) -> str:
80
+ if not isinstance(value, (int, float)):
81
+ return str(value)
82
+ if value < 100_000_000_000:
83
+ return str(value)
84
+ dt_utc = datetime.fromtimestamp(value / 1000, tz=timezone.utc)
85
+ dt_jst = dt_utc.astimezone(JST)
86
+ return (
87
+ f"{dt_jst:%Y-%m-%d %H:%M:%S.%f} JST "
88
+ f"(UTC {dt_utc:%Y-%m-%d %H:%M:%S.%f})"
89
+ )
90
 
91
 
92
+ def _format_json_with_timestamps(value: Any, key: Optional[str] = None) -> Any:
93
+ if isinstance(value, dict):
94
+ return {k: _format_json_with_timestamps(v, k) for k, v in value.items()}
95
+ if isinstance(value, list):
96
+ return [_format_json_with_timestamps(v) for v in value]
97
+ if key and key.lower().endswith("at") and isinstance(value, (int, float)):
98
+ return {"epoch_ms": value, "formatted": _timestamp_to_text(value)}
99
+ return value
100
+
101
+
102
+ def _get_nested(data: Any, *path: str) -> Any:
103
+ current = data
104
+ for key in path:
105
+ if not isinstance(current, dict) or key not in current:
106
+ return None
107
+ current = current[key]
108
+ return current
109
+
110
+
111
+ def _find_first_key(data: Any, key: str) -> Any:
112
+ if isinstance(data, dict):
113
+ if key in data:
114
+ return data[key]
115
+ for value in data.values():
116
+ result = _find_first_key(value, key)
117
+ if result is not None:
118
+ return result
119
+ elif isinstance(data, list):
120
+ for value in data:
121
+ result = _find_first_key(value, key)
122
+ if result is not None:
123
+ return result
124
+ return None
125
+
126
+
127
+ def _parse_socketio_event(packet: str) -> Tuple[Optional[str], Optional[list]]:
128
+ if not packet.startswith("42"):
129
+ return None, None
130
+ payload_text = packet[2:]
131
+ if payload_text.startswith("/"):
132
+ try:
133
+ namespace, payload_text = payload_text.split(",", 1)
134
+ except ValueError:
135
+ return None, None
136
+ else:
137
+ namespace = "/"
138
+ try:
139
+ payload = json.loads(payload_text)
140
+ except json.JSONDecodeError:
141
+ return namespace, None
142
+ if not isinstance(payload, list):
143
+ return namespace, None
144
+ return namespace, payload
145
 
146
 
147
  # ---------------------------------------------------------------------------
148
+ # WebSocket ラッパ(送受信をデバッグログするだけの薄いラッパ)
149
  # ---------------------------------------------------------------------------
150
+ class _LoggingWS:
151
+ def __init__(self, ws: websocket.WebSocket, label: str, messenger: "ChannelMessenger") -> None:
152
  self._ws = ws
153
  self._label = label
154
+ self._messenger = messenger
155
  self._send_count = 0
156
  self._recv_count = 0
157
 
158
+ def send(self, payload: Any) -> None:
159
  self._send_count += 1
160
  text = (
161
  payload.decode("utf-8", errors="replace")
162
  if isinstance(payload, (bytes, bytearray))
163
  else payload
164
  )
165
+ self._messenger._log_block_debug(
166
  f"{self._label}-SEND",
167
  f"#{self._send_count} ({len(text)} bytes)",
168
  text,
169
  )
170
  self._ws.send(payload)
171
 
172
+ def recv(self) -> Optional[str]:
173
  packet = self._ws.recv()
174
  self._recv_count += 1
175
  if packet is None:
176
+ self._messenger._log_debug(
177
+ f"{self._label}-RECV",
178
+ f"#{self._recv_count} (None / 切断)",
179
+ )
180
  return None
181
  text = (
182
  packet.decode("utf-8", errors="replace")
183
  if isinstance(packet, (bytes, bytearray))
184
  else packet
185
  )
186
+ self._messenger._log_block_debug(
187
  f"{self._label}-RECV",
188
  f"#{self._recv_count} ({len(text)} chars)",
189
  text,
 
197
  try:
198
  self._ws.close()
199
  finally:
200
+ self._messenger._log_debug(self._label, "WebSocket をクローズしました。")
201
 
202
  def __getattr__(self, name: str) -> Any:
203
  return getattr(self._ws, name)
204
 
205
 
206
  # ---------------------------------------------------------------------------
207
+ # 本体
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
208
  # ---------------------------------------------------------------------------
209
+ class ChannelMessenger:
210
+ # ログレベル
211
+ LOG_NONE = 0
212
+ LOG_ERROR = 1
213
+ LOG_DEBUG = 2
214
+
215
+ def __init__(
216
+ self,
217
+ client: ChannelWorksClient,
218
+ channel_slug: Optional[str] = None,
219
+ manager_id: Optional[str] = None,
220
+ channel_id: Optional[str] = None,
221
+ group_path: str = DEFAULT_GROUP_PATH,
222
+ log_level: int = LOG_ERROR,
223
+ heartbeat_interval: float = 30.0,
224
+ heartbeat_response_timeout: float = 5.0,
225
+ token_check_interval: float = 1.0,
226
+ request_timeout: float = 30.0,
227
+ channel_recv_timeout: float = 3.0,
228
+ account_recv_timeout: float = 1.0,
229
+ on_push: Optional[Callable[[Dict[str, Any]], None]] = None,
230
+ ) -> None:
231
+ if client is None:
232
+ raise ValueError("client (ChannelWorksClient) は必須です。")
233
+
234
+ self.client = client
235
+ self.channel_slug = channel_slug
236
+ self.manager_id = manager_id
237
+ self.channel_id = channel_id
238
+ self.group_path = group_path
239
+ self.log_level = log_level
240
+
241
+ self.heartbeat_interval = heartbeat_interval
242
+ self.heartbeat_response_timeout = heartbeat_response_timeout
243
+ self.token_check_interval = token_check_interval
244
+ self.request_timeout = request_timeout
245
+ self.channel_recv_timeout = channel_recv_timeout
246
+ self.account_recv_timeout = account_recv_timeout
247
+
248
+ self.on_push = on_push
249
+
250
+ self._lock = threading.RLock()
251
+ self._stop_event = threading.Event()
252
+ self._channel_thread: Optional[threading.Thread] = None
253
+ self._account_thread: Optional[threading.Thread] = None
254
+ # x-account の最新値。channel/account 両ワーカーがこの dict を参照する。
255
+ self._token_holder: Dict[str, Optional[str]] = {"token": None}
256
+
257
+ # ------------------------------------------------------------------
258
+ # ログ
259
+ # ------------------------------------------------------------------
260
+ @staticmethod
261
+ def _now_str() -> str:
262
+ return datetime.now(JST).strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]
263
+
264
+ def _log_debug(self, tag: str, msg: str) -> None:
265
+ if self.log_level >= self.LOG_DEBUG:
266
+ print(f"[{self._now_str()}] [{tag}] {msg}", flush=True)
267
+
268
+ def _log_error(self, tag: str, msg: str) -> None:
269
+ if self.log_level >= self.LOG_ERROR:
270
+ print(f"[{self._now_str()}] [{tag}] {msg}", flush=True)
271
+
272
+ def _log_block_debug(self, tag: str, title: str, body: str) -> None:
273
+ if self.log_level < self.LOG_DEBUG:
274
+ return
275
+ ts = self._now_str()
276
+ print(f"\n[{ts}] [{tag}] ===== {title} =====", flush=True)
277
+ print(body, flush=True)
278
+ print(f"[{ts}] [{tag}] ===== /{title} =====", flush=True)
279
+
280
+ def set_log_level(self, level: int) -> None:
281
+ self.log_level = level
282
+
283
+ # ------------------------------------------------------------------
284
+ # 公開 API
285
+ # ------------------------------------------------------------------
286
+ def start(self) -> None:
287
+ """channel / account の各 WebSocket をバックグラウンドで起動する。"""
288
+ with self._lock:
289
+ if self._is_running():
290
+ self._log_debug("MAIN", "ChannelMessenger は既に起動中です。")
291
+ return
292
+
293
+ tokens = self.client.get_tokens() or {}
294
+ required = ("ch-session-1", "ch-veil-id", "x-account", "x-account-refresh")
295
+ missing = [k for k in required if not tokens.get(k)]
296
+ if missing:
297
+ raise RuntimeError(
298
+ "client から必要なトークンが取得できません。先に login() を実行してください: "
299
+ + ", ".join(missing)
300
+ )
301
 
302
+ if self.manager_id is None or self.channel_id is None:
303
+ self._resolve_ids(tokens)
304
 
305
+ self._token_holder["token"] = tokens["x-account"]
306
+ self._stop_event.clear()
 
 
 
 
 
 
 
307
 
308
+ self._channel_thread = threading.Thread(
309
+ target=self._channel_worker,
310
+ args=(tokens["x-account"], self.channel_id, self.manager_id),
311
+ name="channel-ws",
312
+ daemon=True,
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
313
  )
314
+ self._account_thread = threading.Thread(
315
+ target=self._account_worker,
316
+ name="account-ws",
317
+ daemon=True,
 
 
318
  )
319
+ self._channel_thread.start()
320
+ self._account_thread.start()
321
+ self._log_debug(
322
+ "MAIN",
323
+ f"起動: manager_id={self.manager_id} channel_id={self.channel_id}",
 
324
  )
325
 
326
+ def wait(self, poll_interval: float = 0.5) -> None:
327
+ """両ワーカーが終了するまでブロックする。"""
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
328
  try:
329
+ while self._is_running():
330
+ time.sleep(poll_interval)
331
+ except KeyboardInterrupt:
332
+ self._log_debug("MAIN", "KeyboardInterrupt を受信")
333
+ raise
334
+
335
+ def stop(self, join_timeout: float = 5.0) -> None:
336
+ """ワーカーを停止する。"""
337
+ self._stop_event.set()
338
+ for thread in (self._channel_thread, self._account_thread):
339
+ if thread and thread.is_alive():
340
+ thread.join(timeout=join_timeout)
341
+ self._channel_thread = None
342
+ self._account_thread = None
343
+ self._log_debug("MAIN", "ChannelMessenger を停止しました。")
344
+
345
+ def is_running(self) -> bool:
346
+ return self._is_running()
347
 
348
+ def close(self) -> None:
349
+ self.stop()
 
 
 
350
 
351
+ def __enter__(self) -> "ChannelMessenger":
352
+ return self
 
353
 
354
+ def __exit__(self, exc_type, exc, tb) -> None:
355
+ self.close()
356
 
357
+ # ------------------------------------------------------------------
358
+ # 内部: 状態
359
+ # ------------------------------------------------------------------
360
+ def _is_running(self) -> bool:
361
+ return any(
362
+ t is not None and t.is_alive()
363
+ for t in (self._channel_thread, self._account_thread)
364
  )
 
 
 
 
365
 
366
+ # ------------------------------------------------------------------
367
+ # 内部: HTTP / ID 解決
368
+ # ------------------------------------------------------------------
369
+ def _resolve_ids(self, tokens: Dict[str, str]) -> None:
370
+ room_id = getattr(self.client, "room_id", None)
371
+ if not room_id:
372
+ raise RuntimeError("client.room_id が取得できません。")
373
+
374
+ session = requests.Session()
375
+ session.cookies.update(
376
+ {
377
+ "ch-session-1": tokens.get("ch-session-1") or "",
378
+ "ch-veil-id": tokens.get("ch-veil-id") or "",
379
+ "x-account": tokens.get("x-account") or "",
380
+ "x-account-refresh": tokens.get("x-account-refresh") or "",
381
+ }
382
  )
 
 
 
383
 
384
+ room_url = CHANNEL_API_URL_TEMPLATE.format(room_id=room_id)
385
+
386
+ if self.manager_id is None:
387
+ room_json = self._http_get_json(session, room_url, tokens)
388
+ manager_id = _get_nested(room_json, "manager", "id")
389
+ if manager_id is None:
390
+ manager = _find_first_key(room_json, "manager")
391
+ if isinstance(manager, dict):
392
+ manager_id = manager.get("id")
393
+ if manager_id is None:
394
+ raise RuntimeError(
395
+ f"{room_url} の JSON から manager.id を取得できませんでした。"
396
+ )
397
+ self.manager_id = str(manager_id)
398
 
399
+ if self.channel_id is None:
400
+ if not self.channel_slug:
401
+ raise RuntimeError(
402
+ "channel_id を自動解決するには channel_slug を指定してください。"
403
+ )
404
+ lookup_url = CHANNEL_LOOKUP_URL_TEMPLATE.format(slug=self.channel_slug)
405
+ channel_json = self._http_get_json(session, lookup_url, tokens)
406
+ channel_id = _get_nested(channel_json, "channel", "id")
407
+ if channel_id is None:
408
+ channel_id = _find_first_key(channel_json, "id")
409
+ if channel_id is None:
410
+ raise RuntimeError(
411
+ f"{lookup_url} の JSON から channel.id を取得できませんでした。"
412
+ )
413
+ self.channel_id = str(channel_id)
414
+
415
+ def _http_get_json(
416
+ self, session: requests.Session, url: str, tokens: Dict[str, str]
417
+ ) -> Any:
418
+ self._log_debug("HTTP", f"GET {url}")
419
+ response = session.get(
420
+ url,
421
+ headers={
422
+ "Accept": "application/json",
423
+ "Accept-Language": "ja",
424
+ "Referer": "https://channel.works/",
425
+ "Origin": "https://channel.works",
426
+ "User-Agent": (
427
+ "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
428
+ "AppleWebKit/537.36 (KHTML, like Gecko) "
429
+ "Chrome/152.0.0.0 Safari/537.36"
430
+ ),
431
+ "x-account": tokens["x-account"],
432
+ },
433
+ timeout=self.request_timeout,
434
  )
435
+ self._log_debug("HTTP", f"-> status={response.status_code} ({url})")
436
+ response.raise_for_status()
437
+ try:
438
+ data = response.json()
439
+ except ValueError as exc:
440
+ raise RuntimeError(
441
+ f"JSON を取得できませんでした: {url}\n"
442
+ f"status={response.status_code}\n"
443
+ f"body={response.text[:500]}"
444
+ ) from exc
445
+ self._log_block_debug(
446
+ "HTTP-BODY",
447
+ url,
448
+ json.dumps(data, ensure_ascii=False),
 
449
  )
450
+ return data
451
+
452
+ # ------------------------------------------------------------------
453
+ # 内部: WebSocket 接続
454
+ # ------------------------------------------------------------------
455
+ def _open_channel_ws(self, session_token: str, channel_id: str) -> _LoggingWS:
456
+ raw = websocket.create_connection(
457
+ WS_URL,
458
+ timeout=self.request_timeout,
459
+ origin="https://desk.channel.io",
460
+ host="desk-ws.channel.io",
 
 
 
 
461
  )
462
+ raw.settimeout(self.channel_recv_timeout)
463
+ self._log_debug("WS", "Channel WebSocket 接続完了")
464
+ ws = _LoggingWS(raw, label="CHANNEL-WS", messenger=self)
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
465
 
466
+ deadline = time.monotonic() + self.request_timeout
467
+ connect_sent = False
468
+ while time.monotonic() < deadline:
469
  try:
470
  packet = ws.recv()
471
  except websocket.WebSocketTimeoutException:
472
  continue
 
473
  if packet is None:
474
  raise RuntimeError("Channel WebSocket が切断されました。")
475
+ if packet.startswith("0"):
476
+ self._log_debug("SOCKET", "Channel Engine.IO open 受信")
477
+ self._send_channel_connect(ws, session_token, channel_id)
478
+ connect_sent = True
 
479
  continue
480
+ if packet == "2":
481
+ self._send_pong(ws)
 
482
  continue
483
+ if connect_sent and packet.startswith("40/desk/channel,"):
484
+ self._log_debug("SOCKET", "Channel connect 応答 (sid) 受信")
485
+ self._send_join_group(ws, self.group_path)
486
+ return ws
487
+ self._log_debug("SOCKET", f"Channel 接続待機中の受信: {packet!r}")
488
+ raise TimeoutError("Channel connect 応答を受信できませんでした。")
489
+
490
+ def _open_account_ws(self, session_token: str) -> _LoggingWS:
491
+ raw = websocket.create_connection(
492
+ ACCOUNT_WS_URL,
493
+ timeout=self.request_timeout,
494
+ origin="https://desk.channel.io",
495
+ host="account-ws.channel.io",
496
+ )
497
+ raw.settimeout(self.account_recv_timeout)
498
+ self._log_debug("WS", "Account WebSocket 接続完了")
499
+ ws = _LoggingWS(raw, label="ACCOUNT-WS", messenger=self)
500
 
501
+ deadline = time.monotonic() + self.request_timeout
502
+ while time.monotonic() < deadline:
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
503
  try:
504
  packet = ws.recv()
505
  except websocket.WebSocketTimeoutException:
506
  continue
 
507
  if packet is None:
508
  raise RuntimeError("Account WebSocket が切断されました。")
509
+ if packet.startswith("0"):
510
+ self._log_debug("SOCKET", "Account Engine.IO open 受信")
511
+ self._send_account_connect(ws, session_token)
512
+ return ws
513
  if packet == "2":
514
+ self._send_pong(ws)
515
  continue
516
+ self._log_debug("SOCKET", f"Account 接続待機中の受信: {packet!r}")
517
+ raise TimeoutError("Account Engine.IO open パケットを受信できませんでした。")
518
+
519
+ # ------------------------------------------------------------------
520
+ # 内部: Socket.IO パケット送信
521
+ # ------------------------------------------------------------------
522
+ def _send_join_group(self, ws: _LoggingWS, group_path: str) -> None:
523
+ message = '42/desk/channel,' + json.dumps(
 
 
 
 
 
 
 
 
 
524
  ["join", group_path],
525
  ensure_ascii=False,
526
  separators=(",", ":"),
527
  )
528
+ self._log_debug("SOCKET", f"join 送信: {group_path}")
529
+ ws.send(message)
 
 
 
 
 
 
 
 
 
530
 
531
+ def _send_channel_connect(
532
+ self, ws: _LoggingWS, session_token: str, channel_id: str
533
+ ) -> None:
534
+ message = '40/desk/channel,' + json.dumps(
535
+ {"jwt": session_token, "channelId": channel_id},
536
+ ensure_ascii=False,
537
+ separators=(",", ":"),
538
+ )
539
+ self._log_debug("SOCKET", f"channel connect 送信: channelId={channel_id}")
540
+ ws.send(message)
541
 
542
+ def _send_account_connect(self, ws: _LoggingWS, session_token: str) -> None:
543
+ message = '40/desk/account,' + json.dumps(
544
+ {"jwt": session_token},
545
+ ensure_ascii=False,
546
+ separators=(",", ":"),
547
+ )
548
+ self._log_debug("SOCKET", "account connect 送信")
549
+ ws.send(message)
550
 
551
+ def _send_refresh(self, ws: _LoggingWS, session_token: str) -> None:
552
+ message = '42/desk/account,0' + json.dumps(
553
+ ["refresh", {"jwt": session_token}],
554
+ ensure_ascii=False,
555
+ separators=(",", ":"),
 
 
 
 
 
 
556
  )
557
+ self._log_debug("SOCKET", "refresh 送信(account WS)")
558
+ ws.send(message)
 
 
 
559
 
560
+ def _send_channel_refresh(
561
+ self, ws: _LoggingWS, session_token: str, channel_id: str
562
+ ) -> None:
563
+ message = '42/desk/channel,3' + json.dumps(
564
+ ["refresh", {"jwt": session_token, "channelId": channel_id}],
565
+ ensure_ascii=False,
566
+ separators=(",", ":"),
 
 
 
 
 
 
 
 
 
 
567
  )
568
+ self._log_debug(
569
+ "SOCKET", f"refresh 送信(channel WS, channelId={channel_id})"
 
 
 
570
  )
571
+ ws.send(message)
 
572
 
573
+ def _send_heartbeat(self, ws: _LoggingWS, manager_id: str) -> None:
574
+ message = '42/desk/channel,' + json.dumps(
575
+ ["heartbeat", manager_id],
576
+ ensure_ascii=False,
577
+ separators=(",", ":"),
578
+ )
579
+ self._log_debug("SOCKET", f"heartbeat 送信: manager_id={manager_id}")
580
+ ws.send(message)
581
+
582
+ def _send_pong(self, ws: _LoggingWS) -> None:
583
+ self._log_debug("SOCKET", "Engine.IO ping(2) 受信 -> pong(3) 送信")
584
+ ws.send("3")
585
+
586
+ # ------------------------------------------------------------------
587
+ # 内部: push 処理
588
+ # ------------------------------------------------------------------
589
+ def _handle_push(self, payload: list) -> None:
590
+ if len(payload) < 2:
591
+ return
592
+ wrapper = payload[1]
593
+ if not isinstance(wrapper, dict):
594
+ return
595
+
596
+ if self.log_level >= self.LOG_DEBUG:
597
+ self._log_block_debug(
598
+ "PUSH",
599
+ f"type={wrapper.get('type')} event={wrapper.get('event')}",
600
+ json.dumps(
601
+ _format_json_with_timestamps(wrapper),
602
+ ensure_ascii=False,
603
+ ),
604
+ )
605
+
606
+ if self.on_push is not None:
607
+ try:
608
+ self.on_push(wrapper)
609
+ except Exception as exc:
610
+ self._log_error("PUSH", f"on_push コールバックで例外: {exc!r}")
611
+
612
+ # ------------------------------------------------------------------
613
+ # 内部: channel ワーカー
614
+ # ------------------------------------------------------------------
615
+ def _channel_worker(
616
+ self, session_token: str, channel_id: str, manager_id: str
617
+ ) -> None:
618
  try:
619
+ ws = self._open_channel_ws(session_token, channel_id)
620
+ except Exception as exc:
621
+ self._log_error("CHANNEL", f"接続失敗: {exc!r}")
622
+ return
623
+
624
+ current_token = session_token
625
+ last_token_check = time.monotonic()
626
+
627
+ ready_received = False
628
+ heartbeat_phase = "waiting" # "waiting" / "first_sent" / "second_sent"
629
+ phase_started_at = 0.0
630
+
631
+ try:
632
+ while not self._stop_event.is_set():
633
+ now = time.monotonic()
634
+
635
+ if ready_received:
636
+ if heartbeat_phase == "waiting":
637
+ if now - phase_started_at >= self.heartbeat_interval:
638
+ self._send_heartbeat(ws, manager_id)
639
+ heartbeat_phase = "first_sent"
640
+ phase_started_at = now
641
+ elif heartbeat_phase == "first_sent":
642
+ if now - phase_started_at >= self.heartbeat_response_timeout:
643
+ self._log_debug(
644
+ "SOCKET",
645
+ f"1回目 heartbeat 応答なし ({self.heartbeat_response_timeout}s) → 再送",
646
+ )
647
+ self._send_heartbeat(ws, manager_id)
648
+ phase_started_at = now
649
+ elif heartbeat_phase == "second_sent":
650
+ if now - phase_started_at >= self.heartbeat_response_timeout:
651
+ self._log_debug(
652
+ "SOCKET",
653
+ f"2回目 heartbeat 応答なし ({self.heartbeat_response_timeout}s) → 再送",
654
+ )
655
+ self._send_heartbeat(ws, manager_id)
656
+ phase_started_at = now
657
+
658
+ if now - last_token_check >= self.token_check_interval:
659
+ new_token = self._token_holder.get("token")
660
+ if new_token and new_token != current_token:
661
+ self._log_debug(
662
+ "AUTH", "x-account 変更検出 → Channel refresh 送信"
663
+ )
664
+ self._send_channel_refresh(ws, new_token, channel_id)
665
+ current_token = new_token
666
+ last_token_check = now
667
+
668
+ try:
669
+ packet = ws.recv()
670
+ except websocket.WebSocketTimeoutException:
671
+ continue
672
+
673
+ if packet is None:
674
+ raise RuntimeError("Channel WebSocket が切断されました。")
675
+ if packet == "2":
676
+ self._send_pong(ws)
677
+ continue
678
+ if packet == "3":
679
+ self._log_debug("SOCKET", "Channel Engine.IO pong(3) 受信")
680
+ continue
681
+
682
+ namespace, payload = _parse_socketio_event(packet)
683
+ if namespace is None or payload is None:
684
+ continue
685
+
686
+ if namespace == "/desk/channel" and payload:
687
+ event_name = payload[0]
688
+ self._log_debug("SOCKET", f"event={event_name}")
689
+
690
+ if event_name == "ready":
691
+ ready_received = True
692
+ heartbeat_phase = "waiting"
693
+ phase_started_at = time.monotonic()
694
+ elif event_name == "heartbeat":
695
+ if heartbeat_phase == "first_sent":
696
+ self._log_debug(
697
+ "SOCKET", "1回目の応答受信 → 2回目を送信"
698
+ )
699
+ self._send_heartbeat(ws, manager_id)
700
+ heartbeat_phase = "second_sent"
701
+ phase_started_at = time.monotonic()
702
+ elif heartbeat_phase == "second_sent":
703
+ self._log_debug(
704
+ "SOCKET", "2回目の応答受信 → 次サイクル待機"
705
+ )
706
+ heartbeat_phase = "waiting"
707
+ phase_started_at = time.monotonic()
708
+ elif event_name == "create":
709
+ self._log_block_debug("CREATE", "raw packet", packet)
710
+ elif event_name == "push":
711
+ self._handle_push(payload)
712
+ elif event_name == "expired":
713
+ self._log_error(
714
+ "CHANNEL", "expired を受信 → 接続を終了します。"
715
+ )
716
+ break
717
+ except Exception as exc:
718
+ self._log_error("CHANNEL", f"エラーで終了: {exc!r}")
719
  finally:
720
+ ws.close()
721
+
722
+ # ------------------------------------------------------------------
723
+ # 内部: account ワーカー
724
+ # ------------------------------------------------------------------
725
+ def _account_worker(self) -> None:
726
+ try:
727
+ ws = self._open_account_ws(self._token_holder["token"] or "")
728
+ except Exception as exc:
729
+ self._log_error("ACCOUNT", f"接続失敗: {exc!r}")
730
+ return
731
 
732
+ last_token_check = time.monotonic()
733
+ try:
734
+ while not self._stop_event.is_set():
735
+ now = time.monotonic()
736
+ if now - last_token_check >= self.token_check_interval:
737
+ current_tokens = self.client.get_tokens() or {}
738
+ current_session_token = current_tokens.get("x-account")
739
+ if (
740
+ current_session_token
741
+ and current_session_token != self._token_holder["token"]
742
+ ):
743
+ self._log_debug(
744
+ "AUTH", "x-account 変更検出 → refresh 送信(account WS)"
745
+ )
746
+ self._send_refresh(ws, current_session_token)
747
+ self._token_holder["token"] = current_session_token
748
+ last_token_check = now
749
+
750
+ try:
751
+ packet = ws.recv()
752
+ except websocket.WebSocketTimeoutException:
753
+ continue
754
+
755
+ if packet is None:
756
+ raise RuntimeError("Account WebSocket が切断されました。")
757
+ if packet == "2":
758
+ self._send_pong(ws)
759
+ continue
760
+ if packet == "3":
761
+ self._log_debug("SOCKET", "Account Engine.IO pong(3) 受信")
762
+ continue
763
+
764
+ namespace, payload = _parse_socketio_event(packet)
765
+ if namespace is None or payload is None:
766
+ continue
767
+ self._log_debug(
768
+ "SOCKET", f"Account event: ns={namespace}, payload={payload!r}"
769
+ )
770
+ except Exception as exc:
771
+ self._log_error("ACCOUNT", f"エラーで終了: {exc!r}")
772
+ finally:
773
+ ws.close()
774