use super::native_types::DescribeResponse;
use super::runner::Plugin;
use crate::config::PluginEntry;
use crate::error::Result;
use crate::negotiation::FailurePolicy;
use crate::participant::{Event, EventCtx, Participant, Projection};
use crate::participant_config::{
effective_subscriptions, InvocationOverrides, LocalPluginEntry,
};
use crate::store::Store;
use std::collections::BTreeMap;
pub use crate::plugin::native_proto::NativeOutcome;
use crate::plugin::native_proto::NativeProtocol;
pub struct NativePluginParticipant {
name: String,
plugin: Plugin,
projection: Projection,
subscriptions: Vec<Event>,
failure_policies: BTreeMap<Event, FailurePolicy>,
retry_budget: usize,
wants_context: bool,
}
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,
wants_context: describe.wants_context,
})
}
}
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 = match ctx.post {
Some(t) => t.clone(),
None => match ctx.store.load_task(ctx.task_id) {
Ok(t) => t,
Err(_) => return None,
},
};
let proto = NativeProtocol::new(
&self.plugin,
&self.name,
event,
task,
self.retry_budget,
self.wants_context,
ctx.identity.to_string(),
ctx.task_before.map(|t| Box::new(t.clone())),
ctx.commit.map(str::to_string),
ctx.overrides.to_vec(),
);
Some(proto)
}
}
#[cfg(test)]
#[path = "native_participant_describe_tests.rs"]
mod describe_tests;