mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-07-26 16:07:53 +00:00
* 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.
160 lines
6.6 KiB
Python
160 lines
6.6 KiB
Python
"""Global authentication middleware — fail-closed safety net.
|
|
|
|
Rejects unauthenticated requests to non-public paths with 401. When a
|
|
request passes the cookie check, resolves the JWT payload to a real
|
|
``User`` object and stamps it into both ``request.state.user`` and the
|
|
``deerflow.runtime.user_context`` contextvar so that repository-layer
|
|
owner filtering works automatically via the sentinel pattern.
|
|
|
|
Fine-grained permission checks remain in authz.py decorators.
|
|
"""
|
|
|
|
from collections.abc import Callable
|
|
|
|
from fastapi import HTTPException, Request, Response
|
|
from starlette.middleware.base import BaseHTTPMiddleware
|
|
from starlette.responses import JSONResponse
|
|
from starlette.types import ASGIApp
|
|
|
|
from app.gateway.auth.errors import AuthErrorCode, AuthErrorResponse
|
|
from app.gateway.auth_disabled import (
|
|
AUTH_SOURCE_AUTH_DISABLED,
|
|
AUTH_SOURCE_INTERNAL,
|
|
AUTH_SOURCE_SESSION,
|
|
get_auth_disabled_user,
|
|
is_auth_disabled,
|
|
)
|
|
from app.gateway.authz import _ALL_PERMISSIONS, AuthContext
|
|
from app.gateway.internal_auth import INTERNAL_AUTH_HEADER_NAME, get_internal_user, is_valid_internal_auth_token
|
|
from deerflow.runtime.user_context import reset_current_user, set_current_user
|
|
|
|
# Paths that never require authentication.
|
|
_PUBLIC_PATH_PREFIXES: tuple[str, ...] = (
|
|
"/health",
|
|
"/docs",
|
|
"/redoc",
|
|
"/openapi.json",
|
|
"/api/v1/auth/oauth/",
|
|
"/api/v1/auth/callback/",
|
|
# Inbound webhooks authenticate themselves via provider-specific signatures
|
|
# (e.g. GitHub's X-Hub-Signature-256), not session cookies.
|
|
"/api/webhooks/",
|
|
)
|
|
|
|
# Exact auth paths that are public (login/register/status check).
|
|
# /api/v1/auth/me, /api/v1/auth/change-password etc. are NOT public.
|
|
_PUBLIC_EXACT_PATHS: frozenset[str] = frozenset(
|
|
{
|
|
"/api/v1/auth/login/local",
|
|
"/api/v1/auth/register",
|
|
"/api/v1/auth/logout",
|
|
"/api/v1/auth/setup-status",
|
|
"/api/v1/auth/initialize",
|
|
"/api/v1/auth/providers",
|
|
}
|
|
)
|
|
|
|
|
|
def _is_public(path: str) -> bool:
|
|
stripped = path.rstrip("/")
|
|
if stripped in _PUBLIC_EXACT_PATHS:
|
|
return True
|
|
return any(path.startswith(prefix) for prefix in _PUBLIC_PATH_PREFIXES)
|
|
|
|
|
|
class AuthMiddleware(BaseHTTPMiddleware):
|
|
"""Strict auth gate: reject requests without a valid session.
|
|
|
|
Two-stage check for non-public paths:
|
|
|
|
1. Cookie presence — return 401 NOT_AUTHENTICATED if missing
|
|
2. JWT validation via ``get_optional_user_from_request`` — return 401
|
|
TOKEN_INVALID if the token is absent, malformed, expired, or the
|
|
signed user does not exist / is stale
|
|
|
|
On success, stamps ``request.state.user`` and the
|
|
``deerflow.runtime.user_context`` contextvar so that repository-layer
|
|
owner filters work downstream without every route needing a
|
|
``@require_auth`` decorator. Routes that need per-resource
|
|
authorization (e.g. "user A cannot read user B's thread by guessing
|
|
the URL") should additionally use ``@require_permission(...,
|
|
owner_check=True)`` for explicit enforcement — but authentication
|
|
itself is fully handled here.
|
|
"""
|
|
|
|
def __init__(self, app: ASGIApp) -> None:
|
|
super().__init__(app)
|
|
|
|
async def dispatch(self, request: Request, call_next: Callable) -> Response:
|
|
if _is_public(request.url.path):
|
|
return await call_next(request)
|
|
|
|
internal_user = None
|
|
if is_valid_internal_auth_token(request.headers.get(INTERNAL_AUTH_HEADER_NAME)):
|
|
# Extract the channel owner user ID from the trusted header.
|
|
# When present, the synthetic internal user carries the actual
|
|
# owner identity so that get_effective_user_id() and per-user
|
|
# filesystem paths (custom skills, memory, thread data) resolve
|
|
# to the IM channel user instead of falling back to "default".
|
|
from app.gateway.internal_auth import INTERNAL_OWNER_USER_ID_HEADER_NAME
|
|
|
|
owner_user_id = request.headers.get(INTERNAL_OWNER_USER_ID_HEADER_NAME)
|
|
if owner_user_id:
|
|
owner_user_id = owner_user_id.strip()
|
|
internal_user = get_internal_user(owner_user_id=owner_user_id or None)
|
|
|
|
auth_source = AUTH_SOURCE_SESSION
|
|
access_token = request.cookies.get("access_token")
|
|
|
|
# Non-public path: require session cookie
|
|
if internal_user is not None:
|
|
user = internal_user
|
|
auth_source = AUTH_SOURCE_INTERNAL
|
|
elif access_token:
|
|
# Strict JWT validation: reject junk/expired tokens with 401
|
|
# right here instead of silently passing through. This closes
|
|
# the "junk cookie bypass" gap (AUTH_TEST_PLAN test 7.5.8):
|
|
# without this, non-isolation routes like /api/models would
|
|
# accept any cookie-shaped string as authentication.
|
|
#
|
|
# We call the *strict* resolver so that fine-grained error
|
|
# codes (token_expired, token_invalid, user_not_found, …)
|
|
# propagate from AuthErrorCode, not get flattened into one
|
|
# generic code. BaseHTTPMiddleware doesn't let HTTPException
|
|
# bubble up, so we catch and render it as JSONResponse here.
|
|
from app.gateway.deps import get_current_user_from_request
|
|
|
|
try:
|
|
user = await get_current_user_from_request(request)
|
|
except HTTPException as exc:
|
|
if not is_auth_disabled():
|
|
return JSONResponse(status_code=exc.status_code, content={"detail": exc.detail})
|
|
user = get_auth_disabled_user()
|
|
auth_source = AUTH_SOURCE_AUTH_DISABLED
|
|
elif is_auth_disabled():
|
|
user = get_auth_disabled_user()
|
|
auth_source = AUTH_SOURCE_AUTH_DISABLED
|
|
else:
|
|
return JSONResponse(
|
|
status_code=401,
|
|
content={
|
|
"detail": AuthErrorResponse(
|
|
code=AuthErrorCode.NOT_AUTHENTICATED,
|
|
message="Authentication required",
|
|
).model_dump()
|
|
},
|
|
)
|
|
|
|
# Stamp both request.state.user (for the contextvar pattern)
|
|
# and request.state.auth (so @require_permission's "auth is
|
|
# None" branch short-circuits instead of running the entire
|
|
# JWT-decode + DB-lookup pipeline a second time per request).
|
|
request.state.user = user
|
|
request.state.auth_source = auth_source
|
|
request.state.auth = AuthContext(user=user, permissions=_ALL_PERMISSIONS)
|
|
token = set_current_user(user)
|
|
try:
|
|
return await call_next(request)
|
|
finally:
|
|
reset_current_user(token)
|