Skip to content

Commit 13db2fe

Browse files
committed
PR-7502 Implemented CNM creation
1 parent 267fb5b commit 13db2fe

12 files changed

Lines changed: 319 additions & 34 deletions

File tree

.github/workflows/lint.yml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,5 +25,5 @@ jobs:
2525
- uses: actions/checkout@v4
2626
- uses: astral-sh/ruff-action@v3
2727
with:
28-
version: "~=0.13.3"
28+
version: "~=0.15.13"
2929
args: format --check --diff --output-format=github

ctorm/ctorm.cfg.example

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,22 +1,26 @@
11
[ctorm]
22

3+
calc_md5 = true
4+
granule_goal = 200
35
granules_sqs_queue_url = "https://sqs.us-west-2.amazonaws.com/123456789012/ctorm-dev-granules"
46

57
[[ ctorm.pipelines ]]
68
bucketname = "nisar-bucket-0"
79
ingest_keypair_name = ""
810
ingest_queue_name = "sds-n-cumulus-test-nisar-workflow-queue"
911
ingest_sqs_queue_url = ""
12+
prepare_keypair_name = "NISAR_P"
13+
provider = "FOO"
1014
share = 50
1115
ummg_prefix = "UMMG/FOO_COLLECTION/"
12-
prepare_keypair_name = "NISAR_P"
1316

1417
[[ctorm.pipelines]]
1518
bucketname = "nisar-bucket-2"
1619
ingest_keypair_name = ""
1720
ingest_queue_name = "sds-n-cumulus-test-nisar-workflow-queue"
1821
ingest_sqs_queue_url = ""
1922
prepare_keypair_name = "NISAR"
23+
provider = "FOO"
2024
share = 45
2125
ummg_prefix = "UMMG/BAR_COLLECTION/"
2226

@@ -26,5 +30,6 @@ ingest_keypair_name = ""
2630
ingest_queue_name = "asf-cumulus-test-opera-workflow-queue"
2731
ingest_sqs_queue_url = ""
2832
prepare_keypair_name = "OPERA"
33+
provider = "FOO"
2934
share = 5
3035
ummg_prefix = "UMMG/BAZ_COLLECTION/"

ctorm/ctorm/cnm.py

Lines changed: 131 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,131 @@
1+
import re
2+
import uuid
3+
from collections.abc import Callable
4+
from datetime import datetime, timezone
5+
from logging import getLogger
6+
from pathlib import Path
7+
8+
from ctorm.config import (
9+
MD5_CHECKSUM_PATTERN,
10+
CtormPipeline,
11+
CtormPreparedFile,
12+
CtormPreparedGranule,
13+
)
14+
15+
log = getLogger(__name__)
16+
17+
DATA_TYPE_MAP = {
18+
".context.json": "metadata",
19+
".dataset.json": "metadata",
20+
".h5": "data",
21+
".iso.xml": "metadata",
22+
".jpg": "browse",
23+
".log": "metadata",
24+
".md5": "metadata",
25+
".met.json": "metadata",
26+
".nc": "data",
27+
".pdf": "qa",
28+
".png": "browse",
29+
".qa.h5": "qa",
30+
".rc.yaml": "metadata",
31+
".tif": "data",
32+
".xml": "data",
33+
}
34+
35+
type CtormCnmSGeneratorType = Callable[[CtormPreparedGranule, str], dict]
36+
37+
38+
class CtormCnmSGenerator:
39+
"""input dict like:
40+
```
41+
{
42+
"bm": {
43+
"sds-n-cumulus-prod-nisar-products": "B1",
44+
"sds-n-cumulus-prod-nisar-jpl-private-data": "B2"
45+
},
46+
"g": "NISAR_L0_RRST_VC08_20250821T101036_20250821T101041_P00408_J_001",
47+
"c": "NISAR_L0A_RRST_BETA_V1",
48+
"cv": "1",
49+
"f": [
50+
{
51+
"f": "s3://$B1/NISAR_L0A_RRST_BETA_V1/$G/$G.bin",
52+
"m": "e989430f4c5bfa04eeb26e56f4cd8375",
53+
"s": 1073829760
54+
},
55+
{
56+
"f": "s3://$B1/NISAR_L0A_RRST_BETA_V1/$G/$G.rc.yaml",
57+
"s": 120597,
58+
"m": "11bcaec780f0d246f3996b3f08680410"
59+
},
60+
{
61+
"f": "s3://$B1/NISAR_L0A_RRST_BETA_V1/$G/$G.bin.qa",
62+
"s": 1134,
63+
},
64+
]
65+
},
66+
```
67+
"""
68+
69+
def _assemble_file_dict(
70+
self, file_dict: CtormPreparedFile, g_name: str, bucket_map: dict
71+
):
72+
file_s3uri = file_dict["f"].replace("$G", g_name)
73+
for bucket, replacetoken in bucket_map.items():
74+
file_s3uri = file_s3uri.replace(f"${replacetoken}", bucket)
75+
filename = Path(file_s3uri).name
76+
return {
77+
"name": filename,
78+
"type": self._get_type(filename),
79+
"uri": file_s3uri,
80+
"size": file_dict["s"],
81+
"checksum": self._get_checksum(file_dict),
82+
"checksumType": "md5",
83+
}
84+
85+
def __call__(
86+
self,
87+
granule_dict: CtormPreparedGranule,
88+
provider: str,
89+
) -> dict:
90+
cnm_s = {
91+
"identifier": str(uuid.uuid4()),
92+
"collection": granule_dict["c"],
93+
"version": "1.3",
94+
"submissionTime": datetime.now(tz=timezone.utc).strftime(
95+
"%Y-%m-%dT%H:%M:%S.%fZ"
96+
),
97+
"product": {
98+
"name": granule_dict["g"],
99+
"dataVersion": granule_dict["cv"],
100+
"files": [
101+
self._assemble_file_dict(
102+
file, granule_dict["g"], granule_dict["bm"]
103+
)
104+
for file in granule_dict["f"]
105+
],
106+
},
107+
"provider": provider,
108+
}
109+
110+
return cnm_s
111+
112+
def _get_type(self, filename: str):
113+
114+
suffixes = Path(filename).suffixes
115+
while suffixes:
116+
data_type = DATA_TYPE_MAP.get("".join(suffixes))
117+
118+
if data_type:
119+
return data_type
120+
121+
suffixes = suffixes[1:]
122+
123+
log.debug("suffix from %s is not in map. Using 'data'", filename)
124+
return "data"
125+
126+
def _get_checksum(self, file_dict: CtormPreparedFile):
127+
m = MD5_CHECKSUM_PATTERN.match(file_dict.get("m", "not md5"))
128+
if m:
129+
return m.group(1)
130+
131+
return "00000000000000000000000000000000"

ctorm/ctorm/cnm_sender.py

Lines changed: 22 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,25 @@
1+
from logging import getLogger
2+
from typing import List
3+
4+
from ctorm.cnm import CtormCnmSGenerator
5+
from ctorm.config import CtormConfig, CtormPreparedGranule
6+
7+
log = getLogger(__name__)
8+
9+
110
class CnmSender:
2-
def __init__(self, cfg):
11+
def __init__(
12+
self,
13+
cfg: CtormConfig,
14+
granules: List[CtormPreparedGranule],
15+
provider: str = "TODO: PROVIDER",
16+
):
317
self.cfg = cfg
18+
self.granules = granules
19+
self.provider = provider
20+
self.cnm_s_generator = CtormCnmSGenerator()
421

5-
def send(self):
6-
pass
22+
def send_all(self):
23+
for granule in self.granules:
24+
cnms = self.cnm_s_generator(granule, self.provider)
25+
log.debug("Sending %s", cnms)

ctorm/ctorm/config.py

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
1+
import re
2+
import tomllib
13
from dataclasses import dataclass, fields
4+
from typing import NotRequired, TypedDict
25

3-
import tomllib
46
from boto3.session import Session
57
from mypy_boto3_s3 import S3Client
68

@@ -11,6 +13,21 @@
1113

1214
DEFAULT = "default"
1315
AWS_REGION = "us-west-2" # We will never not want us-west-2
16+
MD5_CHECKSUM_PATTERN = re.compile(r"^([\da-f]{32})$")
17+
18+
19+
class CtormPreparedFile(TypedDict):
20+
f: str
21+
m: NotRequired[str]
22+
s: int
23+
24+
25+
class CtormPreparedGranule(TypedDict):
26+
bm: dict[str, str]
27+
g: str
28+
c: str
29+
cv: str
30+
f: list[CtormPreparedFile]
1431

1532

1633
@dataclass
@@ -33,13 +50,16 @@ class CtormPipeline:
3350
s3_client: S3Client = None
3451
session: Session = None
3552
ummg_prefix: str = "UMMG/"
53+
provider: str = "TODO"
3654

3755

3856
@dataclass
3957
class CtormConfig:
4058
pipelines: list
4159
granules_sqs_queue_url: str
4260

61+
calc_md5: bool = True
62+
4363
# This is the number of granules that will be prepared for test.
4464
granule_goal: int = 2000000
4565

ctorm/ctorm/lambda_run.py

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,21 @@
33
"""
44

55
import json
6+
import logging
67

78
from ctorm.load_tester import lambda_handler
89

10+
log = logging.getLogger()
11+
12+
base_fmt_str = "%(levelname)s: %(message)s (%(filename)s line %(lineno)d/)"
13+
screen_fmt = logging.Formatter(
14+
"%(asctime)s.%(msecs)d " + base_fmt_str, "%Y-%m-%dT%H:%M:%S"
15+
)
16+
screenlog = logging.StreamHandler()
17+
screenlog.setFormatter(screen_fmt)
18+
log.addHandler(screenlog)
19+
20+
921
# mock event data
1022
with open("../../data/lambda_event.json", "r") as f:
1123
mock_event = json.load(f)

ctorm/ctorm/load_tester.py

Lines changed: 27 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,27 @@
1+
import json
12
import logging
23
import os
4+
from typing import List
35

6+
from aws_lambda_typing import context as context_
7+
from aws_lambda_typing import events
48
from cnm_sender import CnmSender
59

10+
from ctorm.config import (
11+
CtormConfig,
12+
CtormPreparedGranule,
13+
)
14+
615
log = logging.getLogger(__name__)
716

817

18+
def configure_cfg():
19+
cfg = CtormConfig.from_file(
20+
cfg_file=os.getenv("CFG_FILE", "./ctorm.cfg"),
21+
)
22+
return cfg
23+
24+
925
def configure_logging() -> None:
1026
level_name = os.environ.get("LOG_LEVEL", "INFO").upper()
1127
level = getattr(logging, level_name, logging.INFO)
@@ -21,17 +37,20 @@ def configure_logging() -> None:
2137
for handler in root_logger.handlers:
2238
handler.setLevel(level)
2339
handler.setFormatter(
24-
logging.Formatter("%(asctime)s %(levelname)s %(name)s %(filename)s:%(lineno)d - %(message)s")
40+
logging.Formatter(
41+
"%(asctime)s %(levelname)s %(name)s %(filename)s:%(lineno)d - %(message)s"
42+
)
2543
)
2644

2745

28-
def load_test():
29-
c_sender = CnmSender({})
30-
c_sender.send()
46+
def load_test(cfg: CtormConfig, granule_list: List[CtormPreparedGranule]):
47+
c_sender = CnmSender(cfg, granule_list)
48+
c_sender.send_all()
3149

3250

33-
def lambda_handler(event, context):
51+
def lambda_handler(event: events.EventBridgeEvent, context: context_.Context):
3452
configure_logging()
53+
cfg = configure_cfg()
3554

3655
log.info(
3756
"Starting CNM sender invocation",
@@ -42,8 +61,9 @@ def lambda_handler(event, context):
4261

4362
try:
4463
log.debug("Received event: %s", event)
45-
46-
load_test()
64+
g_list = event["Records"].pop().get("body")
65+
g_list = json.loads(g_list).get("granules")
66+
load_test(cfg, g_list)
4767

4868
log.info("CNM sender invocation completed")
4969
return {"ok": True}
@@ -54,5 +74,4 @@ def lambda_handler(event, context):
5474

5575

5676
if __name__ == "__main__":
57-
configure_logging()
5877
lambda_handler({}, {})

0 commit comments

Comments
 (0)