From bec99b106ea97164ed2dd6bccac993a31ad41955 Mon Sep 17 00:00:00 2001 From: aran-yogesh Date: Thu, 19 Feb 2026 14:06:06 -0800 Subject: [PATCH 1/6] feat: add multimodal image support for Linear comments and queued run --- .../agent/middleware/check_message_queue.py | 58 ++++++++- apps/agent/agent/utils/multimodal.py | 86 +++++++++++++ apps/agent/agent/webapp.py | 115 ++++++++++++++++-- 3 files changed, 242 insertions(+), 17 deletions(-) create mode 100644 apps/agent/agent/utils/multimodal.py diff --git a/apps/agent/agent/middleware/check_message_queue.py b/apps/agent/agent/middleware/check_message_queue.py index b85248fd..89cc49a3 100644 --- a/apps/agent/agent/middleware/check_message_queue.py +++ b/apps/agent/agent/middleware/check_message_queue.py @@ -8,12 +8,16 @@ human messages before the next model call. from __future__ import annotations import logging +import os from typing import Any +import httpx from langchain.agents.middleware import AgentState, before_model from langgraph.config import get_config, get_store from langgraph.runtime import Runtime +from ..utils.multimodal import fetch_image_block + logger = logging.getLogger(__name__) @@ -80,11 +84,55 @@ async def check_message_queue_before_model( # noqa: PLR0911 thread_id, ) - content_blocks = [ - {"type": "text", "text": msg.get("content", "")} - for msg in queued_messages - if msg.get("content") - ] + async def _build_blocks_from_payload( + payload: dict[str, Any], + ) -> list[dict[str, Any]]: + text = payload.get("text", "") + image_urls = payload.get("image_urls", []) or [] + blocks: list[dict[str, Any]] = [] + if text: + blocks.append({"type": "text", "text": text}) + + if not image_urls: + return blocks + linear_api_key = os.environ.get("LINEAR_API_KEY", "") + async with httpx.AsyncClient() as client: + for image_url in image_urls: + headers = None + if "uploads.linear.app" in image_url: + if linear_api_key: + headers = {"Authorization": linear_api_key} + else: + logger.warning( + "LINEAR_API_KEY not set; cannot authenticate image fetch for %s", + image_url, + ) + image_block = await fetch_image_block( + image_url, client, headers=headers + ) + if image_block: + blocks.append(image_block) + return blocks + + content_blocks: list[dict[str, Any]] = [] + for msg in queued_messages: + content = msg.get("content") + if isinstance(content, dict) and ( + "text" in content or "image_urls" in content + ): + logger.debug("Queued message contains text + image URLs") + blocks = await _build_blocks_from_payload(content) + content_blocks.extend(blocks) + continue + if isinstance(content, list): + logger.debug( + "Queued message contains %d content block(s)", len(content) + ) + content_blocks.extend(content) + continue + if isinstance(content, str) and content: + logger.debug("Queued message contains text content") + content_blocks.append({"type": "text", "text": content}) if not content_blocks: return None diff --git a/apps/agent/agent/utils/multimodal.py b/apps/agent/agent/utils/multimodal.py new file mode 100644 index 00000000..a8cb05cb --- /dev/null +++ b/apps/agent/agent/utils/multimodal.py @@ -0,0 +1,86 @@ +"""Utilities for building multimodal content blocks.""" + +from __future__ import annotations + +import base64 +import logging +import mimetypes +import re +from typing import Any + +import httpx +from langchain_core.messages.content import create_image_block + +logger = logging.getLogger(__name__) + +IMAGE_MARKDOWN_RE = re.compile(r"!\[[^\]]*\]\((https?://[^\s)]+)\)") +IMAGE_URL_RE = re.compile( + r"(https?://[^\s)]+\.(?:png|jpe?g|gif|webp|bmp|tiff)(?:\?[^\s)]+)?)", + re.IGNORECASE, +) + + +def extract_image_urls(text: str) -> list[str]: + """Extract image URLs from markdown image syntax and direct image links.""" + if not text: + return [] + + urls: list[str] = [] + urls.extend(IMAGE_MARKDOWN_RE.findall(text)) + urls.extend(IMAGE_URL_RE.findall(text)) + + deduped = _dedupe_urls(urls) + if deduped: + logger.debug("Extracted %d image URL(s)", len(deduped)) + return deduped + + +def strip_image_urls(text: str, image_urls: list[str]) -> str: + """Remove markdown image links and direct image URLs from text.""" + if not text or not image_urls: + return text + + logger.debug("Stripping %d image URL(s) from text", len(image_urls)) + text = IMAGE_MARKDOWN_RE.sub("", text) + for url in image_urls: + text = text.replace(url, "") + return text + + +async def fetch_image_block( + image_url: str, + client: httpx.AsyncClient, + headers: dict[str, str] | None = None, +) -> dict[str, Any] | None: + """Fetch image bytes and build an image content block.""" + try: + logger.debug("Fetching image from %s", image_url) + response = await client.get(image_url, headers=headers) + response.raise_for_status() + content_type = response.headers.get("Content-Type", "").split(";")[0].strip() + if not content_type: + guessed, _ = mimetypes.guess_type(image_url) + content_type = guessed or "application/octet-stream" + + encoded = base64.b64encode(response.content).decode("ascii") + logger.info( + "Fetched image %s (%s, %d bytes)", + image_url, + content_type, + len(response.content), + ) + return create_image_block(base64=encoded, mime_type=content_type) + except Exception: + logger.exception("Failed to fetch image from %s", image_url) + return None + + +def _dedupe_urls(urls: list[str]) -> list[str]: + seen: set[str] = set() + deduped: list[str] = [] + for url in urls: + if url in seen: + continue + seen.add(url) + deduped.append(url) + return deduped diff --git a/apps/agent/agent/webapp.py b/apps/agent/agent/webapp.py index de40a6da..c3b308c5 100644 --- a/apps/agent/agent/webapp.py +++ b/apps/agent/agent/webapp.py @@ -15,6 +15,7 @@ from langgraph_sdk import get_client # Local import for encryption from .encryption import encrypt_token +from .utils.multimodal import extract_image_urls, fetch_image_block, strip_image_urls logger = logging.getLogger(__name__) @@ -36,6 +37,7 @@ 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: @@ -131,6 +133,8 @@ def get_repo_config_from_team_mapping( return {"owner": "langchain-ai", "name": "langchainplus"} + + 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. @@ -426,7 +430,9 @@ async def is_thread_active(thread_id: str) -> bool: return status == "busy" -async def queue_message_for_thread(thread_id: str, message_content: str) -> bool: +async def queue_message_for_thread( + thread_id: str, message_content: str | list[dict[str, Any]] | dict[str, Any] +) -> bool: """Queue a message for a thread that is currently active. Stores the message in the langgraph store, namespaced to the thread. @@ -435,7 +441,7 @@ async def queue_message_for_thread(thread_id: str, message_content: str) -> bool Args: thread_id: The LangGraph thread ID - message_content: The message content to queue + message_content: The message content to queue (text or content blocks) Returns: True if successfully queued, False otherwise @@ -571,9 +577,20 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 title = full_issue.get("title", "No title") description = full_issue.get("description") or "No description" + image_urls: list[str] = [] + description_image_urls = extract_image_urls(description) + if description_image_urls: + image_urls.extend(description_image_urls) + description = strip_image_urls(description, description_image_urls) + logger.debug( + "Found %d image URL(s) in issue description", + len(description_image_urls), + ) comments = full_issue.get("comments", {}).get("nodes", []) comments_text = "" + triggering_comment = issue_data.get("triggering_comment", "") + triggering_comment_id = issue_data.get("triggering_comment_id", "") bot_message_prefixes = ( "🔐 **GitHub Authentication Required**", @@ -585,32 +602,77 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 "❌ **Agent Error**", ) + comment_ids: set[str] = set() + comment_id_to_index: dict[str, int] = {} if comments: last_bot_comment_idx = -1 for i, comment in enumerate(comments): + comment_id = comment.get("id", "") + if comment_id: + comment_ids.add(comment_id) + comment_id_to_index[comment_id] = i 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 + trigger_index = None + if triggering_comment_id: + trigger_index = comment_id_to_index.get(triggering_comment_id) + if trigger_index is not None: + relevant_comments = comments[trigger_index:] + logger.debug( + "Using triggering comment index %d to build relevant comments", + trigger_index, + ) + else: + 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", "") + body_image_urls = extract_image_urls(body) + if body_image_urls: + image_urls.extend(body_image_urls) + body = strip_image_urls(body, body_image_urls) + logger.debug( + "Found %d image URL(s) in comment by %s", + len(body_image_urls), + author, + ) if any(body.startswith(prefix) for prefix in bot_message_prefixes): continue comments_text += f"\n**{author}:** {body}\n" + if triggering_comment and triggering_comment_id not in comment_ids: + if not comments_text: + comments_text = "\n\n## Comments:\n" + trigger_author = comment_author.get("name", "Unknown") + trigger_body = triggering_comment + trigger_image_urls = extract_image_urls(trigger_body) + if trigger_image_urls: + image_urls.extend(trigger_image_urls) + trigger_body = strip_image_urls(trigger_body, trigger_image_urls) + logger.debug( + "Found %d image URL(s) in triggering comment by %s", + len(trigger_image_urls), + trigger_author, + ) + comments_text += f"\n**{trigger_author}:** {trigger_body}\n" + logger.debug( + "Appended triggering comment %s not present in issue comments list", + triggering_comment_id or "", + ) + prompt = ( f"Please work on the following issue:\n\n" f"## Title: {title}\n\n" @@ -619,6 +681,34 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 "Please analyze this issue and implement the necessary changes. " "When you're done, commit and push your changes." ) + content_blocks: list[dict[str, Any]] = [{"type": "text", "text": prompt}] + if image_urls: + seen_urls: set[str] = set() + deduped_urls: list[str] = [] + for url in image_urls: + if url in seen_urls: + continue + seen_urls.add(url) + deduped_urls.append(url) + image_urls = deduped_urls + logger.info("Preparing %d image(s) for multimodal content", len(image_urls)) + logger.debug("Image URLs: %s", image_urls) + + async with httpx.AsyncClient() as client: + for image_url in image_urls: + headers = None + if "uploads.linear.app" in image_url: + if LINEAR_API_KEY: + headers = {"Authorization": LINEAR_API_KEY} + else: + logger.warning( + "LINEAR_API_KEY not set; cannot authenticate image fetch for %s", + image_url, + ) + image_block = await fetch_image_block(image_url, client, headers=headers) + if image_block: + content_blocks.append(image_block) + logger.info("Built %d content block(s) for prompt", len(content_blocks)) identifier = full_issue.get("identifier", "") or issue_data.get("identifier", "") linear_project_id = "" @@ -652,9 +742,10 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 thread_id, ) + queued_payload = {"text": prompt, "image_urls": image_urls} queued = await queue_message_for_thread( thread_id=thread_id, - message_content=prompt, + message_content=queued_payload, ) if queued: @@ -669,7 +760,7 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 await langgraph_client.runs.create( thread_id, "agent", - input={"messages": [{"role": "user", "content": prompt}]}, + input={"messages": [{"role": "user", "content": content_blocks}]}, config={"configurable": configurable}, if_not_exists="create", ) From 446b659c543d3c966bb4dc670f3d322857f4a7cd Mon Sep 17 00:00:00 2001 From: aran-yogesh Date: Sun, 22 Feb 2026 18:57:49 -0800 Subject: [PATCH 2/6] refactor: keep image URLs in text and centralize image auth --- .../agent/middleware/check_message_queue.py | 52 ++++++++----------- apps/agent/agent/utils/multimodal.py | 26 +++++----- apps/agent/agent/webapp.py | 18 ++----- 3 files changed, 40 insertions(+), 56 deletions(-) diff --git a/apps/agent/agent/middleware/check_message_queue.py b/apps/agent/agent/middleware/check_message_queue.py index 89cc49a3..66178f83 100644 --- a/apps/agent/agent/middleware/check_message_queue.py +++ b/apps/agent/agent/middleware/check_message_queue.py @@ -27,6 +27,28 @@ class LinearNotifyState(AgentState): linear_messages_sent_count: int +async def _build_blocks_from_payload( + payload: dict[str, Any], +) -> list[dict[str, Any]]: + text = payload.get("text", "") + image_urls = payload.get("image_urls", []) or [] + blocks: list[dict[str, Any]] = [] + if text: + blocks.append({"type": "text", "text": text}) + + if not image_urls: + return blocks + linear_api_key = os.environ.get("LINEAR_API_KEY", "") + async with httpx.AsyncClient() as client: + for image_url in image_urls: + image_block = await fetch_image_block( + image_url, client, linear_api_key=linear_api_key + ) + if image_block: + blocks.append(image_block) + return blocks + + @before_model(state_schema=LinearNotifyState) async def check_message_queue_before_model( # noqa: PLR0911 state: LinearNotifyState, # noqa: ARG001 @@ -84,36 +106,6 @@ async def check_message_queue_before_model( # noqa: PLR0911 thread_id, ) - async def _build_blocks_from_payload( - payload: dict[str, Any], - ) -> list[dict[str, Any]]: - text = payload.get("text", "") - image_urls = payload.get("image_urls", []) or [] - blocks: list[dict[str, Any]] = [] - if text: - blocks.append({"type": "text", "text": text}) - - if not image_urls: - return blocks - linear_api_key = os.environ.get("LINEAR_API_KEY", "") - async with httpx.AsyncClient() as client: - for image_url in image_urls: - headers = None - if "uploads.linear.app" in image_url: - if linear_api_key: - headers = {"Authorization": linear_api_key} - else: - logger.warning( - "LINEAR_API_KEY not set; cannot authenticate image fetch for %s", - image_url, - ) - image_block = await fetch_image_block( - image_url, client, headers=headers - ) - if image_block: - blocks.append(image_block) - return blocks - content_blocks: list[dict[str, Any]] = [] for msg in queued_messages: content = msg.get("content") diff --git a/apps/agent/agent/utils/multimodal.py b/apps/agent/agent/utils/multimodal.py index a8cb05cb..4443cc78 100644 --- a/apps/agent/agent/utils/multimodal.py +++ b/apps/agent/agent/utils/multimodal.py @@ -5,6 +5,7 @@ from __future__ import annotations import base64 import logging import mimetypes +import os import re from typing import Any @@ -35,26 +36,27 @@ def extract_image_urls(text: str) -> list[str]: return deduped -def strip_image_urls(text: str, image_urls: list[str]) -> str: - """Remove markdown image links and direct image URLs from text.""" - if not text or not image_urls: - return text - - logger.debug("Stripping %d image URL(s) from text", len(image_urls)) - text = IMAGE_MARKDOWN_RE.sub("", text) - for url in image_urls: - text = text.replace(url, "") - return text - async def fetch_image_block( image_url: str, client: httpx.AsyncClient, - headers: dict[str, str] | None = None, + *, + linear_api_key: str | None = None, ) -> dict[str, Any] | None: """Fetch image bytes and build an image content block.""" try: logger.debug("Fetching image from %s", image_url) + headers = None + if "uploads.linear.app" in image_url: + if linear_api_key is None: + linear_api_key = os.environ.get("LINEAR_API_KEY", "") + if linear_api_key: + headers = {"Authorization": linear_api_key} + else: + logger.warning( + "LINEAR_API_KEY not set; cannot authenticate image fetch for %s", + image_url, + ) response = await client.get(image_url, headers=headers) response.raise_for_status() content_type = response.headers.get("Content-Type", "").split(";")[0].strip() diff --git a/apps/agent/agent/webapp.py b/apps/agent/agent/webapp.py index c3b308c5..1b6c0270 100644 --- a/apps/agent/agent/webapp.py +++ b/apps/agent/agent/webapp.py @@ -15,7 +15,7 @@ from langgraph_sdk import get_client # Local import for encryption from .encryption import encrypt_token -from .utils.multimodal import extract_image_urls, fetch_image_block, strip_image_urls +from .utils.multimodal import extract_image_urls, fetch_image_block logger = logging.getLogger(__name__) @@ -581,7 +581,6 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 description_image_urls = extract_image_urls(description) if description_image_urls: image_urls.extend(description_image_urls) - description = strip_image_urls(description, description_image_urls) logger.debug( "Found %d image URL(s) in issue description", len(description_image_urls), @@ -643,7 +642,6 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 body_image_urls = extract_image_urls(body) if body_image_urls: image_urls.extend(body_image_urls) - body = strip_image_urls(body, body_image_urls) logger.debug( "Found %d image URL(s) in comment by %s", len(body_image_urls), @@ -661,7 +659,6 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 trigger_image_urls = extract_image_urls(trigger_body) if trigger_image_urls: image_urls.extend(trigger_image_urls) - trigger_body = strip_image_urls(trigger_body, trigger_image_urls) logger.debug( "Found %d image URL(s) in triggering comment by %s", len(trigger_image_urls), @@ -696,16 +693,9 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 async with httpx.AsyncClient() as client: for image_url in image_urls: - headers = None - if "uploads.linear.app" in image_url: - if LINEAR_API_KEY: - headers = {"Authorization": LINEAR_API_KEY} - else: - logger.warning( - "LINEAR_API_KEY not set; cannot authenticate image fetch for %s", - image_url, - ) - image_block = await fetch_image_block(image_url, client, headers=headers) + image_block = await fetch_image_block( + image_url, client, linear_api_key=LINEAR_API_KEY + ) if image_block: content_blocks.append(image_block) logger.info("Built %d content block(s) for prompt", len(content_blocks)) From 69e1eeb36f1f5f0406c3db6224553ad64c9369fb Mon Sep 17 00:00:00 2001 From: aran-yogesh Date: Mon, 23 Feb 2026 14:04:39 -0800 Subject: [PATCH 3/6] test: add unit tests for extract_image_urls function --- apps/agent/tests/test_multimodal.py | 101 ++++++++++++++++++++++++++++ 1 file changed, 101 insertions(+) create mode 100644 apps/agent/tests/test_multimodal.py diff --git a/apps/agent/tests/test_multimodal.py b/apps/agent/tests/test_multimodal.py new file mode 100644 index 00000000..66cc6129 --- /dev/null +++ b/apps/agent/tests/test_multimodal.py @@ -0,0 +1,101 @@ +from __future__ import annotations + +from agent.utils.multimodal import extract_image_urls + + +def test_extract_image_urls_empty() -> None: + assert extract_image_urls("") == [] + + +def test_extract_image_urls_markdown_and_direct_dedupes() -> None: + text = ( + "Here is an image ![alt](https://example.com/a.png) and another " + "![https://example.com/b.JPG?size=large plus a repeat https://example.com/a.png" + ) + + assert extract_image_urls(text) == [ + "https://example.com/a.png", + "https://example.com/b.JPG?size=large", + ] + + +def test_extract_image_urls_ignores_non_images() -> None: + text = "Not images: https://example.com/file.pdf and https://example.com/noext" + + assert extract_image_urls(text) == [] + + +def test_extract_image_urls_markdown_syntax() -> None: + text = "Check out this screenshot: ![Screenshot](https://example.com/screenshot.png)" + + assert extract_image_urls(text) == ["https://example.com/screenshot.png"] + + +def test_extract_image_urls_direct_links() -> None: + text = "Direct link: https://example.com/photo.jpg and another https://example.com/image.gif" + + assert extract_image_urls(text) == [ + "https://example.com/photo.jpg", + "https://example.com/image.gif", + ] + + +def test_extract_image_urls_various_formats() -> None: + text = ( + "Multiple formats: " + "https://example.com/image.png " + "https://example.com/photo.jpeg " + "https://example.com/pic.gif " + "https://example.com/img.webp " + "https://example.com/bitmap.bmp " + "https://example.com/scan.tiff" + ) + + assert extract_image_urls(text) == [ + "https://example.com/image.png", + "https://example.com/photo.jpeg", + "https://example.com/pic.gif", + "https://example.com/img.webp", + "https://example.com/bitmap.bmp", + "https://example.com/scan.tiff", + ] + + +def test_extract_image_urls_with_query_params() -> None: + text = "Image with params: https://cdn.example.com/image.png?width=800&height=600" + + assert extract_image_urls(text) == ["https://cdn.example.com/image.png?width=800&height=600"] + + +def test_extract_image_urls_case_insensitive() -> None: + text = "Mixed case: https://example.com/Image.PNG and https://example.com/photo.JpEg" + + assert extract_image_urls(text) == [ + "https://example.com/Image.PNG", + "https://example.com/photo.JpEg", + ] + + +def test_extract_image_urls_deduplication() -> None: + text = ( + "Same URL twice: https://example.com/image.png " + "and again https://example.com/image.png" + ) + + assert extract_image_urls(text) == ["https://example.com/image.png"] + + +def test_extract_image_urls_mixed_markdown_and_direct() -> None: + text = ( + "Markdown: ![alt text](https://example.com/markdown.png) " + "and direct: https://example.com/direct.jpg " + "and another markdown ![](https://example.com/another.gif)" + ) + + result = extract_image_urls(text) + assert set(result) == { + "https://example.com/markdown.png", + "https://example.com/direct.jpg", + "https://example.com/another.gif", + } + assert len(result) == 3 From fb23ef6c159ccbcd83749e7ae3cbe0e65844a81a Mon Sep 17 00:00:00 2001 From: Aran Yogesh Date: Tue, 24 Feb 2026 11:21:32 -0800 Subject: [PATCH 4/6] Apply suggestion from @bracesproul Co-authored-by: Brace Sproul --- apps/agent/agent/webapp.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/apps/agent/agent/webapp.py b/apps/agent/agent/webapp.py index 1b6c0270..84c8eee4 100644 --- a/apps/agent/agent/webapp.py +++ b/apps/agent/agent/webapp.py @@ -133,8 +133,6 @@ def get_repo_config_from_team_mapping( return {"owner": "langchain-ai", "name": "langchainplus"} - - 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. From f82408cd612bc1cba21fe53599cb6cf872faf4ba Mon Sep 17 00:00:00 2001 From: aran-yogesh Date: Tue, 24 Feb 2026 12:11:24 -0800 Subject: [PATCH 5/6] refactor: use content block helpers and shared url dedupe; run format and lint --- apps/agent/agent/integrations/langsmith.py | 12 ++------ .../agent/middleware/check_message_queue.py | 14 ++------- apps/agent/agent/middleware/open_pr.py | 8 ++--- apps/agent/agent/server.py | 10 ++----- apps/agent/agent/tools/commit_and_open_pr.py | 4 +-- apps/agent/agent/utils/github.py | 24 ++++----------- apps/agent/agent/utils/multimodal.py | 28 ++++++++---------- apps/agent/agent/webapp.py | 29 +++++++------------ apps/agent/tests/test_multimodal.py | 5 +--- 9 files changed, 41 insertions(+), 93 deletions(-) diff --git a/apps/agent/agent/integrations/langsmith.py b/apps/agent/agent/integrations/langsmith.py index 3bf23afd..b1679f51 100644 --- a/apps/agent/agent/integrations/langsmith.py +++ b/apps/agent/agent/integrations/langsmith.py @@ -9,7 +9,7 @@ import contextlib import os import time from abc import ABC, abstractmethod -from typing import TYPE_CHECKING, Any +from typing import Any from deepagents.backends.protocol import ( ExecuteResponse, @@ -19,7 +19,6 @@ from deepagents.backends.protocol import ( WriteResult, ) from deepagents.backends.sandbox import BaseSandbox - from langsmith.sandbox import Sandbox, SandboxClient, SandboxTemplate @@ -148,9 +147,7 @@ class LangSmithBackend(BaseSandbox): responses: list[FileDownloadResponse] = [] for path in paths: content = self._sandbox.read(path) - responses.append( - FileDownloadResponse(path=path, content=content, error=None) - ) + responses.append(FileDownloadResponse(path=path, content=content, error=None)) return responses def upload_files(self, files: list[tuple[str, bytes]]) -> list[FileUploadResponse]: @@ -209,10 +206,7 @@ class LangSmithProvider(SandboxProvider): template_name=resolved_template_name, timeout=timeout ) except Exception as e: - msg = ( - f"Failed to create sandbox from template " - f"'{resolved_template_name}': {e}" - ) + msg = f"Failed to create sandbox from template '{resolved_template_name}': {e}" raise RuntimeError(msg) from e # Verify sandbox is ready by polling diff --git a/apps/agent/agent/middleware/check_message_queue.py b/apps/agent/agent/middleware/check_message_queue.py index 66178f83..a8825761 100644 --- a/apps/agent/agent/middleware/check_message_queue.py +++ b/apps/agent/agent/middleware/check_message_queue.py @@ -8,7 +8,6 @@ human messages before the next model call. from __future__ import annotations import logging -import os from typing import Any import httpx @@ -38,12 +37,9 @@ async def _build_blocks_from_payload( if not image_urls: return blocks - linear_api_key = os.environ.get("LINEAR_API_KEY", "") async with httpx.AsyncClient() as client: for image_url in image_urls: - image_block = await fetch_image_block( - image_url, client, linear_api_key=linear_api_key - ) + image_block = await fetch_image_block(image_url, client) if image_block: blocks.append(image_block) return blocks @@ -109,17 +105,13 @@ async def check_message_queue_before_model( # noqa: PLR0911 content_blocks: list[dict[str, Any]] = [] for msg in queued_messages: content = msg.get("content") - if isinstance(content, dict) and ( - "text" in content or "image_urls" in content - ): + if isinstance(content, dict) and ("text" in content or "image_urls" in content): logger.debug("Queued message contains text + image URLs") blocks = await _build_blocks_from_payload(content) content_blocks.extend(blocks) continue if isinstance(content, list): - logger.debug( - "Queued message contains %d content block(s)", len(content) - ) + logger.debug("Queued message contains %d content block(s)", len(content)) content_blocks.extend(content) continue if isinstance(content, str) and content: diff --git a/apps/agent/agent/middleware/open_pr.py b/apps/agent/agent/middleware/open_pr.py index 426bed6d..0a4b05d5 100644 --- a/apps/agent/agent/middleware/open_pr.py +++ b/apps/agent/agent/middleware/open_pr.py @@ -181,16 +181,12 @@ I've {action} pull request to address this issue: logger.info("Changes detected, preparing PR for thread %s", thread_id) - current_branch = await asyncio.to_thread( - git_current_branch, sandbox_backend, repo_dir - ) + current_branch = await asyncio.to_thread(git_current_branch, sandbox_backend, repo_dir) target_branch = f"open-swe/{thread_id}" if current_branch != target_branch: - await asyncio.to_thread( - git_checkout_branch, sandbox_backend, repo_dir, target_branch - ) + await asyncio.to_thread(git_checkout_branch, sandbox_backend, repo_dir, target_branch) await asyncio.to_thread( git_config_user, diff --git a/apps/agent/agent/server.py b/apps/agent/agent/server.py index c08746a9..c162642f 100644 --- a/apps/agent/agent/server.py +++ b/apps/agent/agent/server.py @@ -4,7 +4,6 @@ # Suppress deprecation warnings from langchain_core (e.g., Pydantic V1 on Python 3.14+) # ruff: noqa: E402 import logging -import os import warnings logger = logging.getLogger(__name__) @@ -29,6 +28,7 @@ from deepagents.backends.protocol import SandboxBackendProtocol from langchain_anthropic import ChatAnthropic from .encryption import decrypt_token +from .integrations.langsmith import _create_langsmith_sandbox from .middleware import ( ToolErrorMiddleware, check_message_queue_before_model, @@ -37,8 +37,6 @@ from .middleware import ( ) from .prompt import construct_system_prompt from .tools import commit_and_open_pr, fetch_url, http_request -from .integrations.langsmith import _create_langsmith_sandbox - client = get_client() @@ -46,12 +44,12 @@ SANDBOX_CREATING = "__creating__" SANDBOX_CREATION_TIMEOUT = 180 SANDBOX_POLL_INTERVAL = 1.0 -from .utils.sandbox_state import SANDBOX_BACKENDS, get_sandbox_id_from_metadata from .utils.github import ( git_has_uncommitted_changes, is_valid_git_repo, remove_directory, ) +from .utils.sandbox_state import SANDBOX_BACKENDS, get_sandbox_id_from_metadata async def _clone_or_pull_repo_in_sandbox( # noqa: PLR0915 @@ -87,9 +85,7 @@ async def _clone_or_pull_repo_in_sandbox( # noqa: PLR0915 is_git_repo = await loop.run_in_executor(None, is_valid_git_repo, sandbox_backend, repo_dir) if not is_git_repo: - logger.warning( - "Repo directory missing or not a valid git repo at %s, removing", repo_dir - ) + logger.warning("Repo directory missing or not a valid git repo at %s, removing", repo_dir) try: removed = await loop.run_in_executor(None, remove_directory, sandbox_backend, repo_dir) if not removed: diff --git a/apps/agent/agent/tools/commit_and_open_pr.py b/apps/agent/agent/tools/commit_and_open_pr.py index 2afc3615..78a8c1e3 100644 --- a/apps/agent/agent/tools/commit_and_open_pr.py +++ b/apps/agent/agent/tools/commit_and_open_pr.py @@ -181,9 +181,7 @@ def commit_and_open_pr( "pr_url": None, } - base_branch = asyncio.run( - get_github_default_branch(repo_owner, repo_name, github_token) - ) + base_branch = asyncio.run(get_github_default_branch(repo_owner, repo_name, github_token)) pr_url, _pr_number, pr_existing = asyncio.run( create_github_pr( repo_owner=repo_owner, diff --git a/apps/agent/agent/utils/github.py b/apps/agent/agent/utils/github.py index d966aa5c..b09c6325 100644 --- a/apps/agent/agent/utils/github.py +++ b/apps/agent/agent/utils/github.py @@ -37,24 +37,18 @@ def remove_directory(sandbox_backend: SandboxBackendProtocol, repo_dir: str) -> return result.exit_code == 0 -def git_has_uncommitted_changes( - sandbox_backend: SandboxBackendProtocol, repo_dir: str -) -> bool: +def git_has_uncommitted_changes(sandbox_backend: SandboxBackendProtocol, repo_dir: str) -> bool: """Check whether the repo has uncommitted changes.""" result = _run_git(sandbox_backend, repo_dir, "git status --porcelain") return result.exit_code == 0 and bool(result.output.strip()) -def git_fetch_origin( - sandbox_backend: SandboxBackendProtocol, repo_dir: str -) -> ExecuteResponse: +def git_fetch_origin(sandbox_backend: SandboxBackendProtocol, repo_dir: str) -> ExecuteResponse: """Fetch latest from origin (best-effort).""" return _run_git(sandbox_backend, repo_dir, "git fetch origin 2>/dev/null || true") -def git_has_unpushed_commits( - sandbox_backend: SandboxBackendProtocol, repo_dir: str -) -> bool: +def git_has_unpushed_commits(sandbox_backend: SandboxBackendProtocol, repo_dir: str) -> bool: """Check whether there are commits not pushed to upstream.""" git_log_cmd = ( "git log --oneline @{upstream}..HEAD 2>/dev/null " @@ -75,9 +69,7 @@ def git_checkout_branch( ) -> bool: """Checkout branch, creating it if needed.""" safe_branch = shlex.quote(branch) - checkout_result = _run_git( - sandbox_backend, repo_dir, f"git checkout -b {safe_branch}" - ) + checkout_result = _run_git(sandbox_backend, repo_dir, f"git checkout -b {safe_branch}") if checkout_result.exit_code == 0: return True fallback = _run_git(sandbox_backend, repo_dir, f"git checkout {safe_branch}") @@ -97,9 +89,7 @@ def git_config_user( _run_git(sandbox_backend, repo_dir, f"git config user.email {safe_email}") -def git_add_all( - sandbox_backend: SandboxBackendProtocol, repo_dir: str -) -> ExecuteResponse: +def git_add_all(sandbox_backend: SandboxBackendProtocol, repo_dir: str) -> ExecuteResponse: """Stage all changes.""" return _run_git(sandbox_backend, repo_dir, "git add -A") @@ -112,9 +102,7 @@ def git_commit( return _run_git(sandbox_backend, repo_dir, f"git commit -m {safe_message}") -def git_get_remote_url( - sandbox_backend: SandboxBackendProtocol, repo_dir: str -) -> str | None: +def git_get_remote_url(sandbox_backend: SandboxBackendProtocol, repo_dir: str) -> str | None: """Get the origin remote URL.""" result = _run_git(sandbox_backend, repo_dir, "git remote get-url origin") if result.exit_code != 0: diff --git a/apps/agent/agent/utils/multimodal.py b/apps/agent/agent/utils/multimodal.py index 4443cc78..a7e83cd5 100644 --- a/apps/agent/agent/utils/multimodal.py +++ b/apps/agent/agent/utils/multimodal.py @@ -30,26 +30,22 @@ def extract_image_urls(text: str) -> list[str]: urls.extend(IMAGE_MARKDOWN_RE.findall(text)) urls.extend(IMAGE_URL_RE.findall(text)) - deduped = _dedupe_urls(urls) + deduped = dedupe_urls(urls) if deduped: logger.debug("Extracted %d image URL(s)", len(deduped)) return deduped - async def fetch_image_block( image_url: str, client: httpx.AsyncClient, - *, - linear_api_key: str | None = None, ) -> dict[str, Any] | None: """Fetch image bytes and build an image content block.""" try: logger.debug("Fetching image from %s", image_url) headers = None if "uploads.linear.app" in image_url: - if linear_api_key is None: - linear_api_key = os.environ.get("LINEAR_API_KEY", "") + linear_api_key = os.environ.get("LINEAR_API_KEY", "") if linear_api_key: headers = {"Authorization": linear_api_key} else: @@ -62,7 +58,13 @@ async def fetch_image_block( content_type = response.headers.get("Content-Type", "").split(";")[0].strip() if not content_type: guessed, _ = mimetypes.guess_type(image_url) - content_type = guessed or "application/octet-stream" + if not guessed: + logger.warning( + "Could not determine content type for %s; skipping image", + image_url, + ) + return None + content_type = guessed encoded = base64.b64encode(response.content).decode("ascii") logger.info( @@ -77,12 +79,6 @@ async def fetch_image_block( return None -def _dedupe_urls(urls: list[str]) -> list[str]: - seen: set[str] = set() - deduped: list[str] = [] - for url in urls: - if url in seen: - continue - seen.add(url) - deduped.append(url) - return deduped +def dedupe_urls(urls: list[str]) -> list[str]: + deduped: set[str] = set(urls) + return list(deduped) diff --git a/apps/agent/agent/webapp.py b/apps/agent/agent/webapp.py index 84c8eee4..1949b0e6 100644 --- a/apps/agent/agent/webapp.py +++ b/apps/agent/agent/webapp.py @@ -11,11 +11,12 @@ from typing import Any import httpx import jwt from fastapi import BackgroundTasks, FastAPI, HTTPException, Request +from langchain_core.messages.content import create_text_block from langgraph_sdk import get_client # Local import for encryption from .encryption import encrypt_token -from .utils.multimodal import extract_image_urls, fetch_image_block +from .utils.multimodal import dedupe_urls, extract_image_urls, fetch_image_block logger = logging.getLogger(__name__) @@ -37,7 +38,6 @@ 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: @@ -78,9 +78,11 @@ LINEAR_TEAM_TO_REPO: dict[str, dict[str, Any] | dict[str, str]] = { "open-swe-v3-test": {"owner": "aran-yogesh", "name": "nimedge"}, "open-swe-dev-test": {"owner": "aran-yogesh", "name": "TalkBack"}, }, - "default": {"owner": "aran-yogesh", "name": "TalkBack"} # Fallback for issues without project + "default": { + "owner": "aran-yogesh", + "name": "TalkBack", + }, # Fallback for issues without project }, - "LangChain OSS": { "projects": { "deepagents": {"owner": "langchain-ai", "name": "deepagents"}, @@ -93,9 +95,7 @@ LINEAR_TEAM_TO_REPO: dict[str, dict[str, Any] | dict[str, str]] = { }, "default": {"owner": "langchain-ai", "name": "ai-sdr"}, }, - "Docs": { - "default": {"owner": "langchain-ai", "name": "docs"} - }, + "Docs": {"default": {"owner": "langchain-ai", "name": "docs"}}, } @@ -676,24 +676,15 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 "Please analyze this issue and implement the necessary changes. " "When you're done, commit and push your changes." ) - content_blocks: list[dict[str, Any]] = [{"type": "text", "text": prompt}] + content_blocks: list[dict[str, Any]] = [create_text_block(prompt)] if image_urls: - seen_urls: set[str] = set() - deduped_urls: list[str] = [] - for url in image_urls: - if url in seen_urls: - continue - seen_urls.add(url) - deduped_urls.append(url) - image_urls = deduped_urls + image_urls = dedupe_urls(image_urls) logger.info("Preparing %d image(s) for multimodal content", len(image_urls)) logger.debug("Image URLs: %s", image_urls) async with httpx.AsyncClient() as client: for image_url in image_urls: - image_block = await fetch_image_block( - image_url, client, linear_api_key=LINEAR_API_KEY - ) + image_block = await fetch_image_block(image_url, client) if image_block: content_blocks.append(image_block) logger.info("Built %d content block(s) for prompt", len(content_blocks)) diff --git a/apps/agent/tests/test_multimodal.py b/apps/agent/tests/test_multimodal.py index 66cc6129..5dca4d39 100644 --- a/apps/agent/tests/test_multimodal.py +++ b/apps/agent/tests/test_multimodal.py @@ -77,10 +77,7 @@ def test_extract_image_urls_case_insensitive() -> None: def test_extract_image_urls_deduplication() -> None: - text = ( - "Same URL twice: https://example.com/image.png " - "and again https://example.com/image.png" - ) + text = "Same URL twice: https://example.com/image.png and again https://example.com/image.png" assert extract_image_urls(text) == ["https://example.com/image.png"] From 4d114d2e29c3d5388168dcd96bed8a1c0f62b173 Mon Sep 17 00:00:00 2001 From: aran-yogesh Date: Tue, 24 Feb 2026 12:36:09 -0800 Subject: [PATCH 6/6] fix: preserve image url order when deduping --- apps/agent/agent/utils/multimodal.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/apps/agent/agent/utils/multimodal.py b/apps/agent/agent/utils/multimodal.py index a7e83cd5..bc2b7ff6 100644 --- a/apps/agent/agent/utils/multimodal.py +++ b/apps/agent/agent/utils/multimodal.py @@ -80,5 +80,4 @@ async def fetch_image_block( def dedupe_urls(urls: list[str]) -> list[str]: - deduped: set[str] = set(urls) - return list(deduped) + return list(dict.fromkeys(urls))