Skip to content
Merged
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
24 changes: 24 additions & 0 deletions backend/api/adapter/cognito_adapter.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import time
from typing import List, Dict, Optional

import boto3
Expand All @@ -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:
Expand All @@ -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"]
)
Expand All @@ -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()
)
Expand Down Expand Up @@ -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":
Expand All @@ -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":
Expand All @@ -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 = [
{
Expand Down
25 changes: 24 additions & 1 deletion backend/api/adapter/dynamodb_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down
22 changes: 11 additions & 11 deletions backend/api/application/services/data_service.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import uuid
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
from threading import Thread
from typing import List, Tuple
Expand Down Expand Up @@ -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:
"""
Expand Down
26 changes: 26 additions & 0 deletions backend/test/api/adapter/test_cognito_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
40 changes: 40 additions & 0 deletions backend/test/api/adapter/test_dynamodb_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
21 changes: 21 additions & 0 deletions backend/test/api/application/services/test_data_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
1 change: 1 addition & 0 deletions backend/test/api/test_entry.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
10 changes: 10 additions & 0 deletions docs/changelog/api.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
4 changes: 2 additions & 2 deletions frontend/src/pages/catalog/index.tsx
Original file line number Diff line number Diff line change
@@ -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'
Expand Down Expand Up @@ -207,7 +207,7 @@ function CatalogPage() {
<TableCell>
<LayerChip layer={d.layer} />
</TableCell>
<TableCell sx={{ fontFamily: 'monospace', fontSize: 11 }}>{formatDate(d.last_updated)}</TableCell>
<TableCell sx={{ fontFamily: 'monospace', fontSize: 11 }}>{formatTs(d.last_updated) ?? '—'}</TableCell>
<TableCell sx={{ fontFamily: 'monospace', fontSize: 11 }}>{d.last_uploaded_by ?? '—'}</TableCell>
</TableRow>
))
Expand Down
2 changes: 1 addition & 1 deletion frontend/src/service/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,7 @@ export type Dataset = {
dataset: string
version: number
sensitivity?: string
last_updated?: string
last_updated?: number
last_uploaded_by?: string
}

Expand Down
4 changes: 2 additions & 2 deletions infrastructure/modules/rapid/variables.tf
Original file line number Diff line number Diff line change
Expand Up @@ -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" {
Expand Down
Loading