github-actions[bot] commited on
Commit
958cb1b
·
1 Parent(s): cec9c04

deploy: backend from a793581

Browse files
Files changed (1) hide show
  1. scripts/ingest_curriculum.py +306 -129
scripts/ingest_curriculum.py CHANGED
@@ -1,159 +1,336 @@
1
  from __future__ import annotations
2
 
3
- import argparse
4
  import hashlib
5
  import json
6
- import logging
7
  import os
 
 
8
  import sys
 
9
  from pathlib import Path
10
- from typing import Any, Dict, List
11
 
12
- sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
 
 
 
 
13
 
14
- from rag.vectorstore_loader import (
15
- get_vectorstore_components,
16
- reset_vectorstore_singleton,
17
- )
 
 
 
 
 
 
 
18
 
19
- logger = logging.getLogger(__name__)
20
 
 
 
21
 
22
- def _resolve_data_dir(raw: str | None) -> Path:
23
- if raw:
24
- p = Path(raw)
25
- if p.is_absolute():
26
- return p
27
- p = Path.cwd() / raw
28
- if p.exists():
29
- return p
30
- default = Path(__file__).resolve().parents[1] / "datasets"
31
- return default
32
 
 
 
 
 
33
 
34
- def _iter_json_files(data_dir: Path):
35
- for file in sorted(data_dir.rglob("*")):
36
- if file.suffix not in {".json", ".jsonl"}:
37
- continue
38
- yield file
 
 
 
 
 
 
 
39
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
40
 
41
- def _load_records(file_path: Path) -> List[Dict[str, Any]]:
42
- records: List[Dict[str, Any]] = []
43
  try:
44
- raw = file_path.read_text(encoding="utf-8").strip()
45
- if file_path.suffix == ".jsonl":
46
- for lineno, line in enumerate(raw.splitlines(), start=1):
47
- line = line.strip()
48
- if not line:
49
- continue
50
- try:
51
- records.append(json.loads(line))
52
- except json.JSONDecodeError:
53
- logger.warning("Skipping malformed JSONL line %s:%d", file_path.name, lineno)
54
  else:
55
- parsed = json.loads(raw)
56
- if isinstance(parsed, list):
57
- records.extend(parsed)
58
- elif isinstance(parsed, dict):
59
- records.append(parsed)
60
- except Exception as exc:
61
- logger.warning("Failed to parse %s: %s", file_path.name, exc)
62
- return records
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
63
 
64
 
65
- def _build_id(source_file: str, page: int, content: str) -> str:
66
- key = f"{source_file}::{page}::{content[:120]}"
67
- return hashlib.sha256(key.encode()).hexdigest()[:40]
 
 
 
68
 
 
 
 
 
 
 
 
 
 
 
 
 
 
69
 
70
- def main() -> None:
71
- parser = argparse.ArgumentParser(description="Ingest DepEd SHS curriculum JSON/JSONL into ChromaDB")
72
- parser.add_argument("--data-dir", default=None, help="Directory containing .json/.jsonl files")
73
- parser.add_argument("--reset", action="store_true", help="Reset the vectorstore singleton before ingestion")
74
- args = parser.parse_args()
75
 
76
- data_dir = _resolve_data_dir(args.data_dir)
77
- logger.info("Ingesting from: %s", data_dir)
 
 
 
78
 
79
- if args.reset:
80
- reset_vectorstore_singleton()
81
- _, collection, _ = get_vectorstore_components()
 
 
 
 
 
 
 
 
 
 
82
  try:
83
- collection.delete(ids=collection.get(include=[])["ids"])
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
84
  except Exception:
85
  pass
86
- reset_vectorstore_singleton()
87
-
88
- total_processed = 0
89
- total_upserted = 0
90
- total_errors = 0
91
-
92
- _, collection, embedder = get_vectorstore_components()
93
-
94
- for file_path in _iter_json_files(data_dir):
95
- records = _load_records(file_path)
96
- documents: List[str] = []
97
- metadatas: List[Dict[str, Any]] = []
98
- ids: List[str] = []
99
- embeddings_list: List[List[float]] = []
100
-
101
- for record in records:
102
- total_processed += 1
103
- content = str(record.get("content") or "").strip()
104
- if not content:
105
- logger.debug("Skipping empty content in %s", file_path.name)
106
- continue
107
-
108
- try:
109
- subject = str(record.get("subject") or "unknown")
110
- quarter = int(record.get("quarter") or 0)
111
- page = int(record.get("page") or 0)
112
- content_domain = str(record.get("content_domain") or "unknown")
113
- chunk_type = str(record.get("chunk_type") or "unknown")
114
- source_file = str(record.get("source_file") or file_path.name)
115
-
116
- embedding = embedder.encode(content).tolist()
117
- chunk_id = _build_id(source_file, page, content)
118
-
119
- metadata = {
120
- "subject": subject,
121
- "quarter": quarter,
122
- "content_domain": content_domain,
123
- "chunk_type": chunk_type,
124
- "source_file": source_file,
125
- "page": page,
126
- }
127
-
128
- documents.append(content)
129
- metadatas.append(metadata)
130
- ids.append(chunk_id)
131
- embeddings_list.append(embedding)
132
-
133
- except Exception as exc:
134
- total_errors += 1
135
- logger.warning("Error processing record in %s: %s", file_path.name, exc)
136
-
137
- if documents:
138
- try:
139
- collection.upsert(
140
- ids=ids,
141
- documents=documents,
142
- metadatas=metadatas,
143
- embeddings=embeddings_list,
144
- )
145
- total_upserted += len(documents)
146
- logger.info("Upserted %d chunks from %s", len(documents), file_path.name)
147
- except Exception as exc:
148
- total_errors += len(documents)
149
- logger.warning("Failed to upsert batch from %s: %s", file_path.name, exc)
150
-
151
- print(f"=== Ingestion Summary ===")
152
- print(f"Total records processed: {total_processed}")
153
- print(f"Total chunks upserted: {total_upserted}")
154
- print(f"Total errors: {total_errors}")
155
 
156
 
157
  if __name__ == "__main__":
158
- logging.basicConfig(level=logging.INFO)
159
- main()
 
1
  from __future__ import annotations
2
 
 
3
  import hashlib
4
  import json
 
5
  import os
6
+ import re
7
+ from collections import Counter
8
  import sys
9
+ from datetime import datetime, timezone
10
  from pathlib import Path
11
+ from typing import Dict, Iterable, List
12
 
13
+ BASE_DIR = Path(__file__).resolve().parents[1]
14
+ if str(BASE_DIR) not in sys.path:
15
+ sys.path.insert(0, str(BASE_DIR))
16
+ if str(BASE_DIR / "backend") not in sys.path:
17
+ sys.path.insert(0, str(BASE_DIR / "backend"))
18
 
19
+ try:
20
+ from backend.rag.liteparse_utils import extract_text
21
+ except ImportError:
22
+ from rag.liteparse_utils import extract_text
23
+ CURRICULUM_DIR = Path(os.getenv("CURRICULUM_DIR", BASE_DIR / "datasets" / "curriculum"))
24
+ VECTORSTORE_DIR = Path(os.getenv("VECTORSTORE_DIR", BASE_DIR / "datasets" / "vectorstore"))
25
+ COLLECTION_NAME = "curriculum_chunks"
26
+ EMBED_MODEL_NAME = os.getenv("EMBEDDING_MODEL", "BAAI/bge-small-en-v1.5")
27
+ CURRICULUM_SOURCE_REPO_ID = os.getenv("CURRICULUM_SOURCE_REPO_ID", "").strip()
28
+ CURRICULUM_SOURCE_REPO_TYPE = os.getenv("CURRICULUM_SOURCE_REPO_TYPE", "dataset").strip() or "dataset"
29
+ CURRICULUM_SOURCE_REVISION = os.getenv("CURRICULUM_SOURCE_REVISION", "main").strip() or "main"
30
 
 
31
 
32
+ def _norm(text: str) -> str:
33
+ return re.sub(r"\s+", " ", text.strip().lower())
34
 
 
 
 
 
 
 
 
 
 
 
35
 
36
+ def infer_metadata(path: Path, text: str = "") -> Dict[str, object]:
37
+ parts = [part.lower() for part in path.parts]
38
+ joined = " ".join(parts)
39
+ stem_lower = path.stem.lower()
40
 
41
+ clean_joined = re.sub(r"problems?", "", joined)
42
+ clean_stem = re.sub(r"problems?", "", stem_lower)
43
+ if "stat" in clean_joined or "prob" in clean_joined or "stat" in clean_stem or "prob" in clean_stem:
44
+ subject = "statistics_and_probability"
45
+ elif "finite math 1" in joined or "finite_mathematics_1" in joined:
46
+ subject = "finite_mathematics_1"
47
+ elif "finite math 2" in joined or "finite_mathematics_2" in joined:
48
+ subject = "finite_mathematics_2"
49
+ elif "finite" in joined or "finite" in stem_lower:
50
+ subject = "finite_mathematics"
51
+ else:
52
+ subject = "general_mathematics"
53
 
54
+ name_and_parts = f"{joined} {stem_lower}"
55
+ quarter_match = re.search(
56
+ r"quarter\s*([1-4])|[_\-\b\s]q([1-4])[_\-\b\s.]|^q([1-4])[_\-\b\s.]|[\b_]q([1-4])[\b_]",
57
+ name_and_parts,
58
+ )
59
+ quarter = 0
60
+ if quarter_match:
61
+ quarter = int(next(group for group in quarter_match.groups() if group))
62
+ else:
63
+ mod_match = re.search(r"module\s*([1-4])|[\b_]mod([1-4])[\b_]", name_and_parts)
64
+ if mod_match:
65
+ quarter = int(next(group for group in mod_match.groups() if group))
66
+
67
+ resource_type = "learning_activity_sheet" if "learning activity" in joined or "las" in stem_lower else "lesson_exemplar"
68
+ if "curriculum" in joined:
69
+ resource_type = "curriculum_guide"
70
+ elif "budget" in joined:
71
+ resource_type = "budget_of_work"
72
 
 
 
73
  try:
74
+ storage_path = path.resolve().relative_to((BASE_DIR / "datasets").resolve()).as_posix()
75
+ except ValueError:
76
+ norm_posix = path.as_posix()
77
+ if "datasets/" in norm_posix:
78
+ storage_path = norm_posix.split("datasets/", 1)[1]
79
+ elif "curriculum" in norm_posix:
80
+ storage_path = "curriculum" + norm_posix.split("curriculum", 1)[1]
 
 
 
81
  else:
82
+ storage_path = f"curriculum/{path.name}"
83
+
84
+ if subject == "statistics_and_probability":
85
+ content_domain = "statistics"
86
+ elif "business" in joined or "bus_math" in joined:
87
+ content_domain = "business_math"
88
+ else:
89
+ content_domain = "general"
90
+
91
+ return {
92
+ "subject": subject,
93
+ "quarter": quarter,
94
+ "content_domain": content_domain,
95
+ "resource_type": resource_type,
96
+ "source_file": path.name,
97
+ "source_path": path.as_posix(),
98
+ "storage_path": storage_path,
99
+ }
100
+
101
+
102
+ def discover_curriculum_files(data_dir: Path) -> List[Path]:
103
+ """Prefer LiteParse-generated Markdown and use PDFs only without a twin."""
104
+ markdown_files = sorted(
105
+ file for file in data_dir.rglob("*.md") if file.name.lower() != "readme.md"
106
+ )
107
+ markdown_stems = {file.stem.lower() for file in markdown_files}
108
+ pdf_files = sorted(file for file in data_dir.rglob("*.pdf") if file.stem.lower() not in markdown_stems)
109
+ return markdown_files + pdf_files
110
 
111
 
112
+ def _resolve_source_dir() -> Path:
113
+ if CURRICULUM_DIR.exists():
114
+ return CURRICULUM_DIR
115
+ if not CURRICULUM_SOURCE_REPO_ID:
116
+ raise SystemExit(f"Missing curriculum directory: {CURRICULUM_DIR}")
117
+ from huggingface_hub import snapshot_download
118
 
119
+ source_dir = Path(snapshot_download(
120
+ repo_id=CURRICULUM_SOURCE_REPO_ID,
121
+ repo_type=CURRICULUM_SOURCE_REPO_TYPE,
122
+ revision=CURRICULUM_SOURCE_REVISION,
123
+ allow_patterns=["*.pdf", "**/*.pdf", "*.md", "**/*.md"],
124
+ ))
125
+ CURRICULUM_DIR.mkdir(parents=True, exist_ok=True)
126
+ for source_file in source_dir.rglob("*"):
127
+ if source_file.is_file() and source_file.suffix.lower() in {".pdf", ".md"}:
128
+ target = CURRICULUM_DIR / source_file.relative_to(source_dir)
129
+ target.parent.mkdir(parents=True, exist_ok=True)
130
+ target.write_bytes(source_file.read_bytes())
131
+ return CURRICULUM_DIR
132
 
 
 
 
 
 
133
 
134
+ def _clean_text(text: str) -> str:
135
+ if not text:
136
+ return ""
137
+ # Strip lone surrogate characters that invalidate UTF-8 encoding in Rust tokenizers
138
+ return text.encode("utf-8", "ignore").decode("utf-8")
139
 
140
+
141
+ def chunk_text(text: str) -> List[str]:
142
+ from langchain_text_splitters import RecursiveCharacterTextSplitter
143
+
144
+ cleaned = _clean_text(text)
145
+ splitter = RecursiveCharacterTextSplitter(chunk_size=2000, chunk_overlap=200, separators=["\n\n", "\n", ". ", " ", ""])
146
+ return [_clean_text(chunk.strip()) for chunk in splitter.split_text(cleaned) if chunk.strip()]
147
+
148
+
149
+ def _read_source(path: Path) -> str:
150
+ if path.suffix.lower() == ".md":
151
+ raw = path.read_text(encoding="utf-8", errors="ignore")
152
+ else:
153
  try:
154
+ import pypdf
155
+ reader = pypdf.PdfReader(path)
156
+ extracted = "\n".join(page.extract_text() or "" for page in reader.pages)
157
+ if len(extracted.strip()) >= 100:
158
+ raw = extracted
159
+ else:
160
+ raw = extract_text(path)
161
+ except Exception:
162
+ raw = extract_text(path)
163
+ return _clean_text(raw)
164
+
165
+
166
+ def build_documents(data_dir: Path) -> tuple[List[str], List[Dict[str, object]], List[str]]:
167
+ documents: List[str] = []
168
+ metadatas: List[Dict[str, object]] = []
169
+ ids: List[str] = []
170
+ for source_file in discover_curriculum_files(data_dir):
171
+ text = _read_source(source_file)
172
+ metadata = infer_metadata(source_file, text)
173
+ storage_path = str(metadata.get("storage_path") or source_file.stem)
174
+ path_hash = hashlib.md5(storage_path.encode("utf-8")).hexdigest()[:8]
175
+ for index, chunk in enumerate(chunk_text(text), start=1):
176
+ documents.append(chunk)
177
+ metadatas.append({**metadata, "chunk_index": index})
178
+ ids.append(f"{path_hash}-{source_file.stem}-{index}")
179
+ return documents, metadatas, ids
180
+
181
+
182
+ def main(argv: List[str] | None = None) -> None:
183
+ import argparse
184
+
185
+ parser = argparse.ArgumentParser(description="Ingest the SSHS curriculum corpus into ChromaDB")
186
+ parser.add_argument("--data-dir", type=Path, default=None)
187
+ parser.add_argument("--vectorstore-dir", type=Path, default=None)
188
+ parser.add_argument(
189
+ "--dry-run",
190
+ action="store_true",
191
+ help="Print discovered files, inferred metadata (subject, quarter, storage_path), and estimated chunks without updating Chroma",
192
+ )
193
+ args = parser.parse_args(argv)
194
+
195
+ data_dir = args.data_dir or _resolve_source_dir()
196
+ vectorstore_dir = args.vectorstore_dir or VECTORSTORE_DIR
197
+ files = discover_curriculum_files(data_dir)
198
+ if not files:
199
+ raise SystemExit(f"No Markdown or PDF curriculum files found in {data_dir}")
200
+
201
+ if args.dry_run:
202
+ print(f"=== DRY RUN: Ingesting curriculum from {data_dir} ===")
203
+ print(f"Discovered {len(files)} files:\n")
204
+ total_estimated_chunks = 0
205
+ chunks_by_subject: Dict[str, int] = Counter()
206
+ for idx, source_file in enumerate(files, start=1):
207
+ text = _read_source(source_file)
208
+ metadata = infer_metadata(source_file, text)
209
+ chunks = chunk_text(text)
210
+ chunk_count = len(chunks)
211
+ total_estimated_chunks += chunk_count
212
+ subj = str(metadata["subject"])
213
+ chunks_by_subject[subj] += chunk_count
214
+ print(
215
+ f"[{idx:02d}/{len(files):02d}] {source_file.name}\n"
216
+ f" storage_path: {metadata['storage_path']}\n"
217
+ f" subject: {metadata['subject']}\n"
218
+ f" quarter: {metadata['quarter']}\n"
219
+ f" content_domain: {metadata['content_domain']}\n"
220
+ f" chunks: {chunk_count}\n"
221
+ )
222
+ print("=== DRY RUN SUMMARY ===")
223
+ print(f"Total discovered files: {len(files)}")
224
+ print(f"Total estimated chunks: {total_estimated_chunks}")
225
+ print(f"Chunks per subject: {dict(chunks_by_subject)}")
226
+ return
227
+
228
+ vectorstore_dir.mkdir(parents=True, exist_ok=True)
229
+ documents, metadatas, ids = build_documents(data_dir)
230
+ if not documents:
231
+ raise SystemExit("No text extracted from curriculum files")
232
+ import chromadb
233
+ from sentence_transformers import SentenceTransformer
234
+
235
+ cache_file = vectorstore_dir / "embeddings_cache.npy"
236
+ embeddings: List[List[float]] = []
237
+ if cache_file.exists():
238
+ try:
239
+ import numpy as np
240
+ cached_data = np.load(cache_file)
241
+ if len(cached_data) == len(documents) and cached_data.shape[1] == 384:
242
+ print(f"Loaded {len(cached_data)} cached embeddings from {cache_file}")
243
+ embeddings = cached_data.tolist()
244
+ except Exception:
245
+ embeddings = []
246
+
247
+ if not embeddings:
248
+ embedder = SentenceTransformer(EMBED_MODEL_NAME)
249
+ encoded = embedder.encode(
250
+ documents,
251
+ batch_size=64,
252
+ normalize_embeddings=True,
253
+ show_progress_bar=True,
254
+ )
255
+ try:
256
+ import numpy as np
257
+ np.save(cache_file, encoded)
258
+ except Exception:
259
+ pass
260
+ embeddings = encoded.tolist()
261
+
262
+ import sqlite3
263
+
264
+ client = chromadb.PersistentClient(path=str(vectorstore_dir))
265
+ chunk_count_before = 0
266
+ sqlite_file = vectorstore_dir / "chroma.sqlite3"
267
+ needs_recreate = False
268
+ if sqlite_file.exists():
269
+ try:
270
+ with sqlite3.connect(str(sqlite_file)) as conn:
271
+ row = conn.execute(
272
+ "SELECT dimension FROM collections WHERE name = ?", (COLLECTION_NAME,)
273
+ ).fetchone()
274
+ if row and row[0] is not None and row[0] != len(embeddings[0]):
275
+ print(
276
+ f"Dimension mismatch ({row[0]} != {len(embeddings[0])}); "
277
+ f"recreating collection {COLLECTION_NAME}..."
278
+ )
279
+ needs_recreate = True
280
  except Exception:
281
  pass
282
+
283
+ if needs_recreate:
284
+ try:
285
+ existing_col = client.get_collection(COLLECTION_NAME)
286
+ chunk_count_before = existing_col.count()
287
+ client.delete_collection(COLLECTION_NAME)
288
+ except Exception:
289
+ pass
290
+ collection = client.create_collection(
291
+ name=COLLECTION_NAME, metadata={"hnsw:space": "cosine"}
292
+ )
293
+ else:
294
+ try:
295
+ collection = client.get_collection(COLLECTION_NAME)
296
+ chunk_count_before = collection.count()
297
+ except Exception:
298
+ collection = client.create_collection(
299
+ name=COLLECTION_NAME, metadata={"hnsw:space": "cosine"}
300
+ )
301
+
302
+ print(f"Chunk count before: {chunk_count_before}")
303
+
304
+ existing_ids = set(collection.get(include=[])["ids"])
305
+ new_ids = set(ids)
306
+ stale_ids = list(existing_ids - new_ids)
307
+ if stale_ids:
308
+ print(f"Pruning {len(stale_ids)} stale chunks from collection...")
309
+ for start in range(0, len(stale_ids), 500):
310
+ collection.delete(ids=stale_ids[start : start + 500])
311
+
312
+ for start in range(0, len(ids), 500):
313
+ end = start + 500
314
+ collection.upsert(
315
+ ids=ids[start:end],
316
+ documents=documents[start:end],
317
+ metadatas=metadatas[start:end],
318
+ embeddings=embeddings[start:end],
319
+ )
320
+ chunk_count_after = collection.count()
321
+ print(f"Chunk count after: {chunk_count_after}")
322
+ summary = {
323
+ "lastIngested": datetime.now(timezone.utc).isoformat(),
324
+ "totalChunks": len(documents),
325
+ "chunkCountBefore": chunk_count_before,
326
+ "chunkCountAfter": chunk_count_after,
327
+ "sourceFiles": [path.as_posix() for path in files],
328
+ "chunksPerSubject": dict(Counter(str(meta["subject"]) for meta in metadatas)),
329
+ }
330
+ (vectorstore_dir / "ingest_summary.json").write_text(json.dumps(summary, indent=2), encoding="utf-8")
331
+ print(f"Total chunks: {len(documents)}")
332
+ print(f"Source files: {len(files)}")
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
333
 
334
 
335
  if __name__ == "__main__":
336
+ main()