Skip to content

Commit 26f2114

Browse files
authored
fix(ostool-server): retry failed session releases (#162)
1 parent 3a09db0 commit 26f2114

2 files changed

Lines changed: 196 additions & 67 deletions

File tree

docs/api.md

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -282,7 +282,12 @@ DELETE /api/v1/sessions/{session_id}
282282
}
283283
```
284284

285-
删除成功时服务端返回 `202 Accepted` 且没有响应体;删除时返回 `404`,客户端也将其视为已释放。
285+
删除成功时服务端返回 `202 Accepted` 且没有响应体;这只表示服务端已经接受异步释放请求。
286+
服务端会等待会话任务退出、关闭开发板电源并清理会话文件,全部完成后才删除会话并把开发板
287+
恢复为 `idle`。若释放步骤因串口或继电器短暂不可用而失败,会话会保持 `releasing`
288+
开发板 runtime 保留对应的 `active_session_id``last_release_error`,服务端每 5 秒重试
289+
尚未完成的步骤;重试成功前该开发板不会重新分配。删除不存在的会话时返回 `404`,客户端也
290+
将其视为已释放。
286291

287292
### 获取启动配置
288293

ostool-server/src/state.rs

Lines changed: 190 additions & 66 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,10 @@ const RELEASE_RETRY_ATTEMPTS: usize = 3;
2424
const RELEASE_RETRY_DELAY: Duration = Duration::from_millis(200);
2525
const RELEASE_WAIT_TIMEOUT: Duration = Duration::from_secs(2);
2626
const RELEASE_COMPLETION_TIMEOUT: Duration = Duration::from_secs(12);
27+
#[cfg(not(test))]
28+
const RELEASE_RECOVERY_RETRY_DELAY: Duration = Duration::from_secs(5);
29+
#[cfg(test)]
30+
const RELEASE_RECOVERY_RETRY_DELAY: Duration = Duration::from_millis(100);
2731
const ZHONGSHENG_RELEASE_SETTLE_DELAY: Duration = Duration::from_secs(10);
2832

2933
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
@@ -386,7 +390,11 @@ impl AppState {
386390
Ok(())
387391
}
388392

389-
pub async fn mark_board_idle(&self, board_id: &str, session_id: &str) -> anyhow::Result<()> {
393+
pub async fn complete_session_release(
394+
&self,
395+
board_id: &str,
396+
session_id: &str,
397+
) -> anyhow::Result<()> {
390398
let mut runtimes = self.board_runtimes.write().await;
391399
let runtime = runtimes
392400
.get_mut(board_id)
@@ -396,10 +404,16 @@ impl AppState {
396404
anyhow::bail!("board `{board_id}` is no longer associated with session `{session_id}`");
397405
}
398406

407+
let mut sessions = self.sessions.write().await;
408+
if !sessions.contains_key(session_id) {
409+
anyhow::bail!("session `{session_id}` disappeared before board release completed");
410+
}
411+
399412
runtime.lease_state = BoardLeaseState::Idle;
400413
runtime.active_session_id = None;
401414
runtime.last_release_error = None;
402415
runtime.updated_at = Utc::now();
416+
sessions.remove(session_id);
403417
Ok(())
404418
}
405419

@@ -419,7 +433,6 @@ impl AppState {
419433
}
420434

421435
runtime.lease_state = BoardLeaseState::Error;
422-
runtime.active_session_id = None;
423436
runtime.last_release_error = Some(error);
424437
runtime.updated_at = Utc::now();
425438
Ok(())
@@ -435,10 +448,6 @@ impl AppState {
435448
.map_err(|_| anyhow::anyhow!("release coordinator is not running"))
436449
}
437450

438-
pub async fn remove_session_runtime(&self, session_id: &str) {
439-
self.sessions.write().await.remove(session_id);
440-
}
441-
442451
async fn wait_for_session_removed(
443452
&self,
444453
session_id: &str,
@@ -469,64 +478,86 @@ impl AppState {
469478
job.reason
470479
);
471480

472-
let mut errors = Vec::new();
481+
let mut session_tasks_stopped = false;
482+
let mut board_powered_off = false;
483+
let mut release_settled = false;
484+
let mut tftp_cleaned = false;
473485

474-
if let Err(err) = self
475-
.wait_for_session_tasks_to_stop(&session, RELEASE_WAIT_TIMEOUT)
476-
.await
477-
{
478-
errors.push(err);
479-
}
486+
loop {
487+
let mut errors = Vec::new();
480488

481-
if let Err(err) = retry_release_step(RELEASE_RETRY_ATTEMPTS, RELEASE_RETRY_DELAY, || {
482-
let state = self.clone();
483-
let board = board.clone();
484-
async move {
485-
state
486-
.execute_board_power_action(&board, PowerAction::Off)
489+
if !session_tasks_stopped {
490+
match self
491+
.wait_for_session_tasks_to_stop(&session, RELEASE_WAIT_TIMEOUT)
487492
.await
488-
.map(|_| ())
489-
.map_err(|err| err.to_string())
493+
{
494+
Ok(()) => session_tasks_stopped = true,
495+
Err(err) => errors.push(err),
496+
}
490497
}
491-
})
492-
.await
493-
{
494-
errors.push(format!("power-off failed: {err}"));
495-
}
496498

497-
if errors.is_empty()
498-
&& let Some(delay) = release_settle_delay(&board)
499-
{
500-
tokio::time::sleep(delay).await;
501-
}
499+
if !board_powered_off {
500+
match retry_release_step(RELEASE_RETRY_ATTEMPTS, RELEASE_RETRY_DELAY, || {
501+
let state = self.clone();
502+
let board = board.clone();
503+
async move {
504+
state
505+
.execute_board_power_action(&board, PowerAction::Off)
506+
.await
507+
.map(|_| ())
508+
.map_err(|err| err.to_string())
509+
}
510+
})
511+
.await
512+
{
513+
Ok(()) => board_powered_off = true,
514+
Err(err) => errors.push(format!("power-off failed: {err}")),
515+
}
516+
}
502517

503-
if let Err(err) = retry_release_step(RELEASE_RETRY_ATTEMPTS, RELEASE_RETRY_DELAY, || {
504-
let manager = self.tftp_manager.clone();
505-
let session_id = snapshot.id.clone();
506-
async move {
507-
manager
508-
.read()
509-
.await
510-
.clone()
511-
.remove_session_dir(&session_id)
512-
.await
513-
.map_err(|err| err.to_string())
518+
if !release_settled && session_tasks_stopped && board_powered_off {
519+
if let Some(delay) = release_settle_delay(&board) {
520+
tokio::time::sleep(delay).await;
521+
}
522+
release_settled = true;
514523
}
515-
})
516-
.await
517-
{
518-
errors.push(format!("tftp cleanup failed: {err}"));
519-
}
520524

521-
if errors.is_empty() {
522-
if let Err(err) = self.mark_board_idle(&snapshot.board_id, &snapshot.id).await {
523-
log::warn!(
524-
"failed to mark board `{}` idle after releasing session `{}`: {err:#}",
525-
snapshot.board_id,
526-
snapshot.id
527-
);
525+
if !tftp_cleaned {
526+
match retry_release_step(RELEASE_RETRY_ATTEMPTS, RELEASE_RETRY_DELAY, || {
527+
let manager = self.tftp_manager.clone();
528+
let session_id = snapshot.id.clone();
529+
async move {
530+
manager
531+
.read()
532+
.await
533+
.clone()
534+
.remove_session_dir(&session_id)
535+
.await
536+
.map_err(|err| err.to_string())
537+
}
538+
})
539+
.await
540+
{
541+
Ok(()) => tftp_cleaned = true,
542+
Err(err) => errors.push(format!("tftp cleanup failed: {err}")),
543+
}
528544
}
529-
} else {
545+
546+
if errors.is_empty()
547+
&& session_tasks_stopped
548+
&& board_powered_off
549+
&& release_settled
550+
&& tftp_cleaned
551+
{
552+
match self
553+
.complete_session_release(&snapshot.board_id, &snapshot.id)
554+
.await
555+
{
556+
Ok(()) => return,
557+
Err(err) => errors.push(format!("failed to mark board idle: {err:#}")),
558+
}
559+
}
560+
530561
let message = errors.join("; ");
531562
if let Err(err) = self
532563
.mark_board_error(&snapshot.board_id, &snapshot.id, message.clone())
@@ -539,13 +570,13 @@ impl AppState {
539570
);
540571
}
541572
log::warn!(
542-
"release for session `{}` completed with errors: {}",
573+
"release for session `{}` failed and will retry in {:?}: {}",
543574
snapshot.id,
575+
RELEASE_RECOVERY_RETRY_DELAY,
544576
message
545577
);
578+
tokio::time::sleep(RELEASE_RECOVERY_RETRY_DELAY).await;
546579
}
547-
548-
self.remove_session_runtime(&snapshot.id).await;
549580
}
550581

551582
async fn wait_for_session_tasks_to_stop(
@@ -601,7 +632,15 @@ pub enum TouchSessionError {
601632

602633
#[cfg(test)]
603634
mod tests {
604-
use std::{fs, path::Path, sync::Arc, time::Duration};
635+
use std::{
636+
fs,
637+
path::Path,
638+
sync::{
639+
Arc,
640+
atomic::{AtomicUsize, Ordering},
641+
},
642+
time::Duration,
643+
};
605644

606645
use async_trait::async_trait;
607646
#[cfg(unix)]
@@ -616,7 +655,7 @@ mod tests {
616655
};
617656

618657
use super::{
619-
BoardLeaseState, RELEASE_COMPLETION_TIMEOUT, TouchSessionError,
658+
BoardLeaseState, RELEASE_COMPLETION_TIMEOUT, RELEASE_RETRY_ATTEMPTS, TouchSessionError,
620659
ZHONGSHENG_RELEASE_SETTLE_DELAY, build_app_state, release_settle_delay,
621660
};
622661
use crate::{
@@ -733,6 +772,7 @@ mod tests {
733772

734773
struct FailingRemoveTftpManager {
735774
root_dir: std::path::PathBuf,
775+
failures_remaining: Option<AtomicUsize>,
736776
}
737777

738778
#[async_trait]
@@ -790,6 +830,17 @@ mod tests {
790830
}
791831

792832
async fn remove_session_dir(&self, _session_id: &str) -> anyhow::Result<()> {
833+
if let Some(failures_remaining) = &self.failures_remaining {
834+
if failures_remaining
835+
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |remaining| {
836+
remaining.checked_sub(1)
837+
})
838+
.is_err()
839+
{
840+
return Ok(());
841+
}
842+
anyhow::bail!("simulated transient TFTP cleanup failure");
843+
}
793844
anyhow::bail!("simulated TFTP cleanup failure")
794845
}
795846

@@ -892,7 +943,7 @@ mod tests {
892943
}
893944

894945
#[tokio::test]
895-
async fn remove_session_marks_board_error_when_tftp_cleanup_fails() {
946+
async fn failed_release_retains_session_for_retry() {
896947
let temp = tempdir().unwrap();
897948
let root = temp.path().to_path_buf();
898949
let power_log = root.join("power.log");
@@ -911,6 +962,7 @@ mod tests {
911962
let state = build_app_state(config_path, config, manager).await.unwrap();
912963
*state.tftp_manager.write().await = Arc::new(FailingRemoveTftpManager {
913964
root_dir: root.join("tftp"),
965+
failures_remaining: None,
914966
});
915967

916968
let board = BoardConfig {
@@ -941,12 +993,26 @@ mod tests {
941993
.await
942994
.insert("session-1".into(), session);
943995

944-
let removed = state.remove_session("session-1").await.unwrap();
945-
assert!(removed.is_some());
996+
state
997+
.request_session_stop("session-1", SessionStopReason::ApiDelete)
998+
.await
999+
.unwrap();
1000+
1001+
let runtime = tokio::time::timeout(Duration::from_secs(2), async {
1002+
loop {
1003+
let runtime = state.board_runtime_status("board-1").await.unwrap();
1004+
if runtime.lease_state == BoardLeaseState::Error {
1005+
break runtime;
1006+
}
1007+
tokio::time::sleep(Duration::from_millis(25)).await;
1008+
}
1009+
})
1010+
.await
1011+
.unwrap();
1012+
9461013
assert_eq!(fs::read_to_string(power_log).unwrap(), "off");
947-
assert!(!state.sessions.read().await.contains_key("session-1"));
948-
let runtime = state.board_runtime_status("board-1").await.unwrap();
949-
assert_eq!(runtime.lease_state, BoardLeaseState::Error);
1014+
assert!(state.sessions.read().await.contains_key("session-1"));
1015+
assert_eq!(runtime.active_session_id.as_deref(), Some("session-1"));
9501016
assert!(
9511017
runtime
9521018
.last_release_error
@@ -955,6 +1021,64 @@ mod tests {
9551021
);
9561022
}
9571023

1024+
#[tokio::test]
1025+
async fn release_recovers_after_transient_tftp_cleanup_failure() {
1026+
let temp = tempdir().unwrap();
1027+
let root = temp.path().to_path_buf();
1028+
let state = test_state(&root).await;
1029+
*state.tftp_manager.write().await = Arc::new(FailingRemoveTftpManager {
1030+
root_dir: root.join("tftp"),
1031+
failures_remaining: Some(AtomicUsize::new(RELEASE_RETRY_ATTEMPTS)),
1032+
});
1033+
1034+
let board = BoardConfig {
1035+
id: "board-1".into(),
1036+
board_type: "demo".into(),
1037+
tags: vec![],
1038+
serial: None,
1039+
power_management: PowerManagementConfig::Custom(CustomPowerManagement {
1040+
power_on_cmd: "printf on >/dev/null".into(),
1041+
power_off_cmd: "printf off >/dev/null".into(),
1042+
}),
1043+
boot: BootConfig::Pxe(PxeProfile::default()),
1044+
notes: None,
1045+
disabled: false,
1046+
};
1047+
state
1048+
.boards
1049+
.write()
1050+
.await
1051+
.insert(board.id.clone(), board.clone());
1052+
state.sync_board_runtime_states().await;
1053+
let session =
1054+
SessionState::new_with_actor("session-1".into(), board.clone(), None, state.clone());
1055+
state.claim_board_for_session(&board.id, "session-1").await;
1056+
state
1057+
.sessions
1058+
.write()
1059+
.await
1060+
.insert("session-1".into(), session);
1061+
1062+
state
1063+
.request_session_stop("session-1", SessionStopReason::ApiDelete)
1064+
.await
1065+
.unwrap();
1066+
1067+
tokio::time::timeout(Duration::from_secs(5), async {
1068+
loop {
1069+
let runtime = state.board_runtime_status("board-1").await.unwrap();
1070+
if runtime.lease_state == BoardLeaseState::Idle
1071+
&& state.get_session("session-1").await.is_none()
1072+
{
1073+
break;
1074+
}
1075+
tokio::time::sleep(Duration::from_millis(25)).await;
1076+
}
1077+
})
1078+
.await
1079+
.unwrap();
1080+
}
1081+
9581082
#[tokio::test]
9591083
async fn duplicate_stop_requests_only_power_off_once() {
9601084
let temp = tempdir().unwrap();

0 commit comments

Comments
 (0)