feat: 导入 InterferenceETL 初始源码
This commit is contained in:
@@ -0,0 +1,3 @@
|
|||||||
|
*
|
||||||
|
!Dockerfile
|
||||||
|
!requirements.txt
|
||||||
@@ -0,0 +1,9 @@
|
|||||||
|
__pycache__/
|
||||||
|
*.py[cod]
|
||||||
|
.pytest_cache/
|
||||||
|
.venv/
|
||||||
|
.env
|
||||||
|
runtime_config.py
|
||||||
|
output/
|
||||||
|
dist/
|
||||||
|
aidocs/
|
||||||
@@ -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"]
|
||||||
@@ -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
|
||||||
|
```
|
||||||
@@ -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.
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
openpyxl==3.1.5
|
||||||
@@ -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"
|
||||||
@@ -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())
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
|
||||||
@@ -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()
|
||||||
Reference in New Issue
Block a user