wx-agent/tools/ship_model_db.py
2026-07-15 16:30:58 +08:00

201 lines
7.2 KiB
Python

"""
舷号-型号映射数据库模块
用于管理舰船型号、舷号、舰名的映射关系
数据存储在 PostgreSQL 的 ship_model_mapping 表中
每艘舰一条记录
"""
import json
from typing import Dict, List, Any, Optional
from config import POSTGRES_CONNECTION_STRING
async def init_ship_model_mapping_table():
"""
初始化 ship_model_mapping 表并插入初始数据
每艘舰一条记录,包含型号、舷号、舰名
"""
try:
from psycopg_pool import AsyncConnectionPool
except ImportError:
print("[ship_model_db] psycopg_pool 未安装,无法初始化表")
return
async with AsyncConnectionPool(
POSTGRES_CONNECTION_STRING,
kwargs={"autocommit": True}
) as pool:
async with pool.connection() as conn:
async with conn.cursor() as cur:
await cur.execute("""
CREATE TABLE IF NOT EXISTS ship_model_mapping (
id SERIAL PRIMARY KEY,
model_name TEXT NOT NULL,
hull_number TEXT NOT NULL UNIQUE,
ship_name TEXT NOT NULL DEFAULT ''
)
""")
await cur.execute("""
CREATE INDEX IF NOT EXISTS ship_model_mapping_model_name_idx ON ship_model_mapping(model_name)
""")
await cur.execute("""
CREATE INDEX IF NOT EXISTS ship_model_mapping_hull_number_idx ON ship_model_mapping(hull_number)
""")
print("[ship_model_db] ship_model_mapping 表初始化完成")
await cur.execute("SELECT COUNT(*) FROM ship_model_mapping")
count = (await cur.fetchone())[0]
if count == 0:
initial_data = [
("055型驱逐舰", "102", "拉萨舰"),
("055型驱逐舰", "103", "鞍山舰"),
("055型驱逐舰", "104", "无锡舰"),
("055型驱逐舰", "105", "大连舰"),
("055型驱逐舰", "106", "延安舰"),
("055型驱逐舰", "107", "遵义舰"),
("055型驱逐舰", "108", "咸阳舰"),
("052D型驱逐舰", "172", "昆明舰"),
("052D型驱逐舰", "173", "长沙舰"),
("052D型驱逐舰", "174", "合肥舰"),
("052D型驱逐舰", "175", "银川舰"),
("052D型驱逐舰", "163", "南昌舰"),
("054A型护卫舰", "529", "舟山舰"),
("054A型护卫舰", "530", "徐州舰"),
("054A型护卫舰", "547", "临沂舰"),
("054A型护卫舰", "568", "衡阳舰"),
("054A型护卫舰", "570", "黄山舰"),
]
for model_name, hull_number, ship_name in initial_data:
await cur.execute(
"""
INSERT INTO ship_model_mapping (model_name, hull_number, ship_name)
VALUES (%s, %s, %s)
ON CONFLICT (hull_number) DO NOTHING
""",
(model_name, hull_number, ship_name)
)
print(f"[ship_model_db] 已插入 {len(initial_data)} 条初始数据")
async def load_ship_model_mapping() -> List[Dict[str, str]]:
"""
从数据库加载全部舰船映射数据
Returns:
每艘舰一条记录的列表
[
{"model_name": "055型驱逐舰", "hull_number": "102", "ship_name": "拉萨舰"},
...
]
"""
try:
from psycopg_pool import AsyncConnectionPool
except ImportError:
print("[ship_model_db] psycopg_pool 未安装,返回空映射")
return []
try:
async with AsyncConnectionPool(
POSTGRES_CONNECTION_STRING,
kwargs={"autocommit": True}
) as pool:
async with pool.connection() as conn:
async with conn.cursor() as cur:
await cur.execute(
"SELECT model_name, hull_number, ship_name FROM ship_model_mapping ORDER BY id"
)
rows = await cur.fetchall()
result = []
for row in rows:
model_name, hull_number, ship_name = row
result.append({
"model_name": model_name,
"hull_number": hull_number,
"ship_name": ship_name or "",
})
return result
except Exception as e:
print(f"[ship_model_db] 加载映射数据失败: {str(e)}")
return []
async def add_ship(
model_name: str,
hull_number: str,
ship_name: str = ""
) -> bool:
"""
添加单艘舰船
Args:
model_name: 型号名称
hull_number: 舷号
ship_name: 舰名
Returns:
是否添加成功
"""
try:
from psycopg_pool import AsyncConnectionPool
except ImportError:
print("[ship_model_db] psycopg_pool 未安装,无法添加")
return False
try:
async with AsyncConnectionPool(
POSTGRES_CONNECTION_STRING,
kwargs={"autocommit": True}
) as pool:
async with pool.connection() as conn:
async with conn.cursor() as cur:
await cur.execute(
"""
INSERT INTO ship_model_mapping (model_name, hull_number, ship_name)
VALUES (%s, %s, %s)
ON CONFLICT (hull_number) DO UPDATE SET
model_name = EXCLUDED.model_name,
ship_name = EXCLUDED.ship_name
""",
(model_name, hull_number, ship_name)
)
print(f"[ship_model_db] 舰船 {hull_number}({ship_name}) 已保存")
return True
except Exception as e:
print(f"[ship_model_db] 添加舰船失败: {str(e)}")
return False
async def delete_ship(hull_number: str) -> bool:
"""
删除单艘舰船
Args:
hull_number: 舷号
Returns:
是否删除成功
"""
try:
from psycopg_pool import AsyncConnectionPool
except ImportError:
print("[ship_model_db] psycopg_pool 未安装,无法删除")
return False
try:
async with AsyncConnectionPool(
POSTGRES_CONNECTION_STRING,
kwargs={"autocommit": True}
) as pool:
async with pool.connection() as conn:
async with conn.cursor() as cur:
await cur.execute(
"DELETE FROM ship_model_mapping WHERE hull_number = %s",
(hull_number,)
)
print(f"[ship_model_db] 舰船 {hull_number} 已删除")
return True
except Exception as e:
print(f"[ship_model_db] 删除舰船失败: {str(e)}")
return False