use super::native_types::{DescribeResponse, ProposeConflict, ProposeOk, ProposeResponse};
use super::runner::Plugin;
use crate::config::PluginEntry;
use crate::error::Result;
use crate::negotiation::{AttemptClass, CommitPolicy, FailurePolicy, Protocol};
use crate::participant::{Event, EventCtx, Participant, Projection};
use crate::participant_config::{
effective_subscriptions, InvocationOverrides, LocalPluginEntry,
};
use crate::store::Store;
use crate::task::Task;
use serde_json::Value;
use std::collections::BTreeMap;
pub struct NativePluginParticipant {
name: String,
plugin: Plugin,
projection: Projection,
subscriptions: Vec<Event>,
failure_policies: BTreeMap<Event, FailurePolicy>,
retry_budget: usize,
}
pub const DEFAULT_NATIVE_RETRY_BUDGET: usize = 5;
impl NativePluginParticipant {
pub fn from_describe(
store: &Store,
name: String,
entry: &PluginEntry,
local: Option<&LocalPluginEntry>,
invocation: &InvocationOverrides,
describe: DescribeResponse,
) -> Result<Self> {
let plugin = Plugin::resolve(store, &name, entry);
let resolved = effective_subscriptions(&name, entry, local, invocation);
let declared: std::collections::BTreeSet<Event> =
describe.subscriptions.iter().copied().collect();
let failure_policies: BTreeMap<Event, FailurePolicy> = resolved
.into_iter()
.filter(|(ev, _)| declared.contains(ev))
.collect();
let subscriptions: Vec<Event> = failure_policies.keys().copied().collect();
let projection = describe.projection.into_projection()?;
let retry_budget = describe
.retry_budget
.unwrap_or(DEFAULT_NATIVE_RETRY_BUDGET);
Ok(Self {
name,
plugin,
projection,
subscriptions,
failure_policies,
retry_budget,
})
}
}
#[derive(Debug, Clone)]
pub struct NativeOutcome {
pub task_projection: Value,
pub commit_policy: CommitPolicy,
}
pub struct NativeProtocol<'a> {
plugin: &'a Plugin,
name: &'a str,
event: Event,
task: Box<Task>,
accepted: Option<ProposeOk>,
pending_conflict: Option<ProposeConflict>,
retry_budget: usize,
}
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 Some(resp) = self.plugin.propose(self.event, &self.task)? 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, super::native_types::CommitPolicyWire::into_policy);
NativeOutcome {
task_projection: ok.task,
commit_policy,
}
}
fn retry_budget(&self) -> usize {
self.retry_budget
}
}
fn classify(
resp: ProposeResponse,
accepted: &mut Option<ProposeOk>,
pending_conflict: &mut Option<ProposeConflict>,
name: &str,
) -> AttemptClass {
match (resp.ok, resp.conflict) {
(Some(ok), _) => {
*accepted = Some(ok);
AttemptClass::Ok
}
(None, Some(conflict)) => {
*pending_conflict = Some(conflict);
AttemptClass::Conflict
}
(None, None) => AttemptClass::Other(format!(
"plugin `{name}` propose returned neither ok nor conflict"
)),
}
}
impl Participant for NativePluginParticipant {
type Outcome = NativeOutcome;
type Protocol<'a>
= NativeProtocol<'a>
where
Self: 'a;
fn name(&self) -> &str {
&self.name
}
fn subscriptions(&self) -> &[Event] {
&self.subscriptions
}
fn projection(&self) -> &Projection {
&self.projection
}
fn failure_policy(&self, event: Event) -> FailurePolicy {
self.failure_policies
.get(&event)
.copied()
.unwrap_or(FailurePolicy::BestEffort)
}
fn protocol<'a>(
&'a self,
event: Event,
ctx: EventCtx<'a>,
) -> Option<Self::Protocol<'a>> {
let task = ctx.store.load_task(ctx.task_id).ok()?;
Some(NativeProtocol {
plugin: &self.plugin,
name: &self.name,
event,
task: Box::new(task),
accepted: None,
pending_conflict: None,
retry_budget: self.retry_budget,
})
}
}
#[cfg(test)]
impl<'a> NativeProtocol<'a> {
pub(crate) fn __test_new(
plugin: &'a Plugin, name: &'a str, event: Event, task: Task, retry_budget: usize,
) -> Self {
Self {
plugin, name, event, task: Box::new(task),
accepted: None, pending_conflict: None, retry_budget,
}
}
pub(crate) fn __test_record_ok(&mut self, ok: ProposeOk) { self.accepted = Some(ok); }
pub(crate) fn __test_record_conflict(&mut self, c: ProposeConflict) {
self.pending_conflict = Some(c);
}
pub(crate) fn __test_task_title(&self) -> String { self.task.title.clone() }
pub(crate) fn __test_has_pending_conflict(&self) -> bool {
self.pending_conflict.is_some()
}
}
#[cfg(test)]
pub(crate) fn __test_classify(
resp: ProposeResponse, accepted: &mut Option<ProposeOk>,
pending_conflict: &mut Option<ProposeConflict>, name: &str,
) -> AttemptClass {
classify(resp, accepted, pending_conflict, name)
}
#[cfg(test)]
#[path = "native_participant_test_helpers.rs"]
mod test_helpers;
#[cfg(test)]
#[path = "native_participant_describe_tests.rs"]
mod describe_tests;
#[cfg(test)]
#[path = "native_participant_proto_tests.rs"]
mod proto_tests;