File size: 17,687 Bytes
bdef324
 
 
 
39018b4
 
bdef324
39018b4
 
 
 
 
bdef324
 
 
39018b4
 
 
 
bdef324
 
39018b4
 
 
 
bdef324
 
 
39018b4
 
 
 
 
 
 
 
 
 
 
 
 
bdef324
39018b4
 
 
bdef324
 
 
39018b4
bdef324
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
39018b4
 
bdef324
 
 
 
 
 
39018b4
bdef324
39018b4
bdef324
39018b4
 
bdef324
 
39018b4
bdef324
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
39018b4
 
 
 
bdef324
 
 
 
39018b4
 
bdef324
 
 
39018b4
 
bdef324
39018b4
 
bdef324
 
39018b4
 
bdef324
 
 
 
39018b4
bdef324
 
 
39018b4
 
 
 
 
bdef324
39018b4
 
bdef324
39018b4
bdef324
 
 
 
 
 
 
 
 
 
 
 
39018b4
bdef324
 
 
 
 
287a59e
 
 
 
 
 
bdef324
39018b4
 
 
 
 
bdef324
39018b4
 
 
 
bdef324
39018b4
 
bdef324
 
 
 
39018b4
 
 
 
 
bdef324
39018b4
 
 
 
 
bdef324
 
 
 
 
39018b4
 
bdef324
 
 
 
 
 
39018b4
bdef324
 
 
 
39018b4
bdef324
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
39018b4
 
 
bdef324
 
 
 
 
 
39018b4
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
bdef324
39018b4
 
 
 
 
 
 
 
 
 
 
 
 
 
bdef324
39018b4
 
 
 
 
 
 
 
 
 
 
 
 
 
bdef324
39018b4
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
# 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 ───────────────────────────────────────────────────────────────────
@app.on_event("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.")


@app.on_event("shutdown")
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
# ══════════════════════════════════════════════════════════════════════════════

@app.get("/health")
def health():
    return {"status": "ok", "time": datetime.now().isoformat()}

# ── Auth ──────────────────────────────────────────────────────────────────────
@app.post("/auth/login")
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) ─────────────────────────────────────────────────
@app.get("/status", dependencies=[Depends(verify_token)])
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()),
    }

@app.get("/logs", dependencies=[Depends(verify_token)])
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)

@app.post("/clear-logs", dependencies=[Depends(verify_token)])
def clear_logs():
    """Clear all in-memory logs."""
    memory_handler.clear_logs()
    return {"success": True, "message": "Log berhasil dihapus."}

# ── Komoditas ─────────────────────────────────────────────────────────────────
@app.get("/komoditas", dependencies=[Depends(verify_token)])
def list_komoditas():
    df = get_all_komoditas()
    return df.to_dict('records')

# ── Prediksi Data ─────────────────────────────────────────────────────────────
@app.get("/prediksi/{komoditas_id}", dependencies=[Depends(verify_token)])
def get_prediksi(komoditas_id: int):
    return get_prediksi_data(komoditas_id)

# ── Manual Forecast ───────────────────────────────────────────────────────────
@app.post("/forecast/run", dependencies=[Depends(verify_token)])
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."}

@app.post("/forecast/run-all", dependencies=[Depends(verify_token)])
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."}

@app.post("/forecast/stop", dependencies=[Depends(verify_token)])
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) ───────────────────────────────────────────
@app.get("/schedules", dependencies=[Depends(verify_token)])
def list_schedules():
    return forecast_scheduler.list_schedules()

@app.post("/schedules/add", dependencies=[Depends(verify_token)])
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))

@app.delete("/schedules/{job_id}/remove", dependencies=[Depends(verify_token)])
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))

@app.post("/schedule/start", dependencies=[Depends(verify_token)])
def start_scheduler():
    forecast_scheduler.start_all(default_schedules=runtime_config.default_schedules)
    return {"success": True, "message": "Scheduler forecast diaktifkan."}

@app.post("/schedule/stop", dependencies=[Depends(verify_token)])
def stop_scheduler():
    forecast_scheduler.stop_all()
    return {"success": True, "message": "Scheduler forecast dimatikan."}

# ── Runtime Config ────────────────────────────────────────────────────────────
@app.get("/settings", dependencies=[Depends(verify_token)])
def get_runtime_settings():
    """Get runtime configuration (editable via dashboard)."""
    return runtime_config.to_dict()

@app.post("/settings", dependencies=[Depends(verify_token)])
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 ──────────────────────────────────────────────────────────────
@app.get("/insight/{komoditas_id}", dependencies=[Depends(verify_token)])
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]

@app.put("/insight/{insight_id}", dependencies=[Depends(verify_token)])
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."}

@app.post("/insight/{komoditas_id}", dependencies=[Depends(verify_token)])
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."}

@app.delete("/insight/{insight_id}", dependencies=[Depends(verify_token)])
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 ────────────────────────────────────────────────────────────
@app.put("/ringkasan/{komoditas_id}", dependencies=[Depends(verify_token)])
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 ─────────────────────────────────────────────────────────────
@app.put("/model-ml/{komoditas_id}", dependencies=[Depends(verify_token)])
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")