diff --git a/backend/api/adapter/cognito_adapter.py b/backend/api/adapter/cognito_adapter.py index 5cfb8744..ee0988d2 100644 --- a/backend/api/adapter/cognito_adapter.py +++ b/backend/api/adapter/cognito_adapter.py @@ -1,3 +1,4 @@ +import time from typing import List, Dict, Optional import boto3 @@ -17,12 +18,17 @@ from api.domain.user import UserRequest, UserResponse +SUBJECTS_CACHE_TTL_SECONDS = 300 + + class CognitoAdapter: def __init__( self, cognito_client=boto3.client("cognito-idp", region_name=AWS_REGION) ): self.cognito_client = cognito_client self.placeholder_client_name = "string" + self._subjects_cache = None + self._subjects_cache_expiry = 0.0 def create_client_app(self, client_request: ClientRequest) -> ClientResponse: try: @@ -44,6 +50,7 @@ def create_client_app(self, client_request: ClientRequest) -> ClientResponse: PreventUserExistenceErrors="ENABLED", ) + self._invalidate_subjects_cache() return self._create_client_response( client_request, cognito_response["UserPoolClient"] ) @@ -64,6 +71,7 @@ def create_user(self, user_request: UserRequest) -> UserResponse: "EMAIL", ], ) + self._invalidate_subjects_cache() return self._create_user_response( cognito_response, user_request.get_permissions() ) @@ -92,6 +100,7 @@ def delete_client_app(self, client_id: str): self.cognito_client.delete_user_pool_client( UserPoolId=COGNITO_USER_POOL_ID, ClientId=client_id ) + self._invalidate_subjects_cache() except ClientError as error: AppLogger.info(f"Deleting client {client_id} failed with: {error.response}") if error.response["Error"]["Code"] == "ResourceNotFoundException": @@ -109,6 +118,7 @@ def delete_user(self, username: str): UserPoolId=COGNITO_USER_POOL_ID, Username=username, ) + self._invalidate_subjects_cache() except ClientError as error: AppLogger.info(f"Deleting user {username} failed with: {error.response}") if error.response["Error"]["Code"] == "UserNotFoundException": @@ -118,6 +128,20 @@ def delete_user(self, username: str): ) def get_all_subjects(self) -> List[Dict[str, Optional[str]]]: + now = time.monotonic() + if self._subjects_cache is not None and now < self._subjects_cache_expiry: + return self._subjects_cache + + subjects = self._fetch_all_subjects() + self._subjects_cache = subjects + self._subjects_cache_expiry = now + SUBJECTS_CACHE_TTL_SECONDS + return subjects + + def _invalidate_subjects_cache(self): + self._subjects_cache = None + self._subjects_cache_expiry = 0.0 + + def _fetch_all_subjects(self) -> List[Dict[str, Optional[str]]]: try: clients = [ { diff --git a/backend/api/adapter/dynamodb_adapter.py b/backend/api/adapter/dynamodb_adapter.py index 734fe7ba..a40ab6b2 100644 --- a/backend/api/adapter/dynamodb_adapter.py +++ b/backend/api/adapter/dynamodb_adapter.py @@ -2,7 +2,7 @@ from abc import ABC, abstractmethod from dataclasses import dataclass from functools import reduce -from typing import Any, Callable, Dict, List, Optional, Type +from typing import Any, Callable, Dict, List, Optional, Tuple, Type import boto3 from boto3.dynamodb.conditions import Attr, Key, Or @@ -416,6 +416,29 @@ def get_latest_successful_upload_job(self, dataset: Type[DatasetMetadata]) -> Op AppLogger.warning(f"Error fetching latest upload job for dataset: {error}") return None + def get_latest_successful_upload_jobs(self) -> Dict[Tuple[str, str, str], Dict]: + """ + Get the most recent successful upload job for every dataset in a single query. + Returns a dict keyed by (layer, domain, dataset) with the latest job details. + """ + try: + jobs = self.collect_all_items( + self.service_table.query, + KeyConditionExpression=Key("PK").eq("JOB"), + FilterExpression=Attr("Type").eq("UPLOAD") & Attr("Status").eq("SUCCESS"), + ) + except ClientError as error: + AppLogger.warning(f"Error fetching latest upload jobs: {error}") + return {} + + latest_by_dataset: Dict[Tuple[str, str, str], Dict] = {} + for job in jobs: + key = (job.get("Layer"), job.get("Domain"), job.get("Dataset")) + current = latest_by_dataset.get(key) + if current is None or job.get("CreatedAt", 0) > current.get("CreatedAt", 0): + latest_by_dataset[key] = job + return {key: self._map_job(job) for key, job in latest_by_dataset.items()} + def get_schema(self, dataset: Type[DatasetMetadata]) -> Optional[dict]: try: return self.schema_table.query( diff --git a/backend/api/application/services/data_service.py b/backend/api/application/services/data_service.py index fcc98a7c..baa25de1 100644 --- a/backend/api/application/services/data_service.py +++ b/backend/api/application/services/data_service.py @@ -1,5 +1,4 @@ import uuid -from concurrent.futures import ThreadPoolExecutor from pathlib import Path from threading import Thread from typing import List, Tuple @@ -193,20 +192,21 @@ def get_last_updated_time(self, metadata: DatasetMetadata) -> str: return last_updated or "Never updated" def enrich_datasets_for_ui(self, datasets: List[DatasetMetadata]) -> List[dict]: + latest_jobs = self.job_service.db_adapter.get_latest_successful_upload_jobs() + subject_names = { + subject["subject_id"]: subject["subject_name"] + for subject in self.subject_service.cognito_adapter.get_all_subjects() + } + def enrich(dataset: DatasetMetadata) -> dict: d = dataset.to_dict() - try: - d["last_updated"] = self.get_last_updated_time(dataset) - except Exception: - d["last_updated"] = None - try: - d["last_uploaded_by"] = self.get_last_uploader(dataset) - except Exception: - d["last_uploaded_by"] = None + job = latest_jobs.get((dataset.layer, dataset.domain, dataset.dataset)) + d["last_updated"] = job.get("createdat") if job else None + subject_id = job.get("sk2") if job else None + d["last_uploaded_by"] = subject_names.get(subject_id) if subject_id else None return d - with ThreadPoolExecutor(max_workers=10) as pool: - return list(pool.map(enrich, datasets)) + return [enrich(dataset) for dataset in datasets] def get_last_uploader(self, metadata: DatasetMetadata) -> str: """ diff --git a/backend/test/api/adapter/test_cognito_adapter.py b/backend/test/api/adapter/test_cognito_adapter.py index 1750c92f..d95798f1 100644 --- a/backend/test/api/adapter/test_cognito_adapter.py +++ b/backend/test/api/adapter/test_cognito_adapter.py @@ -472,6 +472,32 @@ def test_gets_all_subjects(self): assert result == expected + def test_caches_subjects_between_calls(self): + self.cognito_boto_client.get_paginator.return_value.paginate.side_effect = [ + [{"UserPoolClients": [{"ClientId": "c1", "ClientName": "client_1"}]}], + [{"Users": []}], + ] + + first = self.cognito_adapter.get_all_subjects() + second = self.cognito_adapter.get_all_subjects() + + assert first == second + assert self.cognito_boto_client.get_paginator.return_value.paginate.call_count == 2 + + def test_invalidates_subjects_cache_on_delete_user(self): + self.cognito_boto_client.get_paginator.return_value.paginate.side_effect = [ + [{"UserPoolClients": [{"ClientId": "c1", "ClientName": "client_1"}]}], + [{"Users": []}], + [{"UserPoolClients": [{"ClientId": "c1", "ClientName": "client_1"}]}], + [{"Users": []}], + ] + + self.cognito_adapter.get_all_subjects() + self.cognito_adapter.delete_user("some-user") + self.cognito_adapter.get_all_subjects() + + assert self.cognito_boto_client.get_paginator.return_value.paginate.call_count == 4 + def test_raises_error_when_listing_clients_fails(self): self.cognito_boto_client.get_paginator.return_value.paginate.side_effect = ( ClientError( diff --git a/backend/test/api/adapter/test_dynamodb_adapter.py b/backend/test/api/adapter/test_dynamodb_adapter.py index 97ee885c..c0c65ba4 100644 --- a/backend/test/api/adapter/test_dynamodb_adapter.py +++ b/backend/test/api/adapter/test_dynamodb_adapter.py @@ -756,6 +756,46 @@ def test_get_jobs(self, mock_time): IndexName="JOB_SUBJECT_ID", ) + def test_get_latest_successful_upload_jobs(self): + self.service_table.query.return_value = { + "Items": [ + { + "PK": "JOB", + "SK": "job-old", + "SK2": "subject-1", + "Type": "UPLOAD", + "Status": "SUCCESS", + "Layer": "raw", + "Domain": "domain1", + "Dataset": "dataset1", + "CreatedAt": 1000, + }, + { + "PK": "JOB", + "SK": "job-new", + "SK2": "subject-2", + "Type": "UPLOAD", + "Status": "SUCCESS", + "Layer": "raw", + "Domain": "domain1", + "Dataset": "dataset1", + "CreatedAt": 2000, + }, + ], + "Count": 2, + } + + result = self.dynamo_adapter.get_latest_successful_upload_jobs() + + assert set(result.keys()) == {("raw", "domain1", "dataset1")} + latest = result[("raw", "domain1", "dataset1")] + assert latest["createdat"] == 2000 + assert latest["sk2"] == "subject-2" + self.service_table.query.assert_called_once_with( + KeyConditionExpression=Key("PK").eq("JOB"), + FilterExpression=Attr("Type").eq("UPLOAD") & Attr("Status").eq("SUCCESS"), + ) + @patch("api.adapter.dynamodb_adapter.time") def test_get_jobs_for_no_jobs_returned(self, mock_time): mock_time.time.return_value = 19821 diff --git a/backend/test/api/application/services/test_data_service.py b/backend/test/api/application/services/test_data_service.py index d20ee636..0876d439 100644 --- a/backend/test/api/application/services/test_data_service.py +++ b/backend/test/api/application/services/test_data_service.py @@ -663,6 +663,27 @@ def test_generates_raw_file_identifier(self): pattern = "[\\d\\w]{8}-[\\d\\w]{4}-[\\d\\w]{4}-[\\d\\w]{4}-[\\d\\w]{12}" assert re.match(pattern, filename) + def test_enrich_datasets_for_ui_uses_bulk_jobs_and_subjects(self): + datasets = [ + DatasetMetadata("raw", "domain1", "dataset1", 1), + DatasetMetadata("raw", "domain2", "dataset2", 1), + ] + self.job_service.db_adapter.get_latest_successful_upload_jobs.return_value = { + ("raw", "domain1", "dataset1"): {"sk2": "subject-123", "createdat": 1700000000}, + } + self.subject_service.cognito_adapter.get_all_subjects.return_value = [ + {"subject_id": "subject-123", "subject_name": "test_user"}, + ] + + result = self.data_service.enrich_datasets_for_ui(datasets) + + self.job_service.db_adapter.get_latest_successful_upload_jobs.assert_called_once_with() + self.subject_service.cognito_adapter.get_all_subjects.assert_called_once_with() + assert result[0]["last_updated"] == 1700000000 + assert result[0]["last_uploaded_by"] == "test_user" + assert result[1]["last_updated"] is None + assert result[1]["last_uploaded_by"] is None + class TestQueryDataset: def setup_method(self): diff --git a/backend/test/api/test_entry.py b/backend/test/api/test_entry.py index ca50182f..82038a93 100644 --- a/backend/test/api/test_entry.py +++ b/backend/test/api/test_entry.py @@ -92,6 +92,7 @@ def test_gets_datasets_for_ui_read( mock_get_authorised_datasets.assert_called_once_with(subject_id, Action.READ) assert response.status_code == 200 + mock_data_service.enrich_datasets_for_ui.assert_called_once() class TestMethodsUI(BaseClientTest): diff --git a/docs/changelog/api.md b/docs/changelog/api.md index aa76de62..38942ff7 100644 --- a/docs/changelog/api.md +++ b/docs/changelog/api.md @@ -1,5 +1,15 @@ # API Changelog +## v8.0.2 - _2026-06-04_ + +See [v8.0.2] changes + +### Features + +- Cached enhanced api calls + +[v8.0.2]: https://github.com/no10ds/rapid/compare/v8.0.0...v8.0.2 + ## v8.0.1 - _2026-06-04_ See [v8.0.1] changes diff --git a/frontend/src/pages/catalog/index.tsx b/frontend/src/pages/catalog/index.tsx index 2ad936bd..669e877f 100644 --- a/frontend/src/pages/catalog/index.tsx +++ b/frontend/src/pages/catalog/index.tsx @@ -1,7 +1,7 @@ import AccountLayout from '@/components/Layout/AccountLayout' import ErrorCard from '@/components/ErrorCard/ErrorCard' import LayerChip from '@/components/Chip/LayerChip' -import { formatDate } from '@/utils/date' +import { formatTs } from '@/utils/date' import { sortByString } from '@/utils/sort' import { getDatasetsUi } from '@/service' import { Dataset } from '@/service/types' @@ -207,7 +207,7 @@ function CatalogPage() { - {formatDate(d.last_updated)} + {formatTs(d.last_updated) ?? '—'} {d.last_uploaded_by ?? '—'} )) diff --git a/frontend/src/service/types.ts b/frontend/src/service/types.ts index 0cf7c051..bb745cea 100644 --- a/frontend/src/service/types.ts +++ b/frontend/src/service/types.ts @@ -166,7 +166,7 @@ export type Dataset = { dataset: string version: number sensitivity?: string - last_updated?: string + last_updated?: number last_uploaded_by?: string } diff --git a/infrastructure/modules/rapid/variables.tf b/infrastructure/modules/rapid/variables.tf index 9eece2ee..59ebfb11 100644 --- a/infrastructure/modules/rapid/variables.tf +++ b/infrastructure/modules/rapid/variables.tf @@ -13,13 +13,13 @@ variable "app-replica-count-max" { variable "application_version" { type = string description = "The version number for the application image (e.g.: v1.0.4, v1.0.x-latest, etc.)" - default = "v8.0.1" + default = "v8.0.2" } variable "ui_version" { type = string description = "The version number for the static ui (e.g.: v1.0.0, etc.)" - default = "v8.0.1" + default = "v8.0.2" } variable "catalog_disabled" {