From d2085d57f84966bd0d1ffd42f7ed9965a3fa1749 Mon Sep 17 00:00:00 2001 From: SergeantPanda Date: Tue, 16 Sep 2025 12:43:21 -0500 Subject: [PATCH] Add sentence transformers to new matching function. --- apps/channels/tasks.py | 697 +++++++++++++++++++++++++++++------------ 1 file changed, 500 insertions(+), 197 deletions(-) diff --git a/apps/channels/tasks.py b/apps/channels/tasks.py index f67fe563..1619843f 100755 --- a/apps/channels/tasks.py +++ b/apps/channels/tasks.py @@ -28,6 +28,113 @@ from urllib.parse import quote logger = logging.getLogger(__name__) +# Lazy loading for ML models - only imported/loaded when needed +_ml_model_cache = { + 'sentence_transformer': None, + 'model_path': os.path.join("/data", "models", "all-MiniLM-L6-v2"), # Use /data for persistence + 'model_name': "sentence-transformers/all-MiniLM-L6-v2" +} + +def get_sentence_transformer(): + """Lazy load the sentence transformer model only when needed""" + if _ml_model_cache['sentence_transformer'] is None: + try: + from sentence_transformers import SentenceTransformer + from sentence_transformers import util + + model_path = _ml_model_cache['model_path'] + model_name = _ml_model_cache['model_name'] + cache_dir = os.path.dirname(model_path) # /data/models + + # Check environment variable to disable downloads + disable_downloads = os.environ.get('DISABLE_ML_DOWNLOADS', 'false').lower() == 'true' + + # Ensure directory exists and is writable + os.makedirs(cache_dir, exist_ok=True) + + # Debug: List what's actually in the cache directory + try: + if os.path.exists(cache_dir): + logger.info(f"Cache directory contents: {os.listdir(cache_dir)}") + for item in os.listdir(cache_dir): + item_path = os.path.join(cache_dir, item) + if os.path.isdir(item_path): + logger.info(f" Subdirectory '{item}' contains: {os.listdir(item_path)}") + except Exception as e: + logger.info(f"Could not list cache directory: {e}") + + # Check if model files exist in our expected location + config_path = os.path.join(model_path, "config.json") + + logger.info(f"Checking for cached model at {model_path}") + logger.info(f"Config exists: {os.path.exists(config_path)}") + + # Also check if the model exists in the sentence-transformers default naming convention + alt_model_name = model_name.replace("/", "_") + alt_model_path = os.path.join(cache_dir, alt_model_name) + alt_config_path = os.path.join(alt_model_path, "config.json") + logger.info(f"Alternative path check - {alt_model_path}, config exists: {os.path.exists(alt_config_path)}") + + # Check for Hugging Face Hub cache format (newer format) + hf_model_name = f"models--{model_name.replace('/', '--')}" + hf_model_path = os.path.join(cache_dir, hf_model_name) + hf_snapshots_path = os.path.join(hf_model_path, "snapshots") + + logger.info(f"Hugging Face cache path check - {hf_model_path}, snapshots exists: {os.path.exists(hf_snapshots_path)}") + + # If HF cache exists, find the latest snapshot + hf_config_exists = False + hf_snapshot_path = None + if os.path.exists(hf_snapshots_path): + try: + snapshots = os.listdir(hf_snapshots_path) + if snapshots: + # Use the first (and likely only) snapshot + hf_snapshot_path = os.path.join(hf_snapshots_path, snapshots[0]) + hf_config_path = os.path.join(hf_snapshot_path, "config.json") + hf_config_exists = os.path.exists(hf_config_path) + logger.info(f"HF snapshot path: {hf_snapshot_path}, config exists: {hf_config_exists}") + except Exception as e: + logger.info(f"Error checking HF cache: {e}") + + # First try to load from our specific path + if os.path.exists(config_path): + logger.info(f"Loading cached sentence transformer from {model_path}") + _ml_model_cache['sentence_transformer'] = SentenceTransformer(model_path) + elif os.path.exists(alt_config_path): + logger.info(f"Loading cached sentence transformer from alternative path {alt_model_path}") + _ml_model_cache['sentence_transformer'] = SentenceTransformer(alt_model_path) + elif hf_config_exists and hf_snapshot_path: + logger.info(f"Loading cached sentence transformer from HF cache {hf_snapshot_path}") + _ml_model_cache['sentence_transformer'] = SentenceTransformer(hf_snapshot_path) + elif disable_downloads: + logger.warning(f"ML model not found and downloads disabled (DISABLE_ML_DOWNLOADS=true). Skipping ML matching.") + return None, None + else: + logger.info(f"Model cache not found, downloading {model_name}") + # Let sentence-transformers handle the download with its cache folder + _ml_model_cache['sentence_transformer'] = SentenceTransformer( + model_name, + cache_folder=cache_dir + ) + logger.info(f"Model downloaded and loaded successfully") + + return _ml_model_cache['sentence_transformer'], util + except ImportError: + logger.warning("sentence-transformers not available - ML-enhanced matching disabled") + return None, None + except Exception as e: + logger.error(f"Failed to load sentence transformer: {e}") + return None, None + else: + from sentence_transformers import util + return _ml_model_cache['sentence_transformer'], util + +# ML matching thresholds (same as original script) +BEST_FUZZY_THRESHOLD = 85 +LOWER_FUZZY_THRESHOLD = 40 +EMBED_SIM_THRESHOLD = 0.65 + # Words we remove to help with fuzzy + embedding matching COMMON_EXTRANEOUS_WORDS = [ "tv", "channel", "network", "television", @@ -50,155 +157,373 @@ def normalize_name(name: str) -> str: norm = name.lower() norm = re.sub(r"\[.*?\]", "", norm) + + # Extract and preserve important call signs from parentheses before removing them + # This captures call signs like (KVLY), (KING), (KARE), etc. + call_sign_match = re.search(r"\(([A-Z]{3,5})\)", name) + preserved_call_sign = "" + if call_sign_match: + preserved_call_sign = " " + call_sign_match.group(1).lower() + + # Now remove all parentheses content norm = re.sub(r"\(.*?\)", "", norm) + + # Add back the preserved call sign + norm = norm + preserved_call_sign + norm = re.sub(r"[^\w\s]", "", norm) tokens = norm.split() tokens = [t for t in tokens if t not in COMMON_EXTRANEOUS_WORDS] norm = " ".join(tokens).strip() return norm +def match_channels_to_epg(channels_data, epg_data, region_code=None, use_ml=True): + """ + EPG matching logic that finds the best EPG matches for channels using + multiple matching strategies including fuzzy matching and ML models. + """ + channels_to_update = [] + matched_channels = [] + + # Try to get ML models if requested (but don't load yet - lazy loading) + st_model, util = None, None + epg_embeddings = None + ml_available = use_ml + + # Process each channel + for chan in channels_data: + normalized_tvg_id = chan.get("tvg_id", "") + fallback_name = chan["tvg_id"].strip() if chan["tvg_id"] else chan["name"] + + # Step 1: Exact TVG ID match + epg_by_tvg_id = next((epg for epg in epg_data if epg["tvg_id"] == normalized_tvg_id), None) + if normalized_tvg_id and epg_by_tvg_id: + chan["epg_data_id"] = epg_by_tvg_id["id"] + channels_to_update.append(chan) + matched_channels.append((chan['id'], fallback_name, epg_by_tvg_id["tvg_id"])) + logger.info(f"Channel {chan['id']} '{fallback_name}' => EPG found by exact tvg_id={epg_by_tvg_id['tvg_id']}") + continue + + # Step 2: Secondary TVG ID check (legacy compatibility) + if chan["tvg_id"]: + epg_match = [epg["id"] for epg in epg_data if epg["tvg_id"] == chan["tvg_id"]] + if epg_match: + chan["epg_data_id"] = epg_match[0] + channels_to_update.append(chan) + matched_channels.append((chan['id'], fallback_name, chan["tvg_id"])) + logger.info(f"Channel {chan['id']} '{chan['name']}' => EPG found by secondary tvg_id={chan['tvg_id']}") + continue + + # Step 3: Name-based fuzzy matching + if not chan["norm_chan"]: + logger.debug(f"Channel {chan['id']} '{chan['name']}' => empty after normalization, skipping") + continue + + best_score = 0 + best_epg = None + + # Debug: show what we're matching against + logger.debug(f"Fuzzy matching '{chan['norm_chan']}' against EPG entries...") + + # Find best fuzzy match + for row in epg_data: + if not row.get("norm_name"): + continue + + base_score = fuzz.ratio(chan["norm_chan"], row["norm_name"]) + bonus = 0 + + # Apply region-based bonus/penalty + if region_code and row.get("tvg_id"): + combined_text = row["tvg_id"].lower() + " " + row["name"].lower() + dot_regions = re.findall(r'\.([a-z]{2})', combined_text) + + if dot_regions: + if region_code in dot_regions: + bonus = 30 # Bigger bonus for matching region + else: + bonus = -15 # Penalty for different region + elif region_code in combined_text: + bonus = 15 + + score = base_score + bonus + + # Debug the best few matches + if score > 50: # Only show decent matches + logger.debug(f" EPG '{row['name']}' (norm: '{row['norm_name']}') => score: {score} (base: {base_score}, bonus: {bonus})") + + if score > best_score: + best_score = score + best_epg = row + + # Log the best score we found + if best_epg: + logger.info(f"Channel {chan['id']} '{chan['name']}' => best match: '{best_epg['name']}' (score: {best_score})") + + # Debug: Show some other potential matches for analysis + def score_epg_entry(epg_row): + base_score = fuzz.ratio(chan["norm_chan"], epg_row.get("norm_name", "")) + bonus = 0 + if region_code and epg_row.get("tvg_id"): + combined_text = epg_row["tvg_id"].lower() + " " + epg_row["name"].lower() + dot_regions = re.findall(r'\.([a-z]{2})', combined_text) + if dot_regions: + if region_code in dot_regions: + bonus = 30 + else: + bonus = -15 + elif region_code in combined_text: + bonus = 15 + return base_score + bonus + + # Check specifically for entries matching the channel's call sign or name parts + channel_keywords = chan["norm_chan"].split() + potential_matches = [] + for keyword in channel_keywords: + if len(keyword) >= 3: # Only check meaningful keywords + matching_entries = [row for row in epg_data if keyword.lower() in row['name'].lower() or keyword.lower() in row['tvg_id'].lower()] + potential_matches.extend(matching_entries) + + # Remove duplicates + unique_matches = [] + seen_ids = set() + for match in potential_matches: + if match['tvg_id'] not in seen_ids: + seen_ids.add(match['tvg_id']) + unique_matches.append(match) + + if unique_matches: + logger.info(f"Found {len(unique_matches)} entries containing channel keywords {channel_keywords}:") + for match_row in unique_matches: + match_score = score_epg_entry(match_row) + # Show original name vs normalized name to debug normalization + logger.info(f" Match: '{match_row['name']}' → normalized: '{match_row.get('norm_name', 'MISSING')}' (tvg_id: {match_row['tvg_id']}) => score: {match_score}") + else: + logger.warning(f"No entries found containing any of the channel keywords: {channel_keywords}") + + sorted_scores = sorted([(row, score_epg_entry(row)) for row in epg_data if row.get("norm_name") and score_epg_entry(row) > 20], key=lambda x: x[1], reverse=True) + + # Remove duplicates based on tvg_id + seen_tvg_ids = set() + unique_sorted_scores = [] + for row, score in sorted_scores: + if row['tvg_id'] not in seen_tvg_ids: + seen_tvg_ids.add(row['tvg_id']) + unique_sorted_scores.append((row, score)) + + logger.debug(f"Channel {chan['id']} '{chan['name']}' => top 10 unique fuzzy matches:") + for i, (epg_row, score) in enumerate(unique_sorted_scores[:10]): + # Highlight entries that contain any of the channel's keywords + channel_keywords = chan["norm_chan"].split() + is_keyword_match = any(keyword in epg_row['name'].lower() or keyword in epg_row['tvg_id'].lower() for keyword in channel_keywords if len(keyword) >= 3) + + if is_keyword_match: + logger.info(f" {i+1}. 🎯 KEYWORD MATCH: '{epg_row['name']}' (tvg_id: {epg_row['tvg_id']}) => score: {score} (norm_name: '{epg_row.get('norm_name', 'MISSING')}')") + else: + logger.debug(f" {i+1}. '{epg_row['name']}' (tvg_id: {epg_row['tvg_id']}) => score: {score}") + else: + logger.debug(f"Channel {chan['id']} '{chan['name']}' => no EPG entries with valid norm_name found") + continue + + # High confidence match - accept immediately + if best_score >= BEST_FUZZY_THRESHOLD: + chan["epg_data_id"] = best_epg["id"] + channels_to_update.append(chan) + matched_channels.append((chan['id'], chan['name'], best_epg["tvg_id"])) + logger.info(f"Channel {chan['id']} '{chan['name']}' => matched tvg_id={best_epg['tvg_id']} (score={best_score})") + + # Medium confidence - use ML if available (lazy load models here) + elif best_score >= LOWER_FUZZY_THRESHOLD and ml_available: + logger.debug(f"Channel {chan['id']} entering ML matching path at {time.time()}") + + # Note: If experiencing 5+ second delays here, check if Celery Beat + # task 'scan_and_process_files' is running too frequently and blocking execution + + # Lazy load ML models only when we actually need them + if st_model is None: + logger.debug(f"Channel {chan['id']} about to load ML model at {time.time()}") + st_model, util = get_sentence_transformer() + logger.debug(f"Channel {chan['id']} finished loading ML model at {time.time()}") + + # Lazy generate embeddings only when we actually need them + if epg_embeddings is None and st_model and any(row.get("norm_name") for row in epg_data): + try: + logger.info("Generating embeddings for EPG data using ML model (lazy loading)") + epg_embeddings = st_model.encode( + [row["norm_name"] for row in epg_data if row.get("norm_name")], + convert_to_tensor=True + ) + except Exception as e: + logger.warning(f"Failed to generate embeddings: {e}") + epg_embeddings = None + + if epg_embeddings is not None and st_model: + try: + # Generate embedding for this channel + chan_embedding = st_model.encode(chan["norm_chan"], convert_to_tensor=True) + + # Calculate similarity with all EPG embeddings + sim_scores = util.cos_sim(chan_embedding, epg_embeddings)[0] + top_index = int(sim_scores.argmax()) + top_value = float(sim_scores[top_index]) + + if top_value >= EMBED_SIM_THRESHOLD: + # Find the EPG entry that corresponds to this embedding index + epg_with_names = [epg for epg in epg_data if epg.get("norm_name")] + matched_epg = epg_with_names[top_index] + + chan["epg_data_id"] = matched_epg["id"] + channels_to_update.append(chan) + matched_channels.append((chan['id'], chan['name'], matched_epg["tvg_id"])) + logger.info(f"Channel {chan['id']} '{chan['name']}' => matched EPG tvg_id={matched_epg['tvg_id']} (fuzzy={best_score}, ML-sim={top_value:.2f})") + else: + logger.info(f"Channel {chan['id']} '{chan['name']}' => fuzzy={best_score}, ML-sim={top_value:.2f} < {EMBED_SIM_THRESHOLD}, trying last resort...") + + # Last resort: try ML with very low fuzzy threshold + if top_value >= 0.45: # Lower ML threshold as last resort + epg_with_names = [epg for epg in epg_data if epg.get("norm_name")] + matched_epg = epg_with_names[top_index] + + chan["epg_data_id"] = matched_epg["id"] + channels_to_update.append(chan) + matched_channels.append((chan['id'], chan['name'], matched_epg["tvg_id"])) + logger.info(f"Channel {chan['id']} '{chan['name']}' => LAST RESORT match EPG tvg_id={matched_epg['tvg_id']} (fuzzy={best_score}, ML-sim={top_value:.2f})") + else: + logger.info(f"Channel {chan['id']} '{chan['name']}' => even last resort ML-sim {top_value:.2f} < 0.45, skipping") + + except Exception as e: + logger.warning(f"ML matching failed for channel {chan['id']}: {e}") + # Fall back to non-ML decision + logger.info(f"Channel {chan['id']} '{chan['name']}' => fuzzy score {best_score} below threshold, skipping") + + # Last resort: Try ML matching even with very low fuzzy scores + elif best_score >= 20 and ml_available: + logger.debug(f"Channel {chan['id']} entering last resort ML matching at {time.time()}") + + # Lazy load ML models for last resort attempts + if st_model is None: + logger.debug(f"Channel {chan['id']} loading ML model for last resort at {time.time()}") + st_model, util = get_sentence_transformer() + logger.debug(f"Channel {chan['id']} finished loading ML model for last resort at {time.time()}") + + # Lazy generate embeddings for last resort attempts + if epg_embeddings is None and st_model and any(row.get("norm_name") for row in epg_data): + try: + logger.info("Generating embeddings for EPG data using ML model (last resort lazy loading)") + epg_embeddings = st_model.encode( + [row["norm_name"] for row in epg_data if row.get("norm_name")], + convert_to_tensor=True + ) + except Exception as e: + logger.warning(f"Failed to generate embeddings for last resort: {e}") + epg_embeddings = None + + if epg_embeddings is not None and st_model: + try: + logger.info(f"Channel {chan['id']} '{chan['name']}' => trying ML as last resort (fuzzy={best_score})") + # Generate embedding for this channel + chan_embedding = st_model.encode(chan["norm_chan"], convert_to_tensor=True) + + # Calculate similarity with all EPG embeddings + sim_scores = util.cos_sim(chan_embedding, epg_embeddings)[0] + top_index = int(sim_scores.argmax()) + top_value = float(sim_scores[top_index]) + + if top_value >= 0.50: # Even lower threshold for last resort + # Find the EPG entry that corresponds to this embedding index + epg_with_names = [epg for epg in epg_data if epg.get("norm_name")] + matched_epg = epg_with_names[top_index] + + chan["epg_data_id"] = matched_epg["id"] + channels_to_update.append(chan) + matched_channels.append((chan['id'], chan['name'], matched_epg["tvg_id"])) + logger.info(f"Channel {chan['id']} '{chan['name']}' => DESPERATE LAST RESORT match EPG tvg_id={matched_epg['tvg_id']} (fuzzy={best_score}, ML-sim={top_value:.2f})") + else: + logger.info(f"Channel {chan['id']} '{chan['name']}' => desperate last resort ML-sim {top_value:.2f} < 0.50, giving up") + except Exception as e: + logger.warning(f"Last resort ML matching failed for channel {chan['id']}: {e}") + logger.info(f"Channel {chan['id']} '{chan['name']}' => best fuzzy score={best_score} < {LOWER_FUZZY_THRESHOLD}, giving up") + else: + # No ML available or very low fuzzy score + logger.info(f"Channel {chan['id']} '{chan['name']}' => best fuzzy score={best_score} < {LOWER_FUZZY_THRESHOLD}, no ML fallback available") + + return { + "channels_to_update": channels_to_update, + "matched_channels": matched_channels + } + @shared_task def match_epg_channels(): """ - Goes through all Channels and tries to find a matching EPGData row by: - 1) If channel.tvg_id is valid in EPGData, skip. - 2) If channel has a tvg_id but not found in EPGData, attempt direct EPGData lookup. - 3) Otherwise, perform name-based fuzzy matching with optional region-based bonus. - 4) If a match is found, we set channel.tvg_id - 5) Summarize and log results. + Uses integrated EPG matching instead of external script. + Provides the same functionality with better performance and maintainability. """ try: - logger.info("Starting EPG matching logic...") + logger.info("Starting integrated EPG matching...") - # Attempt to retrieve a "preferred-region" if configured + # Get region preference try: region_obj = CoreSettings.objects.get(key="preferred-region") region_code = region_obj.value.strip().lower() except CoreSettings.DoesNotExist: region_code = None - matched_channels = [] - channels_to_update = [] - # Get channels that don't have EPG data assigned channels_without_epg = Channel.objects.filter(epg_data__isnull=True) logger.info(f"Found {channels_without_epg.count()} channels without EPG data") - channels_json = [] + channels_data = [] for channel in channels_without_epg: - # Normalize TVG ID - strip whitespace and convert to lowercase normalized_tvg_id = channel.tvg_id.strip().lower() if channel.tvg_id else "" - if normalized_tvg_id: - logger.info(f"Processing channel {channel.id} '{channel.name}' with TVG ID='{normalized_tvg_id}'") - - channels_json.append({ + channels_data.append({ "id": channel.id, "name": channel.name, - "tvg_id": normalized_tvg_id, # Use normalized TVG ID - "original_tvg_id": channel.tvg_id, # Keep original for reference + "tvg_id": normalized_tvg_id, + "original_tvg_id": channel.tvg_id, "fallback_name": normalized_tvg_id if normalized_tvg_id else channel.name, - "norm_chan": normalize_name(normalized_tvg_id if normalized_tvg_id else channel.name) + "norm_chan": normalize_name(channel.name) # Always use channel name for fuzzy matching! }) - # Similarly normalize EPG data TVG IDs - epg_json = [] + # Get all EPG data + epg_data = [] for epg in EPGData.objects.all(): normalized_tvg_id = epg.tvg_id.strip().lower() if epg.tvg_id else "" - epg_json.append({ + epg_data.append({ 'id': epg.id, - 'tvg_id': normalized_tvg_id, # Use normalized TVG ID - 'original_tvg_id': epg.tvg_id, # Keep original for reference + 'tvg_id': normalized_tvg_id, + 'original_tvg_id': epg.tvg_id, 'name': epg.name, 'norm_name': normalize_name(epg.name), 'epg_source_id': epg.epg_source.id if epg.epg_source else None, }) - # Log available EPG data TVG IDs for debugging - unique_epg_tvg_ids = set(e['tvg_id'] for e in epg_json if e['tvg_id']) - logger.info(f"Available EPG TVG IDs: {', '.join(sorted(unique_epg_tvg_ids))}") + logger.info(f"Processing {len(channels_data)} channels against {len(epg_data)} EPG entries") - payload = { - "channels": channels_json, - "epg_data": epg_json, - "region_code": region_code, - } - - with tempfile.NamedTemporaryFile(delete=False) as temp_file: - temp_file.write(json.dumps(payload).encode('utf-8')) - temp_file_path = temp_file.name - - # After writing to the file but before subprocess - # Explicitly delete the large data structures - del payload - gc.collect() - - process = subprocess.Popen( - ['python', '/app/scripts/epg_match.py', temp_file_path], - stdout=subprocess.PIPE, - stderr=subprocess.PIPE, - text=True - ) - - stdout = '' - block_size = 1024 - - while True: - # Monitor stdout and stderr for readability - readable, _, _ = select.select([process.stdout, process.stderr], [], [], 1) # timeout of 1 second - - if not readable: # timeout expired - if process.poll() is not None: # check if process finished - break - else: # process still running, continue - continue - - for stream in readable: - if stream == process.stdout: - stdout += stream.read(block_size) - elif stream == process.stderr: - error = stream.readline() - if error: - logger.info(error.strip()) - - if process.poll() is not None: - break - - process.wait() - os.remove(temp_file_path) - - if process.returncode != 0: - return f"Failed to process EPG matching" - - result = json.loads(stdout) - # This returns lists of dicts, not model objects + # Run EPG matching + result = match_channels_to_epg(channels_data, epg_data, region_code, use_ml=True) channels_to_update_dicts = result["channels_to_update"] matched_channels = result["matched_channels"] - # Explicitly clean up large objects - del stdout, result - gc.collect() - - # Convert your dict-based 'channels_to_update' into real Channel objects + # Update channels in database if channels_to_update_dicts: - # Extract IDs of the channels that need updates channel_ids = [d["id"] for d in channels_to_update_dicts] - - # Fetch them from DB channels_qs = Channel.objects.filter(id__in=channel_ids) channels_list = list(channels_qs) - # Build a map from channel_id -> epg_data_id (or whatever fields you need) - epg_mapping = { - d["id"]: d["epg_data_id"] for d in channels_to_update_dicts - } + # Create mapping from channel_id to epg_data_id + epg_mapping = {d["id"]: d["epg_data_id"] for d in channels_to_update_dicts} - # Populate each Channel object with the updated epg_data_id + # Update each channel with matched EPG data for channel_obj in channels_list: - # The script sets 'epg_data_id' in the returned dict - # We either assign directly, or fetch the EPGData instance if needed. - channel_obj.epg_data_id = epg_mapping.get(channel_obj.id) + epg_data_id = epg_mapping.get(channel_obj.id) + if epg_data_id: + try: + epg_data_obj = EPGData.objects.get(id=epg_data_id) + channel_obj.epg_data = epg_data_obj + except EPGData.DoesNotExist: + logger.error(f"EPG data {epg_data_id} not found for channel {channel_obj.id}") - # Now we have real model objects, so bulk_update will work + # Bulk update all channels Channel.objects.bulk_update(channels_list, ["epg_data"]) total_matched = len(matched_channels) @@ -209,9 +534,9 @@ def match_epg_channels(): else: logger.info("No new channels were matched.") - logger.info("Finished EPG matching logic.") + logger.info("Finished integrated EPG matching.") - # Send update with additional information for refreshing UI + # Send WebSocket update channel_layer = get_channel_layer() associations = [ {"channel_id": chan["id"], "epg_data_id": chan["epg_data_id"]} @@ -225,19 +550,19 @@ def match_epg_channels(): "data": { "success": True, "type": "epg_match", - "refresh_channels": True, # Flag to tell frontend to refresh channels + "refresh_channels": True, "matches_count": total_matched, "message": f"EPG matching complete: {total_matched} channel(s) matched", - "associations": associations # Add the associations data + "associations": associations } } ) return f"Done. Matched {total_matched} channel(s)." + finally: - # Final cleanup + # Memory cleanup gc.collect() - # Use our standardized cleanup function for more thorough memory management from core.utils import cleanup_memory cleanup_memory(log_usage=True, force_collection=True) @@ -245,16 +570,14 @@ def match_epg_channels(): @shared_task def match_single_channel_epg(channel_id): """ - Try to match a single channel with EPG data using optimized logic - that doesn't require loading all EPG data or running the external script. - Returns a dict with match status and message. + Try to match a single channel with EPG data using the integrated matching logic + that includes both fuzzy and ML-enhanced matching. Returns a dict with match status and message. """ try: from apps.channels.models import Channel from apps.epg.models import EPGData - import re - logger.info(f"Starting optimized single channel EPG matching for channel ID {channel_id}") + logger.info(f"Starting integrated single channel EPG matching for channel ID {channel_id}") # Get the channel try: @@ -266,112 +589,92 @@ def match_single_channel_epg(channel_id): if channel.epg_data: return {"matched": False, "message": f"Channel '{channel.name}' already has EPG data assigned"} - # Get region preference - try: - region_obj = CoreSettings.objects.get(key="preferred-region") - region_code = region_obj.value.strip().lower() - except CoreSettings.DoesNotExist: - region_code = None - - # Prepare channel data + # Prepare single channel data for matching (same format as bulk matching) normalized_tvg_id = channel.tvg_id.strip().lower() if channel.tvg_id else "" - normalized_channel_name = normalize_name(channel.name) + channel_data = { + "id": channel.id, + "name": channel.name, + "tvg_id": normalized_tvg_id, + "original_tvg_id": channel.tvg_id, + "fallback_name": normalized_tvg_id if normalized_tvg_id else channel.name, + "norm_chan": normalize_name(channel.name) # Always use channel name for fuzzy matching! + } - logger.info(f"Matching channel '{channel.name}' (TVG ID: '{channel.tvg_id}') against EPG data") + logger.info(f"Channel data prepared: name='{channel.name}', tvg_id='{normalized_tvg_id}', norm_chan='{channel_data['norm_chan']}'") - # Step 1: Try exact TVG ID match first (most efficient) - if normalized_tvg_id: - epg_exact_match = EPGData.objects.filter(tvg_id__iexact=channel.tvg_id).first() - if epg_exact_match: - logger.info(f"Channel '{channel.name}' matched with EPG '{epg_exact_match.name}' by exact TVG ID match") - channel.epg_data = epg_exact_match - channel.save(update_fields=["epg_data"]) - return { - "matched": True, - "message": f"Channel '{channel.name}' matched with EPG '{epg_exact_match.name}' by exact TVG ID match" - } + # Debug: Test what the normalization does to preserve call signs + test_name = "NBC 11 (KVLY) - Fargo" # Example for testing + test_normalized = normalize_name(test_name) + logger.debug(f"DEBUG normalization example: '{test_name}' → '{test_normalized}' (call sign preserved)") - # Step 2: Try case-insensitive TVG ID match - if normalized_tvg_id: - epg_case_match = EPGData.objects.filter(tvg_id__icontains=normalized_tvg_id).first() - if epg_case_match: - logger.info(f"Channel '{channel.name}' matched with EPG '{epg_case_match.name}' by case-insensitive TVG ID match") - channel.epg_data = epg_case_match - channel.save(update_fields=["epg_data"]) - return { - "matched": True, - "message": f"Channel '{channel.name}' matched with EPG '{epg_case_match.name}' by case-insensitive TVG ID match" - } + # Get all EPG data for matching - must include norm_name field + epg_data_list = [] + for epg in EPGData.objects.filter(name__isnull=False).exclude(name=''): + normalized_epg_tvg_id = epg.tvg_id.strip().lower() if epg.tvg_id else "" + epg_data_list.append({ + 'id': epg.id, + 'tvg_id': normalized_epg_tvg_id, + 'original_tvg_id': epg.tvg_id, + 'name': epg.name, + 'norm_name': normalize_name(epg.name), + 'epg_source_id': epg.epg_source.id if epg.epg_source else None, + }) - # Step 3: Fuzzy name matching (only if name-based matching is needed) - if not normalized_channel_name: - return {"matched": False, "message": f"Channel '{channel.name}' has no usable name for matching"} + if not epg_data_list: + return {"matched": False, "message": "No EPG data available for matching"} - # Query EPG data with name filtering to reduce dataset - epg_candidates = EPGData.objects.filter(name__isnull=False).exclude(name='').values('id', 'name', 'tvg_id') - epg_count = epg_candidates.count() - logger.info(f"Fuzzy matching against {epg_count} EPG entries (optimized - not loading all EPG data)") + logger.info(f"Matching single channel '{channel.name}' against {len(epg_data_list)} EPG entries") - best_score = 0 - best_epg = None + # Use the EPG matching function + result = match_channels_to_epg([channel_data], epg_data_list) + channels_to_update = result.get("channels_to_update", []) + matched_channels = result.get("matched_channels", []) - for epg in epg_candidates: - if not epg['name']: - continue + if channels_to_update: + # Find our channel in the results + channel_match = None + for update in channels_to_update: + if update["id"] == channel.id: + channel_match = update + break - epg_normalized_name = normalize_name(epg['name']) - if not epg_normalized_name: - continue + if channel_match: + # Apply the match to the channel + try: + epg_data = EPGData.objects.get(id=channel_match['epg_data_id']) + channel.epg_data = epg_data + channel.save(update_fields=["epg_data"]) - # Calculate base fuzzy score - base_score = fuzz.ratio(normalized_channel_name, epg_normalized_name) - bonus = 0 + # Find match details from matched_channels for better reporting + match_details = None + for match_info in matched_channels: + if match_info[0] == channel.id: # matched_channels format: (channel_id, channel_name, epg_info) + match_details = match_info + break - # Apply region-based bonus/penalty if applicable - if region_code and epg['tvg_id']: - combined_text = epg['tvg_id'].lower() + " " + epg['name'].lower() - dot_regions = re.findall(r'\.([a-z]{2})', combined_text) + success_msg = f"Channel '{channel.name}' matched with EPG '{epg_data.name}'" + if match_details: + success_msg += f" (matched via: {match_details[2]})" - if dot_regions: - if region_code in dot_regions: - bonus = 30 # Bigger bonus for matching region - else: - bonus = -15 # Penalty for different region - elif region_code in combined_text: - bonus = 15 + logger.info(success_msg) - final_score = base_score + bonus + return { + "matched": True, + "message": success_msg, + "epg_name": epg_data.name, + "epg_id": epg_data.id + } + except EPGData.DoesNotExist: + return {"matched": False, "message": "Matched EPG data not found"} - if final_score > best_score: - best_score = final_score - best_epg = epg - - # Apply matching thresholds (same as the ML script) - BEST_FUZZY_THRESHOLD = 85 - - if best_epg and best_score >= BEST_FUZZY_THRESHOLD: - try: - logger.info(f"Channel '{channel.name}' matched with EPG '{best_epg['name']}' (score: {best_score})") - epg_data = EPGData.objects.get(id=best_epg['id']) - channel.epg_data = epg_data - channel.save(update_fields=["epg_data"]) - - return { - "matched": True, - "message": f"Channel '{channel.name}' matched with EPG '{epg_data.name}' (score: {best_score})" - } - except EPGData.DoesNotExist: - return {"matched": False, "message": "Matched EPG data not found"} - - # No good match found - logger.info(f"No suitable EPG match found for channel '{channel.name}' (best score: {best_score})") + # No match found return { "matched": False, - "message": f"No suitable EPG match found for channel '{channel.name}' (best score: {best_score})" + "message": f"No suitable EPG match found for channel '{channel.name}'" } except Exception as e: - logger.error(f"Error in optimized single channel EPG matching: {e}", exc_info=True) + logger.error(f"Error in integrated single channel EPG matching: {e}", exc_info=True) return {"matched": False, "message": f"Error during matching: {str(e)}"}