diff --git a/README.en.md b/README.en.md index 6c4914c..4ceb92e 100644 --- a/README.en.md +++ b/README.en.md @@ -144,6 +144,9 @@ ostool board config # List remote board types ostool board ls +# Connect a specific remote board +ostool board connect -b OrangePi-5-Plus --board-id OrangePi-5-Plus-2 + # Run on a remote board ostool board run @@ -352,6 +355,8 @@ ostool board config `server` should be a complete URL including `http://` or `https://`; the optional `port` overrides the URL port. For legacy LAN configurations, a bare IPv4 or IPv6 address is interpreted as `http://`. The base release's persisted `server_ip` / `port` pair is also migrated to `server` / `port` when read; the next configuration save writes only the new format. Bare host names are not supported. Project-local `.board.toml` `server` / `port` fields still apply to `ostool board run`, with precedence lower than CLI flags and higher than the global config. +`ostool board connect -b ` allocates any available board of that type; pass `--board-id ` as well when you need to connect a specific board. + A `.board.toml` file can declare shared files relative to its own directory with `session_files`. The caller supplies that directory through `BoardRunRequest::with_session_files`; after the board session is created, diff --git a/README.md b/README.md index f881e73..0fa6bb5 100644 --- a/README.md +++ b/README.md @@ -144,6 +144,9 @@ ostool board config # 查看远端开发板类型 ostool board ls +# 连接指定开发板 +ostool board connect -b OrangePi-5-Plus --board-id OrangePi-5-Plus-2 + # 在远端开发板上运行 ostool board run @@ -348,6 +351,8 @@ ostool board config `server` 应使用包含 `http://` 或 `https://` 的完整 URL;可选的 `port` 会覆盖 URL 中的端口。为兼容旧的局域网配置,裸 IPv4 或 IPv6 地址会自动补为 `http://`。基线版本写出的 `server_ip` / `port` 也会在读取时迁移为 `server` / `port`,下一次保存配置时只写新格式;无 scheme 的主机名不支持。项目级 `.board.toml` 中的 `server` / `port` 仍可用于 `ostool board run`,其优先级低于命令行参数,高于全局配置。 +`ostool board connect -b ` 会按类型分配任意空闲开发板;需要连接某一块具体开发板时,可以额外传入 `--board-id `。 + `.board.toml` 可以用 `session_files` 声明相对于配置文件目录的共享文件。调用方通过 `BoardRunRequest::with_session_files` 提供该目录,ostool 会在 board session 建立后按原相对路径上传,并在每个 `shell_check_steps` 的 `shell_cmd` 中展开 diff --git a/docs/api.md b/docs/api.md index 921b3ca..5820c98 100644 --- a/docs/api.md +++ b/docs/api.md @@ -41,7 +41,7 @@ HTTP 客户端不跟随重定向。认证模式下,绝对 WebSocket URL 必须 | `ostool auth status [--server URL] [--port PORT]` | 显示当前 endpoint 的认证状态。 | 显示凭据类型、已知过期时间和 scope,不显示 Token。 | 无网络 API | | `ostool logout [--server URL] [--port PORT]` | 退出登录。 | OAuth 凭据会尝试远端撤销;随后删除本地凭据。PAT 仅删除本地副本。 | `POST /oauth/revoke`(仅 OAuth) | | `ostool board ls [--server URL] [--port PORT]` | 查询按类型聚合的可用开发板信息。 | 调用时携带 Bearer Token。 | `GET /api/v1/board-types` | -| `ostool board connect --board-type TYPE [--server URL] [--port PORT]` | 请求服务端从指定类型中自动分配一块开发板并打开串口终端。 | REST 和 WebSocket 请求均携带 Bearer Token。 | `POST /api/v1/sessions`;`POST /api/v1/sessions/{session_id}/heartbeat`;WebSocket `/api/v1/sessions/{session_id}/serial/ws`;`DELETE /api/v1/sessions/{session_id}` | +| `ostool board connect --board-type TYPE [--board-id BOARD_ID] [--server URL] [--port PORT]` | 请求服务端从指定类型中自动分配一块开发板,或指定一块开发板,并打开串口终端。 | REST 和 WebSocket 请求均携带 Bearer Token。 | `POST /api/v1/sessions`;`POST /api/v1/sessions/{session_id}/heartbeat`;WebSocket `/api/v1/sessions/{session_id}/serial/ws`;`DELETE /api/v1/sessions/{session_id}` | | `ostool board run [--server URL] [--port PORT]` | 构建后请求服务端按 `.board.toml` 的 `board_type` 自动分配开发板并启动。 | REST 和 WebSocket 请求均携带 Bearer Token。 | 始终:`POST /api/v1/sessions`、`POST /api/v1/sessions/{session_id}/heartbeat`、`DELETE /api/v1/sessions/{session_id}`。U-Boot:`GET /boot-profile`、`GET /serial`、`GET /tftp`、`GET /dtb`、`GET /dtb/download`、`PUT /files`、WebSocket `/serial/ws`。HTTP Boot:`GET /boot-profile`、`GET /serial`、`PUT /http-boot/kernel`、WebSocket `/serial/ws`。 | ## OAuth Device Authorization API @@ -627,12 +627,15 @@ Content-Type: application/json { "board_type": "rk3568", + "board_id": "rk3568-02", "required_tags": [], "client_name": "ostool" } ``` -用于 `ostool board connect` 和 `ostool board run`。`board_type` 必填,`required_tags` 由当前 CLI 固定发送空数组,`client_name` 固定为 `ostool`。服务端在满足类型和标签条件的空闲开发板中自动分配,不支持通过当前 `ostool` 指定 `board_id`。 +用于 `ostool board connect` 和 `ostool board run`。`board_type` 必填,`board_id` 可选,`required_tags` 由当前 CLI 固定发送空数组,`client_name` 固定为 `ostool`。未提供 `board_id` 时,服务端在满足类型和标签条件的空闲开发板中自动分配;提供时只分配该 ID 对应的、类型匹配、未禁用且空闲的开发板。CLI 通过 `ostool board connect -b --board-id ` 使用此能力;`--board-type` 仍是必填参数。 + +服务端会对 `board_id` 做首尾空白清理,清理后为空的值返回 `400`。不存在的开发板返回 `404`,开发板类型不匹配返回 `400`,指定开发板不可用时返回 `409`。客户端会将 CLI 参数规范化后发送,并在创建成功后校验响应中的 `board_id`;如果服务端返回了另一块开发板,客户端会尽力使用 `DELETE /api/v1/sessions/{session_id}` 释放刚创建的会话,然后报错,不会静默连接错误设备。 成功时返回 `201 Created`: @@ -653,7 +656,9 @@ Content-Type: application/json 当前 `ostool-server` 的固定会话 TTL 为 10 秒,每次心跳会把到期时间更新为服务端当前时间之后 10 秒;`ostool` 在成功创建会话后每秒发送一次心跳。独立认证后端可以采用不同 TTL,但必须返回真实的 `lease_expires_at` 并在心跳时续租。 -指定类型不存在时返回 `404`;类型存在但没有符合条件的空闲开发板时返回 `409`。只有结构化错误中的 `code` 恰好为 `conflict`,且 `message` 与服务端生成的 `no available board for type …` 完全匹配时,当前客户端才会每秒重试;其他 `409` 会直接返回给调用者。 +指定类型不存在时返回 `404`;类型存在但没有符合条件的空闲开发板时返回 `409`。未指定 `board_id` 时,只有结构化错误中的 `code` 恰好为 `conflict`,且 `message` 与服务端生成的 `no available board for type …` 完全匹配,客户端才会每秒重试。指定 `board_id` 且该板卡忙时,服务端返回与 `board is not available` 对应的结构化错误,客户端同样每秒重试;其他错误直接返回给调用者。 + +`board_id` 是向后兼容的可选 JSON 字段。旧客户端不发送该字段,仍按类型自动分配。新客户端连接会忽略未知 JSON 字段但仍按类型分配的旧服务端时,会通过成功响应中的 `board_id` 检测到服务端未执行指定分配,尽力释放该会话并报错;因此不会因为旧服务端静默忽略字段而连接到另一块真实设备。若旧服务端拒绝未知字段,则直接返回其错误。 ### 查询会话详情 diff --git a/ostool-server/src/api/models.rs b/ostool-server/src/api/models.rs index 2ac91d1..5e6d697 100644 --- a/ostool-server/src/api/models.rs +++ b/ostool-server/src/api/models.rs @@ -38,6 +38,16 @@ pub struct CreateSessionRequest { pub client_name: Option, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CreateSessionRequestWithBoardId { + pub board_type: String, + #[serde(default)] + pub board_id: Option, + #[serde(default)] + pub required_tags: Vec, + pub client_name: Option, +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub struct SessionCreatedResponse { pub session_id: String, diff --git a/ostool-server/src/api/router.rs b/ostool-server/src/api/router.rs index 4d8c238..60981a1 100644 --- a/ostool-server/src/api/router.rs +++ b/ostool-server/src/api/router.rs @@ -30,7 +30,7 @@ use crate::{ AdminServerConfigEditable, AdminServerConfigReadonly, AdminServerConfigResponse, AdminSessionsResponse, AdminTftpConfigResponse, AdminTftpStatusResponse, BoardPowerAction, BoardPowerStatusResponse, BoardRuntimeStatusResponse, - BoardTypeSummary, BootProfileResponse, CreateSessionRequest, + BoardTypeSummary, BootProfileResponse, CreateSessionRequestWithBoardId, CreateVirtualDeviceRequest, DtbFileResponse, HeartbeatResponse, HttpBootFileResponse, KernelPublishResponse, LoaderDeviceSummary, NetworkInterfaceSummary, SerialPortSummary, SerialStatusResponse, SessionCreatedResponse, SessionDetailResponse, @@ -38,7 +38,7 @@ use crate::{ UpdateServerConfigRequest, VirtualDeviceSummary, VirtualDevicesResponse, }, }, - board_pool::BoardAllocationStatus, + board_pool::{BoardAllocationStatus, BoardAllocationWithIdStatus}, config::{ BoardConfig, BootConfig, PowerManagementConfig, ServerConfig, TftpConfig, UbootNetworkMode, }, @@ -1288,28 +1288,62 @@ async fn list_board_types( async fn create_session( State(state): State, - axum::Json(request): axum::Json, + axum::Json(request): axum::Json, ) -> Result<(StatusCode, axum::Json), ApiError> { if request.board_type.trim().is_empty() { return Err(ApiError::bad_request("board_type must not be empty")); } + let board_id = request.board_id.as_deref().map(str::trim); + if board_id == Some("") { + return Err(ApiError::bad_request( + "board_id must not be empty when provided", + )); + } - let session = state - .create_session( - &request.board_type, - &request.required_tags, - request.client_name.clone(), - ) - .await - .map_err(|err| match err { - BoardAllocationStatus::BoardTypeNotFound => { - ApiError::not_found(format!("board type `{}` not found", request.board_type)) - } - BoardAllocationStatus::NoAvailableBoard => ApiError::conflict(format!( - "no available board for type `{}`", - request.board_type - )), - })?; + let session = match board_id { + Some(board_id) => state + .create_session_with_board_id( + &request.board_type, + board_id, + &request.required_tags, + request.client_name.clone(), + ) + .await + .map_err(|err| match err { + BoardAllocationWithIdStatus::BoardTypeNotFound => { + ApiError::not_found(format!("board type `{}` not found", request.board_type)) + } + BoardAllocationWithIdStatus::BoardNotFound => { + ApiError::not_found(format!("board `{board_id}` not found")) + } + BoardAllocationWithIdStatus::BoardTypeMismatch { + board_id, + board_type, + actual_board_type, + } => ApiError::bad_request(format!( + "board `{board_id}` has type `{actual_board_type}`, not `{board_type}`" + )), + BoardAllocationWithIdStatus::NoAvailableBoard => { + ApiError::conflict(format!("board `{board_id}` is not available")) + } + })?, + None => state + .create_session( + &request.board_type, + &request.required_tags, + request.client_name.clone(), + ) + .await + .map_err(|err| match err { + BoardAllocationStatus::BoardTypeNotFound => { + ApiError::not_found(format!("board type `{}` not found", request.board_type)) + } + BoardAllocationStatus::NoAvailableBoard => ApiError::conflict(format!( + "no available board for type `{}`", + request.board_type + )), + })?, + }; let board = state .session_board(&session.id) @@ -5077,6 +5111,76 @@ mod tests { assert_eq!(value["serial_available"], true); } + #[tokio::test] + async fn create_session_allocates_requested_board_id() { + let app = test_router().await; + for board_id in ["demo-01", "demo-02"] { + let mut board = sample_board(board_id); + board.board_type = "demo".to_string(); + assert_eq!( + create_board(&app, serde_json::to_value(board).unwrap()).await, + StatusCode::CREATED + ); + } + + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/sessions") + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from( + json!({ + "board_type": "demo", + "board_id": " demo-02 ", + "required_tags": [], + "client_name": "test", + }) + .to_string(), + )) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::CREATED); + let body = to_bytes(response.into_body(), usize::MAX).await.unwrap(); + let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); + assert_eq!(value["board_id"], "demo-02"); + } + + #[tokio::test] + async fn create_session_rejects_missing_requested_board_id() { + let app = test_router().await; + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/sessions") + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from( + json!({ + "board_type": "demo", + "board_id": "missing-demo-01", + "required_tags": [], + "client_name": "test", + }) + .to_string(), + )) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NOT_FOUND); + let body = to_bytes(response.into_body(), usize::MAX).await.unwrap(); + let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); + assert_eq!(value["code"], "not_found"); + assert_eq!(value["message"], "board `missing-demo-01` not found"); + } + #[tokio::test] async fn create_session_returns_conflict_without_waiting_when_pool_is_busy() { let app = test_router().await; diff --git a/ostool-server/src/board_pool.rs b/ostool-server/src/board_pool.rs index 2af87af..149c0a7 100644 --- a/ostool-server/src/board_pool.rs +++ b/ostool-server/src/board_pool.rs @@ -8,12 +8,84 @@ pub enum BoardAllocationStatus { NoAvailableBoard, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum BoardAllocationWithIdStatus { + BoardTypeNotFound, + BoardNotFound, + BoardTypeMismatch { + board_id: String, + board_type: String, + actual_board_type: String, + }, + NoAvailableBoard, +} + pub fn allocate_board( boards: &BTreeMap, leased_boards: &BTreeSet, board_type: &str, required_tags: &[String], ) -> Result { + allocate_board_impl(boards, leased_boards, board_type, None, required_tags).map_err(|status| { + match status { + BoardAllocationWithIdStatus::BoardTypeNotFound => { + BoardAllocationStatus::BoardTypeNotFound + } + BoardAllocationWithIdStatus::NoAvailableBoard => { + BoardAllocationStatus::NoAvailableBoard + } + BoardAllocationWithIdStatus::BoardNotFound + | BoardAllocationWithIdStatus::BoardTypeMismatch { .. } => { + unreachable!("board ID errors are impossible without a requested board ID") + } + } + }) +} + +pub fn allocate_board_with_board_id( + boards: &BTreeMap, + leased_boards: &BTreeSet, + board_type: &str, + board_id: &str, + required_tags: &[String], +) -> Result { + let board_id = board_id.trim(); + allocate_board_impl( + boards, + leased_boards, + board_type, + Some(board_id), + required_tags, + ) +} + +fn allocate_board_impl( + boards: &BTreeMap, + leased_boards: &BTreeSet, + board_type: &str, + board_id: Option<&str>, + required_tags: &[String], +) -> Result { + if let Some(board_id) = board_id { + let board = boards + .get(board_id) + .ok_or(BoardAllocationWithIdStatus::BoardNotFound)?; + if board.board_type != board_type { + return Err(BoardAllocationWithIdStatus::BoardTypeMismatch { + board_id: board_id.to_string(), + board_type: board_type.to_string(), + actual_board_type: board.board_type.clone(), + }); + } + if board.disabled + || leased_boards.contains(&board.id) + || !required_tags.iter().all(|tag| board.tags.contains(tag)) + { + return Err(BoardAllocationWithIdStatus::NoAvailableBoard); + } + return Ok(board.clone()); + } + let matching_boards = boards .values() .filter(|board| !board.disabled) @@ -26,9 +98,9 @@ pub fn allocate_board( .values() .any(|board| !board.disabled && board.board_type == board_type); return Err(if board_type_exists { - BoardAllocationStatus::NoAvailableBoard + BoardAllocationWithIdStatus::NoAvailableBoard } else { - BoardAllocationStatus::BoardTypeNotFound + BoardAllocationWithIdStatus::BoardTypeNotFound }); } @@ -36,5 +108,74 @@ pub fn allocate_board( .into_iter() .find(|board| !leased_boards.contains(&board.id)) .cloned() - .ok_or(BoardAllocationStatus::NoAvailableBoard) + .ok_or(BoardAllocationWithIdStatus::NoAvailableBoard) +} + +#[cfg(test)] +mod tests { + use std::collections::{BTreeMap, BTreeSet}; + + use crate::{ + board_pool::{BoardAllocationWithIdStatus, allocate_board_with_board_id}, + config::{BoardConfig, BootConfig, CustomPowerManagement, PowerManagementConfig}, + }; + + fn board(id: &str, board_type: &str) -> BoardConfig { + BoardConfig { + id: id.to_string(), + board_type: board_type.to_string(), + tags: Vec::new(), + serial: None, + power_management: PowerManagementConfig::Custom(CustomPowerManagement { + power_on_cmd: "echo on".into(), + power_off_cmd: "echo off".into(), + }), + boot: BootConfig::Uboot(Default::default()), + network_identity: None, + notes: None, + disabled: false, + } + } + + #[test] + fn allocate_board_uses_requested_board_id() { + let boards = BTreeMap::from([ + ("demo-01".to_string(), board("demo-01", "demo")), + ("demo-02".to_string(), board("demo-02", "demo")), + ]); + + let allocated = + allocate_board_with_board_id(&boards, &BTreeSet::new(), "demo", "demo-02", &[]) + .unwrap(); + + assert_eq!(allocated.id, "demo-02"); + } + + #[test] + fn allocate_board_rejects_busy_requested_board_id() { + let boards = BTreeMap::from([("demo-01".to_string(), board("demo-01", "demo"))]); + let leased_boards = BTreeSet::from(["demo-01".to_string()]); + + let err = allocate_board_with_board_id(&boards, &leased_boards, "demo", "demo-01", &[]) + .unwrap_err(); + + assert_eq!(err, BoardAllocationWithIdStatus::NoAvailableBoard); + } + + #[test] + fn allocate_board_rejects_mismatched_requested_board_type() { + let boards = BTreeMap::from([("demo-01".to_string(), board("demo-01", "other"))]); + + let err = allocate_board_with_board_id(&boards, &BTreeSet::new(), "demo", "demo-01", &[]) + .unwrap_err(); + + assert_eq!( + err, + BoardAllocationWithIdStatus::BoardTypeMismatch { + board_id: "demo-01".to_string(), + board_type: "demo".to_string(), + actual_board_type: "other".to_string(), + } + ); + } } diff --git a/ostool-server/src/state.rs b/ostool-server/src/state.rs index b388546..960733e 100644 --- a/ostool-server/src/state.rs +++ b/ostool-server/src/state.rs @@ -12,7 +12,10 @@ use serde::{Deserialize, Serialize}; use tokio::sync::{Mutex, RwLock, mpsc}; use crate::{ - board_pool::{BoardAllocationStatus, allocate_board}, + board_pool::{ + BoardAllocationStatus, BoardAllocationWithIdStatus, allocate_board, + allocate_board_with_board_id, + }, board_store::fs::FileBoardStore, config::{BoardConfig, PowerManagementConfig, ServerConfig}, dtb_store::DtbStore, @@ -187,6 +190,40 @@ impl AppState { required_tags: &[String], client_name: Option, ) -> Result { + self.create_session_inner(board_type, None, required_tags, client_name) + .await + .map_err(|status| match status { + BoardAllocationWithIdStatus::BoardTypeNotFound => { + BoardAllocationStatus::BoardTypeNotFound + } + BoardAllocationWithIdStatus::NoAvailableBoard => { + BoardAllocationStatus::NoAvailableBoard + } + BoardAllocationWithIdStatus::BoardNotFound + | BoardAllocationWithIdStatus::BoardTypeMismatch { .. } => { + unreachable!("board ID errors are impossible without a requested board ID") + } + }) + } + + pub async fn create_session_with_board_id( + &self, + board_type: &str, + board_id: &str, + required_tags: &[String], + client_name: Option, + ) -> Result { + self.create_session_inner(board_type, Some(board_id), required_tags, client_name) + .await + } + + async fn create_session_inner( + &self, + board_type: &str, + board_id: Option<&str>, + required_tags: &[String], + client_name: Option, + ) -> Result { loop { let _inventory_guard = self.board_inventory_gate.lock().await; let boards = self.boards.read().await; @@ -196,7 +233,24 @@ impl AppState { .filter(|(_, runtime)| runtime.lease_state != BoardLeaseState::Idle) .map(|(board_id, _)| board_id.clone()) .collect::>(); - let board = allocate_board(&boards, &unavailable_board_ids, board_type, required_tags)?; + let board = match board_id { + Some(board_id) => allocate_board_with_board_id( + &boards, + &unavailable_board_ids, + board_type, + board_id, + required_tags, + )?, + None => allocate_board(&boards, &unavailable_board_ids, board_type, required_tags) + .map_err(|status| match status { + BoardAllocationStatus::BoardTypeNotFound => { + BoardAllocationWithIdStatus::BoardTypeNotFound + } + BoardAllocationStatus::NoAvailableBoard => { + BoardAllocationWithIdStatus::NoAvailableBoard + } + })?, + }; drop(runtimes); drop(boards); diff --git a/ostool/src/board/client.rs b/ostool/src/board/client.rs index bda037f..0f94590 100644 --- a/ostool/src/board/client.rs +++ b/ostool/src/board/client.rs @@ -54,6 +54,15 @@ pub struct CreateSessionRequest { pub client_name: Option, } +#[derive(Debug, Clone, Serialize)] +pub struct CreateSessionRequestWithBoardId { + pub board_type: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub board_id: Option, + pub required_tags: Vec, + pub client_name: Option, +} + #[derive(Debug, Clone, Deserialize)] pub struct SessionCreatedResponse { pub session_id: String, @@ -218,6 +227,13 @@ impl BoardServerClientError { && self.message == format!("no available board for type `{board_type}`") } + pub fn is_no_available_board_id(&self, board_id: &str) -> bool { + let board_id = board_id.trim(); + self.status == StatusCode::CONFLICT + && self.code.as_deref() == Some("conflict") + && self.message == format!("board `{board_id}` is not available") + } + pub fn is_board_type_not_found_for(&self, board_type: &str) -> bool { self.status == StatusCode::NOT_FOUND && self.code.as_deref() == Some("not_found") @@ -261,19 +277,64 @@ impl BoardServerClient { pub async fn create_session( &self, board_type: &str, + ) -> Result { + self.create_session_inner(board_type, None).await + } + + pub async fn create_session_with_board_id( + &self, + board_type: &str, + board_id: &str, + ) -> Result { + let board_id = board_id.trim(); + if board_id.is_empty() { + return Err(BoardServerClientError { + status: StatusCode::BAD_REQUEST, + code: Some("bad_request".to_string()), + message: "board_id must not be empty when provided".to_string(), + }); + } + self.create_session_inner(board_type, Some(board_id)).await + } + + async fn create_session_inner( + &self, + board_type: &str, + board_id: Option<&str>, ) -> Result { let response = self .request(Method::POST, self.endpoint("/api/v1/sessions")) .await? - .json(&CreateSessionRequest { + .json(&CreateSessionRequestWithBoardId { board_type: board_type.to_string(), + board_id: board_id.map(ToOwned::to_owned), required_tags: vec![], client_name: Some("ostool".to_string()), }) .send() .await .map_err(Self::request_error)?; - self.decode_json(response).await + let session: SessionCreatedResponse = self.decode_json(response).await?; + if let Some(requested_board_id) = board_id + && session.board_id != requested_board_id + { + if let Err(error) = self.delete_session(&session.session_id).await { + log::warn!( + "failed to release mismatched board session `{}` after requesting board `{}`: {error}", + session.session_id, + requested_board_id + ); + } + return Err(BoardServerClientError { + status: StatusCode::CONFLICT, + code: Some("conflict".to_string()), + message: format!( + "server allocated board `{}` while board `{}` was requested", + session.board_id, requested_board_id + ), + }); + } + Ok(session) } pub async fn heartbeat( @@ -655,9 +716,18 @@ impl fmt::Display for BoardTypeSummary { #[cfg(test)] mod tests { + use std::sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }; + use chrono::{DateTime, Utc}; use reqwest::StatusCode; use serde::Serialize; + use tokio::{ + io::{AsyncReadExt as _, AsyncWriteExt as _}, + net::TcpStream, + }; use url::Url; use super::{BoardServerClient, BoardTypeSummary, BootConfig, parse_error_body}; @@ -801,6 +871,89 @@ mod tests { assert!(!error.is_no_available_board_for("rk3568")); } + #[tokio::test] + async fn create_session_releases_mismatched_requested_board_id() { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let saw_delete = Arc::new(AtomicBool::new(false)); + let saw_delete_server = saw_delete.clone(); + + let server = tokio::spawn(async move { + let (mut post_socket, _) = listener.accept().await.unwrap(); + let request = read_http_request(&mut post_socket).await; + assert!(request.starts_with("POST /api/v1/sessions ")); + assert!(request.contains(r#""board_id":"demo-02""#)); + write_http_response( + &mut post_socket, + "201 Created", + r#"{"session_id":"session-1","board_id":"demo-01","lease_expires_at":"2026-09-03T00:00:00Z","serial_available":false,"boot_mode":"uboot","ws_url":null}"#, + ) + .await; + + let (mut delete_socket, _) = listener.accept().await.unwrap(); + let request = read_http_request(&mut delete_socket).await; + assert!(request.starts_with("DELETE /api/v1/sessions/session-1 ")); + saw_delete_server.store(true, Ordering::SeqCst); + write_http_response(&mut delete_socket, "202 Accepted", "").await; + }); + + let endpoint = + BoardEndpoint::new(&format!("http://{address}"), None, AuthMode::Disabled).unwrap(); + let client = BoardServerClient::new_with_endpoint(endpoint).unwrap(); + let error = client + .create_session_with_board_id("demo", " demo-02 ") + .await + .unwrap_err(); + + server.await.unwrap(); + assert!(saw_delete.load(Ordering::SeqCst)); + assert_eq!(error.status, StatusCode::CONFLICT); + assert_eq!(error.code.as_deref(), Some("conflict")); + assert_eq!( + error.message, + "server allocated board `demo-01` while board `demo-02` was requested" + ); + } + + async fn read_http_request(socket: &mut TcpStream) -> String { + let mut bytes = Vec::new(); + let mut buffer = [0_u8; 1024]; + loop { + let read = socket.read(&mut buffer).await.unwrap(); + assert_ne!(read, 0, "client closed connection before request completed"); + bytes.extend_from_slice(&buffer[..read]); + if request_complete(&bytes) { + break; + } + } + String::from_utf8(bytes).unwrap() + } + + fn request_complete(bytes: &[u8]) -> bool { + let Some(header_end) = bytes.windows(4).position(|window| window == b"\r\n\r\n") else { + return false; + }; + let headers = String::from_utf8_lossy(&bytes[..header_end]); + let content_length = headers + .lines() + .find_map(|line| { + let (name, value) = line.split_once(':')?; + name.eq_ignore_ascii_case("content-length") + .then_some(value.trim()) + }) + .and_then(|value| value.parse::().ok()) + .unwrap_or(0); + bytes.len() >= header_end + 4 + content_length + } + + async fn write_http_response(socket: &mut TcpStream, status: &str, body: &str) { + let response = format!( + "HTTP/1.1 {status}\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}", + body.len() + ); + socket.write_all(response.as_bytes()).await.unwrap(); + } + #[test] fn parse_uboot_boot_profile() { let response: super::BootProfileResponse = serde_json::from_str( diff --git a/ostool/src/board/mod.rs b/ostool/src/board/mod.rs index 1b296ae..c3c7663 100644 --- a/ostool/src/board/mod.rs +++ b/ostool/src/board/mod.rs @@ -205,6 +205,21 @@ pub async fn acquire_board_session( Ok((client, session)) } +pub async fn acquire_board_session_with_board_id( + server: &str, + port: u16, + board_type: &str, + board_id: &str, +) -> anyhow::Result<(BoardServerClient, BoardSession)> { + let client = BoardServerClient::new(server, port)?; + let session = BoardSession::acquire_with_board_id(client.clone(), board_type, board_id) + .await + .with_context(|| { + format!("failed to acquire board type `{board_type}` with board id `{board_id}`") + })?; + Ok((client, session)) +} + pub async fn acquire_board_session_endpoint( endpoint: BoardEndpoint, board_type: &str, @@ -216,11 +231,36 @@ pub async fn acquire_board_session_endpoint( Ok((client, session)) } +pub async fn acquire_board_session_endpoint_with_board_id( + endpoint: BoardEndpoint, + board_type: &str, + board_id: &str, +) -> anyhow::Result<(BoardServerClient, BoardSession)> { + let client = BoardServerClient::new_with_endpoint(endpoint)?; + let session = BoardSession::acquire_with_board_id(client.clone(), board_type, board_id) + .await + .with_context(|| { + format!("failed to acquire board type `{board_type}` with board id `{board_id}`") + })?; + Ok((client, session)) +} + pub async fn connect_board(server: &str, port: u16, board_type: &str) -> anyhow::Result<()> { let (client, session) = acquire_board_session(server, port, board_type).await?; connect_allocated_board(client, session, board_type).await } +pub async fn connect_board_with_board_id( + server: &str, + port: u16, + board_type: &str, + board_id: &str, +) -> anyhow::Result<()> { + let (client, session) = + acquire_board_session_with_board_id(server, port, board_type, board_id).await?; + connect_allocated_board(client, session, board_type).await +} + pub async fn connect_board_endpoint( endpoint: BoardEndpoint, board_type: &str, @@ -229,6 +269,16 @@ pub async fn connect_board_endpoint( connect_allocated_board(client, session, board_type).await } +pub async fn connect_board_endpoint_with_board_id( + endpoint: BoardEndpoint, + board_type: &str, + board_id: &str, +) -> anyhow::Result<()> { + let (client, session) = + acquire_board_session_endpoint_with_board_id(endpoint, board_type, board_id).await?; + connect_allocated_board(client, session, board_type).await +} + async fn connect_allocated_board( client: BoardServerClient, session: BoardSession, diff --git a/ostool/src/board/session.rs b/ostool/src/board/session.rs index b5ae322..d4788c3 100644 --- a/ostool/src/board/session.rs +++ b/ostool/src/board/session.rs @@ -39,12 +39,46 @@ pub struct BoardSession { impl BoardSession { pub async fn acquire(client: BoardServerClient, board_type: &str) -> anyhow::Result { - let info = acquire_session_with( - board_type, - || client.create_session(board_type), - |duration| tokio::time::sleep(duration), - ) - .await?; + Self::acquire_inner(client, board_type, None).await + } + + pub async fn acquire_with_board_id( + client: BoardServerClient, + board_type: &str, + board_id: &str, + ) -> anyhow::Result { + let board_id = board_id.trim(); + if board_id.is_empty() { + return Err(anyhow!("board_id must not be empty when provided")); + } + Self::acquire_inner(client, board_type, Some(board_id)).await + } + + async fn acquire_inner( + client: BoardServerClient, + board_type: &str, + board_id: Option<&str>, + ) -> anyhow::Result { + let info = match board_id { + Some(board_id) => { + acquire_session_with( + board_type, + Some(board_id), + || client.create_session_with_board_id(board_type, board_id), + |duration| tokio::time::sleep(duration), + ) + .await? + } + None => { + acquire_session_with( + board_type, + None, + || client.create_session(board_type), + |duration| tokio::time::sleep(duration), + ) + .await? + } + }; let lease_expires_at = Arc::new(RwLock::new(info.lease_expires_at)); let (heartbeat_stop, heartbeat_rx) = watch::channel(false); @@ -196,6 +230,7 @@ async fn run_heartbeat_loop( async fn acquire_session_with( board_type: &str, + board_id: Option<&str>, mut create: CreateFn, mut sleep: SleepFn, ) -> Result @@ -205,9 +240,21 @@ where SleepFn: FnMut(Duration) -> SleepFut, SleepFut: Future, { + let board_id = board_id.map(str::trim); loop { match create().await { Ok(session) => return Ok(session), + Err(err) + if board_id + .map(|board_id| err.is_no_available_board_id(board_id)) + .unwrap_or(false) => + { + println!( + "Board `{}` is not available, retrying in 1s...", + board_id.unwrap() + ); + sleep(Duration::from_secs(1)).await; + } Err(err) if err.is_no_available_board_for(board_type) => { println!("No available board for type `{board_type}`, retrying in 1s..."); sleep(Duration::from_secs(1)).await; @@ -254,6 +301,14 @@ mod tests { } } + fn no_board_id_error(board_id: &str) -> BoardServerClientError { + BoardServerClientError { + status: StatusCode::CONFLICT, + code: Some("conflict".to_string()), + message: format!("board `{board_id}` is not available"), + } + } + #[test] fn shared_file_keeps_the_requested_relative_path() { let uploaded_at = Utc::now(); @@ -286,6 +341,7 @@ mod tests { let session = acquire_session_with( "rk3568", + None, { let responses = responses.clone(); move || { @@ -317,6 +373,41 @@ mod tests { ); } + #[tokio::test] + async fn acquire_session_retries_until_requested_board_id_is_available() { + let responses = Arc::new(Mutex::new(vec![ + Ok(created_session("demo-session")), + Err(no_board_id_error("demo-02")), + ])); + let sleeps = Arc::new(Mutex::new(Vec::new())); + + let session = acquire_session_with( + "rk3568", + Some(" demo-02 "), + { + let responses = responses.clone(); + move || { + let responses = responses.clone(); + async move { responses.lock().unwrap().pop().unwrap() } + } + }, + { + let sleeps = sleeps.clone(); + move |duration| { + let sleeps = sleeps.clone(); + async move { + sleeps.lock().unwrap().push(duration); + } + } + }, + ) + .await + .unwrap(); + + assert_eq!(session.session_id, "demo-session"); + assert_eq!(*sleeps.lock().unwrap(), vec![Duration::from_secs(1)]); + } + #[tokio::test] async fn acquire_session_stops_retrying_on_non_conflict_error() { let error = BoardServerClientError { @@ -327,6 +418,7 @@ mod tests { let result = acquire_session_with( "rk3568", + None, || { let error = error.clone(); async move { Err(error) } @@ -349,6 +441,7 @@ mod tests { let result = acquire_session_with( "rk3568", + None, || { let error = error.clone(); async move { Err(error) } diff --git a/ostool/src/main.rs b/ostool/src/main.rs index b79e845..e0748c6 100644 --- a/ostool/src/main.rs +++ b/ostool/src/main.rs @@ -164,6 +164,9 @@ struct BoardConnectArgs { /// Board type to allocate and connect #[arg(short = 'b', long)] board_type: String, + /// Specific board ID to allocate and connect + #[arg(long)] + board_id: Option, #[command(flatten)] server: BoardServerArgs, } @@ -278,7 +281,17 @@ async fn try_main() -> Result<()> { let global_config = board::load_board_global_config_with_notice()?; let endpoint = global_config .resolve_endpoint(args.server.server.as_deref(), args.server.port)?; - board::connect_board_endpoint(endpoint, &args.board_type).await?; + match args.board_id.as_deref() { + Some(board_id) => { + board::connect_board_endpoint_with_board_id( + endpoint, + &args.board_type, + board_id, + ) + .await? + } + None => board::connect_board_endpoint(endpoint, &args.board_type).await?, + } } BoardSubCommands::Run(args) => { let mut invocation = init_invocation(manifest.clone())?; @@ -822,6 +835,7 @@ mod tests { command: BoardSubCommands::Connect(args), }) => { assert_eq!(args.board_type, "rk3568"); + assert!(args.board_id.is_none()); assert!(args.server.server.is_none()); assert!(args.server.port.is_none()); } @@ -837,6 +851,8 @@ mod tests { "connect", "--board-type", "rk3568", + "--board-id", + "rk3568-2", "--server", "http://10.0.0.2", "--port", @@ -849,6 +865,7 @@ mod tests { command: BoardSubCommands::Connect(args), }) => { assert_eq!(args.board_type, "rk3568"); + assert_eq!(args.board_id.as_deref(), Some("rk3568-2")); assert_eq!(args.server.server.as_deref(), Some("http://10.0.0.2")); assert_eq!(args.server.port, Some(9000)); }