Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ members = [
"crates/persisting-capture",
"crates/persisting-cli",
"crates/persisting-dlcapt",
"crates/persisting-compute",
]

[workspace.package]
Expand Down
1 change: 1 addition & 0 deletions crates/persisting-cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ persisting-capture = { path = "../persisting-capture" }
persisting-dlcapt = { path = "../persisting-dlcapt", optional = true }
persisting-engine = { path = "../persisting-engine" }
persisting-proto = { path = "../persisting-proto" }
persisting-compute = { path = "../persisting-compute", features = ["traj-sink"] }
ron = "0.8"
serde = { version = "1", features = ["derive"] }
serde_json = "1"
Expand Down
165 changes: 165 additions & 0 deletions crates/persisting-cli/src/capture/proxy_config_args.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
//! Shared proxy config CLI flags + env vars for `traj capture` / `traj proxy`.

use std::path::PathBuf;

use anyhow::{Context, Result};
use clap::Args;
use persisting_capture::config::{
env, materialize_proxy_config, parse_capture_level, resolve_proxy_config, CaptureLevel,
ProxyConfig, ProxyConfigInput, ProxyConfigOverrides,
};

/// Proxy settings overridable via CLI flags and environment variables.
#[derive(Debug, Clone, Default, Args)]
pub struct ProxyConfigOverrideArgs {
/// Proxy listen address (overrides TOML `listen`).
#[arg(long, env = env::LISTEN, value_name = "ADDR")]
pub listen: Option<String>,
/// Admin API listen address (overrides TOML `admin_listen`).
#[arg(long, env = env::ADMIN_LISTEN, value_name = "ADDR")]
pub admin_listen: Option<String>,
/// Trajectory agent id segment (overrides TOML `agent_id`).
#[arg(long, env = env::AGENT_ID, value_name = "SEG")]
pub agent_id: Option<String>,
/// Session id HTTP header name (overrides TOML `session_header`).
#[arg(long, env = env::SESSION_HEADER, value_name = "HEADER")]
pub session_header: Option<String>,
/// Capture granularity: `summary`, `dialogue`, or `full`.
#[arg(long, env = env::CAPTURE_LEVEL, value_name = "LEVEL")]
pub capture_level: Option<String>,
/// Model routes as JSON array (overrides all `[[models]]` entries).
#[arg(long, env = env::MODELS_JSON, value_name = "JSON")]
pub models_json: Option<String>,
/// Model routes as TOML `[[models]]` table(s) (overrides all routes).
#[arg(long, env = env::MODELS_TOML, value_name = "TOML")]
pub models_toml: Option<String>,
/// Model route name to create or patch (default: `*`).
#[arg(long, env = env::MODEL, value_name = "NAME")]
pub model: Option<String>,
/// Upstream OpenAI-compatible base URL for `--model` (e.g. `https://api.deepseek.com/v1`).
#[arg(long, env = env::UPSTREAM, value_name = "URL")]
pub upstream: Option<String>,
/// Anthropic-compatible upstream for `/v1/messages`.
#[arg(long, env = env::UPSTREAM_ANTHROPIC, value_name = "URL")]
pub upstream_anthropic: Option<String>,
/// Provider label for the model route: `openai`, `anthropic`, …
#[arg(long, env = env::PROVIDER, value_name = "NAME")]
pub provider: Option<String>,
/// Env var holding the upstream API key for this route.
#[arg(long, env = env::API_KEY_ENV, value_name = "VAR")]
pub api_key_env: Option<String>,
/// Inline upstream API key (prefer `--api-key-env` in scripts).
#[arg(long, env = env::API_KEY, value_name = "KEY")]
pub api_key: Option<String>,
/// Forward client model to another configured route name.
#[arg(long, env = env::FORWARD, value_name = "NAME")]
pub forward: Option<String>,
/// Legacy API path prefix for `--model` route (prefer full prefix in `--upstream`).
#[arg(long, env = env::PATH_PREFIX, value_name = "PREFIX")]
pub path_prefix: Option<String>,
}

impl ProxyConfigOverrideArgs {
pub fn to_overrides(&self) -> Result<ProxyConfigOverrides> {
let capture_level = match &self.capture_level {
Some(s) => Some(parse_capture_level(s)?),
None => None,
};
Ok(ProxyConfigOverrides {
listen: self.listen.clone(),
admin_listen: self.admin_listen.clone(),
agent_id: self.agent_id.clone(),
session_header: self.session_header.clone(),
capture_level,
debug: None,
models_json: self.models_json.clone(),
models_toml: self.models_toml.clone(),
model_name: self.model.clone(),
upstream: self.upstream.clone(),
upstream_anthropic: self.upstream_anthropic.clone(),
provider: self.provider.clone(),
api_key_env: self.api_key_env.clone(),
api_key: self.api_key.clone(),
forward: self.forward.clone(),
path_prefix: self.path_prefix.clone(),
})
}
}

/// Optional proxy TOML path plus CLI/env overrides.
#[derive(Debug, Clone, Default, Args)]
pub struct ProxyConfigArgs {
/// Proxy config TOML (`listen`, `models`, …). Optional when set via env/flags.
#[arg(
long,
short = 'c',
value_name = "FILE",
env = env::CONFIG_FILE,
conflicts_with = "config_toml"
)]
pub config: Option<PathBuf>,
/// Full proxy config as inline TOML (alternative to `-c`; supports every TOML field).
#[arg(long, env = env::CONFIG_TOML, value_name = "TOML", conflicts_with = "config")]
pub config_toml: Option<String>,
#[command(flatten)]
pub overrides: ProxyConfigOverrideArgs,
}

impl ProxyConfigArgs {
pub fn input(&self) -> Result<ProxyConfigInput> {
Ok(ProxyConfigInput {
config_file: self.config.clone(),
config_toml: self.config_toml.clone(),
overrides: self.overrides.to_overrides()?,
})
}

pub fn resolve(&self) -> Result<ProxyConfig> {
resolve_proxy_config(&self.input()?)
}

pub fn resolve_with_debug(&self, cli_debug: bool) -> Result<ProxyConfig> {
let mut input = self.input()?;
if cli_debug {
input.overrides.debug = Some(true);
}
resolve_proxy_config(&input)
}
}

/// Resolved config and on-disk path (explicit `-c` or materialized `{storage}/proxy.toml`).
pub struct ResolvedProxyConfig {
pub config: ProxyConfig,
pub config_path: PathBuf,
}

impl ProxyConfigArgs {
pub fn materialize(&self, storage: &std::path::Path, cli_debug: bool) -> Result<ResolvedProxyConfig> {
let mut input = self.input()?;
if cli_debug {
input.overrides.debug = Some(true);
}
let (config, config_path) =
materialize_proxy_config(storage, &input).context("resolve proxy config")?;
Ok(ResolvedProxyConfig {
config,
config_path,
})
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn override_args_map_capture_level() {
let args = ProxyConfigOverrideArgs {
capture_level: Some("full".into()),
upstream: Some("http://127.0.0.1:1/v1".into()),
..Default::default()
};
let o = args.to_overrides().unwrap();
assert_eq!(o.capture_level, Some(CaptureLevel::Full));
}
}
34 changes: 24 additions & 10 deletions crates/persisting-cli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -296,12 +296,11 @@ struct Cli {
#[derive(Debug, Subcommand)]
enum Command {
Search(SearchArgs),
/// Agent trajectory: capture, proxy, inspect, repair
#[command(long_about = TRAJ_LONG_ABOUT)]
/// Agent trajectory: capture, proxy, inspect, repair(短名 `traj`)
#[command(visible_alias = "traj", long_about = TRAJ_LONG_ABOUT)]
Trajectory(TrajectoryArgs),
/// Short alias for `trajectory`
#[command(name = "traj", long_about = TRAJ_LONG_ABOUT)]
Traj(TrajectoryArgs),
/// Run a compute plan (default). Use `--check` to validate locally first.
Compute(persisting_compute::ComputeArgs),
}

#[derive(Debug, Args)]
Expand Down Expand Up @@ -1035,10 +1034,25 @@ fn engine_lib_names() -> [&'static str; 3] {

fn main() -> Result<()> {
let cli = Cli::parse_from(normalize_cli_args(std::env::args().collect()));
let mut lazy = LazyEngine::new(&cli);
match &cli.command {
Command::Search(args) => run_search(&mut lazy, args)?,
Command::Trajectory(args) | Command::Traj(args) => run_trajectory(&mut lazy, args)?,
Command::Compute(args) => {
persisting_compute::cli::init_tracing_with_verbose(args.verbose);
let args = args.clone();
let code = tokio::runtime::Runtime::new()
.context("tokio runtime")?
.block_on(persisting_compute::run_compute(args))?;
if code != std::process::ExitCode::SUCCESS {
std::process::exit(1);
}
}
Command::Search(args) => {
let mut lazy = LazyEngine::new(&cli);
run_search(&mut lazy, args)?;
}
Command::Trajectory(args) => {
let mut lazy = LazyEngine::new(&cli);
run_trajectory(&mut lazy, args)?;
}
}
Ok(())
}
Expand Down Expand Up @@ -2525,11 +2539,11 @@ mod tests {

fn add_args_from_cli(argv: &[&str]) -> TrajectoryAddArgs {
let cli = Cli::try_parse_from(argv).unwrap();
let Command::Traj(TrajectoryArgs {
let Command::Trajectory(TrajectoryArgs {
command: TrajectoryCommand::Add(args),
}) = cli.command
else {
panic!("expected traj add");
panic!("expected traj/trajectory add");
};
args
}
Expand Down
35 changes: 35 additions & 0 deletions crates/persisting-compute/Cargo.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
[package]
name = "persisting-compute"
version.workspace = true
edition.workspace = true
authors.workspace = true
license.workspace = true
description = "Thin compute control plane: plan script → streamed tasks → Pulsing workers (via `persisting compute`)"

[dependencies]
anyhow = "1"
async-trait = "0.1"
chrono = { version = "0.4", optional = true }
clap = { version = "4", features = ["derive", "env"] }
futures = "0.3"
persisting-capture = { path = "../persisting-capture", optional = true }
persisting-engine = { path = "../persisting-engine", optional = true }
persisting-proto = { path = "../persisting-proto", optional = true }
pulsing-actor = { workspace = true }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
tokio = { version = "1", features = ["macros", "rt-multi-thread", "process", "io-util", "sync", "time", "signal", "fs"] }
tokio-stream = "0.1"
tokio-util = { version = "0.7", features = ["rt"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
uuid = { version = "1", features = ["v4"] }

[features]
default = []
# Append compute results to Vortex via TrajectoryAppend (Tee with JsonlFileSink).
traj-sink = ["dep:chrono", "dep:persisting-capture", "dep:persisting-engine", "dep:persisting-proto"]

[dev-dependencies]
tempfile = "3"
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
Loading
Loading