mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
Simplify workflow version registration client and tests
Move the content-derived id check into create_workflow_version so every caller gets it, drop the redundant expected_id parameter from register_workflow_versions, and collapse the duplicated httpmock setups and ordering machinery in the client tests. Use in-scope imports and the neighbouring reader idiom in the server intent tests. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
1a1c5ab29d
commit
637d6a0d84
2 changed files with 124 additions and 201 deletions
|
|
@ -24,7 +24,7 @@ use fabro_llm::types::{Message as LlmMessage, Request as LlmRequest, TokenCounts
|
|||
use fabro_model::catalog::LlmCatalogSettings;
|
||||
use fabro_model::{Catalog, ModelRef, ProviderId, ReasoningEffort, Speed};
|
||||
use fabro_types::settings::ServerAuthMethod;
|
||||
use fabro_types::settings::run::EnvironmentProvider;
|
||||
use fabro_types::settings::run::{ApprovalMode, EnvironmentProvider};
|
||||
use fabro_types::{
|
||||
AgentBackend, AttrValue, AuthMethod, BlobHash, CommandTermination, FailureCategory,
|
||||
FailureDetail, GitRunTarget, Graph, InterviewQuestionRecord, Node, Outcome, ParallelBranchId,
|
||||
|
|
@ -3807,13 +3807,10 @@ async fn post_runs_run_intent_args_true_override_resolved_settings_without_start
|
|||
assert_eq!(body["lifecycle"]["status"]["kind"], "submitted");
|
||||
let run_store = state.stores.runs.open_run_reader(&run_id).await.unwrap();
|
||||
let projection = run_store.state().await.unwrap();
|
||||
assert_eq!(
|
||||
projection.spec.settings.run.execution.mode,
|
||||
fabro_types::settings::run::RunMode::DryRun
|
||||
);
|
||||
assert_eq!(projection.spec.settings.run.execution.mode, RunMode::DryRun);
|
||||
assert_eq!(
|
||||
projection.spec.settings.run.execution.approval,
|
||||
fabro_types::settings::run::ApprovalMode::Auto
|
||||
ApprovalMode::Auto
|
||||
);
|
||||
assert!(projection.spec.settings.run.environment.lifecycle.preserve);
|
||||
}
|
||||
|
|
@ -3875,32 +3872,28 @@ preserve = true
|
|||
.parse::<RunId>()
|
||||
.unwrap();
|
||||
let omitted_id = omitted["id"].as_str().unwrap().parse::<RunId>().unwrap();
|
||||
let explicit_false = state
|
||||
let explicit_false_store = state
|
||||
.stores
|
||||
.runs
|
||||
.open_run_reader(&explicit_false_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.state()
|
||||
.await
|
||||
.unwrap();
|
||||
let omitted = state
|
||||
let explicit_false = explicit_false_store.state().await.unwrap();
|
||||
let omitted_store = state
|
||||
.stores
|
||||
.runs
|
||||
.open_run_reader(&omitted_id)
|
||||
.await
|
||||
.unwrap()
|
||||
.state()
|
||||
.await
|
||||
.unwrap();
|
||||
let omitted = omitted_store.state().await.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
explicit_false.spec.settings.run.execution.mode,
|
||||
fabro_types::settings::run::RunMode::Normal
|
||||
RunMode::Normal
|
||||
);
|
||||
assert_eq!(
|
||||
explicit_false.spec.settings.run.execution.approval,
|
||||
fabro_types::settings::run::ApprovalMode::Prompt
|
||||
ApprovalMode::Prompt
|
||||
);
|
||||
assert!(
|
||||
!explicit_false
|
||||
|
|
@ -3911,13 +3904,10 @@ preserve = true
|
|||
.lifecycle
|
||||
.preserve
|
||||
);
|
||||
assert_eq!(
|
||||
omitted.spec.settings.run.execution.mode,
|
||||
fabro_types::settings::run::RunMode::DryRun
|
||||
);
|
||||
assert_eq!(omitted.spec.settings.run.execution.mode, RunMode::DryRun);
|
||||
assert_eq!(
|
||||
omitted.spec.settings.run.execution.approval,
|
||||
fabro_types::settings::run::ApprovalMode::Auto
|
||||
ApprovalMode::Auto
|
||||
);
|
||||
assert!(omitted.spec.settings.run.environment.lifecycle.preserve);
|
||||
}
|
||||
|
|
@ -16556,7 +16546,7 @@ level = "debug"
|
|||
"goal should be persisted from the manifest"
|
||||
);
|
||||
assert!(
|
||||
resolved_run.execution.mode == fabro_types::settings::run::RunMode::DryRun,
|
||||
resolved_run.execution.mode == RunMode::DryRun,
|
||||
"run execution mode should inherit from server settings"
|
||||
);
|
||||
assert_eq!(
|
||||
|
|
|
|||
|
|
@ -701,38 +701,39 @@ impl Client {
|
|||
self.submit_create_run(manifest.into()).await
|
||||
}
|
||||
|
||||
/// Registers one workflow version and verifies the server assigned the
|
||||
/// content-derived id, so a mismatched response fails loudly here rather
|
||||
/// than being trusted downstream.
|
||||
pub async fn create_workflow_version(
|
||||
&self,
|
||||
version: &WorkflowVersion,
|
||||
) -> Result<WorkflowVersionId> {
|
||||
let expected_id = version.id()?;
|
||||
let response = self
|
||||
.send_api(|client| {
|
||||
let version = version.clone();
|
||||
async move { client.create_workflow_version().body(version).send().await }
|
||||
})
|
||||
.await?;
|
||||
Ok(response.into_inner().workflow_version_id)
|
||||
let returned_id = response.into_inner().workflow_version_id;
|
||||
if returned_id != expected_id {
|
||||
bail!(
|
||||
"workflow version registration returned {returned_id} for expected {expected_id}"
|
||||
);
|
||||
}
|
||||
Ok(returned_id)
|
||||
}
|
||||
|
||||
/// Registers versions in iteration order, stopping at the first failure.
|
||||
/// Callers must order dependencies before the versions that reference them.
|
||||
pub async fn register_workflow_versions<'a>(
|
||||
&self,
|
||||
versions: impl IntoIterator<Item = (WorkflowVersionId, &'a WorkflowVersion)>,
|
||||
versions: impl IntoIterator<Item = &'a WorkflowVersion>,
|
||||
) -> Result<()> {
|
||||
for (completed, (expected_id, version)) in versions.into_iter().enumerate() {
|
||||
let completed_noun = if completed == 1 { "entry" } else { "entries" };
|
||||
let returned_id = self
|
||||
.create_workflow_version(version)
|
||||
for (index, version) in versions.into_iter().enumerate() {
|
||||
self.create_workflow_version(version)
|
||||
.await
|
||||
.with_context(|| {
|
||||
format!(
|
||||
"failed to register workflow version {expected_id} after {completed} {completed_noun} completed"
|
||||
)
|
||||
})?;
|
||||
if returned_id != expected_id {
|
||||
bail!(
|
||||
"workflow version registration returned {returned_id} for expected {expected_id} after {completed} {completed_noun} completed"
|
||||
);
|
||||
}
|
||||
.with_context(|| format!("failed to register workflow version at index {index}"))?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
|
@ -2320,13 +2321,13 @@ fn add_pr_upgrade_hint(err: anyhow::Error) -> anyhow::Error {
|
|||
mod tests {
|
||||
use std::collections::BTreeMap;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
use chrono::Duration as ChronoDuration;
|
||||
use fabro_types::WorkflowPath;
|
||||
use fabro_util::exit;
|
||||
use httpmock::Method::{GET, POST};
|
||||
use httpmock::{HttpMockResponse, MockServer};
|
||||
use httpmock::MockServer;
|
||||
use serde_json::json;
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
use tokio::net::TcpListener;
|
||||
|
|
@ -2366,10 +2367,10 @@ mod tests {
|
|||
|
||||
fn test_workflow_version(
|
||||
name: &str,
|
||||
workflow_dependencies: BTreeMap<fabro_types::WorkflowPath, fabro_types::WorkflowVersionId>,
|
||||
) -> fabro_types::WorkflowVersion {
|
||||
let entrypoint = fabro_types::WorkflowPath::new("workflow.fabro").unwrap();
|
||||
fabro_types::WorkflowVersion::new(
|
||||
workflow_dependencies: BTreeMap<WorkflowPath, WorkflowVersionId>,
|
||||
) -> WorkflowVersion {
|
||||
let entrypoint = WorkflowPath::new("workflow.fabro").unwrap();
|
||||
WorkflowVersion::new(
|
||||
entrypoint.clone(),
|
||||
BTreeMap::from([(entrypoint, format!("digraph {name} {{}}"))]),
|
||||
workflow_dependencies,
|
||||
|
|
@ -2377,14 +2378,27 @@ mod tests {
|
|||
.unwrap()
|
||||
}
|
||||
|
||||
fn workflow_version_response(
|
||||
workflow_version_id: &fabro_types::WorkflowVersionId,
|
||||
) -> HttpMockResponse {
|
||||
HttpMockResponse::builder()
|
||||
.status(201)
|
||||
/// Mocks `POST /api/v1/workflow-versions` for exactly this version body.
|
||||
async fn mock_create_workflow_version<'a>(
|
||||
server: &'a MockServer,
|
||||
version: &WorkflowVersion,
|
||||
then: impl FnOnce(httpmock::Then) -> httpmock::Then,
|
||||
) -> httpmock::Mock<'a> {
|
||||
let body = serde_json::to_value(version).unwrap();
|
||||
server
|
||||
.mock_async(|when, respond| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(body);
|
||||
then(respond);
|
||||
})
|
||||
.await
|
||||
}
|
||||
|
||||
fn created_workflow_version(then: httpmock::Then, id: WorkflowVersionId) -> httpmock::Then {
|
||||
then.status(201)
|
||||
.header("content-type", "application/json")
|
||||
.body(json!({ "workflow_version_id": workflow_version_id }).to_string())
|
||||
.build()
|
||||
.json_body(json!({ "workflow_version_id": id }))
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
|
@ -2392,16 +2406,10 @@ mod tests {
|
|||
let server = MockServer::start_async().await;
|
||||
let version = test_workflow_version("ExactVersion", BTreeMap::new());
|
||||
let expected_id = version.id().unwrap();
|
||||
let mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&version).unwrap());
|
||||
then.status(201)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({ "workflow_version_id": expected_id }));
|
||||
})
|
||||
.await;
|
||||
let mock = mock_create_workflow_version(&server, &version, |then| {
|
||||
created_workflow_version(then, expected_id)
|
||||
})
|
||||
.await;
|
||||
|
||||
let client = Client::new_no_proxy(&server.url("")).unwrap();
|
||||
let actual_id = client.create_workflow_version(&version).await.unwrap();
|
||||
|
|
@ -2411,49 +2419,52 @@ mod tests {
|
|||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_workflow_versions_preserves_dependency_first_order() {
|
||||
async fn create_workflow_version_rejects_returned_id_mismatch() {
|
||||
let server = MockServer::start_async().await;
|
||||
let version = test_workflow_version("Expected", BTreeMap::new());
|
||||
let expected_id = version.id().unwrap();
|
||||
let returned_id = test_workflow_version("Returned", BTreeMap::new())
|
||||
.id()
|
||||
.unwrap();
|
||||
let mock = mock_create_workflow_version(&server, &version, |then| {
|
||||
created_workflow_version(then, returned_id)
|
||||
})
|
||||
.await;
|
||||
|
||||
let client = Client::new_no_proxy(&server.url("")).unwrap();
|
||||
let message = client
|
||||
.create_workflow_version(&version)
|
||||
.await
|
||||
.unwrap_err()
|
||||
.to_string();
|
||||
|
||||
mock.assert_async().await;
|
||||
assert!(message.contains(&expected_id.to_string()));
|
||||
assert!(message.contains(&returned_id.to_string()));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_workflow_versions_registers_every_entry_in_order() {
|
||||
let server = MockServer::start_async().await;
|
||||
let child = test_workflow_version("Child", BTreeMap::new());
|
||||
let child_id = child.id().unwrap();
|
||||
let parent = test_workflow_version(
|
||||
"Parent",
|
||||
BTreeMap::from([(fabro_types::WorkflowPath::new("child").unwrap(), child_id)]),
|
||||
BTreeMap::from([(WorkflowPath::new("child").unwrap(), child_id)]),
|
||||
);
|
||||
let parent_id = parent.id().unwrap();
|
||||
let child_response_id = child_id;
|
||||
let child_response_completed = Arc::new(AtomicBool::new(false));
|
||||
let child_response_completed_for_child = Arc::clone(&child_response_completed);
|
||||
let child_mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&child).unwrap());
|
||||
then.respond_with(move |_| {
|
||||
let response = workflow_version_response(&child_response_id);
|
||||
child_response_completed_for_child.store(true, Ordering::SeqCst);
|
||||
response
|
||||
});
|
||||
})
|
||||
.await;
|
||||
let parent_response_id = parent_id;
|
||||
let parent_mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&parent).unwrap());
|
||||
then.respond_with(move |_| {
|
||||
assert!(
|
||||
child_response_completed.load(Ordering::SeqCst),
|
||||
"parent registration began before the child response completed"
|
||||
);
|
||||
workflow_version_response(&parent_response_id)
|
||||
});
|
||||
})
|
||||
.await;
|
||||
let child_mock = mock_create_workflow_version(&server, &child, |then| {
|
||||
created_workflow_version(then, child_id)
|
||||
})
|
||||
.await;
|
||||
let parent_mock = mock_create_workflow_version(&server, &parent, |then| {
|
||||
created_workflow_version(then, parent_id)
|
||||
})
|
||||
.await;
|
||||
|
||||
let client = Client::new_no_proxy(&server.url("")).unwrap();
|
||||
client
|
||||
.register_workflow_versions([(child_id, &child), (parent_id, &parent)])
|
||||
.register_workflow_versions([&child, &parent])
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
|
|
@ -2461,113 +2472,46 @@ mod tests {
|
|||
parent_mock.assert_async().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_workflow_versions_rejects_returned_id_mismatch() {
|
||||
let server = MockServer::start_async().await;
|
||||
let first = test_workflow_version("Expected", BTreeMap::new());
|
||||
let expected_id = first.id().unwrap();
|
||||
let returned_id = test_workflow_version("Returned", BTreeMap::new())
|
||||
.id()
|
||||
.unwrap();
|
||||
let later = test_workflow_version("Later", BTreeMap::new());
|
||||
let later_id = later.id().unwrap();
|
||||
let first_mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&first).unwrap());
|
||||
then.status(201)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({ "workflow_version_id": returned_id }));
|
||||
})
|
||||
.await;
|
||||
let later_mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&later).unwrap());
|
||||
then.status(201)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({ "workflow_version_id": later_id }));
|
||||
})
|
||||
.await;
|
||||
|
||||
let client = Client::new_no_proxy(&server.url("")).unwrap();
|
||||
let error = client
|
||||
.register_workflow_versions([(expected_id, &first), (later_id, &later)])
|
||||
.await
|
||||
.unwrap_err();
|
||||
let message = error.to_string();
|
||||
|
||||
first_mock.assert_async().await;
|
||||
later_mock.assert_calls_async(0).await;
|
||||
assert!(message.contains(&expected_id.to_string()));
|
||||
assert!(message.contains(&returned_id.to_string()));
|
||||
assert!(message.contains("0 entries completed"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn register_workflow_versions_stops_after_request_failure() {
|
||||
let server = MockServer::start_async().await;
|
||||
let first = test_workflow_version("First", BTreeMap::new());
|
||||
let first_id = first.id().unwrap();
|
||||
let failing = test_workflow_version("Failing", BTreeMap::new());
|
||||
let failing_id = failing.id().unwrap();
|
||||
let later = test_workflow_version("NeverSent", BTreeMap::new());
|
||||
let later_id = later.id().unwrap();
|
||||
let first_mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&first).unwrap());
|
||||
then.status(201)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({ "workflow_version_id": first_id }));
|
||||
})
|
||||
.await;
|
||||
let failing_mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&failing).unwrap());
|
||||
then.status(422)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({
|
||||
"errors": [{
|
||||
"status": "422",
|
||||
"title": "Unprocessable Entity",
|
||||
"detail": "workflow dependency was not found",
|
||||
"code": "workflow_version_dependency_not_found"
|
||||
}]
|
||||
}));
|
||||
})
|
||||
.await;
|
||||
let later_mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&later).unwrap());
|
||||
then.status(201)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({ "workflow_version_id": later_id }));
|
||||
})
|
||||
.await;
|
||||
let first_mock = mock_create_workflow_version(&server, &first, |then| {
|
||||
created_workflow_version(then, first_id)
|
||||
})
|
||||
.await;
|
||||
let failing_mock = mock_create_workflow_version(&server, &failing, |then| {
|
||||
then.status(422)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({
|
||||
"errors": [{
|
||||
"status": "422",
|
||||
"title": "Unprocessable Entity",
|
||||
"detail": "workflow dependency was not found",
|
||||
"code": "workflow_version_dependency_not_found"
|
||||
}]
|
||||
}))
|
||||
})
|
||||
.await;
|
||||
let later_mock = mock_create_workflow_version(&server, &later, |then| {
|
||||
created_workflow_version(then, later_id)
|
||||
})
|
||||
.await;
|
||||
|
||||
let client = Client::new_no_proxy(&server.url("")).unwrap();
|
||||
let error = client
|
||||
.register_workflow_versions([
|
||||
(first_id, &first),
|
||||
(failing_id, &failing),
|
||||
(later_id, &later),
|
||||
])
|
||||
.register_workflow_versions([&first, &failing, &later])
|
||||
.await
|
||||
.unwrap_err();
|
||||
|
||||
first_mock.assert_async().await;
|
||||
failing_mock.assert_async().await;
|
||||
later_mock.assert_calls_async(0).await;
|
||||
assert!(error.to_string().contains(&failing_id.to_string()));
|
||||
assert!(error.to_string().contains("1 entry completed"));
|
||||
assert!(error.to_string().contains("index 1"));
|
||||
let failure = api_failure_for(&error).expect("API failure metadata should survive context");
|
||||
assert_eq!(failure.status, fabro_http::StatusCode::UNPROCESSABLE_ENTITY);
|
||||
assert_eq!(
|
||||
|
|
@ -2587,13 +2531,8 @@ mod tests {
|
|||
.await;
|
||||
let client = Client::new_no_proxy(&server.url("")).unwrap();
|
||||
|
||||
client
|
||||
.register_workflow_versions(std::iter::empty::<(
|
||||
fabro_types::WorkflowVersionId,
|
||||
&fabro_types::WorkflowVersion,
|
||||
)>())
|
||||
.await
|
||||
.unwrap();
|
||||
let none: [&WorkflowVersion; 0] = [];
|
||||
client.register_workflow_versions(none).await.unwrap();
|
||||
|
||||
unexpected.assert_calls_async(0).await;
|
||||
}
|
||||
|
|
@ -2603,16 +2542,10 @@ mod tests {
|
|||
let server = MockServer::start_async().await;
|
||||
let version = test_workflow_version("Repeated", BTreeMap::new());
|
||||
let expected_id = version.id().unwrap();
|
||||
let mock = server
|
||||
.mock_async(|when, then| {
|
||||
when.method(POST)
|
||||
.path("/api/v1/workflow-versions")
|
||||
.json_body(serde_json::to_value(&version).unwrap());
|
||||
then.status(201)
|
||||
.header("content-type", "application/json")
|
||||
.json_body(json!({ "workflow_version_id": expected_id }));
|
||||
})
|
||||
.await;
|
||||
let mock = mock_create_workflow_version(&server, &version, |then| {
|
||||
created_workflow_version(then, expected_id)
|
||||
})
|
||||
.await;
|
||||
let client = Client::new_no_proxy(&server.url("")).unwrap();
|
||||
|
||||
let first = client.create_workflow_version(&version).await.unwrap();
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue