同步远端分支并保留序列帧去背景修复
Project CI / Backend tests (pull_request) Failing after 14s
Project CI / Repository checks (pull_request) Failing after 14s
Project CI / Frontend tests (pull_request) Successful in 2m39s
Project CI / Native shell tests (pull_request) Successful in 13m41s

This commit is contained in:
2026-08-17 20:22:20 +08:00
11 changed files with 287 additions and 27 deletions
@@ -1,7 +1,7 @@
use std::{future::Future, time::Duration};
use reqwest::{Client, StatusCode, header};
use serde::Deserialize;
use serde::{Deserialize, de::DeserializeOwned};
use url::Url;
use crate::{Config, DeploymentStatus, HealthStatus};
@@ -52,14 +52,13 @@ pub struct PreviewResult {
#[serde(alias = "id")]
pub deployment_id: Option<String>,
pub project_name: Option<String>,
#[serde(alias = "sourceBranch")]
pub branch: Option<String>,
#[serde(alias = "sourceCommit")]
pub resolved_commit: Option<String>,
pub status: Option<DeploymentStatus>,
pub health: Option<HealthStatus>,
pub phase: Option<String>,
pub health_status: Option<String>,
pub web_port: Option<u16>,
pub web_url: Option<String>,
pub message: Option<String>,
}
@@ -205,11 +204,14 @@ impl JenkinsClient {
Fut: Future<Output = ()>,
{
let build_url = loop {
let item_url = reference
let mut item_url = reference
.queue_url
.join("api/json")
.map_err(|error| error.to_string())?;
let item: QueueItem = self.get_json(item_url).await?;
item_url
.query_pairs_mut()
.append_pair("tree", "cancelled,executable[url]");
let item: QueueItem = self.get_json_with_retry(item_url).await?;
if item.cancelled.unwrap_or(false) {
return Ok(JenkinsOutcome {
success: false,
@@ -226,10 +228,13 @@ impl JenkinsClient {
on_build(build_url.as_str()).await;
let successful = loop {
let state_url = build_url
let mut state_url = build_url
.join("api/json")
.map_err(|error| error.to_string())?;
let state: BuildState = self.get_json(state_url).await?;
state_url
.query_pairs_mut()
.append_pair("tree", "building,result");
let state: BuildState = self.get_json_with_retry(state_url).await?;
if state.building {
tokio::time::sleep(self.poll_interval).await;
continue;
@@ -247,7 +252,7 @@ impl JenkinsClient {
let artifact_url = build_url
.join("artifact/preview-result.json")
.map_err(|error| error.to_string())?;
let result = self.get_json(artifact_url).await?;
let result = self.get_json_with_retry(artifact_url).await?;
Ok(JenkinsOutcome {
success: true,
cancelled: false,
@@ -267,18 +272,50 @@ impl JenkinsClient {
if !response.status().is_success() {
return Err(format!("Jenkins 状态返回 HTTP {}", response.status()));
}
response
.json()
let bytes = response
.bytes()
.await
.map_err(|_| "Jenkins 状态响应格式无效".to_string())
.map_err(|error| format!("Jenkins 状态响应正文读取失败: {error}"))?;
serde_json::from_slice(&bytes).map_err(|error| format!("Jenkins 状态响应格式无效: {error}"))
}
async fn get_json_with_retry<T: DeserializeOwned>(&self, url: Url) -> Result<T, String> {
const MAX_ATTEMPTS: usize = 10;
let mut last_error = None;
for attempt in 1..=MAX_ATTEMPTS {
match self.get_json(url.clone()).await {
Ok(value) => return Ok(value),
Err(error) => last_error = Some(error),
}
if attempt < MAX_ATTEMPTS {
tokio::time::sleep(self.poll_interval).await;
}
}
Err(last_error.unwrap_or_else(|| "Jenkins 状态请求失败".to_string()))
}
fn resolve_trusted_url(&self, value: &str) -> Result<Url, String> {
let url = Url::parse(value)
let returned = Url::parse(value)
.or_else(|_| self.job_url.join(value))
.map_err(|_| "Jenkins 返回了无效 URL".to_string())?;
self.ensure_same_origin(&url)?;
Ok(url)
if !returned.username().is_empty()
|| returned.password().is_some()
|| !returned.path().starts_with(self.root_url.path())
{
return Err("Jenkins 返回了非受信源 URL".to_string());
}
let relative_path = returned
.path()
.strip_prefix(self.root_url.path())
.ok_or_else(|| "Jenkins 返回了非受信路径".to_string())?;
let mut trusted = self
.root_url
.join(relative_path)
.map_err(|_| "Jenkins 返回了无效 URL".to_string())?;
trusted.set_query(returned.query());
trusted.set_fragment(None);
self.ensure_same_origin(&trusted)?;
Ok(trusted)
}
fn ensure_same_origin(&self, url: &Url) -> Result<(), String> {
@@ -24,6 +24,7 @@ use tower_http::{
trace::TraceLayer,
};
use tracing::{error, warn};
use url::Url;
use uuid::Uuid;
pub use config::Config;
@@ -86,6 +87,8 @@ pub struct Deployment {
pub resolved_commit: Option<String>,
pub status: DeploymentStatus,
pub health: HealthStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub web_port: Option<u16>,
#[serde(skip_serializing_if = "Option::is_none")]
pub web_url: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
@@ -385,6 +388,7 @@ async fn list_deployments(
.read()
.await
.values()
.filter(|record| record.public.status != DeploymentStatus::Stopped)
.map(|record| record.public.clone())
.collect();
deployments.sort_by(|left, right| right.created_at.cmp(&left.created_at));
@@ -494,6 +498,7 @@ async fn create_deployment(
resolved_commit: None,
status: DeploymentStatus::Queued,
health: HealthStatus::Pending,
web_port: None,
web_url: None,
jenkins_build_url: None,
created_at: now,
@@ -770,8 +775,23 @@ async fn apply_outcome(state: &AppState, id: &str, outcome: JenkinsOutcome) {
record.public.status = DeploymentStatus::Failed;
record.public.health = HealthStatus::Unknown;
record.public.message = Some("Jenkins 返回了无效 Web 地址".to_string());
} else if result
.web_port
.is_some_and(|value| !(8400..=8499).contains(&value))
|| result
.web_url
.as_deref()
.zip(result.web_port)
.is_some_and(|(url, port)| {
Url::parse(url).ok().and_then(|url| url.port()) != Some(port)
})
{
record.public.status = DeploymentStatus::Failed;
record.public.health = HealthStatus::Unknown;
record.public.message = Some("Jenkins 返回了无效 Web 端口".to_string());
} else if record.operation == Operation::Deploy
&& (result.resolved_commit.is_none()
|| result.web_port.is_none()
|| result.web_url.is_none()
|| result.phase.as_deref() != Some("RUNNING")
|| result.health_status.as_deref() != Some("HEALTHY"))
@@ -794,6 +814,9 @@ async fn apply_outcome(state: &AppState, id: &str, outcome: JenkinsOutcome) {
if let Some(value) = result.web_url {
record.public.web_url = Some(value);
}
if let Some(value) = result.web_port {
record.public.web_port = Some(value);
}
if let Some(value) = result.health {
record.public.health = value;
}
@@ -838,6 +861,7 @@ async fn apply_outcome(state: &AppState, id: &str, outcome: JenkinsOutcome) {
Operation::Uninstall => {
record.public.status = DeploymentStatus::Stopped;
record.public.health = HealthStatus::Unknown;
record.public.web_port = None;
record.public.web_url = None;
record.public.can_uninstall = false;
if record.public.message.is_none() {
@@ -877,7 +901,7 @@ fn load_deployments(config: &Config) -> Result<HashMap<String, DeploymentRecord>
return Err("预览部署状态文件 schemaVersion 不受支持".to_string());
}
let mut deployments = HashMap::new();
for record in persisted.deployments {
for mut record in persisted.deployments {
validate_deployment_id(&record.public.id)
.map_err(|_| "状态文件包含无效部署 ID".to_string())?;
let branch = validate_branch(&record.public.branch)
@@ -896,6 +920,27 @@ fn load_deployments(config: &Config) -> Result<HashMap<String, DeploymentRecord>
{
return Err("状态文件包含无效 Web 地址".to_string());
}
let url_port = record
.public
.web_url
.as_deref()
.and_then(|value| Url::parse(value).ok())
.and_then(|url| url.port());
if record
.public
.web_port
.is_some_and(|value| !(8400..=8499).contains(&value))
|| record
.public
.web_port
.zip(url_port)
.is_some_and(|(saved, parsed)| saved != parsed)
{
return Err("状态文件包含无效 Web 端口".to_string());
}
if record.public.web_port.is_none() {
record.public.web_port = url_port;
}
if deployments
.insert(record.public.id.clone(), record)
.is_some()
@@ -1,5 +1,8 @@
use std::{
sync::{Arc, Mutex},
sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
},
time::Duration,
};
@@ -25,6 +28,7 @@ const ORIGIN: &str = "http://preview.internal:8080";
#[derive(Clone, Default)]
struct MockJenkinsState {
requests: Arc<Mutex<Vec<RecordedRequest>>>,
invalid_build_responses: Arc<AtomicUsize>,
}
#[derive(Debug)]
@@ -95,12 +99,27 @@ async fn mock_trigger(
async fn mock_queue(axum::extract::Path(id): axum::extract::Path<u64>) -> Json<Value> {
Json(
json!({"cancelled": false, "executable": {"url": format!("/jenkins/job/shared/job/Genarrative-Preview-Deployer/{id}/")}}),
json!({"cancelled": false, "executable": {"url": format!("http://192.168.35.82:8080/jenkins/job/shared/job/Genarrative-Preview-Deployer/{id}/")}}),
)
}
async fn mock_build() -> Json<Value> {
Json(json!({"building": false, "result": "SUCCESS"}))
async fn mock_build(
State(state): State<MockJenkinsState>,
uri: axum::http::Uri,
) -> impl IntoResponse {
if uri.query() != Some("tree=building%2Cresult") {
return (StatusCode::BAD_REQUEST, "missing bounded tree query").into_response();
}
if state
.invalid_build_responses
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |remaining| {
remaining.checked_sub(1)
})
.is_ok()
{
return (StatusCode::OK, "Jenkins is finalizing the build").into_response();
}
Json(json!({"building": false, "result": "SUCCESS"})).into_response()
}
async fn mock_artifact(axum::extract::Path(build): axum::extract::Path<u64>) -> Json<Value> {
@@ -111,9 +130,10 @@ async fn mock_artifact(axum::extract::Path(build): axum::extract::Path<u64>) ->
"deploymentId": "preview-63d38d3da6bc9b06",
"projectName": "genarrative-preview-63d38d3da6bc9b06",
"branch": "feature/preview-ui",
"sourceCommit": "0123456789abcdef0123456789abcdef01234567",
"resolvedCommit": "0123456789abcdef0123456789abcdef01234567",
"phase": "RUNNING",
"healthStatus": "HEALTHY",
"webPort": 8400,
"webUrl": "http://192.168.35.82:8400",
"message": "预览实例已发布"
}))
@@ -346,6 +366,7 @@ async fn deploy_and_uninstall_use_fixed_job_and_apply_owned_artifacts() {
.clone();
assert_eq!(deployment.status, super::DeploymentStatus::Running);
assert_eq!(deployment.health, super::HealthStatus::Healthy);
assert_eq!(deployment.web_port, Some(8400));
assert_eq!(
deployment.web_url.as_deref(),
Some("http://192.168.35.82:8400")
@@ -402,8 +423,27 @@ async fn deploy_and_uninstall_use_fixed_job_and_apply_owned_artifacts() {
.clone();
assert_eq!(deployment.status, super::DeploymentStatus::Stopped);
assert_eq!(deployment.health, super::HealthStatus::Unknown);
assert_eq!(deployment.web_port, None);
assert_eq!(deployment.web_url, None);
assert!(!deployment.can_uninstall);
let list_request = axum::http::Request::builder()
.method("GET")
.uri("/api/preview-deployer/deployments")
.header(header::HOST, HOST)
.header(header::COOKIE, &cookie)
.body(Body::empty())
.unwrap();
let list_response = app.clone().oneshot(list_request).await.unwrap();
assert_eq!(list_response.status(), StatusCode::OK);
let list_body = list_response
.into_body()
.collect()
.await
.unwrap()
.to_bytes();
let list: Value = serde_json::from_slice(&list_body).unwrap();
assert_eq!(list["deployments"], json!([]));
let state_file = state.config.state_file.clone();
let mut recovered_config = test_config(state.config.jenkins_root_url.clone());
recovered_config.state_file = state_file.clone();
@@ -424,6 +464,44 @@ async fn deploy_and_uninstall_use_fixed_job_and_apply_owned_artifacts() {
std::fs::remove_file(state_file).unwrap();
}
#[tokio::test]
async fn transient_invalid_build_status_is_retried_before_reading_artifact() {
let (jenkins_url, mock) = start_mock_jenkins().await;
mock.invalid_build_responses.store(1, Ordering::SeqCst);
let state = AppState::new(test_config(jenkins_url)).unwrap();
let app = build_router(state.clone());
let cookie = login_cookie(&app).await;
let id = super::derive_deployment_id("feature/preview-ui");
let deploy_request = axum::http::Request::builder()
.method("POST")
.uri("/api/preview-deployer/deployments")
.header(header::HOST, HOST)
.header(header::ORIGIN, ORIGIN)
.header(header::COOKIE, &cookie)
.header(header::CONTENT_TYPE, "application/json")
.body(Body::from(r#"{"branch":"feature/preview-ui"}"#))
.unwrap();
assert_eq!(
app.oneshot(deploy_request).await.unwrap().status(),
StatusCode::ACCEPTED
);
wait_for_status(&state, &id, super::DeploymentStatus::Running)
.await
.expect("transient invalid Jenkins response is retried");
assert_eq!(
state
.deployments
.read()
.await
.get(&id)
.unwrap()
.public
.health,
super::HealthStatus::Healthy
);
std::fs::remove_file(&state.config.state_file).unwrap();
}
#[tokio::test]
async fn duplicate_active_branch_is_rejected_without_second_jenkins_trigger() {
let (jenkins_url, mock) = start_mock_jenkins().await;
@@ -442,6 +520,7 @@ async fn duplicate_active_branch_is_rejected_without_second_jenkins_trigger() {
resolved_commit: None,
status: super::DeploymentStatus::Building,
health: super::HealthStatus::Pending,
web_port: None,
web_url: None,
jenkins_build_url: None,
created_at: now,
@@ -519,3 +598,43 @@ fn branch_commit_and_web_url_validation_are_strict() {
"192.168.35.82"
));
}
#[test]
fn legacy_running_state_recovers_web_port_from_validated_url() {
let config = test_config(Url::parse("http://127.0.0.1:18080/jenkins/").unwrap());
let id = super::derive_deployment_id("feature/legacy-running");
std::fs::write(
&config.state_file,
serde_json::to_vec(&json!({
"schemaVersion": 1,
"deployments": [{
"public": {
"id": id,
"branch": "feature/legacy-running",
"status": "running",
"health": "healthy",
"webUrl": "http://192.168.35.82:8407",
"createdAt": 1,
"updatedAt": 2,
"canUninstall": true
},
"operation": "deploy"
}]
}))
.unwrap(),
)
.unwrap();
let state = AppState::new(config).unwrap();
assert_eq!(
state
.deployments
.blocking_read()
.get(&id)
.unwrap()
.public
.web_port,
Some(8407)
);
std::fs::remove_file(&state.config.state_file).unwrap();
}