jarvis-cloud / backend /main.py.pre_fix_routes
Jarvis2345's picture
Squash history — remove all prior commits (secret hygiene, S4)
a31f556
Raw
History Blame
20.3 kB
import asyncio
import logging
import os
import sys
# --- load .env early: token_manager.py executes at import time, before main() ---
try:
from dotenv import load_dotenv as _load_dotenv
if getattr(sys, 'frozen', False):
_app_data = next(
(sys.argv[i + 1] for i, a in enumerate(sys.argv)
if a == '--app-data-dir' and i + 1 < len(sys.argv)),
None
)
if _app_data:
_load_dotenv(os.path.join(_app_data, '.env'), override=False)
else:
_load_dotenv()
else:
_load_dotenv()
del _load_dotenv
except ImportError:
pass
# -----------------------------------------------------------------------
# Insert root directory to sys.path to allow 'import modules.x' to work natively from Cloud
parent_dir = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
if parent_dir not in sys.path:
sys.path.insert(0, parent_dir)
os.environ["TYPEGUARD_DISABLE"] = "1"
# Torch load monkey patch moved to backend/voice/vad.py for lazy-loading
import inspect
_orig_getsource = inspect.getsource
def _safe_getsource(obj):
try:
return _orig_getsource(obj)
except OSError:
return ""
inspect.getsource = _safe_getsource
_orig_getsourcelines = inspect.getsourcelines
def _safe_getsourcelines(obj):
try:
return _orig_getsourcelines(obj)
except OSError:
return ([""], 0)
inspect.getsourcelines = _safe_getsourcelines
import uvicorn
# Ensure backend module is in sys.path when running from compiled PyInstaller executable
if getattr(sys, 'frozen', False):
sys.path.append(sys._MEIPASS)
from backend.ws.agent_ws import ws_manager
# --- STATIC IMPORTS FOR PYINSTALLER ---
# ----------------------------------------
from backend.db.migrations import run_all_migrations
from backend.routes.agent_routes import router as agent_router
from backend.routes.memory_routes import router as memory_router
from backend.routes.voice_routes import router as voice_router
from backend.routes.system_routes import router as system_router
from backend.routes.automation_routes import router as automation_router
from backend.routes.security_routes import router as security_router
from backend.routes.github_routes import router as github_router
from backend.routes.xr_routes import router as xr_router
from backend.routes.android_routes import router as android_router
from backend.routes.config_routes import router as config_router
from backend.routes.easter_egg_routes import router as easter_egg_router
from backend.routes.omega_routes import router as omega_router
from backend.routes.persona_routes import router as persona_router
from backend.routes.internet_routes import router as internet_router
from backend.routes.sentinel_routes import router as sentinel_router
from backend.routers.gaming_routes import router as gaming_router
from backend.services.voice_service import start_wake_word_loop
from backend.services.usb_monitor import start_usb_monitor
from backend.services.system_monitor import start_system_stats_loop
from backend.services.pc_mic_service import start_pc_mic_loop
from backend.voice.audio_ws import run_in_thread as start_audio_ws
APP_VERSION = "1.0.0"
from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Depends
from backend.security.auth import verify_token
try:
from slowapi import Limiter, _rate_limit_exceeded_handler
from slowapi.util import get_remote_address
from slowapi.errors import RateLimitExceeded
from slowapi.middleware import SlowAPIMiddleware
SLOWAPI_AVAILABLE = True
except ImportError:
logging.warning("slowapi package not found. Rate limiting is disabled.")
SLOWAPI_AVAILABLE = False
try:
import sentry_sdk
sentry_sdk.init(
dsn="https://c48f4234bb502ff74a2c5e518a3a65b2@o4511579717369856.ingest.us.sentry.io/4511579850407936",
traces_sample_rate=1.0,
send_default_pii=True,
)
except Exception:
pass
print(r"""
____. _____ ____________________.__ _________
| | / _ \______ \____ \ \ \ \ \ ___/
| |/ /_\ \| _// | \ \ \ \ \____ \
/\__| / | \ | \ | \ \ \_\ \ \ \
\________\____|__ /____|_ /_______/___/\______/____/
\/ \/
Fast Booting... Lazy Loading Heavy ML Models...
""")
app = FastAPI(title="JARVIS / FRIDAY OS Backend Sidecar")
if SLOWAPI_AVAILABLE:
limiter = Limiter(key_func=get_remote_address, default_limits=["100/minute"])
app.state.limiter = limiter
app.add_exception_handler(RateLimitExceeded, _rate_limit_exceeded_handler)
app.add_middleware(SlowAPIMiddleware)
app.include_router(agent_router, prefix="/agent", dependencies=[Depends(verify_token)])
app.include_router(memory_router, prefix="/memory", dependencies=[Depends(verify_token)])
app.include_router(voice_router, prefix="/voice", dependencies=[Depends(verify_token)])
app.include_router(system_router, prefix="/system", dependencies=[Depends(verify_token)])
app.include_router(automation_router, prefix="/automation", dependencies=[Depends(verify_token)])
app.include_router(security_router, prefix="/security", dependencies=[Depends(verify_token)])
app.include_router(github_router, prefix="/github", dependencies=[Depends(verify_token)])
app.include_router(xr_router, prefix="/xr", dependencies=[Depends(verify_token)])
app.include_router(android_router, prefix="/android", dependencies=[Depends(verify_token)])
app.include_router(config_router, prefix="/config", dependencies=[Depends(verify_token)])
app.include_router(omega_router, prefix="") # Prefix defined in router as /omega
app.include_router(easter_egg_router, prefix="/easter") # Exclude easter egg from auth
app.include_router(persona_router, prefix="/persona", dependencies=[Depends(verify_token)])
app.include_router(internet_router, prefix="/internet", dependencies=[Depends(verify_token)])
app.include_router(sentinel_router, prefix="/sentinel", dependencies=[Depends(verify_token)])
app.include_router(gaming_router, prefix="", dependencies=[Depends(verify_token)]) # gaming routes already have /gaming prefix
@app.get("/")
async def root_health_check():
return {"status": "JARVIS / FRIDAY Cloud OS Online", "message": "All systems nominal."}
@app.get("/wake")
@app.post("/wake")
async def wake_server():
"""Public endpoint — anyone can hit this from anywhere in the world to wake the cloud server.
HF Spaces sleep after inactivity; this wakes them up instantly."""
import time
return {
"status": "online",
"message": "JARVIS / FRIDAY Cloud Brain is awake and operational.",
"timestamp": time.time(),
"servers": {
"hf_space": "online",
"sqlite": "online",
"chroma": "online",
"mongodb_atlas": "online"
}
}
@app.get("/ping")
async def ping():
"""Ultra-lightweight public ping endpoint for keepalive from mobile app."""
return {"pong": True}
from fastapi import Request
from fastapi.responses import JSONResponse
from fastapi.middleware.cors import CORSMiddleware
from backend.dependencies.auth import verify_master_token
# Configure CORS for production (Tauri and Android network)
app.add_middleware(
CORSMiddleware,
allow_origins=["tauri://localhost", "http://localhost", "http://localhost:1420", "*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
@app.exception_handler(Exception)
async def global_exception_handler(request: Request, exc: Exception):
logging.error(f"Global Error on {request.url.path}: {exc}")
return JSONResponse(
status_code=500,
content={"error": "Internal Server Error", "detail": str(exc)},
)
@app.middleware("http")
async def auth_middleware(request: Request, call_next):
allowed_paths = ["/", "/health", "/wake", "/ping", "/android/pair", "/voice/upload_log"]
if request.url.path in allowed_paths or request.url.path.startswith("/ws") or request.url.path.endswith("/ws"):
return await call_next(request)
# Allow local connections without auth
if request.client.host in ["127.0.0.1", "localhost", "::1"]:
return await call_next(request)
token = request.headers.get("Authorization")
if not token or not verify_master_token(token):
return JSONResponse(status_code=401, content={"error": "Unauthorized Access"})
return await call_next(request)
IS_READY = False
audio_ws_server = None
BACKGROUND_TASKS = []
@app.on_event("startup")
async def startup_event():
global IS_READY, audio_ws_server
logging.info("Starting JARVIS / FRIDAY backend sidecar...")
# Start Audio WS in background thread (no-op stub in cloud mode)
audio_ws_server = start_audio_ws()
# Start automations scheduler
from backend.services.automation_service import init_automations
await init_automations()
# Track all asyncio background loops
# pc_mic_loop, wake_word_loop, usb, stats are local-only — skip in cloud (HF has no mic/usb)
cloud_mode = os.environ.get("CLOUD_ENV", "false").lower() == "true"
if not cloud_mode:
BACKGROUND_TASKS.append(asyncio.create_task(start_pc_mic_loop()))
BACKGROUND_TASKS.append(asyncio.create_task(start_wake_word_loop()))
BACKGROUND_TASKS.append(asyncio.create_task(start_usb_monitor()))
BACKGROUND_TASKS.append(asyncio.create_task(start_system_stats_loop()))
# --- JARVIS 10X Universal Gaming Coach ---
try:
from backend.gaming.coach_engine import coach_engine
from backend.gaming.overlay_renderer import run_overlay_in_background
run_overlay_in_background() # Spawns PyQt5 window in a background OS thread
BACKGROUND_TASKS.append(asyncio.create_task(coach_engine.start_loop()))
except Exception as e:
logging.warning(f"JARVIS 10X Coach not active (missing deps / Qt crash?): {e}")
else:
logging.info("[Startup] Cloud mode — skipping local hardware loops (mic, wake, usb, stats).")
# 1. Models Loader Placeholder (Wait for models)
# TTS model loading is offloaded to the Audio WS background thread to prevent blocking Uvicorn.
# 2. mDNS Registration
from backend.services.mdns_discovery import register_mdns_service
port = int(os.environ.get("JARVIS_PORT_BOUND", os.environ.get("JARVIS_PORT", "7474")))
register_mdns_service(port)
# 3. Auto-start PC Relay Client (Local Windows Only)
cloud_env = os.environ.get("CLOUD_ENV", "false").lower() == "true"
if not cloud_env:
try:
import subprocess
parent_dir = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
relay_script = os.path.join(parent_dir, "pc_relay_client.py")
if os.path.exists(relay_script):
si = subprocess.STARTUPINFO()
si.dwFlags |= subprocess.STARTF_USESHOWWINDOW
si.wShowWindow = 0
subprocess.Popen([sys.executable, relay_script], cwd=parent_dir, startupinfo=si, creationflags=0x08000000)
logging.info("Auto-started pc_relay_client.py natively from main server.")
except Exception as e:
logging.error(f"Could not auto-start pc_relay_client.py: {e}")
# 4. §2.3+2.4 — Start Continuous Research Mode (both environments)
from backend.omega.research_engine import start_continuous_research_mode
BACKGROUND_TASKS.append(asyncio.create_task(start_continuous_research_mode("jarvis")))
logging.info("[Startup] Continuous Research Mode activated. JARVIS will monitor and queue proposals.")
# 5. §Token Safety — Start Token Limit Checkpoint & Resume Monitor
from backend.services.token_manager import start_token_refresh_monitor
BACKGROUND_TASKS.append(asyncio.create_task(start_token_refresh_monitor()))
logging.info("[Startup] Token Refresh Monitor activated. All Gemini 3.5 Flash tasks will auto-resume on 429.")
# 6. Broadcast Ready (Health check passes)
IS_READY = True
@app.on_event("shutdown")
async def shutdown_event():
logging.info("Shutting down JARVIS backend sidecar...")
if audio_ws_server:
audio_ws_server.stop()
from backend.services.mdns_discovery import unregister_mdns_service
unregister_mdns_service()
from backend.services.automation_service import shutdown_automations
await shutdown_automations()
from backend.omega.research_engine import stop_continuous_research_mode
stop_continuous_research_mode()
from backend.services.token_manager import stop_token_refresh_monitor
stop_token_refresh_monitor()
# Gracefully cancel all background loops
for task in BACKGROUND_TASKS:
task.cancel()
await asyncio.gather(*BACKGROUND_TASKS, return_exceptions=True)
@app.get("/health")
async def health():
if not IS_READY:
from fastapi import HTTPException
raise HTTPException(status_code=503, detail="Backend starting up")
return {"status": "ok", "version": APP_VERSION}
@app.websocket("/ws")
async def websocket_hub(websocket: WebSocket):
await ws_manager.connect(websocket)
async def heartbeat():
while True:
await asyncio.sleep(30)
try:
await websocket.send_json({"event": "ping"})
except Exception:
logging.warning("Heartbeat failed, closing dead WS connection.")
try:
await websocket.close()
except Exception as e:
import logging; logging.getLogger(__name__).error(f"Swallowed exception: {e}")
break
heartbeat_task = asyncio.create_task(heartbeat())
try:
while True:
data = await websocket.receive_json()
await ws_manager.handle_client_event(websocket, data)
except WebSocketDisconnect:
pass
finally:
heartbeat_task.cancel()
ws_manager.disconnect(websocket)
@app.websocket("/agent/ws")
async def agent_websocket_alias(websocket: WebSocket):
await websocket_hub(websocket)
@app.websocket("/voice/ws")
async def voice_websocket_proxy(websocket: WebSocket):
await websocket.accept()
import websockets
try:
async with websockets.connect("ws://127.0.0.1:8767") as target_ws:
async def forward_to_target():
try:
while True:
msg = await websocket.receive()
if "bytes" in msg:
await target_ws.send(msg["bytes"])
elif "text" in msg:
await target_ws.send(msg["text"])
except Exception as e:
import logging; logging.getLogger(__name__).error(f"Swallowed exception: {e}")
async def forward_to_client():
try:
async for msg in target_ws:
if isinstance(msg, bytes):
await websocket.send_bytes(msg)
else:
await websocket.send_text(msg)
except Exception as e:
import logging; logging.getLogger(__name__).error(f"Swallowed exception: {e}")
await asyncio.gather(forward_to_target(), forward_to_client())
except Exception as e:
logging.error(f"Voice WS Proxy Error: {e}")
@app.websocket("/scene/ws")
async def scene_websocket_proxy(websocket: WebSocket):
await websocket.accept()
import websockets
try:
async with websockets.connect("ws://127.0.0.1:5050") as target_ws:
async def forward_to_target():
try:
while True:
msg = await websocket.receive()
if "bytes" in msg:
await target_ws.send(msg["bytes"])
elif "text" in msg:
await target_ws.send(msg["text"])
except Exception as e:
import logging; logging.getLogger(__name__).error(f"Swallowed exception: {e}")
async def forward_to_client():
try:
async for msg in target_ws:
if isinstance(msg, bytes):
await websocket.send_bytes(msg)
else:
await websocket.send_text(msg)
except Exception as e:
import logging; logging.getLogger(__name__).error(f"Swallowed exception: {e}")
await asyncio.gather(forward_to_target(), forward_to_client())
except Exception as e:
logging.error(f"Scene WS Proxy Error: {e}")
def main():
import argparse
parser = argparse.ArgumentParser()
parser.add_argument("--app-data-dir", default="")
parser.add_argument("--resource-dir", default="")
args, unknown = parser.parse_known_args()
if args.app_data_dir:
os.environ["JARVIS_APP_DATA_DIR"] = args.app_data_dir
if args.resource_dir:
os.environ["JARVIS_RESOURCE_DIR"] = args.resource_dir
import re
import json
class JSONMaskingFormatter(logging.Formatter):
key_patterns = [
re.compile(r'(sk-[a-zA-Z0-9]{20,})'),
re.compile(r'(AIza[0-9A-Za-z-_]{30,})'),
re.compile(r'\b([A-Za-z0-9-_]{32,})\b')
]
def format(self, record):
msg = record.getMessage()
for pattern in self.key_patterns:
def repl(m):
s = m.group(1)
if len(s) > 8:
return '*' * (len(s) - 4) + s[-4:]
return s
msg = pattern.sub(repl, msg)
log_record = {
"timestamp": self.formatTime(record, self.datefmt),
"level": record.levelname,
"name": record.name,
"message": msg,
}
if record.exc_info:
log_record["exception"] = self.formatException(record.exc_info)
return json.dumps(log_record)
root_logger = logging.getLogger()
root_logger.setLevel(logging.INFO)
handler = logging.StreamHandler()
handler.setFormatter(JSONMaskingFormatter())
root_logger.addHandler(handler)
# Run migrations synchronously before mounting Uvicorn
if getattr(sys, 'frozen', False):
db_path = os.path.join(os.environ.get('APPDATA', os.path.expanduser('~')), 'JARVIS_OS', 'memory.db')
else:
db_path = os.path.join(os.path.dirname(__file__), 'memory.db')
asyncio.run(run_all_migrations(db_path))
cloud_env = os.environ.get("CLOUD_ENV", "false").lower() == "true"
host = "0.0.0.0" if cloud_env else "127.0.0.1"
base_port = int(os.environ.get("PORT", "7860")) if cloud_env else int(os.environ.get("JARVIS_PORT", "7474"))
for port in range(base_port, base_port + 11):
try:
os.environ["JARVIS_PORT_BOUND"] = str(port)
print(f"JARVIS_PORT_BOUND={port}", flush=True)
uvicorn.run(app, host=host, port=port, log_config=None, ws_ping_interval=20, ws_ping_timeout=10)
break
except OSError as e:
if "WinError 10048" in str(e) or "address already in use" in str(e).lower():
logging.warning(f"Port {port} in use, trying next...")
continue
raise
if __name__ == "__main__":
main()