Spaces:
Sleeping
Sleeping
| 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) |