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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,8 @@ pub const fn http_application_owner_kind(
| ApplicationSurfaceOperation::Runtime
| ApplicationSurfaceOperation::ProjectList
| ApplicationSurfaceOperation::ProjectSearch
| ApplicationSurfaceOperation::ProjectContext => return None,
| ApplicationSurfaceOperation::ProjectContext
| ApplicationSurfaceOperation::AdminSync => return None,
ApplicationSurfaceOperation::FactStoreCurate
| ApplicationSurfaceOperation::FactStoreAdd
| ApplicationSurfaceOperation::FactStoreSearch
Expand Down
275 changes: 72 additions & 203 deletions crates/tracedecay-cli/src/commands/index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,10 @@ use tracedecay_runtime_core::config::ProfileRoot;

use tracedecay_project::project::TraceDecay;

use tracedecay_contracts::graph_tool::GraphToolResultV1;
use tracedecay_contracts::retrieval::{AdminSyncAdmissionV1, AdminSyncResultV1};
use tracedecay_tool_catalog::ApplicationSurfaceOperation;

use super::daemon::daemon_tool_json;

/// True when the global DB has zero registered projects (or can't be opened
Expand Down Expand Up @@ -217,13 +221,7 @@ async fn brokered_init(
// eventually demand it does not survive a daemon restart. Request the
// reconciliation explicitly so the message below reports something that
// actually happened.
let reconcile = tracedecay::daemon::call_default_tool_awaiting_project_open(
handshake,
"tracedecay_admin_sync",
serde_json::json!({}),
init_deadline,
)
.await;
let reconcile = admin_sync(handshake.clone(), init_deadline).await;
let reconcile = match reconcile {
Err(error) => {
if !code_index_reconciliation_is_optional(project_path, &error).await {
Expand All @@ -237,27 +235,18 @@ async fn brokered_init(
}
Ok(reconcile) => reconcile,
};
match admin_sync_status(&reconcile).as_deref() {
// Status `queued` means the daemon accepted the reconcile demand into
// its pre-mount queue. Init's user-facing confirmation names that
// request (`requested`), matching the brokered-init contract tests and
// dogfood journeys, not the internal queue noun.
Some("queued") => eprintln!(
match reconcile.status {
// `queued` means the daemon accepted the reconcile demand into its
// pre-mount queue. Init's confirmation names that request
// (`requested`), not the internal queue noun.
AdminSyncAdmissionV1::Queued => eprintln!(
"initialized {}; daemon code-index reconciliation requested",
project_path.display()
),
Some("not_applicable") => eprintln!(
AdminSyncAdmissionV1::NotApplicable => eprintln!(
"initialized {}; code indexing does not apply to this non-Git project",
project_path.display()
),
Some(status) => eprintln!(
"initialized {}; daemon code-index reconciliation status is {status}",
project_path.display()
),
None => eprintln!(
"initialized {}; daemon code-index reconciliation returned no status",
project_path.display()
),
}
Ok(())
}
Expand All @@ -275,17 +264,26 @@ fn reject_brokered_folder_options(
})
}

fn admin_sync_status(envelope: &serde_json::Value) -> Option<String> {
let text = envelope
.get("content")?
.as_array()?
.iter()
.find_map(|block| block.get("text").and_then(|text| text.as_str()))?;
serde_json::from_str::<serde_json::Value>(text)
.ok()?
.get("status")?
.as_str()
.map(str::to_owned)
/// Asks the project's owner for the operator's code-index reconcile.
async fn admin_sync(
handshake: tracedecay_daemon_protocol::DaemonHandshake,
deadline: tokio::time::Instant,
) -> tracedecay_domain::errors::Result<AdminSyncResultV1> {
match crate::tool_command::owner_operation_result(
handshake,
ApplicationSurfaceOperation::AdminSync,
serde_json::json!({}),
deadline,
)
.await?
{
GraphToolResultV1::AdminSync(result) => Ok(result),
_ => Err(tracedecay_domain::errors::TraceDecayError::project_route(
"owner_result_mismatch",
false,
"the project owner answered tracedecay_admin_sync with another operation's result",
)),
}
}

async fn code_index_reconciliation_is_optional(
Expand Down Expand Up @@ -423,161 +421,47 @@ mod init_bootstrap_tests {
);
}

use super::daemon_precondition_tests::SocketEnvGuard;

static SOCKET_ENV_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());

type RecordedRequests = std::sync::Arc<std::sync::Mutex<Vec<(String, serde_json::Value)>>>;

/// A daemon on `socket` that authenticates `connections` requests and
/// answers each with `respond(tool_name, request_id)`.
fn spawn_fixture_daemon(
temp: &Path,
socket: &Path,
connections: usize,
respond: fn(&str, serde_json::Value) -> serde_json::Value,
) -> (tokio::task::JoinHandle<()>, RecordedRequests) {
use tokio::io::{AsyncBufReadExt, AsyncWriteExt};

let authority = tracedecay_daemon_identity::authority::DaemonAuthority::acquire(
temp,
&tracedecay_daemon_protocol::DaemonEndpoint::Unix(socket.to_path_buf()),
env!("CARGO_PKG_VERSION"),
)
.expect("publish the fixture daemon's authority record");
let auth_token = authority.auth_token().to_owned();
let listener = tokio::net::UnixListener::bind(socket).unwrap();
let recorded = RecordedRequests::default();
let responder = {
let recorded = std::sync::Arc::clone(&recorded);
tokio::spawn(async move {
let _authority = authority;
for _ in 0..connections {
let (stream, _addr) = listener.accept().await.unwrap();
let (reader, mut writer) = stream.into_split();
let mut lines = tokio::io::BufReader::new(reader).lines();
let preface = lines.next_line().await.unwrap().unwrap();
assert!(
tracedecay_daemon_protocol::DaemonAuthPreface::from_line(preface.trim())
.expect("auth preface")
.authenticate(&auth_token),
"init must present the daemon token"
);
let _handshake_line = lines.next_line().await.unwrap().unwrap();
let request_line = lines.next_line().await.unwrap().unwrap();
let request: serde_json::Value = serde_json::from_str(&request_line).unwrap();
let name = request["params"]["name"]
.as_str()
.unwrap_or_default()
.to_owned();
recorded
.lock()
.unwrap()
.push((name.clone(), request["params"]["arguments"].clone()));
let response = respond(&name, request["id"].clone());
writer
.write_all(serde_json::to_string(&response).unwrap().as_bytes())
.await
.unwrap();
writer.write_all(b"\n").await.unwrap();
writer.shutdown().await.unwrap();
}
})
};
(responder, recorded)
}

fn empty_tool_result(_tool: &str, id: serde_json::Value) -> serde_json::Value {
serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"result": { "content": [] },
})
}

/// Init's "daemon code-index reconciliation requested" must describe a
/// request that actually crossed the wire: admission first, then the
/// explicit `tracedecay_admin_sync` reconcile. Without the second call the
/// first index only starts if the background full-server upgrade survives
/// long enough to demand it, which a daemon restart silently discards.
#[tokio::test]
async fn brokered_init_requests_a_real_code_index_reconciliation() {
let _serialize = SOCKET_ENV_TEST_LOCK.lock().await;
let temp = tempfile::TempDir::new().unwrap();
let project = temp.path().join("project");
let profile = temp.path().join("profile");
std::fs::create_dir_all(&project).unwrap();
let socket = temp.path().join("daemon.sock");
let (responder, recorded) =
spawn_fixture_daemon(temp.path(), &socket, 2, empty_tool_result);
let _socket_env = SocketEnvGuard::set(&socket);

let handshake = test_handshake(&project, &profile);
brokered_init(&project, &[], &[], &handshake)
.await
.expect("brokered init against the fixture daemon");
tokio::time::timeout(std::time::Duration::from_secs(5), responder)
.await
.expect("fixture daemon must observe both requests")
.expect("fixture daemon task");

let recorded = recorded.lock().unwrap();
let names: Vec<&str> = recorded.iter().map(|(name, _)| name.as_str()).collect();
assert_eq!(
names,
["tracedecay_status", "tracedecay_admin_sync"],
"init must request the reconcile it reports, after admission"
);
assert_eq!(
recorded[0].1["admission_only"],
serde_json::json!(true),
"the bootstrap status call stays admission-only"
);
}

/// A worktree parked on a corrupt publication authority refuses the
/// reconcile; the operator reads its cause and remedy as whole fields,
/// not folded into one sentence cut at the diagnostic bound.
#[tokio::test]
async fn brokered_init_prints_a_parked_refusal_as_fields() {
fn refuse_sync(tool: &str, id: serde_json::Value) -> serde_json::Value {
if tool != "tracedecay_admin_sync" {
return empty_tool_result(tool, id);
}
let parked = tracedecay_contracts::code_index_freshness::CodeIndexConvergenceParkedV1 {
reason: format!("the publication authority is corrupt: {}", "x".repeat(600)),
blocked_reason: None,
remediation: "run `tracedecay daemon restart`".to_owned(),
parked_at_micros: 1,
observed_passes: 1,
retries_on_wake: false,
};
serde_json::to_value(tracedecay_mcp::tool_error_response(
id,
tool,
&parked.publication_authority_corrupt_error(),
))
.unwrap()
}

let _serialize = SOCKET_ENV_TEST_LOCK.lock().await;
let temp = tempfile::TempDir::new().unwrap();
let project = temp.path().join("project");
let profile = temp.path().join("profile");
std::fs::create_dir_all(&project).unwrap();
let socket = temp.path().join("daemon.sock");
let (responder, _recorded) = spawn_fixture_daemon(temp.path(), &socket, 2, refuse_sync);
let _socket_env = SocketEnvGuard::set(&socket);

let error = brokered_init(&project, &[], &[], &test_handshake(&project, &profile))
.await
.expect_err("a parked worktree refuses the reconcile");
responder.await.expect("fixture daemon task");
#[test]
fn a_parked_admin_sync_refusal_prints_as_fields() {
let parked = tracedecay_contracts::code_index_freshness::CodeIndexConvergenceParkedV1 {
reason: format!("the publication authority is corrupt: {}", "x".repeat(600)),
blocked_reason: None,
remediation: "run `tracedecay daemon restart`".to_owned(),
parked_at_micros: 1,
observed_passes: 1,
retries_on_wake: false,
};
let detail = parked
.publication_authority_corrupt_error()
.project_route_typed_detail()
.cloned()
.expect("a parked refusal carries its typed detail");
let refusal = tracedecay_mcp::application_output::tool_result::ApplicationRefusal {
operation: ApplicationSurfaceOperation::AdminSync,
binding_id: tracedecay_tool_catalog::BindingId::new("binding.cli.admin_sync.v1")
.unwrap(),
problem: tracedecay_contracts::ApplicationProblemEnvelope::new(
tracedecay_contracts::ResultContractRef::new(
tracedecay_tool_catalog::SchemaId::new(
"schema.application.primitive.admin-sync.result",
)
.unwrap(),
1,
)
.unwrap(),
tracedecay_contracts::RequestId::new("request.cli.admin-sync").unwrap(),
tracedecay_contracts::ApplicationProblem::from_detail(detail),
)
.unwrap(),
};

assert_eq!(
crate::commands::process_error_text(error),
crate::commands::process_error_text(refusal.into_error()),
format!(
"project route error (code_index_publication_authority_corrupt)\n\
"project route error (application.code-index.parked)\n\
Parked cause: the publication authority is corrupt: {}\n\
Parked remedy: run `tracedecay daemon restart`\n\
Retries on wake: false",
Expand Down Expand Up @@ -709,35 +593,20 @@ pub(crate) async fn handle_sync(
)
.await?;
let handshake = super::daemon::client_handshake(profile, Some(&resolved.project_path))?;
let result = tracedecay::daemon::call_default_tool(
&handshake,
"tracedecay_admin_sync",
serde_json::json!({}),
)
.await?;
let deadline = tokio::time::Instant::now() + crate::tool_command::tool_command_deadline()?;
let result = admin_sync(handshake, deadline).await?;
if verbose {
eprintln!(
"{}",
serde_json::to_string_pretty(&result).unwrap_or_default()
);
eprintln!("{}", serde_json::to_string_pretty(&result)?);
}
match admin_sync_status(&result).as_deref() {
Some("queued") => eprintln!(
match result.status {
AdminSyncAdmissionV1::Queued => eprintln!(
"code-index reconciliation queued via daemon for {}",
resolved.project_path.display()
),
Some("not_applicable") => eprintln!(
AdminSyncAdmissionV1::NotApplicable => eprintln!(
"code indexing does not apply to {}",
resolved.project_path.display()
),
Some(status) => eprintln!(
"code-index reconciliation status is {status} for {}",
resolved.project_path.display()
),
None => eprintln!(
"daemon code-index reconciliation returned no status for {}",
resolved.project_path.display()
),
}
if doctor {
tracedecay::doctor::run_doctor(profile, crate::cloud::doctor_network_probes()).await?;
Expand Down
Loading
Loading