Spaces:
Sleeping
Sleeping
| # main.py β FastAPI: Refactored with modular config, logger, scheduler | |
| # Endpoints: forecast, CRUD, scheduler, logs, settings | |
| import os | |
| import threading | |
| from datetime import datetime | |
| from typing import Optional | |
| from fastapi import FastAPI, HTTPException, Depends, BackgroundTasks, Query | |
| from fastapi.staticfiles import StaticFiles | |
| from fastapi.middleware.cors import CORSMiddleware | |
| from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials | |
| from pydantic import BaseModel | |
| from config import get_settings, runtime_config | |
| from logger import get_app_logger, memory_handler | |
| from scheduler import forecast_scheduler | |
| from database import ( | |
| get_all_komoditas, get_harga_harian, get_pasar_list, | |
| get_prediksi_data, save_model_ml, save_hasil_prediksi, | |
| save_ringkasan_prediksi, save_insight_prediksi, | |
| update_komoditas_volatilitas, get_engine, | |
| save_all_model_results, | |
| ) | |
| from ml_pipeline import run_pipeline | |
| from sqlalchemy import text | |
| settings = get_settings() | |
| logger = get_app_logger("main") | |
| app = FastAPI(title="SIKOMO ML API") | |
| app.add_middleware( | |
| CORSMiddleware, | |
| allow_origins=["*"], | |
| allow_methods=["*"], | |
| allow_headers=["*"], | |
| ) | |
| # ββ Auth ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| security = HTTPBearer() | |
| def verify_token(creds: HTTPAuthorizationCredentials = Depends(security)): | |
| if creds.credentials != settings.API_SECRET_KEY: | |
| raise HTTPException(status_code=401, detail="Unauthorized") | |
| return creds.credentials | |
| # ββ State tracking ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| running_tasks: dict = {} # komoditas_id -> True/False | |
| stop_event = threading.Event() # For aborting forecast | |
| def run_forecast_for_komoditas(komoditas_id: int, komoditas_nama: str): | |
| """Run forecast for a single komoditas. Tracks running state.""" | |
| if running_tasks.get(komoditas_id): | |
| logger.warning(f"β οΈ {komoditas_nama} sudah sedang berjalan, skip.") | |
| return | |
| running_tasks[komoditas_id] = True | |
| try: | |
| if stop_event.is_set(): | |
| logger.warning(f"β Forecast dihentikan sebelum {komoditas_nama}.") | |
| return | |
| df = get_harga_harian(komoditas_id) | |
| if df.empty or len(df) < 60: | |
| logger.warning(f"β οΈ {komoditas_nama}: data kurang ({len(df)} baris)") | |
| return | |
| price_cols = [c for c in df.columns if c != 'date'] | |
| target_col = price_cols[0] | |
| pasar_df = get_pasar_list(komoditas_id) | |
| pasar_id = int(pasar_df.iloc[0]['id']) if not pasar_df.empty else 1 | |
| def log_cb(msg): | |
| logger.info(f"[{komoditas_nama}] {msg}") | |
| result = run_pipeline( | |
| komoditas_id=komoditas_id, | |
| komoditas_nama=komoditas_nama, | |
| df=df, | |
| target_market=target_col, | |
| pasar_id=pasar_id, | |
| n_trials=runtime_config.optuna_trials, | |
| forecast_days=runtime_config.forecast_days, | |
| log_cb=log_cb, | |
| ) | |
| # Save results to DB | |
| meta = result['metadata'] | |
| labels = result['labels'] | |
| preds = result['predictions'] | |
| all_model_results = result.get('model_comparison', {}) | |
| # Save all model results (semua model, bukan hanya best) | |
| save_all_model_results(all_model_results, meta, komoditas_id, pasar_id) | |
| # Save best model specifically | |
| model_id = save_model_ml(meta) | |
| save_hasil_prediksi(preds, komoditas_id, pasar_id, model_id, meta) | |
| save_ringkasan_prediksi({ | |
| 'harga_min': labels['harga_min'], | |
| 'harga_max': labels['harga_max'], | |
| 'tren': labels['tren'], | |
| 'confidence_level': labels['confidence_level'], | |
| 'status_analisis': labels['status_analisis']['judul'], | |
| 'deskripsi_status': labels['status_analisis']['deskripsi'], | |
| 'tanggal_mulai': preds[0]['tanggal'], | |
| 'tanggal_akhir': preds[-1]['tanggal'], | |
| }, komoditas_id, model_id) | |
| save_insight_prediksi(labels['insights'], komoditas_id, model_id) | |
| cv = meta.get('data_quality', {}).get('cv', 0) | |
| update_komoditas_volatilitas(komoditas_id, 1 if cv >= 5 else 0, cv) | |
| logger.info(f"β {komoditas_nama} selesai! Best: {meta['nama_model']} MAPE={meta['mape']:.2f}%") | |
| except Exception as e: | |
| logger.error(f"β {komoditas_nama} gagal: {e}") | |
| finally: | |
| running_tasks[komoditas_id] = False | |
| def auto_forecast_all(): | |
| """Run forecast for ALL komoditas (scheduled or manual run-all).""" | |
| logger.info("π Auto forecast semua komoditas dimulai...") | |
| stop_event.clear() | |
| try: | |
| komoditas_list = get_all_komoditas() | |
| total = len(komoditas_list) | |
| for idx, (_, row) in enumerate(komoditas_list.iterrows()): | |
| if stop_event.is_set(): | |
| logger.warning("β Forecast dihentikan oleh user.") | |
| break | |
| logger.info(f"π¦ [{idx+1}/{total}] Memproses {row['nama']}...") | |
| run_forecast_for_komoditas(int(row['id']), row['nama']) | |
| logger.info("π Auto forecast selesai.") | |
| except Exception as e: | |
| logger.error(f"β Auto forecast error: {e}") | |
| # Register forecast job with scheduler | |
| forecast_scheduler.set_forecast_job(auto_forecast_all) | |
| # ββ Startup βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def startup_event(): | |
| logger.info("π Memulai SIKOMO Forecast API Server...") | |
| os.makedirs(runtime_config.model_dir, exist_ok=True) | |
| forecast_scheduler.start_all(default_schedules=runtime_config.default_schedules) | |
| logger.info("π API Server siap menerima permintaan.") | |
| async def shutdown_event(): | |
| forecast_scheduler.stop_all() | |
| logger.info("π API Server dimatikan.") | |
| # ββ Pydantic Models ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| class ForecastRequest(BaseModel): | |
| komoditas_id: int | |
| class InsightUpdate(BaseModel): | |
| konten: str | |
| tipe: str | |
| ikon: str | |
| urutan: int | |
| class RingkasanUpdate(BaseModel): | |
| status_analisis: Optional[str] = None | |
| deskripsi_status: Optional[str] = None | |
| tren: Optional[str] = None | |
| class ModelUpdate(BaseModel): | |
| deskripsi: Optional[str] = None | |
| catatan_validasi: Optional[str] = None | |
| class RuntimeConfigUpdate(BaseModel): | |
| key: str | |
| value: str | |
| class ScheduleAdd(BaseModel): | |
| cron_expression: str | |
| label: Optional[str] = "" | |
| # ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| # ENDPOINTS | |
| # ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def health(): | |
| return {"status": "ok", "time": datetime.now().isoformat()} | |
| # ββ Auth ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def login(body: dict): | |
| if body.get('password') != settings.DASHBOARD_PASSWORD: | |
| raise HTTPException(status_code=401, detail="Password salah") | |
| return {"token": settings.API_SECRET_KEY} | |
| # ββ Status & Logs (realtime) βββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def get_status(): | |
| """Overall system status including scheduler and running tasks.""" | |
| sched = forecast_scheduler.get_status() | |
| return { | |
| **sched, | |
| "running_tasks": {k: v for k, v in running_tasks.items() if v}, | |
| "any_forecast_running": any(running_tasks.values()), | |
| } | |
| def get_logs(limit: int = Query(200, le=500), level: Optional[str] = None): | |
| """Get structured logs with optional level filter.""" | |
| return memory_handler.get_logs(limit=limit, level=level) | |
| def clear_logs(): | |
| """Clear all in-memory logs.""" | |
| memory_handler.clear_logs() | |
| return {"success": True, "message": "Log berhasil dihapus."} | |
| # ββ Komoditas βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def list_komoditas(): | |
| df = get_all_komoditas() | |
| return df.to_dict('records') | |
| # ββ Prediksi Data βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def get_prediksi(komoditas_id: int): | |
| return get_prediksi_data(komoditas_id) | |
| # ββ Manual Forecast βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def run_forecast_manual(req: ForecastRequest, bg: BackgroundTasks): | |
| """Run forecast for a SINGLE komoditas (not all).""" | |
| if running_tasks.get(req.komoditas_id): | |
| raise HTTPException(409, f"Komoditas ID {req.komoditas_id} sedang diproses.") | |
| komoditas_list = get_all_komoditas() | |
| row = komoditas_list[komoditas_list['id'] == req.komoditas_id] | |
| if row.empty: | |
| raise HTTPException(404, "Komoditas tidak ditemukan") | |
| nama = row.iloc[0]['nama'] | |
| stop_event.clear() | |
| bg.add_task(run_forecast_for_komoditas, req.komoditas_id, nama) | |
| return {"message": f"Forecast untuk {nama} dimulai di background."} | |
| def run_forecast_all(bg: BackgroundTasks): | |
| """Run forecast for ALL komoditas.""" | |
| if any(running_tasks.values()): | |
| raise HTTPException(409, "Ada forecast yang sedang berjalan.") | |
| stop_event.clear() | |
| bg.add_task(auto_forecast_all) | |
| return {"message": "Forecast semua komoditas dimulai."} | |
| def stop_forecast(): | |
| """Abort running forecast.""" | |
| logger.warning("π Menerima permintaan penghentian forecast...") | |
| stop_event.set() | |
| return {"success": True, "message": "Permintaan penghentian dikirim."} | |
| # ββ Scheduler (Multiple Schedules) βββββββββββββββββββββββββββββββββββββββββββ | |
| def list_schedules(): | |
| return forecast_scheduler.list_schedules() | |
| def add_schedule(schedule: ScheduleAdd): | |
| try: | |
| job_id = forecast_scheduler.add_schedule( | |
| schedule.cron_expression, label=schedule.label | |
| ) | |
| return { | |
| "success": True, | |
| "job_id": job_id, | |
| "message": f"Jadwal berhasil ditambahkan: {schedule.cron_expression}", | |
| } | |
| except Exception as e: | |
| raise HTTPException(status_code=400, detail=str(e)) | |
| def remove_schedule(job_id: str): | |
| try: | |
| forecast_scheduler.remove_schedule(job_id) | |
| return {"success": True, "message": "Jadwal berhasil dihapus."} | |
| except KeyError: | |
| raise HTTPException(status_code=404, detail="Jadwal tidak ditemukan") | |
| except Exception as e: | |
| raise HTTPException(status_code=400, detail=str(e)) | |
| def start_scheduler(): | |
| forecast_scheduler.start_all(default_schedules=runtime_config.default_schedules) | |
| return {"success": True, "message": "Scheduler forecast diaktifkan."} | |
| def stop_scheduler(): | |
| forecast_scheduler.stop_all() | |
| return {"success": True, "message": "Scheduler forecast dimatikan."} | |
| # ββ Runtime Config ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def get_runtime_settings(): | |
| """Get runtime configuration (editable via dashboard).""" | |
| return runtime_config.to_dict() | |
| def update_runtime_settings(body: RuntimeConfigUpdate): | |
| """Update a single runtime config key.""" | |
| try: | |
| runtime_config.update(body.key, body.value) | |
| return {"message": f"{body.key} berhasil diperbarui ke '{body.value}'."} | |
| except ValueError as e: | |
| raise HTTPException(400, str(e)) | |
| # ββ CRUD Insight ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def get_insights(komoditas_id: int): | |
| engine = get_engine() | |
| with engine.connect() as conn: | |
| rows = conn.execute( | |
| text("SELECT * FROM insight_prediksi WHERE komoditas_id=:kid ORDER BY urutan"), | |
| {'kid': komoditas_id} | |
| ).fetchall() | |
| return [dict(r._mapping) for r in rows] | |
| def update_insight(insight_id: int, body: InsightUpdate): | |
| engine = get_engine() | |
| with engine.begin() as conn: | |
| conn.execute(text(""" | |
| UPDATE insight_prediksi | |
| SET konten=:konten, tipe=:tipe, ikon=:ikon, urutan=:urutan, updated_at=NOW() | |
| WHERE id=:id | |
| """), {**body.dict(), 'id': insight_id}) | |
| return {"message": "Insight diperbarui."} | |
| def add_insight(komoditas_id: int, body: InsightUpdate): | |
| engine = get_engine() | |
| with engine.begin() as conn: | |
| conn.execute(text(""" | |
| INSERT INTO insight_prediksi (komoditas_id, konten, tipe, ikon, urutan, is_active, created_at, updated_at) | |
| VALUES (:kid, :konten, :tipe, :ikon, :urutan, 1, NOW(), NOW()) | |
| """), {**body.dict(), 'kid': komoditas_id}) | |
| return {"message": "Insight ditambahkan."} | |
| def delete_insight(insight_id: int): | |
| engine = get_engine() | |
| with engine.begin() as conn: | |
| conn.execute(text("DELETE FROM insight_prediksi WHERE id=:id"), {'id': insight_id}) | |
| return {"message": "Insight dihapus."} | |
| # ββ CRUD Ringkasan ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def update_ringkasan(komoditas_id: int, body: RingkasanUpdate): | |
| engine = get_engine() | |
| updates = {k: v for k, v in body.dict().items() if v is not None} | |
| if not updates: | |
| raise HTTPException(400, "Tidak ada field yang diupdate.") | |
| set_clause = ', '.join([f"{k}=:{k}" for k in updates]) | |
| with engine.begin() as conn: | |
| conn.execute( | |
| text(f"UPDATE ringkasan_prediksi SET {set_clause} WHERE komoditas_id=:kid ORDER BY created_at DESC LIMIT 1"), | |
| {**updates, 'kid': komoditas_id} | |
| ) | |
| return {"message": "Ringkasan diperbarui."} | |
| # ββ CRUD Model ML βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| def update_model_ml(komoditas_id: int, body: ModelUpdate): | |
| engine = get_engine() | |
| updates = {k: v for k, v in body.dict().items() if v is not None} | |
| if not updates: | |
| raise HTTPException(400, "Tidak ada field yang diupdate.") | |
| set_clause = ', '.join([f"{k}=:{k}" for k in updates]) | |
| with engine.begin() as conn: | |
| conn.execute( | |
| text(f"UPDATE model_ml SET {set_clause}, updated_at=NOW() WHERE komoditas_id=:kid AND is_active=1"), | |
| {**updates, 'kid': komoditas_id} | |
| ) | |
| return {"message": "Model ML diperbarui."} | |
| # ββ Mount Frontend Static βββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| if os.path.isdir("static"): | |
| app.mount("/", StaticFiles(directory="static", html=True), name="frontend") |