From 5cab1cf1e2eb886eddf8fdcbc1f20818f8fa4966 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 18 Jul 2026 18:01:08 +0530 Subject: [PATCH 01/27] fix(harvester): address repository client review feedback --- application/utils/harvester/git_repository_client.py | 1 + 1 file changed, 1 insertion(+) diff --git a/application/utils/harvester/git_repository_client.py b/application/utils/harvester/git_repository_client.py index 1390e6606..a37deb798 100644 --- a/application/utils/harvester/git_repository_client.py +++ b/application/utils/harvester/git_repository_client.py @@ -160,6 +160,7 @@ def checkout(self, reference: str) -> None: "-C", str(self.local_path), "checkout", + "--", reference, ], check=True, From 2f981e21a99fc3e351ba3ab7f3c3964af573004d Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Wed, 8 Jul 2026 13:02:05 +0530 Subject: [PATCH 02/27] feat(harvester): implement incremental change detection --- .../harvester_test/change_detector_test.py | 52 +++++++++++ .../harvester_test/checkpoint_store_test.py | 92 +++++++++++++++++++ .../utils/harvester/change_detector.py | 73 +++++++++++++++ .../utils/harvester/checkpoint_store.py | 54 +++++++++++ application/utils/harvester/models.py | 16 ++++ 5 files changed, 287 insertions(+) create mode 100644 application/tests/harvester_test/change_detector_test.py create mode 100644 application/tests/harvester_test/checkpoint_store_test.py create mode 100644 application/utils/harvester/change_detector.py create mode 100644 application/utils/harvester/checkpoint_store.py create mode 100644 application/utils/harvester/models.py diff --git a/application/tests/harvester_test/change_detector_test.py b/application/tests/harvester_test/change_detector_test.py new file mode 100644 index 000000000..a74590fd8 --- /dev/null +++ b/application/tests/harvester_test/change_detector_test.py @@ -0,0 +1,52 @@ +import unittest +from unittest.mock import MagicMock +from unittest.mock import patch + +from application.utils.harvester.change_detector import ( + ChangeDetector, +) + + +class ChangeDetectorTests(unittest.TestCase): + @patch("application.utils.harvester.change_detector.subprocess.run") + def test_get_modified_files_since(self, mock_run): + mock_run.return_value = MagicMock( + stdout="a.md\nb.md\na.md\n", + ) + + client = MagicMock() + detector = ChangeDetector(client) + + files = detector.get_modified_files_since("abc123") + + self.assertEqual( + files, + [ + "a.md", + "b.md", + ], + ) + + @patch("application.utils.harvester.change_detector.subprocess.run") + def test_get_commits_since(self, mock_run): + mock_run.return_value = MagicMock( + stdout="111\n222\n333\n", + ) + + client = MagicMock() + detector = ChangeDetector(client) + + commits = detector.get_commits_since("abc123") + + self.assertEqual( + commits, + [ + "111", + "222", + "333", + ], + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/tests/harvester_test/checkpoint_store_test.py b/application/tests/harvester_test/checkpoint_store_test.py new file mode 100644 index 000000000..9859b1fed --- /dev/null +++ b/application/tests/harvester_test/checkpoint_store_test.py @@ -0,0 +1,92 @@ +import unittest +from datetime import datetime +from pathlib import Path + +from application.utils.harvester.checkpoint_store import ( + CheckpointStore, +) +from application.utils.harvester.models import ( + RepositoryCheckpoint, +) + + +class CheckpointStoreTests(unittest.TestCase): + def test_save_and_load_checkpoint(self): + tmp_dir = Path(self._testMethodName) + + try: + store = CheckpointStore( + tmp_dir / "checkpoints.json", + ) + + checkpoint = RepositoryCheckpoint( + repository_id="owasp-asvs", + last_processed_commit="abc123", + updated_at=datetime.now(), + ) + + store.save(checkpoint) + + loaded = store.load("owasp-asvs") + + if loaded is None: + self.fail("Checkpoint should have been loaded") + + self.assertEqual( + loaded.last_processed_commit, + "abc123", + ) + + finally: + if tmp_dir.exists(): + import shutil + + shutil.rmtree(tmp_dir) + + def test_load_missing_file(self): + tmp_dir = Path(self._testMethodName) + + try: + store = CheckpointStore( + tmp_dir / "missing.json", + ) + + self.assertIsNone( + store.load("repo"), + ) + + finally: + if tmp_dir.exists(): + import shutil + + shutil.rmtree(tmp_dir) + + def test_load_missing_repository(self): + tmp_dir = Path(self._testMethodName) + + try: + store = CheckpointStore( + tmp_dir / "checkpoint.json", + ) + + store.save( + RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="abc123", + updated_at=datetime.now(), + ) + ) + + self.assertIsNone( + store.load("repo-b"), + ) + + finally: + if tmp_dir.exists(): + import shutil + + shutil.rmtree(tmp_dir) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/utils/harvester/change_detector.py b/application/utils/harvester/change_detector.py new file mode 100644 index 000000000..8b0153492 --- /dev/null +++ b/application/utils/harvester/change_detector.py @@ -0,0 +1,73 @@ +import logging +import subprocess + +from .git_repository_client import GitRepositoryClient + +logger = logging.getLogger(__name__) + + +class ChangeDetector: + def __init__(self, repository_client: GitRepositoryClient): + self.repository_client = repository_client + + def get_modified_files_since(self, commit_sha: str) -> list[str]: + logger.info( + "Detecting changes since commit %s", + commit_sha, + ) + + try: + result = subprocess.run( + [ + "git", + "-C", + str(self.repository_client.get_local_path()), + "diff", + "--name-only", + commit_sha, + "HEAD", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ) + except subprocess.CalledProcessError as exc: + logger.error("Git command failed: %s", exc.stderr) + raise + + files = [ + file_path for file_path in result.stdout.splitlines() if file_path.strip() + ] + + return sorted(set(files)) + + def get_commits_since(self, commit_sha: str) -> list[str]: + try: + result = subprocess.run( + [ + "git", + "-C", + str(self.repository_client.get_local_path()), + "log", + "--format=%H", + f"{commit_sha}..HEAD", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ) + except subprocess.CalledProcessError as exc: + logger.error("Git command failed: %s", exc.stderr) + raise + + commits = [sha for sha in result.stdout.splitlines() if sha.strip()] + + logger.info( + "Detected %s commits since %s", + len(commits), + commit_sha, + ) + + return commits diff --git a/application/utils/harvester/checkpoint_store.py b/application/utils/harvester/checkpoint_store.py new file mode 100644 index 000000000..542d131dd --- /dev/null +++ b/application/utils/harvester/checkpoint_store.py @@ -0,0 +1,54 @@ +import json +from datetime import datetime +from pathlib import Path + +from .models import RepositoryCheckpoint + + +class CheckpointStore: + def __init__(self, checkpoint_file: Path): + self.checkpoint_file = checkpoint_file + + def load(self, repository_id: str) -> RepositoryCheckpoint | None: + if not self.checkpoint_file.exists(): + return None + + data = json.loads( + self.checkpoint_file.read_text( + encoding="utf-8", + ) + ) + + if repository_id not in data: + return None + + checkpoint = data[repository_id] + + return RepositoryCheckpoint( + repository_id=repository_id, + last_processed_commit=checkpoint["last_processed_commit"], + updated_at=datetime.fromisoformat(checkpoint["updated_at"]), + ) + + def save(self, checkpoint: RepositoryCheckpoint) -> None: + data = {} + if self.checkpoint_file.exists(): + data = json.loads(self.checkpoint_file.read_text(encoding="utf-8")) + + data[checkpoint.repository_id] = { + "last_processed_commit": checkpoint.last_processed_commit, + "updated_at": checkpoint.updated_at.isoformat(), + } + + self.checkpoint_file.parent.mkdir( + parents=True, + exist_ok=True, + ) + + self.checkpoint_file.write_text( + json.dumps( + data, + indent=2, + ), + encoding="utf-8", + ) diff --git a/application/utils/harvester/models.py b/application/utils/harvester/models.py new file mode 100644 index 000000000..fe936535a --- /dev/null +++ b/application/utils/harvester/models.py @@ -0,0 +1,16 @@ +from dataclasses import dataclass +from datetime import datetime + + +@dataclass(slots=True) +class RepositoryCheckpoint: + repository_id: str + last_processed_commit: str | None + updated_at: datetime + + +@dataclass(slots=True) +class RepositoryChangeSet: + repository_id: str + commit_sha: str + modified_files: list[str] From eb72939a4dd3f6cc8b2125b3d0055a1cd04c83df Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 18 Jul 2026 18:51:21 +0530 Subject: [PATCH 03/27] fix(harvester): write checkpoints atomically --- application/utils/harvester/checkpoint_store.py | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/application/utils/harvester/checkpoint_store.py b/application/utils/harvester/checkpoint_store.py index 542d131dd..1af13168c 100644 --- a/application/utils/harvester/checkpoint_store.py +++ b/application/utils/harvester/checkpoint_store.py @@ -1,4 +1,5 @@ import json +import os from datetime import datetime from pathlib import Path @@ -27,13 +28,19 @@ def load(self, repository_id: str) -> RepositoryCheckpoint | None: return RepositoryCheckpoint( repository_id=repository_id, last_processed_commit=checkpoint["last_processed_commit"], - updated_at=datetime.fromisoformat(checkpoint["updated_at"]), + updated_at=datetime.fromisoformat( + checkpoint["updated_at"], + ), ) def save(self, checkpoint: RepositoryCheckpoint) -> None: data = {} if self.checkpoint_file.exists(): - data = json.loads(self.checkpoint_file.read_text(encoding="utf-8")) + data = json.loads( + self.checkpoint_file.read_text( + encoding="utf-8", + ) + ) data[checkpoint.repository_id] = { "last_processed_commit": checkpoint.last_processed_commit, @@ -45,10 +52,13 @@ def save(self, checkpoint: RepositoryCheckpoint) -> None: exist_ok=True, ) - self.checkpoint_file.write_text( + temp_file = self.checkpoint_file.with_suffix(".tmp") + temp_file.write_text( json.dumps( data, indent=2, ), encoding="utf-8", ) + + os.replace(temp_file, self.checkpoint_file) From 877bf7a338d2210fc9e235e0beca4c484ad610ba Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Thu, 23 Jul 2026 18:01:46 +0530 Subject: [PATCH 04/27] Use immutable commit ranges for change detection --- .../harvester_test/change_detector_test.py | 136 ++++++++++++++++-- .../git_repository_client_integration_test.py | 66 +++++++++ .../utils/harvester/change_detector.py | 52 +++++-- 3 files changed, 235 insertions(+), 19 deletions(-) diff --git a/application/tests/harvester_test/change_detector_test.py b/application/tests/harvester_test/change_detector_test.py index a74590fd8..29a77cc35 100644 --- a/application/tests/harvester_test/change_detector_test.py +++ b/application/tests/harvester_test/change_detector_test.py @@ -1,5 +1,6 @@ import unittest from unittest.mock import MagicMock +from unittest.mock import call from unittest.mock import patch from application.utils.harvester.change_detector import ( @@ -10,14 +11,21 @@ class ChangeDetectorTests(unittest.TestCase): @patch("application.utils.harvester.change_detector.subprocess.run") def test_get_modified_files_since(self, mock_run): - mock_run.return_value = MagicMock( - stdout="a.md\nb.md\na.md\n", - ) - client = MagicMock() + client.get_local_path.return_value = "/tmp/repo" + + mock_run.side_effect = [ + MagicMock(stdout="resolved_base\n"), + MagicMock(stdout="resolved_target\n"), + MagicMock(stdout="a.md\nb.md\na.md\n"), + ] + detector = ChangeDetector(client) - files = detector.get_modified_files_since("abc123") + files = detector.get_modified_files_since( + "base", + "target", + ) self.assertEqual( files, @@ -27,16 +35,73 @@ def test_get_modified_files_since(self, mock_run): ], ) - @patch("application.utils.harvester.change_detector.subprocess.run") - def test_get_commits_since(self, mock_run): - mock_run.return_value = MagicMock( - stdout="111\n222\n333\n", + mock_run.assert_has_calls( + [ + call( + [ + "git", + "-C", + "/tmp/repo", + "rev-parse", + "--verify", + "--end-of-options", + "base^{commit}", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ), + call( + [ + "git", + "-C", + "/tmp/repo", + "rev-parse", + "--verify", + "--end-of-options", + "target^{commit}", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ), + call( + [ + "git", + "-C", + "/tmp/repo", + "diff", + "--name-only", + "resolved_base", + "resolved_target", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ), + ] ) + @patch("application.utils.harvester.change_detector.subprocess.run") + def test_get_commits_since(self, mock_run): client = MagicMock() + client.get_local_path.return_value = "/tmp/repo" + + mock_run.side_effect = [ + MagicMock(stdout="resolved_base\n"), + MagicMock(stdout="resolved_target\n"), + MagicMock(stdout="111\n222\n333\n"), + ] + detector = ChangeDetector(client) - commits = detector.get_commits_since("abc123") + commits = detector.get_commits_since( + "base", + "target", + ) self.assertEqual( commits, @@ -47,6 +112,57 @@ def test_get_commits_since(self, mock_run): ], ) + self.assertEqual( + mock_run.call_args_list, + [ + call( + [ + "git", + "-C", + "/tmp/repo", + "rev-parse", + "--verify", + "--end-of-options", + "base^{commit}", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ), + call( + [ + "git", + "-C", + "/tmp/repo", + "rev-parse", + "--verify", + "--end-of-options", + "target^{commit}", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ), + call( + [ + "git", + "-C", + "/tmp/repo", + "log", + "--reverse", + "--format=%H", + "resolved_base..resolved_target", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ), + ], + ) + if __name__ == "__main__": unittest.main() diff --git a/application/tests/harvester_test/git_repository_client_integration_test.py b/application/tests/harvester_test/git_repository_client_integration_test.py index 4803e1d67..f2f585098 100644 --- a/application/tests/harvester_test/git_repository_client_integration_test.py +++ b/application/tests/harvester_test/git_repository_client_integration_test.py @@ -7,6 +7,7 @@ from application.utils.harvester.git_repository_client import ( GitRepositoryClient, ) +from application.utils.harvester.change_detector import ChangeDetector class IntegrationGitRepositoryClient(GitRepositoryClient): @@ -201,6 +202,71 @@ def run_sync(client): self.assertEqual((self.cache / "test.txt").read_text(), "v1") + def test_change_detector_uses_captured_target_sha(self): + client = self.create_client() + client.clone() + + detector = ChangeDetector(client) + + base = client.get_current_commit_sha() + + (self.work / "file.txt").write_text("B") + git("add", ".", cwd=self.work) + git("commit", "-m", "second", cwd=self.work) + git("push", "origin", "main", cwd=self.work) + + client.fetch() + + target = client.get_current_commit_sha() + + (self.work / "another.txt").write_text("C") + git("add", ".", cwd=self.work) + git("commit", "-m", "third", cwd=self.work) + git("push", "origin", "main", cwd=self.work) + + client.fetch() + files = detector.get_modified_files_since(base, target) + commits = detector.get_commits_since(base, target) + + self.assertEqual(files, ["file.txt"]) + self.assertEqual(commits, [target]) + + def test_change_detector_returns_commits_oldest_first(self): + client = self.create_client() + client.clone() + + detector = ChangeDetector(client) + base = client.get_current_commit_sha() + + (self.work / "file.txt").write_text("B") + git("add", ".", cwd=self.work) + git("commit", "-m", "B", cwd=self.work) + commit_b = git_output("rev-parse", "HEAD", cwd=self.work) + + (self.work / "file.txt").write_text("C") + git("add", ".", cwd=self.work) + git("commit", "-m", "C", cwd=self.work) + commit_c = git_output("rev-parse", "HEAD", cwd=self.work) + + (self.work / "file.txt").write_text("D") + git("add", ".", cwd=self.work) + git("commit", "-m", "D", cwd=self.work) + commit_d = git_output("rev-parse", "HEAD", cwd=self.work) + + git("push", "origin", "main", cwd=self.work) + + client.fetch() + + commits = detector.get_commits_since(base, commit_d) + self.assertEqual( + commits, + [ + commit_b, + commit_c, + commit_d, + ], + ) + if __name__ == "__main__": unittest.main() diff --git a/application/utils/harvester/change_detector.py b/application/utils/harvester/change_detector.py index 8b0153492..584d43d9c 100644 --- a/application/utils/harvester/change_detector.py +++ b/application/utils/harvester/change_detector.py @@ -10,12 +10,41 @@ class ChangeDetector: def __init__(self, repository_client: GitRepositoryClient): self.repository_client = repository_client - def get_modified_files_since(self, commit_sha: str) -> list[str]: + def _resolve_commit(self, commit_sha: str) -> str: + try: + result = subprocess.run( + [ + "git", + "-C", + str(self.repository_client.get_local_path()), + "rev-parse", + "--verify", + "--end-of-options", + f"{commit_sha}^{{commit}}", + ], + capture_output=True, + text=True, + check=True, + timeout=60, + ) + except subprocess.CalledProcessError as exc: + logger.error("Git command failed: %s", exc.stderr) + raise + + return result.stdout.strip() + + def get_modified_files_since( + self, base_commit: str, target_commit: str + ) -> list[str]: logger.info( - "Detecting changes since commit %s", - commit_sha, + "Detecting changes between %s and %s", + base_commit, + target_commit, ) + base = self._resolve_commit(base_commit) + target = self._resolve_commit(target_commit) + try: result = subprocess.run( [ @@ -24,8 +53,8 @@ def get_modified_files_since(self, commit_sha: str) -> list[str]: str(self.repository_client.get_local_path()), "diff", "--name-only", - commit_sha, - "HEAD", + base, + target, ], capture_output=True, text=True, @@ -42,7 +71,10 @@ def get_modified_files_since(self, commit_sha: str) -> list[str]: return sorted(set(files)) - def get_commits_since(self, commit_sha: str) -> list[str]: + def get_commits_since(self, base_commit: str, target_commit: str) -> list[str]: + base = self._resolve_commit(base_commit) + target = self._resolve_commit(target_commit) + try: result = subprocess.run( [ @@ -50,8 +82,9 @@ def get_commits_since(self, commit_sha: str) -> list[str]: "-C", str(self.repository_client.get_local_path()), "log", + "--reverse", "--format=%H", - f"{commit_sha}..HEAD", + f"{base}..{target}", ], capture_output=True, text=True, @@ -65,9 +98,10 @@ def get_commits_since(self, commit_sha: str) -> list[str]: commits = [sha for sha in result.stdout.splitlines() if sha.strip()] logger.info( - "Detected %s commits since %s", + "Detected %s commits between %s and %s", len(commits), - commit_sha, + base_commit, + target_commit, ) return commits From 967f14636cabd5fb61cddbca8adc3667276dfa94 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Thu, 23 Jul 2026 18:18:48 +0530 Subject: [PATCH 05/27] test: replace mocked /tmp/repo path --- .../tests/harvester_test/change_detector_test.py | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/application/tests/harvester_test/change_detector_test.py b/application/tests/harvester_test/change_detector_test.py index 29a77cc35..5cb33456a 100644 --- a/application/tests/harvester_test/change_detector_test.py +++ b/application/tests/harvester_test/change_detector_test.py @@ -12,7 +12,7 @@ class ChangeDetectorTests(unittest.TestCase): @patch("application.utils.harvester.change_detector.subprocess.run") def test_get_modified_files_since(self, mock_run): client = MagicMock() - client.get_local_path.return_value = "/tmp/repo" + client.get_local_path.return_value = "repo-under-test" mock_run.side_effect = [ MagicMock(stdout="resolved_base\n"), @@ -41,7 +41,7 @@ def test_get_modified_files_since(self, mock_run): [ "git", "-C", - "/tmp/repo", + "repo-under-test", "rev-parse", "--verify", "--end-of-options", @@ -56,7 +56,7 @@ def test_get_modified_files_since(self, mock_run): [ "git", "-C", - "/tmp/repo", + "repo-under-test", "rev-parse", "--verify", "--end-of-options", @@ -71,7 +71,7 @@ def test_get_modified_files_since(self, mock_run): [ "git", "-C", - "/tmp/repo", + "repo-under-test", "diff", "--name-only", "resolved_base", @@ -88,7 +88,7 @@ def test_get_modified_files_since(self, mock_run): @patch("application.utils.harvester.change_detector.subprocess.run") def test_get_commits_since(self, mock_run): client = MagicMock() - client.get_local_path.return_value = "/tmp/repo" + client.get_local_path.return_value = "repo-under-test" mock_run.side_effect = [ MagicMock(stdout="resolved_base\n"), @@ -119,7 +119,7 @@ def test_get_commits_since(self, mock_run): [ "git", "-C", - "/tmp/repo", + "repo-under-test", "rev-parse", "--verify", "--end-of-options", @@ -134,7 +134,7 @@ def test_get_commits_since(self, mock_run): [ "git", "-C", - "/tmp/repo", + "repo-under-test", "rev-parse", "--verify", "--end-of-options", @@ -149,7 +149,7 @@ def test_get_commits_since(self, mock_run): [ "git", "-C", - "/tmp/repo", + "repo-under-test", "log", "--reverse", "--format=%H", From bb90b14fea105ff5babfdc6de85dfd3f7b06a58a Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Thu, 23 Jul 2026 19:33:23 +0530 Subject: [PATCH 06/27] feat(db): add artifact ingestion persistence models --- application/database/db.py | 139 ++++++++++++++++++ application/tests/import_run_test.py | 52 +++++++ ...b3c4d5e_add_artifact_ingest_persistence.py | 82 +++++++++++ 3 files changed, 273 insertions(+) create mode 100644 migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py diff --git a/application/database/db.py b/application/database/db.py index 4ca0dc040..7ecc9feff 100644 --- a/application/database/db.py +++ b/application/database/db.py @@ -257,6 +257,83 @@ class StagedChangeSet(BaseModel): # type: ignore created_at = sqla.Column(sqla.DateTime, nullable=False) +class ArtifactIngestEvent(BaseModel): # type: ignore + """Tracks one harvested artifact persisted per import run.""" + + __tablename__ = "artifact_ingest_event" + id = sqla.Column(sqla.String, primary_key=True, default=generate_uuid) + run_id = sqla.Column( + sqla.String, + sqla.ForeignKey("import_run.id", onupdate="CASCADE", ondelete="CASCADE"), + nullable=False, + ) + artifact_id = sqla.Column(sqla.String, nullable=False) + harvest_mode = sqla.Column(sqla.String, nullable=False) + event_type = sqla.Column(sqla.String, nullable=False) + source_json = sqla.Column(sqla.Text, nullable=False) + locator_json = sqla.Column(sqla.Text, nullable=False) + artifact_json = sqla.Column(sqla.Text, nullable=False) + harvest_json = sqla.Column(sqla.Text, nullable=False) + observed_at = sqla.Column(sqla.DateTime, nullable=False) + created_at = sqla.Column(sqla.DateTime, nullable=False) + + __table_args__ = ( + sqla.UniqueConstraint( + run_id, + artifact_id, + name="uq_artifact_ingest_event_run_artifact", + ), + ) + + +class IngestChunk(BaseModel): # type: ignore + """Tracks every chunk belonging to an artifact ingest event.""" + + __tablename__ = "ingest_chunk" + id = sqla.Column(sqla.String, primary_key=True, default=generate_uuid) + artifact_event_id = sqla.Column( + sqla.String, + sqla.ForeignKey( + "artifact_ingest_event.id", + onupdate="CASCADE", + ondelete="CASCADE", + ), + nullable=False, + ) + chunk_id = sqla.Column(sqla.String, nullable=False) + text = sqla.Column(sqla.Text, nullable=False) + char_count = sqla.Column(sqla.Integer, nullable=False) + span_json = sqla.Column(sqla.Text, nullable=False) + delta_json = sqla.Column(sqla.Text, nullable=True) + created_at = sqla.Column(sqla.DateTime, nullable=False) + + __table_args__ = ( + sqla.UniqueConstraint( + artifact_event_id, + chunk_id, + name="uq_ingest_chunk_artifact_chunk", + ), + ) + + +def _serialize_json_value(value: Any) -> str: + if value is None: + return "null" + if isinstance(value, str): + return value + return flask_json.dumps(value) + + +def _normalize_utc_datetime(value: Any) -> Any: + from datetime import datetime, timezone + + if isinstance(value, datetime): + if value.tzinfo is None: + return value + return value.astimezone(timezone.utc) + return value + + def create_import_run(source: str, version: Optional[str] = None) -> ImportRun: """Create and persist an import run record. Returns the new ImportRun.""" from datetime import datetime, timezone @@ -304,6 +381,68 @@ def get_previous_import_run(source: str, current_run_id: str) -> Optional[Import ) +def create_artifact_ingest_event( + *, + run_id: str, + artifact_id: str, + harvest_mode: str, + event_type: str, + source_json: Any, + locator_json: Any, + artifact_json: Any, + harvest_json: Any, + observed_at: Any, +) -> ArtifactIngestEvent: + from datetime import datetime, timezone + + observed_at = _normalize_utc_datetime(observed_at) + + event = ArtifactIngestEvent( + id=generate_uuid(), + run_id=run_id, + artifact_id=artifact_id, + harvest_mode=harvest_mode, + event_type=event_type, + source_json=_serialize_json_value(source_json), + locator_json=_serialize_json_value(locator_json), + artifact_json=_serialize_json_value(artifact_json), + harvest_json=_serialize_json_value(harvest_json), + observed_at=observed_at, + created_at=_normalize_utc_datetime(datetime.now(timezone.utc)), + ) + sqla.session.add(event) + sqla.session.commit() + return event + + +def create_ingest_chunk( + *, + artifact_event_id: str, + chunk_id: str, + text: str, + char_count: int, + span_json: Any, + delta_json: Optional[Any] = None, +) -> IngestChunk: + from datetime import datetime, timezone + + chunk = IngestChunk( + id=generate_uuid(), + artifact_event_id=artifact_event_id, + chunk_id=chunk_id, + text=text, + char_count=char_count, + span_json=_serialize_json_value(span_json), + delta_json=( + _serialize_json_value(delta_json) if delta_json is not None else None + ), + created_at=_normalize_utc_datetime(datetime.now(timezone.utc)), + ) + sqla.session.add(chunk) + sqla.session.commit() + return chunk + + def persist_standard_snapshot( *, run_id: str, diff --git a/application/tests/import_run_test.py b/application/tests/import_run_test.py index c3e303625..7b03161a6 100644 --- a/application/tests/import_run_test.py +++ b/application/tests/import_run_test.py @@ -1,6 +1,9 @@ """Tests for import run metadata (Step 6).""" +import json import unittest +from datetime import datetime, timezone + from application import create_app, sqla from application.database import db @@ -31,3 +34,52 @@ def test_get_latest_import_run(self) -> None: self.assertIsNotNone(latest) self.assertEqual(latest.id, run2.id) self.assertEqual(latest.version, "2.0") + + def test_create_artifact_ingest_event_and_chunk(self) -> None: + run = db.create_import_run(source="artifact_ingest", version="1.0") + observed_at = datetime.now(timezone.utc) + + event = db.create_artifact_ingest_event( + run_id=run.id, + artifact_id="artifact-1", + harvest_mode="backfill", + event_type="discovered", + source_json={"uri": "https://example.com/source"}, + locator_json={"path": "/tmp/source"}, + artifact_json={"id": "artifact-1"}, + harvest_json={"status": "ok"}, + observed_at=observed_at, + ) + + self.assertIsNotNone(event.id) + self.assertEqual(event.run_id, run.id) + self.assertEqual(event.artifact_id, "artifact-1") + self.assertEqual( + json.loads(event.source_json), {"uri": "https://example.com/source"} + ) + self.assertEqual(json.loads(event.locator_json), {"path": "/tmp/source"}) + self.assertEqual(json.loads(event.artifact_json), {"id": "artifact-1"}) + self.assertEqual(json.loads(event.harvest_json), {"status": "ok"}) + self.assertEqual( + event.observed_at.replace(tzinfo=None), + observed_at.astimezone(timezone.utc).replace(tzinfo=None), + ) + self.assertIsNotNone(event.created_at) + + chunk = db.create_ingest_chunk( + artifact_event_id=event.id, + chunk_id="chunk-1", + text="hello world", + char_count=11, + span_json={"start": 0, "end": 11}, + delta_json={"op": "add"}, + ) + + self.assertIsNotNone(chunk.id) + self.assertEqual(chunk.artifact_event_id, event.id) + self.assertEqual(chunk.chunk_id, "chunk-1") + self.assertEqual(chunk.text, "hello world") + self.assertEqual(chunk.char_count, 11) + self.assertEqual(json.loads(chunk.span_json), {"start": 0, "end": 11}) + self.assertEqual(json.loads(chunk.delta_json), {"op": "add"}) + self.assertIsNotNone(chunk.created_at) diff --git a/migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py b/migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py new file mode 100644 index 000000000..29cdcda5e --- /dev/null +++ b/migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py @@ -0,0 +1,82 @@ +"""add artifact ingest event and chunk tables + +Revision ID: 9f1a2b3c4d5e +Revises: e1f2a3b4c5d6 +Create Date: 2026-07-23 + +""" + +from alembic import op +import sqlalchemy as sa + + +revision = "9f1a2b3c4d5e" +down_revision = "e1f2a3b4c5d6" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "artifact_ingest_event", + sa.Column("id", sa.String(), primary_key=True), + sa.Column("run_id", sa.String(), nullable=False), + sa.Column("artifact_id", sa.String(), nullable=False), + sa.Column("harvest_mode", sa.String(), nullable=False), + sa.Column("event_type", sa.String(), nullable=False), + sa.Column("source_json", sa.Text(), nullable=False), + sa.Column("locator_json", sa.Text(), nullable=False), + sa.Column("artifact_json", sa.Text(), nullable=False), + sa.Column("harvest_json", sa.Text(), nullable=False), + sa.Column("observed_at", sa.DateTime(), nullable=False), + sa.Column("created_at", sa.DateTime(), nullable=False), + sa.ForeignKeyConstraint( + ["run_id"], + ["import_run.id"], + onupdate="CASCADE", + ondelete="CASCADE", + ), + ) + op.create_unique_constraint( + "uq_artifact_ingest_event_run_artifact", + "artifact_ingest_event", + ["run_id", "artifact_id"], + ) + + op.create_table( + "ingest_chunk", + sa.Column("id", sa.String(), primary_key=True), + sa.Column("artifact_event_id", sa.String(), nullable=False), + sa.Column("chunk_id", sa.String(), nullable=False), + sa.Column("text", sa.Text(), nullable=False), + sa.Column("char_count", sa.Integer(), nullable=False), + sa.Column("span_json", sa.Text(), nullable=False), + sa.Column("delta_json", sa.Text(), nullable=True), + sa.Column("created_at", sa.DateTime(), nullable=False), + sa.ForeignKeyConstraint( + ["artifact_event_id"], + ["artifact_ingest_event.id"], + onupdate="CASCADE", + ondelete="CASCADE", + ), + ) + op.create_unique_constraint( + "uq_ingest_chunk_artifact_chunk", + "ingest_chunk", + ["artifact_event_id", "chunk_id"], + ) + + +def downgrade(): + op.drop_constraint( + "uq_ingest_chunk_artifact_chunk", + "ingest_chunk", + type_="unique", + ) + op.drop_table("ingest_chunk") + op.drop_constraint( + "uq_artifact_ingest_event_run_artifact", + "artifact_ingest_event", + type_="unique", + ) + op.drop_table("artifact_ingest_event") From ab615d5d92c0dbfb5e140c7602d0b248967dac35 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Thu, 23 Jul 2026 22:54:20 +0530 Subject: [PATCH 07/27] Address CodeRabbit review feedback --- application/database/db.py | 4 ---- application/utils/harvester/git_repository_client.py | 2 ++ application/utils/harvester/repository_cache.py | 3 +++ application/utils/harvester/repository_lock.py | 2 +- 4 files changed, 6 insertions(+), 5 deletions(-) diff --git a/application/database/db.py b/application/database/db.py index 7ecc9feff..75f983b63 100644 --- a/application/database/db.py +++ b/application/database/db.py @@ -317,10 +317,6 @@ class IngestChunk(BaseModel): # type: ignore def _serialize_json_value(value: Any) -> str: - if value is None: - return "null" - if isinstance(value, str): - return value return flask_json.dumps(value) diff --git a/application/utils/harvester/git_repository_client.py b/application/utils/harvester/git_repository_client.py index a37deb798..468be2665 100644 --- a/application/utils/harvester/git_repository_client.py +++ b/application/utils/harvester/git_repository_client.py @@ -20,6 +20,8 @@ def __init__( branch: str = "main", local_path: Path | None = None, ) -> None: + if branch.startswith("-"): + raise ValueError("Invalid git branch") self.owner = owner self.repository = repository self.branch = branch diff --git a/application/utils/harvester/repository_cache.py b/application/utils/harvester/repository_cache.py index 94b721a60..a8c9b632c 100644 --- a/application/utils/harvester/repository_cache.py +++ b/application/utils/harvester/repository_cache.py @@ -18,6 +18,9 @@ def build_repository_cache_path( if not _VALID_COMPONENT.fullmatch(repository): raise ValueError(f"Invalid repository name: {repository}") + if branch in {".", ".."}: + raise ValueError("Invalid branch name") + encoded_branch = quote(branch, safe="") candidate = CACHE_ROOT / owner.casefold() / repository.casefold() / encoded_branch diff --git a/application/utils/harvester/repository_lock.py b/application/utils/harvester/repository_lock.py index 9890e6e2e..a9779033b 100644 --- a/application/utils/harvester/repository_lock.py +++ b/application/utils/harvester/repository_lock.py @@ -14,7 +14,7 @@ def repository_lock(repository_path: Path): Acquire an exclusive inter-process lock for a repository cache path. """ - lock_path = repository_path.with_suffix(".lock") + lock_path = repository_path.parent / f"{repository_path.name}.lock" lock_path.parent.mkdir(parents=True, exist_ok=True) with lock_path.open("w") as lock_file: From a6929df2f7cae7986b0ee9aa17a3b371554527d6 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 25 Jul 2026 18:50:40 +0530 Subject: [PATCH 08/27] feat(harvester): replace filesystem checkpoint store with postgres --- application/database/db.py | 30 ++ .../harvester_test/checkpoint_store_test.py | 259 ++++++++++++------ .../utils/harvester/checkpoint_store.py | 123 +++++---- application/utils/harvester/models.py | 4 + ...f6a7b8c9_add_harvester_checkpoint_table.py | 41 +++ 5 files changed, 332 insertions(+), 125 deletions(-) create mode 100644 migrations/versions/d4e5f6a7b8c9_add_harvester_checkpoint_table.py diff --git a/application/database/db.py b/application/database/db.py index 75f983b63..1969e9e9b 100644 --- a/application/database/db.py +++ b/application/database/db.py @@ -8,6 +8,7 @@ import time import yaml +from datetime import datetime, timezone from pprint import pprint from collections import Counter, defaultdict @@ -316,6 +317,35 @@ class IngestChunk(BaseModel): # type: ignore ) +class HarvesterCheckpoint(BaseModel): # type: ignore + __tablename__ = "harvester_checkpoint" + repository_id = sqla.Column(sqla.String, primary_key=True) + provider = sqla.Column(sqla.String, nullable=False) + owner = sqla.Column(sqla.String, nullable=False) + repository = sqla.Column(sqla.String, nullable=False) + branch = sqla.Column(sqla.String, nullable=False) + last_processed_commit = sqla.Column(sqla.String, nullable=True) + created_at = sqla.Column( + sqla.DateTime(timezone=True), + nullable=False, + default=lambda: datetime.now(timezone.utc), + ) + updated_at = sqla.Column( + sqla.DateTime(timezone=True), + nullable=False, + default=lambda: datetime.now(timezone.utc), + ) + __table_args__ = ( + sqla.UniqueConstraint( + "provider", + "owner", + "repository", + "branch", + name="uq_harvester_checkpoint_canonical_source", + ), + ) + + def _serialize_json_value(value: Any) -> str: return flask_json.dumps(value) diff --git a/application/tests/harvester_test/checkpoint_store_test.py b/application/tests/harvester_test/checkpoint_store_test.py index 9859b1fed..5d5f1e947 100644 --- a/application/tests/harvester_test/checkpoint_store_test.py +++ b/application/tests/harvester_test/checkpoint_store_test.py @@ -1,91 +1,192 @@ import unittest -from datetime import datetime -from pathlib import Path +from datetime import datetime, timezone -from application.utils.harvester.checkpoint_store import ( - CheckpointStore, -) -from application.utils.harvester.models import ( - RepositoryCheckpoint, -) +from application import create_app, sqla +from application.utils.harvester.checkpoint_store import CheckpointStore +from application.utils.harvester.models import RepositoryCheckpoint class CheckpointStoreTests(unittest.TestCase): - def test_save_and_load_checkpoint(self): - tmp_dir = Path(self._testMethodName) - - try: - store = CheckpointStore( - tmp_dir / "checkpoints.json", - ) - - checkpoint = RepositoryCheckpoint( - repository_id="owasp-asvs", - last_processed_commit="abc123", - updated_at=datetime.now(), - ) - - store.save(checkpoint) - - loaded = store.load("owasp-asvs") - - if loaded is None: - self.fail("Checkpoint should have been loaded") - - self.assertEqual( - loaded.last_processed_commit, - "abc123", - ) + def setUp(self) -> None: + self.app = create_app(mode="test") + self.app_context = self.app.app_context() + self.app_context.push() + sqla.create_all() - finally: - if tmp_dir.exists(): - import shutil + def tearDown(self) -> None: + sqla.session.remove() + sqla.drop_all() + self.app_context.pop() - shutil.rmtree(tmp_dir) - - def test_load_missing_file(self): - tmp_dir = Path(self._testMethodName) - - try: - store = CheckpointStore( - tmp_dir / "missing.json", - ) - - self.assertIsNone( - store.load("repo"), - ) - - finally: - if tmp_dir.exists(): - import shutil - - shutil.rmtree(tmp_dir) + def test_save_and_load_checkpoint(self): + store = CheckpointStore() + checkpoint = RepositoryCheckpoint( + repository_id="owasp-asvs", + last_processed_commit="abc123", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="owasp", + repository="asvs", + branch="main", + ) + + store.save(checkpoint) + loaded = store.load("owasp-asvs") + + self.assertIsNotNone(loaded) + assert loaded is not None + self.assertEqual(loaded.last_processed_commit, "abc123") + self.assertEqual(loaded.provider, "github") + + def test_update_upsert_and_two_repositories_remain_isolated(self): + store = CheckpointStore() + repo_a = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="commit-1", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-a", + branch="main", + ) + repo_b = RepositoryCheckpoint( + repository_id="repo-b", + last_processed_commit="commit-b", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-b", + branch="main", + ) + store.save(repo_a) + store.save(repo_b) + + updated_a = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="commit-2", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-a", + branch="main", + ) + store.save(updated_a) + + loaded_a = store.load("repo-a") + loaded_b = store.load("repo-b") + + self.assertIsNotNone(loaded_a) + self.assertIsNotNone(loaded_b) + assert loaded_a is not None + assert loaded_b is not None + self.assertEqual(loaded_a.last_processed_commit, "commit-2") + self.assertEqual(loaded_b.last_processed_commit, "commit-b") + + def test_duplicate_canonical_source_identity_rejected(self): + store = CheckpointStore() + first = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="commit-1", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="shared", + branch="main", + ) + second = RepositoryCheckpoint( + repository_id="repo-b", + last_processed_commit="commit-2", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="shared", + branch="main", + ) + + store.save(first) + + with self.assertRaises(ValueError): + store.save(second) + + def test_immutable_repository_identity(self): + store = CheckpointStore() + first = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="commit-1", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-a", + branch="main", + ) + store.save(first) + + conflicting = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="commit-2", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-a", + branch="develop", + ) + + with self.assertRaises(ValueError): + store.save(conflicting) + + def test_null_initial_checkpoint(self): + store = CheckpointStore() + checkpoint = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit=None, + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-a", + branch="main", + ) + + store.save(checkpoint) + loaded = store.load("repo-a") + + self.assertIsNotNone(loaded) + assert loaded is not None + self.assertIsNone(loaded.last_processed_commit) + + def test_transaction_rollback_leaves_previous_checkpoint_intact(self): + store = CheckpointStore() + original = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="commit-1", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-a", + branch="main", + ) + store.save(original) + + conflicting = RepositoryCheckpoint( + repository_id="repo-a", + last_processed_commit="commit-2", + updated_at=datetime.now(timezone.utc), + provider="github", + owner="sample", + repository="repo-a", + branch="develop", + ) + + with self.assertRaises(ValueError): + store.save(conflicting) + + loaded = store.load("repo-a") + self.assertIsNotNone(loaded) + assert loaded is not None + self.assertEqual(loaded.last_processed_commit, "commit-1") def test_load_missing_repository(self): - tmp_dir = Path(self._testMethodName) - - try: - store = CheckpointStore( - tmp_dir / "checkpoint.json", - ) - - store.save( - RepositoryCheckpoint( - repository_id="repo-a", - last_processed_commit="abc123", - updated_at=datetime.now(), - ) - ) - - self.assertIsNone( - store.load("repo-b"), - ) - - finally: - if tmp_dir.exists(): - import shutil - - shutil.rmtree(tmp_dir) + store = CheckpointStore() + self.assertIsNone(store.load("repo-b")) if __name__ == "__main__": diff --git a/application/utils/harvester/checkpoint_store.py b/application/utils/harvester/checkpoint_store.py index 1af13168c..e9e3b62f4 100644 --- a/application/utils/harvester/checkpoint_store.py +++ b/application/utils/harvester/checkpoint_store.py @@ -1,64 +1,95 @@ -import json -import os -from datetime import datetime -from pathlib import Path +from typing import Any +from sqlalchemy.exc import IntegrityError + +from application import sqla +from application.database.db import HarvesterCheckpoint from .models import RepositoryCheckpoint class CheckpointStore: - def __init__(self, checkpoint_file: Path): - self.checkpoint_file = checkpoint_file + def __init__(self, session: Any = None) -> None: + self._session = session - def load(self, repository_id: str) -> RepositoryCheckpoint | None: - if not self.checkpoint_file.exists(): - return None + @property + def session(self) -> Any: + return self._session if self._session is not None else sqla.session - data = json.loads( - self.checkpoint_file.read_text( - encoding="utf-8", - ) + def load(self, repository_id: str) -> RepositoryCheckpoint | None: + session = self.session + record = ( + session.query(HarvesterCheckpoint) + .filter_by(repository_id=repository_id) + .first() ) - - if repository_id not in data: + if record is None: return None - - checkpoint = data[repository_id] - return RepositoryCheckpoint( - repository_id=repository_id, - last_processed_commit=checkpoint["last_processed_commit"], - updated_at=datetime.fromisoformat( - checkpoint["updated_at"], - ), + repository_id=record.repository_id, + last_processed_commit=record.last_processed_commit, + updated_at=record.updated_at, + provider=record.provider, + owner=record.owner, + repository=record.repository, + branch=record.branch, ) def save(self, checkpoint: RepositoryCheckpoint) -> None: - data = {} - if self.checkpoint_file.exists(): - data = json.loads( - self.checkpoint_file.read_text( - encoding="utf-8", + session = self.session + existing = ( + session.query(HarvesterCheckpoint) + .filter_by(repository_id=checkpoint.repository_id) + .first() + ) + + if existing is None: + canonical_conflict = ( + session.query(HarvesterCheckpoint) + .filter_by( + provider=checkpoint.provider, + owner=checkpoint.owner, + repository=checkpoint.repository, + branch=checkpoint.branch, ) + .first() ) + if canonical_conflict is not None: + session.rollback() + raise ValueError("duplicate canonical source identity") - data[checkpoint.repository_id] = { - "last_processed_commit": checkpoint.last_processed_commit, - "updated_at": checkpoint.updated_at.isoformat(), - } - - self.checkpoint_file.parent.mkdir( - parents=True, - exist_ok=True, - ) + new_record = HarvesterCheckpoint( + repository_id=checkpoint.repository_id, + provider=checkpoint.provider, + owner=checkpoint.owner, + repository=checkpoint.repository, + branch=checkpoint.branch, + last_processed_commit=checkpoint.last_processed_commit, + updated_at=checkpoint.updated_at, + ) + session.add(new_record) + try: + session.commit() + except IntegrityError: + session.rollback() + raise ValueError("duplicate canonical source identity") + except Exception: + session.rollback() + raise + return - temp_file = self.checkpoint_file.with_suffix(".tmp") - temp_file.write_text( - json.dumps( - data, - indent=2, - ), - encoding="utf-8", - ) + if ( + existing.provider != checkpoint.provider + or existing.owner != checkpoint.owner + or existing.repository != checkpoint.repository + or existing.branch != checkpoint.branch + ): + session.rollback() + raise ValueError("immutable repository identity") - os.replace(temp_file, self.checkpoint_file) + existing.last_processed_commit = checkpoint.last_processed_commit + existing.updated_at = checkpoint.updated_at + try: + session.commit() + except Exception: + session.rollback() + raise diff --git a/application/utils/harvester/models.py b/application/utils/harvester/models.py index fe936535a..2050913c7 100644 --- a/application/utils/harvester/models.py +++ b/application/utils/harvester/models.py @@ -7,6 +7,10 @@ class RepositoryCheckpoint: repository_id: str last_processed_commit: str | None updated_at: datetime + provider: str + owner: str + repository: str + branch: str @dataclass(slots=True) diff --git a/migrations/versions/d4e5f6a7b8c9_add_harvester_checkpoint_table.py b/migrations/versions/d4e5f6a7b8c9_add_harvester_checkpoint_table.py new file mode 100644 index 000000000..11d4d5e49 --- /dev/null +++ b/migrations/versions/d4e5f6a7b8c9_add_harvester_checkpoint_table.py @@ -0,0 +1,41 @@ +"""add harvester_checkpoint table + +Revision ID: d4e5f6a7b8c9 +Revises: 9f1a2b3c4d5e +Create Date: 2026-07-25 + +""" + +from alembic import op +import sqlalchemy as sa + + +revision = "d4e5f6a7b8c9" +down_revision = "9f1a2b3c4d5e" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "harvester_checkpoint", + sa.Column("repository_id", sa.String(), primary_key=True), + sa.Column("provider", sa.String(), nullable=False), + sa.Column("owner", sa.String(), nullable=False), + sa.Column("repository", sa.String(), nullable=False), + sa.Column("branch", sa.String(), nullable=False), + sa.Column("last_processed_commit", sa.String(), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.UniqueConstraint( + "provider", + "owner", + "repository", + "branch", + name="uq_harvester_checkpoint_canonical_source", + ), + ) + + +def downgrade(): + op.drop_table("harvester_checkpoint") From 92f7f0d9d62eb59e25670fa45f6b1a5ba60e5217 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 25 Jul 2026 19:17:14 +0530 Subject: [PATCH 09/27] test(harvester): match checkout invocation --- application/tests/harvester_test/git_repository_client_test.py | 1 + 1 file changed, 1 insertion(+) diff --git a/application/tests/harvester_test/git_repository_client_test.py b/application/tests/harvester_test/git_repository_client_test.py index 774715617..c8c5b6d4c 100644 --- a/application/tests/harvester_test/git_repository_client_test.py +++ b/application/tests/harvester_test/git_repository_client_test.py @@ -114,6 +114,7 @@ def test_checkout_runs_git_command(self, mock_run): "-C", str(client.get_local_path()), "checkout", + "--", "main", ], check=True, From 189d84b89e6edfc80f4f50c4b226094a62983819 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sun, 26 Jul 2026 19:39:05 +0530 Subject: [PATCH 10/27] Address review feedback for Week 3 --- application/tests/harvester_test/git_repository_client_test.py | 1 - application/utils/harvester/git_repository_client.py | 1 - ..._table.py => 6a9d0d62ef41_add_harvester_checkpoint_table.py} | 2 +- .../versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py | 2 +- 4 files changed, 2 insertions(+), 4 deletions(-) rename migrations/versions/{d4e5f6a7b8c9_add_harvester_checkpoint_table.py => 6a9d0d62ef41_add_harvester_checkpoint_table.py} (97%) diff --git a/application/tests/harvester_test/git_repository_client_test.py b/application/tests/harvester_test/git_repository_client_test.py index c8c5b6d4c..774715617 100644 --- a/application/tests/harvester_test/git_repository_client_test.py +++ b/application/tests/harvester_test/git_repository_client_test.py @@ -114,7 +114,6 @@ def test_checkout_runs_git_command(self, mock_run): "-C", str(client.get_local_path()), "checkout", - "--", "main", ], check=True, diff --git a/application/utils/harvester/git_repository_client.py b/application/utils/harvester/git_repository_client.py index 468be2665..bed925d85 100644 --- a/application/utils/harvester/git_repository_client.py +++ b/application/utils/harvester/git_repository_client.py @@ -162,7 +162,6 @@ def checkout(self, reference: str) -> None: "-C", str(self.local_path), "checkout", - "--", reference, ], check=True, diff --git a/migrations/versions/d4e5f6a7b8c9_add_harvester_checkpoint_table.py b/migrations/versions/6a9d0d62ef41_add_harvester_checkpoint_table.py similarity index 97% rename from migrations/versions/d4e5f6a7b8c9_add_harvester_checkpoint_table.py rename to migrations/versions/6a9d0d62ef41_add_harvester_checkpoint_table.py index 11d4d5e49..9674897aa 100644 --- a/migrations/versions/d4e5f6a7b8c9_add_harvester_checkpoint_table.py +++ b/migrations/versions/6a9d0d62ef41_add_harvester_checkpoint_table.py @@ -10,7 +10,7 @@ import sqlalchemy as sa -revision = "d4e5f6a7b8c9" +revision = "6a9d0d62ef41" down_revision = "9f1a2b3c4d5e" branch_labels = None depends_on = None diff --git a/migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py b/migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py index 29cdcda5e..9fad7bed6 100644 --- a/migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py +++ b/migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py @@ -11,7 +11,7 @@ revision = "9f1a2b3c4d5e" -down_revision = "e1f2a3b4c5d6" +down_revision = "c7d8e9f0a1b2" branch_labels = None depends_on = None From 981931342cf52941eead880cd768039d6911c303 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sun, 26 Jul 2026 19:49:51 +0530 Subject: [PATCH 11/27] docs : fix Revision ID in the migration file --- .../versions/6a9d0d62ef41_add_harvester_checkpoint_table.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/migrations/versions/6a9d0d62ef41_add_harvester_checkpoint_table.py b/migrations/versions/6a9d0d62ef41_add_harvester_checkpoint_table.py index 9674897aa..4a3ed2d54 100644 --- a/migrations/versions/6a9d0d62ef41_add_harvester_checkpoint_table.py +++ b/migrations/versions/6a9d0d62ef41_add_harvester_checkpoint_table.py @@ -1,6 +1,6 @@ """add harvester_checkpoint table -Revision ID: d4e5f6a7b8c9 +Revision ID: 6a9d0d62ef41 Revises: 9f1a2b3c4d5e Create Date: 2026-07-25 From 8763f2173ed5299f133aa2f616ce74071e3a83e6 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 18 Jul 2026 18:01:08 +0530 Subject: [PATCH 12/27] fix(harvester): address repository client review feedback --- application/utils/harvester/git_repository_client.py | 1 + 1 file changed, 1 insertion(+) diff --git a/application/utils/harvester/git_repository_client.py b/application/utils/harvester/git_repository_client.py index bed925d85..468be2665 100644 --- a/application/utils/harvester/git_repository_client.py +++ b/application/utils/harvester/git_repository_client.py @@ -162,6 +162,7 @@ def checkout(self, reference: str) -> None: "-C", str(self.local_path), "checkout", + "--", reference, ], check=True, From 1d5a9f0de51750ebf1cf5fa920e1e92bc8ea753e Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Wed, 8 Jul 2026 14:10:50 +0530 Subject: [PATCH 13/27] feat(harvester): implement repository file filtering --- .../tests/harvester_test/file_filter_test.py | 62 +++++++++++++++++++ .../filtering_benchmark_test.py | 49 +++++++++++++++ .../harvester_test/filtering_metrics_test.py | 35 +++++++++++ application/utils/harvester/__init__.py | 10 +++ application/utils/harvester/file_filter.py | 54 ++++++++++++++++ .../utils/harvester/filtering_benchmark.py | 38 ++++++++++++ .../utils/harvester/filtering_metrics.py | 23 +++++++ application/utils/harvester/models.py | 7 +++ 8 files changed, 278 insertions(+) create mode 100644 application/tests/harvester_test/file_filter_test.py create mode 100644 application/tests/harvester_test/filtering_benchmark_test.py create mode 100644 application/tests/harvester_test/filtering_metrics_test.py create mode 100644 application/utils/harvester/file_filter.py create mode 100644 application/utils/harvester/filtering_benchmark.py create mode 100644 application/utils/harvester/filtering_metrics.py diff --git a/application/tests/harvester_test/file_filter_test.py b/application/tests/harvester_test/file_filter_test.py new file mode 100644 index 000000000..ed7958fa9 --- /dev/null +++ b/application/tests/harvester_test/file_filter_test.py @@ -0,0 +1,62 @@ +import unittest + +from application.utils.harvester.file_filter import ( + FileFilter, +) + + +class FileFilterTests(unittest.TestCase): + def test_extension_filtering(self): + file_filter = FileFilter() + + result = file_filter.filter_files( + [ + "README.md", + "image.png", + "script.js", + ] + ) + + self.assertEqual( + result, + ["README.md"], + ) + + def test_regex_filtering(self): + file_filter = FileFilter() + + result = file_filter.filter_files( + [ + ".github/workflows/test.yml", + "docs/setup.md", + ] + ) + + self.assertEqual( + result, + ["docs/setup.md"], + ) + + def test_combined_filtering(self): + file_filter = FileFilter() + + result = file_filter.filter_files( + [ + "README.md", + ".github/workflows/test.yml", + "node_modules/react/index.js", + "docs/setup.md", + ] + ) + + self.assertEqual( + result, + [ + "README.md", + "docs/setup.md", + ], + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/tests/harvester_test/filtering_benchmark_test.py b/application/tests/harvester_test/filtering_benchmark_test.py new file mode 100644 index 000000000..c34105270 --- /dev/null +++ b/application/tests/harvester_test/filtering_benchmark_test.py @@ -0,0 +1,49 @@ +import unittest + +from application.utils.harvester.file_filter import ( + FileFilter, +) +from application.utils.harvester.models import ( + FilteringMetrics, +) + + +class FilteringBenchmarkTests(unittest.TestCase): + def test_filtering_benchmark(self): + files = [ + "README.md", + ".github/workflows/ci.yml", + "docs/guide.md", + "image.png", + "notes.txt", + "package-lock.json", + ] + + file_filter = FileFilter() + + retained_files = file_filter.filter_files(files) + + metrics = FilteringMetrics( + total_files=len(files), + retained_files=len(retained_files), + filtered_files=len(files) - len(retained_files), + ) + + self.assertEqual( + metrics.total_files, + 6, + ) + + self.assertEqual( + metrics.retained_files, + 3, + ) + + self.assertEqual( + metrics.filtered_files, + 3, + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/tests/harvester_test/filtering_metrics_test.py b/application/tests/harvester_test/filtering_metrics_test.py new file mode 100644 index 000000000..252cb0941 --- /dev/null +++ b/application/tests/harvester_test/filtering_metrics_test.py @@ -0,0 +1,35 @@ +import unittest + +from application.utils.harvester.filtering_metrics import ( + FilteringMetricsCollector, +) + + +class FilteringMetricsCollectorTests(unittest.TestCase): + def test_filtering_metrics_collection(self): + collector = FilteringMetricsCollector() + + collector.record_retained() + collector.record_retained() + collector.record_filtered() + + metrics = collector.build() + + self.assertEqual( + metrics.total_files, + 3, + ) + + self.assertEqual( + metrics.retained_files, + 2, + ) + + self.assertEqual( + metrics.filtered_files, + 1, + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/utils/harvester/__init__.py b/application/utils/harvester/__init__.py index 2ac608b3e..eda28b629 100644 --- a/application/utils/harvester/__init__.py +++ b/application/utils/harvester/__init__.py @@ -17,12 +17,22 @@ from .git_repository_client import GitRepositoryClient from .repository_client import RepositoryClient from .repository_cache import build_repository_cache_path +from .file_filter import FileFilter +from .filtering_metrics import FilteringMetricsCollector +from .filtering_benchmark import ( + FilteringBenchmark, + FilteringBenchmarkResult, +) __all__ = [ "build_repository_cache_path", "ChunkingConfig", "ConfigLoaderError", "GitRepositoryClient", + "FileFilter", + "FilteringMetricsCollector", + "FilteringBenchmark", + "FilteringBenchmarkResult", "PathRules", "PollingConfig", "RepositoryClient", diff --git a/application/utils/harvester/file_filter.py b/application/utils/harvester/file_filter.py new file mode 100644 index 000000000..706dffe2c --- /dev/null +++ b/application/utils/harvester/file_filter.py @@ -0,0 +1,54 @@ +import re + +DEFAULT_ALLOWED_EXTENSIONS = { + ".md", + ".mdx", + ".rst", + ".txt", + ".adoc", +} + +DEFAULT_EXCLUDE_PATTERNS = [ + r"^\.github/", + r"^\.git/", + r"^node_modules/", + r"^dist/", + r"^build/", + r"^coverage/", + r"^vendor/", + r".*package-lock\.json$", + r".*yarn\.lock$", + r".*pnpm-lock\.yaml$", +] + + +class FileFilter: + def __init__( + self, + exclude_patterns: list[str] | None = None, + allowed_extensions: set[str] | None = None, + ): + self.exclude_patterns = exclude_patterns or DEFAULT_EXCLUDE_PATTERNS + self.allowed_extensions = allowed_extensions or DEFAULT_ALLOWED_EXTENSIONS + + def is_excluded_by_pattern(self, file_path: str) -> bool: + return any(re.search(pattern, file_path) for pattern in self.exclude_patterns) + + def is_allowed_extension(self, file_path: str) -> bool: + return any( + file_path.endswith(extension) for extension in self.allowed_extensions + ) + + def filter_files(self, files: list[str]) -> list[str]: + filtered_files = [] + + for file_path in files: + if self.is_excluded_by_pattern(file_path): + continue + + if not self.is_allowed_extension(file_path): + continue + + filtered_files.append(file_path) + + return filtered_files diff --git a/application/utils/harvester/filtering_benchmark.py b/application/utils/harvester/filtering_benchmark.py new file mode 100644 index 000000000..2fc095e4e --- /dev/null +++ b/application/utils/harvester/filtering_benchmark.py @@ -0,0 +1,38 @@ +from dataclasses import dataclass + +from .file_filter import FileFilter +from .models import FilteringMetrics + + +@dataclass +class FilteringBenchmarkResult: + total_files: int + retained_files: int + filtered_files: int + retention_rate: float + filtering_rate: float + + +class FilteringBenchmark: + def __init__( + self, + file_filter: FileFilter, + metrics: FilteringMetrics, + ): + self.file_filter = file_filter + self.metrics = metrics + + def run(self, file_paths: list[str]) -> FilteringBenchmarkResult: + retained = self.file_filter.filter_files(file_paths) + + total = len(file_paths) + retained_count = len(retained) + filtered_count = total - retained_count + + return FilteringBenchmarkResult( + total_files=total, + retained_files=retained_count, + filtered_files=filtered_count, + retention_rate=(retained_count / total if total else 0.0), + filtering_rate=(filtered_count / total if total else 0.0), + ) diff --git a/application/utils/harvester/filtering_metrics.py b/application/utils/harvester/filtering_metrics.py new file mode 100644 index 000000000..d496c3e2a --- /dev/null +++ b/application/utils/harvester/filtering_metrics.py @@ -0,0 +1,23 @@ +from .models import FilteringMetrics + + +class FilteringMetricsCollector: + def __init__(self): + self.total_files = 0 + self.retained_files = 0 + self.filtered_files = 0 + + def record_retained(self) -> None: + self.total_files += 1 + self.retained_files += 1 + + def record_filtered(self) -> None: + self.total_files += 1 + self.filtered_files += 1 + + def build(self) -> FilteringMetrics: + return FilteringMetrics( + total_files=self.total_files, + retained_files=self.retained_files, + filtered_files=self.filtered_files, + ) diff --git a/application/utils/harvester/models.py b/application/utils/harvester/models.py index 2050913c7..227c0d64e 100644 --- a/application/utils/harvester/models.py +++ b/application/utils/harvester/models.py @@ -1,5 +1,6 @@ from dataclasses import dataclass from datetime import datetime +from pydantic import BaseModel @dataclass(slots=True) @@ -18,3 +19,9 @@ class RepositoryChangeSet: repository_id: str commit_sha: str modified_files: list[str] + + +class FilteringMetrics(BaseModel): + total_files: int + retained_files: int + filtered_files: int From a13641a43faf589631c394b8621a5784abe79f03 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 18 Jul 2026 21:08:19 +0530 Subject: [PATCH 14/27] fix(harvester): improve filtering benchmark and sync behavior --- .../filtering_benchmark_test.py | 37 +++++-------------- .../git_repository_client_test.py | 18 ++++++++- .../utils/harvester/filtering_benchmark.py | 8 +--- .../utils/harvester/git_repository_client.py | 1 - 4 files changed, 27 insertions(+), 37 deletions(-) diff --git a/application/tests/harvester_test/filtering_benchmark_test.py b/application/tests/harvester_test/filtering_benchmark_test.py index c34105270..a75c435ec 100644 --- a/application/tests/harvester_test/filtering_benchmark_test.py +++ b/application/tests/harvester_test/filtering_benchmark_test.py @@ -1,11 +1,7 @@ import unittest -from application.utils.harvester.file_filter import ( - FileFilter, -) -from application.utils.harvester.models import ( - FilteringMetrics, -) +from application.utils.harvester.file_filter import FileFilter +from application.utils.harvester.filtering_benchmark import FilteringBenchmark class FilteringBenchmarkTests(unittest.TestCase): @@ -19,30 +15,15 @@ def test_filtering_benchmark(self): "package-lock.json", ] - file_filter = FileFilter() + benchmark = FilteringBenchmark(file_filter=FileFilter()) - retained_files = file_filter.filter_files(files) + result = benchmark.run(files) - metrics = FilteringMetrics( - total_files=len(files), - retained_files=len(retained_files), - filtered_files=len(files) - len(retained_files), - ) - - self.assertEqual( - metrics.total_files, - 6, - ) - - self.assertEqual( - metrics.retained_files, - 3, - ) - - self.assertEqual( - metrics.filtered_files, - 3, - ) + self.assertEqual(result.total_files, 6) + self.assertEqual(result.retained_files, 3) + self.assertEqual(result.filtered_files, 3) + self.assertEqual(result.retention_rate, 0.5) + self.assertEqual(result.filtering_rate, 0.5) if __name__ == "__main__": diff --git a/application/tests/harvester_test/git_repository_client_test.py b/application/tests/harvester_test/git_repository_client_test.py index 774715617..1121c973f 100644 --- a/application/tests/harvester_test/git_repository_client_test.py +++ b/application/tests/harvester_test/git_repository_client_test.py @@ -70,7 +70,8 @@ def test_sync_clones_when_repository_missing(self): mock_clone.assert_called_once() - def test_sync_fetches_when_repository_exists(self): + @patch("application.utils.harvester.git_repository_client.subprocess.run") + def test_sync_fetches_when_repository_exists(self, mock_run): client = GitRepositoryClient( owner="OWASP", repository="ASVS", @@ -88,6 +89,21 @@ def test_sync_fetches_when_repository_exists(self): mock_fetch.assert_called_once() + mock_run.assert_called_once_with( + [ + "git", + "-C", + str(client.get_local_path()), + "reset", + "--hard", + "origin/main", + ], + check=True, + capture_output=True, + text=True, + timeout=300, + ) + @patch("application.utils.harvester.git_repository_client.subprocess.run") def test_fetch_runs_git_command(self, mock_run): client = GitRepositoryClient( diff --git a/application/utils/harvester/filtering_benchmark.py b/application/utils/harvester/filtering_benchmark.py index 2fc095e4e..5de2e7b69 100644 --- a/application/utils/harvester/filtering_benchmark.py +++ b/application/utils/harvester/filtering_benchmark.py @@ -1,7 +1,6 @@ from dataclasses import dataclass from .file_filter import FileFilter -from .models import FilteringMetrics @dataclass @@ -14,13 +13,8 @@ class FilteringBenchmarkResult: class FilteringBenchmark: - def __init__( - self, - file_filter: FileFilter, - metrics: FilteringMetrics, - ): + def __init__(self, file_filter: FileFilter): self.file_filter = file_filter - self.metrics = metrics def run(self, file_paths: list[str]) -> FilteringBenchmarkResult: retained = self.file_filter.filter_files(file_paths) diff --git a/application/utils/harvester/git_repository_client.py b/application/utils/harvester/git_repository_client.py index 468be2665..bed925d85 100644 --- a/application/utils/harvester/git_repository_client.py +++ b/application/utils/harvester/git_repository_client.py @@ -162,7 +162,6 @@ def checkout(self, reference: str) -> None: "-C", str(self.local_path), "checkout", - "--", reference, ], check=True, From 81b757e0f351c111acf0d3810b54ae4042a94e57 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Wed, 29 Jul 2026 12:10:39 +0530 Subject: [PATCH 15/27] fix(harvester): isolate filter defaults and preserve empty overrides --- .../tests/harvester_test/file_filter_test.py | 25 +++++++++++++++++++ application/utils/harvester/file_filter.py | 11 ++++++-- 2 files changed, 34 insertions(+), 2 deletions(-) diff --git a/application/tests/harvester_test/file_filter_test.py b/application/tests/harvester_test/file_filter_test.py index ed7958fa9..c6500c621 100644 --- a/application/tests/harvester_test/file_filter_test.py +++ b/application/tests/harvester_test/file_filter_test.py @@ -57,6 +57,31 @@ def test_combined_filtering(self): ], ) + def test_empty_overrides_are_respected(self): + file_filter = FileFilter( + exclude_patterns=[], + allowed_extensions=set(), + ) + + result = file_filter.filter_files( + [ + "README.md", + "image.png", + ] + ) + + self.assertEqual(result, []) + + def test_default_instances_are_isolated(self): + first = FileFilter() + second = FileFilter() + + first.exclude_patterns.append("custom") + first.allowed_extensions.add(".pdf") + + self.assertNotIn("custom", second.exclude_patterns) + self.assertNotIn(".pdf", second.allowed_extensions) + if __name__ == "__main__": unittest.main() diff --git a/application/utils/harvester/file_filter.py b/application/utils/harvester/file_filter.py index 706dffe2c..7bf97b34b 100644 --- a/application/utils/harvester/file_filter.py +++ b/application/utils/harvester/file_filter.py @@ -28,8 +28,15 @@ def __init__( exclude_patterns: list[str] | None = None, allowed_extensions: set[str] | None = None, ): - self.exclude_patterns = exclude_patterns or DEFAULT_EXCLUDE_PATTERNS - self.allowed_extensions = allowed_extensions or DEFAULT_ALLOWED_EXTENSIONS + if exclude_patterns is None: + self.exclude_patterns = list(DEFAULT_EXCLUDE_PATTERNS) + else: + self.exclude_patterns = list(exclude_patterns) + + if allowed_extensions is None: + self.allowed_extensions = set(DEFAULT_ALLOWED_EXTENSIONS) + else: + self.allowed_extensions = set(allowed_extensions) def is_excluded_by_pattern(self, file_path: str) -> bool: return any(re.search(pattern, file_path) for pattern in self.exclude_patterns) From 105ac27d0e93398911a154e1c6606cd0b9987bbe Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Wed, 29 Jul 2026 13:50:49 +0530 Subject: [PATCH 16/27] Use glob-based path filtering for harvester --- .../tests/harvester_test/file_filter_test.py | 46 +++++++++--- .../utils/harvester/exclude_patterns.txt | 3 +- application/utils/harvester/file_filter.py | 70 ++++++++++++------- 3 files changed, 85 insertions(+), 34 deletions(-) diff --git a/application/tests/harvester_test/file_filter_test.py b/application/tests/harvester_test/file_filter_test.py index c6500c621..9594841a2 100644 --- a/application/tests/harvester_test/file_filter_test.py +++ b/application/tests/harvester_test/file_filter_test.py @@ -22,7 +22,7 @@ def test_extension_filtering(self): ["README.md"], ) - def test_regex_filtering(self): + def test_path_exclusion(self): file_filter = FileFilter() result = file_filter.filter_files( @@ -43,8 +43,8 @@ def test_combined_filtering(self): result = file_filter.filter_files( [ "README.md", - ".github/workflows/test.yml", - "node_modules/react/index.js", + ".github/workflows/README.md", + "node_modules/react/README.md", "docs/setup.md", ] ) @@ -72,15 +72,45 @@ def test_empty_overrides_are_respected(self): self.assertEqual(result, []) - def test_default_instances_are_isolated(self): + def test_nested_directory_globs(self): + file_filter = FileFilter() + + result = file_filter.filter_files( + [ + ".github/README.md", + "packages/site/node_modules/README.md", + "docs/archive/old.md", + ".cursor/rules/project.md", + "docs/setup.md", + ] + ) + + self.assertEqual(result, ["docs/setup.md"]) + + def test_explicit_empty_exclusions(self): + file_filter = FileFilter(exclude_patterns=[]) + + result = file_filter.filter_files( + [ + ".github/README.md", + ] + ) + + self.assertEqual( + result, + [".github/README.md"], + ) + + def test_default_instance_isolation(self): first = FileFilter() second = FileFilter() - first.exclude_patterns.append("custom") - first.allowed_extensions.add(".pdf") + first.exclude_patterns.append("**/foo/**") - self.assertNotIn("custom", second.exclude_patterns) - self.assertNotIn(".pdf", second.allowed_extensions) + self.assertNotIn( + "**/foo/**", + second.exclude_patterns, + ) if __name__ == "__main__": diff --git a/application/utils/harvester/exclude_patterns.txt b/application/utils/harvester/exclude_patterns.txt index 499850ae8..2b92e85fc 100644 --- a/application/utils/harvester/exclude_patterns.txt +++ b/application/utils/harvester/exclude_patterns.txt @@ -4,7 +4,8 @@ # to filter non-documentation files during harvesting. -**/.git/* +**/.github/** +**/.git/** **/node_modules/** **/__pycache__/** **/.claude/** diff --git a/application/utils/harvester/file_filter.py b/application/utils/harvester/file_filter.py index 7bf97b34b..086d54b9b 100644 --- a/application/utils/harvester/file_filter.py +++ b/application/utils/harvester/file_filter.py @@ -1,4 +1,6 @@ -import re +from pathlib import PurePosixPath +from pathlib import Path +import pathspec DEFAULT_ALLOWED_EXTENSIONS = { ".md", @@ -8,18 +10,16 @@ ".adoc", } -DEFAULT_EXCLUDE_PATTERNS = [ - r"^\.github/", - r"^\.git/", - r"^node_modules/", - r"^dist/", - r"^build/", - r"^coverage/", - r"^vendor/", - r".*package-lock\.json$", - r".*yarn\.lock$", - r".*pnpm-lock\.yaml$", -] +DEFAULT_EXCLUDE_PATTERNS = tuple( + line.strip() + for line in ( + Path(__file__) + .with_name("exclude_patterns.txt") + .read_text(encoding="utf-8") + .splitlines() + ) + if line.strip() and not line.lstrip().startswith("#") +) class FileFilter: @@ -28,18 +28,38 @@ def __init__( exclude_patterns: list[str] | None = None, allowed_extensions: set[str] | None = None, ): - if exclude_patterns is None: - self.exclude_patterns = list(DEFAULT_EXCLUDE_PATTERNS) - else: - self.exclude_patterns = list(exclude_patterns) + self.exclude_patterns: list[str] = ( + list(DEFAULT_EXCLUDE_PATTERNS) + if exclude_patterns is None + else list(exclude_patterns) + ) + + self.allowed_extensions: set[str] = ( + set(DEFAULT_ALLOWED_EXTENSIONS) + if allowed_extensions is None + else set(allowed_extensions) + ) + + self._validate_patterns() + + try: + self._exclude_spec = pathspec.PathSpec.from_lines( + "gitignore", + self.exclude_patterns, + ) + except Exception as exc: + raise ValueError("Invalid exclude glob") from exc + + def _validate_patterns(self) -> None: + if any(not pattern for pattern in self.exclude_patterns): + raise ValueError("Exclude pattern cannot be empty") - if allowed_extensions is None: - self.allowed_extensions = set(DEFAULT_ALLOWED_EXTENSIONS) - else: - self.allowed_extensions = set(allowed_extensions) + def _normalize_path(self, file_path: str) -> str: + return PurePosixPath(file_path).as_posix() def is_excluded_by_pattern(self, file_path: str) -> bool: - return any(re.search(pattern, file_path) for pattern in self.exclude_patterns) + normalized = self._normalize_path(file_path) + return self._exclude_spec.match_file(normalized) def is_allowed_extension(self, file_path: str) -> bool: return any( @@ -47,7 +67,7 @@ def is_allowed_extension(self, file_path: str) -> bool: ) def filter_files(self, files: list[str]) -> list[str]: - filtered_files = [] + filtered = [] for file_path in files: if self.is_excluded_by_pattern(file_path): @@ -56,6 +76,6 @@ def filter_files(self, files: list[str]) -> list[str]: if not self.is_allowed_extension(file_path): continue - filtered_files.append(file_path) + filtered.append(file_path) - return filtered_files + return filtered From e19394963b25964dbb61d316595f0978be873851 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Wed, 29 Jul 2026 14:10:24 +0530 Subject: [PATCH 17/27] feat(harvester): improve file filtering and exclusions --- .../harvester_test/git_repository_client_test.py | 15 --------------- 1 file changed, 15 deletions(-) diff --git a/application/tests/harvester_test/git_repository_client_test.py b/application/tests/harvester_test/git_repository_client_test.py index 1121c973f..941afd67e 100644 --- a/application/tests/harvester_test/git_repository_client_test.py +++ b/application/tests/harvester_test/git_repository_client_test.py @@ -89,21 +89,6 @@ def test_sync_fetches_when_repository_exists(self, mock_run): mock_fetch.assert_called_once() - mock_run.assert_called_once_with( - [ - "git", - "-C", - str(client.get_local_path()), - "reset", - "--hard", - "origin/main", - ], - check=True, - capture_output=True, - text=True, - timeout=300, - ) - @patch("application.utils.harvester.git_repository_client.subprocess.run") def test_fetch_runs_git_command(self, mock_run): client = GitRepositoryClient( From d5dde16369b0bf4267f97da3f0738e8457ce5620 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Wed, 29 Jul 2026 17:03:36 +0530 Subject: [PATCH 18/27] test: strengthen FileFilter validation --- application/tests/harvester_test/file_filter_test.py | 4 ++++ application/utils/harvester/file_filter.py | 3 +++ 2 files changed, 7 insertions(+) diff --git a/application/tests/harvester_test/file_filter_test.py b/application/tests/harvester_test/file_filter_test.py index 9594841a2..82412df0d 100644 --- a/application/tests/harvester_test/file_filter_test.py +++ b/application/tests/harvester_test/file_filter_test.py @@ -112,6 +112,10 @@ def test_default_instance_isolation(self): second.exclude_patterns, ) + def test_empty_extension_raises(self): + with self.assertRaises(ValueError): + FileFilter(allowed_extensions={""}) + if __name__ == "__main__": unittest.main() diff --git a/application/utils/harvester/file_filter.py b/application/utils/harvester/file_filter.py index 086d54b9b..f1da0ee79 100644 --- a/application/utils/harvester/file_filter.py +++ b/application/utils/harvester/file_filter.py @@ -54,6 +54,9 @@ def _validate_patterns(self) -> None: if any(not pattern for pattern in self.exclude_patterns): raise ValueError("Exclude pattern cannot be empty") + if any(not extension for extension in self.allowed_extensions): + raise ValueError("Allowed extension cannot be empty") + def _normalize_path(self, file_path: str) -> str: return PurePosixPath(file_path).as_posix() From d8cafa9fa3a20d466822a1c68cb26b93a931bb07 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 18 Jul 2026 18:01:08 +0530 Subject: [PATCH 19/27] fix(harvester): address repository client review feedback --- application/utils/harvester/git_repository_client.py | 1 + 1 file changed, 1 insertion(+) diff --git a/application/utils/harvester/git_repository_client.py b/application/utils/harvester/git_repository_client.py index bed925d85..468be2665 100644 --- a/application/utils/harvester/git_repository_client.py +++ b/application/utils/harvester/git_repository_client.py @@ -162,6 +162,7 @@ def checkout(self, reference: str) -> None: "-C", str(self.local_path), "checkout", + "--", reference, ], check=True, From f3c8eb1a0e303dea20ce215cd40213de1e4ec999 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 18 Jul 2026 21:08:19 +0530 Subject: [PATCH 20/27] fix(harvester): improve filtering benchmark and sync behavior --- .../harvester_test/git_repository_client_test.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/application/tests/harvester_test/git_repository_client_test.py b/application/tests/harvester_test/git_repository_client_test.py index 941afd67e..1121c973f 100644 --- a/application/tests/harvester_test/git_repository_client_test.py +++ b/application/tests/harvester_test/git_repository_client_test.py @@ -89,6 +89,21 @@ def test_sync_fetches_when_repository_exists(self, mock_run): mock_fetch.assert_called_once() + mock_run.assert_called_once_with( + [ + "git", + "-C", + str(client.get_local_path()), + "reset", + "--hard", + "origin/main", + ], + check=True, + capture_output=True, + text=True, + timeout=300, + ) + @patch("application.utils.harvester.git_repository_client.subprocess.run") def test_fetch_runs_git_command(self, mock_run): client = GitRepositoryClient( From 3221820f43b8332d8b6bb7894b9e5d2627013e78 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Fri, 10 Jul 2026 13:47:17 +0530 Subject: [PATCH 21/27] feat(harvester): add git diff retrieval pipeline --- .gitignore | 2 + .../harvester_test/diff_retriever_test.py | 49 +++++++++++++++++++ application/utils/harvester/__init__.py | 3 ++ application/utils/harvester/diff_retriever.py | 39 +++++++++++++++ 4 files changed, 93 insertions(+) create mode 100644 application/tests/harvester_test/diff_retriever_test.py create mode 100644 application/utils/harvester/diff_retriever.py diff --git a/.gitignore b/.gitignore index c8a45f86a..4e7c6220c 100644 --- a/.gitignore +++ b/.gitignore @@ -84,3 +84,5 @@ tmp/ cres/* ### Local project management tooling project management scripts/ + +.harvester_cache/ diff --git a/application/tests/harvester_test/diff_retriever_test.py b/application/tests/harvester_test/diff_retriever_test.py new file mode 100644 index 000000000..e716f2452 --- /dev/null +++ b/application/tests/harvester_test/diff_retriever_test.py @@ -0,0 +1,49 @@ +import unittest +from unittest.mock import MagicMock +from unittest.mock import patch + +from application.utils.harvester.diff_retriever import ( + DiffRetriever, +) + + +class DiffRetrieverTests(unittest.TestCase): + @patch("application.utils.harvester.diff_retriever.subprocess.run") + def test_get_diff(self, mock_run): + mock_run.return_value = MagicMock( + stdout="diff --git a/README.md b/README.md\n", + ) + + client = MagicMock() + client.get_local_path.return_value = "/tmp/repo" + + retriever = DiffRetriever(client) + + diff = retriever.get_diff( + "abc123", + "def456", + ) + + self.assertEqual( + diff, + "diff --git a/README.md b/README.md\n", + ) + + mock_run.assert_called_once_with( + [ + "git", + "-C", + "/tmp/repo", + "diff", + "abc123", + "def456", + ], + capture_output=True, + text=True, + check=True, + timeout=300, + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/utils/harvester/__init__.py b/application/utils/harvester/__init__.py index eda28b629..9961aae16 100644 --- a/application/utils/harvester/__init__.py +++ b/application/utils/harvester/__init__.py @@ -19,6 +19,8 @@ from .repository_cache import build_repository_cache_path from .file_filter import FileFilter from .filtering_metrics import FilteringMetricsCollector +from .diff_retriever import DiffRetriever + from .filtering_benchmark import ( FilteringBenchmark, FilteringBenchmarkResult, @@ -28,6 +30,7 @@ "build_repository_cache_path", "ChunkingConfig", "ConfigLoaderError", + "DiffRetriever", "GitRepositoryClient", "FileFilter", "FilteringMetricsCollector", diff --git a/application/utils/harvester/diff_retriever.py b/application/utils/harvester/diff_retriever.py new file mode 100644 index 000000000..6067c9216 --- /dev/null +++ b/application/utils/harvester/diff_retriever.py @@ -0,0 +1,39 @@ +import logging +import subprocess + +from .git_repository_client import GitRepositoryClient + +logger = logging.getLogger(__name__) + + +class DiffRetriever: + def __init__(self, repository_client: GitRepositoryClient) -> None: + self.repository_client = repository_client + + def get_diff(self, base_commit: str, target_commit: str = "HEAD") -> str: + logger.info( + "Retrieving diff between %s and %s", + base_commit, + target_commit, + ) + + try: + result = subprocess.run( + [ + "git", + "-C", + str(self.repository_client.get_local_path()), + "diff", + base_commit, + target_commit, + ], + check=True, + capture_output=True, + text=True, + timeout=300, + ) + except subprocess.CalledProcessError as exc: + logger.error("Failed to retrieve diff: %s", exc.stderr) + raise + + return result.stdout From 2d0489498dea8f5453644d593fc9dfeb1388c0b6 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Fri, 10 Jul 2026 14:00:26 +0530 Subject: [PATCH 22/27] feat(harvester): parse unified git diffs --- .../tests/harvester_test/diff_parser_test.py | 89 +++++++++++++++++++ application/utils/harvester/diff_parser.py | 43 +++++++++ application/utils/harvester/models.py | 6 ++ 3 files changed, 138 insertions(+) create mode 100644 application/tests/harvester_test/diff_parser_test.py create mode 100644 application/utils/harvester/diff_parser.py diff --git a/application/tests/harvester_test/diff_parser_test.py b/application/tests/harvester_test/diff_parser_test.py new file mode 100644 index 000000000..db001df4f --- /dev/null +++ b/application/tests/harvester_test/diff_parser_test.py @@ -0,0 +1,89 @@ +import unittest + +from application.utils.harvester.diff_parser import ( + DiffParser, +) + + +class DiffParserTests(unittest.TestCase): + def test_single_file_diff(self): + parser = DiffParser() + + diff = """diff --git a/test.md b/test.md +--- a/test.md ++++ b/test.md +@@ +-old ++new ++another +""" + + blocks = parser.parse(diff) + + self.assertEqual( + len(blocks), + 1, + ) + + self.assertEqual( + blocks[0].file_path, + "test.md", + ) + + self.assertEqual( + blocks[0].added_lines, + [ + "new", + "another", + ], + ) + + def test_multiple_files(self): + parser = DiffParser() + + diff = """diff --git a/a.md b/a.md +@@ ++one +diff --git a/b.md b/b.md +@@ ++two +""" + + blocks = parser.parse(diff) + + self.assertEqual( + len(blocks), + 2, + ) + + self.assertEqual( + blocks[0].file_path, + "a.md", + ) + + self.assertEqual( + blocks[1].file_path, + "b.md", + ) + + def test_deleted_lines_are_ignored(self): + parser = DiffParser() + + diff = """diff --git a/test.md b/test.md +@@ +-old ++new +""" + + blocks = parser.parse(diff) + + self.assertEqual( + blocks[0].added_lines, + [ + "new", + ], + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/utils/harvester/diff_parser.py b/application/utils/harvester/diff_parser.py new file mode 100644 index 000000000..d9c8e9d32 --- /dev/null +++ b/application/utils/harvester/diff_parser.py @@ -0,0 +1,43 @@ +import re + +from .models import DiffBlock + + +class DiffParser: + def parse(self, diff: str) -> list[DiffBlock]: + blocks: list[DiffBlock] = [] + + current_file: str | None = None + added_lines: list[str] = [] + + for line in diff.splitlines(): + if line.startswith("diff --git"): + if current_file is not None: + blocks.append( + DiffBlock(file_path=current_file, added_lines=added_lines) + ) + + match = re.match(r"diff --git a/(.+?) b/", line) + + if match: + current_file = match.group(1) + added_lines = [] + + continue + + if line.startswith("+++"): + continue + + if line.startswith("---"): + continue + + if line.startswith("@@"): + continue + + if line.startswith("+") and not line.startswith("+++"): + added_lines.append(line[1:]) + + if current_file is not None: + blocks.append(DiffBlock(file_path=current_file, added_lines=added_lines)) + + return blocks diff --git a/application/utils/harvester/models.py b/application/utils/harvester/models.py index 227c0d64e..1a03db0cf 100644 --- a/application/utils/harvester/models.py +++ b/application/utils/harvester/models.py @@ -25,3 +25,9 @@ class FilteringMetrics(BaseModel): total_files: int retained_files: int filtered_files: int + + +@dataclass(slots=True) +class DiffBlock: + file_path: str + added_lines: list[str] From 67658dabb99bbc8c385635b40a911a560e36d68d Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Fri, 10 Jul 2026 14:45:58 +0530 Subject: [PATCH 23/27] feat(harvester): normalize extracted diff content --- .../harvester_test/diff_normalizer_test.py | 105 ++++++++++++++++++ .../utils/harvester/diff_normalizer.py | 33 ++++++ 2 files changed, 138 insertions(+) create mode 100644 application/tests/harvester_test/diff_normalizer_test.py create mode 100644 application/utils/harvester/diff_normalizer.py diff --git a/application/tests/harvester_test/diff_normalizer_test.py b/application/tests/harvester_test/diff_normalizer_test.py new file mode 100644 index 000000000..19d3c8ec6 --- /dev/null +++ b/application/tests/harvester_test/diff_normalizer_test.py @@ -0,0 +1,105 @@ +import unittest + +from application.utils.harvester.diff_normalizer import ( + DiffNormalizer, +) + +from application.utils.harvester.models import ( + DiffBlock, +) + + +class DiffNormalizerTests(unittest.TestCase): + def test_whitespace_normalization(self): + normalizer = DiffNormalizer() + + blocks = [ + DiffBlock( + file_path="README.md", + added_lines=[ + " Hello World ", + "\t\tTabs\t\tEverywhere\t", + "", + " ", + "Unicode\u00a0Space", + "Mix\t of\t tabs and spaces", + " Multiple words together ", + "\u00a0\u00a0Leading unicode spaces\u00a0", + " ## Authentication ", + " - Use MFA ", + " `inline code` ", + " **Important** ", + ], + ) + ] + + result = normalizer.normalize(blocks) + + self.assertEqual( + result[0].added_lines, + [ + "Hello World", + "Tabs Everywhere", + "Unicode Space", + "Mix of tabs and spaces", + "Multiple words together", + "Leading unicode spaces", + "## Authentication", + "- Use MFA", + "`inline code`", + "**Important**", + ], + ) + + def test_remove_empty_lines(self): + normalizer = DiffNormalizer() + + blocks = [ + DiffBlock( + file_path="README.md", + added_lines=[ + "", + " ", + "Hello", + ], + ) + ] + + result = normalizer.normalize(blocks) + + self.assertEqual( + result[0].added_lines, + [ + "Hello", + ], + ) + + def test_multiple_blocks(self): + normalizer = DiffNormalizer() + + blocks = [ + DiffBlock( + file_path="a.md", + added_lines=[" One "], + ), + DiffBlock( + file_path="b.md", + added_lines=[" Two "], + ), + ] + + result = normalizer.normalize(blocks) + + self.assertEqual( + result[0].added_lines, + ["One"], + ) + + self.assertEqual( + result[1].added_lines, + ["Two"], + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/application/utils/harvester/diff_normalizer.py b/application/utils/harvester/diff_normalizer.py new file mode 100644 index 000000000..babdf3508 --- /dev/null +++ b/application/utils/harvester/diff_normalizer.py @@ -0,0 +1,33 @@ +import textacy.preprocessing as prep + +from .models import DiffBlock + + +class DiffNormalizer: + def normalize_line(self, line: str) -> str: + line = prep.normalize.unicode(line) + line = prep.normalize.whitespace(line) + return line.strip() + + def normalize(self, blocks: list[DiffBlock]) -> list[DiffBlock]: + normalized: list[DiffBlock] = [] + + for block in blocks: + cleaned_lines: list[str] = [] + + for line in block.added_lines: + line = self.normalize_line(line) + + if not line: + continue + + cleaned_lines.append(line) + + normalized.append( + DiffBlock( + file_path=block.file_path, + added_lines=cleaned_lines, + ) + ) + + return normalized From eaa6ae93eaa2b42d8ed7e44a6ce804eef964d93e Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Fri, 10 Jul 2026 17:14:07 +0530 Subject: [PATCH 24/27] Enhance diff pipeline with metadata and normalization --- .../harvester_test/diff_normalizer_test.py | 12 +++++ .../tests/harvester_test/diff_parser_test.py | 52 ++++++++++++------- .../harvester_test/diff_pipeline_test.py | 50 ++++++++++++++++++ .../harvester_test/diff_retriever_test.py | 14 +++++ .../utils/harvester/diff_normalizer.py | 14 +++++ application/utils/harvester/diff_parser.py | 33 ++++++++++-- application/utils/harvester/diff_retriever.py | 39 +++++++++++++- application/utils/harvester/models.py | 8 +++ 8 files changed, 199 insertions(+), 23 deletions(-) create mode 100644 application/tests/harvester_test/diff_pipeline_test.py diff --git a/application/tests/harvester_test/diff_normalizer_test.py b/application/tests/harvester_test/diff_normalizer_test.py index 19d3c8ec6..04eb1ce3f 100644 --- a/application/tests/harvester_test/diff_normalizer_test.py +++ b/application/tests/harvester_test/diff_normalizer_test.py @@ -1,4 +1,5 @@ import unittest +from datetime import datetime from application.utils.harvester.diff_normalizer import ( DiffNormalizer, @@ -9,6 +10,13 @@ ) +DIFF_METADATA = { + "repository": "OWASP/ASVS", + "commit_sha": "abc123", + "committed_at": datetime(2026, 1, 1), +} + + class DiffNormalizerTests(unittest.TestCase): def test_whitespace_normalization(self): normalizer = DiffNormalizer() @@ -30,6 +38,7 @@ def test_whitespace_normalization(self): " `inline code` ", " **Important** ", ], + **DIFF_METADATA, ) ] @@ -62,6 +71,7 @@ def test_remove_empty_lines(self): " ", "Hello", ], + **DIFF_METADATA, ) ] @@ -81,10 +91,12 @@ def test_multiple_blocks(self): DiffBlock( file_path="a.md", added_lines=[" One "], + **DIFF_METADATA, ), DiffBlock( file_path="b.md", added_lines=[" Two "], + **DIFF_METADATA, ), ] diff --git a/application/tests/harvester_test/diff_parser_test.py b/application/tests/harvester_test/diff_parser_test.py index db001df4f..a4444e8d9 100644 --- a/application/tests/harvester_test/diff_parser_test.py +++ b/application/tests/harvester_test/diff_parser_test.py @@ -1,9 +1,14 @@ +from datetime import UTC, datetime import unittest from application.utils.harvester.diff_parser import ( DiffParser, ) +TEST_REPOSITORY = "OWASP/ASVS" +TEST_COMMIT_SHA = "abc123" +TEST_COMMITTED_AT = datetime.now(UTC) + class DiffParserTests(unittest.TestCase): def test_single_file_diff(self): @@ -18,13 +23,15 @@ def test_single_file_diff(self): +another """ - blocks = parser.parse(diff) - - self.assertEqual( - len(blocks), - 1, + blocks = parser.parse( + diff, + repository=TEST_REPOSITORY, + commit_sha=TEST_COMMIT_SHA, + committed_at=TEST_COMMITTED_AT, ) + self.assertEqual(len(blocks), 1) + self.assertEqual( blocks[0].file_path, "test.md", @@ -38,6 +45,10 @@ def test_single_file_diff(self): ], ) + self.assertEqual(blocks[0].repository, TEST_REPOSITORY) + self.assertEqual(blocks[0].commit_sha, TEST_COMMIT_SHA) + self.assertEqual(blocks[0].committed_at, TEST_COMMITTED_AT) + def test_multiple_files(self): parser = DiffParser() @@ -49,22 +60,20 @@ def test_multiple_files(self): +two """ - blocks = parser.parse(diff) - - self.assertEqual( - len(blocks), - 2, + blocks = parser.parse( + diff, + repository=TEST_REPOSITORY, + commit_sha=TEST_COMMIT_SHA, + committed_at=TEST_COMMITTED_AT, ) - self.assertEqual( - blocks[0].file_path, - "a.md", - ) + self.assertEqual(len(blocks), 2) - self.assertEqual( - blocks[1].file_path, - "b.md", - ) + self.assertEqual(blocks[0].file_path, "a.md") + self.assertEqual(blocks[1].file_path, "b.md") + + self.assertEqual(blocks[0].repository, TEST_REPOSITORY) + self.assertEqual(blocks[1].repository, TEST_REPOSITORY) def test_deleted_lines_are_ignored(self): parser = DiffParser() @@ -75,7 +84,12 @@ def test_deleted_lines_are_ignored(self): +new """ - blocks = parser.parse(diff) + blocks = parser.parse( + diff, + repository=TEST_REPOSITORY, + commit_sha=TEST_COMMIT_SHA, + committed_at=TEST_COMMITTED_AT, + ) self.assertEqual( blocks[0].added_lines, diff --git a/application/tests/harvester_test/diff_pipeline_test.py b/application/tests/harvester_test/diff_pipeline_test.py new file mode 100644 index 000000000..f27624170 --- /dev/null +++ b/application/tests/harvester_test/diff_pipeline_test.py @@ -0,0 +1,50 @@ +from datetime import UTC, datetime +import time +import unittest + +from application.utils.harvester.diff_normalizer import DiffNormalizer +from application.utils.harvester.diff_parser import DiffParser +from application.utils.harvester.diff_retriever import DiffRetriever +from application.utils.harvester.git_repository_client import GitRepositoryClient + + +class DiffPipelineBenchmark(unittest.TestCase): + """ + Simple benchmark to ensure the complete diff pipeline remains fast. + + This is not intended as a strict performance benchmark, only as a + regression guard against accidental slowdowns. + """ + + def test_pipeline_benchmark(self): + client = GitRepositoryClient( + "OWASP", + "ASVS", + "master", + ) + + retriever = DiffRetriever(client) + parser = DiffParser() + normalizer = DiffNormalizer() + + start = time.perf_counter() + + diff = retriever.get_diff( + "a79c0184", + "122d9e0969465a6041e16c806a0464b35deea444", + ) + + blocks = parser.parse( + diff, + repository="OWASP/ASVS", + commit_sha="122d9e0969465a6041e16c806a0464b35deea444", + committed_at=datetime.now(UTC), + ) + + normalizer.normalize(blocks) + + elapsed = time.perf_counter() - start + + print(f"\nPipeline took {elapsed:.3f}s") + + self.assertLess(elapsed, 5) diff --git a/application/tests/harvester_test/diff_retriever_test.py b/application/tests/harvester_test/diff_retriever_test.py index e716f2452..0e890a07f 100644 --- a/application/tests/harvester_test/diff_retriever_test.py +++ b/application/tests/harvester_test/diff_retriever_test.py @@ -44,6 +44,20 @@ def test_get_diff(self, mock_run): timeout=300, ) + @patch("application.utils.harvester.diff_retriever.subprocess.run") + def test_large_diff_raises(self, mock_run): + mock_run.return_value = MagicMock( + stdout="A" * (51 * 1024 * 1024), + ) + + client = MagicMock() + client.get_local_path.return_value = "/tmp/repo" + + retriever = DiffRetriever(client) + + with self.assertRaises(ValueError): + retriever.get_diff("a", "b") + if __name__ == "__main__": unittest.main() diff --git a/application/utils/harvester/diff_normalizer.py b/application/utils/harvester/diff_normalizer.py index babdf3508..fe8773349 100644 --- a/application/utils/harvester/diff_normalizer.py +++ b/application/utils/harvester/diff_normalizer.py @@ -1,15 +1,26 @@ import textacy.preprocessing as prep +from application.utils.harvester import repository_client from .models import DiffBlock class DiffNormalizer: + """ + Normalizes extracted diff content. + + Whitespace is collapsed, Unicode normalized, + and empty lines removed. + """ + def normalize_line(self, line: str) -> str: line = prep.normalize.unicode(line) line = prep.normalize.whitespace(line) return line.strip() def normalize(self, blocks: list[DiffBlock]) -> list[DiffBlock]: + """ + Normalize every added line in each DiffBlock. + """ normalized: list[DiffBlock] = [] for block in blocks: @@ -27,6 +38,9 @@ def normalize(self, blocks: list[DiffBlock]) -> list[DiffBlock]: DiffBlock( file_path=block.file_path, added_lines=cleaned_lines, + repository=block.repository, + commit_sha=block.commit_sha, + committed_at=block.committed_at, ) ) diff --git a/application/utils/harvester/diff_parser.py b/application/utils/harvester/diff_parser.py index d9c8e9d32..8bdaeb42e 100644 --- a/application/utils/harvester/diff_parser.py +++ b/application/utils/harvester/diff_parser.py @@ -1,10 +1,23 @@ +from datetime import datetime import re from .models import DiffBlock class DiffParser: - def parse(self, diff: str) -> list[DiffBlock]: + """ + Parses unified git diffs into DiffBlock objects. + + Only added lines are extracted. + Deleted lines and diff metadata are ignored. + """ + + def parse( + self, diff: str, repository: str, commit_sha: str, committed_at: datetime + ) -> list[DiffBlock]: + """ + Convert a unified git diff into DiffBlock objects. + """ blocks: list[DiffBlock] = [] current_file: str | None = None @@ -14,7 +27,13 @@ def parse(self, diff: str) -> list[DiffBlock]: if line.startswith("diff --git"): if current_file is not None: blocks.append( - DiffBlock(file_path=current_file, added_lines=added_lines) + DiffBlock( + file_path=current_file, + added_lines=added_lines, + repository=repository, + commit_sha=commit_sha, + committed_at=committed_at, + ) ) match = re.match(r"diff --git a/(.+?) b/", line) @@ -38,6 +57,14 @@ def parse(self, diff: str) -> list[DiffBlock]: added_lines.append(line[1:]) if current_file is not None: - blocks.append(DiffBlock(file_path=current_file, added_lines=added_lines)) + blocks.append( + DiffBlock( + file_path=current_file, + added_lines=added_lines, + repository=repository, + commit_sha=commit_sha, + committed_at=committed_at, + ) + ) return blocks diff --git a/application/utils/harvester/diff_retriever.py b/application/utils/harvester/diff_retriever.py index 6067c9216..78c03bb13 100644 --- a/application/utils/harvester/diff_retriever.py +++ b/application/utils/harvester/diff_retriever.py @@ -7,10 +7,37 @@ class DiffRetriever: + MAX_DIFF_SIZE_BYTES = 50 * 1024 * 1024 + """ + + Retrieves unified git diffs between two commits. + + This class is responsible only for retrieving raw diff text. + + Parsing and normalization are handled by downstream components. + + """ + def __init__(self, repository_client: GitRepositoryClient) -> None: self.repository_client = repository_client def get_diff(self, base_commit: str, target_commit: str = "HEAD") -> str: + """ + Return the unified git diff between two commits. + + Args: + base_commit: + Base commit SHA. + target_commit: + Target commit SHA or branch. + + Raises: + subprocess.CalledProcessError: + If git diff fails. + + ValueError: + If the diff exceeds the configured size limit. + """ logger.info( "Retrieving diff between %s and %s", base_commit, @@ -36,4 +63,14 @@ def get_diff(self, base_commit: str, target_commit: str = "HEAD") -> str: logger.error("Failed to retrieve diff: %s", exc.stderr) raise - return result.stdout + diff = result.stdout + + diff_size = len(diff.encode("utf-8")) + + if diff_size > self.MAX_DIFF_SIZE_BYTES: + raise ValueError( + f"Diff size ({diff_size} bytes) exceeds " + f"maximum supported size ({self.MAX_DIFF_SIZE_BYTES} bytes)." + ) + + return diff diff --git a/application/utils/harvester/models.py b/application/utils/harvester/models.py index 1a03db0cf..0eca718c9 100644 --- a/application/utils/harvester/models.py +++ b/application/utils/harvester/models.py @@ -29,5 +29,13 @@ class FilteringMetrics(BaseModel): @dataclass(slots=True) class DiffBlock: + """ + Intermediate representation of normalized additions + extracted from a repository diff. + """ + file_path: str added_lines: list[str] + repository: str + commit_sha: str + committed_at: datetime | None = None From 79080d737334ddd802497b562d5629d3754c02a5 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Sat, 18 Jul 2026 23:51:05 +0530 Subject: [PATCH 25/27] fix(harvester): address review feedback --- .../harvester_test/diff_pipeline_test.py | 24 ++++++++++++++++--- .../harvester_test/diff_retriever_test.py | 5 ++-- .../git_repository_client_test.py | 16 +------------ .../utils/harvester/checkpoint_store.py | 3 +++ application/utils/harvester/diff_parser.py | 6 ++--- application/utils/harvester/diff_retriever.py | 12 ++++++---- 6 files changed, 37 insertions(+), 29 deletions(-) diff --git a/application/tests/harvester_test/diff_pipeline_test.py b/application/tests/harvester_test/diff_pipeline_test.py index f27624170..5ce0a6fc2 100644 --- a/application/tests/harvester_test/diff_pipeline_test.py +++ b/application/tests/harvester_test/diff_pipeline_test.py @@ -1,4 +1,5 @@ from datetime import UTC, datetime +import subprocess import time import unittest @@ -22,6 +23,23 @@ def test_pipeline_benchmark(self): "ASVS", "master", ) + client.sync() + + head_commit = client.get_current_commit_sha() + + previous_commit = subprocess.run( + [ + "git", + "-C", + str(client.get_local_path()), + "rev-parse", + "HEAD~1", + ], + check=True, + capture_output=True, + text=True, + timeout=300, + ).stdout.strip() retriever = DiffRetriever(client) parser = DiffParser() @@ -30,14 +48,14 @@ def test_pipeline_benchmark(self): start = time.perf_counter() diff = retriever.get_diff( - "a79c0184", - "122d9e0969465a6041e16c806a0464b35deea444", + previous_commit, + head_commit, ) blocks = parser.parse( diff, repository="OWASP/ASVS", - commit_sha="122d9e0969465a6041e16c806a0464b35deea444", + commit_sha=head_commit, committed_at=datetime.now(UTC), ) diff --git a/application/tests/harvester_test/diff_retriever_test.py b/application/tests/harvester_test/diff_retriever_test.py index 0e890a07f..416499e5f 100644 --- a/application/tests/harvester_test/diff_retriever_test.py +++ b/application/tests/harvester_test/diff_retriever_test.py @@ -11,7 +11,7 @@ class DiffRetrieverTests(unittest.TestCase): @patch("application.utils.harvester.diff_retriever.subprocess.run") def test_get_diff(self, mock_run): mock_run.return_value = MagicMock( - stdout="diff --git a/README.md b/README.md\n", + stdout=b"diff --git a/README.md b/README.md\n", ) client = MagicMock() @@ -39,7 +39,6 @@ def test_get_diff(self, mock_run): "def456", ], capture_output=True, - text=True, check=True, timeout=300, ) @@ -47,7 +46,7 @@ def test_get_diff(self, mock_run): @patch("application.utils.harvester.diff_retriever.subprocess.run") def test_large_diff_raises(self, mock_run): mock_run.return_value = MagicMock( - stdout="A" * (51 * 1024 * 1024), + stdout=b"A" * (51 * 1024 * 1024), ) client = MagicMock() diff --git a/application/tests/harvester_test/git_repository_client_test.py b/application/tests/harvester_test/git_repository_client_test.py index 1121c973f..8bbff6ab8 100644 --- a/application/tests/harvester_test/git_repository_client_test.py +++ b/application/tests/harvester_test/git_repository_client_test.py @@ -89,21 +89,6 @@ def test_sync_fetches_when_repository_exists(self, mock_run): mock_fetch.assert_called_once() - mock_run.assert_called_once_with( - [ - "git", - "-C", - str(client.get_local_path()), - "reset", - "--hard", - "origin/main", - ], - check=True, - capture_output=True, - text=True, - timeout=300, - ) - @patch("application.utils.harvester.git_repository_client.subprocess.run") def test_fetch_runs_git_command(self, mock_run): client = GitRepositoryClient( @@ -130,6 +115,7 @@ def test_checkout_runs_git_command(self, mock_run): "-C", str(client.get_local_path()), "checkout", + "--", "main", ], check=True, diff --git a/application/utils/harvester/checkpoint_store.py b/application/utils/harvester/checkpoint_store.py index e9e3b62f4..ed757bb56 100644 --- a/application/utils/harvester/checkpoint_store.py +++ b/application/utils/harvester/checkpoint_store.py @@ -22,8 +22,10 @@ def load(self, repository_id: str) -> RepositoryCheckpoint | None: .filter_by(repository_id=repository_id) .first() ) + if record is None: return None + return RepositoryCheckpoint( repository_id=record.repository_id, last_processed_commit=record.last_processed_commit, @@ -53,6 +55,7 @@ def save(self, checkpoint: RepositoryCheckpoint) -> None: ) .first() ) + if canonical_conflict is not None: session.rollback() raise ValueError("duplicate canonical source identity") diff --git a/application/utils/harvester/diff_parser.py b/application/utils/harvester/diff_parser.py index 8bdaeb42e..d0f124bf9 100644 --- a/application/utils/harvester/diff_parser.py +++ b/application/utils/harvester/diff_parser.py @@ -44,16 +44,16 @@ def parse( continue - if line.startswith("+++"): + if line.startswith("+++ b/") or line.startswith("++/dev/null"): continue - if line.startswith("---"): + if line.startswith("--- a/") or line.startswith("--- /dev/null"): continue if line.startswith("@@"): continue - if line.startswith("+") and not line.startswith("+++"): + if line.startswith("+"): added_lines.append(line[1:]) if current_file is not None: diff --git a/application/utils/harvester/diff_retriever.py b/application/utils/harvester/diff_retriever.py index 78c03bb13..fce640d52 100644 --- a/application/utils/harvester/diff_retriever.py +++ b/application/utils/harvester/diff_retriever.py @@ -56,16 +56,18 @@ def get_diff(self, base_commit: str, target_commit: str = "HEAD") -> str: ], check=True, capture_output=True, - text=True, timeout=300, ) except subprocess.CalledProcessError as exc: - logger.error("Failed to retrieve diff: %s", exc.stderr) + logger.error( + "Failed to retrieve diff: %s", + exc.stderr.decode("utf-8", errors="replace"), + ) raise - diff = result.stdout + diff_bytes = result.stdout - diff_size = len(diff.encode("utf-8")) + diff_size = len(diff_bytes) if diff_size > self.MAX_DIFF_SIZE_BYTES: raise ValueError( @@ -73,4 +75,4 @@ def get_diff(self, base_commit: str, target_commit: str = "HEAD") -> str: f"maximum supported size ({self.MAX_DIFF_SIZE_BYTES} bytes)." ) - return diff + return diff_bytes.decode("utf-8", errors="replace") From 4e6b8a5334a8685b03185441296e01619988d0fb Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Thu, 30 Jul 2026 18:10:49 +0530 Subject: [PATCH 26/27] Harden diff retrieval and isolate network benchmark tests --- .../harvester_test/diff_pipeline_test.py | 5 ++ .../harvester_test/diff_retriever_test.py | 65 +++++++++++++++---- application/utils/harvester/diff_parser.py | 2 +- application/utils/harvester/diff_retriever.py | 21 ++++++ requirements-dev.txt | 1 + 5 files changed, 79 insertions(+), 15 deletions(-) diff --git a/application/tests/harvester_test/diff_pipeline_test.py b/application/tests/harvester_test/diff_pipeline_test.py index 5ce0a6fc2..07160f170 100644 --- a/application/tests/harvester_test/diff_pipeline_test.py +++ b/application/tests/harvester_test/diff_pipeline_test.py @@ -2,6 +2,7 @@ import subprocess import time import unittest +import os from application.utils.harvester.diff_normalizer import DiffNormalizer from application.utils.harvester.diff_parser import DiffParser @@ -18,6 +19,10 @@ class DiffPipelineBenchmark(unittest.TestCase): """ def test_pipeline_benchmark(self): + + if os.getenv("OPENCRE_RUN_NETWORK_TESTS") != "1": + self.skipTest("Network benchmark disabled") + client = GitRepositoryClient( "OWASP", "ASVS", diff --git a/application/tests/harvester_test/diff_retriever_test.py b/application/tests/harvester_test/diff_retriever_test.py index 416499e5f..502ac2d39 100644 --- a/application/tests/harvester_test/diff_retriever_test.py +++ b/application/tests/harvester_test/diff_retriever_test.py @@ -1,6 +1,7 @@ import unittest from unittest.mock import MagicMock from unittest.mock import patch +from unittest.mock import call from application.utils.harvester.diff_retriever import ( DiffRetriever, @@ -10,9 +11,11 @@ class DiffRetrieverTests(unittest.TestCase): @patch("application.utils.harvester.diff_retriever.subprocess.run") def test_get_diff(self, mock_run): - mock_run.return_value = MagicMock( - stdout=b"diff --git a/README.md b/README.md\n", - ) + mock_run.side_effect = [ + MagicMock(stdout="abc123\n"), + MagicMock(stdout="def456\n"), + MagicMock(stdout=b"diff --git a/README.md b/README.md\n"), + ] client = MagicMock() client.get_local_path.return_value = "/tmp/repo" @@ -29,18 +32,52 @@ def test_get_diff(self, mock_run): "diff --git a/README.md b/README.md\n", ) - mock_run.assert_called_once_with( + mock_run.assert_has_calls( [ - "git", - "-C", - "/tmp/repo", - "diff", - "abc123", - "def456", - ], - capture_output=True, - check=True, - timeout=300, + call( + [ + "git", + "-C", + "/tmp/repo", + "rev-parse", + "--verify", + "--end-of-options", + "abc123^{commit}", + ], + check=True, + capture_output=True, + text=True, + timeout=60, + ), + call( + [ + "git", + "-C", + "/tmp/repo", + "rev-parse", + "--verify", + "--end-of-options", + "def456^{commit}", + ], + check=True, + capture_output=True, + text=True, + timeout=60, + ), + call( + [ + "git", + "-C", + "/tmp/repo", + "diff", + "abc123", + "def456", + ], + check=True, + capture_output=True, + timeout=300, + ), + ] ) @patch("application.utils.harvester.diff_retriever.subprocess.run") diff --git a/application/utils/harvester/diff_parser.py b/application/utils/harvester/diff_parser.py index d0f124bf9..e8d98b9d3 100644 --- a/application/utils/harvester/diff_parser.py +++ b/application/utils/harvester/diff_parser.py @@ -44,7 +44,7 @@ def parse( continue - if line.startswith("+++ b/") or line.startswith("++/dev/null"): + if line.startswith("+++ b/") or line.startswith("+++ /dev/null"): continue if line.startswith("--- a/") or line.startswith("--- /dev/null"): diff --git a/application/utils/harvester/diff_retriever.py b/application/utils/harvester/diff_retriever.py index fce640d52..f40f1cbdf 100644 --- a/application/utils/harvester/diff_retriever.py +++ b/application/utils/harvester/diff_retriever.py @@ -44,6 +44,9 @@ def get_diff(self, base_commit: str, target_commit: str = "HEAD") -> str: target_commit, ) + base_commit = self._resolve_commit(base_commit) + target_commit = self._resolve_commit(target_commit) + try: result = subprocess.run( [ @@ -76,3 +79,21 @@ def get_diff(self, base_commit: str, target_commit: str = "HEAD") -> str: ) return diff_bytes.decode("utf-8", errors="replace") + + def _resolve_commit(self, commit: str) -> str: + result = subprocess.run( + [ + "git", + "-C", + str(self.repository_client.get_local_path()), + "rev-parse", + "--verify", + "--end-of-options", + f"{commit}^{{commit}}", + ], + check=True, + capture_output=True, + text=True, + timeout=60, + ) + return result.stdout.strip() diff --git a/requirements-dev.txt b/requirements-dev.txt index 4b8784bdc..283942164 100644 --- a/requirements-dev.txt +++ b/requirements-dev.txt @@ -78,6 +78,7 @@ types-PyYAML typing-inspect pycodestyle pyflakes +textacy # lint / test / typecheck black==24.4.2 From a1b2139ed015eeea3f6c0c152b9ac51c44744598 Mon Sep 17 00:00:00 2001 From: ParthAggarwal16 Date: Thu, 30 Jul 2026 19:01:24 +0530 Subject: [PATCH 27/27] Addressing code rabbit comment --- application/utils/harvester/diff_parser.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/application/utils/harvester/diff_parser.py b/application/utils/harvester/diff_parser.py index e8d98b9d3..12b14784c 100644 --- a/application/utils/harvester/diff_parser.py +++ b/application/utils/harvester/diff_parser.py @@ -38,9 +38,8 @@ def parse( match = re.match(r"diff --git a/(.+?) b/", line) - if match: - current_file = match.group(1) - added_lines = [] + current_file = match.group(1) if match else None + added_lines = [] continue