mirror of
https://github.com/Sea-Haven-Industries/open-swe.git
synced 2026-10-01 09:43:14 +00:00
Add Python agent with LangSmith sandbox support
- Add apps/agent/ with server.py, webapp.py, encryption.py - Configure to use deepagents branch with LangSmith sandbox - Add langgraph.json for deployment - Add pyproject.toml with dependencies
This commit is contained in:
parent
f154609351
commit
b2c466dfb7
17 changed files with 5132 additions and 0 deletions
57
apps/agent/Makefile
Normal file
57
apps/agent/Makefile
Normal file
|
|
@ -0,0 +1,57 @@
|
|||
.PHONY: all format lint test tests integration_tests help run dev
|
||||
|
||||
# Default target executed when no arguments are given to make.
|
||||
all: help
|
||||
|
||||
######################
|
||||
# DEVELOPMENT
|
||||
######################
|
||||
|
||||
dev:
|
||||
langgraph dev
|
||||
|
||||
run:
|
||||
uvicorn agent.webapp:app --reload --port 8000
|
||||
|
||||
install:
|
||||
uv pip install -e .
|
||||
|
||||
######################
|
||||
# TESTING
|
||||
######################
|
||||
|
||||
TEST_FILE ?= tests/
|
||||
|
||||
test tests:
|
||||
uv run pytest -vvv $(TEST_FILE)
|
||||
|
||||
integration_tests:
|
||||
uv run pytest -vvv tests/integration_tests/
|
||||
|
||||
######################
|
||||
# LINTING AND FORMATTING
|
||||
######################
|
||||
|
||||
PYTHON_FILES=.
|
||||
|
||||
lint:
|
||||
uv run ruff check $(PYTHON_FILES)
|
||||
uv run ruff format $(PYTHON_FILES) --diff
|
||||
|
||||
format:
|
||||
uv run ruff format $(PYTHON_FILES)
|
||||
uv run ruff check --fix $(PYTHON_FILES)
|
||||
|
||||
######################
|
||||
# HELP
|
||||
######################
|
||||
|
||||
help:
|
||||
@echo '----'
|
||||
@echo 'dev - run LangGraph dev server'
|
||||
@echo 'run - run webhook server'
|
||||
@echo 'install - install dependencies'
|
||||
@echo 'format - run code formatters'
|
||||
@echo 'lint - run linters'
|
||||
@echo 'test - run unit tests'
|
||||
@echo 'integration_tests - run integration tests'
|
||||
9
apps/agent/README.md
Normal file
9
apps/agent/README.md
Normal file
|
|
@ -0,0 +1,9 @@
|
|||
# Open SWE Agent
|
||||
|
||||
Python-based software engineering agent that automates issue resolution and PR creation.
|
||||
|
||||
## Installation
|
||||
|
||||
```bash
|
||||
pip install -e .
|
||||
```
|
||||
BIN
apps/agent/__pycache__/encryption.cpython-311.pyc
Normal file
BIN
apps/agent/__pycache__/encryption.cpython-311.pyc
Normal file
Binary file not shown.
BIN
apps/agent/__pycache__/server.cpython-311.pyc
Normal file
BIN
apps/agent/__pycache__/server.cpython-311.pyc
Normal file
Binary file not shown.
BIN
apps/agent/__pycache__/server.cpython-314.pyc
Normal file
BIN
apps/agent/__pycache__/server.cpython-314.pyc
Normal file
Binary file not shown.
BIN
apps/agent/__pycache__/webapp.cpython-311.pyc
Normal file
BIN
apps/agent/__pycache__/webapp.cpython-311.pyc
Normal file
Binary file not shown.
3
apps/agent/agent/__init__.py
Normal file
3
apps/agent/agent/__init__.py
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
"""Open SWE Agent - Main Python agent application."""
|
||||
|
||||
__version__ = "0.1.0"
|
||||
BIN
apps/agent/agent/__pycache__/__init__.cpython-311.pyc
Normal file
BIN
apps/agent/agent/__pycache__/__init__.cpython-311.pyc
Normal file
Binary file not shown.
BIN
apps/agent/agent/__pycache__/encryption.cpython-311.pyc
Normal file
BIN
apps/agent/agent/__pycache__/encryption.cpython-311.pyc
Normal file
Binary file not shown.
BIN
apps/agent/agent/__pycache__/server.cpython-311.pyc
Normal file
BIN
apps/agent/agent/__pycache__/server.cpython-311.pyc
Normal file
Binary file not shown.
BIN
apps/agent/agent/__pycache__/webapp.cpython-311.pyc
Normal file
BIN
apps/agent/agent/__pycache__/webapp.cpython-311.pyc
Normal file
Binary file not shown.
74
apps/agent/agent/encryption.py
Normal file
74
apps/agent/agent/encryption.py
Normal file
|
|
@ -0,0 +1,74 @@
|
|||
"""Encryption utilities for sensitive data like tokens."""
|
||||
|
||||
import logging
|
||||
import os
|
||||
|
||||
from cryptography.fernet import Fernet, InvalidToken
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class EncryptionKeyMissingError(ValueError):
|
||||
"""Raised when TOKEN_ENCRYPTION_KEY environment variable is not set."""
|
||||
|
||||
|
||||
def _get_encryption_key() -> bytes:
|
||||
"""Get or derive the encryption key from environment variable.
|
||||
|
||||
Uses TOKEN_ENCRYPTION_KEY env var if set (must be 32 url-safe base64 bytes),
|
||||
otherwise derives a key from LANGSMITH_API_KEY using SHA256.
|
||||
|
||||
Returns:
|
||||
32-byte Fernet-compatible key
|
||||
|
||||
Raises:
|
||||
EncryptionKeyMissingError: If TOKEN_ENCRYPTION_KEY is not set
|
||||
"""
|
||||
explicit_key = os.environ.get("TOKEN_ENCRYPTION_KEY")
|
||||
if not explicit_key:
|
||||
raise EncryptionKeyMissingError
|
||||
|
||||
return explicit_key.encode()
|
||||
|
||||
|
||||
def encrypt_token(token: str) -> str:
|
||||
"""Encrypt a token for safe storage.
|
||||
|
||||
Args:
|
||||
token: The plaintext token to encrypt
|
||||
|
||||
Returns:
|
||||
Base64-encoded encrypted token
|
||||
"""
|
||||
if not token:
|
||||
return ""
|
||||
|
||||
key = _get_encryption_key()
|
||||
f = Fernet(key)
|
||||
encrypted = f.encrypt(token.encode())
|
||||
return encrypted.decode()
|
||||
|
||||
|
||||
def decrypt_token(encrypted_token: str) -> str:
|
||||
"""Decrypt an encrypted token.
|
||||
|
||||
Args:
|
||||
encrypted_token: The base64-encoded encrypted token
|
||||
|
||||
Returns:
|
||||
The plaintext token, or empty string if decryption fails
|
||||
"""
|
||||
if not encrypted_token:
|
||||
return ""
|
||||
|
||||
try:
|
||||
key = _get_encryption_key()
|
||||
f = Fernet(key)
|
||||
decrypted = f.decrypt(encrypted_token.encode())
|
||||
return decrypted.decode()
|
||||
except InvalidToken:
|
||||
logger.warning("Failed to decrypt token: invalid token")
|
||||
return ""
|
||||
except EncryptionKeyMissingError:
|
||||
logger.warning("Failed to decrypt token: encryption key not set")
|
||||
return ""
|
||||
937
apps/agent/agent/server.py
Normal file
937
apps/agent/agent/server.py
Normal file
|
|
@ -0,0 +1,937 @@
|
|||
"""Main entry point and CLI loop for Open SWE agent."""
|
||||
# ruff: noqa: E402
|
||||
|
||||
# Suppress deprecation warnings from langchain_core (e.g., Pydantic V1 on Python 3.14+)
|
||||
# ruff: noqa: E402
|
||||
import logging
|
||||
import os
|
||||
import warnings
|
||||
from collections.abc import Sequence
|
||||
from typing import Any
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
from langchain.agents.middleware import AgentState, after_agent, after_model, before_model
|
||||
from langchain.agents.middleware.types import AgentMiddleware
|
||||
from langchain.tools import BaseTool
|
||||
from langchain_core.language_models import BaseChatModel
|
||||
from langgraph.config import get_config, get_store
|
||||
from langgraph.graph.state import RunnableConfig
|
||||
from langgraph.pregel import Pregel
|
||||
from langgraph.runtime import Runtime
|
||||
|
||||
warnings.filterwarnings("ignore", module="langchain_core._api.deprecation")
|
||||
|
||||
import asyncio
|
||||
|
||||
# Suppress Pydantic v1 compatibility warnings from langchain on Python 3.14+
|
||||
warnings.filterwarnings("ignore", message=".*Pydantic V1.*", category=UserWarning)
|
||||
|
||||
# Now safe to import agent (which imports LangChain modules)
|
||||
# Async wrapper for create_sandbox
|
||||
from contextlib import asynccontextmanager
|
||||
|
||||
from deepagents import create_deep_agent
|
||||
from deepagents.backends.sandbox import SandboxBackendProtocol
|
||||
from deepagents_cli.agent import get_system_prompt
|
||||
from deepagents_cli.config import config, settings
|
||||
from deepagents_cli.integrations.sandbox_factory import create_sandbox
|
||||
from deepagents_cli.tools import fetch_url, http_request, web_search
|
||||
|
||||
# Local import for encryption
|
||||
from .encryption import decrypt_token
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def create_sandbox_async(provider: str, **kwargs):
|
||||
"""Async wrapper around create_sandbox for compatibility."""
|
||||
with create_sandbox(provider, **kwargs) as sandbox:
|
||||
yield sandbox
|
||||
|
||||
|
||||
def create_server_agent(
|
||||
model: str | BaseChatModel | None,
|
||||
assistant_id: str,
|
||||
*,
|
||||
tools: list[BaseTool] | None = None,
|
||||
sandbox: SandboxBackendProtocol | None = None,
|
||||
sandbox_type: str | None = None,
|
||||
system_prompt: str | None = None,
|
||||
auto_approve: bool = True, # noqa: ARG001 - Always True for Open SWE
|
||||
working_dir: str | None = None,
|
||||
middleware: Sequence[AgentMiddleware] = (),
|
||||
) -> Pregel:
|
||||
"""Create a server-mode agent for Open SWE.
|
||||
|
||||
This creates an agent configured for server/cloud deployment with sandbox
|
||||
support and custom middleware. Always runs with auto_approve=True.
|
||||
|
||||
Args:
|
||||
model: LLM model to use. Can be None for introspection-only mode.
|
||||
assistant_id: Agent identifier for memory/state storage
|
||||
tools: Additional tools to provide to agent
|
||||
sandbox: Optional sandbox backend for remote execution (e.g., LangSmithBackend).
|
||||
sandbox_type: Type of sandbox provider ("langsmith").
|
||||
Used for system prompt generation.
|
||||
system_prompt: Override the default system prompt. If None, generates one
|
||||
based on sandbox_type and assistant_id.
|
||||
working_dir: Override the default working directory (e.g., cloned repo path).
|
||||
Used in system prompt to tell the agent where to operate.
|
||||
middleware: Sequence of middleware to apply to the agent.
|
||||
|
||||
Returns:
|
||||
Configured LangGraph Pregel instance ready for execution
|
||||
"""
|
||||
agent_tools = tools or []
|
||||
|
||||
# Get or use custom system prompt
|
||||
if system_prompt is None:
|
||||
if sandbox_type is not None:
|
||||
system_prompt = get_system_prompt(
|
||||
assistant_id=assistant_id,
|
||||
sandbox_type=sandbox_type,
|
||||
working_dir=working_dir,
|
||||
)
|
||||
else:
|
||||
# Only happens when thread_id is None / not actually running
|
||||
system_prompt = ""
|
||||
|
||||
return create_deep_agent(
|
||||
model=model,
|
||||
system_prompt=system_prompt,
|
||||
tools=agent_tools,
|
||||
backend=sandbox,
|
||||
middleware=middleware,
|
||||
interrupt_on={}, # Always auto-approve for Open SWE
|
||||
).with_config(config)
|
||||
|
||||
|
||||
tools = [http_request, fetch_url]
|
||||
if settings.has_tavily:
|
||||
tools.append(web_search)
|
||||
|
||||
from langgraph_sdk import get_client
|
||||
|
||||
client = get_client()
|
||||
|
||||
SANDBOX_CREATING = "__creating__"
|
||||
SANDBOX_CREATION_TIMEOUT = 180
|
||||
SANDBOX_POLL_INTERVAL = 1.0
|
||||
|
||||
# HTTP status codes
|
||||
HTTP_CREATED = 201
|
||||
HTTP_UNPROCESSABLE_ENTITY = 422
|
||||
|
||||
# Message count thresholds
|
||||
MIN_MESSAGES_FOR_PREV_CHECK = 2
|
||||
|
||||
_SANDBOX_BACKENDS: dict[str, Any] = {}
|
||||
|
||||
import httpx
|
||||
|
||||
LINEAR_API_KEY = os.environ.get("LINEAR_API_KEY", "")
|
||||
|
||||
|
||||
async def create_github_pr(
|
||||
repo_owner: str,
|
||||
repo_name: str,
|
||||
github_token: str,
|
||||
title: str,
|
||||
head_branch: str,
|
||||
base_branch: str,
|
||||
body: str,
|
||||
) -> tuple[str | None, int | None]:
|
||||
"""Create a GitHub pull request via the API.
|
||||
|
||||
Args:
|
||||
repo_owner: Repository owner (e.g., "langchain-ai")
|
||||
repo_name: Repository name (e.g., "deepagents")
|
||||
github_token: GitHub access token
|
||||
title: PR title
|
||||
head_branch: Source branch name
|
||||
base_branch: Target branch name
|
||||
body: PR description
|
||||
|
||||
Returns:
|
||||
Tuple of (pr_url, pr_number) if successful, (None, None) otherwise
|
||||
"""
|
||||
pr_payload = {
|
||||
"title": title,
|
||||
"head": head_branch,
|
||||
"base": base_branch,
|
||||
"body": body,
|
||||
}
|
||||
|
||||
logger.info(
|
||||
"Creating PR: head=%s, base=%s, repo=%s/%s",
|
||||
head_branch,
|
||||
base_branch,
|
||||
repo_owner,
|
||||
repo_name,
|
||||
)
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient() as http_client:
|
||||
pr_response = await http_client.post(
|
||||
f"https://api.github.com/repos/{repo_owner}/{repo_name}/pulls",
|
||||
headers={
|
||||
"Authorization": f"Bearer {github_token}",
|
||||
"Accept": "application/vnd.github+json",
|
||||
"X-GitHub-Api-Version": "2022-11-28",
|
||||
},
|
||||
json=pr_payload,
|
||||
)
|
||||
|
||||
pr_data = pr_response.json()
|
||||
|
||||
if pr_response.status_code == HTTP_CREATED:
|
||||
pr_url = pr_data.get("html_url")
|
||||
pr_number = pr_data.get("number")
|
||||
logger.info("PR created successfully: %s", pr_url)
|
||||
return pr_url, pr_number
|
||||
|
||||
if pr_response.status_code == HTTP_UNPROCESSABLE_ENTITY:
|
||||
logger.error("GitHub API validation error (422): %s", pr_data.get("message"))
|
||||
else:
|
||||
logger.error(
|
||||
"GitHub API error (%s): %s",
|
||||
pr_response.status_code,
|
||||
pr_data.get("message"),
|
||||
)
|
||||
|
||||
if "errors" in pr_data:
|
||||
logger.error("GitHub API errors detail: %s", pr_data.get("errors"))
|
||||
|
||||
return None, None
|
||||
|
||||
except httpx.HTTPError:
|
||||
logger.exception("Failed to create PR via GitHub API")
|
||||
return None, None
|
||||
|
||||
|
||||
async def get_github_default_branch(
|
||||
repo_owner: str,
|
||||
repo_name: str,
|
||||
github_token: str,
|
||||
) -> str:
|
||||
"""Get the default branch of a GitHub repository via the API.
|
||||
|
||||
Args:
|
||||
repo_owner: Repository owner (e.g., "langchain-ai")
|
||||
repo_name: Repository name (e.g., "deepagents")
|
||||
github_token: GitHub access token
|
||||
|
||||
Returns:
|
||||
The default branch name (e.g., "main" or "master")
|
||||
"""
|
||||
try:
|
||||
async with httpx.AsyncClient() as http_client:
|
||||
response = await http_client.get(
|
||||
f"https://api.github.com/repos/{repo_owner}/{repo_name}",
|
||||
headers={
|
||||
"Authorization": f"Bearer {github_token}",
|
||||
"Accept": "application/vnd.github+json",
|
||||
"X-GitHub-Api-Version": "2022-11-28",
|
||||
},
|
||||
)
|
||||
|
||||
if response.status_code == 200: # noqa: PLR2004
|
||||
repo_data = response.json()
|
||||
default_branch = repo_data.get("default_branch", "main")
|
||||
logger.debug("Got default branch from GitHub API: %s", default_branch)
|
||||
return default_branch
|
||||
|
||||
logger.warning(
|
||||
"Failed to get repo info from GitHub API (%s), falling back to 'main'",
|
||||
response.status_code,
|
||||
)
|
||||
return "main"
|
||||
|
||||
except httpx.HTTPError:
|
||||
logger.exception("Failed to get default branch from GitHub API, falling back to 'main'")
|
||||
return "main"
|
||||
|
||||
|
||||
async def comment_on_linear_issue(issue_id: str, comment_body: str) -> bool:
|
||||
"""Add a comment to a Linear issue.
|
||||
|
||||
Args:
|
||||
issue_id: The Linear issue ID
|
||||
comment_body: The comment text
|
||||
|
||||
Returns:
|
||||
True if successful, False otherwise
|
||||
"""
|
||||
if not LINEAR_API_KEY:
|
||||
return False
|
||||
|
||||
import httpx
|
||||
|
||||
url = "https://api.linear.app/graphql"
|
||||
|
||||
mutation = """
|
||||
mutation CommentCreate($issueId: String!, $body: String!) {
|
||||
commentCreate(input: { issueId: $issueId, body: $body }) {
|
||||
success
|
||||
comment {
|
||||
id
|
||||
}
|
||||
}
|
||||
}
|
||||
"""
|
||||
|
||||
async with httpx.AsyncClient() as http_client:
|
||||
try:
|
||||
response = await http_client.post(
|
||||
url,
|
||||
headers={
|
||||
"Authorization": LINEAR_API_KEY,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
json={
|
||||
"query": mutation,
|
||||
"variables": {"issueId": issue_id, "body": comment_body},
|
||||
},
|
||||
)
|
||||
response.raise_for_status()
|
||||
result = response.json()
|
||||
return bool(result.get("data", {}).get("commentCreate", {}).get("success"))
|
||||
except Exception: # noqa: BLE001
|
||||
return False
|
||||
|
||||
|
||||
class LinearNotifyState(AgentState):
|
||||
"""Extended agent state for tracking Linear notifications."""
|
||||
|
||||
linear_messages_sent_count: int
|
||||
|
||||
|
||||
@before_model(state_schema=LinearNotifyState)
|
||||
async def check_message_queue_before_model( # noqa: PLR0911
|
||||
state: LinearNotifyState, # noqa: ARG001
|
||||
runtime: Runtime, # noqa: ARG001
|
||||
) -> dict[str, Any] | None:
|
||||
"""Middleware that checks for queued messages before each model call.
|
||||
|
||||
If messages are found in the queue for this thread, it extracts all messages,
|
||||
adds them to the conversation state as new human messages, and clears the queue.
|
||||
Messages are processed in FIFO order (oldest first).
|
||||
|
||||
This enables handling of follow-up comments that arrive while the agent is busy.
|
||||
The agent will see the new messages and can incorporate them into its response.
|
||||
"""
|
||||
try:
|
||||
config = get_config()
|
||||
configurable = config.get("configurable", {})
|
||||
thread_id = configurable.get("thread_id")
|
||||
|
||||
if not thread_id:
|
||||
return None
|
||||
|
||||
try:
|
||||
store = get_store()
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.debug("Could not get store from context: %s", e)
|
||||
return None
|
||||
|
||||
if store is None:
|
||||
return None
|
||||
|
||||
namespace = ("queue", thread_id)
|
||||
|
||||
try:
|
||||
queued_item = await store.aget(namespace, "pending_messages")
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("Failed to get queued item: %s", e)
|
||||
return None
|
||||
|
||||
if queued_item is None:
|
||||
return None
|
||||
|
||||
queued_value = queued_item.value
|
||||
queued_messages = queued_value.get("messages", [])
|
||||
|
||||
# Delete early to prevent duplicate processing if middleware runs again
|
||||
await store.adelete(namespace, "pending_messages")
|
||||
|
||||
if not queued_messages:
|
||||
return None
|
||||
|
||||
logger.info(
|
||||
"Found %d queued message(s) for thread %s, injecting into state",
|
||||
len(queued_messages),
|
||||
thread_id,
|
||||
)
|
||||
|
||||
content_blocks = [
|
||||
{"type": "text", "text": msg.get("content", "")}
|
||||
for msg in queued_messages
|
||||
if msg.get("content")
|
||||
]
|
||||
|
||||
if not content_blocks:
|
||||
return None
|
||||
|
||||
new_message = {
|
||||
"role": "user",
|
||||
"content": content_blocks,
|
||||
}
|
||||
|
||||
logger.info(
|
||||
"Injected %d queued message(s) into state for thread %s",
|
||||
len(content_blocks),
|
||||
thread_id,
|
||||
)
|
||||
|
||||
return {"messages": [new_message]} # noqa: TRY300
|
||||
except Exception:
|
||||
logger.exception("Error in check_message_queue_before_model")
|
||||
return None
|
||||
|
||||
|
||||
@after_model(state_schema=LinearNotifyState)
|
||||
async def post_to_linear_after_model( # noqa: PLR0911, PLR0912
|
||||
state: LinearNotifyState,
|
||||
runtime: Runtime, # noqa: ARG001
|
||||
) -> dict[str, Any] | None:
|
||||
"""Middleware that posts AI responses to Linear after each model call.
|
||||
|
||||
Only posts if:
|
||||
- This is a Linear-triggered conversation (has linear_issue in config)
|
||||
- There's exactly 1 human message (initial request)
|
||||
- The previous message was from human (not a tool result)
|
||||
- The AI response has text content (not just tool calls)
|
||||
- The message hasn't already been sent (tracked via linear_messages_sent_count)
|
||||
"""
|
||||
try:
|
||||
config = get_config()
|
||||
configurable = config.get("configurable", {})
|
||||
|
||||
linear_issue = configurable.get("linear_issue", {})
|
||||
linear_issue_id = linear_issue.get("id")
|
||||
|
||||
if not linear_issue_id:
|
||||
return None
|
||||
|
||||
messages = state.get("messages", [])
|
||||
if not messages:
|
||||
return None
|
||||
|
||||
sent_count = state.get("linear_messages_sent_count", 0)
|
||||
|
||||
human_message_count = 0
|
||||
for msg in messages:
|
||||
if isinstance(msg, dict):
|
||||
role = msg.get("role", "")
|
||||
else:
|
||||
role = getattr(msg, "type", "") or getattr(msg, "role", "")
|
||||
if role in ("human", "user"):
|
||||
human_message_count += 1
|
||||
|
||||
if human_message_count != 1:
|
||||
return None
|
||||
|
||||
last_message = messages[-1]
|
||||
if isinstance(last_message, dict):
|
||||
role = last_message.get("role", "")
|
||||
content = last_message.get("content", "")
|
||||
else:
|
||||
role = getattr(last_message, "type", "") or getattr(last_message, "role", "")
|
||||
content = getattr(last_message, "content", "")
|
||||
|
||||
if role not in ("ai", "assistant"):
|
||||
return None
|
||||
|
||||
ai_message_count = 0
|
||||
for msg in messages:
|
||||
if isinstance(msg, dict):
|
||||
r = msg.get("role", "")
|
||||
else:
|
||||
r = getattr(msg, "type", "") or getattr(msg, "role", "")
|
||||
if r in ("ai", "assistant"):
|
||||
ai_message_count += 1
|
||||
|
||||
if ai_message_count <= sent_count:
|
||||
return None
|
||||
|
||||
if len(messages) >= MIN_MESSAGES_FOR_PREV_CHECK:
|
||||
prev_message = messages[-2]
|
||||
if isinstance(prev_message, dict):
|
||||
prev_role = prev_message.get("role", "")
|
||||
else:
|
||||
prev_role = getattr(prev_message, "type", "") or getattr(prev_message, "role", "")
|
||||
|
||||
if prev_role not in ("human", "user"):
|
||||
return None
|
||||
|
||||
if not content or not isinstance(content, str):
|
||||
return None
|
||||
|
||||
comment = f"""🤖 **Agent Response**
|
||||
|
||||
{content}"""
|
||||
logger.info("Posting AI response to Linear issue %s", linear_issue_id)
|
||||
success = await comment_on_linear_issue(linear_issue_id, comment)
|
||||
|
||||
if success:
|
||||
logger.info("Successfully posted to Linear")
|
||||
return {"linear_messages_sent_count": ai_message_count}
|
||||
logger.warning("Failed to post to Linear")
|
||||
|
||||
except Exception:
|
||||
logger.exception("Error in post_to_linear_after_model")
|
||||
return None
|
||||
|
||||
|
||||
@after_agent
|
||||
async def open_pr_if_needed( # noqa: PLR0912, PLR0915
|
||||
state: AgentState,
|
||||
runtime: Runtime, # noqa: ARG001
|
||||
) -> dict[str, Any] | None:
|
||||
"""Middleware that commits/pushes changes and comments on Linear after agent runs."""
|
||||
logger.info("After-agent middleware started")
|
||||
pr_url = None
|
||||
pr_number = None
|
||||
pr_title = "feat: Open SWE PR"
|
||||
|
||||
try:
|
||||
config = get_config()
|
||||
configurable = config.get("configurable", {})
|
||||
thread_id = configurable.get("thread_id")
|
||||
logger.debug("Middleware running for thread %s", thread_id)
|
||||
|
||||
last_message_content = ""
|
||||
messages = state.get("messages", [])
|
||||
if messages:
|
||||
last_message = messages[-1]
|
||||
if isinstance(last_message, dict):
|
||||
last_message_content = last_message.get("content", "")
|
||||
elif hasattr(last_message, "content"):
|
||||
last_message_content = last_message.content
|
||||
|
||||
linear_issue = configurable.get("linear_issue", {})
|
||||
linear_issue_id = linear_issue.get("id")
|
||||
|
||||
if not thread_id:
|
||||
if linear_issue_id and last_message_content:
|
||||
comment = f"""🤖 **Agent Response**
|
||||
|
||||
{last_message_content}"""
|
||||
await comment_on_linear_issue(linear_issue_id, comment)
|
||||
return None
|
||||
|
||||
repo_config = configurable.get("repo", {})
|
||||
repo_owner = repo_config.get("owner")
|
||||
repo_name = repo_config.get("name")
|
||||
|
||||
sandbox_backend = _SANDBOX_BACKENDS.get(thread_id)
|
||||
|
||||
repo_dir = f"/workspace/{repo_name}"
|
||||
|
||||
if not sandbox_backend or not repo_dir:
|
||||
if linear_issue_id and last_message_content:
|
||||
comment = f"""🤖 **Agent Response**
|
||||
|
||||
{last_message_content}"""
|
||||
await comment_on_linear_issue(linear_issue_id, comment)
|
||||
return None
|
||||
|
||||
result = await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git status --porcelain"
|
||||
)
|
||||
|
||||
has_uncommitted_changes = result.exit_code == 0 and result.output.strip()
|
||||
|
||||
await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git fetch origin 2>/dev/null || true"
|
||||
)
|
||||
git_log_cmd = (
|
||||
f"cd {repo_dir} && git log --oneline @{{upstream}}..HEAD 2>/dev/null "
|
||||
"|| git log --oneline origin/HEAD..HEAD 2>/dev/null || echo ''"
|
||||
)
|
||||
unpushed_result = await asyncio.to_thread(sandbox_backend.execute, git_log_cmd)
|
||||
has_unpushed_commits = unpushed_result.exit_code == 0 and unpushed_result.output.strip()
|
||||
|
||||
has_changes = has_uncommitted_changes or has_unpushed_commits
|
||||
|
||||
if not has_changes:
|
||||
logger.info("No changes detected, skipping PR creation")
|
||||
if linear_issue_id and last_message_content:
|
||||
comment = f"""🤖 **Agent Response**
|
||||
|
||||
{last_message_content}"""
|
||||
await comment_on_linear_issue(linear_issue_id, comment)
|
||||
return None
|
||||
|
||||
logger.info("Changes detected, preparing PR for thread %s", thread_id)
|
||||
|
||||
branch_result = await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git rev-parse --abbrev-ref HEAD"
|
||||
)
|
||||
current_branch = branch_result.output.strip() if branch_result.exit_code == 0 else ""
|
||||
|
||||
target_branch = f"open-swe/{thread_id}"
|
||||
|
||||
if current_branch != target_branch:
|
||||
checkout_result = await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git checkout -b {target_branch}"
|
||||
)
|
||||
if checkout_result.exit_code != 0:
|
||||
await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git checkout {target_branch}"
|
||||
)
|
||||
|
||||
await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git config user.name 'Open SWE[bot]'"
|
||||
)
|
||||
await asyncio.to_thread(
|
||||
sandbox_backend.execute,
|
||||
f"cd {repo_dir} && git config user.email 'Open SWE@users.noreply.github.com'",
|
||||
)
|
||||
|
||||
await asyncio.to_thread(sandbox_backend.execute, f"cd {repo_dir} && git add -A")
|
||||
|
||||
await asyncio.to_thread(
|
||||
sandbox_backend.execute, f'cd {repo_dir} && git commit -m "feat: Open SWE PR"'
|
||||
)
|
||||
|
||||
encrypted_token = configurable.get("github_token_encrypted")
|
||||
if encrypted_token:
|
||||
github_token = decrypt_token(encrypted_token)
|
||||
|
||||
if github_token:
|
||||
remote_result = await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git remote get-url origin"
|
||||
)
|
||||
if remote_result.exit_code == 0:
|
||||
remote_url = remote_result.output.strip()
|
||||
if "github.com" in remote_url and "@" not in remote_url:
|
||||
# Convert https://github.com/owner/repo.git to https://git:token@github.com/owner/repo.git
|
||||
auth_url = remote_url.replace("https://", f"https://git:{github_token}@")
|
||||
await asyncio.to_thread(
|
||||
sandbox_backend.execute,
|
||||
f"cd {repo_dir} && git push {auth_url} {target_branch}",
|
||||
)
|
||||
else:
|
||||
await asyncio.to_thread(
|
||||
sandbox_backend.execute, f"cd {repo_dir} && git push origin {target_branch}"
|
||||
)
|
||||
|
||||
# Get default branch from GitHub API (most reliable method)
|
||||
base_branch = await get_github_default_branch(repo_owner, repo_name, github_token)
|
||||
logger.info("Using base branch: %s", base_branch)
|
||||
|
||||
pr_title = "feat: Open SWE PR"
|
||||
pr_body = "Automated PR created by Open SWE agent."
|
||||
|
||||
pr_url, pr_number = await create_github_pr(
|
||||
repo_owner=repo_owner,
|
||||
repo_name=repo_name,
|
||||
github_token=github_token,
|
||||
title=pr_title,
|
||||
head_branch=target_branch,
|
||||
base_branch=base_branch,
|
||||
body=pr_body,
|
||||
)
|
||||
|
||||
linear_issue = configurable.get("linear_issue", {})
|
||||
linear_issue_id = linear_issue.get("id")
|
||||
|
||||
if linear_issue_id and last_message_content:
|
||||
if pr_url:
|
||||
comment = f"""✅ **Pull Request Created**
|
||||
|
||||
I've created a pull request to address this issue:
|
||||
|
||||
**[PR #{pr_number}: {pr_title}]({pr_url})**
|
||||
|
||||
---
|
||||
|
||||
🤖 **Agent Response**
|
||||
|
||||
{last_message_content}"""
|
||||
else:
|
||||
comment = f"""🤖 **Agent Response**
|
||||
|
||||
{last_message_content}"""
|
||||
await comment_on_linear_issue(linear_issue_id, comment)
|
||||
|
||||
logger.info("After-agent middleware completed successfully")
|
||||
|
||||
except Exception as e:
|
||||
logger.exception("Error in after-agent middleware")
|
||||
try:
|
||||
config = get_config()
|
||||
configurable = config.get("configurable", {})
|
||||
linear_issue = configurable.get("linear_issue", {})
|
||||
linear_issue_id = linear_issue.get("id")
|
||||
if linear_issue_id:
|
||||
error_comment = f"""❌ **Agent Error**
|
||||
|
||||
An error occurred while processing this issue:
|
||||
|
||||
```
|
||||
{type(e).__name__}: {e}
|
||||
```"""
|
||||
await comment_on_linear_issue(linear_issue_id, error_comment)
|
||||
except Exception:
|
||||
logger.exception("Failed to post error comment to Linear")
|
||||
return None
|
||||
|
||||
|
||||
async def _clone_or_pull_repo_in_sandbox( # noqa: PLR0915
|
||||
sandbox_backend: SandboxBackendProtocol,
|
||||
owner: str,
|
||||
repo: str,
|
||||
github_token: str | None = None,
|
||||
) -> str:
|
||||
"""Clone a GitHub repo into the sandbox, or pull if it already exists.
|
||||
|
||||
Args:
|
||||
sandbox_backend: The sandbox backend to execute commands in (LangSmithBackend)
|
||||
owner: GitHub repo owner
|
||||
repo: GitHub repo name
|
||||
github_token: GitHub access token (from agent auth or env var)
|
||||
|
||||
Returns:
|
||||
Path to the cloned/updated repo directory
|
||||
"""
|
||||
logger.info("_clone_or_pull_repo_in_sandbox called for %s/%s", owner, repo)
|
||||
loop = asyncio.get_event_loop()
|
||||
|
||||
token = github_token
|
||||
if not token:
|
||||
msg = "No GitHub token provided"
|
||||
logger.error(msg)
|
||||
raise ValueError(msg)
|
||||
|
||||
repo_dir = f"/workspace/{repo}"
|
||||
|
||||
logger.debug("Checking if repo already exists at %s", repo_dir)
|
||||
try:
|
||||
check_result = await loop.run_in_executor(
|
||||
None, sandbox_backend.execute, f"test -d {repo_dir}/.git && echo exists"
|
||||
)
|
||||
logger.debug(
|
||||
"Check result: exit_code=%s, output=%s",
|
||||
check_result.exit_code,
|
||||
check_result.output[:200] if check_result.output else "",
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to execute check command in sandbox")
|
||||
raise
|
||||
|
||||
if check_result.exit_code == 0 and "exists" in check_result.output:
|
||||
logger.info("Repo already exists at %s, pulling latest changes", repo_dir)
|
||||
try:
|
||||
status_result = await loop.run_in_executor(
|
||||
None, sandbox_backend.execute, f"cd {repo_dir} && git status --porcelain"
|
||||
)
|
||||
logger.debug("Git status result: exit_code=%s", status_result.exit_code)
|
||||
except Exception:
|
||||
logger.exception("Failed to get git status")
|
||||
raise
|
||||
|
||||
# CRITICAL: Ensure remote URL doesn't contain token (clean up from previous runs)
|
||||
clean_url = f"https://github.com/{owner}/{repo}.git"
|
||||
try:
|
||||
await loop.run_in_executor(
|
||||
None,
|
||||
sandbox_backend.execute,
|
||||
f"cd {repo_dir} && git remote set-url origin {clean_url}",
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to set remote URL")
|
||||
raise
|
||||
|
||||
if status_result.exit_code == 0 and not status_result.output.strip():
|
||||
auth_url = f"https://git:{token}@github.com/{owner}/{repo}.git"
|
||||
try:
|
||||
pull_result = await loop.run_in_executor(
|
||||
None, sandbox_backend.execute, f"cd {repo_dir} && git pull {auth_url}"
|
||||
)
|
||||
logger.debug("Git pull result: exit_code=%s", pull_result.exit_code)
|
||||
if pull_result.exit_code != 0:
|
||||
logger.warning(
|
||||
"Git pull failed with exit code %s: %s",
|
||||
pull_result.exit_code,
|
||||
pull_result.output[:200] if pull_result.output else "",
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to execute git pull")
|
||||
raise
|
||||
else:
|
||||
logger.info("Cloning repo %s/%s to %s", owner, repo, repo_dir)
|
||||
clone_url = f"https://git:{token}@github.com/{owner}/{repo}.git"
|
||||
try:
|
||||
result = await loop.run_in_executor(
|
||||
None, sandbox_backend.execute, f"git clone {clone_url} {repo_dir}"
|
||||
)
|
||||
logger.debug("Git clone result: exit_code=%s", result.exit_code)
|
||||
except Exception:
|
||||
logger.exception("Failed to execute git clone")
|
||||
raise
|
||||
|
||||
if result.exit_code != 0:
|
||||
msg = f"Failed to clone repo {owner}/{repo}: {result.output}"
|
||||
logger.error(msg)
|
||||
raise RuntimeError(msg)
|
||||
|
||||
clean_url = f"https://github.com/{owner}/{repo}.git"
|
||||
try:
|
||||
await loop.run_in_executor(
|
||||
None,
|
||||
sandbox_backend.execute,
|
||||
f"cd {repo_dir} && git remote set-url origin {clean_url}",
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to set remote URL after clone")
|
||||
raise
|
||||
|
||||
logger.info("Repo setup complete at %s", repo_dir)
|
||||
return repo_dir
|
||||
|
||||
|
||||
async def _get_sandbox_id_from_metadata(thread_id: str) -> str | None:
|
||||
"""Get sandbox_id from thread metadata."""
|
||||
thread = await client.threads.get(thread_id=thread_id)
|
||||
return thread.get("metadata", {}).get("sandbox_id")
|
||||
|
||||
|
||||
async def _wait_for_sandbox_id(thread_id: str) -> str:
|
||||
"""Wait for sandbox_id to be set in thread metadata.
|
||||
|
||||
Polls thread metadata until sandbox_id is set to a real value
|
||||
(not the creating sentinel).
|
||||
|
||||
Raises:
|
||||
TimeoutError: If sandbox creation takes too long
|
||||
"""
|
||||
elapsed = 0.0
|
||||
while elapsed < SANDBOX_CREATION_TIMEOUT:
|
||||
sandbox_id = await _get_sandbox_id_from_metadata(thread_id)
|
||||
if sandbox_id is not None and sandbox_id != SANDBOX_CREATING:
|
||||
return sandbox_id
|
||||
await asyncio.sleep(SANDBOX_POLL_INTERVAL)
|
||||
elapsed += SANDBOX_POLL_INTERVAL
|
||||
|
||||
msg = f"Timeout waiting for sandbox creation for thread {thread_id}"
|
||||
raise TimeoutError(msg)
|
||||
|
||||
|
||||
def graph_loaded_for_execution(config: RunnableConfig) -> bool:
|
||||
"""Check if the graph is loaded for actual execution vs introspection."""
|
||||
return (
|
||||
config["configurable"].get("__is_for_execution__", False)
|
||||
if "configurable" in config
|
||||
else False
|
||||
)
|
||||
|
||||
|
||||
async def get_agent(config: RunnableConfig) -> Pregel: # noqa: PLR0915
|
||||
"""Get or create an agent with a sandbox for the given thread."""
|
||||
thread_id = config["configurable"].get("thread_id", None)
|
||||
logger.info("get_agent called for thread %s", thread_id)
|
||||
|
||||
repo_config = config["configurable"].get("repo", {})
|
||||
repo_owner = repo_config.get("owner")
|
||||
repo_name = repo_config.get("name")
|
||||
|
||||
encrypted_token = config["configurable"].get("github_token_encrypted")
|
||||
if encrypted_token:
|
||||
github_token = decrypt_token(encrypted_token)
|
||||
logger.debug("Decrypted GitHub token")
|
||||
|
||||
if thread_id is None or not graph_loaded_for_execution(config):
|
||||
logger.info("No thread_id or not for execution, returning agent without sandbox")
|
||||
return create_server_agent(
|
||||
model=None,
|
||||
assistant_id="agent",
|
||||
tools=tools,
|
||||
sandbox=None,
|
||||
sandbox_type=None,
|
||||
auto_approve=True,
|
||||
)
|
||||
|
||||
sandbox_id = await _get_sandbox_id_from_metadata(thread_id)
|
||||
|
||||
if sandbox_id == SANDBOX_CREATING:
|
||||
logger.info("Sandbox creation in progress, waiting...")
|
||||
sandbox_id = await _wait_for_sandbox_id(thread_id)
|
||||
|
||||
if sandbox_id is None:
|
||||
logger.info("Creating new sandbox for thread %s", thread_id)
|
||||
await client.threads.update(thread_id=thread_id, metadata={"sandbox_id": SANDBOX_CREATING})
|
||||
|
||||
try:
|
||||
sandbox_cm = create_sandbox_async("langsmith", cleanup=False)
|
||||
sandbox_backend = await sandbox_cm.__aenter__()
|
||||
logger.info("Sandbox created: %s", sandbox_backend.id)
|
||||
|
||||
# Update metadata immediately after sandbox creation so other callers
|
||||
# can connect to the sandbox while we clone the repo
|
||||
await client.threads.update(
|
||||
thread_id=thread_id,
|
||||
metadata={"sandbox_id": sandbox_backend.id},
|
||||
)
|
||||
|
||||
repo_dir = None
|
||||
if repo_owner and repo_name:
|
||||
logger.info("Cloning repo %s/%s into sandbox", repo_owner, repo_name)
|
||||
repo_dir = await _clone_or_pull_repo_in_sandbox(
|
||||
sandbox_backend, repo_owner, repo_name, github_token
|
||||
)
|
||||
logger.info("Repo cloned to %s", repo_dir)
|
||||
|
||||
await client.threads.update(
|
||||
thread_id=thread_id,
|
||||
metadata={"repo_dir": repo_dir},
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to create sandbox or clone repo")
|
||||
try:
|
||||
await client.threads.update(thread_id=thread_id, metadata={"sandbox_id": None})
|
||||
logger.info("Reset sandbox_id to None for thread %s", thread_id)
|
||||
except Exception:
|
||||
logger.exception("Failed to reset sandbox_id metadata")
|
||||
raise
|
||||
else:
|
||||
logger.info("Connecting to existing sandbox %s", sandbox_id)
|
||||
try:
|
||||
sandbox_cm = create_sandbox_async("langsmith", sandbox_id=sandbox_id, cleanup=False)
|
||||
sandbox_backend = await sandbox_cm.__aenter__()
|
||||
logger.info("Connected to existing sandbox %s", sandbox_id)
|
||||
except Exception:
|
||||
logger.exception("Failed to connect to existing sandbox %s", sandbox_id)
|
||||
raise
|
||||
|
||||
thread = await client.threads.get(thread_id=thread_id)
|
||||
repo_dir = thread.get("metadata", {}).get("repo_dir")
|
||||
|
||||
if repo_owner and repo_name:
|
||||
logger.info("Pulling latest changes for repo %s/%s", repo_owner, repo_name)
|
||||
try:
|
||||
repo_dir = await _clone_or_pull_repo_in_sandbox(
|
||||
sandbox_backend, repo_owner, repo_name, github_token
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to pull repo in existing sandbox")
|
||||
raise
|
||||
|
||||
_SANDBOX_BACKENDS[thread_id] = sandbox_backend
|
||||
|
||||
logger.info("Returning agent with sandbox for thread %s", thread_id)
|
||||
return create_server_agent(
|
||||
model=None,
|
||||
assistant_id="agent",
|
||||
tools=tools,
|
||||
sandbox=sandbox_backend,
|
||||
sandbox_type="langsmith",
|
||||
auto_approve=True,
|
||||
working_dir=repo_dir,
|
||||
middleware=[
|
||||
check_message_queue_before_model,
|
||||
post_to_linear_after_model,
|
||||
open_pr_if_needed,
|
||||
],
|
||||
)
|
||||
762
apps/agent/agent/webapp.py
Normal file
762
apps/agent/agent/webapp.py
Normal file
|
|
@ -0,0 +1,762 @@
|
|||
"""Custom FastAPI routes for LangGraph server."""
|
||||
|
||||
import hashlib
|
||||
import hmac
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
import jwt
|
||||
from fastapi import BackgroundTasks, FastAPI, HTTPException, Request
|
||||
from langgraph_sdk import get_client
|
||||
|
||||
# Local import for encryption
|
||||
from .encryption import encrypt_token
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
app = FastAPI()
|
||||
|
||||
LINEAR_WEBHOOK_SECRET = os.environ.get("LINEAR_WEBHOOK_SECRET", "")
|
||||
|
||||
LANGGRAPH_URL = os.environ.get("LANGGRAPH_URL") or os.environ.get(
|
||||
"LANGGRAPH_URL_PROD", "http://localhost:2024"
|
||||
)
|
||||
|
||||
LANGSMITH_API_KEY = os.environ.get("LANGSMITH_API_KEY") or os.environ.get(
|
||||
"LANGSMITH_API_KEY_PROD", ""
|
||||
)
|
||||
LANGSMITH_API_URL = os.environ.get("LANGSMITH_ENDPOINT", "https://api.smith.langchain.com")
|
||||
|
||||
GITHUB_OAUTH_PROVIDER_ID = os.environ.get("GITHUB_OAUTH_PROVIDER_ID", "")
|
||||
|
||||
LINEAR_API_KEY = os.environ.get("LINEAR_API_KEY", "")
|
||||
|
||||
X_SERVICE_AUTH_JWT_SECRET = os.environ.get("X_SERVICE_AUTH_JWT_SECRET", "")
|
||||
|
||||
|
||||
def get_service_jwt_token_for_user(
|
||||
user_id: str, tenant_id: str, expiration_seconds: int = 300
|
||||
) -> str:
|
||||
"""Create a short-lived service JWT for authenticating as a specific user.
|
||||
|
||||
Args:
|
||||
tenant_id: The LangSmith tenant ID to associate with the token
|
||||
user_id: The LangSmith user ID to associate with the token
|
||||
expiration_seconds: Token expiration time in seconds (default: 5 minutes)
|
||||
|
||||
Returns:
|
||||
JWT token string
|
||||
|
||||
Raises:
|
||||
ValueError: If X_SERVICE_AUTH_JWT_SECRET is not configured
|
||||
"""
|
||||
if not X_SERVICE_AUTH_JWT_SECRET:
|
||||
msg = "X_SERVICE_AUTH_JWT_SECRET is not configured. Cannot generate service keys."
|
||||
raise ValueError(msg)
|
||||
|
||||
exp_datetime = datetime.now(tz=UTC) + timedelta(seconds=expiration_seconds)
|
||||
exp = int(exp_datetime.timestamp())
|
||||
|
||||
payload = {
|
||||
"sub": "unspecified",
|
||||
"exp": exp,
|
||||
"tenant_id": tenant_id,
|
||||
"user_id": user_id,
|
||||
}
|
||||
|
||||
return jwt.encode(payload, X_SERVICE_AUTH_JWT_SECRET, algorithm="HS256")
|
||||
|
||||
|
||||
LINEAR_TEAM_TO_REPO: dict[str, dict[str, str]] = {
|
||||
"Brace's test workspace": {"owner": "langchain-ai", "name": "open-swe"},
|
||||
"Yogesh-dev": {"owner": "aran-yogesh", "name": "nimedge"},
|
||||
}
|
||||
|
||||
|
||||
async def get_ls_user_id_from_email(email: str) -> dict[str, str | None]:
|
||||
"""Get the LangSmith user ID and tenant ID from a user's email.
|
||||
|
||||
Args:
|
||||
email: The user's email address
|
||||
|
||||
Returns:
|
||||
Dict with 'ls_user_id' and 'tenant_id' keys (values may be None if not found)
|
||||
"""
|
||||
if not LANGSMITH_API_KEY:
|
||||
return {"ls_user_id": None, "tenant_id": None}
|
||||
|
||||
url = f"{LANGSMITH_API_URL}/api/v1/workspaces/current/members/active"
|
||||
|
||||
async with httpx.AsyncClient() as client:
|
||||
try:
|
||||
response = await client.get(
|
||||
url,
|
||||
headers={"X-API-Key": LANGSMITH_API_KEY},
|
||||
params={"emails": [email]},
|
||||
)
|
||||
response.raise_for_status()
|
||||
members = response.json()
|
||||
|
||||
if members and len(members) > 0:
|
||||
member = members[0]
|
||||
return {
|
||||
"ls_user_id": member.get("ls_user_id"),
|
||||
"tenant_id": member.get("tenant_id"),
|
||||
}
|
||||
except httpx.HTTPError:
|
||||
logger.debug("HTTP error getting LangSmith user info for email")
|
||||
return {"ls_user_id": None, "tenant_id": None}
|
||||
|
||||
|
||||
LANGSMITH_HOST_API_URL = os.environ.get("LANGSMITH_HOST_API_URL", "https://api.host.langchain.com")
|
||||
|
||||
|
||||
async def get_github_token_for_user(ls_user_id: str, tenant_id: str) -> dict[str, Any]:
|
||||
"""Get GitHub OAuth token for a user via LangSmith agent auth.
|
||||
|
||||
Args:
|
||||
ls_user_id: The LangSmith user ID
|
||||
tenant_id: The LangSmith tenant ID
|
||||
|
||||
Returns:
|
||||
Dict with either 'token' key or 'auth_url' key
|
||||
"""
|
||||
if not GITHUB_OAUTH_PROVIDER_ID:
|
||||
return {"error": "GITHUB_OAUTH_PROVIDER_ID not configured"}
|
||||
|
||||
try:
|
||||
service_token = get_service_jwt_token_for_user(ls_user_id, tenant_id)
|
||||
|
||||
headers = {
|
||||
"X-Service-Key": service_token,
|
||||
"X-Tenant-Id": tenant_id,
|
||||
"X-User-Id": ls_user_id,
|
||||
}
|
||||
|
||||
payload = {
|
||||
"provider": GITHUB_OAUTH_PROVIDER_ID,
|
||||
"scopes": ["repo"],
|
||||
"user_id": ls_user_id,
|
||||
"ls_user_id": ls_user_id,
|
||||
}
|
||||
|
||||
async with httpx.AsyncClient() as client:
|
||||
response = await client.post(
|
||||
f"{LANGSMITH_HOST_API_URL}/v2/auth/authenticate",
|
||||
json=payload,
|
||||
headers=headers,
|
||||
)
|
||||
response.raise_for_status()
|
||||
response_data = response.json()
|
||||
|
||||
token = response_data.get("token")
|
||||
auth_url = response_data.get("url")
|
||||
|
||||
if token:
|
||||
return {"token": token}
|
||||
if auth_url:
|
||||
return {"auth_url": auth_url}
|
||||
return {"error": f"Unexpected auth result: {response_data}"}
|
||||
|
||||
except httpx.HTTPStatusError as e:
|
||||
return {"error": f"HTTP error: {e.response.status_code} - {e.response.text}"}
|
||||
except Exception as e: # noqa: BLE001
|
||||
return {"error": str(e)}
|
||||
|
||||
|
||||
async def react_to_linear_comment(comment_id: str, emoji: str = "👀") -> bool:
|
||||
"""Add an emoji reaction to a Linear comment.
|
||||
|
||||
Args:
|
||||
comment_id: The Linear comment ID
|
||||
emoji: The emoji to react with (default: eyes 👀)
|
||||
|
||||
Returns:
|
||||
True if successful, False otherwise
|
||||
"""
|
||||
if not LINEAR_API_KEY:
|
||||
return False
|
||||
|
||||
url = "https://api.linear.app/graphql"
|
||||
|
||||
mutation = """
|
||||
mutation ReactionCreate($commentId: String!, $emoji: String!) {
|
||||
reactionCreate(input: { commentId: $commentId, emoji: $emoji }) {
|
||||
success
|
||||
}
|
||||
}
|
||||
"""
|
||||
|
||||
async with httpx.AsyncClient() as client:
|
||||
try:
|
||||
response = await client.post(
|
||||
url,
|
||||
headers={
|
||||
"Authorization": LINEAR_API_KEY,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
json={
|
||||
"query": mutation,
|
||||
"variables": {"commentId": comment_id, "emoji": emoji},
|
||||
},
|
||||
)
|
||||
response.raise_for_status()
|
||||
result = response.json()
|
||||
return bool(result.get("data", {}).get("reactionCreate", {}).get("success"))
|
||||
except Exception: # noqa: BLE001
|
||||
return False
|
||||
|
||||
|
||||
async def comment_on_linear_issue(issue_id: str, comment_body: str) -> bool:
|
||||
"""Add a comment to a Linear issue.
|
||||
|
||||
Args:
|
||||
issue_id: The Linear issue ID
|
||||
comment_body: The comment text
|
||||
|
||||
Returns:
|
||||
True if successful, False otherwise
|
||||
"""
|
||||
if not LINEAR_API_KEY:
|
||||
return False
|
||||
|
||||
url = "https://api.linear.app/graphql"
|
||||
|
||||
mutation = """
|
||||
mutation CommentCreate($issueId: String!, $body: String!) {
|
||||
commentCreate(input: { issueId: $issueId, body: $body }) {
|
||||
success
|
||||
comment {
|
||||
id
|
||||
}
|
||||
}
|
||||
}
|
||||
"""
|
||||
|
||||
async with httpx.AsyncClient() as client:
|
||||
try:
|
||||
response = await client.post(
|
||||
url,
|
||||
headers={
|
||||
"Authorization": LINEAR_API_KEY,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
json={
|
||||
"query": mutation,
|
||||
"variables": {"issueId": issue_id, "body": comment_body},
|
||||
},
|
||||
)
|
||||
response.raise_for_status()
|
||||
result = response.json()
|
||||
return bool(result.get("data", {}).get("commentCreate", {}).get("success"))
|
||||
except Exception: # noqa: BLE001
|
||||
return False
|
||||
|
||||
|
||||
async def fetch_linear_issue_details(issue_id: str) -> dict[str, Any] | None:
|
||||
"""Fetch full issue details from Linear API including description and comments.
|
||||
|
||||
Args:
|
||||
issue_id: The Linear issue ID
|
||||
|
||||
Returns:
|
||||
Full issue data dict, or None if fetch failed
|
||||
"""
|
||||
if not LINEAR_API_KEY:
|
||||
return None
|
||||
|
||||
url = "https://api.linear.app/graphql"
|
||||
|
||||
query = """
|
||||
query GetIssue($issueId: String!) {
|
||||
issue(id: $issueId) {
|
||||
id
|
||||
identifier
|
||||
title
|
||||
description
|
||||
url
|
||||
comments {
|
||||
nodes {
|
||||
id
|
||||
body
|
||||
createdAt
|
||||
user {
|
||||
id
|
||||
name
|
||||
email
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
"""
|
||||
|
||||
async with httpx.AsyncClient() as client:
|
||||
try:
|
||||
response = await client.post(
|
||||
url,
|
||||
headers={
|
||||
"Authorization": LINEAR_API_KEY,
|
||||
"Content-Type": "application/json",
|
||||
},
|
||||
json={
|
||||
"query": query,
|
||||
"variables": {"issueId": issue_id},
|
||||
},
|
||||
)
|
||||
response.raise_for_status()
|
||||
result = response.json()
|
||||
|
||||
return result.get("data", {}).get("issue")
|
||||
except httpx.HTTPError:
|
||||
return None
|
||||
|
||||
|
||||
def generate_thread_id_from_issue(issue_id: str) -> str:
|
||||
"""Generate a deterministic thread ID from a Linear issue ID.
|
||||
|
||||
Args:
|
||||
issue_id: The Linear issue ID
|
||||
|
||||
Returns:
|
||||
A UUID-formatted thread ID derived from the issue ID
|
||||
"""
|
||||
hash_bytes = hashlib.sha256(f"linear-issue:{issue_id}".encode()).hexdigest()
|
||||
return (
|
||||
f"{hash_bytes[:8]}-{hash_bytes[8:12]}-{hash_bytes[12:16]}-"
|
||||
f"{hash_bytes[16:20]}-{hash_bytes[20:32]}"
|
||||
)
|
||||
|
||||
|
||||
async def is_thread_active(thread_id: str) -> bool:
|
||||
"""Check if a thread is currently active (has a running run).
|
||||
|
||||
Args:
|
||||
thread_id: The LangGraph thread ID
|
||||
|
||||
Returns:
|
||||
True if the thread status is "busy", False otherwise
|
||||
"""
|
||||
langgraph_client = get_client(url=LANGGRAPH_URL)
|
||||
try:
|
||||
logger.debug("Fetching thread status for %s from %s", thread_id, LANGGRAPH_URL)
|
||||
thread = await langgraph_client.threads.get(thread_id)
|
||||
status = thread.get("status", "idle")
|
||||
logger.info(
|
||||
"Thread %s status check: status=%s, is_busy=%s",
|
||||
thread_id,
|
||||
status,
|
||||
status == "busy",
|
||||
)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning(
|
||||
"Failed to get thread status for %s: %s (type: %s) - assuming not active",
|
||||
thread_id,
|
||||
e,
|
||||
type(e).__name__,
|
||||
)
|
||||
status = "idle"
|
||||
return status == "busy"
|
||||
|
||||
|
||||
async def queue_message_for_thread(thread_id: str, message_content: str) -> bool:
|
||||
"""Queue a message for a thread that is currently active.
|
||||
|
||||
Stores the message in the langgraph store, namespaced to the thread.
|
||||
Supports multiple queued messages by storing them as a list (FIFO order).
|
||||
The before_model middleware will pick them up and inject them into state.
|
||||
|
||||
Args:
|
||||
thread_id: The LangGraph thread ID
|
||||
message_content: The message content to queue
|
||||
|
||||
Returns:
|
||||
True if successfully queued, False otherwise
|
||||
"""
|
||||
langgraph_client = get_client(url=LANGGRAPH_URL)
|
||||
try:
|
||||
namespace = ("queue", thread_id)
|
||||
key = "pending_messages"
|
||||
|
||||
new_message = {"content": message_content}
|
||||
|
||||
existing_messages: list[dict[str, Any]] = []
|
||||
try:
|
||||
existing_item = await langgraph_client.store.get_item(namespace, key)
|
||||
if existing_item and existing_item.get("value"):
|
||||
existing_messages = existing_item["value"].get("messages", [])
|
||||
except Exception: # noqa: BLE001
|
||||
logger.debug("No existing queued messages for thread %s", thread_id)
|
||||
|
||||
existing_messages.append(new_message)
|
||||
value = {"messages": existing_messages}
|
||||
|
||||
logger.info(
|
||||
"Attempting to queue message for thread %s (total queued: %d)",
|
||||
thread_id,
|
||||
len(existing_messages),
|
||||
)
|
||||
await langgraph_client.store.put_item(namespace, key, value)
|
||||
logger.info("Successfully queued message for thread %s", thread_id)
|
||||
return True # noqa: TRY300
|
||||
except Exception:
|
||||
logger.exception("Failed to queue message for thread %s", thread_id)
|
||||
return False
|
||||
|
||||
|
||||
async def process_linear_issue( # noqa: PLR0912, PLR0915
|
||||
issue_data: dict[str, Any], repo_config: dict[str, str]
|
||||
) -> None:
|
||||
"""Process a Linear issue by creating a new LangGraph thread and run.
|
||||
|
||||
Args:
|
||||
issue_data: The Linear issue data from webhook (basic info only).
|
||||
repo_config: The repo configuration with owner and name.
|
||||
"""
|
||||
issue_id = issue_data.get("id", "")
|
||||
logger.info(
|
||||
"Processing Linear issue %s for repo %s/%s",
|
||||
issue_id,
|
||||
repo_config.get("owner"),
|
||||
repo_config.get("name"),
|
||||
)
|
||||
|
||||
triggering_comment_id = issue_data.get("triggering_comment_id", "")
|
||||
if triggering_comment_id:
|
||||
await react_to_linear_comment(triggering_comment_id, "👀")
|
||||
|
||||
thread_id = generate_thread_id_from_issue(issue_id)
|
||||
|
||||
full_issue = await fetch_linear_issue_details(issue_id)
|
||||
if not full_issue:
|
||||
full_issue = issue_data
|
||||
|
||||
user_email = None
|
||||
user_name = None
|
||||
comment_author = issue_data.get("comment_author", {})
|
||||
if comment_author:
|
||||
user_email = comment_author.get("email")
|
||||
user_name = comment_author.get("name")
|
||||
if not user_email:
|
||||
creator = full_issue.get("creator", {})
|
||||
if creator:
|
||||
user_email = creator.get("email")
|
||||
user_name = user_name or creator.get("name")
|
||||
if not user_email:
|
||||
assignee = full_issue.get("assignee", {})
|
||||
if assignee:
|
||||
user_email = assignee.get("email")
|
||||
user_name = user_name or assignee.get("name")
|
||||
|
||||
user_mention = f"@{user_name}" if user_name else ""
|
||||
|
||||
logger.info(
|
||||
"User email: %s, GITHUB_OAUTH_PROVIDER_ID set: %s",
|
||||
user_email,
|
||||
bool(GITHUB_OAUTH_PROVIDER_ID),
|
||||
)
|
||||
|
||||
github_token = None
|
||||
if user_email and GITHUB_OAUTH_PROVIDER_ID:
|
||||
user_info = await get_ls_user_id_from_email(user_email)
|
||||
ls_user_id = user_info.get("ls_user_id")
|
||||
tenant_id = user_info.get("tenant_id")
|
||||
logger.info(
|
||||
"LangSmith user ID for %s: %s, tenant_id: %s", user_email, ls_user_id, tenant_id
|
||||
)
|
||||
|
||||
if ls_user_id and tenant_id:
|
||||
auth_result = await get_github_token_for_user(ls_user_id, tenant_id)
|
||||
logger.info("Auth result keys: %s", list(auth_result.keys()))
|
||||
|
||||
if "token" in auth_result:
|
||||
github_token = auth_result["token"]
|
||||
logger.info("GitHub token obtained for user %s", user_email)
|
||||
elif "auth_url" in auth_result:
|
||||
auth_url = auth_result["auth_url"]
|
||||
logger.info("GitHub auth required for user %s, sending auth URL", user_email)
|
||||
comment = (
|
||||
f"🔐 **GitHub Authentication Required** {user_mention}\n\n"
|
||||
"To allow the Open SWE agent to work on this issue, "
|
||||
"please authenticate with GitHub by clicking the link below:\n\n"
|
||||
f"[Authenticate with GitHub]({auth_url})\n\n"
|
||||
"Once authenticated, reply to this issue mentioning @openswe to retry."
|
||||
)
|
||||
|
||||
await comment_on_linear_issue(issue_id, comment)
|
||||
return
|
||||
else:
|
||||
logger.warning("Auth result has neither token nor auth_url: %s", auth_result)
|
||||
else:
|
||||
logger.warning("User %s not found in LangSmith workspace", user_email)
|
||||
comment = (
|
||||
f"🔐 **GitHub Authentication Required** {user_mention}\n\n"
|
||||
f"Could not find a LangSmith account for **{user_email}**.\n\n"
|
||||
"Please ensure this email is invited to the main LangSmith organization. "
|
||||
"If your Linear account uses a different email than your LangSmith account, "
|
||||
"you may need to update one of them to match.\n\n"
|
||||
"Once your email is added to LangSmith, "
|
||||
"reply to this issue mentioning @openswe to retry."
|
||||
)
|
||||
|
||||
await comment_on_linear_issue(issue_id, comment)
|
||||
return
|
||||
|
||||
title = full_issue.get("title", "No title")
|
||||
description = full_issue.get("description") or "No description"
|
||||
|
||||
comments = full_issue.get("comments", {}).get("nodes", [])
|
||||
comments_text = ""
|
||||
|
||||
bot_message_prefixes = (
|
||||
"🔐 **GitHub Authentication Required**",
|
||||
"✅ **Pull Request Created**",
|
||||
"🤖 **Agent Response**",
|
||||
"❌ **Agent Error**",
|
||||
)
|
||||
|
||||
if comments:
|
||||
last_bot_comment_idx = -1
|
||||
for i, comment in enumerate(comments):
|
||||
body = comment.get("body", "")
|
||||
if any(body.startswith(prefix) for prefix in bot_message_prefixes):
|
||||
last_bot_comment_idx = i
|
||||
|
||||
relevant_comments = []
|
||||
for i, comment in enumerate(comments):
|
||||
if i <= last_bot_comment_idx:
|
||||
continue
|
||||
body = comment.get("body", "")
|
||||
if "@openswe" in body.lower():
|
||||
relevant_comments.append(comment)
|
||||
relevant_comments.extend(comments[i + 1 :])
|
||||
break
|
||||
|
||||
if relevant_comments:
|
||||
comments_text = "\n\n## Comments:\n"
|
||||
for comment in relevant_comments:
|
||||
author = comment.get("user", {}).get("name", "Unknown")
|
||||
body = comment.get("body", "")
|
||||
if any(body.startswith(prefix) for prefix in bot_message_prefixes):
|
||||
continue
|
||||
comments_text += f"\n**{author}:** {body}\n"
|
||||
|
||||
prompt = (
|
||||
f"Please work on the following issue:\n\n"
|
||||
f"## Title: {title}\n\n"
|
||||
f"## Description:\n{description}\n"
|
||||
f"{comments_text}\n\n"
|
||||
"Please analyze this issue and implement the necessary changes. "
|
||||
"When you're done, commit and push your changes."
|
||||
)
|
||||
|
||||
configurable: dict[str, Any] = {
|
||||
"repo": repo_config,
|
||||
"linear_issue": {
|
||||
"id": issue_id,
|
||||
"title": title,
|
||||
"url": full_issue.get("url", "") or issue_data.get("url", ""),
|
||||
"identifier": full_issue.get("identifier", "") or issue_data.get("identifier", ""),
|
||||
},
|
||||
}
|
||||
if github_token:
|
||||
configurable["github_token_encrypted"] = encrypt_token(github_token)
|
||||
|
||||
logger.info("Checking if thread %s is active before creating run", thread_id)
|
||||
thread_active = await is_thread_active(thread_id)
|
||||
logger.info("Thread %s active status: %s", thread_id, thread_active)
|
||||
|
||||
if thread_active:
|
||||
logger.info(
|
||||
"Thread %s is active (busy), will queue message instead of creating run",
|
||||
thread_id,
|
||||
)
|
||||
|
||||
queued = await queue_message_for_thread(
|
||||
thread_id=thread_id,
|
||||
message_content=prompt,
|
||||
)
|
||||
|
||||
if queued:
|
||||
logger.info(
|
||||
"Message queued for thread %s, will be processed by middleware", thread_id
|
||||
)
|
||||
else:
|
||||
logger.error("Failed to queue message for thread %s", thread_id)
|
||||
else:
|
||||
logger.info("Creating LangGraph run for thread %s", thread_id)
|
||||
langgraph_client = get_client(url=LANGGRAPH_URL)
|
||||
await langgraph_client.runs.create(
|
||||
thread_id,
|
||||
"agent",
|
||||
input={"messages": [{"role": "user", "content": prompt}]},
|
||||
config={"configurable": configurable},
|
||||
if_not_exists="create",
|
||||
)
|
||||
logger.info("LangGraph run created successfully for thread %s", thread_id)
|
||||
else:
|
||||
logger.warning("No GitHub token available, cannot create run for issue %s", issue_id)
|
||||
if not user_email:
|
||||
comment = (
|
||||
f"🔐 **GitHub Authentication Required** {user_mention}\n\n"
|
||||
"Could not determine the user email from this issue. "
|
||||
"Please ensure your Linear account has an email address configured.\n\n"
|
||||
"Reply to this issue mentioning @openswe to retry."
|
||||
)
|
||||
elif not GITHUB_OAUTH_PROVIDER_ID:
|
||||
comment = (
|
||||
f"❌ **Configuration Error** {user_mention}\n\n"
|
||||
"The Open SWE agent is not properly configured (missing GitHub OAuth provider).\n\n"
|
||||
"Please contact your administrator."
|
||||
)
|
||||
else:
|
||||
comment = (
|
||||
f"🔐 **GitHub Authentication Required** {user_mention}\n\n"
|
||||
"Unable to authenticate with GitHub. "
|
||||
"Please ensure you have connected your GitHub account in LangSmith.\n\n"
|
||||
"Reply to this issue mentioning @openswe to retry."
|
||||
)
|
||||
|
||||
await comment_on_linear_issue(issue_id, comment)
|
||||
|
||||
|
||||
def verify_linear_signature(body: bytes, signature: str, secret: str) -> bool:
|
||||
"""Verify the Linear webhook signature.
|
||||
|
||||
Args:
|
||||
body: Raw request body bytes
|
||||
signature: The Linear-Signature header value
|
||||
secret: The webhook signing secret
|
||||
|
||||
Returns:
|
||||
True if signature is valid, False otherwise
|
||||
"""
|
||||
if not secret:
|
||||
return True
|
||||
|
||||
expected = hmac.new(secret.encode("utf-8"), body, hashlib.sha256).hexdigest()
|
||||
|
||||
return hmac.compare_digest(expected, signature)
|
||||
|
||||
|
||||
@app.post("/webhooks/linear")
|
||||
async def linear_webhook( # noqa: PLR0911, PLR0912, PLR0915
|
||||
request: Request, background_tasks: BackgroundTasks
|
||||
) -> dict[str, str]:
|
||||
"""Handle Linear webhooks.
|
||||
|
||||
Triggers a new LangGraph run when an issue gets the 'open-swe' label added.
|
||||
"""
|
||||
logger.info("Received Linear webhook")
|
||||
body = await request.body()
|
||||
|
||||
signature = request.headers.get("Linear-Signature", "")
|
||||
if LINEAR_WEBHOOK_SECRET and not verify_linear_signature(
|
||||
body, signature, LINEAR_WEBHOOK_SECRET
|
||||
):
|
||||
logger.warning("Invalid webhook signature")
|
||||
raise HTTPException(status_code=401, detail="Invalid signature")
|
||||
|
||||
try:
|
||||
payload = json.loads(body)
|
||||
except json.JSONDecodeError:
|
||||
logger.exception("Failed to parse webhook JSON")
|
||||
return {"status": "error", "message": "Invalid JSON"}
|
||||
|
||||
if payload.get("type") != "Comment":
|
||||
logger.debug("Ignoring webhook: not a Comment event")
|
||||
return {"status": "ignored", "reason": "Not a Comment event"}
|
||||
|
||||
action = payload.get("action")
|
||||
if action != "create":
|
||||
logger.debug("Ignoring webhook: action is %s, not create", action)
|
||||
return {
|
||||
"status": "ignored",
|
||||
"reason": f"Comment action is '{action}', only processing 'create'",
|
||||
}
|
||||
|
||||
data = payload.get("data", {})
|
||||
|
||||
if data.get("botActor"):
|
||||
logger.debug("Ignoring webhook: comment is from a bot")
|
||||
return {"status": "ignored", "reason": "Comment is from a bot"}
|
||||
|
||||
comment_body = data.get("body", "")
|
||||
bot_message_prefixes = [
|
||||
"🔐 **GitHub Authentication Required**",
|
||||
"✅ **Pull Request Created**",
|
||||
"🤖 **Agent Response**",
|
||||
"❌ **Agent Error**",
|
||||
]
|
||||
for prefix in bot_message_prefixes:
|
||||
if comment_body.startswith(prefix):
|
||||
logger.debug("Ignoring webhook: comment is our own bot message")
|
||||
return {"status": "ignored", "reason": "Comment is our own bot message"}
|
||||
if "@openswe" not in comment_body.lower():
|
||||
logger.debug("Ignoring webhook: comment doesn't mention @openswe")
|
||||
return {"status": "ignored", "reason": "Comment doesn't mention @openswe"}
|
||||
|
||||
issue = data.get("issue", {})
|
||||
if not issue:
|
||||
logger.debug("Ignoring webhook: no issue data in comment")
|
||||
return {"status": "ignored", "reason": "No issue data in comment"}
|
||||
|
||||
team = issue.get("team", {})
|
||||
team_id = team.get("id", "") if team else ""
|
||||
team_name = team.get("name", "") if team else ""
|
||||
|
||||
repo_config = None
|
||||
if team_id and team_id in LINEAR_TEAM_TO_REPO:
|
||||
repo_config = LINEAR_TEAM_TO_REPO[team_id]
|
||||
elif team_name and team_name in LINEAR_TEAM_TO_REPO:
|
||||
repo_config = LINEAR_TEAM_TO_REPO[team_name]
|
||||
|
||||
if not repo_config:
|
||||
for label in issue.get("labels", []):
|
||||
label_name = label.get("name", "")
|
||||
if label_name.startswith("repo:"):
|
||||
repo_ref = label_name[5:] # Remove "repo:" prefix
|
||||
if "/" in repo_ref:
|
||||
owner, name = repo_ref.split("/", 1)
|
||||
repo_config = {"owner": owner, "name": name}
|
||||
break
|
||||
|
||||
if not repo_config:
|
||||
repo_config = {"owner": "langchain-ai", "name": "langchainplus"}
|
||||
|
||||
repo_owner = repo_config["owner"]
|
||||
repo_name = repo_config["name"]
|
||||
|
||||
issue["triggering_comment"] = comment_body
|
||||
issue["triggering_comment_id"] = data.get("id", "")
|
||||
comment_user = data.get("user", {})
|
||||
if comment_user:
|
||||
issue["comment_author"] = comment_user
|
||||
|
||||
logger.info(
|
||||
"Accepted webhook for issue '%s' (%s), scheduling background task",
|
||||
issue.get("title"),
|
||||
issue.get("id"),
|
||||
)
|
||||
background_tasks.add_task(process_linear_issue, issue, repo_config)
|
||||
|
||||
return {
|
||||
"status": "accepted",
|
||||
"message": f"Processing issue '{issue.get('title')}' for repo {repo_owner}/{repo_name}",
|
||||
}
|
||||
|
||||
|
||||
@app.get("/webhooks/linear")
|
||||
async def linear_webhook_verify() -> dict[str, str]:
|
||||
"""Verify endpoint for Linear webhook setup."""
|
||||
return {"status": "ok", "message": "Linear webhook endpoint is active"}
|
||||
|
||||
|
||||
@app.get("/health")
|
||||
async def health_check() -> dict[str, str]:
|
||||
"""Health check endpoint."""
|
||||
return {"status": "healthy"}
|
||||
11
apps/agent/langgraph.json
Normal file
11
apps/agent/langgraph.json
Normal file
|
|
@ -0,0 +1,11 @@
|
|||
{
|
||||
"graphs": {
|
||||
"agent": "agent.server:get_agent"
|
||||
},
|
||||
"dependencies": ["."],
|
||||
"http": {
|
||||
"app": "agent.webapp:app"
|
||||
},
|
||||
"env": ".env"
|
||||
|
||||
}
|
||||
70
apps/agent/pyproject.toml
Normal file
70
apps/agent/pyproject.toml
Normal file
|
|
@ -0,0 +1,70 @@
|
|||
[project]
|
||||
name = "open-swe-agent"
|
||||
version = "0.1.0"
|
||||
description = "Open SWE Agent - Python agent for automating software engineering tasks"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.11"
|
||||
license = { text = "MIT" }
|
||||
dependencies = [
|
||||
# Core deepagents CLI - points to branch with LangSmith sandbox and working_dir support
|
||||
"deepagents-cli @ git+https://github.com/langchain-ai/deepagents.git@yogesh/working_dir_and_template_config#subdirectory=libs/cli",
|
||||
|
||||
# FastAPI for webhook handling
|
||||
"fastapi>=0.104.0",
|
||||
"uvicorn>=0.24.0",
|
||||
|
||||
# HTTP client
|
||||
"httpx>=0.25.0",
|
||||
|
||||
# JWT for service authentication
|
||||
"PyJWT>=2.8.0",
|
||||
|
||||
# Encryption
|
||||
"cryptography>=41.0.0",
|
||||
|
||||
# LangGraph SDK for thread management
|
||||
"langgraph-sdk>=0.1.0",
|
||||
|
||||
# LangChain dependencies (will be pulled in by deepagents-cli but listing for clarity)
|
||||
"langchain>=0.2.0",
|
||||
"langgraph>=0.2.0",
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
dev = [
|
||||
"pytest>=7.0.0",
|
||||
"pytest-asyncio>=0.21.0",
|
||||
"ruff>=0.1.0",
|
||||
]
|
||||
|
||||
[build-system]
|
||||
requires = ["hatchling"]
|
||||
build-backend = "hatchling.build"
|
||||
|
||||
[tool.hatch.metadata]
|
||||
allow-direct-references = true
|
||||
|
||||
[tool.hatch.build.targets.wheel]
|
||||
packages = ["agent"]
|
||||
|
||||
[tool.ruff]
|
||||
line-length = 100
|
||||
target-version = "py311"
|
||||
|
||||
[tool.ruff.lint]
|
||||
select = [
|
||||
"E", # pycodestyle errors
|
||||
"W", # pycodestyle warnings
|
||||
"F", # Pyflakes
|
||||
"I", # isort
|
||||
"B", # flake8-bugbear
|
||||
"C4", # flake8-comprehensions
|
||||
"UP", # pyupgrade
|
||||
]
|
||||
ignore = [
|
||||
"E501", # line too long (handled by formatter)
|
||||
]
|
||||
|
||||
[tool.pytest.ini_options]
|
||||
asyncio_mode = "auto"
|
||||
testpaths = ["tests"]
|
||||
3209
apps/agent/uv.lock
generated
Normal file
3209
apps/agent/uv.lock
generated
Normal file
File diff suppressed because it is too large
Load diff
Loading…
Add table
Reference in a new issue