""" normalize.py — BestOfSaoPaulo.blog Normalization Pipeline Reads wp_bosp_staging rows where normalized = 0, sends each to Gemini for normalization, translation, and neighborhood inference, then writes the result to wp_bosp_events and marks normalized = 1 in staging. One Gemini call per staging row. Temperature 0. Strict JSON output. """ import re import sys import json import time import logging import datetime import unicodedata from pathlib import Path import mysql.connector import google.generativeai as genai import config # ============================================================ # 1. LOGGING # ============================================================ def setup_logger(): log_dir = Path(config.LOGS_DIR) log_dir.mkdir(parents=True, exist_ok=True) timestamp = datetime.datetime.now().strftime('%Y-%m-%d_%H-%M-%S') log_file = log_dir / f'normalize_{timestamp}.log' logger = logging.getLogger('bosp_normalize') logger.setLevel(logging.DEBUG) fh = logging.FileHandler(log_file, encoding='utf-8') fh.setLevel(logging.DEBUG) fmt = logging.Formatter('[%(asctime)s UTC] %(levelname)s — %(message)s', datefmt='%Y-%m-%d %H:%M:%S') fh.setFormatter(fmt) logger.addHandler(fh) return logger log = setup_logger() # ============================================================ # 2. GEMINI SETUP # ============================================================ genai.configure(api_key=config.GEMINI_API_KEY) _gemini_model = genai.GenerativeModel( model_name=config.GEMINI_MODEL, generation_config=genai.types.GenerationConfig( temperature=config.GEMINI_TEMP, max_output_tokens=config.GEMINI_MAX_TOKENS, ), ) class GeminiQuotaExhausted(Exception): """Raised when Gemini daily quota is exhausted — abort the entire run.""" pass def call_gemini(prompt: str, retries: int = 3) -> str | None: """ Send a prompt to Gemini. Returns raw text response or None on failure. Retries up to `retries` times, honoring Gemini's retry-after delay. Raises GeminiQuotaExhausted if daily quota is gone. """ for attempt in range(retries): try: response = _gemini_model.generate_content(prompt) return response.text except Exception as e: err_str = str(e).lower() # Daily quota exhausted — no point retrying at all, abort the run if 'resource_exhausted' in err_str or ('quota' in err_str and '429' in err_str): log.critical('Gemini daily quota exhausted — aborting normalize run: %s', e) raise GeminiQuotaExhausted(str(e)) # Rate limit with retry-after — honor the actual wait time wait = 2 ** attempt if '429' in err_str or 'rate' in err_str: import re as _re m = _re.search(r'retry.?after[^\d]*(\d+)', err_str) if m: wait = int(m.group(1)) log.warning('Gemini rate limit (attempt %d/%d) — honoring retry-after %ds', attempt + 1, retries, wait) else: wait = 60 log.warning('Gemini rate limit (attempt %d/%d) — defaulting wait to %ds', attempt + 1, retries, wait) else: log.warning('Gemini error (attempt %d/%d): %s — retrying in %ds', attempt + 1, retries, e, wait) time.sleep(wait) log.error('Gemini failed after %d attempts', retries) return None # ============================================================ # 3. DATABASE # ============================================================ def get_db(): return mysql.connector.connect( host=config.DB_HOST, port=config.DB_PORT, database=config.DB_NAME, user=config.DB_USER, password=config.DB_PASSWORD, charset='utf8mb4', autocommit=False, ) def load_pending_rows(db) -> list: """Load all staging rows where normalized = 0.""" cur = db.cursor(dictionary=True) cur.execute( "SELECT * FROM `%s` WHERE normalized = 0 ORDER BY id ASC" % config.TABLE_STAGING ) rows = cur.fetchall() cur.close() return rows def mark_normalized(db, staging_id: int): """Set normalized = 1 for a staging row.""" cur = db.cursor() cur.execute( "UPDATE `%s` SET normalized = 1 WHERE id = %%s" % config.TABLE_STAGING, (staging_id,) ) db.commit() cur.close() def insert_event_row(db, row: dict) -> int | None: """ Insert a normalized event into wp_bosp_events. Returns the new row ID on success, None on failure. Skips insert if a row with the same venue + event_date already exists. """ # Dupe check: venue + date if row.get('venue') and row.get('event_date'): cur = db.cursor() cur.execute( "SELECT id FROM `%s` WHERE venue = %%s AND event_date = %%s LIMIT 1" % config.TABLE_EVENTS, (row['venue'], row['event_date']) ) existing = cur.fetchone() cur.close() if existing: log.info('DUPE SKIP (venue+date): %s / %s', row.get('venue'), row.get('event_date')) return None cols = [ 'staging_id', 'name_en', 'name_pt', 'name_es', 'slug_en', 'slug_pt', 'slug_es', 'venue', 'address', 'cep', 'neighborhood', 'event_date', 'event_time', 'price', 'price_color', 'category', 'subcategory', 'description_en', 'description_pt', 'description_es', 'website', 'instagram', 'facebook', 'whatsapp', 'other_link', 'source_url', 'is_admin_posted', 'identity_tags', 'accessibility_tags', 'watch_party', 'watch_party_team', 'sports_gender', ] placeholders = ', '.join(['%s'] * len(cols)) col_names = ', '.join(['`%s`' % c for c in cols]) sql = "INSERT INTO `%s` (%s) VALUES (%s)" % (config.TABLE_EVENTS, col_names, placeholders) values = [row.get(c) for c in cols] cur = db.cursor() try: cur.execute(sql, values) db.commit() new_id = cur.lastrowid return new_id except Exception as e: db.rollback() log.error('DB insert error: %s', e) return None finally: cur.close() # ============================================================ # 4. SLUG GENERATION # ============================================================ def make_slug(text: str) -> str: """ Generate a URL slug from any text. Handles accented characters, punctuation, and whitespace. """ if not text: return '' # Normalize unicode — decompose accented chars nfkd = unicodedata.normalize('NFKD', text) ascii_text = nfkd.encode('ascii', 'ignore').decode('ascii') # Lowercase slug = ascii_text.lower() # Replace non-alphanumeric with hyphens slug = re.sub(r'[^a-z0-9]+', '-', slug) # Strip leading/trailing hyphens slug = slug.strip('-') # Collapse multiple hyphens slug = re.sub(r'-{2,}', '-', slug) return slug def unique_slug(db, base_slug: str, lang: str) -> str: """ Ensure slug is unique in wp_bosp_events for the given language column. Appends -2, -3, etc. if needed. """ col = 'slug_%s' % lang candidate = base_slug suffix = 2 cur = db.cursor() while True: cur.execute( "SELECT id FROM `%s` WHERE `%s` = %%s LIMIT 1" % (config.TABLE_EVENTS, col), (candidate,) ) if not cur.fetchone(): break candidate = '%s-%d' % (base_slug, suffix) suffix += 1 cur.close() return candidate # ============================================================ # 5. PRICE HELPERS # ============================================================ def parse_price(raw_price) -> float | None: """ Parse a raw price value (string or number) to float. Returns 0.0 for free indicators, None if unparseable. """ if raw_price is None: return None text = str(raw_price).strip().lower() if not text: return None free_words = ['grátis', 'gratuito', 'free', 'entrada franca', 'entrada livre'] if any(w in text for w in free_words): return 0.0 # Strip currency symbols, normalize decimal comma cleaned = re.sub(r'[^\d,\.]', ' ', text) cleaned = cleaned.replace(',', '.') nums = re.findall(r'\d+(?:\.\d+)?', cleaned) if nums: val = float(nums[0]) return val return None # ============================================================ # 6. GEMINI NORMALIZATION PROMPT # ============================================================ _CATEGORIES_STR = '\n'.join('- ' + c for c in config.CANONICAL_CATEGORIES) _NEIGHBORHOODS_STR = '\n'.join('- ' + n for n in config.CANONICAL_NEIGHBORHOODS) def build_prompt(row: dict) -> str: """ Build the Gemini normalization prompt for a staging row. Returns a prompt string that instructs Gemini to return strict JSON. """ today = datetime.date.today().strftime('%Y-%m-%d') # Assemble all available raw data into a readable block for Gemini raw_block = '\n'.join([ 'Today\'s date: %s' % today, 'Raw Name: %s' % (row.get('raw_name') or ''), 'Raw Venue: %s' % (row.get('raw_venue') or ''), 'Raw Address: %s' % (row.get('raw_address') or ''), 'Raw CEP: %s' % (row.get('raw_cep') or ''), 'Raw Date: %s' % (row.get('raw_date') or ''), 'Raw Time: %s' % (row.get('raw_time') or ''), 'Raw Price: %s' % (row.get('raw_price') or ''), 'Raw Category: %s' % (row.get('raw_category') or ''), 'Raw Subcategory: %s' % (row.get('raw_subcategory') or ''), 'Raw Description: %s' % (row.get('raw_description') or ''), 'Raw Website: %s' % (row.get('raw_website') or ''), 'Raw Instagram: %s' % (row.get('raw_instagram') or ''), 'Raw Facebook: %s' % (row.get('raw_facebook') or ''), 'Source URL: %s' % (row.get('source_url') or ''), ]) prompt = """You are a data normalizer for BestOfSaoPaulo.blog, a São Paulo events website. You will receive raw scraped event data and must return a single JSON object with normalized fields. Respond with ONLY the JSON object. No explanation, no markdown, no code fences. Pure JSON only. RAW EVENT DATA: {raw_block} INSTRUCTIONS: 1. name_en: Write a clean English event name. If the original is in Portuguese, translate it naturally. 2. name_pt: The event name in Brazilian Portuguese. Keep it natural. 3. name_es: The event name in Spanish (Latin American). Keep it natural. 4. description_en: Write an engaging English description of this event. Use the raw description as source material. 2-4 sentences. Do not invent facts not present in the raw data. 5. description_pt: Translate description_en to Brazilian Portuguese. Literal and accurate. Do not add or remove content. 6. description_es: Translate description_en to Latin American Spanish. Literal and accurate. Do not add or remove content. 7. category: Map to EXACTLY ONE of these canonical categories (return the exact string): {categories} Choose the best match. If nothing fits, use "Other". 8. neighborhood: Infer the São Paulo neighborhood from the raw address and venue text. Map to EXACTLY ONE of these (return the exact string): {neighborhoods} If you cannot determine the neighborhood from the available text, return "Other". Do NOT guess. Only return a specific neighborhood if the address or venue text clearly indicates it. 9. venue: Clean up the raw venue name. Remove extra whitespace, fix obvious typos. Return the venue name only (not the full address). 10. address: Clean up the raw address. Return street + number only if clearly present. Otherwise return empty string "". 11. cep: Return the CEP (Brazilian postal code) if present in the raw data, in format XXXXX-XXX. Otherwise return "". 12. event_date: Return the event date in YYYY-MM-DD format. If Raw Date is provided, convert it. If Raw Date is empty, extract the date from Raw Description using today's date as reference for year inference. If no date can be found anywhere, return "". 13. event_time: Return the time in HH:MM format (24-hour). If not available, return "". 14. price: Return a numeric float representing the price in BRL. Return 0.0 if free. Return null if unknown. Return this exact JSON structure: {{ "name_en": "", "name_pt": "", "name_es": "", "description_en": "", "description_pt": "", "description_es": "", "category": "", "neighborhood": "", "venue": "", "address": "", "cep": "", "event_date": "", "event_time": "", "price": null }}""".format( raw_block=raw_block, categories=_CATEGORIES_STR, neighborhoods=_NEIGHBORHOODS_STR, ) return prompt # ============================================================ # 7. PARSE GEMINI RESPONSE # ============================================================ def parse_gemini_response(text: str) -> dict | None: """ Parse the JSON response from Gemini. Strips any accidental markdown fences before parsing. Returns dict or None on failure. """ if not text: return None # Strip markdown code fences if Gemini adds them despite instructions cleaned = text.strip() cleaned = re.sub(r'^```(?:json)?\s*', '', cleaned, flags=re.IGNORECASE) cleaned = re.sub(r'\s*```$', '', cleaned) cleaned = cleaned.strip() try: data = json.loads(cleaned) return data except json.JSONDecodeError as e: log.error('JSON parse error: %s | raw text: %.200s', e, text) return None def validate_gemini_data(data: dict) -> dict: """ Validate and sanitize the parsed Gemini response. Enforces canonical category. Clamps price. Returns cleaned dict. """ # Category must be one of the 11 canonicals category = str(data.get('category') or '').strip() if category not in config.CANONICAL_CATEGORIES: log.warning('Non-canonical category from Gemini: "%s" — defaulting to Other', category) category = 'Other' data['category'] = category # Neighborhood must be in our list neighborhood = str(data.get('neighborhood') or '').strip() if neighborhood not in config.CANONICAL_NEIGHBORHOODS: log.warning('Non-canonical neighborhood from Gemini: "%s" — defaulting to Other', neighborhood) neighborhood = 'Other' data['neighborhood'] = neighborhood # event_date must be YYYY-MM-DD or empty event_date = str(data.get('event_date') or '').strip() if event_date and not re.match(r'^\d{4}-\d{2}-\d{2}$', event_date): log.warning('Bad date format from Gemini: "%s" — clearing', event_date) event_date = '' data['event_date'] = event_date or None # event_time must be HH:MM or empty event_time = str(data.get('event_time') or '').strip() if event_time and not re.match(r'^\d{2}:\d{2}$', event_time): log.warning('Bad time format from Gemini: "%s" — clearing', event_time) event_time = '' data['event_time'] = event_time or None # price: must be float or None. Clamp to ceiling. price_raw = data.get('price') if price_raw is None: data['price'] = None else: try: price_float = float(price_raw) if price_float > config.MAX_PRICE_BRL: log.warning('Price %.2f exceeds ceiling — clamping to None', price_float) data['price'] = None else: data['price'] = price_float except (TypeError, ValueError): data['price'] = None # Ensure all text fields are strings, not None for field in ('name_en', 'name_pt', 'name_es', 'description_en', 'description_pt', 'description_es', 'venue', 'address', 'cep'): data[field] = str(data.get(field) or '').strip() return data # ============================================================ # 8. BUILD THE EVENTS ROW FROM STAGING + GEMINI DATA # ============================================================ def build_events_row(staging: dict, gemini: dict, db) -> dict: """ Merge staging row and validated Gemini data into a wp_bosp_events row dict. """ # Price: prefer Gemini's parsed price; fall back to parsing raw_price ourselves price = gemini.get('price') if price is None: price = parse_price(staging.get('raw_price')) # Final ceiling check if price is not None and float(price) > config.MAX_PRICE_BRL: price = None # Price color price_color = None if price is not None: price_color = config.get_price_color(float(price)) # Names: Gemini provides EN/PT/ES. Fallback to raw_name if Gemini returned empty. name_en = gemini.get('name_en') or str(staging.get('raw_name') or '').strip() name_pt = gemini.get('name_pt') or name_en name_es = gemini.get('name_es') or name_en # Slugs: generated from names, guaranteed unique per language slug_en_base = make_slug(name_en) slug_pt_base = make_slug(name_pt) slug_es_base = make_slug(name_es) slug_en = unique_slug(db, slug_en_base, 'en') if slug_en_base else '' slug_pt = unique_slug(db, slug_pt_base, 'pt') if slug_pt_base else '' slug_es = unique_slug(db, slug_es_base, 'es') if slug_es_base else '' # Venue: prefer Gemini cleaned; fallback to raw venue = gemini.get('venue') or str(staging.get('raw_venue') or '').strip() # Address: prefer Gemini; fallback to raw address = gemini.get('address') or str(staging.get('raw_address') or '').strip() # CEP: prefer Gemini; fallback to raw cep = gemini.get('cep') or str(staging.get('raw_cep') or '').strip() # Date/time: from Gemini only — already validated to YYYY-MM-DD / HH:MM. # Do NOT fall back to raw_date/raw_time — those may be non-ISO strings # that MySQL DATE/TIME columns will reject or silently mangle. event_date = gemini.get('event_date') or None event_time = gemini.get('event_time') or None # Descriptions description_en = gemini.get('description_en') or str(staging.get('raw_description') or '').strip() description_pt = gemini.get('description_pt') or description_en description_es = gemini.get('description_es') or description_en # Category (always from Gemini — it's been validated against canonical list) category = gemini.get('category') or 'Other' # Subcategory: pass through raw if present (Gemini doesn't assign subcategory) subcategory = str(staging.get('raw_subcategory') or '').strip() or None # Neighborhood: from Gemini neighborhood = gemini.get('neighborhood') or 'Other' # Social links: pass through from staging website = str(staging.get('raw_website') or '').strip() or None instagram = str(staging.get('raw_instagram') or '').strip() or None facebook = str(staging.get('raw_facebook') or '').strip() or None whatsapp = str(staging.get('raw_whatsapp') or '').strip() or None other_link = str(staging.get('raw_other_link') or '').strip() or None # Tags: pass through from staging as-is (JSON strings) identity_tags = staging.get('raw_identity_tags') or None accessibility_tags = staging.get('raw_accessibility_tags') or None # Watch party fields watch_party = int(staging.get('raw_watch_party') or 0) watch_party_team = str(staging.get('raw_watch_party_team') or '').strip() or None sports_gender = str(staging.get('raw_sports_gender') or '').strip() or None return { 'staging_id': staging['id'], 'name_en': name_en or None, 'name_pt': name_pt or None, 'name_es': name_es or None, 'slug_en': slug_en or None, 'slug_pt': slug_pt or None, 'slug_es': slug_es or None, 'venue': venue or None, 'address': address or None, 'cep': cep or None, 'neighborhood': neighborhood, 'event_date': event_date, 'event_time': event_time, 'price': price, 'price_color': price_color, 'category': category, 'subcategory': subcategory, 'description_en': description_en or None, 'description_pt': description_pt or None, 'description_es': description_es or None, 'website': website, 'instagram': instagram, 'facebook': facebook, 'whatsapp': whatsapp, 'other_link': other_link, 'source_url': str(staging.get('source_url') or '').strip() or None, 'is_admin_posted': 0, 'identity_tags': identity_tags, 'accessibility_tags': accessibility_tags, 'watch_party': watch_party, 'watch_party_team': watch_party_team, 'sports_gender': sports_gender, } # ============================================================ # 9. PROCESS ONE STAGING ROW # ============================================================ def process_row(staging: dict, db) -> bool: """ Normalize one staging row: call Gemini, parse response, insert into wp_bosp_events, mark staging normalized. Returns True on success (including dupe skip), False on hard failure. """ staging_id = staging['id'] raw_name = staging.get('raw_name') or '(no name)' log.info('Processing staging id=%d: %s', staging_id, raw_name) # Build and send Gemini prompt prompt = build_prompt(staging) raw_response = call_gemini(prompt) if raw_response is None: log.error('Gemini returned None for staging id=%d — skipping', staging_id) return False # Parse response gemini_data = parse_gemini_response(raw_response) if gemini_data is None: log.error('Could not parse Gemini JSON for staging id=%d — skipping', staging_id) return False # Validate and sanitize gemini_data = validate_gemini_data(gemini_data) # Require at minimum: name + (venue OR address) after normalization name_ok = bool(gemini_data.get('name_en') or staging.get('raw_name')) venue_ok = bool(gemini_data.get('venue') or staging.get('raw_venue')) address_ok = bool(gemini_data.get('address') or staging.get('raw_address')) if not name_ok or not (venue_ok or address_ok): log.warning( 'Staging id=%d failed viability after normalization ' '(name=%s venue/addr=%s) — marking normalized, skipping insert', staging_id, name_ok, venue_ok or address_ok ) # Mark normalized so we don't reprocess it endlessly mark_normalized(db, staging_id) return True # Build the events row events_row = build_events_row(staging, gemini_data, db) # Insert into wp_bosp_events new_id = insert_event_row(db, events_row) if new_id is None: # Either a dupe (logged as DUPE SKIP above) or a DB error (logged as DB insert error above). # Either way: mark normalized so we don't reprocess endlessly, but count as a non-insert. log.warning('staging id=%d: no event row inserted (dupe or error) — marking normalized', staging_id) mark_normalized(db, staging_id) return True # Mark staging row as normalized mark_normalized(db, staging_id) log.info( 'SUCCESS: staging id=%d → events id=%d | "%s" | %s | %s', staging_id, new_id, events_row.get('name_en', ''), events_row.get('event_date', ''), events_row.get('neighborhood', ''), ) return True # ============================================================ # 10. MAIN # ============================================================ def main(): log.info('============================================================') log.info('BOSP Normalize started') log.info('============================================================') # Bootstrap connection — load pending list only, then close immediately. # Each row gets its own fresh connection inside the loop (scraper.py pattern). try: bootstrap_db = get_db() except Exception as e: log.critical('DB connection failed: %s', e) sys.exit(1) try: pending = load_pending_rows(bootstrap_db) except Exception as e: log.critical('Failed to load pending rows: %s', e) bootstrap_db.close() sys.exit(1) bootstrap_db.close() if not pending: log.info('No rows to normalize. Exiting.') return log.info('Found %d row(s) to normalize', len(pending)) success_count = 0 fail_count = 0 for staging in pending: # Fresh DB connection per row — eliminates MySQL timeout risk on long runs try: db = get_db() except Exception as e: log.error('DB connection failed for staging id=%d: %s — skipping', staging.get('id'), e) fail_count += 1 time.sleep(0.5) continue try: ok = process_row(staging, db) if ok: success_count += 1 else: fail_count += 1 except GeminiQuotaExhausted: # Daily quota gone — no point processing remaining rows fail_count += 1 db.close() break except Exception as e: log.error('Unhandled error on staging id=%d: %s', staging.get('id'), e, exc_info=True) fail_count += 1 try: db.rollback() except Exception: pass finally: try: db.close() except Exception: pass # Brief pause between Gemini calls — avoid rate limiting time.sleep(0.5) log.info('============================================================') log.info('BOSP Normalize finished — success: %d | failed: %d', success_count, fail_count) log.info('============================================================') if __name__ == '__main__': main()