mirror of
https://github.com/agentscope-ai/ReMe.git
synced 2026-08-28 05:25:04 +00:00
Some checks are pending
Pre-commit / run (ubuntu-latest) (push) Waiting to run
Tests ReMe / Unit Tests - py3.13 (push) Waiting to run
Tests ReMe / Unit Tests - py3.11 (push) Waiting to run
Tests ReMe / Unit Tests - py3.12 (push) Waiting to run
Windows Smoke / CLI smoke - py3.11 (push) Waiting to run
* feat: add ssh proxy * feat: add ssh proxy * feat: add ssh proxy * feat: add ssh proxy * feat: add prompt * feat: add agent wrapper * feat: add agent wrapper * feat: add agent wrapper * feat: add tushare skill * feat: add tushare skill * feat: add tushare skill * feat: add none stream * chore(deps): update dependency versions in pyproject.toml - Bump claude-agent-sdk from 0.2.123 to 0.2.126 - Upgrade pre-commit to version 4.6.1 or higher - Upgrade pytest to version 9.1.1 or higher * feat(agent_wrapper): add session compaction support and unify session commands - Introduce compact_session method to BaseAgentWrapper and implement it in AsAgentWrapper, CcAgentWrapper, and CodexAgentWrapper - Add session_command module with SessionCommandResult dataclass and handle_session_command function for /clear and /compact commands - Update __init__.py exports to include session_command handlers - Modify DingTalkWaitStep to handle session commands via handle_session_command function - Remove streaming mode from DingTalkWaitStep and simplify reply handling to final Markdown replies only - Add unit tests for session compaction methods and session command handling across wrappers and DingTalk integration - Clean up and remove obsolete streaming and card rendering code from DingTalk wait step - Adjust daily_cookbook.yaml to remove stream and card_update_interval config entries for DingTalk wait step * feat(auto_fin): add Auto Fin simulated portfolio cookbook workflow - Add comprehensive Auto Fin schema exports for multiple models and enums - Implement base class and helpers for Auto Fin analysis steps - Create file, state, and formatting utilities for Auto Fin with atomic file writes and locking - Define Auto Fin pipeline with four analysis agents: backtest, event, portfolio, and US correlation - Register Auto Fin package in cookbook workflows and schema initialization - Add detailed documentation in markdown describing the system design, workflow, and data contracts * feat(outbound_proxy): add application-scoped outbound HTTP proxy components - Introduce BaseOutboundProxy and OutboundProxyEndpoint as core contracts - Implement FixedHttpOutboundProxy for external HTTP proxy integration - Add SshHttpOutboundProxy providing SSH-backed local HTTP proxy tunnels - Register outbound proxy components in component registry and enumeration - Update components package to include outbound_proxy module - Add dependency on pproxy for SSH HTTP proxy bridging - Include comprehensive unit tests covering proxy lifecycle, validation, environment merging, error handling, readiness, and monitoring mechanisms * refactor(network): replace SSH proxy with explicit HTTP outbound proxy - Remove SSH proxy helper implementation and references in codebase - Add support for explicit HTTP proxy URL in arXiv and HuggingFace clients - Modify clients to use async context manager for consistent resource handling - Update daily paper steps to forward outbound proxy configuration explicitly - Change tests to cover new proxy usage model and remove SSH proxy mocks - Add outbound proxy component configuration in daily_cookbook.yaml - Ensure proxy URL usage disables environment trust in HTTP clients - Fix app context component enum access to be defensive against missing keys * feat(agent_wrapper): add managed proxy support for command environments - Introduce BaseOutboundProxy binding in BaseAgentWrapper for outbound proxy management - Add bash_environment and command_proxy_environment properties to apply proxy settings - Update WorkspaceBackend instantiation in AsAgentWrapper to use bash_environment - Inject managed proxy export commands into Claude Code Bash commands via hooks - Enhance CodexAgentWrapper to include managed proxy in shell environment policy - Modify daily_cookbook.yaml steps to specify outbound_proxy as default where needed - Add comprehensive unit tests verifying managed proxy injection and environment isolation - Ensure subprocess_environment remains unchanged while proxy is applied selectively to commands * refactor(memory): replace search job_tools with memory in daily cookbook config - Change workspace_dir default from .reme to reme_workspace - Replace search job_tools with memory across multiple components and jobs - Update descriptions to reflect long-term memory retrieval instead of search - Modify system prompts to instruct using memory for retrieving notes - Adjust unit tests to verify memory job_tools and job presence instead of search - Ensure consistency in configuration and tests for memory backend usage * refactor(config): rename memory to memory_search in daily cookbook config - Change all occurrences of "memory" to "memory_search" in job_tools and job definitions - Update related system prompts to reflect the new memory_search terminology - Modify unit tests to assert the presence of memory_search instead of memory - Ensure consistency across skills, job tools, and backend configurations in multiple components * feat(auto_fin): add deterministic quantitative research and ranking fusion - Introduce new schema models: EtfScore, RankingMetrics, ExtremeAnalysis, DimensionRanking, and FusionRanking to represent deterministic research outputs - Add ranking data to event, backtest, us_correlation, and portfolio analysis outputs - Implement ranking_section renderer to format Top20 scores and diagnostics in Markdown - Develop AutoFinQuantStep for deterministic ETF ranking using TuShare data, Polars, and a custom extremely randomized tree ensemble - Integrate quantitative rankings into backtest and portfolio analysis steps and reports - Extend auto_fin pipeline with new quant_enabled and quant_required config options - Enforce ranking constraints like unique codes, contiguous ranks, and normalized fusion weights - Update analysis YAMLs with rules limiting data freshness, universe, and ranking usage - Incorporate ranking outputs into all major markdown report bodies in Auto Fin pipeline - Add concurrency-limited asynchronous TuShare client to fetch required market data - Introduce cross-sectional rank correlation and NDCG metrics for ranking quality evaluation * feat(auto_fin): implement stage-wise notification and reporting for analysis pipeline - Refactor notification config in daily_cookbook.yaml to support dispatch steps - Update AutoFinNotificationStep to deduplicate notifications per run stage - Add _notify_stage method in pipeline to send notifications for each analysis stage - Implement persistence and notification for event, backtest, US correlation, and portfolio stages - Modify pipeline flow to persist reports and notify after each stage completion - Adjust metadata to track notifications and errors per stage - Update tests to verify stage-wise notification sending and deduplication - Remove older combined report persistence in favor of modular stage handling * feat(auto_fin): add outbound proxy support for Tushare API usage - Introduce BaseOutboundProxy reference in AutoFinPipelineStep and AutoFinQuantStep - Update TushareResearchClient and trade calendar fetch to accept and use proxy URL - Create _ProxiedTushareApi adapter to route Tushare requests via explicit HTTP proxy - Modify create_tushare_api utility to optionally return proxied API client - Add unit tests covering proxy forwarding and client behavior with managed proxies - Ensure proxy usage respects explicit proxy URL over environment fallback - Integrate outbound proxy into data fetching and quantitative research steps * feat(auto_fin): enforce checkpoint time validation and add state models - Introduce AnalysisState base class and specific states for event, backtest, and US correlation analyses - Replace analysis output types with corresponding state classes in run schemas - Add require_checkpoint_reached method to validate decision_at/data_cutoff against current time - Enforce checkpoint time checks before analysis steps in event, backtest, portfolio, and quant analyses - Refactor quant data loading to include adjustment factors and apply price adjustments without fallback - Update analysis YAML docs to require real-time checkpoint validation and forbid using future data - Improve portfolio run serialization by excluding redundant legacy fields and nested proposed actions - Add helper to extract readable sections from persisted checkpoint documents - Fix event analysis output validation to reject events and sources with future timestamps * feat(auto_fin): auto-select latest reached checkpoint if none specified - Extend checkpoint config to accept empty string for auto selection - Add static method to compute latest checkpoint reached by current time - Modify pipeline step to auto-select checkpoint based on trade calendar and time - Adjust force flag default depending on whether checkpoint is explicit or auto - Log details when checkpoint is auto-selected to improve observability - Add comprehensive tests for auto checkpoint selection logic and edge cases - Remove deprecated default and required constraints from force parameter in config * refactor(auto_fin): unify datetime comparison with compare_datetimes utility - Replace direct datetime comparisons with compare_datetimes function calls - Use cmp_to_key with compare_datetimes for sorting datetime tuples and lists - Update validation logic in backtest, event, analysis, and ledger modules for consistent datetime handling - Add unit tests to verify handling of naive and aware datetime comparisons in event and backtest validations - Ensure marked_at and interval_end timestamps are set and compared consistently using compare_datetimes - Improve correctness of ordering and conditional checks related to timestamps throughout auto_fin steps and ledger code * feat(auto_fin): add datetime comparison helper for mixed timezone data - Implement compare_datetimes function to handle naive and aware datetimes - Ensure naive datetime is interpreted in the known timezone of the counterpart - Facilitate comparisons between legacy and timezone-aware Auto Fin data - Add module docstring explaining purpose of the helpers * docs(auto_fin): enforce unique ETF representative per sub-theme in analysis rules - Update backtest.yaml to recommend or highlight only one ETF per sub-theme for ETF analyses - Modify event.yaml to map only one representative ETF per sub-theme, avoiding duplicate recommendations - Revise portfolio.yaml to restrict holdings/buys to a single ETF per sub-theme, preventing repeated buys of highly overlapping ETFs - Adjust us_correlation.yaml to retain only one representative A-share ETF per sub-theme for mapping or recommendation - Add test to verify presence of new sub-theme uniqueness guidance in step prompts * feat(auto_fin): separate draft model and include deterministic fusion ranking - Introduce _PortfolioProposalDraft pydantic model for agent-authored fields before ranking - Discard any "fusion_ranking" data from draft to prevent conflicts with canonical ranking - Modify AutoFinPortfolioStep to receive draft, enrich with fusion_ranking, and produce final output - Update tests to use _PortfolioProposalDraft and validate deterministic fusion ranking propagation - Add async test verifying fusion ranking is correctly set in portfolio output with no errors * refactor(auto_fin): rewrite and simplify Auto Fin schema and steps - Remove legacy Auto Fin analysis step modules and helpers - Replace complex ranking and portfolio models with simplified current-news models - Update schema to focus on news-case workflow with new domain models - Remove A-share decision checkpoints and backtest details from schema - Simplify recommendation and decision output structures - Clean up deprecated state and utility functions - Update Auto Fin steps initialization to new pipeline steps only - Improve uniqueness validation for themes and ETFs in research plan * feat(auto_fin): implement full local cache and analysis workflow for Auto Fin - Add AutoFinDataStep to prepare and cache daily TuShare data with lookback - Add AutoFinAnalysisStep to analyze cached data and generate Markdown report - Implement detailed time window, ETF filtering, and historical case validation - Introduce YAML prompts for planning and decision-making steps - Update .gitignore to include reme_workspace/ - Clean up config and import structure for auto_fin steps - Remove old pipeline.py and consolidate functionality into new modules - Use polars for efficient CSV reading and data processing - Ensure atomic writes and strict JSON serialization for cache files - Enforce rules on news timing, ETF universe, and historical case usage * fix(auto_fin): restrict news data source to '财联社' in analysis and cache - Update analysis templates to specify current news as from '财联社' only - Modify news fetching functions to filter by source '财联社' - Add validation method to check cached news source correctness - Update news caching logic to exclude non-'财联社' news - Enhance unit tests with multiple sources to ensure filtering works - Confirm news API calls include source filter parameter as '财联社' * refactor(auto_fin): convert I/O methods to asynchronous implementations - Change _news, _dataset, and _theme_data methods to async for improved concurrency - Move JSONL and CSV reading operations to asynchronous wrappers using asyncio.to_thread - Remove synchronous _read_jsonl and _read_csv functions, integrate them as static async class methods - Update cache validation methods to async, awaiting I/O operations accordingly - Adjust usage of dataset and news retrieval in analysis step to await asynchronous methods - Add async unit test to validate JSONL reading with unicode line separators - Preserve existing functionality while enabling non-blocking file and data access * fix(nx_file_graph): defer networkx import and improve dependency handling - Move networkx import inside NxFileGraph constructor for lazy loading - Raise ImportError with original exception context if networkx is missing - Remove module-level fallback assignment of nx to None - Expand test to block loading of multiple optional core dependencies eagerly - Change exception type in test from ModuleNotFoundError to AssertionError - Update test comments to reflect broader optional dependency checks * feat(embedding_store): add quota retry delay mechanism for embedding requests - Introduce quota_retry_delay parameter to configure wait time before retry on quota exhaustion - Implement detection of insufficient quota errors in LocalEmbeddingStore without external SDK - Add retry logic with custom delay when quota is insufficient during embedding requests - Update configuration to set max_retries and quota_retry_delay defaults for embedding store - Add unit tests covering quota exhaustion retry behavior with delay and opt-in control - Ensure existing retry behavior remains unchanged if quota_retry_delay is not set * feat(auto_fin): add detailed logging to analysis and data fetching steps - Add _preview static method for bounded diagnostic output in analysis.py - Log prompt start, completion, errors, and validation details in _reply method - Add info logs for major processing steps in execute method of analysis.py - Add debug and info logs for cache validation, data fetching, and pagination in data.py - Log conditions for skipping reports and cache plans in data.py execute method - Log download summaries and cache writes for news and ETF data - Improve error logging with exception details in cache validation functions - Ensure all logs include context such as record counts, paths, and parameters * refactor(auto_fin): overhaul Auto Fin workflow and schema contracts - Replace old Auto Fin schema models with comprehensive new data classes - Remove legacy Auto Fin analysis step in favor of modular agent-based steps - Introduce AutoFinAgentStep for validating structured agent replies - Simplify data cleaning and JSONL writing utilities for news cache - Remove synchronous and asynchronous dataset methods from analysis step - Redefine Auto Fin analysis configuration for 360-day news retention and multi-step pipeline - Remove embedded analysis prompt templates and replace with agent-driven logic - Update __init__.py exports to match new step implementations and remove deprecated classes - Improve error handling and validation in agent step reply processing - Clean up redundant imports and unused code in analysis and data preparation modules * feat(auto_fin): add detailed logging for analysis and data processing steps - Add timing logs to measure agent prompt processing duration in analysis.py - Log news cache hits and news write paths with record counts in data.py - Include detailed info logs for news download start and completion in data.py - Add start, progress, and completion logs with topic and event counts in history.py - Log start and completion of merge step including path and ETF count in merge.py - Add start and done logs with window and news counts in topic.py * feat(auto_fin): enhance schema and steps with detailed ETF and event modeling - Replace and add multiple AutoFin schema classes to support detailed ETF selection, historical research, market analysis, forecast models, and report output with validation - Implement Shanghai timezone normalization and strict validation in schema models - Remove deprecated AutoFin analysis agent step and consolidate reply handling in base step - Introduce AutoFinStep base class with shared helpers for prompt handling, data fetching, logging, and JSONL file operations - Add AutoFinDataStep to manage daily news data complete with schedule validation, caching, and source validation logic - Update cookbook configuration to customize auto_fin step parameters and simplify outbound proxy settings - Refactor imports and clean unused code for better maintainability * feat(auto_fin): introduce detailed historical event resolution and market similarity analysis - Add AutoFinHistoricalEventReference and AutoFinHistoricalSimilarity models for refined event referencing and similarity judgment - Implement validation to ensure non-empty critical fields and uniqueness of historical news IDs - Develop method to resolve Agent-selected historical event references from workspace files with strict path and existence checks - Enrich historical events with market entry and future returns data after resolution - Redesign market step to calculate similarity-weighted ETF forecasts based on matched historical event similarities - Enforce validation on matched historical events for uniqueness and proper weight summation - Simplify merge step output to final Markdown report without YAML frontmatter and redundant fields - Update user instructions for history search, market, and merge steps to reflect new data structures and responsibilities - Adjust test suite to cover new schema and step behavior changes, including enhanced validation and JSON output formats * feat(auto_fin): add new cron jobs and output analysis jsonl - Add new cron jobs auto_fin_1145_cron and auto_fin_1800_cron with auto_fin_steps - Change auto_fin_0930_cron schedule to run Monday to Sunday - Extend merge step to write analysis data to auto_fin_analysis.jsonl - Update unit tests to verify new cron jobs and their steps configuration * fix(auto_fin): improve atomic file write and refresh daily index - Change temporary file naming to include UUID for uniqueness and hidden prefix - Replace atomic write method from using Path.replace to os.replace with safe unlink - Add import and use os.replace for safer file replace operation - Refresh daily index after writing auto finance markdown and JSONL files - Import and call refresh_day_index in merge step to update file index asynchronously * docs(cookbook): add optional SSH proxy configuration in README files - Introduce optional SSH proxy setup in auto-fin and daily_paper cookbooks - Provide instructions to enable outbound proxy via `daily_cookbook.yaml` and environment variables - Add `REME_PROXY_IP` and `REME_PROXY_ACCOUNT` environment variables descriptions in multiple README files - Update English and Chinese README and README_ZH documents with proxy details - Maintain consistent formatting of environment variable tables across documents * fix(file_io): include schema_version in hidden metadata keys - Added "schema_version" to _INDEX_HIDDEN_METADATA_KEYS in _daily_index.py - Updated _render_notes_block to always include additional keys regardless of schema_version fix(deps): move pproxy dependency to later in pyproject.toml - Removed pproxy from early dependencies list - Added pproxy back near the end of dependency list for better ordering fix(outbound_proxy): require pproxy package for ssh_http proxy - Added importlib.util check for pproxy package presence - Raise RuntimeError if pproxy is not installed when using SSH HTTP outbound proxy - Improved error message suggests installing reme-ai with 'core' extra * docs(readme): update News section with new Cookbook workflows - Clarify introduction of optional Cookbooks with Daily Paper and Auto Fin workflows - Update English README to reflect both paper discovery and file-native ETF event research - Revise Chinese README to include financial news and historical market data research capability - Maintain announcement of paper acceptance at Findings of ACL 2026 * feat(auto_fin): add calculation results to final Markdown output - Implement _calculation_results to summarize forecast for each ETF analyzed - Include program-calculated results in the JSON input for the Markdown report - Update YAML template to incorporate calculation results and adjust recommendation rules - Refine recommendation logic to rely on event impact judgments combined with calculation outputs - Modify tests to verify presence of calculation results and updated report content and format * up prompt * fix(keyword_index): ignore non-indexable chunks during keyword sync - Add is_indexable method to base and BM25 keyword index classes to check text tokenizability - Update local file store to exclude non-indexable chunks from expected document IDs to prevent rebuild - Fix JSONL chunker to correctly handle Unicode line separator U+2028 inside JSON strings without splitting - Add test to ensure non-empty but non-indexable chunk does not trigger keyword index rebuild - Add test to verify U+2028 character does not cause incorrect JSONL record splitting
570 lines
23 KiB
Python
570 lines
23 KiB
Python
"""Codex Python SDK backend for the unified agent wrapper."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
from collections.abc import AsyncGenerator
|
|
from contextlib import suppress
|
|
from dataclasses import dataclass
|
|
from functools import partial
|
|
import hashlib
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
import sys
|
|
import tempfile
|
|
from typing import Any, TYPE_CHECKING
|
|
|
|
from .base_agent_wrapper import BaseAgentWrapper
|
|
from ..component_registry import R
|
|
from ...enumeration import ChunkEnum
|
|
from ...schema import StreamChunk
|
|
|
|
if TYPE_CHECKING:
|
|
from openai_codex import AsyncCodex, AsyncThread, CodexConfig, RunInput
|
|
from openai_codex.types import Notification
|
|
else:
|
|
# Keep the optional Codex SDK out of ReMe's package import path. Tests may
|
|
# also replace this value before the SDK is loaded.
|
|
AsyncCodex: Any = None
|
|
|
|
|
|
def _get_async_codex_class():
|
|
"""Return the optional Codex client class, importing its SDK on first use."""
|
|
global AsyncCodex # pylint: disable=global-statement
|
|
if AsyncCodex is None:
|
|
from openai_codex import AsyncCodex as AsyncCodexClass
|
|
|
|
AsyncCodex = AsyncCodexClass
|
|
return AsyncCodex
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class _CodexAuthConfig:
|
|
"""Resolved authentication settings for one Codex app-server."""
|
|
|
|
mode: str
|
|
api_key: str = ""
|
|
base_url: str = ""
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class _CodexLaunchConfig:
|
|
"""Options fixed for the lifetime of one Codex app-server."""
|
|
|
|
auth_mode: str
|
|
api_key: str
|
|
base_url: str
|
|
codex_bin: str | None
|
|
config_overrides: tuple[str, ...]
|
|
experimental_api: bool
|
|
|
|
|
|
@R.register("codex")
|
|
class CodexAgentWrapper(BaseAgentWrapper):
|
|
"""Agent wrapper backed by the Codex Python SDK."""
|
|
|
|
SDK_PACKAGE = "openai-codex"
|
|
_CLIENT_OPTION_NAMES = frozenset(
|
|
{
|
|
"api_key",
|
|
"auth_mode",
|
|
"base_url",
|
|
"codex_bin",
|
|
"codex_home",
|
|
"config_overrides",
|
|
"cwd",
|
|
"experimental_api",
|
|
"launch_args_override",
|
|
},
|
|
)
|
|
|
|
# pylint: disable=too-many-arguments
|
|
def __init__(
|
|
self,
|
|
mcp_config: str | None = None,
|
|
codex_home: str | Path | None = None,
|
|
*,
|
|
auth_mode: str = "auto",
|
|
api_key: str = "",
|
|
base_url: str = "",
|
|
codex_bin: str | None = None,
|
|
config_overrides: list[str] | tuple[str, ...] | None = None,
|
|
experimental_api: bool = True,
|
|
**kwargs,
|
|
) -> None:
|
|
if "launch_args_override" in kwargs:
|
|
raise TypeError("launch_args_override is not supported; configure codex_bin instead")
|
|
super().__init__(**kwargs)
|
|
self.mcp_config = mcp_config
|
|
self._codex_home = codex_home
|
|
self._launch_config = _CodexLaunchConfig(
|
|
auth_mode=auth_mode,
|
|
api_key=api_key,
|
|
base_url=base_url,
|
|
codex_bin=codex_bin,
|
|
config_overrides=tuple(config_overrides or ()),
|
|
experimental_api=experimental_api,
|
|
)
|
|
self._codex: AsyncCodex | None = None
|
|
self._turn_lock = asyncio.Lock()
|
|
self._mcp_snapshot_path: Path | None = None
|
|
self._thread_tool_contexts: dict[str, str] = {}
|
|
|
|
@property
|
|
def session_path(self) -> Path:
|
|
"""Directory used for Codex state and persisted threads."""
|
|
if self._codex_home:
|
|
path = Path(self._codex_home).expanduser()
|
|
return path if path.is_absolute() else self.workspace_path / path
|
|
if self.app_context is None:
|
|
return self.workspace_path / "mem_session" / "codex"
|
|
return self.workspace_path / self.app_context.app_config.mem_session_dir / "codex"
|
|
|
|
@property
|
|
def wrapper_session_path(self) -> Path:
|
|
"""ReMe-owned session data, kept separate from a shared OAuth CODEX_HOME."""
|
|
if self.app_context is None:
|
|
return self.workspace_path / "mem_session" / "codex"
|
|
return self.workspace_path / self.app_context.app_config.mem_session_dir / "codex"
|
|
|
|
def _ensure_skills(self, skills: list[str] | str | None) -> None:
|
|
"""Expose selected project skills through Codex's repo-level directory."""
|
|
sources = self._resolve_project_skills(skills)
|
|
if not sources:
|
|
return
|
|
|
|
target_root = self.project_path / ".agents" / "skills"
|
|
target_root.mkdir(parents=True, exist_ok=True)
|
|
for name, source in sources.items():
|
|
target = target_root / name
|
|
if target.is_symlink():
|
|
if target.resolve() == source.resolve():
|
|
continue
|
|
raise FileExistsError(f"Codex skill conflict: {target} points to {target.resolve(strict=False)}")
|
|
if target.exists():
|
|
raise FileExistsError(f"Codex skill conflict: {target} already exists and was preserved")
|
|
relative_source = os.path.relpath(source, target.parent)
|
|
target.symlink_to(relative_source, target_is_directory=True)
|
|
|
|
def _explicit_mcp_config(self, kwargs: dict[str, Any]) -> str | None:
|
|
value = kwargs.get("mcp_config") if "mcp_config" in kwargs else self.mcp_config
|
|
if value is None:
|
|
return None
|
|
source = Path(str(value)).expanduser()
|
|
if source.suffix in {".yaml", ".yml", ".json"}:
|
|
if not source.is_absolute():
|
|
source = self.workspace_path / source
|
|
return str(source.absolute())
|
|
return str(value)
|
|
|
|
def _effective_config_snapshot(self) -> Path:
|
|
"""Create one private snapshot that remains valid for the client lifetime."""
|
|
if self._mcp_snapshot_path is not None:
|
|
return self._mcp_snapshot_path
|
|
if self.app_context is None:
|
|
raise RuntimeError("Cannot snapshot MCP config without an app_context")
|
|
snapshot_dir = self.wrapper_session_path / "reme-mcp"
|
|
snapshot_dir.mkdir(parents=True, exist_ok=True)
|
|
fd, raw_path = tempfile.mkstemp(prefix="config-", suffix=".json", dir=snapshot_dir)
|
|
try:
|
|
os.fchmod(fd, 0o600)
|
|
with os.fdopen(fd, "w", encoding="utf-8") as stream:
|
|
json.dump(self.app_context.app_config.model_dump(mode="json"), stream)
|
|
except BaseException:
|
|
with suppress(OSError):
|
|
os.close(fd)
|
|
Path(raw_path).unlink(missing_ok=True)
|
|
raise
|
|
self._mcp_snapshot_path = Path(raw_path)
|
|
return self._mcp_snapshot_path
|
|
|
|
def _mcp_config_source(self, kwargs: dict[str, Any]) -> str:
|
|
return self._explicit_mcp_config(kwargs) or str(self._effective_config_snapshot())
|
|
|
|
@staticmethod
|
|
def _resolve_auth_config(auth_mode: str, api_key: str = "", base_url: str = "") -> _CodexAuthConfig:
|
|
"""Resolve login from wrapper options.
|
|
|
|
The Codex child process still inherits the parent process environment.
|
|
"""
|
|
requested_mode = str(auth_mode or "auto").lower()
|
|
if requested_mode not in {"auto", "api_key", "oauth"}:
|
|
raise ValueError("auth_mode must be one of: auto, api_key, oauth")
|
|
if requested_mode == "oauth":
|
|
return _CodexAuthConfig(mode="oauth")
|
|
|
|
api_key = api_key if isinstance(api_key, str) else ""
|
|
if requested_mode == "api_key" and not api_key:
|
|
raise ValueError("auth_mode='api_key' requires a non-empty API key")
|
|
if not api_key:
|
|
return _CodexAuthConfig(mode="oauth")
|
|
|
|
base_url = base_url if isinstance(base_url, str) else ""
|
|
return _CodexAuthConfig(mode="api_key", api_key=api_key, base_url=base_url)
|
|
|
|
def _build_client_config(self, auth: _CodexAuthConfig) -> CodexConfig:
|
|
from openai_codex import CodexConfig
|
|
|
|
env = dict(self.subprocess_environment)
|
|
self.session_path.mkdir(parents=True, exist_ok=True)
|
|
env["CODEX_HOME"] = str(self.session_path)
|
|
|
|
overrides = list(self._launch_config.config_overrides)
|
|
if auth.base_url:
|
|
overrides.append(f"openai_base_url={json.dumps(auth.base_url)}")
|
|
login_method = "api" if auth.mode == "api_key" else "chatgpt"
|
|
overrides.append(f"forced_login_method={json.dumps(login_method)}")
|
|
return CodexConfig(
|
|
codex_bin=self._launch_config.codex_bin,
|
|
config_overrides=tuple(overrides),
|
|
cwd=str(self.cwd),
|
|
env=env,
|
|
client_name="reme",
|
|
client_title="ReMe",
|
|
experimental_api=self._launch_config.experimental_api,
|
|
)
|
|
|
|
@classmethod
|
|
def _reject_client_options(cls, kwargs: dict[str, Any]) -> None:
|
|
invalid = sorted(cls._CLIENT_OPTION_NAMES.intersection(kwargs))
|
|
if invalid:
|
|
names = ", ".join(invalid)
|
|
raise TypeError(f"Codex client options must be configured on the wrapper: {names}")
|
|
|
|
def _mcp_server_config(self, kwargs: dict[str, Any]) -> dict[str, Any] | None:
|
|
from ..job import BackgroundJob, StreamJob
|
|
|
|
job_names = list(dict.fromkeys(kwargs.get("job_tools") or []))
|
|
if not job_names:
|
|
return None
|
|
jobs = self._resolve_job_tools(job_names)
|
|
unsupported = [job.name for job in jobs if isinstance(job, (BackgroundJob, StreamJob))]
|
|
if unsupported:
|
|
raise TypeError(f"Codex job_tools must be non-stream request jobs: {', '.join(unsupported)}")
|
|
|
|
config_source = self._mcp_config_source(kwargs)
|
|
args = [
|
|
"-m",
|
|
"reme.components.agent_wrapper.codex_mcp_server",
|
|
"--config",
|
|
config_source,
|
|
"--workspace",
|
|
str(self.workspace_path),
|
|
]
|
|
for name in job_names:
|
|
args.extend(["--job", name])
|
|
args.extend(["--tool-context-id", str(kwargs.get("tool_context_id") or "")])
|
|
return {
|
|
"command": sys.executable,
|
|
"args": args,
|
|
"cwd": str(self.project_path),
|
|
"required": True,
|
|
"enabled_tools": job_names,
|
|
"startup_timeout_sec": kwargs.get("mcp_startup_timeout", 30),
|
|
"tool_timeout_sec": kwargs.get("mcp_tool_timeout", 300),
|
|
}
|
|
|
|
def _thread_config(self, kwargs: dict[str, Any]) -> dict[str, Any] | None:
|
|
config = dict(kwargs.get("config") or {})
|
|
if proxy_environment := self.command_proxy_environment:
|
|
shell_environment_policy = dict(config.get("shell_environment_policy") or {})
|
|
environment = dict(shell_environment_policy.get("set") or {})
|
|
environment.update(proxy_environment)
|
|
shell_environment_policy["set"] = environment
|
|
config["shell_environment_policy"] = shell_environment_policy
|
|
if server := self._mcp_server_config(kwargs):
|
|
servers = dict(config.get("mcp_servers") or {})
|
|
server_key = hashlib.sha256(json.dumps(server, sort_keys=True).encode()).hexdigest()[:12]
|
|
servers[f"reme_jobs_{server_key}"] = server
|
|
config["mcp_servers"] = servers
|
|
return config or None
|
|
|
|
@staticmethod
|
|
def _enum(enum_cls: Any, value: Any, default: Any = None) -> Any:
|
|
if value is None:
|
|
return default
|
|
return value if isinstance(value, enum_cls) else enum_cls(value)
|
|
|
|
async def _open_thread(self, codex: AsyncCodex, kwargs: dict[str, Any]) -> AsyncThread:
|
|
from openai_codex import ApprovalMode, Sandbox
|
|
from openai_codex.types import Personality, ThreadSource, ThreadStartSource
|
|
|
|
resume = kwargs.get("resume") or ""
|
|
session_id = kwargs.get("session_id") or ""
|
|
if resume and session_id and resume != session_id:
|
|
raise ValueError("resume and session_id must identify the same Codex thread")
|
|
thread_id = resume or session_id
|
|
fork_session = bool(kwargs.get("fork_session", False))
|
|
if fork_session and not thread_id:
|
|
raise ValueError("fork_session=True requires resume or session_id")
|
|
requested_tool_context = str(kwargs.get("tool_context_id") or "")
|
|
if not fork_session and thread_id in self._thread_tool_contexts:
|
|
if requested_tool_context != self._thread_tool_contexts[thread_id]:
|
|
raise ValueError("tool_context_id cannot change when resuming a Codex thread")
|
|
|
|
common = {
|
|
"approval_mode": self._enum(ApprovalMode, kwargs.get("approval_mode"), ApprovalMode.auto_review),
|
|
"base_instructions": kwargs.get("base_instructions"),
|
|
"config": self._thread_config(kwargs),
|
|
"cwd": str(self.cwd),
|
|
"developer_instructions": kwargs.get("system_prompt"),
|
|
"model": kwargs.get("model"),
|
|
"model_provider": kwargs.get("model_provider"),
|
|
"sandbox": self._enum(Sandbox, kwargs.get("sandbox"), Sandbox.full_access),
|
|
"service_tier": kwargs.get("service_tier"),
|
|
}
|
|
personality = self._enum(Personality, kwargs.get("personality"))
|
|
thread_source = self._enum(ThreadSource, kwargs.get("thread_source"))
|
|
if fork_session:
|
|
thread = await codex.thread_fork(
|
|
thread_id,
|
|
ephemeral=kwargs.get("ephemeral"),
|
|
thread_source=thread_source,
|
|
**common,
|
|
)
|
|
elif thread_id:
|
|
thread = await codex.thread_resume(thread_id, personality=personality, **common)
|
|
else:
|
|
thread = await codex.thread_start(
|
|
ephemeral=kwargs.get("ephemeral"),
|
|
personality=personality,
|
|
service_name=kwargs.get("service_name"),
|
|
session_start_source=self._enum(ThreadStartSource, kwargs.get("session_start_source")),
|
|
thread_source=thread_source,
|
|
**common,
|
|
)
|
|
self._thread_tool_contexts[thread.id] = requested_tool_context
|
|
return thread
|
|
|
|
def _turn_kwargs(self, kwargs: dict[str, Any]) -> dict[str, Any]:
|
|
from openai_codex import ApprovalMode, Sandbox
|
|
from openai_codex.types import Personality, ReasoningEffort, ReasoningSummary
|
|
|
|
return {
|
|
"approval_mode": self._enum(ApprovalMode, kwargs.get("approval_mode")),
|
|
"cwd": str(self.cwd),
|
|
"effort": self._enum(ReasoningEffort, kwargs.get("effort")),
|
|
"model": kwargs.get("model"),
|
|
"output_schema": kwargs.get("output_schema"),
|
|
"personality": self._enum(Personality, kwargs.get("personality")),
|
|
"sandbox": self._enum(Sandbox, kwargs.get("sandbox")),
|
|
"service_tier": kwargs.get("service_tier"),
|
|
"summary": self._enum(ReasoningSummary, kwargs.get("summary")),
|
|
}
|
|
|
|
async def _get_codex(self) -> AsyncCodex:
|
|
"""Lazily start the component-owned client from its fixed launch configuration."""
|
|
if self._codex is not None:
|
|
return self._codex
|
|
auth = self._resolve_auth_config(
|
|
self._launch_config.auth_mode,
|
|
self._launch_config.api_key,
|
|
self._launch_config.base_url,
|
|
)
|
|
codex = _get_async_codex_class()(self._build_client_config(auth))
|
|
try:
|
|
if auth.mode == "api_key":
|
|
await codex.login_api_key(auth.api_key)
|
|
else:
|
|
account = await codex.account()
|
|
if account.account is None:
|
|
raise RuntimeError(f"No ChatGPT OAuth login found in CODEX_HOME: {self.session_path}")
|
|
except BaseException:
|
|
await codex.close()
|
|
raise
|
|
self._codex = codex
|
|
return codex
|
|
|
|
async def _close(self) -> None:
|
|
"""Close the persistent app-server and remove its private config snapshot."""
|
|
async with self._turn_lock:
|
|
codex, self._codex = self._codex, None
|
|
self._thread_tool_contexts.clear()
|
|
try:
|
|
if codex is not None:
|
|
await codex.close()
|
|
finally:
|
|
if self._mcp_snapshot_path is not None:
|
|
self._mcp_snapshot_path.unlink(missing_ok=True)
|
|
self._mcp_snapshot_path = None
|
|
|
|
@staticmethod
|
|
def _serialize(value: Any) -> Any:
|
|
from pydantic_core import to_jsonable_python
|
|
|
|
return to_jsonable_python(value, by_alias=True)
|
|
|
|
async def reply(self, inputs: RunInput, **kwargs) -> dict:
|
|
"""Run one Codex turn and return its final response."""
|
|
self._reject_client_options(kwargs)
|
|
kwargs = self._merged_kwargs(kwargs)
|
|
self._ensure_skills(kwargs.get("skills"))
|
|
await self.start()
|
|
async with self._turn_lock:
|
|
codex = await self._get_codex()
|
|
thread = await self._open_thread(codex, kwargs)
|
|
result = await thread.run(inputs, **self._turn_kwargs(kwargs))
|
|
|
|
final_response = result.final_response or ""
|
|
response = {
|
|
"session_id": thread.id,
|
|
"last_message": final_response,
|
|
"result": final_response,
|
|
"turn": self._serialize(result),
|
|
}
|
|
if kwargs.get("output_schema") is not None:
|
|
try:
|
|
response["structured_output"] = json.loads(final_response)
|
|
except json.JSONDecodeError as exc:
|
|
raise ValueError("Codex returned invalid JSON for the requested output_schema") from exc
|
|
return response
|
|
|
|
async def compact_session(self, session_id: str) -> None:
|
|
"""Start native compaction of a Codex thread."""
|
|
await self.start()
|
|
async with self._turn_lock:
|
|
codex = await self._get_codex()
|
|
thread = await codex.thread_resume(session_id)
|
|
await thread.compact()
|
|
|
|
@classmethod
|
|
# pylint: disable=too-many-return-statements
|
|
def _event_to_chunks(cls, event: Notification, session_id: str) -> list[StreamChunk]:
|
|
"""Convert one Codex app-server notification to unified stream chunks."""
|
|
method, payload = event.method, event.payload
|
|
make_chunk = partial(cls._chunk, session_id=session_id)
|
|
if method == "turn/started":
|
|
return [make_chunk(ChunkEnum.REPLY_START, metadata={"turn_id": payload.turn.id})]
|
|
if method == "item/agentMessage/delta":
|
|
return [make_chunk(ChunkEnum.CONTENT, block_id=payload.item_id, chunk=payload.delta)]
|
|
if method in {"item/reasoning/summaryTextDelta", "item/reasoning/textDelta", "item/plan/delta"}:
|
|
return [make_chunk(ChunkEnum.THINK, block_id=payload.item_id, chunk=payload.delta)]
|
|
if method in {"item/commandExecution/outputDelta", "item/fileChange/outputDelta"}:
|
|
return [
|
|
make_chunk(
|
|
ChunkEnum.TOOL_RESULT,
|
|
block_id=payload.item_id,
|
|
tool_call_id=payload.item_id,
|
|
chunk=payload.delta,
|
|
),
|
|
]
|
|
if method == "item/mcpToolCall/progress":
|
|
return [
|
|
make_chunk(
|
|
ChunkEnum.TOOL_RESULT,
|
|
block_id=payload.item_id,
|
|
tool_call_id=payload.item_id,
|
|
chunk=payload.message,
|
|
),
|
|
]
|
|
if method in {"item/autoApprovalReview/started", "item/autoApprovalReview/completed"}:
|
|
action = cls._serialize(payload.action)
|
|
review = cls._serialize(payload.review)
|
|
review_id = payload.review_id
|
|
target_item_id = getattr(payload, "target_item_id", None)
|
|
status = "started" if method.endswith("/started") else "completed"
|
|
decision_source = cls._serialize(getattr(payload, "decision_source", None))
|
|
return [
|
|
make_chunk(
|
|
ChunkEnum.APPROVAL,
|
|
block_id=target_item_id or review_id,
|
|
tool_call_id=target_item_id,
|
|
chunk=action,
|
|
metadata={
|
|
"review_id": review_id,
|
|
"status": status,
|
|
"review": review,
|
|
"decision_source": decision_source,
|
|
"turn_id": payload.turn_id,
|
|
},
|
|
),
|
|
]
|
|
if method in {"item/started", "item/completed"}:
|
|
item = payload.item.root
|
|
item_type, item_id = item.type, item.id
|
|
tool_types = {"commandExecution", "fileChange", "mcpToolCall", "dynamicToolCall", "collabAgentToolCall"}
|
|
if item_type not in tool_types:
|
|
return []
|
|
name = getattr(item, "tool", None) or item_type
|
|
chunk_type = ChunkEnum.TOOL_CALL if method == "item/started" else ChunkEnum.TOOL_RESULT
|
|
return [
|
|
make_chunk(
|
|
chunk_type,
|
|
block_id=item_id,
|
|
tool_call_id=item_id,
|
|
tool_call_name=name,
|
|
chunk=cls._serialize(item),
|
|
),
|
|
]
|
|
if method == "thread/tokenUsage/updated":
|
|
usage = payload.token_usage.last
|
|
data = cls._serialize(usage)
|
|
return [
|
|
make_chunk(
|
|
ChunkEnum.USAGE,
|
|
chunk=data,
|
|
input_tokens=usage.input_tokens,
|
|
output_tokens=usage.output_tokens,
|
|
),
|
|
]
|
|
if method == "error":
|
|
return [
|
|
make_chunk(
|
|
ChunkEnum.ERROR,
|
|
chunk=payload.error.message,
|
|
metadata={"will_retry": payload.will_retry},
|
|
),
|
|
]
|
|
if method == "turn/completed":
|
|
turn = payload.turn
|
|
chunks = []
|
|
if turn.error:
|
|
chunks.append(make_chunk(ChunkEnum.ERROR, chunk=turn.error.message))
|
|
chunks.append(
|
|
make_chunk(
|
|
ChunkEnum.REPLY_END,
|
|
metadata={
|
|
"turn_id": turn.id,
|
|
"status": turn.status.value,
|
|
"duration_ms": turn.duration_ms,
|
|
},
|
|
),
|
|
)
|
|
return chunks
|
|
if getattr(payload, "turn_id", None):
|
|
return [
|
|
make_chunk(
|
|
ChunkEnum.DATA,
|
|
block_id=getattr(payload, "item_id", None),
|
|
chunk=cls._serialize(payload),
|
|
metadata={"codex_method": method},
|
|
),
|
|
]
|
|
return []
|
|
|
|
async def reply_stream(self, inputs: RunInput, **kwargs) -> AsyncGenerator[StreamChunk, None]:
|
|
"""Stream Codex app-server notifications as unified chunks."""
|
|
self._reject_client_options(kwargs)
|
|
kwargs = self._merged_stream_kwargs(kwargs)
|
|
self._ensure_skills(kwargs.get("skills"))
|
|
await self.start()
|
|
async with self._turn_lock:
|
|
codex = await self._get_codex()
|
|
thread = await self._open_thread(codex, kwargs)
|
|
turn = await thread.turn(inputs, **self._turn_kwargs(kwargs))
|
|
stream = turn.stream()
|
|
completed = False
|
|
try:
|
|
async for event in stream:
|
|
if event.method == "turn/completed":
|
|
completed = True
|
|
for chunk in self._event_to_chunks(event, thread.id):
|
|
yield chunk
|
|
finally:
|
|
if not completed:
|
|
try:
|
|
await turn.interrupt()
|
|
except Exception as exc: # pylint: disable=broad-exception-caught
|
|
self.logger.warning(f"Failed to interrupt Codex turn {turn.id}: {exc}")
|
|
await stream.aclose()
|