Fix Phase 1 security and data integrity audit findings (#49)

BLOCK-01: Add API Gateway JWT authorizer with Cognito, route internal
Lambda calls through Function URL to bypass gateway auth
BLOCK-02/03: Prevent proposal number race condition with pg_advisory_xact_lock
and filter revision numbers from max-number query
BLOCK-04: Restrict VendorProposals and GeneratedPdfs to admins/sysadmins
BLOCK-05: Sum all vendor costs instead of overwriting with single vendor
BLOCK-06: Enable ValidateAudience on JWT, add Auth env vars to API Lambda
BLOCK-07: Validate ID token signature in AuthController via OIDC discovery
BLOCK-08: Use batchItemFailures in all Lambda SQS handlers
BLOCK-09: Increase SQS visibility timeout from 180s to 720s
FIX-10: Scope dispatcher queries to own proposals (IDOR fix)
This commit is contained in:
Adam Moussa 2026-05-20 18:51:31 -04:00 • committed by GitHub
parent dfbd0562ee
commit 091c5fcb44
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 182 additions and 53 deletions

View file

@ -6,7 +6,7 @@ Internal proposal management platform for Sea Haven Industries. Dispatchers subm
Monorepo with five primary services: 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 19 Web** -- MUI v7 admin/dispatcher workspace served via CloudFront + S3
- **React Native Mobile** -- iOS-first field app for dispatchers (offline-capable) - **React Native Mobile** -- iOS-first field app for dispatchers (offline-capable)
- **Python Lambdas** -- PDF extraction, PDF generation, library ingestion, AI suggestions, AOSS index provisioning - **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 | | CDK Stack | Key Resources |
|---|---| |---|---|
| `proposal-system-foundation` | RDS PostgreSQL 15 (t4g.small), S3 buckets, SQS queue + DLQ, Cognito user pool, Secrets Manager | | `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) | | `proposal-system-frontend` | CloudFront distribution (S3 OAC) |
| Resource Type | Names | | Resource Type | Names |
|---|---| |---|---|
| S3 Buckets | `proposal-system-uploads`, `proposal-system-generated`, `proposal-system-library`, `seahaven-ios-certificates` | | 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` | | Secrets | `proposal-system/db-credentials`, `proposal-system/internal-api-key` |
## Local Development ## 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_ISSUER_ID` | App Store Connect issuer |
| `ASC_KEY_CONTENT` | App Store Connect API key (base64) | | `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 ## Data Flow
1. Dispatcher submits proposal request (web or mobile) 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 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 4. Suggestions Lambda queries Bedrock KB for similar proposals, generates line items via Claude
5. Admin reviews/edits line items in pricing workspace 5. Admin reviews/edits line items in pricing workspace
6. On approval: `pdf-generate` Lambda creates branded PDF 6. On approval: `pdf-generate` Lambda creates branded PDF
7. On send: `library-ingest` Lambda adds approved proposal to KB for future matching 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.

View file

@ -5,6 +5,8 @@ using System.Text;
using System.Text.Json; using System.Text.Json;
using Microsoft.AspNetCore.Mvc; using Microsoft.AspNetCore.Mvc;
using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore;
using Microsoft.IdentityModel.Protocols;
using Microsoft.IdentityModel.Protocols.OpenIdConnect;
using Microsoft.IdentityModel.Tokens; using Microsoft.IdentityModel.Tokens;
using ProposalSystem.Domain.Entities; using ProposalSystem.Domain.Entities;
using ProposalSystem.Infrastructure.Data; using ProposalSystem.Infrastructure.Data;
@ -40,7 +42,35 @@ public class AuthController : ControllerBase
return BadRequest(new { message = "Failed to exchange authorization code" }); return BadRequest(new { message = "Failed to exchange authorization code" });
var handler = new JwtSecurityTokenHandler(); 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<OpenIdConnectConfiguration>(
$"{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 var sub = idToken.Claims.FirstOrDefault(c => c.Type == "sub")?.Value
?? throw new InvalidOperationException("No sub claim in ID token"); ?? throw new InvalidOperationException("No sub claim in ID token");
@ -71,12 +101,17 @@ public class AuthController : ControllerBase
_db.Users.Add(user); _db.Users.Add(user);
await _db.SaveChangesAsync(ct); await _db.SaveChangesAsync(ct);
} }
else if (user.Email != email || user.DisplayName != name) else
{ {
user.Email = email; var changed = false;
user.DisplayName = name; if (user.Email != email) { user.Email = email; changed = true; }
user.UpdatedAt = DateTime.UtcNow; if (user.DisplayName != name) { user.DisplayName = name; changed = true; }
await _db.SaveChangesAsync(ct); if (user.Role != role) { user.Role = role; changed = true; }
if (changed)
{
user.UpdatedAt = DateTime.UtcNow;
await _db.SaveChangesAsync(ct);
}
} }
return Ok(new AuthResponse( return Ok(new AuthResponse(

View file

@ -9,7 +9,7 @@ namespace ProposalSystem.Api.Controllers;
[ApiController] [ApiController]
[Route("api/generated-pdfs")] [Route("api/generated-pdfs")]
[Authorize] [Authorize(Roles = "admins,sysadmins")]
public class GeneratedPdfsController : ControllerBase public class GeneratedPdfsController : ControllerBase
{ {
private readonly ProposalDbContext _db; private readonly ProposalDbContext _db;

View file

@ -8,7 +8,7 @@ namespace ProposalSystem.Api.Controllers;
[ApiController] [ApiController]
[Route("api/vendor-proposals")] [Route("api/vendor-proposals")]
[Authorize] [Authorize(Roles = "admins,sysadmins")]
public class VendorProposalsController : ControllerBase public class VendorProposalsController : ControllerBase
{ {
private readonly ProposalDbContext _db; private readonly ProposalDbContext _db;
@ -38,14 +38,13 @@ public class VendorProposalsController : ControllerBase
await _db.SaveChangesAsync(ct); 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); proposal.VendorTotalCost = await _db.VendorProposals
if (proposal != null) .Where(v => v.ProposalId == vendor.ProposalId)
{ .SumAsync(v => v.TotalVendorCost, ct);
proposal.VendorTotalCost = vendor.TotalVendorCost; await _db.SaveChangesAsync(ct);
await _db.SaveChangesAsync(ct);
}
} }
return NoContent(); return NoContent();

View file

@ -63,11 +63,14 @@ if (!string.IsNullOrEmpty(cognitoAuthority))
.AddJwtBearer(options => .AddJwtBearer(options =>
{ {
options.Authority = cognitoAuthority; options.Authority = cognitoAuthority;
var webClientId = builder.Configuration["COGNITO_WEB_CLIENT_ID"] ?? "";
var mobileClientId = builder.Configuration["COGNITO_MOBILE_CLIENT_ID"] ?? "";
options.TokenValidationParameters = new TokenValidationParameters options.TokenValidationParameters = new TokenValidationParameters
{ {
ValidateIssuerSigningKey = true, ValidateIssuerSigningKey = true,
ValidateIssuer = true, ValidateIssuer = true,
ValidateAudience = false, ValidateAudience = true,
ValidAudiences = new[] { webClientId, mobileClientId }.Where(s => !string.IsNullOrEmpty(s)).ToList(),
ValidateLifetime = true, ValidateLifetime = true,
RoleClaimType = "cognito:groups", RoleClaimType = "cognito:groups",
}; };

View file

@ -18,8 +18,13 @@ public class ProposalNumberGenerator : IProposalNumberGenerator
var year = DateTime.UtcNow.Year; var year = DateTime.UtcNow.Year;
var prefix = $"SHI-{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 var lastNumber = await _db.Proposals
.Where(p => p.ProposalNumber.StartsWith(prefix)) .Where(p => p.ProposalNumber.StartsWith(prefix))
.Where(p => !p.ProposalNumber.Contains("-R"))
.OrderByDescending(p => p.ProposalNumber) .OrderByDescending(p => p.ProposalNumber)
.Select(p => p.ProposalNumber) .Select(p => p.ProposalNumber)
.FirstOrDefaultAsync(ct); .FirstOrDefaultAsync(ct);

View file

@ -30,6 +30,8 @@ public class ProposalService : IProposalService
public async Task<ProposalResponse> CreateAsync(CreateProposalRequest request, CancellationToken ct = default) public async Task<ProposalResponse> CreateAsync(CreateProposalRequest request, CancellationToken ct = default)
{ {
await using var transaction = await _db.Database.BeginTransactionAsync(ct);
var proposalNumber = await _numberGenerator.GenerateAsync(ct); var proposalNumber = await _numberGenerator.GenerateAsync(ct);
var now = DateTime.UtcNow; var now = DateTime.UtcNow;
@ -53,6 +55,7 @@ public class ProposalService : IProposalService
_db.Proposals.Add(proposal); _db.Proposals.Add(proposal);
await _db.SaveChangesAsync(ct); await _db.SaveChangesAsync(ct);
await transaction.CommitAsync(ct);
await _audit.LogAsync(AuditAction.Submit, proposal.Id, null, ct); await _audit.LogAsync(AuditAction.Submit, proposal.Id, null, ct);
@ -69,7 +72,10 @@ public class ProposalService : IProposalService
.Include(p => p.ApprovedBy) .Include(p => p.ApprovedBy)
.FirstOrDefaultAsync(p => p.Id == id, ct); .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<PagedResponse<ProposalListResponse>> GetAllAsync(ProposalFilterRequest filter, CancellationToken ct = default) public async Task<PagedResponse<ProposalListResponse>> GetAllAsync(ProposalFilterRequest filter, CancellationToken ct = default)
@ -79,7 +85,7 @@ public class ProposalService : IProposalService
.Include(p => p.AssignedAdmin) .Include(p => p.AssignedAdmin)
.AsQueryable(); .AsQueryable();
if (filter.Mine) if (_currentUser.Role == UserRole.Dispatcher || filter.Mine)
{ {
query = query.Where(p => p.SubmittedById == _currentUser.UserId); query = query.Where(p => p.SubmittedById == _currentUser.UserId);
} }

View file

@ -28,6 +28,8 @@ const compute = new ComputeStack(app, 'proposal-system-compute', {
libraryBucket: foundation.libraryBucket, libraryBucket: foundation.libraryBucket,
jobsQueue: foundation.jobsQueue, jobsQueue: foundation.jobsQueue,
userPool: foundation.userPool, userPool: foundation.userPool,
webClientId: foundation.webClientId,
mobileClientId: foundation.mobileClientId,
}); });
new FrontendStack(app, 'proposal-system-frontend', { new FrontendStack(app, 'proposal-system-frontend', {

View file

@ -2,6 +2,7 @@ import * as cdk from 'aws-cdk-lib';
import * as ec2 from 'aws-cdk-lib/aws-ec2'; import * as ec2 from 'aws-cdk-lib/aws-ec2';
import * as lambda from 'aws-cdk-lib/aws-lambda'; import * as lambda from 'aws-cdk-lib/aws-lambda';
import * as apigatewayv2 from 'aws-cdk-lib/aws-apigatewayv2'; 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 apigatewayv2Integrations from 'aws-cdk-lib/aws-apigatewayv2-integrations';
import * as iam from 'aws-cdk-lib/aws-iam'; import * as iam from 'aws-cdk-lib/aws-iam';
import * as s3 from 'aws-cdk-lib/aws-s3'; import * as s3 from 'aws-cdk-lib/aws-s3';
@ -24,6 +25,8 @@ export interface ComputeStackProps extends cdk.StackProps {
libraryBucket: s3.IBucket; libraryBucket: s3.IBucket;
jobsQueue: sqs.IQueue; jobsQueue: sqs.IQueue;
userPool: cognito.IUserPool; userPool: cognito.IUserPool;
webClientId: string;
mobileClientId: string;
} }
export class ComputeStack extends cdk.Stack { export class ComputeStack extends cdk.Stack {
@ -209,6 +212,11 @@ export class ComputeStack extends cdk.Stack {
LIBRARY_BUCKET: props.libraryBucket.bucketName, LIBRARY_BUCKET: props.libraryBucket.bucketName,
JOBS_QUEUE_URL: props.jobsQueue.queueUrl, JOBS_QUEUE_URL: props.jobsQueue.queueUrl,
INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, 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, tracing: lambda.Tracing.ACTIVE,
logRetention: logs.RetentionDays.TWO_MONTHS, logRetention: logs.RetentionDays.TWO_MONTHS,
@ -226,6 +234,11 @@ export class ComputeStack extends cdk.Stack {
resources: [props.userPool.userPoolArn], 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 // API Gateway HTTP API
const httpApi = new apigatewayv2.HttpApi(this, 'HttpApi', { const httpApi = new apigatewayv2.HttpApi(this, 'HttpApi', {
apiName: 'proposal-system-gateway', apiName: 'proposal-system-gateway',
@ -251,10 +264,29 @@ export class ComputeStack extends cdk.Stack {
apiFunction 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({ httpApi.addRoutes({
path: '/{proxy+}', path: '/{proxy+}',
methods: [apigatewayv2.HttpMethod.ANY], methods: [apigatewayv2.HttpMethod.ANY],
integration: apiIntegration, integration: apiIntegration,
authorizer: jwtAuthorizer,
}); });
// Python Lambda: Suggestions Engine // Python Lambda: Suggestions Engine
@ -272,7 +304,7 @@ export class ComputeStack extends cdk.Stack {
environment: { environment: {
KNOWLEDGE_BASE_ID: knowledgeBase.attrKnowledgeBaseId, KNOWLEDGE_BASE_ID: knowledgeBase.attrKnowledgeBaseId,
MODEL_ID: 'us.anthropic.claude-sonnet-4-5-20250929-v1:0', 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, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn,
}, },
logRetention: logs.RetentionDays.TWO_MONTHS, logRetention: logs.RetentionDays.TWO_MONTHS,
@ -303,7 +335,7 @@ export class ComputeStack extends cdk.Stack {
environment: { environment: {
UPLOADS_BUCKET: props.uploadsBucket.bucketName, UPLOADS_BUCKET: props.uploadsBucket.bucketName,
MODEL_ID: 'us.anthropic.claude-sonnet-4-5-20250929-v1:0', 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, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn,
}, },
logRetention: logs.RetentionDays.TWO_MONTHS, logRetention: logs.RetentionDays.TWO_MONTHS,
@ -330,7 +362,7 @@ export class ComputeStack extends cdk.Stack {
securityGroups: [props.lambdaSecurityGroup], securityGroups: [props.lambdaSecurityGroup],
environment: { environment: {
GENERATED_BUCKET: props.generatedBucket.bucketName, GENERATED_BUCKET: props.generatedBucket.bucketName,
API_BASE_URL: httpApi.apiEndpoint, API_BASE_URL: apiFunctionUrl.url,
INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn,
}, },
logRetention: logs.RetentionDays.TWO_MONTHS, logRetention: logs.RetentionDays.TWO_MONTHS,
@ -355,7 +387,7 @@ export class ComputeStack extends cdk.Stack {
LIBRARY_BUCKET: props.libraryBucket.bucketName, LIBRARY_BUCKET: props.libraryBucket.bucketName,
KNOWLEDGE_BASE_ID: knowledgeBase.attrKnowledgeBaseId, KNOWLEDGE_BASE_ID: knowledgeBase.attrKnowledgeBaseId,
DATA_SOURCE_ID: dataSource.attrDataSourceId, DATA_SOURCE_ID: dataSource.attrDataSourceId,
API_BASE_URL: httpApi.apiEndpoint, API_BASE_URL: apiFunctionUrl.url,
INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn, INTERNAL_API_KEY_SECRET_ARN: internalApiKeySecret.secretArn,
}, },
logRetention: logs.RetentionDays.TWO_MONTHS, logRetention: logs.RetentionDays.TWO_MONTHS,
@ -371,6 +403,7 @@ export class ComputeStack extends cdk.Stack {
// SQS Event Sources with message filtering // SQS Event Sources with message filtering
suggestionsFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { suggestionsFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, {
batchSize: 1, batchSize: 1,
reportBatchItemFailures: true,
filters: [ filters: [
lambda.FilterCriteria.filter({ lambda.FilterCriteria.filter({
body: { jobType: lambda.FilterRule.isEqual('suggestions') }, body: { jobType: lambda.FilterRule.isEqual('suggestions') },
@ -380,6 +413,7 @@ export class ComputeStack extends cdk.Stack {
pdfExtractFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { pdfExtractFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, {
batchSize: 1, batchSize: 1,
reportBatchItemFailures: true,
filters: [ filters: [
lambda.FilterCriteria.filter({ lambda.FilterCriteria.filter({
body: { jobType: lambda.FilterRule.isEqual('pdf-extract') }, body: { jobType: lambda.FilterRule.isEqual('pdf-extract') },
@ -389,6 +423,7 @@ export class ComputeStack extends cdk.Stack {
pdfGenerateFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { pdfGenerateFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, {
batchSize: 1, batchSize: 1,
reportBatchItemFailures: true,
filters: [ filters: [
lambda.FilterCriteria.filter({ lambda.FilterCriteria.filter({
body: { jobType: lambda.FilterRule.isEqual('pdf-generate') }, body: { jobType: lambda.FilterRule.isEqual('pdf-generate') },
@ -398,6 +433,7 @@ export class ComputeStack extends cdk.Stack {
libraryIngestFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, { libraryIngestFunction.addEventSource(new lambdaEventSources.SqsEventSource(props.jobsQueue, {
batchSize: 1, batchSize: 1,
reportBatchItemFailures: true,
filters: [ filters: [
lambda.FilterCriteria.filter({ lambda.FilterCriteria.filter({
body: { jobType: lambda.FilterRule.isEqual('library-ingest') }, body: { jobType: lambda.FilterRule.isEqual('library-ingest') },

View file

@ -17,6 +17,8 @@ export class FoundationStack extends cdk.Stack {
public readonly libraryBucket: s3.IBucket; public readonly libraryBucket: s3.IBucket;
public readonly jobsQueue: sqs.IQueue; public readonly jobsQueue: sqs.IQueue;
public readonly userPool: cognito.IUserPool; public readonly userPool: cognito.IUserPool;
public readonly webClientId: string;
public readonly mobileClientId: string;
constructor(scope: Construct, id: string, props?: cdk.StackProps) { constructor(scope: Construct, id: string, props?: cdk.StackProps) {
super(scope, id, props); super(scope, id, props);
@ -156,7 +158,7 @@ export class FoundationStack extends cdk.Stack {
this.jobsQueue = new sqs.Queue(this, 'JobsQueue', { this.jobsQueue = new sqs.Queue(this, 'JobsQueue', {
queueName: 'proposal-system-jobs', queueName: 'proposal-system-jobs',
visibilityTimeout: cdk.Duration.seconds(180), visibilityTimeout: cdk.Duration.seconds(720),
deadLetterQueue: { deadLetterQueue: {
queue: dlq, queue: dlq,
maxReceiveCount: 3, maxReceiveCount: 3,
@ -234,6 +236,8 @@ export class FoundationStack extends cdk.Stack {
}, },
}); });
this.webClientId = webClient.userPoolClientId;
// Mobile App Client (PKCE) // Mobile App Client (PKCE)
const mobileClient = userPool.addClient('MobileClient', { const mobileClient = userPool.addClient('MobileClient', {
userPoolClientName: 'proposal-system-mobile', userPoolClientName: 'proposal-system-mobile',
@ -253,6 +257,8 @@ export class FoundationStack extends cdk.Stack {
}, },
}); });
this.mobileClientId = mobileClient.userPoolClientId;
// CloudWatch Log Groups // CloudWatch Log Groups
const logGroupNames = [ const logGroupNames = [
'proposal-system-api', 'proposal-system-api',

View file

@ -41,12 +41,17 @@ def _get_api_key() -> str:
def handler(event, context): def handler(event, context):
batch_item_failures = []
for record in event.get("Records", []): for record in event.get("Records", []):
body = json.loads(record["body"]) try:
payload = body.get("payload", body) body = json.loads(record["body"])
proposal_id = payload["proposalId"] payload = body.get("payload", body)
process_ingestion(proposal_id) proposal_id = payload["proposalId"]
return {"statusCode": 200} 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): def process_ingestion(proposal_id: str):

View file

@ -41,19 +41,24 @@ def _get_api_key() -> str:
def handler(event, context): def handler(event, context):
batch_item_failures = []
for record in event.get("Records", []): for record in event.get("Records", []):
body = json.loads(record["body"]) try:
payload = body.get("payload", body) body = json.loads(record["body"])
proposal_id = payload["proposalId"] payload = body.get("payload", body)
s3_key = payload.get("s3Key", "") proposal_id = payload["proposalId"]
vendor_proposal_id = payload.get("vendorProposalId", "") s3_key = payload.get("s3Key", "")
vendor_proposal_id = payload.get("vendorProposalId", "")
if not s3_key: if not s3_key:
logger.warning("No s3Key in payload for proposal %s", proposal_id) logger.warning("No s3Key in payload for proposal %s", proposal_id)
continue continue
process_pdf(proposal_id, s3_key, vendor_proposal_id) process_pdf(proposal_id, s3_key, vendor_proposal_id)
return {"statusCode": 200} 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): def process_pdf(proposal_id: str, s3_key: str, vendor_proposal_id: str):

View file

@ -65,12 +65,17 @@ def _get_api_key() -> str:
def handler(event, context): def handler(event, context):
batch_item_failures = []
for record in event.get("Records", []): for record in event.get("Records", []):
body = json.loads(record["body"]) try:
payload = body.get("payload", body) body = json.loads(record["body"])
proposal_id = payload["proposalId"] payload = body.get("payload", body)
generate_pdf(proposal_id) proposal_id = payload["proposalId"]
return {"statusCode": 200} 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): def generate_pdf(proposal_id: str):

View file

@ -38,13 +38,18 @@ def _get_api_key() -> str:
def handler(event, context): def handler(event, context):
batch_item_failures = []
for record in event.get("Records", []): for record in event.get("Records", []):
body = json.loads(record["body"]) try:
payload = body.get("payload", body) body = json.loads(record["body"])
proposal_id = payload["proposalId"] payload = body.get("payload", body)
trigger = payload.get("trigger", "generate") proposal_id = payload["proposalId"]
process_suggestion(proposal_id, trigger) trigger = payload.get("trigger", "generate")
return {"statusCode": 200} 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): def process_suggestion(proposal_id: str, trigger: str):