Skip to content
Merged
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
2 changes: 2 additions & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand All @@ -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]
Expand Down
45 changes: 45 additions & 0 deletions docs/ingestion.md
Original file line number Diff line number Diff line change
@@ -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 `<source>:<source-key>`. 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.
88 changes: 88 additions & 0 deletions src/github_agent_bridge/queue.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
from __future__ import annotations

import hashlib
import json
import sqlite3
from importlib import resources
Expand Down Expand Up @@ -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."""

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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(
Expand All @@ -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:
Expand Down
7 changes: 6 additions & 1 deletion src/github_agent_bridge/reader.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
23 changes: 23 additions & 0 deletions src/github_agent_bridge/sql/schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
7 changes: 6 additions & 1 deletion tests/test_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -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} <notifications@github.com>"),
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} <notifications@github.com>",
),
policy,
)
app = create_app(DashboardConfig(db=db, secret_key="secret", allowed_users={"alice"}, admin_users={"alice"}))
Expand Down
65 changes: 64 additions & 1 deletion tests/test_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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", "<mail@github.com>") == (
"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", "<mail@github.com>") == (
"email:<mail@github.com>"
)


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, "<first@github.com>", BODY1), policy())
duplicate = Notification(
uid=2,
message_id="<second@github.com>",
subject="Re: [gisce/erp] PR",
from_addr="Edu <notifications@github.com>",
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] == [
("<first@github.com>", "accepted", first.id),
("<second@github.com>", "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")
Expand Down
8 changes: 4 additions & 4 deletions tests/test_reader.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Loading