From 0d29451362102f6ee1c2008ef6b4111d42abe219 Mon Sep 17 00:00:00 2001 From: Ed Olivares <34591886+eudoroolivares2016@users.noreply.github.com> Date: Thu, 13 Aug 2026 16:27:28 -0400 Subject: [PATCH 01/12] CMR-11421: Pass the Cmr-Send-Kms-Metadata-Fixer to CMR during writeback npm audit fixing --- package-lock.json | 6 +++--- .../shared/__tests__/writeCorrectedMetadataToCmr.test.js | 7 +++++-- serverless/src/shared/writeCorrectedMetadataToCmr.js | 3 ++- 3 files changed, 10 insertions(+), 6 deletions(-) diff --git a/package-lock.json b/package-lock.json index 9f402bdf..b5c5e595 100644 --- a/package-lock.json +++ b/package-lock.json @@ -11491,9 +11491,9 @@ "license": "MIT" }, "node_modules/nanoid": { - "version": "3.3.16", - "resolved": "https://registry.npmjs.org/nanoid/-/nanoid-3.3.16.tgz", - "integrity": "sha512-bzlKTyNJ7+LdGIIwy8ijFpIqEQIvafahV7eYykJ8Cvh42EdJeODoJ6gUJXpQJvej1BddH8OqTXZNE/KfbWAu8Q==", + "version": "3.3.18", + "resolved": "https://registry.npmjs.org/nanoid/-/nanoid-3.3.18.tgz", + "integrity": "sha512-DTg4MJbGMWkfi6VZFdNt2/caMbQy4Ou+Op/hJQvGEWcnVfoA1QA+xzRKAzw9jD6+GVOOeYr/mIcuDSdug6F6+w==", "funding": [ { "type": "github", diff --git a/serverless/src/shared/__tests__/writeCorrectedMetadataToCmr.test.js b/serverless/src/shared/__tests__/writeCorrectedMetadataToCmr.test.js index ade04daa..d9ef05aa 100644 --- a/serverless/src/shared/__tests__/writeCorrectedMetadataToCmr.test.js +++ b/serverless/src/shared/__tests__/writeCorrectedMetadataToCmr.test.js @@ -102,7 +102,8 @@ describe('when writing corrected metadata to cmr', () => { Authorization: 'Bearer writer-token', 'Client-Id': 'kms-metadata-correction-service', 'Cmr-Validate-Keywords': 'false', - 'Cmr-Validate-Umm-C': 'false' + 'Cmr-Validate-Umm-C': 'false', + 'Cmr-Send-Kms-Metadata-Fixer': 'false' } }) }) @@ -125,7 +126,8 @@ describe('when writing corrected metadata to cmr', () => { Authorization: 'Bearer writer-token', 'Client-Id': 'kms-metadata-correction-service', 'Cmr-Validate-Keywords': 'true', - 'Cmr-Validate-Umm-C': 'true' + 'Cmr-Validate-Umm-C': 'true', + 'Cmr-Send-Kms-Metadata-Fixer': 'false' } })) }) @@ -287,6 +289,7 @@ describe('when writing corrected metadata to cmr', () => { headers: { Authorization: 'Bearer writer-token', 'Client-Id': 'kms-metadata-correction-service', + 'Cmr-Send-Kms-Metadata-Fixer': 'false', 'Cmr-Validate-Keywords': 'false', 'Cmr-Validate-Umm-C': 'false' } diff --git a/serverless/src/shared/writeCorrectedMetadataToCmr.js b/serverless/src/shared/writeCorrectedMetadataToCmr.js index cf091ec1..8dff6a58 100644 --- a/serverless/src/shared/writeCorrectedMetadataToCmr.js +++ b/serverless/src/shared/writeCorrectedMetadataToCmr.js @@ -354,7 +354,8 @@ export const writeCorrectedMetadataToCmr = async ({ Authorization: authorizationToken, 'Client-Id': CMR_WRITEBACK_CLIENT_ID, 'Cmr-Validate-Keywords': getValidationHeaderValue('CMR_WRITEBACK_VALIDATE_KEYWORDS'), - 'Cmr-Validate-Umm-C': getValidationHeaderValue('CMR_WRITEBACK_VALIDATE_UMM_C') + 'Cmr-Validate-Umm-C': getValidationHeaderValue('CMR_WRITEBACK_VALIDATE_UMM_C'), + 'Cmr-Send-Kms-Metadata-Fixer': 'false' // Explicitly disable the CMR metadata fixer for writeback requests to prevent loops } }) From 56930fb81211fb51bd4eb45a5fb1ebfb09d5f807 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 09:03:57 -0400 Subject: [PATCH 02/12] KMS-648: automate RDF export and environment mirroring` --- README.md | 4 + bin/deploy-bamboo.sh | 1 + cdk/app/lib/KmsStack.ts | 1 + cdk/app/lib/helper/KmsLambdaFunctions.ts | 43 ++ cdk/bin/main.ts | 1 + package-lock.json | 83 ++-- package.json | 1 + .../src/exportRdf/__tests__/handler.test.js | 97 +++++ serverless/src/exportRdf/handler.js | 77 ++++ .../src/mirrorRdf/__tests__/handler.test.js | 372 ++++++++++++++++++ serverless/src/mirrorRdf/handler.js | 223 +++++++++++ .../shared/__tests__/exportRdfToS3.test.js | 29 ++ serverless/src/shared/exportRdfToS3.js | 33 +- 13 files changed, 932 insertions(+), 33 deletions(-) create mode 100644 serverless/src/exportRdf/__tests__/handler.test.js create mode 100644 serverless/src/exportRdf/handler.js create mode 100644 serverless/src/mirrorRdf/__tests__/handler.test.js create mode 100644 serverless/src/mirrorRdf/handler.js diff --git a/README.md b/README.md index 95006809..4e4cc52c 100644 --- a/README.md +++ b/README.md @@ -475,6 +475,7 @@ export bamboo_SUBNET_ID_C={subnet #3} export bamboo_VPC_ID={your vpc id} export bamboo_RDF4J_USER_NAME=[your rdfdb user name] export bamboo_RDF4J_PASSWORD=[your rdfdb password] +export bamboo_RDF_MIRROR_SOURCE_ENV=[optional sit|uat|prod source for RDF mirroring] export bamboo_EDL_HOST=[edl host name] export bamboo_EDL_UID=[edl user id] export bamboo_EDL_PASSWORD=[edl password] @@ -501,6 +502,9 @@ Notes: - When configured, `bamboo_CMR_SYSTEM_TOKEN_PARAMETER_NAME` is the primary source for the CMR authorization value. `bamboo_CMR_WRITER_TOKEN` is used only as a fallback and must include the `Bearer` prefix. +- Set `bamboo_RDF_MIRROR_SOURCE_ENV` to `sit`, `uat`, or `prod` to enable the nightly published + and draft RDF mirror from that environment. Leave it empty to disable automatic imports; an + authenticated `POST /rdf/mirror` can also run the configured mirror manually. - Leave `bamboo_CMR_WRITEBACK_PROVIDERS` empty to disable provider rollout for CMR writeback. - Set `bamboo_CMR_WRITEBACK_VALIDATE_KEYWORDS` and `bamboo_CMR_WRITEBACK_VALIDATE_UMM_C` to `true` to reject writebacks that still fail CMR keyword or UMM-C validation. diff --git a/bin/deploy-bamboo.sh b/bin/deploy-bamboo.sh index f7022f7f..ffca0e31 100755 --- a/bin/deploy-bamboo.sh +++ b/bin/deploy-bamboo.sh @@ -61,6 +61,7 @@ dockerRun() { --env "VPC_ID=$bamboo_VPC_ID" \ --env "RDF4J_USER_NAME=$bamboo_RDF4J_USER_NAME" \ --env "RDF4J_PASSWORD=$bamboo_RDF4J_PASSWORD" \ + --env "RDF_MIRROR_SOURCE_ENV=${bamboo_RDF_MIRROR_SOURCE_ENV:-}" \ --env "EDL_PASSWORD=$bamboo_EDL_PASSWORD" \ --env "EDL_CLIENT_ID=$bamboo_EDL_CLIENT_ID" \ --env "CMR_BASE_URL=$bamboo_CMR_BASE_URL" \ diff --git a/cdk/app/lib/KmsStack.ts b/cdk/app/lib/KmsStack.ts index 3aa1cb24..f4ca8c9d 100644 --- a/cdk/app/lib/KmsStack.ts +++ b/cdk/app/lib/KmsStack.ts @@ -42,6 +42,7 @@ export interface KmsStackProps extends cdk.StackProps { AWS_ENDPOINT_URL?: string RDF_BUCKET_NAME: string RDF4J_PASSWORD: string + RDF_MIRROR_SOURCE_ENV?: string RDF4J_SERVICE_URL: string RDF4J_USER_NAME: string } diff --git a/cdk/app/lib/helper/KmsLambdaFunctions.ts b/cdk/app/lib/helper/KmsLambdaFunctions.ts index 8a3463e1..bfec9fc5 100644 --- a/cdk/app/lib/helper/KmsLambdaFunctions.ts +++ b/cdk/app/lib/helper/KmsLambdaFunctions.ts @@ -45,6 +45,7 @@ interface LambdaFunctionsProps { LOG_LEVEL?: string; RDF_BUCKET_NAME: string, RDF4J_PASSWORD: string; + RDF_MIRROR_SOURCE_ENV?: string; RDF4J_SERVICE_URL: string; RDF4J_USER_NAME: string; KEYWORD_EVENTS_TOPIC_ARN?: string; @@ -144,6 +145,7 @@ export class LambdaFunctions { this.createTreeOperationApiLambdas(scope) this.createNightlyCachePrimeCron(scope) this.createCrudOperationApiLambdas(scope) + this.createRdfMirrorApiAndCron(scope) this.createPublishEventBridgeWiring(scope) } @@ -405,6 +407,16 @@ export class LambdaFunctions { true ) + this.createApiLambda( + scope, + 'exportRdf/handler.js', + 'export-rdf', + 'exportRdf', + '/rdf/export', + 'POST', + true + ) + this.createApiLambda( scope, 'rebuildRedisCache/handler.js', @@ -661,6 +673,37 @@ export class LambdaFunctions { ) } + /** + * Creates the manually invokable RDF mirror endpoint and its optional nightly schedule. + * The schedule is omitted when no source environment is configured. + * @param {Construct} scope Construct scope. + * @private + */ + private createRdfMirrorApiAndCron(scope: Construct) { + const mirrorLambda = this.createApiLambda( + scope, + 'mirrorRdf/handler.js', + 'mirror-rdf', + 'mirrorRdf', + '/rdf/mirror', + 'POST', + true, + Duration.minutes(15), + 2048 + ) + + if (!this.props.environment.RDF_MIRROR_SOURCE_ENV) return + + this.setupCronJob( + scope, + mirrorLambda, + // EventBridge cron is UTC; 05:00 UTC is midnight EST (01:00 EDT). + 'cron(0 5 * * ? *)', + {}, + 'NightlyRdfMirror' + ) + } + /** * Sets up a CloudWatch Events Rule to trigger a Lambda function on a schedule. * diff --git a/cdk/bin/main.ts b/cdk/bin/main.ts index 5c2c4b69..821cf010 100644 --- a/cdk/bin/main.ts +++ b/cdk/bin/main.ts @@ -200,6 +200,7 @@ async function main() { : (lbStack?.rdf4jServiceUrl || process.env.RDF4J_SERVICE_URL || 'http://localhost:8081'), RDF4J_USER_NAME: process.env.RDF4J_USER_NAME || 'rdf4j', RDF4J_PASSWORD: process.env.RDF4J_PASSWORD || 'rdf4j', + RDF_MIRROR_SOURCE_ENV: process.env.RDF_MIRROR_SOURCE_ENV || '', RDF_BUCKET_NAME: process.env.RDF_BUCKET_NAME || 'kms-rdf-backup', CMR_BASE_URL: cmrBaseUrl, EDL_PASSWORD: process.env.EDL_PASSWORD || '', diff --git a/package-lock.json b/package-lock.json index b9553901..1ca7367d 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,6 +16,7 @@ "@aws-sdk/client-sns": "^3.997.0", "@aws-sdk/client-sqs": "^3.997.0", "@aws-sdk/client-ssm": "^3.1096.0", + "@aws-sdk/s3-request-presigner": "3.981.0", "@xmldom/xmldom": "^0.8.10", "compact-object-deep": "^1.0.0", "csv": "^6.3.11", @@ -846,16 +847,16 @@ } }, "node_modules/@aws-sdk/core": { - "version": "3.977.1", - "resolved": "https://registry.npmjs.org/@aws-sdk/core/-/core-3.977.1.tgz", - "integrity": "sha512-KVtQRtc00ES/y+Sc3vYXeP6pCIcNlBJCZOwvqSy8ZpVGmbM5+IG+AfhuTKQ2oXmIVqZJewaGMMpzPkywC6xg0w==", + "version": "3.977.7", + "resolved": "https://registry.npmjs.org/@aws-sdk/core/-/core-3.977.7.tgz", + "integrity": "sha512-I88Iov89NVmjSmJLKSv7Cn9M2J+a2942OkA8nZCbz+sl4ZeY4zEOcoLOrbt1GRfQ8zEQKnjAJdXixA3J/p1fDQ==", "license": "Apache-2.0", "dependencies": { - "@aws-sdk/types": "^3.974.2", - "@aws-sdk/xml-builder": "^3.972.37", + "@aws-sdk/types": "^3.974.3", + "@aws-sdk/xml-builder": "^3.972.38", "@aws/lambda-invoke-store": "^0.3.0", - "@smithy/core": "^3.29.8", - "@smithy/signature-v4": "^5.6.9", + "@smithy/core": "^3.31.1", + "@smithy/signature-v4": "^5.6.12", "@smithy/types": "^4.16.1", "bowser": "^2.11.0", "tslib": "^2.6.2" @@ -1326,6 +1327,25 @@ "node": ">=20.0.0" } }, + "node_modules/@aws-sdk/s3-request-presigner": { + "version": "3.981.0", + "resolved": "https://registry.npmjs.org/@aws-sdk/s3-request-presigner/-/s3-request-presigner-3.981.0.tgz", + "integrity": "sha512-p4XPExNNDCDus+65YgjBUvQPQek9zZfxullsIBre22cVa7yq3VBDtXsyVTTfTjKnIWW2K89gVEd7e7S7Xkm6jQ==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/signature-v4-multi-region": "3.981.0", + "@aws-sdk/types": "^3.973.1", + "@aws-sdk/util-format-url": "^3.972.3", + "@smithy/middleware-endpoint": "^4.4.12", + "@smithy/protocol-http": "^5.3.8", + "@smithy/smithy-client": "^4.11.1", + "@smithy/types": "^4.12.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, "node_modules/@aws-sdk/signature-v4-multi-region": { "version": "3.981.0", "resolved": "https://registry.npmjs.org/@aws-sdk/signature-v4-multi-region/-/signature-v4-multi-region-3.981.0.tgz", @@ -1360,9 +1380,9 @@ } }, "node_modules/@aws-sdk/types": { - "version": "3.974.2", - "resolved": "https://registry.npmjs.org/@aws-sdk/types/-/types-3.974.2.tgz", - "integrity": "sha512-3W6IUtSxFbH6X7Wb7DzGCV5QiFQsd0g8bOfntpmDxQlzBoKWUMBu/JPQR0DwkE+Hpnxd6db1tXbOwdeHddG6cA==", + "version": "3.974.3", + "resolved": "https://registry.npmjs.org/@aws-sdk/types/-/types-3.974.3.tgz", + "integrity": "sha512-ECAqfpNsef+7MO8qtR0h9KcFIBAygaE7Cm6UOiQl+ft+uVap+1G7bNEjs4mdJE2OnA4m6k7i8peH8uGIAsOMGw==", "license": "Apache-2.0", "dependencies": { "@smithy/types": "^4.16.1", @@ -1399,6 +1419,19 @@ "node": ">=20.0.0" } }, + "node_modules/@aws-sdk/util-format-url": { + "version": "3.972.44", + "resolved": "https://registry.npmjs.org/@aws-sdk/util-format-url/-/util-format-url-3.972.44.tgz", + "integrity": "sha512-MpVw+TzzP1lMaxDgJMYvWC9gHkQQytA+wPlj7hmshlNlptfae1LdWvnzMUnvQjkVVYTDcdUU1OPVCYpuerZXiw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.977.7", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, "node_modules/@aws-sdk/util-locate-window": { "version": "3.804.0", "resolved": "https://registry.npmjs.org/@aws-sdk/util-locate-window/-/util-locate-window-3.804.0.tgz", @@ -1449,9 +1482,9 @@ } }, "node_modules/@aws-sdk/xml-builder": { - "version": "3.972.37", - "resolved": "https://registry.npmjs.org/@aws-sdk/xml-builder/-/xml-builder-3.972.37.tgz", - "integrity": "sha512-zKq4HQum8JwDyEuyfuI4bbiAcU0KxP6qy+9PR/IsR92IyE/DaBAikzAS50tjxip4bqIIANpCcG+Yyj6CVhXupg==", + "version": "3.972.38", + "resolved": "https://registry.npmjs.org/@aws-sdk/xml-builder/-/xml-builder-3.972.38.tgz", + "integrity": "sha512-grf7mzfVxBS5AlsuTvBN7uDpzqohFww9fRPCO+EBSUdvtsYMcPSKdz54h/7XiscqNcUM1Ae1MF7JLHmiYYuzbQ==", "license": "Apache-2.0", "dependencies": { "@smithy/types": "^4.16.1", @@ -4513,12 +4546,12 @@ } }, "node_modules/@smithy/core": { - "version": "3.31.0", - "resolved": "https://registry.npmjs.org/@smithy/core/-/core-3.31.0.tgz", - "integrity": "sha512-sylYk2l9d7CmRv8ts8p0SDQUr3VO+HMeS1nrjL6+UtbO8ktJHTOeQ1McX+aAyvGGccp5aZX9eNtdcXrSwzoZaw==", + "version": "3.32.0", + "resolved": "https://registry.npmjs.org/@smithy/core/-/core-3.32.0.tgz", + "integrity": "sha512-NAiCSC78fzbNIEWoheoF74Ob5ZorLijCHpMY26Fqvqg/+9LuyIqMfHDg2p8Yk1rqOyowtiL3y7WX0AW+teL6zw==", "license": "Apache-2.0", "dependencies": { - "@smithy/types": "^4.16.1", + "@smithy/types": "^4.17.0", "tslib": "^2.6.2" }, "engines": { @@ -4896,13 +4929,13 @@ } }, "node_modules/@smithy/signature-v4": { - "version": "5.6.11", - "resolved": "https://registry.npmjs.org/@smithy/signature-v4/-/signature-v4-5.6.11.tgz", - "integrity": "sha512-7HsspeiNCZvZHEJ22vV5L/QYuJdTyJvPJvMrYD3AgkM3IJB0pkln4jkjPvtpTWRMkHXbO8WKwNjoVdVlBFwHmw==", + "version": "5.7.0", + "resolved": "https://registry.npmjs.org/@smithy/signature-v4/-/signature-v4-5.7.0.tgz", + "integrity": "sha512-hCynhm22wMJ8wTF9crcwu8mxggtUrSLLJgDcGUvYFBqpofxycYJCGKOMYg4xtPPFtgNiDJSYmhsWLTrcU/g59Q==", "license": "Apache-2.0", "dependencies": { - "@smithy/core": "^3.31.0", - "@smithy/types": "^4.16.1", + "@smithy/core": "^3.32.0", + "@smithy/types": "^4.17.0", "tslib": "^2.6.2" }, "engines": { @@ -4947,9 +4980,9 @@ } }, "node_modules/@smithy/types": { - "version": "4.16.1", - "resolved": "https://registry.npmjs.org/@smithy/types/-/types-4.16.1.tgz", - "integrity": "sha512-0JFs3V2y2M9tKW5na/qxe69Zv+uxLMO7QBbhxF/FHu/Gp2NFZAAL9tWl9PU02xxo07pb3G9FTyjNc6D5uZrJIg==", + "version": "4.17.0", + "resolved": "https://registry.npmjs.org/@smithy/types/-/types-4.17.0.tgz", + "integrity": "sha512-Aw4joiM0ZdErpo39lCj8phT2lxoiKZV+KZzBxnnQhWVtU2Is/WffQSL04uUWRcXUse9Ln8vXZK6V/FwqRVnQpg==", "license": "Apache-2.0", "dependencies": { "tslib": "^2.6.2" diff --git a/package.json b/package.json index 40dfd15b..b20b5c28 100644 --- a/package.json +++ b/package.json @@ -42,6 +42,7 @@ "@aws-sdk/client-sns": "^3.997.0", "@aws-sdk/client-sqs": "^3.997.0", "@aws-sdk/client-ssm": "^3.1096.0", + "@aws-sdk/s3-request-presigner": "3.981.0", "@xmldom/xmldom": "^0.8.10", "compact-object-deep": "^1.0.0", "csv": "^6.3.11", diff --git a/serverless/src/exportRdf/__tests__/handler.test.js b/serverless/src/exportRdf/__tests__/handler.test.js new file mode 100644 index 00000000..f56f25ae --- /dev/null +++ b/serverless/src/exportRdf/__tests__/handler.test.js @@ -0,0 +1,97 @@ +import { GetObjectCommand } from '@aws-sdk/client-s3' +import { getSignedUrl } from '@aws-sdk/s3-request-presigner' +import { + beforeEach, + describe, + expect, + test, + vi +} from 'vitest' + +import { getS3Client } from '@/shared/awsClients' +import { exportRdfToS3 } from '@/shared/exportRdfToS3' +import { logger } from '@/shared/logger' + +import { exportRdf } from '../handler' + +vi.mock('@aws-sdk/client-s3') +vi.mock('@aws-sdk/s3-request-presigner') +vi.mock('@/shared/awsClients') +vi.mock('@/shared/exportRdfToS3') +vi.mock('@/shared/logger') + +describe('exportRdf', () => { + beforeEach(() => { + vi.resetAllMocks() + vi.mocked(getS3Client).mockReturnValue({}) + vi.mocked(exportRdfToS3).mockResolvedValue({ + bucketName: 'kms-rdf-backup-test', + s3Key: '21.4/rdf.xml.gz' + }) + + vi.mocked(getSignedUrl).mockResolvedValue('https://example.com/rdf.xml.gz') + }) + + test('exports the requested graph and returns its temporary download URL', async () => { + const response = await exportRdf({ + queryStringParameters: { + version: 'published' + } + }) + + expect(exportRdfToS3).toHaveBeenCalledWith({ + version: 'published', + archive: true + }) + + expect(GetObjectCommand).toHaveBeenCalledWith({ + Bucket: 'kms-rdf-backup-test', + Key: '21.4/rdf.xml.gz', + ResponseContentDisposition: 'attachment; filename="kms-published-rdf.xml.gz"', + ResponseContentType: 'application/gzip' + }) + + expect(getSignedUrl).toHaveBeenCalledWith( + {}, + expect.any(GetObjectCommand), + { expiresIn: 300 } + ) + + expect(response).toEqual(expect.objectContaining({ + statusCode: 200, + body: JSON.stringify({ + version: 'published', + downloadUrl: 'https://example.com/rdf.xml.gz', + expiresIn: 300 + }) + })) + }) + + test('rejects a missing or unsupported version', async () => { + const unsupportedResponse = await exportRdf({ + queryStringParameters: { + version: 'invalid' + } + }) + const missingResponse = await exportRdf({}) + + expect(unsupportedResponse.statusCode).toBe(400) + expect(missingResponse.statusCode).toBe(400) + expect(exportRdfToS3).not.toHaveBeenCalled() + }) + + test('returns an internal error when the export fails', async () => { + vi.mocked(exportRdfToS3).mockRejectedValue(new Error('RDF4J unavailable')) + + const response = await exportRdf({ + queryStringParameters: { + version: 'draft' + } + }) + + expect(response.statusCode).toBe(500) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-export] Failed to create draft RDF export, error=Error: RDF4J unavailable' + ) + }) +}) diff --git a/serverless/src/exportRdf/handler.js b/serverless/src/exportRdf/handler.js new file mode 100644 index 00000000..5d09c77a --- /dev/null +++ b/serverless/src/exportRdf/handler.js @@ -0,0 +1,77 @@ +import { GetObjectCommand } from '@aws-sdk/client-s3' +import { getSignedUrl } from '@aws-sdk/s3-request-presigner' + +import { getS3Client } from '@/shared/awsClients' +import { exportRdfToS3 } from '@/shared/exportRdfToS3' +import { getApplicationConfig } from '@/shared/getConfig' +import { logger } from '@/shared/logger' + +const DOWNLOAD_URL_EXPIRATION_SECONDS = 300 + +/** + * Exports one RDF graph to private S3 gzip storage and returns a temporary download URL. + * + * @param {object} event - API Gateway event containing `version` as a query parameter. + * @returns {Promise} API Gateway response. + */ +export const exportRdf = async (event) => { + const { defaultResponseHeaders } = getApplicationConfig() + const version = event?.queryStringParameters?.version + + if (!['draft', 'published'].includes(version)) { + return { + statusCode: 400, + headers: defaultResponseHeaders, + body: JSON.stringify({ + error: 'version query parameter must be draft or published' + }) + } + } + + try { + const { + bucketName, + s3Key + } = await exportRdfToS3({ + version, + archive: true + }) + + const fileName = `kms-${version}-rdf.xml.gz` + const downloadUrl = await getSignedUrl( + getS3Client(), + new GetObjectCommand({ + Bucket: bucketName, + Key: s3Key, + ResponseContentDisposition: `attachment; filename="${fileName}"`, + ResponseContentType: 'application/gzip' + }), + { expiresIn: DOWNLOAD_URL_EXPIRATION_SECONDS } + ) + + return { + statusCode: 200, + headers: { + ...defaultResponseHeaders, + 'Content-Type': 'application/json' + }, + body: JSON.stringify({ + version, + downloadUrl, + expiresIn: DOWNLOAD_URL_EXPIRATION_SECONDS + }) + } + } catch (error) { + logger.error(`[rdf-export] Failed to create ${version} RDF export, error=${error.toString()}`) + + return { + statusCode: 500, + headers: defaultResponseHeaders, + body: JSON.stringify({ + error: 'Unable to create RDF export' + }) + } + } +} + +export default exportRdf diff --git a/serverless/src/mirrorRdf/__tests__/handler.test.js b/serverless/src/mirrorRdf/__tests__/handler.test.js new file mode 100644 index 00000000..f273635d --- /dev/null +++ b/serverless/src/mirrorRdf/__tests__/handler.test.js @@ -0,0 +1,372 @@ +import { gzipSync } from 'zlib' + +import { + beforeEach, + describe, + expect, + test, + vi +} from 'vitest' + +import { getCmrSystemToken } from '@/shared/getCmrWriterToken' +import { logger } from '@/shared/logger' + +import { mirrorRdf } from '../handler' + +vi.mock('@/shared/getConfig', () => ({ + getApplicationConfig: vi.fn(() => ({ + defaultResponseHeaders: { + 'Access-Control-Allow-Origin': '*' + } + })) +})) + +vi.mock('@/shared/getCmrWriterToken', () => ({ + getCmrSystemToken: vi.fn() +})) + +vi.mock('@/shared/logger') + +const createArchiveResponse = (rdfXml) => { + const archive = gzipSync(rdfXml) + const arrayBuffer = archive.buffer.slice( + archive.byteOffset, + archive.byteOffset + archive.byteLength + ) + + return { + ok: true, + arrayBuffer: vi.fn().mockResolvedValue(arrayBuffer) + } +} + +const createTextResponse = ({ + ok = true, + status = 200, + text = '' +} = {}) => ({ + ok, + status, + text: vi.fn().mockResolvedValue(text) +}) + +const configureSuccessfulFetch = () => { + vi.mocked(fetch) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('published')) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('draft')) + .mockResolvedValueOnce(createTextResponse()) + .mockResolvedValueOnce(createTextResponse()) + .mockResolvedValueOnce(createTextResponse()) + .mockResolvedValueOnce(createTextResponse()) +} + +describe('mirrorRdf', () => { + beforeEach(() => { + vi.resetAllMocks() + vi.stubGlobal('fetch', vi.fn()) + delete process.env.RDF_MIRROR_SOURCE_ENV + vi.mocked(getCmrSystemToken).mockResolvedValue('system-token') + process.env.RDF4J_SERVICE_URL = 'http://rdf4j.test:8080' + process.env.RDF4J_REPOSITORY_ID = 'kms-test' + process.env.RDF4J_USER_NAME = 'rdf-user' + process.env.RDF4J_PASSWORD = 'rdf-password' + }) + + test('skips the mirror when no source environment is configured', async () => { + const response = await mirrorRdf() + + expect(response.statusCode).toBe(200) + expect(JSON.parse(response.body)).toEqual({ + status: 'skipped', + reason: 'RDF_MIRROR_SOURCE_ENV is not configured' + }) + + expect(fetch).not.toHaveBeenCalled() + }) + + test('downloads both graphs before replacing the destination contexts', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'SIT' + configureSuccessfulFetch() + + const response = await mirrorRdf({ + requestContext: {}, + headers: { authorization: 'Bearer manual-token' } + }) + + expect(response.statusCode).toBe(200) + expect(JSON.parse(response.body)).toEqual({ + status: 'mirrored', + sourceEnvironment: 'sit', + versions: ['published', 'draft'] + }) + + expect(fetch).toHaveBeenNthCalledWith( + 1, + 'https://cmr.sit.earthdata.nasa.gov/kms/rdf/export?version=published', + { + method: 'POST', + headers: { + Accept: 'application/json', + Authorization: 'system-token' + } + } + ) + + expect(fetch).toHaveBeenNthCalledWith( + 3, + 'https://cmr.sit.earthdata.nasa.gov/kms/rdf/export?version=draft', + expect.any(Object) + ) + + const publishedContext = encodeURIComponent( + '' + ) + expect(String(vi.mocked(fetch).mock.calls[4][0])).toContain(`context=${publishedContext}`) + expect(vi.mocked(fetch).mock.calls[4][1]).toEqual({ + method: 'DELETE', + headers: { + Authorization: `Basic ${Buffer.from('rdf-user:rdf-password').toString('base64')}` + } + }) + + expect(vi.mocked(fetch).mock.calls[5][1]).toEqual(expect.objectContaining({ + method: 'POST', + body: 'published' + })) + + expect(logger.info).toHaveBeenCalledWith('[rdf-mirror] Mirrored RDF graphs', { + sourceEnvironment: 'sit', + versions: ['published', 'draft'] + }) + }) + + test('uses the CMR system token for a scheduled invocation', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'uat' + vi.mocked(getCmrSystemToken).mockResolvedValue('scheduled-system-token') + configureSuccessfulFetch() + + await mirrorRdf({ source: 'aws.events' }) + + expect(fetch).toHaveBeenNthCalledWith( + 1, + 'https://cmr.uat.earthdata.nasa.gov/kms/rdf/export?version=published', + expect.objectContaining({ + headers: { + Accept: 'application/json', + Authorization: 'scheduled-system-token' + } + }) + ) + }) + + test('returns an API error without clearing graphs when a source export fails', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(fetch).mockResolvedValueOnce(createTextResponse({ + ok: false, + status: 503, + text: 'unavailable' + })) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(fetch).toHaveBeenCalledTimes(1) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Source published export failed: 503 unavailable' + ) + }) + + test('returns an API error when the source export has no download URL', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(fetch).mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({}) + }) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(fetch).toHaveBeenCalledTimes(1) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Source published export did not return a download URL' + ) + }) + + test('returns an API error when the source gzip download fails', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(fetch) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createTextResponse({ + ok: false, + status: 504 + })) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(fetch).toHaveBeenCalledTimes(2) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Source published download failed: 504' + ) + }) + + test('returns an API error when the source download is not valid gzip', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(fetch) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) + }) + .mockResolvedValueOnce({ + ok: true, + arrayBuffer: vi.fn().mockResolvedValue(Buffer.from('not-gzip')) + }) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(logger.error).toHaveBeenCalledWith( + expect.stringContaining('Source published download is not valid gzip:') + ) + }) + + test('returns an API error when rdf.xml does not contain RDF/XML', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(fetch) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('')) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Source published download does not contain RDF/XML' + ) + }) + + test('returns an API error when clearing a destination graph fails', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(fetch) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('published')) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('draft')) + .mockResolvedValueOnce(createTextResponse({ + ok: false, + status: 500, + text: 'clear failed' + })) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Failed to clear destination published graph: 500 clear failed' + ) + }) + + test('returns an API error when importing a destination graph fails', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(fetch) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('published')) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('draft')) + .mockResolvedValueOnce(createTextResponse()) + .mockResolvedValueOnce(createTextResponse({ + ok: false, + status: 500, + text: 'import failed' + })) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Failed to import destination published graph: 500 import failed' + ) + }) + + test('uses local RDF4J defaults and permits an absent destination context', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + delete process.env.RDF4J_SERVICE_URL + delete process.env.RDF4J_REPOSITORY_ID + delete process.env.RDF4J_USER_NAME + delete process.env.RDF4J_PASSWORD + vi.mocked(fetch) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('published')) + .mockResolvedValueOnce({ + ok: true, + json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) + }) + .mockResolvedValueOnce(createArchiveResponse('draft')) + .mockResolvedValueOnce(createTextResponse({ + ok: false, + status: 404 + })) + .mockResolvedValueOnce(createTextResponse()) + .mockResolvedValueOnce(createTextResponse()) + .mockResolvedValueOnce(createTextResponse()) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(200) + expect(String(vi.mocked(fetch).mock.calls[4][0])).toContain( + 'http://localhost:8081/rdf4j-server/repositories/kms/statements' + ) + + expect(vi.mocked(fetch).mock.calls[4][1]).toEqual({ + method: 'DELETE', + headers: { + Authorization: `Basic ${Buffer.from('rdf4j:rdf4j').toString('base64')}` + } + }) + }) + + test('does not call the source when the CMR system token is unavailable', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(getCmrSystemToken).mockResolvedValue(undefined) + + const response = await mirrorRdf({ requestContext: {} }) + + expect(response.statusCode).toBe(500) + expect(fetch).not.toHaveBeenCalled() + }) + + test('throws a scheduled invocation error so EventBridge can report the failure', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'invalid' + + await expect(mirrorRdf({ source: 'aws.events' })) + .rejects.toThrow('RDF_MIRROR_SOURCE_ENV must be sit, uat, or prod') + }) +}) diff --git a/serverless/src/mirrorRdf/handler.js b/serverless/src/mirrorRdf/handler.js new file mode 100644 index 00000000..2787c563 --- /dev/null +++ b/serverless/src/mirrorRdf/handler.js @@ -0,0 +1,223 @@ +import { promisify } from 'util' +import zlib from 'zlib' + +import { getCmrSystemToken } from '@/shared/getCmrWriterToken' +import { getApplicationConfig } from '@/shared/getConfig' +import { logger } from '@/shared/logger' + +const RDF_VERSIONS = ['published', 'draft'] +const SOURCE_BASE_URLS = { + sit: 'https://cmr.sit.earthdata.nasa.gov/kms', + uat: 'https://cmr.uat.earthdata.nasa.gov/kms', + prod: 'https://cmr.earthdata.nasa.gov/kms' +} +const gunzip = promisify(zlib.gunzip) + +/** + * Resolves the configured source environment to its KMS API base URL. + * + * @returns {{sourceEnvironment: string, sourceBaseUrl: string}|undefined} Source configuration. + * @throws {Error} When the configured environment is unsupported. + */ +const getSourceConfiguration = () => { + const sourceEnvironment = String(process.env.RDF_MIRROR_SOURCE_ENV || '') + .trim() + .toLowerCase() + + if (!sourceEnvironment) return undefined + + const sourceBaseUrl = SOURCE_BASE_URLS[sourceEnvironment] + + if (!sourceBaseUrl) { + throw new Error('RDF_MIRROR_SOURCE_ENV must be sit, uat, or prod') + } + + return { + sourceEnvironment, + sourceBaseUrl + } +} + +/** + * Requests and downloads one gzip-compressed RDF graph from the source KMS environment. + * + * @param {object} params Download parameters. + * @param {string} params.authorization Authorization value forwarded to the source KMS API. + * @param {string} params.sourceBaseUrl Source KMS API base URL. + * @param {'published'|'draft'} params.version RDF graph version. + * @returns {Promise} Uncompressed RDF/XML content. + */ +const downloadRdf = async ({ + authorization, + sourceBaseUrl, + version +}) => { + const headers = { + Accept: 'application/json', + Authorization: authorization + } + const exportResponse = await fetch(`${sourceBaseUrl}/rdf/export?version=${version}`, { + method: 'POST', + headers + }) + + if (!exportResponse.ok) { + const responseText = await exportResponse.text() + throw new Error(`Source ${version} export failed: ${exportResponse.status} ${responseText}`) + } + + const { downloadUrl } = await exportResponse.json() + + if (!downloadUrl) { + throw new Error(`Source ${version} export did not return a download URL`) + } + + const downloadResponse = await fetch(downloadUrl) + + if (!downloadResponse.ok) { + throw new Error(`Source ${version} download failed: ${downloadResponse.status}`) + } + + let rdfXml + try { + const compressedRdf = Buffer.from(await downloadResponse.arrayBuffer()) + rdfXml = (await gunzip(compressedRdf)).toString('utf8') + } catch (error) { + throw new Error(`Source ${version} download is not valid gzip: ${error.message}`) + } + + if (!rdfXml.includes('} + */ +const replaceDestinationGraph = async ({ rdfXml, version }) => { + const serviceUrl = String(process.env.RDF4J_SERVICE_URL || 'http://localhost:8081') + .replace(/\/$/, '') + const repositoryId = process.env.RDF4J_REPOSITORY_ID || 'kms' + const statementsUrl = new URL(`${serviceUrl}/rdf4j-server/repositories/${repositoryId}/statements`) + const graphUri = `https://gcmd.earthdata.nasa.gov/kms/version/${version}` + const credentials = Buffer.from( + `${process.env.RDF4J_USER_NAME || 'rdf4j'}:${process.env.RDF4J_PASSWORD || 'rdf4j'}` + ).toString('base64') + const authorization = `Basic ${credentials}` + + statementsUrl.searchParams.set('context', `<${graphUri}>`) + + const clearResponse = await fetch(statementsUrl, { + method: 'DELETE', + headers: { Authorization: authorization } + }) + + if (!clearResponse.ok && clearResponse.status !== 404) { + const responseText = await clearResponse.text() + throw new Error(`Failed to clear destination ${version} graph: ${clearResponse.status} ${responseText}`) + } + + const importResponse = await fetch(statementsUrl, { + method: 'POST', + headers: { + Authorization: authorization, + 'Content-Type': 'application/rdf+xml' + }, + body: rdfXml + }) + + if (!importResponse.ok) { + const responseText = await importResponse.text() + throw new Error(`Failed to import destination ${version} graph: ${importResponse.status} ${responseText}`) + } +} + +/** + * Downloads published and draft RDF from the configured source environment and mirrors + * both graphs into the local RDF4J repository. + * + * @param {object} event API Gateway or EventBridge invocation event. + * @returns {Promise} API Gateway-compatible result. + */ +export const mirrorRdf = async (event = {}) => { + const { defaultResponseHeaders } = getApplicationConfig() + const isApiRequest = Boolean(event.requestContext) + + try { + const sourceConfiguration = getSourceConfiguration() + + if (!sourceConfiguration) { + return { + statusCode: 200, + headers: defaultResponseHeaders, + body: JSON.stringify({ + status: 'skipped', + reason: 'RDF_MIRROR_SOURCE_ENV is not configured' + }) + } + } + + const authorization = await getCmrSystemToken() + + if (!authorization) { + throw new Error('CMR system token is unavailable') + } + + // Complete every source download before modifying either destination graph. + const rdfByVersion = await RDF_VERSIONS.reduce( + (previousDownloads, version) => previousDownloads.then(async (downloads) => ({ + ...downloads, + [version]: await downloadRdf({ + authorization, + sourceBaseUrl: sourceConfiguration.sourceBaseUrl, + version + }) + })), + Promise.resolve({}) + ) + + await RDF_VERSIONS.reduce( + (previousImport, version) => previousImport.then(() => replaceDestinationGraph({ + rdfXml: rdfByVersion[version], + version + })), + Promise.resolve() + ) + + logger.info('[rdf-mirror] Mirrored RDF graphs', { + sourceEnvironment: sourceConfiguration.sourceEnvironment, + versions: RDF_VERSIONS + }) + + return { + statusCode: 200, + headers: defaultResponseHeaders, + body: JSON.stringify({ + status: 'mirrored', + sourceEnvironment: sourceConfiguration.sourceEnvironment, + versions: RDF_VERSIONS + }) + } + } catch (error) { + logger.error(`[rdf-mirror] Failed to mirror RDF graphs, error=${error.toString()}`) + + if (!isApiRequest) throw error + + return { + statusCode: 500, + headers: defaultResponseHeaders, + body: JSON.stringify({ + error: 'Unable to mirror RDF graphs' + }) + } + } +} + +export default mirrorRdf diff --git a/serverless/src/shared/__tests__/exportRdfToS3.test.js b/serverless/src/shared/__tests__/exportRdfToS3.test.js index 8264af81..a4b344cf 100644 --- a/serverless/src/shared/__tests__/exportRdfToS3.test.js +++ b/serverless/src/shared/__tests__/exportRdfToS3.test.js @@ -1,3 +1,5 @@ +import { gunzipSync } from 'zlib' + import { PutObjectCommand, S3Client } from '@aws-sdk/client-s3' import { beforeEach, @@ -84,6 +86,33 @@ describe('exportRdfToS3', () => { })) }) + test('should compress an RDF export with gzip', async () => { + const result = await exportRdfToS3({ + version: 'published', + archive: true + }) + + const upload = vi.mocked(PutObjectCommand).mock.calls[0][0] + + expect(result).toEqual({ + bucketName: 'kms-rdf-backup-test', + s3Key: '21.4/rdf.xml.gz' + }) + + expect(upload).toEqual(expect.objectContaining({ + ContentType: 'application/gzip' + })) + + expect(gunzipSync(upload.Body).toString('utf8')).toBe('Test RDF Data') + }) + + test('should reject an unsupported RDF version', async () => { + await expect(exportRdfToS3({ version: 'invalid' })) + .rejects.toThrow('Unsupported RDF version: invalid') + + expect(sparqlRequest).not.toHaveBeenCalled() + }) + test('should handle S3 upload failure', async () => { S3Client.prototype.send.mockRejectedValueOnce(new Error('S3 upload failed')) diff --git a/serverless/src/shared/exportRdfToS3.js b/serverless/src/shared/exportRdfToS3.js index 015f711e..9e75738c 100644 --- a/serverless/src/shared/exportRdfToS3.js +++ b/serverless/src/shared/exportRdfToS3.js @@ -1,3 +1,6 @@ +import { promisify } from 'util' +import zlib from 'zlib' + import { PutObjectCommand } from '@aws-sdk/client-s3' import { getS3Client } from '@/shared/awsClients' @@ -6,6 +9,8 @@ import { getApplicationConfig } from '@/shared/getConfig' import { getVersionMetadata } from '@/shared/getVersionMetadata' import { sparqlRequest } from '@/shared/sparqlRequest' +const gzip = promisify(zlib.gzip) + /** * Exports RDF data to an S3 bucket. * @@ -25,11 +30,16 @@ import { sparqlRequest } from '@/shared/sparqlRequest' * @async * @function exportRdfToS3 * @param {Object} params - The parameters for the export. - * @param {string} params.version - The version of the RDF data to export (e.g., 'published', 'draft'). - * @returns {Promise<{s3Key: string}>} A promise that resolves to an object containing the S3 key of the exported file. + * @param {string} params.version - The RDF graph to export (`published` or `draft`). + * @param {boolean} [params.archive=false] - Whether to compress the RDF with gzip. + * @returns {Promise<{bucketName: string, s3Key: string}>} The exported object's location. * @throws Will throw an error if the RDF data fetch fails or if there are issues with S3 operations. */ -export const exportRdfToS3 = async ({ version }) => { +export const exportRdfToS3 = async ({ version, archive = false }) => { + if (!['draft', 'published'].includes(version)) { + throw new Error(`Unsupported RDF version: ${version}`) + } + const { env } = getApplicationConfig() const s3BucketName = `kms-rdf-backup-${env}` const s3Client = getS3Client() @@ -52,28 +62,35 @@ export const exportRdfToS3 = async ({ version }) => { const rdfData = await response.text() + const fileName = archive ? 'rdf.xml.gz' : 'rdf.xml' + // Generate the S3 key based on the current date and version let s3Key if (version === 'published') { const { versionName } = await getVersionMetadata(version) - s3Key = `${versionName}/rdf.xml` + s3Key = `${versionName}/${fileName}` } else { const currentDate = new Date() const year = currentDate.getUTCFullYear() const month = String(currentDate.getUTCMonth() + 1).padStart(2, '0') const day = String(currentDate.getUTCDate()).padStart(2, '0') - s3Key = `${version}/${year}/${month}/${day}/rdf.xml` + s3Key = `${version}/${year}/${month}/${day}/${fileName}` } + const body = archive ? await gzip(rdfData) : rdfData + // Upload RDF data to S3 await s3Client.send(new PutObjectCommand({ Bucket: s3BucketName, Key: s3Key, - Body: rdfData, - ContentType: 'application/rdf+xml' + Body: body, + ContentType: archive ? 'application/gzip' : 'application/rdf+xml' })) console.log(`RDF data for version ${version} exported successfully to ${s3Key}`) - return { s3Key } + return { + bucketName: s3BucketName, + s3Key + } } From b9d46d5bc46d6754289a1c0c5a0a09c3ea7c917f Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 09:19:26 -0400 Subject: [PATCH 03/12] KMS-648: Added script to download exports for testing. --- scripts/local/run_rdf_export_smoke.sh | 98 +++++++++++++++++++++++++++ 1 file changed, 98 insertions(+) create mode 100755 scripts/local/run_rdf_export_smoke.sh diff --git a/scripts/local/run_rdf_export_smoke.sh b/scripts/local/run_rdf_export_smoke.sh new file mode 100755 index 00000000..f304968b --- /dev/null +++ b/scripts/local/run_rdf_export_smoke.sh @@ -0,0 +1,98 @@ +#!/usr/bin/env bash + +set -euo pipefail + +# Downloads and validates the published and draft RDF gzip exports. +# +# Usage: +# KMS_AUTHORIZATION='' \ +# ./scripts/local/run_rdf_export_smoke.sh [outputDirectory] + +usage() { + cat <<'EOF' +Usage: + KMS_AUTHORIZATION='' \ + ./scripts/local/run_rdf_export_smoke.sh [outputDirectory] + +Environment: + KMS_AUTHORIZATION Authorization header value passed through exactly as provided. + KMS_BASE_URL Optional override for the KMS base URL. + +The output directory defaults to /tmp/kms-rdf-export-smoke-. +EOF +} + +if [[ "${1:-}" == "--help" || "${1:-}" == "-h" ]]; then + usage + exit 0 +fi + +ENVIRONMENT="${1:-}" + +if [[ -z "$ENVIRONMENT" ]]; then + usage >&2 + exit 1 +fi + +AUTHORIZATION_VALUE="${KMS_AUTHORIZATION:?Missing KMS_AUTHORIZATION environment variable.}" + +case "$ENVIRONMENT" in + sit) + DEFAULT_BASE_URL="https://cmr.sit.earthdata.nasa.gov/kms" + ;; + uat) + DEFAULT_BASE_URL="https://cmr.uat.earthdata.nasa.gov/kms" + ;; + prod) + DEFAULT_BASE_URL="https://cmr.earthdata.nasa.gov/kms" + ;; + *) + echo "Unsupported environment \"$ENVIRONMENT\". Expected one of: sit, uat, prod" >&2 + exit 1 + ;; +esac + +BASE_URL="${KMS_BASE_URL:-$DEFAULT_BASE_URL}" +OUTPUT_DIRECTORY="${2:-/tmp/kms-rdf-export-smoke-${ENVIRONMENT}}" +mkdir -p "$OUTPUT_DIRECTORY" + +download_export() { + local version="$1" + local response_file="${OUTPUT_DIRECTORY}/${version}.response.json" + local gzip_file="${OUTPUT_DIRECTORY}/${version}.rdf.xml.gz" + local rdf_file="${OUTPUT_DIRECTORY}/${version}.rdf.xml" + + echo "[rdf-export-smoke] POST ${BASE_URL}/rdf/export?version=${version}" >&2 + curl \ + --silent \ + --show-error \ + --fail-with-body \ + --request POST \ + --header "Authorization: ${AUTHORIZATION_VALUE}" \ + --header 'Accept: application/json' \ + "${BASE_URL}/rdf/export?version=${version}" \ + --output "$response_file" + + local download_url + download_url="$(jq --exit-status --raw-output '.downloadUrl' "$response_file")" + + curl \ + --silent \ + --show-error \ + --fail \ + --location \ + "$download_url" \ + --output "$gzip_file" + + gzip --test "$gzip_file" + gzip --decompress --stdout "$gzip_file" > "$rdf_file" + grep --quiet '&2 + echo "[rdf-export-smoke] Extracted RDF/XML: ${rdf_file}" >&2 +} + +download_export published +download_export draft + +echo "[rdf-export-smoke] Passed. Artifacts are in ${OUTPUT_DIRECTORY}" >&2 From 30b2aec19c913d660100bc071b07ef4e9b43bfd9 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 09:56:40 -0400 Subject: [PATCH 04/12] KMS-648: Added script for testing of updating local mirror. --- scripts/local/run_rdf_mirror_smoke.sh | 20 ++++++++++ .../src/mirrorRdf/__tests__/handler.test.js | 39 +++++++++++++------ serverless/src/mirrorRdf/handler.js | 17 +++++++- 3 files changed, 63 insertions(+), 13 deletions(-) create mode 100755 scripts/local/run_rdf_mirror_smoke.sh diff --git a/scripts/local/run_rdf_mirror_smoke.sh b/scripts/local/run_rdf_mirror_smoke.sh new file mode 100755 index 00000000..db182556 --- /dev/null +++ b/scripts/local/run_rdf_mirror_smoke.sh @@ -0,0 +1,20 @@ +#!/usr/bin/env bash + +set -euo pipefail + +# Usage: KMS_AUTHORIZATION='' ./scripts/local/run_rdf_mirror_smoke.sh + +KMS_BASE_URL="${KMS_BASE_URL:-http://127.0.0.1:3013}" +AUTHORIZATION_VALUE="${KMS_AUTHORIZATION:?Missing KMS_AUTHORIZATION environment variable.}" + +curl --silent --show-error --fail-with-body \ + --request POST \ + --header "Authorization: ${AUTHORIZATION_VALUE}" \ + "${KMS_BASE_URL}/rdf/mirror" | jq . + +curl --silent --show-error --fail-with-body \ + "${KMS_BASE_URL}/concept_versions/version_type/all" +printf '\n' + +curl --silent --show-error --fail-with-body "${KMS_BASE_URL}/status" +printf '\n' diff --git a/serverless/src/mirrorRdf/__tests__/handler.test.js b/serverless/src/mirrorRdf/__tests__/handler.test.js index f273635d..9040fe55 100644 --- a/serverless/src/mirrorRdf/__tests__/handler.test.js +++ b/serverless/src/mirrorRdf/__tests__/handler.test.js @@ -68,6 +68,11 @@ const configureSuccessfulFetch = () => { .mockResolvedValueOnce(createTextResponse()) } +const createApiEvent = () => ({ + requestContext: {}, + headers: { Authorization: 'api-token' } +}) + describe('mirrorRdf', () => { beforeEach(() => { vi.resetAllMocks() @@ -115,11 +120,13 @@ describe('mirrorRdf', () => { method: 'POST', headers: { Accept: 'application/json', - Authorization: 'system-token' + Authorization: 'Bearer manual-token' } } ) + expect(getCmrSystemToken).not.toHaveBeenCalled() + expect(fetch).toHaveBeenNthCalledWith( 3, 'https://cmr.sit.earthdata.nasa.gov/kms/rdf/export?version=draft', @@ -175,7 +182,7 @@ describe('mirrorRdf', () => { text: 'unavailable' })) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) expect(fetch).toHaveBeenCalledTimes(1) @@ -191,7 +198,7 @@ describe('mirrorRdf', () => { json: vi.fn().mockResolvedValue({}) }) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) expect(fetch).toHaveBeenCalledTimes(1) @@ -212,7 +219,7 @@ describe('mirrorRdf', () => { status: 504 })) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) expect(fetch).toHaveBeenCalledTimes(2) @@ -233,7 +240,7 @@ describe('mirrorRdf', () => { arrayBuffer: vi.fn().mockResolvedValue(Buffer.from('not-gzip')) }) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) expect(logger.error).toHaveBeenCalledWith( @@ -250,7 +257,7 @@ describe('mirrorRdf', () => { }) .mockResolvedValueOnce(createArchiveResponse('')) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) expect(logger.error).toHaveBeenCalledWith( @@ -277,7 +284,7 @@ describe('mirrorRdf', () => { text: 'clear failed' })) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) expect(logger.error).toHaveBeenCalledWith( @@ -305,7 +312,7 @@ describe('mirrorRdf', () => { text: 'import failed' })) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) expect(logger.error).toHaveBeenCalledWith( @@ -338,7 +345,7 @@ describe('mirrorRdf', () => { .mockResolvedValueOnce(createTextResponse()) .mockResolvedValueOnce(createTextResponse()) - const response = await mirrorRdf({ requestContext: {} }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(200) expect(String(vi.mocked(fetch).mock.calls[4][0])).toContain( @@ -353,14 +360,24 @@ describe('mirrorRdf', () => { }) }) - test('does not call the source when the CMR system token is unavailable', async () => { + test('does not call the source when an API request has no Authorization header', async () => { process.env.RDF_MIRROR_SOURCE_ENV = 'prod' - vi.mocked(getCmrSystemToken).mockResolvedValue(undefined) const response = await mirrorRdf({ requestContext: {} }) expect(response.statusCode).toBe(500) expect(fetch).not.toHaveBeenCalled() + expect(getCmrSystemToken).not.toHaveBeenCalled() + }) + + test('throws when a scheduled invocation has no CMR system token', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'prod' + vi.mocked(getCmrSystemToken).mockResolvedValue(undefined) + + await expect(mirrorRdf({ source: 'aws.events' })) + .rejects.toThrow('CMR system token is unavailable') + + expect(fetch).not.toHaveBeenCalled() }) test('throws a scheduled invocation error so EventBridge can report the failure', async () => { diff --git a/serverless/src/mirrorRdf/handler.js b/serverless/src/mirrorRdf/handler.js index 2787c563..fa5a6d85 100644 --- a/serverless/src/mirrorRdf/handler.js +++ b/serverless/src/mirrorRdf/handler.js @@ -13,6 +13,15 @@ const SOURCE_BASE_URLS = { } const gunzip = promisify(zlib.gunzip) +/** + * Reads the Authorization header from an API Gateway event without modifying its value. + * + * @param {object} event API Gateway invocation event. + * @returns {string|undefined} Incoming authorization value. + */ +const getRequestAuthorization = (event) => Object.entries(event.headers || {}) + .find(([headerName]) => headerName.toLowerCase() === 'authorization')?.[1] + /** * Resolves the configured source environment to its KMS API base URL. * @@ -164,10 +173,14 @@ export const mirrorRdf = async (event = {}) => { } } - const authorization = await getCmrSystemToken() + const authorization = isApiRequest + ? getRequestAuthorization(event) + : await getCmrSystemToken() if (!authorization) { - throw new Error('CMR system token is unavailable') + throw new Error(isApiRequest + ? 'Authorization header is unavailable' + : 'CMR system token is unavailable') } // Complete every source download before modifying either destination graph. From 20adf89c2fc03083be45b7cd9764f6bb08301090 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 12:08:34 -0400 Subject: [PATCH 05/12] KMS-648: Fix local startup for RDF mirror testing --- bin/start-local.sh | 10 ++++++---- cdk/bin/main.ts | 4 +++- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/bin/start-local.sh b/bin/start-local.sh index c6ca5927..ba15e327 100755 --- a/bin/start-local.sh +++ b/bin/start-local.sh @@ -1,5 +1,7 @@ #!/bin/bash +set -e + SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" PROJECT_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" # shellcheck source=bin/env/local_env.sh @@ -37,13 +39,13 @@ fi clearStaleSAMContainers +# Synthesize the CDK stack +cd "${PROJECT_ROOT}/cdk" +npx cdk synth --context useLocalstack="true" --output ./cdk.out > /dev/null + "${PROJECT_ROOT}/scripts/localstack/run_bridge.sh" & LOCAL_BRIDGE_PID=$! -# Synthesize the CDK stack -cd cdk -cdk synth --context useLocalstack="true" --output ./cdk.out > /dev/null 2>&1 - # Start SAM local sam local start-api \ --template-file ./cdk.out/KmsStack.template.json \ diff --git a/cdk/bin/main.ts b/cdk/bin/main.ts index 821cf010..a8c0dbbb 100644 --- a/cdk/bin/main.ts +++ b/cdk/bin/main.ts @@ -85,7 +85,9 @@ async function main() { } const vpcId = useLocalstack ? 'dummy-vpc-id' : process.env.VPC_ID - const cmrBaseUrl = requireEnv('CMR_BASE_URL') + const cmrBaseUrl = useLocalstack + ? process.env.CMR_BASE_URL || '' + : requireEnv('CMR_BASE_URL') if (!vpcId) { throw new Error('VPC_ID environment variable is not set') From bf7e07f8950b2f8c15de9675cc1fb84d45ae6ab5 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 12:25:48 -0400 Subject: [PATCH 06/12] KMS-648: Fix Vite path resolution during local CDK startup --- vite.config.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/vite.config.js b/vite.config.js index 17739eb0..1b7256f7 100644 --- a/vite.config.js +++ b/vite.config.js @@ -19,7 +19,7 @@ function getHandlerEntries(dir) { }, {}) } -const handlerEntries = getHandlerEntries('./serverless/src') +const handlerEntries = getHandlerEntries(path.resolve(__dirname, 'serverless/src')) export default defineConfig({ build: { From 7cc6cce2d982f4c30516dfc0bef5e2e228c34bab Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 12:58:32 -0400 Subject: [PATCH 07/12] KMS-648: Removed volume for logs --- bin/env/local_env.sh | 1 + bin/rdf4j/start.sh | 1 - scripts/local/run_rdf_mirror_smoke.sh | 4 ++- .../src/mirrorRdf/__tests__/handler.test.js | 25 ++++++++++++++++++- serverless/src/mirrorRdf/handler.js | 7 +++++- 5 files changed, 34 insertions(+), 4 deletions(-) diff --git a/bin/env/local_env.sh b/bin/env/local_env.sh index 3c06755c..30f61ba2 100644 --- a/bin/env/local_env.sh +++ b/bin/env/local_env.sh @@ -5,6 +5,7 @@ export RDF4J_SERVICE_URL="${RDF4J_SERVICE_URL:-http://rdf4j-server:8080}" export RDF4J_HOST_SERVICE_URL="${RDF4J_HOST_SERVICE_URL:-http://localhost:8081}" export RDF4J_USER_NAME="${RDF4J_USER_NAME:-rdf4j}" export RDF4J_PASSWORD="${RDF4J_PASSWORD:-rdf4j}" +export RDF_MIRROR_SOURCE_ENV="${RDF_MIRROR_SOURCE_ENV:-local}" export CMR_BASE_URL="${CMR_BASE_URL:-}" export RDF4J_CONTAINER_MEMORY_LIMIT="${RDF4J_CONTAINER_MEMORY_LIMIT:-2048}" export REDIS_ENABLED="${REDIS_ENABLED:-true}" diff --git a/bin/rdf4j/start.sh b/bin/rdf4j/start.sh index 27ad13d7..22fc80a1 100755 --- a/bin/rdf4j/start.sh +++ b/bin/rdf4j/start.sh @@ -27,5 +27,4 @@ docker run \ -e "RDF4J_USER_NAME=${RDF4J_USER_NAME}" \ -e "RDF4J_PASSWORD=${RDF4J_PASSWORD}" \ -e "RDF4J_CONTAINER_MEMORY_LIMIT=${RDF4J_CONTAINER_MEMORY_LIMIT}" \ - -v logs:/usr/local/tomcat/logs \ rdf4j:latest diff --git a/scripts/local/run_rdf_mirror_smoke.sh b/scripts/local/run_rdf_mirror_smoke.sh index db182556..65ad7ad8 100755 --- a/scripts/local/run_rdf_mirror_smoke.sh +++ b/scripts/local/run_rdf_mirror_smoke.sh @@ -2,7 +2,9 @@ set -euo pipefail -# Usage: KMS_AUTHORIZATION='' ./scripts/local/run_rdf_mirror_smoke.sh +# Start KMS with `npm run start-local`, then run: +# KMS_AUTHORIZATION='' ./scripts/local/run_rdf_mirror_smoke.sh +# Local startup defaults RDF_MIRROR_SOURCE_ENV to local, exercising local export and import. KMS_BASE_URL="${KMS_BASE_URL:-http://127.0.0.1:3013}" AUTHORIZATION_VALUE="${KMS_AUTHORIZATION:?Missing KMS_AUTHORIZATION environment variable.}" diff --git a/serverless/src/mirrorRdf/__tests__/handler.test.js b/serverless/src/mirrorRdf/__tests__/handler.test.js index 9040fe55..6e6452ff 100644 --- a/serverless/src/mirrorRdf/__tests__/handler.test.js +++ b/serverless/src/mirrorRdf/__tests__/handler.test.js @@ -78,6 +78,7 @@ describe('mirrorRdf', () => { vi.resetAllMocks() vi.stubGlobal('fetch', vi.fn()) delete process.env.RDF_MIRROR_SOURCE_ENV + delete process.env.AWS_SAM_LOCAL vi.mocked(getCmrSystemToken).mockResolvedValue('system-token') process.env.RDF4J_SERVICE_URL = 'http://rdf4j.test:8080' process.env.RDF4J_REPOSITORY_ID = 'kms-test' @@ -155,6 +156,21 @@ describe('mirrorRdf', () => { }) }) + test('uses the local SAM API as the source during local testing', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'local' + process.env.AWS_SAM_LOCAL = 'true' + configureSuccessfulFetch() + + const response = await mirrorRdf(createApiEvent()) + + expect(response.statusCode).toBe(200) + expect(fetch).toHaveBeenNthCalledWith( + 1, + 'http://host.docker.internal:3013/rdf/export?version=published', + expect.any(Object) + ) + }) + test('uses the CMR system token for a scheduled invocation', async () => { process.env.RDF_MIRROR_SOURCE_ENV = 'uat' vi.mocked(getCmrSystemToken).mockResolvedValue('scheduled-system-token') @@ -384,6 +400,13 @@ describe('mirrorRdf', () => { process.env.RDF_MIRROR_SOURCE_ENV = 'invalid' await expect(mirrorRdf({ source: 'aws.events' })) - .rejects.toThrow('RDF_MIRROR_SOURCE_ENV must be sit, uat, or prod') + .rejects.toThrow('RDF_MIRROR_SOURCE_ENV must be local, sit, uat, or prod') + }) + + test('rejects the local source outside SAM local', async () => { + process.env.RDF_MIRROR_SOURCE_ENV = 'local' + + await expect(mirrorRdf({ source: 'aws.events' })) + .rejects.toThrow('RDF_MIRROR_SOURCE_ENV local is only supported by SAM local') }) }) diff --git a/serverless/src/mirrorRdf/handler.js b/serverless/src/mirrorRdf/handler.js index fa5a6d85..9a5faf97 100644 --- a/serverless/src/mirrorRdf/handler.js +++ b/serverless/src/mirrorRdf/handler.js @@ -7,6 +7,7 @@ import { logger } from '@/shared/logger' const RDF_VERSIONS = ['published', 'draft'] const SOURCE_BASE_URLS = { + local: 'http://host.docker.internal:3013', sit: 'https://cmr.sit.earthdata.nasa.gov/kms', uat: 'https://cmr.uat.earthdata.nasa.gov/kms', prod: 'https://cmr.earthdata.nasa.gov/kms' @@ -38,7 +39,11 @@ const getSourceConfiguration = () => { const sourceBaseUrl = SOURCE_BASE_URLS[sourceEnvironment] if (!sourceBaseUrl) { - throw new Error('RDF_MIRROR_SOURCE_ENV must be sit, uat, or prod') + throw new Error('RDF_MIRROR_SOURCE_ENV must be local, sit, uat, or prod') + } + + if (sourceEnvironment === 'local' && process.env.AWS_SAM_LOCAL !== 'true') { + throw new Error('RDF_MIRROR_SOURCE_ENV local is only supported by SAM local') } return { From d5d0410489ddec6ea02ecacda17482f51206e4d2 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 13:18:05 -0400 Subject: [PATCH 08/12] KMS-648: Updated README to include how ot upgrades sams --- README.md | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/README.md b/README.md index 4e4cc52c..7e349899 100644 --- a/README.md +++ b/README.md @@ -52,6 +52,19 @@ To run local server with SAM watch mode enabled npm run start-local:watch ``` +#### Unsupported `nodejs24.x` runtime + +If local startup reports that `nodejs24.x` is unsupported, upgrade AWS SAM CLI. Local Lambdas are run by SAM, so rebuilding LocalStack will not resolve this error. + +```bash +sam --version +brew update +brew upgrade aws-sam-cli +sam --version +``` + +After upgrading, rerun `npm run start-local`. SAM will download the Node.js 24 Lambda runtime image when it is first needed. + ### Why local uses SAM and LocalStack Local development intentionally splits responsibilities between SAM and LocalStack: From 2a8862ac04ac2790063bdeefe809d8dcc901108d Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 15:35:25 -0400 Subject: [PATCH 09/12] KMS-648: Removed auth on retrieving rdf exports --- cdk/app/lib/helper/KmsLambdaFunctions.ts | 2 +- scripts/local/run_rdf_export_smoke.sh | 12 ++--- .../src/mirrorRdf/__tests__/handler.test.js | 45 +++---------------- serverless/src/mirrorRdf/handler.js | 29 +----------- 4 files changed, 10 insertions(+), 78 deletions(-) diff --git a/cdk/app/lib/helper/KmsLambdaFunctions.ts b/cdk/app/lib/helper/KmsLambdaFunctions.ts index bfec9fc5..8c1f6129 100644 --- a/cdk/app/lib/helper/KmsLambdaFunctions.ts +++ b/cdk/app/lib/helper/KmsLambdaFunctions.ts @@ -414,7 +414,7 @@ export class LambdaFunctions { 'exportRdf', '/rdf/export', 'POST', - true + false ) this.createApiLambda( diff --git a/scripts/local/run_rdf_export_smoke.sh b/scripts/local/run_rdf_export_smoke.sh index f304968b..10f82d88 100755 --- a/scripts/local/run_rdf_export_smoke.sh +++ b/scripts/local/run_rdf_export_smoke.sh @@ -5,18 +5,15 @@ set -euo pipefail # Downloads and validates the published and draft RDF gzip exports. # # Usage: -# KMS_AUTHORIZATION='' \ -# ./scripts/local/run_rdf_export_smoke.sh [outputDirectory] +# ./scripts/local/run_rdf_export_smoke.sh [outputDirectory] usage() { cat <<'EOF' Usage: - KMS_AUTHORIZATION='' \ - ./scripts/local/run_rdf_export_smoke.sh [outputDirectory] + ./scripts/local/run_rdf_export_smoke.sh [outputDirectory] Environment: - KMS_AUTHORIZATION Authorization header value passed through exactly as provided. - KMS_BASE_URL Optional override for the KMS base URL. + KMS_BASE_URL Optional override for the KMS base URL. The output directory defaults to /tmp/kms-rdf-export-smoke-. EOF @@ -34,8 +31,6 @@ if [[ -z "$ENVIRONMENT" ]]; then exit 1 fi -AUTHORIZATION_VALUE="${KMS_AUTHORIZATION:?Missing KMS_AUTHORIZATION environment variable.}" - case "$ENVIRONMENT" in sit) DEFAULT_BASE_URL="https://cmr.sit.earthdata.nasa.gov/kms" @@ -68,7 +63,6 @@ download_export() { --show-error \ --fail-with-body \ --request POST \ - --header "Authorization: ${AUTHORIZATION_VALUE}" \ --header 'Accept: application/json' \ "${BASE_URL}/rdf/export?version=${version}" \ --output "$response_file" diff --git a/serverless/src/mirrorRdf/__tests__/handler.test.js b/serverless/src/mirrorRdf/__tests__/handler.test.js index 6e6452ff..cd1bfabc 100644 --- a/serverless/src/mirrorRdf/__tests__/handler.test.js +++ b/serverless/src/mirrorRdf/__tests__/handler.test.js @@ -8,7 +8,6 @@ import { vi } from 'vitest' -import { getCmrSystemToken } from '@/shared/getCmrWriterToken' import { logger } from '@/shared/logger' import { mirrorRdf } from '../handler' @@ -21,10 +20,6 @@ vi.mock('@/shared/getConfig', () => ({ })) })) -vi.mock('@/shared/getCmrWriterToken', () => ({ - getCmrSystemToken: vi.fn() -})) - vi.mock('@/shared/logger') const createArchiveResponse = (rdfXml) => { @@ -69,8 +64,7 @@ const configureSuccessfulFetch = () => { } const createApiEvent = () => ({ - requestContext: {}, - headers: { Authorization: 'api-token' } + requestContext: {} }) describe('mirrorRdf', () => { @@ -79,7 +73,6 @@ describe('mirrorRdf', () => { vi.stubGlobal('fetch', vi.fn()) delete process.env.RDF_MIRROR_SOURCE_ENV delete process.env.AWS_SAM_LOCAL - vi.mocked(getCmrSystemToken).mockResolvedValue('system-token') process.env.RDF4J_SERVICE_URL = 'http://rdf4j.test:8080' process.env.RDF4J_REPOSITORY_ID = 'kms-test' process.env.RDF4J_USER_NAME = 'rdf-user' @@ -102,10 +95,7 @@ describe('mirrorRdf', () => { process.env.RDF_MIRROR_SOURCE_ENV = 'SIT' configureSuccessfulFetch() - const response = await mirrorRdf({ - requestContext: {}, - headers: { authorization: 'Bearer manual-token' } - }) + const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(200) expect(JSON.parse(response.body)).toEqual({ @@ -120,14 +110,11 @@ describe('mirrorRdf', () => { { method: 'POST', headers: { - Accept: 'application/json', - Authorization: 'Bearer manual-token' + Accept: 'application/json' } } ) - expect(getCmrSystemToken).not.toHaveBeenCalled() - expect(fetch).toHaveBeenNthCalledWith( 3, 'https://cmr.sit.earthdata.nasa.gov/kms/rdf/export?version=draft', @@ -171,9 +158,8 @@ describe('mirrorRdf', () => { ) }) - test('uses the CMR system token for a scheduled invocation', async () => { + test('downloads public source exports for a scheduled invocation', async () => { process.env.RDF_MIRROR_SOURCE_ENV = 'uat' - vi.mocked(getCmrSystemToken).mockResolvedValue('scheduled-system-token') configureSuccessfulFetch() await mirrorRdf({ source: 'aws.events' }) @@ -183,8 +169,7 @@ describe('mirrorRdf', () => { 'https://cmr.uat.earthdata.nasa.gov/kms/rdf/export?version=published', expect.objectContaining({ headers: { - Accept: 'application/json', - Authorization: 'scheduled-system-token' + Accept: 'application/json' } }) ) @@ -376,26 +361,6 @@ describe('mirrorRdf', () => { }) }) - test('does not call the source when an API request has no Authorization header', async () => { - process.env.RDF_MIRROR_SOURCE_ENV = 'prod' - - const response = await mirrorRdf({ requestContext: {} }) - - expect(response.statusCode).toBe(500) - expect(fetch).not.toHaveBeenCalled() - expect(getCmrSystemToken).not.toHaveBeenCalled() - }) - - test('throws when a scheduled invocation has no CMR system token', async () => { - process.env.RDF_MIRROR_SOURCE_ENV = 'prod' - vi.mocked(getCmrSystemToken).mockResolvedValue(undefined) - - await expect(mirrorRdf({ source: 'aws.events' })) - .rejects.toThrow('CMR system token is unavailable') - - expect(fetch).not.toHaveBeenCalled() - }) - test('throws a scheduled invocation error so EventBridge can report the failure', async () => { process.env.RDF_MIRROR_SOURCE_ENV = 'invalid' diff --git a/serverless/src/mirrorRdf/handler.js b/serverless/src/mirrorRdf/handler.js index 9a5faf97..29b07b1d 100644 --- a/serverless/src/mirrorRdf/handler.js +++ b/serverless/src/mirrorRdf/handler.js @@ -1,7 +1,6 @@ import { promisify } from 'util' import zlib from 'zlib' -import { getCmrSystemToken } from '@/shared/getCmrWriterToken' import { getApplicationConfig } from '@/shared/getConfig' import { logger } from '@/shared/logger' @@ -14,15 +13,6 @@ const SOURCE_BASE_URLS = { } const gunzip = promisify(zlib.gunzip) -/** - * Reads the Authorization header from an API Gateway event without modifying its value. - * - * @param {object} event API Gateway invocation event. - * @returns {string|undefined} Incoming authorization value. - */ -const getRequestAuthorization = (event) => Object.entries(event.headers || {}) - .find(([headerName]) => headerName.toLowerCase() === 'authorization')?.[1] - /** * Resolves the configured source environment to its KMS API base URL. * @@ -56,23 +46,17 @@ const getSourceConfiguration = () => { * Requests and downloads one gzip-compressed RDF graph from the source KMS environment. * * @param {object} params Download parameters. - * @param {string} params.authorization Authorization value forwarded to the source KMS API. * @param {string} params.sourceBaseUrl Source KMS API base URL. * @param {'published'|'draft'} params.version RDF graph version. * @returns {Promise} Uncompressed RDF/XML content. */ const downloadRdf = async ({ - authorization, sourceBaseUrl, version }) => { - const headers = { - Accept: 'application/json', - Authorization: authorization - } const exportResponse = await fetch(`${sourceBaseUrl}/rdf/export?version=${version}`, { method: 'POST', - headers + headers: { Accept: 'application/json' } }) if (!exportResponse.ok) { @@ -178,22 +162,11 @@ export const mirrorRdf = async (event = {}) => { } } - const authorization = isApiRequest - ? getRequestAuthorization(event) - : await getCmrSystemToken() - - if (!authorization) { - throw new Error(isApiRequest - ? 'Authorization header is unavailable' - : 'CMR system token is unavailable') - } - // Complete every source download before modifying either destination graph. const rdfByVersion = await RDF_VERSIONS.reduce( (previousDownloads, version) => previousDownloads.then(async (downloads) => ({ ...downloads, [version]: await downloadRdf({ - authorization, sourceBaseUrl: sourceConfiguration.sourceBaseUrl, version }) From 7ff91f1b125d155d6fde5a60adb3c131e5857a04 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Fri, 14 Aug 2026 16:05:06 -0400 Subject: [PATCH 10/12] KMS-648: PR Updates --- cdk/app/lib/helper/IamSetup.ts | 3 +- .../src/mirrorRdf/__tests__/handler.test.js | 136 +++++++----------- serverless/src/mirrorRdf/handler.js | 68 +++++---- .../shared/__tests__/sparqlRequest.test.js | 13 ++ serverless/src/shared/sparqlRequest.js | 3 +- 5 files changed, 109 insertions(+), 114 deletions(-) diff --git a/cdk/app/lib/helper/IamSetup.ts b/cdk/app/lib/helper/IamSetup.ts index 0ae5f1a2..9525e904 100644 --- a/cdk/app/lib/helper/IamSetup.ts +++ b/cdk/app/lib/helper/IamSetup.ts @@ -103,8 +103,7 @@ export class IamSetup { 's3:ListBucket', 's3:PutLifecycleConfiguration', 's3:GetBucketLocation', - 's3:ListAllMyBuckets', - 's3:HeadBucket' + 's3:ListAllMyBuckets' ], resources: [ `arn:aws:s3:::kms-rdf-backup-${stage}`, diff --git a/serverless/src/mirrorRdf/__tests__/handler.test.js b/serverless/src/mirrorRdf/__tests__/handler.test.js index cd1bfabc..42c937d3 100644 --- a/serverless/src/mirrorRdf/__tests__/handler.test.js +++ b/serverless/src/mirrorRdf/__tests__/handler.test.js @@ -9,6 +9,12 @@ import { } from 'vitest' import { logger } from '@/shared/logger' +import { sparqlRequest } from '@/shared/sparqlRequest' +import { + commitTransaction, + rollbackTransaction, + startTransaction +} from '@/shared/transactionHelpers' import { mirrorRdf } from '../handler' @@ -21,6 +27,8 @@ vi.mock('@/shared/getConfig', () => ({ })) vi.mock('@/shared/logger') +vi.mock('@/shared/sparqlRequest') +vi.mock('@/shared/transactionHelpers') const createArchiveResponse = (rdfXml) => { const archive = gzipSync(rdfXml) @@ -57,10 +65,6 @@ const configureSuccessfulFetch = () => { json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) }) .mockResolvedValueOnce(createArchiveResponse('draft')) - .mockResolvedValueOnce(createTextResponse()) - .mockResolvedValueOnce(createTextResponse()) - .mockResolvedValueOnce(createTextResponse()) - .mockResolvedValueOnce(createTextResponse()) } const createApiEvent = () => ({ @@ -77,6 +81,10 @@ describe('mirrorRdf', () => { process.env.RDF4J_REPOSITORY_ID = 'kms-test' process.env.RDF4J_USER_NAME = 'rdf-user' process.env.RDF4J_PASSWORD = 'rdf-password' + vi.mocked(startTransaction).mockResolvedValue('transaction-url') + vi.mocked(sparqlRequest).mockResolvedValue(createTextResponse()) + vi.mocked(commitTransaction).mockResolvedValue() + vi.mocked(rollbackTransaction).mockResolvedValue() }) test('skips the mirror when no source environment is configured', async () => { @@ -121,21 +129,30 @@ describe('mirrorRdf', () => { expect.any(Object) ) - const publishedContext = encodeURIComponent( - '' - ) - expect(String(vi.mocked(fetch).mock.calls[4][0])).toContain(`context=${publishedContext}`) - expect(vi.mocked(fetch).mock.calls[4][1]).toEqual({ - method: 'DELETE', - headers: { - Authorization: `Basic ${Buffer.from('rdf-user:rdf-password').toString('base64')}` + expect(startTransaction).toHaveBeenCalledTimes(2) + expect(sparqlRequest).toHaveBeenNthCalledWith(1, { + method: 'PUT', + body: 'CLEAR GRAPH ', + contentType: 'application/sparql-update', + transaction: { + transactionUrl: 'transaction-url', + action: 'UPDATE' } }) - expect(vi.mocked(fetch).mock.calls[5][1]).toEqual(expect.objectContaining({ - method: 'POST', - body: 'published' - })) + expect(sparqlRequest).toHaveBeenNthCalledWith(2, { + method: 'PUT', + body: 'published', + contentType: 'application/rdf+xml', + version: 'published', + transaction: { + transactionUrl: 'transaction-url', + action: 'ADD' + } + }) + + expect(commitTransaction).toHaveBeenCalledTimes(2) + expect(rollbackTransaction).not.toHaveBeenCalled() expect(logger.info).toHaveBeenCalledWith('[rdf-mirror] Mirrored RDF graphs', { sourceEnvironment: 'sit', @@ -268,97 +285,52 @@ describe('mirrorRdf', () => { test('returns an API error when clearing a destination graph fails', async () => { process.env.RDF_MIRROR_SOURCE_ENV = 'prod' - vi.mocked(fetch) - .mockResolvedValueOnce({ - ok: true, - json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) - }) - .mockResolvedValueOnce(createArchiveResponse('published')) - .mockResolvedValueOnce({ - ok: true, - json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) - }) - .mockResolvedValueOnce(createArchiveResponse('draft')) - .mockResolvedValueOnce(createTextResponse({ - ok: false, - status: 500, - text: 'clear failed' - })) + configureSuccessfulFetch() + vi.mocked(sparqlRequest).mockRejectedValueOnce(new Error('clear failed')) const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) + expect(rollbackTransaction).toHaveBeenCalledWith('transaction-url') + expect(commitTransaction).not.toHaveBeenCalled() expect(logger.error).toHaveBeenCalledWith( - '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Failed to clear destination published graph: 500 clear failed' + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Failed to replace destination published graph: clear failed' ) }) test('returns an API error when importing a destination graph fails', async () => { process.env.RDF_MIRROR_SOURCE_ENV = 'prod' - vi.mocked(fetch) - .mockResolvedValueOnce({ - ok: true, - json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) - }) - .mockResolvedValueOnce(createArchiveResponse('published')) - .mockResolvedValueOnce({ - ok: true, - json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) - }) - .mockResolvedValueOnce(createArchiveResponse('draft')) + configureSuccessfulFetch() + vi.mocked(sparqlRequest) .mockResolvedValueOnce(createTextResponse()) - .mockResolvedValueOnce(createTextResponse({ - ok: false, - status: 500, - text: 'import failed' - })) + .mockRejectedValueOnce(new Error('import failed')) const response = await mirrorRdf(createApiEvent()) expect(response.statusCode).toBe(500) + expect(rollbackTransaction).toHaveBeenCalledWith('transaction-url') + expect(commitTransaction).not.toHaveBeenCalled() expect(logger.error).toHaveBeenCalledWith( - '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Failed to import destination published graph: 500 import failed' + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Failed to replace destination published graph: import failed' ) }) - test('uses local RDF4J defaults and permits an absent destination context', async () => { + test('reports the original replacement error when rollback also fails', async () => { process.env.RDF_MIRROR_SOURCE_ENV = 'prod' - delete process.env.RDF4J_SERVICE_URL - delete process.env.RDF4J_REPOSITORY_ID - delete process.env.RDF4J_USER_NAME - delete process.env.RDF4J_PASSWORD - vi.mocked(fetch) - .mockResolvedValueOnce({ - ok: true, - json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/published.rdf.xml.gz' }) - }) - .mockResolvedValueOnce(createArchiveResponse('published')) - .mockResolvedValueOnce({ - ok: true, - json: vi.fn().mockResolvedValue({ downloadUrl: 'https://download.test/draft.rdf.xml.gz' }) - }) - .mockResolvedValueOnce(createArchiveResponse('draft')) - .mockResolvedValueOnce(createTextResponse({ - ok: false, - status: 404 - })) - .mockResolvedValueOnce(createTextResponse()) - .mockResolvedValueOnce(createTextResponse()) - .mockResolvedValueOnce(createTextResponse()) + configureSuccessfulFetch() + vi.mocked(sparqlRequest).mockRejectedValueOnce(new Error('clear failed')) + vi.mocked(rollbackTransaction).mockRejectedValueOnce(new Error('rollback failed')) const response = await mirrorRdf(createApiEvent()) - expect(response.statusCode).toBe(200) - expect(String(vi.mocked(fetch).mock.calls[4][0])).toContain( - 'http://localhost:8081/rdf4j-server/repositories/kms/statements' + expect(response.statusCode).toBe(500) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to roll back published graph replacement, error=Error: rollback failed' ) - expect(vi.mocked(fetch).mock.calls[4][1]).toEqual({ - method: 'DELETE', - headers: { - Authorization: `Basic ${Buffer.from('rdf4j:rdf4j').toString('base64')}` - } - }) + expect(logger.error).toHaveBeenCalledWith( + '[rdf-mirror] Failed to mirror RDF graphs, error=Error: Failed to replace destination published graph: clear failed' + ) }) test('throws a scheduled invocation error so EventBridge can report the failure', async () => { diff --git a/serverless/src/mirrorRdf/handler.js b/serverless/src/mirrorRdf/handler.js index 29b07b1d..fc513f8b 100644 --- a/serverless/src/mirrorRdf/handler.js +++ b/serverless/src/mirrorRdf/handler.js @@ -3,6 +3,12 @@ import zlib from 'zlib' import { getApplicationConfig } from '@/shared/getConfig' import { logger } from '@/shared/logger' +import { sparqlRequest } from '@/shared/sparqlRequest' +import { + commitTransaction, + rollbackTransaction, + startTransaction +} from '@/shared/transactionHelpers' const RDF_VERSIONS = ['published', 'draft'] const SOURCE_BASE_URLS = { @@ -100,40 +106,44 @@ const downloadRdf = async ({ * @returns {Promise} */ const replaceDestinationGraph = async ({ rdfXml, version }) => { - const serviceUrl = String(process.env.RDF4J_SERVICE_URL || 'http://localhost:8081') - .replace(/\/$/, '') - const repositoryId = process.env.RDF4J_REPOSITORY_ID || 'kms' - const statementsUrl = new URL(`${serviceUrl}/rdf4j-server/repositories/${repositoryId}/statements`) const graphUri = `https://gcmd.earthdata.nasa.gov/kms/version/${version}` - const credentials = Buffer.from( - `${process.env.RDF4J_USER_NAME || 'rdf4j'}:${process.env.RDF4J_PASSWORD || 'rdf4j'}` - ).toString('base64') - const authorization = `Basic ${credentials}` + let transactionUrl - statementsUrl.searchParams.set('context', `<${graphUri}>`) - - const clearResponse = await fetch(statementsUrl, { - method: 'DELETE', - headers: { Authorization: authorization } - }) + try { + transactionUrl = await startTransaction() + + await sparqlRequest({ + method: 'PUT', + body: `CLEAR GRAPH <${graphUri}>`, + contentType: 'application/sparql-update', + transaction: { + transactionUrl, + action: 'UPDATE' + } + }) - if (!clearResponse.ok && clearResponse.status !== 404) { - const responseText = await clearResponse.text() - throw new Error(`Failed to clear destination ${version} graph: ${clearResponse.status} ${responseText}`) - } + await sparqlRequest({ + method: 'PUT', + body: rdfXml, + contentType: 'application/rdf+xml', + version, + transaction: { + transactionUrl, + action: 'ADD' + } + }) - const importResponse = await fetch(statementsUrl, { - method: 'POST', - headers: { - Authorization: authorization, - 'Content-Type': 'application/rdf+xml' - }, - body: rdfXml - }) + await commitTransaction(transactionUrl) + } catch (error) { + if (transactionUrl) { + try { + await rollbackTransaction(transactionUrl) + } catch (rollbackError) { + logger.error(`[rdf-mirror] Failed to roll back ${version} graph replacement, error=${rollbackError.toString()}`) + } + } - if (!importResponse.ok) { - const responseText = await importResponse.text() - throw new Error(`Failed to import destination ${version} graph: ${importResponse.status} ${responseText}`) + throw new Error(`Failed to replace destination ${version} graph: ${error.message}`) } } diff --git a/serverless/src/shared/__tests__/sparqlRequest.test.js b/serverless/src/shared/__tests__/sparqlRequest.test.js index 9c1780cf..d6f9c601 100644 --- a/serverless/src/shared/__tests__/sparqlRequest.test.js +++ b/serverless/src/shared/__tests__/sparqlRequest.test.js @@ -32,6 +32,7 @@ describe('sparqlRequest', () => { afterEach(() => { delete process.env.RDF4J_SERVICE_URL + delete process.env.RDF4J_REPOSITORY_ID delete process.env.RDF4J_USER_NAME delete process.env.RDF4J_PASSWORD }) @@ -249,6 +250,18 @@ describe('sparqlRequest', () => { ) }) + test('should use the configured RDF4J repository', async () => { + process.env.RDF4J_REPOSITORY_ID = 'alternate-repository' + global.fetch.mockResolvedValue({ ok: true }) + + await sparqlRequest({ method: 'GET' }) + + expect(global.fetch).toHaveBeenCalledWith( + 'http://test-server.com/rdf4j-server/repositories/alternate-repository', + expect.any(Object) + ) + }) + test('should append context for statements path when version is set for non-query/update content type', async () => { const mockResponse = { ok: true, diff --git a/serverless/src/shared/sparqlRequest.js b/serverless/src/shared/sparqlRequest.js index 8c872b6e..b3382602 100644 --- a/serverless/src/shared/sparqlRequest.js +++ b/serverless/src/shared/sparqlRequest.js @@ -116,8 +116,9 @@ export const sparqlRequest = async (props) => { */ const getSparqlEndpoint = () => { const baseUrl = process.env.RDF4J_SERVICE_URL || 'http://localhost:8080' + const repositoryId = process.env.RDF4J_REPOSITORY_ID || 'kms' - return `${baseUrl}/rdf4j-server/repositories/kms` + return `${baseUrl}/rdf4j-server/repositories/${repositoryId}` } /** From 2cee155ee095411e18492934f9ffe348a187f266 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Mon, 17 Aug 2026 16:39:51 -0400 Subject: [PATCH 11/12] KMS-648: Added auth header for keyword uuid search. --- cdk/app/lib/CmrEventProcessingStack.ts | 1 + .../helper/CmrKeywordEventsListenerSetup.ts | 19 ++++++++ .../__tests__/handler.test.js | 8 ++++ .../src/cmrKeywordEventsListener/handler.js | 2 + .../src/publisher/__tests__/handler.test.js | 11 +++++ serverless/src/publisher/handler.js | 19 ++++++-- .../getCmrCollectionConceptIds.test.js | 45 +++++++++++++++++-- .../src/shared/getCmrCollectionConceptIds.js | 20 ++++++++- 8 files changed, 117 insertions(+), 8 deletions(-) diff --git a/cdk/app/lib/CmrEventProcessingStack.ts b/cdk/app/lib/CmrEventProcessingStack.ts index 18fc92e0..7e12a0cf 100644 --- a/cdk/app/lib/CmrEventProcessingStack.ts +++ b/cdk/app/lib/CmrEventProcessingStack.ts @@ -88,6 +88,7 @@ export class CmrEventProcessingStack extends cdk.Stack { const listenerSetup = new CmrKeywordEventsListenerSetup(this, 'CmrKeywordEventsListener', { cmrBaseUrl: props.cmrBaseUrl, + cmrSystemTokenParameterName: props.cmrSystemTokenParameterName, prefix: props.prefix, stage: props.stage, keywordEventsTopic: topic, diff --git a/cdk/app/lib/helper/CmrKeywordEventsListenerSetup.ts b/cdk/app/lib/helper/CmrKeywordEventsListenerSetup.ts index 7c05a4ce..bf5ca2a0 100644 --- a/cdk/app/lib/helper/CmrKeywordEventsListenerSetup.ts +++ b/cdk/app/lib/helper/CmrKeywordEventsListenerSetup.ts @@ -2,6 +2,7 @@ import * as path from 'path' import * as cdk from 'aws-cdk-lib' import * as ec2 from 'aws-cdk-lib/aws-ec2' +import * as iam from 'aws-cdk-lib/aws-iam' import * as eventsources from 'aws-cdk-lib/aws-lambda-event-sources' import { NodejsFunction } from 'aws-cdk-lib/aws-lambda-nodejs' import * as sns from 'aws-cdk-lib/aws-sns' @@ -16,6 +17,7 @@ import { NODE_LAMBDA_RUNTIME } from './NodeLambdaRuntime' */ interface CmrKeywordEventsListenerSetupProps { cmrBaseUrl: string + cmrSystemTokenParameterName?: string prefix: string securityGroup: ec2.SecurityGroup stage: string @@ -45,6 +47,7 @@ export class CmrKeywordEventsListenerSetup extends Construct { const { cmrBaseUrl, + cmrSystemTokenParameterName, keywordEventsTopic, metadataCorrectionRequestsTopic, prefix, @@ -72,6 +75,9 @@ export class CmrKeywordEventsListenerSetup extends Construct { memorySize: 1024, environment: { CMR_BASE_URL: cmrBaseUrl, + ...(cmrSystemTokenParameterName + ? { CMR_SYSTEM_TOKEN_PARAMETER_NAME: cmrSystemTokenParameterName } + : {}), METADATA_CORRECTION_REQUESTS_TOPIC_ARN: metadataCorrectionRequestsTopic.topicArn }, depsLockFilePath: path.join(projectRoot, 'package-lock.json'), @@ -92,6 +98,19 @@ export class CmrKeywordEventsListenerSetup extends Construct { this.queue.grantConsumeMessages(this.listenerLambda) metadataCorrectionRequestsTopic.grantPublish(this.listenerLambda) + if (cmrSystemTokenParameterName) { + const systemTokenParameterArn = cdk.Stack.of(this).formatArn({ + service: 'ssm', + resource: 'parameter', + resourceName: cmrSystemTokenParameterName.replace(/^\//, '') + }) + + this.listenerLambda.addToRolePolicy(new iam.PolicyStatement({ + actions: ['ssm:GetParameter'], + resources: [systemTokenParameterArn] + })) + } + this.queueUrlOutput = new cdk.CfnOutput(this, 'CmrKeywordEventsQueueUrl', { description: 'Queue URL for CMR keyword event processing', exportName: `${prefix}-CmrKeywordEventsQueueUrl`, diff --git a/serverless/src/cmrKeywordEventsListener/__tests__/handler.test.js b/serverless/src/cmrKeywordEventsListener/__tests__/handler.test.js index a8f1e7fc..79f23bd5 100644 --- a/serverless/src/cmrKeywordEventsListener/__tests__/handler.test.js +++ b/serverless/src/cmrKeywordEventsListener/__tests__/handler.test.js @@ -169,6 +169,14 @@ describe('when the CMR keyword events processor is invoked', () => { expect.stringContaining('collectionConceptId=C1000000000-PROV') ) + expect(logger.info).toHaveBeenCalledWith( + expect.stringContaining('uuid=1234') + ) + + expect(logger.info).toHaveBeenCalledWith( + expect.stringContaining('eventTimestamp=2026-04-21T00:00:00.000Z') + ) + expect(logger.info).toHaveBeenCalledWith( expect.stringContaining('messageId=metadata-correction-message-123') ) diff --git a/serverless/src/cmrKeywordEventsListener/handler.js b/serverless/src/cmrKeywordEventsListener/handler.js index a915720b..4772479d 100644 --- a/serverless/src/cmrKeywordEventsListener/handler.js +++ b/serverless/src/cmrKeywordEventsListener/handler.js @@ -71,6 +71,8 @@ const publishCollectionCorrectionRequests = async (collectionConceptIds, keyword logger.info( '[consumer] Published metadata correction request ' + `collectionConceptId=${metadataCorrectionRequest.collectionConceptId} ` + + `uuid=${metadataCorrectionRequest.keywordEvent.uuid || 'n/a'} ` + + `eventTimestamp=${metadataCorrectionRequest.keywordEvent.timestamp || 'n/a'} ` + `messageId=${publishResult.messageId || 'n/a'} ` + `topicArn=${publishResult.topicArn || 'n/a'}` ) diff --git a/serverless/src/publisher/__tests__/handler.test.js b/serverless/src/publisher/__tests__/handler.test.js index f65dc63b..6dda78d0 100644 --- a/serverless/src/publisher/__tests__/handler.test.js +++ b/serverless/src/publisher/__tests__/handler.test.js @@ -217,6 +217,17 @@ describe('publisher handler', () => { NewKeywordObject: SCIENCE_PATH_KEYWORD })) + expect(logger.info).toHaveBeenCalledWith('[publisher] Published keyword event', { + versionName: 'v1.0.0', + messageId: 'message-1', + eventType: 'INSERTED', + scheme: 'sciencekeywords', + uuid: 'uuid1', + timestamp: '2023-06-01T00:00:00.000Z', + oldKeywordObject: undefined, + newKeywordObject: SCIENCE_PATH_KEYWORD + }) + expect(eventBridgeMock.commandCalls(PutEventsCommand).length).toBe(1) expect(result).toEqual({ diff --git a/serverless/src/publisher/handler.js b/serverless/src/publisher/handler.js index 9cabaf3b..efc785dc 100644 --- a/serverless/src/publisher/handler.js +++ b/serverless/src/publisher/handler.js @@ -120,13 +120,14 @@ const emitCachingEvent = async ({ versionName, publishDate, keywordEvents }) => * * @async * @param {Array} keywordEvents - Keyword event payloads to publish. + * @param {string} versionName - Unique KMS version being published. * @returns {Promise<{ * attemptedCount: number, * publishedCount: number, * failedEvents: Array<{ keywordEvent: Object, error: string, attempts: number }> * }>} SNS publish summary for the completed batch. */ -const publishKeywordEvents = async (keywordEvents) => keywordEvents.reduce( +const publishKeywordEvents = async (keywordEvents, versionName) => keywordEvents.reduce( async (summaryPromise, keywordEvent) => { const summary = await summaryPromise let lastError @@ -145,7 +146,19 @@ const publishKeywordEvents = async (keywordEvents) => keywordEvents.reduce( } // eslint-disable-next-line no-await-in-loop - await publishKeywordEvent(keywordEvent) + const publishResult = await publishKeywordEvent(keywordEvent) + + logger.info('[publisher] Published keyword event', { + versionName, + messageId: publishResult.messageId || 'n/a', + eventType: keywordEvent.EventType, + scheme: keywordEvent.Scheme, + uuid: keywordEvent.UUID, + timestamp: keywordEvent.Timestamp, + oldKeywordObject: keywordEvent.OldKeywordObject, + newKeywordObject: keywordEvent.NewKeywordObject + }) + summary.publishedCount += 1 lastError = undefined break @@ -326,7 +339,7 @@ export const publisher = async (event) => { postPublishFailures.push(skipMessage) logger.warn(`[publisher] ${skipMessage}`) } else { - const publishSummary = await publishKeywordEvents(keywordEvents) + const publishSummary = await publishKeywordEvents(keywordEvents, versionName) keywordEventsPublished = publishSummary.publishedCount keywordEventPublishFailures = publishSummary.failedEvents.length diff --git a/serverless/src/shared/__tests__/getCmrCollectionConceptIds.test.js b/serverless/src/shared/__tests__/getCmrCollectionConceptIds.test.js index 10fbd510..7a315f76 100644 --- a/serverless/src/shared/__tests__/getCmrCollectionConceptIds.test.js +++ b/serverless/src/shared/__tests__/getCmrCollectionConceptIds.test.js @@ -8,9 +8,14 @@ import { import { cmrGetRequest } from '../cmrGetRequest' import { getCmrCollectionConceptIds } from '../getCmrCollectionConceptIds' +import { getCmrSystemToken } from '../getCmrWriterToken' import { logger } from '../logger' vi.mock('../cmrGetRequest') +vi.mock('../getCmrWriterToken', () => ({ + getCmrSystemToken: vi.fn() +})) + vi.mock('../logger', () => ({ logger: { info: vi.fn(), @@ -34,6 +39,7 @@ describe('getCmrCollectionConceptIds', () => { beforeEach(() => { vi.clearAllMocks() + vi.mocked(getCmrSystemToken).mockResolvedValue('Bearer system-token') }) const mockHeaders = (cmrHits) => ({ @@ -72,8 +78,19 @@ describe('getCmrCollectionConceptIds', () => { ]) expect(cmrGetRequest).toHaveBeenCalledWith({ - path: '/search/collections.umm_json?keyword=1234-5678-9ABC-DEF0&page_size=2000&page_num=1' + path: '/search/collections.umm_json?keyword=1234-5678-9ABC-DEF0&page_size=2000&page_num=1', + headers: { + Authorization: 'Bearer system-token' + } }) + + expect(logger.info).toHaveBeenCalledWith( + '[cmr-lookup] Sending authenticated CMR collection UUID search ' + + 'scheme=sciencekeywords ' + + 'uuid=1234-5678-9ABC-DEF0 ' + + 'pageNumber=1 ' + + 'authorizationPresent=true' + ) }) test('should return unique collection concept ids for any keyword scheme because lookup is UUID driven', async () => { @@ -102,7 +119,10 @@ describe('getCmrCollectionConceptIds', () => { ]) expect(cmrGetRequest).toHaveBeenCalledWith({ - path: '/search/collections.umm_json?keyword=DATA-CENTER-UUID&page_size=2000&page_num=1' + path: '/search/collections.umm_json?keyword=DATA-CENTER-UUID&page_size=2000&page_num=1', + headers: { + Authorization: 'Bearer system-token' + } }) }) @@ -137,11 +157,17 @@ describe('getCmrCollectionConceptIds', () => { }) expect(cmrGetRequest).toHaveBeenNthCalledWith(1, expect.objectContaining({ - path: '/search/collections.umm_json?keyword=1234-5678-9ABC-DEF0&page_size=2000&page_num=1' + path: '/search/collections.umm_json?keyword=1234-5678-9ABC-DEF0&page_size=2000&page_num=1', + headers: { + Authorization: 'Bearer system-token' + } })) expect(cmrGetRequest).toHaveBeenNthCalledWith(2, expect.objectContaining({ - path: '/search/collections.umm_json?keyword=1234-5678-9ABC-DEF0&page_size=2000&page_num=2' + path: '/search/collections.umm_json?keyword=1234-5678-9ABC-DEF0&page_size=2000&page_num=2', + headers: { + Authorization: 'Bearer system-token' + } })) expect(result).toEqual([ @@ -157,6 +183,17 @@ describe('getCmrCollectionConceptIds', () => { })).rejects.toThrow('Missing keyword UUID for CMR concept-id lookup') }) + test('should reject before requesting CMR when the system token is unavailable', async () => { + vi.mocked(getCmrSystemToken).mockResolvedValue(undefined) + + await expect(getCmrCollectionConceptIds({ + scheme: 'sciencekeywords', + uuid: '1234' + })).rejects.toThrow('Missing CMR system token for CMR concept-id lookup') + + expect(cmrGetRequest).not.toHaveBeenCalled() + }) + test('should throw when the CMR response is not ok', async () => { cmrGetRequest.mockResolvedValue({ ok: false, diff --git a/serverless/src/shared/getCmrCollectionConceptIds.js b/serverless/src/shared/getCmrCollectionConceptIds.js index 5de0e4ca..9038aa5d 100644 --- a/serverless/src/shared/getCmrCollectionConceptIds.js +++ b/serverless/src/shared/getCmrCollectionConceptIds.js @@ -1,5 +1,6 @@ import { cmrGetRequest } from './cmrGetRequest' import { formatKeywordObjectForLog } from './formatKeywordObjectForLog' +import { getCmrSystemToken } from './getCmrWriterToken' import { logger } from './logger' /** @@ -77,10 +78,27 @@ export const getCmrCollectionConceptIds = async ({ keywordObject }) + const authorizationToken = await getCmrSystemToken() + + if (!authorizationToken) { + throw new Error('Missing CMR system token for CMR concept-id lookup') + } + // Request one page of concept ids and keep the raw response for pagination headers. const requestPage = async (pageNumber) => { + logger.info( + '[cmr-lookup] Sending authenticated CMR collection UUID search ' + + `scheme=${scheme} ` + + `uuid=${uuid} ` + + `pageNumber=${pageNumber} ` + + `authorizationPresent=${Boolean(authorizationToken)}` + ) + const response = await cmrGetRequest({ - path: CMR_COLLECTION_SEARCH_PATH(uuid, pageNumber) + path: CMR_COLLECTION_SEARCH_PATH(uuid, pageNumber), + headers: { + Authorization: authorizationToken + } }) if (!response.ok) { From 7b6e395d569b772a02fe8c3f5f838ff7d9f75771 Mon Sep 17 00:00:00 2001 From: "Christopher D. Gokey" Date: Thu, 20 Aug 2026 21:07:14 -0400 Subject: [PATCH 12/12] KMS-648: Updated based on PR feedback. --- bin/rdf4j/start.sh | 1 + serverless/src/mirrorRdf/handler.js | 1 + 2 files changed, 2 insertions(+) diff --git a/bin/rdf4j/start.sh b/bin/rdf4j/start.sh index 22fc80a1..27ad13d7 100755 --- a/bin/rdf4j/start.sh +++ b/bin/rdf4j/start.sh @@ -27,4 +27,5 @@ docker run \ -e "RDF4J_USER_NAME=${RDF4J_USER_NAME}" \ -e "RDF4J_PASSWORD=${RDF4J_PASSWORD}" \ -e "RDF4J_CONTAINER_MEMORY_LIMIT=${RDF4J_CONTAINER_MEMORY_LIMIT}" \ + -v logs:/usr/local/tomcat/logs \ rdf4j:latest diff --git a/serverless/src/mirrorRdf/handler.js b/serverless/src/mirrorRdf/handler.js index fc513f8b..36063401 100644 --- a/serverless/src/mirrorRdf/handler.js +++ b/serverless/src/mirrorRdf/handler.js @@ -156,6 +156,7 @@ const replaceDestinationGraph = async ({ rdfXml, version }) => { */ export const mirrorRdf = async (event = {}) => { const { defaultResponseHeaders } = getApplicationConfig() + // Scheduled failures must be thrown for Lambda retry/alarm handling; API calls receive JSON. const isApiRequest = Boolean(event.requestContext) try {