From e5d7ad9b410dc82f9b5725d135d2b60b45906765 Mon Sep 17 00:00:00 2001 From: SergeantPanda Date: Tue, 11 Mar 2025 11:38:26 -0500 Subject: [PATCH] Data has been resturcured to make more sense. --- apps/proxy/ts_proxy/server.py | 84 ++++++++++++++++++++++++----------- apps/proxy/ts_proxy/views.py | 41 +++++++++++++++-- 2 files changed, 95 insertions(+), 30 deletions(-) diff --git a/apps/proxy/ts_proxy/server.py b/apps/proxy/ts_proxy/server.py index 3600513b..31550ed7 100644 --- a/apps/proxy/ts_proxy/server.py +++ b/apps/proxy/ts_proxy/server.py @@ -649,23 +649,20 @@ class ClientManager: if time_since_last < self.heartbeat_interval * 0.8: continue - # Update the client's individual key with new TTL - client_key = f"ts_proxy:client:{self.channel_id}:{client_id}" - pipe.setex(client_key, self.client_ttl, str(current_time)) + # FIXED: Update client hash instead of separate activity key + client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" + + # Only update the last_active field in the hash, preserving other fields + pipe.hset(client_key, "last_active", str(current_time)) + pipe.expire(client_key, self.client_ttl) # Refresh TTL # Keep client in the set with TTL pipe.sadd(self.client_set_key, client_id) - - # Update last activity timestamp in a separate key - activity_key = f"ts_proxy:client:{self.channel_id}:{client_id}:last_active" - pipe.setex(activity_key, self.client_ttl, str(current_time)) + pipe.expire(self.client_set_key, self.client_ttl) # Track last heartbeat locally self.last_heartbeat_time[client_id] = current_time - # Always refresh the TTL on the set itself - pipe.expire(self.client_set_key, self.client_ttl) - # Execute all commands atomically pipe.execute() @@ -674,7 +671,7 @@ class ClientManager: except Exception as e: logging.error(f"Error in client heartbeat thread: {e}") - + thread = threading.Thread(target=heartbeat_task, daemon=True) thread.name = f"client-heartbeat-{self.channel_id}" thread.start() @@ -698,7 +695,7 @@ class ClientManager: except Exception as e: logging.error(f"Error notifying owner of client activity: {e}") - def add_client(self, client_id): + def add_client(self, client_id, user_agent=None): """Add a client to this channel locally and in Redis""" with self.lock: self.clients.add(client_id) @@ -707,32 +704,51 @@ class ClientManager: if self.redis_client: current_time = str(time.time()) - # Add to channel's client set - already standardized + # Add to channel's client set self.redis_client.sadd(self.client_set_key, client_id) self.redis_client.expire(self.client_set_key, self.client_ttl) # STANDARDIZED KEY: Individual client under channel namespace - client_key = f"ts_proxy:channel:{self.channel_id}:client:{client_id}" - self.redis_client.setex(client_key, self.client_ttl, current_time) + client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" - # STANDARDIZED KEY: Client activity under channel namespace - activity_key = f"ts_proxy:channel:{self.channel_id}:client:{client_id}:last_active" - self.redis_client.setex(activity_key, self.client_ttl, current_time) + # Store client info as a hash with all info in one place + client_data = { + "last_active": current_time, + "worker_id": self.worker_id or "unknown", + "connect_time": current_time + } + + # Add user agent if provided + if user_agent: + client_data["user_agent"] = user_agent + + # Use HSET to store client data as a hash + self.redis_client.hset(client_key, mapping=client_data) + self.redis_client.expire(client_key, self.client_ttl) # Clear any initialization timer self.redis_client.delete(f"ts_proxy:channel:{self.channel_id}:init_time") self._notify_owner_of_activity() - # Publish client connected event - event_data = json.dumps({ + # Publish client connected event with user agent + event_data = { "event": "client_connected", "channel_id": self.channel_id, "client_id": client_id, "worker_id": self.worker_id or "unknown", "timestamp": time.time() - }) - self.redis_client.publish(f"ts_proxy:events:{self.channel_id}", event_data) + } + + if user_agent: + event_data["user_agent"] = user_agent + logging.debug(f"Storing user agent '{user_agent}' for client {client_id}") + else: + logging.debug(f"No user agent provided for client {client_id}") + self.redis_client.publish( + f"ts_proxy:events:{self.channel_id}", + json.dumps(event_data) + ) # Get total clients across all workers total_clients = self.get_total_client_count() @@ -758,8 +774,8 @@ class ClientManager: self.redis_client.srem(self.client_set_key, client_id) # STANDARDIZED KEY: Delete individual client keys - client_key = f"ts_proxy:channel:{self.channel_id}:client:{client_id}" - activity_key = f"ts_proxy:channel:{self.channel_id}:client:{client_id}:last_active" + client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" + activity_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}:last_active" self.redis_client.delete(client_key, activity_key) self._notify_owner_of_activity() @@ -804,7 +820,8 @@ class ClientManager: try: # Refresh TTL for all clients belonging to this worker for client_id in self.clients: - client_key = f"ts_proxy:client:{self.channel_id}:{client_id}" + # STANDARDIZED: Use channel namespace for client keys + client_key = f"ts_proxy:channel:{self.channel_id}:clients:{client_id}" self.redis_client.expire(client_key, self.client_ttl) # Refresh TTL on the set itself @@ -1264,12 +1281,25 @@ class ProxyServer: # Check Redis using the standard key pattern if self.redis_client: + # Primary check - look for channel metadata metadata_key = f"ts_proxy:channel:{channel_id}:metadata" - # If metadata exists return true + # If metadata exists, return true if self.redis_client.exists(metadata_key): return True - + + # Additional checks if metadata doesn't exist + additional_keys = [ + f"ts_proxy:channel:{channel_id}:active", + f"ts_proxy:channel:{channel_id}:clients", + f"ts_proxy:channel:{channel_id}:buffer:index", + f"ts_proxy:channel:{channel_id}:owner" + ] + + for key in additional_keys: + if self.redis_client.exists(key): + return True + return False def stop_channel(self, channel_id): diff --git a/apps/proxy/ts_proxy/views.py b/apps/proxy/ts_proxy/views.py index 271ae380..867a9b6d 100644 --- a/apps/proxy/ts_proxy/views.py +++ b/apps/proxy/ts_proxy/views.py @@ -98,11 +98,46 @@ def stream_ts(request, channel_id): chunks_sent = 0 try: - logger.info(f"[{client_id}] New client connected to channel {channel_id}") + # ENHANCED USER AGENT DETECTION - check multiple possible headers + user_agent = None - # Add client to manager + # Try multiple possible header formats + ua_headers = ['HTTP_USER_AGENT', 'User-Agent', 'user-agent', 'User_Agent'] + + for header in ua_headers: + if header in request.META: + user_agent = request.META[header] + logger.debug(f"Found user agent in header: {header}") + break + + # Try request.headers dictionary (Django 2.2+) + if not user_agent and hasattr(request, 'headers'): + for header in ['User-Agent', 'user-agent']: + if header in request.headers: + user_agent = request.headers[header] + logger.debug(f"Found user agent in request.headers: {header}") + break + + # Final fallback - check if in any header with case-insensitive matching + if not user_agent: + for key, value in request.META.items(): + if key.upper().replace('_', '-') == 'USER-AGENT': + user_agent = value + logger.debug(f"Found user agent in alternate header: {key}") + break + + # Log headers for debugging user agent issues + if not user_agent: + # Log all headers to help troubleshoot + headers = {k: v for k, v in request.META.items() if k.startswith('HTTP_')} + logger.debug(f"No user agent found in request. Available headers: {headers}") + user_agent = "Unknown-Client" # Default value instead of None + + logger.info(f"[{client_id}] New client connected to channel {channel_id} with user agent: {user_agent}") + + # Add client to manager with user agent client_manager = proxy_server.client_managers[channel_id] - client_count = client_manager.add_client(client_id) + client_count = client_manager.add_client(client_id, user_agent) # If this is the first client, try to acquire ownership if client_count == 1 and not proxy_server.am_i_owner(channel_id):