#![cfg(unix)]
#[path = "support/broker.rs"]
mod broker_fixture;
#[path = "support/plan.rs"]
mod plan_fixture;
use std::{
fs,
os::unix::fs::PermissionsExt,
process::Command,
time::{SystemTime, UNIX_EPOCH},
};
use shepherd_cli::{
BindRootDispatchRequest, CarrierAttachmentExpectationRequest, DispatchService, DispatchStore,
NativeBroker, PreparePendingDispatchRequest,
shepherd::{
Harness,
dispatch::{AgentId, DispatchState, PendingLaunchState, ProjectId, Role, RunId, SessionId},
registry::Registry,
},
};
#[test]
fn broker_fixture_child() {
broker_fixture::child_main();
}
#[test]
fn direct_start_root_prepares_exact_lane_specialists_through_native_broker() {
let suffix = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock")
.as_nanos();
let root = std::env::temp_dir().join(format!("shepherd-direct-start-{suffix:x}"));
fs::create_dir_all(root.join(".shepherd")).expect("private fixture");
fs::set_permissions(&root, fs::Permissions::from_mode(0o700)).expect("private root");
let root = fs::canonicalize(root).expect("canonical root");
let project_id = "018f47ce-72d7-7f64-9eb1-2f651d521c2a";
fs::write(
root.join(".shepherd/project.json"),
format!("{{\"id\":\"{project_id}\"}}"),
)
.expect("project identity");
fs::create_dir_all(root.join("docs")).expect("task directory");
fs::write(root.join("docs/task.md"), "Bounded direct-start fixture.\n").expect("task");
fs::write(root.join(".gitignore"), "installed/\n").expect("exclude installed carrier");
let installed = root.join("installed");
let compiled = Command::new(env!("CARGO_BIN_EXE_shepherd"))
.args(["compile", "--target", "pi", "--out"])
.arg(&installed)
.current_dir(&root)
.output()
.expect("canonical compiler");
assert!(
compiled.status.success(),
"{}",
String::from_utf8_lossy(&compiled.stderr)
);
assert!(
Command::new("git")
.args(["init", "--quiet"])
.current_dir(&root)
.status()
.expect("fixture git")
.success()
);
let baseline = plan_fixture::open_execution(&root, "v657", &["lane-a"], 3);
let now = i64::try_from(
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock")
.as_millis(),
)
.expect("clock fits i64");
Registry::open_migrated(root.join(".shepherd/shepherd.db")).expect("registry")
.execute(
"INSERT INTO projects (id, name, created_at, updated_at) VALUES (?1, ?2, ?3, ?3) ON CONFLICT(id) DO NOTHING",
(project_id, "direct-start fixture", now),
).expect("registry project");
let store = DispatchStore::new(root.join(".shepherd/runs"));
let service = DispatchService::with_project_root(
store.clone(),
ProjectId::new(project_id).expect("project"),
&root,
)
.with_installed_package(&installed, installed.join(".shepherd-generated.json"));
let binding = BindRootDispatchRequest {
schema: "shepherd.dispatch-request/1".into(),
run: Some("v657".into()),
harness: Harness::Pi,
session_id: "root-session".into(),
role_carrier: "shepherd:shepherd".into(),
mode: shepherd_cli::shepherd::dispatch::RootMode::Execution,
lease_ms: 600_000,
};
service
.bind_root(binding.clone(), now)
.expect("exact native execution root");
let mut planning_root = binding;
planning_root.mode = shepherd_cli::shepherd::dispatch::RootMode::Planting;
assert!(
service.bind_root(planning_root, now).is_err(),
"planning cannot gain execution authority"
);
let endpoint = fs::canonicalize("/tmp")
.expect("short socket root")
.join(format!("ss-{}-{suffix:x}", std::process::id()))
.join("broker.sock");
let broker = NativeBroker::start(service, &endpoint).expect("native broker");
let mut parent = broker.connect().expect("root connection");
parent
.register_parent(
Harness::Pi,
Role::Shepherd,
SessionId::new("root-session").expect("root"),
SessionId::new("root-session").expect("root"),
None,
)
.expect("native root peer");
let run = RunId::new("v657").expect("run");
for (role, work_kind, write_scope) in [
(
Role::Coder,
"production-code",
vec!["src/plan-fixture-0.rs".to_owned()],
),
(
Role::Worker,
"artifact",
vec![".shepherd/runs/v657/lanes/lane-a/workers/worker.md".into()],
),
(Role::Auditor, "review", vec![]),
] {
let result_dir = if role == Role::Worker {
"workers"
} else {
"reports"
};
let request = PreparePendingDispatchRequest {
schema: "shepherd.pending-dispatch-request/2".into(),
run: Some("v657".into()),
role: role.to_string(),
work_kind: work_kind.into(),
lane: Some("lane-a".into()),
parent_dispatch_id: None,
replaces_agent_id: None,
baseline: baseline.clone(),
read_scope: vec!["docs/**".into(), "src/plan-fixture-0.rs".into()],
write_scope,
result_artifact: format!(".shepherd/runs/v657/lanes/lane-a/{result_dir}/{role}.md"),
review_artifact: format!(".shepherd/runs/v657/lanes/lane-a/reviews/{role}.md"),
task_file: "docs/task.md".into(),
child_session_id: format!("session-{role}"),
lease_ms: 60_000,
expected_attachment: CarrierAttachmentExpectationRequest {
target: Harness::Pi,
role: role.to_string(),
agent_id: format!("direct-{role}"),
attachment_kind: "pi-skill-path".into(),
},
};
for (label, mut invalid) in [
("fabricated-lead", request.clone()),
("absent-plan-lane", request.clone()),
("wrong-work-kind", request.clone()),
] {
match label {
"fabricated-lead" => invalid.parent_dispatch_id = Some("invented-conductor".into()),
"absent-plan-lane" => invalid.lane = Some("unplanned-lane".into()),
"wrong-work-kind" => invalid.work_kind = "planning".into(),
_ => unreachable!(),
}
invalid.expected_attachment.agent_id = format!("denied-{role}-{label}");
assert!(
parent.prepare(invalid).is_err(),
"{role} must reject {label}"
);
assert!(
store
.load_pending_for_agent(
&run,
&AgentId::new(format!("denied-{role}-{label}")).expect("agent")
)
.is_err()
);
}
if role != Role::Coder {
let mut production = request.clone();
production.write_scope = vec!["src/plan-fixture-0.rs".into()];
production.expected_attachment.agent_id = format!("denied-{role}-production");
assert!(
parent.prepare(production).is_err(),
"direct root cannot grant {role} source writes"
);
}
let launch = parent.prepare(request.clone()).unwrap_or_else(|error| {
panic!("direct /start must prepare {role} without a Conductor: {error}")
});
let pending = store
.load_pending_for_agent(
&run,
&AgentId::new(format!("direct-{role}")).expect("agent"),
)
.expect("native pending receipt");
assert_eq!(pending.caller_role, Role::Shepherd);
assert_eq!(pending.role, role);
assert!(
pending.parent_dispatch_id.is_none(),
"no invented child lead"
);
assert_eq!(
pending.lane.as_ref().map(|lane| lane.as_str()),
Some("lane-a")
);
assert_eq!(pending.launch_state, PendingLaunchState::Pending);
assert_eq!(pending.expected_attachment.role, role);
let mut provider = broker_fixture::LiveProvider::launch_prepared(
&mut parent,
&endpoint,
&installed,
&root.join("provider-fixtures"),
request,
&launch,
)
.unwrap_or_else(|error| panic!("direct /start must activate actual {role} peer: {error}"));
let active = provider.record();
assert_eq!(active.state, DispatchState::Active);
assert_eq!(active.role, role);
assert!(active.parent_agent_id.is_none());
assert_eq!(active.root_session_id.as_str(), "root-session");
let result = root.join(
active
.result_artifact
.as_deref()
.expect("Native issued result path"),
);
fs::create_dir_all(result.parent().expect("result directory")).expect("result directory");
fs::write(
&result,
format!("Deterministic {role} fixture completed.\n"),
)
.expect("owned result");
let stopped = provider
.complete()
.expect("actual child sends native completion event");
assert_eq!(
stopped.agent_id,
AgentId::new(format!("direct-{role}")).expect("agent")
);
assert_eq!(stopped.state, DispatchState::Stopped);
assert!(stopped.parent_agent_id.is_none());
}
}