/** * Salesforce integration client implementation. * * Executes Salesforce operations through the Superblocks orchestrator * including SOQL queries, single-object CRUD, and bulk operations. */ import type { z } from "zod"; import { QueryValidationError, RestApiValidationError } from "../../errors.js"; import { IntegrationError } from "../../runtime/errors.js"; import type { QueryExecutor, TraceMetadata } from "../registry.js"; import type { IntegrationConfig, IntegrationClientImpl } from "../types.js"; import { describeType } from "../utils.js"; import type { SalesforceClient } from "./types.js"; /** * SOQL action string value expected by the backend. */ const SOQL_ACTION_SOQL = "SOQL_ACTION_SOQL"; /** * CRUD action string values matching proto CrudAction enum. */ const CRUD_ACTION = { CREATE: "CRUD_ACTION_CREATE", UPDATE: "CRUD_ACTION_UPDATE", DELETE: "CRUD_ACTION_DELETE", READ: "CRUD_ACTION_READ", } as const; /** * Bulk action string values matching proto BulkAction enum. */ const BULK_ACTION = { CREATE: "BULK_ACTION_CREATE", UPDATE: "BULK_ACTION_UPDATE", DELETE: "BULK_ACTION_DELETE", UPSERT: "BULK_ACTION_UPSERT", } as const; /** * Internal implementation of the Salesforce client. */ export class SalesforceClientImpl implements SalesforceClient, IntegrationClientImpl { readonly config: IntegrationConfig; private readonly executeQuery: QueryExecutor; constructor(config: IntegrationConfig, executeQuery: QueryExecutor) { this.config = config; this.executeQuery = executeQuery; } get name(): string { return this.config.name; } get pluginId(): string { return this.config.pluginId; } /** * Execute a request and wrap errors. */ private async exec( request: Record, operation: string, metadata?: TraceMetadata, ): Promise { try { return await this.executeQuery(request, undefined, metadata); } catch (error) { if ( error instanceof QueryValidationError || error instanceof RestApiValidationError || error instanceof IntegrationError ) { throw error; } throw new IntegrationError(this.config.name, operation, error); } } /** * Build a CRUD request object. */ private buildCrudRequest( action: string, resourceType: string, options?: { resourceId?: string; resourceBody?: string }, ): Record { const crud: Record = { action, resourceType, }; if (options?.resourceId) { crud.resourceId = options.resourceId; } if (options?.resourceBody) { crud.resourceBody = options.resourceBody; } return { crud }; } /** * Build a bulk request object. */ private buildBulkRequest( action: string, resourceType: string, records: Array>, externalId?: string, ): Record { const bulk: Record = { action, resourceType, resourceBody: JSON.stringify(records), }; if (externalId) { bulk.externalId = externalId; } return { bulk }; } async query( soql: string, schema: z.ZodSchema, metadata?: TraceMetadata, ): Promise { const request = { soql: { action: SOQL_ACTION_SOQL, sqlBody: soql, }, }; const result = await this.exec( request as unknown as Record, "query", metadata, ); if (!Array.isArray(result)) { throw new IntegrationError( this.config.name, "query", `Expected array result from Salesforce query, got: ${describeType(result)}`, ); } const validated: T[] = []; for (let i = 0; i < result.length; i++) { const row = result[i]; const parseResult = schema.safeParse(row); if (!parseResult.success) { throw new QueryValidationError( `Row ${i} validation failed: ${parseResult.error.message}`, { rowIndex: i, errors: parseResult.error.errors.map((err) => ({ path: err.path, message: err.message, })), row, }, ); } validated.push(parseResult.data); } return validated; } async create( resourceType: string, body: Record, metadata?: TraceMetadata, ): Promise { const request = this.buildCrudRequest(CRUD_ACTION.CREATE, resourceType, { resourceBody: JSON.stringify(body), }); return this.exec(request, "create", metadata); } async read( resourceType: string, resourceId: string, schema: z.ZodSchema, metadata?: TraceMetadata, ): Promise { const request = this.buildCrudRequest(CRUD_ACTION.READ, resourceType, { resourceId, }); const result = await this.exec(request, "read", metadata); const parseResult = schema.safeParse(result); if (!parseResult.success) { throw new RestApiValidationError( `Result validation failed: ${parseResult.error.message}`, { zodError: parseResult.error, data: result, }, ); } return parseResult.data; } async update( resourceType: string, resourceId: string, body: Record, metadata?: TraceMetadata, ): Promise { const request = this.buildCrudRequest(CRUD_ACTION.UPDATE, resourceType, { resourceId, resourceBody: JSON.stringify(body), }); return this.exec(request, "update", metadata); } async remove( resourceType: string, resourceId: string, metadata?: TraceMetadata, ): Promise { const request = this.buildCrudRequest(CRUD_ACTION.DELETE, resourceType, { resourceId, }); return this.exec(request, "remove", metadata); } async bulkCreate( resourceType: string, records: Array>, metadata?: TraceMetadata, ): Promise { const request = this.buildBulkRequest( BULK_ACTION.CREATE, resourceType, records, ); return this.exec(request, "bulkCreate", metadata); } async bulkUpdate( resourceType: string, records: Array>, metadata?: TraceMetadata, ): Promise { const request = this.buildBulkRequest( BULK_ACTION.UPDATE, resourceType, records, ); return this.exec(request, "bulkUpdate", metadata); } async bulkDelete( resourceType: string, records: Array>, metadata?: TraceMetadata, ): Promise { const request = this.buildBulkRequest( BULK_ACTION.DELETE, resourceType, records, ); return this.exec(request, "bulkDelete", metadata); } async bulkUpsert( resourceType: string, records: Array>, externalId: string, metadata?: TraceMetadata, ): Promise { const request = this.buildBulkRequest( BULK_ACTION.UPSERT, resourceType, records, externalId, ); return this.exec(request, "bulkUpsert", metadata); } }