From aebb06a1bbdc4c34a98bd8db053e040c7f3a1e0d Mon Sep 17 00:00:00 2001 From: wendrul Date: Thu, 3 Jul 2025 08:00:37 +0200 Subject: [PATCH] Fix errors --- backend/windmill-api/src/workspace_fork.rs | 20 +- backend/windmill-api/src/workspace_merge.rs | 295 +++++++++++--------- 2 files changed, 177 insertions(+), 138 deletions(-) diff --git a/backend/windmill-api/src/workspace_fork.rs b/backend/windmill-api/src/workspace_fork.rs index d75a69d480..78f88af66f 100644 --- a/backend/windmill-api/src/workspace_fork.rs +++ b/backend/windmill-api/src/workspace_fork.rs @@ -7,7 +7,7 @@ */ use crate::db::{ApiAuthed, DB}; -use axum::{extract::Path, routing::{get, post}, Json, Router}; +use axum::{extract::{Extension, Path}, routing::{get, post}, Json, Router}; use serde::{Deserialize, Serialize}; use sqlx::{FromRow, Postgres, Transaction}; use windmill_common::error::{Error, JsonResult, Result}; @@ -306,14 +306,14 @@ pub fn workspaced_service() -> Router { async fn create_fork_handler( authed: ApiAuthed, + Extension(db): Extension, Path(workspace_id): Path, Json(request): Json, ) -> JsonResult { - let db = &authed.db; require_admin(authed.is_admin, &authed.username)?; let fork = create_fork( - db, + &db, &workspace_id, &authed.email, &request.name, @@ -325,19 +325,18 @@ async fn create_fork_handler( async fn get_fork_info_handler( authed: ApiAuthed, + Extension(db): Extension, Path(workspace_id): Path, ) -> JsonResult> { - let db = &authed.db; - - let fork_info = get_fork_info(db, &workspace_id).await?; + let fork_info = get_fork_info(&db, &workspace_id).await?; Ok(Json(fork_info)) } async fn list_forks_handler( authed: ApiAuthed, + Extension(db): Extension, Path(workspace_id): Path, ) -> JsonResult> { - let db = &authed.db; require_admin(authed.is_admin, &authed.username)?; let forks = sqlx::query_as!( @@ -348,7 +347,7 @@ async fn list_forks_handler( ORDER BY created_at DESC", workspace_id ) - .fetch_all(db) + .fetch_all(&db) .await?; Ok(Json(forks)) @@ -356,10 +355,9 @@ async fn list_forks_handler( async fn list_resource_refs_handler( authed: ApiAuthed, + Extension(db): Extension, Path(workspace_id): Path, ) -> JsonResult> { - let db = &authed.db; - - let refs = get_forked_resource_refs(db, &workspace_id).await?; + let refs = get_forked_resource_refs(&db, &workspace_id).await?; Ok(Json(refs)) } \ No newline at end of file diff --git a/backend/windmill-api/src/workspace_merge.rs b/backend/windmill-api/src/workspace_merge.rs index 934d3a61e7..90aa1916ab 100644 --- a/backend/windmill-api/src/workspace_merge.rs +++ b/backend/windmill-api/src/workspace_merge.rs @@ -7,12 +7,16 @@ */ use crate::db::{ApiAuthed, DB}; -use axum::{extract::Path, routing::{get, post}, Json, Router}; +use axum::{ + extract::{Extension, Path}, + routing::{get, post}, + Json, Router, +}; use serde::{Deserialize, Serialize}; use sqlx::{FromRow, Postgres, Transaction}; -use windmill_common::error::{Error, JsonResult, Result}; -use windmill_common::utils::{require_admin}; use std::collections::HashMap; +use windmill_common::error::{Error, JsonResult, Result}; +use windmill_common::utils::require_admin; #[derive(FromRow, Serialize, Deserialize)] pub struct WorkspaceMergeRequest { @@ -92,7 +96,7 @@ pub async fn create_merge_request( auto_merge: bool, ) -> Result { let mut tx = db.begin().await?; - + // Verify that source is a fork of target let fork_info = sqlx::query_scalar!( "SELECT parent_workspace_id FROM workspace_fork WHERE fork_workspace_id = $1", @@ -100,17 +104,16 @@ pub async fn create_merge_request( ) .fetch_optional(&mut *tx) .await?; - - let parent_workspace_id = fork_info.ok_or_else(|| - Error::BadRequest("Source workspace is not a fork".to_string()) - )?; - + + let parent_workspace_id = + fork_info.ok_or_else(|| Error::BadRequest("Source workspace is not a fork".to_string()))?; + if parent_workspace_id != target_workspace_id { return Err(Error::BadRequest( - "Source workspace is not a fork of the target workspace".to_string() + "Source workspace is not a fork of the target workspace".to_string(), )); } - + // Create merge request let merge_request_id = sqlx::query_scalar!( "INSERT INTO workspace_merge_request @@ -126,22 +129,24 @@ pub async fn create_merge_request( ) .fetch_one(&mut *tx) .await?; - + // Analyze changes and create merge changes - let changes = analyze_workspace_changes(&mut tx, source_workspace_id, target_workspace_id, merge_request_id).await?; - + let changes = analyze_workspace_changes( + &mut tx, + source_workspace_id, + target_workspace_id, + merge_request_id, + ) + .await?; + tx.commit().await?; - + // Fetch the created merge request let merge_request = get_merge_request(db, merge_request_id).await?; - + let conflicts_count = changes.iter().filter(|c| c.has_conflict).count(); - - Ok(MergeRequestResponse { - merge_request, - changes, - conflicts_count, - }) + + Ok(MergeRequestResponse { merge_request, changes, conflicts_count }) } /// Analyze changes between fork and parent workspace @@ -152,7 +157,7 @@ async fn analyze_workspace_changes( merge_request_id: i64, ) -> Result> { let mut changes = Vec::new(); - + // Get all resources that have been modified in the fork (is_reference = false) let modified_resources = sqlx::query!( "SELECT resource_type, resource_path FROM forked_resource_refs @@ -161,53 +166,69 @@ async fn analyze_workspace_changes( ) .fetch_all(&mut **tx) .await?; - + for resource in modified_resources { let resource_type = &resource.resource_type; let resource_path = &resource.resource_path; - + // Get content hashes for comparison let (source_hash, target_hash) = match resource_type.as_str() { "script" => { let source = sqlx::query_scalar!( "SELECT encode(sha256(content::bytea), 'hex') FROM script WHERE workspace_id = $1 AND path = $2", - source_workspace_id, resource_path - ).fetch_optional(&mut **tx).await?.flatten(); - + source_workspace_id, + resource_path + ) + .fetch_optional(&mut **tx) + .await? + .flatten(); + let target = sqlx::query_scalar!( "SELECT encode(sha256(content::bytea), 'hex') FROM script WHERE workspace_id = $1 AND path = $2", - target_workspace_id, resource_path - ).fetch_optional(&mut **tx).await?.flatten(); - + target_workspace_id, + resource_path + ) + .fetch_optional(&mut **tx) + .await? + .flatten(); + (source, target) - }, + } "flow" => { let source = sqlx::query_scalar!( "SELECT encode(sha256(value::text::bytea), 'hex') FROM flow WHERE workspace_id = $1 AND path = $2", - source_workspace_id, resource_path - ).fetch_optional(&mut **tx).await?.flatten(); - + source_workspace_id, + resource_path + ) + .fetch_optional(&mut **tx) + .await? + .flatten(); + let target = sqlx::query_scalar!( "SELECT encode(sha256(value::text::bytea), 'hex') FROM flow WHERE workspace_id = $1 AND path = $2", - target_workspace_id, resource_path - ).fetch_optional(&mut **tx).await?.flatten(); - + target_workspace_id, + resource_path + ) + .fetch_optional(&mut **tx) + .await? + .flatten(); + (source, target) - }, + } _ => (None, None), // Add more resource types as needed }; - + let change_type = match (&source_hash, &target_hash) { (Some(_), None) => "added", (Some(s), Some(t)) if s != t => "modified", (None, Some(_)) => "deleted", _ => continue, // No change }; - + // Check for conflicts (if target has been modified since fork point) let fork_point = sqlx::query_scalar!( "SELECT fork_point FROM workspace_fork WHERE fork_workspace_id = $1", @@ -215,32 +236,40 @@ async fn analyze_workspace_changes( ) .fetch_one(&mut **tx) .await?; - + let target_modified_after_fork = match resource_type.as_str() { - "script" => { - sqlx::query_scalar!( - "SELECT created_at > $1 FROM script + "script" => sqlx::query_scalar!( + "SELECT created_at > $1 FROM script WHERE workspace_id = $2 AND path = $3", - fork_point, target_workspace_id, resource_path - ).fetch_optional(&mut **tx).await?.unwrap_or(false) - }, - "flow" => { - sqlx::query_scalar!( - "SELECT edited_at > $1 FROM flow + fork_point, + target_workspace_id, + resource_path + ) + .fetch_optional(&mut **tx) + .await? + .flatten() + .unwrap_or(false), + "flow" => sqlx::query_scalar!( + "SELECT edited_at > $1 FROM flow WHERE workspace_id = $2 AND path = $3", - fork_point, target_workspace_id, resource_path - ).fetch_optional(&mut **tx).await?.unwrap_or(false) - }, + fork_point, + target_workspace_id, + resource_path + ) + .fetch_optional(&mut **tx) + .await? + .flatten() + .unwrap_or(false), _ => false, }; - + let has_conflict = target_modified_after_fork && change_type == "modified"; let conflict_reason = if has_conflict { Some("Resource was modified in both source and target workspace".to_string()) } else { None }; - + // Create merge change record let change_id = sqlx::query_scalar!( "INSERT INTO workspace_merge_change @@ -259,7 +288,7 @@ async fn analyze_workspace_changes( ) .fetch_one(&mut **tx) .await?; - + changes.push(WorkspaceMergeChange { id: change_id, merge_request_id, @@ -277,7 +306,7 @@ async fn analyze_workspace_changes( created_at: chrono::Utc::now(), }); } - + Ok(changes) } @@ -293,25 +322,23 @@ pub async fn get_merge_request(db: &DB, merge_request_id: i64) -> Result Result<()> { +pub async fn execute_merge(db: &DB, merge_request_id: i64, merged_by: &str) -> Result<()> { let mut tx = db.begin().await?; - + // Get merge request details let merge_request = get_merge_request(db, merge_request_id).await?; - + if merge_request.status != "pending" && merge_request.status != "approved" { - return Err(Error::BadRequest("Merge request is not in a mergeable state".to_string())); + return Err(Error::BadRequest( + "Merge request is not in a mergeable state".to_string(), + )); } - + // Get all unresolved conflicts let unresolved_conflicts = sqlx::query_scalar!( "SELECT COUNT(*) FROM workspace_merge_change @@ -319,15 +346,16 @@ pub async fn execute_merge( merge_request_id ) .fetch_one(&mut *tx) - .await?; - + .await? + .unwrap_or(0); + if unresolved_conflicts > 0 { return Err(Error::BadRequest(format!( - "Cannot merge: {} unresolved conflicts remain", + "Cannot merge: {} unresolved conflicts remain", unresolved_conflicts ))); } - + // Apply all changes to target workspace let changes = sqlx::query_as!( WorkspaceMergeChange, @@ -339,16 +367,17 @@ pub async fn execute_merge( ) .fetch_all(&mut *tx) .await?; - + for change in changes { apply_change_to_workspace( - &mut tx, + &mut tx, &merge_request.source_workspace_id, &merge_request.target_workspace_id, - &change - ).await?; + &change, + ) + .await?; } - + // Mark merge request as merged sqlx::query!( "UPDATE workspace_merge_request @@ -359,9 +388,9 @@ pub async fn execute_merge( ) .execute(&mut *tx) .await?; - + tx.commit().await?; - + Ok(()) } @@ -405,7 +434,7 @@ async fn apply_change_to_workspace( ) .execute(&mut **tx) .await?; - }, + } "deleted" => { sqlx::query!( "UPDATE script SET deleted = true @@ -415,14 +444,18 @@ async fn apply_change_to_workspace( ) .execute(&mut **tx) .await?; - }, - _ => return Err(Error::InternalErr(format!("Unknown change type: {}", change.change_type))), + } + _ => { + return Err(Error::InternalErr(format!( + "Unknown change type: {}", + change.change_type + ))) + } } - }, - "flow" => { - match change.change_type.as_str() { - "added" | "modified" => { - sqlx::query!( + } + "flow" => match change.change_type.as_str() { + "added" | "modified" => { + sqlx::query!( "INSERT INTO flow (workspace_id, path, summary, description, value, edited_by, edited_at, archived, schema, extra_perms, dependency_job, draft_only, tag, @@ -443,26 +476,33 @@ async fn apply_change_to_workspace( ) .execute(&mut **tx) .await?; - }, - "deleted" => { - sqlx::query!( - "UPDATE flow SET archived = true + } + "deleted" => { + sqlx::query!( + "UPDATE flow SET archived = true WHERE workspace_id = $1 AND path = $2", - target_workspace_id, - change.resource_path - ) - .execute(&mut **tx) - .await?; - }, - _ => return Err(Error::InternalErr(format!("Unknown change type: {}", change.change_type))), + target_workspace_id, + change.resource_path + ) + .execute(&mut **tx) + .await?; + } + _ => { + return Err(Error::InternalErr(format!( + "Unknown change type: {}", + change.change_type + ))) } }, // Add more resource types as needed (apps, variables, etc.) _ => { - tracing::warn!("Merge not implemented for resource type: {}", change.resource_type); + tracing::warn!( + "Merge not implemented for resource type: {}", + change.resource_type + ); } } - + Ok(()) } @@ -473,45 +513,48 @@ pub fn workspaced_service() -> Router { .route("/merge_request/:id", get(get_merge_request_handler)) .route("/merge_request/:id/execute", post(execute_merge_handler)) .route("/merge_request/:id/changes", get(get_merge_changes_handler)) - .route("/merge_request/:id/resolve_conflict", post(resolve_conflict_handler)) + .route( + "/merge_request/:id/resolve_conflict", + post(resolve_conflict_handler), + ) } async fn create_merge_request_handler( authed: ApiAuthed, + Extension(db): Extension, Path(workspace_id): Path, Json(request): Json, ) -> JsonResult { - let db = &authed.db; require_admin(authed.is_admin, &authed.username)?; - + // Get parent workspace ID for this fork let parent_workspace_id = sqlx::query_scalar!( "SELECT parent_workspace_id FROM workspace_fork WHERE fork_workspace_id = $1", workspace_id ) - .fetch_optional(db) + .fetch_optional(&db) .await? .ok_or_else(|| Error::BadRequest("Workspace is not a fork".to_string()))?; - + let merge_response = create_merge_request( - db, + &db, &workspace_id, &parent_workspace_id, &authed.email, &request.title, request.description.as_deref(), request.auto_merge.unwrap_or(false), - ).await?; - + ) + .await?; + Ok(Json(merge_response)) } async fn list_merge_requests_handler( authed: ApiAuthed, + Extension(db): Extension, Path(workspace_id): Path, ) -> JsonResult> { - let db = &authed.db; - let merge_requests = sqlx::query_as!( WorkspaceMergeRequest, "SELECT id, source_workspace_id, target_workspace_id, created_by, created_at, @@ -522,39 +565,37 @@ async fn list_merge_requests_handler( ORDER BY created_at DESC", workspace_id ) - .fetch_all(db) + .fetch_all(&db) .await?; - + Ok(Json(merge_requests)) } async fn get_merge_request_handler( authed: ApiAuthed, + Extension(db): Extension, Path((workspace_id, merge_request_id)): Path<(String, i64)>, ) -> JsonResult { - let db = &authed.db; - - let merge_request = get_merge_request(db, merge_request_id).await?; + let merge_request = get_merge_request(&db, merge_request_id).await?; Ok(Json(merge_request)) } async fn execute_merge_handler( authed: ApiAuthed, + Extension(db): Extension, Path((workspace_id, merge_request_id)): Path<(String, i64)>, ) -> JsonResult<()> { - let db = &authed.db; require_admin(authed.is_admin, &authed.username)?; - - execute_merge(db, merge_request_id, &authed.email).await?; + + execute_merge(&db, merge_request_id, &authed.email).await?; Ok(Json(())) } async fn get_merge_changes_handler( authed: ApiAuthed, + Extension(db): Extension, Path((workspace_id, merge_request_id)): Path<(String, i64)>, ) -> JsonResult> { - let db = &authed.db; - let changes = sqlx::query_as!( WorkspaceMergeChange, "SELECT id, merge_request_id, resource_type, resource_path, change_type, @@ -564,20 +605,20 @@ async fn get_merge_changes_handler( ORDER BY resource_type, resource_path", merge_request_id ) - .fetch_all(db) + .fetch_all(&db) .await?; - + Ok(Json(changes)) } async fn resolve_conflict_handler( authed: ApiAuthed, + Extension(db): Extension, Path((workspace_id, merge_request_id)): Path<(String, i64)>, Json(request): Json, ) -> JsonResult<()> { - let db = &authed.db; require_admin(authed.is_admin, &authed.username)?; - + sqlx::query!( "UPDATE workspace_merge_change SET resolved = true, resolution_strategy = $1, resolved_by = $2, resolved_at = NOW() @@ -587,8 +628,8 @@ async fn resolve_conflict_handler( request.change_id, merge_request_id ) - .execute(db) + .execute(&db) .await?; - + Ok(Json(())) -} \ No newline at end of file +}