use crate::{
client::{
ApiCapability, Auth, KibanaClient, KibanaVersion, KibanaVersionInfo, parse_kibana_version,
},
etl::{Extractor, Loader, Transformer},
kibana::agents::{AgentEntry, AgentsExtractor, AgentsLoader, AgentsManifest},
kibana::dependencies::{
Dependency, find_agent_dependencies, find_skill_dependencies, find_tool_dependencies,
find_workflow_dependencies,
},
kibana::saved_objects::{SavedObjectsExtractor, SavedObjectsLoader},
kibana::skills::{
SkillEntry, SkillsExtractor, SkillsLoader, SkillsManifest, skill_to_directory,
skill_to_value,
},
kibana::spaces::{SpacesExtractor, SpacesLoader, SpacesManifest},
kibana::tools::{ToolsExtractor, ToolsLoader, ToolsManifest},
kibana::workflows::{
WorkflowEntry, WorkflowsExtractor, WorkflowsLoader, WorkflowsManifest,
workflow_resource_path,
},
storage::{self, DirectoryReader, DirectoryWriter},
transform::{
FieldDropper, FieldEscaper, FieldUnescaper, ManagedFlagAdder, MultilineFieldFormatter,
VegaSpecEscaper, VegaSpecUnescaper,
},
};
use eyre::{Context, Report, Result};
use owo_colors::OwoColorize;
use serde_json::Value;
use std::collections::HashSet;
use std::fmt;
use std::path::Path;
use std::sync::Arc;
use tokio::task::JoinSet;
use url::Url;
const SKILL_FETCH_BATCH_SIZE: usize = 16;
pub fn load_kibana_client(project_dir: impl AsRef<Path>) -> Result<KibanaClient> {
let url_str = std::env::var("KIBANA_URL").context("KIBANA_URL environment variable not set")?;
let url = Url::parse(&url_str).with_context(|| format!("Invalid KIBANA_URL: {}", url_str))?;
let auth = if let Ok(apikey) = std::env::var("KIBANA_APIKEY") {
Auth::Apikey(apikey)
} else if let (Ok(username), Ok(password)) = (
std::env::var("KIBANA_USERNAME"),
std::env::var("KIBANA_PASSWORD"),
) {
Auth::Basic(username, password)
} else {
Auth::None
};
let max_requests = std::env::var("KIBANA_MAX_REQUESTS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(8);
let spaces_manifest_path = project_dir.as_ref().join("spaces.yml");
let spaces = if spaces_manifest_path.exists() {
log::debug!("Loading spaces from {}", spaces_manifest_path.display());
SpacesManifest::read(&spaces_manifest_path)
.context("Failed to load spaces manifest")?
.spaces
.into_iter()
.map(|space| (space.id, space.name))
.collect::<Vec<_>>()
} else {
log::debug!("No spaces manifest found, defaulting to 'default' space");
[("default".to_string(), "Default".to_string())]
.into_iter()
.collect::<Vec<_>>()
};
KibanaClient::builder(url)
.auth(auth)
.max_concurrency(max_requests)
.spaces(spaces)
.build()
.map_err(Report::new)
.context("Failed to create Kibana client")
}
#[derive(Debug)]
pub struct VersionWarning {
message: String,
}
impl VersionWarning {
fn new(message: impl Into<String>) -> Self {
Self {
message: message.into(),
}
}
}
impl fmt::Display for VersionWarning {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.message)
}
}
impl std::error::Error for VersionWarning {}
fn version_warning(message: impl Into<String>) -> Report {
Report::new(VersionWarning::new(message))
}
pub fn version_warning_message(report: &Report) -> Option<&str> {
report
.downcast_ref::<VersionWarning>()
.map(|warning| warning.message.as_str())
}
#[derive(Debug, Clone)]
struct VersionPreflight {
detected: KibanaVersionInfo,
recorded: Option<KibanaVersion>,
}
fn parse_capability(api: &str) -> Option<ApiCapability> {
match api {
"saved_objects" | "saved_object" | "objects" | "object" => {
Some(ApiCapability::SavedObjects)
}
"workflows" | "workflow" => Some(ApiCapability::Workflows),
"agents" | "agent" => Some(ApiCapability::Agents),
"tools" | "tool" => Some(ApiCapability::Tools),
"skills" | "skill" => Some(ApiCapability::Skills),
"spaces" | "space" => Some(ApiCapability::Spaces),
_ => None,
}
}
fn capability_requested(capability: ApiCapability, filter: Option<&[String]>) -> bool {
match filter {
None => true,
Some(values) => values
.iter()
.filter_map(|v| parse_capability(&v.to_lowercase()))
.any(|candidate| candidate == capability),
}
}
fn workflow_file_stem(name: &str) -> String {
let sanitized = storage::sanitize_filename(name);
let stem = sanitized
.split_whitespace()
.collect::<Vec<_>>()
.join("_")
.to_lowercase();
if stem.is_empty() {
"unnamed".to_string()
} else {
stem
}
}
fn is_older_major_minor(target: &KibanaVersion, recorded: &KibanaVersion) -> bool {
target.major < recorded.major
|| (target.major == recorded.major && target.minor < recorded.minor)
}
fn load_recorded_kibana_version(project_dir: &Path) -> Result<Option<KibanaVersion>> {
let spaces_manifest_path = project_dir.join("spaces.yml");
if !spaces_manifest_path.exists() {
return Ok(None);
}
let manifest = SpacesManifest::read(&spaces_manifest_path)?;
match manifest.kibana_version() {
Some(version) => match parse_kibana_version(version) {
Ok(parsed) => Ok(Some(parsed)),
Err(e) => {
log::warn!(
"Ignoring invalid spaces.yml kibana.version '{}': {}",
version,
e
);
Ok(None)
}
},
None => Ok(None),
}
}
fn persist_kibana_version(project_dir: &Path, version: &str) -> Result<()> {
let spaces_manifest_path = project_dir.join("spaces.yml");
let mut manifest = if spaces_manifest_path.exists() {
SpacesManifest::read(&spaces_manifest_path)?
} else {
let mut created = SpacesManifest::new();
created.add_space("default".to_string(), "Default".to_string());
created
};
manifest.set_kibana_version(version.to_string());
manifest.write(&spaces_manifest_path)?;
Ok(())
}
async fn run_version_preflight(
project_dir: &Path,
client: &KibanaClient,
action: &str,
force: bool,
) -> Result<VersionPreflight> {
let detected = client.server_version_info().await?;
let recorded = load_recorded_kibana_version(project_dir)?;
if let Some(recorded_version) = recorded.as_ref()
&& is_older_major_minor(&detected.parsed, recorded_version)
{
let message = format!(
"{} requires --force because target Kibana {} is older than recorded repository version {}",
action, detected.parsed, recorded_version
);
if force {
log::warn!("{}", message);
} else {
return Err(version_warning(message));
}
}
Ok(VersionPreflight { detected, recorded })
}
fn should_process_api(api_name: &str, filter: Option<&[String]>) -> bool {
match filter {
None => true, Some(apis) => apis.iter().any(|api| {
let api = api.to_lowercase();
match api_name {
"saved_objects" => {
api == "saved_objects"
|| api == "objects"
|| api == "object"
|| api == "saved_object"
}
"workflows" => api == "workflows" || api == "workflow",
"agents" => api == "agents" || api == "agent",
"tools" => api == "tools" || api == "tool",
"skills" => api == "skills" || api == "skill",
"spaces" => api == "spaces" || api == "space",
_ => false,
}
}),
}
}
pub async fn pull_saved_objects(
project_dir: impl AsRef<Path>,
space_filter: Option<&[String]>,
api_filter: Option<&[String]>,
force: bool,
) -> Result<usize> {
let project_dir = Arc::<Path>::from(project_dir.as_ref());
let api_filter = api_filter.map(|f| Arc::new(f.to_vec()));
log::info!("Connecting to Kibana...");
let client = load_kibana_client(&*project_dir)?;
let preflight = run_version_preflight(&project_dir, &client, "pull", force).await?;
let detected_version = preflight.detected.parsed;
let mut unsupported = Vec::new();
for capability in [
ApiCapability::Agents,
ApiCapability::Tools,
ApiCapability::Skills,
ApiCapability::Workflows,
] {
if capability_requested(capability, api_filter.as_deref().map(|f| f.as_slice()))
&& !KibanaClient::supports_capability(&detected_version, capability)
{
unsupported.push(KibanaClient::unsupported_capability_reason(
&detected_version,
capability,
));
}
}
let warning_exit = !force && !unsupported.is_empty();
for warning in &unsupported {
log::warn!("{}", warning);
}
if force && !unsupported.is_empty() {
log::warn!("--force enabled: attempting API calls despite unsupported version checks");
}
let can_pull_spaces = !warning_exit
|| KibanaClient::supports_capability(&detected_version, ApiCapability::Spaces)
|| force;
if can_pull_spaces && should_process_api("spaces", api_filter.as_deref().map(|f| f.as_slice()))
{
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
log::info!("Pulling space definitions...");
match pull_spaces_internal(&*project_dir, &client).await {
Ok(space_count) => {
log::info!("✓ Pulled {} space definition(s)", space_count);
}
Err(e) => {
log::warn!("Failed to pull space definitions: {}", e);
}
}
}
} else {
log::debug!("Skipping spaces pull (filtered out)");
}
let target_space_ids = get_target_space_ids(&client, space_filter);
let mut total_saved_objects = 0;
let mut total_workflows = 0;
let mut total_agents = 0;
let mut total_tools = 0;
let mut total_skills = 0;
let mut set = JoinSet::new();
for space_id in target_space_ids {
let client = client.clone();
let project_dir = project_dir.clone();
let api_filter = api_filter.clone();
let detected_version = detected_version.clone();
set.spawn(async move {
log::info!("Processing space: {}", space_id.cyan());
let mut s_so = 0;
let mut s_wf = 0;
let mut s_ag = 0;
let mut s_tl = 0;
let mut s_sk = 0;
let space_client = client.space(&space_id)?;
let api_filter_slice = api_filter.as_deref().map(|f| f.as_slice());
let can_pull_workflows = force
|| KibanaClient::supports_capability(&detected_version, ApiCapability::Workflows);
let can_pull_agents = force
|| KibanaClient::supports_capability(&detected_version, ApiCapability::Agents);
let can_pull_tools =
force || KibanaClient::supports_capability(&detected_version, ApiCapability::Tools);
let can_pull_skills = force
|| KibanaClient::supports_capability(&detected_version, ApiCapability::Skills);
if should_process_api("saved_objects", api_filter_slice)
&& let Ok(count) = pull_space_saved_objects(&project_dir, &space_client).await
{
s_so = count;
}
if can_pull_workflows
&& should_process_api("workflows", api_filter_slice)
&& let Ok(count) = pull_space_workflows(&project_dir, &space_client).await
{
s_wf = count;
log::debug!("Pulled {} workflow(s) for space {}", count, space_id.cyan());
}
if can_pull_agents
&& should_process_api("agents", api_filter_slice)
&& let Ok(count) = pull_space_agents(&project_dir, &space_client).await
{
s_ag = count;
log::debug!("Pulled {} agent(s) for space {}", count, space_id.cyan());
}
if can_pull_tools
&& should_process_api("tools", api_filter_slice)
&& let Ok(count) = pull_space_tools(&project_dir, &space_client).await
{
s_tl = count;
log::debug!("Pulled {} tool(s) for space {}", count, space_id.cyan());
}
if can_pull_skills
&& should_process_api("skills", api_filter_slice)
&& let Ok(count) = pull_space_skills(&project_dir, &space_client).await
{
s_sk = count;
log::debug!("Pulled {} skill(s) for space {}", count, space_id.cyan());
}
Ok::<(usize, usize, usize, usize, usize), eyre::Report>((s_so, s_wf, s_ag, s_tl, s_sk))
});
}
while let Some(res) = set.join_next().await {
match res {
Ok(Ok((so, wf, ag, tl, sk))) => {
total_saved_objects += so;
total_workflows += wf;
total_agents += ag;
total_tools += tl;
total_skills += sk;
}
Ok(Err(e)) => log::error!("Space processing failed: {}", e),
Err(e) => log::error!("Task panicked: {}", e),
}
}
log::info!(
"✓ Pull complete: {} saved object(s), {} workflow(s), {} agent(s), {} tool(s), {} skill(s)",
total_saved_objects,
total_workflows,
total_agents,
total_tools,
total_skills
);
persist_kibana_version(&project_dir, &preflight.detected.raw)?;
if let Some(recorded) = preflight.recorded {
log::debug!("Repository kibana.version before pull: {}", recorded);
}
if warning_exit {
return Err(version_warning(
"One or more requested APIs are unsupported for the target Kibana version".to_string(),
));
}
Ok(total_saved_objects + total_workflows + total_agents + total_tools + total_skills)
}
pub async fn push_saved_objects(
project_dir: impl AsRef<Path>,
managed: bool,
space_filter: Option<&[String]>,
api_filter: Option<&[String]>,
force: bool,
) -> Result<usize> {
let project_dir = Arc::<Path>::from(project_dir.as_ref());
let api_filter = api_filter.map(|f| Arc::new(f.to_vec()));
log::info!("Connecting to Kibana...");
let client = load_kibana_client(&*project_dir)?;
let preflight = run_version_preflight(&project_dir, &client, "push", force).await?;
let detected_version = preflight.detected.parsed;
let mut unsupported = Vec::new();
for capability in [
ApiCapability::Agents,
ApiCapability::Tools,
ApiCapability::Skills,
ApiCapability::Workflows,
] {
if capability_requested(capability, api_filter.as_deref().map(|f| f.as_slice()))
&& !KibanaClient::supports_capability(&detected_version, capability)
{
unsupported.push(KibanaClient::unsupported_capability_reason(
&detected_version,
capability,
));
}
}
let warning_exit = !force && !unsupported.is_empty();
for warning in &unsupported {
log::warn!("{}", warning);
}
if force && !unsupported.is_empty() {
log::warn!("--force enabled: attempting API calls despite unsupported version checks");
}
if should_process_api("spaces", api_filter.as_deref().map(|f| f.as_slice())) {
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
log::info!("Pushing space definitions...");
match push_spaces_internal(&*project_dir, &client).await {
Ok(space_count) => {
log::info!("✓ Pushed {} space definition(s)", space_count);
}
Err(e) => {
log::warn!("Failed to push space definitions: {}", e);
}
}
}
} else {
log::debug!("Skipping spaces push (filtered out)");
}
let target_space_ids = get_target_space_ids(&client, space_filter);
let mut total_saved_objects = 0;
let mut total_workflows = 0;
let mut total_agents = 0;
let mut total_tools = 0;
let mut total_skills = 0;
let mut set = JoinSet::new();
for space_id in target_space_ids {
let client = client.clone();
let project_dir = project_dir.clone();
let api_filter = api_filter.clone();
let detected_version = detected_version.clone();
set.spawn(async move {
log::info!("Processing space: {}", space_id.cyan());
let mut s_so = 0;
let mut s_wf = 0;
let mut s_ag = 0;
let mut s_tl = 0;
let mut s_sk = 0;
let space_client = client.space(&space_id)?;
let api_filter_slice = api_filter.as_deref().map(|f| f.as_slice());
let can_push_workflows = force
|| KibanaClient::supports_capability(&detected_version, ApiCapability::Workflows);
let can_push_agents = force
|| KibanaClient::supports_capability(&detected_version, ApiCapability::Agents);
let can_push_tools =
force || KibanaClient::supports_capability(&detected_version, ApiCapability::Tools);
let can_push_skills = force
|| KibanaClient::supports_capability(&detected_version, ApiCapability::Skills);
if should_process_api("saved_objects", api_filter_slice) {
match push_space_saved_objects(&project_dir, &space_client, managed).await {
Ok(count) => {
s_so = count;
}
Err(e) => {
log::warn!(
"Failed to push saved objects for space {}: {}",
space_id.cyan(),
e
);
}
}
}
if can_push_workflows && should_process_api("workflows", api_filter_slice) {
match push_space_workflows(&project_dir, &space_client).await {
Ok(count) => {
s_wf = count;
}
Err(e) => {
log::warn!(
"Failed to push workflows for space {}: {}",
space_id.cyan(),
e
);
}
}
}
if can_push_tools && should_process_api("tools", api_filter_slice) {
match push_space_tools(&project_dir, &space_client).await {
Ok(count) => {
s_tl = count;
}
Err(e) => {
log::warn!("Failed to push tools for space {}: {}", space_id.cyan(), e);
}
}
}
if can_push_skills && should_process_api("skills", api_filter_slice) {
match push_space_skills(&project_dir, &space_client).await {
Ok(count) => {
s_sk = count;
}
Err(e) => {
log::warn!("Failed to push skills for space {}: {}", space_id.cyan(), e);
}
}
}
if can_push_agents && should_process_api("agents", api_filter_slice) {
match push_space_agents(&project_dir, &space_client).await {
Ok(count) => {
s_ag = count;
}
Err(e) => {
log::warn!("Failed to push agents for space {}: {}", space_id.cyan(), e);
}
}
}
Ok::<(usize, usize, usize, usize, usize), eyre::Report>((s_so, s_wf, s_ag, s_tl, s_sk))
});
}
while let Some(res) = set.join_next().await {
match res {
Ok(Ok((so, wf, ag, tl, sk))) => {
total_saved_objects += so;
total_workflows += wf;
total_agents += ag;
total_tools += tl;
total_skills += sk;
}
Ok(Err(e)) => log::error!("Space processing failed: {}", e),
Err(e) => log::error!("Task panicked: {}", e),
}
}
log::info!(
"✓ Push complete: {} saved object(s), {} workflow(s), {} agent(s), {} tool(s), {} skill(s)",
total_saved_objects,
total_workflows,
total_agents,
total_tools,
total_skills
);
if warning_exit {
return Err(version_warning(
"One or more requested APIs are unsupported for the target Kibana version".to_string(),
));
}
Ok(total_saved_objects + total_workflows + total_agents + total_tools + total_skills)
}
pub async fn bundle_to_ndjson(
project_dir: impl AsRef<Path>,
_output_file: impl AsRef<Path>,
managed: bool,
space_filter: Option<&[String]>,
api_filter: Option<&[String]>,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let recorded_version = load_recorded_kibana_version(project_dir)?;
let target_space_ids = get_target_space_ids_from_manifest(project_dir, space_filter);
let mut total_saved_objects = 0;
let mut total_workflows = 0;
let mut total_agents = 0;
let mut total_tools = 0;
let mut total_skills = 0;
let can_bundle_workflows = recorded_version
.as_ref()
.map(|v| KibanaClient::supports_capability(v, ApiCapability::Workflows))
.unwrap_or(true);
let can_bundle_agents = recorded_version
.as_ref()
.map(|v| KibanaClient::supports_capability(v, ApiCapability::Agents))
.unwrap_or(true);
let can_bundle_tools = recorded_version
.as_ref()
.map(|v| KibanaClient::supports_capability(v, ApiCapability::Tools))
.unwrap_or(true);
let can_bundle_skills = recorded_version
.as_ref()
.map(|v| KibanaClient::supports_capability(v, ApiCapability::Skills))
.unwrap_or(true);
if let Some(version) = recorded_version.as_ref() {
if capability_requested(ApiCapability::Workflows, api_filter) && !can_bundle_workflows {
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(version, ApiCapability::Workflows)
);
}
if capability_requested(ApiCapability::Agents, api_filter) && !can_bundle_agents {
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(version, ApiCapability::Agents)
);
}
if capability_requested(ApiCapability::Tools, api_filter) && !can_bundle_tools {
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(version, ApiCapability::Tools)
);
}
if capability_requested(ApiCapability::Skills, api_filter) && !can_bundle_skills {
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(version, ApiCapability::Skills)
);
}
}
for space_id in &target_space_ids {
log::info!("Bundling space: {}", space_id.cyan());
if should_process_api("saved_objects", api_filter)
&& let Ok(count) = bundle_space_saved_objects(project_dir, space_id, managed).await
{
total_saved_objects += count;
}
if can_bundle_workflows
&& should_process_api("workflows", api_filter)
&& let Ok(count) = bundle_space_workflows(project_dir, space_id).await
{
total_workflows += count;
log::debug!(
"Bundled {} workflow(s) for space {}",
count,
space_id.cyan()
);
}
if can_bundle_tools
&& should_process_api("tools", api_filter)
&& let Ok(count) = bundle_space_tools(project_dir, space_id).await
{
total_tools += count;
log::debug!("Bundled {} tool(s) for space {}", count, space_id.cyan());
}
if can_bundle_skills
&& should_process_api("skills", api_filter)
&& let Ok(count) = bundle_space_skills(project_dir, space_id).await
{
total_skills += count;
log::debug!("Bundled {} skill(s) for space {}", count, space_id.cyan());
}
if can_bundle_agents
&& should_process_api("agents", api_filter)
&& let Ok(count) = bundle_space_agents(project_dir, space_id).await
{
total_agents += count;
log::debug!("Bundled {} agent(s) for space {}", count, space_id.cyan());
}
}
if should_process_api("spaces", api_filter) {
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
log::info!("Bundling space definitions...");
let spaces_output = project_dir.join("bundle/spaces.ndjson");
match bundle_spaces_to_ndjson_internal(project_dir, &spaces_output).await {
Ok(space_count) => {
log::info!("✓ Bundled {} space definition(s)", space_count);
}
Err(e) => {
log::warn!("Failed to bundle space definitions: {}", e);
}
}
}
}
log::info!(
"✓ Bundle complete: {} saved object(s), {} workflow(s), {} agent(s), {} tool(s), {} skill(s)",
total_saved_objects,
total_workflows,
total_agents,
total_tools,
total_skills
);
Ok(total_saved_objects + total_workflows + total_agents + total_tools + total_skills)
}
pub async fn init_from_export(
export_path: impl AsRef<Path>,
manifest_dir: impl AsRef<Path>,
) -> Result<usize> {
use crate::kibana::saved_objects::{SavedObject, SavedObjectsManifest};
use std::io::{BufRead, BufReader};
let export_path = export_path.as_ref();
let manifest_dir = manifest_dir.as_ref();
log::info!("Reading export from {}", export_path.display());
let file = std::fs::File::open(export_path)?;
let reader = BufReader::new(file);
let mut objects = Vec::new();
let mut saved_objects = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let obj: serde_json::Value = serde_json::from_str(&line)?;
if let (Some(obj_type), Some(obj_id)) = (
obj.get("type").and_then(|v| v.as_str()),
obj.get("id").and_then(|v| v.as_str()),
) {
saved_objects.push(SavedObject::new(obj_type, obj_id));
}
objects.push(obj);
}
log::info!("Read {} object(s) from export", objects.len());
let manifest = SavedObjectsManifest::with_objects(saved_objects);
let manifest_path = manifest_dir.join("manifest");
std::fs::create_dir_all(&manifest_path)?;
let manifest_file = manifest_path.join("saved_objects.json");
let manifest_json = serde_json::to_string_pretty(&manifest)?;
std::fs::write(&manifest_file, manifest_json)?;
log::info!("✓ Created manifest with {} object(s)", manifest.count());
log::info!(" Manifest: {}", manifest_file.display());
let objects_dir = manifest_dir.join("objects");
let writer = DirectoryWriter::new_with_options(&objects_dir, true)?;
let drop_fields = FieldDropper::default_kibana_fields();
let unescape = FieldUnescaper::default_kibana_fields();
let vega_unescape = VegaSpecUnescaper::new();
let dropped = drop_fields.transform_many(objects)?;
let unescaped = unescape.transform_many(dropped)?;
let vega_unescaped = vega_unescape.transform_many(unescaped)?;
use crate::etl::Loader;
let count = writer.load(vega_unescaped).await?;
log::info!(" Objects: {}", objects_dir.display());
log::info!("✓ Wrote {} object files", count);
Ok(count)
}
pub async fn add_objects_to_manifest(
project_dir: impl AsRef<Path>,
space_id: &str,
objects_to_add: Option<Vec<String>>,
file_path: Option<impl AsRef<Path>>,
force: bool,
) -> Result<usize> {
use crate::kibana::saved_objects::{SavedObject, SavedObjectsExtractor, SavedObjectsManifest};
let project_dir = project_dir.as_ref();
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
let spaces_manifest = SpacesManifest::read(&spaces_manifest_path)?;
if !spaces_manifest.spaces.iter().any(|s| s.id == space_id) {
eyre::bail!(
"Space {} is not managed. Add it first with: kibob add spaces . --include '^{}$'",
space_id.cyan(),
space_id
);
}
}
log::info!("Loading existing manifest from {}", project_dir.display());
let manifest_path = get_space_saved_objects_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
SavedObjectsManifest::read(&manifest_path)?
} else {
log::info!("No existing manifest found, will create new one");
if let Some(parent) = manifest_path.parent() {
std::fs::create_dir_all(parent)?;
}
SavedObjectsManifest::new()
};
log::info!("Current manifest has {} objects", manifest.count());
let new_objects = if let Some(file) = file_path {
let file_path = file.as_ref();
log::info!("Reading objects from {}", file_path.display());
use std::io::{BufRead, BufReader};
let file = std::fs::File::open(file_path)?;
let reader = BufReader::new(file);
let mut objs = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let obj: serde_json::Value = serde_json::from_str(&line)?;
objs.push(obj);
}
objs
} else if let Some(object_specs) = objects_to_add {
log::info!("Fetching {} object(s) from Kibana", object_specs.len());
let client = load_kibana_client(project_dir)?;
let _preflight = run_version_preflight(project_dir, &client, "add", force).await?;
let space_client = client.space(space_id)?;
let mut saved_objects = Vec::new();
for spec in &object_specs {
let parts: Vec<&str> = spec.split(['=', ':']).collect();
if parts.len() != 2 {
eyre::bail!(
"Invalid object spec: {}. Expected format: type=id or type:id",
spec
);
}
saved_objects.push(SavedObject::new(parts[0], parts[1]));
}
let temp_manifest = SavedObjectsManifest::with_objects(saved_objects);
let extractor = SavedObjectsExtractor::new(space_client, temp_manifest);
extractor.extract().await?
} else {
eyre::bail!("Must specify either --objects or --file");
};
log::info!("Adding {} new object(s)", new_objects.len());
for obj in &new_objects {
if let (Some(obj_type), Some(obj_id)) = (
obj.get("type").and_then(|v| v.as_str()),
obj.get("id").and_then(|v| v.as_str()),
) {
manifest.add_object(SavedObject::new(obj_type, obj_id));
}
}
let manifest_json = serde_json::to_string_pretty(&manifest)?;
std::fs::write(&manifest_path, manifest_json)?;
log::info!("✓ Updated manifest now has {} objects", manifest.count());
let objects_dir = get_space_objects_dir(project_dir, space_id);
let writer = DirectoryWriter::new_with_options(&objects_dir, true)?;
let drop_fields = FieldDropper::default_kibana_fields();
let unescape = FieldUnescaper::default_kibana_fields();
let vega_unescape = VegaSpecUnescaper::new();
let dropped = drop_fields.transform_many(new_objects)?;
let unescaped = unescape.transform_many(dropped)?;
let vega_unescaped = vega_unescape.transform_many(unescaped)?;
let count = writer.load(vega_unescaped).await?;
log::info!("✓ Wrote {} new object file(s)", count);
Ok(count)
}
#[allow(clippy::too_many_arguments)]
pub async fn add_workflows_to_manifest(
project_dir: impl AsRef<Path>,
space_id: &str,
query: Option<String>,
include: Option<String>,
exclude: Option<String>,
file_path: Option<String>,
exclude_dependencies: bool,
force: bool,
) -> Result<usize> {
use crate::kibana::workflows::WorkflowEntry;
let project_dir = project_dir.as_ref();
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
let manifest = SpacesManifest::read(&spaces_manifest_path)?;
if !manifest.spaces.iter().any(|s| s.id == space_id) {
eyre::bail!(
"Space {} is not managed. Add it first with: kibob add spaces . --include '^{}$'",
space_id.cyan(),
space_id
);
}
}
log::info!(
"Loading workflows manifest for space {} from {}",
space_id.cyan(),
project_dir.display()
);
let manifest_path = get_space_workflows_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
WorkflowsManifest::read(&manifest_path)?
} else {
log::info!("No existing manifest found, will create new one");
WorkflowsManifest::new()
};
log::info!("Current manifest has {} workflow(s)", manifest.count());
let requires_kibana_calls = file_path.is_none() || !exclude_dependencies;
let mut preflight_client: Option<KibanaClient> = None;
let mut detected_version: Option<KibanaVersion> = None;
if requires_kibana_calls {
let client = load_kibana_client(project_dir)?;
let preflight = run_version_preflight(project_dir, &client, "add", force).await?;
detected_version = Some(preflight.detected.parsed);
preflight_client = Some(client);
}
if let Some(version) = detected_version
&& !KibanaClient::supports_capability(&version, ApiCapability::Workflows)
{
let reason =
KibanaClient::unsupported_capability_reason(&version, ApiCapability::Workflows);
if force {
log::warn!("{}", reason);
log::warn!("--force enabled: attempting workflows API calls anyway");
} else if file_path.is_none() {
return Err(version_warning(reason));
} else {
log::warn!("{}", reason);
}
}
let new_workflows: Vec<serde_json::Value> = if let Some(file) = file_path {
log::info!("Reading workflows from {}", file);
let file_path = std::path::Path::new(&file);
if !file_path.exists() {
eyre::bail!("File not found: {}", file_path.display());
}
let extension = file_path.extension().and_then(|s| s.to_str()).unwrap_or("");
match extension {
"ndjson" => {
use std::io::{BufRead, BufReader};
let file = std::fs::File::open(file_path)?;
let reader = BufReader::new(file);
let mut workflows = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let workflow: serde_json::Value = serde_json::from_str(&line)?;
workflows.push(workflow);
}
log::info!("Read {} workflow(s) from NDJSON file", workflows.len());
workflows
}
"json" => {
let content = std::fs::read_to_string(file_path)?;
let parsed = storage::from_json5_str(&content)?;
if let Some(results) = parsed.get("results").and_then(|v| v.as_array()) {
log::info!("Read {} workflow(s) from JSON API response", results.len());
results.to_vec()
} else if let Some(arr) = parsed.as_array() {
log::info!("Read {} workflow(s) from JSON array", arr.len());
arr.to_vec()
} else {
log::info!("Read 1 workflow from JSON file");
vec![parsed]
}
}
_ => {
eyre::bail!(
"Unsupported file format: {}. Expected .json or .ndjson",
extension
);
}
}
} else {
log::info!(
"Searching workflows via API in space {}...",
space_id.cyan()
);
let client = preflight_client
.clone()
.unwrap_or(load_kibana_client(project_dir)?);
let space_client = client.space(space_id)?;
let extractor = WorkflowsExtractor::new(space_client, None);
extractor.search_workflows(query.as_deref(), None).await?
};
log::info!("Found {} workflow(s) before filtering", new_workflows.len());
let filtered_workflows: Vec<serde_json::Value> = {
let mut workflows = new_workflows;
if let Some(include_pattern) = &include {
let regex = regex::Regex::new(include_pattern)
.with_context(|| format!("Invalid include regex pattern: {}", include_pattern))?;
workflows.retain(|w| {
w.get("name")
.and_then(|v| v.as_str())
.map(|name| regex.is_match(name))
.unwrap_or(false)
});
log::info!(
"After include filter '{}': {} workflow(s)",
include_pattern,
workflows.len()
);
}
if let Some(exclude_pattern) = &exclude {
let regex = regex::Regex::new(exclude_pattern)
.with_context(|| format!("Invalid exclude regex pattern: {}", exclude_pattern))?;
workflows.retain(|w| {
w.get("name")
.and_then(|v| v.as_str())
.map(|name| !regex.is_match(name))
.unwrap_or(true)
});
log::info!(
"After exclude filter '{}': {} workflow(s)",
exclude_pattern,
workflows.len()
);
}
workflows
};
log::info!(
"Adding {} workflow(s) after filtering",
filtered_workflows.len()
);
let workflows_dir = get_space_workflows_dir(project_dir, space_id);
std::fs::create_dir_all(&workflows_dir)?;
let mut added_count = 0;
let mut all_deps = Vec::new();
for workflow in &filtered_workflows {
let workflow_id = workflow
.get("id")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Workflow missing 'id' field"))?;
let workflow_name = workflow
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Workflow missing 'name' field"))?;
if manifest.add_workflow(WorkflowEntry::new(workflow_id, workflow_name)) {
log::debug!(
"Added workflow to manifest: {} ({})",
workflow_name,
workflow_id
);
let workflow_file =
workflows_dir.join(format!("{}.json", workflow_file_stem(workflow_name)));
let json = storage::to_string_with_multiline(workflow)?;
std::fs::write(&workflow_file, json)?;
log::debug!("Wrote workflow file: {}", workflow_file.display());
added_count += 1;
if !exclude_dependencies {
all_deps.extend(find_workflow_dependencies(workflow));
}
} else {
log::debug!("Workflow already in manifest, skipping: {}", workflow_name);
}
}
let mut dep_summary = DependencySummary::new();
if !all_deps.is_empty() {
let client = preflight_client.unwrap_or(load_kibana_client(project_dir)?);
let version = client.server_version().await?;
dep_summary =
resolve_and_add_dependencies(project_dir, space_id, &client, all_deps, version, force)
.await?;
}
let manifest_dir = get_space_manifest_dir(project_dir, space_id);
std::fs::create_dir_all(&manifest_dir)?;
manifest.write(&manifest_path)?;
log::info!(
"✓ Updated manifest for space {} now has {} workflow(s)",
space_id.cyan(),
manifest.count()
);
if !dep_summary.is_empty() {
log::info!(
"✓ Added {} new workflow(s) (plus dependencies: {})",
added_count,
dep_summary.format_summary()
);
} else {
log::info!("✓ Added {} new workflow(s)", added_count);
}
Ok(added_count + dep_summary.total())
}
pub async fn add_spaces_to_manifest(
project_dir: impl AsRef<Path>,
space_filter: Option<&[String]>,
query: Option<String>,
include: Option<String>,
exclude: Option<String>,
file_path: Option<String>,
force: bool,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading spaces manifest from {}", project_dir.display());
let manifest_path = project_dir.join("spaces.yml");
let mut manifest = if manifest_path.exists() {
SpacesManifest::read(&manifest_path)?
} else {
log::info!("No existing manifest found, will create new one");
SpacesManifest::new()
};
log::info!("Current manifest has {} space(s)", manifest.count());
let new_spaces: Vec<serde_json::Value> = if let Some(file) = file_path {
log::info!("Reading spaces from {}", file);
let file_path = std::path::Path::new(&file);
if !file_path.exists() {
eyre::bail!("File not found: {}", file_path.display());
}
let extension = file_path.extension().and_then(|s| s.to_str()).unwrap_or("");
match extension {
"ndjson" => {
use std::io::{BufRead, BufReader};
let file = std::fs::File::open(file_path)?;
let reader = BufReader::new(file);
let mut spaces = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let space: serde_json::Value = serde_json::from_str(&line)?;
spaces.push(space);
}
log::info!("Read {} space(s) from NDJSON file", spaces.len());
spaces
}
"json" => {
let content = std::fs::read_to_string(file_path)?;
let parsed = storage::from_json5_str(&content)?;
if let Some(arr) = parsed.as_array() {
log::info!("Read {} space(s) from JSON array", arr.len());
arr.to_vec()
} else {
log::info!("Read 1 space from JSON file");
vec![parsed]
}
}
_ => {
eyre::bail!(
"Unsupported file format: {}. Expected .json or .ndjson",
extension
);
}
}
} else {
if query.is_some() {
log::warn!("Spaces API doesn't support query filtering - fetching all spaces");
}
log::info!("Fetching spaces via API...");
let client = load_kibana_client(project_dir)?;
let _preflight = run_version_preflight(project_dir, &client, "add", force).await?;
let extractor = SpacesExtractor::new(client, None);
extractor.search_spaces(query.as_deref()).await?
};
log::info!("Found {} space(s) before filtering", new_spaces.len());
let filtered_spaces: Vec<serde_json::Value> = {
let mut spaces = new_spaces;
if let Some(include_pattern) = &include {
let regex = regex::Regex::new(include_pattern)
.with_context(|| format!("Invalid include regex pattern: {}", include_pattern))?;
spaces.retain(|s| {
s.get("name")
.and_then(|v| v.as_str())
.map(|name| regex.is_match(name))
.unwrap_or(false)
});
log::info!(
"After include filter '{}': {} space(s)",
include_pattern,
spaces.len()
);
}
if let Some(exclude_pattern) = &exclude {
let regex = regex::Regex::new(exclude_pattern)
.with_context(|| format!("Invalid exclude regex pattern: {}", exclude_pattern))?;
spaces.retain(|s| {
s.get("name")
.and_then(|v| v.as_str())
.map(|name| !regex.is_match(name))
.unwrap_or(true)
});
log::info!(
"After exclude filter '{}': {} space(s)",
exclude_pattern,
spaces.len()
);
}
if let Some(filter_ids) = space_filter {
spaces.retain(|s| {
s.get("id")
.and_then(|v| v.as_str())
.map(|id| filter_ids.contains(&id.to_string()))
.unwrap_or(false)
});
log::info!("After space ID filter: {} space(s)", spaces.len());
}
spaces
};
log::info!("Adding {} space(s) after filtering", filtered_spaces.len());
let mut added_count = 0;
for space in &filtered_spaces {
let space_id = space
.get("id")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Space missing 'id' field"))?;
let space_name = space
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Space missing 'name' field"))?;
if manifest.add_space(space_id.to_string(), space_name.to_string()) {
log::debug!("Added space to manifest: {} ({})", space_id, space_name);
let space_file = get_space_file(project_dir, space_id);
if let Some(parent) = space_file.parent() {
std::fs::create_dir_all(parent)?;
}
let json = serde_json::to_string_pretty(space)?;
std::fs::write(&space_file, json)?;
log::debug!("Wrote space file: {}", space_file.display());
added_count += 1;
} else {
log::debug!("Space already in manifest, skipping: {}", space_id);
}
}
manifest.write(&manifest_path)?;
log::info!("✓ Updated manifest now has {} space(s)", manifest.count());
log::info!("✓ Added {} new space(s)", added_count);
Ok(added_count)
}
fn load_spaces_manifest(project_dir: impl AsRef<Path>) -> Result<SpacesManifest> {
let manifest_path = project_dir.as_ref().join("spaces.yml");
if !manifest_path.exists() {
eyre::bail!("Spaces manifest not found: {}", manifest_path.display());
}
Ok(SpacesManifest::read(&manifest_path)?)
}
fn get_target_space_ids(client: &KibanaClient, space_filter: Option<&[String]>) -> Vec<String> {
if let Some(filter) = space_filter {
filter.to_vec()
} else {
client.space_ids().into_iter().map(String::from).collect()
}
}
fn get_target_space_ids_from_manifest(
project_dir: &Path,
space_filter: Option<&[String]>,
) -> Vec<String> {
if let Some(filter) = space_filter {
return filter.to_vec();
}
let manifest_path = project_dir.join("spaces.yml");
if manifest_path.exists()
&& let Ok(manifest) = SpacesManifest::read(&manifest_path)
{
return manifest.spaces.into_iter().map(|s| s.id).collect();
}
vec!["default".to_string()]
}
async fn pull_spaces_internal(
project_dir: impl AsRef<Path>,
client: &KibanaClient,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::debug!("Loading spaces manifest from {}", project_dir.display());
let manifest = load_spaces_manifest(project_dir)?;
log::debug!("Manifest loaded: {} space(s)", manifest.count());
log::debug!("Extracting spaces from Kibana...");
let extractor = SpacesExtractor::new(client.clone(), Some(manifest));
let spaces = extractor.extract().await?;
let mut count = 0;
for space in &spaces {
let space_id = space
.get("id")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Space missing 'id' field"))?;
let space_file = get_space_file(project_dir, space_id);
if let Some(parent) = space_file.parent() {
std::fs::create_dir_all(parent)?;
}
let json = serde_json::to_string_pretty(space)?;
std::fs::write(&space_file, json)?;
log::debug!("Wrote space: {}", space_file.display());
count += 1;
}
Ok(count)
}
pub async fn pull_spaces(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading spaces manifest from {}", project_dir.display());
let manifest = load_spaces_manifest(project_dir)?;
log::info!("Manifest loaded: {} space(s)", manifest.count());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
log::info!("Extracting spaces from Kibana...");
let extractor = SpacesExtractor::new(client, Some(manifest));
let spaces = extractor.extract().await?;
let mut count = 0;
for space in &spaces {
let space_id = space
.get("id")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Space missing 'id' field"))?;
let space_file = get_space_file(project_dir, space_id);
if let Some(parent) = space_file.parent() {
std::fs::create_dir_all(parent)?;
}
let json = serde_json::to_string_pretty(space)?;
std::fs::write(&space_file, json)?;
log::debug!("Wrote space: {}", space_file.display());
count += 1;
}
log::info!("✓ Pulled {} space(s)", count);
Ok(count)
}
async fn push_spaces_internal(
project_dir: impl AsRef<Path>,
client: &KibanaClient,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::debug!("Loading spaces from {}", project_dir.display());
let manifest = load_spaces_manifest(project_dir)?;
let mut spaces = Vec::new();
for entry in &manifest.spaces {
let space_file = get_space_file(project_dir, &entry.id);
if !space_file.exists() {
log::warn!(
"Space file not found for '{}': {}",
entry.id,
space_file.display()
);
continue;
}
let content = std::fs::read_to_string(&space_file)?;
let space = storage::from_json5_str(&content)?;
spaces.push(space);
}
log::debug!("Read {} space(s) from disk", spaces.len());
let loader = SpacesLoader::new(client.clone());
let count = loader.load(spaces).await?;
Ok(count)
}
pub async fn push_spaces(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading spaces from {}", project_dir.display());
let manifest = load_spaces_manifest(project_dir)?;
let mut spaces = Vec::new();
for entry in &manifest.spaces {
let space_file = get_space_file(project_dir, &entry.id);
if !space_file.exists() {
log::warn!(
"Space file not found for '{}': {}",
entry.id,
space_file.display()
);
continue;
}
let content = std::fs::read_to_string(&space_file)?;
let space = storage::from_json5_str(&content)?;
spaces.push(space);
}
log::info!("Read {} space(s) from disk", spaces.len());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
let loader = SpacesLoader::new(client);
let count = loader.load(spaces).await?;
log::info!("✓ Pushed {} space(s) to Kibana", count);
Ok(count)
}
async fn bundle_spaces_to_ndjson_internal(
project_dir: impl AsRef<Path>,
output_file: impl AsRef<Path>,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let output_file = output_file.as_ref();
log::debug!("Loading spaces from {}", project_dir.display());
let manifest = load_spaces_manifest(project_dir)?;
let mut spaces = Vec::new();
for entry in &manifest.spaces {
let space_file = get_space_file(project_dir, &entry.id);
if !space_file.exists() {
log::warn!(
"Space file not found for '{}': {}",
entry.id,
space_file.display()
);
continue;
}
let content = std::fs::read_to_string(&space_file)?;
let space = storage::from_json5_str(&content)?;
spaces.push(space);
}
log::debug!("Read {} space(s) from disk", spaces.len());
use std::io::Write;
let mut file = std::fs::File::create(output_file)?;
for space in &spaces {
let json_line = serde_json::to_string(space)?;
writeln!(file, "{}", json_line)?;
}
Ok(spaces.len())
}
pub async fn bundle_spaces_to_ndjson(
project_dir: impl AsRef<Path>,
output_file: impl AsRef<Path>,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let output_file = output_file.as_ref();
log::info!("Loading spaces from {}", project_dir.display());
let manifest = load_spaces_manifest(project_dir)?;
let mut spaces = Vec::new();
for entry in &manifest.spaces {
let space_file = get_space_file(project_dir, &entry.id);
if !space_file.exists() {
log::warn!(
"Space file not found for '{}': {}",
entry.id,
space_file.display()
);
continue;
}
let content = std::fs::read_to_string(&space_file)?;
let space = storage::from_json5_str(&content)?;
spaces.push(space);
}
log::info!("Read {} space(s) from disk", spaces.len());
use std::io::Write;
let mut file = std::fs::File::create(output_file)?;
for space in &spaces {
let json_line = serde_json::to_string(space)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} space(s) to {}",
spaces.len(),
output_file.display()
);
Ok(spaces.len())
}
fn load_workflows_manifest(project_dir: impl AsRef<Path>) -> Result<WorkflowsManifest> {
let manifest_path = project_dir.as_ref().join("manifest/workflows.yml");
if !manifest_path.exists() {
eyre::bail!("Workflows manifest not found: {}", manifest_path.display());
}
Ok(WorkflowsManifest::read(&manifest_path)?)
}
pub async fn pull_workflows(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading workflows manifest from {}", project_dir.display());
let manifest = load_workflows_manifest(project_dir)?;
log::info!("Manifest loaded: {} workflow(s)", manifest.count());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
let space_client = client.space("default")?;
log::info!("Extracting workflows from Kibana...");
let extractor = WorkflowsExtractor::new(space_client, Some(manifest));
let workflows = extractor.extract().await?;
use crate::etl::Transformer;
use crate::transform::YamlFormatter;
let formatter = YamlFormatter::for_workflows();
let formatted_workflows = workflows
.into_iter()
.map(|w| formatter.transform(w))
.collect::<std::result::Result<Vec<_>, kibana_sync::Error>>();
let formatted_workflows = formatted_workflows?;
let workflows_dir = project_dir.join("workflows");
std::fs::create_dir_all(&workflows_dir)?;
let mut count = 0;
for workflow in &formatted_workflows {
let workflow_name = workflow
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Workflow missing 'name' field"))?;
let workflow_file =
workflows_dir.join(format!("{}.json", workflow_file_stem(workflow_name)));
let json = storage::to_string_with_multiline(workflow)?;
std::fs::write(&workflow_file, json)?;
log::debug!("Wrote workflow: {}", workflow_file.display());
count += 1;
}
log::info!(
"✓ Pulled {} workflow(s) to {}",
count,
workflows_dir.display()
);
Ok(count)
}
pub async fn push_workflows(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading workflows from {}", project_dir.display());
let workflows_dir = project_dir.join("workflows");
if !workflows_dir.exists() {
eyre::bail!("Workflows directory not found: {}", workflows_dir.display());
}
let mut workflows = Vec::new();
for entry in std::fs::read_dir(&workflows_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let workflow = storage::read_json5_file(&path)?;
workflows.push(workflow);
}
}
log::info!("Read {} workflow(s) from disk", workflows.len());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
let space_client = client.space("default")?;
let loader = WorkflowsLoader::new(space_client);
let count = loader.load(workflows).await?;
log::info!("✓ Pushed {} workflow(s) to Kibana", count);
Ok(count)
}
pub async fn bundle_workflows_to_ndjson(
project_dir: impl AsRef<Path>,
output_file: impl AsRef<Path>,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let output_file = output_file.as_ref();
log::info!("Loading workflows from {}", project_dir.display());
let workflows_dir = project_dir.join("workflows");
if !workflows_dir.exists() {
eyre::bail!("Workflows directory not found: {}", workflows_dir.display());
}
let mut workflows = Vec::new();
for entry in std::fs::read_dir(&workflows_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let workflow = storage::read_json5_file(&path)?;
workflows.push(workflow);
}
}
log::info!("Read {} workflow(s) from disk", workflows.len());
use std::io::Write;
let mut file = std::fs::File::create(output_file)?;
for workflow in &workflows {
let json_line = serde_json::to_string(workflow)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} workflow(s) to {}",
workflows.len(),
output_file.display()
);
Ok(workflows.len())
}
fn load_agents_manifest(project_dir: impl AsRef<Path>) -> Result<AgentsManifest> {
let project_dir = project_dir.as_ref();
let manifest_path = project_dir.join("manifest/agents.yml");
if !manifest_path.exists() {
eyre::bail!(
"Agents manifest not found: {}. Run 'kibob add agents' to create it.",
manifest_path.display()
);
}
Ok(AgentsManifest::read(&manifest_path)?)
}
#[allow(clippy::too_many_arguments)]
pub async fn add_agents_to_manifest(
project_dir: impl AsRef<Path>,
space_id: &str,
query: Option<String>,
include: Option<String>,
exclude: Option<String>,
file_path: Option<String>,
exclude_dependencies: bool,
force: bool,
) -> Result<usize> {
use crate::kibana::agents::AgentEntry;
let project_dir = project_dir.as_ref();
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
let manifest = SpacesManifest::read(&spaces_manifest_path)?;
if !manifest.spaces.iter().any(|s| s.id == space_id) {
eyre::bail!(
"Space {} is not managed. Add it first with: kibob add spaces . --include '^{}$'",
space_id.cyan(),
space_id
);
}
}
log::info!(
"Loading agents manifest for space {} from {}",
space_id.cyan(),
project_dir.display()
);
let manifest_path = get_space_agents_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
AgentsManifest::read(&manifest_path)?
} else {
log::info!("No existing manifest found, will create new one");
AgentsManifest::new()
};
log::info!("Current manifest has {} agent(s)", manifest.count());
let requires_kibana_calls = file_path.is_none() || !exclude_dependencies;
let mut preflight_client: Option<KibanaClient> = None;
let mut detected_version: Option<KibanaVersion> = None;
if requires_kibana_calls {
let client = load_kibana_client(project_dir)?;
let preflight = run_version_preflight(project_dir, &client, "add", force).await?;
detected_version = Some(preflight.detected.parsed);
preflight_client = Some(client);
}
if let Some(version) = detected_version
&& !KibanaClient::supports_capability(&version, ApiCapability::Agents)
{
let reason = KibanaClient::unsupported_capability_reason(&version, ApiCapability::Agents);
if force {
log::warn!("{}", reason);
log::warn!("--force enabled: attempting agents API calls anyway");
} else if file_path.is_none() {
return Err(version_warning(reason));
} else {
log::warn!("{}", reason);
}
}
let new_agents: Vec<serde_json::Value> = if let Some(file) = file_path {
log::info!("Reading agents from {}", file);
let file_path = std::path::Path::new(&file);
if !file_path.exists() {
eyre::bail!("File not found: {}", file_path.display());
}
let extension = file_path.extension().and_then(|s| s.to_str()).unwrap_or("");
match extension {
"ndjson" => {
use std::io::{BufRead, BufReader};
let file = std::fs::File::open(file_path)?;
let reader = BufReader::new(file);
let mut agents = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let agent: serde_json::Value = serde_json::from_str(&line)?;
agents.push(agent);
}
log::info!("Read {} agent(s) from NDJSON file", agents.len());
agents
}
"json" => {
let content = std::fs::read_to_string(file_path)?;
let parsed = storage::from_json5_str(&content)?;
if let Some(arr) = parsed.as_array() {
log::info!("Read {} agent(s) from JSON array", arr.len());
arr.to_vec()
} else {
log::info!("Read 1 agent from JSON file");
vec![parsed]
}
}
_ => {
eyre::bail!(
"Unsupported file format: {}. Expected .json or .ndjson",
extension
);
}
}
} else {
log::info!("Searching agents via API in space {}...", space_id.cyan());
let client = preflight_client
.clone()
.unwrap_or(load_kibana_client(project_dir)?);
let space_client = client.space(space_id)?;
let extractor = AgentsExtractor::new(space_client, None);
extractor.search_agents(query.as_deref()).await?
};
log::info!("Found {} agent(s) before filtering", new_agents.len());
let filtered_agents: Vec<serde_json::Value> = {
let mut agents = new_agents;
if let Some(include_pattern) = &include {
let regex = regex::Regex::new(include_pattern)
.with_context(|| format!("Invalid include regex pattern: {}", include_pattern))?;
agents.retain(|a| {
a.get("name")
.and_then(|v| v.as_str())
.map(|name| regex.is_match(name))
.unwrap_or(false)
});
log::info!(
"After include filter '{}': {} agent(s)",
include_pattern,
agents.len()
);
}
if let Some(exclude_pattern) = &exclude {
let regex = regex::Regex::new(exclude_pattern)
.with_context(|| format!("Invalid exclude regex pattern: {}", exclude_pattern))?;
agents.retain(|a| {
a.get("name")
.and_then(|v| v.as_str())
.map(|name| !regex.is_match(name))
.unwrap_or(true)
});
log::info!(
"After exclude filter '{}': {} agent(s)",
exclude_pattern,
agents.len()
);
}
agents
};
log::info!("Adding {} agent(s) after filtering", filtered_agents.len());
let agents_dir = get_space_agents_dir(project_dir, space_id);
std::fs::create_dir_all(&agents_dir)?;
let mut added_count = 0;
let mut all_deps = Vec::new();
for agent in &filtered_agents {
let agent_id = agent
.get("id")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Agent missing 'id' field"))?;
let agent_name = agent
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Agent missing 'name' field"))?;
if manifest.add_agent(AgentEntry::new(agent_id, agent_name)) {
log::debug!("Added agent to manifest: {} ({})", agent_name, agent_id);
let agent_file = agents_dir.join(format!("{}.json", agent_name));
let json = serde_json::to_string_pretty(agent)?;
std::fs::write(&agent_file, json)?;
log::debug!("Wrote agent file: {}", agent_file.display());
added_count += 1;
if !exclude_dependencies {
all_deps.extend(find_agent_dependencies(agent));
}
} else {
log::debug!("Agent already in manifest, skipping: {}", agent_name);
}
}
let mut dep_summary = DependencySummary::new();
if !all_deps.is_empty() {
let client = preflight_client.unwrap_or(load_kibana_client(project_dir)?);
let version = client.server_version().await?;
dep_summary =
resolve_and_add_dependencies(project_dir, space_id, &client, all_deps, version, force)
.await?;
}
let manifest_dir = get_space_manifest_dir(project_dir, space_id);
std::fs::create_dir_all(&manifest_dir)?;
manifest.write(&manifest_path)?;
log::info!(
"✓ Updated manifest for space {} now has {} agent(s)",
space_id.cyan(),
manifest.count()
);
if !dep_summary.is_empty() {
log::info!(
"✓ Added {} new agent(s) (plus dependencies: {})",
added_count,
dep_summary.format_summary()
);
} else {
log::info!("✓ Added {} new agent(s)", added_count);
}
Ok(added_count + dep_summary.total())
}
pub async fn pull_agents(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading agents manifest from {}", project_dir.display());
let manifest = load_agents_manifest(project_dir)?;
log::info!("Manifest loaded: {} agent(s)", manifest.count());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
let space_client = client.space("default")?;
log::info!("Extracting agents from Kibana...");
let extractor = AgentsExtractor::new(space_client, Some(manifest));
let agents = extractor.extract().await?;
let formatter = MultilineFieldFormatter::for_agents();
let agents: Vec<_> = agents
.into_iter()
.map(|agent| formatter.transform(agent))
.collect::<std::result::Result<_, kibana_sync::Error>>()?;
let agents_dir = project_dir.join("agents");
std::fs::create_dir_all(&agents_dir)?;
let mut count = 0;
for agent in &agents {
let agent_name = agent
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Agent missing 'name' field"))?;
let agent_file = agents_dir.join(format!("{}.json", agent_name));
let json = storage::to_string_with_multiline(agent)?;
std::fs::write(&agent_file, json)?;
log::debug!("Wrote agent: {}", agent_file.display());
count += 1;
}
log::info!("✓ Pulled {} agent(s) to {}", count, agents_dir.display());
Ok(count)
}
pub async fn push_agents(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading agents from {}", project_dir.display());
let agents_dir = project_dir.join("agents");
if !agents_dir.exists() {
eyre::bail!("Agents directory not found: {}", agents_dir.display());
}
let mut agents = Vec::new();
for entry in std::fs::read_dir(&agents_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let agent = storage::read_json5_file(&path)?;
agents.push(agent);
}
}
log::info!("Read {} agent(s) from disk", agents.len());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
let space_client = client.space("default")?;
let loader = AgentsLoader::new(space_client);
let count = loader.load(agents).await?;
log::info!("✓ Pushed {} agent(s) to Kibana", count);
Ok(count)
}
pub async fn bundle_agents_to_ndjson(
project_dir: impl AsRef<Path>,
output_file: impl AsRef<Path>,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let output_file = output_file.as_ref();
log::info!("Loading agents from {}", project_dir.display());
let agents_dir = project_dir.join("agents");
if !agents_dir.exists() {
eyre::bail!("Agents directory not found: {}", agents_dir.display());
}
let mut agents = Vec::new();
for entry in std::fs::read_dir(&agents_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let agent = storage::read_json5_file(&path)?;
agents.push(agent);
}
}
log::info!("Read {} agent(s) from disk", agents.len());
use std::io::Write;
let mut file = std::fs::File::create(output_file)?;
for agent in &agents {
let json_line = serde_json::to_string(agent)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} agent(s) to {}",
agents.len(),
output_file.display()
);
Ok(agents.len())
}
fn load_tools_manifest(project_dir: impl AsRef<Path>) -> Result<ToolsManifest> {
let project_dir = project_dir.as_ref();
let manifest_path = project_dir.join("manifest/tools.yml");
if !manifest_path.exists() {
eyre::bail!(
"Tools manifest not found: {}. Run 'kibob add tools' to create it.",
manifest_path.display()
);
}
Ok(ToolsManifest::read(&manifest_path)?)
}
#[allow(clippy::too_many_arguments)]
pub async fn add_tools_to_manifest(
project_dir: impl AsRef<Path>,
space_id: &str,
query: Option<String>,
include: Option<String>,
exclude: Option<String>,
file_path: Option<String>,
exclude_dependencies: bool,
force: bool,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
let client = load_kibana_client(project_dir)?;
if client.space(space_id).is_err() {
eyre::bail!(
"Space {} is not managed. Add it first with: kibob add spaces . --include '^{}$'",
space_id.cyan(),
space_id
);
}
}
log::info!(
"Loading tools manifest for space {} from {}",
space_id.cyan(),
project_dir.display()
);
let manifest_path = get_space_tools_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
ToolsManifest::read(&manifest_path)?
} else {
log::info!("No existing manifest found, will create new one");
ToolsManifest::new()
};
log::info!("Current manifest has {} tool(s)", manifest.count());
let requires_kibana_calls = file_path.is_none() || !exclude_dependencies;
let mut preflight_client: Option<KibanaClient> = None;
let mut detected_version: Option<KibanaVersion> = None;
if requires_kibana_calls {
let client = load_kibana_client(project_dir)?;
let preflight = run_version_preflight(project_dir, &client, "add", force).await?;
detected_version = Some(preflight.detected.parsed);
preflight_client = Some(client);
}
if let Some(version) = detected_version
&& !KibanaClient::supports_capability(&version, ApiCapability::Tools)
{
let reason = KibanaClient::unsupported_capability_reason(&version, ApiCapability::Tools);
if force {
log::warn!("{}", reason);
log::warn!("--force enabled: attempting tools API calls anyway");
} else if file_path.is_none() {
return Err(version_warning(reason));
} else {
log::warn!("{}", reason);
}
}
let new_tools: Vec<serde_json::Value> = if let Some(file) = file_path {
log::info!("Reading tools from {}", file);
let file_path = std::path::Path::new(&file);
if !file_path.exists() {
eyre::bail!("File not found: {}", file_path.display());
}
let extension = file_path.extension().and_then(|s| s.to_str()).unwrap_or("");
match extension {
"ndjson" => {
use std::io::{BufRead, BufReader};
let file = std::fs::File::open(file_path)?;
let reader = BufReader::new(file);
let mut tools = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let tool: serde_json::Value = serde_json::from_str(&line)?;
tools.push(tool);
}
log::info!("Read {} tool(s) from NDJSON file", tools.len());
tools
}
"json" => {
let content = std::fs::read_to_string(file_path)?;
let parsed = storage::from_json5_str(&content)?;
if let Some(arr) = parsed.as_array() {
log::info!("Read {} tool(s) from JSON array", arr.len());
arr.to_vec()
} else {
log::info!("Read 1 tool from JSON file");
vec![parsed]
}
}
_ => {
eyre::bail!(
"Unsupported file format: {}. Expected .json or .ndjson",
extension
);
}
}
} else {
log::info!("Searching tools via API in space {}...", space_id.cyan());
let client = preflight_client
.clone()
.unwrap_or(load_kibana_client(project_dir)?);
let space_client = client.space(space_id)?;
let extractor = ToolsExtractor::new(space_client, None);
extractor.search_tools(query.as_deref()).await?
};
log::info!("Found {} tool(s) before filtering", new_tools.len());
let filtered_tools: Vec<serde_json::Value> = {
let mut tools = new_tools;
if let Some(include_pattern) = &include {
let regex = regex::Regex::new(include_pattern)
.with_context(|| format!("Invalid include regex pattern: {}", include_pattern))?;
tools.retain(|t| {
let filter_field = t
.get("name")
.or_else(|| t.get("id"))
.and_then(|v| v.as_str())
.unwrap_or("");
regex.is_match(filter_field)
});
log::info!(
"After include filter '{}': {} tool(s)",
include_pattern,
tools.len()
);
}
if let Some(exclude_pattern) = &exclude {
let regex = regex::Regex::new(exclude_pattern)
.with_context(|| format!("Invalid exclude regex pattern: {}", exclude_pattern))?;
tools.retain(|t| {
let filter_field = t
.get("name")
.or_else(|| t.get("id"))
.and_then(|v| v.as_str())
.unwrap_or("");
!regex.is_match(filter_field)
});
log::info!(
"After exclude filter '{}': {} tool(s)",
exclude_pattern,
tools.len()
);
}
tools
};
log::info!("Adding {} tool(s) after filtering", filtered_tools.len());
let tools_dir = get_space_tools_dir(project_dir, space_id);
std::fs::create_dir_all(&tools_dir)?;
let mut added_count = 0;
let mut all_deps = Vec::new();
for tool in &filtered_tools {
let tool_id = tool
.get("id")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Tool missing 'id' field"))?;
let filename = tool.get("name").and_then(|v| v.as_str()).unwrap_or(tool_id);
if manifest.add_tool(tool_id.to_string()) {
log::debug!("Added tool to manifest: {}", tool_id);
let tool_file = tools_dir.join(format!("{}.json", filename));
let json = serde_json::to_string_pretty(tool)?;
std::fs::write(&tool_file, json)?;
log::debug!("Wrote tool file: {}", tool_file.display());
added_count += 1;
if !exclude_dependencies {
all_deps.extend(find_tool_dependencies(tool));
}
} else {
log::debug!("Tool already in manifest, skipping: {}", tool_id);
}
}
let mut dep_summary = DependencySummary::new();
if !all_deps.is_empty() {
let client = preflight_client.unwrap_or(load_kibana_client(project_dir)?);
let version = client.server_version().await?;
dep_summary =
resolve_and_add_dependencies(project_dir, space_id, &client, all_deps, version, force)
.await?;
}
let manifest_dir = get_space_manifest_dir(project_dir, space_id);
std::fs::create_dir_all(&manifest_dir)?;
manifest.write(&manifest_path)?;
log::info!(
"✓ Updated manifest for space {} now has {} tool(s)",
space_id.cyan(),
manifest.count()
);
if !dep_summary.is_empty() {
log::info!(
"✓ Added {} new tool(s) (plus dependencies: {})",
added_count,
dep_summary.format_summary()
);
} else {
log::info!("✓ Added {} new tool(s)", added_count);
}
Ok(added_count + dep_summary.total())
}
#[allow(clippy::too_many_arguments)]
pub async fn add_skills_to_manifest(
project_dir: impl AsRef<Path>,
space_id: &str,
query: Option<String>,
include: Option<String>,
exclude: Option<String>,
file_path: Option<String>,
exclude_dependencies: bool,
force: bool,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let spaces_manifest_path = project_dir.join("spaces.yml");
if spaces_manifest_path.exists() {
let manifest = SpacesManifest::read(&spaces_manifest_path)?;
if !manifest.spaces.iter().any(|s| s.id == space_id) {
eyre::bail!(
"Space {} is not managed. Add it first with: kibob add spaces . --include '^{}$'",
space_id.cyan(),
space_id
);
}
}
let requires_kibana_calls = file_path.is_none() || !exclude_dependencies;
let mut preflight_client: Option<KibanaClient> = None;
let mut detected_version: Option<KibanaVersion> = None;
if requires_kibana_calls {
let client = load_kibana_client(project_dir)?;
let preflight = run_version_preflight(project_dir, &client, "add", force).await?;
detected_version = Some(preflight.detected.parsed);
preflight_client = Some(client);
}
if let Some(version) = detected_version
&& !KibanaClient::supports_capability(&version, ApiCapability::Skills)
{
let reason = KibanaClient::unsupported_capability_reason(&version, ApiCapability::Skills);
if force {
log::warn!("{}", reason);
log::warn!("--force enabled: attempting skills API calls anyway");
} else if file_path.is_none() {
return Err(version_warning(reason));
} else {
log::warn!("{}", reason);
}
}
let filters_applied_before_fetch = file_path.is_none() && query.is_none();
let new_skills: Vec<Value> = if let Some(file) = file_path {
read_agent_builder_values_from_file(&file, "skill")?
} else {
let client = preflight_client
.clone()
.unwrap_or(load_kibana_client(project_dir)?);
let space_client = client.space(space_id)?;
let extractor = SkillsExtractor::new(space_client.clone(), None);
if let Some(skill_id) = query.as_deref() {
log::info!(
"Fetching skill {} in space {}...",
skill_id,
space_id.cyan()
);
vec![extractor.fetch_skill(skill_id).await?]
} else {
log::info!("Searching skills via API in space {}...", space_id.cyan());
let listed = extractor.search_skills(false).await?;
log::info!("Found {} skill(s) before filtering", listed.len());
let listed = filter_skills(listed, &include, &exclude)?;
let mut skill_ids = Vec::new();
for (index, skill) in listed.into_iter().enumerate() {
if is_readonly(&skill) {
continue;
}
let Some(skill_id) = skill.get("id").and_then(|v| v.as_str()) else {
log::warn!("Skipping skill list entry without id");
continue;
};
skill_ids.push((index, skill_id.to_string()));
}
fetch_skills_in_order(space_client, skill_ids).await?
}
};
let filtered_skills = if filters_applied_before_fetch {
new_skills
} else {
log::info!("Found {} skill(s) before filtering", new_skills.len());
filter_skills(new_skills, &include, &exclude)?
};
log::info!("Adding {} skill(s) after filtering", filtered_skills.len());
let skills_dir = get_space_skills_dir(project_dir, space_id);
std::fs::create_dir_all(&skills_dir)?;
let manifest_path = get_space_skills_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
SkillsManifest::read(&manifest_path)?
} else {
SkillsManifest::new()
};
let mut added_count = 0;
let mut all_deps = Vec::new();
for skill in &filtered_skills {
if is_readonly(skill) {
log::debug!(
"Skipping readonly skill {}",
skill
.get("id")
.and_then(|v| v.as_str())
.unwrap_or("<unknown>")
);
continue;
}
skill_to_directory(&skills_dir, skill)?;
let added = manifest.add_skill(skill_entry(skill)?);
if added {
added_count += 1;
if !exclude_dependencies {
all_deps.extend(find_skill_dependencies(skill));
}
}
}
manifest.write(&manifest_path)?;
let mut dep_summary = DependencySummary::new();
if !all_deps.is_empty() {
let client = preflight_client.unwrap_or(load_kibana_client(project_dir)?);
let version = client.server_version().await?;
dep_summary =
resolve_and_add_dependencies(project_dir, space_id, &client, all_deps, version, force)
.await?;
}
if !dep_summary.is_empty() {
log::info!(
"✓ Added {} skill(s) (plus dependencies: {})",
added_count,
dep_summary.format_summary()
);
} else {
log::info!("✓ Added {} skill(s)", added_count);
}
Ok(added_count + dep_summary.total())
}
fn read_agent_builder_values_from_file(file: &str, label: &str) -> Result<Vec<Value>> {
log::info!("Reading {}s from {}", label, file);
let file_path = std::path::Path::new(file);
if !file_path.exists() {
eyre::bail!("File not found: {}", file_path.display());
}
let extension = file_path.extension().and_then(|s| s.to_str()).unwrap_or("");
match extension {
"ndjson" => {
use std::io::{BufRead, BufReader};
let file = std::fs::File::open(file_path)?;
let reader = BufReader::new(file);
let mut values = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
values.push(serde_json::from_str(&line)?);
}
log::info!("Read {} {}(s) from NDJSON file", values.len(), label);
Ok(values)
}
"json" => {
let content = std::fs::read_to_string(file_path)?;
let parsed = storage::from_json5_str(&content)?;
if let Some(results) = parsed.get("results").and_then(|v| v.as_array()) {
log::info!("Read {} {}(s) from JSON API response", results.len(), label);
Ok(results.to_vec())
} else if let Some(arr) = parsed.as_array() {
log::info!("Read {} {}(s) from JSON array", arr.len(), label);
Ok(arr.to_vec())
} else {
log::info!("Read 1 {} from JSON file", label);
Ok(vec![parsed])
}
}
_ => {
eyre::bail!(
"Unsupported file format: {}. Expected .json or .ndjson",
extension
);
}
}
}
fn is_readonly(value: &Value) -> bool {
value
.get("readonly")
.and_then(|value| value.as_bool())
.unwrap_or(false)
}
fn skill_filter_field(skill: &Value) -> &str {
skill
.get("name")
.or_else(|| skill.get("id"))
.and_then(|v| v.as_str())
.unwrap_or("")
}
async fn fetch_skills_in_order(
client: KibanaClient,
skill_ids: Vec<(usize, String)>,
) -> Result<Vec<Value>> {
let mut fetched_skills = Vec::new();
for chunk in skill_ids.chunks(SKILL_FETCH_BATCH_SIZE) {
let mut set = JoinSet::new();
for (index, skill_id) in chunk {
let extractor = SkillsExtractor::new(client.clone(), None);
let index = *index;
let skill_id = skill_id.clone();
set.spawn(async move {
match extractor.fetch_skill(&skill_id).await {
Ok(skill) => Ok((index, skill_id, skill)),
Err(err) => Err((skill_id, err.to_string())),
}
});
}
while let Some(result) = set.join_next().await {
match result {
Ok(Ok((index, _skill_id, skill))) if !is_readonly(&skill) => {
fetched_skills.push((index, skill));
}
Ok(Ok((_index, skill_id, _skill))) => {
log::warn!("Skipping readonly skill {}", skill_id);
}
Ok(Err((skill_id, err))) => {
eyre::bail!("Failed to fetch skill {}: {}", skill_id, err)
}
Err(err) => eyre::bail!("Task panicked: {}", err),
}
}
}
fetched_skills.sort_by_key(|(index, _)| *index);
Ok(fetched_skills
.into_iter()
.map(|(_index, skill)| skill)
.collect())
}
fn filter_skills(
mut skills: Vec<Value>,
include: &Option<String>,
exclude: &Option<String>,
) -> Result<Vec<Value>> {
if let Some(include_pattern) = include {
let regex = regex::Regex::new(include_pattern)
.with_context(|| format!("Invalid include regex pattern: {}", include_pattern))?;
skills.retain(|skill| regex.is_match(skill_filter_field(skill)));
log::info!(
"After include filter '{}': {} skill(s)",
include_pattern,
skills.len()
);
}
if let Some(exclude_pattern) = exclude {
let regex = regex::Regex::new(exclude_pattern)
.with_context(|| format!("Invalid exclude regex pattern: {}", exclude_pattern))?;
skills.retain(|skill| !regex.is_match(skill_filter_field(skill)));
log::info!(
"After exclude filter '{}': {} skill(s)",
exclude_pattern,
skills.len()
);
}
Ok(skills)
}
pub async fn pull_tools(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading tools manifest from {}", project_dir.display());
let manifest = load_tools_manifest(project_dir)?;
log::info!("Manifest loaded: {} tool(s)", manifest.count());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
let space_client = client.space("default")?;
log::info!("Extracting tools from Kibana...");
let extractor = ToolsExtractor::new(space_client, Some(manifest));
let tools = extractor.extract().await?;
let formatter = MultilineFieldFormatter::for_tools();
let tools: Vec<_> = tools
.into_iter()
.map(|tool| formatter.transform(tool))
.collect::<std::result::Result<_, kibana_sync::Error>>()?;
let tools_dir = project_dir.join("tools");
std::fs::create_dir_all(&tools_dir)?;
let mut count = 0;
for tool in &tools {
let tool_id = tool
.get("id")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Tool missing 'id' field"))?;
let filename = tool.get("name").and_then(|v| v.as_str()).unwrap_or(tool_id);
let tool_file = tools_dir.join(format!("{}.json", filename));
let json = storage::to_string_with_multiline(tool)?;
std::fs::write(&tool_file, json)?;
log::debug!("Wrote tool: {}", tool_file.display());
count += 1;
}
log::info!("✓ Pulled {} tool(s) to {}", count, tools_dir.display());
Ok(count)
}
pub async fn push_tools(project_dir: impl AsRef<Path>) -> Result<usize> {
let project_dir = project_dir.as_ref();
log::info!("Loading tools from {}", project_dir.display());
let tools_dir = project_dir.join("tools");
if !tools_dir.exists() {
eyre::bail!("Tools directory not found: {}", tools_dir.display());
}
let mut tools = Vec::new();
for entry in std::fs::read_dir(&tools_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let tool = storage::read_json5_file(&path)?;
tools.push(tool);
}
}
log::info!("Read {} tool(s) from disk", tools.len());
log::info!("Connecting to Kibana...");
let client = load_kibana_client(project_dir)?;
let space_client = client.space("default")?;
let loader = ToolsLoader::new(space_client);
let count = loader.load(tools).await?;
log::info!("✓ Pushed {} tool(s) to Kibana", count);
Ok(count)
}
pub async fn bundle_tools_to_ndjson(
project_dir: impl AsRef<Path>,
output_file: impl AsRef<Path>,
) -> Result<usize> {
let project_dir = project_dir.as_ref();
let output_file = output_file.as_ref();
log::info!("Loading tools from {}", project_dir.display());
let tools_dir = project_dir.join("tools");
if !tools_dir.exists() {
eyre::bail!("Tools directory not found: {}", tools_dir.display());
}
let mut tools = Vec::new();
for entry in std::fs::read_dir(&tools_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let tool = storage::read_json5_file(&path)?;
tools.push(tool);
}
}
log::info!("Read {} tool(s) from disk", tools.len());
use std::io::Write;
let mut file = std::fs::File::create(output_file)?;
for tool in &tools {
let json_line = serde_json::to_string(tool)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} tool(s) to {}",
tools.len(),
output_file.display()
);
Ok(tools.len())
}
async fn pull_space_saved_objects(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let manifest_path = get_space_saved_objects_manifest(project_dir, space_id);
if !manifest_path.exists() {
log::debug!(
"No saved objects manifest for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Pulling saved objects for space {}", space_id.cyan());
let manifest = crate::kibana::saved_objects::SavedObjectsManifest::read(&manifest_path)?;
log::debug!("Loaded {} object(s) from manifest", manifest.count());
let extractor = SavedObjectsExtractor::new(client.clone(), manifest);
let drop_fields = FieldDropper::default_kibana_fields();
let unescape = FieldUnescaper::default_kibana_fields();
let vega_unescape = VegaSpecUnescaper::new();
let objects_dir = get_space_objects_dir(project_dir, space_id);
let writer = DirectoryWriter::new_with_options(&objects_dir, true)?;
writer.clear()?;
let objects = extractor.extract().await?;
let dropped = drop_fields.transform_many(objects)?;
let unescaped = unescape.transform_many(dropped)?;
let vega_unescaped = vega_unescape.transform_many(unescaped)?;
let count = writer.load(vega_unescaped).await?;
log::info!(
"✓ Pulled {} saved object(s) for space {}",
count,
space_id.cyan()
);
Ok(count)
}
async fn pull_space_workflows(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let manifest_path = get_space_workflows_manifest(project_dir, space_id);
if !manifest_path.exists() {
log::debug!(
"No workflows manifest for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Pulling workflows for space {}", space_id.cyan());
let manifest = WorkflowsManifest::read(&manifest_path)?;
log::debug!("Loaded {} workflow(s) from manifest", manifest.count());
let extractor = WorkflowsExtractor::new(client.clone(), Some(manifest));
let workflows = extractor.extract().await?;
use crate::etl::Transformer;
use crate::transform::YamlFormatter;
let formatter = YamlFormatter::for_workflows();
let formatted_workflows = workflows
.into_iter()
.map(|w| formatter.transform(w))
.collect::<std::result::Result<Vec<_>, kibana_sync::Error>>();
let formatted_workflows = formatted_workflows?;
let workflows_dir = get_space_workflows_dir(project_dir, space_id);
std::fs::create_dir_all(&workflows_dir)?;
let mut count = 0;
for workflow in &formatted_workflows {
let workflow_name = workflow
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| eyre::eyre!("Workflow missing 'name' field"))?;
let workflow_file =
workflows_dir.join(format!("{}.json", workflow_file_stem(workflow_name)));
let json = storage::to_string_with_multiline(workflow)?;
std::fs::write(&workflow_file, json)?;
count += 1;
}
log::info!(
"✓ Pulled {} workflow(s) for space {}",
count,
space_id.cyan()
);
Ok(count)
}
async fn pull_space_agents(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let manifest_path = get_space_agents_manifest(project_dir, space_id);
if !manifest_path.exists() {
log::debug!("No agents manifest for space {}, skipping", space_id.cyan());
return Ok(0);
}
log::info!("Pulling agents for space {}", space_id.cyan());
let manifest = AgentsManifest::read(&manifest_path)?;
log::debug!("Loaded {} agent(s) from manifest", manifest.count());
let extractor = AgentsExtractor::new(client.clone(), Some(manifest));
let agents = extractor.extract().await?;
let formatter = MultilineFieldFormatter::for_agents();
let agents: Vec<_> = agents
.into_iter()
.map(|agent| formatter.transform(agent))
.collect::<std::result::Result<_, kibana_sync::Error>>()?;
let agents_dir = get_space_agents_dir(project_dir, space_id);
std::fs::create_dir_all(&agents_dir)?;
let mut count = 0;
for agent in &agents {
let agent_name = agent
.get("name")
.and_then(|v| v.as_str())
.or_else(|| agent.get("id").and_then(|v| v.as_str()))
.ok_or_else(|| eyre::eyre!("Agent missing both 'name' and 'id' fields"))?;
let agent_file = agents_dir.join(format!("{}.json", agent_name));
let json = storage::to_string_with_multiline(agent)?;
std::fs::write(&agent_file, json)?;
count += 1;
}
log::info!("✓ Pulled {} agent(s) for space {}", count, space_id.cyan());
Ok(count)
}
async fn pull_space_tools(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let manifest_path = get_space_tools_manifest(project_dir, space_id);
if !manifest_path.exists() {
log::debug!("No tools manifest for space {}, skipping", space_id.cyan());
return Ok(0);
}
log::info!("Pulling tools for space {}", space_id.cyan());
let manifest = ToolsManifest::read(&manifest_path)?;
log::debug!("Loaded {} tool(s) from manifest", manifest.count());
let extractor = ToolsExtractor::new(client.clone(), Some(manifest));
let tools = extractor.extract().await?;
let formatter = MultilineFieldFormatter::for_tools();
let tools: Vec<_> = tools
.into_iter()
.map(|tool| formatter.transform(tool))
.collect::<std::result::Result<_, kibana_sync::Error>>()?;
let tools_dir = get_space_tools_dir(project_dir, space_id);
std::fs::create_dir_all(&tools_dir)?;
let mut count = 0;
for tool in &tools {
let tool_name = tool
.get("name")
.and_then(|v| v.as_str())
.or_else(|| tool.get("id").and_then(|v| v.as_str()))
.ok_or_else(|| eyre::eyre!("Tool missing both 'name' and 'id' fields"))?;
let tool_file = tools_dir.join(format!("{}.json", tool_name));
let json = storage::to_string_with_multiline(tool)?;
std::fs::write(&tool_file, json)?;
count += 1;
}
log::info!("✓ Pulled {} tool(s) for space {}", count, space_id.cyan());
Ok(count)
}
async fn pull_space_skills(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let manifest_path = get_space_skills_manifest(project_dir, space_id);
if !manifest_path.exists() {
log::debug!("No skills manifest for space {}, skipping", space_id.cyan());
return Ok(0);
}
log::info!("Pulling skills for space {}", space_id.cyan());
let manifest = SkillsManifest::read(&manifest_path)?;
log::debug!("Loaded {} skill(s) from manifest", manifest.count());
let skills_dir = get_space_skills_dir(project_dir, space_id);
std::fs::create_dir_all(&skills_dir)?;
let skill_ids = manifest
.skills
.iter()
.enumerate()
.map(|(index, entry)| (index, entry.id.clone()))
.collect::<Vec<_>>();
let skills = fetch_skills_in_order(client.clone(), skill_ids).await?;
let mut count = 0;
for full_skill in skills {
skill_to_directory(&skills_dir, &full_skill)?;
count += 1;
}
log::info!("✓ Pulled {} skill(s) for space {}", count, space_id.cyan());
Ok(count)
}
async fn push_space_saved_objects(
project_dir: &Path,
client: &KibanaClient,
managed: bool,
) -> Result<usize> {
let space_id = client.space_id();
let objects_dir = get_space_objects_dir(project_dir, space_id);
if !objects_dir.exists() {
log::debug!(
"No objects directory for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Pushing saved objects for space {}", space_id.cyan());
let reader = DirectoryReader::new(&objects_dir);
let vega_escaper = VegaSpecEscaper::new();
let escaper = FieldEscaper::default_kibana_fields();
let managed_flag = ManagedFlagAdder::new(managed);
let loader = SavedObjectsLoader::new(client.clone());
let objects = reader.extract().await?;
let vega_escaped = vega_escaper.transform_many(objects)?;
let escaped = escaper.transform_many(vega_escaped)?;
let flagged = managed_flag.transform_many(escaped)?;
let count = loader.load(flagged).await?;
log::info!(
"✓ Pushed {} saved object(s) for space {}",
count,
space_id.cyan()
);
Ok(count)
}
async fn push_space_workflows(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let workflows_dir = get_space_workflows_dir(project_dir, space_id);
if !workflows_dir.exists() {
log::debug!(
"No workflows directory for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Pushing workflows for space {}", space_id.cyan());
let mut workflows = Vec::new();
for entry in std::fs::read_dir(&workflows_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let workflow = storage::read_json5_file(&path)?;
workflows.push(workflow);
}
}
let loader = WorkflowsLoader::new(client.clone());
let count = loader.load(workflows).await?;
log::info!(
"✓ Pushed {} workflow(s) for space {}",
count,
space_id.cyan()
);
Ok(count)
}
async fn push_space_agents(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let agents_dir = get_space_agents_dir(project_dir, space_id);
if !agents_dir.exists() {
log::debug!(
"No agents directory for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Pushing agents for space {}", space_id.cyan());
let mut agents = Vec::new();
for entry in std::fs::read_dir(&agents_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let agent = storage::read_json5_file(&path)?;
agents.push(agent);
}
}
let loader = AgentsLoader::new(client.clone());
let count = loader.load(agents).await?;
log::info!("✓ Pushed {} agent(s) for space {}", count, space_id.cyan());
Ok(count)
}
async fn push_space_tools(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let tools_dir = get_space_tools_dir(project_dir, space_id);
if !tools_dir.exists() {
log::debug!("No tools directory for space {}, skipping", space_id.cyan());
return Ok(0);
}
log::info!("Pushing tools for space {}", space_id.cyan());
let mut tools = Vec::new();
for entry in std::fs::read_dir(&tools_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let tool = storage::read_json5_file(&path)?;
tools.push(tool);
}
}
let loader = ToolsLoader::new(client.clone());
let count = loader.load(tools).await?;
log::info!("✓ Pushed {} tool(s) for space {}", count, space_id.cyan());
Ok(count)
}
async fn push_space_skills(project_dir: &Path, client: &KibanaClient) -> Result<usize> {
let space_id = client.space_id();
let skills_dir = get_space_skills_dir(project_dir, space_id);
let skills = read_skill_values_from_dir(
&skills_dir,
&get_space_skills_manifest(project_dir, space_id),
)?;
if skills.is_empty() {
if skills_dir.exists() {
log::debug!("No skills selected for space {}, skipping", space_id.cyan());
} else {
log::debug!(
"No skills directory for space {}, skipping",
space_id.cyan()
);
}
return Ok(0);
}
log::info!("Pushing skills for space {}", space_id.cyan());
let loader = SkillsLoader::new(client.clone());
let count = loader.load(skills).await?;
log::info!("✓ Pushed {} skill(s) for space {}", count, space_id.cyan());
Ok(count)
}
async fn bundle_space_saved_objects(
project_dir: &Path,
space_id: &str,
managed: bool,
) -> Result<usize> {
let objects_dir = get_space_objects_dir(project_dir, space_id);
if !objects_dir.exists() {
log::debug!(
"No objects directory for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Bundling saved objects for space {}", space_id.cyan());
let reader = DirectoryReader::new(&objects_dir);
let vega_escaper = VegaSpecEscaper::new();
let escaper = FieldEscaper::default_kibana_fields();
let managed_flag = ManagedFlagAdder::new(managed);
let objects = reader.extract().await?;
let vega_escaped = vega_escaper.transform_many(objects)?;
let escaped = escaper.transform_many(vega_escaped)?;
let flagged = managed_flag.transform_many(escaped)?;
let bundle_dir = get_space_bundle_dir(project_dir, space_id);
std::fs::create_dir_all(&bundle_dir)?;
let output_file = bundle_dir.join("saved_objects.ndjson");
use std::io::Write;
let mut file = std::fs::File::create(&output_file)?;
for obj in &flagged {
let json_line = serde_json::to_string(obj)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} saved object(s) for space {} to {}",
flagged.len(),
space_id.cyan(),
output_file.display()
);
Ok(flagged.len())
}
async fn bundle_space_workflows(project_dir: &Path, space_id: &str) -> Result<usize> {
let workflows_dir = get_space_workflows_dir(project_dir, space_id);
if !workflows_dir.exists() {
log::debug!(
"No workflows directory for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Bundling workflows for space {}", space_id.cyan());
let mut workflows = Vec::new();
for entry in std::fs::read_dir(&workflows_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let workflow = storage::read_json5_file(&path)?;
workflows.push(workflow);
}
}
let bundle_dir = get_space_bundle_dir(project_dir, space_id);
std::fs::create_dir_all(&bundle_dir)?;
let output_file = bundle_dir.join("workflows.ndjson");
use std::io::Write;
let mut file = std::fs::File::create(&output_file)?;
for workflow in &workflows {
let json_line = serde_json::to_string(workflow)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} workflow(s) for space {} to {}",
workflows.len(),
space_id.cyan(),
output_file.display()
);
Ok(workflows.len())
}
async fn bundle_space_agents(project_dir: &Path, space_id: &str) -> Result<usize> {
let agents_dir = get_space_agents_dir(project_dir, space_id);
if !agents_dir.exists() {
log::debug!(
"No agents directory for space {}, skipping",
space_id.cyan()
);
return Ok(0);
}
log::info!("Bundling agents for space {}", space_id.cyan());
let mut agents = Vec::new();
for entry in std::fs::read_dir(&agents_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let agent = storage::read_json5_file(&path)?;
agents.push(agent);
}
}
let bundle_dir = get_space_bundle_dir(project_dir, space_id);
std::fs::create_dir_all(&bundle_dir)?;
let output_file = bundle_dir.join("agents.ndjson");
use std::io::Write;
let mut file = std::fs::File::create(&output_file)?;
for agent in &agents {
let json_line = serde_json::to_string(agent)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} agent(s) for space {} to {}",
agents.len(),
space_id.cyan(),
output_file.display()
);
Ok(agents.len())
}
async fn bundle_space_tools(project_dir: &Path, space_id: &str) -> Result<usize> {
let tools_dir = get_space_tools_dir(project_dir, space_id);
if !tools_dir.exists() {
log::debug!("No tools directory for space {}, skipping", space_id.cyan());
return Ok(0);
}
log::info!("Bundling tools for space {}", space_id.cyan());
let mut tools = Vec::new();
for entry in std::fs::read_dir(&tools_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json") {
let tool = storage::read_json5_file(&path)?;
tools.push(tool);
}
}
let bundle_dir = get_space_bundle_dir(project_dir, space_id);
std::fs::create_dir_all(&bundle_dir)?;
let output_file = bundle_dir.join("tools.ndjson");
use std::io::Write;
let mut file = std::fs::File::create(&output_file)?;
for tool in &tools {
let json_line = serde_json::to_string(tool)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} tool(s) for space {} to {}",
tools.len(),
space_id.cyan(),
output_file.display()
);
Ok(tools.len())
}
async fn bundle_space_skills(project_dir: &Path, space_id: &str) -> Result<usize> {
let skills_dir = get_space_skills_dir(project_dir, space_id);
let skills = read_skill_values_from_dir(
&skills_dir,
&get_space_skills_manifest(project_dir, space_id),
)?;
if skills.is_empty() {
if skills_dir.exists() {
log::debug!("No skills selected for space {}, skipping", space_id.cyan());
} else {
log::debug!(
"No skills directory for space {}, skipping",
space_id.cyan()
);
}
return Ok(0);
}
log::info!("Bundling skills for space {}", space_id.cyan());
let bundle_dir = get_space_bundle_dir(project_dir, space_id);
std::fs::create_dir_all(&bundle_dir)?;
let output_file = bundle_dir.join("skills.ndjson");
use std::io::Write;
let mut file = std::fs::File::create(&output_file)?;
for skill in &skills {
let json_line = serde_json::to_string(skill)?;
writeln!(file, "{}", json_line)?;
}
log::info!(
"✓ Bundled {} skill(s) for space {} to {}",
skills.len(),
space_id.cyan(),
output_file.display()
);
Ok(skills.len())
}
fn read_skill_values_from_dir(skills_dir: &Path, manifest_path: &Path) -> Result<Vec<Value>> {
let manifest = if manifest_path.exists() {
Some(SkillsManifest::read(manifest_path)?)
} else {
None
};
if !skills_dir.exists() {
if let Some(manifest) = manifest
&& let Some(entry) = manifest.skills.first()
{
eyre::bail!(
"skill '{}' listed in {} was not found under {}",
entry.id,
manifest_path.display(),
skills_dir.display()
);
}
return Ok(Vec::new());
}
let mut skill_dirs = Vec::new();
for entry in std::fs::read_dir(skills_dir)? {
let entry = entry?;
let path = entry.path();
let metadata = std::fs::symlink_metadata(&path)?;
if metadata.file_type().is_symlink() {
eyre::bail!("skill directory cannot be a symlink: {}", path.display());
}
if metadata.is_dir() && path.join("SKILL.md").exists() {
skill_dirs.push(path);
}
}
skill_dirs.sort();
let mut skills = Vec::new();
for dir in skill_dirs {
skills.push(skill_to_value(&dir, true)?);
}
let Some(manifest) = manifest else {
return Ok(skills);
};
manifest
.skills
.iter()
.map(|entry| {
skills
.iter()
.find(|skill| skill.get("id").and_then(|id| id.as_str()) == Some(&entry.id))
.cloned()
.ok_or_else(|| {
eyre::eyre!(
"skill '{}' listed in {} was not found under {}",
entry.id,
manifest_path.display(),
skills_dir.display()
)
})
})
.collect()
}
fn skill_entry(skill: &Value) -> Result<SkillEntry> {
let id = skill
.get("id")
.and_then(|value| value.as_str())
.ok_or_else(|| eyre::eyre!("Skill missing 'id' field"))?;
let name = skill
.get("name")
.and_then(|value| value.as_str())
.unwrap_or(id);
Ok(SkillEntry::new(id, name))
}
fn get_space_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
project_dir.join(space_id)
}
fn get_space_manifest_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_dir(project_dir, space_id).join("manifest")
}
fn get_space_objects_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_dir(project_dir, space_id).join("objects")
}
fn get_space_workflows_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_dir(project_dir, space_id).join("workflows")
}
fn get_space_agents_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_dir(project_dir, space_id).join("agents")
}
fn get_space_tools_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_dir(project_dir, space_id).join("tools")
}
fn get_space_skills_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_dir(project_dir, space_id).join("skills")
}
fn get_space_skills_manifest(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_manifest_dir(project_dir, space_id).join("skills.yml")
}
fn get_space_saved_objects_manifest(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_manifest_dir(project_dir, space_id).join("saved_objects.json")
}
fn get_space_workflows_manifest(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_manifest_dir(project_dir, space_id).join("workflows.yml")
}
fn get_space_agents_manifest(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_manifest_dir(project_dir, space_id).join("agents.yml")
}
fn get_space_tools_manifest(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_manifest_dir(project_dir, space_id).join("tools.yml")
}
fn get_space_bundle_dir(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
project_dir.join("bundle").join(space_id)
}
fn get_space_file(project_dir: &Path, space_id: &str) -> std::path::PathBuf {
get_space_dir(project_dir, space_id).join("space.json")
}
#[derive(Debug, Default, Clone)]
pub struct DependencySummary {
pub agents: usize,
pub tools: usize,
pub skills: usize,
pub workflows: usize,
}
impl DependencySummary {
pub fn new() -> Self {
Self::default()
}
pub fn is_empty(&self) -> bool {
self.agents == 0 && self.tools == 0 && self.skills == 0 && self.workflows == 0
}
pub fn total(&self) -> usize {
self.agents + self.tools + self.skills + self.workflows
}
pub fn format_summary(&self) -> String {
let mut parts = Vec::new();
if self.agents > 0 {
parts.push(format!("{} agent(s)", self.agents));
}
if self.tools > 0 {
parts.push(format!("{} tool(s)", self.tools));
}
if self.skills > 0 {
parts.push(format!("{} skill(s)", self.skills));
}
if self.workflows > 0 {
parts.push(format!("{} workflow(s)", self.workflows));
}
parts.join(", ")
}
}
async fn resolve_and_add_dependencies(
project_dir: &Path,
space_id: &str,
client: &KibanaClient,
initial_deps: Vec<Dependency>,
detected_version: KibanaVersion,
force: bool,
) -> Result<DependencySummary> {
let client = client.space(space_id)?;
let mut pending_deps = initial_deps;
let mut processed_ids = HashSet::new();
let mut summary = DependencySummary::new();
while let Some(dep) = pending_deps.pop() {
match dep {
Dependency::Agent(id) => {
if !force
&& !KibanaClient::supports_capability(&detected_version, ApiCapability::Agents)
{
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(
&detected_version,
ApiCapability::Agents
)
);
continue;
}
if !processed_ids.insert(format!("agent:{}", id)) {
continue;
}
let manifest_path = get_space_agents_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
AgentsManifest::read(&manifest_path)?
} else {
AgentsManifest::new()
};
if manifest.contains_id(&id) {
continue;
}
log::info!("Automatically adding dependent agent: {}", id.cyan());
let path = format!("api/agent_builder/agents/{}", id);
let response = client.get(&path).await?;
if !response.status().is_success() {
log::warn!(
"Failed to fetch dependent agent {}: {}",
id,
response.status()
);
continue;
}
let agent: Value = response.json().await?;
let name = agent.get("name").and_then(|v| v.as_str()).unwrap_or(&id);
if manifest.add_agent(AgentEntry::new(&id, name)) {
manifest.write(&manifest_path)?;
let agents_dir = get_space_agents_dir(project_dir, space_id);
std::fs::create_dir_all(&agents_dir)?;
let agent_file = agents_dir.join(format!("{}.json", name));
let json = storage::to_string_with_multiline(&agent)?;
std::fs::write(&agent_file, json)?;
summary.agents += 1;
pending_deps.extend(find_agent_dependencies(&agent));
}
}
Dependency::Tool(id) => {
if !force
&& !KibanaClient::supports_capability(&detected_version, ApiCapability::Tools)
{
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(
&detected_version,
ApiCapability::Tools
)
);
continue;
}
if !processed_ids.insert(format!("tool:{}", id)) {
continue;
}
let manifest_path = get_space_tools_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
ToolsManifest::read(&manifest_path)?
} else {
ToolsManifest::new()
};
if manifest.contains(&id) {
continue;
}
log::info!("Automatically adding dependent tool: {}", id.cyan());
let path = format!("api/agent_builder/tools/{}", id);
let response = client.get(&path).await?;
if !response.status().is_success() {
log::warn!(
"Failed to fetch dependent tool {}: {}",
id,
response.status()
);
continue;
}
let tool: Value = response.json().await?;
let name = tool.get("name").and_then(|v| v.as_str()).unwrap_or(&id);
if manifest.add_tool(id.clone()) {
manifest.write(&manifest_path)?;
let tools_dir = get_space_tools_dir(project_dir, space_id);
std::fs::create_dir_all(&tools_dir)?;
let tool_file = tools_dir.join(format!("{}.json", name));
let json = storage::to_string_with_multiline(&tool)?;
std::fs::write(&tool_file, json)?;
summary.tools += 1;
pending_deps.extend(find_tool_dependencies(&tool));
}
}
Dependency::Skill(id) => {
if !force
&& !KibanaClient::supports_capability(&detected_version, ApiCapability::Skills)
{
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(
&detected_version,
ApiCapability::Skills
)
);
continue;
}
if !processed_ids.insert(format!("skill:{}", id)) {
continue;
}
log::info!("Automatically adding dependent skill: {}", id.cyan());
let path = format!("api/agent_builder/skills/{}", id);
let response = client.get(&path).await?;
if !response.status().is_success() {
log::warn!(
"Failed to fetch dependent skill {}: {}",
id,
response.status()
);
continue;
}
let skill: Value = response.json().await?;
if is_readonly(&skill) {
log::debug!("Skipping readonly dependent skill {}", id);
continue;
}
let manifest_path = get_space_skills_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
SkillsManifest::read(&manifest_path)?
} else {
SkillsManifest::new()
};
if manifest.contains_id(&id) {
continue;
}
let skills_dir = get_space_skills_dir(project_dir, space_id);
std::fs::create_dir_all(&skills_dir)?;
skill_to_directory(&skills_dir, &skill)?;
manifest.add_skill(skill_entry(&skill)?);
manifest.write(&manifest_path)?;
summary.skills += 1;
pending_deps.extend(find_skill_dependencies(&skill));
}
Dependency::Workflow(id) => {
if !force
&& !KibanaClient::supports_capability(
&detected_version,
ApiCapability::Workflows,
)
{
log::warn!(
"{}",
KibanaClient::unsupported_capability_reason(
&detected_version,
ApiCapability::Workflows
)
);
continue;
}
if !processed_ids.insert(format!("workflow:{}", id)) {
continue;
}
let manifest_path = get_space_workflows_manifest(project_dir, space_id);
let mut manifest = if manifest_path.exists() {
WorkflowsManifest::read(&manifest_path)?
} else {
WorkflowsManifest::new()
};
if manifest.contains_id(&id) {
continue;
}
log::info!("Automatically adding dependent workflow: {}", id.cyan());
let path = workflow_resource_path(&id);
let response = client.get_internal(&path).await?;
if !response.status().is_success() {
log::warn!(
"Failed to fetch dependent workflow {}: {}",
id,
response.status()
);
continue;
}
let workflow: Value = response.json().await?;
let name = workflow.get("name").and_then(|v| v.as_str()).unwrap_or(&id);
if manifest.add_workflow(WorkflowEntry::new(&id, name)) {
manifest.write(&manifest_path)?;
let workflows_dir = get_space_workflows_dir(project_dir, space_id);
std::fs::create_dir_all(&workflows_dir)?;
let workflow_file =
workflows_dir.join(format!("{}.json", workflow_file_stem(name)));
let json = storage::to_string_with_multiline(&workflow)?;
std::fs::write(&workflow_file, json)?;
summary.workflows += 1;
pending_deps.extend(find_workflow_dependencies(&workflow));
}
}
}
}
Ok(summary)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
#[serial_test::serial]
fn test_load_kibana_client_no_url() {
unsafe {
std::env::remove_var("KIBANA_URL");
std::env::remove_var("KIBANA_USERNAME");
std::env::remove_var("KIBANA_PASSWORD");
std::env::remove_var("KIBANA_APIKEY");
}
let result = load_kibana_client(".");
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("KIBANA_URL"));
}
#[test]
#[serial_test::serial]
fn test_load_kibana_client_with_url() {
unsafe {
std::env::set_var("KIBANA_URL", "http://localhost:5601");
std::env::remove_var("KIBANA_USERNAME");
std::env::remove_var("KIBANA_PASSWORD");
std::env::remove_var("KIBANA_APIKEY");
}
let result = load_kibana_client(".");
assert!(result.is_ok());
unsafe {
std::env::remove_var("KIBANA_URL");
}
}
#[test]
#[serial_test::serial]
fn test_load_kibana_client_invalid_url() {
unsafe {
std::env::set_var("KIBANA_URL", "not-a-valid-url");
}
let result = load_kibana_client(".");
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("Invalid KIBANA_URL")
);
unsafe {
std::env::remove_var("KIBANA_URL");
}
}
#[test]
#[serial_test::serial]
fn test_get_target_space_ids() {
let temp_dir = tempfile::TempDir::new().unwrap();
let project_dir = temp_dir.path();
let url = url::Url::parse("http://localhost:5601").unwrap();
let manifest_path = project_dir.join("spaces.yml");
let manifest = crate::kibana::spaces::SpacesManifest::with_spaces(vec![
crate::kibana::spaces::SpaceEntry::new("default".into(), "Default".into()),
crate::kibana::spaces::SpaceEntry::new("marketing".into(), "Marketing".into()),
]);
manifest.write(&manifest_path).unwrap();
let client = KibanaClient::builder(url)
.auth(Auth::None)
.max_concurrency(8)
.spaces(
manifest
.spaces
.iter()
.map(|space| (space.id.clone(), space.name.clone())),
)
.build()
.unwrap();
let mut ids = get_target_space_ids(&client, None);
ids.sort();
let mut expected = vec!["default".to_string(), "marketing".to_string()];
expected.sort();
assert_eq!(ids, expected);
let ids = get_target_space_ids(&client, Some(&["default".to_string()]));
assert_eq!(ids, vec!["default"]);
let mut ids = get_target_space_ids(
&client,
Some(&["default".to_string(), "marketing".to_string()]),
);
ids.sort();
assert_eq!(ids, expected);
}
#[test]
#[serial_test::serial]
fn test_get_target_space_ids_from_manifest() {
let temp_dir = tempfile::TempDir::new().unwrap();
let project_dir = temp_dir.path();
let ids = get_target_space_ids_from_manifest(project_dir, None);
assert_eq!(ids, vec!["default"]);
let manifest_path = project_dir.join("spaces.yml");
let manifest = crate::kibana::spaces::SpacesManifest::with_spaces(vec![
crate::kibana::spaces::SpaceEntry::new("default".into(), "Default".into()),
crate::kibana::spaces::SpaceEntry::new("eng".into(), "Engineering".into()),
]);
manifest.write(&manifest_path).unwrap();
let mut ids = get_target_space_ids_from_manifest(project_dir, None);
ids.sort();
let mut expected = vec!["default".to_string(), "eng".to_string()];
expected.sort();
assert_eq!(ids, expected);
let ids = get_target_space_ids_from_manifest(project_dir, Some(&["eng".to_string()]));
assert_eq!(ids, vec!["eng"]);
}
#[test]
fn test_dependency_summary_format() {
let mut summary = DependencySummary::new();
assert_eq!(summary.format_summary(), "");
assert_eq!(summary.total(), 0);
summary.agents = 1;
assert_eq!(summary.format_summary(), "1 agent(s)");
assert_eq!(summary.total(), 1);
summary.tools = 2;
assert_eq!(summary.format_summary(), "1 agent(s), 2 tool(s)");
assert_eq!(summary.total(), 3);
summary.skills = 3;
assert_eq!(
summary.format_summary(),
"1 agent(s), 2 tool(s), 3 skill(s)"
);
assert_eq!(summary.total(), 6);
summary.workflows = 3;
assert_eq!(
summary.format_summary(),
"1 agent(s), 2 tool(s), 3 skill(s), 3 workflow(s)"
);
assert_eq!(summary.total(), 9);
}
#[test]
fn test_parse_capability_accepts_skills_aliases() {
assert_eq!(parse_capability("skills"), Some(ApiCapability::Skills));
assert_eq!(parse_capability("skill"), Some(ApiCapability::Skills));
}
#[test]
fn test_should_process_api_supports_skills_filters() {
let skills = vec!["skills".to_string()];
assert!(should_process_api("skills", Some(&skills)));
assert!(!should_process_api("tools", Some(&skills)));
let singular = vec!["skill".to_string()];
assert!(should_process_api("skills", Some(&singular)));
let mixed = vec![
"agents".to_string(),
"skill".to_string(),
"workflows".to_string(),
];
assert!(should_process_api("agents", Some(&mixed)));
assert!(should_process_api("skills", Some(&mixed)));
assert!(should_process_api("workflows", Some(&mixed)));
assert!(!should_process_api("tools", Some(&mixed)));
}
#[test]
fn test_workflow_file_stem_sanitizes_lowercases_and_replaces_spaces() {
assert_eq!(
workflow_file_stem("Daily Workflow: Threat Hunting"),
"daily_workflow__threat_hunting"
);
assert_eq!(workflow_file_stem(" Workflow One "), "workflow_one");
assert_eq!(workflow_file_stem(""), "unnamed");
}
#[test]
fn test_skills_unsupported_warning_names_version_and_maturity() {
let version = parse_kibana_version("9.3.0").unwrap();
let warning = KibanaClient::unsupported_capability_reason(&version, ApiCapability::Skills);
assert!(warning.contains("API 'skills'"));
assert!(warning.contains("9.4.0+"));
assert!(warning.contains("detected 9.3.0"));
assert!(warning.contains("Experimental as of 9.4"));
}
#[test]
fn test_is_older_major_minor() {
let recorded = parse_kibana_version("9.3.2").unwrap();
assert!(is_older_major_minor(
&parse_kibana_version("9.2.9").unwrap(),
&recorded
));
assert!(!is_older_major_minor(
&parse_kibana_version("9.3.0").unwrap(),
&recorded
));
assert!(!is_older_major_minor(
&parse_kibana_version("10.0.0").unwrap(),
&recorded
));
}
#[test]
fn test_version_warning_downcast() {
let report = version_warning("test warning");
assert_eq!(version_warning_message(&report), Some("test warning"));
}
}