erebos 0.2.2 → 0.2.3
raw patch · 26 files changed
+881/−383 lines, 26 filesPVP: major bump suggested
API removals or changes: PVP suggests a major version bump
API changes (from Hackage documentation)
- Erebos.Chatroom: instance Erebos.Service.Service Erebos.Chatroom.ChatroomService
- Erebos.DirectMessage: instance Erebos.Service.Service Erebos.DirectMessage.DirectMessage
- Erebos.Discovery: instance Erebos.Service.Service Erebos.Discovery.DiscoveryService
- Erebos.Invite: instance Erebos.Service.Service Erebos.Invite.InviteService
- Erebos.Pairing: instance Erebos.Pairing.PairingResult a => Erebos.Service.Service (Erebos.Pairing.PairingService a)
- Erebos.Service: instance Control.Monad.Error.Class.MonadError Erebos.Error.ErebosError (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance Control.Monad.IO.Class.MonadIO (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance Control.Monad.Reader.Class.MonadReader (Erebos.Service.ServiceInput s) (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance Control.Monad.State.Class.MonadState (Erebos.Service.ServiceHandlerState s) (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance Control.Monad.Writer.Class.MonadWriter [Erebos.Service.ServiceReply s] (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance Erebos.State.MonadHead Erebos.State.LocalState (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance Erebos.Storable.Internal.MonadStorage (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance GHC.Base.Applicative (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance GHC.Base.Functor (Erebos.Service.ServiceHandler s)
- Erebos.Service: instance GHC.Base.Monad (Erebos.Service.ServiceHandler s)
- Erebos.State: LocalHeadT :: (Storage -> Stored h -> m (a, Stored h)) -> LocalHeadT h (m :: Type -> Type) a
- Erebos.State: [runLocalHeadT] :: LocalHeadT h (m :: Type -> Type) a -> Storage -> Stored h -> m (a, Stored h)
- Erebos.State: instance Control.Monad.IO.Class.MonadIO m => Erebos.Storable.Internal.MonadStorage (Erebos.State.LocalHeadT h m)
- Erebos.State: newtype LocalHeadT h (m :: Type -> Type) a
- Erebos.Storage.Head: instance Control.Monad.IO.Class.MonadIO m => Erebos.Storable.Internal.MonadStorage (Control.Monad.Trans.Reader.ReaderT (Erebos.Storage.Head.Head a) m)
- Erebos.Storage.Head: instance GHC.Show.Show a => GHC.Show.Show (Erebos.Storage.Head.Head a)
- Erebos.Sync: instance Erebos.Service.Service Erebos.Sync.SyncService
+ Erebos.Chatroom: instance Service.Service Erebos.Chatroom.ChatroomService
+ Erebos.Conversation: makeConversation :: ConversationType conv msg => conv -> Conversation
+ Erebos.DirectMessage: instance Service.Service Erebos.DirectMessage.DirectMessage
+ Erebos.Discovery: [discoverySearchForOwner] :: DiscoveryAttributes -> Bool
+ Erebos.Discovery: instance Erebos.Object.Internal.Storable Erebos.Discovery.DiscoveryVia
+ Erebos.Discovery: instance GHC.Base.Monoid Erebos.Discovery.ResultValue
+ Erebos.Discovery: instance GHC.Base.Semigroup Erebos.Discovery.ResultValue
+ Erebos.Discovery: instance Service.Service Erebos.Discovery.DiscoveryService
+ Erebos.Invite: instance Service.Service Erebos.Invite.InviteService
+ Erebos.Network: getNextPeerChangeChan :: Server -> IO (TChan Peer)
+ Erebos.Network: instance Erebos.Network.Protocol.ProtocolAddressType Erebos.Network.PeerAddress
+ Erebos.Network: isReliableTransport :: PeerAddressType addr => proxy addr -> Bool
+ Erebos.Pairing: instance Erebos.Pairing.PairingResult a => Service.Service (Erebos.Pairing.PairingService a)
+ Erebos.Service: globalStorageWatcher :: (Service s, Eq a) => (Stored LocalState -> a) -> (Server -> a -> ExceptT ErebosError IO ()) -> SomeStorageWatcher s
+ Erebos.Service: globalStorageWatcherH :: (Service s, Eq a) => (Head LocalState -> a) -> (Server -> a -> ExceptT ErebosError IO ()) -> SomeStorageWatcher s
+ Erebos.Service: someStorageWatcher :: (Service s, Eq a) => (Stored LocalState -> a) -> (a -> ServiceHandler s ()) -> SomeStorageWatcher s
+ Erebos.Service: someStorageWatcherHC :: (Service s, Eq a) => (Stored LocalState -> HeadCacheType LocalState -> a) -> (a -> ServiceHandler s ()) -> SomeStorageWatcher s
+ Erebos.State: ($dmmloadKeyMb) :: (MonadStorage m, KeyPair sec pub, MonadIO m) => Stored pub -> m (Maybe sec)
+ Erebos.State: ($dmmstoreKey) :: (MonadStorage m, KeyPair sec pub, MonadIO m) => sec -> m ()
+ Erebos.State: data LocalHeadT h (m :: Type -> Type) a
+ Erebos.State: getLocalHeadCache :: MonadHead a m => proxy a -> m (HeadCacheType a)
+ Erebos.State: instance Control.Monad.IO.Class.MonadIO m => Erebos.Storage.Monad.MonadStorage (Erebos.State.LocalHeadT h m)
+ Erebos.State: lookupSharedValueH :: SharedType a => Head LocalState -> a
+ Erebos.State: lookupSharedValueHC :: SharedType a => Stored LocalState -> HeadCacheType LocalState -> a
+ Erebos.State: lookupSharedValueM :: (SharedType a, MonadHead LocalState m) => m a
+ Erebos.State: mloadKeyMb :: (MonadStorage m, KeyPair sec pub) => Stored pub -> m (Maybe sec)
+ Erebos.State: mstoreKey :: (MonadStorage m, KeyPair sec pub) => sec -> m ()
+ Erebos.State: runLocalHead :: HeadType h => LocalHead h a -> Storage -> Stored h -> (a, Stored h)
+ Erebos.State: runLocalHead' :: HeadType h => LocalHead h a -> Head h -> (a, Stored h)
+ Erebos.State: runLocalHeadT :: (HeadType h, Monad m) => LocalHeadT h m a -> Storage -> Stored h -> m (a, Stored h)
+ Erebos.State: runLocalHeadT' :: (HeadType h, Monad m) => LocalHeadT h m a -> Head h -> m (a, Stored h)
+ Erebos.State: type LocalHead h a = LocalHeadT h Identity a
+ Erebos.Storable: ($dmmloadKeyMb) :: (MonadStorage m, KeyPair sec pub, MonadIO m) => Stored pub -> m (Maybe sec)
+ Erebos.Storable: ($dmmstoreKey) :: (MonadStorage m, KeyPair sec pub, MonadIO m) => sec -> m ()
+ Erebos.Storable: mloadKey :: (KeyPair sec pub, MonadStorage m, MonadError e m, FromErebosError e) => Stored pub -> m sec
+ Erebos.Storable: mloadKeyMb :: (MonadStorage m, KeyPair sec pub) => Stored pub -> m (Maybe sec)
+ Erebos.Storable: mstoreKey :: (MonadStorage m, KeyPair sec pub) => sec -> m ()
+ Erebos.Storage: ($dmmloadKeyMb) :: (MonadStorage m, KeyPair sec pub, MonadIO m) => Stored pub -> m (Maybe sec)
+ Erebos.Storage: ($dmmstoreKey) :: (MonadStorage m, KeyPair sec pub, MonadIO m) => sec -> m ()
+ Erebos.Storage: mloadKey :: (KeyPair sec pub, MonadStorage m, MonadError e m, FromErebosError e) => Stored pub -> m sec
+ Erebos.Storage: mloadKeyMb :: (MonadStorage m, KeyPair sec pub) => Stored pub -> m (Maybe sec)
+ Erebos.Storage: mstoreKey :: (MonadStorage m, KeyPair sec pub) => sec -> m ()
+ Erebos.Storage: updateHead' :: (HeadType a, MonadIO m) => Head a -> (Head a -> m (Stored a, b)) -> m (Maybe (Head a), b)
+ Erebos.Storage.Graph: type StoredTips a = [Stored a]
+ Erebos.Storage.Head: ($dmheadCacheInit) :: (HeadType a, HeadCacheType a ~ ()) => Stored a -> HeadCacheType a
+ Erebos.Storage.Head: ($dmheadCacheUpdate) :: (HeadType a, HeadCacheType a ~ ()) => Stored a -> HeadCacheType a -> HeadCacheType a
+ Erebos.Storage.Head: headCache :: Head a -> HeadCacheType a
+ Erebos.Storage.Head: headCacheInit :: HeadType a => Stored a -> HeadCacheType a
+ Erebos.Storage.Head: headCacheUpdate :: HeadType a => Stored a -> HeadCacheType a -> HeadCacheType a
+ Erebos.Storage.Head: instance Control.Monad.IO.Class.MonadIO m => Erebos.Storage.Monad.MonadStorage (Control.Monad.Trans.Reader.ReaderT (Erebos.Storage.Head.Head a) m)
+ Erebos.Storage.Head: instance GHC.Show.Show (Erebos.Storage.Head.Head a)
+ Erebos.Storage.Head: type HeadCacheType a = ();
+ Erebos.Storage.Head: updateHead' :: (HeadType a, MonadIO m) => Head a -> (Head a -> m (Stored a, b)) -> m (Maybe (Head a), b)
+ Erebos.Storage.Head: }
+ Erebos.Storage.Merge: type StoredTips a = [Stored a]
+ Erebos.Sync: instance Service.Service Erebos.Sync.SyncService
- Erebos.Chatroom: chatroomSetSubscribe :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Bool -> m ()
+ Erebos.Chatroom: chatroomSetSubscribe :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Bool -> m ()
- Erebos.Chatroom: deleteChatroomByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m ()
+ Erebos.Chatroom: deleteChatroomByStateData :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m ()
- Erebos.Chatroom: joinChatroom :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => ChatroomState -> m ()
+ Erebos.Chatroom: joinChatroom :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => ChatroomState -> m ()
- Erebos.Chatroom: joinChatroomAs :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => UnifiedIdentity -> ChatroomState -> m ()
+ Erebos.Chatroom: joinChatroomAs :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => UnifiedIdentity -> ChatroomState -> m ()
- Erebos.Chatroom: joinChatroomAsByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => UnifiedIdentity -> Stored ChatroomStateData -> m ()
+ Erebos.Chatroom: joinChatroomAsByStateData :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => UnifiedIdentity -> Stored ChatroomStateData -> m ()
- Erebos.Chatroom: joinChatroomByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m ()
+ Erebos.Chatroom: joinChatroomByStateData :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m ()
- Erebos.Chatroom: leaveChatroom :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => ChatroomState -> m ()
+ Erebos.Chatroom: leaveChatroom :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => ChatroomState -> m ()
- Erebos.Chatroom: leaveChatroomByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m ()
+ Erebos.Chatroom: leaveChatroomByStateData :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m ()
- Erebos.Chatroom: sendChatroomMessage :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => ChatroomState -> Text -> m ()
+ Erebos.Chatroom: sendChatroomMessage :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => ChatroomState -> Text -> m ()
- Erebos.Chatroom: sendChatroomMessageByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Text -> m ()
+ Erebos.Chatroom: sendChatroomMessageByStateData :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Text -> m ()
- Erebos.Chatroom: updateChatroomByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Maybe Text -> Maybe Text -> m (Maybe ChatroomState)
+ Erebos.Chatroom: updateChatroomByStateData :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Maybe Text -> Maybe Text -> m (Maybe ChatroomState)
- Erebos.Conversation: sendMessage :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Conversation -> Text -> m ()
+ Erebos.Conversation: sendMessage :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Conversation -> Text -> m ()
- Erebos.DirectMessage: sendDirectMessage :: forall (f :: Type -> Type) m. (Foldable f, Applicative f, MonadHead LocalState m) => Identity f -> Text -> m ()
+ Erebos.DirectMessage: sendDirectMessage :: forall (f :: Type -> Type) m. (Foldable f, Applicative f, MonadHead LocalState m, MonadIO m) => Identity f -> Text -> m ()
- Erebos.Discovery: DiscoveryAttributes :: Maybe Word16 -> Maybe Text -> Maybe Word16 -> Maybe Text -> (Peer -> PeerAddress -> Bool) -> Bool -> DiscoveryAttributes
+ Erebos.Discovery: DiscoveryAttributes :: Maybe Word16 -> Maybe Text -> Maybe Word16 -> Maybe Text -> (Peer -> PeerAddress -> Bool) -> Bool -> Bool -> DiscoveryAttributes
- Erebos.Discovery: DiscoveryResult :: Either Ref RefDigest -> [DiscoveryAddress] -> DiscoveryService
+ Erebos.Discovery: DiscoveryResult :: Either Ref RefDigest -> [DiscoveryAddress] -> [DiscoveryVia] -> DiscoveryService
- Erebos.Identity: mergeIdentity :: forall m e (f :: Type -> Type). (MonadStorage m, MonadError e m, FromErebosError e, MonadIO m) => Identity f -> m UnifiedIdentity
+ Erebos.Identity: mergeIdentity :: forall m e (f :: Type -> Type). (MonadStorage m, MonadError e m, FromErebosError e) => Identity f -> m UnifiedIdentity
- Erebos.Invite: createSingleContactInvite :: MonadHead LocalState m => Text -> m Invite
+ Erebos.Invite: createSingleContactInvite :: (MonadHead LocalState m, MonadIO m) => Text -> m Invite
- Erebos.Service: runServiceHandler :: Service s => Head LocalState -> ServiceInput s -> ServiceState s -> ServiceGlobalState s -> ServiceHandler s () -> IO ([ServiceReply s], (ServiceState s, ServiceGlobalState s))
+ Erebos.Service: runServiceHandler :: Service s => Head LocalState -> ServiceInput s -> ServiceState s -> ServiceGlobalState s -> ServiceHandler s () -> IO ([ServiceReply s], (ServiceState s, ServiceGlobalState s, Head LocalState))
- Erebos.State: class (MonadIO m, MonadStorage m) => MonadHead a (m :: Type -> Type)
+ Erebos.State: class (MonadStorage m, HeadType a) => MonadHead a (m :: Type -> Type)
- Erebos.Storage.Head: class Storable a => HeadType a
+ Erebos.Storage.Head: class Storable a => HeadType a where {
Files
- CHANGELOG.md +8/−0
- erebos.cabal +3/−1
- main/Main.hs +14/−6
- main/State.hs +1/−1
- main/Test.hs +13/−6
- main/WebSocket.hs +2/−0
- src/Erebos/Chatroom.hs +20/−21
- src/Erebos/Conversation.hs +12/−3
- src/Erebos/DirectMessage.hs +8/−8
- src/Erebos/Discovery.hs +242/−59
- src/Erebos/Identity.hs +2/−2
- src/Erebos/Invite.hs +7/−8
- src/Erebos/Network.hs +59/−23
- src/Erebos/Network/Protocol.hs +22/−14
- src/Erebos/Object/Internal.hs +6/−3
- src/Erebos/Service.hs +11/−170
- src/Erebos/State.hs +85/−14
- src/Erebos/Storable.hs +2/−1
- src/Erebos/Storable/Internal.hs +5/−18
- src/Erebos/Storage.hs +3/−3
- src/Erebos/Storage/Graph.hs +1/−0
- src/Erebos/Storage/Head.hs +98/−21
- src/Erebos/Storage/Merge.hs +1/−0
- src/Erebos/Storage/Monad.hs +37/−0
- src/Erebos/Sync.hs +1/−1
- src/Service.hs +218/−0
CHANGELOG.md view
@@ -1,5 +1,13 @@ # Revision history for erebos +## 0.2.3 -- 2026-08-19++* Discovery service handles situation when search target joins later.+* Automatically search for other devices of the same owner when using the discovery service.+* Added API to mark a custom underlying transport layer as reliable.+* Fixed thread diff on first callback from `watchDirectMessageThreads` for given peer.+* Fixed discovery when forced to bind not-selected ephemeral port.+ ## 0.2.2 -- 2026-06-25 * Initial support for `chunked` and `dir` object types.
erebos.cabal view
@@ -1,7 +1,7 @@ Cabal-Version: 3.0 Name: erebos-Version: 0.2.2+Version: 0.2.3 Synopsis: Decentralized messaging and synchronization Description: Library and simple CLI interface implementing the Erebos identity@@ -128,10 +128,12 @@ Erebos.Storage.Disk Erebos.Storage.Internal Erebos.Storage.Memory+ Erebos.Storage.Monad Erebos.Storage.Platform Erebos.TextFormat.Types Erebos.UUID Erebos.Util+ Service c-sources: src/Erebos/Network/ifaddrs.c
main/Main.hs view
@@ -51,6 +51,7 @@ import Erebos.State import Erebos.Storable import Erebos.Storage+import Erebos.Storage.Head import Erebos.Storage.Merge import Erebos.Sync import Erebos.TextFormat@@ -467,7 +468,7 @@ chatroomSetVar <- liftIO $ newEmptyMVar let autoSubscribe = optChatroomAutoSubscribe opts- chatroomList = fromSetBy (comparing roomStateData) . lookupSharedValue . lsShared . headObject $ erebosHead+ chatroomList = fromSetBy (comparing roomStateData) . lookupSharedValueH $ erebosHead watched <- if isJust autoSubscribe || any roomStateSubscribe chatroomList then do fmap Just $ liftIO $ watchChatroomsForCli tui extPrintLn erebosHead@@ -628,11 +629,18 @@ instance MonadHead LocalState CommandM where updateLocalHead f = do- h <- gets csHead- (Just h', x) <- maybe (fail "failed to reload head") (flip updateHead f) =<< reloadHead h+ h <- maybe (fail "failed to reload head") return =<< reloadHead =<< gets csHead+ ( Just h', x ) <- updateHead' h $ \h' -> do+ modify $ \s -> s { csHead = h' }+ f (headStoredObject h') modify $ \s -> s { csHead = h' } return x + getLocalHeadCache _ = do+ h <- maybe (fail "failed to reload head") return =<< reloadHead =<< gets csHead+ modify $ \s -> s { csHead = h }+ return $ headCache h+ type Command = CommandM () getSelectedPeer :: CommandM Peer@@ -1040,7 +1048,7 @@ cmdContacts = do args <- words <$> asks ciLine ehead <- gets csHead- let contacts = fromSetBy (comparing contactName) $ lookupSharedValue $ lsShared $ headObject ehead+ let contacts = fromSetBy (comparing contactName) $ lookupSharedValueH ehead verbose = "-v" `elem` args set <- asks ciSetContextOptions set WatchContacts $ map SelectedContact contacts@@ -1073,7 +1081,7 @@ _ -> liftIO $ do setPrompt term "Name: " getInputLine term $ KeepPrompt . maybe T.empty T.pack- (lookupSharedValue . lsShared . fromStored <$> getLocalHead) >>= \case+ lookupSharedValueM >>= \case Just (self :: ComposedIdentity) -> do invite <- createSingleContactInvite name dgst : _ <- return $ refDigest . storedRef <$> idDataF self@@ -1161,7 +1169,7 @@ NoContext -> cmdPutStrLn "nothing selected" where printContactOrIdentityDetails cid = do- contacts <- fromSetBy (comparing contactName) . lookupSharedValue . lsShared . fromStored <$> getLocalHead+ contacts <- fromSetBy (comparing contactName) <$> lookupSharedValueM case find (maybe False (sameIdentity cid) . contactIdentity) contacts of Just contact -> printContactDetails contact Nothing -> printIdentityDetails cid
main/State.hs view
@@ -89,7 +89,7 @@ createSingleIdentity owner name = createIdentity name (Just owner) -updateSharedIdentity :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Terminal -> m ()+updateSharedIdentity :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Terminal -> m () updateSharedIdentity term = updateLocalState_ $ updateSharedState_ $ \case Just identity -> do Just . toComposedIdentity <$> interactiveIdentityUpdate term identity
main/Test.hs view
@@ -282,11 +282,18 @@ instance MonadHead LocalState CommandM where updateLocalHead f = do- Just h <- gets tsHead- (Just h', x) <- maybe (fail "failed to reload head") (flip updateHead f) =<< reloadHead h+ h <- maybe (fail "failed to reload head") return =<< reloadHead =<< getOrLoadHead+ ( Just h', x ) <- updateHead' h $ \h' -> do+ modify $ \s -> s { tsHead = Just h' }+ f (headStoredObject h') modify $ \s -> s { tsHead = Just h' } return x + getLocalHeadCache _ = do+ h <- maybe (fail "failed to reload head") return =<< reloadHead =<< getOrLoadHead+ modify $ \s -> s { tsHead = Just h }+ return $ headCache h+ type Command = CommandM () commands :: [ ( Text, Command ) ]@@ -847,7 +854,7 @@ Nothing <- gets tsWatchedSharedIdentity out <- asks tiOutput- w <- liftIO $ watchHeadWith h (lookupSharedValue . lsShared . headObject) $ \case+ w <- liftIO $ watchHeadWith h lookupSharedValueH $ \case Just (idt :: ComposedIdentity) -> do outLine out $ unwords $ "shared-identity" : map (maybe "<unnamed>" T.unpack . idName) (unfoldOwners idt) Nothing -> do@@ -917,7 +924,7 @@ cmdContactList :: Command cmdContactList = do h <- getHead- let contacts = fromSetBy (comparing contactName) . lookupSharedValue . lsShared . headObject $ h+ let contacts = fromSetBy (comparing contactName) . lookupSharedValueH $ h forM_ contacts $ \c -> do r:_ <- return $ filterAncestors $ concatMap storedRoots $ toComponents c cmdOut $ concat@@ -932,7 +939,7 @@ getContact :: Text -> CommandM Contact getContact cid = do h <- getHead- let contacts = fromSetBy (comparing contactName) . lookupSharedValue . lsShared . headObject $ h+ let contacts = fromSetBy (comparing contactName) . lookupSharedValueH $ h [contact] <- flip filterM contacts $ \c -> do r:_ <- return $ filterAncestors $ concatMap storedRoots $ toComponents c return $ T.pack (show $ refDigest $ storedRef r) == cid@@ -973,7 +980,7 @@ dmList :: Foldable f => Identity f -> Command dmList peer = do- threads <- dmThreadList . lookupSharedValue . lsShared . headObject <$> getHead+ threads <- dmThreadList . lookupSharedValueH <$> getHead case find (sameIdentity peer . msgPeer) threads of Just thread -> do forM_ (reverse $ dmThreadToListUnread thread) $ \( DirectMessage {..}, new ) -> cmdOut $ "dm-list-item"
main/WebSocket.hs view
@@ -39,6 +39,8 @@ | Just WS.ConnectionClosed <- fromException e -> return () | otherwise -> throwIO e + isReliableTransport _ = True+ data WebSocketOptions = WebSocketOptions { wsAddress :: String
src/Erebos/Chatroom.hs view
@@ -201,17 +201,17 @@ cmpView msg = (zonedTimeToUTC $ mdTime $ fromSigned msg, msg) sendChatroomMessage- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => ChatroomState -> Text -> m () sendChatroomMessage rstate msg = sendChatroomMessageByStateData (head $ roomStateData rstate) msg sendChatroomMessageByStateData- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Text -> m () sendChatroomMessageByStateData lookupData msg = sendRawChatroomMessageByStateData lookupData Nothing Nothing (Just msg) False sendRawChatroomMessageByStateData- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Maybe UnifiedIdentity -> Maybe (Stored (Signed ChatMessageData)) -> Maybe Text -> Bool -> m () sendRawChatroomMessageByStateData lookupData mbIdentity mdReplyTo mdText mdLeave = void $ findAndUpdateChatroomState $ \cstate -> do guard $ any (lookupData `precedesOrEquals`) $ roomStateData cstate@@ -220,7 +220,7 @@ | Just identity <- mbIdentity -> return identity | Just identity <- roomStateIdentity cstate -> return identity | otherwise -> localIdentity . fromStored <$> getLocalHead- secret <- loadKey $ idKeyMessage mdFrom+ secret <- mloadKey $ idKeyMessage mdFrom mdTime <- liftIO getZonedTime let mdPrev = roomStateMessageData cstate mdRoom = if null (roomStateMessageData cstate)@@ -317,7 +317,7 @@ (, cstate) <$> storeSetAdd cstate rooms findAndUpdateChatroomState- :: (MonadStorage m, MonadHead LocalState m)+ :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => (ChatroomState -> Maybe (m ChatroomState)) -> m (Maybe ChatroomState) findAndUpdateChatroomState f = do@@ -335,7 +335,7 @@ [] -> return (roomSet, Nothing) deleteChatroomByStateData- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m () deleteChatroomByStateData lookupData = void $ findAndUpdateChatroomState $ \cstate -> do guard $ any (lookupData `precedesOrEquals`) $ roomStateData cstate@@ -346,7 +346,7 @@ } updateChatroomByStateData- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Maybe Text -> Maybe Text@@ -355,7 +355,7 @@ guard $ any (lookupData `precedesOrEquals`) $ roomStateData cstate room <- roomStateRoom cstate Just $ do- secret <- loadKey $ roomKey room+ secret <- mloadKey $ roomKey room rdata <- mstore =<< sign secret =<< mstore ChatroomData { rdPrev = roomData room , rdName = newName@@ -371,12 +371,11 @@ listChatrooms :: MonadHead LocalState m => m [ChatroomState] listChatrooms = filter (not . roomStateDeleted) .- fromSetBy (comparing $ roomName <=< roomStateRoom) .- lookupSharedValue . lsShared . fromStored <$> getLocalHead+ fromSetBy (comparing $ roomName <=< roomStateRoom) <$> lookupSharedValueM findChatroom :: MonadHead LocalState m => (ChatroomState -> Bool) -> m (Maybe ChatroomState) findChatroom p = do- list <- map snd . chatroomSetToList . lookupSharedValue . lsShared . fromStored <$> getLocalHead+ list <- map snd . chatroomSetToList <$> lookupSharedValueM return $ find p list findChatroomByRoomData :: MonadHead LocalState m => Stored (Signed ChatroomData) -> m (Maybe ChatroomState)@@ -387,7 +386,7 @@ findChatroomByStateData cdata = findChatroom $ any (cdata `precedesOrEquals`) . roomStateData chatroomSetSubscribe- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> Bool -> m () chatroomSetSubscribe lookupData subscribe = void $ findAndUpdateChatroomState $ \cstate -> do guard $ any (lookupData `precedesOrEquals`) $ roomStateData cstate@@ -408,32 +407,32 @@ toList $ ancestors $ roomStateMessageData joinChatroom- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => ChatroomState -> m () joinChatroom rstate = joinChatroomByStateData (head $ roomStateData rstate) joinChatroomByStateData- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m () joinChatroomByStateData lookupData = sendRawChatroomMessageByStateData lookupData Nothing Nothing Nothing False joinChatroomAs- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => UnifiedIdentity -> ChatroomState -> m () joinChatroomAs identity rstate = joinChatroomAsByStateData identity (head $ roomStateData rstate) joinChatroomAsByStateData- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => UnifiedIdentity -> Stored ChatroomStateData -> m () joinChatroomAsByStateData identity lookupData = sendRawChatroomMessageByStateData lookupData (Just identity) Nothing Nothing False leaveChatroom- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => ChatroomState -> m () leaveChatroom rstate = leaveChatroomByStateData (head $ roomStateData rstate) leaveChatroomByStateData- :: (MonadStorage m, MonadHead LocalState m, MonadError e m, FromErebosError e)+ :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Stored ChatroomStateData -> m () leaveChatroomByStateData lookupData = sendRawChatroomMessageByStateData lookupData Nothing Nothing Nothing True @@ -453,7 +452,7 @@ watchChatrooms :: MonadIO m => Head LocalState -> (Set ChatroomState -> Maybe [ChatroomSetChange] -> IO ()) -> m WatchedHead watchChatrooms h f = liftIO $ do lastVar <- newIORef Nothing- watchHeadWith h (lookupSharedValue . lsShared . headObject) $ \cur -> do+ watchHeadWith h lookupSharedValueH $ \cur -> do let curList = chatroomSetToList cur mbLast <- readIORef lastVar writeIORef lastVar $ Just curList@@ -537,7 +536,7 @@ svcModify $ \s -> s { psSendRoomUpdates = True } when (not previouslyUpdated) $ do- syncChatroomsToPeer . lookupSharedValue . lsShared . fromStored =<< getLocalHead+ syncChatroomsToPeer =<< lookupSharedValueM when chatRoomQuery $ do rooms <- listChatrooms@@ -615,7 +614,7 @@ replyPacket emptyPacket { chatRoomQuery = True } serviceStorageWatchers _ = (:[]) $- SomeStorageWatcher (lookupSharedValue . lsShared . fromStored) syncChatroomsToPeer+ someStorageWatcher (lookupSharedValue . lsShared . fromStored) syncChatroomsToPeer syncChatroomsToPeer :: Set ChatroomState -> ServiceHandler ChatroomService () syncChatroomsToPeer set = do
src/Erebos/Conversation.hs view
@@ -11,6 +11,7 @@ formatMessageFT, Conversation,+ makeConversation, isSameConversation, directMessageConversation, chatroomConversation,@@ -32,6 +33,7 @@ ) where import Control.Monad.Except+import Control.Monad.IO.Class import Data.List import Data.Maybe@@ -39,6 +41,7 @@ import Data.Text qualified as T import Data.Time.Format import Data.Time.LocalTime+import Data.Typeable import Erebos.Chatroom import Erebos.Conversation.Class@@ -89,6 +92,12 @@ = DirectMessageConversation DirectMessageThread | ChatroomConversation ChatroomState +makeConversation :: ConversationType conv msg => conv -> Conversation+makeConversation conv+ | Just x <- cast conv = DirectMessageConversation x+ | Just x <- cast conv = ChatroomConversation x+ | otherwise = error "unhandled conversation type"+ withConversation :: (forall conv msg. ConversationType conv msg => conv -> a) -> Conversation -> a withConversation f (DirectMessageConversation conv) = f conv withConversation f (ChatroomConversation conv) = f conv@@ -107,7 +116,7 @@ directMessageConversation :: MonadHead LocalState m => ComposedIdentity -> m Conversation directMessageConversation peer = do createOrUpdateDirectMessagePeer peer- (find (sameIdentity peer . msgPeer) . dmThreadList . lookupSharedValue . lsShared . fromStored <$> getLocalHead) >>= \case+ (find (sameIdentity peer . msgPeer) . dmThreadList <$> lookupSharedValueM) >>= \case Just thread -> return $ DirectMessageConversation thread Nothing -> return $ DirectMessageConversation $ dmEmptyThread peer @@ -127,7 +136,7 @@ fromMaybe cur <$> chatroomConversation rstate lookupConversations :: MonadHead LocalState m => m [ Conversation ]-lookupConversations = map DirectMessageConversation . dmThreadList . lookupSharedValue . lsShared . fromStored <$> getLocalHead+lookupConversations = map DirectMessageConversation . dmThreadList <$> lookupSharedValueM lookupConversationByRef :: MonadHead LocalState m => RefDigest -> m (Maybe Conversation) lookupConversationByRef dgst = find ((dgst ==) . conversationReference) <$> lookupConversations@@ -151,7 +160,7 @@ conversationHistoryChange = withConversations $ \since -> fmap (map (uncurry Message)) . convMessageListSince since -sendMessage :: (MonadHead LocalState m, MonadError e m, FromErebosError e) => Conversation -> Text -> m ()+sendMessage :: (MonadHead LocalState m, MonadIO m, MonadError e m, FromErebosError e) => Conversation -> Text -> m () sendMessage (DirectMessageConversation thread) text = sendDirectMessage (msgPeer thread) text sendMessage (ChatroomConversation rstate) text = sendChatroomMessage rstate text
src/Erebos/DirectMessage.hs view
@@ -120,8 +120,8 @@ let msg = fromStored smsg powner <- asks $ finalOwner . svcPeerIdentity erb <- svcGetLocal- let DirectMessageThreads prev _ = lookupSharedValue $ lsShared $ fromStored erb- sent = concat $ propertyValue $ findMsgProperty powner msSent prev+ DirectMessageThreads prev _ <- lookupSharedValueM+ let sent = concat $ propertyValue $ findMsgProperty powner msSent prev received = concat $ propertyValue $ findMsgProperty powner msReceived prev received' = filterAncestors $ smsg : received if powner `sameIdentity` msgFrom msg ||@@ -146,14 +146,14 @@ else join $ asks $ dmOwnerMismatch . svcAttributes serviceNewPeer = do- syncDirectMessageToPeer . lookupSharedValue . lsShared . fromStored =<< svcGetLocal+ syncDirectMessageToPeer =<< lookupSharedValueM serviceUpdatedPeer = do updateDirectMessagePeer . finalOwner =<< asks svcPeerIdentity serviceStorageWatchers _ =- [ SomeStorageWatcher (lookupSharedValue . lsShared . fromStored) syncDirectMessageToPeer- , GlobalStorageWatcher (lookupSharedValue . lsShared . fromStored) findMissingPeers+ [ someStorageWatcherHC lookupSharedValueHC syncDirectMessageToPeer+ , globalStorageWatcherH lookupSharedValueH findMissingPeers ] @@ -213,7 +213,7 @@ findMsgPropertyUpdate prev mss = findPropertyUpdate prev mss -sendDirectMessage :: (Foldable f, Applicative f, MonadHead LocalState m)+sendDirectMessage :: (Foldable f, Applicative f, MonadHead LocalState m, MonadIO m) => Identity f -> Text -> m () sendDirectMessage pid text = updateLocalState_ $ \ls -> do let self = localIdentity $ fromStored ls@@ -424,7 +424,7 @@ watchDirectMessageThreads :: Head LocalState -> (DirectMessageThread -> DirectMessageThread -> IO ()) -> IO WatchedHead watchDirectMessageThreads h callback = do prevVar <- newMVar Nothing- watchHeadWith h (lookupSharedValue . lsShared . headObject) $ \(DirectMessageThreads sms _) -> do+ watchHeadWith h lookupSharedValueH $ \(DirectMessageThreads sms _) -> do modifyMVar_ prevVar $ \case Just ( prev, prevPeers ) -> do let addPeer (p : ps) p'@@ -437,7 +437,7 @@ | p `sameIdentity` p' = let x' = f x in ( ( x, x' ), ( p', x' ) : ps ) | otherwise = (px :) <$> updatePeer ps p' f updatePeer [] p' f =- let x = dmEmptyThread p'; x' = f x+ let x = messageThreadFor (dmEmptyThread p') prev; x' = f x in ( ( x, x' ), [ ( p', x' ) ] ) peers <- (\f -> foldM f prevPeers changedPeers) $ \peers peer -> do
src/Erebos/Discovery.hs view
@@ -26,6 +26,8 @@ import Data.Text qualified as T import Data.Word +import System.Clock+ import Text.Read #ifdef ENABLE_ICE_SUPPORT@@ -37,6 +39,7 @@ import Erebos.Object import Erebos.Service import Erebos.Service.Stream+import Erebos.State import Erebos.Storable @@ -51,7 +54,7 @@ = DiscoverySelf [ DiscoveryAddress ] (Maybe Int) | DiscoveryAcknowledged [ DiscoveryAddress ] (Maybe Text) (Maybe Word16) (Maybe Text) (Maybe Word16) | DiscoverySearch (Either Ref RefDigest)- | DiscoveryResult (Either Ref RefDigest) [ DiscoveryAddress ]+ | DiscoveryResult (Either Ref RefDigest) [ DiscoveryAddress ] [ DiscoveryVia ] | DiscoveryConnectionRequest DiscoveryConnection | DiscoveryConnectionResponse DiscoveryConnection @@ -61,6 +64,11 @@ | DiscoveryTunnel | DiscoveryOther Text +data DiscoveryVia = DiscoveryVia+ { viaIdentity :: RefDigest+ , viaAddress :: [ DiscoveryAddress ]+ }+ data DiscoveryAttributes = DiscoveryAttributes { discoveryStunPort :: Maybe Word16 , discoveryStunServer :: Maybe Text@@ -68,6 +76,7 @@ , discoveryTurnServer :: Maybe Text , discoveryProvideTunnel :: Peer -> PeerAddress -> Bool , discoveryDebugLog :: Bool+ , discoverySearchForOwner :: Bool } defaultDiscoveryAttributes :: DiscoveryAttributes@@ -78,6 +87,7 @@ , discoveryTurnServer = Nothing , discoveryProvideTunnel = \_ _ -> False , discoveryDebugLog = False+ , discoverySearchForOwner = True } data DiscoveryConnection = DiscoveryConnection@@ -109,9 +119,10 @@ storeMbText "turn-server" turnServer storeMbInt "turn-port" turnPort DiscoverySearch edgst -> either (storeRawRef "search") (storeRawWeak "search") edgst- DiscoveryResult edgst addr -> do+ DiscoveryResult edgst addr via -> do either (storeRawRef "result") (storeRawWeak "result") edgst mapM_ (storeText "address") addr+ mapM_ (storeRef "via") via DiscoveryConnectionRequest conn -> storeConnection "request" conn DiscoveryConnectionResponse conn -> storeConnection "response" conn @@ -150,6 +161,7 @@ , Right <$> loadRawWeak "result" ] <*> loadTexts "address"+ <*> loadRefs "via" , loadConnection "request" DiscoveryConnectionRequest , loadConnection "response" DiscoveryConnectionResponse ]@@ -192,7 +204,16 @@ | otherwise -> DiscoveryOther str +instance Storable DiscoveryVia where+ store' DiscoveryVia {..} = storeRec $ do+ storeRawWeak "id" viaIdentity+ mapM_ (storeText "address") viaAddress+ load' = loadRec $ do+ viaIdentity <- loadRawWeak "id"+ viaAddress <- loadTexts "address"+ return DiscoveryVia {..} + data DiscoveryPeer = DiscoveryPeer { dpPriority :: Int , dpPeer :: Maybe Peer@@ -208,8 +229,27 @@ , dpIceSession = Nothing } +viaFromPeer :: DiscoveryAttributes -> Peer -> DiscoveryPeer -> IO (Maybe DiscoveryVia)+viaFromPeer attrs peer DiscoveryPeer {..}+ | Just dpeer <- dpPeer+ , _ : _ <- dpAddress+ = do+ getPeerIdentity dpeer >>= \case+ PeerIdentityFull pid -> do+ offerTunnel <- offerTunnelBetween attrs peer dpeer+ let viaAddress = (if offerTunnel then (++ [ DiscoveryTunnel ]) else id) dpAddress+ let viaIdentity = refDigest $ storedRef $ idData pid+ return $ Just DiscoveryVia {..}+ _ -> return Nothing+ | otherwise+ = return Nothing+++ data DiscoveryPeerState = DiscoveryPeerState- { dpsOurTunnelRequests :: [ ( RefDigest, StreamWriter ) ]+ { dpsWeAskedFor :: Map RefDigest SearchStatus+ , dpsPeerSearchingFor :: Map RefDigest SearchStatus+ , dpsOurTunnelRequests :: [ ( RefDigest, StreamWriter ) ] -- ( original target, our write stream ) , dpsRelayedTunnelRequests :: [ ( RefDigest, ( StreamReader, StreamWriter )) ] -- ( original source, ( from source, to target ))@@ -219,10 +259,31 @@ } data DiscoveryGlobalState = DiscoveryGlobalState- { dgsPeers :: Map RefDigest DiscoveryPeer+ { dgsPeers :: Map RefDigest ResultValue , dgsSearchingFor :: Set RefDigest } +data ResultValue = ResultValue+ { rvDirect :: Maybe DiscoveryPeer+ , rvVia :: [ DiscoveryPeer ]+ }++instance Semigroup ResultValue where+ new <> old = ResultValue+ { rvDirect = if (dpPriority <$> rvDirect new) > (dpPriority <$> rvDirect old)+ then rvDirect new else rvDirect old+ , rvVia = rvVia new ++ rvVia old+ }++instance Monoid ResultValue where+ mempty = ResultValue+ { rvDirect = Nothing+ , rvVia = []+ }++data SearchStatus+ = SearchingSince TimeSpec+ instance Service DiscoveryService where serviceID _ = mkServiceID "dd59c89c-69cc-4703-b75b-4ddcd4b3c23c" @@ -231,7 +292,9 @@ type ServiceState DiscoveryService = DiscoveryPeerState emptyServiceState _ = DiscoveryPeerState- { dpsOurTunnelRequests = []+ { dpsWeAskedFor = M.empty+ , dpsPeerSearchingFor = M.empty+ , dpsOurTunnelRequests = [] , dpsRelayedTunnelRequests = [] , dpsStunServer = Nothing , dpsTurnServer = Nothing@@ -257,9 +320,6 @@ , show paddrs ] - let insertHelper new old | dpPriority new > dpPriority old = new- | otherwise = old- let matchedAddrs = flip filter addrs $ \case DiscoveryICE -> True DiscoveryIP ipaddr port ->@@ -273,7 +333,11 @@ , dpAddress = matchedAddrs , dpIceSession = Nothing }- svcModifyGlobal $ \s -> s { dgsPeers = M.insertWith insertHelper (refDigest $ storedRef sdata) dp $ dgsPeers s }+ rv = ResultValue+ { rvDirect = Just dp+ , rvVia = [ dp ]+ }+ svcModifyGlobal $ \s -> s { dgsPeers = M.insertWith (<>) (refDigest $ storedRef sdata) rv $ dgsPeers s } attrs <- asks svcAttributes replyPacket $ DiscoveryAcknowledged matchedAddrs (discoveryStunServer attrs)@@ -281,6 +345,44 @@ (discoveryTurnServer attrs) (discoveryTurnPort attrs) + server <- asks svcServer+ afterCommit $ void $ forkIO $ do+ peers <- getCurrentPeerList server+ let dgsts = identityDigests pid+ forM_ (filter (peer /=) peers) $ \sp -> do+ runPeerService @DiscoveryService sp $ do+ peerSearchingFor <- dpsPeerSearchingFor <$> svcGet+ when (any (`M.member` peerSearchingFor) dgsts) $ do+ let peerSearchingFor' = foldl' (flip M.delete) peerSearchingFor dgsts+ svcModify $ \s -> s { dpsPeerSearchingFor = peerSearchingFor' }+ spid <- asks svcPeerIdentity+ st <- getStorage+ forM_ dgsts $ \dgst -> do+ when (dgst `M.member` peerSearchingFor) $ do+ offerTunnel <- offerTunnelBetween attrs peer sp >>= return . \case+ True -> (++ [ DiscoveryTunnel ])+ False -> id+ let discoveryAddrs = offerTunnel matchedAddrs+ let ( results, via )+ | dgst == (refDigest $ storedRef $ idData pid)+ = ( discoveryAddrs, [] )+ | otherwise+ -- Results should be empty for this case (not searching exactly for the device id),+ -- but keep compatibility for now.+ = ( discoveryAddrs, [ DiscoveryVia (refDigest $ storedRef $ idData pid) discoveryAddrs ] )++ debugLog $+ "found for " <> show (refDigest $ storedRef $ idData spid) <>+ " dgst " <> show dgst <>+ " result [" <> T.unpack (T.intercalate "," $ map toText results) <> "]" <>+ " via " <> show (map (\v -> ( viaIdentity v, map toText $ viaAddress v )) via)+ -- Try to promote weak ref to normal one for older peers:+ edgst <- maybe (Right dgst) Left <$> liftIO (refFromDigest st dgst)+ replyPacket $ DiscoveryResult edgst results via+ debugLog $+ "remains asked by " <> show (refDigest $ storedRef $ idData spid) <>+ ": " <> show (M.keys peerSearchingFor')+ DiscoveryAcknowledged _ stunServer stunPort turnServer turnPort -> do paddr <- asks svcPeerAddress >>= return . \case (DatagramAddress saddr) -> T.pack . show . fst <$> inetFromSockAddr saddr@@ -297,41 +399,77 @@ } DiscoverySearch edgst -> do- pid <- asks svcPeerIdentity- dpeer <- M.lookup (either refDigest id edgst) . dgsPeers <$> svcGetGlobal+ let dgst = either refDigest id edgst peer <- asks svcPeer- paddr <- asks svcPeerAddress- attrs <- asks svcAttributes- let offerTunnel- | discoveryProvideTunnel attrs peer paddr = (++ [ DiscoveryTunnel ])- | otherwise = id- let results = maybe [] (offerTunnel . dpAddress) dpeer- replyPacket $ DiscoveryResult edgst results- debugLog $ "search by " <> show (refDigest $ storedRef $ idData pid) <>- " for " <> show (either refDigest id edgst) <>- " result [" <> T.unpack (T.intercalate "," $ map toText results) <> "]"+ pid <- asks svcPeerIdentity+ (M.lookup dgst . dgsPeers <$> svcGetGlobal) >>= \case+ Just rv+ | direct <- (\p -> Just p <* guard (p /= peer)) =<< dpPeer =<< rvDirect rv+ -- Direct results should be empty unless searching exactly for the device id,+ -- but keep compatibility for now.+ , rvia <- filter ((Just peer /=) . dpPeer) $ rvVia rv+ , isJust direct || not (null rvia)+ -> do+ attrs <- asks svcAttributes+ offerTunnel <- case direct of+ Just dpeer -> offerTunnelBetween attrs peer dpeer >>= return . \case+ True -> (++ [ DiscoveryTunnel ])+ False -> id+ Nothing -> return id+ let results+ | isJust direct = offerTunnel $ maybe [] dpAddress $ rvDirect rv+ | otherwise = []+ via <- liftIO $ fmap catMaybes $ mapM (viaFromPeer attrs peer) rvia+ replyPacket $ DiscoveryResult edgst results via+ debugLog $ "search by " <> show (refDigest $ storedRef $ idData pid) <>+ " for " <> show (either refDigest id edgst) <>+ " result [" <> T.unpack (T.intercalate "," $ map toText results) <> "]" <>+ " via " <> show (map (\v -> ( viaIdentity v, map toText $ viaAddress v )) via) - DiscoveryResult _ [] -> do- -- not found- return ()+ _ -> do+ now <- liftIO $ getTime Monotonic+ searchingFor <- dpsPeerSearchingFor <$> svcGet+ let seachingFor' = M.insert dgst (SearchingSince now) searchingFor+ svcModify $ \s -> s { dpsPeerSearchingFor = seachingFor' }+ debugLog $ "search by " <> show (refDigest $ storedRef $ idData pid) <>+ " for " <> show (either refDigest id edgst) <>+ " not found"+ debugLog $ "peer " <> show (refDigest $ storedRef $ idData pid) <>+ " searching for " <> show (M.keys seachingFor') - DiscoveryResult edgst addrs -> do+ DiscoveryResult edgst addrs via -> do let dgst = either refDigest id edgst- -- TODO: check if we really requested that server <- asks svcServer st <- getStorage self <- svcSelf discoveryPeer <- asks svcPeer- let runAsService = runPeerService @DiscoveryService discoveryPeer+ pid <- asks svcPeerIdentity + weAskedFor <- dpsWeAskedFor <$> svcGet+ let askedFor = M.member dgst weAskedFor+ debugLog $+ "result from " <> show (refDigest $ storedRef $ idData pid) <>+ " for " <> show dgst <> ": [" <> T.unpack (T.intercalate "," $ map toText addrs) <> "]" <>+ " via " <> show (map (\v -> ( viaIdentity v, map toText $ viaAddress v )) via) <>+ (if askedFor then "" else " (not asked for)")++ when askedFor $ do+ let weAskedFor' = M.delete dgst weAskedFor+ svcModify $ \s -> s { dpsWeAskedFor = weAskedFor' }+ debugLog $+ "remains asked " <> show (refDigest $ storedRef $ idData pid) <>+ " for " <> show (M.keys weAskedFor')++ let runAsService = runPeerService @DiscoveryService @IO discoveryPeer+ let tryAddresses = \case DiscoveryIP ipaddr port : _ -> do void $ liftIO $ forkIO $ do let saddr = inetToSockAddr ( ipaddr, port ) peer <- serverPeer server saddr runAsService $ do- let upd dp = dp { dpPeer = Just peer }- svcModifyGlobal $ \s -> s { dgsPeers = M.alter (Just . upd . fromMaybe emptyPeer) dgst $ dgsPeers s }+ let upd rv = rv { rvDirect = Just $ (fromMaybe emptyPeer $ rvDirect rv) { dpPeer = Just peer } }+ svcModifyGlobal $ \s -> s { dgsPeers = M.alter (Just . upd . fromMaybe mempty) dgst $ dgsPeers s } DiscoveryICE : rest -> do #ifdef ENABLE_ICE_SUPPORT@@ -356,8 +494,8 @@ Left err -> printOp $ "Discovery: failed to send connection request: " ++ err runAsService $ do- let upd dp = dp { dpIceSession = Just ice }- svcModifyGlobal $ \s -> s { dgsPeers = M.alter (Just . upd . fromMaybe emptyPeer) dgst $ dgsPeers s }+ let upd rv = rv { rvDirect = Just $ (fromMaybe emptyPeer $ rvDirect rv) { dpIceSession = Just ice } }+ svcModifyGlobal $ \s -> s { dgsPeers = M.alter (Just . upd . fromMaybe mempty) dgst $ dgsPeers s } Nothing -> do #endif@@ -367,12 +505,18 @@ discoverySetupTunnelResponse dgst addr : rest -> do- svcPrint $ "Discovery: unsupported address in result: " ++ T.unpack (toText addr)+ debugLog $ "unsupported address in result: " ++ T.unpack (toText addr) tryAddresses rest - [] -> svcPrint $ "Discovery: no (supported) address received for " <> show dgst+ [] -> debugLog $ "no (supported) address received for " <> show dgst - tryAddresses addrs+ when askedFor $ do+ tryAddresses $ concat+ -- ignore direct connections for self/owner+ [ if dgst `elem` identityDigests self then [] else addrs+ ] +++ -- ignore connections via ourselves+ concat (map viaAddress $ filter ((refDigest (storedRef (idData self)) /=) . viaIdentity) via) DiscoveryConnectionRequest conn -> do self <- svcSelf@@ -423,17 +567,16 @@ else do -- request to some of our peers, relay peer <- asks svcPeer- paddr <- asks svcPeerAddress- mbdp <- M.lookup (either refDigest id $ dconnTarget conn) . dgsPeers <$> svcGetGlobal+ mbrv <- M.lookup (either refDigest id $ dconnTarget conn) . dgsPeers <$> svcGetGlobal streams <- receivedStreams- case mbdp of+ case rvDirect =<< mbrv of Nothing -> replyPacket $ DiscoveryConnectionResponse rconn Just dp | Just dpeer <- dpPeer dp -> if- | dconnTunnel conn -> if- | not (discoveryProvideTunnel attrs peer paddr) -> do+ | dconnTunnel conn -> offerTunnelBetween attrs peer dpeer >>= \case+ False -> do replyPacket $ DiscoveryConnectionResponse rconn- | fromSource : _ <- streams -> do+ True | fromSource : _ <- streams -> do void $ liftIO $ forkIO $ runPeerService @DiscoveryService dpeer $ do debugLog $ "setting up tunnel from " <> show (either refDigest id $ dconnSource conn) <> " to " <> show (either refDigest id $ dconnTarget conn)@@ -441,7 +584,7 @@ svcModify $ \s -> s { dpsRelayedTunnelRequests = ( either refDigest id $ dconnSource conn, ( fromSource, toTarget )) : dpsRelayedTunnelRequests s } replyPacket $ DiscoveryConnectionRequest conn- | otherwise -> do+ _ | otherwise -> do svcPrint $ "Discovery: missing stream on tunnel request (relay)" | otherwise -> do sendToPeer dpeer $ DiscoveryConnectionRequest conn@@ -464,9 +607,9 @@ -> do let saddr = inetToSockAddr ( ipaddr, port ) peer <- liftIO $ serverPeer server saddr- let upd dp = dp { dpPeer = Just peer }+ let upd rv = rv { rvDirect = Just $ (fromMaybe emptyPeer $ rvDirect rv) { dpPeer = Just peer } } svcModifyGlobal $ \s -> s- { dgsPeers = M.alter (Just . upd . fromMaybe emptyPeer) (either refDigest id $ dconnTarget conn) $ dgsPeers s }+ { dgsPeers = M.alter (Just . upd . fromMaybe mempty) (either refDigest id $ dconnTarget conn) $ dgsPeers s } | dconnTunnel conn , Just tunnelWriter <- lookup (either refDigest id (dconnTarget conn)) (dpsOurTunnelRequests dps)@@ -490,8 +633,8 @@ liftIO $ closeStream tunnelWriter #ifdef ENABLE_ICE_SUPPORT- | Just dp <- M.lookup (either refDigest id $ dconnTarget conn) dpeers- , Just ice <- dpIceSession dp+ | Just rv <- M.lookup (either refDigest id $ dconnTarget conn) dpeers+ , Just ice <- dpIceSession =<< rvDirect rv , Just rinfo <- dconnIceInfo conn -> do liftIO $ iceConnect ice rinfo $ void $ serverPeerIce server ice #endif@@ -504,7 +647,7 @@ filter ((either refDigest id (dconnSource conn) /=) . fst) (dpsRelayedTunnelRequests s) } case M.lookup (either refDigest id $ dconnSource conn) dpeers of- Just dp | Just dpeer <- dpPeer dp -> if+ Just ResultValue { rvDirect = Just dp } | Just dpeer <- dpPeer dp -> if -- successful tunnel request | dconnTunnel conn , Just ( fromSource, toTarget ) <- lookup (either refDigest id (dconnSource conn)) (dpsRelayedTunnelRequests dps)@@ -548,18 +691,44 @@ let searchingFor = foldl' (flip S.delete) (dgsSearchingFor gs) (identityDigests pid) svcModifyGlobal $ \s -> s { dgsSearchingFor = searchingFor } + searchForOwner <- asks (discoverySearchForOwner . svcAttributes) >>= \case+ True -> do+ lookupSharedValueM >>= \case+ Just (self :: ComposedIdentity) -> do+ return $ S.fromList $ map (refDigest . storedRef) $ idDataF self+ Nothing -> do+ return S.empty+ False -> return S.empty+ let searchingFor' = searchingFor `S.union` searchForOwner+ when (not $ null addrs) $ do sendToPeer peer $ DiscoverySelf addrs Nothing- forM_ searchingFor $ \dgst -> do- sendToPeer peer $ DiscoverySearch (Right dgst) + when (not $ null searchingFor') $ do+ forM_ searchingFor' $ \dgst -> do+ sendToPeer peer $ DiscoverySearch (Right dgst)++ now <- liftIO $ getTime Monotonic+ let weAskedFor' = M.fromAscList $ map (, SearchingSince now) $ S.toAscList searchingFor'+ svcModify $ \s -> s { dpsWeAskedFor = weAskedFor' }+ debugLog $+ "we asked new peer " <> show (refDigest $ storedRef $ idData pid) <>+ " for " <> show (M.keys weAskedFor')+ serviceUpdatedPeer = do pid <- asks svcPeerIdentity peer <- asks svcPeer isPeerDropped peer >>= \case True -> do peers <- dgsPeers <$> svcGetGlobal- let peers' = M.filter ((Just peer /=) . dpPeer) peers+ let peers' = M.mapMaybe removePeer peers+ removePeer rv =+ let rv' = rv { rvDirect = if (dpPeer =<< rvDirect rv) == Just peer then Nothing else rvDirect rv+ , rvVia = filter ((Just peer /=) . dpPeer) $ rvVia rv+ }+ in if isJust (rvDirect rv') || not (null (rvVia rv'))+ then Just rv'+ else Nothing svcModifyGlobal $ \s -> s { dgsPeers = peers' } debugLog $ "dropped peer " <> show [ refDigest $ storedRef $ idData pid, refDigest $ storedRef $ idExtData pid ] <> ", map size " <> show (M.size peers) <> " -> " <> show (M.size peers')@@ -622,16 +791,20 @@ return $ dgst `elem` identityDigests pid _ -> return False when (not $ or match) $ do- alreadySearching <- modifyServiceGlobalState server (Proxy @DiscoveryService) $ \s ->- let alreadySearching = S.member dgst $ dgsSearchingFor s- in ( if alreadySearching then s else s- { dgsSearchingFor = S.insert dgst $ dgsSearchingFor s- }- , alreadySearching- )- when (not alreadySearching) $ do- forM_ peers $ \peer -> do- sendToPeer peer $ DiscoverySearch $ Right dgst+ _ <- modifyServiceGlobalState server (Proxy @DiscoveryService) $ \s ->+ ( s { dgsSearchingFor = S.insert dgst $ dgsSearchingFor s }, () )+ now <- liftIO $ getTime Monotonic+ forM_ peers $ \peer -> do+ runPeerService peer $ do+ weAskedFor <- dpsWeAskedFor <$> svcGet+ when (not $ M.member dgst weAskedFor) $ do+ let weAskedFor' = M.insert dgst (SearchingSince now) weAskedFor+ svcModify $ \s -> s { dpsWeAskedFor = weAskedFor' }+ pid <- asks svcPeerIdentity+ debugLog $+ "we asked " <> show (refDigest $ storedRef $ idData pid) <>+ " for " <> show dgst <> " " <> show (M.keys weAskedFor')+ replyPacket $ DiscoverySearch $ Right dgst data TunnelAddress = TunnelAddress@@ -665,6 +838,16 @@ connectionToAddressClosed TunnelAddress {..} = do closeStream tunnelWriter++offerTunnelBetween :: MonadIO m => DiscoveryAttributes -> Peer -> Peer -> m Bool+offerTunnelBetween attrs p1 p2 =+ offerTunnelFor p1 >>= \case+ True -> return True+ False -> offerTunnelFor p2+ where+ offerTunnelFor peer = do+ addrs <- getPeerAddresses peer+ return $ any (discoveryProvideTunnel attrs peer) addrs relayStream :: StreamReader -> StreamWriter -> IO () relayStream r w = do
src/Erebos/Identity.hs view
@@ -325,7 +325,7 @@ findResult [] = Nothing findResult xs = sel $ fromSigned $ minimum xs -mergeIdentity :: (MonadStorage m, MonadError e m, FromErebosError e, MonadIO m) => Identity f -> m UnifiedIdentity+mergeIdentity :: (MonadStorage m, MonadError e m, FromErebosError e) => Identity f -> m UnifiedIdentity mergeIdentity idt | Just idt' <- toUnifiedIdentity idt = return idt' mergeIdentity idt@Identity {..} = do (owner, ownerData) <- case idOwner_ of@@ -335,7 +335,7 @@ return (Just owner, Just $ idData owner) let public = idKeyIdentity idt- secret <- loadKey public+ secret <- mloadKey public unifiedBaseData <- case toList $ idDataF idt of
src/Erebos/Invite.hs view
@@ -165,7 +165,7 @@ sharedTypeID _ = mkSharedTypeID "b1ebf228-4892-476b-ba04-0c26320139b1" -createSingleContactInvite :: MonadHead LocalState m => Text -> m Invite+createSingleContactInvite :: (MonadHead LocalState m, MonadIO m) => Text -> m Invite createSingleContactInvite name = do time <- liftIO getZonedTime token <- liftIO $ getRandomBytes 32@@ -190,8 +190,7 @@ -> m () acceptInvite from token = do prev <- find ((Just token ==) . acceptedInviteToken)- . fromSetBy (comparing acceptedInviteToken) . lookupSharedValue . lsShared . fromStored- <$> getLocalHead+ . fromSetBy (comparing acceptedInviteToken) <$> lookupSharedValueM accepted <- mergeSorted @AcceptedInvite . (: []) <$> mstore AcceptedInviteData { aidPrev = maybe [] acceptedInviteData prev , aidToken = Just token@@ -261,7 +260,7 @@ serviceHandler = fromStored >>> \case AcceptInvite token -> do- invites <- fromSetBy (comparing inviteToken) . lookupSharedValue . lsShared . fromStored <$> getLocalHead+ invites <- fromSetBy (comparing inviteToken) <$> lookupSharedValueM case find ((Just token ==) . inviteToken) invites of Just invite | Just name <- inviteContact invite@@ -296,7 +295,7 @@ asks (inviteHookReplyInvalid . svcAttributes) >>= ($ token) svcModify $ filter (/= token) - accepted <- fromSetBy (comparing acceptedInviteToken) . lookupSharedValue . lsShared . fromStored <$> getLocalHead+ accepted <- fromSetBy (comparing acceptedInviteToken) <$> lookupSharedValueM case find ((Just token ==) . acceptedInviteToken) accepted of Just invite -> do aidata <- mstore AcceptedInviteData@@ -316,7 +315,7 @@ then do svcSet $ filter (/= token) waitingTokens - accepted <- fromSetBy (comparing acceptedInviteToken) . lookupSharedValue . lsShared . fromStored <$> getLocalHead+ accepted <- fromSetBy (comparing acceptedInviteToken) <$> lookupSharedValueM case find ((Just token ==) . acceptedInviteToken) accepted of Just invite -> do aidata <- mstore AcceptedInviteData@@ -343,7 +342,7 @@ svcPrint $ "Received unknown invite packet" serviceNewPeer = do- invites <- fromSetBy (comparing acceptedInviteToken) . lookupSharedValue . lsShared . fromStored <$> getLocalHead+ invites <- fromSetBy (comparing acceptedInviteToken) <$> lookupSharedValueM peerDigests <- asks $ identityOwnerDigests . svcPeerIdentity forM_ invites $ \case AcceptedInvite@@ -356,7 +355,7 @@ _ -> return () serviceStorageWatchers _ = (:[]) $- GlobalStorageWatcher (lookupSharedValue . lsShared . fromStored) sendAcceptedInvites+ globalStorageWatcher (lookupSharedValue . lsShared . fromStored) sendAcceptedInvites sendAcceptedInvites :: Server -> Set AcceptedInvite -> ExceptT ErebosError IO ()
src/Erebos/Network.hs view
@@ -4,6 +4,7 @@ stopServer, getCurrentPeerList, getNextPeerChange,+ getNextPeerChangeChan, getServerAddresses, ServerOptions(..), serverIdentity, defaultServerOptions, @@ -75,10 +76,13 @@ import Erebos.State import Erebos.Storable.Internal import Erebos.Storage+import Erebos.Storage.Head import Erebos.Storage.Key import Erebos.Storage.Merge +import Service + discoveryPort :: PortNumber discoveryPort = 29665 @@ -93,6 +97,7 @@ { serverStorage :: Storage , serverOptions :: ServerOptions , serverOrigHead :: Head LocalState+ , serverCurrentHead :: MVar (Head LocalState) , serverIdentity_ :: MVar UnifiedIdentity , serverThreads :: MVar [ThreadId] , serverSocket :: MVar Socket@@ -117,6 +122,9 @@ getNextPeerChange :: Server -> IO Peer getNextPeerChange = atomically . readTChan . serverChanPeer +getNextPeerChangeChan :: Server -> IO (TChan Peer)+getNextPeerChangeChan = atomically . dupTChan . serverChanPeer+ data ServerOptions = ServerOptions { serverPort :: PortNumber , serverRetryUnspecifiedPort :: Bool@@ -180,6 +188,8 @@ class (Eq addr, Ord addr, Show addr, Typeable addr) => PeerAddressType addr where sendBytesToAddress :: addr -> ByteString -> IO () connectionToAddressClosed :: addr -> IO ()+ isReliableTransport :: proxy addr -> Bool+ isReliableTransport _ = False data PeerAddress = forall addr. PeerAddressType addr => CustomPeerAddress addr@@ -208,7 +218,12 @@ compare (DatagramAddress addr) (DatagramAddress addr') = compare addr addr' +instance ProtocolAddressType PeerAddress where+ addrIsReliableTransport = \case+ CustomPeerAddress (_ :: a) -> isReliableTransport (Proxy @a)+ DatagramAddress _ -> False + data PeerIdentity = PeerIdentityUnknown (TVar [ UnifiedIdentity -> ExceptT ErebosError IO () ]) | PeerIdentityRef WaitingRef (TVar [ UnifiedIdentity -> ExceptT ErebosError IO () ])@@ -263,6 +278,7 @@ startServer :: ServerOptions -> Head LocalState -> (String -> IO ()) -> [SomeService] -> IO Server startServer serverOptions serverOrigHead logd' serverServices = do let serverStorage = headStorage serverOrigHead+ serverCurrentHead <- newMVar serverOrigHead serverIdentity_ <- newMVar $ headLocalIdentity serverOrigHead serverThreads <- newMVar [] serverSocket <- newEmptyMVar@@ -339,10 +355,21 @@ PeerIdentityFull _ -> writeTQueue serverIOActions $ do runPeerService peer $ act . sel =<< svcGetLocal _ -> return ()+ SomeStorageWatcherHC sel act -> do+ watchHeadWith serverOrigHead (\h -> sel (headStoredObject h) (headCache h)) $ \_ -> do+ withMVar serverPeers $ mapM_ $ \peer -> atomically $ do+ readTVar (peerIdentityVar peer) >>= \case+ PeerIdentityFull _ -> writeTQueue serverIOActions $ do+ runPeerService peer $ act =<< (sel <$> getLocalHead <*> getLocalHeadCache @LocalState Proxy)+ _ -> return () GlobalStorageWatcher sel act -> do watchHeadWith serverOrigHead (sel . headStoredObject) $ \x -> do atomically $ writeTQueue serverIOActions $ do act server x+ GlobalStorageWatcherH sel act -> do+ watchHeadWith serverOrigHead sel $ \x -> do+ atomically $ writeTQueue serverIOActions $ do+ act server x forkServerThread server "receiver" $ forever $ do (msg, saddr) <- S.recvFrom sock 4096@@ -1058,19 +1085,22 @@ , svcPrintOp = atomically . logd , svcNewStreams = newStreams }- reloadHead (serverOrigHead server) >>= \case- Nothing -> atomically $ do- logd $ "current head deleted"- putTMVar (peerServiceState peer) svcs- putTMVar (serverServiceStates server) global- Just h -> do- (rsp, (s', gs')) <- runServiceHandler h inp ps gs handler- moveKeys (peerStorage peer) (serverStorage server)- when (not (null rsp)) $ do- sendToPeerList peer rsp- atomically $ do- putTMVar (peerServiceState peer) $ M.insert svc (SomeServiceState proxy s') svcs- putTMVar (serverServiceStates server) $ M.insert svc (SomeServiceGlobalState proxy gs') global+ modifyMVar_ (serverCurrentHead server) $ \ph -> do+ reloadHead ph >>= \case+ Nothing -> atomically $ do+ logd $ "current head deleted"+ putTMVar (peerServiceState peer) svcs+ putTMVar (serverServiceStates server) global+ return ph+ Just h -> do+ ( rsp, ( s', gs', h' ) ) <- runServiceHandler h inp ps gs handler+ moveKeys (peerStorage peer) (serverStorage server)+ when (not (null rsp)) $ do+ sendToPeerList peer rsp+ atomically $ do+ putTMVar (peerServiceState peer) $ M.insert svc (SomeServiceState proxy s') svcs+ putTMVar (serverServiceStates server) $ M.insert svc (SomeServiceGlobalState proxy gs') global+ return h' _ -> do atomically $ logd $ "can't run service handler on peer with incomplete identity " ++ show paddr @@ -1129,16 +1159,22 @@ getServerAddresses :: Server -> IO [ SockAddr ] getServerAddresses Server {..} = do- alloca $ \pcount -> do- ptr <- cLocalAddresses pcount- if ptr == nullPtr- then do- return []- else do- count <- fromIntegral <$> peek pcount- res <- peekArray count ptr- cFree ptr- return $ map (inetToSockAddr . (, serverPort serverOptions )) res+ tryReadMVar serverSocket >>= \case+ Just sock -> do+ getSocketName sock >>= \case+ SockAddrInet6 port _ _ _ -> do+ alloca $ \pcount -> do+ ptr <- cLocalAddresses pcount+ if ptr == nullPtr+ then do+ return []+ else do+ count <- fromIntegral <$> peek pcount+ res <- peekArray count ptr+ cFree ptr+ return $ map (inetToSockAddr . (, port )) res+ _ -> return []+ Nothing -> return [] getBroadcastAddresses :: PortNumber -> IO [SockAddr] getBroadcastAddresses port = do
src/Erebos/Network/Protocol.hs view
@@ -17,6 +17,7 @@ ControlMessage(..), erebosNetworkProtocol, + ProtocolAddressType(..), Connection, connAddress, connData,@@ -210,6 +211,9 @@ transportFromObject _ = Nothing +class ProtocolAddressType addr where+ addrIsReliableTransport :: addr -> Bool+ data GlobalState addr = (Eq addr, Show addr) => GlobalState { gIdentity :: TVar (UnifiedIdentity, [UnifiedIdentity]) , gConnections :: TVar [Connection addr]@@ -246,6 +250,7 @@ , cNextKeepAlive :: TVar (Maybe TimeSpec) , cInStreams :: TVar [(Word8, Stream)] , cOutStreams :: TVar [(Word8, Stream)]+ , cIsReliable :: Bool } instance Eq (Connection addr) where@@ -517,7 +522,7 @@ | ReceivedAnnounce addr RefDigest -erebosNetworkProtocol :: (Eq addr, Ord addr, Show addr)+erebosNetworkProtocol :: (Eq addr, Ord addr, Show addr, ProtocolAddressType addr) => UnifiedIdentity -> (String -> STM ()) -> (String -> STM ())@@ -562,7 +567,7 @@ catch io $ \(e :: SomeException) -> atomically $ gLog $ "exception during network protocol handling: " <> show e -getConnection :: GlobalState addr -> addr -> STM (Connection addr)+getConnection :: ProtocolAddressType addr => GlobalState addr -> addr -> STM (Connection addr) getConnection gs addr = do maybe (newConnection gs addr) return =<< findConnection gs addr @@ -570,11 +575,12 @@ findConnection GlobalState {..} addr = do find ((addr==) . cAddress) <$> readTVar gConnections -newConnection :: GlobalState addr -> addr -> STM (Connection addr)+newConnection :: ProtocolAddressType addr => GlobalState addr -> addr -> STM (Connection addr) newConnection cGlobalState@GlobalState {..} addr = do conns <- readTVar gConnections let cAddress = addr+ cIsReliable = addrIsReliableTransport addr (cDataUp, cDataInternal) <- newFlow cChannel <- newTVar ChannelNone cCookie <- newTVar Nothing@@ -598,7 +604,7 @@ writeFlow cDataInternal (Just up) return $ return () -processIncoming :: GlobalState addr -> STM (IO ())+processIncoming :: ProtocolAddressType addr => GlobalState addr -> STM (IO ()) processIncoming gs@GlobalState {..} = do guard =<< isEmptyTMVar gNextUp guard =<< canWriteFlow gControlFlow@@ -697,7 +703,7 @@ Left err -> do atomically $ gLog $ show addr <> ": failed to parse packet: " <> showErebosError err -processPacket :: GlobalState addr -> Either addr (Connection addr) -> Bool -> TransportPacket a -> IO (Maybe (Connection addr, Maybe (TransportPacket a)))+processPacket :: ProtocolAddressType addr => GlobalState addr -> Either addr (Connection addr) -> Bool -> TransportPacket a -> IO (Maybe (Connection addr, Maybe (TransportPacket a))) processPacket gs@GlobalState {..} econn secure packet@(TransportPacket (TransportHeader header) _) = if -- Established secure communication | Right conn <- econn, secure@@ -738,18 +744,20 @@ _ -> return Nothing -- Initiation packet- | _:_ <- mapMaybe (\case Initiation x -> Just x; _ -> Nothing) header+ | _ : _ <- mapMaybe (\case Initiation x -> Just x; _ -> Nothing) header+ , pids <- mapMaybe (\case AnnounceSelf x -> Just x; _ -> Nothing) header , Just ver <- version -> do cookie <- createCookie gs addr atomically $ do identity <- fst <$> readTVar gIdentity- let reply = BL.toStrict $ serializeObject $ transportToObject gStorage $ TransportHeader- [ CookieSet cookie- , AnnounceSelf $ refDigest $ storedRef $ idData identity- , ProtocolVersion ver- ]- writeFlow gDataFlow (addr, reply)+ when ((refDigest $ storedRef $ idData identity) `notElem` pids) $ do+ let reply = BL.toStrict $ serializeObject $ transportToObject gStorage $ TransportHeader+ [ CookieSet cookie+ , AnnounceSelf $ refDigest $ storedRef $ idData identity+ , ProtocolVersion ver+ ]+ writeFlow gDataFlow (addr, reply) return Nothing -- Announce packet outside any connection@@ -873,7 +881,7 @@ writeTVar cNextKeepAlive $ Just next -processOutgoing :: forall addr. GlobalState addr -> STM (IO ())+processOutgoing :: forall addr. ProtocolAddressType addr => GlobalState addr -> STM (IO ()) processOutgoing gs@GlobalState {..} = do let sendNextPacket :: Connection addr -> STM (IO ())@@ -1024,7 +1032,7 @@ conns <- readTVar gConnections msum $ concat $- [ map retransmitPacket conns+ [ map retransmitPacket $ filter (not . cIsReliable) conns , map sendNextPacket conns , [ handleControlRequests ] , map sendKeepAlive conns
src/Erebos/Object/Internal.hs view
@@ -30,7 +30,7 @@ storeRecItems, Load, LoadRec,- evalLoad,+ evalLoad, evalLoadWithObject, loadCurrentRef, loadCurrentObject, loadRecCurrentRef, loadRecItems, @@ -454,8 +454,11 @@ deriving (Functor, Applicative, Alternative, Monad, MonadPlus, MonadError ErebosError) evalLoad :: Load a -> Ref -> a-evalLoad (Load f) ref = either (error {- TODO throw -} . ((BC.unpack (showRef ref) ++ ": ") ++) . showErebosError) id $- runExcept $ runReaderT f (ref, lazyLoadObject ref)+evalLoad act ref = evalLoadWithObject act ref (lazyLoadObject ref)++evalLoadWithObject :: Load a -> Ref -> Object -> a+evalLoadWithObject (Load f) ref obj = either (error {- TODO throw -} . ((BC.unpack (showRef ref) ++ ": ") ++) . showErebosError) id $+ runExcept $ runReaderT f ( ref, obj ) loadCurrentRef :: Load Ref loadCurrentRef = Load $ asks fst
src/Erebos/Service.hs view
@@ -3,7 +3,8 @@ SomeService(..), someService, someServiceAttr, someServiceID, SomeServiceState(..), fromServiceState, someServiceEmptyState, SomeServiceGlobalState(..), fromServiceGlobalState, someServiceEmptyGlobalState,- SomeStorageWatcher(..),+ SomeStorageWatcher(SomeStorageWatcher, GlobalStorageWatcher),+ someStorageWatcher, someStorageWatcherHC, globalStorageWatcher, globalStorageWatcherH, ServiceID, mkServiceID, ServiceHandler,@@ -23,184 +24,24 @@ ) where import Control.Monad.Except-import Control.Monad.Reader-import Control.Monad.State-import Control.Monad.Writer -import Data.Kind-import Data.Typeable--import Erebos.Identity import {-# SOURCE #-} Erebos.Network import Erebos.Network.Protocol import Erebos.State import Erebos.Storable import Erebos.Storage.Head-import Erebos.UUID qualified as U -class (- Typeable s, Storable s,- Typeable (ServiceAttributes s),- Typeable (ServiceState s),- Typeable (ServiceGlobalState s)- ) => Service s where-- serviceID :: proxy s -> ServiceID- serviceHandler :: Stored s -> ServiceHandler s ()-- serviceNewPeer :: ServiceHandler s ()- serviceNewPeer = return ()-- serviceUpdatedPeer :: ServiceHandler s ()- serviceUpdatedPeer = return ()-- type ServiceAttributes s = attr | attr -> s- type ServiceAttributes s = Proxy s- defaultServiceAttributes :: proxy s -> ServiceAttributes s- default defaultServiceAttributes :: ServiceAttributes s ~ Proxy s => proxy s -> ServiceAttributes s- defaultServiceAttributes _ = Proxy-- type ServiceState s :: Type- type ServiceState s = ()- emptyServiceState :: proxy s -> ServiceState s- default emptyServiceState :: ServiceState s ~ () => proxy s -> ServiceState s- emptyServiceState _ = ()-- type ServiceGlobalState s :: Type- type ServiceGlobalState s = ()- emptyServiceGlobalState :: proxy s -> ServiceGlobalState s- default emptyServiceGlobalState :: ServiceGlobalState s ~ () => proxy s -> ServiceGlobalState s- emptyServiceGlobalState _ = ()-- serviceStorageWatchers :: proxy s -> [SomeStorageWatcher s]- serviceStorageWatchers _ = []-- serviceStopServer :: proxy s -> Server -> ServiceGlobalState s -> [ ( Peer, ServiceState s ) ] -> IO ()- serviceStopServer _ _ _ _ = return ()---data SomeService = forall s. Service s => SomeService (Proxy s) (ServiceAttributes s)--someService :: forall s proxy. Service s => proxy s -> SomeService-someService _ = SomeService @s Proxy (defaultServiceAttributes @s Proxy)--someServiceAttr :: forall s. Service s => ServiceAttributes s -> SomeService-someServiceAttr attr = SomeService @s Proxy attr--someServiceID :: SomeService -> ServiceID-someServiceID (SomeService s _) = serviceID s--data SomeServiceState = forall s. Service s => SomeServiceState (Proxy s) (ServiceState s)--fromServiceState :: Service s => proxy s -> SomeServiceState -> Maybe (ServiceState s)-fromServiceState _ (SomeServiceState _ s) = cast s--someServiceEmptyState :: SomeService -> SomeServiceState-someServiceEmptyState (SomeService p _) = SomeServiceState p (emptyServiceState p)--data SomeServiceGlobalState = forall s. Service s => SomeServiceGlobalState (Proxy s) (ServiceGlobalState s)--fromServiceGlobalState :: Service s => proxy s -> SomeServiceGlobalState -> Maybe (ServiceGlobalState s)-fromServiceGlobalState _ (SomeServiceGlobalState _ s) = cast s--someServiceEmptyGlobalState :: SomeService -> SomeServiceGlobalState-someServiceEmptyGlobalState (SomeService p _) = SomeServiceGlobalState p (emptyServiceGlobalState p)---data SomeStorageWatcher s- = forall a. Eq a => SomeStorageWatcher (Stored LocalState -> a) (a -> ServiceHandler s ())- | forall a. Eq a => GlobalStorageWatcher (Stored LocalState -> a) (Server -> a -> ExceptT ErebosError IO ())+import Service -mkServiceID :: String -> ServiceID-mkServiceID = maybe (error "Invalid service ID") ServiceID . U.fromString--data ServiceInput s = ServiceInput- { svcAttributes :: ServiceAttributes s- , svcPeer :: Peer- , svcPeerAddress :: PeerAddress- , svcPeerIdentity :: UnifiedIdentity- , svcServer :: Server- , svcPrintOp :: String -> IO ()- , svcNewStreams :: [ RawStreamReader ]- }--data ServiceReply s- = ServiceReply (Either s (Stored s)) Bool- | ServiceOpenStream (RawStreamWriter -> IO ())- | ServiceFinally (IO ())--data ServiceHandlerState s = ServiceHandlerState- { svcValue :: ServiceState s- , svcGlobal :: ServiceGlobalState s- , svcLocal :: Stored LocalState- }--newtype ServiceHandler s a = ServiceHandler (ReaderT (ServiceInput s) (WriterT [ServiceReply s] (StateT (ServiceHandlerState s) (ExceptT ErebosError IO))) a)- deriving (Functor, Applicative, Monad, MonadReader (ServiceInput s), MonadWriter [ServiceReply s], MonadState (ServiceHandlerState s), MonadError ErebosError, MonadIO)--instance MonadStorage (ServiceHandler s) where- getStorage = asks $ peerStorage . svcPeer--instance MonadHead LocalState (ServiceHandler s) where- updateLocalHead f = do- (ls, x) <- f =<< gets svcLocal- modify $ \s -> s { svcLocal = ls }- return x--runServiceHandler :: Service s => Head LocalState -> ServiceInput s -> ServiceState s -> ServiceGlobalState s -> ServiceHandler s () -> IO ([ServiceReply s], (ServiceState s, ServiceGlobalState s))-runServiceHandler h input svc global shandler = do- let sstate = ServiceHandlerState { svcValue = svc, svcGlobal = global, svcLocal = headStoredObject h }- ServiceHandler handler = shandler- (runExceptT $ flip runStateT sstate $ execWriterT $ flip runReaderT input $ handler) >>= \case- Left err -> do- svcPrintOp input $ "service failed: " ++ showErebosError err- return ([], (svc, global))- Right (rsp, sstate')- | svcLocal sstate' == svcLocal sstate -> return (rsp, (svcValue sstate', svcGlobal sstate'))- | otherwise -> replaceHead h (svcLocal sstate') >>= \case- Left (Just h') -> runServiceHandler h' input svc global shandler- _ -> return (rsp, (svcValue sstate', svcGlobal sstate'))--svcGet :: ServiceHandler s (ServiceState s)-svcGet = gets svcValue--svcSet :: ServiceState s -> ServiceHandler s ()-svcSet x = modify $ \st -> st { svcValue = x }--svcModify :: (ServiceState s -> ServiceState s) -> ServiceHandler s ()-svcModify f = modify $ \st -> st { svcValue = f (svcValue st) }--svcGetGlobal :: ServiceHandler s (ServiceGlobalState s)-svcGetGlobal = gets svcGlobal--svcSetGlobal :: ServiceGlobalState s -> ServiceHandler s ()-svcSetGlobal x = modify $ \st -> st { svcGlobal = x }--svcModifyGlobal :: (ServiceGlobalState s -> ServiceGlobalState s) -> ServiceHandler s ()-svcModifyGlobal f = modify $ \st -> st { svcGlobal = f (svcGlobal st) }--svcGetLocal :: ServiceHandler s (Stored LocalState)-svcGetLocal = gets svcLocal--svcSetLocal :: Stored LocalState -> ServiceHandler s ()-svcSetLocal x = modify $ \st -> st { svcLocal = x }--svcSelf :: ServiceHandler s UnifiedIdentity-svcSelf = maybe (throwOtherError "failed to validate own identity") return .- validateExtendedIdentity . lsIdentity . fromStored =<< svcGetLocal--svcPrint :: String -> ServiceHandler s ()-svcPrint str = afterCommit . ($ str) =<< asks svcPrintOp--replyPacket :: Service s => s -> ServiceHandler s ()-replyPacket x = tell [ServiceReply (Left x) True]+someStorageWatcher :: forall s a. (Service s, Eq a) => (Stored LocalState -> a) -> (a -> ServiceHandler s ()) -> SomeStorageWatcher s+someStorageWatcher = SomeStorageWatcher -replyStored :: Service s => Stored s -> ServiceHandler s ()-replyStored x = tell [ServiceReply (Right x) True]+someStorageWatcherHC :: forall s a. (Service s, Eq a) => (Stored LocalState -> HeadCacheType LocalState -> a) -> (a -> ServiceHandler s ()) -> SomeStorageWatcher s+someStorageWatcherHC = SomeStorageWatcherHC -replyStoredRef :: Service s => Stored s -> ServiceHandler s ()-replyStoredRef x = tell [ServiceReply (Right x) False]+globalStorageWatcher :: forall s a. (Service s, Eq a) => (Stored LocalState -> a) -> (Server -> a -> ExceptT ErebosError IO ()) -> SomeStorageWatcher s+globalStorageWatcher = GlobalStorageWatcher -afterCommit :: IO () -> ServiceHandler s ()-afterCommit x = tell [ServiceFinally x]+globalStorageWatcherH :: forall s a. (Service s, Eq a) => (Head LocalState -> a) -> (Server -> a -> ExceptT ErebosError IO ()) -> SomeStorageWatcher s+globalStorageWatcherH = GlobalStorageWatcherH
src/Erebos/State.hs view
@@ -6,10 +6,12 @@ MonadStorage(..), MonadHead(..), updateLocalHead_,- LocalHeadT(..),+ LocalHeadT, runLocalHeadT, runLocalHeadT',+ LocalHead, runLocalHead, runLocalHead', updateLocalState, updateLocalState_, updateSharedState, updateSharedState_,+ lookupSharedValueH, lookupSharedValueHC, lookupSharedValueM, lookupSharedValue, makeSharedStateUpdate, localIdentity,@@ -20,11 +22,15 @@ import Control.Monad import Control.Monad.Except+import Control.Monad.Identity qualified as CMI import Control.Monad.Reader import Data.Bifunctor import Data.ByteString (ByteString) import Data.ByteString.Char8 qualified as BC+import Data.Foldable+import Data.Map.Strict qualified as MS+import Data.Maybe import Data.Typeable import Erebos.Identity@@ -35,11 +41,12 @@ import Erebos.Storage.Merge import Erebos.UUID (UUID) import Erebos.UUID qualified as U+import Erebos.Util data LocalState = LocalState { lsPrev :: Maybe RefDigest , lsIdentity :: Stored (Signed ExtendedIdentityData)- , lsShared :: [Stored SharedState]+ , lsShared :: [ Stored SharedState ] , lsOther :: [ ( ByteString, RecItem ) ] } @@ -74,7 +81,22 @@ instance HeadType LocalState where headTypeID _ = mkHeadTypeID "1d7491a9-7bcb-4eaa-8f13-c8c4c4087e4e"+ type HeadCacheType LocalState = LocalStateCache+ headCacheInit sls =+ let LocalState {..} = fromStored sls+ sharedCache = MS.fromAscList $ map (\sid -> ( sid, lookupSharedValueObjects sid [] lsShared )) $ collectSharedTypeIDs [] lsShared+ in LocalStateCache lsShared sharedCache+ headCacheUpdate sls prev =+ let LocalState {..} = fromStored sls+ sids = collectSharedTypeIDs (lscSharedTips prev) lsShared+ upd cache sid = MS.insertWith (\ss ss' -> filterAncestors (ss ++ ss')) sid (lookupSharedValueObjects sid (lscSharedTips prev) lsShared) cache+ in LocalStateCache lsShared (foldl' upd (lscSharedCache prev) sids) +data LocalStateCache = LocalStateCache+ { lscSharedTips :: [ Stored SharedState ]+ , lscSharedCache :: MS.Map SharedTypeID (StoredTips SharedState)+ }+ instance Storable SharedState where store' st = storeRec $ do mapM_ (storeRef "PREV") $ ssPrev st@@ -90,10 +112,12 @@ sharedTypeID _ = mkSharedTypeID "0c6c1fe0-f2d7-4891-926b-c332449f7871" -class (MonadIO m, MonadStorage m) => MonadHead a m where+class (MonadStorage m, HeadType a) => MonadHead a m where updateLocalHead :: (Stored a -> m (Stored a, b)) -> m b getLocalHead :: m (Stored a) getLocalHead = updateLocalHead $ \x -> return (x, x)+ getLocalHeadCache :: proxy a -> m (HeadCacheType a)+ getLocalHeadCache _ = headCacheInit @a <$> getLocalHead updateLocalHead_ :: MonadHead a m => (Stored a -> m (Stored a)) -> m () updateLocalHead_ f = updateLocalHead (fmap (,()) . f)@@ -101,11 +125,27 @@ instance (HeadType a, MonadIO m) => MonadHead a (ReaderT (Head a) m) where updateLocalHead f = do h <- ask- snd <$> updateHead h f+ snd <$> updateHead' h (\h' -> local (const h') (f $ headStoredObject h'))+ getLocalHeadCache _ = asks headCache -newtype LocalHeadT h m a = LocalHeadT { runLocalHeadT :: Storage -> Stored h -> m ( a, Stored h ) }+newtype LocalHeadT h m a = LocalHeadT { runLocalHeadT_ :: Storage -> ( Stored h, HeadCacheType h ) -> m ( a, ( Stored h, HeadCacheType h ) ) } +runLocalHeadT :: forall h m a. (HeadType h, Monad m) => LocalHeadT h m a -> Storage -> Stored h -> m ( a, Stored h )+runLocalHeadT act st h = fmap fst <$> runLocalHeadT_ act st ( h, headCacheInit h )++runLocalHeadT' :: forall h m a. (HeadType h, Monad m) => LocalHeadT h m a -> Head h -> m ( a, Stored h )+runLocalHeadT' act h = fmap fst <$> runLocalHeadT_ act (headStorage h) ( headStoredObject h, headCache h )++type LocalHead h a = LocalHeadT h CMI.Identity a++runLocalHead :: forall h a. HeadType h => LocalHead h a -> Storage -> Stored h -> ( a, Stored h )+runLocalHead act st h = CMI.runIdentity $ runLocalHeadT act st h++runLocalHead' :: forall h a. HeadType h => LocalHead h a -> Head h -> ( a, Stored h )+runLocalHead' act h = CMI.runIdentity $ runLocalHeadT' act h++ instance Functor m => Functor (LocalHeadT h m) where fmap f (LocalHeadT act) = LocalHeadT $ \st h -> first f <$> act st h @@ -127,12 +167,15 @@ getStorage = LocalHeadT $ \st h -> return ( st, h ) instance (HeadType h, MonadIO m) => MonadHead h (LocalHeadT h m) where- updateLocalHead f = LocalHeadT $ \st h -> do+ updateLocalHead f = LocalHeadT $ \st ( h, c ) -> do let LocalHeadT act = f h- ( ( h', x ), _ ) <- act st h- return ( x, h' )+ ( ( h'', x ), ( _, c' ) ) <- act st ( h, c )+ let c'' = headCacheUpdate h'' c'+ return ( x, ( h'', c'' ) ) + getLocalHeadCache _ = LocalHeadT $ \_ hc@( _, c ) -> return ( c, hc ) + localIdentity :: LocalState -> UnifiedIdentity localIdentity ls = maybe (error "failed to verify local identity") (updateOwners $ maybe [] idExtDataF $ lookupSharedValue $ lsShared ls)@@ -159,18 +202,46 @@ updateSharedState :: forall a b m. (SharedType a, MonadHead LocalState m) => (a -> m (a, b)) -> Stored LocalState -> m (Stored LocalState, b) updateSharedState f = \ls -> do let shared = lsShared $ fromStored ls- val = lookupSharedValue shared+ val <- lookupSharedValueM (val', x) <- f val (,x) <$> if toComponents val' == toComponents val then return ls else do shared' <- makeSharedStateUpdate val' shared mstore (fromStored ls) { lsShared = [shared'] } -lookupSharedValue :: forall a. SharedType a => [Stored SharedState] -> a-lookupSharedValue = mergeSorted . filterAncestors . map wrappedLoad . concatMap (ssValue . fromStored) . filterAncestors . helper- where helper (x:xs) | Just sid <- ssType (fromStored x), sid == sharedTypeID @a Proxy = x : helper xs- | otherwise = helper $ ssPrev (fromStored x) ++ xs- helper [] = []++collectSharedTypeIDs :: [ Stored SharedState ] -> [ Stored SharedState ] -> [ SharedTypeID ]+collectSharedTypeIDs since (x : xs)+ | any (x `precedesOrEquals`) since+ = collectSharedTypeIDs since xs+ | otherwise+ = maybeToList (ssType (fromStored x)) `mergeUniq` collectSharedTypeIDs since (ssPrev (fromStored x) ++ xs)+collectSharedTypeIDs _ [] = []++lookupSharedValueObjects :: SharedTypeID -> [ Stored SharedState ] -> [ Stored SharedState ] -> StoredTips SharedState+lookupSharedValueObjects sid since = filterAncestors . helper+ where+ helper (x : xs)+ | any (x `precedesOrEquals`) since = helper xs+ | ssType (fromStored x) == Just sid = x : helper xs+ | otherwise = helper $ ssPrev (fromStored x) ++ xs+ helper [] = []++lookupSharedValueC :: forall a. SharedType a => HeadCacheType LocalState -> a+lookupSharedValueC = mergeSorted . filterAncestors . map wrappedLoad . concatMap (ssValue . fromStored) .+ fromMaybe [] . MS.lookup (sharedTypeID @a Proxy) . lscSharedCache++lookupSharedValueH :: forall a. SharedType a => Head LocalState -> a+lookupSharedValueH = lookupSharedValueC . headCache++lookupSharedValueHC :: forall a. SharedType a => Stored LocalState -> HeadCacheType LocalState -> a+lookupSharedValueHC _ = lookupSharedValueC++lookupSharedValueM :: forall a m. (SharedType a, MonadHead LocalState m) => m a+lookupSharedValueM = lookupSharedValueC <$> getLocalHeadCache @LocalState Proxy++lookupSharedValue :: forall a. SharedType a => [ Stored SharedState ] -> a+lookupSharedValue = mergeSorted . filterAncestors . map wrappedLoad . concatMap (ssValue . fromStored) . lookupSharedValueObjects (sharedTypeID @a Proxy) [] makeSharedStateUpdate :: forall a m. (SharedType a, MonadStorage m) => a -> [ Stored SharedState ] -> m (Stored SharedState) makeSharedStateUpdate val prev = mstore SharedState
src/Erebos/Storable.hs view
@@ -36,7 +36,7 @@ copyStored, unsafeMapStored, - Storage, MonadStorage(..),+ Storage, MonadStorage(..), mloadKey, module Erebos.Error, ) where@@ -44,3 +44,4 @@ import Erebos.Error import Erebos.Object.Internal import Erebos.Storable.Internal+import Erebos.Storage.Monad
src/Erebos/Storable/Internal.hs view
@@ -5,13 +5,11 @@ Stored(..), fromStored, storedRef, storedStorage,- wrappedStore, wrappedLoad,+ wrappedStore, wrappedLoad, loadFromObject, copyStored, unsafeMapStored, collectObjects, collectStoredObjects,-- MonadStorage(..), ) where import Control.Monad.Reader@@ -20,9 +18,8 @@ import Data.Set (Set) import Data.Set qualified as S -import Erebos.Storage.Internal- import Erebos.Object.Internal+import Erebos.Storage.Internal data Stored a = Stored@@ -61,6 +58,9 @@ wrappedLoad :: Storable a => Ref -> Stored a wrappedLoad ref = Stored ref (load ref) +loadFromObject :: Storable a => Stored Object -> Stored a+loadFromObject (Stored ref obj) = Stored ref (evalLoadWithObject load' ref obj)+ copyStored :: forall m a. MonadIO m => Storage -> Stored a -> m (Stored a) copyStored st (Stored ref' x) = liftIO $ returnLoadResult . fmap (\r -> Stored r x) <$> copyRef' st ref' @@ -85,16 +85,3 @@ in ((o : xs') ++ xs, s') helper _ ( xs, s ) = ( xs, s ) collectOtherStored seen _ = ( [], seen )---class Monad m => MonadStorage m where- getStorage :: m Storage- mstore :: Storable a => a -> m (Stored a)-- default mstore :: MonadIO m => Storable a => a -> m (Stored a)- mstore x = do- st <- getStorage- wrappedStore st x--instance MonadIO m => MonadStorage (ReaderT Storage m) where- getStorage = ask
src/Erebos/Storage.hs view
@@ -14,17 +14,17 @@ HeadID, HeadTypeID, headId, headStorage, headRef, headObject, headStoredObject, loadHeads, loadHead, reloadHead,- storeHead, replaceHead, updateHead, updateHead_,+ storeHead, replaceHead, updateHead, updateHead', updateHead_, WatchedHead, watchHead, watchHeadWith, unwatchHead, watchHeadRaw, - MonadStorage(..),+ MonadStorage(..), mloadKey, ) where import Erebos.Object.Internal-import Erebos.Storable.Internal import Erebos.Storage.Disk import Erebos.Storage.Head import Erebos.Storage.Memory+import Erebos.Storage.Monad
src/Erebos/Storage/Graph.hs view
@@ -4,6 +4,7 @@ compareGeneration, generationMax, storedGeneration, + StoredTips, generations, generationsBy, ancestors, precedes,
src/Erebos/Storage/Head.hs view
@@ -9,11 +9,13 @@ -- * Head type and accessors Head, HeadType(..), HeadID, HeadTypeID, mkHeadTypeID,- headId, headStorage, headRef, headObject, headStoredObject,+ headId, headStorage,+ headRef, headObject, headStoredObject,+ headCache, -- * Loading and storing heads loadHeads, loadHead, reloadHead,- storeHead, replaceHead, updateHead, updateHead_,+ storeHead, replaceHead, updateHead, updateHead', updateHead_, loadHeadRaw, storeHeadRaw, replaceHeadRaw, -- * Watching heads@@ -27,12 +29,14 @@ import Control.Monad.Reader import Data.Bifunctor+import Data.Kind import Data.Typeable import Erebos.Object.Internal import Erebos.Storable.Internal import Erebos.Storage.Backend import Erebos.Storage.Internal+import Erebos.Storage.Monad import Erebos.UUID qualified as U @@ -42,9 +46,16 @@ -- Each possible head type has associated unique ID, represented as -- `HeadTypeID'. For each type, there can be multiple individual heads in given -- storage, each also identified by unique ID (`HeadID').-data Head a = Head HeadID (Stored a)- deriving (Eq, Show)+data Head a = Head HeadID (Stored a) !(HeadCacheType a) +instance Eq (Head a) where+ h == h' = headId h == headId h' &&+ headStorage h == headStorage h' &&+ headStoredObject h == headStoredObject h'++instance Show (Head a) where+ show h = "Head " <> show (headId h) <> " " <> show (storedRef $ headStoredObject h)+ -- | Instances of this class can be used as objects pointed to by heads in -- Erebos storage. Each such type must be `Storable' and have a unique ID. --@@ -57,13 +68,28 @@ headTypeID :: proxy a -> HeadTypeID -- ^ Get the ID of the given head type; must be unique for each `HeadType' instance. + type HeadCacheType a :: Type+ type HeadCacheType a = ()++ headCacheInit :: Stored a -> HeadCacheType a+ default headCacheInit :: HeadCacheType a ~ () => Stored a -> HeadCacheType a+ headCacheInit _ = ()++ headCacheUpdate :: Stored a -> HeadCacheType a -> HeadCacheType a+ default headCacheUpdate :: HeadCacheType a ~ () => Stored a -> HeadCacheType a -> HeadCacheType a+ headCacheUpdate _ _ = ()+ instance MonadIO m => MonadStorage (ReaderT (Head a) m) where getStorage = asks $ headStorage +initHead :: HeadType a => HeadID -> Stored a -> Head a+initHead hid obj = Head hid obj (headCacheInit obj)++ -- | Get `HeadID' associated with given `Head'. headId :: Head a -> HeadID-headId (Head uuid _) = uuid+headId (Head uuid _ _) = uuid -- | Get storage from which the `Head' was loaded. headStorage :: Head a -> Storage@@ -71,16 +97,19 @@ -- | Get `Ref' of the `Head'\'s associated object. headRef :: Head a -> Ref-headRef (Head _ sx) = storedRef sx+headRef (Head _ sx _) = storedRef sx -- | Get the object the `Head' pointed to when it was loaded. headObject :: Head a -> a-headObject (Head _ sx) = fromStored sx+headObject (Head _ sx _) = fromStored sx -- | Get the object the `Head' pointed to when it was loaded as a `Stored' value. headStoredObject :: Head a -> Stored a-headStoredObject (Head _ sx) = sx+headStoredObject (Head _ sx _) = sx +headCache :: Head a -> HeadCacheType a+headCache (Head _ _ cache) = cache+ -- | Create `HeadTypeID' from string representation of UUID. mkHeadTypeID :: String -> HeadTypeID mkHeadTypeID = maybe (error "Invalid head type ID") HeadTypeID . U.fromString@@ -89,7 +118,7 @@ -- | Load all `Head's of type @a@ from storage. loadHeads :: forall a m. MonadIO m => HeadType a => Storage -> m [Head a] loadHeads st@Storage {..} =- map (uncurry Head . fmap (wrappedLoad . Ref st))+ map (uncurry initHead . fmap (wrappedLoad . Ref st)) <$> liftIO (backendLoadHeads stBackend (headTypeID @a Proxy)) -- | Try to load a `Head' of type @a@ from storage.@@ -98,7 +127,7 @@ => Storage -- ^ Storage from which to load the head -> HeadID -- ^ ID of the particular head -> m (Maybe (Head a)) -- ^ Head object, or `Nothing' if not found-loadHead st hid = fmap (Head hid . wrappedLoad) <$> loadHeadRaw st (headTypeID @a Proxy) hid+loadHead st hid = fmap (initHead hid . wrappedLoad) <$> loadHeadRaw st (headTypeID @a Proxy) hid -- | Try to load `Head' using a raw head and type IDs, getting `Ref' if found. loadHeadRaw@@ -112,8 +141,15 @@ -- | Reload the given head from storage, returning `Head' with updated object, -- or `Nothing' if there is no longer head with the particular ID in storage.-reloadHead :: (HeadType a, MonadIO m) => Head a -> m (Maybe (Head a))-reloadHead (Head hid val) = loadHead (storedStorage val) hid+reloadHead :: forall a m. (HeadType a, MonadIO m) => Head a -> m (Maybe (Head a))+reloadHead orig@(Head hid val cache) = do+ loadHeadRaw (storedStorage val) (headTypeID @a Proxy) hid >>= \case+ Just ref+ | ref == storedRef val -> return $ Just orig+ | otherwise -> do+ let obj = wrappedLoad ref+ return $ Just $ Head hid obj (headCacheUpdate obj cache)+ Nothing -> return Nothing -- | Store a new `Head' of type 'a' in the storage. storeHead :: forall a m. MonadIO m => HeadType a => Storage -> a -> m (Head a)@@ -121,7 +157,7 @@ let tid = headTypeID @a Proxy stored <- wrappedStore st obj hid <- storeHeadRaw st tid (storedRef stored)- return $ Head hid stored+ return $ initHead hid stored -- | Store a new `Head' in the storage, using the raw `HeadTypeID' and `Ref', -- the function returns the assigned `HeadID' of the new head.@@ -149,12 +185,18 @@ -- Head value was updated in storage, the new head is @h@ (which is -- the same as first parameter with associated object replaced by -- the second parameter).-replaceHead prev@(Head hid pobj) stored' = liftIO $ do+replaceHead prev@(Head hid pobj cache) stored' = liftIO $ do let st = headStorage prev tid = headTypeID @a Proxy stored <- copyStored st stored'- bimap (fmap $ Head hid . wrappedLoad) (const $ Head hid stored) <$>- replaceHeadRaw st tid hid (storedRef pobj) (storedRef stored)+ replaceHeadRaw st tid hid (storedRef pobj) (storedRef stored) >>= \case+ Left Nothing -> do+ return $ Left Nothing+ Left (Just r) -> do+ let obj = wrappedLoad r+ return $ Left $ Just $ Head hid obj (headCacheUpdate obj cache)+ Right _ -> do+ return $ Right $ Head hid stored (headCacheUpdate stored cache) -- | Try to replace existing head using raw IDs and `Ref's. replaceHeadRaw@@ -192,12 +234,28 @@ -- ^ First element contains either the new head as @`Just' h@, or -- `Nothing' in case the head no longer exists in storage. Second -- element is the value from last call to the update function.-updateHead h f = do- (o, x) <- f $ headStoredObject h+updateHead h f = updateHead' h (f . headStoredObject)++-- | Update existing existing `Head' of type @a@ in the storage, using a given+-- function. The update function may be called multiple times in case the head+-- content changes concurrently during evaluation.+updateHead'+ :: (HeadType a, MonadIO m)+ => Head a -- ^ Existing head to be updated+ -> (Head a -> m ( Stored a, b ))+ -- ^ Function that gets current value of the head and returns updated+ -- value, along with a custom extra value to be returned from+ -- `updateHead' call. The function may be called multiple times.+ -> m ( Maybe (Head a), b )+ -- ^ First element contains either the new head as @`Just' h@, or+ -- `Nothing' in case the head no longer exists in storage. Second+ -- element is the value from last call to the update function.+updateHead' h f = do+ (o, x) <- f h replaceHead h o >>= \case Right h' -> return (Just h', x) Left Nothing -> return (Nothing, x)- Left (Just h') -> updateHead h' f+ Left (Just h') -> updateHead' h' f -- | Update existing existing `Head' of type @a@ in the storage, using a given -- function. The update function may be called multiple times in case the head@@ -232,8 +290,27 @@ -> (Head a -> b) -- ^ Selector function -> (b -> IO ()) -- ^ Callback -> IO WatchedHead -- ^ Watched head handle-watchHeadWith (Head hid val) sel cb = do- watchHeadRaw (storedStorage val) (headTypeID @a Proxy) hid (sel . Head hid . wrappedLoad) cb+watchHeadWith orig@(Head hid val _) sel cb = do+ let st = storedStorage val++ memo <- newEmptyMVar+ let cb' dgst = do+ let obj = wrappedLoad $ Ref st dgst+ modifyMVar_ memo $ \( cache, prev ) -> do+ let cache' = maybe (headCacheInit obj) (headCacheUpdate obj) cache+ let x = sel $ Head hid obj cache'+ when (Just x /= prev) $ cb x+ return ( Just cache', Just x )+ wid <- case storedStorage val of+ Storage {..} -> backendWatchHead stBackend (headTypeID @a Proxy) hid cb'++ h <- reloadHead orig+ let cur = sel <$> h+ maybe (return ()) cb cur+ putMVar memo ( headCache <$> h, cur )++ return $ WatchedHead st wid memo+ -- | Watch the given head using raw IDs and a selector from `Ref'. watchHeadRaw :: forall b. Eq b => Storage -> HeadTypeID -> HeadID -> (Ref -> b) -> (b -> IO ()) -> IO WatchedHead
src/Erebos/Storage/Merge.hs view
@@ -7,6 +7,7 @@ compareGeneration, generationMax, storedGeneration, + StoredTips, generations, generationsBy, ancestors, precedes,
+ src/Erebos/Storage/Monad.hs view
@@ -0,0 +1,37 @@+module Erebos.Storage.Monad (+ MonadStorage(..),+ mloadKey,+) where++import Control.Monad.Except+import Control.Monad.Reader++import Erebos.Error+import Erebos.Object.Internal+import Erebos.Storable.Internal+import Erebos.Storage.Key+++class Monad m => MonadStorage m where+ getStorage :: m Storage+ mstore :: Storable a => a -> m (Stored a)++ default mstore :: MonadIO m => Storable a => a -> m (Stored a)+ mstore x = do+ st <- getStorage+ wrappedStore st x++ mstoreKey :: KeyPair sec pub => sec -> m ()+ default mstoreKey :: (KeyPair sec pub, MonadIO m) => sec -> m ()+ mstoreKey = liftIO . storeKey++ mloadKeyMb :: KeyPair sec pub => Stored pub -> m (Maybe sec)+ default mloadKeyMb :: (KeyPair sec pub, MonadIO m) => Stored pub -> m (Maybe sec)+ mloadKeyMb = loadKeyMb++mloadKey :: (KeyPair sec pub, MonadStorage m, MonadError e m, FromErebosError e) => Stored pub -> m sec+mloadKey pub = maybe (throwOtherError $ "secret key not found for " <> show (storedRef pub)) return =<< mloadKeyMb pub+++instance MonadIO m => MonadStorage (ReaderT Storage m) where+ getStorage = ask
src/Erebos/Sync.hs view
@@ -32,7 +32,7 @@ serviceNewPeer = notifyPeer . lsShared . fromStored =<< svcGetLocal serviceUpdatedPeer = serviceNewPeer- serviceStorageWatchers _ = (:[]) $ SomeStorageWatcher (lsShared . fromStored) notifyPeer+ serviceStorageWatchers _ = (:[]) $ someStorageWatcher (lsShared . fromStored) notifyPeer instance Storable SyncService where store' (SyncPacket smsg) = store' smsg
+ src/Service.hs view
@@ -0,0 +1,218 @@+module Service (+ Service(..),+ SomeService(..), someService, someServiceAttr, someServiceID,+ SomeServiceState(..), fromServiceState, someServiceEmptyState,+ SomeServiceGlobalState(..), fromServiceGlobalState, someServiceEmptyGlobalState,+ SomeStorageWatcher(..),+ ServiceID, mkServiceID,++ ServiceHandler,+ ServiceInput(..),+ ServiceReply(..),+ runServiceHandler,++ svcGet, svcSet, svcModify,+ svcGetGlobal, svcSetGlobal, svcModifyGlobal,+ svcGetLocal, svcSetLocal,++ svcSelf,+ svcPrint,++ replyPacket, replyStored, replyStoredRef,+ afterCommit,+) where++import Control.Monad.Except+import Control.Monad.Reader+import Control.Monad.State+import Control.Monad.Writer++import Data.Kind+import Data.Typeable++import Erebos.Identity+import {-# SOURCE #-} Erebos.Network+import Erebos.Network.Protocol+import Erebos.State+import Erebos.Storable+import Erebos.Storage.Head+import Erebos.UUID qualified as U++class (+ Typeable s, Storable s,+ Typeable (ServiceAttributes s),+ Typeable (ServiceState s),+ Typeable (ServiceGlobalState s)+ ) => Service s where++ serviceID :: proxy s -> ServiceID+ serviceHandler :: Stored s -> ServiceHandler s ()++ serviceNewPeer :: ServiceHandler s ()+ serviceNewPeer = return ()++ serviceUpdatedPeer :: ServiceHandler s ()+ serviceUpdatedPeer = return ()++ type ServiceAttributes s = attr | attr -> s+ type ServiceAttributes s = Proxy s+ defaultServiceAttributes :: proxy s -> ServiceAttributes s+ default defaultServiceAttributes :: ServiceAttributes s ~ Proxy s => proxy s -> ServiceAttributes s+ defaultServiceAttributes _ = Proxy++ type ServiceState s :: Type+ type ServiceState s = ()+ emptyServiceState :: proxy s -> ServiceState s+ default emptyServiceState :: ServiceState s ~ () => proxy s -> ServiceState s+ emptyServiceState _ = ()++ type ServiceGlobalState s :: Type+ type ServiceGlobalState s = ()+ emptyServiceGlobalState :: proxy s -> ServiceGlobalState s+ default emptyServiceGlobalState :: ServiceGlobalState s ~ () => proxy s -> ServiceGlobalState s+ emptyServiceGlobalState _ = ()++ serviceStorageWatchers :: proxy s -> [SomeStorageWatcher s]+ serviceStorageWatchers _ = []++ serviceStopServer :: proxy s -> Server -> ServiceGlobalState s -> [ ( Peer, ServiceState s ) ] -> IO ()+ serviceStopServer _ _ _ _ = return ()+++data SomeService = forall s. Service s => SomeService (Proxy s) (ServiceAttributes s)++someService :: forall s proxy. Service s => proxy s -> SomeService+someService _ = SomeService @s Proxy (defaultServiceAttributes @s Proxy)++someServiceAttr :: forall s. Service s => ServiceAttributes s -> SomeService+someServiceAttr attr = SomeService @s Proxy attr++someServiceID :: SomeService -> ServiceID+someServiceID (SomeService s _) = serviceID s++data SomeServiceState = forall s. Service s => SomeServiceState (Proxy s) (ServiceState s)++fromServiceState :: Service s => proxy s -> SomeServiceState -> Maybe (ServiceState s)+fromServiceState _ (SomeServiceState _ s) = cast s++someServiceEmptyState :: SomeService -> SomeServiceState+someServiceEmptyState (SomeService p _) = SomeServiceState p (emptyServiceState p)++data SomeServiceGlobalState = forall s. Service s => SomeServiceGlobalState (Proxy s) (ServiceGlobalState s)++fromServiceGlobalState :: Service s => proxy s -> SomeServiceGlobalState -> Maybe (ServiceGlobalState s)+fromServiceGlobalState _ (SomeServiceGlobalState _ s) = cast s++someServiceEmptyGlobalState :: SomeService -> SomeServiceGlobalState+someServiceEmptyGlobalState (SomeService p _) = SomeServiceGlobalState p (emptyServiceGlobalState p)+++data SomeStorageWatcher s+ = forall a. Eq a => SomeStorageWatcher (Stored LocalState -> a) (a -> ServiceHandler s ())+ | forall a. Eq a => SomeStorageWatcherHC (Stored LocalState -> HeadCacheType LocalState -> a) (a -> ServiceHandler s ())+ | forall a. Eq a => GlobalStorageWatcher (Stored LocalState -> a) (Server -> a -> ExceptT ErebosError IO ())+ | forall a. Eq a => GlobalStorageWatcherH (Head LocalState -> a) (Server -> a -> ExceptT ErebosError IO ())+++mkServiceID :: String -> ServiceID+mkServiceID = maybe (error "Invalid service ID") ServiceID . U.fromString++data ServiceInput s = ServiceInput+ { svcAttributes :: ServiceAttributes s+ , svcPeer :: Peer+ , svcPeerAddress :: PeerAddress+ , svcPeerIdentity :: UnifiedIdentity+ , svcServer :: Server+ , svcPrintOp :: String -> IO ()+ , svcNewStreams :: [ RawStreamReader ]+ }++data ServiceReply s+ = ServiceReply (Either s (Stored s)) Bool+ | ServiceOpenStream (RawStreamWriter -> IO ())+ | ServiceFinally (IO ())++data ServiceHandlerState s = ServiceHandlerState+ { svcValue :: ServiceState s+ , svcGlobal :: ServiceGlobalState s+ , svcLocal :: Stored LocalState+ , svcLocalCache :: HeadCacheType LocalState+ }++newtype ServiceHandler s a = ServiceHandler (ReaderT (ServiceInput s) (WriterT [ServiceReply s] (StateT (ServiceHandlerState s) (ExceptT ErebosError IO))) a)+ deriving (Functor, Applicative, Monad, MonadReader (ServiceInput s), MonadWriter [ServiceReply s], MonadState (ServiceHandlerState s), MonadError ErebosError, MonadIO)++instance MonadStorage (ServiceHandler s) where+ getStorage = asks $ peerStorage . svcPeer++instance MonadHead LocalState (ServiceHandler s) where+ updateLocalHead f = do+ ls <- gets svcLocal+ c <- gets svcLocalCache+ ( ls', x ) <- f ls+ let c' = headCacheUpdate ls' c+ modify $ \s -> s { svcLocal = ls', svcLocalCache = c' }+ return x+ getLocalHeadCache _ = gets svcLocalCache++runServiceHandler+ :: Service s+ => Head LocalState -> ServiceInput s -> ServiceState s -> ServiceGlobalState s+ -> ServiceHandler s ()+ -> IO ( [ ServiceReply s], ( ServiceState s, ServiceGlobalState s, Head LocalState ) )+runServiceHandler h input svc global shandler = do+ let sstate = ServiceHandlerState { svcValue = svc, svcGlobal = global, svcLocal = headStoredObject h, svcLocalCache = headCache h }+ ServiceHandler handler = shandler+ (runExceptT $ flip runStateT sstate $ execWriterT $ flip runReaderT input $ handler) >>= \case+ Left err -> do+ svcPrintOp input $ "service failed: " ++ showErebosError err+ return ( [], ( svc, global, h ) )+ Right (rsp, sstate')+ | svcLocal sstate' == svcLocal sstate -> return ( rsp, ( svcValue sstate', svcGlobal sstate', h ))+ | otherwise -> replaceHead h (svcLocal sstate') >>= \case+ Left (Just h') -> runServiceHandler h' input svc global shandler+ Left Nothing -> return ( rsp, ( svcValue sstate', svcGlobal sstate', h ) )+ Right h' -> return ( rsp, ( svcValue sstate', svcGlobal sstate', h' ) )++svcGet :: ServiceHandler s (ServiceState s)+svcGet = gets svcValue++svcSet :: ServiceState s -> ServiceHandler s ()+svcSet x = modify $ \st -> st { svcValue = x }++svcModify :: (ServiceState s -> ServiceState s) -> ServiceHandler s ()+svcModify f = modify $ \st -> st { svcValue = f (svcValue st) }++svcGetGlobal :: ServiceHandler s (ServiceGlobalState s)+svcGetGlobal = gets svcGlobal++svcSetGlobal :: ServiceGlobalState s -> ServiceHandler s ()+svcSetGlobal x = modify $ \st -> st { svcGlobal = x }++svcModifyGlobal :: (ServiceGlobalState s -> ServiceGlobalState s) -> ServiceHandler s ()+svcModifyGlobal f = modify $ \st -> st { svcGlobal = f (svcGlobal st) }++svcGetLocal :: ServiceHandler s (Stored LocalState)+svcGetLocal = gets svcLocal++svcSetLocal :: Stored LocalState -> ServiceHandler s ()+svcSetLocal x = modify $ \st -> st { svcLocal = x }++svcSelf :: ServiceHandler s UnifiedIdentity+svcSelf = maybe (throwOtherError "failed to validate own identity") return .+ validateExtendedIdentity . lsIdentity . fromStored =<< svcGetLocal++svcPrint :: String -> ServiceHandler s ()+svcPrint str = afterCommit . ($ str) =<< asks svcPrintOp++replyPacket :: Service s => s -> ServiceHandler s ()+replyPacket x = tell [ServiceReply (Left x) True]++replyStored :: Service s => Stored s -> ServiceHandler s ()+replyStored x = tell [ServiceReply (Right x) True]++replyStoredRef :: Service s => Stored s -> ServiceHandler s ()+replyStoredRef x = tell [ServiceReply (Right x) False]++afterCommit :: IO () -> ServiceHandler s ()+afterCommit x = tell [ServiceFinally x]