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
93 changes: 70 additions & 23 deletions apollo/server/routes/admin_workflows.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@
from fastapi.responses import HTMLResponse, RedirectResponse, JSONResponse
from typing import List, Optional

from temporalio.client import RPCError

from apollo.server.utils import templates, admin_user_scheme
from apollo.server.services.workflow_service import WorkflowService
from apollo.server.services.database_service import DatabaseService
Expand All @@ -23,6 +25,9 @@ async def admin_workflows(request: Request, user: User = Depends(admin_user_sche
env_info = await db_service.get_environment_info()
index_state = await db_service.get_last_indexed_at()

tracked_wf_id = request.session.pop("tracked_workflow_id", None)
tracked_wf_name = request.session.pop("tracked_workflow_name", None)

return templates.TemplateResponse(
"admin_workflows.jinja", {
"request": request,
Expand All @@ -32,6 +37,8 @@ async def admin_workflows(request: Request, user: User = Depends(admin_user_sche
"reset_allowed": env_info["reset_allowed"],
"last_indexed_at": index_state.get("last_indexed_at_iso"),
"last_indexed_exists": index_state.get("exists", False),
"tracked_workflow_id": tracked_wf_id,
"tracked_workflow_name": tracked_wf_name,
}
)

Expand All @@ -44,30 +51,50 @@ async def trigger_rh_matcher(
):
"""Trigger RHMatcherWorkflow from admin web interface"""
try:
# Convert string versions to integers if provided
int_versions = None
if major_versions:
int_versions = [int(v) for v in major_versions if v.isdigit()]

service = WorkflowService()
workflow_id = await service.trigger_rh_matcher_workflow(int_versions)

Logger().info(f"Admin user {user.email} triggered RhMatcherWorkflow {workflow_id} with major_versions: {int_versions}")

# Store success message in session
request.session["workflow_message"] = f"RhMatcherWorkflow triggered successfully: {workflow_id}"

Logger().info(
f"Admin user {user.email} triggered RhMatcherWorkflow "
f"{workflow_id} with major_versions: {int_versions}"
)

request.session["workflow_message"] = (
f"RhMatcherWorkflow triggered successfully: {workflow_id}"
)
request.session["workflow_type"] = "success"

request.session["tracked_workflow_id"] = workflow_id
request.session["tracked_workflow_name"] = "RH Matcher"

except RPCError as e:
if "already running" in str(e).lower():
request.session["workflow_message"] = (
"A workflow for this version is already running"
)
request.session["workflow_type"] = "error"
else:
Logger().error(f"Temporal error triggering workflow: {e}")
request.session["workflow_message"] = (
f"Workflow service error: {e}"
)
request.session["workflow_type"] = "error"

except ValueError as e:
Logger().error(f"Validation error triggering RhMatcher workflow: {str(e)}")
request.session["workflow_message"] = f"Error: {str(e)}"
Logger().error(f"Validation error triggering workflow: {e}")
request.session["workflow_message"] = f"Error: {e}"
request.session["workflow_type"] = "error"

except Exception as e:
Logger().error(f"Error triggering RhMatcher workflow: {str(e)}")
request.session["workflow_message"] = f"Error triggering workflow: {str(e)}"
Logger().error(f"Error triggering RhMatcher workflow: {e}")
request.session["workflow_message"] = (
f"Error triggering workflow: {e}"
)
request.session["workflow_type"] = "error"

return RedirectResponse(url="/admin/workflows", status_code=303)


Expand All @@ -80,21 +107,41 @@ async def trigger_poll_rhcsaf(
try:
service = WorkflowService()
workflow_id = await service.trigger_poll_rhcsaf_workflow()

Logger().info(f"Admin user {user.email} triggered PollRHCSAFAdvisoriesWorkflow {workflow_id}")

# Store success message in session
request.session["workflow_message"] = f"PollRHCSAFAdvisoriesWorkflow triggered successfully: {workflow_id}"

Logger().info(
f"Admin user {user.email} triggered "
f"PollRHCSAFAdvisoriesWorkflow {workflow_id}"
)

request.session["workflow_message"] = (
f"PollRHCSAFAdvisoriesWorkflow triggered successfully: "
f"{workflow_id}"
)
request.session["workflow_type"] = "success"

request.session["tracked_workflow_id"] = workflow_id
request.session["tracked_workflow_name"] = "Poll RHCSAF"

except Exception as e:
Logger().error(f"Error triggering PollRHCSAF workflow: {str(e)}")
request.session["workflow_message"] = f"Error triggering workflow: {str(e)}"
Logger().error(f"Error triggering PollRHCSAF workflow: {e}")
request.session["workflow_message"] = (
f"Error triggering workflow: {e}"
)
request.session["workflow_type"] = "error"

return RedirectResponse(url="/admin/workflows", status_code=303)


@router.get("/workflows/{workflow_id}/status")
async def admin_workflow_status(
workflow_id: str,
user: User = Depends(admin_user_scheme),
):
"""Return workflow status as JSON for the admin polling UI."""
service = WorkflowService()
status_info = await service.get_workflow_status(workflow_id)
return JSONResponse(status_info)


@router.post("/workflows/update-index-timestamp")
async def update_index_timestamp(
request: Request,
Expand Down
45 changes: 30 additions & 15 deletions apollo/server/routes/api_workflows.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@
from fastapi import APIRouter, Depends, HTTPException, status
from typing import List

from temporalio.client import RPCError

from apollo.server.models.workflow import (
ProductListResponse, ProductInfo, WorkflowTriggerRequest,
WorkflowTriggerResponse, WorkflowStatusResponse, WorkflowListResponse
Expand Down Expand Up @@ -51,40 +53,53 @@ async def trigger_rh_matcher_workflow(
"""
try:
service = WorkflowService()
workflow_id = await service.trigger_rh_matcher_workflow(request.major_versions)

workflow_id = await service.trigger_rh_matcher_workflow(
request.major_versions,
)

logger = Logger()
logger.info(f"User {user.email} triggered RhMatcherWorkflow {workflow_id} with major_versions: {request.major_versions}")

logger.info(
f"User {user.email} triggered RhMatcherWorkflow "
f"{workflow_id} with major_versions: {request.major_versions}"
)

return WorkflowTriggerResponse(
workflow_id=workflow_id,
status="started",
message="RhMatcherWorkflow triggered successfully",
filtered_major_versions=request.major_versions
filtered_major_versions=request.major_versions,
)

except ValueError as e:
# Handle validation errors (invalid major versions)

except RPCError as e:
if "already running" in str(e).lower():
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail="A workflow for this version is already running",
)
logger = Logger()
logger.error(f"Validation error triggering workflow: {str(e)}")
logger.error(f"Temporal RPC error triggering workflow: {e}")
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Workflow service unavailable",
)
except ValueError as e:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=str(e)
detail=str(e),
)
except RuntimeError as e:
# Handle Temporal client errors
logger = Logger()
logger.error(f"Runtime error triggering workflow: {str(e)}")
logger.error(f"Runtime error triggering workflow: {e}")
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Workflow service unavailable"
detail="Workflow service unavailable",
)
except Exception as e:
logger = Logger()
logger.error(f"Error triggering workflow: {str(e)}")
logger.error(f"Error triggering workflow: {e}")
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail="Failed to trigger workflow"
detail="Failed to trigger workflow",
)


Expand Down
Loading
Loading