anoderb
feat: integrate MySQL user login, SQLite configs, AI Center, and dynamic ML hyperparameter tuning
d8fb724
Raw
History Blame Contribute Delete
27.5 kB
# main.py β€” FastAPI: Refactored with modular config, logger, scheduler
import os
import threading
import jwt
from datetime import datetime
from typing import Optional, List
from fastapi import FastAPI, HTTPException, Depends, BackgroundTasks, Query
from fastapi.staticfiles import StaticFiles
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
from pydantic import BaseModel
# Modular App Imports
from app.config import get_settings, runtime_config
from app.logger import get_app_logger, get_forecast_logger, scraper_memory_handler, forecast_memory_handler
from app.scheduler import scraper_scheduler, forecast_scheduler
from app.scraper import scraping_hari_ini, upload_to_mysql, stop_event as scraper_stop_event
from app.database import (
get_db_engine,
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,
save_all_model_results,
)
from app.ml_pipeline import run_pipeline
from sqlalchemy import text
settings = get_settings()
logger = get_app_logger("main")
forecast_logger = get_forecast_logger("main")
app = FastAPI(title="Sikomo Scraper & ML Portal API", version="2.0.0")
app.add_middleware(
CORSMiddleware,
allow_origins=["*"], # Allow all for HF Spaces compatibility
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
# ── Authentication Dependency ──────────────────────────────────────────────────
security = HTTPBearer()
def verify_token(creds: HTTPAuthorizationCredentials = Depends(security)):
token = creds.credentials
if token == settings.API_SECRET_KEY:
return {"sub": "system", "email": "system@sikomo.com", "name": "System API", "role": "super_admin"}
try:
payload = jwt.decode(token, settings.JWT_SECRET_KEY, algorithms=[settings.JWT_ALGORITHM])
return payload
except jwt.ExpiredSignatureError:
raise HTTPException(status_code=401, detail="Token has expired. Please login again.")
except jwt.InvalidTokenError:
raise HTTPException(status_code=401, detail="Invalid token. Authentication failed.")
# ── ML Forecasting State ──────────────────────────────────────────────────────
running_forecast_tasks: dict = {} # komoditas_id -> True/False
forecast_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_forecast_tasks.get(komoditas_id):
forecast_logger.warning(f"⚠️ {komoditas_nama} sudah sedang berjalan, skip.")
return
running_forecast_tasks[komoditas_id] = True
try:
if forecast_stop_event.is_set():
forecast_logger.warning(f"β›” Forecast dihentikan sebelum {komoditas_nama}.")
return
import app.settings_store as store
cfg = store.get_resolved_config(komoditas_id)
training_days = cfg.get("training_days", 730)
df = get_harga_harian(komoditas_id, days=training_days)
if df.empty or len(df) < 30: # Allow at least 30 rows of training data
forecast_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):
forecast_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,
config=cfg,
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
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)
forecast_logger.info(f"βœ… {komoditas_nama} selesai! Best: {meta['nama_model']} MAPE={meta['mape']:.2f}%")
except Exception as e:
forecast_logger.error(f"❌ {komoditas_nama} gagal: {e}")
finally:
running_forecast_tasks[komoditas_id] = False
def auto_forecast_all():
"""Run forecast for ALL komoditas (scheduled or manual run-all)."""
forecast_logger.info("πŸš€ Auto forecast semua komoditas dimulai...")
forecast_stop_event.clear()
try:
komoditas_list = get_all_komoditas()
total = len(komoditas_list)
for idx, (_, row) in enumerate(komoditas_list.iterrows()):
if forecast_stop_event.is_set():
forecast_logger.warning("β›” Forecast dihentikan oleh user.")
break
forecast_logger.info(f"πŸ“¦ [{idx+1}/{total}] Memproses {row['nama']}...")
run_forecast_for_komoditas(int(row['id']), row['nama'])
forecast_logger.info("🏁 Auto forecast selesai.")
except Exception as e:
forecast_logger.error(f"❌ Auto forecast error: {e}")
# Register forecast job with scheduler
forecast_scheduler.set_forecast_job(auto_forecast_all)
# ── API Startup & Shutdown ─────────────────────────────────────────────────────
@app.on_event("startup")
async def startup_event():
logger.info("πŸŽ‰ Memulai Sikomo Scraper & Forecast API Server...")
os.makedirs(runtime_config.model_dir, exist_ok=True)
scraper_scheduler.start_all()
forecast_scheduler.start_all(default_schedules=runtime_config.default_schedules)
logger.info("πŸš€ API Server siap menerima permintaan.")
@app.on_event("shutdown")
async def shutdown_event():
scraper_scheduler.stop_all()
forecast_scheduler.stop_all()
logger.info("πŸ›‘ API Server dimatikan.")
# ── Request Validation Models ──────────────────────────────────────────────────
class ScheduleAdd(BaseModel):
cron_expression: str
label: Optional[str] = ""
class CustomScrapeRequest(BaseModel):
date: str # Format YYYY-MM-DD
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 CommodityConfigInput(BaseModel):
use_template: int
training_days: int
test_split_ratio: float
forecast_days: int
enable_optuna_tuning: int
optuna_trials: int
enabled_models: List[str]
hyperparams: dict
class TemplateInput(BaseModel):
id: Optional[int] = None
name: str
is_default: int
training_days: int
test_split_ratio: float
forecast_days: int
enable_optuna_tuning: int
optuna_trials: int
enabled_models: List[str]
hyperparams: dict
# ══════════════════════════════════════════════════════════════════════════════
# ENDPOINTS
# ══════════════════════════════════════════════════════════════════════════════
# ── Auth ──────────────────────────────────────────────────────────────────────
@app.post("/auth/login")
def login(body: dict):
email = body.get('email')
password = body.get('password')
if not email or not password:
raise HTTPException(status_code=400, detail="Email dan password harus diisi")
try:
from app.database import verify_user_mysql
user_info = verify_user_mysql(email, password)
# Issue JWT Token
from datetime import datetime, timedelta
payload = {
"sub": str(user_info["id"]),
"email": user_info["email"],
"name": user_info["name"],
"role": user_info["role"],
"exp": datetime.utcnow() + timedelta(hours=24)
}
token = jwt.encode(payload, settings.JWT_SECRET_KEY, algorithm=settings.JWT_ALGORITHM)
return {
"token": token,
"user": {
"name": user_info["name"],
"email": user_info["email"],
"role": user_info["role"]
}
}
except Exception as e:
raise HTTPException(status_code=401, detail=str(e))
# ── Portal Scraper API ────────────────────────────────────────────────────────
@app.get("/scraper/status", dependencies=[Depends(verify_token)])
async def get_scraper_status():
return scraper_scheduler.get_status()
@app.get("/scraper/logs", dependencies=[Depends(verify_token)])
async def get_scraper_logs(limit: int = Query(100, le=500), level: Optional[str] = None):
return scraper_memory_handler.get_logs(limit=limit, level=level)
@app.get("/scraper/schedules", dependencies=[Depends(verify_token)])
async def get_scraper_schedules():
return scraper_scheduler.list_schedules()
@app.post("/scraper/schedules/add", dependencies=[Depends(verify_token)])
async def add_scraper_schedule(schedule: ScheduleAdd):
try:
job_id = scraper_scheduler.add_schedule(schedule.cron_expression, label=schedule.label)
return {
"success": True,
"job_id": job_id,
"message": f"Jadwal scraper berhasil ditambahkan: {schedule.cron_expression}"
}
except Exception as e:
raise HTTPException(status_code=400, detail=str(e))
@app.delete("/scraper/schedules/{job_id}/remove", dependencies=[Depends(verify_token)])
async def remove_scraper_schedule(job_id: str):
try:
scraper_scheduler.remove_schedule(job_id)
return {"success": True, "message": "Jadwal scraper 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))
@app.post("/scraper/run/manual", dependencies=[Depends(verify_token)])
def manual_scrape():
logger.info("⚑ Permintaan manual scraping dipicu via API...")
try:
df = scraping_hari_ini()
engine = get_db_engine()
stats = upload_to_mysql(df, engine)
return {"success": True, "stats": stats}
except Exception as e:
logger.error(f"❌ Manual scraping gagal: {e}")
raise HTTPException(status_code=500, detail=str(e))
@app.post("/scraper/run/custom", dependencies=[Depends(verify_token)])
def custom_scrape(req: CustomScrapeRequest):
logger.info(f"⚑ Permintaan scraping custom untuk tanggal {req.date} dipicu via API...")
try:
from datetime import datetime
target_date = datetime.strptime(req.date, "%Y-%m-%d")
df = scraping_hari_ini(target_date=target_date)
engine = get_db_engine()
stats = upload_to_mysql(df, engine)
return {"success": True, "stats": stats}
except ValueError:
raise HTTPException(status_code=400, detail="Format tanggal salah. Gunakan YYYY-MM-DD.")
except Exception as e:
logger.error(f"❌ Scraping custom gagal: {e}")
raise HTTPException(status_code=500, detail=str(e))
@app.post("/scraper/run/stop", dependencies=[Depends(verify_token)])
def stop_manual_scrape():
logger.warning("πŸ›‘ Menerima permintaan penghentian scraping manual...")
scraper_stop_event.set()
return {"success": True, "message": "Permintaan penghentian scraper dikirim"}
@app.post("/scraper/schedule/start", dependencies=[Depends(verify_token)])
async def start_scraper_scheduler():
scraper_scheduler.start_all()
return {"success": True, "message": "Scheduler scraper diaktifkan"}
@app.post("/scraper/schedule/stop", dependencies=[Depends(verify_token)])
async def stop_scraper_scheduler():
scraper_scheduler.stop_all()
return {"success": True, "message": "Scheduler scraper dimatikan"}
# ── Portal Forecast API ───────────────────────────────────────────────────────
@app.get("/forecast/status", dependencies=[Depends(verify_token)])
def get_forecast_status():
sched = forecast_scheduler.get_status()
return {
**sched,
"running_tasks": {k: v for k, v in running_forecast_tasks.items() if v},
"any_forecast_running": any(running_forecast_tasks.values()),
}
@app.get("/forecast/logs", dependencies=[Depends(verify_token)])
def get_forecast_logs(limit: int = Query(200, le=500), level: Optional[str] = None):
return forecast_memory_handler.get_logs(limit=limit, level=level)
@app.post("/forecast/clear-logs", dependencies=[Depends(verify_token)])
def clear_forecast_logs():
forecast_memory_handler.clear_logs()
return {"success": True, "message": "Log forecast berhasil dihapus."}
@app.get("/forecast/schedules", dependencies=[Depends(verify_token)])
def list_forecast_schedules():
return forecast_scheduler.list_schedules()
@app.post("/forecast/schedules/add", dependencies=[Depends(verify_token)])
def add_forecast_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 forecast berhasil ditambahkan: {schedule.cron_expression}",
}
except Exception as e:
raise HTTPException(status_code=400, detail=str(e))
@app.delete("/forecast/schedules/{job_id}/remove", dependencies=[Depends(verify_token)])
def remove_forecast_schedule(job_id: str):
try:
forecast_scheduler.remove_schedule(job_id)
return {"success": True, "message": "Jadwal forecast 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))
@app.post("/forecast/schedule/start", dependencies=[Depends(verify_token)])
def start_forecast_scheduler():
forecast_scheduler.start_all(default_schedules=runtime_config.default_schedules)
return {"success": True, "message": "Scheduler forecast diaktifkan."}
@app.post("/forecast/schedule/stop", dependencies=[Depends(verify_token)])
def stop_forecast_scheduler():
forecast_scheduler.stop_all()
return {"success": True, "message": "Scheduler forecast dimatikan."}
@app.get("/forecast/komoditas", dependencies=[Depends(verify_token)])
def list_forecast_komoditas():
df = get_all_komoditas()
return df.to_dict('records')
@app.get("/forecast/active-models-summary", dependencies=[Depends(verify_token)])
def get_active_models_summary_endpoint():
from app.database import get_active_models_summary
try:
return get_active_models_summary()
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.get("/forecast/prediksi/{komoditas_id}", dependencies=[Depends(verify_token)])
def get_forecast_prediksi(komoditas_id: int):
return get_prediksi_data(komoditas_id)
@app.post("/forecast/run", dependencies=[Depends(verify_token)])
def run_forecast_manual(req: ForecastRequest, bg: BackgroundTasks):
"""Run forecast for a SINGLE komoditas."""
if running_forecast_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']
forecast_stop_event.clear()
bg.add_task(run_forecast_for_komoditas, req.komoditas_id, nama)
return {"message": f"Forecast untuk {nama} dimulai di background."}
@app.post("/forecast/run-all", dependencies=[Depends(verify_token)])
def run_forecast_all(bg: BackgroundTasks):
"""Run forecast for ALL komoditas."""
if any(running_forecast_tasks.values()):
raise HTTPException(409, "Ada forecast yang sedang berjalan.")
forecast_stop_event.clear()
bg.add_task(auto_forecast_all)
return {"message": "Forecast semua komoditas dimulai."}
@app.post("/forecast/stop", dependencies=[Depends(verify_token)])
def stop_forecast():
"""Abort running forecast."""
forecast_logger.warning("πŸ›‘ Menerima permintaan penghentian forecast...")
forecast_stop_event.set()
return {"success": True, "message": "Permintaan penghentian dikirim."}
# ── SQLite Settings & Templates Endpoints ─────────────────────────────────────
@app.get("/forecast/configs/{komoditas_id}", dependencies=[Depends(verify_token)])
def get_commodity_config_endpoint(komoditas_id: int):
import app.settings_store as store
return store.get_commodity_config(komoditas_id)
@app.post("/forecast/configs/{komoditas_id}", dependencies=[Depends(verify_token)])
def save_commodity_config_endpoint(komoditas_id: int, body: CommodityConfigInput):
import app.settings_store as store
try:
store.save_commodity_config(
komoditas_id=komoditas_id,
use_template=body.use_template,
training_days=body.training_days,
test_split_ratio=body.test_split_ratio,
forecast_days=body.forecast_days,
enable_tuning=body.enable_optuna_tuning,
optuna_trials=body.optuna_trials,
enabled_models=body.enabled_models,
hyperparams=body.hyperparams
)
return {"success": True, "message": "Konfigurasi komoditas berhasil disimpan."}
except Exception as e:
raise HTTPException(status_code=400, detail=str(e))
@app.get("/forecast/templates", dependencies=[Depends(verify_token)])
def get_templates_endpoint():
import app.settings_store as store
return store.get_templates()
@app.post("/forecast/templates", dependencies=[Depends(verify_token)])
def save_template_endpoint(body: TemplateInput):
import app.settings_store as store
try:
if body.id:
store.update_template(
id_val=body.id,
name=body.name,
training_days=body.training_days,
test_split_ratio=body.test_split_ratio,
forecast_days=body.forecast_days,
enable_tuning=body.enable_optuna_tuning,
optuna_trials=body.optuna_trials,
enabled_models=body.enabled_models,
hyperparams=body.hyperparams,
is_default=body.is_default
)
return {"success": True, "message": "Template berhasil diperbarui.", "id": body.id}
else:
new_id = store.add_template(
name=body.name,
training_days=body.training_days,
test_split_ratio=body.test_split_ratio,
forecast_days=body.forecast_days,
enable_tuning=body.enable_optuna_tuning,
optuna_trials=body.optuna_trials,
enabled_models=body.enabled_models,
hyperparams=body.hyperparams,
is_default=body.is_default
)
return {"success": True, "message": "Template baru berhasil ditambahkan.", "id": new_id}
except Exception as e:
raise HTTPException(status_code=400, detail=str(e))
@app.post("/forecast/templates/{template_id}/set-default", dependencies=[Depends(verify_token)])
def set_default_template_endpoint(template_id: int):
import app.settings_store as store
try:
store.set_default_template(template_id)
return {"success": True, "message": "Template default berhasil diubah."}
except Exception as e:
raise HTTPException(status_code=400, detail=str(e))
@app.delete("/forecast/templates/{template_id}", dependencies=[Depends(verify_token)])
def delete_template_endpoint(template_id: int):
import app.settings_store as store
try:
store.delete_template(template_id)
return {"success": True, "message": "Template berhasil dihapus."}
except ValueError as ve:
raise HTTPException(status_code=400, detail=str(ve))
except Exception as e:
raise HTTPException(status_code=400, detail=str(e))
@app.get("/forecast/settings", dependencies=[Depends(verify_token)])
def get_runtime_settings():
import app.settings_store as store
default_tmpl = store.get_default_template()
return default_tmpl or {}
@app.post("/forecast/settings", dependencies=[Depends(verify_token)])
def update_runtime_settings(body: RuntimeConfigUpdate):
return {"message": "Endpoint ini tidak digunakan lagi. Konfigurasi sekarang menggunakan SQLite."}
# ── CRUD Forecast Insights, Ringkasan, Model ──────────────────────────────────
@app.get("/forecast/insight/{komoditas_id}", dependencies=[Depends(verify_token)])
def get_insights(komoditas_id: int):
engine = get_db_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]
@app.put("/forecast/insight/{insight_id}", dependencies=[Depends(verify_token)])
def update_insight(insight_id: int, body: InsightUpdate):
engine = get_db_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."}
@app.post("/forecast/insight/{komoditas_id}", dependencies=[Depends(verify_token)])
def add_insight(komoditas_id: int, body: InsightUpdate):
engine = get_db_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."}
@app.delete("/forecast/insight/{insight_id}", dependencies=[Depends(verify_token)])
def delete_insight(insight_id: int):
engine = get_db_engine()
with engine.begin() as conn:
conn.execute(text("DELETE FROM insight_prediksi WHERE id=:id"), {'id': insight_id})
return {"message": "Insight dihapus."}
@app.put("/forecast/ringkasan/{komoditas_id}", dependencies=[Depends(verify_token)])
def update_ringkasan(komoditas_id: int, body: RingkasanUpdate):
engine = get_db_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."}
@app.put("/forecast/model-ml/{komoditas_id}", dependencies=[Depends(verify_token)])
def update_model_ml(komoditas_id: int, body: ModelUpdate):
engine = get_db_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."}
# ── Debug Endpoints ───────────────────────────────────────────────────────────
@app.get("/debug/static", dependencies=[Depends(verify_token)])
def debug_static():
files = []
for root, dirs, filenames in os.walk("/app/static"):
for f in filenames:
files.append(os.path.relpath(os.path.join(root, f), "/app/static"))
return {"files": sorted(files)}
# ── SPA Serving fallback ──────────────────────────────────────────────────────
static_dirs = [
"/app/static",
os.path.join(os.path.dirname(os.path.dirname(__file__)), "static"),
os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(__file__))), "frontend", "out")
]
mounted = False
for sdir in static_dirs:
if os.path.exists(sdir):
app.mount("/", StaticFiles(directory=sdir, html=True), name="spa")
mounted = True
break
if not mounted:
@app.get("/")
async def root_fallback():
return {
"message": "Sikomo Scraper & Forecast API Server beroperasi normal.",
"note": "Frontend static files belum di-build/mount ke direktori static."
}