diff --git a/docs/README.md b/docs/README.md index bdbfc2c..aa7c5db 100644 --- a/docs/README.md +++ b/docs/README.md @@ -9,6 +9,7 @@ A compact map of the `github-agent-bridge` documentation set. | Understand what this project is | [`../README.md`](../README.md) | Overview | | Install a deployment | [`installation.md`](installation.md) | How-to | | Understand system design | [`architecture.md`](architecture.md) | Explanation | +| Understand event identity and transport migration | [`ingestion.md`](ingestion.md) | Explanation | | Configure policy | [`policy-reference.md`](policy-reference.md) | Reference | | Roll out safely | [`shadow-canary.md`](shadow-canary.md) | How-to | | Operate production | [`operations.md`](operations.md) | How-to | @@ -26,6 +27,7 @@ A compact map of the `github-agent-bridge` documentation set. flowchart TD A[Scope] --> B[Architecture] B --> I[Installation] + B --> J[Event ingestion] I --> C[Policy] C --> D[Shadow/canary rollout] D --> E[Operations] diff --git a/docs/ingestion.md b/docs/ingestion.md new file mode 100644 index 0000000..fa87c2c --- /dev/null +++ b/docs/ingestion.md @@ -0,0 +1,45 @@ +# Event ingestion and transport migration + +The bridge separates three identities: + +- A **receipt** identifies one delivery from one transport. Email uses + `Message-ID`; a future webhook transport will use `X-GitHub-Delivery`. +- A **GitHub event** identifies the underlying action independently of its + transport. It is derived only from immutable GitHub object IDs exposed by + the notification, such as a comment, review, or workflow-run ID. +- A **work key** (`owner/repo#number`) serializes work in one thread. It must + not deduplicate distinct events in that thread. + +`ingest_receipts` makes retries from a transport idempotent. `github_events` +implements first-writer-wins across transports and points every duplicate +receipt at the winning job. Both records and the job are written in one SQLite +transaction, before the IMAP high-water mark advances. + +When a notification lacks enough immutable data to prove event identity, the +event key falls back to `:`. The bridge deliberately +accepts a possible duplicate rather than risk dropping a legitimate action. + +## Transport comparison + +| Concern | IMAP | GitHub App webhook | +| --- | --- | --- | +| Trust | Auth headers and GitHub sender checks | Mandatory HMAC-SHA256 signature over raw bytes, plus installation/repository policy | +| Idempotency | `Message-ID` receipt | `X-GitHub-Delivery` receipt | +| Cross-source identity | IDs parsed from GitHub URLs | IDs in the structured payload | +| Latency | Polling and mail delivery | Immediate delivery with retries | +| Operations | Mailbox cursor, credentials, formatting drift | Public TLS endpoint, secret rotation, delivery monitoring | + +## Gradual rollout + +1. **Phase 0 (implemented):** route IMAP through the common transactional + ingestor while preserving the existing queue and dispatch behavior. +2. **Shadow webhook:** verify signatures and persist receipts/events, but do + not create jobs. Compare coverage and canonical keys with IMAP. +3. **Canary dual ingest:** allow webhook enqueue only for `enabledRepos`. + The unique event key guarantees that the first source wins. +4. **Webhook primary:** keep IMAP as a delayed fallback until a complete + operational cycle has no unexplained IMAP-only actionable events. + +Webhook ingestion must not be enabled until raw-body signature verification, +payload-size limits, secret rotation, recovery of persisted-but-unprocessed +receipts, and metrics for source divergence are in place. diff --git a/src/github_agent_bridge/queue.py b/src/github_agent_bridge/queue.py index 091428f..52453cd 100644 --- a/src/github_agent_bridge/queue.py +++ b/src/github_agent_bridge/queue.py @@ -1,5 +1,6 @@ from __future__ import annotations +import hashlib import json import sqlite3 from importlib import resources @@ -41,6 +42,33 @@ def semantic_event_identity( return (ctx.work_key, action, ctx.target_kind, target_id, (trigger_actor or "").lower()) +def canonical_event_key( + action: str, + ctx: GitHubContext, + source: str, + source_key: str, +) -> str: + """Identify a GitHub event across transports when immutable IDs prove identity. + + A source-specific fallback is deliberately used when an email does not expose + enough immutable GitHub data. False negatives are safer than merging distinct + user actions. + """ + repo = (ctx.repo or "").lower() + identities = ( + ("issue_comment", ctx.comment_id), + ("pull_request_review_comment", ctx.review_comment_id), + ("pull_request_review", ctx.review_id), + ("commit_comment", ctx.commit_comment_id), + ) + for event_type, target_id in identities: + if repo and target_id: + return f"{event_type}:created:{repo}:{target_id}" + if repo and ctx.workflow_run_id: + return f"workflow_run:{action}:{repo}:{ctx.workflow_run_id}" + return f"{source}:{source_key}" + + class ClosingConnection(sqlite3.Connection): """Commit or roll back a context-managed connection, then close it.""" @@ -80,6 +108,18 @@ def init(self) -> None: self._ensure_indexes(con) def enqueue(self, n: Notification, policy: Policy) -> tuple[Job | None, str]: + """Backward-compatible email enqueue entrypoint.""" + return self.ingest(n, policy, source="email", source_key=n.message_id) + + def ingest( + self, + n: Notification, + policy: Policy, + *, + source: str = "email", + source_key: str | None = None, + ) -> tuple[Job | None, str]: + source_key = source_key or n.message_id ctx = extract_github_context(n.body) action = classify_github_action( n.subject, @@ -126,11 +166,43 @@ def enqueue(self, n: Notification, policy: Policy) -> tuple[Job | None, str]: status = {"auto": "done", "ask": "waiting_approval", "deny": "denied"}.get(decision, "pending") now = utc_now() trigger_actor = trigger_actor_details_for_enqueue(n, ctx) + event_key = canonical_event_key(action, ctx, source, source_key) + payload_hash = hashlib.sha256(n.body.encode("utf-8")).hexdigest() if trigger_actor and trigger_actor.user_id: metadata["trigger_actor_id"] = trigger_actor.user_id with self.connect() as con: con.execute("BEGIN IMMEDIATE") try: + try: + con.execute( + "INSERT INTO ingest_receipts(source,source_key,payload_hash,event_key,status,created_at,updated_at) VALUES(?,?,?,?,?,?,?)", + (source, source_key, payload_hash, event_key, "received", now, now), + ) + receipt_id = int(con.execute("SELECT last_insert_rowid()").fetchone()[0]) + except sqlite3.IntegrityError: + receipt = con.execute( + "SELECT job_id FROM ingest_receipts WHERE source=? AND source_key=?", + (source, source_key), + ).fetchone() + con.commit() + job = self.get(int(receipt["job_id"])) if receipt and receipt["job_id"] else None + return job, "duplicate" + event = con.execute( + "SELECT job_id FROM github_events WHERE event_key=?", + (event_key,), + ).fetchone() + if event is not None: + con.execute( + "UPDATE ingest_receipts SET status='duplicate',job_id=?,updated_at=? WHERE id=?", + (event["job_id"], now, receipt_id), + ) + con.commit() + job = self.get(int(event["job_id"])) if event["job_id"] else None + return job, "duplicate" + con.execute( + "INSERT INTO github_events(event_key,first_source,context_json,created_at,updated_at) VALUES(?,?,?,?,?)", + (event_key, source, ctx.to_json(), now, now), + ) existing = con.execute( f"SELECT * FROM jobs WHERE work_key=? AND status IN ({','.join('?' for _ in COALESCE_STATUSES)}) ORDER BY id LIMIT 1", (ctx.work_key, *COALESCE_STATUSES), @@ -169,6 +241,14 @@ def enqueue(self, n: Notification, policy: Policy) -> tuple[Job | None, str]: else: con.execute("UPDATE jobs SET coalesced_count=coalesced_count+1, uid=?, message_id=message_id, subject=?, context_json=?, updated_at=? WHERE id=?", (n.uid, n.subject, ctx.to_json(), now, existing["id"])) self._log(con, existing["id"], ctx.work_key, "coalesced", "Notification coalesced into active job", n.message_id) + con.execute( + "UPDATE github_events SET job_id=?,updated_at=? WHERE event_key=?", + (existing["id"], now, event_key), + ) + con.execute( + "UPDATE ingest_receipts SET status='accepted',job_id=?,updated_at=? WHERE id=?", + (existing["id"], now, receipt_id), + ) con.commit() if policy.feedback_learning.enabled and existing["message_id"] != n.message_id: feedback.capture_feedback( @@ -187,6 +267,14 @@ def enqueue(self, n: Notification, policy: Policy) -> tuple[Job | None, str]: (ctx.work_key, ctx.repo, ctx.issue_number, status, action, decision, intent, n.subject, n.message_id, n.uid, trigger_actor.login if trigger_actor else None, trigger_actor.avatar_url if trigger_actor else None, ctx.to_json(), json.dumps(metadata), now, now), ) job_id = int(con.execute("SELECT last_insert_rowid()").fetchone()[0]) + con.execute( + "UPDATE github_events SET job_id=?,updated_at=? WHERE event_key=?", + (job_id, now, event_key), + ) + con.execute( + "UPDATE ingest_receipts SET status='accepted',job_id=?,updated_at=? WHERE id=?", + (job_id, now, receipt_id), + ) self._log(con, job_id, ctx.work_key, "queued" if status == "pending" else status, f"decision={decision} action={action}", n.message_id) con.commit() if policy.feedback_learning.enabled: diff --git a/src/github_agent_bridge/reader.py b/src/github_agent_bridge/reader.py index 7aab967..5723740 100644 --- a/src/github_agent_bridge/reader.py +++ b/src/github_agent_bridge/reader.py @@ -70,7 +70,12 @@ def _fetch_once(self) -> int: if is_github_notification_message(msg, from_addr): n = Notification(uid=uid, message_id=message_id, subject=subject, from_addr=from_addr, body=extract_body_text(msg), auth=parse_auth_results(msg)) try: - self.queue.enqueue(n, self.policy) + self.queue.ingest( + n, + self.policy, + source="email", + source_key=n.message_id, + ) except sqlite3.Error: raise except Exception as exc: diff --git a/src/github_agent_bridge/sql/schema.sql b/src/github_agent_bridge/sql/schema.sql index b934417..c7894b0 100644 --- a/src/github_agent_bridge/sql/schema.sql +++ b/src/github_agent_bridge/sql/schema.sql @@ -39,6 +39,29 @@ CREATE INDEX IF NOT EXISTS idx_jobs_dashboard_order ON jobs( COALESCE(finished_at, started_at, updated_at, created_at) DESC, id DESC ); +CREATE TABLE IF NOT EXISTS ingest_receipts ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + source TEXT NOT NULL, + source_key TEXT NOT NULL, + payload_hash TEXT NOT NULL, + event_key TEXT, + job_id INTEGER REFERENCES jobs(id) ON DELETE SET NULL, + status TEXT NOT NULL CHECK(status IN ('received','accepted','duplicate')), + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + UNIQUE(source, source_key) +); +CREATE INDEX IF NOT EXISTS idx_ingest_receipts_event_key ON ingest_receipts(event_key); +CREATE INDEX IF NOT EXISTS idx_ingest_receipts_job_id ON ingest_receipts(job_id); +CREATE TABLE IF NOT EXISTS github_events ( + event_key TEXT PRIMARY KEY, + job_id INTEGER REFERENCES jobs(id) ON DELETE SET NULL, + first_source TEXT NOT NULL, + context_json TEXT NOT NULL, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_github_events_job_id ON github_events(job_id); CREATE TABLE IF NOT EXISTS coalesced_notifications ( id INTEGER PRIMARY KEY AUTOINCREMENT, job_id INTEGER NOT NULL REFERENCES jobs(id) ON DELETE CASCADE, diff --git a/tests/test_backend.py b/tests/test_backend.py index 27cf0c6..241112d 100644 --- a/tests/test_backend.py +++ b/tests/test_backend.py @@ -1319,7 +1319,12 @@ def test_dashboard_mcp_users_include_all_job_actors(tmp_path): policy = Policy(trusted_orgs=["pilipilisbot"]) for index in range(205): q.enqueue( - notif(uid=index + 1, mid=f"<{index + 1}@github.com>", from_addr=f"user{index:03d} "), + notif( + uid=index + 1, + mid=f"<{index + 1}@github.com>", + body=f"@pilipilisbot https://github.com/gisce/erp/pull/1#issuecomment-{index + 10}", + from_addr=f"user{index:03d} ", + ), policy, ) app = create_app(DashboardConfig(db=db, secret_key="secret", allowed_users={"alice"}, admin_users={"alice"})) diff --git a/tests/test_queue.py b/tests/test_queue.py index 29162fb..8c538ed 100644 --- a/tests/test_queue.py +++ b/tests/test_queue.py @@ -2,7 +2,8 @@ import pytest -from github_agent_bridge.models import Notification +from github_agent_bridge.models import GitHubContext, Notification +from github_agent_bridge.queue import canonical_event_key from github_agent_bridge.intent_classifier import IntentClassification from github_agent_bridge.policy import FeedbackLearning, IntentClassifier, Policy from github_agent_bridge.queue import JobQueue @@ -118,6 +119,68 @@ def test_enqueue_and_coalesce_same_work_key(tmp_path, monkeypatch): assert job1.trigger_actor_avatar_url == "https://github.com/Edu.png?size=80" +def test_canonical_event_key_uses_immutable_comment_id_across_sources(): + ctx = GitHubContext( + urls=["https://github.com/gisce/erp/issues/42#issuecomment-123"], + repo="gisce/erp", + issue_number=42, + comment_id=123, + target_kind="issue", + ) + + assert canonical_event_key("reply_comment", ctx, "email", "") == ( + "issue_comment:created:gisce/erp:123" + ) + assert canonical_event_key("reply_comment", ctx, "webhook", "delivery-1") == ( + "issue_comment:created:gisce/erp:123" + ) + + +def test_canonical_event_key_falls_back_to_source_receipt_when_identity_is_uncertain(): + ctx = GitHubContext( + urls=["https://github.com/gisce/erp/issues/42"], + repo="gisce/erp", + issue_number=42, + target_kind="issue", + ) + + assert canonical_event_key("mention", ctx, "email", "") == ( + "email:" + ) + + +def test_ingest_records_receipt_and_event_and_deduplicates_same_event(tmp_path, monkeypatch): + monkeypatch.setattr("github_agent_bridge.actors.github_actor_details_for_context", lambda ctx, *, gh_bin="gh": None) + q = JobQueue(tmp_path / "q.sqlite3") + + first, first_state = q.ingest(notif(1, "", BODY1), policy()) + duplicate = Notification( + uid=2, + message_id="", + subject="Re: [gisce/erp] PR", + from_addr="Edu ", + body=BODY1, + auth={"spf": True, "dkim": True, "dmarc": True}, + ) + second, second_state = q.ingest(duplicate, policy()) + + assert first_state == "enqueued" + assert second_state == "duplicate" + assert second.id == first.id + with q.connect() as con: + receipts = con.execute( + "SELECT source_key,status,job_id FROM ingest_receipts ORDER BY id" + ).fetchall() + events = con.execute("SELECT event_key,job_id FROM github_events").fetchall() + assert [(row["source_key"], row["status"], row["job_id"]) for row in receipts] == [ + ("", "accepted", first.id), + ("", "duplicate", first.id), + ] + assert len(events) == 1 + assert events[0]["event_key"] == "issue_comment:created:gisce/erp:10" + assert events[0]["job_id"] == first.id + + def test_equivalent_open_issue_notification_coalesces_after_claim(tmp_path, monkeypatch): monkeypatch.setattr("github_agent_bridge.actors.github_actor_details_for_context", lambda ctx, *, gh_bin="gh": None) q = JobQueue(tmp_path / "q.sqlite3") diff --git a/tests/test_reader.py b/tests/test_reader.py index 9603029..b61b3d9 100644 --- a/tests/test_reader.py +++ b/tests/test_reader.py @@ -161,17 +161,17 @@ def test_fetch_once_retries_transient_enqueue_storage_failure(monkeypatch, tmp_p "@pilipilisbot https://github.com/gisce/erp/issues/42#issuecomment-99", ), }) - original_enqueue = queue.enqueue + original_ingest = queue.ingest attempts = {"count": 0} - def fail_once_then_enqueue(notification, policy): + def fail_once_then_ingest(notification, policy, **kwargs): attempts["count"] += 1 if attempts["count"] == 1: raise sqlite3.OperationalError("database is locked") - return original_enqueue(notification, policy) + return original_ingest(notification, policy, **kwargs) monkeypatch.setattr(imaplib, "IMAP4_SSL", lambda *args: mailbox) - monkeypatch.setattr(queue, "enqueue", fail_once_then_enqueue) + monkeypatch.setattr(queue, "ingest", fail_once_then_ingest) monkeypatch.setattr("github_agent_bridge.actors.github_actor_details_for_context", lambda ctx, *, gh_bin="gh": None) reader = ImapReader(