mod diag;
mod dispatch;
mod limits;
mod native_participant;
mod native_types;
mod participant;
mod runner;
mod types;
pub use dispatch::{dispatch_push, dispatch_sync};
pub use native_participant::{NativeOutcome, NativePluginParticipant};
pub use native_types::{
CommitPolicyWire, DescribeResponse, ProjectionWire, ProposeConflict, ProposeOk,
ProposeResponse,
};
pub use participant::{LegacyOutcome, LegacyPluginParticipant};
pub use runner::Plugin;
pub use types::{PushResponse, SyncCreate, SyncDelete, SyncReport, SyncUpdate};
use crate::commit_policy::{plan, Contribution, PlanOp};
use crate::error::Result;
use crate::negotiation::{CommitPolicy, FailurePolicy};
use crate::participant::Projection;
use crate::store::{task_lock, Store};
use crate::task::Task;
use chrono::Utc;
use serde_json::{Map, Value};
pub enum ContributionPayload {
Legacy(PushResponse),
Native(Value),
}
pub struct PushContribution {
pub name: String,
pub projection: Projection,
pub payload: ContributionPayload,
pub failure_policy: FailurePolicy,
pub commit_policy: CommitPolicy,
}
pub fn apply_push_contributions(
store: &Store,
task_id: &str,
contributions: &[PushContribution],
) -> Result<()> {
if contributions.is_empty() {
return Ok(());
}
let plan_steps = plan(
&contributions
.iter()
.map(|c| Contribution {
name: c.name.clone(),
failure_policy: c.failure_policy,
commit_policy: c.commit_policy.clone(),
})
.collect::<Vec<_>>(),
&format!("balls: update external for {task_id}"),
)?;
let _g = task_lock(store, task_id)?;
let mut task = store.load_task(task_id)?;
let now = Utc::now();
for op in plan_steps {
match op {
PlanOp::Apply(i) => {
apply_one(&mut task, &contributions[i], now);
store.save_task(&task)?;
}
PlanOp::Commit(msg) => {
store.commit_task(task_id, &msg)?;
}
}
}
Ok(())
}
fn apply_one(task: &mut Task, c: &PushContribution, now: chrono::DateTime<chrono::Utc>) {
let payload_obj = payload_as_map(&c.name, &c.payload);
project_overlay(task, &payload_obj, &c.projection);
task.synced_at.insert(c.name.clone(), now);
task.touch();
}
fn payload_as_map(name: &str, payload: &ContributionPayload) -> Map<String, Value> {
match payload {
ContributionPayload::Legacy(resp) => {
let mut external = Map::new();
external.insert(name.to_string(), Value::Object(resp.0.clone()));
let mut root = Map::new();
root.insert("external".into(), Value::Object(external));
root
}
ContributionPayload::Native(v) => match v {
Value::Object(map) => map.clone(),
_ => Map::new(),
},
}
}
fn project_overlay(task: &mut Task, payload: &Map<String, Value>, projection: &Projection) {
if !projection.owns.is_empty() {
let mut current = task_to_object(task);
for field in &projection.owns {
let key = crate::plugin::native_types::field_wire_name(*field);
if let Some(v) = payload.get(key) {
current.insert(key.to_string(), v.clone());
}
}
if let Ok(merged) = serde_json::from_value::<Task>(Value::Object(current)) {
*task = merged;
}
}
if let Some(Value::Object(payload_external)) = payload.get("external") {
for prefix in &projection.external_prefixes {
if let Some(slice) = payload_external.get(prefix) {
task.external.insert(prefix.clone(), slice.clone());
}
}
}
}
fn task_to_object(task: &Task) -> Map<String, Value> {
serde_json::to_value(task).unwrap().as_object().unwrap().clone()
}
#[cfg(test)]
#[path = "mod_apply_tests.rs"]
mod apply_tests;