commit 1c0f0eff62c10df916e5b2d5901882f2902c8554 Author: Nixevol Date: Fri Jul 31 12:21:05 2026 +0800 feat: 实现干扰小时指标自动处理 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..75e842e --- /dev/null +++ b/.gitignore @@ -0,0 +1,7 @@ +__pycache__/ +*.py[cod] +.pytest_cache/ +.venv/ +.env +output/ +dist/ diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..3acb17e --- /dev/null +++ b/Dockerfile @@ -0,0 +1,7 @@ +FROM python:3.13.11-slim + +RUN pip install --no-cache-dir openpyxl==3.1.5 + +WORKDIR /workspace + +CMD ["python", "main.py"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..10adae1 --- /dev/null +++ b/README.md @@ -0,0 +1,88 @@ +# InterferenceETL + +InterferenceETL processes the newest complete hour of interference KPI files from Metrix storage. It reads seven ZIP/XLSX source types, converts each `Sheet0` to CSV, and creates one merged interference summary. + +## Output + +Each run writes one window directory: + +```text +output/2026073110001100/ +├── converted/ +│ ├── 5G下FDD干扰监控_2026073110001100.csv +│ ├── 5G干扰监控_2026073110001100.csv +│ └── ... seven source CSV files +├── interference_summary_2026073110001100.csv +└── manifest.json +``` + +The summary columns are: + +```text +hour_start,hour_end,source_type,cgi,cell_name,interference_dbm,source_path +``` + +The script never modifies source storage and does not write a database yet. Pending decisions are recorded in `docs/assumptions.md`. + +## 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 +``` + +Use `--window 2026073110001100` to require one exact source window. Without it, the script selects the newest window containing all seven source types and falls back if the newest observed window is incomplete. + +## Metrix Script Management + +The target server is offline, so build the runtime image on a connected machine and import it through Metrix Container Management: + +```powershell +docker build -t interference-etl-runtime:1.0 . +New-Item -ItemType Directory -Force dist | Out-Null +docker save -o dist\interference-etl-runtime-1.0.tar interference-etl-runtime:1.0 +``` + +The local `dist/` directory is intentionally excluded from Git. Transfer the generated TAR as a deployment artifact rather than source code. + +Create the Metrix script project with: + +```text +Name: InterferenceETL +Language: python +Base image: interference-etl-runtime:1.0 +Network: bridge +Run command: python main.py +Timeout: 1800 seconds +``` + +Configure project environment variables in Metrix instead of committing secrets: + +```json +{ + "METRIX_API_BASE_URL": "http://172.17.0.1:8000", + "METRIX_API_TOKEN": "", + "METRIX_STORAGE_ID": "stg_4d9a910d72", + "INTERFERENCE_SOURCE_ROOT": "/网优日常优化数据文档/(勿删)干扰定时小时指标", + "INTERFERENCE_OUTPUT_DIR": "/workspace/output", + "INTERFERENCE_LOOKBACK_DAYS": "3", + "INTERFERENCE_PLMN": "460-00" +} +``` + +`172.17.0.1` is the current Linux Docker default-bridge gateway used to reach the Metrix host port. Verify it before creating the online schedule because Docker bridge configuration can differ by server. + +## 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 Require an exact 16-digit source window +--lookback-days N Number of newest date directories scanned, default 3 +--api-base URL Metrix API base URL +--api-token TOKEN Metrix API token +--storage-id ID Metrix storage connection ID +--root PATH Source directory in Metrix storage +--plmn PLMN LTE/SDR CGI PLMN prefix, default 460-00 +``` diff --git a/aidocs/project_context.md b/aidocs/project_context.md new file mode 100644 index 0000000..34e4afc --- /dev/null +++ b/aidocs/project_context.md @@ -0,0 +1,14 @@ +# InterferenceETL project context + +## 2026-07-31: Initial hourly interference pipeline + +- This repository is developed as the `InterferenceETL` submodule under Metrix and is intended to run later in Metrix Script Management on an hourly schedule. +- `main.py` reads the configured Metrix SFTP storage through Metrix API Token authentication or a local directory for tests. It selects the newest hour containing all seven known interference source types and falls back from a newer incomplete hour. +- Each selected ZIP must contain exactly one XLSX. The script strictly validates the known `Sheet0` header, writes one full UTF-8-BOM CSV per source type, and writes one merged CSV containing `hour_start`, `hour_end`, `source_type`, `cgi`, `cell_name`, `interference_dbm`, and `source_path`. +- LTE/SDR CGI currently uses `460-00-{node}-{cell}`; NR uses `masterOperatorId` unchanged. This is an explicit assumption pending user confirmation. +- Runs are idempotent at the output-window directory level. Generation happens in a scoped temporary directory, then replaces only the same window below the configured output root. Source storage is never modified. +- Database persistence is intentionally not implemented because the target schema and credentials are pending. `manifest.json` records input paths, sizes, SHA-256 hashes, row counts, warnings, and generated files for later ingestion auditing. +- The Metrix script container needs `openpyxl==3.1.5`, bridge networking, `python main.py`, a reachable `METRIX_API_BASE_URL`, and `METRIX_API_TOKEN` injected through project environment settings. Secrets must not be committed. +- Read-only validation against the current Metrix storage selected window `2026073110001100`, processed all seven source ZIP files, and produced 1,831 summary rows. Per-source row counts were `7 / 804 / 45 / 86 / 700 / 164 / 25` in `EXPECTED_TYPES` order; sampled CGI, cell name, and interference values matched the source workbooks. +- The Linux/amd64 runtime image is `interference-etl-runtime:1.0`. Its offline archive is generated locally at `dist/interference-etl-runtime-1.0.tar` (46,828,032 bytes, SHA-256 `B65B94941BB304CA955224B885582E36BE05E390CBD76FC4C069494352E77746`) and remains excluded from Git. +- Mock tests pass on Windows Python and inside the runtime image. They cover complete-hour fallback, strict schema rejection, summary extraction, and cross-midnight window parsing. diff --git a/docs/assumptions.md b/docs/assumptions.md new file mode 100644 index 0000000..0466cff --- /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. +- LTE/SDR CGI is derived as `{PLMN}-{eNodeBId}-{cellId}` with default PLMN `460-00`. +- NR CGI uses the source `masterOperatorId` unchanged. Confirm whether the final database should retain this source format or normalize it. +- Output CSV files use UTF-8 with BOM so they open correctly in Excel. +- Source files are read-only. The script writes only below its configured output directory. + +## Pending user decisions + +- Target database connection, database name, table name, and credentials. +- Whether the seven converted full CSV files must be retained after successful database import. +- Whether the `指标(计数器)` metadata sheet must also be stored. +- Final CGI normalization rules, especially for NR `masterOperatorId`. +- 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 complete-hour fallback, schema validation, CSV conversion, CGI extraction, summary merging, and cross-midnight windows. No mock credentials or fake database writes are present in production code. diff --git a/main.py b/main.py new file mode 100644 index 0000000..f12af30 --- /dev/null +++ b/main.py @@ -0,0 +1,491 @@ +from __future__ import annotations + +import argparse +import csv +import hashlib +import io +import json +import os +from dataclasses import dataclass +from datetime import datetime, timedelta, 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 Request, urlopen +import warnings +import zipfile + +from openpyxl import load_workbook + + +EXPECTED_TYPES = ( + "5G下FDD干扰监控", + "5G干扰监控", + "700M下FDD干扰监控", + "700M干扰监控", + "SDR_FDD干扰监控", + "SDR_TDD干扰监控", + "反开RD干扰监控", +) + +DEFAULT_STORAGE_ID = "stg_4d9a910d72" +DEFAULT_ROOT = "/网优日常优化数据文档/(勿删)干扰定时小时指标" +FILE_RE = re.compile(r"^(?P.+)_LWP_每小时_过滤110_(?P\d{16})\.zip$", re.IGNORECASE) + +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 = ( + "hour_start", + "hour_end", + "source_type", + "cgi", + "cell_name", + "interference_dbm", + "source_path", +) + + +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: ... + + +class ApiSource: + def __init__(self, base_url: str, token: str, storage_id: str, root: str) -> None: + if not token: + raise ProcessingError("METRIX_API_TOKEN is required in API mode") + self.base_url = base_url.rstrip("/") + self.token = token + 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._request(endpoint, {"path": path}, timeout=120) + + def _list_dir(self, path: str) -> list[dict[str, object]]: + endpoint = f"/api/storages/{quote(self.storage_id, safe='')}/files" + payload = self._request(endpoint, {"path": path, "recursive": "false"}) + return json.loads(payload.decode("utf-8")).get("entries", []) + + def _request(self, endpoint: str, query: dict[str, object], timeout: int = 30) -> bytes: + url = f"{self.base_url}{endpoint}?{urlencode(query)}" + request = Request(url, headers={"Authorization": f"Bearer {self.token}", "Accept": "application/json"}) + last_error: Exception | None = None + for attempt in range(3): + try: + with urlopen(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}") + + +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 parse_candidate(path: str, size: int = 0) -> Candidate | None: + name = posixpath.basename(path.replace("\\", "/")) + match = FILE_RE.fullmatch(name) + if not match or match.group("source_type") not in EXPECTED_TYPES: + return None + return Candidate(match.group("source_type"), match.group("window"), path, size) + + +def select_window(candidates: list[Candidate], requested: str = "") -> tuple[str, dict[str, Candidate], list[str]]: + grouped: dict[str, dict[str, Candidate]] = {} + duplicates: list[str] = [] + for candidate in candidates: + window_group = grouped.setdefault(candidate.window, {}) + if candidate.source_type in window_group: + duplicates.append(f"{candidate.window}/{candidate.source_type}") + window_group[candidate.source_type] = candidate + if duplicates: + raise ProcessingError(f"Duplicate source files: {', '.join(sorted(duplicates))}") + complete = sorted(window for window, items in grouped.items() if all(name in items for name in EXPECTED_TYPES)) + if requested: + if requested not in complete: + present = sorted(grouped.get(requested, {})) + missing = [name for name in EXPECTED_TYPES if name not in present] + raise ProcessingError(f"Requested window is incomplete: {requested}; missing={missing}") + selected = requested + elif complete: + selected = complete[-1] + else: + raise ProcessingError("No hour contains all seven interference source types") + warnings_out: list[str] = [] + latest_seen = max(grouped) if grouped else "" + if latest_seen and latest_seen != selected: + missing = [name for name in EXPECTED_TYPES if name not in grouped[latest_seen]] + warnings_out.append(f"Latest observed window {latest_seen} is incomplete; using {selected}; missing={missing}") + return selected, grouped[selected], warnings_out + + +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 = "", plmn: str = "460-00") -> Path: + candidates = source.candidates(lookback_days) + window, selected, warnings_out = select_window(candidates, requested_window) + temp_dir = output_root.resolve() / f".{window}.tmp-{os.getpid()}" + final_dir = output_root.resolve() / window + 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]] = [] + manifest_files: list[dict[str, object]] = [] + 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}_{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_rows.append(summary_record(record, source_type, candidate.path, window, plmn, row_number)) + manifest_files.append( + { + "source_type": source_type, + "source_path": candidate.path, + "source_size": candidate.size, + "sha256": hashlib.sha256(raw_zip).hexdigest(), + "xlsx_member": member_name, + "rows": len(rows), + "converted_csv": f"converted/{csv_name}", + } + ) + + summary_rows.sort(key=lambda item: (EXPECTED_TYPES.index(item["source_type"]), item["cgi"], item["cell_name"])) + summary_name = f"interference_summary_{window}.csv" + write_dict_csv(temp_dir / summary_name, SUMMARY_HEADER, summary_rows) + manifest = { + "generated_at": datetime.now(timezone.utc).isoformat(), + "window": window, + "hour_start": window_bounds(window)[0], + "hour_end": window_bounds(window)[1], + "source_types": list(EXPECTED_TYPES), + "source_file_count": len(manifest_files), + "summary_rows": len(summary_rows), + "warnings": warnings_out, + "files": manifest_files, + "summary_csv": summary_name, + } + (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(f"selected_window={window}") + print(f"source_files={len(manifest_files)}") + print(f"summary_rows={len(summary_rows)}") + for message in warnings_out: + print(f"warning={message}") + print(f"output={final_dir}") + return final_dir + + +def summary_record( + record: dict[str, object], source_type: str, source_path: str, window: str, plmn: str, row_number: int +) -> dict[str, str]: + hour_start = normalize_cell(record["开始时间"]) + expected_start, hour_end = window_bounds(window) + if hour_start != expected_start: + raise ProcessingError(f"{source_path}: row {row_number} time {hour_start!r} does not match {expected_start!r}") + + if source_type in ("5G干扰监控", "700M干扰监控"): + cgi = required(record, "masterOperatorId", 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(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(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(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) + + return { + "hour_start": hour_start, + "hour_end": hour_end, + "source_type": source_type, + "cgi": cgi, + "cell_name": cell_name, + "interference_dbm": interference, + "source_path": source_path, + } + + +def build_cgi(plmn: str, record: dict[str, object], node_column: str, cell_column: str, path: str, row: int) -> str: + return f"{plmn}-{required(record, node_column, path, row)}-{required(record, cell_column, path, row)}" + + +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 window_bounds(window: str) -> tuple[str, str]: + if not re.fullmatch(r"\d{16}", window): + raise ProcessingError(f"Invalid window: {window}") + start = datetime.strptime(window[:12], "%Y%m%d%H%M") + end_hour = int(window[12:14]) + end_minute = int(window[14:16]) + end = start.replace(hour=end_hour, minute=end_minute) + if end <= start: + end += timedelta(days=1) + return start.strftime("%Y-%m-%d %H:%M:%S"), end.strftime("%Y-%m-%d %H:%M:%S") + + +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 write_csv(path: Path, header: tuple[str, ...], rows: list[tuple[object, ...]]) -> None: + with path.open("w", encoding="utf-8-sig", 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: + with path.open("w", encoding="utf-8-sig", newline="") as file: + writer = csv.DictWriter(file, fieldnames=header, extrasaction="raise") + writer.writeheader() + writer.writerows(rows) + + +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 hour of interference KPI files") + 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 exact 16-digit source window") + parser.add_argument("--lookback-days", type=int, default=int(os.getenv("INTERFERENCE_LOOKBACK_DAYS", "3"))) + parser.add_argument("--api-base", default=os.getenv("METRIX_API_BASE_URL", "http://172.17.0.1:8000")) + parser.add_argument("--api-token", default=os.getenv("METRIX_API_TOKEN", "")) + 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("--plmn", default=os.getenv("INTERFERENCE_PLMN", "460-00")) + 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") + source: Source = LocalSource(args.source_dir) if args.source_dir else ApiSource(args.api_base, args.api_token, args.storage_id, args.root) + process(source, args.output_dir, args.lookback_days, args.window, args.plmn) + 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/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..cbaa947 --- /dev/null +++ b/tests/test_pipeline.py @@ -0,0 +1,118 @@ +from __future__ import annotations + +import csv +from datetime import datetime +import io +import json +from pathlib import Path +import tempfile +import unittest +import zipfile + +from openpyxl import Workbook + +import main + + +COMPLETE_WINDOW = "2026073110001100" +INCOMPLETE_WINDOW = "2026073111001200" + + +class PipelineTest(unittest.TestCase): + def test_latest_complete_window_is_converted_and_merged(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + output = Path(temp) / "output" + for source_type in main.EXPECTED_TYPES: + create_archive(root, source_type, COMPLETE_WINDOW) + create_archive(root, main.EXPECTED_TYPES[0], INCOMPLETE_WINDOW) + + result = main.process(main.LocalSource(root), output, lookback_days=3) + + self.assertEqual(result.name, COMPLETE_WINDOW) + converted = sorted((result / "converted").glob("*.csv")) + self.assertEqual(len(converted), 7) + with (result / f"interference_summary_{COMPLETE_WINDOW}.csv").open(encoding="utf-8-sig", newline="") as file: + rows = list(csv.DictReader(file)) + self.assertEqual(len(rows), 7) + self.assertEqual({row["source_type"] for row in rows}, set(main.EXPECTED_TYPES)) + self.assertIn("460-00-100-1", {row["cgi"] for row in rows}) + self.assertIn("46000-100-1", {row["cgi"] for row in rows}) + self.assertTrue(all(row["interference_dbm"] == "-100.5" for row in rows)) + + manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8")) + self.assertEqual(manifest["source_file_count"], 7) + self.assertEqual(manifest["summary_rows"], 7) + self.assertEqual(len(manifest["warnings"]), 1) + self.assertIn(INCOMPLETE_WINDOW, manifest["warnings"][0]) + + def test_requested_incomplete_window_is_rejected(self) -> None: + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) / "source" + create_archive(root, main.EXPECTED_TYPES[0], INCOMPLETE_WINDOW) + + with self.assertRaisesRegex(main.ProcessingError, "Requested window is incomplete"): + main.process(main.LocalSource(root), Path(temp) / "output", 3, INCOMPLETE_WINDOW) + + 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_cross_midnight_window(self) -> None: + self.assertEqual( + main.window_bounds("2026073023000000"), + ("2026-07-30 23:00:00", "2026-07-31 00:00:00"), + ) + + +def create_archive(root: Path, source_type: str, window: str, bad_header: bool = False) -> 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}_LWP_每小时_过滤110_{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, window)) + 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) -> 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, + "cellId": 1, + "小区ID": 1, + "masterOperatorId": "46000-100-1", + "E-UTRAN FDD小区名称": f"{source_type}-小区", + "E-UTRAN TDD小区名称": f"{source_type}-小区", + "CU小区配置名称": f"{source_type}-小区", + "小区名称": f"{source_type}-小区", + "载波平均噪声干扰(dBm)": -100.5, + "小区上行平均干扰电平(dBm)": -100.5, + } + ) + return [values[column] for column in header] + + +if __name__ == "__main__": + unittest.main()