tlpt 0.8.0

A set of protocols for building local-first distributed systems
Documentation
use std::collections::HashMap;

use crew_rs::{CredentialSource, Crew, CrewChange, CrewID, MakeStatement, MemberCredential, READER_ROLE, READ_INVITATION_ROLE, WRITER_ROLE, WRITE_INVITATION_ROLE, ADMIN_INVITATION_ROLE, ADMIN_ROLE};
use endr::{BlobDiff, Diff};
use futures::{future, Future, StreamExt, Stream};
use litl::{impl_debug_as_litl, impl_nested_tagged_data_serde, NestedTaggedData, Val};
use ridl::{signing::{SignerID, Signed}, symm_encr::{DecryptionError, KeySecret}};
use thiserror::Error;
use tracing::{debug, error};

use crate::{
    conventions::{named_doc_path, NAMED_DOC_PREFIX, INTRODUCTION_PREFIX},
    invitation::{Invitation, InvitationPrivatePart, InvitationToken},
    node::WeakTlpt,
    Doc, DocID, Tlpt,
};

#[derive(Clone)]
pub struct Team {
    tlpt: Tlpt,
    pub(crate) crew: Crew,
}

impl Team {
    pub fn id(&self) -> TeamID {
        TeamID(self.crew.id())
    }
}

impl Tlpt {
    pub async fn create_team<I: IntoIterator<Item = TeamID>>(&self, parent_teams: I) -> (Team, MemberCredential) {
        let (team, credential) = self.create_team_without_intro(parent_teams).await;

        self.ensure_introduced_in(&team).await;

        (team, credential)
    }

    pub(crate) async fn create_team_without_intro<I: IntoIterator<Item = TeamID>>(&self, parent_teams: I) -> (Team, MemberCredential) {
        let (crew, credential) = self.croo().create_crew_with_parents(parent_teams.into_iter().map(|t| t.0), false).await;

        let team = Team {
            tlpt: self.clone(),
            crew,
        };

        (team, credential)
    }

    pub async fn load_team(&self, team_id: &TeamID) -> Team {
        Team {
            tlpt: self.clone(),
            crew: self.croo().load_crew(team_id.0).await,
        }
    }

    pub async fn load_team_with_credential_source(
        &self,
        team_id: &TeamID,
        credential_source: impl CredentialSource + 'static,
    ) -> Team {
        Team {
            tlpt: self.clone(),
            crew: self
                .croo()
                .load_crew_with_credential_source(team_id.0, credential_source)
                .await,
        }
    }

    pub fn get_loaded_team(&self, team_id: &TeamID) -> Option<Team> {
        self.croo().get_loaded_crew(team_id.0).map(|crew| Team {
            tlpt: self.clone(),
            crew,
        })
    }

    pub async fn join_team(
        &self,
        invitation_token: InvitationToken,
    ) -> Result<(Team, Invitation, InvitationPrivatePart), JoinTeamError> {
        let blob_diffs = self.endr().diffs(
            invitation_token.invitation_id,
            format!("load-{:?}", invitation_token.invitation_id),
        );

        debug!(
            "Loading invitation data {:?}",
            invitation_token.invitation_id
        );

        let invitation_data = blob_diffs
            .filter_map(|diff| match diff {
                Diff::Blob(BlobDiff { data, .. }) => future::ready(data),
                _ => panic!("Unexpected blob diff"),
            })
            .next()
            .await
            .ok_or(JoinTeamError::CouldntLoadInvitation)?;

        debug!("Got invitation data: {:?}", invitation_data);

        let invitation: Invitation = litl::from_val(invitation_data)?;

        let private_part = invitation_token.secret.decrypt(&invitation.private)?;

        let (team_crew, credential) = self.croo().join_crew(&private_part.inner).await;

        let team = Team {
            tlpt: self.clone(),
            crew: team_crew.clone(),
        };

        debug!("About to ensure introduced in");

        let roles = team_crew
            .current_state()
            .unwrap()
            .roles_of(&credential.signer().pub_id());

        if roles.contains("writer") || roles.contains("admin") {
            self.ensure_introduced_in(&team).await;
        }

        debug!("Ensured introduced in");

        Ok((team, invitation, private_part))
    }
}

impl Team {
    // TODO(design): introduce some kind of limited-use per invitation
    pub async fn create_invitation(
        &self,
        kind: &str,
        doc: Option<DocID>,
        public_meta: Option<Val>,
        private_meta: Option<Val>,
        expires_at: ti64::MsSinceEpoch,
    ) -> Result<InvitationToken, String> {
        let secret = KeySecret::new_random();

        let croo_invitation = self.crew.create_invitation(match kind {
            crew_rs::READER_ROLE => READ_INVITATION_ROLE,
            crew_rs::WRITER_ROLE => WRITE_INVITATION_ROLE,
            crew_rs::ADMIN_ROLE => ADMIN_INVITATION_ROLE,
            _ => panic!("Unknown invitation kind"),
        }).await.unwrap();

        let invitation = Invitation {
            public_meta,
            private: secret.encrypt(&InvitationPrivatePart {
                inner: croo_invitation,
                doc,
                private_meta,
            })
        };

        let invitation_id = self.tlpt.endr().create_blob(invitation).await;

        Ok(InvitationToken {
            invitation_id,
            secret,
        })
    }

    pub async fn set_named_doc(&self, name: &str, doc_id: DocID) {
        self.crew
            .make_changes([CrewChange::MakeStatement(MakeStatement {
                path: named_doc_path(name),
                value: litl::to_val(doc_id).unwrap(),
            })])
            .await
            .unwrap();
    }

    pub fn get_named_doc(&self, name: &str) -> Option<DocID> {
        self.crew
            .current_state()?
            .statements
            .iter()
            .find_map(|(path, val, _)| {
                if path.starts_with(&NAMED_DOC_PREFIX) && path.replace(NAMED_DOC_PREFIX, "") == name
                {
                    Some(litl::from_val::<DocID>(val.clone()).unwrap())
                } else {
                    None
                }
            })
    }

    pub async fn wait_for_named_doc(&self, name: &str) -> DocID {
        self.crew
            .wait_for_state(|state| {
                state.statements.iter().any(|(path, val, _)| {
                    path.starts_with(&NAMED_DOC_PREFIX)
                        && path.replace(NAMED_DOC_PREFIX, "") == name
                })
            })
            .await;

        self.get_named_doc(name).unwrap()
    }

    pub fn all_named_docs(&self) -> HashMap<String, DocID> {
        self.crew
            .current_state()
            .unwrap()
            .statements
            .iter()
            .filter_map(|(path, val, _)| {
                if path.starts_with(&NAMED_DOC_PREFIX) {
                    let name = path.replace(NAMED_DOC_PREFIX, "");
                    Some((name, litl::from_val(val.clone()).unwrap()))
                } else {
                    None
                }
            })
            .collect()
    }

    pub fn stream_named_docs(&self) -> impl Stream<Item = HashMap<String, DocID>> {
        self.crew
            .updates("named_docs".to_owned())
            .filter_map(|update| future::ready(update.current_state().cloned()))
            .filter_map(|state| {
                let mut named_docs = HashMap::new();
                for (path, val, _) in state.statements {
                    if path.starts_with(&NAMED_DOC_PREFIX) {
                        let name = path.replace(NAMED_DOC_PREFIX, "");
                        named_docs.insert(name, litl::from_val(val).unwrap());
                    }
                }
                future::ready(if named_docs.is_empty() {
                    None
                } else {
                    Some(named_docs)
                })
            })
    }

    pub async fn wait_for_signer_profile(&self, signer: SignerID) -> Option<Doc> {
        let signed_profile_id = self.crew.updates("find_signer_profile".to_owned()).filter_map(|update|
            future::ready(update.current_state().and_then(|state|
                state.statements.iter().find_map(|(path, val, by)| {
                    if path.starts_with(&INTRODUCTION_PREFIX) && by == &signer {
                        let signed_profile_id = litl::from_val::<Signed<DocID>>(val.clone()).unwrap();

                        Some(signed_profile_id)
                    } else {
                        None
                    }
                })
            ))
        ).next().await?;

        let tlpt = self.tlpt.clone();
        match tlpt.verify_profile(signed_profile_id).await {
            Ok(profile) => Some(profile),
            Err(err) => {
                error!("Couldn't verify profile: {:?}", err);
                None
            }
        }
    }
}

#[derive(Error, Debug)]
pub enum JoinTeamError {
    #[error("Couldn't load invitation")]
    CouldntLoadInvitation,
    #[error("Couldn't deserialize invitation")]
    CouldntDeserializeInvitation(#[from] litl::ValDeserializerError),
    #[error("Couldn't decrypt invitation")]
    CouldntDecryptInvitation(#[from] DecryptionError),
}

#[derive(Copy, Clone, PartialEq, Eq, Hash)]
pub struct TeamID(pub CrewID);

impl_nested_tagged_data_serde!(TeamID);
impl_debug_as_litl!(TeamID);

impl NestedTaggedData for TeamID {
    const TAG: &'static str = "team";

    type Inner = CrewID;

    fn as_inner(&self) -> &Self::Inner {
        &self.0
    }

    fn from_inner(inner: Self::Inner) -> Self
    where
        Self: Sized,
    {
        TeamID(inner)
    }
}