http.ts 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157
  1. /**
  2. * One HTTP client per connection.
  3. *
  4. * Attaches the credentials, refreshes an OAuth token before it expires and
  5. * again on a 401, retries 429 and 5xx with backoff, and logs failures. Token
  6. * refresh runs under a per-connection lock so two jobs cannot both refresh
  7. * and invalidate each other's token, which is how QuickBooks-style rotating
  8. * refresh tokens get lost.
  9. */
  10. import { db } from '@/lib/db'
  11. import { type OAuth2Spec, needsRefresh, refreshToken, resolveClient } from './oauth'
  12. import {
  13. type ConnectorHttp,
  14. ConnectorHttpError,
  15. type LogLevel,
  16. type OAuthCredentials,
  17. } from './types'
  18. import { openCredentials, sealCredentials } from './vault'
  19. const MAX_RETRIES = 3
  20. const refreshLocks = new Map<string, Promise<OAuthCredentials>>()
  21. export interface HttpAuth {
  22. /** OAuth: refresh through the vendor. Static: nothing to refresh. */
  23. oauth?: OAuth2Spec
  24. /** Header builder for non-OAuth auth. Defaults to Bearer with credentials.accessToken. */
  25. headers?: (credentials: Record<string, unknown>) => Record<string, string>
  26. }
  27. async function sleep(ms: number): Promise<void> {
  28. await new Promise((resolve) => setTimeout(resolve, ms))
  29. }
  30. /** The vendor's request id, in the header Intuit and most others use for it. */
  31. function requestIdOf(res: Response): string | null {
  32. return res.headers.get('intuit_tid') ?? res.headers.get('x-request-id') ?? null
  33. }
  34. function retryDelayMs(attempt: number, res: Response | null): number {
  35. const retryAfter = res?.headers.get('retry-after')
  36. if (retryAfter) {
  37. const seconds = Number(retryAfter)
  38. if (Number.isFinite(seconds)) return Math.min(seconds, 60) * 1000
  39. }
  40. const base = [500, 2000, 5000][attempt] ?? 5000
  41. return base + Math.floor(Math.random() * 500)
  42. }
  43. /**
  44. * Refresh under a lock. A second caller waiting on the same connection gets
  45. * the same promise, so the vendor sees one refresh. The database is reread
  46. * first, in case another process refreshed already.
  47. */
  48. export async function refreshUnderLock(
  49. connectionId: string,
  50. spec: OAuth2Spec
  51. ): Promise<OAuthCredentials> {
  52. const inflight = refreshLocks.get(connectionId)
  53. if (inflight) return inflight
  54. const run = (async () => {
  55. const row = await db.integrationConnection.findUnique({
  56. where: { id: connectionId },
  57. select: { credentials: true },
  58. })
  59. const current = openCredentials(row?.credentials) as unknown as OAuthCredentials
  60. if (!needsRefresh(current)) return current
  61. const client = resolveClient(spec, current as unknown as Record<string, unknown>)
  62. if (!client) throw new Error('No OAuth client configured for this connection')
  63. const next = await refreshToken({ spec, client, credentials: current })
  64. await db.integrationConnection.update({
  65. where: { id: connectionId },
  66. data: { credentials: sealCredentials(next as unknown as Record<string, unknown>) },
  67. })
  68. return next
  69. })()
  70. refreshLocks.set(connectionId, run)
  71. try {
  72. return await run
  73. } finally {
  74. refreshLocks.delete(connectionId)
  75. }
  76. }
  77. export function createConnectorHttp(input: {
  78. connectionId: string
  79. credentials: Record<string, unknown>
  80. auth: HttpAuth
  81. log: (level: LogLevel, message: string, details?: Record<string, unknown>) => Promise<void>
  82. }): ConnectorHttp {
  83. let credentials = input.credentials
  84. const authHeaders = async (force: boolean): Promise<Record<string, string>> => {
  85. if (input.auth.oauth) {
  86. const oauth = credentials as unknown as OAuthCredentials
  87. if (force || needsRefresh(oauth)) {
  88. if (force) {
  89. // Make the lock refresh even when the stored expiry looks fine.
  90. credentials = { ...credentials, expiresAt: 0 }
  91. }
  92. const fresh = await refreshUnderLock(input.connectionId, input.auth.oauth)
  93. credentials = fresh as unknown as Record<string, unknown>
  94. }
  95. return { Authorization: `Bearer ${(credentials as unknown as OAuthCredentials).accessToken}` }
  96. }
  97. if (input.auth.headers) return input.auth.headers(credentials)
  98. const token = typeof credentials.accessToken === 'string' ? credentials.accessToken : ''
  99. return token ? { Authorization: `Bearer ${token}` } : {}
  100. }
  101. const doFetch = async (url: string, init?: RequestInit): Promise<Response> => {
  102. let forceRefresh = false
  103. let last: Response | null = null
  104. for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) {
  105. const headers = new Headers(init?.headers)
  106. for (const [k, v] of Object.entries(await authHeaders(forceRefresh))) headers.set(k, v)
  107. forceRefresh = false
  108. // A form carries its own type, with the boundary only fetch can write.
  109. if (init?.body && !(init.body instanceof FormData) && !headers.has('Content-Type'))
  110. headers.set('Content-Type', 'application/json')
  111. let res: Response
  112. try {
  113. res = await fetch(url, { ...init, headers })
  114. } catch (err) {
  115. if (attempt === MAX_RETRIES) throw err
  116. await sleep(retryDelayMs(attempt, null))
  117. continue
  118. }
  119. if (res.status === 401 && input.auth.oauth && attempt === 0) {
  120. forceRefresh = true
  121. continue
  122. }
  123. if ((res.status === 429 || res.status >= 500) && attempt < MAX_RETRIES) {
  124. last = res
  125. await input.log('warn', `Retrying ${res.status} from ${new URL(url).host}`, {
  126. status: res.status,
  127. attempt: attempt + 1,
  128. })
  129. await sleep(retryDelayMs(attempt, res))
  130. continue
  131. }
  132. return res
  133. }
  134. return last as Response
  135. }
  136. return {
  137. fetch: doFetch,
  138. async json<T>(url: string, init?: RequestInit): Promise<T> {
  139. const res = await doFetch(url, init)
  140. const text = await res.text()
  141. if (!res.ok) throw new ConnectorHttpError(res.status, text, url, requestIdOf(res))
  142. if (!text) return undefined as T
  143. return JSON.parse(text) as T
  144. },
  145. }
  146. }