import os
import pathlib
import requests
from recap.converters.avro import AvroConverter
from recap.converters.protobuf import ProtobufConverter
from recap.types import to_dict
GABLE_API_ENDPOINT = os.environ["GABLE_API_ENDPOINT"]
GABLE_API_KEY = os.environ["GABLE_API_KEY"]
COMPONENT_ID = os.environ["COMPONENT_ID"]
TABLE_ID = os.environ["TABLE_ID"]
TABLE_NAME = os.environ["TABLE_NAME"]
ROOT = pathlib.Path(os.environ.get("SCHEMA_ROOT", "schemas"))
def extract_fields(path: pathlib.Path):
'''
Uses Gable's recap-core library to parse Avro and Proto files
and return their fields
'''
text = path.read_text()
if path.suffix == ".avsc":
recap_type = AvroConverter().to_recap(text)
elif path.suffix == ".proto":
recap_type = ProtobufConverter().to_recap(text)
else:
return []
return to_dict(recap_type)["fields"]
def start_run():
start_run_request = {
"action": "upload",
"type": "DATA_STORE",
"collection_mechanism": "BYOL",
"code_info": {
"namespace": "qa",
"repo_uri": f"{os.environ['GITHUB_SERVER_URL']}/{os.environ['GITHUB_REPOSITORY']}",
"repo_branch": os.environ["GITHUB_REF_NAME"],
"repo_commit": os.environ["GITHUB_SHA"],
"project_root": "/",
"repo_name": os.environ["GITHUB_REPOSITORY"].split("/")[-1],
"job_trigger": "MANUAL",
"external_component_id": COMPONENT_ID,
},
}
r = requests.post(
f"{GABLE_API_ENDPOINT}/v0/sca/start-run",
json=start_run_request,
headers={"x-api-key": GABLE_API_KEY},
timeout=30,
)
r.raise_for_status()
return r.json()["runId"]
def publish(payload):
r = requests.post(
f"{GABLE_API_ENDPOINT}/v0/sca/results",
json=payload,
headers={"x-api-key": GABLE_API_KEY},
timeout=30,
)
r.raise_for_status()
def main():
run_id = start_run()
fields = []
# Parse each file using recap and combine their schemas
for path in sorted(ROOT.rglob("*")):
if path.suffix not in (".avsc", ".proto"):
continue
fields.extend(extract_fields(path))
print(f"collected {path}")
# Construct the component's payload
payload = {
"external_component_id": COMPONENT_ID,
"type": "DATA_STORE",
"run_id": run_id,
"external_table_id": TABLE_ID,
"table_metadata": {
"type": "dynamodb",
"table_name": TABLE_NAME,
},
"schema": {"fields": fields},
}
publish(payload)
if __name__ == "__main__":
main()