backbone-integrations 0.6.1

Integration registry: connectors, integration accounts and an idempotent inbound event lane, with one OAuth flow (HMAC-bound state, PKCE)
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
//! `refresh_oauth_credentials` — the refresh-before-expiry scheduled job
//! (hand-authored, user-owned).
//!
//! The declaration of record is the `scheduled_jobs.refresh_oauth_credentials`
//! block in `schema/hooks/index.hook.yaml`; this file is the handler it names.
//! Its declared posture (ADR-0020 vocabulary):
//!
//! - **`posture: pull`** — an interval-driven scan (`*/15 * * * *`). The
//!   interval is the FLOOR, never the contract: the account row's honest
//!   `expires_at` decides when a refresh is due, so a schedule gap degrades
//!   latency, never correctness. Request-path lazy refresh (refresh-on-use
//!   inside the expiry window) rides the same due-window computation and is
//!   available to future consumers without a second code path.
//! - **`pickup_lock: true`** — the claim is `SELECT ... FOR UPDATE SKIP
//!   LOCKED`, one account per short transaction. Two concurrent replicas
//!   (or a manual overlap with the interval) take disjoint accounts instead
//!   of double-refreshing.
//! - **`commit_policy: commit_per_batch`** — each account's refresh commits
//!   independently: claim, exchange, rotate, mirror, commit. One provider
//!   outage rolls back exactly its own account; the rest of the batch is not
//!   held hostage to it. The at-least-once window this opens (ADR-0017) is
//!   bounded by rotate-lineage semantics — a retried refresh rotates again,
//!   it never forks.
//!
//! Per-account flow:
//!
//! 1. **Claim** — the oldest account `status = 'active'` whose mirrored
//!    `expires_at` lands inside the refresh window
//!    (`expires_at < now + refresh_window_seconds`), locked `FOR UPDATE
//!    SKIP LOCKED` in its own transaction. The lock IS the claim: it is held
//!    for the account's full processing, then committed.
//! 2. **Read** — the current token bundle crosses the credential port
//!    ([`OAuthCredentialStore::read_token`]). A scope with no readable
//!    credential (never issued / revoked / expired past honesty) is account
//!    drift: the account moves to `expired` — the reconnect surface — and
//!    the store is left to its own lazy expiry, exactly as an
//!    `invalid_grant` below.
//! 3. **Exchange** — `grant_type=refresh_token` at the provider's VALIDATED
//!    token endpoint through [`OAuthTransport`]. A provider answer with no
//!    `expires_in` is refused as unstoreable (an honest lifetime is the only
//!    kind the flow stores); the account stays due and the next tick
//!    retries. `invalid_grant` (400) is the dead-refresh-token signal: the
//!    account moves to `expired` and the user must reconnect — the
//!    self-heal-to-reconnect behavior, minus the mid-transaction commit of
//!    the code this port replaces.
//! 4. **Rotate** — the successor bundle goes to the store through
//!    [`OAuthCredentialStore::rotate`] (lineage preserved; providers that
//!    hand back a fresh refresh token on every exchange — Microsoft — and
//!    providers that do not — Google — are handled by the same call: the
//!    successor keeps the previous refresh token only when the provider
//!    returned none).
//! 5. **Mirror** — the account row's advisory `expires_at` mirror moves to
//!    the successor's honest expiry and `last_refreshed_at` to now; commit.
//!    A crash between rotate and mirror leaves the mirror stale — the next
//!    tick re-claims the account and rotates again (lineage tolerates it);
//!    expiry TRUTH lives in the store, the mirror only drives scheduling.
//!
//! **Per-scope handler**: the module's own tables carry no company column
//! (ADR-0029) — the composing service's tenancy decorator owns org scoping —
//! but the credential STORE across the port is still company-scoped, so the
//! sweep's `company_id` parameter is the store's scope key: the HOST names it
//! (the job cannot learn it from a row), one sweep per scope via
//! [`refresh_oauth_credentials_for_companies`]. Each sweep transaction binds
//! the composing service's ambient org scope when one is resolved, so a
//! decorated deployment's org fence still applies to the claim, the mirror,
//! and the expire writes; with no scope bound they run unscoped and a
//! decorated deployment refuses them (fail closed).

use std::collections::HashSet;

use chrono::{DateTime, Utc};
use sqlx::{PgPool, Row};
use tracing::{debug, warn};
use uuid::Uuid;

use backbone_orm::org_scope;

use crate::application::service::integrations_oauth_ports::{
    OAuthCredentialFailure, OAuthCredentialStore, TokenBundle, PURPOSE_OAUTH_TOKEN,
};
use crate::infrastructure::http::endpoint_guard::{
    EndpointOverrides, OAuthClientConfigs, OAuthTransport, ProviderRegistry, TokenRequestForm,
    TransportFailureKind, ValidatedEndpoints,
};

/// The due-window and batch knobs (the `oauth.refresh_window_seconds` /
/// `oauth.refresh_batch_size` configuration values). Defaults: refresh when
/// ten minutes remain; at most one hundred accounts per run.
#[derive(Debug, Clone, PartialEq)]
pub struct RefreshSchedule {
    /// Refresh when `expires_at` is less than this many seconds away.
    pub refresh_window_seconds: i64,
    /// Upper bound on accounts refreshed in one run. Not a loss: unclaimed
    /// accounts are still due, so the next tick (manual or interval) picks
    /// them up while their window still matches.
    pub refresh_batch_size: i64,
}

impl Default for RefreshSchedule {
    fn default() -> Self {
        Self { refresh_window_seconds: 600, refresh_batch_size: 100 }
    }
}

/// One run's counters.
#[derive(Debug, Clone, Default, PartialEq)]
pub struct RefreshReport {
    /// Accounts claimed, refreshed, rotated, and mirrored this run.
    pub refreshed: usize,
    /// Accounts moved to `expired` (invalid_grant, or no readable
    /// credential — the reconnect surface).
    pub expired: usize,
    /// Accounts skipped on a retryable failure (store or transport
    /// unreachable, unstoreable provider answer). Left due; next tick
    /// retries.
    pub skipped: usize,
}

/// One due account, as claimed.
struct DueAccount {
    id: Uuid,
    provider: String,
    account_ref: String,
}

/// Build the refresh grant form for one account: the provider's client
/// credentials plus the current refresh token. Omitting `scope` keeps the
/// originally granted scopes (the provider's refresh semantics).
fn refresh_form(
    client_id: &str,
    client_secret: Option<&str>,
    refresh_token: &str,
) -> TokenRequestForm {
    TokenRequestForm {
        grant_type: "refresh_token".into(),
        code: None,
        refresh_token: Some(refresh_token.to_string()),
        redirect_uri: None,
        code_verifier: None,
        client_id: client_id.to_string(),
        client_secret: client_secret.map(str::to_string),
        scope: None,
    }
}

/// Run the refresh sweep for ONE credential-store scope. `company_id` is the
/// legacy tenancy twin (ADR-0029): the module's own tables are unfenced, so
/// it is NOT a claim predicate — it keys the store across the port on every
/// read/rotate below. The composing service's ambient org scope, when
/// resolved, is bound on each claim transaction so a decorated deployment's
/// org fence applies.
pub async fn refresh_oauth_credentials(
    pool: &PgPool,
    company_id: Uuid,
    registry: &ProviderRegistry,
    overrides: &EndpointOverrides,
    clients: &OAuthClientConfigs,
    store: &dyn OAuthCredentialStore,
    transport: &dyn OAuthTransport,
    schedule: &RefreshSchedule,
) -> Result<RefreshReport, sqlx::Error> {
    let mut report = RefreshReport::default();
    // Accounts this run already attempted (refreshed, expired, or skipped).
    // A skipped account is rolled back and stays due — excluding it here
    // means one run attempts each account at most once instead of spinning
    // its whole batch on the same stuck row; the NEXT tick retries it.
    let mut attempted: Vec<Uuid> = Vec::new();
    for _ in 0..schedule.refresh_batch_size.max(0) {
        let mut tx = pool.begin().await?;
        if let Some(scope) = org_scope::current_org_scope() {
            org_scope::bind_org_scope_on(&mut tx, &scope).await?;
        }
        let claimed = sqlx::query(
            r#"SELECT id, provider::text AS provider, account_ref
                 FROM integrations.integration_accounts
                WHERE status = 'active'
                  AND expires_at IS NOT NULL
                  AND expires_at < now() + make_interval(secs => $1)
                  AND id <> ALL($2::uuid[])
                ORDER BY expires_at
                LIMIT 1
                FOR UPDATE SKIP LOCKED"#,
        )
        .bind(schedule.refresh_window_seconds)
        .bind(&attempted)
        .fetch_optional(&mut *tx)
        .await?;
        let row = match claimed {
            Some(row) => row,
            None => {
                // Nothing due (or everything claimable is locked by a
                // concurrent run) — the run is complete.
                tx.rollback().await?;
                break;
            }
        };
        let account = DueAccount {
            id: row.get("id"),
            provider: row.get("provider"),
            account_ref: row.get("account_ref"),
        };
        attempted.push(account.id);
        refresh_one(
            tx, company_id, registry, overrides, clients, store, transport, &account, &mut report,
        )
        .await?;
    }
    Ok(report)
}

/// The host-driven fan-out: run the sweep once per named credential-store
/// scope. The scopes are named by the HOST (the store across the port is
/// still company-scoped, and a job cannot self-enumerate under its fence); a
/// failure for one scope is reported, not fatal to the rest.
pub async fn refresh_oauth_credentials_for_companies(
    pool: &PgPool,
    registry: &ProviderRegistry,
    overrides: &EndpointOverrides,
    clients: &OAuthClientConfigs,
    store: &dyn OAuthCredentialStore,
    transport: &dyn OAuthTransport,
    schedule: &RefreshSchedule,
    companies: &[Uuid],
) -> Vec<(Uuid, Result<RefreshReport, sqlx::Error>)> {
    let mut out = Vec::with_capacity(companies.len());
    for company_id in companies {
        let r = refresh_oauth_credentials(
            pool, *company_id, registry, overrides, clients, store, transport, schedule,
        )
        .await;
        out.push((*company_id, r));
    }
    out
}

/// Refresh ONE claimed account inside its own transaction. The row lock is
/// held for the whole call; every exit path either commits the account's
/// outcome (refreshed / expired) or rolls back (skipped — the account stays
/// due and the next tick retries). This is the `commit_per_batch` grain: one
/// account, one commit. `company_id` is the credential-store scope key (the
/// legacy tenancy twin, ADR-0029).
async fn refresh_one(
    mut tx: sqlx::Transaction<'_, sqlx::Postgres>,
    company_id: Uuid,
    registry: &ProviderRegistry,
    overrides: &EndpointOverrides,
    clients: &OAuthClientConfigs,
    store: &dyn OAuthCredentialStore,
    transport: &dyn OAuthTransport,
    account: &DueAccount,
    report: &mut RefreshReport,
) -> Result<(), sqlx::Error> {
    // The validated endpoints for this provider — the guard runs here too:
    // a drifted or malicious override refuses the account (skipped, still
    // due) before any bytes leave the process.
    let endpoints: ValidatedEndpoints = match ValidatedEndpoints::resolve(
        registry,
        &account.provider,
        overrides,
    ) {
        Ok(endpoints) => endpoints,
        Err(e) => {
            warn!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                provider = %account.provider,
                "endpoint guard refused the provider config; account left due: {e}"
            );
            report.skipped += 1;
            tx.rollback().await?;
            return Ok(());
        }
    };
    let client = match clients.get(&account.provider) {
        Some(client) => client,
        None => {
            warn!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                provider = %account.provider,
                "no OAuth client configured for provider; account left due"
            );
            report.skipped += 1;
            tx.rollback().await?;
            return Ok(());
        }
    };

    // The current bundle through the port. A scope the store cannot read is
    // account drift (or an honestly-expired credential): move the account to
    // `expired` and leave the store to its lazy expiry.
    let bundle: TokenBundle = match store
        .read_token(company_id, &account.provider, &account.account_ref)
        .await
    {
        Ok(bundle) => bundle,
        Err(e) if e.is_transport() => {
            // The store itself was unreachable — retryable, zero writes.
            debug!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                "credential store unreachable; account left due ({e})"
            );
            report.skipped += 1;
            tx.rollback().await?;
            return Ok(());
        }
        Err(OAuthCredentialFailure { code, message }) => {
            warn!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                provider = %account.provider,
                "credential unreadable ({code}: {message}); account moves to expired (reconnect required)"
            );
            expire_account(tx, account.id).await?;
            report.expired += 1;
            return Ok(());
        }
    };

    let refresh_token = match bundle.refresh_token() {
        Some(token) => token.to_string(),
        None => {
            // An access token with no refresh token cannot be renewed — the
            // connection ages out at expiry. Reconnect is the only path.
            warn!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                provider = %account.provider,
                "credential carries no refresh token; account moves to expired (reconnect required)"
            );
            expire_account(tx, account.id).await?;
            report.expired += 1;
            return Ok(());
        }
    };

    // The refresh grant at the validated endpoint.
    let form = refresh_form(&client.client_id, client.client_secret.as_deref(), &refresh_token);
    let now = Utc::now();
    let response = match transport.exchange(&endpoints.token, &form).await {
        Ok(response) => response,
        Err(e) if e.kind == TransportFailureKind::InvalidGrant => {
            // Dead refresh token — the self-heal-to-reconnect path. The
            // credential is left to the store's lazy expiry (no store write
            // from here).
            warn!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                provider = %account.provider,
                "provider rejected the refresh grant (invalid_grant); account moves to expired (reconnect required)"
            );
            expire_account(tx, account.id).await?;
            report.expired += 1;
            return Ok(());
        }
        Err(e) => {
            // Provider refusal / network / guard refusal — retryable, zero
            // writes, account stays due.
            debug!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                "refresh exchange failed; account left due ({e})"
            );
            report.skipped += 1;
            tx.rollback().await?;
            return Ok(());
        }
    };

    // Honest lifetimes only: a response with no usable expiry is unstoreable.
    let expires_at: DateTime<Utc> = match response.expires_at(now) {
        Some(expires_at) => expires_at,
        None => {
            warn!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                provider = %account.provider,
                "provider returned no expires_in; unstoreable response refused, account left due"
            );
            report.skipped += 1;
            tx.rollback().await?;
            return Ok(());
        }
    };

    // The successor bundle: a provider-returned refresh token wins; when the
    // provider returns none (Google), the previous refresh token persists.
    let successor = TokenBundle::new(
        response.access_token.clone(),
        response
            .refresh_token
            .clone()
            .or_else(|| bundle.refresh_token().map(str::to_string)),
        expires_at,
        response.scope.clone().or_else(|| bundle.scope().map(str::to_string)),
    );

    // Rotate through the store (lineage preserved), then mirror.
    match store
        .rotate(
            company_id,
            &account.provider,
            &account.account_ref,
            successor,
            expires_at,
        )
        .await
    {
        Ok(_) => {}
        Err(e) if e.is_transport() => {
            debug!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                "credential store unreachable at rotate; account left due ({e})"
            );
            report.skipped += 1;
            tx.rollback().await?;
            return Ok(());
        }
        Err(OAuthCredentialFailure { code, message }) => {
            warn!(
                target: "integrations.oauth.refresh",
                account_id = %account.id,
                provider = %account.provider,
                "rotate refused ({code}: {message}); account moves to expired"
            );
            expire_account(tx, account.id).await?;
            report.expired += 1;
            return Ok(());
        }
    }

    let mirrored = sqlx::query(
        r#"UPDATE integrations.integration_accounts
              SET expires_at = $2,
                  last_refreshed_at = now()
            WHERE id = $1"#,
    )
    .bind(account.id)
    .bind(expires_at)
    .execute(&mut *tx)
    .await?;
    if mirrored.rows_affected() != 1 {
        warn!(
            target: "integrations.oauth.refresh",
            account_id = %account.id,
            "account mirror affected no rows; rolled back"
        );
        report.skipped += 1;
        tx.rollback().await?;
        return Ok(());
    }
    tx.commit().await?;
    report.refreshed += 1;
    debug!(
        target: "integrations.oauth.refresh",
        account_id = %account.id,
        provider = %account.provider,
        purpose = PURPOSE_OAUTH_TOKEN,
        "refreshed before expiry; successor mirrored"
    );
    Ok(())
}

/// Move an account to the terminal-but-reconnectable `expired` status and
/// commit its transaction.
async fn expire_account(
    mut tx: sqlx::Transaction<'_, sqlx::Postgres>,
    account_id: Uuid,
) -> Result<(), sqlx::Error> {
    sqlx::query(
        r#"UPDATE integrations.integration_accounts
              SET status = 'expired'
            WHERE id = $1"#,
    )
    .bind(account_id)
    .execute(&mut *tx)
    .await?;
    tx.commit().await?;
    Ok(())
}