Skip to content
← Public packages

@kody/salesforce

Salesforce REST helpers for identity, SOQL, multi-org routing, and error investigation.

src/request.ts

398 lines · 11.2 KB · TypeScript
import { createAuthenticatedFetch, kody } from 'kody:runtime'
import type {
	JsonRecord,
	SalesforceHttpMethod,
	SalesforceRequestInput,
	SalesforceRequestResult,
	SalesforceUserInfo,
} from './types.ts'
import {
	optionalBoolean,
	optionalString,
	requireApiVersion,
	requireRecord,
	requireString,
} from './types.ts'

export const SALESFORCE_INTEGRATION = 'salesforce'
export const DEFAULT_API_VERSION = 'v61.0'
export const DEFAULT_LOGIN_ORIGIN = 'https://login.salesforce.com'

const READ_ONLY_PATH_PREFIXES = [
	'/services/oauth2/userinfo',
	'/services/data/',
]

const authenticatedFetches = new Map<string, typeof fetch>()
const cachedLoginOrigins = new Map<string, string>()
const cachedInstanceUrls = new Map<string, string>()

/**
 * Call a Salesforce REST endpoint with a saved OAuth integration (`integrationName`, default `salesforce`). Unknown and mutating endpoints require `confirm: true`; use `dryRun: true` to preview.
 * @param input.path - REST path relative to the data API.
 * @param input.method - HTTP method (default GET).
 * @param input.confirm - Required for unknown/mutating calls.
 * @returns Parsed JSON body or dry-run preview.
 * @example
 * import request from 'kody:@kody/salesforce/request'
 * const result = await request({
 *   path: 'query',
 *   query: { q: 'SELECT Id, Name FROM Account LIMIT 5' },
 * })
 */
export default async function requestEntrypoint(
	params: Partial<SalesforceRequestInput> & Record<string, unknown> = {},
) {
	const input = requireRecord(params, 'request')
	requireString(input.path, 'path')
	const method = optionalString(input.method, 'method')
	if (
		method !== undefined &&
		method !== 'GET' &&
		method !== 'POST' &&
		method !== 'PATCH' &&
		method !== 'PUT' &&
		method !== 'DELETE'
	) {
		throw new Error('method must be GET, POST, PATCH, PUT, or DELETE.')
	}
	return request({
		...(input as SalesforceRequestInput),
		method: (method as SalesforceHttpMethod | undefined) ?? 'GET',
		integrationName: optionalString(input.integrationName, 'integrationName'),
		confirm: optionalBoolean(input.confirm, 'confirm'),
		dryRun: optionalBoolean(input.dryRun, 'dryRun'),
	})
}

export function resolveIntegrationName(value?: string): string {
	const name = optionalString(value, 'integrationName')
	return name && name.length > 0 ? name : SALESFORCE_INTEGRATION
}

export class SalesforceApiError extends Error {
	readonly status: number
	readonly statusText: string
	readonly details: unknown

	constructor(
		message: string,
		options: {
			status: number
			statusText: string
			details: unknown
		},
	) {
		super(message)
		this.name = 'SalesforceApiError'
		this.status = options.status
		this.statusText = options.statusText
		this.details = options.details
	}
}

async function getAuthenticatedFetch(
	integrationName: string,
): Promise<typeof fetch> {
	let authedFetch = authenticatedFetches.get(integrationName)
	if (!authedFetch) {
		authedFetch = await createAuthenticatedFetch(integrationName)
		authenticatedFetches.set(integrationName, authedFetch)
	}
	return authedFetch
}

export async function getLoginOrigin(
	integrationName: string = SALESFORCE_INTEGRATION,
): Promise<string> {
	const cached = cachedLoginOrigins.get(integrationName)
	if (cached) return cached

	const integration = await kody.integrationGet({ name: integrationName })
	const tokenUrl = integration?.tokenUrl || integration?.integration?.tokenUrl
	if (typeof tokenUrl === 'string' && tokenUrl.length > 0) {
		const origin = new URL(tokenUrl).origin
		cachedLoginOrigins.set(integrationName, origin)
		return origin
	}

	cachedLoginOrigins.set(integrationName, DEFAULT_LOGIN_ORIGIN)
	return DEFAULT_LOGIN_ORIGIN
}

function normalizeInstanceUrl(value: string): string {
	const url = new URL(requireString(value, 'instanceUrl'))
	if (url.protocol !== 'https:') {
		throw new Error('instanceUrl must be an https URL.')
	}
	return url.origin
}

export async function getInstanceUrl(
	override?: string,
	integrationName: string = SALESFORCE_INTEGRATION,
): Promise<string> {
	if (override) {
		return normalizeInstanceUrl(override)
	}
	const cached = cachedInstanceUrls.get(integrationName)
	if (cached) return cached

	const userinfo = await fetchUserInfo(integrationName)
	const restUrl = userinfo.urls?.rest
	if (typeof restUrl === 'string' && restUrl.length > 0) {
		const origin = normalizeInstanceUrl(restUrl)
		cachedInstanceUrls.set(integrationName, origin)
		return origin
	}

	const customDomain = userinfo.urls?.custom_domain
	if (typeof customDomain === 'string' && customDomain.length > 0) {
		const origin = normalizeInstanceUrl(customDomain)
		cachedInstanceUrls.set(integrationName, origin)
		return origin
	}

	throw new Error(
		`Salesforce userinfo did not include an instance URL. Reconnect the "${integrationName}" OAuth integration.`,
	)
}

export async function fetchUserInfo(
	integrationName: string = SALESFORCE_INTEGRATION,
): Promise<SalesforceUserInfo> {
	const loginOrigin = await getLoginOrigin(integrationName)
	const result = await salesforceRequestWithResponse({
		path: `${loginOrigin}/services/oauth2/userinfo`,
		method: 'GET',
		integrationName,
	})
	return requireRecord(result.body, 'userinfo') as SalesforceUserInfo
}

export function isKnownReadOnlyRequest(method: SalesforceHttpMethod, path: string): boolean {
	if (method === 'GET') return true
	if (method !== 'POST') return false

	const normalized = path.split('?')[0] ?? path
	return (
		normalized.includes('/query') ||
		normalized.includes('/queryAll') ||
		normalized.includes('/search') ||
		normalized.endsWith('/describe') ||
		normalized.includes('/describe/')
	)
}

function buildUrl(
	path: string,
	options: {
		instanceUrl: string
		apiVersion: string
		query?: SalesforceRequestInput['query']
	},
): string {
	const trimmed = requireString(path, 'path')
	let url: URL

	if (/^https?:\/\//i.test(trimmed)) {
		url = new URL(trimmed)
	} else if (trimmed.startsWith('/services/')) {
		url = new URL(trimmed, options.instanceUrl)
	} else if (trimmed.startsWith('/')) {
		url = new URL(trimmed, options.instanceUrl)
	} else {
		const relative = trimmed.replace(/^\/+/, '')
		url = new URL(`/services/data/${options.apiVersion}/${relative}`, options.instanceUrl)
	}

	if (options.query) {
		for (const [key, value] of Object.entries(options.query)) {
			if (value === undefined || value === null) continue
			url.searchParams.set(key, String(value))
		}
	}

	return url.toString()
}

function pathForSafetyCheck(absoluteUrl: string, instanceUrl: string, loginOrigin: string): string {
	const url = new URL(absoluteUrl)
	if (url.origin === instanceUrl || url.origin === loginOrigin) {
		return `${url.pathname}${url.search}`
	}
	return absoluteUrl
}

export async function rawSalesforceRequest(input: SalesforceRequestInput): Promise<Response> {
	const method = input.method ?? 'GET'
	const apiVersion = requireApiVersion(input.apiVersion ?? DEFAULT_API_VERSION)
	const integrationName = resolveIntegrationName(input.integrationName)
	const loginOrigin = await getLoginOrigin(integrationName)
	const usesLoginHost =
		/^https?:\/\//i.test(input.path) && new URL(input.path).origin === loginOrigin
	const instanceUrl = usesLoginHost
		? loginOrigin
		: await getInstanceUrl(
				optionalString(input.instanceUrl, 'instanceUrl'),
				integrationName,
			)
	const url = buildUrl(input.path, {
		instanceUrl,
		apiVersion,
		query: input.query,
	})

	const headers = new Headers(input.headers)
	if (!headers.has('accept')) headers.set('accept', 'application/json')

	let body: BodyInit | undefined
	if (input.body !== undefined && method !== 'GET' && method !== 'DELETE') {
		if (!headers.has('content-type')) {
			headers.set('content-type', 'application/json')
		}
		body = typeof input.body === 'string' ? input.body : JSON.stringify(input.body)
	}

	const authedFetch = await getAuthenticatedFetch(integrationName)
	return authedFetch(url, {
		method,
		headers,
		body,
	})
}

export async function salesforceRequestWithResponse(
	input: SalesforceRequestInput,
): Promise<SalesforceRequestResult> {
	const response = await rawSalesforceRequest(input)
	const headers = selectResponseHeaders(response.headers)
	const body = await parseBody(response)

	if (!response.ok) {
		throw new SalesforceApiError(
			`Salesforce API request failed: ${response.status} ${response.statusText}`,
			{
				status: response.status,
				statusText: response.statusText,
				details: body,
			},
		)
	}

	return {
		ok: true,
		status: response.status,
		statusText: response.statusText,
		headers,
		body,
	}
}

export async function salesforceJson<T>(input: SalesforceRequestInput): Promise<T> {
	const result = await salesforceRequestWithResponse(input)
	return result.body as T
}

export async function request(
	input: SalesforceRequestInput,
): Promise<SalesforceRequestResult | JsonRecord> {
	const method = input.method ?? 'GET'
	const apiVersion = requireApiVersion(input.apiVersion ?? DEFAULT_API_VERSION)
	const integrationName = resolveIntegrationName(input.integrationName)
	const loginOrigin = await getLoginOrigin(integrationName).catch(
		() => DEFAULT_LOGIN_ORIGIN,
	)
	const usesLoginHost =
		/^https?:\/\//i.test(input.path) && new URL(input.path).origin === loginOrigin

	if (input.dryRun) {
		let instanceUrl = optionalString(input.instanceUrl, 'instanceUrl')
		if (!instanceUrl) {
			if (usesLoginHost) {
				instanceUrl = loginOrigin
			} else {
				try {
					instanceUrl = await getInstanceUrl(undefined, integrationName)
				} catch {
					instanceUrl = 'https://your-org.my.salesforce.com'
				}
			}
		} else {
			instanceUrl = normalizeInstanceUrl(instanceUrl)
		}

		const url = buildUrl(input.path, {
			instanceUrl,
			apiVersion,
			query: input.query,
		})
		const safetyPath = pathForSafetyCheck(url, instanceUrl, loginOrigin)
		return {
			dryRun: true,
			method,
			path: input.path,
			url,
			apiVersion,
			instanceUrl,
			integrationName,
			body: input.body ?? null,
			query: input.query ?? null,
			authorization: `[managed ${integrationName} OAuth integration]`,
			readOnly: isKnownReadOnlyRequest(method, safetyPath),
		}
	}

	const instanceUrl = usesLoginHost
		? loginOrigin
		: await getInstanceUrl(
				optionalString(input.instanceUrl, 'instanceUrl'),
				integrationName,
			)
	const url = buildUrl(input.path, {
		instanceUrl,
		apiVersion,
		query: input.query,
	})
	const safetyPath = pathForSafetyCheck(url, instanceUrl, loginOrigin)

	const apparentlyReadOnly =
		isKnownReadOnlyRequest(method, safetyPath) ||
		READ_ONLY_PATH_PREFIXES.some((prefix) => safetyPath.startsWith(prefix) && method === 'GET')

	if (!apparentlyReadOnly && input.confirm !== true) {
		throw new Error(
			'This Salesforce request may mutate data. Pass confirm: true only after explicit user approval, or use dryRun: true.',
		)
	}

	return salesforceRequestWithResponse({ ...input, method, apiVersion })
}

async function parseBody(response: Response): Promise<unknown> {
	if (response.status === 204) return null
	const text = await response.text()
	if (!text) return null
	try {
		return JSON.parse(text)
	} catch {
		return text
	}
}

function selectResponseHeaders(headers: Headers): Record<string, string> {
	const selected = [
		'content-type',
		'content-length',
		'sforce-limit-info',
		'location',
		'etag',
	]
	const output: Record<string, string> = {}
	for (const key of selected) {
		const value = headers.get(key)
		if (value) output[key] = value
	}
	return output
}