Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 11 additions & 1 deletion jobs/autoindexer/data_source_handler/git_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ def __init__(self, index_name: str, config: dict[str, Any], rag_client: KAITORAG
self.exclude_matcher = None
self.include_matcher = None
self.last_indexed_commit = self.config.get("lastIndexedCommit", "")
self.conditions = self.config.get("conditions", [])

factory = get_factory(MatcherImplementation.PURE_PYTHON)
if self.exclude_paths:
Expand Down Expand Up @@ -112,13 +113,22 @@ def update_index(self) -> list[str]:

# Clone or fetch repository
self._setup_repository()

# Check for previous error conditions in AutoIndexer status
last_indexing_had_errors = any(
condition.get("type") == "AutoIndexerError" and condition.get("status") == "True"
for condition in self.conditions
)

if last_indexing_had_errors:
logger.warning("Previous indexing had errors, performing full indexing to recover")

# Determine indexing strategy based on configuration
if self.commit:
# Specific commit requested - index all files at that commit
logger.info(f"Indexing specific commit: {self.commit}")
self._index_all_files()
elif self.last_indexed_commit:
elif self.last_indexed_commit and not last_indexing_had_errors:
# Incremental indexing - process diff since last indexed commit
logger.info(f"Incremental indexing since commit: {self.last_indexed_commit}")
self._index_diff_files()
Expand Down
10 changes: 9 additions & 1 deletion jobs/autoindexer/data_source_handler/kusto_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ def __init__(self, index_name: str, config: dict[str, Any], rag_client: KAITORAG
self.language = self.config.get("language")
self.initial_query = self.config.get("initialQuery")
self.incremental_query = self.config.get("incrementalQuery")
self.conditions = self.config.get("conditions", [])

self.errors = []
self.total_time = None
Expand Down Expand Up @@ -98,8 +99,15 @@ def _build_query(self) -> str:
For incremental queries, replaces $LAST_INDEXING_TIMESTAMP with the actual timestamp.
"""
last_timestamp = self._get_last_checkpoint_time()

last_indexing_had_errors = any(
condition.get("type") == "AutoIndexerError" and condition.get("status") == "True"
for condition in self.conditions
)
if last_indexing_had_errors:
logger.warning("Previous indexing had errors, using initial query to recover")

if last_timestamp is None:
if last_timestamp is None or last_indexing_had_errors:
# First run: use initial query
logger.info("🆕 FIRST RUN: Using initialQuery")
query = self.initial_query
Expand Down
5 changes: 4 additions & 1 deletion jobs/autoindexer/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -231,14 +231,16 @@ def _apply_crd_config(self, crd_config: dict[str, Any]):
"paths": git_config.get("paths", []),
"excludePaths": git_config.get("excludePaths", []),
"lastIndexedCommit": crd_config.get("status", {}).get("lastIndexedCommit", ""),
"conditions": crd_config.get("status", {}).get("conditions", [])
})
logger.info("Updated Git data source configuration from CRD")

elif ds_config.get("static") and ds_config["type"] == "Static":
static_config = ds_config["static"]
self.datasource_config.update({
"autoindexer_name": autoindexer_full_name,
"urls": static_config.get("urls", [])
"urls": static_config.get("urls", []),
"conditions": crd_config.get("status", {}).get("conditions", [])
})
logger.info("Updated Static data source configuration from CRD")

Expand All @@ -249,6 +251,7 @@ def _apply_crd_config(self, crd_config: dict[str, Any]):
"language": database_config.get("language"),
"initialQuery": database_config.get("initialQuery"),
"incrementalQuery": database_config.get("incrementalQuery"),
"conditions": crd_config.get("status", {}).get("conditions", [])
})
logger.info(f"Updated Database data source configuration from CRD (language: {database_config.get('language')})")

Expand Down
Loading
Loading