upwork-grant-project / fetch_api.py
ravi2814's picture
Rename fetch api.py to fetch_api.py
eddca54 verified
Raw
History Blame Contribute Delete
7.21 kB
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)