mirror of
https://github.com/Dispatcharr/Dispatcharr.git
synced 2026-07-20 16:51:10 +00:00
Add sentence transformers to new matching function.
This commit is contained in:
parent
f6be6bc3a9
commit
d2085d57f8
1 changed files with 500 additions and 197 deletions
|
|
@ -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)}"}
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue