commit 2eda0a0887e91f730e11bdab658c561ec6e15e66 Author: Nixevol Date: Thu Sep 24 06:16:41 2026 +0800 feat: 导入 InterferenceETL 初始源码 diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..c8a679b --- /dev/null +++ b/.dockerignore @@ -0,0 +1,3 @@ +* +!Dockerfile +!requirements.txt diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..bf89f68 --- /dev/null +++ b/.gitignore @@ -0,0 +1,9 @@ +__pycache__/ +*.py[cod] +.pytest_cache/ +.venv/ +.env +runtime_config.py +output/ +dist/ +aidocs/ diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..4edffa1 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,8 @@ +FROM python:3.13.11-slim + +COPY requirements.txt /tmp/requirements.txt +RUN pip install --no-cache-dir -r /tmp/requirements.txt + +WORKDIR /workspace + +CMD ["python", "main.py"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..2b0b72c --- /dev/null +++ b/README.md @@ -0,0 +1,123 @@ +# InterferenceETL + +InterferenceETL follows the newest interference KPI time group in Metrix storage. It reads seven ZIP/XLSX source types, converts each `Sheet0` to CSV, and creates one merged high-interference summary. + +Source file names are identified by one of the seven fixed prefixes plus a final `_YYYYMMDDHHMMHHMM.zip`; text between the prefix and time is unrestricted. Only the first `YYYYMMDDHHMM` start time is used for grouping. Start times are rounded down to natural 15-minute boundaries, so `12:00` through `12:14` belong to `12:00`. This also works unchanged when all providers switch together to hourly or daily delivery. + +Every run follows only the group containing the newest source start time. It never falls back or backfills an older group. If any source type is missing, the script prints `status=waiting` with the missing types and exits successfully. When one source has multiple files in the group, its latest start time wins. + +## Output + +Each run writes one window directory: + +```text +output/20260731100000/ +├── converted/ +│ ├── 5G下FDD干扰监控_2026073110001100.csv +│ ├── 5G干扰监控_2026073110001100.csv +│ └── ... seven source CSV files +├── interference_summary_20260731100000.csv +└── manifest.json +``` + +For every processed group, the script also reads the latest filename-dated XLSX from each configured CellData directory. It builds CellData CGI as `460-00-{eNB/gNB}-{CI}` and adds coordinates plus direction angle to matching interference rows. Unmatched rows keep empty coordinates and use `0` for azimuth. + +The summary columns are: + +```text +metric_time,network_type,cgi,cell_name,interference_dbm,prev_interference_dbm,longitude,latitude,azimuth,nearby_count,nearby_26g,nearby_700m,nearby_tdd,nearby_fdd,prev_nearby_count,prev_nearby_26g,prev_nearby_700m,prev_nearby_tdd,prev_nearby_fdd +``` + +All selected workbook rows must belong to the selected natural 15-minute group. Every summary row receives the same normalized `metric_time`; a row outside the group fails the run before database or history output. + +`network_type` is derived from the seven source types: `5G干扰监控` is `2.6G`, `700M干扰监控` is `700M`, `SDR_TDD干扰监控` and `反开RD干扰监控` are `TDD`, and `SDR_FDD干扰监控`, `5G下FDD干扰监控`, and `700M下FDD干扰监控` are `FDD`. The summary and database keep only high-interference rows. Thresholds keep the original directory grouping and are independent of the `network_type` label: the two `5G...` sources use `>= -107 dBm`, and the other five sources use `>= -110 dBm`. Converted source CSV files remain full source conversions. + +`nearby_count` is the number of other high-interference cells from the same selected time group within 1 km, across all network types. The cell itself is excluded, the 1 km boundary is included, and rows without coordinates use `0`. Candidate cells are found through a 1 km spatial grid index and confirmed with exact Haversine distance. + +`nearby_26g`, `nearby_700m`, `nearby_tdd`, and `nearby_fdd` split the same neighbors by the neighbor's `network_type`. Each field defaults to `0`, and the four values always add up to `nearby_count`. + +When multiple retained source rows have the same CGI, the summary keeps the row with the numerically largest interference value. Equal values keep the first row encountered. Deduplication happens before nearby-cell counting and database insertion; converted source CSVs remain complete. + +`prev_interference_dbm`, `prev_nearby_count`, and the four `prev_nearby_26g/700m/tdd/fdd` fields come from the previously retained database time, matched by CGI before the current batch replaces old rows. They represent the same cell's previous-period interference value and nearby high-interference counts. On the first run, or when a CGI did not exist in the previous period, these fields are empty in CSV/history output and `NULL` in MySQL. + +The script never modifies or deletes source storage files. Before downloading source ZIP files or CellData, it compares the normalized target time with the database. A target already present, or older than the database, is not processed again. + +The database table contains normalized `metric_time DATETIME`, `network_type VARCHAR(16)`, `interference_dbm DECIMAL(10,3)`, nullable `prev_interference_dbm DECIMAL(10,3)`, nullable `longitude` and `latitude`, `azimuth DECIMAL(6,2) NOT NULL DEFAULT 0`, the five `NOT NULL DEFAULT 0` INT count columns (`nearby_count`, `nearby_26g`, `nearby_700m`, `nearby_tdd`, `nearby_fdd`), and the five nullable INT previous-period count columns (`prev_nearby_count`, `prev_nearby_26g`, `prev_nearby_700m`, `prev_nearby_tdd`, `prev_nearby_fdd`). A successful transaction replaces the target batch and deletes every other database time, so the table retains only the latest processed group. The previous-period values are read before this replacement transaction. + +After the database succeeds, the same nineteen result columns are uploaded as GBK CSV to: + +```text +/网优日常优化数据文档/(勿删)干扰定时小时指标/干扰历史数据/YYYY-MM-DD/干扰数据处理结果_YYYYMMDDHHMMSS.csv +``` + +The filename time comes from normalized `metric_time`, never the server clock. If the database already contains the target time but its history CSV is missing or empty, the script exports that time from the database and repairs the history file without reprocessing source data. + +Only the final history CSV uploaded to Metrix Storage uses GBK without a UTF-8 BOM. Source non-breaking spaces (`U+00A0`) are normalized to ordinary spaces in this history file because GBK cannot encode them; other unsupported characters still fail visibly. The seven local converted CSV files and local merged summary remain UTF-8 with BOM, and `manifest.json` remains UTF-8 JSON. + +CGI is generated with fixed rules: + +- NR (`5G干扰监控`, `700M干扰监控`): `{gNBplmn}-{gNBId}-{cellId}`. +- 4G (the other five source types): `460-00-{eNodeBID}-{小区ID}`, using each source schema's actual equivalent column names. + +The Metrix API address, API Token, and database connection ID are constants in `runtime_config.py`; they are not project environment variables. Copy `runtime_config.example.py` to `runtime_config.py` and fill in the actual Token and the `conn_id` of `ShareMySQL`. The fixed database is `interference_etl`, the fixed table is `interference_hourly_summary`, and the script creates both automatically when they do not exist. + +The five CellData directories are fixed in `CELL_DATA_DIRECTORIES`. In each directory, only XLSX files ending with a valid `YYYYMMDD.xlsx` date are considered, and the latest date is selected. CellData column `方向角` maps to summary/database column `azimuth`; an empty source value becomes `0`. + +## Local development + +```powershell +python -m pip install -r requirements.txt +python -m unittest discover -s tests -v +python main.py --source-dir C:\path\to\mock-or-exported-tree --output-dir output --no-database +``` + +Use `--window 2026073110001100` to select the natural 15-minute group containing that source start time. Without it, the newest observed source start time determines the group. Incomplete groups wait and never fall back. + +## Metrix Script Management + +The target server is offline. Build the dependency image locally, then generate a self-contained Script Management upload ZIP: + +```powershell +docker build -t interference-etl-runtime:1.1 . +python scripts\build_offline_package.py +``` + +The generated `dist/InterferenceETL-offline-1.1.zip` contains `main.py`, `runtime_config.py`, a ZIP execution entry point, and a `vendor/` directory containing `openpyxl` and `et_xmlfile`. The builder verifies imports and ZIP execution using the standard `python:3.13.11-slim` image, so the uploaded workspace does not need an online `pip install` or a custom runtime image. The build fails when `runtime_config.py` is missing. + +Create the project, then upload and extract the ZIP in its Script Management workspace. The normal file tree should contain `main.py` at `/workspace/main.py`. Use these project settings: + +```text +Name: InterferenceETL +Language: python +Base image: python:3.13.11-slim +Network: bridge +Run command: python main.py +Timeout: 1800 seconds +``` + +If the server keeps the ZIP as one file instead of extracting it, use `python InterferenceETL-offline-1.1.zip` as the run command. The ZIP contains `__main__.py` for this mode. + +Configure only run-specific source and output settings in the project environment: + +```json +{ + "METRIX_STORAGE_ID": "stg_4d9a910d72", + "INTERFERENCE_SOURCE_ROOT": "/网优日常优化数据文档/(勿删)干扰定时小时指标", + "INTERFERENCE_OUTPUT_DIR": "/workspace/output", + "INTERFERENCE_LOOKBACK_DAYS": "3" +} +``` + +Both storage reads and database writes use the Metrix API at `http://188.5.127.115:18271`. SSH remains the deployment and operational channel for uploading files, starting runs, and reading logs. + +## CLI options + +```text +--source-dir PATH Read a local directory tree instead of Metrix API +--output-dir PATH Output root, default output or INTERFERENCE_OUTPUT_DIR +--window WINDOW Select the natural 15-minute group for a 16-digit source window +--lookback-days N Number of newest date directories scanned, default 3 +--storage-id ID Metrix storage connection ID +--root PATH Source directory in Metrix storage +--no-database Generate CSV files without writing the database +``` diff --git a/docs/assumptions.md b/docs/assumptions.md new file mode 100644 index 0000000..07d2775 --- /dev/null +++ b/docs/assumptions.md @@ -0,0 +1,26 @@ +# Assumptions and pending decisions + +## Implemented assumptions + +- The source directory contains seven expected source types for each complete hour. +- The newest usable hour is the newest 16-digit filename window for which all seven types exist. A newer incomplete hour is logged and skipped. +- Each ZIP contains exactly one XLSX member. The member is selected by extension because the two 700M ZIP series use a legacy filename encoding. +- Only `Sheet0` is converted and merged. The `指标(计数器)` sheet is metadata and is not included in the summary. +- NR CGI is derived as `{gNBplmn}-{gNBId}-{cellId}`. +- 4G CGI is derived as `460-00-{eNodeBID}-{小区ID}`, using the equivalent node and cell column names in each source schema. +- Only history CSV files uploaded to Metrix Storage use GBK without a UTF-8 BOM. Local converted and summary CSV files remain UTF-8 with BOM; `manifest.json` remains UTF-8 JSON. +- Source files are read-only. The script writes only below its configured output directory. +- MySQL is the scheduled-run output. The table uses one `metric_time DATETIME` column containing the source KPI start time. +- A successful run transactionally refreshes the selected hour and deletes rows for all other hours. A selected hour older than the newest database hour is rejected to prevent data rollback. + +## Pending user decisions + +- Whether the seven converted full CSV files must be retained after successful database import. +- Whether the `指标(计数器)` metadata sheet must also be stored. +- Schedule minute and expected source-file arrival delay. A safe initial suggestion is hourly at minute 20. +- Historical backfill range and retention policy for generated CSV files. +- Whether an incomplete newest hour should only warn, fail the run, or wait and retry. + +## Mock boundary + +Automated tests generate seven in-memory XLSX/ZIP sources plus a newer incomplete hour. The mock verifies both CGI rules, complete-hour fallback, schema validation, CSV conversion, summary merging, cross-midnight windows, transactional latest-hour replacement, and database rollback protection. diff --git a/main.py b/main.py new file mode 100644 index 0000000..f8d11bc --- /dev/null +++ b/main.py @@ -0,0 +1,1235 @@ +from __future__ import annotations + +import argparse +import csv +from decimal import Decimal, InvalidOperation +import hashlib +import io +import json +import math +import os +from dataclasses import dataclass +from datetime import datetime, timezone +from pathlib import Path +import posixpath +import re +import shutil +import sys +import time +from typing import Protocol +from urllib.error import HTTPError, URLError +from urllib.parse import quote, urlencode +from urllib.request import ProxyHandler, Request, build_opener +import warnings +import zipfile + + +VENDOR_DIR = Path(__file__).resolve().parent / "vendor" +sys.path.insert(0, str(VENDOR_DIR)) + +from openpyxl import load_workbook + +try: + from runtime_config import METRIX_API_BASE_URL, METRIX_API_TOKEN, METRIX_DATABASE_CONNECTION_ID +except ImportError: + METRIX_API_BASE_URL = "http://188.5.127.115:18271" + METRIX_API_TOKEN = "" + METRIX_DATABASE_CONNECTION_ID = "" + + +EXPECTED_TYPES = ( + "5G下FDD干扰监控", + "5G干扰监控", + "700M下FDD干扰监控", + "700M干扰监控", + "SDR_FDD干扰监控", + "SDR_TDD干扰监控", + "反开RD干扰监控", +) + +NETWORK_TYPE_BY_SOURCE = { + "5G下FDD干扰监控": "FDD", + "5G干扰监控": "2.6G", + "700M下FDD干扰监控": "FDD", + "700M干扰监控": "700M", + "SDR_FDD干扰监控": "FDD", + "SDR_TDD干扰监控": "TDD", + "反开RD干扰监控": "TDD", +} + +# Thresholds keep the original band grouping (the two 5G-directory sources use +# -107, all others -110); they are independent of the network_type label. +MIN_INTERFERENCE_BY_SOURCE = { + "5G下FDD干扰监控": Decimal("-107"), + "5G干扰监控": Decimal("-107"), + "700M下FDD干扰监控": Decimal("-110"), + "700M干扰监控": Decimal("-110"), + "SDR_FDD干扰监控": Decimal("-110"), + "SDR_TDD干扰监控": Decimal("-110"), + "反开RD干扰监控": Decimal("-110"), +} + +EARTH_RADIUS_KM = 6371.0088 +NEARBY_RADIUS_KM = 1.0 + +NEARBY_COLUMN_BY_NETWORK_TYPE = { + "2.6G": "nearby_26g", + "700M": "nearby_700m", + "TDD": "nearby_tdd", + "FDD": "nearby_fdd", +} + +DEFAULT_STORAGE_ID = "stg_4d9a910d72" +DEFAULT_ROOT = "/网优日常优化数据文档/(勿删)干扰定时小时指标" +CELL_DATA_DIRECTORIES = ( + "/网优日常优化数据文档/日常性能报表/2026年/5G/5G小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/700M/700M小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/5G_反开/反开小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/TDD_LTE/LTE基础信息数据/LTE_小区信息表", + "/网优日常优化数据文档/日常性能报表/2026年/FDD_LTE/FDD基础信息数据/FDD小区信息表", +) +FILE_TIME_RE = re.compile(r"_(?P\d{16})\.zip$", re.IGNORECASE) +CELL_DATA_FILE_RE = re.compile(r"(?P20\d{6})\.xlsx$", re.IGNORECASE) + +LTE_PLMN = "460-00" +DATABASE_NAME = "interference_etl" +DATABASE_TABLE = "interference_hourly_summary" +DEFAULT_HISTORY_ROOT = f"{DEFAULT_ROOT}/干扰历史数据" +LOCAL_CSV_ENCODING = "utf-8-sig" +HISTORY_CSV_ENCODING = "gbk" + +HEADER_FDD = ( + "开始时间", + "粒度", + "子网ID", + "子网名称", + "网元ID", + "管理网元", + "eNodeB CUID", + "eNodeB CU名称", + "LTEID", + "LTE名称", + "E-UTRAN FDD小区ID", + "E-UTRAN FDD小区名称", + "cellId", + "eNodeBId", + "载波平均噪声干扰(dBm)", + "集团-上下行总业务量(GB)", + "RRC连接建立最大用户数", +) + +HEADER_NR = ( + "开始时间", + "粒度", + "子网ID", + "子网名称", + "网元ID", + "管理网元", + "gNB CU-CP功能配置ID", + "gNB CU-CP功能配置名称", + "CU小区配置ID", + "CU小区配置名称", + "cellId", + "duMeMoId", + "gNBId", + "gNBIdLength", + "gNBplmn", + "masterOperatorId", + "nrCarrierGroupId", + "nrPhysicalCellDUId", + "小区上行平均干扰电平(dBm)", + "5G上下行总流量(上行PDCP PDU数据量+下行PDCP成功发送数据量)(GB)", + "RRC连接最大连接用户数", +) + +HEADER_SDR = ( + "开始时间", + "粒度", + "子网ID", + "子网名称", + "网元ID", + "管理网元", + "eNodeBID", + "eNodeB名称", + "小区ID", + "小区名称", + "载波平均噪声干扰(dBm)", + "集团-上下行总业务量(GB)", + "RRC连接建立最大用户数", +) + +HEADER_RD = ( + "开始时间", + "粒度", + "子网ID", + "子网名称", + "网元ID", + "管理网元", + "eNodeB CUID", + "eNodeB CU名称", + "LTEID", + "LTE名称", + "E-UTRAN TDD小区ID", + "E-UTRAN TDD小区名称", + "cellId", + "eNodeBId", + "载波平均噪声干扰(dBm)", + "上下行总业务量(GB)", + "RRC连接建立最大用户数", +) + +EXPECTED_HEADERS = { + "5G下FDD干扰监控": HEADER_FDD, + "5G干扰监控": HEADER_NR, + "700M下FDD干扰监控": HEADER_FDD, + "700M干扰监控": HEADER_NR, + "SDR_FDD干扰监控": HEADER_SDR, + "SDR_TDD干扰监控": HEADER_SDR, + "反开RD干扰监控": HEADER_RD, +} + +SUMMARY_HEADER = ( + "metric_time", + "network_type", + "cgi", + "cell_name", + "interference_dbm", + "prev_interference_dbm", + "longitude", + "latitude", + "azimuth", + "nearby_count", + "nearby_26g", + "nearby_700m", + "nearby_tdd", + "nearby_fdd", + "prev_nearby_count", + "prev_nearby_26g", + "prev_nearby_700m", + "prev_nearby_tdd", + "prev_nearby_fdd", +) + +class ProcessingError(RuntimeError): + pass + + +@dataclass(frozen=True) +class Candidate: + source_type: str + window: str + path: str + size: int = 0 + + +class Source(Protocol): + def candidates(self, lookback_days: int) -> list[Candidate]: ... + + def download(self, path: str) -> bytes: ... + + def cell_data_workbooks(self) -> list[tuple[str, bytes]]: ... + + +class SummaryStore(Protocol): + def latest_metric_time(self) -> datetime | None: ... + + def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: ... + + def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]: ... + + +class HistoryStore(Protocol): + def exists(self, metric_time: datetime) -> bool: ... + + def upload(self, metric_time: datetime, payload: bytes) -> str: ... + + +class MetrixApiClient: + def __init__(self, base_url: str, token: str) -> None: + if not token: + raise ProcessingError("METRIX_API_TOKEN is required for Metrix API access") + self.base_url = base_url.rstrip("/") + self.token = token + self.opener = build_opener(ProxyHandler({})) + + def get_bytes(self, endpoint: str, query: dict[str, object] | None = None, timeout: int = 30) -> bytes: + return self._request("GET", endpoint, query=query, timeout=timeout) + + def get_json(self, endpoint: str, query: dict[str, object] | None = None, timeout: int = 30) -> object: + return self._decode_json(self.get_bytes(endpoint, query, timeout)) + + def post_json(self, endpoint: str, payload: dict[str, object], timeout: int = 30) -> object: + body = json.dumps(payload, ensure_ascii=False).encode("utf-8") + return self._decode_json( + self._request("POST", endpoint, body=body, content_type="application/json", timeout=timeout) + ) + + def post_file( + self, + endpoint: str, + query: dict[str, object], + filename: str, + payload: bytes, + timeout: int = 120, + ) -> object: + boundary = f"----MetrixBoundary{time.time_ns()}" + disposition = f'Content-Disposition: form-data; name="file"; filename="{filename}"\r\n' + body = ( + f"--{boundary}\r\n".encode() + + disposition.encode("utf-8") + + b"Content-Type: text/csv\r\n\r\n" + + payload + + f"\r\n--{boundary}--\r\n".encode() + ) + return self._decode_json( + self._request( + "POST", + endpoint, + query=query, + body=body, + content_type=f"multipart/form-data; boundary={boundary}", + timeout=timeout, + ) + ) + + def _request( + self, + method: str, + endpoint: str, + query: dict[str, object] | None = None, + body: bytes | None = None, + content_type: str = "", + timeout: int = 30, + ) -> bytes: + url = f"{self.base_url}{endpoint}" + if query: + url = f"{url}?{urlencode(query)}" + headers = {"Authorization": f"Bearer {self.token}", "Accept": "application/json"} + if content_type: + headers["Content-Type"] = content_type + request = Request(url, data=body, headers=headers, method=method) + last_error: Exception | None = None + for attempt in range(3): + try: + with self.opener.open(request, timeout=timeout) as response: + return response.read() + except HTTPError as exc: + detail = exc.read().decode("utf-8", "replace")[:1000] + if exc.code < 500: + raise ProcessingError(f"Metrix API {exc.code}: {detail}") from exc + last_error = exc + except (URLError, TimeoutError) as exc: + last_error = exc + if attempt < 2: + time.sleep(2**attempt) + raise ProcessingError(f"Metrix API request failed: {last_error}") + + @staticmethod + def _decode_json(payload: bytes) -> object: + try: + return json.loads(payload.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError) as exc: + raise ProcessingError("Metrix API returned invalid JSON") from exc + + +class ApiSummaryStore: + def __init__(self, client: MetrixApiClient, connection_id: str) -> None: + if not connection_id: + raise ProcessingError("METRIX_DATABASE_CONNECTION_ID is required for database output") + self.client = client + self.connection_id = connection_id + self.table = DATABASE_TABLE + + def latest_metric_time(self) -> datetime | None: + self._ensure_schema() + payload = self._query( + f"SELECT MAX(metric_time) AS latest_time FROM `{self.table}`", + database=DATABASE_NAME, + ) + rows = payload.get("rows", []) + value = rows[0].get("latest_time") if rows else None + return datetime.fromisoformat(str(value)) if value else None + + def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: + if not rows: + raise ProcessingError("Refusing to replace database data with an empty batch") + metric_times = {datetime.strptime(row["metric_time"], "%Y-%m-%d %H:%M:%S") for row in rows} + if len(metric_times) != 1: + raise ProcessingError("Database batch must contain exactly one metric time") + metric_time = next(iter(metric_times)) + latest_time = self.latest_metric_time() + if latest_time is not None and latest_time > metric_time: + raise ProcessingError( + f"Database already contains newer metric time {latest_time:%Y-%m-%d %H:%M:%S}; " + f"refusing to replace it with {metric_time:%Y-%m-%d %H:%M:%S}" + ) + + metric_literal = f"'{metric_time:%Y-%m-%d %H:%M:%S}'" + values = ",\n".join(self._row_values(metric_literal, row) for row in rows) + script_payload = self.client.post_json( + self._endpoint("run-script"), + { + "content": f""" + START TRANSACTION; + DELETE FROM `{self.table}` WHERE metric_time = {metric_literal}; + INSERT INTO `{self.table}` + (metric_time, network_type, cgi, cell_name, interference_dbm, prev_interference_dbm, longitude, latitude, azimuth, nearby_count, nearby_26g, nearby_700m, nearby_tdd, nearby_fdd, prev_nearby_count, prev_nearby_26g, prev_nearby_700m, prev_nearby_tdd, prev_nearby_fdd) + VALUES + {values}; + DELETE FROM `{self.table}` WHERE metric_time <> {metric_literal}; + COMMIT; + """, + "database": DATABASE_NAME, + "stop_on_error": True, + "single_session": True, + }, + timeout=120, + ) + if not isinstance(script_payload, dict) or not isinstance(script_payload.get("results"), list): + raise ProcessingError("Metrix Database API returned an invalid script result") + results = script_payload["results"] + failed = next((item for item in results if not item.get("ok")), None) + if failed is not None: + raise ProcessingError(f"Database script failed at statement {failed.get('index')}: {failed.get('message', '')}") + if len(results) != 5: + raise ProcessingError(f"Database script returned {len(results)} results; expected 5") + + return { + "enabled": True, + "table": self.table, + "metric_time": metric_time.strftime("%Y-%m-%d %H:%M:%S"), + "inserted_rows": int(results[2].get("affected_rows") or 0), + "refreshed_rows": int(results[1].get("affected_rows") or 0), + "old_rows_deleted": int(results[3].get("affected_rows") or 0), + } + + def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]: + self._ensure_schema() + metric_literal = f"'{metric_time:%Y-%m-%d %H:%M:%S}'" + sql = ( + f"SELECT {', '.join(SUMMARY_HEADER)} " + f"FROM `{self.table}` " + f"WHERE metric_time = {metric_literal} ORDER BY cgi, cell_name" + ) + result: list[dict[str, str]] = [] + page = 1 + while True: + payload = self._query(sql, database=DATABASE_NAME, page=page, page_size=1000) + rows = payload.get("rows", []) + if not isinstance(rows, list): + raise ProcessingError("Metrix Database API returned invalid rows") + for row in rows: + normalized = {column: normalize_cell(row.get(column)) for column in SUMMARY_HEADER} + normalized["metric_time"] = datetime.fromisoformat(normalized["metric_time"]).strftime("%Y-%m-%d %H:%M:%S") + result.append(normalized) + total = int(payload.get("total") or len(result)) + if len(result) >= total or not rows: + return result + page += 1 + + def _ensure_schema(self) -> None: + self._query( + f"CREATE DATABASE IF NOT EXISTS `{DATABASE_NAME}` " + "CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci" + ) + self._query( + f""" + CREATE TABLE IF NOT EXISTS `{self.table}` ( + metric_time DATETIME NOT NULL COMMENT '指标开始时间', + network_type VARCHAR(16) NOT NULL DEFAULT '', + cgi VARCHAR(128) NOT NULL, + cell_name VARCHAR(255) NOT NULL, + interference_dbm DECIMAL(10,3) NOT NULL, + prev_interference_dbm DECIMAL(10,3) NULL, + longitude DECIMAL(10,6) NULL, + latitude DECIMAL(10,6) NULL, + azimuth DECIMAL(6,2) NOT NULL DEFAULT 0, + nearby_count INT NOT NULL DEFAULT 0, + nearby_26g INT NOT NULL DEFAULT 0, + nearby_700m INT NOT NULL DEFAULT 0, + nearby_tdd INT NOT NULL DEFAULT 0, + nearby_fdd INT NOT NULL DEFAULT 0, + prev_nearby_count INT NULL, + prev_nearby_26g INT NULL, + prev_nearby_700m INT NULL, + prev_nearby_tdd INT NULL, + prev_nearby_fdd INT NULL, + PRIMARY KEY (metric_time, cgi) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 + """, + database=DATABASE_NAME, + ) + columns = self.client.get_json( + self._endpoint("columns"), + {"database": DATABASE_NAME, "table": self.table}, + ) + if not isinstance(columns, list): + raise ProcessingError("Metrix Database API returned invalid column metadata") + existing_columns = {str(item.get("name")) for item in columns if isinstance(item, dict)} + migrations = { + "network_type": "VARCHAR(16) NOT NULL DEFAULT ''", + "prev_interference_dbm": "DECIMAL(10,3) NULL", + "longitude": "DECIMAL(10,6) NULL", + "latitude": "DECIMAL(10,6) NULL", + "azimuth": "DECIMAL(6,2) NOT NULL DEFAULT 0", + "nearby_count": "INT NOT NULL DEFAULT 0", + "nearby_26g": "INT NOT NULL DEFAULT 0", + "nearby_700m": "INT NOT NULL DEFAULT 0", + "nearby_tdd": "INT NOT NULL DEFAULT 0", + "nearby_fdd": "INT NOT NULL DEFAULT 0", + "prev_nearby_count": "INT NULL", + "prev_nearby_26g": "INT NULL", + "prev_nearby_700m": "INT NULL", + "prev_nearby_tdd": "INT NULL", + "prev_nearby_fdd": "INT NULL", + } + for column, definition in migrations.items(): + if column not in existing_columns: + self._query( + f"ALTER TABLE `{self.table}` ADD COLUMN `{column}` {definition}", + database=DATABASE_NAME, + ) + + def _endpoint(self, action: str) -> str: + return f"/api/databases/{quote(self.connection_id, safe='')}/{action}" + + def _query(self, sql: str, database: str = "", page: int = 1, page_size: int = 100) -> dict[str, object]: + payload = self.client.post_json( + self._endpoint("query"), + {"sql": sql, "database": database, "page": page, "page_size": page_size}, + timeout=120, + ) + if not isinstance(payload, dict): + raise ProcessingError("Metrix Database API returned an invalid query result") + return payload + + @staticmethod + def _row_values(metric_literal: str, row: dict[str, str]) -> str: + return "(" + ", ".join( + ( + metric_literal, + sql_text_literal(row["network_type"]), + sql_text_literal(row["cgi"]), + sql_text_literal(row["cell_name"]), + sql_decimal_literal(row["interference_dbm"], "interference_dbm"), + sql_decimal_literal(row["prev_interference_dbm"], "prev_interference_dbm", nullable=True), + sql_decimal_literal(row["longitude"], "longitude", nullable=True), + sql_decimal_literal(row["latitude"], "latitude", nullable=True), + sql_decimal_literal(row["azimuth"], "azimuth"), + str(int(row["nearby_count"])), + str(int(row["nearby_26g"])), + str(int(row["nearby_700m"])), + str(int(row["nearby_tdd"])), + str(int(row["nearby_fdd"])), + sql_decimal_literal(row["prev_nearby_count"], "prev_nearby_count", nullable=True), + sql_decimal_literal(row["prev_nearby_26g"], "prev_nearby_26g", nullable=True), + sql_decimal_literal(row["prev_nearby_700m"], "prev_nearby_700m", nullable=True), + sql_decimal_literal(row["prev_nearby_tdd"], "prev_nearby_tdd", nullable=True), + sql_decimal_literal(row["prev_nearby_fdd"], "prev_nearby_fdd", nullable=True), + ) + ) + ")" + + +class ApiHistoryStore: + def __init__(self, client: MetrixApiClient, storage_id: str, root: str = DEFAULT_HISTORY_ROOT) -> None: + self.client = client + self.storage_id = storage_id + self.root = root.rstrip("/") + + def exists(self, metric_time: datetime) -> bool: + directory, filename = self.target(metric_time) + base_entries = self._list_dir(self.root) + if not any(item.get("is_dir") and item.get("path") == directory for item in base_entries): + return False + return any( + not item.get("is_dir") and item.get("name") == filename and int(item.get("size") or 0) > 0 + for item in self._list_dir(directory) + ) + + def upload(self, metric_time: datetime, payload: bytes) -> str: + directory, filename = self.target(metric_time) + base_entries = self._list_dir(self.root) + if not any(item.get("is_dir") and item.get("path") == directory for item in base_entries): + response = self.client.post_json( + self._endpoint("mkdir"), + {"path": directory}, + timeout=120, + ) + if not isinstance(response, dict) or response.get("path") != directory: + raise ProcessingError("Metrix Storage API returned an invalid mkdir result") + response = self.client.post_file( + self._endpoint("upload"), + {"path": directory}, + filename, + payload, + ) + expected_path = posixpath.join(directory, filename) + if not isinstance(response, dict) or response.get("path") != expected_path: + raise ProcessingError("Metrix Storage API returned an invalid upload result") + return expected_path + + def target(self, metric_time: datetime) -> tuple[str, str]: + directory = posixpath.join(self.root, metric_time.strftime("%Y-%m-%d")) + filename = f"干扰数据处理结果_{metric_time:%Y%m%d%H%M%S}.csv" + return directory, filename + + def _endpoint(self, action: str) -> str: + return f"/api/storages/{quote(self.storage_id, safe='')}/{action}" + + def _list_dir(self, path: str) -> list[dict[str, object]]: + payload = self.client.get_json(self._endpoint("files"), {"path": path, "recursive": "false"}) + if not isinstance(payload, dict) or not isinstance(payload.get("entries"), list): + raise ProcessingError("Metrix Storage API returned an invalid file list") + return payload["entries"] + + +class ApiSource: + def __init__(self, client: MetrixApiClient, storage_id: str, root: str) -> None: + self.client = client + self.storage_id = storage_id + self.root = root.rstrip("/") or "/" + + def candidates(self, lookback_days: int) -> list[Candidate]: + root_entries = self._list_dir(self.root) + date_dirs = sorted( + (item for item in root_entries if item.get("is_dir") and re.fullmatch(r"\d{4}-\d{2}-\d{2}", item.get("name", ""))), + key=lambda item: item["name"], + ) + selected_dirs = date_dirs[-max(1, lookback_days) :] + result: list[Candidate] = [] + for directory in selected_dirs: + for item in self._list_dir(directory["path"]): + if item.get("is_dir"): + continue + candidate = parse_candidate(item["path"], int(item.get("size") or 0)) + if candidate is not None: + result.append(candidate) + return result + + def download(self, path: str) -> bytes: + endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/download" + return self.client.get_bytes(endpoint, {"path": path}, timeout=120) + + def cell_data_workbooks(self) -> list[tuple[str, bytes]]: + workbooks: list[tuple[str, bytes]] = [] + for directory in CELL_DATA_DIRECTORIES: + item = select_latest_cell_data_file(self._list_dir(directory), directory) + path = str(item["path"]) + workbooks.append((path, self.download(path))) + return workbooks + + def _list_dir(self, path: str) -> list[dict[str, object]]: + endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/files" + payload = self.client.get_json(endpoint, {"path": path, "recursive": "false"}) + if not isinstance(payload, dict) or not isinstance(payload.get("entries"), list): + raise ProcessingError("Metrix Storage API returned an invalid file list") + return payload["entries"] + + +class LocalSource: + def __init__(self, root: Path) -> None: + self.root = root.resolve() + if not self.root.is_dir(): + raise ProcessingError(f"Local source directory does not exist: {self.root}") + + def candidates(self, lookback_days: int) -> list[Candidate]: + del lookback_days + result: list[Candidate] = [] + for path in self.root.rglob("*.zip"): + candidate = parse_candidate(str(path.resolve()), path.stat().st_size) + if candidate is not None: + result.append(candidate) + return result + + def download(self, path: str) -> bytes: + return Path(path).read_bytes() + + def cell_data_workbooks(self) -> list[tuple[str, bytes]]: + directories = [self.root / Path(directory.lstrip("/")) for directory in CELL_DATA_DIRECTORIES] + existing = [directory for directory in directories if directory.is_dir()] + if not existing: + return [] + if len(existing) != len(directories): + raise ProcessingError("Local CellData source must contain all five configured directories") + workbooks: list[tuple[str, bytes]] = [] + for directory in directories: + entries = [ + {"name": path.name, "path": str(path.resolve()), "is_dir": path.is_dir()} + for path in directory.iterdir() + ] + item = select_latest_cell_data_file(entries, str(directory)) + path = str(item["path"]) + workbooks.append((path, self.download(path))) + return workbooks + + +def parse_candidate(path: str, size: int = 0) -> Candidate | None: + name = posixpath.basename(path.replace("\\", "/")) + source_type = next((item for item in EXPECTED_TYPES if name.startswith(f"{item}_")), "") + match = FILE_TIME_RE.search(name) + if not source_type or not match: + return None + window = match.group("window") + try: + window_start(window) + except ProcessingError: + return None + return Candidate(source_type, window, path, size) + + +def select_latest_cell_data_file(entries: list[dict[str, object]], directory: str) -> dict[str, object]: + dated_files: list[tuple[datetime, str, dict[str, object]]] = [] + for item in entries: + if item.get("is_dir"): + continue + name = str(item.get("name") or "") + match = CELL_DATA_FILE_RE.search(name) + if not match: + continue + try: + file_date = datetime.strptime(match.group("date"), "%Y%m%d") + except ValueError: + continue + dated_files.append((file_date, name, item)) + if not dated_files: + raise ProcessingError(f"No dated CellData XLSX found in {directory}") + return max(dated_files, key=lambda entry: (entry[0], entry[1]))[2] + + +def parse_cell_data_workbook(raw_xlsx: bytes, path: str) -> dict[str, tuple[str, str, str]]: + try: + with warnings.catch_warnings(): + warnings.filterwarnings("ignore", message="Workbook contains no default style") + workbook = load_workbook(io.BytesIO(raw_xlsx), read_only=True, data_only=True) + except Exception as exc: + raise ProcessingError(f"Invalid CellData workbook {path}: {exc}") from exc + try: + if "小区信息表" not in workbook.sheetnames: + raise ProcessingError(f"CellData workbook does not contain 小区信息表: {path}") + rows = workbook["小区信息表"].iter_rows(values_only=True) + try: + header = tuple(normalize_cell(value) for value in next(rows)) + except StopIteration as exc: + raise ProcessingError(f"CellData workbook is empty: {path}") from exc + required_columns = ("eNB/gNB", "CI", "经度", "纬度", "方向角") + missing = [column for column in required_columns if column not in header] + if missing: + raise ProcessingError(f"CellData workbook missing columns {missing}: {path}") + indexes = {column: header.index(column) for column in required_columns} + metadata: dict[str, tuple[str, str, str]] = {} + for row_number, row in enumerate(rows, start=2): + node = normalize_cell(row[indexes["eNB/gNB"]]) + cell = normalize_cell(row[indexes["CI"]]) + longitude = normalize_cell(row[indexes["经度"]]) + latitude = normalize_cell(row[indexes["纬度"]]) + azimuth = normalize_cell(row[indexes["方向角"]]) or "0" + if not node or not cell: + continue + try: + azimuth_number = Decimal(azimuth) + except InvalidOperation as exc: + raise ProcessingError(f"Invalid CellData azimuth {azimuth!r} in {path}, row {row_number}") from exc + if not azimuth_number.is_finite(): + raise ProcessingError(f"Invalid CellData azimuth {azimuth!r} in {path}, row {row_number}") + cgi = f"{LTE_PLMN}-{node}-{cell}" + value = (longitude, latitude, format(azimuth_number, "f")) + existing = metadata.get(cgi) + if existing is not None and existing != value: + raise ProcessingError(f"Conflicting CellData metadata for {cgi} in {path}, row {row_number}") + metadata[cgi] = value + return metadata + finally: + workbook.close() + + +def load_cell_metadata(workbooks: list[tuple[str, bytes]]) -> dict[str, tuple[str, str, str]]: + metadata: dict[str, tuple[str, str, str]] = {} + for path, raw_xlsx in workbooks: + for cgi, value in parse_cell_data_workbook(raw_xlsx, path).items(): + existing = metadata.get(cgi) + if existing is not None and existing != value: + raise ProcessingError(f"Conflicting CellData metadata for {cgi} across workbooks") + metadata[cgi] = value + return metadata + + +def select_window(candidates: list[Candidate], requested: str = "") -> tuple[datetime | None, dict[str, Candidate], list[str]]: + if not candidates: + return None, {}, list(EXPECTED_TYPES) + target_time = metric_time_for_window(requested) if requested else metric_time_for_window(max(candidates, key=candidate_key).window) + grouped = [candidate for candidate in candidates if metric_time_for_window(candidate.window) == target_time] + selected: dict[str, Candidate] = {} + for source_type in EXPECTED_TYPES: + matching = [candidate for candidate in grouped if candidate.source_type == source_type] + if matching: + selected[source_type] = max(matching, key=candidate_key) + missing = [source_type for source_type in EXPECTED_TYPES if source_type not in selected] + return target_time, selected, missing + + +def candidate_key(candidate: Candidate) -> tuple[datetime, str, str]: + return window_start(candidate.window), candidate.window, candidate.path + + +def window_start(window: str) -> datetime: + if not re.fullmatch(r"\d{16}", window): + raise ProcessingError(f"Invalid window: {window}") + try: + start = datetime.strptime(window[:12], "%Y%m%d%H%M") + datetime.strptime(window[12:], "%H%M") + except ValueError as exc: + raise ProcessingError(f"Invalid window: {window}") from exc + return start + + +def metric_time_for_window(window: str) -> datetime: + start = window_start(window) + return start.replace(minute=(start.minute // 15) * 15, second=0, microsecond=0) + + +def parse_workbook(raw_zip: bytes, source_type: str) -> tuple[tuple[str, ...], list[tuple[object, ...]], str]: + try: + with zipfile.ZipFile(io.BytesIO(raw_zip)) as archive: + bad_member = archive.testzip() + if bad_member: + raise ProcessingError(f"ZIP CRC check failed: {bad_member}") + xlsx_members = [item for item in archive.infolist() if not item.is_dir() and item.filename.lower().endswith(".xlsx")] + if len(xlsx_members) != 1: + raise ProcessingError(f"Expected one XLSX member, found {len(xlsx_members)}") + member = xlsx_members[0] + workbook_bytes = archive.read(member) + except zipfile.BadZipFile as exc: + raise ProcessingError("Invalid ZIP archive") from exc + + with warnings.catch_warnings(): + warnings.filterwarnings("ignore", message="Workbook contains no default style") + workbook = load_workbook(io.BytesIO(workbook_bytes), read_only=True, data_only=True) + try: + if "Sheet0" not in workbook.sheetnames: + raise ProcessingError("Workbook does not contain Sheet0") + rows = workbook["Sheet0"].iter_rows(values_only=True) + try: + header = tuple(normalize_cell(value) for value in next(rows)) + except StopIteration as exc: + raise ProcessingError("Sheet0 is empty") from exc + expected = EXPECTED_HEADERS[source_type] + if header != expected: + raise ProcessingError(f"Unexpected Sheet0 header for {source_type}: {header}") + data = [tuple(row) for row in rows if any(value is not None and normalize_cell(value) != "" for value in row)] + return header, data, member.filename + finally: + workbook.close() + + +def process( + source: Source, + output_root: Path, + lookback_days: int, + requested_window: str = "", + store: SummaryStore | None = None, + history: HistoryStore | None = None, +) -> Path | None: + candidates = source.candidates(lookback_days) + metric_time, selected, missing = select_window(candidates, requested_window) + if metric_time is None: + print("status=waiting") + print("reason=no_source_files") + return None + metric_time_text = metric_time.strftime("%Y-%m-%d %H:%M:%S") + if missing: + print("status=waiting") + print(f"target_time={metric_time_text}") + print(f"missing_sources={','.join(missing)}") + return None + + latest_time: datetime | None = None + previous_rows: list[dict[str, str]] = [] + if store is not None: + latest_time = store.latest_metric_time() + if latest_time is not None and latest_time >= metric_time: + if latest_time == metric_time and history is not None and not history.exists(metric_time): + stored_rows = store.rows_for_time(metric_time) + if not stored_rows: + raise ProcessingError(f"Database contains no rows for {metric_time_text}") + history_path = history.upload( + metric_time, + dict_csv_bytes(SUMMARY_HEADER, stored_rows, HISTORY_CSV_ENCODING), + ) + print("status=history_repaired") + print(f"target_time={metric_time_text}") + print(f"history_output={history_path}") + else: + print("status=skipped") + print(f"target_time={metric_time_text}") + print(f"database_metric_time={latest_time:%Y-%m-%d %H:%M:%S}") + return None + if latest_time is not None: + previous_rows = store.rows_for_time(latest_time) + + cell_data_workbooks = source.cell_data_workbooks() + cell_metadata = load_cell_metadata(cell_data_workbooks) + output_key = metric_time.strftime("%Y%m%d%H%M%S") + temp_dir = output_root.resolve() / f".{output_key}.tmp-{os.getpid()}" + final_dir = output_root.resolve() / output_key + ensure_scoped(output_root.resolve(), temp_dir) + if temp_dir.exists(): + shutil.rmtree(temp_dir) + converted_dir = temp_dir / "converted" + converted_dir.mkdir(parents=True) + + summary_rows: list[dict[str, str]] = [] + threshold_filtered_rows = 0 + manifest_files: list[dict[str, object]] = [] + database_result: dict[str, object] = {"enabled": False} + history_result: dict[str, object] = {"enabled": False} + try: + for source_type in EXPECTED_TYPES: + candidate = selected[source_type] + raw_zip = source.download(candidate.path) + header, rows, member_name = parse_workbook(raw_zip, source_type) + csv_name = f"{source_type}_{candidate.window}.csv" + write_csv(converted_dir / csv_name, header, rows) + for row_number, row in enumerate(rows, start=2): + record = dict(zip(header, row, strict=True)) + summary = summary_record(record, source_type, candidate.path, metric_time, row_number, cell_metadata) + if passes_interference_threshold(summary, source_type): + summary_rows.append(summary) + else: + threshold_filtered_rows += 1 + manifest_files.append( + { + "source_size": candidate.size, + "sha256": hashlib.sha256(raw_zip).hexdigest(), + "source_window": candidate.window, + "xlsx_member": member_name, + "rows": len(rows), + "converted_csv": f"converted/{csv_name}", + } + ) + + duplicate_cgi_rows = deduplicate_by_cgi(summary_rows) + populate_nearby_counts(summary_rows) + populate_previous_values(summary_rows, previous_rows) + summary_rows.sort(key=lambda item: (item["cgi"], item["cell_name"])) + summary_name = f"interference_summary_{output_key}.csv" + write_dict_csv(temp_dir / summary_name, SUMMARY_HEADER, summary_rows) + if store is not None: + database_result = store.replace_latest(summary_rows) + if history is not None: + history_path = history.upload( + metric_time, + dict_csv_bytes(SUMMARY_HEADER, summary_rows, HISTORY_CSV_ENCODING), + ) + history_result = {"enabled": True, "path": history_path} + matched_coordinates = sum(bool(row["longitude"] and row["latitude"]) for row in summary_rows) + manifest = { + "generated_at": datetime.now(timezone.utc).isoformat(), + "metric_time": metric_time_text, + "source_file_count": len(manifest_files), + "cell_data_file_count": len(cell_data_workbooks), + "summary_rows": len(summary_rows), + "threshold_filtered_rows": threshold_filtered_rows, + "duplicate_cgi_rows": duplicate_cgi_rows, + "coordinate_matched_rows": matched_coordinates, + "coordinate_unmatched_rows": len(summary_rows) - matched_coordinates, + "files": manifest_files, + "summary_csv": summary_name, + "database": database_result, + "history": history_result, + } + (temp_dir / "manifest.json").write_text(json.dumps(manifest, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") + output_root.resolve().mkdir(parents=True, exist_ok=True) + if final_dir.exists(): + ensure_scoped(output_root.resolve(), final_dir) + shutil.rmtree(final_dir) + temp_dir.replace(final_dir) + except Exception: + shutil.rmtree(temp_dir, ignore_errors=True) + raise + + print("status=processed") + print(f"target_time={metric_time_text}") + print(f"source_files={len(manifest_files)}") + print(f"cell_data_files={len(cell_data_workbooks)}") + print(f"summary_rows={len(summary_rows)}") + print(f"threshold_filtered_rows={threshold_filtered_rows}") + print(f"duplicate_cgi_rows={duplicate_cgi_rows}") + print(f"coordinate_matched_rows={matched_coordinates}") + print(f"coordinate_unmatched_rows={len(summary_rows) - matched_coordinates}") + if database_result["enabled"]: + print(f"database_metric_time={database_result['metric_time']}") + print(f"database_inserted_rows={database_result['inserted_rows']}") + print(f"database_old_rows_deleted={database_result['old_rows_deleted']}") + if history_result["enabled"]: + print(f"history_output={history_result['path']}") + print(f"output={final_dir}") + return final_dir + + +def summary_record( + record: dict[str, object], + source_type: str, + source_path: str, + metric_time: datetime, + row_number: int, + cell_metadata: dict[str, tuple[str, str, str]] | None = None, +) -> dict[str, str]: + row_time_text = normalize_cell(record["开始时间"]) + try: + row_time = datetime.strptime(row_time_text, "%Y-%m-%d %H:%M:%S") + except ValueError as exc: + raise ProcessingError(f"{source_path}: row {row_number} has invalid time {row_time_text!r}") from exc + row_metric_time = row_time.replace(minute=(row_time.minute // 15) * 15, second=0, microsecond=0) + if row_metric_time != metric_time: + raise ProcessingError( + f"{source_path}: row {row_number} time {row_time_text!r} is outside metric group " + f"{metric_time:%Y-%m-%d %H:%M:%S}" + ) + + if source_type in ("5G干扰监控", "700M干扰监控"): + cgi = build_cgi("", record, ("gNBplmn", "gNBId", "cellId"), source_path, row_number) + cell_name = required(record, "CU小区配置名称", source_path, row_number) + interference = required(record, "小区上行平均干扰电平(dBm)", source_path, row_number) + elif source_type in ("SDR_FDD干扰监控", "SDR_TDD干扰监控"): + cgi = build_cgi(LTE_PLMN, record, ("eNodeBID", "小区ID"), source_path, row_number) + cell_name = required(record, "小区名称", source_path, row_number) + interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number) + elif source_type == "反开RD干扰监控": + cgi = build_cgi(LTE_PLMN, record, ("eNodeBId", "cellId"), source_path, row_number) + cell_name = required(record, "E-UTRAN TDD小区名称", source_path, row_number) + interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number) + else: + cgi = build_cgi(LTE_PLMN, record, ("eNodeBId", "cellId"), source_path, row_number) + cell_name = required(record, "E-UTRAN FDD小区名称", source_path, row_number) + interference = required(record, "载波平均噪声干扰(dBm)", source_path, row_number) + + longitude, latitude, azimuth = (cell_metadata or {}).get(cgi, ("", "", "0")) + record_out = { + "metric_time": metric_time.strftime("%Y-%m-%d %H:%M:%S"), + "network_type": NETWORK_TYPE_BY_SOURCE[source_type], + "cgi": cgi, + "cell_name": cell_name, + "interference_dbm": interference, + "prev_interference_dbm": "", + "longitude": longitude, + "latitude": latitude, + "azimuth": azimuth, + "nearby_count": "0", + "prev_nearby_count": "", + } + for column in NEARBY_COLUMN_BY_NETWORK_TYPE.values(): + record_out[column] = "0" + record_out[f"prev_{column}"] = "" + return record_out + + +def passes_interference_threshold(record: dict[str, str], source_type: str) -> bool: + try: + interference = Decimal(record["interference_dbm"]) + except InvalidOperation as exc: + raise ProcessingError(f"Invalid interference value: {record['interference_dbm']}") from exc + if not interference.is_finite(): + raise ProcessingError(f"Invalid interference value: {record['interference_dbm']}") + return interference >= MIN_INTERFERENCE_BY_SOURCE[source_type] + + +def populate_nearby_counts(rows: list[dict[str, str]]) -> None: + type_counts = [dict.fromkeys(NEARBY_COLUMN_BY_NETWORK_TYPE.values(), 0) for _ in rows] + spatial_index: dict[tuple[int, int, int], list[tuple[int, float, float]]] = {} + for index, row in enumerate(rows): + if not row["longitude"] or not row["latitude"]: + continue + try: + longitude = float(row["longitude"]) + latitude = float(row["latitude"]) + except ValueError as exc: + raise ProcessingError(f"Invalid coordinates for CGI {row['cgi']}") from exc + if not math.isfinite(longitude) or not math.isfinite(latitude) or not -180 <= longitude <= 180 or not -90 <= latitude <= 90: + raise ProcessingError(f"Invalid coordinates for CGI {row['cgi']}") + + column = NEARBY_COLUMN_BY_NETWORK_TYPE[row["network_type"]] + bucket = spatial_bucket(longitude, latitude) + for x_offset in (-1, 0, 1): + for y_offset in (-1, 0, 1): + for z_offset in (-1, 0, 1): + nearby_bucket = (bucket[0] + x_offset, bucket[1] + y_offset, bucket[2] + z_offset) + for other_index, other_longitude, other_latitude in spatial_index.get(nearby_bucket, ()): + if haversine_distance_km(longitude, latitude, other_longitude, other_latitude) <= NEARBY_RADIUS_KM: + type_counts[index][NEARBY_COLUMN_BY_NETWORK_TYPE[rows[other_index]["network_type"]]] += 1 + type_counts[other_index][column] += 1 + spatial_index.setdefault(bucket, []).append((index, longitude, latitude)) + + for row, counts in zip(rows, type_counts, strict=True): + for column, count in counts.items(): + row[column] = str(count) + row["nearby_count"] = str(sum(counts.values())) + + +def deduplicate_by_cgi(rows: list[dict[str, str]]) -> int: + selected: dict[str, dict[str, str]] = {} + removed = 0 + for row in rows: + previous = selected.get(row["cgi"]) + if previous is None: + selected[row["cgi"]] = row + continue + if Decimal(row["interference_dbm"]) > Decimal(previous["interference_dbm"]): + selected[row["cgi"]] = row + removed += 1 + rows[:] = selected.values() + return removed + + +def populate_previous_values(rows: list[dict[str, str]], previous_rows: list[dict[str, str]]) -> None: + previous_by_cgi = {row["cgi"]: row for row in previous_rows} + copied_columns = ("interference_dbm", "nearby_count", *NEARBY_COLUMN_BY_NETWORK_TYPE.values()) + for row in rows: + previous = previous_by_cgi.get(row["cgi"]) + for column in copied_columns: + row[f"prev_{column}"] = previous[column] if previous is not None else "" + + +def spatial_bucket(longitude: float, latitude: float) -> tuple[int, int, int]: + longitude = math.radians(longitude) + latitude = math.radians(latitude) + radius_at_latitude = EARTH_RADIUS_KM * math.cos(latitude) + return ( + math.floor(radius_at_latitude * math.cos(longitude) / NEARBY_RADIUS_KM), + math.floor(radius_at_latitude * math.sin(longitude) / NEARBY_RADIUS_KM), + math.floor(EARTH_RADIUS_KM * math.sin(latitude) / NEARBY_RADIUS_KM), + ) + + +def haversine_distance_km(longitude1: float, latitude1: float, longitude2: float, latitude2: float) -> float: + longitude1, latitude1, longitude2, latitude2 = map( + math.radians, + (longitude1, latitude1, longitude2, latitude2), + ) + longitude_delta = longitude2 - longitude1 + latitude_delta = latitude2 - latitude1 + haversine = math.sin(latitude_delta / 2) ** 2 + ( + math.cos(latitude1) * math.cos(latitude2) * math.sin(longitude_delta / 2) ** 2 + ) + return 2 * EARTH_RADIUS_KM * math.asin(math.sqrt(min(1.0, haversine))) + + +def build_cgi(prefix: str, record: dict[str, object], columns: tuple[str, ...], path: str, row: int) -> str: + parts = [required(record, column, path, row) for column in columns] + return "-".join(([prefix] if prefix else []) + parts) + + +def required(record: dict[str, object], column: str, path: str, row: int) -> str: + value = normalize_cell(record[column]) + if not value: + raise ProcessingError(f"{path}: row {row} has empty {column}") + return value + + +def normalize_cell(value: object) -> str: + if value is None: + return "" + if isinstance(value, datetime): + return value.strftime("%Y-%m-%d %H:%M:%S") + if isinstance(value, float) and value.is_integer(): + return str(int(value)) + return str(value).strip() + + +def normalize_csv_cell(value: object) -> str: + return normalize_cell(value).replace("\u00a0", " ") + + +def sql_text_literal(value: str) -> str: + encoded = value.encode("utf-8").hex() + return "''" if not encoded else f"CONVERT(0x{encoded} USING utf8mb4)" + + +def sql_decimal_literal(value: str, field: str, nullable: bool = False) -> str: + normalized = value.strip() + if not normalized and nullable: + return "NULL" + try: + number = Decimal(normalized) + except InvalidOperation as exc: + raise ProcessingError(f"Invalid decimal value for {field}: {value}") from exc + if not number.is_finite(): + raise ProcessingError(f"Invalid decimal value for {field}: {value}") + return format(number, "f") + + +def write_csv(path: Path, header: tuple[str, ...], rows: list[tuple[object, ...]]) -> None: + with path.open("w", encoding=LOCAL_CSV_ENCODING, newline="") as file: + writer = csv.writer(file) + writer.writerow(header) + for row in rows: + writer.writerow(normalize_cell(value) for value in row) + + +def write_dict_csv(path: Path, header: tuple[str, ...], rows: list[dict[str, str]]) -> None: + path.write_bytes(dict_csv_bytes(header, rows)) + + +def dict_csv_bytes( + header: tuple[str, ...], + rows: list[dict[str, str]], + encoding: str = LOCAL_CSV_ENCODING, +) -> bytes: + output = io.StringIO(newline="") + writer = csv.DictWriter(output, fieldnames=header, extrasaction="raise") + writer.writeheader() + for row in rows: + if encoding == HISTORY_CSV_ENCODING: + writer.writerow({column: normalize_csv_cell(row[column]) for column in header}) + else: + writer.writerow(row) + return output.getvalue().encode(encoding) + + +def ensure_scoped(root: Path, target: Path) -> None: + if target == root or root not in target.parents: + raise ProcessingError(f"Refusing to modify path outside output root: {target}") + + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser(description="Process the latest complete interference KPI time group") + parser.add_argument("--source-dir", type=Path, help="Use a local source tree instead of the Metrix storage API") + parser.add_argument("--output-dir", type=Path, default=Path(os.getenv("INTERFERENCE_OUTPUT_DIR", "output"))) + parser.add_argument( + "--window", + default=os.getenv("INTERFERENCE_WINDOW", ""), + help="Optional 16-digit source window used to select a natural 15-minute group", + ) + parser.add_argument("--lookback-days", type=int, default=int(os.getenv("INTERFERENCE_LOOKBACK_DAYS", "3"))) + parser.add_argument("--storage-id", default=os.getenv("METRIX_STORAGE_ID", DEFAULT_STORAGE_ID)) + parser.add_argument("--root", default=os.getenv("INTERFERENCE_SOURCE_ROOT", DEFAULT_ROOT)) + parser.add_argument("--no-database", action="store_true", help="Generate files without writing the database") + return parser + + +def main(argv: list[str] | None = None) -> int: + args = build_parser().parse_args(argv) + if args.lookback_days < 1: + raise ProcessingError("lookback-days must be at least 1") + client: MetrixApiClient | None = None + if args.source_dir: + source: Source = LocalSource(args.source_dir) + else: + client = MetrixApiClient(METRIX_API_BASE_URL, METRIX_API_TOKEN) + source = ApiSource(client, args.storage_id, args.root) + if args.no_database: + store = None + history = None + else: + client = client or MetrixApiClient(METRIX_API_BASE_URL, METRIX_API_TOKEN) + store = ApiSummaryStore(client, METRIX_DATABASE_CONNECTION_ID) + history = ApiHistoryStore(client, args.storage_id, f"{args.root.rstrip('/')}/干扰历史数据") + process(source, args.output_dir, args.lookback_days, args.window, store=store, history=history) + return 0 + + +if __name__ == "__main__": + try: + raise SystemExit(main()) + except ProcessingError as exc: + print(f"error={exc}", file=sys.stderr) + raise SystemExit(1) from exc diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..61d7845 --- /dev/null +++ b/requirements.txt @@ -0,0 +1 @@ +openpyxl==3.1.5 diff --git a/runtime_config.example.py b/runtime_config.example.py new file mode 100644 index 0000000..f4c48d2 --- /dev/null +++ b/runtime_config.example.py @@ -0,0 +1,3 @@ +METRIX_API_BASE_URL = "http://188.5.127.115:18271" +METRIX_API_TOKEN = "replace-with-metrix-api-token" +METRIX_DATABASE_CONNECTION_ID = "replace-with-database-connection-id" diff --git a/scripts/build_offline_package.py b/scripts/build_offline_package.py new file mode 100644 index 0000000..417655d --- /dev/null +++ b/scripts/build_offline_package.py @@ -0,0 +1,122 @@ +from __future__ import annotations + +import argparse +import hashlib +from pathlib import Path +import shutil +import subprocess +import tempfile +import zipfile + + +PROJECT_ROOT = Path(__file__).resolve().parents[1] +DEFAULT_IMAGE = "interference-etl-runtime:1.1" +DEFAULT_BASE_IMAGE = "python:3.13.11-slim" +DEFAULT_OUTPUT = PROJECT_ROOT / "dist" / "InterferenceETL-offline-1.1.zip" +PACKAGE_FILES = ("main.py", "requirements.txt", "README.md", "runtime_config.py") +PACKAGE_ENTRYPOINT = "from main import main\nraise SystemExit(main())\n" + +COPY_DEPENDENCIES = r""" +from importlib.metadata import distribution +from pathlib import Path +import shutil + +target = Path('/package/vendor') +for name in ('openpyxl', 'et-xmlfile'): + package = distribution(name) + for item in package.files or (): + if '__pycache__' in item.parts or item.suffix in ('.pyc', '.pyo'): + continue + source = Path(package.locate_file(item)) + if not source.is_file(): + continue + destination = target / item + destination.parent.mkdir(parents=True, exist_ok=True) + shutil.copy2(source, destination) +""" + + +def build(output: Path, runtime_image: str, base_image: str) -> Path: + output = output.resolve() + output.parent.mkdir(parents=True, exist_ok=True) + missing = [name for name in PACKAGE_FILES if not (PROJECT_ROOT / name).is_file()] + if missing: + raise FileNotFoundError(f"Missing package files: {', '.join(missing)}") + with tempfile.TemporaryDirectory(prefix="interference-etl-package-") as temp: + staging = Path(temp) + subprocess.run( + [ + "docker", + "run", + "--rm", + "-v", + f"{staging}:/package", + runtime_image, + "python", + "-c", + COPY_DEPENDENCIES, + ], + check=True, + ) + for name in PACKAGE_FILES: + shutil.copy2(PROJECT_ROOT / name, staging / name) + (staging / "__main__.py").write_text(PACKAGE_ENTRYPOINT, encoding="utf-8") + subprocess.run( + [ + "docker", + "run", + "--rm", + "-e", + "PYTHONDONTWRITEBYTECODE=1", + "-v", + f"{staging}:/workspace:ro", + "-w", + "/workspace", + base_image, + "python", + "-c", + "import main, openpyxl; print('offline_package_imports=ok')", + ], + check=True, + ) + if output.exists(): + output.unlink() + with zipfile.ZipFile(output, "w", compression=zipfile.ZIP_DEFLATED, compresslevel=9) as archive: + for path in sorted(staging.rglob("*")): + if path.is_file(): + archive.write(path, path.relative_to(staging).as_posix()) + subprocess.run( + [ + "docker", + "run", + "--rm", + "-e", + "PYTHONDONTWRITEBYTECODE=1", + "-v", + f"{output.parent}:/dist:ro", + base_image, + "python", + f"/dist/{output.name}", + "--help", + ], + check=True, + ) + return output + + +def main() -> int: + parser = argparse.ArgumentParser(description="Build a self-contained Metrix Script Management upload ZIP") + parser.add_argument("--runtime-image", default=DEFAULT_IMAGE) + parser.add_argument("--base-image", default=DEFAULT_BASE_IMAGE) + parser.add_argument("--output", type=Path, default=DEFAULT_OUTPUT) + args = parser.parse_args() + output = build(args.output, args.runtime_image, args.base_image) + digest = hashlib.sha256(output.read_bytes()).hexdigest().upper() + print(f"output={output}") + print(f"size={output.stat().st_size}") + print(f"sha256={digest}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/__init__.py b/tests/__init__.py new file mode 100644 index 0000000..8b13789 --- /dev/null +++ b/tests/__init__.py @@ -0,0 +1 @@ + diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py new file mode 100644 index 0000000..50e8935 --- /dev/null +++ b/tests/test_pipeline.py @@ -0,0 +1,659 @@ +from __future__ import annotations + +import csv +from datetime import datetime +import io +import json +from pathlib import Path +import tempfile +import unittest +import zipfile + +import main +from openpyxl import Workbook + + +COMPLETE_WINDOW = "2026073110001100" +INCOMPLETE_WINDOW = "2026073111001200" +TARGET_KEY = "20260731100000" + + +class PipelineTest(unittest.TestCase): + def test_latest_group_is_converted_merged_stored_and_archived(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + output = Path(temp) / "output" + windows = ["2026073110001100", "2026073110141100", "2026073110051100"] + for index, source_type in enumerate(main.EXPECTED_TYPES): + create_archive(root, source_type, windows[index % len(windows)], middle="任意粒度") + create_archive(root, main.EXPECTED_TYPES[0], "2026073110091100", middle="更新版本") + create_cell_data_sources(root) + store = RecordingStore() + history = RecordingHistory() + + result = main.process(main.LocalSource(root), output, lookback_days=3, store=store, history=history) + + self.assertIsNotNone(result) + assert result is not None + self.assertEqual(result.name, TARGET_KEY) + converted = sorted((result / "converted").glob("*.csv")) + self.assertEqual(len(converted), 7) + summary_path = result / f"interference_summary_{TARGET_KEY}.csv" + summary_bytes = summary_path.read_bytes() + self.assertTrue(summary_bytes.startswith(b"\xef\xbb\xbf")) + with summary_path.open(encoding="utf-8-sig", newline="") as file: + rows = list(csv.DictReader(file)) + self.assertEqual(len(rows), 2) + self.assertEqual( + list(rows[0]), + [ + "metric_time", + "network_type", + "cgi", + "cell_name", + "interference_dbm", + "prev_interference_dbm", + "longitude", + "latitude", + "azimuth", + "nearby_count", + "nearby_26g", + "nearby_700m", + "nearby_tdd", + "nearby_fdd", + "prev_nearby_count", + "prev_nearby_26g", + "prev_nearby_700m", + "prev_nearby_tdd", + "prev_nearby_fdd", + ], + ) + rows_by_cgi = {row["cgi"]: row for row in rows} + self.assertEqual(rows_by_cgi["460-00-200-1"]["network_type"], "2.6G") + self.assertEqual(rows_by_cgi["460-00-100-1"]["network_type"], "FDD") + self.assertTrue(all(row["metric_time"] == "2026-07-31 10:00:00" for row in rows)) + self.assertEqual(rows_by_cgi["460-00-200-1"]["cell_name"], "5G干扰监控-小区") + self.assertTrue(all(row["interference_dbm"] == "-100.5" for row in rows)) + self.assertTrue(all(row["prev_interference_dbm"] == "" for row in rows)) + self.assertTrue(all(row["prev_nearby_count"] == "" for row in rows)) + self.assertEqual(rows_by_cgi["460-00-200-1"]["longitude"], "113.123456") + self.assertEqual(rows_by_cgi["460-00-200-1"]["latitude"], "22.654321") + self.assertEqual(rows_by_cgi["460-00-200-1"]["azimuth"], "30") + self.assertEqual(rows_by_cgi["460-00-200-1"]["nearby_count"], "0") + self.assertEqual(rows_by_cgi["460-00-100-1"]["nearby_count"], "0") + self.assertFalse(rows_by_cgi["460-00-100-1"]["longitude"]) + self.assertEqual(rows_by_cgi["460-00-100-1"]["azimuth"], "0") + + manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8")) + self.assertEqual(manifest["source_file_count"], 7) + self.assertEqual(manifest["cell_data_file_count"], 5) + self.assertEqual(manifest["summary_rows"], 2) + self.assertEqual(manifest["threshold_filtered_rows"], 0) + self.assertEqual(manifest["duplicate_cgi_rows"], 5) + self.assertEqual(manifest["coordinate_matched_rows"], 1) + self.assertEqual(manifest["coordinate_unmatched_rows"], 1) + self.assertEqual(manifest["metric_time"], "2026-07-31 10:00:00") + self.assertNotIn("hour_start", manifest) + self.assertNotIn("hour_end", manifest) + self.assertNotIn("source_types", manifest) + self.assertTrue(all("source_type" not in item and "source_path" not in item for item in manifest["files"])) + self.assertIn("2026073110091100", [item["source_window"] for item in manifest["files"]]) + self.assertEqual(manifest["database"]["metric_time"], "2026-07-31 10:00:00") + self.assertEqual(len(store.rows), 2) + self.assertEqual(history.path, "/history/2026-07-31/干扰数据处理结果_20260731100000.csv") + archived_rows = list(csv.DictReader(io.StringIO(history.payload.decode("gbk")))) + self.assertEqual(list(archived_rows[0]), list(main.SUMMARY_HEADER)) + self.assertEqual(len(archived_rows), 2) + + def test_interference_thresholds_remove_only_lower_values(self) -> None: + for source_type, threshold in ( + ("5G干扰监控", "-107"), + ("5G下FDD干扰监控", "-107"), + ("700M干扰监控", "-110"), + ("700M下FDD干扰监控", "-110"), + ("SDR_FDD干扰监控", "-110"), + ("SDR_TDD干扰监控", "-110"), + ("反开RD干扰监控", "-110"), + ): + row = database_row() + row["interference_dbm"] = threshold + self.assertTrue(main.passes_interference_threshold(row, source_type)) + row["interference_dbm"] = str(float(threshold) - 0.1) + self.assertFalse(main.passes_interference_threshold(row, source_type)) + + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + interference = -107.1 if source_type == main.EXPECTED_TYPES[0] else -100.5 + create_archive(root, source_type, COMPLETE_WINDOW, interference_dbm=interference) + create_cell_data_sources(root) + + result = main.process(main.LocalSource(root), Path(temp) / "output", 3) + manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8")) + + self.assertEqual(manifest["summary_rows"], 2) + self.assertEqual(manifest["threshold_filtered_rows"], 1) + self.assertEqual(manifest["duplicate_cgi_rows"], 4) + + def test_latest_incomplete_group_waits_without_falling_back(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + create_archive(root, source_type, COMPLETE_WINDOW) + create_archive(root, main.EXPECTED_TYPES[0], INCOMPLETE_WINDOW) + store = RecordingStore() + + result = main.process(main.LocalSource(root), Path(temp) / "output", 3, store=store) + + self.assertIsNone(result) + self.assertEqual(store.latest_calls, 0) + self.assertEqual(store.rows, []) + + def test_no_source_files_waits_successfully(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + root.mkdir() + store = RecordingStore() + + result = main.process(main.LocalSource(root), Path(temp) / "output", 3, store=store) + + self.assertIsNone(result) + self.assertEqual(store.latest_calls, 0) + + def test_schema_change_is_rejected(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + create_archive(root, source_type, COMPLETE_WINDOW, bad_header=source_type == main.EXPECTED_TYPES[0]) + + with self.assertRaisesRegex(main.ProcessingError, "Unexpected Sheet0 header"): + main.process(main.LocalSource(root), Path(temp) / "output", 3) + + def test_workbook_row_outside_selected_quarter_is_rejected(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + create_archive( + root, + source_type, + COMPLETE_WINDOW, + row_window="2026073110151100" if source_type == main.EXPECTED_TYPES[0] else "", + ) + create_cell_data_sources(root) + + with self.assertRaisesRegex(main.ProcessingError, "is outside metric group"): + main.process(main.LocalSource(root), Path(temp) / "output", 3) + + def test_file_name_middle_is_flexible_and_time_uses_natural_quarter(self) -> None: + candidate = main.parse_candidate( + "/source/700M干扰监控_任意描述_2026073112141300.zip", + 123, + ) + + self.assertIsNotNone(candidate) + assert candidate is not None + self.assertEqual(candidate.source_type, "700M干扰监控") + self.assertEqual(candidate.size, 123) + self.assertEqual(main.metric_time_for_window(candidate.window), datetime(2026, 7, 31, 12, 0)) + self.assertEqual(main.metric_time_for_window("2026073112151300"), datetime(2026, 7, 31, 12, 15)) + + def test_latest_dated_cell_data_file_is_selected(self) -> None: + entries = [ + {"name": "江门5G小区信息表20260727.xlsx", "path": "/old.xlsx", "is_dir": False}, + {"name": "江门5G小区信息表20260728.xlsx", "path": "/latest.xlsx", "is_dir": False}, + {"name": "说明.txt", "path": "/说明.txt", "is_dir": False}, + ] + + selected = main.select_latest_cell_data_file(entries, "/cell-data") + + self.assertEqual(selected["path"], "/latest.xlsx") + + def test_cell_data_schema_change_is_rejected(self) -> None: + raw = create_cell_data_workbook(include_latitude=False) + + with self.assertRaisesRegex(main.ProcessingError, "missing columns.*纬度"): + main.parse_cell_data_workbook(raw, "bad.xlsx") + + def test_conflicting_cell_data_metadata_are_rejected(self) -> None: + first = create_cell_data_workbook(longitude=113.1) + second = create_cell_data_workbook(longitude=113.2) + + with self.assertRaisesRegex(main.ProcessingError, "Conflicting CellData metadata"): + main.load_cell_metadata([("first.xlsx", first), ("second.xlsx", second)]) + + def test_empty_cell_data_azimuth_defaults_to_zero(self) -> None: + metadata = main.parse_cell_data_workbook(create_cell_data_workbook(azimuth=None), "cell-data.xlsx") + + self.assertEqual(metadata["460-00-200-1"], ("113.123456", "22.654321", "0")) + + def test_nearby_count_uses_high_interference_rows_with_coordinates(self) -> None: + rows = [ + nearby_row("a", "113.000000", "22.000000", "2.6G"), + nearby_row("b", "113.005000", "22.000000", "700M"), + nearby_row("c", "113.020000", "22.000000", "TDD"), + nearby_row("d", "113.000000", "22.000000", "FDD"), + nearby_row("e", "", "", "TDD"), + ] + + main.populate_nearby_counts(rows) + + self.assertEqual([row["nearby_count"] for row in rows], ["2", "2", "0", "2", "0"]) + by_cgi = {row["cgi"]: row for row in rows} + self.assertEqual( + [by_cgi["a"][column] for column in ("nearby_26g", "nearby_700m", "nearby_tdd", "nearby_fdd")], + ["0", "1", "0", "1"], + ) + self.assertEqual( + [by_cgi["b"][column] for column in ("nearby_26g", "nearby_700m", "nearby_tdd", "nearby_fdd")], + ["1", "0", "0", "1"], + ) + self.assertEqual( + [by_cgi["d"][column] for column in ("nearby_26g", "nearby_700m", "nearby_tdd", "nearby_fdd")], + ["1", "1", "0", "0"], + ) + for row in rows: + per_type_total = sum( + int(row[column]) for column in ("nearby_26g", "nearby_700m", "nearby_tdd", "nearby_fdd") + ) + self.assertEqual(per_type_total, int(row["nearby_count"])) + self.assertLess(main.haversine_distance_km(113.0, 22.0, 113.005, 22.0), 1.0) + self.assertGreater(main.haversine_distance_km(113.0, 22.0, 113.02, 22.0), 1.0) + + def test_previous_values_are_matched_by_cgi(self) -> None: + rows = [database_row(), database_row()] + rows[0]["cgi"] = "matched" + rows[1]["cgi"] = "new" + previous = database_row() + previous["cgi"] = "matched" + previous["interference_dbm"] = "-106.25" + previous["nearby_count"] = "8" + previous["nearby_26g"] = "3" + previous["nearby_700m"] = "2" + previous["nearby_tdd"] = "2" + previous["nearby_fdd"] = "1" + + main.populate_previous_values(rows, [previous]) + + self.assertEqual(rows[0]["prev_interference_dbm"], "-106.25") + self.assertEqual(rows[0]["prev_nearby_count"], "8") + self.assertEqual(rows[0]["prev_nearby_26g"], "3") + self.assertEqual(rows[0]["prev_nearby_700m"], "2") + self.assertEqual(rows[0]["prev_nearby_tdd"], "2") + self.assertEqual(rows[0]["prev_nearby_fdd"], "1") + self.assertEqual(rows[1]["prev_interference_dbm"], "") + self.assertEqual(rows[1]["prev_nearby_count"], "") + self.assertEqual(rows[1]["prev_nearby_26g"], "") + self.assertEqual(rows[1]["prev_nearby_700m"], "") + self.assertEqual(rows[1]["prev_nearby_tdd"], "") + self.assertEqual(rows[1]["prev_nearby_fdd"], "") + + def test_duplicate_cgi_keeps_maximum_interference(self) -> None: + rows = [database_row(), database_row(), database_row()] + rows[0]["cgi"] = "duplicate" + rows[0]["interference_dbm"] = "-108.5" + rows[1]["cgi"] = "duplicate" + rows[1]["interference_dbm"] = "-101.25" + rows[2]["cgi"] = "unique" + + removed = main.deduplicate_by_cgi(rows) + + self.assertEqual(removed, 1) + self.assertEqual(len(rows), 2) + self.assertEqual({row["cgi"]: row["interference_dbm"] for row in rows}, { + "duplicate": "-101.25", + "unique": "-100.5", + }) + + def test_gbk_history_csv_normalizes_non_breaking_spaces(self) -> None: + payload = main.dict_csv_bytes( + ("cell_name",), + [{"cell_name": "测试\u00a0小区"}], + main.HISTORY_CSV_ENCODING, + ) + + self.assertEqual(payload.decode("gbk"), "cell_name\r\n测试 小区\r\n") + + def test_api_store_replaces_same_hour_and_deletes_other_hours(self) -> None: + client = FakeApiClient("2026-07-31T10:00:00") + store = main.ApiSummaryStore(client, "db_share_mysql") + + result = store.replace_latest([database_row()]) + + self.assertEqual(result["metric_time"], "2026-07-31 10:00:00") + self.assertEqual(result["inserted_rows"], 1) + self.assertEqual(result["refreshed_rows"], 7) + self.assertEqual(result["old_rows_deleted"], 20) + queries = [payload["sql"] for endpoint, payload in client.posts if endpoint.endswith("/query")] + self.assertTrue(any("CREATE DATABASE IF NOT EXISTS `interference_etl`" in query for query in queries)) + self.assertTrue(any("ADD COLUMN `network_type`" in query for query in queries)) + self.assertTrue(any("ADD COLUMN `longitude`" in query for query in queries)) + self.assertTrue(any("ADD COLUMN `latitude`" in query for query in queries)) + self.assertTrue(any("ADD COLUMN `azimuth`" in query for query in queries)) + self.assertTrue(any("ADD COLUMN `nearby_count`" in query for query in queries)) + for column in ("nearby_26g", "nearby_700m", "nearby_tdd", "nearby_fdd"): + self.assertTrue(any(f"ADD COLUMN `{column}`" in query for query in queries)) + self.assertTrue(any(f"ADD COLUMN `prev_{column}`" in query for query in queries)) + self.assertTrue(any("ADD COLUMN `prev_interference_dbm`" in query for query in queries)) + self.assertTrue(any("ADD COLUMN `prev_nearby_count`" in query for query in queries)) + script = next(payload for endpoint, payload in client.posts if endpoint.endswith("/run-script")) + self.assertTrue(script["single_session"]) + self.assertIn("START TRANSACTION", script["content"]) + self.assertIn("NULL, NULL, NULL, 0, 0, 0, 0, 0, 0, NULL, NULL, NULL, NULL, NULL", script["content"]) + self.assertIn("CONVERT(0x", script["content"]) + self.assertNotIn("source_type", script["content"]) + self.assertNotIn("source_path", script["content"]) + + def test_api_store_refuses_to_replace_a_newer_hour(self) -> None: + client = FakeApiClient("2026-07-31T11:00:00") + store = main.ApiSummaryStore(client, "db_share_mysql") + + with self.assertRaisesRegex(main.ProcessingError, "already contains newer metric time"): + store.replace_latest([database_row()]) + + self.assertFalse(any(endpoint.endswith("/run-script") for endpoint, _ in client.posts)) + + def test_api_store_propagates_script_failure(self) -> None: + client = FakeApiClient(None, fail_script=True) + store = main.ApiSummaryStore(client, "db_share_mysql") + + with self.assertRaisesRegex(main.ProcessingError, "statement 3: insert failed"): + store.replace_latest([database_row()]) + + def test_api_store_rejects_invalid_decimal(self) -> None: + client = FakeApiClient(None) + store = main.ApiSummaryStore(client, "db_share_mysql") + row = database_row() + row["interference_dbm"] = "not-a-number" + + with self.assertRaisesRegex(main.ProcessingError, "Invalid decimal value for interference_dbm"): + store.replace_latest([row]) + + def test_api_store_exports_normalized_rows_for_history(self) -> None: + client = FakeApiClient(None) + store = main.ApiSummaryStore(client, "db_share_mysql") + + rows = store.rows_for_time(datetime(2026, 7, 31, 10, 0)) + + self.assertEqual(rows, [database_row()]) + export_query = next( + payload for endpoint, payload in client.posts + if endpoint.endswith("/query") and str(payload["sql"]).startswith("SELECT metric_time") + ) + self.assertEqual(export_query["page_size"], 1000) + + def test_existing_database_time_repairs_missing_history_without_reprocessing(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + for source_type in main.EXPECTED_TYPES: + create_archive(root, source_type, COMPLETE_WINDOW) + stored_row = database_row() + store = RecordingStore(datetime(2026, 7, 31, 10, 0), [stored_row]) + history = RecordingHistory(exists=False) + + result = main.process(main.LocalSource(root), Path(temp) / "output", 3, store=store, history=history) + + self.assertIsNone(result) + self.assertEqual(store.rows_for_time_calls, 1) + self.assertEqual(history.path, "/history/2026-07-31/干扰数据处理结果_20260731100000.csv") + archived_rows = list(csv.DictReader(io.StringIO(history.payload.decode("gbk")))) + self.assertEqual(archived_rows, [stored_row]) + + def test_api_history_store_creates_date_directory_and_uploads_csv(self) -> None: + client = FakeHistoryApiClient() + history = main.ApiHistoryStore(client, "stg_test", "/history") + metric_time = datetime(2026, 8, 6, 12, 0) + + self.assertFalse(history.exists(metric_time)) + path = history.upload(metric_time, b"csv-data") + + self.assertEqual(path, "/history/2026-08-06/干扰数据处理结果_20260806120000.csv") + self.assertEqual(client.mkdir_paths, ["/history/2026-08-06"]) + self.assertEqual(client.uploads, [("/history/2026-08-06", "干扰数据处理结果_20260806120000.csv", b"csv-data")]) + + +class RecordingStore: + def __init__(self, latest_time: datetime | None = None, stored_rows: list[dict[str, str]] | None = None) -> None: + self.rows: list[dict[str, str]] = [] + self.latest_time = latest_time + self.stored_rows = stored_rows or [] + self.latest_calls = 0 + self.rows_for_time_calls = 0 + + def latest_metric_time(self) -> datetime | None: + self.latest_calls += 1 + return self.latest_time + + def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: + self.rows = list(rows) + return { + "enabled": True, + "table": main.DATABASE_TABLE, + "metric_time": rows[0]["metric_time"], + "inserted_rows": len(rows), + "refreshed_rows": 0, + "old_rows_deleted": 0, + } + + def rows_for_time(self, metric_time: datetime) -> list[dict[str, str]]: + self.rows_for_time_calls += 1 + self.asserted_metric_time = metric_time + return list(self.stored_rows) + + +class RecordingHistory: + def __init__(self, exists: bool = False) -> None: + self.exists_value = exists + self.path = "" + self.payload = b"" + + def exists(self, metric_time: datetime) -> bool: + self.checked_metric_time = metric_time + return self.exists_value + + def upload(self, metric_time: datetime, payload: bytes) -> str: + self.payload = payload + self.path = f"/history/{metric_time:%Y-%m-%d}/干扰数据处理结果_{metric_time:%Y%m%d%H%M%S}.csv" + return self.path + + +class FakeApiClient: + def __init__(self, latest_time: str | None, fail_script: bool = False) -> None: + self.latest_time = latest_time + self.fail_script = fail_script + self.posts: list[tuple[str, dict[str, object]]] = [] + + def get_json(self, endpoint: str, query: dict[str, object] | None = None, timeout: int = 30) -> object: + del endpoint, query, timeout + return [ + {"name": "metric_time"}, + {"name": "cgi"}, + {"name": "cell_name"}, + {"name": "interference_dbm"}, + ] + + def post_json(self, endpoint: str, payload: dict[str, object], timeout: int = 30) -> object: + del timeout + self.posts.append((endpoint, payload)) + if endpoint.endswith("/query"): + sql = str(payload["sql"]).lstrip() + if sql.startswith("SELECT MAX"): + return {"rows": [{"latest_time": self.latest_time}]} + if sql.startswith("SELECT metric_time"): + row = database_row() + row["metric_time"] = "2026-07-31T10:00:00" + return {"rows": [row], "total": 1} + return {"affected_rows": 0} + if self.fail_script: + return { + "results": [ + {"index": 1, "ok": True, "affected_rows": 0}, + {"index": 2, "ok": True, "affected_rows": 7}, + {"index": 3, "ok": False, "message": "insert failed", "affected_rows": 0}, + ], + "stopped": True, + } + return { + "results": [ + {"index": 1, "ok": True, "affected_rows": 0}, + {"index": 2, "ok": True, "affected_rows": 7}, + {"index": 3, "ok": True, "affected_rows": 1}, + {"index": 4, "ok": True, "affected_rows": 20}, + {"index": 5, "ok": True, "affected_rows": 0}, + ], + "stopped": False, + } + + +class FakeHistoryApiClient: + def __init__(self) -> None: + self.mkdir_paths: list[str] = [] + self.uploads: list[tuple[str, str, bytes]] = [] + + def get_json(self, endpoint: str, query: dict[str, object] | None = None, timeout: int = 30) -> object: + del endpoint, timeout + path = str((query or {}).get("path") or "") + if path == "/history": + return {"entries": []} + return {"entries": []} + + def post_json(self, endpoint: str, payload: dict[str, object], timeout: int = 30) -> object: + del endpoint, timeout + path = str(payload["path"]) + self.mkdir_paths.append(path) + return {"path": path} + + def post_file( + self, + endpoint: str, + query: dict[str, object], + filename: str, + payload: bytes, + timeout: int = 120, + ) -> object: + del endpoint, timeout + directory = str(query["path"]) + self.uploads.append((directory, filename, payload)) + return {"path": f"{directory}/{filename}"} + + +def database_row() -> dict[str, str]: + return { + "metric_time": "2026-07-31 10:00:00", + "network_type": "2.6G", + "cgi": "460-00-200-1", + "cell_name": "测试小区", + "interference_dbm": "-100.5", + "prev_interference_dbm": "", + "longitude": "", + "latitude": "", + "azimuth": "0", + "nearby_count": "0", + "nearby_26g": "0", + "nearby_700m": "0", + "nearby_tdd": "0", + "nearby_fdd": "0", + "prev_nearby_count": "", + "prev_nearby_26g": "", + "prev_nearby_700m": "", + "prev_nearby_tdd": "", + "prev_nearby_fdd": "", + } + + +def nearby_row(cgi: str, longitude: str, latitude: str, network_type: str = "2.6G") -> dict[str, str]: + row = database_row() + row["cgi"] = cgi + row["longitude"] = longitude + row["latitude"] = latitude + row["network_type"] = network_type + return row + + +def create_archive( + root: Path, + source_type: str, + window: str, + bad_header: bool = False, + interference_dbm: float = -100.5, + middle: str = "LWP_每小时_过滤110", + row_window: str = "", +) -> None: + date_dir = root / f"{window[:4]}-{window[4:6]}-{window[6:8]}" + date_dir.mkdir(parents=True, exist_ok=True) + filename = f"{source_type}_{middle}_{window}" + workbook = Workbook() + sheet = workbook.active + sheet.title = "Sheet0" + header = list(main.EXPECTED_HEADERS[source_type]) + if bad_header: + header[-1] = "unexpected" + sheet.append(header) + sheet.append(mock_row(source_type, row_window or window, interference_dbm)) + metadata = workbook.create_sheet("指标(计数器)") + metadata.append(["指标或计数器", "指标或计数器描述", "指标公式", "指标或计数器状态"]) + content = io.BytesIO() + workbook.save(content) + workbook.close() + with zipfile.ZipFile(date_dir / f"{filename}.zip", "w", zipfile.ZIP_DEFLATED) as archive: + archive.writestr(f"{filename}.xlsx", content.getvalue()) + + +def mock_row(source_type: str, window: str, interference_dbm: float = -100.5) -> list[object]: + header = main.EXPECTED_HEADERS[source_type] + values: dict[str, object] = {column: "mock" for column in header} + values.update( + { + "开始时间": datetime.strptime(window[:12], "%Y%m%d%H%M"), + "粒度": "1 小时", + "eNodeBId": 100, + "eNodeBID": 100, + "gNBId": 200, + "gNBplmn": "460-00", + "cellId": 1, + "小区ID": 1, + "masterOperatorId": "unused-source-value", + "E-UTRAN FDD小区名称": f"{source_type}-小区", + "E-UTRAN TDD小区名称": f"{source_type}-小区", + "CU小区配置名称": f"{source_type}-小区", + "小区名称": f"{source_type}-小区", + "载波平均噪声干扰(dBm)": interference_dbm, + "小区上行平均干扰电平(dBm)": interference_dbm, + } + ) + return [values[column] for column in header] + + +def create_cell_data_sources(root: Path) -> None: + raw = create_cell_data_workbook() + for index, configured_directory in enumerate(main.CELL_DATA_DIRECTORIES, start=1): + directory = root / Path(configured_directory.lstrip("/")) + directory.mkdir(parents=True, exist_ok=True) + (directory / f"江门小区信息表{index}-20260728.xlsx").write_bytes(raw) + + +def create_cell_data_workbook( + longitude: float = 113.123456, + include_latitude: bool = True, + azimuth: float | None = 30, +) -> bytes: + workbook = Workbook() + sheet = workbook.active + sheet.title = "小区信息表" + header = ["小区名称", "eNB/gNB", "CI", "经度"] + if include_latitude: + header.append("纬度") + header.append("方向角") + sheet.append(header) + row: list[object] = ["测试小区", 200, 1, longitude] + if include_latitude: + row.append(22.654321) + row.append(azimuth) + sheet.append(row) + content = io.BytesIO() + workbook.save(content) + workbook.close() + return content.getvalue() + + +if __name__ == "__main__": + unittest.main()