import sys import os import time import json import requests from huggingface_hub import HfApi, hf_hub_download from datetime import datetime REPO_ID = "ravi2814/grant-data-storage" KERNEL_SLUG = "grant-data-updater-worker" CODE_DIR = "kaggle_worker" os.environ["HF_HUB_DISABLE_PROGRESS_BARS"] = "1" def print_ui(msg): print(msg, flush=True) class RemoteController: def __init__(self): self.username = os.getenv("KAGGLE_USERNAME") self.hf_token = os.getenv("HF_TOKEN") self.hf_api = HfApi(token=self.hf_token) from kaggle.api.kaggle_api_extended import KaggleApi self.api = KaggleApi() self.api.authenticate() if not os.path.exists(CODE_DIR): os.makedirs(CODE_DIR) def reset_logs(self): try: self.hf_api.delete_file( path_in_repo="live_log.txt", repo_id=REPO_ID, repo_type="dataset", commit_message="Reset logs" ) except: pass def check_for_updates(self): print_ui("Checking for updates...") gov_date = None try: api_url = "https://open.canada.ca/data/api/3/action/package_show?id=432527ab-7aac-45b5-81d6-7597107a7013" resp = requests.get(api_url, timeout=10) result = resp.json()['result'] resource = next((r for r in result['resources'] if r['id'] == '1d15a62f-5656-49ad-8c88-f40ce689d831'), None) if resource: gov_str = resource.get('last_modified') or result.get('metadata_modified') gov_date = datetime.fromisoformat(gov_str.replace('Z', '+00:00')).replace(tzinfo=None) print_ui(f"Government Data Date: {gov_date}") else: return True, None, None try: hf_hub_download(repo_id=REPO_ID, filename="last_metadata.json", local_dir=".", token=self.hf_token, force_download=True, repo_type="dataset") with open("last_metadata.json", "r") as f: meta = json.load(f) my_str = meta.get('resource_modified') if not my_str: return True, gov_date, meta my_date = datetime.fromisoformat(my_str.replace('Z', '+00:00')).replace(tzinfo=None) print_ui(f"Your Database Date: {my_date}") if gov_date > my_date: print_ui("New data available.") return True, gov_date, meta else: return False, gov_date, meta except: return True, gov_date, None except Exception as e: print_ui(f"Check error: {e}. Forcing run.") return True, None, None def update_timestamp(self, meta): print_ui("Status: Up to date.") if meta: meta['checked_at'] = datetime.now().isoformat() with open("last_metadata.json", "w") as f: json.dump(meta, f) self.hf_api.upload_file( path_or_fileobj="last_metadata.json", path_in_repo="last_metadata.json", repo_id=REPO_ID, repo_type="dataset", commit_message="Timestamp update" ) print_ui("Timestamp updated.") def run_remote_worker(self): print_ui("Preparing worker...") with open("worker_script.py", "r") as f: script_content = f.read() if "HF_TOKEN_PLACEHOLDER" in script_content: script_content = script_content.replace("HF_TOKEN_PLACEHOLDER", self.hf_token) with open(os.path.join(CODE_DIR, "worker.py"), "w") as f: f.write(script_content) meta = { "id": f"{self.username}/{KERNEL_SLUG}", "title": "Grant Data Updater Worker", "code_file": "worker.py", "language": "python", "kernel_type": "script", "is_private": "true", "enable_gpu": "true", "enable_internet": "true", "dataset_sources": [], "kernel_sources": [], "competition_sources": [] } with open(os.path.join(CODE_DIR, "kernel-metadata.json"), "w") as f: json.dump(meta, f) print_ui("Deploying to Kaggle...") try: self.api.kernels_push(CODE_DIR) print_ui("Code pushed.") except Exception as e: if "already queued" in str(e).lower(): print_ui("Worker already queued.") else: raise e def monitor_progress(self): print_ui("Connecting to worker...") start_time = time.time() last_log_content = "" kaggle_confirmed = False time.sleep(5) # Give API a moment while True: try: stat = self.api.kernels_status(f"{self.username}/{KERNEL_SLUG}") status = stat['status'] if isinstance(stat, dict) else getattr(stat, 'status', 'unknown') except: status = "unknown" # --- KEY CHANGE: Trigger UI update immediately --- if not kaggle_confirmed and status in ['queued', 'running', 'starting']: kaggle_confirmed = True print_ui("STATUS: KAGGLE_STARTED") try: hf_hub_download(repo_id=REPO_ID, filename="live_log.txt", local_dir=".", token=self.hf_token, force_download=True, repo_type="dataset") if os.path.exists("live_log.txt"): with open("live_log.txt", "r") as f: content = f.read() if len(content) > len(last_log_content): new_text = content[len(last_log_content):] print(new_text, end='', flush=True) last_log_content = content except: pass if status == 'complete': print_ui("\nWorker finished successfully.") return True elif status == 'error': print_ui("\nWorker failed.") return False time.sleep(5) if time.time() - start_time > 3600: print_ui("Timeout: 60 minutes limit.") return False if __name__ == "__main__": try: if not os.getenv("HF_TOKEN") or not os.getenv("KAGGLE_USERNAME"): print_ui("Error: Missing credentials") sys.exit(1) controller = RemoteController() needs_update, gov_date, meta = controller.check_for_updates() if needs_update: controller.reset_logs() controller.run_remote_worker() success = controller.monitor_progress() if success: print_ui("Sync Complete") else: sys.exit(1) else: controller.update_timestamp(meta) except Exception as e: print_ui(f"Error: {e}") sys.exit(1)