use super::runner::Plugin;
use super::types::{PushResponse, SyncReport};
use crate::config::PluginEntry;
use crate::error::Result;
use crate::negotiation::{AttemptClass, 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 std::collections::BTreeMap;
pub struct LegacyPluginParticipant {
name: String,
plugin: Plugin,
projection: Projection,
subscriptions: Vec<Event>,
failure_policies: BTreeMap<Event, FailurePolicy>,
sync_filter: Option<String>,
}
impl LegacyPluginParticipant {
pub fn from_entry(
store: &Store,
name: String,
entry: &PluginEntry,
sync_filter: Option<String>,
) -> Self {
Self::resolved(store, name, entry, None, &InvocationOverrides::default(), sync_filter)
}
pub fn resolved(
store: &Store,
name: String,
entry: &PluginEntry,
local: Option<&LocalPluginEntry>,
invocation: &InvocationOverrides,
sync_filter: Option<String>,
) -> Self {
let plugin = Plugin::resolve(store, &name, entry);
let failure_policies = effective_subscriptions(&name, entry, local, invocation);
let subscriptions = failure_policies.keys().copied().collect();
let projection = Projection::external_only(name.clone());
Self {
name,
plugin,
projection,
subscriptions,
failure_policies,
sync_filter,
}
}
}
#[derive(Debug)]
pub enum LegacyOutcome {
Push(Option<PushResponse>),
Sync(Option<SyncReport>),
}
pub enum LegacyProtocol<'a> {
Push {
plugin: &'a Plugin,
task: Box<Task>,
response: Option<PushResponse>,
},
Sync {
plugin: &'a Plugin,
tasks: Vec<Task>,
filter: Option<&'a str>,
report: Option<SyncReport>,
},
}
impl Protocol for LegacyProtocol<'_> {
type Outcome = LegacyOutcome;
fn propose(&mut self) -> Result<AttemptClass> {
match self {
LegacyProtocol::Push { plugin, task, response } => {
if !plugin.auth_check() {
return Ok(AttemptClass::Other("auth check failed".into()));
}
Ok(match plugin.push(task).ok().flatten() {
Some(r) => {
*response = Some(r);
AttemptClass::Ok
}
None => AttemptClass::Other("plugin push returned no data".into()),
})
}
LegacyProtocol::Sync { plugin, tasks, filter, report } => {
if !plugin.auth_check() {
return Ok(AttemptClass::Other("auth check failed".into()));
}
Ok(match plugin.sync(tasks, *filter).ok().flatten() {
Some(r) => {
*report = Some(r);
AttemptClass::Ok
}
None => AttemptClass::Other("plugin sync returned no data".into()),
})
}
}
}
fn fetch_remote_view(&mut self) -> Result<()> {
Ok(())
}
fn pushed(&mut self) -> Self::Outcome {
match self {
LegacyProtocol::Push { response, .. } => LegacyOutcome::Push(response.take()),
LegacyProtocol::Sync { report, .. } => LegacyOutcome::Sync(report.take()),
}
}
fn retry_budget(&self) -> usize {
1
}
}
impl Participant for LegacyPluginParticipant {
type Outcome = LegacyOutcome;
type Protocol<'a>
= LegacyProtocol<'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>> {
if matches!(event, Event::Sync) {
let tasks = ctx.store.all_tasks().ok()?;
return Some(LegacyProtocol::Sync {
plugin: &self.plugin,
tasks,
filter: self.sync_filter.as_deref(),
report: None,
});
}
let task = match ctx.post {
Some(t) => t.clone(),
None => ctx.store.load_task(ctx.task_id).ok()?,
};
Some(LegacyProtocol::Push {
plugin: &self.plugin,
task: Box::new(task),
response: None,
})
}
}
#[cfg(test)]
#[path = "participant_tests.rs"]
mod tests;