From b66672c8c665bd69cdedda9851386fe537308c7e Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Thu, 20 Aug 2026 09:23:26 -0700 Subject: [PATCH 1/8] sql table for cnm status --- .../daac_archiver/sql_mws/__init__.py | 0 .../sql_mws/catalia_status_db.py | 150 ++++++++++++++++++ requirements.txt | 1 + setup.py | 3 +- 4 files changed, 153 insertions(+), 1 deletion(-) create mode 100644 cumulus_lambda_functions/daac_archiver/sql_mws/__init__.py create mode 100644 cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py diff --git a/cumulus_lambda_functions/daac_archiver/sql_mws/__init__.py b/cumulus_lambda_functions/daac_archiver/sql_mws/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py new file mode 100644 index 00000000..eb3f1039 --- /dev/null +++ b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py @@ -0,0 +1,150 @@ +""" +SQL (SQLAlchemy) based re-implementation of +cumulus_lambda_functions/daac_archiver/ddb_mws/catalia_status_db.py + +DynamoDB is difficult to run ad-hoc/small analytics queries against, so this module +moves the same "status per archiving identifier" data into a relational table instead. + +The public API (class-level column-name constants + `get`/`add` method signatures) is +kept identical to the DynamoDB version so existing callers (daac_receiver.py, +services/status_update_svc.py, catalya_uds_api/granules_archive_api.py) can swap the +import without any other code changes. + +-------------------------------------------------------------------------------------- +Suggested DDL (PostgreSQL dialect) for the table this class reads/writes. Replace +``uds_ctla_daac_status`` below with whatever ``table_name`` is passed into +``CataliaStatusDb()``. + +The DynamoDB table used ``identifier`` as its partition/primary key and ``datetime`` +as its sort/secondary key. A relational table needs a single-column primary key, so a +surrogate auto-increment ``id`` column is used instead, and the original +``(identifier, datetime)`` pairing is preserved as a UNIQUE constraint -- this keeps +the same "one row per status update, no silent overwrites" behavior that the DynamoDB +version got from `replace=False` + a ConditionExpression. + + CREATE TABLE uds_ctla_daac_status ( + id SERIAL PRIMARY KEY, + identifier VARCHAR(255) NOT NULL, + datetime VARCHAR(64) NOT NULL, + collection VARCHAR(255), + target_collection VARCHAR(255), + name VARCHAR(255), + status VARCHAR(64), + "errorCode" VARCHAR(255), + "errorMessage" VARCHAR(2048), + href VARCHAR(2048), + CONSTRAINT uq_uds_ctla_daac_status_identifier_datetime UNIQUE (identifier, datetime) + ); + CREATE INDEX ix_uds_ctla_daac_status_identifier ON uds_ctla_daac_status (identifier); + +Note: `errorCode`/`errorMessage` are quoted above because they're camelCase -- +Postgres folds unquoted identifiers to lowercase. The SQLAlchemy `Table` defined in +this module manages that quoting automatically. +-------------------------------------------------------------------------------------- +""" +import logging +import os +from typing import Optional + +from sqlalchemy import ( + Column, + Integer, + MetaData, + String, + Table, + UniqueConstraint, + create_engine, + insert, + select, +) +from sqlalchemy.engine import Engine +from sqlalchemy.exc import IntegrityError + +logger = logging.getLogger(__name__) + + +class CataliaStatusDb: + identifier = 'identifier' + collection = 'collection' + target_collection = 'target_collection' + name_str = 'name' + status = 'status' + error_code = 'errorCode' + error_message = 'errorMessage' + href_str = 'href' + datetime_str = 'datetime' + + def __init__(self, table_name: str, db_url: Optional[str] = None): + """ + :param table_name: name of the SQL table to read/write, analogous to the DDB table name. + :param db_url: SQLAlchemy connection URL. Falls back to the CATALYA_SQL_DB_URL + env var, then to a local sqlite file if neither is provided (useful for + local development/tests without a real database provisioned). + """ + self.__engine: Engine = create_engine( + db_url or os.getenv('CATALYA_SQL_DB_URL', 'sqlite:///catalia_status.db'), + pool_pre_ping=True, + ) + self.__metadata = MetaData() + self.__data_columns = [ + self.identifier, self.datetime_str, self.collection, self.target_collection, + self.name_str, self.status, self.error_code, self.error_message, self.href_str, + ] + self.__table = self.__build_table(table_name) + + def __build_table(self, table_name: str) -> Table: + return Table( + table_name, self.__metadata, + Column('id', Integer, primary_key=True, autoincrement=True), + Column(self.identifier, String(255), nullable=False, index=True), + Column(self.datetime_str, String(64), nullable=False), + Column(self.collection, String(255)), + Column(self.target_collection, String(255)), + Column(self.name_str, String(255)), + Column(self.status, String(64)), + Column(self.error_code, String(255)), + Column(self.error_message, String(2048)), + Column(self.href_str, String(2048)), + UniqueConstraint(self.identifier, self.datetime_str, name=f'uq_{table_name}_{self.identifier}_{self.datetime_str}'), + ) + + def create_table_if_missing(self): + """ + Creates the underlying table (see DDL in the module docstring) if it doesn't + already exist. Not called automatically from __init__ so that schema creation + stays an explicit, deliberate action (e.g. once during deployment, or from + tests using a throwaway sqlite DB) rather than a side-effect of instantiation. + """ + self.__metadata.create_all(self.__engine, tables=[self.__table], checkfirst=True) + return self + + def get(self, identifier: str): + columns = [self.__table.c[col_name] for col_name in self.__data_columns] + stmt = select(*columns).where( + self.__table.c[self.identifier] == identifier + ).order_by(self.__table.c[self.datetime_str]) + with self.__engine.connect() as conn: + rows = conn.execute(stmt).mappings().all() + return [dict(row) for row in rows] + + def add(self, identifier: str, collection: str, name_str: str, status: str, datetime_str: str, error_code: str=None, error_message: str=None, href_str: str=None, target_collection: str=None): + item1 = { + self.identifier: identifier, + self.datetime_str: datetime_str, + self.name_str: name_str, + self.collection: collection, + self.status: status, + self.error_code: error_code, + self.error_message: error_message, + self.href_str: href_str, + self.target_collection: target_collection, + } + item1 = {k: v for k, v in item1.items() if v is not None} + stmt = insert(self.__table).values(**item1) + try: + with self.__engine.begin() as conn: + conn.execute(stmt) + except IntegrityError as e: + logger.warning(f'item already exists for identifier={identifier}, datetime={datetime_str}') + raise RuntimeError('Item exists. Unable to overwrite') from e + return diff --git a/requirements.txt b/requirements.txt index b4e2af79..7c22c588 100644 --- a/requirements.txt +++ b/requirements.txt @@ -31,6 +31,7 @@ requests-aws4auth==1.2.3 rpds-py==0.20.0 six==1.16.0 sniffio==1.3.1 +SQLAlchemy==2.0.36 starlette==0.38.6 tenacity==8.2.3 typing_extensions==4.12.2 diff --git a/setup.py b/setup.py index 243463ec..4880e132 100644 --- a/setup.py +++ b/setup.py @@ -7,7 +7,8 @@ 'mangum', 'uvicorn', 'pygeofilter', - 'python-dotenv' + 'python-dotenv', + 'sqlalchemy' ] setup( From cde0dea6e03a9c5bba5824872d3349f3a98cc471 Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Thu, 20 Aug 2026 14:32:25 -0700 Subject: [PATCH 2/8] fix: prep for report + analysis --- .../sql_mws/catalia_status_db.py | 91 ++++++++-- requirements.txt | 1 + setup.py | 3 +- tf-module/daac_delivery_analysis/main.tf | 34 ++++ tf-module/daac_delivery_analysis/rds.tf | 155 ++++++++++++++++++ .../terraform.tf.example | 7 + tf-module/daac_delivery_analysis/variables.tf | 31 ++++ 7 files changed, 308 insertions(+), 14 deletions(-) create mode 100644 tf-module/daac_delivery_analysis/main.tf create mode 100644 tf-module/daac_delivery_analysis/rds.tf create mode 100644 tf-module/daac_delivery_analysis/terraform.tf.example create mode 100644 tf-module/daac_delivery_analysis/variables.tf diff --git a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py index eb3f1039..be1cd748 100644 --- a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py +++ b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py @@ -8,7 +8,10 @@ The public API (class-level column-name constants + `get`/`add` method signatures) is kept identical to the DynamoDB version so existing callers (daac_receiver.py, services/status_update_svc.py, catalya_uds_api/granules_archive_api.py) can swap the -import without any other code changes. +import without any other code changes. A `search()` method is added on top for +time-range + collection/target_collection/status filtering to support analytics/report +use cases (funnel counts, success/failure rates, etc.) -- it returns raw matching rows +and leaves any aggregation to the caller. -------------------------------------------------------------------------------------- Suggested DDL (PostgreSQL dialect) for the table this class reads/writes. Replace @@ -22,10 +25,16 @@ the same "one row per status update, no silent overwrites" behavior that the DynamoDB version got from `replace=False` + a ConditionExpression. +``datetime`` is stored as epoch-milliseconds (UTC, BIGINT) rather than a formatted +string -- cheap to index/compare and unambiguous across timezones -- instead of the +RFC3339 string used on the DynamoDB side. Values passed into `add()`/`search()` may +still be given as RFC3339 strings (e.g. what `TimeUtils.get_current_time()` produces); +they're converted to epoch-milliseconds internally. + CREATE TABLE uds_ctla_daac_status ( id SERIAL PRIMARY KEY, identifier VARCHAR(255) NOT NULL, - datetime VARCHAR(64) NOT NULL, + datetime BIGINT NOT NULL, -- epoch milliseconds, UTC collection VARCHAR(255), target_collection VARCHAR(255), name VARCHAR(255), @@ -35,7 +44,11 @@ href VARCHAR(2048), CONSTRAINT uq_uds_ctla_daac_status_identifier_datetime UNIQUE (identifier, datetime) ); - CREATE INDEX ix_uds_ctla_daac_status_identifier ON uds_ctla_daac_status (identifier); + CREATE INDEX ix_uds_ctla_daac_status_identifier ON uds_ctla_daac_status (identifier); + CREATE INDEX ix_uds_ctla_daac_status_datetime ON uds_ctla_daac_status (datetime); + CREATE INDEX ix_uds_ctla_daac_status_collection ON uds_ctla_daac_status (collection); + CREATE INDEX ix_uds_ctla_daac_status_target_collection ON uds_ctla_daac_status (target_collection); + CREATE INDEX ix_uds_ctla_daac_status_status ON uds_ctla_daac_status (status); Note: `errorCode`/`errorMessage` are quoted above because they're camelCase -- Postgres folds unquoted identifiers to lowercase. The SQLAlchemy `Table` defined in @@ -44,15 +57,18 @@ """ import logging import os -from typing import Optional +from datetime import datetime, timezone +from typing import Optional, Union from sqlalchemy import ( + BigInteger, Column, Integer, MetaData, String, Table, UniqueConstraint, + and_, create_engine, insert, select, @@ -97,17 +113,34 @@ def __build_table(self, table_name: str) -> Table: table_name, self.__metadata, Column('id', Integer, primary_key=True, autoincrement=True), Column(self.identifier, String(255), nullable=False, index=True), - Column(self.datetime_str, String(64), nullable=False), - Column(self.collection, String(255)), - Column(self.target_collection, String(255)), + Column(self.datetime_str, BigInteger, nullable=False, index=True), + Column(self.collection, String(255), index=True), + Column(self.target_collection, String(255), index=True), Column(self.name_str, String(255)), - Column(self.status, String(64)), + Column(self.status, String(64), index=True), Column(self.error_code, String(255)), Column(self.error_message, String(2048)), Column(self.href_str, String(2048)), UniqueConstraint(self.identifier, self.datetime_str, name=f'uq_{table_name}_{self.identifier}_{self.datetime_str}'), ) + @staticmethod + def _to_epoch_millis(value: Union[int, float, str]) -> int: + """ + Accepts either an epoch-millisecond int/float, or an RFC3339 string + (e.g. "2026-08-20T00:00:00.000000Z", as produced by + `f'{TimeUtils.get_current_time()}Z'`), and returns the epoch-millisecond + int representation stored in / queried from the DB. Naive strings (no + UTC offset) are assumed to already be in UTC. + """ + if isinstance(value, (int, float)): + return int(value) + normalized = value[:-1] + '+00:00' if value.endswith('Z') else value + parsed = datetime.fromisoformat(normalized) + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return int(parsed.timestamp() * 1000) + def create_table_if_missing(self): """ Creates the underlying table (see DDL in the module docstring) if it doesn't @@ -119,18 +152,50 @@ def create_table_if_missing(self): return self def get(self, identifier: str): + return self.search(identifier=identifier) + + def search(self, start_datetime: Union[int, float, str] = None, end_datetime: Union[int, float, str] = None, + collection: str = None, target_collection: str = None, status: str = None, identifier: str = None): + """ + Returns raw rows matching the given filters, ordered by datetime ascending. + No aggregation is done here -- callers are expected to compute + success/failure rates, funnel counts, latency, etc. from the returned rows. + + :param start_datetime: inclusive lower bound (epoch-millis int, or RFC3339 string) + :param end_datetime: inclusive upper bound (epoch-millis int, or RFC3339 string) + :param collection: exact-match filter on the source collection + :param target_collection: exact-match filter on the DAAC target collection + :param status: exact-match filter on status (e.g. 'cnm-receive-success') + :param identifier: exact-match filter on the archiving identifier + """ columns = [self.__table.c[col_name] for col_name in self.__data_columns] - stmt = select(*columns).where( - self.__table.c[self.identifier] == identifier - ).order_by(self.__table.c[self.datetime_str]) + conditions = [] + if start_datetime is not None: + conditions.append(self.__table.c[self.datetime_str] >= self._to_epoch_millis(start_datetime)) + if end_datetime is not None: + conditions.append(self.__table.c[self.datetime_str] <= self._to_epoch_millis(end_datetime)) + if collection is not None: + conditions.append(self.__table.c[self.collection] == collection) + if target_collection is not None: + conditions.append(self.__table.c[self.target_collection] == target_collection) + if status is not None: + conditions.append(self.__table.c[self.status] == status) + if identifier is not None: + conditions.append(self.__table.c[self.identifier] == identifier) + + stmt = select(*columns) + if conditions: + stmt = stmt.where(and_(*conditions)) + stmt = stmt.order_by(self.__table.c[self.datetime_str]) + with self.__engine.connect() as conn: rows = conn.execute(stmt).mappings().all() return [dict(row) for row in rows] - def add(self, identifier: str, collection: str, name_str: str, status: str, datetime_str: str, error_code: str=None, error_message: str=None, href_str: str=None, target_collection: str=None): + def add(self, identifier: str, collection: str, name_str: str, status: str, datetime_str: Union[int, float, str], error_code: str=None, error_message: str=None, href_str: str=None, target_collection: str=None): item1 = { self.identifier: identifier, - self.datetime_str: datetime_str, + self.datetime_str: self._to_epoch_millis(datetime_str), self.name_str: name_str, self.collection: collection, self.status: status, diff --git a/requirements.txt b/requirements.txt index 7c22c588..b116e489 100644 --- a/requirements.txt +++ b/requirements.txt @@ -16,6 +16,7 @@ jsonschema-specifications==2023.12.1 lark==0.12.0 mangum==0.18.0 mdps-ds-lib==1.2.0.dev400 +psycopg2-binary==2.9.10 pydantic==2.9.2 pydantic_core==2.23.4 pygeofilter==0.2.4 diff --git a/setup.py b/setup.py index 4880e132..2f1aa3e7 100644 --- a/setup.py +++ b/setup.py @@ -8,7 +8,8 @@ 'uvicorn', 'pygeofilter', 'python-dotenv', - 'sqlalchemy' + 'sqlalchemy', + 'psycopg2-binary' ] setup( diff --git a/tf-module/daac_delivery_analysis/main.tf b/tf-module/daac_delivery_analysis/main.tf new file mode 100644 index 00000000..164bf708 --- /dev/null +++ b/tf-module/daac_delivery_analysis/main.tf @@ -0,0 +1,34 @@ + provider "aws" { + region = var.aws_region + ignore_tags { + key_prefixes = ["gsfc-ngap"] + } +} +data "aws_caller_identity" "current" {} + +locals { + account_id = data.aws_caller_identity.current.account_id + lambda_file_name = "${path.module}/build/cumulus_lambda_functions_deployment.zip" + security_group_ids_set = var.security_group_ids != null + lambda_role_arn = data.aws_iam_role.lambda_processing.arn + lambda_python_runtime = "python3.10" +} + +variable "buckets" { + description = "Map identifying the buckets for the deployment" + type = map(object({ name = string, type = string })) + default = {} +} +## resources = [for k, v in var.dynamo_tables : "${v.arn}/stream/*"] +#variable "dynamo_tables" { +# type = map(object({ name = string, arn = string })) +#} + +data "aws_security_group" "uds_lambda_sg_no_ingress_all_egress" { + name = "${var.prefix}-uds_lambda_sg_no_ingress_all_egress" +} + +data "aws_iam_role" "lambda_processing" { + # count = var.create_lambda_role ? 1 : 0 + name = "${var.prefix}-lambda-processing" +} diff --git a/tf-module/daac_delivery_analysis/rds.tf b/tf-module/daac_delivery_analysis/rds.tf new file mode 100644 index 00000000..e6d7cdea --- /dev/null +++ b/tf-module/daac_delivery_analysis/rds.tf @@ -0,0 +1,155 @@ +# Create an Autora V2 database and output the url / port / username. +# create a parameter store, secret text mode, and store the url / port / username / password so that applications can pull it to connect them. + +terraform { + required_providers { + random = { + source = "hashicorp/random" + version = "~> 3.6" + } + } +} + +variable "rds_master_username" { + type = string + default = "uds_ctla_admin" + description = "Master username for the Aurora Serverless v2 cluster" +} + +variable "rds_database_name" { + type = string + default = "daac_delivery_analysis" + description = "Initial database name created in the Aurora Serverless v2 cluster" +} + +variable "rds_engine_version" { + type = string + description = "Aurora PostgreSQL engine version. Must be a version that supports Serverless v2 (e.g. \"16.4\")" + default = "17.9" +} + +variable "rds_min_acu" { + type = number + default = 0.5 + description = "Minimum Aurora Capacity Units (ACUs) for Serverless v2 scaling" +} + +variable "rds_max_acu" { + type = number + default = 1 + description = "Maximum Aurora Capacity Units (ACUs) for Serverless v2 scaling" +} + +variable "rds_instance_count" { + type = number + default = 1 + description = "Number of Aurora Serverless v2 instances to create in the cluster" +} + +variable "rds_skip_final_snapshot" { + type = bool + default = true + description = "Whether to skip taking a final DB snapshot when the cluster is destroyed" +} + +variable "rds_deletion_protection" { + type = bool + default = false + description = "Whether to enable deletion protection on the Aurora cluster" +} + +resource "random_password" "aurora_master" { + length = 32 + special = true + override_special = "!#$%^&*()-_=+[]{}<>:?" +} + +resource "aws_db_subnet_group" "daac_delivery_analysis" { + name = "${var.prefix}-daac-delivery-analysis" + subnet_ids = var.cumulus_lambda_subnet_ids + tags = var.tags +} + +resource "aws_security_group" "daac_delivery_analysis_rds" { + name = "${var.prefix}-daac-delivery-analysis-rds" + vpc_id = var.cumulus_lambda_vpc_id + ingress { + from_port = 5432 + to_port = 5432 + protocol = "tcp" + security_groups = local.security_group_ids_set ? var.security_group_ids : [data.aws_security_group.uds_lambda_sg_no_ingress_all_egress.id] + } + egress { + from_port = 0 + to_port = 0 + protocol = "-1" + cidr_blocks = ["0.0.0.0/0"] + } + tags = var.tags +} + +resource "aws_rds_cluster" "daac_delivery_analysis" { + cluster_identifier = "${var.prefix}-daac-delivery-analysis" + engine = "aurora-postgresql" + engine_mode = "provisioned" + engine_version = var.rds_engine_version + + database_name = var.rds_database_name + master_username = var.rds_master_username + master_password = random_password.aurora_master.result + + db_subnet_group_name = aws_db_subnet_group.daac_delivery_analysis.name + vpc_security_group_ids = [aws_security_group.daac_delivery_analysis_rds.id] + + storage_encrypted = true + + skip_final_snapshot = var.rds_skip_final_snapshot + final_snapshot_identifier = "${var.prefix}-daac-delivery-analysis-final" + deletion_protection = var.rds_deletion_protection + + serverlessv2_scaling_configuration { + min_capacity = var.rds_min_acu + max_capacity = var.rds_max_acu + } + + lifecycle { + ignore_changes = [master_password] + } + + tags = var.tags +} + +resource "aws_rds_cluster_instance" "daac_delivery_analysis" { + count = var.rds_instance_count + cluster_identifier = aws_rds_cluster.daac_delivery_analysis.id + instance_class = "db.serverless" + engine = aws_rds_cluster.daac_delivery_analysis.engine + engine_version = aws_rds_cluster.daac_delivery_analysis.engine_version + tags = var.tags +} + +resource "aws_ssm_parameter" "daac_delivery_analysis_db_credentials" { + name = "/${var.prefix}/daac-delivery-analysis/rds_credentials" + type = "SecureString" + value = jsonencode({ + URL = aws_rds_cluster.daac_delivery_analysis.endpoint + PORT = aws_rds_cluster.daac_delivery_analysis.port + USERNAME = var.rds_master_username + PASSWORD = random_password.aurora_master.result + DBNAME = var.rds_database_name + }) + description = "Secure connection credentials for the DAAC delivery analysis Aurora Serverless v2 database" + tags = var.tags +} + +output "daac_delivery_analysis_db_url" { + value = aws_rds_cluster.daac_delivery_analysis.endpoint +} + +output "daac_delivery_analysis_db_port" { + value = aws_rds_cluster.daac_delivery_analysis.port +} + +output "daac_delivery_analysis_db_username" { + value = var.rds_master_username +} diff --git a/tf-module/daac_delivery_analysis/terraform.tf.example b/tf-module/daac_delivery_analysis/terraform.tf.example new file mode 100644 index 00000000..442b0a57 --- /dev/null +++ b/tf-module/daac_delivery_analysis/terraform.tf.example @@ -0,0 +1,7 @@ +terraform { + backend "s3" { + region = "us-west-2" + bucket = "catalya-app-catalog" + key = "catalya-uds-dev/archive_system/daac_delivery_analysis/terraform.tfstate" + } +} diff --git a/tf-module/daac_delivery_analysis/variables.tf b/tf-module/daac_delivery_analysis/variables.tf new file mode 100644 index 00000000..9b5d0216 --- /dev/null +++ b/tf-module/daac_delivery_analysis/variables.tf @@ -0,0 +1,31 @@ +variable "prefix" { + type = string +} +variable "aws_region" { + type = string + default = "us-west-2" +} + +variable "account_id" { + type = string + description = "AWS Account ID" +} + +variable "tags" { + description = "Tags to be applied to Cumulus resources that support tags" + type = map(string) + default = {} +} +variable "cumulus_lambda_vpc_id" { + type = string +} +variable "security_group_ids" { + description = "Security Group IDs for Lambdas" + type = list(string) + default = null +} +variable "cumulus_lambda_subnet_ids" { + description = "Subnet IDs for Lambdas" + type = list(string) + default = null +} From d05700e8ba48ff3c9cce7bf000189737377f0ad8 Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Thu, 20 Aug 2026 15:50:16 -0700 Subject: [PATCH 3/8] fix: replacing with status db --- .../catalya_uds_api/granules_archive_api.py | 7 ++-- .../daac_archiver/daac_receiver.py | 8 +++-- .../services/status_update_svc.py | 11 ++++--- .../sql_mws/catalia_status_db.py | 33 +++++++++++++++---- 4 files changed, 42 insertions(+), 17 deletions(-) diff --git a/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py b/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py index ab742ca6..a1f71896 100644 --- a/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py +++ b/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py @@ -5,7 +5,7 @@ from cumulus_lambda_functions.daac_archiver.daac_archiver_catalia_2 import DaacArchiverCatalia -from cumulus_lambda_functions.daac_archiver.ddb_mws.catalia_status_db import CataliaStatusDb +from cumulus_lambda_functions.daac_archiver.sql_mws.catalia_status_db import CataliaStatusDb from cumulus_lambda_functions.lib.lambda_logger_generator import LambdaLoggerGenerator from cumulus_lambda_functions.lib.uds_fast_api.internal_ddb_connector import InternalDDBConnector from cumulus_lambda_functions.lib.uds_fast_api.web_service_constants import WebServiceConstants @@ -323,8 +323,9 @@ async def archive_entire_collection_actual(request: Request, collection_id: str) @router.get("/{operation_id}/") async def get_archive_status(request: Request, operation_id: str): LOGGER.debug(f'started get_archive_status with operation_id: {operation_id}') - status_ddb = CataliaStatusDb(os.getenv('CATALYA_STATUS_DB', None)) - existing_statuses = status_ddb.get(operation_id) + uds_api_creds = json.loads(AwsParamStore().get_param(os.getenv('CATALYA_RDS_CREDS', 'NA'))) + status_db = CataliaStatusDb(os.getenv('CATALYA_STATUS_DB'), uds_api_creds) + existing_statuses = status_db.get(operation_id) if len(existing_statuses) < 1: raise HTTPException(status_code=404, detail=f'STATUS DB does not have any entry for {operation_id}') return {'status_list': existing_statuses} diff --git a/cumulus_lambda_functions/daac_archiver/daac_receiver.py b/cumulus_lambda_functions/daac_archiver/daac_receiver.py index b34fa48a..f5858ff2 100644 --- a/cumulus_lambda_functions/daac_archiver/daac_receiver.py +++ b/cumulus_lambda_functions/daac_archiver/daac_receiver.py @@ -2,11 +2,12 @@ import os import requests from mdps_ds_lib.lib.aws.aws_message_transformers import AwsMessageTransformers +from mdps_ds_lib.lib.aws.aws_param_store import AwsParamStore from mdps_ds_lib.lib.utils.json_validator import JsonValidator from cumulus_lambda_functions.daac_archiver.cnm_plugins.cnm_plugin_processor import CnmPluginProcessor from cumulus_lambda_functions.daac_archiver.cnm_plugins.cnm_plugin_abstract import CnmPluginAbstract -from cumulus_lambda_functions.daac_archiver.ddb_mws.catalia_status_db import CataliaStatusDb +from cumulus_lambda_functions.daac_archiver.sql_mws.catalia_status_db import CataliaStatusDb from cumulus_lambda_functions.lib.lambda_logger_generator import LambdaLoggerGenerator from cumulus_lambda_functions.lib.uds_db.uds_collections import UdsCollections @@ -36,8 +37,9 @@ def update_stac(self, cnm_notification_msg): raise ValueError(f"missing ARCHIVAL_STATUS_MECHANISM environment variable or value is not {['UDS', 'FAST_STAC']}") if update_type == 'UDS': return self.update_stac_uds(cnm_notification_msg) - status_ddb = CataliaStatusDb(os.getenv('CATALYA_STATUS_DB', None)) - existing_statuses = status_ddb.get(cnm_notification_msg['identifier']) + uds_api_creds = json.loads(AwsParamStore().get_param(os.getenv('CATALYA_RDS_CREDS', 'NA'))) + status_db = CataliaStatusDb(os.getenv('CATALYA_STATUS_DB'), uds_api_creds) + existing_statuses = status_db.get(cnm_notification_msg['identifier']) if len(existing_statuses) < 1: raise ValueError(f'unknown collection & granule: {cnm_notification_msg}') plugin_processor_params = { diff --git a/cumulus_lambda_functions/daac_archiver/services/status_update_svc.py b/cumulus_lambda_functions/daac_archiver/services/status_update_svc.py index 78ffc536..21877b9d 100644 --- a/cumulus_lambda_functions/daac_archiver/services/status_update_svc.py +++ b/cumulus_lambda_functions/daac_archiver/services/status_update_svc.py @@ -1,12 +1,13 @@ import json import os +from mdps_ds_lib.lib.aws.aws_param_store import AwsParamStore from mdps_ds_lib.lib.aws.aws_s3 import AwsS3 from mdps_ds_lib.lib.utils.time_utils import TimeUtils from cumulus_lambda_functions.daac_archiver.ddb_mws.catalia_archiving_traces import CataliaArchivingTraces from cumulus_lambda_functions.daac_archiver.services.sfa_client_mw import SfaClientMw +from cumulus_lambda_functions.daac_archiver.sql_mws.catalia_status_db import CataliaStatusDb from cumulus_lambda_functions.lib.lambda_logger_generator import LambdaLoggerGenerator -from cumulus_lambda_functions.daac_archiver.ddb_mws.catalia_status_db import CataliaStatusDb LOGGER = LambdaLoggerGenerator.get_logger(__name__, LambdaLoggerGenerator.get_level_from_env()) @@ -54,7 +55,9 @@ class StatusUpdateSvc: def __init__(self): self.__uds_ctla_archiving_traces = CataliaArchivingTraces(os.getenv('CATALYA_TRACING_DB', None)) - self.__status_ddb = CataliaStatusDb(os.getenv('CATALYA_STATUS_DB', None)) + uds_api_creds = json.loads(AwsParamStore().get_param(os.getenv('CATALYA_RDS_CREDS', 'NA'))) + self.__status_db = CataliaStatusDb(os.getenv('CATALYA_STATUS_DB'), uds_api_creds) + self.__archiving_granules_stac = None self.__identifier, self.__collection, self.__target_collection, self.__granule = None, None, None, None @@ -85,7 +88,7 @@ def validate_status(self, archival_status): def update_status_wrapper(self, cnm_notification_msg: dict): if any([k is None for k in [self.__identifier, self.__collection, self.__target_collection, self.__granule]]): - existing_statuses = self.__status_ddb.get(cnm_notification_msg['identifier']) + existing_statuses = self.__status_db.get(cnm_notification_msg['identifier']) if len(existing_statuses) < 1: raise ValueError(f'unknown collection & granule: {cnm_notification_msg}') self.__identifier, self.__collection, self.__target_collection, self.__granule = cnm_notification_msg['identifier'], existing_statuses[0][CataliaStatusDb.collection], existing_statuses[0][CataliaStatusDb.target_collection], existing_statuses[0][CataliaStatusDb.name_str] @@ -107,7 +110,7 @@ def update_status_ddb(self, archival_status): if any([k is None for k in [self.__identifier, self.__collection, self.__granule]]): raise ValueError(f'missing identifier, collection, or granule ID') try: - self.__status_ddb.add(self.__identifier, self.__collection, self.__granule, archival_status['status'], + self.__status_db.add(self.__identifier, self.__collection, self.__granule, archival_status['status'], archival_status['datetime'], archival_status['errorCode'] if 'errorCode' in archival_status else None, archival_status['errorMessage'] if 'errorMessage' in archival_status else None, diff --git a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py index be1cd748..2e85ef95 100644 --- a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py +++ b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py @@ -56,7 +56,6 @@ -------------------------------------------------------------------------------------- """ import logging -import os from datetime import datetime, timezone from typing import Optional, Union @@ -73,7 +72,7 @@ insert, select, ) -from sqlalchemy.engine import Engine +from sqlalchemy.engine import Engine, URL from sqlalchemy.exc import IntegrityError logger = logging.getLogger(__name__) @@ -90,15 +89,22 @@ class CataliaStatusDb: href_str = 'href' datetime_str = 'datetime' - def __init__(self, table_name: str, db_url: Optional[str] = None): + db_host_key = 'URL' + db_port_key = 'PORT' + db_user_key = 'USERNAME' + db_password_key = 'PASSWORD' + db_name_key = 'DBNAME' + + def __init__(self, table_name: str, db_config: dict): """ :param table_name: name of the SQL table to read/write, analogous to the DDB table name. - :param db_url: SQLAlchemy connection URL. Falls back to the CATALYA_SQL_DB_URL - env var, then to a local sqlite file if neither is provided (useful for - local development/tests without a real database provisioned). + :param db_config: dict used to build a PostgreSQL connection URL, matching the JSON + shape stored in the `/${prefix}/daac-delivery-analysis/rds_credentials` SSM + parameter (see tf-module/daac_delivery_analysis/rds.tf). Must contain all of the + following keys: `URL`, `PORT`, `USERNAME`, `PASSWORD`, `DBNAME`. """ self.__engine: Engine = create_engine( - db_url or os.getenv('CATALYA_SQL_DB_URL', 'sqlite:///catalia_status.db'), + self.__build_db_url(db_config), pool_pre_ping=True, ) self.__metadata = MetaData() @@ -108,6 +114,19 @@ def __init__(self, table_name: str, db_url: Optional[str] = None): ] self.__table = self.__build_table(table_name) + def __build_db_url(self, db_config: dict) -> URL: + missing_keys = [k for k in (self.db_host_key, self.db_port_key, self.db_user_key, self.db_password_key, self.db_name_key) if k not in db_config] + if missing_keys: + raise ValueError(f'db_config is missing required key(s): {missing_keys}') + return URL.create( + drivername='postgresql', + username=db_config[self.db_user_key], + password=db_config[self.db_password_key], + host=db_config[self.db_host_key], + port=int(db_config[self.db_port_key]), + database=db_config[self.db_name_key], + ) + def __build_table(self, table_name: str) -> Table: return Table( table_name, self.__metadata, From 2ec8e965fb26438a477d7ef3d99ac1ff0bc2ef81 Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Tue, 25 Aug 2026 05:21:47 -0700 Subject: [PATCH 4/8] fix: add missing env var --- tf-module/uds_catalia/daac_archiver.tf | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tf-module/uds_catalia/daac_archiver.tf b/tf-module/uds_catalia/daac_archiver.tf index cdf80354..449f23b1 100644 --- a/tf-module/uds_catalia/daac_archiver.tf +++ b/tf-module/uds_catalia/daac_archiver.tf @@ -50,6 +50,8 @@ resource "aws_lambda_function" "uds_daac_archiver_response" { CNM_PLUG_IN_NAMES = var.CNM_PLUG_IN_NAMES CNM_STORAGE_BUCKET = var.CNM_STORAGE_BUCKET CNM_STORAGE_PREFIX = var.CNM_STORAGE_PREFIX + CNM_STORAGE_CLASS = var.CNM_STORAGE_CLASS + } } From 3654f6abb7e5bc049facdd105b9f2221ae23a531 Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Tue, 25 Aug 2026 05:51:24 -0700 Subject: [PATCH 5/8] fix: add rds env var --- tf-module/uds_catalia/daac_archiver.tf | 1 + tf-module/uds_catalia/uds_api_lambda.tf | 1 + tf-module/uds_catalia/variables.tf | 4 ++++ 3 files changed, 6 insertions(+) diff --git a/tf-module/uds_catalia/daac_archiver.tf b/tf-module/uds_catalia/daac_archiver.tf index 449f23b1..04774c88 100644 --- a/tf-module/uds_catalia/daac_archiver.tf +++ b/tf-module/uds_catalia/daac_archiver.tf @@ -51,6 +51,7 @@ resource "aws_lambda_function" "uds_daac_archiver_response" { CNM_STORAGE_BUCKET = var.CNM_STORAGE_BUCKET CNM_STORAGE_PREFIX = var.CNM_STORAGE_PREFIX CNM_STORAGE_CLASS = var.CNM_STORAGE_CLASS + CATALYA_RDS_CREDS = var.CATALYA_RDS_CREDS_PARAM_PATH } } diff --git a/tf-module/uds_catalia/uds_api_lambda.tf b/tf-module/uds_catalia/uds_api_lambda.tf index 4c0d178a..a72fe885 100644 --- a/tf-module/uds_catalia/uds_api_lambda.tf +++ b/tf-module/uds_catalia/uds_api_lambda.tf @@ -30,6 +30,7 @@ resource "aws_lambda_function" "uds_api_1" { CNM_STORAGE_CLASS = var.CNM_STORAGE_CLASS CNM_STORAGE_BUCKET = var.CNM_STORAGE_BUCKET CNM_STORAGE_PREFIX = var.CNM_STORAGE_PREFIX + CATALYA_RDS_CREDS = var.CATALYA_RDS_CREDS_PARAM_PATH } } diff --git a/tf-module/uds_catalia/variables.tf b/tf-module/uds_catalia/variables.tf index 661c3b0b..a677bb17 100644 --- a/tf-module/uds_catalia/variables.tf +++ b/tf-module/uds_catalia/variables.tf @@ -299,4 +299,8 @@ variable "CNM_STORAGE_CLASS" { } variable "daac_account_ids" { type = list(string) +} +variable "CATALYA_RDS_CREDS_PARAM_PATH" { + type = string + description = "PARAMETER Store Path where it has a JSON dictionary of RDS / MARIA DB / Aurora V2 connection details: {\"DBNAME\":\"xxx\",\"PASSWORD\":\"xxx\",\"PORT\":xxx,\"URL\":\"xxx\",\"USERNAME\":\"xxx\"} " } \ No newline at end of file From 7fe1e39a635b2e7c6335ed89e6f9892569cb3af1 Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Tue, 25 Aug 2026 06:40:47 -0700 Subject: [PATCH 6/8] fix: update sns policy to accept data from daac --- tf-module/uds_catalia/daac_archiver.tf | 7 +++---- .../uds_catalia/daac_archiver_sns_policy.json | 18 ++++++------------ tf-module/uds_catalia/variables.tf | 6 +++++- 3 files changed, 14 insertions(+), 17 deletions(-) diff --git a/tf-module/uds_catalia/daac_archiver.tf b/tf-module/uds_catalia/daac_archiver.tf index 04774c88..4f1e8786 100644 --- a/tf-module/uds_catalia/daac_archiver.tf +++ b/tf-module/uds_catalia/daac_archiver.tf @@ -113,17 +113,16 @@ resource "aws_sns_topic" "uds_daac_archiver_response" { // TODO add access policy to be pushed from DAAC / other AWS account } -resource "aws_sns_topic_policy" "daac_archiver_response_policy" { +resource "aws_sns_topic_policy" "uds_daac_archiver_response_policy" { arn = aws_sns_topic.uds_daac_archiver_response.arn policy = templatefile("${path.module}/daac_archiver_sns_policy.json", { + sns_arn: aws_sns_topic.uds_daac_archiver_response.arn, + daac_lambda_2_sns_role: var.DAAC_LAMBDA_2_SNS_ROLE, region: var.aws_region, accountId: local.account_id, - snsName: "${var.prefix}-daac_archiver_response", - daacAccountIds: jsonencode(formatlist("arn:aws:iam::%s:root", var.daac_account_ids)) }) } - module "daac_archiver_response" { source = "../sqs--sns-lambda-connector" diff --git a/tf-module/uds_catalia/daac_archiver_sns_policy.json b/tf-module/uds_catalia/daac_archiver_sns_policy.json index b0c50a04..d2ae0e2e 100644 --- a/tf-module/uds_catalia/daac_archiver_sns_policy.json +++ b/tf-module/uds_catalia/daac_archiver_sns_policy.json @@ -18,23 +18,17 @@ "SNS:ListSubscriptionsByTopic", "SNS:Publish" ], - "Resource": "arn:aws:sns:${region}:${accountId}:${snsName}" + "Resource": "${sns_arn}" }, + { - "Sid": "2", + "Sid": "AllowCumulusUat3LambdaProcessingPublish", "Effect": "Allow", "Principal": { - "AWS": ${daacAccountIds} + "AWS": "${daac_lambda_2_sns_role}" }, - "Action": [ - "SNS:Publish" - ], - "Resource": "arn:aws:sns:${region}:${accountId}:${snsName}", - "Condition": { - "StringEquals": { - "AWS:SourceOwner": "${accountId}" - } - } + "Action": "SNS:Publish", + "Resource": "${sns_arn}" } ] } diff --git a/tf-module/uds_catalia/variables.tf b/tf-module/uds_catalia/variables.tf index a677bb17..ce3337da 100644 --- a/tf-module/uds_catalia/variables.tf +++ b/tf-module/uds_catalia/variables.tf @@ -303,4 +303,8 @@ variable "daac_account_ids" { variable "CATALYA_RDS_CREDS_PARAM_PATH" { type = string description = "PARAMETER Store Path where it has a JSON dictionary of RDS / MARIA DB / Aurora V2 connection details: {\"DBNAME\":\"xxx\",\"PASSWORD\":\"xxx\",\"PORT\":xxx,\"URL\":\"xxx\",\"USERNAME\":\"xxx\"} " -} \ No newline at end of file +} +variable "DAAC_LAMBDA_2_SNS_ROLE" { + type = string + description = "arn:aws:iam:::role/" +} From e6d80168f743e0dd3556b2344e525c2d4cfeb7b3 Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Tue, 25 Aug 2026 10:08:38 -0700 Subject: [PATCH 7/8] fix: adding create table endpoint as admin --- .../catalya_uds_api/auth_admin_api.py | 35 +++++++++++++++++++ .../catalya_uds_api/granules_archive_api.py | 19 ++++++++++ .../sql_mws/catalia_status_db.py | 9 +++++ 3 files changed, 63 insertions(+) diff --git a/cumulus_lambda_functions/catalya_uds_api/auth_admin_api.py b/cumulus_lambda_functions/catalya_uds_api/auth_admin_api.py index b45a929c..a1f4a58f 100644 --- a/cumulus_lambda_functions/catalya_uds_api/auth_admin_api.py +++ b/cumulus_lambda_functions/catalya_uds_api/auth_admin_api.py @@ -1,10 +1,13 @@ +import json from typing import Union from cumulus_lambda_functions.daac_archiver.ddb_mws.catalia_auth_db import CataliaAuthDb +from cumulus_lambda_functions.daac_archiver.sql_mws.catalia_status_db import CataliaStatusDb from cumulus_lambda_functions.lib.lambda_logger_generator import LambdaLoggerGenerator from cumulus_lambda_functions.lib.uds_fast_api.fast_api_utils import FastApiUtils from cumulus_lambda_functions.lib.uds_fast_api.web_service_constants import WebServiceConstants from fastapi import APIRouter, HTTPException, Request +from mdps_ds_lib.lib.aws.aws_param_store import AwsParamStore LOGGER = LambdaLoggerGenerator.get_logger(__name__, LambdaLoggerGenerator.get_level_from_env()) @@ -248,3 +251,35 @@ async def list_auth_mappings(request: Request, tenant: Union[str, None]=None, ve if query_result['statusCode'] == 200: return query_result['body'] raise HTTPException(status_code=query_result['statusCode'], detail=query_result['body']) + +@router.post("/status-table") +@router.post("/status-table/") +async def build_status_table(request: Request): + """ + Ensures the status db table (CATALYA_STATUS_DB) exists in the RDS Postgres + database, creating it -- along with its indexes/unique constraint -- via + CataliaStatusDb.create_table_if_missing() if it doesn't. Meant to be called once + on startup/deployment so the archiving Lambdas never hit a missing-table error + at runtime. + """ + LOGGER.debug('started build_status_table') + auth_info = FastApiUtils.get_authorization_info(request) + auth_crud = AuthCrud(auth_info, {}) + is_admin_result = auth_crud.is_admin() + if is_admin_result['statusCode'] != 200: + raise HTTPException(status_code=is_admin_result['statusCode'], detail=is_admin_result['body']) + + required_env = ['CATALYA_RDS_CREDS', 'CATALYA_STATUS_DB'] + if not all([k in os.environ for k in required_env]): + LOGGER.error(f'one or more missing env: {required_env}') + raise HTTPException(status_code=500, detail=f'one or more missing env: {required_env}') + + table_name = os.getenv('CATALYA_STATUS_DB') + db_config = json.loads(AwsParamStore().get_param(os.getenv('CATALYA_RDS_CREDS'))) + status_db = CataliaStatusDb(table_name, db_config) + + already_existed = status_db.table_exists() + if not already_existed: + status_db.create_table_if_missing() + LOGGER.info(f'created status db table: {table_name}') + return {'table': table_name, 'already_existed': already_existed, 'created': not already_existed} \ No newline at end of file diff --git a/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py b/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py index a1f71896..09f63479 100644 --- a/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py +++ b/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py @@ -329,3 +329,22 @@ async def get_archive_status(request: Request, operation_id: str): if len(existing_statuses) < 1: raise HTTPException(status_code=404, detail=f'STATUS DB does not have any entry for {operation_id}') return {'status_list': existing_statuses} + + +""" +Add a new endpoint for report. +Ask for thse 2 + + collection VARCHAR(255), + target_collection VARCHAR(255), + +They are mandatory. +Ask for start and end time which are optional. +So, it could be +1. for all time, +2. after start time +3. before end time +4. from start to end. + +Use table to write a query +""" \ No newline at end of file diff --git a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py index 2e85ef95..e5ff7190 100644 --- a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py +++ b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py @@ -69,6 +69,7 @@ UniqueConstraint, and_, create_engine, + inspect, insert, select, ) @@ -160,6 +161,14 @@ def _to_epoch_millis(value: Union[int, float, str]) -> int: parsed = parsed.replace(tzinfo=timezone.utc) return int(parsed.timestamp() * 1000) + def table_exists(self) -> bool: + """ + Returns True if the underlying table already exists in the database, False + otherwise. Useful for callers that want to know/report whether + create_table_if_missing() is actually about to create something. + """ + return inspect(self.__engine).has_table(self.__table.name, schema=self.__table.schema) + def create_table_if_missing(self): """ Creates the underlying table (see DDL in the module docstring) if it doesn't From 88a145b1fd2f58c528e15c65560cb43d222c4fb6 Mon Sep 17 00:00:00 2001 From: Wai Phyo Date: Tue, 25 Aug 2026 10:51:56 -0700 Subject: [PATCH 8/8] feat: add report endpoint --- .../catalya_uds_api/granules_archive_api.py | 72 +++++++++++++------ .../sql_mws/catalia_status_db.py | 30 ++++++++ 2 files changed, 82 insertions(+), 20 deletions(-) diff --git a/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py b/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py index 09f63479..9245e3ac 100644 --- a/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py +++ b/cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py @@ -5,6 +5,7 @@ from cumulus_lambda_functions.daac_archiver.daac_archiver_catalia_2 import DaacArchiverCatalia +from cumulus_lambda_functions.daac_archiver.services.status_update_svc import StatusUpdateSvc from cumulus_lambda_functions.daac_archiver.sql_mws.catalia_status_db import CataliaStatusDb from cumulus_lambda_functions.lib.lambda_logger_generator import LambdaLoggerGenerator from cumulus_lambda_functions.lib.uds_fast_api.internal_ddb_connector import InternalDDBConnector @@ -319,6 +320,56 @@ async def archive_entire_collection_actual(request: Request, collection_id: str) return {'message': 'archive initiated'} +@router.get("/report") +@router.get("/report/") +async def get_archive_report(request: Request, collection: str, target_collection: str, + start_datetime: Optional[str] = None, end_datetime: Optional[str] = None): + """ + Summary status counts (submit/response success & failure, plus the full + per-status breakdown) for a given collection + target_collection pair. + + collection & target_collection are mandatory. start_datetime/end_datetime are + optional (either, both, or neither may be given), covering all-time / since / + until / between reporting windows. Values may be epoch-millis ints or RFC3339 + strings (same formats CataliaStatusDb.search() already accepts). + + Status names come from StatusUpdateSvc.archival_status_schema's enum instead of + being hardcoded here, so this stays in sync if that list changes. + """ + LOGGER.debug(f'started get_archive_report for collection={collection}, target_collection={target_collection}, ' + f'start_datetime={start_datetime}, end_datetime={end_datetime}') + uds_api_creds = json.loads(AwsParamStore().get_param(os.getenv('CATALYA_RDS_CREDS', 'NA'))) + status_db = CataliaStatusDb(os.getenv('CATALYA_STATUS_DB'), uds_api_creds) + + status_counts = status_db.count_by_status(collection, target_collection, start_datetime, end_datetime) + + known_statuses = StatusUpdateSvc.archival_status_schema['properties']['status']['enum'] + full_counts = {status: status_counts.get(status, 0) for status in known_statuses} + + def find_status(stage_keyword: str, result_keyword: str) -> Optional[str]: + matches = [s for s in known_statuses if stage_keyword in s and result_keyword in s] + return matches[0] if matches else None + + submit_success_status = find_status('submit', 'success') + submit_failed_status = find_status('submit', 'failed') + response_success_status = find_status('receive', 'success') + response_failed_status = find_status('receive', 'failed') + + return { + 'collection': collection, + 'target_collection': target_collection, + 'start_datetime': start_datetime, + 'end_datetime': end_datetime, + 'status_counts': full_counts, + 'summary': { + 'submit_success': full_counts.get(submit_success_status, 0), + 'submit_failed': full_counts.get(submit_failed_status, 0), + 'response_success': full_counts.get(response_success_status, 0), + 'response_failed': full_counts.get(response_failed_status, 0), + }, + } + + @router.get("/{operation_id}") @router.get("/{operation_id}/") async def get_archive_status(request: Request, operation_id: str): @@ -328,23 +379,4 @@ async def get_archive_status(request: Request, operation_id: str): existing_statuses = status_db.get(operation_id) if len(existing_statuses) < 1: raise HTTPException(status_code=404, detail=f'STATUS DB does not have any entry for {operation_id}') - return {'status_list': existing_statuses} - - -""" -Add a new endpoint for report. -Ask for thse 2 - - collection VARCHAR(255), - target_collection VARCHAR(255), - -They are mandatory. -Ask for start and end time which are optional. -So, it could be -1. for all time, -2. after start time -3. before end time -4. from start to end. - -Use table to write a query -""" \ No newline at end of file + return {'status_list': existing_statuses} \ No newline at end of file diff --git a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py index e5ff7190..c6ad0de5 100644 --- a/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py +++ b/cumulus_lambda_functions/daac_archiver/sql_mws/catalia_status_db.py @@ -69,6 +69,7 @@ UniqueConstraint, and_, create_engine, + func, inspect, insert, select, @@ -220,6 +221,35 @@ def search(self, start_datetime: Union[int, float, str] = None, end_datetime: Un rows = conn.execute(stmt).mappings().all() return [dict(row) for row in rows] + def count_by_status(self, collection: str, target_collection: str, + start_datetime: Union[int, float, str] = None, end_datetime: Union[int, float, str] = None) -> dict: + """ + Returns {status: count} for the given collection/target_collection (both + mandatory), optionally bounded by [start_datetime, end_datetime] -- either, + both, or neither may be given, covering "all time" / "since" / "until" / + "between" reporting windows. Statuses with zero matching rows are simply + absent from the result (callers wanting a fixed set of statuses -- e.g. from + StatusUpdateSvc.archival_status_schema -- should fill in 0 for any missing). + """ + conditions = [ + self.__table.c[self.collection] == collection, + self.__table.c[self.target_collection] == target_collection, + ] + if start_datetime is not None: + conditions.append(self.__table.c[self.datetime_str] >= self._to_epoch_millis(start_datetime)) + if end_datetime is not None: + conditions.append(self.__table.c[self.datetime_str] <= self._to_epoch_millis(end_datetime)) + + status_col = self.__table.c[self.status] + stmt = ( + select(status_col, func.count().label('status_count')) + .where(and_(*conditions)) + .group_by(status_col) + ) + with self.__engine.connect() as conn: + rows = conn.execute(stmt).all() + return {row[0]: row[1] for row in rows} + def add(self, identifier: str, collection: str, name_str: str, status: str, datetime_str: Union[int, float, str], error_code: str=None, error_message: str=None, href_str: str=None, target_collection: str=None): item1 = { self.identifier: identifier,