File size: 2,549 Bytes
464c149
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
"""Background thread that periodically inserts new flow/metric rows so the
demo DB feels like it's reflecting a live network, without needing a real
ingestion pipeline."""

import random
import sqlite3
import threading
import time
from datetime import datetime
from pathlib import Path

DB_PATH = Path(__file__).parent / "network.db"
PROTOCOLS = ["TCP", "UDP", "ICMP"]


def _tick(conn):
    cur = conn.cursor()
    cur.execute("SELECT device_id FROM devices")
    device_ids = [r[0] for r in cur.fetchall()]
    if not device_ids:
        return

    now = datetime.utcnow().isoformat(timespec="seconds")
    did = random.choice(device_ids)

    cur.execute(
        "INSERT INTO traffic_flows (ts, src_ip, dst_ip, src_port, dst_port, protocol, "
        "bytes, packets, duration_ms, device_id, label) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
        (
            now,
            f"10.0.0.{random.randint(1, 254)}",
            f"192.168.{random.randint(0,255)}.{random.randint(1,254)}",
            random.randint(1024, 65535),
            random.choice([80, 443, 22, 53]),
            random.choice(PROTOCOLS),
            random.randint(200, 2_000_000),
            random.randint(1, 3000),
            random.randint(1, 5000),
            did,
            "benign",
        ),
    )

    cur.execute("SELECT interface_id, device_id FROM interfaces WHERE device_id = ?", (did,))
    ifaces = cur.fetchall()
    if ifaces:
        iface_id, iface_device_id = random.choice(ifaces)
        metric_name = random.choice(["latency_ms", "packet_loss_pct", "bandwidth_util_pct"])
        value = {
            "latency_ms": round(random.uniform(1, 180), 2),
            "packet_loss_pct": round(random.uniform(0, 6), 2),
            "bandwidth_util_pct": round(random.uniform(2, 98), 2),
        }[metric_name]
        cur.execute(
            "INSERT INTO metrics_timeseries (ts, device_id, interface_id, metric_name, value) "
            "VALUES (?, ?, ?, ?, ?)",
            (now, iface_device_id, iface_id, metric_name, value),
        )

    conn.commit()


def _loop(interval_seconds):
    conn = sqlite3.connect(DB_PATH, check_same_thread=False)
    while True:
        try:
            _tick(conn)
        except Exception as e:
            print(f"[simulator] tick failed: {e}")
        time.sleep(interval_seconds)


def start(interval_seconds: int = 5):
    """Starts the simulator in a daemon thread. Call once at app startup."""
    t = threading.Thread(target=_loop, args=(interval_seconds,), daemon=True)
    t.start()
    return t