CI/CD Pipeline - Northern Thailand Ping River Monitor / Test Suite (3.11) (push) Failing after 23s
CI/CD Pipeline - Northern Thailand Ping River Monitor / Build Docker Image (push) Skipped
CI/CD Pipeline - Northern Thailand Ping River Monitor / Integration Test with Services (push) Skipped
CI/CD Pipeline - Northern Thailand Ping River Monitor / Deploy to Staging (push) Skipped
CI/CD Pipeline - Northern Thailand Ping River Monitor / Deploy to Production (push) Skipped
CI/CD Pipeline - Northern Thailand Ping River Monitor / Performance Test (push) Skipped
CI/CD Pipeline - Northern Thailand Ping River Monitor / Code Quality (push) Successful in 14s
Documentation / Validate Documentation (push) Failing after 8s
Documentation / Generate API Documentation (push) Successful in 9s
Documentation / Build Sphinx Documentation (push) Successful in 17s
CI/CD Pipeline - Northern Thailand Ping River Monitor / Cleanup (push) Successful in 1s
Documentation / Documentation Summary (push) Successful in 3s
New SQLAdapter.get_database_stats() aggregates totals, station count, date range, and hourly-slot coverage in one query per dialect. The endpoint follows the existing 503-guard/to_thread/TTL-cache pattern; the dashboard gains a five-tile stats strip on the existing refresh cadence. Coverage denominator is hour-truncated so off-hour endpoints cannot push it past 100%; MySQL slot expression avoids % characters that would break under pyformat bind interpolation.
1027 lines
39 KiB
Python
1027 lines
39 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Database adapters for different storage backends
|
|
"""
|
|
|
|
import datetime
|
|
import logging
|
|
from abc import ABC, abstractmethod
|
|
from typing import Dict, List, Optional
|
|
|
|
|
|
# Base adapter interface
|
|
class DatabaseAdapter(ABC):
|
|
@abstractmethod
|
|
def connect(self):
|
|
pass
|
|
|
|
@abstractmethod
|
|
def save_measurements(self, measurements: List[Dict]) -> bool:
|
|
pass
|
|
|
|
@abstractmethod
|
|
def get_latest_measurements(self, limit: int = 100) -> List[Dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
def get_measurements_by_timerange(
|
|
self,
|
|
start_time: datetime.datetime,
|
|
end_time: datetime.datetime,
|
|
station_codes: Optional[List[str]] = None,
|
|
) -> List[Dict]:
|
|
pass
|
|
|
|
@abstractmethod
|
|
def get_measurements_for_date(self, target_date: datetime.datetime) -> List[Dict]:
|
|
pass
|
|
|
|
def get_measurement_date_range(
|
|
self,
|
|
) -> Optional[tuple]:
|
|
"""Return (min_timestamp, max_timestamp) of stored measurements.
|
|
|
|
Returns None when the backend has no data or does not support the query.
|
|
"""
|
|
return None
|
|
|
|
def get_recorded_hours_by_day(
|
|
self, start_date: datetime.date, end_date: datetime.date
|
|
) -> Optional[Dict[datetime.date, set]]:
|
|
"""Map each day in [start_date, end_date] to the set of hours (0-23)
|
|
that have at least one measurement.
|
|
|
|
Returns None when the backend does not support hour-granular gap
|
|
detection (callers should fall back to day-granular checks).
|
|
"""
|
|
return None
|
|
|
|
def get_database_stats(self) -> Optional[Dict]:
|
|
"""Summary statistics over stored measurements: total count, distinct
|
|
stations, first/last timestamp, and hourly-slot coverage.
|
|
|
|
Returns None when the backend has no data or does not support the query.
|
|
"""
|
|
return None
|
|
|
|
|
|
# InfluxDB Adapter
|
|
class InfluxDBAdapter(DatabaseAdapter):
|
|
def __init__(
|
|
self,
|
|
host: str = "localhost",
|
|
port: int = 8086,
|
|
database: str = "water_monitoring",
|
|
username: str = None,
|
|
password: str = None,
|
|
):
|
|
self.host = host
|
|
self.port = port
|
|
self.database = database
|
|
self.username = username
|
|
self.password = password
|
|
self.client = None
|
|
|
|
def connect(self):
|
|
try:
|
|
from influxdb import InfluxDBClient
|
|
|
|
self.client = InfluxDBClient(
|
|
host=self.host,
|
|
port=self.port,
|
|
username=self.username,
|
|
password=self.password,
|
|
database=self.database,
|
|
)
|
|
|
|
# Create database if it doesn't exist
|
|
databases = self.client.get_list_database()
|
|
if not any(db["name"] == self.database for db in databases):
|
|
self.client.create_database(self.database)
|
|
logging.info(f"Created InfluxDB database: {self.database}")
|
|
|
|
# Create retention policy (keep data for 2 years, downsample after 30 days)
|
|
retention_policies = self.client.get_list_retention_policies(self.database)
|
|
if not any(rp["name"] == "water_data_policy" for rp in retention_policies):
|
|
self.client.create_retention_policy(
|
|
"water_data_policy",
|
|
"730d", # 2 years
|
|
"1", # replication factor
|
|
database=self.database,
|
|
default=True,
|
|
)
|
|
|
|
logging.info("Connected to InfluxDB successfully")
|
|
return True
|
|
|
|
except ImportError:
|
|
logging.error("InfluxDB client not installed. Run: pip install influxdb")
|
|
return False
|
|
except Exception as e:
|
|
logging.error(f"Failed to connect to InfluxDB: {e}")
|
|
return False
|
|
|
|
def save_measurements(self, measurements: List[Dict]) -> bool:
|
|
if not self.client:
|
|
logging.error("InfluxDB client not connected")
|
|
return False
|
|
|
|
try:
|
|
points = []
|
|
for measurement in measurements:
|
|
point = {
|
|
"measurement": "water_data",
|
|
"tags": {
|
|
"station_code": measurement["station_code"],
|
|
"station_name_en": measurement["station_name_en"],
|
|
"station_name_th": measurement["station_name_th"],
|
|
},
|
|
"time": measurement["timestamp"].isoformat(),
|
|
"fields": {
|
|
"water_level": float(measurement["water_level"]),
|
|
"discharge": float(measurement["discharge"])
|
|
if measurement.get("discharge") is not None
|
|
else None,
|
|
"discharge_percent": float(measurement["discharge_percent"])
|
|
if measurement.get("discharge_percent")
|
|
else None,
|
|
},
|
|
}
|
|
points.append(point)
|
|
|
|
success = self.client.write_points(points)
|
|
if success:
|
|
logging.info(f"Successfully wrote {len(points)} points to InfluxDB")
|
|
return success
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error writing to InfluxDB: {e}")
|
|
return False
|
|
|
|
def get_latest_measurements(self, limit: int = 100) -> List[Dict]:
|
|
if not self.client:
|
|
return []
|
|
|
|
try:
|
|
# Cast limit to int so it can never carry an injection payload.
|
|
limit = int(limit)
|
|
query = f"""
|
|
SELECT last("water_level") as water_level,
|
|
last("discharge") as discharge,
|
|
last("discharge_percent") as discharge_percent
|
|
FROM "water_data"
|
|
GROUP BY "station_code", "station_name_en", "station_name_th"
|
|
LIMIT {limit}
|
|
"""
|
|
|
|
result = self.client.query(query)
|
|
measurements = []
|
|
|
|
for point in result.get_points():
|
|
measurements.append(
|
|
{
|
|
"timestamp": point["time"],
|
|
"station_code": point.get("station_code"),
|
|
"station_name_en": point.get("station_name_en"),
|
|
"station_name_th": point.get("station_name_th"),
|
|
"water_level": point.get("water_level"),
|
|
"discharge": point.get("discharge"),
|
|
"discharge_percent": point.get("discharge_percent"),
|
|
}
|
|
)
|
|
|
|
return measurements
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error querying InfluxDB: {e}")
|
|
return []
|
|
|
|
def get_measurements_by_timerange(
|
|
self,
|
|
start_time: datetime.datetime,
|
|
end_time: datetime.datetime,
|
|
station_codes: Optional[List[str]] = None,
|
|
) -> List[Dict]:
|
|
if not self.client:
|
|
return []
|
|
|
|
try:
|
|
# start_time/end_time are datetime objects (fixed isoformat, injection-safe).
|
|
# station_codes are untrusted strings -> bind them as parameters.
|
|
bind_params = {}
|
|
where_clause = f"time >= '{start_time.isoformat()}' AND time <= '{end_time.isoformat()}'"
|
|
if station_codes:
|
|
placeholders = []
|
|
for i, code in enumerate(station_codes):
|
|
key = f"sc{i}"
|
|
bind_params[key] = code
|
|
placeholders.append(f"station_code = ${key}")
|
|
where_clause += " AND (" + " OR ".join(placeholders) + ")"
|
|
|
|
query = f"""
|
|
SELECT "water_level", "discharge", "discharge_percent", "station_code", "station_name_en", "station_name_th"
|
|
FROM "water_data"
|
|
WHERE {where_clause}
|
|
ORDER BY time DESC
|
|
"""
|
|
|
|
result = self.client.query(query, bind_params=bind_params)
|
|
measurements = []
|
|
|
|
for point in result.get_points():
|
|
measurements.append(
|
|
{
|
|
"timestamp": point["time"],
|
|
"station_code": point.get("station_code"),
|
|
"station_name_en": point.get("station_name_en"),
|
|
"station_name_th": point.get("station_name_th"),
|
|
"water_level": point.get("water_level"),
|
|
"discharge": point.get("discharge"),
|
|
"discharge_percent": point.get("discharge_percent"),
|
|
}
|
|
)
|
|
|
|
return measurements
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error querying InfluxDB: {e}")
|
|
return []
|
|
|
|
|
|
# MySQL/PostgreSQL Adapter
|
|
class SQLAdapter(DatabaseAdapter):
|
|
def __init__(self, connection_string: str, db_type: str = "mysql"):
|
|
self.connection_string = connection_string
|
|
self.db_type = db_type.lower()
|
|
self.engine = None
|
|
|
|
# Add SQLite-specific connection parameters for better concurrency
|
|
if self.db_type == "sqlite":
|
|
if "?" not in connection_string:
|
|
self.connection_string += "?timeout=30&check_same_thread=False"
|
|
else:
|
|
self.connection_string += "&timeout=30&check_same_thread=False"
|
|
|
|
def connect(self):
|
|
try:
|
|
from sqlalchemy import create_engine
|
|
|
|
self.engine = create_engine(self.connection_string, pool_pre_ping=True)
|
|
|
|
# Create tables
|
|
self._create_tables()
|
|
|
|
logging.info(f"Connected to {self.db_type.upper()} successfully")
|
|
return True
|
|
|
|
except ImportError:
|
|
logging.error(
|
|
"SQLAlchemy not installed. Run: pip install sqlalchemy pymysql"
|
|
)
|
|
return False
|
|
except Exception as e:
|
|
logging.error(f"Failed to connect to {self.db_type.upper()}: {e}")
|
|
return False
|
|
|
|
def _create_tables(self):
|
|
from sqlalchemy import text
|
|
|
|
# Stations table - adjust for different databases
|
|
if self.db_type == "sqlite":
|
|
stations_sql = """
|
|
CREATE TABLE IF NOT EXISTS stations (
|
|
id INTEGER PRIMARY KEY,
|
|
station_code TEXT UNIQUE NOT NULL,
|
|
thai_name TEXT NOT NULL,
|
|
english_name TEXT NOT NULL,
|
|
latitude REAL,
|
|
longitude REAL,
|
|
geohash TEXT,
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
|
|
)
|
|
"""
|
|
|
|
measurements_sql = """
|
|
CREATE TABLE IF NOT EXISTS water_measurements (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
timestamp DATETIME NOT NULL,
|
|
station_id INTEGER NOT NULL,
|
|
water_level REAL,
|
|
discharge REAL,
|
|
discharge_percent REAL,
|
|
status TEXT DEFAULT 'active',
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
FOREIGN KEY (station_id) REFERENCES stations(id),
|
|
UNIQUE(timestamp, station_id)
|
|
)
|
|
"""
|
|
|
|
# Create indexes separately for SQLite
|
|
index_sql = [
|
|
"CREATE INDEX IF NOT EXISTS idx_timestamp ON water_measurements(timestamp)",
|
|
"CREATE INDEX IF NOT EXISTS idx_station_timestamp ON water_measurements(station_id, timestamp)",
|
|
]
|
|
|
|
elif self.db_type == "postgresql":
|
|
stations_sql = """
|
|
CREATE TABLE IF NOT EXISTS stations (
|
|
id SERIAL PRIMARY KEY,
|
|
station_code VARCHAR(10) UNIQUE NOT NULL,
|
|
thai_name VARCHAR(255) NOT NULL,
|
|
english_name VARCHAR(255) NOT NULL,
|
|
latitude DECIMAL(10,8),
|
|
longitude DECIMAL(11,8),
|
|
geohash VARCHAR(20),
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
|
|
)
|
|
"""
|
|
|
|
measurements_sql = """
|
|
CREATE TABLE IF NOT EXISTS water_measurements (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
timestamp TIMESTAMP NOT NULL,
|
|
station_id INTEGER NOT NULL,
|
|
water_level NUMERIC(10,3),
|
|
discharge NUMERIC(10,2),
|
|
discharge_percent NUMERIC(5,2),
|
|
status VARCHAR(20) DEFAULT 'active',
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
FOREIGN KEY (station_id) REFERENCES stations(id),
|
|
UNIQUE(timestamp, station_id)
|
|
)
|
|
"""
|
|
|
|
index_sql = [
|
|
"CREATE INDEX IF NOT EXISTS idx_timestamp ON water_measurements(timestamp)",
|
|
"CREATE INDEX IF NOT EXISTS idx_station_timestamp ON water_measurements(station_id, timestamp DESC)",
|
|
]
|
|
|
|
else: # MySQL
|
|
stations_sql = """
|
|
CREATE TABLE IF NOT EXISTS stations (
|
|
id INT PRIMARY KEY,
|
|
station_code VARCHAR(10) UNIQUE NOT NULL,
|
|
thai_name VARCHAR(255) NOT NULL,
|
|
english_name VARCHAR(255) NOT NULL,
|
|
latitude DECIMAL(10,8),
|
|
longitude DECIMAL(11,8),
|
|
geohash VARCHAR(20),
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
|
|
)
|
|
"""
|
|
|
|
measurements_sql = """
|
|
CREATE TABLE IF NOT EXISTS water_measurements (
|
|
id BIGINT AUTO_INCREMENT PRIMARY KEY,
|
|
timestamp DATETIME NOT NULL,
|
|
station_id INT NOT NULL,
|
|
water_level DECIMAL(10,3),
|
|
discharge DECIMAL(10,2),
|
|
discharge_percent DECIMAL(5,2),
|
|
status VARCHAR(20) DEFAULT 'active',
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
FOREIGN KEY (station_id) REFERENCES stations(id),
|
|
UNIQUE KEY unique_measurement (timestamp, station_id),
|
|
INDEX idx_timestamp (timestamp),
|
|
INDEX idx_station_timestamp (station_id, timestamp)
|
|
)
|
|
"""
|
|
index_sql = []
|
|
|
|
with self.engine.begin() as conn:
|
|
conn.execute(text(stations_sql))
|
|
conn.execute(text(measurements_sql))
|
|
|
|
# Create indexes for SQLite and PostgreSQL
|
|
for index in index_sql:
|
|
conn.execute(text(index))
|
|
|
|
# Transaction is automatically committed when context manager exits
|
|
|
|
def save_measurements(self, measurements: List[Dict]) -> bool:
|
|
if not self.engine:
|
|
return False
|
|
|
|
try:
|
|
from sqlalchemy import text
|
|
|
|
with self.engine.begin() as conn:
|
|
# Insert/update stations
|
|
for measurement in measurements:
|
|
if self.db_type == "sqlite":
|
|
station_sql = """
|
|
INSERT OR REPLACE INTO stations
|
|
(id, station_code, thai_name, english_name,
|
|
latitude, longitude, geohash, updated_at)
|
|
VALUES
|
|
(:station_id, :station_code, :thai_name, :english_name,
|
|
:latitude, :longitude, :geohash, CURRENT_TIMESTAMP)
|
|
"""
|
|
elif self.db_type == "postgresql":
|
|
station_sql = """
|
|
INSERT INTO stations
|
|
(id, station_code, thai_name, english_name,
|
|
latitude, longitude, geohash, updated_at)
|
|
VALUES
|
|
(:station_id, :station_code, :thai_name, :english_name,
|
|
:latitude, :longitude, :geohash, NOW())
|
|
ON CONFLICT (id) DO UPDATE SET
|
|
thai_name = EXCLUDED.thai_name,
|
|
english_name = EXCLUDED.english_name,
|
|
latitude = EXCLUDED.latitude,
|
|
longitude = EXCLUDED.longitude,
|
|
geohash = EXCLUDED.geohash,
|
|
updated_at = NOW()
|
|
"""
|
|
else: # MySQL
|
|
station_sql = """
|
|
INSERT INTO stations
|
|
(id, station_code, thai_name, english_name,
|
|
latitude, longitude, geohash, updated_at)
|
|
VALUES
|
|
(:station_id, :station_code, :thai_name, :english_name,
|
|
:latitude, :longitude, :geohash, NOW())
|
|
ON DUPLICATE KEY UPDATE
|
|
thai_name = VALUES(thai_name),
|
|
english_name = VALUES(english_name),
|
|
latitude = VALUES(latitude),
|
|
longitude = VALUES(longitude),
|
|
geohash = VALUES(geohash),
|
|
updated_at = NOW()
|
|
"""
|
|
|
|
conn.execute(
|
|
text(station_sql),
|
|
{
|
|
"station_id": measurement["station_id"],
|
|
"station_code": measurement["station_code"],
|
|
"thai_name": measurement["station_name_th"],
|
|
"english_name": measurement["station_name_en"],
|
|
"latitude": measurement.get("latitude"),
|
|
"longitude": measurement.get("longitude"),
|
|
"geohash": measurement.get("geohash"),
|
|
},
|
|
)
|
|
|
|
# Insert measurements
|
|
for measurement in measurements:
|
|
if self.db_type == "sqlite":
|
|
measurement_sql = """
|
|
INSERT OR REPLACE INTO water_measurements
|
|
(timestamp, station_id, water_level, discharge, discharge_percent, status)
|
|
VALUES (:timestamp, :station_id, :water_level, :discharge, :discharge_percent, :status)
|
|
"""
|
|
elif self.db_type == "postgresql":
|
|
measurement_sql = """
|
|
INSERT INTO water_measurements
|
|
(timestamp, station_id, water_level, discharge, discharge_percent, status)
|
|
VALUES (:timestamp, :station_id, :water_level, :discharge, :discharge_percent, :status)
|
|
ON CONFLICT (timestamp, station_id) DO UPDATE SET
|
|
water_level = EXCLUDED.water_level,
|
|
discharge = EXCLUDED.discharge,
|
|
discharge_percent = EXCLUDED.discharge_percent,
|
|
status = EXCLUDED.status
|
|
"""
|
|
else: # MySQL
|
|
measurement_sql = """
|
|
INSERT INTO water_measurements
|
|
(timestamp, station_id, water_level, discharge, discharge_percent, status)
|
|
VALUES (:timestamp, :station_id, :water_level, :discharge, :discharge_percent, :status)
|
|
ON DUPLICATE KEY UPDATE
|
|
water_level = VALUES(water_level),
|
|
discharge = VALUES(discharge),
|
|
discharge_percent = VALUES(discharge_percent),
|
|
status = VALUES(status)
|
|
"""
|
|
|
|
conn.execute(
|
|
text(measurement_sql),
|
|
{
|
|
"timestamp": measurement["timestamp"],
|
|
"station_id": measurement["station_id"],
|
|
"water_level": measurement["water_level"],
|
|
"discharge": measurement["discharge"],
|
|
"discharge_percent": measurement["discharge_percent"],
|
|
"status": measurement["status"],
|
|
},
|
|
)
|
|
|
|
# Transaction is automatically committed when context manager exits
|
|
logging.info(
|
|
f"Successfully saved {len(measurements)} measurements to {self.db_type.upper()}"
|
|
)
|
|
return True
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error saving to {self.db_type.upper()}: {e}")
|
|
return False
|
|
|
|
def get_latest_measurements(self, limit: int = 100) -> List[Dict]:
|
|
if not self.engine:
|
|
return []
|
|
|
|
try:
|
|
from sqlalchemy import text
|
|
|
|
query = """
|
|
SELECT m.timestamp, s.station_code, s.english_name, s.thai_name,
|
|
m.water_level, m.discharge, m.discharge_percent, m.status
|
|
FROM water_measurements m
|
|
JOIN stations s ON m.station_id = s.id
|
|
INNER JOIN (
|
|
SELECT station_id, MAX(timestamp) as max_timestamp
|
|
FROM water_measurements
|
|
GROUP BY station_id
|
|
) latest ON m.station_id = latest.station_id AND m.timestamp = latest.max_timestamp
|
|
ORDER BY s.station_code
|
|
LIMIT :limit
|
|
"""
|
|
|
|
with self.engine.connect() as conn:
|
|
result = conn.execute(text(query), {"limit": limit})
|
|
measurements = []
|
|
|
|
for row in result:
|
|
measurements.append(
|
|
{
|
|
"timestamp": row[0],
|
|
"station_code": row[1],
|
|
"station_name_en": row[2],
|
|
"station_name_th": row[3],
|
|
"water_level": float(row[4])
|
|
if row[4] is not None
|
|
else None,
|
|
"discharge": float(row[5]) if row[5] is not None else None,
|
|
"discharge_percent": float(row[6])
|
|
if row[6] is not None
|
|
else None,
|
|
"status": row[7],
|
|
}
|
|
)
|
|
|
|
return measurements
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error querying {self.db_type.upper()}: {e}")
|
|
return []
|
|
|
|
def get_measurements_by_timerange(
|
|
self,
|
|
start_time: datetime.datetime,
|
|
end_time: datetime.datetime,
|
|
station_codes: Optional[List[str]] = None,
|
|
) -> List[Dict]:
|
|
if not self.engine:
|
|
return []
|
|
|
|
try:
|
|
from sqlalchemy import text
|
|
|
|
where_clause = "m.timestamp BETWEEN :start_time AND :end_time"
|
|
params = {"start_time": start_time, "end_time": end_time}
|
|
|
|
if station_codes:
|
|
placeholders = ",".join(
|
|
[f":station_{i}" for i in range(len(station_codes))]
|
|
)
|
|
where_clause += f" AND s.station_code IN ({placeholders})"
|
|
for i, code in enumerate(station_codes):
|
|
params[f"station_{i}"] = code
|
|
|
|
query = f"""
|
|
SELECT m.timestamp, s.station_code, s.english_name, s.thai_name,
|
|
m.water_level, m.discharge, m.discharge_percent, m.status
|
|
FROM water_measurements m
|
|
JOIN stations s ON m.station_id = s.id
|
|
WHERE {where_clause}
|
|
ORDER BY m.timestamp DESC, s.station_code
|
|
"""
|
|
|
|
with self.engine.connect() as conn:
|
|
result = conn.execute(text(query), params)
|
|
measurements = []
|
|
|
|
for row in result:
|
|
measurements.append(
|
|
{
|
|
"timestamp": row[0],
|
|
"station_code": row[1],
|
|
"station_name_en": row[2],
|
|
"station_name_th": row[3],
|
|
"water_level": float(row[4])
|
|
if row[4] is not None
|
|
else None,
|
|
"discharge": float(row[5]) if row[5] is not None else None,
|
|
"discharge_percent": float(row[6])
|
|
if row[6] is not None
|
|
else None,
|
|
"status": row[7],
|
|
}
|
|
)
|
|
|
|
return measurements
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error querying {self.db_type.upper()}: {e}")
|
|
return []
|
|
|
|
def get_measurements_for_date(self, target_date: datetime.datetime) -> List[Dict]:
|
|
"""Get all measurements for a specific date"""
|
|
if not self.engine:
|
|
return []
|
|
|
|
try:
|
|
from sqlalchemy import text
|
|
|
|
# Get start and end of the target date
|
|
start_of_day = target_date.replace(
|
|
hour=0, minute=0, second=0, microsecond=0
|
|
)
|
|
end_of_day = target_date.replace(
|
|
hour=23, minute=59, second=59, microsecond=999999
|
|
)
|
|
|
|
query = """
|
|
SELECT m.timestamp, m.station_id, s.station_code, s.thai_name,
|
|
m.water_level, m.discharge, m.discharge_percent, m.status
|
|
FROM water_measurements m
|
|
LEFT JOIN stations s ON m.station_id = s.id
|
|
WHERE m.timestamp >= :start_time AND m.timestamp <= :end_time
|
|
ORDER BY m.timestamp DESC
|
|
"""
|
|
|
|
with self.engine.connect() as conn:
|
|
result = conn.execute(
|
|
text(query), {"start_time": start_of_day, "end_time": end_of_day}
|
|
)
|
|
|
|
measurements = []
|
|
for row in result:
|
|
measurements.append(
|
|
{
|
|
"timestamp": row[0],
|
|
"station_id": row[1],
|
|
"station_code": row[2] or f"Station_{row[1]}",
|
|
"station_name_th": row[3] or f"Station {row[1]}",
|
|
"water_level": float(row[4])
|
|
if row[4] is not None
|
|
else None,
|
|
"discharge": float(row[5]) if row[5] is not None else None,
|
|
"discharge_percent": float(row[6])
|
|
if row[6] is not None
|
|
else None,
|
|
"status": row[7],
|
|
}
|
|
)
|
|
|
|
return measurements
|
|
|
|
except Exception as e:
|
|
logging.error(
|
|
f"Error querying {self.db_type.upper()} for date {target_date.date()}: {e}"
|
|
)
|
|
return []
|
|
|
|
@staticmethod
|
|
def _coerce_date(value) -> Optional[datetime.date]:
|
|
"""Normalize a DB-returned day value (str/date/datetime) to a date."""
|
|
if value is None:
|
|
return None
|
|
if isinstance(value, datetime.datetime):
|
|
return value.date()
|
|
if isinstance(value, datetime.date):
|
|
return value
|
|
# SQLite returns strings, e.g. '2024-09-15'
|
|
return datetime.datetime.strptime(str(value)[:10], "%Y-%m-%d").date()
|
|
|
|
def get_measurement_date_range(self) -> Optional[tuple]:
|
|
if not self.engine:
|
|
return None
|
|
|
|
try:
|
|
from sqlalchemy import text
|
|
|
|
query = "SELECT MIN(timestamp), MAX(timestamp) FROM water_measurements"
|
|
with self.engine.connect() as conn:
|
|
row = conn.execute(text(query)).fetchone()
|
|
|
|
if not row or row[0] is None:
|
|
return None
|
|
|
|
def to_datetime(value):
|
|
if isinstance(value, datetime.datetime):
|
|
return value
|
|
return datetime.datetime.fromisoformat(str(value)[:19])
|
|
|
|
return (to_datetime(row[0]), to_datetime(row[1]))
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error querying {self.db_type.upper()} date range: {e}")
|
|
return None
|
|
|
|
def get_recorded_hours_by_day(
|
|
self, start_date: datetime.date, end_date: datetime.date
|
|
) -> Optional[Dict[datetime.date, set]]:
|
|
if not self.engine:
|
|
return None
|
|
|
|
try:
|
|
from sqlalchemy import text
|
|
|
|
if self.db_type == "sqlite":
|
|
day_expr = "DATE(timestamp)"
|
|
hour_expr = "CAST(strftime('%H', timestamp) AS INTEGER)"
|
|
elif self.db_type == "postgresql":
|
|
day_expr = "CAST(timestamp AS DATE)"
|
|
hour_expr = "CAST(EXTRACT(HOUR FROM timestamp) AS INTEGER)"
|
|
else: # MySQL
|
|
day_expr = "DATE(timestamp)"
|
|
hour_expr = "HOUR(timestamp)"
|
|
|
|
query = f"""
|
|
SELECT {day_expr} AS day, {hour_expr} AS hour
|
|
FROM water_measurements
|
|
WHERE timestamp >= :start_time AND timestamp < :end_time
|
|
GROUP BY {day_expr}, {hour_expr}
|
|
"""
|
|
|
|
start_time = datetime.datetime.combine(start_date, datetime.time.min)
|
|
end_time = datetime.datetime.combine(
|
|
end_date + datetime.timedelta(days=1), datetime.time.min
|
|
)
|
|
|
|
hours_by_day: Dict[datetime.date, set] = {}
|
|
with self.engine.connect() as conn:
|
|
result = conn.execute(
|
|
text(query), {"start_time": start_time, "end_time": end_time}
|
|
)
|
|
for row in result:
|
|
day = self._coerce_date(row[0])
|
|
if day is None:
|
|
continue
|
|
hours_by_day.setdefault(day, set()).add(int(row[1]))
|
|
|
|
return hours_by_day
|
|
|
|
except Exception as e:
|
|
logging.error(
|
|
f"Error querying {self.db_type.upper()} recorded hours: {e}"
|
|
)
|
|
return None
|
|
|
|
def get_database_stats(self) -> Optional[Dict]:
|
|
if not self.engine:
|
|
return None
|
|
|
|
try:
|
|
from sqlalchemy import text
|
|
|
|
if self.db_type == "sqlite":
|
|
slot_expr = "strftime('%Y-%m-%d %H', timestamp)"
|
|
elif self.db_type == "postgresql":
|
|
slot_expr = "TO_CHAR(timestamp, 'YYYY-MM-DD HH24')"
|
|
else: # MySQL
|
|
# %-free expression: a bare % inside text() breaks as soon as the
|
|
# query gains a bind parameter (pyformat interpolation)
|
|
slot_expr = "CONCAT(DATE(timestamp), ' ', HOUR(timestamp))"
|
|
|
|
query = f"""
|
|
SELECT COUNT(*),
|
|
COUNT(DISTINCT station_id),
|
|
MIN(timestamp),
|
|
MAX(timestamp),
|
|
COUNT(DISTINCT {slot_expr})
|
|
FROM water_measurements
|
|
"""
|
|
|
|
with self.engine.connect() as conn:
|
|
row = conn.execute(text(query)).fetchone()
|
|
|
|
if not row or not row[0]:
|
|
return None
|
|
|
|
def to_datetime(value):
|
|
if isinstance(value, datetime.datetime):
|
|
return value
|
|
return datetime.datetime.fromisoformat(str(value)[:19])
|
|
|
|
first_ts = to_datetime(row[2])
|
|
last_ts = to_datetime(row[3])
|
|
# Truncate to the hour before differencing so the slot count matches
|
|
# the DISTINCT day-hour slots and coverage cannot exceed 100%
|
|
first_slot = first_ts.replace(minute=0, second=0, microsecond=0)
|
|
last_slot = last_ts.replace(minute=0, second=0, microsecond=0)
|
|
expected_hours = (
|
|
int((last_slot - first_slot).total_seconds() // 3600) + 1
|
|
)
|
|
recorded_hours = int(row[4])
|
|
coverage_percent = round(100.0 * recorded_hours / expected_hours, 1)
|
|
|
|
return {
|
|
"total_measurements": int(row[0]),
|
|
"station_count": int(row[1]),
|
|
"first_timestamp": first_ts,
|
|
"last_timestamp": last_ts,
|
|
"recorded_hours": recorded_hours,
|
|
"expected_hours": expected_hours,
|
|
"coverage_percent": coverage_percent,
|
|
}
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error querying {self.db_type.upper()} stats: {e}")
|
|
return None
|
|
|
|
|
|
# VictoriaMetrics Adapter (using Prometheus format)
|
|
class VictoriaMetricsAdapter(DatabaseAdapter):
|
|
def __init__(self, host: str = "localhost", port: int = 8428):
|
|
self.host = host
|
|
self.port = port
|
|
|
|
# Handle HTTPS URLs and reverse proxy configurations
|
|
if host.startswith(("http://", "https://")):
|
|
self.base_url = host
|
|
if port != 80 and port != 443 and not host.endswith(f":{port}"):
|
|
# Only add port if it's not standard and not already in URL
|
|
if "://" in host and ":" not in host.split("://")[1]:
|
|
self.base_url = f"{host}:{port}"
|
|
else:
|
|
# Default to HTTP for localhost, HTTPS for remote hosts
|
|
protocol = (
|
|
"https"
|
|
if host != "localhost" and not host.startswith("127.")
|
|
else "http"
|
|
)
|
|
if (protocol == "https" and port == 443) or (
|
|
protocol == "http" and port == 80
|
|
):
|
|
self.base_url = f"{protocol}://{host}"
|
|
else:
|
|
self.base_url = f"{protocol}://{host}:{port}"
|
|
|
|
@staticmethod
|
|
def _escape_label(value) -> str:
|
|
"""Escape a Prometheus label value per the exposition format spec.
|
|
|
|
Station names include arbitrary Thai text (and could be set via the API),
|
|
so backslashes, double-quotes and newlines must be escaped to avoid
|
|
producing malformed or injected exposition lines.
|
|
"""
|
|
return str(value).replace("\\", "\\\\").replace('"', '\\"').replace("\n", "\\n")
|
|
|
|
@staticmethod
|
|
def _metric_value(value) -> Optional[float]:
|
|
"""Coerce a numeric field to float, or None if it isn't a valid number."""
|
|
try:
|
|
return float(value)
|
|
except (TypeError, ValueError):
|
|
return None
|
|
|
|
def connect(self):
|
|
try:
|
|
import requests
|
|
|
|
# Test connection with SSL verification and timeout
|
|
response = requests.get(
|
|
f"{self.base_url}/api/v1/status/config",
|
|
timeout=10,
|
|
verify=True, # Enable SSL verification for HTTPS
|
|
)
|
|
if response.status_code == 200:
|
|
logging.info(
|
|
f"Connected to VictoriaMetrics successfully at {self.base_url}"
|
|
)
|
|
return True
|
|
else:
|
|
logging.error(
|
|
f"VictoriaMetrics connection failed: {response.status_code}"
|
|
)
|
|
return False
|
|
except requests.exceptions.SSLError as e:
|
|
logging.error(f"SSL error connecting to VictoriaMetrics: {e}")
|
|
return False
|
|
except requests.exceptions.ConnectionError as e:
|
|
logging.error(f"Connection error to VictoriaMetrics: {e}")
|
|
return False
|
|
except Exception as e:
|
|
logging.error(f"Failed to connect to VictoriaMetrics: {e}")
|
|
return False
|
|
|
|
def save_measurements(self, measurements: List[Dict]) -> bool:
|
|
try:
|
|
import requests
|
|
|
|
# Convert to Prometheus format
|
|
metrics_data = []
|
|
timestamp_ms = int(datetime.datetime.now().timestamp() * 1000)
|
|
|
|
for measurement in measurements:
|
|
# Escape label values once per measurement (untrusted Thai/English names).
|
|
labels = (
|
|
f'station_code="{self._escape_label(measurement["station_code"])}",'
|
|
f'station_name_en="{self._escape_label(measurement["station_name_en"])}",'
|
|
f'station_name_th="{self._escape_label(measurement["station_name_th"])}"'
|
|
)
|
|
|
|
# Water level metric
|
|
water_level = self._metric_value(measurement.get("water_level"))
|
|
if water_level is not None:
|
|
metrics_data.append(
|
|
f"water_level{{{labels}}} {water_level} {timestamp_ms}"
|
|
)
|
|
|
|
# Discharge metric
|
|
discharge = self._metric_value(measurement.get("discharge"))
|
|
if discharge is not None:
|
|
metrics_data.append(
|
|
f"water_discharge{{{labels}}} {discharge} {timestamp_ms}"
|
|
)
|
|
|
|
# Discharge percentage metric
|
|
discharge_percent = self._metric_value(
|
|
measurement.get("discharge_percent")
|
|
)
|
|
if discharge_percent is not None:
|
|
metrics_data.append(
|
|
f"water_discharge_percent{{{labels}}} {discharge_percent} {timestamp_ms}"
|
|
)
|
|
|
|
# Send to VictoriaMetrics
|
|
data = "\n".join(metrics_data)
|
|
response = requests.post(
|
|
f"{self.base_url}/api/v1/import/prometheus",
|
|
data=data,
|
|
headers={"Content-Type": "text/plain"},
|
|
timeout=30,
|
|
)
|
|
|
|
if response.status_code == 204:
|
|
logging.info(
|
|
f"Successfully sent {len(measurements)} measurements to VictoriaMetrics"
|
|
)
|
|
return True
|
|
else:
|
|
logging.error(
|
|
f"VictoriaMetrics import failed: {response.status_code} - {response.text}"
|
|
)
|
|
return False
|
|
|
|
except Exception as e:
|
|
logging.error(f"Error sending to VictoriaMetrics: {e}")
|
|
return False
|
|
|
|
def get_latest_measurements(self, limit: int = 100) -> List[Dict]:
|
|
# VictoriaMetrics queries would be implemented here
|
|
# This is a simplified version
|
|
logging.warning(
|
|
"get_latest_measurements not fully implemented for VictoriaMetrics"
|
|
)
|
|
return []
|
|
|
|
def get_measurements_by_timerange(
|
|
self,
|
|
start_time: datetime.datetime,
|
|
end_time: datetime.datetime,
|
|
station_codes: Optional[List[str]] = None,
|
|
) -> List[Dict]:
|
|
# VictoriaMetrics range queries would be implemented here
|
|
logging.warning(
|
|
"get_measurements_by_timerange not fully implemented for VictoriaMetrics"
|
|
)
|
|
return []
|
|
|
|
def get_measurements_for_date(self, target_date: datetime.datetime) -> List[Dict]:
|
|
"""Get all measurements for a specific date"""
|
|
logging.warning(
|
|
"get_measurements_for_date not fully implemented for VictoriaMetrics"
|
|
)
|
|
return []
|
|
|
|
|
|
# Factory function to create appropriate adapter
|
|
def create_database_adapter(db_type: str, **kwargs) -> DatabaseAdapter:
|
|
"""
|
|
Factory function to create database adapter
|
|
|
|
Args:
|
|
db_type: 'influxdb', 'mysql', 'postgresql', 'sqlite', or 'victoriametrics'
|
|
**kwargs: Database-specific connection parameters
|
|
"""
|
|
db_type = db_type.lower()
|
|
|
|
if db_type == "influxdb":
|
|
return InfluxDBAdapter(**kwargs)
|
|
elif db_type == "mysql":
|
|
return SQLAdapter(db_type="mysql", **kwargs)
|
|
elif db_type == "postgresql":
|
|
return SQLAdapter(db_type="postgresql", **kwargs)
|
|
elif db_type == "sqlite":
|
|
return SQLAdapter(db_type="sqlite", **kwargs)
|
|
elif db_type == "victoriametrics":
|
|
return VictoriaMetricsAdapter(**kwargs)
|
|
else:
|
|
raise ValueError(f"Unsupported database type: {db_type}")
|