diff --git a/tools/caretaker-agent/cloudrun/triage-worker/.dockerignore b/tools/caretaker-agent/cloudrun/triage-worker/.dockerignore new file mode 100644 index 00000000000..705caee9b23 --- /dev/null +++ b/tools/caretaker-agent/cloudrun/triage-worker/.dockerignore @@ -0,0 +1,9 @@ +__pycache__/ +*.pyc +*.pyo +*.pyd +.pytest_cache/ +venv/ +experimental/ +tests/ +.env diff --git a/tools/caretaker-agent/cloudrun/triage-worker/Dockerfile b/tools/caretaker-agent/cloudrun/triage-worker/Dockerfile new file mode 100644 index 00000000000..35341c4bd4e --- /dev/null +++ b/tools/caretaker-agent/cloudrun/triage-worker/Dockerfile @@ -0,0 +1,11 @@ +FROM python:3.13-slim +WORKDIR /app +ENV PYTHONUNBUFFERED=1 +RUN apt-get update && apt-get install -y \ + git \ + curl \ + && rm -rf /var/lib/apt/lists/* +RUN git clone https://github.com/google-gemini/gemini-cli.git /opt/gemini-cli +COPY . . +RUN pip install --no-cache-dir -r requirements.txt +CMD ["python", "main.py"] \ No newline at end of file diff --git a/tools/caretaker-agent/cloudrun/triage-worker/main.py b/tools/caretaker-agent/cloudrun/triage-worker/main.py new file mode 100644 index 00000000000..a562dfeabdb --- /dev/null +++ b/tools/caretaker-agent/cloudrun/triage-worker/main.py @@ -0,0 +1,152 @@ +import os +import json +import base64 +import sys + +from google.cloud import firestore +from triage_orchestrator import process_issue_triage +from utils.validator import validate_triage_result +from utils.egress import send_label_action, send_comment_action +from db.issues_store import IssuesStore, ClaimAction, ReleaseAction + + +def main() -> None: + """ + Main entrypoint for Caretaker Triage Cloud Run Job worker. + + Assumptions: + - Assumes ISSUE_DETAILS env var contains base64-encoded JSON payload. + - Assumes WORKFLOW_EXECUTION_ID contains unique lock holder ID. + """ + # Cloud Run Jobs inject data via environment variables + encoded_data = os.environ.get("ISSUE_DETAILS") + + if not encoded_data: + print("[PROD] Error: No data provided in ISSUE_DETAILS.") + sys.exit(1) + + try: + payload = json.loads(base64.b64decode(encoded_data)) + except Exception as e: + print(f"[PROD] Error decoding payload: {e}") + sys.exit(1) + + try: + issue_number = int(payload.get("issue_number")) + except (TypeError, ValueError): + print("[PROD] Error: issue_number is not a valid number. Exiting.") + sys.exit(1) + + try: + owner, repo = payload.get("repository", "").split("/") + if not owner or not repo: + raise ValueError + except (TypeError, ValueError): + print("[PROD] Error: Malformed repository format (expected 'owner/repo'). Exiting.") + sys.exit(1) + + lock_holder = os.environ.get("WORKFLOW_EXECUTION_ID", "local-exec") + + # Initialize Firestore Client & IssuesStore + project_id = os.environ.get("PROJECT_ID") + db_id = os.environ.get("FIRESTORE_DATABASE") + collection_name = os.environ.get("FIRESTORE_COLLECTION", "issues") + + db_client = firestore.Client(project=project_id, database=db_id) + store = IssuesStore(db_client, collection_name) + + # Claim the lock + claim_action = store.acquire_lock(owner, repo, issue_number, lock_holder) + + if claim_action == ClaimAction.SKIP: + print( + f"[WORKER] Issue #{issue_number} already handled or active lock " + "present. Exiting." + ) + sys.exit(0) + elif claim_action == ClaimAction.NEEDS_HUMAN: + print(f"[WORKER] Issue #{issue_number} requires human review. Exiting.") + sys.exit(0) + + print(f"[WORKER] Starting triage for issue #{issue_number}...") + try: + success, raw_output = process_issue_triage(payload) + except Exception as e: + print(f"[WORKER] Triage process failed with exception: {e}") + success, raw_output = False, "" + + if success: + try: + triage_result = json.loads(raw_output) + validate_triage_result(triage_result) + + quality = triage_result.get("triage_metadata", {}).get("quality") + workable_spec = triage_result.get("workable_spec", {}) + + if quality in ["SPAM", "EMPTY", "FEATURE"]: + print(f"[WORKER] Quality: {quality}. Applying needs_human label.") + send_label_action(owner, repo, issue_number, ["needs_human"]) + store.release_lock( + owner, + repo, + issue_number, + lock_holder, + success=True, + status="LOW_QUALITY", + ) + sys.exit(0) + elif quality == "NEEDS_INFO": + print(f"[WORKER] Quality: NEEDS_INFO. Leaving comment.") + comment_body = ( + triage_result.get("triage_metadata", {}) + .get("comment", "") + .strip() + ) + send_comment_action(owner, repo, issue_number, comment_body) + store.release_lock( + owner, + repo, + issue_number, + lock_holder, + success=True, + status="NEEDS_INFO", + ) + sys.exit(0) + else: + effort = triage_result.get("triage_metadata", {}).get( + "effort_estimate" + ) + print( + f"[WORKER] Quality: OK. Effort: {effort}. Applying " + "effort label." + ) + send_label_action( + owner, repo, issue_number, [f"effort/{effort.lower()}"] + ) + store.release_lock( + owner, + repo, + issue_number, + lock_holder, + success=True, + status="TRIAGED", + workable_spec=workable_spec, + ) + print(f"[WORKER] Triage success.") + sys.exit(0) + + except Exception as e: + print(f"[WORKER] Validation failed: {e}") + success = False + + # If an exception happens in json.loads or validate_triage_result + # If LLM inference itself fails inside process_issue_triage + if not success: + release_action = store.release_lock( + owner, repo, issue_number, lock_holder, success=False + ) + sys.exit(1 if release_action == ReleaseAction.RETRY else 0) + + +if __name__ == "__main__": + main() diff --git a/tools/caretaker-agent/cloudrun/triage-worker/requirements.txt b/tools/caretaker-agent/cloudrun/triage-worker/requirements.txt index 4a38e337d5f..fa7b37ec269 100644 --- a/tools/caretaker-agent/cloudrun/triage-worker/requirements.txt +++ b/tools/caretaker-agent/cloudrun/triage-worker/requirements.txt @@ -1 +1,5 @@ google-cloud-firestore>=2.15.0, <3.0.0 +google-cloud-storage +google-cloud-pubsub +google-antigravity>=0.1.0 +python-dotenv \ No newline at end of file diff --git a/tools/caretaker-agent/cloudrun/triage-worker/tests/test_integration_main.py b/tools/caretaker-agent/cloudrun/triage-worker/tests/test_integration_main.py new file mode 100644 index 00000000000..00765849252 --- /dev/null +++ b/tools/caretaker-agent/cloudrun/triage-worker/tests/test_integration_main.py @@ -0,0 +1,215 @@ +""" +Integration tests for main.py. + +Verifies workflow execution across main.py and issues_store.py. +External network boundaries (Firestore database client and LLM inference) are mocked. +""" + +import unittest +from unittest.mock import patch, MagicMock +import os +import json +import base64 +from db.issues_store import IssuesStore, ClaimAction, ReleaseAction +import main as main_module +from main import main + +VALID_WORKABLE_SPEC = { + "issue_id": "owner/repo#42", + "summary": {"problem": "p", "root_cause": "r", "context": "c"}, + "implementation_plan": {"files_to_modify": ["src/app.ts"], "steps": ["Fix bug"]}, + "testing_strategy": {"test_file": "tests/app.test.ts", "expected_behavior": "Pass", "verification_steps": ["Check"], "framework": "Vitest"} +} + +INTEGRATION_OK_PAYLOAD = { + "triage_metadata": { + "quality": "OK", + "reasoning": "Actionable bug report.", + "comment": "", + "effort_estimate": "SMALL", + "effort_reasoning": "Easy fix." + }, + "workable_spec": VALID_WORKABLE_SPEC +} + +INTEGRATION_NEEDS_INFO_PAYLOAD = { + "triage_metadata": { + "quality": "NEEDS_INFO", + "reasoning": "The issue reports a crash on startup, but lacks any actual details.", + "comment": "Hi! Thanks for commenting on this issue, we need more information to triage the bug.", + "effort_estimate": "", + "effort_reasoning": "" + }, + "workable_spec": {} +} + +INTEGRATION_INVALID_EFFORT_PAYLOAD = { + "triage_metadata": { + "quality": "OK", + "reasoning": "Some reasoning.", + "comment": "", + "effort_estimate": "HUGE", + "effort_reasoning": "This will take a while to fix." + }, + "workable_spec": VALID_WORKABLE_SPEC +} + +class TestIntegrationMain(unittest.TestCase): + + def setUp(self): + # Mock environment variables + self.env_patcher = patch.dict(os.environ, { + "ISSUE_DETAILS": base64.b64encode(json.dumps({ + "issue_number": 42, + "repository": "owner/repo", + "title": "Fix crash", + "body": "App crashes on start" + }).encode("utf-8")).decode("utf-8"), + "WORKFLOW_EXECUTION_ID": "test-workflow-exec-101", + "PROJECT_ID": "test-gcp-project", + "EGRESS_TOPIC_ID": "test-egress-actions" + }) + self.env_patcher.start() + + # Mock the Firestore database client at the network boundary + self.mock_db = MagicMock() + self.db_patcher = patch("main.firestore.Client", return_value=self.mock_db) + self.db_patcher.start() + + self.mock_doc_ref = MagicMock() + self.mock_snapshot = MagicMock() + self.mock_transaction = MagicMock() + + self.mock_db.collection.return_value.document.return_value = self.mock_doc_ref + self.mock_db.transaction.return_value = self.mock_transaction + self.mock_doc_ref.get.return_value = self.mock_snapshot + self.mock_snapshot.exists = True + + # In-memory document state simulation + self.stored_data = {} + self.mock_snapshot.to_dict.side_effect = lambda: self.stored_data + + def mock_update(doc_ref, updates): + if "status" in updates: + self.stored_data["status"] = updates["status"] + if "workable_spec" in updates: + self.stored_data["workable_spec"] = updates["workable_spec"] + if "lock.holder" in updates: + if "lock" not in self.stored_data: + self.stored_data["lock"] = {} + self.stored_data["lock"]["holder"] = updates["lock.holder"] + + # Bind mock_update to execute whenever transaction.update is invoked + self.mock_transaction.update.side_effect = mock_update + + # Mock IssuesStore instance + self.mock_store = MagicMock() + self.store_patcher = patch("main.IssuesStore", return_value=self.mock_store) + self.store_patcher.start() + + # Wire mock_store methods to execute real store logic against mock_db + real_store = IssuesStore(self.mock_db, "issues") + self.mock_store.acquire_lock.side_effect = real_store.acquire_lock + self.mock_store.release_lock.side_effect = real_store.release_lock + + def tearDown(self): + self.store_patcher.stop() + self.db_patcher.stop() + self.env_patcher.stop() + + @patch("main.process_issue_triage") + @patch("main.send_label_action") + def test_main_ok_quality_flow(self, mock_send_label, mock_triage): + """Verifies end-to-end flow for OK quality issues.""" + self.stored_data = { + "status": "UNTRIAGED", + "triage_attempts": 0, + "lock": {"holder": None, "expires_at": None} + } + mock_triage.return_value = (True, json.dumps(INTEGRATION_OK_PAYLOAD)) + + with self.assertRaises(SystemExit) as ctx: + main() + + self.assertEqual(ctx.exception.code, 0) + self.mock_store.acquire_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101") + self.mock_store.release_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101", success=True, status="TRIAGED", workable_spec=INTEGRATION_OK_PAYLOAD["workable_spec"]) + mock_send_label.assert_called_once_with("owner", "repo", 42, ["effort/small"]) + + # Verify state transition in store data + self.assertEqual(self.stored_data["status"], "TRIAGED") + self.assertEqual(self.stored_data["workable_spec"], VALID_WORKABLE_SPEC) + self.assertIsNone(self.stored_data["lock"]["holder"]) + + @patch("main.process_issue_triage") + @patch("main.send_comment_action") + def test_main_needs_info_flow(self, mock_send_comment, mock_triage): + """Verifies end-to-end flow for NEEDS_INFO issues.""" + self.stored_data = { + "status": "UNTRIAGED", + "triage_attempts": 0, + "lock": {"holder": None, "expires_at": None} + } + mock_triage.return_value = (True, json.dumps(INTEGRATION_NEEDS_INFO_PAYLOAD)) + + with self.assertRaises(SystemExit) as ctx: + main() + + self.assertEqual(ctx.exception.code, 0) + self.mock_store.acquire_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101") + self.mock_store.release_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101", success=True, status="NEEDS_INFO") + mock_send_comment.assert_called_once_with("owner", "repo", 42, INTEGRATION_NEEDS_INFO_PAYLOAD["triage_metadata"]["comment"]) + + self.assertEqual(self.stored_data["status"], "NEEDS_INFO") + self.assertIsNone(self.stored_data["lock"]["holder"]) + + @patch("main.process_issue_triage") + @patch("main.send_label_action") + def test_main_low_quality_flows(self, mock_send_label, mock_triage): + """Verifies end-to-end flow for low quality issues.""" + for quality in ["SPAM", "EMPTY", "FEATURE"]: + self.mock_store.acquire_lock.reset_mock() + self.mock_store.release_lock.reset_mock() + mock_send_label.reset_mock() + mock_triage.reset_mock() + + self.stored_data = { + "status": "UNTRIAGED", + "triage_attempts": 0, + "lock": {"holder": None, "expires_at": None} + } + payload = {"triage_metadata": {"quality": quality}} + mock_triage.return_value = (True, json.dumps(payload)) + + with self.assertRaises(SystemExit) as ctx: + main() + + self.assertEqual(ctx.exception.code, 0) + self.mock_store.acquire_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101") + self.mock_store.release_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101", success=True, status="LOW_QUALITY") + mock_send_label.assert_called_once_with("owner", "repo", 42, ["needs_human"]) + + self.assertEqual(self.stored_data["status"], "LOW_QUALITY") + self.assertIsNone(self.stored_data["lock"]["holder"]) + + @patch("main.process_issue_triage") + def test_main_validation_failure_triggers_retry(self, mock_triage): + """Verifies retry state transition when validation fails.""" + self.stored_data = { + "status": "UNTRIAGED", + "triage_attempts": 0, + "lock": {"holder": None, "expires_at": None} + } + mock_triage.return_value = (True, json.dumps(INTEGRATION_INVALID_EFFORT_PAYLOAD)) + + with self.assertRaises(SystemExit) as ctx: + main() + + self.assertEqual(ctx.exception.code, 1) + self.mock_store.acquire_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101") + self.mock_store.release_lock.assert_called_once_with("owner", "repo", 42, "test-workflow-exec-101", success=False) + self.assertEqual(self.stored_data["status"], "UNTRIAGED") + self.assertIsNone(self.stored_data["lock"]["holder"]) + +if __name__ == "__main__": + unittest.main() diff --git a/tools/caretaker-agent/cloudrun/triage-worker/triage_orchestrator.py b/tools/caretaker-agent/cloudrun/triage-worker/triage_orchestrator.py new file mode 100644 index 00000000000..7c5f45d2522 --- /dev/null +++ b/tools/caretaker-agent/cloudrun/triage-worker/triage_orchestrator.py @@ -0,0 +1,74 @@ +import os +import asyncio +from utils.agent_logger import upload_to_bucket, upload_agent_run_logs, extract_final_output +from google.antigravity import Agent, LocalAgentConfig +from google.antigravity.hooks.policy import allow, deny + +def process_issue_triage(payload: dict) -> tuple[bool, str]: + """ + LLM inference via Antigravity SDK. + """ + issue_num = payload.get("issue_number") + title = payload.get("title", "") + body = payload.get("body", "") + repo_name = payload.get("repository", "") + + current_dir = os.path.dirname(os.path.abspath(__file__)) + system_prompt_path = os.path.join(current_dir, ".gemini", "triage_orchestrator.md") + target_cwd = "/opt/gemini-cli" + + + policies = [ + # Deny all tools by default + deny("*"), + + # Whitelist specific read-only and skill tools + allow("view_file"), + allow("list_directory"), + allow("find_file"), + allow("search_directory"), + allow("activate_skill"), + allow("finish") + ] + + with open(system_prompt_path, "r") as f: + system_instructions = f.read() + + skills_dir = os.path.join(current_dir, ".gemini", "skills") + + prompt = f"Repository: {repo_name}\nIssue Number: {issue_num}\nTitle: {title}\nDescription: {body}" + + async def run_triage(): + config = LocalAgentConfig( + system_instructions=system_instructions, + skills_paths=[skills_dir], + api_key=os.environ.get("GEMINI_API_KEY"), + workspaces=[target_cwd, current_dir], + policies=policies, + ) + + print("[LOGIC] Initializing Antigravity Agent...") + async with Agent(config) as agent: + print("[LOGIC] Sending triage request to Agent...") + response = await agent.chat(prompt) + + # Resolve all execution chunks (thoughts, tool calls, and tool results) + resolved_chunks = await response.resolve() + + # Extract the final step's output + text_output = extract_final_output(resolved_chunks) + + upload_agent_run_logs(repo_name, issue_num, resolved_chunks) + + print(f"[LOGIC] Agent Response:\n{text_output}") + + return True, text_output + + try: + success, raw_output = asyncio.run(run_triage()) + return success, raw_output + except Exception as e: + error_msg = f"Error during Antigravity Agent run: {e}" + print(f"[LOGIC] {error_msg}") + upload_to_bucket(repo_name, issue_num, error_msg) + return False, error_msg diff --git a/tools/caretaker-agent/cloudrun/triage-worker/utils/agent_logger.py b/tools/caretaker-agent/cloudrun/triage-worker/utils/agent_logger.py new file mode 100644 index 00000000000..b08a7152483 --- /dev/null +++ b/tools/caretaker-agent/cloudrun/triage-worker/utils/agent_logger.py @@ -0,0 +1,68 @@ +import json +import datetime +import os +from google.cloud import storage +from google.antigravity.types import Text + +BUCKET_NAME = os.environ.get("TRIAGE_DEBUG_LOGS_BUCKET") + +def upload_to_bucket(repository: str, issue_number: str | int, payload: str) -> None: + """ + Uploads a string payload directly to the triage debug logs GCS bucket. + """ + if not payload or not BUCKET_NAME: + if not BUCKET_NAME: + print("[LOGIC] Warning: Missing TRIAGE_DEBUG_LOGS_BUCKET, skipping GCS upload.") + return + try: + storage_client = storage.Client() + bucket = storage_client.bucket(BUCKET_NAME) + + timestamp = datetime.datetime.now(datetime.timezone.utc).strftime("%Y%m%d_%H%M%S") + safe_repo = repository.replace("/", "_") + blob_name = f"{safe_repo}/issue_{issue_number}_{timestamp}_debug.log" + + blob = bucket.blob(blob_name) + blob.upload_from_string(payload, content_type="text/plain") + print(f"[LOGIC] Uploaded debug logs to gs://{BUCKET_NAME}/{blob_name}") + except Exception as e: + print(f"[LOGIC] Error uploading debug logs to GCS: {e}") + +def upload_agent_run_logs(repository: str, issue_number: str | int, resolved_chunks: list) -> None: + """ + Serializes the Antigravity Agent resolved stream chunks to JSON and uploads them to GCS. + """ + if not resolved_chunks: + return + # Format chunks into a serializable format + try: + serializable = [] + for chunk in resolved_chunks: + try: + dumped = chunk.model_dump() + except AttributeError: + dumped = chunk.dict() + dumped["chunk_type"] = chunk.__class__.__name__ + serializable.append(dumped) + + debug_log_str = json.dumps(serializable, indent=2, default=str) + upload_to_bucket(repository, issue_number, debug_log_str) + except Exception as e: + print(f"[LOGIC] Error serializing/uploading agent debug logs: {e}") + +def extract_final_output(resolved_chunks: list) -> str: + """ + Extracts the final response text generated by the agent during the last execution step. + + Since the agent outputs conversational planning text during its tool-calling turns, + this function isolates the final turn (highest step_index) and filters for text chunks, + stripping away any intermediate planning text, thoughts, or markdown backticks around the JSON. + """ + # Filter for Text chunks (skip Thought and ToolCall chunks) + text_chunks = [c for c in resolved_chunks if isinstance(c, Text)] + if not text_chunks: + return "" + + # Final output is in the last step + last_step = text_chunks[-1].step_index + return "".join(c.text for c in text_chunks if c.step_index == last_step).strip() \ No newline at end of file diff --git a/tools/caretaker-agent/cloudrun/triage-worker/utils/egress.py b/tools/caretaker-agent/cloudrun/triage-worker/utils/egress.py new file mode 100644 index 00000000000..d6192c1fb90 --- /dev/null +++ b/tools/caretaker-agent/cloudrun/triage-worker/utils/egress.py @@ -0,0 +1,63 @@ +import os +import json +from google.cloud import pubsub_v1 + +def _publish_egress_action(egress_event: dict) -> None: + """ + [Internal] Publishes an EgressEvent JSON payload to Pub/Sub. + """ + project_id = os.environ.get("PROJECT_ID") + egress_topic_id = os.environ.get("EGRESS_TOPIC_ID") + + if not project_id or not egress_topic_id: + print( + f"[WORKER] Warning: Missing PROJECT_ID ({project_id}) or " + f"EGRESS_TOPIC_ID ({egress_topic_id}), skipping egress." + ) + return + try: + publisher = pubsub_v1.PublisherClient() + topic_path = publisher.topic_path(project_id, egress_topic_id) + data = json.dumps(egress_event).encode("utf-8") + future = publisher.publish(topic_path, data) + message_id = future.result() + print( + f"[WORKER] Published egress action to Pub/Sub ({egress_topic_id}). " + f"Message ID: {message_id}" + ) + except Exception as e: + print(f"[WORKER] Error publishing to Pub/Sub: {e}") + + +def send_label_action( + owner: str, repo: str, issue_number: int, labels: list[str] +) -> None: + """ + Helper to publish a LABEL action to egress. + """ + _publish_egress_action({ + "action": "LABEL", + "payload": { + "owner": owner, + "repo": repo, + "issueNumber": issue_number, + "labels": labels, + }, + }) + + +def send_comment_action( + owner: str, repo: str, issue_number: int, comment_body: str +) -> None: + """ + Helper to publish a COMMENT action to egress. + """ + _publish_egress_action({ + "action": "COMMENT", + "payload": { + "owner": owner, + "repo": repo, + "issueNumber": issue_number, + "commentBody": comment_body, + }, + })