Zheng Feng dcb2e687d5
feat(channels): add GitHub as a webhook-driven channel (#3754)
* feat(channels): add GitHub event-driven agents (#3754)

Add a webhook-driven GitHub channel with fail-closed webhook routing, deterministic per-agent PR/issue threads, mention-gated trigger fan-out, GitHub App token injection for sandboxed gh/git commands, and backend/AGENTS.md documentation.

* fix(llm-middleware): classify bare IndexError as transient

Upstream chat providers occasionally return 200 OK with an empty
generations list (observed against Volces "coding" on
ark.cn-beijing.volces.com). When that happens,
langchain_core.language_models.chat_models.ainvoke raises
``IndexError: list index out of range`` at
``llm_result.generations[0][0].message`` and kills the run.

Treat a bare IndexError reaching the middleware as a transient
upstream-payload glitch and route it through the existing
retry/backoff path instead of failing the whole agent run. The
retry budget and backoff schedule are unchanged.

Adds three regression tests covering the classifier and both the
recover-on-retry and exhausted-retries paths.

* fix(runtime): ignore stale LLM fallback markers from prior runs

When a run on a thread ends with the LLM-error-handling middleware emitting
a `deerflow_error_fallback`-marked AIMessage (e.g. after the IndexError
empty-generations classification fix lands), that message is persisted to
the thread's checkpoint as part of the messages channel. LangGraph replays
the full message history in `stream_mode="values"` chunks, so every
subsequent run on the same thread re-streams the stale fallback marker —
and the worker's chunk scanner faithfully picks it up, flipping
`RunStatus.success` to `RunStatus.error` for runs that themselves had
no LLM failure at all.

Snapshot the set of pre-existing message ids from the pre-run checkpoint
and thread it through `_extract_llm_error_fallback_message` /
`_try_extract_from_message` as a filter. Markers on history messages are
ignored; markers on fresh messages produced during this run still trip
the error path. Falls back to an empty set when the checkpointer is
absent or the snapshot can't be captured, preserving the prior behavior
on first-run / no-state paths.

Adds unit tests for the new filter (helper-level and `_collect_pre_existing_message_ids`)
plus an integration test exercising the full `run_agent` path with a stale
history checkpointer.

* fix(channels): make github channel fire-and-forget to avoid httpx.ReadTimeout on long runs

GitHub agent runs (clone -> edit -> test -> push -> PR) routinely exceed
the langgraph_sdk default 300s read deadline. The manager's runs.wait
call kept an HTTP stream open for the entire run lifetime, so the long
run blew up with httpx.ReadTimeout and the outer except branch then
released the dedupe key and emitted a false 'internal error' outbound.

The GitHub channel's outbound send is log-only by design: agents post to
the issue/PR via the gh CLI in the sandbox when they choose to comment
or create a PR. There is nothing for the manager to ferry back, so the
long-poll was pure overhead.

This change adds ChannelRunPolicy.fire_and_forget (default False) and
sets it True for the github channel. When fire_and_forget is True,
_handle_chat dispatches via client.runs.create (short POST, returns
once the run is pending) instead of client.runs.wait, and skips the
response-extraction + outbound-publish block. ConflictError on a busy
thread still trips the standard THREAD_BUSY_MESSAGE path so behavior on
the busy case is preserved for any future non-github fire-and-forget
channel.

Other (non-github) channels are unchanged: their policy defaults
fire_and_forget=False and they continue to dispatch via runs.wait.

Adds 6 regression tests in tests/test_channels.py::TestGithubFireAndForget:
- Default ChannelRunPolicy.fire_and_forget is False.
- The github policy registers fire_and_forget=True.
- github inbound calls runs.create, not runs.wait, with the right kwargs.
- github inbound publishes no outbound on success.
- ConflictError from runs.create still emits THREAD_BUSY_MESSAGE.
- Non-github channels (slack) still dispatch via runs.wait.

* test(lead-agent): accept user_id kwarg in skill-policy test stubs

The two GitHub-channel tests added in #3754 stubbed
_load_enabled_skills_for_tool_policy with a lambda that only accepted
`available_skills` and `app_config`, but the real function (and its call
site in agent.py) also passes `user_id`. This raised TypeError on every
run, failing backend-unit-tests.

Add `user_id=None` to match the three sibling stubs in the same file.

* refactor(gateway): disambiguate context-key set names

The two frozensets _INTERNAL_ONLY_CONTEXT_KEYS and _CONTEXT_ONLY_KEYS
shared a confusable "CONTEXT_ONLY" token in different orders, and the
first broke the _CONTEXT_<X>_KEYS pattern of its sibling
_CONTEXT_CONFIGURABLE_KEYS. Rename to make the distinct axes explicit:

  _CONTEXT_INTERNAL_CALLER_KEYS  - WHO: internal callers (scheduler) only
  _CONTEXT_RUNTIME_ONLY_KEYS     - WHERE: runtime context only, never configurable

Pure rename, no behavior change.
2026-07-04 22:56:24 +08:00

248 lines
11 KiB
Python

"""Build the GitHub webhook → agent registry.
Walks every user's custom-agent directory under ``{base_dir}/users/`` plus
the legacy shared layout at ``{base_dir}/agents/`` and indexes every agent
that declares a ``github:`` block by the ``(repo, event)`` pairs it
declares an interest in.
The dispatcher calls :func:`build_github_agent_registry` once per webhook
delivery. We avoid re-parsing every ``config.yaml`` on each call via a
small mtime-keyed cache: the directory listing + ``stat()`` per config
file is cheap (~µs), while ``yaml.safe_load`` is the dominant cost
(~hundreds of µs per file). The cache key is the sorted tuple of
``(user_id, agent_name, config.yaml mtime)`` triples; any mtime change,
addition, or deletion invalidates the cache transparently. Operators
who hand-edit ``config.yaml`` see the change on the next webhook.
Cache invalidation caveat: mtime granularity on macOS HFS+ / APFS is
~1 µs but on some filesystems (FAT, network shares with caching) it's
1 s. Two edits inside the same coarse-tick would look identical. For
the dispatch path that's fine — webhooks are rare relative to operator
edits, and the next non-coincident write reconciles. If we ever land
operator tooling that batches sub-second edits, we can extend the
signature with file size or a content hash.
"""
from __future__ import annotations
import logging
import threading
from dataclasses import dataclass
from pathlib import Path
from app.gateway.github.triggers import _resolved_trigger
from deerflow.config.agents_config import (
AgentConfig,
GitHubTriggerConfig,
load_agent_config,
)
from deerflow.config.paths import get_paths
from deerflow.runtime.user_context import DEFAULT_USER_ID
logger = logging.getLogger(__name__)
@dataclass(frozen=True)
class GitHubAgentMatch:
"""One ``(user, agent, _resolved_trigger)`` row in the ``(repo, event)`` index.
The trigger is the binding override merged with per-event field defaults
(see :func:`app.gateway.github.triggers._resolved_trigger`), so the
dispatcher does not have to re-resolve it at fan-out time. Pre-resolving
here also folds the per-binding lookup out of the hot path: the registry
already chose the right binding for this ``(repo, event)``, so the
dispatcher's old "find the binding whose ``.repo`` matches" loop —
which silently dropped events when an agent had multiple bindings on
one repo (PR feedback R3) — disappears entirely. Single-binding-per-repo
is enforced upstream by :class:`GitHubAgentConfig`'s validator, so
each ``(repo, event)`` resolves to exactly one trigger per agent.
The ``github:`` block is read off ``agent.github`` (always non-None
here — the rebuild filters agents without one before constructing a
match), so we don't carry a separate ``github`` field.
"""
user_id: str
agent: AgentConfig
trigger: GitHubTriggerConfig
# Cache: (signature, registry). ``signature`` is a tuple of
# ``(user_id, agent_name, mtime)`` triples. Identical signature → registry
# is still valid, skip the YAML parses.
_Signature = tuple[tuple[str, str, float], ...]
_Registry = dict[tuple[str, str], list[GitHubAgentMatch]]
_cache: tuple[_Signature, _Registry] | None = None
# Threading lock (not asyncio): build_github_agent_registry is invoked
# from asyncio.to_thread in the dispatcher, so the lock is acquired on
# the worker thread. A plain Lock is the right primitive here.
_cache_lock = threading.Lock()
def _discover_user_ids() -> list[str]:
"""Return all user-id directories under ``base_dir/users/``.
Falls back to ``[DEFAULT_USER_ID]`` so the no-auth dev setup (which
keeps everything in ``users/default/``) is always covered even before
the directory has been created on disk.
"""
paths = get_paths()
users_dir: Path = paths.base_dir / "users"
if not users_dir.exists():
return [DEFAULT_USER_ID]
found: list[str] = []
for entry in sorted(users_dir.iterdir()):
if entry.is_dir() and (entry / "agents").exists():
found.append(entry.name)
if DEFAULT_USER_ID not in found:
found.append(DEFAULT_USER_ID)
return found
def _gather_agent_signature() -> tuple[_Signature, list[tuple[str, str]]]:
"""Return (signature, [(user_id, agent_name)]) for every agent on disk.
The signature lets us skip the YAML parse on warm hits; the
discovered list lets the rebuilder process exactly the agents that
the signature covers. Doing iterdir + stat is cheap (~µs each); the
full cost we avoid is the ``yaml.safe_load`` per config.
Includes the legacy shared layout at ``{base_dir}/agents/`` under the
:data:`DEFAULT_USER_ID` bucket so unmigrated installations still
receive webhook fan-out. Per-user entries shadow legacy entries with
the same name (matching :func:`list_custom_agents`' precedence), so
once an install runs ``migrate_user_isolation.py`` the legacy entry
is silently superseded rather than producing duplicate rows.
"""
paths = get_paths()
sig: list[tuple[str, str, float]] = []
discovered: list[tuple[str, str]] = []
seen: set[tuple[str, str]] = set()
for user_id in _discover_user_ids():
agent_root = paths.user_agents_dir(user_id)
if not agent_root.exists():
continue
for entry in sorted(agent_root.iterdir()):
config = entry / "config.yaml"
if not entry.is_dir() or not config.exists():
continue
try:
mtime = config.stat().st_mtime
except OSError:
# Vanished between iterdir and stat — racing operator
# edit. Drop from this round; next call picks it up.
continue
sig.append((user_id, entry.name, mtime))
discovered.append((user_id, entry.name))
seen.add((user_id, entry.name))
# Legacy shared layout: {base_dir}/agents/{name}/. CLAUDE.md commits
# to this as a read-only fallback for unmigrated installs, and
# load_agent_config() / list_custom_agents() honour it — the webhook
# path must too, or an unmigrated install with a ``github:`` block on
# a shared agent silently fans out to nothing. Legacy entries map
# onto the DEFAULT_USER_ID bucket because that is the user-id
# ``load_agent_config(name)`` resolves them under at run-time.
legacy_root = paths.agents_dir
if legacy_root.exists():
for entry in sorted(legacy_root.iterdir()):
config = entry / "config.yaml"
if not entry.is_dir() or not config.exists():
continue
# Per-user shadow: if users/default/agents/{name} already
# exists, skip the legacy entry so we don't index the same
# agent twice with conflicting trigger sets.
if (DEFAULT_USER_ID, entry.name) in seen:
continue
try:
mtime = config.stat().st_mtime
except OSError:
continue
sig.append((DEFAULT_USER_ID, entry.name, mtime))
discovered.append((DEFAULT_USER_ID, entry.name))
seen.add((DEFAULT_USER_ID, entry.name))
return tuple(sig), discovered
def _rebuild(discovered: list[tuple[str, str]]) -> _Registry:
"""Parse every agent's config.yaml and build the (repo, event) index.
Each ``(repo, event)`` slot stores :class:`GitHubAgentMatch` rows — the
user_id + AgentConfig + the trigger already resolved (binding override
merged with per-event field defaults). The dispatcher then only needs
to apply the trigger; it never re-walks ``bindings`` to find the right
one. Single-binding-per-repo is enforced by
:class:`GitHubAgentConfig`'s validator, so a duplicate-repo config
fails to load here (logged as a skip) instead of producing duplicate
rows in this index.
"""
index: _Registry = {}
for user_id, agent_name in discovered:
try:
cfg = load_agent_config(agent_name, user_id=user_id)
except Exception as exc: # noqa: BLE001 — one bad agent must not kill the scan
logger.warning("github_registry: skipping agent %s/%s: %s", user_id, agent_name, exc)
continue
if cfg is None or cfg.github is None:
continue
for binding in cfg.github.bindings:
for event, override in binding.triggers.items():
resolved = _resolved_trigger(event, {event: override})
if resolved is None:
# ``_resolved_trigger`` only returns None when the
# event is not in the dict we passed — by construction
# it is here, so this branch is unreachable. Keep the
# guard for type-checker happiness.
continue
index.setdefault((binding.repo, event), []).append(GitHubAgentMatch(user_id=user_id, agent=cfg, trigger=resolved))
return index
def build_github_agent_registry() -> _Registry:
"""Return ``{(repo, event): [GitHubAgentMatch, ...]}`` across all users.
Each agent appears in the index once per ``(repo, declared_event)`` pair,
with the per-event trigger pre-resolved by merging the binding override
with :data:`app.gateway.github.triggers.DEFAULT_TRIGGERS`. Events are
opt-in per binding: an agent only registers for the events it explicitly
lists under ``github.bindings[].triggers``. An agent that declares an
empty ``triggers:`` map (or omits it) registers for nothing and the
dispatcher will never fan a webhook out to it.
Warm path (no agents added/removed/edited since the last call) costs
only the iterdir + stat pass — no YAML parsing. Cold path parses
every config.yaml and refreshes the cache. The result is shared
across callers (returned by reference) since :class:`GitHubAgentMatch`
is frozen and the registry is intended as read-only.
"""
global _cache
with _cache_lock:
signature, discovered = _gather_agent_signature()
if _cache is not None and _cache[0] == signature:
return _cache[1]
registry = _rebuild(discovered)
_cache = (signature, registry)
return registry
def _invalidate_cache() -> None:
"""Drop the cached registry. Test-only helper."""
global _cache
with _cache_lock:
_cache = None
def lookup_agents(
registry: _Registry,
repo: str,
event: str,
) -> list[GitHubAgentMatch]:
"""Convenience: return the list of agent matches for ``(repo, event)``.
Each match carries the user, AgentConfig (with ``.github`` attached),
and the pre-resolved trigger config for this specific event, so the
caller does not need to walk the agent's ``bindings`` again.
"""
return registry.get((repo, event), [])