/** * DynamoDB client implementation. * * Uses the native DynamoDB plugin type from @superblocksteam/types. * All operations pass parameters as a JSON string in the `body` field, * which the orchestrator parses and forwards directly to the AWS SDK. */ import type { z } from "zod"; import { 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 type { DynamoDBClient, DynamoDBAttributeValue, DynamoDBScanOptions, } from "./types.js"; const DYNAMODB_SCAN_OPTION_KEYS = new Set([ "filterExpression", "expressionAttributeValues", "expressionAttributeNames", "exclusiveStartKey", "limit", "projectionExpression", "indexName", "segment", "totalSegments", ]); /** * Internal implementation of DynamoDBClient. * * The orchestrator's DynamoDB handler only uses two fields from the proto: * - `action`: which AWS SDK command to execute * - `body`: JSON string of parameters passed directly to that command * * All other proto fields (table, filterBy, newValues, etc.) are ignored. */ export class DynamoDBClientImpl implements DynamoDBClient, 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; } /** * Build a request with action and JSON body. * * @param action - AWS SDK command name (e.g., 'executeStatement', 'getItem') * @param params - Parameters to serialize as JSON body */ private buildRequest( action: string, params: Record, ): Record { return { action, body: JSON.stringify(params), }; } /** * Validate a result against a schema. */ private validateResult( result: unknown, schema: z.ZodSchema, operation: string, ): T { 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; } /** * Execute a request and wrap errors. */ private async executeWithErrorHandling( request: Record, operation: string, metadata?: TraceMetadata, ): Promise { try { return await this.executeQuery(request, undefined, metadata); } catch (error) { if ( error instanceof RestApiValidationError || error instanceof IntegrationError ) { throw error; } throw new IntegrationError(this.config.name, operation, error); } } async query( statement: string, schema: z.ZodSchema, params?: DynamoDBAttributeValue[], metadata?: TraceMetadata, ): Promise { const body: Record = { Statement: statement }; if (params && params.length > 0) { body.Parameters = params; } const request = this.buildRequest("executeStatement", body); const result = await this.executeWithErrorHandling( request, "query", metadata, ); return this.validateResult(result, schema, "query"); } async getItem( table: string, key: Record, schema: z.ZodSchema, metadata?: TraceMetadata, ): Promise { const request = this.buildRequest("getItem", { TableName: table, Key: key, }); const result = await this.executeWithErrorHandling( request, "getItem", metadata, ); return this.validateResult(result, schema, "getItem"); } async putItem( table: string, item: Record, metadata?: TraceMetadata, ): Promise { const request = this.buildRequest("putItem", { TableName: table, Item: item, }); return this.executeWithErrorHandling(request, "putItem", metadata); } async updateItem( table: string, key: Record, updateExpression: string, expressionAttributeValues: Record, expressionAttributeNames?: Record, metadata?: TraceMetadata, ): Promise { const params: Record = { TableName: table, Key: key, UpdateExpression: updateExpression, ExpressionAttributeValues: expressionAttributeValues, }; if (expressionAttributeNames) { params.ExpressionAttributeNames = expressionAttributeNames; } const request = this.buildRequest("updateItem", params); return this.executeWithErrorHandling(request, "updateItem", metadata); } async deleteItem( table: string, key: Record, metadata?: TraceMetadata, ): Promise { const request = this.buildRequest("deleteItem", { TableName: table, Key: key, }); return this.executeWithErrorHandling(request, "deleteItem", metadata); } /** Copy optional AWS Scan fields onto the request body. */ private applyScanParams( params: Record, options: DynamoDBScanOptions, ): void { if (options.filterExpression) { params.FilterExpression = options.filterExpression; } if ( options.expressionAttributeValues && Object.keys(options.expressionAttributeValues).length > 0 ) { params.ExpressionAttributeValues = options.expressionAttributeValues; } if ( options.expressionAttributeNames && Object.keys(options.expressionAttributeNames).length > 0 ) { params.ExpressionAttributeNames = options.expressionAttributeNames; } if (options.exclusiveStartKey) { params.ExclusiveStartKey = options.exclusiveStartKey; } if (options.limit !== undefined) { params.Limit = options.limit; } if (options.projectionExpression) { params.ProjectionExpression = options.projectionExpression; } if (options.indexName) { params.IndexName = options.indexName; } if (options.segment !== undefined) { params.Segment = options.segment; } if (options.totalSegments !== undefined) { params.TotalSegments = options.totalSegments; } } private isScanOptions(value: unknown): value is DynamoDBScanOptions { return ( typeof value === "object" && value !== null && Object.keys(value).every((key) => DYNAMODB_SCAN_OPTION_KEYS.has(key)) ); } private isTraceMetadata(value: unknown): value is TraceMetadata { return ( typeof value === "object" && value !== null && Object.keys(value).every( (key) => key === "label" || key === "description", ) && Object.values(value).every( (entry) => entry === undefined || typeof entry === "string", ) ); } async scan( table: string, schema: z.ZodSchema, filterExpressionOrOptions?: string | DynamoDBScanOptions, expressionAttributeValuesOrMetadata?: | Record | TraceMetadata, expressionAttributeNames?: Record, metadata?: TraceMetadata, ): Promise { const params: Record = { TableName: table }; let resolvedMetadata = metadata; if ( typeof filterExpressionOrOptions === "object" && filterExpressionOrOptions !== null ) { if (!this.isScanOptions(filterExpressionOrOptions)) { throw new Error( `Invalid DynamoDB scan options: ${Object.keys(filterExpressionOrOptions).join(", ")}`, ); } this.applyScanParams(params, filterExpressionOrOptions); if (expressionAttributeValuesOrMetadata !== undefined) { if (!this.isTraceMetadata(expressionAttributeValuesOrMetadata)) { throw new Error("Invalid DynamoDB scan trace metadata"); } resolvedMetadata = expressionAttributeValuesOrMetadata; } } else { if (filterExpressionOrOptions) { params.FilterExpression = filterExpressionOrOptions; } if ( expressionAttributeValuesOrMetadata && Object.keys(expressionAttributeValuesOrMetadata).length > 0 ) { params.ExpressionAttributeValues = expressionAttributeValuesOrMetadata; } if ( expressionAttributeNames && Object.keys(expressionAttributeNames).length > 0 ) { params.ExpressionAttributeNames = expressionAttributeNames; } } const request = this.buildRequest("scan", params); const result = await this.executeWithErrorHandling( request, "scan", resolvedMetadata, ); return this.validateResult(result, schema, "scan"); } async queryTable( table: string, keyConditionExpression: string, expressionAttributeValues: Record, schema: z.ZodSchema, expressionAttributeNames?: Record, metadata?: TraceMetadata, ): Promise { const params: Record = { TableName: table, KeyConditionExpression: keyConditionExpression, ExpressionAttributeValues: expressionAttributeValues, }; if ( expressionAttributeNames && Object.keys(expressionAttributeNames).length > 0 ) { params.ExpressionAttributeNames = expressionAttributeNames; } const request = this.buildRequest("query", params); const result = await this.executeWithErrorHandling( request, "queryTable", metadata, ); return this.validateResult(result, schema, "queryTable"); } async batchWriteItem( requestItems: Record>>, metadata?: TraceMetadata, ): Promise { const request = this.buildRequest("batchWriteItem", { RequestItems: requestItems, }); return this.executeWithErrorHandling(request, "batchWriteItem", metadata); } async listTables( schema: z.ZodSchema, metadata?: TraceMetadata, ): Promise { const request = this.buildRequest("listTables", {}); const result = await this.executeWithErrorHandling( request, "listTables", metadata, ); return this.validateResult(result, schema, "listTables"); } async describeTable( table: string, schema: z.ZodSchema, metadata?: TraceMetadata, ): Promise { const request = this.buildRequest("describeTable", { TableName: table, }); const result = await this.executeWithErrorHandling( request, "describeTable", metadata, ); return this.validateResult(result, schema, "describeTable"); } async deleteTable(table: string, metadata?: TraceMetadata): Promise { const request = this.buildRequest("deleteTable", { TableName: table, }); return this.executeWithErrorHandling(request, "deleteTable", metadata); } }