Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 35 additions & 0 deletions cumulus_lambda_functions/catalya_uds_api/auth_admin_api.py
Original file line number Diff line number Diff line change
@@ -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())

Expand Down Expand Up @@ -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}
60 changes: 56 additions & 4 deletions cumulus_lambda_functions/catalya_uds_api/granules_archive_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@


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.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
from cumulus_lambda_functions.lib.uds_fast_api.web_service_constants import WebServiceConstants
Expand Down Expand Up @@ -319,12 +320,63 @@ 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):
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}
return {'status_list': existing_statuses}
8 changes: 5 additions & 3 deletions cumulus_lambda_functions/daac_archiver/daac_receiver.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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 = {
Expand Down
Original file line number Diff line number Diff line change
@@ -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())

Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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]
Expand All @@ -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,
Expand Down
Empty file.
Loading
Loading