diff --git a/README.md b/README.md index 11bd291..de193c5 100644 --- a/README.md +++ b/README.md @@ -6,7 +6,7 @@ Internal proposal management platform for Sea Haven Industries. Dispatchers subm Monorepo with five primary services: -- **.NET 8 API** -- Clean Architecture REST API hosted on Lambda behind API Gateway +- **.NET 8 API** -- Clean Architecture REST API hosted on Lambda behind API Gateway (JWT-authorized) with Function URL for internal access - **React 19 Web** -- MUI v7 admin/dispatcher workspace served via CloudFront + S3 - **React Native Mobile** -- iOS-first field app for dispatchers (offline-capable) - **Python Lambdas** -- PDF extraction, PDF generation, library ingestion, AI suggestions, AOSS index provisioning @@ -46,13 +46,13 @@ All resources are in **us-east-1** (account 328440206208). | CDK Stack | Key Resources | |---|---| | `proposal-system-foundation` | RDS PostgreSQL 15 (t4g.small), S3 buckets, SQS queue + DLQ, Cognito user pool, Secrets Manager | -| `proposal-system-compute` | API Gateway HTTP API, .NET 8 API Lambda, Python Lambdas (pdf-extract, pdf-generate, library-ingest, suggestions, oss-index-creator), OpenSearch Serverless collection, Bedrock KB | +| `proposal-system-compute` | API Gateway HTTP API (JWT authorizer), .NET 8 API Lambda + Function URL, Python Lambdas (pdf-extract, pdf-generate, library-ingest, suggestions, oss-index-creator), OpenSearch Serverless collection, Bedrock KB | | `proposal-system-frontend` | CloudFront distribution (S3 OAC) | | Resource Type | Names | |---|---| | S3 Buckets | `proposal-system-uploads`, `proposal-system-generated`, `proposal-system-library`, `seahaven-ios-certificates` | -| SQS | `proposal-system-jobs` + `proposal-system-jobs-dlq` (message body filtering by jobType) | +| SQS | `proposal-system-jobs` (720s visibility, reportBatchItemFailures) + `proposal-system-jobs-dlq` (message body filtering by jobType) | | Secrets | `proposal-system/db-credentials`, `proposal-system/internal-api-key` | ## Local Development @@ -149,12 +149,29 @@ Build and upload to TestFlight is handled by the `cd-mobile-ios.yaml` reusable w | `ASC_ISSUER_ID` | App Store Connect issuer | | `ASC_KEY_CONTENT` | App Store Connect API key (base64) | +## Authentication & Authorization + +Two-layer auth architecture with defense-in-depth: + +| Path | Authorizer | Authentication | +|---|---|---| +| External clients → API Gateway `/{proxy+}` | Cognito JWT authorizer (web + mobile client IDs) | .NET JWT middleware (ValidateAudience=true) | +| `/api/health` | None (public) | None | +| `/api/auth/callback`, `/api/auth/dev-login` | None (unauthenticated) | None (pre-auth endpoints) | +| Internal Lambdas → Function URL | None (NONE auth type) | Internal API key (`X-Internal-Api-Key` header, value from Secrets Manager) | + +**Role-based access:** Cognito groups (`dispatchers`, `admins`, `sysadmins`) map to API roles via `cognito:groups` claim. Dispatchers can only see their own proposals (ownership enforced in service layer). VendorProposals and GeneratedPdfs endpoints restricted to admins/sysadmins. + +**Internal API key:** Python Lambdas call the .NET API via a Lambda Function URL (bypasses API Gateway JWT check). The `InternalApiKeyMiddleware` validates the key and assigns the `admins` role to the synthetic identity. + ## Data Flow 1. Dispatcher submits proposal request (web or mobile) -2. API creates proposal record, publishes SQS message +2. API creates proposal record (with advisory-locked number generation), publishes SQS message 3. If vendor PDF attached: `pdf-extract` Lambda parses and structures data 4. Suggestions Lambda queries Bedrock KB for similar proposals, generates line items via Claude 5. Admin reviews/edits line items in pricing workspace 6. On approval: `pdf-generate` Lambda creates branded PDF 7. On send: `library-ingest` Lambda adds approved proposal to KB for future matching + +Failed SQS messages are reported via `batchItemFailures` and retried up to 3 times before moving to the DLQ. diff --git a/api/src/ProposalSystem.Api/Controllers/AuthController.cs b/api/src/ProposalSystem.Api/Controllers/AuthController.cs index 19dcf12..37463ee 100644 --- a/api/src/ProposalSystem.Api/Controllers/AuthController.cs +++ b/api/src/ProposalSystem.Api/Controllers/AuthController.cs @@ -5,6 +5,8 @@ using System.Text; using System.Text.Json; using Microsoft.AspNetCore.Mvc; using Microsoft.EntityFrameworkCore; +using Microsoft.IdentityModel.Protocols; +using Microsoft.IdentityModel.Protocols.OpenIdConnect; using Microsoft.IdentityModel.Tokens; using ProposalSystem.Domain.Entities; using ProposalSystem.Infrastructure.Data; @@ -40,7 +42,35 @@ public class AuthController : ControllerBase return BadRequest(new { message = "Failed to exchange authorization code" }); var handler = new JwtSecurityTokenHandler(); - var idToken = handler.ReadJwtToken(tokenResponse.IdToken); + + var authority = _config["Auth:Authority"]; + JwtSecurityToken idToken; + if (!string.IsNullOrEmpty(authority)) + { + var configManager = new ConfigurationManager( + $"{authority}/.well-known/openid-configuration", + new OpenIdConnectConfigurationRetriever(), + new HttpDocumentRetriever()); + var oidcConfig = await configManager.GetConfigurationAsync(ct); + + var validationParams = new TokenValidationParameters + { + ValidateIssuerSigningKey = true, + IssuerSigningKeys = oidcConfig.SigningKeys, + ValidateIssuer = true, + ValidIssuer = authority, + ValidateAudience = true, + ValidAudience = clientId, + ValidateLifetime = true, + }; + + handler.ValidateToken(tokenResponse.IdToken, validationParams, out var validatedToken); + idToken = (JwtSecurityToken)validatedToken; + } + else + { + idToken = handler.ReadJwtToken(tokenResponse.IdToken); + } var sub = idToken.Claims.FirstOrDefault(c => c.Type == "sub")?.Value ?? throw new InvalidOperationException("No sub claim in ID token"); @@ -71,12 +101,17 @@ public class AuthController : ControllerBase _db.Users.Add(user); await _db.SaveChangesAsync(ct); } - else if (user.Email != email || user.DisplayName != name) + else { - user.Email = email; - user.DisplayName = name; - user.UpdatedAt = DateTime.UtcNow; - await _db.SaveChangesAsync(ct); + var changed = false; + if (user.Email != email) { user.Email = email; changed = true; } + if (user.DisplayName != name) { user.DisplayName = name; changed = true; } + if (user.Role != role) { user.Role = role; changed = true; } + if (changed) + { + user.UpdatedAt = DateTime.UtcNow; + await _db.SaveChangesAsync(ct); + } } return Ok(new AuthResponse( diff --git a/api/src/ProposalSystem.Api/Controllers/GeneratedPdfsController.cs b/api/src/ProposalSystem.Api/Controllers/GeneratedPdfsController.cs index 8f3d181..e06c409 100644 --- a/api/src/ProposalSystem.Api/Controllers/GeneratedPdfsController.cs +++ b/api/src/ProposalSystem.Api/Controllers/GeneratedPdfsController.cs @@ -9,7 +9,7 @@ namespace ProposalSystem.Api.Controllers; [ApiController] [Route("api/generated-pdfs")] -[Authorize] +[Authorize(Roles = "admins,sysadmins")] public class GeneratedPdfsController : ControllerBase { private readonly ProposalDbContext _db; diff --git a/api/src/ProposalSystem.Api/Controllers/VendorProposalsController.cs b/api/src/ProposalSystem.Api/Controllers/VendorProposalsController.cs index f311a99..9b3b455 100644 --- a/api/src/ProposalSystem.Api/Controllers/VendorProposalsController.cs +++ b/api/src/ProposalSystem.Api/Controllers/VendorProposalsController.cs @@ -8,7 +8,7 @@ namespace ProposalSystem.Api.Controllers; [ApiController] [Route("api/vendor-proposals")] -[Authorize] +[Authorize(Roles = "admins,sysadmins")] public class VendorProposalsController : ControllerBase { private readonly ProposalDbContext _db; @@ -38,14 +38,13 @@ public class VendorProposalsController : ControllerBase await _db.SaveChangesAsync(ct); - if (vendor.TotalVendorCost > 0) + var proposal = await _db.Proposals.FindAsync(new object[] { vendor.ProposalId }, ct); + if (proposal != null) { - var proposal = await _db.Proposals.FindAsync(new object[] { vendor.ProposalId }, ct); - if (proposal != null) - { - proposal.VendorTotalCost = vendor.TotalVendorCost; - await _db.SaveChangesAsync(ct); - } + proposal.VendorTotalCost = await _db.VendorProposals + .Where(v => v.ProposalId == vendor.ProposalId) + .SumAsync(v => v.TotalVendorCost, ct); + await _db.SaveChangesAsync(ct); } return NoContent(); diff --git a/api/src/ProposalSystem.Api/Program.cs b/api/src/ProposalSystem.Api/Program.cs index 088f998..2024dda 100644 --- a/api/src/ProposalSystem.Api/Program.cs +++ b/api/src/ProposalSystem.Api/Program.cs @@ -63,11 +63,14 @@ if (!string.IsNullOrEmpty(cognitoAuthority)) .AddJwtBearer(options => { options.Authority = cognitoAuthority; + var webClientId = builder.Configuration["COGNITO_WEB_CLIENT_ID"] ?? ""; + var mobileClientId = builder.Configuration["COGNITO_MOBILE_CLIENT_ID"] ?? ""; options.TokenValidationParameters = new TokenValidationParameters { ValidateIssuerSigningKey = true, ValidateIssuer = true, - ValidateAudience = false, + ValidateAudience = true, + ValidAudiences = new[] { webClientId, mobileClientId }.Where(s => !string.IsNullOrEmpty(s)).ToList(), ValidateLifetime = true, RoleClaimType = "cognito:groups", }; diff --git a/api/src/ProposalSystem.Infrastructure/Services/ProposalNumberGenerator.cs b/api/src/ProposalSystem.Infrastructure/Services/ProposalNumberGenerator.cs index 2b14af2..ab7b1b7 100644 --- a/api/src/ProposalSystem.Infrastructure/Services/ProposalNumberGenerator.cs +++ b/api/src/ProposalSystem.Infrastructure/Services/ProposalNumberGenerator.cs @@ -18,8 +18,13 @@ public class ProposalNumberGenerator : IProposalNumberGenerator var year = DateTime.UtcNow.Year; var prefix = $"SHI-{year}-"; + // Advisory lock prevents concurrent number generation within the same transaction + await _db.Database.ExecuteSqlRawAsync( + "SELECT pg_advisory_xact_lock(hashtext('proposal_number_gen'))", ct); + var lastNumber = await _db.Proposals .Where(p => p.ProposalNumber.StartsWith(prefix)) + .Where(p => !p.ProposalNumber.Contains("-R")) .OrderByDescending(p => p.ProposalNumber) .Select(p => p.ProposalNumber) .FirstOrDefaultAsync(ct); diff --git a/api/src/ProposalSystem.Infrastructure/Services/ProposalService.cs b/api/src/ProposalSystem.Infrastructure/Services/ProposalService.cs index e807934..e41273c 100644 --- a/api/src/ProposalSystem.Infrastructure/Services/ProposalService.cs +++ b/api/src/ProposalSystem.Infrastructure/Services/ProposalService.cs @@ -30,6 +30,8 @@ public class ProposalService : IProposalService public async Task CreateAsync(CreateProposalRequest request, CancellationToken ct = default) { + await using var transaction = await _db.Database.BeginTransactionAsync(ct); + var proposalNumber = await _numberGenerator.GenerateAsync(ct); var now = DateTime.UtcNow; @@ -53,6 +55,7 @@ public class ProposalService : IProposalService _db.Proposals.Add(proposal); await _db.SaveChangesAsync(ct); + await transaction.CommitAsync(ct); await _audit.LogAsync(AuditAction.Submit, proposal.Id, null, ct); @@ -69,7 +72,10 @@ public class ProposalService : IProposalService .Include(p => p.ApprovedBy) .FirstOrDefaultAsync(p => p.Id == id, ct); - return proposal == null ? null : MapToResponse(proposal); + if (proposal == null) return null; + if (_currentUser.Role == UserRole.Dispatcher && proposal.SubmittedById != _currentUser.UserId) + return null; + return MapToResponse(proposal); } public async Task> GetAllAsync(ProposalFilterRequest filter, CancellationToken ct = default) @@ -79,7 +85,7 @@ public class ProposalService : IProposalService .Include(p => p.AssignedAdmin) .AsQueryable(); - if (filter.Mine) + if (_currentUser.Role == UserRole.Dispatcher || filter.Mine) { query = query.Where(p => p.SubmittedById == _currentUser.UserId); } diff --git a/infra/bin/app.ts b/infra/bin/app.ts index 759b825..ce76967 100644 --- a/infra/bin/app.ts +++ b/infra/bin/app.ts @@ -28,6 +28,8 @@ const compute = new ComputeStack(app, 'proposal-system-compute', { libraryBucket: foundation.libraryBucket, jobsQueue: foundation.jobsQueue, userPool: foundation.userPool, + webClientId: foundation.webClientId, + mobileClientId: foundation.mobileClientId, }); new FrontendStack(app, 'proposal-system-frontend', { diff --git a/infra/lib/compute-stack.ts b/infra/lib/compute-stack.ts index 98e624e..f0b46bd 100644 --- a/infra/lib/compute-stack.ts +++ b/infra/lib/compute-stack.ts @@ -2,6 +2,7 @@ import * as cdk from 'aws-cdk-lib'; import * as ec2 from 'aws-cdk-lib/aws-ec2'; import * as lambda from 'aws-cdk-lib/aws-lambda'; import * as apigatewayv2 from 'aws-cdk-lib/aws-apigatewayv2'; +import * as apigatewayv2Authorizers from 'aws-cdk-lib/aws-apigatewayv2-authorizers'; import * as apigatewayv2Integrations from 'aws-cdk-lib/aws-apigatewayv2-integrations'; import * as iam from 'aws-cdk-lib/aws-iam'; import * as s3 from 'aws-cdk-lib/aws-s3'; @@ -24,6 +25,8 @@ export interface ComputeStackProps extends cdk.StackProps { libraryBucket: s3.IBucket; jobsQueue: sqs.IQueue; userPool: cognito.IUserPool; + webClientId: string; + mobileClientId: string; } export class ComputeStack extends cdk.Stack { @@ -209,6 +212,11 @@ export class ComputeStack extends cdk.Stack { LIBRARY_BUCKET: props.libraryBucket.bucketName, JOBS_QUEUE_URL: props.jobsQueue.queueUrl, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, + Auth__Authority: `https://cognito-idp.${this.region}.amazonaws.com/${props.userPool.userPoolId}`, + Auth__ClientId: props.webClientId, + Auth__CognitoDomain: `proposal-system-seahaven.auth.${this.region}.amazoncognito.com`, + COGNITO_WEB_CLIENT_ID: props.webClientId, + COGNITO_MOBILE_CLIENT_ID: props.mobileClientId, }, tracing: lambda.Tracing.ACTIVE, logRetention: logs.RetentionDays.TWO_MONTHS, @@ -226,6 +234,11 @@ export class ComputeStack extends cdk.Stack { resources: [props.userPool.userPoolArn], })); + // Function URL for internal Lambda-to-API calls (bypasses API Gateway JWT authorizer) + const apiFunctionUrl = apiFunction.addFunctionUrl({ + authType: lambda.FunctionUrlAuthType.NONE, + }); + // API Gateway HTTP API const httpApi = new apigatewayv2.HttpApi(this, 'HttpApi', { apiName: 'proposal-system-gateway', @@ -251,10 +264,29 @@ export class ComputeStack extends cdk.Stack { apiFunction ); + const jwtAuthorizer = new apigatewayv2Authorizers.HttpJwtAuthorizer( + 'CognitoAuthorizer', + `https://cognito-idp.${this.region}.amazonaws.com/${props.userPool.userPoolId}`, + { jwtAudience: [props.webClientId, props.mobileClientId] }, + ); + + httpApi.addRoutes({ + path: '/api/health', + methods: [apigatewayv2.HttpMethod.GET], + integration: apiIntegration, + }); + + httpApi.addRoutes({ + path: '/api/auth/{proxy+}', + methods: [apigatewayv2.HttpMethod.POST], + integration: apiIntegration, + }); + httpApi.addRoutes({ path: '/{proxy+}', methods: [apigatewayv2.HttpMethod.ANY], integration: apiIntegration, + authorizer: jwtAuthorizer, }); // Python Lambda: Suggestions Engine @@ -272,7 +304,7 @@ export class ComputeStack extends cdk.Stack { environment: { KNOWLEDGE_BASE_ID: knowledgeBase.attrKnowledgeBaseId, MODEL_ID: 'us.anthropic.claude-sonnet-4-5-20250929-v1:0', - API_BASE_URL: httpApi.apiEndpoint, + API_BASE_URL: apiFunctionUrl.url, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, }, logRetention: logs.RetentionDays.TWO_MONTHS, @@ -303,7 +335,7 @@ export class ComputeStack extends cdk.Stack { environment: { UPLOADS_BUCKET: props.uploadsBucket.bucketName, MODEL_ID: 'us.anthropic.claude-sonnet-4-5-20250929-v1:0', - API_BASE_URL: httpApi.apiEndpoint, + API_BASE_URL: apiFunctionUrl.url, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, }, logRetention: logs.RetentionDays.TWO_MONTHS, @@ -330,7 +362,7 @@ export class ComputeStack extends cdk.Stack { securityGroups: [props.lambdaSecurityGroup], environment: { GENERATED_BUCKET: props.generatedBucket.bucketName, - API_BASE_URL: httpApi.apiEndpoint, + API_BASE_URL: apiFunctionUrl.url, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, }, logRetention: logs.RetentionDays.TWO_MONTHS, @@ -355,7 +387,7 @@ export class ComputeStack extends cdk.Stack { LIBRARY_BUCKET: props.libraryBucket.bucketName, KNOWLEDGE_BASE_ID: knowledgeBase.attrKnowledgeBaseId, DATA_SOURCE_ID: dataSource.attrDataSourceId, - API_BASE_URL: httpApi.apiEndpoint, + API_BASE_URL: apiFunctionUrl.url, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, }, logRetention: logs.RetentionDays.TWO_MONTHS, @@ -371,6 +403,7 @@ export class ComputeStack extends cdk.Stack { // SQS Event Sources with message filtering suggestionsFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { batchSize: 1, + reportBatchItemFailures: true, filters: [ lambda.FilterCriteria.filter({ body: { jobType: lambda.FilterRule.isEqual('suggestions') }, @@ -380,6 +413,7 @@ export class ComputeStack extends cdk.Stack { pdfExtractFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { batchSize: 1, + reportBatchItemFailures: true, filters: [ lambda.FilterCriteria.filter({ body: { jobType: lambda.FilterRule.isEqual('pdf-extract') }, @@ -389,6 +423,7 @@ export class ComputeStack extends cdk.Stack { pdfGenerateFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { batchSize: 1, + reportBatchItemFailures: true, filters: [ lambda.FilterCriteria.filter({ body: { jobType: lambda.FilterRule.isEqual('pdf-generate') }, @@ -398,6 +433,7 @@ export class ComputeStack extends cdk.Stack { libraryIngestFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { batchSize: 1, + reportBatchItemFailures: true, filters: [ lambda.FilterCriteria.filter({ body: { jobType: lambda.FilterRule.isEqual('library-ingest') }, diff --git a/infra/lib/foundation-stack.ts b/infra/lib/foundation-stack.ts index 0f2ad3e..f04e40a 100644 --- a/infra/lib/foundation-stack.ts +++ b/infra/lib/foundation-stack.ts @@ -17,6 +17,8 @@ export class FoundationStack extends cdk.Stack { public readonly libraryBucket: s3.IBucket; public readonly jobsQueue: sqs.IQueue; public readonly userPool: cognito.IUserPool; + public readonly webClientId: string; + public readonly mobileClientId: string; constructor(scope: Construct, id: string, props?: cdk.StackProps) { super(scope, id, props); @@ -156,7 +158,7 @@ export class FoundationStack extends cdk.Stack { this.jobsQueue = new sqs.Queue(this, 'JobsQueue', { queueName: 'proposal-system-jobs', - visibilityTimeout: cdk.Duration.seconds(180), + visibilityTimeout: cdk.Duration.seconds(720), deadLetterQueue: { queue: dlq, maxReceiveCount: 3, @@ -234,6 +236,8 @@ export class FoundationStack extends cdk.Stack { }, }); + this.webClientId = webClient.userPoolClientId; + // Mobile App Client (PKCE) const mobileClient = userPool.addClient('MobileClient', { userPoolClientName: 'proposal-system-mobile', @@ -253,6 +257,8 @@ export class FoundationStack extends cdk.Stack { }, }); + this.mobileClientId = mobileClient.userPoolClientId; + // CloudWatch Log Groups const logGroupNames = [ 'proposal-system-api', diff --git a/lambdas/library-ingest/app.py b/lambdas/library-ingest/app.py index eca3195..0533df8 100644 --- a/lambdas/library-ingest/app.py +++ b/lambdas/library-ingest/app.py @@ -41,12 +41,17 @@ def _get_api_key() -> str: def handler(event, context): + batch_item_failures = [] for record in event.get("Records", []): - body = json.loads(record["body"]) - payload = body.get("payload", body) - proposal_id = payload["proposalId"] - process_ingestion(proposal_id) - return {"statusCode": 200} + try: + body = json.loads(record["body"]) + payload = body.get("payload", body) + proposal_id = payload["proposalId"] + process_ingestion(proposal_id) + except Exception as e: + logger.error("Failed to process record %s: %s", record.get("messageId"), e) + batch_item_failures.append({"itemIdentifier": record["messageId"]}) + return {"batchItemFailures": batch_item_failures} def process_ingestion(proposal_id: str): diff --git a/lambdas/pdf-extract/app.py b/lambdas/pdf-extract/app.py index f1bc293..26a0ec5 100644 --- a/lambdas/pdf-extract/app.py +++ b/lambdas/pdf-extract/app.py @@ -41,19 +41,24 @@ def _get_api_key() -> str: def handler(event, context): + batch_item_failures = [] for record in event.get("Records", []): - body = json.loads(record["body"]) - payload = body.get("payload", body) - proposal_id = payload["proposalId"] - s3_key = payload.get("s3Key", "") - vendor_proposal_id = payload.get("vendorProposalId", "") + try: + body = json.loads(record["body"]) + payload = body.get("payload", body) + proposal_id = payload["proposalId"] + s3_key = payload.get("s3Key", "") + vendor_proposal_id = payload.get("vendorProposalId", "") - if not s3_key: - logger.warning("No s3Key in payload for proposal %s", proposal_id) - continue + if not s3_key: + logger.warning("No s3Key in payload for proposal %s", proposal_id) + continue - process_pdf(proposal_id, s3_key, vendor_proposal_id) - return {"statusCode": 200} + process_pdf(proposal_id, s3_key, vendor_proposal_id) + except Exception as e: + logger.error("Failed to process record %s: %s", record.get("messageId"), e) + batch_item_failures.append({"itemIdentifier": record["messageId"]}) + return {"batchItemFailures": batch_item_failures} def process_pdf(proposal_id: str, s3_key: str, vendor_proposal_id: str): diff --git a/lambdas/pdf-generate/app.py b/lambdas/pdf-generate/app.py index 37e75db..62ce665 100644 --- a/lambdas/pdf-generate/app.py +++ b/lambdas/pdf-generate/app.py @@ -65,12 +65,17 @@ def _get_api_key() -> str: def handler(event, context): + batch_item_failures = [] for record in event.get("Records", []): - body = json.loads(record["body"]) - payload = body.get("payload", body) - proposal_id = payload["proposalId"] - generate_pdf(proposal_id) - return {"statusCode": 200} + try: + body = json.loads(record["body"]) + payload = body.get("payload", body) + proposal_id = payload["proposalId"] + generate_pdf(proposal_id) + except Exception as e: + logger.error("Failed to process record %s: %s", record.get("messageId"), e) + batch_item_failures.append({"itemIdentifier": record["messageId"]}) + return {"batchItemFailures": batch_item_failures} def generate_pdf(proposal_id: str): diff --git a/lambdas/suggestions/app.py b/lambdas/suggestions/app.py index 9e68f83..9d7ede2 100644 --- a/lambdas/suggestions/app.py +++ b/lambdas/suggestions/app.py @@ -38,13 +38,18 @@ def _get_api_key() -> str: def handler(event, context): + batch_item_failures = [] for record in event.get("Records", []): - body = json.loads(record["body"]) - payload = body.get("payload", body) - proposal_id = payload["proposalId"] - trigger = payload.get("trigger", "generate") - process_suggestion(proposal_id, trigger) - return {"statusCode": 200} + try: + body = json.loads(record["body"]) + payload = body.get("payload", body) + proposal_id = payload["proposalId"] + trigger = payload.get("trigger", "generate") + process_suggestion(proposal_id, trigger) + except Exception as e: + logger.error("Failed to process record %s: %s", record.get("messageId"), e) + batch_item_failures.append({"itemIdentifier": record["messageId"]}) + return {"batchItemFailures": batch_item_failures} def process_suggestion(proposal_id: str, trigger: str):