Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
63843a9
feat: connectors — one-click OAuth integrations, per-user accounts, t…
nihalashetty Aug 14, 2026
b1332d1
fix(connectors): Google id arguments accept a pasted link, not a gues…
nihalashetty Aug 14, 2026
32c0da2
fix(tools): interpolated lists/dicts render as JSON, and Sheets write…
nihalashetty Aug 14, 2026
58159aa
fix(connectors): audit every REST action against vendor docs
nihalashetty Aug 14, 2026
b550e24
fix: address code-review findings on the connectors branch
nihalashetty Aug 14, 2026
a77d0eb
fix: second review pass — escaping regression, cache race, and six sm…
nihalashetty Aug 14, 2026
0e52b55
fix: third review pass - MCP credential pooling, connection lifecycle…
nihalashetty Aug 14, 2026
3d7ba53
fix(connectors): a losing concurrent install no longer blanks the win…
nihalashetty Aug 16, 2026
1e424c8
fix(connectors): Google connectors ask for offline access, so they su…
nihalashetty Aug 16, 2026
d05cfb5
refactor(oauth): one mechanism for extra authorize params, not two
nihalashetty Aug 16, 2026
45af818
fix(mcp): retired transports close on a timer, and never under the bu…
nihalashetty Aug 16, 2026
4a78b73
fix(connectors): uninstall gives back the egress hosts the install added
nihalashetty Aug 16, 2026
717b024
perf(connectors): the installed list resolves every token bundle in o…
nihalashetty Aug 16, 2026
5d32f86
fix(connectors): every route reports the same, live action count
nihalashetty Aug 16, 2026
b7827b6
docs(connectors): strike the two decisions the branch overruled
nihalashetty Aug 16, 2026
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
26 changes: 26 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,32 @@ FORGE_EGRESS_ALLOW_PRIVATE_HOSTS=[] # e.g. ["localhost","127.0.0.1"] (Dock
# resolves to a different host per deploy (dev/qa/prod). A template referencing a key NOT in this
# map fails the call loudly. Blank/unset = {}.
FORGE_TOOL_VARS= # e.g. {"api_base":"https://api.example.com"}
# --- Connectors: the vendor OAuth apps this deployment signs people in with ---
# This is the ONLY place a bundled connector's credentials live. Register each vendor app once,
# put it here, restart - and from then on every user just clicks Connect, signs in with their own
# account on the vendor's page, and approves. Nobody types a client secret into Forge's UI, and
# a connector whose group is missing here says so instead of showing a form.
#
# Keyed by CREDENTIAL GROUP, which is the vendor, not the connector: one "google" entry covers
# Gmail, Calendar, Drive and Sheets; one "microsoft" entry covers Outlook.
# google Gmail · Calendar · Drive · Sheets microsoft Outlook
# github GitHub hubspot HubSpot
# airtable Airtable
# Slack, Notion, Linear and Atlassian need NO entry: they publish OAuth metadata and Forge
# registers a client with them automatically (RFC 7591) the first time someone connects.
#
# Redirect/callback URL to register with every vendor (must match exactly):
# <FORGE_PUBLIC_BASE_URL>/v1/oauth/callback dev: http://localhost:8000/v1/oauth/callback
#
# Google: create a "Web application" OAuth client and paste the callback above into its
# Authorized redirect URIs. It must be a Web application client, not a Desktop one: Desktop
# clients have no redirect-URI field at all and accept only loopback addresses, so Forge's
# callback can never be registered against one and every sign-in fails with
# redirect_uri_mismatch. Note the "Authorized JavaScript origins" field strips the path - the
# full URL belongs in "Authorized redirect URIs". Until the app passes Google verification it
# is capped at 100 users and shows an "unverified app" warning; Gmail read/modify are
# RESTRICTED scopes needing a CASA audit.
FORGE_CONNECTOR_OAUTH_APPS= # {"google":{"client_id":"…","client_secret":"…"},"github":{…}}
# Deployment-wide fallback for a per-user auth provider's `token_ctx_key`: the run-context key an
# integration forwards its per-user token under (via X-Forge-Context) when a provider doesn't set
# its own. Empty = off. e.g. user_token
Expand Down
115 changes: 115 additions & 0 deletions apps/api/forge/auth_providers/oauth_flow.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
"""Shared authorization-code connect flow (PKCE + signed state).

Both the Auth Providers screen and the Connectors screen start the same 3-legged OAuth dance
against the same `/v1/oauth/callback`. Keeping the URL construction here means the security
properties - PKCE S256, a signed short-lived state, and carrying only the provider's declared
per-user dims - are implemented once instead of drifting between two call sites.
"""

from __future__ import annotations

import base64
import hashlib
import secrets as _secrets
from urllib.parse import urlencode

from forge.config import settings
from forge.models import AuthProvider
from forge.secrets.store import SecretStore
from forge.security import create_state_token


class OAuthNotConfigured(ValueError):
"""The provider is missing something the authorize step needs (client id, authorize URL)."""


def redirect_uri(cfg: dict) -> str:
return cfg.get("redirect_uri") or f"{settings.public_base_url.rstrip('/')}/v1/oauth/callback"


def token_request(cfg: dict, data: dict) -> tuple[dict, dict]:
"""Form body + headers for a token-endpoint POST, given a body that already carries
`client_id`/`client_secret`. Returns `(data, headers)`.

Two things a naive POST gets wrong against real servers:

* `Accept: application/json`. GitHub's token endpoint answers `application/x-www-form-
urlencoded` without it, so a perfectly successful exchange then dies in `.json()`.
* client_secret_basic. RFC 6749 §2.3.1 makes HTTP Basic the method a server MUST support
and form-post the one it MAY; a few (Airtable) accept only Basic. `token_auth: "basic"`
on the provider config moves the secret into the Authorization header. `client_id` stays
in the body, which every server tolerates and some still require.
"""
headers = {"Accept": "application/json"}
if cfg.get("token_auth") != "basic":
return {k: v for k, v in data.items() if v is not None}, headers
client_id, client_secret = data.get("client_id") or "", data.get("client_secret") or ""
creds = base64.b64encode(f"{client_id}:{client_secret}".encode()).decode()
headers["Authorization"] = "Basic " + creds
return {k: v for k, v in data.items() if v is not None and k != "client_secret"}, headers


async def build_authorize_url(
ap: AuthProvider, *, tenant_id: str, project_id: str, context: dict | None = None,
secrets: SecretStore | None = None,
) -> str:
"""The provider's authorize URL, carrying a signed state that binds the callback back to
this tenant/project/provider (and, for a per-user provider, to the end user being connected).

The PKCE verifier rides inside the SIGNED state. That is safe for a CONFIDENTIAL client -
which is what Forge always registers, since it holds the client secret server-side in its
own encrypted store - because the token exchange also requires that secret. A public client
(no secret) would need the verifier held server-side instead; see routers/oauth.py.
"""
cfg = ap.config or {}
store = secrets or SecretStore()
client_id = None
if cfg.get("client_id_ref"):
try:
client_id = await store.read_ref(tenant_id=tenant_id, project_id=project_id, ref=cfg["client_id_ref"])
except Exception as e: # noqa: BLE001 - surface as a configuration problem, not a 500
raise OAuthNotConfigured("client_id secret is not set") from e
if not client_id:
raise OAuthNotConfigured("client_id secret is not set")
if not cfg.get("authorize_url"):
raise OAuthNotConfigured("authorize_url is not configured")

verifier = _secrets.token_urlsafe(64)
challenge = base64.urlsafe_b64encode(hashlib.sha256(verifier.encode()).digest()).decode().rstrip("=")
claims = {"tid": tenant_id, "pid": project_id, "ap": ap.id, "cv": verifier}
per_user = cfg.get("per_user_context_keys") or []
ctx = context or {}
user_ctx = {k: ctx[k] for k in per_user if k in ctx}
if user_ctx:
claims["ctx"] = user_ctx

q = {
"response_type": "code",
"client_id": str(client_id),
"redirect_uri": redirect_uri(cfg),
"state": create_state_token(claims),
"code_challenge": challenge,
"code_challenge_method": "S256",
}
if cfg.get("scope"):
q["scope"] = cfg["scope"]
# RFC 8707: bind the issued token to the resource it is for, so a token minted for one MCP
# server can't be replayed against another. Only sent when the provider names a resource
# (connectors set it to the MCP server URL); omitted otherwise to avoid upsetting servers
# that reject unknown parameters.
if cfg.get("resource"):
q["resource"] = cfg["resource"]
# Vendor-specific authorize parameters - Google's `access_type`/`prompt` for a refresh token,
# Slack's `user_scope`, HubSpot's optional scopes. `authorize_params` is the ONE mechanism;
# there is deliberately no top-level special case for individual keys. There used to be one
# for `access_type`/`prompt`, and because nothing ever wrote them at the top level it was
# dead on the connector path - which is precisely how four Google manifests came to ship
# without asking for offline access while a block of code appeared to be handling it.
#
# Applied last but never over the protocol parameters above: a manifest must not be able to
# redirect the callback or weaken PKCE by declaring `redirect_uri` or `code_challenge_method`
# as an "extra".
for key, value in (cfg.get("authorize_params") or {}).items():
if key not in q:
q[key] = str(value)
return f"{cfg['authorize_url']}?{urlencode(q)}"
7 changes: 6 additions & 1 deletion apps/api/forge/auth_providers/resolver.py
Original file line number Diff line number Diff line change
Expand Up @@ -274,9 +274,14 @@ async def _refresh_oauth(self, provider, cfg: dict, read, bundle: dict, tenant_i
"client_secret": await read(cfg.get("client_secret_ref")),
}
client = client or shared_async_client()
# Same client-auth method and Accept header the initial exchange used - a provider that
# only speaks client_secret_basic rejects the refresh just as readily as the exchange.
from forge.auth_providers.oauth_flow import token_request

form, headers = token_request(cfg, data)
r = await guarded_request(
client, "POST", cfg["token_url"],
data={k: v for k, v in data.items() if v is not None}, timeout=30, follow_redirects=True,
data=form, headers=headers, timeout=30, follow_redirects=True,
)
r.raise_for_status()
body = r.json()
Expand Down
105 changes: 88 additions & 17 deletions apps/api/forge/auth_providers/templates.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@

from __future__ import annotations

import json
import re
from collections.abc import Collection
from typing import Any
Expand Down Expand Up @@ -42,55 +43,125 @@ def _lookup(path: str, vars: dict, strict_ns: Collection[str] = ()) -> Any:
return cur


def _sub_one(mm: re.Match, vars: dict, strict_ns: Collection[str] = ()) -> str:
def _sub_one(mm: re.Match, vars: dict, strict_ns: Collection[str] = (),
escape_json: bool = False) -> str:
# Embedded token (not a whole-string match): stringify the resolved value. Only a missing
# value (None) becomes empty - a falsy-but-real value like 0 or False must render as "0"/
# "False", not "" (an `x or ""` here would silently drop legitimate zeros/booleans).
v = _lookup(mm.group(1), vars, strict_ns)
return "" if v is None else str(v)


def render_template(s: str, vars: dict, *, strict_ns: Collection[str] = ()) -> Any:
if v is None:
return ""
# Inside a JSON body template the token sits between quotes, so the CONTENT has to be
# JSON-escaped: a newline, a double quote or a backslash in model-written text (a note, an
# email body, a search string) otherwise terminates the string early and the whole body stops
# being JSON. dumps()[1:-1] escapes the content without adding a second pair of quotes.
if escape_json and isinstance(v, str):
return json.dumps(v, ensure_ascii=False)[1:-1]
# A list/dict interpolated into a body template is being placed into JSON - `str()` there
# yields Python's repr (single quotes, True/None), which is not JSON, so the body fails to
# parse and gets sent as raw text. `[[1, 2]]` happens to be valid JSON and `[['a', 'b']]` is
# not, which is why this only ever showed up on string data. Scalars keep str(): a bare
# `{{input.flag}}` in a query string renders "False", which is what that context wants.
if isinstance(v, (list, dict)):
return json.dumps(v, ensure_ascii=False)
return str(v)


def render_template(s: str, vars: dict, *, strict_ns: Collection[str] = (),
escape_json: bool = False) -> Any:
# `strict_ns` names namespaces whose missing keys raise MissingTemplateVar instead of
# rendering empty (used for {{env.*}}). Defaults to lenient for all namespaces.
# Whole-string single token -> preserve native type (numbers, objects).
# `escape_json` JSON-escapes substituted strings, for a template that is a JSON document
# (the REST body path opts in). Whole-string single token -> preserve native type (numbers,
# objects), which needs no escaping because the value is never re-parsed from text.
m = _TOKEN.fullmatch(s.strip())
if m:
return _lookup(m.group(1), vars, strict_ns)
return _TOKEN.sub(lambda mm: _sub_one(mm, vars, strict_ns), s)
return _TOKEN.sub(lambda mm: _sub_one(mm, vars, strict_ns, escape_json), s)


#: Directives that need STRUCTURAL rendering (parse the JSON, then walk it) rather than plain
#: string substitution, because they produce a value the surrounding text can't express.
DIRECTIVES = ("$each", "$mime")

def has_each_directive(obj: Any) -> bool:
"""True if `obj` (a parsed JSON structure) contains a `$each` loop directive anywhere - i.e.
a dict that has "$each" as a KEY. Used to decide whether a body template needs structural

def has_structural_directive(obj: Any) -> bool:
"""True if `obj` (a parsed JSON structure) contains a `$each` or `$mime` directive anywhere -
i.e. a dict that has one as a KEY. Used to decide whether a body template needs structural
rendering; a literal "$each" appearing inside a string value is NOT a directive and must not
trigger it (that would silently change type coercion for unrelated templates)."""
if isinstance(obj, dict):
if "$each" in obj:
if any(d in obj for d in DIRECTIVES):
return True
return any(has_each_directive(v) for v in obj.values())
return any(has_structural_directive(v) for v in obj.values())
if isinstance(obj, list):
return any(has_each_directive(v) for v in obj)
return any(has_structural_directive(v) for v in obj)
return False


def render_value(obj: Any, vars: dict, *, allow_each: bool = False, strict_ns: Collection[str] = ()) -> Any:
"""Walk a parsed JSON structure, rendering `{{token}}` leaves. `$each` loop directives are
honored ONLY when `allow_each=True` (the REST body-template path opts in); every other caller
- auth token_fetch/extract rules, data-node payloads - passes the default False, so a literal
object key named "$each" stays an ordinary key instead of being reinterpreted as a loop.
"""Walk a parsed JSON structure, rendering `{{token}}` leaves. Structural directives (`$each`,
`$mime`) are honored ONLY when `allow_each=True` (the REST body-template path opts in); every
other caller - auth token_fetch/extract rules, data-node payloads - passes the default False,
so a literal object key named "$each" stays an ordinary key instead of being reinterpreted.
`strict_ns` propagates the fail-loud namespaces (e.g. env) to every leaf."""
if isinstance(obj, str):
return render_template(obj, vars, strict_ns=strict_ns)
if isinstance(obj, dict):
if allow_each and "$each" in obj:
return _render_each(obj, vars, strict_ns=strict_ns)
if allow_each and "$mime" in obj:
return _render_mime(obj["$mime"], vars, strict_ns=strict_ns)
return {k: render_value(v, vars, allow_each=allow_each, strict_ns=strict_ns) for k, v in obj.items()}
if isinstance(obj, list):
return [render_value(v, vars, allow_each=allow_each, strict_ns=strict_ns) for v in obj]
return obj


#: Header name -> the key a `$mime` spec uses for it.
_MIME_HEADERS = (
("To", "to"), ("Cc", "cc"), ("Bcc", "bcc"), ("From", "from"), ("Reply-To", "reply_to"),
("Subject", "subject"), ("In-Reply-To", "in_reply_to"), ("References", "references"),
)


def _render_mime(spec: Any, vars: dict, *, strict_ns: Collection[str] = ()) -> str:
"""Build an RFC 2822 message from `{to, subject, text, ...}` and return it base64url-encoded.

This exists because Gmail's send endpoint accepts ONLY a base64url-encoded MIME message.
Declaring that as a tool argument means asking a language model to base64-encode by hand -
which it cannot do reliably - and the malformed result comes back as an opaque HTTP 400 with
nothing in it to debug. So the model supplies the fields a person would type and the encoding
happens here, where stdlib `email` gets header encoding, non-ASCII subjects, address lists
and CRLF line endings right.

Address fields accept a string or a list. Supply `text` (and optionally `html` for a
multipart/alternative). Empty/absent headers are omitted rather than sent blank.
"""
import base64
from email.message import EmailMessage
from email.policy import SMTP

if not isinstance(spec, dict):
raise ValueError("$mime expects an object, e.g. {\"to\": \"…\", \"subject\": \"…\", \"text\": \"…\"}")
resolved = {k: render_value(v, vars, strict_ns=strict_ns) for k, v in spec.items()}

msg = EmailMessage()
for header, key in _MIME_HEADERS:
val = resolved.get(key)
if isinstance(val, (list, tuple)):
val = ", ".join(str(v).strip() for v in val if str(v).strip())
if val is None or str(val).strip() == "":
continue
msg[header] = str(val).strip()
msg.set_content(str(resolved.get("text") or resolved.get("body") or ""))
html = resolved.get("html")
if html:
msg.add_alternative(str(html), subtype="html")
# SMTP policy => CRLF line endings, as RFC 2822 requires.
return base64.urlsafe_b64encode(msg.as_bytes(policy=SMTP)).decode()


def _render_each(directive: dict, vars: dict, *, strict_ns: Collection[str] = ()) -> list:
"""Expand a `{"$each": "{{input.rows}}", "$as": "row", "$do": {...}}` loop directive into a
list: render `$do` once per item of the array `$each` resolves to, with the item bound under
Expand Down
Loading
Loading