|
| 1 | +import base64 |
| 2 | +import json |
| 3 | +import logging |
| 4 | +import random |
| 5 | +import string |
| 6 | +import threading |
| 7 | +import time |
| 8 | +import urllib.request |
| 9 | +import urllib.error |
| 10 | +from typing import Callable, Optional, List, Dict, Any |
| 11 | + |
| 12 | +logger = logging.getLogger(__name__) |
| 13 | + |
| 14 | + |
| 15 | +class RendezqueueClient: |
| 16 | + def __init__( |
| 17 | + self, |
| 18 | + url: str, |
| 19 | + key: str, |
| 20 | + hue: str, |
| 21 | + on_data: Optional[Callable[[List[bytes]], None]] = None, |
| 22 | + on_error: Optional[Callable[[Exception], None]] = None, |
| 23 | + poll_interval_ms: int = 2000, |
| 24 | + ) -> None: |
| 25 | + self.url = url |
| 26 | + self.key = key |
| 27 | + self.hue = hue |
| 28 | + self.on_data = on_data |
| 29 | + self.on_error = on_error or (lambda e: logger.error(f"Rendezqueue error: {e}")) |
| 30 | + self.poll_interval_ms = poll_interval_ms |
| 31 | + |
| 32 | + self.sid_counter = 1 |
| 33 | + self.sid = self._generate_sid() |
| 34 | + self.offset = 0 |
| 35 | + self.outgoing_queue: List[bytes] = [] |
| 36 | + self.lock = threading.Lock() |
| 37 | + |
| 38 | + self.is_polling = False |
| 39 | + self.is_stopped = True |
| 40 | + self.poll_thread: Optional[threading.Thread] = None |
| 41 | + |
| 42 | + def _generate_sid(self) -> str: |
| 43 | + random_part = "".join( |
| 44 | + random.choices(string.ascii_lowercase + string.digits, k=7) |
| 45 | + ) |
| 46 | + return f"{self.hue}-{self.sid_counter}-{random_part}" |
| 47 | + |
| 48 | + def start(self) -> None: |
| 49 | + if not self.is_stopped: |
| 50 | + return |
| 51 | + self.is_stopped = False |
| 52 | + self.poll_thread = threading.Thread(target=self._poll_loop, daemon=True) |
| 53 | + self.poll_thread.start() |
| 54 | + |
| 55 | + def stop(self) -> None: |
| 56 | + if self.is_stopped: |
| 57 | + return |
| 58 | + self.is_stopped = True |
| 59 | + |
| 60 | + def send(self, value: Any) -> None: |
| 61 | + if isinstance(value, str): |
| 62 | + value = value.encode("utf-8") |
| 63 | + with self.lock: |
| 64 | + self.outgoing_queue.append(value) |
| 65 | + |
| 66 | + def _start_new_session(self) -> None: |
| 67 | + with self.lock: |
| 68 | + self.sid_counter += 1 |
| 69 | + self.sid = self._generate_sid() |
| 70 | + self.offset = 0 |
| 71 | + |
| 72 | + def _poll_loop(self) -> None: |
| 73 | + # Initial poll |
| 74 | + self._poll() |
| 75 | + while not self.is_stopped: |
| 76 | + time.sleep(self.poll_interval_ms / 1000.0) |
| 77 | + if self.is_stopped: |
| 78 | + break |
| 79 | + self._poll() |
| 80 | + |
| 81 | + def _poll(self) -> None: |
| 82 | + if self.is_polling: |
| 83 | + return |
| 84 | + self.is_polling = True |
| 85 | + try: |
| 86 | + with self.lock: |
| 87 | + req_sid = self.sid |
| 88 | + req_offset = self.offset |
| 89 | + # Take snapshot of queue to send |
| 90 | + snapshot_queue = list(self.outgoing_queue) |
| 91 | + b64_values = [ |
| 92 | + base64.b64encode(v).decode("ascii") for v in snapshot_queue |
| 93 | + ] |
| 94 | + |
| 95 | + request_body = { |
| 96 | + "key": self.key, |
| 97 | + "sid": req_sid, |
| 98 | + "offset": req_offset, |
| 99 | + "values": b64_values, |
| 100 | + "b64": 1, |
| 101 | + } |
| 102 | + |
| 103 | + data = json.dumps(request_body).encode("utf-8") |
| 104 | + headers = { |
| 105 | + "Content-Type": "application/json", |
| 106 | + "User-Agent": "Rendezqueue-Python-Client/0.0.0", |
| 107 | + } |
| 108 | + req = urllib.request.Request(self.url, data=data, headers=headers) |
| 109 | + |
| 110 | + try: |
| 111 | + with urllib.request.urlopen(req, timeout=10) as response: |
| 112 | + if response.status != 200: |
| 113 | + text = response.read().decode("utf-8") |
| 114 | + self.on_error( |
| 115 | + Exception(f"Server error: {response.status} {text}") |
| 116 | + ) |
| 117 | + return |
| 118 | + |
| 119 | + resp_body = response.read().decode("utf-8") |
| 120 | + msg = json.loads(resp_body) |
| 121 | + data = self._decode_response(msg) |
| 122 | + |
| 123 | + received_values = data.get("values", []) |
| 124 | + # Session ended if we got values OR (server offset > 0 and no ttl/keepalive) |
| 125 | + session_has_ended = len(received_values) > 0 or ( |
| 126 | + data.get("offset", 0) > 0 and "ttl" not in data |
| 127 | + ) |
| 128 | + |
| 129 | + with self.lock: |
| 130 | + if self.sid != req_sid: |
| 131 | + # Session changed, ignore response |
| 132 | + return |
| 133 | + |
| 134 | + if session_has_ended: |
| 135 | + # Remove sent messages |
| 136 | + sent_count = len(snapshot_queue) |
| 137 | + del self.outgoing_queue[:sent_count] |
| 138 | + |
| 139 | + # Note: self.offset is reset by _start_new_session anyway |
| 140 | + else: |
| 141 | + server_offset = data.get("offset", 0) |
| 142 | + accepted_count = server_offset - self.offset |
| 143 | + if accepted_count > 0: |
| 144 | + del self.outgoing_queue[:accepted_count] |
| 145 | + self.offset = server_offset |
| 146 | + |
| 147 | + if session_has_ended: |
| 148 | + if self.on_data: |
| 149 | + self.on_data(received_values) |
| 150 | + self._start_new_session() |
| 151 | + |
| 152 | + except urllib.error.HTTPError as e: |
| 153 | + text = e.read().decode("utf-8") |
| 154 | + self.on_error(Exception(f"Server error: {e.code} {text}")) |
| 155 | + except Exception as e: |
| 156 | + self.on_error(e) |
| 157 | + |
| 158 | + except Exception as e: |
| 159 | + self.on_error(e) |
| 160 | + finally: |
| 161 | + self.is_polling = False |
| 162 | + |
| 163 | + def _decode_response(self, msg: Dict[str, Any]) -> Dict[str, Any]: |
| 164 | + b64_flags = msg.get("b64", 0) |
| 165 | + if b64_flags & 4: |
| 166 | + msg["key"] = self._b64decode_padded(msg["key"]).decode("utf-8") |
| 167 | + if b64_flags & 2: |
| 168 | + msg["sid"] = self._b64decode_padded(msg["sid"]).decode("utf-8") |
| 169 | + if "values" in msg and (b64_flags & 1): |
| 170 | + msg["values"] = [self._b64decode_padded(v) for v in msg["values"]] |
| 171 | + return msg |
| 172 | + |
| 173 | + def _b64decode_padded(self, s: str) -> bytes: |
| 174 | + missing_padding = len(s) % 4 |
| 175 | + if missing_padding: |
| 176 | + s += "=" * (4 - missing_padding) |
| 177 | + return base64.b64decode(s, validate=False) |
0 commit comments