feat(connect): hosted connector-key registry + validate RPC + entitlements endpoint; instance sync writes connector grants hourly (#1748)
* feat(entitlements): partition the self-host bypass so connector capabilities fall through to grants; capability_grants.source accepts 'connector' Sovereign plan WS3 PR3: ships dark, nothing changes for hosted. - lib/entitlements/keys.ts: CONNECTOR_CAPABILITIES = bank_sync, skatteverket, org_lookup, migration (services Accounted operates that a self-hosted instance cannot provide itself) + isConnectorCapability(). Separate from PAID_CAPABILITIES and outside the trial-seed trigger on purpose: a hosted company can never hold a connector grant. - lib/entitlements/has-capability.ts: isPaywallBypassed() -> isBypassedFor(key). Hosted: byte-identical (dev / DISABLE_PAYWALL bypass, FORCE_PAYWALL wins, else the grant lookup). Self-host: local capabilities always on (FORCE_PAYWALL included, as the existing test demands); connector capabilities behave like hosted, i.e. dev bypass, FORCE_PAYWALL, else the grant lookup where the connector sync will write source='connector' rows. getCompanyEntitlements on a self-host: local paid keys + active connector keys, state 'paid' with an active connector grant else 'none' (never the hosted trial copy). - Migration 20260820122000: capability_grants.source CHECK gains 'connector', found through pg_constraint (the CHECK was declared inline and auto-named; Postgres stores IN as = ANY, matched accordingly). pg-real test: connector accepted, unknown source rejected, upsert on the (scope, key, source) identity, trial seed writes no connector rows. - Tests: self-hosted connector matrix (local all-on without DB, connector gated by grant/expiry, dev bypass all-on, FORCE_PAYWALL gates connector keys only, bulk resolution, entitlements shape); two pre-existing tests that asserted the old "self-host holds connector keys" contract updated to the new one. Verified: full unit suite green, pg-real suite for lib/entitlements green against a local supabase/postgres with every migration applied, lint ratchet, guards. Deferred to the instance-wiring PR: adding the connector extensions to the self-host Docker preset (dead-end upsells until a key can be issued). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * refactor(entitlements): fold the self-host branch into the existing grants query One .or(scopeFilter), not two: the duplicated helper pushed the no-phantom-columns unresolvable-expression count to 380/379. Behaviour is unchanged; the self-host matrix tests still pass. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * feat(connect): hosted connector-key registry + validate RPC + entitlements endpoint; instance sync writes connector grants hourly Sovereign plan WS3 PR4 ("key infra enabling manual sales"), stacked on the entitlement partition (#1747). Nothing is purchasable yet; this is the plumbing both ends need before the first manually issued key. Hosted side: - Migration 20260820123000: connector_keys (SHA-256 key_hash, prefix, org_number, pinned instance_url, scopes, status, Stripe ids, current_period_end, per-minute rate limit, active_company_count, last_seen/synced) and connector_usage_events (per-request metering, separate from metered_events whose company_id references hosted companies). RLS on, NO policies: service role only. RPC validate_and_increment_connector_key copies the api_keys pattern (FOR UPDATE, minute window, suspended reported not counted, revoked = no row) and is REVOKEd from PUBLIC/anon/authenticated, GRANTed to service_role. pg-real test covers validate/count, unknown+revoked, suspended, rate limit, execute privileges per role, RLS invisibility, usage cascade. - lib/connect/contract.ts (shared wire types), lib/connect/hosted/keys.ts (generate/hash/validate -> 401/403/429 mapping), with-connector-auth.ts (Bearer or X-Connector-Key, one usage row per request, 500 envelope on handler throw), /api/connect/entitlements GET + POST (records active_company_count, pins instance_url on first report, never moves a pinned one), scripts/issue-connector-key.ts (dry run unless --confirm, prints the key once + the .env lines). Instance side: - lib/connect/instance/config.ts (GNUBOK_CONNECTOR_KEY, GNUBOK_CONNECT_URL default https://app.gnubok.se), sync.ts: reports the active company count and writes source='connector' grants for every company x covered scope, expires_at = min(now+72h, period_end+3d); 401/403 or a non-active status deletes them (freeze-and-retain); network/5xx/429 leave them alone. /api/connector/sync/cron (hourly) runs it; not_configured without a key. - Crontab generator gains EXTRA_JOBS (variant-only jobs not in vercel.json, with reasons) + drift tests; docker/crontab.self-hosted regenerated with the hourly sync. Docs (SELF-HOSTING connector section, env templates), DECISIONS. Tests: 52 new unit tests (keys, auth wrapper, route, config, sync outcomes and grant arithmetic, cron route, crontab EXTRA_JOBS) + 7 pg-real tests run locally against supabase/postgres with every migration applied. no-phantom-columns ceiling +1 with a reason (the bulk grant upsert). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * chore(connect): update pg test to re-versioned migration 20260831190000 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzNkSsR18pLFitJdYn8QEb * fix(connect): RPC errors answer 503 not 401; X-Connector-Key wins over Authorization A hosted DB error mapped to 401 made the instance sync treat a pooler blip as key revocation and delete its entire connector grant cache, zeroing the 72h offline grace. 503 lands in the sync's keep-grants branch (already test-pinned). Bearer-first extraction hashed the upstream token on dual-header proxied calls, 401ing the exact shape X-Connector-Key exists for. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzNkSsR18pLFitJdYn8QEb * fix(connect): sync deletes grants only on a body-proven connector rejection, never bare 401/403 A WAF challenge page, edge deployment protection, or an egress proxy answers 401/403 without the hosted app ever running; trusting status alone wiped the instance's 72h offline grant cache within the hour. Deletion now requires the hosted route's own rejection code in the JSON body; codeless 401/403 keeps grants (server_error branch). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzNkSsR18pLFitJdYn8QEb * fix(connect): PR #1748 review batch: https-only connect URL, atomic pin, prefix-gated Bearer, deferred metering, entitlements validation, integer months - GNUBOK_CONNECT_URL must be https (http only for loopback); invalid or plaintext URLs disable the connector instead of sending the key. - instance_url pin update filters on IS NULL; a lost race re-reads and reports the winner's pin. - extractConnectorKey: a Bearer is the connector credential only with the gnubok_ck_ prefix; upstream Bearer falls through to X-Connector-Key. - Usage metering runs via after() off the response path (inline outside a request scope). - Sync validates entitlements shape: unknown status or malformed current_period_end keeps grants (server_error), never deletes. - issue-connector-key rejects fractional --months. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzNkSsR18pLFitJdYn8QEb --------- Co-authored-by: Jakob Wennberg <311770904+jakobwennberg-oss@users.noreply.github.com> Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: Emil <emilmattsson14@gmail.com>
This commit is contained in:
co-authored by
Claude Fable 5
Jakob Wennberg
Emil
parent
cfce2de925
commit
0ff1b05553
@@ -0,0 +1,45 @@
|
||||
/**
|
||||
* The wire contract between a self-hosted Accounted instance and the hosted
|
||||
* connector service (app.gnubok.se/api/connect/*). Shared by both sides so
|
||||
* the instance sync and the hosted endpoint cannot drift apart.
|
||||
*
|
||||
* Background: a self-hosted instance runs everything itself except the
|
||||
* services only Accounted can operate (bank sync via our PSD2/AISP
|
||||
* credentials, the Skatteverket API client, TIC org lookup, the migration
|
||||
* gateway). A connector key (`gnubok_ck_...`) is the subscription token for
|
||||
* those; this is the Nabu Casa model: everything local is free AGPL, the key
|
||||
* buys access to the hosted connectors, and enforcement is key auth at the
|
||||
* hosted proxy, never a licence check inside the instance.
|
||||
*/
|
||||
|
||||
export const CONNECTOR_KEY_PREFIX = 'gnubok_ck_'
|
||||
|
||||
/** Header alternative to `Authorization: Bearer`, for proxied calls where Authorization carries an upstream token. */
|
||||
export const CONNECTOR_KEY_HEADER = 'x-connector-key'
|
||||
|
||||
export const CONNECTOR_ENTITLEMENTS_PATH = '/api/connect/entitlements'
|
||||
|
||||
/** Default hosted origin. app.gnubok.se stays the machine-facing host for API traffic. */
|
||||
export const DEFAULT_CONNECT_BASE_URL = 'https://app.gnubok.se'
|
||||
|
||||
export type ConnectorKeyStatus = 'active' | 'suspended' | 'revoked'
|
||||
|
||||
/** What the hosted service tells an instance about its key. */
|
||||
export interface ConnectorEntitlements {
|
||||
status: ConnectorKeyStatus
|
||||
/** Capability keys the subscription covers (subset of CONNECTOR_CAPABILITIES). */
|
||||
scopes: string[]
|
||||
/** End of the paid period, ISO; null for an open-ended (manually issued) key. */
|
||||
current_period_end: string | null
|
||||
org_number: string
|
||||
/** The instance origin this key is pinned to; null until the first sync claims it. */
|
||||
instance_url: string | null
|
||||
server_time: string
|
||||
}
|
||||
|
||||
/** What an instance reports on every sync (quantity billing input). */
|
||||
export interface ConnectorSyncReport {
|
||||
active_company_count: number
|
||||
instance_url?: string
|
||||
app_version?: string
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
import { describe, it, expect, vi } from 'vitest'
|
||||
import type { SupabaseClient } from '@supabase/supabase-js'
|
||||
import { generateConnectorKey, hashConnectorKey, isConnectorKeyFormat, validateConnectorKey } from '../keys'
|
||||
|
||||
function supabaseWithRpc(result: { data?: unknown; error?: unknown }): { supabase: SupabaseClient; rpc: ReturnType<typeof vi.fn> } {
|
||||
const rpc = vi.fn().mockResolvedValue({ data: result.data ?? null, error: result.error ?? null })
|
||||
return { supabase: { rpc } as unknown as SupabaseClient, rpc }
|
||||
}
|
||||
|
||||
const ROW = {
|
||||
connector_key_id: '11111111-1111-4111-8111-111111111111',
|
||||
org_number: '5561234567',
|
||||
instance_url: 'https://bokforing.example.se',
|
||||
scopes: ['bank_sync', 'skatteverket'],
|
||||
status: 'active',
|
||||
current_period_end: '2027-01-01T00:00:00.000Z',
|
||||
rate_limited: false,
|
||||
}
|
||||
|
||||
describe('connector key primitives', () => {
|
||||
it('generates a gnubok_ck_ key with a display prefix and a SHA-256 hash', () => {
|
||||
const { key, hash, prefix } = generateConnectorKey()
|
||||
expect(key.startsWith('gnubok_ck_')).toBe(true)
|
||||
expect(key.length).toBeGreaterThan(40)
|
||||
expect(prefix).toBe(key.slice(0, 18))
|
||||
expect(hash).toBe(hashConnectorKey(key))
|
||||
expect(hash).toMatch(/^[0-9a-f]{64}$/)
|
||||
expect(generateConnectorKey().key).not.toBe(key)
|
||||
})
|
||||
|
||||
it('recognises the key format', () => {
|
||||
expect(isConnectorKeyFormat(generateConnectorKey().key)).toBe(true)
|
||||
expect(isConnectorKeyFormat('gnubok_sk_abcdefghijklmnopqrstuvwxyz')).toBe(false)
|
||||
expect(isConnectorKeyFormat('gnubok_ck_short')).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
describe('validateConnectorKey', () => {
|
||||
it('rejects a malformed key without touching the database', async () => {
|
||||
const { supabase, rpc } = supabaseWithRpc({ data: [ROW] })
|
||||
expect(await validateConnectorKey('nope', supabase)).toMatchObject({ ok: false, status: 401, code: 'CONNECTOR_KEY_INVALID' })
|
||||
expect(rpc).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('hashes the key and calls the atomic RPC', async () => {
|
||||
const { key, hash } = generateConnectorKey()
|
||||
const { supabase, rpc } = supabaseWithRpc({ data: [ROW] })
|
||||
const result = await validateConnectorKey(key, supabase)
|
||||
expect(rpc).toHaveBeenCalledWith('validate_and_increment_connector_key', { p_key_hash: hash })
|
||||
expect(result).toEqual({
|
||||
ok: true,
|
||||
key: {
|
||||
id: ROW.connector_key_id,
|
||||
orgNumber: '5561234567',
|
||||
instanceUrl: 'https://bokforing.example.se',
|
||||
scopes: ['bank_sync', 'skatteverket'],
|
||||
status: 'active',
|
||||
currentPeriodEnd: '2027-01-01T00:00:00.000Z',
|
||||
},
|
||||
})
|
||||
})
|
||||
|
||||
it('maps no row (unknown/revoked) to 401, but an RPC error to 503', async () => {
|
||||
const { key } = generateConnectorKey()
|
||||
expect(await validateConnectorKey(key, supabaseWithRpc({ data: [] }).supabase)).toMatchObject({ ok: false, status: 401 })
|
||||
// NEVER 401 on a database error: the instance sync deletes its whole
|
||||
// connector grant cache on 401/403, so a hosted pooler blip answered as
|
||||
// 401 would destroy a paying instance's 72h offline grace. 503 lands in
|
||||
// the sync's keep-grants branch.
|
||||
expect(await validateConnectorKey(key, supabaseWithRpc({ error: { message: 'boom' } }).supabase)).toMatchObject({
|
||||
ok: false,
|
||||
status: 503,
|
||||
code: 'CONNECTOR_VALIDATION_UNAVAILABLE',
|
||||
})
|
||||
})
|
||||
|
||||
it('maps a suspended key to 403 and a rate-limited one to 429', async () => {
|
||||
const { key } = generateConnectorKey()
|
||||
expect(
|
||||
await validateConnectorKey(key, supabaseWithRpc({ data: [{ ...ROW, status: 'suspended' }] }).supabase),
|
||||
).toMatchObject({ ok: false, status: 403, code: 'CONNECTOR_KEY_SUSPENDED' })
|
||||
expect(
|
||||
await validateConnectorKey(key, supabaseWithRpc({ data: [{ ...ROW, rate_limited: true }] }).supabase),
|
||||
).toMatchObject({ ok: false, status: 429, code: 'CONNECTOR_RATE_LIMITED' })
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,108 @@
|
||||
import { describe, it, expect, vi, beforeEach } from 'vitest'
|
||||
import { NextResponse } from 'next/server'
|
||||
|
||||
const validateMock = vi.fn()
|
||||
vi.mock('../keys', () => ({
|
||||
validateConnectorKey: (...args: unknown[]) => validateMock(...args),
|
||||
}))
|
||||
|
||||
const usageInsert = vi.fn()
|
||||
const from = vi.fn(() => ({ insert: usageInsert }))
|
||||
vi.mock('@/lib/auth/api-keys', () => ({
|
||||
createServiceClientNoCookies: () => ({ from }),
|
||||
}))
|
||||
|
||||
import { extractConnectorKey, withConnectorAuth } from '../with-connector-auth'
|
||||
|
||||
const VALID = {
|
||||
ok: true,
|
||||
key: {
|
||||
id: '11111111-1111-4111-8111-111111111111',
|
||||
orgNumber: '5561234567',
|
||||
instanceUrl: null,
|
||||
scopes: ['bank_sync'],
|
||||
status: 'active',
|
||||
currentPeriodEnd: null,
|
||||
},
|
||||
}
|
||||
|
||||
function req(headers: Record<string, string> = {}): Request {
|
||||
return new Request('https://app.gnubok.se/api/connect/entitlements', { headers })
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
usageInsert.mockResolvedValue({ error: null })
|
||||
})
|
||||
|
||||
describe('extractConnectorKey', () => {
|
||||
it('takes a connector-prefixed Bearer, else X-Connector-Key, else the raw Bearer for the format-check 401', () => {
|
||||
expect(extractConnectorKey(req({ authorization: 'Bearer gnubok_ck_a' }))).toBe('gnubok_ck_a')
|
||||
expect(extractConnectorKey(req({ 'x-connector-key': 'gnubok_ck_b' }))).toBe('gnubok_ck_b')
|
||||
// A connector-prefixed Bearer is unambiguous and wins.
|
||||
expect(extractConnectorKey(req({ authorization: 'Bearer gnubok_ck_a', 'x-connector-key': 'gnubok_ck_b' }))).toBe('gnubok_ck_a')
|
||||
// The proxied-call shape (SKV data proxy): Authorization is the user's
|
||||
// upstream SKV token, X-Connector-Key authenticates the instance. The
|
||||
// connector key MUST win or every such request 401s on a hashed upstream
|
||||
// token.
|
||||
expect(extractConnectorKey(req({ authorization: 'Bearer upstream-skv-token', 'x-connector-key': 'gnubok_ck_b' }))).toBe('gnubok_ck_b')
|
||||
// Non-prefixed Bearer alone still reaches the format check (401).
|
||||
expect(extractConnectorKey(req({ authorization: 'Bearer not-a-connector-key' }))).toBe('not-a-connector-key')
|
||||
expect(extractConnectorKey(req())).toBeNull()
|
||||
expect(extractConnectorKey(req({ authorization: 'Basic xyz' }))).toBeNull()
|
||||
})
|
||||
})
|
||||
|
||||
describe('withConnectorAuth', () => {
|
||||
const handler = vi.fn(async (_req: Request, _ctx: { key: { id: string } }) => NextResponse.json({ data: 'ok' }))
|
||||
const wrapped = withConnectorAuth('connect.entitlements', handler)
|
||||
|
||||
it('401 without a key, and never calls the handler', async () => {
|
||||
const res = await wrapped(req())
|
||||
expect(res.status).toBe(401)
|
||||
expect(await res.json()).toMatchObject({ code: 'CONNECTOR_KEY_MISSING' })
|
||||
expect(handler).not.toHaveBeenCalled()
|
||||
expect(validateMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('passes the validation failure through (401/403/429) with its code', async () => {
|
||||
for (const failure of [
|
||||
{ ok: false, status: 401, code: 'CONNECTOR_KEY_INVALID', error: 'Invalid connector key' },
|
||||
{ ok: false, status: 403, code: 'CONNECTOR_KEY_SUSPENDED', error: 'Connector key is suspended' },
|
||||
{ ok: false, status: 429, code: 'CONNECTOR_RATE_LIMITED', error: 'Rate limit exceeded' },
|
||||
]) {
|
||||
validateMock.mockResolvedValueOnce(failure)
|
||||
const res = await wrapped(req({ authorization: 'Bearer gnubok_ck_x' }))
|
||||
expect(res.status).toBe(failure.status)
|
||||
expect(await res.json()).toEqual({ error: failure.error, code: failure.code })
|
||||
}
|
||||
expect(handler).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('runs the handler with the validated key and records one usage event', async () => {
|
||||
validateMock.mockResolvedValueOnce(VALID)
|
||||
const res = await wrapped(req({ authorization: 'Bearer gnubok_ck_x' }))
|
||||
expect(res.status).toBe(200)
|
||||
expect(res.headers.get('X-Request-Id')).toMatch(/^conn_/)
|
||||
expect(handler).toHaveBeenCalledTimes(1)
|
||||
expect(handler.mock.calls[0][1].key.id).toBe(VALID.key.id)
|
||||
expect(from).toHaveBeenCalledWith('connector_usage_events')
|
||||
expect(usageInsert).toHaveBeenCalledWith({
|
||||
connector_key_id: VALID.key.id,
|
||||
service: 'entitlements',
|
||||
endpoint: '/api/connect/entitlements',
|
||||
status_code: 200,
|
||||
})
|
||||
})
|
||||
|
||||
it('turns a throwing handler into a 500 envelope and still meters it', async () => {
|
||||
validateMock.mockResolvedValueOnce(VALID)
|
||||
const boom = withConnectorAuth('connect.entitlements', async () => {
|
||||
throw new Error('db down')
|
||||
})
|
||||
const res = await boom(req({ authorization: 'Bearer gnubok_ck_x' }))
|
||||
expect(res.status).toBe(500)
|
||||
expect(await res.json()).toEqual({ error: 'Internal error', code: 'INTERNAL_ERROR' })
|
||||
expect(usageInsert).toHaveBeenCalledWith(expect.objectContaining({ status_code: 500 }))
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,101 @@
|
||||
import crypto from 'node:crypto'
|
||||
import type { SupabaseClient } from '@supabase/supabase-js'
|
||||
import { createServiceClientNoCookies } from '@/lib/auth/api-keys'
|
||||
import { CONNECTOR_KEY_PREFIX, type ConnectorKeyStatus } from '../contract'
|
||||
|
||||
/**
|
||||
* Hosted-side connector key primitives. Mirrors lib/auth/api-keys.ts for
|
||||
* `gnubok_sk_` API keys: 32 CSPRNG bytes, SHA-256 at rest (the hash is the
|
||||
* lookup; a slow KDF adds nothing on a 256-bit random secret and would sit on
|
||||
* the hot path of every proxied request), atomic validate + rate-limit in a
|
||||
* SECURITY DEFINER RPC that only service_role may execute.
|
||||
*/
|
||||
|
||||
export function generateConnectorKey(): { key: string; hash: string; prefix: string } {
|
||||
const random = crypto.randomBytes(32).toString('base64url')
|
||||
const key = `${CONNECTOR_KEY_PREFIX}${random}`
|
||||
return { key, hash: hashConnectorKey(key), prefix: key.slice(0, CONNECTOR_KEY_PREFIX.length + 8) }
|
||||
}
|
||||
|
||||
export function hashConnectorKey(key: string): string {
|
||||
return crypto.createHash('sha256').update(key).digest('hex')
|
||||
}
|
||||
|
||||
export function isConnectorKeyFormat(key: string): boolean {
|
||||
return key.startsWith(CONNECTOR_KEY_PREFIX) && key.length > CONNECTOR_KEY_PREFIX.length + 16
|
||||
}
|
||||
|
||||
export interface ValidatedConnectorKey {
|
||||
id: string
|
||||
orgNumber: string
|
||||
instanceUrl: string | null
|
||||
scopes: string[]
|
||||
status: ConnectorKeyStatus
|
||||
currentPeriodEnd: string | null
|
||||
}
|
||||
|
||||
export type ConnectorKeyValidation =
|
||||
| { ok: true; key: ValidatedConnectorKey }
|
||||
| { ok: false; status: 401 | 403 | 429; code: 'CONNECTOR_KEY_INVALID' | 'CONNECTOR_KEY_SUSPENDED' | 'CONNECTOR_RATE_LIMITED'; error: string }
|
||||
| { ok: false; status: 503; code: 'CONNECTOR_VALIDATION_UNAVAILABLE'; error: string }
|
||||
|
||||
/**
|
||||
* Validate a presented key: format check, RPC lookup (atomic rate-limit
|
||||
* increment), status mapping. Never throws on a bad key.
|
||||
*
|
||||
* A database/RPC error maps to 503, NEVER 401: the instance-side sync treats
|
||||
* 401/403 as key revocation and deletes its entire connector grant cache
|
||||
* (lib/connect/instance/sync.ts), so answering a hosted pooler blip with 401
|
||||
* would let a transient hosted incident destroy a paying instance's 72h
|
||||
* offline grace. 503 lands in the sync's keep-grants branch. Only a genuine
|
||||
* empty result (unknown or revoked key) is 401. This deliberately diverges
|
||||
* from the api-keys precedent, where a spurious 401 costs one request.
|
||||
*/
|
||||
export async function validateConnectorKey(
|
||||
key: string,
|
||||
supabase: SupabaseClient = createServiceClientNoCookies(),
|
||||
): Promise<ConnectorKeyValidation> {
|
||||
if (!isConnectorKeyFormat(key)) {
|
||||
return { ok: false, status: 401, code: 'CONNECTOR_KEY_INVALID', error: 'Invalid connector key' }
|
||||
}
|
||||
const { data, error } = await supabase.rpc('validate_and_increment_connector_key', {
|
||||
p_key_hash: hashConnectorKey(key),
|
||||
})
|
||||
if (error) {
|
||||
return {
|
||||
ok: false,
|
||||
status: 503,
|
||||
code: 'CONNECTOR_VALIDATION_UNAVAILABLE',
|
||||
error: 'Connector key validation temporarily unavailable',
|
||||
}
|
||||
}
|
||||
if (!data || (Array.isArray(data) && data.length === 0)) {
|
||||
return { ok: false, status: 401, code: 'CONNECTOR_KEY_INVALID', error: 'Invalid connector key' }
|
||||
}
|
||||
const row = (Array.isArray(data) ? data[0] : data) as {
|
||||
connector_key_id: string
|
||||
org_number: string
|
||||
instance_url: string | null
|
||||
scopes: string[] | null
|
||||
status: string
|
||||
current_period_end: string | null
|
||||
rate_limited: boolean
|
||||
}
|
||||
if (row.status !== 'active') {
|
||||
return { ok: false, status: 403, code: 'CONNECTOR_KEY_SUSPENDED', error: 'Connector key is suspended' }
|
||||
}
|
||||
if (row.rate_limited) {
|
||||
return { ok: false, status: 429, code: 'CONNECTOR_RATE_LIMITED', error: 'Rate limit exceeded' }
|
||||
}
|
||||
return {
|
||||
ok: true,
|
||||
key: {
|
||||
id: row.connector_key_id,
|
||||
orgNumber: row.org_number,
|
||||
instanceUrl: row.instance_url,
|
||||
scopes: row.scopes ?? [],
|
||||
status: 'active',
|
||||
currentPeriodEnd: row.current_period_end,
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,115 @@
|
||||
import type { SupabaseClient } from '@supabase/supabase-js'
|
||||
import { NextResponse, after } from 'next/server'
|
||||
import { createServiceClientNoCookies } from '@/lib/auth/api-keys'
|
||||
import { createLogger, type Logger } from '@/lib/logger'
|
||||
import { CONNECTOR_KEY_HEADER, CONNECTOR_KEY_PREFIX } from '../contract'
|
||||
import { validateConnectorKey, type ValidatedConnectorKey } from './keys'
|
||||
|
||||
/**
|
||||
* Wrapper for the hosted connector endpoints (app/api/connect/*), the
|
||||
* connector twin of lib/api/v1/with-api-v1.ts.
|
||||
*
|
||||
* 1. Extracts the key from `Authorization: Bearer gnubok_ck_...` or, for
|
||||
* proxied calls where Authorization carries an upstream token, from
|
||||
* `X-Connector-Key`.
|
||||
* 2. Validates it through the atomic RPC (rate-limited): 401 unknown /
|
||||
* revoked, 403 suspended, 429 over the per-minute limit.
|
||||
* 3. Runs the handler with a service-role client and the validated key.
|
||||
* 4. Records one connector_usage_events row (metering, never blocking).
|
||||
*
|
||||
* These routes are reached by self-hosted instances with a key, never by a
|
||||
* browser session, so there is no cookie/MFA handling here; the proxy
|
||||
* middleware already lets /api/* through for bearer callers.
|
||||
*/
|
||||
|
||||
export interface ConnectorContext {
|
||||
requestId: string
|
||||
log: Logger
|
||||
supabase: SupabaseClient
|
||||
key: ValidatedConnectorKey
|
||||
}
|
||||
|
||||
type ConnectorHandler = (request: Request, ctx: ConnectorContext) => Promise<NextResponse | Response>
|
||||
|
||||
export function extractConnectorKey(request: Request): string | null {
|
||||
// A Bearer value is the connector credential only when it looks like one
|
||||
// (gnubok_ck_ prefix); otherwise it is an UPSTREAM token on a proxied call
|
||||
// and the connector key rides in X-Connector-Key. Hashing the upstream
|
||||
// token instead would 401 every such request. A non-prefixed Bearer with
|
||||
// no X-Connector-Key still falls through to the format check's 401.
|
||||
const auth = request.headers.get('authorization')
|
||||
const bearer = auth?.startsWith('Bearer ') ? auth.slice(7).trim() || null : null
|
||||
if (bearer?.startsWith(CONNECTOR_KEY_PREFIX)) return bearer
|
||||
const header = request.headers.get(CONNECTOR_KEY_HEADER)
|
||||
if (header?.trim()) return header.trim()
|
||||
return bearer
|
||||
}
|
||||
|
||||
export function withConnectorAuth(
|
||||
operation: string,
|
||||
handler: ConnectorHandler,
|
||||
options: { service?: string } = {},
|
||||
): (request: Request) => Promise<Response> {
|
||||
const service = options.service ?? operation.split('.')[1] ?? operation
|
||||
return async function wrapped(request: Request): Promise<Response> {
|
||||
const requestId = `conn_${crypto.randomUUID()}`
|
||||
const log = createLogger(`connect/${operation}`, { requestId, operation })
|
||||
const supabase = createServiceClientNoCookies()
|
||||
|
||||
const presented = extractConnectorKey(request)
|
||||
if (!presented) {
|
||||
return NextResponse.json(
|
||||
{ error: 'Missing connector key', code: 'CONNECTOR_KEY_MISSING' },
|
||||
{ status: 401, headers: { 'X-Request-Id': requestId } },
|
||||
)
|
||||
}
|
||||
const validation = await validateConnectorKey(presented, supabase)
|
||||
if (!validation.ok) {
|
||||
log.warn('connector auth rejected', { code: validation.code })
|
||||
return NextResponse.json(
|
||||
{ error: validation.error, code: validation.code },
|
||||
{ status: validation.status, headers: { 'X-Request-Id': requestId } },
|
||||
)
|
||||
}
|
||||
|
||||
let response: Response
|
||||
try {
|
||||
response = await handler(request, { requestId, log, supabase, key: validation.key })
|
||||
} catch (err) {
|
||||
log.error('connector handler failed', err)
|
||||
response = NextResponse.json(
|
||||
{ error: 'Internal error', code: 'INTERNAL_ERROR' },
|
||||
{ status: 500 },
|
||||
)
|
||||
}
|
||||
response.headers.set('X-Request-Id', requestId)
|
||||
|
||||
// Metering: one row per request, never on the critical path.
|
||||
const endpoint = (() => {
|
||||
try {
|
||||
return new URL(request.url).pathname
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
})()
|
||||
const recordUsage = async (): Promise<void> => {
|
||||
const { error: usageError } = await supabase.from('connector_usage_events').insert({
|
||||
connector_key_id: validation.key.id,
|
||||
service,
|
||||
endpoint,
|
||||
status_code: response.status,
|
||||
})
|
||||
if (usageError) log.warn('usage event not recorded', { err: usageError.message })
|
||||
}
|
||||
try {
|
||||
// Off the response path: the caller should not wait on metering.
|
||||
// Same pattern as lib/webhooks/dispatch-kick.ts.
|
||||
after(() => recordUsage())
|
||||
} catch {
|
||||
// Outside a request scope (unit tests): record inline.
|
||||
await recordUsage()
|
||||
}
|
||||
|
||||
return response
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
import { describe, it, expect, afterEach, vi } from 'vitest'
|
||||
import { getConnectorConfig, isConnectorConfigured } from '../config'
|
||||
|
||||
afterEach(() => vi.unstubAllEnvs())
|
||||
|
||||
describe('getConnectorConfig', () => {
|
||||
it('is null without a key (hosted, or a self-host without a subscription)', () => {
|
||||
vi.stubEnv('GNUBOK_CONNECTOR_KEY', '')
|
||||
expect(getConnectorConfig()).toBeNull()
|
||||
expect(isConnectorConfigured()).toBe(false)
|
||||
})
|
||||
|
||||
it('defaults the hosted origin to app.gnubok.se and strips trailing slashes from an override', () => {
|
||||
vi.stubEnv('GNUBOK_CONNECTOR_KEY', 'gnubok_ck_x')
|
||||
vi.stubEnv('GNUBOK_CONNECT_URL', '')
|
||||
expect(getConnectorConfig()).toEqual({ key: 'gnubok_ck_x', baseUrl: 'https://app.gnubok.se' })
|
||||
vi.stubEnv('GNUBOK_CONNECT_URL', 'https://connect.example.se/')
|
||||
expect(getConnectorConfig()?.baseUrl).toBe('https://connect.example.se')
|
||||
})
|
||||
|
||||
it('rejects non-https and malformed GNUBOK_CONNECT_URL (fail closed: the key is never sent in plaintext)', () => {
|
||||
vi.stubEnv('GNUBOK_CONNECTOR_KEY', 'gnubok_ck_x')
|
||||
for (const bad of ['http://connect.example.se', 'ftp://connect.example.se', 'not a url', 'connect.example.se']) {
|
||||
vi.stubEnv('GNUBOK_CONNECT_URL', bad)
|
||||
expect(getConnectorConfig(), bad).toBeNull()
|
||||
expect(isConnectorConfigured(), bad).toBe(false)
|
||||
}
|
||||
})
|
||||
|
||||
it('allows plain http for loopback development hosts only', () => {
|
||||
vi.stubEnv('GNUBOK_CONNECTOR_KEY', 'gnubok_ck_x')
|
||||
for (const ok of ['http://localhost:3000', 'http://127.0.0.1:3000']) {
|
||||
vi.stubEnv('GNUBOK_CONNECT_URL', ok)
|
||||
expect(getConnectorConfig()?.baseUrl, ok).toBe(ok)
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,254 @@
|
||||
import { describe, it, expect, vi, beforeEach } from 'vitest'
|
||||
import type { SupabaseClient } from '@supabase/supabase-js'
|
||||
import { connectorGrantExpiry, syncConnectorEntitlements, CONNECTOR_GRANT_TTL_MS } from '../sync'
|
||||
|
||||
/**
|
||||
* A purpose-built Supabase mock: `companies` answers the paginated read,
|
||||
* `capability_grants` records upserts and deletes (with their filters), so the
|
||||
* grant arithmetic can be asserted exactly.
|
||||
*/
|
||||
function makeSupabase(companyIds: string[]) {
|
||||
const upserts: Array<{ rows: unknown[]; opts: unknown }> = []
|
||||
const deletes: Array<{ filters: Array<[string, ...unknown[]]> }> = []
|
||||
let deleteCount = 0
|
||||
|
||||
const grantsChain = () => {
|
||||
const deleteRecord: { filters: Array<[string, ...unknown[]]> } = { filters: [] }
|
||||
const chain: Record<string, unknown> = {}
|
||||
chain.upsert = (rows: unknown[], opts: unknown) => {
|
||||
upserts.push({ rows, opts })
|
||||
return Promise.resolve({ error: null })
|
||||
}
|
||||
chain.delete = () => {
|
||||
deletes.push(deleteRecord)
|
||||
const dchain: Record<string, unknown> = {
|
||||
eq: (...a: unknown[]) => {
|
||||
deleteRecord.filters.push(['eq', ...a])
|
||||
return dchain
|
||||
},
|
||||
not: (...a: unknown[]) => {
|
||||
deleteRecord.filters.push(['not', ...a])
|
||||
return dchain
|
||||
},
|
||||
then: (resolve: (v: unknown) => void) => resolve({ error: null, count: deleteCount }),
|
||||
}
|
||||
return dchain
|
||||
}
|
||||
return chain
|
||||
}
|
||||
const companiesChain = () => {
|
||||
let rangeFrom = 0
|
||||
const chain: Record<string, unknown> = {
|
||||
select: () => chain,
|
||||
is: () => chain,
|
||||
order: () => chain,
|
||||
range: (from: number) => {
|
||||
rangeFrom = from
|
||||
return chain
|
||||
},
|
||||
then: (resolve: (v: unknown) => void) =>
|
||||
resolve({ data: companyIds.slice(rangeFrom, rangeFrom + 1000).map((id) => ({ id })), error: null }),
|
||||
}
|
||||
return chain
|
||||
}
|
||||
const supabase = {
|
||||
from: (table: string) => (table === 'companies' ? companiesChain() : grantsChain()),
|
||||
} as unknown as SupabaseClient
|
||||
return {
|
||||
supabase,
|
||||
upserts,
|
||||
deletes,
|
||||
setDeleteCount: (n: number) => {
|
||||
deleteCount = n
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
const CONFIG = { key: 'gnubok_ck_test', baseUrl: 'https://app.gnubok.se' }
|
||||
const NOW = new Date('2026-08-20T12:00:00.000Z')
|
||||
const C1 = '11111111-1111-4111-8111-111111111111'
|
||||
const C2 = '22222222-2222-4222-8222-222222222222'
|
||||
|
||||
function jsonResponse(status: number, body: unknown): Response {
|
||||
return new Response(JSON.stringify(body), { status, headers: { 'content-type': 'application/json' } })
|
||||
}
|
||||
|
||||
beforeEach(() => vi.clearAllMocks())
|
||||
|
||||
describe('connectorGrantExpiry', () => {
|
||||
it('is now + 72h without a period end, and the earlier of that and period_end + 3d otherwise', () => {
|
||||
expect(connectorGrantExpiry(NOW, null)).toBe(new Date(NOW.getTime() + CONNECTOR_GRANT_TTL_MS).toISOString())
|
||||
// period ends in 10 days: 72h wins
|
||||
expect(connectorGrantExpiry(NOW, '2026-08-30T00:00:00.000Z')).toBe('2026-08-23T12:00:00.000Z')
|
||||
// period ended yesterday: period_end + 3d wins (grace), earlier than 72h
|
||||
expect(connectorGrantExpiry(NOW, '2026-08-19T12:00:00.000Z')).toBe('2026-08-22T12:00:00.000Z')
|
||||
})
|
||||
})
|
||||
|
||||
describe('syncConnectorEntitlements', () => {
|
||||
it('is a no-op without a key', async () => {
|
||||
const { supabase, upserts } = makeSupabase([C1])
|
||||
const fetchImpl = vi.fn()
|
||||
const result = await syncConnectorEntitlements(supabase, { config: null, fetchImpl })
|
||||
expect(result.outcome).toBe('not_configured')
|
||||
expect(fetchImpl).not.toHaveBeenCalled()
|
||||
expect(upserts).toHaveLength(0)
|
||||
})
|
||||
|
||||
it('reports the active company count with the key and writes one connector grant per company and scope', async () => {
|
||||
const { supabase, upserts, deletes } = makeSupabase([C1, C2])
|
||||
const fetchImpl = vi.fn().mockResolvedValue(
|
||||
jsonResponse(200, {
|
||||
data: {
|
||||
status: 'active',
|
||||
scopes: ['bank_sync', 'skatteverket', 'not_a_connector_scope'],
|
||||
current_period_end: '2027-01-01T00:00:00.000Z',
|
||||
org_number: '5561234567',
|
||||
instance_url: 'https://bokforing.example.se',
|
||||
server_time: NOW.toISOString(),
|
||||
},
|
||||
}),
|
||||
)
|
||||
const result = await syncConnectorEntitlements(supabase, {
|
||||
config: CONFIG,
|
||||
fetchImpl,
|
||||
now: NOW,
|
||||
instanceUrl: 'https://bokforing.example.se',
|
||||
appVersion: '1.0.0',
|
||||
})
|
||||
|
||||
expect(fetchImpl).toHaveBeenCalledTimes(1)
|
||||
const [url, init] = fetchImpl.mock.calls[0] as [string, RequestInit]
|
||||
expect(url).toBe('https://app.gnubok.se/api/connect/entitlements')
|
||||
expect(init.method).toBe('POST')
|
||||
expect((init.headers as Record<string, string>).Authorization).toBe('Bearer gnubok_ck_test')
|
||||
expect(JSON.parse(String(init.body))).toEqual({
|
||||
active_company_count: 2,
|
||||
instance_url: 'https://bokforing.example.se',
|
||||
app_version: '1.0.0',
|
||||
})
|
||||
|
||||
expect(result).toMatchObject({ outcome: 'synced', companies: 2, grantsUpserted: 4, scopes: ['bank_sync', 'skatteverket'] })
|
||||
expect(result.expiresAt).toBe('2026-08-23T12:00:00.000Z')
|
||||
expect(upserts).toHaveLength(1)
|
||||
expect(upserts[0].opts).toEqual({ onConflict: 'company_id,team_id,capability_key,source' })
|
||||
expect(upserts[0].rows).toEqual(
|
||||
expect.arrayContaining([
|
||||
{ company_id: C1, team_id: null, capability_key: 'bank_sync', source: 'connector', expires_at: '2026-08-23T12:00:00.000Z' },
|
||||
{ company_id: C2, team_id: null, capability_key: 'skatteverket', source: 'connector', expires_at: '2026-08-23T12:00:00.000Z' },
|
||||
]),
|
||||
)
|
||||
// scopes no longer covered are dropped: delete source=connector NOT IN (kept)
|
||||
expect(deletes).toHaveLength(1)
|
||||
expect(deletes[0].filters).toEqual([
|
||||
['eq', 'source', 'connector'],
|
||||
['not', 'capability_key', 'in', '(bank_sync,skatteverket)'],
|
||||
])
|
||||
})
|
||||
|
||||
it('removes every connector grant on a 401/403 carrying a connector rejection code (freeze-and-retain)', async () => {
|
||||
for (const [status, code] of [[401, 'CONNECTOR_KEY_INVALID'], [403, 'CONNECTOR_KEY_SUSPENDED'], [401, 'CONNECTOR_KEY_MISSING']] as const) {
|
||||
const { supabase, upserts, deletes, setDeleteCount } = makeSupabase([C1])
|
||||
setDeleteCount(4)
|
||||
const result = await syncConnectorEntitlements(supabase, {
|
||||
config: CONFIG,
|
||||
fetchImpl: vi.fn().mockResolvedValue(jsonResponse(status, { error: 'rejected', code })),
|
||||
now: NOW,
|
||||
})
|
||||
expect(result).toMatchObject({ outcome: 'revoked', httpStatus: status, grantsDeleted: 4 })
|
||||
expect(upserts).toHaveLength(0)
|
||||
expect(deletes[0].filters).toEqual([['eq', 'source', 'connector']])
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps every grant on a 401/403 WITHOUT a connector rejection code (WAF challenge, edge protection, egress proxy)', async () => {
|
||||
// A bare-status 401/403 can come from layers where the hosted app never
|
||||
// ran. Only the hosted route's own rejection code proves revocation;
|
||||
// anything else must not destroy the 72h offline cache.
|
||||
const bodies: Array<[number, () => Response]> = [
|
||||
[403, () => new Response('<html>Attack challenge</html>', { status: 403, headers: { 'content-type': 'text/html' } })],
|
||||
[401, () => new Response('Authentication Required', { status: 401 })],
|
||||
[401, () => jsonResponse(401, { error: 'edge auth', code: 'SOME_OTHER_CODE' })],
|
||||
[403, () => jsonResponse(403, { message: 'forbidden by proxy' })],
|
||||
]
|
||||
for (const [status, make] of bodies) {
|
||||
const { supabase, upserts, deletes } = makeSupabase([C1])
|
||||
const result = await syncConnectorEntitlements(supabase, {
|
||||
config: CONFIG,
|
||||
fetchImpl: vi.fn().mockResolvedValue(make()),
|
||||
now: NOW,
|
||||
})
|
||||
expect(result).toMatchObject({ outcome: 'server_error', httpStatus: status, grantsDeleted: 0 })
|
||||
expect(upserts).toHaveLength(0)
|
||||
expect(deletes).toHaveLength(0)
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps grants on an UNKNOWN status or malformed current_period_end (contract drift lands in server_error, never delete)', async () => {
|
||||
const shapes = [
|
||||
{ status: 'past_due', scopes: ['bank_sync'], current_period_end: null, org_number: 'x', instance_url: null, server_time: 'x' },
|
||||
{ status: 'active', scopes: ['bank_sync'], current_period_end: 'not-a-date', org_number: 'x', instance_url: null, server_time: 'x' },
|
||||
{ status: 'active', scopes: [42], current_period_end: null, org_number: 'x', instance_url: null, server_time: 'x' },
|
||||
]
|
||||
for (const data of shapes) {
|
||||
const { supabase, upserts, deletes } = makeSupabase([C1])
|
||||
const result = await syncConnectorEntitlements(supabase, {
|
||||
config: CONFIG,
|
||||
fetchImpl: vi.fn().mockResolvedValue(jsonResponse(200, { data })),
|
||||
now: NOW,
|
||||
})
|
||||
expect(result.outcome, JSON.stringify(data)).toBe('server_error')
|
||||
expect(upserts).toHaveLength(0)
|
||||
expect(deletes).toHaveLength(0)
|
||||
}
|
||||
})
|
||||
|
||||
it('removes every connector grant when the key is not active', async () => {
|
||||
const { supabase, deletes } = makeSupabase([C1])
|
||||
const result = await syncConnectorEntitlements(supabase, {
|
||||
config: CONFIG,
|
||||
fetchImpl: vi.fn().mockResolvedValue(
|
||||
jsonResponse(200, { data: { status: 'suspended', scopes: ['bank_sync'], current_period_end: null, org_number: 'x', instance_url: null, server_time: 'x' } }),
|
||||
),
|
||||
now: NOW,
|
||||
})
|
||||
expect(result).toMatchObject({ outcome: 'revoked', status: 'suspended' })
|
||||
expect(deletes).toHaveLength(1)
|
||||
})
|
||||
|
||||
// The offline grace: a hosted outage must not touch the grants.
|
||||
it('leaves grants alone on a network error, a 5xx and a 429', async () => {
|
||||
const { supabase: s1, upserts: u1, deletes: d1 } = makeSupabase([C1])
|
||||
const r1 = await syncConnectorEntitlements(s1, { config: CONFIG, fetchImpl: vi.fn().mockRejectedValue(new Error('ECONNREFUSED')), now: NOW })
|
||||
expect(r1.outcome).toBe('network_error')
|
||||
expect(u1).toHaveLength(0)
|
||||
expect(d1).toHaveLength(0)
|
||||
|
||||
for (const status of [500, 503, 429]) {
|
||||
const { supabase, upserts, deletes } = makeSupabase([C1])
|
||||
const r = await syncConnectorEntitlements(supabase, { config: CONFIG, fetchImpl: vi.fn().mockResolvedValue(jsonResponse(status, {})), now: NOW })
|
||||
expect(r).toMatchObject({ outcome: 'server_error', httpStatus: status })
|
||||
expect(upserts).toHaveLength(0)
|
||||
expect(deletes).toHaveLength(0)
|
||||
}
|
||||
})
|
||||
|
||||
it('treats an unreadable 200 payload as a server error, not a revocation', async () => {
|
||||
const { supabase, deletes } = makeSupabase([C1])
|
||||
const r = await syncConnectorEntitlements(supabase, { config: CONFIG, fetchImpl: vi.fn().mockResolvedValue(jsonResponse(200, { data: { nope: true } })), now: NOW })
|
||||
expect(r.outcome).toBe('server_error')
|
||||
expect(deletes).toHaveLength(0)
|
||||
})
|
||||
|
||||
it('with an empty scope list writes nothing and drops every connector grant', async () => {
|
||||
const { supabase, upserts, deletes } = makeSupabase([C1])
|
||||
const r = await syncConnectorEntitlements(supabase, {
|
||||
config: CONFIG,
|
||||
fetchImpl: vi.fn().mockResolvedValue(jsonResponse(200, { data: { status: 'active', scopes: [], current_period_end: null, org_number: 'x', instance_url: null, server_time: 'x' } })),
|
||||
now: NOW,
|
||||
})
|
||||
expect(r).toMatchObject({ outcome: 'synced', grantsUpserted: 0 })
|
||||
expect(upserts).toHaveLength(0)
|
||||
expect(deletes[0].filters).toEqual([['eq', 'source', 'connector']])
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,52 @@
|
||||
import { DEFAULT_CONNECT_BASE_URL } from '../contract'
|
||||
import { createLogger } from '@/lib/logger'
|
||||
|
||||
const log = createLogger('connect/config')
|
||||
|
||||
/**
|
||||
* Instance-side connector configuration (a self-hosted deployment).
|
||||
*
|
||||
* GNUBOK_CONNECTOR_KEY the `gnubok_ck_...` key issued for this instance
|
||||
* GNUBOK_CONNECT_URL hosted origin, default https://app.gnubok.se
|
||||
*
|
||||
* Unset on hosted and on a self-host without a subscription: then the
|
||||
* connector sync is a no-op and the connector capabilities stay gated.
|
||||
*
|
||||
* GNUBOK_CONNECT_URL must be https: the hourly sync sends the long-lived
|
||||
* connector key as a Bearer header to this origin, so an http:// typo would
|
||||
* ship the credential in plaintext. Plain http is allowed only for loopback
|
||||
* hosts (local development against a dev server). An invalid or non-https
|
||||
* URL disables the connector entirely (fail closed, nothing is sent).
|
||||
*/
|
||||
export interface ConnectorConfig {
|
||||
key: string
|
||||
baseUrl: string
|
||||
}
|
||||
|
||||
const LOOPBACK_HOSTS = new Set(['localhost', '127.0.0.1', '[::1]'])
|
||||
|
||||
export function getConnectorConfig(): ConnectorConfig | null {
|
||||
const key = process.env.GNUBOK_CONNECTOR_KEY?.trim()
|
||||
if (!key) return null
|
||||
const raw = (process.env.GNUBOK_CONNECT_URL?.trim() || DEFAULT_CONNECT_BASE_URL).replace(/\/+$/, '')
|
||||
let url: URL
|
||||
try {
|
||||
url = new URL(raw)
|
||||
} catch {
|
||||
log.warn('GNUBOK_CONNECT_URL is not a valid URL; connector disabled', { value: raw })
|
||||
return null
|
||||
}
|
||||
const loopback = LOOPBACK_HOSTS.has(url.hostname) || LOOPBACK_HOSTS.has(`[${url.hostname}]`)
|
||||
if (url.protocol !== 'https:' && !(url.protocol === 'http:' && loopback)) {
|
||||
log.warn('GNUBOK_CONNECT_URL must be https (http only for loopback); connector disabled', {
|
||||
protocol: url.protocol,
|
||||
host: url.hostname,
|
||||
})
|
||||
return null
|
||||
}
|
||||
return { key, baseUrl: raw }
|
||||
}
|
||||
|
||||
export function isConnectorConfigured(): boolean {
|
||||
return getConnectorConfig() !== null
|
||||
}
|
||||
@@ -0,0 +1,236 @@
|
||||
import type { SupabaseClient } from '@supabase/supabase-js'
|
||||
import { CONNECTOR_CAPABILITIES } from '@/lib/entitlements/keys'
|
||||
import { createLogger } from '@/lib/logger'
|
||||
import { fetchAllRows } from '@/lib/supabase/fetch-all'
|
||||
import {
|
||||
CONNECTOR_ENTITLEMENTS_PATH,
|
||||
type ConnectorEntitlements,
|
||||
type ConnectorSyncReport,
|
||||
} from '../contract'
|
||||
import { getConnectorConfig, type ConnectorConfig } from './config'
|
||||
|
||||
const log = createLogger('connector-sync')
|
||||
|
||||
/**
|
||||
* Grant rows ARE the offline cache. Each sync re-stamps the connector grants
|
||||
* to expire at min(now + 72h, period_end + 3d): the hosted service only has
|
||||
* to be reachable once in three days, and a lapsed subscription freezes the
|
||||
* connector capabilities within days even if the sync keeps running. Uses
|
||||
* the existing expiry check in lib/entitlements, zero new cache code.
|
||||
*/
|
||||
export const CONNECTOR_GRANT_TTL_MS = 72 * 60 * 60 * 1000
|
||||
export const CONNECTOR_PERIOD_GRACE_MS = 3 * 24 * 60 * 60 * 1000
|
||||
const UPSERT_CHUNK = 500
|
||||
|
||||
export type ConnectorSyncOutcome =
|
||||
| 'not_configured'
|
||||
| 'synced'
|
||||
| 'revoked'
|
||||
| 'network_error'
|
||||
| 'server_error'
|
||||
|
||||
export interface ConnectorSyncResult {
|
||||
outcome: ConnectorSyncOutcome
|
||||
companies: number
|
||||
grantsUpserted: number
|
||||
grantsDeleted: number
|
||||
status?: string
|
||||
httpStatus?: number
|
||||
scopes?: string[]
|
||||
expiresAt?: string
|
||||
message?: string
|
||||
}
|
||||
|
||||
export interface ConnectorSyncOptions {
|
||||
fetchImpl?: typeof fetch
|
||||
now?: Date
|
||||
config?: ConnectorConfig | null
|
||||
instanceUrl?: string | null
|
||||
appVersion?: string | null
|
||||
}
|
||||
|
||||
export function connectorGrantExpiry(now: Date, currentPeriodEnd: string | null): string {
|
||||
const ttl = now.getTime() + CONNECTOR_GRANT_TTL_MS
|
||||
if (!currentPeriodEnd) return new Date(ttl).toISOString()
|
||||
const periodGrace = new Date(currentPeriodEnd).getTime() + CONNECTOR_PERIOD_GRACE_MS
|
||||
return new Date(Math.min(ttl, periodGrace)).toISOString()
|
||||
}
|
||||
|
||||
async function deleteConnectorGrants(supabase: SupabaseClient, keepKeys: string[] | null): Promise<number> {
|
||||
let query = supabase.from('capability_grants').delete({ count: 'exact' }).eq('source', 'connector')
|
||||
if (keepKeys && keepKeys.length > 0) {
|
||||
query = query.not('capability_key', 'in', `(${keepKeys.join(',')})`)
|
||||
}
|
||||
const { error, count } = await query
|
||||
if (error) throw new Error(`Failed to delete connector grants: ${error.message}`)
|
||||
return count ?? 0
|
||||
}
|
||||
|
||||
/**
|
||||
* One sync run: report the active company count to the hosted service, then
|
||||
* translate the answer into source='connector' capability grants for every
|
||||
* company on this instance.
|
||||
*
|
||||
* 200 active -> upsert grants for scopes, drop grants for scopes no
|
||||
* longer covered (synced)
|
||||
* 200 non-active -> delete all connector grants (freeze-and-retain) (revoked)
|
||||
* 401 / 403 -> delete all connector grants (revoked)
|
||||
* network error,
|
||||
* 429, 5xx, other -> leave grants alone, they expire on their own (network_error / server_error)
|
||||
*/
|
||||
export async function syncConnectorEntitlements(
|
||||
supabase: SupabaseClient,
|
||||
options: ConnectorSyncOptions = {},
|
||||
): Promise<ConnectorSyncResult> {
|
||||
const config = options.config === undefined ? getConnectorConfig() : options.config
|
||||
if (!config) return { outcome: 'not_configured', companies: 0, grantsUpserted: 0, grantsDeleted: 0 }
|
||||
|
||||
const now = options.now ?? new Date()
|
||||
const fetchImpl = options.fetchImpl ?? fetch
|
||||
|
||||
const companies = await fetchAllRows<{ id: string }>(({ from, to }) =>
|
||||
supabase
|
||||
.from('companies')
|
||||
.select('id')
|
||||
.is('archived_at', null)
|
||||
.order('id', { ascending: true })
|
||||
.range(from, to),
|
||||
)
|
||||
const companyIds = companies.map((c) => c.id)
|
||||
|
||||
const report: ConnectorSyncReport = {
|
||||
active_company_count: companyIds.length,
|
||||
...(options.instanceUrl ? { instance_url: options.instanceUrl } : {}),
|
||||
...(options.appVersion ? { app_version: options.appVersion } : {}),
|
||||
}
|
||||
|
||||
let response: Response
|
||||
try {
|
||||
response = await fetchImpl(`${config.baseUrl}${CONNECTOR_ENTITLEMENTS_PATH}`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${config.key}`,
|
||||
'Content-Type': 'application/json',
|
||||
Accept: 'application/json',
|
||||
},
|
||||
body: JSON.stringify(report),
|
||||
})
|
||||
} catch (err) {
|
||||
const message = err instanceof Error ? err.message : String(err)
|
||||
log.warn('connector sync: hosted service unreachable, keeping existing grants', { message })
|
||||
return { outcome: 'network_error', companies: companyIds.length, grantsUpserted: 0, grantsDeleted: 0, message }
|
||||
}
|
||||
|
||||
if (response.status === 401 || response.status === 403) {
|
||||
// Delete the grant cache ONLY on a genuine connector rejection, proven by
|
||||
// the hosted route's own JSON body code (with-connector-auth always sends
|
||||
// one). A bare-status 401/403 also comes from layers where the app never
|
||||
// ran: a Vercel WAF challenge page, edge deployment protection, an egress
|
||||
// proxy at the self-host. Trusting the status alone let any of those wipe
|
||||
// a paying instance's 72h offline grace within the hour: the same failure
|
||||
// class as the RPC-error-to-401 mapping, one layer up. Anything without a
|
||||
// known rejection code keeps the grants and expires naturally.
|
||||
let rejectionCode: string | null = null
|
||||
try {
|
||||
const body = (await response.clone().json()) as { code?: string }
|
||||
if (typeof body.code === 'string') rejectionCode = body.code
|
||||
} catch {
|
||||
// Non-JSON body (challenge/error page): not a connector rejection.
|
||||
}
|
||||
const isConnectorRejection =
|
||||
rejectionCode === 'CONNECTOR_KEY_MISSING' ||
|
||||
rejectionCode === 'CONNECTOR_KEY_INVALID' ||
|
||||
rejectionCode === 'CONNECTOR_KEY_SUSPENDED'
|
||||
if (isConnectorRejection) {
|
||||
const grantsDeleted = await deleteConnectorGrants(supabase, null)
|
||||
log.warn('connector sync: key rejected by hosted service, connector grants removed', {
|
||||
httpStatus: response.status,
|
||||
code: rejectionCode,
|
||||
grantsDeleted,
|
||||
})
|
||||
return { outcome: 'revoked', companies: companyIds.length, grantsUpserted: 0, grantsDeleted, httpStatus: response.status }
|
||||
}
|
||||
log.warn('connector sync: 401/403 without a connector rejection code (edge/proxy?), keeping existing grants', {
|
||||
httpStatus: response.status,
|
||||
code: rejectionCode,
|
||||
})
|
||||
return { outcome: 'server_error', companies: companyIds.length, grantsUpserted: 0, grantsDeleted: 0, httpStatus: response.status }
|
||||
}
|
||||
if (!response.ok) {
|
||||
log.warn('connector sync: hosted service error, keeping existing grants', { httpStatus: response.status })
|
||||
return { outcome: 'server_error', companies: companyIds.length, grantsUpserted: 0, grantsDeleted: 0, httpStatus: response.status }
|
||||
}
|
||||
|
||||
let entitlements: ConnectorEntitlements
|
||||
try {
|
||||
const body = (await response.json()) as { data?: ConnectorEntitlements }
|
||||
// Full shape validation, not just presence: an UNKNOWN status string must
|
||||
// land in keep-grants (server_error), not fall through to the
|
||||
// "status !== 'active'" delete below (contract drift or a tampering
|
||||
// middlebox must never wipe the cache), and a malformed
|
||||
// current_period_end would make connectorGrantExpiry throw mid-write.
|
||||
if (
|
||||
!body.data ||
|
||||
!['active', 'suspended', 'revoked'].includes(body.data.status as string) ||
|
||||
!Array.isArray(body.data.scopes) ||
|
||||
!body.data.scopes.every((s) => typeof s === 'string') ||
|
||||
!(
|
||||
body.data.current_period_end === null ||
|
||||
body.data.current_period_end === undefined ||
|
||||
(typeof body.data.current_period_end === 'string' &&
|
||||
Number.isFinite(new Date(body.data.current_period_end).getTime()))
|
||||
)
|
||||
) {
|
||||
throw new Error('unexpected entitlements payload')
|
||||
}
|
||||
entitlements = body.data
|
||||
} catch (err) {
|
||||
const message = err instanceof Error ? err.message : String(err)
|
||||
log.warn('connector sync: unreadable entitlements payload, keeping existing grants', { message })
|
||||
return { outcome: 'server_error', companies: companyIds.length, grantsUpserted: 0, grantsDeleted: 0, httpStatus: response.status, message }
|
||||
}
|
||||
|
||||
if (entitlements.status !== 'active') {
|
||||
const grantsDeleted = await deleteConnectorGrants(supabase, null)
|
||||
log.warn('connector sync: key not active, connector grants removed', { status: entitlements.status, grantsDeleted })
|
||||
return { outcome: 'revoked', companies: companyIds.length, grantsUpserted: 0, grantsDeleted, status: entitlements.status }
|
||||
}
|
||||
|
||||
const scopes = entitlements.scopes.filter((s) => (CONNECTOR_CAPABILITIES as readonly string[]).includes(s))
|
||||
const expiresAt = connectorGrantExpiry(now, entitlements.current_period_end)
|
||||
|
||||
let grantsUpserted = 0
|
||||
if (scopes.length > 0 && companyIds.length > 0) {
|
||||
const rows = companyIds.flatMap((companyId) =>
|
||||
scopes.map((capabilityKey) => ({
|
||||
company_id: companyId,
|
||||
team_id: null,
|
||||
capability_key: capabilityKey,
|
||||
source: 'connector',
|
||||
expires_at: expiresAt,
|
||||
})),
|
||||
)
|
||||
for (let i = 0; i < rows.length; i += UPSERT_CHUNK) {
|
||||
const chunk = rows.slice(i, i + UPSERT_CHUNK)
|
||||
const { error } = await supabase
|
||||
.from('capability_grants')
|
||||
.upsert(chunk, { onConflict: 'company_id,team_id,capability_key,source' })
|
||||
if (error) throw new Error(`Failed to upsert connector grants: ${error.message}`)
|
||||
grantsUpserted += chunk.length
|
||||
}
|
||||
}
|
||||
// Scopes the subscription no longer covers (or none at all): drop them.
|
||||
const grantsDeleted = await deleteConnectorGrants(supabase, scopes.length > 0 ? scopes : null)
|
||||
|
||||
log.info('connector sync complete', { companies: companyIds.length, scopes, grantsUpserted, grantsDeleted, expiresAt })
|
||||
return {
|
||||
outcome: 'synced',
|
||||
companies: companyIds.length,
|
||||
grantsUpserted,
|
||||
grantsDeleted,
|
||||
status: entitlements.status,
|
||||
scopes,
|
||||
expiresAt,
|
||||
httpStatus: response.status,
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user