Skip to content
Open
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
37 changes: 30 additions & 7 deletions backend/api/plugins/github/plugin.py
Original file line number Diff line number Diff line change
Expand Up @@ -303,16 +303,39 @@ async def fetch_installation_repositories(self, installation_token: str) -> list
"Accept": "application/vnd.github.v3+json",
}

try:
response = await self.http.get("/installation/repositories", headers=headers)
all_repos: list[dict] = []
page = 1
per_page = 100

if response.status_code != 200:
raise HTTPException(
status_code=response.status_code, detail=f"GitHub API error: {response.text}"
try:
while True:
response = await self.http.get(
"/installation/repositories",
headers=headers,
params={"page": page, "per_page": per_page},
)

data = response.json()
return data.get("repositories", [])
if response.status_code != 200:
raise HTTPException(
status_code=response.status_code,
detail=f"GitHub API error: {response.text}",
)

data = response.json()
repos = data.get("repositories", [])

if not repos:
break

all_repos.extend(repos)

if len(repos) < per_page:
break

page += 1

logger.info(f"Fetched {len(all_repos)} total repositories for installation")
return all_repos

except HTTPException as e:
raise e
Expand Down
142 changes: 142 additions & 0 deletions backend/api/routers/sources/gitlab/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
- Sync repositories when a group token is added
- Sync merge requests for each repository
- Sync group members
- Handle project_create system hook events
"""

import logging
Expand All @@ -16,7 +17,9 @@
from api.database import (
db_create_repository,
db_create_repository_membership,
db_get_all_organizations,
db_get_or_create_provider_identity,
db_get_organization_by_external_id,
db_get_organization_by_id,
db_get_repository_by_external_id,
db_get_repository_by_id,
Expand All @@ -26,6 +29,7 @@
from api.plugins.faststream import get_faststream_broker
from api.plugins.faststream.config import STREAM_MAXLEN
from api.routers.sources.gitlab.schemas import (
GitLabProjectCreatedMessage,
GitLabSyncGroupMessage,
GitLabSyncGroupRepositoriesMessage,
GitLabSyncMembersMessage,
Expand Down Expand Up @@ -372,3 +376,141 @@ async def sync_group_members(message: GitLabSyncMembersMessage) -> None:
)

logger.info(f"Synced {len(members)} members for group {group_id}")


@router.subscriber(
stream=StreamSub(
"reviewate.events.gitlab.project_created", group="reviewate", consumer="worker-1"
)
)
async def handle_project_created(message: GitLabProjectCreatedMessage) -> None:
"""Handle a GitLab project_create system hook event.

The system hook payload includes project_namespace_id, which is the numeric
group ID — matching directly against external_org_id in the DB, no API call needed.
Projects in subgroups (where namespace_id != root group id) are skipped for now.

Args:
message: Project created message from system hook
"""
project_id = str(message.project_id)
namespace_id = str(message.namespace_id)
path_with_namespace = message.path_with_namespace

logger.info(f"Processing project_create for {path_with_namespace} (id={project_id})")

app = get_current_app()
gitlab_plugin = app.gitlab
encryptor = get_encryptor()

# Phase 1: Look up org by namespace_id.
# For root-group projects this matches directly (no API call).
# For subgroup projects we resolve the root group first using any available token.
with app.database.session() as db:
org = db_get_organization_by_external_id(db, namespace_id, provider="gitlab")

if not org:
# Subgroup case: walk up to root using any available token
any_encrypted = next(
(o.gitlab_access_token_encrypted for o in db_get_all_organizations(db)
if o.provider == "gitlab" and o.gitlab_access_token_encrypted),
None,
)
if not any_encrypted:
logger.debug("No GitLab token available to resolve subgroup — skipping")
return
resolve_token = encryptor.decrypt(any_encrypted)

else:
resolve_token = None # not needed, already found

if not org:
try:
root_group = await gitlab_plugin.resolve_root_group(resolve_token, namespace_id)
root_group_id = str(root_group["id"])
except Exception as e:
logger.debug(f"Could not resolve root group for namespace {namespace_id}: {e}")
return

with app.database.session() as db:
org = db_get_organization_by_external_id(db, root_group_id, provider="gitlab")

if not org:
logger.debug(f"Namespace {namespace_id} not registered in Reviewate — skipping {path_with_namespace}")
return

with app.database.session() as db:
existing = db_get_repository_by_external_id(db, project_id)
if existing:
logger.debug(f"Repository {path_with_namespace} (id={project_id}) already exists")
return

if not org.gitlab_access_token_encrypted:
logger.warning(f"Org {org.id} has no GitLab token — cannot fetch project {path_with_namespace}")
return

org_id = org.id
org_provider_url = org.provider_url or "https://gitlab.com"
access_token = encryptor.decrypt(org.gitlab_access_token_encrypted)

# Phase 2: Fetch full project info (need web_url, avatar, etc.)
try:
project_info = await gitlab_plugin.fetch_project(access_token, project_id)
except Exception as e:
logger.error(
f"Failed to fetch project info for {path_with_namespace} (id={project_id}): {e}",
exc_info=True,
)
return

project_avatar_url = project_info.get("avatar_url") or project_info.get(
"namespace", {}
).get("avatar_url")

# Phase 3: Create repo record and publish SSE
with app.database.session() as db:
repository = db_create_repository(
db=db,
organization_id=org_id,
external_repo_id=project_id,
name=project_info["name"],
web_url=project_info["web_url"],
provider="gitlab",
provider_url=org_provider_url,
avatar_url=project_avatar_url,
)

logger.info(f"Created repository: {project_info['name']} ({repository.id})")

try:
await publish_repository_event(
organization_id=str(org_id),
action="created",
repository={
"id": str(repository.id),
"name": repository.name,
"external_repo_id": repository.external_repo_id,
"web_url": repository.web_url,
"provider": repository.provider,
"created_at": repository.created_at.isoformat(),
"updated_at": repository.updated_at.isoformat(),
},
)
except Exception as e:
logger.error(f"Failed to publish repository SSE event: {e}", exc_info=True)

# Phase 4: Queue MR sync
try:
broker = get_faststream_broker()
sync_message = GitLabSyncRepositoryMRsMessage(
repository_id=str(repository.id),
organization_id=str(org_id),
)
await broker.publish(
sync_message,
stream="reviewate.events.gitlab.sync_repository_mrs",
maxlen=STREAM_MAXLEN,
)
logger.debug(f"Queued MR sync for new project {path_with_namespace}")
except Exception as e:
logger.error(f"Failed to queue MR sync for {path_with_namespace}: {e}")
11 changes: 11 additions & 0 deletions backend/api/routers/sources/gitlab/schemas.py
Original file line number Diff line number Diff line change
Expand Up @@ -108,3 +108,14 @@ class GitLabSyncMembersMessage(BaseModel):
default=None,
description="Encrypted access token override (for subgroup tokens)",
)


class GitLabProjectCreatedMessage(BaseModel):
"""Message schema for handling a GitLab project_create system hook event.

Published when a GitLab system hook fires for a new project.
"""

project_id: int = Field(description="GitLab project ID")
namespace_id: int = Field(description="Namespace (group) ID that owns the project")
path_with_namespace: str = Field(description="Full project path (e.g. myorg/myproject)")
28 changes: 26 additions & 2 deletions backend/api/routers/webhooks/github/installations.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
GitHubAppInstallationRepositoriesEvent,
GitHubSyncInstallationMessage,
GitHubSyncMembersMessage,
GitHubSyncRepositoryPRsMessage,
)

logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -264,7 +265,10 @@ async def handle_repositories_added(
if existing_repo:
continue # Skip if already exists

# Extract avatar URL from owner (repo's owner avatar)
# repositories_added items are minimal repo objects (id, name, full_name, private)
# — html_url and owner are not included, so we derive them from full_name.
full_name = repo_info.get("full_name", repo_info["name"])
web_url = repo_info.get("html_url") or f"https://github.com/{full_name}"
repo_avatar_url = repo_info.get("owner", {}).get("avatar_url")

# Create repository
Expand All @@ -273,7 +277,7 @@ async def handle_repositories_added(
organization_id=org.id,
external_repo_id=external_repo_id,
name=repo_info["name"],
web_url=repo_info.get("html_url", ""),
web_url=web_url,
provider="github",
provider_url="https://github.com",
avatar_url=repo_avatar_url,
Expand Down Expand Up @@ -308,6 +312,26 @@ async def handle_repositories_added(
except Exception as e:
logger.error(f"Failed to publish repository SSE event: {e}", exc_info=True)

# Queue PR sync for the newly added repository
try:
owner, _, repo_name = full_name.partition("/")
if owner and repo_name:
broker = get_faststream_broker()
sync_message = GitHubSyncRepositoryPRsMessage(
repository_id=str(repository.id),
installation_id=installation_id,
owner=owner,
repo_name=repo_name,
)
await broker.publish(
sync_message,
stream="reviewate.events.github.sync_repository_prs",
maxlen=STREAM_MAXLEN,
)
logger.debug(f"Queued PR sync for repository {full_name}")
except Exception as e:
logger.error(f"Failed to queue PR sync for {repo_info.get('name')}: {e}")

return WebhookResponse(
message=f"Added {added_count} repositories to organization {org.name}",
processed=True,
Expand Down
42 changes: 38 additions & 4 deletions backend/api/routers/webhooks/gitlab/handlers.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
"""GitLab webhook handlers.

This module handles all GitLab webhook events.
This module handles all GitLab webhook events, including system hooks.
"""

import json
Expand All @@ -10,6 +10,9 @@
from sqlalchemy.orm import Session

from api.database import get_session
from api.plugins.faststream import get_faststream_broker
from api.plugins.faststream.config import STREAM_MAXLEN
from api.routers.sources.gitlab.schemas import GitLabProjectCreatedMessage

from ..utils import WebhookResponse
from .dependencies import verify_gitlab_webhook
Expand All @@ -18,6 +21,7 @@
from .schemas import (
GitLabMergeRequestEvent,
GitLabNoteEvent,
GitLabProjectCreatedSystemHookEvent,
)

logger = logging.getLogger(__name__)
Expand All @@ -34,7 +38,8 @@
name="gitlab_webhook",
summary="GitLab webhook router",
description=(
"Unified webhook endpoint for all GitLab events. Currently handles merge request events."
"Unified webhook endpoint for all GitLab events. Handles merge request events, "
"note events, and system hook events (project_create)."
),
response_model=WebhookResponse,
status_code=202,
Expand All @@ -45,6 +50,9 @@ async def gitlab_webhook(
) -> WebhookResponse:
"""Handle GitLab webhook events and route to appropriate handlers.

Supports both group/project webhooks (object_kind-based) and system hooks
(event_name-based). System hooks are admin-configured and fire globally.

Args:
request: FastAPI request object
db: Database session
Expand All @@ -60,8 +68,34 @@ async def gitlab_webhook(
# Parse event to determine type
event_data = json.loads(body)
object_kind = event_data.get("object_kind", "")
event_name = event_data.get("event_name", "")

# System hooks use event_name instead of object_kind
if event_name == "project_create":
event = GitLabProjectCreatedSystemHookEvent(**event_data)

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[CRITICAL] Pydantic ValidationError not caught — schema instantiation outside error handler

The schema instantiation event = GitLabProjectCreatedSystemHookEvent(**event_data) (line 75) is not wrapped in the try/except Exception block that follows (lines 76–92). Any required field mismatch or validation error from the schema will propagate as HTTP 500 instead of being caught and logged.

Fix: Wrap the instantiation:

if event_name == "project_create":
    try:
        event = GitLabProjectCreatedSystemHookEvent(**event_data)
    except Exception as e:
        logger.warning(f"Invalid project_create payload: {e}")
        return WebhookResponse(message="Invalid payload", processed=False)
    try:
        broker = get_faststream_broker()
        ...

try:
broker = get_faststream_broker()
msg = GitLabProjectCreatedMessage(
project_id=event.project_id,
namespace_id=event.project_namespace_id,
path_with_namespace=event.path_with_namespace,
)
await broker.publish(
msg,
stream="reviewate.events.gitlab.project_created",
maxlen=STREAM_MAXLEN,
)
logger.info(
f"Queued project_create sync for {event.path_with_namespace} (id={event.project_id})"
)
except Exception as e:
logger.error(f"Failed to queue project_create sync: {e}", exc_info=True)
return WebhookResponse(
message=f"Project creation queued for sync: {event.path_with_namespace}",
processed=True,
)

# Route based on event type
# Group/project webhooks use object_kind
if object_kind == "merge_request":
event = GitLabMergeRequestEvent(**event_data)
return await handle_merge_request_event(event, db)
Expand All @@ -70,6 +104,6 @@ async def gitlab_webhook(
return await handle_note_event(note_event, db)
else:
return WebhookResponse(
message=f"Event type '{object_kind}' not supported",
message=f"Event type '{object_kind or event_name}' not supported",
processed=False,
)
16 changes: 16 additions & 0 deletions backend/api/routers/webhooks/gitlab/schemas.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,22 @@ class GitLabNoteEvent(BaseModel):
repository: dict[str, Any] | None = Field(description="Repository object", default=None)


class GitLabProjectCreatedSystemHookEvent(BaseModel):
"""GitLab system hook payload for project_create events.

Fired by GitLab system hooks (admin-configured) when a new project is created
anywhere on the instance.
"""

event_name: str = Field(description="Event name (project_create)")
project_id: int = Field(description="GitLab project ID")
name: str = Field(description="Project name")
path: str = Field(description="Project path slug")
path_with_namespace: str = Field(description="Full path including namespace (e.g. myorg/myproject)")
project_namespace_id: int = Field(description="Namespace (group) ID that owns the project")
project_visibility: str | None = Field(default=None, description="Visibility level")


class GitLabFeedbackSignalMessage(BaseModel):
"""Message schema for GitLab feedback signal processing."""

Expand Down
Loading
Loading