Skip to content

Commit 05490d1

Browse files
authored
Merge pull request #81 from govtechmy/feat/SSD-1229-fileVersion-field-in-DatasetStatus
Feat/ssd 1229 file version field in dataset status
2 parents b5ab815 + 1bfe514 commit 05490d1

6 files changed

Lines changed: 84 additions & 20 deletions

File tree

src/core/gsheet.py

Lines changed: 52 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,59 @@
11
import requests
2+
import logging
3+
import re
4+
import urllib.parse
25

6+
logger = logging.getLogger(__name__)
37

4-
def fetch_csv_data(gsheet_id: str, gid: str) -> bytes:
8+
def _extract_filename(content_disp: str) -> str | None:
9+
if not content_disp:
10+
return None
11+
12+
# RFC 5987 format: filename*=UTF-8''encoded_name.csv
13+
match_star = re.search(r"filename\*=UTF-8''([^;]+)(?:;|$)", content_disp, re.IGNORECASE,)
14+
if match_star:
15+
filename = urllib.parse.unquote(match_star.group(1))
16+
logger.info("Extracted filename: %s", filename)
17+
return filename
18+
19+
# Standard format: filename="file.csv" OR filename=file.csv
20+
match = re.search(r'filename=(?:"([^"]*)"|([^;]+))(?:;|$)', content_disp, re.IGNORECASE,)
21+
if match:
22+
filename = match.group(1) if match.group(1) is not None else match.group(2).strip()
23+
logger.info("Extracted filename (standard): %s", filename)
24+
return filename
25+
26+
logger.warning("Failed to extract filename from Content-Disposition: %s", content_disp)
27+
return None
28+
29+
def _extract_file_version(file_name: str | None) -> str | None:
30+
if not file_name:
31+
return None
32+
33+
name = file_name.strip().rsplit(".", 1)[0]
34+
name = name.split(" - ")[0]
35+
parts = name.split("_")
36+
37+
if len(parts) >= 2:
38+
file_version = parts[-1]
39+
logger.info("Extracted fileVersion: %s", file_version)
40+
return file_version
41+
42+
logger.warning("Failed to extract fileVersion from filename: %s", file_name)
43+
return None
44+
45+
def fetch_csv_data(gsheet_id: str, gid: str) -> tuple[bytes, str | None]:
546
url = f"https://docs.google.com/spreadsheets/d/{gsheet_id}/export?format=csv&gid={gid}"
6-
print(f'Fetching CSV data from URL: {url}')
47+
logger.info("Fetching CSV data from URL: %s", url)
748

849
response = requests.get(url)
950
response.raise_for_status()
10-
return response.content
51+
52+
content_type = response.headers.get("Content-Type", "")
53+
if "text/csv" not in content_type:
54+
logger.warning("Google Sheet not accessible as CSV. Content-Type: %s", content_type)
55+
56+
content_disp = response.headers.get("Content-Disposition", "")
57+
filename = _extract_filename(content_disp)
58+
59+
return response.content, filename

src/core/s3.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,16 +6,18 @@
66
from botocore.exceptions import ClientError, ResponseStreamingError
77

88
from src.core.aws import get_s3_client, get_s3_bucket_name
9+
from src.core.gsheet import _extract_file_version
910

1011
s3 = get_s3_client()
1112
logger = logging.getLogger(__name__)
1213

13-
def _upload_to_s3(csv_bytes: bytes, bucket: str, prefix: str) -> str:
14+
def _upload_to_s3(csv_bytes: bytes, bucket: str, prefix: str, source_filename: str | None = None,) -> str:
1415
if not bucket:
1516
bucket = get_s3_bucket_name()
1617

1718
timestamp = int(time.time())
18-
s3_key = f"{prefix}/{timestamp}.csv"
19+
version = _extract_file_version(source_filename) or "unknown"
20+
s3_key = f"{prefix}/{version}/{timestamp}.csv"
1921

2022
s3.put_object(
2123
Bucket=bucket,

src/models/dataset_status.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ class DatasetStatus(BaseModel):
2121

2222
id: str = Field(..., alias="_id", description="Dataset identifier (e.g. sekolah, institusi, analitik)")
2323
lastUpdatedAt: datetime = Field(default_factory=_utc_now, description="UTC timestamp of last successful ingestion")
24+
fileVersion: str | None = Field(default=None, description="Cleaned version of the file name e.g SenaraiSekolahWeb_Mac2026")
2425

2526
def to_document(self) -> dict:
2627
return self.model_dump(by_alias=True)

src/pipeline/dataset_status.py

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,24 +5,29 @@
55
from src.config.settings import Settings
66
from src.core.time import _utc_now
77
from src.utils.db.get_db_collection import get_db_collection
8+
from typing import Optional
89

910
logger = logging.getLogger(__name__)
1011

1112

12-
def upsert_dataset_status(dataset_name: str, settings: Settings) -> None:
13+
def upsert_dataset_status(dataset_name: str, settings: Settings, file_version: Optional[str]) -> None:
1314
"""Upsert the lastUpdatedAt timestamp for a dataset into DatasetStatus.
1415
1516
The document shape is:
16-
{"_id": <dataset_name>, "lastUpdatedAt": <UTC datetime>}
17+
{"_id": <dataset_name>, "lastUpdatedAt": <UTC datetime>, "fileVersion": <file_version>}
18+
fileVersion is only set if provided (not None).
1719
"""
1820
if not dataset_name:
1921
raise ValueError("dataset_name must be a non-empty string")
2022

2123
try:
2224
collection = get_db_collection(settings, name=settings.dataset_status_collection)
25+
update_fields = {"lastUpdatedAt": _utc_now()}
26+
if file_version is not None:
27+
update_fields["fileVersion"] = file_version
2328
collection.update_one(
2429
{"_id": dataset_name},
25-
{"$set": {"lastUpdatedAt": _utc_now()}},
30+
{"$set": update_fields},
2631
upsert=True,
2732
)
2833
except Exception:

src/pipeline/ingestion.py

Lines changed: 18 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717

1818
from src.config import Settings, get_settings
1919
from src.models import Sekolah
20-
from src.core.gsheet import fetch_csv_data
20+
from src.core.gsheet import fetch_csv_data, _extract_file_version
2121
from src.core.s3 import (_upload_to_s3, _latest_csv_from_s3, _read_csv_from_s3)
2222
from src.models.sekolah import SekolahStatus
2323
from src.pipeline.status_sync import sync_entiti_statuses
@@ -91,22 +91,22 @@ def _read_csv(path: str) -> Iterable[Dict[str, Any]]:
9191
def _read_google_sheet(sheet_id: str, gid: str) -> Iterable[Dict[str, Any]]:
9292
logger.info("Scraping Google Sheet directly (sheet_id=%s, gid=%s)", sheet_id, gid,)
9393

94-
csv_bytes = fetch_csv_data(sheet_id, gid)
94+
csv_bytes, _ = fetch_csv_data(sheet_id, gid)
9595
df = pd.read_csv(io.BytesIO(csv_bytes), dtype=str).fillna("")
9696

9797
logger.info("Google Sheet loaded: %d rows, %d columns", df.shape[0], df.shape[1])
9898
return df.to_dict(orient="records")
9999

100100

101-
def _load_rows(settings: Settings) -> Iterable[Dict[str, Any]]:
102-
csv_bytes = fetch_csv_data(settings.gsheet_id, settings.gsheet_gid)
101+
def _load_rows(settings: Settings) -> tuple[Iterable[Dict[str, Any]], str | None]:
102+
csv_bytes, file_name = fetch_csv_data(settings.gsheet_id, settings.gsheet_gid)
103103
logger.info("Uploading CSV data to S3 bucket %s", settings.s3_bucket_dataproc)
104-
s3_key = _upload_to_s3(csv_bytes, settings.s3_bucket_dataproc, settings.s3_prefix_sekolah)
104+
s3_key = _upload_to_s3(csv_bytes, settings.s3_bucket_dataproc, settings.s3_prefix_sekolah, file_name)
105105
logger.info("CSV uploaded to S3 at key: %s", s3_key)
106106

107107
df = _read_csv_from_s3(settings.s3_bucket_dataproc, s3_key)
108108
logger.info("CSV loaded from S3: %d rows, %d columns", df.shape[0], df.shape[1])
109-
return df.to_dict(orient="records")
109+
return df.to_dict(orient="records"), file_name
110110

111111

112112
def _chunked(rows: Iterable[Dict[str, Any]], size: int) -> Iterator[list[Dict[str, Any]]]:
@@ -142,13 +142,14 @@ def _format_validation_messages(exc: ValidationError) -> list[str]:
142142
def _collect_documents(
143143
settings: Settings,
144144
kodSekolah_madani: Set[str],
145-
) -> tuple[list[dict[str, Any]], list[dict[str, Any]], int, set[Any]]:
145+
) -> tuple[list[dict[str, Any]], list[dict[str, Any]], int, set[Any], str | None]:
146146
documents: list[dict[str, Any]] = []
147147
errors: list[dict[str, Any]] = []
148148
total = 0
149149
present_identifiers: set[Any] = set()
150150

151-
for index, row in enumerate(_load_rows(settings), start=1):
151+
rows, file_name = _load_rows(settings)
152+
for index, row in enumerate(rows, start=1):
152153
total += 1
153154

154155
raw_kod = str(row.get("KODSEKOLAH", "")).strip()
@@ -169,7 +170,7 @@ def _collect_documents(
169170
document["isSekolahAngkatMADANI"] = canonical_kod in kodSekolah_madani
170171
documents.append(document)
171172

172-
return documents, errors, total, present_identifiers
173+
return documents, errors, total, present_identifiers, file_name
173174

174175

175176
def _load_kodSekolah_madani(database, settings: Settings) -> Set[str]:
@@ -326,7 +327,7 @@ def run(settings: Settings) -> dict[str, Any]:
326327
sekolah_collection = database[Sekolah.collection_name]
327328
entiti_collection = database[settings.entiti_sekolah_collection]
328329
kodSekolah_madani = _load_kodSekolah_madani(database, settings)
329-
documents, errors, total, present_identifiers = _collect_documents(settings, kodSekolah_madani)
330+
documents, errors, total, present_identifiers, file_name = _collect_documents(settings, kodSekolah_madani)
330331

331332
for document in documents:
332333
# All schools present in raw file are ACTIVE
@@ -378,7 +379,13 @@ def run(settings: Settings) -> dict[str, Any]:
378379
"inactivated": inactivated,
379380
"entiti_synced": entiti_synced,
380381
}
381-
upsert_dataset_status("sekolah", settings)
382+
383+
file_version = _extract_file_version(file_name)
384+
385+
if not file_version:
386+
logger.warning("File name not found or failed to extract fileVersion. filename=%s", file_name)
387+
388+
upsert_dataset_status("sekolah", settings, file_version)
382389
return summary
383390

384391
def run_with_overrides(**overrides: Any) -> dict[str, Any]:

src/pipeline/institusi.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ def _load_rows(settings: Settings) -> Iterable[Dict[str, Any]]:
7070
if not settings.institusi_gsheet_id or not settings.institusi_gsheet_gid:
7171
raise RuntimeError("GSHEET_ID and INSTITUSI_GSHEET_GID must be set")
7272

73-
csv_bytes = fetch_csv_data(settings.institusi_gsheet_id, settings.institusi_gsheet_gid)
73+
csv_bytes, _ = fetch_csv_data(settings.institusi_gsheet_id, settings.institusi_gsheet_gid)
7474
logger.info("Uploading Institusi CSV data to S3 bucket %s", settings.s3_bucket_dataproc)
7575
s3_key = _upload_to_s3(csv_bytes, settings.s3_bucket_dataproc, settings.s3_prefix_institusi)
7676
logger.info("Institusi CSV uploaded to S3 at key: %s", s3_key)

0 commit comments

Comments
 (0)