graphrecords-query 0.5.0

High-performance graph-based data records
Documentation
use crate::{
    Bare, BareValueDomain, Definite, EvaluateOperand, Explain, Failure, IndexDomain, Indexed,
    Labeled, Multiple, Operand, QueryResult, Single, Unordered, ValueDomain,
    error::grouping::MissingGroupAggregate,
    execution::EvaluationCache,
    index::GroupKey,
    operands::{OperandHandle, Partition},
    operations::{Apply, GroupKernel, Operation, OperationContext, Prepare},
    optimizer::{Estimate, OperationInputs, OptimizerHints, PlanIdentity, PlanInputs, Stats},
    registry::operation_manifest,
    traits::Broadcast,
};
use graphrecords_core::GraphRecord;

#[derive(Clone, Explain, Operation, OperationInputs, OptimizerHints, PlanIdentity, PlanInputs)]
#[operation(scope = Group)]
#[explain(label = "Broadcast")]
#[plan(optimizer_hints(empty = if_all))]
pub struct BroadcastOperation;

impl Prepare for BroadcastOperation {
    type Prepared<'a> = ();

    fn prepare<'a>(
        &'a self,
        _graphrecord: &'a GraphRecord,
        _cache: &'a EvaluationCache<'a>,
    ) -> QueryResult<Self::Prepared<'a>> {
        Ok(())
    }
}

impl<M: IndexDomain, K: GroupKey, I: IndexDomain, V: ValueDomain>
    GroupKernel<M, K, OperandHandle<Indexed<I, V>, Single>> for BroadcastOperation
{
    type Output = OperandHandle<Indexed<M, V>, Multiple<Unordered>>;

    fn execute<'a>(
        _graphrecord: &'a GraphRecord,
        partition: Partition<'a, M, K, OperandHandle<Indexed<I, V>, Single>>,
        _prepared: Self::Prepared<'a>,
    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
        let (buckets, key_failures) = partition.into_parts();

        Ok(Box::new(
            buckets
                .into_iter()
                .flat_map(|(_, members, payload)| {
                    let aggregate = match payload {
                        Ok(Some((_, outcome))) => Some(outcome),
                        Ok(None) => None,
                        Err(failure) => Some(Err(failure)),
                    };

                    members.into_iter().map(move |member| {
                        let outcome = aggregate.clone().unwrap_or_else(|| {
                            Err(Failure::new_at::<M, _>(
                                Self::LABEL,
                                MissingGroupAggregate,
                                &member,
                            ))
                        });

                        (member, outcome)
                    })
                })
                .chain(
                    key_failures
                        .into_iter()
                        .map(|(member, failure)| (member, Err(failure))),
                ),
        ))
    }

    fn estimate(&self, input: Estimate, _stats: &Stats) -> Estimate {
        Estimate {
            distinct: input.elements,
            ..Estimate::UNKNOWN
        }
    }
}

impl<M: IndexDomain, K: GroupKey, V: BareValueDomain>
    GroupKernel<M, K, OperandHandle<Bare<V>, Single>> for BroadcastOperation
{
    type Output = OperandHandle<Indexed<M, V>, Multiple<Unordered>>;

    fn execute<'a>(
        _graphrecord: &'a GraphRecord,
        partition: Partition<'a, M, K, OperandHandle<Bare<V>, Single>>,
        _prepared: Self::Prepared<'a>,
    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
        let (buckets, key_failures) = partition.into_parts();

        Ok(Box::new(
            buckets
                .into_iter()
                .flat_map(|(_, members, payload)| {
                    let aggregate = match payload {
                        Ok(Some(outcome)) => Some(outcome),
                        Ok(None) => None,
                        Err(failure) => Some(Err(failure)),
                    };

                    members.into_iter().map(move |member| {
                        let outcome = aggregate.clone().unwrap_or_else(|| {
                            Err(Failure::new_at::<M, _>(
                                Self::LABEL,
                                MissingGroupAggregate,
                                &member,
                            ))
                        });

                        (member, outcome)
                    })
                })
                .chain(
                    key_failures
                        .into_iter()
                        .map(|(member, failure)| (member, Err(failure))),
                ),
        ))
    }

    fn estimate(&self, input: Estimate, _stats: &Stats) -> Estimate {
        Estimate {
            distinct: input.elements,
            ..Estimate::UNKNOWN
        }
    }
}

impl<M: IndexDomain, K: GroupKey, I: IndexDomain, V: ValueDomain>
    GroupKernel<M, K, OperandHandle<Indexed<I, V>, Definite>> for BroadcastOperation
{
    type Output = OperandHandle<Indexed<M, V>, Multiple<Unordered>>;

    fn execute<'a>(
        _graphrecord: &'a GraphRecord,
        partition: Partition<'a, M, K, OperandHandle<Indexed<I, V>, Definite>>,
        _prepared: Self::Prepared<'a>,
    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
        let (buckets, key_failures) = partition.into_parts();

        Ok(Box::new(
            buckets
                .into_iter()
                .flat_map(|(_, members, payload)| {
                    let aggregate = match payload {
                        Ok((_, outcome)) => outcome,
                        Err(failure) => Err(failure),
                    };

                    members
                        .into_iter()
                        .map(move |member| (member, aggregate.clone()))
                })
                .chain(
                    key_failures
                        .into_iter()
                        .map(|(member, failure)| (member, Err(failure))),
                ),
        ))
    }

    fn estimate(&self, input: Estimate, _stats: &Stats) -> Estimate {
        Estimate {
            distinct: input.elements,
            ..Estimate::UNKNOWN
        }
    }
}

impl<M: IndexDomain, K: GroupKey, V: BareValueDomain>
    GroupKernel<M, K, OperandHandle<Bare<V>, Definite>> for BroadcastOperation
{
    type Output = OperandHandle<Indexed<M, V>, Multiple<Unordered>>;

    fn execute<'a>(
        _graphrecord: &'a GraphRecord,
        partition: Partition<'a, M, K, OperandHandle<Bare<V>, Definite>>,
        _prepared: Self::Prepared<'a>,
    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
        let (buckets, key_failures) = partition.into_parts();

        Ok(Box::new(
            buckets
                .into_iter()
                .flat_map(|(_, members, payload)| {
                    let aggregate = match payload {
                        Ok(outcome) => outcome,
                        Err(failure) => Err(failure),
                    };

                    members
                        .into_iter()
                        .map(move |member| (member, aggregate.clone()))
                })
                .chain(
                    key_failures
                        .into_iter()
                        .map(|(member, failure)| (member, Err(failure))),
                ),
        ))
    }

    fn estimate(&self, input: Estimate, _stats: &Stats) -> Estimate {
        Estimate {
            distinct: input.elements,
            ..Estimate::UNKNOWN
        }
    }
}

impl<O: Apply<BroadcastOperation>> Broadcast for O {
    type ReturnOperand = O::Output;

    fn broadcast(&self) -> Self::ReturnOperand {
        Self::ReturnOperand::new(OperationContext::new(self.clone(), BroadcastOperation))
    }
}

operation_manifest! {
    BroadcastOperation {
        method: Broadcast::broadcast;
        scope: group;

        kernel {
            group: <M: IndexDomain, K: GroupKey>;
            parameters: <I: IndexDomain, V: ValueDomain>;
            input: OperandHandle<Indexed<I, V>, Single>;
            output: OperandHandle<Indexed<M, V>, Multiple<Unordered>>;
        }
        kernel {
            group: <M: IndexDomain, K: GroupKey>;
            parameters: <V: BareValueDomain>;
            input: OperandHandle<Bare<V>, Single>;
            output: OperandHandle<Indexed<M, V>, Multiple<Unordered>>;
        }
        kernel {
            group: <M: IndexDomain, K: GroupKey>;
            parameters: <I: IndexDomain, V: ValueDomain>;
            input: OperandHandle<Indexed<I, V>, Definite>;
            output: OperandHandle<Indexed<M, V>, Multiple<Unordered>>;
        }
        kernel {
            group: <M: IndexDomain, K: GroupKey>;
            parameters: <V: BareValueDomain>;
            input: OperandHandle<Bare<V>, Definite>;
            output: OperandHandle<Indexed<M, V>, Multiple<Unordered>>;
        }
    }
}