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
1 change: 1 addition & 0 deletions Cargo.lock

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

4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -111,8 +111,8 @@ Core repo flow:
omnigraph init --schema ./schema.pg ./repo.omni
omnigraph load --data ./data.jsonl --mode overwrite ./repo.omni
omnigraph snapshot ./repo.omni --branch main --json
omnigraph read ./repo.omni --query ./queries.gq --name get_person --params '{"name":"Alice"}'
omnigraph change ./repo.omni --query ./queries.gq --name insert_person --params '{"name":"Mina","age":28}'
omnigraph read --uri ./repo.omni --query ./queries.gq --name get_person --params '{"name":"Alice"}'
omnigraph change --uri ./repo.omni --query ./queries.gq --name insert_person --params '{"name":"Mina","age":28}'
omnigraph branch create --uri ./repo.omni --from main feature-x
omnigraph branch merge --uri ./repo.omni feature-x --into main
```
Expand Down
165 changes: 90 additions & 75 deletions crates/omnigraph-cli/src/main.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
use std::fs;
use std::io::{self, Write};
use std::path::Path;
use std::path::PathBuf;

Expand Down Expand Up @@ -125,7 +126,7 @@ enum Command {
config: Option<PathBuf>,
#[arg(long)]
branch: Option<String>,
#[arg(long)]
#[arg(long, hide = true)]
jsonl: bool,
#[arg(long = "type")]
type_names: Vec<String>,
Expand All @@ -145,7 +146,10 @@ enum Command {
/// Execute a read query against a branch or snapshot
Read {
/// Repo URI
#[arg(long)]
uri: Option<String>,
#[arg(hide = true)]
legacy_uri: Option<String>,
#[arg(long)]
target: Option<String>,
#[arg(long)]
Expand All @@ -172,7 +176,10 @@ enum Command {
/// Execute a graph change query against a branch
Change {
/// Repo URI
#[arg(long)]
uri: Option<String>,
#[arg(hide = true)]
legacy_uri: Option<String>,
#[arg(long)]
target: Option<String>,
#[arg(long)]
Expand Down Expand Up @@ -347,6 +354,7 @@ enum CommitCommand {
/// Show a graph commit
Show {
/// Repo URI
#[arg(long)]
uri: Option<String>,
#[arg(long)]
target: Option<String>,
Expand Down Expand Up @@ -428,6 +436,8 @@ struct LoadOutput<'a> {
mode: &'a str,
nodes_loaded: usize,
edges_loaded: usize,
node_types_loaded: usize,
edge_types_loaded: usize,
}

#[derive(Debug, Serialize)]
Expand Down Expand Up @@ -644,31 +654,6 @@ async fn remote_json<T: DeserializeOwned>(
Ok(serde_json::from_str(&text)?)
}

async fn remote_text(
client: &reqwest::Client,
method: Method,
url: String,
body: Option<Value>,
bearer_token: Option<&str>,
) -> Result<String> {
let request = apply_bearer_token(client.request(method, url), bearer_token);
let request = if let Some(body) = body {
request.json(&body)
} else {
request
};
let response = request.send().await?;
let status = response.status();
let text = response.text().await?;
if !status.is_success() {
if let Ok(error) = serde_json::from_str::<ErrorOutput>(&text) {
bail!(error.error);
}
bail!("server returned {}: {}", status, text);
}
Ok(text)
}

fn resolve_uri(
config: &OmnigraphConfig,
cli_uri: Option<String>,
Expand Down Expand Up @@ -784,14 +769,18 @@ fn print_load_human(
mode: CliLoadMode,
nodes_loaded: usize,
edges_loaded: usize,
node_types_loaded: usize,
edge_types_loaded: usize,
) {
println!(
"loaded {} on branch {} with {}: {} node types, {} edge types",
"loaded {} on branch {} with {}: {} nodes across {} node types, {} edges across {} edge types",
uri,
branch,
mode.as_str(),
nodes_loaded,
edges_loaded
node_types_loaded,
edges_loaded,
edge_types_loaded
);
}

Expand Down Expand Up @@ -1148,23 +1137,17 @@ fn resolve_alias<'a>(
Ok(Some((alias_name, alias)))
}

fn normalize_alias_args(
fn normalize_legacy_alias_uri(
uri: Option<String>,
target: Option<&str>,
default_target_present: bool,
target_available: bool,
alias_name: Option<&str>,
mut alias_args: Vec<String>,
) -> (Option<String>, Vec<String>) {
let Some(candidate) = uri else {
return (None, alias_args);
};

if alias_name.is_some()
&& (target.is_some() || default_target_present)
&& !is_remote_uri(&candidate)
&& !candidate.contains(std::path::MAIN_SEPARATOR)
&& !Path::new(&candidate).exists()
{
if alias_name.is_some() && target_available {
alias_args.insert(0, candidate);
return (None, alias_args);
}
Expand Down Expand Up @@ -1392,36 +1375,53 @@ async fn execute_change_remote(
.await
}

async fn execute_export(
async fn execute_export_to_writer<W: Write>(
uri: &str,
branch: &str,
type_names: &[String],
table_keys: &[String],
) -> Result<String> {
writer: &mut W,
) -> Result<()> {
let db = Omnigraph::open(uri).await?;
Ok(db.export_jsonl(branch, type_names, table_keys).await?)
db.export_jsonl_to_writer(branch, type_names, table_keys, writer)
.await?;
writer.flush()?;
Ok(())
}

async fn execute_export_remote(
async fn execute_export_remote_to_writer<W: Write>(
client: &reqwest::Client,
uri: &str,
branch: &str,
type_names: &[String],
table_keys: &[String],
bearer_token: Option<&str>,
) -> Result<String> {
remote_text(
client,
Method::POST,
remote_url(uri, "/export"),
Some(serde_json::to_value(ExportRequest {
branch: Some(branch.to_string()),
type_names: type_names.to_vec(),
table_keys: table_keys.to_vec(),
})?),
writer: &mut W,
) -> Result<()> {
let request = apply_bearer_token(
client.request(Method::POST, remote_url(uri, "/export")),
bearer_token,
)
.await
.json(&ExportRequest {
branch: Some(branch.to_string()),
type_names: type_names.to_vec(),
table_keys: table_keys.to_vec(),
});
let mut response = request.send().await?;
let status = response.status();
if !status.is_success() {
let text = response.text().await?;
if let Ok(error) = serde_json::from_str::<ErrorOutput>(&text) {
bail!(error.error);
}
bail!("server returned {}: {}", status, text);
}

while let Some(chunk) = response.chunk().await? {
writer.write_all(&chunk)?;
}
writer.flush()?;
Ok(())
}

#[tokio::main]
Expand Down Expand Up @@ -1479,8 +1479,10 @@ async fn main() -> Result<()> {
uri: &uri,
branch: &branch,
mode: mode.as_str(),
nodes_loaded: result.nodes_loaded.len(),
edges_loaded: result.edges_loaded.len(),
nodes_loaded: result.nodes_loaded.values().sum(),
edges_loaded: result.edges_loaded.values().sum(),
node_types_loaded: result.nodes_loaded.len(),
edge_types_loaded: result.edges_loaded.len(),
};
if json {
print_json(&payload)?;
Expand All @@ -1491,6 +1493,8 @@ async fn main() -> Result<()> {
mode,
payload.nodes_loaded,
payload.edges_loaded,
payload.node_types_loaded,
payload.edge_types_loaded,
);
}
}
Expand Down Expand Up @@ -1831,7 +1835,7 @@ async fn main() -> Result<()> {
target,
config,
branch,
jsonl: _,
jsonl,
type_names,
table_keys,
} => {
Expand All @@ -1840,20 +1844,27 @@ async fn main() -> Result<()> {
resolve_remote_bearer_token(&config, uri.as_deref(), target.as_deref())?;
let uri = resolve_uri(&config, uri, target.as_deref())?;
let branch = resolve_branch(&config, branch, None, "main");
let output = if is_remote_uri(&uri) {
execute_export_remote(
if jsonl {
eprintln!("warning: --jsonl is deprecated; `omnigraph export` always emits JSONL");
}

let stdout = io::stdout();
let mut stdout = stdout.lock();
if is_remote_uri(&uri) {
execute_export_remote_to_writer(
&http_client,
&uri,
&branch,
&type_names,
&table_keys,
bearer_token.as_deref(),
&mut stdout,
)
.await?
.await?;
} else {
execute_export(&uri, &branch, &type_names, &table_keys).await?
};
print!("{output}");
execute_export_to_writer(&uri, &branch, &type_names, &table_keys, &mut stdout)
.await?;
}
}
Command::Run { command } => match command {
RunCommand::List {
Expand Down Expand Up @@ -1984,6 +1995,7 @@ async fn main() -> Result<()> {
},
Command::Read {
uri,
legacy_uri,
target,
config,
alias,
Expand All @@ -2004,13 +2016,14 @@ async fn main() -> Result<()> {
let alias = resolve_alias(&config, alias.as_deref(), AliasCommand::Read)?;
let alias_name = alias.as_ref().map(|(name, _)| *name);
let alias_config = alias.as_ref().map(|(_, alias)| *alias);
let (uri, alias_args) = normalize_alias_args(
uri,
target.as_deref(),
config.cli_target_name().is_some(),
alias_name,
alias_args,
);
let target_available = target.is_some()
|| alias_config
.and_then(|alias| alias.target.as_deref())
.is_some()
|| config.cli_target_name().is_some();
let (legacy_uri, alias_args) =
normalize_legacy_alias_uri(legacy_uri, target_available, alias_name, alias_args);
let uri = uri.or(legacy_uri);
let target_name = target
.as_deref()
.or_else(|| alias_config.and_then(|alias| alias.target.as_deref()));
Expand Down Expand Up @@ -2067,6 +2080,7 @@ async fn main() -> Result<()> {
}
Command::Change {
uri,
legacy_uri,
target,
config,
alias,
Expand All @@ -2085,13 +2099,14 @@ async fn main() -> Result<()> {
let alias = resolve_alias(&config, alias.as_deref(), AliasCommand::Change)?;
let alias_name = alias.as_ref().map(|(name, _)| *name);
let alias_config = alias.as_ref().map(|(_, alias)| *alias);
let (uri, alias_args) = normalize_alias_args(
uri,
target.as_deref(),
config.cli_target_name().is_some(),
alias_name,
alias_args,
);
let target_available = target.is_some()
|| alias_config
.and_then(|alias| alias.target.as_deref())
.is_some()
|| config.cli_target_name().is_some();
let (legacy_uri, alias_args) =
normalize_legacy_alias_uri(legacy_uri, target_available, alias_name, alias_args);
let uri = uri.or(legacy_uri);
let target_name = target
.as_deref()
.or_else(|| alias_config.and_then(|alias| alias.target.as_deref()));
Expand Down
Loading
Loading