use super::native_participant::{NativeOutcome, NativePluginParticipant};
use super::participant::{LegacyOutcome, LegacyPluginParticipant};
use super::types::SyncReport;
use super::{ContributionPayload, PushContribution};
use crate::error::Result;
use crate::negotiation::{Accepted, NegotiationResult};
use crate::participant::{self, Event, EventCtx, Participant, Projection};
use crate::participant_config::InvocationOverrides;
use crate::store::Store;
use crate::task::Task;
pub struct DispatchInput<'a> {
pub store: &'a Store,
pub task_before: Option<&'a Task>,
pub task: &'a Task,
pub event: Event,
pub identity: &'a str,
pub commit: Option<&'a str>,
pub overrides: &'a InvocationOverrides,
pub override_tokens: &'a [String],
}
#[derive(Debug, Default)]
pub struct DispatchOutcome {
pub skipped: Vec<(String, String)>,
}
enum DispatchItem {
Contribution(PushContribution),
Skipped(String, String),
Inert,
}
pub fn dispatch_push(input: &DispatchInput) -> Result<DispatchOutcome> {
let &DispatchInput {
store,
task_before,
task,
event,
identity,
commit,
overrides,
override_tokens,
} = input;
debug_assert!(matches!(
event,
Event::Create | Event::Claim | Event::Review | Event::Close | Event::Update
));
let cfg = store.load_config()?;
let mut contributions = Vec::new();
let mut skipped = Vec::new();
for (name, entry) in cfg.plugins.iter().filter(|(_, e)| e.enabled) {
let ctx = EventCtx::new(event, store, &task.id, identity).with_context(
Some(task),
task_before,
commit,
override_tokens,
);
match dispatch_one_push(store, name, entry, event, ctx, overrides)? {
DispatchItem::Contribution(c) => contributions.push(c),
DispatchItem::Skipped(n, reason) => skipped.push((n, reason)),
DispatchItem::Inert => {}
}
}
super::apply_push_contributions(store, &task.id, &contributions, &skipped)?;
Ok(DispatchOutcome { skipped })
}
fn dispatch_one_push(
store: &Store,
name: &str,
entry: &crate::config::PluginEntry,
event: Event,
ctx: EventCtx<'_>,
overrides: &InvocationOverrides,
) -> Result<DispatchItem> {
let plugin = super::Plugin::resolve(store, name, entry);
if let Some(describe) = plugin.describe()? {
return run_native(store, name, entry, event, ctx, overrides, describe);
}
run_legacy(store, name, entry, event, ctx, overrides)
}
fn run_native(
store: &Store,
name: &str,
entry: &crate::config::PluginEntry,
event: Event,
ctx: EventCtx<'_>,
overrides: &InvocationOverrides,
describe: super::native_types::DescribeResponse,
) -> Result<DispatchItem> {
let participant = NativePluginParticipant::from_describe(
store,
name.to_string(),
entry,
None,
overrides,
describe,
)?;
if !participant.subscriptions().contains(&event) {
return Ok(DispatchItem::Inert);
}
let failure_policy = participant.failure_policy(event);
let projection = participant.projection().clone();
match participant::run(&participant, event, ctx)? {
NegotiationResult::Ok(Accepted {
outcome: NativeOutcome { task_projection, commit_policy },
..
}) => Ok(DispatchItem::Contribution(PushContribution {
name: name.to_string(),
projection,
payload: ContributionPayload::Native(task_projection),
failure_policy,
commit_policy,
})),
NegotiationResult::Skipped(reason) => {
Ok(DispatchItem::Skipped(name.to_string(), reason))
}
NegotiationResult::Staged(_) => Ok(DispatchItem::Inert),
}
}
fn run_legacy(
store: &Store,
name: &str,
entry: &crate::config::PluginEntry,
event: Event,
ctx: EventCtx<'_>,
overrides: &InvocationOverrides,
) -> Result<DispatchItem> {
let participant = LegacyPluginParticipant::resolved(
store,
name.to_string(),
entry,
None,
overrides,
None,
);
if !participant.subscriptions().contains(&event) {
return Ok(DispatchItem::Inert);
}
let failure_policy = participant.failure_policy(event);
let projection = Projection::external_only(name);
if let NegotiationResult::Ok(Accepted {
outcome: LegacyOutcome::Push(Some(r)),
commit_policy,
}) = participant::run(&participant, event, ctx)?
{
return Ok(DispatchItem::Contribution(PushContribution {
name: name.to_string(),
projection,
payload: ContributionPayload::Legacy(r),
failure_policy,
commit_policy,
}));
}
Ok(DispatchItem::Inert)
}
pub fn dispatch_sync(
store: &Store,
filter: Option<&str>,
identity: &str,
) -> Result<Vec<(String, SyncReport)>> {
let cfg = store.load_config()?;
let mut reports = Vec::new();
for (name, entry) in cfg.plugins.iter().filter(|(_, e)| e.enabled) {
let participant = LegacyPluginParticipant::from_entry(
store,
name.clone(),
entry,
filter.map(str::to_string),
);
let ctx = EventCtx::new(Event::Sync, store, filter.unwrap_or(""), identity);
if let Ok(NegotiationResult::Ok(Accepted {
outcome: LegacyOutcome::Sync(Some(r)),
..
})) = participant::run(&participant, Event::Sync, ctx)
{
reports.push((name.clone(), r));
}
}
Ok(reports)
}
pub fn dispatch_drop(store: &Store, task: &Task, identity: &str) -> Result<()> {
let cfg = store.load_config()?;
let overrides = InvocationOverrides::default();
for (name, entry) in cfg.plugins.iter().filter(|(_, e)| e.enabled) {
let ctx = EventCtx::new(Event::Drop, store, &task.id, identity);
let _ = dispatch_one_push(store, name, entry, Event::Drop, ctx, &overrides);
}
Ok(())
}
#[cfg(test)]
#[path = "dispatch_tests.rs"]
mod tests;