diff --git a/src/poller/forwarder.py b/src/poller/forwarder.py index e996b61..7289791 100644 --- a/src/poller/forwarder.py +++ b/src/poller/forwarder.py @@ -17,9 +17,36 @@ def github_delivery_id(event_type: str, *identity: object) -> str: return f"poller-{event_type}-{digest}" -def jira_delivery_id() -> str: - """Build a unique delivery ID for one synthetic Jira event.""" - return f"poller-jira-{uuid4()}" +def jira_delivery_id(payload: dict[str, Any] | None = None) -> str: + """Build a replay-stable ID when a Jira provider identity is available. + + The random fallback is retained for legacy callers that do not provide a + webhook-shaped payload. Provider comment IDs and issue ``updated`` values + are stable across poller restarts and therefore make the forwarded + delivery interchangeable with a native Jira webhook at Forge. + """ + if payload is None: + return f"poller-jira-{uuid4()}" + issue = payload.get("issue", {}) + key = issue.get("key", "") if isinstance(issue, dict) else "" + comment = payload.get("comment", {}) + if isinstance(comment, dict) and comment.get("id") is not None: + identity = ("comment", key, comment["id"]) + else: + fields = issue.get("fields", {}) if isinstance(issue, dict) else {} + updated = fields.get("updated") if isinstance(fields, dict) else None + if isinstance(updated, str) and updated: + identity = ("issue", key, updated) + else: + changelog = payload.get("changelog", {}) + items = changelog.get("items") if isinstance(changelog, dict) else None + if items: + identity = ("changelog", key, items) + else: + return f"poller-jira-{uuid4()}" + raw_identity = "\x1f".join(str(part) for part in identity) + digest = hashlib.sha256(raw_identity.encode()).hexdigest()[:24] + return f"poller-jira-{digest}" async def forward_jira( @@ -27,7 +54,7 @@ async def forward_jira( ) -> None: settings = get_settings() url = f"{settings.forge_gateway_url}/api/v1/webhooks/jira" - delivery_id = delivery_id or jira_delivery_id() + delivery_id = delivery_id or jira_delivery_id(payload) headers = { "Content-Type": "application/json", "X-Atlassian-Webhook-Identifier": delivery_id, diff --git a/src/poller/jira.py b/src/poller/jira.py index e9d1c60..b192e6e 100644 --- a/src/poller/jira.py +++ b/src/poller/jira.py @@ -24,7 +24,7 @@ def __init__(self) -> None: } async def get_issue(self, key: str) -> dict[str, Any]: - fields = "summary,issuetype,status,labels,comment" + fields = "summary,issuetype,status,labels,comment,updated" async with httpx.AsyncClient(headers=self._headers) as client: r = await client.get( f"{self._base}/rest/api/3/issue/{key}", diff --git a/src/poller/models.py b/src/poller/models.py index 2d56484..6e0d014 100644 --- a/src/poller/models.py +++ b/src/poller/models.py @@ -26,6 +26,10 @@ class TicketState: summary: str labels: set[str] last_comment_id: str | None + # Native Jira issue revision from fields.updated. This is delivery + # metadata for generating equivalent webhook-shaped observations, not a + # workflow checkpoint or poll cursor. + issue_updated: str | None = None prs: list[PrState] = field(default_factory=list) # PRD proposals PR tracking prd_pr_repo: str | None = None diff --git a/src/poller/payloads.py b/src/poller/payloads.py index 8fb23ee..7005f13 100644 --- a/src/poller/payloads.py +++ b/src/poller/payloads.py @@ -8,17 +8,23 @@ def label_changed( summary: str, old_labels: set[str], new_labels: set[str], + updated: str | None = None, ) -> dict[str, Any]: + fields: dict[str, Any] = { + "issuetype": {"name": issue_type}, + "status": {"name": status}, + "summary": summary, + "labels": sorted(new_labels), + } + if updated: + # Jira's issue updated timestamp is the native revision available to + # the REST poller for label/status changes. + fields["updated"] = updated return { "webhookEvent": "jira:issue_updated", "issue": { "key": ticket_key, - "fields": { - "issuetype": {"name": issue_type}, - "status": {"name": status}, - "summary": summary, - "labels": sorted(new_labels), - }, + "fields": fields, }, "changelog": { "items": [ @@ -43,26 +49,37 @@ def comment_created( author_account_id: str, author_display_name: str, author_email: str = "", + comment_id: str | int | None = None, + issue_updated: str | None = None, + comment_created: str | None = None, ) -> dict[str, Any]: + fields: dict[str, Any] = { + "issuetype": {"name": issue_type}, + "status": {"name": status}, + "summary": summary, + "labels": sorted(labels), + } + if issue_updated: + fields["updated"] = issue_updated + comment: dict[str, Any] = { + "body": body, + "author": { + "accountId": author_account_id, + "displayName": author_display_name, + "emailAddress": author_email, + }, + } + if comment_id is not None: + comment["id"] = comment_id + if comment_created: + comment["created"] = comment_created return { "webhookEvent": "comment_created", "issue": { "key": ticket_key, - "fields": { - "issuetype": {"name": issue_type}, - "status": {"name": status}, - "summary": summary, - "labels": sorted(labels), - }, - }, - "comment": { - "body": body, - "author": { - "accountId": author_account_id, - "displayName": author_display_name, - "emailAddress": author_email, - }, + "fields": fields, }, + "comment": comment, "user": {"accountId": author_account_id, "displayName": author_display_name}, } @@ -96,19 +113,32 @@ def pr_review_submitted( review_state: str, review_body: str, reviewer_login: str, + review_id: object | None = None, + head_sha: str = "", + base_branch: str = "", + pr_body: str = "", + pr_state: str = "open", + draft: bool = False, ) -> dict[str, Any]: + review: dict[str, Any] = { + "state": review_state.lower(), + "body": review_body, + "user": {"login": reviewer_login}, + } + if review_id is not None: + review["id"] = review_id return { "action": "submitted", - "review": { - "state": review_state.lower(), - "body": review_body, - }, + "review": review, "pull_request": { "number": pr_number, "title": pr_title, - "state": "open", + "body": pr_body, + "state": pr_state, "html_url": pr_url, - "head": {"ref": branch}, + "head": {"ref": branch, "sha": head_sha}, + "base": {"ref": base_branch}, + "draft": draft, }, "repository": {"full_name": repo}, "sender": {"login": reviewer_login}, @@ -120,13 +150,28 @@ def issue_comment( pr_number: int, comment_body: str, sender_login: str, + comment_id: object | None = None, + issue_title: str = "", + issue_body: str = "", + issue_state: str = "open", + issue_url: str = "", ) -> dict[str, Any]: + comment: dict[str, Any] = {"body": comment_body} + if comment_id is not None: + comment["id"] = comment_id + issue: dict[str, Any] = { + "number": pr_number, + "title": issue_title, + "body": issue_body, + "state": issue_state, + "html_url": issue_url, + } + if issue_url: + issue["pull_request"] = {"html_url": issue_url} return { "action": "created", - "issue": {"number": pr_number}, - "comment": { - "body": comment_body, - }, + "issue": issue, + "comment": comment, "repository": {"full_name": repo}, "sender": {"login": sender_login}, } @@ -138,6 +183,9 @@ def pr_merged( pr_number: int, pr_title: str, pr_url: str, + head_sha: str = "", + pr_body: str = "", + base_branch: str = "", ) -> dict[str, Any]: return { "action": "closed", @@ -145,9 +193,11 @@ def pr_merged( "number": pr_number, "merged": True, "title": pr_title, + "body": pr_body, "state": "closed", "html_url": pr_url, - "head": {"ref": branch}, + "head": {"ref": branch, "sha": head_sha}, + "base": {"ref": base_branch}, }, "repository": {"full_name": repo}, "sender": {"login": "poller"}, diff --git a/src/poller/watcher.py b/src/poller/watcher.py index f308a1f..ce2d209 100644 --- a/src/poller/watcher.py +++ b/src/poller/watcher.py @@ -63,6 +63,7 @@ async def add(self, ticket_key: str) -> None: summary=state.summary, old_labels=state.labels - {"forge:managed"}, new_labels=state.labels, + updated=state.issue_updated, ) ) async with self._lock: @@ -255,6 +256,7 @@ async def _snapshot(self, ticket_key: str) -> TicketState: labels = set(fields.get("labels", [])) comments = fields.get("comment", {}).get("comments", []) last_comment_id = comments[-1]["id"] if comments else None + issue_updated = fields.get("updated") issue_type = fields.get("issuetype", {}).get("name", "") status = fields.get("status", {}).get("name", "") @@ -304,6 +306,7 @@ async def _snapshot(self, ticket_key: str) -> TicketState: summary=summary, labels=labels, last_comment_id=last_comment_id, + issue_updated=issue_updated if isinstance(issue_updated, str) else None, prs=prs, ) @@ -400,6 +403,9 @@ def _result() -> dict: pr_number=prd_pr_number, pr_title=pr_data.get("title", ""), pr_url=pr_data.get("html_url", ""), + head_sha=pr_data.get("head", {}).get("sha", ""), + pr_body=pr_data.get("body", "") or "", + base_branch=pr_data.get("base", {}).get("ref", ""), ), event_type="pull_request", delivery_id=forwarder.github_delivery_id( @@ -431,6 +437,12 @@ def _result() -> dict: review_state=rev.get("state", ""), review_body=rev.get("body", "") or "", reviewer_login=reviewer, + review_id=rev.get("id"), + head_sha=pr_data.get("head", {}).get("sha", ""), + base_branch=pr_data.get("base", {}).get("ref", ""), + pr_body=pr_data.get("body", "") or "", + pr_state=pr_data.get("state", "open"), + draft=bool(pr_data.get("draft", False)), ), event_type="pull_request_review", delivery_id=forwarder.github_delivery_id( @@ -460,6 +472,11 @@ def _result() -> dict: pr_number=prd_pr_number, comment_body=c.get("body", ""), sender_login=sender, + comment_id=c.get("id"), + issue_title=prd_title, + issue_body=pr_data.get("body", "") or "", + issue_state=pr_data.get("state", "open"), + issue_url=prd_url, ), event_type="issue_comment", delivery_id=forwarder.github_delivery_id( @@ -548,6 +565,9 @@ def _result() -> dict: pr_number=spec_pr_number, pr_title=pr_data.get("title", ""), pr_url=pr_data.get("html_url", ""), + head_sha=pr_data.get("head", {}).get("sha", ""), + pr_body=pr_data.get("body", "") or "", + base_branch=pr_data.get("base", {}).get("ref", ""), ), event_type="pull_request", delivery_id=forwarder.github_delivery_id( @@ -578,6 +598,12 @@ def _result() -> dict: review_state=rev.get("state", ""), review_body=rev.get("body", "") or "", reviewer_login=reviewer, + review_id=rev.get("id"), + head_sha=pr_data.get("head", {}).get("sha", ""), + base_branch=pr_data.get("base", {}).get("ref", ""), + pr_body=pr_data.get("body", "") or "", + pr_state=pr_data.get("state", "open"), + draft=bool(pr_data.get("draft", False)), ), event_type="pull_request_review", delivery_id=forwarder.github_delivery_id( @@ -606,6 +632,11 @@ def _result() -> dict: pr_number=spec_pr_number, comment_body=c.get("body", ""), sender_login=sender, + comment_id=c.get("id"), + issue_title=spec_title, + issue_body=pr_data.get("body", "") or "", + issue_state=pr_data.get("state", "open"), + issue_url=spec_url, ), event_type="issue_comment", delivery_id=forwarder.github_delivery_id( @@ -651,6 +682,9 @@ async def _poll(self, ticket_key: str) -> None: summary=new_summary, old_labels=state.labels, new_labels=new_labels, + updated=fields.get("updated") + if isinstance(fields.get("updated"), str) + else None, ) ) @@ -675,6 +709,13 @@ async def _poll(self, ticket_key: str) -> None: author_account_id=author.get("accountId", ""), author_display_name=author.get("displayName", ""), author_email=author.get("emailAddress", ""), + comment_id=comment.get("id"), + issue_updated=fields.get("updated") + if isinstance(fields.get("updated"), str) + else None, + comment_created=comment.get("created") + if isinstance(comment.get("created"), str) + else None, ) ) @@ -728,6 +769,8 @@ async def _poll(self, ticket_key: str) -> None: pr_number=pr.pr_number, pr_title=pr.pr_title or "", pr_url=pr.pr_url or "", + head_sha=pr.head_sha or "", + base_branch=pr_data.get("base", {}).get("ref", ""), ), event_type="pull_request", delivery_id=forwarder.github_delivery_id( @@ -813,6 +856,12 @@ async def _poll(self, ticket_key: str) -> None: review_state=rev.get("state", ""), review_body=rev.get("body", "") or "", reviewer_login=reviewer, + review_id=rev.get("id"), + head_sha=pr_data.get("head", {}).get("sha", ""), + base_branch=pr_data.get("base", {}).get("ref", ""), + pr_body=pr_data.get("body", "") or "", + pr_state=pr_data.get("state", "open"), + draft=bool(pr_data.get("draft", False)), ), event_type="pull_request_review", delivery_id=forwarder.github_delivery_id( @@ -846,6 +895,11 @@ async def _poll(self, ticket_key: str) -> None: pr_number=pr.pr_number, comment_body=body, sender_login=sender, + comment_id=c.get("id"), + issue_title=pr.pr_title or "", + issue_body=pr_data.get("body", "") or "", + issue_state=pr_data.get("state", "open"), + issue_url=pr.pr_url or "", ), event_type="issue_comment", delivery_id=forwarder.github_delivery_id( @@ -877,6 +931,9 @@ async def _poll(self, ticket_key: str) -> None: summary=new_summary, labels=new_labels, last_comment_id=new_last_comment_id, + issue_updated=fields.get("updated") + if isinstance(fields.get("updated"), str) + else state.issue_updated, prs=prs, prd_pr_repo=prd_updates.get("prd_pr_repo", state.prd_pr_repo), prd_pr_number=prd_updates.get("prd_pr_number", state.prd_pr_number), diff --git a/tests/test_observation_contract.py b/tests/test_observation_contract.py new file mode 100644 index 0000000..cc1ae37 --- /dev/null +++ b/tests/test_observation_contract.py @@ -0,0 +1,99 @@ +"""Checks for the provider identity fields Forge normalizes from poller payloads.""" + +from poller import forwarder, payloads + + +def test_synthetic_github_review_preserves_provider_review_id() -> None: + payload = payloads.pr_review_submitted( + repo="acme/api", + branch="feature", + pr_number=42, + pr_title="Change", + pr_url="https://github.com/acme/api/pull/42", + review_state="APPROVED", + review_body="Looks good", + reviewer_login="octocat", + review_id=1234, + head_sha="abc123", + base_branch="main", + ) + + assert payload["review"]["id"] == 1234 + assert payload["review"]["user"]["login"] == "octocat" + assert payload["pull_request"]["head"]["sha"] == "abc123" + assert payload["pull_request"]["base"]["ref"] == "main" + assert forwarder.github_delivery_id("pull_request_review", "acme/api", 42, 1234) + + +def test_synthetic_github_comment_preserves_provider_comment_id() -> None: + payload = payloads.issue_comment( + repo="acme/api", + pr_number=42, + comment_body="Please update this", + sender_login="octocat", + comment_id=5678, + issue_title="Change", + issue_url="https://github.com/acme/api/pull/42", + ) + + assert payload["comment"]["id"] == 5678 + assert payload["issue"]["pull_request"]["html_url"].endswith("/42") + + +def test_synthetic_github_merge_preserves_provider_head_revision() -> None: + payload = payloads.pr_merged( + repo="acme/api", + branch="feature", + pr_number=42, + pr_title="Change", + pr_url="https://github.com/acme/api/pull/42", + head_sha="abc123", + base_branch="main", + ) + + assert payload["pull_request"]["head"]["sha"] == "abc123" + assert payload["pull_request"]["base"]["ref"] == "main" + + +def test_provider_delivery_ids_are_replay_stable_but_event_specific() -> None: + first = forwarder.github_delivery_id("issue_comment", "acme/api", 42, 5678) + replay = forwarder.github_delivery_id("issue_comment", "acme/api", 42, 5678) + other = forwarder.github_delivery_id("issue_comment", "acme/api", 42, 5679) + + assert first == replay + assert first != other + + +def test_synthetic_jira_comment_preserves_provider_comment_identity() -> None: + payload = payloads.comment_created( + ticket_key="FORGE-17", + issue_type="Bug", + status="In Progress", + summary="A bug", + labels={"forge:managed"}, + body="Please investigate", + author_account_id="alice", + author_display_name="Alice", + comment_id="10042", + issue_updated="2026-08-27T10:01:00.000+0000", + comment_created="2026-08-27T10:00:59.000+0000", + ) + + assert payload["comment"]["id"] == "10042" + assert payload["issue"]["fields"]["updated"].startswith("2026-08-27") + assert forwarder.jira_delivery_id(payload) == forwarder.jira_delivery_id(payload) + + +def test_synthetic_jira_label_revision_is_replay_stable() -> None: + payload = payloads.label_changed( + ticket_key="FORGE-17", + issue_type="Bug", + status="In Progress", + summary="A bug", + old_labels={"forge:managed"}, + new_labels={"forge:managed", "forge:approved"}, + updated="2026-08-27T10:02:00.000+0000", + ) + + assert payload["issue"]["fields"]["updated"] == "2026-08-27T10:02:00.000+0000" + assert forwarder.jira_delivery_id(payload) == forwarder.jira_delivery_id(payload)