From 1f6aff68a0cf89162a7582c7a77946d1b2419885 Mon Sep 17 00:00:00 2001 From: "warp-agent-staging[bot]" <240773466+warp-agent-staging[bot]@users.noreply.github.com> Date: Fri, 4 Sep 2026 15:09:40 +0000 Subject: [PATCH 1/5] [REV-2383] Scope environment selection by team Co-Authored-By: Warp Agent --- app/src/ai/agent_sdk/ambient.rs | 21 ++-- app/src/ai/agent_sdk/common.rs | 23 +++-- app/src/ai/agent_sdk/common_tests.rs | 137 ++++++++++++++++++++++++++- app/src/ai/agent_sdk/environment.rs | 41 +++++--- app/src/ai/agent_sdk/integration.rs | 49 ++++++---- app/src/ai/agent_sdk/mod.rs | 6 +- app/src/ai/agent_sdk/schedule.rs | 40 ++++---- crates/warp_cli/src/environment.rs | 9 +- crates/warp_cli/src/lib_tests.rs | 37 ++++++++ 9 files changed, 289 insertions(+), 74 deletions(-) diff --git a/app/src/ai/agent_sdk/ambient.rs b/app/src/ai/agent_sdk/ambient.rs index 44b48a8ce26..05eade6231b 100644 --- a/app/src/ai/agent_sdk/ambient.rs +++ b/app/src/ai/agent_sdk/ambient.rs @@ -386,12 +386,14 @@ impl AmbientAgentRunner { vec![] }; - if let Err(err) = - super::common::validate_team_scope(&args.scope.team_selection, ctx) - { - super::report_fatal_error(err, ctx); - return; - } + let team_scope = + match super::common::resolve_team_scope(&args.scope.team_selection, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + super::report_fatal_error(err, ctx); + return; + } + }; let mut environment_args = args.environment; if environment_args.environment.is_none() && !environment_args.no_environment @@ -402,8 +404,11 @@ impl AmbientAgentRunner { environment_args.environment = Some(environment_id); } - let environment_id = match EnvironmentChoice::resolve_for_create(environment_args, ctx) - { + let environment_id = match EnvironmentChoice::resolve_for_create( + environment_args, + &team_scope, + ctx, + ) { Ok(EnvironmentChoice::None) => { eprintln!("Agent will run without an environment."); None diff --git a/app/src/ai/agent_sdk/common.rs b/app/src/ai/agent_sdk/common.rs index 227d231fb92..4273ac042bc 100644 --- a/app/src/ai/agent_sdk/common.rs +++ b/app/src/ai/agent_sdk/common.rs @@ -29,7 +29,7 @@ use crate::workspaces::update_manager::TeamUpdateManager; use crate::workspaces::user_workspaces::team_workspace_settings::{ NotATeamMemberError, TeamScopeForCli, TeamScopeForCliError, }; -use crate::workspaces::user_workspaces::{SoleTeamError, TeamScope as _, UserWorkspaces}; +use crate::workspaces::user_workspaces::{SoleTeamError, TeamScope, UserWorkspaces}; /// How long to wait for workspace metadata to refresh. pub const WORKSPACE_METADATA_REFRESH_TIMEOUT: Duration = Duration::from_secs(10); @@ -147,8 +147,8 @@ fn describe_team_resolution_error(error: TeamScopeForCliError, ctx: &AppContext) } } -/// The team a CLI command's policy reads are scoped to. -fn resolve_team_scope( +/// The team a CLI command acts within. +pub(super) fn resolve_team_scope( team_selection: &TeamSelection, ctx: &AppContext, ) -> anyhow::Result { @@ -213,13 +213,14 @@ pub fn resolve_owner(scope: &ObjectScope, ctx: &AppContext) -> anyhow::Result anyhow::Result<()> { - if !team_selection.is_team() { - return Ok(()); +pub(super) fn environment_is_visible_to_scope( + environment: &CloudAmbientAgentEnvironment, + team_scope: &(impl TeamScope + ?Sized), +) -> bool { + match environment.permissions().owner { + Owner::User { .. } => true, + Owner::Team { team_uid } => team_scope.team_uid() == Some(team_uid), } - resolve_team_scope(team_selection, ctx).map(|_| ()) } /// Refresh workspace metadata before executing an operation. @@ -324,10 +325,11 @@ pub enum EnvironmentChoice { } impl EnvironmentChoice { - /// Resolve the environment to use when creating an agent integration. + /// Resolve the environment to use when creating an agent operation. /// Warp Drive *must* have been synced first. pub fn resolve_for_create( args: EnvironmentCreateArgs, + team_scope: &(impl TeamScope + ?Sized), ctx: &AppContext, ) -> Result { if args.no_environment { @@ -339,6 +341,7 @@ impl EnvironmentChoice { let mut synced_environments: Vec<(ServerId, &CloudAmbientAgentEnvironment)> = all_environments .iter() + .filter(|env| environment_is_visible_to_scope(env, team_scope)) .filter_map(|env| { if let SyncId::ServerId(server_id) = env.sync_id() { Some((server_id, env)) diff --git a/app/src/ai/agent_sdk/common_tests.rs b/app/src/ai/agent_sdk/common_tests.rs index 8c04db7db08..d0a6bd7e0af 100644 --- a/app/src/ai/agent_sdk/common_tests.rs +++ b/app/src/ai/agent_sdk/common_tests.rs @@ -1,11 +1,16 @@ use std::collections::HashMap; +use warp_cli::environment::EnvironmentCreateArgs; use warpui::App; use super::{ - classify_agent_mode_base_model_id, parse_ambient_task_id, validate_agent_mode_base_model_id, + EnvironmentChoice, classify_agent_mode_base_model_id, environment_is_visible_to_scope, + parse_ambient_task_id, validate_agent_mode_base_model_id, }; use crate::LaunchMode; +use crate::ai::cloud_environments::{ + AmbientAgentEnvironment, CloudAmbientAgentEnvironment, CloudAmbientAgentEnvironmentModel, +}; use crate::ai::execution_profiles::profiles::AIExecutionProfilesModel; use crate::ai::llms::{ AvailableLLMs, LLMContextWindow, LLMId, LLMInfo, LLMPreferences, LLMProvider, LLMUsageMetadata, @@ -15,13 +20,141 @@ use crate::ai::mcp::TemplatableMCPServerManager; use crate::auth::AuthStateProvider; use crate::auth::auth_manager::AuthManager; use crate::cloud_object::model::persistence::CloudModel; +use crate::cloud_object::{CloudObjectMetadata, CloudObjectPermissions, Owner}; use crate::network::NetworkStatus; use crate::server::cloud_objects::update_manager::UpdateManager; +use crate::server::ids::{ServerId, SyncId}; use crate::server::server_api::ServerApiProvider; use crate::server::sync_queue::SyncQueue; use crate::test_util::settings::initialize_settings_for_tests; use crate::workspaces::team_tester::TeamTesterStatus; -use crate::workspaces::user_workspaces::{TeamlessScopeForTest, UserWorkspaces}; +use crate::workspaces::user_workspaces::{ + TeamContextForOperation, TeamlessScopeForTest, UserWorkspaces, +}; + +fn environment_with_owner( + sync_id: SyncId, + name: &str, + owner: Owner, +) -> CloudAmbientAgentEnvironment { + let environment = AmbientAgentEnvironment::new( + name.to_string(), + None, + Vec::new(), + "ubuntu:latest".to_string(), + Vec::new(), + ); + let mut permissions = CloudObjectPermissions::mock_personal(); + permissions.owner = owner; + CloudAmbientAgentEnvironment::new( + sync_id, + CloudAmbientAgentEnvironmentModel::new(environment), + CloudObjectMetadata::mock(), + permissions, + ) +} + +#[test] +fn environment_scope_includes_personal_and_matching_team_environments() { + let selected_team_uid = ServerId::from(123); + let other_team_uid = ServerId::from(456); + let selected_scope = TeamContextForOperation::new_for_test(selected_team_uid); + let personal_environment = environment_with_owner( + SyncId::ServerId(ServerId::from(1)), + "Personal", + Owner::mock_current_user(), + ); + let selected_team_environment = environment_with_owner( + SyncId::ServerId(ServerId::from(2)), + "Selected team", + Owner::Team { + team_uid: selected_team_uid, + }, + ); + let other_team_environment = environment_with_owner( + SyncId::ServerId(ServerId::from(3)), + "Other team", + Owner::Team { + team_uid: other_team_uid, + }, + ); + + assert!(environment_is_visible_to_scope( + &personal_environment, + &selected_scope + )); + assert!(environment_is_visible_to_scope( + &selected_team_environment, + &selected_scope + )); + assert!(!environment_is_visible_to_scope( + &other_team_environment, + &selected_scope + )); +} + +#[test] +fn teamless_environment_scope_includes_only_personal_environments() { + let personal_environment = environment_with_owner( + SyncId::ServerId(ServerId::from(1)), + "Personal", + Owner::mock_current_user(), + ); + let team_environment = environment_with_owner( + SyncId::ServerId(ServerId::from(2)), + "Team", + Owner::Team { + team_uid: ServerId::from(123), + }, + ); + + assert!(environment_is_visible_to_scope( + &personal_environment, + &TeamlessScopeForTest + )); + assert!(!environment_is_visible_to_scope( + &team_environment, + &TeamlessScopeForTest + )); +} + +#[test] +fn explicit_environment_id_remains_resource_authoritative() { + App::test((), |mut app| async move { + let cloud_model = app.add_singleton_model(CloudModel::mock); + let server_id = ServerId::from(123); + let sync_id = SyncId::ServerId(server_id); + let team_environment = environment_with_owner( + sync_id, + "Other team", + Owner::Team { + team_uid: ServerId::from(456), + }, + ); + cloud_model.update(&mut app, |model, ctx| { + model.create_object(sync_id, team_environment, ctx); + }); + + let choice = app.update(|ctx| { + EnvironmentChoice::resolve_for_create( + EnvironmentCreateArgs { + environment: Some(server_id.to_string()), + no_environment: false, + }, + &TeamlessScopeForTest, + ctx, + ) + }); + + assert_eq!( + choice.unwrap(), + EnvironmentChoice::Environment { + id: server_id.to_string(), + name: "Other team".to_string(), + } + ); + }); +} #[test] fn parse_ambient_task_id_accepts_valid_ids() { diff --git a/app/src/ai/agent_sdk/environment.rs b/app/src/ai/agent_sdk/environment.rs index 2af6f353483..677cb511251 100644 --- a/app/src/ai/agent_sdk/environment.rs +++ b/app/src/ai/agent_sdk/environment.rs @@ -2,13 +2,14 @@ use std::collections::HashSet; use comfy_table::Cell; use cynic::QueryBuilder; +use futures::future; use inquire::error::InquireError; use inquire::{Confirm, Select}; use serde::Serialize; use warp_cli::GlobalOptions; use warp_cli::agent::OutputFormat; use warp_cli::environment::{EnvironmentCommand, ImageCommand}; -use warp_cli::scope::ObjectScope; +use warp_cli::scope::{ObjectScope, TeamSelection}; use warp_graphql::queries::get_oauth_connect_tx_status::OauthConnectTxStatus; use warp_graphql::queries::list_warp_dev_images::{ ListWarpDevImages, ListWarpDevImagesResult, ListWarpDevImagesVariables, @@ -63,8 +64,10 @@ pub fn run( ) -> anyhow::Result<()> { let runner = ctx.add_singleton_model(|_ctx| EnvironmentCommandRunner); match command { - EnvironmentCommand::List => { - runner.update(ctx, |runner, ctx| runner.list(global_options, ctx)); + EnvironmentCommand::List { team_selection } => { + runner.update(ctx, |runner, ctx| { + runner.list(global_options, team_selection, ctx) + }); Ok(()) } EnvironmentCommand::Create { @@ -184,24 +187,36 @@ impl EnvironmentCommandRunner { }); } - fn list(&self, global_options: GlobalOptions, ctx: &mut ModelContext) { - let initial_sync = UpdateManager::as_ref(ctx) - .initial_load_complete() - .with_timeout(WARP_DRIVE_SYNC_TIMEOUT); + fn list( + &self, + global_options: GlobalOptions, + team_selection: TeamSelection, + ctx: &mut ModelContext, + ) { + let refresh_future = super::common::refresh_workspace_metadata(ctx); + let warp_drive_sync_future = super::common::refresh_warp_drive(ctx); + let setup_future = future::try_join(refresh_future, warp_drive_sync_future); - ctx.spawn(initial_sync, move |_, result, ctx| { - if result.is_err() { - super::report_fatal_error( - anyhow::anyhow!("Timed out waiting for Warp Drive to sync"), - ctx, - ); + ctx.spawn(setup_future, move |_, result, ctx| { + if let Err(err) = result { + super::report_fatal_error(err, ctx); return; } + let team_scope = match super::common::resolve_team_scope(&team_selection, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + super::report_fatal_error(err, ctx); + return; + } + }; let environments = CloudAmbientAgentEnvironment::get_all(ctx); let environment_infos: Vec<_> = environments .iter() + .filter(|environment| { + super::common::environment_is_visible_to_scope(environment, &team_scope) + }) .map(|environment| { let name = environment.model().string_model.name.clone(); let description = environment.model().string_model.description.clone(); diff --git a/app/src/ai/agent_sdk/integration.rs b/app/src/ai/agent_sdk/integration.rs index 681527869c5..c0d077c5def 100644 --- a/app/src/ai/agent_sdk/integration.rs +++ b/app/src/ai/agent_sdk/integration.rs @@ -2,6 +2,7 @@ use futures::future; use warp_cli::GlobalOptions; use warp_cli::integration::{CreateIntegrationArgs, IntegrationCommand, UpdateIntegrationArgs}; use warp_cli::provider::ProviderType; +use warp_cli::scope::TeamSelection; use warp_graphql::mutations::create_simple_integration::CreateSimpleIntegrationOutput; use warp_graphql::queries::get_oauth_connect_tx_status::OauthConnectTxStatus; use warp_graphql::queries::get_simple_integrations::SimpleIntegrationsOutput; @@ -73,6 +74,14 @@ impl IntegrationCommandRunner { ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); return; } + let team_scope = + match super::common::resolve_team_scope(&TeamSelection { team: None }, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); + return; + } + }; let loaded_file = match args.config_file.file.as_deref() { Some(path) => match super::config_file::load_config_file(path) { @@ -154,26 +163,26 @@ impl IntegrationCommandRunner { environment_args.environment = merged_config.environment_id.take(); } - let environment_uid = match EnvironmentChoice::resolve_for_create(environment_args, ctx) - { - Ok(EnvironmentChoice::None) => { - eprintln!("Creating integration without an environment."); - None - } - Ok(EnvironmentChoice::Environment { id, .. }) => { - eprintln!("Creating integration with environment {id}."); - Some(id) - } - Err(ResolveConfigurationError::Canceled) => { - eprintln!("Integration creation canceled."); - ctx.terminate_app(TerminationMode::ForceTerminate, None); - return; - } - Err(err) => { - super::report_fatal_error(anyhow::anyhow!(err), ctx); - return; - } - }; + let environment_uid = + match EnvironmentChoice::resolve_for_create(environment_args, &team_scope, ctx) { + Ok(EnvironmentChoice::None) => { + eprintln!("Creating integration without an environment."); + None + } + Ok(EnvironmentChoice::Environment { id, .. }) => { + eprintln!("Creating integration with environment {id}."); + Some(id) + } + Err(ResolveConfigurationError::Canceled) => { + eprintln!("Integration creation canceled."); + ctx.terminate_app(TerminationMode::ForceTerminate, None); + return; + } + Err(err) => { + super::report_fatal_error(anyhow::anyhow!(err), ctx); + return; + } + }; runner.start_create_or_update_flow( ctx, diff --git a/app/src/ai/agent_sdk/mod.rs b/app/src/ai/agent_sdk/mod.rs index eeadd8dff33..fd6b905e86c 100644 --- a/app/src/ai/agent_sdk/mod.rs +++ b/app/src/ai/agent_sdk/mod.rs @@ -1583,7 +1583,7 @@ fn command_requires_auth(command: &CliCommand) -> bool { AgentCommand::Skills(_) => true, }, CliCommand::Environment(environment_cmd) => match environment_cmd { - EnvironmentCommand::List => true, + EnvironmentCommand::List { .. } => true, EnvironmentCommand::Create { .. } => true, EnvironmentCommand::Delete { .. } => true, EnvironmentCommand::Update { .. } => true, @@ -1809,7 +1809,9 @@ fn command_to_telemetry_event(command: &CliCommand) -> CliTelemetryEvent { CliCommand::Agent(AgentCommand::Update(_)) => CliTelemetryEvent::AgentUpdate, CliCommand::Agent(AgentCommand::Delete(_)) => CliTelemetryEvent::AgentDelete, CliCommand::Agent(AgentCommand::Skills(_)) => CliTelemetryEvent::AgentSkills, - CliCommand::Environment(EnvironmentCommand::List) => CliTelemetryEvent::EnvironmentList, + CliCommand::Environment(EnvironmentCommand::List { .. }) => { + CliTelemetryEvent::EnvironmentList + } CliCommand::Environment(EnvironmentCommand::Create { .. }) => { CliTelemetryEvent::EnvironmentCreate } diff --git a/app/src/ai/agent_sdk/schedule.rs b/app/src/ai/agent_sdk/schedule.rs index fec13daa9e7..fecd9b5e0bd 100644 --- a/app/src/ai/agent_sdk/schedule.rs +++ b/app/src/ai/agent_sdk/schedule.rs @@ -51,6 +51,14 @@ fn create(ctx: &mut AppContext, args: CreateScheduleArgs) -> anyhow::Result<()> super::report_fatal_error(err, ctx); return; } + let team_scope = + match super::common::resolve_team_scope(&args.scope.team_selection, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + super::report_fatal_error(err, ctx); + return; + } + }; let loaded_file = match args.config_file.file.as_deref() { Some(path) => match super::config_file::load_config_file(path) { @@ -73,22 +81,22 @@ fn create(ctx: &mut AppContext, args: CreateScheduleArgs) -> anyhow::Result<()> environment_args.environment = Some(environment_id); } - let environment_id = match EnvironmentChoice::resolve_for_create(environment_args, ctx) - { - Ok(EnvironmentChoice::None) => { - eprintln!("Scheduling agent to run without an environment."); - None - } - Ok(EnvironmentChoice::Environment { id, .. }) => Some(id), - Err(ResolveConfigurationError::Canceled) => { - ctx.terminate_app(TerminationMode::ForceTerminate, None); - return; - } - Err(err) => { - super::report_fatal_error(anyhow::anyhow!(err), ctx); - return; - } - }; + let environment_id = + match EnvironmentChoice::resolve_for_create(environment_args, &team_scope, ctx) { + Ok(EnvironmentChoice::None) => { + eprintln!("Scheduling agent to run without an environment."); + None + } + Ok(EnvironmentChoice::Environment { id, .. }) => Some(id), + Err(ResolveConfigurationError::Canceled) => { + ctx.terminate_app(TerminationMode::ForceTerminate, None); + return; + } + Err(err) => { + super::report_fatal_error(anyhow::anyhow!(err), ctx); + return; + } + }; let owner = match super::common::resolve_owner(&args.scope, ctx) { Ok(owner) => owner, diff --git a/crates/warp_cli/src/environment.rs b/crates/warp_cli/src/environment.rs index 31d8b67e9dd..ea70c83613f 100644 --- a/crates/warp_cli/src/environment.rs +++ b/crates/warp_cli/src/environment.rs @@ -1,6 +1,6 @@ use clap::{ArgAction, ArgGroup, Args, Subcommand}; -use crate::scope::ObjectScope; +use crate::scope::{ObjectScope, TeamSelection}; /// Maximum length for environment descriptions. const MAX_DESCRIPTION_LENGTH: usize = 240; @@ -24,7 +24,10 @@ fn validate_description(s: &str) -> Result { #[command(visible_alias = "e")] pub enum EnvironmentCommand { /// List cloud environments. - List, + List { + #[command(flatten)] + team_selection: TeamSelection, + }, /// Manage base images for cloud environments. #[command(subcommand)] Image(ImageCommand), @@ -104,7 +107,7 @@ pub enum EnvironmentCommand { impl EnvironmentCommand { pub(crate) fn as_str_for_tracing(&self) -> &'static str { match self { - EnvironmentCommand::List => "environment list", + EnvironmentCommand::List { .. } => "environment list", EnvironmentCommand::Image(_) => "environment image", EnvironmentCommand::Create { .. } => "environment create", EnvironmentCommand::Delete { .. } => "environment delete", diff --git a/crates/warp_cli/src/lib_tests.rs b/crates/warp_cli/src/lib_tests.rs index c6547c08a90..b249413f56d 100644 --- a/crates/warp_cli/src/lib_tests.rs +++ b/crates/warp_cli/src/lib_tests.rs @@ -2262,6 +2262,43 @@ fn environment_image_list_parses() { assert!(matches!(image_cmd, ImageCommand::List)); } +fn parse_environment_list(args: &[&str]) -> crate::scope::TeamSelection { + let full_args = std::iter::once("warp") + .chain(["environment", "list"]) + .chain(args.iter().copied()); + let args = Args::try_parse_from(full_args).expect("environment list args should parse"); + + let Some(Command::CommandLine(boxed_cmd)) = args.command else { + panic!("Expected `warp environment list` command"); + }; + let CliCommand::Environment(EnvironmentCommand::List { team_selection }) = boxed_cmd.as_ref() + else { + panic!("Expected `warp environment list` command"); + }; + + team_selection.clone() +} + +#[test] +fn environment_list_defaults_to_implicit_team_selection() { + let team_selection = parse_environment_list(&[]); + + assert_eq!(team_selection.team, None); +} + +#[test] +fn environment_list_accepts_bare_team_selection() { + let team_selection = parse_environment_list(&["--team"]); + + assert_eq!(team_selection.team, Some(None)); +} + +#[test] +fn environment_list_accepts_explicit_team_selection() { + let team_selection = parse_environment_list(&["--team=123"]); + + assert_eq!(team_selection.team, Some(Some("123".to_string()))); +} #[test] fn environment_create_accepts_description() { From 6f2563c9b0d3484614bdd04d154409a238c7eec0 Mon Sep 17 00:00:00 2001 From: "warp-agent-staging[bot]" <240773466+warp-agent-staging[bot]@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:31:21 +0000 Subject: [PATCH 2/5] fix(cli): honor personal environment scope Co-Authored-By: Warp Agent --- app/src/ai/agent_sdk/ambient.rs | 15 ++-- app/src/ai/agent_sdk/common.rs | 10 +++ app/src/ai/agent_sdk/common_tests.rs | 87 ++++++++++++++----- app/src/ai/agent_sdk/schedule.rs | 15 ++-- .../team_workspace_settings.rs | 7 ++ 5 files changed, 95 insertions(+), 39 deletions(-) diff --git a/app/src/ai/agent_sdk/ambient.rs b/app/src/ai/agent_sdk/ambient.rs index 05eade6231b..4ac5ed52227 100644 --- a/app/src/ai/agent_sdk/ambient.rs +++ b/app/src/ai/agent_sdk/ambient.rs @@ -386,14 +386,13 @@ impl AmbientAgentRunner { vec![] }; - let team_scope = - match super::common::resolve_team_scope(&args.scope.team_selection, ctx) { - Ok(team_scope) => team_scope, - Err(err) => { - super::report_fatal_error(err, ctx); - return; - } - }; + let team_scope = match super::common::resolve_environment_team_scope(&args.scope, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + super::report_fatal_error(err, ctx); + return; + } + }; let mut environment_args = args.environment; if environment_args.environment.is_none() && !environment_args.no_environment diff --git a/app/src/ai/agent_sdk/common.rs b/app/src/ai/agent_sdk/common.rs index 4273ac042bc..4389bbd36e8 100644 --- a/app/src/ai/agent_sdk/common.rs +++ b/app/src/ai/agent_sdk/common.rs @@ -156,6 +156,16 @@ pub(super) fn resolve_team_scope( .team_scope_for_cli(team_selection) .map_err(|err| describe_team_resolution_error(err, ctx)) } +pub(super) fn resolve_environment_team_scope( + scope: &ObjectScope, + ctx: &AppContext, +) -> anyhow::Result { + if scope.personal { + Ok(TeamScopeForCli::personal()) + } else { + resolve_team_scope(&scope.team_selection, ctx) + } +} /// [`validate_agent_mode_base_model_id`], also rejecting a model `scope`'s team does not let this /// member use. diff --git a/app/src/ai/agent_sdk/common_tests.rs b/app/src/ai/agent_sdk/common_tests.rs index d0a6bd7e0af..0838ecb6188 100644 --- a/app/src/ai/agent_sdk/common_tests.rs +++ b/app/src/ai/agent_sdk/common_tests.rs @@ -1,11 +1,12 @@ use std::collections::HashMap; use warp_cli::environment::EnvironmentCreateArgs; +use warp_cli::scope::{ObjectScope, TeamSelection}; use warpui::App; use super::{ EnvironmentChoice, classify_agent_mode_base_model_id, environment_is_visible_to_scope, - parse_ambient_task_id, validate_agent_mode_base_model_id, + parse_ambient_task_id, resolve_environment_team_scope, validate_agent_mode_base_model_id, }; use crate::LaunchMode; use crate::ai::cloud_environments::{ @@ -26,10 +27,11 @@ use crate::server::cloud_objects::update_manager::UpdateManager; use crate::server::ids::{ServerId, SyncId}; use crate::server::server_api::ServerApiProvider; use crate::server::sync_queue::SyncQueue; +use crate::settings::PrivacySettings; use crate::test_util::settings::initialize_settings_for_tests; use crate::workspaces::team_tester::TeamTesterStatus; use crate::workspaces::user_workspaces::{ - TeamContextForOperation, TeamlessScopeForTest, UserWorkspaces, + TeamContextForOperation, TeamScope, TeamlessScopeForTest, UserWorkspaces, }; fn environment_with_owner( @@ -94,28 +96,67 @@ fn environment_scope_includes_personal_and_matching_team_environments() { } #[test] -fn teamless_environment_scope_includes_only_personal_environments() { - let personal_environment = environment_with_owner( - SyncId::ServerId(ServerId::from(1)), - "Personal", - Owner::mock_current_user(), - ); - let team_environment = environment_with_owner( - SyncId::ServerId(ServerId::from(2)), - "Team", - Owner::Team { - team_uid: ServerId::from(123), - }, - ); +fn multi_team_personal_scope_includes_only_personal_environments() { + App::test((), |mut app| async move { + initialize_settings_for_tests(&mut app); + app.add_singleton_model(PrivacySettings::mock); + let user_workspaces = app.add_singleton_model(UserWorkspaces::default_mock); + user_workspaces.update(&mut app, |user_workspaces, ctx| { + user_workspaces.setup_test_workspace(ctx); + user_workspaces.update_current_workspace( + |workspace| { + let mut second_team = workspace.teams[0].clone(); + second_team.uid = ServerId::from(456); + second_team.name = "Second team".to_string(); + workspace.teams.push(second_team); + }, + ctx, + ); + }); + let implicit_scope = app.read(|ctx| { + resolve_environment_team_scope( + &ObjectScope { + team_selection: TeamSelection { team: None }, + personal: false, + }, + ctx, + ) + }); + let personal_scope = app + .read(|ctx| { + resolve_environment_team_scope( + &ObjectScope { + team_selection: TeamSelection { team: None }, + personal: true, + }, + ctx, + ) + }) + .expect("explicit personal scope should not require a sole team"); + let personal_environment = environment_with_owner( + SyncId::ServerId(ServerId::from(1)), + "Personal", + Owner::mock_current_user(), + ); + let team_environment = environment_with_owner( + SyncId::ServerId(ServerId::from(2)), + "Team", + Owner::Team { + team_uid: ServerId::from(123), + }, + ); - assert!(environment_is_visible_to_scope( - &personal_environment, - &TeamlessScopeForTest - )); - assert!(!environment_is_visible_to_scope( - &team_environment, - &TeamlessScopeForTest - )); + assert!(implicit_scope.is_err()); + assert_eq!(personal_scope.team_uid(), None); + assert!(environment_is_visible_to_scope( + &personal_environment, + &personal_scope + )); + assert!(!environment_is_visible_to_scope( + &team_environment, + &personal_scope + )); + }); } #[test] diff --git a/app/src/ai/agent_sdk/schedule.rs b/app/src/ai/agent_sdk/schedule.rs index fecd9b5e0bd..62f0fcd2a1a 100644 --- a/app/src/ai/agent_sdk/schedule.rs +++ b/app/src/ai/agent_sdk/schedule.rs @@ -51,14 +51,13 @@ fn create(ctx: &mut AppContext, args: CreateScheduleArgs) -> anyhow::Result<()> super::report_fatal_error(err, ctx); return; } - let team_scope = - match super::common::resolve_team_scope(&args.scope.team_selection, ctx) { - Ok(team_scope) => team_scope, - Err(err) => { - super::report_fatal_error(err, ctx); - return; - } - }; + let team_scope = match super::common::resolve_environment_team_scope(&args.scope, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + super::report_fatal_error(err, ctx); + return; + } + }; let loaded_file = match args.config_file.file.as_deref() { Some(path) => match super::config_file::load_config_file(path) { diff --git a/app/src/workspaces/user_workspaces/team_workspace_settings.rs b/app/src/workspaces/user_workspaces/team_workspace_settings.rs index 92ebfdfd361..acb98aa29a0 100644 --- a/app/src/workspaces/user_workspaces/team_workspace_settings.rs +++ b/app/src/workspaces/user_workspaces/team_workspace_settings.rs @@ -101,6 +101,13 @@ impl TeamScope for TeamContext<'_> { /// memberships instead of from a window. #[cfg(not(target_family = "wasm"))] pub struct TeamScopeForCli(Option); +#[cfg(not(target_family = "wasm"))] +impl TeamScopeForCli { + /// A teamless scope selected explicitly with a personal CLI flag. + pub(crate) fn personal() -> Self { + Self(None) + } +} #[cfg(not(target_family = "wasm"))] impl sealed::Sealed for TeamScopeForCli {} From 6e4cc7858d8a4b6a3d231d0854fc21789f81ae0c Mon Sep 17 00:00:00 2001 From: "warp-agent-staging[bot]" <240773466+warp-agent-staging[bot]@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:10:31 +0000 Subject: [PATCH 3/5] [REV-2383] Scope integrations and provider setup by team --- app/src/ai/agent_sdk/integration.rs | 137 +++++++++++++----- app/src/ai/agent_sdk/integration_tests.rs | 27 ++++ app/src/ai/agent_sdk/mod.rs | 2 +- app/src/ai/agent_sdk/provider.rs | 21 +-- app/src/server/server_api.rs | 20 +++ app/src/server/server_api/integrations.rs | 13 +- crates/warp_cli/src/integration.rs | 12 +- crates/warp_cli/src/lib_tests.rs | 101 +++++++++++++ crates/warp_cli/src/provider.rs | 13 +- .../warp_server_client/src/graphql_helpers.rs | 31 +++- .../src/graphql_helpers_tests.rs | 29 +++- 11 files changed, 340 insertions(+), 66 deletions(-) create mode 100644 app/src/ai/agent_sdk/integration_tests.rs diff --git a/app/src/ai/agent_sdk/integration.rs b/app/src/ai/agent_sdk/integration.rs index c0d077c5def..944056532ce 100644 --- a/app/src/ai/agent_sdk/integration.rs +++ b/app/src/ai/agent_sdk/integration.rs @@ -13,6 +13,7 @@ use super::common::{EnvironmentChoice, ResolveConfigurationError}; use super::integration_output; use super::oauth_flow::poll_oauth_until_terminal; use crate::server::server_api::ServerApiProvider; +use crate::server::team_scope::RequestTeamScope; pub fn run( ctx: &mut AppContext, @@ -27,41 +28,86 @@ pub fn run( IntegrationCommand::Update(args) => { runner.update(ctx, |runner, ctx| runner.update(args, ctx)); } - IntegrationCommand::List => { - runner.update(ctx, |runner, ctx| runner.list(global_options, ctx)); + IntegrationCommand::List { team_selection } => { + runner.update(ctx, |runner, ctx| { + runner.list(global_options, team_selection, ctx) + }); } } Ok(()) } struct IntegrationCommandRunner; +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct IntegrationRetryState { + request_team_scope: RequestTeamScope, + attempt: u32, +} -impl IntegrationCommandRunner { - fn list(&self, global_options: GlobalOptions, ctx: &mut ModelContext) { - // Hardcoded set of providers that this client knows how to render. - let providers = vec![ProviderType::Linear, ProviderType::Slack]; - let provider_slugs: Vec = providers.into_iter().map(|p| p.slug()).collect(); - - let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); +impl IntegrationRetryState { + fn new(request_team_scope: RequestTeamScope) -> Self { + Self { + request_team_scope, + attempt: 1, + } + } - let list_future = async move { - integrations_client - .list_simple_integrations(provider_slugs) - .await - }; + fn next(self) -> Self { + Self { + attempt: self.attempt + 1, + ..self + } + } +} - ctx.spawn( - list_future, - move |_, result: anyhow::Result, ctx| match result { - Ok(output) => { - integration_output::print_integrations(&output, global_options.output_format); - ctx.terminate_app(TerminationMode::ForceTerminate, None); - } +impl IntegrationCommandRunner { + fn list( + &self, + global_options: GlobalOptions, + team_selection: TeamSelection, + ctx: &mut ModelContext, + ) { + let refresh_future = super::common::refresh_workspace_metadata(ctx); + ctx.spawn(refresh_future, move |_, result, ctx| { + if let Err(err) = result { + super::report_fatal_error(err, ctx); + return; + } + let team_scope = match super::common::resolve_team_scope(&team_selection, ctx) { + Ok(team_scope) => team_scope, Err(err) => { - ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); + super::report_fatal_error(err, ctx); + return; } - }, - ); + }; + let request_team_scope = RequestTeamScope::from_scope(&team_scope); + let provider_slugs = [ProviderType::Linear, ProviderType::Slack] + .into_iter() + .map(|provider| provider.slug()) + .collect(); + let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); + let list_future = async move { + integrations_client + .list_simple_integrations(request_team_scope, provider_slugs) + .await + }; + + ctx.spawn( + list_future, + move |_, result: anyhow::Result, ctx| match result { + Ok(output) => { + integration_output::print_integrations( + &output, + global_options.output_format, + ); + ctx.terminate_app(TerminationMode::ForceTerminate, None); + } + Err(err) => { + ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); + } + }, + ); + }); } fn create(&self, args: CreateIntegrationArgs, ctx: &mut ModelContext) { @@ -74,14 +120,14 @@ impl IntegrationCommandRunner { ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); return; } - let team_scope = - match super::common::resolve_team_scope(&TeamSelection { team: None }, ctx) { - Ok(team_scope) => team_scope, - Err(err) => { - ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); - return; - } - }; + let team_scope = match super::common::resolve_team_scope(&args.team_selection, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); + return; + } + }; + let request_team_scope = RequestTeamScope::from_scope(&team_scope); let loaded_file = match args.config_file.file.as_deref() { Some(path) => match super::config_file::load_config_file(path) { @@ -186,6 +232,7 @@ impl IntegrationCommandRunner { runner.start_create_or_update_flow( ctx, + IntegrationRetryState::new(request_team_scope), integration_type, environment_uid, base_prompt, @@ -195,7 +242,6 @@ impl IntegrationCommandRunner { worker_host, enabled, is_update, - 1, ); }); } @@ -204,6 +250,7 @@ impl IntegrationCommandRunner { fn start_create_or_update_flow( &self, ctx: &mut ModelContext, + retry_state: IntegrationRetryState, integration_type: String, environment_uid: Option, base_prompt: Option, @@ -213,12 +260,10 @@ impl IntegrationCommandRunner { worker_host: Option, enabled: bool, is_update: bool, - attempt: u32, ) { const MAX_CREATE_ATTEMPTS: u32 = 8; let action = if is_update { "update" } else { "creation" }; - - if attempt > MAX_CREATE_ATTEMPTS { + if retry_state.attempt > MAX_CREATE_ATTEMPTS { ctx.terminate_app( TerminationMode::ForceTerminate, Some(Err(anyhow::anyhow!( @@ -243,6 +288,7 @@ impl IntegrationCommandRunner { let create_future = async move { integrations_client .create_or_update_simple_integration( + retry_state.request_team_scope, future_integration_type, future_is_update, future_environment_uid, @@ -288,7 +334,6 @@ impl IntegrationCommandRunner { let next_worker_host = worker_host.clone(); let next_enabled = enabled; let next_is_update = is_update; - let next_attempt = attempt + 1; ctx.spawn( poll_future, @@ -299,6 +344,7 @@ impl IntegrationCommandRunner { // This may happen multiple times if the user needs to authorize multiple services. runner.start_create_or_update_flow( ctx, + retry_state.next(), next_integration_type, next_environment_uid, next_base_prompt, @@ -308,7 +354,6 @@ impl IntegrationCommandRunner { next_worker_host, next_enabled, next_is_update, - next_attempt, ); } Ok(OauthConnectTxStatus::Failed) => { @@ -395,6 +440,14 @@ impl IntegrationCommandRunner { ctx.terminate_app(TerminationMode::ForceTerminate, Some(Err(err))); return; } + let team_scope = match super::common::resolve_team_scope(&args.team_selection, ctx) { + Ok(team_scope) => team_scope, + Err(err) => { + super::report_fatal_error(err, ctx); + return; + } + }; + let request_team_scope = RequestTeamScope::from_scope(&team_scope); let loaded_file = match args.config_file.file.as_deref() { Some(path) => match super::config_file::load_config_file(path) { @@ -499,6 +552,7 @@ impl IntegrationCommandRunner { // Explicitly requested to update without an environment. runner.start_create_or_update_flow( ctx, + IntegrationRetryState::new(request_team_scope), integration_type, Some(String::new()), base_prompt, @@ -508,7 +562,6 @@ impl IntegrationCommandRunner { worker_host, enabled, is_update, - 1, ); return; } @@ -517,6 +570,7 @@ impl IntegrationCommandRunner { runner.start_create_or_update_flow( ctx, + IntegrationRetryState::new(request_team_scope), integration_type, environment_uid, base_prompt, @@ -526,7 +580,6 @@ impl IntegrationCommandRunner { worker_host, enabled, is_update, - 1, ); }); } @@ -536,3 +589,7 @@ impl warpui::Entity for IntegrationCommandRunner { type Event = (); } impl SingletonEntity for IntegrationCommandRunner {} + +#[cfg(test)] +#[path = "integration_tests.rs"] +mod tests; diff --git a/app/src/ai/agent_sdk/integration_tests.rs b/app/src/ai/agent_sdk/integration_tests.rs new file mode 100644 index 00000000000..9d764cbfb4a --- /dev/null +++ b/app/src/ai/agent_sdk/integration_tests.rs @@ -0,0 +1,27 @@ +use super::IntegrationRetryState; +use crate::server::ids::ServerId; +use crate::server::team_scope::RequestTeamScope; +use crate::workspaces::user_workspaces::{TeamContextForOperation, TeamlessScopeForTest}; + +#[test] +fn oauth_retry_retains_the_initiating_team_scope() { + let team_uid = ServerId::from(7); + let scope = TeamContextForOperation::new_for_test(team_uid); + let retry_state = IntegrationRetryState::new(RequestTeamScope::from_scope(&scope)); + + let retry_state = retry_state.next(); + + assert_eq!(retry_state.attempt, 2); + assert_eq!(retry_state.request_team_scope.team_uid(), Some(team_uid)); +} + +#[test] +fn oauth_retry_retains_a_teamless_initiating_scope() { + let retry_state = + IntegrationRetryState::new(RequestTeamScope::from_scope(&TeamlessScopeForTest)); + + let retry_state = retry_state.next(); + + assert_eq!(retry_state.attempt, 2); + assert_eq!(retry_state.request_team_scope.team_uid(), None); +} diff --git a/app/src/ai/agent_sdk/mod.rs b/app/src/ai/agent_sdk/mod.rs index fd6b905e86c..1d20c5c97ae 100644 --- a/app/src/ai/agent_sdk/mod.rs +++ b/app/src/ai/agent_sdk/mod.rs @@ -1876,7 +1876,7 @@ fn command_to_telemetry_event(command: &CliCommand) -> CliTelemetryEvent { CliCommand::Integration(integration_cmd) => match integration_cmd { IntegrationCommand::Create(_) => CliTelemetryEvent::IntegrationCreate, IntegrationCommand::Update(_) => CliTelemetryEvent::IntegrationUpdate, - IntegrationCommand::List => CliTelemetryEvent::IntegrationList, + IntegrationCommand::List { .. } => CliTelemetryEvent::IntegrationList, }, CliCommand::Schedule(c) => match c.subcommand() { None | Some(ScheduleSubcommand::Create(_)) => CliTelemetryEvent::ScheduleCreate, diff --git a/app/src/ai/agent_sdk/provider.rs b/app/src/ai/agent_sdk/provider.rs index 9c89d23f653..dd05eb23ac6 100644 --- a/app/src/ai/agent_sdk/provider.rs +++ b/app/src/ai/agent_sdk/provider.rs @@ -3,13 +3,14 @@ use comfy_table::Cell; use serde::Serialize; use warp_cli::GlobalOptions; use warp_cli::provider::{ProviderCommand, ProviderType}; +use warp_cli::scope::ObjectScope; use warp_core::channel::ChannelState; use warpui::platform::TerminationMode; use warpui::{AppContext, ModelContext, SingletonEntity}; use crate::ai::agent_sdk::common::describe_sole_team_error; use crate::ai::agent_sdk::output::{self, TableFormat}; -use crate::workspaces::user_workspaces::UserWorkspaces; +use crate::workspaces::user_workspaces::{SoleTeamError, TeamScope}; /// Handle provider-related CLI commands. pub fn run( @@ -20,7 +21,7 @@ pub fn run( let runner = ctx.add_singleton_model(|_ctx| ProviderCommandRunner); match command { ProviderCommand::Setup(args) => runner.update(ctx, |runner, ctx| { - runner.setup(args.provider_type, args.team, args.personal, ctx) + runner.setup(args.provider_type, args.scope, ctx) }), ProviderCommand::List => runner.update(ctx, |runner, ctx| runner.list(global_options, ctx)), } @@ -34,15 +35,14 @@ impl ProviderCommandRunner { fn setup( &self, provider_type: ProviderType, - team: bool, - personal: bool, + scope: ObjectScope, ctx: &mut ModelContext, ) -> anyhow::Result<()> { // Construct the OAuth connect URL let server_url = ChannelState::server_root_url(); - let mut use_team_auth = team; - if !team && !personal { + let mut use_team_auth = scope.is_team(); + if !scope.is_team() && !scope.personal { if provider_type.allowed_in_team_context() && provider_type.allowed_in_personal_context() { @@ -52,16 +52,17 @@ impl ProviderCommandRunner { )); } use_team_auth = provider_type.allowed_in_team_context(); - } else if personal { + } else if scope.personal { use_team_auth = false; } // TODO(bens): initiate the OAuth flow and use the login-less auth URL let slug = provider_type.slug(); let url = if use_team_auth { - let team_uid = UserWorkspaces::as_ref(ctx) - .sole_team_uid() - .map_err(|err| describe_sole_team_error(err, ctx))?; + let team_scope = super::common::resolve_team_scope(&scope.team_selection, ctx)?; + let team_uid = team_scope + .team_uid() + .ok_or_else(|| describe_sole_team_error(SoleTeamError::NoTeam, ctx))?; format!("{server_url}/oauth/connect/{slug}?principalType=team&principalId={team_uid}") } else { format!("{server_url}/oauth/connect/{slug}") diff --git a/app/src/server/server_api.rs b/app/src/server/server_api.rs index 5d5821597d2..389bc88b10f 100644 --- a/app/src/server/server_api.rs +++ b/app/src/server/server_api.rs @@ -566,6 +566,26 @@ impl ServerApi { timeout, ) } + pub fn send_team_scoped_graphql_request< + 'a, + QF, + O: warp_graphql::client::Operation + Send + 'a, + >( + &'a self, + operation: O, + timeout: Option, + team_scope: RequestTeamScope, + ) -> BoxFuture<'a, Result> + where + QF: 'a, + { + warp_server_client::graphql_helpers::send_team_scoped_graphql_request( + &self.base_client, + operation, + timeout, + team_scope.team_uid().map(|team_uid| team_uid.to_string()), + ) + } /// Opens an SSE stream to the agent event-push endpoint. /// diff --git a/app/src/server/server_api/integrations.rs b/app/src/server/server_api/integrations.rs index b85abac1bdf..e48f6c2094b 100644 --- a/app/src/server/server_api/integrations.rs +++ b/app/src/server/server_api/integrations.rs @@ -37,6 +37,7 @@ use super::ServerApi; use crate::channel::ChannelState; use crate::features::FeatureFlag; use crate::server::graphql::{get_request_context, get_user_facing_error_message}; +use crate::server::team_scope::RequestTeamScope; #[cfg(not(target_family = "wasm"))] pub trait IntegrationsClientBounds: Send + Sync {} @@ -79,6 +80,7 @@ pub trait IntegrationsClient: 'static + IntegrationsClientBounds { #[allow(clippy::too_many_arguments)] async fn create_or_update_simple_integration( &self, + team_scope: RequestTeamScope, integration_type: String, is_update: bool, environment_uid: Option, @@ -96,6 +98,7 @@ pub trait IntegrationsClient: 'static + IntegrationsClientBounds { /// regardless of whether the connection or integration currently exists. async fn list_simple_integrations( &self, + team_scope: RequestTeamScope, providers: Vec, ) -> Result; @@ -167,6 +170,7 @@ impl IntegrationsClient for ServerApi { #[allow(clippy::too_many_arguments)] async fn create_or_update_simple_integration( &self, + team_scope: RequestTeamScope, integration_type: String, is_update: bool, environment_uid: Option, @@ -193,7 +197,9 @@ impl IntegrationsClient for ServerApi { }; let operation = CreateSimpleIntegration::build(variables); - let response = self.send_graphql_request(operation, None).await?; + let response = self + .send_team_scoped_graphql_request(operation, None, team_scope) + .await?; match response.create_simple_integration { CreateSimpleIntegrationResult::CreateSimpleIntegrationOutput(output) => Ok(output), CreateSimpleIntegrationResult::UserFacingError(error) => { @@ -232,6 +238,7 @@ impl IntegrationsClient for ServerApi { async fn list_simple_integrations( &self, + team_scope: RequestTeamScope, providers: Vec, ) -> Result { let variables = SimpleIntegrationsVariables { @@ -240,7 +247,9 @@ impl IntegrationsClient for ServerApi { }; let operation = SimpleIntegrations::build(variables); - let response = self.send_graphql_request(operation, None).await?; + let response = self + .send_team_scoped_graphql_request(operation, None, team_scope) + .await?; match response.simple_integrations { SimpleIntegrationsResult::SimpleIntegrationsOutput(output) => Ok(output), diff --git a/crates/warp_cli/src/integration.rs b/crates/warp_cli/src/integration.rs index bcf95985d19..8075394f97a 100644 --- a/crates/warp_cli/src/integration.rs +++ b/crates/warp_cli/src/integration.rs @@ -5,6 +5,7 @@ use crate::environment::{EnvironmentCreateArgs, EnvironmentUpdateArgs}; use crate::mcp::MCPSpec; use crate::model::ModelArgs; use crate::provider::ProviderType; +use crate::scope::TeamSelection; /// Integration-related subcommands. #[derive(Debug, Clone, Subcommand)] @@ -15,7 +16,10 @@ pub enum IntegrationCommand { /// Update an integration. Update(UpdateIntegrationArgs), /// List simple integrations and their connection status. - List, + List { + #[command(flatten)] + team_selection: TeamSelection, + }, } impl IntegrationCommand { @@ -23,7 +27,7 @@ impl IntegrationCommand { match self { IntegrationCommand::Create(_) => "integration create", IntegrationCommand::Update(_) => "integration update", - IntegrationCommand::List => "integration list", + IntegrationCommand::List { .. } => "integration list", } } } @@ -33,6 +37,8 @@ pub struct CreateIntegrationArgs { /// Provider to create the integration for. #[arg(value_enum)] pub provider: ProviderType, + #[command(flatten)] + pub team_selection: TeamSelection, #[command(flatten)] pub model: ModelArgs, @@ -68,6 +74,8 @@ pub struct UpdateIntegrationArgs { /// Provider to update the integration for. #[arg(value_enum)] pub provider: ProviderType, + #[command(flatten)] + pub team_selection: TeamSelection, #[command(flatten)] pub model: ModelArgs, diff --git a/crates/warp_cli/src/lib_tests.rs b/crates/warp_cli/src/lib_tests.rs index b249413f56d..213a066c074 100644 --- a/crates/warp_cli/src/lib_tests.rs +++ b/crates/warp_cli/src/lib_tests.rs @@ -9,6 +9,7 @@ use crate::environment::{EnvironmentCommand, ImageCommand}; use crate::harness_support::{HarnessSupportCommand, TaskStatus}; use crate::integration::IntegrationCommand; use crate::memory_store::{MemoryCommand, MemoryStoreCommand}; +use crate::provider::ProviderCommand; use crate::schedule::ScheduleSubcommand; use crate::secret::{CodexMethod, CreateProvider, SecretCommand}; use crate::task::{MessageCommand, TaskCommand}; @@ -23,6 +24,106 @@ fn identifies_worker_subcommands() { assert!(!is_worker_invocation("--prompt")); } +#[test] +fn integration_commands_parse_team_selection() { + let cases = [ + ( + vec!["warp", "integration", "list", "--team=team-list"], + "team-list", + ), + ( + vec![ + "warp", + "integration", + "create", + "slack", + "--team=team-create", + ], + "team-create", + ), + ( + vec![ + "warp", + "integration", + "update", + "linear", + "--team=team-update", + ], + "team-update", + ), + ]; + + for (command, expected_team_uid) in cases { + let args = Args::try_parse_from(command).expect("integration team scope should parse"); + let Some(Command::CommandLine(boxed_cmd)) = args.command else { + panic!("Expected an integration command"); + }; + let team_selection = match boxed_cmd.as_ref() { + CliCommand::Integration(IntegrationCommand::List { team_selection }) => team_selection, + CliCommand::Integration(IntegrationCommand::Create(args)) => &args.team_selection, + CliCommand::Integration(IntegrationCommand::Update(args)) => &args.team_selection, + _ => panic!("Expected an integration command"), + }; + + assert_eq!(team_selection.requested_team_uid(), Some(expected_team_uid)); + } +} + +#[test] +fn integration_list_distinguishes_default_and_bare_team_selection() { + for (command, expected_team) in [ + (vec!["warp", "integration", "list"], None), + (vec!["warp", "integration", "list", "--team"], Some(None)), + ] { + let args = Args::try_parse_from(command).expect("integration list scope should parse"); + let Some(Command::CommandLine(boxed_cmd)) = args.command else { + panic!("Expected `warp integration list` command"); + }; + let CliCommand::Integration(IntegrationCommand::List { team_selection }) = + boxed_cmd.as_ref() + else { + panic!("Expected `warp integration list` command"); + }; + + assert_eq!(team_selection.team, expected_team); + } +} + +#[test] +fn provider_setup_parses_object_scope() { + let team = Args::try_parse_from(["warp", "provider", "setup", "slack", "--team=team-provider"]) + .expect("provider team scope should parse"); + let Some(Command::CommandLine(boxed_cmd)) = team.command else { + panic!("Expected `warp provider setup` command"); + }; + let CliCommand::Provider(ProviderCommand::Setup(args)) = boxed_cmd.as_ref() else { + panic!("Expected `warp provider setup` command"); + }; + assert_eq!(args.scope.requested_team_uid(), Some("team-provider")); + assert!(!args.scope.personal); + + let personal = Args::try_parse_from(["warp", "provider", "setup", "slack", "--personal"]) + .expect("provider personal scope should parse"); + let Some(Command::CommandLine(boxed_cmd)) = personal.command else { + panic!("Expected `warp provider setup` command"); + }; + let CliCommand::Provider(ProviderCommand::Setup(args)) = boxed_cmd.as_ref() else { + panic!("Expected `warp provider setup` command"); + }; + assert!(args.scope.personal); + assert!(!args.scope.is_team()); + + Args::try_parse_from([ + "warp", + "provider", + "setup", + "slack", + "--personal", + "--team=team-provider", + ]) + .expect_err("provider scopes should be mutually exclusive"); +} + /// Pins that each pair of constants names the same variable under both prefixes. A typo in /// either half would otherwise go unnoticed until a consumer read the wrong name. #[test] diff --git a/crates/warp_cli/src/provider.rs b/crates/warp_cli/src/provider.rs index caebd29b7ef..bd3fe6b9762 100644 --- a/crates/warp_cli/src/provider.rs +++ b/crates/warp_cli/src/provider.rs @@ -1,4 +1,6 @@ -use clap::{ArgGroup, Args, Subcommand, ValueEnum}; +use clap::{Args, Subcommand, ValueEnum}; + +use crate::scope::ObjectScope; /// Provider-related subcommands. #[derive(Debug, Clone, Subcommand)] @@ -53,15 +55,10 @@ impl ProviderType { } #[derive(Debug, Clone, Args)] -#[command(group(ArgGroup::new("scope").required(false)))] pub struct SetupArgs { /// The type of provider to setup. pub provider_type: ProviderType, - /// Setup provider for a team - #[arg(long, group = "scope")] - pub team: bool, - /// Setup provider for a personal account - #[arg(long, conflicts_with = "team", group = "scope")] - pub personal: bool, + #[command(flatten)] + pub scope: ObjectScope, } diff --git a/crates/warp_server_client/src/graphql_helpers.rs b/crates/warp_server_client/src/graphql_helpers.rs index a5f0e09d026..b81ea96667d 100644 --- a/crates/warp_server_client/src/graphql_helpers.rs +++ b/crates/warp_server_client/src/graphql_helpers.rs @@ -8,7 +8,7 @@ use warp_graphql::client::{GraphQLError, Operation, RequestOptions}; use warpui_core::r#async::BoxFuture; use crate::auth::AuthEvent; -use crate::base_client::BaseClient; +use crate::base_client::{BaseClient, TEAM_UID_HEADER}; /// Sends a GraphQL operation through a base client supplied by the application. /// @@ -28,6 +28,35 @@ where }) } +fn apply_request_team_scope( + mut options: RequestOptions, + team_uid: Option, +) -> RequestOptions { + if let Some(team_uid) = team_uid { + options + .headers + .insert(TEAM_UID_HEADER.to_string(), team_uid); + } + options +} +pub fn send_team_scoped_graphql_request<'a, QF: 'a, O>( + base_client: &'a BaseClient, + operation: O, + timeout: Option, + team_uid: Option, +) -> BoxFuture<'a, Result> +where + O: Operation + Send + 'a, +{ + Box::pin(async move { + let options = apply_request_team_scope( + base_client.graphql_request_options(timeout).await?, + team_uid, + ); + send_graphql_request_with_options(base_client, operation, options).await + }) +} + async fn send_graphql_request_with_options( base_client: &BaseClient, operation: O, diff --git a/crates/warp_server_client/src/graphql_helpers_tests.rs b/crates/warp_server_client/src/graphql_helpers_tests.rs index 1607ff6ee37..199bbe32921 100644 --- a/crates/warp_server_client/src/graphql_helpers_tests.rs +++ b/crates/warp_server_client/src/graphql_helpers_tests.rs @@ -10,9 +10,11 @@ use http::StatusCode; use warp_graphql::client::{GraphQLError, RequestOptions}; use warp_server_auth::auth_state::AuthState; -use super::send_graphql_request; +use super::{apply_request_team_scope, send_graphql_request}; use crate::auth::AuthEvent; -use crate::base_client::{AuthenticatedGraphqlConfig, BaseClient, GraphqlRoutingConfig}; +use crate::base_client::{ + AuthenticatedGraphqlConfig, BaseClient, GraphqlRoutingConfig, TEAM_UID_HEADER, +}; fn base_client(auth_state: AuthState) -> (BaseClient, async_channel::Receiver) { let (event_sender, event_receiver) = async_channel::unbounded(); @@ -30,6 +32,29 @@ fn base_client(auth_state: AuthState) -> (BaseClient, async_channel::Receiver Date: Fri, 4 Sep 2026 20:08:51 +0000 Subject: [PATCH 4/5] [REV-2383] Verify scoped integration retries --- app/src/ai/agent_sdk/integration.rs | 28 +++-- app/src/ai/agent_sdk/integration_tests.rs | 110 +++++++++++++++--- app/src/ai/agent_sdk/provider.rs | 40 ++++++- app/src/ai/agent_sdk/provider_tests.rs | 75 ++++++++++++ crates/warp_server_client/src/base_client.rs | 8 ++ .../src/graphql_helpers_tests.rs | 92 ++++++++++----- 6 files changed, 296 insertions(+), 57 deletions(-) create mode 100644 app/src/ai/agent_sdk/provider_tests.rs diff --git a/app/src/ai/agent_sdk/integration.rs b/app/src/ai/agent_sdk/integration.rs index 944056532ce..79c37389e48 100644 --- a/app/src/ai/agent_sdk/integration.rs +++ b/app/src/ai/agent_sdk/integration.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use futures::future; use warp_cli::GlobalOptions; use warp_cli::integration::{CreateIntegrationArgs, IntegrationCommand, UpdateIntegrationArgs}; @@ -13,6 +15,7 @@ use super::common::{EnvironmentChoice, ResolveConfigurationError}; use super::integration_output; use super::oauth_flow::poll_oauth_until_terminal; use crate::server::server_api::ServerApiProvider; +use crate::server::server_api::integrations::IntegrationsClient; use crate::server::team_scope::RequestTeamScope; pub fn run( @@ -20,7 +23,9 @@ pub fn run( global_options: GlobalOptions, command: IntegrationCommand, ) -> anyhow::Result<()> { - let runner = ctx.add_singleton_model(|_ctx| IntegrationCommandRunner); + let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); + let runner = + ctx.add_singleton_model(move |_ctx| IntegrationCommandRunner::new(integrations_client)); match command { IntegrationCommand::Create(args) => { runner.update(ctx, |runner, ctx| runner.create(args, ctx)); @@ -37,7 +42,9 @@ pub fn run( Ok(()) } -struct IntegrationCommandRunner; +struct IntegrationCommandRunner { + integrations_client: Arc, +} #[derive(Clone, Copy, Debug, PartialEq, Eq)] struct IntegrationRetryState { request_team_scope: RequestTeamScope, @@ -61,6 +68,12 @@ impl IntegrationRetryState { } impl IntegrationCommandRunner { + fn new(integrations_client: Arc) -> Self { + Self { + integrations_client, + } + } + fn list( &self, global_options: GlobalOptions, @@ -68,7 +81,7 @@ impl IntegrationCommandRunner { ctx: &mut ModelContext, ) { let refresh_future = super::common::refresh_workspace_metadata(ctx); - ctx.spawn(refresh_future, move |_, result, ctx| { + ctx.spawn(refresh_future, move |runner, result, ctx| { if let Err(err) = result { super::report_fatal_error(err, ctx); return; @@ -85,7 +98,7 @@ impl IntegrationCommandRunner { .into_iter() .map(|provider| provider.slug()) .collect(); - let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); + let integrations_client = runner.integrations_client.clone(); let list_future = async move { integrations_client .list_simple_integrations(request_team_scope, provider_slugs) @@ -274,7 +287,7 @@ impl IntegrationCommandRunner { return; } - let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); + let integrations_client = self.integrations_client.clone(); let future_integration_type = integration_type.clone(); let future_environment_uid = environment_uid.clone(); @@ -304,7 +317,7 @@ impl IntegrationCommandRunner { ctx.spawn( create_future, - move |_runner, result: anyhow::Result, ctx| { + move |runner, result: anyhow::Result, ctx| { match result { Ok(output) => { println!("{}", output.message); @@ -318,8 +331,7 @@ impl IntegrationCommandRunner { println!("Authorize the provider here: {auth_url}\n"); ctx.open_url(&auth_url); - let integrations_client = ServerApiProvider::as_ref(ctx) - .get_integrations_client(); + let integrations_client = runner.integrations_client.clone(); let tx_id = tx_id.into_inner(); let poll_future = diff --git a/app/src/ai/agent_sdk/integration_tests.rs b/app/src/ai/agent_sdk/integration_tests.rs index 9d764cbfb4a..d669fd7b8ca 100644 --- a/app/src/ai/agent_sdk/integration_tests.rs +++ b/app/src/ai/agent_sdk/integration_tests.rs @@ -1,27 +1,105 @@ -use super::IntegrationRetryState; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::Duration; + +use warp_graphql::mutations::create_simple_integration::CreateSimpleIntegrationOutput; +use warp_graphql::queries::get_oauth_connect_tx_status::OauthConnectTxStatus; +use warpui::App; +use warpui::r#async::Timer; + +use super::{IntegrationCommandRunner, IntegrationRetryState}; use crate::server::ids::ServerId; +use crate::server::server_api::integrations::MockIntegrationsClient; use crate::server::team_scope::RequestTeamScope; use crate::workspaces::user_workspaces::{TeamContextForOperation, TeamlessScopeForTest}; #[test] -fn oauth_retry_retains_the_initiating_team_scope() { - let team_uid = ServerId::from(7); - let scope = TeamContextForOperation::new_for_test(team_uid); - let retry_state = IntegrationRetryState::new(RequestTeamScope::from_scope(&scope)); - - let retry_state = retry_state.next(); - - assert_eq!(retry_state.attempt, 2); - assert_eq!(retry_state.request_team_scope.team_uid(), Some(team_uid)); +fn oauth_retry_uses_the_initiating_team_scope() { + let scope = TeamContextForOperation::new_for_test(ServerId::from(7)); + run_oauth_retry(RequestTeamScope::from_scope(&scope), false); } #[test] -fn oauth_retry_retains_a_teamless_initiating_scope() { - let retry_state = - IntegrationRetryState::new(RequestTeamScope::from_scope(&TeamlessScopeForTest)); +fn oauth_retry_uses_the_initiating_teamless_scope() { + run_oauth_retry(RequestTeamScope::from_scope(&TeamlessScopeForTest), true); +} + +fn run_oauth_retry(expected_scope: RequestTeamScope, is_update: bool) { + App::test((), move |mut app| async move { + let request_count = Arc::new(AtomicUsize::new(0)); + let observed_request_count = request_count.clone(); + let mut integrations_client = MockIntegrationsClient::new(); + integrations_client + .expect_create_or_update_simple_integration() + .times(2) + .returning( + move |request_scope, + integration_type, + request_is_update, + _environment_uid, + _base_prompt, + _model_id, + _mcp_servers_json, + _remove_mcp_server_names, + _worker_host, + enabled| { + assert_eq!(request_scope, expected_scope); + assert_eq!(integration_type, "slack"); + assert_eq!(request_is_update, is_update); + assert!(enabled); + let request_index = observed_request_count.fetch_add(1, Ordering::SeqCst); + Ok(if request_index == 0 { + CreateSimpleIntegrationOutput { + auth_url: Some("https://example.com/oauth".to_string()), + success: false, + message: "Authorization required".to_string(), + tx_id: Some(cynic::Id::new("oauth-tx")), + } + } else { + CreateSimpleIntegrationOutput { + auth_url: None, + success: true, + message: "Integration saved".to_string(), + tx_id: None, + } + }) + }, + ); + integrations_client + .expect_poll_oauth_connect_status() + .times(1) + .returning(|tx_id| { + assert_eq!(tx_id, "oauth-tx"); + Ok(OauthConnectTxStatus::Completed) + }); + + let runner = + app.add_model(move |_| IntegrationCommandRunner::new(Arc::new(integrations_client))); + runner.update(&mut app, |runner, ctx| { + runner.start_create_or_update_flow( + ctx, + IntegrationRetryState::new(expected_scope), + "slack".to_string(), + None, + None, + None, + None, + None, + None, + true, + is_update, + ); + }); - let retry_state = retry_state.next(); + for _ in 0..120 { + if request_count.load(Ordering::SeqCst) == 2 { + break; + } + Timer::after(Duration::from_millis(50)).await; + } - assert_eq!(retry_state.attempt, 2); - assert_eq!(retry_state.request_team_scope.team_uid(), None); + assert_eq!(request_count.load(Ordering::SeqCst), 2); + futures_lite::future::yield_now().await; + assert!(app.termination_result().is_none()); + }); } diff --git a/app/src/ai/agent_sdk/provider.rs b/app/src/ai/agent_sdk/provider.rs index dd05eb23ac6..bbc1bbb9454 100644 --- a/app/src/ai/agent_sdk/provider.rs +++ b/app/src/ai/agent_sdk/provider.rs @@ -1,4 +1,6 @@ //! Provider command for linking third-party services. +use std::future::Future; + use comfy_table::Cell; use serde::Serialize; use warp_cli::GlobalOptions; @@ -20,9 +22,12 @@ pub fn run( ) -> anyhow::Result<()> { let runner = ctx.add_singleton_model(|_ctx| ProviderCommandRunner); match command { - ProviderCommand::Setup(args) => runner.update(ctx, |runner, ctx| { - runner.setup(args.provider_type, args.scope, ctx) - }), + ProviderCommand::Setup(args) => { + runner.update(ctx, |runner, ctx| { + runner.setup(args.provider_type, args.scope, ctx) + }); + Ok(()) + } ProviderCommand::List => runner.update(ctx, |runner, ctx| runner.list(global_options, ctx)), } } @@ -32,7 +37,30 @@ struct ProviderCommandRunner; impl ProviderCommandRunner { // This shouldn't need to be done, it's usually done as part of create - fn setup( + fn setup(&self, provider_type: ProviderType, scope: ObjectScope, ctx: &mut ModelContext) { + let refresh_future = super::common::refresh_workspace_metadata(ctx); + self.setup_after_workspace_metadata_refresh(refresh_future, provider_type, scope, ctx); + } + + fn setup_after_workspace_metadata_refresh( + &self, + refresh_future: impl Future> + Send + 'static, + provider_type: ProviderType, + scope: ObjectScope, + ctx: &mut ModelContext, + ) { + ctx.spawn(refresh_future, move |runner, result, ctx| { + if let Err(err) = result { + super::report_fatal_error(err, ctx); + return; + } + if let Err(err) = runner.finish_setup(provider_type, scope, ctx) { + super::report_fatal_error(err, ctx); + } + }); + } + + fn finish_setup( &self, provider_type: ProviderType, scope: ObjectScope, @@ -154,3 +182,7 @@ impl TableFormat for ProviderInfo { ] } } + +#[cfg(test)] +#[path = "provider_tests.rs"] +mod tests; diff --git a/app/src/ai/agent_sdk/provider_tests.rs b/app/src/ai/agent_sdk/provider_tests.rs new file mode 100644 index 00000000000..3af1648416b --- /dev/null +++ b/app/src/ai/agent_sdk/provider_tests.rs @@ -0,0 +1,75 @@ +use std::sync::{Arc, Mutex}; + +use warp_cli::provider::ProviderType; +use warp_cli::scope::{ObjectScope, TeamSelection}; +use warpui::App; +use warpui::r#async::Timer; + +use super::ProviderCommandRunner; +use crate::server::ids::ServerId; +use crate::settings::PrivacySettings; +use crate::workspaces::user_workspaces::UserWorkspaces; + +#[test] +fn setup_resolves_team_scope_after_workspace_metadata_refresh() { + App::test((), |mut app| async move { + app.add_singleton_model(PrivacySettings::mock); + let user_workspaces = app.add_singleton_model(UserWorkspaces::default_mock); + let opened_urls = Arc::new(Mutex::new(Vec::new())); + let captured_urls = opened_urls.clone(); + app.update(|ctx| { + ctx.set_before_open_url(move |url, _| { + captured_urls.lock().unwrap().push(url.to_string()); + url.to_string() + }); + }); + + let (refresh_sender, refresh_receiver) = async_channel::bounded(1); + let runner = app.add_model(|_| ProviderCommandRunner); + runner.update(&mut app, |runner, ctx| { + runner.setup_after_workspace_metadata_refresh( + async move { + refresh_receiver.recv().await.unwrap(); + Ok(()) + }, + ProviderType::Slack, + ObjectScope { + team_selection: TeamSelection { team: Some(None) }, + personal: false, + }, + ctx, + ); + }); + + for _ in 0..3 { + futures_lite::future::yield_now().await; + } + assert!(opened_urls.lock().unwrap().is_empty()); + + user_workspaces.update(&mut app, |workspaces, ctx| { + workspaces.setup_test_workspace(ctx); + }); + refresh_sender.send(()).await.unwrap(); + + for _ in 0..20 { + if !opened_urls.lock().unwrap().is_empty() { + break; + } + Timer::after(std::time::Duration::from_millis(10)).await; + } + + let opened_urls = opened_urls.lock().unwrap(); + assert_eq!(opened_urls.len(), 1); + let expected_url_suffix = format!( + "/oauth/connect/slack?principalType=team&principalId={}", + ServerId::from(2) + ); + assert!( + opened_urls[0].ends_with(&expected_url_suffix), + "opened URL: {}", + opened_urls[0] + ); + drop(opened_urls); + assert!(app.termination_result().is_none()); + }); +} diff --git a/crates/warp_server_client/src/base_client.rs b/crates/warp_server_client/src/base_client.rs index 99aa70e90a7..ac00b50ddc4 100644 --- a/crates/warp_server_client/src/base_client.rs +++ b/crates/warp_server_client/src/base_client.rs @@ -253,6 +253,14 @@ impl BaseClient { *self.ambient_agent_task_id.write() = task_id; } + #[cfg(test)] + pub(crate) fn set_ambient_workload_token_for_test(&self, token: &str) { + *self.ambient_workload_token.lock() = Some(warp_isolation_platform::WorkloadToken { + token: token.to_string(), + expires_at: None, + }); + } + /// Returns an ambient agent workload token when the current runtime can issue one. pub async fn get_or_create_ambient_workload_token(&self) -> Result> { if cfg!(target_family = "wasm") { diff --git a/crates/warp_server_client/src/graphql_helpers_tests.rs b/crates/warp_server_client/src/graphql_helpers_tests.rs index 199bbe32921..e0cc3a3a1b0 100644 --- a/crates/warp_server_client/src/graphql_helpers_tests.rs +++ b/crates/warp_server_client/src/graphql_helpers_tests.rs @@ -10,7 +10,7 @@ use http::StatusCode; use warp_graphql::client::{GraphQLError, RequestOptions}; use warp_server_auth::auth_state::AuthState; -use super::{apply_request_team_scope, send_graphql_request}; +use super::{send_graphql_request, send_team_scoped_graphql_request}; use crate::auth::AuthEvent; use crate::base_client::{ AuthenticatedGraphqlConfig, BaseClient, GraphqlRoutingConfig, TEAM_UID_HEADER, @@ -18,41 +18,55 @@ use crate::base_client::{ fn base_client(auth_state: AuthState) -> (BaseClient, async_channel::Receiver) { let (event_sender, event_receiver) = async_channel::unbounded(); - ( - BaseClient::new( - Arc::new(http_client::Client::new()), - Arc::new(auth_state), - event_sender, - None, - GraphqlRoutingConfig::default(), - AuthenticatedGraphqlConfig::default(), - None, - ), - event_receiver, - ) + let base_client = BaseClient::new( + Arc::new(http_client::Client::new()), + Arc::new(auth_state), + event_sender, + None, + GraphqlRoutingConfig::default(), + AuthenticatedGraphqlConfig::default(), + None, + ); + base_client.set_ambient_workload_token_for_test("test-workload-token"); + (base_client, event_receiver) } #[test] -fn team_scoped_request_options_preserve_authentication_and_set_team_header() { - let options = RequestOptions { - auth_token: Some("daemon-token".to_string()), - ..RequestOptions::default() - }; - - let options = apply_request_team_scope(options, Some("team-uid-123".to_string())); - - assert_eq!(options.auth_token.as_deref(), Some("daemon-token")); - assert_eq!( - options.headers.get(TEAM_UID_HEADER).map(String::as_str), - Some("team-uid-123") - ); +fn team_scoped_graphql_request_preserves_authentication_and_sets_team_header() { + let (base_client, event_receiver) = externally_authenticated_base_client("daemon-token"); + let send_count = Arc::new(AtomicUsize::new(0)); + + block_on(send_team_scoped_graphql_request( + &base_client, + FakeGraphqlOperation::successful_with_team( + Some("daemon-token"), + Some("team-uid-123"), + send_count.clone(), + ), + None, + Some("team-uid-123".to_string()), + )) + .unwrap(); + + assert_eq!(send_count.load(Ordering::SeqCst), 1); + assert_no_events(&event_receiver); } #[test] -fn teamless_request_options_omit_team_header() { - let options = apply_request_team_scope(RequestOptions::default(), None); +fn teamless_graphql_request_preserves_authentication_and_omits_team_header() { + let (base_client, event_receiver) = externally_authenticated_base_client("daemon-token"); + let send_count = Arc::new(AtomicUsize::new(0)); + + block_on(send_team_scoped_graphql_request( + &base_client, + FakeGraphqlOperation::successful_with_team(Some("daemon-token"), None, send_count.clone()), + None, + None, + )) + .unwrap(); - assert!(!options.headers.contains_key(TEAM_UID_HEADER)); + assert_eq!(send_count.load(Ordering::SeqCst), 1); + assert_no_events(&event_receiver); } #[test] @@ -105,6 +119,7 @@ fn assert_user_disabled_event(event_receiver: &async_channel::Receiver, + expected_team_uid: Option, send_count: Arc, result: FakeGraphqlResult, } @@ -119,6 +134,19 @@ impl FakeGraphqlOperation { fn successful(expected_auth_token: Option<&str>, send_count: Arc) -> Self { Self { expected_auth_token: expected_auth_token.map(ToOwned::to_owned), + expected_team_uid: None, + send_count, + result: FakeGraphqlResult::Success, + } + } + fn successful_with_team( + expected_auth_token: Option<&str>, + expected_team_uid: Option<&str>, + send_count: Arc, + ) -> Self { + Self { + expected_auth_token: expected_auth_token.map(ToOwned::to_owned), + expected_team_uid: expected_team_uid.map(ToOwned::to_owned), send_count, result: FakeGraphqlResult::Success, } @@ -131,6 +159,7 @@ impl FakeGraphqlOperation { ) -> Self { Self { expected_auth_token: expected_auth_token.map(ToOwned::to_owned), + expected_team_uid: None, send_count, result: FakeGraphqlResult::Rejected(status), } @@ -143,6 +172,7 @@ impl FakeGraphqlOperation { ) -> Self { Self { expected_auth_token: expected_auth_token.map(ToOwned::to_owned), + expected_team_uid: None, send_count, result: FakeGraphqlResult::ResponseErrors(messages), } @@ -170,6 +200,10 @@ impl warp_graphql::client::Operation<()> for FakeGraphqlOperation { { Box::pin(async move { assert_eq!(options.auth_token, self.expected_auth_token); + assert_eq!( + options.headers.get(TEAM_UID_HEADER), + self.expected_team_uid.as_ref() + ); self.send_count.fetch_add(1, Ordering::SeqCst); match self.result { FakeGraphqlResult::Success => Ok(GraphQlResponse { From 048a19a5133392bc6042bd34c20643d4da50fa14 Mon Sep 17 00:00:00 2001 From: "warp-agent-staging[bot]" <240773466+warp-agent-staging[bot]@users.noreply.github.com> Date: Fri, 4 Sep 2026 20:56:21 +0000 Subject: [PATCH 5/5] [REV-2383] Exercise production scope boundaries --- app/src/ai/agent_sdk/integration.rs | 73 ++++++--------- app/src/ai/agent_sdk/integration_tests.rs | 101 +++------------------ app/src/ai/agent_sdk/provider.rs | 11 --- app/src/ai/agent_sdk/provider_tests.rs | 104 +++++++++++++--------- 4 files changed, 100 insertions(+), 189 deletions(-) diff --git a/app/src/ai/agent_sdk/integration.rs b/app/src/ai/agent_sdk/integration.rs index 79c37389e48..7c65cad9b6b 100644 --- a/app/src/ai/agent_sdk/integration.rs +++ b/app/src/ai/agent_sdk/integration.rs @@ -1,5 +1,3 @@ -use std::sync::Arc; - use futures::future; use warp_cli::GlobalOptions; use warp_cli::integration::{CreateIntegrationArgs, IntegrationCommand, UpdateIntegrationArgs}; @@ -15,7 +13,6 @@ use super::common::{EnvironmentChoice, ResolveConfigurationError}; use super::integration_output; use super::oauth_flow::poll_oauth_until_terminal; use crate::server::server_api::ServerApiProvider; -use crate::server::server_api::integrations::IntegrationsClient; use crate::server::team_scope::RequestTeamScope; pub fn run( @@ -23,9 +20,7 @@ pub fn run( global_options: GlobalOptions, command: IntegrationCommand, ) -> anyhow::Result<()> { - let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); - let runner = - ctx.add_singleton_model(move |_ctx| IntegrationCommandRunner::new(integrations_client)); + let runner = ctx.add_singleton_model(|_ctx| IntegrationCommandRunner); match command { IntegrationCommand::Create(args) => { runner.update(ctx, |runner, ctx| runner.create(args, ctx)); @@ -42,9 +37,7 @@ pub fn run( Ok(()) } -struct IntegrationCommandRunner { - integrations_client: Arc, -} +struct IntegrationCommandRunner; #[derive(Clone, Copy, Debug, PartialEq, Eq)] struct IntegrationRetryState { request_team_scope: RequestTeamScope, @@ -65,15 +58,26 @@ impl IntegrationRetryState { ..self } } -} -impl IntegrationCommandRunner { - fn new(integrations_client: Arc) -> Self { - Self { - integrations_client, + fn continue_after_oauth( + self, + poll_result: anyhow::Result, + ) -> anyhow::Result { + match poll_result { + Ok(OauthConnectTxStatus::Completed) => Ok(self.next()), + Ok(OauthConnectTxStatus::Failed) => Err(anyhow::anyhow!("OAuth authorization failed.")), + Ok(OauthConnectTxStatus::Expired) => { + Err(anyhow::anyhow!("OAuth authorization expired.")) + } + Ok(OauthConnectTxStatus::Pending) | Ok(OauthConnectTxStatus::InProgress) => Err( + anyhow::anyhow!("Unexpected non-terminal OAuth status returned"), + ), + Err(err) => Err(anyhow::anyhow!("Error polling OAuth status: {err}")), } } +} +impl IntegrationCommandRunner { fn list( &self, global_options: GlobalOptions, @@ -81,7 +85,7 @@ impl IntegrationCommandRunner { ctx: &mut ModelContext, ) { let refresh_future = super::common::refresh_workspace_metadata(ctx); - ctx.spawn(refresh_future, move |runner, result, ctx| { + ctx.spawn(refresh_future, move |_, result, ctx| { if let Err(err) = result { super::report_fatal_error(err, ctx); return; @@ -98,7 +102,7 @@ impl IntegrationCommandRunner { .into_iter() .map(|provider| provider.slug()) .collect(); - let integrations_client = runner.integrations_client.clone(); + let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); let list_future = async move { integrations_client .list_simple_integrations(request_team_scope, provider_slugs) @@ -287,7 +291,7 @@ impl IntegrationCommandRunner { return; } - let integrations_client = self.integrations_client.clone(); + let integrations_client = ServerApiProvider::as_ref(ctx).get_integrations_client(); let future_integration_type = integration_type.clone(); let future_environment_uid = environment_uid.clone(); @@ -317,7 +321,7 @@ impl IntegrationCommandRunner { ctx.spawn( create_future, - move |runner, result: anyhow::Result, ctx| { + move |_runner, result: anyhow::Result, ctx| { match result { Ok(output) => { println!("{}", output.message); @@ -331,7 +335,8 @@ impl IntegrationCommandRunner { println!("Authorize the provider here: {auth_url}\n"); ctx.open_url(&auth_url); - let integrations_client = runner.integrations_client.clone(); + let integrations_client = + ServerApiProvider::as_ref(ctx).get_integrations_client(); let tx_id = tx_id.into_inner(); let poll_future = @@ -350,13 +355,11 @@ impl IntegrationCommandRunner { ctx.spawn( poll_future, move |runner, poll_result, ctx| { - match poll_result { - Ok(OauthConnectTxStatus::Completed) => { - // Inner loop done; try create or update again (outer loop). - // This may happen multiple times if the user needs to authorize multiple services. + match retry_state.continue_after_oauth(poll_result) { + Ok(next_retry_state) => { runner.start_create_or_update_flow( ctx, - retry_state.next(), + next_retry_state, next_integration_type, next_environment_uid, next_base_prompt, @@ -368,30 +371,10 @@ impl IntegrationCommandRunner { next_is_update, ); } - Ok(OauthConnectTxStatus::Failed) => { - ctx.terminate_app( - TerminationMode::ForceTerminate, - Some(Err(anyhow::anyhow!("OAuth authorization failed."))), - ); - } - Ok(OauthConnectTxStatus::Expired) => { - ctx.terminate_app( - TerminationMode::ForceTerminate, - Some(Err(anyhow::anyhow!("OAuth authorization expired."))), - ); - } - Ok(OauthConnectTxStatus::Pending) - | Ok(OauthConnectTxStatus::InProgress) => { - // Should not be returned by poll_oauth_until_terminal. - ctx.terminate_app( - TerminationMode::ForceTerminate, - Some(Err(anyhow::anyhow!("Unexpected non-terminal OAuth status returned"))), - ); - } Err(err) => { ctx.terminate_app( TerminationMode::ForceTerminate, - Some(Err(anyhow::anyhow!("Error polling OAuth status: {err}"))), + Some(Err(err)), ); } } diff --git a/app/src/ai/agent_sdk/integration_tests.rs b/app/src/ai/agent_sdk/integration_tests.rs index d669fd7b8ca..d4d485c764e 100644 --- a/app/src/ai/agent_sdk/integration_tests.rs +++ b/app/src/ai/agent_sdk/integration_tests.rs @@ -1,105 +1,26 @@ -use std::sync::Arc; -use std::sync::atomic::{AtomicUsize, Ordering}; -use std::time::Duration; - -use warp_graphql::mutations::create_simple_integration::CreateSimpleIntegrationOutput; use warp_graphql::queries::get_oauth_connect_tx_status::OauthConnectTxStatus; -use warpui::App; -use warpui::r#async::Timer; -use super::{IntegrationCommandRunner, IntegrationRetryState}; +use super::IntegrationRetryState; use crate::server::ids::ServerId; -use crate::server::server_api::integrations::MockIntegrationsClient; use crate::server::team_scope::RequestTeamScope; use crate::workspaces::user_workspaces::{TeamContextForOperation, TeamlessScopeForTest}; #[test] -fn oauth_retry_uses_the_initiating_team_scope() { +fn completed_oauth_continuation_preserves_the_initiating_team_scope() { let scope = TeamContextForOperation::new_for_test(ServerId::from(7)); - run_oauth_retry(RequestTeamScope::from_scope(&scope), false); + assert_completed_oauth_continuation(RequestTeamScope::from_scope(&scope)); } #[test] -fn oauth_retry_uses_the_initiating_teamless_scope() { - run_oauth_retry(RequestTeamScope::from_scope(&TeamlessScopeForTest), true); +fn completed_oauth_continuation_preserves_the_initiating_teamless_scope() { + assert_completed_oauth_continuation(RequestTeamScope::from_scope(&TeamlessScopeForTest)); } -fn run_oauth_retry(expected_scope: RequestTeamScope, is_update: bool) { - App::test((), move |mut app| async move { - let request_count = Arc::new(AtomicUsize::new(0)); - let observed_request_count = request_count.clone(); - let mut integrations_client = MockIntegrationsClient::new(); - integrations_client - .expect_create_or_update_simple_integration() - .times(2) - .returning( - move |request_scope, - integration_type, - request_is_update, - _environment_uid, - _base_prompt, - _model_id, - _mcp_servers_json, - _remove_mcp_server_names, - _worker_host, - enabled| { - assert_eq!(request_scope, expected_scope); - assert_eq!(integration_type, "slack"); - assert_eq!(request_is_update, is_update); - assert!(enabled); - let request_index = observed_request_count.fetch_add(1, Ordering::SeqCst); - Ok(if request_index == 0 { - CreateSimpleIntegrationOutput { - auth_url: Some("https://example.com/oauth".to_string()), - success: false, - message: "Authorization required".to_string(), - tx_id: Some(cynic::Id::new("oauth-tx")), - } - } else { - CreateSimpleIntegrationOutput { - auth_url: None, - success: true, - message: "Integration saved".to_string(), - tx_id: None, - } - }) - }, - ); - integrations_client - .expect_poll_oauth_connect_status() - .times(1) - .returning(|tx_id| { - assert_eq!(tx_id, "oauth-tx"); - Ok(OauthConnectTxStatus::Completed) - }); - - let runner = - app.add_model(move |_| IntegrationCommandRunner::new(Arc::new(integrations_client))); - runner.update(&mut app, |runner, ctx| { - runner.start_create_or_update_flow( - ctx, - IntegrationRetryState::new(expected_scope), - "slack".to_string(), - None, - None, - None, - None, - None, - None, - true, - is_update, - ); - }); - - for _ in 0..120 { - if request_count.load(Ordering::SeqCst) == 2 { - break; - } - Timer::after(Duration::from_millis(50)).await; - } +fn assert_completed_oauth_continuation(expected_scope: RequestTeamScope) { + let retry_state = IntegrationRetryState::new(expected_scope) + .continue_after_oauth(Ok(OauthConnectTxStatus::Completed)) + .unwrap(); - assert_eq!(request_count.load(Ordering::SeqCst), 2); - futures_lite::future::yield_now().await; - assert!(app.termination_result().is_none()); - }); + assert_eq!(retry_state.attempt, 2); + assert_eq!(retry_state.request_team_scope, expected_scope); } diff --git a/app/src/ai/agent_sdk/provider.rs b/app/src/ai/agent_sdk/provider.rs index bbc1bbb9454..f1aabf31a99 100644 --- a/app/src/ai/agent_sdk/provider.rs +++ b/app/src/ai/agent_sdk/provider.rs @@ -1,5 +1,4 @@ //! Provider command for linking third-party services. -use std::future::Future; use comfy_table::Cell; use serde::Serialize; @@ -39,16 +38,6 @@ impl ProviderCommandRunner { // This shouldn't need to be done, it's usually done as part of create fn setup(&self, provider_type: ProviderType, scope: ObjectScope, ctx: &mut ModelContext) { let refresh_future = super::common::refresh_workspace_metadata(ctx); - self.setup_after_workspace_metadata_refresh(refresh_future, provider_type, scope, ctx); - } - - fn setup_after_workspace_metadata_refresh( - &self, - refresh_future: impl Future> + Send + 'static, - provider_type: ProviderType, - scope: ObjectScope, - ctx: &mut ModelContext, - ) { ctx.spawn(refresh_future, move |runner, result, ctx| { if let Err(err) = result { super::report_fatal_error(err, ctx); diff --git a/app/src/ai/agent_sdk/provider_tests.rs b/app/src/ai/agent_sdk/provider_tests.rs index 3af1648416b..d96bb0a785d 100644 --- a/app/src/ai/agent_sdk/provider_tests.rs +++ b/app/src/ai/agent_sdk/provider_tests.rs @@ -1,37 +1,76 @@ -use std::sync::{Arc, Mutex}; +use std::sync::Arc; use warp_cli::provider::ProviderType; use warp_cli::scope::{ObjectScope, TeamSelection}; use warpui::App; -use warpui::r#async::Timer; use super::ProviderCommandRunner; +use crate::auth::AuthStateProvider; +use crate::network::NetworkStatus; use crate::server::ids::ServerId; +use crate::server::server_api::team::MockTeamClient; use crate::settings::PrivacySettings; -use crate::workspaces::user_workspaces::UserWorkspaces; +use crate::workspaces::team::Team; +use crate::workspaces::team_tester::TeamTesterStatus; +use crate::workspaces::update_manager::TeamUpdateManager; +use crate::workspaces::user_workspaces::{ + UserWorkspaces, WorkspacesMetadataResponse, WorkspacesMetadataWithPricing, +}; +use crate::workspaces::workspace::{Workspace, WorkspaceUid}; #[test] -fn setup_resolves_team_scope_after_workspace_metadata_refresh() { +fn setup_resolves_team_scope_from_refreshed_workspace_metadata() { App::test((), |mut app| async move { + let workspace_uid = WorkspaceUid::from(ServerId::from(1)); + let team_uid = ServerId::from(2); + let workspace = Workspace::from_local_cache( + workspace_uid, + "Test Workspace".to_string(), + Some(vec![Team::from_local_cache( + team_uid, + "Test Team".to_string(), + None, + None, + None, + None, + )]), + None, + ); + let mut team_client = MockTeamClient::new(); + team_client + .expect_workspaces_metadata() + .times(1) + .return_once(|| { + Ok(WorkspacesMetadataWithPricing { + metadata: WorkspacesMetadataResponse { + workspaces: vec![workspace], + joinable_teams: vec![], + experiments: None, + ai_credit_availability: None, + user_purchase_policy: None, + }, + pricing_info: None, + }) + }); + + app.add_singleton_model(|_| NetworkStatus::new()); + app.add_singleton_model(TeamTesterStatus::new); + app.add_singleton_model(|_| AuthStateProvider::new_for_test()); app.add_singleton_model(PrivacySettings::mock); - let user_workspaces = app.add_singleton_model(UserWorkspaces::default_mock); - let opened_urls = Arc::new(Mutex::new(Vec::new())); - let captured_urls = opened_urls.clone(); + app.add_singleton_model(UserWorkspaces::default_mock); + app.add_singleton_model(|ctx| TeamUpdateManager::new(Arc::new(team_client), None, ctx)); + + let (opened_url_sender, opened_url_receiver) = async_channel::bounded(1); app.update(|ctx| { ctx.set_before_open_url(move |url, _| { - captured_urls.lock().unwrap().push(url.to_string()); + opened_url_sender.try_send(url.to_string()).unwrap(); url.to_string() }); }); - let (refresh_sender, refresh_receiver) = async_channel::bounded(1); let runner = app.add_model(|_| ProviderCommandRunner); runner.update(&mut app, |runner, ctx| { - runner.setup_after_workspace_metadata_refresh( - async move { - refresh_receiver.recv().await.unwrap(); - Ok(()) - }, + runner.setup( ProviderType::Slack, ObjectScope { team_selection: TeamSelection { team: Some(None) }, @@ -41,35 +80,14 @@ fn setup_resolves_team_scope_after_workspace_metadata_refresh() { ); }); - for _ in 0..3 { - futures_lite::future::yield_now().await; - } - assert!(opened_urls.lock().unwrap().is_empty()); - - user_workspaces.update(&mut app, |workspaces, ctx| { - workspaces.setup_test_workspace(ctx); - }); - refresh_sender.send(()).await.unwrap(); - - for _ in 0..20 { - if !opened_urls.lock().unwrap().is_empty() { - break; - } - Timer::after(std::time::Duration::from_millis(10)).await; - } - - let opened_urls = opened_urls.lock().unwrap(); - assert_eq!(opened_urls.len(), 1); - let expected_url_suffix = format!( - "/oauth/connect/slack?principalType=team&principalId={}", - ServerId::from(2) - ); - assert!( - opened_urls[0].ends_with(&expected_url_suffix), - "opened URL: {}", - opened_urls[0] - ); - drop(opened_urls); + assert!(matches!( + opened_url_receiver.try_recv(), + Err(async_channel::TryRecvError::Empty) + )); + let opened_url = opened_url_receiver.recv().await.unwrap(); + assert!(opened_url.ends_with(&format!( + "/oauth/connect/slack?principalType=team&principalId={team_uid}" + ))); assert!(app.termination_result().is_none()); }); }