diff --git a/.vscode/settings.json b/.vscode/settings.json index 3ad0b35..4ca12c2 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -1,5 +1,5 @@ { - "python.defaultInterpreterPath": "${workspaceFolder}/scraping-service/.venv/bin/python3", + "python.defaultInterpreterPath": "${workspaceFolder}/scraping/.venv/bin/python3", "python.analysis.extraPaths": ["${workspaceFolder}/server/src"], "python.terminal.activateEnvironment": true, "cSpell.words": [ diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..e37eef3 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,9 @@ +# Dev Environment Setup + +Use bun instead of node + +# Code Styling + +# Testing Preferences + +For golang use testify for tests diff --git a/client/src/lib/components/onboarding/onboarding-chips-field.svelte b/client/src/lib/components/onboarding/onboarding-chips-field.svelte index ce360d0..746aa09 100644 --- a/client/src/lib/components/onboarding/onboarding-chips-field.svelte +++ b/client/src/lib/components/onboarding/onboarding-chips-field.svelte @@ -62,7 +62,11 @@ onclick={() => inputEl?.focus()} > {#each values as item, i (item + i)} - diff --git a/client/src/lib/components/onboarding/onboarding-form.svelte b/client/src/lib/components/onboarding/onboarding-form.svelte index 7ea3498..9c25072 100644 --- a/client/src/lib/components/onboarding/onboarding-form.svelte +++ b/client/src/lib/components/onboarding/onboarding-form.svelte @@ -27,7 +27,9 @@ let step = $state(1); function applyFieldErrors( - currentErrors: Parameters[0] extends (arg: infer T) => unknown ? T : never, + currentErrors: Parameters[0] extends (arg: infer T) => unknown + ? T + : never, fieldErrors: Record ) { const mutableErrors = currentErrors as Record; diff --git a/client/src/lib/components/onboarding/onboarding-step-nav.svelte b/client/src/lib/components/onboarding/onboarding-step-nav.svelte index 1447e80..ab3e12d 100644 --- a/client/src/lib/components/onboarding/onboarding-step-nav.svelte +++ b/client/src/lib/components/onboarding/onboarding-step-nav.svelte @@ -1,5 +1,5 @@ @@ -43,7 +41,7 @@ let title = $state(''); Choose a new name and term
-
+
Any: + with path.open("r", encoding="utf-8") as f: + return json.load(f) + + +def build_artifacts_payload(data_dir: pathlib.Path = DATA_DIR) -> dict[str, Any]: + rmp_data = load_json(data_dir / "rmp_data.json") + return { + "courses": load_json(data_dir / "all_courses.json"), + "programs": load_json(data_dir / "all_programs_with_requirements.json"), + "teachers": rmp_data.get("professors", []), + "schedules": load_json(data_dir / "all_possible_schedules.json"), + } + + +def build_json_request(url: str, token: str, payload: dict[str, Any]) -> urllib.request.Request: + body = json.dumps(payload).encode("utf-8") + return urllib.request.Request( + url, + data=body, + headers={ + "Authorization": f"Bearer {token}", + "Content-Type": "application/json", + }, + method="POST", + ) + + +def post_json(url: str, token: str, payload: dict[str, Any]) -> dict[str, Any]: + request = build_json_request(url, token, payload) + try: + with urllib.request.urlopen(request) as response: + return json.loads(response.read().decode("utf-8")) + except urllib.error.HTTPError as e: + detail = e.read().decode("utf-8") + raise RuntimeError(f"POST {url} failed with {e.code}: {detail}") from e + + +def run_publish( + api_base_url: str, + token: str, + data_dir: pathlib.Path = DATA_DIR, + source: str = "scraping", + promote: bool = True, +) -> dict[str, Any]: + if not token: + raise ValueError("INTERNAL_SERVICE_TOKEN is required") + + api_base_url = api_base_url.rstrip("/") + run = post_json( + f"{api_base_url}/internal/scrape-runs", + token, + {"source": source, "metadata": {"publisher": "publish_scrape_run.py"}}, + ) + run_id = run["ID"] if "ID" in run else run["id"] + + artifacts = build_artifacts_payload(data_dir) + staged = post_json( + f"{api_base_url}/internal/scrape-runs/{run_id}/artifacts", + token, + artifacts, + ) + + result: dict[str, Any] = {"run": run, "staged": staged} + if promote: + result["promoted"] = post_json( + f"{api_base_url}/internal/scrape-runs/{run_id}/promote", + token, + {}, + ) + return result + + +def parse_args(argv: list[str]) -> argparse.Namespace: + parser = argparse.ArgumentParser(description="Publish scraper artifacts to the Go ingest API.") + parser.add_argument("--api-base-url", default=os.getenv("GO_API_BASE_URL", "http://localhost:8000")) + parser.add_argument("--token", default=os.getenv("INTERNAL_SERVICE_TOKEN", "")) + parser.add_argument("--data-dir", type=pathlib.Path, default=DATA_DIR) + parser.add_argument("--source", default="scraping") + parser.add_argument("--no-promote", action="store_true", help="Stage artifacts without promoting them.") + return parser.parse_args(argv) + + +def main(argv: list[str] | None = None) -> int: + args = parse_args(argv or sys.argv[1:]) + result = run_publish( + api_base_url=args.api_base_url, + token=args.token, + data_dir=args.data_dir, + source=args.source, + promote=not args.no_promote, + ) + print(json.dumps(result, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scraping/scripts/seed_db.py b/scraping/scripts/seed_db.py index 23e9e9c..459587b 100644 --- a/scraping/scripts/seed_db.py +++ b/scraping/scripts/seed_db.py @@ -133,6 +133,10 @@ def normalize_term(term_str): def seed(): + print( + "WARNING: seed_db.py is the legacy direct-DB seeder. " + "Prefer scripts/publish_scrape_run.py so the Go API owns validation and promotion." + ) conn = get_db_connection() cur = conn.cursor() diff --git a/scraping/scripts/test_get_all_programs.py b/scraping/scripts/tests/test_get_all_programs.py similarity index 100% rename from scraping/scripts/test_get_all_programs.py rename to scraping/scripts/tests/test_get_all_programs.py diff --git a/scraping/scripts/test_parse_antirequisites.py b/scraping/scripts/tests/test_parse_antirequisites.py similarity index 100% rename from scraping/scripts/test_parse_antirequisites.py rename to scraping/scripts/tests/test_parse_antirequisites.py diff --git a/scraping/scripts/tests/test_publish_scrape_run.py b/scraping/scripts/tests/test_publish_scrape_run.py new file mode 100644 index 0000000..ee829e7 --- /dev/null +++ b/scraping/scripts/tests/test_publish_scrape_run.py @@ -0,0 +1,55 @@ +import importlib.util +import json +import pathlib +import tempfile +import unittest + + +SCRIPT_PATH = pathlib.Path(__file__).resolve().parent / "publish_scrape_run.py" +SPEC = importlib.util.spec_from_file_location("publish_scrape_run", SCRIPT_PATH) +assert SPEC is not None +assert SPEC.loader is not None +PUBLISHER = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(PUBLISHER) + + +class BuildArtifactsPayloadTests(unittest.TestCase): + def test_builds_payload_from_existing_scraper_files(self): + with tempfile.TemporaryDirectory() as tmp: + data_dir = pathlib.Path(tmp) + (data_dir / "all_courses.json").write_text( + json.dumps([{"course_name": "COMPSCI 1DM3 - Discrete Mathematics"}]) + ) + (data_dir / "all_programs_with_requirements.json").write_text( + json.dumps([{"program_name": "Computer Science", "requirements": ["COMPSCI 1DM3"]}]) + ) + (data_dir / "rmp_data.json").write_text( + json.dumps({"professors": [{"id": "rmp_1", "name": "Jane Doe"}]}) + ) + (data_dir / "all_possible_schedules.json").write_text( + json.dumps([{"term": "Fall 2026", "courses": []}]) + ) + + payload = PUBLISHER.build_artifacts_payload(data_dir) + + self.assertEqual(payload["courses"][0]["course_name"], "COMPSCI 1DM3 - Discrete Mathematics") + self.assertEqual(payload["programs"][0]["program_name"], "Computer Science") + self.assertEqual(payload["teachers"][0]["id"], "rmp_1") + self.assertEqual(payload["schedules"][0]["term"], "Fall 2026") + + +class RequestTests(unittest.TestCase): + def test_builds_internal_authorization_request(self): + request = PUBLISHER.build_json_request( + "http://localhost:8000/internal/scrape-runs", + "secret", + {"source": "manual"}, + ) + + self.assertEqual(request.get_method(), "POST") + self.assertEqual(request.headers["Authorization"], "Bearer secret") + self.assertEqual(request.headers["Content-type"], "application/json") + + +if __name__ == "__main__": + unittest.main() diff --git a/scraping/scripts/test_scrape_mosaic_schedules.py b/scraping/scripts/tests/test_scrape_mosaic_schedules.py similarity index 96% rename from scraping/scripts/test_scrape_mosaic_schedules.py rename to scraping/scripts/tests/test_scrape_mosaic_schedules.py index 313585e..47b5621 100644 --- a/scraping/scripts/test_scrape_mosaic_schedules.py +++ b/scraping/scripts/tests/test_scrape_mosaic_schedules.py @@ -7,18 +7,14 @@ sys.modules.setdefault("dotenv", types.SimpleNamespace(load_dotenv=lambda *_args, **_kwargs: None)) sys.modules.setdefault("psycopg2", types.SimpleNamespace(connect=lambda *_args, **_kwargs: None)) -sys.modules.setdefault( - "playwright", - types.SimpleNamespace(async_api=types.SimpleNamespace(async_playwright=None)), +sys.modules["playwright"] = types.SimpleNamespace( + async_api=types.SimpleNamespace(async_playwright=None) ) -sys.modules.setdefault( - "playwright.async_api", - types.SimpleNamespace( - async_playwright=None, - Page=object, - Frame=object, - Browser=object, - ), +sys.modules["playwright.async_api"] = types.SimpleNamespace( + async_playwright=None, + Page=object, + Frame=object, + Browser=object, ) diff --git a/scraping/scripts/test_seed_db.py b/scraping/scripts/tests/test_seed_db.py similarity index 100% rename from scraping/scripts/test_seed_db.py rename to scraping/scripts/tests/test_seed_db.py diff --git a/scraping/test_fetch.py b/scraping/test_fetch.py deleted file mode 100644 index 462f793..0000000 --- a/scraping/test_fetch.py +++ /dev/null @@ -1,33 +0,0 @@ -import asyncio -from playwright.async_api import async_playwright - -async def main(): - async with async_playwright() as p: - browser = await p.chromium.launch() - context = await browser.new_context() - page = await context.new_page() - - url = "https://academiccalendars.romcmaster.ca/content.php?catoid=65&navoid=14802&filter%5Bitem_type%5D=3&filter%5Bonly_active%5D=1&filter%5B3%5D=1&filter%5Bcpage%5D=28#acalog_template_course_filter" - await page.goto(url) - - course_links = await page.evaluate(""" - () => Array.from(document.querySelectorAll("a[onclick^='showCourse']")).map(a => { - const match = a.getAttribute('onclick').match(/showCourse\('(\d+)',\s*'(\d+)'/); - return { text: a.innerText.trim(), coid: match ? match[2] : null }; - }).filter(item => item.coid !== null) - """) - - print(f"Found {len(course_links)} courses") - if course_links: - first = course_links[0] - print(f"First course: {first}") - - preview_url = f"https://academiccalendars.romcmaster.ca/ajax/preview_course.php?catoid=65&show&coid={first['coid']}" - resp = await context.request.get(preview_url) - print(f"Status: {resp.status}") - text = await resp.text() - print(f"Response: {text[:200]}") - - await browser.close() - -asyncio.run(main()) diff --git a/scraping/test_fetch2.py b/scraping/test_fetch2.py deleted file mode 100644 index f45cfd3..0000000 --- a/scraping/test_fetch2.py +++ /dev/null @@ -1,32 +0,0 @@ -import asyncio -from playwright.async_api import async_playwright - -async def main(): - async with async_playwright() as p: - browser = await p.chromium.launch() - context = await browser.new_context(ignore_https_errors=True) - page = await context.new_page() - - url = "https://academiccalendars.romcmaster.ca/content.php?catoid=65&navoid=14802&filter%5Bitem_type%5D=3&filter%5Bonly_active%5D=1&filter%5B3%5D=1&filter%5Bcpage%5D=28#acalog_template_course_filter" - await page.goto(url) - - course_links = await page.evaluate(""" - () => Array.from(document.querySelectorAll("a[onclick^='showCourse']")).map(a => { - const match = a.getAttribute('onclick').match(/showCourse\\('(\\d+)',\\s*'(\\d+)'/); - return { text: a.innerText.trim(), coid: match ? match[2] : null }; - }).filter(item => item.coid !== null) - """) - - if course_links: - first = course_links[0] - print(f"First course: {first}") - - preview_url = f"https://academiccalendars.romcmaster.ca/ajax/preview_course.php?catoid=65&show&coid={first['coid']}" - resp = await context.request.get(preview_url, ignore_https_errors=True) - print(f"Status: {resp.status}") - text = await resp.text() - print(f"Response: {text[:50]}") - - await browser.close() - -asyncio.run(main()) diff --git a/server/db/migrations/20260614143000_scrape_ingestion_pipeline.sql b/server/db/migrations/20260614143000_scrape_ingestion_pipeline.sql new file mode 100644 index 0000000..32cbefa --- /dev/null +++ b/server/db/migrations/20260614143000_scrape_ingestion_pipeline.sql @@ -0,0 +1,80 @@ +-- +goose Up + +CREATE TABLE IF NOT EXISTS scrape_runs ( + id UUID PRIMARY KEY DEFAULT uuidv7(), + source TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'received', + metadata JSONB NOT NULL DEFAULT '{}'::jsonb, + error_message TEXT, + course_count INTEGER NOT NULL DEFAULT 0, + program_count INTEGER NOT NULL DEFAULT 0, + teacher_count INTEGER NOT NULL DEFAULT 0, + schedule_term_count INTEGER NOT NULL DEFAULT 0, + started_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + staged_at TIMESTAMPTZ, + promoted_at TIMESTAMPTZ, + failed_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + CONSTRAINT scrape_runs_status_check CHECK ( + status IN ('received', 'staged', 'promoting', 'succeeded', 'failed') + ) +); + +CREATE TABLE IF NOT EXISTS staging_courses ( + run_id UUID NOT NULL REFERENCES scrape_runs(id) ON DELETE CASCADE, + course_code TEXT NOT NULL, + payload JSONB NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (run_id, course_code) +); + +CREATE TABLE IF NOT EXISTS staging_programs ( + run_id UUID NOT NULL REFERENCES scrape_runs(id) ON DELETE CASCADE, + program_name TEXT NOT NULL, + payload JSONB NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (run_id, program_name) +); + +CREATE TABLE IF NOT EXISTS staging_teachers ( + run_id UUID NOT NULL REFERENCES scrape_runs(id) ON DELETE CASCADE, + rmp_id TEXT NOT NULL, + payload JSONB NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (run_id, rmp_id) +); + +CREATE TABLE IF NOT EXISTS staging_schedules ( + run_id UUID NOT NULL REFERENCES scrape_runs(id) ON DELETE CASCADE, + term TEXT NOT NULL, + course_code TEXT NOT NULL, + payload JSONB NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (run_id, term, course_code) +); + +CREATE INDEX IF NOT EXISTS idx_staging_courses_run_id + ON staging_courses (run_id); + +CREATE INDEX IF NOT EXISTS idx_staging_programs_run_id + ON staging_programs (run_id); + +CREATE INDEX IF NOT EXISTS idx_staging_teachers_run_id + ON staging_teachers (run_id); + +CREATE INDEX IF NOT EXISTS idx_staging_schedules_run_id_term + ON staging_schedules (run_id, term); + +-- +goose Down + +DROP INDEX IF EXISTS idx_staging_schedules_run_id_term; +DROP INDEX IF EXISTS idx_staging_teachers_run_id; +DROP INDEX IF EXISTS idx_staging_programs_run_id; +DROP INDEX IF EXISTS idx_staging_courses_run_id; + +DROP TABLE IF EXISTS staging_schedules; +DROP TABLE IF EXISTS staging_teachers; +DROP TABLE IF EXISTS staging_programs; +DROP TABLE IF EXISTS staging_courses; +DROP TABLE IF EXISTS scrape_runs; diff --git a/server/db/query/scrape_ingest.sql b/server/db/query/scrape_ingest.sql new file mode 100644 index 0000000..c909bf3 --- /dev/null +++ b/server/db/query/scrape_ingest.sql @@ -0,0 +1,209 @@ +-- name: CreateScrapeRun :one +INSERT INTO scrape_runs (source, metadata) +VALUES ($1, $2) +RETURNING id, source, status, metadata, error_message, course_count, program_count, teacher_count, schedule_term_count, started_at, staged_at, promoted_at, failed_at, created_at, updated_at; + +-- name: GetScrapeRun :one +SELECT id, source, status, metadata, error_message, course_count, program_count, teacher_count, schedule_term_count, started_at, staged_at, promoted_at, failed_at, created_at, updated_at +FROM scrape_runs +WHERE id = $1 +LIMIT 1; + +-- name: MarkScrapeRunStaged :exec +UPDATE scrape_runs +SET status = 'staged', + course_count = $2, + program_count = $3, + teacher_count = $4, + schedule_term_count = $5, + staged_at = NOW(), + updated_at = NOW(), + error_message = NULL +WHERE id = $1; + +-- name: MarkScrapeRunPromoting :exec +UPDATE scrape_runs +SET status = 'promoting', + updated_at = NOW(), + error_message = NULL +WHERE id = $1; + +-- name: MarkScrapeRunSucceeded :exec +UPDATE scrape_runs +SET status = 'succeeded', + promoted_at = NOW(), + updated_at = NOW(), + error_message = NULL +WHERE id = $1; + +-- name: MarkScrapeRunFailed :exec +UPDATE scrape_runs +SET status = 'failed', + failed_at = NOW(), + updated_at = NOW(), + error_message = $2 +WHERE id = $1; + +-- name: ClearScrapeRunStagedSchedules :exec +DELETE FROM staging_schedules WHERE run_id = $1; + +-- name: ClearScrapeRunStagedTeachers :exec +DELETE FROM staging_teachers WHERE run_id = $1; + +-- name: ClearScrapeRunStagedPrograms :exec +DELETE FROM staging_programs WHERE run_id = $1; + +-- name: ClearScrapeRunStagedCourses :exec +DELETE FROM staging_courses WHERE run_id = $1; + +-- name: StageCoursePayload :exec +INSERT INTO staging_courses (run_id, course_code, payload) +VALUES ($1, $2, $3) +ON CONFLICT (run_id, course_code) DO UPDATE SET + payload = EXCLUDED.payload; + +-- name: StageProgramPayload :exec +INSERT INTO staging_programs (run_id, program_name, payload) +VALUES ($1, $2, $3) +ON CONFLICT (run_id, program_name) DO UPDATE SET + payload = EXCLUDED.payload; + +-- name: StageTeacherPayload :exec +INSERT INTO staging_teachers (run_id, rmp_id, payload) +VALUES ($1, $2, $3) +ON CONFLICT (run_id, rmp_id) DO UPDATE SET + payload = EXCLUDED.payload; + +-- name: StageSchedulePayload :exec +INSERT INTO staging_schedules (run_id, term, course_code, payload) +VALUES ($1, $2, $3, $4) +ON CONFLICT (run_id, term, course_code) DO UPDATE SET + payload = EXCLUDED.payload; + +-- name: ListStagedCoursePayloads :many +SELECT payload +FROM staging_courses +WHERE run_id = $1 +ORDER BY course_code; + +-- name: ListStagedProgramPayloads :many +SELECT payload +FROM staging_programs +WHERE run_id = $1 +ORDER BY program_name; + +-- name: ListStagedTeacherPayloads :many +SELECT payload +FROM staging_teachers +WHERE run_id = $1 +ORDER BY rmp_id; + +-- name: ListStagedSchedulePayloads :many +SELECT payload +FROM staging_schedules +WHERE run_id = $1 +ORDER BY term, course_code; + +-- name: ListStagedScheduleTerms :many +SELECT DISTINCT term +FROM staging_schedules +WHERE run_id = $1 +ORDER BY term; + +-- name: UpsertScrapeCourse :exec +INSERT INTO course (code, name, description, restrictions, prerequisites, units, level_number) +VALUES ($1, $2, $3, $4, $5, $6, $7) +ON CONFLICT (code) DO UPDATE SET + name = EXCLUDED.name, + description = EXCLUDED.description, + restrictions = EXCLUDED.restrictions, + prerequisites = EXCLUDED.prerequisites, + units = EXCLUDED.units, + level_number = EXCLUDED.level_number; + +-- name: ListScrapeCourseIDs :many +SELECT code, id +FROM course +ORDER BY code; + +-- name: UpsertScrapeTeacher :exec +INSERT INTO teacher (name, avg_rating, avg_difficulty, department, rmp_id, num_ratings) +VALUES ($1, $2, $3, $4, $5, $6) +ON CONFLICT (rmp_id) DO UPDATE SET + name = EXCLUDED.name, + avg_rating = EXCLUDED.avg_rating, + avg_difficulty = EXCLUDED.avg_difficulty, + department = EXCLUDED.department, + num_ratings = EXCLUDED.num_ratings; + +-- name: ListScrapeTeacherIDsByName :many +SELECT name, id +FROM teacher +ORDER BY name; + +-- name: UpsertScrapeProgram :one +WITH updated AS ( + UPDATE program + SET source_url = sqlc.arg('source_url'), + requirement_codes = sqlc.arg('requirement_codes'), + requirements_by_level = sqlc.arg('requirements_by_level') + WHERE program.name = sqlc.arg('program_name') + RETURNING id +), +inserted AS ( + INSERT INTO program (name, source_url, requirement_codes, requirements_by_level) + SELECT sqlc.arg('program_name'), + sqlc.arg('source_url'), + sqlc.arg('requirement_codes'), + sqlc.arg('requirements_by_level') + WHERE NOT EXISTS (SELECT 1 FROM updated) + RETURNING id +) +SELECT id FROM updated +UNION ALL +SELECT id FROM inserted +LIMIT 1; + +-- name: DeleteScrapeProgramCourses :exec +DELETE FROM program_courses +WHERE program_id = $1; + +-- name: InsertScrapeProgramCourse :exec +INSERT INTO program_courses (program_id, course_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING; + +-- name: InsertScrapeCourseTeacher :exec +INSERT INTO course_teachers (course_id, teacher_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING; + +-- name: DeleteScrapeSectionsForTerms :exec +DELETE FROM section +WHERE term = ANY($1::text[]); + +-- name: CreateScrapeSection :one +INSERT INTO section (course_id, name, type, term, mode, is_in_person) +VALUES ($1, $2, $3, $4, $5, $6) +RETURNING id; + +-- name: CreateScrapeSectionMeeting :exec +INSERT INTO section_meeting (section_id, days, start_time, end_time, building, room) +VALUES ($1, $2, $3, $4, $5, $6); + +-- name: InsertScrapeSectionTeacher :exec +INSERT INTO section_teachers (section_id, teacher_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING; + +-- name: InsertScrapeSectionReference :exec +INSERT INTO section_references (parent_section_id, child_section_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING; + +-- name: InsertScrapeTeacherProgramLinks :exec +INSERT INTO teacher_programs (teacher_id, program_id) +SELECT DISTINCT ct.teacher_id, pc.program_id +FROM course_teachers ct +JOIN program_courses pc ON pc.course_id = ct.course_id +ON CONFLICT DO NOTHING; diff --git a/server/internal/app/routes.go b/server/internal/app/routes.go index d0db3c7..a4b0114 100644 --- a/server/internal/app/routes.go +++ b/server/internal/app/routes.go @@ -63,5 +63,14 @@ func addRoutes(r *chi.Mux, cfg *config.Config, s *Services) { r.Get("/", handlers.HandleGetAllPlans(s.Planner)) }) + r.Route("/internal", func(r chi.Router) { + r.Use(mw.RequireInternalToken(cfg.InternalServiceToken)) + + r.Post("/scrape-runs", handlers.HandleCreateScrapeRun(s.ScrapeIngest)) + r.Get("/scrape-runs/{run_id}", handlers.HandleGetScrapeRun(s.ScrapeIngest)) + r.Post("/scrape-runs/{run_id}/artifacts", handlers.HandleStageScrapeArtifacts(s.ScrapeIngest)) + r.Post("/scrape-runs/{run_id}/promote", handlers.HandlePromoteScrapeRun(s.ScrapeIngest)) + }) + r.With(mw.RequireCSRF(cfg.SecretKey)).Post("/logout", handlers.HandleLogout(cfg, s.Auth)) } diff --git a/server/internal/app/services.go b/server/internal/app/services.go index 9c85dc2..20a4abd 100644 --- a/server/internal/app/services.go +++ b/server/internal/app/services.go @@ -7,17 +7,19 @@ import ( "github.com/twitocode/pathweave/go-api/internal/config" "github.com/twitocode/pathweave/go-api/internal/db" "github.com/twitocode/pathweave/go-api/internal/service" + "github.com/twitocode/pathweave/go-api/internal/service/scraping" ) type Services struct { - Auth *service.AuthService - User *service.UserService - Onboarding *service.OnboardingService - Program *service.ProgramService - Course *service.CourseService - Embedding *service.EmbeddingService - AI *service.AIService - Planner *service.PlannerService + Auth *service.AuthService + User *service.UserService + Onboarding *service.OnboardingService + Program *service.ProgramService + Course *service.CourseService + Embedding *service.EmbeddingService + AI *service.AIService + Planner *service.PlannerService + ScrapeIngest *scraping.ScrapeIngestService } func NewServices(cfg *config.Config, pool *pgxpool.Pool, log *zap.Logger) *Services { @@ -27,13 +29,14 @@ func NewServices(cfg *config.Config, pool *pgxpool.Pool, log *zap.Logger) *Servi ais := service.NewAIService(cfg, queries, log) return &Services{ - Auth: service.NewAuthService(cfg, queries, log), - User: service.NewUserService(queries, log), - Onboarding: service.NewOnboardingService(queries, log), - Course: service.NewCourseService(queries, log, es, ais), - Program: service.NewProgramService(queries, log), - Embedding: es, - AI: ais, - Planner: service.NewPlannerService(queries, log), + Auth: service.NewAuthService(cfg, queries, log), + User: service.NewUserService(queries, log), + Onboarding: service.NewOnboardingService(queries, log), + Course: service.NewCourseService(queries, log, es, ais), + Program: service.NewProgramService(queries, log), + Embedding: es, + AI: ais, + Planner: service.NewPlannerService(queries, log), + ScrapeIngest: scraping.NewScrapeIngestService(queries, pool, log), } } diff --git a/server/internal/db/models.go b/server/internal/db/models.go index ce71276..ffdfea7 100644 --- a/server/internal/db/models.go +++ b/server/internal/db/models.go @@ -88,6 +88,24 @@ type ProgramRequirementLevel struct { SortOrder int32 } +type ScrapeRun struct { + ID uuid.UUID + Source string + Status string + Metadata []byte + ErrorMessage pgtype.Text + CourseCount int32 + ProgramCount int32 + TeacherCount int32 + ScheduleTermCount int32 + StartedAt pgtype.Timestamptz + StagedAt pgtype.Timestamptz + PromotedAt pgtype.Timestamptz + FailedAt pgtype.Timestamptz + CreatedAt pgtype.Timestamptz + UpdatedAt pgtype.Timestamptz +} + type Section struct { ID int32 CourseID int64 @@ -118,6 +136,35 @@ type SectionTeacher struct { TeacherID int32 } +type StagingCourse struct { + RunID uuid.UUID + CourseCode string + Payload []byte + CreatedAt pgtype.Timestamptz +} + +type StagingProgram struct { + RunID uuid.UUID + ProgramName string + Payload []byte + CreatedAt pgtype.Timestamptz +} + +type StagingSchedule struct { + RunID uuid.UUID + Term string + CourseCode string + Payload []byte + CreatedAt pgtype.Timestamptz +} + +type StagingTeacher struct { + RunID uuid.UUID + RmpID string + Payload []byte + CreatedAt pgtype.Timestamptz +} + type Teacher struct { ID int32 Name string diff --git a/server/internal/db/scrape_ingest.sql.go b/server/internal/db/scrape_ingest.sql.go new file mode 100644 index 0000000..72d5624 --- /dev/null +++ b/server/internal/db/scrape_ingest.sql.go @@ -0,0 +1,728 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.30.0 +// source: scrape_ingest.sql + +package db + +import ( + "context" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgtype" +) + +const clearScrapeRunStagedCourses = `-- name: ClearScrapeRunStagedCourses :exec +DELETE FROM staging_courses WHERE run_id = $1 +` + +func (q *Queries) ClearScrapeRunStagedCourses(ctx context.Context, runID uuid.UUID) error { + _, err := q.db.Exec(ctx, clearScrapeRunStagedCourses, runID) + return err +} + +const clearScrapeRunStagedPrograms = `-- name: ClearScrapeRunStagedPrograms :exec +DELETE FROM staging_programs WHERE run_id = $1 +` + +func (q *Queries) ClearScrapeRunStagedPrograms(ctx context.Context, runID uuid.UUID) error { + _, err := q.db.Exec(ctx, clearScrapeRunStagedPrograms, runID) + return err +} + +const clearScrapeRunStagedSchedules = `-- name: ClearScrapeRunStagedSchedules :exec +DELETE FROM staging_schedules WHERE run_id = $1 +` + +func (q *Queries) ClearScrapeRunStagedSchedules(ctx context.Context, runID uuid.UUID) error { + _, err := q.db.Exec(ctx, clearScrapeRunStagedSchedules, runID) + return err +} + +const clearScrapeRunStagedTeachers = `-- name: ClearScrapeRunStagedTeachers :exec +DELETE FROM staging_teachers WHERE run_id = $1 +` + +func (q *Queries) ClearScrapeRunStagedTeachers(ctx context.Context, runID uuid.UUID) error { + _, err := q.db.Exec(ctx, clearScrapeRunStagedTeachers, runID) + return err +} + +const createScrapeRun = `-- name: CreateScrapeRun :one +INSERT INTO scrape_runs (source, metadata) +VALUES ($1, $2) +RETURNING id, source, status, metadata, error_message, course_count, program_count, teacher_count, schedule_term_count, started_at, staged_at, promoted_at, failed_at, created_at, updated_at +` + +type CreateScrapeRunParams struct { + Source string + Metadata []byte +} + +func (q *Queries) CreateScrapeRun(ctx context.Context, arg CreateScrapeRunParams) (ScrapeRun, error) { + row := q.db.QueryRow(ctx, createScrapeRun, arg.Source, arg.Metadata) + var i ScrapeRun + err := row.Scan( + &i.ID, + &i.Source, + &i.Status, + &i.Metadata, + &i.ErrorMessage, + &i.CourseCount, + &i.ProgramCount, + &i.TeacherCount, + &i.ScheduleTermCount, + &i.StartedAt, + &i.StagedAt, + &i.PromotedAt, + &i.FailedAt, + &i.CreatedAt, + &i.UpdatedAt, + ) + return i, err +} + +const createScrapeSection = `-- name: CreateScrapeSection :one +INSERT INTO section (course_id, name, type, term, mode, is_in_person) +VALUES ($1, $2, $3, $4, $5, $6) +RETURNING id +` + +type CreateScrapeSectionParams struct { + CourseID int64 + Name string + Type string + Term string + Mode string + IsInPerson bool +} + +func (q *Queries) CreateScrapeSection(ctx context.Context, arg CreateScrapeSectionParams) (int32, error) { + row := q.db.QueryRow(ctx, createScrapeSection, + arg.CourseID, + arg.Name, + arg.Type, + arg.Term, + arg.Mode, + arg.IsInPerson, + ) + var id int32 + err := row.Scan(&id) + return id, err +} + +const createScrapeSectionMeeting = `-- name: CreateScrapeSectionMeeting :exec +INSERT INTO section_meeting (section_id, days, start_time, end_time, building, room) +VALUES ($1, $2, $3, $4, $5, $6) +` + +type CreateScrapeSectionMeetingParams struct { + SectionID int32 + Days string + StartTime pgtype.Time + EndTime pgtype.Time + Building string + Room string +} + +func (q *Queries) CreateScrapeSectionMeeting(ctx context.Context, arg CreateScrapeSectionMeetingParams) error { + _, err := q.db.Exec(ctx, createScrapeSectionMeeting, + arg.SectionID, + arg.Days, + arg.StartTime, + arg.EndTime, + arg.Building, + arg.Room, + ) + return err +} + +const deleteScrapeProgramCourses = `-- name: DeleteScrapeProgramCourses :exec +DELETE FROM program_courses +WHERE program_id = $1 +` + +func (q *Queries) DeleteScrapeProgramCourses(ctx context.Context, programID int64) error { + _, err := q.db.Exec(ctx, deleteScrapeProgramCourses, programID) + return err +} + +const deleteScrapeSectionsForTerms = `-- name: DeleteScrapeSectionsForTerms :exec +DELETE FROM section +WHERE term = ANY($1::text[]) +` + +func (q *Queries) DeleteScrapeSectionsForTerms(ctx context.Context, dollar_1 []string) error { + _, err := q.db.Exec(ctx, deleteScrapeSectionsForTerms, dollar_1) + return err +} + +const getScrapeRun = `-- name: GetScrapeRun :one +SELECT id, source, status, metadata, error_message, course_count, program_count, teacher_count, schedule_term_count, started_at, staged_at, promoted_at, failed_at, created_at, updated_at +FROM scrape_runs +WHERE id = $1 +LIMIT 1 +` + +func (q *Queries) GetScrapeRun(ctx context.Context, id uuid.UUID) (ScrapeRun, error) { + row := q.db.QueryRow(ctx, getScrapeRun, id) + var i ScrapeRun + err := row.Scan( + &i.ID, + &i.Source, + &i.Status, + &i.Metadata, + &i.ErrorMessage, + &i.CourseCount, + &i.ProgramCount, + &i.TeacherCount, + &i.ScheduleTermCount, + &i.StartedAt, + &i.StagedAt, + &i.PromotedAt, + &i.FailedAt, + &i.CreatedAt, + &i.UpdatedAt, + ) + return i, err +} + +const insertScrapeCourseTeacher = `-- name: InsertScrapeCourseTeacher :exec +INSERT INTO course_teachers (course_id, teacher_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING +` + +type InsertScrapeCourseTeacherParams struct { + CourseID int64 + TeacherID int32 +} + +func (q *Queries) InsertScrapeCourseTeacher(ctx context.Context, arg InsertScrapeCourseTeacherParams) error { + _, err := q.db.Exec(ctx, insertScrapeCourseTeacher, arg.CourseID, arg.TeacherID) + return err +} + +const insertScrapeProgramCourse = `-- name: InsertScrapeProgramCourse :exec +INSERT INTO program_courses (program_id, course_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING +` + +type InsertScrapeProgramCourseParams struct { + ProgramID int64 + CourseID int64 +} + +func (q *Queries) InsertScrapeProgramCourse(ctx context.Context, arg InsertScrapeProgramCourseParams) error { + _, err := q.db.Exec(ctx, insertScrapeProgramCourse, arg.ProgramID, arg.CourseID) + return err +} + +const insertScrapeSectionReference = `-- name: InsertScrapeSectionReference :exec +INSERT INTO section_references (parent_section_id, child_section_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING +` + +type InsertScrapeSectionReferenceParams struct { + ParentSectionID int32 + ChildSectionID int32 +} + +func (q *Queries) InsertScrapeSectionReference(ctx context.Context, arg InsertScrapeSectionReferenceParams) error { + _, err := q.db.Exec(ctx, insertScrapeSectionReference, arg.ParentSectionID, arg.ChildSectionID) + return err +} + +const insertScrapeSectionTeacher = `-- name: InsertScrapeSectionTeacher :exec +INSERT INTO section_teachers (section_id, teacher_id) +VALUES ($1, $2) +ON CONFLICT DO NOTHING +` + +type InsertScrapeSectionTeacherParams struct { + SectionID int32 + TeacherID int32 +} + +func (q *Queries) InsertScrapeSectionTeacher(ctx context.Context, arg InsertScrapeSectionTeacherParams) error { + _, err := q.db.Exec(ctx, insertScrapeSectionTeacher, arg.SectionID, arg.TeacherID) + return err +} + +const insertScrapeTeacherProgramLinks = `-- name: InsertScrapeTeacherProgramLinks :exec +INSERT INTO teacher_programs (teacher_id, program_id) +SELECT DISTINCT ct.teacher_id, pc.program_id +FROM course_teachers ct +JOIN program_courses pc ON pc.course_id = ct.course_id +ON CONFLICT DO NOTHING +` + +func (q *Queries) InsertScrapeTeacherProgramLinks(ctx context.Context) error { + _, err := q.db.Exec(ctx, insertScrapeTeacherProgramLinks) + return err +} + +const listScrapeCourseIDs = `-- name: ListScrapeCourseIDs :many +SELECT code, id +FROM course +ORDER BY code +` + +type ListScrapeCourseIDsRow struct { + Code string + ID int64 +} + +func (q *Queries) ListScrapeCourseIDs(ctx context.Context) ([]ListScrapeCourseIDsRow, error) { + rows, err := q.db.Query(ctx, listScrapeCourseIDs) + if err != nil { + return nil, err + } + defer rows.Close() + var items []ListScrapeCourseIDsRow + for rows.Next() { + var i ListScrapeCourseIDsRow + if err := rows.Scan(&i.Code, &i.ID); err != nil { + return nil, err + } + items = append(items, i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const listScrapeTeacherIDsByName = `-- name: ListScrapeTeacherIDsByName :many +SELECT name, id +FROM teacher +ORDER BY name +` + +type ListScrapeTeacherIDsByNameRow struct { + Name string + ID int32 +} + +func (q *Queries) ListScrapeTeacherIDsByName(ctx context.Context) ([]ListScrapeTeacherIDsByNameRow, error) { + rows, err := q.db.Query(ctx, listScrapeTeacherIDsByName) + if err != nil { + return nil, err + } + defer rows.Close() + var items []ListScrapeTeacherIDsByNameRow + for rows.Next() { + var i ListScrapeTeacherIDsByNameRow + if err := rows.Scan(&i.Name, &i.ID); err != nil { + return nil, err + } + items = append(items, i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const listStagedCoursePayloads = `-- name: ListStagedCoursePayloads :many +SELECT payload +FROM staging_courses +WHERE run_id = $1 +ORDER BY course_code +` + +func (q *Queries) ListStagedCoursePayloads(ctx context.Context, runID uuid.UUID) ([][]byte, error) { + rows, err := q.db.Query(ctx, listStagedCoursePayloads, runID) + if err != nil { + return nil, err + } + defer rows.Close() + var items [][]byte + for rows.Next() { + var payload []byte + if err := rows.Scan(&payload); err != nil { + return nil, err + } + items = append(items, payload) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const listStagedProgramPayloads = `-- name: ListStagedProgramPayloads :many +SELECT payload +FROM staging_programs +WHERE run_id = $1 +ORDER BY program_name +` + +func (q *Queries) ListStagedProgramPayloads(ctx context.Context, runID uuid.UUID) ([][]byte, error) { + rows, err := q.db.Query(ctx, listStagedProgramPayloads, runID) + if err != nil { + return nil, err + } + defer rows.Close() + var items [][]byte + for rows.Next() { + var payload []byte + if err := rows.Scan(&payload); err != nil { + return nil, err + } + items = append(items, payload) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const listStagedSchedulePayloads = `-- name: ListStagedSchedulePayloads :many +SELECT payload +FROM staging_schedules +WHERE run_id = $1 +ORDER BY term, course_code +` + +func (q *Queries) ListStagedSchedulePayloads(ctx context.Context, runID uuid.UUID) ([][]byte, error) { + rows, err := q.db.Query(ctx, listStagedSchedulePayloads, runID) + if err != nil { + return nil, err + } + defer rows.Close() + var items [][]byte + for rows.Next() { + var payload []byte + if err := rows.Scan(&payload); err != nil { + return nil, err + } + items = append(items, payload) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const listStagedScheduleTerms = `-- name: ListStagedScheduleTerms :many +SELECT DISTINCT term +FROM staging_schedules +WHERE run_id = $1 +ORDER BY term +` + +func (q *Queries) ListStagedScheduleTerms(ctx context.Context, runID uuid.UUID) ([]string, error) { + rows, err := q.db.Query(ctx, listStagedScheduleTerms, runID) + if err != nil { + return nil, err + } + defer rows.Close() + var items []string + for rows.Next() { + var term string + if err := rows.Scan(&term); err != nil { + return nil, err + } + items = append(items, term) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const listStagedTeacherPayloads = `-- name: ListStagedTeacherPayloads :many +SELECT payload +FROM staging_teachers +WHERE run_id = $1 +ORDER BY rmp_id +` + +func (q *Queries) ListStagedTeacherPayloads(ctx context.Context, runID uuid.UUID) ([][]byte, error) { + rows, err := q.db.Query(ctx, listStagedTeacherPayloads, runID) + if err != nil { + return nil, err + } + defer rows.Close() + var items [][]byte + for rows.Next() { + var payload []byte + if err := rows.Scan(&payload); err != nil { + return nil, err + } + items = append(items, payload) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const markScrapeRunFailed = `-- name: MarkScrapeRunFailed :exec +UPDATE scrape_runs +SET status = 'failed', + failed_at = NOW(), + updated_at = NOW(), + error_message = $2 +WHERE id = $1 +` + +type MarkScrapeRunFailedParams struct { + ID uuid.UUID + ErrorMessage pgtype.Text +} + +func (q *Queries) MarkScrapeRunFailed(ctx context.Context, arg MarkScrapeRunFailedParams) error { + _, err := q.db.Exec(ctx, markScrapeRunFailed, arg.ID, arg.ErrorMessage) + return err +} + +const markScrapeRunPromoting = `-- name: MarkScrapeRunPromoting :exec +UPDATE scrape_runs +SET status = 'promoting', + updated_at = NOW(), + error_message = NULL +WHERE id = $1 +` + +func (q *Queries) MarkScrapeRunPromoting(ctx context.Context, id uuid.UUID) error { + _, err := q.db.Exec(ctx, markScrapeRunPromoting, id) + return err +} + +const markScrapeRunStaged = `-- name: MarkScrapeRunStaged :exec +UPDATE scrape_runs +SET status = 'staged', + course_count = $2, + program_count = $3, + teacher_count = $4, + schedule_term_count = $5, + staged_at = NOW(), + updated_at = NOW(), + error_message = NULL +WHERE id = $1 +` + +type MarkScrapeRunStagedParams struct { + ID uuid.UUID + CourseCount int32 + ProgramCount int32 + TeacherCount int32 + ScheduleTermCount int32 +} + +func (q *Queries) MarkScrapeRunStaged(ctx context.Context, arg MarkScrapeRunStagedParams) error { + _, err := q.db.Exec(ctx, markScrapeRunStaged, + arg.ID, + arg.CourseCount, + arg.ProgramCount, + arg.TeacherCount, + arg.ScheduleTermCount, + ) + return err +} + +const markScrapeRunSucceeded = `-- name: MarkScrapeRunSucceeded :exec +UPDATE scrape_runs +SET status = 'succeeded', + promoted_at = NOW(), + updated_at = NOW(), + error_message = NULL +WHERE id = $1 +` + +func (q *Queries) MarkScrapeRunSucceeded(ctx context.Context, id uuid.UUID) error { + _, err := q.db.Exec(ctx, markScrapeRunSucceeded, id) + return err +} + +const stageCoursePayload = `-- name: StageCoursePayload :exec +INSERT INTO staging_courses (run_id, course_code, payload) +VALUES ($1, $2, $3) +ON CONFLICT (run_id, course_code) DO UPDATE SET + payload = EXCLUDED.payload +` + +type StageCoursePayloadParams struct { + RunID uuid.UUID + CourseCode string + Payload []byte +} + +func (q *Queries) StageCoursePayload(ctx context.Context, arg StageCoursePayloadParams) error { + _, err := q.db.Exec(ctx, stageCoursePayload, arg.RunID, arg.CourseCode, arg.Payload) + return err +} + +const stageProgramPayload = `-- name: StageProgramPayload :exec +INSERT INTO staging_programs (run_id, program_name, payload) +VALUES ($1, $2, $3) +ON CONFLICT (run_id, program_name) DO UPDATE SET + payload = EXCLUDED.payload +` + +type StageProgramPayloadParams struct { + RunID uuid.UUID + ProgramName string + Payload []byte +} + +func (q *Queries) StageProgramPayload(ctx context.Context, arg StageProgramPayloadParams) error { + _, err := q.db.Exec(ctx, stageProgramPayload, arg.RunID, arg.ProgramName, arg.Payload) + return err +} + +const stageSchedulePayload = `-- name: StageSchedulePayload :exec +INSERT INTO staging_schedules (run_id, term, course_code, payload) +VALUES ($1, $2, $3, $4) +ON CONFLICT (run_id, term, course_code) DO UPDATE SET + payload = EXCLUDED.payload +` + +type StageSchedulePayloadParams struct { + RunID uuid.UUID + Term string + CourseCode string + Payload []byte +} + +func (q *Queries) StageSchedulePayload(ctx context.Context, arg StageSchedulePayloadParams) error { + _, err := q.db.Exec(ctx, stageSchedulePayload, + arg.RunID, + arg.Term, + arg.CourseCode, + arg.Payload, + ) + return err +} + +const stageTeacherPayload = `-- name: StageTeacherPayload :exec +INSERT INTO staging_teachers (run_id, rmp_id, payload) +VALUES ($1, $2, $3) +ON CONFLICT (run_id, rmp_id) DO UPDATE SET + payload = EXCLUDED.payload +` + +type StageTeacherPayloadParams struct { + RunID uuid.UUID + RmpID string + Payload []byte +} + +func (q *Queries) StageTeacherPayload(ctx context.Context, arg StageTeacherPayloadParams) error { + _, err := q.db.Exec(ctx, stageTeacherPayload, arg.RunID, arg.RmpID, arg.Payload) + return err +} + +const upsertScrapeCourse = `-- name: UpsertScrapeCourse :exec +INSERT INTO course (code, name, description, restrictions, prerequisites, units, level_number) +VALUES ($1, $2, $3, $4, $5, $6, $7) +ON CONFLICT (code) DO UPDATE SET + name = EXCLUDED.name, + description = EXCLUDED.description, + restrictions = EXCLUDED.restrictions, + prerequisites = EXCLUDED.prerequisites, + units = EXCLUDED.units, + level_number = EXCLUDED.level_number +` + +type UpsertScrapeCourseParams struct { + Code string + Name string + Description string + Restrictions string + Prerequisites []string + Units int32 + LevelNumber pgtype.Int4 +} + +func (q *Queries) UpsertScrapeCourse(ctx context.Context, arg UpsertScrapeCourseParams) error { + _, err := q.db.Exec(ctx, upsertScrapeCourse, + arg.Code, + arg.Name, + arg.Description, + arg.Restrictions, + arg.Prerequisites, + arg.Units, + arg.LevelNumber, + ) + return err +} + +const upsertScrapeProgram = `-- name: UpsertScrapeProgram :one +WITH updated AS ( + UPDATE program + SET source_url = $1, + requirement_codes = $2, + requirements_by_level = $3 + WHERE program.name = $4 + RETURNING id +), +inserted AS ( + INSERT INTO program (name, source_url, requirement_codes, requirements_by_level) + SELECT $4, + $1, + $2, + $3 + WHERE NOT EXISTS (SELECT 1 FROM updated) + RETURNING id +) +SELECT id FROM updated +UNION ALL +SELECT id FROM inserted +LIMIT 1 +` + +type UpsertScrapeProgramParams struct { + SourceUrl pgtype.Text + RequirementCodes []string + RequirementsByLevel []byte + ProgramName string +} + +func (q *Queries) UpsertScrapeProgram(ctx context.Context, arg UpsertScrapeProgramParams) (int64, error) { + row := q.db.QueryRow(ctx, upsertScrapeProgram, + arg.SourceUrl, + arg.RequirementCodes, + arg.RequirementsByLevel, + arg.ProgramName, + ) + var id int64 + err := row.Scan(&id) + return id, err +} + +const upsertScrapeTeacher = `-- name: UpsertScrapeTeacher :exec +INSERT INTO teacher (name, avg_rating, avg_difficulty, department, rmp_id, num_ratings) +VALUES ($1, $2, $3, $4, $5, $6) +ON CONFLICT (rmp_id) DO UPDATE SET + name = EXCLUDED.name, + avg_rating = EXCLUDED.avg_rating, + avg_difficulty = EXCLUDED.avg_difficulty, + department = EXCLUDED.department, + num_ratings = EXCLUDED.num_ratings +` + +type UpsertScrapeTeacherParams struct { + Name string + AvgRating pgtype.Numeric + AvgDifficulty pgtype.Numeric + Department string + RmpID string + NumRatings int32 +} + +func (q *Queries) UpsertScrapeTeacher(ctx context.Context, arg UpsertScrapeTeacherParams) error { + _, err := q.db.Exec(ctx, upsertScrapeTeacher, + arg.Name, + arg.AvgRating, + arg.AvgDifficulty, + arg.Department, + arg.RmpID, + arg.NumRatings, + ) + return err +} diff --git a/server/internal/handlers/scrape_ingest.go b/server/internal/handlers/scrape_ingest.go new file mode 100644 index 0000000..5ef8b1c --- /dev/null +++ b/server/internal/handlers/scrape_ingest.go @@ -0,0 +1,101 @@ +package handlers + +import ( + "net/http" + + "github.com/go-chi/chi/v5" + "github.com/google/uuid" + "github.com/twitocode/pathweave/go-api/internal/common" + "github.com/twitocode/pathweave/go-api/internal/middleware" + "github.com/twitocode/pathweave/go-api/internal/service/scraping" + "go.uber.org/zap" +) + +func HandleCreateScrapeRun(s *scraping.ScrapeIngestService) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + var body scraping.CreateScrapeRunRequest + if err := common.DecodeJSON(r, &body); err != nil { + common.WriteError(w, http.StatusBadRequest, "invalid json in request") + return + } + + run, err := s.CreateRun(r.Context(), body) + if err != nil { + middleware.Logger(r).Error("could not create scrape run", zap.Error(err)) + common.WriteError(w, http.StatusBadRequest, err.Error()) + return + } + + common.WriteJSON(w, http.StatusCreated, run) + } +} + +func HandleGetScrapeRun(s *scraping.ScrapeIngestService) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + runID, ok := scrapeRunIDParam(w, r) + if !ok { + return + } + + run, err := s.GetRun(r.Context(), runID) + if err != nil { + middleware.Logger(r).Error("could not get scrape run", zap.String("run_id", runID.String()), zap.Error(err)) + common.WriteError(w, http.StatusNotFound, "scrape run not found") + return + } + + common.WriteJSON(w, http.StatusOK, run) + } +} + +func HandleStageScrapeArtifacts(s *scraping.ScrapeIngestService) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + runID, ok := scrapeRunIDParam(w, r) + if !ok { + return + } + + var body scraping.StageScrapeArtifactsRequest + if err := common.DecodeJSON(r, &body); err != nil { + common.WriteError(w, http.StatusBadRequest, "invalid json in request") + return + } + + result, err := s.StageArtifacts(r.Context(), runID, body) + if err != nil { + middleware.Logger(r).Error("could not stage scrape artifacts", zap.String("run_id", runID.String()), zap.Error(err)) + common.WriteError(w, http.StatusBadRequest, err.Error()) + return + } + + common.WriteJSON(w, http.StatusOK, result) + } +} + +func HandlePromoteScrapeRun(s *scraping.ScrapeIngestService) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + runID, ok := scrapeRunIDParam(w, r) + if !ok { + return + } + + result, err := s.PromoteRun(r.Context(), runID) + if err != nil { + middleware.Logger(r).Error("could not promote scrape run", zap.String("run_id", runID.String()), zap.Error(err)) + common.WriteError(w, http.StatusInternalServerError, err.Error()) + return + } + + common.WriteJSON(w, http.StatusOK, result) + } +} + +func scrapeRunIDParam(w http.ResponseWriter, r *http.Request) (uuid.UUID, bool) { + value := chi.URLParam(r, "run_id") + runID, err := uuid.Parse(value) + if err != nil { + common.WriteError(w, http.StatusBadRequest, "invalid scrape run id") + return uuid.Nil, false + } + return runID, true +} diff --git a/server/internal/middleware/internal.go b/server/internal/middleware/internal.go new file mode 100644 index 0000000..7b3878a --- /dev/null +++ b/server/internal/middleware/internal.go @@ -0,0 +1,27 @@ +package middleware + +import ( + "net/http" + "strings" + + "github.com/twitocode/pathweave/go-api/internal/common" +) + +func RequireInternalToken(token string) func(http.Handler) http.Handler { + return func(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if token == "" { + common.WriteError(w, http.StatusUnauthorized, "internal service token is not configured") + return + } + + authHeader := strings.TrimSpace(r.Header.Get("Authorization")) + if authHeader != "Bearer "+token { + common.WriteError(w, http.StatusUnauthorized, "invalid internal service token") + return + } + + next.ServeHTTP(w, r) + }) + } +} diff --git a/server/internal/middleware/internal_test.go b/server/internal/middleware/internal_test.go new file mode 100644 index 0000000..a8a4a3b --- /dev/null +++ b/server/internal/middleware/internal_test.go @@ -0,0 +1,39 @@ +package middleware + +import ( + "net/http" + "net/http/httptest" + "testing" + + "github.com/stretchr/testify/require" +) + +func TestRequireInternalTokenAllowsMatchingBearerToken(t *testing.T) { + called := false + handler := RequireInternalToken("secret")(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + called = true + w.WriteHeader(http.StatusNoContent) + })) + + req := httptest.NewRequest(http.MethodPost, "/internal/scrape-runs", nil) + req.Header.Set("Authorization", "Bearer secret") + res := httptest.NewRecorder() + + handler.ServeHTTP(res, req) + + require.True(t, called) + require.Equal(t, http.StatusNoContent, res.Code) +} + +func TestRequireInternalTokenRejectsMissingToken(t *testing.T) { + handler := RequireInternalToken("secret")(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + t.Fatal("handler should not be called") + })) + + req := httptest.NewRequest(http.MethodPost, "/internal/scrape-runs", nil) + res := httptest.NewRecorder() + + handler.ServeHTTP(res, req) + + require.Equal(t, http.StatusUnauthorized, res.Code) +} diff --git a/server/internal/service/scraping/scrape_ingest.go b/server/internal/service/scraping/scrape_ingest.go new file mode 100644 index 0000000..08202df --- /dev/null +++ b/server/internal/service/scraping/scrape_ingest.go @@ -0,0 +1,45 @@ +package scraping + +import ( + "context" + "encoding/json" + "errors" + "strings" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/twitocode/pathweave/go-api/internal/db" + "go.uber.org/zap" +) + +type ScrapeIngestService struct { + db *db.Queries + pool *pgxpool.Pool + log *zap.Logger +} + +func NewScrapeIngestService(queries *db.Queries, pool *pgxpool.Pool, log *zap.Logger) *ScrapeIngestService { + return &ScrapeIngestService{db: queries, pool: pool, log: log} +} + +func (s *ScrapeIngestService) CreateRun(ctx context.Context, req CreateScrapeRunRequest) (db.ScrapeRun, error) { + source := strings.TrimSpace(req.Source) + if source == "" { + source = "manual" + } + metadata := req.Metadata + if len(metadata) == 0 { + metadata = json.RawMessage(`{}`) + } + if !json.Valid(metadata) { + return db.ScrapeRun{}, errors.New("metadata must be valid JSON") + } + return s.db.CreateScrapeRun(ctx, db.CreateScrapeRunParams{ + Source: source, + Metadata: []byte(metadata), + }) +} + +func (s *ScrapeIngestService) GetRun(ctx context.Context, runID uuid.UUID) (db.ScrapeRun, error) { + return s.db.GetScrapeRun(ctx, runID) +} diff --git a/server/internal/service/scraping/scrape_ingest_normalize.go b/server/internal/service/scraping/scrape_ingest_normalize.go new file mode 100644 index 0000000..3ebdb5a --- /dev/null +++ b/server/internal/service/scraping/scrape_ingest_normalize.go @@ -0,0 +1,309 @@ +package scraping + +import ( + "errors" + "fmt" + "regexp" + "sort" + "strings" + "time" + + "github.com/jackc/pgx/v5/pgtype" +) + +func normalizeCourseRecord(course rawCoursePayload) (normalizedCourseRecord, error) { + code := strings.TrimSpace(course.Code) + name := strings.TrimSpace(course.Name) + if code == "" || name == "" { + code, name = splitCourseName(course.CourseName) + } + if code == "" { + return normalizedCourseRecord{}, errors.New("course code is required") + } + level := extractCourseLevelNumber(code) + return normalizedCourseRecord{ + Code: code, + Name: name, + Description: strings.TrimSpace(course.Description), + Restrictions: strings.TrimSpace(course.Restrictions), + Prerequisites: course.Prerequisites, + Units: parseUnits(course.Units), + LevelNumber: level, + }, nil +} + +func splitCourseName(courseName string) (string, string) { + courseName = strings.TrimSpace(courseName) + if courseName == "" { + return "", "" + } + if code, title, ok := strings.Cut(courseName, " - "); ok { + return strings.TrimSpace(code), strings.TrimSpace(title) + } + return "", courseName +} + +var courseCodePattern = regexp.MustCompile(`\b([A-Z]{2,10}\s\d[A-Z0-9]{2,4}(?:\s+A/B)?)\b`) + +func normalizeProgramRequirementCodes(requirements []string) []string { + codes := make([]string, 0, len(requirements)) + seen := make(map[string]struct{}, len(requirements)) + for _, requirement := range requirements { + match := courseCodePattern.FindStringSubmatch(strings.TrimSpace(requirement)) + if len(match) < 2 { + continue + } + code := strings.TrimSpace(match[1]) + if _, ok := seen[code]; ok { + continue + } + seen[code] = struct{}{} + codes = append(codes, code) + } + return codes +} + +func parseUnits(units string) int32 { + match := regexp.MustCompile(`(\d+)`).FindStringSubmatch(units) + if len(match) < 2 { + return 0 + } + var value int32 + _, _ = fmt.Sscanf(match[1], "%d", &value) + return value +} + +func extractCourseLevelNumber(courseCode string) *int32 { + match := regexp.MustCompile(`\b[A-Z]{2,10}\s(\d)`).FindStringSubmatch(courseCode) + if len(match) < 2 { + return nil + } + var value int32 + _, _ = fmt.Sscanf(match[1], "%d", &value) + return &value +} + +func normalizeTerm(term string) string { + term = strings.TrimSpace(term) + if term == "" { + return "Unknown" + } + match := regexp.MustCompile(`^(\d{4})\s+(.+)$`).FindStringSubmatch(term) + if len(match) == 3 { + return strings.TrimSpace(match[2]) + " " + match[1] + } + return term +} + +func parseSectionName(raw string) (string, string) { + raw = strings.TrimSpace(raw) + if raw == "" { + return "", "" + } + match := regexp.MustCompile(`^([A-Z0-9]+)\s*-\s*([A-Z]+)(?:\s+\(\d+\))?$`).FindStringSubmatch(strings.ToUpper(raw)) + if len(match) == 3 { + return match[2] + " " + match[1], match[2] + } + parts := strings.Fields(raw) + if len(parts) > 0 { + return raw, strings.ToUpper(parts[0]) + } + return raw, "" +} + +func parseScrapeTime(value string) (string, error) { + value = strings.TrimSpace(value) + if value == "" || strings.EqualFold(value, "TBA") { + return "", errors.New("time is empty") + } + for _, layout := range []string{"3:04PM", "3:04 PM"} { + parsed, err := time.Parse(layout, value) + if err == nil { + return parsed.Format("15:04:05"), nil + } + } + return "", fmt.Errorf("invalid time %q", value) +} + +func parseNullableScrapeTime(value string) (pgtype.Time, error) { + value = strings.TrimSpace(value) + if value == "" || strings.EqualFold(value, "TBA") { + return pgtype.Time{}, nil + } + normalized, err := parseScrapeTime(value) + if err != nil { + return pgtype.Time{}, err + } + parsed, err := time.Parse("15:04:05", normalized) + if err != nil { + return pgtype.Time{}, err + } + return pgtype.Time{ + Microseconds: int64(parsed.Hour()*3600+parsed.Minute()*60+parsed.Second()) * 1_000_000, + Valid: true, + }, nil +} + +func parseScrapeLocation(value string) (string, string) { + raw := strings.TrimSpace(value) + if raw == "" { + return "", "" + } + upper := strings.ToUpper(raw) + if upper == "IN PERSON" || upper == "IN-PERSON" { + return "", "" + } + if strings.Contains(upper, "TBD") || strings.Contains(upper, "TBA") || + strings.Contains(upper, "ANNOUNCED") || strings.Contains(upper, "SEE CLASS NOTES") { + return "TBD", "TBD" + } + if strings.Contains(upper, "ONLINE") || strings.Contains(upper, "VIRTUAL") { + return "Online", "Online" + } + parts := strings.Fields(raw) + if len(parts) >= 2 && regexp.MustCompile(`^[A-Z0-9]+$`).MatchString(parts[0]) { + room := strings.TrimSpace(strings.Join(parts[1:], " ")) + room = strings.TrimSpace(regexp.MustCompile(`(?i)lab`).ReplaceAllString(room, "")) + return parts[0], room + } + return raw, "" +} + +func getInstructorNames(value string) []string { + cleaned := strings.ReplaceAll(value, "\u00a0", " ") + names := make([]string, 0) + for _, line := range strings.Split(cleaned, "\n") { + for _, part := range strings.Split(line, ",") { + name := strings.Join(strings.Fields(part), " ") + if name != "" { + names = append(names, name) + } + } + } + return names +} + +func getAllInstructorNames(section rawSectionPayload) []string { + seen := make(map[string]struct{}) + for _, detail := range section.Details { + for _, name := range getInstructorNames(detail.Instructor) { + seen[name] = struct{}{} + } + } + names := make([]string, 0, len(seen)) + for name := range seen { + names = append(names, name) + } + sort.Strings(names) + return names +} + +func getSectionInstructorSet(section rawSectionPayload) map[string]struct{} { + names := make(map[string]struct{}) + for _, name := range getAllInstructorNames(section) { + if !strings.EqualFold(name, "Staff") { + names[name] = struct{}{} + } + } + return names +} + +func detectDeliveryMode(section rawSectionPayload) (string, bool) { + hasOnline := false + hasInPerson := false + for _, detail := range section.Details { + room := strings.ToUpper(strings.TrimSpace(detail.Room)) + switch { + case strings.Contains(room, "ONLINE") || strings.Contains(room, "VIRTUAL"): + hasOnline = true + case room == "IN PERSON" || room == "IN-PERSON": + hasInPerson = true + case room != "" && room != "TBA" && room != "TBD": + hasInPerson = true + } + } + switch { + case hasOnline && hasInPerson: + return "Blended", true + case hasOnline: + return "Online", false + case hasInPerson: + return "In Person", true + default: + return "Unknown", false + } +} + +func buildSectionReferences(sections []normalizedSection) []sectionReference { + parents := make([]normalizedSection, 0) + children := make([]normalizedSection, 0) + for _, section := range sections { + switch section.Type { + case "LEC", "SEM": + parents = append(parents, section) + case "LAB", "TUT": + children = append(children, section) + } + } + refs := make([]sectionReference, 0) + for _, parent := range parents { + for _, child := range children { + if len(child.InstructorSet) == 0 || len(parent.InstructorSet) == 0 || intersects(parent.InstructorSet, child.InstructorSet) { + refs = append(refs, sectionReference{ParentID: parent.ID, ChildID: child.ID}) + } + } + } + return refs +} + +func intersects(a, b map[string]struct{}) bool { + for key := range a { + if _, ok := b[key]; ok { + return true + } + } + return false +} + +func collectScheduleTeacherNames(schedules []rawScheduleCoursePayload) map[string]struct{} { + names := make(map[string]struct{}) + for _, course := range schedules { + for _, section := range course.Sections { + for _, name := range getAllInstructorNames(section) { + names[name] = struct{}{} + } + } + } + return names +} + +func resolveScheduleCourseCode(course rawScheduleCoursePayload, titleCodeMap map[string]string) string { + code := strings.TrimSpace(course.CourseCode) + if courseCodePattern.MatchString(code) { + return code + } + if resolved := titleCodeMap[strings.TrimSpace(course.CourseTitle)]; resolved != "" { + return resolved + } + return code +} + +func nonRMPID(name string) string { + id := strings.ToLower(strings.TrimSpace(name)) + id = regexp.MustCompile(`\s+`).ReplaceAllString(id, "_") + return "non_rmp_" + id +} + +func nullableText(value string) pgtype.Text { + value = strings.TrimSpace(value) + if value == "" { + return pgtype.Text{} + } + return pgtype.Text{String: value, Valid: true} +} + +func numericFromFloat64(value float64) pgtype.Numeric { + var numeric pgtype.Numeric + _ = numeric.Scan(fmt.Sprintf("%g", value)) + return numeric +} diff --git a/server/internal/service/scraping/scrape_ingest_payloads.go b/server/internal/service/scraping/scrape_ingest_payloads.go new file mode 100644 index 0000000..713d61f --- /dev/null +++ b/server/internal/service/scraping/scrape_ingest_payloads.go @@ -0,0 +1,106 @@ +package scraping + +import ( + "encoding/json" + + "github.com/google/uuid" +) + +type CreateScrapeRunRequest struct { + Source string `json:"source"` + Metadata json.RawMessage `json:"metadata"` +} + +type StageScrapeArtifactsRequest struct { + Courses []rawCoursePayload `json:"courses"` + Programs []rawProgramPayload `json:"programs"` + Teachers []rawTeacherPayload `json:"teachers"` + Schedules []rawScheduleTermPayload `json:"schedules"` +} + +type StageScrapeArtifactsResult struct { + RunID uuid.UUID `json:"run_id"` + CourseCount int `json:"course_count"` + ProgramCount int `json:"program_count"` + TeacherCount int `json:"teacher_count"` + ScheduleTermCount int `json:"schedule_term_count"` +} + +type PromoteScrapeRunResult struct { + RunID uuid.UUID `json:"run_id"` + Status string `json:"status"` + Promoted bool `json:"promoted"` +} + +type rawCoursePayload struct { + Code string `json:"code"` + Name string `json:"name"` + CourseName string `json:"course_name"` + Units string `json:"units"` + Description string `json:"description"` + Restrictions string `json:"restrictions"` + Prerequisites []string `json:"prerequisites"` +} + +type rawProgramPayload struct { + ProgramName string `json:"program_name"` + URL string `json:"url"` + Requirements []string `json:"requirements"` + RequirementsByLevel json.RawMessage `json:"requirements_by_level"` +} + +type rawTeacherPayload struct { + ID string `json:"id"` + Name string `json:"name"` + AvgRating float64 `json:"avgRating"` + AvgDifficulty float64 `json:"avgDifficulty"` + Department string `json:"department"` + NumRatings int32 `json:"numRatings"` + Courses []string `json:"courses"` +} + +type rawScheduleTermPayload struct { + Term string `json:"term"` + Courses []rawScheduleCoursePayload `json:"courses"` +} + +type rawScheduleCoursePayload struct { + CourseCode string `json:"course_code"` + CourseTitle string `json:"course_title"` + Term string `json:"term"` + Sections []rawSectionPayload `json:"sections"` +} + +type rawSectionPayload struct { + SectionName string `json:"section_name"` + Details []rawMeetingPayload `json:"details"` +} + +type rawMeetingPayload struct { + Days string `json:"days"` + StartTime string `json:"start_time"` + EndTime string `json:"end_time"` + Room string `json:"room"` + Instructor string `json:"instructor"` +} + +type normalizedCourseRecord struct { + Code string + Name string + Description string + Restrictions string + Prerequisites []string + Units int32 + LevelNumber *int32 +} + +type normalizedSection struct { + ID int32 + Type string + InstructorSet map[string]struct{} +} + +type sectionReference struct { + ParentID int32 + ChildID int32 +} diff --git a/server/internal/service/scraping/scrape_ingest_promote.go b/server/internal/service/scraping/scrape_ingest_promote.go new file mode 100644 index 0000000..af86446 --- /dev/null +++ b/server/internal/service/scraping/scrape_ingest_promote.go @@ -0,0 +1,449 @@ +package scraping + +import ( + "context" + "encoding/json" + "errors" + "strings" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgtype" + "github.com/twitocode/pathweave/go-api/internal/db" + "go.uber.org/zap" +) + +func (s *ScrapeIngestService) PromoteRun(ctx context.Context, runID uuid.UUID) (PromoteScrapeRunResult, error) { + log := s.log.With(zap.String("run_id", runID.String())) + log.Info("scrape promote started") + + tx, err := s.pool.Begin(ctx) + if err != nil { + log.Error("scrape promote failed to begin transaction", zap.Error(err)) + return PromoteScrapeRunResult{}, err + } + + qtx := s.db.WithTx(tx) + if err := qtx.MarkScrapeRunPromoting(ctx, runID); err != nil { + _ = tx.Rollback(ctx) + log.Error("scrape promote failed to mark run promoting", zap.Error(err)) + return PromoteScrapeRunResult{}, err + } + if err := s.promoteRunInTx(ctx, qtx, runID, log); err != nil { + _ = tx.Rollback(ctx) + _ = s.db.MarkScrapeRunFailed(ctx, db.MarkScrapeRunFailedParams{ + ID: runID, + ErrorMessage: pgtype.Text{String: err.Error(), Valid: true}, + }) + log.Error("scrape promote failed", zap.Error(err)) + return PromoteScrapeRunResult{}, err + } + if err := qtx.MarkScrapeRunSucceeded(ctx, runID); err != nil { + _ = tx.Rollback(ctx) + log.Error("scrape promote failed to mark run succeeded", zap.Error(err)) + return PromoteScrapeRunResult{}, err + } + if err := tx.Commit(ctx); err != nil { + log.Error("scrape promote failed to commit transaction", zap.Error(err)) + return PromoteScrapeRunResult{}, err + } + log.Info("scrape promote completed") + return PromoteScrapeRunResult{RunID: runID, Status: "succeeded", Promoted: true}, nil +} + +func (s *ScrapeIngestService) promoteRunInTx(ctx context.Context, qtx *db.Queries, runID uuid.UUID, log *zap.Logger) error { + log.Info("scrape promote loading staged payloads") + courses, err := loadPayloads[rawCoursePayload](ctx, qtx.ListStagedCoursePayloads, runID) + if err != nil { + return err + } + programs, err := loadPayloads[rawProgramPayload](ctx, qtx.ListStagedProgramPayloads, runID) + if err != nil { + return err + } + teachers, err := loadPayloads[rawTeacherPayload](ctx, qtx.ListStagedTeacherPayloads, runID) + if err != nil { + return err + } + schedules, err := loadPayloads[rawScheduleCoursePayload](ctx, qtx.ListStagedSchedulePayloads, runID) + if err != nil { + return err + } + terms, err := qtx.ListStagedScheduleTerms(ctx, runID) + if err != nil { + return err + } + log.Info("scrape promote loaded staged payloads", + zap.Int("courses", len(courses)), + zap.Int("programs", len(programs)), + zap.Int("teachers", len(teachers)), + zap.Int("schedule_courses", len(schedules)), + zap.Strings("terms", terms), + ) + + log.Info("scrape promote upserting courses", zap.Int("count", len(courses))) + normalizedCourses, titleCodeMap, err := promoteCourses(ctx, qtx, courses) + if err != nil { + return err + } + log.Info("scrape promote upserted courses", zap.Int("count", len(normalizedCourses))) + + log.Info("scrape promote loading course ids") + courseIDs, err := loadCourseIDs(ctx, qtx) + if err != nil { + return err + } + log.Info("scrape promote loaded course ids", zap.Int("count", len(courseIDs))) + + scheduleTeacherNames := collectScheduleTeacherNames(schedules) + log.Info("scrape promote upserting teachers", + zap.Int("rmp_teachers", len(teachers)), + zap.Int("schedule_teacher_names", len(scheduleTeacherNames)), + ) + if err := promoteTeachers(ctx, qtx, teachers, scheduleTeacherNames); err != nil { + return err + } + log.Info("scrape promote upserted teachers") + + log.Info("scrape promote loading teacher ids") + teacherIDs, err := loadTeacherIDsByName(ctx, qtx) + if err != nil { + return err + } + log.Info("scrape promote loaded teacher ids", zap.Int("count", len(teacherIDs))) + + log.Info("scrape promote upserting programs", zap.Int("count", len(programs))) + if err := promotePrograms(ctx, qtx, programs, courseIDs); err != nil { + return err + } + log.Info("scrape promote upserted programs", zap.Int("count", len(programs))) + + log.Info("scrape promote linking course teachers") + if err := promoteCourseTeacherLinks(ctx, qtx, teachers, schedules, courseIDs, teacherIDs, titleCodeMap); err != nil { + return err + } + log.Info("scrape promote linked course teachers") + + log.Info("scrape promote deleting sections for terms", zap.Strings("terms", terms)) + if err := deleteSectionsForTerms(ctx, qtx, terms); err != nil { + return err + } + log.Info("scrape promote deleted sections for terms") + + log.Info("scrape promote upserting sections", zap.Int("schedule_courses", len(schedules))) + if err := promoteSections(ctx, qtx, schedules, courseIDs, teacherIDs, titleCodeMap); err != nil { + return err + } + log.Info("scrape promote upserted sections") + + log.Info("scrape promote linking teacher programs") + if err := promoteTeacherProgramLinks(ctx, qtx); err != nil { + return err + } + log.Info("scrape promote linked teacher programs") + return nil +} + +func clearRunStaging(ctx context.Context, q *db.Queries, runID uuid.UUID) error { + if err := q.ClearScrapeRunStagedSchedules(ctx, runID); err != nil { + return err + } + if err := q.ClearScrapeRunStagedTeachers(ctx, runID); err != nil { + return err + } + if err := q.ClearScrapeRunStagedPrograms(ctx, runID); err != nil { + return err + } + return q.ClearScrapeRunStagedCourses(ctx, runID) +} + +func loadPayloads[T any](ctx context.Context, list func(context.Context, uuid.UUID) ([][]byte, error), runID uuid.UUID) ([]T, error) { + payloads, err := list(ctx, runID) + if err != nil { + return nil, err + } + items := make([]T, 0, len(payloads)) + for _, payload := range payloads { + var item T + if err := json.Unmarshal(payload, &item); err != nil { + return nil, err + } + items = append(items, item) + } + return items, nil +} + +func promoteCourses(ctx context.Context, q *db.Queries, courses []rawCoursePayload) ([]normalizedCourseRecord, map[string]string, error) { + records := make([]normalizedCourseRecord, 0, len(courses)) + titleCodeMap := make(map[string]string, len(courses)) + for _, course := range courses { + record, err := normalizeCourseRecord(course) + if err != nil { + return nil, nil, err + } + records = append(records, record) + if record.Name != "" { + titleCodeMap[record.Name] = record.Code + } + var level pgtype.Int4 + if record.LevelNumber != nil { + level = pgtype.Int4{Int32: *record.LevelNumber, Valid: true} + } + if err := q.UpsertScrapeCourse(ctx, db.UpsertScrapeCourseParams{ + Code: record.Code, + Name: record.Name, + Description: record.Description, + Restrictions: record.Restrictions, + Prerequisites: record.Prerequisites, + Units: record.Units, + LevelNumber: level, + }); err != nil { + return nil, nil, err + } + } + return records, titleCodeMap, nil +} + +func promoteTeachers(ctx context.Context, q *db.Queries, teachers []rawTeacherPayload, scheduleTeacherNames map[string]struct{}) error { + knownNames := make(map[string]struct{}, len(teachers)) + for _, teacher := range teachers { + name := strings.TrimSpace(teacher.Name) + if name == "" { + continue + } + rmpID := strings.TrimSpace(teacher.ID) + if rmpID == "" { + return errors.New("teacher id is required") + } + department := strings.TrimSpace(teacher.Department) + if department == "" { + department = "Unknown" + } + if err := q.UpsertScrapeTeacher(ctx, db.UpsertScrapeTeacherParams{ + Name: name, + AvgRating: numericFromFloat64(teacher.AvgRating), + AvgDifficulty: numericFromFloat64(teacher.AvgDifficulty), + Department: department, + RmpID: rmpID, + NumRatings: teacher.NumRatings, + }); err != nil { + return err + } + knownNames[strings.ToLower(name)] = struct{}{} + } + + for name := range scheduleTeacherNames { + if _, ok := knownNames[strings.ToLower(name)]; ok { + continue + } + rmpID := nonRMPID(name) + if err := q.UpsertScrapeTeacher(ctx, db.UpsertScrapeTeacherParams{ + Name: name, + AvgRating: numericFromFloat64(0), + AvgDifficulty: numericFromFloat64(0), + Department: "Unknown", + RmpID: rmpID, + NumRatings: 0, + }); err != nil { + return err + } + } + return nil +} + +func promotePrograms(ctx context.Context, q *db.Queries, programs []rawProgramPayload, courseIDs map[string]int64) error { + for _, program := range programs { + name := strings.TrimSpace(program.ProgramName) + if name == "" { + return errors.New("program_name is required") + } + requirementsByLevel := program.RequirementsByLevel + if len(requirementsByLevel) == 0 { + requirementsByLevel = json.RawMessage(`[]`) + } + requirementCodes := normalizeProgramRequirementCodes(program.Requirements) + programID, err := q.UpsertScrapeProgram(ctx, db.UpsertScrapeProgramParams{ + ProgramName: name, + SourceUrl: nullableText(program.URL), + RequirementCodes: requirementCodes, + RequirementsByLevel: []byte(requirementsByLevel), + }) + if err != nil { + return err + } + if err := q.DeleteScrapeProgramCourses(ctx, programID); err != nil { + return err + } + for _, code := range requirementCodes { + courseID, ok := courseIDs[code] + if !ok { + continue + } + if err := q.InsertScrapeProgramCourse(ctx, db.InsertScrapeProgramCourseParams{ + ProgramID: programID, + CourseID: courseID, + }); err != nil { + return err + } + } + } + return nil +} + +func promoteCourseTeacherLinks(ctx context.Context, q *db.Queries, teachers []rawTeacherPayload, schedules []rawScheduleCoursePayload, courseIDs map[string]int64, teacherIDs map[string]int32, titleCodeMap map[string]string) error { + for _, teacher := range teachers { + teacherID, ok := teacherIDs[strings.TrimSpace(teacher.Name)] + if !ok { + continue + } + for _, code := range teacher.Courses { + courseID, ok := courseIDs[strings.TrimSpace(code)] + if !ok { + continue + } + if err := q.InsertScrapeCourseTeacher(ctx, db.InsertScrapeCourseTeacherParams{ + CourseID: courseID, + TeacherID: teacherID, + }); err != nil { + return err + } + } + } + + for _, course := range schedules { + courseID, ok := courseIDs[resolveScheduleCourseCode(course, titleCodeMap)] + if !ok { + continue + } + for _, section := range course.Sections { + for _, name := range getAllInstructorNames(section) { + teacherID, ok := teacherIDs[name] + if !ok { + continue + } + if err := q.InsertScrapeCourseTeacher(ctx, db.InsertScrapeCourseTeacherParams{ + CourseID: courseID, + TeacherID: teacherID, + }); err != nil { + return err + } + } + } + } + return nil +} + +func deleteSectionsForTerms(ctx context.Context, q *db.Queries, terms []string) error { + if len(terms) == 0 { + return nil + } + return q.DeleteScrapeSectionsForTerms(ctx, terms) +} + +func promoteSections(ctx context.Context, q *db.Queries, schedules []rawScheduleCoursePayload, courseIDs map[string]int64, teacherIDs map[string]int32, titleCodeMap map[string]string) error { + for _, course := range schedules { + term := normalizeTerm(course.Term) + courseID, ok := courseIDs[resolveScheduleCourseCode(course, titleCodeMap)] + if !ok { + continue + } + insertedSections := make([]normalizedSection, 0, len(course.Sections)) + for _, section := range course.Sections { + name, sectionType := parseSectionName(section.SectionName) + if name == "" { + continue + } + mode, isInPerson := detectDeliveryMode(section) + sectionID, err := q.CreateScrapeSection(ctx, db.CreateScrapeSectionParams{ + CourseID: courseID, + Name: name, + Type: sectionType, + Term: term, + Mode: mode, + IsInPerson: isInPerson, + }) + if err != nil { + return err + } + for _, detail := range section.Details { + startTime, err := parseNullableScrapeTime(detail.StartTime) + if err != nil { + return err + } + endTime, err := parseNullableScrapeTime(detail.EndTime) + if err != nil { + return err + } + building, room := parseScrapeLocation(detail.Room) + if err := q.CreateScrapeSectionMeeting(ctx, db.CreateScrapeSectionMeetingParams{ + SectionID: sectionID, + Days: strings.TrimSpace(detail.Days), + StartTime: startTime, + EndTime: endTime, + Building: building, + Room: room, + }); err != nil { + return err + } + } + linkedTeacherIDs := make(map[int32]struct{}) + for _, instructor := range getAllInstructorNames(section) { + teacherID, ok := teacherIDs[instructor] + if !ok { + continue + } + if _, seen := linkedTeacherIDs[teacherID]; seen { + continue + } + linkedTeacherIDs[teacherID] = struct{}{} + if err := q.InsertScrapeSectionTeacher(ctx, db.InsertScrapeSectionTeacherParams{ + SectionID: sectionID, + TeacherID: teacherID, + }); err != nil { + return err + } + } + insertedSections = append(insertedSections, normalizedSection{ + ID: sectionID, + Type: sectionType, + InstructorSet: getSectionInstructorSet(section), + }) + } + for _, ref := range buildSectionReferences(insertedSections) { + if err := q.InsertScrapeSectionReference(ctx, db.InsertScrapeSectionReferenceParams{ + ParentSectionID: ref.ParentID, + ChildSectionID: ref.ChildID, + }); err != nil { + return err + } + } + } + return nil +} + +func promoteTeacherProgramLinks(ctx context.Context, q *db.Queries) error { + return q.InsertScrapeTeacherProgramLinks(ctx) +} + +func loadCourseIDs(ctx context.Context, q *db.Queries) (map[string]int64, error) { + rows, err := q.ListScrapeCourseIDs(ctx) + if err != nil { + return nil, err + } + ids := make(map[string]int64, len(rows)) + for _, row := range rows { + ids[row.Code] = row.ID + } + return ids, nil +} + +func loadTeacherIDsByName(ctx context.Context, q *db.Queries) (map[string]int32, error) { + rows, err := q.ListScrapeTeacherIDsByName(ctx) + if err != nil { + return nil, err + } + ids := make(map[string]int32, len(rows)) + for _, row := range rows { + ids[row.Name] = row.ID + } + return ids, nil +} diff --git a/server/internal/service/scraping/scrape_ingest_stage.go b/server/internal/service/scraping/scrape_ingest_stage.go new file mode 100644 index 0000000..c12cbb0 --- /dev/null +++ b/server/internal/service/scraping/scrape_ingest_stage.go @@ -0,0 +1,127 @@ +package scraping + +import ( + "context" + "encoding/json" + "errors" + "strings" + + "github.com/google/uuid" + "github.com/twitocode/pathweave/go-api/internal/db" +) + +func (s *ScrapeIngestService) StageArtifacts(ctx context.Context, runID uuid.UUID, req StageScrapeArtifactsRequest) (StageScrapeArtifactsResult, error) { + tx, err := s.pool.Begin(ctx) + if err != nil { + return StageScrapeArtifactsResult{}, err + } + defer func() { _ = tx.Rollback(ctx) }() + + qtx := s.db.WithTx(tx) + if err := clearRunStaging(ctx, qtx, runID); err != nil { + return StageScrapeArtifactsResult{}, err + } + + for _, course := range req.Courses { + record, err := normalizeCourseRecord(course) + if err != nil { + return StageScrapeArtifactsResult{}, err + } + payload, err := json.Marshal(course) + if err != nil { + return StageScrapeArtifactsResult{}, err + } + if err := qtx.StageCoursePayload(ctx, db.StageCoursePayloadParams{ + RunID: runID, + CourseCode: record.Code, + Payload: payload, + }); err != nil { + return StageScrapeArtifactsResult{}, err + } + } + + for _, program := range req.Programs { + name := strings.TrimSpace(program.ProgramName) + if name == "" { + return StageScrapeArtifactsResult{}, errors.New("program_name is required") + } + payload, err := json.Marshal(program) + if err != nil { + return StageScrapeArtifactsResult{}, err + } + if err := qtx.StageProgramPayload(ctx, db.StageProgramPayloadParams{ + RunID: runID, + ProgramName: name, + Payload: payload, + }); err != nil { + return StageScrapeArtifactsResult{}, err + } + } + + for _, teacher := range req.Teachers { + rmpID := strings.TrimSpace(teacher.ID) + if rmpID == "" { + return StageScrapeArtifactsResult{}, errors.New("teacher id is required") + } + payload, err := json.Marshal(teacher) + if err != nil { + return StageScrapeArtifactsResult{}, err + } + if err := qtx.StageTeacherPayload(ctx, db.StageTeacherPayloadParams{ + RunID: runID, + RmpID: rmpID, + Payload: payload, + }); err != nil { + return StageScrapeArtifactsResult{}, err + } + } + + terms := make(map[string]struct{}) + for _, termPayload := range req.Schedules { + term := normalizeTerm(termPayload.Term) + if term == "Unknown" { + return StageScrapeArtifactsResult{}, errors.New("schedule term is required") + } + terms[term] = struct{}{} + for _, course := range termPayload.Courses { + if course.Term == "" { + course.Term = term + } + courseCode := strings.TrimSpace(course.CourseCode) + if courseCode == "" { + return StageScrapeArtifactsResult{}, errors.New("schedule course_code is required") + } + payload, err := json.Marshal(course) + if err != nil { + return StageScrapeArtifactsResult{}, err + } + if err := qtx.StageSchedulePayload(ctx, db.StageSchedulePayloadParams{ + RunID: runID, + Term: term, + CourseCode: courseCode, + Payload: payload, + }); err != nil { + return StageScrapeArtifactsResult{}, err + } + } + } + + result := StageScrapeArtifactsResult{ + RunID: runID, + CourseCount: len(req.Courses), + ProgramCount: len(req.Programs), + TeacherCount: len(req.Teachers), + ScheduleTermCount: len(terms), + } + if err := qtx.MarkScrapeRunStaged(ctx, db.MarkScrapeRunStagedParams{ + ID: runID, + CourseCount: int32(result.CourseCount), + ProgramCount: int32(result.ProgramCount), + TeacherCount: int32(result.TeacherCount), + ScheduleTermCount: int32(result.ScheduleTermCount), + }); err != nil { + return StageScrapeArtifactsResult{}, err + } + + return result, tx.Commit(ctx) +} diff --git a/server/internal/service/scraping/scrape_ingest_test.go b/server/internal/service/scraping/scrape_ingest_test.go new file mode 100644 index 0000000..523aec0 --- /dev/null +++ b/server/internal/service/scraping/scrape_ingest_test.go @@ -0,0 +1,51 @@ +package scraping + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestScrapeIngestNormalizeCourseRecord(t *testing.T) { + record, err := normalizeCourseRecord(rawCoursePayload{ + CourseName: "COMPSCI 1DM3 - Discrete Mathematics for Computer Science", + Units: "3 unit(s)", + Description: "Sets and logic", + Restrictions: "Antirequisite(s): MATH 1DM3", + Prerequisites: []string{"MATH 1ZA3"}, + }) + require.NoError(t, err) + require.Equal(t, "COMPSCI 1DM3", record.Code) + require.Equal(t, "Discrete Mathematics for Computer Science", record.Name) + require.Equal(t, int32(3), record.Units) + require.NotNil(t, record.LevelNumber) + require.Equal(t, int32(1), *record.LevelNumber) +} + +func TestScrapeIngestParseMeetingDetails(t *testing.T) { + start, err := parseScrapeTime("4:30PM") + require.NoError(t, err) + require.Equal(t, "16:30:00", start) + + building, room := parseScrapeLocation("BSB B156") + require.Equal(t, "BSB", building) + require.Equal(t, "B156", room) + + building, room = parseScrapeLocation("Online") + require.Equal(t, "Online", building) + require.Equal(t, "Online", room) +} + +func TestScrapeIngestBuildSectionReferences(t *testing.T) { + sections := []normalizedSection{ + {ID: 1, Type: "LEC", InstructorSet: map[string]struct{}{"Jane Doe": {}}}, + {ID: 2, Type: "LAB", InstructorSet: map[string]struct{}{"Jane Doe": {}}}, + {ID: 3, Type: "TUT", InstructorSet: map[string]struct{}{}}, + {ID: 4, Type: "LAB", InstructorSet: map[string]struct{}{"Other Person": {}}}, + } + + refs := buildSectionReferences(sections) + require.Len(t, refs, 2) + require.Equal(t, sectionReference{ParentID: 1, ChildID: 2}, refs[0]) + require.Equal(t, sectionReference{ParentID: 1, ChildID: 3}, refs[1]) +}