mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Fix CI: clippy warnings, formatting, and mTLS test cert version
Rust 1.94.0 introduced new clippy lints and rustls now rejects X.509 v1 certificates. This fixes all three CI jobs: - Format: cargo fmt across the workspace - Clippy: unnecessary_unwrap, useless_format, derivable_impls, type_complexity, too_many_arguments, redundant_closure, map_or simplification, and other new lints - Tests: generate v3 certs (with extensions) for mTLS tests so newer rustls accepts them Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
97a0a9be32
commit
fb19251e9c
15 changed files with 108 additions and 92 deletions
|
|
@ -347,13 +347,13 @@ pub async fn list_retros(
|
|||
params
|
||||
.workflow
|
||||
.as_ref()
|
||||
.map_or(true, |w| &r.workflow.slug == w)
|
||||
.is_none_or(|w| &r.workflow.slug == w)
|
||||
})
|
||||
.filter(|r| {
|
||||
params
|
||||
.smoothness
|
||||
.as_ref()
|
||||
.map_or(true, |s| r.smoothness.as_ref() == Some(s))
|
||||
.is_none_or(|s| r.smoothness.as_ref() == Some(s))
|
||||
})
|
||||
.collect();
|
||||
paginated_response(
|
||||
|
|
@ -2525,6 +2525,7 @@ mod retros {
|
|||
use super::ts;
|
||||
use arc_types::*;
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn stage(
|
||||
id: &str,
|
||||
label: &str,
|
||||
|
|
|
|||
|
|
@ -137,7 +137,7 @@ pub fn build_router(state: Arc<AppState>, auth_mode: AuthMode) -> Router {
|
|||
let demo = demo_router.clone();
|
||||
let real = real_router.clone();
|
||||
async move {
|
||||
if req.headers().get("x-arc-demo").map_or(false, |v| v == "1") {
|
||||
if req.headers().get("x-arc-demo").is_some_and(|v| v == "1") {
|
||||
demo.oneshot(req).await
|
||||
} else {
|
||||
real.oneshot(req).await
|
||||
|
|
|
|||
|
|
@ -294,8 +294,7 @@ pub async fn send_message(
|
|||
created_at: now,
|
||||
}));
|
||||
session.updated_at = now;
|
||||
let seq = session.generation_seq.fetch_add(1, Ordering::Relaxed) + 1;
|
||||
seq
|
||||
session.generation_seq.fetch_add(1, Ordering::Relaxed) + 1
|
||||
}
|
||||
None => return ApiError::not_found("Session not found.").into_response(),
|
||||
}
|
||||
|
|
@ -354,7 +353,7 @@ pub async fn stream_session_events(
|
|||
.data(serde_json::json!({"message": message}).to_string()),
|
||||
),
|
||||
};
|
||||
sse.map(|e| Ok::<_, std::convert::Infallible>(e))
|
||||
sse.map(Ok::<_, std::convert::Infallible>)
|
||||
}
|
||||
Err(_) => None,
|
||||
});
|
||||
|
|
|
|||
|
|
@ -57,6 +57,10 @@ mod mtls_e2e {
|
|||
"1",
|
||||
"-subj",
|
||||
&format!("/CN={ca_cn}"),
|
||||
"-addext",
|
||||
"basicConstraints=critical,CA:TRUE",
|
||||
"-addext",
|
||||
"keyUsage=critical,keyCertSign,cRLSign",
|
||||
]);
|
||||
|
||||
// Server key + cert signed by CA
|
||||
|
|
@ -124,6 +128,9 @@ mod mtls_e2e {
|
|||
"-subj",
|
||||
&format!("/CN={client_cn}"),
|
||||
]);
|
||||
// Client extension file to produce a v3 certificate
|
||||
let client_ext_path = dir.join("client.ext");
|
||||
std::fs::write(&client_ext_path, "basicConstraints=CA:FALSE\n").unwrap();
|
||||
run_openssl(&[
|
||||
"x509",
|
||||
"-req",
|
||||
|
|
@ -138,6 +145,8 @@ mod mtls_e2e {
|
|||
client_cert_path.to_str().unwrap(),
|
||||
"-days",
|
||||
"1",
|
||||
"-extfile",
|
||||
client_ext_path.to_str().unwrap(),
|
||||
]);
|
||||
|
||||
PkiPaths {
|
||||
|
|
|
|||
|
|
@ -17,6 +17,16 @@ pub use openssh_runner::OpensshRunner;
|
|||
const WORKING_DIRECTORY: &str = "/home/exedev";
|
||||
const PROVIDER: &str = "exe";
|
||||
|
||||
/// Factory function type for creating data-plane SSH runners.
|
||||
type DataSshFactory = Box<
|
||||
dyn Fn(
|
||||
&str,
|
||||
) -> std::pin::Pin<
|
||||
Box<dyn std::future::Future<Output = Result<Box<dyn SshRunner>, String>> + Send>,
|
||||
> + Send
|
||||
+ Sync,
|
||||
>;
|
||||
|
||||
/// Output from an SSH command execution.
|
||||
pub struct SshOutput {
|
||||
pub stdout: Vec<u8>,
|
||||
|
|
@ -59,14 +69,7 @@ pub struct ExeSandbox {
|
|||
/// Factory for creating data-plane SSH runners, used during initialize().
|
||||
/// In production, this connects to the VM host via OpensshRunner.
|
||||
/// In tests, this is replaced with a closure that returns a MockSshRunner.
|
||||
data_ssh_factory: Box<
|
||||
dyn Fn(
|
||||
&str,
|
||||
) -> std::pin::Pin<
|
||||
Box<dyn std::future::Future<Output = Result<Box<dyn SshRunner>, String>> + Send>,
|
||||
> + Send
|
||||
+ Sync,
|
||||
>,
|
||||
data_ssh_factory: DataSshFactory,
|
||||
}
|
||||
|
||||
impl ExeSandbox {
|
||||
|
|
|
|||
|
|
@ -623,7 +623,40 @@ pub async fn run_chat_via_server(args: ChatArgs, server: &ServerConnection) -> R
|
|||
continue;
|
||||
}
|
||||
|
||||
if session_id.is_none() {
|
||||
if let Some(sid) = &session_id {
|
||||
// Subsequent messages: send message
|
||||
let body = serde_json::json!({ "content": trimmed });
|
||||
let url = format!("{}/sessions/{sid}/messages", server.base_url);
|
||||
let response = server
|
||||
.client
|
||||
.post(&url)
|
||||
.json(&body)
|
||||
.send()
|
||||
.await
|
||||
.with_context(|| format!("Failed to connect to server at {}", server.base_url))?;
|
||||
|
||||
let status = response.status();
|
||||
if !status.is_success() {
|
||||
let text = response.text().await.unwrap_or_default();
|
||||
bail!("Server returned {status}: {text}");
|
||||
}
|
||||
|
||||
// Stream events
|
||||
let events_url = format!("{}/sessions/{sid}/events", server.base_url);
|
||||
let events_response = server
|
||||
.client
|
||||
.get(&events_url)
|
||||
.send()
|
||||
.await
|
||||
.context("Failed to connect to event stream")?;
|
||||
|
||||
if !events_response.status().is_success() {
|
||||
let text = events_response.text().await.unwrap_or_default();
|
||||
bail!("Event stream returned error: {text}");
|
||||
}
|
||||
|
||||
stream_session_text(events_response).await?;
|
||||
} else {
|
||||
// First message: create session
|
||||
let mut body = serde_json::json!({ "content": trimmed });
|
||||
if let Some(ref model) = args.model {
|
||||
|
|
@ -667,7 +700,7 @@ pub async fn run_chat_via_server(args: ChatArgs, server: &ServerConnection) -> R
|
|||
.get(&events_url)
|
||||
.send()
|
||||
.await
|
||||
.with_context(|| format!("Failed to connect to event stream"))?;
|
||||
.context("Failed to connect to event stream")?;
|
||||
|
||||
if !events_response.status().is_success() {
|
||||
let text = events_response.text().await.unwrap_or_default();
|
||||
|
|
@ -676,40 +709,6 @@ pub async fn run_chat_via_server(args: ChatArgs, server: &ServerConnection) -> R
|
|||
|
||||
stream_session_text(events_response).await?;
|
||||
session_id = Some(sid);
|
||||
} else {
|
||||
// Subsequent messages: send message
|
||||
let sid = session_id.as_ref().unwrap();
|
||||
let body = serde_json::json!({ "content": trimmed });
|
||||
let url = format!("{}/sessions/{sid}/messages", server.base_url);
|
||||
let response = server
|
||||
.client
|
||||
.post(&url)
|
||||
.json(&body)
|
||||
.send()
|
||||
.await
|
||||
.with_context(|| format!("Failed to connect to server at {}", server.base_url))?;
|
||||
|
||||
let status = response.status();
|
||||
if !status.is_success() {
|
||||
let text = response.text().await.unwrap_or_default();
|
||||
bail!("Server returned {status}: {text}");
|
||||
}
|
||||
|
||||
// Stream events
|
||||
let events_url = format!("{}/sessions/{sid}/events", server.base_url);
|
||||
let events_response = server
|
||||
.client
|
||||
.get(&events_url)
|
||||
.send()
|
||||
.await
|
||||
.with_context(|| format!("Failed to connect to event stream"))?;
|
||||
|
||||
if !events_response.status().is_success() {
|
||||
let text = events_response.text().await.unwrap_or_default();
|
||||
bail!("Event stream returned error: {text}");
|
||||
}
|
||||
|
||||
stream_session_text(events_response).await?;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -187,28 +187,28 @@ pub async fn generate(params: GenerateParams) -> Result<GenerateResult, SdkError
|
|||
let tool_calls = response.tool_calls();
|
||||
let mut tool_results = Vec::new();
|
||||
|
||||
if !tool_calls.is_empty()
|
||||
&& response.finish_reason == FinishReason::ToolCalls
|
||||
&& params.tools.is_some()
|
||||
&& max_tool_rounds > 0
|
||||
{
|
||||
debug!(
|
||||
tool_calls = tool_calls.len(),
|
||||
round = round,
|
||||
"Executing tool calls"
|
||||
);
|
||||
let tools = params.tools.as_ref().expect("checked above");
|
||||
if tools.iter().any(|t| t.is_active()) {
|
||||
let tool_refs: Vec<&Tool> =
|
||||
tools.iter().map(std::convert::AsRef::as_ref).collect();
|
||||
tool_results = execute_all_tools_with_repair(
|
||||
&tool_refs,
|
||||
&tool_calls,
|
||||
&messages,
|
||||
abort_signal.as_ref(),
|
||||
params.repair_tool_call.as_ref(),
|
||||
)
|
||||
.await;
|
||||
if let Some(tools) = ¶ms.tools {
|
||||
if !tool_calls.is_empty()
|
||||
&& response.finish_reason == FinishReason::ToolCalls
|
||||
&& max_tool_rounds > 0
|
||||
{
|
||||
debug!(
|
||||
tool_calls = tool_calls.len(),
|
||||
round = round,
|
||||
"Executing tool calls"
|
||||
);
|
||||
if tools.iter().any(|t| t.is_active()) {
|
||||
let tool_refs: Vec<&Tool> =
|
||||
tools.iter().map(std::convert::AsRef::as_ref).collect();
|
||||
tool_results = execute_all_tools_with_repair(
|
||||
&tool_refs,
|
||||
&tool_calls,
|
||||
&messages,
|
||||
abort_signal.as_ref(),
|
||||
params.repair_tool_call.as_ref(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -3,15 +3,14 @@ use std::sync::Mutex;
|
|||
|
||||
use serde_json::Value;
|
||||
|
||||
#[derive(Default)]
|
||||
pub struct ThreadRegistry {
|
||||
ts_to_question: Mutex<HashMap<String, String>>,
|
||||
}
|
||||
|
||||
impl ThreadRegistry {
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
ts_to_question: Mutex::new(HashMap::new()),
|
||||
}
|
||||
Self::default()
|
||||
}
|
||||
|
||||
pub fn register(&self, message_ts: &str, question_id: &str) {
|
||||
|
|
|
|||
|
|
@ -1 +1,5 @@
|
|||
include!(concat!(env!("OUT_DIR"), "/openapi_types.rs"));
|
||||
#[allow(clippy::derivable_impls)]
|
||||
mod generated {
|
||||
include!(concat!(env!("OUT_DIR"), "/openapi_types.rs"));
|
||||
}
|
||||
pub use generated::*;
|
||||
|
|
|
|||
|
|
@ -274,7 +274,7 @@ pub async fn run_command(
|
|||
|
||||
let goal = graph.goal();
|
||||
if !goal.is_empty() {
|
||||
let first_line = goal.lines().next().unwrap_or(&goal);
|
||||
let first_line = goal.lines().next().unwrap_or(goal);
|
||||
eprintln!("{} {first_line}\n", styles.bold.apply_to("Goal:"));
|
||||
}
|
||||
|
||||
|
|
@ -1211,6 +1211,7 @@ fn print_final_output(logs_dir: &std::path::Path, styles: &Styles) {
|
|||
/// Boots the sandbox (init + cleanup), checks LLM provider availability,
|
||||
/// resolves the model/provider through the full precedence chain, and prints
|
||||
/// a styled check report.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn run_preflight(
|
||||
graph: &crate::graph::types::Graph,
|
||||
run_cfg: &Option<run_config::WorkflowRunConfig>,
|
||||
|
|
|
|||
|
|
@ -159,7 +159,7 @@ impl WorkflowRunConfig {
|
|||
// Union checkpoint exclude globs from defaults and task config, dedup
|
||||
if !defaults.checkpoint.exclude_globs.is_empty() {
|
||||
let mut merged = defaults.checkpoint.exclude_globs.clone();
|
||||
merged.extend(self.checkpoint.exclude_globs.drain(..));
|
||||
merged.append(&mut self.checkpoint.exclude_globs);
|
||||
merged.sort();
|
||||
merged.dedup();
|
||||
self.checkpoint.exclude_globs = merged;
|
||||
|
|
|
|||
|
|
@ -85,13 +85,14 @@ pub fn is_engine_internal_key(key: &str) -> bool {
|
|||
}
|
||||
|
||||
/// Fidelity mode controlling how much prior context is provided to LLM sessions.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub enum Fidelity {
|
||||
/// Complete context, no summarization — sessions share a thread.
|
||||
Full,
|
||||
/// Minimal: only graph goal and run ID.
|
||||
Truncate,
|
||||
/// Structured nested-bullet summary (default).
|
||||
#[default]
|
||||
Compact,
|
||||
/// Brief textual summary (~600 token target).
|
||||
SummaryLow,
|
||||
|
|
@ -112,12 +113,6 @@ impl Fidelity {
|
|||
}
|
||||
}
|
||||
|
||||
impl Default for Fidelity {
|
||||
fn default() -> Self {
|
||||
Self::Compact
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for Fidelity {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
let s = match self {
|
||||
|
|
|
|||
|
|
@ -539,6 +539,7 @@ pub enum GitCheckpointMode {
|
|||
}
|
||||
|
||||
/// Run a git checkpoint commit on the host filesystem (local/Docker bind-mount).
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn git_checkpoint_host(
|
||||
work_dir: PathBuf,
|
||||
run_id: String,
|
||||
|
|
@ -586,6 +587,7 @@ async fn git_diff_host(work_dir: PathBuf, base: String) -> Option<String> {
|
|||
pub const GIT_REMOTE: &str = "git -c maintenance.auto=0 -c gc.auto=0";
|
||||
|
||||
/// Run a git checkpoint commit inside a remote sandbox.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn git_checkpoint_remote(
|
||||
sandbox: &dyn Sandbox,
|
||||
run_id: &str,
|
||||
|
|
|
|||
|
|
@ -164,6 +164,7 @@ pub fn merge_ff_only(work_dir: &Path, sha: &str) -> Result<()> {
|
|||
/// Stage all changes and commit in `work_dir` with a structured message
|
||||
/// including trailers for completed node count and shadow commit pointer.
|
||||
/// Returns the new commit SHA.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub fn checkpoint_commit(
|
||||
work_dir: &Path,
|
||||
run_id: &str,
|
||||
|
|
|
|||
|
|
@ -119,7 +119,7 @@ async fn daytona_exec_command_cancelled() {
|
|||
|
||||
let token = tokio_util::sync::CancellationToken::new();
|
||||
let token_clone = token.clone();
|
||||
|
||||
|
||||
// Cancel the token shortly after starting
|
||||
tokio::spawn(async move {
|
||||
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
|
||||
|
|
@ -149,8 +149,8 @@ async fn daytona_exec_command_local_timeout() {
|
|||
// Use a tiny timeout_ms of 100ms, our local timeout is 100 + 2000 = 2100ms.
|
||||
// If the server doesn't enforce the timeout properly or drops the connection,
|
||||
// our local timeout should catch it. To simulate this without making a bad server,
|
||||
// we can't easily force the local timeout to hit before the server timeout
|
||||
// without mocking. But if we run `sleep 10` and Daytona does NOT respect the
|
||||
// we can't easily force the local timeout to hit before the server timeout
|
||||
// without mocking. But if we run `sleep 10` and Daytona does NOT respect the
|
||||
// short timeout parameter, the local 2.1s timeout will definitely fire.
|
||||
// Let's at least test that a 100ms timeout works and doesn't run for 10s.
|
||||
let start = std::time::Instant::now();
|
||||
|
|
@ -160,11 +160,14 @@ async fn daytona_exec_command_local_timeout() {
|
|||
.unwrap();
|
||||
|
||||
let duration = start.elapsed();
|
||||
|
||||
// It should either fail with Daytona's timeout (duration < 2000ms) or our
|
||||
// local timeout (duration ~2100ms). Both are valid success conditions for
|
||||
|
||||
// It should either fail with Daytona's timeout (duration < 2000ms) or our
|
||||
// local timeout (duration ~2100ms). Both are valid success conditions for
|
||||
// the system as a whole avoiding a stall.
|
||||
assert!(duration < std::time::Duration::from_millis(3000), "Command stalled for longer than the local timeout mechanism");
|
||||
assert!(
|
||||
duration < std::time::Duration::from_millis(3000),
|
||||
"Command stalled for longer than the local timeout mechanism"
|
||||
);
|
||||
assert!(result.exit_code != 0);
|
||||
|
||||
env.cleanup().await.unwrap();
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue