fix(sync): align customer backend behavior
This commit is contained in:
@@ -14,6 +14,7 @@ from app.services.scheme_management import (
|
||||
store_scheme_info,
|
||||
)
|
||||
from app.services.tjnetwork import get_all_scada_info
|
||||
from app.services.time_api import extract_date, parse_utc_time, utc_now
|
||||
|
||||
|
||||
def run_burst_detection(
|
||||
@@ -241,7 +242,7 @@ def list_burst_detection_schemes(
|
||||
network: str,
|
||||
query_date: datetime | str | None = None,
|
||||
) -> list[dict[str, Any]]:
|
||||
parsed_date = _to_datetime(query_date).date() if query_date is not None else None
|
||||
parsed_date = extract_date(query_date, field_name="query_date") if query_date is not None else None
|
||||
return query_burst_detection_schemes(
|
||||
name=network,
|
||||
network=network,
|
||||
@@ -269,7 +270,7 @@ def _store_burst_detection_scheme(
|
||||
if scheme_name_exists(network, scheme_name):
|
||||
raise ValueError(f"方案名称已存在: {scheme_name}")
|
||||
|
||||
now_iso = datetime.now().isoformat()
|
||||
now_iso = utc_now().isoformat()
|
||||
scheme_detail = {
|
||||
"network": network,
|
||||
"sensor_nodes": payload.get("sensor_nodes", []),
|
||||
@@ -426,6 +427,4 @@ def _build_observed_pressure_from_scada(
|
||||
|
||||
|
||||
def _to_datetime(value: datetime | str) -> datetime:
|
||||
if isinstance(value, datetime):
|
||||
return value
|
||||
return datetime.fromisoformat(value)
|
||||
return parse_utc_time(value)
|
||||
|
||||
@@ -15,6 +15,7 @@ from app.services.scheme_management import (
|
||||
store_scheme_info,
|
||||
)
|
||||
from app.services.tjnetwork import dump_inp, get_all_scada_info
|
||||
from app.services.time_api import extract_date, parse_utc_time, utc_now
|
||||
|
||||
SeriesInput = pd.Series | dict[str, Any] | list[dict[str, Any]]
|
||||
FLOW_SCADA_TYPES = {"pipe_flow", "flow", "demand"}
|
||||
@@ -301,7 +302,7 @@ def run_burst_location_by_network(
|
||||
def list_burst_location_schemes(
|
||||
network: str, query_date: datetime | str | None = None
|
||||
) -> list[dict[str, Any]]:
|
||||
parsed_date = _to_datetime(query_date).date() if query_date is not None else None
|
||||
parsed_date = extract_date(query_date, field_name="query_date") if query_date is not None else None
|
||||
return query_burst_location_schemes(
|
||||
name=network, network=network, query_date=parsed_date
|
||||
)
|
||||
@@ -327,7 +328,7 @@ def _store_burst_scheme(
|
||||
if scheme_name_exists(network, scheme_name):
|
||||
raise ValueError(f"方案名称已存在: {scheme_name}")
|
||||
|
||||
now_iso = datetime.now().isoformat()
|
||||
now_iso = utc_now().isoformat()
|
||||
scheme_detail = {
|
||||
"network": network,
|
||||
"pressure_scada_ids": payload.get("pressure_scada_ids", []),
|
||||
@@ -641,9 +642,7 @@ def _dedupe_ids(ids: list[str] | None) -> list[str]:
|
||||
|
||||
|
||||
def _to_datetime(value: datetime | str) -> datetime:
|
||||
if isinstance(value, datetime):
|
||||
return value
|
||||
return datetime.fromisoformat(value)
|
||||
return parse_utc_time(value)
|
||||
|
||||
|
||||
def _prepare_burst_inp(network: str) -> str:
|
||||
|
||||
@@ -23,6 +23,7 @@ from app.services.tjnetwork import (
|
||||
get_network_link_nodes,
|
||||
get_network_node_coords,
|
||||
)
|
||||
from app.services.time_api import extract_date, parse_utc_time, utc_now
|
||||
|
||||
DEFAULT_N_WORKERS = max(1, min((os.cpu_count() or 1) - 1, 4))
|
||||
|
||||
@@ -119,7 +120,7 @@ def run_leakage_identification(
|
||||
scheme_start_time = (
|
||||
_to_datetime(scada_start).isoformat()
|
||||
if scada_start is not None
|
||||
else datetime.now().isoformat()
|
||||
else utc_now().isoformat()
|
||||
)
|
||||
scheme_detail = {
|
||||
"network": network,
|
||||
@@ -177,7 +178,7 @@ def run_leakage_identification(
|
||||
def list_leakage_identify_schemes(
|
||||
network: str, query_date: datetime | str | None = None
|
||||
) -> list[dict[str, Any]]:
|
||||
parsed_date = _to_datetime(query_date).date() if query_date is not None else None
|
||||
parsed_date = extract_date(query_date, field_name="query_date") if query_date is not None else None
|
||||
return query_leakage_identify_schemes(
|
||||
name=network, network=network, query_date=parsed_date
|
||||
)
|
||||
@@ -509,9 +510,7 @@ def _build_observed_pressure_from_scada(
|
||||
|
||||
|
||||
def _to_datetime(value: datetime | str) -> datetime:
|
||||
if isinstance(value, datetime):
|
||||
return value
|
||||
return datetime.fromisoformat(value)
|
||||
return parse_utc_time(value)
|
||||
|
||||
|
||||
def _prepare_leakage_inp(network: str) -> str:
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
import os
|
||||
from app.core.config import settings
|
||||
|
||||
# 从环境变量 NETWORK_NAME 读取
|
||||
name = os.getenv("NETWORK_NAME")
|
||||
name = settings.NETWORK_NAME
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import ast
|
||||
import json
|
||||
from datetime import date
|
||||
from datetime import date, datetime
|
||||
|
||||
import geopandas as gpd
|
||||
import pandas as pd
|
||||
@@ -8,6 +8,7 @@ import psycopg
|
||||
from sqlalchemy import create_engine
|
||||
|
||||
from app.core.config import get_pgconn_string
|
||||
from app.services.time_api import parse_utc_time
|
||||
|
||||
|
||||
# 2025/03/23
|
||||
@@ -89,7 +90,7 @@ def store_scheme_info(
|
||||
scheme_name: str,
|
||||
scheme_type: str,
|
||||
username: str,
|
||||
scheme_start_time: str,
|
||||
scheme_start_time: datetime | str,
|
||||
scheme_detail: dict,
|
||||
):
|
||||
"""
|
||||
@@ -112,13 +113,16 @@ def store_scheme_info(
|
||||
"""
|
||||
# 将字典转换为 JSON 字符串
|
||||
scheme_detail_json = json.dumps(scheme_detail)
|
||||
normalized_scheme_start_time = parse_utc_time(
|
||||
scheme_start_time, field_name="scheme_start_time"
|
||||
)
|
||||
cur.execute(
|
||||
sql,
|
||||
(
|
||||
scheme_name,
|
||||
scheme_type,
|
||||
username,
|
||||
scheme_start_time,
|
||||
normalized_scheme_start_time,
|
||||
scheme_detail_json,
|
||||
),
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user