#![cfg(unix)]
#[path = "support/broker.rs"]
mod broker_fixture;
#[path = "support/plan.rs"]
mod plan_fixture;
use serde_json::json;
use sha2::{Digest, Sha256};
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, DispatchRecord, DispatchState, PendingLaunchState, ProjectId, Role, RunId,
SessionId,
},
registry::Registry,
},
};
fn digest(bytes: &[u8]) -> String {
Sha256::digest(bytes)
.iter()
.map(|byte| format!("{byte:02x}"))
.collect()
}
#[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 fixture = tempfile::Builder::new()
.prefix("shepherd-direct-start-")
.tempdir()
.expect("create fixture root");
let root = fs::canonicalize(fixture.path()).expect("canonical root");
fs::create_dir_all(root.join(".shepherd")).expect("private fixture");
fs::set_permissions(&root, fs::Permissions::from_mode(0o700)).expect("private root");
let project_id = "018f47ce-72d7-7f64-9eb1-2f651d521c2a";
plan_fixture::write_project_identity(&root, project_id, 1000);
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");
let mut completed_subject: Option<DispatchRecord> = None;
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");
let pending_task = store
.load_pending_for_agent(&run, &active.agent_id)
.expect("active pending task");
let task_digest = pending_task
.task_sha256
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
let evidence_path =
format!(".shepherd/runs/v657/lanes/lane-a/reports/direct-{role}-evidence.txt");
let evidence_bytes = format!("Deterministic {role} evidence.\n");
fs::write(root.join(&evidence_path), evidence_bytes.as_bytes()).expect("evidence");
let startup = active
.startup_attachment
.as_ref()
.expect("startup attachment");
let typed = match role {
Role::Coder => json!({
"schema":"shepherd.coder-result/1", "run":"v657", "lane":"lane-a",
"node":"direct-coder", "role":"coder", "outcome":"deterministic Coder fixture",
"task_digest":task_digest, "startup_skill":startup.skill,
"skill_bundle_digest":startup.bundle_digest, "worktree":root.display().to_string(),
"baseline_commit":pending_task.baseline_commit, "owned_paths":["src/plan-fixture-0.rs"],
"changed_paths":["src/plan-fixture-0.rs"], "result_artifact":active.result_artifact,
"evidence":[{"path":evidence_path,"sha256":digest(evidence_bytes.as_bytes()),"bytes":evidence_bytes.len(),"command":"cargo test","exit_status":0}],
"status":"green"
}),
Role::Worker => json!({
"schema":"shepherd.worker-result/1", "run":"v657", "lane":"lane-a",
"node":"direct-worker", "role":"worker", "work_kind":"artifact",
"deliverable":"deterministic Worker fixture", "source_paths":["docs/task.md"],
"owned_scope":[active.result_artifact], "budget":{"tool_calls":1,"seconds":1},
"output_shape":"typed worker result", "result_artifact":active.result_artifact,
"evidence":[{"path":evidence_path,"sha256":digest(evidence_bytes.as_bytes()),"bytes":evidence_bytes.len(),"command":"cargo test","exit_status":0}],
"status":"complete", "task_digest":task_digest, "startup_skill":startup.skill,
"skill_bundle_digest":startup.bundle_digest
}),
Role::Auditor => {
let subject = completed_subject.as_ref().expect("Coder subject");
let subject_pending = store
.load_pending_for_agent(&run, &subject.agent_id)
.expect("subject pending");
let subject_digest = subject.result_sha256.clone().expect("subject digest");
json!({
"schema":"shepherd.review-result/1", "run":"v657", "lane":"lane-a",
"mode":"auditor-posthoc", "reviewer_role":"auditor", "candidate_commit":"0123456789abcdef0123456789abcdef01234567",
"input_digest":subject_digest, "startup_skill":"reviewing", "skill_bundle_digest":startup.bundle_digest,
"result_channel":"native-result", "review_artifact":active.review_artifact,
"subject_result_artifact":subject.result_artifact, "subject_task_digest":subject_pending.task_sha256.iter().map(|byte| format!("{byte:02x}")).collect::<String>(),
"verdict":"pass", "findings":[], "report_path":"reports/direct-auditor.md"
})
}
_ => unreachable!("direct-start fixture role is bounded"),
};
let typed_bytes = serde_json::to_vec(&typed).expect("typed result");
fs::write(&result, &typed_bytes).expect("owned result");
if role == Role::Auditor {
let review = root.join(active.review_artifact.as_ref().expect("review artifact"));
fs::create_dir_all(review.parent().expect("review directory"))
.expect("review directory");
fs::write(review, &typed_bytes).expect("review 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());
if role == Role::Coder {
completed_subject = Some(stopped);
}
}
}