skrapyx-api / firebase_manager.py
usmanovrustam's picture
Deploy Qwen 3.5 Intelligence Engine
c9e5599
Raw History Blame
15.2 kB
import os
import json
import gzip
from cryptography.fernet import Fernet
import firebase_admin
from firebase_admin import credentials, firestore, auth, storage
from google.cloud.firestore_v1.base_query import FieldFilter
from dotenv import load_dotenv
load_dotenv()
class FirebaseManager:
_instance = None
def __new__(cls):
if cls._instance is None:
cls._instance = super(FirebaseManager, cls).__new__(cls)
cls._instance.initialized = False
cls._instance.db = None
cls._instance.bucket = None
cls._instance._init_firebase()
cls._instance._init_vault()
return cls._instance
def _init_vault(self):
"""Initialize encryption engine"""
self.key = os.getenv("MASTER_ENCRYPTION_KEY")
if self.key:
try:
self.cipher = Fernet(self.key.encode())
print("[Vault] Neural Encryption Engine: Online")
except Exception as e:
print(f"[Vault] Encryption Init Failed: {e}")
self.cipher = None
else:
print("[Vault] WARNING: No Encryption Key found. Scans will be unencrypted.")
self.cipher = None
def _init_firebase(self):
# 1. Try to load from environment variable (Production Best Practice)
env_cred = os.getenv("FIREBASE_SERVICE_ACCOUNT")
if env_cred:
try:
cred_dict = json.loads(env_cred)
cred = credentials.Certificate(cred_dict)
firebase_admin.initialize_app(cred, {
'storageBucket': os.getenv("FIREBASE_STORAGE_BUCKET", f"{cred.project_id}.appspot.com")
})
self.db = firestore.client()
self.bucket = storage.bucket()
self.initialized = True
print("[Firebase] Successfully initialized from Environment Variable")
return
except Exception as e:
print(f"[Firebase] Env Init Failed: {e}")
# 2. Fallback to Dynamic Key Detection (Local Development)
search_paths = [
"config/store-screapper-firebase-adminsdk-fbsvc-5374dc9fc8.json",
"data/serviceAccountKey.json"
]
cred_path = None
for p in search_paths:
if os.path.exists(p):
cred_path = p
break
if not cred_path and os.path.exists("config"):
for f in os.listdir("config"):
if f.endswith(".json") and "firebase-adminsdk" in f:
cred_path = os.path.join("config", f)
break
if cred_path:
try:
cred = credentials.Certificate(cred_path)
# Dynamic project discovery
project_id = cred.project_id
# Try appspot.com first, fallback to firebasestorage.app
suffixes = [".appspot.com", ".firebasestorage.app"]
bucket_found = False
for suffix in suffixes:
bucket_name = f"{project_id}{suffix}"
try:
if firebase_admin._apps:
firebase_admin.delete_app(firebase_admin.get_app())
# Re-initialize with current candidate
firebase_admin.initialize_app(cred, {'storageBucket': bucket_name})
test_bucket = storage.bucket()
# CRITICAL: .exists() returns bool, must check it explicitly
if test_bucket.exists():
self.db = firestore.client()
self.bucket = test_bucket
self.initialized = True
bucket_found = True
print(f"[Firebase] Connected to Storage Bucket: {bucket_name}")
break
else:
print(f"[Firebase] Bucket {bucket_name} reported as non-existent.")
except Exception as e:
print(f"[Firebase] {bucket_name} initialization error: {e}")
if not bucket_found:
print("[Firebase] WARNING: No valid Storage Bucket found. Data will NOT be persisted to Cloud Storage.")
# Still initialize Firestore if possible
if not firebase_admin._apps:
firebase_admin.initialize_app(cred)
self.db = firestore.client()
self.initialized = True
self.bucket = None
except Exception as e:
print(f"[Firebase] Initialization fatal error: {e}. Local mode possible.")
else:
print("[Firebase] No service account key found in config/ or data/. Operating in Local Mode.")
def get_projects(self, user_id):
if not self.initialized: return {}
# Personalized sub-collection access
docs = self.db.collection("users").document(user_id).collection("projects").stream()
return {doc.id: doc.to_dict() for doc in docs}
def get_or_create_user(self, uid, email, name):
if not self.initialized: return {"user_id": uid, "email": email, "credits": 1}
doc_ref = self.db.collection('users').document(uid)
doc = doc_ref.get()
if doc.exists:
profile = doc.to_dict()
# Migration: Add username if missing
if "username" not in profile:
profile["username"] = email.split('@')[0] if email else "operator"
doc_ref.set(profile, merge=True)
return profile
else:
profile = {
"user_id": uid,
"email": email,
"name": name,
"username": email.split('@')[0] if email else "operator",
"credits": 1,
"is_pro": False,
"subscription_expiry": None,
"subscription_plan": None,
"subscription_id": None,
"created_at": firestore.SERVER_TIMESTAMP
}
doc_ref.set(profile)
return profile
def upload_scan(self, scan_id, data):
"""Uploads compressed & encrypted scan JSON to Firebase Storage"""
if not self.initialized or not self.bucket: return False
try:
# 1. Serialize
json_data = json.dumps(data).encode('utf-8')
# 2. Compress (Lossless Gzip)
compressed_data = gzip.compress(json_data, compresslevel=9)
# 3. Encrypt (Optional but recommended)
final_data = compressed_data
content_type = 'application/x-gzip'
if self.cipher:
final_data = self.cipher.encrypt(compressed_data)
content_type = 'application/octet-stream'
blob = self.bucket.blob(f"scans/{scan_id}.bin") # Change extension to .bin for vault files
blob.upload_from_string(
data=final_data,
content_type=content_type
)
return True
except Exception as e:
print(f"[Firebase] Vault Upload failed: {e}")
return False
def download_scan(self, scan_id):
"""Downloads, decrypts and decompresses scan from Firebase Storage"""
if not self.initialized or not self.bucket: return None
try:
# Try new .bin format first
blob = self.bucket.blob(f"scans/{scan_id}.bin")
if not blob.exists():
# Fallback to legacy .json format
blob = self.bucket.blob(f"scans/{scan_id}.json")
if not blob.exists(): return None
print(f"[Vault] Legacy Scan Detected: {scan_id}")
return json.loads(blob.download_as_text())
raw_bytes = blob.download_as_bytes()
# 1. Decrypt if possible
processed_data = raw_bytes
if self.cipher:
try:
processed_data = self.cipher.decrypt(raw_bytes)
except Exception:
# If decryption fails, maybe it's just compressed or legacy
pass
# 2. Decompress
try:
decompressed = gzip.decompress(processed_data)
return json.loads(decompressed.decode('utf-8'))
except Exception:
# Fallback for plain binary or failed decompression
return json.loads(processed_data.decode('utf-8'))
except Exception as e:
print(f"[Firebase] Vault Download failed: {e}")
return None
def delete_scan(self, scan_id):
"""Deletes scan JSON from Firebase Storage"""
if not self.initialized or not self.bucket: return False
try:
blob = self.bucket.blob(f"scans/{scan_id}.json")
if blob.exists(): blob.delete()
return True
except Exception as e:
print(f"[Firebase] Delete failed: {e}")
return False
def deduct_credit(self, user_id):
if not self.initialized: return True
user_ref = self.db.collection("users").document(user_id)
@firestore.transactional
def update_in_transaction(transaction, user_ref):
snapshot = user_ref.get(transaction=transaction)
if not snapshot.exists: return False
# Skip deduction only for active Pro users
if snapshot.get("is_pro") == True:
return True
credits = snapshot.get("credits") or 0
if credits > 0:
transaction.update(user_ref, {"credits": credits - 1})
return True
return False
transaction = self.db.transaction()
return update_in_transaction(transaction, user_ref)
def add_credits(self, user_id, amount):
"""Atomically increment user credits in Firestore"""
if not self.initialized: return False
try:
user_ref = self.db.collection("users").document(user_id)
@firestore.transactional
def update_in_transaction(transaction, user_ref):
snapshot = user_ref.get(transaction=transaction)
if not snapshot.exists: return False
current = snapshot.get("credits") or 0
transaction.update(user_ref, {"credits": current + amount})
return True
transaction = self.db.transaction()
return update_in_transaction(transaction, user_ref)
except Exception as e:
print(f"[Firebase] Credit addition failed: {e}")
return False
def upgrade_to_pro(self, user_id, status=True, expiry_date=None, plan_type=None, subscription_id=None):
"""Enable Pro status for a user with subscription metadata"""
if not self.initialized: return False
try:
user_ref = self.db.collection("users").document(user_id)
update_data = {
"is_pro": status,
"credits": 99 if status else 0,
"subscription_expiry": expiry_date,
"subscription_plan": plan_type,
"subscription_id": subscription_id
}
user_ref.update({k: v for k, v in update_data.items() if v is not None or not status})
return True
except Exception as e:
print(f"[Firebase] Pro upgrade failed: {e}")
return False
def enroll_in_metered(self, user_id, subscription_id=None):
"""Enroll user in the metered plan (pay-as-you-go) without granting bulk credits or pro status"""
if not self.initialized: return False
try:
user_ref = self.db.collection("users").document(user_id)
update_data = {
"is_pro": False,
"subscription_plan": "metered",
"subscription_id": subscription_id
}
# We explicitly do NOT touch "credits" so they keep their 0 balance (triggering the $0.50 UI prompt)
user_ref.update(update_data)
return True
except Exception as e:
print(f"[Firebase] Metered enrollment failed: {e}")
return False
def save_project(self, user_id, project_id, data):
if not self.initialized: return
self.db.collection("users").document(user_id).collection("projects").document(project_id).set(data)
def delete_project(self, user_id, project_id):
if not self.initialized: return
self.db.collection("users").document(user_id).collection("projects").document(project_id).delete()
# Job Engine Integration (Personalized)
def get_jobs(self, user_id):
if not self.initialized: return {}
docs = self.db.collection("users").document(user_id).collection("jobs").stream()
return {doc.id: doc.to_dict() for doc in docs}
def update_job(self, user_id, job_id, data):
if not self.initialized: return
self.db.collection("users").document(user_id).collection("jobs").document(job_id).set(data, merge=True)
# Transaction Ledger
def record_transaction(self, user_id, session_id, data):
"""Store transaction record in user sub-collection"""
if not self.initialized: return False
try:
tx_ref = self.db.collection("users").document(user_id).collection("transactions").document(session_id)
tx_ref.set(data)
return True
except Exception as e:
print(f"[Firebase] Transaction record failed: {e}")
return False
def save_purchase(self, user_id, session_id, data):
"""Store minimal purchase marker for idempotency in user sub-collection"""
if not self.initialized: return False
try:
p_ref = self.db.collection("users").document(user_id).collection("purchases").document(session_id)
p_ref.set(data)
return True
except Exception as e:
print(f"[Firebase] Purchase save failed: {e}")
return False
def get_purchase(self, user_id, session_id):
"""Check for existing purchase marker"""
if not self.initialized: return None
doc = self.db.collection("users").document(user_id).collection("purchases").document(session_id).get()
return doc.to_dict() if doc.exists else None
def get_transactions(self, user_id):
"""Fetch transaction history for user"""
if not self.initialized: return []
docs = self.db.collection("users").document(user_id).collection("transactions").order_by("created_at", direction=firestore.Query.DESCENDING).stream()
return [doc.to_dict() for doc in docs]
firebase_manager = FirebaseManager()