Skip to content

rs_workflows/utils/catalog.md

<< Back to index

Helper task to interact with the rs-catalog.

get_catalog_items(flow_env, item_ids, collections) async

Get items from a set of rs-catalog collections

Source code in docs/rs-client-libraries/rs_workflows/utils/catalog.py
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
@task(name="Retrieve rs-catalog items from collections")
async def get_catalog_items(flow_env: FlowEnv, item_ids: list[str], collections: list[str]) -> ItemCollection | None:
    """
    Get items from a set of rs-catalog collections
    """
    get_run_logger().info(
        f"Search items {', '.join(item_ids)} in the collections {', '.join(collections)} from the rs-catalog.",
    )
    size = len(item_ids)
    return flow_env.rs_client.get_catalog_client().search(
        method="POST",
        collections=collections,
        ids=item_ids,
        max_items=size,
        limit=size,
    )

get_single_catalog_item(flow_env, item_id, collections) async

Get an item from a set of rs-catalog collections

Source code in docs/rs-client-libraries/rs_workflows/utils/catalog.py
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
@task(name="Retrieve rs-catalog item from collection")
async def get_single_catalog_item(flow_env: FlowEnv, item_id: str, collections: list[str]) -> Item | None:
    """
    Get an item from a set of rs-catalog collections
    """
    logger = get_run_logger()
    result: Item | None = None

    # Try to retrieve the item in the collections
    item_collection: ItemCollection | None = await get_catalog_items(flow_env, [item_id], collections)

    count = 0
    if item_collection is not None:
        count = len(item_collection.items)
    if count == 1:
        # One item was found on the rs-catalog
        logger.info(
            f"✅ The STAC item 🧊 '{item_id}' has been found on the rs-catalog collections {', '.join(collections)}.",
        )
        if item_collection is not None:
            result = item_collection.items[0]
    else:
        logger.warning(
            f"⚠️ The STAC item 🧊 '{item_id}' was not found on the rs-catalog collections {', '.join(collections)}.",
        )

    return result

is_published(item)

Check if the item is published.

Source code in docs/rs-client-libraries/rs_workflows/utils/catalog.py
88
89
90
91
92
93
94
95
96
def is_published(item: Item) -> tuple[bool, datetime | None]:
    """
    Check if the item is published.
    """
    if published_date_str := item.properties.get("published"):
        published_date = datetime.fromisoformat(published_date_str.replace("Z", "+00:00"))
        return published_date <= datetime.now(timezone.utc), published_date

    return False, None

is_unpublished(item)

Check if the item is unpublished.

Source code in docs/rs-client-libraries/rs_workflows/utils/catalog.py
77
78
79
80
81
82
83
84
85
def is_unpublished(item: Item) -> tuple[bool, datetime | None]:
    """
    Check if the item is unpublished.
    """
    if unpublished_date_str := item.properties.get("unpublished"):
        unpublished_date = datetime.fromisoformat(unpublished_date_str.replace("Z", "+00:00"))
        return unpublished_date <= datetime.now(timezone.utc), unpublished_date

    return False, None

published_stac_item(flow_env, item, collection_name) async

" Push a STAC item into the rs-catalog.

Source code in docs/rs-client-libraries/rs_workflows/utils/catalog.py
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
@task
async def published_stac_item(flow_env: FlowEnv, item: Item, collection_name: str) -> None:
    """ "
    Push a STAC item into the rs-catalog.
    """
    logger = get_run_logger()
    logger.info(f"The STAC item 🧊 '{item.id}' will be published on the collection '{collection_name}'.")
    items_metadata: list[DprProcessedItemMetadata] = []
    publish_mapping: list[FlowGeneratedProduct] = []

    items_metadata.append(
        DprProcessedItemMetadata(
            output_product_id=item.id,
            product_type=item.properties.get("product:type"),
            stac_item=item,
        ),
    )

    publish_mapping.append(
        FlowGeneratedProduct(
            name=item.id,
            product_type=str(item.properties.get("product:type")),
            collection_name=collection_name,
        ),
    )

    logger.debug(f"items_metadata = {items_metadata}")
    logger.debug(f"publish_mapping = {publish_mapping}")
    await publish(
        flow_env.serialize(),
        publish_mapping,
        items_metadata,
    )