426 lines
15 KiB
Rust
Raw Normal View History

//! 组织树 API
//!
//! 两个独立接口:
//! 1. /api/organization-tree — 返回多级组织->项目树
//! 2. /api/cabinets/filter — 根据组织/项目筛选设备列表
use axum::extract::{State, Query};
use axum::Json;
use serde::Deserialize;
use serde_json::{json, Value};
use super::organizations::{CabinetRow, OrganizationRow, ProjectRow, TreeNode};
use crate::error::AppError;
use crate::middleware::auth::{self, CurrentUser};
use crate::commands::AppState;
/// SQL查询专用结构体只包含数据库字段
#[derive(Debug, Clone, sqlx::FromRow)]
struct CabinetSqlRow {
id: i64,
project_id: Option<i64>,
abstract_id: Option<String>,
imei: String,
iccid: Option<String>,
name: Option<String>,
address: Option<String>,
status: i8,
board_count: Option<i64>,
channel_count: Option<i64>,
idle_count: Option<i64>,
charging_count: Option<i64>,
full_count: Option<i64>,
fault_count: Option<i64>,
}
impl From<CabinetSqlRow> for CabinetRow {
fn from(row: CabinetSqlRow) -> Self {
CabinetRow {
id: row.id,
project_id: row.project_id,
abstract_id: row.abstract_id,
imei: row.imei,
iccid: row.iccid,
name: row.name,
address: row.address,
status: row.status,
board_count: row.board_count,
channel_count: row.channel_count,
idle_count: row.idle_count,
charging_count: row.charging_count,
full_count: row.full_count,
fault_count: row.fault_count,
rssi: None,
pow_fail_dc: None,
pow_fail_ac: None,
is_online: None,
}
}
}
/// 组织树接口 — 返回多级组织->项目树(含设备数量)
pub async fn organization_tree(
user: CurrentUser,
State(state): State<AppState>,
) -> Result<Json<Value>, AppError> {
auth::check_permission(&user, "device:view")?;
let db = &state.mysql;
// 获取所有组织
let orgs: Vec<OrganizationRow> = sqlx::query_as(
"SELECT id, name, parent_id FROM organizations ORDER BY id"
)
.fetch_all(db)
.await?;
// 获取所有项目
let projects: Vec<ProjectRow> = sqlx::query_as(
"SELECT id, organization_id, name FROM projects ORDER BY id"
)
.fetch_all(db)
.await?;
// 获取每个项目的设备数量
let project_counts: Vec<(i64, i64)> = sqlx::query_as(
"SELECT project_id, COUNT(*) as cnt FROM cabinets GROUP BY project_id"
)
.fetch_all(db)
.await?;
let project_count_map: std::collections::HashMap<i64, i64> = project_counts.into_iter().collect();
// 递归构建树
fn build_tree(
orgs: &[OrganizationRow],
projects: &[ProjectRow],
project_count_map: &std::collections::HashMap<i64, i64>,
parent_id: Option<i64>,
) -> Vec<TreeNode> {
orgs
.iter()
.filter(|o| o.parent_id == parent_id)
.map(|org| {
// 递归获取子组织
let children = build_tree(orgs, projects, project_count_map, Some(org.id));
// 获取该组织下的项目
let org_projects: Vec<TreeNode> = projects
.iter()
.filter(|p| p.organization_id == Some(org.id))
.map(|p| {
let count = project_count_map.get(&p.id).copied().unwrap_or(0);
TreeNode {
id: p.id,
name: format!("{} ({})", p.name, count),
children: None,
abstract_id: None,
imei: None,
status: None,
}
})
.collect();
// 合并子组织和项目
let mut all_children = children;
all_children.extend(org_projects);
// 计算本组织下的总设备数量
let total_count: i64 = all_children.iter()
.map(|c| {
// 从项目名称中提取数量
if let Some(start) = c.name.rfind('(') {
if let Some(end) = c.name.rfind(')') {
c.name[start+1..end].parse().unwrap_or(0)
} else { 0 }
} else { 0 }
})
.sum();
TreeNode {
id: org.id,
name: format!("{} ({})", org.name, total_count),
children: if all_children.is_empty() { None } else { Some(all_children) },
abstract_id: None,
imei: None,
status: None,
}
})
.collect()
}
let tree = build_tree(&orgs, &projects, &project_count_map, None);
Ok(Json(json!(tree)))
}
/// 设备列表筛选参数
#[derive(Deserialize)]
pub struct CabinetFilterParams {
pub organization_id: Option<i64>,
pub project_id: Option<i64>,
pub page: Option<i64>,
pub page_size: Option<i64>,
pub bound: Option<String>,
}
/// 设备列表接口 — 根据组织/项目筛选
pub async fn filter_cabinets(
user: CurrentUser,
State(state): State<AppState>,
Query(params): Query<CabinetFilterParams>,
) -> Result<Json<Value>, AppError> {
auth::check_permission(&user, "device:view")?;
let db = &state.mysql;
let redis = &state.redis;
let page = params.page.unwrap_or(1).max(1);
let page_size = params.page_size.unwrap_or(12).clamp(1, 100);
let offset = (page - 1) * page_size;
// 基础SQL联表查询获取通道状态统计
let base_sql = "
SELECT
c.id, c.project_id, c.abstract_id, c.imei, c.iccid, c.name, c.address, c.status,
(SELECT COUNT(*) FROM cabin_boards cb WHERE cb.cabinet_id = c.id) as board_count,
(SELECT COUNT(*) FROM cabin_boards cb2
JOIN compartments comp ON comp.cabin_board_id = cb2.id
WHERE cb2.cabinet_id = c.id) as channel_count,
(SELECT COUNT(*) FROM cabin_boards cb3
JOIN compartments comp2 ON comp2.cabin_board_id = cb3.id
WHERE cb3.cabinet_id = c.id AND comp2.status = 0) as idle_count,
(SELECT COUNT(*) FROM cabin_boards cb4
JOIN compartments comp3 ON comp3.cabin_board_id = cb4.id
WHERE cb4.cabinet_id = c.id AND comp3.status = 1) as charging_count,
(SELECT COUNT(*) FROM cabin_boards cb5
JOIN compartments comp4 ON comp4.cabin_board_id = cb5.id
WHERE cb5.cabinet_id = c.id AND comp4.status = 2) as full_count,
(SELECT COUNT(*) FROM cabin_boards cb6
JOIN compartments comp5 ON comp5.cabin_board_id = cb6.id
WHERE cb6.cabinet_id = c.id AND comp5.status = 3) as fault_count
FROM cabinets c
";
let (mut where_clause, binds): (String, Vec<String>) = if let Some(proj_id) = params.project_id {
("WHERE c.project_id = ?".into(), vec![proj_id.to_string()])
} else if let Some(org_id) = params.organization_id {
("JOIN projects p ON c.project_id = p.id
WHERE p.organization_id IN (
SELECT id FROM organizations WHERE id = ? OR parent_id = ?
)".into(), vec![org_id.to_string(), org_id.to_string()])
} else if user.role_level >= 2 {
(String::new(), vec![])
} else if let Some(org_id) = user.organization_id {
("JOIN projects p ON c.project_id = p.id
WHERE p.organization_id = ?".into(), vec![org_id.to_string()])
} else {
("WHERE 1=0".into(), vec![])
};
// 已绑定/未绑定过滤
if let Some(bound) = params.bound {
match bound.as_str() {
"1" => {
if where_clause.is_empty() {
where_clause = "WHERE c.abstract_id IS NOT NULL".into();
} else {
where_clause.push_str(" AND c.abstract_id IS NOT NULL");
}
}
"0" => {
if where_clause.is_empty() {
where_clause = "WHERE c.abstract_id IS NULL".into();
} else {
where_clause.push_str(" AND c.abstract_id IS NULL");
}
}
_ => {}
}
}
// 总数量
let count_sql = format!("SELECT COUNT(*) FROM cabinets c {}", where_clause);
let mut count_query = sqlx::query_scalar::<_, i64>(&count_sql);
for b in &binds {
count_query = count_query.bind(b);
}
let total: i64 = count_query.fetch_one(db).await.unwrap_or(0);
if total == 0 || offset >= total {
return Ok(Json(json!({"total": 0, "data": []})));
}
// 分页数据
let data_sql = format!("{} {} ORDER BY c.id LIMIT ? OFFSET ?", base_sql, where_clause);
let mut data_query = sqlx::query_as::<_, CabinetSqlRow>(&data_sql);
for b in &binds {
data_query = data_query.bind(b);
}
let sql_rows: Vec<CabinetSqlRow> = data_query
.bind(page_size)
.bind(offset)
.fetch_all(db).await?;
// 转换为CabinetRow
let mut rows: Vec<CabinetRow> = sql_rows.into_iter().map(CabinetRow::from).collect();
// 从Redis获取实时状态Pipeline批量HGET1次网络往返
{
let client = redis::Client::open(state.redis_url.as_str()).map_err(|e| {
AppError::Internal(format!("Redis客户端创建失败: {}", e))
})?;
let mut conn = client.get_async_connection().await.map_err(|e| {
AppError::Internal(format!("Redis连接失败: {}", e))
})?;
// 用Pipeline一次发所有HGET命令每个设备1次HGET取4个字段
let mut pipe = redis::pipe();
for row in &rows {
let key = format!("device:{}", row.imei);
pipe.cmd("HGET").arg(&key).arg("online");
pipe.cmd("HGET").arg(&key).arg("rssi");
pipe.cmd("HGET").arg(&key).arg("pow_fail_dc");
pipe.cmd("HGET").arg(&key).arg("pow_fail_ac");
}
// 一次发送,获取所有结果
let results: Vec<Option<String>> = pipe.query_async(&mut conn).await.unwrap_or_default();
// 每4个一组解析
for (i, row) in rows.iter_mut().enumerate() {
let base = i * 4;
row.is_online = results.get(base).and_then(|v| v.as_ref()).map(|v| v == "1");
row.rssi = results.get(base + 1).and_then(|v| v.as_ref()).and_then(|v| v.parse::<i32>().ok());
row.pow_fail_dc = results.get(base + 2).and_then(|v| v.as_ref()).map(|v| v == "1");
row.pow_fail_ac = results.get(base + 3).and_then(|v| v.as_ref()).map(|v| v == "1");
}
}
Ok(Json(json!({"total": total, "data": rows})))
}
/// 设备实时状态参数
#[derive(Deserialize)]
pub struct DeviceStatusParams {
pub imei: String,
}
/// 设备实时状态接口 — 返回板级+通道级状态
pub async fn device_realtime_status(
user: CurrentUser,
State(state): State<AppState>,
Query(params): Query<DeviceStatusParams>,
) -> Result<Json<Value>, AppError> {
auth::check_permission(&user, "device:view")?;
let client = redis::Client::open(state.redis_url.as_str()).map_err(|e| {
AppError::Internal(format!("Redis客户端创建失败: {}", e))
})?;
let mut conn = client.get_async_connection().await.map_err(|e| {
AppError::Internal(format!("Redis连接失败: {}", e))
})?;
let device_key = format!("device:{}", params.imei);
// 获取设备级字段
let fields = vec!["online", "rssi", "pow_fail_dc", "pow_fail_ac"];
let values: Vec<Option<String>> = redis::cmd("HMGET")
.arg(&device_key)
.arg(&fields)
.query_async(&mut conn)
.await
.unwrap_or_default();
let device_status = json!({
"online": values.get(0).and_then(|v| v.as_ref()).map(|v| v == "1").unwrap_or(false),
"rssi": values.get(1).and_then(|v| v.as_ref()).and_then(|v| v.parse::<i32>().ok()),
"pow_fail_dc": values.get(2).and_then(|v| v.as_ref()).map(|v| v == "1").unwrap_or(false),
"pow_fail_ac": values.get(3).and_then(|v| v.as_ref()).map(|v| v == "1").unwrap_or(false),
});
// 扫描板级key: device:{imei}:board:*
let board_pattern = format!("device:{}:board:*", params.imei);
let board_keys: Vec<String> = redis::cmd("KEYS")
.arg(&board_pattern)
.query_async(&mut conn)
.await
.unwrap_or_default();
let mut boards = json!({});
for board_key in &board_keys {
// 提取board_idx
let board_idx = board_key.split(":").last().unwrap_or("");
if board_idx.is_empty() {
continue;
}
// 获取板级字段
let board_fields = vec!["status_hex"];
let board_values: Vec<Option<String>> = redis::cmd("HMGET")
.arg(board_key)
.arg(&board_fields)
.query_async(&mut conn)
.await
.unwrap_or_default();
let mut board_data = json!({
"status_hex": board_values.get(0).and_then(|v| v.as_ref()).cloned().unwrap_or_default(),
});
// 扫描通道级key: device:{imei}:board:{idx}:ch:*
let ch_pattern = format!("{}:ch:*", board_key);
let ch_keys: Vec<String> = redis::cmd("KEYS")
.arg(&ch_pattern)
.query_async(&mut conn)
.await
.unwrap_or_default();
let mut channels = json!({});
for ch_key in &ch_keys {
let ch_idx = ch_key.split(":").last().unwrap_or("");
if ch_idx.is_empty() {
continue;
}
let ch_fields = vec!["on", "full", "fault", "hex"];
let ch_values: Vec<Option<String>> = redis::cmd("HMGET")
.arg(ch_key)
.arg(&ch_fields)
.query_async(&mut conn)
.await
.unwrap_or_default();
// 状态映射fault > full > on > idle
let fault = ch_values.get(0).and_then(|v| v.as_ref()).map(|v| v == "1").unwrap_or(false);
let full = ch_values.get(1).and_then(|v| v.as_ref()).map(|v| v == "1").unwrap_or(false);
let on = ch_values.get(2).and_then(|v| v.as_ref()).map(|v| v == "1").unwrap_or(false);
let status = if fault {
"fault"
} else if full {
"full"
} else if on {
"charging"
} else {
"idle"
};
channels[ch_idx] = json!({
"status": status,
"on": on,
"full": full,
"fault": fault,
"hex": ch_values.get(3).and_then(|v| v.as_ref()).cloned().unwrap_or_default(),
});
}
board_data["channels"] = channels;
boards[board_idx] = board_data;
}
Ok(Json(json!({
"device": device_status,
"boards": boards,
})))
}