def start_tick_collection(self, tokens: List[int], duration: Optional[int] = 15): """ Fine-tuned WebSocket Handler Args: tokens: List of instrument tokens duration: If None, stays open for trading. If int, collects data for X seconds. """ if not self.access_token: logger.error("❌ Cannot start WebSocket: No Access Token") return [] if hasattr(self, 'kws') and self.kws is not None: try: logger.info("Closing previous KiteTicker before creating new one...") self.kws.close() except Exception as e: logger.debug(f"Error closing previous kws: {e}") self.kws = None #self.kws = KiteTicker(self.api_key, self.access_token) self.kws = KiteTicker( self.api_key, self.access_token, reconnect = False, #reconnect_max_tries = 5, # ← was default 300 #reconnect_max_delay = 10, # ← was default 60s connect_timeout = 30, ) collected_ticks = [] collection_complete = threading.Event() # Initialize your threading event for live mode tracking if not hasattr(self, '_ws_connected_event'): self._ws_connected_event = threading.Event() self._ws_connected_event.clear() # Reset flag state before starting connection # Initialize Managers for t in tokens: if t not in self.candle_managers: self.candle_managers[t] = CandleManager(interval="5min") def on_ticks(ws, ticks): # ccv = self._ccv for tick in ticks: token = tick['instrument_token'] #ccv.on_tick(tick) if ccv else None v_traded = tick.get('last_quantity', 0) # ✅ Feed every tick to external callback (CCV, TickVWAP, etc.) if self.tick_callback: try: self.tick_callback(tick) except Exception as e: pass # never let a callback crash the WebSocket # Store for collection mode if duration is not None: tick_data = { 'instrument_token': tick.get('instrument_token'), 'lp': tick.get('last_price', 0), 'volume': tick.get('volume_traded', tick.get('volume', 0)), # cumulative day volume 'oi': tick.get('oi', 0), 'timestamp': datetime.now(IST).isoformat() } # print(f"Collected Tick: {tick_data}") collected_ticks.append(tick_data) # Update Live Candle Managers if token in self.candle_managers: self.candle_managers[token].add_tick(tick) df = self.candle_managers[token].get_candles() if not df.empty: last_candle = df.iloc[-1] # You can trigger your strategy here using last_candle def on_connect(ws, response): logger.info(f"📡 WebSocket connected. Subscribing to {len(tokens)} tokens.") ws.subscribe(tokens) ws.set_mode(ws.MODE_FULL, tokens) self._ws_connected_event.set() def on_error(ws, code, reason): logger.error(f"🔴 {code} - {reason}".format(code=code, reason=reason)) def on_close(ws, code, reason): logger.warning(f"Connection closed: {code} - {reason}".format(code=code, reason=reason)) collection_complete.set() self._ws_connected_event.clear() # reset so next call retries # Assign Callbacks self.kws.on_ticks = on_ticks self.kws.on_connect = on_connect self.kws.on_error = on_error self.kws.on_close = on_close # 🚀 Start Connection self.kws.connect(threaded=True) # Logic for Data Collection Mode (Reference Points) if duration: # logger.info(f"⏳ Collecting reference ticks for {duration}s...") time.sleep(duration) try: self.kws.close() except Exception: pass time.sleep(2) # give Zerodha server time to release the session self.kws = None # ← clear reference so ensure_websocket creates fresh one return collected_ticks # Logic for Live Trading Mode else: connection_established = self._ws_connected_event.wait(timeout=30.0) if connection_established: logger.info("✅ Connection established safely. Handing control back to main thread.") return [] else: logger.warning("❌ WebSocket connection timeout in start_tick_collection (Handshake took too long)") try: self.kws.close() # Stop background retry thread to prevent memory leakage except Exception: pass return [] # logger.info("🚀 WebSocket running in LIVE MODE (Persistent)") # We wait a few seconds to ensure connection is stable # Wait up to 15s for stable connection (30 × 0.5s) # timeout = 120 # while timeout > 0: # try: # if self.kws.is_connected(): # return [] # except Exception: # pass # time.sleep(0.5) # timeout -= 1 # logger.warning("WebSocket connection timeout in start_tick_collection") # try: # self.kws.close() # ← stop retry thread, prevent zombie pileup # except Exception: # pass # return [] def ensure_websocket(self, tokens: List[int]) -> None: """ Ensure a live persistent WebSocket is running and subscribed to tokens. Safe to call multiple times — reuses the connection if already up. """ try: if hasattr(self, 'kws') and self.kws and self.kws.is_connected(): self.kws.subscribe(tokens) self.kws.set_mode(self.kws.MODE_FULL, tokens) # logger.info(f"WebSocket already connected — subscribed to tokens {tokens}") return except Exception: pass # ── Already starting — just wait for it ────────────────────────────── with self._ws_lock: if self._ws_starting: logger.info("WebSocket already started — waiting...") connected = self._ws_connected_event.wait(timeout=30.0) if connected: try: self.kws.subscribe(tokens) self.kws.set_mode(self.kws.MODE_FULL, tokens) except Exception: pass return else: logger.warning("⚠️ The parallel initialization attempt timed out.") return # ── Start new connection ────────────────────────────────────── self._ws_starting = True self._ws_connected_event.clear() def _start(): try: self.start_tick_collection(tokens, duration=None) except Exception as e: logger.error(f"WebSocket thread error: {e}") finally: with self._ws_lock: self._ws_starting = False threading.Thread( target = _start, daemon = True, name = "live_ws_ccv", ).start() # ── Wait up to 15s for connection ───────────────────────────────────── logger.info("📡 Waiting up to 30s for background WebSocket handshake...") connected = self._ws_connected_event.wait(timeout=30.0) # 5. Fallback verification step if event flag processing got interrupted if not connected: logger.warning("⚠️ Event flag not set within 30s. Checking socket handle connection state directly...") for _ in range(6): try: if hasattr(self, 'kws') and self.kws and self.kws.is_connected(): connected = True self._ws_connected_event.set() break except Exception: pass time.sleep(1.0) if connected: logger.info(f"🚀 WebSocket successfully stabilized and feeding ticks for tokens: {tokens}") else: logger.error( "❌ WebSocket completely failed to connect on the server. " "Falling back to polling mode for strategy updates." ) try: if hasattr(self, 'kws') and self.kws: self.kws.close() # Prevent dead socket leak except Exception: pass