mirror of
https://github.com/Sea-Haven-Industries/open-swe.git
synced 2026-09-30 17:23:15 +00:00
216 lines
8.9 KiB
Python
216 lines
8.9 KiB
Python
|
|
"""Confluence Atlassian Connect webhook: lifecycle + comment-created handler.
|
||
|
|
|
||
|
|
Lifecycle (installed/uninstalled) implements the verify-before-overwrite guard;
|
||
|
|
the comment handler mirrors ``webhooks/jira.py:process_jira_issue`` with the
|
||
|
|
Phase-2 lesson baked in: the JWT-signed webhook body is only a pointer, so the
|
||
|
|
triggering comment's real author, text, and container are re-fetched server-side
|
||
|
|
via the Basic-auth service account before anything security-relevant is derived.
|
||
|
|
"""
|
||
|
|
|
||
|
|
import os
|
||
|
|
from typing import Any
|
||
|
|
|
||
|
|
from langchain_core.messages.content import create_text_block
|
||
|
|
|
||
|
|
from agent import webapp
|
||
|
|
from agent.utils import atlassian_connect as ac
|
||
|
|
|
||
|
|
# The Connect app's own Confluence service-account accountId. When set, comments
|
||
|
|
# authored by it are ignored (self-trigger loop guard, like the Linear botActor
|
||
|
|
# / Jira comment_author_is_bot early-outs).
|
||
|
|
CONFLUENCE_BOT_ACCOUNT_ID = os.environ.get("CONFLUENCE_BOT_ACCOUNT_ID", "")
|
||
|
|
|
||
|
|
|
||
|
|
async def process_install(request: Any, body: dict[str, Any]) -> tuple[int, str]:
|
||
|
|
"""Handle POST /connect/installed. Returns (status_code, detail).
|
||
|
|
|
||
|
|
With signed-install, Atlassian RS256-signs every install callback (including
|
||
|
|
the first), so both first-install and re-install are verified against
|
||
|
|
Atlassian's published keys — there is no trust-on-first-use, and the
|
||
|
|
re-install path cannot be used to rotate our secret without a valid
|
||
|
|
Atlassian signature.
|
||
|
|
"""
|
||
|
|
client_key = body.get("clientKey") or ""
|
||
|
|
shared_secret = body.get("sharedSecret") or ""
|
||
|
|
base_url = body.get("baseUrl", "") or ""
|
||
|
|
product_type = body.get("productType", "") or ""
|
||
|
|
if not client_key or not shared_secret:
|
||
|
|
return 400, "Missing clientKey/sharedSecret"
|
||
|
|
|
||
|
|
claims = await ac.verify_asymmetric_install_jwt(request, expected_client_key=client_key)
|
||
|
|
if claims is None:
|
||
|
|
webapp.logger.warning("Rejecting Connect install for %s: signature unverified", client_key)
|
||
|
|
return 401, "Install verification failed"
|
||
|
|
|
||
|
|
# Mandatory tenant binding: signed-install proves the caller is *an*
|
||
|
|
# Atlassian tenant, not *ours*, so only installs from an allowlisted
|
||
|
|
# (signature-verified) clientKey are accepted. To bootstrap, add the
|
||
|
|
# clientKey logged here to CONNECT_EXPECTED_CLIENT_KEYS and re-install.
|
||
|
|
if not ac.client_key_allowed(client_key):
|
||
|
|
webapp.logger.warning(
|
||
|
|
"Rejecting Connect install: clientKey %s not in CONNECT_EXPECTED_CLIENT_KEYS",
|
||
|
|
client_key,
|
||
|
|
)
|
||
|
|
return 403, "clientKey not allowed"
|
||
|
|
|
||
|
|
# Defense-in-depth (only enforced when configured): the callback's baseUrl
|
||
|
|
# host must be our Confluence site.
|
||
|
|
if ac.CONNECT_EXPECTED_BASE_URL_HOSTS and not ac.base_url_host_allowed(base_url):
|
||
|
|
webapp.logger.warning("Rejecting Connect install: baseUrl %s not allowed", base_url)
|
||
|
|
return 403, "baseUrl host not allowed"
|
||
|
|
|
||
|
|
existing = await ac.get_installation(client_key)
|
||
|
|
await ac.put_installation(
|
||
|
|
client_key, shared_secret, base_url, product_type, first_install=existing is None
|
||
|
|
)
|
||
|
|
webapp.logger.info(
|
||
|
|
"Connect %s verified and stored for %s",
|
||
|
|
"first-install" if existing is None else "re-install",
|
||
|
|
client_key,
|
||
|
|
)
|
||
|
|
return 204, ""
|
||
|
|
|
||
|
|
|
||
|
|
async def process_uninstall(request: Any, body: dict[str, Any]) -> tuple[int, str]:
|
||
|
|
"""Handle POST /connect/uninstalled. Returns (status_code, detail)."""
|
||
|
|
client_key = body.get("clientKey") or ""
|
||
|
|
if not client_key:
|
||
|
|
return 400, "Missing clientKey"
|
||
|
|
claims = await ac.verify_asymmetric_install_jwt(request, expected_client_key=client_key)
|
||
|
|
if claims is None:
|
||
|
|
webapp.logger.warning("Rejecting Connect uninstall for %s: unverified", client_key)
|
||
|
|
return 401, "Uninstall verification failed"
|
||
|
|
existing = await ac.get_installation(client_key)
|
||
|
|
if existing is None:
|
||
|
|
return 204, "" # idempotent
|
||
|
|
await ac.delete_installation(client_key)
|
||
|
|
webapp.logger.info("Connect uninstall verified for %s", client_key)
|
||
|
|
return 204, ""
|
||
|
|
|
||
|
|
|
||
|
|
def _extract_comment_id(payload: dict[str, Any]) -> str:
|
||
|
|
comment = payload.get("comment")
|
||
|
|
if isinstance(comment, dict) and comment.get("id"):
|
||
|
|
return str(comment["id"])
|
||
|
|
if payload.get("commentId"):
|
||
|
|
return str(payload["commentId"])
|
||
|
|
content = payload.get("content")
|
||
|
|
if isinstance(content, dict) and content.get("id"):
|
||
|
|
return str(content["id"])
|
||
|
|
return ""
|
||
|
|
|
||
|
|
|
||
|
|
async def process_confluence_comment(payload: dict[str, Any], client_key: str = "") -> None:
|
||
|
|
"""Corroborate a comment_created event server-side and dispatch a run."""
|
||
|
|
comment_id = _extract_comment_id(payload)
|
||
|
|
if not comment_id:
|
||
|
|
webapp.logger.debug("Ignoring Confluence webhook: no comment id in payload")
|
||
|
|
return
|
||
|
|
|
||
|
|
server_comment = await webapp.fetch_confluence_comment(comment_id)
|
||
|
|
if not server_comment:
|
||
|
|
webapp.logger.warning(
|
||
|
|
"Rejecting Confluence webhook: comment %s could not be corroborated", comment_id
|
||
|
|
)
|
||
|
|
return
|
||
|
|
|
||
|
|
author = server_comment.get("author") or {}
|
||
|
|
account_id = author.get("account_id") or ""
|
||
|
|
display_name = author.get("name") or ""
|
||
|
|
body_text = server_comment.get("body") or ""
|
||
|
|
page_id = server_comment.get("page_id") or ""
|
||
|
|
space_key = server_comment.get("space_key") or ""
|
||
|
|
|
||
|
|
# Self-trigger loop guard: ignore the app's own comments (its confluence_comment
|
||
|
|
# replies can echo "@openswe" and otherwise re-trigger).
|
||
|
|
if CONFLUENCE_BOT_ACCOUNT_ID and account_id == CONFLUENCE_BOT_ACCOUNT_ID:
|
||
|
|
webapp.logger.debug("Ignoring Confluence webhook: comment authored by the bot account")
|
||
|
|
return
|
||
|
|
for prefix in webapp._GITHUB_BOT_MESSAGE_PREFIXES:
|
||
|
|
if body_text.startswith(prefix):
|
||
|
|
webapp.logger.debug("Ignoring Confluence webhook: comment is our own bot message")
|
||
|
|
return
|
||
|
|
if "@openswe" not in body_text.lower():
|
||
|
|
webapp.logger.debug("Ignoring Confluence webhook: comment doesn't mention @openswe")
|
||
|
|
return
|
||
|
|
|
||
|
|
actor_email = await webapp.get_confluence_user_email(account_id) if account_id else None
|
||
|
|
|
||
|
|
repo_config = webapp.extract_repo_from_text(body_text, default_owner=webapp.DEFAULT_REPO_OWNER)
|
||
|
|
if not repo_config:
|
||
|
|
repo_config = webapp.get_repo_config_from_confluence_mapping(space_key)
|
||
|
|
if not repo_config:
|
||
|
|
repo_config = await webapp.get_team_default_repo()
|
||
|
|
if not repo_config:
|
||
|
|
webapp.logger.info("Ignoring Confluence webhook: no repo resolved for space %s", space_key)
|
||
|
|
return
|
||
|
|
if not webapp._is_repo_allowed(repo_config):
|
||
|
|
webapp.logger.warning(
|
||
|
|
"Rejecting Confluence webhook: repo '%s/%s' not in allowlist",
|
||
|
|
repo_config.get("owner"),
|
||
|
|
repo_config.get("name"),
|
||
|
|
)
|
||
|
|
return
|
||
|
|
|
||
|
|
mapped_login = await webapp.resolve_login_from_email_async(actor_email) if actor_email else None
|
||
|
|
if mapped_login and not webapp.is_login_mapped(mapped_login):
|
||
|
|
webapp.logger.info(
|
||
|
|
"Confluence actor login %s is not an active mapping; running unattributed", mapped_login
|
||
|
|
)
|
||
|
|
mapped_login = None
|
||
|
|
|
||
|
|
thread_id = webapp.generate_thread_id_from_confluence_comment(client_key, comment_id)
|
||
|
|
page = await webapp.fetch_confluence_page(page_id) if page_id else None
|
||
|
|
page_title = (page or {}).get("title", "") or "Confluence page"
|
||
|
|
page_url = (page or {}).get("url", "")
|
||
|
|
|
||
|
|
triggered_by = f"## Triggered by: {display_name}\n\n" if display_name else ""
|
||
|
|
prompt = (
|
||
|
|
f"Please act on the following Confluence comment:\n\n"
|
||
|
|
f"## Repository: {repo_config.get('owner')}/{repo_config.get('name')}\n\n"
|
||
|
|
f"## Confluence page: {page_title} ({space_key}) - Page ID: {page_id}\n\n"
|
||
|
|
f"{triggered_by}"
|
||
|
|
f"## Comment:\n{body_text}\n\n"
|
||
|
|
f"Please analyze this and implement the necessary changes. When you're done, commit and "
|
||
|
|
f"push your changes."
|
||
|
|
)
|
||
|
|
content_blocks: list[dict[str, Any]] = [create_text_block(prompt)]
|
||
|
|
|
||
|
|
configurable: dict[str, Any] = {
|
||
|
|
"repo": repo_config,
|
||
|
|
"confluence": {
|
||
|
|
"comment_id": comment_id,
|
||
|
|
"page_id": page_id,
|
||
|
|
"space_key": space_key,
|
||
|
|
"url": page_url,
|
||
|
|
"triggering_user_name": display_name or "",
|
||
|
|
},
|
||
|
|
"user_email": actor_email,
|
||
|
|
"source": "confluence",
|
||
|
|
}
|
||
|
|
if mapped_login:
|
||
|
|
configurable["github_login"] = mapped_login
|
||
|
|
|
||
|
|
await webapp.upsert_agent_thread_owner_metadata(
|
||
|
|
thread_id,
|
||
|
|
source="confluence",
|
||
|
|
repo_config=repo_config,
|
||
|
|
github_login=mapped_login or "",
|
||
|
|
user_email=actor_email or "",
|
||
|
|
title=page_title,
|
||
|
|
source_context={"confluence": configurable["confluence"]},
|
||
|
|
)
|
||
|
|
|
||
|
|
run = await webapp.dispatch_agent_run(
|
||
|
|
thread_id,
|
||
|
|
content_blocks,
|
||
|
|
configurable,
|
||
|
|
source="confluence",
|
||
|
|
metadata=webapp._AGENT_VERSION_METADATA,
|
||
|
|
)
|
||
|
|
webapp.logger.info(
|
||
|
|
"LangGraph run dispatched for Confluence thread %s (run=%s)",
|
||
|
|
thread_id,
|
||
|
|
run.get("run_id") if isinstance(run, dict) else None,
|
||
|
|
)
|