mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-10 08:07:03 +00:00
* feat: let a native trigger be disabled without deleting it Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: show and control the native trigger pause outside the flow editor Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: create a native trigger already paused instead of pausing it after Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * docs: correct the native trigger enabled comments for create-time init Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
975 lines
30 KiB
Rust
975 lines
30 KiB
Rust
/*!
|
|
* Integration tests for the native trigger system (Google).
|
|
*
|
|
* Tests cover 4 business-logic areas:
|
|
* 1. Resource path change — cleanup old path, recreate at new path
|
|
* 2. Config loading — workspace-level, instance-level, token update
|
|
* 3. Channel expiration renewal — should_renew_channel pure logic
|
|
* 4. Delete workspace integration — full cascade, cleanup preserves triggers, parse_stop_channel_params
|
|
*/
|
|
|
|
use serde_json::json;
|
|
use sqlx::{Pool, Postgres};
|
|
|
|
use axum::http::StatusCode;
|
|
use windmill_api_auth::ApiAuthed;
|
|
use windmill_common::{
|
|
error::Error,
|
|
variables::{build_crypt, encrypt},
|
|
};
|
|
use windmill_native_triggers::{
|
|
classify_read_failure, decrypt_oauth_data, delete_native_trigger, delete_workspace_integration,
|
|
get_workspace_integration,
|
|
github::GitHub,
|
|
google::{parse_stop_channel_params, should_renew_channel},
|
|
grant_refused, http_error_status, list_native_triggers, map_external_error,
|
|
native_trigger_is_enabled,
|
|
nextcloud::NextCloud,
|
|
require_native_integration_use, set_native_trigger_enabled, store_native_trigger,
|
|
store_workspace_integration, update_native_trigger, External, ExternalReadFailure,
|
|
HttpRequestError, NativeTriggerConfig, OAuthConfig, ServiceName,
|
|
};
|
|
|
|
// ============================================================================
|
|
// Helpers
|
|
// ============================================================================
|
|
|
|
async fn insert_test_script(db: &Pool<Postgres>, path: &str) -> anyhow::Result<i64> {
|
|
let hash: i64 = rand::random::<i64>().unsigned_abs() as i64;
|
|
sqlx::query(
|
|
"INSERT INTO script (workspace_id, hash, path, summary, description, content,
|
|
created_by, language, kind, lock)
|
|
VALUES ('test-workspace', $1, $2, '', '', 'def main(): pass',
|
|
'test-user', 'python3', 'script', '')",
|
|
)
|
|
.bind(hash)
|
|
.bind(path)
|
|
.execute(db)
|
|
.await?;
|
|
Ok(hash)
|
|
}
|
|
|
|
fn test_authed() -> ApiAuthed {
|
|
ApiAuthed {
|
|
email: "test@windmill.dev".to_string(),
|
|
username: "test-user".to_string(),
|
|
is_admin: true,
|
|
is_operator: false,
|
|
groups: vec!["all".to_string()],
|
|
folders: vec![],
|
|
scopes: None,
|
|
username_override: None,
|
|
username_override_is_token_label: false,
|
|
is_session_token: false,
|
|
token_prefix: None,
|
|
read_only: false,
|
|
job_id: None,
|
|
credential_expiry: None,
|
|
}
|
|
}
|
|
|
|
/// Set up a complete workspace integration with account+variable+resource.
|
|
/// Returns (resource_path, account_id).
|
|
async fn setup_oauth_integration(
|
|
db: &Pool<Postgres>,
|
|
service_name: ServiceName,
|
|
resource_path: &str,
|
|
access_token: &str,
|
|
refresh_token: &str,
|
|
oauth_data_override: Option<serde_json::Value>,
|
|
) -> anyhow::Result<i32> {
|
|
// 1. Create account with is_workspace_integration=true
|
|
let account_id: i32 = sqlx::query_scalar!(
|
|
"INSERT INTO account (workspace_id, client, expires_at, refresh_token, is_workspace_integration)
|
|
VALUES ('test-workspace', $1, now() + interval '1 hour', $2, true)
|
|
RETURNING id",
|
|
service_name.as_str(),
|
|
refresh_token,
|
|
)
|
|
.fetch_one(db)
|
|
.await?;
|
|
|
|
// 2. Encrypt and create variable
|
|
let mc = build_crypt(db, "test-workspace").await?;
|
|
let encrypted = encrypt(&mc, access_token);
|
|
|
|
sqlx::query!(
|
|
"INSERT INTO variable (workspace_id, path, value, is_secret, description, account, is_oauth)
|
|
VALUES ('test-workspace', $1, $2, true, 'test oauth token', $3, true)",
|
|
resource_path,
|
|
encrypted,
|
|
account_id,
|
|
)
|
|
.execute(db)
|
|
.await?;
|
|
|
|
// 3. Create resource
|
|
let resource_value = json!({ "token": format!("$var:{}", resource_path) });
|
|
sqlx::query!(
|
|
"INSERT INTO resource (workspace_id, path, value, resource_type, extra_perms, created_by)
|
|
VALUES ('test-workspace', $1, $2, $3, '{}'::jsonb, 'test-user')",
|
|
resource_path,
|
|
resource_value,
|
|
service_name.resource_type(),
|
|
)
|
|
.execute(db)
|
|
.await?;
|
|
|
|
// 4. Store workspace integration with resource_path
|
|
let oauth_data = oauth_data_override.unwrap_or_else(|| {
|
|
json!({
|
|
"client_id": "test-client-id",
|
|
"client_secret": "test-client-secret",
|
|
"base_url": "https://example.com",
|
|
"resource_path": resource_path,
|
|
})
|
|
});
|
|
|
|
let authed = test_authed();
|
|
let mut tx = db.begin().await?;
|
|
store_workspace_integration(
|
|
&mut *tx,
|
|
&authed,
|
|
"test-workspace",
|
|
service_name,
|
|
oauth_data,
|
|
Some(resource_path),
|
|
)
|
|
.await?;
|
|
tx.commit().await?;
|
|
|
|
Ok(account_id)
|
|
}
|
|
|
|
fn now_ms() -> i64 {
|
|
std::time::SystemTime::now()
|
|
.duration_since(std::time::UNIX_EPOCH)
|
|
.unwrap()
|
|
.as_millis() as i64
|
|
}
|
|
|
|
// ============================================================================
|
|
// 1. Resource Path Change
|
|
// ============================================================================
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_resource_path_change(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
let path_a = "u/test-user/native_gworkspace";
|
|
setup_oauth_integration(
|
|
&db,
|
|
ServiceName::Google,
|
|
path_a,
|
|
"token-a",
|
|
"refresh-a",
|
|
None,
|
|
)
|
|
.await?;
|
|
|
|
// Verify decrypt works at path A
|
|
let config: OAuthConfig =
|
|
decrypt_oauth_data(&db, "test-workspace", ServiceName::Google).await?;
|
|
assert_eq!(config.access_token, "token-a");
|
|
|
|
// Cleanup old path
|
|
let mut tx = db.begin().await?;
|
|
windmill_native_triggers::workspace_integrations::cleanup_oauth_resource(
|
|
&mut *tx,
|
|
"test-workspace",
|
|
ServiceName::Google,
|
|
)
|
|
.await;
|
|
tx.commit().await?;
|
|
|
|
// Recreate at path B
|
|
let path_b = "u/test-user/native_gworkspace_v2";
|
|
setup_oauth_integration(
|
|
&db,
|
|
ServiceName::Google,
|
|
path_b,
|
|
"token-b",
|
|
"refresh-b",
|
|
None,
|
|
)
|
|
.await?;
|
|
|
|
// Path A resources should be gone
|
|
let var_count: i64 = sqlx::query_scalar!(
|
|
"SELECT count(*) FROM variable WHERE workspace_id = 'test-workspace' AND path = $1",
|
|
path_a,
|
|
)
|
|
.fetch_one(&db)
|
|
.await?
|
|
.unwrap_or(0);
|
|
assert_eq!(var_count, 0, "variable at old path should be deleted");
|
|
|
|
let res_count: i64 = sqlx::query_scalar!(
|
|
"SELECT count(*) FROM resource WHERE workspace_id = 'test-workspace' AND path = $1",
|
|
path_a,
|
|
)
|
|
.fetch_one(&db)
|
|
.await?
|
|
.unwrap_or(0);
|
|
assert_eq!(res_count, 0, "resource at old path should be deleted");
|
|
|
|
// Path B should work
|
|
let config: OAuthConfig =
|
|
decrypt_oauth_data(&db, "test-workspace", ServiceName::Google).await?;
|
|
assert_eq!(config.access_token, "token-b");
|
|
assert_eq!(config.refresh_token.as_deref(), Some("refresh-b"));
|
|
|
|
Ok(())
|
|
}
|
|
|
|
// ============================================================================
|
|
// 2. Config Loading — workspace vs instance + token update
|
|
// ============================================================================
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_decrypt_workspace_level(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
let resource_path = "u/test-user/native_gworkspace";
|
|
setup_oauth_integration(
|
|
&db,
|
|
ServiceName::Google,
|
|
resource_path,
|
|
"ws-access-token",
|
|
"ws-refresh-token",
|
|
None,
|
|
)
|
|
.await?;
|
|
|
|
let config: OAuthConfig =
|
|
decrypt_oauth_data(&db, "test-workspace", ServiceName::Google).await?;
|
|
|
|
assert_eq!(config.access_token, "ws-access-token");
|
|
assert_eq!(config.refresh_token.as_deref(), Some("ws-refresh-token"));
|
|
assert_eq!(config.client_id, "test-client-id");
|
|
assert_eq!(config.client_secret, "test-client-secret");
|
|
assert_eq!(config.base_url, "https://example.com");
|
|
|
|
Ok(())
|
|
}
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_decrypt_instance_level(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
// Insert instance-level credentials into global_settings
|
|
sqlx::query!(
|
|
"INSERT INTO global_settings (name, value) VALUES ('oauths', $1)
|
|
ON CONFLICT (name) DO UPDATE SET value = $1",
|
|
json!({
|
|
"gworkspace": {
|
|
"id": "instance-client-id",
|
|
"secret": "instance-client-secret"
|
|
}
|
|
}),
|
|
)
|
|
.execute(&db)
|
|
.await?;
|
|
|
|
let resource_path = "u/test-user/native_gworkspace";
|
|
let oauth_data = json!({
|
|
"instance_shared": true,
|
|
"base_url": "https://accounts.google.com",
|
|
"resource_path": resource_path,
|
|
});
|
|
|
|
setup_oauth_integration(
|
|
&db,
|
|
ServiceName::Google,
|
|
resource_path,
|
|
"inst-access-token",
|
|
"inst-refresh-token",
|
|
Some(oauth_data),
|
|
)
|
|
.await?;
|
|
|
|
let config: OAuthConfig =
|
|
decrypt_oauth_data(&db, "test-workspace", ServiceName::Google).await?;
|
|
|
|
assert_eq!(config.client_id, "instance-client-id");
|
|
assert_eq!(config.client_secret, "instance-client-secret");
|
|
assert_eq!(config.access_token, "inst-access-token");
|
|
|
|
Ok(())
|
|
}
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_token_update_persists(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
let resource_path = "u/test-user/native_gworkspace";
|
|
let account_id = setup_oauth_integration(
|
|
&db,
|
|
ServiceName::Google,
|
|
resource_path,
|
|
"old-access-token",
|
|
"old-refresh-token",
|
|
None,
|
|
)
|
|
.await?;
|
|
|
|
// Verify old tokens
|
|
let config: OAuthConfig =
|
|
decrypt_oauth_data(&db, "test-workspace", ServiceName::Google).await?;
|
|
assert_eq!(config.access_token, "old-access-token");
|
|
|
|
// Simulate token refresh: update variable + account
|
|
let mc = build_crypt(&db, "test-workspace").await?;
|
|
let new_encrypted = encrypt(&mc, "new-access-token");
|
|
sqlx::query!(
|
|
"UPDATE variable SET value = $1 WHERE workspace_id = 'test-workspace' AND path = $2",
|
|
new_encrypted,
|
|
resource_path,
|
|
)
|
|
.execute(&db)
|
|
.await?;
|
|
|
|
sqlx::query!(
|
|
"UPDATE account SET refresh_token = $1 WHERE workspace_id = 'test-workspace' AND id = $2",
|
|
"new-refresh-token",
|
|
account_id,
|
|
)
|
|
.execute(&db)
|
|
.await?;
|
|
|
|
// Verify new tokens
|
|
let config: OAuthConfig =
|
|
decrypt_oauth_data(&db, "test-workspace", ServiceName::Google).await?;
|
|
assert_eq!(config.access_token, "new-access-token");
|
|
assert_eq!(config.refresh_token.as_deref(), Some("new-refresh-token"));
|
|
|
|
Ok(())
|
|
}
|
|
|
|
// ============================================================================
|
|
// 3. Channel Expiration Renewal — should_renew_channel
|
|
// ============================================================================
|
|
|
|
#[test]
|
|
fn test_require_native_integration_use_blocks_operators() {
|
|
// Regression: the integration *use* routes (calendar/drive/repo/event pickers)
|
|
// must reject read-only operators, who cannot create native triggers and so
|
|
// must not be able to drive the admin-configured integration's upstream API.
|
|
let mut operator = test_authed();
|
|
operator.is_admin = false;
|
|
operator.is_operator = true;
|
|
assert!(require_native_integration_use(&operator).is_err());
|
|
|
|
// A regular non-admin author (the population that configures triggers) is allowed.
|
|
let mut author = test_authed();
|
|
author.is_admin = false;
|
|
author.is_operator = false;
|
|
assert!(require_native_integration_use(&author).is_ok());
|
|
|
|
// Admins are allowed.
|
|
assert!(require_native_integration_use(&test_authed()).is_ok());
|
|
}
|
|
|
|
#[test]
|
|
fn test_should_renew_drive_channel_expired() {
|
|
let config = json!({
|
|
"triggerType": "drive",
|
|
"expiration": (now_ms() - 1000).to_string(),
|
|
});
|
|
assert!(should_renew_channel(&config));
|
|
}
|
|
|
|
#[test]
|
|
fn test_should_renew_drive_channel_within_window() {
|
|
// 30 minutes remaining — within the 1-hour Drive renewal window
|
|
let config = json!({
|
|
"triggerType": "drive",
|
|
"expiration": (now_ms() + 30 * 60 * 1000).to_string(),
|
|
});
|
|
assert!(should_renew_channel(&config));
|
|
}
|
|
|
|
#[test]
|
|
fn test_should_renew_drive_channel_not_yet() {
|
|
// 2 hours remaining — outside the 1-hour Drive renewal window
|
|
let config = json!({
|
|
"triggerType": "drive",
|
|
"expiration": (now_ms() + 2 * 60 * 60 * 1000).to_string(),
|
|
});
|
|
assert!(!should_renew_channel(&config));
|
|
}
|
|
|
|
#[test]
|
|
fn test_should_renew_calendar_channel_within_window() {
|
|
// 12 hours remaining — within the 1-day Calendar renewal window
|
|
let config = json!({
|
|
"triggerType": "calendar",
|
|
"expiration": (now_ms() + 12 * 60 * 60 * 1000).to_string(),
|
|
});
|
|
assert!(should_renew_channel(&config));
|
|
}
|
|
|
|
#[test]
|
|
fn test_should_renew_calendar_channel_not_yet() {
|
|
// 2 days remaining — outside the 1-day Calendar renewal window
|
|
let config = json!({
|
|
"triggerType": "calendar",
|
|
"expiration": (now_ms() + 2 * 24 * 60 * 60 * 1000).to_string(),
|
|
});
|
|
assert!(!should_renew_channel(&config));
|
|
}
|
|
|
|
#[test]
|
|
fn test_should_renew_channel_zero_expiration() {
|
|
let config = json!({
|
|
"triggerType": "drive",
|
|
"expiration": "0",
|
|
});
|
|
assert!(!should_renew_channel(&config));
|
|
}
|
|
|
|
#[test]
|
|
fn test_should_renew_channel_missing_fields() {
|
|
assert!(!should_renew_channel(&json!({})));
|
|
}
|
|
|
|
// ============================================================================
|
|
// 4. Delete Workspace Integration
|
|
// ============================================================================
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_delete_integration_full_cascade(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
let resource_path = "u/test-user/native_gworkspace";
|
|
let account_id = setup_oauth_integration(
|
|
&db,
|
|
ServiceName::Google,
|
|
resource_path,
|
|
"token",
|
|
"refresh",
|
|
None,
|
|
)
|
|
.await?;
|
|
|
|
// Add a native trigger linked to this integration
|
|
insert_test_script(&db, "f/test/handler").await?;
|
|
let trigger_config = NativeTriggerConfig {
|
|
script_path: "f/test/handler".to_string(),
|
|
is_flow: false,
|
|
webhook_token: "abcdefghij1234567890".to_string(),
|
|
};
|
|
store_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Google,
|
|
"ext-1",
|
|
&trigger_config,
|
|
json!({"triggerType": "drive"}),
|
|
None,
|
|
true,
|
|
)
|
|
.await?;
|
|
|
|
// Step 1: Delete triggers
|
|
let deleted =
|
|
delete_native_trigger(&db, "test-workspace", ServiceName::Google, "ext-1").await?;
|
|
assert!(deleted);
|
|
|
|
// Step 2: Cleanup OAuth resources
|
|
let mut tx = db.begin().await?;
|
|
windmill_native_triggers::workspace_integrations::cleanup_oauth_resource(
|
|
&mut *tx,
|
|
"test-workspace",
|
|
ServiceName::Google,
|
|
)
|
|
.await;
|
|
tx.commit().await?;
|
|
|
|
// Step 3: Delete workspace integration
|
|
let mut tx = db.begin().await?;
|
|
let deleted =
|
|
delete_workspace_integration(&mut *tx, "test-workspace", ServiceName::Google).await?;
|
|
tx.commit().await?;
|
|
assert!(deleted);
|
|
|
|
// Verify everything is gone
|
|
let var_count: i64 = sqlx::query_scalar!(
|
|
"SELECT count(*) FROM variable WHERE workspace_id = 'test-workspace' AND path = $1",
|
|
resource_path,
|
|
)
|
|
.fetch_one(&db)
|
|
.await?
|
|
.unwrap_or(0);
|
|
assert_eq!(var_count, 0);
|
|
|
|
let acc_count: i64 = sqlx::query_scalar!(
|
|
"SELECT count(*) FROM account WHERE workspace_id = 'test-workspace' AND id = $1",
|
|
account_id,
|
|
)
|
|
.fetch_one(&db)
|
|
.await?
|
|
.unwrap_or(0);
|
|
assert_eq!(acc_count, 0);
|
|
|
|
let res_count: i64 = sqlx::query_scalar!(
|
|
"SELECT count(*) FROM resource WHERE workspace_id = 'test-workspace' AND path = $1",
|
|
resource_path,
|
|
)
|
|
.fetch_one(&db)
|
|
.await?
|
|
.unwrap_or(0);
|
|
assert_eq!(res_count, 0);
|
|
|
|
assert!(
|
|
get_workspace_integration(&db, "test-workspace", ServiceName::Google)
|
|
.await
|
|
.is_err()
|
|
);
|
|
|
|
Ok(())
|
|
}
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_cleanup_preserves_triggers(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
let resource_path = "u/test-user/native_gworkspace";
|
|
setup_oauth_integration(
|
|
&db,
|
|
ServiceName::Google,
|
|
resource_path,
|
|
"token",
|
|
"refresh",
|
|
None,
|
|
)
|
|
.await?;
|
|
|
|
// Create a trigger
|
|
insert_test_script(&db, "f/test/handler").await?;
|
|
let trigger_config = NativeTriggerConfig {
|
|
script_path: "f/test/handler".to_string(),
|
|
is_flow: false,
|
|
webhook_token: "abcdefghij1234567890".to_string(),
|
|
};
|
|
store_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Google,
|
|
"ext-1",
|
|
&trigger_config,
|
|
json!({"triggerType": "drive"}),
|
|
None,
|
|
true,
|
|
)
|
|
.await?;
|
|
|
|
// Cleanup OAuth only — should NOT remove the trigger
|
|
let mut tx = db.begin().await?;
|
|
windmill_native_triggers::workspace_integrations::cleanup_oauth_resource(
|
|
&mut *tx,
|
|
"test-workspace",
|
|
ServiceName::Google,
|
|
)
|
|
.await;
|
|
tx.commit().await?;
|
|
|
|
// OAuth resources gone
|
|
let var_count: i64 = sqlx::query_scalar!(
|
|
"SELECT count(*) FROM variable WHERE workspace_id = 'test-workspace' AND path = $1",
|
|
resource_path,
|
|
)
|
|
.fetch_one(&db)
|
|
.await?
|
|
.unwrap_or(0);
|
|
assert_eq!(var_count, 0);
|
|
|
|
// Trigger still exists
|
|
let trigger_count: i64 = sqlx::query_scalar!(
|
|
"SELECT count(*) FROM native_trigger WHERE workspace_id = 'test-workspace' AND service_name = 'google'"
|
|
)
|
|
.fetch_one(&db)
|
|
.await?
|
|
.unwrap_or(0);
|
|
assert_eq!(trigger_count, 1, "trigger should survive OAuth cleanup");
|
|
|
|
Ok(())
|
|
}
|
|
|
|
// ============================================================================
|
|
// 5. Runnable rename
|
|
// ============================================================================
|
|
|
|
/// A rename has to carry the trigger row onto the new path and report it as moved: listings only
|
|
/// return rows whose runnable still exists, so one left behind on the old path disappears from the
|
|
/// UI for good, and one not reported keeps a webhook aimed at the old path.
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_rename_moves_native_trigger(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
insert_test_script(&db, "f/test/before").await?;
|
|
store_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"ext-1",
|
|
&NativeTriggerConfig {
|
|
script_path: "f/test/before".to_string(),
|
|
is_flow: false,
|
|
webhook_token: "abcdefghij1234567890".to_string(),
|
|
},
|
|
json!({"event": "OCP\\Files\\Events\\Node\\NodeCreatedEvent"}),
|
|
None,
|
|
true,
|
|
)
|
|
.await?;
|
|
// An unrelated trigger already sitting on the target path must not be reported as moved.
|
|
insert_test_script(&db, "f/test/after").await?;
|
|
store_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"ext-2",
|
|
&NativeTriggerConfig {
|
|
script_path: "f/test/after".to_string(),
|
|
is_flow: false,
|
|
webhook_token: "0987654321jihgfedcba".to_string(),
|
|
},
|
|
json!({"event": "OCP\\Files\\Events\\Node\\NodeCreatedEvent"}),
|
|
None,
|
|
true,
|
|
)
|
|
.await?;
|
|
|
|
let mut tx = db.begin().await?;
|
|
sqlx::query!(
|
|
"UPDATE script SET path = $1 WHERE workspace_id = 'test-workspace' AND path = $2",
|
|
"f/test/after",
|
|
"f/test/before",
|
|
)
|
|
.execute(&mut *tx)
|
|
.await?;
|
|
let moved = windmill_common::triggers::update_triggers_script_path(
|
|
&mut tx,
|
|
"f/test/after",
|
|
"f/test/before",
|
|
"test-workspace",
|
|
false,
|
|
)
|
|
.await?;
|
|
tx.commit().await?;
|
|
|
|
assert_eq!(
|
|
moved
|
|
.iter()
|
|
.map(|t| (t.service_name.as_str(), t.external_id.as_str()))
|
|
.collect::<Vec<_>>(),
|
|
vec![("nextcloud", "ext-1")]
|
|
);
|
|
|
|
let triggers = list_native_triggers(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
None,
|
|
None,
|
|
Some("f/test/after"),
|
|
Some(false),
|
|
)
|
|
.await?;
|
|
assert_eq!(
|
|
triggers.len(),
|
|
2,
|
|
"the moved trigger should be listed under the new path"
|
|
);
|
|
|
|
Ok(())
|
|
}
|
|
|
|
// --- parse_stop_channel_params ---
|
|
|
|
#[test]
|
|
fn test_parse_stop_channel_params_drive() {
|
|
let config = json!({
|
|
"triggerType": "drive",
|
|
"googleChannelId": "chan-abc",
|
|
"googleResourceId": "res-123",
|
|
});
|
|
let (channel_id, resource_id, url) = parse_stop_channel_params(&config);
|
|
assert_eq!(channel_id.as_deref(), Some("chan-abc"));
|
|
assert_eq!(resource_id, "res-123");
|
|
assert!(
|
|
url.contains("googleapis.com/drive/v3/channels/stop"),
|
|
"url={}",
|
|
url
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn test_parse_stop_channel_params_calendar() {
|
|
let config = json!({
|
|
"triggerType": "calendar",
|
|
"googleChannelId": "chan-xyz",
|
|
"googleResourceId": "res-456",
|
|
});
|
|
let (channel_id, resource_id, url) = parse_stop_channel_params(&config);
|
|
assert_eq!(channel_id.as_deref(), Some("chan-xyz"));
|
|
assert_eq!(resource_id, "res-456");
|
|
assert!(
|
|
url.contains("googleapis.com/calendar/v3/channels/stop"),
|
|
"url={}",
|
|
url
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn test_parse_stop_channel_params_default() {
|
|
// Missing triggerType defaults to Drive; missing googleChannelId yields None.
|
|
let config = json!({ "googleResourceId": "res-789" });
|
|
let (channel_id, resource_id, url) = parse_stop_channel_params(&config);
|
|
assert!(channel_id.is_none());
|
|
assert_eq!(resource_id, "res-789");
|
|
assert!(url.contains("drive/v3/channels/stop"), "url={}", url);
|
|
}
|
|
|
|
#[test]
|
|
fn test_parse_stop_channel_params_missing_resource_id() {
|
|
let config = json!({ "triggerType": "drive" });
|
|
let (channel_id, resource_id, _url) = parse_stop_channel_params(&config);
|
|
assert!(channel_id.is_none());
|
|
assert_eq!(resource_id, "");
|
|
}
|
|
|
|
// --- provider error reporting ---
|
|
|
|
fn nextcloud_error(status: StatusCode, body: &str) -> Error {
|
|
NextCloud.external_api_error(HttpRequestError::ApiError { status, body: body.to_string() })
|
|
}
|
|
|
|
/// A rejection has to reach the user as the service's own sentence plus what to do about it,
|
|
/// never as an internal error carrying the raw envelope.
|
|
#[test]
|
|
fn test_provider_rejection_is_readable_and_not_internal() {
|
|
let err = nextcloud_error(
|
|
StatusCode::FORBIDDEN,
|
|
r#"{"ocs":{"meta":{"status":"failure","statuscode":403,"message":"Logged in account must be an admin, a sub admin or gotten special right to access this setting"},"data":[]}}"#,
|
|
);
|
|
|
|
let message = map_external_error(err).to_string();
|
|
assert!(
|
|
message.contains("Logged in account must be an admin"),
|
|
"message={message}"
|
|
);
|
|
assert!(
|
|
!message.contains("\"ocs\""),
|
|
"the envelope should not reach the user: {message}"
|
|
);
|
|
assert!(
|
|
message.contains("Workspace settings > Integrations"),
|
|
"the hint should say what to do: {message}"
|
|
);
|
|
}
|
|
|
|
/// Both the "trigger is gone on the service" path and the delete that tolerates an
|
|
/// already-removed webhook branch on this status.
|
|
#[test]
|
|
fn test_provider_404_is_recognized() {
|
|
let err = nextcloud_error(StatusCode::NOT_FOUND, "{}");
|
|
assert_eq!(http_error_status(&err), Some(StatusCode::NOT_FOUND));
|
|
assert!(
|
|
matches!(map_external_error(err), Error::NotFound(_)),
|
|
"a missing external trigger must map to NotFound"
|
|
);
|
|
assert!(matches!(
|
|
classify_read_failure(nextcloud_error(StatusCode::NOT_FOUND, "{}")),
|
|
ExternalReadFailure::Missing
|
|
));
|
|
}
|
|
|
|
/// Sending a user to reconnect their integration is only right when the token endpoint refused
|
|
/// the grant; a busy or broken endpoint has them fix credentials that are fine.
|
|
#[test]
|
|
fn test_refresh_failures_blame_only_the_grant_they_refuse() {
|
|
let refused = |code: u16| grant_refused(Some(StatusCode::from_u16(code).unwrap()), "");
|
|
for code in [400, 401, 403] {
|
|
assert!(refused(code), "{code} refuses the grant");
|
|
}
|
|
for code in [404, 408, 429, 500, 503] {
|
|
assert!(!refused(code), "{code} says nothing about the grant");
|
|
}
|
|
assert!(!grant_refused(None, ""));
|
|
|
|
// GitHub answers `bad_refresh_token` with HTTP 200, so the body is the only tell.
|
|
let ok = Some(StatusCode::OK);
|
|
assert!(grant_refused(ok, r#"{"error":"bad_refresh_token"}"#));
|
|
assert!(grant_refused(ok, r#"{"error":"invalid_grant"}"#));
|
|
assert!(!grant_refused(
|
|
ok,
|
|
r#"{"access_token":"t","token_type":"bearer"}"#
|
|
));
|
|
}
|
|
|
|
/// A service that is busy or broken has not refused anything, and callers react differently to
|
|
/// the two. GitHub and Google spend a 403 on throttling, where advice about permissions sends
|
|
/// the reader after a problem they do not have.
|
|
#[test]
|
|
fn test_transient_service_failures_are_not_refusals() {
|
|
for transient in [408, 429, 503] {
|
|
let err = nextcloud_error(StatusCode::from_u16(transient).unwrap(), "{}");
|
|
assert!(
|
|
matches!(map_external_error(err), Error::BadGateway(_)),
|
|
"{transient} should read as the service failing to serve, not refusing"
|
|
);
|
|
}
|
|
|
|
// GitHub words its throttle two ways, and neither is a permission problem.
|
|
for wording in [
|
|
"API rate limit exceeded for user ID 1.",
|
|
"You have exceeded a secondary rate limit.",
|
|
"You have triggered an abuse detection mechanism.",
|
|
] {
|
|
let throttled = GitHub.external_api_error(HttpRequestError::ApiError {
|
|
status: StatusCode::FORBIDDEN,
|
|
body: format!(r#"{{"message":"{wording}"}}"#),
|
|
});
|
|
let throttled = map_external_error(throttled);
|
|
assert!(
|
|
matches!(throttled, Error::BadGateway(_)),
|
|
"a throttled 403 is the service failing to serve: {throttled:?}"
|
|
);
|
|
assert!(
|
|
!throttled.to_string().contains("admin rights"),
|
|
"a throttled 403 must not advise about permissions: {throttled}"
|
|
);
|
|
}
|
|
|
|
let refused = GitHub.external_api_error(HttpRequestError::ApiError {
|
|
status: StatusCode::FORBIDDEN,
|
|
body: r#"{"message":"Must have admin rights to Repository."}"#.to_string(),
|
|
});
|
|
assert!(
|
|
map_external_error(refused)
|
|
.to_string()
|
|
.contains("admin rights"),
|
|
"a real 403 keeps its guidance"
|
|
);
|
|
}
|
|
|
|
/// A service read degrades to the stored configuration, but only for the service's own
|
|
/// failures: `External::get` also runs queries, and reporting one of those as the service's
|
|
/// word would hide a Windmill outage behind a 200.
|
|
#[test]
|
|
fn test_only_service_failures_degrade_the_read() {
|
|
assert!(matches!(
|
|
classify_read_failure(nextcloud_error(StatusCode::FORBIDDEN, "{}")),
|
|
ExternalReadFailure::Unreadable(_)
|
|
));
|
|
assert!(matches!(
|
|
classify_read_failure(Error::internal_err("connection pool timed out")),
|
|
ExternalReadFailure::Internal(_)
|
|
));
|
|
let internal = Error::internal_err("connection pool timed out");
|
|
assert!(
|
|
matches!(map_external_error(internal), Error::InternalErrLoc { .. }),
|
|
"a non-provider error must pass through unmapped"
|
|
);
|
|
}
|
|
|
|
/// The pause switch a webhook delivery is gated on. A trigger arrives enabled, survives an
|
|
/// unrelated edit, and an unknown one reads as enabled so a delivery Windmill cannot place is
|
|
/// never silently dropped.
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base"))]
|
|
async fn test_native_trigger_enabled_toggle(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
insert_test_script(&db, "f/test/handler").await?;
|
|
let config = NativeTriggerConfig {
|
|
script_path: "f/test/handler".to_string(),
|
|
is_flow: false,
|
|
webhook_token: "abcdefghij1234567890".to_string(),
|
|
};
|
|
store_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"ext-1",
|
|
&config,
|
|
json!({"event": "OCA\\Files\\Event\\LoadAdditionalScriptsEvent"}),
|
|
None,
|
|
true,
|
|
)
|
|
.await?;
|
|
|
|
assert!(
|
|
native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-1").await?,
|
|
"a new trigger fires"
|
|
);
|
|
|
|
assert!(
|
|
set_native_trigger_enabled(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"ext-1",
|
|
false
|
|
)
|
|
.await?
|
|
);
|
|
assert!(
|
|
!native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-1").await?
|
|
);
|
|
|
|
// Saving a configuration must not resume a trigger someone paused.
|
|
update_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"ext-1",
|
|
&config,
|
|
None,
|
|
Some("edited"),
|
|
)
|
|
.await?;
|
|
assert!(
|
|
!native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-1").await?,
|
|
"an edit leaves the pause in place"
|
|
);
|
|
|
|
// A recreate registers a fresh trigger and must be able to come up already paused, in one
|
|
// write, rather than being enabled for as long as it takes a second call to arrive.
|
|
store_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"ext-2",
|
|
&config,
|
|
json!({}),
|
|
None,
|
|
false,
|
|
)
|
|
.await?;
|
|
assert!(
|
|
!native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-2").await?
|
|
);
|
|
|
|
// The conflict branch is a re-registration of a trigger that already exists, so it carries no
|
|
// opinion about the pause.
|
|
store_native_trigger(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"ext-2",
|
|
&config,
|
|
json!({}),
|
|
None,
|
|
true,
|
|
)
|
|
.await?;
|
|
assert!(
|
|
!native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-2").await?,
|
|
"re-registering leaves the pause in place"
|
|
);
|
|
|
|
assert!(
|
|
!set_native_trigger_enabled(
|
|
&db,
|
|
"test-workspace",
|
|
ServiceName::Nextcloud,
|
|
"unknown",
|
|
true
|
|
)
|
|
.await?,
|
|
"nothing to toggle"
|
|
);
|
|
assert!(
|
|
native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "unknown").await?,
|
|
"a trigger Windmill has no row for is not treated as paused"
|
|
);
|
|
|
|
Ok(())
|
|
}
|