use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::process::{Command, ExitStatus, Stdio};
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use serde_json::{json, Map, Value};
use crate::edits::Operation;
use crate::event::Source;
use crate::graph::NodeStatus;
use crate::ledger::{LaunchRecord, RunPaths};
use crate::plan::Node;
use crate::projection::RunState;
use crate::taskgraph::{QualifiedId, BINARY_ENV};
const SHADOW_SOURCE: &str = "onepipeline-writeback";
const COMMAND_LIMIT: Duration = Duration::from_secs(60);
const RETRY_AFTER: Duration = Duration::from_millis(250);
const CLOSEOUT_WAIT: Duration = Duration::from_millis(2_250);
#[derive(Clone, PartialEq)]
struct Snapshot {
project: QualifiedId,
dir: PathBuf,
nodes: BTreeMap<String, Node>,
statuses: BTreeMap<String, NodeStatus>,
settlements: BTreeMap<String, Value>,
project_metadata: BTreeMap<String, Value>,
}
#[derive(Default, PartialEq, Eq)]
enum WorkerState {
#[default]
Idle,
Working,
StopRequested,
}
#[derive(Default)]
struct Pending {
latest: Option<Snapshot>,
last_success: Option<Snapshot>,
worker: WorkerState,
}
impl Pending {
fn queue(&mut self, snapshot: Snapshot) -> bool {
if self.latest.as_ref() == Some(&snapshot) {
return false;
}
if self.worker == WorkerState::Idle
&& self.latest.is_none()
&& self.last_success.as_ref() == Some(&snapshot)
{
return false;
}
self.latest = Some(snapshot);
true
}
}
pub struct Writeback {
pending: Arc<(Mutex<Pending>, Condvar)>,
}
impl Writeback {
pub fn start(binary: PathBuf, paths: &RunPaths, launch: &LaunchRecord) -> Option<Self> {
let pending = Arc::new((Mutex::new(Pending::default()), Condvar::new()));
let worker_pending = Arc::clone(&pending);
let run_dir = paths.dir.clone();
let launch_dir = if launch.dir.as_os_str().is_empty() {
PathBuf::from(".")
} else {
launch.dir.clone()
};
std::thread::Builder::new()
.name(format!("writeback-{}", paths.run))
.spawn(move || worker(binary, launch_dir, run_dir, worker_pending))
.ok()?;
let writer = Self { pending };
Some(writer)
}
pub fn publish(&self, paths: &RunPaths, launch: &LaunchRecord, state: &RunState) {
let Ok(project) = launch.project.parse() else {
return;
};
let snapshot = Snapshot {
project,
dir: paths.dir.join("writeback"),
nodes: all_nodes(paths, state),
statuses: state.statuses(),
settlements: settlements(paths),
project_metadata: state
.plan
.as_ref()
.map(|plan| {
let mut metadata = BTreeMap::from([
(
"onepipeline.schema_version".into(),
json!(plan.schema_version),
),
("onepipeline.concurrency".into(), json!(plan.concurrency)),
]);
if let Some(goal) = &plan.goal {
metadata.insert("onepipeline.goal".into(), json!(goal));
}
if let Some(name) = &plan.name {
metadata.insert("onepipeline.name".into(), json!(name));
}
metadata
})
.unwrap_or_default(),
};
let (lock, ready) = &*self.pending;
if let Ok(mut pending) = lock.lock() {
if pending.queue(snapshot) {
ready.notify_one();
}
}
}
pub fn wait_briefly(&self) {
let deadline = Instant::now() + CLOSEOUT_WAIT;
let (lock, ready) = &*self.pending;
let Ok(mut pending) = lock.lock() else { return };
while (pending.latest.is_some() || pending.worker == WorkerState::Working)
&& Instant::now() < deadline
{
let wait = deadline.saturating_duration_since(Instant::now());
let Ok((next, _)) = ready.wait_timeout(pending, wait) else {
return;
};
pending = next;
}
}
}
impl Drop for Writeback {
fn drop(&mut self) {
let (lock, ready) = &*self.pending;
if let Ok(mut pending) = lock.lock() {
pending.worker = WorkerState::StopRequested;
ready.notify_one();
}
}
}
fn worker(
binary: PathBuf,
launch_dir: PathBuf,
run_dir: PathBuf,
pending: Arc<(Mutex<Pending>, Condvar)>,
) {
let mut failing = false;
loop {
let snapshot = {
let (lock, ready) = &*pending;
let mut state = match lock.lock() {
Ok(state) => state,
Err(_) => return,
};
while state.latest.is_none() && state.worker != WorkerState::StopRequested {
state = match ready.wait(state) {
Ok(state) => state,
Err(_) => return,
};
}
if state.worker == WorkerState::StopRequested && state.latest.is_none() {
return;
}
if state.worker == WorkerState::Idle {
state.worker = WorkerState::Working;
}
state
.latest
.take()
.expect("the worker was woken by a snapshot")
};
match project(&binary, &launch_dir, &run_dir, &snapshot) {
Ok(()) => {
if failing {
eprintln!(
"onetaskgraph write-back recovered for '{}'",
snapshot.project
);
}
failing = false;
let (lock, _) = &*pending;
if let Ok(mut state) = lock.lock() {
state.last_success = Some(snapshot.clone());
if state.worker == WorkerState::StopRequested && state.latest.is_none() {
return;
}
}
}
Err(error) => {
if !failing {
eprintln!(
"onetaskgraph write-back failed for '{}': {error}; retrying",
snapshot.project
);
}
failing = true;
std::thread::sleep(RETRY_AFTER);
let (lock, ready) = &*pending;
let Ok(mut state) = lock.lock() else { return };
if state.worker == WorkerState::StopRequested {
return;
}
if state.latest.is_none() {
state.latest = Some(snapshot);
ready.notify_one();
}
}
}
let (lock, ready) = &*pending;
if let Ok(mut state) = lock.lock() {
if state.worker == WorkerState::Working {
state.worker = WorkerState::Idle;
}
ready.notify_all();
}
}
}
fn project(
binary: &Path,
launch_dir: &Path,
run_dir: &Path,
snapshot: &Snapshot,
) -> Result<(), String> {
let destination_project = destination_project(binary, launch_dir, run_dir, snapshot)?;
let origins = destination_origins(binary, launch_dir, run_dir, snapshot)?;
write_shadow(snapshot, &origins, &destination_project)?;
let root = snapshot.dir.to_string_lossy().into_owned();
let shadow_project = format!("{SHADOW_SOURCE}:{}", project_file(&snapshot.project));
let args = [
"project",
"copy",
&shadow_project,
"--to",
snapshot.project.source(),
"--json",
"--set",
&format!("sources.{SHADOW_SOURCE}.plugin=local-md"),
"--set",
&format!("sources.{SHADOW_SOURCE}.config.root={root}"),
];
let output = bounded_output(binary, launch_dir, run_dir, "project-copy", &args)?;
if output.status.success() {
Ok(())
} else {
Err(format!(
"copy exited {}: {}",
exit(&output.status),
String::from_utf8_lossy(&output.stderr).trim()
))
}
}
fn destination_project(
binary: &Path,
launch_dir: &Path,
run_dir: &Path,
snapshot: &Snapshot,
) -> Result<DestinationProjectItem, String> {
let args = ["project", "show", snapshot.project.as_str(), "--json"];
let output = bounded_output(binary, launch_dir, run_dir, "project-show", &args)?;
if !output.status.success() {
return Err(format!(
"project show exited {}: {}",
exit(&output.status),
String::from_utf8_lossy(&output.stderr).trim()
));
}
let response: ProjectPage =
serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
if !response.errors.is_empty() {
return Err("project show returned partial results".to_owned());
}
let mut items = response.items.into_iter();
let project = items
.next()
.ok_or_else(|| format!("project '{}' was not found", snapshot.project))?;
if items.next().is_some() || project.id != snapshot.project {
return Err(format!(
"project show returned the wrong project for '{}'",
snapshot.project
));
}
let _ = (response.next, response.plan);
Ok(project.item)
}
fn destination_origins(
binary: &Path,
launch_dir: &Path,
run_dir: &Path,
snapshot: &Snapshot,
) -> Result<BTreeMap<String, String>, String> {
let mut origins = BTreeMap::new();
let mut page: Option<String> = None;
let mut cursors = BTreeSet::new();
loop {
let mut args = vec![
"task".to_owned(),
"list".to_owned(),
"--project".to_owned(),
snapshot.project.as_str().to_owned(),
"--limit".to_owned(),
"10000".to_owned(),
"--json".to_owned(),
];
if let Some(token) = &page {
args.extend(["--page".to_owned(), token.clone()]);
}
let output = bounded_output(binary, launch_dir, run_dir, "task-list", &args)?;
if !output.status.success() {
return Err(format!(
"task list exited {}: {}",
output.status,
String::from_utf8_lossy(&output.stderr).trim()
));
}
let response: TaskPage =
serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
if !response.errors.is_empty() {
return Err("task list returned partial results".to_owned());
}
for task in response.items {
let node = task
.item
.metadata
.get("onepipeline.id")
.and_then(Value::as_str)
.filter(|id| !id.is_empty())
.ok_or_else(|| format!("task '{}' has no onepipeline.id", task.id.as_str()))?;
if origins
.insert(node.to_owned(), task.id.as_str().to_owned())
.is_some()
{
return Err(format!("project has more than one task for node '{node}'"));
}
}
let _ = response.plan;
let Some(next) = response.next else { break };
if next.is_empty() {
return Err("task list returned an empty next-page cursor".to_owned());
}
if !cursors.insert(next.clone()) {
return Err("task list repeated a next-page cursor".to_owned());
}
page = Some(next);
}
Ok(origins)
}
struct Output {
status: ExitStatus,
stdout: Vec<u8>,
stderr: Vec<u8>,
}
fn bounded_output<S: AsRef<std::ffi::OsStr>>(
binary: &Path,
launch_dir: &Path,
run_dir: &Path,
name: &str,
args: &[S],
) -> Result<Output, String> {
let stdout = run_dir.join(format!("writeback-{name}.stdout"));
let stderr = run_dir.join(format!("writeback-{name}.stderr"));
let stdout_file = std::fs::File::create(&stdout).map_err(|error| error.to_string())?;
let stderr_file = std::fs::File::create(&stderr).map_err(|error| error.to_string())?;
let mut child = Command::new(binary)
.current_dir(launch_dir)
.args(args)
.env_remove(BINARY_ENV)
.stdout(Stdio::from(stdout_file))
.stderr(Stdio::from(stderr_file))
.spawn()
.map_err(|error| format!("cannot run {}: {error}", binary.display()))?;
let started = Instant::now();
let status = loop {
match child.try_wait() {
Ok(Some(status)) => break status,
Ok(None) if started.elapsed() < COMMAND_LIMIT => {
std::thread::sleep(Duration::from_millis(25));
}
Ok(None) => {
let _ = child.kill();
let _ = child.wait();
return Err(format!(
"{name} exceeded {} seconds",
COMMAND_LIMIT.as_secs()
));
}
Err(error) => return Err(format!("cannot wait for {name}: {error}")),
}
};
Ok(Output {
status,
stdout: std::fs::read(stdout).map_err(|error| error.to_string())?,
stderr: std::fs::read(stderr).map_err(|error| error.to_string())?,
})
}
fn exit(status: &ExitStatus) -> String {
status
.code()
.map_or_else(|| "on a signal".into(), |code| code.to_string())
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct TaskPage {
items: Vec<DestinationTask>,
next: Option<String>,
plan: Value,
errors: Vec<Value>,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct ProjectPage {
items: Vec<DestinationProject>,
next: Option<String>,
plan: Value,
errors: Vec<Value>,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct DestinationProject {
id: QualifiedId,
item: DestinationProjectItem,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct DestinationProjectItem {
#[serde(rename = "id")]
_id: String,
#[serde(rename = "title")]
_title: String,
content: Option<String>,
#[serde(rename = "status")]
_status: DestinationStatus,
#[serde(rename = "labels")]
_labels: Vec<DestinationLabel>,
#[serde(rename = "url")]
_url: Option<String>,
#[serde(rename = "created_at")]
_created_at: Option<String>,
#[serde(rename = "updated_at")]
_updated_at: Option<String>,
metadata: BTreeMap<String, Value>,
#[serde(rename = "repositories")]
_repositories: Vec<String>,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct DestinationTask {
id: QualifiedId,
item: DestinationTaskItem,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct DestinationTaskItem {
#[serde(rename = "id")]
_id: String,
#[serde(rename = "title")]
_title: String,
#[serde(rename = "content")]
_content: Option<String>,
#[serde(rename = "status")]
_status: DestinationStatus,
#[serde(rename = "labels")]
_labels: Vec<DestinationLabel>,
#[serde(rename = "project")]
_project: Option<String>,
#[serde(rename = "url")]
_url: Option<String>,
#[serde(rename = "created_at")]
_created_at: Option<String>,
#[serde(rename = "updated_at")]
_updated_at: Option<String>,
metadata: BTreeMap<String, Value>,
#[serde(rename = "repositories")]
_repositories: Vec<String>,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct DestinationStatus {
#[serde(rename = "category")]
_category: DestinationStatusCategory,
#[serde(rename = "name")]
_name: String,
}
#[derive(Deserialize)]
#[serde(rename_all = "kebab-case")]
enum DestinationStatusCategory {
Backlog,
Todo,
InProgress,
Done,
Cancelled,
Unknown,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct DestinationLabel {
#[serde(rename = "id")]
_id: String,
#[serde(rename = "name")]
_name: String,
#[serde(rename = "color")]
_color: Option<String>,
}
fn write_shadow(
snapshot: &Snapshot,
origins: &BTreeMap<String, String>,
destination_project: &DestinationProjectItem,
) -> Result<(), String> {
let projects = snapshot.dir.join("projects");
let tasks = snapshot
.dir
.join("tasks")
.join(project_file(&snapshot.project));
std::fs::create_dir_all(&projects).map_err(|e| e.to_string())?;
std::fs::create_dir_all(&tasks).map_err(|e| e.to_string())?;
let mut project_metadata = destination_project.metadata.clone();
for (key, value) in &snapshot.project_metadata {
if project_metadata.contains_key(key) {
project_metadata.insert(key.clone(), value.clone());
}
}
project_metadata.insert(
"onetaskgraph.origin".into(),
json!(snapshot.project.as_str()),
);
document(
&projects.join(format!("{}.md", project_file(&snapshot.project))),
&json!({
"title": snapshot.project.native(), "metadata": project_metadata
}),
destination_project.content.as_deref().unwrap_or_default(),
)?;
for (id, node) in &snapshot.nodes {
let mut wire = serde_json::to_value(node)
.map_err(|e| e.to_string())?
.as_object()
.cloned()
.ok_or_else(|| "node did not serialize as a mapping".to_owned())?;
let title = wire
.remove("title")
.and_then(|v| v.as_str().map(str::to_owned))
.unwrap_or_default();
let content = wire
.remove("task")
.and_then(|v| v.as_str().map(str::to_owned))
.unwrap_or_default();
let deps = wire
.remove("deps")
.and_then(|v| v.as_array().cloned())
.unwrap_or_default();
let repo = wire
.remove("repo")
.and_then(|v| v.as_str().map(str::to_owned));
wire.remove("id");
let mut metadata = Map::new();
metadata.insert("onepipeline.id".into(), json!(id));
if let Some(origin) = origins.get(id) {
metadata.insert("onetaskgraph.origin".into(), json!(origin));
}
for (key, value) in wire {
metadata.insert(format!("onepipeline.{key}"), value);
}
if let Some(settlement) = snapshot.settlements.get(id) {
metadata.insert(crate::taskgraph::SETTLEMENT_KEY.into(), settlement.clone());
}
let local_deps: Vec<String> = deps
.iter()
.filter_map(Value::as_str)
.filter(|dep| !crate::graph::is_cross_dag(dep))
.map(task_file)
.collect();
let cross: Vec<String> = deps
.iter()
.filter_map(Value::as_str)
.filter(|dep| crate::graph::is_cross_dag(dep))
.map(str::to_owned)
.collect();
if !cross.is_empty() {
metadata.insert("onepipeline.deps".into(), json!(cross));
}
let status = snapshot
.statuses
.get(id)
.copied()
.unwrap_or(NodeStatus::Cancelled);
let mut front = Map::new();
front.insert("title".into(), json!(title));
front.insert("project".into(), json!(project_file(&snapshot.project)));
front.insert("status".into(), json!(category(status)));
front.insert("depends_on".into(), json!(local_deps));
front.insert("metadata".into(), Value::Object(metadata));
if let Some(repo) = repo {
if repo.starts_with("github.com/") {
front.insert("repositories".into(), json!([repo]));
} else if let Some(Value::Object(metadata)) = front.get_mut("metadata") {
metadata.insert("onepipeline.repo".into(), json!(repo));
}
}
document(
&tasks.join(format!("{}.md", task_file(id))),
&Value::Object(front),
&content,
)?;
}
Ok(())
}
fn document(path: &Path, front: &Value, body: &str) -> Result<(), String> {
let yaml = serde_norway::to_string(front).map_err(|e| e.to_string())?;
std::fs::write(path, format!("---\n{yaml}---\n{body}")).map_err(|e| e.to_string())
}
#[derive(Serialize)]
enum TaskCategory {
#[serde(rename = "in progress")]
InProgress,
#[serde(rename = "done")]
Done,
#[serde(rename = "cancelled")]
Cancelled,
#[serde(rename = "todo")]
Todo,
}
fn category(status: NodeStatus) -> TaskCategory {
match status {
NodeStatus::Running => TaskCategory::InProgress,
NodeStatus::Done | NodeStatus::Failed => TaskCategory::Done,
NodeStatus::Parked | NodeStatus::Cancelled | NodeStatus::Skipped => TaskCategory::Cancelled,
NodeStatus::Pending | NodeStatus::Ready | NodeStatus::Waiting | NodeStatus::Blocked => {
TaskCategory::Todo
}
}
}
fn all_nodes(paths: &RunPaths, state: &RunState) -> BTreeMap<String, Node> {
let mut nodes: BTreeMap<String, Node> = state
.plan
.as_ref()
.into_iter()
.flat_map(|plan| plan.tasks.iter())
.map(|node| (node.id.clone(), node.clone()))
.collect();
for event in crate::journal::read(&paths.journal()) {
if event.source != Source::Pipeline
|| event.kind.0 != crate::journal::PipelineKind::EditCommitted.as_str()
{
continue;
}
let Some(operations) = event
.payload
.get("operations")
.and_then(|v| serde_json::from_value::<Vec<Operation>>(v.clone()).ok())
else {
continue;
};
for operation in operations {
if let Operation::NodeAdded { node, .. } = operation {
nodes.insert(node.id.clone(), *node);
}
}
}
for node in state.graph.iter() {
nodes.insert(node.id.clone(), node.clone());
}
nodes
}
fn settlements(paths: &RunPaths) -> BTreeMap<String, Value> {
let mut found = BTreeMap::new();
for event in crate::journal::read(&paths.journal()) {
if event.source == Source::Pipeline
&& event.kind.0 == crate::journal::PipelineKind::NodeSettled.as_str()
{
if let Some(node) = event.labels.node.as_ref() {
found.insert(node.clone(), Value::Object(event.payload));
}
}
}
found
}
fn project_file(project: &QualifiedId) -> String {
encoded(project.as_str())
}
fn task_file(id: &str) -> String {
encoded(id)
}
fn encoded(value: &str) -> String {
value
.as_bytes()
.iter()
.map(|byte| format!("{byte:02x}"))
.collect()
}
#[cfg(test)]
mod tests {
use super::{Pending, Snapshot, TaskCategory, WorkerState};
use crate::graph::NodeStatus;
use std::collections::BTreeMap;
use std::path::PathBuf;
fn snapshot(status: NodeStatus) -> Snapshot {
Snapshot {
project: "plans:deduplication".parse().expect("a qualified project"),
dir: PathBuf::from("writeback"),
nodes: BTreeMap::new(),
statuses: BTreeMap::from([("node".to_owned(), status)]),
settlements: BTreeMap::new(),
project_metadata: BTreeMap::new(),
}
}
#[test]
fn returning_to_the_last_success_supersedes_a_different_pending_snapshot() {
let first = snapshot(NodeStatus::Pending);
let superseded = snapshot(NodeStatus::Running);
let mut pending = Pending {
latest: Some(superseded),
last_success: Some(first.clone()),
worker: WorkerState::Working,
};
assert!(pending.queue(first.clone()));
assert!(pending.latest.as_ref() == Some(&first));
}
#[test]
fn task_categories_remain_named_by_the_approved_contract() {
let contract = include_str!("../docs/contract.md");
for category in [
TaskCategory::Todo,
TaskCategory::InProgress,
TaskCategory::Done,
TaskCategory::Cancelled,
] {
let native = serde_json::to_value(category)
.expect("a task category serializes")
.as_str()
.expect("a task category is a string")
.replace(' ', "-");
assert!(
contract.contains(&format!("`{native}`")),
"docs/contract.md no longer names the projected status category `{native}`"
);
}
}
}