mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-21 08:02:38 +00:00
fix: mint fresh Google channel IDs on update/renew (#8870)
Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
f655d4f295
commit
ad2e855a83
@@ -554,9 +554,11 @@ async fn test_cleanup_preserves_triggers(db: Pool<Postgres>) -> anyhow::Result<(
|
||||
fn test_parse_stop_channel_params_drive() {
|
||||
let config = json!({
|
||||
"triggerType": "drive",
|
||||
"googleChannelId": "chan-abc",
|
||||
"googleResourceId": "res-123",
|
||||
});
|
||||
let (resource_id, url) = parse_stop_channel_params(&config);
|
||||
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"),
|
||||
@@ -569,9 +571,11 @@ fn test_parse_stop_channel_params_drive() {
|
||||
fn test_parse_stop_channel_params_calendar() {
|
||||
let config = json!({
|
||||
"triggerType": "calendar",
|
||||
"googleChannelId": "chan-xyz",
|
||||
"googleResourceId": "res-456",
|
||||
});
|
||||
let (resource_id, url) = parse_stop_channel_params(&config);
|
||||
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"),
|
||||
@@ -582,9 +586,10 @@ fn test_parse_stop_channel_params_calendar() {
|
||||
|
||||
#[test]
|
||||
fn test_parse_stop_channel_params_default() {
|
||||
// Missing triggerType defaults to Drive
|
||||
// Missing triggerType defaults to Drive; missing googleChannelId yields None.
|
||||
let config = json!({ "googleResourceId": "res-789" });
|
||||
let (resource_id, url) = parse_stop_channel_params(&config);
|
||||
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);
|
||||
}
|
||||
@@ -592,6 +597,7 @@ fn test_parse_stop_channel_params_default() {
|
||||
#[test]
|
||||
fn test_parse_stop_channel_params_missing_resource_id() {
|
||||
let config = json!({ "triggerType": "drive" });
|
||||
let (resource_id, _url) = parse_stop_channel_params(&config);
|
||||
let (channel_id, resource_id, _url) = parse_stop_channel_params(&config);
|
||||
assert!(channel_id.is_none());
|
||||
assert_eq!(resource_id, "");
|
||||
}
|
||||
|
||||
@@ -47,8 +47,10 @@ impl External for Google {
|
||||
db: &DB,
|
||||
_tx: &mut PgConnection,
|
||||
) -> Result<Self::CreateResponse> {
|
||||
// At creation time, channel_id also becomes the trigger's external_id
|
||||
// (see external_id_and_metadata_from_response).
|
||||
let channel_id = uuid::Uuid::new_v4().to_string();
|
||||
self.create_watch_channel(w_id, &channel_id, webhook_token, data, db)
|
||||
self.create_watch_channel(w_id, &channel_id, &channel_id, webhook_token, data, db)
|
||||
.await
|
||||
}
|
||||
|
||||
@@ -65,9 +67,11 @@ impl External for Google {
|
||||
// Google doesn't support updating watch channels — delete old, create new.
|
||||
let _ = self.delete(w_id, oauth_data, external_id, db, tx).await;
|
||||
|
||||
// Reuse the same channel ID so external_id stays permanent
|
||||
// Google rejects reused channel IDs, so we must mint a fresh one.
|
||||
// external_id stays the stable routing key used in the webhook URL.
|
||||
let channel_id = uuid::Uuid::new_v4().to_string();
|
||||
let resp = self
|
||||
.create_watch_channel(w_id, external_id, webhook_token, data, db)
|
||||
.create_watch_channel(w_id, external_id, &channel_id, webhook_token, data, db)
|
||||
.await?;
|
||||
|
||||
self.service_config_from_create_response(data, &resp)
|
||||
@@ -106,11 +110,14 @@ impl External for Google {
|
||||
}
|
||||
let config = config.unwrap();
|
||||
|
||||
let (google_resource_id, url) = super::parse_stop_channel_params(&config);
|
||||
let (google_channel_id, google_resource_id, url) =
|
||||
super::parse_stop_channel_params(&config);
|
||||
// Fall back to external_id for triggers created before googleChannelId was tracked.
|
||||
let channel_id = google_channel_id.unwrap_or_else(|| external_id.to_string());
|
||||
|
||||
if !google_resource_id.is_empty() {
|
||||
let stop_request =
|
||||
StopChannelRequest { id: external_id.to_string(), resource_id: google_resource_id };
|
||||
StopChannelRequest { id: channel_id.clone(), resource_id: google_resource_id };
|
||||
|
||||
// Stop the channel (ignore errors - channel may have already expired)
|
||||
let result: std::result::Result<serde_json::Value, _> = self
|
||||
@@ -118,7 +125,7 @@ impl External for Google {
|
||||
.await;
|
||||
|
||||
if let Err(e) = result {
|
||||
tracing::warn!("Failed to stop Google channel {}: {}", external_id, e);
|
||||
tracing::warn!("Failed to stop Google channel {}: {}", channel_id, e);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -169,6 +176,7 @@ impl External for Google {
|
||||
resp: &Self::CreateResponse,
|
||||
) -> (String, Option<serde_json::Value>) {
|
||||
let metadata = serde_json::json!({
|
||||
"googleChannelId": resp.id,
|
||||
"googleResourceId": resp.resource_id,
|
||||
"expiration": resp.expiration,
|
||||
});
|
||||
@@ -181,6 +189,7 @@ impl External for Google {
|
||||
resp: &Self::CreateResponse,
|
||||
) -> Option<serde_json::Value> {
|
||||
let mut config = data.service_config.clone();
|
||||
config.google_channel_id = Some(resp.id.clone());
|
||||
config.google_resource_id = Some(resp.resource_id.clone());
|
||||
config.expiration = Some(resp.expiration.clone());
|
||||
serde_json::to_value(&config).ok()
|
||||
@@ -194,10 +203,14 @@ impl External for Google {
|
||||
// Helper methods for creating trigger type-specific watches
|
||||
impl Google {
|
||||
/// Build a webhook URL and watch request, then register the channel with Google.
|
||||
/// Used by both `create()` (new UUID) and `update()` (reuse existing external_id).
|
||||
/// `external_id` is the trigger's stable routing key (used in the webhook URL).
|
||||
/// `channel_id` is the identifier sent to Google — it must be globally unique
|
||||
/// across the trigger's lifetime, so callers should mint a fresh UUID per call
|
||||
/// on update/renew.
|
||||
async fn create_watch_channel(
|
||||
&self,
|
||||
w_id: &str,
|
||||
external_id: &str,
|
||||
channel_id: &str,
|
||||
webhook_token: &str,
|
||||
data: &NativeTriggerData<GoogleServiceConfig>,
|
||||
@@ -209,15 +222,16 @@ impl Google {
|
||||
w_id,
|
||||
&data.script_path,
|
||||
data.is_flow,
|
||||
Some(channel_id),
|
||||
Some(external_id),
|
||||
ServiceName::Google,
|
||||
webhook_token,
|
||||
);
|
||||
|
||||
tracing::info!(
|
||||
"Creating Google {} watch channel '{}' with webhook URL: {}",
|
||||
"Creating Google {} watch channel '{}' (external_id={}) with webhook URL: {}",
|
||||
data.service_config.trigger_type,
|
||||
channel_id,
|
||||
external_id,
|
||||
webhook_url
|
||||
);
|
||||
|
||||
@@ -310,7 +324,8 @@ impl Google {
|
||||
|
||||
/// Renew an expiring Google watch channel.
|
||||
/// Rotates the webhook token (creating a new one with the same label),
|
||||
/// stops the old channel and creates a new one with the same channel ID.
|
||||
/// stops the old channel and creates a new one with a fresh channel ID
|
||||
/// (Google rejects reused channel IDs with `channelIdNotUnique`).
|
||||
/// Returns (new_service_config, new_plaintext_token, old_token_hash).
|
||||
/// Callers should delete old_token_hash after successfully updating the trigger.
|
||||
pub async fn renew_channel(
|
||||
@@ -336,42 +351,46 @@ impl Google {
|
||||
}
|
||||
};
|
||||
let base_url = &**BASE_URL.load();
|
||||
// Reuse the same channel ID so external_id stays permanent
|
||||
let channel_id = trigger.external_id.clone();
|
||||
// external_id is our stable routing key; channel_id must be fresh so Google
|
||||
// doesn't reject the new channel with `channelIdNotUnique`.
|
||||
let new_channel_id = uuid::Uuid::new_v4().to_string();
|
||||
// Old channel_id comes from service_config; fall back to external_id for
|
||||
// triggers created before googleChannelId was tracked.
|
||||
let old_channel_id = config
|
||||
.google_channel_id
|
||||
.clone()
|
||||
.unwrap_or_else(|| trigger.external_id.clone());
|
||||
let webhook_url = generate_webhook_service_url(
|
||||
base_url,
|
||||
w_id,
|
||||
&trigger.script_path,
|
||||
trigger.is_flow,
|
||||
Some(&channel_id),
|
||||
Some(&trigger.external_id),
|
||||
ServiceName::Google,
|
||||
&rotated.new_token,
|
||||
);
|
||||
|
||||
tracing::info!(
|
||||
"Renewing Google {} watch channel '{}' with webhook URL: {}",
|
||||
"Renewing Google {} watch channel for '{}': old channel_id={}, new channel_id={}, webhook URL: {}",
|
||||
config.trigger_type,
|
||||
channel_id,
|
||||
trigger.external_id,
|
||||
old_channel_id,
|
||||
new_channel_id,
|
||||
webhook_url
|
||||
);
|
||||
|
||||
let expiration_ms = chrono::Utc::now().timestamp_millis()
|
||||
+ (config.max_expiration_hours() as i64 * 3600 * 1000);
|
||||
let mut watch_request = WatchRequest::new(channel_id.clone(), webhook_url);
|
||||
let mut watch_request = WatchRequest::new(new_channel_id.clone(), webhook_url);
|
||||
watch_request.expiration = Some(expiration_ms);
|
||||
|
||||
// Best-effort stop old channel before creating a new one
|
||||
let old_google_resource_id = trigger
|
||||
.service_config
|
||||
.as_ref()
|
||||
.and_then(|c| c.get("googleResourceId"))
|
||||
.and_then(|r| r.as_str())
|
||||
.unwrap_or_default();
|
||||
let old_google_resource_id = config.google_resource_id.clone().unwrap_or_default();
|
||||
|
||||
if !old_google_resource_id.is_empty() {
|
||||
let stop_request = StopChannelRequest {
|
||||
id: channel_id.clone(),
|
||||
resource_id: old_google_resource_id.to_string(),
|
||||
id: old_channel_id.clone(),
|
||||
resource_id: old_google_resource_id,
|
||||
};
|
||||
let url = match config.trigger_type {
|
||||
GoogleTriggerType::Calendar => {
|
||||
@@ -387,13 +406,13 @@ impl Google {
|
||||
if let Err(e) = result {
|
||||
tracing::warn!(
|
||||
"Failed to stop old Google channel {} during renewal: {}",
|
||||
channel_id,
|
||||
old_channel_id,
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Create new watch channel with the same channel ID
|
||||
// Create new watch channel with a fresh channel ID
|
||||
let resp = match config.trigger_type {
|
||||
GoogleTriggerType::Drive => {
|
||||
self.create_drive_watch(w_id, &config, &watch_request, db)
|
||||
@@ -405,8 +424,9 @@ impl Google {
|
||||
}
|
||||
};
|
||||
|
||||
// Build the updated service_config with new expiration
|
||||
// Build the updated service_config with new channel_id, resource_id, expiration
|
||||
let mut new_config = config;
|
||||
new_config.google_channel_id = Some(resp.id);
|
||||
new_config.google_resource_id = Some(resp.resource_id);
|
||||
new_config.expiration = Some(resp.expiration);
|
||||
|
||||
|
||||
@@ -25,9 +25,16 @@ pub mod routes;
|
||||
|
||||
pub use external::should_renew_channel;
|
||||
|
||||
/// Extracts `(google_resource_id, stop_url)` from a native trigger's service_config JSON.
|
||||
/// Used by the `delete` method and tested independently.
|
||||
pub fn parse_stop_channel_params(config: &serde_json::Value) -> (String, String) {
|
||||
/// Extracts `(google_channel_id, google_resource_id, stop_url)` from a native trigger's
|
||||
/// service_config JSON. `google_channel_id` falls back to `None` for triggers created
|
||||
/// before the field was introduced — callers should use `external_id` in that case.
|
||||
/// Used by the `delete`/renewal paths and tested independently.
|
||||
pub fn parse_stop_channel_params(config: &serde_json::Value) -> (Option<String>, String, String) {
|
||||
let google_channel_id = config
|
||||
.get("googleChannelId")
|
||||
.and_then(|r| r.as_str())
|
||||
.map(String::from);
|
||||
|
||||
let google_resource_id = config
|
||||
.get("googleResourceId")
|
||||
.and_then(|r| r.as_str())
|
||||
@@ -44,7 +51,7 @@ pub fn parse_stop_channel_params(config: &serde_json::Value) -> (String, String)
|
||||
_ => format!("{}/channels/stop", endpoints::DRIVE_API_BASE),
|
||||
};
|
||||
|
||||
(google_resource_id, stop_url)
|
||||
(google_channel_id, google_resource_id, stop_url)
|
||||
}
|
||||
|
||||
/// Handler struct for Google triggers (stateless, used for routing)
|
||||
@@ -96,6 +103,12 @@ pub struct GoogleServiceConfig {
|
||||
/// The resource ID assigned by Google for the watch channel
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub google_resource_id: Option<String>,
|
||||
/// The channel ID currently registered with Google. Regenerated on every
|
||||
/// create/update/renew because Google rejects reused channel IDs with
|
||||
/// `channelIdNotUnique`. Absent on triggers created before this field existed —
|
||||
/// fall back to `external_id` in that case.
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub google_channel_id: Option<String>,
|
||||
/// Channel expiration time (Unix timestamp in milliseconds, as string)
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub expiration: Option<String>,
|
||||
|
||||
Reference in New Issue
Block a user