use super::native_types::{
CommitPolicyWire, EventCtxWire, ProposeConflict, ProposeOk, ProposeResponse,
};
use super::runner::{event_subcommand_arg, Plugin};
use crate::error::Result;
use crate::negotiation::{AttemptClass, CommitPolicy, Protocol};
use crate::participant::Event;
use crate::task::Task;
use serde_json::Value;
#[derive(Debug, Clone)]
pub struct NativeOutcome {
pub task_projection: Value,
pub commit_policy: CommitPolicy,
}
pub struct NativeProtocol<'a> {
pub(crate) plugin: &'a Plugin,
pub(crate) name: &'a str,
pub(crate) event: Event,
pub(crate) task: Box<Task>,
pub(crate) accepted: Option<ProposeOk>,
pub(crate) pending_conflict: Option<ProposeConflict>,
pub(crate) retry_budget: usize,
pub(crate) wants_context: bool,
pub(crate) identity: String,
pub(crate) task_before: Option<Box<Task>>,
pub(crate) commit: Option<String>,
pub(crate) override_tokens: Vec<String>,
}
impl<'a> NativeProtocol<'a> {
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
plugin: &'a Plugin,
name: &'a str,
event: Event,
task: Task,
retry_budget: usize,
wants_context: bool,
identity: String,
task_before: Option<Box<Task>>,
commit: Option<String>,
override_tokens: Vec<String>,
) -> Self {
Self {
plugin,
name,
event,
task: Box::new(task),
accepted: None,
pending_conflict: None,
retry_budget,
wants_context,
identity,
task_before,
commit,
override_tokens,
}
}
}
impl Protocol for NativeProtocol<'_> {
type Outcome = NativeOutcome;
fn propose(&mut self) -> Result<AttemptClass> {
if !self.plugin.auth_check() {
return Ok(AttemptClass::Other(format!(
"plugin `{}` auth-check failed",
self.name
)));
}
let ctx_json = if self.wants_context {
Some(EventCtxWire::for_event(
event_subcommand_arg(self.event),
&self.identity,
self.task.repo.clone(),
self.task_before.as_deref(),
self.commit.as_deref(),
&self.override_tokens,
)?)
} else {
None
};
let Some(resp) =
self.plugin
.propose(self.event, &self.task, ctx_json.as_deref())?
else {
return Ok(AttemptClass::Other(format!(
"plugin `{}` propose returned no usable response",
self.name
)));
};
Ok(classify(
resp,
&mut self.accepted,
&mut self.pending_conflict,
self.name,
))
}
fn fetch_remote_view(&mut self) -> Result<()> {
self.pending_conflict = None;
Ok(())
}
fn pushed(&mut self) -> Self::Outcome {
let ok = self.accepted.take().unwrap_or(ProposeOk {
task: Value::Null,
commit_policy: None,
});
let commit_policy = ok
.commit_policy
.map_or_else(CommitPolicy::default, CommitPolicyWire::into_policy);
NativeOutcome {
task_projection: ok.task,
commit_policy,
}
}
fn retry_budget(&self) -> usize {
self.retry_budget
}
}
pub(crate) fn classify(
resp: ProposeResponse,
accepted: &mut Option<ProposeOk>,
pending_conflict: &mut Option<ProposeConflict>,
name: &str,
) -> AttemptClass {
let ProposeResponse { ok, conflict, reject, extra } = resp;
if let Some(ok) = ok {
*accepted = Some(ok);
return AttemptClass::Ok;
}
if let Some(conflict) = conflict {
*pending_conflict = Some(conflict);
return AttemptClass::Conflict;
}
if let Some(reject) = reject {
return AttemptClass::Reject(format!(
"plugin `{name}` rejected: {}",
reject.reason
));
}
if !extra.is_empty() {
let v = extra.keys().cloned().collect::<Vec<_>>().join(", ");
return AttemptClass::Other(format!(
"plugin `{name}` propose returned unknown variant(s): {v}"
));
}
AttemptClass::Other(format!(
"plugin `{name}` propose returned neither ok nor conflict"
))
}
#[cfg(test)]
#[path = "native_proto_test_helpers.rs"]
mod test_helpers;
#[cfg(test)]
#[path = "native_participant_proto_tests.rs"]
mod proto_tests;