mirror of
https://github.com/Sea-Haven-Industries/open-swe.git
synced 2026-09-30 10:23:14 +00:00
Plan step C4 (docs/upstream-sync/domain-reorg/reorg-build-plan.md, approved
decisions 1-2): split the 2,590-line agent/webapp.py monolith into
agent/webhooks/common.py (shared verify/dispatch helpers), agent/api/app.py
(composition), agent/api/health.py (/health + /webhooks/run-complete), and
per-source {github,linear,slack,jira,confluence}_routes.py. Atlassian
Connect lifecycle + descriptor routes (/connect/*) fold into
confluence_routes.py; webapp.py becomes the upstream-shaped compatibility
shim (from .api.app import app). langgraph.json http.app stays
agent.webapp:app via the shim.
Fork content, upstream layout: linear/slack route files verified
content-identical to upstream 8356eb34 and taken verbatim; github_routes is
upstream + the fork's CI auto-fix trigger wiring; jira/confluence routes are
fork-only, transformed to the same common.X / service.X module-attribute
style. All signature verification (GitHub HMAC, Slack, Linear
timestamp-freshness, verify_jira_secret + opt-in HMAC/timestamp/IP
allowlist, Connect JWT/qsh), token-attribution gating, TID-COLLIDE-01 repo
binding, _is_repo_auto_review_enabled gates, and public-repo org gate move
unchanged.
Handlers rewired from webapp.X to common.X; test monkeypatch sites across
26 files + conftest.py + e2e/harness.py retargeted to
webhook_common/handler/route modules per upstream's pattern. Residual
agent.webapp importers: only the shim, langgraph.json http.app, Makefile
uvicorn target, and docs (doc-path updates land in C7).
Gates: ruff check + format, pytest --co, full unit (1637 passed), full
Playwright E2E vs real langgraph dev (9/9), residual-importer sweep.
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.utils import atlassian_connect as ac
|
|
|
|
from . import common
|
|
|
|
# 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:
|
|
common.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):
|
|
common.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):
|
|
common.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
|
|
)
|
|
common.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:
|
|
common.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)
|
|
common.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:
|
|
common.logger.debug("Ignoring Confluence webhook: no comment id in payload")
|
|
return
|
|
|
|
server_comment = await common.fetch_confluence_comment(comment_id)
|
|
if not server_comment:
|
|
common.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:
|
|
common.logger.debug("Ignoring Confluence webhook: comment authored by the bot account")
|
|
return
|
|
for prefix in common._GITHUB_BOT_MESSAGE_PREFIXES:
|
|
if body_text.startswith(prefix):
|
|
common.logger.debug("Ignoring Confluence webhook: comment is our own bot message")
|
|
return
|
|
if "@openswe" not in body_text.lower():
|
|
common.logger.debug("Ignoring Confluence webhook: comment doesn't mention @openswe")
|
|
return
|
|
|
|
actor_email = await common.get_confluence_user_email(account_id) if account_id else None
|
|
|
|
repo_config = common.extract_repo_from_text(body_text, default_owner=common.DEFAULT_REPO_OWNER)
|
|
if not repo_config:
|
|
repo_config = common.get_repo_config_from_confluence_mapping(space_key)
|
|
if not repo_config:
|
|
repo_config = await common.get_team_default_repo()
|
|
if not repo_config:
|
|
common.logger.info("Ignoring Confluence webhook: no repo resolved for space %s", space_key)
|
|
return
|
|
if not common._is_repo_allowed(repo_config):
|
|
common.logger.warning(
|
|
"Rejecting Confluence webhook: repo '%s/%s' not in allowlist",
|
|
repo_config.get("owner"),
|
|
repo_config.get("name"),
|
|
)
|
|
return
|
|
|
|
mapped_login = await common.resolve_login_from_email_async(actor_email) if actor_email else None
|
|
if mapped_login and not common.is_login_mapped(mapped_login):
|
|
common.logger.info(
|
|
"Confluence actor login %s is not an active mapping; running unattributed", mapped_login
|
|
)
|
|
mapped_login = None
|
|
|
|
thread_id = common.generate_thread_id_from_confluence_comment(client_key, comment_id)
|
|
page = await common.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 common.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 common.dispatch_agent_run(
|
|
thread_id,
|
|
content_blocks,
|
|
configurable,
|
|
source="confluence",
|
|
metadata=common._AGENT_VERSION_METADATA,
|
|
)
|
|
common.logger.info(
|
|
"LangGraph run dispatched for Confluence thread %s (run=%s)",
|
|
thread_id,
|
|
run.get("run_id") if isinstance(run, dict) else None,
|
|
)
|