diff --git a/README.md b/README.md index a90861f..5714624 100644 --- a/README.md +++ b/README.md @@ -26,6 +26,13 @@ The script never modifies or deletes source storage files. By default it writes The database table contains one time column, `metric_time DATETIME`, which is the source KPI start time. A same-hour rerun replaces that whole batch. After a successful insert, rows for every other hour are deleted in the same transaction. The script refuses to replace a newer database hour with an older source hour. +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. + +MySQL connection values, database `metrix`, and table `interference_hourly_summary` are constants at the top of `main.py`; they are not project environment variables. + ## Local development ```powershell @@ -45,7 +52,7 @@ 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` plus `openpyxl`, `et_xmlfile`, and `PyMySQL`. The builder verifies imports 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 ZIP is 325,932 bytes with SHA-256 `D57930E36B3C70DF7D209C4A83546FEA3B73F87022910D1FDB548D6BED71675D`. +The generated `dist/InterferenceETL-offline-1.1.zip` contains `main.py` plus `openpyxl`, `et_xmlfile`, and `PyMySQL`. The builder verifies imports 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 ZIP is 325,489 bytes with SHA-256 `4F655D35F89A826F28A41578C9EE3E347CE31A27619D3D06AE0FD9587660DCAF`. Create the project, then upload and extract the ZIP in its Script Management workspace. Use these project settings: @@ -58,7 +65,7 @@ Run command: python main.py Timeout: 1800 seconds ``` -Configure project environment variables in Metrix instead of committing secrets: +Configure only the Metrix source settings in the project environment: ```json { @@ -67,18 +74,11 @@ Configure project environment variables in Metrix instead of committing secrets: "METRIX_STORAGE_ID": "stg_4d9a910d72", "INTERFERENCE_SOURCE_ROOT": "/网优日常优化数据文档/(勿删)干扰定时小时指标", "INTERFERENCE_OUTPUT_DIR": "/workspace/output", - "INTERFERENCE_LOOKBACK_DAYS": "3", - "INTERFERENCE_PLMN": "460-00", - "INTERFERENCE_DB_HOST": "", - "INTERFERENCE_DB_PORT": "3306", - "INTERFERENCE_DB_USER": "", - "INTERFERENCE_DB_PASSWORD": "", - "INTERFERENCE_DB_NAME": "", - "INTERFERENCE_DB_TABLE": "interference_hourly_summary" + "INTERFERENCE_LOOKBACK_DAYS": "3" } ``` -`172.17.0.1` is the current Linux Docker default-bridge gateway used to reach the Metrix host port. Verify both API and MySQL addresses before creating the online schedule because Docker network configuration can differ by server. A container name such as `ShareMySQL` resolves only when the script container joins the same Docker network. +`172.17.0.1` is the current Linux Docker default-bridge gateway used by the script to reach both the Metrix API and the host-published MySQL port. `ShareMySQL` itself does not resolve from the script project's bridge network. ## CLI options @@ -91,12 +91,5 @@ Configure project environment variables in Metrix instead of committing secrets: --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 --no-database Generate CSV files without writing MySQL ---db-host HOST MySQL host or INTERFERENCE_DB_HOST ---db-port PORT MySQL port, default 3306 ---db-user USER MySQL user ---db-password VALUE MySQL password; prefer the environment variable ---db-name NAME MySQL database name ---db-table NAME Table name, default interference_hourly_summary ``` diff --git a/aidocs/project_context.md b/aidocs/project_context.md index e5aebd1..acb4fba 100644 --- a/aidocs/project_context.md +++ b/aidocs/project_context.md @@ -5,18 +5,18 @@ - 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. +- NR CGI uses `{gNBplmn}-{gNBId}-{cellId}`. All 4G sources use `460-00-{node}-{cell}`, mapped to the actual node and cell column names in each workbook schema. - 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. - `manifest.json` records input paths, sizes, SHA-256 hashes, row counts, warnings, generated files, and the database write result for ingestion auditing. -- The Metrix script container needs `openpyxl==3.1.5`, `PyMySQL==1.1.2`, bridge networking, `python main.py`, reachable API/MySQL addresses, and credentials injected through project environment settings. Secrets must not be committed. +- The Metrix script container needs `openpyxl==3.1.5`, `PyMySQL==1.1.2`, bridge networking, and `python main.py`. Only Metrix Storage API settings are injected through project environment settings. - 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. - `interference-etl-runtime:1.1` is now a local build helper only. Script projects use the standard `python:3.13.11-slim` image and upload `dist/InterferenceETL-offline-1.1.zip`, which vendors `openpyxl`, `et_xmlfile`, and `PyMySQL` at the workspace root. - Six Mock tests pass on Windows Python and inside the runtime image. They cover complete-hour fallback, strict schema rejection, summary extraction, cross-midnight window parsing, transactional database replacement, and rollback protection. ## 2026-07-31: Latest-hour MySQL retention -- Scheduled runs write to MySQL by default. File-only development runs must explicitly pass `--no-database`; connection values come from `INTERFERENCE_DB_*` environment variables. -- The default table is `interference_hourly_summary`. Its only time column is `metric_time DATETIME`, populated from the source KPI `开始时间`, so users can identify the hour represented by every row. +- Scheduled runs write to MySQL by default. File-only development runs must explicitly pass `--no-database`; the MySQL host, port, account, password, database, and table are constants at the top of `main.py`. +- The fixed database is `metrix` and the fixed table is `interference_hourly_summary`. Its only time column is `metric_time DATETIME`, populated from the source KPI `开始时间`, so users can identify the hour represented by every row. - One transaction deletes the selected hour for idempotent refresh, inserts its complete batch, then deletes every other database hour. Any failure rolls back the data changes, and an older selected hour cannot replace a newer hour already stored. - Database retention never deletes or modifies Metrix Storage/SFTP source ZIP or XLSX files. Generated CSV retention remains a separate pending decision. @@ -24,4 +24,11 @@ - `scripts/build_offline_package.py` copies the three runtime distributions from the verified dependency image, adds the application files, and creates a ZIP that Metrix can extract without preserving executable bits or symlinks. - The builder runs an import check inside `python:3.13.11-slim`; the package therefore needs neither a custom server image nor online package installation. -- Current artifact: `dist/InterferenceETL-offline-1.1.zip`, 325,932 bytes, SHA-256 `D57930E36B3C70DF7D209C4A83546FEA3B73F87022910D1FDB548D6BED71675D`. The builder excludes bytecode caches, and the ignored `dist/` directory is the only location for this generated package. +- Current artifact: `dist/InterferenceETL-offline-1.1.zip`, 325,489 bytes, SHA-256 `4F655D35F89A826F28A41578C9EE3E347CE31A27619D3D06AE0FD9587660DCAF`. The builder excludes bytecode caches, and the ignored `dist/` directory is the only location for this generated package. + +## 2026-07-31: Fixed CGI and database configuration + +- Replaced the NR `masterOperatorId` passthrough with `gNBplmn-gNBId-cellId`. The remaining five 4G sources use the fixed `460-00` prefix plus their schema-specific base-station and cell ID columns. +- Removed CGI PLMN and MySQL connection command-line/environment options. Database configuration now lives as a small constant block at the top of `main.py`; the verified bridge-network address is `172.17.0.1:3306` because `ShareMySQL` DNS is unavailable from the default script network. +- Read-only online validation selected complete window `2026073115001600` and checked all 2,275 generated CGI values against the seven converted source CSV files. Per-source row counts were `6 / 1,198 / 48 / 91 / 731 / 179 / 22` in `EXPECTED_TYPES` order. +- A connection-only check using the fixed constants returned database `metrix`. It did not create the table, insert rows, or delete existing data; all remote validation files were removed afterward. diff --git a/docs/assumptions.md b/docs/assumptions.md index 3541261..be75dbd 100644 --- a/docs/assumptions.md +++ b/docs/assumptions.md @@ -6,8 +6,8 @@ - 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. +- 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. - 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. - MySQL is the scheduled-run output. The table uses one `metric_time DATETIME` column containing the source KPI start time. @@ -15,14 +15,12 @@ ## Pending user decisions -- Target database host, database name, user, and credentials. The default table name is `interference_hourly_summary` and can be overridden. - 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, cross-midnight windows, transactional latest-hour replacement, and database rollback protection. No mock credentials or fake database writes are present in production code. +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 index abe53e6..a85a23c 100644 --- a/main.py +++ b/main.py @@ -38,6 +38,14 @@ DEFAULT_STORAGE_ID = "stg_4d9a910d72" DEFAULT_ROOT = "/网优日常优化数据文档/(勿删)干扰定时小时指标" FILE_RE = re.compile(r"^(?P.+)_LWP_每小时_过滤110_(?P\d{16})\.zip$", re.IGNORECASE) +LTE_PLMN = "460-00" +DATABASE_HOST = "172.17.0.1" +DATABASE_PORT = 3306 +DATABASE_USER = "root" +DATABASE_PASSWORD = "OSp!jmgm@26" +DATABASE_NAME = "metrix" +DATABASE_TABLE = "interference_hourly_summary" + HEADER_FDD = ( "开始时间", "粒度", @@ -138,10 +146,6 @@ SUMMARY_HEADER = ( "source_path", ) -DEFAULT_DB_TABLE = "interference_hourly_summary" -DB_IDENTIFIER_RE = re.compile(r"^[A-Za-z0-9_]+$") - - class ProcessingError(RuntimeError): pass @@ -165,24 +169,8 @@ class SummaryStore(Protocol): class MySQLSummaryStore: - def __init__( - self, - host: str, - port: int, - user: str, - password: str, - database: str, - table: str = DEFAULT_DB_TABLE, - connect_factory: Callable[[], object] | None = None, - ) -> None: - if not DB_IDENTIFIER_RE.fullmatch(table): - raise ProcessingError(f"Invalid database table name: {table!r}") - self.host = host - self.port = port - self.user = user - self.password = password - self.database = database - self.table = table + def __init__(self, connect_factory: Callable[[], object] | None = None) -> None: + self.table = DATABASE_TABLE self.connect_factory = connect_factory def replace_latest(self, rows: list[dict[str, str]]) -> dict[str, object]: @@ -273,11 +261,11 @@ class MySQLSummaryStore: raise ProcessingError("PyMySQL is required for database output") from exc try: return pymysql.connect( - host=self.host, - port=self.port, - user=self.user, - password=self.password, - database=self.database, + host=DATABASE_HOST, + port=DATABASE_PORT, + user=DATABASE_USER, + password=DATABASE_PASSWORD, + database=DATABASE_NAME, charset="utf8mb4", autocommit=False, connect_timeout=15, @@ -438,7 +426,6 @@ def process( output_root: Path, lookback_days: int, requested_window: str = "", - plmn: str = "460-00", store: SummaryStore | None = None, ) -> Path: candidates = source.candidates(lookback_days) @@ -463,7 +450,7 @@ def process( 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)) + summary_rows.append(summary_record(record, source_type, candidate.path, window, row_number)) manifest_files.append( { "source_type": source_type, @@ -518,7 +505,7 @@ def process( def summary_record( - record: dict[str, object], source_type: str, source_path: str, window: str, plmn: str, row_number: int + record: dict[str, object], source_type: str, source_path: str, window: str, row_number: int ) -> dict[str, str]: hour_start = normalize_cell(record["开始时间"]) expected_start, hour_end = window_bounds(window) @@ -526,19 +513,19 @@ def summary_record( 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) + 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(plmn, record, "eNodeBID", "小区ID", source_path, row_number) + 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(plmn, record, "eNodeBId", "cellId", source_path, row_number) + 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(plmn, record, "eNodeBId", "cellId", source_path, row_number) + 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) @@ -553,8 +540,9 @@ def summary_record( } -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 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: @@ -616,46 +604,17 @@ def build_parser() -> argparse.ArgumentParser: 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")) parser.add_argument("--no-database", action="store_true", help="Generate files without writing MySQL") - parser.add_argument("--db-host", default=os.getenv("INTERFERENCE_DB_HOST", "")) - parser.add_argument("--db-port", type=int, default=int(os.getenv("INTERFERENCE_DB_PORT", "3306"))) - parser.add_argument("--db-user", default=os.getenv("INTERFERENCE_DB_USER", "")) - parser.add_argument("--db-password", default=os.getenv("INTERFERENCE_DB_PASSWORD", "")) - parser.add_argument("--db-name", default=os.getenv("INTERFERENCE_DB_NAME", "")) - parser.add_argument("--db-table", default=os.getenv("INTERFERENCE_DB_TABLE", DEFAULT_DB_TABLE)) return parser -def mysql_store_from_args(args: argparse.Namespace) -> MySQLSummaryStore: - required_values = { - "INTERFERENCE_DB_HOST": args.db_host, - "INTERFERENCE_DB_USER": args.db_user, - "INTERFERENCE_DB_PASSWORD": args.db_password, - "INTERFERENCE_DB_NAME": args.db_name, - } - missing = [name for name, value in required_values.items() if not value] - if missing: - raise ProcessingError(f"Missing database settings: {', '.join(missing)}; use --no-database for file-only runs") - if not 1 <= args.db_port <= 65535: - raise ProcessingError("db-port must be between 1 and 65535") - return MySQLSummaryStore( - host=args.db_host, - port=args.db_port, - user=args.db_user, - password=args.db_password, - database=args.db_name, - table=args.db_table, - ) - - 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) - store = None if args.no_database else mysql_store_from_args(args) - process(source, args.output_dir, args.lookback_days, args.window, args.plmn, store) + store = None if args.no_database else MySQLSummaryStore() + process(source, args.output_dir, args.lookback_days, args.window, store=store) return 0 diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index a24540a..87d4804 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -37,8 +37,11 @@ class PipelineTest(unittest.TestCase): 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}) + rows_by_type = {row["source_type"]: row for row in rows} + self.assertEqual(rows_by_type["5G干扰监控"]["cgi"], "460-00-200-1") + self.assertEqual(rows_by_type["700M干扰监控"]["cgi"], "460-00-200-1") + for source_type in set(main.EXPECTED_TYPES) - {"5G干扰监控", "700M干扰监控"}: + self.assertEqual(rows_by_type[source_type]["cgi"], "460-00-100-1") self.assertTrue(all(row["interference_dbm"] == "-100.5" for row in rows)) manifest = json.loads((result / "manifest.json").read_text(encoding="utf-8")) @@ -74,14 +77,7 @@ class PipelineTest(unittest.TestCase): def test_mysql_store_replaces_same_hour_and_deletes_other_hours(self) -> None: connection = FakeConnection(datetime(2026, 7, 31, 10), same_hour_rows=7, old_rows=20) - store = main.MySQLSummaryStore( - host="", - port=3306, - user="", - password="", - database="", - connect_factory=lambda: connection, - ) + store = main.MySQLSummaryStore(connect_factory=lambda: connection) result = store.replace_latest([database_row()]) @@ -96,14 +92,7 @@ class PipelineTest(unittest.TestCase): def test_mysql_store_refuses_to_replace_a_newer_hour(self) -> None: connection = FakeConnection(datetime(2026, 7, 31, 11)) - store = main.MySQLSummaryStore( - host="", - port=3306, - user="", - password="", - database="", - connect_factory=lambda: connection, - ) + store = main.MySQLSummaryStore(connect_factory=lambda: connection) with self.assertRaisesRegex(main.ProcessingError, "already contains newer metric time"): store.replace_latest([database_row()]) @@ -122,7 +111,7 @@ class RecordingStore: self.rows = list(rows) return { "enabled": True, - "table": main.DEFAULT_DB_TABLE, + "table": main.DATABASE_TABLE, "metric_time": rows[0]["hour_start"], "inserted_rows": len(rows), "refreshed_rows": 0, @@ -183,7 +172,7 @@ def database_row() -> dict[str, str]: "hour_start": "2026-07-31 10:00:00", "hour_end": "2026-07-31 11:00:00", "source_type": "5G干扰监控", - "cgi": "46000-100-1", + "cgi": "460-00-200-1", "cell_name": "测试小区", "interference_dbm": "-100.5", "source_path": "/source.zip", @@ -220,9 +209,11 @@ def mock_row(source_type: str, window: str) -> list[object]: "粒度": "1 小时", "eNodeBId": 100, "eNodeBID": 100, + "gNBId": 200, + "gNBplmn": "460-00", "cellId": 1, "小区ID": 1, - "masterOperatorId": "46000-100-1", + "masterOperatorId": "unused-source-value", "E-UTRAN FDD小区名称": f"{source_type}-小区", "E-UTRAN TDD小区名称": f"{source_type}-小区", "CU小区配置名称": f"{source_type}-小区",